hive-runtime, hive-agent: model/effort picker for ACP agents from configOptions
An ACP agent's model and effort pickers now list what its session offers (its `model` and `thought_level` config options) instead of the claude model list and EFFORT_LEVELS. A pick goes through the same Bus::set_model / Bus::set_effort -> Config.model / Config.effort path as on claude; before each prompt the ACP runtime sets it with `session/set_config_option`, model first, and only when the session offers that value. Options are re-read from the set response and from `config_option_update`, so the effort picker disappears when the chosen model offers no effort levels, and is hidden while a newly picked model waits for the next turn. The session's options reach the web UI through a `Choices` handle from a new `Runtime::choices`, registered on the bus the way `canceller` is. /api/model and /api/effort accept only offered values on ACP. The claude path is unchanged. Refs #4391
This commit is contained in:
parent
abd547f600
commit
f189724a4c
13 changed files with 559 additions and 44 deletions
|
|
@ -53,7 +53,7 @@ hand-maintained per-file tree drifts out of sync with the code.
|
||||||
launch-config layer (tool-group/capability → `--allowedTools`,
|
launch-config layer (tool-group/capability → `--allowedTools`,
|
||||||
`--mcp-config` render).
|
`--mcp-config` render).
|
||||||
- **`hive-runtime/`** — the runtime an agent's turns run on: one
|
- **`hive-runtime/`** — the runtime an agent's turns run on: one
|
||||||
`Runtime` interface (`run`/`compact`/`archive`/`canceller`) with a `claude` backend
|
`Runtime` interface (`run`/`compact`/`archive`/`canceller`/`choices`) with a `claude` backend
|
||||||
(a pass-through to `hive-claude`) and an `acp` backend (any Agent Client
|
(a pass-through to `hive-claude`) and an `acp` backend (any Agent Client
|
||||||
Protocol agent, spawned from the command/args/env nix hands it — no agent
|
Protocol agent, spawned from the command/args/env nix hands it — no agent
|
||||||
is named in Rust). The ACP backend translates its stream into claude's
|
is named in Rust). The ACP backend translates its stream into claude's
|
||||||
|
|
|
||||||
|
|
@ -307,8 +307,9 @@ On `acp`:
|
||||||
system prompt.
|
system prompt.
|
||||||
- `interrupt` sends `session/cancel`, and `force` changes nothing. The
|
- `interrupt` sends `session/cancel`, and `force` changes nothing. The
|
||||||
runtime kills an agent that ignores the cancel for 10s.
|
runtime kills an agent that ignores the cancel for 10s.
|
||||||
- `model` and `effort` do nothing: the agent runs the model it's
|
- the runtime sets `model` and `effort` on the session when the agent
|
||||||
configured with.
|
offers that value for its `model` or `thought_level` config option. Any
|
||||||
|
other value does nothing (`hive_runtime::AcpRuntime`).
|
||||||
- the agent's permission requests get the same answers a claude subagent's
|
- the agent's permission requests get the same answers a claude subagent's
|
||||||
tool list gives: MCP tools and file tools, `fetch` with `web_tools`,
|
tool list gives: MCP tools and file tools, `fetch` with `web_tools`,
|
||||||
`execute` with `execution` (`session::acp_permits`).
|
`execute` with `execution` (`session::acp_permits`).
|
||||||
|
|
|
||||||
|
|
@ -61,6 +61,11 @@ structurally rather than for one specific trigger. Two columns:
|
||||||
display) this rewrite exists to fix.
|
display) this rewrite exists to fix.
|
||||||
- **Effort badge** (`effort · <level> ▾`): same shape, `/api/effort`,
|
- **Effort badge** (`effort · <level> ▾`): same shape, `/api/effort`,
|
||||||
shown when `state.available_efforts` is non-empty.
|
shown when `state.available_efforts` is non-empty.
|
||||||
|
- On an ACP agent both pickers list what its session offers (its
|
||||||
|
`model` and `thought_level` config options), and are empty until the
|
||||||
|
first turn has attached a session. The effort picker is hidden when the
|
||||||
|
model offers no effort levels, and while a newly picked model waits for
|
||||||
|
the next turn, since the levels come with the model.
|
||||||
- **Ctx / cost badges**: `ctx · 142k` (last inference's prompt size,
|
- **Ctx / cost badges**: `ctx · 142k` (last inference's prompt size,
|
||||||
tooltip shows % of context window) / `cost · 1.3M` (cumulative
|
tooltip shows % of context window) / `cost · 1.3M` (cumulative
|
||||||
tokens billed across every inference in the last turn).
|
tokens billed across every inference in the last turn).
|
||||||
|
|
@ -303,6 +308,9 @@ shaped).
|
||||||
rejected rather than forwarded to `claude --effort`. Persists via
|
rejected rather than forwarded to `claude --effort`. Persists via
|
||||||
`Bus::set_effort`, which emits `EffortChanged`. Applies on the next
|
`Bus::set_effort`, which emits `EffortChanged`. Applies on the next
|
||||||
session start (no mid-session swap).
|
session start (no mid-session swap).
|
||||||
|
- On an ACP agent, `/api/model` and `/api/effort` accept only values the
|
||||||
|
session offers, and the runtime sets them with `session/set_config_option`
|
||||||
|
before the next prompt.
|
||||||
- `POST /api/new-session` — arm a one-shot for the next turn to
|
- `POST /api/new-session` — arm a one-shot for the next turn to
|
||||||
drop `--continue`. Emits a `LiveEvent::Note`.
|
drop `--continue`. Emits a `LiveEvent::Note`.
|
||||||
- `POST /api/logout` — three-step teardown that re-uses the
|
- `POST /api/logout` — three-step teardown that re-uses the
|
||||||
|
|
|
||||||
|
|
@ -335,6 +335,11 @@ pub struct Bus {
|
||||||
/// by the serve loop when it builds the session; unset on claude, whose
|
/// by the serve loop when it builds the session; unset on claude, whose
|
||||||
/// turn `/api/cancel` stops by signalling the `claude` process.
|
/// turn `/api/cancel` stops by signalling the `claude` process.
|
||||||
turn_canceller: Arc<OnceLock<hive_runtime::Canceller>>,
|
turn_canceller: Arc<OnceLock<hive_runtime::Canceller>>,
|
||||||
|
/// The model and effort the session offers, for a runtime that reads them
|
||||||
|
/// from its session (ACP). Set once by the serve loop when it builds the
|
||||||
|
/// session; unset on claude, whose pickers list the configured models and
|
||||||
|
/// `EFFORT_LEVELS`.
|
||||||
|
session_choices: Arc<OnceLock<hive_runtime::Choices>>,
|
||||||
/// Current fresh-claude-session id (FK to `sessions.id`). Set by the
|
/// Current fresh-claude-session id (FK to `sessions.id`). Set by the
|
||||||
/// bin loop after minting a session row on a fresh start; stamped onto
|
/// bin loop after minting a session row on a fresh start; stamped onto
|
||||||
/// every `turn_stats` row until the next fresh session. `None` before
|
/// every `turn_stats` row until the next fresh session. `None` before
|
||||||
|
|
@ -427,6 +432,7 @@ impl Bus {
|
||||||
compact_pending: Arc::new(Mutex::new(None)),
|
compact_pending: Arc::new(Mutex::new(None)),
|
||||||
post_compact_wake: Arc::new(Mutex::new(None)),
|
post_compact_wake: Arc::new(Mutex::new(None)),
|
||||||
turn_canceller: Arc::default(),
|
turn_canceller: Arc::default(),
|
||||||
|
session_choices: Arc::default(),
|
||||||
session_id: Arc::new(Mutex::new(None)),
|
session_id: Arc::new(Mutex::new(None)),
|
||||||
fresh_session: Arc::new(AtomicBool::new(false)),
|
fresh_session: Arc::new(AtomicBool::new(false)),
|
||||||
tool_calls: Arc::new(Mutex::new(std::collections::HashMap::new())),
|
tool_calls: Arc::new(Mutex::new(std::collections::HashMap::new())),
|
||||||
|
|
@ -463,6 +469,19 @@ impl Bus {
|
||||||
self.turn_canceller.get()
|
self.turn_canceller.get()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Record where the session's model and effort choices are read. Only the
|
||||||
|
/// first call takes effect; there is one session per harness.
|
||||||
|
pub fn set_session_choices(&self, choices: hive_runtime::Choices) {
|
||||||
|
let _ = self.session_choices.set(choices);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Where the session's model and effort choices are read, if the runtime
|
||||||
|
/// offers them (see [`Self::set_session_choices`]).
|
||||||
|
#[must_use]
|
||||||
|
pub fn session_choices(&self) -> Option<&hive_runtime::Choices> {
|
||||||
|
self.session_choices.get()
|
||||||
|
}
|
||||||
|
|
||||||
/// Request a session reset (operator `POST /api/new-session`). Deferred:
|
/// Request a session reset (operator `POST /api/new-session`). Deferred:
|
||||||
/// the flag is consumed at the next turn boundary by `drive_turn`, which
|
/// the flag is consumed at the next turn boundary by `drive_turn`, which
|
||||||
/// archives the current session so no claude process is mid-write when the
|
/// archives the current session so no claude process is mid-write when the
|
||||||
|
|
|
||||||
|
|
@ -635,6 +635,16 @@ async fn serve_main<S: Surface>(socket: &Path, poll_ms: u64) -> Result<()> {
|
||||||
.await
|
.await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Hand the bus the session's handles the web UI reads from outside a turn.
|
||||||
|
fn share_session_handles(bus: &Bus, session: &turn::AgentSession) {
|
||||||
|
if let Some(canceller) = hive_runtime::Runtime::canceller(session) {
|
||||||
|
bus.set_turn_canceller(canceller);
|
||||||
|
}
|
||||||
|
if let Some(choices) = hive_runtime::Runtime::choices(session) {
|
||||||
|
bus.set_session_choices(choices);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// The long-running message loop. Long-polls the broker via
|
/// The long-running message loop. Long-polls the broker via
|
||||||
/// `S::recv_next`, drives a turn per message, parks on auth-failed,
|
/// `S::recv_next`, drives a turn per message, parks on auth-failed,
|
||||||
/// otherwise retries.
|
/// otherwise retries.
|
||||||
|
|
@ -662,9 +672,7 @@ async fn serve_loop<S: Surface>(
|
||||||
// The durable agent session, built once and reused for every turn +
|
// The durable agent session, built once and reused for every turn +
|
||||||
// idle compaction below (it's effectively stateless).
|
// idle compaction below (it's effectively stateless).
|
||||||
let session = turn::make_session(&bus)?;
|
let session = turn::make_session(&bus)?;
|
||||||
if let Some(canceller) = hive_runtime::Runtime::canceller(&session) {
|
share_session_handles(&bus, &session);
|
||||||
bus.set_turn_canceller(canceller);
|
|
||||||
}
|
|
||||||
// Tracks the last observed pause state so the transitions get logged
|
// Tracks the last observed pause state so the transitions get logged
|
||||||
// once each instead of twelve lines a minute while parked.
|
// once each instead of twelve lines a minute while parked.
|
||||||
let mut was_paused = false;
|
let mut was_paused = false;
|
||||||
|
|
|
||||||
|
|
@ -136,7 +136,8 @@ pub(super) struct ModelForm {
|
||||||
/// `recv` cycle. Empty / whitespace-only inputs are rejected. No
|
/// `recv` cycle. Empty / whitespace-only inputs are rejected. No
|
||||||
/// claude-side validation — we just hand the string through to
|
/// claude-side validation — we just hand the string through to
|
||||||
/// `claude --model <name>`; an unknown model surfaces as a turn
|
/// `claude --model <name>`; an unknown model surfaces as a turn
|
||||||
/// failure in the live panel and the operator can revert.
|
/// failure in the live panel and the operator can revert. On ACP the name
|
||||||
|
/// must be one the session offers, since the runtime sets no other.
|
||||||
pub(super) async fn post_set_model(
|
pub(super) async fn post_set_model(
|
||||||
State(state): State<AppState>,
|
State(state): State<AppState>,
|
||||||
Form(form): Form<ModelForm>,
|
Form(form): Form<ModelForm>,
|
||||||
|
|
@ -145,10 +146,17 @@ pub(super) async fn post_set_model(
|
||||||
if name.is_empty() {
|
if name.is_empty() {
|
||||||
return error_response(StatusCode::BAD_REQUEST, "model: name required");
|
return error_response(StatusCode::BAD_REQUEST, "model: name required");
|
||||||
}
|
}
|
||||||
|
let text = if state.bus.session_choices().is_some() {
|
||||||
|
let offered = super::state::pickers(&state.bus).available_models;
|
||||||
|
if !offered.iter().any(|m| m == name) {
|
||||||
|
return not_offered("model", &offered);
|
||||||
|
}
|
||||||
|
format!("operator: /model — model set to '{name}' from the next turn")
|
||||||
|
} else {
|
||||||
|
format!("operator: /model — claude model set to '{name}' for future turns")
|
||||||
|
};
|
||||||
state.bus.set_model(name);
|
state.bus.set_model(name);
|
||||||
state.bus.emit(crate::events::LiveEvent::Note {
|
state.bus.emit(crate::events::LiveEvent::Note { text });
|
||||||
text: format!("operator: /model — claude model set to '{name}' for future turns"),
|
|
||||||
});
|
|
||||||
tracing::info!(%name, "operator set model");
|
tracing::info!(%name, "operator set model");
|
||||||
(axum::http::StatusCode::OK, "ok").into_response()
|
(axum::http::StatusCode::OK, "ok").into_response()
|
||||||
}
|
}
|
||||||
|
|
@ -163,13 +171,22 @@ pub(super) struct EffortForm {
|
||||||
/// server-side against [`crate::harness_state::EFFORT_LEVELS`] — an out-of-set
|
/// server-side against [`crate::harness_state::EFFORT_LEVELS`] — an out-of-set
|
||||||
/// value is rejected rather than handed to `claude --effort`, since an
|
/// value is rejected rather than handed to `claude --effort`, since an
|
||||||
/// unknown level would fail every subsequent launch. Applies on the next
|
/// unknown level would fail every subsequent launch. Applies on the next
|
||||||
/// session start (no mid-session swap).
|
/// session start (no mid-session swap). On ACP the level must instead be one
|
||||||
|
/// the session offers for its model, and applies from the next turn.
|
||||||
pub(super) async fn post_set_effort(
|
pub(super) async fn post_set_effort(
|
||||||
State(state): State<AppState>,
|
State(state): State<AppState>,
|
||||||
Form(form): Form<EffortForm>,
|
Form(form): Form<EffortForm>,
|
||||||
) -> Response {
|
) -> Response {
|
||||||
let level = form.effort.trim();
|
let level = form.effort.trim();
|
||||||
if !crate::harness_state::is_valid_effort(level) {
|
let text = if state.bus.session_choices().is_some() {
|
||||||
|
let offered = super::state::pickers(&state.bus).available_efforts;
|
||||||
|
if !offered.iter().any(|e| e == level) {
|
||||||
|
return not_offered("effort", &offered);
|
||||||
|
}
|
||||||
|
format!("operator: /effort — effort set to '{level}' from the next turn")
|
||||||
|
} else if crate::harness_state::is_valid_effort(level) {
|
||||||
|
format!("operator: /effort — claude effort set to '{level}' for future sessions")
|
||||||
|
} else {
|
||||||
return error_response(
|
return error_response(
|
||||||
StatusCode::BAD_REQUEST,
|
StatusCode::BAD_REQUEST,
|
||||||
&format!(
|
&format!(
|
||||||
|
|
@ -177,15 +194,23 @@ pub(super) async fn post_set_effort(
|
||||||
crate::harness_state::EFFORT_LEVELS.join(", ")
|
crate::harness_state::EFFORT_LEVELS.join(", ")
|
||||||
),
|
),
|
||||||
);
|
);
|
||||||
}
|
};
|
||||||
state.bus.set_effort(level);
|
state.bus.set_effort(level);
|
||||||
state.bus.emit(crate::events::LiveEvent::Note {
|
state.bus.emit(crate::events::LiveEvent::Note { text });
|
||||||
text: format!("operator: /effort — claude effort set to '{level}' for future sessions"),
|
|
||||||
});
|
|
||||||
tracing::info!(%level, "operator set effort");
|
tracing::info!(%level, "operator set effort");
|
||||||
(axum::http::StatusCode::OK, "ok").into_response()
|
(axum::http::StatusCode::OK, "ok").into_response()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Refuses a pick the ACP session does not offer.
|
||||||
|
fn not_offered(what: &str, offered: &[String]) -> Response {
|
||||||
|
let message = if offered.is_empty() {
|
||||||
|
format!("{what}: the agent's session offers no choice now")
|
||||||
|
} else {
|
||||||
|
format!("{what}: the agent's session offers {}", offered.join(", "))
|
||||||
|
};
|
||||||
|
error_response(StatusCode::BAD_REQUEST, &message)
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Deserialize)]
|
#[derive(Deserialize)]
|
||||||
pub(super) struct MarkTodosDoneForm {
|
pub(super) struct MarkTodosDoneForm {
|
||||||
/// Comma-separated todo ids. Same "one field, JS joins the checked
|
/// Comma-separated todo ids. Same "one field, JS joins the checked
|
||||||
|
|
|
||||||
|
|
@ -46,7 +46,7 @@ pub(super) async fn api_state(State(state): State<AppState>) -> axum::Json<State
|
||||||
let context_window_tokens = state.bus.effective_context_window(&model);
|
let context_window_tokens = state.bus.effective_context_window(&model);
|
||||||
let ctx_usage = state.bus.last_ctx_usage();
|
let ctx_usage = state.bus.last_ctx_usage();
|
||||||
let cost_usage = state.bus.last_cost_usage();
|
let cost_usage = state.bus.last_cost_usage();
|
||||||
let effort = state.bus.effort();
|
let pickers = pickers(&state.bus);
|
||||||
axum::Json(StateSnapshot {
|
axum::Json(StateSnapshot {
|
||||||
seq,
|
seq,
|
||||||
label: state.label.clone(),
|
label: state.label.clone(),
|
||||||
|
|
@ -57,7 +57,7 @@ pub(super) async fn api_state(State(state): State<AppState>) -> axum::Json<State
|
||||||
inbox,
|
inbox,
|
||||||
turn_state,
|
turn_state,
|
||||||
turn_state_since,
|
turn_state_since,
|
||||||
model,
|
model: pickers.model,
|
||||||
resolved_model,
|
resolved_model,
|
||||||
context_window_tokens,
|
context_window_tokens,
|
||||||
ctx_usage,
|
ctx_usage,
|
||||||
|
|
@ -68,12 +68,9 @@ pub(super) async fn api_state(State(state): State<AppState>) -> axum::Json<State
|
||||||
.filter(|s| !s.is_empty()),
|
.filter(|s| !s.is_empty()),
|
||||||
hive_name: crate::identity::hive_name(),
|
hive_name: crate::identity::hive_name(),
|
||||||
swarm_name: crate::identity::swarm_name(),
|
swarm_name: crate::identity::swarm_name(),
|
||||||
available_models: available_models(),
|
available_models: pickers.available_models,
|
||||||
effort,
|
effort: pickers.effort,
|
||||||
available_efforts: crate::harness_state::EFFORT_LEVELS
|
available_efforts: pickers.available_efforts,
|
||||||
.iter()
|
|
||||||
.map(ToString::to_string)
|
|
||||||
.collect(),
|
|
||||||
paused: crate::paths::paused_marker().exists(),
|
paused: crate::paths::paused_marker().exists(),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
@ -190,6 +187,7 @@ pub(super) struct StateSnapshot {
|
||||||
/// absent or empty. The frontend model quick-picker renders one button
|
/// absent or empty. The frontend model quick-picker renders one button
|
||||||
/// per entry in this list, so operators can add new models or drop
|
/// per entry in this list, so operators can add new models or drop
|
||||||
/// ones they don't want without touching the frontend code.
|
/// ones they don't want without touching the frontend code.
|
||||||
|
/// On ACP, the models the session offers instead (see [`pickers`]).
|
||||||
available_models: Vec<String>,
|
available_models: Vec<String>,
|
||||||
/// Currently-active claude effort level. Reflected on the page so the
|
/// Currently-active claude effort level. Reflected on the page so the
|
||||||
/// operator's effort picker shows the live selection. Mutable at
|
/// operator's effort picker shows the live selection. Mutable at
|
||||||
|
|
@ -199,6 +197,8 @@ pub(super) struct StateSnapshot {
|
||||||
/// (`low`, `medium`, `high`, `xhigh`, `max`) — sourced from
|
/// (`low`, `medium`, `high`, `xhigh`, `max`) — sourced from
|
||||||
/// [`crate::harness_state::EFFORT_LEVELS`], not operator-configurable like
|
/// [`crate::harness_state::EFFORT_LEVELS`], not operator-configurable like
|
||||||
/// `available_models`. The frontend renders one button per entry.
|
/// `available_models`. The frontend renders one button per entry.
|
||||||
|
/// On ACP, the levels the session offers instead, empty when it offers
|
||||||
|
/// none, which hides the picker (see [`pickers`]).
|
||||||
available_efforts: Vec<String>,
|
available_efforts: Vec<String>,
|
||||||
/// Whether this agent's turn loop is currently parked (the harness
|
/// Whether this agent's turn loop is currently parked (the harness
|
||||||
/// keeps serving this page + its MCP daemons but drives no turns).
|
/// keeps serving this page + its MCP daemons but drives no turns).
|
||||||
|
|
@ -443,3 +443,123 @@ fn available_models() -> Vec<String> {
|
||||||
models
|
models
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The model and effort pickers: the value each shows and the values it
|
||||||
|
/// offers.
|
||||||
|
#[derive(Debug, PartialEq, Eq)]
|
||||||
|
pub(super) struct Pickers {
|
||||||
|
pub(super) model: String,
|
||||||
|
pub(super) available_models: Vec<String>,
|
||||||
|
pub(super) effort: String,
|
||||||
|
pub(super) available_efforts: Vec<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// On claude, the requested model and effort, the configured models and
|
||||||
|
/// [`crate::harness_state::EFFORT_LEVELS`]. On ACP, what the session offers
|
||||||
|
/// (see [`offered_pickers`]).
|
||||||
|
pub(super) fn pickers(bus: &crate::events::Bus) -> Pickers {
|
||||||
|
let (model, effort) = (bus.model(), bus.effort());
|
||||||
|
match bus.session_choices() {
|
||||||
|
Some(choices) => offered_pickers(model, effort, choices.get()),
|
||||||
|
None => Pickers {
|
||||||
|
model,
|
||||||
|
available_models: available_models(),
|
||||||
|
effort,
|
||||||
|
available_efforts: crate::harness_state::EFFORT_LEVELS
|
||||||
|
.iter()
|
||||||
|
.map(ToString::to_string)
|
||||||
|
.collect(),
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Pickers for what an ACP session offers. The runtime sets a requested value
|
||||||
|
/// on the next turn, and only if the session offers it, so each picker shows
|
||||||
|
/// the requested value when offered and the session's own otherwise. The
|
||||||
|
/// effort levels on offer come with the model, so while a new model waits
|
||||||
|
/// for the next turn the effort picker is empty.
|
||||||
|
fn offered_pickers(
|
||||||
|
model: String,
|
||||||
|
effort: String,
|
||||||
|
offered: hive_runtime::SessionChoices,
|
||||||
|
) -> Pickers {
|
||||||
|
let shown = |wanted: String, choice: Option<&hive_runtime::Choice>| match choice {
|
||||||
|
Some(c) if !c.values.contains(&wanted) => c.current.clone(),
|
||||||
|
_ => wanted,
|
||||||
|
};
|
||||||
|
let model_pending = offered
|
||||||
|
.model
|
||||||
|
.as_ref()
|
||||||
|
.is_some_and(|c| c.current != model && c.values.contains(&model));
|
||||||
|
let effort_choice = offered.effort.filter(|_| !model_pending);
|
||||||
|
Pickers {
|
||||||
|
model: shown(model, offered.model.as_ref()),
|
||||||
|
effort: shown(effort, effort_choice.as_ref()),
|
||||||
|
available_models: offered.model.map(|c| c.values).unwrap_or_default(),
|
||||||
|
available_efforts: effort_choice.map(|c| c.values).unwrap_or_default(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use hive_runtime::{Choice, SessionChoices};
|
||||||
|
|
||||||
|
use super::{Pickers, offered_pickers};
|
||||||
|
|
||||||
|
fn choice(current: &str, values: &[&str]) -> Choice {
|
||||||
|
Choice {
|
||||||
|
current: current.to_owned(),
|
||||||
|
values: values.iter().map(|v| (*v).to_owned()).collect(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn strings(values: &[&str]) -> Vec<String> {
|
||||||
|
values.iter().map(|v| (*v).to_owned()).collect()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn acp_pickers_list_what_the_session_offers() {
|
||||||
|
let offered = SessionChoices {
|
||||||
|
model: Some(choice("m/think", &["m/think", "m/plain"])),
|
||||||
|
effort: Some(choice("low", &["low", "high"])),
|
||||||
|
};
|
||||||
|
// `haiku` and `max` are not on offer: the session's own values show.
|
||||||
|
assert_eq!(
|
||||||
|
offered_pickers("haiku".into(), "max".into(), offered.clone()),
|
||||||
|
Pickers {
|
||||||
|
model: "m/think".into(),
|
||||||
|
available_models: strings(&["m/think", "m/plain"]),
|
||||||
|
effort: "low".into(),
|
||||||
|
available_efforts: strings(&["low", "high"]),
|
||||||
|
}
|
||||||
|
);
|
||||||
|
let picked = offered_pickers("m/think".into(), "high".into(), offered);
|
||||||
|
assert_eq!(picked.effort, "high");
|
||||||
|
assert_eq!(picked.available_efforts, strings(&["low", "high"]));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn acp_effort_picker_hides_without_effort_or_while_a_model_is_pending() {
|
||||||
|
let plain = SessionChoices {
|
||||||
|
model: Some(choice("m/plain", &["m/think", "m/plain"])),
|
||||||
|
effort: None,
|
||||||
|
};
|
||||||
|
let shown = offered_pickers("m/plain".into(), "high".into(), plain);
|
||||||
|
assert_eq!(shown.model, "m/plain");
|
||||||
|
assert!(shown.available_efforts.is_empty());
|
||||||
|
|
||||||
|
let think = SessionChoices {
|
||||||
|
model: Some(choice("m/think", &["m/think", "m/plain"])),
|
||||||
|
effort: Some(choice("low", &["low", "high"])),
|
||||||
|
};
|
||||||
|
let pending = offered_pickers("m/plain".into(), "high".into(), think);
|
||||||
|
assert_eq!(pending.model, "m/plain");
|
||||||
|
assert!(pending.available_efforts.is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn acp_pickers_are_empty_before_a_session_offers_anything() {
|
||||||
|
let none = offered_pickers("haiku".into(), "high".into(), SessionChoices::default());
|
||||||
|
assert!(none.available_models.is_empty() && none.available_efforts.is_empty());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -1,7 +1,7 @@
|
||||||
# hive-runtime
|
# hive-runtime
|
||||||
|
|
||||||
The layer an agent's turns are driven through: one `Runtime` interface
|
The layer an agent's turns are driven through: one `Runtime` interface
|
||||||
(`run`, `compact`, `archive`, `canceller`) with a backend per runtime.
|
(`run`, `compact`, `archive`, `canceller`, `choices`) with a backend per runtime.
|
||||||
|
|
||||||
- **claude** — `claude --print` through the `hive-claude` crate's
|
- **claude** — `claude --print` through the `hive-claude` crate's
|
||||||
`InfiniteSession`. A pass-through: same spawn, same session handling, same
|
`InfiniteSession`. A pass-through: same spawn, same session handling, same
|
||||||
|
|
@ -56,3 +56,14 @@ past the watermark, or on `compact` (the operator's `/compact`, the agent's
|
||||||
window if that is shorter), is cancelled and handled as the case above; the
|
window if that is shorter), is cancelled and handled as the case above; the
|
||||||
checkpoint turn is not run a second time. So a compaction always leaves a
|
checkpoint turn is not run a second time. So a compaction always leaves a
|
||||||
smaller session behind, and a failed one is not retried on the next turn.
|
smaller session behind, and a failed one is not retried on the next turn.
|
||||||
|
|
||||||
|
## ACP backend: model and effort
|
||||||
|
|
||||||
|
Before each prompt, `Config::model` and then `Config::effort` are set on the
|
||||||
|
session with `session/set_config_option`, as the value of its first `model`
|
||||||
|
and `thought_level` config option. A value the session does not offer, or
|
||||||
|
already holds, is not sent, so a claude alias such as `haiku` does nothing
|
||||||
|
on an ACP agent. Model goes first because the effort levels on offer depend
|
||||||
|
on it. The `Choices` handle (`Runtime::choices`) reads what the session
|
||||||
|
offers now, including after a `config_option_update`. It is empty until a
|
||||||
|
session is attached.
|
||||||
|
|
|
||||||
|
|
@ -11,8 +11,8 @@ mod stream;
|
||||||
use std::collections::VecDeque;
|
use std::collections::VecDeque;
|
||||||
use std::path::{Path, PathBuf};
|
use std::path::{Path, PathBuf};
|
||||||
use std::pin::Pin;
|
use std::pin::Pin;
|
||||||
use std::sync::Arc;
|
|
||||||
use std::sync::atomic::{AtomicBool, Ordering};
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
|
use std::sync::{Arc, PoisonError};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
use hive_claude::{CompactionPolicy, Config, Progress, Sink};
|
use hive_claude::{CompactionPolicy, Config, Progress, Sink};
|
||||||
|
|
@ -54,6 +54,10 @@ const COMPACT_COMMAND: &str = "compact";
|
||||||
/// session is attached, or never.
|
/// session is attached, or never.
|
||||||
const COMMANDS_WAIT: Duration = SETTLE_MAX;
|
const COMMANDS_WAIT: Duration = SETTLE_MAX;
|
||||||
|
|
||||||
|
/// The config option categories [`Config::model`] and [`Config::effort`] set.
|
||||||
|
const MODEL_CATEGORY: &str = "model";
|
||||||
|
const EFFORT_CATEGORY: &str = "thought_level";
|
||||||
|
|
||||||
/// One `session/request_permission` request, as a [`PermissionPolicy`] sees it.
|
/// One `session/request_permission` request, as a [`PermissionPolicy`] sees it.
|
||||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
pub struct PermissionAsk<'a> {
|
pub struct PermissionAsk<'a> {
|
||||||
|
|
@ -149,6 +153,58 @@ impl Canceller {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A value the session lets the client pick: the one it has, and those it
|
||||||
|
/// accepts.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub struct Choice {
|
||||||
|
pub current: String,
|
||||||
|
pub values: Vec<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The model and effort the loaded session lets the client pick, from its
|
||||||
|
/// first `model` and `thought_level` config options. `None` where it offers
|
||||||
|
/// none; which effort levels a session offers can change with its model.
|
||||||
|
#[derive(Debug, Clone, Default, PartialEq, Eq)]
|
||||||
|
pub struct SessionChoices {
|
||||||
|
pub model: Option<Choice>,
|
||||||
|
pub effort: Option<Choice>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Reads the [`SessionChoices`] of an [`AcpRuntime`]'s session from outside
|
||||||
|
/// the call driving it. Cheap to clone.
|
||||||
|
#[derive(Clone, Default)]
|
||||||
|
pub struct Choices(Arc<std::sync::Mutex<Vec<stream::ConfigOption>>>);
|
||||||
|
|
||||||
|
impl Choices {
|
||||||
|
/// What the session offers now. Empty until a session is attached, and
|
||||||
|
/// after the agent is respawned until one is again.
|
||||||
|
#[must_use]
|
||||||
|
pub fn get(&self) -> SessionChoices {
|
||||||
|
let choice = |category| {
|
||||||
|
self.find(category).map(|o| Choice {
|
||||||
|
current: o.current,
|
||||||
|
values: o.values,
|
||||||
|
})
|
||||||
|
};
|
||||||
|
SessionChoices {
|
||||||
|
model: choice(MODEL_CATEGORY),
|
||||||
|
effort: choice(EFFORT_CATEGORY),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn find(&self, category: &str) -> Option<stream::ConfigOption> {
|
||||||
|
let options = self.0.lock().unwrap_or_else(PoisonError::into_inner);
|
||||||
|
options
|
||||||
|
.iter()
|
||||||
|
.find(|o| o.category.as_deref() == Some(category))
|
||||||
|
.cloned()
|
||||||
|
}
|
||||||
|
|
||||||
|
fn replace(&self, options: Vec<stream::ConfigOption>) {
|
||||||
|
*self.0.lock().unwrap_or_else(PoisonError::into_inner) = options;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Marks a turn in flight for [`Canceller::cancel`] while it lives.
|
/// Marks a turn in flight for [`Canceller::cancel`] while it lives.
|
||||||
struct InTurn<'a>(&'a AtomicBool);
|
struct InTurn<'a>(&'a AtomicBool);
|
||||||
|
|
||||||
|
|
@ -183,7 +239,9 @@ enum TurnKind {
|
||||||
/// since ACP has no system prompt) and `idle_timeout` (with no
|
/// since ACP has no system prompt) and `idle_timeout` (with no
|
||||||
/// `session/update` for that long the turn is cancelled, and fails with
|
/// `session/update` for that long the turn is cancelled, and fails with
|
||||||
/// [`AcpError::IdleTimeout`]). Everything else in it is claude's and is
|
/// [`AcpError::IdleTimeout`]). Everything else in it is claude's and is
|
||||||
/// ignored.
|
/// ignored, except `model` and `effort`: before each prompt they are set as
|
||||||
|
/// the session's `model` and then `thought_level` config option, when the
|
||||||
|
/// session offers that value. [`Runtime::choices`] reads what it offers.
|
||||||
///
|
///
|
||||||
/// Compaction, proactive once `policy` says so after a turn or on
|
/// Compaction, proactive once `policy` says so after a turn or on
|
||||||
/// [`Runtime::compact`], runs the agent's advertised `compact` command as a
|
/// [`Runtime::compact`], runs the agent's advertised `compact` command as a
|
||||||
|
|
@ -199,6 +257,7 @@ pub struct AcpRuntime<P: CompactionPolicy> {
|
||||||
cancel: Arc<CancelState>,
|
cancel: Arc<CancelState>,
|
||||||
cancel_grace: Duration,
|
cancel_grace: Duration,
|
||||||
compact_idle: Duration,
|
compact_idle: Duration,
|
||||||
|
choices: Choices,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The running agent process and the session loaded into it.
|
/// The running agent process and the session loaded into it.
|
||||||
|
|
@ -211,6 +270,20 @@ struct Live {
|
||||||
/// The commands the agent last advertised for `loaded`; `None` until it
|
/// The commands the agent last advertised for `loaded`; `None` until it
|
||||||
/// has.
|
/// has.
|
||||||
commands: Option<Vec<String>>,
|
commands: Option<Vec<String>>,
|
||||||
|
/// The config options `loaded` last reported.
|
||||||
|
choices: Choices,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Live {
|
||||||
|
fn offer(&mut self, options: Vec<stream::ConfigOption>) {
|
||||||
|
if let Some(model) = options
|
||||||
|
.iter()
|
||||||
|
.find(|o| o.category.as_deref() == Some(MODEL_CATEGORY))
|
||||||
|
{
|
||||||
|
self.model = Some(model.current.clone());
|
||||||
|
}
|
||||||
|
self.choices.replace(options);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<P: CompactionPolicy> AcpRuntime<P> {
|
impl<P: CompactionPolicy> AcpRuntime<P> {
|
||||||
|
|
@ -233,6 +306,7 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
||||||
cancel: Arc::default(),
|
cancel: Arc::default(),
|
||||||
cancel_grace: CANCEL_GRACE,
|
cancel_grace: CANCEL_GRACE,
|
||||||
compact_idle: COMPACT_IDLE,
|
compact_idle: COMPACT_IDLE,
|
||||||
|
choices: Choices::default(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -266,12 +340,14 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
||||||
return Err(AcpError::NoHttpMcp.into());
|
return Err(AcpError::NoHttpMcp.into());
|
||||||
}
|
}
|
||||||
tracing::info!(agent = %init["agentInfo"], "ACP agent initialized");
|
tracing::info!(agent = %init["agentInfo"], "ACP agent initialized");
|
||||||
|
self.choices.replace(Vec::new());
|
||||||
Ok(Live {
|
Ok(Live {
|
||||||
conn,
|
conn,
|
||||||
load_session: caps["loadSession"] == Value::Bool(true),
|
load_session: caps["loadSession"] == Value::Bool(true),
|
||||||
loaded: None,
|
loaded: None,
|
||||||
model: None,
|
model: None,
|
||||||
commands: None,
|
commands: None,
|
||||||
|
choices: self.choices.clone(),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -303,6 +379,7 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
||||||
match live.conn.request("session/load", params).await {
|
match live.conn.request("session/load", params).await {
|
||||||
Ok(response) => {
|
Ok(response) => {
|
||||||
live.model = stream::session_model(&response);
|
live.model = stream::session_model(&response);
|
||||||
|
live.offer(stream::config_options(&response).unwrap_or_default());
|
||||||
live.loaded = Some(id.clone());
|
live.loaded = Some(id.clone());
|
||||||
live.commands = None;
|
live.commands = None;
|
||||||
discard_stale(live)?;
|
discard_stale(live)?;
|
||||||
|
|
@ -324,6 +401,7 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
||||||
})?
|
})?
|
||||||
.to_owned();
|
.to_owned();
|
||||||
live.model = stream::session_model(&response);
|
live.model = stream::session_model(&response);
|
||||||
|
live.offer(stream::config_options(&response).unwrap_or_default());
|
||||||
live.loaded = Some(id.clone());
|
live.loaded = Some(id.clone());
|
||||||
live.commands = None;
|
live.commands = None;
|
||||||
write_id(&self.pending_file(), &id)?;
|
write_id(&self.pending_file(), &id)?;
|
||||||
|
|
@ -347,6 +425,7 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
||||||
let cwd = session_cwd(config);
|
let cwd = session_cwd(config);
|
||||||
let (servers, names) = mcp_servers(config)?;
|
let (servers, names) = mcp_servers(config)?;
|
||||||
let (session, created) = self.attach(live, &cwd, &servers).await?;
|
let (session, created) = self.attach(live, &cwd, &servers).await?;
|
||||||
|
choose_model_and_effort(live, &session, config, sink).await?;
|
||||||
let text = prompt_text(config, prompt, created);
|
let text = prompt_text(config, prompt, created);
|
||||||
let params = json!({ "sessionId": session, "prompt": [{ "type": "text", "text": text }] });
|
let params = json!({ "sessionId": session, "prompt": [{ "type": "text", "text": text }] });
|
||||||
discard_stale(live)?;
|
discard_stale(live)?;
|
||||||
|
|
@ -593,6 +672,10 @@ impl<P: CompactionPolicy> Runtime for AcpRuntime<P> {
|
||||||
Some(Canceller(self.cancel.clone()))
|
Some(Canceller(self.cancel.clone()))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn choices(&self) -> Option<Choices> {
|
||||||
|
Some(self.choices.clone())
|
||||||
|
}
|
||||||
|
|
||||||
/// Moves the session file aside; the agent keeps its own copy of the
|
/// Moves the session file aside; the agent keeps its own copy of the
|
||||||
/// session, so nothing is lost. A new session whose first prompt was never
|
/// session, so nothing is lost. A new session whose first prompt was never
|
||||||
/// answered is dropped too.
|
/// answered is dropped too.
|
||||||
|
|
@ -625,9 +708,7 @@ fn deliver(
|
||||||
match incoming {
|
match incoming {
|
||||||
Incoming::Update(params) => {
|
Incoming::Update(params) => {
|
||||||
if params["sessionId"].as_str() == agent.loaded.as_deref() {
|
if params["sessionId"].as_str() == agent.loaded.as_deref() {
|
||||||
if let Some(advertised) = stream::advertised_commands(¶ms["update"]) {
|
keep_advertised(agent, ¶ms["update"]);
|
||||||
agent.commands = Some(advertised);
|
|
||||||
}
|
|
||||||
for event in mapper.push(¶ms["update"]) {
|
for event in mapper.push(¶ms["update"]) {
|
||||||
sink.on_event(&event);
|
sink.on_event(&event);
|
||||||
}
|
}
|
||||||
|
|
@ -646,10 +727,22 @@ fn deliver(
|
||||||
true
|
true
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Keep the commands and config options an `update` for the loaded session
|
||||||
|
/// advertises.
|
||||||
|
fn keep_advertised(agent: &mut Live, update: &Value) {
|
||||||
|
if let Some(advertised) = stream::advertised_commands(update) {
|
||||||
|
agent.commands = Some(advertised);
|
||||||
|
}
|
||||||
|
if let Some(options) = stream::updated_config_options(update) {
|
||||||
|
agent.offer(options);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Drop what the agent sent outside a turn: the history a `session/load`
|
/// Drop what the agent sent outside a turn: the history a `session/load`
|
||||||
/// replays (the caller already has it), and anything that arrived after the
|
/// replays (the caller already has it), and anything that arrived after the
|
||||||
/// previous turn settled, which must not be shown as part of the next one.
|
/// previous turn settled, which must not be shown as part of the next one.
|
||||||
/// The commands the agent advertises for the loaded session are kept.
|
/// The commands and config options the agent advertises for the loaded
|
||||||
|
/// session are kept.
|
||||||
fn discard_stale(live: &mut Live) -> std::result::Result<(), AcpError> {
|
fn discard_stale(live: &mut Live) -> std::result::Result<(), AcpError> {
|
||||||
while let Ok(incoming) = live.conn.incoming.try_recv() {
|
while let Ok(incoming) = live.conn.incoming.try_recv() {
|
||||||
between_turns(live, incoming)?;
|
between_turns(live, incoming)?;
|
||||||
|
|
@ -660,10 +753,8 @@ fn discard_stale(live: &mut Live) -> std::result::Result<(), AcpError> {
|
||||||
fn between_turns(agent: &mut Live, incoming: Incoming) -> std::result::Result<(), AcpError> {
|
fn between_turns(agent: &mut Live, incoming: Incoming) -> std::result::Result<(), AcpError> {
|
||||||
match incoming {
|
match incoming {
|
||||||
Incoming::Update(params) => {
|
Incoming::Update(params) => {
|
||||||
if params["sessionId"].as_str() == agent.loaded.as_deref()
|
if params["sessionId"].as_str() == agent.loaded.as_deref() {
|
||||||
&& let Some(advertised) = stream::advertised_commands(¶ms["update"])
|
keep_advertised(agent, ¶ms["update"]);
|
||||||
{
|
|
||||||
agent.commands = Some(advertised);
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
Incoming::Stdout(line) | Incoming::Stderr(line) => {
|
Incoming::Stdout(line) | Incoming::Stderr(line) => {
|
||||||
|
|
@ -692,6 +783,59 @@ async fn advertises_compact(live: &mut Live) -> std::result::Result<bool, AcpErr
|
||||||
.any(|name| name == COMPACT_COMMAND))
|
.any(|name| name == COMPACT_COMMAND))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Set `config`'s model, then its effort, on `session`: the effort levels
|
||||||
|
/// on offer depend on the model.
|
||||||
|
async fn choose_model_and_effort(
|
||||||
|
live: &mut Live,
|
||||||
|
session: &str,
|
||||||
|
config: &Config,
|
||||||
|
sink: &impl Sink,
|
||||||
|
) -> std::result::Result<(), AcpError> {
|
||||||
|
choose(live, session, MODEL_CATEGORY, config.model.as_deref(), sink).await?;
|
||||||
|
choose(
|
||||||
|
live,
|
||||||
|
session,
|
||||||
|
EFFORT_CATEGORY,
|
||||||
|
config.effort.as_deref(),
|
||||||
|
sink,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Set `session`'s config option of `category` to `wanted`, if the session
|
||||||
|
/// offers that value and holds another. A refusal is reported to `sink`, and
|
||||||
|
/// the turn goes ahead on the value the session holds.
|
||||||
|
async fn choose(
|
||||||
|
live: &mut Live,
|
||||||
|
session: &str,
|
||||||
|
category: &str,
|
||||||
|
wanted: Option<&str>,
|
||||||
|
sink: &impl Sink,
|
||||||
|
) -> std::result::Result<(), AcpError> {
|
||||||
|
let Some(option) = live.choices.find(category) else {
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
let Some(wanted) =
|
||||||
|
wanted.filter(|w| *w != option.current && option.values.iter().any(|v| v == w))
|
||||||
|
else {
|
||||||
|
return Ok(());
|
||||||
|
};
|
||||||
|
let params = json!({ "sessionId": session, "configId": option.id, "value": wanted });
|
||||||
|
match live.conn.request("session/set_config_option", params).await {
|
||||||
|
Ok(response) => {
|
||||||
|
if let Some(options) = stream::config_options(&response) {
|
||||||
|
live.offer(options);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Err(e @ AcpError::Rpc { .. }) => {
|
||||||
|
tracing::warn!(error = %e, category, wanted, "ACP session/set_config_option failed");
|
||||||
|
sink.on_stderr_line(&format!("ACP agent refused {category} {wanted}: {e}"));
|
||||||
|
}
|
||||||
|
Err(e) => return Err(e),
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
/// The prompt as sent: the first one of a session carries the system prompt.
|
/// The prompt as sent: the first one of a session carries the system prompt.
|
||||||
fn prompt_text(config: &Config, prompt: &str, first: bool) -> String {
|
fn prompt_text(config: &Config, prompt: &str, first: bool) -> String {
|
||||||
match (first, &config.system_prompt_file) {
|
match (first, &config.system_prompt_file) {
|
||||||
|
|
@ -758,7 +902,7 @@ mod tests {
|
||||||
|
|
||||||
use hive_claude::{Config, NoopSink, PercentPolicy, Sink};
|
use hive_claude::{Config, NoopSink, PercentPolicy, Sink};
|
||||||
|
|
||||||
use super::{AcpError, AcpRuntime, TurnKind};
|
use super::{AcpError, AcpRuntime, Choice, SessionChoices, TurnKind};
|
||||||
use crate::{AcpCommand, Error, Runtime};
|
use crate::{AcpCommand, Error, Runtime};
|
||||||
|
|
||||||
/// An ACP agent in plain `sh`. It answers `initialize`, numbers its
|
/// An ACP agent in plain `sh`. It answers `initialize`, numbers its
|
||||||
|
|
@ -784,8 +928,12 @@ mod tests {
|
||||||
/// tokens used on each prompt to session `s1`. With `STALL_COMPACT` set, it
|
/// tokens used on each prompt to session `s1`. With `STALL_COMPACT` set, it
|
||||||
/// answers `/compact` only when cancelled, as `silent` does. With `BLANK`
|
/// answers `/compact` only when cancelled, as `silent` does. With `BLANK`
|
||||||
/// set, it answers a prompt whose text is exactly that as `blank` does.
|
/// set, it answers a prompt whose text is exactly that as `blank` does.
|
||||||
|
/// With `OPTIONS` set, its sessions offer models `m/think` (the default)
|
||||||
|
/// and `m/plain`, and effort levels `low` (the default) and `high` on
|
||||||
|
/// `m/think` only; it appends each `session/set_config_option` to
|
||||||
|
/// `<log>.sets` as `id=value`.
|
||||||
const AGENT: &str = r#"
|
const AGENT: &str = r#"
|
||||||
n=0 prompt=
|
n=0 prompt= model=m/think effort=low
|
||||||
printf 'start\n' >> "$1.methods"
|
printf 'start\n' >> "$1.methods"
|
||||||
update() {
|
update() {
|
||||||
printf '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"%s","update":%s}}\n' "$1" "$2"
|
printf '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"%s","update":%s}}\n' "$1" "$2"
|
||||||
|
|
@ -799,6 +947,16 @@ ended() {
|
||||||
bare() {
|
bare() {
|
||||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$1"
|
printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$1"
|
||||||
}
|
}
|
||||||
|
options() {
|
||||||
|
[ -n "$OPTIONS" ] || return 0
|
||||||
|
o="{\"id\":\"model\",\"name\":\"Model\",\"category\":\"model\",\"type\":\"select\",\"currentValue\":\"$model\",\"options\":[{\"value\":\"m/think\",\"name\":\"Think\"},{\"value\":\"m/plain\",\"name\":\"Plain\"}]}"
|
||||||
|
[ "$model" = m/think ] && o="$o,{\"id\":\"effort\",\"name\":\"Effort\",\"category\":\"thought_level\",\"type\":\"select\",\"currentValue\":\"$effort\",\"options\":[{\"value\":\"low\",\"name\":\"Low\"},{\"value\":\"high\",\"name\":\"High\"}]}"
|
||||||
|
printf ',"configOptions":[%s]' "$o"
|
||||||
|
}
|
||||||
|
result() {
|
||||||
|
o=$(options)
|
||||||
|
printf '{"jsonrpc":"2.0","id":%s,"result":{%s}}\n' "$1" "${o#,}"
|
||||||
|
}
|
||||||
while IFS= read -r line; do
|
while IFS= read -r line; do
|
||||||
id=${line#*\"id\":}; id=${id%%[,\}]*}
|
id=${line#*\"id\":}; id=${id%%[,\}]*}
|
||||||
m=${line#*\"method\":\"}; m=${m%%\"*}
|
m=${line#*\"method\":\"}; m=${m%%\"*}
|
||||||
|
|
@ -809,11 +967,17 @@ while IFS= read -r line; do
|
||||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"protocolVersion":1,"agentCapabilities":{"loadSession":true,"mcpCapabilities":{"http":true}}}}\n' "$id" ;;
|
printf '{"jsonrpc":"2.0","id":%s,"result":{"protocolVersion":1,"agentCapabilities":{"loadSession":true,"mcpCapabilities":{"http":true}}}}\n' "$id" ;;
|
||||||
session/new)
|
session/new)
|
||||||
n=$((n+1))
|
n=$((n+1))
|
||||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"sessionId":"s%s"}}\n' "$id" "$n"
|
printf '{"jsonrpc":"2.0","id":%s,"result":{"sessionId":"s%s"%s}}\n' "$id" "$n" "$(options)"
|
||||||
advertise "s$n" ;;
|
advertise "s$n" ;;
|
||||||
session/load)
|
session/load)
|
||||||
printf '{"jsonrpc":"2.0","id":%s,"result":{}}\n' "$id"
|
result "$id"
|
||||||
advertise "$sid" ;;
|
advertise "$sid" ;;
|
||||||
|
session/set_config_option)
|
||||||
|
cid=${line#*\"configId\":\"}; cid=${cid%%\"*}
|
||||||
|
val=${line#*\"value\":\"}; val=${val%%\"*}
|
||||||
|
printf '%s=%s\n' "$cid" "$val" >> "$1.sets"
|
||||||
|
case $cid in model) model=$val ;; effort) effort=$val ;; esac
|
||||||
|
result "$id" ;;
|
||||||
session/prompt)
|
session/prompt)
|
||||||
printf '%s\n' "$line" >> "$1"
|
printf '%s\n' "$line" >> "$1"
|
||||||
[ -n "$USED" ] && [ "$sid" = s1 ] && update "$sid" "{\"sessionUpdate\":\"usage_update\",\"used\":$USED,\"size\":1000}"
|
[ -n "$USED" ] && [ "$sid" = s1 ] && update "$sid" "{\"sessionUpdate\":\"usage_update\",\"used\":$USED,\"size\":1000}"
|
||||||
|
|
@ -1291,4 +1455,60 @@ done
|
||||||
);
|
);
|
||||||
assert_eq!(recorded(dir.path()), "s2");
|
assert_eq!(recorded(dir.path()), "s2");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn choice(current: &str, values: &[&str]) -> Choice {
|
||||||
|
Choice {
|
||||||
|
current: current.to_owned(),
|
||||||
|
values: values.iter().map(|v| (*v).to_owned()).collect(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Each `session/set_config_option` the agent was sent, as `id=value`.
|
||||||
|
fn sets(dir: &Path) -> Vec<String> {
|
||||||
|
std::fs::read_to_string(dir.join("prompts.sets"))
|
||||||
|
.unwrap_or_default()
|
||||||
|
.lines()
|
||||||
|
.map(str::to_owned)
|
||||||
|
.collect()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn model_and_effort_are_set_and_effort_goes_with_a_model_without_it() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let runtime = agent(dir.path(), "ok", &[("OPTIONS", "1")], policy(0));
|
||||||
|
let choices = runtime.choices().unwrap();
|
||||||
|
assert_eq!(choices.get(), SessionChoices::default());
|
||||||
|
let mut config = Config {
|
||||||
|
model: Some("haiku".into()),
|
||||||
|
effort: Some("high".into()),
|
||||||
|
..config(dir.path())
|
||||||
|
};
|
||||||
|
|
||||||
|
runtime.run(&config, "one", &NoopSink).await.unwrap();
|
||||||
|
// `haiku` is not on offer, so only the effort is set.
|
||||||
|
assert_eq!(sets(dir.path()), ["effort=high"]);
|
||||||
|
assert_eq!(
|
||||||
|
choices.get(),
|
||||||
|
SessionChoices {
|
||||||
|
model: Some(choice("m/think", &["m/think", "m/plain"])),
|
||||||
|
effort: Some(choice("high", &["low", "high"])),
|
||||||
|
}
|
||||||
|
);
|
||||||
|
|
||||||
|
config.model = Some("m/plain".into());
|
||||||
|
let done = runtime.run(&config, "two", &NoopSink).await.unwrap();
|
||||||
|
assert_eq!(sets(dir.path()), ["effort=high", "model=m/plain"]);
|
||||||
|
assert_eq!(
|
||||||
|
choices.get(),
|
||||||
|
SessionChoices {
|
||||||
|
model: Some(choice("m/plain", &["m/think", "m/plain"])),
|
||||||
|
effort: None,
|
||||||
|
}
|
||||||
|
);
|
||||||
|
assert_eq!(done.telemetry.model.as_deref(), Some("m/plain"));
|
||||||
|
|
||||||
|
// Values the session already holds are not set again.
|
||||||
|
runtime.run(&config, "three", &NoopSink).await.unwrap();
|
||||||
|
assert_eq!(sets(dir.path()).len(), 2);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -212,6 +212,61 @@ pub(super) fn advertised_commands(update: &Value) -> Option<Vec<String>> {
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// One select config option a session advertises.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub(super) struct ConfigOption {
|
||||||
|
pub(super) id: String,
|
||||||
|
pub(super) category: Option<String>,
|
||||||
|
pub(super) current: String,
|
||||||
|
pub(super) values: Vec<String>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The select options in `value`'s `configOptions`, as a session response, a
|
||||||
|
/// `session/set_config_option` response and a `config_option_update` carry
|
||||||
|
/// them; `None` if it has none. Grouped values are flattened; boolean options
|
||||||
|
/// are left out.
|
||||||
|
pub(super) fn config_options(value: &Value) -> Option<Vec<ConfigOption>> {
|
||||||
|
let options = value.get("configOptions")?.as_array()?;
|
||||||
|
Some(
|
||||||
|
options
|
||||||
|
.iter()
|
||||||
|
.filter_map(|option| {
|
||||||
|
let entries = option.get("options").and_then(Value::as_array);
|
||||||
|
let values = entries
|
||||||
|
.into_iter()
|
||||||
|
.flatten()
|
||||||
|
.flat_map(
|
||||||
|
|entry| match entry.get("options").and_then(Value::as_array) {
|
||||||
|
Some(group) => group.iter().collect(),
|
||||||
|
None => vec![entry],
|
||||||
|
},
|
||||||
|
)
|
||||||
|
.filter_map(|entry| entry.get("value").and_then(Value::as_str))
|
||||||
|
.map(str::to_owned)
|
||||||
|
.collect();
|
||||||
|
Some(ConfigOption {
|
||||||
|
id: option.get("id")?.as_str()?.to_owned(),
|
||||||
|
category: option
|
||||||
|
.get("category")
|
||||||
|
.and_then(Value::as_str)
|
||||||
|
.map(str::to_owned),
|
||||||
|
current: option.get("currentValue")?.as_str()?.to_owned(),
|
||||||
|
values,
|
||||||
|
})
|
||||||
|
})
|
||||||
|
.collect(),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The options a `config_option_update` sets, or `None` for any other
|
||||||
|
/// `update`. Each replaces the session's previous options.
|
||||||
|
pub(super) fn updated_config_options(update: &Value) -> Option<Vec<ConfigOption>> {
|
||||||
|
if update.get("sessionUpdate").and_then(Value::as_str) != Some("config_option_update") {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
Some(config_options(update).unwrap_or_default())
|
||||||
|
}
|
||||||
|
|
||||||
/// Convert a claude `--mcp-config` document into ACP's `mcpServers` list.
|
/// Convert a claude `--mcp-config` document into ACP's `mcpServers` list.
|
||||||
/// Returns the list and the server names, in the same order.
|
/// Returns the list and the server names, in the same order.
|
||||||
pub(super) fn mcp_servers(config: &Value) -> (Vec<Value>, Vec<String>) {
|
pub(super) fn mcp_servers(config: &Value) -> (Vec<Value>, Vec<String>) {
|
||||||
|
|
@ -346,7 +401,9 @@ fn tool_result(id: &str, output: &str, is_error: bool) -> Value {
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{StreamMapper, canonical_tool_name, mcp_servers, session_model};
|
use super::{
|
||||||
|
ConfigOption, StreamMapper, canonical_tool_name, config_options, mcp_servers, session_model,
|
||||||
|
};
|
||||||
use serde_json::{Value, json};
|
use serde_json::{Value, json};
|
||||||
|
|
||||||
fn mapper() -> StreamMapper {
|
fn mapper() -> StreamMapper {
|
||||||
|
|
@ -547,4 +604,31 @@ mod tests {
|
||||||
assert_eq!(session_model(&resp).as_deref(), Some("p/m-2"));
|
assert_eq!(session_model(&resp).as_deref(), Some("p/m-2"));
|
||||||
assert_eq!(session_model(&json!({})), None);
|
assert_eq!(session_model(&json!({})), None);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn config_options_flatten_groups_and_skip_booleans() {
|
||||||
|
let resp = json!({
|
||||||
|
"configOptions": [
|
||||||
|
{ "id": "model", "name": "Model", "category": "model", "type": "select",
|
||||||
|
"currentValue": "p/a",
|
||||||
|
"options": [
|
||||||
|
{ "group": "p", "name": "P", "options": [
|
||||||
|
{ "value": "p/a", "name": "A" }, { "value": "p/b", "name": "B" },
|
||||||
|
] },
|
||||||
|
{ "value": "q/c", "name": "C" },
|
||||||
|
] },
|
||||||
|
{ "id": "fast", "name": "Fast", "type": "boolean", "currentValue": true },
|
||||||
|
],
|
||||||
|
});
|
||||||
|
assert_eq!(
|
||||||
|
config_options(&resp),
|
||||||
|
Some(vec![ConfigOption {
|
||||||
|
id: "model".into(),
|
||||||
|
category: Some("model".into()),
|
||||||
|
current: "p/a".into(),
|
||||||
|
values: vec!["p/a".into(), "p/b".into(), "q/c".into()],
|
||||||
|
}])
|
||||||
|
);
|
||||||
|
assert_eq!(config_options(&json!({})), None);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -4,7 +4,7 @@ use std::path::PathBuf;
|
||||||
|
|
||||||
use hive_claude::{CompactionPolicy, Config, InfiniteSession, Progress, SessionStore, Sink};
|
use hive_claude::{CompactionPolicy, Config, InfiniteSession, Progress, SessionStore, Sink};
|
||||||
|
|
||||||
use crate::{Canceller, Result, Runtime};
|
use crate::{Canceller, Choices, Result, Runtime};
|
||||||
|
|
||||||
/// `claude --print` turns on a titled [`InfiniteSession`], which owns
|
/// `claude --print` turns on a titled [`InfiniteSession`], which owns
|
||||||
/// resume-or-create and compaction.
|
/// resume-or-create and compaction.
|
||||||
|
|
@ -43,4 +43,8 @@ impl<P: CompactionPolicy> Runtime for ClaudeRuntime<P> {
|
||||||
fn canceller(&self) -> Option<Canceller> {
|
fn canceller(&self) -> Option<Canceller> {
|
||||||
None
|
None
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn choices(&self) -> Option<Choices> {
|
||||||
|
None
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -21,7 +21,10 @@ mod acp;
|
||||||
mod claude;
|
mod claude;
|
||||||
mod spec;
|
mod spec;
|
||||||
|
|
||||||
pub use acp::{AcpError, AcpRuntime, Canceller, PermissionAsk, PermissionPolicy};
|
pub use acp::{
|
||||||
|
AcpError, AcpRuntime, Canceller, Choice, Choices, PermissionAsk, PermissionPolicy,
|
||||||
|
SessionChoices,
|
||||||
|
};
|
||||||
pub use claude::ClaudeRuntime;
|
pub use claude::ClaudeRuntime;
|
||||||
pub use hive_claude::{
|
pub use hive_claude::{
|
||||||
CompactionPolicy, Config, PercentPolicy, Progress, SessionStore, Sink, Telemetry, TokenUsage,
|
CompactionPolicy, Config, PercentPolicy, Progress, SessionStore, Sink, Telemetry, TokenUsage,
|
||||||
|
|
@ -57,6 +60,11 @@ pub trait Runtime {
|
||||||
/// `None` for a runtime without one: a claude turn is stopped by
|
/// `None` for a runtime without one: a claude turn is stopped by
|
||||||
/// signalling its `claude` process.
|
/// signalling its `claude` process.
|
||||||
fn canceller(&self) -> Option<Canceller>;
|
fn canceller(&self) -> Option<Canceller>;
|
||||||
|
|
||||||
|
/// A handle reading the model and effort the session lets the client
|
||||||
|
/// pick. `None` for a runtime whose session offers none to read: claude's
|
||||||
|
/// `--model` and `--effort` are passed through unchecked.
|
||||||
|
fn choices(&self) -> Option<Choices>;
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The runtime an agent was configured with, chosen at startup from a
|
/// The runtime an agent was configured with, chosen at startup from a
|
||||||
|
|
@ -94,6 +102,13 @@ impl<P: CompactionPolicy> Runtime for AgentRuntime<P> {
|
||||||
Self::Acp(r) => r.canceller(),
|
Self::Acp(r) => r.canceller(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn choices(&self) -> Option<Choices> {
|
||||||
|
match self {
|
||||||
|
Self::Claude(r) => r.choices(),
|
||||||
|
Self::Acp(r) => r.choices(),
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Why a runtime operation did not complete.
|
/// Why a runtime operation did not complete.
|
||||||
|
|
|
||||||
|
|
@ -82,8 +82,8 @@
|
||||||
//! session id is kept in a per-name file under the harness dir
|
//! session id is kept in a per-name file under the harness dir
|
||||||
//! (`acp_session_file`): the agent process lives for one run, `interrupt`
|
//! (`acp_session_file`): the agent process lives for one run, `interrupt`
|
||||||
//! sends `session/cancel` rather than a signal, a failed turn is `Failed`
|
//! sends `session/cancel` rather than a signal, a failed turn is `Failed`
|
||||||
//! and never `Killed`, and `model`/`effort` are not applied — the ACP agent
|
//! and never `Killed`, and `model`/`effort` apply only as values the ACP
|
||||||
//! runs whatever model it is configured with.
|
//! session offers.
|
||||||
|
|
||||||
use std::collections::HashMap;
|
use std::collections::HashMap;
|
||||||
use std::os::unix::process::ExitStatusExt as _;
|
use std::os::unix::process::ExitStatusExt as _;
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue