From 1d7420eb51be3c0bb7c4a5bf3f375b380d13d83f Mon Sep 17 00:00:00 2001 From: Tim Disney Date: Mon, 20 Jul 2026 09:12:44 -0700 Subject: [PATCH] Phase 3 review fixes (round 3, daemon control + egress) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - #6 cli: reconcileOrphans transitions orphaned running rows via markCrashed (attempts preserved, cooldown/gave_up honored) instead of reset; tested - #7 dispatch: close awaiting_input stale-index race — an awaiting_input row is re-dispatchable only once the index shows the assignee's question (caught up) AND a reply (no longer awaitingInput); tested - #11 config: memory limit threaded config->dispatcher->TurnInput (4g floor in turn.ts) - #3 egress: verify a pre-existing internal network is actually --internal, else fail closed - #14 egress: poll proxy readiness (State.Running) before declaring egress up, else fail closed - #9 proxy: CONNECT restricted to port 443 - #12 egress/proxy/config: explicit empty allowlist = deny-all end to end (--allow-none) Co-Authored-By: Claude Fable 5 --- packages/daemon/src/cli.ts | 28 ++++- packages/daemon/src/config.ts | 6 + packages/daemon/src/dispatch.ts | 24 ++++ packages/daemon/src/egress.ts | 54 ++++++++- packages/daemon/src/proxy.ts | 46 ++++++-- packages/daemon/test/cli.test.mjs | 29 +++++ packages/daemon/test/config.test.mjs | 20 ++++ packages/daemon/test/dispatch.test.mjs | 92 ++++++++++++++++ packages/daemon/test/egress.test.mjs | 140 ++++++++++++++++++++++-- packages/daemon/test/proxy.test.mjs | 75 ++++++++++++- packages/daemon/test/reconcile.test.mjs | 63 +++++++++++ 11 files changed, 548 insertions(+), 29 deletions(-) create mode 100644 packages/daemon/test/reconcile.test.mjs diff --git a/packages/daemon/src/cli.ts b/packages/daemon/src/cli.ts index 2d51bda..5d89885 100644 --- a/packages/daemon/src/cli.ts +++ b/packages/daemon/src/cli.ts @@ -127,6 +127,9 @@ export interface ResolvedRunConfig { proxyImage: string allowlist: string[] gitSchemes: string[] + /** Container memory limit (docker `--memory` syntax, e.g. `4g`). No default at this layer — + * `turn.ts` supplies the `4g` floor when unset (a turn must never run unbounded). */ + memory?: string } const DEFAULT_INTERNAL_NETWORK = 'radial-internal' @@ -171,6 +174,7 @@ export function resolveRunConfig(run: DaemonRunConfig): ResolvedRunConfig { proxyImage: run.proxyImage ?? DEFAULT_PROXY_IMAGE, allowlist: run.allowlist ?? ['api.anthropic.com'], gitSchemes: run.gitSchemes ?? ['https'], + ...(run.memory !== undefined ? { memory: run.memory } : {}), } } @@ -236,6 +240,19 @@ export async function reconcileOrphan( console.log(`reconciled orphaned turn ${requestUri} (label ${label ?? 'unknown'})`) } +/** Startup orphan reconciliation: any row still "running" belongs to a prior process that died + * (or was killed) mid-turn. Its container (if any survived) is killed via `reconcileOrphan`, and + * the row is transitioned with `ledger.markCrashed` — NOT `ledger.reset` — so a daemon that dies + * mid-turn on every restart still counts each restart as a failed attempt (preserving `attempts`, + * applying the cooldown, and eventually `gave_up`) instead of silently resetting the retry counter + * and bypassing the retry bound forever. Exported so tests can drive it directly. */ +export async function reconcileOrphans(runner: ContainerRunner, ledger: TurnLedger): Promise { + for (const row of ledger.running()) { + await reconcileOrphan(runner, row.containerLabel, row.requestUri) + ledger.markCrashed(row.requestUri) + } +} + async function runCommand(args: string[]): Promise { const { config, run } = await requireRunConfig(values(args, '--config')[0]) const interval = Number(values(args, '--interval')[0] ?? 5_000) @@ -257,12 +274,10 @@ async function runCommand(args: string[]): Promise { const runner = new DockerRunner() // Startup orphan reconciliation: any row still "running" belongs to a prior process that died - // (or was killed) mid-turn — its container (if any survived) is killed and the row cleared so - // the request is dispatchable again on the next tick. - for (const row of ledger.running()) { - await reconcileOrphan(runner, row.containerLabel, row.requestUri) - ledger.reset(row.requestUri) - } + // (or was killed) mid-turn — its container (if any survived) is killed and the row transitioned + // to "crashed"/"gave_up" (attempts preserved) so a repeatedly-crashing daemon still respects the + // retry bound instead of resetting it on every restart. See `reconcileOrphans`. + await reconcileOrphans(runner, ledger) // Egress enforcement (design §13, "enforced NOW, no bridge fallback"): on by default (opt out // with `run.egress: false` in radial.json). There are only two allowed outcomes: egress is @@ -317,6 +332,7 @@ async function runCommand(args: string[]): Promise { timeoutMs: run.timeoutMs, ...(dispatcherNetwork !== undefined ? { network: dispatcherNetwork } : {}), allowedSchemes: run.gitSchemes, + ...(run.memory !== undefined ? { memory: run.memory } : {}), runTurn: boundRunTurn, log: (message) => console.log(message), }) diff --git a/packages/daemon/src/config.ts b/packages/daemon/src/config.ts index 66aaab5..750c55d 100644 --- a/packages/daemon/src/config.ts +++ b/packages/daemon/src/config.ts @@ -49,6 +49,10 @@ export interface DaemonRunConfig { proxyImage?: string allowlist?: string[] gitSchemes?: string[] + /** Container memory limit (docker `--memory` syntax, e.g. `4g`/`512m`). Unset means "use the + * daemon's default floor" (see `turn.ts`, which floors at `4g` — a turn must never run + * unbounded); this only overrides that default, it never widens it above what an operator sets. */ + memory?: string } const object = (value: unknown): value is Record => @@ -133,6 +137,7 @@ export function parseRunConfig(value: unknown): DaemonRunConfig { const proxyImage = text(value.proxyImage, 'run.proxyImage') const allowlist = strings(value.allowlist, 'run.allowlist') const gitSchemes = strings(value.gitSchemes, 'run.gitSchemes') + const memory = text(value.memory, 'run.memory') return { spaces, ...(stateDir !== undefined ? { stateDir } : {}), @@ -147,6 +152,7 @@ export function parseRunConfig(value: unknown): DaemonRunConfig { ...(proxyImage !== undefined ? { proxyImage } : {}), ...(allowlist !== undefined ? { allowlist } : {}), ...(gitSchemes !== undefined ? { gitSchemes } : {}), + ...(memory !== undefined ? { memory } : {}), } } diff --git a/packages/daemon/src/dispatch.ts b/packages/daemon/src/dispatch.ts index bff50fd..2038785 100644 --- a/packages/daemon/src/dispatch.ts +++ b/packages/daemon/src/dispatch.ts @@ -57,6 +57,26 @@ export function selectDispatchable( const actor = actors.select(assignee, 'plan') if (!actor) continue + // Close the awaiting_input stale-index race: a ledger row can flip to `awaiting_input` (and + // become `eligible()`) before the next index snapshot has ingested the assignee's question. + // If the index hasn't caught up yet, `goalView.awaitingInput` is still empty and this request + // would otherwise be re-selected and redispatched before any human ever sees the question. + // Require that the index actually contains the assignee's (non-declining) question for this + // exact request/cid before treating an `awaiting_input` row as re-dispatchable at all; the + // existing `awaitingInput` check below still blocks the case where the question IS ingested + // but nobody has replied yet. Net effect: an awaiting_input request is dispatchable again only + // once the index shows both the question (caught up) and a reply (no longer awaitingInput). + if (ledger.get(request.uri)?.state === 'awaiting_input') { + const questionInIndex = goalView.messages.some( + (m) => + m.did === assignee && + m.value.re && + m.value.re.uri === request.uri && + m.value.re.cid === request.cid && + m.value.declines !== true, + ) + if (!questionInIndex) continue // stale index: don't redispatch ahead of the ingested question + } if (goalView.awaitingInput.includes(request.uri)) continue if (!ledger.eligible(request.uri, now)) continue @@ -82,6 +102,9 @@ export interface DispatcherDeps { timeoutMs: number network?: string allowedSchemes: string[] + /** Container memory limit (docker `--memory` syntax, e.g. `4g`), threaded straight into + * `TurnInput.memory`. Unset lets `turn.ts` apply its own `4g` floor. */ + memory?: string /** Injected wrapper around turn.ts's runTurn that binds TurnDeps (runner/harness/secrets/etc). */ runTurn: (input: TurnInput) => Promise log?: (message: string) => void @@ -140,6 +163,7 @@ export class TurnDispatcher { timeoutMs: this.#deps.timeoutMs, ...(this.#deps.network !== undefined ? { network: this.#deps.network } : {}), allowedSchemes: this.#deps.allowedSchemes, + ...(this.#deps.memory !== undefined ? { memory: this.#deps.memory } : {}), } const promise = this.#deps diff --git a/packages/daemon/src/egress.ts b/packages/daemon/src/egress.ts index 21b289f..045d4ef 100644 --- a/packages/daemon/src/egress.ts +++ b/packages/daemon/src/egress.ts @@ -56,15 +56,41 @@ export function proxyRunArgs(input: { }): string[] { const args = ['run', '-d', '--name', input.name, '--network', input.egressNetwork, '--label', input.name] args.push(input.image) - for (const host of input.allowlist) args.push('--allow', host) + // An explicit empty allowlist means deny-all and must be passed through as such: zero `--allow` + // args would be indistinguishable from "no allowlist configured at all", which proxy.ts's CLI + // parser treats as "fall back to the default allowlist" — silently turning deny-all into + // allow-Anthropic. `--allow-none` makes deny-all explicit end to end (see proxy.ts). + if (input.allowlist.length === 0) { + args.push('--allow-none') + } else { + for (const host of input.allowlist) args.push('--allow', host) + } return args } +/** `docker network inspect -f '{{.Internal}}' ` trimmed to a boolean. Used to verify a + * pre-existing network really is internal before trusting it (see `ensureNetwork` below). */ +async function isNetworkInternal(exec: DockerExec, name: string): Promise { + const result = await exec(['network', 'inspect', '-f', '{{.Internal}}', name]) + return result.code === 0 && result.stdout.trim() === 'true' +} + async function ensureNetwork(exec: DockerExec, name: string, internal: boolean): Promise { const result = await exec(networkCreateArgs(name, internal)) if (result.code !== 0 && !/already exists/i.test(result.stderr)) { throw new Error(`docker network create ${name} failed (exit ${result.code}): ${result.stderr}`) } + if (result.code !== 0 && internal) { + // "Already exists" for the internal network: a pre-existing bridge network under this name + // would silently give turns unrestricted egress while the daemon still reports enforcement. + // Fail closed unless docker itself confirms the existing network really is internal. + if (!(await isNetworkInternal(exec, name))) { + throw new Error( + `docker network "${name}" already exists but is not internal; refusing to reuse it for ` + + `turn egress enforcement (a non-internal network would give turns unrestricted egress)`, + ) + } + } } /** @@ -83,8 +109,32 @@ export function createEgressManager(input: { proxyName: string allowlist: string[] exec?: DockerExec + /** Bounded readiness-poll knobs for the post-`docker run` check (see below). Small defaults; + * overridable so tests can run this fast without waiting on real container start latency. */ + readyAttempts?: number + readyDelayMs?: number }): EgressManager { const exec = input.exec ?? defaultExec + const readyAttempts = input.readyAttempts ?? 10 + const readyDelayMs = input.readyDelayMs ?? 200 + + /** Polls `docker inspect -f '{{.State.Running}}' ` until it reports `true`, up to a small + * bounded number of attempts. A successful `docker run -d` only means the container was created + * — it may still exit immediately (bad image, crash on boot) or never bind; without this check + * turns would dispatch onto a proxy that isn't actually there and burn retries with no egress at + * all. Throws (fail closed) if the proxy never reports ready. */ + async function waitUntilRunning(name: string): Promise { + for (let attempt = 1; attempt <= readyAttempts; attempt += 1) { + const result = await exec(['inspect', '-f', '{{.State.Running}}', name]) + if (result.code === 0 && result.stdout.trim() === 'true') return + if (attempt < readyAttempts) await new Promise((resolvePromise) => setTimeout(resolvePromise, readyDelayMs)) + } + throw new Error( + `radial-proxy container "${name}" did not report running after ${readyAttempts} check(s); ` + + `refusing to treat egress as enforced`, + ) + } + return { async start(): Promise { await ensureNetwork(exec, input.internalNetwork, true) @@ -111,6 +161,8 @@ export function createEgressManager(input: { ) } + await waitUntilRunning(input.proxyName) + return { network: input.internalNetwork, httpsProxy: `http://${input.proxyName}:8080` } }, diff --git a/packages/daemon/src/proxy.ts b/packages/daemon/src/proxy.ts index d19fb5f..85a0824 100644 --- a/packages/daemon/src/proxy.ts +++ b/packages/daemon/src/proxy.ts @@ -28,13 +28,17 @@ function targetFromConnect(url: string): { host: string; port: number } | undefi return host.length > 0 && Number.isFinite(port) ? { host, port } : undefined } +type Connect = typeof connect + class HttpEgressProxy implements EgressProxy { readonly #allowlist: string[] + readonly #connect: Connect #server: ReturnType | undefined #port = 0 - constructor(allowlist: string[]) { + constructor(allowlist: string[], connectFn: Connect = connect) { this.#allowlist = allowlist + this.#connect = connectFn } get port(): number { @@ -68,11 +72,15 @@ class HttpEgressProxy implements EgressProxy { #handleConnect(req: IncomingMessage, clientSocket: Socket, head: Uint8Array): void { const target = targetFromConnect(req.url ?? '') - if (!target || !isHostAllowed(target.host, this.#allowlist)) { + // 443-only by default (design §13: the allowlist is for HTTPS API egress, not arbitrary TCP). + // `targetFromConnect` parses both host and port from the CONNECT target; without this check an + // allowlisted host on a non-443 port (e.g. `api.anthropic.com:22`) would still be tunneled. + // Per-host port extension could be added later by widening this check if it's ever needed. + if (!target || target.port !== 443 || !isHostAllowed(target.host, this.#allowlist)) { clientSocket.end('HTTP/1.1 403 Forbidden\r\n\r\n') return } - const upstream = connect({ host: target.host, port: target.port }, () => { + const upstream = this.#connect({ host: target.host, port: target.port }, () => { clientSocket.write('HTTP/1.1 200 Connection Established\r\n\r\n') if (head.length > 0) upstream.write(head) // pipe() (rather than manual data listeners) forwards end-of-stream in both directions, so @@ -87,34 +95,50 @@ class HttpEgressProxy implements EgressProxy { /** * An allowlisting HTTP CONNECT proxy: the only egress path out of the run sandbox (design §13). - * Any CONNECT target whose host isn't in `allowlist` is rejected with 403 before a socket to it - * is ever opened; an allowed target gets a raw byte-for-byte tunnel. This is what - * `docker/proxy.Dockerfile` runs standalone in its own container (see the CLI entry below). + * Any CONNECT target whose host isn't in `allowlist`, or whose port isn't `443`, is rejected with + * 403 before a socket to it is ever opened; an allowed `host:443` target gets a raw byte-for-byte + * tunnel. This is what `docker/proxy.Dockerfile` runs standalone in its own container (see the CLI + * entry below). `connectFn` is an optional injection point (defaults to `node:net`'s `connect`, + * exactly as in production) purely so tests can exercise the real allowlist/port gate against a + * literal `:443` CONNECT target while redirecting the actual upstream dial to a local, unprivileged + * test server — binding a real listener on port 443 requires root on most systems. */ -export function createEgressProxy(allowlist: string[]): EgressProxy { - return new HttpEgressProxy(allowlist) +export function createEgressProxy(allowlist: string[], connectFn?: Connect): EgressProxy { + return new HttpEgressProxy(allowlist, connectFn) } -function parseCliArgs(args: string[]): { allow: string[]; port: number } { +/** + * Standalone-CLI arg parsing (exported for unit testing). Deny-all must survive here: an explicit + * `--allow-none` (emitted by `egress.ts`'s `proxyRunArgs` for an explicit empty allowlist) yields + * `allow: []`. Only when NEITHER any `--allow` NOR `--allow-none` is given — i.e. no allowlist + * decision was ever communicated to this process — does it fall back to the default allowlist. + * Without this distinction, deny-all (`[]`) and "unset" would both produce zero `--allow` flags and + * be indistinguishable, silently turning deny-all into allow-Anthropic. + */ +export function parseProxyArgs(args: string[]): { allow: string[]; port: number } { const allow: string[] = [] let port = 8080 + let sawAllowFlag = false for (let index = 0; index < args.length; index += 1) { const flag = args[index] const value = args[index + 1] if (flag === '--allow' && value !== undefined) { allow.push(value) + sawAllowFlag = true index += 1 + } else if (flag === '--allow-none') { + sawAllowFlag = true } else if (flag === '--port' && value !== undefined) { port = Number(value) index += 1 } } - return { allow: allow.length > 0 ? allow : ['api.anthropic.com'], port } + return { allow: sawAllowFlag ? allow : ['api.anthropic.com'], port } } const entry = argv[1] if (entry && import.meta.url === pathToFileURL(realpathSync(entry)).href) { - const { allow, port } = parseCliArgs(argv.slice(2)) + const { allow, port } = parseProxyArgs(argv.slice(2)) const proxy = createEgressProxy(allow) // Bind 0.0.0.0: this standalone entry is what `docker/proxy.Dockerfile` runs as the // radial-proxy container's entrypoint, and it must be reachable from the turn container diff --git a/packages/daemon/test/cli.test.mjs b/packages/daemon/test/cli.test.mjs index 8d2c35d..5927ba5 100644 --- a/packages/daemon/test/cli.test.mjs +++ b/packages/daemon/test/cli.test.mjs @@ -32,3 +32,32 @@ it('resolveRunConfig keeps an explicit run.network even with egress enforced, in assert.equal(resolved.egress, true) assert.equal(resolved.network, 'operator-custom-net') }) + +// FIX #12: an explicit `allowlist: []` (deny-all) is not the same as "unset", and `?? []`-style +// fallbacks must not conflate them — `[]` is non-nullish, so `run.allowlist ?? [default]` already +// preserves it correctly; this pins that down against a regression. +it('resolveRunConfig preserves an explicit empty allowlist (deny-all) rather than defaulting it to Anthropic', () => { + const resolved = resolveRunConfig({ + spaces: ['at://did:plc:human/com.disnetdev.radial.space/space1'], + allowlist: [], + }) + assert.deepEqual(resolved.allowlist, []) +}) + +it('resolveRunConfig defaults the allowlist to Anthropic only when truly unset', () => { + const resolved = resolveRunConfig({ spaces: ['at://did:plc:human/com.disnetdev.radial.space/space1'] }) + assert.deepEqual(resolved.allowlist, ['api.anthropic.com']) +}) + +// FIX #11: run.memory threads through resolveRunConfig with no default applied at this layer +// (turn.ts supplies the 4g floor when unset). +it('resolveRunConfig threads run.memory through unchanged, and leaves it unset when not configured', () => { + const withMemory = resolveRunConfig({ + spaces: ['at://did:plc:human/com.disnetdev.radial.space/space1'], + memory: '2g', + }) + assert.equal(withMemory.memory, '2g') + + const withoutMemory = resolveRunConfig({ spaces: ['at://did:plc:human/com.disnetdev.radial.space/space1'] }) + assert.equal('memory' in withoutMemory, false) +}) diff --git a/packages/daemon/test/config.test.mjs b/packages/daemon/test/config.test.mjs index efe3324..549d5ef 100644 --- a/packages/daemon/test/config.test.mjs +++ b/packages/daemon/test/config.test.mjs @@ -12,6 +12,7 @@ import { writeConfig, loadConfig, parseConfig, + parseRunConfig, passwordEnvName, resolveAgentInits, } from '../dist/index.js' @@ -165,6 +166,25 @@ it('scaffolds a config that round-trips straight into agent registration', async for (const profile of Object.keys(loaded.agents)) assertProfileName(profile) }) +// --- FIX #11 / FIX #12: run.memory and run.allowlist parsing -------------- + +it('parseRunConfig validates and threads through run.memory', () => { + const parsed = parseRunConfig({ spaces: ['at://did:plc:human/com.disnetdev.radial.space/space1'], memory: '2g' }) + assert.equal(parsed.memory, '2g') + assert.throws(() => parseRunConfig({ spaces: ['at://x'], memory: 4 }), /run\.memory must be a non-empty string/) +}) + +it('parseRunConfig leaves run.memory unset when not configured (no layer-injected default)', () => { + const parsed = parseRunConfig({ spaces: ['at://did:plc:human/com.disnetdev.radial.space/space1'] }) + assert.equal('memory' in parsed, false) +}) + +it('parseRunConfig preserves an explicit empty run.allowlist (deny-all) rather than dropping it', () => { + const parsed = parseRunConfig({ spaces: ['at://did:plc:human/com.disnetdev.radial.space/space1'], allowlist: [] }) + assert.deepEqual(parsed.allowlist, []) + assert.equal('allowlist' in parsed, true) +}) + it('refuses to clobber an existing config unless forced', async () => { const directory = await mkdtemp(join(tmpdir(), 'radial-clobber-')) const path = join(directory, 'radial.json') diff --git a/packages/daemon/test/dispatch.test.mjs b/packages/daemon/test/dispatch.test.mjs index 26e86c4..ec9d7ff 100644 --- a/packages/daemon/test/dispatch.test.mjs +++ b/packages/daemon/test/dispatch.test.mjs @@ -283,6 +283,70 @@ it('rejects a request the ledger considers ineligible: running, cooling down, an gaveUp.close() }) +// --- FIX #7: awaiting_input stale-index race -------------------------------- +// A ledger row in `awaiting_input` is immediately `eligible()`. If the next index snapshot has not +// yet ingested the assignee's question, `goalView.awaitingInput` is still empty too, so without the +// extra index-caught-up check the request would be re-selected and redispatched before any human +// ever sees the question. These three cases pin the intended matrix. + +it('awaiting_input + index has NOT ingested the assignee question yet: NOT selected (stale index)', () => { + const { records, request } = buildScenario() + const index = materialize(store(records), { spaceUri: SPACE_URI }) + const ledger = new TurnLedger() + ledger.markAwaitingInput(request.uri) + const actors = registryFor([actorFor(AGENT)]) + // Sanity: the ledger alone considers awaiting_input immediately eligible again. + assert.equal(ledger.eligible(request.uri), true) + assert.deepEqual(selectDispatchable(index, actors, ledger), []) + ledger.close() +}) + +it('awaiting_input + question ingested but unanswered (still in goalView.awaitingInput): NOT selected', () => { + const { records, request } = buildScenario() + const question = mk(AGENT, COLLECTIONS.message, 'question-1', 'cid-question-1', { + $type: COLLECTIONS.message, + goal: request.value.goal, + body: 'which approach do you want?', + mentions: [], + re: ref(request), + createdAt: '2026-01-01T00:03:10Z', + }) + const index = materialize(store([...records, question]), { spaceUri: SPACE_URI }) + assert.ok(index.goals[0].awaitingInput.includes(request.uri)) + const ledger = new TurnLedger() + ledger.markAwaitingInput(request.uri) + const actors = registryFor([actorFor(AGENT)]) + assert.deepEqual(selectDispatchable(index, actors, ledger), []) + ledger.close() +}) + +it('awaiting_input + question ingested AND answered (request absent from awaitingInput): selected', () => { + const { records, request } = buildScenario() + const question = mk(AGENT, COLLECTIONS.message, 'question-1', 'cid-question-1', { + $type: COLLECTIONS.message, + goal: request.value.goal, + body: 'which approach do you want?', + mentions: [], + re: ref(request), + createdAt: '2026-01-01T00:03:10Z', + }) + const answer = mk(HUMAN, COLLECTIONS.message, 'answer-1', 'cid-answer-1', { + $type: COLLECTIONS.message, + goal: request.value.goal, + body: 'go with approach A', + mentions: [], + parent: ref(question), + createdAt: '2026-01-01T00:03:20Z', + }) + const index = materialize(store([...records, question, answer]), { spaceUri: SPACE_URI }) + assert.ok(!index.goals[0].awaitingInput.includes(request.uri)) + const ledger = new TurnLedger() + ledger.markAwaitingInput(request.uri) + const actors = registryFor([actorFor(AGENT)]) + assert.deepEqual(uris(selectDispatchable(index, actors, ledger)), [request.uri]) + ledger.close() +}) + // --- Turn-loop integration via TurnDispatcher ----------------------------- function sendLine(socketPath, obj) { @@ -492,6 +556,34 @@ it('awaiting_input outcome is not relaunched, and once the question lands in the }) }) +// --- FIX #11: config-overridable memory limit ------------------------------- + +it('a configured DispatcherDeps.memory reaches the captured TurnInput.memory', async () => { + const { records, request } = buildScenario() + const index = materialize(store(records), { spaceUri: SPACE_URI }) + const actors = registryFor([actorFor(AGENT)]) + const ledger = new TurnLedger() + let captured + const dispatcher = new TurnDispatcher({ + ledger, + concurrency: 1, + runDirFor: (uri) => join('/tmp', 'radial-memory-test', uri.replace(/[^a-z0-9]/gi, '_')), + image: 'radial-turn:test', + timeoutMs: 30_000, + allowedSchemes: ['https'], + memory: '2g', + runTurn: async (input) => { + captured = input + return { outcome: 'fulfilled', acceptedRef: { uri: 'at://x/y/z', cid: 'bafy' }, label: 'radial.turn.x' } + }, + }) + dispatcher.pump(index, actors) + await dispatcher.drain() + assert.equal(captured.memory, '2g') + assert.equal(ledger.get(request.uri).state, 'fulfilled') + ledger.close() +}) + // --- Orphan reconciliation -------------------------------------------------- it('startup orphan reconciliation kills the tracked container and resets the ledger row to eligible', async () => { diff --git a/packages/daemon/test/egress.test.mjs b/packages/daemon/test/egress.test.mjs index 6c3c358..b106b1a 100644 --- a/packages/daemon/test/egress.test.mjs +++ b/packages/daemon/test/egress.test.mjs @@ -37,7 +37,7 @@ it('proxyRunArgs attaches the proxy to the egress network and passes the allowli assert.deepEqual(args.slice(imageIndex + 1), ['--allow', 'api.anthropic.com', '--allow', 'example.test']) }) -it('proxyRunArgs allows an empty allowlist (image falls back to its own default)', () => { +it('proxyRunArgs emits --allow-none for an explicit empty allowlist (deny-all), not zero --allow flags', () => { const args = proxyRunArgs({ name: 'radial-proxy', image: 'radial-proxy:latest', @@ -46,7 +46,21 @@ it('proxyRunArgs allows an empty allowlist (image falls back to its own default) allowlist: [], }) const imageIndex = args.indexOf('radial-proxy:latest') - assert.deepEqual(args.slice(imageIndex + 1), []) + assert.deepEqual(args.slice(imageIndex + 1), ['--allow-none']) + assert.ok(!args.includes('--allow')) +}) + +it('proxyRunArgs emits one --allow per host and no --allow-none for a non-empty allowlist', () => { + const args = proxyRunArgs({ + name: 'radial-proxy', + image: 'radial-proxy:latest', + internalNetwork: 'radial-internal', + egressNetwork: 'radial-egress', + allowlist: ['api.anthropic.com', 'example.test'], + }) + const imageIndex = args.indexOf('radial-proxy:latest') + assert.deepEqual(args.slice(imageIndex + 1), ['--allow', 'api.anthropic.com', '--allow', 'example.test']) + assert.ok(!args.includes('--allow-none')) }) function recordingExec(results) { @@ -59,8 +73,14 @@ function recordingExec(results) { return { exec, calls } } -it('createEgressManager.start() creates both networks, runs the proxy, joins the internal network, and returns the plan', async () => { - const { exec, calls } = recordingExec([]) +it('createEgressManager.start() creates both networks, runs the proxy, joins the internal network, verifies readiness, and returns the plan', async () => { + const { exec, calls } = recordingExec([ + { code: 0, stdout: '', stderr: '' }, // network create --internal + { code: 0, stdout: '', stderr: '' }, // network create (egress) + { code: 0, stdout: '', stderr: '' }, // run -d + { code: 0, stdout: '', stderr: '' }, // network connect + { code: 0, stdout: 'true\n', stderr: '' }, // docker inspect -f '{{.State.Running}}' radial-proxy + ]) const manager = createEgressManager({ internalNetwork: 'radial-internal', egressNetwork: 'radial-egress', @@ -78,11 +98,19 @@ it('createEgressManager.start() creates both networks, runs the proxy, joins the ['network', 'create', 'radial-egress'], ['run', '-d', '--name', 'radial-proxy', '--network', 'radial-egress', '--label', 'radial-proxy', 'radial-proxy:latest', '--allow', 'api.anthropic.com'], ['network', 'connect', 'radial-internal', 'radial-proxy'], + ['inspect', '-f', '{{.State.Running}}', 'radial-proxy'], ]) }) -it('createEgressManager.start() tolerates "network already exists" but rethrows other network-create failures', async () => { - const alreadyExists = recordingExec([{ code: 1, stdout: '', stderr: 'Error: network with name radial-internal already exists' }]) +it('createEgressManager.start() tolerates "network already exists" (verifying it really is internal) but rethrows other network-create failures', async () => { + const alreadyExists = recordingExec([ + { code: 1, stdout: '', stderr: 'Error: network with name radial-internal already exists' }, // network create --internal + { code: 0, stdout: 'true\n', stderr: '' }, // docker network inspect -f '{{.Internal}}' radial-internal + { code: 0, stdout: '', stderr: '' }, // network create (egress) + { code: 0, stdout: '', stderr: '' }, // run -d + { code: 0, stdout: '', stderr: '' }, // network connect + { code: 0, stdout: 'true\n', stderr: '' }, // readiness: docker inspect -f '{{.State.Running}}' + ]) const manager = createEgressManager({ internalNetwork: 'radial-internal', egressNetwork: 'radial-egress', @@ -92,7 +120,8 @@ it('createEgressManager.start() tolerates "network already exists" but rethrows exec: alreadyExists.exec, }) await manager.start() - assert.equal(alreadyExists.calls.length, 4) // did not stop short on the "already exists" result + assert.equal(alreadyExists.calls.length, 6) // did not stop short on the "already exists" result + assert.deepEqual(alreadyExists.calls[1], ['network', 'inspect', '-f', '{{.Internal}}', 'radial-internal']) const otherFailure = recordingExec([{ code: 1, stdout: '', stderr: 'permission denied' }]) const failing = createEgressManager({ @@ -106,6 +135,103 @@ it('createEgressManager.start() tolerates "network already exists" but rethrows await assert.rejects(() => failing.start(), /permission denied/) }) +// --- FIX #3: verify a pre-existing "already exists" internal network is actually internal ----- + +it('createEgressManager.start() fails closed when a pre-existing internal-network name is NOT actually internal', async () => { + const notInternal = recordingExec([ + { code: 1, stdout: '', stderr: 'Error: network with name radial-internal already exists' }, // network create --internal + { code: 0, stdout: 'false\n', stderr: '' }, // docker network inspect reports NOT internal + ]) + const manager = createEgressManager({ + internalNetwork: 'radial-internal', + egressNetwork: 'radial-egress', + proxyImage: 'radial-proxy:latest', + proxyName: 'radial-proxy', + allowlist: [], + exec: notInternal.exec, + }) + await assert.rejects(() => manager.start(), /not internal/) + // Fails closed before ever running the proxy container. + assert.deepEqual( + notInternal.calls, + [ + ['network', 'create', '--internal', 'radial-internal'], + ['network', 'inspect', '-f', '{{.Internal}}', 'radial-internal'], + ], + ) +}) + +it('createEgressManager.start() proceeds when a pre-existing internal-network name really is internal', async () => { + const isInternal = recordingExec([ + { code: 1, stdout: '', stderr: 'Error: network with name radial-internal already exists' }, + { code: 0, stdout: 'true\n', stderr: '' }, + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: 'true\n', stderr: '' }, + ]) + const manager = createEgressManager({ + internalNetwork: 'radial-internal', + egressNetwork: 'radial-egress', + proxyImage: 'radial-proxy:latest', + proxyName: 'radial-proxy', + allowlist: [], + exec: isInternal.exec, + }) + const plan = await manager.start() + assert.deepEqual(plan, { network: 'radial-internal', httpsProxy: 'http://radial-proxy:8080' }) +}) + +// --- FIX #14: verify proxy readiness before returning the EgressPlan, else fail closed --------- + +it('createEgressManager.start() resolves once docker inspect reports the proxy container Running', async () => { + const { exec } = recordingExec([ + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: 'true\n', stderr: '' }, // readiness check succeeds on the first attempt + ]) + const manager = createEgressManager({ + internalNetwork: 'radial-internal', + egressNetwork: 'radial-egress', + proxyImage: 'radial-proxy:latest', + proxyName: 'radial-proxy', + allowlist: [], + exec, + readyAttempts: 3, + readyDelayMs: 1, + }) + const plan = await manager.start() + assert.deepEqual(plan, { network: 'radial-internal', httpsProxy: 'http://radial-proxy:8080' }) +}) + +it('createEgressManager.start() rejects (fails closed) when the proxy never reports Running', async () => { + const { exec, calls } = recordingExec([ + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: '', stderr: '' }, + { code: 0, stdout: 'false\n', stderr: '' }, + { code: 0, stdout: 'false\n', stderr: '' }, + { code: 1, stdout: '', stderr: 'no such container' }, + ]) + const manager = createEgressManager({ + internalNetwork: 'radial-internal', + egressNetwork: 'radial-egress', + proxyImage: 'radial-proxy:latest', + proxyName: 'radial-proxy', + allowlist: [], + exec, + readyAttempts: 3, + readyDelayMs: 1, // keep the test fast; bounded attempts are what matter, not real wall time + }) + await assert.rejects(() => manager.start(), /did not report running/) + // Exactly the bounded number of readiness attempts were made — not an unbounded retry loop. + const inspectCalls = calls.filter((c) => c[0] === 'inspect') + assert.equal(inspectCalls.length, 3) +}) + it('createEgressManager.start() surfaces a failed proxy run or a failed network connect', async () => { const runFails = recordingExec([ { code: 0, stdout: '', stderr: '' }, diff --git a/packages/daemon/test/proxy.test.mjs b/packages/daemon/test/proxy.test.mjs index cc4d4ab..62136b2 100644 --- a/packages/daemon/test/proxy.test.mjs +++ b/packages/daemon/test/proxy.test.mjs @@ -1,7 +1,7 @@ import assert from 'node:assert/strict' import { connect, createServer } from 'node:net' import { it } from 'node:test' -import { createEgressProxy, isHostAllowed } from '../dist/index.js' +import { createEgressProxy, isHostAllowed, parseProxyArgs } from '../dist/index.js' it('isHostAllowed matches exactly, case-insensitively, and strips a trailing port', () => { assert.equal(isHostAllowed('api.anthropic.com', ['api.anthropic.com']), true) @@ -12,6 +12,24 @@ it('isHostAllowed matches exactly, case-insensitively, and strips a trailing por assert.equal(isHostAllowed('', ['api.anthropic.com']), false) }) +// --- FIX #12: deny-all (--allow-none) must survive CLI arg parsing distinctly from "unset" ----- + +it('parseProxyArgs(["--allow-none"]) is explicit deny-all: allow: []', () => { + assert.deepEqual(parseProxyArgs(['--allow-none']), { allow: [], port: 8080 }) +}) + +it('parseProxyArgs([]) with no --allow/--allow-none at all falls back to the default allowlist', () => { + assert.deepEqual(parseProxyArgs([]), { allow: ['api.anthropic.com'], port: 8080 }) +}) + +it('parseProxyArgs(["--allow", "x.test"]) uses exactly the given allowlist', () => { + assert.deepEqual(parseProxyArgs(['--allow', 'x.test']), { allow: ['x.test'], port: 8080 }) +}) + +it('parseProxyArgs honors --port alongside --allow', () => { + assert.deepEqual(parseProxyArgs(['--allow', 'x.test', '--port', '9090']), { allow: ['x.test'], port: 9090 }) +}) + function connectRaw(port) { return new Promise((resolvePromise, reject) => { const socket = connect({ host: '127.0.0.1', port }) @@ -41,20 +59,29 @@ it('denies a CONNECT to a host that is not on the allowlist', async () => { } }) -it('tunnels an allowlisted CONNECT target byte-for-byte', async () => { +// `createEgressProxy`'s optional third argument lets a test override the upstream `connect` used +// once a CONNECT target clears the allowlist/port gate. This is purely a testing seam (production +// always uses the real `node:net#connect`, unmodified) that redirects the *actual* TCP dial to a +// local, unprivileged test target while still exercising the real gate logic against a literal +// `:443` CONNECT target — binding a real listener on port 443 itself requires root on most systems. +function connectToRealTargetInsteadOf443(realPort) { + return (_options, listener) => connect({ host: '127.0.0.1', port: realPort }, listener) +} + +it('tunnels an allowlisted :443 CONNECT target byte-for-byte', async () => { const target = createServer((socket) => { socket.on('data', (chunk) => socket.write(chunk)) }) await new Promise((resolvePromise) => target.listen(0, '127.0.0.1', resolvePromise)) const targetPort = target.address().port - const proxy = createEgressProxy(['127.0.0.1']) + const proxy = createEgressProxy(['127.0.0.1'], connectToRealTargetInsteadOf443(targetPort)) const proxyPort = await proxy.listen(0) try { const socket = await connectRaw(proxyPort) const connectResponsePromise = readOnce(socket) - socket.write(`CONNECT 127.0.0.1:${targetPort} HTTP/1.1\r\nHost: 127.0.0.1:${targetPort}\r\n\r\n`) + socket.write('CONNECT 127.0.0.1:443 HTTP/1.1\r\nHost: 127.0.0.1:443\r\n\r\n') const connectResponse = await connectResponsePromise assert.match(connectResponse, /200/) @@ -70,3 +97,43 @@ it('tunnels an allowlisted CONNECT target byte-for-byte', async () => { await new Promise((resolvePromise) => target.close(resolvePromise)) } }) + +// --- FIX #9: CONNECT is restricted to port 443, even for an allowlisted host -------------------- + +it('denies a CONNECT to an allowlisted host on a non-443 port', async () => { + const proxy = createEgressProxy(['allowed.test']) + const port = await proxy.listen(0) + try { + const socket = await connectRaw(port) + const responsePromise = readOnce(socket) + socket.write('CONNECT allowed.test:22 HTTP/1.1\r\nHost: allowed.test:22\r\n\r\n') + const response = await responsePromise + assert.match(response, /403/) + await new Promise((resolvePromise) => socket.on('close', resolvePromise)) + } finally { + await proxy.close() + } +}) + +it('tunnels the same allowlisted host on port 443 (denied on 22 above, allowed here)', async () => { + const target = createServer((socket) => { + socket.on('data', (chunk) => socket.write(chunk)) + }) + await new Promise((resolvePromise) => target.listen(0, '127.0.0.1', resolvePromise)) + const targetPort = target.address().port + + const proxy = createEgressProxy(['allowed.test'], connectToRealTargetInsteadOf443(targetPort)) + const proxyPort = await proxy.listen(0) + try { + const socket = await connectRaw(proxyPort) + const connectResponsePromise = readOnce(socket) + socket.write('CONNECT allowed.test:443 HTTP/1.1\r\nHost: allowed.test:443\r\n\r\n') + const connectResponse = await connectResponsePromise + assert.match(connectResponse, /200/) + socket.end() + await new Promise((resolvePromise) => socket.on('close', resolvePromise)) + } finally { + await proxy.close() + await new Promise((resolvePromise) => target.close(resolvePromise)) + } +}) diff --git a/packages/daemon/test/reconcile.test.mjs b/packages/daemon/test/reconcile.test.mjs new file mode 100644 index 0000000..dd3db00 --- /dev/null +++ b/packages/daemon/test/reconcile.test.mjs @@ -0,0 +1,63 @@ +import assert from 'node:assert/strict' +import { it } from 'node:test' +import { FakeContainerRunner, TurnLedger } from '../dist/index.js' +import { reconcileOrphans } from '../dist/cli.js' + +// FIX #6: orphan reconciliation must preserve attempts (markCrashed, not reset). A daemon that +// repeatedly crashes mid-turn must still respect the retry bound across restarts, instead of +// silently resetting the retry counter (and therefore never giving up) every time it comes back up. + +const REQUEST_URI = 'at://did:plc:human/com.disnetdev.radial.artifactRequest/request-1' +const REQUEST_CID = 'cid-request-1' +const LABEL = 'radial.turn.orphan' + +it('reconcileOrphans kills the tracked container and transitions the row to crashed with attempts preserved (not deleted)', async () => { + const ledger = new TurnLedger(':memory:', { retryBound: 5, cooldownMs: 60_000 }) + // Seed one prior failed attempt, then a fresh "running" row (as if the daemon restarted right + // after redispatching), so attempts is already > 0 going into this reconciliation. + ledger.markRunning(REQUEST_URI, REQUEST_CID, { containerLabel: LABEL, checkoutPath: '/tmp/x' }) + ledger.markCrashed(REQUEST_URI) + assert.equal(ledger.get(REQUEST_URI).attempts, 1) + ledger.markRunning(REQUEST_URI, REQUEST_CID, { containerLabel: LABEL, checkoutPath: '/tmp/x' }) + assert.equal(ledger.get(REQUEST_URI).state, 'running') + + const runner = new FakeContainerRunner(async () => ({ exitCode: 0, timedOut: false })) + // Seed a started "container" under the label the way a real orphan would be tracked. + await runner.run({ label: LABEL, image: 'x', argv: [], env: {}, mounts: [], timeoutMs: 1000 }) + const killed = [] + const originalKill = runner.kill.bind(runner) + runner.kill = async (idOrName) => { + killed.push(idOrName) + return originalKill(idOrName) + } + + await reconcileOrphans(runner, ledger) + + assert.ok(killed.length > 0) + assert.deepEqual(await runner.listByLabel(LABEL), []) + + const row = ledger.get(REQUEST_URI) + assert.ok(row, 'the row must still exist — reconciliation must not delete it') + assert.ok(row.state === 'crashed' || row.state === 'gave_up', `expected crashed/gave_up, got ${row.state}`) + assert.equal(row.attempts, 2) // the pre-seeded attempt plus this reconciliation's attempt + assert.deepEqual(ledger.running(), []) // no longer "running": the orphan is fully reconciled + + ledger.close() +}) + +it('reconcileOrphans applied repeatedly across "restarts" eventually gives up instead of resetting the counter', async () => { + const retryBound = 3 + const ledger = new TurnLedger(':memory:', { retryBound, cooldownMs: 0 }) + const runner = new FakeContainerRunner(async () => ({ exitCode: 0, timedOut: false })) + + for (let restart = 0; restart < retryBound; restart += 1) { + ledger.markRunning(REQUEST_URI, REQUEST_CID, { containerLabel: LABEL, checkoutPath: '/tmp/x' }) + await runner.run({ label: LABEL, image: 'x', argv: [], env: {}, mounts: [], timeoutMs: 1000 }) + await reconcileOrphans(runner, ledger) + } + + const row = ledger.get(REQUEST_URI) + assert.equal(row.state, 'gave_up') + assert.equal(row.attempts, retryBound) + ledger.close() +}) -- 2.51.2