Files
Eric Allam 78e2d0e9f9
🚀 Publish Trigger.dev Docker / ʦ TypeScript (push) Has been cancelled
🚀 Publish Trigger.dev Docker / Unit Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / e2e Tests (push) Has been cancelled
🚀 Publish Trigger.dev Docker / publish (push) Has been cancelled
Fix SSE causing app crashes when inner async loop takes too long
2023-08-14 16:21:45 +01:00

71 lines
1.6 KiB
TypeScript

import { eventStream } from "remix-utils";
import { logger } from "~/services/logger.server";
type SseProps = {
request: Request;
pingInterval?: number;
updateInterval?: number;
run: (send: (event: Event) => void, stop: () => void) => void;
};
type Event = {
/**
* @default "update"
*/
event?: string;
data: string;
};
export function sse({ request, pingInterval = 1000, updateInterval = 348, run }: SseProps) {
let pinger: NodeJS.Timer | undefined = undefined;
let updater: NodeJS.Timer | undefined = undefined;
const abort = () => {
if (pinger) {
clearInterval(pinger);
}
if (updater) {
clearInterval(updater);
}
};
return eventStream(request.signal, (send) => {
const safeSend = (args: { event?: string; data: string }) => {
try {
send(args);
} catch (error) {
if (error instanceof Error) {
if (error.name !== "TypeError") {
logger.debug("Error sending SSE, aborting", {
error: {
name: error.name,
message: error.message,
stack: error.stack,
},
args,
});
}
} else {
logger.debug("Uknown error sending SSE, aborting", {
error,
args,
});
}
abort();
}
};
pinger = setInterval(() => {
safeSend({ event: "ping", data: new Date().toISOString() });
}, pingInterval);
updater = setInterval(async () => {
run(safeSend, abort);
}, updateInterval);
return abort;
});
}