Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
@@ -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.
Original file line number Diff line number Diff line change
@@ -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.
49 changes: 40 additions & 9 deletions packages/cli/src/utils/mail.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand All @@ -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<string>();

export function mkdirMailDirectory(path: string): void {
Expand Down Expand Up @@ -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
Expand All @@ -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();
Expand Down Expand Up @@ -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);
}

Expand Down
Loading
Loading