diff --git a/src/orchestrator-loop.ts b/src/orchestrator-loop.ts index 526ad75f..a399b3bd 100644 --- a/src/orchestrator-loop.ts +++ b/src/orchestrator-loop.ts @@ -103,6 +103,9 @@ export class OrchestratorLoop extends EventEmitter { /** Phase-level timeout timer */ private phaseTimeoutTimer: NodeJS.Timeout | null = null; + /** Post-phase delay timer before verification */ + private postPhaseTimer: NodeJS.Timeout | null = null; + /** Session completion listener (bound for cleanup) */ private sessionCompletionListener: ((sessionId: string, phrase: string) => void) | null = null; @@ -559,6 +562,10 @@ export class OrchestratorLoop extends EventEmitter { clearTimeout(this.phaseTimeoutTimer); this.phaseTimeoutTimer = null; } + if (this.postPhaseTimer) { + clearTimeout(this.postPhaseTimer); + this.postPhaseTimer = null; + } } private pollPhaseStatus(phase: OrchestratorPhase): void { @@ -615,8 +622,9 @@ export class OrchestratorLoop extends EventEmitter { // Phase has failed tasks this.handlePhaseError(phase, 'One or more tasks failed'); } else { - // All tasks completed — run verification - setTimeout(() => { + // All tasks completed — run verification after brief delay + this.postPhaseTimer = setTimeout(() => { + this.postPhaseTimer = null; this.verifyCurrentPhase().catch((err) => this.handleError(err)); }, POST_PHASE_DELAY_MS); } @@ -765,10 +773,15 @@ export class OrchestratorLoop extends EventEmitter { this.persist(); + // Set up handlers so task completion is tracked + this.setupTaskHandlers(); + // Assign to a session const sessions = this.sessionManager.getIdleSessions(); if (sessions.length === 0) { - console.warn('[Orchestrator] No idle sessions for replan — task queued, waiting'); + console.warn('[Orchestrator] No idle sessions for replan — task queued, will pick up on next poll'); + // Start polling so the task gets assigned when a session becomes idle + this.startPhasePoll(phase); return; } diff --git a/src/web/routes/orchestrator-routes.ts b/src/web/routes/orchestrator-routes.ts index b7bfe5ed..019626fb 100644 --- a/src/web/routes/orchestrator-routes.ts +++ b/src/web/routes/orchestrator-routes.ts @@ -38,10 +38,10 @@ export function registerOrchestratorRoutes(app: FastifyInstance, ctx: Orchestrat return loop; } - let eventForwardingAttached = false; + let forwardingLoop: import('../../orchestrator-loop.js').OrchestratorLoop | null = null; function setupEventForwarding(loop: import('../../orchestrator-loop.js').OrchestratorLoop) { - if (eventForwardingAttached) return; - eventForwardingAttached = true; + if (forwardingLoop === loop) return; // Already attached to this loop instance + forwardingLoop = loop; loop.on('stateChanged', (state, prevState) => { ctx.broadcast(SseEvent.OrchestratorStateChanged, { state, prevState }); }); diff --git a/test/orchestrator-loop.test.ts b/test/orchestrator-loop.test.ts index 7ebe2e80..06479ddc 100644 --- a/test/orchestrator-loop.test.ts +++ b/test/orchestrator-loop.test.ts @@ -547,6 +547,32 @@ describe('OrchestratorLoop', () => { // After pause, listener should be removed expect(mockSessionManager.listenerCount('sessionCompletion')).toBeLessThan(listenerCount); }); + + it('cancels pending verify timer on pause', async () => { + const plan = createTestPlan(); + plan.phases = [plan.phases[0]]; // No verification criteria + mockPlannerInstance.generatePlan.mockResolvedValue(plan); + + // Tasks complete immediately — triggers post-phase delay timer + mockTaskQueue.addTask.mockImplementation((options) => { + const task = createMockTask(options); + task.complete(); + mockTaskQueue._tasks.set(task.id, task); + return task; + }); + + await loop.start('Goal'); + await loop.approve(); + + // Wait for poll to detect completion (2s) but pause before verify runs (1s delay) + await new Promise((resolve) => setTimeout(resolve, 2500)); + loop.pause(); + expect(loop.state).toBe('paused'); + + // Wait past where verify would have fired — state should still be paused + await new Promise((resolve) => setTimeout(resolve, 2000)); + expect(loop.state).toBe('paused'); + }); }); describe('resume()', () => {