New eveal #32

Merged
Sirttas merged 115 commits from new-eveal into main 2026-07-14 17:15:11 +02:00
Showing only changes of commit 5a974bac8f - Show all commits
+6 -21
View File
@@ -1,8 +1,6 @@
import log from "loglevel"; import log from "loglevel";
import {getAccessToken} from "@/auth/token"; import {activityApi} from "@/mammon";
import {mammonUrl, refreshAccessToken} from "@/mammon";
const STREAM_URL = mammonUrl + "activities/processed";
const RECONNECT_DELAY_MILLIS = 3_000; const RECONNECT_DELAY_MILLIS = 3_000;
const FRAME_SEPARATOR = "\n\n"; const FRAME_SEPARATOR = "\n\n";
const PROCESSED_EVENT = "event:processed"; const PROCESSED_EVENT = "event:processed";
@@ -48,33 +46,20 @@ const consume = async (body: ReadableStream<Uint8Array>) => {
} }
}; };
const connect = async (retried = false): Promise<void> => { const connect = async (): Promise<void> => {
if (listeners.size === 0) { if (listeners.size === 0) {
return; return;
} }
const signal = controller?.signal; const signal = controller?.signal;
const token = getAccessToken();
try { try {
const response = await fetch(STREAM_URL, { const response = await activityApi.streamProcessing({
headers: { responseType: "stream",
Accept: "text/event-stream", adapter: "fetch",
...(token ? {Authorization: `Bearer ${token}`} : {}),
},
signal, signal,
}); });
if (response.status === 401) { await consume(response.data as unknown as ReadableStream<Uint8Array>);
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);
scheduleReconnect(); scheduleReconnect();
} catch (error) { } catch (error) {
if (listeners.size === 0 || signal?.aborted) { if (listeners.size === 0 || signal?.aborted) {