返回 CodeWhale
fleet-client.test.mjs
根目录 / npm / runtime-sdk / test / fleet-client.test.mjs
1 import assert from "node:assert/strict";
2 import { once } from "node:events";
3 import { createServer } from "node:http";
4 import test from "node:test";
5 import {
6 CodeWhaleRuntimeClient,
7 RuntimeApiError,
8 RuntimeCapabilityError,
9 createRuntimeClient,
10 } from "../index.js";
11
12 function jsonResponse(body, init = {}) {
13 return new Response(JSON.stringify(body), {
14 status: init.status ?? 200,
15 headers: { "content-type": "application/json", ...(init.headers ?? {}) },
16 });
17 }
18
19 function fakeFetch(responseFactory) {
20 const calls = [];
21 const fetch = async (url, init) => {
22 calls.push({ url: url.toString(), init });
23 return responseFactory(url, init, calls.length);
24 };
25 fetch.calls = calls;
26 return fetch;
27 }
28
29 test("createRuntimeClient returns CodeWhaleRuntimeClient instance", () => {
30 const fetch = fakeFetch(() => jsonResponse({}));
31 const client = createRuntimeClient({ fetch });
32
33 assert.ok(client instanceof CodeWhaleRuntimeClient);
34 assert.equal(client.baseUrl, "http://127.0.0.1:7878/");
35 });
36
37 test("listFleetRuns calls the Runtime API with bearer auth", async () => {
38 const fetch = fakeFetch(() =>
39 jsonResponse({
40 status: { runs: 1, workers: {} },
41 runs: [{ id: "run-1", name: "smoke", tasks: [], labels: {} }],
42 }),
43 );
44 const client = createRuntimeClient({
45 baseUrl: "http://127.0.0.1:7878",
46 token: "token-1",
47 fetch,
48 });
49
50 const response = await client.listFleetRuns();
51
52 assert.equal(response.runs[0].id, "run-1");
53 assert.equal(fetch.calls[0].url, "http://127.0.0.1:7878/v1/fleet/runs");
54 assert.equal(fetch.calls[0].init.method, "GET");
55 assert.equal(fetch.calls[0].init.headers.get("authorization"), "Bearer token-1");
56 assert.equal(fetch.calls[0].init.redirect, "error");
57 });
58
59 test("fleet event paths cannot send runtime credentials outside the configured origin", async () => {
60 const fetch = fakeFetch(() => jsonResponse({ events: [] }));
61 const client = createRuntimeClient({ baseUrl: "http://127.0.0.1:7878/", token: "fixture-token", fetch });
62 for (const path of [
63 "https://example.invalid/events", "//example.invalid/events", "\\\\example.invalid/events",
64 "http://127.0.0.1:7879/events", "https://127.0.0.1:7878/events",
65 "data:application/json,%7B%7D", "file:///tmp/events", "http://user:password@127.0.0.1:7878/events",
66 ]) {
67 await assert.rejects(client.fleetEvents("run-1", { path }).next(), /configured HTTP\(S\) origin/);
68 }
69 assert.equal(fetch.calls.length, 0, "no rejected destination reaches fetch");
70 });
71
72 test("fleet event paths preserve safe relative and same-origin URL semantics", async () => {
73 const fetch = fakeFetch(() => jsonResponse({ events: [] }));
74 const client = createRuntimeClient({ baseUrl: "http://127.0.0.1:7878/runtime/", token: "fixture-token", fetch });
75 for (const path of ["events", "/events", "http://127.0.0.1:7878/events", "//127.0.0.1:7878/events"]) {
76 assert.deepEqual(await client.fleetEvents("run-1", { path, after: "cursor/a", limit: 5 }).next(), { value: undefined, done: true });
77 }
78 assert.deepEqual(fetch.calls.map(({ url }) => new URL(url).pathname), ["/runtime/events", "/events", "/events", "/events"]);
79 for (const { url, init } of fetch.calls) {
80 assert.equal(new URL(url).searchParams.get("after"), "cursor/a");
81 assert.equal(new URL(url).searchParams.get("limit"), "5");
82 assert.equal(init.headers.get("authorization"), "Bearer fixture-token");
83 assert.equal(init.redirect, "error");
84 }
85 });
86
87 test("authenticated runtime GET and POST reject actual redirects before a second request", async (t) => {
88 const escaped = [];
89 const destination = createServer((req, res) => {
90 escaped.push({ url: req.url, authorization: req.headers.authorization });
91 res.writeHead(200, { "content-type": "application/json" }).end("{}");
92 });
93 destination.listen(0, "127.0.0.1");
94 await once(destination, "listening");
95 t.after(() => new Promise(resolve => { destination.close(resolve); destination.closeAllConnections(); }));
96 const received = [];
97 const source = createServer((req, res) => {
98 received.push({ method: req.method, url: req.url, authorization: req.headers.authorization });
99 if (req.url !== "/v1/fleet/runs") {
100 escaped.push({ url: req.url, authorization: req.headers.authorization });
101 res.writeHead(200, { "content-type": "application/json" }).end("{}");
102 return;
103 }
104 res.writeHead(307, { location: req.method === "GET" ? "/unexpected" : `http://127.0.0.1:${destination.address().port}/unexpected` }).end();
105 });
106 source.listen(0, "127.0.0.1");
107 await once(source, "listening");
108 t.after(() => new Promise(resolve => { source.close(resolve); source.closeAllConnections(); }));
109 const client = createRuntimeClient({ baseUrl: `http://127.0.0.1:${source.address().port}`, token: "fixture-token" });
110 await assert.rejects(client.listFleetRuns(), TypeError);
111 await assert.rejects(client.createFleetRun({ name: "local-only" }), TypeError);
112 assert.deepEqual(received.map(({ method }) => method), ["GET", "POST"]);
113 assert.ok(received.every(({ authorization }) => authorization === "Bearer fixture-token"));
114 assert.deepEqual(escaped, [], "neither a same-origin nor an off-origin redirect is followed");
115 });
116
117 test("worker and run actions use POST endpoints", async () => {
118 const fetch = fakeFetch((url) =>
119 jsonResponse(
120 url.pathname.endsWith("/stop")
121 ? {
122 action: "stop",
123 run_id: "run-1",
124 stopped: 1,
125 status: { runs: 1, workers: {} },
126 }
127 : {
128 action: url.pathname.endsWith("/restart")
129 ? "restart"
130 : url.pathname.endsWith("/stop")
131 ? "stop"
132 : "interrupt",
133 worker: { worker_id: "w1", artifacts: [] },
134 },
135 ),
136 );
137 const client = new CodeWhaleRuntimeClient({ fetch });
138
139 await client.interruptWorker("w1");
140 await client.stopWorker("w1");
141 await client.restartWorker("w1");
142 await client.startFleetRun("run-1");
143 await client.stopFleetRun("run-1");
144
145 assert.deepEqual(
146 fetch.calls.map((call) => [new URL(call.url).pathname, call.init.method]),
147 [
148 ["/v1/fleet/workers/w1/interrupt", "POST"],
149 ["/v1/fleet/workers/w1/stop", "POST"],
150 ["/v1/fleet/workers/w1/restart", "POST"],
151 ["/v1/fleet/runs/run-1/start", "POST"],
152 ["/v1/fleet/runs/run-1/stop", "POST"],
153 ],
154 );
155 });
156
157 test("managed Fleet helpers send explicit launch metadata and reconnect cursors", async () => {
158 const fetch = fakeFetch((url) =>
159 jsonResponse(
160 url.pathname.endsWith("/events/replay")
161 ? { run_id: "run-1", events: [], has_more: false, history_truncated: false }
162 : { execution: "awaiting_start", run: { id: "run-1" }, warnings: [] },
163 ),
164 );
165 const client = new CodeWhaleRuntimeClient({ fetch });
166 const spec = {
167 target: "this_computer",
168 roles: [{ name: "reviewer" }],
169 workflow: {
170 id: "review",
171 kind: "parallel",
172 tasks: [
173 {
174 id: "review",
175 name: "Review",
176 instructions: "Review.",
177 worker: { role: "reviewer" },
178 budget: { max_steps: 0 },
179 },
180 ],
181 },
182 };
183
184 await client.createFleetRun(spec);
185 await client.replayFleetEvents("run-1", { after: "fev1_cursor_worker", limit: 25 });
186
187 assert.deepEqual(JSON.parse(fetch.calls[0].init.body), spec);
188 assert.equal(JSON.parse(fetch.calls[0].init.body).workflow.tasks[0].budget.max_steps, 0);
189 const replayUrl = new URL(fetch.calls[1].url);
190 assert.equal(replayUrl.pathname, "/v1/fleet/runs/run-1/events/replay");
191 assert.equal(replayUrl.searchParams.get("after"), "fev1_cursor_worker");
192 assert.equal(replayUrl.searchParams.get("limit"), "25");
193 });
194
195 test("unsupported fleet capabilities raise typed errors", async () => {
196 const fetch = fakeFetch(() => jsonResponse({ error: "not found" }, { status: 404 }));
197 const client = new CodeWhaleRuntimeClient({ fetch });
198
199 await assert.rejects(
200 () => client.createFleetRun({ name: "future" }),
201 (error) =>
202 error instanceof RuntimeCapabilityError &&
203 error.capability === "fleet_run_create" &&
204 error.status === 404,
205 );
206
207 await assert.rejects(
208 async () => {
209 for await (const _event of client.fleetEvents("run-1")) {
210 throw new Error("unexpected event");
211 }
212 },
213 (error) =>
214 error instanceof RuntimeCapabilityError &&
215 error.capability === "fleet_event_stream" &&
216 error.status === 404,
217 );
218 });
219
220 test("fleetEvents can replay JSON event fixtures when the API exposes them", async () => {
221 const fetch = fakeFetch(() =>
222 jsonResponse({
223 events: [
224 {
225 seq: 1,
226 run_id: "run-1",
227 worker_id: "w1",
228 task_id: "task-1",
229 timestamp: "2026-06-13T00:00:00Z",
230 label: "running",
231 payload: { state: "running" },
232 },
233 ],
234 }),
235 );
236 const client = new CodeWhaleRuntimeClient({ fetch });
237
238 const events = [];
239 for await (const event of client.fleetEvents("run-1", { path: "/v1/fleet/runs/run-1/events" })) {
240 events.push(event);
241 }
242
243 assert.equal(events.length, 1);
244 assert.equal(events[0].payload.state, "running");
245 });
246
247 test("fleetEvents parses text/event-stream frames", async () => {
248 const encoder = new TextEncoder();
249 const body = new ReadableStream({
250 start(controller) {
251 controller.enqueue(
252 encoder.encode(
253 'id: fev1_heartbeat_worker\nevent: fleet.worker.heartbeat\ndata: {"cursor":"fev1_heartbeat_worker","event":"fleet.worker.heartbeat","run_id":"run-1","worker_id":"w1","task_id":"task-1","timestamp":"2026-06-13T00:00:01Z","worker_seq":2,"payload":{"state":"heartbeat","memory_mb":128}}\n\n',
254 ),
255 );
256 controller.close();
257 },
258 });
259 const fetch = fakeFetch(
260 () =>
261 new Response(body, {
262 status: 200,
263 headers: { "content-type": "text/event-stream" },
264 }),
265 );
266 const client = new CodeWhaleRuntimeClient({ fetch });
267
268 const events = [];
269 for await (const event of client.fleetEvents("run-1", { after: "fev1_previous", limit: 10 })) {
270 events.push(event);
271 }
272
273 assert.equal(events.length, 1);
274 assert.equal(events[0].payload.state, "heartbeat");
275 assert.equal(events[0].payload.memory_mb, 128);
276 const eventUrl = new URL(fetch.calls[0].url);
277 assert.equal(eventUrl.searchParams.get("after"), "fev1_previous");
278 assert.equal(eventUrl.searchParams.get("limit"), "10");
279 assert.equal(fetch.calls[0].init.headers.get("accept"), "text/event-stream");
280 });
281
282 test("fleetEvents preserves SSE control event names", async () => {
283 const encoder = new TextEncoder();
284 const body = new ReadableStream({
285 start(controller) {
286 controller.enqueue(
287 encoder.encode(
288 'event: fleet.replay.cursor_unavailable\r\ndata: {"run_id":"run-1","reload_projection":true}\r\n\r\n',
289 ),
290 );
291 controller.close();
292 },
293 });
294 const client = new CodeWhaleRuntimeClient({
295 fetch: fakeFetch(
296 () =>
297 new Response(body, {
298 status: 200,
299 headers: { "content-type": "text/event-stream" },
300 }),
301 ),
302 });
303
304 const events = [];
305 for await (const event of client.fleetEvents("run-1")) {
306 events.push(event);
307 }
308
309 assert.deepEqual(events, [
310 {
311 event: "fleet.replay.cursor_unavailable",
312 run_id: "run-1",
313 reload_projection: true,
314 },
315 ]);
316 });
317
318 test("ordinary HTTP errors remain RuntimeApiError", async () => {
319 const fetch = fakeFetch(() => jsonResponse({ error: "bad" }, { status: 500 }));
320 const client = new CodeWhaleRuntimeClient({ fetch });
321
322 await assert.rejects(
323 () => client.getFleetRun("run-1"),
324 (error) =>
325 error instanceof RuntimeApiError &&
326 !(error instanceof RuntimeCapabilityError) &&
327 error.status === 500,
328 );
329 });
330
330 lines Plain Text