diff --git a/packages/core/src/__tests__/loop.test.ts b/packages/core/src/__tests__/loop.test.ts new file mode 100644 index 0000000..a05c44b --- /dev/null +++ b/packages/core/src/__tests__/loop.test.ts @@ -0,0 +1,284 @@ +import { describe, it, expect, vi, afterEach } from "vitest"; +import { loop } from "../loop.js"; +import type { LoopStep } from "../types.js"; +import type { SessionBase, ToolResult } from "@codespar/types"; + +const ok = (tool: string, data: unknown = {}): ToolResult => ({ + success: true, + data, + error: null, + duration: 0, + server: "test", + tool, +}); + +const fail = (tool: string, error = "boom"): ToolResult => ({ + success: false, + data: null, + error, + duration: 0, + server: "test", + tool, +}); + +/** + * Minimal SessionBase whose execute() is driven by a per-call handler. + * loop() only ever calls session.execute, so the rest are inert stubs. + */ +function recordingSession( + handler: (tool: string, params: Record) => Promise | ToolResult, +): { session: SessionBase; calls: Array<{ tool: string; params: Record }> } { + const calls: Array<{ tool: string; params: Record }> = []; + const session: SessionBase = { + id: "ses_loop", + status: "active", + async execute(tool, params) { + calls.push({ tool, params }); + return handler(tool, params); + }, + async send() { + return { message: "", tool_calls: [], iterations: 0 }; + }, + async *sendStream() {}, + async connections() { + return []; + }, + async close() {}, + }; + return { session, calls }; +} + +afterEach(() => { + vi.useRealTimers(); + vi.restoreAllMocks(); +}); + +describe("loop() — happy path", () => { + it("runs all steps in order and reports completion", async () => { + const { session, calls } = recordingSession((tool) => ok(tool)); + const steps: LoopStep[] = [ + { tool: "a", params: {} }, + { tool: "b", params: {} }, + ]; + + const result = await loop(session, { steps }); + + expect(result.success).toBe(true); + expect(calls.map((c) => c.tool)).toEqual(["a", "b"]); + expect(result.results.map((r) => r.tool)).toEqual(["a", "b"]); + expect(result.completedSteps).toBe(2); + expect(result.totalSteps).toBe(2); + expect(result.duration).toBeGreaterThanOrEqual(0); + }); + + it("invokes onStepComplete per successful step with (step, result, index)", async () => { + const { session } = recordingSession((tool) => ok(tool)); + const onStepComplete = vi.fn(); + const steps: LoopStep[] = [ + { tool: "a", params: {} }, + { tool: "b", params: {} }, + ]; + + await loop(session, { steps, onStepComplete }); + + expect(onStepComplete).toHaveBeenCalledTimes(2); + expect(onStepComplete).toHaveBeenNthCalledWith(1, steps[0], expect.objectContaining({ tool: "a" }), 0); + expect(onStepComplete).toHaveBeenNthCalledWith(2, steps[1], expect.objectContaining({ tool: "b" }), 1); + }); +}); + +describe("loop() — dynamic params and conditional steps", () => { + it("calls step.params(prevResults) with prior results and forwards the return to execute", async () => { + const { session, calls } = recordingSession((tool) => ok(tool, { id: `${tool}-1` })); + const steps: LoopStep[] = [ + { tool: "create", params: { name: "Maria" } }, + { + tool: "charge", + params: (prev) => ({ ref: (prev[0]!.data as { id: string }).id }), + }, + ]; + + await loop(session, { steps }); + + expect(calls[1]).toEqual({ tool: "charge", params: { ref: "create-1" } }); + }); + + it("skips a step whose when() returns false", async () => { + const { session, calls } = recordingSession((tool) => ok(tool)); + const steps: LoopStep[] = [ + { tool: "a", params: {} }, + { tool: "b", params: {}, when: () => false }, + { tool: "c", params: {} }, + ]; + + const result = await loop(session, { steps }); + + expect(calls.map((c) => c.tool)).toEqual(["a", "c"]); + expect(result.results.map((r) => r.tool)).toEqual(["a", "c"]); + expect(result.totalSteps).toBe(3); + expect(result.completedSteps).toBe(2); + }); + + it("passes prevResults to when() so it can branch on earlier output", async () => { + const { session, calls } = recordingSession((tool) => ok(tool, { skip: tool === "a" })); + const steps: LoopStep[] = [ + { tool: "a", params: {} }, + { + tool: "b", + params: {}, + when: (prev) => !(prev[0]!.data as { skip: boolean }).skip, + }, + ]; + + await loop(session, { steps }); + + expect(calls.map((c) => c.tool)).toEqual(["a"]); + }); +}); + +describe("loop() — failure handling", () => { + it("aborts on first failure by default and reports it", async () => { + const onStepError = vi.fn(); + const { session, calls } = recordingSession((tool) => (tool === "b" ? fail(tool) : ok(tool))); + const steps: LoopStep[] = [ + { tool: "a", params: {} }, + { tool: "b", params: {} }, + { tool: "c", params: {} }, + ]; + + const result = await loop(session, { steps, onStepError }); + + expect(result.success).toBe(false); + expect(calls.map((c) => c.tool)).toEqual(["a", "b"]); + expect(result.completedSteps).toBe(1); + expect(result.totalSteps).toBe(3); + expect(onStepError).toHaveBeenCalledWith(steps[1], expect.any(Error), 1); + expect(onStepError.mock.calls[0]![1].message).toBe("boom"); + }); + + it("continues remaining steps when abortOnError is false", async () => { + const { session, calls } = recordingSession((tool) => (tool === "b" ? fail(tool) : ok(tool))); + const steps: LoopStep[] = [ + { tool: "a", params: {} }, + { tool: "b", params: {} }, + { tool: "c", params: {} }, + ]; + + const result = await loop(session, { steps, abortOnError: false }); + + expect(calls.map((c) => c.tool)).toEqual(["a", "b", "c"]); + expect(result.success).toBe(false); + expect(result.completedSteps).toBe(2); + expect(result.results).toHaveLength(3); + }); + + it("catches a thrown error from execute and surfaces it via onStepError", async () => { + const onStepError = vi.fn(); + const { session } = recordingSession((tool) => { + if (tool === "a") throw new Error("network down"); + return ok(tool); + }); + const steps: LoopStep[] = [{ tool: "a", params: {} }]; + + const result = await loop(session, { steps, onStepError }); + + expect(result.success).toBe(false); + expect(onStepError).toHaveBeenCalledWith(steps[0], expect.any(Error), 0); + expect(onStepError.mock.calls[0]![1].message).toBe("network down"); + }); +}); + +describe("loop() — retry policy", () => { + it("retries up to maxRetries then succeeds, with linear backoff delays", async () => { + vi.useFakeTimers(); + const delays: number[] = []; + vi.spyOn(globalThis, "setTimeout").mockImplementation(((cb: () => void, ms?: number) => { + delays.push(ms ?? 0); + cb(); + return 0 as unknown as ReturnType; + }) as typeof setTimeout); + + let attempts = 0; + const { session } = recordingSession((tool) => { + attempts += 1; + return attempts < 3 ? fail(tool) : ok(tool); + }); + + const result = await loop(session, { + steps: [{ tool: "a", params: {} }], + retryPolicy: { maxRetries: 3, backoff: "linear", baseDelay: 100 }, + }); + + expect(result.success).toBe(true); + expect(attempts).toBe(3); + // delay = baseDelay * (attempt + 1): 100 after attempt 0, 200 after attempt 1 + expect(delays).toEqual([100, 200]); + }); + + it("uses exponential backoff when configured", async () => { + vi.useFakeTimers(); + const delays: number[] = []; + vi.spyOn(globalThis, "setTimeout").mockImplementation(((cb: () => void, ms?: number) => { + delays.push(ms ?? 0); + cb(); + return 0 as unknown as ReturnType; + }) as typeof setTimeout); + + const { session } = recordingSession((tool) => fail(tool)); + + const result = await loop(session, { + steps: [{ tool: "a", params: {} }], + retryPolicy: { maxRetries: 3, backoff: "exponential", baseDelay: 50 }, + }); + + expect(result.success).toBe(false); + // delay = baseDelay * 2^attempt: 50, 100, 200 across attempts 0..2 (3 retries) + expect(delays).toEqual([50, 100, 200]); + }); + + it("does not retry when maxRetries is unset (single attempt)", async () => { + let attempts = 0; + const { session } = recordingSession((tool) => { + attempts += 1; + return fail(tool); + }); + + const result = await loop(session, { steps: [{ tool: "a", params: {} }] }); + + expect(attempts).toBe(1); + expect(result.success).toBe(false); + }); +}); + +describe("loop() — edge cases (current behavior, see docs/fix-core.md R4)", () => { + // CHARACTERIZATION OF A SUSPECTED BUG, not an endorsement. + // results.every() on an empty array is vacuously true, so a run where + // every step is skipped (or there are no steps) reports success:true + // with completedSteps:0. Documented in docs/fix-core.md (item R4 / C2 + // pattern). If loop() is later changed to treat "nothing ran" as a + // non-success, update these expectations deliberately. + it("reports success:true with zero steps (vacuous truth)", async () => { + const { session } = recordingSession((tool) => ok(tool)); + + const result = await loop(session, { steps: [] }); + + expect(result.success).toBe(true); + expect(result.completedSteps).toBe(0); + expect(result.totalSteps).toBe(0); + }); + + it("reports success:true when every step is skipped by when() (vacuous truth)", async () => { + const { session, calls } = recordingSession((tool) => ok(tool)); + const steps: LoopStep[] = [ + { tool: "a", params: {}, when: () => false }, + { tool: "b", params: {}, when: () => false }, + ]; + + const result = await loop(session, { steps }); + + expect(calls).toHaveLength(0); + expect(result.success).toBe(true); + expect(result.completedSteps).toBe(0); + expect(result.totalSteps).toBe(2); + }); +}); diff --git a/packages/core/src/__tests__/wait-for-connections.test.ts b/packages/core/src/__tests__/wait-for-connections.test.ts new file mode 100644 index 0000000..c0a1cae --- /dev/null +++ b/packages/core/src/__tests__/wait-for-connections.test.ts @@ -0,0 +1,109 @@ +import { describe, it, expect, vi, afterEach } from "vitest"; +import { CodeSpar } from "../index.js"; + +/** + * Regression coverage for the `manageConnections.waitForConnections` + * gate (docs/fix-core.md C2). The bug: `connections()` swallows a + * failing/empty response into `[]`, and `[].every(...)` is vacuously + * true, so the wait loop breaks on the first poll and createSession + * resolves while *zero* servers are connected — the opposite of what + * waitForConnections promises. + */ + +const SESSION_RESPONSE = { + ok: true, + status: 201, + text: async () => "", + json: async () => ({ + id: "ses_wfc", + org_id: "org_t", + user_id: "u1", + servers: ["zoop"], + status: "active", + created_at: new Date().toISOString(), + closed_at: null, + }), +}; + +function connectionsResponse(servers: Array<{ connected: boolean }>) { + return { + ok: true, + status: 200, + text: async () => "", + json: async () => ({ servers, tools: [] }), + }; +} + +const originalFetch = globalThis.fetch; + +afterEach(() => { + globalThis.fetch = originalFetch; + vi.useRealTimers(); + vi.restoreAllMocks(); +}); + +/** Advances fake timers until the pending promise settles. */ +async function settle(p: Promise): Promise { + let done = false; + void p.then(() => { + done = true; + }); + await vi.advanceTimersByTimeAsync(0); + for (let i = 0; i < 20 && !done; i++) { + await vi.advanceTimersByTimeAsync(1000); + } + return p; +} + +describe("createSession — waitForConnections", () => { + it("does not treat an empty/failed connections response as 'all connected'", async () => { + vi.useFakeTimers(); + let connCalls = 0; + globalThis.fetch = vi.fn(async (url: string) => { + if (String(url).endsWith("/connections")) { + connCalls += 1; + // Empty servers list — the vacuous-truth trap. + return connectionsResponse([]) as unknown as Response; + } + return SESSION_RESPONSE as unknown as Response; + }) as unknown as typeof fetch; + + const cs = new CodeSpar({ apiKey: "csk_live_t", baseUrl: "https://api.example.com" }); + const session = await settle( + cs.create("u1", { + servers: ["zoop"], + manageConnections: { waitForConnections: true, timeout: 5000 }, + }), + ); + + expect(session.id).toBe("ses_wfc"); + // Buggy code breaks on the first poll (connCalls === 1). Correct + // behavior keeps polling until the timeout elapses. + expect(connCalls).toBeGreaterThan(1); + }); + + it("resolves as soon as every server reports connected", async () => { + vi.useFakeTimers(); + let connCalls = 0; + globalThis.fetch = vi.fn(async (url: string) => { + if (String(url).endsWith("/connections")) { + connCalls += 1; + const connected = connCalls >= 3; + return connectionsResponse([{ connected }]) as unknown as Response; + } + return SESSION_RESPONSE as unknown as Response; + }) as unknown as typeof fetch; + + const cs = new CodeSpar({ apiKey: "csk_live_t", baseUrl: "https://api.example.com" }); + const session = await settle( + cs.create("u1", { + servers: ["zoop"], + manageConnections: { waitForConnections: true, timeout: 60000 }, + }), + ); + + expect(session.id).toBe("ses_wfc"); + // Polled until the 3rd response flipped connected:true, then stopped. + expect(connCalls).toBe(3); + }); +}); diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index 04ca194..b625b99 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -404,7 +404,10 @@ export async function createSession( const start = Date.now(); while (Date.now() - start < timeout) { const conns = await session.connections(); - if (conns.every((c) => c.connected)) break; + // `connections()` returns [] on a failed/empty response; an empty + // list must NOT count as "all connected" — [].every() is vacuously + // true and would resolve the gate with zero servers connected. + if (conns.length > 0 && conns.every((c) => c.connected)) break; await new Promise((resolve) => setTimeout(resolve, 1000)); } }