From 4dd27b44683b73abe1ade882ad81ef0504da8ff3 Mon Sep 17 00:00:00 2001 From: Eric Allam Date: Thu, 12 Dec 2024 15:28:59 +0000 Subject: [PATCH] Turn on the relay realtime stream service --- .../app/routes/realtime.v1.streams.$runId.$streamId.ts | 5 +++-- .../app/services/realtime/relayRealtimeStreams.server.ts | 3 ++- 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/apps/webapp/app/routes/realtime.v1.streams.$runId.$streamId.ts b/apps/webapp/app/routes/realtime.v1.streams.$runId.$streamId.ts index c370869d3..99e9cdb8d 100644 --- a/apps/webapp/app/routes/realtime.v1.streams.$runId.$streamId.ts +++ b/apps/webapp/app/routes/realtime.v1.streams.$runId.$streamId.ts @@ -1,6 +1,7 @@ import { ActionFunctionArgs } from "@remix-run/server-runtime"; import { z } from "zod"; import { $replica } from "~/db.server"; +import { relayRealtimeStreams } from "~/services/realtime/relayRealtimeStreams.server"; import { v1RealtimeStreams } from "~/services/realtime/v1StreamsGlobal.server"; import { createLoaderApiRoute } from "~/services/routeBuilders/apiBuilder.server"; @@ -16,7 +17,7 @@ export async function action({ request, params }: ActionFunctionArgs) { return new Response("No body provided", { status: 400 }); } - return v1RealtimeStreams.ingestData(request.body, $params.runId, $params.streamId); + return relayRealtimeStreams.ingestData(request.body, $params.runId, $params.streamId); } export const loader = createLoaderApiRoute( @@ -51,7 +52,7 @@ export const loader = createLoaderApiRoute( }, }, async ({ params, request, resource: run, authentication }) => { - return v1RealtimeStreams.streamResponse( + return relayRealtimeStreams.streamResponse( request, run.friendlyId, params.streamId, diff --git a/apps/webapp/app/services/realtime/relayRealtimeStreams.server.ts b/apps/webapp/app/services/realtime/relayRealtimeStreams.server.ts index 43e5c81b9..61b98a3b9 100644 --- a/apps/webapp/app/services/realtime/relayRealtimeStreams.server.ts +++ b/apps/webapp/app/services/realtime/relayRealtimeStreams.server.ts @@ -27,7 +27,7 @@ export class RelayRealtimeStreams implements StreamIngestor, StreamResponder { private waitForBufferInterval: number; constructor(private options: RelayRealtimeStreamsOptions) { - this.waitForBufferTimeout = options.waitForBufferTimeout ?? 5000; + this.waitForBufferTimeout = options.waitForBufferTimeout ?? 1200; this.waitForBufferInterval = options.waitForBufferInterval ?? 50; // Periodic cleanup @@ -106,6 +106,7 @@ export class RelayRealtimeStreams implements StreamIngestor, StreamResponder { "Content-Type": "text/event-stream", "Cache-Control": "no-cache", Connection: "keep-alive", + "x-relay-realtime-streams": "true", }, }); }