diff --git a/lib/core/extensions/rpc/json_rpc_payload_limits.dart b/lib/core/extensions/rpc/json_rpc_payload_limits.dart index 274dcd7..3095e5c 100644 --- a/lib/core/extensions/rpc/json_rpc_payload_limits.dart +++ b/lib/core/extensions/rpc/json_rpc_payload_limits.dart @@ -11,6 +11,14 @@ const int kDefaultJsonRpcMaxLineBytes = 32 * 1024 * 1024; /// Lines above this UTF-8 length are `jsonDecode`d off the UI isolate. const int kJsonRpcOffIsolateDecodeThresholdBytes = 64 * 1024; +/// Lines waiting to be handled before the plugin's stdout is paused. +const int kDefaultJsonRpcMaxBufferedLines = 256; + +/// Total size of lines waiting to be handled before stdout is paused. A single +/// larger line is still accepted (up to the line limit); it just pauses reading +/// until it has been handled. +const int kDefaultJsonRpcMaxBufferedBytes = 16 * 1024 * 1024; + /// Thrown when a plugin emits a newline-delimited JSON line larger than the /// configured maximum. class JsonRpcPayloadTooLargeException implements Exception { @@ -108,6 +116,10 @@ class _BoundedUtf8LineSplitter cancelOnError: true, ); + // Pausing the consumer must pause the byte source too; otherwise the + // controller buffers every decoded line while the plugin keeps writing. + controller.onPause = sub.pause; + controller.onResume = sub.resume; controller.onCancel = () => sub.cancel(); return controller.stream; } diff --git a/lib/core/extensions/rpc/json_rpc_stdio_client.dart b/lib/core/extensions/rpc/json_rpc_stdio_client.dart index d1cb953..df66eaa 100644 --- a/lib/core/extensions/rpc/json_rpc_stdio_client.dart +++ b/lib/core/extensions/rpc/json_rpc_stdio_client.dart @@ -18,6 +18,8 @@ class JsonRpcStdioClient { required IOSink stdin, this.requestTimeout = const Duration(seconds: 10), this.maxLineBytes = kDefaultJsonRpcMaxLineBytes, + this.maxBufferedLines = kDefaultJsonRpcMaxBufferedLines, + this.maxBufferedBytes = kDefaultJsonRpcMaxBufferedBytes, }) : _stdin = stdin, _lines = stdout.transform( boundedUtf8LineSplitter(maxLineBytes: maxLineBytes), @@ -35,6 +37,13 @@ class JsonRpcStdioClient { final Duration requestTimeout; final int maxLineBytes; + /// Backpressure: while more than this many lines (or [maxBufferedBytes]) are + /// waiting to be handled, the plugin's stdout is paused, so a driver that + /// streams faster than we decode fills the OS pipe instead of our memory. + /// Reading resumes once the backlog drops to half of either limit. + final int maxBufferedLines; + final int maxBufferedBytes; + final Map> _pending = {}; var _nextId = 1; var _closed = false; @@ -50,6 +59,16 @@ class JsonRpcStdioClient { /// Serializes async line handling so large-line isolate decode stays ordered. Future _lineChain = Future.value(); + var _bufferedLines = 0; + var _bufferedBytes = 0; + var _readPaused = false; + + /// Lines received but not yet handled (the backlog backpressure bounds). + int get bufferedLineCount => _bufferedLines; + + /// True while stdout reading is paused because the backlog is too large. + bool get isReadPaused => _readPaused; + /// Sends a JSON-RPC request and waits for the matching response. Future sendRequest( String method, [ @@ -116,7 +135,24 @@ class JsonRpcStdioClient { } void _onLine(String line) { - _lineChain = _lineChain.then((_) => _handleLine(line)); + _bufferedLines++; + _bufferedBytes += line.length; + if (!_readPaused && + (_bufferedLines > maxBufferedLines || + _bufferedBytes > maxBufferedBytes)) { + _readPaused = true; + _subscription?.pause(); + } + _lineChain = _lineChain.then((_) => _handleLine(line)).whenComplete(() { + _bufferedLines--; + _bufferedBytes -= line.length; + if (_readPaused && + _bufferedLines <= maxBufferedLines ~/ 2 && + _bufferedBytes <= maxBufferedBytes ~/ 2) { + _readPaused = false; + _subscription?.resume(); + } + }); } Future _handleLine(String line) async { @@ -162,12 +198,17 @@ class JsonRpcStdioClient { void _onDone() { _fatalError ??= StateError('Plugin stdout closed'); - for (final pending in _pending.values) { - if (!pending.isCompleted) { - pending.completeError(_fatalError!); + // Replies already received may still be queued (a large line is decoded in + // an isolate). Fail what is left only after the backlog has been handled, + // so a plugin that answers and exits does not lose its last responses. + _lineChain = _lineChain.whenComplete(() { + for (final pending in _pending.values) { + if (!pending.isCompleted) { + pending.completeError(_fatalError!); + } } - } - _pending.clear(); + _pending.clear(); + }); } } diff --git a/test/core/extensions/rpc/json_rpc_stdio_client_test.dart b/test/core/extensions/rpc/json_rpc_stdio_client_test.dart index b8d454e..107f739 100644 --- a/test/core/extensions/rpc/json_rpc_stdio_client_test.dart +++ b/test/core/extensions/rpc/json_rpc_stdio_client_test.dart @@ -1,6 +1,7 @@ import 'dart:async'; import 'dart:convert'; import 'dart:io'; +import 'dart:math' as math; import 'package:flutter_test/flutter_test.dart'; import 'package:querya_desktop/core/extensions/rpc/json_rpc_payload_limits.dart'; @@ -32,6 +33,208 @@ void main() { }); }); + /// IOSink allows one flush at a time, so requests are issued one per turn. + Future>> sendRequests( + JsonRpcStdioClient client, + int count, + ) async { + final replies = >[]; + for (var i = 1; i <= count; i++) { + replies.add(client.sendRequest('q$i')); + await Future.delayed(Duration.zero); + } + return replies; + } + + group('boundedUtf8LineSplitter backpressure', () { + test('pausing the consumer pauses the byte source and resume resumes it', + () async { + var paused = 0; + var resumed = 0; + final source = StreamController>( + onPause: () => paused++, + onResume: () => resumed++, + ); + final lines = []; + final sub = source.stream + .transform(boundedUtf8LineSplitter(maxLineBytes: 1024)) + .listen(lines.add); + + sub.pause(); + await Future.delayed(Duration.zero); + expect(paused, 1); + expect(source.isPaused, isTrue); + + sub.resume(); + await Future.delayed(Duration.zero); + expect(resumed, 1); + expect(source.isPaused, isFalse); + + source.add(utf8.encode('after\n')); + await Future.delayed(Duration.zero); + expect(lines, ['after']); + await sub.cancel(); + await source.close(); + }); + }); + + group('JsonRpcStdioClient backpressure', () { + // A reply bigger than kJsonRpcOffIsolateDecodeThresholdBytes is decoded in + // an isolate, so handling is slower than arrival. + String bigReply(int id) => jsonEncode({ + 'jsonrpc': '2.0', + 'id': id, + 'result': {'blob': 'x' * (kJsonRpcOffIsolateDecodeThresholdBytes + 1024)}, + }); + + test('pauses stdout while a backlog is waiting and resumes after it drains', + () async { + const requests = 24; + const limit = 4; + final stdout = StreamController>(); + final stdin = StreamController>(); + final client = JsonRpcStdioClient( + stdout: stdout.stream, + stdin: IOSink(stdin.sink), + maxBufferedLines: limit, + requestTimeout: const Duration(seconds: 30), + ); + stdin.stream.listen((_) {}); + + final replies = await sendRequests(client, requests); + + var sawPaused = false; + var peak = 0; + final watcher = Timer.periodic(const Duration(milliseconds: 1), (_) { + sawPaused = sawPaused || client.isReadPaused; + peak = math.max(peak, client.bufferedLineCount); + }); + + for (var id = 1; id <= requests; id++) { + stdout.add(utf8.encode('${bigReply(id)}\n')); + } + final results = await Future.wait(replies); + watcher.cancel(); + + expect(results, hasLength(requests)); + expect(sawPaused, isTrue, reason: 'stdout must be paused under backlog'); + // Lines already handed over when the pause lands may still arrive. + expect(peak, lessThanOrEqualTo(limit + 2)); + expect(client.bufferedLineCount, 0); + expect(client.isReadPaused, isFalse); + expect(stdout.isPaused, isFalse); + + await client.close(); + await stdout.close(); + await stdin.close(); + }); + + test('a fast producer cannot push the backlog past the limit', () async { + const requests = 60; + const limit = 3; + var produced = 0; + final start = Completer(); + // An async* generator honours pause, like a real pipe would. It starts + // only once every request is registered. + Stream> plugin() async* { + await start.future; + for (var id = 1; id <= requests; id++) { + produced++; + yield utf8.encode('${bigReply(id)}\n'); + } + } + + final stdin = StreamController>(); + stdin.stream.listen((_) {}); + final client = JsonRpcStdioClient( + stdout: plugin(), + stdin: IOSink(stdin.sink), + maxBufferedLines: limit, + requestTimeout: const Duration(seconds: 30), + ); + + // Register the requests the plugin answers; ids are assigned in order. + final replies = await sendRequests(client, requests); + start.complete(); + + var peakAhead = 0; + var handled = 0; + for (final r in replies) { + r.then((_) => handled++); + } + final watcher = Timer.periodic(const Duration(milliseconds: 1), (_) { + peakAhead = math.max(peakAhead, produced - handled); + }); + await Future.wait(replies); + watcher.cancel(); + + expect(produced, requests); + // Produced-but-unhandled lines stay near the limit, not near 60. + expect(peakAhead, lessThanOrEqualTo(limit + 6)); + + await client.close(); + await stdin.close(); + }); + + test('replies queued when stdout closes are still delivered', () async { + final stdout = StreamController>(); + final stdin = StreamController>(); + stdin.stream.listen((_) {}); + final client = JsonRpcStdioClient( + stdout: stdout.stream, + stdin: IOSink(stdin.sink), + requestTimeout: const Duration(seconds: 30), + ); + final replies = await sendRequests(client, 3); + + // Large replies are decoded off-isolate, then the plugin exits at once. + for (var id = 1; id <= 3; id++) { + stdout.add(utf8.encode('${bigReply(id)}\n')); + } + await stdout.close(); + + final results = await Future.wait(replies); + expect(results, hasLength(3)); + expect((results.first as Map)['blob'], isA()); + + // Whatever is still unanswered after the close fails as before. + await expectLater( + client.sendRequest('late'), + throwsA(isA()), + ); + await client.close(); + await stdin.close(); + }); + + test('backlog is also bounded by bytes', () async { + final stdout = StreamController>(); + final stdin = StreamController>(); + stdin.stream.listen((_) {}); + final client = JsonRpcStdioClient( + stdout: stdout.stream, + stdin: IOSink(stdin.sink), + maxBufferedBytes: 1024, + requestTimeout: const Duration(seconds: 30), + ); + final replies = await sendRequests(client, 6); + + var sawPaused = false; + final watcher = Timer.periodic(const Duration(milliseconds: 1), (_) { + sawPaused = sawPaused || client.isReadPaused; + }); + for (var id = 1; id <= 6; id++) { + stdout.add(utf8.encode('${bigReply(id)}\n')); + } + await Future.wait(replies); + watcher.cancel(); + expect(sawPaused, isTrue); + + await client.close(); + await stdout.close(); + await stdin.close(); + }); + }); + group('JsonRpcStdioClient payload bounds', () { test('completes pending request with payload-too-large error', () async { final stdout = StreamController>();