diff --git a/.changelog/unreleased/fixed-535-branch-receiver-before-ack.md b/.changelog/unreleased/fixed-535-branch-receiver-before-ack.md index 6252fab2..0a47e521 100644 --- a/.changelog/unreleased/fixed-535-branch-receiver-before-ack.md +++ b/.changelog/unreleased/fixed-535-branch-receiver-before-ack.md @@ -1 +1 @@ -- **Inbox-routed branch relay deliveries are recorded before ACK and reused or republished on resend.** Inbox payload conflicts are refused without ACK. Reply, forward and drop handler outcomes are outside this guarantee. Promotion rejects repeated signed-envelope messageIds. +- **Inbox-routed branch relay deliveries are recorded before ACK.** Inbox payload conflicts are refused without ACK. Reply, forward and drop handler outcomes are outside this guarantee. Promotion rejects repeated signed-envelope messageIds. diff --git a/.changelog/unreleased/fixed-573-relay-resend-after-ack.md b/.changelog/unreleased/fixed-573-relay-resend-after-ack.md new file mode 100644 index 00000000..2386f8b1 --- /dev/null +++ b/.changelog/unreleased/fixed-573-relay-resend-after-ack.md @@ -0,0 +1 @@ +- **Deduplicate relay resends after local ACK while the receipt remains.** A fixed set of lock stripes serializes acceptance; a failed acceptance that throws rolls back what the attempt wrote, and incomplete rollback refuses delivery with a possible-duplicate warning; the receipt TTL also applies to flat per-branch markers. diff --git a/packages/cli/src/utils/mail.ts b/packages/cli/src/utils/mail.ts index deb6dfae..40e8ebbd 100644 --- a/packages/cli/src/utils/mail.ts +++ b/packages/cli/src/utils/mail.ts @@ -394,6 +394,9 @@ export class MailSyncError extends Error { } } +export class MailSendInputError extends Error {} +export class MailInboxFullError extends Error {} + export function syncMailFile(path: string): void { try { const fd = openSync(path, "r"); @@ -405,6 +408,21 @@ export function syncMailDirectory(path: string): void { syncMailFile(path); } +/** Remove a file, confirm it is gone, and sync its directory; throws if removal cannot be confirmed. */ +export function removeMailFileConfirmed(path: string): void { + try { unlinkSync(path); } + catch (error) { if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error; } + try { lstatSync(path); } + catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code === "ENOENT" || code === "ENOTDIR") { syncMailDirectory(dirname(path)); return; } + throw error; + } + throw new Error("file still present after removal"); +} + +export class DeadLetterCleanupError extends Error {} + const pendingDirectorySyncs = new Set(); export function mkdirMailDirectory(path: string): void { @@ -568,10 +586,14 @@ export function findRelayedRecord(agent: string, delivery: { branchId: string; i } export function sendMessage(to: string, body: string, from?: string, relayDelivery?: { branchId: string; id: string }, senderTimestamp?: string, wireRecipient = to, inboxRoot?: string): MailMessage & { filePath: string } { - assertValidAgentId(to); const sender = from || "unknown"; - assertValidAgentId(sender); - assertValidBody(body); + try { + assertValidAgentId(to); + assertValidAgentId(sender); + assertValidBody(body); + } catch (error) { + throw new MailSendInputError(error instanceof Error ? error.message : String(error)); + } // Guard: when running in test mode or when the caller has explicitly // opted in, refuse to write to the default ~/.tps/mail/ directory @@ -588,7 +610,7 @@ export function sendMessage(to: string, body: string, from?: string, relayDelive const inbox = inboxRoot === undefined ? getInbox(to) : inboxAtRoot(inboxRoot); const quotaCount = readdirSync(inbox.fresh).filter((f) => f.endsWith(".json")).length; if (quotaCount >= MAX_INBOX_MESSAGES) { - throw new Error(inboxFullMessage(to, quotaCount)); + throw new MailInboxFullError(inboxFullMessage(to, quotaCount)); } const timestamp = new Date().toISOString(); @@ -731,11 +753,20 @@ export function deadLetterUndelivered( const filename = `${safeTs}-${record.id}-${randomUUID()}.json`; const tmpPath = join(inbox.tmp, filename); writeFileSync(tmpPath, JSON.stringify({ ...record, read: false, ...(relayDelivery ? { relayDelivery, relayPayload: { from: record.from, to: record.to, body: record.body, timestamp: record.timestamp }, receivedAt: new Date().toISOString() } : {}) }, null, 2), "utf-8"); - writeReasonSidecar(inbox.dlq, filename, cls, reason); - if (relayDelivery) { - syncMailFile(join(inbox.dlq, `${filename}.reason`)); - publishRelayedRecord(tmpPath, join(inbox.dlq, filename)); - } else renameSync(tmpPath, join(inbox.dlq, filename)); + const sidecar = join(inbox.dlq, `${filename}.reason`); + try { + writeReasonSidecar(inbox.dlq, filename, cls, reason); + if (relayDelivery) { + syncMailFile(sidecar); + publishRelayedRecord(tmpPath, join(inbox.dlq, filename)); + } else renameSync(tmpPath, join(inbox.dlq, filename)); + } catch (error) { + if (relayDelivery && !existsSync(join(inbox.dlq, filename))) { + try { removeMailFileConfirmed(sidecar); } + catch { throw new DeadLetterCleanupError("dead-letter sidecar removal could not be confirmed"); } + } + throw error; + } return join(inbox.dlq, filename); } diff --git a/packages/cli/src/utils/relay.ts b/packages/cli/src/utils/relay.ts index 8245ad3b..9520bfa3 100644 --- a/packages/cli/src/utils/relay.ts +++ b/packages/cli/src/utils/relay.ts @@ -1,9 +1,9 @@ -import { appendFileSync, existsSync, mkdirSync, readdirSync, readFileSync, renameSync, writeFileSync } from "node:fs"; +import { appendFileSync, existsSync, lstatSync, mkdirSync, readdirSync, readFileSync, renameSync, rmSync, statSync, unlinkSync, writeFileSync } from "node:fs"; import { dirname, join, resolve, sep } from "node:path"; import { homedir } from "node:os"; -import { randomUUID } from "node:crypto"; +import { createHash, randomUUID } from "node:crypto"; import { sanitizeIdentifier } from "../schema/sanitizer.js"; -import { countInboxMessages, deadLetterUndelivered, findRelayedRecord, getMailDir, mkdirMailDirectory, MailSyncError, relayAcceptRoot, relayAcceptLockRoot, syncMailFile, syncMailDirectory, inboxFullMessage, MAX_INBOX_MESSAGES, sendMessage, type PromoteRejectClass } from "./mail.js"; +import { countInboxMessages, DeadLetterCleanupError, deadLetterUndelivered, findRelayedRecord, removeMailFileConfirmed, getMailDir, mkdirMailDirectory, MailInboxFullError, MailSendInputError, relayAcceptRoot, syncMailFile, syncMailDirectory, inboxFullMessage, MAX_INBOX_MESSAGES, sendMessage, type PromoteRejectClass } from "./mail.js"; import { acquireMailLockSync } from "./mail-lock.js"; import { LoopDetector } from "./loop-detector.js"; import { FileSystemTransport, resolveTransport, TransportRegistry, type TransportChannel, type TpsMessage } from "./transport.js"; @@ -16,7 +16,7 @@ import { registerServiceProxyHandler } from "./service-proxy-host.js"; import { clearHostState, writeHostState, type HostConnectionState, type ServiceHealth } from "./connection-state.js"; import { listServices } from "./service-registry.js"; import snooplogg from "snooplogg"; -import type { ZodError } from "zod"; +import { z, type ZodError } from "zod"; const { log: slog, warn: swarn, error: serror } = snooplogg("tps:relay"); @@ -371,27 +371,206 @@ export function handleIncomingMail(branchId: string, msg: TpsMessage): void { }); } -function recordAcceptance(acceptedDir: string, marker: string): void { - mkdirMailDirectory(acceptedDir); +/** + * The payload an acceptance receipt binds to its branch+id: the relay payload + * of the delivery that was accepted, so a later resend can be judged against it + * after the inbox record is gone. + */ +const RelayAcceptReceiptSchema = z.object({ + from: MailDeliverBodySchema.shape.from, + to: MailDeliverBodySchema.shape.to, + body: z.string(), + timestamp: MailDeliverBodySchema.shape.timestamp, +}); +type RelayAcceptReceipt = z.infer; + +function acceptReceipt(body: MailDeliverBody): RelayAcceptReceipt { + return { from: body.from, to: body.to, body: body.content, timestamp: body.timestamp }; +} + +function refusalText(error: unknown, content: string): string { + const message = error instanceof Error ? error.message : String(error); + return content ? message.replaceAll(content, "[redacted]") : message; +} + +/** A receipt stops blocking resends after this age. */ +const RELAY_ACCEPT_RECEIPT_TTL_MS = 7 * 24 * 60 * 60 * 1000; + +function relayAcceptReceiptTtlMs(): number { + const raw = Number(process.env.TPS_RELAY_ACCEPT_RECEIPT_TTL_MS); + return Number.isFinite(raw) && raw > 0 ? raw : RELAY_ACCEPT_RECEIPT_TTL_MS; +} + +function receiptBucket(now = Date.now()): string { + return new Date(now).toISOString().slice(0, 10); +} + +export function relayAcceptanceReceiptPath(branchId: string, id: string, now = Date.now()): string { + return join(getMailDir(), ".relay-accepted", "by-branch", branchId, receiptBucket(now), id); +} + +/** Remove old day buckets and expired flat per-branch markers. */ +export function pruneRelayAcceptanceReceipts(acceptedDir: string, now = Date.now(), ttlMs = relayAcceptReceiptTtlMs()): number { + if (!existsSync(acceptedDir)) return 0; + let removed = 0; + for (const entry of readdirSync(acceptedDir, { withFileTypes: true })) { + const path = join(acceptedDir, entry.name); + if (entry.isDirectory() && /^\d{4}-\d{2}-\d{2}$/.test(entry.name)) { + const end = Date.parse(`${entry.name}T00:00:00.000Z`) + 24 * 60 * 60 * 1000; + if (!Number.isFinite(end) || end > now - ttlMs) continue; + try { rmSync(path, { recursive: true }); removed++; } catch {} + } else if (entry.isFile() && MailDeliverBodySchema.shape.id.safeParse(entry.name).success) { + try { + if (now - statSync(path).mtimeMs <= ttlMs) continue; + unlinkSync(path); + removed++; + } catch {} + } + } + if (removed > 0) syncMailDirectory(acceptedDir); + return removed; +} + +function startReceiptPrune(branchId?: string): () => void { + if (branchId && !/^[a-zA-Z0-9_-]+$/.test(branchId)) throw new Error("invalid branch id for receipt prune"); + const pass = () => { + const root = join(getMailDir(), ".relay-accepted", "by-branch"); + const prune = (dir: string) => { + try { pruneRelayAcceptanceReceipts(dir); } + catch { console.error("[relay] acceptance receipt prune failed"); } + }; + try { + if (branchId) prune(join(root, branchId)); + else if (existsSync(root)) { + for (const entry of readdirSync(root, { withFileTypes: true })) { + if (entry.isDirectory()) prune(join(root, entry.name)); + } + } + } catch { console.error("[relay] acceptance receipt prune failed"); } + }; + pass(); + const configured = Number(process.env.TPS_RELAY_ACCEPT_PRUNE_INTERVAL_MS); + const timer = setInterval(pass, Number.isFinite(configured) && configured > 0 ? configured : 60_000); + timer.unref(); + return () => clearInterval(timer); +} + +class ReceiptRecoveryIncompleteError extends Error {} + +function receiptMayExist(path: string): boolean { + try { lstatSync(path); return true; } + catch (error) { + const code = (error as NodeJS.ErrnoException).code; + return code !== "ENOENT" && code !== "ENOTDIR"; + } +} + +function recordAcceptance(acceptedDir: string, marker: string, receipt: RelayAcceptReceipt): void { + mkdirMailDirectory(dirname(marker)); const tmp = `${marker}.tmp`; - writeFileSync(tmp, "", { mode: 0o600 }); - syncMailFile(tmp); - renameSync(tmp, marker); - syncMailDirectory(acceptedDir); + try { + writeFileSync(tmp, JSON.stringify(receipt), { mode: 0o600, flag: "wx" }); + syncMailFile(tmp); + renameSync(tmp, marker); + syncMailDirectory(dirname(marker)); + } catch (error) { + if (receiptMayExist(marker)) { + try { unlinkSync(marker); syncMailDirectory(dirname(marker)); } + catch { throw new ReceiptRecoveryIncompleteError("relay receipt recovery incomplete; retry may duplicate"); } + } + try { removeMailFileConfirmed(tmp); } + catch { + if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw new ReceiptRecoveryIncompleteError("relay receipt recovery incomplete; retry may duplicate"); + } + throw error; + } + try { + pruneRelayAcceptanceReceipts(acceptedDir); + } catch { + console.error("[relay] acceptance receipt prune failed"); + } +} + +/** + * Find a marker unless its age is known to exceed the receipt TTL. + */ +function existingAcceptanceMarker(acceptedDir: string, id: string, legacyMarker: string): string | undefined { + if (receiptMayExist(acceptedDir)) { + for (const bucket of readdirSync(acceptedDir, { withFileTypes: true }).filter((entry) => entry.isDirectory()).map((entry) => entry.name).sort().reverse()) { + const path = join(acceptedDir, bucket, id); + if (!receiptMayExist(path)) continue; + try { + if (Date.now() - statSync(path).mtimeMs > relayAcceptReceiptTtlMs()) continue; + } catch { + if (!receiptMayExist(path)) continue; + } + return path; + } + } + const flatMarker = join(acceptedDir, id); + if (receiptMayExist(flatMarker)) { + try { if (Date.now() - statSync(flatMarker).mtimeMs <= relayAcceptReceiptTtlMs()) return flatMarker; } + catch { if (receiptMayExist(flatMarker)) return flatMarker; } + } + if (!receiptMayExist(legacyMarker)) return undefined; + try { if (Date.now() - statSync(legacyMarker).mtimeMs > relayAcceptReceiptTtlMs()) return undefined; } + catch { if (!receiptMayExist(legacyMarker)) return undefined; } + return legacyMarker; +} + +/** + * Whether a readable receipt matches the payload; legacy empty markers cannot be compared. + */ +function acceptanceReceiptMatches(path: string, expected: RelayAcceptReceipt): boolean { + const raw = readFileSync(path, "utf-8"); + let parsed: unknown; + try { + parsed = JSON.parse(raw); + } catch { + return false; + } + const receipt = RelayAcceptReceiptSchema.safeParse(parsed); + if (!receipt.success) return false; + const r = receipt.data; + return r.from === expected.from && r.to === expected.to && r.body === expected.body && r.timestamp === expected.timestamp; } -/** The shared wait budget for relay acceptance's mailbox locks. */ +function undoNewRelayRecord(root: string, delivery: { branchId: string; id: string }): void { + for (const dir of ["tmp", "new", "dlq"]) { + const parent = join(root, dir); + if (!existsSync(parent)) continue; + for (const entry of readdirSync(parent, { withFileTypes: true })) { + if (!entry.isFile() || !entry.name.endsWith(".json")) continue; + const path = join(parent, entry.name); + let record: { relayDelivery?: { branchId?: string; id?: string } }; + try { record = JSON.parse(readFileSync(path, "utf8")); } + catch { throw new Error("relay record recovery incomplete; delivery may be retried"); } + if (record.relayDelivery?.branchId !== delivery.branchId || record.relayDelivery.id !== delivery.id) continue; + unlinkSync(path); + if (dir === "dlq" && existsSync(`${path}.reason`)) unlinkSync(`${path}.reason`); + syncMailDirectory(parent); + } + } +} + +/** Wait budget for relay acceptance locks. */ const RELAY_ACCEPT_LOCK_TIMEOUT_MS = 2000; +export const RELAY_ACCEPT_LOCK_STRIPES = 64; + +export function relayAcceptanceLockRoot(branchId: string, id: string): string { + const stripe = createHash("sha256").update(branchId).update("\0").update(id).digest()[0]! % RELAY_ACCEPT_LOCK_STRIPES; + return join(getMailDir(), ".relay-accept-locks", String(stripe)); +} function relayAcceptLockTimeoutMs(): number { const raw = Number(process.env.TPS_RELAY_ACCEPT_LOCK_TIMEOUT_MS); return Number.isFinite(raw) && raw > 0 ? raw : RELAY_ACCEPT_LOCK_TIMEOUT_MS; } -/** The named refusal when an acceptance mailbox lock times out. */ +/** The named refusal when relay acceptance waits too long for a lock. */ export class RelayAcceptLockTimeoutError extends Error { constructor(recipient: string) { - super(`relay acceptance timed out waiting for the mailbox lock for ${recipient}`); + super(`relay acceptance timed out waiting for a lock for ${recipient}`); this.name = "RelayAcceptLockTimeoutError"; } } @@ -414,30 +593,49 @@ export function deliverRelayedToLocal(branchId: string, body: MailDeliverBody): if (!lock) throw new RelayAcceptLockTimeoutError(body.to); return lock; }; - const acceptanceRoot = relayAcceptLockRoot(body.to); + const acceptanceRoot = relayAcceptanceLockRoot(branchId, body.id); mkdirMailDirectory(acceptanceRoot); const acceptanceLock = acquireLock(acceptanceRoot); let mailboxLock: ReturnType | undefined; try { + // Hold the stripe before each mailbox lock; release mailbox locks before the stripe. const recipientRoot = relayAcceptRoot(body.to); mkdirMailDirectory(recipientRoot); if (recipientRoot !== acceptanceRoot) mailboxLock = acquireLock(recipientRoot); relayAcceptTestHook?.(recipientRoot); // The marker path includes the branch: a 64-hex id is deterministic, so two branches can send the same one. const acceptedDir = join(getMailDir(), ".relay-accepted", "by-branch", branchId); - const marker = join(acceptedDir, body.id); + const marker = relayAcceptanceReceiptPath(branchId, body.id); + const flatMarker = join(acceptedDir, body.id); const legacyMarker = join(getMailDir(), ".relay-accepted", body.id); - const existingMarker = existsSync(marker) ? marker : existsSync(legacyMarker) ? legacyMarker : undefined; + const receipt = acceptReceipt(body); + const existingMarker = existingAcceptanceMarker(acceptedDir, body.id, legacyMarker); const delivery = { branchId, id: body.id }; const existingRecord = findRelayedRecord(body.to, delivery, { from: body.from, to: body.to, body: body.content, timestamp: body.timestamp }, body.to, { heldRoot: recipientRoot, heldAcceptanceRoot: acceptanceRoot, acquireLock }); - if (existingMarker && existingRecord) { - syncMailFile(existingMarker); - syncMailDirectory(dirname(existingMarker)); + if (existingRecord) { + if (existingMarker && existingMarker !== legacyMarker && existingMarker !== flatMarker && !acceptanceReceiptMatches(existingMarker, receipt)) { + throw new Error(`relayed delivery conflict for branch ${branchId} message ${body.id}`); + } + if (existingMarker && acceptanceReceiptMatches(existingMarker, receipt)) { + syncMailFile(existingMarker); + syncMailDirectory(dirname(existingMarker)); + } else { + try { recordAcceptance(acceptedDir, marker, receipt); } + catch { throw new Error(`relay receipt recovery incomplete for message ${body.id}; retry may duplicate`); } + } return false; } - if (!existingMarker && existingRecord) { - recordAcceptance(acceptedDir, marker); - return false; + if (existingMarker) { + // A marker without a record blocks publication unless its receipt matches. + if (acceptanceReceiptMatches(existingMarker, receipt)) { + syncMailFile(existingMarker); + syncMailDirectory(dirname(existingMarker)); + return false; + } + if (existingMarker === flatMarker || existingMarker === legacyMarker) { + throw new Error(`relayed delivery conflict for branch ${branchId} message ${body.id}: accepted before receipts existed; payload cannot be compared; sender must not retry`); + } + throw new Error(`relayed delivery conflict for branch ${branchId} message ${body.id}`); } let delivered: boolean; @@ -445,30 +643,39 @@ export function deliverRelayedToLocal(branchId: string, body: MailDeliverBody): sendMessage(body.to, body.content, body.from, delivery, body.timestamp, body.to, recipientRoot); delivered = true; } catch (e: unknown) { - if (e instanceof MailSyncError) throw e; + try { undoNewRelayRecord(recipientRoot, delivery); } + catch { throw new Error(`relay record recovery incomplete for message ${body.id}; retry may duplicate`); } const reason = e instanceof Error ? e.message : String(e); - const cls: PromoteRejectClass = /inbox full/i.test(reason) + const cls: PromoteRejectClass | null = e instanceof MailInboxFullError ? "inbox-full" - : /^(Invalid agent id|Message body)/.test(reason) ? "invalid" : "storage-unavailable"; - console.error(`[relay] local delivery failed for message ${body.id} to ${body.to}: ${reason}`); + : e instanceof MailSendInputError ? "invalid" : null; + if (!cls) throw new Error("relay inbox write failed; retry delivery"); + console.error(`[relay] local delivery failed for message ${body.id} to ${body.to}: ${refusalText(e, body.content)}`); try { deadLetterUndelivered( body.to, { id: body.id, from: body.from, to: body.to, body: body.content, timestamp: body.timestamp }, cls, - reason, + refusalText(e, body.content), delivery, recipientRoot, ); - } catch (dlqErr: unknown) { - console.error( - `[relay] dead-letter failed for message ${body.id} to ${body.to}: ${dlqErr instanceof Error ? dlqErr.message : String(dlqErr)}`, - ); - throw dlqErr; + } catch (deadLetterError) { + if (deadLetterError instanceof DeadLetterCleanupError) throw new Error(`relay record recovery incomplete for message ${body.id}; retry may duplicate`); + try { undoNewRelayRecord(recipientRoot, delivery); } + catch { throw new Error(`relay record recovery incomplete for message ${body.id}; retry may duplicate`); } + console.error(`[relay] dead-letter failed for message ${body.id} to ${body.to}`); + throw new Error("relay dead-letter failed; retry delivery"); } delivered = false; } - recordAcceptance(acceptedDir, marker); + try { recordAcceptance(acceptedDir, marker, receipt); } + catch (error) { + if (error instanceof ReceiptRecoveryIncompleteError || receiptMayExist(marker)) throw new Error(`relay receipt recovery incomplete for message ${body.id}; retry may duplicate`); + try { undoNewRelayRecord(recipientRoot, delivery); } + catch { throw new Error(`relay record recovery incomplete for message ${body.id}; retry may duplicate`); } + throw new Error("relay acceptance receipt write failed; retry delivery"); + } return delivered; } finally { mailboxLock?.release(); @@ -484,12 +691,7 @@ async function acceptRelayedMail( onAccepted?: () => void, ): Promise { const delivered = deliverRelayedToLocal(branchId, body); - // Fire the caller's accepted-delivery hook only when this call published a - // new inbox record. `delivered` is false for a resend whose record already - // exists, and for an inbox write failure that deliverRelayedToLocal handled - // by dead-lettering the message; both are still ACKed below. A same-id - // conflict throws before this point, as does a sync failure or a failed - // dead-letter, so none of those is announced or ACKed. + // A delivery with delivered=false (dead-lettered or receipt-backed duplicate) is still ACKed below. if (delivered) { const reportError = () => { console.error(`[relay] onAccepted failed for message ${body.id} to ${body.to}`); @@ -515,6 +717,7 @@ function logRefusedDelivery(branchId: string, error: ZodError): void { export function startRelay(agentId: string): () => void { assertAgent(agentId); + const stopPrune = startReceiptPrune(); let remoteCleanup: (() => Promise) | null = null; connectRemoteBranches(transportRegistry, (branchId, msg) => { @@ -539,6 +742,7 @@ export function startRelay(agentId: string): () => void { const stop = () => { clearInterval(timer); + stopPrune(); remoteCleanup?.().catch(() => {}); }; return stop; @@ -720,6 +924,7 @@ export async function syncRemoteBranch(branchId: string): Promise<{ received: nu const hostKp = await loadHostIdentity(); const transport = transportType === "ws" ? new WsNoiseTransport(hostKp) : new NoiseIkTransport(hostKp); const channel = await transport.connect({ host, port, branchId, hostPublicKey: branch.encryptionKey }); + const stopPrune = startReceiptPrune(branchId); let received = 0; try { @@ -741,13 +946,14 @@ export async function syncRemoteBranch(branchId: string): Promise<{ received: nu void acceptRelayedMail(channel, branchId, msg, parsed.data).then((delivered) => { if (delivered) received++; }).catch((error: unknown) => { - console.error(`[relay] acceptance failed for message ${parsed.data.id} to ${parsed.data.to}: ${String(error)}`); + console.error(`[relay] acceptance failed for message ${parsed.data.id} to ${parsed.data.to}: ${refusalText(error, parsed.data.content)}`); }); }; channel.onMessage(handler); }); } finally { + stopPrune(); await channel.close().catch(() => {}); } @@ -789,6 +995,7 @@ export async function connectAndKeepAlive( const branch = lookupBranch(branchId); if (!branch?.encryptionKey) throw new Error(`Branch '${branchId}' missing encryption key`); const hostKp = await loadHostIdentity(); + const stopPrune = startReceiptPrune(branchId); const loop = async () => { let backoff = RECONNECT_BASE_MS; @@ -857,7 +1064,7 @@ export async function connectAndKeepAlive( const parsed = MailDeliverBodySchema.safeParse(msg.body); if (parsed.success) { void acceptRelayedMail(channel, branchId, msg, parsed.data, () => opts.onAccepted?.(msg)).catch((error: unknown) => { - console.error(`[relay] acceptance failed for message ${parsed.data.id} to ${parsed.data.to}: ${String(error)}`); + console.error(`[relay] acceptance failed for message ${parsed.data.id} to ${parsed.data.to}: ${refusalText(error, parsed.data.content)}`); }); } else { logRefusedDelivery(branchId, parsed.error); @@ -891,6 +1098,7 @@ export async function connectAndKeepAlive( return async () => { stopped = true; + stopPrune(); try { await currentChannel?.close(); } catch {} clearHostState(branchId); }; diff --git a/packages/cli/test/office-connect-announce.test.ts b/packages/cli/test/office-connect-announce.test.ts index 65cfbf24..2271a8b5 100644 --- a/packages/cli/test/office-connect-announce.test.ts +++ b/packages/cli/test/office-connect-announce.test.ts @@ -445,7 +445,7 @@ describe("office connect announcement follows local acceptance", () => { expect(announced).toEqual([1]); }); - test("a delivery whose inbox write fails is dead-lettered and is not announced", async () => { + test("an inbox write failure is unacked and unannounced", async () => { const inbox = getInbox("local"); const write = fs.writeFileSync; let failed = false; @@ -463,9 +463,9 @@ describe("office connect announcement follows local acceptance", () => { fault.mockRestore(); } expect(failed).toBe(true); - expect(acks.length).toBe(1); + expect(acks.length).toBe(0); expect(jsonFiles(inbox.fresh)).toEqual([]); - expect(jsonFiles(inbox.dlq).length).toBe(1); + expect(jsonFiles(inbox.dlq)).toEqual([]); expect(announced).toEqual([]); }); }); diff --git a/packages/cli/test/relay-accept-invariants.test.ts b/packages/cli/test/relay-accept-invariants.test.ts new file mode 100644 index 00000000..f28deda4 --- /dev/null +++ b/packages/cli/test/relay-accept-invariants.test.ts @@ -0,0 +1,221 @@ +import { afterEach, beforeEach, expect, mock, spyOn, test } from "bun:test"; +import * as fs from "node:fs"; +import { randomUUID } from "node:crypto"; +import { tmpdir } from "node:os"; +import { dirname, join } from "node:path"; +import { ackMessageAtPath, getInbox, sendMessage } from "../src/utils/mail.js"; +import { deliverRelayedToLocal, pruneRelayAcceptanceReceipts, relayAcceptanceReceiptPath, startRelay } from "../src/utils/relay.js"; + +let root: string; +let saved: Record; + +beforeEach(() => { + root = fs.mkdtempSync(join(tmpdir(), "relay-accept-invariants-")); + saved = Object.fromEntries(["HOME", "TPS_MAIL_DIR", "TPS_RELAY_ACCEPT_PRUNE_INTERVAL_MS"].map((key) => [key, process.env[key]])); + process.env.HOME = root; + process.env.TPS_MAIL_DIR = join(root, "mail"); +}); + +afterEach(() => { + for (const [key, value] of Object.entries(saved)) { + if (value === undefined) delete process.env[key]; + else process.env[key] = value; + } + fs.rmSync(root, { recursive: true, force: true }); +}); + +afterEach(() => { mock.restore(); }); + +function body() { + return { id: randomUUID(), from: "remote", to: "local", content: "private-payload-text", timestamp: new Date().toISOString() }; +} + +function records(): string[] { + const inbox = getInbox("local"); + return fs.readdirSync(inbox.fresh).filter((file) => file.endsWith(".json")); +} + +test("receipt write failure removes the new record before retry, ACK, and resend", () => { + const item = body(); + const marker = relayAcceptanceReceiptPath("remote", item.id); + const errors = spyOn(console, "error").mockImplementation(() => {}); + const write = fs.writeFileSync; + const fault = spyOn(fs, "writeFileSync").mockImplementation((path, data, options) => { + if (String(path) === `${marker}.tmp`) throw new Error(`receipt failure ${item.content}`); + return write(path, data, options); + }); + let refusal = ""; + try { + try { deliverRelayedToLocal("remote", item); } + catch (error) { refusal = String(error); } + } + finally { fault.mockRestore(); } + expect(refusal).toContain("relay acceptance receipt write failed; retry delivery"); + expect(refusal).not.toContain(item.content); + expect(records()).toEqual([]); + expect(fs.existsSync(marker)).toBe(false); + expect(deliverRelayedToLocal("remote", item)).toBe(true); + const [file] = records(); + ackMessageAtPath(join(getInbox("local").fresh, file!)); + expect(deliverRelayedToLocal("remote", item)).toBe(false); + expect(records()).toEqual([]); + expect(fs.statSync(marker).mode & 0o777).toBe(0o600); + expect(errors.mock.calls.flat().join("\n")).not.toContain(item.content); + errors.mockRestore(); +}); + +test("record publication failure leaves no receipt and resend publishes once", () => { + const item = body(); + const marker = relayAcceptanceReceiptPath("remote", item.id); + const rename = fs.renameSync; + const fault = spyOn(fs, "renameSync").mockImplementation((from, to) => { + if (String(to).startsWith(getInbox("local").fresh)) throw new Error("injected record publication failure"); + return rename(from, to); + }); + try { expect(() => deliverRelayedToLocal("remote", item)).toThrow(); } + finally { fault.mockRestore(); } + expect(fs.existsSync(marker)).toBe(false); + expect(records()).toEqual([]); + expect(deliverRelayedToLocal("remote", item)).toBe(true); + expect(deliverRelayedToLocal("remote", item)).toBe(false); + expect(records()).toHaveLength(1); +}); + +test("inbox write failures remain unacked and retryable", () => { + const item = body(); + const marker = relayAcceptanceReceiptPath("remote", item.id); + const tmp = getInbox("local").tmp; + const write = fs.writeFileSync; + const fault = spyOn(fs, "writeFileSync").mockImplementation((path, data, options) => { + if (String(path).startsWith(tmp)) throw new Error("Inbox full: injected write failure"); + return write(path, data, options); + }); + try { expect(() => deliverRelayedToLocal("remote", item)).toThrow("relay inbox write failed; retry delivery"); } + finally { fault.mockRestore(); } + expect(fs.existsSync(marker)).toBe(false); + expect(records()).toEqual([]); + expect(deliverRelayedToLocal("remote", item)).toBe(true); + expect(records()).toHaveLength(1); +}); + +test("a receipt rollback sync failure preserves the record and warns of a possible duplicate", () => { + const item = body(); + const marker = relayAcceptanceReceiptPath("remote", item.id); + const open = fs.openSync; + const sync = fs.fsyncSync; + const paths = new Map(); + const opened = spyOn(fs, "openSync").mockImplementation((path, flags, mode) => { + const fd = open(path, flags, mode); + paths.set(fd, String(path)); + return fd; + }); + const fault = spyOn(fs, "fsyncSync").mockImplementation((fd) => { + if (paths.get(fd) === dirname(marker)) throw new Error("injected receipt directory sync failure"); + return sync(fd); + }); + let refusal = ""; + try { + try { deliverRelayedToLocal("remote", item); } + catch (error) { refusal = String(error); } + } finally { + fault.mockRestore(); + opened.mockRestore(); + } + expect(refusal).toContain("retry may duplicate"); + expect(refusal).not.toContain(item.content); + expect(fs.existsSync(marker)).toBe(false); + expect(records()).toHaveLength(1); +}); + +test("old day buckets with more than one former pass expire together", () => { + const accepted = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote"); + for (const day of ["2020-01-01", "2020-01-02"]) { + const bucket = join(accepted, day); + fs.mkdirSync(bucket, { recursive: true }); + for (let i = 0; i < 2050; i++) fs.writeFileSync(join(bucket, `receipt-${i}`), "x"); + } + expect(pruneRelayAcceptanceReceipts(accepted, Date.now(), 1000)).toBe(2); + expect(fs.readdirSync(accepted)).toEqual([]); +}); + +test("a flat per-branch marker with a live record remains an in-flight duplicate", () => { + const item = body(); + const original = sendMessage(item.to, item.content, item.from, { branchId: "remote", id: item.id }, item.timestamp); + const flatMarker = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", item.id); + fs.mkdirSync(dirname(flatMarker), { recursive: true }); + fs.writeFileSync(flatMarker, ""); + expect(deliverRelayedToLocal("remote", item)).toBe(false); + expect(records()).toHaveLength(1); + expect(fs.existsSync(relayAcceptanceReceiptPath("remote", item.id))).toBe(true); + ackMessageAtPath(original.filePath); + expect(deliverRelayedToLocal("remote", item)).toBe(false); + expect(records()).toEqual([]); +}); + +test("a receipt in an ended day bucket younger than the TTL survives prune and still dedups", () => { + const item = body(); + expect(deliverRelayedToLocal("remote", item)).toBe(true); + const [file] = records(); + ackMessageAtPath(join(getInbox("local").fresh, file!)); + const marker = relayAcceptanceReceiptPath("remote", item.id); + const accepted = dirname(dirname(marker)); + const ended = join(accepted, new Date(Date.now() - 3 * 24 * 60 * 60 * 1000).toISOString().slice(0, 10)); + fs.renameSync(dirname(marker), ended); + expect(pruneRelayAcceptanceReceipts(accepted, Date.now(), 7 * 24 * 60 * 60 * 1000)).toBe(0); + expect(fs.existsSync(join(ended, item.id))).toBe(true); + expect(deliverRelayedToLocal("remote", item)).toBe(false); + expect(records()).toEqual([]); +}); + +test("a receipt that vanishes between listing and stat is treated as absent", () => { + const item = body(); + expect(deliverRelayedToLocal("remote", item)).toBe(true); + const [file] = records(); + ackMessageAtPath(join(getInbox("local").fresh, file!)); + const marker = relayAcceptanceReceiptPath("remote", item.id); + const stat = fs.statSync; + const fault = spyOn(fs, "statSync").mockImplementation(((path: fs.PathLike, options?: unknown) => { + if (String(path) === marker) fs.unlinkSync(marker); + return (stat as (p: fs.PathLike, o?: unknown) => unknown)(path, options); + }) as typeof fs.statSync); + try { expect(deliverRelayedToLocal("remote", item)).toBe(true); } + finally { fault.mockRestore(); } + expect(records()).toHaveLength(1); +}); + +test("flat per-branch markers expire with receipts", () => { + const accepted = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote"); + fs.mkdirSync(accepted, { recursive: true }); + const old = join(accepted, randomUUID()); + const current = join(accepted, randomUUID()); + fs.writeFileSync(old, ""); + fs.writeFileSync(current, ""); + fs.utimesSync(old, new Date(0), new Date(0)); + expect(pruneRelayAcceptanceReceipts(accepted, Date.now(), 1000)).toBe(1); + expect(fs.existsSync(old)).toBe(false); + expect(fs.existsSync(current)).toBe(true); +}); + +test("relay timer retries a failed prune without a new delivery", async () => { + const accepted = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote"); + const old = join(accepted, "2020-01-01"); + fs.mkdirSync(old, { recursive: true }); + fs.writeFileSync(join(old, "receipt"), "x"); + process.env.TPS_RELAY_ACCEPT_PRUNE_INTERVAL_MS = "10"; + const remove = fs.rmSync; + let calls = 0; + const fault = spyOn(fs, "rmSync").mockImplementation((path, options) => { + if (String(path) === old && calls++ === 0) throw new Error("injected prune failure"); + return remove(path, options); + }); + const stop = startRelay("local"); + try { + const deadline = Date.now() + 1000; + while (fs.existsSync(old) && Date.now() < deadline) await Bun.sleep(10); + expect(fs.existsSync(old)).toBe(false); + expect(calls).toBeGreaterThanOrEqual(2); + } finally { + stop(); + fault.mockRestore(); + } +}); diff --git a/packages/cli/test/relay-accept-single-lock.test.ts b/packages/cli/test/relay-accept-single-lock.test.ts index 859ef58c..175db917 100644 --- a/packages/cli/test/relay-accept-single-lock.test.ts +++ b/packages/cli/test/relay-accept-single-lock.test.ts @@ -4,11 +4,11 @@ import { afterEach, beforeEach, describe, expect, test } from "bun:test"; import { spawn } from "node:child_process"; import { randomUUID } from "node:crypto"; -import { existsSync, mkdirSync, mkdtempSync, readdirSync, readFileSync, rmSync, writeFileSync } from "node:fs"; +import { existsSync, mkdirSync, mkdtempSync, readdirSync, readFileSync, rmSync, statSync, writeFileSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; -import { getInbox } from "../src/utils/mail.js"; -import { RelayAcceptLockTimeoutError, deliverRelayedToLocal } from "../src/utils/relay.js"; +import { ackMessageAtPath, getInbox } from "../src/utils/mail.js"; +import { RELAY_ACCEPT_LOCK_STRIPES, RelayAcceptLockTimeoutError, deliverRelayedToLocal, relayAcceptanceLockRoot, relayAcceptanceReceiptPath } from "../src/utils/relay.js"; const BRANCH = "remote"; const RECIPIENT = "local"; @@ -127,9 +127,9 @@ function deliverChild(prefix: string, ids: string[], env: Record } /** Spawn a real process that holds a mailbox lock until released. */ -function holdLock(agent = RECIPIENT): Promise<{ release: () => Promise }> { +function holdLock(lockRoot = join(mail, RECIPIENT)): Promise<{ release: () => Promise }> { const child = spawn("bun", [childScript], { - env: { ...process.env, RELAY_CHILD_MODE: "hold", RELAY_CHILD_ROOT: join(mail, agent) }, + env: { ...process.env, RELAY_CHILD_MODE: "hold", RELAY_CHILD_ROOT: lockRoot }, stdio: ["pipe", "pipe", "inherit"], }); const exited = new Promise((resolve) => child.once("exit", resolve)); @@ -154,6 +154,23 @@ function recordsByDeliveryId(): Map { return byId; } +/** The stored receipt carries the delivered body, and a resend of that body dedups after ACK. */ +function expectReceiptKeepsBody(id: string, content: string): void { + const receipt = JSON.parse(readFileSync(relayAcceptanceReceiptPath(BRANCH, id), "utf8")) as { body: string }; + expect(receipt.body).toBe(content); + expect(deliverRelayedToLocal(BRANCH, { id, from: FROM, to: RECIPIENT, content, timestamp: TIMESTAMP })).toBe(false); + expect(recordsByDeliveryId().has(id)).toBe(false); +} + +function expectReceiptMatchesDeliveredRecord(id: string, bodyPattern: RegExp): void { + const [file] = recordsByDeliveryId().get(id)!; + const path = join(getInbox(RECIPIENT).fresh, file!); + const record = JSON.parse(readFileSync(path, "utf8")) as { body: string }; + expect(record.body).toMatch(bodyPattern); + ackMessageAtPath(path); + expectReceiptKeepsBody(id, record.body); +} + async function waitForFile(path: string): Promise { const deadline = Date.now() + 10000; while (!existsSync(path)) { @@ -163,6 +180,27 @@ async function waitForFile(path: string): Promise { } describe("relay acceptance (cli#561)", () => { + test("a second recipient waits for the same delivery id before any publication", async () => { + const id = randomUUID(); + const pause = join(root, "cross-recipient-pause"); + const started = join(root, "cross-recipient-started"); + const first = deliverChild("same-", [id], { RELAY_CHILD_PAUSE: pause }); + let second: Promise | undefined; + try { + await waitForFile(pause + ".ready"); + second = deliverChild("same-", [id], { RELAY_CHILD_TO: "other", RELAY_CHILD_STARTED: started }); + await waitForFile(started); + await new Promise((resolve) => setTimeout(resolve, 100)); + expect(jsonFiles(getInbox("other").fresh)).toEqual([]); + expect(existsSync(relayAcceptanceReceiptPath(BRANCH, id))).toBe(false); + } finally { writeFileSync(pause + ".release", ""); } + const [a, b] = await Promise.all([first, second]); + expect(a).toEqual({ delivered: 1, duplicate: 0, refused: 0 }); + expect(b).toEqual({ delivered: 0, duplicate: 0, refused: 1 }); + expect(jsonFiles(getInbox(RECIPIENT).fresh)).toHaveLength(1); + expect(jsonFiles(getInbox("other").fresh)).toEqual([]); + }, 60_000); + test.each(["new", "dlq"])("routing change before %s publication", async (destination) => { const id = randomUUID(); const hostRoot = getInbox(RECIPIENT).root; @@ -182,7 +220,7 @@ describe("relay acceptance (cli#561)", () => { await waitForFile(started); await new Promise((resolve) => setTimeout(resolve, 100)); expect(jsonFiles(join(branchRoot, "new"))).toEqual([]); - expect(existsSync(join(mail, ".relay-accepted", "by-branch", BRANCH, id))).toBe(false); + expect(existsSync(relayAcceptanceReceiptPath(BRANCH, id))).toBe(false); } finally { writeFileSync(pause + ".release", ""); results = await Promise.all([first, second]); @@ -219,6 +257,110 @@ describe("relay acceptance (cli#561)", () => { const byId = recordsByDeliveryId(); expect([...byId.keys()].sort()).toEqual([...ids].sort()); for (const id of ids) expect(byId.get(id)).toHaveLength(1); + + for (const id of ids) expectReceiptMatchesDeliveredRecord(id, /^[AB]-/); + }, 60_000); + + test("a differing payload waiting on the stripe is refused and the receipt keeps the delivered body", async () => { + const id = randomUUID(); + const pause = join(root, "differing-pause"); + const started = join(root, "differing-started"); + const first = deliverChild("A-", [id], { RELAY_CHILD_PAUSE: pause }); + let second: Promise | undefined; + try { + await waitForFile(pause + ".ready"); + second = deliverChild("B-", [id], { RELAY_CHILD_STARTED: started }); + await waitForFile(started); + await new Promise((resolve) => setTimeout(resolve, 200)); + } finally { writeFileSync(pause + ".release", ""); } + // ACK each record the moment it appears, so a waiter that took the lock late finds no record. + const acked: string[] = []; + let settled = false; + const both = Promise.all([first, second]).finally(() => { settled = true; }); + while (!settled) { + for (const file of jsonFiles(getInbox(RECIPIENT).fresh)) { + const path = join(getInbox(RECIPIENT).fresh, file); + try { + acked.push((JSON.parse(readFileSync(path, "utf8")) as { body: string }).body); + ackMessageAtPath(path); + } catch {} + } + await new Promise((resolve) => setImmediate(resolve)); + } + const [a, b] = await both; + expect(a).toEqual({ delivered: 1, duplicate: 0, refused: 0 }); + expect(b).toEqual({ delivered: 0, duplicate: 0, refused: 1 }); + expect(acked).toEqual([`A-${id}`]); + expectReceiptKeepsBody(id, `A-${id}`); + }, 60_000); + + test("lock directories stay bounded across distinct deliveries", () => { + for (let i = 0; i < RELAY_ACCEPT_LOCK_STRIPES * 2 + 1; i++) { + const id = randomUUID(); + expect(deliverRelayedToLocal(BRANCH, { id, from: FROM, to: RECIPIENT, content: id, timestamp: TIMESTAMP })).toBe(true); + const [file] = jsonFiles(getInbox(RECIPIENT).fresh); + ackMessageAtPath(join(getInbox(RECIPIENT).fresh, file!)); + } + const lockRoot = join(mail, ".relay-accept-locks"); + const countDirectories = (dir: string): number => readdirSync(dir, { withFileTypes: true }).reduce( + (total, entry) => total + (entry.isDirectory() ? 1 + countDirectories(join(dir, entry.name)) : 0), 0, + ); + expect(countDirectories(lockRoot)).toBeLessThanOrEqual(RELAY_ACCEPT_LOCK_STRIPES); + }, 60_000); + + test("different ids on one stripe each publish once across processes", async () => { + const seen = new Map(); + let pair: [string, string] | undefined; + for (let i = 0; i <= RELAY_ACCEPT_LOCK_STRIPES; i++) { + const id = randomUUID(); + const stripe = relayAcceptanceLockRoot(BRANCH, id); + const prior = seen.get(stripe); + if (prior) { pair = [prior, id]; break; } + seen.set(stripe, id); + } + expect(pair).toBeDefined(); + const [a, b] = await Promise.all([deliverChild("same-", pair!), deliverChild("same-", pair!)]); + expect(a.delivered + b.delivered).toBe(2); + expect(a.duplicate + b.duplicate).toBe(2); + expect(a.refused + b.refused).toBe(0); + const byId = recordsByDeliveryId(); + for (const id of pair!) expect(byId.get(id)).toHaveLength(1); + }, 20_000); + + test("two recipients resending the same id around local ACK keep the original receipt", async () => { + const id = randomUUID(); + const body = { id, from: FROM, to: RECIPIENT, content: `same-${id}`, timestamp: TIMESTAMP }; + expect(deliverRelayedToLocal(BRANCH, body)).toBe(true); + const marker = relayAcceptanceReceiptPath(BRANCH, id); + const original = readFileSync(marker, "utf8"); + const inode = statSync(marker).ino; + const holder = await holdLock(relayAcceptanceLockRoot(BRANCH, id)); + const local = deliverChild("same-", [id]); + const other = deliverChild("same-", [id], { RELAY_CHILD_TO: "other" }); + try { + const [file] = jsonFiles(getInbox(RECIPIENT).fresh); + ackMessageAtPath(join(getInbox(RECIPIENT).fresh, file!)); + } finally { await holder.release(); } + const [a, b] = await Promise.all([local, other]); + expect(a).toEqual({ delivered: 0, duplicate: 1, refused: 0 }); + expect(b).toEqual({ delivered: 0, duplicate: 0, refused: 1 }); + expect(readFileSync(marker, "utf8")).toBe(original); + expect(statSync(marker).ino).toBe(inode); + expect(jsonFiles(getInbox(RECIPIENT).fresh)).toEqual([]); + expect(jsonFiles(getInbox("other").fresh)).toEqual([]); + }, 60_000); + + test("the stripe lock blocks every recipient for that id", async () => { + const id = randomUUID(); + const holder = await holdLock(relayAcceptanceLockRoot(BRANCH, id)); + try { + process.env.TPS_RELAY_ACCEPT_LOCK_TIMEOUT_MS = "75"; + for (const to of [RECIPIENT, "other"]) { + expect(() => deliverRelayedToLocal(BRANCH, { id, from: FROM, to, content: "held", timestamp: TIMESTAMP })).toThrow(RelayAcceptLockTimeoutError); + expect(jsonFiles(getInbox(to).fresh)).toEqual([]); + } + expect(existsSync(relayAcceptanceReceiptPath(BRANCH, id))).toBe(false); + } finally { await holder.release(); } }, 60_000); test("a recipient lock timeout refuses by name, publishing no mail record or acceptance marker", async () => { @@ -235,7 +377,7 @@ describe("relay acceptance (cli#561)", () => { }, 60_000); test("a different mailbox's lock uses the acceptance deadline and refuses by name, publishing no mail record or acceptance marker", async () => { - const holder = await holdLock("other"); + const holder = await holdLock(join(mail, "other")); try { process.env.TPS_RELAY_ACCEPT_LOCK_TIMEOUT_MS = "75"; const body = { id: randomUUID(), from: FROM, to: RECIPIENT, content: "timed out", timestamp: TIMESTAMP }; diff --git a/packages/cli/test/relay-attempt-artifacts.test.ts b/packages/cli/test/relay-attempt-artifacts.test.ts new file mode 100644 index 00000000..3edaabfe --- /dev/null +++ b/packages/cli/test/relay-attempt-artifacts.test.ts @@ -0,0 +1,131 @@ +import { afterEach, beforeEach, expect, mock, spyOn, test } from "bun:test"; +import * as fs from "node:fs"; +import { randomUUID } from "node:crypto"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { getInbox, MAX_INBOX_MESSAGES, sendMessage } from "../src/utils/mail.js"; +import { deliverRelayedToLocal, relayAcceptanceReceiptPath } from "../src/utils/relay.js"; + +let root: string; +let saved: Record; + +beforeEach(() => { + root = fs.mkdtempSync(join(tmpdir(), "relay-attempt-artifacts-")); + saved = Object.fromEntries(["HOME", "TPS_MAIL_DIR"].map((key) => [key, process.env[key]])); + process.env.HOME = root; + process.env.TPS_MAIL_DIR = join(root, "mail"); +}); + +afterEach(() => { mock.restore(); }); + +afterEach(() => { + for (const [key, value] of Object.entries(saved)) { + if (value === undefined) delete process.env[key]; + else process.env[key] = value; + } + fs.rmSync(root, { recursive: true, force: true }); +}); + +function body() { + return { id: randomUUID(), from: "remote", to: "local", content: "payload", timestamp: new Date().toISOString() }; +} + +function records(): string[] { + return fs.readdirSync(getInbox("local").fresh).filter((file) => file.endsWith(".json")); +} + +function failOn(name: "renameSync" | "unlinkSync", match: (path: string) => boolean, message: string) { + const real = fs[name] as (...args: any[]) => unknown; + return spyOn(fs, name).mockImplementation(((...args: any[]) => { + if (match(String(args[name === "renameSync" ? 1 : 0]))) throw new Error(message); + return real(...args); + }) as any); +} + +function refusalOf(item: ReturnType): string { + try { deliverRelayedToLocal("remote", item); } catch (error) { return String(error); } + return ""; +} + +test("receipt temp: rename fails and the temp delete fails -> incomplete-recovery refusal", () => { + const item = body(); + const marker = relayAcceptanceReceiptPath("remote", item.id); + const rename = failOn("renameSync", (path) => path === marker, "injected receipt rename failure"); + const unlink = failOn("unlinkSync", (path) => path === `${marker}.tmp`, "injected temp delete failure"); + let refusal = ""; + try { refusal = refusalOf(item); } + finally { unlink.mockRestore(); rename.mockRestore(); } + expect(refusal).toContain("relay receipt recovery incomplete"); + expect(refusal).toContain("retry may duplicate"); + expect(refusal).not.toContain("retry delivery"); +}); + +test("receipt temp: rename fails and the temp delete succeeds -> retryable refusal, no temp left", () => { + const item = body(); + const marker = relayAcceptanceReceiptPath("remote", item.id); + const rename = failOn("renameSync", (path) => path === marker, "injected receipt rename failure"); + let refusal = ""; + try { refusal = refusalOf(item); } + finally { rename.mockRestore(); } + expect(refusal).toContain("relay acceptance receipt write failed; retry delivery"); + expect(fs.existsSync(`${marker}.tmp`)).toBe(false); + expect(records()).toEqual([]); +}); + +function fillInbox(): void { + for (let i = 0; i < MAX_INBOX_MESSAGES; i++) sendMessage("local", `filler-${i}`, "seeder"); +} + +test("dead-letter sidecar: publish fails and the sidecar delete fails -> incomplete-recovery refusal", () => { + const item = body(); + fillInbox(); + const dlq = getInbox("local").dlq; + const rename = failOn("renameSync", (path) => path.startsWith(dlq), "injected dead-letter publish failure"); + const unlink = failOn("unlinkSync", (path) => path.endsWith(".reason"), "injected sidecar delete failure"); + let refusal = ""; + try { refusal = refusalOf(item); } + finally { unlink.mockRestore(); rename.mockRestore(); } + expect(refusal).toContain("relay record recovery incomplete"); + expect(refusal).toContain("retry may duplicate"); + expect(refusal).not.toContain("retry delivery"); +}); + +test("dead-letter sidecar: publish fails and the sidecar delete succeeds -> no sidecar left", () => { + const item = body(); + fillInbox(); + const dlq = getInbox("local").dlq; + const rename = failOn("renameSync", (path) => path.startsWith(dlq), "injected dead-letter publish failure"); + let refusal = ""; + try { refusal = refusalOf(item); } + finally { rename.mockRestore(); } + expect(refusal).toContain("relay dead-letter failed; retry delivery"); + expect(fs.readdirSync(dlq).filter((file) => file.endsWith(".reason"))).toEqual([]); + expect(records()).toHaveLength(MAX_INBOX_MESSAGES); + expect(fs.readdirSync(getInbox("local").tmp).filter((file) => file.endsWith(".json"))).toEqual([]); +}); + +test("dead-letter temp record: publish fails and the temp delete succeeds -> no tmp record left, retryable refusal", () => { + const item = body(); + fillInbox(); + const inbox = getInbox("local"); + const rename = failOn("renameSync", (path) => path.startsWith(inbox.dlq), "injected dead-letter publish failure"); + let refusal = ""; + try { refusal = refusalOf(item); } + finally { rename.mockRestore(); } + expect(refusal).toContain("relay dead-letter failed; retry delivery"); + expect(fs.readdirSync(inbox.tmp).filter((file) => file.endsWith(".json"))).toEqual([]); +}); + +test("dead-letter temp record: publish fails and the temp delete fails -> incomplete-recovery refusal", () => { + const item = body(); + fillInbox(); + const inbox = getInbox("local"); + const rename = failOn("renameSync", (path) => path.startsWith(inbox.dlq), "injected dead-letter publish failure"); + const unlink = failOn("unlinkSync", (path) => path.startsWith(inbox.tmp) && path.endsWith(".json"), "injected temp delete failure"); + let refusal = ""; + try { refusal = refusalOf(item); } + finally { unlink.mockRestore(); rename.mockRestore(); } + expect(refusal).toContain("relay record recovery incomplete"); + expect(refusal).toContain("retry may duplicate"); + expect(refusal).not.toContain("retry delivery"); +}); diff --git a/packages/cli/test/relay-delivery-loss.test.ts b/packages/cli/test/relay-delivery-loss.test.ts index 4fc23b01..05695a2b 100644 --- a/packages/cli/test/relay-delivery-loss.test.ts +++ b/packages/cli/test/relay-delivery-loss.test.ts @@ -4,10 +4,10 @@ import * as fs from "node:fs"; import { randomUUID } from "node:crypto"; import { join } from "node:path"; import { tmpdir } from "node:os"; -import { gcMessages, getInbox, MAX_INBOX_MESSAGES, sendMessage } from "../src/utils/mail.js"; +import { gcMessages, getInbox, MAX_INBOX_MESSAGES, sendMessage, ackMessageAtPath } from "../src/utils/mail.js"; import { runBranch, writeBranchConf } from "../src/commands/branch.js"; import { runMail } from "../src/commands/mail.js"; -import { syncRemoteBranch, connectAndKeepAlive, deliverRelayedToLocal } from "../src/utils/relay.js"; +import { RELAY_ACCEPT_LOCK_STRIPES, syncRemoteBranch, connectAndKeepAlive, deliverRelayedToLocal, relayAcceptanceReceiptPath, pruneRelayAcceptanceReceipts } from "../src/utils/relay.js"; import * as ws from "../src/utils/ws-noise-transport.js"; import { generateKeyPair, initHostIdentity, registerBranch, saveKeyPair } from "../src/utils/identity.js"; import { drainOutbox, OUTBOX_RESEND_BASE_MS, queueOutboxMessage } from "../src/utils/outbox.js"; @@ -169,7 +169,7 @@ for (const entry of ["sync", "connect"] as const) { expect(rejected.some((r) => r.id === bodies[2]!.id)).toBe(true); }); - for (const fault of ["sidecar", "record", "crash-before-publish"] as const) { + for (const fault of ["sidecar", "record", "publish-throws"] as const) { test(`${fault} failure leaves the branch source unacked and retryable`, async () => { fillInbox(); const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "retry me", SEEDS))); @@ -178,9 +178,9 @@ for (const entry of ["sync", "connect"] as const) { await start(); const write = fs.writeFileSync; const rename = fs.renameSync; - const injected = fault === "crash-before-publish" + const injected = fault === "publish-throws" ? spyOn(fs, "renameSync").mockImplementation((src, dst) => { - if (String(dst).startsWith(inbox.dlq)) throw new Error("simulated crash before publish"); + if (String(dst).startsWith(inbox.dlq)) throw new Error("simulated publish failure"); return rename(src, dst); }) : spyOn(fs, "writeFileSync").mockImplementation((path, data, opts) => { @@ -191,11 +191,11 @@ for (const entry of ["sync", "connect"] as const) { expect(acks).toEqual([]); expect(drainOutbox(false).map((m) => m.id)).toEqual([body.id]); expect(jsonFiles(inbox.dlq)).toEqual([]); - if (fault === "crash-before-publish") expect(fs.readdirSync(inbox.dlq).some((f) => f.endsWith(".reason"))).toBe(true); + if (fault === "publish-throws") expect(fs.readdirSync(inbox.dlq).some((f) => f.endsWith(".reason"))).toBe(false); const logs = errors.mock.calls.flat().join("\n"); expect(logs).toContain(body.id); expect(logs).toContain("local"); - expect(logs).toContain(fault === "crash-before-publish" ? "simulated crash" : `injected ${fault}`); + expect(logs).toContain("relay dead-letter failed; retry delivery"); await emit(); expect(acks).toEqual([]); const later = Date.now() + OUTBOX_RESEND_BASE_MS; @@ -207,7 +207,7 @@ for (const entry of ["sync", "connect"] as const) { }); } - test("a one-shot local write failure is dead-lettered retryable and the next mail check delivers it", async () => { + test("a local write failure leaves the source unacked and its resend delivers once", async () => { const env = buildSignedEnvelope("remote", "local", "write fault", SEEDS); queue(JSON.stringify(env)); const inbox = getInbox("local"); @@ -224,11 +224,16 @@ for (const entry of ["sync", "connect"] as const) { }); try { await emit(); } finally { injected.mockRestore(); } expect(failed).toBe(true); - expect(acks.length).toBe(1); - expect(drainOutbox(false)).toEqual([]); - expect(jsonFiles(inbox.dlq).length).toBe(1); - for (const file of jsonFiles(inbox.dlq)) expect(fs.readFileSync(join(inbox.dlq, `${file}.reason`), "utf8")).toContain("class: storage-unavailable"); + expect(acks).toEqual([]); + expect(drainOutbox(false)).toHaveLength(1); + expect(jsonFiles(inbox.dlq)).toEqual([]); expect(jsonFiles(inbox.cur)).toEqual([]); + const later = Date.now() + OUTBOX_RESEND_BASE_MS; + const clock = spyOn(Date, "now").mockReturnValue(later); + try { await emit(); } finally { clock.mockRestore(); } + expect(acks).toHaveLength(1); + expect(drainOutbox(false)).toEqual([]); + expect(jsonFiles(inbox.fresh)).toHaveLength(1); const output = spyOn(console, "log").mockImplementation(() => {}); await runMail({ action: "check", agent: "local", json: true }); const delivered = JSON.parse(String(output.mock.calls.at(-1)![0])); @@ -313,9 +318,6 @@ for (const entry of ["sync", "connect"] as const) { const source = join(dir, "incomplete.json"); const incomplete = JSON.stringify({ relayDelivery: { branchId: "remote", id: body.id } }); fs.writeFileSync(source, incomplete); - const marker = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id); - fs.mkdirSync(join(marker, ".."), { recursive: true }); - fs.writeFileSync(marker, ""); const errors = spyOn(console, "error").mockImplementation(() => {}); await start(); const send = channel.send; @@ -371,7 +373,7 @@ for (const entry of ["sync", "connect"] as const) { const [file] = jsonFiles(dir).filter((name) => JSON.parse(fs.readFileSync(join(dir, name), "utf8")).relayDelivery?.id === body.id); const source = join(dir, file); const before = fs.readFileSync(source, "utf8"); - const marker = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id); + const marker = relayAcceptanceReceiptPath("remote", body.id); const beforeMarker = fs.readFileSync(marker, "utf8"); const errors = spyOn(console, "error").mockImplementation(() => {}); await start(); @@ -418,7 +420,7 @@ for (const entry of ["sync", "connect"] as const) { if (destination === "cur") { record.read = true; record.ackedAt = body.timestamp; } fs.writeFileSync(join(initialDir, file), JSON.stringify(record)); if (destination === "cur") fs.renameSync(join(initialDir, file), join(inbox.cur, file)); - if (prior === "no-marker") fs.rmSync(join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id)); + if (prior === "no-marker") fs.rmSync(relayAcceptanceReceiptPath("remote", body.id)); } fs.writeFileSync(join(inbox.dlq, "expired-canary.json"), JSON.stringify({ id: "expired", from: "remote", to: "local", body: "expired", timestamp: body.timestamp, read: false })); await start(); @@ -453,38 +455,29 @@ for (const entry of ["sync", "connect"] as const) { } } - for (const consumed of [false, true]) { - test(`marker failure then redelivery keeps one record in ${consumed ? "cur" : "new"}`, async () => { + test("receipt failure removes the new record before redelivery", async () => { const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "accept once", SEEDS))); const inbox = getInbox("local"); - const marker = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id); + const marker = relayAcceptanceReceiptPath("remote", body.id); fs.mkdirSync(`${marker}.tmp`, { recursive: true }); spyOn(console, "error").mockImplementation(() => {}); await start(); const msg: TpsMessage = { type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }; await deliverDirect(msg); expect(acks).toEqual([]); - expect(jsonFiles(inbox.fresh).length).toBe(1); + expect(jsonFiles(inbox.fresh)).toEqual([]); expect(fs.existsSync(marker)).toBe(false); expect(drainOutbox(false).map((m) => m.id)).toEqual([body.id]); - if (consumed) { - const output = spyOn(console, "log").mockImplementation(() => {}); - await runMail({ action: "check", agent: "local", json: true }); - expect(JSON.parse(String(output.mock.calls.at(-1)![0])).length).toBe(1); - expect(jsonFiles(inbox.fresh)).toEqual([]); - expect(jsonFiles(inbox.cur).length).toBe(1); - } fs.rmSync(`${marker}.tmp`, { recursive: true, force: true }); await deliverDirect(msg); expect(acks.map((ack) => (ack.body as { id: string }).id)).toEqual([body.id]); expect(jsonFiles(inbox.fresh).length + jsonFiles(inbox.cur).length).toBe(1); expect(drainOutbox(false)).toEqual([]); }); - } test("ACK failure then redelivery keeps one inbox record", async () => { const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "ACK retry", SEEDS))); - const marker = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id); + const marker = relayAcceptanceReceiptPath("remote", body.id); spyOn(console, "error").mockImplementation(() => {}); await start(); const msg: TpsMessage = { type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }; @@ -504,14 +497,125 @@ for (const entry of ["sync", "connect"] as const) { expect(drainOutbox(false)).toEqual([]); }); - for (const state of ["removed-unread", "consumed", "legacy"] as const) { - test(`existing ${state} marker with no record republishes before ACK`, async () => { + async function ackOnlyRecord(): Promise { + const inbox = getInbox("local"); + const [file] = jsonFiles(inbox.fresh); + ackMessageAtPath(join(inbox.fresh, file)); + expect(jsonFiles(inbox.fresh)).toEqual([]); + } + + test("an identical resend after ACK is a duplicate: one record, no second delivered", async () => { + const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "once", SEEDS))); + await start(); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }); + expect(acks.map((ack) => ack.body)).toEqual([{ id: body.id, accepted: true }]); + await ackOnlyRecord(); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 2, ts: new Date().toISOString(), body }); + expect(acks.map((ack) => ack.body)).toEqual([{ id: body.id, accepted: true }, { id: body.id, accepted: true }]); + expect(jsonFiles(getInbox("local").fresh)).toEqual([]); + expect(jsonFiles(getInbox("local").cur)).toEqual([]); + expect(drainOutbox(false)).toEqual([]); + }); + + test("a differing resend after ACK is refused without a second delivery", async () => { + const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "once", SEEDS))); + const errors = spyOn(console, "error").mockImplementation(() => {}); + await start(); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }); + await ackOnlyRecord(); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 2, ts: new Date().toISOString(), body: { ...body, content: "changed payload" } }); + expect(acks.map((ack) => ack.body)).toEqual([{ id: body.id, accepted: true }]); + expect(jsonFiles(getInbox("local").fresh)).toEqual([]); + expect(jsonFiles(getInbox("local").dlq)).toEqual([]); + expect(errors.mock.calls.flat().join("\n")).toContain("conflict"); + }); + + test("a flat per-branch marker from before receipts refuses resend after local ACK without ACK", async () => { + const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "legacy acceptance", SEEDS))); + const original = sendMessage("local", body.content, body.from, { branchId: "remote", id: body.id }, body.timestamp); + const marker = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id); + fs.mkdirSync(join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote"), { recursive: true }); + fs.writeFileSync(marker, ""); + const errors = spyOn(console, "error").mockImplementation(() => {}); + await start(); + ackMessageAtPath(original.filePath); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 2, ts: new Date().toISOString(), body }); + expect(jsonFiles(getInbox("local").fresh)).toEqual([]); + expect(jsonFiles(getInbox("local").dlq)).toEqual([]); + expect(acks).toEqual([]); + expect(drainOutbox(false).map((item) => item.id)).toContain(body.id); + expect(fs.existsSync(relayAcceptanceReceiptPath("remote", body.id))).toBe(false); + expect(errors.mock.calls.flat().join("\n")).toContain(`relayed delivery conflict for branch remote message ${body.id}: accepted before receipts existed; payload cannot be compared; sender must not retry`); + }); + + test("a receipt older than the prune bound no longer blocks a resend", async () => { + const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "late", SEEDS))); + await start(); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }); + await ackOnlyRecord(); + const marker = relayAcceptanceReceiptPath("remote", body.id); + process.env.TPS_RELAY_ACCEPT_RECEIPT_TTL_MS = "1000"; + const old = new Date(Date.now() - 5000); + fs.utimesSync(marker, old, old); + try { + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 2, ts: new Date().toISOString(), body }); + } finally { + delete process.env.TPS_RELAY_ACCEPT_RECEIPT_TTL_MS; + } + expect(acks.map((ack) => ack.body)).toEqual([{ id: body.id, accepted: true }, { id: body.id, accepted: true }]); + expect(jsonFiles(getInbox("local").fresh)).toHaveLength(1); + }); + + test("a whole old receipt bucket is removed when a new receipt is written", async () => { + const first = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "first accepted", SEEDS))); + await start(); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body: first }); + const accepted = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote"); + const firstMarker = relayAcceptanceReceiptPath("remote", first.id); + expect(fs.existsSync(firstMarker)).toBe(true); + process.env.TPS_RELAY_ACCEPT_RECEIPT_TTL_MS = "1000"; + const oldBucket = join(accepted, "2020-01-01"); + fs.mkdirSync(oldBucket); + fs.renameSync(firstMarker, join(oldBucket, first.id)); + const second = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "second accepted", SEEDS))); + try { + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 2, ts: new Date().toISOString(), body: second }); + } finally { + delete process.env.TPS_RELAY_ACCEPT_RECEIPT_TTL_MS; + } + expect(fs.existsSync(oldBucket)).toBe(false); + expect(fs.existsSync(relayAcceptanceReceiptPath("remote", second.id))).toBe(true); + expect(jsonFiles(getInbox("local").fresh)).toHaveLength(2); + }); + + test("a marker read failure refuses the delivery without an ACK", async () => { + const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "marker read", SEEDS))); + const errors = spyOn(console, "error").mockImplementation(() => {}); + await start(); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }); + await ackOnlyRecord(); + const marker = relayAcceptanceReceiptPath("remote", body.id); + const read = fs.readFileSync; + const fault = spyOn(fs, "readFileSync").mockImplementation((path, options) => { + if (String(path) === marker) throw Object.assign(new Error("injected marker read denied"), { code: "EACCES" }); + return read(path, options as BufferEncoding); + }); + try { + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 2, ts: new Date().toISOString(), body }); + } finally { + fault.mockRestore(); + } + expect(acks.length).toBe(1); + expect(jsonFiles(getInbox("local").fresh)).toEqual([]); + expect(errors.mock.calls.flat().join("\n")).toContain("injected marker read denied"); + }); + + for (const state of ["removed-unread", "consumed"] as const) { + test(`an existing ${state} receipt with no record dedups an identical resend and is ACKed`, async () => { const envelope = buildSignedEnvelope("remote", "local", state, SEEDS); const body = queue(JSON.stringify(envelope)); const inbox = getInbox("local"); - const marker = state === "legacy" - ? join(process.env.TPS_MAIL_DIR!, ".relay-accepted", body.id) - : join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id); + const marker = relayAcceptanceReceiptPath("remote", body.id); const client = new MailClient(process.env.TPS_MAIL_DIR!, undefined, "local", { async getAgent(id) { const seed = SEEDS[id as keyof typeof SEEDS]; @@ -521,54 +625,61 @@ for (const entry of ["sync", "connect"] as const) { spyOn(console, "error").mockImplementation(() => {}); await start(); const msg: TpsMessage = { type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }; - if (state === "legacy") { - fs.mkdirSync(join(process.env.TPS_MAIL_DIR!, ".relay-accepted"), { recursive: true }); - fs.writeFileSync(marker, ""); - } else { - const send = channel.send; - channel.send = async () => { throw new Error("injected lost ACK"); }; - await deliverDirect(msg); - channel.send = send; - expect(acks).toEqual([]); - expect(fs.readFileSync(marker, "utf8")).toBe(""); - expect(jsonFiles(inbox.fresh)).toHaveLength(1); - if (state === "consumed") { - expect(await client.checkNewMail()).toHaveLength(1); - expect(hasCommittedMessageId(join(process.env.TPS_MAIL_DIR!, "local"), envelope.messageId)).toBe(true); - expect(envelope.messageId).not.toBe(body.id); - } else { - expect(hasCommittedMessageId(join(process.env.TPS_MAIL_DIR!, "local"), envelope.messageId)).toBe(false); - } - const clock = spyOn(Date, "now").mockReturnValue(Date.now() + 2000); - try { expect(gcMessages("local", "24h", undefined, "1s")).toBe(1); } - finally { clock.mockRestore(); } - expect(jsonFiles(inbox.fresh)).toEqual([]); - expect(jsonFiles(inbox.cur)).toEqual([]); + const send = channel.send; + channel.send = async () => { throw new Error("injected lost ACK"); }; + await deliverDirect(msg); + channel.send = send; + expect(acks).toEqual([]); + expect(jsonFiles(inbox.fresh)).toHaveLength(1); + if (state === "consumed") { + expect(await client.checkNewMail()).toHaveLength(1); + expect(hasCommittedMessageId(join(process.env.TPS_MAIL_DIR!, "local"), envelope.messageId)).toBe(true); } + const clock = spyOn(Date, "now").mockReturnValue(Date.now() + 2000); + try { expect(gcMessages("local", "24h", undefined, "1s")).toBe(1); } + finally { clock.mockRestore(); } + expect(jsonFiles(inbox.fresh)).toEqual([]); + expect(jsonFiles(inbox.cur)).toEqual([]); + const recordsAtAck: number[] = []; - const send = channel.send; + const ackSend = channel.send; channel.send = async (ack) => { if (ack.type === MSG_MAIL_ACK) recordsAtAck.push(jsonFiles(inbox.fresh).length); - return send(ack); + return ackSend(ack); }; await deliverDirect(msg); - expect(recordsAtAck).toEqual([1]); - expect(jsonFiles(inbox.fresh)).toHaveLength(1); + expect(recordsAtAck).toEqual([0]); + expect(jsonFiles(inbox.fresh)).toEqual([]); expect(jsonFiles(inbox.cur)).toEqual([]); expect(jsonFiles(inbox.dlq)).toEqual([]); expect(acks.map((ack) => ack.body)).toEqual([{ id: body.id, accepted: true }]); expect(drainOutbox(false)).toEqual([]); - const promoted = await client.checkNewMail(); - expect(promoted).toHaveLength(state === "consumed" ? 0 : 1); - if (state === "consumed") { - const [file] = jsonFiles(inbox.dlq); - expect(fs.readFileSync(join(inbox.dlq, `${file}.reason`), "utf8")).toContain("class: replay"); - } + expect(await client.checkNewMail()).toEqual([]); + expect(fs.existsSync(marker)).toBe(true); }); } + test("an empty pre-receipt marker with no record refuses without an ACK", async () => { + const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "legacy", SEEDS))); + const inbox = getInbox("local"); + const legacy = join(process.env.TPS_MAIL_DIR!, ".relay-accepted"); + const marker = join(legacy, body.id); + const errors = spyOn(console, "error").mockImplementation(() => {}); + await start(); + fs.mkdirSync(legacy, { recursive: true }); + fs.writeFileSync(marker, ""); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }); + expect(acks).toEqual([]); + expect(errors.mock.calls.flat().join("\n")).toContain("conflict"); + expect(jsonFiles(inbox.fresh)).toEqual([]); + expect(jsonFiles(inbox.cur)).toEqual([]); + expect(jsonFiles(inbox.dlq)).toEqual([]); + expect(drainOutbox(false).map((item) => item.id)).toEqual([body.id]); + expect(fs.readFileSync(marker, "utf8")).toBe(""); + }); + for (const signature of ["valid", "invalid"] as const) { - test(`resend with an unrelated consumed id and ${signature} signature republishes before ACK`, async () => { + test(`a delivery whose envelope reuses a consumed message id is refused at promotion (${signature} signature)`, async () => { const original = buildSignedEnvelope("remote", "local", "unrelated", SEEDS); sendMessage("local", JSON.stringify(original), "remote"); const client = new MailClient(process.env.TPS_MAIL_DIR!, undefined, "local", { @@ -582,29 +693,9 @@ for (const entry of ["sync", "connect"] as const) { if (signature === "invalid") envelope.body = "tampered"; const body = queue(JSON.stringify(envelope)); const inbox = getInbox("local"); - const marker = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id); spyOn(console, "error").mockImplementation(() => {}); await start(); - const msg: TpsMessage = { type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }; - const send = channel.send; - channel.send = async () => { throw new Error("injected lost ACK"); }; - await deliverDirect(msg); - channel.send = send; - expect(acks).toEqual([]); - expect(fs.readFileSync(marker, "utf8")).toBe(""); - const clock = spyOn(Date, "now").mockReturnValue(Date.now() + 2000); - try { expect(gcMessages("local", "24h", undefined, "1s")).toBe(2); } - finally { clock.mockRestore(); } - expect(jsonFiles(inbox.fresh)).toEqual([]); - expect(jsonFiles(inbox.cur)).toEqual([]); - expect(jsonFiles(inbox.dlq)).toEqual([]); - const recordsAtAck: number[] = []; - channel.send = async (ack) => { - if (ack.type === MSG_MAIL_ACK) recordsAtAck.push(jsonFiles(inbox.fresh).length); - return send(ack); - }; - await deliverDirect(msg); - expect(recordsAtAck).toEqual([1]); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }); expect(acks.map((ack) => ack.body)).toEqual([{ id: body.id, accepted: true }]); expect(drainOutbox(false)).toEqual([]); expect(await client.checkNewMail()).toEqual([]); @@ -633,7 +724,7 @@ for (const entry of ["sync", "connect"] as const) { expect(acks).toEqual([]); expect(fs.readFileSync(source, "utf8")).toBe(before); expect(jsonFiles(inbox.fresh).length + jsonFiles(inbox.cur).length + jsonFiles(inbox.dlq).length).toBe(1); - expect(fs.existsSync(join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id))).toBe(false); + expect(fs.existsSync(relayAcceptanceReceiptPath("remote", body.id))).toBe(false); expect(errors.mock.calls.flat().join("\n")).toContain(`relayed record read failed: ${source}`); expect(drainOutbox(false).map((item) => item.id)).toContain(body.id); if (dir === "dlq") fs.writeFileSync(`${source}.reason`, "class: inbox-full\n"); @@ -647,7 +738,7 @@ for (const entry of ["sync", "connect"] as const) { test("inbox and DLQ write failures leave no marker or ACK, then retry writes one record", async () => { const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "write retry", SEEDS))); const inbox = getInbox("local"); - const marker = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id); + const marker = relayAcceptanceReceiptPath("remote", body.id); spyOn(console, "error").mockImplementation(() => {}); await start(); const msg: TpsMessage = { type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }; @@ -675,13 +766,21 @@ for (const entry of ["sync", "connect"] as const) { spyOn(console, "error").mockImplementation(() => {}); await start(); fs.mkdirSync(join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote"), { recursive: true }); + for (let stripe = 0; stripe < RELAY_ACCEPT_LOCK_STRIPES; stripe++) { + fs.mkdirSync(join(process.env.TPS_MAIL_DIR!, ".relay-accept-locks", String(stripe)), { recursive: true }); + } const sync = fs.fsyncSync; const trace: number[] = []; const capture = spyOn(fs, "fsyncSync").mockImplementation((fd) => { trace.push(fd); return sync(fd); }); const initial = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "sync trace", SEEDS))); - try { await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body: initial }); } + try { + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body: initial }); + trace.length = 0; + const warmed = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "warmed sync trace", SEEDS))); + await deliverDirect({ type: MSG_MAIL_DELIVER, seq: 2, ts: new Date().toISOString(), body: warmed }); + } finally { capture.mockRestore(); } - expect(acks.length).toBe(1); + expect(acks.length).toBe(2); expect(trace.length).toBeGreaterThanOrEqual(destination === "inbox" ? 4 : 5); for (let step = 1; step <= trace.length; step++) { const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", `sync fault ${step}`, SEEDS))); @@ -694,8 +793,10 @@ for (const entry of ["sync", "connect"] as const) { return sync(fd); }); try { await deliverDirect(msg); } finally { fault.mockRestore(); } - expect(calls).toBe(step); + expect(calls).toBeGreaterThanOrEqual(step); expect(acks.length).toBe(beforeAcks); + expect(fs.existsSync(relayAcceptanceReceiptPath("remote", body.id))).toBe(false); + expect(jsonFiles(inbox.fresh).length + jsonFiles(inbox.dlq).length).toBe(beforeRecords); expect(drainOutbox(false).map((m) => m.id)).toContain(body.id); await deliverDirect(msg); expect(acks.length).toBe(beforeAcks + 1); @@ -705,18 +806,18 @@ for (const entry of ["sync", "connect"] as const) { }); } - test("DLQ record before marker failure is reused on redelivery", async () => { + test("DLQ record is removed after receipt failure and written on redelivery", async () => { fillInbox(); const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "DLQ retry", SEEDS))); const inbox = getInbox("local"); - const marker = join(process.env.TPS_MAIL_DIR!, ".relay-accepted", "by-branch", "remote", body.id); + const marker = relayAcceptanceReceiptPath("remote", body.id); fs.mkdirSync(`${marker}.tmp`, { recursive: true }); spyOn(console, "error").mockImplementation(() => {}); await start(); const msg: TpsMessage = { type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body }; await deliverDirect(msg); expect(acks).toEqual([]); - expect(jsonFiles(inbox.dlq).length).toBe(1); + expect(jsonFiles(inbox.dlq)).toEqual([]); fs.rmSync(`${marker}.tmp`, { recursive: true }); await deliverDirect(msg); expect(acks.length).toBe(1); diff --git a/packages/cli/test/relay-review.test.ts b/packages/cli/test/relay-review.test.ts index d45f85a7..ee949c93 100644 --- a/packages/cli/test/relay-review.test.ts +++ b/packages/cli/test/relay-review.test.ts @@ -4,7 +4,7 @@ import { randomUUID } from "node:crypto"; import { tmpdir } from "node:os"; import { dirname, join } from "node:path"; import { checkMessages, getInbox, listMessages, promote, sendMessage } from "../src/utils/mail.js"; -import { deliverRelayedToLocal } from "../src/utils/relay.js"; +import { deliverRelayedToLocal, relayAcceptanceReceiptPath } from "../src/utils/relay.js"; import { runMail } from "../src/commands/mail.js"; import { buildSignedEnvelope, pubkeyFromSeed, writeKeyFile } from "./helpers/stub-flair.js"; @@ -76,7 +76,7 @@ describe("relay review regressions", () => { const [file] = fs.readdirSync(inbox.fresh); expect(JSON.parse(fs.readFileSync(join(inbox.fresh, file!), "utf8")).body).toBe(message.content); const directories = synced.filter((path) => !path.endsWith(".json") && !path.endsWith(message.id) && !path.endsWith(".tmp")); - expect(new Set(directories)).toEqual(new Set([inbox.tmp, inbox.fresh, accepted])); + for (const path of [inbox.tmp, inbox.fresh, join(accepted, new Date().toISOString().slice(0, 10))]) expect(directories).toContain(path); }); test("relay publication syncs parents of newly created directories", () => { @@ -86,10 +86,7 @@ describe("relay review regressions", () => { const inbox = getInbox("local"); const acceptedRoot = join(mail, ".relay-accepted"); const directories = synced.filter((path) => !path.endsWith(".json") && !path.endsWith(message.id) && !path.endsWith(".tmp")); - expect(new Set(directories)).toEqual(new Set([ - root, mail, inbox.root, inbox.tmp, inbox.fresh, - acceptedRoot, join(acceptedRoot, "by-branch"), join(acceptedRoot, "by-branch", "remote"), - ])); + for (const path of [root, mail, inbox.root, inbox.tmp, inbox.fresh, join(acceptedRoot, "by-branch", "remote", new Date().toISOString().slice(0, 10))]) expect(directories).toContain(path); }); test("a retry syncs directory entries left pending by a creation sync failure", () => { @@ -97,16 +94,21 @@ describe("relay review regressions", () => { const acceptedRoot = join(mail, ".relay-accepted"); const synced = traceSync(acceptedRoot, true); const message = body(); - expect(() => deliverRelayedToLocal("remote", message)).toThrow("mail fsync failed"); + expect(() => deliverRelayedToLocal("remote", message)).toThrow("relay acceptance receipt write failed; retry delivery"); synced.length = 0; - expect(deliverRelayedToLocal("remote", message)).toBe(false); + expect(deliverRelayedToLocal("remote", message)).toBe(true); expect(synced).toContain(acceptedRoot); expect(synced).toContain(mail); - expect(fs.readFileSync(join(acceptedRoot, "by-branch", "remote", message.id), "utf8")).toBe(""); + expect(JSON.parse(fs.readFileSync(relayAcceptanceReceiptPath("remote", message.id), "utf8"))).toEqual({ + from: message.from, + to: message.to, + body: message.content, + timestamp: message.timestamp, + }); }); for (const dir of ["new", "cur", "dlq"] as const) { - test(`a truncated ${dir} record is preserved and reported while a marked delivery is republished`, () => { + test(`a truncated ${dir} record is preserved and reported; a marked delivery with no usable receipt is refused`, () => { const inbox = getInbox("local"); const message = body(); const corrupt = join(inbox.root, dir, "truncated.json"); @@ -114,9 +116,10 @@ describe("relay review regressions", () => { fs.writeFileSync(corrupt, raw); const accepted = join(mail, ".relay-accepted", "by-branch", "remote"); fs.mkdirSync(accepted, { recursive: true }); - fs.writeFileSync(join(accepted, message.id), ""); + fs.mkdirSync(join(accepted, new Date().toISOString().slice(0, 10)), { recursive: true }); + fs.writeFileSync(relayAcceptanceReceiptPath("remote", message.id), ""); const errors = spyOn(console, "error").mockImplementation(() => {}); - expect(deliverRelayedToLocal("remote", message)).toBe(true); + expect(() => deliverRelayedToLocal("remote", message)).toThrow(`relayed delivery conflict for branch remote message ${message.id}`); const quarantine = join(inbox.root, "quarantine"); const files = fs.readdirSync(quarantine).filter((file) => file.endsWith(".json")); expect(files).toHaveLength(1); @@ -125,7 +128,8 @@ describe("relay review regressions", () => { expect(fs.readFileSync(`${path}.reason`, "utf8")).toContain(corrupt); expect(errors.mock.calls.flat().join("\n")).toContain(corrupt); expect(fs.readdirSync(join(inbox.root, dir))).not.toContain("truncated.json"); - expect(deliverRelayedToLocal("remote", message)).toBe(false); + fs.rmSync(relayAcceptanceReceiptPath("remote", message.id)); + expect(deliverRelayedToLocal("remote", message)).toBe(true); const records = fs.readdirSync(inbox.fresh).filter((file) => file.endsWith(".json")); expect(records).toHaveLength(1); expect(JSON.parse(fs.readFileSync(join(inbox.fresh, records[0]!), "utf8")).body).toBe(message.content); @@ -150,7 +154,7 @@ describe("relay review regressions", () => { expect(errors.mock.calls.flat().join("\n")).toContain(source); const records = fs.readdirSync(inbox.fresh).filter((file) => file !== "unreadable.json" && file.endsWith(".json")); expect(records).toHaveLength(0); - expect(fs.existsSync(join(mail, ".relay-accepted", "by-branch", "remote", message.id))).toBe(false); + expect(fs.existsSync(relayAcceptanceReceiptPath("remote", message.id))).toBe(false); }); test("mail list and CLI list use receipt time with a timestamp fallback", async () => {