From c9396b3e23f2ae081b197b25542727a338fe95b2 Mon Sep 17 00:00:00 2001 From: Tim Disney Date: Tue, 28 Jul 2026 11:17:36 -0700 Subject: [PATCH] make several radiald instances on one machine explicit (#20) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two agent identities under one daemon already worked — profiles carry their own `identifier`. What was singular was the LOCATION: `new FileSessionStore()` with no argument, movable only by `$RADIAL_DATA_DIR`, with `sessions.json` keyed by profile name. Two instances sharing a data directory and reusing a profile name clobbered each other silently; the loser then ran under the wrong DID with no complaint. - `--data-dir` beside `--config` on every command, and a config-relative `"dataDir"` in radial.json, so `--config ` alone selects an instance. New `instance.ts` owns the precedence (flag → env → config → default) and reports when the environment shadowed the file. - `run.stateDir` defaults to `/run`, so one setting separates the ledgers, index and per-turn directories too. - `radiald init` refuses to re-register a profile under a different DID (`--replace-identity` overrides); `loadActors` refuses a configured DID that does not match its stored session, and warns on a drifted handle or a profile with no session. - `radiald run` prints the config, data dir, sessions file, state dir and instance id it resolved, plus the DID and handle per profile. - New `StateDirLock`: one daemon per state directory, enforced by the same cross-process SQLite lock the session store uses. `index reset` takes it too. - Container labels carry an instance segment, since docker's namespace is machine-global — no `--name` collisions, no reconciling away another instance's live containers. - `radial --data-dir` for the human CLI; docs and CHANGELOG. New tests: instance.test.mjs, state-lock.test.mjs, and two-instances.test.mjs — the offline drill with two identities and the same profile names on one machine. Co-authored-by: claudebot.disnetdev.com (did:plc:n6ku5xddiuguwze3f356evla) --- CHANGELOG.md | 44 + docs/operators.md | 43 + docs/radial-json.md | 81 +- docs/running-an-agent.md | 7 + packages/atproto/src/node.ts | 6 +- packages/daemon/src/actors.ts | 57 +- packages/daemon/src/check-dispatch.ts | 6 +- packages/daemon/src/check-runner.ts | 22 +- packages/daemon/src/cli.ts | 796 +++++++++++-------- packages/daemon/src/config.ts | 18 +- packages/daemon/src/dispatch.ts | 6 +- packages/daemon/src/index.ts | 2 + packages/daemon/src/init.ts | 20 + packages/daemon/src/instance.ts | 153 ++++ packages/daemon/src/state-lock.ts | 60 ++ packages/daemon/src/turn.ts | 27 +- packages/daemon/test/actors.test.mjs | 103 ++- packages/daemon/test/check-dispatch.test.mjs | 27 +- packages/daemon/test/check-runner.test.mjs | 9 + packages/daemon/test/cli.test.mjs | 131 ++- packages/daemon/test/config.test.mjs | 26 + packages/daemon/test/dispatch.test.mjs | 32 + packages/daemon/test/init.test.mjs | 84 ++ packages/daemon/test/instance.test.mjs | 105 +++ packages/daemon/test/state-lock.test.mjs | 64 ++ packages/daemon/test/turn.test.mjs | 17 + packages/daemon/test/two-instances.test.mjs | 162 ++++ packages/sidecar/src/cli.ts | 10 +- 28 files changed, 1750 insertions(+), 368 deletions(-) create mode 100644 packages/daemon/src/instance.ts create mode 100644 packages/daemon/src/state-lock.ts create mode 100644 packages/daemon/test/instance.test.mjs create mode 100644 packages/daemon/test/state-lock.test.mjs create mode 100644 packages/daemon/test/two-instances.test.mjs diff --git a/CHANGELOG.md b/CHANGELOG.md index 6e6aa00..5b2476d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -91,8 +91,52 @@ including multi-operator hardening. `clocks` in `radiald index --digest`. Advisory only — nothing in the fold compares two clocks. +- **Several `radiald` instances on one machine, spelled rather than remembered.** + `--data-dir ` sits beside `--config` on every command that takes one, and + a top-level `"dataDir"` in `radial.json` — resolved relative to the config file, + so the same command from two directories is the same instance — makes `--config` + select an instance outright: sessions, run state and all. `radiald config init + --data-dir` writes it into the scaffold and prints follow-up commands that carry + `--config`. Precedence is `--data-dir` → `$RADIAL_DATA_DIR` → `"dataDir"` → + the platform default, and the environment shadowing a different `"dataDir"` is + warned about by name. `radiald run` prints the config, data dir, sessions file, + state dir and instance id it resolved, plus the DID *and* handle every profile is + running as. Two identities under **one** daemon never needed any of this and are + unchanged; see `docs/radial-json.md`, "More than one daemon on one machine". +- **A run state directory holds a lock.** `radiald run` takes a cross-process + SQLite lock on `/run.lock.db` for its lifetime, so a second daemon on + one state directory is refused at startup instead of becoming a second writer on + the record index and three ledgers. `radiald index reset` takes the same lock, + which replaces the "stop radiald run first" line in `--help` with an actual check. + The OS releases it when the holder exits, so a crashed daemon leaves nothing stale. +- **Container names are scoped to the instance.** Turn and check labels are now + `radial.turn..` / `radial.check..`, where + `` is derived from the run state directory. Docker's namespace is + machine-global, so two daemons serving one space could previously collide on + `docker run --name` and kill each other's live containers during startup orphan + reconciliation. Reconciliation is unaffected across the upgrade: it kills the + label stored in its own ledger row. + ### Changed +- **`run.stateDir` defaults under the resolved data directory** (`/run`) + rather than the process-wide one. Adding `"dataDir"` to an existing config + therefore moves the state directory unless `run.stateDir` is set explicitly: the + new one starts empty, which means a full re-ingest and fresh ledgers (attempts and + cooldowns reset). Nothing is lost protocol-side — a fulfilled request is closed by + the fold, not by the ledger — and the resolved path is printed at startup. Move the + old directory, or set `run.stateDir`, to keep it. +- **`radiald init` refuses to re-register a profile under a different DID.** + `sessions.json` is keyed by profile name, so a second instance sharing a data + directory and reusing a profile name used to replace the first one's identity in + silence — after which the other daemon ran as somebody else with no error. The + refusal names both DIDs and the file; `--replace-identity` is the override for + deliberately re-pointing a profile. Re-authenticating the same DID is unchanged. +- **`radiald run` fails when a configured DID does not match its stored session**, + rather than running under whichever identity is in the sessions file. A configured + *handle* that has drifted only warns — handles can legitimately be renamed after + init. A configured profile with no session is still skipped, but now says so and + names the file it looked in. - **The web app asks for granular atproto OAuth scopes instead of `transition:generic`.** The old coarse scope let a signed-in tab write *anything* into the human's repo — their posts and follows included. The diff --git a/docs/operators.md b/docs/operators.md index 004a55a..d6650e0 100644 --- a/docs/operators.md +++ b/docs/operators.md @@ -58,6 +58,49 @@ No shared database, no shared queue, no coordinator, no agreement about intervals or images, no network path between the two daemons. Neither operator can see the other's ledgers, and neither needs to. +### 1a. The second daemon on *your* machine + +§1 is about the second operator. The same-machine case — two daemons, two agent +identities, one host — is a different question, and the answer is not "run two +daemons": several identities under **one** daemon already work, because each +profile may carry its own `identifier`, and one process serves them all. + +Run a second daemon when the two need to differ in something a profile cannot +express: a separate lifecycle (restart one without the other), a different GitHub +account, a different provider key, a different poll interval. + +Then give each instance its own config naming its own data directory, and let +`--config` select it: + +```sh +radiald config init agent-b.example --config ~/radial/b.json --data-dir ~/radial/b-state +printf '%s\n' "$B_PASSWORD" | radiald init --config ~/radial/b.json --password-stdin +radiald run --config ~/radial/b.json +``` + +Under systemd, each unit carries its own environment — this is the part no config +file can hold, because it is credentials: + +```ini +[Service] +ExecStart=/usr/bin/radiald run --config /home/op/radial/b.json +Environment=GH_TOKEN=… +Environment=ANTHROPIC_API_KEY=… +``` + +`gh auth token` is per *user*, not per instance, so two daemons acting as two +GitHub accounts each need `GH_TOKEN` in their own process — the daemon prefers it +over `gh` when it is set. Profile names need not differ between instances; data +directories must. `radiald init` refuses to re-register a profile under a different +DID (naming both, and `--replace-identity` if you meant it), `radiald run` refuses +to start when a configured DID does not match its stored session, and only one +daemon may hold a run state directory at a time. + +The field reference for all of this is +[radial-json.md, "More than one daemon on one machine"](radial-json.md#more-than-one-daemon-on-one-machine). +Note that two daemons polling one space duplicate the ingestion traffic — inherent +to running two, not a cost of separating their state. + --- ## 2. Routing: which of your agents answers what diff --git a/docs/radial-json.md b/docs/radial-json.md index 909758a..5469d90 100644 --- a/docs/radial-json.md +++ b/docs/radial-json.md @@ -19,9 +19,25 @@ In precedence order: else `~/.config/radial/radial.json` A `radial.json` beside a checkout is useful for testing or a second set of -identities. Sessions written by `radiald init` live separately, under -`$RADIAL_DATA_DIR` (else the platform data directory) — the config never holds a -credential. +identities. Sessions written by `radiald init` live separately, in the **data +directory** — the config never holds a credential. + +## The data directory + +Everything one daemon owns locally: `sessions.json`, and (unless `run.stateDir` +says otherwise) the run state directory holding the record index, checkpoints and +the turn/check/claim ledgers. In precedence order: + +1. `--data-dir `, resolved against the working directory +2. `$RADIAL_DATA_DIR` +3. `"dataDir"` in this file, resolved against **the directory of this file** +4. `$XDG_STATE_HOME/radial`, else `~/.local/state/radial` + +Rung 3 is config-relative on purpose: a working-directory-relative `dataDir` +would mean the same command run from two directories runs under two identities, +which is exactly the mistake this setting exists to prevent. `radiald run` prints +the directory it resolved and which rung chose it, and warns when +`$RADIAL_DATA_DIR` shadowed a different `dataDir` in the file. ## Shape @@ -59,7 +75,8 @@ credential. | Field | Meaning | | --- | --- | -| `identifier` | The agent account every profile uses unless it overrides `identifier` itself. A handle or a DID. | +| `identifier` | The agent account every profile uses unless it overrides `identifier` itself. A handle or a DID. Written as a DID it is *checked*: `radiald run` refuses to start when the stored session is a different identity, rather than running as whoever is in `sessions.json`. A handle only warns, since a handle can be renamed after `radiald init`. | +| `dataDir` | This instance's data directory (sessions, and by default the run state dir), resolved relative to this file. Naming it here is what lets `--config ` alone select one of several daemons on a machine — see "More than one daemon on one machine". | | `pds` | Override for PDS resolution. Normally omitted: the PDS is resolved from the identity — a handle through its `_atproto` DNS TXT record (falling back to `https:///.well-known/atproto-did`), then the DID document's `AtprotoPersonalDataServer` endpoint. Set it against a local PDS in testing. | | `harness` | Default agent harness: `claude` (the [Claude Code](https://claude.com/claude-code) CLI) or `pi` (the [pi coding agent](https://pi.dev)). Inherited by every profile that does not name its own. Validated at `radiald init` and again at `radiald run` startup — an unknown name is refused, never silently downgraded to another harness. | | `models` | Default model preference order, `"name=costHint"` or `{ "name", "costHint" }`. The turn hands these to the harness, which takes the first; an empty list lets the harness choose. The `name` reaches the harness's CLI verbatim — see "Choosing a harness and a provider" for each one's syntax. | @@ -90,7 +107,7 @@ into the file. | Field | Default | Meaning | | --- | --- | --- | | `spaces` | — | The `at://` URIs of the spaces to poll. | -| `stateDir` | under the Radial data directory | Records, checkpoints, turn and check ledgers. | +| `stateDir` | `/run` | Records, checkpoints, turn/check/claim ledgers, per-turn directories and turn sockets. A **relative** path here is resolved against the working directory, not the config — `radiald run` warns when it sees one. Only one daemon may hold a state directory at a time; a second is refused at startup. | | `image` / `checkImage` | `radial-turn:latest` / falls back to `image` | Turn and check container images. | | `concurrency` / `checkConcurrency` | 1 / 1 | Turns and checks run on separate budgets, so a slow check never blocks a turn. | | `timeoutMs` / `checkTimeoutMs` | 60 min / 15 min | Whole-run wall clock caps, the check one covering its host-side checkout too. Outliving one is a retryable crash. | @@ -142,6 +159,60 @@ A stream event wakes the run loop immediately rather than waiting out the interval. Record **deletes** are counted and never applied — Radial has no record tombstones (see `docs/operators.md`). +## More than one daemon on one machine + +Two agent identities under one daemon needs none of this: profiles already carry +their own `identifier`, and one `radiald run` serves them all. What this section +is for is two *daemons* — separate processes, separate lifecycles, separate +GitHub accounts or provider keys. + +Give each one a config that names its own data directory, and `--config` selects +it outright: + +```sh +radiald config init agent-a.example --config ~/radial/a.json --data-dir ~/radial/a-state +radiald config init agent-b.example --config ~/radial/b.json --data-dir ~/radial/b-state + +printf '%s\n' "$A_PASSWORD" | radiald init --config ~/radial/a.json --password-stdin +printf '%s\n' "$B_PASSWORD" | radiald init --config ~/radial/b.json --password-stdin + +radiald run --config ~/radial/a.json & +radiald run --config ~/radial/b.json & +``` + +No environment variable appears anywhere above, which is the point: `$RADIAL_CONFIG` +and `$RADIAL_DATA_DIR` still work, and still take precedence, but forgetting one on +a single invocation is how an instance ends up reading the wrong `sessions.json`. + +**Profile names do not have to differ.** Both instances may call a profile +`planner`: a profile name is a record key in the agent's own repo, and the two +repos are different. What must differ is the data directory — profiles are keyed +by name *within* a `sessions.json`, so two instances sharing one file would +overwrite each other's identity. `radiald init` refuses to do that (it names both +DIDs and the file; `--replace-identity` is the override), and `radiald run` refuses +to start when a configured DID does not match its stored session. + +Per instance: the config, the data directory, `sessions.json`, the run state +directory and everything in it (record index, checkpoints, the three ledgers, +per-turn directories, tangled push keys), the container name prefix, and anything +taken from the process environment — `GH_TOKEN`, provider API keys, `--interval`. +Under systemd that is one `Environment=` block per unit. + +Shared: the docker daemon and its images, and the `gh` CLI's own configuration. +`gh auth token` is per user, not per instance, so two daemons that must act as two +GitHub accounts each need their own `GH_TOKEN` exported in their own process. + +Two things worth knowing before you split: + +- **Two daemons polling one space duplicate the ingestion traffic.** That is + inherent to running two of them, not a cost of separating their state. +- **Adding `dataDir` to an existing config moves the state directory** (unless + `run.stateDir` is set explicitly). The new one starts empty: a full re-ingest, and + fresh ledgers, so attempts and cooldowns reset. Nothing is lost protocol-side — + a fulfilled request is closed by the fold, not by the ledger — but if you want the + old state, move it or set `run.stateDir` to where it already is. The startup banner + prints the resolved state directory, so the change is visible on the first run. + ## Choosing a harness and a provider Three harnesses ship today, and `harness` is resolved **per profile**, so one daemon can run all of diff --git a/docs/running-an-agent.md b/docs/running-an-agent.md index 2298049..1b2699b 100644 --- a/docs/running-an-agent.md +++ b/docs/running-an-agent.md @@ -66,6 +66,13 @@ override the directory, alongside the `$RADIAL_DATA_DIR` holding sessions). A useful for testing or a second set of identities; `--config ` or `$RADIAL_CONFIG` overrides both. +`--data-dir ` sits beside `--config` on every command that takes one, and +chooses where this daemon's `sessions.json` and run state live. Setting `"dataDir"` +in the config itself (resolved relative to the config file) is usually better: +`--config` then selects an instance outright, which is what running **two** daemons +on one machine needs — see +[radial-json.md, "More than one daemon on one machine"](radial-json.md#more-than-one-daemon-on-one-machine). + Profiles share one agent identity by default and inherit the top-level `pds`, `harness`, and `models`; each publishes its own agent record keyed by profile name, so one identity can offer several distinct capability sets: diff --git a/packages/atproto/src/node.ts b/packages/atproto/src/node.ts index d45307e..7f673c4 100644 --- a/packages/atproto/src/node.ts +++ b/packages/atproto/src/node.ts @@ -137,8 +137,10 @@ export class FileSessionStore { } } -export function defaultDataDirectory(): string { - const env = process.env +/** The operator's Radial state directory. `env` is a parameter rather than a read of `process.env` + * so a caller resolving a data directory for one *instance* (daemon `instance.ts`) can decide the + * precedence itself and still land on the same default. */ +export function defaultDataDirectory(env: Record = process.env): string { if (env.RADIAL_DATA_DIR) return env.RADIAL_DATA_DIR if (env.XDG_STATE_HOME) return join(env.XDG_STATE_HOME, 'radial') if (!env.HOME) throw new Error('Set RADIAL_DATA_DIR or HOME') diff --git a/packages/daemon/src/actors.ts b/packages/daemon/src/actors.ts index 9cbdf5a..cbe8fda 100644 --- a/packages/daemon/src/actors.ts +++ b/packages/daemon/src/actors.ts @@ -27,21 +27,58 @@ export interface ActorRegistry { /** * Loads only profiles that have already been `radiald init`'d (a stored session exists); a - * configured-but-uninitialized profile is silently skipped rather than treated as an error, so an - * operator can add profiles to radial.json ahead of running init for them. + * configured-but-uninitialized profile is skipped rather than treated as an error, so an operator + * can add profiles to radial.json ahead of running init for them — but the skip is now reported + * through `warn`, naming the sessions file it looked in. + * + * It also cross-checks the identity: a config whose `identifier` is a DID that differs from the + * stored session's is the failure that made several instances on one machine dangerous — the daemon + * would run under whichever identity happened to be in `sessions.json` and say nothing — so that one + * is fatal. A HANDLE that differs is only a warning: handles can legitimately be renamed after init. */ export async function loadActors( config: RadialConfig, store: FileSessionStore, - options: { fetcher?: FetchLike } = {}, + options: { + fetcher?: FetchLike + /** Operator-facing advisories. Absent means the caller does not want them (tests, tooling). */ + warn?: (message: string) => void + /** Named in those advisories, so an operator knows which file to edit. */ + configPath?: string + } = {}, ): Promise { const profiles = Object.keys(config.agents).sort() const all: LoadedActor[] = [] + const uninitialized: string[] = [] + const sessionsPath = store.path for (const profile of profiles) { const session = await store.get(profile) - if (!session) continue const agent = config.agents[profile] + if (!session) { + if (agent) uninitialized.push(profile) + continue + } if (!agent) continue + const configured = agent.identifier ?? config.identifier + if (configured?.startsWith('did:') && configured !== session.did) { + throw new Error( + `Profile "${profile}" is configured as ${configured} in ${options.configPath ?? 'the config'} ` + + `but ${sessionsPath} holds a session for ${session.did}. This daemon would run under the wrong ` + + 'identity: point it at its own data directory (--data-dir, or "dataDir" in radial.json), or ' + + 're-run "radiald init" for it.', + ) + } + if ( + configured && + !configured.startsWith('did:') && + configured.toLowerCase() !== session.handle.toLowerCase() + ) { + options.warn?.( + `profile "${profile}" is configured as ${configured} but its stored session is ${session.handle} ` + + `(${session.did}) in ${sessionsPath}. A handle can be renamed after init, so this is not fatal — ` + + 'write the DID in "identifier" (or re-run "radiald init") to make it checkable.', + ) + } // Same precedence `radiald init` used when it published the agent record. A profile with a // stored session and no harness anywhere cannot be run at all, and skipping it silently is // exactly the class of bug this resolution exists to prevent — so it is a named hard error. @@ -63,6 +100,18 @@ export async function loadActors( client, }) } + // Said once, and only when SOME profile did load: a run where none did is an uninitialized (or + // misdirected) instance, and its caller has a better error for that than a list of every profile. + if (uninitialized.length > 0 && all.length > 0) { + const where = options.configPath ? ` in ${options.configPath}` : '' + options.warn?.( + `${uninitialized.length === 1 ? 'profile' : 'profiles'} ${uninitialized.join(', ')} ` + + `${uninitialized.length === 1 ? 'is' : 'are'} configured${where} but ${uninitialized.length === 1 ? 'has' : 'have'} ` + + `no session in ${sessionsPath}, so ${uninitialized.length === 1 ? 'it' : 'they'} will not run. ` + + `Initialize ${uninitialized.length === 1 ? 'it' : 'them'}: radiald init ${uninitialized.join(' ')}` + + (options.configPath ? ` --config ${options.configPath}` : ''), + ) + } const byDid = new Map() for (const actor of all) { const existing = byDid.get(actor.did) diff --git a/packages/daemon/src/check-dispatch.ts b/packages/daemon/src/check-dispatch.ts index 04296ae..94d83cd 100644 --- a/packages/daemon/src/check-dispatch.ts +++ b/packages/daemon/src/check-dispatch.ts @@ -137,6 +137,9 @@ export interface CheckDispatcherDeps { allowedSchemes: string[] /** Container memory limit (docker `--memory`); unset lets check-runner apply its `4g` floor. */ memory?: string + /** This daemon instance's id (`instanceId(stateDir)`), mixed into every container label — the + * check-side half of the same machine-global docker namespace problem turns have. */ + instance?: string /** Injected wrapper around check-runner's runCheckRun (binds runner/checkout/now). */ runCheckRun: (input: CheckRunInput) => Promise log?: (message: string) => void @@ -191,7 +194,7 @@ export class CheckDispatcher { const keyStr = checkKeyString(key) if (this.#inFlight.has(keyStr)) continue - const label = checkContainerLabel(key.artifactUri, key.artifactCid, key.commit) + const label = checkContainerLabel(key.artifactUri, key.artifactCid, key.commit, this.#deps.instance) this.#deps.ledger.markRunning(key, { containerLabel: label }) this.#deps.log?.(`dispatching check ${key.artifactUri} @ ${key.commit} (label ${label})`) @@ -206,6 +209,7 @@ export class CheckDispatcher { timeoutMs: this.#deps.timeoutMs, allowedSchemes: this.#deps.allowedSchemes, ...(this.#deps.memory !== undefined ? { memory: this.#deps.memory } : {}), + ...(this.#deps.instance !== undefined ? { instance: this.#deps.instance } : {}), } const promise = this.#deps diff --git a/packages/daemon/src/check-runner.ts b/packages/daemon/src/check-runner.ts index abf31b8..0006930 100644 --- a/packages/daemon/src/check-runner.ts +++ b/packages/daemon/src/check-runner.ts @@ -30,11 +30,18 @@ export function checkrunRkey(artifactUri: string, artifactCid: string, commit: s return `check-${base32Encode(digest).slice(0, 24)}` } -/** `radial.check.`: the checkrun rkey with its `check-` prefix swapped for `radial.check.`, so - * it is a legal docker `--label`/`--name` and killable by that same string (mirrors - * `turnContainerLabel`). */ -export function checkContainerLabel(artifactUri: string, artifactCid: string, commit: string): string { - return `radial.check.${checkrunRkey(artifactUri, artifactCid, commit).replace(/^check-/, '')}` +/** `radial.check..`: the checkrun rkey with its `check-` prefix swapped for + * `radial.check.`, so it is a legal docker `--label`/`--name` and killable by that same string + * (mirrors `turnContainerLabel`, including the instance segment that keeps two daemons on one + * machine from colliding on a container name). */ +export function checkContainerLabel( + artifactUri: string, + artifactCid: string, + commit: string, + instance?: string, +): string { + const hash = checkrunRkey(artifactUri, artifactCid, commit).replace(/^check-/, '') + return `radial.check.${instance ? `${instance}.` : ''}${hash}` } const utf8Encoder = new TextEncoder() @@ -103,6 +110,9 @@ export interface CheckRunInput { memory?: string /** Optional daemon-owned GitHub credential for a private, validated github.com checkout. */ githubToken?: string + /** This daemon instance's id (`instanceId(stateDir)`), mixed into the container label/name — see + * `checkContainerLabel`. */ + instance?: string } export interface CheckRunDeps { @@ -238,7 +248,7 @@ export async function runCheckRun(input: CheckRunInput, deps: CheckRunDeps): Pro await writeFile(join(checksDir, 'run.sh'), orchestratorScript(input.checks.length)) const spec: ContainerSpec = { - label: checkContainerLabel(input.artifact.uri, input.artifact.cid, input.commit), + label: checkContainerLabel(input.artifact.uri, input.artifact.cid, input.commit, input.instance), image: input.image, argv: ['sh', '/checks/run.sh'], // Env is deliberately EMPTY (design §13): no ANTHROPIC_API_KEY, no RADIAL_* socket/token, no diff --git a/packages/daemon/src/cli.ts b/packages/daemon/src/cli.ts index ca87613..31c98dc 100644 --- a/packages/daemon/src/cli.ts +++ b/packages/daemon/src/cli.ts @@ -4,7 +4,7 @@ import { argv, stdin, stdout } from 'node:process' import { createHash } from 'node:crypto' import { realpathSync } from 'node:fs' import { mkdir, rm } from 'node:fs/promises' -import { join } from 'node:path' +import { dirname, isAbsolute, join } from 'node:path' import { pathToFileURL } from 'node:url' import { ClockSkewObserver, @@ -35,8 +35,6 @@ import { defaultConfigPath, extraIdentities, writeConfig, - findConfigPath, - loadConfig, parseModel, passwordEnvName, resolveAgentInits, @@ -45,8 +43,16 @@ import { type DaemonRunConfig, type Environment, type ForgeConfig, - type RadialConfig, } from './config.js' +import { + ConfigNotFoundError, + describeDataDirSource, + instanceId, + loadInstance, + resolveDataDir, + type Instance, +} from './instance.js' +import { StateDirLock } from './state-lock.js' import { FetchBobbin } from './bobbin.js' import { remoteBranchExists } from './bundle-writer.js' import { CheckDispatcher } from './check-dispatch.js' @@ -77,16 +83,20 @@ import { import { runTurn, turnContainerLabel, type TurnInput } from './turn.js' const usage = `Usage: - radiald config init [handle-or-did] [--config ] [--force] + radiald config init [handle-or-did] [--config ] [--data-dir ] [--force] writes a starter ${CONFIG_FILENAME} to ~/.config/radial/, prompting for the - agent handle when it is not given - radiald init [profile]... [--config ] [--password-stdin] [--update] + agent handle when it is not given. --data-dir writes a "dataDir" into it, so + that one --config selects this instance — sessions, state and all + radiald init [profile]... [--config ] [--data-dir ] + [--password-stdin] [--update] [--replace-identity] initializes every profile in ${CONFIG_FILENAME} (or only the named ones), read from the nearest ${CONFIG_FILENAME} or ~/.config/radial/${CONFIG_FILENAME}; the agent identity's app password is read from stdin or $RADIAL_PASSWORD. --update republishes a profile whose "artifactTypes" (or name, harness, or models) have changed since it was published — the agent record is what an - assignee menu offers, so editing ${CONFIG_FILENAME} alone changes nothing + assignee menu offers, so editing ${CONFIG_FILENAME} alone changes nothing. + --replace-identity re-registers a profile under a different DID, which is + otherwise refused (it is how one instance silently takes another's identity) radiald init --pds --identifier --name --harness <${HARNESSES.map((harness) => harness.name).join('|')}> [--model ]... --artifact-type ... --password-stdin [--update] @@ -97,27 +107,31 @@ const usage = `Usage: per section; two operators compare them to find a sync gap (design §12). It also prints "clocks": every author that has dated a record in this machine's future, so two operators comparing digests compare clocks too - radiald index reset [--config ] + radiald index reset [--config ] [--data-dir ] clears the local synced-record index and its checkpoints under "run.stateDir"; - stop radiald run before using this command (turn and check ledgers are preserved). + refuses while a radiald run holds that state dir (turn and check ledgers are + preserved). Re-ingesting also restarts every live claim's lease on THIS machine: a lease is honoured for its declared duration from first sighting, and after a reset every record is being seen for the first time again - radiald run [--config ] [--interval ] + radiald run [--config ] [--data-dir ] [--interval ] the long-lived daemon: polls every space in the config's "run" block, dispatches eligible plan requests to sandboxed containers, and durably - tracks each attempt in a turn ledger under "run.stateDir" - radiald turn list [--state ] [--config ] + tracks each attempt in a turn ledger under "run.stateDir". It prints the + config, data dir, sessions file and state dir it resolved, plus the DID + every profile is running as, and holds a lock on the state dir so a second + daemon cannot open the same ledgers + radiald turn list [--state ] [--config ] [--data-dir ] prints the turn ledger, newest first: state, when it last moved, the request uri, and attempts/retry/branch. --state narrows it — use "--state awaiting_input" to see every request parked on a question - radiald turn reset [--config ] + radiald turn reset [--config ] [--data-dir ] clears a request's ledger row (including a gave-up/cooling-down one) so it becomes dispatchable again - radiald check reset [--config ] + radiald check reset [--config ] [--data-dir ] clears every check-ledger row for an artifact uri (all versions/commits) so its checks run again - radiald claim reset [--config ] + radiald claim reset [--config ] [--data-dir ] clears this daemon's claim rows for an open request, so it may claim it again; one that has burned through every claim generation gets a fresh batch of them, starting above the records it already wrote` @@ -145,6 +159,7 @@ function one(args: string[], flag: string): string { const VALUELESS_FLAGS = new Set([ '--password-stdin', '--update', + '--replace-identity', '--force', '--skip-github-auth', '--watch', @@ -180,19 +195,24 @@ async function prompt(question: string): Promise { async function configInit(args: string[]): Promise { const path = values(args, '--config')[0] ?? defaultConfigPath() + const dataDir = values(args, '--data-dir')[0] let identifier = positional(args)[0] if (!identifier) { if (!stdin.isTTY) throw new Error(`Missing agent handle\n${usage}`) identifier = await prompt('Agent handle or DID: ') } if (!identifier) throw new Error('An agent handle or DID is required') - const config = buildDefaultConfig(identifier) + const config = buildDefaultConfig(identifier, dataDir ? { dataDir } : {}) await writeConfig(path, config, { force: args.includes('--force') }) const profiles = Object.keys(config.agents).join(', ') console.log(`Wrote ${path}`) console.log(` identity ${identifier}, profiles: ${profiles}`) + if (dataDir) console.log(` data dir ${dataDir} (resolved relative to ${dirname(path)})`) console.log(`Edit it, then register the profiles:`) - console.log(` printf '%s\\n' "$APP_PASSWORD" | radiald init --password-stdin`) + // The config path is spelled out rather than left implicit: with several instances on one machine, + // an init that discovers "the nearest radial.json" is exactly how the wrong one gets initialized. + console.log(` printf '%s\\n' "$APP_PASSWORD" | radiald init --config ${path} --password-stdin`) + console.log(` radiald run --config ${path}`) } export interface ResolvedRunConfig { @@ -244,14 +264,23 @@ function defaultRunStateDir(): string { } } -/** `config.js` validates the raw `run` block; the daemon applies runtime defaults. */ -export function resolveRunConfig(run: DaemonRunConfig): ResolvedRunConfig { +/** + * `config.js` validates the raw `run` block; the daemon applies runtime defaults. + * + * `options.dataDir` is this instance's resolved data directory: with it, `stateDir` defaults to + * `/run`, so one setting separates *everything* an instance owns. Without it the historic + * default stands, which is what keeps every existing call site (and test) meaning what it meant. + */ +export function resolveRunConfig( + run: DaemonRunConfig, + options: { dataDir?: string } = {}, +): ResolvedRunConfig { const turnTransport = run.turnTransport ?? 'unix' if (run.network?.toLowerCase() === 'host' || /^container:/i.test(run.network ?? '')) throw new Error('run.network must not use host or container networking') const mergePollInterval = Math.max(1, run.mergePollIntervalMs ?? 60_000) return { spaces: run.spaces, - stateDir: run.stateDir ?? defaultRunStateDir(), + stateDir: run.stateDir ?? (options.dataDir ? join(options.dataDir, 'run') : defaultRunStateDir()), image: run.image ?? 'radial-turn:latest', concurrency: run.concurrency ?? 1, timeoutMs: run.timeoutMs ?? 60 * 60_000, @@ -417,18 +446,33 @@ function checkRunDir(stateDir: string, key: { artifactUri: string; artifactCid: return join(stateDir, 'checks', hex.slice(0, 12)) } +/** `loadInstance` with this CLI's `--config`/`--data-dir` flags and its usage text on the one error + * an operator can fix by reading it. Every command that takes a config now takes a data dir too. */ +async function requireInstance(args: string[]): Promise { + const config = values(args, '--config')[0] + const dataDir = values(args, '--data-dir')[0] + try { + return await loadInstance({ config, dataDir }) + } catch (error) { + if (error instanceof ConfigNotFoundError) throw new Error(`${error.message}\n${usage}`) + throw error + } +} + +/** + * The whole argument vector rather than just `--config`, so `--data-dir` reaches every command that + * reads the `run` block — and so each of them resolves its state directory the same way `radiald + * run` does. A `turn list` that looked at a different ledger than the daemon writes would be worse + * than no `turn list` at all. + */ async function requireRunConfig( - configPath: string | undefined, -): Promise<{ configPath: string; config: RadialConfig; run: ResolvedRunConfig }> { - const path = await findConfigPath(configPath) - if (!path) { - throw new Error( - `No ${CONFIG_FILENAME} found in this directory or any parent, nor at ${defaultConfigPath()}\n${usage}`, - ) + args: string[], +): Promise<{ instance: Instance; run: ResolvedRunConfig }> { + const instance = await requireInstance(args) + if (!instance.config.run) { + throw new Error(`${instance.configPath} has no "run" block; add one before using "radiald run"\n${usage}`) } - const config = await loadConfig(path) - if (!config.run) throw new Error(`${path} has no "run" block; add one before using "radiald run"\n${usage}`) - return { configPath: path, config, run: resolveRunConfig(config.run) } + return { instance, run: resolveRunConfig(instance.config.run, { dataDir: instance.dataDir.path }) } } const TURN_STATES: TurnState[] = ['pending', 'running', 'awaiting_input', 'fulfilled', 'crashed', 'gave_up'] @@ -445,7 +489,7 @@ async function turnListCommand(args: string[]): Promise { if (state && !TURN_STATES.includes(state)) { throw new Error(`Unknown turn state '${state}'; one of ${TURN_STATES.join(', ')}\n${usage}`) } - const { run } = await requireRunConfig(values(args, '--config')[0]) + const { run } = await requireRunConfig(args) const ledger = new TurnLedger(join(run.stateDir, 'ledger.db'), { retryBound: run.retryBound, cooldownMs: run.cooldownMs, @@ -473,7 +517,7 @@ async function turnListCommand(args: string[]): Promise { async function turnResetCommand(args: string[]): Promise { const requestUri = args[0] if (!requestUri) throw new Error(`Missing \n${usage}`) - const { run } = await requireRunConfig(values(args, '--config')[0]) + const { run } = await requireRunConfig(args) const ledger = new TurnLedger(join(run.stateDir, 'ledger.db'), { retryBound: run.retryBound, cooldownMs: run.cooldownMs, @@ -494,7 +538,7 @@ async function turnResetCommand(args: string[]): Promise { async function claimResetCommand(args: string[]): Promise { const requestUri = args[0] if (!requestUri) throw new Error(`Missing \n${usage}`) - const { run } = await requireRunConfig(values(args, '--config')[0]) + const { run } = await requireRunConfig(args) const ledger = new ClaimLedger(join(run.stateDir, 'claim-ledger.db')) try { ledger.reset(requestUri, { generations: MAX_CLAIM_GENERATIONS }) @@ -507,7 +551,7 @@ async function claimResetCommand(args: string[]): Promise { async function checkResetCommand(args: string[]): Promise { const artifactUri = args[0] if (!artifactUri) throw new Error(`Missing \n${usage}`) - const { run } = await requireRunConfig(values(args, '--config')[0]) + const { run } = await requireRunConfig(args) const ledger = new CheckLedger(join(run.stateDir, 'check-ledger.db'), { retryBound: run.checkRetryBound, cooldownMs: run.checkCooldownMs, @@ -534,8 +578,16 @@ export async function resetRunIndex(stateDir: string): Promise { } async function indexResetCommand(args: string[]): Promise { - const { run } = await requireRunConfig(values(args, '--config')[0]) - await resetRunIndex(run.stateDir) + const { run } = await requireRunConfig(args) + // "Stop radiald run first" used to be a line in `usage` and nothing else — deleting the databases + // out from under a live daemon is exactly what the state-dir lock exists to make impossible. + await mkdir(run.stateDir, { recursive: true }) + const lock = StateDirLock.acquire(run.stateDir) + try { + await resetRunIndex(run.stateDir) + } finally { + lock.release() + } console.log(`reset local index under ${run.stateDir}`) } @@ -702,8 +754,63 @@ export function coverageAdvisories( return advisories } +/** + * The paths this invocation actually resolved, as the lines `radiald run` prints. + * + * Every failure this whole feature is about is an operator who believed a different set of paths: + * the wrong sessions file, a state directory shared with another daemon, a `$RADIAL_DATA_DIR` left + * over in a shell. None of them are visible from the command line that was typed, so the daemon says + * them out loud once, at the top, including WHICH rule chose the data directory. Pure and exported + * so the wording is testable, the same way `forgeAdvisory` and `coverageAdvisories` are. + */ +export function describeInstance(instance: Instance, run?: ResolvedRunConfig): string[] { + const lines = [ + ` config: ${instance.configPath}`, + ` data: ${instance.dataDir.path} (${describeDataDirSource(instance.dataDir)})`, + ` sessions: ${instance.sessions.path}`, + ] + if (run) lines.push(` state: ${run.stateDir} (instance ${instanceId(run.stateDir)})`) + if (instance.dataDir.shadowed) { + lines.push( + ` ⚠️ $RADIAL_DATA_DIR overrode "dataDir": ${instance.dataDir.shadowed} in ${instance.configPath} ` + + 'is not the directory in use. Unset it to run the instance the config describes.', + ) + } + // Unlike `dataDir`, a relative `run.stateDir` is resolved against the WORKING DIRECTORY (§5 of the + // plan: changing that would silently relocate the state of every operator using the documented + // `"stateDir": ".radial-state"`). Two instances started from two directories therefore get two + // state dirs from one config, which is worth saying out loud. + if (run && !isAbsolute(run.stateDir)) { + lines.push( + ` ⚠️ "run.stateDir" is relative, so it resolves against the working directory (${process.cwd()}). ` + + 'Make it absolute, or drop it and let it follow the data directory.', + ) + } + return lines +} + +/** What `radiald run` prints for each loaded profile: the identity it is actually running as, which + * is the fact a shared `sessions.json` used to get wrong in silence. */ +export function describeActors(actors: ActorRegistry): string[] { + const width = Math.max(...actors.all.map((actor) => actor.profile.length), 0) + return actors.all.map( + (actor) => ` ${actor.profile.padEnd(width)} ${actor.did} ${actor.session.handle} ${actor.harness}`, + ) +} + +/** The empty-registry error, naming the file that was actually read and what the config expected to + * find in it — "run radiald init" is no help when the question is *which* sessions file. */ +export function emptyRegistryError(instance: Instance): string { + const configured = Object.keys(instance.config.agents).join(', ') || 'none' + return ( + `No initialized agent profiles found for ${instance.configPath} in ${instance.sessions.path} ` + + `(config declares: ${configured}). Run: radiald init --config ${instance.configPath}` + ) +} + async function runCommand(args: string[]): Promise { - const { configPath, config, run } = await requireRunConfig(values(args, '--config')[0]) + const { instance, run } = await requireRunConfig(args) + const { configPath, config } = instance // `run.spaces` may be empty in the file (the init scaffold writes `spaces: []` so the block it // emits is loadable); a daemon with nothing to poll is only an error here, where polling starts. if (run.spaces.length === 0) { @@ -719,8 +826,9 @@ async function runCommand(args: string[]): Promise { console.log( `radiald starting: ${run.spaces.length} space(s), interval ${interval}ms, concurrency ${run.concurrency}, ` + - `open egress\n config: ${configPath}\n state: ${run.stateDir}`, + 'open egress', ) + for (const line of describeInstance(instance, run)) console.log(line) console.log( `claims: lease ${run.claims.leaseMs}ms, renew every ${run.claims.renewIntervalMs}ms, ` + `confirm after ${run.claims.confirmCycles} ingestion cycle(s), at most ${run.claims.maxOutstanding} outstanding, ` + @@ -740,312 +848,323 @@ async function runCommand(args: string[]): Promise { await mkdir(join(run.stateDir, 'checks'), { recursive: true }) await mkdir(join(run.stateDir, 'auto-review'), { recursive: true }) - const store = new FileSessionStore() - // Every PDS response already carries a `Date` header, and reading it is the only way this daemon - // can detect a clock that runs BEHIND its peers — the direction record timestamps cannot - // distinguish from ingestion lag, and the direction that quietly wins claim races. It rides on - // requests that had to happen anyway and only ever logs. - const clocks = new ClockSkewObserver() - const probingFetch = observingFetch((input, init) => fetch(input, init), clocks) - const actors = await loadActors(config, store, { fetcher: probingFetch }) - if (actors.all.length === 0) { - throw new Error('No initialized agent profiles found for this config; run "radiald init" first') - } - console.log( - `loaded ${actors.all.length} agent profile(s): ` + - actors.all.map((actor) => `${actor.profile} (${actor.harness})`).join(', '), - ) + // One daemon per state directory. Two are two writers on `records.db`, `sync.db` and three + // ledgers, and nothing downstream of here would notice. Released in the `finally` at the bottom. + const stateLock = StateDirLock.acquire(run.stateDir) + const instanceLabel = instanceId(run.stateDir) - // Fail closed at startup rather than per dispatch: an unrunnable harness or a profile with no - // provider credential is knowable now, and finding out per-request costs a ledger attempt and a - // cooldown each time. Names only in the log — never a value. - const modelEnv = resolveModelEnvironment(actors.all, run.modelEnv, process.env) - console.log(`model credentials forwarded to turns: ${Object.keys(modelEnv).join(', ') || '(none)'}`) + try { + const store = instance.sessions + // Every PDS response already carries a `Date` header, and reading it is the only way this daemon + // can detect a clock that runs BEHIND its peers — the direction record timestamps cannot + // distinguish from ingestion lag, and the direction that quietly wins claim races. It rides on + // requests that had to happen anyway and only ever logs. + const clocks = new ClockSkewObserver() + const probingFetch = observingFetch((input, init) => fetch(input, init), clocks) + const actors = await loadActors(config, store, { + fetcher: probingFetch, + configPath, + warn: (message) => console.warn(`⚠️ ${message}`), + }) + if (actors.all.length === 0) throw new Error(emptyRegistryError(instance)) + console.log(`loaded ${actors.all.length} agent profile(s):`) + for (const line of describeActors(actors)) console.log(line) - const ledger = new TurnLedger(join(run.stateDir, 'ledger.db'), { - retryBound: run.retryBound, - cooldownMs: run.cooldownMs, - }) - const checkLedger = new CheckLedger(join(run.stateDir, 'check-ledger.db'), { - retryBound: run.checkRetryBound, - cooldownMs: run.checkCooldownMs, - }) - // Its own file: claims and turns are independent workstreams with independent lifecycles, and a - // daemon that restarts mid-lease reads this to renew rather than claim twice. - const claimLedger = new ClaimLedger(join(run.stateDir, 'claim-ledger.db')) - 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 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) - await reconcileCheckOrphans(runner, checkLedger) - - - // One adapter per configured forge, selected per project by `gitUrl` host. GitHub authentication - // is resolved once but no longer gates the whole daemon: a tangled-only operator has no use for - // it, and a space that holds both kinds of project must be able to serve both. - const githubToken = await resolveGitHubToken() - const keyStore = new TangledKeyStore(run.stateDir) - const adapters: ForgeAdapter[] = [] - for (const configured of run.forges) { - if (configured.kind === 'github') { - if (!githubToken) { + // Fail closed at startup rather than per dispatch: an unrunnable harness or a profile with no + // provider credential is knowable now, and finding out per-request costs a ledger attempt and a + // cooldown each time. Names only in the log — never a value. + const modelEnv = resolveModelEnvironment(actors.all, run.modelEnv, process.env) + console.log(`model credentials forwarded to turns: ${Object.keys(modelEnv).join(', ') || '(none)'}`) + + const ledger = new TurnLedger(join(run.stateDir, 'ledger.db'), { + retryBound: run.retryBound, + cooldownMs: run.cooldownMs, + }) + const checkLedger = new CheckLedger(join(run.stateDir, 'check-ledger.db'), { + retryBound: run.checkRetryBound, + cooldownMs: run.checkCooldownMs, + }) + // Its own file: claims and turns are independent workstreams with independent lifecycles, and a + // daemon that restarts mid-lease reads this to renew rather than claim twice. + const claimLedger = new ClaimLedger(join(run.stateDir, 'claim-ledger.db')) + 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 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) + await reconcileCheckOrphans(runner, checkLedger) + + + // One adapter per configured forge, selected per project by `gitUrl` host. GitHub authentication + // is resolved once but no longer gates the whole daemon: a tangled-only operator has no use for + // it, and a space that holds both kinds of project must be able to serve both. + const githubToken = await resolveGitHubToken() + const keyStore = new TangledKeyStore(run.stateDir) + const adapters: ForgeAdapter[] = [] + for (const configured of run.forges) { + if (configured.kind === 'github') { + if (!githubToken) { + console.warn( + 'run.forge is "github" but GitHub authentication is unavailable; merge observation is disabled.', + ) + } + adapters.push(new GitHubForge({ fetch, ...(githubToken ? { token: githubToken } : {}) })) + if (githubToken) console.log('forge: github observation adapter enabled') + continue + } + // The tangled adapter writes as the SAME agent identity that signs the artifact — no forge + // account, no operator token in the middle. The implementation-producing actor is the one that + // authors its pull records, chosen the same deterministic way the merge poller picks a writer. + const writer = (): (typeof actors.all)[number] | undefined => + [...actors.all].sort((a, b) => compareCodePoints(a.did, b.did)).find((actor) => actor.artifactTypes.includes('implementation')) ?? + [...actors.all].sort((a, b) => compareCodePoints(a.did, b.did))[0] + const knownHosts = configured.knownHosts ?? [] + // Pull state comes from Bobbin, tangled's public read-only appview, unless the operator turned + // it off — see `bobbin.ts` for why an appview answers this better than a PDS scan can. + const bobbin = + configured.api === false ? undefined : new FetchBobbin({ fetch, ...(configured.api ? { url: configured.api } : {}) }) + const tangled = new TangledForge({ + reader: new FetchRepoTransport(), + ...(bobbin ? { bobbin } : {}), + ...(configured.hosts ? { hosts: configured.hosts } : {}), + resolveHandle: (handle) => resolveHandleDid(handle, { resolveTxt: nodeResolveTxt }), + actor: () => { + const actor = writer() + if (!actor) return undefined + return { + did: actor.did, + createForeign: (collection, value) => actor.client.createForeign(collection, value), + putForeign: (collection, uri, value, options) => + actor.client.putForeign(collection, uri, value, options), + uploadBlob: (bytes, contentType) => actor.client.uploadBlob(bytes, contentType), + } + }, + // Push material is only assembled when a turn actually needs it, so a daemon that only + // OBSERVES tangled pulls needs neither a key nor pinned host keys. + push: async () => { + const actor = writer() + if (!actor || knownHosts.length === 0) return undefined + const privateKey = await keyStore.read(actor.did) + if (!privateKey) return undefined + return { privateKey, knownHosts: knownHostsFile(knownHosts) } + }, + log: (message) => console.log(message), + }) + adapters.push(tangled) + console.log(`forge: tangled adapter enabled for ${(configured.hosts ?? ['tangled.org']).join(', ')}`) + if (knownHosts.length === 0) { console.warn( - 'run.forge is "github" but GitHub authentication is unavailable; merge observation is disabled.', + 'the tangled forge has no "knownHosts": implementation turns cannot push until you add ' + + `ssh-keyscan output for its hosts to ${configPath}. Radial will not disable host-key checking.`, ) } - adapters.push(new GitHubForge({ fetch, ...(githubToken ? { token: githubToken } : {}) })) - if (githubToken) console.log('forge: github observation adapter enabled') - continue } - // The tangled adapter writes as the SAME agent identity that signs the artifact — no forge - // account, no operator token in the middle. The implementation-producing actor is the one that - // authors its pull records, chosen the same deterministic way the merge poller picks a writer. - const writer = (): (typeof actors.all)[number] | undefined => - [...actors.all].sort((a, b) => compareCodePoints(a.did, b.did)).find((actor) => actor.artifactTypes.includes('implementation')) ?? - [...actors.all].sort((a, b) => compareCodePoints(a.did, b.did))[0] - const knownHosts = configured.knownHosts ?? [] - // Pull state comes from Bobbin, tangled's public read-only appview, unless the operator turned - // it off — see `bobbin.ts` for why an appview answers this better than a PDS scan can. - const bobbin = - configured.api === false ? undefined : new FetchBobbin({ fetch, ...(configured.api ? { url: configured.api } : {}) }) - const tangled = new TangledForge({ - reader: new FetchRepoTransport(), - ...(bobbin ? { bobbin } : {}), - ...(configured.hosts ? { hosts: configured.hosts } : {}), - resolveHandle: (handle) => resolveHandleDid(handle, { resolveTxt: nodeResolveTxt }), - actor: () => { - const actor = writer() - if (!actor) return undefined - return { - did: actor.did, - createForeign: (collection, value) => actor.client.createForeign(collection, value), - putForeign: (collection, uri, value, options) => - actor.client.putForeign(collection, uri, value, options), - uploadBlob: (bytes, contentType) => actor.client.uploadBlob(bytes, contentType), - } - }, - // Push material is only assembled when a turn actually needs it, so a daemon that only - // OBSERVES tangled pulls needs neither a key nor pinned host keys. - push: async () => { - const actor = writer() - if (!actor || knownHosts.length === 0) return undefined - const privateKey = await keyStore.read(actor.did) - if (!privateKey) return undefined - return { privateKey, knownHosts: knownHostsFile(knownHosts) } - }, - log: (message) => console.log(message), + const forges = new ForgeRegistry(adapters) + // Legacy single-adapter shape, still what the dispatcher's `forge` and the merge poller consume + // where a URL alone has to be routed. Undefined when nothing can observe. + const forge = adapters.some((adapter) => adapter.canObserve) ? routingForgeState(forges) : undefined + const advisory = forgeAdvisory({ + adapter: Boolean(forge), + configured: run.forges.length > 0, + producesImplementation: actors.all.some((actor) => actor.artifactTypes.includes('implementation')), + configPath, }) - adapters.push(tangled) - console.log(`forge: tangled adapter enabled for ${(configured.hosts ?? ['tangled.org']).join(', ')}`) - if (knownHosts.length === 0) { - console.warn( - 'the tangled forge has no "knownHosts": implementation turns cannot push until you add ' + - `ssh-keyscan output for its hosts to ${configPath}. Radial will not disable host-key checking.`, - ) - } - } - const forges = new ForgeRegistry(adapters) - // Legacy single-adapter shape, still what the dispatcher's `forge` and the merge poller consume - // where a URL alone has to be routed. Undefined when nothing can observe. - const forge = adapters.some((adapter) => adapter.canObserve) ? routingForgeState(forges) : undefined - const advisory = forgeAdvisory({ - adapter: Boolean(forge), - configured: run.forges.length > 0, - producesImplementation: actors.all.some((actor) => actor.artifactTypes.includes('implementation')), - configPath, - }) - if (advisory) console.warn(advisory) - - const boundRunTurn = (input: TurnInput) => - runTurn(input, { - runner, - harnesses: { select: selectHarness }, - modelEnv, - ...(githubToken ? { githubToken } : {}), + if (advisory) console.warn(advisory) + + const boundRunTurn = (input: TurnInput) => + runTurn(input, { + runner, + harnesses: { select: selectHarness }, + modelEnv, + ...(githubToken ? { githubToken } : {}), + forges, + log: (message) => console.log(message), + }) + + const dispatcher = new TurnDispatcher({ + ledger, + concurrency: run.concurrency, + runDirFor: (requestUri) => turnRunDir(run.stateDir, requestUri), + image: run.image, + timeoutMs: run.timeoutMs, + ...(run.network !== undefined ? { network: run.network } : {}), + allowedSchemes: run.gitSchemes, + ...(run.memory !== undefined ? { memory: run.memory } : {}), + turnTransport: run.turnTransport, + instance: instanceLabel, + ...(forge ? { forge } : {}), forges, + // The branch probe the no-forge reuse fallback leans on, carrying the daemon's token so it can + // answer for a private repository too. + branchExists: ({ gitUrl, branch }) => + remoteBranchExists({ + gitUrl, + branch, + allowedSchemes: run.gitSchemes, + ...(githubToken ? { githubToken } : {}), + }), + implementationEnabled: !!githubToken, + runTurn: boundRunTurn, + // Unassigned requests dispatch only for a claim this daemon wrote and has watched win. + heldClaims: () => claimLedger.heldRequests(), + kill: (label) => runner.kill(label), log: (message) => console.log(message), }) - const dispatcher = new TurnDispatcher({ - ledger, - concurrency: run.concurrency, - runDirFor: (requestUri) => turnRunDir(run.stateDir, requestUri), - image: run.image, - timeoutMs: run.timeoutMs, - ...(run.network !== undefined ? { network: run.network } : {}), - allowedSchemes: run.gitSchemes, - ...(run.memory !== undefined ? { memory: run.memory } : {}), - turnTransport: run.turnTransport, - ...(forge ? { forge } : {}), - forges, - // The branch probe the no-forge reuse fallback leans on, carrying the daemon's token so it can - // answer for a private repository too. - branchExists: ({ gitUrl, branch }) => - remoteBranchExists({ - gitUrl, - branch, - allowedSchemes: run.gitSchemes, - ...(githubToken ? { githubToken } : {}), - }), - implementationEnabled: !!githubToken, - runTurn: boundRunTurn, - // Unassigned requests dispatch only for a claim this daemon wrote and has watched win. - heldClaims: () => claimLedger.heldRequests(), - kill: (label) => runner.kill(label), - log: (message) => console.log(message), - }) - - // Checks have an empty container environment. A GitHub token, when available, is used solely by - // the daemon-host checkout helper for private repositories and is never added to the container. - const boundRunCheckRun = (input: CheckRunInput) => - runCheckRun(githubToken ? { ...input, githubToken } : input, { runner }) - const checkDispatcher = new CheckDispatcher({ - ledger: checkLedger, - concurrency: run.checkConcurrency, - ourDids: new Set(actors.all.map((actor) => actor.did)), - runDirFor: (key) => checkRunDir(run.stateDir, key), - image: run.checkImage, - timeoutMs: run.checkTimeoutMs, - allowedSchemes: run.gitSchemes, - ...(run.memory !== undefined ? { memory: run.memory } : {}), - runCheckRun: boundRunCheckRun, - log: (message) => console.log(message), - }) + // Checks have an empty container environment. A GitHub token, when available, is used solely by + // the daemon-host checkout helper for private repositories and is never added to the container. + const boundRunCheckRun = (input: CheckRunInput) => + runCheckRun(githubToken ? { ...input, githubToken } : input, { runner }) + const checkDispatcher = new CheckDispatcher({ + ledger: checkLedger, + concurrency: run.checkConcurrency, + ourDids: new Set(actors.all.map((actor) => actor.did)), + runDirFor: (key) => checkRunDir(run.stateDir, key), + image: run.checkImage, + timeoutMs: run.checkTimeoutMs, + allowedSchemes: run.gitSchemes, + ...(run.memory !== undefined ? { memory: run.memory } : {}), + instance: instanceLabel, + runCheckRun: boundRunCheckRun, + log: (message) => console.log(message), + }) - const spaces = run.spaces.map((spaceUri) => ({ - spaceUri, - runtime: new SpaceSyncRuntime( + const spaces = run.spaces.map((spaceUri) => ({ spaceUri, - { - records: join(run.stateDir, 'records.db'), - checkpoints: join(run.stateDir, 'sync.db'), - }, - new FetchRepoTransport(undefined, probingFetch), - ), - })) - - // Both halves of skew detection, reported on the same fresh index the dispatchers see: the - // symmetric `Date`-header probe above, and the one-sided "somebody wrote a record dated in our - // future" heuristic. Advisory only — the fold's rules are durations and ingestion cycles, never a - // comparison of two clocks (see docs/operators.md §3, "Claims and clocks"). - const skew = new SkewWatch({ - thresholdMs: clockSkewThresholdMs(run.claims.leaseMs), - samples: () => clocks.samples(), - log: (message) => console.warn(`⚠️ ${message}`), - }) + runtime: new SpaceSyncRuntime( + spaceUri, + { + records: join(run.stateDir, 'records.db'), + checkpoints: join(run.stateDir, 'sync.db'), + }, + new FetchRepoTransport(undefined, probingFetch), + ), + })) - // Claim manager (design §11): the only thing that lets this daemon touch an OPEN (unassigned) - // request. It writes a claim, waits to see it win the fold, and only then does `heldClaims` let - // the dispatcher launch. Assigned requests never pass through it and gain no latency. - const claimManager = new ClaimManager({ - ledger: claimLedger, - turns: ledger, - timing: run.claims, - now: () => Date.now(), - capacity: () => run.concurrency - dispatcher.inFlight, - abandon: (requestUri, reason) => dispatcher.abandon(requestUri, reason), - log: (message) => console.log(message), - }) + // Both halves of skew detection, reported on the same fresh index the dispatchers see: the + // symmetric `Date`-header probe above, and the one-sided "somebody wrote a record dated in our + // future" heuristic. Advisory only — the fold's rules are durations and ingestion cycles, never a + // comparison of two clocks (see docs/operators.md §3, "Claims and clocks"). + const skew = new SkewWatch({ + thresholdMs: clockSkewThresholdMs(run.claims.leaseMs), + samples: () => clocks.samples(), + log: (message) => console.warn(`⚠️ ${message}`), + }) - // Jetstream (design §5), opt-in. It shortens ingestion latency and nothing else: polling stays - // the authority, so an endpoint that is unreachable degrades to exactly the pre-Phase-7 behaviour. - let wakeLoop: (() => void) | undefined - if (run.jetstream) { - const endpoint = run.jetstream.endpoint - console.log( - `jetstream: subscribing to ${endpoint} (polling continues as the authority and backfill, ` + - `every ${run.jetstream.backfillIntervalMs}ms while the stream is healthy)`, - ) - for (const { spaceUri, runtime } of spaces) { - runtime.startJetstream({ - endpoint, - degradedAfterMs: 2 * interval, - onRecord: () => wakeLoop?.(), - onStatus: (status, detail) => - console.log(`jetstream ${spaceUri}: ${status}${detail ? ` (${detail})` : ''}`), - log: (message) => console.log(message), - }) + // Claim manager (design §11): the only thing that lets this daemon touch an OPEN (unassigned) + // request. It writes a claim, waits to see it win the fold, and only then does `heldClaims` let + // the dispatcher launch. Assigned requests never pass through it and gain no latency. + const claimManager = new ClaimManager({ + ledger: claimLedger, + turns: ledger, + timing: run.claims, + now: () => Date.now(), + capacity: () => run.concurrency - dispatcher.inFlight, + abandon: (requestUri, reason) => dispatcher.abandon(requestUri, reason), + log: (message) => console.log(message), + }) + + // Jetstream (design §5), opt-in. It shortens ingestion latency and nothing else: polling stays + // the authority, so an endpoint that is unreachable degrades to exactly the pre-Phase-7 behaviour. + let wakeLoop: (() => void) | undefined + if (run.jetstream) { + const endpoint = run.jetstream.endpoint + console.log( + `jetstream: subscribing to ${endpoint} (polling continues as the authority and backfill, ` + + `every ${run.jetstream.backfillIntervalMs}ms while the stream is healthy)`, + ) + for (const { spaceUri, runtime } of spaces) { + runtime.startJetstream({ + endpoint, + degradedAfterMs: 2 * interval, + onRecord: () => wakeLoop?.(), + onStatus: (status, detail) => + console.log(`jetstream ${spaceUri}: ${status}${detail ? ` (${detail})` : ''}`), + log: (message) => console.log(message), + }) + } } - } - // Merge-observation poller (design §10): the GitHub adapter's `getPullRequestState` structurally - // satisfies MergePoller's `ForgeStateSource`, so the same `forge` drives it. Only runs when the - // forge is configured — no forge, no polling. - const mergePoller = forge - ? new MergePoller({ - adapter: forge, - now: () => Date.now(), - pollIntervalMs: run.mergePollIntervalMs, - backoffMaxMs: run.mergePollBackoffMaxMs, - log: (message) => console.log(message), - }) - : undefined + // Merge-observation poller (design §10): the GitHub adapter's `getPullRequestState` structurally + // satisfies MergePoller's `ForgeStateSource`, so the same `forge` drives it. Only runs when the + // forge is configured — no forge, no polling. + const mergePoller = forge + ? new MergePoller({ + adapter: forge, + now: () => Date.now(), + pollIntervalMs: run.mergePollIntervalMs, + backoffMaxMs: run.mergePollBackoffMaxMs, + log: (message) => console.log(message), + }) + : undefined - // Auto-review trigger (design §7): writes a `type: "review"` request the moment an artifact whose - // effective auto-review config is on lands. No forge dependency — it runs whenever actors exist; - // its observed-subjects ledger lives under `/auto-review/`. - const autoReviewTrigger = new AutoReviewTrigger({ - stateDir: run.stateDir, - now: () => new Date().toISOString(), - log: (message) => console.log(message), - }) + // Auto-review trigger (design §7): writes a `type: "review"` request the moment an artifact whose + // effective auto-review config is on lands. No forge dependency — it runs whenever actors exist; + // its observed-subjects ledger lives under `/auto-review/`. + const autoReviewTrigger = new AutoReviewTrigger({ + stateDir: run.stateDir, + now: () => new Date().toISOString(), + log: (message) => console.log(message), + }) - const controller = new AbortController() - const stop = (): void => controller.abort() - process.on('SIGTERM', stop) - process.on('SIGINT', stop) + const controller = new AbortController() + const stop = (): void => controller.abort() + process.on('SIGTERM', stop) + process.on('SIGINT', stop) - try { - await runDaemon({ - spaces, - dispatcher, - checkDispatcher, - ...(mergePoller ? { mergePoller } : {}), - autoReviewTrigger, - claimManager, - skew, - actors, - // While every space's stream is healthy, polling drops to the backfill rate; the moment one - // is not, the short interval is back. Read afresh each tick, so it tracks the stream live. - intervalMs: run.jetstream - ? () => - spaces.every(({ runtime }) => runtime.streaming) - ? (run.jetstream as { backfillIntervalMs: number }).backfillIntervalMs - : interval - : interval, - onWake: (wake) => { - wakeLoop = wake - }, - // The registry lives in the space, so this is the first moment the daemon can tell an operator - // that a type nothing it runs accepts has been sitting in it. - onFirstIndex: (spaceUri, index) => { - for (const advisory of coverageAdvisories(index, actors, { spaceUri, configPath })) { - console.warn(`⚠️ ${advisory}`) - } - }, - signal: controller.signal, - log: (message) => console.error(message), - }) - } finally { - // Graceful shutdown: kill every in-flight container, then wait for its turn promise to - // settle (each already never rejects — see TurnDispatcher.pump) so the ledger reflects a - // clean "crashed" row rather than leaving it stuck at "running" for next startup's orphan - // reconciliation to redo. Every step here is best-effort — shutdown must not hang. - for (const label of [...dispatcher.runningLabels(), ...checkDispatcher.runningLabels()]) { - await runner.kill(label).catch(() => {}) + try { + await runDaemon({ + spaces, + dispatcher, + checkDispatcher, + ...(mergePoller ? { mergePoller } : {}), + autoReviewTrigger, + claimManager, + skew, + actors, + // While every space's stream is healthy, polling drops to the backfill rate; the moment one + // is not, the short interval is back. Read afresh each tick, so it tracks the stream live. + intervalMs: run.jetstream + ? () => + spaces.every(({ runtime }) => runtime.streaming) + ? (run.jetstream as { backfillIntervalMs: number }).backfillIntervalMs + : interval + : interval, + onWake: (wake) => { + wakeLoop = wake + }, + // The registry lives in the space, so this is the first moment the daemon can tell an operator + // that a type nothing it runs accepts has been sitting in it. + onFirstIndex: (spaceUri, index) => { + for (const advisory of coverageAdvisories(index, actors, { spaceUri, configPath })) { + console.warn(`⚠️ ${advisory}`) + } + }, + signal: controller.signal, + log: (message) => console.error(message), + }) + } finally { + // Graceful shutdown: kill every in-flight container, then wait for its turn promise to + // settle (each already never rejects — see TurnDispatcher.pump) so the ledger reflects a + // clean "crashed" row rather than leaving it stuck at "running" for next startup's orphan + // reconciliation to redo. Every step here is best-effort — shutdown must not hang. + for (const label of [...dispatcher.runningLabels(), ...checkDispatcher.runningLabels()]) { + await runner.kill(label).catch(() => {}) + } + await dispatcher.drain().catch(() => {}) + await checkDispatcher.drain().catch(() => {}) + await autoReviewTrigger.drain().catch(() => {}) + await claimManager.drain().catch(() => {}) + for (const { runtime } of spaces) runtime.close() + ledger.close() + checkLedger.close() + claimLedger.close() } - await dispatcher.drain().catch(() => {}) - await checkDispatcher.drain().catch(() => {}) - await autoReviewTrigger.drain().catch(() => {}) - await claimManager.drain().catch(() => {}) - for (const { runtime } of spaces) runtime.close() - ledger.close() - checkLedger.close() - claimLedger.close() + } finally { + stateLock.release() } } @@ -1114,8 +1233,13 @@ export async function main(args = argv.slice(2)): Promise { return } if (args[0] !== 'init') throw new Error(usage) + const replaceIdentity = args.includes('--replace-identity') if (args.includes('--identifier')) { if (!args[1] || !args.includes('--password-stdin')) throw new Error(usage) + // The flag-driven path has no config to read a `dataDir` from, but `--data-dir` (and + // `$RADIAL_DATA_DIR`) still decide where the session lands — it used to take the process default + // silently, which is the one place an instance could write to the wrong sessions file. + const store = new FileSessionStore(resolveDataDir({ flag: values(args, '--data-dir')[0] }).path) const result = await initializeAgent({ profile: args[1], service: one(args, '--pds'), @@ -1127,25 +1251,23 @@ export async function main(args = argv.slice(2)): Promise { harness: selectHarness(one(args, '--harness')).name, models: values(args, '--model').map(parseModel), artifactTypes: values(args, '--artifact-type'), - }, { update: args.includes('--update') }) + }, { store, update: args.includes('--update'), replaceIdentity }) console.log(JSON.stringify(result)) + console.log(`sessions: ${store.path}`) await ensureGitHubAuth({ skip: args.includes('--skip-github-auth') }) return } - const path = await findConfigPath(values(args, '--config')[0]) - if (!path) { - throw new Error( - `No ${CONFIG_FILENAME} found in this directory or any parent, nor at ${defaultConfigPath()}\n${usage}`, - ) - } - const config = await loadConfig(path) + const instance = await requireInstance(args) + const { configPath: path, config } = instance // Which forges this operator serves decides what `init` has to set up. A tangled forge needs an // ed25519 push key published under the agent DID; GitHub needs `gh` to be logged in. An operator // who serves only one should not be walked through the other's login flow. const forgeKinds = new Set(configuredForges(config.run).map((forge) => forge.kind)) const wantsGitHub = forgeKinds.size === 0 || forgeKinds.has('github') + // The key belongs to the instance being initialized, so it follows the resolved data directory + // rather than the process-wide default — the same rule `radiald run` uses for `run.stateDir`. const tangledKeys = forgeKinds.has('tangled') - ? new TangledKeyStore(config.run?.stateDir ?? defaultRunStateDir()) + ? new TangledKeyStore(config.run?.stateDir ?? join(instance.dataDir.path, 'run')) : undefined const profiles = positional(args.slice(1)) const byIdentifier: Record = {} @@ -1162,7 +1284,12 @@ export async function main(args = argv.slice(2)): Promise { const width = Math.max(...inits.map((init) => init.profile.length), 0) for (const init of inits) { try { - const result = await initializeAgent(init, { update, ...(tangledKeys ? { tangledKeys } : {}) }) + const result = await initializeAgent(init, { + store: instance.sessions, + update, + replaceIdentity, + ...(tangledKeys ? { tangledKeys } : {}), + }) // The published artifact types, not just the DID: this line is the only place an operator sees // what the space will offer this profile for, and a profile that has not caught up with a // registry type is invisible otherwise. @@ -1186,6 +1313,9 @@ export async function main(args = argv.slice(2)): Promise { console.error(`✗ ${init.profile.padEnd(width)} ${error instanceof Error ? error.message : error}`) } } + // Said once, after the per-profile lines: which file these identities were written to is the fact + // that decides whether `radiald run` picks them up, and it is invisible from the command typed. + console.log(`sessions: ${instance.sessions.path} (${describeDataDirSource(instance.dataDir)})`) if (failed) process.exitCode = 1 else if (wantsGitHub) await ensureGitHubAuth({ skip: args.includes('--skip-github-auth') }) } diff --git a/packages/daemon/src/config.ts b/packages/daemon/src/config.ts index 355f528..1dd106b 100644 --- a/packages/daemon/src/config.ts +++ b/packages/daemon/src/config.ts @@ -19,6 +19,13 @@ export interface AgentConfig { export interface RadialConfig { identifier?: string | undefined + /** + * Where this instance's sessions and (by default) its run state live, resolved relative to the + * DIRECTORY OF THIS FILE — see `resolveDataDir` in `instance.ts`, which owns the precedence. + * Naming it here is what lets `--config` alone select one of several daemons on a machine, with + * no `$RADIAL_DATA_DIR` discipline on every invocation. + */ + dataDir?: string | undefined pds?: string | undefined harness?: string | undefined models?: Array | undefined @@ -425,6 +432,7 @@ export function parseConfig(value: unknown): RadialConfig { } return { identifier: text(value.identifier, 'identifier'), + dataDir: text(value.dataDir, 'dataDir'), pds: text(value.pds, 'pds'), harness: text(value.harness, 'harness'), models: models(value.models, 'models'), @@ -580,9 +588,16 @@ export function extraIdentities(config: RadialConfig, profiles: string[]): strin * updating the one under review. `spaces` is left empty for the operator to fill in — every other * `run` default is computed at runtime and deliberately not frozen into the file. */ -export function buildDefaultConfig(identifier: string): RadialConfig { +export function buildDefaultConfig( + identifier: string, + options: { dataDir?: string } = {}, +): RadialConfig { return { identifier, + // Emitted only when asked for. A scaffold that names its own data directory is what makes a + // SECOND instance on this machine one `--config` away: see docs/radial-json.md, "More than one + // daemon on one machine". + ...(options.dataDir ? { dataDir: options.dataDir } : {}), // `claude` stays the scaffold default: it is the tested path, and changing what a new operator // gets is a separate decision from making a second harness available. Switching a profile (or // the whole config) to `pi` — and picking a non-Anthropic provider — is one field plus one @@ -627,6 +642,7 @@ export async function writeConfig( // that got a fresh pull request for every implementation v2. const body = { identifier: config.identifier, + ...(config.dataDir ? { dataDir: config.dataDir } : {}), ...(config.pds ? { pds: config.pds } : {}), ...(config.harness ? { harness: config.harness } : {}), ...(config.models ? { models: config.models } : {}), diff --git a/packages/daemon/src/dispatch.ts b/packages/daemon/src/dispatch.ts index 0058739..0282242 100644 --- a/packages/daemon/src/dispatch.ts +++ b/packages/daemon/src/dispatch.ts @@ -546,6 +546,9 @@ export interface DispatcherDeps { memory?: string /** Threaded into `TurnInput.turnTransport`. Unset lets `turn.ts` default to `unix`. */ turnTransport?: 'unix' | 'tcp' + /** This daemon instance's id (`instanceId(stateDir)`), mixed into every container label so two + * daemons on one machine never derive the same docker name for one request. */ + instance?: string /** Injected wrapper around turn.ts's runTurn that binds TurnDeps (runner/harness/secrets/etc). */ runTurn: (input: TurnInput) => Promise /** The request URIs whose claim this daemon has confirmed, read fresh on every pump — supplied by @@ -724,7 +727,7 @@ export class TurnDispatcher { } this.#contended.delete(uri) - const label = turnContainerLabel(uri, cid) + const label = turnContainerLabel(uri, cid, this.#deps.instance) const runDir = this.#deps.runDirFor(uri) this.#deps.ledger.markRunning(uri, cid, { containerLabel: label, checkoutPath: join(runDir, 'checkout') }) this.#deps.log?.(`dispatching turn ${uri} (label ${label})`) @@ -863,6 +866,7 @@ export class TurnDispatcher { allowedSchemes: this.#deps.allowedSchemes, ...(this.#deps.memory !== undefined ? { memory: this.#deps.memory } : {}), ...(this.#deps.turnTransport !== undefined ? { turnTransport: this.#deps.turnTransport } : {}), + ...(this.#deps.instance !== undefined ? { instance: this.#deps.instance } : {}), ...(checkoutRef !== undefined ? { checkoutRef } : {}), // `prev` is scope- and type-agnostic: any v2 artifact continues its predecessor's chain. Only // `branch` remains implementation-specific. diff --git a/packages/daemon/src/index.ts b/packages/daemon/src/index.ts index 9b69ebc..6617385 100644 --- a/packages/daemon/src/index.ts +++ b/packages/daemon/src/index.ts @@ -17,8 +17,10 @@ export * from './forge-tangled.js' export * from './github-auth.js' export * from './harness.js' export * from './init.js' +export * from './instance.js' export * from './ledger.js' export * from './merge-poll.js' export * from './runtime.js' +export * from './state-lock.js' export * from './turn.js' export * from './turn-socket.js' diff --git a/packages/daemon/src/init.ts b/packages/daemon/src/init.ts index 33c1950..c63ca63 100644 --- a/packages/daemon/src/init.ts +++ b/packages/daemon/src/init.ts @@ -148,6 +148,15 @@ export async function initializeAgent( tangledKeys?: TangledKeyStore /** Foreign-record reads for the key-publication check. Defaults to a `FetchRepoTransport`. */ reader?: ForeignRecordReader + /** + * Overwrite a stored profile that belongs to a DIFFERENT DID, rather than refusing. + * + * The refusal it turns off is the whole point of the guard: `sessions.json` is keyed by profile + * name, so a second instance sharing a data directory and reusing a profile name silently + * replaces the first one's identity — and the other daemon then runs under the wrong DID with no + * error. Re-authenticating the SAME DID is untouched; that is the ordinary session refresh. + */ + replaceIdentity?: boolean } = {}, ): Promise { const store = options.store ?? new FileSessionStore() @@ -165,6 +174,17 @@ export async function initializeAgent( if (input.identifier.startsWith('did:') && input.identifier !== session.did) { throw new Error(`Authenticated DID ${session.did} does not match ${input.identifier}`) } + // Before ANY record is written: a stored profile that belongs to another DID is an instance about + // to be clobbered, and the only symptom afterwards is the other daemon running as somebody else. + const registered = await store.get(input.profile) + if (registered && registered.did !== session.did && !options.replaceIdentity) { + throw new Error( + `Profile "${input.profile}" in ${store.path} is already registered to ${registered.did}; ` + + `this init authenticated as ${session.did}. Give this instance its own data directory ` + + '(--data-dir, or "dataDir" in radial.json), or pass --replace-identity to overwrite it.', + ) + } + let pending = session const client = new CredentialClient(session, options.fetcher, async (rotated) => { pending = rotated diff --git a/packages/daemon/src/instance.ts b/packages/daemon/src/instance.ts new file mode 100644 index 0000000..dfb9400 --- /dev/null +++ b/packages/daemon/src/instance.ts @@ -0,0 +1,153 @@ +// One daemon instance: a config file, the data directory that holds its sessions, and the identity +// that pair resolves to. This is the file that makes "several radiald on one machine" a thing an +// operator can spell rather than a thing they have to remember. +// +// `config.ts` parses; this resolves. The separation matters because the data directory has a +// PRECEDENCE (flag, environment, config file, default) and one of its rungs is resolved relative to +// the config file it came from — neither of which is a property of the config's shape. + +import { createHash } from 'node:crypto' +import { dirname, isAbsolute, join, resolve } from 'node:path' +import { FileSessionStore, defaultDataDirectory } from '@radial/atproto/node' +import { + CONFIG_FILENAME, + defaultConfigPath, + findConfigPath, + loadConfig, + type Environment, + type RadialConfig, +} from './config.js' + +/** Which rung of the precedence decided the data directory — printed at startup, because the whole + * failure this exists to prevent is an operator who believes a different one won. */ +export type DataDirSource = 'flag' | 'env' | 'config' | 'default' + +export interface ResolvedDataDir { + path: string + source: DataDirSource + /** Set when `$RADIAL_DATA_DIR` overrode a DIFFERENT `dataDir` in the config file: the stale export + * in a shell that quietly wins over the file is the shadowing an operator cannot see otherwise. */ + shadowed?: string +} + +/** Thrown when config discovery finds nothing. A named class so `cli.ts` can append its usage text + * to this one case without pattern-matching on a message. */ +export class ConfigNotFoundError extends Error {} + +/** + * Where this instance's state lives, and why. + * + * Precedence: `--data-dir` → `$RADIAL_DATA_DIR` → `dataDir` in the config → the operator default. + * The environment beats the file for the same reason `$RADIAL_CONFIG` does — an explicit invocation + * beats a file — but it is reported (`shadowed`) rather than silent. + * + * A `dataDir` in the config is resolved against the DIRECTORY OF THE CONFIG, not the working + * directory. That is the load-bearing choice: a cwd-relative `dataDir` would mean the same command + * run from two directories runs under two identities, which is the failure being fixed. `--data-dir` + * is resolved against the cwd, because that is what a flag typed in a shell means. + */ +export function resolveDataDir(input: { + flag?: string | undefined + env?: Environment + config?: RadialConfig | undefined + configPath?: string | undefined + cwd?: string +}): ResolvedDataDir { + const env = input.env ?? process.env + const cwd = input.cwd ?? process.cwd() + const fromConfig = input.config?.dataDir + const configRelative = (value: string): string => + isAbsolute(value) ? value : resolve(input.configPath ? dirname(input.configPath) : cwd, value) + + if (input.flag) return { path: resolve(cwd, input.flag), source: 'flag' } + if (env.RADIAL_DATA_DIR) { + const path = resolve(cwd, env.RADIAL_DATA_DIR) + const declared = fromConfig ? configRelative(fromConfig) : undefined + return { + path, + source: 'env', + ...(declared && declared !== path ? { shadowed: declared } : {}), + } + } + if (fromConfig) return { path: configRelative(fromConfig), source: 'config' } + let fallback: string + try { + fallback = defaultDataDirectory(env) + } catch { + // Same no-HOME escape hatch the run state directory has always had: a project-local directory + // beats refusing to run at all in a minimal environment. + fallback = join(cwd, '.radial-state') + } + return { path: fallback, source: 'default' } +} + +/** `radiald --data-dir ~/radial/b-state` in the operator's words, for a startup line. */ +export function describeDataDirSource(dataDir: ResolvedDataDir): string { + switch (dataDir.source) { + case 'flag': + return 'from --data-dir' + case 'env': + return 'from $RADIAL_DATA_DIR' + case 'config': + return 'from "dataDir" in the config' + case 'default': + return 'the default' + } +} + +/** + * A short, stable id for this instance, derived from its run state directory. + * + * It exists because docker's namespace is machine-global while everything else here is per data + * directory: two instances serving one space would otherwise derive the SAME container name from the + * same request, colliding on `docker run --name` and killing each other's live containers during + * startup orphan reconciliation. Derived rather than configured so an operator gains one without + * editing anything, and printed at startup so it is identifiable in `docker ps`. + */ +export function instanceId(stateDir: string): string { + return createHash('sha256').update(resolve(stateDir)).digest('hex').slice(0, 6) +} + +export interface Instance { + configPath: string + config: RadialConfig + dataDir: ResolvedDataDir + /** Sessions for THIS instance: `/sessions.json`, never the process-wide default. */ + sessions: FileSessionStore +} + +/** + * The one call every command uses: find the config, resolve the data directory, and open the session + * store that pair implies. `radiald run --config ~/radial/b.json` fully describes an instance when + * that config carries a `dataDir` — which is the "easy" this whole change is for. + */ +export async function loadInstance(args: { + config?: string | undefined + dataDir?: string | undefined + env?: Environment + cwd?: string +}): Promise { + const env = args.env ?? process.env + const cwd = args.cwd ?? process.cwd() + const configPath = await findConfigPath(args.config, cwd, env) + if (!configPath) { + let fallback: string + try { + fallback = defaultConfigPath(env) + } catch { + fallback = `$HOME/.config/radial/${CONFIG_FILENAME}` + } + throw new ConfigNotFoundError( + `No ${CONFIG_FILENAME} found in this directory or any parent, nor at ${fallback}`, + ) + } + const config = await loadConfig(configPath) + const dataDir = resolveDataDir({ + ...(args.dataDir !== undefined ? { flag: args.dataDir } : {}), + env, + config, + configPath, + cwd, + }) + return { configPath, config, dataDir, sessions: new FileSessionStore(dataDir.path) } +} diff --git a/packages/daemon/src/state-lock.ts b/packages/daemon/src/state-lock.ts new file mode 100644 index 0000000..0df9061 --- /dev/null +++ b/packages/daemon/src/state-lock.ts @@ -0,0 +1,60 @@ +import { DatabaseSync } from 'node:sqlite' +import { join } from 'node:path' + +/** + * One `radiald run` per run state directory, enforced. + * + * Separate data directories are only a convention until something checks them, and two daemons on + * one `stateDir` are two writers on `records.db`, `sync.db` and three ledgers — with no complaint + * from either. This is the same cross-process advisory lock `FileSessionStore` already leans on: an + * open `BEGIN IMMEDIATE` on a SQLite file, held for the process lifetime. It has the property a + * pidfile does not — the OS releases it when the holder crashes, so a killed daemon does not leave a + * stale lock its replacement has to be told to ignore. + */ +export class StateDirLock { + readonly path: string + #database: DatabaseSync | undefined + + private constructor(path: string, database: DatabaseSync) { + this.path = path + this.#database = database + } + + /** Acquires the lock for `stateDir`, or throws the operator-facing "already in use" error. The + * directory must already exist (`radiald run` creates it before locking). */ + static acquire(stateDir: string): StateDirLock { + const path = join(stateDir, 'run.lock.db') + const database = new DatabaseSync(path) + try { + database.exec(` + PRAGMA busy_timeout = 1000; + CREATE TABLE IF NOT EXISTS run_lock (id INTEGER PRIMARY KEY) STRICT; + BEGIN IMMEDIATE; + `) + } catch (error) { + database.close() + throw new Error( + `${stateDir} is already in use by another radiald. Give this instance its own data ` + + 'directory (--data-dir, or "dataDir" in radial.json), or stop the other daemon.' + + (error instanceof Error && !/SQLITE_BUSY|database is locked/i.test(error.message) + ? ` (${error.message})` + : ''), + ) + } + return new StateDirLock(path, database) + } + + /** Releases it. Idempotent, so a `finally` may call it after an early failure. */ + release(): void { + const database = this.#database + if (!database) return + this.#database = undefined + try { + database.exec('COMMIT') + } catch { + // Best effort: closing the handle releases the lock either way, and shutdown must not hang. + } finally { + database.close() + } + } +} diff --git a/packages/daemon/src/turn.ts b/packages/daemon/src/turn.ts index 04ae570..81e4fba 100644 --- a/packages/daemon/src/turn.ts +++ b/packages/daemon/src/turn.ts @@ -159,6 +159,9 @@ export interface TurnInput { * binds the turn socket on 0.0.0.0 instead and points the container at * `host.docker.internal`. */ turnTransport?: 'unix' | 'tcp' + /** This daemon instance's id (`instanceId(stateDir)`), mixed into the container label/name so two + * daemons on one machine cannot collide on a name — or reconcile away each other's containers. */ + instance?: string } /** Resolves a profile's declared harness name to the harness that runs it. `selectHarness` @@ -331,13 +334,21 @@ async function saveTurnError(input: TurnInput, error: unknown, secrets: string[] ) } -/** `radial.turn.`, where `` is the deterministic plan-artifact rkey with its `plan-` - * prefix stripped: dot-separated and `=`-free, so it doubles as a valid `docker --label`/`--name` - * and is killable by that same string. (A deliberate, minor deviation from the design's literal - * `radial.turn=` — `=` is illegal in a docker `--name`, and `dockerRunArgs` reuses the label - * as the container name.) */ -export function turnContainerLabel(requestUri: string, requestCid: string): string { - return `radial.turn.${planArtifactRkey(requestUri, requestCid).replace(/^plan-/, '')}` +/** `radial.turn..`, where `` is the deterministic plan-artifact rkey with its + * `plan-` prefix stripped: dot-separated and `=`-free, so it doubles as a valid + * `docker --label`/`--name` and is killable by that same string. (A deliberate, minor deviation from + * the design's literal `radial.turn=` — `=` is illegal in a docker `--name`, and + * `dockerRunArgs` reuses the label as the container name.) + * + * `instance` (`instanceId(stateDir)`) is what keeps two daemons on one machine out of each other's + * way: docker's namespace is machine-global, so two instances serving the same space would otherwise + * derive the same name from the same request — colliding on `docker run --name`, and killing each + * other's live containers during startup orphan reconciliation. Omitted, the label is the + * pre-instance one, which is what every caller with no state directory to hash (tests, one-off + * tooling) gets. */ +export function turnContainerLabel(requestUri: string, requestCid: string, instance?: string): string { + const hash = planArtifactRkey(requestUri, requestCid).replace(/^plan-/, '') + return `radial.turn.${instance ? `${instance}.` : ''}${hash}` } /** @@ -349,7 +360,7 @@ export function turnContainerLabel(requestUri: string, requestCid: string): stri * dispatcher treats as a crash). */ export async function runTurn(input: TurnInput, deps: TurnDeps): Promise { - const label = turnContainerLabel(input.request.uri, input.request.cid) + const label = turnContainerLabel(input.request.uri, input.request.cid, input.instance) const token = deps.makeToken?.() ?? crypto.randomUUID() const isImpl = input.artifactType.name === IMPLEMENTATION_TYPE const isReview = input.artifactType.name === REVIEW_TYPE_NAME diff --git a/packages/daemon/test/actors.test.mjs b/packages/daemon/test/actors.test.mjs index ff5c527..52152f1 100644 --- a/packages/daemon/test/actors.test.mjs +++ b/packages/daemon/test/actors.test.mjs @@ -11,9 +11,10 @@ const session = (profile, did) => ({ profile, }) -/** Minimal FileSessionStore stand-in: loadActors only reads `get(profile)`. */ -function storeWith(sessions) { +/** Minimal FileSessionStore stand-in: loadActors reads `get(profile)` and names `path` in advice. */ +function storeWith(sessions, path = '/state/b/sessions.json') { return { + path, get: async (profile) => sessions[profile], coordinateRefresh: async (_profile, current) => current, } @@ -79,6 +80,104 @@ it('an initialized profile with no harness anywhere is a named error, never a si ) }) +// ── the identity a profile is actually running as ──────────────────────────────────────────────── +// A daemon pointed at the wrong `sessions.json` used to load whatever identity happened to be in it +// and say nothing. A configured DID is checkable, so a mismatch is fatal; a handle is not (it can be +// renamed after init), so that one only warns. + +it('a configured DID that does not match the stored session is fatal, naming both and both files', async () => { + const config = { + identifier: 'did:plc:aaa', + harness: 'claude', + agents: { planner: { artifactTypes: ['plan'] } }, + } + await assert.rejects( + loadActors(config, storeWith({ planner: session('planner', 'did:plc:bbb') }), { + configPath: '/home/op/radial/b.json', + }), + (error) => { + assert.match(error.message, /did:plc:aaa/) + assert.match(error.message, /did:plc:bbb/) + assert.match(error.message, /\/home\/op\/radial\/b\.json/) + assert.match(error.message, /\/state\/b\/sessions\.json/) + assert.match(error.message, /--data-dir/) + return true + }, + ) +}) + +it('a per-profile DID override is checked against that profile, not the config default', async () => { + const config = { + identifier: 'did:plc:aaa', + harness: 'claude', + agents: { + planner: { artifactTypes: ['plan'] }, + reviewer: { artifactTypes: ['review'], identifier: 'did:plc:ccc' }, + }, + } + const sessions = { + planner: session('planner', 'did:plc:aaa'), + reviewer: session('reviewer', 'did:plc:ccc'), + } + const actors = await loadActors(config, storeWith(sessions)) + assert.deepEqual(actors.all.map((actor) => actor.did), ['did:plc:aaa', 'did:plc:ccc']) + + sessions.reviewer = session('reviewer', 'did:plc:aaa') + await assert.rejects(loadActors(config, storeWith(sessions)), /Profile "reviewer" is configured as did:plc:ccc/) +}) + +it('a renamed handle warns and still loads; a matching one says nothing', async () => { + const config = { + identifier: 'old-name.example', + harness: 'claude', + agents: { planner: { artifactTypes: ['plan'] } }, + } + const warnings = [] + const actors = await loadActors(config, storeWith({ planner: session('planner', 'did:plc:aaa') }), { + warn: (message) => warnings.push(message), + }) + // Advisory, never fatal: a handle can legitimately be renamed after `radiald init`. + assert.equal(actors.all.length, 1) + assert.equal(warnings.length, 1) + assert.match(warnings[0], /configured as old-name.example but its stored session is planner.test/) + + const quiet = [] + await loadActors( + { ...config, identifier: 'PLANNER.test' }, + storeWith({ planner: session('planner', 'did:plc:aaa') }), + { warn: (message) => quiet.push(message) }, + ) + assert.deepEqual(quiet, []) +}) + +it('uninitialized profiles are named once, with the file looked in — and never when none loaded', async () => { + const config = { + identifier: 'planner.test', + harness: 'claude', + agents: { planner: { artifactTypes: ['plan'] }, reviewer: { artifactTypes: ['review'] } }, + } + const warnings = [] + const actors = await loadActors(config, storeWith({ planner: session('planner', 'did:plc:aaa') }), { + warn: (message) => warnings.push(message), + configPath: '/home/op/radial/b.json', + }) + assert.deepEqual(actors.all.map((actor) => actor.profile), ['planner']) + assert.equal(warnings.length, 1) + assert.match(warnings[0], /profile reviewer is configured in \/home\/op\/radial\/b\.json/) + assert.match(warnings[0], /no session in \/state\/b\/sessions\.json/) + assert.match(warnings[0], /radiald init reviewer --config \/home\/op\/radial\/b\.json/) + + // Nothing initialized at all is a different failure — an uninitialized or misdirected instance — + // and `radiald run` has a better error for it than one line per configured profile. + const none = [] + const empty = await loadActors(config, storeWith({}), { warn: (message) => none.push(message) }) + assert.equal(empty.all.length, 0) + assert.deepEqual(none, []) + + // A caller that wants no advisories (tests, tooling) passes no `warn` and gets the old silence. + assert.equal((await loadActors(config, storeWith({ planner: session('planner', 'did:plc:aaa') }))).all.length, 1) +}) + it('a config with no models at all loads actors with an empty model list', async () => { const config = { identifier: 'agents.test', diff --git a/packages/daemon/test/check-dispatch.test.mjs b/packages/daemon/test/check-dispatch.test.mjs index 3297770..c8949e9 100644 --- a/packages/daemon/test/check-dispatch.test.mjs +++ b/packages/daemon/test/check-dispatch.test.mjs @@ -1,7 +1,7 @@ import assert from 'node:assert/strict' import { it } from 'node:test' import { COLLECTIONS, MemoryRecordStore, materialize } from '../../core/dist/index.js' -import { CheckDispatcher, CheckLedger, selectCheckable } from '../dist/index.js' +import { CheckDispatcher, CheckLedger, checkContainerLabel, selectCheckable } from '../dist/index.js' const ROOT = 'did:plc:root' const HUMAN = 'did:plc:human' @@ -330,7 +330,7 @@ it('skips when there is no loaded actor to author a checkrun', () => { // --- CheckDispatcher ------------------------------------------------------ -function dispatcherFor(records, { concurrency = 1, runCheckRun } = {}) { +function dispatcherFor(records, { concurrency = 1, runCheckRun, instance } = {}) { const ledger = new CheckLedger() const dispatcher = new CheckDispatcher({ ledger, @@ -340,11 +340,34 @@ function dispatcherFor(records, { concurrency = 1, runCheckRun } = {}) { image: 'radial-check:latest', timeoutMs: 60_000, allowedSchemes: ['https'], + ...(instance ? { instance } : {}), runCheckRun, }) return { ledger, dispatcher } } +it('stamps this daemon instance on the check label it tracks and the one the run uses', async () => { + // Same machine-global docker namespace as turns: the ledger row and the run have to agree, and + // two daemons must not derive one name for one artifact+commit. + const { records, artifact } = buildScenario() + const key = { artifactUri: artifact.uri, artifactCid: artifact.cid, commit: SHA } + let captured + const { ledger, dispatcher } = dispatcherFor(records, { + instance: '4f2a91', + runCheckRun: async (input) => { + captured = input + return { outcome: 'done', ref: { uri: 'at://x', cid: 'c' } } + }, + }) + dispatcher.pump(indexOf(records), registryFor([actorFor(AGENT)])) + const expected = checkContainerLabel(key.artifactUri, key.artifactCid, key.commit, '4f2a91') + assert.equal(ledger.get(key).containerLabel, expected) + assert.equal(captured.instance, '4f2a91') + assert.notEqual(expected, checkContainerLabel(key.artifactUri, key.artifactCid, key.commit)) + await dispatcher.drain() + ledger.close() +}) + it('dispatcher marks a done check done in the ledger and never re-dispatches it', async () => { const { records, artifact } = buildScenario() const key = { artifactUri: artifact.uri, artifactCid: artifact.cid, commit: SHA } diff --git a/packages/daemon/test/check-runner.test.mjs b/packages/daemon/test/check-runner.test.mjs index 84a7bf4..186b2cf 100644 --- a/packages/daemon/test/check-runner.test.mjs +++ b/packages/daemon/test/check-runner.test.mjs @@ -21,6 +21,15 @@ const ARTIFACT = { uri: 'at://did:plc:human/com.disnetdev.radial.artifact/a1', c const COMMIT = '1234567890abcdef1234567890abcdef12345678' const GIT_URL = 'https://example.test/radial-ng.git' +// The check-side half of the machine-global docker namespace: same rule as `turnContainerLabel`. +it('a check container label is scoped to the instance that launched it', () => { + const a = checkContainerLabel(ARTIFACT.uri, ARTIFACT.cid, COMMIT, '4f2a91') + const b = checkContainerLabel(ARTIFACT.uri, ARTIFACT.cid, COMMIT, 'aa11bb') + assert.notEqual(a, b) + assert.ok(a.startsWith('radial.check.4f2a91.')) + assert.match(a, /^[a-zA-Z0-9][a-zA-Z0-9_.-]*$/) +}) + async function withTmp(run) { const dir = await mkdtemp(join(tmpdir(), 'radial-checkrun-')) try { diff --git a/packages/daemon/test/cli.test.mjs b/packages/daemon/test/cli.test.mjs index cdae37e..1e9c51d 100644 --- a/packages/daemon/test/cli.test.mjs +++ b/packages/daemon/test/cli.test.mjs @@ -7,6 +7,9 @@ import { COLLECTIONS } from '../../core/dist/index.js' import { collectModelEnv, coverageAdvisories, + describeActors, + describeInstance, + emptyRegistryError, forgeAdvisory, main, modelEnvNames, @@ -16,7 +19,7 @@ import { claimTimingAdvisory, routingForgeState, } from '../dist/cli.js' -import { ForgeRegistry, TurnLedger } from '../dist/index.js' +import { ForgeRegistry, TurnLedger, instanceId, loadInstance } from '../dist/index.js' // `egress.start()` itself is Docker-only and is not exercised here (see docker-smoke.test.mjs, // which is env-gated behind RADIAL_DOCKER_TESTS). What *is* unit-testable without Docker is that @@ -43,6 +46,104 @@ it('resolveRunConfig defaults merge poll timing and clamps invalid values', () = assert.equal(clamped.mergePollIntervalMs, 120_000) assert.equal(clamped.mergePollBackoffMaxMs, 120_000) }) +it('run.stateDir defaults under the resolved data directory, and an explicit one still wins', () => { + // One setting separates everything an instance owns: point `dataDir` somewhere and the ledgers, + // checkpoints and per-turn directories follow it, with no second thing to remember. + const SPACE = 'at://did:plc:human/com.disnetdev.radial.space/space1' + assert.equal(resolveRunConfig({ spaces: [SPACE] }, { dataDir: '/srv/radial/b' }).stateDir, '/srv/radial/b/run') + assert.equal( + resolveRunConfig({ spaces: [SPACE], stateDir: '/var/lib/radial-b' }, { dataDir: '/srv/radial/b' }).stateDir, + '/var/lib/radial-b', + ) + // No data directory supplied: the historic default, so every existing call site means what it did. + assert.ok(resolveRunConfig({ spaces: [SPACE] }).stateDir.endsWith('/run')) +}) + +// The startup banner. Every failure this feature is about is an operator who believed a different +// set of paths, so the wording is asserted rather than left to drift. +const instanceFixture = (dataDir, run) => ({ + configPath: '/home/op/radial/b.json', + config: { identifier: 'agent-b.example', agents: { planner: { artifactTypes: ['plan'] } } }, + dataDir, + sessions: { path: join(dataDir.path, 'sessions.json') }, + ...(run ? { run } : {}), +}) + +it('describeInstance names every resolved path, and which rule chose the data directory', () => { + const lines = describeInstance( + instanceFixture({ path: '/home/op/radial/b-state', source: 'config' }), + resolveRunConfig({ spaces: [] }, { dataDir: '/home/op/radial/b-state' }), + ) + assert.match(lines[0], /config: {3}\/home\/op\/radial\/b\.json/) + assert.match(lines[1], /data: {5}\/home\/op\/radial\/b-state {3}\(from "dataDir" in the config\)/) + assert.match(lines[2], /sessions: \/home\/op\/radial\/b-state\/sessions\.json/) + assert.match(lines[3], /state: {4}\/home\/op\/radial\/b-state\/run/) + // The instance id is printed because it is the segment that appears in `docker ps`. + assert.ok(lines[3].includes(instanceId('/home/op/radial/b-state/run'))) + + // Each source names itself, so "why is it reading THAT file" is answerable from the log alone. + const source = (dataDir) => describeInstance(instanceFixture(dataDir))[1] + assert.match(source({ path: '/x', source: 'flag' }), /\(from --data-dir\)/) + assert.match(source({ path: '/x', source: 'env' }), /\(from \$RADIAL_DATA_DIR\)/) + assert.match(source({ path: '/x', source: 'default' }), /\(the default\)/) +}) + +it('describeInstance warns about a shadowing $RADIAL_DATA_DIR and a relative run.stateDir', () => { + const shadowed = describeInstance( + instanceFixture({ path: '/env/state', source: 'env', shadowed: '/home/op/radial/b-state' }), + ) + assert.match(shadowed.at(-1), /\$RADIAL_DATA_DIR overrode "dataDir": \/home\/op\/radial\/b-state/) + + // A relative `stateDir` stays cwd-relative (changing that would silently relocate the state of + // every operator using the documented `".radial-state"`), so it is warned about instead. + const relative = describeInstance( + instanceFixture({ path: '/home/op/radial/b-state', source: 'config' }), + resolveRunConfig({ spaces: [], stateDir: '.radial-state' }), + ) + assert.match(relative.at(-1), /"run\.stateDir" is relative/) + + // Neither advisory fires for the ordinary case. + const clean = describeInstance( + instanceFixture({ path: '/home/op/radial/b-state', source: 'config' }), + resolveRunConfig({ spaces: [] }, { dataDir: '/home/op/radial/b-state' }), + ) + assert.equal(clean.length, 4) +}) + +it('the per-profile startup line carries the DID each profile is actually running as', () => { + // Identical profile names on two instances is the supported case, so the profile name alone says + // nothing; the DID is the fact a shared sessions.json used to get wrong in silence. + const actor = (profile, did, handle, harness) => ({ + profile, + did, + harness, + session: { handle }, + }) + const lines = describeActors({ + all: [ + actor('planner', 'did:plc:bbb', 'agent-b.example', 'claude'), + actor('implementer', 'did:plc:bbb', 'agent-b.example', 'claude'), + actor('reviewer', 'did:plc:ccc', 'reviewer-b.example', 'pi'), + ], + }) + assert.deepEqual(lines, [ + ' planner did:plc:bbb agent-b.example claude', + ' implementer did:plc:bbb agent-b.example claude', + ' reviewer did:plc:ccc reviewer-b.example pi', + ]) +}) + +it('the empty-registry error names the sessions file it read and what the config expected in it', () => { + // "run radiald init" is no help when the question is *which* sessions file was read. + const message = emptyRegistryError( + instanceFixture({ path: '/home/op/radial/b-state', source: 'config' }), + ) + assert.match(message, /No initialized agent profiles found for \/home\/op\/radial\/b\.json/) + assert.match(message, /in \/home\/op\/radial\/b-state\/sessions\.json/) + assert.match(message, /config declares: planner/) + assert.match(message, /radiald init --config \/home\/op\/radial\/b\.json/) +}) + it('resolveRunConfig defaults turnTransport to "unix"', () => { const resolved = resolveRunConfig({ spaces: ['at://did:plc:human/com.disnetdev.radial.space/space1'] }) assert.equal(resolved.turnTransport, 'unix') @@ -233,6 +334,34 @@ it('names a profile whose config has moved on from the record it published', () ) }) +// `radiald config init --data-dir` — the one-command scaffold for a SECOND instance. What it writes +// has to be enough that `--config` alone selects the instance afterwards. +it('config init writes a dataDir and prints follow-up commands that name the config', async () => { + const dir = await mkdtemp(join(tmpdir(), 'radial-config-init-')) + const lines = [] + const log = console.log + console.log = (line) => lines.push(String(line)) + try { + const configPath = join(dir, 'b.json') + await main(['config', 'init', 'agent-b.example', '--config', configPath, '--data-dir', join(dir, 'b-state')]) + const written = JSON.parse(await readFile(configPath, 'utf8')) + assert.equal(written.dataDir, join(dir, 'b-state')) + assert.equal(written.identifier, 'agent-b.example') + // Both follow-ups carry `--config`: with several instances on a machine, an init that discovers + // "the nearest radial.json" is exactly how the wrong one gets initialized. + assert.ok(lines.some((line) => line.includes(`radiald init --config ${configPath} --password-stdin`))) + assert.ok(lines.some((line) => line.includes(`radiald run --config ${configPath}`))) + + // And what it wrote round-trips into a full instance: the sessions file follows the config. + const instance = await loadInstance({ config: configPath, env: { HOME: dir }, cwd: '/' }) + assert.equal(instance.sessions.path, join(dir, 'b-state', 'sessions.json')) + assert.equal(instance.dataDir.source, 'config') + } finally { + console.log = log + await rm(dir, { recursive: true, force: true }) + } +}) + // `radiald turn list` — the command docs/plan.md sends an operator to before enabling replies. It // reads the same ledger `turn reset` writes, so the two agree about which rows exist. it('turn list prints the ledger and narrows to a state', async () => { diff --git a/packages/daemon/test/config.test.mjs b/packages/daemon/test/config.test.mjs index 688a5d0..c7cf2e6 100644 --- a/packages/daemon/test/config.test.mjs +++ b/packages/daemon/test/config.test.mjs @@ -207,6 +207,32 @@ it('scaffolds a run block with the forge already wired, and preserves it across assert.deepEqual((await loadConfig(again)).run, { spaces: [], forge: { kind: 'github' } }) }) +it('carries a "dataDir" through parse, scaffold, and a write/read round trip', async () => { + // `dataDir` is what lets one `--config` select an instance outright — sessions, state and all — + // so it has to survive anything that rewrites the file, exactly as `run.forge` does. + assert.equal(parseConfig({ dataDir: 'b-state', agents: { p: { artifactTypes: ['plan'] } } }).dataDir, 'b-state') + assert.equal(parseConfig({ agents: { p: { artifactTypes: ['plan'] } } }).dataDir, undefined) + assert.throws( + () => parseConfig({ dataDir: '', agents: { p: { artifactTypes: ['plan'] } } }), + /dataDir must be a non-empty string/, + ) + + const directory = await mkdtemp(join(tmpdir(), 'radial-scaffold-datadir-')) + const path = join(directory, 'radial.json') + await writeConfig(path, buildDefaultConfig('agent.example', { dataDir: '/srv/radial/b-state' })) + const raw = JSON.parse(await readFile(path, 'utf8')) + assert.equal(raw.dataDir, '/srv/radial/b-state') + assert.deepEqual(Object.keys(raw), ['identifier', 'dataDir', 'harness', 'models', 'run', 'agents']) + + const again = join(directory, 'again.json') + await writeConfig(again, await loadConfig(path)) + assert.equal((await loadConfig(again)).dataDir, '/srv/radial/b-state') + + // Absent unless asked for: a scaffold that pinned a data directory would freeze the default into + // every operator's file, which is what `run`'s other settings deliberately avoid. + assert.equal(buildDefaultConfig('agent.example').dataDir, undefined) +}) + it('parses an empty run.spaces (the scaffold) but still demands an array', () => { // A scaffold nothing can load is a scaffold nothing can edit: `radiald init`, `turn reset` and // friends all have to work on a file whose spaces the operator has not filled in yet. Only diff --git a/packages/daemon/test/dispatch.test.mjs b/packages/daemon/test/dispatch.test.mjs index f39050a..73643a5 100644 --- a/packages/daemon/test/dispatch.test.mjs +++ b/packages/daemon/test/dispatch.test.mjs @@ -20,6 +20,7 @@ import { reviewRkey, runTurn, selectDispatchable, + turnContainerLabel, } from '../dist/index.js' import { reconcileOrphan, reconcileOrphans } from '../dist/cli.js' @@ -1931,6 +1932,37 @@ it('dispatcher threads prev into a non-implementation turn and never a branch', ledger.close() }) +it('stamps this daemon instance on the container label it tracks and the one the turn runs under', async () => { + // The ledger row and the turn have to agree about the label, because that string is what shutdown + // and startup orphan reconciliation kill by — and with two daemons on one machine, killing by the + // OTHER instance's label is exactly the accident the instance segment prevents. + const { records, request } = buildScenario() + const index = materialize(store(records), { spaceUri: SPACE_URI }) + const actors = registryFor([actorFor(AGENT, ['plan'])]) + const ledger = new TurnLedger() + let captured + const dispatcher = new TurnDispatcher({ + ledger, + concurrency: 1, + runDirFor: () => '/tmp/radial-instance-label', + image: 'radial-turn:test', + timeoutMs: 30_000, + allowedSchemes: ['https'], + instance: '4f2a91', + runTurn: async (input) => { + captured = input + return { outcome: 'fulfilled', acceptedRef: { uri: 'at://x', cid: 'c' }, label: 'l' } + }, + }) + dispatcher.pump(index, actors) + const expected = turnContainerLabel(request.uri, request.cid, '4f2a91') + assert.equal(ledger.get(request.uri).containerLabel, expected) + assert.equal(captured.instance, '4f2a91') + assert.notEqual(expected, turnContainerLabel(request.uri, request.cid)) + await dispatcher.drain() + ledger.close() +}) + it('rejects an ambiguous project-scoped request by anchoring the explanation to a candidate', async () => { const base = buildScenario() const a = priorArtifact({ project: ref(base.project) }, 'architecture', 'arch-a') diff --git a/packages/daemon/test/init.test.mjs b/packages/daemon/test/init.test.mjs index f289b77..45fe83a 100644 --- a/packages/daemon/test/init.test.mjs +++ b/packages/daemon/test/init.test.mjs @@ -254,6 +254,90 @@ it('writes nothing when the published record already says what the config says', assert.deepEqual(writes, []) }) +// ── one profile name, two identities ──────────────────────────────────────────────────────────── +// `sessions.json` is keyed by profile name, so a second instance sharing a data directory and +// reusing a profile name replaces the first one's identity — and the other daemon then runs as +// somebody else with no error at all. That is the failure separate data directories exist to +// prevent, and this is the guard that makes getting it wrong loud instead of silent. + +/** A PDS stub for `did`, with no agent record published yet. */ +const identityPds = (did, handle) => async (url, init = {}) => { + const target = String(url) + if (target.endsWith('createSession')) { + return Response.json({ did, handle, accessJwt: 'access', refreshJwt: 'refresh' }) + } + if (target.includes('getRecord')) { + return Response.json({ error: 'RecordNotFound', message: 'missing' }, { status: 400 }) + } + const body = JSON.parse(init.body) + return Response.json({ uri: `at://${did}/${body.collection}/${body.rkey}`, cid: `cid-${did}` }) +} + +const initAs = (identifier) => ({ + profile: 'planner', + service: 'https://pds.test', + identifier, + password: 'app-password', + handleName: 'planner', + harness: 'claude', + models: [], + artifactTypes: ['plan'], + now: '2026-07-26T00:00:00Z', +}) + +it('refuses to re-register a profile under a different DID, and leaves the stored session alone', async () => { + const directory = await mkdtemp(join(tmpdir(), 'radial-clobber-')) + const store = new FileSessionStore(directory) + await initializeAgent(initAs('a.example'), { store, fetcher: identityPds('did:plc:aaa', 'a.example') }) + + await assert.rejects( + initializeAgent(initAs('b.example'), { store, fetcher: identityPds('did:plc:bbb', 'b.example') }), + (error) => { + // Both DIDs, the file, and the two ways out — an operator hitting this cannot see any of them. + assert.match(error.message, /did:plc:aaa/) + assert.match(error.message, /did:plc:bbb/) + assert.ok(error.message.includes(store.path)) + assert.match(error.message, /--data-dir/) + assert.match(error.message, /--replace-identity/) + return true + }, + ) + assert.equal((await store.get('planner')).did, 'did:plc:aaa') + await rm(directory, { recursive: true, force: true }) +}) + +it('--replace-identity overwrites it, and re-authenticating the same DID never needed a flag', async () => { + const directory = await mkdtemp(join(tmpdir(), 'radial-replace-')) + const store = new FileSessionStore(directory) + await initializeAgent(initAs('a.example'), { store, fetcher: identityPds('did:plc:aaa', 'a.example') }) + + // The ordinary session refresh: same identity, same profile, no flag, no complaint. + await initializeAgent(initAs('a.example'), { store, fetcher: identityPds('did:plc:aaa', 'a.example') }) + assert.equal((await store.get('planner')).did, 'did:plc:aaa') + + await initializeAgent(initAs('b.example'), { + store, + fetcher: identityPds('did:plc:bbb', 'b.example'), + replaceIdentity: true, + }) + assert.equal((await store.get('planner')).did, 'did:plc:bbb') + await rm(directory, { recursive: true, force: true }) +}) + +it('two data directories hold the same profile name under two identities — the whole point', async () => { + const a = await mkdtemp(join(tmpdir(), 'radial-instance-a-')) + const b = await mkdtemp(join(tmpdir(), 'radial-instance-b-')) + const storeA = new FileSessionStore(a) + const storeB = new FileSessionStore(b) + await initializeAgent(initAs('a.example'), { store: storeA, fetcher: identityPds('did:plc:aaa', 'a.example') }) + await initializeAgent(initAs('b.example'), { store: storeB, fetcher: identityPds('did:plc:bbb', 'b.example') }) + assert.equal((await storeA.get('planner')).did, 'did:plc:aaa') + assert.equal((await storeB.get('planner')).did, 'did:plc:bbb') + assert.notEqual(storeA.path, storeB.path) + await rm(a, { recursive: true, force: true }) + await rm(b, { recursive: true, force: true }) +}) + it('describes drift field by field, and says nothing about fields that match', () => { const configured = { ...published, diff --git a/packages/daemon/test/instance.test.mjs b/packages/daemon/test/instance.test.mjs new file mode 100644 index 0000000..e2cf654 --- /dev/null +++ b/packages/daemon/test/instance.test.mjs @@ -0,0 +1,105 @@ +// Data-directory resolution: the one thing that decides WHICH instance a `radiald` invocation is. +// Every assertion here is about a way an operator can end up running the wrong one. + +import assert from 'node:assert/strict' +import { mkdtemp, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join, resolve } from 'node:path' +import { it } from 'node:test' +import { describeDataDirSource, instanceId, loadInstance, resolveDataDir } from '../dist/index.js' + +const CONFIG = { + identifier: 'agent.example', + harness: 'claude', + agents: { planner: { artifactTypes: ['plan'] } }, +} + +it('resolves the data directory by precedence: flag, env, config, default', () => { + const env = { HOME: '/home/op' } + const configPath = '/home/op/radial/b.json' + const config = { ...CONFIG, dataDir: 'b-state' } + + assert.deepEqual( + resolveDataDir({ flag: '/flag/state', env, config, configPath, cwd: '/cwd' }), + { path: '/flag/state', source: 'flag' }, + ) + // The flag beats the environment, which beats the file — and the shadowed file value is REPORTED + // rather than silently dropped, because a stale export in a shell is invisible otherwise. + assert.deepEqual( + resolveDataDir({ env: { ...env, RADIAL_DATA_DIR: '/env/state' }, config, configPath, cwd: '/cwd' }), + { path: '/env/state', source: 'env', shadowed: '/home/op/radial/b-state' }, + ) + assert.deepEqual(resolveDataDir({ env, config, configPath, cwd: '/cwd' }), { + path: '/home/op/radial/b-state', + source: 'config', + }) + assert.deepEqual(resolveDataDir({ env, config: CONFIG, configPath, cwd: '/cwd' }), { + path: '/home/op/.local/state/radial', + source: 'default', + }) +}) + +it('an env var that agrees with the config is not reported as shadowing it', () => { + const resolved = resolveDataDir({ + env: { HOME: '/home/op', RADIAL_DATA_DIR: '/home/op/radial/b-state' }, + config: { ...CONFIG, dataDir: 'b-state' }, + configPath: '/home/op/radial/b.json', + cwd: '/cwd', + }) + assert.equal(resolved.source, 'env') + assert.equal(resolved.shadowed, undefined) +}) + +it('a config "dataDir" is relative to the CONFIG, and a --data-dir flag is relative to the cwd', () => { + // The load-bearing asymmetry: the same `radiald run --config …` from two working directories must + // be the same instance, or the whole feature is another way to get the wrong identity. A flag + // typed in a shell means what a shell means by it. + const config = { ...CONFIG, dataDir: './b-state' } + const configPath = '/home/op/radial/b.json' + const fromHome = resolveDataDir({ env: { HOME: '/home/op' }, config, configPath, cwd: '/home/op' }) + const fromElsewhere = resolveDataDir({ env: { HOME: '/home/op' }, config, configPath, cwd: '/tmp/somewhere' }) + assert.equal(fromHome.path, '/home/op/radial/b-state') + assert.deepEqual(fromHome, fromElsewhere) + + assert.equal( + resolveDataDir({ flag: 'b-state', env: { HOME: '/home/op' }, cwd: '/tmp/somewhere' }).path, + '/tmp/somewhere/b-state', + ) + // An absolute `dataDir` is left alone. + assert.equal( + resolveDataDir({ env: { HOME: '/home/op' }, config: { ...CONFIG, dataDir: '/srv/b' }, configPath }).path, + '/srv/b', + ) +}) + +it('names the rule that chose the directory, for the startup banner', () => { + const say = (source) => describeDataDirSource({ path: '/x', source }) + assert.equal(say('flag'), 'from --data-dir') + assert.equal(say('env'), 'from $RADIAL_DATA_DIR') + assert.equal(say('config'), 'from "dataDir" in the config') + assert.equal(say('default'), 'the default') +}) + +it('loadInstance opens the sessions file the config points at, without any environment discipline', async () => { + const dir = await mkdtemp(join(tmpdir(), 'radial-instance-')) + try { + const configPath = join(dir, 'b.json') + await writeFile(configPath, JSON.stringify({ ...CONFIG, dataDir: 'b-state' })) + // Deliberately NO $RADIAL_DATA_DIR: `--config` alone has to fully describe the instance. + const instance = await loadInstance({ config: configPath, env: { HOME: dir }, cwd: '/' }) + assert.equal(instance.configPath, configPath) + assert.equal(instance.dataDir.source, 'config') + assert.equal(instance.dataDir.path, join(dir, 'b-state')) + assert.equal(instance.sessions.path, join(dir, 'b-state', 'sessions.json')) + } finally { + await rm(dir, { recursive: true, force: true }) + } +}) + +it('instanceId is stable per state directory, and differs between two of them', () => { + assert.equal(instanceId('/srv/a/run'), instanceId('/srv/a/run/')) + assert.equal(instanceId('/srv/a/run'), instanceId(resolve('/srv/a', 'run'))) + assert.notEqual(instanceId('/srv/a/run'), instanceId('/srv/b/run')) + // Short enough to read in `docker ps`, and legal in a docker container name. + assert.match(instanceId('/srv/a/run'), /^[0-9a-f]{6}$/) +}) diff --git a/packages/daemon/test/state-lock.test.mjs b/packages/daemon/test/state-lock.test.mjs new file mode 100644 index 0000000..d76e008 --- /dev/null +++ b/packages/daemon/test/state-lock.test.mjs @@ -0,0 +1,64 @@ +// The state-dir lock: what turns "give each instance its own data directory" from a convention into +// something the daemon enforces. Two `radiald run` on one state dir are two writers on `records.db`, +// `sync.db` and three ledgers, and nothing downstream notices. +// +// Cross-process is covered by the same SQLite semantics `FileSessionStore` already relies on; two +// connections in one process take the same lock, which is what these drive. + +import assert from 'node:assert/strict' +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { it } from 'node:test' +import { StateDirLock } from '../dist/index.js' + +it('a second daemon on one state dir is refused, by name, and told what to do about it', async () => { + const dir = await mkdtemp(join(tmpdir(), 'radial-state-lock-')) + try { + const held = StateDirLock.acquire(dir) + try { + assert.throws( + () => StateDirLock.acquire(dir), + (error) => { + assert.match(error.message, /already in use by another radiald/) + // The fix, in the message: an operator hitting this has two daemons and one directory. + assert.match(error.message, /--data-dir/) + assert.ok(error.message.includes(dir)) + return true + }, + ) + } finally { + held.release() + } + // Released, so the replacement daemon starts — the ordinary restart path. + StateDirLock.acquire(dir).release() + } finally { + await rm(dir, { recursive: true, force: true }) + } +}) + +it('two state dirs do not contend, which is the whole point of separating them', async () => { + const a = await mkdtemp(join(tmpdir(), 'radial-state-lock-a-')) + const b = await mkdtemp(join(tmpdir(), 'radial-state-lock-b-')) + try { + const lockA = StateDirLock.acquire(a) + const lockB = StateDirLock.acquire(b) + lockA.release() + lockB.release() + } finally { + await rm(a, { recursive: true, force: true }) + await rm(b, { recursive: true, force: true }) + } +}) + +it('release is idempotent, so a failed startup may release in a finally', async () => { + const dir = await mkdtemp(join(tmpdir(), 'radial-state-lock-idem-')) + try { + const lock = StateDirLock.acquire(dir) + lock.release() + lock.release() + StateDirLock.acquire(dir).release() + } finally { + await rm(dir, { recursive: true, force: true }) + } +}) diff --git a/packages/daemon/test/turn.test.mjs b/packages/daemon/test/turn.test.mjs index ad543d4..ef028b8 100644 --- a/packages/daemon/test/turn.test.mjs +++ b/packages/daemon/test/turn.test.mjs @@ -20,8 +20,25 @@ import { reviewRkey, runTurn, selectHarness, + turnContainerLabel, } from '../dist/index.js' +// Docker's namespace is machine-global while everything else an instance owns is per data +// directory, so the label carries the instance: two daemons serving one space would otherwise +// derive the same `--name` from the same request, collide on `docker run`, and reconcile away each +// other's live containers at startup. +it('a turn container label is scoped to the instance that launched it', () => { + const request = { uri: 'at://did:plc:human/com.disnetdev.radial.artifactRequest/r1', cid: 'cid-r1' } + const bare = turnContainerLabel(request.uri, request.cid) + const a = turnContainerLabel(request.uri, request.cid, '4f2a91') + const b = turnContainerLabel(request.uri, request.cid, 'aa11bb') + assert.notEqual(a, b) + assert.ok(a.startsWith('radial.turn.4f2a91.')) + assert.ok(a.endsWith(bare.slice('radial.turn.'.length))) + // Still a legal docker container name (`dockerRunArgs` reuses the label as `--name`). + assert.match(a, /^[a-zA-Z0-9][a-zA-Z0-9_.-]*$/) +}) + /** The real registry: `runTurn` resolves each profile's declared harness through it, fail-closed. */ const HARNESSES_FOR_TEST = { select: selectHarness } diff --git a/packages/daemon/test/two-instances.test.mjs b/packages/daemon/test/two-instances.test.mjs new file mode 100644 index 0000000..d714f41 --- /dev/null +++ b/packages/daemon/test/two-instances.test.mjs @@ -0,0 +1,162 @@ +// The two-instance drill: two `radiald` on one machine, two agent identities, THE SAME profile +// names — the arrangement that used to require invisible `$RADIAL_CONFIG` / `$RADIAL_DATA_DIR` +// discipline on every invocation and, when it was got wrong, ran a daemon under the wrong DID +// without a word. +// +// In the spirit of multi-operator.test.mjs: entirely offline, no Docker, no wall-clock waits. Two +// configs in two directories, two `LocalPds`, and every assertion is a claim the goal makes. + +import assert from 'node:assert/strict' +import { mkdir, mkdtemp, readFile, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { it } from 'node:test' +import { LocalPds } from '../../atproto/test/local-pds.mjs' +import { + StateDirLock, + initializeAgent, + instanceId, + loadActors, + loadInstance, + resolveAgentInits, + turnContainerLabel, +} from '../dist/index.js' +import { resolveRunConfig } from '../dist/cli.js' + +const REQUEST = { + uri: 'at://did:plc:human/com.disnetdev.radial.artifactRequest/shared', + cid: 'cid-shared-request', +} +const SPACE = 'at://did:plc:human/com.disnetdev.radial.space/one' + +/** One instance on disk: `/.json` naming `/-state` as its data directory. */ +async function scaffold(root, name, pds) { + const configPath = join(root, `${name}.json`) + await writeFile( + configPath, + JSON.stringify({ + identifier: pds.handle, + // Config-relative, so the same command from any working directory is the same instance. + dataDir: `${name}-state`, + pds: pds.service, + harness: 'claude', + // Deliberately IDENTICAL profile names on both sides: distinct names must not be a + // requirement, because a profile name is a record key in the agent's own repo. + agents: { + planner: { artifactTypes: ['plan'] }, + implementer: { artifactTypes: ['implementation'] }, + }, + run: { spaces: [SPACE] }, + }), + ) + // No environment at all beyond HOME: `--config` alone has to select the instance. + const instance = await loadInstance({ config: configPath, env: { HOME: root }, cwd: '/' }) + const run = resolveRunConfig(instance.config.run, { dataDir: instance.dataDir.path }) + return { instance, run, pds } +} + +/** What `radiald init --config ` does, with the network replaced by a LocalPds. */ +async function init(instance, pds, options = {}) { + const inits = resolveAgentInits(instance.config, [], { default: 'app-password' }) + const results = [] + for (const input of inits) { + results.push( + await initializeAgent(input, { + store: instance.sessions, + fetcher: (url, init) => pds.fetch(url, init), + ...options, + }), + ) + } + return results +} + +it('two instances on one machine: two identities, the same profile names, nothing shared', async () => { + const root = await mkdtemp(join(tmpdir(), 'radial-two-instances-')) + try { + const a = await scaffold(root, 'a', new LocalPds('did:plc:agenta')) + const b = await scaffold(root, 'b', new LocalPds('did:plc:agentb')) + + await init(a.instance, a.pds) + await init(b.instance, b.pds) + + // 1. Two sessions files, each holding its own DID under the same two profile names. Neither + // init rewrote the other's — the failure that made this arrangement unsafe. + assert.notEqual(a.instance.sessions.path, b.instance.sessions.path) + for (const [instance, did] of [ + [a.instance, 'did:plc:agenta'], + [b.instance, 'did:plc:agentb'], + ]) { + const stored = JSON.parse(await readFile(instance.sessions.path, 'utf8')) + assert.deepEqual(Object.keys(stored.profiles).sort(), ['implementer', 'planner']) + assert.deepEqual( + Object.values(stored.profiles).map((profile) => profile.did), + [did, did], + ) + } + + // 2. Each daemon loads the right DID for every profile. + for (const [instance, did] of [ + [a.instance, 'did:plc:agenta'], + [b.instance, 'did:plc:agentb'], + ]) { + const actors = await loadActors(instance.config, instance.sessions, { + configPath: instance.configPath, + }) + assert.deepEqual(actors.all.map((actor) => actor.profile), ['implementer', 'planner']) + assert.deepEqual(actors.all.map((actor) => actor.did), [did, did]) + } + + // 3. State dirs, instance ids, and container labels for ONE request all differ — the last one + // because docker's namespace is machine-global even when everything else is separated. + assert.notEqual(a.run.stateDir, b.run.stateDir) + assert.equal(a.run.stateDir, join(a.instance.dataDir.path, 'run')) + const idA = instanceId(a.run.stateDir) + const idB = instanceId(b.run.stateDir) + assert.notEqual(idA, idB) + assert.notEqual( + turnContainerLabel(REQUEST.uri, REQUEST.cid, idA), + turnContainerLabel(REQUEST.uri, REQUEST.cid, idB), + ) + + // 4. Negative: pointing instance B's init at A's data directory is refused, naming both DIDs and + // the file — instead of silently taking over A's identity. + const crossed = await loadInstance({ + config: b.instance.configPath, + dataDir: a.instance.dataDir.path, + env: { HOME: root }, + cwd: '/', + }) + assert.equal(crossed.sessions.path, a.instance.sessions.path) + await assert.rejects(init(crossed, b.pds), (error) => { + assert.match(error.message, /did:plc:agenta/) + assert.match(error.message, /did:plc:agentb/) + assert.ok(error.message.includes(a.instance.sessions.path)) + return true + }) + // A's stored identity is untouched by the refusal. + assert.equal((await a.instance.sessions.get('planner')).did, 'did:plc:agenta') + + // 5. Negative: a config that names A's DID while reading B's sessions refuses at load, rather + // than running under the wrong identity. + await assert.rejects( + loadActors( + { ...b.instance.config, identifier: 'did:plc:agenta' }, + b.instance.sessions, + { configPath: b.instance.configPath }, + ), + /would run under the wrong identity/, + ) + + // 6. Negative: two daemons on one state dir conflict; two on their own do not. + await mkdir(a.run.stateDir, { recursive: true }) + await mkdir(b.run.stateDir, { recursive: true }) + const lockA = StateDirLock.acquire(a.run.stateDir) + const lockB = StateDirLock.acquire(b.run.stateDir) + assert.throws(() => StateDirLock.acquire(a.run.stateDir), /already in use by another radiald/) + lockA.release() + lockB.release() + } finally { + await rm(root, { recursive: true, force: true }) + } +}) diff --git a/packages/sidecar/src/cli.ts b/packages/sidecar/src/cli.ts index 521f96a..b15ee16 100644 --- a/packages/sidecar/src/cli.ts +++ b/packages/sidecar/src/cli.ts @@ -60,7 +60,9 @@ record: "artifact submit", "review submit", "answer submit" (the reply an answer turn was commissioned for) and "message post" (the can't-complete question). References may be bare at:// URIs or pinned at://...#CID locators. -Use --profile NAME to select an actor and --json for structured output.` +Use --profile NAME to select an actor and --json for structured output. +Use --data-dir PATH (else $RADIAL_DATA_DIR) to read a different instance's +sessions.json, when this machine runs more than one radiald.` async function readStdin(): Promise { let value = '' @@ -89,7 +91,11 @@ export async function main(args = argv.slice(2)): Promise { return } - const store = new FileSessionStore() + // `--data-dir` selects WHICH set of identities this CLI speaks for, the same way it does for + // `radiald`: one machine may hold several instances, each with its own `sessions.json`. Read here + // rather than in `runCli`, which only ever sees flags it knows and ignores the rest. + const dataDir = value(args, '--data-dir') + const store = dataDir ? new FileSessionStore(dataDir) : new FileSessionStore() if (args[0] === 'auth' && args[1] === 'login') { const profile = value(args, '--profile') const identifier = value(args, '--identifier') -- 2.51.2