diff --git a/README.md b/README.md index eb5c9f8..1244728 100644 --- a/README.md +++ b/README.md @@ -50,6 +50,7 @@ OrchestrationGet id="1" - Cron, event, hybrid, dynamic goal, and opt-in workflow loops - Idle-safe agent re-wakes with dynamic-loop restart/session-switch recovery +- Observable seven-day recurring-loop expiry through `LoopList`, `loops:expired`, and a hidden Pi notification - Background command monitoring with buffered output, `onDone` wakes, and renewable inactivity alerts - Optional `pi-tasks` integration and a native task fallback - Session-scoped, bounded async subagent orchestration through protocol-v2 `pi-subagents` diff --git a/docs/REFERENCE.md b/docs/REFERENCE.md index 75c2da4..9627f87 100644 --- a/docs/REFERENCE.md +++ b/docs/REFERENCE.md @@ -48,9 +48,9 @@ Project scope shares durable state but does not yet elect one scheduler owner ac - hybrid: `{type:"hybrid", cron, event, debounceMs}` - dynamic: `{type:"dynamic"}` -`LoopCreate` creates ordinary controllers. Use dynamic loops for one evolving goal without named phase/outcome routing; use workflows for ordered phases, conditional outcomes, rework, or durable handoff; use standalone tasks for independently completable backlog items. `LoopUpdate` is only for dynamic controllers that are neither workflow nor orchestration owned and must persist `continue` after empty or unchanged iterations while work remains. `LoopDelete` pauses or removes ordinary/workflow controllers and cancellation-fences orchestration before stopping its workers. Recurring loops expire after seven days unless recreated; fire limits bound repeated execution. +`LoopCreate` creates ordinary controllers. Use dynamic loops for one evolving goal without named phase/outcome routing; use workflows for ordered phases, conditional outcomes, rework, or durable handoff; use standalone tasks for independently completable backlog items. `LoopUpdate` is only for dynamic controllers that are neither workflow nor orchestration owned and must persist `continue` after empty or unchanged iterations while work remains. `LoopDelete` pauses or removes ordinary/workflow controllers and cancellation-fences orchestration before stopping its workers. Recurring loops expire after seven days unless explicitly recreated; `LoopList` exposes the ISO `expiresAt` boundary. Seven-day expiry and stale event/hybrid retirement during session recovery emit `loops:expired` plus a hidden notification with `deleted` or `paused` disposition. Fire limits bound repeated execution. -Wake delivery is idle-driven. A due timer or event mutates loop state, emits `loop:fire`, buffers a generation-tagged notification, and sends a hidden Pi message when delivery is safe. Stale extension contexts are probed before fire mutation. +Wake delivery is idle-driven. A due timer or event mutates loop state, emits `loop:fire`, buffers a generation-tagged notification, and sends a hidden Pi message when delivery is safe. Retirement follows the same generation-fenced notification path after the store mutation and emits a typed `loops:expired` payload whose reason distinguishes `expires_at` from `resume_event_stale`. Pending fire and retirement notifications are memory-only; they are not a durable event ledger across process death. Stale extension contexts are probed before fire mutation. ## Subagent orchestration model diff --git a/docs/USAGE_GUIDE.md b/docs/USAGE_GUIDE.md index f9ce9c7..4519a3f 100644 --- a/docs/USAGE_GUIDE.md +++ b/docs/USAGE_GUIDE.md @@ -15,7 +15,7 @@ LoopCreate trigger="0 9 * * 1-5" prompt="Review weekday alerts" maxFires=10 Intervals such as `5m`, `2h`, and `1d` are converted to cron expressions. Full five-field cron expressions are also accepted. Cron and hybrid loops track their next fire time and deliver only when the agent is idle. -Use `maxFires` for polling or other bounded work so a loop cannot run indefinitely. Recurring loops expire after seven days. +Use `maxFires` for polling or other bounded work so a loop cannot run indefinitely. Recurring loops expire after seven days. `LoopList` exposes each controller's exact `expiresAt` boundary. Seven-day expiry and stale event/hybrid retirement during session recovery emit `loops:expired` and a hidden Pi wake that reports whether the controller was deleted or paused. Renewal is explicit: recreate the loop only when its schedule is still required. ### Event loops @@ -244,6 +244,11 @@ Monitor events: - `monitor:done` - `monitor:error` +Loop lifecycle events: + +- `loops:expired` — `{loopId, prompt, trigger, recurring, createdAt, expiresAt, expiredAt, disposition, source, reason}` where `reason` is `expires_at` or `resume_event_stale` +- `loops:autodeleted` + Native task lifecycle events: - `tasks:created` @@ -254,7 +259,6 @@ Native task lifecycle events: - `tasks:updated` - `tasks:deleted` - `tasks:backlog_empty` -- `loops:autodeleted` Task event payloads include `previousStatus`. Transition events report the status before the transition; details-only `tasks:updated` events report the status current at edit time. diff --git a/src/api.ts b/src/api.ts index c009bd9..a9b1e4c 100644 --- a/src/api.ts +++ b/src/api.ts @@ -35,12 +35,16 @@ export { rpcCall, rpcProbe, } from "./rpc/cross-extension-rpc.js"; +export type { LoopExpiredPayload } from "./runtime/loop-events.js"; export { NATIVE_TASKS_PROVIDER } from "./runtime/native-task-rpc.js"; export { resolveLoopStorePath, resolveTaskStorePath } from "./runtime/scope.js"; export type { TaskClaimInput, TaskClaimResult } from "./task-store.js"; export { TaskStore } from "./task-store.js"; export type { TaskClaim, TaskEntry, TaskStatus, TaskStoreData } from "./task-types.js"; export type { + LoopExpiryDisposition, + LoopExpiryReason, + LoopExpirySource, LoopPauseKind, LoopPauseRecord, MonitorOutcome, diff --git a/src/commands/loop-command.ts b/src/commands/loop-command.ts index 4cf5758..88fcc42 100644 --- a/src/commands/loop-command.ts +++ b/src/commands/loop-command.ts @@ -169,7 +169,12 @@ export function registerLoopCommand(options: LoopCommandOptions): void { if (entry) { const actions = ["x Delete"]; if (entry.status === "active") actions.unshift("- Pause"); - else if (entry.status === "paused" && !entry.orchestration && !isTerminalWorkflowRun(entry.workflow)) actions.unshift("* Resume"); + else if ( + entry.status === "paused" + && Date.now() < entry.expiresAt + && !entry.orchestration + && !isTerminalWorkflowRun(entry.workflow) + ) actions.unshift("* Resume"); actions.push("< Back"); const detail = entry.workflow diff --git a/src/index.ts b/src/index.ts index 5551dc2..a69617e 100644 --- a/src/index.ts +++ b/src/index.ts @@ -21,6 +21,7 @@ import { registerLoopCommand } from "./commands/loop-command.js"; import { atMaxFires } from "./loop-reducer.js"; import { MonitorManager } from "./monitor-manager.js"; import { rpcCall, rpcProbe } from "./rpc/cross-extension-rpc.js"; +import { buildLoopExpiredPayload } from "./runtime/loop-events.js"; import { createMonitorOnDoneRuntime } from "./runtime/monitor-ondone-runtime.js"; import { createNotificationRuntime, @@ -40,7 +41,7 @@ import { registerMonitorTools } from "./tools/monitor-tools.js"; import { registerSubagentOrchestrationTools } from "./tools/subagent-orchestration-tools.js"; import { registerWorkflowTools } from "./tools/workflow-tools.js"; import { TriggerSystem } from "./trigger-system.js"; -import type { LoopEntry, LoopFireOrigin, MonitorEntry, Trigger } from "./types.js"; +import type { LoopEntry, LoopExpiryDisposition, LoopExpiryReason, LoopExpirySource, LoopFireOrigin, MonitorEntry, Trigger } from "./types.js"; import { LoopWidget } from "./ui/widget.js"; import { atWorkflowStateFireLimit, getActiveWorkflowStateLoop, isTerminalWorkflowRun } from "./workflow-reducer.js"; @@ -73,7 +74,16 @@ export default function (pi: ExtensionAPI) { // call), so stale monitors don't linger in the count between turns. monitorManager.setOnChange(() => widget.update()); - scheduler = new CronScheduler(store, (entry, origin) => onLoopFire(entry, undefined, origin)); + function createScheduler(loopStore: LoopStore): CronScheduler { + return new CronScheduler( + loopStore, + (entry, origin) => onLoopFire(entry, undefined, origin), + (entry, disposition) => emitLoopExpired(entry, disposition, "scheduler", "expires_at"), + isCurrentExtensionContext, + ); + } + + scheduler = createScheduler(store); triggerSystem = new TriggerSystem(pi, scheduler, store, (entry, origin) => onLoopFire(entry, undefined, origin)); let taskProvider: TaskProviderRuntime | undefined; @@ -202,6 +212,25 @@ export default function (pi: ExtensionAPI) { } } + function emitLoopExpired( + entry: LoopEntry, + disposition: LoopExpiryDisposition, + source: LoopExpirySource, + reason: LoopExpiryReason, + generation = sessionGeneration, + ): void { + if (generation !== sessionGeneration || !isCurrentExtensionContext()) return; + triggerSystem.remove(entry.id); + const payload = buildLoopExpiredPayload(entry, disposition, source, reason, Date.now()); + try { + pi.events.emit("loops:expired", payload); + } catch (error) { + debug(`loops:expired #${entry.id} — event listener failed`, error); + } + void notificationRuntime.queueOrDeliverLoopExpired({ ...payload, sessionGeneration: generation }) + .catch((error) => debug(`loops:expired #${entry.id} — notification failed`, error)); + } + function emitLoopFire(entry: LoopEntry, monitor?: MonitorEntry, orchestrationWakeSequence?: number): void { pi.events.emit("loop:fire", { loopId: entry.id, @@ -332,7 +361,7 @@ export default function (pi: ExtensionAPI) { memoryLoopStores.set(sessionId, store); } widget.setStore(store); - scheduler = new CronScheduler(store, (entry, origin) => onLoopFire(entry, undefined, origin)); + scheduler = createScheduler(store); triggerSystem = new TriggerSystem(pi, scheduler, store, (entry, origin) => onLoopFire(entry, undefined, origin)); }, clearAllLoops: () => { @@ -365,6 +394,10 @@ export default function (pi: ExtensionAPI) { shutdownMonitors: () => monitorManager.shutdown(), hasPendingTasks, cleanDoneTasks, + isContextCurrent: isCurrentExtensionContext, + emitLoopExpired: (entry, disposition, reason, generation) => { + emitLoopExpired(entry, disposition, "session_recovery", reason, generation); + }, }); // ── Loop fire handler — queues an in-memory notification, then injects a custom message when delivery is safe ── diff --git a/src/runtime/loop-events.ts b/src/runtime/loop-events.ts index 4a0220b..a6cb87c 100644 --- a/src/runtime/loop-events.ts +++ b/src/runtime/loop-events.ts @@ -1,4 +1,4 @@ -import type { LoopEntry } from "../types.js"; +import type { LoopEntry, LoopExpiryDisposition, LoopExpiryReason, LoopExpirySource } from "../types.js"; export type LoopAutoDeleteReason = "task_backlog_empty"; @@ -20,6 +20,19 @@ export interface LoopAutodeletedPayload { pendingCount: number; } +export interface LoopExpiredPayload { + loopId: string; + prompt: string; + trigger: LoopEntry["trigger"]; + recurring: boolean; + createdAt: number; + expiresAt: number; + expiredAt: number; + disposition: LoopExpiryDisposition; + source: LoopExpirySource; + reason: LoopExpiryReason; +} + export interface TaskBacklogEmptyPayload { pendingCount: 0; deletedLoopIds: string[]; @@ -49,6 +62,27 @@ export function buildLoopAutodeletedPayload( }; } +export function buildLoopExpiredPayload( + entry: LoopEntry, + disposition: LoopExpiryDisposition, + source: LoopExpirySource, + reason: LoopExpiryReason, + expiredAt: number, +): LoopExpiredPayload { + return { + loopId: entry.id, + prompt: entry.prompt, + trigger: entry.trigger, + recurring: entry.recurring, + createdAt: entry.createdAt, + expiresAt: entry.expiresAt, + expiredAt, + disposition, + source, + reason, + }; +} + export function buildTaskBacklogEmptyPayload( deletedLoopIds: string[], ): TaskBacklogEmptyPayload { diff --git a/src/runtime/notification-runtime.ts b/src/runtime/notification-runtime.ts index 66b9716..f4f1e21 100644 --- a/src/runtime/notification-runtime.ts +++ b/src/runtime/notification-runtime.ts @@ -15,6 +15,7 @@ import { import { getOrchestrationCounts } from "../orchestration-reducer.js"; import type { DynamicLoopState, MonitorOutcome, OrchestrationState, Trigger, WorkflowRunState } from "../types.js"; import { getWorkflowOutcomeAvailability } from "../workflow-reducer.js"; +import type { LoopExpiredPayload } from "./loop-events.js"; import { TASK_BACKLOG_ACTION_CONTRACT } from "./task-backlog-runtime.js"; const MAX_ORCHESTRATION_WAKE_CHARS = 12_288; @@ -64,6 +65,7 @@ export interface NotificationRuntimeOptions { export interface NotificationRuntime { syncRuntimeState(options?: { agentRunning?: boolean; hasPendingMessages?: boolean }): void; queueOrDeliverNotification(data: LoopFireEvent): Promise; + queueOrDeliverLoopExpired(data: LoopExpiredPayload & { sessionGeneration?: number }): Promise; queueOrDeliverMonitorStarted(data: MonitorStartedEvent): Promise; discardMonitorStarted(monitorId: string): void; flushPendingNotifications(options?: { ignorePendingMessages?: boolean }): Promise; @@ -294,6 +296,31 @@ export function createNotificationRuntime(options: NotificationRuntimeOptions): }; } + function buildLoopExpiredNotification( + data: LoopExpiredPayload & { sessionGeneration?: number }, + ): PendingNotification { + const isStaleEvent = data.reason === "resume_event_stale"; + return { + sessionGeneration: data.sessionGeneration ?? sessionGeneration, + loopId: data.loopId, + prompt: data.prompt, + trigger: data.trigger, + timestamp: data.expiredAt, + recurring: data.recurring, + key: `loop:${data.loopId}:expired:${data.expiresAt}`, + message: [ + isStaleEvent + ? `[pi-loop] Loop #${data.loopId} retired during session recovery and was ${data.disposition}.` + : `[pi-loop] Loop #${data.loopId} expired and was ${data.disposition}.`, + data.prompt, + isStaleEvent + ? "Event and hybrid subscriptions do not resume across sessions." + : `Expiry boundary: ${new Date(data.expiresAt).toISOString()}`, + "Recreate it explicitly if this controller is still required; retirement does not imply consent to renew indefinitely.", + ].join("\n"), + }; + } + function buildMonitorStartedNotification(data: MonitorStartedEvent): PendingNotification { const label = data.description ?? data.command.slice(0, 80); return { @@ -402,6 +429,25 @@ export function createNotificationRuntime(options: NotificationRuntimeOptions): await flushPendingNotifications(); } + async function queueOrDeliverLoopExpired( + data: LoopExpiredPayload & { sessionGeneration?: number }, + ): Promise { + if (data.sessionGeneration !== undefined && data.sessionGeneration !== sessionGeneration) { + debug?.(`loops:expired #${data.loopId} — stale session generation, dropping wake`); + return; + } + const notification = buildLoopExpiredNotification(data); + applyNotificationEvent({ + type: "NOTIFICATION_QUEUED", + at: notification.timestamp, + source: "system", + entityType: "notification", + entityId: notification.key, + payload: { notification }, + }); + await flushPendingNotifications(); + } + async function queueOrDeliverMonitorStarted(data: MonitorStartedEvent): Promise { if (data.sessionGeneration !== undefined && data.sessionGeneration !== sessionGeneration) { debug?.(`monitor:started #${data.monitorId} — stale session generation, dropping wake`); @@ -455,6 +501,7 @@ export function createNotificationRuntime(options: NotificationRuntimeOptions): return { syncRuntimeState, queueOrDeliverNotification, + queueOrDeliverLoopExpired, queueOrDeliverMonitorStarted, discardMonitorStarted, flushPendingNotifications, diff --git a/src/runtime/session-runtime.ts b/src/runtime/session-runtime.ts index 2ad5311..5e4c7df 100644 --- a/src/runtime/session-runtime.ts +++ b/src/runtime/session-runtime.ts @@ -1,5 +1,6 @@ import type { ExtensionAPI, ExtensionContext } from "@earendil-works/pi-coding-agent"; import { LoopStore } from "../store.js"; +import type { LoopEntry, LoopExpiryDisposition, LoopExpiryReason } from "../types.js"; import type { NotificationRuntime } from "./notification-runtime.js"; import type { LoopScope } from "./scope.js"; @@ -38,6 +39,13 @@ export interface SessionRuntimeOptions { shutdownMonitors: () => Promise; hasPendingTasks: () => Promise; cleanDoneTasks: () => Promise; + isContextCurrent: () => boolean; + emitLoopExpired: ( + entry: LoopEntry, + disposition: LoopExpiryDisposition, + reason: LoopExpiryReason, + generation: number, + ) => void; } export function registerSessionRuntimeHooks(options: SessionRuntimeOptions): void { @@ -68,6 +76,8 @@ export function registerSessionRuntimeHooks(options: SessionRuntimeOptions): voi shutdownMonitors, hasPendingTasks, cleanDoneTasks, + isContextCurrent, + emitLoopExpired, } = options; let storeUpgraded = false; @@ -116,9 +126,19 @@ export function registerSessionRuntimeHooks(options: SessionRuntimeOptions): voi migrateTaskBacklogLoops(); if (!isCurrentGeneration(generation)) return; clearWorkflowMonitorWaits(); + if (!isContextCurrent()) return; const store = getStore(); - store.clearExpired(); - store.expireEventLoops(sessionStartedAt); + const expired = store.expireEntries(sessionStartedAt); + for (const record of expired) { + if (!isCurrentGeneration(generation)) return; + emitLoopExpired(record.entry, record.disposition, record.reason, generation); + } + if (!isCurrentGeneration(generation) || !isContextCurrent()) return; + const staleEventLoops = store.expireEventLoopEntries(sessionStartedAt); + for (const record of staleEventLoops) { + if (!isCurrentGeneration(generation)) return; + emitLoopExpired(record.entry, record.disposition, record.reason, generation); + } await recoverOrchestrations(); if (!isCurrentGeneration(generation)) return; const triggerSystem = getTriggerSystem(); diff --git a/src/scheduler.ts b/src/scheduler.ts index b0bfd13..51c3e93 100644 --- a/src/scheduler.ts +++ b/src/scheduler.ts @@ -1,6 +1,6 @@ import { computeJitter, cronToNextFire } from "./loop-parse.js"; import type { LoopStore } from "./store.js"; -import type { LoopEntry, LoopFireOrigin } from "./types.js"; +import type { LoopEntry, LoopExpiryDisposition, LoopFireOrigin } from "./types.js"; import { atWorkflowStateFireLimit, getActiveWorkflowStateLoop, isTerminalWorkflowRun } from "./workflow-reducer.js"; function computeNextFire(entry: LoopEntry): Date { @@ -17,43 +17,60 @@ function computeNextFire(entry: LoopEntry): Date { export class CronScheduler { private fireTimes = new Map(); + private expiryTimes = new Map(); constructor( private store: LoopStore, private onFire: (entry: LoopEntry, origin: LoopFireOrigin) => void, + private onExpired?: (entry: LoopEntry, disposition: LoopExpiryDisposition) => void, + private canExpire: () => boolean = () => true, ) {} start(): void { for (const storedEntry of this.store.list()) { let entry = storedEntry; if (entry.status !== "active" || entry.orchestration || entry.workflow?.waitingMonitor || isTerminalWorkflowRun(entry.workflow)) continue; - if (entry.trigger.type === "cron" || entry.trigger.type === "hybrid" || entry.trigger.type === "dynamic") { - if (entry.trigger.type === "dynamic" && entry.dynamic?.awaitingUpdate && !this.fireTimes.has(entry.id)) { - entry = this.store.updateDynamic(entry.id, { - dynamic: { - awaitingUpdate: false, - nextWakeAt: undefined, - lastUpdatedAt: Date.now(), - }, - }) ?? entry; - } - this.armTimer(entry); + if (entry.trigger.type === "event") { + this.expiryTimes.set(entry.id, entry.expiresAt); + continue; + } + if (entry.trigger.type === "dynamic" && entry.dynamic?.awaitingUpdate && !this.fireTimes.has(entry.id)) { + entry = this.store.updateDynamic(entry.id, { + dynamic: { + awaitingUpdate: false, + nextWakeAt: undefined, + lastUpdatedAt: Date.now(), + }, + }) ?? entry; } + this.armTimer(entry); } } stop(): void { this.fireTimes.clear(); + this.expiryTimes.clear(); } add(entry: LoopEntry): void { - if (!entry.orchestration && !entry.workflow?.waitingMonitor && !isTerminalWorkflowRun(entry.workflow) && (entry.trigger.type === "cron" || entry.trigger.type === "hybrid" || entry.trigger.type === "dynamic")) { - this.armTimer(entry); + if (entry.orchestration || entry.workflow?.waitingMonitor || isTerminalWorkflowRun(entry.workflow)) return; + if (entry.trigger.type === "event") { + this.expiryTimes.set(entry.id, entry.expiresAt); + return; } + this.armTimer(entry); + } + + expire(entry: LoopEntry, now = Date.now()): boolean { + const current = this.store.get(entry.id); + if (current?.status !== "active") return true; + if (now < current.expiresAt) return false; + return this.retireExpired(current, now); } remove(id: string): void { this.fireTimes.delete(id); + this.expiryTimes.delete(id); } nextFire(id: string): number | undefined { @@ -63,7 +80,24 @@ export class CronScheduler { private retire(entry: LoopEntry): void { if (entry.workflow || entry.taskBacklog) this.store.pause(entry.id, "controller_limit", "scheduler fire cap reached"); else this.store.delete(entry.id); - this.fireTimes.delete(entry.id); + this.remove(entry.id); + } + + private retireExpired(entry: LoopEntry, now = Date.now()): boolean { + if (!this.canExpire()) return false; + const record = this.store.expireEntry(entry.id, now); + if (!record) { + const current = this.store.get(entry.id); + if (current?.status === "active") { + this.expiryTimes.set(entry.id, current.expiresAt); + return false; + } + this.remove(entry.id); + return true; + } + this.remove(entry.id); + this.onExpired?.(record.entry, record.disposition); + return true; } private armTimer(entry: LoopEntry): void { @@ -80,15 +114,31 @@ export class CronScheduler { } const fireTime = nextFire.getTime() + jitter; - if (fireTime > entry.expiresAt) { - this.retire(entry); + if (fireTime >= entry.expiresAt) { + this.fireTimes.delete(entry.id); + this.expiryTimes.set(entry.id, entry.expiresAt); return; } + this.expiryTimes.delete(entry.id); this.fireTimes.set(entry.id, fireTime); } pump(now: number, filter?: (entry: LoopEntry) => boolean): void { + for (const [id, expiryTime] of this.expiryTimes) { + if (now < expiryTime) continue; + const entry = this.store.get(id); + if (entry?.status !== "active") { + this.remove(id); + continue; + } + if (now < entry.expiresAt) { + this.expiryTimes.set(id, entry.expiresAt); + continue; + } + this.expire(entry, now); + } + for (const [id, fireTime] of this.fireTimes) { if (now < fireTime) continue; @@ -103,7 +153,7 @@ export class CronScheduler { if (filter && !filter(entry)) continue; if (now >= entry.expiresAt) { - this.retire(entry); + this.retireExpired(entry, now); continue; } diff --git a/src/store.ts b/src/store.ts index 7b6ed99..8ae6a66 100644 --- a/src/store.ts +++ b/src/store.ts @@ -3,7 +3,7 @@ import { join } from "node:path"; import { type LoopReducerEffect, type LoopReducerEvent, type LoopReducerState, reduceLoopState } from "./loop-reducer.js"; import { applyOrchestrationEvent, type OrchestrationEvent, validateOrchestrationDefinition, validatePersistedOrchestration } from "./orchestration-reducer.js"; import { ReducerBackedStore } from "./reducer-backed-store.js"; -import type { DynamicLoopState, LoopDeletionTombstone, LoopDeletionTombstoneInput, LoopEntry, LoopFireOrigin, LoopPauseKind, LoopPauseRecord, LoopStoreData, OrchestrationActor, OrchestrationDefinitionInput, Trigger, WorkflowDefinition, WorkflowMonitorWait, WorkflowRevisionFailure, WorkflowRunState, WorkflowRuntimeActor, WorkflowTerminalStatus } from "./types.js"; +import type { DynamicLoopState, LoopDeletionTombstone, LoopDeletionTombstoneInput, LoopEntry, LoopExpiryDisposition, LoopExpiryReason, LoopFireOrigin, LoopPauseKind, LoopPauseRecord, LoopStoreData, OrchestrationActor, OrchestrationDefinitionInput, Trigger, WorkflowDefinition, WorkflowMonitorWait, WorkflowRevisionFailure, WorkflowRunState, WorkflowRuntimeActor, WorkflowTerminalStatus } from "./types.js"; import { validateWorkflowDefinition } from "./workflow-definition.js"; import { isTerminalWorkflowRun, transitionWorkflowRun, validateWorkflowAdmissionRecord, type WorkflowTransitionFailure, type WorkflowTransitionInput } from "./workflow-reducer.js"; import { validatePersistedWorkflowRevision, type WorkflowRevisionInput, type WorkflowRevisionSummary } from "./workflow-revision.js"; @@ -93,6 +93,12 @@ function normalizeLoopEntry(entry: LoopEntry): LoopEntry { }; } +export interface ExpiredLoopRecord { + entry: LoopEntry; + disposition: LoopExpiryDisposition; + reason: LoopExpiryReason; +} + const MAX_LOOPS = 25; const TOMBSTONE_TTL_MS = 10 * 60 * 1000; @@ -196,7 +202,7 @@ export class LoopStore extends ReducerBackedStore { const entry = this.entries.get(id); - if (!entry || isTerminalWorkflowRun(entry.workflow)) return undefined; + if (!entry || Date.now() >= entry.expiresAt || isTerminalWorkflowRun(entry.workflow)) return undefined; this.applyReducerEvent({ type: "LOOP_RESUMED", at: Date.now(), @@ -299,6 +305,7 @@ export class LoopStore extends ReducerBackedStore= entry.expiresAt) return undefined; if (entry.status === "paused") { this.applyReducerEvent({ type: "LOOP_RESUMED", @@ -598,38 +605,57 @@ export class LoopStore extends ReducerBackedStore this.expireEntryUnlocked(id, now)); + } + + expireEntries(now = Date.now()): ExpiredLoopRecord[] { return this.withLock(() => { - const now = Date.now(); - let count = 0; - for (const [id, entry] of [...this.entries.entries()]) { - if (now < entry.expiresAt) continue; - this.applyReducerEvent(entry.workflow || entry.orchestration - ? { - type: "LOOP_PAUSED", - at: now, - source: "system", - entityType: "loop", - entityId: id, - payload: { id, kind: "controller_limit", reason: "loop expiry reached" }, - } - : { - type: "LOOP_EXPIRED", - at: now, - source: "system", - entityType: "loop", - entityId: id, - payload: { id, reason: "expires_at" }, - }); - count++; + const expired: ExpiredLoopRecord[] = []; + for (const id of [...this.entries.keys()]) { + const record = this.expireEntryUnlocked(id, now); + if (record) expired.push(record); } - return count; + return expired; }); } - expireEventLoops(sessionStartedAt: number): number { + clearExpired(): number { + return this.expireEntries().length; + } + + expireEventLoopEntries(sessionStartedAt: number): ExpiredLoopRecord[] { return this.withLock(() => { - let count = 0; + const expired: ExpiredLoopRecord[] = []; for (const [id, entry] of [...this.entries.entries()]) { if (entry.status !== "active") continue; if (entry.trigger.type !== "event" && entry.trigger.type !== "hybrid") continue; @@ -644,12 +670,16 @@ export class LoopStore extends ReducerBackedStore { const entries = [...this.entries.values()]; diff --git a/src/tools/loop-tools.ts b/src/tools/loop-tools.ts index 30ee6e3..9b868a7 100644 --- a/src/tools/loop-tools.ts +++ b/src/tools/loop-tools.ts @@ -150,9 +150,12 @@ function continueDynamicLoop( store: LoopStoreLike, triggerSystem: TriggerSystemLike, ): { applied: boolean; message: string } { + if (Date.now() >= entry.expiresAt) { + return { applied: false, message: `Loop #${params.id} has expired; recreate it explicitly if work remains.` }; + } const { nextWakeAt, error } = resolveNextWakeAt(params.nextInterval); if (error) return { applied: false, message: error }; - if (nextWakeAt !== undefined && nextWakeAt > entry.expiresAt) { + if (nextWakeAt !== undefined && nextWakeAt >= entry.expiresAt) { return { applied: false, message: `nextInterval exceeds loop #${params.id}'s remaining lifetime.` }; } @@ -390,6 +393,7 @@ export function registerLoopTools(options: LoopToolsOptions): void { const statusIcon = entry.status === "active" ? "*" : entry.status === "paused" ? "-" : "x"; let line = `${statusIcon} #${entry.id} [${entry.status}] ${entry.prompt.slice(0, 60)}`; line += ` (${triggerDesc})`; + line += ` expiresAt: ${new Date(entry.expiresAt).toISOString()}`; if (nextFire) { const remaining = Math.max(0, nextFire - Date.now()); line += ` next: ${formatRemaining(remaining)}`; diff --git a/src/trigger-system.ts b/src/trigger-system.ts index ce51952..c4fa9db 100644 --- a/src/trigger-system.ts +++ b/src/trigger-system.ts @@ -39,9 +39,7 @@ export class TriggerSystem { add(entry: LoopEntry): void { if (isTerminalWorkflowRun(entry.workflow)) return; - if (entry.trigger.type === "cron" || entry.trigger.type === "hybrid" || entry.trigger.type === "dynamic") { - this.scheduler.add(entry); - } + this.scheduler.add(entry); if (entry.trigger.type === "event" || entry.trigger.type === "hybrid") { const ev = entry.trigger.type === "hybrid" ? entry.trigger.event : entry.trigger; this.subscribeEvent(entry, ev.source, ev.filter); @@ -117,7 +115,13 @@ export class TriggerSystem { return; } - this.lastFireTime.set(current.id, Date.now()); + const now = Date.now(); + if (now >= current.expiresAt) { + if (this.scheduler.expire(current, now)) this.remove(current.id); + return; + } + + this.lastFireTime.set(current.id, now); this.onFire(current, "event"); const fresh = this.store.get(entry.id); diff --git a/src/types.ts b/src/types.ts index 7ecea99..442484f 100644 --- a/src/types.ts +++ b/src/types.ts @@ -11,6 +11,9 @@ export interface LoopDeletionTombstone { export type LoopDeletionTombstoneInput = Omit; export type LoopStatus = "active" | "paused"; +export type LoopExpiryDisposition = "deleted" | "paused"; +export type LoopExpiryReason = "expires_at" | "resume_event_stale"; +export type LoopExpirySource = "scheduler" | "session_recovery"; export type LoopPauseKind = "administrative" | "controller_limit" | "semantic_terminal" | "orchestration_settlement"; diff --git a/test/index.test.ts b/test/index.test.ts index a16340a..17c2a9c 100644 --- a/test/index.test.ts +++ b/test/index.test.ts @@ -172,6 +172,47 @@ describe("native task fallback", () => { vi.useRealTimers(); }); + it("emits and delivers observable expiry for a recurring cron loop", async () => { + const { pi, toolMap, extensionHandlers, emittedEvents, sentMessages } = createMockPi(); + extension(pi as any); + await vi.advanceTimersByTimeAsync(6100); + vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); + const ctx = createCtx({ sessionId: "expiry-session" }); + + for (const handler of extensionHandlers.get("turn_start") ?? []) await handler(null, ctx); + await toolMap.get("LoopCreate")!.execute!("create-expiring", { + trigger: "0 8 * * *", + prompt: "Daily release check", + triggerType: "cron", + recurring: true, + }); + pi.events.on("loops:expired", () => { + throw new Error("consumer failed"); + }); + + vi.setSystemTime(new Date("2026-01-08T00:00:00Z")); + for (const handler of extensionHandlers.get("turn_start") ?? []) await handler(null, ctx); + await vi.advanceTimersByTimeAsync(0); + + expect(emittedEvents).toContainEqual({ + name: "loops:expired", + payload: expect.objectContaining({ + loopId: "1", + prompt: "Daily release check", + expiresAt: Date.parse("2026-01-08T00:00:00Z"), + disposition: "deleted", + source: "scheduler", + reason: "expires_at", + }), + }); + + for (const handler of extensionHandlers.get("agent_end") ?? []) await handler(null, ctx); + await vi.advanceTimersByTimeAsync(0); + expect(sentMessages.some((item) => item.message.content.includes("Loop #1 expired and was deleted"))).toBe(true); + const list = await toolMap.get("LoopList")!.execute!("list-expired", {}); + expect(list.content[0].text).toContain("No loops configured"); + }); + it("registers native task tools when pi-tasks is unavailable", async () => { const { pi, toolMap, commandMap } = createMockPi(); diff --git a/test/loop-tools.test.ts b/test/loop-tools.test.ts index 82ba137..68e127e 100644 --- a/test/loop-tools.test.ts +++ b/test/loop-tools.test.ts @@ -231,6 +231,19 @@ describe("LoopList", () => { expect(out).toContain("cron:"); }); + it("shows the stable expiry boundary for recurring loops", async () => { + const h = setup(); + vi.useFakeTimers(); + vi.setSystemTime(new Date("2026-01-01T00:00:00Z")); + try { + await h.text("LoopCreate", { trigger: "5m", prompt: "build check", triggerType: "cron" }); + const out = await h.text("LoopList", {}); + expect(out).toContain("expiresAt: 2026-01-08T00:00:00.000Z"); + } finally { + vi.useRealTimers(); + } + }); + it("shows wall-clock age for active loops", async () => { const h = setup(); vi.useFakeTimers(); @@ -392,6 +405,17 @@ describe("LoopUpdate", () => { expect(h.store.get("1")).toEqual(before); }); + it("rejects continuing an expired dynamic controller", async () => { + h.store.get("1")!.expiresAt = Date.now(); + h.triggerSystem.add.mockClear(); + + const out = await h.text("LoopUpdate", { id: "1", status: "continue" }); + + expect(out).toContain("has expired; recreate it explicitly"); + expect(h.store.get("1")?.dynamic?.iteration).toBe(0); + expect(h.triggerSystem.add).not.toHaveBeenCalled(); + }); + it("resumes a paused dynamic loop when it continues", async () => { await h.text("LoopUpdate", { id: "1", status: "paused" }); h.triggerSystem.add.mockClear(); diff --git a/test/notification-runtime.test.ts b/test/notification-runtime.test.ts index f2fcd12..d7dc66d 100644 --- a/test/notification-runtime.test.ts +++ b/test/notification-runtime.test.ts @@ -4,6 +4,89 @@ import { createNotificationRuntime } from "../src/runtime/notification-runtime.j import { createMockPi } from "./helpers/mock-pi.js"; describe("notification runtime session boundary", () => { + it("delivers an explicit recurring-loop expiry notification", async () => { + const { pi, sentMessages } = createMockPi(); + const runtime = createNotificationRuntime({ + pi, + hasPendingTasks: vi.fn(async () => 0), + cleanDoneTasks: vi.fn(async () => {}), + getHasPendingMessages: () => false, + }); + runtime.syncRuntimeState({ agentRunning: false, hasPendingMessages: false }); + + await runtime.queueOrDeliverLoopExpired({ + loopId: "7", + prompt: "Daily release check", + trigger: { type: "cron", schedule: "0 8 * * *" }, + recurring: true, + createdAt: 100, + expiresAt: 200, + expiredAt: 201, + disposition: "deleted", + source: "scheduler", + reason: "expires_at", + }); + + expect(sentMessages).toHaveLength(1); + expect(sentMessages[0].message.content).toContain("Loop #7 expired and was deleted"); + expect(sentMessages[0].message.content).toContain("Daily release check"); + expect(sentMessages[0].message.content).toContain("Recreate it explicitly if this controller is still required"); + }); + + it("explains stale event-loop retirement during session recovery", async () => { + const { pi, sentMessages } = createMockPi(); + const runtime = createNotificationRuntime({ + pi, + hasPendingTasks: vi.fn(async () => 0), + cleanDoneTasks: vi.fn(async () => {}), + getHasPendingMessages: () => false, + }); + + await runtime.queueOrDeliverLoopExpired({ + loopId: "9", + prompt: "Wait for deploy", + trigger: { type: "event", source: "deploy:finished" }, + recurring: true, + createdAt: 100, + expiresAt: 10_000, + expiredAt: 201, + disposition: "deleted", + source: "session_recovery", + reason: "resume_event_stale", + }); + + expect(sentMessages[0].message.content).toContain("retired during session recovery and was deleted"); + expect(sentMessages[0].message.content).toContain("Event and hybrid subscriptions do not resume across sessions"); + expect(sentMessages[0].message.content).not.toContain("Expiry boundary"); + }); + + it("drops an expiry notification from a stale session generation", async () => { + const { pi, sentMessages } = createMockPi(); + const runtime = createNotificationRuntime({ + pi, + hasPendingTasks: vi.fn(async () => 0), + cleanDoneTasks: vi.fn(async () => {}), + getHasPendingMessages: () => false, + }); + runtime.clear("session_switch"); + + await runtime.queueOrDeliverLoopExpired({ + loopId: "7", + prompt: "Old schedule", + trigger: { type: "cron", schedule: "0 8 * * *" }, + recurring: true, + createdAt: 100, + expiresAt: 200, + expiredAt: 201, + disposition: "deleted", + source: "scheduler", + reason: "expires_at", + sessionGeneration: 0, + }); + + expect(sentMessages).toEqual([]); + }); + it("drops a wake whose task lookup completes after session shutdown", async () => { const { pi, sentMessages } = createMockPi(); let resolvePending: ((pending: number) => void) | undefined; diff --git a/test/scheduler.test.ts b/test/scheduler.test.ts index dcec97c..aa67d8c 100644 --- a/test/scheduler.test.ts +++ b/test/scheduler.test.ts @@ -9,13 +9,17 @@ describe("CronScheduler", () => { let store: LoopStore; let scheduler: CronScheduler; let fired: string[]; + let expired: Array<{ id: string; disposition: string }>; beforeEach(() => { vi.useFakeTimers({ shouldAdvanceTime: true }); store = new LoopStore(); fired = []; + expired = []; scheduler = new CronScheduler(store, (entry) => { fired.push(entry.id); + }, (entry, disposition) => { + expired.push({ id: entry.id, disposition }); }); }); @@ -451,14 +455,49 @@ describe("CronScheduler", () => { expect(fired).toHaveLength(0); }); - it("deletes expired entries on pump", () => { - const entry = store.create(cronTrigger, "expired", { recurring: false }); - entry.expiresAt = Date.now() - 1; + it("tracks event-only loops to the same observable expiry boundary", () => { + const entry = store.create({ type: "event", source: "deploy:finished" }, "event expiry", { + recurring: true, + }); + entry.expiresAt = Date.now() + 60_000; scheduler.add(entry); - vi.advanceTimersByTime(10 * 60 * 1000); + vi.advanceTimersByTime(60_000); + scheduler.pump(Date.now()); + + expect(store.get(entry.id)).toBeUndefined(); + expect(expired).toEqual([{ id: entry.id, disposition: "deleted" }]); + }); + + it("does not consume expiry under a stale extension context", () => { + const entry = store.create(cronTrigger, "expired", { recurring: true }); + entry.expiresAt = Date.now() + 60_000; + scheduler = new CronScheduler(store, vi.fn(), vi.fn(), () => false); + scheduler.add(entry); + + vi.advanceTimersByTime(60_000); scheduler.pump(Date.now()); + + expect(store.get(entry.id)).toBeDefined(); + }); + + it("keeps recurring entries observable until expiry and reports retirement once", () => { + vi.setSystemTime(new Date("2026-01-01T00:00:01Z")); + const entry = store.create(cronTrigger, "expired", { recurring: true }); + entry.expiresAt = Date.now() + 60_000; + scheduler.add(entry); + + expect(store.get(entry.id)).toBeDefined(); + expect(scheduler.nextFire(entry.id)).toBeUndefined(); + expect(expired).toEqual([]); + + vi.advanceTimersByTime(60_000); + scheduler.pump(Date.now()); + expect(fired).toHaveLength(0); expect(store.get(entry.id)).toBeUndefined(); + expect(expired).toEqual([{ id: entry.id, disposition: "deleted" }]); + scheduler.pump(Date.now() + 60_000); + expect(expired).toHaveLength(1); }); }); diff --git a/test/session-runtime.test.ts b/test/session-runtime.test.ts index f414ec4..685e1b2 100644 --- a/test/session-runtime.test.ts +++ b/test/session-runtime.test.ts @@ -14,7 +14,11 @@ function setup(overrides: Partial = {}) { advanceSessionGeneration: () => ++sessionGeneration, recreateSessionStore: vi.fn(), clearAllLoops: vi.fn(), - getStore: () => ({ list: () => [], clearExpired: vi.fn(), expireEventLoops: vi.fn() }) as any, + getStore: () => ({ + list: () => [], + expireEntries: vi.fn(() => []), + expireEventLoopEntries: vi.fn(() => []), + }) as any, getScheduler: () => scheduler as any, getTriggerSystem: () => ({ start: vi.fn(), stop: vi.fn() }), setLatestCtx: vi.fn(), @@ -23,6 +27,7 @@ function setup(overrides: Partial = {}) { notificationRuntime: { syncRuntimeState: vi.fn(), queueOrDeliverNotification: vi.fn(async () => {}), + queueOrDeliverLoopExpired: vi.fn(async () => {}), queueOrDeliverMonitorStarted: vi.fn(async () => {}), discardMonitorStarted: vi.fn(), flushPendingNotifications: vi.fn(async () => {}), @@ -40,6 +45,8 @@ function setup(overrides: Partial = {}) { shutdownMonitors: vi.fn(async () => {}), hasPendingTasks: vi.fn(async () => 0), cleanDoneTasks: vi.fn(async () => {}), + isContextCurrent: () => true, + emitLoopExpired: vi.fn(), ...overrides, }; registerSessionRuntimeHooks(options); @@ -83,8 +90,8 @@ describe("session-runtime heartbeat lifecycle", () => { recoverOrchestrations: vi.fn(async () => { calls.push("recover orchestrations"); }), getStore: () => ({ list: () => [{ id: "8", status: "active" }], - clearExpired: vi.fn(() => { calls.push("clear expired"); }), - expireEventLoops: vi.fn(() => { calls.push("expire events"); }), + expireEntries: vi.fn(() => { calls.push("clear expired"); return []; }), + expireEventLoopEntries: vi.fn(() => { calls.push("expire events"); return []; }), }) as any, getTriggerSystem: () => triggerSystem, }); @@ -95,6 +102,78 @@ describe("session-runtime heartbeat lifecycle", () => { expect(calls).toEqual(["migrate", "clear monitor waits", "clear expired", "expire events", "recover orchestrations", "start"]); }); + it("does not consume persisted expiry under a stale extension context", async () => { + const expireEntries = vi.fn(() => []); + const { drive } = setup({ + getStore: () => ({ + list: () => [], + expireEntries, + expireEventLoopEntries: vi.fn(() => []), + }) as any, + isContextCurrent: () => false, + }); + + await drive("session_start"); + + expect(expireEntries).not.toHaveBeenCalled(); + }); + + it("emits recovered expiries after store cleanup", async () => { + const expired = { + id: "8", + prompt: "Daily release check", + trigger: { type: "cron", schedule: "0 8 * * *" }, + status: "active", + recurring: true, + createdAt: 1, + updatedAt: 1, + expiresAt: 2, + }; + const emitLoopExpired = vi.fn(); + const { drive } = setup({ + getStore: () => ({ + list: () => [], + expireEntries: vi.fn(() => [{ entry: expired, disposition: "deleted", reason: "expires_at" }]), + expireEventLoopEntries: vi.fn(() => []), + }) as any, + emitLoopExpired, + }); + + await drive("session_start"); + + expect(emitLoopExpired).toHaveBeenCalledWith(expired, "deleted", "expires_at", 0); + }); + + it("emits stale event-loop retirement found during session recovery", async () => { + const stale = { + id: "9", + prompt: "Wait for deploy", + trigger: { type: "event", source: "deploy:finished" }, + status: "active", + recurring: true, + createdAt: 1, + updatedAt: 1, + expiresAt: Date.now() + 60_000, + }; + const emitLoopExpired = vi.fn(); + const { drive } = setup({ + getStore: () => ({ + list: () => [], + expireEntries: vi.fn(() => []), + expireEventLoopEntries: vi.fn(() => [{ + entry: stale, + disposition: "deleted", + reason: "resume_event_stale", + }]), + }) as any, + emitLoopExpired, + }); + + await drive("session_start"); + + expect(emitLoopExpired).toHaveBeenCalledWith(stale, "deleted", "resume_event_stale", 0); + }); + it("repaints the widget on session_start after the harness resets extension UI", async () => { const widget = { setUICtx: vi.fn(), update: vi.fn() }; const setSessionId = vi.fn(); @@ -164,8 +243,8 @@ describe("session-runtime heartbeat lifecycle", () => { setSessionId, getStore: () => ({ list: () => [{ id: "8", status: "active" }], - clearExpired: vi.fn(), - expireEventLoops: vi.fn(), + expireEntries: vi.fn(() => []), + expireEventLoopEntries: vi.fn(() => []), }) as any, getTriggerSystem: () => triggerSystem, }); @@ -241,8 +320,8 @@ describe("session-runtime heartbeat lifecycle", () => { autoTask: true, trigger: { type: "cron", schedule: "*/5 * * * *" }, }], - clearExpired: vi.fn(), - expireEventLoops: vi.fn(), + expireEntries: vi.fn(() => []), + expireEventLoopEntries: vi.fn(() => []), }) as any, getScheduler: () => scheduler as any, hasPendingTasks: vi.fn(() => { diff --git a/test/store.test.ts b/test/store.test.ts index 609f5eb..3cc0979 100644 --- a/test/store.test.ts +++ b/test/store.test.ts @@ -102,6 +102,15 @@ describe("LoopStore (in-memory)", () => { expect(entry?.pause).toBeUndefined(); }); + it("rejects resuming a controller after its expiry boundary", () => { + const created = store.create(cronTrigger, "expired", { recurring: true }); + const paused = store.pause(created.id)!; + paused.expiresAt = Date.now(); + + expect(store.resume(created.id)).toBeUndefined(); + expect(store.get(created.id)?.status).toBe("paused"); + }); + it("rejects a paused terminal transition without trusted admission", () => { store.create({ type: "dynamic" }, "Investigate", { recurring: true, @@ -247,6 +256,35 @@ describe("LoopStore (in-memory)", () => { expect(store2.list()).toHaveLength(0); }); + it("returns bounded retirement records for recovered expiries", () => { + const ordinary = store.create(cronTrigger, "ordinary", { recurring: true }); + const workflow = store.create({ type: "dynamic" }, "workflow", { + recurring: true, + workflow: { + version: 1, + initialState: "work", + states: { work: { prompt: "Work.", on: { done: "done" } }, done: { prompt: "Done.", terminal: "completed" } }, + }, + }); + const backlog = store.create({ type: "event", source: "tasks:created" }, "backlog", { + recurring: true, + taskBacklog: true, + }); + ordinary.expiresAt = 10; + workflow.expiresAt = 10; + backlog.expiresAt = 10; + + expect(store.expireEntries(10)).toEqual([ + { entry: ordinary, disposition: "deleted", reason: "expires_at" }, + { entry: workflow, disposition: "paused", reason: "expires_at" }, + { entry: backlog, disposition: "paused", reason: "expires_at" }, + ]); + expect(store.get(ordinary.id)).toBeUndefined(); + expect(store.get(workflow.id)?.status).toBe("paused"); + expect(store.get(backlog.id)?.status).toBe("paused"); + expect(store.expireEntries(11)).toEqual([]); + }); + it("pauses expired workflows instead of deleting them", () => { const workflow = store.create({ type: "dynamic" }, "workflow", { recurring: true, @@ -290,14 +328,17 @@ describe("LoopStore (in-memory)", () => { const eventTrigger = { type: "event" as const, source: "monitor:done" }; const cronT = { type: "cron" as const, schedule: "*/5 * * * *" }; - s.create(eventTrigger, "event loop", { recurring: false }); + const first = s.create(eventTrigger, "event loop", { recurring: false }); s.create(cronT, "cron loop", { recurring: true }); - s.create(eventTrigger, "another event", { recurring: true }); + const second = s.create(eventTrigger, "another event", { recurring: true }); s.create({ type: "event", source: "tasks:created" }, "backlog worker", { recurring: true, taskBacklog: true }); // sessionStartedAt is set after creation — simulating loop persisted from prior session const sessionStartedAt = Date.now() + 1; - expect(s.expireEventLoops(sessionStartedAt)).toBe(2); + expect(s.expireEventLoopEntries(sessionStartedAt)).toEqual([ + { entry: first, disposition: "deleted", reason: "resume_event_stale" }, + { entry: second, disposition: "deleted", reason: "resume_event_stale" }, + ]); expect(s.get("2")!.status).toBe("active"); // cron loop untouched expect(s.get("1")).toBeUndefined(); // ordinary event loops deleted @@ -383,6 +424,15 @@ describe("LoopStore (in-memory)", () => { expect(store.get(entry.id)?.dynamic?.state).toBe("newer"); }); + it("rejects dynamic continuation at the authoritative expiry boundary", () => { + const entry = store.create({ type: "dynamic" }, "ship", { recurring: true }); + entry.expiresAt = Date.now(); + + expect(store.continueDynamic(entry.id, { dynamic: { state: "too late", iteration: 1 } })).toBeUndefined(); + expect(store.get(entry.id)?.dynamic?.state).toBeUndefined(); + expect(store.get(entry.id)?.status).toBe("active"); + }); + it("defaults dynamic goal to the prompt", () => { const l = store.create({ type: "dynamic" }, "ship the fix", { recurring: true }); expect(l.dynamic?.goal).toBe("ship the fix"); @@ -455,6 +505,29 @@ describe("LoopStore (file-backed)", () => { expect(loops[0].prompt).toBe("persist test"); }); + it("settles expiry once across stores sharing a project snapshot", () => { + const store1 = new LoopStore(filePath); + const workflow = store1.create({ type: "dynamic" }, "workflow", { + recurring: true, + workflow: { + version: 1, + initialState: "work", + states: { work: { prompt: "Work.", on: { done: "done" } }, done: { prompt: "Done.", terminal: "completed" } }, + }, + }); + const store2 = new LoopStore(filePath); + + expect(store1.expireEntry(workflow.id, workflow.expiresAt)).toMatchObject({ + entry: { id: workflow.id }, + disposition: "paused", + }); + expect(store2.expireEntry(workflow.id, workflow.expiresAt)).toBeUndefined(); + expect(store2.get(workflow.id)).toMatchObject({ + status: "paused", + pause: { kind: "controller_limit", reason: "loop expiry reached" }, + }); + }); + it("persists dynamic loop state to disk", () => { const store1 = new LoopStore(filePath); store1.create({ type: "dynamic" }, "finish dynamic loop", { diff --git a/test/trigger-system.test.ts b/test/trigger-system.test.ts index fcae677..f804ce9 100644 --- a/test/trigger-system.test.ts +++ b/test/trigger-system.test.ts @@ -147,6 +147,42 @@ describe("TriggerSystem", () => { }); }); + it("retires and unsubscribes an event loop at its expiry boundary", () => { + const eventTrigger: Trigger = { type: "event", source: "expired_event" }; + const entry = store.create(eventTrigger, "expired event", { recurring: true }); + entry.expiresAt = Date.now(); + system.add(entry); + const remove = vi.spyOn(system, "remove"); + + pi.events.emit("expired_event", {}); + + expect(store.get(entry.id)).toBeUndefined(); + expect(remove).toHaveBeenCalledWith(entry.id); + const fireCalls = (pi.events.emit as any).mock.calls.filter( + (call: string[]) => call[0] === "loop:fire", + ); + expect(fireCalls).toEqual([]); + }); + + it("keeps an expired event subscription when stale context denies settlement", () => { + const staleScheduler = new CronScheduler(store, vi.fn(), undefined, () => false); + const staleSystem = new TriggerSystem(pi, staleScheduler, store, vi.fn()); + const entry = store.create({ type: "event", source: "stale_expired_event" }, "stale expiry", { + recurring: true, + }); + entry.expiresAt = Date.now(); + staleSystem.add(entry); + const remove = vi.spyOn(staleSystem, "remove"); + + try { + pi.events.emit("stale_expired_event", {}); + expect(store.get(entry.id)?.status).toBe("active"); + expect(remove).not.toHaveBeenCalled(); + } finally { + staleSystem.stop(); + } + }); + it("deletes one-shot event loops immediately after the first fire", () => { const eventTrigger: Trigger = { type: "event", source: "fire_once" }; const entry = store.create(eventTrigger, "one-shot", { recurring: false });