(runId: string, draftId: string, draftTimestamp: string)
| 805 | } |
| 806 | |
| 807 | function startReplyStream(runId: string, draftId: string, draftTimestamp: string) { |
| 808 | closeReplyStream(); |
| 809 | const stream = new EventSource(`/api/codex-threads/${encodeURIComponent(threadId)}/runs/${encodeURIComponent(runId)}/events`); |
| 810 | replyStreamRef.current = stream; |
| 811 | |
| 812 | const handleSnapshot = (event: MessageEvent<string>) => { |
| 813 | try { |
| 814 | const snapshot = JSON.parse(event.data) as ReplyRunSnapshot; |
| 815 | applyRunSnapshot(snapshot, draftId, draftTimestamp); |
| 816 | } catch { |
| 817 | // Ignore malformed SSE chunks. |
| 818 | } |
| 819 | }; |
| 820 | |
| 821 | stream.addEventListener('snapshot', handleSnapshot); |
| 822 | stream.addEventListener('started', handleSnapshot); |
| 823 | stream.addEventListener('assistant', handleSnapshot); |
| 824 | stream.addEventListener('commentary', handleSnapshot); |
| 825 | stream.addEventListener('done', handleSnapshot); |
| 826 | stream.addEventListener('failed', handleSnapshot); |
| 827 | stream.onerror = () => { |
| 828 | if (stream.readyState !== EventSource.CLOSED) { |
| 829 | return; |
| 830 | } |
| 831 | replaceStreamFailureMessage(draftId, draftTimestamp, '连接已中断,请稍后重试。'); |
| 832 | setReplyError('连接已中断,请稍后重试。'); |
| 833 | setSendingReply(false); |
| 834 | closeReplyStream(); |
| 835 | }; |
| 836 | } |
| 837 | |
| 838 | useEffect(() => { |
| 839 | setThreadMessages(initialThreadMessages); |
no test coverage detected