import assert from 'node:assert/strict'; import test from 'node:test'; import { createInMemoryExecutorJobStore } from './executor-gateway.mjs'; import { createExecutorWorkerCoordinator } from './executor-worker-protocol.mjs'; function queuedJob(overrides = {}) { return { version: 'executor-job-state-v1', jobId: 'job-1', idempotencyKey: 'idem-1', requestFingerprint: 'a'.repeat(64), executor: 'goosed', status: 'queued', attempts: 0, maxAttempts: 2, createdAt: 1000, updatedAt: 1000, completedAt: null, ...overrides, }; } async function seed(store, job = queuedJob()) { await store.createIfAbsent(job); } test('worker claims, starts, heartbeats and completes a queued job with lease fencing', async () => { let now = 1000; let token = 0; const store = createInMemoryExecutorJobStore(); await seed(store); const coordinator = createExecutorWorkerCoordinator({ store, nowMs: () => now, randomUUID: () => `lease-${++token}`, }); const claimed = await coordinator.claim({ workerId: 'worker-1', executors: ['goosed'], leaseDurationMs: 5000, }); assert.equal(claimed.version, 'executor-worker-v1'); assert.equal(claimed.leaseToken, 'lease-1'); assert.equal(claimed.job.status, 'leased'); assert.equal(claimed.job.attempts, 1); assert.equal(await coordinator.claim({ workerId: 'worker-2', executors: ['goosed'], }), null); await assert.rejects( coordinator.start('job-1', { leaseToken: 'wrong-token' }), (error) => error.code === 'EXECUTOR_LEASE_LOST', ); const started = await coordinator.start('job-1', { leaseToken: 'lease-1' }); assert.equal(started.status, 'running'); assert.equal(started.startedAt, 1000); now = 3000; const heartbeat = await coordinator.heartbeat('job-1', { leaseToken: 'lease-1', leaseDurationMs: 5000, }); assert.equal(heartbeat.lease.heartbeatAt, 3000); assert.equal(heartbeat.lease.expiresAt, 8000); await coordinator.progress('job-1', { leaseToken: 'lease-1', type: 'executor_job_output_chunk', data: { characters: 12 }, }); now = 4000; const completed = await coordinator.complete('job-1', { leaseToken: 'lease-1', result: { summary: 'done', artifactRefs: [{ kind: 'patch', id: 'artifact-1' }], }, }); assert.equal(completed.status, 'succeeded'); assert.equal(completed.lease, null); assert.deepEqual(completed.result.artifactRefs, [{ kind: 'patch', id: 'artifact-1' }]); assert.deepEqual( (await store.listEvents('job-1')).map((event) => event.type), [ 'executor_job_leased', 'executor_job_started', 'executor_job_heartbeat', 'executor_job_output_chunk', 'executor_job_succeeded', ], ); }); test('retryable failure is delayed and a later terminal failure consumes max attempts', async () => { let now = 1000; let token = 0; const store = createInMemoryExecutorJobStore(); await seed(store); const coordinator = createExecutorWorkerCoordinator({ store, nowMs: () => now, randomUUID: () => `lease-${++token}`, }); const first = await coordinator.claim({ workerId: 'worker-1', executors: ['goosed'] }); await coordinator.start('job-1', { leaseToken: first.leaseToken }); const retryable = await coordinator.fail('job-1', { leaseToken: first.leaseToken, error: { code: 'UPSTREAM_BUSY', message: 'busy', retryable: true }, retryDelayMs: 5000, }); assert.equal(retryable.status, 'retryable'); assert.equal(retryable.nextAttemptAt, 6000); assert.equal(await coordinator.claim({ workerId: 'worker-2', executors: ['goosed'], }), null); now = 6000; const second = await coordinator.claim({ workerId: 'worker-2', executors: ['goosed'] }); assert.equal(second.job.attempts, 2); await coordinator.start('job-1', { leaseToken: second.leaseToken }); const failed = await coordinator.fail('job-1', { leaseToken: second.leaseToken, error: { code: 'UPSTREAM_BUSY', message: 'still busy', retryable: true }, }); assert.equal(failed.status, 'failed'); assert.equal(failed.completedAt, 6000); }); test('expired leases are recovered with fencing and eventually time out', async () => { let now = 1000; let token = 0; const store = createInMemoryExecutorJobStore(); await seed(store, queuedJob({ maxAttempts: 1 })); const coordinator = createExecutorWorkerCoordinator({ store, nowMs: () => now, randomUUID: () => `lease-${++token}`, }); const claimed = await coordinator.claim({ workerId: 'worker-crashed', executors: ['goosed'], leaseDurationMs: 5000, }); await coordinator.start('job-1', { leaseToken: claimed.leaseToken }); now = 6001; const recovery = await coordinator.recoverExpired(); assert.equal(recovery.inspected, 1); assert.equal(recovery.recovered[0].status, 'timed_out'); await assert.rejects( coordinator.heartbeat('job-1', { leaseToken: claimed.leaseToken }), (error) => error.code === 'EXECUTOR_JOB_STATE_CONFLICT', ); const stats = await coordinator.stats(); assert.equal(stats.counts.timed_out, 1); assert.equal(stats.expiredLeases, 0); }); test('worker protocol bounds identifiers and progress payloads', async () => { const store = createInMemoryExecutorJobStore(); await seed(store); const coordinator = createExecutorWorkerCoordinator({ store, randomUUID: () => 'lease-1', }); await assert.rejects( coordinator.claim({ workerId: '../worker', executors: ['goosed'] }), (error) => error.code === 'EXECUTOR_WORKER_ID_INVALID', ); const claimed = await coordinator.claim({ workerId: 'worker-1', executors: ['goosed'] }); await assert.rejects( coordinator.progress('job-1', { leaseToken: claimed.leaseToken, data: { output: 'x'.repeat(20_000) }, }), (error) => error.code === 'EXECUTOR_EVENT_DATA_TOO_LARGE', ); });