diff --git a/.server-changes/streams-version-s2-guard.md b/.server-changes/streams-version-s2-guard.md index 2c607652e..6114ee44e 100644 --- a/.server-changes/streams-version-s2-guard.md +++ b/.server-changes/streams-version-s2-guard.md @@ -3,4 +3,4 @@ area: webapp type: fix --- -Stop creating runs against a realtime streams backend the deployment cannot serve. +Runs no longer end up with realtime streams that cannot be read or written. diff --git a/apps/webapp/app/services/realtime/v1StreamsGlobal.server.ts b/apps/webapp/app/services/realtime/v1StreamsGlobal.server.ts index 8bd1f9d1a..435198e32 100644 --- a/apps/webapp/app/services/realtime/v1StreamsGlobal.server.ts +++ b/apps/webapp/app/services/realtime/v1StreamsGlobal.server.ts @@ -96,27 +96,51 @@ function streamPrefixFor(environment: AuthenticatedEnvironment, basin: string): return segments.join("/"); } -/** - * Resolve the streams version to stamp on a run, falling back to - * `REALTIME_STREAMS_DEFAULT_VERSION` when the caller expresses no preference. - * - * v2 is only ever returned when S2 is actually configured. A run stamped v2 on - * a deployment without S2 is unusable: `getRealtimeStreamInstance` throws for - * the life of the run, and no read or write against its streams can succeed. - * v1 is a working backend, so an unsatisfiable v2 degrades to it. - */ -export function determineRealtimeStreamsVersion(streamVersion?: string): "v1" | "v2" { - const requested = streamVersion ?? env.REALTIME_STREAMS_DEFAULT_VERSION; +export type RealtimeStreamsVersionConfig = { + defaultVersion: "v1" | "v2"; + basin?: string; + accessToken?: string; + skipAccessTokens: boolean; + perOrgBasinsEnabled: boolean; +}; - if ( - requested === "v2" && - env.REALTIME_STREAMS_S2_BASIN && - (env.REALTIME_STREAMS_S2_ACCESS_TOKEN || env.REALTIME_STREAMS_S2_SKIP_ACCESS_TOKENS === "true") - ) { - return "v2"; +/** + * Resolve the streams version to stamp on a run, falling back to the + * deployment default when the caller expresses no preference. + * + * v2 is only ever returned when S2 can actually serve it. A run stamped v2 on a + * deployment without S2 is unusable: `getRealtimeStreamInstance` throws for the + * life of the run, and no read or write against its streams can succeed. v1 is + * a working backend, so an unsatisfiable v2 degrades to it. + * + * A basin can come from the global setting or from per-org provisioning, so + * either satisfies the basin requirement. This mirrors {@link resolveStreamBasin}, + * which resolves run, session and organization basins ahead of the global one. + */ +export function resolveRealtimeStreamsVersion( + streamVersion: string | undefined, + config: RealtimeStreamsVersionConfig +): "v1" | "v2" { + const requested = streamVersion ?? config.defaultVersion; + + if (requested !== "v2") { + return "v1"; } - return "v1"; + const hasCredentials = Boolean(config.accessToken) || config.skipAccessTokens; + const hasBasin = Boolean(config.basin) || config.perOrgBasinsEnabled; + + return hasCredentials && hasBasin ? "v2" : "v1"; +} + +export function determineRealtimeStreamsVersion(streamVersion?: string): "v1" | "v2" { + return resolveRealtimeStreamsVersion(streamVersion, { + defaultVersion: env.REALTIME_STREAMS_DEFAULT_VERSION, + basin: env.REALTIME_STREAMS_S2_BASIN, + accessToken: env.REALTIME_STREAMS_S2_ACCESS_TOKEN, + skipAccessTokens: env.REALTIME_STREAMS_S2_SKIP_ACCESS_TOKENS === "true", + perOrgBasinsEnabled: env.REALTIME_STREAMS_PER_ORG_BASINS_ENABLED === "true", + }); } const s2RealtimeStreamsCache = singleton( diff --git a/apps/webapp/test/determineRealtimeStreamsVersion.test.ts b/apps/webapp/test/determineRealtimeStreamsVersion.test.ts index 6d7307085..ec5b78a97 100644 --- a/apps/webapp/test/determineRealtimeStreamsVersion.test.ts +++ b/apps/webapp/test/determineRealtimeStreamsVersion.test.ts @@ -1,70 +1,87 @@ -import { beforeEach, describe, expect, it, vi } from "vitest"; +import { describe, expect, it } from "vitest"; +import { + resolveRealtimeStreamsVersion, + type RealtimeStreamsVersionConfig, +} from "~/services/realtime/v1StreamsGlobal.server"; -const envMock = vi.hoisted(() => ({ - REALTIME_STREAMS_DEFAULT_VERSION: "v1" as "v1" | "v2", - REALTIME_STREAMS_S2_BASIN: undefined as string | undefined, - REALTIME_STREAMS_S2_ACCESS_TOKEN: undefined as string | undefined, - REALTIME_STREAMS_S2_SKIP_ACCESS_TOKENS: "false", -})); +const NO_S2: RealtimeStreamsVersionConfig = { + defaultVersion: "v1", + basin: undefined, + accessToken: undefined, + skipAccessTokens: false, + perOrgBasinsEnabled: false, +}; -vi.mock("~/env.server", () => ({ env: envMock })); +const GLOBAL_BASIN: RealtimeStreamsVersionConfig = { + ...NO_S2, + basin: "a-basin", + accessToken: "a-token", +}; -import { determineRealtimeStreamsVersion } from "~/services/realtime/v1StreamsGlobal.server"; +const PER_ORG_BASINS: RealtimeStreamsVersionConfig = { + ...NO_S2, + accessToken: "a-token", + perOrgBasinsEnabled: true, +}; -function configureS2() { - envMock.REALTIME_STREAMS_S2_BASIN = "a-basin"; - envMock.REALTIME_STREAMS_S2_ACCESS_TOKEN = "a-token"; -} - -beforeEach(() => { - envMock.REALTIME_STREAMS_DEFAULT_VERSION = "v1"; - envMock.REALTIME_STREAMS_S2_BASIN = undefined; - envMock.REALTIME_STREAMS_S2_ACCESS_TOKEN = undefined; - envMock.REALTIME_STREAMS_S2_SKIP_ACCESS_TOKENS = "false"; -}); - -describe("determineRealtimeStreamsVersion", () => { - it("honours an explicit v2 when S2 is configured", () => { - configureS2(); - expect(determineRealtimeStreamsVersion("v2")).toBe("v2"); +describe("resolveRealtimeStreamsVersion", () => { + it("honours an explicit v2 when a global basin is configured", () => { + expect(resolveRealtimeStreamsVersion("v2", GLOBAL_BASIN)).toBe("v2"); }); - it("accepts a skip-tokens deployment as configured", () => { - envMock.REALTIME_STREAMS_S2_BASIN = "a-basin"; - envMock.REALTIME_STREAMS_S2_SKIP_ACCESS_TOKENS = "true"; - expect(determineRealtimeStreamsVersion("v2")).toBe("v2"); + it("honours an explicit v2 when only per-org basins are configured", () => { + expect(resolveRealtimeStreamsVersion("v2", PER_ORG_BASINS)).toBe("v2"); + }); + + it("accepts a skip-tokens deployment as credentialed", () => { + expect( + resolveRealtimeStreamsVersion("v2", { + ...NO_S2, + basin: "a-basin", + skipAccessTokens: true, + }) + ).toBe("v2"); }); it("degrades an explicit v2 to v1 when S2 is not configured", () => { - expect(determineRealtimeStreamsVersion("v2")).toBe("v1"); + expect(resolveRealtimeStreamsVersion("v2", NO_S2)).toBe("v1"); }); it("falls back to the default version when the caller expresses no preference", () => { - configureS2(); - envMock.REALTIME_STREAMS_DEFAULT_VERSION = "v2"; - expect(determineRealtimeStreamsVersion()).toBe("v2"); + expect( + resolveRealtimeStreamsVersion(undefined, { ...GLOBAL_BASIN, defaultVersion: "v2" }) + ).toBe("v2"); }); it("degrades a v2 default to v1 when S2 is not configured", () => { - envMock.REALTIME_STREAMS_DEFAULT_VERSION = "v2"; - expect(determineRealtimeStreamsVersion()).toBe("v1"); + expect(resolveRealtimeStreamsVersion(undefined, { ...NO_S2, defaultVersion: "v2" })).toBe("v1"); }); - it("requires a basin, not just a token", () => { - envMock.REALTIME_STREAMS_S2_ACCESS_TOKEN = "a-token"; - envMock.REALTIME_STREAMS_DEFAULT_VERSION = "v2"; - expect(determineRealtimeStreamsVersion()).toBe("v1"); - expect(determineRealtimeStreamsVersion("v2")).toBe("v1"); + it("keeps a v2 default on v2 when only per-org basins are configured", () => { + expect( + resolveRealtimeStreamsVersion(undefined, { ...PER_ORG_BASINS, defaultVersion: "v2" }) + ).toBe("v2"); + }); + + it("requires credentials, not just a basin", () => { + const basinOnly = { ...NO_S2, basin: "a-basin", defaultVersion: "v2" as const }; + expect(resolveRealtimeStreamsVersion(undefined, basinOnly)).toBe("v1"); + expect(resolveRealtimeStreamsVersion("v2", basinOnly)).toBe("v1"); + }); + + it("requires a basin, not just credentials", () => { + const tokenOnly = { ...NO_S2, accessToken: "a-token", defaultVersion: "v2" as const }; + expect(resolveRealtimeStreamsVersion(undefined, tokenOnly)).toBe("v1"); + expect(resolveRealtimeStreamsVersion("v2", tokenOnly)).toBe("v1"); }); it("keeps an explicit v1 on v1 even where S2 is available", () => { - configureS2(); - envMock.REALTIME_STREAMS_DEFAULT_VERSION = "v2"; - expect(determineRealtimeStreamsVersion("v1")).toBe("v1"); + expect(resolveRealtimeStreamsVersion("v1", { ...GLOBAL_BASIN, defaultVersion: "v2" })).toBe( + "v1" + ); }); it("treats an unrecognised version as v1", () => { - configureS2(); - expect(determineRealtimeStreamsVersion("v3")).toBe("v1"); + expect(resolveRealtimeStreamsVersion("v3", GLOBAL_BASIN)).toBe("v1"); }); });