Skip to content

Commit eaa9050

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
fix(webapp): narrow realtime session credentials to one stream
Realtime session authorization is now limited to the requested channel and the operations it needs. Mono-RevId: 00b06dd590600adcb2fb2d6adbebde81fd3dd05f
1 parent 2d569f6 commit eaa9050

3 files changed

Lines changed: 262 additions & 63 deletions

File tree

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: fix
4+
---
5+
6+
Realtime session writers now receive authorization limited to the requested session channel.
Lines changed: 186 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,186 @@
1+
import { createCache, createMemoryStore, DefaultStatefulContext, Namespace } from "@internal/cache";
2+
import { createServer, type Server } from "node:http";
3+
import type { AddressInfo } from "node:net";
4+
import { afterAll, beforeAll, beforeEach, describe, expect, it } from "vitest";
5+
import { S2RealtimeStreams } from "./s2realtimeStreams.server";
6+
7+
const BASIN = "test-basin";
8+
const STREAM_PREFIX = "org/org_123/env/dev/env_123";
9+
10+
type IssueAccessTokenRequest = {
11+
id: string;
12+
scope: {
13+
basins: { exact: string };
14+
ops: string[];
15+
streams: { exact?: string; prefix?: string };
16+
};
17+
expires_at: string;
18+
auto_prefix_streams: boolean;
19+
};
20+
21+
let accountServer: Server;
22+
let accountUrl: string;
23+
let issuedRequests: IssueAccessTokenRequest[] = [];
24+
25+
beforeAll(async () => {
26+
accountServer = createServer(async (request, response) => {
27+
if (request.method !== "POST" || request.url !== "/v1/access-tokens") {
28+
response.writeHead(404).end();
29+
return;
30+
}
31+
32+
const chunks: Buffer[] = [];
33+
for await (const chunk of request) {
34+
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk));
35+
}
36+
issuedRequests.push(JSON.parse(Buffer.concat(chunks).toString("utf8")));
37+
38+
response.writeHead(201, { "content-type": "application/json" });
39+
response.end(JSON.stringify({ access_token: `issued-token-${issuedRequests.length}` }));
40+
});
41+
42+
await new Promise<void>((resolve) => accountServer.listen(0, "127.0.0.1", resolve));
43+
const address = accountServer.address() as AddressInfo;
44+
accountUrl = `http://127.0.0.1:${address.port}/v1`;
45+
});
46+
47+
afterAll(async () => {
48+
await new Promise<void>((resolve, reject) => {
49+
accountServer.close((error) => (error ? reject(error) : resolve()));
50+
});
51+
});
52+
53+
beforeEach(() => {
54+
issuedRequests = [];
55+
});
56+
57+
function createAccessTokenCache() {
58+
const context = new DefaultStatefulContext();
59+
return createCache({
60+
accessToken: new Namespace<string>(context, {
61+
stores: [createMemoryStore(100, 0)],
62+
fresh: 60_000,
63+
stale: 120_000,
64+
}),
65+
});
66+
}
67+
68+
function createStreams(options?: {
69+
cache?: ReturnType<typeof createAccessTokenCache>;
70+
skipAccessTokens?: boolean;
71+
accessToken?: string;
72+
}) {
73+
return new S2RealtimeStreams({
74+
basin: BASIN,
75+
accessToken: options?.accessToken ?? "account-token",
76+
streamPrefix: STREAM_PREFIX,
77+
accountUrl,
78+
basinUrl: `${accountUrl}/basins/{basin}`,
79+
cache: options?.cache,
80+
skipAccessTokens: options?.skipAccessTokens,
81+
});
82+
}
83+
84+
describe("S2RealtimeStreams access-token scopes", () => {
85+
it("scopes the ordinary session input initializer to its exact full stream name", async () => {
86+
const streams = createStreams();
87+
88+
const result = await streams.initializeSessionStream("session_123", "in");
89+
90+
const streamName = `${STREAM_PREFIX}/sessions/session_123/in`;
91+
expect(result.responseHeaders).toMatchObject({
92+
"X-S2-Access-Token": "issued-token-1",
93+
"X-S2-Stream-Name": streamName,
94+
"X-S2-Basin": BASIN,
95+
});
96+
expect(issuedRequests).toHaveLength(1);
97+
expect(issuedRequests[0]).toMatchObject({
98+
scope: {
99+
basins: { exact: BASIN },
100+
ops: ["append", "create-stream"],
101+
streams: { exact: streamName },
102+
},
103+
auto_prefix_streams: false,
104+
});
105+
expect(issuedRequests[0]!.scope.streams).not.toHaveProperty("prefix");
106+
});
107+
108+
it("scopes the named session input initializer to only that channel", async () => {
109+
const streams = createStreams();
110+
111+
const result = await streams.initializeSessionStream("chat-room", "in", "steering");
112+
113+
const streamName = `${STREAM_PREFIX}/sessions/chat-room/channels/steering/in`;
114+
expect(result.responseHeaders?.["X-S2-Stream-Name"]).toBe(streamName);
115+
expect(issuedRequests[0]).toMatchObject({
116+
scope: {
117+
ops: ["append", "create-stream"],
118+
streams: { exact: streamName },
119+
},
120+
auto_prefix_streams: false,
121+
});
122+
});
123+
124+
it("keeps trim authority on an exact private session output stream", async () => {
125+
const streams = createStreams();
126+
127+
const result = await streams.initializeSessionStream("session_123", "out", "updates");
128+
129+
const streamName = `${STREAM_PREFIX}/sessions/session_123/channels/updates/out`;
130+
expect(result.responseHeaders?.["X-S2-Stream-Name"]).toBe(streamName);
131+
expect(issuedRequests[0]).toMatchObject({
132+
scope: {
133+
ops: ["append", "create-stream", "trim"],
134+
streams: { exact: streamName },
135+
},
136+
auto_prefix_streams: false,
137+
});
138+
});
139+
140+
it("keeps the relative, auto-prefixed contract for run stream writers", async () => {
141+
const streams = createStreams();
142+
143+
const result = await streams.initializeStream("run_123", "progress");
144+
145+
expect(result.responseHeaders?.["X-S2-Stream-Name"]).toBe("/runs/run_123/progress");
146+
expect(issuedRequests[0]).toMatchObject({
147+
scope: {
148+
ops: ["append", "create-stream"],
149+
streams: { prefix: STREAM_PREFIX },
150+
},
151+
auto_prefix_streams: true,
152+
});
153+
});
154+
155+
it("caches each exact scope separately and ignores the previous broad cache namespace", async () => {
156+
const cache = createAccessTokenCache();
157+
const previousCacheKey = `${BASIN}:${STREAM_PREFIX}:append,create-stream,trim`;
158+
await cache.accessToken.set(previousCacheKey, "broad-token");
159+
const streams = createStreams({ cache });
160+
161+
const first = await streams.initializeSessionStream("session_123", "in");
162+
const repeated = await streams.initializeSessionStream("session_123", "in");
163+
const different = await streams.initializeSessionStream("session_456", "in");
164+
165+
expect(first.responseHeaders?.["X-S2-Access-Token"]).toBe("issued-token-1");
166+
expect(repeated.responseHeaders?.["X-S2-Access-Token"]).toBe("issued-token-1");
167+
expect(different.responseHeaders?.["X-S2-Access-Token"]).toBe("issued-token-2");
168+
expect(issuedRequests).toHaveLength(2);
169+
});
170+
171+
it("returns full stream names when token issuance is disabled", async () => {
172+
const streams = createStreams({ skipAccessTokens: true, accessToken: "" });
173+
174+
const ordinary = await streams.initializeSessionStream("session_123", "in");
175+
const named = await streams.initializeSessionStream("session_123", "in", "steering");
176+
177+
expect(ordinary.responseHeaders).toMatchObject({
178+
"X-S2-Access-Token": "s2-skip-access-tokens",
179+
"X-S2-Stream-Name": `${STREAM_PREFIX}/sessions/session_123/in`,
180+
});
181+
expect(named.responseHeaders?.["X-S2-Stream-Name"]).toBe(
182+
`${STREAM_PREFIX}/sessions/session_123/channels/steering/in`
183+
);
184+
expect(issuedRequests).toHaveLength(0);
185+
});
186+
});

0 commit comments

Comments
 (0)