diff --git a/rust/src/providers/opencodego/console.rs b/rust/src/providers/opencodego/console.rs new file mode 100644 index 0000000000..e37d972bd5 --- /dev/null +++ b/rust/src/providers/opencodego/console.rs @@ -0,0 +1,369 @@ +use chrono::{DateTime, TimeZone, Utc}; +use reqwest::Client; +use serde_json::Value; +use std::time::Duration; + +use crate::core::{ProviderError, RateWindow, UsageSnapshot}; + +use super::USER_AGENT; + +const CONSOLE_WORKSPACES_URL: &str = "https://opencode.ai/console/api/orgs"; +const CONSOLE_GO_STATUS_URL: &str = "https://opencode.ai/console/api/go/status"; +const CONSOLE_BILLING_STATUS_URL: &str = "https://opencode.ai/console/api/billing/status"; +const CONSOLE_WORKSPACE_HEADER: &str = "x-org-id"; +const BILLING_SCALE: f64 = 100_000_000.0; + +pub(super) enum ConsoleUsage { + Snapshot(Box), + NoSubscription, +} + +pub(super) fn normalize_workspace_id(raw: Option<&str>) -> Option { + let raw = raw?.trim(); + if is_workspace_id(raw) { + return Some(raw.to_string()); + } + let url = reqwest::Url::parse(raw).ok()?; + if url.scheme() != "https" || url.host_str() != Some("opencode.ai") { + return None; + } + let segments: Vec<_> = url.path_segments()?.collect(); + let candidate = segments + .windows(2) + .find_map(|parts| matches!(parts[0], "console" | "workspace").then_some(parts[1]))?; + is_workspace_id(candidate).then(|| candidate.to_string()) +} + +fn is_workspace_id(value: &str) -> bool { + let Some(suffix) = value + .strip_prefix("wrk_") + .or_else(|| value.strip_prefix("org_")) + else { + return false; + }; + !suffix.is_empty() + && suffix + .bytes() + .all(|byte| byte.is_ascii_alphanumeric() || byte == b'_' || byte == b'-') +} + +pub(super) async fn fetch_workspace_id( + client: &Client, + cookie_header: &str, + timeout: Duration, +) -> Result { + let text = fetch_text( + client, + CONSOLE_WORKSPACES_URL, + None, + cookie_header, + timeout, + "workspace list", + ) + .await?; + parse_workspace_ids(&text) + .into_iter() + .next() + .ok_or_else(|| ProviderError::Parse("Missing OpenCode Console workspace ID".to_string())) +} + +pub(super) async fn fetch_usage( + client: &Client, + workspace_id: &str, + cookie_header: &str, + timeout: Duration, +) -> Result { + let text = fetch_text( + client, + CONSOLE_GO_STATUS_URL, + Some(workspace_id), + cookie_header, + timeout, + "Go status", + ) + .await?; + parse_usage(&text, Utc::now()) +} + +pub(super) async fn fetch_balance( + client: &Client, + workspace_id: &str, + cookie_header: &str, + timeout: Duration, +) -> Result, ProviderError> { + let text = fetch_text( + client, + CONSOLE_BILLING_STATUS_URL, + Some(workspace_id), + cookie_header, + timeout, + "billing status", + ) + .await?; + parse_balance(&text) +} + +async fn fetch_text( + client: &Client, + url: &str, + workspace_id: Option<&str>, + cookie_header: &str, + timeout: Duration, + what: &str, +) -> Result { + let mut request = client + .get(url) + .timeout(timeout) + .header("Cookie", cookie_header) + .header("User-Agent", USER_AGENT) + .header("Accept", "application/json"); + if let Some(workspace_id) = workspace_id { + request = request.header(CONSOLE_WORKSPACE_HEADER, workspace_id); + } + let response = request.send().await?; + let status = response.status(); + if status.as_u16() == 401 { + return Err(ProviderError::AuthRequired); + } + if !status.is_success() { + return Err(ProviderError::Other(format!( + "OpenCode Console {what} returned {status}" + ))); + } + response.text().await.map_err(ProviderError::Network) +} + +fn parse_workspace_ids(text: &str) -> Vec { + let Ok(Value::Array(rows)) = serde_json::from_str(text) else { + return Vec::new(); + }; + rows.iter() + .filter_map(|row| row.get("id")?.as_str()) + .filter(|id| is_workspace_id(id)) + .map(str::to_string) + .collect() +} + +fn parse_usage(text: &str, _now: DateTime) -> Result { + let root: Value = serde_json::from_str(text) + .map_err(|_| ProviderError::Parse("Invalid OpenCode Console usage payload".to_string()))?; + if root.is_null() || root.get("access").is_some_and(Value::is_null) { + return Ok(ConsoleUsage::NoSubscription); + } + let access = root + .get("access") + .and_then(Value::as_object) + .ok_or_else(|| { + ProviderError::Parse("Invalid OpenCode Console usage payload".to_string()) + })?; + let meters = access + .get("meters") + .and_then(Value::as_object) + .ok_or_else(|| { + ProviderError::Parse("Invalid OpenCode Console usage payload".to_string()) + })?; + let rolling = meters.get("fiveHour").ok_or_else(|| { + ProviderError::Parse("Invalid OpenCode Console usage payload".to_string()) + })?; + let ends_at = access.get("endsAt").and_then(date_value); + let primary = meter_window(rolling, Some(300), None)?; + let mut snapshot = UsageSnapshot::new(primary).with_login_method("OpenCode Go"); + + if let Some(weekly) = meters.get("week") { + snapshot = snapshot.with_secondary(meter_window(weekly, Some(10_080), None)?); + } + if let Some(monthly) = meters.get("month") { + let resets_at = monthly.get("resetsAt").and_then(date_value).or(ends_at); + let minutes = RateWindow::monthly_window_minutes(resets_at).or(Some(43_200)); + snapshot = snapshot.with_tertiary(meter_window(monthly, minutes, resets_at)?); + } + if let Some(ends_at) = ends_at { + snapshot = snapshot.with_extra_rate_window( + "renewal", + "Renews", + RateWindow::with_details(0.0, None, Some(ends_at), None), + ); + } + Ok(ConsoleUsage::Snapshot(Box::new(snapshot))) +} + +fn meter_window( + meter: &Value, + window_minutes: Option, + reset_override: Option>, +) -> Result { + let object = meter + .as_object() + .ok_or_else(|| ProviderError::Parse("Invalid OpenCode Console usage meter".to_string()))?; + let used = numeric_value(object.get("usedMicroCents")) + .ok_or_else(|| ProviderError::Parse("Missing OpenCode Console usage amount".to_string()))?; + let limit = numeric_value(object.get("limitMicroCents")) + .filter(|value| *value > 0.0) + .ok_or_else(|| ProviderError::Parse("Missing OpenCode Console usage limit".to_string()))?; + let resets_at = reset_override.or_else(|| object.get("resetsAt").and_then(date_value)); + Ok(RateWindow::with_details( + ((used / limit) * 100.0).clamp(0.0, 100.0), + window_minutes, + resets_at, + None, + )) +} + +fn parse_balance(text: &str) -> Result, ProviderError> { + let root: Value = serde_json::from_str(text).map_err(|_| { + ProviderError::Parse("Invalid OpenCode Console billing payload".to_string()) + })?; + let billing_mode = root + .get("billingMode") + .and_then(Value::as_str) + .filter(|mode| matches!(*mode, "prepaid" | "legacy" | "seat" | "credit")) + .ok_or_else(|| { + ProviderError::Parse("Invalid OpenCode Console billing payload".to_string()) + })?; + let mode = root + .get("mode") + .and_then(Value::as_str) + .filter(|mode| matches!(*mode, "pay-as-you-go" | "invoiceable")) + .ok_or_else(|| { + ProviderError::Parse("Invalid OpenCode Console billing payload".to_string()) + })?; + if billing_mode != "prepaid" || mode != "pay-as-you-go" { + return Ok(None); + } + let raw = root + .get("balanceMicroCents") + .and_then(Value::as_str) + .ok_or_else(|| ProviderError::Parse("Missing OpenCode Console balance".to_string()))?; + let digits = raw.strip_prefix('-').unwrap_or(raw); + if digits.is_empty() || !digits.bytes().all(|byte| byte.is_ascii_digit()) { + return Err(ProviderError::Parse( + "Invalid OpenCode Console balance".to_string(), + )); + } + let balance = raw + .parse::() + .map_err(|_| ProviderError::Parse("Invalid OpenCode Console balance".to_string()))?; + if !balance.is_finite() { + return Err(ProviderError::Parse( + "Invalid OpenCode Console balance".to_string(), + )); + } + Ok(Some(balance / BILLING_SCALE)) +} + +fn numeric_value(value: Option<&Value>) -> Option { + match value? { + Value::Number(number) => number.as_f64().filter(|value| value.is_finite()), + Value::String(text) => text.parse::().ok().filter(|value| value.is_finite()), + _ => None, + } +} + +#[allow( + clippy::cast_possible_truncation, + clippy::cast_sign_loss, + reason = "validated Unix timestamps are narrowed only after selecting seconds and nanoseconds" +)] +fn date_value(value: &Value) -> Option> { + if let Some(number) = numeric_value(Some(value)) { + let seconds = if number > 1_000_000_000_000.0 { + number / 1000.0 + } else { + number + }; + let whole = seconds.trunc() as i64; + let nanos = ((seconds.fract()) * 1_000_000_000.0).round() as u32; + return Utc.timestamp_opt(whole, nanos).single(); + } + DateTime::parse_from_rfc3339(value.as_str()?) + .ok() + .map(|date| date.with_timezone(&Utc)) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn now() -> DateTime { + DateTime::parse_from_rfc3339("2026-09-20T00:00:00Z") + .unwrap() + .with_timezone(&Utc) + } + + fn usage_json(five_hour_reset: &str) -> String { + format!( + r#"{{"access":{{"endsAt":"2026-10-19T00:00:00Z","meters":{{"fiveHour":{{"resetsAt":{five_hour_reset},"limitMicroCents":"1200000000","usedMicroCents":"300000000"}},"week":{{"resetsAt":"2026-09-21T00:00:00Z","limitMicroCents":"3000000000","usedMicroCents":"1200000000"}},"month":{{"limitMicroCents":"6000000000","usedMicroCents":"600000000"}}}}}}}}"# + ) + } + + #[test] + fn parses_console_workspace_ids_and_urls() { + assert_eq!( + parse_workspace_ids(r#"[{"id":"wrk_ONE"},{"id":"org_TWO"},{"id":"acc_BAD"}]"#), + vec!["wrk_ONE", "org_TWO"] + ); + assert_eq!( + normalize_workspace_id(Some("https://opencode.ai/console/org_TWO/go")).as_deref(), + Some("org_TWO") + ); + assert_eq!( + normalize_workspace_id(Some("https://opencode.ai/workspace/wrk_ONE/go")).as_deref(), + Some("wrk_ONE") + ); + assert_eq!( + normalize_workspace_id(Some("https://example.com/console/org_TWO/go")), + None + ); + } + + #[test] + fn parses_console_microcent_usage_and_nullable_resets() { + let ConsoleUsage::Snapshot(snapshot) = parse_usage(&usage_json("null"), now()).unwrap() + else { + panic!("expected usage snapshot"); + }; + assert!((snapshot.primary.used_percent - 25.0).abs() < 0.001); + assert_eq!(snapshot.primary.resets_at, None); + assert!((snapshot.secondary.unwrap().used_percent - 40.0).abs() < 0.001); + let monthly = snapshot.tertiary.unwrap(); + assert!((monthly.used_percent - 10.0).abs() < 0.001); + assert_eq!( + monthly.resets_at.unwrap().to_rfc3339(), + "2026-10-19T00:00:00+00:00" + ); + } + + #[test] + fn distinguishes_no_subscription_from_malformed_usage() { + assert!(matches!( + parse_usage("null", now()), + Ok(ConsoleUsage::NoSubscription) + )); + assert!(matches!( + parse_usage(r#"{"access":null}"#, now()), + Ok(ConsoleUsage::NoSubscription) + )); + assert!(matches!( + parse_usage(r#"{"access":{}}"#, now()), + Err(ProviderError::Parse(_)) + )); + } + + #[test] + fn parses_only_explicit_prepaid_payg_balances() { + let payload = + r#"{"billingMode":"prepaid","mode":"pay-as-you-go","balanceMicroCents":"2786781005"}"#; + let balance = parse_balance(payload).unwrap().unwrap(); + assert!((balance - 27.867_810_05).abs() < 0.000_001); + assert_eq!( + parse_balance(r#"{"billingMode":"seat","mode":"pay-as-you-go"}"#).unwrap(), + None + ); + assert!( + parse_balance( + r#"{"billingMode":"prepaid","mode":"pay-as-you-go","balanceMicroCents":"1.5"}"# + ) + .is_err() + ); + } +} diff --git a/rust/src/providers/opencodego/legacy.rs b/rust/src/providers/opencodego/legacy.rs new file mode 100644 index 0000000000..c876c1b8c5 --- /dev/null +++ b/rust/src/providers/opencodego/legacy.rs @@ -0,0 +1,490 @@ +use chrono::Utc; +use reqwest::Client; +use std::time::Duration; +use uuid::Uuid; + +use crate::core::{ProviderError, RateWindow, UsageSnapshot}; + +use super::{BASE_URL, USER_AGENT}; + +const SERVER_URL: &str = "https://opencode.ai/_server"; +const WORKSPACES_SERVER_ID: &str = + "def39973159c7f0483d8793a822b8dbb10d067e12c65455fcb4608459ba0234f"; +const BILLING_SERVER_ID: &str = "c83b78a614689c38ebee981f9b39a8b377716db85c1fd7dbab604adc02d3313d"; + +/// Parsed legacy usage response returned to the provider orchestrator. +pub(super) struct LegacyUsage { + pub(super) usage: UsageSnapshot, + pub(super) embedded_balance: Option, +} + +/// Discover the workspace owned by the legacy cookie session. The provider +/// orchestrator resolves this once and carries the resulting typed session +/// through every legacy operation. +pub(super) async fn discover_workspace_id( + client: &Client, + cookie_header: &str, +) -> Result { + let text = fetch_server_text( + client, + cookie_header, + WORKSPACES_SERVER_ID, + None, + BASE_URL, + None, + "workspace API", + ) + .await?; + parse_workspace_ids(&text) + .into_iter() + .next() + .ok_or_else(|| ProviderError::Parse("No workspace ID found".to_string())) +} + +pub(super) async fn fetch_usage( + client: &Client, + cookie_header: &str, + workspace_id: &str, +) -> Result { + let url = format!("{BASE_URL}/workspace/{workspace_id}/go"); + let page = fetch_page_text(client, cookie_header, &url, None, "usage page").await?; + Ok(LegacyUsage { + usage: parse_usage_text(&page)?, + embedded_balance: parse_zen_balance(&page), + }) +} + +pub(super) async fn fetch_balance( + client: &Client, + cookie_header: &str, + workspace_id: &str, + timeout: Duration, +) -> Result, ProviderError> { + let referer = format!("{BASE_URL}/workspace/{workspace_id}"); + let page = fetch_page_text( + client, + cookie_header, + &referer, + Some(timeout), + "Zen dashboard page", + ) + .await?; + if let Some(balance) = parse_zen_balance(&page) { + return Ok(Some(balance)); + } + + let args = serde_json::json!([workspace_id]).to_string(); + let billing = fetch_server_text( + client, + cookie_header, + BILLING_SERVER_ID, + Some(&args), + &referer, + Some(timeout), + "billing API", + ) + .await?; + Ok(parse_billing_server_balance(&billing)) +} + +async fn fetch_page_text( + client: &Client, + cookie_header: &str, + url: &str, + timeout: Option, + what: &str, +) -> Result { + let mut request = client + .get(url) + .header("Cookie", cookie_header) + .header("User-Agent", USER_AGENT) + .header("Referer", BASE_URL) + .header( + "Accept", + "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8", + ); + if let Some(timeout) = timeout { + request = request.timeout(timeout); + } + let response = request.send().await?; + let status = response.status(); + if !status.is_success() { + if status.as_u16() == 401 || status.as_u16() == 403 { + return Err(ProviderError::AuthRequired); + } + return Err(ProviderError::Other(format!( + "OpenCode Go {what} returned {status}" + ))); + } + let text = response.text().await?; + if looks_signed_out(&text) { + return Err(ProviderError::AuthRequired); + } + Ok(text) +} + +async fn fetch_server_text( + client: &Client, + cookie_header: &str, + server_id: &str, + args: Option<&str>, + referer: &str, + timeout: Option, + what: &str, +) -> Result { + let mut url = reqwest::Url::parse(SERVER_URL) + .map_err(|error| ProviderError::Parse(format!("Invalid OpenCode server URL: {error}")))?; + url.query_pairs_mut().append_pair("id", server_id); + if let Some(args) = args { + url.query_pairs_mut().append_pair("args", args); + } + let mut request = client + .get(url) + .header("Cookie", cookie_header) + .header("X-Server-Id", server_id) + .header("X-Server-Instance", format!("server-fn:{}", Uuid::new_v4())) + .header("User-Agent", USER_AGENT) + .header("Origin", BASE_URL) + .header("Referer", referer) + .header( + "Accept", + "text/javascript, application/json;q=0.9, */*;q=0.8", + ); + if let Some(timeout) = timeout { + request = request.timeout(timeout); + } + let response = request.send().await?; + let status = response.status(); + if !status.is_success() { + if status.as_u16() == 401 || status.as_u16() == 403 { + return Err(ProviderError::AuthRequired); + } + return Err(ProviderError::Other(format!( + "OpenCode Go {what} returned {status}" + ))); + } + let text = response.text().await?; + if looks_signed_out(&text) { + return Err(ProviderError::AuthRequired); + } + Ok(text) +} + +fn parse_workspace_ids(text: &str) -> Vec { + let Ok(re) = regex_lite::Regex::new(r#"(wrk_[A-Za-z0-9_-]+)"#) else { + return Vec::new(); + }; + let mut seen = Vec::new(); + for captures in re.captures_iter(text) { + if let Some(value) = captures.get(1) { + let value = value.as_str().to_string(); + if !seen.contains(&value) { + seen.push(value); + } + } + } + seen +} + +fn looks_signed_out(text: &str) -> bool { + let lower = text.to_lowercase(); + lower.contains("auth/authorize") + || lower.contains("\"signin\"") + || lower.contains("please sign in") +} + +fn parse_usage_text(text: &str) -> Result { + let now = Utc::now(); + let rolling = extract_window(text, &["rollingUsage", "rolling_usage", "rolling"]) + .ok_or_else(|| ProviderError::Parse("Missing rolling usage window".to_string()))?; + let weekly = extract_window(text, &["weeklyUsage", "weekly_usage", "weekly"]); + let monthly = extract_window(text, &["monthlyUsage", "monthly_usage", "monthly"]); + + let primary = RateWindow::with_details( + rolling.0, + Some(300), + Some(now + chrono::Duration::seconds(rolling.1)), + None, + ); + let mut snapshot = UsageSnapshot::new(primary).with_login_method("OpenCode Go"); + if let Some((percent, reset)) = weekly { + snapshot = snapshot.with_secondary(RateWindow::with_details( + percent, + Some(10080), + Some(now + chrono::Duration::seconds(reset)), + None, + )); + } + if let Some((percent, reset)) = monthly { + let resets_at = now + chrono::Duration::seconds(reset); + snapshot = snapshot.with_tertiary(RateWindow::with_details( + percent, + RateWindow::monthly_window_minutes(Some(resets_at)).or(Some(43200)), + Some(resets_at), + None, + )); + } + if let Some(renews_at) = super::super::extract_renewal(text) { + snapshot = snapshot.with_extra_rate_window( + "renewal", + "Renews", + RateWindow::with_details(0.0, None, Some(renews_at), None), + ); + } + Ok(snapshot) +} + +fn extract_window(text: &str, names: &[&str]) -> Option<(f64, i64)> { + for name in names { + let percent_pattern = format!( + r#"{}[^}}]*?(?:usagePercent|usedPercent|percentUsed|percent)\s*[:=]\s*([0-9]+(?:\.[0-9]+)?)"#, + name + ); + let reset_pattern = format!( + r#"{}[^}}]*?(?:resetInSec|resetInSeconds|resetSeconds|resetSec)\s*[:=]\s*([0-9]+)"#, + name + ); + if let Some(percent) = super::super::extract_number(&percent_pattern, text) { + #[allow( + clippy::cast_possible_truncation, + reason = "resetInSec values are whole-second counts scraped as integral numbers" + )] + let reset = super::super::extract_number(&reset_pattern, text) + .map(|number| number as i64) + .unwrap_or(0); + return Some((percent.clamp(0.0, 100.0), reset.max(0))); + } + + let used_pattern = format!( + r#"{}[^}}]*?(?:used|usage|consumed)\s*[:=]\s*([0-9]+(?:\.[0-9]+)?)"#, + name + ); + let limit_pattern = format!( + r#"{}[^}}]*?(?:limit|total|allowance)\s*[:=]\s*([0-9]+(?:\.[0-9]+)?)"#, + name + ); + if let (Some(used), Some(limit)) = ( + super::super::extract_number(&used_pattern, text), + super::super::extract_number(&limit_pattern, text), + ) && limit > 0.0 + { + #[allow( + clippy::cast_possible_truncation, + reason = "resetInSec values are whole-second counts scraped as integral numbers" + )] + let reset = super::super::extract_number(&reset_pattern, text) + .map(|number| number as i64) + .unwrap_or(0); + return Some((((used / limit) * 100.0).clamp(0.0, 100.0), reset.max(0))); + } + } + None +} + +fn parse_zen_balance(text: &str) -> Option { + if let Ok(json) = serde_json::from_str::(text) + && let Some(value) = find_balance_value(&json) + { + return Some(value); + } + let patterns = [ + r#"(?i)(?:current\s+balance|zen\s+balance|現在の残高)[^$]{0,80}\$\s*([0-9][0-9,]*(?:\.[0-9]+)?)"#, + r#"(?i)(?:balance|残高)[\s\S]{0,120}?\$\s*([0-9][0-9,]*(?:\.[0-9]+)?)"#, + ]; + patterns.iter().find_map(|pattern| { + let re = regex_lite::Regex::new(pattern).ok()?; + let raw = re.captures(text)?.get(1)?.as_str().replace(',', ""); + raw.parse::().ok() + }) +} + +fn find_balance_value(value: &serde_json::Value) -> Option { + match value { + serde_json::Value::Object(map) => { + for (key, value) in map { + let normalized: String = key + .to_lowercase() + .chars() + .filter(|character| character.is_ascii_alphanumeric()) + .collect(); + if matches!( + normalized.as_str(), + "zenbalance" + | "zencurrentbalance" + | "currentbalance" + | "currentbalanceusd" + | "balanceusd" + | "usdbalance" + ) { + if let Some(number) = value.as_f64() { + return Some(number); + } + if let Some(text) = value.as_str() + && let Ok(number) = text.trim().replace(',', "").parse() + { + return Some(number); + } + } + if let Some(found) = find_balance_value(value) { + return Some(found); + } + } + None + } + serde_json::Value::Array(items) => items.iter().find_map(find_balance_value), + _ => None, + } +} + +fn parse_billing_server_balance(text: &str) -> Option { + const BILLING_SCALE: f64 = 100_000_000.0; + if let Ok(json) = serde_json::from_str::(text) + && let Some(raw) = find_raw_billing_balance(&json) + { + return Some(raw / BILLING_SCALE); + } + let customer_re = regex_lite::Regex::new( + r#"(?:\"customerID\"|customerID)\s*:\s*(?:\$R\[\d+\]\s*=\s*)?\"[^\"]+\""#, + ) + .ok()?; + customer_re.find(text)?; + let balance_re = regex_lite::Regex::new( + r#"(?:\"balance\"|balance)\s*:\s*(?:\$R\[\d+\]\s*=\s*)?(-?[0-9]+(?:\.[0-9]+)?)"#, + ) + .ok()?; + let raw: f64 = balance_re + .captures(text)? + .get(1)? + .as_str() + .replace(',', "") + .parse() + .ok()?; + Some(raw / BILLING_SCALE) +} + +fn find_raw_billing_balance(value: &serde_json::Value) -> Option { + match value { + serde_json::Value::Object(map) => { + if let Some(balance) = map.get("balance") { + let customer_ok = map + .get("customerID") + .and_then(|value| value.as_str()) + .is_some_and(|id| !id.is_empty()); + if !customer_ok { + return None; + } + return billing_numeric_value(balance); + } + map.values().find_map(find_raw_billing_balance) + } + serde_json::Value::Array(items) => items.iter().find_map(find_raw_billing_balance), + _ => None, + } +} + +fn billing_numeric_value(value: &serde_json::Value) -> Option { + match value { + serde_json::Value::Number(number) => number.as_f64(), + serde_json::Value::String(text) => text.trim().replace(',', "").parse().ok(), + _ => None, + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn parses_workspace_ids_without_duplicates() { + let text = r#"{ id: "wrk_abc123", name: "x" } { id: "wrk_def456" } { id: "wrk_abc123" }"#; + assert_eq!( + parse_workspace_ids(text), + vec!["wrk_abc123".to_string(), "wrk_def456".to_string()] + ); + } + + #[test] + fn parses_usage_blocks() { + let text = r#" + rollingUsage: { usagePercent: 42.5, resetInSec: 3600 } + weeklyUsage: { usagePercent: 13, resetInSec: 86400 } + monthlyUsage: { usagePercent: 7, resetInSec: 2592000 } + "#; + let snapshot = parse_usage_text(text).unwrap(); + assert!((snapshot.primary.used_percent - 42.5).abs() < 0.001); + assert!((snapshot.secondary.unwrap().used_percent - 13.0).abs() < 0.001); + assert!((snapshot.tertiary.unwrap().used_percent - 7.0).abs() < 0.001); + } + + #[test] + fn sub_one_percent_computed_used_limit_is_not_rescaled() { + let text = r#" + rollingUsage: { used: 1, limit: 100, resetInSec: 600 } + weeklyUsage: { used: 1, limit: 200, resetInSec: 86400 } + "#; + let snapshot = parse_usage_text(text).unwrap(); + assert!((snapshot.primary.used_percent - 1.0).abs() < 0.001); + assert!((snapshot.secondary.unwrap().used_percent - 0.5).abs() < 0.001); + } + + #[test] + fn direct_percent_one_is_not_rescaled() { + let text = r#"rollingUsage:$R[34]={status:"ok",resetInSec:13631,usagePercent:1} weeklyUsage:$R[35]={status:"ok",resetInSec:53863,usagePercent:15}"#; + let snapshot = parse_usage_text(text).unwrap(); + assert!((snapshot.primary.used_percent - 1.0).abs() < 0.001); + assert!((snapshot.secondary.unwrap().used_percent - 15.0).abs() < 0.001); + } + + #[test] + fn parses_renewal_window() { + let text = r#" + rollingUsage: { usagePercent: 42.5, resetInSec: 3600 } + renewAt: "2026-06-01T12:00:00Z" + "#; + let snapshot = parse_usage_text(text).unwrap(); + let renewal = snapshot + .extra_rate_windows + .iter() + .find(|window| window.id == "renewal") + .expect("renewal window"); + assert_eq!( + renewal.window.resets_at.unwrap().to_rfc3339(), + "2026-06-01T12:00:00+00:00" + ); + } + + #[test] + fn billing_server_balance_needs_customer_marker() { + assert_eq!( + parse_billing_server_balance(r#"{"balance": 1500000000, "customerID": "cus_123"}"#), + Some(15.0) + ); + assert_eq!( + parse_billing_server_balance(r#"{"balance": 1500000000}"#), + None + ); + assert_eq!( + parse_billing_server_balance(r#"{"balance": true, "customerID": "cus_123"}"#), + None + ); + } + + #[test] + fn billing_server_balance_handles_nesting_and_rsc_fragments() { + assert_eq!( + parse_billing_server_balance(r#"{"balance": "1,000,000,000", "customerID": "cus_1"}"#), + Some(10.0) + ); + assert_eq!( + parse_billing_server_balance( + r#"{"data": {"rows": [{"balance": 250000000, "customerID": "cus_2"}]}}"# + ), + Some(2.5) + ); + assert_eq!( + parse_billing_server_balance(r#"customerID:$R[1] = "cus_9"; "balance": -500000000"#), + Some(-5.0) + ); + assert_eq!(parse_billing_server_balance(r#"customerID: "cus_9""#), None); + } +} diff --git a/rust/src/providers/opencodego/mod.rs b/rust/src/providers/opencodego/mod.rs index 0e3f49d38f..d676e3dd94 100644 --- a/rust/src/providers/opencodego/mod.rs +++ b/rust/src/providers/opencodego/mod.rs @@ -5,14 +5,19 @@ //! unless a workspace override scopes the fetch to web first; Web is cookie //! scrape only; Cli is local-only. +mod console; +mod legacy; pub(crate) mod local; +mod transport; mod usage_api; use async_trait::async_trait; use chrono::Utc; use reqwest::Client; +use std::sync::Arc; use std::time::Duration; -use uuid::Uuid; + +use transport::{HttpWebTransport, WebTransport}; use crate::core::{ CostSnapshot, FetchContext, Provider, ProviderError, ProviderFetchResult, ProviderId, @@ -20,12 +25,8 @@ use crate::core::{ }; const BASE_URL: &str = "https://opencode.ai"; -const SERVER_URL: &str = "https://opencode.ai/_server"; /// Source label for quota values reconstructed from the device-local SQLite history. pub const LOCAL_ESTIMATE_SOURCE_LABEL: &str = "local estimate"; -const WORKSPACES_SERVER_ID: &str = - "def39973159c7f0483d8793a822b8dbb10d067e12c65455fcb4608459ba0234f"; -const BILLING_SERVER_ID: &str = "c83b78a614689c38ebee981f9b39a8b377716db85c1fd7dbab604adc02d3313d"; const USER_AGENT: &str = "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0.0.0 Safari/537.36"; // Upstream 0.48.0 #2583 (F15) optional-Zen-balance bounds. @@ -41,6 +42,76 @@ pub struct OpenCodeGoProvider { client: Client, } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +struct CookieCapabilities { + console: bool, + legacy: bool, +} + +#[derive(Clone, Debug)] +struct WebCookieSession { + header: String, + capabilities: CookieCapabilities, +} + +#[derive(Clone, Copy, Debug, PartialEq)] +enum OptionalZenBalance { + Resolved(Option), + LegacyBalanceRequired, +} + +impl WebCookieSession { + fn new(header: &str) -> Self { + let has_cookie = |names: &[&str]| { + header.split(';').any(|part| { + let Some((name, value)) = part.trim().split_once('=') else { + return false; + }; + names.contains(&name.trim()) && !value.trim().is_empty() + }) + }; + Self { + header: header.to_string(), + capabilities: CookieCapabilities { + console: has_cookie(&["__Host-console_session"]), + legacy: has_cookie(&["auth", "__Host-auth"]), + }, + } + } + + fn header(&self) -> &str { + &self.header + } + + fn can_recover_with_legacy(&self, error: &ProviderError) -> bool { + self.capabilities.legacy + && (matches!( + error, + ProviderError::AuthRequired + | ProviderError::Parse(_) + | ProviderError::Other(_) + | ProviderError::Timeout + ) || error.is_transport_failure()) + } + + fn select_legacy_result( + &self, + console_error: ProviderError, + legacy_result: Result, + ) -> Result { + match legacy_result { + Ok(value) => Ok(value), + Err(_legacy_error) + if self.capabilities.console + && !matches!(console_error, ProviderError::AuthRequired) => + { + Err(console_error) + } + Err(legacy_error) => Err(legacy_error), + } + } +} + impl OpenCodeGoProvider { pub fn new() -> Self { Self { @@ -64,386 +135,46 @@ impl OpenCodeGoProvider { } } - fn workspace_id_from_context(workspace_id: Option<&str>) -> Option<&str> { - workspace_id.filter(|id| !id.is_empty()) - } - - async fn fetch_workspace_id( - client: &Client, - cookie_header: &str, - ) -> Result { - let url = format!("{}?id={}", SERVER_URL, WORKSPACES_SERVER_ID); - let response = client - .get(&url) - .header("Cookie", cookie_header) - .header("X-Server-Id", WORKSPACES_SERVER_ID) - .header("X-Server-Instance", format!("server-fn:{}", Uuid::new_v4())) - .header("User-Agent", USER_AGENT) - .header("Origin", BASE_URL) - .header("Referer", BASE_URL) - .header( - "Accept", - "text/javascript, application/json;q=0.9, */*;q=0.8", - ) - .send() - .await?; - - let status = response.status(); - if !status.is_success() { - if status.as_u16() == 401 || status.as_u16() == 403 { - return Err(ProviderError::AuthRequired); - } - return Err(ProviderError::Other(format!( - "OpenCode workspace API returned {}", - status - ))); - } - - let text = response.text().await?; - if Self::looks_signed_out(&text) { - return Err(ProviderError::AuthRequired); - } - - let ids = Self::parse_workspace_ids(&text); - ids.into_iter() - .next() - .ok_or_else(|| ProviderError::Parse("No workspace ID found".to_string())) - } - - async fn fetch_usage_page( - client: &Client, - workspace_id: &str, - cookie_header: &str, - ) -> Result { - let url = format!("{}/workspace/{}/go", BASE_URL, workspace_id); - Self::fetch_page_text(client, &url, cookie_header, None, "usage page").await - } - - /// GET a page with the standard browser-ish headers; `timeout` overrides - /// the client default when the caller is inside a smaller budget. - async fn fetch_page_text( - client: &Client, - url: &str, - cookie_header: &str, - timeout: Option, - what: &str, - ) -> Result { - let mut request = client - .get(url) - .header("Cookie", cookie_header) - .header("User-Agent", USER_AGENT) - .header("Referer", BASE_URL) - .header( - "Accept", - "text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8", - ); - if let Some(timeout) = timeout { - request = request.timeout(timeout); - } - let response = request.send().await?; - - let status = response.status(); - if !status.is_success() { - if status.as_u16() == 401 || status.as_u16() == 403 { - return Err(ProviderError::AuthRequired); - } - return Err(ProviderError::Other(format!( - "OpenCode Go {what} returned {status}" - ))); - } - - let text = response.text().await?; - if Self::looks_signed_out(&text) { - return Err(ProviderError::AuthRequired); - } - Ok(text) - } - - fn parse_usage_text(text: &str) -> Result { - let now = Utc::now(); - - let rolling = Self::extract_window(text, &["rollingUsage", "rolling_usage", "rolling"]) - .ok_or_else(|| ProviderError::Parse("Missing rolling usage window".to_string()))?; - let weekly = Self::extract_window(text, &["weeklyUsage", "weekly_usage", "weekly"]); - let monthly = Self::extract_window(text, &["monthlyUsage", "monthly_usage", "monthly"]); - - let primary = RateWindow::with_details( - rolling.0, - Some(300), - Some(now + chrono::Duration::seconds(rolling.1)), - None, - ); - let mut snap = UsageSnapshot::new(primary).with_login_method("OpenCode Go"); - - if let Some((pct, reset)) = weekly { - snap = snap.with_secondary(RateWindow::with_details( - pct, - Some(10080), - Some(now + chrono::Duration::seconds(reset)), - None, - )); - } - - if let Some((pct, reset)) = monthly { - let resets_at = now + chrono::Duration::seconds(reset); - snap = snap.with_tertiary(RateWindow::with_details( - pct, - RateWindow::monthly_window_minutes(Some(resets_at)).or(Some(43200)), - Some(resets_at), - None, - )); - } - - if let Some(renews_at) = super::extract_renewal(text) { - snap = snap.with_extra_rate_window( - "renewal", - "Renews", - RateWindow::with_details(0.0, None, Some(renews_at), None), - ); - } - - Ok(snap) - } - - /// Extract `(percent, resetInSec)` for a usage block by name. - fn extract_window(text: &str, names: &[&str]) -> Option<(f64, i64)> { - for name in names { - let percent_pattern = format!( - r#"{}[^}}]*?(?:usagePercent|usedPercent|percentUsed|percent)\s*[:=]\s*([0-9]+(?:\.[0-9]+)?)"#, - name - ); - let reset_pattern = format!( - r#"{}[^}}]*?(?:resetInSec|resetInSeconds|resetSeconds|resetSec)\s*[:=]\s*([0-9]+)"#, - name - ); - - let percent = super::extract_number(&percent_pattern, text); - if let Some(p) = percent { - #[allow( - clippy::cast_possible_truncation, - reason = "resetInSec values are whole-second counts scraped as integral numbers" - )] - let reset = super::extract_number(&reset_pattern, text) - .map(|n| n as i64) - .unwrap_or(0); - // Direct percent fields arrive as integer percent in the serialized payload; no fraction scaling (upstream parseSubscription parity; win-fork #247). - return Some((p.clamp(0.0, 100.0), reset.max(0))); - } - - // Computed used/limit (already 0..100) — do not apply fraction *100. - let used_pattern = format!( - r#"{}[^}}]*?(?:used|usage|consumed)\s*[:=]\s*([0-9]+(?:\.[0-9]+)?)"#, - name - ); - let limit_pattern = format!( - r#"{}[^}}]*?(?:limit|total|allowance)\s*[:=]\s*([0-9]+(?:\.[0-9]+)?)"#, - name - ); - if let (Some(used), Some(limit)) = ( - super::extract_number(&used_pattern, text), - super::extract_number(&limit_pattern, text), - ) && limit > 0.0 - { - #[allow( - clippy::cast_possible_truncation, - reason = "resetInSec values are whole-second counts scraped as integral numbers" - )] - let reset = super::extract_number(&reset_pattern, text) - .map(|n| n as i64) - .unwrap_or(0); - let p = (used / limit) * 100.0; - return Some((p.clamp(0.0, 100.0), reset.max(0))); - } - } - None - } - - fn parse_workspace_ids(text: &str) -> Vec { - let pattern = r#"(wrk_[A-Za-z0-9_-]+)"#; - let re = match regex_lite::Regex::new(pattern) { - Ok(r) => r, - Err(_) => return vec![], - }; - let mut seen = Vec::new(); - for caps in re.captures_iter(text) { - if let Some(m) = caps.get(1) { - let s = m.as_str().to_string(); - if !seen.contains(&s) { - seen.push(s); - } - } - } - seen - } - - fn looks_signed_out(text: &str) -> bool { - let lower = text.to_lowercase(); - lower.contains("auth/authorize") - || lower.contains("\"signin\"") - || lower.contains("please sign in") - } - - fn parse_zen_balance(text: &str) -> Option { - if let Ok(json) = serde_json::from_str::(text) - && let Some(value) = Self::find_balance_value(&json) - { - return Some(value); - } - let patterns = [ - r#"(?i)(?:current\s+balance|zen\s+balance|現在の残高)[^$]{0,80}\$\s*([0-9][0-9,]*(?:\.[0-9]+)?)"#, - r#"(?i)(?:balance|残高)[\s\S]{0,120}?\$\s*([0-9][0-9,]*(?:\.[0-9]+)?)"#, - ]; - patterns.iter().find_map(|pattern| { - let re = regex_lite::Regex::new(pattern).ok()?; - let raw = re.captures(text)?.get(1)?.as_str().replace(',', ""); - raw.parse::().ok() - }) - } - - fn find_balance_value(value: &serde_json::Value) -> Option { - match value { - serde_json::Value::Object(map) => { - for (key, value) in map { - let normalized: String = key - .to_lowercase() - .chars() - .filter(|c| c.is_ascii_alphanumeric()) - .collect(); - if matches!( - normalized.as_str(), - "zenbalance" - | "zencurrentbalance" - | "currentbalance" - | "currentbalanceusd" - | "balanceusd" - | "usdbalance" - ) { - if let Some(number) = value.as_f64() { - return Some(number); - } - if let Some(text) = value.as_str() - && let Ok(number) = text.trim().replace(',', "").parse() - { - return Some(number); - } - } - if let Some(found) = Self::find_balance_value(value) { - return Some(found); - } - } - None - } - serde_json::Value::Array(items) => items.iter().find_map(Self::find_balance_value), - _ => None, - } - } - - // ── Zen balance (upstream 0.48.0 #2583) ───────────────────────────────── - - /// Zen dashboard page URL (upstream `zenDashboardURL`). - fn zen_dashboard_url(workspace_id: &str) -> String { - format!("{BASE_URL}/workspace/{workspace_id}") - } - - /// Fetch a server-fn endpoint (same call shape as `fetch_workspace_id`: - /// `GET {SERVER_URL}?id=…&args=…` with the server-fn headers). - async fn fetch_server_text( - client: &Client, - server_id: &str, - args: Option<&str>, - referer: &str, - cookie_header: &str, - timeout: Option, - ) -> Result { - let mut url = reqwest::Url::parse(SERVER_URL) - .map_err(|e| ProviderError::Parse(format!("Invalid OpenCode server URL: {e}")))?; - url.query_pairs_mut().append_pair("id", server_id); - if let Some(args) = args { - url.query_pairs_mut().append_pair("args", args); - } - let mut request = client - .get(url) - .header("Cookie", cookie_header) - .header("X-Server-Id", server_id) - .header("X-Server-Instance", format!("server-fn:{}", Uuid::new_v4())) - .header("User-Agent", USER_AGENT) - .header("Origin", BASE_URL) - .header("Referer", referer) - .header( - "Accept", - "text/javascript, application/json;q=0.9, */*;q=0.8", - ); - if let Some(timeout) = timeout { - request = request.timeout(timeout); - } - let response = request.send().await?; - let status = response.status(); - if !status.is_success() { - if status.as_u16() == 401 || status.as_u16() == 403 { - return Err(ProviderError::AuthRequired); - } - return Err(ProviderError::Other(format!( - "OpenCode Go server returned {status}" - ))); - } - let text = response.text().await?; - if Self::looks_signed_out(&text) { - return Err(ProviderError::AuthRequired); - } - Ok(text) + fn workspace_id_from_context(workspace_id: Option<&str>) -> Option { + console::normalize_workspace_id(workspace_id) } /// Upstream `fetchZenBalance`: the dashboard HTML embeds the balance for /// some page states; the dedicated billing server-fn report (raw 1e-8 USD /// units behind a customerID marker) is the fallback. Optional enrichment /// — every failure degrades to `None`, never to a fetch error. - async fn fetch_zen_balance( - client: &Client, + async fn fetch_zen_balance( + transport: &T, workspace_id: &str, - cookie_header: &str, + cookies: &WebCookieSession, timeout: Duration, - ) -> Option { + ) -> OptionalZenBalance { let request_timeout = timeout.min(ZEN_BALANCE_TIMEOUT); - let referer = Self::zen_dashboard_url(workspace_id); - - let page = Self::fetch_page_text( - client, - &referer, - cookie_header, - Some(request_timeout), - "Zen dashboard page", - ) - .await - .ok()?; - if let Some(balance) = Self::parse_zen_balance(&page) { - return Some(balance); + match transport + .fetch_console_balance(workspace_id, cookies.header(), request_timeout) + .await + { + Ok(balance) => OptionalZenBalance::Resolved(balance), + Err(error) if cookies.can_recover_with_legacy(&error) => { + OptionalZenBalance::LegacyBalanceRequired + } + Err(_) => OptionalZenBalance::Resolved(None), } - - let args = serde_json::json!([workspace_id]).to_string(); - let billing = Self::fetch_server_text( - client, - BILLING_SERVER_ID, - Some(&args), - &referer, - cookie_header, - Some(request_timeout), - ) - .await - .ok()?; - parse_billing_server_balance(&billing) } /// Spawn the optional Zen balance task (25 ms start delay so the usage /// fetch gets the head start, per upstream). Resolves the workspace id /// inside the task when no override is pinned. - fn spawn_zen_balance_task( - &self, - cookie_header: &str, + fn spawn_zen_balance_task( + transport: Arc, + cookies: &WebCookieSession, workspace_id_override: Option<&str>, web_timeout: u64, - ) -> (tokio::task::JoinHandle>, std::time::Instant) { - let client = self.client.clone(); - let cookie_header = cookie_header.to_string(); + ) -> ( + tokio::task::JoinHandle, + std::time::Instant, + ) { + let cookies = cookies.clone(); let workspace_id_override = workspace_id_override.map(str::to_string); let timeout = Duration::from_secs(web_timeout.max(1)); let started_at = std::time::Instant::now(); @@ -451,11 +182,40 @@ impl OpenCodeGoProvider { tokio::time::sleep(ZEN_BALANCE_START_DELAY).await; let workspace_id = match workspace_id_override { Some(id) => id, - None => Self::fetch_workspace_id(&client, &cookie_header) + None => match transport + .fetch_workspace_id(cookies.header(), Duration::from_secs(30)) .await - .ok()?, + { + Ok(id) => id, + Err(error) if cookies.can_recover_with_legacy(&error) => { + return OptionalZenBalance::LegacyBalanceRequired; + } + Err(_) => return OptionalZenBalance::Resolved(None), + }, }; - Self::fetch_zen_balance(&client, &workspace_id, &cookie_header, timeout).await + Self::fetch_zen_balance(transport.as_ref(), &workspace_id, &cookies, timeout).await + }); + (task, started_at) + } + + fn spawn_legacy_balance_task( + transport: Arc, + session: T::LegacySession, + web_timeout: u64, + ) -> ( + tokio::task::JoinHandle, + std::time::Instant, + ) { + let timeout = Duration::from_secs(web_timeout.max(1)).min(ZEN_BALANCE_TIMEOUT); + let started_at = std::time::Instant::now(); + let task = tokio::spawn(async move { + tokio::time::sleep(ZEN_BALANCE_START_DELAY).await; + OptionalZenBalance::Resolved( + transport + .fetch_legacy_balance(&session, timeout) + .await + .unwrap_or(None), + ) }); (task, started_at) } @@ -463,13 +223,13 @@ impl OpenCodeGoProvider { /// Join the optional Zen balance task within the policy budget. A budget /// expiry cancels the in-flight HTTP work instead of leaking it. async fn join_zen_balance( - mut task: tokio::task::JoinHandle>, + mut task: tokio::task::JoinHandle, started_at: std::time::Instant, requires_optional_usage_completeness: bool, - ) -> Option { + ) -> Option { let budget = zen_balance_join_budget(started_at, requires_optional_usage_completeness); match tokio::time::timeout(budget, &mut task).await { - Ok(Ok(balance)) => balance, + Ok(Ok(balance)) => Some(balance), Ok(Err(_)) | Err(_) => { task.abort(); None @@ -477,43 +237,139 @@ impl OpenCodeGoProvider { } } + async fn resolve_legacy_balance( + transport: Arc, + cookies: &WebCookieSession, + started_at: std::time::Instant, + requires_optional_usage_completeness: bool, + ) -> Option { + let budget = zen_balance_join_budget(started_at, requires_optional_usage_completeness); + if budget.is_zero() { + return None; + } + tokio::time::timeout(budget, async { + let session = transport.resolve_legacy_session(cookies).await.ok()?; + transport + .fetch_legacy_balance(&session, budget.min(ZEN_BALANCE_TIMEOUT)) + .await + .unwrap_or(None) + }) + .await + .ok() + .flatten() + } + + async fn finish_zen_balance( + transport: Arc, + cookies: &WebCookieSession, + task: tokio::task::JoinHandle, + started_at: std::time::Instant, + requires_optional_usage_completeness: bool, + ) -> Option { + match Self::join_zen_balance(task, started_at, requires_optional_usage_completeness).await? + { + OptionalZenBalance::Resolved(balance) => balance, + OptionalZenBalance::LegacyBalanceRequired => { + Self::resolve_legacy_balance( + transport, + cookies, + started_at, + requires_optional_usage_completeness, + ) + .await + } + } + } + async fn fetch_with_cookies( &self, ctx: &FetchContext, cookie_header: &str, ) -> Result { + let transport = Arc::new(HttpWebTransport::new(self.client.clone())); + Self::fetch_with_transport(ctx, cookie_header, transport).await + } + + async fn fetch_with_transport( + ctx: &FetchContext, + cookie_header: &str, + transport: Arc, + ) -> Result { + let cookies = WebCookieSession::new(cookie_header); let workspace_id = match Self::workspace_id_from_context(ctx.workspace_id.as_deref()) { - Some(workspace_id) => workspace_id.to_string(), - None => Self::fetch_workspace_id(&self.client, cookie_header).await?, + Some(workspace_id) => workspace_id, + None => match transport + .fetch_workspace_id(cookies.header(), Duration::from_secs(30)) + .await + { + Ok(workspace_id) => workspace_id, + Err(console_error) if cookies.can_recover_with_legacy(&console_error) => { + return Self::fetch_legacy_with_cookies( + ctx, + &cookies, + console_error, + transport, + ) + .await; + } + Err(error) => return Err(error), + }, }; // F15 (#2583): start the optional Zen balance fetch in parallel with the // usage page and bound the join from task creation, so a slow balance // still lands in CLI/serve usage reads without stacking a second wait. - let (zen_task, zen_started) = - self.spawn_zen_balance_task(cookie_header, Some(&workspace_id), ctx.web_timeout); - let page = match Self::fetch_usage_page(&self.client, &workspace_id, cookie_header).await { - Ok(page) => page, - Err(err) => { + let (zen_task, zen_started) = Self::spawn_zen_balance_task( + Arc::clone(&transport), + &cookies, + Some(&workspace_id), + ctx.web_timeout, + ); + let console_result = transport + .fetch_console_usage( + &workspace_id, + cookies.header(), + Duration::from_secs(ctx.web_timeout.max(1)), + ) + .await; + let (usage, embedded_balance) = match console_result { + Ok(console::ConsoleUsage::Snapshot(usage)) => (*usage, None), + Ok(console::ConsoleUsage::NoSubscription) => { zen_task.abort(); - return Err(err); + return Err(ProviderError::Parse( + "No OpenCode Go subscription is available".to_string(), + )); } - }; - let usage = match Self::parse_usage_text(&page) { - Ok(usage) => usage, - Err(err) => { + Err(console_error) if cookies.can_recover_with_legacy(&console_error) => { zen_task.abort(); - return Err(err); + let session = match transport.resolve_legacy_session(&cookies).await { + Ok(session) => session, + Err(error) => { + return cookies.select_legacy_result(console_error, Err(error)); + } + }; + return Self::fetch_legacy_with_session( + ctx, + &cookies, + console_error, + transport, + session, + ) + .await; + } + Err(error) => { + zen_task.abort(); + return Err(error); } }; - // The /go page states embed the balance for some deployments — the - // zero-cost parse wins over the dedicated fetch when it works. - let balance = match Self::parse_zen_balance(&page) { + let balance = match embedded_balance { Some(balance) => { zen_task.abort(); Some(balance) } None => { - Self::join_zen_balance( + Self::finish_zen_balance( + transport, + &cookies, zen_task, zen_started, ctx.requires_optional_usage_completeness, @@ -524,6 +380,61 @@ impl OpenCodeGoProvider { Ok(Self::with_zen_balance(usage, "web", balance)) } + async fn fetch_legacy_with_cookies( + ctx: &FetchContext, + cookies: &WebCookieSession, + console_error: ProviderError, + transport: Arc, + ) -> Result { + let session = match transport.resolve_legacy_session(cookies).await { + Ok(session) => session, + Err(error) => return cookies.select_legacy_result(console_error, Err(error)), + }; + Self::fetch_legacy_with_session(ctx, cookies, console_error, transport, session).await + } + + async fn fetch_legacy_with_session( + ctx: &FetchContext, + cookies: &WebCookieSession, + console_error: ProviderError, + transport: Arc, + session: T::LegacySession, + ) -> Result { + let (zen_task, zen_started) = Self::spawn_legacy_balance_task( + Arc::clone(&transport), + session.clone(), + ctx.web_timeout, + ); + let result = match cookies + .select_legacy_result(console_error, transport.fetch_legacy_usage(&session).await) + { + Ok(result) => result, + Err(error) => { + zen_task.abort(); + return Err(error); + } + }; + let balance = match result.embedded_balance { + Some(balance) => { + zen_task.abort(); + Some(balance) + } + None => { + match Self::join_zen_balance( + zen_task, + zen_started, + ctx.requires_optional_usage_completeness, + ) + .await + { + Some(OptionalZenBalance::Resolved(balance)) => balance, + Some(OptionalZenBalance::LegacyBalanceRequired) | None => None, + } + } + }; + Ok(Self::with_zen_balance(result.usage, "web", balance)) + } + /// Attach an optional Zen balance to the snapshot: informational extra /// window plus the cost row (existing bridge shape). fn with_zen_balance( @@ -686,13 +597,22 @@ impl OpenCodeGoProvider { let Some(cookie_header) = cookie_header else { return Ok(result); }; - let (task, started) = self.spawn_zen_balance_task( - &cookie_header, + let cookies = WebCookieSession::new(&cookie_header); + let transport = Arc::new(HttpWebTransport::new(self.client.clone())); + let (task, started) = Self::spawn_zen_balance_task( + Arc::clone(&transport), + &cookies, ctx.workspace_id.as_deref(), ctx.web_timeout, ); - if let Some(balance) = - Self::join_zen_balance(task, started, ctx.requires_optional_usage_completeness).await + if let Some(balance) = Self::finish_zen_balance( + transport, + &cookies, + task, + started, + ctx.requires_optional_usage_completeness, + ) + .await { result = Self::with_zen_balance(result.usage, &result.source_label.clone(), Some(balance)); @@ -729,254 +649,5 @@ fn zen_balance_join_budget( ZEN_BALANCE_TIMEOUT.saturating_sub(started_at.elapsed()) } -/// Upstream `parseBillingServerResponse`: the billing server-fn reports the -/// balance in raw 1e-8 USD units behind a customerID marker; JSON tree first, -/// RSC-streamed text fragment second. -fn parse_billing_server_balance(text: &str) -> Option { - const BILLING_SCALE: f64 = 100_000_000.0; - if let Ok(json) = serde_json::from_str::(text) - && let Some(raw) = find_raw_billing_balance(&json) - { - return Some(raw / BILLING_SCALE); - } - let customer_re = regex_lite::Regex::new( - r#"(?:"customerID"|customerID)\s*:\s*(?:\$R\[\d+\]\s*=\s*)?"[^"]+""#, - ) - .ok()?; - customer_re.find(text)?; - let balance_re = regex_lite::Regex::new( - r#"(?:"balance"|balance)\s*:\s*(?:\$R\[\d+\]\s*=\s*)?(-?[0-9]+(?:\.[0-9]+)?)"#, - ) - .ok()?; - let raw: f64 = balance_re - .captures(text)? - .get(1)? - .as_str() - .replace(',', "") - .parse() - .ok()?; - Some(raw / BILLING_SCALE) -} - -/// Find a `balance` value guarded by a non-empty `customerID` sibling -/// (upstream `findRawBillingBalance`). An object holding a `balance` key -/// decides terminally — no deeper search below it once the guard fails the -/// value. Booleans are excluded like upstream's `doubleValue`. -fn find_raw_billing_balance(value: &serde_json::Value) -> Option { - match value { - serde_json::Value::Object(map) => { - if let Some(balance) = map.get("balance") { - let customer_ok = map - .get("customerID") - .and_then(|v| v.as_str()) - .is_some_and(|id| !id.is_empty()); - if !customer_ok { - return None; - } - return billing_numeric_value(balance); - } - map.values().find_map(find_raw_billing_balance) - } - serde_json::Value::Array(items) => items.iter().find_map(find_raw_billing_balance), - _ => None, - } -} - -fn billing_numeric_value(value: &serde_json::Value) -> Option { - match value { - serde_json::Value::Number(n) => n.as_f64(), - serde_json::Value::String(s) => s.trim().replace(',', "").parse().ok(), - _ => None, - } -} - #[cfg(test)] -mod tests { - use super::*; - - #[test] - fn parses_workspace_ids() { - let text = r#"{ id: "wrk_abc123", name: "x" } { id: "wrk_def456" }"#; - let ids = OpenCodeGoProvider::parse_workspace_ids(text); - assert_eq!( - ids, - vec!["wrk_abc123".to_string(), "wrk_def456".to_string()] - ); - } - - #[test] - fn uses_context_workspace_id_before_discovery() { - assert_eq!( - OpenCodeGoProvider::workspace_id_from_context(Some("wrk_override")), - Some("wrk_override") - ); - assert_eq!( - OpenCodeGoProvider::workspace_id_from_context(Some("")), - None - ); - } - - #[test] - fn sub_one_percent_computed_used_limit_is_not_rescaled() { - let text = r#" - rollingUsage: { used: 1, limit: 100, resetInSec: 600 } - weeklyUsage: { used: 1, limit: 200, resetInSec: 86400 } - "#; - let snap = OpenCodeGoProvider::parse_usage_text(text).unwrap(); - assert!((snap.primary.used_percent - 1.0).abs() < 0.001); - assert!((snap.secondary.as_ref().unwrap().used_percent - 0.5).abs() < 0.001); - } - - #[test] - fn selected_token_auth_failure_does_not_fall_back_to_local_estimate() { - let mut selected = FetchContext { - auto_prefer_web: true, - ..FetchContext::default() - }; - assert!(!OpenCodeGoProvider::web_error_allows_local_fallback( - &selected, - &ProviderError::AuthRequired - )); - - selected.auto_prefer_web = false; - selected.workspace_id = Some("wrk_example".to_string()); - assert!(OpenCodeGoProvider::web_error_allows_local_fallback( - &selected, - &ProviderError::AuthRequired - )); - } - - #[test] - fn parses_usage_blocks() { - let text = r#" - rollingUsage: { usagePercent: 42.5, resetInSec: 3600 } - weeklyUsage: { usagePercent: 13, resetInSec: 86400 } - monthlyUsage: { usagePercent: 7, resetInSec: 2592000 } - "#; - let snap = OpenCodeGoProvider::parse_usage_text(text).unwrap(); - assert!((snap.primary.used_percent - 42.5).abs() < 0.001); - let secondary = snap.secondary.expect("weekly"); - // usagePercent: 13 is a direct integer percent → 13% - assert!((secondary.used_percent - 13.0).abs() < 0.001); - let tertiary = snap.tertiary.expect("monthly"); - assert!((tertiary.used_percent - 7.0).abs() < 0.001); - let expected = RateWindow::monthly_window_minutes(tertiary.resets_at).or(Some(43200)); - assert_eq!(tertiary.window_minutes, expected); - assert!(tertiary.resets_at.is_some()); - } - - #[test] - fn direct_percent_one_is_not_rescaled() { - let text = r#"rollingUsage:$R[34]={status:"ok",resetInSec:13631,usagePercent:1} weeklyUsage:$R[35]={status:"ok",resetInSec:53863,usagePercent:15}"#; - let snap = OpenCodeGoProvider::parse_usage_text(text).unwrap(); - assert!((snap.primary.used_percent - 1.0).abs() < 0.001); - assert!((snap.secondary.as_ref().unwrap().used_percent - 15.0).abs() < 0.001); - } - - #[test] - fn parses_renewal_window() { - let text = r#" - rollingUsage: { usagePercent: 42.5, resetInSec: 3600 } - weeklyUsage: { usagePercent: 50, resetInSec: 86400 } - renewAt: "2026-06-01T12:00:00Z" - "#; - let snap = OpenCodeGoProvider::parse_usage_text(text).unwrap(); - let renewal = snap - .extra_rate_windows - .iter() - .find(|window| window.id == "renewal") - .expect("renewal window"); - assert_eq!(renewal.title, "Renews"); - assert_eq!( - renewal.window.resets_at.unwrap().to_rfc3339(), - "2026-06-01T12:00:00+00:00" - ); - } - - // ── F15: bounded optional Zen balance wait (upstream #2583) ─── - - #[test] - fn zen_join_budget_grace_vs_completeness() { - let started = std::time::Instant::now(); - // Background/UI reads keep the short join grace. - assert_eq!( - zen_balance_join_budget(started, false), - Duration::from_millis(250) - ); - // Completeness reads get the remainder of the 5 s optional-balance - // budget measured from task creation. - let budget = zen_balance_join_budget(started, true); - assert!(budget <= ZEN_BALANCE_TIMEOUT, "{budget:?}"); - assert!(budget > Duration::from_secs(4), "{budget:?}"); - // An already-exhausted budget joins immediately. - let stale = std::time::Instant::now() - .checked_sub(Duration::from_secs(60)) - .unwrap(); - assert_eq!(zen_balance_join_budget(stale, true), Duration::ZERO); - } - - #[tokio::test] - async fn slow_zen_task_is_abandoned_within_grace() { - let started = std::time::Instant::now(); - let task = tokio::spawn(async { - tokio::time::sleep(Duration::from_secs(30)).await; - Some(42.5) - }); - // UI grace (250 ms) never waits out a 30 s balance fetch. - let balance = OpenCodeGoProvider::join_zen_balance(task, started, false).await; - assert_eq!(balance, None); - } - - #[tokio::test] - async fn fast_zen_task_lands_in_completeness_budget() { - let started = std::time::Instant::now(); - let task = tokio::spawn(async { - tokio::time::sleep(Duration::from_millis(10)).await; - Some(42.5) - }); - let balance = OpenCodeGoProvider::join_zen_balance(task, started, true).await; - assert_eq!(balance, Some(42.5)); - } - - #[test] - fn billing_server_balance_needs_customer_marker() { - // Raw 1e-8-scaled balance behind a customerID → USD. - assert_eq!( - parse_billing_server_balance(r#"{"balance": 1500000000, "customerID": "cus_123"}"#), - Some(15.0) - ); - // Same shape without the marker is not a billing payload. - assert_eq!( - parse_billing_server_balance(r#"{"balance": 1500000000}"#), - None - ); - // Non-numeric balance with marker → no result (upstream terminal guard). - assert_eq!( - parse_billing_server_balance(r#"{"balance": true, "customerID": "cus_123"}"#), - None - ); - } - - #[test] - fn billing_server_balance_handles_strings_nesting_and_rsc_fragments() { - // Numeric strings coerce. - assert_eq!( - parse_billing_server_balance(r#"{"balance": "1,000,000,000", "customerID": "cus_1"}"#), - Some(10.0) - ); - // Nested containers search through. - assert_eq!( - parse_billing_server_balance( - r#"{"data": {"rows": [{"balance": 250000000, "customerID": "cus_2"}]}}"# - ), - Some(2.5) - ); - // RSC-streamed fragment: marker plus plain balance pair. - assert_eq!( - parse_billing_server_balance(r#"customerID:$R[1] = "cus_9"; "balance": -500000000"#), - Some(-5.0) - ); - // Marker alone is not enough. - assert_eq!(parse_billing_server_balance(r#"customerID: "cus_9""#), None); - } -} +mod tests; diff --git a/rust/src/providers/opencodego/tests.rs b/rust/src/providers/opencodego/tests.rs new file mode 100644 index 0000000000..880fd87e2b --- /dev/null +++ b/rust/src/providers/opencodego/tests.rs @@ -0,0 +1,307 @@ +use super::*; +use std::sync::{ + Arc, Mutex, + atomic::{AtomicUsize, Ordering}, +}; + +#[derive(Clone, Debug, Eq, PartialEq)] +struct FakeLegacySession { + workspace_id: String, +} + +struct FakeWebTransport { + console_failures: tokio::sync::Barrier, + resolver_calls: AtomicUsize, + console_timeouts: Mutex>, + usage_sessions: Mutex>, + balance_sessions: Mutex>, + legacy_balance_timeouts: Mutex>, + fail_legacy_resolution: bool, +} + +impl FakeWebTransport { + fn new(fail_legacy_resolution: bool) -> Arc { + Arc::new(Self { + console_failures: tokio::sync::Barrier::new(2), + resolver_calls: AtomicUsize::new(0), + console_timeouts: Mutex::new(Vec::new()), + usage_sessions: Mutex::new(Vec::new()), + balance_sessions: Mutex::new(Vec::new()), + legacy_balance_timeouts: Mutex::new(Vec::new()), + fail_legacy_resolution, + }) + } +} + +#[async_trait] +impl WebTransport for FakeWebTransport { + type LegacySession = FakeLegacySession; + + async fn fetch_workspace_id( + &self, + _cookie_header: &str, + _timeout: Duration, + ) -> Result { + panic!("workspace override should bypass discovery") + } + + async fn fetch_console_usage( + &self, + _workspace_id: &str, + _cookie_header: &str, + timeout: Duration, + ) -> Result { + self.console_timeouts.lock().unwrap().push(timeout); + self.console_failures.wait().await; + Err(ProviderError::Other( + "console usage unavailable".to_string(), + )) + } + + async fn fetch_console_balance( + &self, + _workspace_id: &str, + _cookie_header: &str, + timeout: Duration, + ) -> Result, ProviderError> { + self.console_timeouts.lock().unwrap().push(timeout); + self.console_failures.wait().await; + Err(ProviderError::Other( + "console balance unavailable".to_string(), + )) + } + + async fn resolve_legacy_session( + &self, + _cookies: &WebCookieSession, + ) -> Result { + self.resolver_calls.fetch_add(1, Ordering::SeqCst); + if self.fail_legacy_resolution { + return Err(ProviderError::AuthRequired); + } + Ok(FakeLegacySession { + workspace_id: "wrk_legacy".to_string(), + }) + } + + async fn fetch_legacy_usage( + &self, + session: &Self::LegacySession, + ) -> Result { + self.usage_sessions + .lock() + .unwrap() + .push(session.workspace_id.clone()); + Ok(legacy::LegacyUsage { + usage: UsageSnapshot::new(RateWindow::with_details(25.0, None, None, None)), + embedded_balance: None, + }) + } + + async fn fetch_legacy_balance( + &self, + session: &Self::LegacySession, + timeout: Duration, + ) -> Result, ProviderError> { + self.legacy_balance_timeouts.lock().unwrap().push(timeout); + self.balance_sessions + .lock() + .unwrap() + .push(session.workspace_id.clone()); + Ok(Some(42.5)) + } +} + +#[test] +fn cookie_session_classifies_transport_capabilities_once() { + let both = WebCookieSession::new("auth=legacy; __Host-console_session=console"); + assert_eq!( + both.capabilities, + CookieCapabilities { + console: true, + legacy: true, + } + ); + assert_eq!( + WebCookieSession::new("__Host-console_session=console").capabilities, + CookieCapabilities { + console: true, + legacy: false, + } + ); + assert_eq!( + WebCookieSession::new("auth=; __Host-console_session=").capabilities, + CookieCapabilities { + console: false, + legacy: false, + } + ); +} + +#[test] +fn legacy_recovery_and_error_precedence_follow_cookie_capabilities() { + let both = WebCookieSession::new("auth=legacy; __Host-console_session=console"); + assert!(both.can_recover_with_legacy(&ProviderError::AuthRequired)); + assert!( + !WebCookieSession::new("__Host-console_session=console") + .can_recover_with_legacy(&ProviderError::AuthRequired) + ); + assert!(!both.can_recover_with_legacy(&ProviderError::NotInstalled("terminal".to_string()))); + + let error = both + .select_legacy_result::<()>( + ProviderError::Other("console unavailable".to_string()), + Err(ProviderError::AuthRequired), + ) + .unwrap_err(); + assert!(matches!(error, ProviderError::Other(message) if message == "console unavailable")); + + let error = both + .select_legacy_result::<()>( + ProviderError::AuthRequired, + Err(ProviderError::Parse("legacy payload missing".to_string())), + ) + .unwrap_err(); + assert!(matches!(error, ProviderError::Parse(message) if message == "legacy payload missing")); +} + +#[test] +fn uses_context_workspace_id_before_discovery() { + assert_eq!( + OpenCodeGoProvider::workspace_id_from_context(Some("wrk_override")), + Some("wrk_override".to_string()) + ); + assert_eq!( + OpenCodeGoProvider::workspace_id_from_context(Some("")), + None + ); +} + +#[test] +fn selected_token_auth_failure_does_not_fall_back_to_local_estimate() { + let mut selected = FetchContext { + auto_prefer_web: true, + ..FetchContext::default() + }; + assert!(!OpenCodeGoProvider::web_error_allows_local_fallback( + &selected, + &ProviderError::AuthRequired + )); + + selected.auto_prefer_web = false; + selected.workspace_id = Some("wrk_example".to_string()); + assert!(OpenCodeGoProvider::web_error_allows_local_fallback( + &selected, + &ProviderError::AuthRequired + )); +} + +// ── F15: bounded optional Zen balance wait (upstream #2583) ─── + +#[test] +fn zen_join_budget_grace_vs_completeness() { + let started = std::time::Instant::now(); + // Background/UI reads keep the short join grace. + assert_eq!( + zen_balance_join_budget(started, false), + Duration::from_millis(250) + ); + // Completeness reads get the remainder of the 5 s optional-balance + // budget measured from task creation. + let budget = zen_balance_join_budget(started, true); + assert!(budget <= ZEN_BALANCE_TIMEOUT, "{budget:?}"); + assert!(budget > Duration::from_secs(4), "{budget:?}"); + // An already-exhausted budget joins immediately. + let stale = std::time::Instant::now() + .checked_sub(Duration::from_secs(60)) + .unwrap(); + assert_eq!(zen_balance_join_budget(stale, true), Duration::ZERO); +} + +#[tokio::test] +async fn slow_zen_task_is_abandoned_within_grace() { + let started = std::time::Instant::now(); + let task = tokio::spawn(async { + tokio::time::sleep(Duration::from_secs(30)).await; + OptionalZenBalance::Resolved(Some(42.5)) + }); + // UI grace (250 ms) never waits out a 30 s balance fetch. + let balance = OpenCodeGoProvider::join_zen_balance(task, started, false).await; + assert_eq!(balance, None); +} + +#[tokio::test] +async fn fast_zen_task_lands_in_completeness_budget() { + let started = std::time::Instant::now(); + let task = tokio::spawn(async { + tokio::time::sleep(Duration::from_millis(10)).await; + OptionalZenBalance::Resolved(Some(42.5)) + }); + let balance = OpenCodeGoProvider::join_zen_balance(task, started, true).await; + assert_eq!(balance, Some(OptionalZenBalance::Resolved(Some(42.5)))); +} + +#[tokio::test] +async fn simultaneous_console_failures_resolve_one_legacy_session() { + let transport = FakeWebTransport::new(false); + let ctx = FetchContext { + workspace_id: Some("wrk_console".to_string()), + web_timeout: 1, + requires_optional_usage_completeness: true, + ..FetchContext::default() + }; + + let result = OpenCodeGoProvider::fetch_with_transport( + &ctx, + "auth=legacy; __Host-console_session=console", + Arc::clone(&transport), + ) + .await + .unwrap(); + + assert_eq!(transport.resolver_calls.load(Ordering::SeqCst), 1); + assert_eq!( + transport.usage_sessions.lock().unwrap().as_slice(), + ["wrk_legacy"] + ); + assert_eq!( + transport.balance_sessions.lock().unwrap().as_slice(), + ["wrk_legacy"] + ); + assert_eq!( + transport.console_timeouts.lock().unwrap().as_slice(), + [Duration::from_secs(1), Duration::from_secs(1)] + ); + assert_eq!( + transport.legacy_balance_timeouts.lock().unwrap().as_slice(), + [Duration::from_secs(1)] + ); + assert_eq!(result.cost.unwrap().used, 42.5); +} + +#[tokio::test] +async fn production_fallback_preserves_console_error_precedence() { + let transport = FakeWebTransport::new(true); + let ctx = FetchContext { + workspace_id: Some("wrk_console".to_string()), + web_timeout: 1, + requires_optional_usage_completeness: true, + ..FetchContext::default() + }; + + let error = OpenCodeGoProvider::fetch_with_transport( + &ctx, + "auth=legacy; __Host-console_session=console", + Arc::clone(&transport), + ) + .await + .unwrap_err(); + + assert!(matches!( + error, + ProviderError::Other(message) if message == "console usage unavailable" + )); + assert_eq!(transport.resolver_calls.load(Ordering::SeqCst), 1); + assert!(transport.usage_sessions.lock().unwrap().is_empty()); + assert!(transport.balance_sessions.lock().unwrap().is_empty()); +} diff --git a/rust/src/providers/opencodego/transport.rs b/rust/src/providers/opencodego/transport.rs new file mode 100644 index 0000000000..54b72d6f12 --- /dev/null +++ b/rust/src/providers/opencodego/transport.rs @@ -0,0 +1,128 @@ +use async_trait::async_trait; +use reqwest::Client; +use std::time::Duration; + +use crate::core::ProviderError; + +use super::{WebCookieSession, console, legacy}; + +#[async_trait] +pub(super) trait WebTransport: Send + Sync + 'static { + type LegacySession: Clone + Send + Sync + 'static; + + async fn fetch_workspace_id( + &self, + cookie_header: &str, + timeout: Duration, + ) -> Result; + + async fn fetch_console_usage( + &self, + workspace_id: &str, + cookie_header: &str, + timeout: Duration, + ) -> Result; + + async fn fetch_console_balance( + &self, + workspace_id: &str, + cookie_header: &str, + timeout: Duration, + ) -> Result, ProviderError>; + + async fn resolve_legacy_session( + &self, + cookies: &WebCookieSession, + ) -> Result; + + async fn fetch_legacy_usage( + &self, + session: &Self::LegacySession, + ) -> Result; + + async fn fetch_legacy_balance( + &self, + session: &Self::LegacySession, + timeout: Duration, + ) -> Result, ProviderError>; +} + +#[derive(Clone)] +pub(super) struct HttpWebTransport { + client: Client, +} + +impl HttpWebTransport { + pub(super) fn new(client: Client) -> Self { + Self { client } + } +} + +#[derive(Clone)] +pub(super) struct HttpLegacySession { + cookie_header: String, + workspace_id: String, +} + +#[async_trait] +impl WebTransport for HttpWebTransport { + type LegacySession = HttpLegacySession; + + async fn fetch_workspace_id( + &self, + cookie_header: &str, + timeout: Duration, + ) -> Result { + console::fetch_workspace_id(&self.client, cookie_header, timeout).await + } + + async fn fetch_console_usage( + &self, + workspace_id: &str, + cookie_header: &str, + timeout: Duration, + ) -> Result { + console::fetch_usage(&self.client, workspace_id, cookie_header, timeout).await + } + + async fn fetch_console_balance( + &self, + workspace_id: &str, + cookie_header: &str, + timeout: Duration, + ) -> Result, ProviderError> { + console::fetch_balance(&self.client, workspace_id, cookie_header, timeout).await + } + + async fn resolve_legacy_session( + &self, + cookies: &WebCookieSession, + ) -> Result { + let workspace_id = legacy::discover_workspace_id(&self.client, cookies.header()).await?; + Ok(HttpLegacySession { + cookie_header: cookies.header.clone(), + workspace_id, + }) + } + + async fn fetch_legacy_usage( + &self, + session: &Self::LegacySession, + ) -> Result { + legacy::fetch_usage(&self.client, &session.cookie_header, &session.workspace_id).await + } + + async fn fetch_legacy_balance( + &self, + session: &Self::LegacySession, + timeout: Duration, + ) -> Result, ProviderError> { + legacy::fetch_balance( + &self.client, + &session.cookie_header, + &session.workspace_id, + timeout, + ) + .await + } +}