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
3 changes: 3 additions & 0 deletions .changelog/unreleased/fixed-508-relay-hex-delivery-ids.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
- **Relay delivery and ACK accept the 64-hex delivery ids that GitHub-webhook outbox records carry.**

`tps office sync` and `tps office connect` log a relayed delivery that fails wire-schema validation, naming the branch and the invalid fields.
34 changes: 25 additions & 9 deletions packages/cli/src/utils/relay.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,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";
const { log: slog, warn: swarn, error: serror } = snooplogg("tps:relay");


Expand Down Expand Up @@ -369,10 +370,14 @@ export function handleIncomingMail(branchId: string, msg: TpsMessage): void {
});
}

export function deliverRelayedToLocal(body: MailDeliverBody): boolean {
export function deliverRelayedToLocal(branchId: string, body: MailDeliverBody): boolean {
MailDeliverBodySchema.shape.id.parse(body.id);
const marker = join(getMailDir(), ".relay-accepted", body.id);
if (existsSync(marker)) return false;
if (!/^[a-zA-Z0-9_-]+$/.test(branchId)) throw new Error(`invalid branch id for relayed message ${body.id}`);
// 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);
// A marker in the earlier unscoped layout sits directly in .relay-accepted/.
if (existsSync(marker) || existsSync(join(getMailDir(), ".relay-accepted", body.id))) return false;
let delivered: boolean;
try {
sendMessage(body.to, body.content, body.from);
Expand All @@ -399,20 +404,26 @@ export function deliverRelayedToLocal(body: MailDeliverBody): boolean {
delivered = false;
}
try {
mkdirSync(join(getMailDir(), ".relay-accepted"), { recursive: true });
mkdirSync(acceptedDir, { recursive: true });
writeFileSync(marker, "", "utf-8");
} catch (e: unknown) {
console.error(`[relay] could not record message ${body.id} as accepted: ${e instanceof Error ? e.message : String(e)}`);
}
return delivered;
}

async function acceptRelayedMail(channel: TransportChannel, msg: TpsMessage, body: MailDeliverBody): Promise<boolean> {
const delivered = deliverRelayedToLocal(body);
async function acceptRelayedMail(channel: TransportChannel, branchId: string, msg: TpsMessage, body: MailDeliverBody): Promise<boolean> {
const delivered = deliverRelayedToLocal(branchId, body);
await channel.send({ type: MSG_MAIL_ACK, seq: msg.seq, ts: new Date().toISOString(), body: { id: body.id, accepted: true } });
return delivered;
}

/** Logs a MAIL_DELIVER the schema refused, naming the branch and the failing fields, never their values. */
function logRefusedDelivery(branchId: string, error: ZodError): void {
const fields = [...new Set(error.issues.map((issue) => issue.path.map(String).join(".") || "body"))].join(", ");
console.error(`[relay] refused a MAIL_DELIVER from branch ${branchId}: invalid ${fields}`);
}

export function startRelay(agentId: string): () => void {
assertAgent(agentId);

Expand Down Expand Up @@ -634,8 +645,11 @@ export async function syncRemoteBranch(branchId: string): Promise<{ received: nu
const handler = (msg: TpsMessage) => {
if (msg.type !== MSG_MAIL_DELIVER) return;
const parsed = MailDeliverBodySchema.safeParse(msg.body);
if (!parsed.success) return;
void acceptRelayedMail(channel, msg, parsed.data).then((delivered) => {
if (!parsed.success) {
logRefusedDelivery(branchId, parsed.error);
return;
}
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)}`);
Expand Down Expand Up @@ -753,9 +767,11 @@ export async function connectAndKeepAlive(
state.lastHeartbeatAck = now; // any traffic = alive
const parsed = MailDeliverBodySchema.safeParse(msg.body);
if (parsed.success) {
void acceptRelayedMail(channel, msg, parsed.data).catch((error: unknown) => {
void acceptRelayedMail(channel, branchId, msg, parsed.data).catch((error: unknown) => {
console.error(`[relay] acceptance failed for message ${parsed.data.id} to ${parsed.data.to}: ${String(error)}`);
});
} else {
logRefusedDelivery(branchId, parsed.error);
}
} else {
state.lastHeartbeatAck = now;
Expand Down
7 changes: 5 additions & 2 deletions packages/cli/src/utils/wire-mail.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,11 @@ export const MSG_HTTP_RESPONSE = 0x21;

const SAFE_ID = /^[a-zA-Z0-9._-]{1,64}$/;

/** A relayed delivery id: a UUID, or the 64-hex id of a GitHub-webhook outbox record (outbox.ts). */
const DeliveryIdSchema = z.union([z.string().uuid(), z.string().regex(/^[a-f0-9]{64}$/)]);

export const MailDeliverBodySchema = z.object({
id: z.string().uuid(),
id: DeliveryIdSchema,
from: z.string().regex(SAFE_ID, "Invalid sender identifier"),
to: z.string().regex(SAFE_ID, "Invalid recipient identifier"),
content: z.string(),
Expand All @@ -19,7 +22,7 @@ export const MailDeliverBodySchema = z.object({
export type MailDeliverBody = z.infer<typeof MailDeliverBodySchema>;

export const MailAckBodySchema = z.object({
id: z.string().uuid(),
id: DeliveryIdSchema,
accepted: z.boolean(),
reason: z.string().optional(),
});
Expand Down
75 changes: 75 additions & 0 deletions packages/cli/test/relay-delivery-loss.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import { describe, expect, test, beforeEach, afterEach, spyOn, mock } from "bun:test";
import * as fs from "node:fs";
import { randomUUID } from "node:crypto";
import { join } from "node:path";
import { tmpdir } from "node:os";
import { getInbox, MAX_INBOX_MESSAGES, sendMessage } from "../src/utils/mail.js";
Expand Down Expand Up @@ -265,5 +266,79 @@ for (const entry of ["sync", "connect"] as const) {
expect(acks.length).toBe(1);
expect(drainOutbox(false)).toEqual([]);
});

for (const shape of ["uuid", "64-hex"] as const) {
test(`a ${shape} delivery id produces one inbox record, one ACK and an empty outbox`, async () => {
const env = buildSignedEnvelope("remote", "local", `${shape} delivery`, SEEDS);
const id = shape === "64-hex" ? "ab".repeat(32) : queue(JSON.stringify(env)).id;
if (shape === "64-hex") queueOutboxMessage("local", JSON.stringify(env), "remote", id);
await start();
await emit();
expect(jsonFiles(getInbox("local").fresh).length).toBe(1);
expect(acks.map((ack) => (ack.body as { id: string }).id)).toEqual([id]);
expect(drainOutbox(false)).toEqual([]);
const replay: TpsMessage = { type: MSG_MAIL_DELIVER, seq: 2, ts: new Date().toISOString(), body: { id, from: "remote", to: "local", content: JSON.stringify(env), timestamp: new Date().toISOString() } };
for (const handler of handlers) handler(replay);
await Bun.sleep(0);
expect(jsonFiles(getInbox("local").fresh).length).toBe(1);
expect(acks.length).toBe(2);
});

test(`the same ${shape} delivery id from two branches is delivered for each`, async () => {
const kp = generateKeyPair();
registerBranch("remote-b", kp.signing.publicKey, undefined, kp.encryption.publicKey);
const dirB = join(root, ".tps", "branch-office", "remote-b");
fs.mkdirSync(dirB, { recursive: true });
fs.writeFileSync(join(dirB, "remote.json"), JSON.stringify({ host: "unused", port: 1, transport: "ws" }));
await start();
const before = handlers.size;
if (entry === "sync") {
completion = Promise.all([completion, syncRemoteBranch("remote-b")]);
} else {
const stopA = stop;
const stopB = await connectAndKeepAlive("remote-b");
stop = async () => { await stopB(); await stopA?.(); };
}
for (let i = 0; handlers.size === before && i < 100; i++) await Bun.sleep(10);
expect(handlers.size).toBeGreaterThan(before);
const id = shape === "64-hex" ? "cd".repeat(32) : randomUUID();
const content = JSON.stringify(buildSignedEnvelope("remote", "local", "same id", SEEDS));
const msg: TpsMessage = { type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body: { id, from: "remote", to: "local", content, timestamp: new Date().toISOString() } };
for (const handler of handlers) handler(msg);
await Bun.sleep(0);
expect(jsonFiles(getInbox("local").fresh).length).toBe(2);
expect(acks.length).toBe(2);
});
}

test("a malformed delivery id is refused and logged without the payload", async () => {
const errors = spyOn(console, "error").mockImplementation(() => {});
await start();
const malformed = ["not-a-delivery-id", "AB".repeat(32), "ab".repeat(32) + "a"];
for (const id of malformed) {
const msg: TpsMessage = { type: MSG_MAIL_DELIVER, seq: 1, ts: new Date().toISOString(), body: { id, from: "remote", to: "local", content: "payload-text", timestamp: new Date().toISOString() } };
for (const handler of handlers) handler(msg);
}
await Bun.sleep(0);
expect(acks).toEqual([]);
expect(jsonFiles(getInbox("local").fresh)).toEqual([]);
const refusals = errors.mock.calls.flat().map(String).filter((line) => line.includes("[relay] refused a MAIL_DELIVER from branch remote: invalid id"));
expect(refusals.length).toBe(malformed.length);
const logs = errors.mock.calls.flat().map(String).join("\n");
expect(logs).not.toContain("payload-text");
for (const id of malformed) expect(logs).not.toContain(id);
});

test("a UUID accepted under the earlier unscoped marker layout is acknowledged without a second inbox record", async () => {
const body = queue(JSON.stringify(buildSignedEnvelope("remote", "local", "accepted before upgrade", SEEDS)));
const legacy = join(process.env.TPS_MAIL_DIR!, ".relay-accepted");
fs.mkdirSync(legacy, { recursive: true });
fs.writeFileSync(join(legacy, body.id), "");
await start();
await emit();
expect(jsonFiles(getInbox("local").fresh)).toEqual([]);
expect(acks.length).toBe(1);
expect(drainOutbox(false)).toEqual([]);
});
});
}
Loading