Reduce run polling and workspace inspection amplification (#13174)

## Thinking Path

> - Paperclip manages agent work and shows run progress to operators.
> - Run lists, live events, transcripts, and workspace details must
remain responsive as usage grows.
> - Run-list redaction rereads the full context for every run. Hidden
tabs can still trigger requests through live events and manual timers.
> - Workspace detail reads repeat Git inspection even when concurrent
callers request the same state.
> - This pull request batches registry reads, pauses hidden-tab
refreshes, and caches Git inspection for display.
> - Cleanup keeps fresh Git checks, and redaction keeps company and run
boundaries.

## Linked Issues

**What happened?**
Run-list responses perform one extra database read per run and parse
full context JSON to obtain small secret registries. Hidden tabs
continue transcript reads and event-triggered refetches. Workspace
detail requests repeat Git scans.

**Expected behavior**
A run list reads registries once. Hidden tabs stop recurring run reads
and reconcile when visible. Concurrent workspace detail reads share a
short-lived Git result.

**Steps to reproduce**
1. Open run lists and task transcripts in several tabs while agents run.
2. Hide some tabs and observe transcript and event-triggered requests.
3. Request a 200-run list and count redaction database queries.
4. Request the same workspace detail concurrently and count Git
inspections.

Related: #5255 adjusts polling cadence. This change addresses hidden-tab
lifecycle, batched registry reads, and workspace inspection reuse. No
duplicate with this scope was found.

## What Changed

- Batch heartbeat and live-run redaction into one company-scoped
registry query. Select only registry JSON for run and issue redaction.
- Resolve duplicate secret values once per request. Preserve each run's
registry and remove registry material from responses.
- Suspend company event sockets and transcript reads while hidden.
Refresh active queries and resume transcript offsets on return.
- Prevent queued event invalidations and developer health polling from
fetching in hidden tabs. Gate legacy run-log readers in both UI
variants.
- Exclude legacy plugin placeholder connections from remote health
probes. Select only due connection IDs in SQL before the sweep limit.
Preserve existing plugin records.
- Cache concurrent Git display inspections for five seconds, with at
most 256 entries. Leave close-readiness and cleanup checks uncached.
- Add regression coverage and document the performance behavior.
- Stabilize the existing Rust descendant-lineage fixture: allow a
bounded 30 seconds for 300 durable notifications under concurrent test
load, retaining every correctness assertion and adding timeout
diagnostics.

## Verification

- Regression coverage verifies one registry query for 200 runs, per-run
isolation, request-local secret resolution, decryption failures, Git
cache expiry/bounds, hidden-tab pause, and visibility recovery.
- Real PostgreSQL redaction/run-route suites passed all 57 tests;
workspace-service coverage passed. The health-sweep regression verifies
plugin placeholders and chat connections remain untouched and do not
consume the sweep limit.
- Both legacy transcript viewers retain history and resume their byte
offset after visibility changes. The related visibility/progress/chunk
suites passed all 29 tests. Other focused UI suites and token gates
passed.
- Full `pnpm -r typecheck` and `pnpm build` passed. Affected-package
typechecks/builds passed after review fixes. The concurrent Rust
provider suite passed 84 tests (two ignored), and Rust formatting
passed.
- Full local `pnpm test:run` stopped after the general-server group:
10,538 passed, 65 skipped, four failed. Fresh chat-delivery and
health-sweep reruns passed; building the debug runner fixture cleared
the native-event test. One unchanged native-session recovery assertion
still fails locally with a semantic-digest error instead of the expected
settled-session message. The full local command is therefore not green.
CI runs the later groups separately and skips the two native-session
tests requiring a prebuilt runner binary (confirmed in its 37-test
native-session suite).
- All CI gates pass on final head `ee610e737`: typechecking, general and
serialized tests, browser tests, runner verification, build, and canary
dry run. One server shard passed on its single retry after exposure
fixtures encountered port 42001 where they assumed 42000; that suite
also passed locally (25 passed, three platform-specific skips).
- Greptile reviewed the final head at 5/5 with no actionable findings.

## Risks

- Workspace delivery display can lag local Git changes by five seconds.
Destructive operations still inspect current state.
- Hidden tabs do not receive company live-event notifications until
visible. Active queries refresh on return.
- This change preserves legacy plugin records and does not repair
instance-specific workspace rows. There is no database migration.

## Model Used

OpenAI Codex, GPT-6 family, with reasoning, repository tools, code
execution, and browser inspection. The exact model identifier and
context-window size are not exposed in this session.

## 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 (targeted regressions;
full-suite limitation documented above)
- [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:
Dotta 2026-09-11 08:34:24 -05:00 committed by GitHub
parent 932c8bec56
commit a05b828bcd
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
24 changed files with 576 additions and 85 deletions

View File

@ -1308,3 +1308,35 @@ Networking behavior for this smoke script:
### GitHub identity for shared agents
See [execution GitHub identity](execution-github-identity.md) for the operation-time credential contract, continuation rules, runtime rollout, and acceptance-test requirements.
### Investigating polling load
The company heartbeat-run and live-run lists load secret registries in one
company-scoped query per response. Registry reads project only
`paperclipSecretRedactions` from the run context. They do not load the full
prompt/context JSON. Decrypted values live only for that request and each run
uses its own registry.
Hidden browser tabs suspend the company live-events connection and transcript
log reads. Returning to a visible tab refreshes active queries once and resumes
transcript reads from their retained offsets. A queued live-event invalidation
that flushes after the tab hides marks data stale without starting a refetch.
The developer-server health poll also stops in hidden tabs.
Workspace detail responses share concurrent Git inspections and reuse their
results for up to five seconds after completion. The cache holds at most 256
entries. Close-readiness checks, the terminal-workspace reaper, and the final
cleanup validation still inspect Git afresh. A display result never authorizes
worktree removal.
The connection-health sweep selects only due IDs in SQL before applying its
limit. Legacy `paperclip_plugin` placeholder connections are excluded: their
tools run in plugin workers and do not have remote MCP endpoints. These rows
remain available; the sweep does not disable or delete plugin connections.
When investigating an overloaded instance, distinguish request amplification
from stored configuration problems. Verify connection transport and endpoint
fields before disabling a connection. Verify workspace ownership, active runs,
Git state, and runtime-service readiness before closing a workspace. A missing
URL or old workspace timestamp alone does not prove that a row is disposable.

View File

@ -5812,7 +5812,9 @@ fn durable_descendant_lineage_survives_capacity_and_provider_restoration() {
json!({"text": "Read test context."}),
))
.unwrap();
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
// Persisting 300 descendant notifications can exceed five seconds while
// the other provider tests contend for disk and CPU on a shared runner.
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(30);
let mut completed = false;
let mut children = std::collections::BTreeSet::new();
while std::time::Instant::now() < deadline && !completed {
@ -5830,7 +5832,11 @@ fn durable_descendant_lineage_survives_capacity_and_provider_restoration() {
}
std::thread::sleep(std::time::Duration::from_millis(1));
}
assert!(completed);
assert!(
completed,
"descendant run did not complete; observed {} of 300 children",
children.len()
);
assert_eq!(children.len(), 300);
first.shutdown().unwrap();
drop(first);

View File

@ -37,6 +37,7 @@ const mockInstanceSettingsService = vi.hoisted(() => ({
}));
const mockRunSecretRedactionRegistry = vi.hoisted(() => ({
redactForRuns: vi.fn(async (_companyId: string, values: unknown[]) => values),
redactForRun: vi.fn(
async (_companyId: string, _runId: string, value: unknown) => value,
),
@ -615,6 +616,8 @@ describe("agent live run routes", () => {
expect(res.status, JSON.stringify(res.body)).toBe(200);
expect(limit).toHaveBeenCalledWith(50);
expect(res.body).toHaveLength(50);
expect(mockRunSecretRedactionRegistry.redactForRuns).toHaveBeenCalledTimes(1);
expect(mockRunSecretRedactionRegistry.redactForRun).not.toHaveBeenCalled();
expect(mockHeartbeatService.buildRunOutputSilence).toHaveBeenCalledTimes(
50,
);
@ -659,6 +662,8 @@ describe("agent live run routes", () => {
expect(res.status, JSON.stringify(res.body)).toBe(200);
expect(limit).toHaveBeenCalledWith(50);
expect(res.body).toHaveLength(50);
expect(mockRunSecretRedactionRegistry.redactForRuns).toHaveBeenCalledTimes(1);
expect(mockRunSecretRedactionRegistry.redactForRun).not.toHaveBeenCalled();
});
it("does not pad with recent runs when no minCount is requested", async () => {

View File

@ -19,6 +19,7 @@ import {
secretAccessEvents,
} from "@paperclipai/db";
import { LOW_TRUST_REVIEW_PRESET, type AgentApiKeyScope } from "@paperclipai/shared";
import { REDACTED_EVENT_VALUE } from "../redaction.js";
import { errorHandler } from "../middleware/error-handler.js";
import { secretRoutes } from "../routes/secrets.js";
import { secretService } from "../services/secrets.js";
@ -248,6 +249,24 @@ describeEmbeddedPostgres("agent secret routes", () => {
expect((run.contextSnapshot as { paperclipSecretRedactions: unknown[] }).paperclipSecretRedactions).toHaveLength(1);
});
it("redacts batched runs from projected registries and enforces company scope", async () => {
const first = await seedAgentRun();
const foreign = await seedAgentRun();
const registry = createRunSecretRedactionRegistry(db);
await registry.register(first.companyId, first.heartbeatRunId, "first-secret-value");
await registry.register(foreign.companyId, foreign.heartbeatRunId, "foreign-secret-value");
const runs = [
{ id: first.heartbeatRunId, text: "first-secret-value foreign-secret-value", createdAt: new Date() },
{ id: foreign.heartbeatRunId, text: "foreign-secret-value", createdAt: new Date() },
];
const redacted = await registry.redactForRuns(first.companyId, runs);
expect(redacted[0].text).toBe(`${REDACTED_EVENT_VALUE} foreign-secret-value`);
expect(redacted[0].createdAt).toEqual(runs[0].createdAt);
expect(redacted[1].text).toBe("foreign-secret-value");
expect(await registry.redactForRun(first.companyId, first.heartbeatRunId, runs[0].text))
.toBe(redacted[0].text);
});
it("denies low-trust, task-bridge, and skill-test callers on both routes", async () => {
const lowTrust = await seedAgentRun({
trustPreset: LOW_TRUST_REVIEW_PRESET,

View File

@ -1,6 +1,8 @@
import { describe, expect, it } from "vitest";
import { beforeEach, describe, expect, it, vi } from "vitest";
import { REDACTED_EVENT_VALUE } from "../redaction.js";
import { redactRegisteredSecretValues } from "../services/run-secret-redaction.js";
import type { Db } from "@paperclipai/db";
import { PgDialect } from "drizzle-orm/pg-core";
import { createRunSecretRedactionRegistry, redactRegisteredSecretValues } from "../services/run-secret-redaction.js";
const secret = "q2a-exact-secret-value";
@ -76,3 +78,50 @@ describe("registered run secret redaction", () => {
expect(result.createdAt.toISOString()).toBe("2026-08-06T12:00:00.000Z");
});
});
const { resolveVersion } = vi.hoisted(() => ({ resolveVersion: vi.fn(async ({ material }) => material.value as string) }));
vi.mock("../secrets/provider-registry.js", () => ({ getSecretProvider: () => ({ resolveVersion }) }));
describe("batched run secret redaction", () => {
beforeEach(() => { resolveVersion.mockClear(); });
function fixture(rows: unknown[]) {
const where = vi.fn(async (_predicate: import("drizzle-orm").SQL | undefined) => rows);
const select = vi.fn((_columns: { contextSnapshot: import("drizzle-orm").SQL }) => ({ from: () => ({ where }) }));
return { registry: createRunSecretRedactionRegistry({ select } as unknown as Db), select, where };
}
it("reads only registry JSON once for 200 runs and resolves shared secrets once", async () => {
const contextSnapshot = { paperclipSecretRedactions: [{ fingerprintSha256: "shared", material: { value: secret } }] };
const rows = Array.from({ length: 200 }, (_, i) => ({ id: `run-${i}`, contextSnapshot }));
const { registry, select, where } = fixture(rows);
const result = await registry.redactForRuns("company-1", rows.map(row => ({ ...row, stdoutExcerpt: secret })));
expect(select).toHaveBeenCalledTimes(1);
expect(resolveVersion).toHaveBeenCalledTimes(1);
expect(result.every(run => run.stdoutExcerpt === REDACTED_EVENT_VALUE)).toBe(true);
expect(result[0].contextSnapshot).toEqual({});
const dialect = new PgDialect();
const predicate = dialect.sqlToQuery(where.mock.calls[0][0]);
expect(predicate.params).toContain("company-1");
expect(predicate.sql).toContain('"company_id"');
expect(dialect.sqlToQuery(select.mock.calls[0][0].contextSnapshot).sql).toContain("-> 'paperclipSecretRedactions'");
});
it("keeps each run's registry separate and observes new registrations on the next request", async () => {
const rows = [{ id: "a", contextSnapshot: { paperclipSecretRedactions: [{ fingerprintSha256: "one", material: { value: secret } }] } }];
const { registry } = fixture(rows);
expect(await registry.redactForRuns("company", [{ id: "a", text: secret }, { id: "b", text: secret }]))
.toEqual([{ id: "a", text: REDACTED_EVENT_VALUE }, { id: "b", text: secret }]);
rows[0].contextSnapshot.paperclipSecretRedactions.push({ fingerprintSha256: "two", material: { value: "new-secret" } });
expect(await registry.redactForRuns("company", [{ id: "a", text: "new-secret" }]))
.toEqual([{ id: "a", text: REDACTED_EVENT_VALUE }]);
});
it("does not query for an empty list and fails closed on decryption failure", async () => {
const { registry, select } = fixture([{ id: "a", contextSnapshot: { paperclipSecretRedactions: [{ fingerprintSha256: "one", material: {} }] } }]);
expect(await registry.redactForRuns("company", [])).toEqual([]);
expect(select).not.toHaveBeenCalled();
resolveVersion.mockRejectedValueOnce(new Error("unavailable"));
await expect(registry.redactForRuns("company", [{ id: "a", text: secret }])).rejects.toThrow("unavailable");
});
});

View File

@ -17094,6 +17094,28 @@ describeEmbeddedPostgres("tool access service", () => {
lastHealthAt: new Date(0),
})
.returning();
const [pluginApplication] = await db.insert(toolApplications).values({
companyId: company.id,
applicationKey: `paperclip_plugin:fixture-${randomUUID()}`,
name: "Plugin placeholder",
type: "paperclip_plugin",
status: "active",
metadata: { source: "plugin_backfill" },
}).returning();
const [pluginConnection] = await db.insert(toolConnections).values({
companyId: company.id,
applicationId: pluginApplication!.id,
name: "Plugin placeholder",
uid: `plugin-${randomUUID()}`,
connectionKind: "managed",
transport: "mcp_remote",
status: "active",
enabled: true,
config: { type: "paperclip_plugin" },
transportConfig: { type: "paperclip_plugin" },
healthStatus: "ok",
healthCheckedAt: null,
}).returning();
const connection = await service.createConnection(company.id, {
name: "Swept remote",
transport: "mcp_remote",
@ -17102,7 +17124,7 @@ describeEmbeddedPostgres("tool access service", () => {
status: "active",
});
const sweep = await service.sweepConnectionHealth({ staleAfterMs: 0 });
const sweep = await service.sweepConnectionHealth({ staleAfterMs: 0, limit: 1 });
const [updatedConnection] = await db
.select()
.from(toolConnections)
@ -17112,6 +17134,10 @@ describeEmbeddedPostgres("tool access service", () => {
.from(toolConnections)
.where(eq(toolConnections.id, chatConnection!.id));
const [untouchedPlugin] = await db.select().from(toolConnections)
.where(eq(toolConnections.id, pluginConnection!.id));
expect(untouchedPlugin).toMatchObject({ enabled: true, healthStatus: "ok", healthCheckedAt: null });
expect(sweep).toMatchObject({
checked: 1,
healthy: 0,

View File

@ -0,0 +1,44 @@
import type { ExecutionWorkspace } from "@paperclipai/shared";
import { afterEach, expect, it, vi } from "vitest";
import { createWorkspaceGitInspectionCache } from "../services/workspace-git-inspection-cache.js";
const workspace = { id: "workspace", companyId: "company", cwd: "/repo", baseRef: "master" } as ExecutionWorkspace;
afterEach(() => vi.useRealTimers());
it("coalesces concurrent display reads and expires after five seconds", async () => {
vi.useFakeTimers();
const inspect = vi.fn(async () => ({ dirty: false }));
const read = createWorkspaceGitInspectionCache(inspect);
await Promise.all(Array.from({ length: 100 }, () => read(workspace)));
expect(inspect).toHaveBeenCalledTimes(1);
await read(workspace);
expect(inspect).toHaveBeenCalledTimes(1);
vi.advanceTimersByTime(5_000);
await read(workspace);
expect(inspect).toHaveBeenCalledTimes(2);
// Callers that authorize cleanup retain the uncached inspector.
await inspect();
expect(inspect).toHaveBeenCalledTimes(3);
});
it("does not share results across companies, paths, base refs or workspace revisions", async () => {
const inspect = vi.fn(async () => null);
const read = createWorkspaceGitInspectionCache(inspect);
await read(workspace);
await read({ ...workspace, companyId: "other" });
await read({ ...workspace, cwd: "/other" });
await read({ ...workspace, baseRef: "other" });
await read({ ...workspace, updatedAt: new Date() });
expect(inspect).toHaveBeenCalledTimes(5);
});
it("retries failed inspections and bounds retained entries", async () => {
const inspect = vi.fn(async () => null).mockRejectedValueOnce(new Error("failed"));
const read = createWorkspaceGitInspectionCache(inspect);
await expect(read(workspace)).rejects.toThrow("failed");
await read(workspace);
expect(inspect).toHaveBeenCalledTimes(2);
for (let i = 0; i < 256; i++) await read({ ...workspace, id: String(i) });
await read(workspace);
expect(inspect).toHaveBeenCalledTimes(259);
});

View File

@ -6219,7 +6219,7 @@ export function agentRoutes(
const limit = limitParam ? Math.max(1, Math.min(1000, parseInt(limitParam, 10) || 200)) : undefined;
const summary = req.query.summary === "true" || req.query.summary === "1";
const runs = await heartbeat.list(companyId, agentId, limit, { summary });
res.json(await Promise.all(runs.map((run) => runRedactions.redactForRun(companyId, run.id, run))));
res.json(await runRedactions.redactForRuns(companyId, runs));
});
router.get("/companies/:companyId/provider-traces", async (req, res) => {
@ -6326,20 +6326,20 @@ export function agentRoutes(
const rows = [...liveRuns, ...recentRuns];
const projections = await executionProjectionsForRuns(db, companyId, rows.map(run => run.id));
res.json(await Promise.all(rows.map(async (run) => runRedactions.redactForRun(companyId, run.id, {
res.json(await runRedactions.redactForRuns(companyId, await Promise.all(rows.map(async (run) => ({
...heartbeat.decorateActiveRunStatus(run),
execution: projections.get(run.id) ?? null,
outputSilence: await heartbeat.buildRunOutputSilence(run),
}))));
})))));
return;
}
const projections = await executionProjectionsForRuns(db, companyId, liveRuns.map(run => run.id));
res.json(await Promise.all(liveRuns.map(async (run) => runRedactions.redactForRun(companyId, run.id, {
res.json(await runRedactions.redactForRuns(companyId, await Promise.all(liveRuns.map(async (run) => ({
...heartbeat.decorateActiveRunStatus(run),
execution: projections.get(run.id) ?? null,
outputSilence: await heartbeat.buildRunOutputSilence(run),
}))));
})))));
});
router.get("/heartbeat-runs/:runId", async (req, res) => {

View File

@ -1,3 +1,4 @@
import { createWorkspaceGitInspectionCache } from "./workspace-git-inspection-cache.js";
import { execFile } from "node:child_process";
import { createHash } from "node:crypto";
import fs from "node:fs/promises";
@ -1268,7 +1269,12 @@ type WorkspaceOverviewIssueRow = WorkspaceOverviewLinkedIssue & {
executionWorkspaceId: string;
};
const inspectGitForDisplay = createWorkspaceGitInspectionCache(inspectGitCloseReadiness);
export function executionWorkspaceService(db: Db, opts: ExecutionWorkspaceServiceOptions = {}) {
const inspectDisplay = opts.inspectGitCloseReadiness
? createWorkspaceGitInspectionCache(opts.inspectGitCloseReadiness)
: inspectGitForDisplay;
const recoveryActionsSvc = issueRecoveryActionService(db);
const resolvePullRequestDetails = opts.resolvePullRequestDetails ?? createPullRequestMergeDetailsResolver(db);
const now = opts.now ?? (() => new Date());
@ -1487,7 +1493,7 @@ export function executionWorkspaceService(db: Db, opts: ExecutionWorkspaceServic
async function hydrateWorkspace(row: ExecutionWorkspaceRow, runtimeServices: WorkspaceRuntimeService[] = []) {
const workspace = toExecutionWorkspace(row, runtimeServices);
const { git } = await (opts.inspectGitCloseReadiness ?? inspectGitCloseReadiness)(workspace);
const { git } = await inspectDisplay(workspace);
const assessment = await assessDelivery(row, git);
return toExecutionWorkspace(row, runtimeServices, assessment.deliveryState);
}

View File

@ -1,5 +1,5 @@
import { createHash } from "node:crypto";
import { and, eq, or, sql } from "drizzle-orm";
import { and, eq, inArray, or, sql } from "drizzle-orm";
import type { Db } from "@paperclipai/db";
import { heartbeatRuns } from "@paperclipai/db";
import { REDACTED_EVENT_VALUE } from "../redaction.js";
@ -7,6 +7,8 @@ import { getSecretProvider } from "../secrets/provider-registry.js";
import type { StoredSecretVersionMaterial } from "../secrets/types.js";
const REGISTRY_KEY = "paperclipSecretRedactions";
// Project only the registry: run contexts can contain megabytes of prompt data.
const registrySnapshot = sql`jsonb_build_object('paperclipSecretRedactions', ${heartbeatRuns.contextSnapshot} -> 'paperclipSecretRedactions')`;
type RegistryEntry = {
fingerprintSha256: string;
@ -70,14 +72,14 @@ export function createRunSecretRedactionRegistry(db: Db) {
}
async function valuesForRun(companyId: string, runId: string) {
const rows = await db.select({ contextSnapshot: heartbeatRuns.contextSnapshot })
const rows = await db.select({ contextSnapshot: registrySnapshot })
.from(heartbeatRuns)
.where(and(eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.id, runId)));
return valuesForRuns(rows);
}
async function valuesForIssue(companyId: string, issueId: string) {
const rows = await db.select({ contextSnapshot: heartbeatRuns.contextSnapshot })
const rows = await db.select({ contextSnapshot: registrySnapshot })
.from(heartbeatRuns)
.where(and(
eq(heartbeatRuns.companyId, companyId),
@ -116,6 +118,27 @@ export function createRunSecretRedactionRegistry(db: Db) {
.where(and(eq(heartbeatRuns.companyId, companyId), eq(heartbeatRuns.id, runId)));
});
},
redactForRuns: async <T extends { id: string }>(companyId: string, runs: T[]): Promise<T[]> => {
if (runs.length === 0) return [];
const rows = await db.select({ id: heartbeatRuns.id, contextSnapshot: registrySnapshot })
.from(heartbeatRuns)
.where(and(eq(heartbeatRuns.companyId, companyId), inArray(heartbeatRuns.id, runs.map((run) => run.id))));
// Resolve each encrypted value once per request, but apply only each run's
// own registry. Do not retain plaintext secrets across requests.
const resolved = new Map<string, Promise<string>>();
const valuesByRun = new Map(await Promise.all(rows.map(async (row) => {
const values = await Promise.all(registryEntries(row.contextSnapshot).map((entry) => {
let value = resolved.get(entry.fingerprintSha256);
if (!value) {
value = provider.resolveVersion({ material: entry.material, externalRef: null });
resolved.set(entry.fingerprintSha256, value);
}
return value;
}));
return [row.id, values.sort((a, b) => b.length - a.length)] as const;
})));
return runs.map((run) => redactRegisteredSecretValues(run, valuesByRun.get(run.id) ?? []));
},
redactForRun: async <T>(companyId: string, runId: string, value: T): Promise<T> =>
redactRegisteredSecretValues(value, await valuesForRun(companyId, runId)),
redactForIssue: async <T>(companyId: string, issueId: string, value: T): Promise<T> =>

View File

@ -15,8 +15,10 @@ import {
isNotNull,
isNull,
lt,
lte,
max,
ne,
or,
sql,
} from "drizzle-orm";
import type { Db } from "@paperclipai/db";
@ -8015,26 +8017,32 @@ export function toolAccessService(
const staleAfterMs = input.staleAfterMs ?? 15 * 60 * 1000;
const limit = input.limit ?? 25;
const cutoff = new Date(generatedAt.getTime() - staleAfterMs);
const connections = await db
.select()
// Legacy plugin backfills use a remote transport as a placeholder, but
// their tools run in the plugin worker and have no remote MCP endpoint.
// Select only due IDs in SQL so each scheduler tick does not decode every
// active connection's config and credential metadata.
const due = await db
.select({ id: toolConnections.id })
.from(toolConnections)
.innerJoin(toolApplications, and(
eq(toolApplications.id, toolConnections.applicationId),
eq(toolApplications.companyId, toolConnections.companyId),
))
.where(
and(
eq(toolConnections.enabled, true),
eq(toolConnections.status, "active"),
ne(toolConnections.transport, "chat_sdk"),
ne(toolApplications.type, "paperclip_plugin"),
or(isNull(toolConnections.healthCheckedAt), lte(toolConnections.healthCheckedAt, cutoff)),
),
)
.orderBy(
asc(toolConnections.healthCheckedAt),
sql`${toolConnections.healthCheckedAt} asc nulls first`,
asc(toolConnections.createdAt),
);
const due = connections
.filter(
(connection) =>
!connection.healthCheckedAt || connection.healthCheckedAt <= cutoff,
asc(toolConnections.id),
)
.slice(0, limit);
.limit(limit);
let healthy = 0;
let failed = 0;
const failedConnectionIds: string[] = [];

View File

@ -0,0 +1,30 @@
import type { ExecutionWorkspace } from "@paperclipai/shared";
/** Short-lived display cache only. Destructive operations must inspect afresh. */
export function createWorkspaceGitInspectionCache<T>(inspect: (workspace: ExecutionWorkspace) => Promise<T>) {
const entries = new Map<string, { expiresAt: number; promise: Promise<T> }>();
return (workspace: ExecutionWorkspace): Promise<T> => {
const key = JSON.stringify([
workspace.companyId, workspace.id, workspace.updatedAt, workspace.providerType,
workspace.providerRef, workspace.cwd, workspace.repoUrl, workspace.baseRef,
workspace.branchName, workspace.metadata,
]);
const now = Date.now();
const existing = entries.get(key);
if (existing && existing.expiresAt > now) return existing.promise;
for (const [candidate, entry] of entries) {
if (entry.expiresAt <= now) entries.delete(candidate);
}
if (entries.size >= 256) entries.delete(entries.keys().next().value!);
const entry = { expiresAt: Number.POSITIVE_INFINITY, promise: Promise.resolve().then(() => inspect(workspace)) };
entries.set(key, entry);
entry.promise = entry.promise.then((result) => {
entry.expiresAt = Date.now() + 5_000;
return result;
}, (error) => {
if (entries.get(key) === entry) entries.delete(key);
throw error;
});
return entry.promise;
};
}

View File

@ -248,7 +248,7 @@ export function Layout() {
{ devServer?: { enabled?: boolean } } | undefined;
return data?.devServer?.enabled ? 2000 : false;
},
refetchIntervalInBackground: true,
refetchIntervalInBackground: false,
});
const keyboardShortcutsEnabled =
useQuery({

View File

@ -235,7 +235,7 @@ export function Layout() {
const data = query.state.data as { devServer?: { enabled?: boolean } } | undefined;
return data?.devServer?.enabled ? 2000 : false;
},
refetchIntervalInBackground: true,
refetchIntervalInBackground: false,
});
const keyboardShortcutsEnabled = useQuery({
queryKey: queryKeys.instance.generalSettings,

View File

@ -78,6 +78,7 @@ describe("useLiveRunTranscripts", () => {
const OriginalWebSocket = globalThis.WebSocket;
beforeEach(() => {
vi.spyOn(document, "visibilityState", "get").mockReturnValue("visible");
FakeWebSocket.instances = [];
useQueryMock.mockClear();
logMock.mockReset();
@ -88,6 +89,45 @@ describe("useLiveRunTranscripts", () => {
afterEach(() => {
globalThis.WebSocket = OriginalWebSocket;
vi.restoreAllMocks();
vi.useRealTimers();
});
it("pauses hidden-tab reads and resumes at the retained log offset", async () => {
vi.useFakeTimers();
const visibility = vi.spyOn(document, "visibilityState", "get").mockReturnValue("hidden");
logMock.mockResolvedValue({ runId: "run-1", store: "memory", logRef: "log-1", content: "", nextOffset: 42 });
const runs = [{ id: "run-1", status: "running", adapterType: "codex_local" }];
function Harness() {
useLiveRunTranscripts({ companyId: "company-1", runs, enableRealtimeUpdates: false });
return null;
}
const container = document.createElement("div");
const root = createRoot(container);
try {
await act(async () => root.render(<Harness />));
await act(async () => vi.advanceTimersByTimeAsync(10_000));
expect(logMock).not.toHaveBeenCalled();
expect(FakeWebSocket.instances).toHaveLength(0);
await act(async () => {
visibility.mockReturnValue("visible");
document.dispatchEvent(new Event("visibilitychange"));
});
expect(logMock).toHaveBeenCalledTimes(1);
await act(async () => {
visibility.mockReturnValue("hidden");
document.dispatchEvent(new Event("visibilitychange"));
});
await act(async () => vi.advanceTimersByTimeAsync(10_000));
expect(logMock).toHaveBeenCalledTimes(1);
await act(async () => {
visibility.mockReturnValue("visible");
document.dispatchEvent(new Event("visibilitychange"));
});
expect(logMock).toHaveBeenLastCalledWith("run-1", 42, 256_000, expect.anything());
} finally {
await act(async () => root.unmount());
}
});
it("waits for a connecting socket to open before closing it during cleanup", async () => {

View File

@ -1,3 +1,4 @@
import { usePageVisibility } from "../../lib/page-visibility";
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { readTranscriptRequest } from "./read-transcript-request";
import { useQuery } from "@tanstack/react-query";
@ -102,6 +103,7 @@ export function useLiveRunTranscripts({
}: UseLiveRunTranscriptsOptions) {
// Ticker consumers opt into the silent chunk-count cap; full task views use a
// byte budget that collapses (not discards) the oldest output when exceeded.
const { visible } = usePageVisibility();
const retentionBudget: ChunkRetentionBudget = useMemo(
() =>
typeof maxChunksPerRun === "number"
@ -293,6 +295,7 @@ export function useLiveRunTranscripts({
}, [normalizedRuns, pruneTick]);
useEffect(() => {
if (!visible) return;
const readableRuns = normalizedRuns.filter(canReadPersistedLog);
if (readableRuns.length === 0) return;
@ -383,10 +386,10 @@ export function useLiveRunTranscripts({
controller.abort();
if (interval !== null) window.clearInterval(interval);
};
}, [enableRealtimeUpdates, logPollIntervalMs, logReadLimitBytes, normalizedRuns, runIdsKey, retryGeneration]);
}, [visible, enableRealtimeUpdates, logPollIntervalMs, logReadLimitBytes, normalizedRuns, runIdsKey, retryGeneration]);
useEffect(() => {
if (!enableRealtimeUpdates) return;
if (!visible || !enableRealtimeUpdates) return;
if (!companyId || activeRunIds.size === 0) return;
let closed = false;
@ -515,7 +518,7 @@ export function useLiveRunTranscripts({
}
}
};
}, [activeRunIds, companyId, enableRealtimeUpdates, runById]);
}, [visible, activeRunIds, companyId, enableRealtimeUpdates, runById]);
const transcriptByRun = useMemo(() => {
const next = new Map<string, TranscriptEntry[]>();

View File

@ -1,3 +1,4 @@
import { usePageVisibility } from "@/lib/page-visibility";
import { useCallback, useEffect, useMemo, useRef, useState } from "react";
import { useQuery } from "@tanstack/react-query";
import type { LiveEvent, SummarySlotIssueRef } from "@paperclipai/shared";
@ -62,6 +63,7 @@ export function useSummaryDraftStream(
companyId: string | null | undefined,
generatingIssue: SummarySlotIssueRef | null,
): SummaryDraftStream {
const { visible } = usePageVisibility();
const issueId = generatingIssue?.id ?? null;
const [runId, setRunId] = useState<string | null>(null);
const [chunks, setChunks] = useState<RunLogChunk[]>([]);
@ -118,9 +120,14 @@ export function useSummaryDraftStream(
if (fallbackRunId) setRunId((current) => current ?? fallbackRunId);
}, [fallbackRunId]);
useEffect(() => {
logOffsetRef.current = 0;
pendingLogRowsRef.current = new Map();
}, [runId]);
// Live token deltas over the shared company-events socket.
useCompanyLiveEvent((event: LiveEvent) => {
if (!runId) return;
if (!visible || !runId) return;
if (event.type !== "heartbeat.run.log") return;
const payload = event.payload ?? {};
if (payload.runId !== runId) return;
@ -136,12 +143,12 @@ export function useSummaryDraftStream(
// Hydrate already-emitted output and fill any gaps from the persisted run log.
useEffect(() => {
if (!runId) return;
logOffsetRef.current = 0;
pendingLogRowsRef.current = new Map();
if (!visible || !runId) return;
let cancelled = false;
let reading = false;
const read = async () => {
if (reading || cancelled) return;
reading = true;
try {
const result = await heartbeatsApi.log(runId, logOffsetRef.current, LOG_READ_LIMIT_BYTES);
if (cancelled) return;
@ -153,6 +160,8 @@ export function useSummaryDraftStream(
}
} catch {
// Ignore transient/404 reads (log not yet flushed, run just started).
} finally {
reading = false;
}
};
@ -162,7 +171,7 @@ export function useSummaryDraftStream(
cancelled = true;
window.clearInterval(interval);
};
}, [runId, appendChunks]);
}, [visible, runId, appendChunks]);
const parse = useMemo(() => parseSummaryDraftStream(extractAssistantOutputText(chunks)), [chunks]);

View File

@ -181,6 +181,24 @@ describe("LiveUpdatesProvider socket run notification scope", () => {
})));
}
it("disconnects while hidden and reconciles active queries once on return", async () => {
await receiveStatus({ runId: "child-run", agentId: "child-agent", status: "running" });
const invalidate = vi.spyOn(queryClient, "invalidateQueries");
const visibility = vi.spyOn(document, "visibilityState", "get");
await reactAct(async () => {
visibility.mockReturnValue("hidden");
document.dispatchEvent(new Event("visibilitychange"));
});
expect(sockets[0].onmessage).toBeNull();
invalidate.mockClear();
await reactAct(async () => {
visibility.mockReturnValue("visible");
document.dispatchEvent(new Event("visibilitychange"));
});
await vi.waitFor(() => expect(sockets).toHaveLength(2));
expect(invalidate).toHaveBeenCalledExactlyOnceWith({ type: "active" }, { cancelRefetch: false });
});
it.each(["parent-agent", "child-agent"])("shows an unrelated retryable failure without issueId for %s", async (agentId) => {
// Match the retryable broadcast from execution-status-delivery.ts: it has
// exact run identity but deliberately omits issueId and provider output.

View File

@ -1,3 +1,4 @@
import { getPageVisibility, usePageVisibility } from "../lib/page-visibility";
import {
createContext,
useCallback,
@ -1787,6 +1788,8 @@ export const __liveUpdatesTestUtils = {
};
export function LiveUpdatesProvider({ children }: { children: ReactNode }) {
const { visible } = usePageVisibility();
const wasHidden = useRef(!visible);
const { selectedCompanyId, selectedCompany } = useCompany();
const queryClient = useQueryClient();
const { pushToast } = useToastActions();
@ -1853,7 +1856,17 @@ export function LiveUpdatesProvider({ children }: { children: ReactNode }) {
}, [currentUserId]);
useEffect(() => {
if (!visible) {
wasHidden.current = true;
invalidationBatcher.dispose();
return;
}
if (!canConnectSocket || !liveCompanyId) return;
if (wasHidden.current) {
wasHidden.current = false;
// Reconcile events missed while hidden, including completed runs/issues.
void queryClient.invalidateQueries({ type: "active" }, { cancelRefetch: false });
}
let closed = false;
let reconnectAttempt = 0;
@ -1905,6 +1918,7 @@ export function LiveUpdatesProvider({ children }: { children: ReactNode }) {
};
nextSocket.onmessage = (message) => {
if (!getPageVisibility().visible) return;
const raw = typeof message.data === "string" ? message.data : "";
if (!raw) return;
@ -1961,6 +1975,9 @@ export function LiveUpdatesProvider({ children }: { children: ReactNode }) {
closeSocketQuietly(activeSocket, "provider_unmount");
};
}, [
visible,
invalidationBatcher,
queryClient,
coalescingClient,
liveCompanyId,
pushToast,

View File

@ -18,7 +18,17 @@ function fakeClient() {
describe("createInvalidationBatcher", () => {
beforeEach(() => vi.useFakeTimers());
afterEach(() => vi.useRealTimers());
afterEach(() => { vi.useRealTimers(); vi.unstubAllGlobals(); });
it("marks queries stale without refetching when the tab hides before flush", async () => {
const { client } = fakeClient();
const batcher = createInvalidationBatcher(client);
const pending = batcher.schedule({ queryKey: ["dashboard", "c1"] });
vi.stubGlobal("document", { visibilityState: "hidden" });
await batcher.flush();
await pending;
expect(client.invalidateQueries).toHaveBeenCalledExactlyOnceWith({ queryKey: ["dashboard", "c1"], refetchType: "none" });
});
it("coalesces repeated invalidations of the same key into one call per window", () => {
const { client } = fakeClient();
@ -111,7 +121,7 @@ describe("createInvalidationBatcher", () => {
describe("createCoalescingQueryClient", () => {
beforeEach(() => vi.useFakeTimers());
afterEach(() => vi.useRealTimers());
afterEach(() => { vi.useRealTimers(); vi.unstubAllGlobals(); });
it("batches invalidateQueries but passes other methods straight through", () => {
const setQueryData = vi.fn();

View File

@ -1,3 +1,4 @@
import { getPageVisibility } from "./page-visibility";
import type { InvalidateQueryFilters, QueryClient } from "@tanstack/react-query";
/**
@ -62,7 +63,9 @@ export function createInvalidationBatcher(
const filtersList = [...pending.values()];
pending.clear();
try {
await Promise.all(filtersList.map((filters) => queryClient.invalidateQueries(filters)));
await Promise.all(filtersList.map((filters) => queryClient.invalidateQueries(
getPageVisibility().visible ? filters : { ...filters, refetchType: "none" },
)));
} finally {
deferred?.resolve();
}

View File

@ -0,0 +1,63 @@
// @vitest-environment jsdom
import { act } from "react";
import { createRoot } from "react-dom/client";
import type { HeartbeatRun } from "@paperclipai/shared";
import { afterEach, expect, it, vi } from "vitest";
import { LogViewer } from "./AgentDetail";
import { LogViewer as ProductionLogViewer } from "./AgentDetail.production";
const { log, empty } = vi.hoisted(() => ({ log: vi.fn(), empty: [] }));
vi.mock("../api/heartbeats", () => ({ heartbeatsApi: { log } }));
vi.mock("@tanstack/react-query", async (original) => ({
...await original<typeof import("@tanstack/react-query")>(),
useQuery: () => ({ data: empty }),
}));
vi.mock("../adapters", () => ({
getUIAdapter: () => null,
onAdapterChange: () => () => {},
buildTranscript: (lines: unknown[]) => lines,
}));
vi.mock("../components/transcript/RunTranscriptView", () => ({
RunTranscriptView: ({ entries }: { entries: Array<{ chunk: string }> }) => <div>{entries.map(line => line.chunk).join(" ")}</div>,
}));
(globalThis as typeof globalThis & { IS_REACT_ACT_ENVIRONMENT?: boolean }).IS_REACT_ACT_ENVIRONMENT = true;
afterEach(() => { vi.restoreAllMocks(); log.mockReset(); });
it.each([LogViewer, ProductionLogViewer])("retains legacy history and reads only the next offset on visibility recovery (%#)", async (Viewer) => {
const visibility = vi.spyOn(document, "visibilityState", "get").mockReturnValue("visible");
const row = (seq: number, chunk: string) => JSON.stringify({ seq, ts: `2026-09-10T12:00:0${seq}Z`, stream: "stdout", chunk }) + "\n";
const first = row(1, "retained history");
const second = row(2, "new output");
log.mockResolvedValueOnce({ content: first, nextOffset: first.length });
const run = { id: "run-1", companyId: "company-1", agentId: "agent-1", status: "succeeded", logRef: "log" } as HeartbeatRun;
const container = document.createElement("div");
document.body.append(container);
const root = createRoot(container);
try {
await act(async () => root.render(<Viewer run={run} adapterType="codex_local" />));
expect(container.textContent).toContain("retained history");
await act(async () => {
visibility.mockReturnValue("hidden");
document.dispatchEvent(new Event("visibilitychange"));
});
expect(log).toHaveBeenCalledTimes(1);
expect(container.textContent).toContain("retained history");
let complete!: (value: { content: string; nextOffset: number }) => void;
log.mockImplementationOnce(() => new Promise(resolve => { complete = resolve; }));
await act(async () => {
visibility.mockReturnValue("visible");
document.dispatchEvent(new Event("visibilitychange"));
});
expect(container.textContent).toContain("retained history");
expect(log).toHaveBeenLastCalledWith("run-1", first.length, expect.any(Number));
await act(async () => complete({ content: second, nextOffset: first.length + second.length }));
expect(container.textContent).toContain("retained history new output");
log.mockResolvedValueOnce({ content: first, nextOffset: first.length });
await act(async () => root.render(<Viewer run={{ ...run, id: "run-2" }} adapterType="codex_local" />));
expect(log).toHaveBeenLastCalledWith("run-2", 0, expect.any(Number));
expect(container.textContent).not.toContain("new output");
} finally {
await act(async () => root.unmount());
container.remove();
}
});

View File

@ -1,3 +1,5 @@
import { mergeRunLogChunks, readChunkSeq } from "../lib/run-log-chunks";
import { getPageVisibility, usePageVisibility } from "../lib/page-visibility";
import { useCallback, useEffect, useMemo, useState, useRef } from "react";
import { useParams, useNavigate, Link, Navigate, useBeforeUnload, type NavigateFunction } from "@/lib/router";
import { useQuery, useMutation, useQueryClient, type QueryClient } from "@tanstack/react-query";
@ -383,6 +385,7 @@ function runMetrics(run: HeartbeatRun) {
}
export type RunLogChunk = {
seq?: number;
ts: string;
stream: "stdout" | "stderr" | "system";
chunk: string;
@ -3720,13 +3723,20 @@ function RunDetail({ run: initialRun, agentRouteId, adapterType, adapterConfig }
/* ---- Log Viewer ---- */
function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: string }) {
export function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: string }) {
const { visible } = usePageVisibility();
const [events, setEvents] = useState<HeartbeatRunEvent[]>([]);
const [logLines, setLogLines] = useState<Array<{ ts: string; stream: "stdout" | "stderr" | "system"; chunk: string }>>([]);
const [logLines, setLogLines] = useState<RunLogChunk[]>([]);
const [loading, setLoading] = useState(true);
const [logLoading, setLogLoading] = useState(!!run.logRef);
const [logError, setLogError] = useState<string | null>(null);
const [logOffset, setLogOffset] = useState(0);
const [logOffset, setLogOffsetState] = useState(0);
const logOffsetRef = useRef(0);
const setLogOffset = useCallback((next: number | ((previous: number) => number)) => {
logOffsetRef.current = typeof next === "function" ? next(logOffsetRef.current) : next;
setLogOffsetState(logOffsetRef.current);
}, []);
const logMergeRefs = useRef({ seenChunkKeys: new Set<string>(), trimmedSeqFloorByRun: new Map<string, number>() });
const [hasMoreLog, setHasMoreLog] = useState(false);
const [loadingMoreLog, setLoadingMoreLog] = useState(false);
const [isFollowing, setIsFollowing] = useState(false);
@ -3753,6 +3763,12 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
return err instanceof ApiError && err.status === 404;
}
function appendLogLines(incoming: RunLogChunk[]) {
setLogLines((previous) => mergeRunLogChunks(run.id, previous, incoming.map((line) => ({
...line, dedupeKey: `log:${run.id}:${line.ts}:${line.stream}:${line.chunk}`,
})), logMergeRefs.current, isLive ? MAX_LIVE_LOG_LINES : Number.POSITIVE_INFINITY).chunks);
}
function appendLogContent(content: string, finalize = false) {
if (!content && !finalize) return;
const combined = `${pendingLogLineRef.current}${content}`;
@ -3763,18 +3779,18 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
pendingLogLineRef.current = "";
}
const parsed: Array<{ ts: string; stream: "stdout" | "stderr" | "system"; chunk: string }> = [];
const parsed: RunLogChunk[] = [];
for (const line of split) {
const trimmed = line.trim();
if (!trimmed) continue;
try {
const raw = JSON.parse(trimmed) as { ts?: unknown; stream?: unknown; chunk?: unknown };
const raw = JSON.parse(trimmed) as { ts?: unknown; stream?: unknown; chunk?: unknown; seq?: unknown };
const stream =
raw.stream === "stderr" || raw.stream === "system" ? raw.stream : "stdout";
const chunk = typeof raw.chunk === "string" ? raw.chunk : "";
const ts = typeof raw.ts === "string" ? raw.ts : new Date().toISOString();
if (!chunk) continue;
parsed.push({ ts, stream, chunk });
parsed.push({ ts, stream, chunk, seq: readChunkSeq(raw.seq) });
} catch {
// ignore malformed lines
}
@ -3783,9 +3799,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
if (parsed.length > 0) {
// Live runs stream forever, so cap the retained tail. Terminated runs are
// paginated by the user via "Load more log" and keep their full history.
setLogLines((prev) =>
isLive ? appendCapped(prev, parsed, MAX_LIVE_LOG_LINES) : [...prev, ...parsed],
);
appendLogLines(parsed);
}
}
@ -3883,17 +3897,23 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
setIsFollowing((prev) => (prev ? prev : true));
}, [events.length, logLines.length, isLive, getScrollContainer]);
// Fetch persisted shell log
// Reset only when the log source changes, never when visibility changes.
useEffect(() => {
let cancelled = false;
pendingLogLineRef.current = "";
logMergeRefs.current = { seenChunkKeys: new Set(), trimmedSeqFloorByRun: new Map() };
seenProgressLogLineKeysRef.current = new Set();
setLogLines([]);
setLogOffset(0);
setHasMoreLog(false);
setLoadingMoreLog(false);
setLogError(null);
}, [run.id, run.logRef, setLogOffset]);
// Fetch persisted shell log, retaining partial rows and offsets across hides.
useEffect(() => {
if (!visible) return;
let cancelled = false;
const offset = logOffsetRef.current;
if (!run.logRef && !shouldPollShellLog) {
setLogLoading(false);
return () => {
@ -3904,10 +3924,10 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
setLogLoading(true);
const load = async () => {
try {
const result = await heartbeatsApi.log(run.id, 0, RUN_LOG_PAGE_BYTES);
const result = await heartbeatsApi.log(run.id, offset, RUN_LOG_PAGE_BYTES);
if (cancelled) return;
appendLogContent(result.content, result.nextOffset === undefined);
const next = result.nextOffset ?? result.content.length;
const next = result.nextOffset ?? offset + result.content.length;
setLogOffset(next);
setHasMoreLog(!shouldPollShellLog && result.nextOffset !== undefined);
} catch (err) {
@ -3927,7 +3947,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
return () => {
cancelled = true;
};
}, [run.id, run.logRef, run.logBytes, shouldPollShellLog]);
}, [visible, run.id, run.logRef, run.logBytes, shouldPollShellLog]);
async function loadMorePersistedLog() {
if (loadingMoreLog || !hasMoreLog) return;
@ -3948,27 +3968,42 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
// Poll for live updates
useEffect(() => {
if (!isLive || isStreamingConnected) return;
if (!visible || !isLive || isStreamingConnected) return;
let pending = false;
let cancelled = false;
const interval = setInterval(async () => {
if (pending || cancelled || !getPageVisibility().visible) return;
pending = true;
const maxSeq = events.length > 0 ? Math.max(...events.map((e) => e.seq)) : 0;
try {
const newEvents = await heartbeatsApi.events(run.id, maxSeq, 100);
if (cancelled) return;
if (newEvents.length > 0) {
setEvents((prev) => appendCapped(prev, newEvents, MAX_LIVE_EVENTS));
}
} catch {
// ignore polling errors
} finally {
pending = false;
}
}, 2000);
return () => clearInterval(interval);
}, [run.id, isLive, isStreamingConnected, events]);
return () => {
cancelled = true;
clearInterval(interval);
};
}, [visible, run.id, isLive, isStreamingConnected, events]);
// Poll shell log for running runs
useEffect(() => {
if (!shouldPollShellLog || isStreamingConnected) return;
if (!visible || !shouldPollShellLog || isStreamingConnected) return;
let pending = false;
let cancelled = false;
const interval = setInterval(async () => {
if (pending || cancelled || !getPageVisibility().visible) return;
pending = true;
try {
const result = await heartbeatsApi.log(run.id, logOffset, 256_000);
if (cancelled) return;
if (result.content) {
appendLogContent(result.content, result.nextOffset === undefined);
}
@ -3980,14 +4015,19 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
} catch (err) {
if (isRunLogUnavailable(err)) return;
// ignore polling errors
} finally {
pending = false;
}
}, 2000);
return () => clearInterval(interval);
}, [run.id, shouldPollShellLog, isStreamingConnected, logOffset]);
return () => {
cancelled = true;
clearInterval(interval);
};
}, [visible, run.id, shouldPollShellLog, isStreamingConnected, logOffset]);
// Stream live updates from websocket (primary path for running runs).
useEffect(() => {
if (!isLive) return;
if (!visible || !isLive) return;
let closed = false;
let reconnectTimer: number | null = null;
@ -4031,7 +4071,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
const streamRaw = asNonEmptyString(payload.stream);
const stream = streamRaw === "stderr" || streamRaw === "system" ? streamRaw : "stdout";
const ts = asNonEmptyString((payload as Record<string, unknown>).ts) ?? event.createdAt;
setLogLines((prev) => appendCapped(prev, [{ ts, stream, chunk }], MAX_LIVE_LOG_LINES));
appendLogLines([{ ts, stream, chunk, seq: readChunkSeq(payload.seq) }]);
return;
}
@ -4041,7 +4081,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
const key = heartbeatProgressLogLineKey(line);
if (seenProgressLogLineKeysRef.current.has(key)) return;
seenProgressLogLineKeysRef.current.add(key);
setLogLines((prev) => appendCapped(prev, [line], MAX_LIVE_LOG_LINES));
appendLogLines([line]);
return;
}
@ -4106,7 +4146,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
socket.close(1000, "run_detail_unmount");
}
};
}, [isLive, run.companyId, run.id, run.agentId]);
}, [visible, isLive, run.companyId, run.id, run.agentId]);
const censorUsernameInLogs = useQuery({
queryKey: queryKeys.instance.generalSettings,

View File

@ -1,3 +1,5 @@
import { mergeRunLogChunks, readChunkSeq } from "../lib/run-log-chunks";
import { getPageVisibility, usePageVisibility } from "../lib/page-visibility";
import { useCallback, useEffect, useMemo, useState, useRef } from "react";
import { useParams, useNavigate, Link, Navigate, useBeforeUnload, type NavigateFunction } from "@/lib/router";
import { useQuery, useMutation, useQueryClient, type QueryClient } from "@tanstack/react-query";
@ -384,6 +386,7 @@ function runMetrics(run: HeartbeatRun) {
}
export type RunLogChunk = {
seq?: number;
ts: string;
stream: "stdout" | "stderr" | "system";
chunk: string;
@ -3796,13 +3799,20 @@ function RunDetail({ run: initialRun, agentRouteId, adapterType, adapterConfig }
/* ---- Log Viewer ---- */
function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: string }) {
export function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: string }) {
const { visible } = usePageVisibility();
const [events, setEvents] = useState<HeartbeatRunEvent[]>([]);
const [logLines, setLogLines] = useState<Array<{ ts: string; stream: "stdout" | "stderr" | "system"; chunk: string }>>([]);
const [logLines, setLogLines] = useState<RunLogChunk[]>([]);
const [loading, setLoading] = useState(true);
const [logLoading, setLogLoading] = useState(!!run.logRef);
const [logError, setLogError] = useState<string | null>(null);
const [logOffset, setLogOffset] = useState(0);
const [logOffset, setLogOffsetState] = useState(0);
const logOffsetRef = useRef(0);
const setLogOffset = useCallback((next: number | ((previous: number) => number)) => {
logOffsetRef.current = typeof next === "function" ? next(logOffsetRef.current) : next;
setLogOffsetState(logOffsetRef.current);
}, []);
const logMergeRefs = useRef({ seenChunkKeys: new Set<string>(), trimmedSeqFloorByRun: new Map<string, number>() });
const [hasMoreLog, setHasMoreLog] = useState(false);
const [loadingMoreLog, setLoadingMoreLog] = useState(false);
const [isFollowing, setIsFollowing] = useState(false);
@ -3829,6 +3839,12 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
return err instanceof ApiError && err.status === 404;
}
function appendLogLines(incoming: RunLogChunk[]) {
setLogLines((previous) => mergeRunLogChunks(run.id, previous, incoming.map((line) => ({
...line, dedupeKey: `log:${run.id}:${line.ts}:${line.stream}:${line.chunk}`,
})), logMergeRefs.current, isLive ? MAX_LIVE_LOG_LINES : Number.POSITIVE_INFINITY).chunks);
}
function appendLogContent(content: string, finalize = false) {
if (!content && !finalize) return;
const combined = `${pendingLogLineRef.current}${content}`;
@ -3839,18 +3855,18 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
pendingLogLineRef.current = "";
}
const parsed: Array<{ ts: string; stream: "stdout" | "stderr" | "system"; chunk: string }> = [];
const parsed: RunLogChunk[] = [];
for (const line of split) {
const trimmed = line.trim();
if (!trimmed) continue;
try {
const raw = JSON.parse(trimmed) as { ts?: unknown; stream?: unknown; chunk?: unknown };
const raw = JSON.parse(trimmed) as { ts?: unknown; stream?: unknown; chunk?: unknown; seq?: unknown };
const stream =
raw.stream === "stderr" || raw.stream === "system" ? raw.stream : "stdout";
const chunk = typeof raw.chunk === "string" ? raw.chunk : "";
const ts = typeof raw.ts === "string" ? raw.ts : new Date().toISOString();
if (!chunk) continue;
parsed.push({ ts, stream, chunk });
parsed.push({ ts, stream, chunk, seq: readChunkSeq(raw.seq) });
} catch {
// ignore malformed lines
}
@ -3859,9 +3875,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
if (parsed.length > 0) {
// Live runs stream forever, so cap the retained tail. Terminated runs are
// paginated by the user via "Load more log" and keep their full history.
setLogLines((prev) =>
isLive ? appendCapped(prev, parsed, MAX_LIVE_LOG_LINES) : [...prev, ...parsed],
);
appendLogLines(parsed);
}
}
@ -3959,17 +3973,23 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
setIsFollowing((prev) => (prev ? prev : true));
}, [events.length, logLines.length, isLive, getScrollContainer]);
// Fetch persisted shell log
// Reset only when the log source changes, never when visibility changes.
useEffect(() => {
let cancelled = false;
pendingLogLineRef.current = "";
logMergeRefs.current = { seenChunkKeys: new Set(), trimmedSeqFloorByRun: new Map() };
seenProgressLogLineKeysRef.current = new Set();
setLogLines([]);
setLogOffset(0);
setHasMoreLog(false);
setLoadingMoreLog(false);
setLogError(null);
}, [run.id, run.logRef, setLogOffset]);
// Fetch persisted shell log, retaining partial rows and offsets across hides.
useEffect(() => {
if (!visible) return;
let cancelled = false;
const offset = logOffsetRef.current;
if (!run.logRef && !shouldPollShellLog) {
setLogLoading(false);
return () => {
@ -3980,10 +4000,10 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
setLogLoading(true);
const load = async () => {
try {
const result = await heartbeatsApi.log(run.id, 0, RUN_LOG_PAGE_BYTES);
const result = await heartbeatsApi.log(run.id, offset, RUN_LOG_PAGE_BYTES);
if (cancelled) return;
appendLogContent(result.content, result.nextOffset === undefined);
const next = result.nextOffset ?? result.content.length;
const next = result.nextOffset ?? offset + result.content.length;
setLogOffset(next);
setHasMoreLog(!shouldPollShellLog && result.nextOffset !== undefined);
} catch (err) {
@ -4003,7 +4023,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
return () => {
cancelled = true;
};
}, [run.id, run.logRef, run.logBytes, shouldPollShellLog]);
}, [visible, run.id, run.logRef, run.logBytes, shouldPollShellLog]);
async function loadMorePersistedLog() {
if (loadingMoreLog || !hasMoreLog) return;
@ -4024,27 +4044,42 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
// Poll for live updates
useEffect(() => {
if (!isLive || isStreamingConnected) return;
if (!visible || !isLive || isStreamingConnected) return;
let pending = false;
let cancelled = false;
const interval = setInterval(async () => {
if (pending || cancelled || !getPageVisibility().visible) return;
pending = true;
const maxSeq = events.length > 0 ? Math.max(...events.map((e) => e.seq)) : 0;
try {
const newEvents = await heartbeatsApi.events(run.id, maxSeq, 100);
if (cancelled) return;
if (newEvents.length > 0) {
setEvents((prev) => appendCapped(prev, newEvents, MAX_LIVE_EVENTS));
}
} catch {
// ignore polling errors
} finally {
pending = false;
}
}, 2000);
return () => clearInterval(interval);
}, [run.id, isLive, isStreamingConnected, events]);
return () => {
cancelled = true;
clearInterval(interval);
};
}, [visible, run.id, isLive, isStreamingConnected, events]);
// Poll shell log for running runs
useEffect(() => {
if (!shouldPollShellLog || isStreamingConnected) return;
if (!visible || !shouldPollShellLog || isStreamingConnected) return;
let pending = false;
let cancelled = false;
const interval = setInterval(async () => {
if (pending || cancelled || !getPageVisibility().visible) return;
pending = true;
try {
const result = await heartbeatsApi.log(run.id, logOffset, 256_000);
if (cancelled) return;
if (result.content) {
appendLogContent(result.content, result.nextOffset === undefined);
}
@ -4056,14 +4091,19 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
} catch (err) {
if (isRunLogUnavailable(err)) return;
// ignore polling errors
} finally {
pending = false;
}
}, 2000);
return () => clearInterval(interval);
}, [run.id, shouldPollShellLog, isStreamingConnected, logOffset]);
return () => {
cancelled = true;
clearInterval(interval);
};
}, [visible, run.id, shouldPollShellLog, isStreamingConnected, logOffset]);
// Stream live updates from websocket (primary path for running runs).
useEffect(() => {
if (!isLive) return;
if (!visible || !isLive) return;
let closed = false;
let reconnectTimer: number | null = null;
@ -4107,7 +4147,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
const streamRaw = asNonEmptyString(payload.stream);
const stream = streamRaw === "stderr" || streamRaw === "system" ? streamRaw : "stdout";
const ts = asNonEmptyString((payload as Record<string, unknown>).ts) ?? event.createdAt;
setLogLines((prev) => appendCapped(prev, [{ ts, stream, chunk }], MAX_LIVE_LOG_LINES));
appendLogLines([{ ts, stream, chunk, seq: readChunkSeq(payload.seq) }]);
return;
}
@ -4117,7 +4157,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
const key = heartbeatProgressLogLineKey(line);
if (seenProgressLogLineKeysRef.current.has(key)) return;
seenProgressLogLineKeysRef.current.add(key);
setLogLines((prev) => appendCapped(prev, [line], MAX_LIVE_LOG_LINES));
appendLogLines([line]);
return;
}
@ -4182,7 +4222,7 @@ function LogViewer({ run, adapterType }: { run: HeartbeatRun; adapterType: strin
socket.close(1000, "run_detail_unmount");
}
};
}, [isLive, run.companyId, run.id, run.agentId]);
}, [visible, isLive, run.companyId, run.id, run.agentId]);
const censorUsernameInLogs = useQuery({
queryKey: queryKeys.instance.generalSettings,