diff --git a/server/src/__tests__/low-trust-red-team-routes.test.ts b/server/src/__tests__/low-trust-red-team-routes.test.ts index fdf5574386..b6028d9bd4 100644 --- a/server/src/__tests__/low-trust-red-team-routes.test.ts +++ b/server/src/__tests__/low-trust-red-team-routes.test.ts @@ -79,6 +79,29 @@ function isHeartbeatCleanupFkError(error: unknown) { ); } +// Await every background heartbeat run until the run table is quiescent. A route +// dispatches a wakeup fire-and-forget (void heartbeat.wakeup(...) in +// routes/issues.ts). Such a wakeup, or a run it dispatches, can write issues, +// issue_comments, and heartbeat_runs rows during teardown and race the deletes +// below (a heartbeat_runs delete deadlocks on the ON DELETE SET NULL cascade to +// issues; an issue_comments insert breaks the later delete of issues). +// drainActiveRunExecutions() awaits both in-flight wakeup promises and in-flight +// run executions, so it now also waits for a wakeup that is still before run +// registration. Re-check the run table after the drain as a backstop, and give a +// late run a macrotask before the next attempt, until no run is queued or running. +async function drainHeartbeatRunsToQuiescence( + db: Db, + heartbeat: ReturnType, +) { + for (let attempt = 0; attempt < 50; attempt += 1) { + await heartbeat.drainActiveRunExecutions(); + const runs = await db.select({ status: heartbeatRuns.status }).from(heartbeatRuns); + const hasPending = runs.some((run) => run.status === "queued" || run.status === "running"); + if (!hasPending) return; + await new Promise((resolve) => setTimeout(resolve, 10)); + } +} + async function deleteHeartbeatRunsAndWakeupsAfterActivityLogDrains(db: Db) { for (let attempt = 0; attempt < 10; attempt += 1) { await db.delete(heartbeatRunEvents); @@ -702,6 +725,11 @@ describeEmbeddedPostgres("low-trust red-team HTTP route regression suite", () => }, 20_000); afterEach(async () => { + // Await every in-flight background heartbeat run to quiescence before the + // deletes below. A route dispatches a wakeup fire-and-forget, so a run can + // still be writing issues, issue_comments, and heartbeat_runs rows when + // teardown starts and would race the deletes. + await drainHeartbeatRunsToQuiescence(db, heartbeatService(db)); await db.delete(issueThreadInteractions); await db.delete(issueApprovals); await db.delete(approvals); diff --git a/server/src/services/heartbeat.ts b/server/src/services/heartbeat.ts index 7d10535635..968f1ffe7e 100644 --- a/server/src/services/heartbeat.ts +++ b/server/src/services/heartbeat.ts @@ -579,6 +579,16 @@ const activeRunExecutions = new Set(); // that must guarantee no run write is still in flight (graceful shutdown, and // tests tearing down a shared database) can await drainActiveRunExecutions(). const activeRunExecutionPromises = new Set>(); +// Routes dispatch a wakeup fire-and-forget (void heartbeat.wakeup(...)). The +// wakeup promise stays pending through its asynchronous prologue, and it +// resolves only after it inserts the queued run and registers the run +// execution in activeRunExecutionPromises. Before that point neither +// activeRunExecutionPromises nor the run table shows the pending run, so a +// caller cannot observe the wake. Track each wakeup promise here — shared +// across service instances like the two sets above — so drainActiveRunExecutions +// can await a wake that is still before run registration. A caller that tears +// down a shared database (a test afterEach) then cannot race a late wake. +const activeWakeupPromises = new Set>(); const INLINE_BASE64_IMAGE_DATA_RE = /("type":"image","source":\{"type":"base64","data":")([A-Za-z0-9+/=]{1024,})(")/g; type RuntimeConfigSecretResolver = Pick< @@ -12097,12 +12107,37 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {}) // before the parent promise settles, so we loop until the set is empty rather // than snapshotting once. Callers use this to guarantee no run is still // writing rows/events (graceful shutdown, deterministic test teardown). + // + // Await in-flight wakeup promises first. A wakeup resolves only after it + // registers its run execution, so a wake that is still before run registration + // is invisible to activeRunExecutionPromises alone. Awaiting the wakeup promise + // closes that window: once it settles, any run it dispatched is already in + // activeRunExecutionPromises, and the second await drains that run. A wakeup or + // a run can add more entries as it settles, so loop until both sets are empty. async function drainActiveRunExecutions() { - while (activeRunExecutionPromises.size > 0) { + while (activeWakeupPromises.size > 0 || activeRunExecutionPromises.size > 0) { + await Promise.allSettled([...activeWakeupPromises]); await Promise.all([...activeRunExecutionPromises]); } } + // Public wakeup entry point. Callers dispatch it fire-and-forget, so register + // the promise in activeWakeupPromises before it starts its asynchronous + // prologue. drainActiveRunExecutions can then await a wake that is still before + // run registration. Internal callers reference enqueueWakeup directly and + // already await it, so they do not need this registration. + function trackWakeup( + agentId: string, + opts: WakeupOptions = {}, + ): ReturnType { + const promise = enqueueWakeup(agentId, opts); + activeWakeupPromises.add(promise); + void promise.catch(() => {}).finally(() => { + activeWakeupPromises.delete(promise); + }); + return promise; + } + async function executeRun(runId: string) { if ((await getSchedulingSuppression()).suppressed) return; @@ -17256,7 +17291,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {}) triggerDetail: "manual" | "ping" | "callback" | "system" = "manual", actor?: { actorType?: "user" | "agent" | "system"; actorId?: string | null }, ) => - enqueueWakeup(agentId, { + trackWakeup(agentId, { source, triggerDetail, contextSnapshot, @@ -17264,7 +17299,7 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {}) requestedByActorId: actor?.actorId ?? null, }), - wakeup: enqueueWakeup, + wakeup: trackWakeup, triggerIssueMonitor, reportRunActivity: clearDetachedRunWarning,