| 'use client' |
|
|
| import { EventStreamContentType, fetchEventSource } from '@microsoft/fetch-event-source' |
|
|
| import { getGetCurrentLlmQueryKey, getGetSceneJsonQueryKey } from '@/lib/api/default/default' |
| import type { AppEvent } from '@/lib/api/schemas' |
| import { queryClient } from '@/lib/queryClient' |
| import { useDownloadsStore } from '@/lib/stores/downloadsStore' |
| import { useEditorUiStore } from '@/lib/stores/editorUiStore' |
| import { useEventsStore } from '@/lib/stores/eventsStore' |
| import { useJobsStore } from '@/lib/stores/jobsStore' |
|
|
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| |
| export function connectEvents(baseUrl = '/api/v1'): () => void { |
| const controller = new AbortController() |
| const store = useEventsStore |
|
|
| store.getState().setStatus('connecting') |
|
|
| fetchEventSource(`${baseUrl}/events`, { |
| signal: controller.signal, |
| openWhenHidden: true, |
| headers: { Accept: 'text/event-stream' }, |
| async onopen(res) { |
| if (res.ok && res.headers.get('content-type')?.includes(EventStreamContentType)) { |
| store.getState().setStatus('open') |
| return |
| } |
| if (isFatalStatus(res.status)) { |
| |
| throw new FatalSseError(`SSE rejected: ${res.status} ${res.statusText}`) |
| } |
| |
| |
| throw new RetryableSseError(`SSE not ready: ${res.status}`) |
| }, |
| onmessage(ev) { |
| store.getState().onMessage(ev.id || null) |
| if (!ev.data) return |
| let parsed: AppEvent |
| try { |
| parsed = JSON.parse(ev.data) as AppEvent |
| } catch { |
| console.warn('[sse] malformed frame', ev.data) |
| return |
| } |
| if (process.env.NODE_ENV !== 'production') { |
| console.debug('[sse]', parsed.event, parsed) |
| } |
| dispatch(parsed) |
| }, |
| onerror(err) { |
| store.getState().onError(err instanceof Error ? err.message : String(err)) |
| if (err instanceof FatalSseError) { |
| |
| |
| store.getState().setStatus('error') |
| throw err |
| } |
| |
| |
| const attempt = store.getState().retryAttempt |
| return backoffMs(attempt) |
| }, |
| onclose() { |
| |
| |
| store.getState().setStatus('reconnecting') |
| }, |
| }).catch((err) => { |
| if ((err as { name?: string })?.name === 'AbortError') return |
| console.warn('[sse] fatal', err) |
| store.getState().setStatus('error') |
| }) |
|
|
| return () => { |
| controller.abort() |
| store.getState().reset() |
| } |
| } |
|
|
| |
| |
| |
|
|
| |
| |
| |
| |
| |
| |
| |
| const lastPageByJob = new Map<string, number>() |
|
|
| function invalidateScene(): void { |
| void queryClient.invalidateQueries({ queryKey: getGetSceneJsonQueryKey() }) |
| } |
|
|
| function dispatch(event: AppEvent): void { |
| switch (event.event) { |
| case 'snapshot': |
| |
| |
| useJobsStore.getState().setSnapshot(event.jobs) |
| useDownloadsStore.getState().setSnapshot(event.downloads) |
| lastPageByJob.clear() |
| return |
|
|
| case 'jobStarted': |
| useJobsStore.getState().started(event.id, event.kind) |
| lastPageByJob.set(event.id, -1) |
| return |
|
|
| case 'jobProgress': |
| useJobsStore.getState().progress(event) |
| |
| |
| |
| |
| { |
| const prev = lastPageByJob.get(event.jobId) ?? -1 |
| if (event.currentPage !== prev) { |
| lastPageByJob.set(event.jobId, event.currentPage) |
| if (prev >= 0) invalidateScene() |
| } |
| } |
| return |
|
|
| case 'jobWarning': |
| useJobsStore.getState().warning(event) |
| return |
|
|
| case 'jobFinished': |
| useJobsStore.getState().finished(event.id, event.status, event.error) |
| if (event.status === 'failed' && event.error) { |
| useEditorUiStore.getState().showError(event.error) |
| } |
| lastPageByJob.delete(event.id) |
| |
| |
| invalidateScene() |
| return |
|
|
| case 'downloadProgress': |
| useDownloadsStore.getState().progress(event) |
| return |
|
|
| case 'llmLoading': |
| case 'llmLoaded': |
| case 'llmFailed': |
| case 'llmUnloaded': |
| |
| |
| |
| void queryClient.invalidateQueries({ queryKey: getGetCurrentLlmQueryKey() }) |
| return |
|
|
| |
| |
| |
| } |
| } |
|
|
| |
| |
| |
|
|
| class FatalSseError extends Error {} |
| class RetryableSseError extends Error {} |
|
|
| function isFatalStatus(status: number): boolean { |
| |
| |
| if (status === 408 || status === 429) return false |
| return status >= 400 && status < 500 |
| } |
|
|
| |
| |
| |
| |
| function backoffMs(attempt: number): number { |
| const base = Math.min(10_000, 200 * 2 ** Math.min(attempt, 6)) |
| const jitter = base * 0.2 * (Math.random() * 2 - 1) |
| return Math.max(100, Math.round(base + jitter)) |
| } |
|
|