From 5a974bac8f95d039236e25a4d3eb5e5f6abd1173 Mon Sep 17 00:00:00 2001 From: Sirttas Date: Tue, 7 Jul 2026 21:56:44 +0200 Subject: [PATCH] refactor use axios for fetch --- src/activity/activitiesProcessed.ts | 27 ++++++--------------------- 1 file changed, 6 insertions(+), 21 deletions(-) diff --git a/src/activity/activitiesProcessed.ts b/src/activity/activitiesProcessed.ts index cd57e79..116dc35 100644 --- a/src/activity/activitiesProcessed.ts +++ b/src/activity/activitiesProcessed.ts @@ -1,8 +1,6 @@ import log from "loglevel"; -import {getAccessToken} from "@/auth/token"; -import {mammonUrl, refreshAccessToken} from "@/mammon"; +import {activityApi} from "@/mammon"; -const STREAM_URL = mammonUrl + "activities/processed"; const RECONNECT_DELAY_MILLIS = 3_000; const FRAME_SEPARATOR = "\n\n"; const PROCESSED_EVENT = "event:processed"; @@ -48,33 +46,20 @@ const consume = async (body: ReadableStream) => { } }; -const connect = async (retried = false): Promise => { +const connect = async (): Promise => { if (listeners.size === 0) { return; } const signal = controller?.signal; - const token = getAccessToken(); try { - const response = await fetch(STREAM_URL, { - headers: { - Accept: "text/event-stream", - ...(token ? {Authorization: `Bearer ${token}`} : {}), - }, + const response = await activityApi.streamProcessing({ + responseType: "stream", + adapter: "fetch", signal, }); - if (response.status === 401) { - if (!retried && await refreshAccessToken() && listeners.size > 0) { - return connect(true); - } - return; - } - - if (!response.ok || !response.body) { - throw new Error(`Activity processing SSE stream failed with ${response.status}`); - } - await consume(response.body); + await consume(response.data as unknown as ReadableStream); scheduleReconnect(); } catch (error) { if (listeners.size === 0 || signal?.aborted) {