diff --git a/README.md b/README.md index 361625b..df12f97 100644 --- a/README.md +++ b/README.md @@ -78,11 +78,19 @@ aliases.on("message", (m) => { // time, give it a reply-able name (idempotent; a no-op if already active). await aliases.ensure("alice"); +// Deliver alice's message to the relayed session from alice's own alias +// (starting it if needed), so the session sees alice as the sender and its +// own reply-to-sender comes back on the "message" handler above rather than +// to the relay. Rejects if the alias cannot start, or if it cannot send. +const { msgId } = await aliases.send("alice", { pid: relayedSessionPid }, "ping"); + // …later, once a correspondent is no longer relevant: await aliases.retire("alice"); await aliases.stopAll(); ``` +Sending from the alias rather than from the relay's own peer is what makes a relayed session's natural reply reach the right correspondent: a session replies to whoever sent it a message, so a message delivered from the relay comes back to the relay with nothing to say which correspondent it answers. + ## Limitations - **Same-process constraint**: receipts and idle notices only reach the process that owns the peer's listening socket (the protocol verifies return addresses via kernel peer-pids). Do not split `CcPeer` listening and sending across processes or differently-owned workers. diff --git a/src/adapters/node/alias-worker.ts b/src/adapters/node/alias-worker.ts index c266be5..6cea3ce 100644 --- a/src/adapters/node/alias-worker.ts +++ b/src/adapters/node/alias-worker.ts @@ -4,7 +4,11 @@ import process from "node:process"; import { CcPeer, type InboundMessage } from "../../cc-peer.js"; -import { AliasCommandSchema } from "../../schemas/alias-ipc.js"; +import { CcPeerError } from "../../errors.js"; +import { + AliasCommandSchema, + type AliasSendCommand, +} from "../../schemas/alias-ipc.js"; let peer: CcPeer | undefined; @@ -19,6 +23,10 @@ async function handleCommand(raw: unknown): Promise { process.exit(0); return; } + if (raw.type === "send") { + await handleSend(raw); + return; + } let created: CcPeer; try { created = await CcPeer.create({ @@ -37,3 +45,28 @@ async function handleCommand(raw: unknown): Promise { }); process.send?.({ type: "started" }); } + +async function handleSend(command: Readonly): Promise { + const active = peer; + if (active === undefined) { + process.send?.({ + type: "send_failed", + requestId: command.requestId, + code: "NOT_STARTED", + message: "alias peer has not started", + }); + return; + } + try { + const { msgId } = await active.send(command.target, command.body); + process.send?.({ type: "sent", requestId: command.requestId, msgId }); + } catch (error) { + process.send?.({ + type: "send_failed", + requestId: command.requestId, + code: error instanceof CcPeerError ? error.code : "UNKNOWN", + message: + error instanceof Error ? error.message : "unknown alias send failure", + }); + } +} diff --git a/src/adapters/node/forked-alias-process.integration.test.ts b/src/adapters/node/forked-alias-process.integration.test.ts index ff5634c..d10c003 100644 --- a/src/adapters/node/forked-alias-process.integration.test.ts +++ b/src/adapters/node/forked-alias-process.integration.test.ts @@ -1,4 +1,4 @@ -import { describe, expect, test } from "vitest"; +import { describe, expect, test, vi } from "vitest"; import { fork, type ChildProcess } from "node:child_process"; import { fileURLToPath } from "node:url"; import { mkdtemp, mkdir, readdir, readFile, writeFile } from "node:fs/promises"; @@ -9,7 +9,8 @@ import { once } from "node:events"; import { randomUUID } from "node:crypto"; import { ForkedAliasProcess } from "./forked-alias-process.js"; -import { AliasStartError } from "../../errors.js"; +import { CcPeer, type InboundMessage } from "../../cc-peer.js"; +import { AliasSendError, AliasStartError } from "../../errors.js"; import { REAL_PROCESS_SPAWN_TEST_TIMEOUT_MS } from "../../test/timeouts.js"; import type { RegistryEntry } from "../../schemas/registry.js"; import type { PeerKeyFile } from "../../schemas/keyfile.js"; @@ -116,6 +117,65 @@ describe("ForkedAliasProcess default fork() fallback", () => { ); }); +/** A pid high enough that no live session in the temp home has published an inbox for it, so resolving it reaches the "no auth key" failure rather than a real socket. */ +const PID_WITH_NO_PUBLISHED_INBOX = 999_999; + +describe("ForkedAliasProcess.send against a real receiving peer", () => { + test( + "delivers a message the receiver attributes to the alias, not to the relay", + async () => { + const homeDir = await tempHome(); + const socketDir = join(homeDir, "socks"); + const target = await CcPeer.create({ + name: "send-target", + homeDir, + socketDir, + }); + const received: InboundMessage[] = []; + target.on("message", (message: InboundMessage) => { + received.push(message); + }); + const proc = makeAliasProcess(); + await proc.start({ name: "dana-send-test", homeDir, socketDir }); + const { msgId } = await proc.send( + { pid: process.pid }, + "pong from the alias", + ); + expect(msgId).not.toBe(""); + await vi.waitFor(() => { + expect(received).toHaveLength(1); + }); + expect(received[0]).toEqual( + expect.objectContaining({ + body: "pong from the alias", + fromName: "dana-send-test", + }), + ); + await proc.stop(); + await target.stop(); + }, + REAL_PROCESS_SPAWN_TEST_TIMEOUT_MS, + ); + + test( + "rejects with the alias peer's own failure code when the target has no live inbox", + async () => { + const homeDir = await tempHome(); + const socketDir = join(homeDir, "socks"); + const proc = makeAliasProcess(); + await proc.start({ name: "erin-send-failure-test", homeDir, socketDir }); + const failing = proc.send( + { pid: PID_WITH_NO_PUBLISHED_INBOX }, + "pong into the void", + ); + await expect(failing).rejects.toThrow(AliasSendError); + await expect(failing).rejects.toThrow(/NO_LIVE_INBOX/); + await proc.stop(); + }, + REAL_PROCESS_SPAWN_TEST_TIMEOUT_MS, + ); +}); + describe("ForkedAliasProcess against the real alias-worker source", () => { test( "registers a discoverable peer and relays an inbound reply", diff --git a/src/adapters/node/forked-alias-process.ts b/src/adapters/node/forked-alias-process.ts index 5e6dcd9..2dc1b4a 100644 --- a/src/adapters/node/forked-alias-process.ts +++ b/src/adapters/node/forked-alias-process.ts @@ -7,11 +7,20 @@ import type { } from "../../ports/alias-process.js"; import { AliasMessageEventSchema, + AliasSendFailedEventSchema, + AliasSentEventSchema, AliasStartedEventSchema, type AliasMessageEvent, } from "../../schemas/alias-ipc.js"; -import { AliasStartError } from "../../errors.js"; -import type { InboundMessage } from "../../cc-peer.js"; +import { AliasSendError, AliasStartError } from "../../errors.js"; +import { newMsgId } from "../../domain/ids.js"; +import type { InboundMessage, PeerRef } from "../../cc-peer.js"; + +/** One in-flight send command, settled when its own acknowledgement arrives or the worker exits. */ +interface PendingSend { + resolve: (result: Readonly<{ msgId: string }>) => void; + reject: (error: Error) => void; +} export interface ForkedAliasProcessDeps { /** @@ -28,6 +37,7 @@ export class ForkedAliasProcess implements AliasProcess { readonly events = new EventEmitter(); private readonly forkFn: typeof fork; private readonly workerPath: string; + private readonly pendingSends = new Map(); private child: ChildProcess | undefined; constructor(deps: Readonly) { @@ -41,9 +51,14 @@ export class ForkedAliasProcess implements AliasProcess { child.on("message", (raw: unknown) => { if (AliasMessageEventSchema.is(raw)) { this.events.emit("message", toInboundMessage(raw)); + return; } + this.settleSend(raw); }); child.on("exit", () => { + this.failPendingSends( + "alias worker exited before the send was acknowledged", + ); this.events.emit("exit"); }); await new Promise((resolve, reject) => { @@ -78,10 +93,59 @@ export class ForkedAliasProcess implements AliasProcess { }); } - async stop(): Promise { + async send( + target: Readonly, + body: string, + ): Promise<{ msgId: string }> { + const child = this.runningChild(); + if (child === undefined) { + throw new AliasSendError("alias worker is not running"); + } + const requestId = newMsgId(); + return new Promise<{ msgId: string }>((resolve, reject) => { + this.pendingSends.set(requestId, { resolve, reject }); + child.send({ type: "send", requestId, target, body }); + }); + } + + /** Settles the one in-flight send an acknowledgement names. An acknowledgement for an unknown request is ignored: the send it belongs to was already settled by the worker exiting. */ + private settleSend(raw: unknown): void { + if (AliasSentEventSchema.is(raw)) { + this.takePendingSend(raw.requestId)?.resolve({ msgId: raw.msgId }); + return; + } + if (AliasSendFailedEventSchema.is(raw)) { + this.takePendingSend(raw.requestId)?.reject( + new AliasSendError(`${raw.code}: ${raw.message}`), + ); + } + } + + private takePendingSend(requestId: string): PendingSend | undefined { + const pending = this.pendingSends.get(requestId); + this.pendingSends.delete(requestId); + return pending; + } + + private failPendingSends(reason: string): void { + const pending = [...this.pendingSends.values()]; + this.pendingSends.clear(); + for (const send of pending) { + send.reject(new AliasSendError(reason)); + } + } + + /** The forked child while it is still running, or undefined before start() and once it has ended: a command written past that point reaches a channel nothing is reading. */ + private runningChild(): ChildProcess | undefined { const child = this.child; + if (child === undefined) return undefined; + if (child.exitCode !== null || child.signalCode !== null) return undefined; + return child; + } + + async stop(): Promise { + const child = this.runningChild(); if (child === undefined) return; - if (child.exitCode !== null || child.signalCode !== null) return; await new Promise((resolve) => { child.once("exit", () => { resolve(); diff --git a/src/adapters/node/forked-alias-process.unit.test.ts b/src/adapters/node/forked-alias-process.unit.test.ts index e7b45e1..be9a4f6 100644 --- a/src/adapters/node/forked-alias-process.unit.test.ts +++ b/src/adapters/node/forked-alias-process.unit.test.ts @@ -3,7 +3,8 @@ import { EventEmitter } from "node:events"; import type { ChildProcess } from "node:child_process"; import { ForkedAliasProcess } from "./forked-alias-process.js"; -import { AliasStartError } from "../../errors.js"; +import { AliasSendError, AliasStartError } from "../../errors.js"; +import { AliasSendCommandSchema } from "../../schemas/alias-ipc.js"; /** A fake ChildProcess: just enough of the EventEmitter + send() surface for the adapter's own protocol logic, driven manually by each test. */ class FakeChild extends EventEmitter { @@ -196,6 +197,135 @@ describe("ForkedAliasProcess message relay", () => { }); }); +/** The requestId the adapter generated for its most recent send command, read back off the fake child's own IPC log so tests can answer the exact request. */ +function lastSendRequestId(child: FakeChild): string { + const command = child.sent.at(-1); + if (!AliasSendCommandSchema.is(command)) { + throw new Error("the last IPC command was not a send command"); + } + return command.requestId; +} + +async function startedProcess(): Promise<{ + child: FakeChild; + proc: ForkedAliasProcess; +}> { + const child = new FakeChild(); + const { process: proc } = makeForkedAliasProcess(child); + const pending = proc.start({ name: "alice" }); + child.emit("message", { type: "started" }); + await pending; + return { child, proc }; +} + +describe("ForkedAliasProcess.send", () => { + test("sends a send command carrying the target and body", async () => { + const { child, proc } = await startedProcess(); + const pending = proc.send({ pid: 42 }, "pong"); + expect(child.sent.at(-1)).toEqual({ + type: "send", + requestId: lastSendRequestId(child), + target: { pid: 42 }, + body: "pong", + }); + child.emit("message", { + type: "sent", + requestId: lastSendRequestId(child), + msgId: "m1", + }); + await expect(pending).resolves.toEqual({ msgId: "m1" }); + }); + + test("gives each send its own request id", async () => { + const { child, proc } = await startedProcess(); + const first = proc.send({ pid: 42 }, "one"); + const firstRequestId = lastSendRequestId(child); + const second = proc.send({ pid: 42 }, "two"); + const secondRequestId = lastSendRequestId(child); + expect(secondRequestId).not.toBe(firstRequestId); + child.emit("message", { + type: "sent", + requestId: secondRequestId, + msgId: "m2", + }); + child.emit("message", { + type: "sent", + requestId: firstRequestId, + msgId: "m1", + }); + await expect(first).resolves.toEqual({ msgId: "m1" }); + await expect(second).resolves.toEqual({ msgId: "m2" }); + }); + + test("rejects with the worker's own failure code and message", async () => { + const { child, proc } = await startedProcess(); + const pending = proc.send({ name: "bob" }, "pong"); + child.emit("message", { + type: "send_failed", + requestId: lastSendRequestId(child), + code: "NO_LIVE_INBOX", + message: "no auth key published", + }); + await expect(pending).rejects.toThrow(AliasSendError); + await expect(pending).rejects.toThrow( + "NO_LIVE_INBOX: no auth key published", + ); + }); + + test("ignores an acknowledgement naming a request it is not waiting on", async () => { + const { child, proc } = await startedProcess(); + const pending = proc.send({ pid: 42 }, "pong"); + const requestId = lastSendRequestId(child); + child.emit("message", { + type: "sent", + requestId: "someone-elses-request", + msgId: "stray", + }); + child.emit("message", { + type: "send_failed", + requestId: "someone-elses-request", + code: "TRANSPORT", + message: "stray failure", + }); + child.emit("message", { type: "sent", requestId, msgId: "m1" }); + await expect(pending).resolves.toEqual({ msgId: "m1" }); + }); + + test("rejects every in-flight send when the worker exits", async () => { + const { child, proc } = await startedProcess(); + const first = proc.send({ pid: 42 }, "one"); + const second = proc.send({ pid: 42 }, "two"); + child.emitExit(0); + await expect(first).rejects.toThrow(/exited before the send/); + await expect(second).rejects.toThrow(AliasSendError); + }); + + test("throws before the worker has been started", async () => { + const child = new FakeChild(); + const { process: proc } = makeForkedAliasProcess(child); + await expect(proc.send({ pid: 42 }, "pong")).rejects.toThrow( + AliasSendError, + ); + expect(child.sent).toEqual([]); + }); + + test("throws once the worker has exited", async () => { + const { child, proc } = await startedProcess(); + child.emitExit(0); + await expect(proc.send({ pid: 42 }, "pong")).rejects.toThrow( + "alias worker is not running", + ); + }); + + test("throws once the worker has been killed by a signal", async () => { + const { child, proc } = await startedProcess(); + child.signalCode = "SIGKILL"; + await expect(proc.send({ pid: 42 }, "pong")).rejects.toThrow( + AliasSendError, + ); + }); +}); + describe("ForkedAliasProcess.stop", () => { test("sends a stop command and resolves once the child exits", async () => { const child = new FakeChild(); diff --git a/src/alias-pool.ts b/src/alias-pool.ts index 6ac0de3..8c50446 100644 --- a/src/alias-pool.ts +++ b/src/alias-pool.ts @@ -4,7 +4,7 @@ import { fileURLToPath } from "node:url"; import { ForkedAliasProcess } from "./adapters/node/forked-alias-process.js"; import type { PathConfig } from "./adapters/node/paths.js"; import type { AliasProcess } from "./ports/alias-process.js"; -import type { InboundMessage } from "./cc-peer.js"; +import type { InboundMessage, PeerRef } from "./cc-peer.js"; export interface AliasPoolOptions extends PathConfig { logger?: (message: string) => void; @@ -44,7 +44,7 @@ interface Deps { */ export class AliasPool extends EventEmitter { private readonly active = new Map(); - private readonly pending = new Map>(); + private readonly pending = new Map>(); private readonly log: (message: string) => void; constructor( @@ -67,22 +67,37 @@ export class AliasPool extends EventEmitter { /** Idempotent: a no-op if the alias is already active, and dedup'd if another ensure() for the same name is already in flight. */ async ensure(name: string): Promise { - if (this.active.has(name)) return; + await this.acquire(name); + } + + /** + * Sends `body` to `target` from the named alias's own peer identity, starting that alias first if it is not already running, so the recipient sees the correspondent rather than the relay as the sender. Resolves with the id the alias gave the message; rejects with an `AliasStartError` if the alias could not be started, or an `AliasSendError` if it refused the send or its process ended before acknowledging it. + */ + async send( + name: string, + target: Readonly, + body: string, + ): Promise<{ msgId: string }> { + const proc = await this.acquire(name); + return proc.send(target, body); + } + + /** The single start path behind ensure() and send(): returns the running process for the name, starting one only if neither an active nor an in-flight process already exists. */ + private async acquire(name: string): Promise { + const active = this.active.get(name); + if (active !== undefined) return active; const inFlight = this.pending.get(name); - if (inFlight !== undefined) { - await inFlight; - return; - } + if (inFlight !== undefined) return inFlight; const started = this.startAlias(name); this.pending.set(name, started); try { - await started; + return await started; } finally { this.pending.delete(name); } } - private async startAlias(name: string): Promise { + private async startAlias(name: string): Promise { const proc = this.deps.spawn(); proc.events.on("message", (message: InboundMessage) => { this.emit("message", { alias: name, ...message } satisfies AliasMessage); @@ -102,6 +117,7 @@ export class AliasPool extends EventEmitter { }); this.active.set(name, proc); this.log(`alias ${name} active`); + return proc; } /** Waits for any in-flight ensure() of the same name to settle first, so a retire() issued while an alias is still starting stops it once (and if) it becomes active. A no-op for a name that is neither active nor pending. */ diff --git a/src/alias-pool.unit.test.ts b/src/alias-pool.unit.test.ts index 0cf1260..7394295 100644 --- a/src/alias-pool.unit.test.ts +++ b/src/alias-pool.unit.test.ts @@ -2,7 +2,9 @@ import { describe, expect, test, vi } from "vitest"; import { EventEmitter } from "node:events"; import { AliasPool, workerExtensionFor } from "./alias-pool.js"; +import { AliasSendError } from "./errors.js"; import type { AliasProcess } from "./ports/alias-process.js"; +import type { PeerRef } from "./cc-peer.js"; describe("workerExtensionFor", () => { test("returns .cjs for a module URL ending in .cjs", () => { @@ -22,9 +24,13 @@ describe("workerExtensionFor", () => { class FakeAliasProcess implements AliasProcess { readonly events = new EventEmitter(); startCalls: unknown[] = []; + sendCalls: { target: Readonly; body: string }[] = []; stopCalls = 0; private resolveStart: (() => void) | undefined; private rejectStart: ((error: Error) => void) | undefined; + private resolveSend: + ((result: Readonly<{ msgId: string }>) => void) | undefined; + private rejectSend: ((error: Error) => void) | undefined; async start(options: unknown): Promise { this.startCalls.push(options); @@ -34,6 +40,17 @@ class FakeAliasProcess implements AliasProcess { }); } + async send( + target: Readonly, + body: string, + ): Promise<{ msgId: string }> { + this.sendCalls.push({ target, body }); + return new Promise<{ msgId: string }>((resolve, reject) => { + this.resolveSend = resolve; + this.rejectSend = reject; + }); + } + async stop(): Promise { this.stopCalls += 1; return Promise.resolve(); @@ -46,6 +63,14 @@ class FakeAliasProcess implements AliasProcess { failStart(error: Error): void { this.rejectStart?.(error); } + + finishSend(msgId: string): void { + this.resolveSend?.({ msgId }); + } + + failSend(error: Error): void { + this.rejectSend?.(error); + } } function makePool(procs: readonly FakeAliasProcess[]): { @@ -186,6 +211,76 @@ describe("AliasPool message and exit relay", () => { }); }); +describe("AliasPool.send", () => { + test("starts the alias on demand and sends from it, resolving with the message id", async () => { + const proc = new FakeAliasProcess(); + const { pool, spawnCount } = makePool([proc]); + const pending = pool.send("alice", { pid: 42 }, "pong"); + proc.finishStart(); + await vi.waitFor(() => { + expect(proc.sendCalls).toHaveLength(1); + }); + proc.finishSend("m1"); + await expect(pending).resolves.toEqual({ msgId: "m1" }); + expect(proc.sendCalls).toEqual([{ target: { pid: 42 }, body: "pong" }]); + expect(spawnCount()).toBe(1); + expect(pool.activeAliases()).toEqual(["alice"]); + }); + + test("reuses an already-active alias rather than spawning a second process", async () => { + const proc = new FakeAliasProcess(); + const { pool, spawnCount } = makePool([proc]); + const ensuring = pool.ensure("alice"); + proc.finishStart(); + await ensuring; + const pending = pool.send("alice", { name: "bob" }, "pong"); + await vi.waitFor(() => { + expect(proc.sendCalls).toHaveLength(1); + }); + proc.finishSend("m2"); + await expect(pending).resolves.toEqual({ msgId: "m2" }); + expect(spawnCount()).toBe(1); + }); + + test("waits for an in-flight ensure() of the same name instead of spawning again", async () => { + const proc = new FakeAliasProcess(); + const { pool, spawnCount } = makePool([proc]); + const ensuring = pool.ensure("alice"); + const pending = pool.send("alice", { address: "uds:/tmp/9.sock" }, "pong"); + proc.finishStart(); + await ensuring; + await vi.waitFor(() => { + expect(proc.sendCalls).toHaveLength(1); + }); + proc.finishSend("m3"); + await expect(pending).resolves.toEqual({ msgId: "m3" }); + expect(spawnCount()).toBe(1); + }); + + test("rejects when the alias cannot be started", async () => { + const proc = new FakeAliasProcess(); + const { pool } = makePool([proc]); + const pending = pool.send("alice", { pid: 42 }, "pong"); + proc.failStart(new Error("boom")); + await expect(pending).rejects.toThrow("boom"); + expect(proc.sendCalls).toEqual([]); + expect(pool.activeAliases()).toEqual([]); + }); + + test("propagates the alias process's own send failure", async () => { + const proc = new FakeAliasProcess(); + const { pool } = makePool([proc]); + const pending = pool.send("alice", { pid: 42 }, "pong"); + proc.finishStart(); + await vi.waitFor(() => { + expect(proc.sendCalls).toHaveLength(1); + }); + proc.failSend(new AliasSendError("NO_LIVE_INBOX: no auth key published")); + await expect(pending).rejects.toThrow(AliasSendError); + await expect(pending).rejects.toThrow(/no auth key published/); + }); +}); + describe("AliasPool.retire", () => { test("stops an active alias and removes it", async () => { const proc = new FakeAliasProcess(); diff --git a/src/errors.ts b/src/errors.ts index a42e153..4e686dc 100644 --- a/src/errors.ts +++ b/src/errors.ts @@ -65,3 +65,11 @@ export class AliasStartError extends CcPeerError { this.name = "AliasStartError"; } } + +/** A message could not be sent from a reply alias: its own peer refused the send, or its worker process was gone (or went away) before acknowledging the command. */ +export class AliasSendError extends CcPeerError { + constructor(message: string) { + super("ALIAS_SEND_FAILED", message); + this.name = "AliasSendError"; + } +} diff --git a/src/errors.unit.test.ts b/src/errors.unit.test.ts index 2b1d79d..2d73910 100644 --- a/src/errors.unit.test.ts +++ b/src/errors.unit.test.ts @@ -9,6 +9,7 @@ import { NotStartedError, ProtocolError, AliasStartError, + AliasSendError, } from "./errors.js"; describe("error taxonomy", () => { @@ -22,6 +23,7 @@ describe("error taxonomy", () => { [new NotStartedError("s"), "NOT_STARTED"], [new ProtocolError("p"), "PROTOCOL"], [new AliasStartError("a"), "ALIAS_START_FAILED"], + [new AliasSendError("a"), "ALIAS_SEND_FAILED"], ] as const; for (const [error, code] of cases) { expect(error).toBeInstanceOf(CcPeerError); diff --git a/src/ports/alias-process.ts b/src/ports/alias-process.ts index 8404121..d5fd223 100644 --- a/src/ports/alias-process.ts +++ b/src/ports/alias-process.ts @@ -1,5 +1,6 @@ import type { EventEmitter } from "node:events"; import type { PathConfig } from "../adapters/node/paths.js"; +import type { PeerRef } from "../cc-peer.js"; export interface AliasStartOptions extends PathConfig { name: string; @@ -12,5 +13,9 @@ export interface AliasStartOptions extends PathConfig { export interface AliasProcess { readonly events: EventEmitter; start: (options: Readonly) => Promise; + /** + * Sends `body` to `target` from the alias's own peer identity, so the recipient sees the alias as the sender and replies to it natively. Resolves with the id the alias's peer gave the message; rejects with an `AliasSendError` when the alias refuses the send or its process ends before acknowledging it. + */ + send: (target: Readonly, body: string) => Promise<{ msgId: string }>; stop: () => Promise; } diff --git a/src/schemas/alias-ipc.ts b/src/schemas/alias-ipc.ts index 89974f0..a590316 100644 --- a/src/schemas/alias-ipc.ts +++ b/src/schemas/alias-ipc.ts @@ -20,8 +20,35 @@ export const AliasStopCommandSchema = defineSchema( ); export type AliasStopCommand = z.infer; +/** + * How an alias addresses an outbound send, mirroring `CcPeer`'s own `PeerRef` union. Each variant is strict so a target carrying fields from more than one variant is rejected outright rather than resolved by whichever branch happens to match first. + */ +export const AliasSendTargetSchema = defineSchema( + z.union([ + z.strictObject({ pid: z.number() }), + z.strictObject({ name: z.string().min(1) }), + z.strictObject({ address: z.string().min(1) }), + ]), +); +export type AliasSendTarget = z.infer; + +/** Asks the alias's own CcPeer to send `body` to `target`; the worker answers with a sent or send_failed event carrying the same `requestId`. */ +export const AliasSendCommandSchema = defineSchema( + z.object({ + type: z.literal("send"), + requestId: z.string().min(1), + target: AliasSendTargetSchema, + body: z.string(), + }), +); +export type AliasSendCommand = z.infer; + export const AliasCommandSchema = defineSchema( - z.union([AliasStartCommandSchema, AliasStopCommandSchema]), + z.union([ + AliasStartCommandSchema, + AliasStopCommandSchema, + AliasSendCommandSchema, + ]), ); export type AliasCommand = z.infer; @@ -46,7 +73,33 @@ export const AliasMessageEventSchema = defineSchema( ); export type AliasMessageEvent = z.infer; +/** Acknowledges the send command with the same `requestId`, carrying the id the alias's own peer gave the message. */ +export const AliasSentEventSchema = defineSchema( + z.object({ + type: z.literal("sent"), + requestId: z.string().min(1), + msgId: z.string().min(1), + }), +); +export type AliasSentEvent = z.infer; + +/** Reports a send the alias could not make; `code` is the originating `CcPeerError`'s own machine-readable code where one was thrown. */ +export const AliasSendFailedEventSchema = defineSchema( + z.object({ + type: z.literal("send_failed"), + requestId: z.string().min(1), + code: z.string().min(1), + message: z.string(), + }), +); +export type AliasSendFailedEvent = z.infer; + export const AliasEventSchema = defineSchema( - z.union([AliasStartedEventSchema, AliasMessageEventSchema]), + z.union([ + AliasStartedEventSchema, + AliasMessageEventSchema, + AliasSentEventSchema, + AliasSendFailedEventSchema, + ]), ); export type AliasEvent = z.infer; diff --git a/src/schemas/alias-ipc.unit.test.ts b/src/schemas/alias-ipc.unit.test.ts index 9b336b8..e77a535 100644 --- a/src/schemas/alias-ipc.unit.test.ts +++ b/src/schemas/alias-ipc.unit.test.ts @@ -6,6 +6,10 @@ import { AliasStartedEventSchema, AliasMessageEventSchema, AliasEventSchema, + AliasSendTargetSchema, + AliasSendCommandSchema, + AliasSentEventSchema, + AliasSendFailedEventSchema, } from "./alias-ipc.js"; describe("AliasStartCommandSchema", () => { @@ -48,10 +52,84 @@ describe("AliasStopCommandSchema", () => { }); }); +const VALID_SEND_COMMAND = { + type: "send", + requestId: "req-1", + target: { pid: 4242 }, + body: "pong", +}; + +describe("AliasSendTargetSchema", () => { + test("accepts each of the three ways a peer can be addressed", () => { + expect(AliasSendTargetSchema.is({ pid: 4242 })).toBe(true); + expect(AliasSendTargetSchema.is({ name: "alice" })).toBe(true); + expect( + AliasSendTargetSchema.is({ address: "uds:/tmp/cc-socks/9.sock" }), + ).toBe(true); + }); + + test("rejects an empty name or address", () => { + expect(AliasSendTargetSchema.is({ name: "" })).toBe(false); + expect(AliasSendTargetSchema.is({ address: "" })).toBe(false); + }); + + test("rejects a target mixing two ways of addressing the same peer", () => { + expect(AliasSendTargetSchema.is({ pid: 4242, name: "alice" })).toBe(false); + }); + + test("rejects a target addressing nothing at all", () => { + expect(AliasSendTargetSchema.is({})).toBe(false); + }); +}); + +describe("AliasSendCommandSchema", () => { + test("accepts a well-formed send command", () => { + expect(AliasSendCommandSchema.is(VALID_SEND_COMMAND)).toBe(true); + }); + + test("accepts an empty body", () => { + expect(AliasSendCommandSchema.is({ ...VALID_SEND_COMMAND, body: "" })).toBe( + true, + ); + }); + + test("rejects an empty requestId", () => { + expect( + AliasSendCommandSchema.is({ ...VALID_SEND_COMMAND, requestId: "" }), + ).toBe(false); + }); + + test("rejects a missing body", () => { + expect( + AliasSendCommandSchema.is({ + type: "send", + requestId: "req-1", + target: { pid: 4242 }, + }), + ).toBe(false); + }); + + test("rejects an unaddressable target", () => { + expect( + AliasSendCommandSchema.is({ + ...VALID_SEND_COMMAND, + target: { socket: "/tmp/cc-socks/9.sock" }, + }), + ).toBe(false); + }); + + test("rejects a wrong type discriminant", () => { + expect( + AliasSendCommandSchema.is({ ...VALID_SEND_COMMAND, type: "start" }), + ).toBe(false); + }); +}); + describe("AliasCommandSchema", () => { - test("accepts either command shape", () => { + test("accepts every command shape", () => { expect(AliasCommandSchema.is({ type: "start", name: "alice" })).toBe(true); expect(AliasCommandSchema.is({ type: "stop" })).toBe(true); + expect(AliasCommandSchema.is(VALID_SEND_COMMAND)).toBe(true); }); test("rejects an unrelated shape", () => { @@ -118,10 +196,77 @@ describe("AliasMessageEventSchema", () => { }); }); +const VALID_SENT_EVENT = { + type: "sent", + requestId: "req-1", + msgId: "msg-1", +}; + +const VALID_SEND_FAILED_EVENT = { + type: "send_failed", + requestId: "req-1", + code: "NO_LIVE_INBOX", + message: "no auth key published", +}; + +describe("AliasSentEventSchema", () => { + test("accepts a well-formed acknowledgement", () => { + expect(AliasSentEventSchema.is(VALID_SENT_EVENT)).toBe(true); + }); + + test("rejects an empty requestId or msgId", () => { + expect( + AliasSentEventSchema.is({ ...VALID_SENT_EVENT, requestId: "" }), + ).toBe(false); + expect(AliasSentEventSchema.is({ ...VALID_SENT_EVENT, msgId: "" })).toBe( + false, + ); + }); + + test("rejects a wrong type discriminant", () => { + expect(AliasSentEventSchema.is({ ...VALID_SENT_EVENT, type: "send" })).toBe( + false, + ); + }); +}); + +describe("AliasSendFailedEventSchema", () => { + test("accepts a well-formed failure", () => { + expect(AliasSendFailedEventSchema.is(VALID_SEND_FAILED_EVENT)).toBe(true); + }); + + test("accepts an empty message, since a code alone still identifies the failure", () => { + expect( + AliasSendFailedEventSchema.is({ + ...VALID_SEND_FAILED_EVENT, + message: "", + }), + ).toBe(true); + }); + + test("rejects an empty code", () => { + expect( + AliasSendFailedEventSchema.is({ ...VALID_SEND_FAILED_EVENT, code: "" }), + ).toBe(false); + }); + + test("rejects a missing requestId", () => { + expect( + AliasSendFailedEventSchema.is({ + type: "send_failed", + code: "TRANSPORT", + message: "gone", + }), + ).toBe(false); + }); +}); + describe("AliasEventSchema", () => { - test("accepts either event shape", () => { + test("accepts every event shape", () => { expect(AliasEventSchema.is({ type: "started" })).toBe(true); expect(AliasEventSchema.is(VALID_MESSAGE_EVENT)).toBe(true); + expect(AliasEventSchema.is(VALID_SENT_EVENT)).toBe(true); + expect(AliasEventSchema.is(VALID_SEND_FAILED_EVENT)).toBe(true); }); test("rejects an unrelated shape", () => {