hive-agent: export ACP-reported cost and context fill over OTLP
An ACP agent's `usage_update` carries `cost.{amount,currency}`, the
session's running total (opencode sums every assistant message in the
session). hive-runtime now reads it and turns the running total into
what each report added: a new session counts from zero, a session loaded
into a freshly started agent only baselines on its first report, and a
falling total adds nothing. The spend is held on the runtime until
`Runtime::take_reported_cost` drains it; claude's runtime reports none,
since the claude binary already exports `claude_code.cost.usage`.
hive-agent's existing turn-metrics meter records three new instruments:
- `hyperhive.agent.cost.usage` (counter, `model` + `currency`), ACP only;
- `hyperhive.agent.context.used` / `.size` (gauges, no attributes), for
every backend: the two numbers the web UI's ctx% divides.
The `hyperhive · agents` dashboard gets ACP cost panels on its cost tab
and a context-fill panel on its health tab.
Refs #4845
This commit is contained in:
parent
5d7042e655
commit
c2bdf30e05
8 changed files with 605 additions and 26 deletions
|
|
@ -8,7 +8,7 @@
|
|||
mod rpc;
|
||||
mod stream;
|
||||
|
||||
use std::collections::{HashMap, VecDeque};
|
||||
use std::collections::{BTreeMap, HashMap, VecDeque};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::pin::Pin;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
|
|
@ -229,6 +229,67 @@ impl Choices {
|
|||
}
|
||||
}
|
||||
|
||||
/// An amount of money an agent reported, in an ISO 4217 `currency`.
|
||||
#[derive(Debug, Clone, PartialEq)]
|
||||
pub struct ReportedCost {
|
||||
pub amount: f64,
|
||||
pub currency: String,
|
||||
}
|
||||
|
||||
/// What the agent has reported spending, shared between the runtime and its
|
||||
/// [`Live`] agent so an amount outlives a respawn until it is taken.
|
||||
#[derive(Clone, Default)]
|
||||
struct Spend(Arc<std::sync::Mutex<Spent>>);
|
||||
|
||||
#[derive(Default)]
|
||||
struct Spent {
|
||||
/// The attached session's running total at its last report; `None` until
|
||||
/// it reports one.
|
||||
total: Option<ReportedCost>,
|
||||
/// The attached session is new, so its total started at zero.
|
||||
fresh: bool,
|
||||
/// Spent since the last [`Spend::take`], per currency.
|
||||
untaken: BTreeMap<String, f64>,
|
||||
}
|
||||
|
||||
impl Spend {
|
||||
/// Start following a newly attached session. A loaded session's total
|
||||
/// before its first report is unknown, so that report only sets the
|
||||
/// baseline.
|
||||
fn attach(&self, fresh: bool) {
|
||||
let mut spent = self.lock();
|
||||
spent.total = None;
|
||||
spent.fresh = fresh;
|
||||
}
|
||||
|
||||
/// Fold in a running total the attached session reported. A total lower
|
||||
/// than the last one means the agent dropped part of the session's
|
||||
/// history, so it adds nothing.
|
||||
fn observe(&self, total: ReportedCost) {
|
||||
let mut spent = self.lock();
|
||||
let added = match &spent.total {
|
||||
Some(last) if last.currency == total.currency => total.amount - last.amount,
|
||||
None if spent.fresh => total.amount,
|
||||
_ => 0.0,
|
||||
};
|
||||
if added > 0.0 {
|
||||
*spent.untaken.entry(total.currency.clone()).or_default() += added;
|
||||
}
|
||||
spent.total = Some(total);
|
||||
}
|
||||
|
||||
fn take(&self) -> Vec<ReportedCost> {
|
||||
std::mem::take(&mut self.lock().untaken)
|
||||
.into_iter()
|
||||
.map(|(currency, amount)| ReportedCost { amount, currency })
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn lock(&self) -> std::sync::MutexGuard<'_, Spent> {
|
||||
self.0.lock().unwrap_or_else(PoisonError::into_inner)
|
||||
}
|
||||
}
|
||||
|
||||
fn find<'a>(
|
||||
options: &'a [stream::ConfigOption],
|
||||
category: &str,
|
||||
|
|
@ -291,6 +352,7 @@ pub struct AcpRuntime<P: CompactionPolicy> {
|
|||
cancel_grace: Duration,
|
||||
compact_idle: Duration,
|
||||
choices: Choices,
|
||||
spend: Spend,
|
||||
}
|
||||
|
||||
/// The running agent process and the session loaded into it.
|
||||
|
|
@ -305,6 +367,7 @@ struct Live {
|
|||
commands: Option<Vec<String>>,
|
||||
/// The config options `loaded` last reported.
|
||||
choices: Choices,
|
||||
spend: Spend,
|
||||
}
|
||||
|
||||
impl Live {
|
||||
|
|
@ -340,6 +403,7 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
|||
cancel_grace: CANCEL_GRACE,
|
||||
compact_idle: COMPACT_IDLE,
|
||||
choices: Choices::default(),
|
||||
spend: Spend::default(),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -381,6 +445,7 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
|||
model: None,
|
||||
commands: None,
|
||||
choices: self.choices.clone(),
|
||||
spend: self.spend.clone(),
|
||||
})
|
||||
}
|
||||
|
||||
|
|
@ -414,6 +479,7 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
|||
live.model = stream::session_model(&response);
|
||||
live.offer(stream::config_options(&response).unwrap_or_default());
|
||||
live.loaded = Some(id.clone());
|
||||
live.spend.attach(false);
|
||||
live.commands = None;
|
||||
discard_stale(live)?;
|
||||
return Ok((id, new));
|
||||
|
|
@ -436,6 +502,7 @@ impl<P: CompactionPolicy> AcpRuntime<P> {
|
|||
live.model = stream::session_model(&response);
|
||||
live.offer(stream::config_options(&response).unwrap_or_default());
|
||||
live.loaded = Some(id.clone());
|
||||
live.spend.attach(true);
|
||||
live.commands = None;
|
||||
write_id(&self.pending_file(), &id)?;
|
||||
Ok((id, true))
|
||||
|
|
@ -709,6 +776,10 @@ impl<P: CompactionPolicy> Runtime for AcpRuntime<P> {
|
|||
Some(self.choices.clone())
|
||||
}
|
||||
|
||||
fn take_reported_cost(&self) -> Vec<ReportedCost> {
|
||||
self.spend.take()
|
||||
}
|
||||
|
||||
/// 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
|
||||
/// answered is dropped too.
|
||||
|
|
@ -761,8 +832,11 @@ fn deliver(
|
|||
}
|
||||
|
||||
/// Keep the commands and config options an `update` for the loaded session
|
||||
/// advertises.
|
||||
/// advertises, and the cost it reports.
|
||||
fn keep_advertised(agent: &mut Live, update: &Value) {
|
||||
if let Some(cost) = stream::reported_cost(update) {
|
||||
agent.spend.observe(cost);
|
||||
}
|
||||
if let Some(advertised) = stream::advertised_commands(update) {
|
||||
agent.commands = Some(advertised);
|
||||
}
|
||||
|
|
@ -775,7 +849,7 @@ fn keep_advertised(agent: &mut Live, update: &Value) {
|
|||
/// 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.
|
||||
/// The commands and config options the agent advertises for the loaded
|
||||
/// session are kept.
|
||||
/// session, and the cost it reports, are kept.
|
||||
fn discard_stale(live: &mut Live) -> std::result::Result<(), AcpError> {
|
||||
while let Ok(incoming) = live.conn.incoming.try_recv() {
|
||||
between_turns(live, incoming)?;
|
||||
|
|
@ -941,7 +1015,7 @@ mod tests {
|
|||
|
||||
use hive_claude::{Config, NoopSink, PercentPolicy, Sink};
|
||||
|
||||
use super::{AcpError, AcpRuntime, Choice, SessionChoices, TurnKind};
|
||||
use super::{AcpError, AcpRuntime, Choice, ReportedCost, SessionChoices, Spend, TurnKind};
|
||||
use crate::{AcpCommand, Error, Runtime};
|
||||
|
||||
/// An ACP agent in plain `sh`. It answers `initialize`, numbers its
|
||||
|
|
@ -971,7 +1045,9 @@ mod tests {
|
|||
/// 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`. With `REFUSE` set, it answers the first
|
||||
/// `session/set_config_option` to that value with an error.
|
||||
/// `session/set_config_option` to that value with an error. With `COST`
|
||||
/// set, each prompt adds that many dollars to its session's total (from zero,
|
||||
/// kept in `<log>.cost.<session>` across restarts) and reports it.
|
||||
const AGENT: &str = r#"
|
||||
n=0 prompt= model=m/think effort=low
|
||||
printf 'start\n' >> "$1.methods"
|
||||
|
|
@ -1007,6 +1083,7 @@ while IFS= read -r line; do
|
|||
printf '{"jsonrpc":"2.0","id":%s,"result":{"protocolVersion":1,"agentCapabilities":{"loadSession":true,"mcpCapabilities":{"http":true}}}}\n' "$id" ;;
|
||||
session/new)
|
||||
n=$((n+1))
|
||||
rm -f "$1.cost.s$n"
|
||||
printf '{"jsonrpc":"2.0","id":%s,"result":{"sessionId":"s%s"%s}}\n' "$id" "$n" "$(options)"
|
||||
advertise "s$n" ;;
|
||||
session/load)
|
||||
|
|
@ -1026,6 +1103,11 @@ while IFS= read -r line; do
|
|||
session/prompt)
|
||||
printf '%s\n' "$line" >> "$1"
|
||||
[ -n "$USED" ] && [ "$sid" = s1 ] && update "$sid" "{\"sessionUpdate\":\"usage_update\",\"used\":$USED,\"size\":1000}"
|
||||
if [ -n "$COST" ]; then
|
||||
total=$(( $(cat "$1.cost.$sid" 2>/dev/null || echo 0) + COST ))
|
||||
printf '%s' "$total" > "$1.cost.$sid"
|
||||
update "$sid" "{\"sessionUpdate\":\"usage_update\",\"cost\":{\"amount\":$total,\"currency\":\"USD\"}}"
|
||||
fi
|
||||
case $line in
|
||||
*'"text":"/compact"'*) [ -n "$STALL_COMPACT" ] && prompt=$id && continue ;;
|
||||
esac
|
||||
|
|
@ -1639,4 +1721,93 @@ done
|
|||
Some(choice("m/think", &["m/think", "m/plain"]))
|
||||
);
|
||||
}
|
||||
|
||||
fn usd(amount: f64) -> ReportedCost {
|
||||
ReportedCost {
|
||||
amount,
|
||||
currency: "USD".into(),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_new_session_s_first_report_is_all_spent() {
|
||||
let spend = Spend::default();
|
||||
spend.attach(true);
|
||||
spend.observe(usd(0.5));
|
||||
spend.observe(usd(1.25));
|
||||
assert_eq!(spend.take(), [usd(1.25)]);
|
||||
assert_eq!(spend.take(), []);
|
||||
spend.observe(usd(2.0));
|
||||
assert_eq!(spend.take(), [usd(0.75)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_loaded_session_s_first_report_only_sets_the_baseline() {
|
||||
let spend = Spend::default();
|
||||
spend.attach(false);
|
||||
spend.observe(usd(10.0));
|
||||
assert_eq!(spend.take(), []);
|
||||
spend.observe(usd(10.5));
|
||||
assert_eq!(spend.take(), [usd(0.5)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_falling_total_adds_nothing_and_becomes_the_baseline() {
|
||||
let spend = Spend::default();
|
||||
spend.attach(true);
|
||||
spend.observe(usd(3.0));
|
||||
spend.observe(usd(1.0));
|
||||
spend.observe(usd(1.5));
|
||||
assert_eq!(spend.take(), [usd(3.5)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn untaken_spend_survives_attaching_another_session() {
|
||||
let spend = Spend::default();
|
||||
spend.attach(true);
|
||||
spend.observe(usd(1.0));
|
||||
spend.attach(true);
|
||||
spend.observe(usd(0.25));
|
||||
assert_eq!(spend.take(), [usd(1.25)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_changed_currency_only_sets_a_new_baseline() {
|
||||
let spend = Spend::default();
|
||||
spend.attach(true);
|
||||
spend.observe(usd(1.0));
|
||||
let eur = |amount| ReportedCost {
|
||||
amount,
|
||||
currency: "EUR".into(),
|
||||
};
|
||||
spend.observe(eur(0.75));
|
||||
spend.observe(eur(1.25));
|
||||
assert_eq!(spend.take(), [eur(0.5), usd(1.0)]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn reported_cost_is_what_the_session_added_since_last_taken() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let config = config(dir.path());
|
||||
let env = [("COST", "2")];
|
||||
|
||||
let before = agent(dir.path(), "ok", &env, PercentPolicy::default());
|
||||
before.run(&config, "one", &NoopSink).await.unwrap();
|
||||
before.run(&config, "two", &NoopSink).await.unwrap();
|
||||
assert_eq!(before.take_reported_cost(), [usd(4.0)]);
|
||||
assert_eq!(before.take_reported_cost(), []);
|
||||
drop(before);
|
||||
|
||||
let after = agent(dir.path(), "ok", &env, PercentPolicy::default());
|
||||
after.run(&config, "three", &NoopSink).await.unwrap();
|
||||
assert_eq!(after.take_reported_cost(), []);
|
||||
after.run(&config, "four", &NoopSink).await.unwrap();
|
||||
assert_eq!(after.take_reported_cost(), [usd(2.0)]);
|
||||
|
||||
after.archive().unwrap();
|
||||
after.run(&config, "five", &NoopSink).await.unwrap();
|
||||
assert_eq!(after.take_reported_cost(), [usd(2.0)]);
|
||||
assert_eq!(count(dir.path(), "session/load"), 1);
|
||||
assert_eq!(count(dir.path(), "session/new"), 2);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -7,6 +7,8 @@ use std::collections::HashMap;
|
|||
use hive_claude::{Telemetry, TokenUsage};
|
||||
use serde_json::{Map, Value, json};
|
||||
|
||||
use super::ReportedCost;
|
||||
|
||||
/// Turns one prompt's `session/update` notifications into claude
|
||||
/// `stream-json` events.
|
||||
///
|
||||
|
|
@ -258,6 +260,19 @@ pub(super) fn config_options(value: &Value) -> Option<Vec<ConfigOption>> {
|
|||
)
|
||||
}
|
||||
|
||||
/// The running total cost a `usage_update` reports for its session, or
|
||||
/// `None` for any other `update` and for one that reports no cost.
|
||||
pub(super) fn reported_cost(update: &Value) -> Option<ReportedCost> {
|
||||
if update.get("sessionUpdate").and_then(Value::as_str) != Some("usage_update") {
|
||||
return None;
|
||||
}
|
||||
let cost = update.get("cost")?;
|
||||
Some(ReportedCost {
|
||||
amount: cost.get("amount")?.as_f64()?,
|
||||
currency: cost.get("currency")?.as_str()?.to_owned(),
|
||||
})
|
||||
}
|
||||
|
||||
/// 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>> {
|
||||
|
|
@ -402,7 +417,8 @@ fn tool_result(id: &str, output: &str, is_error: bool) -> Value {
|
|||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{
|
||||
ConfigOption, StreamMapper, canonical_tool_name, config_options, mcp_servers, session_model,
|
||||
ConfigOption, ReportedCost, StreamMapper, canonical_tool_name, config_options, mcp_servers,
|
||||
reported_cost, session_model,
|
||||
};
|
||||
use serde_json::{Value, json};
|
||||
|
||||
|
|
@ -537,6 +553,33 @@ mod tests {
|
|||
assert_eq!(t.model.as_deref(), Some("provider/model"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_usage_update_reports_its_cost_when_it_carries_one() {
|
||||
assert_eq!(
|
||||
reported_cost(&json!({ "sessionUpdate": "usage_update", "used": 5000,
|
||||
"size": 262_144, "cost": { "amount": 0.045, "currency": "USD" } })),
|
||||
Some(ReportedCost {
|
||||
amount: 0.045,
|
||||
currency: "USD".into()
|
||||
})
|
||||
);
|
||||
assert_eq!(
|
||||
reported_cost(&json!({ "sessionUpdate": "usage_update", "used": 5000,
|
||||
"size": 262_144 })),
|
||||
None
|
||||
);
|
||||
assert_eq!(
|
||||
reported_cost(&json!({ "sessionUpdate": "usage_update",
|
||||
"cost": { "amount": 0.045 } })),
|
||||
None
|
||||
);
|
||||
assert_eq!(
|
||||
reported_cost(&json!({ "sessionUpdate": "plan",
|
||||
"cost": { "amount": 0.045, "currency": "USD" } })),
|
||||
None
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn updates_with_nothing_to_show_emit_nothing() {
|
||||
let out = feed(
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ use std::path::PathBuf;
|
|||
|
||||
use hive_claude::{CompactionPolicy, Config, InfiniteSession, Progress, SessionStore, Sink};
|
||||
|
||||
use crate::{Canceller, Choices, Result, Runtime};
|
||||
use crate::{Canceller, Choices, ReportedCost, Result, Runtime};
|
||||
|
||||
/// `claude --print` turns on a titled [`InfiniteSession`], which owns
|
||||
/// resume-or-create and compaction.
|
||||
|
|
@ -47,4 +47,8 @@ impl<P: CompactionPolicy> Runtime for ClaudeRuntime<P> {
|
|||
fn choices(&self) -> Option<Choices> {
|
||||
None
|
||||
}
|
||||
|
||||
fn take_reported_cost(&self) -> Vec<ReportedCost> {
|
||||
Vec::new()
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -23,7 +23,7 @@ mod spec;
|
|||
|
||||
pub use acp::{
|
||||
AcpError, AcpRuntime, Canceller, Choice, Choices, PermissionAsk, PermissionPolicy,
|
||||
SessionChoices,
|
||||
ReportedCost, SessionChoices,
|
||||
};
|
||||
pub use claude::ClaudeRuntime;
|
||||
pub use hive_claude::{
|
||||
|
|
@ -65,6 +65,11 @@ pub trait Runtime {
|
|||
/// 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>;
|
||||
|
||||
/// What the agent reported spending since the last call, one entry per
|
||||
/// currency. Empty for a runtime whose agent reports no cost: claude
|
||||
/// exports its own.
|
||||
fn take_reported_cost(&self) -> Vec<ReportedCost>;
|
||||
}
|
||||
|
||||
/// The runtime an agent was configured with, chosen at startup from a
|
||||
|
|
@ -109,6 +114,13 @@ impl<P: CompactionPolicy> Runtime for AgentRuntime<P> {
|
|||
Self::Acp(r) => r.choices(),
|
||||
}
|
||||
}
|
||||
|
||||
fn take_reported_cost(&self) -> Vec<ReportedCost> {
|
||||
match self {
|
||||
Self::Claude(r) => r.take_reported_cost(),
|
||||
Self::Acp(r) => r.take_reported_cost(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Why a runtime operation did not complete.
|
||||
|
|
|
|||
Loading…
Reference in a new issue