diff --git a/src/activity/activitiesProcessed.ts b/src/activity/activitiesProcessed.ts new file mode 100644 index 0000000..cd57e79 --- /dev/null +++ b/src/activity/activitiesProcessed.ts @@ -0,0 +1,110 @@ +import log from "loglevel"; +import {getAccessToken} from "@/auth/token"; +import {mammonUrl, refreshAccessToken} from "@/mammon"; + +const STREAM_URL = mammonUrl + "activities/processed"; +const RECONNECT_DELAY_MILLIS = 3_000; +const FRAME_SEPARATOR = "\n\n"; +const PROCESSED_EVENT = "event:processed"; + +const listeners = new Set<() => void>(); + +let controller: AbortController | undefined; +let reconnectTimer: ReturnType | undefined; + +const notifyProcessed = () => listeners.forEach(listener => listener()); + +const scheduleReconnect = () => { + if (listeners.size === 0) { + return; + } + reconnectTimer = setTimeout(() => connect(), RECONNECT_DELAY_MILLIS); +}; + +const consume = async (body: ReadableStream) => { + const reader = body.getReader(); + const decoder = new TextDecoder(); + let buffer = ""; + + for (; ;) { + const {value, done} = await reader.read(); + + if (done) { + return; + } + buffer += decoder.decode(value, {stream: true}); + + let boundary = buffer.indexOf(FRAME_SEPARATOR); + + while (boundary !== -1) { + const frame = buffer.slice(0, boundary); + buffer = buffer.slice(boundary + FRAME_SEPARATOR.length); + + if (frame.split("\n").some(line => line.trim() === PROCESSED_EVENT)) { + notifyProcessed(); + } + boundary = buffer.indexOf(FRAME_SEPARATOR); + } + } +}; + +const connect = async (retried = false): 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}`} : {}), + }, + 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); + scheduleReconnect(); + } catch (error) { + if (listeners.size === 0 || signal?.aborted) { + return; + } + log.debug("Activity processing SSE error, reconnecting", error); + scheduleReconnect(); + } +}; + +const close = () => { + if (reconnectTimer) { + clearTimeout(reconnectTimer); + reconnectTimer = undefined; + } + controller?.abort(); + controller = undefined; +}; + +export const onActivitiesProcessed = (listener: () => void): (() => void) => { + listeners.add(listener); + + if (!controller) { + controller = new AbortController(); + connect(); + } + + return () => { + if (listeners.delete(listener) && listeners.size === 0) { + close(); + } + }; +}; diff --git a/src/activity/index.ts b/src/activity/index.ts index 389d889..b7b91ea 100644 --- a/src/activity/index.ts +++ b/src/activity/index.ts @@ -1 +1,3 @@ export {default as SourceLabel} from './SourceLabel.vue'; +export {onActivitiesProcessed} from './activitiesProcessed'; +export {useActivitiesProcessed} from './useActivitiesProcessed'; diff --git a/src/activity/useActivitiesProcessed.ts b/src/activity/useActivitiesProcessed.ts new file mode 100644 index 0000000..0bdc2fa --- /dev/null +++ b/src/activity/useActivitiesProcessed.ts @@ -0,0 +1,8 @@ +import {onUnmounted} from "vue"; +import {onActivitiesProcessed} from "./activitiesProcessed"; + +export const useActivitiesProcessed = (listener: () => void): (() => void) => { + const unsubscribe = onActivitiesProcessed(listener); + onUnmounted(unsubscribe); + return unsubscribe; +}; diff --git a/src/market/acquisition/acquisition.ts b/src/market/acquisition/acquisition.ts index fc41681..1e1b586 100644 --- a/src/market/acquisition/acquisition.ts +++ b/src/market/acquisition/acquisition.ts @@ -2,6 +2,7 @@ import {defineStore} from "pinia"; import {computed, ref} from "vue"; import {acquisitionApi, activityApi} from "@/mammon"; import {AcquisitionResponse, ActivitySourceResponse} from "@/generated/mammon"; +import {onActivitiesProcessed} from "@/activity"; export type RawAcquiredType = { id: string; @@ -66,5 +67,7 @@ export const useAcquiredTypesStore = defineStore('market-acquisition', () => { refresh(); + onActivitiesProcessed(refresh); + return { acquiredTypes: types, addAcquiredType, removeAcquiredType, processNewActivities, refresh }; });