diff --git a/.changelog/unreleased/fixed-492-cur-record-stamp-reconcile.md b/.changelog/unreleased/fixed-492-cur-record-stamp-reconcile.md new file mode 100644 index 00000000..1d97ed2b --- /dev/null +++ b/.changelog/unreleased/fixed-492-cur-record-stamp-reconcile.md @@ -0,0 +1 @@ +- **Report failed cur/ stamps**. Retry failed ack/nack stamp writes. diff --git a/packages/cli/src/utils/mail.ts b/packages/cli/src/utils/mail.ts index 12f28911..2784bcdb 100644 --- a/packages/cli/src/utils/mail.ts +++ b/packages/cli/src/utils/mail.ts @@ -11,7 +11,7 @@ import { } from "@tpsdev-ai/agent"; import { sanitizeIdentifier } from "../schema/sanitizer.js"; import { logEvent } from "./archive.js"; -import { acquireMailLock, acquireMailLockSync, type MailLock } from "./mail-lock.js"; +import { acquireMailLock, acquireMailLockSync, mailLockPath, type MailLock } from "./mail-lock.js"; import { createMailVerifyClient, type MailVerifyConfig } from "./mail-verify.js"; // cli#429: the ONE id shape rule, re-exported so the openclaw-tps-mail plugin @@ -189,10 +189,16 @@ export function updateExistingRecord( mutate: (record: T) => T | null, options: { snapshot?: T; afterWrite?: (record: T) => void; nonBlocking?: boolean } = {}, ): UpdateExistingResult { - const lock = acquireMailLockSync(dirname(dirname(path)), options.nonBlocking ? { timeoutMs: 0 } : {}); + const root = dirname(dirname(path)); + let lock: MailLock | null; + try { + lock = acquireMailLockSync(root, options.nonBlocking ? { timeoutMs: 0 } : {}); + } catch (err: any) { + throw Object.assign(err, { path: err?.path ?? mailLockPath(root) }); + } if (!lock) { if (options.nonBlocking) return { status: "busy" }; - throw new Error(`mail lock contention timeout for ${path}`); + throw Object.assign(new Error(`mail lock contention timeout for ${path}`), { path: mailLockPath(root) }); } const scratchPath = join(dirname(path), `.ack-${randomUUID()}.tmp`); let fd: number | undefined; @@ -215,7 +221,7 @@ export function updateExistingRecord( options.afterWrite?.(updated); return { status: "updated", record: updated }; } catch (err: any) { - if (!replaced && err?.code === "ENOENT") return { status: "gone" }; + if (!replaced && err?.code === "ENOENT" && err?.path === path) return { status: "gone" }; throw err; } finally { if (fd !== undefined) { try { closeSync(fd); } catch {} } diff --git a/plugins/openclaw-tps-mail/src/diagnostics.ts b/plugins/openclaw-tps-mail/src/diagnostics.ts new file mode 100644 index 00000000..36edc69b --- /dev/null +++ b/plugins/openclaw-tps-mail/src/diagnostics.ts @@ -0,0 +1,29 @@ +export type ObligationPresence = "retained" | "none" | "unknown"; + +export interface StampDiagnosticFacts { + kind: string; + actor: string; + id?: string; + path: string; + code?: string; + obligation: ObligationPresence; + retriesExhausted?: boolean; +} + +export function formatStampDiagnostic(facts: StampDiagnosticFacts): string { + const fields = [ + `tps-mail: ${facts.kind}:`, + ...(facts.id === undefined ? [] : [facts.id]), + `actor=${facts.actor}`, + `path=${facts.path}`, + ...(facts.code === undefined ? [] : [`code=${facts.code}`]), + ]; + const state = facts.obligation === "retained" ? "obligation retained" + : facts.obligation === "unknown" ? "state unknown" : undefined; + return [ + fields.join(" "), + ...(state ? [state] : []), + ...(facts.retriesExhausted ? ["no retries left"] : []), + "resolve the failure and restart the account", + ].join("; "); +} diff --git a/plugins/openclaw-tps-mail/src/index.ts b/plugins/openclaw-tps-mail/src/index.ts index f78f4e67..a57564da 100644 --- a/plugins/openclaw-tps-mail/src/index.ts +++ b/plugins/openclaw-tps-mail/src/index.ts @@ -54,9 +54,9 @@ import { randomUUID } from "node:crypto"; import { existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFileSync, watch as fsWatch, type FSWatcher } from "node:fs"; import { homedir } from "node:os"; -import { basename, resolve } from "node:path"; +import { basename, dirname, resolve } from "node:path"; import type { Envelope, ChainEntry } from "@tpsdev-ai/agent"; -import { signEnvelope, verifyEnvelope, verifiedMailTier } from "@tpsdev-ai/agent"; +import { signEnvelope, verifyEnvelope, verifiedMailTier, mailLockPath } from "@tpsdev-ai/agent"; import { readAgentPrivateKey } from "@tpsdev-ai/cli/utils/agent-keys"; import { signForDelivery } from "@tpsdev-ai/cli/utils/mail-producer"; import { isValidEnvelopeId, mailRootForRecordPath, updateExistingRecord, promote, recoverPromoted, verifyRecordForMailbox, sweepStrandedPromoteScratch } from "@tpsdev-ai/cli/utils/mail"; @@ -85,6 +85,7 @@ import { type ReceiptScanDirs, } from "./obligations.js"; import type { OpenClawPluginApi } from "openclaw/plugin-sdk"; +import { formatStampDiagnostic, type ObligationPresence } from "./diagnostics.js"; import { detectHostOpenClawVersion, evaluateHostSilentReplyGuard } from "./host-version.js"; import type { ChannelPlugin } from "openclaw/plugin-sdk/core"; import type { @@ -262,11 +263,17 @@ function findBoundAgents(cfg: any, accountId: string): string[] { // ─── TPS mail envelope helpers ─────────────────────────────────────────────── -function readMailFile(filePath: string): TpsMailBody | null { +function readMailFile(filePath: string, onReadError?: (path: string, code?: string) => void): TpsMailBody | null { try { const raw = readFileSync(filePath, "utf-8"); - return JSON.parse(raw) as TpsMailBody; - } catch { + const record: unknown = JSON.parse(raw); + if (onReadError && (!record || typeof record !== "object" || Array.isArray(record) || typeof (record as TpsMailBody).id !== "string")) { + onReadError?.(filePath); + return null; + } + return record as TpsMailBody; + } catch (err: any) { + onReadError?.(filePath, err?.code); return null; } } @@ -618,12 +625,77 @@ function routeFor(mailDir: string, cfg: any, accountId: string, to: string): Mai return resolveMailRoute({ to, mailDir, localAgents: findBoundAgents(cfg, accountId) }); } -export function patchMailFile(path: string, patch: Partial): boolean { +type CurRecordUpdate = + | { ok: true } + | { ok: false; reason: "record-missing" | "write-failed"; path: string; code?: string }; + +export function patchMailFile(path: string, patch: Partial, stampKey: "ackedAt" | "nackedAt" | undefined, inboundId: string): CurRecordUpdate { try { - return updateExistingRecord(path, (current) => Object.assign(current, patch)).status === "updated"; - } catch { - return false; + let alreadyStamped = false; + const r = updateExistingRecord(path, (current) => { + if (current.id !== inboundId) { + throw Object.assign(new Error("cur identity changed"), { code: "ID_MISMATCH" }); + } + if (stampKey && current[stampKey]) { + alreadyStamped = true; + return null; + } + return Object.assign(current, patch); + }); + if (r.status === "changed" && alreadyStamped) return { ok: true }; + if (r.status === "updated") return { ok: true }; + if (r.status === "gone") return { ok: false, reason: "record-missing", path, code: "ENOENT" }; + return { ok: false, reason: "write-failed", path, code: undefined }; + } catch (err: any) { + return { ok: false, reason: "write-failed", path: typeof err?.path === "string" && (err.path === mailLockPath(dirname(dirname(path))) || err.path.startsWith(`${mailLockPath(dirname(dirname(path)))}.`)) + ? mailLockPath(dirname(dirname(path))) : path, code: typeof err?.code === "string" ? err.code : undefined }; + } +} + +function stampObligationPresence(mailDir: string, agent: string, inboundId: string): ObligationPresence { + const result = readObligationResult(mailDir, agent, inboundId); + return result.status === "found" ? "retained" : result.status === "missing" ? "none" : "unknown"; +} + +export function reconcileTerminalCurStamps(mailDir: string, agent: string, log: any, records?: ObligationRecord[], unknownInbounds = new Set()): Set { + const reportReadError = (obligation: ObligationPresence, id: string) => (path: string, code?: string) => { + if (unknownInbounds.has(id)) return; + unknownInbounds.add(id); + log?.warn?.( + formatStampDiagnostic({ kind: "stamp-reconcile-read-failed", actor: agent, id, path, code, obligation }), + ); + }; + const onObligationReadError = (_path: string, _code: string | undefined, ids: string[]) => { + for (const id of ids) unknownInbounds.add(id); + }; + for (const rec of records ?? listObligations(mailDir, agent, onObligationReadError)) { + if (unknownInbounds.has("*") || unknownInbounds.has(rec.inboundId)) continue; + if (rec.state !== "acked" && rec.state !== "failed") continue; + let unreadable = false; + const onReadError = (path: string, code?: string) => { + unreadable = true; + reportReadError(stampObligationPresence(mailDir, agent, rec.inboundId), rec.inboundId)(path, code); + }; + const kind = rec.state === "acked" ? "ack" : "nack"; + const key = kind === "ack" ? "ackedAt" : "nackedAt"; + const curPath = findCurPath(mailDir, agent, rec.inboundId, onReadError); + if (unreadable || !curPath) continue; + const cur = readMailFile(curPath, onReadError); + if (!cur) continue; + if (cur[key]) continue; + const patch = kind === "ack" + ? { ackedAt: new Date().toISOString(), read: true } + : { nackedAt: new Date().toISOString(), nackReason: rec.failure ?? "failed" }; + const r = patchMailFile(curPath, patch, key, rec.inboundId); + if (r.ok) log?.info?.(`tps-mail: reconciled the ${kind} stamp for ${rec.inboundId} at ${curPath}`); + else { + unknownInbounds.add(rec.inboundId); + log?.warn?.( + formatStampDiagnostic({ kind: `${kind}-stamp-reconcile-failed`, actor: agent, id: rec.inboundId, path: r.path, code: r.code, obligation: stampObligationPresence(mailDir, agent, rec.inboundId) }), + ); + } } + return unknownInbounds; } /** @@ -872,6 +944,40 @@ function readObligationForCleanup(ctx: YieldContext, obligationId: string) { return result; } +/** Delays before each in-process retry of a failed terminal stamp (cli#492). */ +let stampRetryDelaysMs = [1000, 4000, 16000]; + +/** Test-only: shorten the stamp retry delays. Returns the previous delays. */ +export function setStampRetryDelaysForTests(delays: number[]): number[] { + const prev = stampRetryDelaysMs; + stampRetryDelaysMs = delays; + return prev; +} + +function stampTerminalCur( + ctx: YieldContext, + kind: "ack" | "nack", + patch: Partial, + attempt = 0, +): void { + const key = kind === "ack" ? "ackedAt" : "nackedAt"; + const stamped = patchMailFile(ctx.curPath, patch, key, ctx.inboundId); + if (stamped.ok) { + if (attempt > 0) ctx.log?.info?.(`tps-mail: ${kind}-stamp-retry-ok: ${ctx.inboundId} at ${ctx.curPath} (retry ${attempt})`); + return; + } + const delay = stamped.reason === "record-missing" ? undefined : stampRetryDelaysMs[attempt]; + ctx.log?.warn?.( + formatStampDiagnostic({ kind: `${kind}-stamp-failed`, actor: ctx.agent, id: ctx.inboundId, + path: stamped.path, code: stamped.code, obligation: stampObligationPresence(ctx.mailDir, ctx.agent, ctx.inboundId), + retriesExhausted: stamped.reason !== "record-missing" && delay === undefined }), + ); + if (delay === undefined || !isLiveContext(ctx)) return; + const timer = accountTimer(ctx, () => stampTerminalCur(ctx, kind, patch, attempt + 1), delay); + if (typeof (timer as any).unref === "function") (timer as any).unref(); +} + + function ackObligation(ctx: YieldContext, obligationId: string, why: string): void { if (!isLiveContext(ctx)) return; // Only stamp the inbound when the ACK TRANSITION actually landed. A terminal @@ -889,7 +995,7 @@ function ackObligation(ctx: YieldContext, obligationId: string, why: string): vo ); return; } - patchMailFile(ctx.curPath, { ackedAt: new Date().toISOString(), read: true }); + stampTerminalCur(ctx, "ack", { ackedAt: new Date().toISOString(), read: true }); ctx.log?.info?.(`tps-mail: acked ${ctx.inboundId} — ${why}`); releaseObligationState(obligationId); } @@ -1060,7 +1166,7 @@ async function settleObligation( // handed off (cli#403). releaseObligationState(obligationId); if (!s.alreadyStamped) { - patchMailFile(ctx.curPath, { nackedAt: new Date().toISOString(), nackReason: failedOn }); + stampTerminalCur(ctx, "nack", { nackedAt: new Date().toISOString(), nackReason: failedOn }); } // cli#389 round 10, item 1: AWAIT the one send, and record its outcome on the // record. The mail is owed until `nackSentAt` says otherwise (at-least-once). @@ -1534,18 +1640,31 @@ function installYieldSubscription(api: any): boolean { return true; } -/** The cur/ path for an inbound id (cur filenames are timestamp-id, not the id). */ -function findCurPath(mailDir: string, agent: string, inboundId: string): string | null { +function findCurPath(mailDir: string, agent: string, inboundId: string, onReadError?: (path: string, code?: string) => void): string | null { const curDir = resolve(mailDir, agent, "cur"); try { - for (const name of readdirSync(curDir)) { + const names = readdirSync(curDir); + if (onReadError && names.length > 4096) { + onReadError(curDir, "SCAN_LIMIT"); + return null; + } + let match: string | null = null; + for (const name of names) { if (!name.endsWith(".json")) continue; const p = resolve(curDir, name); - const rec = readMailFile(p); - if (rec?.id === inboundId) return p; + const rec = readMailFile(p, onReadError); + if (rec?.id !== inboundId) continue; + if (!onReadError) return p; + if (match) { + onReadError(curDir, "AMBIGUOUS_ID"); + return null; + } + match = p; } - } catch { - // no cur dir yet + if (!match) onReadError?.(curDir, "ENOENT"); + return match; + } catch (err: any) { + onReadError?.(curDir, err?.code); } return null; } @@ -2262,6 +2381,8 @@ const gateway: ChannelGatewayAdapter = { } } catch { /* ignore */ } + const unknownInbounds = reconcileTerminalCurStamps(account.mailDir, agentId, log); + // Crash recovery (at-least-once): re-dispatch cur/ records that were // promoted but never acked/nacked. cur/ is a DESTINATION, so the record // must PROVE it came through promote() (envelopeId + stored signed @@ -2278,7 +2399,7 @@ const gateway: ChannelGatewayAdapter = { const curPath = resolve(curDir, filename); if (seenFiles.has(curPath)) continue; const record = readMailFile(curPath); - if (!record || record.ackedAt || record.nackedAt) continue; + if (!record || record.ackedAt || record.nackedAt || unknownInbounds.has("*") || unknownInbounds.has(record.id)) continue; void recoverUnackedCurRecord(agentId, curPath, record); } } @@ -2341,13 +2462,14 @@ const gateway: ChannelGatewayAdapter = { // a route cannot pin a record forever. try { const retentionDays = resolveObligationRetentionDays(pluginConfig, (cfg as any)?.channels?.[CHANNEL_ID]); - sweepTerminalObligations( + if (!unknownInbounds.has("*")) sweepTerminalObligations( account.mailDir, agentId, retentionDays, log, Date.now(), resolveObligationNackHoldDays(retentionDays, pluginConfig, (cfg as any)?.channels?.[CHANNEL_ID], log), + unknownInbounds, ); } catch (err: any) { log?.warn?.(`tps-mail: obligation retention sweep failed (ignored): ${err?.message ?? String(err)}`); @@ -2358,7 +2480,7 @@ const gateway: ChannelGatewayAdapter = { // and RE-ARM the deadline where work is still outstanding. for (const rec of listObligations(account.mailDir, agentId)) { if (!isLive()) break; - if (TERMINAL_STATES.has(rec.state)) continue; + if (TERMINAL_STATES.has(rec.state) || unknownInbounds.has("*") || unknownInbounds.has(rec.inboundId)) continue; const recCurPath = findCurPath(account.mailDir, agentId, rec.inboundId); const ctx = makeYieldCtx( account.mailDir, diff --git a/plugins/openclaw-tps-mail/src/obligations.ts b/plugins/openclaw-tps-mail/src/obligations.ts index d5207db2..b41c109b 100644 --- a/plugins/openclaw-tps-mail/src/obligations.ts +++ b/plugins/openclaw-tps-mail/src/obligations.ts @@ -200,20 +200,29 @@ export function readObligationResult(mailDir: string, agent: string, inboundId: } } -export function listObligations(mailDir: string, agent: string): ObligationRecord[] { +export function listObligations(mailDir: string, agent: string, onReadError?: (path: string, code: string | undefined, ids: string[], kind: "directory" | "file") => void): ObligationRecord[] { const dir = obligationsDir(mailDir, agent); let names: string[]; try { names = readdirSync(dir); - } catch { + } catch (err: any) { + if (err?.code !== "ENOENT") onReadError?.(dir, err?.code, ["*"], "directory"); return []; } const out: ObligationRecord[] = []; for (const name of names) { if (!name.endsWith(".json") || name.startsWith(".")) continue; try { - out.push(JSON.parse(readFileSync(resolve(dir, name), "utf-8")) as ObligationRecord); - } catch { + const record = JSON.parse(readFileSync(resolve(dir, name), "utf-8")) as ObligationRecord; + if (onReadError && (!record || typeof record.inboundId !== "string" || `${record.inboundId}.json` !== name || !ALL_STATES.has(record.state))) { + const ids = [name.slice(0, -5)]; + if (typeof record?.inboundId === "string") ids.push(record.inboundId); + onReadError(resolve(dir, name), undefined, ids, "file"); + continue; + } + out.push(record); + } catch (err: any) { + onReadError?.(resolve(dir, name), err?.code, [name.slice(0, -5)], "file"); // A torn record is not readable truth; skip it rather than crash recovery. } } @@ -502,6 +511,7 @@ export function sweepTerminalObligations( log?: ObligationLog, nowMs: number = Date.now(), nackHoldDays: number = retentionDays * DEFAULT_NACK_HOLD_MULTIPLE, + unresolved = new Set(), ): RetentionResult { const res: RetentionResult = { removed: 0, @@ -576,6 +586,13 @@ export function sweepTerminalObligations( res.left++; // pending / delivering / posted / yielded are never deletable continue; } + const heldInboundId = (record as { inboundId?: unknown }).inboundId; + if (unresolved.has(name.slice(0, -5)) || (typeof heldInboundId === "string" && unresolved.has(heldInboundId))) { + const obligationId = (record as { obligationId?: unknown }).obligationId; + if (typeof obligationId === "string") liveObligationIds.add(obligationId); + res.heldForRecovery++; + continue; + } const snapshotObligationId = (record as { obligationId?: unknown }).obligationId; if (typeof snapshotObligationId === "string") terminalObligationIds.add(snapshotObligationId); // cli#389 round 11, item 1: a record still OWING its nack mail is not diff --git a/plugins/openclaw-tps-mail/test/cur-record-write.test.ts b/plugins/openclaw-tps-mail/test/cur-record-write.test.ts new file mode 100644 index 00000000..b87c4d0c --- /dev/null +++ b/plugins/openclaw-tps-mail/test/cur-record-write.test.ts @@ -0,0 +1,557 @@ +import { describe, expect, it, beforeEach, afterEach, mock, spyOn } from "bun:test"; +import * as fs from "node:fs"; +import { chmodSync, mkdirSync, mkdtempSync, readdirSync, readFileSync, rmSync, writeFileSync } from "node:fs"; +import { join, resolve } from "node:path"; +import { tmpdir } from "node:os"; +import * as ed from "@noble/ed25519"; +import { createHash } from "node:crypto"; +import { signEnvelope, type ChainEntry } from "@tpsdev-ai/agent"; +import { FileSystemTransport } from "../../../packages/cli/src/utils/transport.js"; + +import { hashes } from "@noble/ed25519"; +hashes.sha512 = (m: Uint8Array) => new Uint8Array(createHash("sha512").update(m).digest()); + +const realFs = { ...fs }; +let removeOnTerminal: { path: string; state: string } | undefined; +let replaceOnTerminal: { path: string; state: string; bytes: string } | undefined; +let stampError: Error | undefined; +let obligationReadFailure: string | undefined; +let obligationReads = 0; +let failObligationReads = Infinity; +let curReadFailure: string | undefined; +let curReads = 0; +let curReadsUntilFailure = 1; +mock.module("node:fs", () => ({ + ...realFs, + readFileSync: (...args: any[]) => { + if (args[0] === curReadFailure && ++curReads >= curReadsUntilFailure) { + throw Object.assign(new Error("injected cur read failure"), { code: "EACCES" }); + } + if (args[0] === obligationReadFailure && ++obligationReads <= failObligationReads) { + throw Object.assign(new Error("injected obligation read failure"), { code: "EACCES" }); + } + return (realFs.readFileSync as any)(...args); + }, + openSync: (...args: any[]) => { + if (stampError && String(args[0]).includes(".ack-")) throw stampError; + return (realFs.openSync as any)(...args); + }, + writeFileSync: (...args: any[]) => { + const result = (realFs.writeFileSync as any)(...args); + if (removeOnTerminal && String(args[0]).includes(".obligations/") && typeof args[1] === "string") { + if (JSON.parse(args[1]).state === removeOnTerminal.state) { + realFs.unlinkSync(removeOnTerminal.path); + removeOnTerminal = undefined; + } + } + if (replaceOnTerminal && String(args[0]).includes(".obligations/") && typeof args[1] === "string") { + if (JSON.parse(args[1]).state === replaceOnTerminal.state) { + realFs.writeFileSync(replaceOnTerminal.path, replaceOnTerminal.bytes); + replaceOnTerminal = undefined; + } + } + return result; + }, +})); + +const FLINT_SEED = Buffer.alloc(32, 0x01); +const ANVIL_SEED = Buffer.alloc(32, 0x02); +const pubkeyFromSeed = (s: Buffer): Buffer => Buffer.from(ed.getPublicKey(new Uint8Array(s))); + +import pluginModule, { setStampRetryDelaysForTests } from "../src/index.js"; + +let capturedPlugin: any; +const mockApi: any = { + registerChannel: ({ plugin }: { plugin: any }) => { capturedPlugin = plugin; }, + registerAgentEventSubscription: () => {}, + logger: { info: () => {}, warn: () => {}, error: () => {} }, +}; +pluginModule.register(mockApi); + +let mailDir: string; +let keysDir: string; +let home: string; +let origHome: string | undefined; +let origKeys: string | undefined; +let prevDelays: number[]; + +beforeEach(() => { + mailDir = mkdtempSync(join(tmpdir(), "tps-492-mail-")); + keysDir = mkdtempSync(join(tmpdir(), "tps-492-keys-")); + home = mkdtempSync(join(tmpdir(), "tps-492-home-")); + writeFileSync(join(keysDir, "anvil.key"), ANVIL_SEED); + writeFileSync(join(keysDir, "flint.key"), FLINT_SEED); + prevDelays = setStampRetryDelaysForTests([100, 100, 100]); + origHome = process.env.HOME; process.env.HOME = home; + origKeys = process.env.TPS_TEST_KEYS_DIR; process.env.TPS_TEST_KEYS_DIR = keysDir; + mock.module("@tpsdev-ai/cli/utils/mail-verify", () => ({ + createMailVerifyClient: async () => ({ + async getAgent(name: string) { + if (name === "flint") return { publicKey: pubkeyFromSeed(FLINT_SEED) }; + if (name === "anvil") return { publicKey: pubkeyFromSeed(ANVIL_SEED) }; + return null; + }, + }), + })); +}); + +afterEach(() => { + removeOnTerminal = undefined; + replaceOnTerminal = undefined; + stampError = undefined; + obligationReadFailure = undefined; + obligationReads = 0; + failObligationReads = Infinity; + curReadFailure = undefined; + curReads = 0; + curReadsUntilFailure = 1; + setStampRetryDelaysForTests(prevDelays); + if (origHome === undefined) delete process.env.HOME; else process.env.HOME = origHome; + if (origKeys === undefined) delete process.env.TPS_TEST_KEYS_DIR; else process.env.TPS_TEST_KEYS_DIR = origKeys; + for (const d of [mailDir, keysDir, home]) { try { rmSync(d, { recursive: true, force: true }); } catch { /* best effort */ } } +}); + +const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms)); +async function pollUntil(pred: () => boolean, ms = 4000): Promise { + const t = Date.now(); + while (Date.now() - t < ms) { if (pred()) return true; await sleep(10); } + return pred(); +} + +function signedBody(from: string, to: string, body: string, seed: Buffer): string { + const chain: ChainEntry[] = [ + { agent: "system", kind: "human", timestamp: new Date().toISOString(), rationale: "originates", signature: null }, + { agent: from, kind: "agent", timestamp: new Date().toISOString(), rationale: `agent ${from}`, signature: null }, + ]; + return JSON.stringify(signEnvelope( + { v: 1, from, to, body, messageId: `env-${Math.random().toString(36).slice(2, 10)}`, timestamp: new Date().toISOString(), delegationChain: chain }, + { [from]: seed }, + )); +} + +/** The anvil cur/ record (the first parseable .json file in anvil/cur). */ +function readCur(): { path: string; record: any } | null { + const dir = resolve(mailDir, "anvil", "cur"); + let names: string[]; + try { names = readdirSync(dir); } catch { return null; } + for (const n of names) { + if (!n.endsWith(".json")) continue; + const p = join(dir, n); + try { return { path: p, record: JSON.parse(readFileSync(p, "utf-8")) }; } catch { /* torn */ } + } + return null; +} +function obligation(id: string): any | null { + try { return JSON.parse(readFileSync(resolve(mailDir, "anvil", ".obligations", `${id}.json`), "utf-8")); } catch { return null; } +} + +interface Boot { + inboundId: string; + logs: string[]; + dispatch: () => any; + deliver: (text: string) => Promise; + skip: (reason?: string) => void; + settle: () => void; + stop: () => Promise; +} + +/** Start one account. `withInbound` writes a new inbound so a turn is dispatched. */ +async function boot(withInbound: boolean, useTransport = false): Promise { + mkdirSync(resolve(mailDir, "flint", "new"), { recursive: true }); + mkdirSync(resolve(mailDir, "anvil", "new"), { recursive: true }); + const inboundId = `msg-${Math.random().toString(36).slice(2, 10)}`; + if (withInbound && useTransport) { + const result = await new FileSystemTransport(() => resolve(mailDir, "anvil")).deliver({ + from: "flint", to: "anvil", + body: Buffer.from(signedBody("flint", "anvil", "inbound", FLINT_SEED)), + headers: { "x-tps-id": inboundId }, + }); + expect(result.delivered).toBe(true); + expect(result.path?.endsWith(`-${inboundId}.json`)).toBe(false); + } else if (withInbound) { + writeFileSync(resolve(mailDir, "anvil", "new", `2026-05-26T00-00-00-${inboundId}.json`), JSON.stringify({ + id: inboundId, from: "flint", to: "anvil", body: signedBody("flint", "anvil", "inbound", FLINT_SEED), + timestamp: new Date().toISOString(), headers: { "X-TPS-Trust": "agent", "X-TPS-Surface": "tps-mail" }, deliveryAttempts: 0, + }, null, 2), "utf-8"); + } + let dispatchedArgs: any = null; + let settleFn: (() => void) | null = null; + const controller = new AbortController(); + const logs: string[] = []; + const channelRuntime = { + routing: { buildAgentSessionKey: (p: any) => `agent:${p.agentId}:tps-mail:default:${p.peer.id}` }, + reply: { + finalizeInboundContext: async (c: any) => ({ ...c, CommandAuthorized: false }), + dispatchReplyWithBufferedBlockDispatcher: async (args: any) => { + dispatchedArgs = args; + await new Promise((res) => { settleFn = res; }); + return { failedCounts: 0 }; + }, + }, + }; + const cfg = { bindings: [{ agentId: "anvil", match: { channel: "tps-mail", accountId: "default" } }] }; + const ctx = { + account: { accountId: "default", mailDir, enabled: true }, + cfg, + log: { info: (m: string) => logs.push(String(m)), warn: (m: string) => logs.push(String(m)), error: (m: string) => logs.push(String(m)) }, + channelRuntime, + abortSignal: controller.signal, + }; + const startPromise = capturedPlugin.gateway.startAccount(ctx); + return { + inboundId, logs, + dispatch: () => dispatchedArgs, + deliver: async (text: string) => { await dispatchedArgs.dispatcherOptions.deliver({ text }, { kind: "final" }); }, + skip: (reason = "empty") => dispatchedArgs.dispatcherOptions.onSkip?.({ text: "" }, { kind: "final", reason }), + settle: () => settleFn?.(), + stop: async () => { controller.abort(); try { await startPromise; } catch { /* aborted */ } }, + }; +} + +const failedLogs = (h: Boot, tag: string) => h.logs.filter((m) => m.includes(tag) && m.includes(h.inboundId)); + +/** Make the cur/ record and directory unwritable so the stamp write fails with EACCES. */ +function breakCur(path: string): () => void { + const curDir = resolve(mailDir, "anvil", "cur"); + chmodSync(path, 0o444); + chmodSync(curDir, 0o555); + return () => { chmodSync(curDir, 0o755); chmodSync(path, 0o644); }; +} + +describe("cli#492 — a failed cur/ stamp write is surfaced and retried", () => { + for (const kind of ["ack", "nack"] as const) { + const stampKey = kind === "ack" ? "ackedAt" : "nackedAt"; + const state = kind === "ack" ? "acked" : "failed"; + const finish = (h: Boot) => { + if (kind === "ack") return h.deliver("verdict").then(() => h.settle()); + h.skip("empty"); // an empty/silent final → the named failure path + h.settle(); + return Promise.resolve(); + }; + + it(`${kind}: startup stamps a promoted FileSystemTransport record after a failed live stamp`, async () => { + const h = await boot(true, true); + try { + expect(await pollUntil(() => h.dispatch() !== null)).toBe(true); + const cur = readCur()!; + expect(cur.record.id).toBe(h.inboundId); + expect(cur.path.endsWith(`-${h.inboundId}.json`)).toBe(false); + expect(cur.record.envelopeId).toBeDefined(); + stampError = Object.assign(new Error("injected stamp failure"), { code: "EACCES" }); + await finish(h); + expect(await pollUntil(() => failedLogs(h, `${kind}-stamp-failed`).length > 0)).toBe(true); + expect(obligation(h.inboundId)?.state).toBe(state); + expect(readCur()?.record?.[stampKey]).toBeUndefined(); + } finally { + await h.stop(); + stampError = undefined; + } + const next = await boot(false); + try { + expect(await pollUntil(() => !!readCur()?.record?.[stampKey])).toBe(true); + expect(next.logs.some((m) => m.includes(`reconciled the ${kind} stamp for ${h.inboundId}`))).toBe(true); + expect(obligation(h.inboundId)?.state).toBe(state); + } finally { + await next.stop(); + } + }); + + it(`${kind}: a transient failed ${stampKey} write is logged by id/path/code and fixed by the in-process retry, no restart`, async () => { + const h = await boot(true); + expect(await pollUntil(() => h.dispatch() !== null, 4000), "dispatch started").toBe(true); + const cur = readCur(); + expect(cur, "the inbound was promoted to cur/").not.toBeNull(); + const restore = breakCur(cur!.path); + let restored = false; + try { + await finish(h); + expect(await pollUntil(() => obligation(h.inboundId)?.state === state, 4000), "the terminal transition is durable").toBe(true); + expect(readCur()?.record?.[stampKey], "the stamp did NOT land").toBeUndefined(); + expect( + await pollUntil( + () => failedLogs(h, `${kind}-stamp-failed`).some((m) => m.includes(`path=${readCur()!.path}`) && m.includes("EACCES")), + 4000, + ), + "the failed write is logged by id, path and code", + ).toBe(true); + } finally { + restore(); + restored = true; + } + expect(restored).toBe(true); + expect(await pollUntil(() => !!readCur()?.record?.[stampKey], 4000), "the in-process retry stamped the record").toBe(true); + expect(obligation(h.inboundId)?.state).toBe(state); + await h.stop(); + }, 20000); + + it(`${kind}: persistent failures exhaust live retries`, async () => { + const h = await boot(true); + expect(await pollUntil(() => h.dispatch() !== null, 4000), "dispatch started").toBe(true); + const cur = readCur(); + expect(cur, "the inbound was promoted to cur/").not.toBeNull(); + const restore = breakCur(cur!.path); + try { + await finish(h); + // initial attempt + 3 retries = 4 logged failures, then no more. + expect(await pollUntil(() => failedLogs(h, `${kind}-stamp-failed`).length >= 4, 4000), "initial attempt + 3 retries").toBe(true); + await sleep(400); + const logged = failedLogs(h, `${kind}-stamp-failed`); + expect(logged.length, "no attempt beyond the bound").toBe(4); + expect(logged[3]).toContain("no retries left"); + expect(readCur()?.record?.[stampKey]).toBeUndefined(); + } finally { + restore(); + } + await sleep(300); + expect(readCur()?.record?.[stampKey], "nothing retries after the bound").toBeUndefined(); + await h.stop(); + + // NEXT ACCOUNT START: the durable terminal state is re-stamped onto the record. + const h2 = await boot(false); + expect(await pollUntil(() => !!readCur()?.record?.[stampKey], 4000), "re-stamped at the next account start").toBe(true); + if (kind === "ack") expect(readCur()?.record?.read).toBe(true); + expect(obligation(h.inboundId)?.state).toBe(state); + await h2.stop(); + }, 20000); + + for (const stage of ["lock", "scratch"] as const) { + it(`${kind}: ${stage} acquisition reports the failing path in live and startup stamps`, async () => { + const h = await boot(true); + expect(await pollUntil(() => h.dispatch() !== null)).toBe(true); + const cur = readCur()!; + const blocked = stage === "lock" ? resolve(mailDir, "anvil") : resolve(mailDir, "anvil", "cur"); + const restore = () => realFs.chmodSync(blocked, 0o755); + realFs.chmodSync(blocked, 0o555); + try { + await finish(h); + expect(await pollUntil(() => failedLogs(h, `${kind}-stamp-failed`).length > 0)).toBe(true); + const diagnostic = failedLogs(h, `${kind}-stamp-failed`)[0]; + expect(diagnostic).toContain(stage === "lock" ? `path=${blocked}/.mail-lock code=EACCES` : `path=${cur.path} code=EACCES`); + if (stage === "lock") expect(diagnostic).not.toContain(`path=${cur.path}`); + expect(readCur()?.record?.[stampKey]).toBeUndefined(); + } finally { + await h.stop(); + restore(); + } + realFs.chmodSync(blocked, 0o555); + const next = await boot(false); + try { + expect(await pollUntil(() => next.logs.some((m) => m.includes(`${kind}-stamp-reconcile-failed`)))).toBe(true); + const diagnostic = next.logs.find((m) => m.includes(`${kind}-stamp-reconcile-failed`))!; + expect(diagnostic).toContain(stage === "lock" ? `path=${blocked}/.mail-lock code=EACCES` : `path=${cur.path} code=EACCES`); + if (stage === "lock") expect(diagnostic).not.toContain(`path=${cur.path}`); + expect(diagnostic).toContain("obligation retained"); + expect(readCur()?.record?.[stampKey]).toBeUndefined(); + } finally { + await next.stop(); + restore(); + } + }); + } + + it(`${kind}: a missing record is reported after the durable transition`, async () => { + const h = await boot(true); + try { + expect(await pollUntil(() => h.dispatch() !== null)).toBe(true); + const cur = readCur()!; + removeOnTerminal = { path: cur.path, state }; + await finish(h); + expect(await pollUntil(() => failedLogs(h, `${kind}-stamp-failed`).length > 0)).toBe(true); + expect(failedLogs(h, `${kind}-stamp-failed`)[0]).toContain(`actor=anvil path=${cur.path} code=ENOENT`); + expect(failedLogs(h, `${kind}-stamp-failed`)[0]).toContain("obligation retained; resolve the failure and restart the account"); + expect(obligation(h.inboundId)?.state).toBe(state); + await sleep(400); + expect(failedLogs(h, `${kind}-stamp-failed`).length).toBe(1); + expect(realFs.existsSync(cur.path)).toBe(false); + } finally { + await h.stop(); + } + }, 15000); + + for (const alreadyStamped of [false, true]) { + it(`${kind}: a replacement identity is refused on the live stamp path (already stamped: ${alreadyStamped})`, async () => { + const h = await boot(true); + try { + expect(await pollUntil(() => h.dispatch() !== null)).toBe(true); + const cur = readCur()!; + const replacement = { ...cur.record, id: "replacement", ...(alreadyStamped ? { ackedAt: "existing-ack", nackedAt: "existing-nack" } : {}) }; + const bytes = JSON.stringify(replacement); + replaceOnTerminal = { path: cur.path, state, bytes }; + await finish(h); + expect(await pollUntil(() => failedLogs(h, `${kind}-stamp-failed`).some((m) => m.includes("code=ID_MISMATCH")))).toBe(true); + expect(obligation(h.inboundId)?.state).toBe(state); + await sleep(500); + expect(realFs.readFileSync(cur.path, "utf8")).toBe(bytes); + expect(failedLogs(h, `${kind}-stamp-failed`)).toHaveLength(4); + expect(h.logs.some((m) => m.includes(`${kind}-stamp-retry-ok`))).toBe(false); + } finally { await h.stop(); } + }, 15000); + + } + + it(`${kind}: an exception without a code is reported and remains retryable`, async () => { + const h = await boot(true); + try { + expect(await pollUntil(() => h.dispatch() !== null)).toBe(true); + stampError = new Error("injected stamp failure"); + await finish(h); + expect(await pollUntil(() => failedLogs(h, `${kind}-stamp-failed`).length > 0)).toBe(true); + const diagnostic = failedLogs(h, `${kind}-stamp-failed`)[0]; + expect(diagnostic).toContain(`actor=anvil`); + expect(diagnostic).not.toContain("code="); + expect(diagnostic).toContain("obligation retained"); + expect(diagnostic).not.toContain("no retries left"); + stampError = undefined; + expect(await pollUntil(() => !!readCur()?.record?.[stampKey])).toBe(true); + expect(obligation(h.inboundId)?.state).toBe(state); + } finally { + stampError = undefined; + await h.stop(); + } + }, 15000); + + it(`${kind}: stopping the account cancels the pending stamp retries`, async () => { + setStampRetryDelaysForTests([731, 731, 731]); + const realSetTimeout = globalThis.setTimeout; + let scheduled = 0; + let fired = 0; + const timerSpy = spyOn(globalThis, "setTimeout").mockImplementation(((fn: any, delay: number, ...args: any[]) => { + if (delay !== 731) return realSetTimeout(fn, delay, ...args); + scheduled++; + return realSetTimeout(() => { fired++; fn(...args); }, delay); + }) as typeof setTimeout); + const h = await boot(true); + expect(await pollUntil(() => h.dispatch() !== null, 4000), "dispatch started").toBe(true); + const cur = readCur(); + expect(cur, "the inbound was promoted to cur/").not.toBeNull(); + const restore = breakCur(cur!.path); + try { + await finish(h); + expect(await pollUntil(() => failedLogs(h, `${kind}-stamp-failed`).length >= 1, 4000), "first failure logged").toBe(true); + expect(scheduled).toBe(1); + await h.stop(); + } finally { + restore(); + } + const before = failedLogs(h, `${kind}-stamp-failed`).length; + try { + await sleep(1200); + expect(fired, "the pending timer callback never runs after stop").toBe(0); + } finally { + timerSpy.mockRestore(); + } + expect(failedLogs(h, `${kind}-stamp-failed`).length, "no retry ran after stop").toBe(before); + expect(readCur()?.record?.[stampKey], "the cancelled retry never stamped").toBeUndefined(); + }, 20000); + } + + it("the successful write path is unchanged: a normal ack stamps ackedAt with no diagnostic", async () => { + const h = await boot(true); + expect(await pollUntil(() => h.dispatch() !== null, 4000), "dispatch started").toBe(true); + await h.deliver("verdict"); + h.settle(); + expect(await pollUntil(() => !!readCur()?.record?.ackedAt, 4000)).toBe(true); + expect(readCur()?.record?.read).toBe(true); + expect(h.logs.some((m) => m.includes("ack-stamp-failed"))).toBe(false); + await h.stop(); + }, 15000); +}); + +for (const state of ["acked", "failed"] as const) { + for (const malformed of [true, false]) { + it(`startup retention ${malformed ? "keeps an aged terminal obligation missing inboundId and its receipt" : "sweeps an aged resolved obligation and its receipt"} (${state})`, async () => { + const first = await boot(true, true); + try { + expect(await pollUntil(() => first.dispatch() !== null)).toBe(true); + if (state === "acked") await first.deliver("verdict"); + else first.skip("empty"); + first.settle(); + expect(await pollUntil(() => !!readCur()?.record?.[state === "acked" ? "ackedAt" : "nackedAt"])).toBe(true); + } finally { await first.stop(); } + const cur = readCur()!; + const path = resolve(mailDir, "anvil", ".obligations", `${first.inboundId}.json`); + const terminal = JSON.parse(realFs.readFileSync(path, "utf8")); + expect(terminal.state).toBe(state); + terminal.lastTransitionAt = new Date(0).toISOString(); + if (malformed) { + delete terminal.inboundId; + delete cur.record.ackedAt; + delete cur.record.nackedAt; + realFs.writeFileSync(cur.path, JSON.stringify(cur.record)); + } + realFs.writeFileSync(path, JSON.stringify(terminal)); + const receiptDir = resolve(mailDir, "anvil", ".obligations", "receipts"); + realFs.mkdirSync(receiptDir, { recursive: true }); + const receiptPath = resolve(receiptDir, `${terminal.obligationId}.json`); + const receipt = realFs.existsSync(receiptPath) + ? JSON.parse(realFs.readFileSync(receiptPath, "utf8")) + : { obligationId: terminal.obligationId }; + receipt.ts = new Date(0).toISOString(); + realFs.writeFileSync(receiptPath, JSON.stringify(receipt)); + const obligationBytes = realFs.readFileSync(path, "utf8"); + const receiptBytes = realFs.readFileSync(receiptPath, "utf8"); + const curBytes = realFs.readFileSync(cur.path, "utf8"); + const next = await boot(false); + try { + expect(await pollUntil(() => next.logs.some((m) => m.includes("obligation retention: removed")))).toBe(true); + expect(next.dispatch()).toBeNull(); + expect(realFs.readFileSync(cur.path, "utf8")).toBe(curBytes); + if (malformed) { + expect(realFs.readFileSync(path, "utf8")).toBe(obligationBytes); + expect(realFs.readFileSync(receiptPath, "utf8")).toBe(receiptBytes); + expect(next.logs.some((m) => m.includes("held 1 for unresolved cur/ recovery"))).toBe(true); + } else { + expect(realFs.existsSync(path)).toBe(false); + expect(realFs.existsSync(receiptPath)).toBe(false); + } + } finally { await next.stop(); } + }, 15000); + } + + for (const stage of ["lookup-read", "reread", "locked-read", "stamp-write"] as const) { + it(`startup ${state} ${stage} holds aged obligations across restarts without redispatch`, async () => { + const first = await boot(true); + try { + expect(await pollUntil(() => first.dispatch() !== null)).toBe(true); + if (state === "acked") await first.deliver("verdict"); + else first.skip("empty"); + first.settle(); + expect(await pollUntil(() => obligation(first.inboundId)?.state === state)).toBe(true); + expect(await pollUntil(() => !!readCur()?.record?.[state === "acked" ? "ackedAt" : "nackedAt"])).toBe(true); + } finally { await first.stop(); } + const cur = readCur()!; + delete cur.record.ackedAt; + delete cur.record.nackedAt; + realFs.writeFileSync(cur.path, JSON.stringify(cur.record)); + const path = resolve(mailDir, "anvil", ".obligations", `${first.inboundId}.json`); + const terminal = JSON.parse(realFs.readFileSync(path, "utf8")); + terminal.lastTransitionAt = new Date(0).toISOString(); + realFs.writeFileSync(path, JSON.stringify(terminal)); + const bytes = realFs.readFileSync(path, "utf8"); + for (let restart = 0; restart < 2; restart++) { + if (stage === "stamp-write") stampError = new Error("injected stamp failure"); + else { + curReadFailure = cur.path; + curReads = 0; + curReadsUntilFailure = stage === "lookup-read" ? 1 : stage === "reread" ? 2 : 3; + } + const next = await boot(false); + try { + expect(await pollUntil(() => next.logs.some((m) => m.includes(`actor=anvil`) && m.includes(first.inboundId)))).toBe(true); + await sleep(100); + expect(next.dispatch()).toBeNull(); + expect(next.logs.filter((m) => m.includes("actor=anvil") && m.includes(first.inboundId))).toHaveLength(1); + expect(realFs.readFileSync(path, "utf8")).toBe(bytes); + expect(JSON.parse(realFs.readFileSync(cur.path, "utf8"))[state === "acked" ? "ackedAt" : "nackedAt"]).toBeUndefined(); + } finally { await next.stop(); } + } + stampError = undefined; + curReadFailure = undefined; + const repaired = await boot(false); + try { + expect(await pollUntil(() => !!readCur()?.record?.[state === "acked" ? "ackedAt" : "nackedAt"])).toBe(true); + expect(repaired.dispatch()).toBeNull(); + } finally { await repaired.stop(); } + }, 15000); + } +} diff --git a/plugins/openclaw-tps-mail/test/diagnostics.test.ts b/plugins/openclaw-tps-mail/test/diagnostics.test.ts new file mode 100644 index 00000000..0d446a8e --- /dev/null +++ b/plugins/openclaw-tps-mail/test/diagnostics.test.ts @@ -0,0 +1,26 @@ +import { expect, test } from "bun:test"; +import { formatStampDiagnostic } from "../src/diagnostics.js"; + +for (const obligation of ["retained", "none", "unknown"] as const) { + for (const retriesExhausted of [false, true]) { + test(`${obligation}, exhausted=${retriesExhausted}`, () => { + expect(formatStampDiagnostic({ kind: "ack-stamp-failed", actor: "anvil", id: "inbound", + path: "/mail/anvil/cur/inbound.json", code: "EACCES", obligation, retriesExhausted })).toBe( + "tps-mail: ack-stamp-failed: inbound actor=anvil path=/mail/anvil/cur/inbound.json code=EACCES; " + + (obligation === "retained" ? "obligation retained; " : obligation === "unknown" ? "state unknown; " : "") + + (retriesExhausted ? "no retries left; " : "") + "resolve the failure and restart the account"); + }); + } +} + +test("optional id and retry flag", () => { + expect(formatStampDiagnostic({ kind: "stamp-reconcile-read-failed", actor: "anvil", path: "/mail/anvil/cur/inbound.json", + code: "EACCES", obligation: "unknown" })).toBe("tps-mail: stamp-reconcile-read-failed: actor=anvil path=/mail/anvil/cur/inbound.json code=EACCES; state unknown; resolve the failure and restart the account"); +}); + +test("an uncoded stamp failure omits code", () => { + const message = formatStampDiagnostic({ kind: "ack-stamp-failed", actor: "anvil", + path: "/mail/anvil/cur/inbound.json", obligation: "unknown" }); + expect(message).not.toContain("code="); + expect(message).toContain("state unknown"); +}); diff --git a/plugins/openclaw-tps-mail/test/obligation-retention.test.ts b/plugins/openclaw-tps-mail/test/obligation-retention.test.ts index 439ea7eb..e59dc7ff 100644 --- a/plugins/openclaw-tps-mail/test/obligation-retention.test.ts +++ b/plugins/openclaw-tps-mail/test/obligation-retention.test.ts @@ -1,10 +1,8 @@ /** * obligation-retention.test.ts — cli#401: the obligation-record retention policy. * - * (a) a terminal record older than N days → removed at startup; a younger one kept. * (b) pending / posted / yielded at ANY age → never removed. * (c) a malformed record older than N → kept + logged once. - * (d) the config key changes N (N=1 removes a 2-day-old terminal record the default would keep). * (e) replay after a sweep: a replayed inbound id whose record was swept opens a FRESH obligation. * * Terminal records only; aged by the record's OWN `lastTransitionAt` (falling @@ -70,8 +68,13 @@ function writeRecordFor( writeFileSync(join(dir, `${inboundId}.json`), JSON.stringify(rec, null, 2), "utf-8"); } -function writeRecord(id: string, state: string, lastTransitionAt: string | null): void { +function writeRecord(id: string, state: string, lastTransitionAt: string | null, stampedCur = false): void { writeRecordFor(AGENT, id, `ob-${id}`, state, lastTransitionAt); + if (stampedCur) { + const curDir = join(mailDir, AGENT, "cur"); + mkdirSync(curDir, { recursive: true }); + writeFileSync(join(curDir, `${id}.json`), JSON.stringify({ id, [state === "acked" ? "ackedAt" : "nackedAt"]: daysAgo(1) })); + } } /** The receipts dir for an agent — inside that agent's own obligation store. */ @@ -127,10 +130,10 @@ afterEach(() => { }); describe("cli#401 — obligation retention", () => { - it("(a) at startup: a terminal record older than N days is REMOVED; a younger one is KEPT", async () => { - writeRecord("old-acked", "acked", daysAgo(10)); - writeRecord("young-acked", "acked", daysAgo(1)); - writeRecord("old-failed", "failed", daysAgo(10)); + it("(a) startup sweeps aged terminal obligations with stamped cur records", async () => { + writeRecord("old-acked", "acked", daysAgo(10), true); + writeRecord("young-acked", "acked", daysAgo(1), true); + writeRecord("old-failed", "failed", daysAgo(10), true); await runStartup({}, () => !existsSync(obligationPath(mailDir, AGENT, "old-acked")) && !existsSync(obligationPath(mailDir, AGENT, "old-failed"))); @@ -175,7 +178,7 @@ describe("cli#401 — obligation retention", () => { expect(existsSync(p)).toBe(true); }); - it("(d) the config key changes N — a 2-day-old terminal record the default keeps is removed at N=1", async () => { + it("(d) startup uses the configured retention for a stamped terminal record", async () => { // resolution: the plugin config key wins; unset/invalid falls back to 7 expect(resolveObligationRetentionDays({ obligationRetentionDays: 1 }, undefined)).toBe(1); expect(resolveObligationRetentionDays({ obligationRetentionDays: "3" }, undefined)).toBe(3); @@ -183,7 +186,7 @@ describe("cli#401 — obligation retention", () => { expect(resolveObligationRetentionDays({}, undefined)).toBe(7); expect(resolveObligationRetentionDays({ obligationRetentionDays: "nope" }, undefined)).toBe(7); - writeRecord("two-days", "acked", daysAgo(2)); + writeRecord("two-days", "acked", daysAgo(2), true); // default (7): kept sweepTerminalObligations(mailDir, AGENT, 7, { info: () => {}, warn: () => {} }); expect(existsSync(obligationPath(mailDir, AGENT, "two-days")), "default keeps a 2-day-old record").toBe(true); @@ -227,7 +230,7 @@ describe("cli#401 — obligation retention", () => { // OpenClaw passes exactly `plugins.entries[id].config` to the plugin as api.pluginConfig: const receivedPluginConfig = openclawConfig.plugins.entries["openclaw-tps-mail"].config; expect(resolveObligationRetentionDays(receivedPluginConfig, undefined)).toBe(1); - writeRecord("doc-two-days", "acked", daysAgo(2)); + writeRecord("doc-two-days", "acked", daysAgo(2), true); await runStartup(receivedPluginConfig, () => !existsSync(obligationPath(mailDir, AGENT, "doc-two-days"))); expect(existsSync(obligationPath(mailDir, AGENT, "doc-two-days")), "the documented path must drive the sweep").toBe(false); }); diff --git a/plugins/openclaw-tps-mail/test/patch-mail-file.test.ts b/plugins/openclaw-tps-mail/test/patch-mail-file.test.ts index a39238bb..de179ae1 100644 --- a/plugins/openclaw-tps-mail/test/patch-mail-file.test.ts +++ b/plugins/openclaw-tps-mail/test/patch-mail-file.test.ts @@ -5,11 +5,42 @@ import { join } from "node:path"; const realFs = { ...fs }; let removeAfterRead: string | undefined; +let readFailure: string | undefined; +let readsUntilFailure = 1; +let listFailure: string | undefined; +let stampOnRead: { path: string; count: number } | undefined; +let readsBeforeRemoval = 1; +let uncodedWriteFailure = false; +let scratchErrorCode: string | undefined; +let renameFailure: string | undefined; +let replaceIdentity = false; mock.module("node:fs", () => ({ ...realFs, + readdirSync: (...args: any[]) => { + if (args[0] === listFailure) throw Object.assign(new Error("injected list failure"), { code: "EACCES" }); + return (realFs.readdirSync as any)(...args); + }, + renameSync: (...args: any[]) => { + if (args[1] === renameFailure) throw Object.assign(new Error("injected rename failure"), { code: "EACCES" }); + return (realFs.renameSync as any)(...args); + }, + openSync: (...args: any[]) => { + if (uncodedWriteFailure && String(args[0]).includes(".ack-")) throw new Error("injected write failure"); + if (scratchErrorCode && String(args[0]).includes(".ack-")) throw Object.assign(new Error("injected scratch failure"), { code: scratchErrorCode, path: args[0] }); + return (realFs.openSync as any)(...args); + }, readFileSync: (...args: any[]) => { + if (args[0] === readFailure && --readsUntilFailure <= 0) throw Object.assign(new Error("injected read failure"), { code: "EACCES" }); const result = (realFs.readFileSync as any)(...args); - if (args[0] === removeAfterRead) { + if (stampOnRead?.path === args[0] && --stampOnRead.count === 0) { + const record = JSON.parse(String(result)); + if (replaceIdentity) record.id = "replacement"; + record.ackedAt = "concurrent-ack"; + record.nackedAt = "concurrent-nack"; + realFs.writeFileSync(stampOnRead.path, JSON.stringify(record)); + stampOnRead = undefined; + } + if (args[0] === removeAfterRead && --readsBeforeRemoval === 0) { realFs.unlinkSync(removeAfterRead!); removeAfterRead = undefined; } @@ -17,34 +48,248 @@ mock.module("node:fs", () => ({ }, })); -const { patchMailFile } = await import("../src/index.js"); +const { acquireMailLockSync, mailLockPath } = await import("@tpsdev-ai/agent"); +const { patchMailFile, reconcileTerminalCurStamps } = await import("../src/index.js"); +const { sweepTerminalObligations } = await import("../src/obligations.js"); const root = realFs.mkdtempSync(join(tmpdir(), "patch-mail-")); const path = join(root, "record.json"); afterEach(() => { removeAfterRead = undefined; + readsBeforeRemoval = readsUntilFailure = 1; + readFailure = listFailure = undefined; + stampOnRead = undefined; + uncodedWriteFailure = replaceIdentity = false; + scratchErrorCode = undefined; + renameFailure = undefined; + realFs.rmSync(join(root, "anvil"), { recursive: true, force: true }); realFs.rmSync(path, { force: true }); }); afterAll(() => realFs.rmSync(root, { recursive: true, force: true })); describe("patchMailFile", () => { + test("a nested lock failure names the lock path", () => { + const f = terminalFixture("acked"); + const lock = acquireMailLockSync(join(root, "anvil"))!; + try { + expect(patchMailFile(f.curPath, { ackedAt: "done" }, "ackedAt", "inbound")).toMatchObject({ ok: false, reason: "write-failed", path: mailLockPath(join(root, "anvil")) }); + } finally { + lock.release(); + } + }); + + test("a lock timeout names the lock path", () => { + const f = terminalFixture("acked"); + const lockPath = mailLockPath(join(root, "anvil")); + realFs.mkdirSync(lockPath); + expect(patchMailFile(f.curPath, { ackedAt: "done" }, "ackedAt", "inbound")).toMatchObject({ ok: false, reason: "write-failed", path: lockPath }); + }); + + for (const code of ["EACCES", "ENOENT"]) { + test(`a scratch ${code} failure names the target path and preserves the record`, () => { + const f = terminalFixture("acked"); + scratchErrorCode = code; + const result = patchMailFile(f.curPath, { ackedAt: "done" }, "ackedAt", "inbound"); + expect(result).toMatchObject({ ok: false, reason: "write-failed", code }); + if (result.ok) throw new Error("expected scratch failure"); + expect(result.path).toBe(f.curPath); + f.reconcile(); + expect(f.logs[0]).toBe(`tps-mail: ack-stamp-reconcile-failed: inbound actor=anvil path=${f.curPath} code=${code}; obligation retained; resolve the failure and restart the account`); + expect(JSON.parse(realFs.readFileSync(f.curPath, "utf8")).ackedAt).toBeUndefined(); + }); + } + test("patches an existing record", () => { realFs.writeFileSync(path, JSON.stringify({ id: "inbound", body: "hello" })); - expect(patchMailFile(path, { read: true })).toBe(true); + expect(patchMailFile(path, { read: true }, undefined, "inbound")).toEqual({ ok: true }); expect(JSON.parse(realFs.readFileSync(path, "utf-8"))).toEqual({ id: "inbound", body: "hello", read: true }); }); test("reports a missing record without creating it", () => { - expect(patchMailFile(path, { read: true })).toBe(false); + expect(patchMailFile(path, { read: true }, undefined, "inbound")).toMatchObject({ ok: false, reason: "record-missing" }); expect(realFs.existsSync(path)).toBe(false); }); test("reports a record removed after the read without recreating it", () => { realFs.writeFileSync(path, JSON.stringify({ id: "inbound", body: "hello" })); removeAfterRead = path; - expect(patchMailFile(path, { ackedAt: "receipt", read: true })).toBe(false); + expect(patchMailFile(path, { ackedAt: "receipt", read: true }, "ackedAt", "inbound")).toMatchObject({ ok: false, reason: "record-missing" }); expect(removeAfterRead).toBeUndefined(); expect(realFs.existsSync(path)).toBe(false); }); }); + + +function terminalFixture(state: "acked" | "failed") { + const obligationDir = join(root, "anvil", ".obligations"); + const curDir = join(root, "anvil", "cur"); + realFs.mkdirSync(obligationDir, { recursive: true }); + realFs.mkdirSync(curDir, { recursive: true }); + const obligationPath = join(obligationDir, "inbound.json"); + const curPath = join(curDir, "timestamp-inbound.json"); + realFs.writeFileSync(obligationPath, JSON.stringify({ obligationId: "ob-inbound", inboundId: "inbound", state, failure: "empty", lastTransitionAt: new Date(0).toISOString() })); + realFs.writeFileSync(curPath, JSON.stringify({ id: "inbound", body: "hello" })); + const logs: string[] = []; + const reconcile = () => reconcileTerminalCurStamps(root, "anvil", { warn: (message: string) => logs.push(message) }); + return { obligationDir, obligationPath, curDir, curPath, logs, reconcile }; +} + +describe("terminal stamp reconciliation", () => { + for (const state of ["acked", "failed"] as const) { + const kind = state === "acked" ? "ack" : "nack"; + const key = state === "acked" ? "ackedAt" : "nackedAt"; + + test(`${kind}: a record disappearing during the locked write is reported`, () => { + const f = terminalFixture(state); + removeAfterRead = f.curPath; + readsBeforeRemoval = 3; + f.reconcile(); + expect(removeAfterRead).toBeUndefined(); + expect(f.logs).toHaveLength(1); + expect(f.logs[0]).toContain(`${kind}-stamp-reconcile-failed: inbound actor=anvil path=${f.curPath} code=ENOENT`); + expect(f.logs[0]).toContain("obligation retained; resolve the failure and restart the account"); + expect(realFs.existsSync(f.curPath)).toBe(false); + expect(JSON.parse(realFs.readFileSync(f.obligationPath, "utf8")).state).toBe(state); + }); + + test(`${kind}: an uncoded write failure is reported and the next run can stamp`, () => { + const f = terminalFixture(state); + uncodedWriteFailure = true; + f.reconcile(); + expect(f.logs[0]).toContain(`actor=anvil path=${f.curPath}`); + expect(f.logs[0]).toContain("obligation retained; resolve the failure and restart the account"); + uncodedWriteFailure = false; + f.reconcile(); + expect(JSON.parse(realFs.readFileSync(f.curPath, "utf8"))[key]).toBeDefined(); + }); + + test(`${kind}: a second successful reconciliation preserves the stamp`, () => { + const f = terminalFixture(state); + f.reconcile(); + const bytes = realFs.readFileSync(f.curPath, "utf8"); + const stat = realFs.statSync(f.curPath); + f.reconcile(); + expect(realFs.readFileSync(f.curPath, "utf8")).toBe(bytes); + expect(realFs.statSync(f.curPath).mtimeMs).toBe(stat.mtimeMs); + expect(realFs.statSync(f.curPath).ino).toBe(stat.ino); + expect(f.logs).toEqual([]); + }); + + test(`${kind}: a concurrent stamp between the precheck and lock is preserved`, () => { + const f = terminalFixture(state); + stampOnRead = { path: f.curPath, count: 2 }; + f.reconcile(); + expect(stampOnRead).toBeUndefined(); + expect(JSON.parse(realFs.readFileSync(f.curPath, "utf8"))[key]).toBe(`concurrent-${kind}`); + }); + } + + for (const stage of ["cur-read", "cur-reread"] as const) { + test(`${stage}: a read failure is diagnosed and remains recoverable`, () => { + const f = terminalFixture("acked"); + const target = f.curPath; + readFailure = target; + if (stage === "cur-reread") readsUntilFailure = 2; + f.reconcile(); + expect(f.logs[0]).toContain("stamp-reconcile-read-failed"); + expect(f.logs[0]).toContain("actor=anvil"); + expect(f.logs[0]).toContain(`path=${target} code=EACCES`); + expect(f.logs[0]).toContain("resolve the failure and restart the account"); + listFailure = readFailure = undefined; + f.reconcile(); + expect(JSON.parse(realFs.readFileSync(f.curPath, "utf8")).ackedAt).toBeDefined(); + }); + } + + test("an absent obligation directory has no read failure", () => { + const logs: string[] = []; + reconcileTerminalCurStamps(root, "anvil", { warn: (m: string) => logs.push(m) }); + expect(logs).toEqual([]); + }); + +}); + +for (const state of ["acked", "failed"] as const) { + for (const stage of ["cur-lookup-read", "cur-reread", "locked-read", "stamp-write", "stamp-rename", "identity-change"] as const) { + test(`${state}: ${stage} holds an aged timestamped inbound and continues healthy records`, () => { + const f = terminalFixture(state); + if (stage.endsWith("read")) { + readFailure = f.curPath; + readsUntilFailure = stage === "cur-reread" ? 2 : stage === "locked-read" ? 3 : 1; + } + if (stage === "stamp-write") uncodedWriteFailure = true; + if (stage === "stamp-rename") renameFailure = f.curPath; + if (stage === "identity-change") { + replaceIdentity = true; + stampOnRead = { path: f.curPath, count: 2 }; + } + const bytes = realFs.readFileSync(f.obligationPath, "utf8"); + const unresolved = f.reconcile(); + expect(unresolved.has("inbound")).toBe(true); + const receiptDir = join(f.obligationDir, "receipts"); + realFs.mkdirSync(receiptDir); + const receiptPath = join(receiptDir, "ob-inbound.json"); + realFs.writeFileSync(receiptPath, JSON.stringify({ obligationId: "ob-inbound", ts: new Date(0).toISOString() })); + realFs.writeFileSync(join(f.obligationDir, "healthy.json"), JSON.stringify({ inboundId: "healthy", obligationId: "ob-healthy", state: "acked", lastTransitionAt: new Date(0).toISOString() })); + listFailure = readFailure = renameFailure = undefined; + uncodedWriteFailure = false; + realFs.writeFileSync(join(f.curDir, "timestamp-healthy.json"), JSON.stringify({ id: "healthy", ackedAt: "done" })); + sweepTerminalObligations(root, "anvil", 7, { warn: (m) => f.logs.push(m) }, Date.now(), 28, unresolved); + expect(realFs.readFileSync(f.obligationPath, "utf8")).toBe(bytes); + expect(realFs.existsSync(receiptPath)).toBe(true); + expect(realFs.existsSync(join(f.obligationDir, "healthy.json"))).toBe(false); + expect(f.logs.filter((m) => m.includes("actor=anvil"))).toHaveLength(1); + expect(f.logs.find((m) => m.includes("actor=anvil"))).toContain(`actor=anvil`); + expect(f.logs.find((m) => m.includes("actor=anvil"))).toContain("restart the account"); + replaceIdentity = false; + if (stage === "identity-change") realFs.writeFileSync(f.curPath, JSON.stringify({ id: "inbound" })); + expect(f.reconcile().size).toBe(0); + expect(JSON.parse(realFs.readFileSync(f.curPath, "utf8"))[state === "acked" ? "ackedAt" : "nackedAt"]).toBeDefined(); + }); + } +} + +for (const presence of ["missing", "unreadable"] as const) { + test(`stamp reconciliation reports ${presence} obligation state`, () => { + const f = terminalFixture("acked"); + const record = JSON.parse(realFs.readFileSync(f.obligationPath, "utf8")); + if (presence === "missing") realFs.unlinkSync(f.obligationPath); + else readFailure = f.obligationPath; + scratchErrorCode = "EACCES"; + reconcileTerminalCurStamps(root, "anvil", { warn: (message: string) => f.logs.push(message) }, [record]); + expect(f.logs[0]).toContain(`ack-stamp-reconcile-failed: inbound actor=anvil path=${f.curPath} code=EACCES`); + expect(f.logs[0]).not.toContain("obligation retained"); + if (presence === "unreadable") expect(f.logs[0]).toContain("state unknown"); + }); +} + +test("stamp reconciliation holds an obligation when an unreadable record cannot be attributed", () => { + const f = terminalFixture("acked"); + realFs.unlinkSync(f.curPath); + const unrelated = join(f.curDir, "timestamp-other.json"); + realFs.writeFileSync(unrelated, JSON.stringify({ id: "other" })); + readFailure = unrelated; + expect(f.reconcile().has("inbound")).toBe(true); + expect(f.logs[0]).toContain(`path=${unrelated} code=EACCES`); + expect(f.logs[0]).toContain("obligation retained"); +}); + +for (const state of ["acked", "failed"] as const) { + for (const lookup of ["missing", "ambiguous", "bounded"] as const) { + test(`${state}: ${lookup} cur lookup holds the obligation`, () => { + const f = terminalFixture(state); + if (lookup === "missing") realFs.unlinkSync(f.curPath); + if (lookup === "ambiguous") realFs.writeFileSync(join(f.curDir, "independent.json"), JSON.stringify({ id: "inbound" })); + if (lookup === "bounded") { + for (let i = 0; i < 4096; i++) realFs.writeFileSync(join(f.curDir, `${i}.json`), JSON.stringify({ id: `other-${i}` })); + } + expect(f.reconcile().has("inbound")).toBe(true); + const code = lookup === "missing" ? "ENOENT" : lookup === "ambiguous" ? "AMBIGUOUS_ID" : "SCAN_LIMIT"; + expect(f.logs[0]).toContain(`path=${f.curDir} code=${code}`); + expect(f.logs[0]).toContain("obligation retained"); + expect(realFs.existsSync(f.obligationPath)).toBe(true); + if (lookup !== "missing") expect(JSON.parse(realFs.readFileSync(f.curPath, "utf8"))[state === "acked" ? "ackedAt" : "nackedAt"]).toBeUndefined(); + }); + } +}