fix(buzz-agent): classify refresh 4xx by OAuth error, add cross-process auth tests

Refresh classification now keys on the OAuth error body, not the bare HTTP
status class. Per RFC 6749 §5.2 only `error == "invalid_grant"` means the
refresh token is dead — the one failure a browser sign-in can repair. Every
other 4xx (`invalid_request`, `invalid_client`, `unsupported_grant_type`,
`invalid_scope`, 408, 429), an unparseable error body, and all 5xx stay in the
infrastructural bucket as `NetworkUnavailable`, so a rate limit or a
misconfigured request can no longer pop a needless browser.

Prove the cross-process single-flight contract with a real second process. The
in-memory `INFLIGHT` registry coalesces same-key callers within one process
before the file lock, so two in-process handles cannot exercise the
cross-process protocol. A new `auth-worker` test binary runs the public
coordinator API against a shared temp cache and a scripted opener, driven by
barrier-marker files, covering (a) a `UserInitiated` denial in one process
shared with an already-waiting `Auto` in another and (b) two coordinator
processes racing to one grant and one cache artifact. Rename the two tests that
falsely claimed to be cross-process to reflect the in-process single-flight
they actually exercise. Run the coordinator integration suite on the Windows CI
job so `LockFileEx` contention and crash release execute rather than compile
only.

Gate the first provider response in the steer fold test so round 1 cannot
complete until the steer is sent and observed accepted, removing a
nextest-scheduling race in which round 2's boundary could drain an empty steer
queue before the steer was dispatched.

Co-authored-by: Will Pfleger <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
Duncan
2026-08-11 10:00:49 -04:00
co-authored by Will Pfleger
parent 373128bb75
commit c581d2eca9
7 changed files with 781 additions and 41 deletions
+7
View File
@@ -1024,6 +1024,13 @@ jobs:
# Serial: windows_resolver_tests mutate process-global env
# (BUZZ_SHELL/GIT_BASH/SystemRoot) that SharedState::new reads.
run: cargo test -p buzz-dev-mcp --target $env:TARGET -- --test-threads=1
- name: Test (buzz-agent auth coordinator)
# The auth coordinator single-flights on an OS advisory lock, which is
# LockFileEx on Windows; this integration suite drives real second
# processes on the same lock file, so it only exercises the Windows
# lock runtime (contention, crash release, cross-process cache) if it
# runs ON Windows. Every other job compiles it but never executes it.
run: cargo test -p buzz-agent --target $env:TARGET --test databricks_auth_coordinator
# Smoke-test the new host-prereq contract: Git for Windows (which provides
# bash) is available on the runner, a shell command round-trips, and bash
# does NOT resolve from System32 (so WSL's launcher is never picked up).
+10
View File
@@ -32,6 +32,16 @@ path = "tests/bin/fake_mcp.rs"
name = "lock-holder"
path = "tests/bin/lock_holder.rs"
# Test-only auth worker: a real second process that runs the PUBLIC auth
# coordinator API (`acquire_with_intent`) with a scripted browser opener and a
# shared temp cache, so the auth tests can prove the cross-process single-flight
# contract end-to-end — durable cooldown sharing and one-grant/one-cache races
# across a genuine process boundary, not two in-process handles. Only used by
# the databricks auth integration tests.
[[bin]]
name = "auth-worker"
path = "tests/bin/auth_worker.rs"
[dependencies]
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "io-std", "io-util", "sync", "process", "time", "net"] }
serde = { workspace = true }
+25 -13
View File
@@ -325,14 +325,15 @@ struct OidcEndpoints {
/// actual credential rejection from a transient fault.
///
/// - [`Refreshed`](Self::Refreshed): a fresh token — success.
/// - [`Rejected`](Self::Rejected): the token endpoint rejected the *grant*
/// (dead/rotated refresh token). This is the only outcome that becomes
/// [`AuthError::RefreshRejected`] for `Headless` or drives a browser
/// fallback for interactive intents.
/// - [`Network`](Self::Network): transport error, timeout, 5xx, or an
/// undecodable/malformed response — infrastructural, never a credential
/// decision, so it surfaces as [`AuthError::NetworkUnavailable`] and never
/// pops a browser.
/// - [`Rejected`](Self::Rejected): the token endpoint returned an
/// `invalid_grant` error (dead/rotated refresh token). This is the only
/// outcome that becomes [`AuthError::RefreshRejected`] for `Headless` or
/// drives a browser fallback for interactive intents.
/// - [`Network`](Self::Network): transport error, timeout, 5xx, any 4xx that
/// is not `invalid_grant` (e.g. `invalid_request`, `invalid_client`, 429),
/// an unparseable error body, or an undecodable/malformed success body —
/// infrastructural or misconfiguration, never a credential decision, so it
/// surfaces as [`AuthError::NetworkUnavailable`] and never pops a browser.
enum RefreshOutcome {
Refreshed(CachedToken),
Rejected,
@@ -511,14 +512,25 @@ impl PkceOAuthTokenSource {
let status = resp.status();
if !status.is_success() {
let body = resp.text().await.unwrap_or_default();
// A 4xx is the token endpoint rejecting the grant (dead/rotated
// refresh token). A 5xx is a provider-side fault — transient, not a
// credential decision — so it stays in the infrastructural bucket.
if status.is_client_error() {
// Per RFC 6749 §5.2 only `error == "invalid_grant"` means the
// refresh token itself is dead (expired/revoked) — the one failure
// a browser sign-in can repair. Every other 4xx (`invalid_request`,
// `invalid_client`, `unsupported_grant_type`, `invalid_scope`, 408,
// 429, …), an unparseable error body, and all 5xx are
// infrastructural or misconfiguration: a browser can't fix them, so
// they stay in the non-credential bucket and surface as
// `NetworkUnavailable` without ever popping a browser.
if status.is_client_error()
&& serde_json::from_str::<Value>(&body)
.ok()
.and_then(|v| v.get("error").and_then(Value::as_str).map(str::to_owned))
.as_deref()
== Some("invalid_grant")
{
tracing::warn!(status = %status, body = %body, "oauth refresh grant rejected");
return RefreshOutcome::Rejected;
}
tracing::warn!(status = %status, body = %body, "oauth refresh server error");
tracing::warn!(status = %status, body = %body, "oauth refresh not repairable by browser");
return RefreshOutcome::Network;
}
let v: Value = match resp.json().await {
+198
View File
@@ -0,0 +1,198 @@
//! Test-only helper: a real second process that runs the PUBLIC auth
//! coordinator (`PkceOAuthTokenSource::acquire_with_intent`) against a shared
//! temp cache, so the auth tests can prove the *cross-process* single-flight
//! contract end-to-end rather than with two in-process handles.
//!
//! The in-process `INFLIGHT` registry coalesces same-key callers within one
//! process before they ever reach the file lock, so two `PkceOAuthTokenSource`
//! instances in one test do NOT exercise the cross-process protocol (the OS
//! advisory lock and the on-disk cache re-read). This binary is a genuine
//! second process: it contends on the same `flock`/`LockFileEx` and reads/writes
//! the same private cache file the parent coordinator does.
//!
//! The browser step is scripted (no real window): the opener drives the
//! loopback callback exactly as a real browser would, and its launch count is
//! reported back so a test can assert "exactly one browser across processes".
//!
//! Env contract (all required unless noted):
//! AUTH_WORKER_DISCOVERY_URL — OIDC discovery URL (the parent stub).
//! AUTH_WORKER_CACHE_DIR — shared cache dir (`cache_dir_override`).
//! AUTH_WORKER_NAMESPACE — cache namespace.
//! AUTH_WORKER_CLIENT_ID — OAuth client id.
//! AUTH_WORKER_SCOPES — comma-separated scopes.
//! AUTH_WORKER_INTENT — auto | userinitiated | headless.
//! AUTH_WORKER_SCRIPT — approve | deny | failopen.
//! AUTH_WORKER_RESULT — path to write the JSON outcome to.
//! AUTH_WORKER_READY_MARKER — (optional) written once the source is built,
//! before acquisition, so the parent can release
//! several workers into a genuine lock race.
//! AUTH_WORKER_START_MARKER — (optional) acquisition blocks until this file
//! exists, so multiple workers begin together.
//! AUTH_WORKER_LAUNCHED_MARKER — (optional) written when the browser opener
//! fires (i.e. this process holds the lock and is
//! mid-flow), so the parent can queue behind it.
//! AUTH_WORKER_PROCEED_MARKER — (optional) the scripted callback is withheld
//! until this file exists, so the parent can
//! confirm another process is already waiting on
//! the lock before this one resolves.
//!
//! Result JSON: `{ "result": "ok"|"<error_code>", "bearer": <string|null>,
//! "launches": <u64> }`.
use std::fs;
use std::io::{Read, Write};
use std::net::TcpStream;
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Duration;
use buzz_agent::auth::{AuthIntent, BrowserOpener, PkceOAuthConfig, PkceOAuthTokenSource};
/// What the scripted "user" does when the coordinator opens a browser.
#[derive(Clone, Copy)]
enum Script {
Approve,
Deny,
FailToOpen,
}
/// A [`BrowserOpener`] that counts launches and drives the loopback callback on
/// a background thread — the same technique as the in-crate test opener, but
/// with two optional cross-process barriers so the parent can order events:
/// `launched_marker` announces that this process holds the lock and has opened
/// the browser, and `proceed_marker` withholds the callback until the parent
/// signals it has queued another process behind the lock.
struct WorkerOpener {
script: Script,
calls: Arc<AtomicU64>,
launched_marker: Option<PathBuf>,
proceed_marker: Option<PathBuf>,
}
impl BrowserOpener for WorkerOpener {
fn open(&self, url: &str) -> Result<(), String> {
self.calls.fetch_add(1, Ordering::SeqCst);
if let Some(marker) = &self.launched_marker {
fs::write(marker, b"launched").expect("write launched marker");
}
let query = match self.script {
Script::FailToOpen => return Err("no browser available".into()),
Script::Approve => "code=scripted-code",
Script::Deny => "error=access_denied",
};
let parsed = url::Url::parse(url).expect("authorize URL must parse");
let redirect = parsed
.query_pairs()
.find(|(k, _)| k == "redirect_uri")
.map(|(_, v)| v.into_owned())
.expect("authorize URL carries redirect_uri");
let state = parsed
.query_pairs()
.find(|(k, _)| k == "state")
.map(|(_, v)| v.into_owned())
.expect("authorize URL carries state");
let redirect = url::Url::parse(&redirect).expect("redirect_uri must parse");
let port = redirect.port().expect("loopback redirect carries a port");
let request = format!(
"GET /?{query}&state={state} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\nConnection: close\r\n\r\n"
);
let proceed = self.proceed_marker.clone();
std::thread::spawn(move || {
// Hold the callback until the parent has confirmed another process
// is already queued behind the lock (bounded so a missing signal
// can't wedge the test past the browser timeout).
if let Some(marker) = proceed {
for _ in 0..6000 {
if marker.exists() {
break;
}
std::thread::sleep(Duration::from_millis(10));
}
}
if let Ok(mut sock) = TcpStream::connect(("127.0.0.1", port)) {
let _ = sock.write_all(request.as_bytes());
let _ = sock.flush();
let mut discard = Vec::new();
let _ = sock.read_to_end(&mut discard);
}
});
Ok(())
}
}
fn env(key: &str) -> String {
std::env::var(key).unwrap_or_else(|_| panic!("{key} set"))
}
#[tokio::main]
async fn main() {
let intent = match env("AUTH_WORKER_INTENT").as_str() {
"auto" => AuthIntent::Auto,
"userinitiated" => AuthIntent::UserInitiated,
"headless" => AuthIntent::Headless,
other => panic!("unknown AUTH_WORKER_INTENT: {other}"),
};
let script = match env("AUTH_WORKER_SCRIPT").as_str() {
"approve" => Script::Approve,
"deny" => Script::Deny,
"failopen" => Script::FailToOpen,
other => panic!("unknown AUTH_WORKER_SCRIPT: {other}"),
};
let result_path = PathBuf::from(env("AUTH_WORKER_RESULT"));
let start_marker = std::env::var("AUTH_WORKER_START_MARKER")
.ok()
.map(PathBuf::from);
let ready_marker = std::env::var("AUTH_WORKER_READY_MARKER")
.ok()
.map(PathBuf::from);
let calls = Arc::new(AtomicU64::new(0));
let opener = WorkerOpener {
script,
calls: calls.clone(),
launched_marker: std::env::var("AUTH_WORKER_LAUNCHED_MARKER")
.ok()
.map(PathBuf::from),
proceed_marker: std::env::var("AUTH_WORKER_PROCEED_MARKER")
.ok()
.map(PathBuf::from),
};
let cfg = PkceOAuthConfig {
discovery_url: env("AUTH_WORKER_DISCOVERY_URL"),
client_id: env("AUTH_WORKER_CLIENT_ID"),
scopes: env("AUTH_WORKER_SCOPES")
.split(',')
.map(str::to_owned)
.collect(),
cache_namespace: env("AUTH_WORKER_NAMESPACE"),
cache_dir_override: Some(PathBuf::from(env("AUTH_WORKER_CACHE_DIR"))),
};
let src = PkceOAuthTokenSource::new_with(cfg, Arc::new(opener)).expect("build token source");
// Announce readiness, then wait for the parent's release so several workers
// hit the lock together — a genuine race rather than staggered spawns.
if let Some(marker) = &ready_marker {
fs::write(marker, b"ready").expect("write ready marker");
}
if let Some(marker) = start_marker {
for _ in 0..6000 {
if marker.exists() {
break;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
}
let (result, bearer) = match src.acquire_with_intent(intent, None).await {
Ok(token) => ("ok".to_owned(), Some(token)),
Err(e) => (e.code().to_owned(), None),
};
let body = serde_json::json!({
"result": result,
"bearer": bearer,
"launches": calls.load(Ordering::SeqCst),
});
fs::write(&result_path, serde_json::to_vec(&body).unwrap()).expect("write result file");
}
@@ -1,17 +1,23 @@
//! Concurrency-matrix tests for the Databricks auth coordinator.
//!
//! The coordinator single-flights OAuth acquisition per cache key using an OS
//! advisory lock, so one browser dance is shared and failures are coalesced
//! The coordinator single-flights OAuth acquisition per cache key. Within one
//! process, same-key callers coalesce on an in-memory `INFLIGHT` registry
//! *before* the file lock; across processes, they serialize on an OS advisory
//! lock and share success through the on-disk cache, with failures coalesced
//! through a durable cooldown sidecar. These tests drive the public API
//! (`acquire_with_intent`, `interactive_login`) with an injected
//! [`BrowserOpener`] that scripts the localhost callback instead of popping a
//! real window — the browser step becomes deterministic and countable.
//!
//! Two `PkceOAuthTokenSource` instances sharing one cache path model two
//! processes: `File::try_lock` is per open-file-description, so distinct
//! handles contend whether or not they live in the same process. The
//! lock-primitive, crash-release, and lock-timeout edges live in the in-crate
//! `auth::tests` module where the private helpers are reachable.
//! Two `PkceOAuthTokenSource` instances in ONE process do not model two
//! processes: the `INFLIGHT` registry intercepts them before the file lock, so
//! same-process tests exercise the in-memory single-flight, not the
//! cross-process protocol. The genuinely cross-process claims — lock
//! contention, crash release, cooldown sharing across a process boundary, and
//! one-grant/one-cache under a real race — are proved with the `lock-holder`
//! and `auth-worker` helper binaries, each a real second process on the same
//! lock file and cache. The lock-primitive and lock-timeout edges live in the
//! in-crate `auth::tests` module where the private helpers are reachable.
use std::io::Write;
use std::net::{SocketAddr, TcpStream};
@@ -141,6 +147,11 @@ enum RefreshMode {
/// `500` — a provider-side fault, transient rather than a credential
/// decision.
ServerError,
/// A 4xx with the given OAuth `error` code in the body. Lets a test assert
/// the coordinator treats `invalid_grant` (any 4xx) as a dead grant, but
/// every other error code — and any non-`invalid_grant` status like `429`
/// — as infrastructural rather than a credential rejection.
ClientError(axum::http::StatusCode, &'static str),
/// Sleep `d` before answering, so the caller's per-request HTTP timeout
/// elapses first (a transport timeout, not a verdict from the provider).
Hang(Duration),
@@ -211,6 +222,9 @@ async fn spawn_stub_with(mode: RefreshMode) -> Stub {
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
Json(json!({ "error": "temporarily_unavailable" })),
),
RefreshMode::ClientError(status, error) => {
(status, Json(json!({ "error": error })))
}
RefreshMode::Succeed | RefreshMode::Hang(_) => (
axum::http::StatusCode::OK,
Json(json!({
@@ -301,7 +315,11 @@ async fn test_same_key_concurrent_callers_share_one_browser_attempt() {
let cache = TempDir::new().unwrap();
let opener = ScriptedOpener::new(Script::Approve);
// Two independent sources on the SAME cache key = two processes racing.
// Two independent sources on the same key in ONE process. The in-memory
// INFLIGHT registry coalesces them before the file lock, so this proves the
// in-process single-flight — one leader runs the browser flow, the other
// joins its published result. The genuine cross-process race is
// `test_crossprocess_two_coordinators_race_to_one_grant_and_cache`.
let a = PkceOAuthTokenSource::new_with(
config(&stub, "/disco/a", cache.path()),
Arc::new(opener.clone()),
@@ -371,42 +389,45 @@ async fn test_denied_then_auto_reads_cooldown_without_second_launch() {
}
#[tokio::test]
async fn test_userinitiated_denial_is_visible_to_crossprocess_auto() {
async fn test_denial_sidecar_is_read_by_later_auto_in_same_process() {
let stub = spawn_stub(false).await;
let cache = TempDir::new().unwrap();
let opener = ScriptedOpener::new(Script::Deny);
// Process A: an explicit UserInitiated attempt is denied.
let proc_a = PkceOAuthTokenSource::new_with(
// An explicit UserInitiated attempt is denied and records the durable
// cooldown sidecar.
let denier = PkceOAuthTokenSource::new_with(
config(&stub, "/disco/a", cache.path()),
Arc::new(opener.clone()),
)
.unwrap();
let denied = proc_a
let denied = denier
.acquire_with_intent(AuthIntent::UserInitiated, None)
.await;
assert_eq!(denied, Err(AuthError::Denied));
assert_eq!(opener.call_count(), 1);
// Process B: a passive Auto caller (distinct instance = distinct process)
// reads the durable sidecar A wrote and does not launch a second browser.
// This is the cross-policy edge: the sidecar is written for ANY failed
// interactive attempt, only the reader policy differs.
let proc_b = PkceOAuthTokenSource::new_with(
// A later passive Auto caller reads the sidecar and does not launch a
// second browser. This is the cross-policy read edge — the sidecar is
// written for ANY failed interactive attempt, only the reader policy
// differs. It is a SEQUENTIAL, same-process read; the genuinely
// cross-process, already-waiting variant is
// `test_crossprocess_userinitiated_denial_shared_with_waiting_auto`.
let later = PkceOAuthTokenSource::new_with(
config(&stub, "/disco/a", cache.path()),
Arc::new(opener.clone()),
)
.unwrap();
let auto = proc_b.acquire_with_intent(AuthIntent::Auto, None).await;
let auto = later.acquire_with_intent(AuthIntent::Auto, None).await;
assert_eq!(
auto,
Err(AuthError::Denied),
"cross-process Auto reads the UserInitiated failure sidecar"
"a later Auto reads the UserInitiated failure sidecar"
);
assert_eq!(
opener.call_count(),
1,
"no second browser across the policy/process boundary"
"no second browser once the denial sidecar is recorded"
);
}
@@ -842,7 +863,131 @@ async fn test_refresh_server_error_is_network_unavailable_not_rejected() {
assert_eq!(stub.refresh_grants.load(Ordering::SeqCst), 1);
}
// ---- in-process joiner shares the leader's FAILURE result ----------------
// ---- 4xx classification: only `invalid_grant` is a dead refresh token -----
//
// RFC 6749 §5.2 uses 400/401 token responses for several `error` codes, but
// only `invalid_grant` means the refresh token is dead. Every other 4xx —
// `invalid_request`, `invalid_client`, `unsupported_grant_type`,
// `invalid_scope`, `408`, `429` — is a request/config/transient fault a
// browser cannot repair, so it must stay infrastructural (`NetworkUnavailable`)
// and never pop a browser. The classifier keys on the OAuth error body, not
// the bare status class.
#[tokio::test]
async fn test_refresh_400_invalid_grant_is_dead_grant_not_network() {
// A 400 (not just 401) carrying `invalid_grant` is still a dead refresh
// token, so a Headless caller must classify it terminally as
// RefreshRejected — proving the decision is the body error, not the status.
let stub = spawn_stub_with(RefreshMode::ClientError(
axum::http::StatusCode::BAD_REQUEST,
"invalid_grant",
))
.await;
let cache = TempDir::new().unwrap();
let opener = ScriptedOpener::new(Script::Approve);
let cfg = config(&stub, "/disco/a", cache.path());
seed_cache(
&cfg,
cache.path(),
json!({
"access_token": "stale",
"refresh_token": "dead-refresh",
"expires_at": 1u64,
}),
);
let src = PkceOAuthTokenSource::new_with(cfg, Arc::new(opener.clone())).unwrap();
let result = src.acquire_with_intent(AuthIntent::Headless, None).await;
assert_eq!(
result,
Err(AuthError::RefreshRejected),
"a 400 invalid_grant is a dead refresh token, not infrastructural"
);
assert_eq!(opener.call_count(), 0, "Headless never opens a browser");
assert_eq!(stub.refresh_grants.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn test_refresh_400_invalid_request_is_network_unavailable_not_rejected() {
// A 400 `invalid_request` is a malformed/misconfigured request, not a dead
// credential. An interactive intent that COULD open a browser must NOT: a
// browser cannot repair it, so it surfaces as NetworkUnavailable.
let stub = spawn_stub_with(RefreshMode::ClientError(
axum::http::StatusCode::BAD_REQUEST,
"invalid_request",
))
.await;
let cache = TempDir::new().unwrap();
let opener = ScriptedOpener::new(Script::Approve);
let cfg = config(&stub, "/disco/a", cache.path());
seed_cache(
&cfg,
cache.path(),
json!({
"access_token": "stale",
"refresh_token": "misconfigured-refresh",
"expires_at": 1u64,
}),
);
let src = PkceOAuthTokenSource::new_with(cfg, Arc::new(opener.clone())).unwrap();
let result = src
.acquire_with_intent(AuthIntent::UserInitiated, None)
.await;
assert_eq!(
result,
Err(AuthError::NetworkUnavailable),
"a non-invalid_grant 4xx is infrastructural, not a credential rejection"
);
assert_eq!(
opener.call_count(),
0,
"a browser cannot repair invalid_request, so none is opened"
);
assert_eq!(stub.refresh_grants.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn test_refresh_429_is_network_unavailable_not_rejected() {
// A 429 rate limit is a transient 4xx: retry later, don't sign in again.
// An interactive intent must not pop a browser off it.
let stub = spawn_stub_with(RefreshMode::ClientError(
axum::http::StatusCode::TOO_MANY_REQUESTS,
"slow_down",
))
.await;
let cache = TempDir::new().unwrap();
let opener = ScriptedOpener::new(Script::Approve);
let cfg = config(&stub, "/disco/a", cache.path());
seed_cache(
&cfg,
cache.path(),
json!({
"access_token": "stale",
"refresh_token": "rate-limited-refresh",
"expires_at": 1u64,
}),
);
let src = PkceOAuthTokenSource::new_with(cfg, Arc::new(opener.clone())).unwrap();
let result = src
.acquire_with_intent(AuthIntent::UserInitiated, None)
.await;
assert_eq!(
result,
Err(AuthError::NetworkUnavailable),
"a 429 rate limit is transient, not a credential rejection"
);
assert_eq!(
opener.call_count(),
0,
"a rate limit must not trigger an interactive browser fallback"
);
assert_eq!(stub.refresh_grants.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn test_two_concurrent_userinitiated_denials_share_one_browser() {
@@ -966,3 +1111,226 @@ async fn test_crossprocess_lock_holder_blocks_then_crash_release_lets_successor_
"Headless successor recovers via refresh without a browser"
);
}
// ---- genuine cross-process coordinator races -----------------------------
//
// The `auth-worker` helper is a real second process running the PUBLIC
// coordinator API against the shared cache. Unlike two in-process handles
// (which the `INFLIGHT` registry coalesces before the file lock), these
// workers contend on the OS advisory lock and share success through the
// on-disk cache exactly as two Buzz processes on one machine would.
/// A spawned `auth-worker`: its child handle plus the file it writes its JSON
/// outcome to.
struct Worker {
child: tokio::process::Child,
result_path: std::path::PathBuf,
}
#[derive(Deserialize)]
struct WorkerOutcome {
result: String,
bearer: Option<String>,
launches: u64,
}
impl Worker {
/// Block until the worker exits, then parse its outcome file.
async fn join(mut self) -> WorkerOutcome {
let status = self.child.wait().await.expect("auth-worker joins");
assert!(
status.success(),
"auth-worker exited with failure: {status}"
);
let body = std::fs::read(&self.result_path).expect("auth-worker wrote its outcome");
serde_json::from_slice(&body).expect("auth-worker outcome parses")
}
}
/// Spawn an `auth-worker` child against `cfg`'s shared cache. `extra` sets the
/// optional barrier-marker env vars ((name, path) pairs) a scenario needs to
/// order events across processes.
fn spawn_worker(
cfg: &PkceOAuthConfig,
cache_dir: &std::path::Path,
intent: &str,
script: &str,
tag: &str,
extra: &[(&str, &std::path::Path)],
) -> Worker {
let result_path = cache_dir.join(format!("{tag}.result.json"));
let mut cmd = tokio::process::Command::new(env!("CARGO_BIN_EXE_auth-worker"));
cmd.env("AUTH_WORKER_DISCOVERY_URL", &cfg.discovery_url)
.env("AUTH_WORKER_CACHE_DIR", cache_dir)
.env("AUTH_WORKER_NAMESPACE", &cfg.cache_namespace)
.env("AUTH_WORKER_CLIENT_ID", &cfg.client_id)
.env("AUTH_WORKER_SCOPES", cfg.scopes.join(","))
.env("AUTH_WORKER_INTENT", intent)
.env("AUTH_WORKER_SCRIPT", script)
.env("AUTH_WORKER_RESULT", &result_path)
.kill_on_drop(true);
for (key, path) in extra {
cmd.env(key, path);
}
let child = cmd.spawn().expect("spawn the auth-worker helper process");
Worker { child, result_path }
}
async fn wait_for_marker(path: &std::path::Path, what: &str) {
for _ in 0..1000 {
if path.exists() {
return;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
panic!("timed out waiting for {what} ({})", path.display());
}
#[tokio::test]
async fn test_crossprocess_userinitiated_denial_shared_with_waiting_auto() {
// Two real processes on one key. The child runs a UserInitiated flow that
// is denied; while it holds the lock and its browser is open, the parent's
// Auto coordinator is already WAITING on the cross-process lock. The child
// must be released only once the parent is queued, so the denial the child
// records is what the waiting Auto observes — one launch total, durable
// Denied for both, across a genuine process boundary.
let stub = spawn_stub(false).await;
let cache = TempDir::new().unwrap();
let cfg = config(&stub, "/disco/a", cache.path());
let launched = cache.path().join("child.launched");
let proceed = cache.path().join("child.proceed");
let child = spawn_worker(
&cfg,
cache.path(),
"userinitiated",
"deny",
"denier",
&[
("AUTH_WORKER_LAUNCHED_MARKER", launched.as_path()),
("AUTH_WORKER_PROCEED_MARKER", proceed.as_path()),
],
);
// Wait until the child holds the lock and has opened its (scripted)
// browser; its callback is withheld until we create `proceed`.
wait_for_marker(&launched, "child browser launch").await;
// The parent's Auto coordinator now contends for the same lock. It cannot
// proceed while the child holds it, so it is a genuine cross-process
// waiter.
let parent = PkceOAuthTokenSource::new_with(
config(&stub, "/disco/a", cache.path()),
Arc::new(ScriptedOpener::new(Script::Approve)),
)
.unwrap();
let auto =
tokio::spawn(async move { parent.acquire_with_intent(AuthIntent::Auto, None).await });
tokio::time::sleep(Duration::from_millis(300)).await;
assert!(
!auto.is_finished(),
"parent Auto must block while the child process holds the lock"
);
// Release the child's callback: it finishes the denial and writes the
// cooldown sidecar, then drops the lock.
std::fs::write(&proceed, b"go").unwrap();
let child_outcome = child.join().await;
assert_eq!(
child_outcome.result, "denied",
"child UserInitiated is denied"
);
assert_eq!(child_outcome.launches, 1, "child opens exactly one browser");
let auto_result = auto.await.expect("parent Auto task joins");
assert_eq!(
auto_result,
Err(AuthError::Denied),
"the already-waiting Auto reads the child's durable denial"
);
assert_eq!(
stub.code_grants.load(Ordering::SeqCst),
0,
"a denied flow never reaches the code exchange"
);
}
#[tokio::test]
async fn test_crossprocess_two_coordinators_race_to_one_grant_and_cache() {
// Two real coordinator processes race on one key from a cold cache. They
// are released together (via a shared start marker) so both contend for the
// lock. Exactly one wins the browser flow and performs the single code
// grant; the other serializes behind the lock and adopts the winner's token
// from the shared cache. Both must observe the same bearer, and the private
// cache must hold exactly one parseable token artifact.
let stub = spawn_stub(false).await;
let cache = TempDir::new().unwrap();
let cfg = config(&stub, "/disco/a", cache.path());
let ready_a = cache.path().join("a.ready");
let ready_b = cache.path().join("b.ready");
let start = cache.path().join("start");
let worker_a = spawn_worker(
&cfg,
cache.path(),
"userinitiated",
"approve",
"a",
&[
("AUTH_WORKER_READY_MARKER", ready_a.as_path()),
("AUTH_WORKER_START_MARKER", start.as_path()),
],
);
let worker_b = spawn_worker(
&cfg,
cache.path(),
"userinitiated",
"approve",
"b",
&[
("AUTH_WORKER_READY_MARKER", ready_b.as_path()),
("AUTH_WORKER_START_MARKER", start.as_path()),
],
);
// Both processes are built and about to acquire; release them together.
wait_for_marker(&ready_a, "worker A ready").await;
wait_for_marker(&ready_b, "worker B ready").await;
std::fs::write(&start, b"go").unwrap();
let (out_a, out_b) = tokio::join!(worker_a.join(), worker_b.join());
assert_eq!(out_a.result, "ok", "worker A authenticates");
assert_eq!(out_b.result, "ok", "worker B authenticates");
let bearer_a = out_a.bearer.expect("worker A returns a bearer");
let bearer_b = out_b.bearer.expect("worker B returns a bearer");
assert_eq!(
bearer_a, bearer_b,
"both processes observe the same bearer from the shared cache"
);
// Exactly one browser launch and one code exchange across both processes.
assert_eq!(
out_a.launches + out_b.launches,
1,
"exactly one browser launch across the two coordinator processes"
);
assert_eq!(
stub.code_grants.load(Ordering::SeqCst),
1,
"exactly one authorization-code exchange across both processes"
);
// The private cache holds exactly one parseable token artifact carrying the
// shared bearer.
let cache_path = cache_file_path(&cfg, cache.path());
let raw = std::fs::read(&cache_path).expect("cache file exists");
let cached: serde_json::Value =
serde_json::from_slice(&raw).expect("cache holds one parseable token artifact");
assert_eq!(
cached.get("access_token").and_then(|v| v.as_str()),
Some(bearer_a.as_str()),
"the cached token is the shared bearer"
);
}
+141 -7
View File
@@ -164,6 +164,108 @@ async fn spawn_capturing_fake_llm_with_statuses(
(url, captures)
}
/// A capturing fake LLM whose FIRST provider response is withheld until
/// `gate` fires. Later responses are served immediately. Used to make
/// round-boundary races deterministic: hold round 1 open until a client action
/// (e.g. a steer) is confirmed, so the second round observes it. Request bodies
/// are recorded into `captures` exactly as `spawn_capturing_fake_llm` does.
async fn spawn_gated_capturing_fake_llm(
responses: Vec<CannedResponse>,
captures: Arc<Mutex<Vec<Value>>>,
gate: Arc<Mutex<Option<tokio::sync::oneshot::Receiver<()>>>>,
) -> (String, Arc<Mutex<Vec<Value>>>) {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
let queue = Arc::new(Mutex::new(VecDeque::from(responses)));
let captures_clone = captures.clone();
tokio::spawn(async move {
let mut request_num = 0usize;
loop {
let (mut sock, _) = match listener.accept().await {
Ok(p) => p,
Err(_) => return,
};
let queue = queue.clone();
let captures = captures_clone.clone();
let gate = gate.clone();
request_num += 1;
let req_num = request_num;
tokio::spawn(async move {
// Read headers.
let mut buf = Vec::new();
let mut tmp = [0u8; 4096];
while !buf.windows(4).any(|w| w == b"\r\n\r\n") {
match sock.read(&mut tmp).await {
Ok(0) | Err(_) => return,
Ok(n) => buf.extend_from_slice(&tmp[..n]),
}
if buf.len() > 2_000_000 {
return;
}
}
let header_end = buf.windows(4).position(|w| w == b"\r\n\r\n").unwrap() + 4;
let header_str = String::from_utf8_lossy(&buf[..header_end]);
let content_length: usize = header_str
.lines()
.find_map(|line| {
let lower = line.to_lowercase();
if lower.starts_with("content-length:") {
lower
.trim_start_matches("content-length:")
.trim()
.parse()
.ok()
} else {
None
}
})
.unwrap_or(0);
let mut body_buf = buf[header_end..].to_vec();
while body_buf.len() < content_length {
match sock.read(&mut tmp).await {
Ok(0) | Err(_) => break,
Ok(n) => body_buf.extend_from_slice(&tmp[..n]),
}
}
if let Ok(parsed) =
serde_json::from_slice::<Value>(&body_buf[..content_length.min(body_buf.len())])
{
captures.lock().await.push(parsed);
}
// Hold the first request's response until the gate opens.
if req_num == 1 {
let rx = gate.lock().await.take();
if let Some(rx) = rx {
let _ = rx.await;
}
}
let response = queue.lock().await.pop_front().unwrap_or(CannedResponse {
status: 500,
body: json!({ "error": "no canned response" }),
});
let body_s = serde_json::to_string(&response.body).unwrap();
let reason = if response.status == 200 {
"OK"
} else {
"Error"
};
let resp = format!(
"HTTP/1.1 {} {}\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}",
response.status,
reason,
body_s.len(),
body_s,
);
let _ = sock.write_all(resp.as_bytes()).await;
let _ = sock.shutdown().await;
});
}
});
(url, captures)
}
struct Harness {
child: tokio::process::Child,
stdin: tokio::process::ChildStdin,
@@ -776,14 +878,37 @@ async fn recv_active_run_id(h: &mut Harness) -> String {
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn steer_folds_into_active_turn_without_cancelling() {
use tokio::sync::oneshot;
// A two-round turn (tool call → text). A steer sent once the run is live
// must (a) be accepted with the matching runId, (b) NOT cancel the turn —
// it still ends with end_turn — and (c) reach the provider as a user turn.
let (url, captures) = spawn_capturing_fake_llm(vec![
openai_tool_call("call_steer", "fake__noop", json!({})),
openai_text("acknowledged the steer"),
])
.await;
//
// The steer is drained only at a round boundary (before the next provider
// request), so it must be enqueued before round 2 begins. Without
// synchronization a fast worker can complete round 1, drain an empty steer
// queue at the round-2 boundary, and dispatch round 2 before the steer is
// even sent — the steer then lands after the turn ends and never reaches
// the provider. To make this deterministic, the FIRST provider response is
// gated: it is withheld until the steer has been sent AND observed
// accepted, so round 1 cannot complete (and round 2 cannot start its drain)
// until the steer is already queued.
let (gate_tx, gate_rx) = oneshot::channel::<()>();
let gate_rx = Arc::new(Mutex::new(Some(gate_rx)));
let responses = vec![
CannedResponse {
status: 200,
body: openai_tool_call("call_steer", "fake__noop", json!({})),
},
CannedResponse {
status: 200,
body: openai_text("acknowledged the steer"),
},
];
let captures: Arc<Mutex<Vec<Value>>> = Arc::new(Mutex::new(Vec::new()));
let (url, _) = spawn_gated_capturing_fake_llm(responses, captures.clone(), gate_rx).await;
let mut h = Harness::spawn(&url).await;
let sid = init_session(&mut h).await;
@@ -797,7 +922,8 @@ async fn steer_folds_into_active_turn_without_cancelling() {
)
.await;
// Learn the run id, then steer into it before the turn finishes.
// Learn the run id (advertised before the gated round-1 request), then steer
// into the live turn while round 1 is still held.
let run_id = recv_active_run_id(&mut h).await;
let steer_text = "STEER-CANARY: also consider the edge case";
let s_id = h
@@ -811,9 +937,12 @@ async fn steer_folds_into_active_turn_without_cancelling() {
)
.await;
// Steer is accepted and echoes the run id it landed in.
// Steer is accepted and echoes the run id it landed in. Only after this
// confirmation do we release the gate, so the steer is guaranteed queued
// before round 2's boundary drains it.
let mut steer_ok = false;
let mut end_turn = false;
let mut gate = Some(gate_tx);
for _ in 0..40 {
let v = h.recv().await;
if v["id"] == json!(s_id) {
@@ -829,6 +958,11 @@ async fn steer_folds_into_active_turn_without_cancelling() {
"steer reply carries a messageId"
);
steer_ok = true;
// Steer accepted — release round 1 so the turn proceeds to round 2,
// whose boundary now drains the queued steer.
if let Some(tx) = gate.take() {
let _ = tx.send(());
}
} else if v["id"] == json!(p_id) {
// The turn was NOT cancelled — it completed normally.
assert_eq!(v["result"]["stopReason"], "end_turn");
+11
View File
@@ -1022,6 +1022,7 @@ dependencies = [
"axum",
"base64 0.22.1",
"dirs",
"fs2",
"getrandom 0.4.3",
"hex",
"nix 0.31.3",
@@ -2999,6 +3000,16 @@ dependencies = [
"percent-encoding",
]
[[package]]
name = "fs2"
version = "0.4.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9564fc758e15025b46aa6643b1b77d047d1a56a1aea6e01002ac0c7026876213"
dependencies = [
"libc",
"winapi",
]
[[package]]
name = "fs_extra"
version = "1.3.0"