decodeSseJson function
Stream<Object>
decodeSseJson(
- Stream<
List< body, {int> > - AcpSseLimits limits = const AcpSseLimits(),
- AcpRemoteDiagnosticHandler? onDiagnostic,
Decodes JSON objects and arrays from a byte-oriented SSE response body.
Comments and non-data fields are ignored. Malformed JSON and primitive
payloads emit a redacted diagnostic and are skipped so the JSON-RPC layer
remains responsible for message-shape validation.
Implementation
Stream<Object> decodeSseJson(
Stream<List<int>> body, {
AcpSseLimits limits = const AcpSseLimits(),
AcpRemoteDiagnosticHandler? onDiagnostic,
}) async* {
limits.validate();
final lines = LineBuffer(maximumLineBytes: limits.maximumLineBytes);
final dataLines = <String>[];
var dataBytes = 0;
void diagnose(AcpRemoteDiagnostic diagnostic) {
try {
onDiagnostic?.call(diagnostic);
} on Object {
// Diagnostics are observational and must not alter stream lifecycle.
}
}
Object? takeEvent() {
if (dataLines.isEmpty) {
return null;
}
final data = dataLines.join('\n');
dataLines.clear();
dataBytes = 0;
if (data.trim().isEmpty) {
return null;
}
final Object? decoded;
try {
decoded = decodeBoundedJson(
data,
maximumNestingDepth: limits.maximumJsonNestingDepth,
);
} on FormatException {
diagnose(
const AcpRemoteDiagnostic('Skipping malformed SSE JSON payload'),
);
return null;
}
if (decoded is Map<Object?, Object?> || decoded is List<Object?>) {
return decoded;
}
diagnose(const AcpRemoteDiagnostic('Skipping primitive SSE JSON payload'));
return null;
}
void addEventLine(String line) {
if (!line.startsWith('data:')) {
return;
}
var value = line.substring(5);
if (value.startsWith(' ')) {
value = value.substring(1);
}
final nextLength =
dataBytes + utf8.encode(value).length + (dataLines.isEmpty ? 0 : 1);
if (nextLength > limits.maximumEventBytes) {
throw AcpSseLimitException(
length: nextLength,
maximum: limits.maximumEventBytes,
);
}
dataBytes = nextLength;
dataLines.add(value);
}
String decodeLine(Uint8List bytes) {
var line = utf8.decode(bytes);
if (line.endsWith('\r')) {
line = line.substring(0, line.length - 1);
}
return line;
}
await for (final chunk in body) {
for (final bytes in lines.add(chunk)) {
final line = decodeLine(bytes);
if (line.isEmpty) {
final event = takeEvent();
if (event != null) {
yield event;
}
} else {
addEventLine(line);
}
}
}
final finalBytes = lines.flush();
if (finalBytes != null) {
final line = decodeLine(finalBytes);
if (line.isNotEmpty) {
addEventLine(line);
}
}
final event = takeEvent();
if (event != null) {
yield event;
}
}