diff --git a/apps/webapp/app/services/realtime/relayRealtimeStreams.server.ts b/apps/webapp/app/services/realtime/relayRealtimeStreams.server.ts index 61b98a3b9..46e16ff5e 100644 --- a/apps/webapp/app/services/realtime/relayRealtimeStreams.server.ts +++ b/apps/webapp/app/services/realtime/relayRealtimeStreams.server.ts @@ -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; 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, });