Improved the relay realtime cleanup
🚀 Publish Trigger.dev Docker / units (push) Failing after 13m59s
🚀 Publish Trigger.dev Docker / typecheck (push) Failing after 14m0s
🚀 Publish Trigger.dev Docker / publish-webapp (push) Has been skipped
🚀 Publish Trigger.dev Docker / publish-worker (push) Has been skipped

This commit is contained in:
Eric Allam
2024-12-12 21:05:27 +00:00
parent 3d7e78c265
commit bb54d09f9b
@@ -7,6 +7,7 @@ import { singleton } from "~/utils/singleton";
export type RelayRealtimeStreamsOptions = {
ttl: number;
cleanupInterval: number;
fallbackIngestor: StreamIngestor;
fallbackResponder: StreamResponder;
waitForBufferTimeout?: number; // Time to wait for buffer in ms (default: 500ms)
@@ -17,6 +18,7 @@ interface RelayedStreamRecord {
stream: ReadableStream<Uint8Array>;
createdAt: number;
lastAccessed: number;
locked: boolean;
finalized: boolean;
}
@@ -33,7 +35,7 @@ export class RelayRealtimeStreams implements StreamIngestor, StreamResponder {
// Periodic cleanup
this.cleanupInterval = setInterval(() => {
this.cleanup();
}, this.options.ttl).unref();
}, this.options.cleanupInterval).unref();
}
async streamResponse(
@@ -76,6 +78,23 @@ export class RelayRealtimeStreams implements StreamIngestor, StreamResponder {
}
}
// Only 1 reader of the stream can use the relayed stream, the rest should use the fallback
if (record.locked) {
logger.debug("[RelayRealtimeStreams][streamResponse] Stream already locked, using fallback", {
streamId,
runId,
});
return this.options.fallbackResponder.streamResponse(
request,
runId,
streamId,
environment,
signal
);
}
record.locked = true;
record.lastAccessed = Date.now();
logger.debug("[RelayRealtimeStreams][streamResponse] Streaming from ephemeral record", {
@@ -106,7 +125,7 @@ export class RelayRealtimeStreams implements StreamIngestor, StreamResponder {
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache",
Connection: "keep-alive",
"x-relay-realtime-streams": "true",
"x-trigger-relay-realtime-streams": "true",
},
});
}
@@ -157,6 +176,7 @@ export class RelayRealtimeStreams implements StreamIngestor, StreamResponder {
createdAt: Date.now(),
lastAccessed: Date.now(),
finalized: false,
locked: false,
};
this._buffers.set(bufferKey, record);
} else {
@@ -167,12 +187,21 @@ export class RelayRealtimeStreams implements StreamIngestor, StreamResponder {
private cleanup() {
const now = Date.now();
logger.debug("[RelayRealtimeStreams][cleanup] Cleaning up old buffers", {
bufferCount: this._buffers.size,
});
for (const [key, record] of this._buffers.entries()) {
// If last accessed is older than ttl, clean up
if (now - record.lastAccessed > this.options.ttl) {
this.deleteBuffer(key);
}
}
logger.debug("[RelayRealtimeStreams][cleanup] Cleaned up old buffers", {
bufferCount: this._buffers.size,
});
}
private deleteBuffer(bufferKey: string) {
@@ -216,6 +245,7 @@ export class RelayRealtimeStreams implements StreamIngestor, StreamResponder {
function initializeRelayRealtimeStreams() {
return new RelayRealtimeStreams({
ttl: 1000 * 60 * 5, // 5 minutes
cleanupInterval: 1000 * 60, // 1 minute
fallbackIngestor: v1RealtimeStreams,
fallbackResponder: v1RealtimeStreams,
});