hive-runtime: compact ACP sessions through the agent's compact command
An ACP agent's session is now compacted like a claude one: proactively once a turn crosses the percent-of-window watermark, and on the operator's /compact or the agent's compact tool. Before, the ACP backend's compact returned Unsupported and no watermark applied to it. - AcpRuntime takes the same CompactionPolicy as ClaudeRuntime; make_session builds one PercentPolicy (with CHECKPOINT_PROMPT) and hands it to whichever backend runs. - The runtime keeps the commands each session advertises in available_commands_update. If `compact` is among them, compaction sends the prompt `/compact` on the same session, which is how the ACP spec runs an advertised command. A proactive compaction runs the checkpoint turn first, as InfiniteSession does. - With no compact command, the checkpoint turn runs, the session is archived, and the next turn starts a new one with the system prompt. - Error::Unsupported had no producer left, so it and drive_turn's "/compact skipped" arm are gone. Refs #4391
This commit is contained in:
parent
3bef1dfab6
commit
cec35bfbb1
6 changed files with 356 additions and 76 deletions
|
|
@ -161,7 +161,8 @@ model-reported context fill reaches `HIVE_COMPACT_WATERMARK_PERCENT`
|
||||||
for the window on turns the model didn't report one. `0` disables proactive
|
for the window on turns the model didn't report one. `0` disables proactive
|
||||||
compaction (the reactive path always applies). The proactive path is
|
compaction (the reactive path always applies). The proactive path is
|
||||||
best-effort — a failed checkpoint or `/compact` never fails the turn that
|
best-effort — a failed checkpoint or `/compact` never fails the turn that
|
||||||
already succeeded.
|
already succeeded. An ACP agent gets the same policy; how it compacts is in
|
||||||
|
[`hive-runtime/README.md`](../../hive-runtime/README.md#acp-backend-compaction).
|
||||||
|
|
||||||
The operator can force a compaction any time via `POST /api/compact`. It's
|
The operator can force a compaction any time via `POST /api/compact`. It's
|
||||||
**deferred**: the handler sets `Bus::request_compact()` and returns
|
**deferred**: the handler sets `Bus::request_compact()` and returns
|
||||||
|
|
|
||||||
|
|
@ -305,15 +305,16 @@ fn compact_percent() -> u8 {
|
||||||
const ACP_SESSION_FILE: &str = "hyperhive-acp-session";
|
const ACP_SESSION_FILE: &str = "hyperhive-acp-session";
|
||||||
|
|
||||||
/// The agent's durable session, on the runtime its environment selects
|
/// The agent's durable session, on the runtime its environment selects
|
||||||
/// (`hive_runtime::RuntimeSpec`). On claude it is the constant-title
|
/// (`hive_runtime::RuntimeSpec`), with hyperhive's percent-of-window
|
||||||
/// [`hive_claude::InfiniteSession`] with hyperhive's percent-of-window
|
/// compaction policy. On claude it is the constant-title
|
||||||
/// compaction policy. Built once by the serve loop (see [`make_session`]) and
|
/// [`hive_claude::InfiniteSession`]. Built once by the serve loop (see
|
||||||
/// threaded through the turns, rather than rebuilt each time.
|
/// [`make_session`]) and threaded through the turns, rather than rebuilt each
|
||||||
|
/// time.
|
||||||
pub type AgentSession = AgentRuntime<PercentPolicy>;
|
pub type AgentSession = AgentRuntime<PercentPolicy>;
|
||||||
|
|
||||||
/// Construct the agent's durable session. On claude: constant title + on-disk
|
/// Construct the agent's durable session, with a percent-of-window compaction
|
||||||
/// store + a percent-of-window compaction policy that checkpoints
|
/// policy that checkpoints (`CHECKPOINT_PROMPT`) before compacting; on claude
|
||||||
/// (`CHECKPOINT_PROMPT`) before compacting. Called once at serve-loop start.
|
/// also a constant title and the on-disk store. Called once at serve-loop start.
|
||||||
/// `percent` comes from a boot-time env var and `default_window` is only a
|
/// `percent` comes from a boot-time env var and `default_window` is only a
|
||||||
/// fallback for turns where the model didn't report a window, so a single
|
/// fallback for turns where the model didn't report a window, so a single
|
||||||
/// build at startup is fine.
|
/// build at startup is fine.
|
||||||
|
|
@ -322,20 +323,20 @@ pub type AgentSession = AgentRuntime<PercentPolicy>;
|
||||||
///
|
///
|
||||||
/// Returns an error if the runtime selection in the environment is invalid.
|
/// Returns an error if the runtime selection in the environment is invalid.
|
||||||
pub fn make_session(bus: &Bus) -> Result<AgentSession> {
|
pub fn make_session(bus: &Bus) -> Result<AgentSession> {
|
||||||
|
let policy = PercentPolicy {
|
||||||
|
percent: compact_percent(),
|
||||||
|
default_window: Some(effective_context_window(bus)),
|
||||||
|
checkpoint_prompt: Some(CHECKPOINT_PROMPT.to_string()),
|
||||||
|
};
|
||||||
Ok(match RuntimeSpec::from_env()? {
|
Ok(match RuntimeSpec::from_env()? {
|
||||||
RuntimeSpec::Claude => AgentRuntime::Claude(ClaudeRuntime::new(
|
RuntimeSpec::Claude => {
|
||||||
session_title(),
|
AgentRuntime::Claude(ClaudeRuntime::new(session_title(), session_store(), policy))
|
||||||
session_store(),
|
}
|
||||||
PercentPolicy {
|
|
||||||
percent: compact_percent(),
|
|
||||||
default_window: Some(effective_context_window(bus)),
|
|
||||||
checkpoint_prompt: Some(CHECKPOINT_PROMPT.to_string()),
|
|
||||||
},
|
|
||||||
)),
|
|
||||||
RuntimeSpec::Acp(command) => AgentRuntime::Acp(Box::new(AcpRuntime::new(
|
RuntimeSpec::Acp(command) => AgentRuntime::Acp(Box::new(AcpRuntime::new(
|
||||||
command,
|
command,
|
||||||
crate::paths::harness_dir().join(ACP_SESSION_FILE),
|
crate::paths::harness_dir().join(ACP_SESSION_FILE),
|
||||||
acp_permission_policy(),
|
acp_permission_policy(),
|
||||||
|
policy,
|
||||||
))),
|
))),
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
@ -460,15 +461,7 @@ pub async fn drive_turn(
|
||||||
// Reflect `Compacting` in the UI like the idle path (`run_pending_compact`)
|
// Reflect `Compacting` in the UI like the idle path (`run_pending_compact`)
|
||||||
// does; the serve loop resets to `Idle` once this turn returns.
|
// does; the serve loop resets to `Idle` once this turn returns.
|
||||||
bus.set_state(crate::events::TurnState::Compacting);
|
bus.set_state(crate::events::TurnState::Compacting);
|
||||||
let compacted = match session.compact(&config, &sink).await {
|
let _ = session.compact(&config, &sink).await;
|
||||||
Err(e @ hive_runtime::Error::Unsupported(_)) => {
|
|
||||||
bus.emit(LiveEvent::Note {
|
|
||||||
text: format!("/compact skipped: {e}"),
|
|
||||||
});
|
|
||||||
false
|
|
||||||
}
|
|
||||||
_ => true,
|
|
||||||
};
|
|
||||||
// If the compact call asked to be woken (the agent's own `compact`
|
// If the compact call asked to be woken (the agent's own `compact`
|
||||||
// tool with a `wake_prompt`), stash it — the serve loop reads it
|
// tool with a `wake_prompt`), stash it — the serve loop reads it
|
||||||
// back after this turn returns and drives a synthetic follow-up
|
// back after this turn returns and drives a synthetic follow-up
|
||||||
|
|
@ -477,7 +470,7 @@ pub async fn drive_turn(
|
||||||
if let Some(prompt) = request.wake_prompt {
|
if let Some(prompt) = request.wake_prompt {
|
||||||
bus.set_post_compact_wake(prompt);
|
bus.set_post_compact_wake(prompt);
|
||||||
}
|
}
|
||||||
return Ok(compacted);
|
return Ok(true);
|
||||||
}
|
}
|
||||||
outcome
|
outcome
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -40,6 +40,15 @@ with `AcpError::IdleTimeout` (`IdleKilled` if the agent was killed). This is
|
||||||
the only way a provider error the agent retries without reporting it, such as
|
the only way a provider error the agent retries without reporting it, such as
|
||||||
an HTTP 429, ends the turn.
|
an HTTP 429, ends the turn.
|
||||||
|
|
||||||
## ACP backend: not yet
|
## ACP backend: compaction
|
||||||
|
|
||||||
`compact` returns `Error::Unsupported`.
|
The same `CompactionPolicy` as the claude backend decides when: after a turn
|
||||||
|
past the watermark, or on `compact` (the operator's `/compact`, the agent's
|
||||||
|
`compact` tool).
|
||||||
|
|
||||||
|
- If the agent advertises a `compact` command (`available_commands_update`),
|
||||||
|
it is sent as the prompt `/compact` on the same session, as the ACP spec
|
||||||
|
runs any advertised command. A proactive compaction runs the policy's
|
||||||
|
checkpoint turn first, as on claude.
|
||||||
|
- Otherwise the policy's checkpoint turn runs, then the session is archived,
|
||||||
|
and the next turn starts a new one, carrying the system prompt again.
|
||||||
|
|
|
||||||
|
|
@ -11,7 +11,7 @@ use std::sync::Arc;
|
||||||
use std::sync::atomic::{AtomicBool, Ordering};
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
use hive_claude::{Config, Progress, Sink};
|
use hive_claude::{CompactionPolicy, Config, Progress, Sink};
|
||||||
use serde_json::{Value, json};
|
use serde_json::{Value, json};
|
||||||
use tokio::sync::futures::Notified;
|
use tokio::sync::futures::Notified;
|
||||||
use tokio::sync::{Mutex, Notify};
|
use tokio::sync::{Mutex, Notify};
|
||||||
|
|
@ -37,6 +37,14 @@ const SETTLE_MAX: Duration = Duration::from_secs(5);
|
||||||
/// killed.
|
/// killed.
|
||||||
const CANCEL_GRACE: Duration = Duration::from_secs(10);
|
const CANCEL_GRACE: Duration = Duration::from_secs(10);
|
||||||
|
|
||||||
|
/// The command an agent advertises to compact its session.
|
||||||
|
const COMPACT_COMMAND: &str = "compact";
|
||||||
|
|
||||||
|
/// How long to wait, on a session whose agent has not advertised its
|
||||||
|
/// commands yet, for it to do so. An agent may send them only after the
|
||||||
|
/// session is attached, or never.
|
||||||
|
const COMMANDS_WAIT: Duration = SETTLE_MAX;
|
||||||
|
|
||||||
/// 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> {
|
||||||
|
|
@ -155,10 +163,17 @@ enum Stop {
|
||||||
/// `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.
|
||||||
pub struct AcpRuntime {
|
///
|
||||||
|
/// Compaction, proactive once `policy` says so after a turn or on
|
||||||
|
/// [`Runtime::compact`], runs the agent's advertised `compact` command as a
|
||||||
|
/// prompt on the same session. An agent advertising none gets the policy's
|
||||||
|
/// checkpoint turn instead, and the session is archived, so the next turn
|
||||||
|
/// starts a new one.
|
||||||
|
pub struct AcpRuntime<P: CompactionPolicy> {
|
||||||
command: AcpCommand,
|
command: AcpCommand,
|
||||||
session_file: PathBuf,
|
session_file: PathBuf,
|
||||||
permit: PermissionPolicy,
|
permit: PermissionPolicy,
|
||||||
|
policy: P,
|
||||||
live: Mutex<Option<Live>>,
|
live: Mutex<Option<Live>>,
|
||||||
cancel: Arc<CancelState>,
|
cancel: Arc<CancelState>,
|
||||||
cancel_grace: Duration,
|
cancel_grace: Duration,
|
||||||
|
|
@ -171,17 +186,27 @@ struct Live {
|
||||||
/// The recorded session, or a new one whose first prompt is unanswered.
|
/// The recorded session, or a new one whose first prompt is unanswered.
|
||||||
loaded: Option<String>,
|
loaded: Option<String>,
|
||||||
model: Option<String>,
|
model: Option<String>,
|
||||||
|
/// The commands the agent last advertised for `loaded`; `None` until it
|
||||||
|
/// has.
|
||||||
|
commands: Option<Vec<String>>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl AcpRuntime {
|
impl<P: CompactionPolicy> AcpRuntime<P> {
|
||||||
/// A runtime spawning `command`, keeping its session id in
|
/// A runtime spawning `command`, keeping its session id in
|
||||||
/// `session_file`, and answering its permission requests with `permit`.
|
/// `session_file`, answering its permission requests with `permit`, and
|
||||||
|
/// compacting per `policy`.
|
||||||
#[must_use]
|
#[must_use]
|
||||||
pub fn new(command: AcpCommand, session_file: PathBuf, permit: PermissionPolicy) -> Self {
|
pub fn new(
|
||||||
|
command: AcpCommand,
|
||||||
|
session_file: PathBuf,
|
||||||
|
permit: PermissionPolicy,
|
||||||
|
policy: P,
|
||||||
|
) -> Self {
|
||||||
Self {
|
Self {
|
||||||
command,
|
command,
|
||||||
session_file,
|
session_file,
|
||||||
permit,
|
permit,
|
||||||
|
policy,
|
||||||
live: Mutex::new(None),
|
live: Mutex::new(None),
|
||||||
cancel: Arc::default(),
|
cancel: Arc::default(),
|
||||||
cancel_grace: CANCEL_GRACE,
|
cancel_grace: CANCEL_GRACE,
|
||||||
|
|
@ -223,6 +248,7 @@ impl AcpRuntime {
|
||||||
load_session: caps["loadSession"] == Value::Bool(true),
|
load_session: caps["loadSession"] == Value::Bool(true),
|
||||||
loaded: None,
|
loaded: None,
|
||||||
model: None,
|
model: None,
|
||||||
|
commands: None,
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -253,9 +279,10 @@ impl AcpRuntime {
|
||||||
let params = json!({ "sessionId": id, "cwd": cwd, "mcpServers": servers });
|
let params = json!({ "sessionId": id, "cwd": cwd, "mcpServers": servers });
|
||||||
match live.conn.request("session/load", params).await {
|
match live.conn.request("session/load", params).await {
|
||||||
Ok(response) => {
|
Ok(response) => {
|
||||||
discard_stale(&mut live.conn)?;
|
|
||||||
live.model = stream::session_model(&response);
|
live.model = stream::session_model(&response);
|
||||||
live.loaded = Some(id.clone());
|
live.loaded = Some(id.clone());
|
||||||
|
live.commands = None;
|
||||||
|
discard_stale(live)?;
|
||||||
return Ok((id, new));
|
return Ok((id, new));
|
||||||
}
|
}
|
||||||
Err(e @ AcpError::Rpc { .. }) => {
|
Err(e @ AcpError::Rpc { .. }) => {
|
||||||
|
|
@ -275,6 +302,7 @@ impl AcpRuntime {
|
||||||
.to_owned();
|
.to_owned();
|
||||||
live.model = stream::session_model(&response);
|
live.model = stream::session_model(&response);
|
||||||
live.loaded = Some(id.clone());
|
live.loaded = Some(id.clone());
|
||||||
|
live.commands = None;
|
||||||
write_id(&self.pending_file(), &id)?;
|
write_id(&self.pending_file(), &id)?;
|
||||||
Ok((id, true))
|
Ok((id, true))
|
||||||
}
|
}
|
||||||
|
|
@ -297,7 +325,7 @@ impl AcpRuntime {
|
||||||
let (session, created) = self.attach(live, &cwd, &servers).await?;
|
let (session, created) = self.attach(live, &cwd, &servers).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(&mut live.conn)?;
|
discard_stale(live)?;
|
||||||
let mut reply = live.conn.send("session/prompt", params).await?;
|
let mut reply = live.conn.send("session/prompt", params).await?;
|
||||||
let mut mapper = StreamMapper::new(names);
|
let mut mapper = StreamMapper::new(names);
|
||||||
let mut stderr_tail = VecDeque::new();
|
let mut stderr_tail = VecDeque::new();
|
||||||
|
|
@ -315,7 +343,7 @@ impl AcpRuntime {
|
||||||
if matches!(incoming, Incoming::Update(_)) {
|
if matches!(incoming, Incoming::Update(_)) {
|
||||||
quiet_until = idle.map(|window| Instant::now() + window);
|
quiet_until = idle.map(|window| Instant::now() + window);
|
||||||
}
|
}
|
||||||
if !deliver(incoming, &session, &mut mapper, &mut stderr_tail, sink) {
|
if !deliver(incoming, &mut mapper, live, &mut stderr_tail, sink) {
|
||||||
let stderr_tail = Vec::from(stderr_tail).join("\n");
|
let stderr_tail = Vec::from(stderr_tail).join("\n");
|
||||||
return Err(AcpError::Exited { stderr_tail }.into());
|
return Err(AcpError::Exited { stderr_tail }.into());
|
||||||
}
|
}
|
||||||
|
|
@ -360,7 +388,7 @@ impl AcpRuntime {
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
{
|
{
|
||||||
if !deliver(incoming, &session, &mut mapper, &mut stderr_tail, sink) {
|
if !deliver(incoming, &mut mapper, live, &mut stderr_tail, sink) {
|
||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -391,17 +419,33 @@ impl AcpRuntime {
|
||||||
progress.telemetry = mapper.telemetry(&response, live.model.as_deref());
|
progress.telemetry = mapper.telemetry(&response, live.model.as_deref());
|
||||||
Ok(progress)
|
Ok(progress)
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
impl Runtime for AcpRuntime {
|
/// The running agent, spawned if there is none or it has exited.
|
||||||
async fn run(&self, config: &Config, prompt: &str, sink: &impl Sink) -> Result<Progress> {
|
async fn running<'g>(
|
||||||
let mut guard = self.live.lock().await;
|
&self,
|
||||||
if let Some(live) = guard.as_mut()
|
guard: &'g mut Option<Live>,
|
||||||
&& live.conn.exited()
|
config: &Config,
|
||||||
{
|
) -> Result<&'g mut Live> {
|
||||||
|
let mut live = guard.take();
|
||||||
|
if live.as_mut().is_some_and(|live| live.conn.exited()) {
|
||||||
tracing::warn!("ACP agent has exited; respawning");
|
tracing::warn!("ACP agent has exited; respawning");
|
||||||
*guard = None;
|
live = None;
|
||||||
}
|
}
|
||||||
|
let live = match live {
|
||||||
|
Some(live) => live,
|
||||||
|
None => self.start(config).await?,
|
||||||
|
};
|
||||||
|
Ok(guard.insert(live))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Run `prompt` as one turn on the durable session.
|
||||||
|
async fn prompt(
|
||||||
|
&self,
|
||||||
|
guard: &mut Option<Live>,
|
||||||
|
config: &Config,
|
||||||
|
prompt: &str,
|
||||||
|
sink: &impl Sink,
|
||||||
|
) -> Result<Progress> {
|
||||||
// Registered before the turn is marked in flight, so a cancel is
|
// Registered before the turn is marked in flight, so a cancel is
|
||||||
// never missed; one that comes before the prompt is sent stops it
|
// never missed; one that comes before the prompt is sent stops it
|
||||||
// right after.
|
// right after.
|
||||||
|
|
@ -410,10 +454,7 @@ impl Runtime for AcpRuntime {
|
||||||
cancelled.as_mut().enable();
|
cancelled.as_mut().enable();
|
||||||
self.cancel.in_turn.store(true, Ordering::SeqCst);
|
self.cancel.in_turn.store(true, Ordering::SeqCst);
|
||||||
let _in_turn = InTurn(&self.cancel.in_turn);
|
let _in_turn = InTurn(&self.cancel.in_turn);
|
||||||
let live = match guard.as_mut() {
|
let live = self.running(guard, config).await?;
|
||||||
Some(live) => live,
|
|
||||||
None => guard.insert(self.start(config).await?),
|
|
||||||
};
|
|
||||||
let result = self.turn(live, config, prompt, sink, cancelled).await;
|
let result = self.turn(live, config, prompt, sink, cancelled).await;
|
||||||
if let Err(Error::Acp(
|
if let Err(Error::Acp(
|
||||||
AcpError::Closed
|
AcpError::Closed
|
||||||
|
|
@ -429,8 +470,65 @@ impl Runtime for AcpRuntime {
|
||||||
result
|
result
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn compact(&self, _config: &Config, _sink: &impl Sink) -> Result<()> {
|
/// Compact the recorded session, if there is one. `checkpoint` runs the
|
||||||
Err(Error::Unsupported("compact"))
|
/// policy's checkpoint turn before an advertised `compact` command too;
|
||||||
|
/// without one it always runs, since archiving keeps nothing of the
|
||||||
|
/// session.
|
||||||
|
async fn compact_session(
|
||||||
|
&self,
|
||||||
|
guard: &mut Option<Live>,
|
||||||
|
config: &Config,
|
||||||
|
sink: &impl Sink,
|
||||||
|
checkpoint: bool,
|
||||||
|
) -> Result<()> {
|
||||||
|
if read_id(&self.session_file).is_none() {
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
let live = self.running(guard, config).await?;
|
||||||
|
let (servers, _) = mcp_servers(config)?;
|
||||||
|
self.attach(live, &session_cwd(config), &servers).await?;
|
||||||
|
if advertises_compact(live).await? {
|
||||||
|
if checkpoint {
|
||||||
|
self.checkpoint(guard, config, sink).await;
|
||||||
|
}
|
||||||
|
self.prompt(guard, config, &format!("/{COMPACT_COMMAND}"), sink)
|
||||||
|
.await?;
|
||||||
|
return Ok(());
|
||||||
|
}
|
||||||
|
tracing::info!("ACP agent advertises no compact command; starting a new session");
|
||||||
|
self.checkpoint(guard, config, sink).await;
|
||||||
|
self.archive()?;
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The policy's checkpoint turn, if it has one. Best-effort: compaction
|
||||||
|
/// goes ahead if it fails.
|
||||||
|
async fn checkpoint(&self, guard: &mut Option<Live>, config: &Config, sink: &impl Sink) {
|
||||||
|
let Some(prompt) = self.policy.checkpoint_prompt() else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
if let Err(e) = self.prompt(guard, config, prompt, sink).await {
|
||||||
|
tracing::warn!(error = %e, "ACP checkpoint turn failed");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<P: CompactionPolicy> Runtime for AcpRuntime<P> {
|
||||||
|
async fn run(&self, config: &Config, prompt: &str, sink: &impl Sink) -> Result<Progress> {
|
||||||
|
let mut guard = self.live.lock().await;
|
||||||
|
let mut progress = self.prompt(&mut guard, config, prompt, sink).await?;
|
||||||
|
if self.policy.should_compact(progress.telemetry.usage()) {
|
||||||
|
match self.compact_session(&mut guard, config, sink, true).await {
|
||||||
|
Ok(()) => progress.compacted = true,
|
||||||
|
Err(e) => tracing::warn!(error = %e, "ACP compaction failed"),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Ok(progress)
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn compact(&self, config: &Config, sink: &impl Sink) -> Result<()> {
|
||||||
|
let mut guard = self.live.lock().await;
|
||||||
|
self.compact_session(&mut guard, config, sink, false).await
|
||||||
}
|
}
|
||||||
|
|
||||||
fn canceller(&self) -> Option<Canceller> {
|
fn canceller(&self) -> Option<Canceller> {
|
||||||
|
|
@ -456,18 +554,22 @@ impl Runtime for AcpRuntime {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Hand one incoming item to the sink. Returns `false` once the agent has
|
/// Hand one incoming item to the sink, and keep the commands the agent
|
||||||
|
/// advertises for the loaded session. Returns `false` once the agent has
|
||||||
/// closed its output.
|
/// closed its output.
|
||||||
fn deliver(
|
fn deliver(
|
||||||
incoming: Incoming,
|
incoming: Incoming,
|
||||||
session: &str,
|
|
||||||
mapper: &mut StreamMapper,
|
mapper: &mut StreamMapper,
|
||||||
|
agent: &mut Live,
|
||||||
stderr_tail: &mut VecDeque<String>,
|
stderr_tail: &mut VecDeque<String>,
|
||||||
sink: &impl Sink,
|
sink: &impl Sink,
|
||||||
) -> bool {
|
) -> bool {
|
||||||
match incoming {
|
match incoming {
|
||||||
Incoming::Update(params) => {
|
Incoming::Update(params) => {
|
||||||
if params["sessionId"].as_str() == Some(session) {
|
if params["sessionId"].as_str() == agent.loaded.as_deref() {
|
||||||
|
if let Some(advertised) = stream::advertised_commands(¶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);
|
||||||
}
|
}
|
||||||
|
|
@ -489,19 +591,49 @@ fn deliver(
|
||||||
/// 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.
|
||||||
fn discard_stale(conn: &mut Connection) -> std::result::Result<(), AcpError> {
|
/// The commands the agent advertises for the loaded session are kept.
|
||||||
while let Ok(incoming) = conn.incoming.try_recv() {
|
fn discard_stale(live: &mut Live) -> std::result::Result<(), AcpError> {
|
||||||
match incoming {
|
while let Ok(incoming) = live.conn.incoming.try_recv() {
|
||||||
Incoming::Update(_) => {}
|
between_turns(live, incoming)?;
|
||||||
Incoming::Stdout(line) | Incoming::Stderr(line) => {
|
|
||||||
tracing::info!(line = %line, "ACP agent output between turns");
|
|
||||||
}
|
|
||||||
Incoming::Closed => return Err(AcpError::Closed),
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn between_turns(agent: &mut Live, incoming: Incoming) -> std::result::Result<(), AcpError> {
|
||||||
|
match incoming {
|
||||||
|
Incoming::Update(params) => {
|
||||||
|
if params["sessionId"].as_str() == agent.loaded.as_deref()
|
||||||
|
&& let Some(advertised) = stream::advertised_commands(¶ms["update"])
|
||||||
|
{
|
||||||
|
agent.commands = Some(advertised);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Incoming::Stdout(line) | Incoming::Stderr(line) => {
|
||||||
|
tracing::info!(line = %line, "ACP agent output between turns");
|
||||||
|
}
|
||||||
|
Incoming::Closed => return Err(AcpError::Closed),
|
||||||
|
}
|
||||||
|
Ok(())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Whether the agent advertises a `compact` command for the loaded session,
|
||||||
|
/// waiting up to [`COMMANDS_WAIT`] for its commands if it has sent none yet.
|
||||||
|
async fn advertises_compact(live: &mut Live) -> std::result::Result<bool, AcpError> {
|
||||||
|
discard_stale(live)?;
|
||||||
|
let deadline = Instant::now() + COMMANDS_WAIT;
|
||||||
|
while live.commands.is_none() {
|
||||||
|
match tokio::time::timeout_at(deadline, live.conn.incoming.recv()).await {
|
||||||
|
Ok(incoming) => between_turns(live, incoming.unwrap_or(Incoming::Closed))?,
|
||||||
|
Err(_) => break,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Ok(live
|
||||||
|
.commands
|
||||||
|
.iter()
|
||||||
|
.flatten()
|
||||||
|
.any(|name| name == COMPACT_COMMAND))
|
||||||
|
}
|
||||||
|
|
||||||
/// 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) {
|
||||||
|
|
@ -566,7 +698,7 @@ mod tests {
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
use std::time::Duration;
|
use std::time::Duration;
|
||||||
|
|
||||||
use hive_claude::{Config, NoopSink, Sink};
|
use hive_claude::{Config, NoopSink, PercentPolicy, Sink};
|
||||||
|
|
||||||
use super::{AcpError, AcpRuntime};
|
use super::{AcpError, AcpRuntime};
|
||||||
use crate::{AcpCommand, Error, Runtime};
|
use crate::{AcpCommand, Error, Runtime};
|
||||||
|
|
@ -583,24 +715,39 @@ mod tests {
|
||||||
/// reply an agent gives a cancelled prompt it was retrying against a
|
/// reply an agent gives a cancelled prompt it was retrying against a
|
||||||
/// rate-limited (HTTP 429) provider without reporting it;
|
/// rate-limited (HTTP 429) provider without reporting it;
|
||||||
/// - `deaf`: nothing, and `session/cancel` is ignored;
|
/// - `deaf`: nothing, and `session/cancel` is ignored;
|
||||||
/// - `trickle`: five text chunks 100ms apart, then `end_turn`.
|
/// - `trickle`: five text chunks 100ms apart, then `end_turn`;
|
||||||
|
/// - `ok`: `end_turn`, like the rest.
|
||||||
|
///
|
||||||
|
/// With `COMMANDS` set, it advertises that one command on each session it
|
||||||
|
/// creates or loads. With `USED` set, it reports that many of 1000 context
|
||||||
|
/// tokens used on each prompt.
|
||||||
const AGENT: &str = r#"
|
const AGENT: &str = r#"
|
||||||
n=0 prompt=
|
n=0 prompt=
|
||||||
printf 'start\n' >> "$1.methods"
|
printf 'start\n' >> "$1.methods"
|
||||||
|
update() {
|
||||||
|
printf '{"jsonrpc":"2.0","method":"session/update","params":{"sessionId":"%s","update":%s}}\n' "$1" "$2"
|
||||||
|
}
|
||||||
|
advertise() {
|
||||||
|
[ -n "$COMMANDS" ] && update "$1" "{\"sessionUpdate\":\"available_commands_update\",\"availableCommands\":[{\"name\":\"$COMMANDS\",\"description\":\"\"}]}"
|
||||||
|
}
|
||||||
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%%\"*}
|
||||||
|
sid=${line#*\"sessionId\":\"}; sid=${sid%%\"*}
|
||||||
printf '%s\n' "$m" >> "$1.methods"
|
printf '%s\n' "$m" >> "$1.methods"
|
||||||
case $m in
|
case $m in
|
||||||
initialize)
|
initialize)
|
||||||
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"}}\n' "$id" "$n"
|
||||||
|
advertise "s$n" ;;
|
||||||
session/load)
|
session/load)
|
||||||
printf '{"jsonrpc":"2.0","id":%s,"result":{}}\n' "$id" ;;
|
printf '{"jsonrpc":"2.0","id":%s,"result":{}}\n' "$id"
|
||||||
|
advertise "$sid" ;;
|
||||||
session/prompt)
|
session/prompt)
|
||||||
printf '%s\n' "$line" >> "$1"
|
printf '%s\n' "$line" >> "$1"
|
||||||
|
[ -n "$USED" ] && update "$sid" "{\"sessionUpdate\":\"usage_update\",\"used\":$USED,\"size\":1000}"
|
||||||
if [ -e "$1.once" ]; then
|
if [ -e "$1.once" ]; then
|
||||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$id"
|
printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$id"
|
||||||
else
|
else
|
||||||
|
|
@ -608,6 +755,8 @@ while IFS= read -r line; do
|
||||||
case $2 in
|
case $2 in
|
||||||
fail)
|
fail)
|
||||||
printf '{"jsonrpc":"2.0","id":%s,"error":{"code":-32603,"message":"provider down"}}\n' "$id" ;;
|
printf '{"jsonrpc":"2.0","id":%s,"error":{"code":-32603,"message":"provider down"}}\n' "$id" ;;
|
||||||
|
ok)
|
||||||
|
printf '{"jsonrpc":"2.0","id":%s,"result":{"stopReason":"end_turn"}}\n' "$id" ;;
|
||||||
silent) prompt=$id ;;
|
silent) prompt=$id ;;
|
||||||
trickle)
|
trickle)
|
||||||
for _ in 1 2 3 4 5; do
|
for _ in 1 2 3 4 5; do
|
||||||
|
|
@ -626,7 +775,16 @@ while IFS= read -r line; do
|
||||||
done
|
done
|
||||||
"#;
|
"#;
|
||||||
|
|
||||||
fn runtime(dir: &Path, mode: &str) -> AcpRuntime {
|
fn runtime(dir: &Path, mode: &str) -> AcpRuntime<PercentPolicy> {
|
||||||
|
agent(dir, mode, &[], PercentPolicy::default())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn agent(
|
||||||
|
dir: &Path,
|
||||||
|
mode: &str,
|
||||||
|
env: &[(&str, &str)],
|
||||||
|
policy: PercentPolicy,
|
||||||
|
) -> AcpRuntime<PercentPolicy> {
|
||||||
let script = dir.join("agent.sh");
|
let script = dir.join("agent.sh");
|
||||||
std::fs::write(&script, AGENT).unwrap();
|
std::fs::write(&script, AGENT).unwrap();
|
||||||
let command = AcpCommand {
|
let command = AcpCommand {
|
||||||
|
|
@ -636,9 +794,43 @@ done
|
||||||
dir.join("prompts").display().to_string(),
|
dir.join("prompts").display().to_string(),
|
||||||
mode.to_owned(),
|
mode.to_owned(),
|
||||||
],
|
],
|
||||||
env: BTreeMap::new(),
|
env: env
|
||||||
|
.iter()
|
||||||
|
.map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
|
||||||
|
.collect::<BTreeMap<_, _>>(),
|
||||||
};
|
};
|
||||||
AcpRuntime::new(command, dir.join("session"), Arc::new(|_: &_| false))
|
AcpRuntime::new(
|
||||||
|
command,
|
||||||
|
dir.join("session"),
|
||||||
|
Arc::new(|_: &_| false),
|
||||||
|
policy,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Compacts at `percent` of the window, after a `CHECKPOINT` turn.
|
||||||
|
fn policy(percent: u8) -> PercentPolicy {
|
||||||
|
PercentPolicy {
|
||||||
|
percent,
|
||||||
|
default_window: None,
|
||||||
|
checkpoint_prompt: Some("CHECKPOINT".to_owned()),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Each prompt the agent was sent: its session and its text.
|
||||||
|
fn prompts(dir: &Path) -> Vec<(String, String)> {
|
||||||
|
std::fs::read_to_string(dir.join("prompts"))
|
||||||
|
.unwrap()
|
||||||
|
.lines()
|
||||||
|
.map(|line| {
|
||||||
|
let request: serde_json::Value = serde_json::from_str(line).unwrap();
|
||||||
|
let params = &request["params"];
|
||||||
|
let text = params["prompt"][0]["text"].as_str().unwrap();
|
||||||
|
(
|
||||||
|
params["sessionId"].as_str().unwrap().to_owned(),
|
||||||
|
text.to_owned(),
|
||||||
|
)
|
||||||
|
})
|
||||||
|
.collect()
|
||||||
}
|
}
|
||||||
|
|
||||||
fn config(dir: &Path) -> Config {
|
fn config(dir: &Path) -> Config {
|
||||||
|
|
@ -838,4 +1030,75 @@ done
|
||||||
runtime.run(&config, "two", &NoopSink).await.unwrap();
|
runtime.run(&config, "two", &NoopSink).await.unwrap();
|
||||||
assert_eq!(count(dir.path(), "start"), 2);
|
assert_eq!(count(dir.path(), "start"), 2);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn sent(session: &str, text: &str) -> (String, String) {
|
||||||
|
(session.to_owned(), text.to_owned())
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn an_advertised_compact_command_runs_on_the_same_session() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let runtime = agent(dir.path(), "ok", &[("COMMANDS", "compact")], policy(0));
|
||||||
|
let config = config(dir.path());
|
||||||
|
|
||||||
|
runtime.run(&config, "one", &NoopSink).await.unwrap();
|
||||||
|
runtime.compact(&config, &NoopSink).await.unwrap();
|
||||||
|
let next = runtime.run(&config, "two", &NoopSink).await.unwrap();
|
||||||
|
|
||||||
|
assert!(!next.created);
|
||||||
|
assert_eq!(
|
||||||
|
prompts(dir.path()),
|
||||||
|
[
|
||||||
|
sent("s1", "SYSTEM PROMPT\n\none"),
|
||||||
|
sent("s1", "/compact"),
|
||||||
|
sent("s1", "two"),
|
||||||
|
]
|
||||||
|
);
|
||||||
|
assert_eq!(count(dir.path(), "session/new"), 1);
|
||||||
|
assert_eq!(recorded(dir.path()), "s1");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn with_no_compact_command_notes_are_written_then_a_new_session_starts() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let runtime = agent(dir.path(), "ok", &[("COMMANDS", "init")], policy(0));
|
||||||
|
let config = config(dir.path());
|
||||||
|
|
||||||
|
runtime.run(&config, "one", &NoopSink).await.unwrap();
|
||||||
|
runtime.compact(&config, &NoopSink).await.unwrap();
|
||||||
|
let next = runtime.run(&config, "two", &NoopSink).await.unwrap();
|
||||||
|
|
||||||
|
assert!(next.created);
|
||||||
|
assert_eq!(
|
||||||
|
prompts(dir.path()),
|
||||||
|
[
|
||||||
|
sent("s1", "SYSTEM PROMPT\n\none"),
|
||||||
|
sent("s1", "CHECKPOINT"),
|
||||||
|
sent("s2", "SYSTEM PROMPT\n\ntwo"),
|
||||||
|
]
|
||||||
|
);
|
||||||
|
assert_eq!(recorded(dir.path()), "s2");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn a_turn_past_the_watermark_checkpoints_then_compacts() {
|
||||||
|
let dir = tempfile::tempdir().unwrap();
|
||||||
|
let env = [("COMMANDS", "compact"), ("USED", "800")];
|
||||||
|
let runtime = agent(dir.path(), "ok", &env, policy(75));
|
||||||
|
|
||||||
|
let done = runtime
|
||||||
|
.run(&config(dir.path()), "one", &NoopSink)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
assert!(done.compacted);
|
||||||
|
assert_eq!(
|
||||||
|
prompts(dir.path()),
|
||||||
|
[
|
||||||
|
sent("s1", "SYSTEM PROMPT\n\none"),
|
||||||
|
sent("s1", "CHECKPOINT"),
|
||||||
|
sent("s1", "/compact"),
|
||||||
|
]
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -185,6 +185,23 @@ pub(super) fn session_model(response: &Value) -> Option<String> {
|
||||||
.map(str::to_owned)
|
.map(str::to_owned)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The command names an `available_commands_update` advertises, or `None`
|
||||||
|
/// for any other `update`. Each replaces the session's previous list.
|
||||||
|
pub(super) fn advertised_commands(update: &Value) -> Option<Vec<String>> {
|
||||||
|
if update.get("sessionUpdate").and_then(Value::as_str) != Some("available_commands_update") {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
let commands = update.get("availableCommands").and_then(Value::as_array);
|
||||||
|
Some(
|
||||||
|
commands
|
||||||
|
.into_iter()
|
||||||
|
.flatten()
|
||||||
|
.filter_map(|c| c.get("name").and_then(Value::as_str))
|
||||||
|
.map(str::to_owned)
|
||||||
|
.collect(),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
/// 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>) {
|
||||||
|
|
|
||||||
|
|
@ -44,8 +44,8 @@ pub trait Runtime {
|
||||||
sink: &impl Sink,
|
sink: &impl Sink,
|
||||||
) -> impl Future<Output = Result<Progress>>;
|
) -> impl Future<Output = Result<Progress>>;
|
||||||
|
|
||||||
/// Compact the durable session now. A backend that cannot returns
|
/// Compact the durable session now. With no session there is nothing to
|
||||||
/// [`Error::Unsupported`].
|
/// compact, and this succeeds without doing anything.
|
||||||
fn compact(&self, config: &Config, sink: &impl Sink) -> impl Future<Output = Result<()>>;
|
fn compact(&self, config: &Config, sink: &impl Sink) -> impl Future<Output = Result<()>>;
|
||||||
|
|
||||||
/// Set the durable session aside so the next turn starts a fresh one.
|
/// Set the durable session aside so the next turn starts a fresh one.
|
||||||
|
|
@ -63,7 +63,7 @@ pub trait Runtime {
|
||||||
/// [`RuntimeSpec`].
|
/// [`RuntimeSpec`].
|
||||||
pub enum AgentRuntime<P: CompactionPolicy> {
|
pub enum AgentRuntime<P: CompactionPolicy> {
|
||||||
Claude(ClaudeRuntime<P>),
|
Claude(ClaudeRuntime<P>),
|
||||||
Acp(Box<AcpRuntime>),
|
Acp(Box<AcpRuntime<P>>),
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<P: CompactionPolicy> Runtime for AgentRuntime<P> {
|
impl<P: CompactionPolicy> Runtime for AgentRuntime<P> {
|
||||||
|
|
@ -106,9 +106,6 @@ pub enum Error {
|
||||||
/// From the ACP backend.
|
/// From the ACP backend.
|
||||||
#[error(transparent)]
|
#[error(transparent)]
|
||||||
Acp(#[from] AcpError),
|
Acp(#[from] AcpError),
|
||||||
/// The backend does not implement this operation.
|
|
||||||
#[error("{0} is not supported by this runtime")]
|
|
||||||
Unsupported(&'static str),
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Result alias for this crate.
|
/// Result alias for this crate.
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue