Watch
0
0
Fork
You've already forked hyperhive
0

hive-agent: publish Claude subscription usage (5h/7d %) as metrics

A detached task polls GET https://api.anthropic.com/api/oauth/usage every
5 minutes with the OAuth access token from ~/.claude/.credentials.json and
records, per window the response names (five_hour, seven_day,
seven_day_sonnet, ...):

- hyperhive.agent.claude_usage.percent   (%, 0-100)
- hyperhive.agent.claude_usage.resets_at (s, unix seconds)

both labelled window=<name>. Endpoint, the anthropic-beta:
oauth-2025-04-20 header and the {utilization, resets_at} window shape are
taken from the claude-code 2.1.283 bundle's own /usage fetch.

The token is only read, never refreshed: claude owns refresh-token
rotation and a second refresher can log the agent out. An expired token
is skipped until claude's next turn refreshes it. API-key agents
(HIVE_USE_API_KEY, the ACP default) and agents with no credentials file
skip quietly, and the task does nothing when OTEL is not configured.
Request failures and non-2xx statuses warn with the status or transport
error only, never the body or the token.
This commit is contained in:
atlas 2026-09-30 20:26:09 +02:00 • committed by mara
commit 8a735ddcb3
4 changed files with 303 additions and 9 deletions

View file

@ -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 <token>` 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<OauthEntry>,
}
#[derive(Deserialize)]
struct OauthEntry {
#[serde(rename = "accessToken")]
access_token: Option<String>,
/// Unix milliseconds.
#[serde(rename = "expiresAt")]
expires_at: Option<i64>,
}
/// 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<u64>,
}
/// 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::<CredentialsFile>(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<Value, String> {
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::<Value>()
.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<Window> {
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));
}
}

View file

@ -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<S: Surface>(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

View file

@ -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<u64>,
/// Recorded from [`crate::claude_usage_watch`]'s own tick, keyed by
/// `window`.
claude_usage_percent: Gauge<f64>,
claude_usage_resets_at: Gauge<u64>,
}
/// 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<u64>) {
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<Instruments> {
if !enabled() {
tracing::debug!("otel turn-metrics: no endpoint configured, exporter disabled");
@ -133,6 +153,14 @@ fn build() -> Option<Instruments> {
.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<SdkMeterProvider> {
/// 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())
}