diff --git a/bin/lib/coordination.js b/bin/lib/coordination.js index 69ba909..09fede0 100644 --- a/bin/lib/coordination.js +++ b/bin/lib/coordination.js @@ -20,6 +20,7 @@ class AsyncCommandCoordinator { this.commandSequences.set(requestId, { phases, currentPhase: 0, + executing: false, phaseData: {}, startTime: Date.now(), requestId @@ -46,6 +47,7 @@ class AsyncCommandCoordinator { currentPhase: sequence.currentPhase }); + sequence.executing = true; try { const result = await phase.handler(sequence.phaseData); @@ -76,7 +78,7 @@ class AsyncCommandCoordinator { elapsed: Date.now() - sequence.startTime }); - this.commandSequences.delete(requestId); + this.removeSequence(requestId); return { completed: true, result: sequence.phaseData }; } else { logger.info('Advancing to next phase', { @@ -94,8 +96,10 @@ class AsyncCommandCoordinator { error: error.message }); - this.commandSequences.delete(requestId); + this.removeSequence(requestId); throw error; + } finally { + sequence.executing = false; } } @@ -118,7 +122,7 @@ class AsyncCommandCoordinator { if (sequence.currentPhase >= sequence.phases.length) { // All phases complete - this.commandSequences.delete(requestId); + this.removeSequence(requestId); return { completed: true, result: sequence.phaseData }; } else { // Continue to next phase @@ -162,16 +166,34 @@ class AsyncCommandCoordinator { cleanupTimedOutSequences(maxAge = 60000) { const now = Date.now(); for (const [requestId, sequence] of this.commandSequences.entries()) { - if (now - sequence.startTime > maxAge) { + const timeout = sequence.phaseData.timeout; + // Native handlers own their timeout and listener cleanup. Do not + // discard a running phase or impose a limit on an unlimited wait. + if (sequence.executing || timeout === 0) continue; + + // Allow the operation's full budget, then a grace period for PHP + // to finish its callback and collect the result (including errors). + const expiryAge = typeof timeout === 'number' && Number.isFinite(timeout) && timeout > 0 + ? timeout + maxAge + : maxAge; + if (now - sequence.startTime > expiryAge) { logger.warn('Cleaning up timed out sequence', { requestId, elapsed: now - sequence.startTime }); - this.commandSequences.delete(requestId); + this.removeSequence(requestId); } } } + removeSequence(requestId) { + this.commandSequences.delete(requestId); + if (this.commandSequences.size === 0 && this.cleanupInterval !== null) { + clearInterval(this.cleanupInterval); + this.cleanupInterval = null; + } + } + generateId(prefix = 'coord') { return `${prefix}_${Date.now()}_${Math.floor(Math.random() * 1000)}`; } @@ -183,12 +205,6 @@ class AsyncCommandCoordinator { if (this.cleanupInterval === null) { this.cleanupInterval = setInterval(() => { this.cleanupTimedOutSequences(); - - // Stop interval when no sequences are active - if (this.commandSequences.size === 0) { - clearInterval(this.cleanupInterval); - this.cleanupInterval = null; - } }, 60000); } } diff --git a/tests/Fixtures/coordinator-expiration.cjs b/tests/Fixtures/coordinator-expiration.cjs new file mode 100644 index 0000000..154d2ca --- /dev/null +++ b/tests/Fixtures/coordinator-expiration.cjs @@ -0,0 +1,107 @@ +const assert = require('node:assert/strict'); +const { test } = require('node:test'); +const { AsyncCommandCoordinator } = require('../../bin/lib/coordination'); + +async function withClock(run) { + const originalNow = Date.now; + let now = originalNow(); + Date.now = () => now; + const coordinator = new AsyncCommandCoordinator(); + + try { + await run(coordinator, (milliseconds) => { now += milliseconds; }); + } finally { + Date.now = originalNow; + clearInterval(coordinator.cleanupInterval); + } +} + +async function waitForCallback(coordinator, timeout) { + coordinator.registerAsyncCommand('request', [ + { name: 'listen', handler: async () => ({}), waitForCallback: true }, + { name: 'finish', handler: async () => ({ done: true }) }, + ]); + await coordinator.executeNextPhase('request', { timeout }); +} + +test('a callback can continue after 60 seconds within its operation timeout', () => withClock(async (coordinator, advance) => { + await waitForCallback(coordinator, 180000); + advance(120001); + coordinator.cleanupTimedOutSequences(); + + assert.equal(coordinator.isWaitingForCallback('request'), true); + const result = await coordinator.continueAfterCallback('request'); + assert.equal(result.completed, true); + assert.equal(result.result.done, true); +})); + +test('an unlimited operation is not expired by the fallback age', () => withClock(async (coordinator, advance) => { + await waitForCallback(coordinator, 0); + advance(3600000); + coordinator.cleanupTimedOutSequences(); + + assert.equal(coordinator.isWaitingForCallback('request'), true); + assert.equal((await coordinator.continueAfterCallback('request')).completed, true); +})); + +test('an abandoned callback expires after its operation timeout and callback grace', () => withClock(async (coordinator, advance) => { + await waitForCallback(coordinator, 180000); + advance(180001); + coordinator.cleanupTimedOutSequences(); + assert.equal(coordinator.isWaitingForCallback('request'), true); + + advance(60000); + coordinator.cleanupTimedOutSequences(); + assert.deepEqual(coordinator.getActiveSequences(), []); + assert.equal(coordinator.cleanupInterval, null); +})); + +test('commands without an operation timeout retain the fallback expiration', () => withClock(async (coordinator, advance) => { + await waitForCallback(coordinator, undefined); + advance(60001); + coordinator.cleanupTimedOutSequences(); + + assert.deepEqual(coordinator.getActiveSequences(), []); + assert.equal(coordinator.cleanupInterval, null); +})); + +test('cleanup does not delete a phase whose handler is still running', () => withClock(async (coordinator, advance) => { + let finish; + coordinator.registerAsyncCommand('request', [{ + name: 'wait', + handler: () => new Promise((resolve) => { finish = resolve; }), + }]); + const pending = coordinator.executeNextPhase('request', { timeout: 1000 }); + advance(120001); + coordinator.cleanupTimedOutSequences(); + + try { + assert.equal(coordinator.getActiveSequences().length, 1); + } finally { + finish({ done: true }); + await pending; + } + assert.deepEqual(coordinator.getActiveSequences(), []); + assert.equal(coordinator.cleanupInterval, null); +})); + +test('the cleanup timer stops after completion and restarts for new work', () => withClock(async (coordinator) => { + await waitForCallback(coordinator, 1000); + assert.notEqual(coordinator.cleanupInterval, null); + await coordinator.continueAfterCallback('request'); + assert.equal(coordinator.cleanupInterval, null); + + await waitForCallback(coordinator, 1000); + assert.notEqual(coordinator.cleanupInterval, null); + await coordinator.continueAfterCallback('request'); + assert.equal(coordinator.cleanupInterval, null); +})); + +test('a failed phase releases its sequence and cleanup timer', () => withClock(async (coordinator) => { + const failure = new Error('Native operation failed'); + coordinator.registerAsyncCommand('request', [{ name: 'fail', handler: async () => { throw failure; } }]); + + await assert.rejects(coordinator.executeNextPhase('request'), (error) => error === failure); + assert.deepEqual(coordinator.getActiveSequences(), []); + assert.equal(coordinator.cleanupInterval, null); +})); diff --git a/tests/Integration/Transport/AsyncCommandCoordinatorTest.php b/tests/Integration/Transport/AsyncCommandCoordinatorTest.php new file mode 100644 index 0000000..207eff1 --- /dev/null +++ b/tests/Integration/Transport/AsyncCommandCoordinatorTest.php @@ -0,0 +1,38 @@ +find('node'); + if (null === $node) { + $this->markTestSkipped('Node.js executable not found.'); + } + + $process = new Process([$node, '--test', __DIR__.'/../../Fixtures/coordinator-expiration.cjs']); + $process->setTimeout(10); + $process->run(); + + $this->assertSame(0, $process->getExitCode(), $process->getOutput().$process->getErrorOutput()); + } +}