fix(runtime): stop Live badge and Working shimmer after task teardown (#10985)

## Thinking Path

> - Paperclip is the open source app people use to manage AI agents for
work.
> - The task list and the chat views show a Live badge and a Working
shimmer for an issue that has an active run.
> - A finished task kept the Live badge and the Working shimmer after
the run ended and the sandbox stopped.
> - The user interface reads run liveness from the
`heartbeat_runs.status` row. The run finalizer writes the terminal
status in a step that is separate from the agent `status=done` update.
When the sandbox or the run process stops between the two steps,
`heartbeat_runs.status` stays `running` forever.
> - A run row that stays `running` makes a finished task look
perpetually Live, and the user interface has no guard for an issue that
already reached a terminal status.
> - This pull request closes the invariant "environment lease released
implies the run is terminal" on the server, and adds a user interface
guard that suppresses live state for a terminal issue.
> - The benefit is that a finished task stops showing Live and Working,
both at the source (the run row) and at the surface (the badge and the
shimmer).

## Linked Issues or Issue Description

**Bug description**

- A completed task kept the Live badge and the Working shimmer after its
run ended and the sandbox was torn down.

**Steps to reproduce**

- Run an agent task to completion. Let the sandbox tear down while the
run finalizer is between the `status=done` update and the terminal
run-status write.
- Open the task list or the chat view for the finished task.

**Expected behavior**

- A finished task shows no Live badge and no Working shimmer.

**Actual behavior (before this change)**

- The finished task showed the Live badge and the Working shimmer
because its `heartbeat_runs.status` row stayed `running`.

This pull request supersedes the two separate pull requests #10954
(frontend) and #10955 (backend). It carries all of their changes for the
same race.

## What Changed

Server:

- Run teardown terminalizes a still-running or still-queued run before
it releases the environment lease. It writes `succeeded` when the issue
already reached `done`, `cancelled` when the issue is `cancelled`, and
`interrupted` otherwise. It never overwrites a status that another path
already made terminal.
- The recovery stale-lock sweep terminalizes an orphaned running run to
`interrupted` after it confirms the process and the sandbox are both
gone. It requires recorded process metadata, so it never terminalizes a
live run, a queued run, or a scheduled retry.
- Each terminal transition writes a run event.
- The stale-lock sweep continues and clears the lock when the audit
write fails. It logs the failure loudly.
- New server tests cover both invariants.

User interface:

- A shared guard suppresses the Live badge and the Working shimmer when
the issue status is terminal.
- The guard keeps non-terminal `queued` and `running` issues live.
- The guard prefers the newest issue live-status snapshot.
- New user interface tests cover the guard and the snapshot preference.

## Verification

Server:

- `pnpm --filter @paperclipai/plugin-sdk ensure-build-deps && tsc
--noEmit` in `server/` — 0 errors.
- `pnpm --filter server test
heartbeat-run-lease-release-terminalization.test.ts
recovery-stale-issue-lock-sweep.test.ts` — 12 tests pass.

User interface:

- `pnpm --filter @paperclipai/ui typecheck` — 0 errors.
- `pnpm exec vitest run ui/src/lib/liveIssueIds.test.ts
ui/src/lib/issue-chat-messages.test.ts` — 40 tests pass.

## Risks

- Low risk. The server change only forces a still-live run row to a
terminal status when the lease releases or when the recovery sweep
confirms the process is dead. It never overwrites an existing terminal
status, and it guards the recovery path with process metadata to avoid
terminalizing a live run.
- The user interface change is additive. The guard only suppresses live
state for a terminal issue and keeps queued and running issues live.
- No database migration. No change to any external endpoint.

## Model Used

- Claude Opus 4.8 (`claude-opus-4-8`), extended thinking, tool use, and
code execution.

## Checklist

- [x] I have included a thinking path that traces from project context
to this change
- [x] I have specified the model used (with version and capability
details)
- [x] I have checked ROADMAP.md and confirmed this PR does not duplicate
planned core work
- [x] I have searched GitHub for duplicate or related PRs and linked
them above
- [x] I have either (a) linked existing issues with `Fixes: #` /
`Closes: #` / `Refs: #` OR (b) described the issue in-PR following the
relevant issue template
- [x] I have not referenced internal/instance-local Paperclip issues or
links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip`
URLs)
- [x] My branch name describes the change (e.g. `docs/...`, `fix/...`)
and contains no internal Paperclip ticket id or instance-derived details
- [x] I have run tests locally and they pass
- [x] I have added or updated tests where applicable
- [x] I have updated relevant documentation to reflect my changes
- [x] I have considered and documented any risks above
- [x] All Paperclip CI gates are green
- [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups
- [x] I will address all Greptile and reviewer comments before
requesting merge

---------

Co-authored-by: Paperclip <noreply@paperclip.ing>
This commit is contained in:
Nicky Leach 2026-08-07 07:23:56 -07:00 committed by GitHub
parent 4683f26c97
commit 01b51dc0e5
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
17 changed files with 1113 additions and 31 deletions

View File

@ -0,0 +1,212 @@
import { randomUUID } from "node:crypto";
import { eq } from "drizzle-orm";
import { afterAll, afterEach, beforeAll, describe, expect, it, vi } from "vitest";
import {
agents,
companies,
createDb,
heartbeatRunEvents,
heartbeatRuns,
issues,
} from "@paperclipai/db";
import {
getEmbeddedPostgresTestSupport,
startEmbeddedPostgresTestDatabase,
} from "./helpers/embedded-postgres.js";
const mockTelemetryClient = vi.hoisted(() => ({ track: vi.fn() }));
vi.mock("../telemetry.ts", () => ({ getTelemetryClient: () => mockTelemetryClient }));
import { heartbeatService } from "../services/heartbeat.ts";
const embeddedPostgresSupport = await getEmbeddedPostgresTestSupport();
const describeEmbeddedPostgres = embeddedPostgresSupport.supported ? describe : describe.skip;
if (!embeddedPostgresSupport.supported) {
console.warn(
`Skipping embedded Postgres lease-release terminalization tests on this host: ${
embeddedPostgresSupport.reason ?? "unsupported environment"
}`,
);
}
describeEmbeddedPostgres("heartbeat terminalizeRunOnLeaseRelease", () => {
let db!: ReturnType<typeof createDb>;
let tempDb: Awaited<ReturnType<typeof startEmbeddedPostgresTestDatabase>> | null = null;
beforeAll(async () => {
tempDb = await startEmbeddedPostgresTestDatabase("paperclip-lease-release-terminal-");
db = createDb(tempDb.connectionString);
}, 20_000);
afterEach(async () => {
await db.delete(heartbeatRunEvents);
await db.delete(issues);
await db.delete(heartbeatRuns);
await db.delete(agents);
await db.delete(companies);
});
afterAll(async () => {
await tempDb?.cleanup();
});
async function seed(input: { issueStatus: string; runStatus: string }) {
const companyId = randomUUID();
const agentId = randomUUID();
const issueId = randomUUID();
const runId = randomUUID();
await db.insert(companies).values({
id: companyId,
name: "Paperclip",
issuePrefix: `T${companyId.replace(/-/g, "").slice(0, 6).toUpperCase()}`,
requireBoardApprovalForNewAgents: false,
});
await db.insert(agents).values({
id: agentId,
companyId,
name: "Coder",
role: "engineer",
status: "active",
adapterType: "codex_local",
adapterConfig: {},
runtimeConfig: {},
permissions: {},
});
await db.insert(issues).values({
id: issueId,
companyId,
title: "Terminalize on lease release",
status: input.issueStatus,
priority: "high",
assigneeAgentId: agentId,
});
await db.insert(heartbeatRuns).values({
id: runId,
companyId,
agentId,
status: input.runStatus,
invocationSource: "manual",
startedAt: new Date(),
contextSnapshot: { issueId },
});
const run = await db
.select()
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runId))
.then((rows) => rows[0]!);
return { companyId, agentId, issueId, runId, run };
}
it("forces a still-running run to succeeded when the issue already reached done", async () => {
// This reproduces the defect: the agent PATCHed the issue to done, but the
// teardown released the environment lease before the run-terminal write.
const { issueId, runId, run } = await seed({ issueStatus: "done", runStatus: "running" });
const heartbeat = heartbeatService(db);
const terminal = await heartbeat.terminalizeRunOnLeaseRelease(run);
expect(terminal.status).toBe("succeeded");
const runStatus = await db
.select({ status: heartbeatRuns.status })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runId))
.then((rows) => rows[0]?.status);
expect(runStatus).toBe("succeeded");
const issueStatus = await db
.select({ status: issues.status })
.from(issues)
.where(eq(issues.id, issueId))
.then((rows) => rows[0]?.status);
expect(issueStatus).toBe("done");
const event = await db
.select({ message: heartbeatRunEvents.message, payload: heartbeatRunEvents.payload })
.from(heartbeatRunEvents)
.where(eq(heartbeatRunEvents.runId, runId))
.then((rows) => rows[0]);
expect(event?.message).toContain("lease release");
expect((event?.payload as { terminalStatus?: string } | null)?.terminalStatus).toBe("succeeded");
});
it("forces a still-running run to interrupted when the issue is not terminal", async () => {
const { runId, run } = await seed({ issueStatus: "in_progress", runStatus: "running" });
const heartbeat = heartbeatService(db);
const terminal = await heartbeat.terminalizeRunOnLeaseRelease(run);
expect(terminal.status).toBe("interrupted");
const row = await db
.select({ status: heartbeatRuns.status, errorCode: heartbeatRuns.errorCode })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runId))
.then((rows) => rows[0]);
expect(row?.status).toBe("interrupted");
expect(row?.errorCode).toBe("lease_released_before_terminal");
});
it("forces a still-queued run to interrupted when the lease releases before it starts", async () => {
// A queued run holds a lease but never reached "running". The teardown
// released the lease, so the run must not stay queued and show a phantom
// live run. A running-only update would miss it.
const { runId, run } = await seed({ issueStatus: "in_progress", runStatus: "queued" });
const heartbeat = heartbeatService(db);
const terminal = await heartbeat.terminalizeRunOnLeaseRelease(run);
expect(terminal.status).toBe("interrupted");
const row = await db
.select({ status: heartbeatRuns.status, errorCode: heartbeatRuns.errorCode })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runId))
.then((rows) => rows[0]);
expect(row?.status).toBe("interrupted");
expect(row?.errorCode).toBe("lease_released_before_terminal");
});
it("forces a still-queued run to succeeded when the issue already reached done", async () => {
const { runId, run } = await seed({ issueStatus: "done", runStatus: "queued" });
const heartbeat = heartbeatService(db);
const terminal = await heartbeat.terminalizeRunOnLeaseRelease(run);
expect(terminal.status).toBe("succeeded");
const runStatus = await db
.select({ status: heartbeatRuns.status })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runId))
.then((rows) => rows[0]?.status);
expect(runStatus).toBe("succeeded");
});
it("keeps an already-terminal run authoritative and records no new event", async () => {
const { runId, run } = await seed({ issueStatus: "done", runStatus: "failed" });
const heartbeat = heartbeatService(db);
const terminal = await heartbeat.terminalizeRunOnLeaseRelease(run);
expect(terminal.status).toBe("failed");
const runStatus = await db
.select({ status: heartbeatRuns.status })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runId))
.then((rows) => rows[0]?.status);
expect(runStatus).toBe("failed");
const eventCount = await db
.select({ id: heartbeatRunEvents.id })
.from(heartbeatRunEvents)
.where(eq(heartbeatRunEvents.runId, runId))
.then((rows) => rows.length);
expect(eventCount).toBe(0);
});
});

View File

@ -6,6 +6,7 @@ import {
agents,
companies,
createDb,
heartbeatRunEvents,
heartbeatRuns,
issueComments,
issueRelations,
@ -46,6 +47,7 @@ describeEmbeddedPostgres("recovery sweepStaleIssueLocks", () => {
await db.delete(issueRelations);
await db.delete(activityLog);
await db.delete(issues);
await db.delete(heartbeatRunEvents);
await db.delete(heartbeatRuns);
await db.delete(agents);
await db.delete(companies);
@ -223,4 +225,369 @@ describeEmbeddedPostgres("recovery sweepStaleIssueLocks", () => {
expect(first.cleared).toBe(1);
expect(second.cleared).toBe(0);
});
it("terminalizes an orphaned running run whose process is gone, then clears the lock", async () => {
const { companyId, agentId, runningRunId } = await seed();
// The run recorded a pid, but the process and its sandbox are gone. A pid
// this large never maps to a live process, so isPidAlive returns false.
// The issue is not terminal, so only the process-death authority applies.
await db
.update(heartbeatRuns)
.set({ processPid: 2_000_000_000 })
.where(eq(heartbeatRuns.id, runningRunId));
const issueId = randomUUID();
await db.insert(issues).values({
id: issueId,
companyId,
title: "Orphaned running run — terminalize then clear",
status: "in_progress",
priority: "high",
assigneeAgentId: agentId,
checkoutRunId: runningRunId,
executionRunId: runningRunId,
executionLockedAt: new Date(),
});
const heartbeat = heartbeatService(db);
const result = await heartbeat.sweepStaleIssueLocks();
expect(result.terminalizedRunIds).toEqual([runningRunId]);
expect(result.cleared).toBe(1);
const run = await db
.select({ status: heartbeatRuns.status, errorCode: heartbeatRuns.errorCode })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runningRunId))
.then((rows) => rows[0]);
// Process died, outcome unknown, so the backstop uses "interrupted".
expect(run?.status).toBe("interrupted");
expect(run?.errorCode).toBe("orphaned_running_run");
const lock = await db
.select({ checkoutRunId: issues.checkoutRunId, executionRunId: issues.executionRunId })
.from(issues)
.where(eq(issues.id, issueId))
.then((rows) => rows[0]);
expect(lock).toEqual({ checkoutRunId: null, executionRunId: null });
const event = await db
.select({ message: heartbeatRunEvents.message })
.from(heartbeatRunEvents)
.where(eq(heartbeatRunEvents.runId, runningRunId))
.then((rows) => rows[0]);
expect(event?.message).toContain("process and sandbox gone");
});
it("terminalizes a running run whose issue is terminal, even while the process stays alive (reuse-lease path)", async () => {
// Reuse Lease ON stops the sandbox but keeps the server process alive, so
// the in-memory handle and the recorded pid can both persist. The
// process-death authority misses this case. The issue-terminal authority
// catches it: the issue reached "done" while the run row stayed "running".
const { companyId, agentId, runningRunId } = await seed();
// process.pid is the live test process, so isPidAlive returns true.
await db
.update(heartbeatRuns)
.set({ processPid: process.pid })
.where(eq(heartbeatRuns.id, runningRunId));
const issueId = randomUUID();
await db.insert(issues).values({
id: issueId,
companyId,
title: "Reused sandbox stopped — issue done, run still running",
status: "done",
priority: "high",
assigneeAgentId: agentId,
checkoutRunId: runningRunId,
executionRunId: runningRunId,
executionLockedAt: new Date(),
});
const heartbeat = heartbeatService(db);
const result = await heartbeat.sweepStaleIssueLocks();
expect(result.terminalizedRunIds).toEqual([runningRunId]);
expect(result.cleared).toBe(1);
const run = await db
.select({ status: heartbeatRuns.status, errorCode: heartbeatRuns.errorCode })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runningRunId))
.then((rows) => rows[0]);
// The issue is "done", so the terminal run status is "succeeded". A
// succeeded run carries no error code.
expect(run?.status).toBe("succeeded");
expect(run?.errorCode).toBeNull();
const lock = await db
.select({ checkoutRunId: issues.checkoutRunId, executionRunId: issues.executionRunId })
.from(issues)
.where(eq(issues.id, issueId))
.then((rows) => rows[0]);
expect(lock).toEqual({ checkoutRunId: null, executionRunId: null });
const event = await db
.select({ message: heartbeatRunEvents.message })
.from(heartbeatRunEvents)
.where(eq(heartbeatRunEvents.runId, runningRunId))
.then((rows) => rows[0]);
expect(event?.message).toContain("issue reached a terminal status");
});
it("terminalizes a running run to cancelled when its issue is cancelled (reuse-lease path)", async () => {
const { companyId, agentId, runningRunId } = await seed();
await db
.update(heartbeatRuns)
.set({ processPid: process.pid })
.where(eq(heartbeatRuns.id, runningRunId));
const issueId = randomUUID();
await db.insert(issues).values({
id: issueId,
companyId,
title: "Reused sandbox stopped — issue cancelled, run still running",
status: "cancelled",
priority: "high",
assigneeAgentId: agentId,
checkoutRunId: runningRunId,
executionRunId: runningRunId,
executionLockedAt: new Date(),
});
const heartbeat = heartbeatService(db);
const result = await heartbeat.sweepStaleIssueLocks();
expect(result.terminalizedRunIds).toEqual([runningRunId]);
const runStatus = await db
.select({ status: heartbeatRuns.status })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runningRunId))
.then((rows) => rows[0]?.status);
expect(runStatus).toBe("cancelled");
});
it("does not terminalize a running run whose process is alive and whose issue is not terminal", async () => {
const { companyId, agentId, runningRunId } = await seed();
// process.pid is the live test process, so isPidAlive returns true.
await db
.update(heartbeatRuns)
.set({ processPid: process.pid })
.where(eq(heartbeatRuns.id, runningRunId));
const issueId = randomUUID();
await db.insert(issues).values({
id: issueId,
companyId,
title: "Live run — preserve",
status: "in_progress",
priority: "high",
assigneeAgentId: agentId,
checkoutRunId: runningRunId,
executionRunId: runningRunId,
executionLockedAt: new Date(),
});
const heartbeat = heartbeatService(db);
const result = await heartbeat.sweepStaleIssueLocks();
expect(result.terminalizedRunIds).toEqual([]);
expect(result.cleared).toBe(0);
const runStatus = await db
.select({ status: heartbeatRuns.status })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runningRunId))
.then((rows) => rows[0]?.status);
expect(runStatus).toBe("running");
});
it("does not terminalize a live run that a terminal issue and an active issue both reference", async () => {
// A stale lock on a terminal issue and the real lock on an active issue can
// point at the same running run. The terminal reference alone must not
// terminalize the run, because the run is still live for the active issue.
const { companyId, agentId, runningRunId } = await seed();
// process.pid is the live test process, so isPidAlive returns true.
await db
.update(heartbeatRuns)
.set({ processPid: process.pid })
.where(eq(heartbeatRuns.id, runningRunId));
const terminalIssueId = randomUUID();
const activeIssueId = randomUUID();
await db.insert(issues).values([
{
id: terminalIssueId,
companyId,
title: "Terminal issue holds a stale lock on the shared run",
status: "done",
priority: "high",
assigneeAgentId: agentId,
checkoutRunId: runningRunId,
executionRunId: null,
},
{
id: activeIssueId,
companyId,
title: "Active issue owns the live shared run",
status: "in_progress",
priority: "high",
assigneeAgentId: agentId,
checkoutRunId: runningRunId,
executionRunId: runningRunId,
executionLockedAt: new Date(),
},
]);
const heartbeat = heartbeatService(db);
const result = await heartbeat.sweepStaleIssueLocks();
// The run stays live, so the sweep terminalizes nothing and clears nothing.
expect(result.terminalizedRunIds).toEqual([]);
expect(result.cleared).toBe(0);
const runStatus = await db
.select({ status: heartbeatRuns.status })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runningRunId))
.then((rows) => rows[0]?.status);
expect(runStatus).toBe("running");
const activeLock = await db
.select({ checkoutRunId: issues.checkoutRunId, executionRunId: issues.executionRunId })
.from(issues)
.where(eq(issues.id, activeIssueId))
.then((rows) => rows[0]);
expect(activeLock).toEqual({ checkoutRunId: runningRunId, executionRunId: runningRunId });
});
it("does not terminalize a shared live run when its context snapshot names the terminal issue", async () => {
// The run context snapshot names the terminal issue. The context-snapshot
// fallback in terminalizeOrphanedRunningRun could read that terminal status
// and terminalize the run. An active issue still owns the run, so the sweep
// must suppress the fallback and keep the run live.
const { companyId, agentId, runningRunId } = await seed();
// process.pid is the live test process, so isPidAlive returns true.
await db
.update(heartbeatRuns)
.set({ processPid: process.pid })
.where(eq(heartbeatRuns.id, runningRunId));
const terminalIssueId = randomUUID();
const activeIssueId = randomUUID();
// The run context snapshot names the terminal issue. This is the path the
// shared-run guard must still block.
await db
.update(heartbeatRuns)
.set({ contextSnapshot: { issueId: terminalIssueId } })
.where(eq(heartbeatRuns.id, runningRunId));
await db.insert(issues).values([
{
id: terminalIssueId,
companyId,
title: "Terminal issue named in the run context snapshot",
status: "done",
priority: "high",
assigneeAgentId: agentId,
checkoutRunId: runningRunId,
executionRunId: null,
},
{
id: activeIssueId,
companyId,
title: "Active issue owns the live shared run",
status: "in_progress",
priority: "high",
assigneeAgentId: agentId,
checkoutRunId: runningRunId,
executionRunId: runningRunId,
executionLockedAt: new Date(),
},
]);
const heartbeat = heartbeatService(db);
const result = await heartbeat.sweepStaleIssueLocks();
// The run stays live, so the sweep terminalizes nothing and clears nothing.
expect(result.terminalizedRunIds).toEqual([]);
expect(result.cleared).toBe(0);
const runStatus = await db
.select({ status: heartbeatRuns.status })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runningRunId))
.then((rows) => rows[0]?.status);
expect(runStatus).toBe("running");
const activeLock = await db
.select({ checkoutRunId: issues.checkoutRunId, executionRunId: issues.executionRunId })
.from(issues)
.where(eq(issues.id, activeIssueId))
.then((rows) => rows[0]);
expect(activeLock).toEqual({ checkoutRunId: runningRunId, executionRunId: runningRunId });
});
it("still clears the lock when the audit write fails after terminalization", async () => {
const { companyId, agentId, runningRunId } = await seed();
// The run recorded a pid that never maps to a live process, so the sweep
// decides to terminalize it. The issue is not terminal, so the
// process-death authority drives the terminalization here.
await db
.update(heartbeatRuns)
.set({ processPid: 2_000_000_000 })
.where(eq(heartbeatRuns.id, runningRunId));
const issueId = randomUUID();
await db.insert(issues).values({
id: issueId,
companyId,
title: "Audit write fails — still clear the lock",
status: "in_progress",
priority: "high",
assigneeAgentId: agentId,
checkoutRunId: runningRunId,
executionRunId: runningRunId,
executionLockedAt: new Date(),
});
// Make only the audit-event insert fail. The run update commits the
// terminal status first, so the audit write is best-effort. The sweep must
// catch the failure and still clear the lock.
const realInsert = db.insert.bind(db);
const insertSpy = vi.spyOn(db, "insert").mockImplementation((table) => {
if (table === heartbeatRunEvents) {
throw new Error("simulated audit write failure");
}
return realInsert(table);
});
try {
const heartbeat = heartbeatService(db);
const result = await heartbeat.sweepStaleIssueLocks();
expect(result.terminalizedRunIds).toEqual([runningRunId]);
expect(result.cleared).toBe(1);
} finally {
insertSpy.mockRestore();
}
// The run reached its terminal status even though the audit write failed.
const runStatus = await db
.select({ status: heartbeatRuns.status })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, runningRunId))
.then((rows) => rows[0]?.status);
expect(runStatus).toBe("interrupted");
// The sweep cleared the lock in the same pass.
const lock = await db
.select({ checkoutRunId: issues.checkoutRunId, executionRunId: issues.executionRunId })
.from(issues)
.where(eq(issues.id, issueId))
.then((rows) => rows[0]);
expect(lock).toEqual({ checkoutRunId: null, executionRunId: null });
// The audit write failed, so no run event exists for this run.
const events = await db
.select({ id: heartbeatRunEvents.id })
.from(heartbeatRunEvents)
.where(eq(heartbeatRunEvents.runId, runningRunId));
expect(events).toEqual([]);
});
});

View File

@ -8825,11 +8825,25 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
runId: string,
status: string,
patch?: Partial<typeof heartbeatRuns.$inferInsert>,
) {
return setRunStatusFromLive(runId, status, ["running"], patch);
}
// Move a run to a new status only when its current status is one of
// `fromStatuses`. The compare-and-set is a single conditional update, so a
// concurrent path can win the race. When this update matches nothing, the
// function reads the current row and reports updated=false, so the caller can
// keep the terminal outcome that another path already wrote.
async function setRunStatusFromLive(
runId: string,
status: string,
fromStatuses: string[],
patch?: Partial<typeof heartbeatRuns.$inferInsert>,
) {
const updated = await db
.update(heartbeatRuns)
.set({ status, ...patch, updatedAt: new Date() })
.where(and(eq(heartbeatRuns.id, runId), eq(heartbeatRuns.status, "running")))
.where(and(eq(heartbeatRuns.id, runId), inArray(heartbeatRuns.status, fromStatuses)))
.returning()
.then((rows) => rows[0] ?? null);
@ -8855,6 +8869,75 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
return { run: current, updated: false as const };
}
// Invariant: when a run releases its environment lease, the run row must be
// terminal. The finalizer writes the terminal status in a step that is
// separate from the agent status=done PATCH. If the sandbox or the run
// process stops between the two steps, heartbeat_runs.status stays "running".
// The UI reads liveness from that row, so a finished task shows "Live"
// forever. This function closes the gap in the run teardown: when the run is
// still running or queued, it forces a terminal status before the lease is
// released. It never overwrites a status that another path already made
// terminal.
async function terminalizeRunOnLeaseRelease(
run: typeof heartbeatRuns.$inferSelect,
): Promise<typeof heartbeatRuns.$inferSelect> {
if (isHeartbeatRunTerminalStatus(run.status)) return run;
if (run.status !== "running" && run.status !== "queued") return run;
// Choose the terminal status that reflects the true outcome. When the issue
// already reached a terminal status, the run reached its goal, so use the
// matching terminal run status. Otherwise the teardown cut the run short,
// so use "interrupted".
const issueId = readNonEmptyString(parseObject(run.contextSnapshot).issueId);
let terminalStatus: "succeeded" | "cancelled" | "interrupted" = "interrupted";
if (issueId) {
const issueStatus = await db
.select({ status: issues.status })
.from(issues)
.where(eq(issues.id, issueId))
.then((rows) => rows[0]?.status ?? null);
if (issueStatus === "done") terminalStatus = "succeeded";
else if (issueStatus === "cancelled") terminalStatus = "cancelled";
}
const message =
`run terminalized on environment lease release: heartbeat_runs.status was still ${run.status} at teardown`;
// Match both "running" and "queued". A queued run has released its lease but
// never reached "running", so a running-only update would miss it and leave
// a phantom live run behind.
const write = await setRunStatusFromLive(run.id, terminalStatus, ["running", "queued"], {
finishedAt: run.finishedAt ?? new Date(),
error: run.error ?? (terminalStatus === "interrupted" ? message : null),
errorCode: run.errorCode ?? (terminalStatus === "interrupted" ? "lease_released_before_terminal" : null),
});
if (!write.updated) {
// Another path already finalized the run. Keep that terminal outcome.
return write.run ?? run;
}
const terminalRun = write.run;
if (terminalRun) {
await appendRunEvent(terminalRun, await nextRunEventSeq(terminalRun.id), {
eventType: "lifecycle",
stream: "system",
level: terminalStatus === "interrupted" ? "warn" : "info",
message,
payload: {
previousStatus: run.status,
terminalStatus,
reason: "environment_lease_release",
...(issueId ? { issueId } : {}),
},
}).catch((eventErr) => {
logger.warn(
{ err: eventErr, runId: run.id },
"failed to append run event for lease-release terminalization",
);
});
}
return terminalRun ?? run;
}
function publishRunLifecyclePluginEvent(run: typeof heartbeatRuns.$inferSelect) {
const eventType =
run.status === "running"
@ -16193,7 +16276,20 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
}
}
} finally {
const latestRun = await getRun(run.id).catch(() => null);
let latestRun = await getRun(run.id).catch(() => null);
// Close the invariant "environment lease released implies the run is
// terminal". When the teardown reaches this point with the run still
// running or queued, force a terminal status before the lease is
// released, so the UI never shows a finished task as "Live".
if (latestRun) {
latestRun = await terminalizeRunOnLeaseRelease(latestRun).catch((terminalizeErr) => {
logger.error(
{ err: terminalizeErr, runId: run.id },
"failed to terminalize run before environment lease release",
);
return latestRun;
});
}
await releaseEnvironmentLeasesForRun({
runId: run.id,
companyId: run.companyId,
@ -18935,6 +19031,8 @@ export function heartbeatService(db: Db, options: HeartbeatServiceOptions = {})
reconcileStrandedAssignedIssues,
terminalizeRunOnLeaseRelease,
sweepStaleIssueLocks,
buildIssueGraphLivenessAutoRecoveryPreview,

View File

@ -5477,22 +5477,178 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup })
return Math.max(1, Math.floor(asNumber(raw, fallback)));
}
// Backstop reconciler: terminalizes a "running" run that can no longer reach a
// terminal status on its own. The run finalizer writes the terminal status in
// a step that is separate from the agent status=done PATCH. When the teardown
// stops between the two steps, heartbeat_runs.status stays "running" forever.
// The UI reads liveness from that row, so the task shows "Live" forever. This
// function forces the run to a terminal status and records a run event, so the
// state is auditable. It never overwrites a status that another path already
// made terminal.
//
// Two independent authorities terminalize the run. Either one is enough:
//
// - Issue-terminal authority: the run's issue already reached a terminal
// status (done or cancelled), but the run row is still "running". A healthy
// run always terminalizes its own row before or just after the issue reaches
// a terminal status, so a lasting "running" row under a terminal issue is
// orphaned. This authority does not depend on process death. It is the only
// authority that catches the reuse-lease path: the release stops the sandbox
// but keeps the server process alive, so the in-memory handle and the
// recorded pid can both persist.
// - Process-death authority: the run has no in-memory handle and its recorded
// process and process group are both gone. This catches a hard server crash
// that skipped the graceful teardown, even when the issue is not terminal.
async function terminalizeOrphanedRunningRun(
run: typeof heartbeatRuns.$inferSelect,
options?: {
// The terminal run status implied by a referencing issue. The caller
// passes it when it already knows the issue that holds the run in a lock
// column. It maps issue "done" to "succeeded" and issue "cancelled" to
// "cancelled". A null value means the referencing issue is not terminal.
referencingIssueTerminalStatus?: "succeeded" | "cancelled" | null;
// True when an active (non-terminal) issue still holds this run in a lock
// column. The run is live for that active issue, so the caller forbids the
// issue-terminal authority. This flag also suppresses the context-snapshot
// fallback below. Without it, a terminal issue named in the run context
// snapshot would still terminalize the shared run and defeat the guard.
runReferencedByActiveIssue?: boolean;
},
): Promise<{ terminalized: boolean; status: string }> {
// Act only on a run in "running" status. A "queued" run has no process yet,
// and a "scheduled_retry" run has no process on purpose because it waits to
// retry. Neither is orphaned, so this function must not terminalize them.
if (run.status !== "running") return { terminalized: false, status: run.status };
const pid = run.processPid ?? null;
const processGroupId = run.processGroupId ?? null;
// Issue-terminal authority. When the run's issue is terminal, the run row is
// orphaned regardless of process or handle state. Prefer the referencing
// issue status that the caller passed, because a lock column is the direct
// link from the stuck "Live" issue to this run. Fall back to the issue id in
// the run context snapshot when the caller passed nothing. Skip the fallback
// when an active issue still references the run. The run is live for that
// active issue, so a terminal issue named in the context snapshot must not
// terminalize it.
let issueTerminalStatus: "succeeded" | "cancelled" | null =
options?.referencingIssueTerminalStatus ?? null;
const issueId = issueIdFromRunContext(run.contextSnapshot);
if (!issueTerminalStatus && !options?.runReferencedByActiveIssue && issueId) {
const issueStatus = await db
.select({ status: issues.status })
.from(issues)
.where(eq(issues.id, issueId))
.then((rows) => rows[0]?.status ?? null);
if (issueStatus === "done") issueTerminalStatus = "succeeded";
else if (issueStatus === "cancelled") issueTerminalStatus = "cancelled";
}
// Process-death authority. The run is live only when a process still backs
// it. Check the in-memory handle first, then the recorded pid and process
// group. Require recorded process metadata, so this authority never fires on
// a run that has not yet stored its pid.
let processGone = false;
if (!runningProcesses.get(run.id)) {
if (typeof pid === "number" || typeof processGroupId === "number") {
const processAlive =
(typeof pid === "number" && isPidAlive(pid)) ||
(typeof processGroupId === "number" && isProcessGroupAlive(processGroupId));
processGone = !processAlive;
}
}
// Neither authority applies. The run is still live, so leave it alone.
if (!issueTerminalStatus && !processGone) {
return { terminalized: false, status: run.status };
}
const authority = issueTerminalStatus ? "issue_terminal" : "process_gone";
const terminalStatus = issueTerminalStatus ?? "interrupted";
const errorCode = issueTerminalStatus
? "orphaned_running_run_issue_terminal"
: "orphaned_running_run";
const message =
authority === "issue_terminal"
? "run terminalized by recovery backstop: issue reached a terminal status while heartbeat_runs.status stayed live"
: "run terminalized by recovery backstop: process and sandbox gone while heartbeat_runs.status stayed live";
const now = new Date();
const updated = await db
.update(heartbeatRuns)
.set({
status: terminalStatus,
finishedAt: run.finishedAt ?? now,
error: run.error ?? (terminalStatus === "interrupted" ? message : null),
errorCode: run.errorCode ?? (terminalStatus === "interrupted" ? errorCode : null),
updatedAt: now,
})
.where(and(eq(heartbeatRuns.id, run.id), eq(heartbeatRuns.status, "running")))
.returning()
.then((rows) => rows[0] ?? null);
if (!updated) {
// Another path finalized the run between the read and this write. Keep
// that terminal outcome authoritative.
const [current] = await db
.select({ status: heartbeatRuns.status })
.from(heartbeatRuns)
.where(eq(heartbeatRuns.id, run.id));
return { terminalized: false, status: current?.status ?? run.status };
}
runningProcesses.delete(run.id);
// The run update above already committed the terminal status. The audit
// event is best-effort: if the insert fails, the caller must still treat
// the run as terminalized and clear the lock in the same sweep. So catch
// the failure, log it, and continue. A thrown error here would abort the
// sweep and leave the stale lock in place.
try {
await appendRecoveryRunEvent(updated, {
level: "warn",
message,
payload: {
source: "recovery.sweep_stale_issue_locks",
authority,
previousStatus: run.status,
terminalStatus,
...(issueId ? { issueId } : {}),
pid,
processGroupId,
},
});
} catch (error) {
logger.error(
{ err: error, runId: run.id, previousStatus: run.status },
"failed to append recovery run event after terminalizing orphaned run; run stays terminal and the sweep clears the lock",
);
}
logger.warn(
{ runId: run.id, authority, previousStatus: run.status, terminalStatus, issueId, pid, processGroupId },
"terminalized orphaned running heartbeat run in stale-lock sweep",
);
return { terminalized: true, status: updated.status };
}
// Backstop sweeper: clears stale lock columns on issues whose checkoutRunId
// or executionRunId points at a heartbeat_runs row that is either missing or
// in a terminal status. Provides self-heal for stale locks that fell outside
// releaseIssueExecutionAndPromote / clearCheckoutRunIfTerminal / adoption.
// Idempotent and safe: clears at most one row's worth of lock columns per
// candidate, and only when the referenced run row is unambiguously terminal.
// Before it evaluates cleanability, it terminalizes any referenced run that
// still claims to be live but can no longer reach a terminal status on its
// own, so a stuck "running" run can no longer block the sweep. Idempotent and
// safe: clears at most one row's worth of lock columns per candidate.
async function sweepStaleIssueLocks() {
const result = {
cleared: 0,
issueIds: [] as string[],
terminalizedRunIds: [] as string[],
};
const candidates = await db
.select({
id: issues.id,
companyId: issues.companyId,
status: issues.status,
checkoutRunId: issues.checkoutRunId,
executionRunId: issues.executionRunId,
})
@ -5511,13 +5667,61 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup })
const runRows =
referencedRunIds.length > 0
? await db
.select({ id: heartbeatRuns.id, status: heartbeatRuns.status })
.select()
.from(heartbeatRuns)
.where(inArray(heartbeatRuns.id, referencedRunIds))
: [];
const runStatusById = new Map<string, string>();
for (const row of runRows) runStatusById.set(row.id, row.status);
// Collect the runs that a non-terminal issue still references. Such a run is
// the live run of an active issue. A different, terminal issue can also hold
// the same run id in a stale lock column. The terminal reference alone must
// not terminalize a run that an active issue still owns, so exclude these
// runs from the issue-terminal authority below.
const runIdsReferencedByActiveIssue = new Set<string>();
for (const issue of candidates) {
if (issue.status === "done" || issue.status === "cancelled") continue;
for (const runId of [issue.checkoutRunId, issue.executionRunId]) {
if (runId) runIdsReferencedByActiveIssue.add(runId);
}
}
// Map each referenced run to the terminal run status implied by its
// referencing issue. When a terminal issue still holds the run in a lock
// column, that run is orphaned: the issue is the stuck "Live" task the UI
// shows. A "done" issue implies "succeeded"; a "cancelled" issue implies
// "cancelled". Skip a run that an active issue also references, because that
// run is still live for the active issue.
const issueTerminalStatusByRunId = new Map<string, "succeeded" | "cancelled">();
for (const issue of candidates) {
const implied =
issue.status === "done"
? "succeeded"
: issue.status === "cancelled"
? "cancelled"
: null;
if (!implied) continue;
for (const runId of [issue.checkoutRunId, issue.executionRunId]) {
if (runId && !runIdsReferencedByActiveIssue.has(runId)) {
issueTerminalStatusByRunId.set(runId, implied);
}
}
}
// Pre-pass: terminalize any referenced run that still claims to be live but
// can no longer reach a terminal status on its own. This lets the sweep
// clear the lock in the same pass instead of waiting for the run to reach a
// terminal status by another route.
for (const row of runRows) {
const outcome = await terminalizeOrphanedRunningRun(row, {
referencingIssueTerminalStatus: issueTerminalStatusByRunId.get(row.id) ?? null,
runReferencedByActiveIssue: runIdsReferencedByActiveIssue.has(row.id),
});
runStatusById.set(row.id, outcome.status);
if (outcome.terminalized) result.terminalizedRunIds.push(row.id);
}
const isCleanable = (runId: string | null) => {
if (!runId) return true;
const status = runStatusById.get(runId);
@ -5576,9 +5780,13 @@ export function recoveryService(db: Db, deps: { enqueueWakeup: RecoveryWakeup })
});
}
if (result.cleared > 0) {
if (result.cleared > 0 || result.terminalizedRunIds.length > 0) {
logger.warn(
{ cleared: result.cleared, issueIds: result.issueIds },
{
cleared: result.cleared,
issueIds: result.issueIds,
terminalizedRunIds: result.terminalizedRunIds,
},
"swept stale issue lock columns",
);
}

View File

@ -71,6 +71,7 @@ import type {
SuggestTasksInteraction,
} from "../lib/issue-thread-interactions";
import { buildIssueThreadInteractionSummary, isIssueThreadInteraction } from "../lib/issue-thread-interactions";
import { isLiveIssueRun } from "../lib/liveIssueIds";
import { resolveIssueChatTranscriptRuns } from "../lib/issueChatTranscriptRuns";
import {
formatTimelineWorkspaceLabel,
@ -4438,9 +4439,10 @@ export function IssueChatThread({
const displayLiveRuns = useMemo(() => {
const deduped = new Map<string, LiveRunForIssue>();
for (const run of liveRuns) {
if (!isLiveIssueRun(run, issueStatus)) continue;
deduped.set(run.id, run);
}
if (activeRun) {
if (activeRun && isLiveIssueRun(activeRun, issueStatus)) {
deduped.set(activeRun.id, {
id: activeRun.id,
status: activeRun.status,
@ -4471,7 +4473,7 @@ export function IssueChatThread({
});
}
return [...deduped.values()].sort((a, b) => new Date(a.createdAt).getTime() - new Date(b.createdAt).getTime());
}, [activeRun, liveRuns]);
}, [activeRun, issueStatus, liveRuns]);
const transcriptRuns = useMemo(() => {
return resolveIssueChatTranscriptRuns({
linkedRuns,
@ -4489,8 +4491,8 @@ export function IssueChatThread({
return ids;
}, [displayLiveRuns]);
const hasActiveRun = useMemo(
() => displayLiveRuns.some((run) => run.status === "running") || activeRun?.status === "running",
[displayLiveRuns, activeRun],
() => displayLiveRuns.some((run) => run.status === "running"),
[displayLiveRuns],
);
// Real-time view of the handoff: a run that starts after the issue payload
// was fetched must quiet the missing-disposition warnings without waiting
@ -4536,6 +4538,7 @@ export function IssueChatThread({
agentMap,
currentUserId,
userLabelMap,
issueStatus,
}),
[
comments,
@ -4552,6 +4555,7 @@ export function IssueChatThread({
agentMap,
currentUserId,
userLabelMap,
issueStatus,
],
);
const stableMessagesRef = useRef<readonly ThreadMessage[]>([]);

View File

@ -32,6 +32,7 @@ import { useSidebar } from "@/context/SidebarContext";
import { cn } from "@/lib/utils";
import { useIssuePlanDocument } from "@/hooks/useIssuePlanDocument";
import { latestSameRunHandoffTimestamp, type IssueChatComment } from "@/lib/issue-chat-messages";
import { isLiveIssueRun, isTerminalIssueStatus } from "@/lib/liveIssueIds";
import { workModeInEffectAt } from "@/lib/issue-timeline-events";
import { workModeMetaFor } from "@/lib/work-mode-meta";
@ -184,9 +185,10 @@ export function TaskChatThread(props: TaskChatThreadProps) {
// The single in-flight run whose turn we stream live (non-terminal).
const liveRun = useMemo(() => {
if (activeRun && !isTerminalRunStatus(activeRun.status)) return activeRun;
return (liveRuns ?? []).find((r) => !isTerminalRunStatus(r.status)) ?? null;
}, [activeRun, liveRuns]);
if (isTerminalIssueStatus(issueStatus)) return null;
if (activeRun && isLiveIssueRun(activeRun, issueStatus)) return activeRun;
return (liveRuns ?? []).find((r) => isLiveIssueRun(r, issueStatus)) ?? null;
}, [activeRun, issueStatus, liveRuns]);
// Runs observed non-terminal while mounted: their turns ANIMATE the fold when
// they settle. Runs already terminal at mount collapse instantly.

View File

@ -624,6 +624,44 @@ describe("buildIssueChatMessages", () => {
});
});
it("suppresses live-run Working messages for terminal issues", () => {
const liveRun: LiveRunForIssue = {
id: "run-live-terminal",
status: "running",
invocationSource: "manual",
triggerDetail: null,
startedAt: "2026-04-06T12:04:00.000Z",
finishedAt: null,
createdAt: "2026-04-06T12:04:00.000Z",
agentId: "agent-1",
agentName: "CodexCoder",
adapterType: "codex_local",
};
const terminalMessages = buildIssueChatMessages({
comments: [],
timelineEvents: [],
linkedRuns: [],
liveRuns: [liveRun],
issueStatus: "done",
currentUserId: "user-1",
});
const liveMessages = buildIssueChatMessages({
comments: [],
timelineEvents: [],
linkedRuns: [],
liveRuns: [liveRun],
issueStatus: "in_progress",
currentUserId: "user-1",
});
expect(terminalMessages.find((message) => message.id === "run-assistant:run-live-terminal")).toBeUndefined();
expect(liveMessages.find((message) => message.id === "run-assistant:run-live-terminal")).toMatchObject({
status: { type: "running" },
metadata: { custom: { waitingText: "Working..." } },
});
});
it("merges thread interactions into the same chronological feed as comments and runs", () => {
const messages = buildIssueChatMessages({
comments: [

View File

@ -16,6 +16,7 @@ import {
type IssueThreadInteraction,
} from "./issue-thread-interactions";
import type { IssueTimelineEvent } from "./issue-timeline-events";
import { isLiveIssueRun } from "./liveIssueIds";
import {
summarizeNotice,
} from "./transcriptPresentation";
@ -954,13 +955,15 @@ export function buildAssistantPartsFromTranscript(entries: readonly IssueChatTra
function normalizeLiveRuns(
liveRuns: readonly LiveRunForIssue[],
activeRun: ActiveRunForIssue | null | undefined,
issueId?: string,
issueId: string | undefined,
issueStatus: string | null | undefined,
) {
const deduped = new Map<string, LiveRunForIssue>();
for (const run of liveRuns) {
if (!isLiveIssueRun(run, issueStatus)) continue;
deduped.set(run.id, run);
}
if (activeRun) {
if (activeRun && isLiveIssueRun(activeRun, issueStatus)) {
deduped.set(activeRun.id, {
id: activeRun.id,
status: activeRun.status,
@ -1056,6 +1059,7 @@ export function buildIssueChatMessages(args: {
agentMap?: Map<string, Agent>;
currentUserId?: string | null;
userLabelMap?: ReadonlyMap<string, string> | null;
issueStatus?: string | null;
}) {
const {
comments,
@ -1073,6 +1077,7 @@ export function buildIssueChatMessages(args: {
agentMap,
currentUserId,
userLabelMap,
issueStatus,
} = args;
const orderedMessages: MessageWithOrder[] = [];
@ -1139,7 +1144,7 @@ export function buildIssueChatMessages(args: {
});
}
for (const run of normalizeLiveRuns(liveRuns, activeRun, issueId)) {
for (const run of normalizeLiveRuns(liveRuns, activeRun, issueId, issueStatus)) {
orderedMessages.push({
createdAtMs: toTimestamp(run.startedAt ?? run.createdAt),
order: 3,

View File

@ -2,6 +2,23 @@ import { describe, expect, it } from "vitest";
import type { LiveRunForIssue } from "../api/heartbeats";
import { collectLiveIssueIds, collectSubtreeLiveCounts } from "./liveIssueIds";
function liveRun(overrides: Partial<LiveRunForIssue>): LiveRunForIssue {
return {
id: "run",
status: "running",
invocationSource: "scheduler",
triggerDetail: null,
startedAt: "2026-04-20T10:00:00.000Z",
finishedAt: null,
createdAt: "2026-04-20T10:00:00.000Z",
agentId: "agent",
agentName: "Agent",
adapterType: "codex_local",
issueId: "issue",
...overrides,
};
}
describe("collectLiveIssueIds", () => {
it("keeps only runs linked to issues", () => {
const liveRuns: LiveRunForIssue[] = [
@ -74,6 +91,72 @@ describe("collectLiveIssueIds", () => {
expect([...collectLiveIssueIds(liveRuns)]).toEqual(["issue-1", "issue-2"]);
});
it("suppresses live ids for terminal issues while keeping non-terminal issues live", () => {
const liveRuns: LiveRunForIssue[] = [
{
id: "run-terminal",
status: "running",
invocationSource: "scheduler",
triggerDetail: null,
startedAt: "2026-04-20T10:00:00.000Z",
finishedAt: null,
createdAt: "2026-04-20T10:00:00.000Z",
agentId: "agent-1",
agentName: "Coder",
adapterType: "codex_local",
issueId: "issue-done",
},
{
id: "run-live",
status: "queued",
invocationSource: "scheduler",
triggerDetail: null,
startedAt: null,
finishedAt: null,
createdAt: "2026-04-20T10:01:00.000Z",
agentId: "agent-2",
agentName: "Builder",
adapterType: "codex_local",
issueId: "issue-open",
},
];
expect([...collectLiveIssueIds(liveRuns, [
{ id: "issue-done", status: "done" },
{ id: "issue-open", status: "in_progress" },
])]).toEqual(["issue-open"]);
});
it("keeps newer terminal snapshots authoritative when stale non-terminal snapshots appear later", () => {
const liveRuns: LiveRunForIssue[] = [
liveRun({ id: "run-done", issueId: "issue-done", status: "running" }),
liveRun({ id: "run-cancelled", issueId: "issue-cancelled", status: "queued" }),
liveRun({ id: "run-open", issueId: "issue-open", status: "running" }),
];
expect([...collectLiveIssueIds(liveRuns, [
{ id: "issue-done", status: "done", updatedAt: "2026-04-20T10:02:00.000Z" },
{ id: "issue-cancelled", status: "cancelled", updatedAt: "2026-04-20T10:02:00.000Z" },
{ id: "issue-open", status: "in_progress", updatedAt: "2026-04-20T10:02:00.000Z" },
{ id: "issue-done", status: "in_progress", updatedAt: "2026-04-20T10:01:00.000Z" },
{ id: "issue-cancelled", status: "todo", updatedAt: "2026-04-20T10:01:00.000Z" },
])]).toEqual(["issue-open"]);
});
it("allows a newer non-terminal snapshot to reopen an issue with a stale terminal snapshot", () => {
const liveRuns: LiveRunForIssue[] = [
liveRun({ id: "run-reopened", issueId: "issue-reopened", status: "running" }),
liveRun({ id: "run-terminal", issueId: "issue-terminal", status: "queued" }),
];
expect([...collectLiveIssueIds(liveRuns, [
{ id: "issue-reopened", status: "done", updatedAt: "2026-04-20T10:01:00.000Z" },
{ id: "issue-terminal", status: "in_progress", updatedAt: "2026-04-20T10:01:00.000Z" },
{ id: "issue-reopened", status: "in_progress", updatedAt: "2026-04-20T10:02:00.000Z" },
{ id: "issue-terminal", status: "done", updatedAt: "2026-04-20T10:02:00.000Z" },
])]).toEqual(["issue-reopened"]);
});
});
describe("collectSubtreeLiveCounts", () => {

View File

@ -1,13 +1,71 @@
import type { IssueStatus } from "@paperclipai/shared";
import type { LiveRunForIssue } from "../api/heartbeats";
function isLiveRunStatus(status: string): boolean {
return status === "queued" || status === "running";
}
export function collectLiveIssueIds(liveRuns: readonly LiveRunForIssue[] | null | undefined): Set<string> {
const TERMINAL_ISSUE_STATUSES = new Set<IssueStatus>(["done", "cancelled"]);
export function isTerminalIssueStatus(status: string | null | undefined): status is IssueStatus {
return TERMINAL_ISSUE_STATUSES.has(status as IssueStatus);
}
export function isLiveIssueRun(
run: Pick<LiveRunForIssue, "status">,
issueStatus?: string | null,
): boolean {
return isLiveRunStatus(run.status) && !isTerminalIssueStatus(issueStatus);
}
export interface LiveIssueStatusNode {
id: string;
status: IssueStatus | string;
updatedAt?: Date | string | number | null;
}
function collectIssueStatusById(issues: readonly LiveIssueStatusNode[] | null | undefined): Map<string, string> {
const snapshotByIssueId = new Map<string, { status: string; updatedAtMs: number | null }>();
for (const issue of issues ?? []) {
const candidate = {
status: issue.status,
updatedAtMs: issueUpdatedAtMs(issue.updatedAt),
};
const existing = snapshotByIssueId.get(issue.id);
if (!existing || shouldReplaceIssueStatusSnapshot(existing, candidate)) snapshotByIssueId.set(issue.id, candidate);
}
return new Map([...snapshotByIssueId].map(([issueId, snapshot]) => [issueId, snapshot.status]));
}
function issueUpdatedAtMs(updatedAt: LiveIssueStatusNode["updatedAt"]): number | null {
if (updatedAt === null || updatedAt === undefined) return null;
const timestamp = updatedAt instanceof Date ? updatedAt.getTime() : new Date(updatedAt).getTime();
return Number.isFinite(timestamp) ? timestamp : null;
}
function shouldReplaceIssueStatusSnapshot(
existing: { status: string; updatedAtMs: number | null },
candidate: { status: string; updatedAtMs: number | null },
): boolean {
if (candidate.updatedAtMs !== null && existing.updatedAtMs !== null) {
if (candidate.updatedAtMs !== existing.updatedAtMs) return candidate.updatedAtMs > existing.updatedAtMs;
} else if (candidate.updatedAtMs !== null) {
return true;
} else if (existing.updatedAtMs !== null) {
return false;
}
return !isTerminalIssueStatus(existing.status) && isTerminalIssueStatus(candidate.status);
}
export function collectLiveIssueIds(
liveRuns: readonly LiveRunForIssue[] | null | undefined,
issues?: readonly LiveIssueStatusNode[] | null,
): Set<string> {
const ids = new Set<string>();
const statusByIssueId = collectIssueStatusById(issues);
for (const run of liveRuns ?? []) {
if (run.issueId && isLiveRunStatus(run.status)) ids.add(run.issueId);
if (run.issueId && isLiveIssueRun(run, statusByIssueId.get(run.issueId))) ids.add(run.issueId);
}
return ids;
}

View File

@ -542,7 +542,7 @@ function ExecutionWorkspaceIssuesList({
});
usePublishSharedQueryData(sharedLiveRuns, liveRuns, liveRunsUpdatedAt);
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns), [liveRuns]);
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns, issues), [issues, liveRuns]);
const updateIssue = useMutation({
mutationFn: ({ id, data }: { id: string; data: Record<string, unknown> }) => issuesApi.update(id, data),

View File

@ -945,7 +945,6 @@ export function Inbox() {
refetchInterval: sharedLiveRuns.refetchInterval,
});
usePublishSharedQueryData(sharedLiveRuns, liveRuns, liveRunsUpdatedAt);
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns), [liveRuns]);
const { data: companyMembers } = useQuery({
queryKey: queryKeys.access.companyUserDirectory(selectedCompanyId!),
queryFn: () => accessApi.listUserDirectory(selectedCompanyId!),
@ -1001,6 +1000,15 @@ export function Inbox() {
enabled: shouldUseIssueSearchSupplement,
placeholderData: (previousData) => previousData,
});
const liveIssueIds = useMemo(
() => collectLiveIssueIds(liveRuns, [
...(issues ?? []),
...mineIssuesRaw,
...touchedIssuesRaw,
...remoteIssueSearchResults,
]),
[issues, liveRuns, mineIssuesRaw, remoteIssueSearchResults, touchedIssuesRaw],
);
const inboxIssueIdsForExternalObjectSummaries = useMemo(() => {
const issueIds = new Set<string>();
for (const issue of mineIssues) issueIds.add(issue.id);

View File

@ -1980,7 +1980,10 @@ export function IssueDetail() {
},
[issue?.id, rawChildIssues],
);
const liveIssueIds = useMemo(() => collectLiveIssueIds(companyLiveRuns), [companyLiveRuns]);
const liveIssueIds = useMemo(
() => collectLiveIssueIds(companyLiveRuns, issue ? [issue, ...childIssues] : childIssues),
[childIssues, companyLiveRuns, issue],
);
const issuePanelKey = useMemo(
() => buildIssuePropertiesPanelKey(issue ?? null, childIssues),
[childIssues, issue],

View File

@ -115,8 +115,6 @@ export function Issues() {
});
usePublishSharedQueryData(sharedLiveRuns, liveRuns, liveRunsUpdatedAt);
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns), [liveRuns]);
const issueLinkState = useMemo(
() =>
createIssueDetailLocationState(
@ -169,6 +167,7 @@ export function Issues() {
});
const issues = useMemo(() => mergeIssuePagesStable(issuePages?.pages ?? []) as Issue[], [issuePages]);
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns, issues), [issues, liveRuns]);
const hasMoreServerIssues = syncedSearch.trim().length === 0
&& hasNextPage === true;
const loadMoreServerIssues = useCallback(() => {

View File

@ -259,13 +259,12 @@ function ProjectIssuesList({ projectId, companyId }: { projectId: string; compan
enabled: !!companyId,
});
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns), [liveRuns]);
const { data: issues, isLoading, error } = useQuery({
queryKey: queryKeys.issues.listByProject(companyId, projectId),
queryFn: () => issuesApi.list(companyId, { projectId }),
enabled: !!companyId,
});
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns, issues), [issues, liveRuns]);
const updateIssue = useMutation({
mutationFn: ({ id, data }: { id: string; data: Record<string, unknown> }) =>
@ -330,13 +329,12 @@ function ProjectPluginOperationsList({
refetchInterval: sharedLiveRuns.refetchInterval,
});
usePublishSharedQueryData(sharedLiveRuns, liveRuns, liveRunsUpdatedAt);
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns), [liveRuns]);
const { data: issues, isLoading, error } = useQuery({
queryKey: queryKeys.issues.listPluginOperationsByProject(companyId, projectId, originKindPrefix),
queryFn: () => issuesApi.list(companyId, { projectId, originKindPrefix }),
enabled: !!companyId && !!projectId,
});
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns, issues), [issues, liveRuns]);
const updateIssue = useMutation({
mutationFn: ({ id, data }: { id: string; data: Record<string, unknown> }) =>

View File

@ -655,7 +655,7 @@ export function Routines() {
() => new Map((routineFolders?.folders ?? []).map((folder) => [folder.id, folder])),
[routineFolders],
);
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns), [liveRuns]);
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns, routineExecutionIssues), [liveRuns, routineExecutionIssues]);
const visibleRoutines = useMemo(
() => (routines ?? []).filter((routine) => routine.status !== "archived"),
[routines],

View File

@ -289,13 +289,12 @@ function PluginSdkIssuesList({
refetchInterval: sharedLiveRuns.refetchInterval,
});
usePublishSharedQueryData(sharedLiveRuns, liveRuns, liveRunsUpdatedAt);
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns), [liveRuns]);
const { data: issues, isLoading, error } = useQuery({
queryKey: issuesQueryKey,
queryFn: () => issuesApi.list(companyId!, issueFilters),
enabled: !!companyId,
});
const liveIssueIds = useMemo(() => collectLiveIssueIds(liveRuns, issues), [issues, liveRuns]);
const updateIssue = useMutation({
mutationFn: ({ id, data }: { id: string; data: Record<string, unknown> }) =>