( stream: ReadableStream<Uint8Array>, schema: DescMessage, onMessage: (chunk: T) => void, )
| 28 | } |
| 29 | |
| 30 | export const processStream = async <T>( |
| 31 | stream: ReadableStream<Uint8Array>, |
| 32 | schema: DescMessage, |
| 33 | onMessage: (chunk: T) => void, |
| 34 | ) => { |
| 35 | const { slowStreamingMode } = useDevtoolStore.getState(); |
| 36 | const reader = stream.getReader(); |
| 37 | const decoder = new TextDecoder("utf-8"); |
| 38 | let buffer = ""; |
| 39 | |
| 40 | while (true) { |
| 41 | const { done, value } = await reader.read(); |
| 42 | if (done) break; |
| 43 | |
| 44 | buffer += decoder.decode(value, { stream: true }); |
| 45 | |
| 46 | let boundary; |
| 47 | while ((boundary = buffer.indexOf("\n")) !== -1) { |
| 48 | const message = buffer.slice(0, boundary); |
| 49 | buffer = buffer.slice(boundary + 1); |
| 50 | |
| 51 | try { |
| 52 | const parsedValue = JSON.parse(message); |
| 53 | const messageData = parsedValue.result || parsedValue; |
| 54 | onMessage(fromJson(schema, messageData) as T); |
| 55 | } catch (err) { |
| 56 | logError("Error parsing message from stream", err, message); |
| 57 | } |
| 58 | } |
| 59 | |
| 60 | if (import.meta.env.DEV && slowStreamingMode) { |
| 61 | await new Promise((resolve) => setTimeout(resolve, 500)); |
| 62 | } |
| 63 | } |
| 64 | }; |
no test coverage detected