Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions changelog.d/durable-wake-intents.fixed.md
Original file line number Diff line number Diff line change
@@ -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.
6 changes: 6 additions & 0 deletions src/framework.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1746,6 +1746,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();

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Failed startup can erase wakes If a persisted wake is still in the future, recovery arms its timer before MCPL configuration is validated. If validation or initialization then fails, the abandoned framework does not dispose of that timer. When it fires, it clears the saved wake and queues inference on a framework whose loop never starts, so a later successful startup cannot recover the promised wake.

Knowledge Base Used: Persistent state and recovery

Prompt To Fix With AI
This is a comment left during a code review.
Path: src/framework.ts
Line: 1661

Comment:
**Failed startup can erase wakes** If a persisted wake is still in the future, recovery arms its timer before MCPL configuration is validated. If validation or initialization then fails, the abandoned framework does not dispose of that timer. When it fires, it clears the saved wake and queues inference on a framework whose loop never starts, so a later successful startup cannot recover the promised wake.

**Knowledge Base Used:** [Persistent state and recovery](https://app.greptile.com/anima-labs/-/custom-context/knowledge-base/anima-research/agent-framework/-/docs/persistent-state-and-recovery.md)

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.

Comment on lines +1749 to +1753

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Restored batches bypass sleep

When a resident has a saved future sleep and a deferred MCPL push batch was pending at shutdown, startup restores the sleep but then queues the batch’s inference request without checking the gate. Starting the framework can run the resident before its promised wake deadline, even without new server input.

Knowledge Base Used: Event gate pipeline

Prompt To Fix With AI
This is a comment left during a code review.
Path: src/framework.ts
Line: 1749-1750

Comment:
**Restored batches bypass sleep**

When a resident has a saved future sleep and a deferred MCPL push batch was pending at shutdown, startup restores the sleep but then queues the batch’s inference request without checking the gate. Starting the framework can run the resident before its promised wake deadline, even without new server input.

**Knowledge Base Used:** [Event gate pipeline](https://app.greptile.com/anima-labs/-/custom-context/knowledge-base/anima-research/agent-framework/-/docs/event-gates.md)

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.


// Initialize MCPL subsystems if configured
if (config.mcplServers && config.mcplServers.length > 0) {
// Validate tool prefixes: no collisions with module names or between servers
Expand Down
228 changes: 218 additions & 10 deletions src/gate/event-gate.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -561,13 +576,19 @@ export class EventGate {
private sleepSuppressed = 0;
private sleepTimer: ReturnType<typeof setTimeout> | null = null;
private sleepAgent: string | undefined;
private sleepArmedAt = 0;
/** 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,
// a self-wake suppresses NOTHING — it means "if nothing else wakes me by
// then, wake me". Any turn start for the agent cancels its pending
// 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<string, ReturnType<typeof setTimeout>>();
private selfWakeIntents = new Map<string, PersistedWakeIntent>();

// Privileged users who may wake the agent through sleep. Hot-reloaded from
// privilegedUsersPath on change.
Expand Down Expand Up @@ -608,8 +629,12 @@ 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');
this.wakeStatePath = join(dirname(this.configPath), `${gateStem}.wake-intents.json`);
this.privilegedUsersPath = opts.privilegedUsersPath;
this.emitTrace = opts.emitTrace;
this.addMessageFn = opts.addMessage;
Expand Down Expand Up @@ -671,6 +696,7 @@ export class EventGate {
opts.scriptTimeoutMs ?? 50,
this.now,
);
if (opts.autoRecover === true) this.recoverWakeIntents();
}

// =========================================================================
Expand Down Expand Up @@ -1565,6 +1591,8 @@ export class EventGate {
if (selfWake) {
clearTimeout(selfWake);
this.selfWakeTimers.delete(agentName);
this.selfWakeIntents.delete(agentName);
this.tryPersistWakeState('cancelling superseded self-wake');
}
}

Expand Down Expand Up @@ -1652,15 +1680,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',
Expand Down Expand Up @@ -1688,8 +1718,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',
Expand Down Expand Up @@ -1731,11 +1771,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;
}

Expand All @@ -1746,14 +1789,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',
Expand All @@ -1767,7 +1827,151 @@ 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<string, unknown>;
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<string, unknown>;
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,
Comment thread
greptile-apps[bot] marked this conversation as resolved.
},
} : {}),
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 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 {
Comment thread
greptile-apps[bot] marked this conversation as resolved.
const state = this.readWakeState();
if (!state) return;
const now = this.now();
if (state.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) {
this.selfWakeIntents.set(intent.agentName!, intent);
if (intent.wakeAt > now) this.armRestoredSelfWake(intent);
Comment on lines +1913 to +1915

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Repeated recovery duplicates self-wakes If an embedder enables autoRecover and later calls the public recoverWakeIntents() method, this path arms the same future self-wake again without cancelling its first timer. Both timers can fire and each sends a notice and requests inference, turning one fallback wake into duplicate requests.

Knowledge Base Used: Event gate policies and scripts

Prompt To Fix With AI
This is a comment left during a code review.
Path: src/gate/event-gate.ts
Line: 1911-1913

Comment:
**Repeated recovery duplicates self-wakes** If an embedder enables `autoRecover` and later calls the public `recoverWakeIntents()` method, this path arms the same future self-wake again without cancelling its first timer. Both timers can fire and each sends a notice and requests inference, turning one fallback wake into duplicate requests.

**Knowledge Base Used:** [Event gate policies and scripts](https://app.greptile.com/anima-labs/-/custom-context/knowledge-base/anima-research/agent-framework/-/docs/event-gate-policies.md)

---

For each issue above, determine whether it is valid and should be fixed. If so, fix it directly.

}
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<string>();
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);
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);
}
}
}
}

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;
Expand Down Expand Up @@ -1925,7 +2129,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();
Expand Down
Loading