Skip to content

nodus.agent

View Markdown

Agents: durable runs on the journal engine (resources.md §4.1 to §4.3, ADR-048, ADR-101).

This module is the client side: defining and deploying an Agent, submitting AgentRuns, waiting for results, and fanning out through AgentGroups. Inside a run, the agent runtime (nodus._runtime.agents) executes the entrypoint with a RunContext and journals every @agent.step; outside a run a step is a plain function call.

The API serves an Agent’s image and perRunMaxCostUSD, an AgentRun’s input, deadline and group membership (group, taskKey, dependsOn), a run’s answer, steps, messages and cancel, and an AgentGroup’s agent, limits, maxCostUSD, sealed, state and bounded Environment evaluation (ADR-119 defers the rest). Anything else a call names raises errors.Unsupported before anything is sent, because the API rejects an unknown field outright; nodus.ClaudeAgent defines an agent on the served kind.

class Agent(name: str, *, image: Any = None, source: str | Path | dict[str, Any] | None = None, entrypoint: str | None = None, setup: str | None = None, secrets: list[Any] | None = None, env: dict[str, str] | None = None, network: _spec.Egress | None = None, models: list[str] | None = None, min_workers: int | None = None, max_workers: int | None = None, scaledown_window: Any = None, per_run_max_cost: Any = None, max_cost: Any = None, cpu: Any = None, memory: Any = None, app: Any = None, project: str | None = None) -> None

A durable agent definition; deploy it, then submit runs with .remote(), .spawn() or .map().

deploy() -> View

Create or update the Agent; every accepted change is a new revision, and new runs pin it.

entrypoint(fn: Callable[..., Any]) -> Callable[..., Any]

Mark fn(ctx, input) as the run entrypoint; its module is uploaded as the Agent’s source.

from_name(name: str, project: str | None = None) -> _Agent

A deployed Agent, submitted to without redeploying it.

local(input: Any = None) -> Any

Run the entrypoint in this process with an in-memory journal (needs the agent runtime).

map(inputs: Iterable[Any], *, max_active: int | None = None, order_outputs: bool = True, return_exceptions: bool = False) -> AsyncIterator[Any]

One run per input in an ephemeral AgentGroup; yields each run’s answer text as the run finishes.

Answers come in input order unless order_outputs=False. A run that does not succeed raises AgentRunFailed, or is yielded as that exception with return_exceptions=True.

Type: str

remote(input: Any = None, **kwargs: Any) -> Any

Submit, wait and return the run’s answer text; raises AgentRunFailed.

spawn(input: Any = None, **kwargs: Any) -> _AgentRun

The same as submit.

step(_fn: Callable[..., Any] | None = None, *, effect: str = 'pure', name: str | None = None) -> Any

Journal a function as a step: pure, idempotent or external (an unknown outcome needs resolution).

submit(input: Any = None, *, idempotency_key: str | None = None, session_key: str | None = None, deadline: str | None = None, group: str | None = None, hold_worker: str | None = None, name: str | None = None) -> _AgentRun

Start a run and return its handle. A name derived from an event makes redeliveries idempotent.

group is refused (submit group runs through AgentGroup.submit_many), and so are session_key and a hold_worker other than "auto", which the API does not serve yet.

class AgentGroup(obj: Obj) -> None

Runs of one Agent under a shared concurrency limit and cost cap, with task dependencies and one cancel.

cancel() -> None

Cancel every unfinished run in the group.

create(name: str, agent: _Agent | str, *, max_active: int | None = None, max_pending: int | None = None, max_held: int | None = None, max_cost: Any = None, evaluation: dict[str, Any] | None = None, project: str | None = None) -> _AgentGroup

Create the group; max_active runs go at once and runs are released while max_cost has room for them.

max_cost caps the runs’ Claude usage (model and routing calls) together; their sandboxes are billed apart.

evaluation={"environment": "nodus/arithmetic-v2@2.0.0", "tasks": 10, "seed": 42} generates and scores a fixed batch automatically; split defaults to test, repetitions to 1, and tasks times repetitions is at most 100. timeout defaults to “30m” (1m to 24h) from group creation, including queued time. Evaluation and agent Sandbox compute is billed separately from the model-only max_cost; use a project Budget to cap total spend. max_held remains unsupported.

delete() -> None

Delete the group and, with it, its runs.

from_name(name: str, project: str | None = None) -> _AgentGroup

A group that exists already, such as one created from YAML.

Type: str

results() -> list[View]

Per-case evaluation outcomes and grader evidence; no hidden answers or task payloads.

runs() -> list[_AgentRun]

The group’s member runs.

seal() -> None

Close the group to new runs; it finishes once every run is terminal.

submit_many(tasks: list[dict[str, Any]]) -> list[_AgentRun]

Create one run per task {key, input, depends_on?} and return them in the caller’s order.

A task starts after the tasks in depends_on, which name tasks of this call or of an earlier one. Runs go out in AgentRunList batches (each batch is all-or-nothing); a cycle or a repeated key raises Invalid.

wait(timeout: float | None = None) -> View

Return the terminal status; submitted batches must be sealed, while evaluations close automatically.

class AgentRun(obj: Obj) -> None

One AgentRun: wait(), result() (the answer text), answer(), steps(), send() and cancel().

answer() -> str

The run’s full answer text.

cancel() -> None
children() -> list[_AgentRun]
from_name(name: str, project: str | None = None) -> _AgentRun
logs(follow: bool = False) -> AsyncIterator[str]

Type: str

Type: _Outputs

resolve(step_id: str, decision: str, evidence: dict[str, str] | None = None, result: Any = None, checkpoint_seq: int | None = None, expected_revision: int | None = None) -> View

Resolve an external step with an unknown outcome: Completed, NoEffect or Cancelled.

result() -> str

The run’s answer text (waits for the run first); raises AgentRunFailed unless it succeeded.

resume() -> None
retry() -> None
send(name: str, payload: Any, message_key: str | None = None) -> View
steps() -> list[View]
suspend() -> None
wait(timeout: float | None = None) -> View

Block until the run is terminal; raises AgentRunFailed unless it succeeded.

class RunContext(Protocol)

What an agent entrypoint receives as ctx; the agent runtime provides the implementation.

child_output(child: Any, name: str) -> Any
continue_as_new(input: Any) -> None
gather(children: list[Any], return_exceptions: bool = False) -> list[Any]

Type: str

map(inputs: Iterable[Any]) -> list[Any]
save_output(name: str, path: str) -> None
send(target: str, name: str, payload: Any) -> None
sleep(duration: str | float) -> None
sleep_until(time: str) -> None
spawn(input: Any, key: str, permissions: Any = None, deadline: str | None = None) -> Any

Type: Path

step(name: str, fn: Callable[[], Any], effect: str = 'pure') -> Any
wait_for_message(name: str, timeout: str | float | None = None) -> Any