import assert from 'node:assert/strict'; import test from 'node:test'; import { createExecutorWorkerRuntime, createWorkerAdapterRegistry, } from './executor-worker-runtime.mjs'; function claim(overrides = {}) { return { leaseToken: 'lease-1', job: { jobId: 'job-1', executor: 'goosed', fallbackExecutor: 'aider', fallbackAvailable: true, timeoutMs: 30_000, request: { executor: 'goosed', task: { instruction: 'do it', workspaceRef: { kind: 'workspace-alias', id: 'canary' }, }, }, ...overrides, }, }; } function fakeClient(nextClaim = claim()) { const calls = []; let claimed = false; return { calls, async claim(input) { calls.push(['claim', input]); if (claimed) return null; claimed = true; return nextClaim; }, async start(jobId, input) { calls.push(['start', jobId, input]); }, async heartbeat(jobId, input) { calls.push(['heartbeat', jobId, input]); }, async progress(jobId, input) { calls.push(['progress', jobId, input]); }, async complete(jobId, input) { calls.push(['complete', jobId, input]); }, async fail(jobId, input) { calls.push(['fail', jobId, input]); }, async recoverExpired(input) { calls.push(['recover', input]); return { recovered: [] }; }, }; } function adapter(id, submit) { return { id, capabilities: [], submit }; } test('worker runtime claims, executes, emits progress and completes through lease protocol', async () => { const client = fakeClient(); const registry = createWorkerAdapterRegistry([ adapter('goosed', async (_request, { emit }) => { await emit('executor_job_adapter_event', { eventType: 'finish' }); return { outcome: 'completed', artifactRefs: [{ kind: 'goosed-session', id: 'session-1' }], }; }), ]); const worker = createExecutorWorkerRuntime({ client, registry, workerId: 'worker-1', }); const result = await worker.runOnce(); assert.equal(result.status, 'succeeded'); assert.deepEqual( client.calls.map((call) => call[0]), ['claim', 'start', 'progress', 'complete'], ); assert.deepEqual(client.calls[0][1].adapterHealth, { goosed: { ok: true } }); assert.equal(client.calls.at(-1)[2].leaseToken, 'lease-1'); assert.deepEqual(worker.status().activeJobs, []); }); test('worker runtime falls back to an installed adapter and records the transition', async () => { const client = fakeClient(); const registry = createWorkerAdapterRegistry([ adapter('goosed', async () => { const error = new Error('primary failed'); error.code = 'PRIMARY_FAILED'; throw error; }), adapter('aider', async (request) => ({ outcome: 'completed', summary: request.executor, artifactRefs: [{ kind: 'patch', id: 'patch-1' }], })), ]); const worker = createExecutorWorkerRuntime({ client, registry, workerId: 'worker-1', }); const result = await worker.runOnce(); assert.equal(result.status, 'succeeded'); const fallback = client.calls.find( (call) => call[0] === 'progress' && call[2].type === 'executor_job_fallback_started', ); assert.deepEqual(fallback[2].data, { from: 'goosed', to: 'aider', reason: 'PRIMARY_FAILED', }); assert.equal(client.calls.at(-1)[0], 'complete'); }); test('worker runtime reports bounded adapter failures without completing the job', async () => { const client = fakeClient(claim({ fallbackAvailable: false })); const registry = createWorkerAdapterRegistry([ adapter('goosed', async () => { const error = new Error('upstream unavailable'); error.code = 'UPSTREAM_UNAVAILABLE'; error.retryable = true; throw error; }), ]); const worker = createExecutorWorkerRuntime({ client, registry, workerId: 'worker-1', }); const result = await worker.runOnce(); assert.equal(result.status, 'failed'); assert.equal(client.calls.at(-1)[0], 'fail'); assert.equal(client.calls.at(-1)[2].error.retryable, true); });