diff --git a/apps/server/src/engine/service.ts b/apps/server/src/engine/service.ts index 3f2af061..b4ef7b04 100644 --- a/apps/server/src/engine/service.ts +++ b/apps/server/src/engine/service.ts @@ -1038,7 +1038,7 @@ export class AgentService { owner, "monitors", monitor.id, - { status: "active" }, + { status: "active", checks: monitor.checks }, { checks: monitor.checks + 1, lastCheckedAt: date(), diff --git a/tests/monitor-recovery.test.ts b/tests/monitor-recovery.test.ts index 8d40eb23..ec039003 100644 --- a/tests/monitor-recovery.test.ts +++ b/tests/monitor-recovery.test.ts @@ -6,6 +6,7 @@ import { join } from "node:path"; import { after, before, test } from "node:test"; import { createApp } from "../apps/server/src/app.ts"; import { createStore, type Store } from "../apps/server/src/db.ts"; +import { LostLeaseError } from "../apps/server/src/engine/worker.ts"; import type { AgentNotification, AgentTask, Monitor } from "../packages/domain/src/agent.ts"; let db: Store, server: Awaited>, directory: string, token: string; @@ -611,3 +612,40 @@ test("a check keeps the baseline from a run that finishes during the request", a assert.equal(task?.status, "queued"); assert.equal(task?.state.lastHash, afterRun); }); + +test("a second observe from a stale snapshot loses its monitor commit instead of clobbering it", async () => { + await read("/sample-page", { text: "No tables available" }); + const monitor = await createMonitor("Stale observe"); + await server.agent.worker.tick(); + const base = await db.get(owner, "monitors", monitor.id); + assert.equal(base?.checks, 1); + + // Two workers read the same monitor snapshot (lease overlap), then both run + // observe() to completion. Freeze the monitor read at the pre-run snapshot. + const originalGet = db.get.bind(db); + db.get = (async (o: string, kind: string, id: string) => { + if (kind === "monitors" && id === monitor.id) return base; + return originalGet(o, kind, id); + }) as Store["get"]; + const task = await originalGet(owner, "tasks", monitor.taskId); + assert.ok(task); + const context = { + signal: new AbortController().signal, + guard: async () => {}, + checkpoint: async (_patch: Partial) => task, + event: async () => {}, + }; + const observe = ( + server.agent as unknown as { + observe(owner: string, task: AgentTask, ctx: typeof context): Promise>; + } + ).observe.bind(server.agent); + try { + await observe(owner, task, context); + await assert.rejects(observe(owner, task, context), LostLeaseError); + } finally { + db.get = originalGet; + } + // The loser's commit was rejected: exactly one increment landed. + assert.equal((await db.get(owner, "monitors", monitor.id))?.checks, 2); +});