Files
PaperClipAI/server/src/services/run-log-store.ts
T
Devin Foley 187a90b7bc Add opt-in in-flight run-log mirroring with graceful-shutdown flush (#10512)
## Thinking Path

> - Paperclip is the open source app people use to manage AI agents for
work
> - The run-log store records each agent run's output and can mirror
completed logs to S3-compatible object storage
> - The mirror uploads only on finalize, so a server restart mid-run
loses the whole in-flight log
> - Deployments and crashes are routine on ephemeral hosts, and lost run
output makes failed runs impossible to debug
> - This pull request adds an opt-in throttled mirror for still-running
logs plus a graceful-shutdown flush
> - The benefit is that a restart mid-run keeps the log tail up to the
last mirror interval, and an orderly restart keeps everything

## Linked Issues or Issue Description

No public issue exists — describing the feature inline (per the feature
request template).

**Subsystem affected**
server/ — REST API & orchestration services

**Problem or motivation**
`RUN_LOG_S3_BUCKET` gives finished run logs durability, but the mirror
uploads only on finalize. A run that is still writing when the server
restarts leaves nothing in object storage. On hosts with ephemeral disks
the local file is gone too, so the run's output is lost end to end and
failed runs cannot be debugged.

**Proposed solution**
Mirror the in-flight log to the same object key on a throttled cadence
(`RUN_LOG_S3_INFLIGHT_MIRROR_SECONDS`), and flush dirty tails during
graceful shutdown. Keep it opt-in so existing deployments see zero new
upload traffic unless they ask for it.

**Alternatives considered**
Per-append uploads (rejected: one PUT per output chunk is hostile to S3
endpoints and run latency). Chunked part objects with read-time
stitching (rejected: complicates the read path, and S3 multipart minimum
part sizes do not fit small tails). Persistent volumes (rejected
upstream already: the data dir is deliberately an emptyDir in hardened
cloud_tenant deployments).

**Roadmap alignment**
Not on ROADMAP.md; extends the existing run-log durability mirror
without changing any default behavior.

**Additional context**
Ranged reads already serve partial objects like a live tail, so the read
path needs no change; finalize overwrites the mirror with the complete
file.

**Related PRs (dedup search):** the finalize-only S3 mirror landed
previously and this extends it; no duplicate or competing PR found for
in-flight run-log mirroring.

## What Changed

- `server/src/services/run-log-store.ts`: new opt-in `inflightMirrorMs`
on the S3 options (`RUN_LOG_S3_INFLIGHT_MIRROR_SECONDS` env). When set,
appends schedule at most one upload of the current file per interval, to
the same key finalize uses. Ranged reads already serve that key, so a
partial object behaves like a live tail and needs no read-path change.
Finalize retires the in-flight bookkeeping and waits out an upload
already on the wire, so a stale partial can never overwrite a finalized
log. Upload failures warn, re-mark the tail dirty, and retry at most
once per interval.
- `server/src/services/run-log-store.ts`: new `flushInflightMirrors()`
on the store and a module-level `flushInFlightRunLogMirrors()` for the
shutdown path. Both are no-ops when the mirror is off.
- `server/src/index.ts`: graceful shutdown flushes dirty in-flight tails
after the heartbeat run drain, so runs the drain did not finalize
(timeouts, the hot-restart skip path) still persist their output.
- `server/src/services/run-log-store.test.ts`: five new tests —
off-by-default (no uploads before finalize), tail preserved after a wipe
without finalize, throttle coalescing with a single flush upload,
finalize superseding the in-flight mirror and retiring its timer, and
upload failures never breaking appends with recovery on the next flush.

## Verification

- `pnpm vitest run server/src/services/run-log-store.test.ts` — 13
passed (8 existing + 5 new).
- `pnpm vitest run server/src/__tests__/heartbeat-run-log.test.ts
server/src/__tests__/heartbeat-active-run-output-watchdog.test.ts` — 21
passed (consumers of the store, unchanged behavior).
- `pnpm -C server run typecheck` — clean.
- Self-hosted behavior is unchanged unless
`RUN_LOG_S3_INFLIGHT_MIRROR_SECONDS` is set: with the variable unset
there are zero new uploads and the finalize-only mirroring is
byte-identical (asserted by the off-by-default test).

## Risks

- Low. The feature is opt-in; unset env preserves today's behavior
exactly. When enabled, worst case is one extra PUT per interval per
active run, and every upload is best-effort — a failing endpoint warns
and never breaks appends, finalization, or shutdown. The finalize path
awaits any in-flight upload before writing the complete file, closing
the only overwrite race the design introduces. Timers are `unref`ed so
the mirror never keeps the process alive.

## Model Used

Claude Fable 5 (`claude-fable-5`, Anthropic; Claude Code CLI with
extended thinking and tool use; tests executed locally via Vitest).



## 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 (no duplicates found for in-flight run-log mirroring; the
finalize-only mirror landed previously and this extends it)
- [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
2026-07-30 14:01:02 -07:00

430 lines
18 KiB
TypeScript

import { createReadStream, promises as fs } from "node:fs";
import path from "node:path";
import { createHash } from "node:crypto";
import { notFound } from "../errors.js";
import { resolvePaperclipInstanceRoot } from "../home-paths.js";
import { createS3StorageProvider } from "../storage/s3-provider.js";
import type { StorageProvider } from "../storage/types.js";
export type RunLogStoreType = "local_file";
export interface RunLogHandle {
store: RunLogStoreType;
logRef: string;
}
export interface RunLogReadOptions {
offset?: number;
limitBytes?: number;
}
export interface RunLogReadResult {
content: string;
nextOffset?: number;
}
export interface RunLogFinalizeSummary {
bytes: number;
sha256?: string;
compressed: boolean;
}
export interface RunLogStore {
begin(input: { companyId: string; agentId: string; runId: string }): Promise<RunLogHandle>;
append(
handle: RunLogHandle,
event: { stream: "stdout" | "stderr" | "system"; chunk: string; ts: string; seq?: number },
): Promise<number>;
finalize(handle: RunLogHandle): Promise<RunLogFinalizeSummary>;
read(handle: RunLogHandle, opts?: RunLogReadOptions): Promise<RunLogReadResult>;
// Optional so existing fakes/fixtures keep compiling: uploads every dirty
// in-flight mirror immediately (graceful-shutdown path). No-op when the
// in-flight mirror is not enabled.
flushInflightMirrors?(): Promise<void>;
}
function safeSegments(...segments: string[]) {
return segments.map((segment) => segment.replace(/[^a-zA-Z0-9._-]/g, "_"));
}
function resolveWithin(basePath: string, relativePath: string) {
const resolved = path.resolve(basePath, relativePath);
const base = path.resolve(basePath) + path.sep;
if (!resolved.startsWith(base) && resolved !== path.resolve(basePath)) {
throw new Error("Invalid log path");
}
return resolved;
}
function normalizeKeyPrefix(prefix: string | undefined): string {
if (!prefix) return "";
return prefix.trim().replace(/^\/+/, "").replace(/\/+$/, "");
}
export interface DurableRunLogStoreOptions {
basePath: string;
// When provided, completed logs are mirrored to object storage on finalize and
// served from there on read whenever the local file is missing (e.g. the pod
// rolled and wiped the emptyDir). When omitted, the store is local-only (the
// historical behaviour: a restart loses the log).
s3?: {
provider: StorageProvider;
keyPrefix?: string;
// When > 0, ALSO mirror the still-running log to the same object key at
// most once per this interval (plus a flush hook for graceful shutdown),
// so a crash mid-run loses at most one interval's tail instead of the
// whole log. Off (undefined/0) preserves the historical finalize-only
// mirroring: no extra PUT traffic unless explicitly opted in.
inflightMirrorMs?: number;
};
}
// Run-log store with TRANSPARENT durability. The store id stays "local_file" so
// nothing downstream (feedback.ts, the heartbeat read cast, fixtures) changes;
// the S3 mirror is keyed by the same logRef and is purely an implementation
// detail. Live append/tail stays on the pod-local file (fast, no per-chunk PUT);
// on finalize the complete .ndjson is uploaded to object storage; on read we try
// local first and fall back to S3 when the local file is gone. This is the fix
// for "Run log not found" after a deploy/restart (the /paperclip data dir is an
// emptyDir in cloud_tenant mode -- persistence is disabled to avoid the
// operator's privileged selinux-relabel init container in our hardened ns).
//
// Optionally (inflightMirrorMs > 0) the still-running log is ALSO mirrored to
// the same key at a throttled cadence and flushed on graceful shutdown, so a
// restart mid-run preserves the tail up to the last mirror instead of losing
// the whole in-flight log. Finalize retires the in-flight bookkeeping (waiting
// out any upload already on the wire) before writing the complete file, so a
// stale partial can never overwrite a finalized log.
export function createDurableRunLogStore(options: DurableRunLogStoreOptions): RunLogStore {
const { basePath } = options;
const s3 = options.s3;
const s3Prefix = normalizeKeyPrefix(s3?.keyPrefix);
const inflightMirrorMs = s3?.inflightMirrorMs && s3.inflightMirrorMs > 0 ? s3.inflightMirrorMs : 0;
function s3Key(logRef: string): string {
return s3Prefix ? `${s3Prefix}/${logRef}` : logRef;
}
// In-flight mirror bookkeeping, keyed by logRef. The mirror uploads the
// CURRENT (partial) file to the SAME key finalize uses: readers already
// range-read that key, so a partial object is served exactly like a live
// tail, and finalize simply overwrites it with the complete file. One
// entry exists only between the first post-interval-eligible append and
// finalize.
interface InflightMirrorEntry {
dirty: boolean;
lastMirrorAt: number;
timer: NodeJS.Timeout | null;
upload: Promise<boolean> | null;
}
const inflightMirrors = new Map<string, InflightMirrorEntry>();
function mirrorInflightNow(logRef: string, entry: InflightMirrorEntry): Promise<boolean> {
entry.dirty = false;
const upload = (async () => {
const absPath = resolveWithin(basePath, logRef);
const stat = await fs.stat(absPath);
if (stat.size === 0) return true;
await s3!.provider.putObject({
objectKey: s3Key(logRef),
// Bound the stream to the stat'ed size: the run is still appending,
// and an unbounded stream that grows past stat.size would violate
// the declared contentLength and fail (or truncate) the upload.
// Bytes appended after the stat stay dirty and ride the next mirror.
body: createReadStream(absPath, { start: 0, end: stat.size - 1 }),
contentType: "application/x-ndjson",
contentLength: stat.size,
});
return true;
})().catch((err) => {
// Best-effort like the finalize mirror: a failing upload must never
// break the run, but a persistently broken mirror should be visible.
console.warn(
`[run-log-store] Failed to mirror in-flight run log to object storage (key: ${s3Key(logRef)}):`,
err,
);
// Re-dirty so the tail retries next interval even without new appends;
// the lastMirrorAt stamp below bounds retries to one per interval.
entry.dirty = true;
return false;
}).finally(() => {
// Stamp AFTER the attempt so a slow or failing endpoint self-throttles
// to one attempt per interval instead of hot-looping.
entry.lastMirrorAt = Date.now();
entry.upload = null;
if (entry.dirty) scheduleInflightMirror(logRef, entry);
});
entry.upload = upload;
return upload;
}
function scheduleInflightMirror(logRef: string, entry: InflightMirrorEntry): void {
if (entry.timer || entry.upload) return;
const delay = Math.max(0, inflightMirrorMs - (Date.now() - entry.lastMirrorAt));
entry.timer = setTimeout(() => {
entry.timer = null;
void mirrorInflightNow(logRef, entry);
}, delay);
// Never keep the process alive just to mirror a tail.
entry.timer.unref?.();
}
function noteInflightAppend(logRef: string): void {
if (!s3 || inflightMirrorMs <= 0) return;
let entry = inflightMirrors.get(logRef);
if (!entry) {
// First mirror lands one full interval after the first append: a run
// that finalizes sooner is covered by the finalize upload, and this
// keeps the steady-state cost at one PUT per interval per active run.
entry = { dirty: false, lastMirrorAt: Date.now(), timer: null, upload: null };
inflightMirrors.set(logRef, entry);
}
entry.dirty = true;
scheduleInflightMirror(logRef, entry);
}
async function retireInflightMirror(logRef: string): Promise<void> {
const entry = inflightMirrors.get(logRef);
if (!entry) return;
inflightMirrors.delete(logRef);
if (entry.timer) {
clearTimeout(entry.timer);
entry.timer = null;
}
// An upload still in flight could otherwise finish AFTER finalize's
// complete-file upload and overwrite it with a stale partial.
if (entry.upload) await entry.upload;
}
async function ensureDir(relativeDir: string) {
const dir = resolveWithin(basePath, relativeDir);
await fs.mkdir(dir, { recursive: true });
}
async function readLocalRange(
filePath: string,
offset: number,
limitBytes: number,
): Promise<RunLogReadResult | null> {
const stat = await fs.stat(filePath).catch(() => null);
if (!stat) return null;
const start = Math.max(0, Math.min(offset, stat.size));
const end = Math.max(start, Math.min(start + limitBytes - 1, stat.size - 1));
if (start > end) return { content: "", nextOffset: start };
const chunks: Buffer[] = [];
try {
await new Promise<void>((resolve, reject) => {
const stream = createReadStream(filePath, { start, end });
stream.on("data", (chunk) => chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)));
stream.on("error", reject);
stream.on("end", () => resolve());
});
} catch (err) {
// File deleted between stat() and open (pod-roll cleanup racing a read):
// treat as missing so the caller falls through to the S3 mirror instead
// of surfacing the very "Run log not found" this store exists to prevent.
if ((err as NodeJS.ErrnoException | null)?.code === "ENOENT") return null;
throw err;
}
const content = Buffer.concat(chunks).toString("utf8");
const nextOffset = end + 1 < stat.size ? end + 1 : undefined;
return { content, nextOffset };
}
async function readS3Range(
logRef: string,
offset: number,
limitBytes: number,
): Promise<RunLogReadResult> {
if (!s3) throw notFound("Run log not found");
const key = s3Key(logRef);
const head = await s3.provider.headObject({ objectKey: key });
if (!head.exists) throw notFound("Run log not found");
const total = head.contentLength ?? 0;
const start = Math.max(0, Math.min(offset, total));
const end = Math.max(start, Math.min(start + limitBytes - 1, total - 1));
if (start > end || total === 0) return { content: "", nextOffset: start < total ? start : undefined };
const result = await s3.provider.getObject({ objectKey: key, range: { start, end } });
const chunks: Buffer[] = [];
await new Promise<void>((resolve, reject) => {
result.stream.on("data", (chunk) => chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)));
result.stream.on("error", reject);
result.stream.on("end", () => resolve());
});
const content = Buffer.concat(chunks).toString("utf8");
const nextOffset = end + 1 < total ? end + 1 : undefined;
return { content, nextOffset };
}
async function sha256File(filePath: string): Promise<string> {
return new Promise<string>((resolve, reject) => {
const hash = createHash("sha256");
const stream = createReadStream(filePath);
stream.on("data", (chunk) => hash.update(chunk));
stream.on("error", reject);
stream.on("end", () => resolve(hash.digest("hex")));
});
}
return {
async begin(input) {
const [companyId, agentId] = safeSegments(input.companyId, input.agentId);
const runId = safeSegments(input.runId)[0]!;
const relDir = path.join(companyId, agentId);
const relPath = path.join(relDir, `${runId}.ndjson`);
await ensureDir(relDir);
const absPath = resolveWithin(basePath, relPath);
await fs.writeFile(absPath, "", "utf8");
await retireInflightMirror(relPath);
return { store: "local_file", logRef: relPath };
},
async append(handle, event) {
if (handle.store !== "local_file") return 0;
const absPath = resolveWithin(basePath, handle.logRef);
const line = JSON.stringify({
ts: event.ts,
stream: event.stream,
chunk: event.chunk,
// Monotonic per-run sequence so readers can dedupe and order records
// even when several identical chunks share the same millisecond ts
// (common for ACP-style token deltas).
...(typeof event.seq === "number" && Number.isFinite(event.seq) ? { seq: event.seq } : {}),
});
const persisted = `${line}\n`;
await fs.appendFile(absPath, persisted, "utf8");
noteInflightAppend(handle.logRef);
return Buffer.byteLength(persisted, "utf8");
},
async finalize(handle) {
if (handle.store !== "local_file") return { bytes: 0, compressed: false };
await retireInflightMirror(handle.logRef);
const absPath = resolveWithin(basePath, handle.logRef);
const stat = await fs.stat(absPath).catch(() => null);
if (!stat) throw notFound("Run log not found");
const hash = await sha256File(absPath);
// Mirror the completed log to object storage so it survives a pod roll.
// Best-effort upload failures must NOT fail run finalization (which also
// records cost/usage); the local copy still serves reads until the pod
// rolls, and a failed mirror only loses durability for that one run.
if (s3) {
try {
// Stream from disk instead of buffering the whole .ndjson in the
// heap; long agent sessions can produce large logs. The file is
// complete at this point, so stat.size is the exact content length.
await s3.provider.putObject({
objectKey: s3Key(handle.logRef),
body: createReadStream(absPath),
contentType: "application/x-ndjson",
contentLength: stat.size,
});
} catch (err) {
// Best-effort: finalization must not break, but a persistently
// failing mirror (bad creds/bucket/endpoint) should be visible to
// operators before a pod roll makes the logs unreadable.
console.warn(
`[run-log-store] Failed to mirror run log to object storage (key: ${s3Key(handle.logRef)}):`,
err,
);
}
}
return { bytes: stat.size, sha256: hash, compressed: false };
},
async read(handle, opts) {
if (handle.store !== "local_file") throw notFound("Run log not found");
const absPath = resolveWithin(basePath, handle.logRef);
const offset = opts?.offset ?? 0;
const limitBytes = opts?.limitBytes ?? 256_000;
const local = await readLocalRange(absPath, offset, limitBytes);
if (local) return local;
// Local file gone (pod rolled) -> serve from the S3 mirror if configured.
return readS3Range(handle.logRef, offset, limitBytes);
},
async flushInflightMirrors() {
if (!s3 || inflightMirrorMs <= 0) return;
const flushEntry = async (logRef: string, entry: InflightMirrorEntry) => {
// Loop until the entry is clean: an append that lands while an
// upload is on the wire re-dirties the entry, and its follow-up
// mirror sits on an unref'ed timer that would never fire once the
// process exits — so re-check after every await instead of trusting
// a single pass. A FAILED attempt ends the loop instead of retrying:
// hot-looping a down endpoint at shutdown would spin forever, and
// the flush is best-effort by design.
for (;;) {
if (entry.timer) {
clearTimeout(entry.timer);
entry.timer = null;
}
if (entry.upload) {
await entry.upload;
continue;
}
if (!entry.dirty) return;
const uploaded = await mirrorInflightNow(logRef, entry);
if (!uploaded) {
if (entry.timer) {
clearTimeout(entry.timer);
entry.timer = null;
}
return;
}
}
};
await Promise.all([...inflightMirrors].map(([logRef, entry]) => flushEntry(logRef, entry)));
},
};
}
// Build the run-log S3 mirror from dedicated RUN_LOG_S3_* env. Deliberately
// separate from PAPERCLIP_STORAGE_PROVIDER so enabling durable run logs does
// NOT redirect the product's workspace/file storage (smaller blast radius).
// Unset RUN_LOG_S3_BUCKET -> no mirror -> local-only (safe degrade). Creds come
// from the standard AWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEY chain.
function resolveRunLogS3(): DurableRunLogStoreOptions["s3"] {
const bucket = process.env.RUN_LOG_S3_BUCKET?.trim();
if (!bucket) return undefined;
const provider = createS3StorageProvider({
bucket,
region: process.env.RUN_LOG_S3_REGION?.trim() || "us-east-1",
endpoint: process.env.RUN_LOG_S3_ENDPOINT?.trim() || undefined,
prefix: undefined, // prefixing is handled by keyPrefix below (kept off the provider)
forcePathStyle: process.env.RUN_LOG_S3_FORCE_PATH_STYLE
? process.env.RUN_LOG_S3_FORCE_PATH_STYLE === "true"
: true, // Cubbit (and most S3-compatible endpoints) need path-style
});
// Opt-in in-flight tail mirroring: at most one partial upload per interval
// per active run, so a crash loses at most one interval's tail. Unset/0
// keeps the historical finalize-only mirroring.
const inflightSeconds = Number.parseFloat(process.env.RUN_LOG_S3_INFLIGHT_MIRROR_SECONDS ?? "");
return {
provider,
keyPrefix: process.env.RUN_LOG_S3_PREFIX?.trim() || "run-logs",
inflightMirrorMs:
Number.isFinite(inflightSeconds) && inflightSeconds > 0 ? Math.round(inflightSeconds * 1000) : undefined,
};
}
let cachedStore: RunLogStore | null = null;
export function getRunLogStore() {
if (cachedStore) return cachedStore;
const basePath = process.env.RUN_LOG_BASE_PATH ?? path.resolve(resolvePaperclipInstanceRoot(), "data", "run-logs");
cachedStore = createDurableRunLogStore({ basePath, s3: resolveRunLogS3() });
return cachedStore;
}
// Graceful-shutdown hook: upload every dirty in-flight run-log tail before
// the process exits, so an orderly restart (deploy, SIGTERM) loses nothing
// even for runs that never reach finalize. No-op when the store was never
// created or in-flight mirroring is off.
export async function flushInFlightRunLogMirrors(): Promise<void> {
await cachedStore?.flushInflightMirrors?.();
}