feat: consume activity processed SSE
This commit is contained in:
@@ -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<typeof setTimeout> | 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<Uint8Array>) => {
|
||||
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<void> => {
|
||||
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();
|
||||
}
|
||||
};
|
||||
};
|
||||
Reference in New Issue
Block a user