From 36686e436c260559ecb6573c7a95e571f107f322 Mon Sep 17 00:00:00 2001 From: NessZerra <90105158+Finesssee@users.noreply.github.com> Date: Tue, 22 Sep 2026 17:59:34 +0700 Subject: [PATCH 1/5] Restore OpenCode Go Console quotas --- rust/src/providers/opencodego/console.rs | 395 +++++++++++++++++++++++ rust/src/providers/opencodego/mod.rs | 153 +++++++-- 2 files changed, 527 insertions(+), 21 deletions(-) create mode 100644 rust/src/providers/opencodego/console.rs diff --git a/rust/src/providers/opencodego/console.rs b/rust/src/providers/opencodego/console.rs new file mode 100644 index 0000000000..fc0de30845 --- /dev/null +++ b/rust/src/providers/opencodego/console.rs @@ -0,0 +1,395 @@ +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 has_legacy_cookie(cookie_header: &str) -> bool { + has_cookie(cookie_header, &["auth", "__Host-auth"]) +} + +pub(super) fn has_console_cookie(cookie_header: &str) -> bool { + has_cookie(cookie_header, &["__Host-console_session"]) +} + +fn has_cookie(cookie_header: &str, names: &[&str]) -> bool { + cookie_header.split(';').any(|part| { + let Some((name, value)) = part.trim().split_once('=') else { + return false; + }; + names.contains(&name.trim()) && !value.trim().is_empty() + }) +} + +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 detects_independent_console_and_legacy_sessions() { + let header = "auth=legacy; __Host-console_session=console"; + assert!(has_legacy_cookie(header)); + assert!(has_console_cookie(header)); + assert!(!has_legacy_cookie("__Host-console_session=console")); + assert!(!has_console_cookie("auth=legacy")); + } + + #[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/mod.rs b/rust/src/providers/opencodego/mod.rs index 0e3f49d38f..259f9cfc57 100644 --- a/rust/src/providers/opencodego/mod.rs +++ b/rust/src/providers/opencodego/mod.rs @@ -5,6 +5,7 @@ //! unless a workspace override scopes the fetch to web first; Web is cookie //! scrape only; Cli is local-only. +mod console; pub(crate) mod local; mod usage_api; @@ -64,13 +65,29 @@ impl OpenCodeGoProvider { } } - fn workspace_id_from_context(workspace_id: Option<&str>) -> Option<&str> { - workspace_id.filter(|id| !id.is_empty()) + fn workspace_id_from_context(workspace_id: Option<&str>) -> Option { + console::normalize_workspace_id(workspace_id) } async fn fetch_workspace_id( client: &Client, cookie_header: &str, + ) -> Result { + let console_result = + console::fetch_workspace_id(client, cookie_header, Duration::from_secs(30)).await; + match console_result { + Ok(workspace_id) => Ok(workspace_id), + Err(console_error) if Self::should_try_legacy(cookie_header, &console_error) => { + let legacy_result = Self::fetch_legacy_workspace_id(client, cookie_header).await; + Self::select_legacy_result(cookie_header, console_error, legacy_result) + } + Err(error) => Err(error), + } + } + + async fn fetch_legacy_workspace_id( + client: &Client, + cookie_header: &str, ) -> Result { let url = format!("{}?id={}", SERVER_URL, WORKSPACES_SERVER_ID); let response = client @@ -119,6 +136,36 @@ impl OpenCodeGoProvider { Self::fetch_page_text(client, &url, cookie_header, None, "usage page").await } + fn should_try_legacy(cookie_header: &str, error: &ProviderError) -> bool { + if !console::has_legacy_cookie(cookie_header) { + return false; + } + matches!( + error, + ProviderError::AuthRequired + | ProviderError::Parse(_) + | ProviderError::Other(_) + | ProviderError::Timeout + ) || error.is_transport_failure() + } + + fn select_legacy_result( + cookie_header: &str, + console_error: ProviderError, + legacy_result: Result, + ) -> Result { + match legacy_result { + Ok(value) => Ok(value), + Err(_legacy_error) + if console::has_console_cookie(cookie_header) + && !matches!(console_error, ProviderError::AuthRequired) => + { + Err(console_error) + } + Err(legacy_error) => Err(legacy_error), + } + } + /// 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( @@ -404,6 +451,11 @@ impl OpenCodeGoProvider { timeout: Duration, ) -> Option { let request_timeout = timeout.min(ZEN_BALANCE_TIMEOUT); + match console::fetch_balance(client, workspace_id, cookie_header, request_timeout).await { + Ok(balance) => return balance, + Err(error) if Self::should_try_legacy(cookie_header, &error) => {} + Err(_) => return None, + } let referer = Self::zen_dashboard_url(workspace_id); let page = Self::fetch_page_text( @@ -483,7 +535,7 @@ impl OpenCodeGoProvider { cookie_header: &str, ) -> Result { let workspace_id = match Self::workspace_id_from_context(ctx.workspace_id.as_deref()) { - Some(workspace_id) => workspace_id.to_string(), + Some(workspace_id) => workspace_id, None => Self::fetch_workspace_id(&self.client, cookie_header).await?, }; // F15 (#2583): start the optional Zen balance fetch in parallel with the @@ -491,23 +543,45 @@ impl OpenCodeGoProvider { // 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) => { - zen_task.abort(); - return Err(err); - } - }; - let usage = match Self::parse_usage_text(&page) { - Ok(usage) => usage, - Err(err) => { - zen_task.abort(); - return Err(err); - } - }; - // 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 console_result = console::fetch_usage( + &self.client, + &workspace_id, + cookie_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(ProviderError::Parse( + "No OpenCode Go subscription is available".to_string(), + )); + } + Err(console_error) if Self::should_try_legacy(cookie_header, &console_error) => { + let legacy_result = + match Self::fetch_usage_page(&self.client, &workspace_id, cookie_header) + .await + { + Ok(page) => Self::parse_usage_text(&page) + .map(|usage| (usage, Self::parse_zen_balance(&page))), + Err(error) => Err(error), + }; + match Self::select_legacy_result(cookie_header, console_error, legacy_result) { + Ok(result) => result, + Err(error) => { + zen_task.abort(); + return Err(error); + } + } + } + Err(error) => { + zen_task.abort(); + return Err(error); + } + }; + let balance = match embedded_balance { Some(balance) => { zen_task.abort(); Some(balance) @@ -808,7 +882,7 @@ mod tests { fn uses_context_workspace_id_before_discovery() { assert_eq!( OpenCodeGoProvider::workspace_id_from_context(Some("wrk_override")), - Some("wrk_override") + Some("wrk_override".to_string()) ); assert_eq!( OpenCodeGoProvider::workspace_id_from_context(Some("")), @@ -846,6 +920,43 @@ mod tests { )); } + #[test] + fn console_recovery_requires_an_independent_legacy_cookie() { + assert!(OpenCodeGoProvider::should_try_legacy( + "auth=legacy; __Host-console_session=console", + &ProviderError::AuthRequired + )); + assert!(!OpenCodeGoProvider::should_try_legacy( + "__Host-console_session=console", + &ProviderError::AuthRequired + )); + assert!(!OpenCodeGoProvider::should_try_legacy( + "auth=legacy", + &ProviderError::NotInstalled("terminal".to_string()) + )); + } + + #[test] + fn failed_legacy_read_preserves_non_auth_console_failure() { + let error = OpenCodeGoProvider::select_legacy_result::<()>( + "auth=legacy; __Host-console_session=console", + ProviderError::Other("console unavailable".to_string()), + Err(ProviderError::AuthRequired), + ) + .unwrap_err(); + assert!(matches!(error, ProviderError::Other(message) if message == "console unavailable")); + + let error = OpenCodeGoProvider::select_legacy_result::<()>( + "auth=legacy; __Host-console_session=console", + 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 parses_usage_blocks() { let text = r#" From b90120b17e1e4be16c256929aac806647fa51a1f Mon Sep 17 00:00:00 2001 From: NessZerra <90105158+Finesssee@users.noreply.github.com> Date: Tue, 22 Sep 2026 20:35:10 +0700 Subject: [PATCH 2/5] Isolate OpenCode legacy fallback --- rust/src/providers/opencodego/legacy.rs | 559 ++++++++++++++++++ rust/src/providers/opencodego/mod.rs | 743 ++++-------------------- 2 files changed, 663 insertions(+), 639 deletions(-) create mode 100644 rust/src/providers/opencodego/legacy.rs diff --git a/rust/src/providers/opencodego/legacy.rs b/rust/src/providers/opencodego/legacy.rs new file mode 100644 index 0000000000..9b7f9938bf --- /dev/null +++ b/rust/src/providers/opencodego/legacy.rs @@ -0,0 +1,559 @@ +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"; + +/// Authenticated legacy transport. Workspace discovery stays inside this +/// session so Console workspace IDs never cross into the independent legacy +/// cookie session. +pub(super) struct LegacySession<'a> { + client: &'a Client, + cookie_header: &'a str, +} + +/// Parsed legacy usage response returned to the provider orchestrator. +pub(super) struct LegacyUsage { + pub(super) usage: UsageSnapshot, + pub(super) embedded_balance: Option, +} + +/// Whether a failed Console request is eligible for the legacy route. +pub(super) fn can_recover(cookie_header: &str, error: &ProviderError) -> bool { + if !super::console::has_legacy_cookie(cookie_header) { + return false; + } + matches!( + error, + ProviderError::AuthRequired + | ProviderError::Parse(_) + | ProviderError::Other(_) + | ProviderError::Timeout + ) || error.is_transport_failure() +} + +/// Resolve competing route failures without hiding a useful Console error. +pub(super) fn select_result( + cookie_header: &str, + console_error: ProviderError, + legacy_result: Result, +) -> Result { + match legacy_result { + Ok(value) => Ok(value), + Err(_legacy_error) + if super::console::has_console_cookie(cookie_header) + && !matches!(console_error, ProviderError::AuthRequired) => + { + Err(console_error) + } + Err(legacy_error) => Err(legacy_error), + } +} + +impl<'a> LegacySession<'a> { + pub(super) fn new(client: &'a Client, cookie_header: &'a str) -> Self { + Self { + client, + cookie_header, + } + } + + /// Discover the workspace owned by this legacy cookie session. + pub(super) async fn discover_workspace_id(&self) -> Result { + let text = self + .fetch_server_text(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())) + } + + /// Execute the complete legacy usage route, including route-local + /// workspace discovery. + pub(super) async fn fetch_usage(&self) -> Result { + let workspace_id = self.discover_workspace_id().await?; + let url = format!("{BASE_URL}/workspace/{workspace_id}/go"); + let page = self.fetch_page_text(&url, None, "usage page").await?; + Ok(LegacyUsage { + usage: parse_usage_text(&page)?, + embedded_balance: parse_zen_balance(&page), + }) + } + + /// Execute the complete legacy balance route with a workspace discovered + /// from the legacy cookie rather than the Console session. + pub(super) async fn fetch_balance( + &self, + timeout: Duration, + ) -> Result, ProviderError> { + let workspace_id = self.discover_workspace_id().await?; + let referer = format!("{BASE_URL}/workspace/{workspace_id}"); + let page = self + .fetch_page_text(&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 = self + .fetch_server_text( + BILLING_SERVER_ID, + Some(&args), + &referer, + Some(timeout), + "billing API", + ) + .await?; + Ok(parse_billing_server_balance(&billing)) + } + + async fn fetch_page_text( + &self, + url: &str, + timeout: Option, + what: &str, + ) -> Result { + let mut request = self + .client + .get(url) + .header("Cookie", self.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( + &self, + 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 = self + .client + .get(url) + .header("Cookie", self.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 console_recovery_requires_an_independent_legacy_cookie() { + assert!(can_recover( + "auth=legacy; __Host-console_session=console", + &ProviderError::AuthRequired + )); + assert!(!can_recover( + "__Host-console_session=console", + &ProviderError::AuthRequired + )); + assert!(!can_recover( + "auth=legacy", + &ProviderError::NotInstalled("terminal".to_string()) + )); + } + + #[test] + fn failed_legacy_read_preserves_non_auth_console_failure() { + let error = select_result::<()>( + "auth=legacy; __Host-console_session=console", + ProviderError::Other("console unavailable".to_string()), + Err(ProviderError::AuthRequired), + ) + .unwrap_err(); + assert!(matches!(error, ProviderError::Other(message) if message == "console unavailable")); + + let error = select_result::<()>( + "auth=legacy; __Host-console_session=console", + 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 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 259f9cfc57..d387713d86 100644 --- a/rust/src/providers/opencodego/mod.rs +++ b/rust/src/providers/opencodego/mod.rs @@ -6,6 +6,7 @@ //! scrape only; Cli is local-only. mod console; +mod legacy; pub(crate) mod local; mod usage_api; @@ -13,7 +14,6 @@ use async_trait::async_trait; use chrono::Utc; use reqwest::Client; use std::time::Duration; -use uuid::Uuid; use crate::core::{ CostSnapshot, FetchContext, Provider, ProviderError, ProviderFetchResult, ProviderId, @@ -21,12 +21,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. @@ -73,371 +69,7 @@ impl OpenCodeGoProvider { client: &Client, cookie_header: &str, ) -> Result { - let console_result = - console::fetch_workspace_id(client, cookie_header, Duration::from_secs(30)).await; - match console_result { - Ok(workspace_id) => Ok(workspace_id), - Err(console_error) if Self::should_try_legacy(cookie_header, &console_error) => { - let legacy_result = Self::fetch_legacy_workspace_id(client, cookie_header).await; - Self::select_legacy_result(cookie_header, console_error, legacy_result) - } - Err(error) => Err(error), - } - } - - async fn fetch_legacy_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 - } - - fn should_try_legacy(cookie_header: &str, error: &ProviderError) -> bool { - if !console::has_legacy_cookie(cookie_header) { - return false; - } - matches!( - error, - ProviderError::AuthRequired - | ProviderError::Parse(_) - | ProviderError::Other(_) - | ProviderError::Timeout - ) || error.is_transport_failure() - } - - fn select_legacy_result( - cookie_header: &str, - console_error: ProviderError, - legacy_result: Result, - ) -> Result { - match legacy_result { - Ok(value) => Ok(value), - Err(_legacy_error) - if console::has_console_cookie(cookie_header) - && !matches!(console_error, ProviderError::AuthRequired) => - { - Err(console_error) - } - Err(legacy_error) => Err(legacy_error), - } - } - - /// 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) + console::fetch_workspace_id(client, cookie_header, Duration::from_secs(30)).await } /// Upstream `fetchZenBalance`: the dashboard HTML embeds the balance for @@ -453,36 +85,13 @@ impl OpenCodeGoProvider { let request_timeout = timeout.min(ZEN_BALANCE_TIMEOUT); match console::fetch_balance(client, workspace_id, cookie_header, request_timeout).await { Ok(balance) => return balance, - Err(error) if Self::should_try_legacy(cookie_header, &error) => {} + Err(error) if legacy::can_recover(cookie_header, &error) => {} Err(_) => return None, } - 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); - } - - 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) + legacy::LegacySession::new(client, cookie_header) + .fetch_balance(request_timeout) + .await + .unwrap_or(None) } /// Spawn the optional Zen balance task (25 ms start delay so the usage @@ -503,15 +112,41 @@ 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) - .await - .ok()?, + None => match Self::fetch_workspace_id(&client, &cookie_header).await { + Ok(id) => id, + Err(error) if legacy::can_recover(&cookie_header, &error) => { + return legacy::LegacySession::new(&client, &cookie_header) + .fetch_balance(timeout.min(ZEN_BALANCE_TIMEOUT)) + .await + .unwrap_or(None); + } + Err(_) => return None, + }, }; Self::fetch_zen_balance(&client, &workspace_id, &cookie_header, timeout).await }); (task, started_at) } + fn spawn_legacy_balance_task( + &self, + cookie_header: &str, + web_timeout: u64, + ) -> (tokio::task::JoinHandle>, std::time::Instant) { + let client = self.client.clone(); + let cookie_header = cookie_header.to_string(); + 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; + legacy::LegacySession::new(&client, &cookie_header) + .fetch_balance(timeout) + .await + .unwrap_or(None) + }); + (task, started_at) + } + /// 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( @@ -536,7 +171,15 @@ impl OpenCodeGoProvider { ) -> Result { let workspace_id = match Self::workspace_id_from_context(ctx.workspace_id.as_deref()) { Some(workspace_id) => workspace_id, - None => Self::fetch_workspace_id(&self.client, cookie_header).await?, + None => match Self::fetch_workspace_id(&self.client, cookie_header).await { + Ok(workspace_id) => workspace_id, + Err(console_error) if legacy::can_recover(cookie_header, &console_error) => { + return self + .fetch_legacy_with_cookies(ctx, cookie_header, console_error) + .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 @@ -550,37 +193,32 @@ impl OpenCodeGoProvider { 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(ProviderError::Parse( - "No OpenCode Go subscription is available".to_string(), - )); - } - Err(console_error) if Self::should_try_legacy(cookie_header, &console_error) => { - let legacy_result = - match Self::fetch_usage_page(&self.client, &workspace_id, cookie_header) - .await - { - Ok(page) => Self::parse_usage_text(&page) - .map(|usage| (usage, Self::parse_zen_balance(&page))), - Err(error) => Err(error), - }; - match Self::select_legacy_result(cookie_header, console_error, legacy_result) { - Ok(result) => result, - Err(error) => { - zen_task.abort(); - return Err(error); - } + let (usage, embedded_balance) = match console_result { + Ok(console::ConsoleUsage::Snapshot(usage)) => (*usage, None), + Ok(console::ConsoleUsage::NoSubscription) => { + zen_task.abort(); + return Err(ProviderError::Parse( + "No OpenCode Go subscription is available".to_string(), + )); + } + Err(console_error) if legacy::can_recover(cookie_header, &console_error) => { + let legacy_result = legacy::LegacySession::new(&self.client, cookie_header) + .fetch_usage() + .await + .map(|result| (result.usage, result.embedded_balance)); + match legacy::select_result(cookie_header, console_error, legacy_result) { + Ok(result) => result, + Err(error) => { + zen_task.abort(); + return Err(error); } } - Err(error) => { - zen_task.abort(); - return Err(error); - } - }; + } + Err(error) => { + zen_task.abort(); + return Err(error); + } + }; let balance = match embedded_balance { Some(balance) => { zen_task.abort(); @@ -598,6 +236,41 @@ impl OpenCodeGoProvider { Ok(Self::with_zen_balance(usage, "web", balance)) } + async fn fetch_legacy_with_cookies( + &self, + ctx: &FetchContext, + cookie_header: &str, + console_error: ProviderError, + ) -> Result { + let (zen_task, zen_started) = + self.spawn_legacy_balance_task(cookie_header, ctx.web_timeout); + let legacy_result = legacy::LegacySession::new(&self.client, cookie_header) + .fetch_usage() + .await; + let result = match legacy::select_result(cookie_header, console_error, legacy_result) { + 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 => { + Self::join_zen_balance( + zen_task, + zen_started, + ctx.requires_optional_usage_completeness, + ) + .await + } + }; + 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( @@ -803,81 +476,10 @@ 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!( @@ -890,17 +492,6 @@ mod tests { ); } - #[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 { @@ -920,90 +511,6 @@ mod tests { )); } - #[test] - fn console_recovery_requires_an_independent_legacy_cookie() { - assert!(OpenCodeGoProvider::should_try_legacy( - "auth=legacy; __Host-console_session=console", - &ProviderError::AuthRequired - )); - assert!(!OpenCodeGoProvider::should_try_legacy( - "__Host-console_session=console", - &ProviderError::AuthRequired - )); - assert!(!OpenCodeGoProvider::should_try_legacy( - "auth=legacy", - &ProviderError::NotInstalled("terminal".to_string()) - )); - } - - #[test] - fn failed_legacy_read_preserves_non_auth_console_failure() { - let error = OpenCodeGoProvider::select_legacy_result::<()>( - "auth=legacy; __Host-console_session=console", - ProviderError::Other("console unavailable".to_string()), - Err(ProviderError::AuthRequired), - ) - .unwrap_err(); - assert!(matches!(error, ProviderError::Other(message) if message == "console unavailable")); - - let error = OpenCodeGoProvider::select_legacy_result::<()>( - "auth=legacy; __Host-console_session=console", - 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 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] @@ -1048,46 +555,4 @@ mod tests { 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); - } } From 7cafc99194996bc274d96e154e90c999f647e7b2 Mon Sep 17 00:00:00 2001 From: NessZerra <90105158+Finesssee@users.noreply.github.com> Date: Tue, 22 Sep 2026 23:03:45 +0700 Subject: [PATCH 3/5] Refactor OpenCode legacy session --- rust/src/providers/opencodego/console.rs | 26 -- rust/src/providers/opencodego/legacy.rs | 337 +++++++++-------------- rust/src/providers/opencodego/mod.rs | 238 ++++++++++++---- 3 files changed, 326 insertions(+), 275 deletions(-) diff --git a/rust/src/providers/opencodego/console.rs b/rust/src/providers/opencodego/console.rs index fc0de30845..e37d972bd5 100644 --- a/rust/src/providers/opencodego/console.rs +++ b/rust/src/providers/opencodego/console.rs @@ -18,23 +18,6 @@ pub(super) enum ConsoleUsage { NoSubscription, } -pub(super) fn has_legacy_cookie(cookie_header: &str) -> bool { - has_cookie(cookie_header, &["auth", "__Host-auth"]) -} - -pub(super) fn has_console_cookie(cookie_header: &str) -> bool { - has_cookie(cookie_header, &["__Host-console_session"]) -} - -fn has_cookie(cookie_header: &str, names: &[&str]) -> bool { - cookie_header.split(';').any(|part| { - let Some((name, value)) = part.trim().split_once('=') else { - return false; - }; - names.contains(&name.trim()) && !value.trim().is_empty() - }) -} - pub(super) fn normalize_workspace_id(raw: Option<&str>) -> Option { let raw = raw?.trim(); if is_workspace_id(raw) { @@ -333,15 +316,6 @@ mod tests { ); } - #[test] - fn detects_independent_console_and_legacy_sessions() { - let header = "auth=legacy; __Host-console_session=console"; - assert!(has_legacy_cookie(header)); - assert!(has_console_cookie(header)); - assert!(!has_legacy_cookie("__Host-console_session=console")); - assert!(!has_console_cookie("auth=legacy")); - } - #[test] fn parses_console_microcent_usage_and_nullable_resets() { let ConsoleUsage::Snapshot(snapshot) = parse_usage(&usage_json("null"), now()).unwrap() diff --git a/rust/src/providers/opencodego/legacy.rs b/rust/src/providers/opencodego/legacy.rs index 9b7f9938bf..c876c1b8c5 100644 --- a/rust/src/providers/opencodego/legacy.rs +++ b/rust/src/providers/opencodego/legacy.rs @@ -12,194 +12,162 @@ const WORKSPACES_SERVER_ID: &str = "def39973159c7f0483d8793a822b8dbb10d067e12c65455fcb4608459ba0234f"; const BILLING_SERVER_ID: &str = "c83b78a614689c38ebee981f9b39a8b377716db85c1fd7dbab604adc02d3313d"; -/// Authenticated legacy transport. Workspace discovery stays inside this -/// session so Console workspace IDs never cross into the independent legacy -/// cookie session. -pub(super) struct LegacySession<'a> { - client: &'a Client, - cookie_header: &'a str, -} - /// Parsed legacy usage response returned to the provider orchestrator. pub(super) struct LegacyUsage { pub(super) usage: UsageSnapshot, pub(super) embedded_balance: Option, } -/// Whether a failed Console request is eligible for the legacy route. -pub(super) fn can_recover(cookie_header: &str, error: &ProviderError) -> bool { - if !super::console::has_legacy_cookie(cookie_header) { - return false; - } - matches!( - error, - ProviderError::AuthRequired - | ProviderError::Parse(_) - | ProviderError::Other(_) - | ProviderError::Timeout - ) || error.is_transport_failure() +/// 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())) } -/// Resolve competing route failures without hiding a useful Console error. -pub(super) fn select_result( +pub(super) async fn fetch_usage( + client: &Client, cookie_header: &str, - console_error: ProviderError, - legacy_result: Result, -) -> Result { - match legacy_result { - Ok(value) => Ok(value), - Err(_legacy_error) - if super::console::has_console_cookie(cookie_header) - && !matches!(console_error, ProviderError::AuthRequired) => - { - Err(console_error) - } - Err(legacy_error) => Err(legacy_error), - } + 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), + }) } -impl<'a> LegacySession<'a> { - pub(super) fn new(client: &'a Client, cookie_header: &'a str) -> Self { - Self { - client, - cookie_header, - } - } - - /// Discover the workspace owned by this legacy cookie session. - pub(super) async fn discover_workspace_id(&self) -> Result { - let text = self - .fetch_server_text(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())) - } - - /// Execute the complete legacy usage route, including route-local - /// workspace discovery. - pub(super) async fn fetch_usage(&self) -> Result { - let workspace_id = self.discover_workspace_id().await?; - let url = format!("{BASE_URL}/workspace/{workspace_id}/go"); - let page = self.fetch_page_text(&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)); } - /// Execute the complete legacy balance route with a workspace discovered - /// from the legacy cookie rather than the Console session. - pub(super) async fn fetch_balance( - &self, - timeout: Duration, - ) -> Result, ProviderError> { - let workspace_id = self.discover_workspace_id().await?; - let referer = format!("{BASE_URL}/workspace/{workspace_id}"); - let page = self - .fetch_page_text(&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)) +} - let args = serde_json::json!([workspace_id]).to_string(); - let billing = self - .fetch_server_text( - 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); } - - async fn fetch_page_text( - &self, - url: &str, - timeout: Option, - what: &str, - ) -> Result { - let mut request = self - .client - .get(url) - .header("Cookie", self.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) { + 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); } - Ok(text) + 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( - &self, - 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 = self - .client - .get(url) - .header("Cookie", self.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) { +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); } - Ok(text) + 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 { @@ -426,43 +394,6 @@ fn billing_numeric_value(value: &serde_json::Value) -> Option { mod tests { use super::*; - #[test] - fn console_recovery_requires_an_independent_legacy_cookie() { - assert!(can_recover( - "auth=legacy; __Host-console_session=console", - &ProviderError::AuthRequired - )); - assert!(!can_recover( - "__Host-console_session=console", - &ProviderError::AuthRequired - )); - assert!(!can_recover( - "auth=legacy", - &ProviderError::NotInstalled("terminal".to_string()) - )); - } - - #[test] - fn failed_legacy_read_preserves_non_auth_console_failure() { - let error = select_result::<()>( - "auth=legacy; __Host-console_session=console", - ProviderError::Other("console unavailable".to_string()), - Err(ProviderError::AuthRequired), - ) - .unwrap_err(); - assert!(matches!(error, ProviderError::Other(message) if message == "console unavailable")); - - let error = select_result::<()>( - "auth=legacy; __Host-console_session=console", - 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 parses_workspace_ids_without_duplicates() { let text = r#"{ id: "wrk_abc123", name: "x" } { id: "wrk_def456" } { id: "wrk_abc123" }"#; diff --git a/rust/src/providers/opencodego/mod.rs b/rust/src/providers/opencodego/mod.rs index d387713d86..0549348391 100644 --- a/rust/src/providers/opencodego/mod.rs +++ b/rust/src/providers/opencodego/mod.rs @@ -38,6 +38,102 @@ 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, +} + +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), + } + } +} + +#[derive(Clone)] +struct LegacyWorkspaceSession { + client: Client, + cookie_header: String, + workspace_id: String, +} + +impl LegacyWorkspaceSession { + async fn resolve(client: &Client, cookies: &WebCookieSession) -> Result { + let workspace_id = legacy::discover_workspace_id(client, cookies.header()).await?; + Ok(Self { + client: client.clone(), + cookie_header: cookies.header.clone(), + workspace_id, + }) + } + + async fn fetch_usage(&self) -> Result { + legacy::fetch_usage(&self.client, &self.cookie_header, &self.workspace_id).await + } + + async fn fetch_balance(&self, timeout: Duration) -> Result, ProviderError> { + legacy::fetch_balance( + &self.client, + &self.cookie_header, + &self.workspace_id, + timeout, + ) + .await + } +} + impl OpenCodeGoProvider { pub fn new() -> Self { Self { @@ -79,16 +175,19 @@ impl OpenCodeGoProvider { async fn fetch_zen_balance( client: &Client, workspace_id: &str, - cookie_header: &str, + cookies: &WebCookieSession, timeout: Duration, ) -> Option { let request_timeout = timeout.min(ZEN_BALANCE_TIMEOUT); - match console::fetch_balance(client, workspace_id, cookie_header, request_timeout).await { + match console::fetch_balance(client, workspace_id, cookies.header(), request_timeout).await + { Ok(balance) => return balance, - Err(error) if legacy::can_recover(cookie_header, &error) => {} + Err(error) if cookies.can_recover_with_legacy(&error) => {} Err(_) => return None, } - legacy::LegacySession::new(client, cookie_header) + LegacyWorkspaceSession::resolve(client, cookies) + .await + .ok()? .fetch_balance(request_timeout) .await .unwrap_or(None) @@ -99,12 +198,12 @@ impl OpenCodeGoProvider { /// inside the task when no override is pinned. fn spawn_zen_balance_task( &self, - cookie_header: &str, + 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(); + 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(); @@ -112,10 +211,12 @@ impl OpenCodeGoProvider { tokio::time::sleep(ZEN_BALANCE_START_DELAY).await; let workspace_id = match workspace_id_override { Some(id) => id, - None => match Self::fetch_workspace_id(&client, &cookie_header).await { + None => match Self::fetch_workspace_id(&client, cookies.header()).await { Ok(id) => id, - Err(error) if legacy::can_recover(&cookie_header, &error) => { - return legacy::LegacySession::new(&client, &cookie_header) + Err(error) if cookies.can_recover_with_legacy(&error) => { + return LegacyWorkspaceSession::resolve(&client, &cookies) + .await + .ok()? .fetch_balance(timeout.min(ZEN_BALANCE_TIMEOUT)) .await .unwrap_or(None); @@ -123,26 +224,20 @@ impl OpenCodeGoProvider { Err(_) => return None, }, }; - Self::fetch_zen_balance(&client, &workspace_id, &cookie_header, timeout).await + Self::fetch_zen_balance(&client, &workspace_id, &cookies, timeout).await }); (task, started_at) } fn spawn_legacy_balance_task( - &self, - cookie_header: &str, + session: LegacyWorkspaceSession, web_timeout: u64, ) -> (tokio::task::JoinHandle>, std::time::Instant) { - let client = self.client.clone(); - let cookie_header = cookie_header.to_string(); 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; - legacy::LegacySession::new(&client, &cookie_header) - .fetch_balance(timeout) - .await - .unwrap_or(None) + session.fetch_balance(timeout).await.unwrap_or(None) }); (task, started_at) } @@ -169,13 +264,14 @@ impl OpenCodeGoProvider { ctx: &FetchContext, cookie_header: &str, ) -> 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, - None => match Self::fetch_workspace_id(&self.client, cookie_header).await { + None => match Self::fetch_workspace_id(&self.client, cookies.header()).await { Ok(workspace_id) => workspace_id, - Err(console_error) if legacy::can_recover(cookie_header, &console_error) => { + Err(console_error) if cookies.can_recover_with_legacy(&console_error) => { return self - .fetch_legacy_with_cookies(ctx, cookie_header, console_error) + .fetch_legacy_with_cookies(ctx, &cookies, console_error) .await; } Err(error) => return Err(error), @@ -185,11 +281,11 @@ impl OpenCodeGoProvider { // 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); + self.spawn_zen_balance_task(&cookies, Some(&workspace_id), ctx.web_timeout); let console_result = console::fetch_usage( &self.client, &workspace_id, - cookie_header, + cookies.header(), Duration::from_secs(ctx.web_timeout.max(1)), ) .await; @@ -201,18 +297,11 @@ impl OpenCodeGoProvider { "No OpenCode Go subscription is available".to_string(), )); } - Err(console_error) if legacy::can_recover(cookie_header, &console_error) => { - let legacy_result = legacy::LegacySession::new(&self.client, cookie_header) - .fetch_usage() - .await - .map(|result| (result.usage, result.embedded_balance)); - match legacy::select_result(cookie_header, console_error, legacy_result) { - Ok(result) => result, - Err(error) => { - zen_task.abort(); - return Err(error); - } - } + Err(console_error) if cookies.can_recover_with_legacy(&console_error) => { + zen_task.abort(); + return self + .fetch_legacy_with_cookies(ctx, &cookies, console_error) + .await; } Err(error) => { zen_task.abort(); @@ -239,15 +328,17 @@ impl OpenCodeGoProvider { async fn fetch_legacy_with_cookies( &self, ctx: &FetchContext, - cookie_header: &str, + cookies: &WebCookieSession, console_error: ProviderError, ) -> Result { + let session = match LegacyWorkspaceSession::resolve(&self.client, cookies).await { + Ok(session) => session, + Err(error) => return cookies.select_legacy_result(console_error, Err(error)), + }; let (zen_task, zen_started) = - self.spawn_legacy_balance_task(cookie_header, ctx.web_timeout); - let legacy_result = legacy::LegacySession::new(&self.client, cookie_header) - .fetch_usage() - .await; - let result = match legacy::select_result(cookie_header, console_error, legacy_result) { + Self::spawn_legacy_balance_task(session.clone(), ctx.web_timeout); + let result = match cookies.select_legacy_result(console_error, session.fetch_usage().await) + { Ok(result) => result, Err(error) => { zen_task.abort(); @@ -433,11 +524,9 @@ impl OpenCodeGoProvider { let Some(cookie_header) = cookie_header else { return Ok(result); }; - let (task, started) = self.spawn_zen_balance_task( - &cookie_header, - ctx.workspace_id.as_deref(), - ctx.web_timeout, - ); + let cookies = WebCookieSession::new(&cookie_header); + let (task, started) = + self.spawn_zen_balance_task(&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 { @@ -480,6 +569,63 @@ fn zen_balance_join_budget( mod tests { use super::*; + #[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!( From 2cbef3a6e167124a2e6aa58f91e36ceab55d0e8b Mon Sep 17 00:00:00 2001 From: NessZerra <90105158+Finesssee@users.noreply.github.com> Date: Tue, 22 Sep 2026 23:43:59 +0700 Subject: [PATCH 4/5] Fix OpenCode legacy fallback race --- rust/src/providers/opencodego/mod.rs | 179 ++++++++++++++++++++++----- 1 file changed, 147 insertions(+), 32 deletions(-) diff --git a/rust/src/providers/opencodego/mod.rs b/rust/src/providers/opencodego/mod.rs index 0549348391..8c860c8f00 100644 --- a/rust/src/providers/opencodego/mod.rs +++ b/rust/src/providers/opencodego/mod.rs @@ -50,6 +50,12 @@ struct WebCookieSession { capabilities: CookieCapabilities, } +#[derive(Clone, Copy, Debug, PartialEq)] +enum OptionalZenBalance { + Resolved(Option), + LegacyBalanceRequired, +} + impl WebCookieSession { fn new(header: &str) -> Self { let has_cookie = |names: &[&str]| { @@ -177,20 +183,16 @@ impl OpenCodeGoProvider { workspace_id: &str, cookies: &WebCookieSession, timeout: Duration, - ) -> Option { + ) -> OptionalZenBalance { let request_timeout = timeout.min(ZEN_BALANCE_TIMEOUT); match console::fetch_balance(client, workspace_id, cookies.header(), request_timeout).await { - Ok(balance) => return balance, - Err(error) if cookies.can_recover_with_legacy(&error) => {} - Err(_) => return None, + Ok(balance) => OptionalZenBalance::Resolved(balance), + Err(error) if cookies.can_recover_with_legacy(&error) => { + OptionalZenBalance::LegacyBalanceRequired + } + Err(_) => OptionalZenBalance::Resolved(None), } - LegacyWorkspaceSession::resolve(client, cookies) - .await - .ok()? - .fetch_balance(request_timeout) - .await - .unwrap_or(None) } /// Spawn the optional Zen balance task (25 ms start delay so the usage @@ -201,7 +203,10 @@ impl OpenCodeGoProvider { cookies: &WebCookieSession, workspace_id_override: Option<&str>, web_timeout: u64, - ) -> (tokio::task::JoinHandle>, std::time::Instant) { + ) -> ( + tokio::task::JoinHandle, + std::time::Instant, + ) { let client = self.client.clone(); let cookies = cookies.clone(); let workspace_id_override = workspace_id_override.map(str::to_string); @@ -214,14 +219,9 @@ impl OpenCodeGoProvider { None => match Self::fetch_workspace_id(&client, cookies.header()).await { Ok(id) => id, Err(error) if cookies.can_recover_with_legacy(&error) => { - return LegacyWorkspaceSession::resolve(&client, &cookies) - .await - .ok()? - .fetch_balance(timeout.min(ZEN_BALANCE_TIMEOUT)) - .await - .unwrap_or(None); + return OptionalZenBalance::LegacyBalanceRequired; } - Err(_) => return None, + Err(_) => return OptionalZenBalance::Resolved(None), }, }; Self::fetch_zen_balance(&client, &workspace_id, &cookies, timeout).await @@ -232,26 +232,37 @@ impl OpenCodeGoProvider { fn spawn_legacy_balance_task( session: LegacyWorkspaceSession, web_timeout: u64, - ) -> (tokio::task::JoinHandle>, std::time::Instant) { + ) -> ( + 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; - session.fetch_balance(timeout).await.unwrap_or(None) + OptionalZenBalance::Resolved(session.fetch_balance(timeout).await.unwrap_or(None)) }); (task, started_at) } + async fn abort_optional_balance_and_resolve( + task: tokio::task::JoinHandle, + resolution: impl std::future::Future>, + ) -> Result { + task.abort(); + resolution.await + } + /// 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 @@ -259,6 +270,52 @@ impl OpenCodeGoProvider { } } + async fn resolve_legacy_balance( + client: &Client, + 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 = LegacyWorkspaceSession::resolve(client, cookies) + .await + .ok()?; + session + .fetch_balance(budget.min(ZEN_BALANCE_TIMEOUT)) + .await + .unwrap_or(None) + }) + .await + .ok() + .flatten() + } + + async fn finish_zen_balance( + client: &Client, + 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( + client, + cookies, + started_at, + requires_optional_usage_completeness, + ) + .await + } + } + } + async fn fetch_with_cookies( &self, ctx: &FetchContext, @@ -298,9 +355,19 @@ impl OpenCodeGoProvider { )); } Err(console_error) if cookies.can_recover_with_legacy(&console_error) => { - zen_task.abort(); + let session = match Self::abort_optional_balance_and_resolve( + zen_task, + LegacyWorkspaceSession::resolve(&self.client, &cookies), + ) + .await + { + Ok(session) => session, + Err(error) => { + return cookies.select_legacy_result(console_error, Err(error)); + } + }; return self - .fetch_legacy_with_cookies(ctx, &cookies, console_error) + .fetch_legacy_with_session(ctx, &cookies, console_error, session) .await; } Err(error) => { @@ -314,7 +381,9 @@ impl OpenCodeGoProvider { Some(balance) } None => { - Self::join_zen_balance( + Self::finish_zen_balance( + &self.client, + &cookies, zen_task, zen_started, ctx.requires_optional_usage_completeness, @@ -335,6 +404,17 @@ impl OpenCodeGoProvider { Ok(session) => session, Err(error) => return cookies.select_legacy_result(console_error, Err(error)), }; + self.fetch_legacy_with_session(ctx, cookies, console_error, session) + .await + } + + async fn fetch_legacy_with_session( + &self, + ctx: &FetchContext, + cookies: &WebCookieSession, + console_error: ProviderError, + session: LegacyWorkspaceSession, + ) -> Result { let (zen_task, zen_started) = Self::spawn_legacy_balance_task(session.clone(), ctx.web_timeout); let result = match cookies.select_legacy_result(console_error, session.fetch_usage().await) @@ -351,12 +431,16 @@ impl OpenCodeGoProvider { Some(balance) } None => { - Self::join_zen_balance( + 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)) @@ -527,8 +611,14 @@ impl OpenCodeGoProvider { let cookies = WebCookieSession::new(&cookie_header); let (task, started) = self.spawn_zen_balance_task(&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( + &self.client, + &cookies, + task, + started, + ctx.requires_optional_usage_completeness, + ) + .await { result = Self::with_zen_balance(result.usage, &result.source_label.clone(), Some(balance)); @@ -568,6 +658,10 @@ fn zen_balance_join_budget( #[cfg(test)] mod tests { use super::*; + use std::sync::{ + Arc, + atomic::{AtomicUsize, Ordering}, + }; #[test] fn cookie_session_classifies_transport_capabilities_once() { @@ -684,7 +778,7 @@ mod tests { let started = std::time::Instant::now(); let task = tokio::spawn(async { tokio::time::sleep(Duration::from_secs(30)).await; - Some(42.5) + 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; @@ -696,9 +790,30 @@ mod tests { let started = std::time::Instant::now(); let task = tokio::spawn(async { tokio::time::sleep(Duration::from_millis(10)).await; - Some(42.5) + OptionalZenBalance::Resolved(Some(42.5)) }); let balance = OpenCodeGoProvider::join_zen_balance(task, started, true).await; - assert_eq!(balance, Some(42.5)); + assert_eq!(balance, Some(OptionalZenBalance::Resolved(Some(42.5)))); + } + + #[tokio::test] + async fn simultaneous_console_failures_resolve_one_legacy_session() { + let resolver_calls = Arc::new(AtomicUsize::new(0)); + let balance_task = tokio::spawn(async { + tokio::task::yield_now().await; + OptionalZenBalance::LegacyBalanceRequired + }); + let calls = Arc::clone(&resolver_calls); + + let resolved = + OpenCodeGoProvider::abort_optional_balance_and_resolve(balance_task, async move { + calls.fetch_add(1, Ordering::SeqCst); + Ok::<_, ProviderError>("shared legacy session") + }) + .await + .unwrap(); + + assert_eq!(resolved, "shared legacy session"); + assert_eq!(resolver_calls.load(Ordering::SeqCst), 1); } } From 838f4cac84a85ac25fe8948cad05f804db27649c Mon Sep 17 00:00:00 2001 From: NessZerra <90105158+Finesssee@users.noreply.github.com> Date: Wed, 23 Sep 2026 00:04:50 +0700 Subject: [PATCH 5/5] Test OpenCode fallback orchestration --- rust/src/providers/opencodego/mod.rs | 366 ++++++--------------- rust/src/providers/opencodego/tests.rs | 307 +++++++++++++++++ rust/src/providers/opencodego/transport.rs | 128 +++++++ 3 files changed, 535 insertions(+), 266 deletions(-) create mode 100644 rust/src/providers/opencodego/tests.rs create mode 100644 rust/src/providers/opencodego/transport.rs diff --git a/rust/src/providers/opencodego/mod.rs b/rust/src/providers/opencodego/mod.rs index 8c860c8f00..d676e3dd94 100644 --- a/rust/src/providers/opencodego/mod.rs +++ b/rust/src/providers/opencodego/mod.rs @@ -8,13 +8,17 @@ 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 transport::{HttpWebTransport, WebTransport}; + use crate::core::{ CostSnapshot, FetchContext, Provider, ProviderError, ProviderFetchResult, ProviderId, ProviderMetadata, RateWindow, SourceMode, UsageSnapshot, @@ -108,38 +112,6 @@ impl WebCookieSession { } } -#[derive(Clone)] -struct LegacyWorkspaceSession { - client: Client, - cookie_header: String, - workspace_id: String, -} - -impl LegacyWorkspaceSession { - async fn resolve(client: &Client, cookies: &WebCookieSession) -> Result { - let workspace_id = legacy::discover_workspace_id(client, cookies.header()).await?; - Ok(Self { - client: client.clone(), - cookie_header: cookies.header.clone(), - workspace_id, - }) - } - - async fn fetch_usage(&self) -> Result { - legacy::fetch_usage(&self.client, &self.cookie_header, &self.workspace_id).await - } - - async fn fetch_balance(&self, timeout: Duration) -> Result, ProviderError> { - legacy::fetch_balance( - &self.client, - &self.cookie_header, - &self.workspace_id, - timeout, - ) - .await - } -} - impl OpenCodeGoProvider { pub fn new() -> Self { Self { @@ -167,25 +139,20 @@ impl OpenCodeGoProvider { console::normalize_workspace_id(workspace_id) } - async fn fetch_workspace_id( - client: &Client, - cookie_header: &str, - ) -> Result { - console::fetch_workspace_id(client, cookie_header, Duration::from_secs(30)).await - } - /// 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, cookies: &WebCookieSession, timeout: Duration, ) -> OptionalZenBalance { let request_timeout = timeout.min(ZEN_BALANCE_TIMEOUT); - match console::fetch_balance(client, workspace_id, cookies.header(), request_timeout).await + 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) => { @@ -198,8 +165,8 @@ impl OpenCodeGoProvider { /// 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, + fn spawn_zen_balance_task( + transport: Arc, cookies: &WebCookieSession, workspace_id_override: Option<&str>, web_timeout: u64, @@ -207,7 +174,6 @@ impl OpenCodeGoProvider { tokio::task::JoinHandle, std::time::Instant, ) { - let client = self.client.clone(); let cookies = cookies.clone(); let workspace_id_override = workspace_id_override.map(str::to_string); let timeout = Duration::from_secs(web_timeout.max(1)); @@ -216,7 +182,10 @@ impl OpenCodeGoProvider { tokio::time::sleep(ZEN_BALANCE_START_DELAY).await; let workspace_id = match workspace_id_override { Some(id) => id, - None => match Self::fetch_workspace_id(&client, cookies.header()).await { + None => match transport + .fetch_workspace_id(cookies.header(), Duration::from_secs(30)) + .await + { Ok(id) => id, Err(error) if cookies.can_recover_with_legacy(&error) => { return OptionalZenBalance::LegacyBalanceRequired; @@ -224,13 +193,14 @@ impl OpenCodeGoProvider { Err(_) => return OptionalZenBalance::Resolved(None), }, }; - Self::fetch_zen_balance(&client, &workspace_id, &cookies, timeout).await + Self::fetch_zen_balance(transport.as_ref(), &workspace_id, &cookies, timeout).await }); (task, started_at) } - fn spawn_legacy_balance_task( - session: LegacyWorkspaceSession, + fn spawn_legacy_balance_task( + transport: Arc, + session: T::LegacySession, web_timeout: u64, ) -> ( tokio::task::JoinHandle, @@ -240,19 +210,16 @@ impl OpenCodeGoProvider { let started_at = std::time::Instant::now(); let task = tokio::spawn(async move { tokio::time::sleep(ZEN_BALANCE_START_DELAY).await; - OptionalZenBalance::Resolved(session.fetch_balance(timeout).await.unwrap_or(None)) + OptionalZenBalance::Resolved( + transport + .fetch_legacy_balance(&session, timeout) + .await + .unwrap_or(None), + ) }); (task, started_at) } - async fn abort_optional_balance_and_resolve( - task: tokio::task::JoinHandle, - resolution: impl std::future::Future>, - ) -> Result { - task.abort(); - resolution.await - } - /// 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( @@ -270,8 +237,8 @@ impl OpenCodeGoProvider { } } - async fn resolve_legacy_balance( - client: &Client, + async fn resolve_legacy_balance( + transport: Arc, cookies: &WebCookieSession, started_at: std::time::Instant, requires_optional_usage_completeness: bool, @@ -281,11 +248,9 @@ impl OpenCodeGoProvider { return None; } tokio::time::timeout(budget, async { - let session = LegacyWorkspaceSession::resolve(client, cookies) - .await - .ok()?; - session - .fetch_balance(budget.min(ZEN_BALANCE_TIMEOUT)) + let session = transport.resolve_legacy_session(cookies).await.ok()?; + transport + .fetch_legacy_balance(&session, budget.min(ZEN_BALANCE_TIMEOUT)) .await .unwrap_or(None) }) @@ -294,8 +259,8 @@ impl OpenCodeGoProvider { .flatten() } - async fn finish_zen_balance( - client: &Client, + async fn finish_zen_balance( + transport: Arc, cookies: &WebCookieSession, task: tokio::task::JoinHandle, started_at: std::time::Instant, @@ -306,7 +271,7 @@ impl OpenCodeGoProvider { OptionalZenBalance::Resolved(balance) => balance, OptionalZenBalance::LegacyBalanceRequired => { Self::resolve_legacy_balance( - client, + transport, cookies, started_at, requires_optional_usage_completeness, @@ -320,16 +285,32 @@ impl OpenCodeGoProvider { &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, - None => match Self::fetch_workspace_id(&self.client, cookies.header()).await { + 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) - .await; + return Self::fetch_legacy_with_cookies( + ctx, + &cookies, + console_error, + transport, + ) + .await; } Err(error) => return Err(error), }, @@ -337,15 +318,19 @@ impl OpenCodeGoProvider { // 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(&cookies, Some(&workspace_id), ctx.web_timeout); - let console_result = console::fetch_usage( - &self.client, - &workspace_id, - cookies.header(), - Duration::from_secs(ctx.web_timeout.max(1)), - ) - .await; + 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) => { @@ -355,20 +340,21 @@ impl OpenCodeGoProvider { )); } Err(console_error) if cookies.can_recover_with_legacy(&console_error) => { - let session = match Self::abort_optional_balance_and_resolve( - zen_task, - LegacyWorkspaceSession::resolve(&self.client, &cookies), - ) - .await - { + zen_task.abort(); + 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, session) - .await; + return Self::fetch_legacy_with_session( + ctx, + &cookies, + console_error, + transport, + session, + ) + .await; } Err(error) => { zen_task.abort(); @@ -382,7 +368,7 @@ impl OpenCodeGoProvider { } None => { Self::finish_zen_balance( - &self.client, + transport, &cookies, zen_task, zen_started, @@ -394,30 +380,33 @@ impl OpenCodeGoProvider { Ok(Self::with_zen_balance(usage, "web", balance)) } - async fn fetch_legacy_with_cookies( - &self, + async fn fetch_legacy_with_cookies( ctx: &FetchContext, cookies: &WebCookieSession, console_error: ProviderError, + transport: Arc, ) -> Result { - let session = match LegacyWorkspaceSession::resolve(&self.client, cookies).await { + 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, session) - .await + Self::fetch_legacy_with_session(ctx, cookies, console_error, transport, session).await } - async fn fetch_legacy_with_session( - &self, + async fn fetch_legacy_with_session( ctx: &FetchContext, cookies: &WebCookieSession, console_error: ProviderError, - session: LegacyWorkspaceSession, + transport: Arc, + session: T::LegacySession, ) -> Result { - let (zen_task, zen_started) = - Self::spawn_legacy_balance_task(session.clone(), ctx.web_timeout); - let result = match cookies.select_legacy_result(console_error, session.fetch_usage().await) + 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) => { @@ -609,10 +598,15 @@ impl OpenCodeGoProvider { return Ok(result); }; let cookies = WebCookieSession::new(&cookie_header); - let (task, started) = - self.spawn_zen_balance_task(&cookies, ctx.workspace_id.as_deref(), ctx.web_timeout); + 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::finish_zen_balance( - &self.client, + transport, &cookies, task, started, @@ -656,164 +650,4 @@ fn zen_balance_join_budget( } #[cfg(test)] -mod tests { - use super::*; - use std::sync::{ - Arc, - atomic::{AtomicUsize, Ordering}, - }; - - #[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 resolver_calls = Arc::new(AtomicUsize::new(0)); - let balance_task = tokio::spawn(async { - tokio::task::yield_now().await; - OptionalZenBalance::LegacyBalanceRequired - }); - let calls = Arc::clone(&resolver_calls); - - let resolved = - OpenCodeGoProvider::abort_optional_balance_and_resolve(balance_task, async move { - calls.fetch_add(1, Ordering::SeqCst); - Ok::<_, ProviderError>("shared legacy session") - }) - .await - .unwrap(); - - assert_eq!(resolved, "shared legacy session"); - assert_eq!(resolver_calls.load(Ordering::SeqCst), 1); - } -} +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 + } +}