Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .server-changes/session-run-streams-version.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,6 @@
---
area: webapp
type: fix
---

Run-scoped realtime streams written inside a chat session run now use the same streams backend as the session itself, instead of falling back to the older one.
8 changes: 8 additions & 0 deletions apps/webapp/app/services/realtime/sessionRunManager.server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import { logger } from "~/services/logger.server";
import { CancelTaskRunService } from "~/v3/services/cancelTaskRun.server";
import { TriggerTaskService } from "~/v3/services/triggerTask.server";
import { isFinalRunStatus } from "~/v3/taskStatus";
import { determineRealtimeStreamsVersion } from "./v1StreamsGlobal.server";

/**
* Schema for `Session.triggerConfig` (stored as JSONB). The wire-format
Expand Down Expand Up @@ -275,6 +276,12 @@ export async function ensureRunForSession(
* Trigger a single run for a session. Builds `TriggerTaskRequestBody`
* by shallow-merging `payloadOverrides` over `config.basePayload` and
* threading `config`'s machine/queue/tags through the trigger options.
*
* A session's own channels are always v2, so the run is stamped to match
* rather than inheriting the `realtimeStreamsVersion` column default. Without
* this, run-scoped `streams.*` calls inside a session run resolve to v1 while
* the session it belongs to is on v2. `determineRealtimeStreamsVersion`
* degrades to v1 where v2 streams are not configured.
*/
async function triggerSessionRun(params: {
session: Pick<Session, "id" | "taskIdentifier">;
Expand Down Expand Up @@ -310,6 +317,7 @@ async function triggerSessionRun(params: {
const result = await service.call(session.taskIdentifier, environment, body, {
triggerSource: "session",
triggerAction: "trigger",
realtimeStreamsVersion: determineRealtimeStreamsVersion("v2"),
Comment thread
matt-aitken marked this conversation as resolved.
Comment thread
matt-aitken marked this conversation as resolved.
Comment thread
matt-aitken marked this conversation as resolved.
});

if (!result) {
Expand Down
1 change: 1 addition & 0 deletions apps/webapp/test/realtimeServices.replicaLag.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -412,6 +412,7 @@ describe("realtime-svc — replica-lag guards", () => {
// previousRunId forwarded to the triggered run is the calling run's cuid (documented fallback).
expect(triggerState.calls).toHaveLength(1);
expect(triggerState.calls[0]!.body.payload.previousRunId).toBe(callingRunId);
expect(triggerState.calls[0]!.options.realtimeStreamsVersion).toBeDefined();
expect(replica.wasHit("taskRun")).toBe(true);

// Proof the null was lag-induced: the primary holds the resolvable friendlyId (≠ the cuid).
Expand Down
147 changes: 147 additions & 0 deletions apps/webapp/test/sessionRunStreamsBackend.e2e.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,147 @@
/**
* Full-stack e2e for which realtime streams backend a Session's run lands on.
*
* Boots the real webapp + Postgres + Redis + s2-lite (via
* startSessionStreamTestServer), creates a Session through the public API so
* the run is triggered by the real `sessionRunManager` path, then appends to a
* run-scoped stream exactly as `streams.append()` does and checks where the
* bytes actually went.
*
* The harness starts the webapp with `REALTIME_STREAMS_DEFAULT_VERSION: "v2"`
* and a live S2, so a run landing on v1 here is not a configuration gap. It
* means the trigger path never asked, and fell through to the
* `realtimeStreamsVersion` column default.
*
* Requires a pre-built webapp: pnpm run build --filter webapp
*/
import { randomBytes } from "crypto";
import Redis from "ioredis";
import { afterAll, beforeAll, describe, expect, it, vi } from "vitest";
import type { SessionStreamTestServer } from "@internal/testcontainers/webapp";
import { startSessionStreamTestServer } from "@internal/testcontainers/webapp";
import { seedTestEnvironment } from "./helpers/seedTestEnvironment";

vi.setConfig({ testTimeout: 120_000, hookTimeout: 180_000 });

let server: SessionStreamTestServer;

beforeAll(async () => {
server = await startSessionStreamTestServer();
}, 180_000);

afterAll(async () => {
await server?.stop();
}, 120_000);

const STREAM_ID = "frames";
const PART_ID = "part";
const FRAME_BYTES = 250 * 1024;
const FRAME_COUNT = 8;

/** Mirrors `S2RealtimeStreams.toStreamName` on the shared-basin prefix. */
function runStreamName(p: {
orgId: string;
envSlug: string;
envId: string;
runId: string;
streamId: string;
}): string {
return `org/${p.orgId}/env/${p.envSlug}/${p.envId}/runs/${p.runId}/${p.streamId}`;
}

/** Mirrors the `keyPrefix` + key shape in `v1StreamsGlobal` / `RedisRealtimeStreams`. */
function redisStreamKey(runId: string, streamId: string): string {
return `tr:realtime:streams:stream:${runId}:${streamId}`;
}

function framesFound(body: string): number {
return Array.from({ length: FRAME_COUNT }, (_, i) => `${PART_ID}-${i}`).filter((id) =>
body.includes(id)
).length;
}

async function s2Body(streamName: string): Promise<string> {
const qs = new URLSearchParams({ seq_num: "0", clamp: "true", wait: "0" });
const res = await fetch(
`${server.s2.endpoint}/v1/streams/${encodeURIComponent(streamName)}/records?${qs}`,
{
headers: {
Authorization: "Bearer ignored",
Accept: "text/event-stream",
"S2-Format": "raw",
"S2-Basin": server.s2.basin,
},
}
);

if (res.status === 404) return "";
expect(res.ok).toBe(true);

return res.text();
}

describe("session runs and the realtime streams backend", () => {
it("stamps the run v2 and routes a run-scoped stream to S2, not Redis", async () => {
const { organization, environment, apiKey } = await seedTestEnvironment(server.prisma);

const createRes = await fetch(`${server.webapp.baseUrl}/api/v1/sessions`, {
method: "POST",
headers: { Authorization: `Bearer ${apiKey}`, "Content-Type": "application/json" },
body: JSON.stringify({
type: "chat.agent",
externalId: `e2e-${randomBytes(6).toString("hex")}`,
taskIdentifier: "e2e-browser-agent",
triggerConfig: { basePayload: {} },
}),
});

expect(createRes.ok).toBe(true);
const created = (await createRes.json()) as { runId: string };
expect(created.runId).toBeTruthy();

const run = await server.prisma.taskRun.findFirstOrThrow({
where: { friendlyId: created.runId },
select: { friendlyId: true, realtimeStreamsVersion: true, streamBasinName: true },
});

const appendStatuses: number[] = [];
for (let i = 0; i < FRAME_COUNT; i++) {
const res = await fetch(
`${server.webapp.baseUrl}/realtime/v1/streams/${created.runId}/self/${STREAM_ID}/append`,
{
method: "POST",
headers: {
Authorization: `Bearer ${apiKey}`,
"Content-Type": "text/plain",
"X-Part-Id": `${PART_ID}-${i}`,
},
body: JSON.stringify({ i, frame: "a".repeat(FRAME_BYTES) }),
}
);
appendStatuses.push(res.status);
}

expect(appendStatuses).toEqual(Array.from({ length: FRAME_COUNT }, () => 200));

const streamName = runStreamName({
orgId: organization.id,
envSlug: environment.slug,
envId: environment.id,
runId: created.runId,
streamId: STREAM_ID,
});
const redis = new Redis({ host: server.redis.host, port: server.redis.port });
let observed: { version: string; framesInS2: number; keyInRedis: boolean };
try {
observed = {
version: run.realtimeStreamsVersion,
framesInS2: framesFound(await s2Body(streamName)),
keyInRedis: (await redis.exists(redisStreamKey(created.runId, STREAM_ID))) === 1,
};
} finally {
redis.disconnect();
}

expect(observed).toEqual({ version: "v2", framesInS2: FRAME_COUNT, keyInRedis: false });
});
});
16 changes: 15 additions & 1 deletion internal-packages/testcontainers/src/webapp.ts
Original file line number Diff line number Diff line change
Expand Up @@ -272,6 +272,12 @@ export type { StartedS2Container } from "./s2";
export interface SessionStreamTestServer extends TestServer {
s2: StartedS2Container;
minio: StartedMinIOContainer;
/**
* Mapped connection for the same Redis the webapp under test uses. Lets a
* test assert which backend a stream actually landed on, rather than
* inferring it from the absence of records in S2.
*/
redis: { host: string; port: number };
}

/**
Expand Down Expand Up @@ -348,5 +354,13 @@ export async function startSessionStreamTestServer(): Promise<SessionStreamTestS
await network.stop().catch((err) => console.error("network.stop failed:", err));
};

return { webapp, prisma: prisma!, databaseUrl: pgUrl!, s2: s2!, minio: minio!, stop };
return {
webapp,
prisma: prisma!,
databaseUrl: pgUrl!,
s2: s2!,
minio: minio!,
redis: { host: redisContainer!.getHost(), port: redisContainer!.getPort() },
stop,
};
}
Loading