Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
e08dd50
fix(openclaw-tps-mail): surface and reconcile a failed cur/ record wr…
tps-flint Oct 3, 2026
80e8c7b
fix(openclaw-tps-mail): a failed terminal stamp is retried in-process…
tps-flint Oct 3, 2026
6e4f587
Merge origin/main (43795383) into the terminal-stamp retry PR; the st…
tps-flint Oct 3, 2026
ac143ab
Merge origin/main (f0205f11) into fix/492-patchmail-terminal-write
tps-flint Oct 3, 2026
58f9a56
fix(openclaw-tps-mail): drop the unused mail-lock subpath export; bod…
tps-flint Oct 3, 2026
6941773
fix(openclaw-tps-mail): failed stamps and unreadable obligations are …
tps-flint Oct 4, 2026
e686921
fix(openclaw-tps-mail): an unreadable obligation is never erased or r…
tps-flint Oct 4, 2026
8aee8e9
fix(openclaw-tps-mail): every startup reconcile failure leaves the in…
tps-flint Oct 4, 2026
5476ff1
fix(openclaw-tps-mail): live stamps require the expected inbound id; …
tps-flint Oct 4, 2026
9082eef
fix(openclaw-tps-mail): distinct remedies for an id mismatch and an u…
tps-flint Oct 4, 2026
73e57ab
fix(openclaw-tps-mail): a directory read failure is an explicit kind,…
tps-flint Oct 4, 2026
a2eedca
Merge main into fix/468 (#500): accept main's removal of packages/pi-…
tps-flint Oct 4, 2026
b99bdcf
fix(openclaw-tps-mail): diagnostics name the path, code and state and…
tps-flint Oct 4, 2026
2b77f91
fix(openclaw-tps-mail): the exhausted-retry diagnostic states no retr…
tps-flint Oct 4, 2026
202761c
fix(openclaw-tps-mail): diagnostics name the path that failed; tests …
tps-flint Oct 4, 2026
977776b
fix(openclaw-tps-mail): a missing record's diagnostic names the recor…
tps-flint Oct 4, 2026
f994643
refactor(openclaw-tps-mail): one diagnostic builder fed with facts fo…
tps-flint Oct 4, 2026
2ba3799
fix(openclaw-tps-mail): diagnostics narrowed to the failed stamp writ…
tps-flint Oct 4, 2026
cab4372
fix(openclaw-tps-mail): reconciliation finds the record by what it is…
tps-flint Oct 4, 2026
b143687
fix(openclaw-tps-mail): retention never sweeps an obligation reconcil…
tps-flint Oct 4, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **Report failed cur/ stamps**. Retry failed ack/nack stamp writes.
14 changes: 10 additions & 4 deletions packages/cli/src/utils/mail.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -189,10 +189,16 @@ export function updateExistingRecord<T extends object>(
mutate: (record: T) => T | null,
options: { snapshot?: T; afterWrite?: (record: T) => void; nonBlocking?: boolean } = {},
): UpdateExistingResult<T> {
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;
Expand All @@ -215,7 +221,7 @@ export function updateExistingRecord<T extends object>(
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 {} }
Expand Down
29 changes: 29 additions & 0 deletions plugins/openclaw-tps-mail/src/diagnostics.ts
Original file line number Diff line number Diff line change
@@ -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("; ");
}
164 changes: 143 additions & 21 deletions plugins/openclaw-tps-mail/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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;
}
}
Expand Down Expand Up @@ -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<TpsMailBody>): boolean {
type CurRecordUpdate =
| { ok: true }
| { ok: false; reason: "record-missing" | "write-failed"; path: string; code?: string };

export function patchMailFile(path: string, patch: Partial<TpsMailBody>, stampKey: "ackedAt" | "nackedAt" | undefined, inboundId: string): CurRecordUpdate {
try {
return updateExistingRecord<TpsMailBody>(path, (current) => Object.assign(current, patch)).status === "updated";
} catch {
return false;
let alreadyStamped = false;
const r = updateExistingRecord<TpsMailBody>(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<string>()): Set<string> {
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;
}

/**
Expand Down Expand Up @@ -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<TpsMailBody>,
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
Expand All @@ -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);
}
Expand Down Expand Up @@ -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).
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -2262,6 +2381,8 @@ const gateway: ChannelGatewayAdapter<TpsMailAccount> = {
}
} 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
Expand All @@ -2278,7 +2399,7 @@ const gateway: ChannelGatewayAdapter<TpsMailAccount> = {
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);
}
}
Expand Down Expand Up @@ -2341,13 +2462,14 @@ const gateway: ChannelGatewayAdapter<TpsMailAccount> = {
// 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)}`);
Expand All @@ -2358,7 +2480,7 @@ const gateway: ChannelGatewayAdapter<TpsMailAccount> = {
// 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,
Expand Down
25 changes: 21 additions & 4 deletions plugins/openclaw-tps-mail/src/obligations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.
}
}
Expand Down Expand Up @@ -502,6 +511,7 @@ export function sweepTerminalObligations(
log?: ObligationLog,
nowMs: number = Date.now(),
nackHoldDays: number = retentionDays * DEFAULT_NACK_HOLD_MULTIPLE,
unresolved = new Set<string>(),
): RetentionResult {
const res: RetentionResult = {
removed: 0,
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading