fix(webapp): count per-org basins as S2 being configured

Gating v2 on the global basin alone would have degraded every run to v1 on a
deployment that provisions a basin per organization and sets no global one,
even though S2 is fully working there. `resolveStreamBasin` already resolves
run, session and organization basins ahead of the global setting, so either
source now satisfies the basin requirement.

Splits the pure resolver out from the env lookup so the version matrix can be
tested without reaching for `env.server`.
This commit is contained in:
Matt Aitken
2026-08-11 11:15:28 +01:00
parent 2eb19be5fc
commit dab295751b
3 changed files with 105 additions and 64 deletions
+1 -1
View File
@@ -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.
@@ -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(
@@ -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");
});
});