| 1 | import { fetchEventSource } from "@microsoft/fetch-event-source" |
| 2 | import { useAuthStore } from "@/stores/auth" |
| 3 | import { HttpClient } from "./httpClient" |
| 4 | |
| 5 | const TRAILING_SLASH_REGEX = /\/$/ |
| 6 | |
| 7 | /** Converte Record<string, string | number> in query string, escludendo undefined/null */ |
| 8 | function paramsToQueryString(params?: Record<string, string | number | undefined>): string { |
| 9 | if (!params || Object.keys(params).length === 0) { |
| 10 | return "" |
| 11 | } |
| 12 | const search = new URLSearchParams() |
| 13 | for (const [key, value] of Object.entries(params)) { |
| 14 | if (value !== undefined && value !== null) { |
| 15 | search.append(key, String(value)) |
| 16 | } |
| 17 | } |
| 18 | const qs = search.toString() |
| 19 | return qs ? `?${qs}` : "" |
| 20 | } |
| 21 | |
| 22 | export interface SSEClientOptions { |
| 23 | /** Path dell'endpoint (es. "/sca/overview/stream") */ |
| 24 | path: string |
| 25 | /** HTTP method (default: "GET") */ |
| 26 | method?: "GET" | "POST" | "PUT" | "PATCH" | "DELETE" |
| 27 | /** Base URL (default: HttpClient.baseURL ovvero "/api") */ |
| 28 | baseURL?: string |
| 29 | /** Parametri convertiti automaticamente in query string */ |
| 30 | params?: Record<string, string | number | undefined> |
| 31 | /** Handler per tipo di evento. Chiave = nome evento SSE (es. "start", "agent_result") + onOpen, onMessage, onError */ |
| 32 | handlers: Record<string, (data: unknown) => void> & { |
| 33 | onOpen?: (response: Response) => void | Promise<void> |
| 34 | onMessage?: (event: { event?: string; data: string }) => void |
| 35 | onError?: (error: unknown) => void |
| 36 | } |
| 37 | /** AbortController per cancellare lo stream */ |
| 38 | signal?: AbortSignal |
| 39 | } |
| 40 | |
| 41 | export async function createSSEStream(options: SSEClientOptions): Promise<void> { |
| 42 | const { path, params, handlers, signal } = options |
| 43 | const method = options.method ?? "GET" |
| 44 | const baseURL = options.baseURL ?? HttpClient.defaults.baseURL ?? "/api" |
| 45 | const authStore = useAuthStore() |
| 46 | |
| 47 | const base = (baseURL ?? "").replace(TRAILING_SLASH_REGEX, "") |
| 48 | const cleanPath = path.startsWith("/") ? path : `/${path}` |
| 49 | const queryString = paramsToQueryString(params) |
| 50 | const url = `${base}${cleanPath}${queryString}` |
| 51 | |
| 52 | const headers: Record<string, string> = {} |
| 53 | if (authStore.userToken) { |
| 54 | headers.Authorization = `Bearer ${authStore.userToken}` |
| 55 | } |
| 56 | |
| 57 | await fetchEventSource(url, { |
| 58 | method, |
| 59 | headers, |
| 60 | signal, |
| 61 | onopen(response) { |
| 62 | if (!response.ok) { |
| 63 | throw new Error(`Failed to connect: ${response.status} ${response.statusText}`) |
| 64 | } |
| 65 | handlers.onOpen?.(response) |
| 66 | return Promise.resolve() |
| 67 | }, |
| 68 | onmessage(event) { |
| 69 | handlers.onMessage?.({ event: event.event, data: event.data ?? "" }) |
| 70 | |
| 71 | if (!event.data) { |
| 72 | return |
| 73 | } |
| 74 | |
| 75 | try { |
| 76 | const data = JSON.parse(event.data) |
| 77 | const eventName = event.event || "message" |
| 78 | const lifecycleKeys = ["onOpen", "onMessage", "onError"] // nomi handler lifecycle, non eventi SSE |
| 79 | if (!lifecycleKeys.includes(eventName)) { |
| 80 | const handler = handlers[eventName] |
| 81 | if (handler) { |
| 82 | handler(data) |
| 83 | } |
| 84 | } |
| 85 | } catch (e) { |
| 86 | console.error("Failed to parse SSE data:", e) |
| 87 | } |
| 88 | }, |
| 89 | onerror(err) { |
| 90 | handlers.onError?.(err) |
| 91 | throw err |
| 92 | } |
| 93 | }) |
| 94 | } |