Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 37 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,43 @@ appear in a patch rather than inflating the version toward 1.0 on a crate still
shape. **Where that happens the entry says so at the top**, because a version number that
under-signals is only acceptable if the changelog over-signals to compensate.

## 0.5.0

The crate now runs on nagoya instead of tokio, and every entry point that spawns a CLI
takes the I/O reactor it should use.

### Breaking

- **Spawning entry points take a `&nagoya::reactor::Handle`.** The caller starts a
`nagoya::reactor::Reactor`, keeps it alive while its runs are in flight, and passes
`&reactor.handle()`. The crate never starts a reactor of its own and holds no global
one. Changed signatures:
- `run(request: &Request, reactor: &Handle) -> Result<Outcome>`
- `stream(request: &Request, reactor: &Handle) -> Result<Run>`
- `interrupt(request: &Request, reactor: &Handle) -> Result<bool>`
- `Probe::run(agent: Agent, reactor: &Handle) -> Result<Probe>`
- `Probe::run_bin(agent: Agent, bin: &str, reactor: &Handle) -> Result<Probe>`
- `AuthStatus::check(agent: Agent, reactor: &Handle) -> Result<AuthStatus>`
- `AuthStatus::check_bin(agent: Agent, bin: &str, reactor: &Handle) -> Result<AuthStatus>`
- `Agent::account_usage(self, reactor: &Handle) -> Result<AccountUsage>`
- `Agent::discover_models(&self, reactor: &Handle) -> Result<Vec<Model>>`
- `nagoya` is re-exported as `agent_abstraction::nagoya`, so a caller can name
`Reactor` and `Handle` without a second copy of the dependency.

### Changed

- **tokio is gone**, from dependencies and dev-dependencies. Child processes come from
`nagoya::process`, the driver and stderr reader run on nagoya's shared pool, and the
channels, `select_biased!` and I/O extension traits come from `futures`. Every future
this crate returns still works under any executor, tokio included.
- **`stream` no longer needs an ambient runtime.** It used to return
`Error::NoRuntime` when called outside a tokio runtime; nagoya's pool starts on first
use and the reactor is passed in, so that variant is never returned now. It stays in the enum so existing matches
compile.
- **A closed control channel no longer spins the Codex and Grok drivers.** Under tokio a
detached interactive run polled its closed channel until tokio's cooperative budget
forced a yield; the arm is now skipped once every sender is gone.

## 0.4.20

### Added
Expand Down
16 changes: 10 additions & 6 deletions Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "agent-abstraction"
version = "0.4.21"
version = "0.5.0"
edition = "2024"
# The floor edition 2024 requires, and where the strictest dependencies (uuid,
# getrandom) sit. Derived from the dependency graph rather than compile-tested.
Expand Down Expand Up @@ -32,10 +32,15 @@ name = "agent_abstraction"
path = "src/lib.rs"

[dependencies]
# `process` gives an async child with piped stdio; `io-util` the line reader that
# turns a JSONL stream into events; `time` the run timeout; `sync` the event
# channel. No `full`: the consumer picks its own runtime features.
tokio = { version = "1", features = ["process", "io-util", "sync", "time", "rt", "macros"] }
# The runtime: `spawn` for the driver task, `sleep` and `timeout` for the run
# deadline, and `process` for an async child with piped stdio. `process` needs
# the reactor, which is why both are named. The returned futures and pipes work
# under any executor, so a consumer on another runtime can still await a `Run`.
nagoya = { version = "^0.1", features = ["reactor", "process"] }
# What nagoya deliberately does not carry: the oneshot and mpsc channels, the
# `select_biased!` that stands in for tokio's biased `select!`, and the
# `AsyncRead`/`AsyncBufRead` extension traits the line reader is built on.
futures = "^0.3"
serde = { version = "1", features = ["derive"] }
serde_json = "1"
thiserror = "2"
Expand All @@ -53,7 +58,6 @@ uuid = { version = "1", features = ["v4", "serde"] }
libc = "0.2"

[dev-dependencies]
tokio = { version = "1", features = ["rt-multi-thread", "macros", "time"] }
# The live tests mint UUIDs to prove a caller-assigned session id round-trips.
uuid = { version = "1", features = ["v4"] }

Expand Down
41 changes: 26 additions & 15 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,19 +20,29 @@ It is a **library, not a CLI**. Your program links it and spawns the agent direc
nothing marshals a request through a command line and back out of stdout twice.

```rust
use agent_abstraction::nagoya::reactor::Reactor;
use agent_abstraction::{Agent, Permission, Request, run};

// The I/O driver for the child's pipes. Yours to start and keep alive while
// runs are in flight; the crate never starts one of its own.
let reactor = Reactor::start()?;

let outcome = run(
&Request::new(Agent::Claude, "Reply with the single word: pong")
.model("haiku")
.permission(Permission::ReadOnly),
&reactor.handle(),
)
.await?;

println!("{}", outcome.text); // "pong"
println!("{:?}", outcome.usage.cost_usd);
```

Every entry point that spawns a CLI (`run`, `stream`, `interrupt`, `Probe::run`,
`AuthStatus::check`, `Agent::account_usage`, `Agent::discover_models`) takes that
`&reactor.handle()`. The examples below assume a `reactor` like this one is in scope.

## What each agent can actually do

Verified live, against `claude 2.1.205`, `codex-cli 0.146.0` and `GitHub Copilot CLI 1.0.78`
Expand Down Expand Up @@ -60,7 +70,7 @@ Both, depending on the agent. Verified by round-trip, not from `--help`:
// Claude and Copilot: the id is yours to pick, so it can match a thread id
// your app already has, with no mapping table in between.
let mine = uuid::Uuid::new_v4().to_string();
let outcome = run(&Request::new(Agent::Claude, "hi").session_id(&mine)).await?;
let outcome = run(&Request::new(Agent::Claude, "hi").session_id(&mine), &reactor.handle()).await?;
assert_eq!(outcome.session.as_deref(), Some(mine.as_str()));
```

Expand All @@ -83,7 +93,7 @@ conversation it meant to branch.
## Streaming

```rust
let mut running = stream(&Request::new(Agent::Claude, "audit this repo"))?;
let mut running = stream(&Request::new(Agent::Claude, "audit this repo"), &reactor.handle())?;
while let Some(event) = running.recv().await {
match event {
Event::Text(text) => print!("{text}"),
Expand Down Expand Up @@ -123,7 +133,7 @@ let turn = Request::new(Agent::Claude, "what did I ask you to remember?")
.session(&store, ".", "thread-42", /* fork */ false)?;

assert_eq!(turn.session_phase(), Some(Phase::Continue));
let outcome = run(&turn).await?;
let outcome = run(&turn, &reactor.handle()).await?;
```

Records live at `<dir>/<project-slug>/<name>.json`, partitioned by project so the same name
Expand Down Expand Up @@ -206,7 +216,7 @@ established here; the entry says which.
Where a CLI can be asked directly, prefer that:

```rust
let models = Agent::Codex.discover_models().await?; // reflects the installed binary
let models = Agent::Codex.discover_models(&reactor.handle()).await?; // reflects the installed binary
```

`discover_models` returns `Error::Unsupported` on Claude and Copilot rather than silently
Expand Down Expand Up @@ -241,7 +251,8 @@ back parsed instead of guessing at formatting the model never promised:
let outcome = run(&Request::new(Agent::Codex, "Alice is 30 years old.")
.schema(r#"{"type":"object",
"properties":{"name":{"type":"string"},"age":{"type":"integer"}},
"required":["name","age"],"additionalProperties":false}"#))
"required":["name","age"],"additionalProperties":false}"#),
&reactor.handle())
.await?;

assert_eq!(outcome.structured.unwrap()["name"], "Alice");
Expand All @@ -265,7 +276,7 @@ cancelling a request should stop the work, not leave an agent running invisibly,
quota and writing files with nobody watching.

```rust
let running = stream(&request)?;
let running = stream(&request, &reactor.handle())?;
drop(running); // agent and its children are killed
running.cancel().await?; // cooperative: returns only once the tree has exited
running.detach(); // opt out: keep running unsupervised
Expand Down Expand Up @@ -294,7 +305,7 @@ a run already under way:

```rust
let request = Request::new(Agent::Claude, prompt).interactive();
let mut run = stream(&request)?;
let mut run = stream(&request, &reactor.handle())?;

// Keep input independent from the task continuously draining run.recv().
let control = run.control();
Expand Down Expand Up @@ -353,7 +364,7 @@ use agent_abstraction::{Agent, Command, Compaction, Event, Request, stream};
let request = Request::command(Agent::Claude, &Command::Compact { instructions: None })
.resume(&session_id);

let mut run = stream(&request)?;
let mut run = stream(&request, &reactor.handle())?;
while let Some(event) = run.recv().await {
if let Event::Compaction(Compaction::Finished { ok, error }) = event {
// `ok: false` with a reason is an answer, not an error.
Expand Down Expand Up @@ -402,7 +413,7 @@ let request = Request::new(Agent::Claude, prompt)
.permission(Permission::Edit)
.approvals();

let mut run = stream(&request)?;
let mut run = stream(&request, &reactor.handle())?;
while let Some(event) = run.recv().await {
if let Event::ApprovalRequest(approval) = event {
// approval.tool is "Bash"; approval.input carries the actual command
Expand Down Expand Up @@ -484,7 +495,7 @@ Without spending a request:

```rust
for agent in Agent::ALL {
let status = AuthStatus::check(agent).await?;
let status = AuthStatus::check(agent, &reactor.handle()).await?;
println!("{agent}: {}", status.summary());
}
```
Expand Down Expand Up @@ -613,7 +624,7 @@ These are `Error::AgentError`, carrying the agent's own wording and the provider
where one was reported:

```rust
match run(&request).await {
match run(&request, &reactor.handle()).await {
Err(Error::AgentError { status: Some(404), message, .. }) => {
// Typically a model the account cannot reach. `message` is the agent's wording.
eprintln!("{message}");
Expand All @@ -634,14 +645,14 @@ passed through untouched.

## Usage and quota

Two questions with two answers. `Outcome::usage` measures the run; `Agent::account_usage()`
Two questions with two answers. `Outcome::usage` measures the run; `Agent::account_usage`
measures the plan behind it. Everything is a value, never a formatted string or a rendered
bar, so a host presents it however it likes.

### Per run, and per session

```rust
let outcome = run(&request).await?;
let outcome = run(&request, &reactor.handle()).await?;
let used = outcome.usage.context_used(); // Option<f64>, 0.0 to 1.0
```

Expand Down Expand Up @@ -710,7 +721,7 @@ until the turn completes, so it never emits this.

```rust
if agent.reports_account_usage() {
let account = agent.account_usage().await?;
let account = agent.account_usage(&reactor.handle()).await?;
for window in &account.windows {
// window.used_percent, window.window_minutes, window.resets_at
}
Expand Down Expand Up @@ -770,7 +781,7 @@ A Rust port of [nickderobertis/oneharness](https://github.com/nickderobertis/one
two JSON round-trips to ask a question.
- **The shell scripts are gone**, 39 of them, mostly CI gates and per-harness e2e drivers.
- **Five harnesses are gone** (OpenCode, Goose, Qwen, Crush, Cursor).
- **Async throughout.** oneharness runs blocking; this streams over tokio, which is what a
- **Async throughout.** oneharness runs blocking; this streams over nagoya, which is what a
Tauri front end needs to render a run as it happens.

Some findings did not survive re-verification against the current CLIs. oneharness models
Expand Down
6 changes: 3 additions & 3 deletions docs/host-integration.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ A user who types a correction mid-turn must not have to wait for the turn to end

```rust
let request = Request::new(Agent::Claude, prompt).interactive();
let mut run = stream(&request)?;
let mut run = stream(&request, &reactor.handle())?;

// Later, from the UI, the moment the user hits enter:
run.send("actually, skip the tests and just fix the parser").await?;
Expand Down Expand Up @@ -70,7 +70,7 @@ let request = Request::new(Agent::Claude, prompt)
.permission(Permission::Edit) // not ReadOnly, see below
.approvals();

let mut run = stream(&request)?;
let mut run = stream(&request, &reactor.handle())?;
while let Some(event) = run.recv().await {
if let Event::ApprovalRequest(approval) = event {
let decision = if user_approves(&approval) {
Expand Down Expand Up @@ -119,7 +119,7 @@ let request = Request::new(Agent::Claude, prompt)
.session(&store, &project, "chat")? // so the conversation continues across turns
.approvals(); // implies .interactive()

let mut run = stream(&request)?;
let mut run = stream(&request, &reactor.handle())?;
while let Some(event) = run.recv().await {
match event {
Event::Text(chunk) => ui.append_assistant(&chunk),
Expand Down
34 changes: 22 additions & 12 deletions src/account.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,9 +36,10 @@
use std::process::Stdio;
use std::time::Duration;

use futures::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use futures::stream::StreamExt as _;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};

use crate::agent::Agent;
use crate::error::{Error, Result};
Expand Down Expand Up @@ -144,6 +145,9 @@ impl Agent {

/// Ask the agent what the account has spent and what remains.
///
/// `reactor` drives the child's pipes; the caller keeps its
/// [`nagoya::reactor::Reactor`] alive until this returns.
///
/// # Errors
/// [`Error::Unsupported`] where the agent has no headless way to answer,
/// which today is Claude and Copilot; check
Expand All @@ -152,9 +156,9 @@ impl Agent {
/// cannot be run, [`Error::Timeout`] if it does not reply,
/// [`Error::AgentError`] if it replies with a refusal, and
/// [`Error::Parse`] if the reply is not the expected shape.
pub async fn account_usage(self) -> Result<AccountUsage> {
pub async fn account_usage(self, reactor: &nagoya::reactor::Handle) -> Result<AccountUsage> {
match self {
Agent::Codex => codex_account_usage(self.bin()).await,
Agent::Codex => codex_account_usage(self.bin(), reactor).await,
// Deliberately an error rather than a half-answer assembled from a
// past run's rate-limit event: that would be neither current nor
// account-wide, and would read as though it were both.
Expand All @@ -174,15 +178,15 @@ impl Agent {
/// names the method, rather than as silence.
///
/// Verified against codex-cli 0.145.0 on 2026-07-29.
async fn codex_account_usage(bin: &str) -> Result<AccountUsage> {
let mut child = tokio::process::Command::new(bin)
async fn codex_account_usage(bin: &str, reactor: &nagoya::reactor::Handle) -> Result<AccountUsage> {
let mut child = nagoya::process::Command::new(bin)
.arg("app-server")
.stdin(Stdio::piped())
.stdout(Stdio::piped())
// Silenced rather than captured: the server logs progress here and none
// of it belongs in an error about usage.
.stderr(Stdio::null())
.spawn()
.spawn(reactor)
.map_err(|source| {
if source.kind() == std::io::ErrorKind::NotFound {
Error::NotInstalled {
Expand All @@ -199,7 +203,7 @@ async fn codex_account_usage(bin: &str) -> Result<AccountUsage> {
})?;

let exchange = codex_exchange(&mut child);
let result = match tokio::time::timeout(QUERY_TIMEOUT, exchange).await {
let result = match nagoya::timeout(QUERY_TIMEOUT, exchange).await {
Ok(result) => result,
Err(_) => Err(Error::Timeout {
bin: bin.to_string(),
Expand All @@ -214,7 +218,7 @@ async fn codex_account_usage(bin: &str) -> Result<AccountUsage> {
}

/// Drive the three requests and collect their replies.
async fn codex_exchange(child: &mut tokio::process::Child) -> Result<AccountUsage> {
async fn codex_exchange(child: &mut nagoya::process::Child) -> Result<AccountUsage> {
const ACCOUNT: i64 = 2;
const LIMITS: i64 = 3;
const USAGE: i64 = 4;
Expand Down Expand Up @@ -266,7 +270,9 @@ async fn codex_exchange(child: &mut tokio::process::Child) -> Result<AccountUsag
while outstanding > 0 {
// A stream that ends before every reply arrives leaves whatever was
// collected in place rather than discarding it.
let Ok(Some(line)) = lines.next_line().await else {
// futures' `Lines` is a stream, so tokio's `Ok(Some(line))` from
// `next_line` is `Some(Ok(line))` here; a read error still ends it.
let Some(Ok(line)) = lines.next().await else {
break;
};
if line.len() > MAX_REPLY_BYTES {
Expand Down Expand Up @@ -470,12 +476,16 @@ mod tests {

/// The capability is answerable without spawning anything, so a host can
/// decide whether to build the panel at all.
#[tokio::test]
async fn agents_that_cannot_report_say_so_without_being_asked_twice() {
#[test]
fn agents_that_cannot_report_say_so_without_being_asked_twice() {
let reactor = nagoya::reactor::Reactor::start().expect("reactor");
for agent in [Agent::Claude, Agent::Copilot] {
assert!(!agent.reports_account_usage(), "{agent}");
assert!(
matches!(agent.account_usage().await, Err(Error::Unsupported { .. })),
matches!(
nagoya::block_on(agent.account_usage(&reactor.handle())),
Err(Error::Unsupported { .. })
),
"{agent} should refuse rather than assemble a partial answer"
);
}
Expand Down
Loading
Loading