decodeSseJson function

Stream<Object> decodeSseJson(
  1. Stream<List<int>> body, {
  2. AcpSseLimits limits = const AcpSseLimits(),
  3. 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;
  }
}