Skip to content

Commit d9b9207

Browse files
Merge pull request #434 from QueryaHub/issue/419-rpc-payload-bounds
perf(extensions): bound JSON-RPC payload size (#419)
2 parents ce4af04 + bbc2777 commit d9b9207

6 files changed

Lines changed: 265 additions & 13 deletions

File tree

‎docs/tz-block-c-rpc-bridge.md‎

Lines changed: 12 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -46,7 +46,15 @@ final result = await rpcClient.sendRequest('db.connect', credentialsMap);
4646

4747
---
4848

49-
## 4. Контракт Ошибок (Error Mapping)
50-
Плагин должен возвращать ошибки согласно спецификации JSON-RPC. RPC Bridge должен уметь парсить эти ошибки и превращать их в понятные Dart-exceptions:
51-
- Ошибка подключения (Timeout, Wrong Password) -> Показывается в UI в красном Snackbar.
52-
- Синтаксическая ошибка SQL -> Выделяется красным в SQL редакторе.
49+
## 5. Лимиты полезной нагрузки (NDJSON)
50+
51+
Каждый ответ — **одна JSON-строка** на `stdout` (newline-delimited). Хост (`JsonRpcStdioClient`) применяет:
52+
53+
| Лимит | Значение по умолчанию | Поведение |
54+
|-------|----------------------|-----------|
55+
| Макс. длина одной строки ответа | **32 MiB** UTF-8 | Fail closed: `JsonRpcPayloadTooLargeException`, все pending RPC завершаются ошибкой |
56+
| Декод больших строк | **> 64 KiB** | `jsonDecode` уходит в isolate |
57+
58+
Для `db.query` хост всегда передаёт `params.limit` (Preferences → Max rows in results), чтобы драйвер обрезал результат **до** сериализации. Драйверы обязаны уважать `limit`.
59+
60+
Чанкованный / бинарный framing для очень больших выборок — follow-up; до него bound + `limit` обязательны.

‎lib/core/extensions/extension_driver_session.dart‎

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import 'package:querya_desktop/core/extensions/rpc/plugin_rpc_bridge.dart';
1414
import 'package:querya_desktop/core/extensions/sandbox/sandbox_os_isolation.dart';
1515
import 'package:querya_desktop/core/extensions/sandbox/unsandboxed_launch_consent_gate.dart';
1616
import 'package:querya_desktop/core/sdui/sdui_tree_schema.dart';
17+
import 'package:querya_desktop/core/storage/app_settings.dart';
1718
import 'package:querya_desktop/core/storage/connection_secrets_store.dart';
1819
import 'package:querya_desktop/core/storage/local_db.dart';
1920

@@ -276,16 +277,21 @@ class ExtensionDriverSession {
276277
}
277278

278279
/// Executes SQL through the plugin (`db.query`) and returns the raw result.
280+
///
281+
/// When [limit] is omitted, Preferences **Max rows in results** is sent so
282+
/// drivers can bound the NDJSON response before it hits the host.
279283
Future<ExtensionQueryResult> query(
280284
ConnectionRow row,
281285
String sql, {
282286
int? limit,
283287
}) async {
284288
final bridge = await ensureConnected(row);
289+
final effectiveLimit =
290+
limit ?? await AppSettings.instance.getSqlResultMaxRows();
285291
final result = await bridge.sendRequest('db.query', {
286292
'connectionId': row.id,
287293
'sql': sql,
288-
if (limit != null) 'limit': limit,
294+
'limit': effectiveLimit,
289295
});
290296
return compute(_parseExtensionQueryResultRpc, result);
291297
}
Lines changed: 114 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,114 @@
1+
import 'dart:async';
2+
import 'dart:convert';
3+
import 'dart:typed_data';
4+
5+
/// Default max UTF-8 byte length of one JSON-RPC stdout line (NDJSON).
6+
///
7+
/// Large `db.query` results are one object per line; without a bound the host
8+
/// can OOM before Preferences row caps apply. Drivers should honor `limit`.
9+
const int kDefaultJsonRpcMaxLineBytes = 32 * 1024 * 1024;
10+
11+
/// Lines above this UTF-8 length are `jsonDecode`d off the UI isolate.
12+
const int kJsonRpcOffIsolateDecodeThresholdBytes = 64 * 1024;
13+
14+
/// Thrown when a plugin emits a newline-delimited JSON line larger than the
15+
/// configured maximum.
16+
class JsonRpcPayloadTooLargeException implements Exception {
17+
JsonRpcPayloadTooLargeException({
18+
required this.maxLineBytes,
19+
required this.receivedBytes,
20+
});
21+
22+
final int maxLineBytes;
23+
final int receivedBytes;
24+
25+
@override
26+
String toString() =>
27+
'JsonRpcPayloadTooLargeException: JSON-RPC line is $receivedBytes bytes '
28+
'(max $maxLineBytes). Reduce result size or pass a smaller `limit`.';
29+
}
30+
31+
/// Splits a byte stream into UTF-8 lines, failing closed if any line exceeds
32+
/// [maxLineBytes] (counted before decode).
33+
StreamTransformer<List<int>, String> boundedUtf8LineSplitter({
34+
int maxLineBytes = kDefaultJsonRpcMaxLineBytes,
35+
}) {
36+
return _BoundedUtf8LineSplitter(maxLineBytes: maxLineBytes);
37+
}
38+
39+
class _BoundedUtf8LineSplitter
40+
extends StreamTransformerBase<List<int>, String> {
41+
_BoundedUtf8LineSplitter({required this.maxLineBytes});
42+
43+
final int maxLineBytes;
44+
45+
@override
46+
Stream<String> bind(Stream<List<int>> stream) {
47+
final controller = StreamController<String>(sync: true);
48+
final pending = BytesBuilder(copy: false);
49+
late final StreamSubscription<List<int>> sub;
50+
51+
void fail(Object error, [StackTrace? st]) {
52+
if (!controller.isClosed) {
53+
controller.addError(error, st);
54+
controller.close();
55+
}
56+
sub.cancel();
57+
}
58+
59+
void emitLine() {
60+
var bytes = pending.takeBytes();
61+
if (bytes.isNotEmpty && bytes.last == 0x0d) {
62+
bytes = Uint8List.sublistView(bytes, 0, bytes.length - 1);
63+
}
64+
if (bytes.isEmpty) return;
65+
try {
66+
controller.add(utf8.decode(bytes));
67+
} catch (e, st) {
68+
fail(e, st);
69+
}
70+
}
71+
72+
sub = stream.listen(
73+
(chunk) {
74+
for (var i = 0; i < chunk.length; i++) {
75+
final b = chunk[i];
76+
if (b == 0x0a) {
77+
emitLine();
78+
continue;
79+
}
80+
if (pending.length >= maxLineBytes) {
81+
fail(
82+
JsonRpcPayloadTooLargeException(
83+
maxLineBytes: maxLineBytes,
84+
receivedBytes: pending.length + 1,
85+
),
86+
);
87+
return;
88+
}
89+
pending.addByte(b);
90+
}
91+
},
92+
onError: fail,
93+
onDone: () {
94+
if (pending.length > 0) {
95+
if (pending.length > maxLineBytes) {
96+
fail(
97+
JsonRpcPayloadTooLargeException(
98+
maxLineBytes: maxLineBytes,
99+
receivedBytes: pending.length,
100+
),
101+
);
102+
return;
103+
}
104+
emitLine();
105+
}
106+
controller.close();
107+
},
108+
cancelOnError: true,
109+
);
110+
111+
controller.onCancel = () => sub.cancel();
112+
return controller.stream;
113+
}
114+
}

‎lib/core/extensions/rpc/json_rpc_stdio_client.dart‎

Lines changed: 32 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,31 +1,49 @@
11
import 'dart:async';
22
import 'dart:convert';
33
import 'dart:io';
4+
import 'dart:isolate';
5+
6+
import 'package:querya_desktop/core/extensions/rpc/json_rpc_payload_limits.dart';
47

58
/// Minimal JSON-RPC 2.0 client over newline-delimited JSON on stdio.
69
///
710
/// Enough for Block E credential injection and later Block C methods without
811
/// pulling `json_rpc_2` yet. One JSON object per line on stdin/stdout.
12+
///
13+
/// Incoming lines are bounded by [maxLineBytes] (see
14+
/// [kDefaultJsonRpcMaxLineBytes]); oversized payloads fail closed.
915
class JsonRpcStdioClient {
1016
JsonRpcStdioClient({
1117
required Stream<List<int>> stdout,
1218
required IOSink stdin,
1319
this.requestTimeout = const Duration(seconds: 10),
20+
this.maxLineBytes = kDefaultJsonRpcMaxLineBytes,
1421
}) : _stdin = stdin,
15-
_lines = utf8.decoder.bind(stdout).transform(const LineSplitter()) {
16-
_subscription = _lines.listen(_onLine, onError: _onError, onDone: _onDone);
22+
_lines = stdout.transform(
23+
boundedUtf8LineSplitter(maxLineBytes: maxLineBytes),
24+
) {
25+
_subscription = _lines.listen(
26+
_onLine,
27+
onError: _onError,
28+
onDone: _onDone,
29+
cancelOnError: false,
30+
);
1731
}
1832

1933
final IOSink _stdin;
2034
final Stream<String> _lines;
2135
final Duration requestTimeout;
36+
final int maxLineBytes;
2237

2338
final Map<int, Completer<Object?>> _pending = {};
2439
var _nextId = 1;
2540
var _closed = false;
2641
StreamSubscription<String>? _subscription;
2742
Object? _fatalError;
2843

44+
/// Serializes async line handling so large-line isolate decode stays ordered.
45+
Future<void> _lineChain = Future<void>.value();
46+
2947
/// Sends a JSON-RPC request and waits for the matching response.
3048
Future<Object?> sendRequest(
3149
String method, [
@@ -82,12 +100,21 @@ class JsonRpcStdioClient {
82100
}
83101

84102
void _onLine(String line) {
103+
_lineChain = _lineChain.then((_) => _handleLine(line));
104+
}
105+
106+
Future<void> _handleLine(String line) async {
85107
if (line.trim().isEmpty) return;
86108
late final Map<String, dynamic> message;
87109
try {
88-
final decoded = jsonDecode(line);
89-
if (decoded is! Map<String, dynamic>) return;
90-
message = decoded;
110+
final Object decoded;
111+
if (line.length > kJsonRpcOffIsolateDecodeThresholdBytes) {
112+
decoded = await Isolate.run(() => jsonDecode(line));
113+
} else {
114+
decoded = jsonDecode(line);
115+
}
116+
if (decoded is! Map) return;
117+
message = Map<String, dynamic>.from(decoded);
91118
} catch (_) {
92119
return;
93120
}

‎lib/features/extensions/extension_sql_workspace.dart‎

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -50,6 +50,7 @@ class _ExtensionSqlWorkspaceState
5050
String? _statusLine;
5151

5252
int _historyMaxEntries = kDefaultSqlHistoryMaxEntries;
53+
int _resultMaxRows = kDefaultSqlResultMaxRows;
5354
double _editorFontSize = kDefaultSqlEditorFontSize;
5455

5556
static const _previewRowLimit = 200;
@@ -89,10 +90,12 @@ class _ExtensionSqlWorkspaceState
8990

9091
Future<void> _loadWorkspaceSettings() async {
9192
final hist = await AppSettings.instance.getSqlHistoryMaxEntries();
93+
final rows = await AppSettings.instance.getSqlResultMaxRows();
9294
final font = await AppSettings.instance.getSqlEditorFontSize();
9395
if (!mounted) return;
9496
setState(() {
9597
_historyMaxEntries = hist;
98+
_resultMaxRows = rows;
9699
_editorFontSize = font;
97100
});
98101
}
@@ -124,8 +127,11 @@ class _ExtensionSqlWorkspaceState
124127
});
125128

126129
try {
127-
final result = await ExtensionDriverSession.instance
128-
.query(widget.connectionRow, userSql);
130+
final result = await ExtensionDriverSession.instance.query(
131+
widget.connectionRow,
132+
userSql,
133+
limit: _resultMaxRows,
134+
);
129135
if (!mounted) return;
130136

131137
setState(() {
@@ -136,7 +142,10 @@ class _ExtensionSqlWorkspaceState
136142
} else {
137143
final elapsed =
138144
result.elapsedMs != null ? ' in ${result.elapsedMs}ms' : '';
139-
_statusLine = '${result.rows.length} row(s)$elapsed.';
145+
final capped = result.rows.length >= _resultMaxRows;
146+
_statusLine = capped
147+
? 'Showing first $_resultMaxRows row(s)$elapsed (result capped).'
148+
: '${result.rows.length} row(s)$elapsed.';
140149
}
141150
_running = false;
142151
});
Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,88 @@
1+
import 'dart:async';
2+
import 'dart:convert';
3+
import 'dart:io';
4+
5+
import 'package:flutter_test/flutter_test.dart';
6+
import 'package:querya_desktop/core/extensions/rpc/json_rpc_payload_limits.dart';
7+
import 'package:querya_desktop/core/extensions/rpc/json_rpc_stdio_client.dart';
8+
9+
void main() {
10+
group('boundedUtf8LineSplitter', () {
11+
test('splits lines and strips CR', () async {
12+
final lines = await Stream<List<int>>.fromIterable([
13+
utf8.encode('one\r\n'),
14+
utf8.encode('two\n'),
15+
]).transform(boundedUtf8LineSplitter(maxLineBytes: 1024)).toList();
16+
expect(lines, ['one', 'two']);
17+
});
18+
19+
test('fails closed when line exceeds max bytes', () async {
20+
final controller = StreamController<List<int>>();
21+
final errors = <Object>[];
22+
final sub = controller.stream
23+
.transform(boundedUtf8LineSplitter(maxLineBytes: 8))
24+
.listen((_) {}, onError: errors.add);
25+
26+
controller.add(utf8.encode('123456789')); // 9 bytes, no newline yet
27+
await Future<void>.delayed(Duration.zero);
28+
expect(errors, isNotEmpty);
29+
expect(errors.first, isA<JsonRpcPayloadTooLargeException>());
30+
await sub.cancel();
31+
await controller.close();
32+
});
33+
});
34+
35+
group('JsonRpcStdioClient payload bounds', () {
36+
test('completes pending request with payload-too-large error', () async {
37+
final stdout = StreamController<List<int>>();
38+
final stdin = StreamController<List<int>>();
39+
final client = JsonRpcStdioClient(
40+
stdout: stdout.stream,
41+
stdin: IOSink(stdin.sink),
42+
maxLineBytes: 32,
43+
requestTimeout: const Duration(seconds: 2),
44+
);
45+
46+
final pending = client.sendRequest('db.query', {'sql': 'SELECT 1'});
47+
// Drain request line from fake stdin.
48+
await stdin.stream.first;
49+
50+
// Oversized reply line (no newline until after overflow).
51+
stdout.add(List<int>.filled(40, 0x61)); // 'a' * 40
52+
await expectLater(pending, throwsA(isA<JsonRpcPayloadTooLargeException>()));
53+
54+
await client.close();
55+
await stdout.close();
56+
await stdin.close();
57+
});
58+
59+
test('decodes normal response', () async {
60+
final stdout = StreamController<List<int>>();
61+
final stdin = StreamController<List<int>>();
62+
final client = JsonRpcStdioClient(
63+
stdout: stdout.stream,
64+
stdin: IOSink(stdin.sink),
65+
requestTimeout: const Duration(seconds: 2),
66+
);
67+
68+
final pending = client.sendRequest('ping');
69+
await stdin.stream.first;
70+
stdout.add(
71+
utf8.encode(
72+
'${jsonEncode({
73+
'jsonrpc': '2.0',
74+
'id': 1,
75+
'result': {'ok': true},
76+
})}\n',
77+
),
78+
);
79+
final result = await pending;
80+
expect(result, isA<Map>());
81+
expect((result as Map)['ok'], true);
82+
83+
await client.close();
84+
await stdout.close();
85+
await stdin.close();
86+
});
87+
});
88+
}

0 commit comments

Comments
 (0)