api_streaming.dart 2.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112
  1. import 'dart:async';
  2. import 'dart:convert';
  3. class ServerSentEvent {
  4. final String id;
  5. final String event;
  6. final String data;
  7. late final dynamic _jsonData = _tryJsonDecode(data);
  8. final int? retry;
  9. dynamic get jsonData {
  10. return _jsonData;
  11. }
  12. dynamic _tryJsonDecode(String dataString) {
  13. try {
  14. return jsonDecode(dataString);
  15. } catch (_) {
  16. return null;
  17. }
  18. }
  19. ServerSentEvent({
  20. required this.id,
  21. required this.event,
  22. required this.data,
  23. this.retry,
  24. });
  25. }
  26. class ResponseStreamMessage {
  27. final String message;
  28. late final ServerSentEvent _serverSentEvent = _parseSseString(message);
  29. ServerSentEvent get serverSentEvent {
  30. return _serverSentEvent;
  31. }
  32. ResponseStreamMessage({
  33. required this.message,
  34. });
  35. ServerSentEvent _parseSseString(String sseString) {
  36. String id = '';
  37. String event = '';
  38. String data = '';
  39. int? retry;
  40. // Split the input string by lines
  41. final lines = sseString.split('\n');
  42. for (var line in lines) {
  43. if (line.startsWith('id:')) {
  44. id = line.substring(3).trim();
  45. } else if (line.startsWith('event:')) {
  46. event = line.substring(6).trim();
  47. } else if (line.startsWith('data:')) {
  48. data += line.substring(5).trim() + '\n';
  49. } else if (line.startsWith('retry:')) {
  50. retry = int.tryParse(line.substring(6).trim());
  51. }
  52. }
  53. // Remove the trailing newline character from data
  54. if (data.endsWith('\n')) {
  55. data = data.substring(0, data.length - 1);
  56. }
  57. return ServerSentEvent(id: id, event: event, data: data, retry: retry);
  58. }
  59. }
  60. class ServerSentEventLineTransformer
  61. extends StreamTransformerBase<String, String> {
  62. @override
  63. Stream<String> bind(Stream<String> stream) {
  64. return Stream<String>.eventTransformed(
  65. stream,
  66. (EventSink<String> sink) => _NewlineEventSink(sink),
  67. );
  68. }
  69. }
  70. class _NewlineEventSink implements EventSink<String> {
  71. final EventSink<String> _sink;
  72. String _buffer = '';
  73. _NewlineEventSink(this._sink);
  74. @override
  75. void add(String data) {
  76. if (data.isEmpty) {
  77. _sink.add(_buffer);
  78. _buffer = '';
  79. } else {
  80. _buffer += (_buffer.isNotEmpty ? '\n' : '') + data;
  81. }
  82. }
  83. @override
  84. void addError(Object error, [StackTrace? stackTrace]) {
  85. _sink.addError(error, stackTrace);
  86. }
  87. @override
  88. void close() {
  89. if (_buffer.isNotEmpty) {
  90. _sink.add(_buffer);
  91. _buffer = '';
  92. }
  93. _sink.close();
  94. }
  95. }