(chatId, query, history, docs, datasets, messageOffset, callbacks)
| 80 | * Returns a promise that resolves when streaming is complete. |
| 81 | */ |
| 82 | export async function sendMessageStream(chatId, query, history, docs, datasets, messageOffset, callbacks) { |
| 83 | const { onStep, onToken, onDocuments, onDone, onError } = callbacks; |
| 84 | |
| 85 | const res = await fetch(`${API_BASE}/chats/${chatId}/message/stream`, { |
| 86 | method: 'POST', |
| 87 | headers: getHeaders(), |
| 88 | body: JSON.stringify({ query, history, docs, datasets, messageOffset }), |
| 89 | }); |
| 90 | |
| 91 | if (!res.ok) { |
| 92 | const errText = await res.text(); |
| 93 | throw new Error(errText || 'Chat stream failed'); |
| 94 | } |
| 95 | |
| 96 | const reader = res.body.getReader(); |
| 97 | const decoder = new TextDecoder(); |
| 98 | let buffer = ''; |
| 99 | |
| 100 | while (true) { |
| 101 | const { done, value } = await reader.read(); |
| 102 | if (done) break; |
| 103 | |
| 104 | buffer += decoder.decode(value, { stream: true }); |
| 105 | |
| 106 | // Process complete SSE messages (separated by double newline) |
| 107 | let boundary = buffer.indexOf('\n\n'); |
| 108 | while (boundary !== -1) { |
| 109 | const message = buffer.substring(0, boundary); |
| 110 | buffer = buffer.substring(boundary + 2); |
| 111 | |
| 112 | // Parse the SSE message |
| 113 | let eventType = 'message'; |
| 114 | let dataStr = ''; |
| 115 | for (const line of message.split('\n')) { |
| 116 | if (line.startsWith('event: ')) { |
| 117 | eventType = line.substring(7).trim(); |
| 118 | } else if (line.startsWith('data: ')) { |
| 119 | dataStr += line.substring(6); |
| 120 | } |
| 121 | } |
| 122 | |
| 123 | if (dataStr) { |
| 124 | try { |
| 125 | const parsed = JSON.parse(dataStr); |
| 126 | switch (eventType) { |
| 127 | case 'step': |
| 128 | onStep?.(parsed.step); |
| 129 | break; |
| 130 | case 'token': |
| 131 | onToken?.(parsed.token); |
| 132 | break; |
| 133 | case 'documents': |
| 134 | onDocuments?.(parsed.documents); |
| 135 | break; |
| 136 | case 'done': |
| 137 | onDone?.(parsed); |
| 138 | break; |
| 139 | case 'error': |
no test coverage detected