Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions lib/core/extensions/rpc/json_rpc_payload_limits.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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;
}
Expand Down
53 changes: 47 additions & 6 deletions lib/core/extensions/rpc/json_rpc_stdio_client.dart
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand All @@ -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<int, Completer<Object?>> _pending = {};
var _nextId = 1;
var _closed = false;
Expand All @@ -50,6 +59,16 @@ class JsonRpcStdioClient {
/// Serializes async line handling so large-line isolate decode stays ordered.
Future<void> _lineChain = Future<void>.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<Object?> sendRequest(
String method, [
Expand Down Expand Up @@ -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<void> _handleLine(String line) async {
Expand Down Expand Up @@ -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();
});
}
}

Expand Down
203 changes: 203 additions & 0 deletions test/core/extensions/rpc/json_rpc_stdio_client_test.dart
Original file line number Diff line number Diff line change
@@ -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';
Expand Down Expand Up @@ -32,6 +33,208 @@ void main() {
});
});

/// IOSink allows one flush at a time, so requests are issued one per turn.
Future<List<Future<Object?>>> sendRequests(
JsonRpcStdioClient client,
int count,
) async {
final replies = <Future<Object?>>[];
for (var i = 1; i <= count; i++) {
replies.add(client.sendRequest('q$i'));
await Future<void>.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<List<int>>(
onPause: () => paused++,
onResume: () => resumed++,
);
final lines = <String>[];
final sub = source.stream
.transform(boundedUtf8LineSplitter(maxLineBytes: 1024))
.listen(lines.add);

sub.pause();
await Future<void>.delayed(Duration.zero);
expect(paused, 1);
expect(source.isPaused, isTrue);

sub.resume();
await Future<void>.delayed(Duration.zero);
expect(resumed, 1);
expect(source.isPaused, isFalse);

source.add(utf8.encode('after\n'));
await Future<void>.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<List<int>>();
final stdin = StreamController<List<int>>();
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<void>();
// An async* generator honours pause, like a real pipe would. It starts
// only once every request is registered.
Stream<List<int>> plugin() async* {
await start.future;
for (var id = 1; id <= requests; id++) {
produced++;
yield utf8.encode('${bigReply(id)}\n');
}
}

final stdin = StreamController<List<int>>();
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<List<int>>();
final stdin = StreamController<List<int>>();
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<String>());

// Whatever is still unanswered after the close fails as before.
await expectLater(
client.sendRequest('late'),
throwsA(isA<StateError>()),
);
await client.close();
await stdin.close();
});

test('backlog is also bounded by bytes', () async {
final stdout = StreamController<List<int>>();
final stdin = StreamController<List<int>>();
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<List<int>>();
Expand Down
Loading