mirror of
https://github.com/stablyai/orca.git
synced 2026-10-09 08:02:35 +00:00
Advance mocked scheduler clocks without real sleeps (#26135)
This commit is contained in:
@@ -172,6 +172,7 @@ beforeEach(() => {
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
vi.useRealTimers()
|
||||
vi.clearAllMocks()
|
||||
vi.restoreAllMocks()
|
||||
vi.unstubAllEnvs()
|
||||
@@ -202,6 +203,7 @@ describe('reconcileCodexHooks', () => {
|
||||
})
|
||||
|
||||
it('asks Codex and converts an older build only once PATH is hydrated, with launches waiting on it', async () => {
|
||||
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
|
||||
writeHooks({ Stop: [olderBuildStop()] })
|
||||
let hydrate!: () => void
|
||||
const started = start(
|
||||
@@ -211,7 +213,7 @@ describe('reconcileCodexHooks', () => {
|
||||
)
|
||||
const launch = reconcileCodexHooksForLaunch()
|
||||
// Why past the reconcile's 500 ms answer wait: an early run would write its stopgap by then.
|
||||
await new Promise((resolve) => setTimeout(resolve, 600))
|
||||
await vi.advanceTimersByTimeAsync(600)
|
||||
expect(mocks.probeCodexVersion).not.toHaveBeenCalled()
|
||||
expect(mocks.realHomeRuns).toBe(0)
|
||||
|
||||
@@ -572,14 +574,17 @@ describe('an app start whose lookup has not answered within its wait', () => {
|
||||
await start()
|
||||
const settled = snapshot(codexHome())
|
||||
// A restart: in-process answers gone, the saved one kept, and `codex --version` slow.
|
||||
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
|
||||
_internals.resetForTesting()
|
||||
lookupInternals.resetForTesting()
|
||||
mocks.probeCodexVersion.mockImplementation(
|
||||
() => new Promise((resolve) => setTimeout(() => resolve('codex-cli 0.149.0'), 900))
|
||||
)
|
||||
|
||||
await start()
|
||||
await new Promise((resolve) => setTimeout(resolve, 1_000))
|
||||
const restarted = start()
|
||||
await vi.advanceTimersByTimeAsync(500)
|
||||
await restarted
|
||||
await vi.advanceTimersByTimeAsync(1_000)
|
||||
await _internals.settledForTesting()
|
||||
|
||||
expect(snapshot(codexHome())).toEqual(settled)
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { afterEach, describe, expect, it } from 'vitest'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
import { OrchestrationDb } from './db'
|
||||
import { reconcileLifecycleMessage } from './lifecycle-reconciliation'
|
||||
import { Coordinator } from './coordinator'
|
||||
@@ -112,11 +112,39 @@ function insertWorkerDone(
|
||||
})
|
||||
}
|
||||
|
||||
/** Drive the mocked scheduler through its existing polling intervals. */
|
||||
async function settleCoordinatorRun<T>(run: Promise<T>): Promise<T> {
|
||||
let settled = false
|
||||
void run.then(
|
||||
() => {
|
||||
settled = true
|
||||
},
|
||||
() => {
|
||||
settled = true
|
||||
}
|
||||
)
|
||||
for (let elapsed = 0; !settled && elapsed < 30_000; elapsed += 50) {
|
||||
await vi.advanceTimersByTimeAsync(50)
|
||||
}
|
||||
if (!settled) {
|
||||
throw new Error('Coordinator did not settle within its test clock budget')
|
||||
}
|
||||
return run
|
||||
}
|
||||
|
||||
describe('Coordinator', () => {
|
||||
let db: OrchestrationDb
|
||||
|
||||
beforeEach(() => {
|
||||
vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout'] })
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
db?.close()
|
||||
try {
|
||||
db?.close()
|
||||
} finally {
|
||||
vi.useRealTimers()
|
||||
}
|
||||
})
|
||||
|
||||
it('throws if no tasks exist', async () => {
|
||||
@@ -151,14 +179,12 @@ describe('Coordinator', () => {
|
||||
const runPromise = coordinator.run()
|
||||
|
||||
// Wait for dispatch to happen
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Simulate the worker completing
|
||||
insertWorkerDone(db, { taskId: task.id, filesModified: ['a.ts'] })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
expect(result.completedTasks).toContain(task.id)
|
||||
expect(runtime.sentMessages.length).toBeGreaterThan(0)
|
||||
@@ -184,15 +210,13 @@ describe('Coordinator', () => {
|
||||
pollIntervalMs: 50
|
||||
})
|
||||
const runPromise = coordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
expect(db.getDispatchContext(task.id)?.assignee_pane_key).toBe('tab_a:leaf_a')
|
||||
expect(db.getDispatchContext(task.id)?.process_incarnation).toBeNull()
|
||||
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
await runPromise
|
||||
await settleCoordinatorRun(runPromise)
|
||||
})
|
||||
|
||||
it('records authenticated process authority for automatic dispatch', async () => {
|
||||
@@ -219,9 +243,7 @@ describe('Coordinator', () => {
|
||||
pollIntervalMs: 50
|
||||
})
|
||||
const runPromise = coordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
expect(db.getDispatchContext(task.id)).toMatchObject({
|
||||
assignee_pane_key: 'tab_a:leaf_a',
|
||||
@@ -230,7 +252,7 @@ describe('Coordinator', () => {
|
||||
})
|
||||
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
await runPromise
|
||||
await settleCoordinatorRun(runPromise)
|
||||
})
|
||||
|
||||
it('records completedTasks when send reconciled worker_done before coordinator read', async () => {
|
||||
@@ -258,7 +280,7 @@ describe('Coordinator', () => {
|
||||
coordinatorHandle: 'coord',
|
||||
pollIntervalMs: 20
|
||||
})
|
||||
const result = await coordinator.run()
|
||||
const result = await settleCoordinatorRun(coordinator.run())
|
||||
|
||||
expect(result.status).toBe('completed')
|
||||
expect(result.completedTasks).toContain(task.id)
|
||||
@@ -302,7 +324,7 @@ describe('Coordinator', () => {
|
||||
coordinatorHandle: 'coord',
|
||||
pollIntervalMs: 20
|
||||
})
|
||||
const result = await coordinator.run()
|
||||
const result = await settleCoordinatorRun(coordinator.run())
|
||||
|
||||
expect(result.status).toBe('completed')
|
||||
expect(result.completedTasks.filter((id) => id === task.id)).toHaveLength(1)
|
||||
@@ -322,9 +344,7 @@ describe('Coordinator', () => {
|
||||
|
||||
const runPromise = coordinator.run()
|
||||
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
expect(runtime.createdTerminals.length).toBe(1)
|
||||
expect(runtime.createdTerminalOptions[0]).not.toHaveProperty('presentation')
|
||||
@@ -332,7 +352,7 @@ describe('Coordinator', () => {
|
||||
// Complete the task
|
||||
insertWorkerDone(db, { taskId: task.id, from: runtime.createdTerminals[0] })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
})
|
||||
|
||||
@@ -359,9 +379,7 @@ describe('Coordinator', () => {
|
||||
|
||||
// Send 3 escalations to trigger circuit breaker
|
||||
for (let i = 0; i < 3; i++) {
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
const dispatch = db.getDispatchContext(task.id)
|
||||
expect(dispatch).toBeDefined()
|
||||
db.insertMessage({
|
||||
@@ -374,7 +392,7 @@ describe('Coordinator', () => {
|
||||
})
|
||||
}
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('failed')
|
||||
expect(result.failedTasks).toContain(task.id)
|
||||
})
|
||||
@@ -397,7 +415,7 @@ describe('Coordinator', () => {
|
||||
pollIntervalMs: 10
|
||||
})
|
||||
|
||||
const result = await coordinator.run()
|
||||
const result = await settleCoordinatorRun(coordinator.run())
|
||||
|
||||
expect(result.status).toBe('failed')
|
||||
expect(result.failedTasks).toContain(task.id)
|
||||
@@ -423,9 +441,7 @@ describe('Coordinator', () => {
|
||||
const runPromise = coordinator.run()
|
||||
|
||||
// Wait for dispatch
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Worker sends decision gate
|
||||
const dispatch = db.getDispatchContext(task.id)
|
||||
@@ -444,9 +460,7 @@ describe('Coordinator', () => {
|
||||
})
|
||||
})
|
||||
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Verify task is blocked
|
||||
const blocked = db.getTask(task.id)
|
||||
@@ -459,13 +473,11 @@ describe('Coordinator', () => {
|
||||
db.resolveGate(gates[0].id, 'yes')
|
||||
|
||||
// Wait for re-dispatch and simulate completion
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 200)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(200)
|
||||
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
expect(result.completedTasks).toContain(task.id)
|
||||
})
|
||||
@@ -493,9 +505,7 @@ describe('Coordinator', () => {
|
||||
const runPromise = coordinator.run()
|
||||
|
||||
// Wait for t1 dispatch
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
// t2 should still be pending
|
||||
expect(db.getTask(t2.id)?.status).toBe('pending')
|
||||
@@ -504,9 +514,7 @@ describe('Coordinator', () => {
|
||||
insertWorkerDone(db, { taskId: t1.id })
|
||||
|
||||
// Wait for t2 to be promoted and dispatched
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 200)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(200)
|
||||
|
||||
// t2 should now be dispatched
|
||||
const t2Status = db.getTask(t2.id)?.status
|
||||
@@ -515,7 +523,7 @@ describe('Coordinator', () => {
|
||||
// Complete t2
|
||||
insertWorkerDone(db, { taskId: t2.id })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
expect(result.completedTasks).toContain(t1.id)
|
||||
expect(result.completedTasks).toContain(t2.id)
|
||||
@@ -543,9 +551,7 @@ describe('Coordinator', () => {
|
||||
|
||||
const runPromise = coordinator.run()
|
||||
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
// Only 2 should be dispatched
|
||||
const dispatched = db.listTasks({ status: 'dispatched' })
|
||||
@@ -554,12 +560,10 @@ describe('Coordinator', () => {
|
||||
// Complete all tasks
|
||||
for (const task of [t1, t2, t3]) {
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
}
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
})
|
||||
|
||||
@@ -591,11 +595,9 @@ describe('Coordinator', () => {
|
||||
|
||||
// Drive one tick then stop — we only need the stale warning to have fired.
|
||||
const runPromise = coordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 80)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(80)
|
||||
coordinator.stop()
|
||||
await runPromise
|
||||
await settleCoordinatorRun(runPromise)
|
||||
|
||||
expect(logs.some((l) => /has not sent a heartbeat/.test(l) && l.includes(task.id))).toBe(true)
|
||||
// Task status must NOT have been auto-failed — logging only.
|
||||
@@ -627,16 +629,14 @@ describe('Coordinator', () => {
|
||||
payload: JSON.stringify({ taskId: task.id, dispatchId: ctx.id, phase: 'implementing' })
|
||||
})
|
||||
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 80)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(80)
|
||||
|
||||
expect(db.getDispatchContext(task.id)?.last_heartbeat_at).toBeTruthy()
|
||||
|
||||
// Complete the task so the coordinator run finishes cleanly.
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
})
|
||||
|
||||
@@ -673,11 +673,9 @@ describe('Coordinator', () => {
|
||||
onLog: (m) => logs.push(m)
|
||||
})
|
||||
const staleRun = staleCoordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 80)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(80)
|
||||
staleCoordinator.stop()
|
||||
await staleRun
|
||||
await settleCoordinatorRun(staleRun)
|
||||
|
||||
expect(db.getTask(task.id)?.status).toBe('dispatched')
|
||||
expect(db.getDispatchContextById(staleCtx.id)?.status).toBe('failed')
|
||||
@@ -694,7 +692,7 @@ describe('Coordinator', () => {
|
||||
coordinatorHandle: 'coord',
|
||||
pollIntervalMs: 20
|
||||
})
|
||||
const result = await completionCoordinator.run()
|
||||
const result = await settleCoordinatorRun(completionCoordinator.run())
|
||||
|
||||
expect(result.status).toBe('completed')
|
||||
expect(db.getTask(task.id)?.status).toBe('completed')
|
||||
@@ -729,7 +727,7 @@ describe('Coordinator', () => {
|
||||
pollIntervalMs: 20,
|
||||
onLog: (m) => logs.push(m)
|
||||
})
|
||||
const result = await coordinator.run()
|
||||
const result = await settleCoordinatorRun(coordinator.run())
|
||||
|
||||
expect(result.status).toBe('completed')
|
||||
expect(db.getTask(task.id)?.status).toBe('completed')
|
||||
@@ -750,12 +748,10 @@ describe('Coordinator', () => {
|
||||
|
||||
const runPromise = coordinator.run()
|
||||
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
coordinator.stop()
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('failed')
|
||||
})
|
||||
|
||||
@@ -783,13 +779,11 @@ describe('Coordinator', () => {
|
||||
})
|
||||
|
||||
const runPromise = coordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
expect(runtime.probeDriftCalls).toContain('wt1')
|
||||
const sent = runtime.sentMessages.find((m) => m.handle === 'term_a')
|
||||
@@ -822,11 +816,9 @@ describe('Coordinator', () => {
|
||||
})
|
||||
|
||||
const runPromise = coordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 250)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(250)
|
||||
coordinator.stop()
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
|
||||
// Why: silent-skip must NOT burn the circuit-breaker budget. Task must
|
||||
// stay in `ready`; failDispatch must NOT be called; no prompt injection
|
||||
@@ -862,13 +854,11 @@ allow-stale-base: true`
|
||||
})
|
||||
|
||||
const runPromise = coordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
const sent = runtime.sentMessages.find((m) => m.handle === 'term_a')
|
||||
expect(sent).toBeDefined()
|
||||
@@ -898,13 +888,11 @@ allow-stale-base: true`
|
||||
})
|
||||
|
||||
const runPromise = coordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
const sent = runtime.sentMessages.find((m) => m.handle === 'term_a')
|
||||
expect(sent).toBeDefined()
|
||||
@@ -931,13 +919,11 @@ allow-stale-base: true`
|
||||
})
|
||||
|
||||
const runPromise = coordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
expect(runtime.probeDriftCalls).toHaveLength(0)
|
||||
expect(logs.some((m) => m.includes('stale-base guard inert'))).toBe(true)
|
||||
@@ -964,13 +950,11 @@ allow-stale-base: true`
|
||||
})
|
||||
|
||||
const runPromise = coordinator.run()
|
||||
await new Promise((r) => {
|
||||
setTimeout(r, 100)
|
||||
})
|
||||
await vi.advanceTimersByTimeAsync(100)
|
||||
|
||||
insertWorkerDone(db, { taskId: task.id })
|
||||
|
||||
const result = await runPromise
|
||||
const result = await settleCoordinatorRun(runPromise)
|
||||
expect(result.status).toBe('completed')
|
||||
const sent = runtime.sentMessages.find((m) => m.handle === 'term_a')
|
||||
expect(sent!.text).not.toContain('--- BASE DRIFT ---')
|
||||
|
||||
Reference in New Issue
Block a user