diff --git a/docs/scheduler/observability.md b/docs/scheduler/observability.md index bf63cd39..3766afea 100644 --- a/docs/scheduler/observability.md +++ b/docs/scheduler/observability.md @@ -361,14 +361,24 @@ here; that's already covered by Claude's own export. -| Metric | Unit | Kind | Attributes | -| -------------------------------------- | ---- | --------- | ----------------------------------------------------------------------------------------------------------- | -| `hyperhive.agent.turn.duration` | `ms` | histogram | `wake_from`, `result_kind`, `model` | -| `hyperhive.agent.turn.count` | — | counter | `wake_from`, `result_kind`, `model` | -| `hyperhive.agent.session.count` | — | counter | `model` (incremented once per fresh, non-`--continue`'d session) | -| `hyperhive.agent.loose_ends.threads` | — | gauge | none | -| `hyperhive.agent.loose_ends.reminders` | — | gauge | none | -| `hyperhive.agent.claude_md.lines` | — | gauge | none — recorded from the `CLAUDE.md`-size watch's own ~15-minute tick, **not** per turn like the rows above | +| Metric | Unit | Kind | Attributes | +| ---------------------------------------- | ---- | --------- | ------------------------------------------------------------------------------------------------------------------ | +| `hyperhive.agent.turn.duration` | `ms` | histogram | `wake_from`, `result_kind`, `model` | +| `hyperhive.agent.turn.count` | — | counter | `wake_from`, `result_kind`, `model` | +| `hyperhive.agent.session.count` | — | counter | `model` (incremented once per fresh, non-`--continue`'d session) | +| `hyperhive.agent.loose_ends.threads` | — | gauge | none | +| `hyperhive.agent.loose_ends.reminders` | — | gauge | none | +| `hyperhive.agent.claude_md.lines` | — | gauge | none — recorded from the `CLAUDE.md`-size watch's own ~15-minute tick, **not** per turn like the rows above | +| `hyperhive.agent.claude_usage.percent` | `%` | gauge | `window` (`five_hour`, `seven_day`, … as the usage endpoint names them) — polled every 5 minutes, **not** per turn | +| `hyperhive.agent.claude_usage.resets_at` | `s` | gauge | `window` — unix seconds at which that window resets; same 5-minute poll | + +For the two `claude_usage` gauges the harness polls the Claude subscription +usage endpoint (`GET /api/oauth/usage`, the one behind claude's own `/usage`) +with the OAuth access token from the agent's own +`~/.claude/.credentials.json`. The harness only reads that token and never +refreshes it — claude owns refresh-token rotation. An agent with no OAuth +session (API-key backends, not yet logged in) or an expired token skips the +poll, so a gauge keeps its last value until the next successful poll. Resource attributes (`service.name`, `agent`, `hive`, `swarm`) come from the same container-wide `OTEL_RESOURCE_ATTRIBUTES` as everything else in this diff --git a/hive-agent/src/claude_usage_watch.rs b/hive-agent/src/claude_usage_watch.rs new file mode 100644 index 00000000..f806222f --- /dev/null +++ b/hive-agent/src/claude_usage_watch.rs @@ -0,0 +1,254 @@ +//! Claude subscription usage watch: polls `GET /api/oauth/usage` with the +//! agent's own OAuth access token and records each window's utilization +//! (`five_hour`, `seven_day`, …) as an OTEL gauge. +//! +//! Endpoint, headers and response shape are the ones the bundled `claude` +//! CLI uses for its own `/usage` view: `Authorization: Bearer ` plus +//! `anthropic-beta: oauth-2025-04-20`; the body is an object whose window +//! entries are `{ utilization: 0-100 | null, resets_at: ISO 8601 | null }`. +//! +//! ⚠️ Read-only use of the current access token. This task never refreshes +//! it: `claude` owns refresh-token rotation, and a second refresher racing it +//! can invalidate the session and log the agent out. An expired token is +//! skipped until `claude`'s next turn refreshes it. +//! +//! Agents without an OAuth session (API-key backends, not yet logged in) +//! have nothing to poll with and skip quietly. + +use std::path::Path; +use std::time::Duration; + +use serde::Deserialize; +use serde_json::Value; + +const POLL_INTERVAL: Duration = Duration::from_mins(5); +const REQUEST_TIMEOUT: Duration = Duration::from_secs(10); +const USAGE_URL: &str = "https://api.anthropic.com/api/oauth/usage"; +const OAUTH_BETA: &str = "oauth-2025-04-20"; +/// The credentials file `claude` itself reads its OAuth session from. +const CREDENTIALS_FILE: &str = ".credentials.json"; + +/// No `Debug`/`Display`, so the token cannot end up in a log line by +/// accident. +struct AccessToken(String); + +enum Token { + /// No usable OAuth session in the credentials file. + Absent, + /// Past `expiresAt`; a request would only 401. + Expired, + Usable(AccessToken), +} + +#[derive(Deserialize)] +struct CredentialsFile { + #[serde(rename = "claudeAiOauth")] + oauth: Option, +} + +#[derive(Deserialize)] +struct OauthEntry { + #[serde(rename = "accessToken")] + access_token: Option, + /// Unix milliseconds. + #[serde(rename = "expiresAt")] + expires_at: Option, +} + +/// One usage window as reported by the endpoint. +#[derive(Debug, PartialEq)] +struct Window { + name: String, + percent: f64, + /// Unix seconds; `None` when the server sent no (or an unparsable) reset. + resets_at: Option, +} + +/// Background loop. Detached task — runs for the harness's lifetime; every +/// failure is logged and retried on the next tick, never fatal. +pub async fn run() { + if crate::login::using_api_key() { + tracing::debug!("claude usage watch: api-key agent, no OAuth token to poll with"); + return; + } + if !crate::otel_turn_metrics::enabled() { + return; + } + let client = match reqwest::Client::builder().timeout(REQUEST_TIMEOUT).build() { + Ok(c) => c, + Err(e) => { + tracing::warn!(error = %e, "claude usage watch: failed to build http client"); + return; + } + }; + let creds = crate::paths::claude_dir().join(CREDENTIALS_FILE); + let mut skip_logged = false; + let mut tick = tokio::time::interval(POLL_INTERVAL); + tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + loop { + tick.tick().await; + let token = match read_token(&creds).await { + Token::Usable(t) => t, + Token::Absent | Token::Expired if skip_logged => continue, + Token::Absent => { + tracing::debug!("claude usage watch: no OAuth access token, skipping"); + skip_logged = true; + continue; + } + Token::Expired => { + tracing::debug!("claude usage watch: OAuth access token expired, skipping"); + skip_logged = true; + continue; + } + }; + skip_logged = false; + match fetch(&client, &token).await { + Ok(body) => { + for w in windows(&body) { + crate::otel_turn_metrics::record_claude_usage(&w.name, w.percent, w.resets_at); + } + } + Err(e) => tracing::warn!("claude usage watch: {e}"), + } + } +} + +async fn read_token(path: &Path) -> Token { + let Ok(bytes) = tokio::fs::read(path).await else { + return Token::Absent; + }; + let now_ms = chrono::Utc::now().timestamp_millis(); + token_from(&bytes, now_ms) +} + +/// Parse errors are deliberately dropped: a serde message can quote the +/// offending value, and this file holds secrets. +fn token_from(bytes: &[u8], now_ms: i64) -> Token { + let Ok(file) = serde_json::from_slice::(bytes) else { + return Token::Absent; + }; + let Some(entry) = file.oauth else { + return Token::Absent; + }; + let Some(token) = entry.access_token.filter(|t| !t.is_empty()) else { + return Token::Absent; + }; + if entry.expires_at.is_some_and(|exp| exp <= now_ms) { + return Token::Expired; + } + Token::Usable(AccessToken(token)) +} + +/// Errors carry the status or the transport error only — never the response +/// body, never the token. +async fn fetch(client: &reqwest::Client, token: &AccessToken) -> Result { + let resp = client + .get(USAGE_URL) + .bearer_auth(&token.0) + .header("anthropic-beta", OAUTH_BETA) + .header(reqwest::header::CONTENT_TYPE, "application/json") + .send() + .await + .map_err(|e| format!("request failed: {e}"))?; + let status = resp.status(); + if !status.is_success() { + return Err(format!("request rejected: {status}")); + } + resp.json::() + .await + .map_err(|e| format!("unreadable response: {e}")) +} + +/// Window entries are the top-level objects carrying both `utilization` and +/// `resets_at`; other top-level keys (`extra_usage`, `limits`, …) are not +/// windows. A window whose `utilization` is `null` is not reported. +fn windows(body: &Value) -> Vec { + let Some(obj) = body.as_object() else { + return Vec::new(); + }; + obj.iter() + .filter_map(|(name, entry)| { + let entry = entry.as_object()?; + let resets_at = entry.get("resets_at")?; + let percent = entry.get("utilization")?.as_f64()?; + let resets_at = resets_at + .as_str() + .and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok()) + .and_then(|t| u64::try_from(t.timestamp()).ok()); + Some(Window { + name: name.clone(), + percent, + resets_at, + }) + }) + .collect() +} + +#[cfg(test)] +mod tests { + use super::*; + + /// Shape of a real response, field names from the `claude` CLI bundle. + const FIXTURE: &str = r#"{ + "five_hour": { "utilization": 42.0, "resets_at": "2026-09-30T18:00:00.123456+00:00" }, + "seven_day": { "utilization": 7, "resets_at": "2026-10-04T09:00:00Z" }, + "seven_day_oauth_apps": null, + "seven_day_opus": { "utilization": null, "resets_at": null }, + "seven_day_sonnet": { "utilization": 3.5, "resets_at": null }, + "extra_usage": { "is_enabled": false, "utilization": null }, + "limits": [ { "kind": "session", "percent": 42.0, "resets_at": null } ] + }"#; + + #[test] + fn parses_windows_and_skips_non_windows() { + let body: Value = serde_json::from_str(FIXTURE).expect("fixture"); + assert_eq!( + windows(&body), + vec![ + Window { + name: "five_hour".into(), + percent: 42.0, + resets_at: Some(1_790_791_200), + }, + Window { + name: "seven_day".into(), + percent: 7.0, + resets_at: Some(1_791_104_400), + }, + Window { + name: "seven_day_sonnet".into(), + percent: 3.5, + resets_at: None, + }, + ] + ); + } + + #[test] + fn non_object_body_yields_nothing() { + assert!(windows(&Value::Null).is_empty()); + assert!(windows(&serde_json::json!([1, 2])).is_empty()); + } + + #[test] + fn token_usable_before_expiry() { + let f = br#"{"claudeAiOauth":{"accessToken":"t","expiresAt":2000,"refreshToken":"r"}}"#; + assert!(matches!(token_from(f, 1000), Token::Usable(t) if t.0 == "t")); + } + + #[test] + fn token_expired_is_not_used() { + let f = br#"{"claudeAiOauth":{"accessToken":"t","expiresAt":1000}}"#; + assert!(matches!(token_from(f, 1000), Token::Expired)); + } + + #[test] + fn missing_or_malformed_credentials_are_absent() { + assert!(matches!(token_from(b"{}", 0), Token::Absent)); + assert!(matches!( + token_from(br#"{"claudeAiOauth":{"accessToken":""}}"#, 0), + Token::Absent + )); + assert!(matches!(token_from(b"not json", 0), Token::Absent)); + } +} diff --git a/hive-agent/src/main.rs b/hive-agent/src/main.rs index ddb5307d..f118ccdd 100644 --- a/hive-agent/src/main.rs +++ b/hive-agent/src/main.rs @@ -10,6 +10,7 @@ //! before lib + bin were collapsed into one) plus the serve loop. mod claude_md_watch; +mod claude_usage_watch; mod db_migrate; mod disk_watch; mod events; @@ -552,6 +553,7 @@ async fn serve_main(socket: &Path, poll_ms: u64) -> Result<()> { // bash-task files + verbose event rows). Runs here, not host-side in // hive-c0re, because the files are agent-owned — see `vacuum` module docs. tokio::spawn(crate::vacuum::run()); + tokio::spawn(claude_usage_watch::run()); // Log web_ui::serve's error instead of dropping it. A bare // `tokio::spawn(web_ui::serve(...))` discards the JoinHandle, so // any Err (e.g. EACCES from `bind_unix` when HIVE_WEB_SOCKET points diff --git a/hive-agent/src/otel_turn_metrics.rs b/hive-agent/src/otel_turn_metrics.rs index 86f89a7b..2717becd 100644 --- a/hive-agent/src/otel_turn_metrics.rs +++ b/hive-agent/src/otel_turn_metrics.rs @@ -22,6 +22,8 @@ //! exporter setup rather than [`record`]'s per-turn call site — `CLAUDE.md` //! size *is* continuously live (like hive-c0re's cgroup values), so //! [`crate::claude_md_watch`] records it from its own periodic tick instead. +//! [`record_claude_usage`] is the same kind of exception, fed from +//! [`crate::claude_usage_watch`]'s poll of the subscription usage endpoint. //! //! Resource attributes (`service.name`, `agent`, `hive`, `swarm`) are picked //! up automatically by the SDK from the container-wide `OTEL_RESOURCE_ATTRIBUTES` @@ -56,6 +58,10 @@ struct Instruments { /// hive-c0re's container gauges), not a once-per-turn one, so it has /// its own entry point rather than riding `record`'s per-turn call. claude_md_lines: Gauge, + /// Recorded from [`crate::claude_usage_watch`]'s own tick, keyed by + /// `window`. + claude_usage_percent: Gauge, + claude_usage_resets_at: Gauge, } /// Lazily built on the first call to [`record`]. `None` when OTEL isn't @@ -109,6 +115,20 @@ pub fn record_claude_md_lines(lines: u64) { inst.claude_md_lines.record(lines, &[]); } +/// Record one subscription usage window (`five_hour`, `seven_day`, …): +/// `percent` is 0-100, `resets_at` unix seconds. No-op when OTEL isn't +/// configured, same as [`record`]. +pub fn record_claude_usage(window: &str, percent: f64, resets_at: Option) { + let Some(inst) = INSTRUMENTS.get_or_init(build).as_ref() else { + return; + }; + let attrs = [KeyValue::new("window", window.to_owned())]; + inst.claude_usage_percent.record(percent, &attrs); + if let Some(resets_at) = resets_at { + inst.claude_usage_resets_at.record(resets_at, &attrs); + } +} + fn build() -> Option { if !enabled() { tracing::debug!("otel turn-metrics: no endpoint configured, exporter disabled"); @@ -133,6 +153,14 @@ fn build() -> Option { .u64_gauge("hyperhive.agent.loose_ends.reminders") .build(), claude_md_lines: meter.u64_gauge("hyperhive.agent.claude_md.lines").build(), + claude_usage_percent: meter + .f64_gauge("hyperhive.agent.claude_usage.percent") + .with_unit("%") + .build(), + claude_usage_resets_at: meter + .u64_gauge("hyperhive.agent.claude_usage.resets_at") + .with_unit("s") + .build(), _provider: provider, }) } @@ -168,7 +196,7 @@ fn build_provider(interval: Duration) -> anyhow::Result { /// hyperhive-specific one: it's also the variable the SDK itself reads to /// build the exporter's URL, so "configured" and "where it goes" cannot /// disagree). -fn enabled() -> bool { +pub fn enabled() -> bool { std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").is_ok_and(|s| !s.trim().is_empty()) }