From e679f1ac44d449992b63507dc7c5756659ae31d6 Mon Sep 17 00:00:00 2001 From: qa-engineer Date: Tue, 29 Sep 2026 22:01:03 +0000 Subject: [PATCH 1/3] fix(gate): persist sleep and self-wake intent --- changelog.d/durable-wake-intents.fixed.md | 5 + src/framework.ts | 4 + src/gate/event-gate.ts | 224 ++++++++++++++++++++- test/event-gate-selfwake.test.ts | 232 +++++++++++++++++++++- 4 files changed, 454 insertions(+), 11 deletions(-) create mode 100644 changelog.d/durable-wake-intents.fixed.md diff --git a/changelog.d/durable-wake-intents.fixed.md b/changelog.d/durable-wake-intents.fixed.md new file mode 100644 index 00000000..6fd2d7b2 --- /dev/null +++ b/changelog.d/durable-wake-intents.fixed.md @@ -0,0 +1,5 @@ +- `sleep` and `skip_reply(wake_in_seconds)` persist their armed wake intent so + process restarts re-arm future deadlines or visibly reconcile overdue and + overdue wakes exactly once, including across graceful restarts. +- Sleep timers re-arm after backward wall-clock steps instead of silently + consuming the only wake callback before the promised deadline. diff --git a/src/framework.ts b/src/framework.ts index 519e3042..d834a7e9 100644 --- a/src/framework.ts +++ b/src/framework.ts @@ -1851,6 +1851,10 @@ export class AgentFramework { // SIGUSR2 not available on this platform — non-fatal. } + // Reconcile wake promises only after every module/server/startup invariant + // above succeeded. A failed create must leave due intents durable. + framework.eventGate?.recoverWakeIntents(); + return framework; } diff --git a/src/gate/event-gate.ts b/src/gate/event-gate.ts index 6db0a737..bde43de9 100644 --- a/src/gate/event-gate.ts +++ b/src/gate/event-gate.ts @@ -11,8 +11,8 @@ * inference tokens, not whether the agent should know about the event. */ -import { existsSync, readFileSync, statSync, writeFileSync, mkdirSync } from 'node:fs'; -import { dirname, join } from 'node:path'; +import { existsSync, readFileSync, statSync, writeFileSync, mkdirSync, renameSync } from 'node:fs'; +import { basename, dirname, join } from 'node:path'; import { GateScript } from './gate-script.js'; import type { ToolDefinition } from '../types/events.js'; @@ -506,6 +506,21 @@ export function formatShadowWarning(w: ShadowWarning): string { // EventGate // --------------------------------------------------------------------------- +interface PersistedWakeIntent { + agentName?: string; + armedAt: number; + wakeAt: number; + source: string; + note?: string; + suppressed?: number; +} + +interface PersistedWakeState { + version: 1; + sleep?: PersistedWakeIntent; + selfWakes: PersistedWakeIntent[]; +} + export class EventGate { private configPath: string; private config: GateConfig; @@ -561,6 +576,9 @@ export class EventGate { private sleepSuppressed = 0; private sleepTimer: ReturnType | null = null; private sleepAgent: string | undefined; + private sleepArmedAt = 0; + /** Branch-independent wake-intent journal, sibling to gate.json. */ + private wakeStatePath: string; // Self-wake timers per agent (skip_reply's wake_in_seconds). Unlike sleep, // a self-wake suppresses NOTHING — it means "if nothing else wakes me by @@ -568,6 +586,7 @@ export class EventGate { // self-wake (an external wake supersedes it; the timer would otherwise // land as a redundant empty wake right after the turn). private selfWakeTimers = new Map>(); + private selfWakeIntents = new Map(); // Privileged users who may wake the agent through sleep. Hot-reloaded from // privilegedUsersPath on change. @@ -610,6 +629,8 @@ export class EventGate { scriptTimeoutMs?: number; }) { this.configPath = opts.configPath; + const gateStem = basename(this.configPath, '.json'); + this.wakeStatePath = join(dirname(this.configPath), `${gateStem}.wake-intents.json`); this.privilegedUsersPath = opts.privilegedUsersPath; this.emitTrace = opts.emitTrace; this.addMessageFn = opts.addMessage; @@ -1565,6 +1586,8 @@ export class EventGate { if (selfWake) { clearTimeout(selfWake); this.selfWakeTimers.delete(agentName); + this.selfWakeIntents.delete(agentName); + this.tryPersistWakeState('cancelling superseded self-wake'); } } @@ -1652,15 +1675,17 @@ export class EventGate { /** Begin (or extend) a sleep window of `seconds`. Returns the wake time. */ setSleep(seconds: number, note?: string, agentName?: string): { until: number } { const ms = Math.max(0, Math.floor(seconds * 1000)); - this.sleepUntil = this.now() + ms; + const armedAt = this.now(); + this.sleepUntil = armedAt + ms; + this.sleepArmedAt = armedAt; this.sleepNote = note; this.sleepAgent = agentName; this.sleepSuppressed = 0; - // Arm the self-wake: when the window elapses, clear sleep and request an - // inference so the agent resumes on its own (not dependent on a heartbeat - // tick or an incoming message). Cancelled by clearSleep()/`wake`. - if (this.sleepTimer) clearTimeout(this.sleepTimer); - this.sleepTimer = setTimeout(() => this.wakeFromSleep('sleep-expired'), ms); + // Arm first: a persistence failure must never leave sleep set without a + // process-local wake. Persist before returning success whenever storage is + // available; failure is loud but the live timer remains authoritative. + this.armSleepTimer(); + this.tryPersistWakeState('arming sleep'); this.reloadPrivilegedIfChanged(); this.emitTrace({ type: 'gate:decision', @@ -1688,8 +1713,18 @@ export class EventGate { const ms = Math.max(1_000, Math.min(3_600_000, Math.floor(seconds * 1000))); const prev = this.selfWakeTimers.get(agentName); if (prev) clearTimeout(prev); + const intent: PersistedWakeIntent = { + agentName, + armedAt: this.now(), + wakeAt: this.now() + ms, + source, + }; + this.selfWakeIntents.set(agentName, intent); + this.tryPersistWakeState('arming self-wake'); const timer = setTimeout(() => { this.selfWakeTimers.delete(agentName); + this.selfWakeIntents.delete(agentName); + this.tryPersistWakeState('firing self-wake'); this.emitTrace({ type: 'gate:decision', eventType: 'selfwake:fired', @@ -1731,11 +1766,14 @@ export class EventGate { /** End sleep immediately. Returns true if the agent was asleep. */ clearSleep(): boolean { + const hadIntent = this.sleepUntil > 0; const wasAsleep = this.sleepUntil > this.now(); this.sleepUntil = 0; + this.sleepArmedAt = 0; this.sleepNote = undefined; this.sleepAgent = undefined; if (this.sleepTimer) { clearTimeout(this.sleepTimer); this.sleepTimer = null; } + if (hadIntent) this.tryPersistWakeState('clearing sleep'); return wasAsleep; } @@ -1746,14 +1784,31 @@ export class EventGate { * than just a suppression window. No-op if sleep was already cleared (via * `wake`) or re-armed to a later time by a fresh `sleep` call. */ + private armSleepTimer(): void { + if (this.sleepTimer) clearTimeout(this.sleepTimer); + const remaining = Math.max(0, this.sleepUntil - this.now()); + // Node clamps larger delays to ~1ms. Chunk long sleeps and check the wall + // deadline on each callback instead of firing years early. + const delay = Math.min(remaining, 0x7fffffff); + this.sleepTimer = setTimeout(() => this.wakeFromSleep('sleep-expired'), delay); + } + private wakeFromSleep(reason: string): void { this.sleepTimer = null; - if (this.sleepUntil === 0 || this.sleepUntil > this.now()) return; + if (this.sleepUntil === 0) return; + // A relative timer may fire before its wall deadline after a backward clock + // step. Re-arm for the remaining wall time instead of losing the promise. + if (this.sleepUntil > this.now()) { + this.armSleepTimer(); + return; + } const note = this.sleepNote; const agent = this.sleepAgent; this.sleepUntil = 0; + this.sleepArmedAt = 0; this.sleepNote = undefined; this.sleepAgent = undefined; + this.tryPersistWakeState('expiring sleep'); this.emitTrace({ type: 'gate:decision', eventType: 'sleep:expired', @@ -1767,7 +1822,152 @@ export class EventGate { for (const a of targets) this.requestInferenceFn(a, detail, 'sleep'); } - /** Current sleep state, or null if awake. */ + private validWakeIntent(value: unknown, allowAnonymous: boolean): value is PersistedWakeIntent { + if (!value || typeof value !== 'object') return false; + const v = value as Record; + return (allowAnonymous ? v.agentName === undefined || typeof v.agentName === 'string' : typeof v.agentName === 'string') + && typeof v.armedAt === 'number' && Number.isFinite(v.armedAt) + && typeof v.wakeAt === 'number' && Number.isFinite(v.wakeAt) + && typeof v.source === 'string' + && (v.note === undefined || typeof v.note === 'string') + && (v.suppressed === undefined || (typeof v.suppressed === 'number' && Number.isFinite(v.suppressed))); + } + + private readWakeState(): PersistedWakeState | null { + if (!existsSync(this.wakeStatePath)) return null; + try { + const raw = JSON.parse(readFileSync(this.wakeStatePath, 'utf8')) as Record; + if (raw.version !== 1 || !Array.isArray(raw.selfWakes)) throw new Error('unsupported wake journal shape'); + const selfWakes = raw.selfWakes.filter((entry) => { + const valid = this.validWakeIntent(entry, false); + if (!valid) console.error(`[gate] skipping invalid self-wake journal entry: ${JSON.stringify(entry)}`); + return valid; + }) as PersistedWakeIntent[]; + let sleep: PersistedWakeIntent | undefined; + if (raw.sleep !== undefined) { + if (this.validWakeIntent(raw.sleep, true)) sleep = raw.sleep; + else console.error(`[gate] skipping invalid sleep journal entry: ${JSON.stringify(raw.sleep)}`); + } + return { version: 1, ...(sleep ? { sleep } : {}), selfWakes }; + } catch (error) { + const quarantine = `${this.wakeStatePath}.corrupt-${this.now()}`; + try { renameSync(this.wakeStatePath, quarantine); } + catch (renameError) { + console.error(`[gate] wake journal unreadable and quarantine failed: ${renameError instanceof Error ? renameError.message : renameError}`); + } + console.error(`[gate] wake journal unreadable; quarantined as ${quarantine}: ${error instanceof Error ? error.message : error}`); + return null; + } + } + + private persistWakeState(): void { + const state: PersistedWakeState = { + version: 1, + ...(this.sleepUntil > 0 ? { + sleep: { + ...(this.sleepAgent ? { agentName: this.sleepAgent } : {}), + armedAt: this.sleepArmedAt, + wakeAt: this.sleepUntil, + source: 'sleep', + ...(this.sleepNote ? { note: this.sleepNote } : {}), + // Suppressed is a current-process diagnostic. Persisting every + // increment would turn each ignored event into a filesystem write. + suppressed: this.sleepSuppressed, + }, + } : {}), + selfWakes: [...this.selfWakeIntents.values()], + }; + mkdirSync(dirname(this.wakeStatePath), { recursive: true }); + const tmp = `${this.wakeStatePath}.tmp-${process.pid}`; + writeFileSync(tmp, JSON.stringify(state, null, 2) + '\n', { mode: 0o600 }); + renameSync(tmp, this.wakeStatePath); + } + + private tryPersistWakeState(action: string): boolean { + try { this.persistWakeState(); return true; } + catch (error) { + console.error(`[gate] failed to persist wake intent while ${action}: ${error instanceof Error ? error.message : error}`); + return false; + } + } + + /** Reconcile durable intents after AgentFramework initialization succeeds. */ + recoverWakeIntents(): void { + const state = this.readWakeState(); + if (!state) return; + const now = this.now(); + const due: Array<{ intent: PersistedWakeIntent; kind: 'sleep' | 'self-wake' }> = []; + + if (state.sleep) { + if (state.sleep.wakeAt > now) { + this.sleepUntil = state.sleep.wakeAt; + this.sleepArmedAt = state.sleep.armedAt; + this.sleepNote = state.sleep.note; + this.sleepAgent = state.sleep.agentName; + // Suppressed is intentionally process-local; a recovered process starts + // counting again from zero (also reflected by getSleepState()). + this.sleepSuppressed = 0; + this.armSleepTimer(); + } else due.push({ intent: state.sleep, kind: 'sleep' }); + } + for (const intent of state.selfWakes) { + if (intent.wakeAt > now) { + this.selfWakeIntents.set(intent.agentName!, intent); + this.armRestoredSelfWake(intent); + } else due.push({ intent, kind: 'self-wake' }); + } + if (due.length === 0) return; + + // Recovery is called only after full framework initialization. Clear each + // due item after requestInference returns (the framework callback enqueued + // it); if delivery throws, retain the durable item for the next boot. + const woken = new Set(); + for (const item of due) { + const targets = item.intent.agentName ? [item.intent.agentName] : this.getAgentNamesFn(); + for (const target of targets) { + const overdueMs = Math.max(0, now - item.intent.wakeAt); + const note = item.intent.note ? ` Note: ${item.intent.note}` : ''; + this.addMessageFn('user', [{ type: 'text', text: + `[wake-recovery] your ${item.kind === 'self-wake' ? item.intent.source : item.kind} wake ` + + `scheduled for ${new Date(item.intent.wakeAt).toISOString()} became overdue while the host was offline by ${overdueMs}ms.${note}` }], + { source: 'gate:wake-recovery', kind: item.kind }, target); + if (!woken.has(target)) { + this.requestInferenceFn(target, `wake intent recovery${item.intent.note ? `: ${item.intent.note}` : ''}`, 'wake-recovery'); + woken.add(target); + } + } + if (item.kind === 'sleep') { + this.sleepUntil = 0; this.sleepArmedAt = 0; this.sleepAgent = undefined; this.sleepNote = undefined; + } else this.selfWakeIntents.delete(item.intent.agentName!); + this.tryPersistWakeState('admitting overdue wake recovery'); + } + } + + private armRestoredSelfWake(intent: PersistedWakeIntent): void { + const remaining = Math.max(0, intent.wakeAt - this.now()); + const delay = Math.min(remaining, 0x7fffffff); + const timer = setTimeout(() => { + if (intent.wakeAt > this.now()) { this.armRestoredSelfWake(intent); return; } + this.fireRestoredSelfWake(intent); + }, delay); + this.selfWakeTimers.set(intent.agentName!, timer); + } + + private fireRestoredSelfWake(intent: PersistedWakeIntent): void { + this.selfWakeTimers.delete(intent.agentName!); + this.selfWakeIntents.delete(intent.agentName!); + this.tryPersistWakeState('firing restored self-wake'); + const nowIso = new Date(this.now()).toISOString().slice(0, 19) + 'Z'; + this.addMessageFn('user', [{ type: 'text', text: + `[self-wake] your ${intent.source} timer (${Math.round((intent.wakeAt - intent.armedAt) / 1000)}s) elapsed — now ${nowIso}` }], + { source: 'gate:self-wake' }, intent.agentName); + this.requestInferenceFn(intent.agentName!, + `self-scheduled wake (${intent.source}, ${Math.round((intent.wakeAt - intent.armedAt) / 1000)}s)`, + 'self-wake'); + } + + /** Current sleep state, or null if awake. `suppressed` counts only events + * observed by the current process; it resets to zero after recovery. */ getSleepState(): { until: number; remainingMs: number; note?: string; suppressed: number } | null { const remainingMs = this.sleepUntil - this.now(); if (remainingMs <= 0) return null; @@ -1925,7 +2125,11 @@ export class EventGate { for (const timer of this.selfWakeTimers.values()) { clearTimeout(timer); } + // Deliberate shutdown preserves the durable intent exactly as armed. + // Only process-local timer handles are released; startup follows the same + // future/re-arm or overdue/reconcile path as an abrupt process loss. this.selfWakeTimers.clear(); + if (this.sleepTimer) { clearTimeout(this.sleepTimer); this.sleepTimer = null; } this.rateLimitBuckets.clear(); this.passiveSampleCounters.clear(); this.rateLimitDenied.clear(); diff --git a/test/event-gate-selfwake.test.ts b/test/event-gate-selfwake.test.ts index 37cc4533..24d6e88f 100644 --- a/test/event-gate-selfwake.test.ts +++ b/test/event-gate-selfwake.test.ts @@ -14,9 +14,11 @@ import { describe, it, beforeEach, afterEach } from 'node:test'; import assert from 'node:assert'; -import { existsSync, mkdirSync, rmSync } from 'node:fs'; +import { existsSync, mkdirSync, rmSync, readFileSync, writeFileSync } from 'node:fs'; import { join } from 'node:path'; import { EventGate } from '../src/gate/event-gate.js'; +import { AgentFramework } from '../src/framework.js'; +import { MockMembrane } from './helpers/mock-membrane.js'; const TMP_DIR = join(import.meta.dirname, '../.test-tmp-gate-selfwake'); @@ -133,3 +135,231 @@ describe('EventGate self-wake', () => { assert.strictEqual(inferenceRequests.length, 0); }); }); + +describe('durable wake intent recovery', () => { + beforeEach(() => { + if (existsSync(TMP_DIR)) rmSync(TMP_DIR, { recursive: true }); + mkdirSync(TMP_DIR, { recursive: true }); + }); + afterEach(() => { + if (existsSync(TMP_DIR)) rmSync(TMP_DIR, { recursive: true }); + }); + + function harness(now: () => number) { + const inferenceRequests: Array<{ agentName: string; reason: string; source: string }> = []; + const messages: string[] = []; + const gate = new EventGate({ + configPath: join(TMP_DIR, 'gate.json'), now, + emitTrace: () => {}, + addMessage: (_p, c) => { messages.push(c.map((b) => b.text).join('\n')); return ''; }, + requestInference: (agentName, reason, source) => inferenceRequests.push({ agentName, reason, source }), + getAgentNames: () => ['agent'], + }); + gate.recoverWakeIntents(); + return { gate, inferenceRequests, messages }; + } + + it('persists sleep before setSleep reports success and re-arms it after restart', () => { + let now = 1_000; + const a = harness(() => now); + const { until } = a.gate.setSleep(60, 'resume work', 'agent'); + const state = JSON.parse(readFileSync(join(TMP_DIR, 'gate.wake-intents.json'), 'utf8')); + assert.strictEqual(state.sleep.wakeAt, until); + assert.strictEqual(state.sleep.note, 'resume work'); + // Simulate abrupt process loss: stop only the old JS timer, without dispose. + const ai = a.gate as unknown as { sleepTimer: ReturnType | null }; + if (ai.sleepTimer) clearTimeout(ai.sleepTimer); + + now += 10_000; + const b = harness(() => now); + assert.strictEqual(b.gate.getSleepState()!.until, until); + b.gate.dispose(); + }); + + it('an overdue sleep emits one recovery marker and one inference', async () => { + let now = 1_000; + const a = harness(() => now); + a.gate.setSleep(5, 'check result', 'agent'); + const ai = a.gate as unknown as { sleepTimer: ReturnType | null }; + if (ai.sleepTimer) clearTimeout(ai.sleepTimer); + + now = 10_000; + const b = harness(() => now); + await Promise.resolve(); + assert.strictEqual(b.inferenceRequests.length, 1); + assert.strictEqual(b.inferenceRequests[0]!.source, 'wake-recovery'); + assert.strictEqual(b.messages.length, 1); + assert.match(b.messages[0]!, /wake-recovery.*overdue while the host was offline/); + b.gate.dispose(); + + const c = harness(() => now); + await Promise.resolve(); + assert.strictEqual(c.inferenceRequests.length, 0, 'recovery is admitted at most once'); + c.gate.dispose(); + }); + + it('graceful shutdown preserves a future sleep and restart re-arms it', async () => { + let now = 1_000; + const a = harness(() => now); + const { until } = a.gate.setSleep(60, 'resume after deploy', 'agent'); + a.gate.dispose(); + + const persisted = JSON.parse(readFileSync(join(TMP_DIR, 'gate.wake-intents.json'), 'utf8')); + assert.strictEqual(persisted.sleep.wakeAt, until, 'dispose must leave the durable promise intact'); + const b = harness(() => now); + await Promise.resolve(); + assert.strictEqual(b.inferenceRequests.length, 0, 'a planned restart must not wake early'); + assert.strictEqual(b.messages.length, 0, 'a planned restart must not invent a cancellation marker'); + assert.strictEqual(b.gate.getSleepState()!.until, until, 'startup re-arms the original deadline'); + b.gate.dispose(); + }); + + it('persists and re-arms skip_reply self-wake across abrupt restart', async () => { + const a = harness(Date.now); + a.gate.armSelfWake('agent', 1, 'skip_reply'); + const ai = a.gate as unknown as { selfWakeTimers: Map> }; + for (const timer of ai.selfWakeTimers.values()) clearTimeout(timer); + + const b = harness(Date.now); + await sleep(1150); + assert.strictEqual(b.inferenceRequests.length, 1); + assert.strictEqual(b.inferenceRequests[0]!.source, 'self-wake'); + assert.match(b.messages[0]!, /your skip_reply timer/); + b.gate.dispose(); + }); + + + it('framework invokes recovery after resident agents are registered', async () => { + const now = Date.now(); + const configPath = join(TMP_DIR, 'gate.json'); + writeFileSync(join(TMP_DIR, 'gate.wake-intents.json'), JSON.stringify({ + version: 1, + sleep: { agentName: 'agent', armedAt: now - 10_000, wakeAt: now - 1_000, source: 'sleep' }, + selfWakes: [], + })); + const membrane = new MockMembrane(); + const framework = await AgentFramework.create({ + storePath: join(TMP_DIR, 'store'), + membrane: membrane.asMembrane(), + agents: [{ name: 'agent', model: 'test-model', systemPrompt: 'test' }], + modules: [], + gate: { configPath }, + }); + await Promise.resolve(); + + const messages = framework.getAgent('agent')!.getContextManager().getAllMessages(); + const text = messages.flatMap((m) => m.content) + .filter((b): b is { type: 'text'; text: string } => b.type === 'text') + .map((b) => b.text).join('\n'); + assert.match(text, /wake-recovery.*overdue while the host was offline/, + 'recovery marker must reach the already-registered target agent'); + const pending = (framework as unknown as { pendingRequests: Array<{ agentName: string; source: string }> }) + .pendingRequests; + assert.ok(pending.some((r) => r.agentName === 'agent' && r.source === 'wake-recovery')); + await framework.stop(); + }); + + it('failed Framework.create leaves a due intent in the journal', async () => { + const now = Date.now(); + const configPath = join(TMP_DIR, 'gate.json'); + const journalPath = join(TMP_DIR, 'gate.wake-intents.json'); + writeFileSync(journalPath, JSON.stringify({ + version: 1, + sleep: { agentName: 'agent', armedAt: now - 10_000, wakeAt: now - 1_000, source: 'sleep' }, + selfWakes: [], + })); + const membrane = new MockMembrane(); + await assert.rejects(AgentFramework.create({ + storePath: join(TMP_DIR, 'failed-store'), + membrane: membrane.asMembrane(), + agents: [{ name: 'agent', model: 'test-model', systemPrompt: 'test' }], + modules: [], gate: { configPath }, + toolResultInlineMaxChars: 10, // validated after gate construction + }), /toolResultInlineMaxChars/); + const state = JSON.parse(readFileSync(journalPath, 'utf8')); + assert.ok(state.sleep, 'failed create must not consume the due wake'); + }); + + it('quarantines a corrupt journal instead of overwriting it', () => { + const journalPath = join(TMP_DIR, 'gate.wake-intents.json'); + writeFileSync(journalPath, '{bad json'); + const h = harness(() => 12345); + assert.ok(existsSync(`${journalPath}.corrupt-12345`)); + assert.strictEqual(existsSync(journalPath), false, 'corrupt original was moved aside'); + h.gate.dispose(); + }); + + it('chunks timers longer than the Node timeout maximum', () => { + let now = 1_000; + const h = harness(() => now); + h.gate.setSleep((0x7fffffff + 60_000) / 1000, undefined, 'agent'); + const timer = (h.gate as unknown as { sleepTimer: { _idleTimeout?: number } }).sleepTimer; + assert.strictEqual(timer._idleTimeout, 0x7fffffff); + h.gate.dispose(); + }); + + it('skips invalid journal entries while recovering valid ones', () => { + const journalPath = join(TMP_DIR, 'gate.wake-intents.json'); + writeFileSync(journalPath, JSON.stringify({ + version: 1, + sleep: { agentName: 7, armedAt: 'bad', wakeAt: null, source: 4 }, + selfWakes: [ + { agentName: false, armedAt: 1, wakeAt: 2, source: 'skip_reply' }, + { agentName: 'agent', armedAt: 1_000, wakeAt: 61_000, source: 'skip_reply' }, + ], + })); + const h = harness(() => 2_000); + assert.strictEqual(h.gate.getSleepState(), null); + const intents = (h.gate as unknown as { selfWakeIntents: Map }).selfWakeIntents; + assert.deepStrictEqual([...intents.keys()], ['agent']); + h.gate.dispose(); + }); + + it('constructor with no journal performs no wake-journal write', () => { + const path = join(TMP_DIR, 'gate.wake-intents.json'); + const h = harness(() => 1_000); + assert.strictEqual(existsSync(path), false); + h.gate.dispose(); + assert.strictEqual(existsSync(path), false); + }); + + it('recovers an anonymous overdue sleep to every registered agent and includes its note', () => { + const messages: string[] = []; + const requests: string[] = []; + writeFileSync(join(TMP_DIR, 'gate.wake-intents.json'), JSON.stringify({ + version: 1, + sleep: { armedAt: 1, wakeAt: 2, source: 'sleep', note: 'check all residents' }, + selfWakes: [], + })); + const gate = new EventGate({ + configPath: join(TMP_DIR, 'gate.json'), now: () => 10, + emitTrace: () => {}, + addMessage: (_p, c, _m, agent) => { messages.push(`${agent}:${c[0]!.text}`); return ''; }, + requestInference: (agent, reason) => requests.push(`${agent}:${reason}`), + getAgentNames: () => ['a', 'b'], + }); + gate.recoverWakeIntents(); + assert.strictEqual(messages.length, 2); + assert.ok(messages.every((m) => m.includes('Note: check all residents'))); + assert.deepStrictEqual(requests.map((r) => r.split(':')[0]), ['a', 'b']); + assert.ok(requests.every((r) => r.includes('check all residents'))); + gate.dispose(); + }); + + it('re-arms a sleep timer that fires before its wall-clock deadline', () => { + let now = 10_000; + const h = harness(() => now); + h.gate.setSleep(10, undefined, 'agent'); + const internals = h.gate as unknown as { + sleepTimer: ReturnType | null; + wakeFromSleep(reason: string): void; + }; + if (internals.sleepTimer) clearTimeout(internals.sleepTimer); + now = 5_000; // backward wall-clock step + internals.wakeFromSleep('sleep-expired'); + assert.ok(internals.sleepTimer, 'early callback must arm a replacement timer'); + assert.strictEqual(h.inferenceRequests.length, 0); + if (internals.sleepTimer) clearTimeout(internals.sleepTimer); + h.gate.dispose(); + }); +}); From 2adb3921336c956ef372fe993e4690788ba33c17 Mon Sep 17 00:00:00 2001 From: qa-engineer Date: Wed, 30 Sep 2026 20:21:02 +0000 Subject: [PATCH 2/3] fix(gate): unref sleep and self-wake timers --- src/framework.ts | 10 +++-- src/gate/event-gate.ts | 54 +++++++++++++------------ test/event-gate-selfwake.test.ts | 68 ++++++++++++++++++++++++++++++++ 3 files changed, 102 insertions(+), 30 deletions(-) diff --git a/src/framework.ts b/src/framework.ts index d834a7e9..4bc9a6cc 100644 --- a/src/framework.ts +++ b/src/framework.ts @@ -1768,6 +1768,12 @@ export class AgentFramework { } } + // Restore sleep suppression before MCPL startup can release buffered input. + // Recovery inference is only queued here: processInferenceRequests runs from + // start()/the event loop after create() finishes, so the recovered turn sees + // every MCPL tool registered by initializeMcpl below. + framework.eventGate?.recoverWakeIntents(); + // Initialize MCPL subsystems if configured if (config.mcplServers && config.mcplServers.length > 0) { // Validate tool prefixes: no collisions with module names or between servers @@ -1851,10 +1857,6 @@ export class AgentFramework { // SIGUSR2 not available on this platform — non-fatal. } - // Reconcile wake promises only after every module/server/startup invariant - // above succeeded. A failed create must leave due intents durable. - framework.eventGate?.recoverWakeIntents(); - return framework; } diff --git a/src/gate/event-gate.ts b/src/gate/event-gate.ts index bde43de9..3ffe5038 100644 --- a/src/gate/event-gate.ts +++ b/src/gate/event-gate.ts @@ -627,6 +627,8 @@ export class EventGate { now?: () => number; /** Per-event timeout (ms) for the optional gate.js script. Default 50. */ scriptTimeoutMs?: number; + /** Recover durable wake intents during construction (default false). */ + autoRecover?: boolean; }) { this.configPath = opts.configPath; const gateStem = basename(this.configPath, '.json'); @@ -692,6 +694,7 @@ export class EventGate { opts.scriptTimeoutMs ?? 50, this.now, ); + if (opts.autoRecover === true) this.recoverWakeIntents(); } // ========================================================================= @@ -1891,38 +1894,41 @@ export class EventGate { } } - /** Reconcile durable intents after AgentFramework initialization succeeds. */ + /** + * Reconcile durable wake intents once an embedding host can accept messages + * and inference requests. Hosts should call this before releasing inbound + * traffic; `autoRecover` is available for standalone embedders. + */ recoverWakeIntents(): void { const state = this.readWakeState(); if (!state) return; const now = this.now(); - const due: Array<{ intent: PersistedWakeIntent; kind: 'sleep' | 'self-wake' }> = []; - if (state.sleep) { - if (state.sleep.wakeAt > now) { - this.sleepUntil = state.sleep.wakeAt; - this.sleepArmedAt = state.sleep.armedAt; - this.sleepNote = state.sleep.note; - this.sleepAgent = state.sleep.agentName; - // Suppressed is intentionally process-local; a recovered process starts - // counting again from zero (also reflected by getSleepState()). - this.sleepSuppressed = 0; - this.armSleepTimer(); - } else due.push({ intent: state.sleep, kind: 'sleep' }); + this.sleepUntil = state.sleep.wakeAt; this.sleepArmedAt = state.sleep.armedAt; + this.sleepNote = state.sleep.note; this.sleepAgent = state.sleep.agentName; this.sleepSuppressed = 0; + if (state.sleep.wakeAt > now) this.armSleepTimer(); } for (const intent of state.selfWakes) { - if (intent.wakeAt > now) { - this.selfWakeIntents.set(intent.agentName!, intent); - this.armRestoredSelfWake(intent); - } else due.push({ intent, kind: 'self-wake' }); + this.selfWakeIntents.set(intent.agentName!, intent); + if (intent.wakeAt > now) this.armRestoredSelfWake(intent); } - if (due.length === 0) return; - - // Recovery is called only after full framework initialization. Clear each - // due item after requestInference returns (the framework callback enqueued - // it); if delivery throws, retain the durable item for the next boot. + const due: Array<{ intent: PersistedWakeIntent; kind: 'sleep' | 'self-wake' }> = []; + if (state.sleep && state.sleep.wakeAt <= now) due.push({ intent: state.sleep, kind: 'sleep' }); + for (const intent of state.selfWakes) if (intent.wakeAt <= now) due.push({ intent, kind: 'self-wake' }); const woken = new Set(); for (const item of due) { + if (item.kind === 'sleep') { + this.sleepUntil = 0; this.sleepArmedAt = 0; this.sleepAgent = undefined; this.sleepNote = undefined; + } else this.selfWakeIntents.delete(item.intent.agentName!); + // Journal admission precedes side effects. Failure restores the intent + // in memory and leaves its durable copy for the next startup. + if (!this.tryPersistWakeState('admitting overdue wake recovery')) { + if (item.kind === 'sleep') { + this.sleepUntil = item.intent.wakeAt; this.sleepArmedAt = item.intent.armedAt; + this.sleepAgent = item.intent.agentName; this.sleepNote = item.intent.note; + } else this.selfWakeIntents.set(item.intent.agentName!, item.intent); + continue; + } const targets = item.intent.agentName ? [item.intent.agentName] : this.getAgentNamesFn(); for (const target of targets) { const overdueMs = Math.max(0, now - item.intent.wakeAt); @@ -1936,10 +1942,6 @@ export class EventGate { woken.add(target); } } - if (item.kind === 'sleep') { - this.sleepUntil = 0; this.sleepArmedAt = 0; this.sleepAgent = undefined; this.sleepNote = undefined; - } else this.selfWakeIntents.delete(item.intent.agentName!); - this.tryPersistWakeState('admitting overdue wake recovery'); } } diff --git a/test/event-gate-selfwake.test.ts b/test/event-gate-selfwake.test.ts index 24d6e88f..4ae44b5f 100644 --- a/test/event-gate-selfwake.test.ts +++ b/test/event-gate-selfwake.test.ts @@ -131,6 +131,12 @@ describe('EventGate self-wake', () => { const { gate, inferenceRequests } = makeGate(); gate.armSelfWake('agent', 1); gate.dispose(); + const internals = gate as unknown as { + sleepTimer: ReturnType | null; + selfWakeTimers: Map>; + }; + assert.strictEqual(internals.sleepTimer, null); + assert.strictEqual(internals.selfWakeTimers.size, 0); await sleep(1150); assert.strictEqual(inferenceRequests.length, 0); }); @@ -315,6 +321,18 @@ describe('durable wake intent recovery', () => { h.gate.dispose(); }); + it('autoRecover lets standalone embedders restore during construction', () => { + writeFileSync(join(TMP_DIR, 'gate.wake-intents.json'), JSON.stringify({ + version: 1, sleep: { agentName: 'agent', armedAt: 1, wakeAt: 60_000, source: 'sleep' }, selfWakes: [], + })); + const gate = new EventGate({ + configPath: join(TMP_DIR, 'gate.json'), now: () => 1_000, autoRecover: true, + emitTrace: () => {}, addMessage: () => '', requestInference: () => {}, getAgentNames: () => ['agent'], + }); + assert.strictEqual(gate.getSleepState()!.until, 60_000); + gate.dispose(); + }); + it('constructor with no journal performs no wake-journal write', () => { const path = join(TMP_DIR, 'gate.wake-intents.json'); const h = harness(() => 1_000); @@ -363,3 +381,53 @@ describe('durable wake intent recovery', () => { h.gate.dispose(); }); }); + +describe('wake recovery startup ordering', () => { + it('AgentFramework.create recovers wake intent before MCPL initialization', async () => { + const order: string[] = []; + const gateProto = EventGate.prototype as unknown as { recoverWakeIntents(): void }; + const frameworkProto = AgentFramework.prototype as unknown as { + initializeMcpl(configs: unknown[], routing?: unknown): Promise; + }; + const originalRecover = gateProto.recoverWakeIntents; + const originalMcpl = frameworkProto.initializeMcpl; + gateProto.recoverWakeIntents = function () { order.push('recover'); }; + frameworkProto.initializeMcpl = async function () { order.push('mcpl'); }; + const dir = join(TMP_DIR, 'ordering'); + try { + const framework = await AgentFramework.create({ + storePath: join(dir, 'store'), membrane: new MockMembrane().asMembrane(), + agents: [{ name: 'agent', model: 'test', systemPrompt: 'test' }], modules: [], + gate: { configPath: join(dir, 'gate.json') }, + mcplServers: [{ id: 'fixture', command: 'unused', args: [] }], + }); + assert.deepStrictEqual(order, ['recover', 'mcpl']); + await framework.stop(); + } finally { + gateProto.recoverWakeIntents = originalRecover; + frameworkProto.initializeMcpl = originalMcpl; + } + }); + + it('does not emit an overdue marker or inference when journal admission fails', () => { + const messages: string[] = []; + const requests: string[] = []; + mkdirSync(TMP_DIR, { recursive: true }); + writeFileSync(join(TMP_DIR, 'gate.wake-intents.json'), JSON.stringify({ + version: 1, + sleep: { agentName: 'agent', armedAt: 1, wakeAt: 2, source: 'sleep' }, selfWakes: [], + })); + const gate = new EventGate({ + configPath: join(TMP_DIR, 'gate.json'), now: () => 10, + emitTrace: () => {}, addMessage: (_p, c) => { messages.push(c[0]!.text); return ''; }, + requestInference: (_a, reason) => requests.push(reason), getAgentNames: () => ['agent'], + }); + (gate as unknown as { persistWakeState(): void }).persistWakeState = () => { throw new Error('disk full'); }; + gate.recoverWakeIntents(); + assert.deepStrictEqual(messages, []); + assert.deepStrictEqual(requests, []); + const state = JSON.parse(readFileSync(join(TMP_DIR, 'gate.wake-intents.json'), 'utf8')); + assert.ok(state.sleep, 'durable due intent remains for next startup'); + gate.dispose(); + }); +}); From 61ca6cd346162a600ee883555643c71cde0bb539 Mon Sep 17 00:00:00 2001 From: qa-engineer Date: Mon, 5 Oct 2026 19:35:23 +0000 Subject: [PATCH 3/3] docs(gate): clarify wake journal agent scope --- src/gate/event-gate.ts | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/gate/event-gate.ts b/src/gate/event-gate.ts index 3ffe5038..663a4c92 100644 --- a/src/gate/event-gate.ts +++ b/src/gate/event-gate.ts @@ -577,7 +577,9 @@ export class EventGate { private sleepTimer: ReturnType | null = null; private sleepAgent: string | undefined; private sleepArmedAt = 0; - /** Branch-independent wake-intent journal, sibling to gate.json. */ + /** Branch-independent wake-intent journal, sibling to gate.json. Sleep is + * gate-global (unchanged semantics): its agentName is informational for the + * eventual wake target; self-wake agentName is authoritative per-agent state. */ private wakeStatePath: string; // Self-wake timers per agent (skip_reply's wake_in_seconds). Unlike sleep,