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
8 changes: 8 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
35 changes: 34 additions & 1 deletion src/adapters/node/alias-worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -19,6 +23,10 @@ async function handleCommand(raw: unknown): Promise<void> {
process.exit(0);
return;
}
if (raw.type === "send") {
await handleSend(raw);
return;
}
let created: CcPeer;
try {
created = await CcPeer.create({
Expand All @@ -37,3 +45,28 @@ async function handleCommand(raw: unknown): Promise<void> {
});
process.send?.({ type: "started" });
}

async function handleSend(command: Readonly<AliasSendCommand>): Promise<void> {
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",
});
}
}
64 changes: 62 additions & 2 deletions src/adapters/node/forked-alias-process.integration.test.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand All @@ -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";
Expand Down Expand Up @@ -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",
Expand Down
72 changes: 68 additions & 4 deletions src/adapters/node/forked-alias-process.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
/**
Expand All @@ -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<string, PendingSend>();
private child: ChildProcess | undefined;

constructor(deps: Readonly<ForkedAliasProcessDeps>) {
Expand All @@ -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<void>((resolve, reject) => {
Expand Down Expand Up @@ -78,10 +93,59 @@ export class ForkedAliasProcess implements AliasProcess {
});
}

async stop(): Promise<void> {
async send(
target: Readonly<PeerRef>,
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<void> {
const child = this.runningChild();
if (child === undefined) return;
if (child.exitCode !== null || child.signalCode !== null) return;
await new Promise<void>((resolve) => {
child.once("exit", () => {
resolve();
Expand Down
Loading
Loading