process-sse-stream.ts 3.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122
  1. import type { StatelessEvent } from '@/lib/types/chat';
  2. import type { StreamBuffer } from '@/lib/buffer/stream-buffer';
  3. import { createLogger } from '@/lib/logger';
  4. const log = createLogger('SSEStream');
  5. /**
  6. * Thin SSE parser — reads the /api/chat response stream and pushes
  7. * typed events into a StreamBuffer. All pacing, state management,
  8. * and UI updates are handled by the buffer's tick loop and callbacks.
  9. */
  10. export async function processSSEStream(
  11. response: Response,
  12. sessionId: string,
  13. buffer: StreamBuffer,
  14. signal?: AbortSignal,
  15. ): Promise<void> {
  16. const reader = response.body?.getReader();
  17. if (!reader) {
  18. throw new Error('No response body');
  19. }
  20. const decoder = new TextDecoder();
  21. let sseBuffer = '';
  22. let currentMessageId: string | null = null;
  23. try {
  24. while (true) {
  25. const { done, value } = await reader.read();
  26. if (done) break;
  27. const chunk = decoder.decode(value, { stream: true });
  28. sseBuffer += chunk;
  29. // Process complete SSE events (split on double newline)
  30. const events = sseBuffer.split('\n\n');
  31. sseBuffer = events.pop() || '';
  32. for (const eventStr of events) {
  33. const line = eventStr.trim();
  34. if (!line.startsWith('data: ')) continue;
  35. let sseError: Error | null = null;
  36. try {
  37. const event: StatelessEvent = JSON.parse(line.slice(6));
  38. switch (event.type) {
  39. case 'agent_start': {
  40. const { messageId, agentId, agentName, agentAvatar, agentColor } = event.data;
  41. currentMessageId = messageId;
  42. buffer.pushAgentStart({
  43. messageId,
  44. agentId,
  45. agentName,
  46. avatar: agentAvatar,
  47. color: agentColor,
  48. });
  49. break;
  50. }
  51. case 'agent_end': {
  52. buffer.pushAgentEnd({
  53. messageId: event.data.messageId,
  54. agentId: event.data.agentId,
  55. });
  56. break;
  57. }
  58. case 'text_delta': {
  59. const targetId = event.data.messageId ?? currentMessageId;
  60. if (!targetId) break;
  61. buffer.pushText(targetId, event.data.content);
  62. break;
  63. }
  64. case 'action': {
  65. const targetId = event.data.messageId ?? currentMessageId;
  66. if (!targetId) break;
  67. if (signal?.aborted) break;
  68. buffer.pushAction({
  69. messageId: targetId,
  70. actionId: event.data.actionId,
  71. actionName: event.data.actionName,
  72. params: event.data.params,
  73. agentId: event.data.agentId,
  74. });
  75. break;
  76. }
  77. case 'thinking': {
  78. buffer.pushThinking(event.data);
  79. break;
  80. }
  81. case 'cue_user': {
  82. buffer.pushCueUser(event.data);
  83. break;
  84. }
  85. case 'done': {
  86. buffer.pushDone(event.data);
  87. break;
  88. }
  89. case 'error': {
  90. sseError = new Error(event.data.message);
  91. buffer.pushError(event.data.message);
  92. break;
  93. }
  94. }
  95. } catch (parseError) {
  96. log.warn('[SSE] Parse error:', parseError);
  97. }
  98. if (sseError) throw sseError;
  99. }
  100. }
  101. } finally {
  102. reader.releaseLock();
  103. }
  104. }