Skip to content
File

Blob: tests/worker/dispatch/queue/execute-run-steps.test.ts

typescript390 lines
1import { describe, expect, it, vi } from "vitest";
2 
3import { UnixTimestampMs } from "@/contracts";
4import { expectTrusted, type RepoConfig } from "@/worker/contracts";
5import {
6 type PreparedExecutionEnvironment,
7 type RunExecutionContext,
8} from "@/worker/dispatch/shared/run-execution-context";
9import { executeRunSteps } from "@/worker/dispatch/shared/run-steps";
10import { createLogBatcher } from "@/worker/dispatch/shared/run-steps/logging";
11 
12import {
13 createQueueLeaseStub,
14 createQueueRunLogsStub,
15 createQueueRunStoreStub,
16 createQueueScope,
17 createQueueState,
18} from "../../../helpers/dispatch/shared";
19 
20type LogBatcherContext = Parameters<typeof createLogBatcher>[0];
21type StepContext = Parameters<typeof executeRunSteps>[0];
22 
23const makeScope = (): RunExecutionContext["scope"] =>
24 createQueueScope({
25 startedAt: expectTrusted(UnixTimestampMs, Date.now(), "UnixTimestampMs"),
26 });
27 
28const makeState = (): RunExecutionContext["state"] => createQueueState();
29 
30const makePrepared = (): PreparedExecutionEnvironment => ({
31 repoConfig: {
32 version: 1,
33 checkout: {
34 depth: 1,
35 },
36 run: {
37 workingDirectory: ".",
38 timeoutSeconds: 60,
39 steps: [
40 {
41 name: "test",
42 run: "npm test",
43 },
44 ],
45 },
46 } as RepoConfig,
47 workingDirectory: "/workspace/repo",
48});
49 
50const makePreparedWithSteps = (...commands: string[]): PreparedExecutionEnvironment => ({
51 repoConfig: {
52 version: 1,
53 checkout: {
54 depth: 1,
55 },
56 run: {
57 workingDirectory: ".",
58 timeoutSeconds: 60,
59 steps: commands.map((command, index) => ({
60 name: `step-${index + 1}`,
61 run: command,
62 })),
63 },
64 } as RepoConfig,
65 workingDirectory: "/workspace/repo",
66});
67 
68const toSseStream = (...events: string[]): ReadableStream<Uint8Array> =>
69 new ReadableStream({
70 start(controller) {
71 const encoder = new TextEncoder();
72 for (const event of events) {
73 controller.enqueue(encoder.encode(event));
74 }
75 controller.close();
76 },
77 });
78 
79const createExecStream = (
80 ...events: Array<
81 | { type: "start"; command?: string; pid?: number }
82 | { type: "stdout" | "stderr"; data: string }
83 | { type: "complete"; exitCode: number }
84 | { type: "error"; error: string }
85 >
86): ReadableStream<Uint8Array> =>
87 toSseStream(
88 ...events.map((event) => `data: ${JSON.stringify({ timestamp: "2026-03-24T12:00:00.000Z", ...event })}\n\n`),
89 );
90 
91describe("queue execute run steps", () => {
92 it("flushes buffered logs through the bound run store", async () => {
93 const runStore = createQueueRunStoreStub({
94 appendLogs: vi.fn(async () => {}),
95 });
96 const batcher = createLogBatcher({
97 runStore,
98 scope: makeScope(),
99 } satisfies LogBatcherContext);
100 
101 batcher.push("stdout", "x".repeat(4096));
102 await batcher.flush();
103 
104 expect(runStore.appendLogs).toHaveBeenCalledWith([
105 expect.objectContaining({
106 stream: "stdout",
107 chunk: "x".repeat(4096),
108 }),
109 ]);
110 });
111 
112 it("summarizes a failed step while preserving stderr logs", async () => {
113 const scope = makeScope();
114 const state = makeState();
115 const runStore = createQueueRunStoreStub({
116 updateState: vi.fn(async () => {}),
117 updateStepState: vi.fn(async () => {}),
118 appendLogs: vi.fn(async () => {}),
119 });
120 const logs = createQueueRunLogsStub({
121 redactMessage: vi.fn((message: string) => `redacted:${message}`),
122 });
123 const session = {
124 execStream: vi.fn(async () =>
125 createExecStream(
126 { type: "stdout", data: "running build\n" },
127 { type: "stderr", data: "boom" },
128 { type: "complete", exitCode: 5 },
129 ),
130 ),
131 };
132 state.session = session as never;
133 
134 const context = {
135 scope,
136 state,
137 runStore,
138 logs,
139 } satisfies StepContext;
140 
141 const outcome = await executeRunSteps(context, createQueueLeaseStub(), makePrepared());
142 
143 expect(outcome).toEqual({
144 kind: "failed",
145 exitCode: 5,
146 errorMessage: 'redacted:Step "test" failed with exit code 5.',
147 });
148 expect(runStore.updateState).toHaveBeenCalledTimes(2);
149 expect(runStore.updateStepState).toHaveBeenNthCalledWith(
150 1,
151 expect.objectContaining({
152 status: "running",
153 }),
154 );
155 expect(runStore.updateStepState).toHaveBeenNthCalledWith(
156 2,
157 expect.objectContaining({
158 status: "failed",
159 exitCode: 5,
160 }),
161 );
162 expect(runStore.appendLogs).toHaveBeenCalledTimes(1);
163 expect(runStore.appendLogs).toHaveBeenCalledWith([
164 expect.objectContaining({
165 stream: "stdout",
166 chunk: "running build\n",
167 }),
168 expect.objectContaining({
169 stream: "stderr",
170 chunk: "boom",
171 }),
172 ]);
173 expect(logs.redactMessage).toHaveBeenCalledWith('Step "test" failed with exit code 5.');
174 expect(session.execStream).toHaveBeenCalledWith(
175 "npm test",
176 expect.objectContaining({
177 cwd: "/workspace/repo",
178 timeout: expect.any(Number),
179 }),
180 );
181 });
182 
183 it("passes a step when execStream emits a successful complete event", async () => {
184 const scope = makeScope();
185 const state = makeState();
186 const runStore = createQueueRunStoreStub({
187 updateState: vi.fn(async () => {}),
188 updateStepState: vi.fn(async () => {}),
189 appendLogs: vi.fn(async () => {}),
190 });
191 const session = {
192 execStream: vi.fn(async () =>
193 createExecStream(
194 { type: "stdout", data: "hello\n" },
195 { type: "stderr", data: "warn\n" },
196 { type: "complete", exitCode: 0 },
197 ),
198 ),
199 };
200 state.session = session as never;
201 
202 const context = {
203 scope,
204 state,
205 runStore,
206 logs: createQueueRunLogsStub(),
207 } satisfies StepContext;
208 
209 const outcome = await executeRunSteps(context, createQueueLeaseStub(), makePrepared());
210 
211 expect(outcome).toEqual({
212 kind: "passed",
213 exitCode: 0,
214 });
215 expect(runStore.updateStepState).toHaveBeenNthCalledWith(
216 2,
217 expect.objectContaining({
218 status: "passed",
219 exitCode: 0,
220 }),
221 );
222 });
223 
224 it("tracks and clears the live streamed process when the start event includes a pid", async () => {
225 const scope = makeScope();
226 const state = makeState();
227 const runStore = createQueueRunStoreStub({
228 updateState: vi.fn(async () => {}),
229 updateStepState: vi.fn(async () => {}),
230 appendLogs: vi.fn(async () => {}),
231 });
232 const liveProcess = {
233 id: "proc-1",
234 pid: 4242,
235 command: "npm test",
236 status: "running",
237 startTime: new Date("2026-03-24T12:00:00.000Z"),
238 };
239 const session = {
240 execStream: vi.fn(async () =>
241 createExecStream(
242 { type: "start", pid: 4242 },
243 { type: "stdout", data: "hello\n" },
244 { type: "complete", exitCode: 0 },
245 ),
246 ),
247 listProcesses: vi.fn(async () => [liveProcess]),
248 };
249 state.session = session as never;
250 
251 const context = {
252 scope,
253 state,
254 runStore,
255 logs: createQueueRunLogsStub(),
256 } satisfies StepContext;
257 
258 const outcome = await executeRunSteps(context, createQueueLeaseStub(), makePrepared());
259 
260 expect(outcome).toEqual({
261 kind: "passed",
262 exitCode: 0,
263 });
264 expect(session.listProcesses).toHaveBeenCalledTimes(1);
265 expect(state.currentProcess).toBeNull();
266 });
267 
268 it("replays repo-defined commands from the first step even when currentStepPosition is already set", async () => {
269 const scope = makeScope();
270 const state = makeState();
271 state.currentStepPosition = 2 as RunExecutionContext["state"]["currentStepPosition"];
272 const runStore = createQueueRunStoreStub({
273 updateState: vi.fn(async () => {}),
274 updateStepState: vi.fn(async () => {}),
275 appendLogs: vi.fn(async () => {}),
276 });
277 const execStream = vi.fn(async (_command: string, _options?: unknown) =>
278 createExecStream({ type: "complete", exitCode: 0 }),
279 );
280 const session = {
281 execStream,
282 };
283 state.session = session as never;
284 
285 const context = {
286 scope,
287 state,
288 runStore,
289 logs: createQueueRunLogsStub(),
290 } satisfies StepContext;
291 
292 const outcome = await executeRunSteps(
293 context,
294 createQueueLeaseStub(),
295 makePreparedWithSteps("echo first", "echo second"),
296 );
297 
298 expect(outcome).toEqual({
299 kind: "passed",
300 exitCode: 0,
301 });
302 expect(execStream.mock.calls.map(([command]) => command)).toEqual(["echo first", "echo second"]);
303 expect(runStore.updateStepState).toHaveBeenNthCalledWith(
304 1,
305 expect.objectContaining({
306 position: 1,
307 status: "running",
308 }),
309 );
310 });
311 
312 it("fails a step when execStream emits a terminal error event", async () => {
313 const scope = makeScope();
314 const state = makeState();
315 const runStore = createQueueRunStoreStub({
316 updateState: vi.fn(async () => {}),
317 updateStepState: vi.fn(async () => {}),
318 appendLogs: vi.fn(async () => {}),
319 });
320 const logs = createQueueRunLogsStub({
321 redactMessage: vi.fn((message: string) => `redacted:${message}`),
322 });
323 const session = {
324 execStream: vi.fn(async () =>
325 createExecStream({ type: "stdout", data: "still running\n" }, { type: "error", error: "stream failed" }),
326 ),
327 };
328 state.session = session as never;
329 
330 const context = {
331 scope,
332 state,
333 runStore,
334 logs,
335 } satisfies StepContext;
336 
337 const outcome = await executeRunSteps(context, createQueueLeaseStub(), makePrepared());
338 
339 expect(outcome).toMatchObject({
340 kind: "failed",
341 exitCode: 1,
342 });
343 expect(logs.redactMessage).toHaveBeenCalledWith(expect.stringContaining("stream failed"));
344 expect(runStore.updateStepState).toHaveBeenNthCalledWith(
345 2,
346 expect.objectContaining({
347 status: "failed",
348 }),
349 );
350 });
351 
352 it("prefers the abnormal stream error over stderr when the stream never completes", async () => {
353 const scope = makeScope();
354 const state = makeState();
355 const runStore = createQueueRunStoreStub({
356 updateState: vi.fn(async () => {}),
357 updateStepState: vi.fn(async () => {}),
358 appendLogs: vi.fn(async () => {}),
359 });
360 const logs = createQueueRunLogsStub({
361 redactMessage: vi.fn((message: string) => `redacted:${message}`),
362 });
363 const session = {
364 execStream: vi.fn(async () =>
365 createExecStream(
366 { type: "stderr", data: "npm warn deprecated something\n" },
367 { type: "stderr", data: "npm warn deprecated else\n" },
368 ),
369 ),
370 };
371 state.session = session as never;
372 
373 const context = {
374 scope,
375 state,
376 runStore,
377 logs,
378 } satisfies StepContext;
379 
380 const outcome = await executeRunSteps(context, createQueueLeaseStub(), makePrepared());
381 
382 expect(outcome).toEqual({
383 kind: "failed",
384 exitCode: null,
385 errorMessage: "redacted:Command stream ended without a terminal event.",
386 });
387 expect(logs.redactMessage).toHaveBeenCalledWith("Command stream ended without a terminal event.");
388 });
389});