From 7665ca03102bff3a6581d4802e52173d402c2a0d Mon Sep 17 00:00:00 2001 From: mohammed naji Date: Thu, 13 Aug 2026 13:06:12 +0400 Subject: [PATCH 1/2] test(watch): make adaptive backoff assertions deterministic The idle-backoff test drove a real watcher for a fixed 190 ms window and then required the 80 ms rung to have been reached while also capping the reconciliation count at five. That contract is constrained from both sides by the host scheduler, so it cannot be repaired by widening the window: a wider window helps the slow side and worsens the count ceiling. Separate the policy from the integration path with the smallest seam that the evidence supports: - extract the idle-backoff rule into a pure exported nextReconciliationIntervalMs and call it from the watch loop, an expression-identical refactor with no behavior change; - export the existing internal createWatchLoopSignal so its wait/wake/ abort/cleanup semantics can be asserted directly under fake timers, without making it injectable or touching any call site. Policy rungs, the cap, the activity reset, the floor and the loop-signal lifecycle now live in deterministic tests with no wall-clock dependency. The integration test keeps a real watcher but waits on an observable barrier and validates every emitted transition against the pure policy, so it holds regardless of host speed or concurrent filesystem activity. The Git-visibility cache test loses its 100 ms window for the same reason and now asserts that observed polls shared fewer Git snapshots than there were reconciliations. Refs #689. Related parent: #654. #654 remains open. --- src/infrastructure/watch.ts | 28 +++- tests/unit/watch-backoff-policy.test.ts | 193 ++++++++++++++++++++++++ tests/unit/watch.test.ts | 112 ++++++++++++-- 3 files changed, 317 insertions(+), 16 deletions(-) create mode 100644 tests/unit/watch-backoff-policy.test.ts diff --git a/src/infrastructure/watch.ts b/src/infrastructure/watch.ts index c051aef5..9bb95601 100644 --- a/src/infrastructure/watch.ts +++ b/src/infrastructure/watch.ts @@ -137,7 +137,8 @@ export interface GraphAutoRefreshController { completed: Promise } -interface WatchLoopSignal { +/** @internal Exported for deterministic testing of wait/wake/abort semantics. */ +export interface WatchLoopSignal { wait(delayMs: number, signal?: AbortSignal): Promise wake(): void } @@ -155,7 +156,21 @@ function resolveWatchPath(watchPath: string): string { return resolve(watchPath) } -function createWatchLoopSignal(): WatchLoopSignal { +/** Activity resets to the minimum; an idle reconciliation doubles the current interval, clamped to [minimum, maximum]. */ +export function nextReconciliationIntervalMs(input: { + currentIntervalMs: number + minimumIntervalMs: number + maximumIntervalMs: number + changedCount: number +}): number { + const { currentIntervalMs, minimumIntervalMs, maximumIntervalMs, changedCount } = input + return changedCount > 0 + ? minimumIntervalMs + : Math.min(maximumIntervalMs, Math.max(minimumIntervalMs, currentIntervalMs * 2)) +} + +/** @internal Exported for deterministic testing of wait/wake/abort semantics. */ +export function createWatchLoopSignal(): WatchLoopSignal { let wakePending = false let wakeResolver: (() => void) | null = null @@ -1052,9 +1067,12 @@ export async function watch(watchPath: string, debounce = 3, options: WatchOptio ) const changedBatch = diffSnapshots(previousSnapshot.fingerprints, nextSnapshot.fingerprints) previousSnapshot = nextSnapshot - currentIntervalMs = changedBatch.length > 0 - ? minimumIntervalMs - : Math.min(maximumIntervalMs, Math.max(minimumIntervalMs, currentIntervalMs * 2)) + currentIntervalMs = nextReconciliationIntervalMs({ + currentIntervalMs, + minimumIntervalMs, + maximumIntervalMs, + changedCount: changedBatch.length, + }) nextReconciliationAt = Date.now() + currentIntervalMs recordSuccessfulReconciliation(state, nextSnapshot, currentIntervalMs, nextReconciliationAt) updateWatcherPolicyState(state, resolvedWatchPath, options, gitVisibilityCache) diff --git a/tests/unit/watch-backoff-policy.test.ts b/tests/unit/watch-backoff-policy.test.ts new file mode 100644 index 00000000..d1aa5f46 --- /dev/null +++ b/tests/unit/watch-backoff-policy.test.ts @@ -0,0 +1,193 @@ +import { afterEach, beforeEach, describe, expect, test, vi } from 'vitest' + +import { + createWatchLoopSignal, + nextReconciliationIntervalMs, +} from '../../src/infrastructure/watch.js' + +describe('nextReconciliationIntervalMs', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + + afterEach(() => { + try { + expect(vi.getTimerCount()).toBe(0) + } finally { + vi.useRealTimers() + } + }) + + const nextInterval = (currentIntervalMs: number, changedCount = 0, maximumIntervalMs = 80): number => ( + nextReconciliationIntervalMs({ + currentIntervalMs, + minimumIntervalMs: 20, + maximumIntervalMs, + changedCount, + }) + ) + + test('backs off from 20 to 40 on the first idle transition', () => { + expect(nextInterval(20)).toBe(40) + }) + + test('backs off from 40 to 80 on the next idle transition', () => { + expect(nextInterval(40)).toBe(80) + }) + + test('holds at the maximum across repeated idle transitions', () => { + let currentIntervalMs = 80 + + for (let step = 0; step < 3; step += 1) { + currentIntervalMs = nextInterval(currentIntervalMs) + expect(currentIntervalMs).toBe(80) + } + }) + + test.each([20, 40, 80])('resets the %i ms rung after activity', (currentIntervalMs) => { + expect(nextInterval(currentIntervalMs, 1)).toBe(20) + }) + + test('raises a current interval below the minimum to the floor', () => { + expect(nextInterval(5)).toBe(20) + }) + + test('caps a doubling step at a non-power-of-two maximum', () => { + expect(nextInterval(40, 0, 50)).toBe(50) + }) + + test.each([ + { initialIntervalMs: 20, idleSteps: 4, expected: [20, 40, 80, 80, 80] }, + ])('emits the exact idle rung sequence from $initialIntervalMs ms', ({ initialIntervalMs, idleSteps, expected }) => { + const emitted = [initialIntervalMs] + let currentIntervalMs = initialIntervalMs + + for (let step = 0; step < idleSteps; step += 1) { + currentIntervalMs = nextInterval(currentIntervalMs) + emitted.push(currentIntervalMs) + } + + expect(emitted).toEqual(expected) + }) +}) + +describe('createWatchLoopSignal', () => { + beforeEach(() => { + vi.useFakeTimers() + }) + + afterEach(() => { + try { + expect(vi.getTimerCount()).toBe(0) + } finally { + vi.restoreAllMocks() + vi.useRealTimers() + } + }) + + test('wait resolves after its delay', async () => { + const controller = new AbortController() + const addEventListener = vi.spyOn(controller.signal, 'addEventListener') + const removeEventListener = vi.spyOn(controller.signal, 'removeEventListener') + const loopSignal = createWatchLoopSignal() + const resolved = vi.fn() + const waiting = loopSignal.wait(100, controller.signal).then(resolved) + + expect(vi.getTimerCount()).toBe(1) + await vi.advanceTimersByTimeAsync(99) + expect(resolved).not.toHaveBeenCalled() + + await vi.advanceTimersByTimeAsync(1) + await waiting + + expect(resolved).toHaveBeenCalledTimes(1) + expect(vi.getTimerCount()).toBe(0) + expect(addEventListener).toHaveBeenCalledTimes(1) + expect(removeEventListener).toHaveBeenCalledTimes(1) + expect(removeEventListener).toHaveBeenCalledWith('abort', addEventListener.mock.calls[0]?.[1]) + }) + + test('wake resolves a pending wait early', async () => { + const controller = new AbortController() + const addEventListener = vi.spyOn(controller.signal, 'addEventListener') + const removeEventListener = vi.spyOn(controller.signal, 'removeEventListener') + const loopSignal = createWatchLoopSignal() + const waiting = loopSignal.wait(100, controller.signal) + + expect(vi.getTimerCount()).toBe(1) + loopSignal.wake() + await waiting + + expect(vi.getTimerCount()).toBe(0) + expect(addEventListener).toHaveBeenCalledTimes(1) + expect(removeEventListener).toHaveBeenCalledTimes(1) + expect(removeEventListener).toHaveBeenCalledWith('abort', addEventListener.mock.calls[0]?.[1]) + }) + + test('a wake before a wait is consumed immediately by the next wait', async () => { + const controller = new AbortController() + const addEventListener = vi.spyOn(controller.signal, 'addEventListener') + const removeEventListener = vi.spyOn(controller.signal, 'removeEventListener') + const loopSignal = createWatchLoopSignal() + + loopSignal.wake() + await loopSignal.wait(100, controller.signal) + + expect(vi.getTimerCount()).toBe(0) + expect(addEventListener).not.toHaveBeenCalled() + expect(removeEventListener).not.toHaveBeenCalled() + }) + + test('an already-aborted signal resolves immediately without scheduling a timer', async () => { + const controller = new AbortController() + controller.abort() + const addEventListener = vi.spyOn(controller.signal, 'addEventListener') + const removeEventListener = vi.spyOn(controller.signal, 'removeEventListener') + const loopSignal = createWatchLoopSignal() + + await loopSignal.wait(100, controller.signal) + + expect(vi.getTimerCount()).toBe(0) + expect(addEventListener).not.toHaveBeenCalled() + expect(removeEventListener).not.toHaveBeenCalled() + }) + + test('aborting during a pending wait resolves it', async () => { + const controller = new AbortController() + const addEventListener = vi.spyOn(controller.signal, 'addEventListener') + const removeEventListener = vi.spyOn(controller.signal, 'removeEventListener') + const loopSignal = createWatchLoopSignal() + const waiting = loopSignal.wait(100, controller.signal) + + expect(vi.getTimerCount()).toBe(1) + controller.abort() + await waiting + + expect(vi.getTimerCount()).toBe(0) + expect(addEventListener).toHaveBeenCalledTimes(1) + expect(removeEventListener).toHaveBeenCalledTimes(1) + expect(removeEventListener).toHaveBeenCalledWith('abort', addEventListener.mock.calls[0]?.[1]) + }) + + test('two sequential waits with the same deadline resolve in call order', async () => { + const controller = new AbortController() + const addEventListener = vi.spyOn(controller.signal, 'addEventListener') + const removeEventListener = vi.spyOn(controller.signal, 'removeEventListener') + const loopSignal = createWatchLoopSignal() + const resolutionOrder: string[] = [] + + const first = loopSignal.wait(100, controller.signal).then(() => resolutionOrder.push('first')) + const second = loopSignal.wait(100, controller.signal).then(() => resolutionOrder.push('second')) + + expect(vi.getTimerCount()).toBe(2) + await vi.advanceTimersByTimeAsync(100) + await Promise.all([first, second]) + + expect(resolutionOrder).toEqual(['first', 'second']) + expect(vi.getTimerCount()).toBe(0) + expect(addEventListener).toHaveBeenCalledTimes(2) + expect(removeEventListener).toHaveBeenCalledTimes(2) + expect(removeEventListener).toHaveBeenNthCalledWith(1, 'abort', addEventListener.mock.calls[0]?.[1]) + expect(removeEventListener).toHaveBeenNthCalledWith(2, 'abort', addEventListener.mock.calls[1]?.[1]) + }) +}) diff --git a/tests/unit/watch.test.ts b/tests/unit/watch.test.ts index 0f6a9179..c66156e2 100644 --- a/tests/unit/watch.test.ts +++ b/tests/unit/watch.test.ts @@ -6,7 +6,7 @@ import { setTimeout as delay } from 'node:timers/promises' import { describe, expect, test, vi } from 'vitest' -import { WATCHED_EXTENSIONS, hasNonCode, notifyOnly, rebuildCode, startGraphAutoRefresh, watch, type WatchReconciliationMetrics } from '../../src/infrastructure/watch.js' +import { WATCHED_EXTENSIONS, hasNonCode, nextReconciliationIntervalMs, notifyOnly, rebuildCode, startGraphAutoRefresh, watch, type WatchReconciliationMetrics } from '../../src/infrastructure/watch.js' import { generateGraph } from '../../src/infrastructure/generate.js' import { parseGenerationPolicy } from '../../src/contracts/generation-policy.js' import { readWatcherStateForGraph } from '../../src/infrastructure/watcher-state.js' @@ -469,18 +469,101 @@ describe('watch', () => { logger: { log() {}, error() {} }, }) - await delay(190) - controller.abort() - await watcher + try { + await waitFor(() => reconciliations.some((metrics) => metrics.nextIntervalMs === 80), 10_000) + } finally { + controller.abort() + await watcher + } expect(reconciliations[0]).toMatchObject({ trigger: 'initial', fileCount: 1, nextIntervalMs: 20 }) - expect(reconciliations.some((metrics) => metrics.nextIntervalMs === 40)).toBe(true) - expect(reconciliations.some((metrics) => metrics.nextIntervalMs === 80)).toBe(true) - expect(reconciliations.length).toBeLessThanOrEqual(5) + + // Every observed transition must equal what the pure policy prescribes for the + // observed activity. This ties the integration path to the deterministic policy + // tests and holds regardless of host speed or concurrent filesystem activity. + for (let index = 1; index < reconciliations.length; index += 1) { + const previous = reconciliations[index - 1] + const current = reconciliations[index] + if (!previous || !current || current.trigger === 'post-rebuild') { + continue + } + expect({ at: index, nextIntervalMs: current.nextIntervalMs }).toEqual({ + at: index, + nextIntervalMs: nextReconciliationIntervalMs({ + currentIntervalMs: previous.nextIntervalMs, + minimumIntervalMs: 20, + maximumIntervalMs: 80, + changedCount: current.changedCount, + }), + }) + } + + const intervals = reconciliations.map((metrics) => metrics.nextIntervalMs) + expect(intervals).toContain(40) + expect(intervals).toContain(80) + expect(intervals.every((interval) => interval >= 20 && interval <= 80)).toBe(true) expect(reconciliations.every((metrics) => metrics.durationMs >= 0 && metrics.directoryCount >= 1)).toBe(true) }) }) + test('resets the reconciliation interval after activity', async () => { + await withTempDirAsync(async (tempDir) => { + writeFileSync(join(tempDir, 'main.ts'), 'export const idle = true\n', 'utf8') + const controller = new AbortController() + const reconciliations: WatchReconciliationMetrics[] = [] + const watcher = watch(tempDir, 0.02, { + signal: controller.signal, + pollIntervalMs: 20, + maxPollIntervalMs: 80, + onReconciliation: (metrics) => reconciliations.push(metrics), + logger: { log() {}, error() {} }, + }) + + try { + await waitFor(() => reconciliations.some((metrics) => metrics.nextIntervalMs === 80), 10_000) + writeFileSync(join(tempDir, 'activity.ts'), 'export const activity = true\n', 'utf8') + await waitFor(() => reconciliations.some((metrics) => ( + metrics.changedCount > 0 && metrics.nextIntervalMs === 20 + )), 10_000) + } finally { + controller.abort() + await watcher + } + + expect(reconciliations.some((metrics) => ( + metrics.changedCount > 0 && metrics.nextIntervalMs === 20 + ))).toBe(true) + }) + }) + + test('emits no further reconciliation after the watcher stops', async () => { + await withTempDirAsync(async (tempDir) => { + writeFileSync(join(tempDir, 'main.ts'), 'export const idle = true\n', 'utf8') + const controller = new AbortController() + const reconciliations: WatchReconciliationMetrics[] = [] + const watcher = watch(tempDir, 0.02, { + signal: controller.signal, + pollIntervalMs: 20, + maxPollIntervalMs: 80, + onReconciliation: (metrics) => reconciliations.push(metrics), + logger: { log() {}, error() {} }, + }) + + try { + await waitFor(() => reconciliations.some((metrics) => metrics.nextIntervalMs === 80), 10_000) + } finally { + controller.abort() + await watcher + } + + const reconciliationCountAfterStop = reconciliations.length + await delay(150) + + expect(reconciliations).toHaveLength(reconciliationCountAfterStop) + expect(readWatcherStateForGraph(join(tempDir, 'out', 'graph.json'))?.status).toBe('stopped') + }) + }) + test('persists pending and stopped watcher health without answering silently stale', async () => { await withTempDirAsync(async (tempDir) => { writeFileSync(join(tempDir, 'main.ts'), 'export const initial = true\n', 'utf8') @@ -729,18 +812,25 @@ describe('watch', () => { try { const { watch: watchWithMockedGit } = await import('../../src/infrastructure/watch.js') const controller = new AbortController() + const reconciliations: WatchReconciliationMetrics[] = [] const watcher = watchWithMockedGit(tempDir, 0.02, { signal: controller.signal, pollIntervalMs: 10, respectGitignore: true, + onReconciliation: (metrics) => reconciliations.push(metrics), logger: { log() {}, error() {} }, }) - await delay(100) - controller.abort() - await watcher + try { + await waitFor(() => reconciliations.length >= 3, 5_000) + } finally { + controller.abort() + await watcher + } - expect(collectGitVisibleFiles).toHaveBeenCalledTimes(1) + expect(collectGitVisibleFiles).toHaveBeenCalled() + // Fewer Git snapshots than reconciliation polls proves cache reuse without depending on crossing its 500 ms window. + expect(collectGitVisibleFiles.mock.calls.length).toBeLessThan(reconciliations.length) } finally { vi.doUnmock('../../src/shared/git.js') vi.resetModules() From df5df05c705740733c564e2f97088f74e6a87997 Mon Sep 17 00:00:00 2001 From: mohammed naji Date: Thu, 13 Aug 2026 13:33:14 +0400 Subject: [PATCH 2/2] docs(watch): mark the backoff policy seam @internal nextReconciliationIntervalMs was the only one of the three exported test seam symbols without an @internal marker, while WatchLoopSignal and createWatchLoopSignal both carried one. Annotate it consistently and keep the sentence describing the rule. The package has no exports map, so this module was already deep-importable in full before these three symbols existed; enforcing that boundary is a separate package-compatibility change and is deliberately not attempted here. Refs #689. Related parent: #654. #654 remains open. --- src/infrastructure/watch.ts | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/src/infrastructure/watch.ts b/src/infrastructure/watch.ts index 9bb95601..95ed8c51 100644 --- a/src/infrastructure/watch.ts +++ b/src/infrastructure/watch.ts @@ -156,7 +156,10 @@ function resolveWatchPath(watchPath: string): string { return resolve(watchPath) } -/** Activity resets to the minimum; an idle reconciliation doubles the current interval, clamped to [minimum, maximum]. */ +/** + * @internal Exported for deterministic testing of the adaptive backoff policy. + * Activity resets to the minimum; an idle reconciliation doubles the current interval, clamped to [minimum, maximum]. + */ export function nextReconciliationIntervalMs(input: { currentIntervalMs: number minimumIntervalMs: number