Skip to content
Merged
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
14 changes: 14 additions & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ function debug(...args: unknown[]) {

export default function (pi: ExtensionAPI) {
const runtimeId = randomUUID();
const runtimeProbeEvent = `pi-loop:runtime-probe:${runtimeId}`;
const piLoopEnv = process.env.PI_LOOP;
const piLoopScope = process.env.PI_LOOP_SCOPE as "memory" | "session" | "project" | undefined;
let loopScope: "memory" | "session" | "project" = piLoopScope ?? "session";
Expand Down Expand Up @@ -87,6 +88,7 @@ export default function (pi: ExtensionAPI) {
store.delete(id);
},
onLoopFire,
isContextCurrent: isCurrentExtensionContext,
completeWorkflowMonitorWait: (id, expected) => store.completeWorkflowMonitorWait(id, expected),
rearmWorkflow: (entry) => {
triggerSystem.add(entry);
Expand Down Expand Up @@ -178,6 +180,17 @@ export default function (pi: ExtensionAPI) {

// ── Loop fire handler ──

function isCurrentExtensionContext(): boolean {
try {
pi.events.emit(runtimeProbeEvent, { runtimeId, sessionGeneration });
return true;
} catch (error) {
if (!isStaleExtensionContextError(error)) throw error;
debug("extension context went stale, dropping runtime callback");
return false;
}
}

function emitLoopFire(entry: LoopEntry, monitor?: MonitorEntry): void {
pi.events.emit("loop:fire", {
loopId: entry.id,
Expand Down Expand Up @@ -209,6 +222,7 @@ export default function (pi: ExtensionAPI) {
monitor?: MonitorEntry,
origin: LoopFireOrigin = monitor ? "monitor" : "dynamic",
): void {
if (!isCurrentExtensionContext()) return;
debug(`loop:fire #${entry.id}`, { prompt: entry.prompt.slice(0, 50) });
const current = store.get(entry.id);
if (current?.status !== "active" || isTerminalWorkflowRun(current?.workflow)) {
Expand Down
5 changes: 5 additions & 0 deletions src/runtime/monitor-ondone-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ export interface MonitorOnDoneRuntimeOptions {
getLoop: (id: string) => LoopEntry | undefined;
deleteLoop: (id: string) => void;
onLoopFire: (entry: LoopEntry) => void;
isContextCurrent: () => boolean;
completeWorkflowMonitorWait: (id: string, expected: WorkflowMonitorWait) => LoopEntry | undefined;
rearmWorkflow: (entry: LoopEntry) => void;
wakeWorkflow: (entry: LoopEntry, monitor: MonitorEntry | undefined) => void;
Expand Down Expand Up @@ -56,6 +57,7 @@ export function createMonitorOnDoneRuntime(options: MonitorOnDoneRuntimeOptions)
getLoop,
deleteLoop,
onLoopFire,
isContextCurrent,
completeWorkflowMonitorWait,
rearmWorkflow,
wakeWorkflow,
Expand All @@ -71,6 +73,7 @@ export function createMonitorOnDoneRuntime(options: MonitorOnDoneRuntimeOptions)
reducers: [monitorCompletionReducerHandler],
effectHandlers: {
DELIVER_MONITOR_ONDONE_WAKE: (effect: ReducerEffect) => {
if (!isContextCurrent()) return;
const { loopId, monitorId, monitor } = effect.payload as {
loopId: string;
monitorId: string;
Expand All @@ -91,6 +94,7 @@ export function createMonitorOnDoneRuntime(options: MonitorOnDoneRuntimeOptions)
function register(doneLoop: LoopEntry, monitorId: string): void {
const timeoutAlert = isTimeoutAlertLoop(doneLoop);
const deliver = (monitor?: MonitorEntry) => {
if (!isContextCurrent()) return;
const outcome = monitor ?? monitorManager.get(monitorId);
if (timeoutAlert && !timedOut(outcome)) {
debug?.(`timeout alert loop #${doneLoop.id} — monitor #${monitorId} ended without timing out, expiring`);
Expand Down Expand Up @@ -128,6 +132,7 @@ export function createMonitorOnDoneRuntime(options: MonitorOnDoneRuntimeOptions)
if (!wait) return;

const deliver = (monitor?: MonitorEntry) => {
if (!isContextCurrent()) return;
const resumed = completeWorkflowMonitorWait(entry.id, wait);
if (!resumed) return;
if (resumed.status !== "active") return;
Expand Down
61 changes: 61 additions & 0 deletions test/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1244,6 +1244,39 @@ describe("native task fallback", () => {
const result = await loopList!.execute?.("2", {});
expect(result.content[0].text).toBe("No loops configured. Use LoopCreate to set up a schedule.");
});

it("does not consume a final fire when a stale runtime cannot dispatch its wake", async () => {
const { pi, toolMap, sentMessages } = createMockPi();

extension(pi as any);
await vi.advanceTimersByTimeAsync(6100);
await Promise.resolve();

const loopCreate = toolMap.get("LoopCreate");
const loopList = toolMap.get("LoopList");
await loopCreate!.execute?.("1", {
trigger: "stale-runtime:test:event",
prompt: "Preserve this fire for the current runtime",
triggerType: "event",
recurring: true,
maxFires: 1,
});

const dispatch = pi.events.emit.getMockImplementation();
pi.events.emit.mockImplementation((name: string, payload: unknown) => {
if (name !== "stale-runtime:test:event") {
throw new Error("This extension ctx is stale after session replacement or reload.");
}
return dispatch?.(name, payload);
});

expect(() => pi.events.emit("stale-runtime:test:event", {})).not.toThrow();
await Promise.resolve();

const result = await loopList!.execute?.("2", {});
expect(result.content[0].text).toContain("* #1 [active] Preserve this fire for the current runtime");
expect(sentMessages).toHaveLength(0);
});
});

describe("dynamic loop pump", () => {
Expand Down Expand Up @@ -1964,6 +1997,34 @@ describe("monitor tool wrappers", () => {
);
}, 10000);

it("drops an onDone wake when monitor completion reaches a stale extension context", async () => {
const { pi, toolMap, sentMessages: sentCustomMessages } = createMockPi();

extension(pi as any);
await vi.advanceTimersByTimeAsync(6100);
vi.useRealTimers();

const monitorCreate = toolMap.get("MonitorCreate");
expect(monitorCreate?.execute).toBeDefined();

const result = await monitorCreate!.execute?.("1", {
command: "node -e \"setTimeout(() => {}, 100)\"",
onDone: "This stale wake must be dropped",
});
expect(result.content[0].text).toContain("Completion wake loop");

pi.events.emit.mockImplementation(() => {
throw new Error("This extension ctx is stale after session replacement or reload.");
});

await new Promise(r => setTimeout(r, 500));

expect(pi.events.emit.mock.calls.some(([name]: [string]) => name === "loop:fire")).toBe(false);
expect(sentCustomMessages).toHaveLength(0);
const loops = await toolMap.get("LoopList")!.execute!("list", {});
expect(loops.content[0].text).toContain("* #1 [active] This stale wake must be dropped");
}, 10000);

it("onDone monitor completion does not rely on monitor:done event dispatch", async () => {
const { pi, toolMap, sentMessages: sentCustomMessages } = createMockPi({ suppressMonitorDoneDispatch: true });

Expand Down
19 changes: 18 additions & 1 deletion test/monitor-ondone-runtime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,14 +30,15 @@ function mockManager(config: { onCompleteReturns: boolean; onTerminalReturns?: b
};
}

function setup(manager: ReturnType<typeof mockManager>) {
function setup(manager: ReturnType<typeof mockManager>, isContextCurrent = () => true) {
const onLoopFire = vi.fn();
const deleteLoop = vi.fn();
const runtime = createMonitorOnDoneRuntime({
monitorManager: manager as any,
getLoop: (id: string) => (id === doneLoop.id ? doneLoop : undefined),
deleteLoop,
onLoopFire,
isContextCurrent,
completeWorkflowMonitorWait: vi.fn(),
rearmWorkflow: vi.fn(),
wakeWorkflow: vi.fn(),
Expand All @@ -62,6 +63,7 @@ describe("monitor-ondone-runtime", () => {
getLoop: (id) => (id === timeoutLoop.id ? timeoutLoop : undefined),
deleteLoop,
onLoopFire,
isContextCurrent: () => true,
completeWorkflowMonitorWait: vi.fn(),
rearmWorkflow: vi.fn(),
wakeWorkflow: vi.fn(),
Expand Down Expand Up @@ -106,6 +108,20 @@ describe("monitor-ondone-runtime", () => {
expect(deleteLoop).toHaveBeenCalledWith("5");
});

it("does not mutate loop state when completion belongs to a stale context", async () => {
const manager = mockManager({ onCompleteReturns: true });
const isContextCurrent = vi.fn(() => false);
const { runtime, onLoopFire, deleteLoop } = setup(manager, isContextCurrent);

runtime.register(doneLoop, "3");
manager.fireCaptured();
await flush();

expect(isContextCurrent).toHaveBeenCalledTimes(1);
expect(onLoopFire).not.toHaveBeenCalled();
expect(deleteLoop).not.toHaveBeenCalled();
});

it("delivers immediately when the monitor is already completed", async () => {
const manager = mockManager({ onCompleteReturns: false, status: "completed" });
const { runtime, onLoopFire, deleteLoop } = setup(manager);
Expand Down Expand Up @@ -198,6 +214,7 @@ describe("monitor-ondone-runtime", () => {
getLoop: () => undefined,
deleteLoop: vi.fn(),
onLoopFire: vi.fn(),
isContextCurrent: () => true,
completeWorkflowMonitorWait,
rearmWorkflow,
wakeWorkflow,
Expand Down
Loading