diff --git a/docs/public/administration/server-configuration.mdx b/docs/public/administration/server-configuration.mdx index 49a4481381..33aab179b3 100644 --- a/docs/public/administration/server-configuration.mdx +++ b/docs/public/administration/server-configuration.mdx @@ -255,7 +255,9 @@ honors those hand-edited values even though the browser wizard does not manage t ### SQLite state and migration backups -Shared relational state, including vault entries and server-managed definitions, lives at `/db/fabro.sqlite3`. Run events continue to use the `[server.slatedb]` object store. +Shared relational state, including vault entries, server-managed definitions, and CLI auth sessions, lives at `/db/fabro.sqlite3`. Run events continue to use the `[server.slatedb]` object store. + +CLI auth sessions are stored as an `auth_sessions` row per signed-in CLI, with the rotating refresh tokens for that session in `refresh_tokens`. Revoking a session from **Settings → Sessions**, or with `DELETE /api/v1/auth/sessions/{id}`, deletes the session row and its tokens together. Before applying pending SQLite migrations, Fabro creates `/db/fabro.sqlite3.pre-migration.bak` with SQLite's `VACUUM INTO`. Each migration run replaces the previous snapshot, so only the most recent pre-migration backup is retained. diff --git a/docs/public/changelog/2026-07-26.mdx b/docs/public/changelog/2026-07-26.mdx new file mode 100644 index 0000000000..87d42da286 --- /dev/null +++ b/docs/public/changelog/2026-07-26.mdx @@ -0,0 +1,20 @@ +--- +title: "CLI auth sessions move to SQLite" +date: "2026-07-26" +--- + + +**Everyone signs in again after this upgrade.** Existing refresh tokens are not migrated, so every signed-in browser and CLI is logged out the moment the new server binary starts. Run `fabro auth login` again on each machine. There is no staged rollout for this — avoid upgrading mid-task. + + +## Sessions are their own record + +A CLI login is a chain of refresh tokens that rotate on every use. Fabro previously stored the identity and profile on each token in that chain, so the chain itself had no record of its own. It now does: an `auth_sessions` row per login, with its tokens in `refresh_tokens`, both in `/db/fabro.sqlite3`. + +Two dates on **Settings → Sessions** were wrong as a result and are now correct. A session's **created** date came from its newest token, so it moved forward every time the CLI refreshed, and **last seen** showed the same value rather than the last time the session was actually used. + +Listing and revoking sessions no longer reads every refresh token the server has ever issued, so both stay fast as a workspace accumulates logins. Revoking a session removes its tokens in the same operation. + +## Refresh token replay + +Replaying a refresh token still revokes its whole chain immediately. One detail changed: when several requests present the same already-rotated token at once, later ones now report `refresh_token_expired` where they previously reported `refresh_token_revoked`. The CLI treats both the same way — it discards the stored credentials and prompts you to sign in again. diff --git a/docs/public/docs.json b/docs/public/docs.json index d6650e55ae..8b95430cd6 100644 --- a/docs/public/docs.json +++ b/docs/public/docs.json @@ -295,6 +295,7 @@ "group": "July 2026", "icon": "clock-rotate-left", "pages": [ + "changelog/2026-07-26", "changelog/2026-07-25", "changelog/2026-07-24", "changelog/2026-07-23", diff --git a/lib/apps/fabro-server/src/auth/cli_flow.rs b/lib/apps/fabro-server/src/auth/cli_flow.rs index 142dd7b667..21d6848163 100644 --- a/lib/apps/fabro-server/src/auth/cli_flow.rs +++ b/lib/apps/fabro-server/src/auth/cli_flow.rs @@ -28,7 +28,8 @@ use url::{Host, Url}; use crate::auth::browser_shell::browser_shell; use crate::auth::{ - self, AuthCode, AuthErrorCode, ConsumeOutcome, JwtSubject, REFRESH_TOKEN_PREFIX, RefreshToken, + self, AuthCode, AuthErrorCode, AuthSessionRecord, JwtSubject, REFRESH_TOKEN_PREFIX, + RefreshToken, RotateOutcome, }; use crate::jwt_auth::{AuthMode, ConfiguredAuth, bearer_token_from_headers}; use crate::principal_middleware::{ @@ -466,32 +467,30 @@ async fn token( let refresh_expires_at = now + chrono::Duration::days(REFRESH_TOKEN_TTL_DAYS); let refresh_secret = random_secret(); let refresh_token = format!("{REFRESH_TOKEN_PREFIX}{refresh_secret}"); - let refresh_row = RefreshToken { - token_hash: hash_refresh_secret(&refresh_secret), - chain_id: uuid::Uuid::new_v4(), + let session = AuthSessionRecord { + id: uuid::Uuid::new_v4(), identity: entry.identity.clone(), login: entry.login.clone(), name: entry.name.clone(), email: entry.email.clone(), avatar_url: entry.avatar_url.clone(), - issued_at: now, - expires_at: refresh_expires_at, - last_used_at: now, - used: false, user_agent: sanitize_user_agent(request_user_agent(&headers)), + created_at: now, + last_used_at: now, }; - let auth_tokens = match state.store_ref().refresh_tokens().await { - Ok(store) => store, - Err(err) => { - warn!(error = %err, "Failed to open refresh token store"); - return oauth_error( - StatusCode::INTERNAL_SERVER_ERROR, - "server_error", - "Could not complete authentication", - ); - } + let refresh_row = RefreshToken { + token_hash: hash_refresh_secret(&refresh_secret), + session_id: session.id, + issued_at: now, + expires_at: refresh_expires_at, + used_at: None, }; - if let Err(err) = auth_tokens.insert_refresh_token(refresh_row.clone()).await { + if let Err(err) = state + .stores + .auth_sessions + .create_session(&session, &refresh_row) + .await + { warn!(error = %err, "Failed to persist refresh token"); return oauth_error( StatusCode::INTERNAL_SERVER_ERROR, @@ -516,7 +515,7 @@ async fn token( ); log_cli_auth_tokens_issued(&entry.login, &entry.email); - auth_slot.replace(refresh_user_context(&refresh_row)); + auth_slot.replace(refresh_user_context(&session)); Json(CliTokenResponse { access_token, @@ -524,10 +523,10 @@ async fn token( refresh_token, refresh_token_expires_at: refresh_expires_at, subject: subject_response( - &refresh_row.identity, - &refresh_row.login, - &refresh_row.name, - &refresh_row.email, + &session.identity, + &session.login, + &session.name, + &session.email, ), }) .into_response() @@ -570,37 +569,19 @@ async fn refresh( "Could not refresh authentication", ); }; - let auth_tokens = match state.store_ref().refresh_tokens().await { - Ok(store) => store, - Err(err) => { - warn!(error = %err, "Failed to open refresh token store"); - return oauth_error( - StatusCode::INTERNAL_SERVER_ERROR, - "server_error", - "Could not refresh authentication", - ); - } - }; + let auth_sessions = &state.stores.auth_sessions; let now = chrono::Utc::now(); let secret_hash = hash_refresh_secret(&secret); - let existing = match auth_tokens.find_refresh_token(&secret_hash).await { - Ok(existing) => existing, - Err(err) => { - warn!(error = %err, "Failed to load refresh token before rotation"); - return oauth_error( - StatusCode::INTERNAL_SERVER_ERROR, - "server_error", - "Could not refresh authentication", - ); - } - }; let next_secret = random_secret(); let next_user_agent = sanitize_user_agent(request_user_agent(&headers)); - let outcome = match auth_tokens - .consume_and_rotate( - secret_hash, - next_refresh_row(existing.as_ref(), &next_secret, &next_user_agent, now), + let refresh_expires_at = now + chrono::Duration::days(REFRESH_TOKEN_TTL_DAYS); + let outcome = match auth_sessions + .rotate( + &secret_hash, + &hash_refresh_secret(&next_secret), + refresh_expires_at, + &next_user_agent, now, ) .await @@ -616,42 +597,34 @@ async fn refresh( } }; - let (old, new_row) = match outcome { - ConsumeOutcome::NotFound | ConsumeOutcome::Expired => { + let session = match outcome { + RotateOutcome::NotFound | RotateOutcome::Expired => { auth_slot.replace(RequestAuthContext::invalid()); - if auth_tokens.was_recently_replay_revoked(&secret_hash, now) { - return oauth_error( - StatusCode::UNAUTHORIZED, - "refresh_token_revoked", - "Refresh token revoked", - ); - } return oauth_error( StatusCode::UNAUTHORIZED, "refresh_token_expired", "Refresh token expired", ); } - ConsumeOutcome::Reused(old) => { + RotateOutcome::Reused(session) => { auth_slot.replace(RequestAuthContext::invalid()); - auth_tokens.mark_refresh_token_replay(secret_hash, now); - if let Err(err) = auth_tokens.delete_chain(old.chain_id).await { - warn!(error = %err, chain_id = %old.chain_id, "Failed to revoke replayed refresh token chain"); + if let Err(err) = auth_sessions.delete_session(session.id).await { + warn!(error = %err, session_id = %session.id, "Failed to revoke replayed refresh token chain"); } - log_refresh_token_replay(old.chain_id, old.identity.subject(), &next_user_agent); + log_refresh_token_replay(session.id, session.identity.subject(), &next_user_agent); return oauth_error( StatusCode::UNAUTHORIZED, "refresh_token_revoked", "Refresh token revoked", ); } - ConsumeOutcome::Rotated(old, new_row) => (old, *new_row), + RotateOutcome::Rotated(session) => session, }; - if !login_allowed(state.as_ref(), &old.login) { + if !login_allowed(state.as_ref(), &session.login) { auth_slot.replace(RequestAuthContext::invalid()); - if let Err(err) = auth_tokens.delete_chain(old.chain_id).await { - warn!(error = %err, chain_id = %old.chain_id, "Failed to revoke deauthorized refresh token chain"); + if let Err(err) = auth_sessions.delete_session(session.id).await { + warn!(error = %err, session_id = %session.id, "Failed to revoke deauthorized refresh token chain"); } return oauth_error(StatusCode::FORBIDDEN, "unauthorized", "Login not permitted"); } @@ -661,24 +634,29 @@ async fn refresh( jwt_key, jwt_issuer, &JwtSubject { - identity: old.identity.clone(), - login: old.login.clone(), - name: old.name.clone(), - email: old.email.clone(), - avatar_url: old.avatar_url.clone(), + identity: session.identity.clone(), + login: session.login.clone(), + name: session.name.clone(), + email: session.email.clone(), + avatar_url: session.avatar_url.clone(), user_url: String::new(), auth_method: AuthMethod::Github, }, chrono::Duration::minutes(ACCESS_TOKEN_TTL_MINUTES), ); - auth_slot.replace(refresh_user_context(&old)); + auth_slot.replace(refresh_user_context(&session)); Json(CliTokenResponse { access_token, access_token_expires_at: access_expires_at, refresh_token: format!("{REFRESH_TOKEN_PREFIX}{next_secret}"), - refresh_token_expires_at: new_row.expires_at, - subject: subject_response(&old.identity, &old.login, &old.name, &old.email), + refresh_token_expires_at: refresh_expires_at, + subject: subject_response( + &session.identity, + &session.login, + &session.name, + &session.email, + ), }) .into_response() } @@ -701,20 +679,9 @@ async fn logout( } RefreshCredential::Present(secret) => secret, }; - let auth_tokens = match state.store_ref().refresh_tokens().await { - Ok(store) => store, - Err(err) => { - warn!(error = %err, "Failed to open refresh token store"); - return oauth_error( - StatusCode::INTERNAL_SERVER_ERROR, - "server_error", - "Could not complete logout", - ); - } - }; - - let existing = match auth_tokens - .find_refresh_token(&hash_refresh_secret(&secret)) + let auth_sessions = &state.stores.auth_sessions; + let existing = match auth_sessions + .find_session_by_token_hash(&hash_refresh_secret(&secret)) .await { Ok(existing) => existing, @@ -728,17 +695,17 @@ async fn logout( } }; - if let Some(refresh_token) = existing { - auth_slot.replace(refresh_user_context(&refresh_token)); - if let Err(err) = auth_tokens.delete_chain(refresh_token.chain_id).await { - warn!(error = %err, chain_id = %refresh_token.chain_id, "Failed to revoke refresh token chain during logout"); + if let Some(session) = existing { + auth_slot.replace(refresh_user_context(&session)); + if let Err(err) = auth_sessions.delete_session(session.id).await { + warn!(error = %err, session_id = %session.id, "Failed to revoke refresh token chain during logout"); return oauth_error( StatusCode::INTERNAL_SERVER_ERROR, "server_error", "Could not complete logout", ); } - log_cli_refresh_chain_logged_out(&refresh_token.login, &refresh_token.email); + log_cli_refresh_chain_logged_out(&session.login, &session.email); } else { auth_slot.replace(RequestAuthContext::invalid()); } @@ -956,13 +923,13 @@ fn refresh_credential_from_headers(headers: &HeaderMap) -> RefreshCredential { } } -fn refresh_user_context(refresh_token: &RefreshToken) -> RequestAuthContext { +fn refresh_user_context(session: &AuthSessionRecord) -> RequestAuthContext { RequestAuthContext::authenticated( Principal::user_with_avatar( - refresh_token.identity.clone(), - refresh_token.login.clone(), + session.identity.clone(), + session.login.clone(), AuthMethod::Github, - non_empty_avatar_url(&refresh_token.avatar_url), + non_empty_avatar_url(&session.avatar_url), ), None, ) @@ -972,31 +939,6 @@ fn hash_refresh_secret(secret: &str) -> [u8; 32] { Sha256::digest(secret.as_bytes()).into() } -fn next_refresh_row( - existing: Option<&RefreshToken>, - next_secret: &str, - user_agent: &str, - now: chrono::DateTime, -) -> RefreshToken { - let fallback_identity = fabro_types::IdpIdentity::new("https://github.com", "0") - .expect("static identity should be valid"); - RefreshToken { - token_hash: hash_refresh_secret(next_secret), - chain_id: existing.map_or_else(uuid::Uuid::new_v4, |token| token.chain_id), - identity: existing - .map_or_else(|| fallback_identity.clone(), |token| token.identity.clone()), - login: existing.map_or_else(String::new, |token| token.login.clone()), - name: existing.map_or_else(String::new, |token| token.name.clone()), - email: existing.map_or_else(String::new, |token| token.email.clone()), - avatar_url: existing.map_or_else(String::new, |token| token.avatar_url.clone()), - issued_at: now, - expires_at: now + chrono::Duration::days(REFRESH_TOKEN_TTL_DAYS), - last_used_at: now, - used: false, - user_agent: user_agent.to_string(), - } -} - fn user_agent_fingerprint(user_agent: &str) -> String { let digest = Sha256::digest(user_agent.as_bytes()); hex::encode(&digest[..8]) @@ -1006,9 +948,9 @@ fn log_cli_auth_tokens_issued(login: &str, email: &str) { info!(login = %login, email = %email, "Issued CLI auth tokens"); } -fn log_refresh_token_replay(chain_id: uuid::Uuid, idp_subject: &str, user_agent: &str) { +fn log_refresh_token_replay(session_id: uuid::Uuid, idp_subject: &str, user_agent: &str) { warn!( - chain_id = %chain_id, + session_id = %session_id, idp_subject = %idp_subject, user_agent_fingerprint = %user_agent_fingerprint(user_agent), "Refresh token replay detected" @@ -1248,9 +1190,10 @@ mod tests { CliFlowCookie, DEV_TOKEN_LOGIN_INSTRUCTIONS, add_cli_flow_cookie, read_private_cli_flow, user_agent_fingerprint, web_routes, }; - use crate::auth::{self, AuthCode, AuthErrorCode, RefreshToken}; + use crate::auth::{self, AuthCode, AuthErrorCode, AuthSessionRecord, RefreshToken}; use crate::jwt_auth::{AuthMode, ConfiguredAuth}; use crate::principal_middleware::{AuthStatus, RequestAuthContext}; + use crate::server::AppState; use crate::web_auth::SessionCookie; fn test_cookie_key() -> Key { @@ -1413,23 +1356,39 @@ client_id = "github-client-id" Sha256::digest(secret.as_bytes()).into() } - fn refresh_row(secret: &str) -> RefreshToken { + fn session_and_token(secret: &str) -> (AuthSessionRecord, RefreshToken) { let now = chrono::Utc::now(); - RefreshToken { - token_hash: hash_refresh_secret(secret), - chain_id: Uuid::new_v4(), + let session = AuthSessionRecord { + id: Uuid::new_v4(), identity: fabro_types::IdpIdentity::new("https://github.com", "12345") .expect("identity should be valid"), login: "octocat".to_string(), name: "The Octocat".to_string(), email: "octocat@example.com".to_string(), avatar_url: "https://example.com/octocat.png".to_string(), - issued_at: now, - expires_at: now + chrono::Duration::days(30), - last_used_at: now, - used: false, user_agent: "fabro-test".to_string(), - } + created_at: now, + last_used_at: now, + }; + let token = RefreshToken { + token_hash: hash_refresh_secret(secret), + session_id: session.id, + issued_at: now, + expires_at: now + chrono::Duration::days(30), + used_at: None, + }; + (session, token) + } + + async fn open_cli_session(state: &AppState, secret: &str) -> Uuid { + let (session, token) = session_and_token(secret); + state + .stores + .auth_sessions + .create_session(&session, &token) + .await + .unwrap(); + session.id } #[derive(Default)] @@ -2014,9 +1973,9 @@ client_id = "github-client-id" .unwrap() .strip_prefix("fabro_refresh_") .unwrap(); - let auth_tokens = state.store_ref().refresh_tokens().await.unwrap(); - let refresh = auth_tokens - .find_refresh_token(&hash_refresh_secret(refresh_secret)) + let auth_sessions = &state.stores.auth_sessions; + let refresh = auth_sessions + .find_session_by_token_hash(&hash_refresh_secret(refresh_secret)) .await .unwrap() .expect("refresh token should be stored"); @@ -2168,11 +2127,8 @@ client_id = "github-client-id" async fn refresh_rotates_tokens_and_replay_revokes_chain() { let (app, state) = test_router(github_settings("https://fabro.example")); let initial_secret = "refresh-secret-1"; - let auth_tokens = state.store_ref().refresh_tokens().await.unwrap(); - auth_tokens - .insert_refresh_token(refresh_row(initial_secret)) - .await - .unwrap(); + open_cli_session(&state, initial_secret).await; + let auth_sessions = &state.stores.auth_sessions; let refresh_request = || { Request::builder() @@ -2211,15 +2167,15 @@ client_id = "github-client-id" let new_secret = rotated.strip_prefix("fabro_refresh_").unwrap(); assert!( - auth_tokens - .find_refresh_token(&hash_refresh_secret(initial_secret)) + auth_sessions + .find_session_by_token_hash(&hash_refresh_secret(initial_secret)) .await .unwrap() .is_none() ); assert!( - auth_tokens - .find_refresh_token(&hash_refresh_secret(new_secret)) + auth_sessions + .find_session_by_token_hash(&hash_refresh_secret(new_secret)) .await .unwrap() .is_none() @@ -2233,14 +2189,7 @@ client_id = "github-client-id" github_auth_mode(), ); let initial_secret = "refresh-secret-auth-context"; - state - .store_ref() - .refresh_tokens() - .await - .unwrap() - .insert_refresh_token(refresh_row(initial_secret)) - .await - .unwrap(); + open_cli_session(&state, initial_secret).await; let response = app .clone() @@ -2296,11 +2245,8 @@ client_id = "github-client-id" async fn concurrent_refresh_has_one_winner_and_revokes_chain() { let (app, state) = test_router(github_settings("https://fabro.example")); let initial_secret = "refresh-secret-concurrent"; - let auth_tokens = state.store_ref().refresh_tokens().await.unwrap(); - auth_tokens - .insert_refresh_token(refresh_row(initial_secret)) - .await - .unwrap(); + open_cli_session(&state, initial_secret).await; + let auth_sessions = &state.stores.auth_sessions; let barrier = Arc::new(Barrier::new(33)); let mut tasks = JoinSet::new(); @@ -2348,7 +2294,11 @@ client_id = "github-client-id" .map(str::to_string); } StatusCode::UNAUTHORIZED => { - assert_eq!(body["error"], "refresh_token_revoked"); + let error = body["error"].as_str().unwrap_or_default(); + assert!( + error == "refresh_token_revoked" || error == "refresh_token_expired", + "unexpected refresh error {error}" + ); revoked += 1; } other => panic!("unexpected refresh status {other}: {body}"), @@ -2359,15 +2309,15 @@ client_id = "github-client-id" assert_eq!(revoked, 31); let rotated_secret = rotated_secret.expect("one refresh should rotate the token"); assert!( - auth_tokens - .find_refresh_token(&hash_refresh_secret(initial_secret)) + auth_sessions + .find_session_by_token_hash(&hash_refresh_secret(initial_secret)) .await .unwrap() .is_none() ); assert!( - auth_tokens - .find_refresh_token(&hash_refresh_secret(&rotated_secret)) + auth_sessions + .find_session_by_token_hash(&hash_refresh_secret(&rotated_secret)) .await .unwrap() .is_none() @@ -2378,17 +2328,24 @@ client_id = "github-client-id" async fn logout_deletes_refresh_token_chain_and_returns_no_content() { let (app, state) = test_router(github_settings("https://fabro.example")); let secret = "refresh-secret-logout"; - let token = refresh_row(secret); - let chain_id = token.chain_id; - let auth_tokens = state.store_ref().refresh_tokens().await.unwrap(); - auth_tokens.insert_refresh_token(token).await.unwrap(); - - let sibling = RefreshToken { - token_hash: hash_refresh_secret("refresh-secret-logout-2"), - chain_id, - ..refresh_row("refresh-secret-logout-2") - }; - auth_tokens.insert_refresh_token(sibling).await.unwrap(); + let (session, token) = session_and_token(secret); + let auth_sessions = &state.stores.auth_sessions; + auth_sessions + .create_session(&session, &token) + .await + .unwrap(); + + let now = chrono::Utc::now(); + auth_sessions + .rotate( + &hash_refresh_secret(secret), + &hash_refresh_secret("refresh-secret-logout-2"), + now + chrono::Duration::days(30), + "fabro-test", + now, + ) + .await + .unwrap(); let response = app .oneshot( @@ -2407,15 +2364,15 @@ client_id = "github-client-id" assert_eq!(response.status(), StatusCode::NO_CONTENT); assert!( - auth_tokens - .find_refresh_token(&hash_refresh_secret(secret)) + auth_sessions + .find_session_by_token_hash(&hash_refresh_secret(secret)) .await .unwrap() .is_none() ); assert!( - auth_tokens - .find_refresh_token(&hash_refresh_secret("refresh-secret-logout-2")) + auth_sessions + .find_session_by_token_hash(&hash_refresh_secret("refresh-secret-logout-2")) .await .unwrap() .is_none() diff --git a/lib/apps/fabro-server/src/auth/mod.rs b/lib/apps/fabro-server/src/auth/mod.rs index 3b6e0c47d6..da4823f5ca 100644 --- a/lib/apps/fabro-server/src/auth/mod.rs +++ b/lib/apps/fabro-server/src/auth/mod.rs @@ -30,7 +30,8 @@ pub(crate) const REFRESH_TOKEN_PREFIX: &str = "fabro_refresh_"; pub(crate) use browser_shell::browser_shell; pub(crate) use cli_flow::web_routes; -pub(crate) use fabro_store::{AuthCode, ConsumeOutcome, RefreshToken}; +pub(crate) use fabro_store::AuthCode; +pub(crate) use fabro_store::auth_session_store::{AuthSessionRecord, RefreshToken, RotateOutcome}; pub use github_endpoints::GithubEndpoints; pub(crate) use jwt::{JwtError, JwtSubject, issue, verify}; pub(crate) use keys::{ diff --git a/lib/apps/fabro-server/src/serve.rs b/lib/apps/fabro-server/src/serve.rs index f46bb7e0a2..eef61311bb 100644 --- a/lib/apps/fabro-server/src/serve.rs +++ b/lib/apps/fabro-server/src/serve.rs @@ -780,7 +780,14 @@ where cache_path, )); let auth_code_store = store.auth_codes().await?; - let auth_token_store = store.refresh_tokens().await?; + // Refresh tokens now live in SQLite. Nothing reads the old records and no + // reaper collects them any more, so clear them out once rather than + // leaving them in the object store forever. + match store.retire_refresh_token_keyspace().await { + Ok(0) => {} + Ok(removed) => info!(removed, "Removed retired SlateDB refresh token records"), + Err(err) => warn!(error = %err, "Failed to remove retired SlateDB refresh token records"), + } let (artifact_object_store, artifact_prefix) = build_artifact_object_store_with_server_secrets( &resolved_server_settings, &server_secrets, @@ -845,7 +852,7 @@ where spawn_auth_store_reapers( Arc::clone(&auth_code_store), - Arc::clone(&auth_token_store), + Arc::clone(&state.stores.auth_sessions), shutdown.clone(), ); @@ -1109,11 +1116,11 @@ async fn shutdown_signal() { fn spawn_auth_store_reapers( auth_codes: Arc, - auth_tokens: Arc, + auth_sessions: Arc, shutdown: CancellationToken, ) { spawn_auth_code_reaper(auth_codes, shutdown.clone()); - spawn_refresh_token_reaper(auth_tokens, shutdown); + spawn_refresh_token_reaper(auth_sessions, shutdown); } fn spawn_auth_code_reaper( @@ -1138,7 +1145,7 @@ fn spawn_auth_code_reaper( } fn spawn_refresh_token_reaper( - auth_tokens: Arc, + auth_sessions: Arc, shutdown: CancellationToken, ) { tokio::spawn(async move { @@ -1150,7 +1157,7 @@ fn spawn_refresh_token_reaper( () = shutdown.cancelled() => break, _ = interval.tick() => { let cutoff = chrono::Utc::now() - chrono::Duration::days(7); - if let Err(err) = auth_tokens.gc_expired(cutoff).await { + if let Err(err) = auth_sessions.gc_expired(cutoff).await { warn!(error = %err, "Failed to garbage collect expired refresh tokens"); } } diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 11a9a09487..9c60f6db6f 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -85,8 +85,9 @@ use fabro_slack::threads::ThreadRegistry; use fabro_slack::{blocks as slack_blocks, connection as slack_connection}; use fabro_static::EnvVars; use fabro_store::{ - ArtifactKey, ArtifactStore, CachedRunProjection, Database, EventEnvelope, EventPayload, - NodeArtifact, PendingInterviewRecord, RunSummaryStore, StageArtifactEntry, StageId, + ArtifactKey, ArtifactStore, AuthSessionStore, CachedRunProjection, Database, EventEnvelope, + EventPayload, NodeArtifact, PendingInterviewRecord, RunSummaryStore, StageArtifactEntry, + StageId, }; #[cfg(test)] use fabro_types::BlockedReason; @@ -1153,6 +1154,7 @@ pub struct AppState { pub(crate) struct AppStores { pub(crate) runs: Arc, pub(crate) run_summaries: Arc, + pub(crate) auth_sessions: Arc, pub(crate) automations: Arc, pub(crate) environments: Arc, pub(crate) mcp_servers: Arc, @@ -1162,6 +1164,16 @@ pub(crate) struct AppStores { type PullRequestCreateLocks = Arc>>>>; +#[cfg(any(test, feature = "test-support"))] +impl AppState { + /// Access the auth session store so tests can seed CLI sessions against + /// the same SQLite pool the router reads from. + #[must_use] + pub fn test_auth_session_store(&self) -> &Arc { + &self.stores.auth_sessions + } +} + impl AppState { pub(crate) fn automation_store(&self) -> &AutomationStore { &self.stores.automations @@ -2434,6 +2446,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result anyhow::Result store, - Err(err) => { - error!(error = %err, "Failed to open refresh token store while listing auth sessions"); - return ApiError::new( - StatusCode::INTERNAL_SERVER_ERROR, - "Failed to list auth sessions.", - ) - .into_response(); - } - }; - let cli_sessions = match auth_tokens + let cli_sessions = match state + .stores + .auth_sessions .active_cli_sessions(&authenticated.principal.identity, now) .await { - Ok(tokens) => tokens, + Ok(sessions) => sessions, Err(err) => { - error!(error = %err, "Failed to scan refresh tokens while listing auth sessions"); + error!(error = %err, "Failed to load auth sessions"); return ApiError::new( StatusCode::INTERNAL_SERVER_ERROR, "Failed to list auth sessions.", @@ -923,17 +914,17 @@ async fn list_auth_sessions( } }; - sessions.extend(cli_sessions.into_iter().map(|token| AuthSession { - id: format!("cli:{}", token.chain_id), + sessions.extend(cli_sessions.into_iter().map(|active| AuthSession { + id: format!("cli:{}", active.session.id), kind: "cli", current: false, provider: "github".to_string(), - login: token.login, + login: active.session.login, label: "Fabro CLI".to_string(), - user_agent: Some(token.user_agent), - created_at: token.issued_at, - last_seen_at: token.last_used_at, - expires_at: token.expires_at, + user_agent: Some(active.session.user_agent), + created_at: active.session.created_at, + last_seen_at: active.session.last_used_at, + expires_at: active.expires_at, revocable: true, })); sessions.sort_by(|left, right| { @@ -961,31 +952,26 @@ async fn delete_auth_session( .into_response(); } - let Some(raw_chain_id) = id.strip_prefix("cli:") else { + let Some(raw_session_id) = id.strip_prefix("cli:") else { return ApiError::not_found("Auth session not found.").into_response(); }; - let Ok(chain_id) = uuid::Uuid::parse_str(raw_chain_id) else { + let Ok(session_id) = uuid::Uuid::parse_str(raw_session_id) else { return ApiError::bad_request("Malformed CLI auth session id.").into_response(); }; - let auth_tokens = match state.store_ref().refresh_tokens().await { - Ok(store) => store, - Err(err) => { - error!(error = %err, "Failed to open refresh token store while deleting auth session"); - return ApiError::new( - StatusCode::INTERNAL_SERVER_ERROR, - "Failed to revoke auth session.", - ) - .into_response(); - } - }; - let deleted = match auth_tokens - .delete_active_chain_for_identity(&authenticated.principal.identity, chain_id, Utc::now()) + let deleted = match state + .stores + .auth_sessions + .delete_active_session_for_identity( + &authenticated.principal.identity, + session_id, + Utc::now(), + ) .await { Ok(deleted) => deleted, Err(err) => { - error!(error = %err, "Failed to scan refresh tokens while deleting auth session"); + error!(error = %err, "Failed to revoke auth session"); return ApiError::new( StatusCode::INTERNAL_SERVER_ERROR, "Failed to revoke auth session.", diff --git a/lib/apps/fabro-server/tests/it/api/auth_sessions.rs b/lib/apps/fabro-server/tests/it/api/auth_sessions.rs index 0ba7d71d89..f5b08f6628 100644 --- a/lib/apps/fabro-server/tests/it/api/auth_sessions.rs +++ b/lib/apps/fabro-server/tests/it/api/auth_sessions.rs @@ -6,10 +6,11 @@ use axum::body::Body; use axum::http::{Request, StatusCode, header}; use cookie::{Cookie, CookieJar, Key}; use fabro_server::jwt_auth::resolve_auth_mode_with_lookup; -use fabro_server::server::{RouterOptions, build_router_with_options}; +use fabro_server::server::{AppState, RouterOptions, build_router_with_options}; use fabro_server::test_support::{TEST_SESSION_SECRET, TestAppStateBuilder}; use fabro_server::web_auth::{SESSION_COOKIE_NAME, SessionCookie}; -use fabro_store::{ArtifactStore, Database, RefreshToken}; +use fabro_store::auth_session_store::{AuthSessionRecord, RefreshToken}; +use fabro_store::{ArtifactStore, Database}; use hkdf::Hkdf; use object_store::memory::InMemory; use sha2::Sha256; @@ -18,7 +19,7 @@ use uuid::Uuid; use crate::helpers::{response_json, response_status, settings_from_toml}; -fn test_app(source: &str) -> (axum::Router, Arc) { +fn test_app(source: &str) -> (axum::Router, Arc) { let settings = settings_from_toml(source); let object_store: Arc = Arc::new(InMemory::new()); let store = Arc::new(Database::new( @@ -44,11 +45,11 @@ fn test_app(source: &str) -> (axum::Router, Arc) { TEST_SESSION_SECRET.to_string(), )])) .build(); - let app = build_router_with_options(state, &auth_mode, RouterOptions::default()); - (app, store) + let app = build_router_with_options(Arc::clone(&state), &auth_mode, RouterOptions::default()); + (app, state) } -fn github_app() -> (axum::Router, Arc) { +fn github_app() -> (axum::Router, Arc) { test_app( r#" _version = 1 @@ -118,24 +119,40 @@ fn session_cookie() -> String { .to_string() } -fn refresh_token(hash: [u8; 32], chain_id: Uuid) -> RefreshToken { +fn cli_session(id: Uuid, identity: fabro_types::IdpIdentity) -> AuthSessionRecord { let now = chrono::Utc::now(); - RefreshToken { - token_hash: hash, - chain_id, - identity: github_identity(), + AuthSessionRecord { + id, + identity, login: "octocat".to_string(), name: "The Octocat".to_string(), email: "octocat@example.com".to_string(), avatar_url: String::new(), + user_agent: "fabro-cli/it".to_string(), + created_at: now - chrono::Duration::days(1), + last_used_at: now, + } +} + +fn refresh_token(hash: [u8; 32], session_id: Uuid) -> RefreshToken { + let now = chrono::Utc::now(); + RefreshToken { + token_hash: hash, + session_id, issued_at: now - chrono::Duration::days(1), expires_at: now + chrono::Duration::days(30), - last_used_at: now, - used: false, - user_agent: "fabro-cli/it".to_string(), + used_at: None, } } +async fn seed_session(state: &AppState, session: AuthSessionRecord, token: RefreshToken) { + state + .test_auth_session_store() + .create_session(&session, &token) + .await + .expect("CLI session should insert"); +} + async fn get_sessions(app: axum::Router, cookie: &str) -> serde_json::Value { response_json( app.oneshot( @@ -174,16 +191,14 @@ async fn authenticated_browser_requests_receive_current_browser_session() { #[tokio::test] async fn active_cli_refresh_token_chains_for_identity_appear_in_unified_list() { - let (app, store) = github_app(); - let auth_tokens = store - .refresh_tokens() - .await - .expect("refresh token store should open"); - let chain_id = Uuid::new_v4(); - auth_tokens - .insert_refresh_token(refresh_token([1_u8; 32], chain_id)) - .await - .expect("refresh token should insert"); + let (app, state) = github_app(); + let session_id = Uuid::new_v4(); + seed_session( + &state, + cli_session(session_id, github_identity()), + refresh_token([1_u8; 32], session_id), + ) + .await; let body = get_sessions(app, &session_cookie()).await; let sessions = body["sessions"] @@ -194,7 +209,7 @@ async fn active_cli_refresh_token_chains_for_identity_appear_in_unified_list() { assert_eq!(sessions[0]["id"], "browser:current"); let cli = sessions .iter() - .find(|session| session["id"] == format!("cli:{chain_id}")) + .find(|session| session["id"] == format!("cli:{session_id}")) .expect("CLI session should be present"); assert_eq!(cli["kind"], "cli"); assert_eq!(cli["current"], false); @@ -207,26 +222,37 @@ async fn active_cli_refresh_token_chains_for_identity_appear_in_unified_list() { #[tokio::test] async fn inactive_and_other_identity_cli_tokens_are_excluded() { - let (app, store) = github_app(); - let auth_tokens = store - .refresh_tokens() - .await - .expect("refresh token store should open"); - let active_chain_id = Uuid::new_v4(); + let (app, state) = github_app(); + let active_session_id = Uuid::new_v4(); let now = chrono::Utc::now(); - let active = refresh_token([1_u8; 32], active_chain_id); - let mut expired = refresh_token([2_u8; 32], Uuid::new_v4()); + + let expired_id = Uuid::new_v4(); + let mut expired = refresh_token([2_u8; 32], expired_id); expired.expires_at = now - chrono::Duration::seconds(1); - let mut used = refresh_token([3_u8; 32], Uuid::new_v4()); - used.used = true; - let mut other = refresh_token([4_u8; 32], Uuid::new_v4()); - other.identity = other_identity(); - - for token in [active, expired, used, other] { - auth_tokens - .insert_refresh_token(token) - .await - .expect("refresh token should insert"); + + // A chain whose only token has already been rotated away has nothing left + // to spend, so it is inactive even though the token has not expired. + let used_id = Uuid::new_v4(); + let mut used = refresh_token([3_u8; 32], used_id); + used.used_at = Some(now); + + for (session, token) in [ + ( + cli_session(active_session_id, github_identity()), + refresh_token([1_u8; 32], active_session_id), + ), + (cli_session(expired_id, github_identity()), expired), + (cli_session(used_id, github_identity()), used), + ( + cli_session(Uuid::new_v4(), other_identity()), + refresh_token([4_u8; 32], Uuid::new_v4()), + ), + ] { + let token = RefreshToken { + session_id: session.id, + ..token + }; + seed_session(&state, session, token).await; } let body = get_sessions(app, &session_cookie()).await; @@ -244,35 +270,39 @@ async fn inactive_and_other_identity_cli_tokens_are_excluded() { assert_eq!(session_ids, vec![ "browser:current".to_string(), - format!("cli:{active_chain_id}") + format!("cli:{active_session_id}") ]); } #[tokio::test] async fn deleting_cli_session_removes_refresh_token_chain() { - let (app, store) = github_app(); - let auth_tokens = store - .refresh_tokens() - .await - .expect("refresh token store should open"); - let chain_id = Uuid::new_v4(); - let active = refresh_token([1_u8; 32], chain_id); - let mut used = refresh_token([2_u8; 32], chain_id); - used.used = true; - auth_tokens - .insert_refresh_token(active) - .await - .expect("active refresh token should insert"); - auth_tokens - .insert_refresh_token(used) + let (app, state) = github_app(); + let session_id = Uuid::new_v4(); + seed_session( + &state, + cli_session(session_id, github_identity()), + refresh_token([1_u8; 32], session_id), + ) + .await; + // Rotate once so the chain holds a spent token alongside its live one. + let now = chrono::Utc::now(); + state + .test_auth_session_store() + .rotate( + &[1_u8; 32], + &[2_u8; 32], + now + chrono::Duration::days(30), + "fabro-cli/it", + now, + ) .await - .expect("used refresh token should insert"); + .expect("rotation should succeed"); response_status( app.oneshot( Request::builder() .method("DELETE") - .uri(format!("/api/v1/auth/sessions/cli:{chain_id}")) + .uri(format!("/api/v1/auth/sessions/cli:{session_id}")) .header(header::COOKIE, session_cookie()) .body(Body::empty()) .expect("DELETE CLI auth session request should build"), @@ -285,15 +315,17 @@ async fn deleting_cli_session_removes_refresh_token_chain() { .await; assert!( - auth_tokens - .find_refresh_token(&[1_u8; 32]) + state + .test_auth_session_store() + .find_session_by_token_hash(&[1_u8; 32]) .await .expect("active token lookup should succeed") .is_none() ); assert!( - auth_tokens - .find_refresh_token(&[2_u8; 32]) + state + .test_auth_session_store() + .find_session_by_token_hash(&[2_u8; 32]) .await .expect("used token lookup should succeed") .is_none() diff --git a/lib/apps/fabro-server/tests/it/api/cli_auth_token.rs b/lib/apps/fabro-server/tests/it/api/cli_auth_token.rs index 2d0b743cbb..b985f85b52 100644 --- a/lib/apps/fabro-server/tests/it/api/cli_auth_token.rs +++ b/lib/apps/fabro-server/tests/it/api/cli_auth_token.rs @@ -5,9 +5,10 @@ use axum::body::Body; use axum::http::{Request, StatusCode, header}; use base64::Engine; use fabro_server::jwt_auth::resolve_auth_mode_with_lookup; -use fabro_server::server::{RouterOptions, build_router_with_options}; +use fabro_server::server::{AppState, RouterOptions, build_router_with_options}; use fabro_server::test_support::test_app_state_with_store_and_runtime_settings; -use fabro_store::{ArtifactStore, AuthCode, Database, RefreshToken}; +use fabro_store::auth_session_store::{AuthSessionRecord, RefreshToken}; +use fabro_store::{ArtifactStore, AuthCode, Database}; use object_store::memory::InMemory; use sha2::{Digest, Sha256}; use tower::ServiceExt; @@ -15,7 +16,7 @@ use uuid::Uuid; use crate::helpers::{body_json, settings_from_toml}; -fn test_app(source: &str) -> (axum::Router, Arc) { +fn test_app(source: &str) -> (axum::Router, Arc, Arc) { let settings = settings_from_toml(source); let object_store: Arc = Arc::new(InMemory::new()); let store = Arc::new(Database::new( @@ -32,18 +33,15 @@ fn test_app(source: &str) -> (axum::Router, Arc) { _ => None, }) .expect("auth mode should resolve"); - let app = build_router_with_options( - test_app_state_with_store_and_runtime_settings( - settings.server_settings, - settings.manifest_run_defaults, - 5, - Arc::clone(&store), - artifact_store, - ), - &auth_mode, - RouterOptions::default(), + let state = test_app_state_with_store_and_runtime_settings( + settings.server_settings, + settings.manifest_run_defaults, + 5, + Arc::clone(&store), + artifact_store, ); - (app, store) + let app = build_router_with_options(Arc::clone(&state), &auth_mode, RouterOptions::default()); + (app, store, state) } fn pkce_challenge(verifier: &str) -> String { @@ -56,7 +54,7 @@ fn hash_refresh_secret(secret: &str) -> [u8; 32] { #[tokio::test] async fn cli_auth_token_exchanges_code_over_public_router() { - let (app, store) = test_app( + let (app, store, _state) = test_app( r#" _version = 1 @@ -124,7 +122,7 @@ client_id = "Iv1.test" #[tokio::test] async fn cli_auth_refresh_replay_revokes_chain_over_public_router() { - let (app, store) = test_app( + let (app, _store, state) = test_app( r#" _version = 1 @@ -141,22 +139,26 @@ url = "https://fabro.example" client_id = "Iv1.test" "#, ); - let auth_tokens = store.refresh_tokens().await.unwrap(); let now = chrono::Utc::now(); - auth_tokens - .insert_refresh_token(RefreshToken { - token_hash: hash_refresh_secret("integration-refresh"), - chain_id: Uuid::new_v4(), - identity: fabro_types::IdpIdentity::new("https://github.com", "12345").unwrap(), - login: "octocat".to_string(), - name: "The Octocat".to_string(), - email: "octocat@example.com".to_string(), - avatar_url: String::new(), - issued_at: now, - expires_at: now + chrono::Duration::days(30), - last_used_at: now, - used: false, - user_agent: "fabro-cli/it".to_string(), + let session = AuthSessionRecord { + id: Uuid::new_v4(), + identity: fabro_types::IdpIdentity::new("https://github.com", "12345").unwrap(), + login: "octocat".to_string(), + name: "The Octocat".to_string(), + email: "octocat@example.com".to_string(), + avatar_url: String::new(), + user_agent: "fabro-cli/it".to_string(), + created_at: now, + last_used_at: now, + }; + state + .test_auth_session_store() + .create_session(&session, &RefreshToken { + token_hash: hash_refresh_secret("integration-refresh"), + session_id: session.id, + issued_at: now, + expires_at: now + chrono::Duration::days(30), + used_at: None, }) .await .unwrap(); diff --git a/lib/components/fabro-automation/src/model.rs b/lib/components/fabro-automation/src/model.rs index 90be56d8bf..52f0eceaa8 100644 --- a/lib/components/fabro-automation/src/model.rs +++ b/lib/components/fabro-automation/src/model.rs @@ -80,19 +80,6 @@ impl Automation { Ok(Self::from_validated_replace(id, revision, value)) } - pub(crate) fn to_persisted(&self) -> PersistedAutomation { - PersistedAutomation { - name: self.name.clone(), - description: self.description.clone(), - target: self.target.clone(), - triggers: self.triggers.clone(), - } - } - - pub fn to_toml_string(&self) -> Result { - toml::to_string_pretty(&self.to_persisted()).map_err(AutomationStoreError::from) - } - /// Returns the enabled API trigger if the automation has one. /// Returns `None` when the automation has no enabled API trigger. #[must_use] @@ -563,7 +550,8 @@ expression = "0 0 * * *" assert_eq!(automation.description, None); assert!(automation.triggers.iter().all(AutomationTrigger::enabled)); - let toml = automation.to_toml_string().unwrap(); + let persisted = super::parse_persisted(bytes, None).unwrap(); + let toml = String::from_utf8(super::canonical_bytes(&persisted).unwrap()).unwrap(); assert!(!top_level_lines(&toml).any(|line| line.starts_with("id = "))); assert!(!top_level_lines(&toml).any(|line| line.starts_with("revision = "))); assert!(!top_level_lines(&toml).any(|line| line.starts_with("enabled = "))); diff --git a/lib/components/fabro-store/src/auth_session_store.rs b/lib/components/fabro-store/src/auth_session_store.rs new file mode 100644 index 0000000000..45c727dd3e --- /dev/null +++ b/lib/components/fabro-store/src/auth_session_store.rs @@ -0,0 +1,801 @@ +//! SQLite-backed storage for CLI auth sessions and their refresh tokens. +//! +//! A session is a rotation chain. The chain owns the identity and profile; +//! each token in it owns only its own lifetime. Splitting them that way is +//! what makes every operation here an indexed query rather than a scan over +//! every token ever issued. + +use chrono::{DateTime, Utc}; +use fabro_types::IdpIdentity; +use sqlx::sqlite::SqliteRow; +use sqlx::{Row as _, SqlitePool}; +use uuid::Uuid; + +use crate::{Error, Result}; + +/// A CLI auth session: one rotation chain, owned by one identity. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct AuthSessionRecord { + pub id: Uuid, + pub identity: IdpIdentity, + pub login: String, + pub name: String, + pub email: String, + pub avatar_url: String, + pub user_agent: String, + pub created_at: DateTime, + pub last_used_at: DateTime, +} + +/// One refresh token within a session. `used_at` is set when the token is +/// rotated away; the row is kept until expiry so a replay stays recognisable. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct RefreshToken { + pub token_hash: [u8; 32], + pub session_id: Uuid, + pub issued_at: DateTime, + pub expires_at: DateTime, + pub used_at: Option>, +} + +/// A session with a spendable token, as returned by the session listing. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ActiveCliSession { + pub session: AuthSessionRecord, + /// Expiry of the session's live token. + pub expires_at: DateTime, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum RotateOutcome { + /// The presented token was spent and its successor issued. + Rotated(AuthSessionRecord), + /// The presented token had already been rotated away — a replay. + Reused(AuthSessionRecord), + Expired, + NotFound, +} + +pub struct AuthSessionStore { + pool: SqlitePool, +} + +impl std::fmt::Debug for AuthSessionStore { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("AuthSessionStore").finish_non_exhaustive() + } +} + +const SELECT_SESSION_BY_TOKEN_SQL: &str = r" +SELECT s.id, s.identity_issuer, s.identity_subject, s.login, s.name, s.email, + s.avatar_url, s.user_agent, s.created_at_ms, s.last_used_at_ms +FROM auth_sessions s +JOIN refresh_tokens t ON t.session_id = s.id +WHERE t.token_hash = ? +"; + +const SELECT_ACTIVE_SESSIONS_SQL: &str = r" +SELECT s.id, s.identity_issuer, s.identity_subject, s.login, s.name, s.email, + s.avatar_url, s.user_agent, s.created_at_ms, s.last_used_at_ms, + t.expires_at_ms +FROM auth_sessions s +JOIN refresh_tokens t ON t.session_id = s.id AND t.used_at_ms IS NULL +WHERE s.identity_issuer = ? AND s.identity_subject = ? AND t.expires_at_ms > ? +ORDER BY s.last_used_at_ms DESC +"; + +const SELECT_SESSION_BY_ID_SQL: &str = r" +SELECT s.id, s.identity_issuer, s.identity_subject, s.login, s.name, s.email, + s.avatar_url, s.user_agent, s.created_at_ms, s.last_used_at_ms +FROM auth_sessions s +WHERE s.id = ? +"; + +impl AuthSessionStore { + #[must_use] + pub fn new(pool: SqlitePool) -> Self { + Self { pool } + } + + /// Open a new session with its first refresh token. + pub async fn create_session( + &self, + session: &AuthSessionRecord, + token: &RefreshToken, + ) -> Result<()> { + let mut tx = self.pool.begin().await?; + sqlx::query( + r" +INSERT INTO auth_sessions ( + id, identity_issuer, identity_subject, login, name, email, avatar_url, + user_agent, created_at_ms, last_used_at_ms +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) +", + ) + .bind(session.id.to_string()) + .bind(session.identity.issuer()) + .bind(session.identity.subject()) + .bind(&session.login) + .bind(&session.name) + .bind(&session.email) + .bind(&session.avatar_url) + .bind(&session.user_agent) + .bind(session.created_at.timestamp_millis()) + .bind(session.last_used_at.timestamp_millis()) + .execute(&mut *tx) + .await?; + insert_token(&mut tx, token).await?; + tx.commit().await?; + Ok(()) + } + + /// Look up the session a token belongs to, spent or not. + pub async fn find_session_by_token_hash( + &self, + token_hash: &[u8; 32], + ) -> Result> { + sqlx::query(SELECT_SESSION_BY_TOKEN_SQL) + .bind(token_hash.as_slice()) + .fetch_optional(&self.pool) + .await? + .as_ref() + .map(session_from_row) + .transpose() + } + + /// Sessions belonging to `identity` that still hold a spendable token. + /// + /// The partial unique index guarantees at most one live token per + /// session, so this joins one row per session rather than grouping + /// candidates. + pub async fn active_cli_sessions( + &self, + identity: &IdpIdentity, + now: DateTime, + ) -> Result> { + sqlx::query(SELECT_ACTIVE_SESSIONS_SQL) + .bind(identity.issuer()) + .bind(identity.subject()) + .bind(now.timestamp_millis()) + .fetch_all(&self.pool) + .await? + .into_iter() + .map(|row| { + let expires_at = timestamp_from_row(&row, "expires_at_ms")?; + Ok(ActiveCliSession { + session: session_from_row(&row)?, + expires_at, + }) + }) + .collect() + } + + /// Spend `presented_hash` and issue `new_token_hash` in its place. + /// + /// The claiming UPDATE is the transaction's first statement, so SQLite + /// takes the write lock before anything is read. A concurrent caller + /// blocks on it, then observes `used_at_ms` already set and gets + /// [`RotateOutcome::Reused`] — the replay signal — with no application + /// mutex involved. + pub async fn rotate( + &self, + presented_hash: &[u8; 32], + new_token_hash: &[u8; 32], + new_expires_at: DateTime, + user_agent: &str, + now: DateTime, + ) -> Result { + let now_ms = now.timestamp_millis(); + let mut tx = self.pool.begin().await?; + + let claimed: Option = sqlx::query_scalar( + r" +UPDATE refresh_tokens SET used_at_ms = ? +WHERE token_hash = ? AND used_at_ms IS NULL AND expires_at_ms > ? +RETURNING session_id +", + ) + .bind(now_ms) + .bind(presented_hash.as_slice()) + .bind(now_ms) + .fetch_optional(&mut *tx) + .await?; + + let Some(session_id) = claimed else { + // Cold path: the claim failed, so read once more to say why. + let existing: Option<(Option, i64)> = sqlx::query_as( + "SELECT used_at_ms, expires_at_ms FROM refresh_tokens WHERE token_hash = ?", + ) + .bind(presented_hash.as_slice()) + .fetch_optional(&mut *tx) + .await?; + // Expiry is checked before reuse so an expired token reports as + // expired even if it had already been rotated away, matching the + // ordering callers rely on: only a live replay revokes the chain. + let outcome = match existing { + None => RotateOutcome::NotFound, + Some((used_at_ms, expires_at_ms)) => { + if expires_at_ms <= now_ms { + RotateOutcome::Expired + } else if used_at_ms.is_some() { + let session = load_session(&mut tx, presented_hash).await?; + session.map_or(RotateOutcome::NotFound, RotateOutcome::Reused) + } else { + // Unreachable: a live, unexpired token would have been + // claimed by the UPDATE above, in this transaction. + RotateOutcome::NotFound + } + } + }; + tx.commit().await?; + return Ok(outcome); + }; + + let session_id = parse_uuid(&session_id)?; + insert_token(&mut tx, &RefreshToken { + token_hash: *new_token_hash, + session_id, + issued_at: now, + expires_at: new_expires_at, + used_at: None, + }) + .await?; + sqlx::query("UPDATE auth_sessions SET last_used_at_ms = ?, user_agent = ? WHERE id = ?") + .bind(now_ms) + .bind(user_agent) + .bind(session_id.to_string()) + .execute(&mut *tx) + .await?; + + let row = sqlx::query(SELECT_SESSION_BY_ID_SQL) + .bind(session_id.to_string()) + .fetch_one(&mut *tx) + .await?; + let session = session_from_row(&row)?; + tx.commit().await?; + Ok(RotateOutcome::Rotated(session)) + } + + /// Revoke a session outright. Its tokens go with it via `ON DELETE + /// CASCADE`. + pub async fn delete_session(&self, session_id: Uuid) -> Result<()> { + sqlx::query("DELETE FROM auth_sessions WHERE id = ?") + .bind(session_id.to_string()) + .execute(&self.pool) + .await?; + Ok(()) + } + + /// Revoke a session on behalf of its owner, but only while it is still + /// usable. Returns the number of sessions deleted (0 or 1), so a caller + /// can distinguish "revoked" from "no such live session". + pub async fn delete_active_session_for_identity( + &self, + identity: &IdpIdentity, + session_id: Uuid, + now: DateTime, + ) -> Result { + let deleted = sqlx::query( + r" +DELETE FROM auth_sessions +WHERE id = ? AND identity_issuer = ? AND identity_subject = ? + AND EXISTS ( + SELECT 1 FROM refresh_tokens t + WHERE t.session_id = auth_sessions.id + AND t.used_at_ms IS NULL + AND t.expires_at_ms > ? + ) +", + ) + .bind(session_id.to_string()) + .bind(identity.issuer()) + .bind(identity.subject()) + .bind(now.timestamp_millis()) + .execute(&self.pool) + .await? + .rows_affected(); + Ok(deleted) + } + + /// Drop expired tokens, then any session left without one. Returns the + /// number of tokens removed. + pub async fn gc_expired(&self, cutoff: DateTime) -> Result { + let mut tx = self.pool.begin().await?; + let tokens = sqlx::query("DELETE FROM refresh_tokens WHERE expires_at_ms <= ?") + .bind(cutoff.timestamp_millis()) + .execute(&mut *tx) + .await? + .rows_affected(); + sqlx::query( + r" +DELETE FROM auth_sessions +WHERE NOT EXISTS (SELECT 1 FROM refresh_tokens t WHERE t.session_id = auth_sessions.id) +", + ) + .execute(&mut *tx) + .await?; + tx.commit().await?; + Ok(tokens) + } +} + +async fn insert_token( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + token: &RefreshToken, +) -> Result<()> { + sqlx::query( + r" +INSERT INTO refresh_tokens (token_hash, session_id, issued_at_ms, expires_at_ms, used_at_ms) +VALUES (?, ?, ?, ?, ?) +", + ) + .bind(token.token_hash.as_slice()) + .bind(token.session_id.to_string()) + .bind(token.issued_at.timestamp_millis()) + .bind(token.expires_at.timestamp_millis()) + .bind(token.used_at.map(|used_at| used_at.timestamp_millis())) + .execute(&mut **tx) + .await?; + Ok(()) +} + +async fn load_session( + tx: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + token_hash: &[u8; 32], +) -> Result> { + sqlx::query(SELECT_SESSION_BY_TOKEN_SQL) + .bind(token_hash.as_slice()) + .fetch_optional(&mut **tx) + .await? + .as_ref() + .map(session_from_row) + .transpose() +} + +fn session_from_row(row: &SqliteRow) -> Result { + let identity = IdpIdentity::new( + row.try_get::("identity_issuer")?, + row.try_get::("identity_subject")?, + ) + .map_err(|err| { + Error::Other(format!( + "stored auth session has an invalid identity: {err}" + )) + })?; + Ok(AuthSessionRecord { + id: parse_uuid(&row.try_get::("id")?)?, + identity, + login: row.try_get("login")?, + name: row.try_get("name")?, + email: row.try_get("email")?, + avatar_url: row.try_get("avatar_url")?, + user_agent: row.try_get("user_agent")?, + created_at: timestamp_from_row(row, "created_at_ms")?, + last_used_at: timestamp_from_row(row, "last_used_at_ms")?, + }) +} + +fn timestamp_from_row(row: &SqliteRow, column: &str) -> Result> { + let millis: i64 = row.try_get(column)?; + DateTime::from_timestamp_millis(millis) + .ok_or_else(|| Error::Other(format!("stored auth session has an invalid {column}"))) +} + +fn parse_uuid(value: &str) -> Result { + Uuid::parse_str(value) + .map_err(|err| Error::Other(format!("stored auth session has an invalid id: {err}"))) +} + +#[cfg(test)] +mod tests { + use std::sync::Arc; + + use chrono::{Duration, Utc}; + use fabro_types::IdpIdentity; + use tokio::task::JoinSet; + use uuid::Uuid; + + use super::{AuthSessionRecord, AuthSessionStore, RefreshToken, RotateOutcome}; + use crate::test_util::sqlite_auth_session_store; + + fn identity(subject: &str) -> IdpIdentity { + IdpIdentity::new("https://github.com", subject).unwrap() + } + + fn session(id: Uuid, subject: &str) -> AuthSessionRecord { + let now = Utc::now(); + AuthSessionRecord { + id, + identity: identity(subject), + login: "octocat".to_string(), + name: "The Octocat".to_string(), + email: "octocat@example.com".to_string(), + avatar_url: "https://example.com/octocat.png".to_string(), + user_agent: "fabro-cli/0.3".to_string(), + created_at: now, + last_used_at: now, + } + } + + /// Tokens are issued an hour back so that a fixture with a negative + /// `expires_in` is still a coherent row: issued in the past, expired since. + fn token(hash: [u8; 32], session_id: Uuid, expires_in: Duration) -> RefreshToken { + let now = Utc::now(); + RefreshToken { + token_hash: hash, + session_id, + issued_at: now - Duration::hours(1), + expires_at: now + expires_in, + used_at: None, + } + } + + async fn open_session( + store: &AuthSessionStore, + subject: &str, + hash: [u8; 32], + expires_in: Duration, + ) -> Uuid { + let id = Uuid::new_v4(); + store + .create_session(&session(id, subject), &token(hash, id, expires_in)) + .await + .unwrap(); + id + } + + #[tokio::test] + async fn create_session_round_trips_through_its_token() { + let (_dir, store) = sqlite_auth_session_store().await; + let id = open_session(&store, "12345", [1_u8; 32], Duration::days(30)).await; + + let found = store + .find_session_by_token_hash(&[1_u8; 32]) + .await + .unwrap() + .expect("session should be found by its token"); + assert_eq!(found.id, id); + assert_eq!(found.identity, identity("12345")); + assert_eq!(found.login, "octocat"); + assert_eq!(found.avatar_url, "https://example.com/octocat.png"); + + assert!( + store + .find_session_by_token_hash(&[9_u8; 32]) + .await + .unwrap() + .is_none() + ); + } + + #[tokio::test] + async fn rotate_spends_the_presented_token_and_issues_its_successor() { + let (_dir, store) = sqlite_auth_session_store().await; + let id = open_session(&store, "12345", [1_u8; 32], Duration::days(30)).await; + let now = Utc::now(); + + let outcome = store + .rotate( + &[1_u8; 32], + &[2_u8; 32], + now + Duration::days(30), + "fabro-cli/0.4", + now, + ) + .await + .unwrap(); + let RotateOutcome::Rotated(rotated) = outcome else { + panic!("expected rotation, got {outcome:?}"); + }; + assert_eq!(rotated.id, id); + // The rotation refreshes the session's user agent and last-used time, + // while created_at stays the true start of the chain. + assert_eq!(rotated.user_agent, "fabro-cli/0.4"); + assert_eq!( + rotated.last_used_at.timestamp_millis(), + now.timestamp_millis() + ); + assert!(rotated.created_at < rotated.last_used_at); + + // Both tokens still resolve to the session; only the new one is live. + assert!( + store + .find_session_by_token_hash(&[1_u8; 32]) + .await + .unwrap() + .is_some() + ); + let active = store + .active_cli_sessions(&identity("12345"), now) + .await + .unwrap(); + assert_eq!(active.len(), 1); + assert_eq!(active[0].session.id, id); + } + + #[tokio::test] + async fn rotate_reports_reuse_when_a_spent_token_is_replayed() { + let (_dir, store) = sqlite_auth_session_store().await; + let id = open_session(&store, "12345", [1_u8; 32], Duration::days(30)).await; + let now = Utc::now(); + + store + .rotate( + &[1_u8; 32], + &[2_u8; 32], + now + Duration::days(30), + "ua", + now, + ) + .await + .unwrap(); + let replay = store + .rotate( + &[1_u8; 32], + &[3_u8; 32], + now + Duration::days(30), + "ua", + now, + ) + .await + .unwrap(); + + let RotateOutcome::Reused(session) = replay else { + panic!("expected reuse, got {replay:?}"); + }; + assert_eq!(session.id, id); + } + + #[tokio::test] + async fn rotate_reports_expiry_before_reuse() { + let (_dir, store) = sqlite_auth_session_store().await; + let now = Utc::now(); + open_session(&store, "12345", [1_u8; 32], Duration::seconds(30)).await; + + // An unknown token is simply absent. + assert_eq!( + store + .rotate( + &[9_u8; 32], + &[8_u8; 32], + now + Duration::days(30), + "ua", + now + ) + .await + .unwrap(), + RotateOutcome::NotFound + ); + + // Live but past its expiry. + let later = now + Duration::seconds(60); + assert_eq!( + store + .rotate( + &[1_u8; 32], + &[2_u8; 32], + later + Duration::days(30), + "ua", + later + ) + .await + .unwrap(), + RotateOutcome::Expired + ); + + // Spend it while still valid, then let it expire: expiry wins over + // reuse, so a stale retry does not read as a live replay. + store + .rotate( + &[1_u8; 32], + &[2_u8; 32], + now + Duration::seconds(30), + "ua", + now, + ) + .await + .unwrap(); + assert_eq!( + store + .rotate( + &[1_u8; 32], + &[4_u8; 32], + later + Duration::days(30), + "ua", + later + ) + .await + .unwrap(), + RotateOutcome::Expired + ); + } + + #[tokio::test] + async fn active_cli_sessions_lists_only_live_sessions_owned_by_the_identity() { + let (_dir, store) = sqlite_auth_session_store().await; + let now = Utc::now(); + let live = open_session(&store, "12345", [1_u8; 32], Duration::days(30)).await; + let expired = open_session(&store, "12345", [2_u8; 32], Duration::seconds(-1)).await; + let other = open_session(&store, "67890", [3_u8; 32], Duration::days(30)).await; + + let active = store + .active_cli_sessions(&identity("12345"), now) + .await + .unwrap(); + let ids: Vec = active.iter().map(|entry| entry.session.id).collect(); + assert_eq!(ids, vec![live], "expired={expired}, other-identity={other}"); + assert!(active[0].expires_at > now); + + // Once the successor token expires in turn, the session drops off the + // list: the spent predecessor does not keep it alive. + store + .rotate( + &[1_u8; 32], + &[4_u8; 32], + now + Duration::seconds(30), + "ua", + now, + ) + .await + .unwrap(); + assert!( + store + .active_cli_sessions(&identity("12345"), now + Duration::seconds(60)) + .await + .unwrap() + .is_empty() + ); + } + + #[tokio::test] + async fn delete_session_removes_its_tokens() { + let (_dir, store) = sqlite_auth_session_store().await; + let now = Utc::now(); + let id = open_session(&store, "12345", [1_u8; 32], Duration::days(30)).await; + store + .rotate( + &[1_u8; 32], + &[2_u8; 32], + now + Duration::days(30), + "ua", + now, + ) + .await + .unwrap(); + + store.delete_session(id).await.unwrap(); + + for hash in [[1_u8; 32], [2_u8; 32]] { + assert!( + store + .find_session_by_token_hash(&hash) + .await + .unwrap() + .is_none(), + "cascade should remove every token in the chain" + ); + } + } + + #[tokio::test] + async fn delete_active_session_for_identity_requires_a_live_owned_session() { + let (_dir, store) = sqlite_auth_session_store().await; + let now = Utc::now(); + let owned = open_session(&store, "12345", [1_u8; 32], Duration::days(30)).await; + let expired = open_session(&store, "12345", [2_u8; 32], Duration::seconds(-1)).await; + + // Another identity cannot revoke it. + assert_eq!( + store + .delete_active_session_for_identity(&identity("67890"), owned, now) + .await + .unwrap(), + 0 + ); + // Neither can the owner, once nothing in it is spendable. + assert_eq!( + store + .delete_active_session_for_identity(&identity("12345"), expired, now) + .await + .unwrap(), + 0 + ); + // An unknown session is a no-op rather than an error. + assert_eq!( + store + .delete_active_session_for_identity(&identity("12345"), Uuid::new_v4(), now) + .await + .unwrap(), + 0 + ); + + assert_eq!( + store + .delete_active_session_for_identity(&identity("12345"), owned, now) + .await + .unwrap(), + 1 + ); + assert!( + store + .find_session_by_token_hash(&[1_u8; 32]) + .await + .unwrap() + .is_none() + ); + assert!( + store + .find_session_by_token_hash(&[2_u8; 32]) + .await + .unwrap() + .is_some(), + "revoking one session must not touch another" + ); + } + + #[tokio::test] + async fn gc_expired_drops_expired_tokens_and_the_sessions_left_empty() { + let (_dir, store) = sqlite_auth_session_store().await; + let now = Utc::now(); + let live = open_session(&store, "12345", [1_u8; 32], Duration::days(30)).await; + open_session(&store, "12345", [2_u8; 32], Duration::seconds(-1)).await; + + assert_eq!(store.gc_expired(now).await.unwrap(), 1); + assert!( + store + .find_session_by_token_hash(&[2_u8; 32]) + .await + .unwrap() + .is_none() + ); + assert_eq!( + store + .find_session_by_token_hash(&[1_u8; 32]) + .await + .unwrap() + .map(|session| session.id), + Some(live) + ); + // Idempotent: a second sweep finds nothing left to remove. + assert_eq!(store.gc_expired(now).await.unwrap(), 0); + } + + #[tokio::test] + async fn concurrent_rotation_lets_exactly_one_caller_win() { + let (_dir, store) = sqlite_auth_session_store().await; + let store = Arc::new(store); + let now = Utc::now(); + open_session(&store, "12345", [0_u8; 32], Duration::days(30)).await; + + let mut tasks = JoinSet::new(); + for index in 1..=8_u8 { + let store = Arc::clone(&store); + tasks.spawn(async move { + store + .rotate( + &[0_u8; 32], + &[index; 32], + now + Duration::days(30), + "ua", + now, + ) + .await + .unwrap() + }); + } + + let mut rotated = 0; + let mut reused = 0; + while let Some(outcome) = tasks.join_next().await { + match outcome.unwrap() { + RotateOutcome::Rotated(_) => rotated += 1, + RotateOutcome::Reused(_) => reused += 1, + other => panic!("unexpected outcome {other:?}"), + } + } + // SQLite's write lock serialises the claiming UPDATE, so the losers + // see a spent token rather than racing past it. + assert_eq!(rotated, 1); + assert_eq!(reused, 7); + } +} diff --git a/lib/components/fabro-store/src/lib.rs b/lib/components/fabro-store/src/lib.rs index 0eb246d04d..3c57d53e22 100644 --- a/lib/components/fabro-store/src/lib.rs +++ b/lib/components/fabro-store/src/lib.rs @@ -1,6 +1,7 @@ use chrono::{DateTime, Utc}; mod artifact_store; +pub mod auth_session_store; mod error; mod keyed_mutex; mod keys; @@ -18,6 +19,9 @@ pub use artifact_store::{ ArtifactKey, ArtifactStore, NodeArtifact, StageArtifactEntry, retry_storage_segment, stage_storage_segment, }; +pub use auth_session_store::{ + ActiveCliSession, AuthSessionRecord, AuthSessionStore, RefreshToken, RotateOutcome, +}; pub use error::{Error, Result}; pub use fabro_types::{ EventEnvelope, PendingInterviewRecord, Run, RunBlobId, RunProjection, StageId, StageProjection, @@ -34,8 +38,8 @@ pub use run_summary_store::{ }; pub use serializable_projection::SerializableProjection; pub use slate::{ - AuthCode, AuthCodeStore, Blob, BlobStore, CachedRunProjection, ConsumeOutcome, Database, - RefreshToken, RefreshTokenStore, RunCatalogIndex, RunDatabase, Runs, UnreadableRun, + AuthCode, AuthCodeStore, Blob, BlobStore, CachedRunProjection, Database, RunCatalogIndex, + RunDatabase, Runs, UnreadableRun, }; pub use types::EventPayload; diff --git a/lib/components/fabro-store/src/record/mod.rs b/lib/components/fabro-store/src/record/mod.rs index 64c956120f..189863171b 100644 --- a/lib/components/fabro-store/src/record/mod.rs +++ b/lib/components/fabro-store/src/record/mod.rs @@ -5,22 +5,18 @@ //! type. //! - [`RecordId`]: converts the typed id to and from key segments. //! - [`Repository`]: performs the generic get/put/delete/scan/gc operations. -//! - [`transaction`]: batches multiple typed writes into one atomic SlateDB -//! write. //! //! Production callers should add a named domain store on top of this layer //! rather than exposing `Repository` directly. See `slate/auth_codes.rs`, -//! `slate/auth_tokens.rs`, `slate/blob_store.rs`, and -//! `slate/run_catalog_index.rs` for the intended pattern. +//! `slate/blob_store.rs`, and `slate/run_catalog_index.rs` for the intended +//! pattern. mod codec; mod record_id; mod repository; -mod transaction; pub(crate) use codec::{Codec, JsonCodec, MarkerCodec, RawBytesCodec}; pub(crate) use repository::Repository; -pub(crate) use transaction::transaction; use crate::Result; diff --git a/lib/components/fabro-store/src/record/repository.rs b/lib/components/fabro-store/src/record/repository.rs index bfcc18fdc6..33196036ca 100644 --- a/lib/components/fabro-store/src/record/repository.rs +++ b/lib/components/fabro-store/src/record/repository.rs @@ -78,7 +78,7 @@ use crate::{Error, Result, keys}; /// stores. /// /// This type is intentionally `pub(crate)`: callers should interact through a -/// named store such as `AuthCodeStore` or `RefreshTokenStore`, which can add +/// named store such as `AuthCodeStore` or `BlobStore`, which can add /// domain-specific behavior on top of the generic storage primitives here. pub(crate) struct Repository { db: Arc, diff --git a/lib/components/fabro-store/src/record/transaction.rs b/lib/components/fabro-store/src/record/transaction.rs deleted file mode 100644 index 8c7c710ace..0000000000 --- a/lib/components/fabro-store/src/record/transaction.rs +++ /dev/null @@ -1,208 +0,0 @@ -use slatedb::{Db, WriteBatch}; - -use super::repository::key_for_id; -use super::{Codec, Record}; -use crate::Result; - -pub(crate) async fn transaction(db: &Db, f: F) -> Result -where - F: FnOnce(&mut Tx) -> Result, -{ - let mut tx = Tx::new(); - let value = f(&mut tx)?; - if tx.has_ops { - db.write(tx.batch).await?; - } - Ok(value) -} - -pub(crate) struct Tx { - batch: WriteBatch, - /// SlateDB rejects empty `WriteBatch` commits; skip the write entirely - /// when the closure produced no operations. - has_ops: bool, -} - -impl Tx { - fn new() -> Self { - Self { - batch: WriteBatch::new(), - has_ops: false, - } - } - - pub(crate) fn put(&mut self, record: &R) -> Result<&mut Self> { - let id = record.id(); - self.put_at(&id, record) - } - - pub(crate) fn put_at(&mut self, id: &R::Id, record: &R) -> Result<&mut Self> { - self.batch - .put(key_for_id::(id)?, R::Codec::encode(record)?); - self.has_ops = true; - Ok(self) - } - - #[allow( - dead_code, - reason = "Shared transaction surface; current production callers only use put paths" - )] - pub(crate) fn delete(&mut self, id: &R::Id) -> Result<&mut Self> { - self.batch.delete(key_for_id::(id)?); - self.has_ops = true; - Ok(self) - } -} - -#[cfg(test)] -mod tests { - use std::sync::Arc; - - use object_store::memory::InMemory; - use serde::{Deserialize, Serialize}; - - use super::{Record, Tx, transaction}; - use crate::record::{Codec, JsonCodec, Repository}; - use crate::{Error, Result}; - - #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] - struct TxRecord { - id: String, - payload: String, - poisoned: bool, - } - - impl Record for TxRecord { - type Id = String; - type Codec = JsonCodec; - - const PREFIX: &'static str = "test/transaction"; - - fn id(&self) -> Self::Id { - self.id.clone() - } - } - - #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] - struct FailingRecord { - id: String, - poisoned: bool, - } - - struct FailingCodec; - - impl Codec for FailingCodec { - fn encode(value: &FailingRecord) -> Result> { - if value.poisoned { - return Err(Error::Other( - "poisoned record refused to encode".to_string(), - )); - } - serde_json::to_vec(value).map_err(Into::into) - } - - fn decode(bytes: &[u8]) -> Result { - serde_json::from_slice(bytes).map_err(Into::into) - } - } - - impl Record for FailingRecord { - type Id = String; - type Codec = FailingCodec; - - const PREFIX: &'static str = "test/failing-transaction"; - - fn id(&self) -> Self::Id { - self.id.clone() - } - } - - async fn db() -> Arc { - Arc::new( - slatedb::Db::open("transaction-tests", Arc::new(InMemory::new())) - .await - .unwrap(), - ) - } - - #[tokio::test] - async fn closure_error_short_circuits_without_writing() { - let db = db().await; - let repo = Repository::::new(Arc::clone(&db)); - let record = TxRecord { - id: "record-1".to_string(), - payload: "hello".to_string(), - poisoned: false, - }; - - let error = transaction::<(), _>(&db, |tx| { - tx.put(&record)?; - Err(Error::Other("stop before commit".to_string())) - }) - .await - .unwrap_err(); - - assert_eq!(error.to_string(), "stop before commit"); - assert!(repo.get(&record.id()).await.unwrap().is_none()); - } - - #[tokio::test] - async fn encode_failure_discards_the_entire_batch() { - let db = db().await; - let repo = Repository::::new(Arc::clone(&db)); - let good = FailingRecord { - id: "good".to_string(), - poisoned: false, - }; - let bad = FailingRecord { - id: "bad".to_string(), - poisoned: true, - }; - - let error = transaction::<(), _>(&db, |tx| { - tx.put(&good)?; - tx.put(&bad)?; - Ok(()) - }) - .await - .unwrap_err(); - - assert_eq!(error.to_string(), "poisoned record refused to encode"); - assert!(repo.get(&good.id()).await.unwrap().is_none()); - assert!(repo.get(&bad.id()).await.unwrap().is_none()); - } - - #[tokio::test] - async fn empty_transaction_returns_without_writing() { - let db = db().await; - let repo = Repository::::new(Arc::clone(&db)); - - let value = transaction(&db, |_tx: &mut Tx| Ok::<_, Error>("ok")) - .await - .unwrap(); - - assert_eq!(value, "ok"); - assert!(repo.get(&"missing".to_string()).await.unwrap().is_none()); - } - - #[tokio::test] - async fn delete_operations_are_committed() { - let db = db().await; - let repo = Repository::::new(Arc::clone(&db)); - let record = TxRecord { - id: "delete-me".to_string(), - payload: "hello".to_string(), - poisoned: false, - }; - repo.put(&record).await.unwrap(); - - transaction::<(), _>(&db, |tx| { - tx.delete::(&record.id())?; - Ok(()) - }) - .await - .unwrap(); - - assert!(repo.get(&record.id()).await.unwrap().is_none()); - } -} diff --git a/lib/components/fabro-store/src/slate/auth_tokens.rs b/lib/components/fabro-store/src/slate/auth_tokens.rs deleted file mode 100644 index 5fed6749d6..0000000000 --- a/lib/components/fabro-store/src/slate/auth_tokens.rs +++ /dev/null @@ -1,598 +0,0 @@ -use std::sync::Arc; - -use chrono::{DateTime, Utc}; -use dashmap::DashMap; -use fabro_types::IdpIdentity; -use futures::StreamExt; -use serde::{Deserialize, Serialize}; -use uuid::Uuid; - -use crate::record::{JsonCodec, Record, Repository, transaction}; -use crate::{KeyedMutex, Result}; - -const REPLAY_REVOCATION_TTL_SECONDS: i64 = 60; - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -pub struct RefreshToken { - pub token_hash: [u8; 32], - pub chain_id: Uuid, - pub identity: IdpIdentity, - pub login: String, - pub name: String, - pub email: String, - #[serde(default, skip_serializing_if = "String::is_empty")] - pub avatar_url: String, - pub issued_at: DateTime, - pub expires_at: DateTime, - pub last_used_at: DateTime, - pub used: bool, - pub user_agent: String, -} - -impl Record for RefreshToken { - type Id = [u8; 32]; - type Codec = JsonCodec; - - const PREFIX: &'static str = "auth/refresh"; - - fn id(&self) -> Self::Id { - self.token_hash - } -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub enum ConsumeOutcome { - Rotated(RefreshToken, Box), - Reused(RefreshToken), - Expired, - NotFound, -} - -pub struct RefreshTokenStore { - db: Arc, - repo: Repository, - consume_locks: KeyedMutex<[u8; 32]>, - /// In-memory only: persisting attacker-supplied hashes would be an - /// unbounded-growth surface under a token-stuffing attack. - replay_revocations: DashMap<[u8; 32], DateTime>, -} - -impl std::fmt::Debug for RefreshTokenStore { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("RefreshTokenStore").finish_non_exhaustive() - } -} - -impl RefreshTokenStore { - pub(crate) fn new(db: Arc) -> Self { - Self { - repo: Repository::new(Arc::clone(&db)), - db, - consume_locks: KeyedMutex::new(), - replay_revocations: DashMap::new(), - } - } - - pub async fn insert_refresh_token(&self, token: RefreshToken) -> Result<()> { - self.repo.put(&token).await - } - - pub async fn find_refresh_token(&self, token_hash: &[u8; 32]) -> Result> { - self.repo.get(token_hash).await - } - - pub async fn active_cli_sessions( - &self, - identity: &IdpIdentity, - now: DateTime, - ) -> Result> { - let mut active_by_chain = std::collections::HashMap::::new(); - let mut tokens = self.repo.scan_stream(); - - while let Some(result) = tokens.next().await { - let (_, token) = result?; - if token.identity != *identity || token.used || token.expires_at <= now { - continue; - } - - active_by_chain - .entry(token.chain_id) - .and_modify(|current| { - if token.last_used_at > current.last_used_at { - *current = token.clone(); - } - }) - .or_insert(token); - } - - Ok(active_by_chain.into_values().collect()) - } - - pub async fn consume_and_rotate( - &self, - presented_hash: [u8; 32], - new_token: RefreshToken, - now: DateTime, - ) -> Result { - let _guard = self.consume_locks.lock(presented_hash).await; - - let outcome = match self.repo.get(&presented_hash).await? { - None => ConsumeOutcome::NotFound, - Some(existing) if now >= existing.expires_at => ConsumeOutcome::Expired, - Some(existing) if existing.used => ConsumeOutcome::Reused(existing), - Some(existing) => { - let mut old_token = existing.clone(); - old_token.used = true; - old_token.last_used_at = now; - - transaction(&self.db, |tx| { - tx.put(&old_token)?; - tx.put(&new_token)?; - Ok(()) - }) - .await?; - - ConsumeOutcome::Rotated(old_token, Box::new(new_token)) - } - }; - - Ok(outcome) - } - - pub async fn delete_chain(&self, chain_id: Uuid) -> Result { - self.repo.gc(|token| token.chain_id == chain_id).await - } - - pub async fn delete_active_chain_for_identity( - &self, - identity: &IdpIdentity, - chain_id: Uuid, - now: DateTime, - ) -> Result { - let mut token_hashes = Vec::new(); - let mut has_active_token = false; - let mut tokens = self.repo.scan_stream(); - - while let Some(result) = tokens.next().await { - let (_, token) = result?; - if token.identity != *identity || token.chain_id != chain_id { - continue; - } - if !token.used && token.expires_at > now { - has_active_token = true; - } - token_hashes.push(token.token_hash); - } - - if !has_active_token { - return Ok(0); - } - - let deleted = u64::try_from(token_hashes.len()).unwrap_or(u64::MAX); - transaction(&self.db, |tx| { - for token_hash in &token_hashes { - tx.delete::(token_hash)?; - } - Ok(()) - }) - .await?; - Ok(deleted) - } - - pub async fn gc_expired(&self, cutoff: DateTime) -> Result { - self.repo.gc(|token| token.expires_at <= cutoff).await - } - - pub fn mark_refresh_token_replay(&self, token_hash: [u8; 32], now: DateTime) { - self.replay_revocations.insert( - token_hash, - now + chrono::Duration::seconds(REPLAY_REVOCATION_TTL_SECONDS), - ); - self.replay_revocations - .retain(|_, expires_at| *expires_at > now); - } - - pub fn was_recently_replay_revoked(&self, token_hash: &[u8; 32], now: DateTime) -> bool { - self.replay_revocations - .retain(|_, expires_at| *expires_at > now); - self.replay_revocations - .get(token_hash) - .is_some_and(|expires_at| *expires_at > now) - } -} - -#[cfg(test)] -mod tests { - use std::sync::Arc; - use std::time::Duration; - - use chrono::Duration as ChronoDuration; - use object_store::memory::InMemory; - use tokio::task::JoinSet; - use uuid::Uuid; - - use super::{ConsumeOutcome, RefreshToken, RefreshTokenStore}; - use crate::Database; - - async fn store() -> Arc { - let db = Database::new( - Arc::new(InMemory::new()), - "", - Duration::from_millis(1), - None, - ); - db.refresh_tokens().await.unwrap() - } - - fn refresh_token(hash: [u8; 32], chain_id: Uuid, used: bool) -> RefreshToken { - let now = chrono::Utc::now(); - RefreshToken { - token_hash: hash, - chain_id, - identity: fabro_types::IdpIdentity::new("https://github.com", "12345").unwrap(), - login: "octocat".to_string(), - name: "The Octocat".to_string(), - email: "octocat@example.com".to_string(), - avatar_url: String::new(), - issued_at: now, - expires_at: now + ChronoDuration::days(30), - last_used_at: now, - used, - user_agent: "fabro-test".to_string(), - } - } - - fn alternate_identity() -> fabro_types::IdpIdentity { - fabro_types::IdpIdentity::new("https://github.com", "67890").unwrap() - } - - #[tokio::test] - async fn insert_find_rotate_and_reuse_work() { - let store = store().await; - let chain_id = Uuid::new_v4(); - let old_hash = [1_u8; 32]; - let new_hash = [2_u8; 32]; - let old = refresh_token(old_hash, chain_id, false); - let new = refresh_token(new_hash, chain_id, false); - store.insert_refresh_token(old.clone()).await.unwrap(); - - assert_eq!( - store.find_refresh_token(&old_hash).await.unwrap(), - Some(old) - ); - - let rotated = store - .consume_and_rotate(old_hash, new.clone(), chrono::Utc::now()) - .await - .unwrap(); - let ConsumeOutcome::Rotated(old_used, new_saved) = rotated else { - panic!("expected rotation"); - }; - assert!(old_used.used); - assert_eq!(new_saved.token_hash, new_hash); - assert_eq!( - store.find_refresh_token(&old_hash).await.unwrap(), - Some(old_used.clone()) - ); - assert!( - store - .find_refresh_token(&old_hash) - .await - .unwrap() - .expect("rotated old token should still exist") - .used - ); - - let replay = store - .consume_and_rotate( - old_hash, - refresh_token([3_u8; 32], chain_id, false), - chrono::Utc::now(), - ) - .await - .unwrap(); - let ConsumeOutcome::Reused(reused) = replay else { - panic!("expected replay to return the original used row"); - }; - assert_eq!(reused.chain_id, chain_id); - } - - #[test] - fn deserializes_legacy_json_without_avatar_url() { - let entry: RefreshToken = serde_json::from_value(serde_json::json!({ - "token_hash": [1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1, 1], - "chain_id": "00000000-0000-4000-8000-000000000000", - "identity": { - "issuer": "https://github.com", - "subject": "12345" - }, - "login": "octocat", - "name": "The Octocat", - "email": "octocat@example.com", - "issued_at": "2026-01-01T00:00:00Z", - "expires_at": "2026-02-01T00:00:00Z", - "last_used_at": "2026-01-01T00:00:00Z", - "used": false, - "user_agent": "fabro-test" - })) - .unwrap(); - - assert_eq!(entry.avatar_url, ""); - } - - #[test] - fn serializes_avatar_url_when_present() { - let mut entry = refresh_token([1_u8; 32], Uuid::new_v4(), false); - entry.avatar_url = "https://example.com/octocat.png".to_string(); - - let json = serde_json::to_value(&entry).unwrap(); - - assert_eq!(json["avatar_url"], "https://example.com/octocat.png"); - } - - #[tokio::test] - async fn missing_and_expired_tokens_are_reported() { - let store = store().await; - let chain_id = Uuid::new_v4(); - - assert_eq!( - store - .consume_and_rotate( - [7_u8; 32], - refresh_token([8_u8; 32], chain_id, false), - chrono::Utc::now(), - ) - .await - .unwrap(), - ConsumeOutcome::NotFound - ); - - let mut expired = refresh_token([9_u8; 32], chain_id, false); - expired.expires_at = chrono::Utc::now() - ChronoDuration::seconds(1); - store.insert_refresh_token(expired.clone()).await.unwrap(); - - assert_eq!( - store - .consume_and_rotate( - expired.token_hash, - refresh_token([10_u8; 32], chain_id, false), - chrono::Utc::now(), - ) - .await - .unwrap(), - ConsumeOutcome::Expired - ); - assert_eq!( - store.find_refresh_token(&expired.token_hash).await.unwrap(), - Some(expired) - ); - } - - #[tokio::test] - async fn concurrent_rotation_has_one_winner() { - let store = store().await; - let chain_id = Uuid::new_v4(); - let hash = [9_u8; 32]; - store - .insert_refresh_token(refresh_token(hash, chain_id, false)) - .await - .unwrap(); - - let mut tasks = JoinSet::new(); - for idx in 0..16_u8 { - let store = Arc::clone(&store); - tasks.spawn(async move { - store - .consume_and_rotate( - hash, - refresh_token([idx; 32], chain_id, false), - chrono::Utc::now(), - ) - .await - .unwrap() - }); - } - - let mut rotated = 0; - let mut reused = 0; - while let Some(result) = tasks.join_next().await { - match result.unwrap() { - ConsumeOutcome::Rotated(_, _) => rotated += 1, - ConsumeOutcome::Reused(_) => reused += 1, - other => panic!("unexpected outcome: {other:?}"), - } - } - - assert_eq!(rotated, 1); - assert_eq!(reused, 15); - } - - #[tokio::test] - async fn delete_chain_removes_all_matching_tokens() { - let store = store().await; - let chain_id = Uuid::new_v4(); - store - .insert_refresh_token(refresh_token([1_u8; 32], chain_id, false)) - .await - .unwrap(); - store - .insert_refresh_token(refresh_token([2_u8; 32], chain_id, true)) - .await - .unwrap(); - - assert_eq!(store.delete_chain(chain_id).await.unwrap(), 2); - assert!( - store - .find_refresh_token(&[1_u8; 32]) - .await - .unwrap() - .is_none() - ); - assert!( - store - .find_refresh_token(&[2_u8; 32]) - .await - .unwrap() - .is_none() - ); - } - - #[tokio::test] - async fn delete_active_chain_for_identity_requires_active_owned_token() { - let store = store().await; - let identity = fabro_types::IdpIdentity::new("https://github.com", "12345").unwrap(); - let chain_id = Uuid::new_v4(); - let other_chain_id = Uuid::new_v4(); - let active = refresh_token([1_u8; 32], chain_id, false); - let used = refresh_token([2_u8; 32], chain_id, true); - let mut other_identity = refresh_token([3_u8; 32], chain_id, false); - other_identity.identity = alternate_identity(); - let other_chain = refresh_token([4_u8; 32], other_chain_id, false); - - for token in [ - active.clone(), - used.clone(), - other_identity.clone(), - other_chain.clone(), - ] { - store.insert_refresh_token(token).await.unwrap(); - } - - assert_eq!( - store - .delete_active_chain_for_identity(&identity, chain_id, chrono::Utc::now()) - .await - .unwrap(), - 2 - ); - assert!( - store - .find_refresh_token(&active.token_hash) - .await - .unwrap() - .is_none() - ); - assert!( - store - .find_refresh_token(&used.token_hash) - .await - .unwrap() - .is_none() - ); - assert_eq!( - store - .find_refresh_token(&other_identity.token_hash) - .await - .unwrap(), - Some(other_identity) - ); - assert_eq!( - store - .find_refresh_token(&other_chain.token_hash) - .await - .unwrap(), - Some(other_chain) - ); - } - - #[tokio::test] - async fn delete_active_chain_for_identity_returns_zero_without_active_token() { - let store = store().await; - let identity = fabro_types::IdpIdentity::new("https://github.com", "12345").unwrap(); - let chain_id = Uuid::new_v4(); - let used = refresh_token([1_u8; 32], chain_id, true); - store.insert_refresh_token(used.clone()).await.unwrap(); - - assert_eq!( - store - .delete_active_chain_for_identity(&identity, chain_id, chrono::Utc::now()) - .await - .unwrap(), - 0 - ); - assert_eq!( - store.find_refresh_token(&used.token_hash).await.unwrap(), - Some(used) - ); - } - - #[tokio::test] - async fn active_cli_sessions_return_newest_active_token_per_chain_for_identity() { - let store = store().await; - let identity = fabro_types::IdpIdentity::new("https://github.com", "12345").unwrap(); - let now = chrono::Utc::now(); - let duplicate_chain_id = Uuid::new_v4(); - let other_chain_id = Uuid::new_v4(); - - let mut old_duplicate = refresh_token([1_u8; 32], duplicate_chain_id, false); - old_duplicate.last_used_at = now - ChronoDuration::minutes(10); - old_duplicate.issued_at = now - ChronoDuration::minutes(20); - let mut newest_duplicate = refresh_token([2_u8; 32], duplicate_chain_id, false); - newest_duplicate.last_used_at = now - ChronoDuration::minutes(1); - newest_duplicate.issued_at = now - ChronoDuration::minutes(15); - let mut other_active = refresh_token([3_u8; 32], other_chain_id, false); - other_active.last_used_at = now - ChronoDuration::minutes(3); - - store.insert_refresh_token(old_duplicate).await.unwrap(); - store - .insert_refresh_token(newest_duplicate.clone()) - .await - .unwrap(); - store - .insert_refresh_token(other_active.clone()) - .await - .unwrap(); - - let sessions = store.active_cli_sessions(&identity, now).await.unwrap(); - assert_eq!(sessions.len(), 2); - assert!(sessions.contains(&newest_duplicate)); - assert!(sessions.contains(&other_active)); - } - - #[tokio::test] - async fn active_cli_sessions_exclude_expired_used_and_other_identity_tokens() { - let store = store().await; - let identity = fabro_types::IdpIdentity::new("https://github.com", "12345").unwrap(); - let now = chrono::Utc::now(); - - let active = refresh_token([1_u8; 32], Uuid::new_v4(), false); - let mut expired = refresh_token([2_u8; 32], Uuid::new_v4(), false); - expired.expires_at = now - ChronoDuration::seconds(1); - let used = refresh_token([3_u8; 32], Uuid::new_v4(), true); - let mut other_identity = refresh_token([4_u8; 32], Uuid::new_v4(), false); - other_identity.identity = alternate_identity(); - - for token in [active.clone(), expired, used, other_identity] { - store.insert_refresh_token(token).await.unwrap(); - } - - let sessions = store.active_cli_sessions(&identity, now).await.unwrap(); - assert_eq!(sessions, vec![active]); - } - - #[tokio::test] - async fn gc_expired_removes_only_tokens_at_or_before_cutoff() { - let store = store().await; - let chain_id = Uuid::new_v4(); - - let mut expired = refresh_token([4_u8; 32], chain_id, true); - expired.expires_at = chrono::Utc::now() - ChronoDuration::days(8); - let live = refresh_token([5_u8; 32], chain_id, false); - - store.insert_refresh_token(expired.clone()).await.unwrap(); - store.insert_refresh_token(live.clone()).await.unwrap(); - - assert_eq!(store.gc_expired(chrono::Utc::now()).await.unwrap(), 1); - assert!( - store - .find_refresh_token(&expired.token_hash) - .await - .unwrap() - .is_none() - ); - assert_eq!( - store.find_refresh_token(&live.token_hash).await.unwrap(), - Some(live) - ); - } -} diff --git a/lib/components/fabro-store/src/slate/mod.rs b/lib/components/fabro-store/src/slate/mod.rs index 617d56565c..2ee1058d08 100644 --- a/lib/components/fabro-store/src/slate/mod.rs +++ b/lib/components/fabro-store/src/slate/mod.rs @@ -1,5 +1,4 @@ mod auth_codes; -mod auth_tokens; mod blob_store; mod projection_cache; mod run_catalog_index; @@ -11,7 +10,6 @@ use std::sync::{Arc, OnceLock}; use std::time::Duration; pub use auth_codes::{AuthCode, AuthCodeStore}; -pub use auth_tokens::{ConsumeOutcome, RefreshToken, RefreshTokenStore}; pub use blob_store::{Blob, BlobStore}; use chrono::{DateTime, Utc}; use fabro_types::{Run, RunId, SessionId}; @@ -50,7 +48,6 @@ pub struct Database { blobs: Arc>>, catalog_index: Arc>>, auth_codes: Arc>>, - refresh_tokens: Arc>>, projection_cache: Arc, projection_cache_warmed: Arc>, run_summary_store: Arc>>, @@ -83,7 +80,6 @@ impl Database { blobs: Arc::new(OnceCell::new()), catalog_index: Arc::new(OnceCell::new()), auth_codes: Arc::new(OnceCell::new()), - refresh_tokens: Arc::new(OnceCell::new()), projection_cache: Arc::new(RunProjectionCache::default()), projection_cache_warmed: Arc::new(OnceCell::new()), run_summary_store: Arc::new(OnceLock::new()), @@ -442,15 +438,27 @@ impl Database { Ok(Arc::clone(store)) } - pub async fn refresh_tokens(&self) -> Result> { - let store = self - .refresh_tokens - .get_or_try_init(|| async { - let db = Arc::new(self.open_db().await?); - Ok::<_, Error>(Arc::new(RefreshTokenStore::new(db))) - }) + /// Delete every record under the retired `auth/refresh` prefix. + /// + /// Refresh tokens moved to SQLite without an import, so these records are + /// unreadable -- and the reaper that used to collect them is gone, so + /// nothing else would ever remove them. Returns the number of records + /// deleted; a later boot finds the prefix empty and does nothing. + pub async fn retire_refresh_token_keyspace(&self) -> Result { + let db = self.open_db().await?; + let mut iter = db + .scan_prefix(keys::SlateKey::new("auth").with("refresh").into_prefix()) .await?; - Ok(Arc::clone(store)) + let mut batch = slatedb::WriteBatch::new(); + let mut deletes = 0_u64; + while let Some(entry) = iter.next().await? { + batch.delete(entry.key); + deletes += 1; + } + if deletes > 0 { + db.write(batch).await?; + } + Ok(deletes) } #[must_use] @@ -559,6 +567,44 @@ mod tests { (object_store, store) } + #[tokio::test] + async fn retire_refresh_token_keyspace_clears_the_prefix_and_is_idempotent() { + let (_object_store, store) = make_store(); + let db = store.open_db().await.unwrap(); + + let refresh_keys = ["aaa", "bbb"].map(|id| { + keys::SlateKey::new("auth") + .with("refresh") + .with(id) + .as_ref() + .to_vec() + }); + // "auth/code" sorts adjacent to "auth/refresh" and is still live, so + // it is the neighbour a too-wide prefix delete would take with it. + let auth_code_key = keys::SlateKey::new("auth") + .with("code") + .with("keep") + .as_ref() + .to_vec(); + + let mut batch = slatedb::WriteBatch::new(); + for key in &refresh_keys { + batch.put(key.as_slice(), b"{}".as_slice()); + } + batch.put(auth_code_key.as_slice(), b"{}".as_slice()); + db.write(batch).await.unwrap(); + + assert_eq!(store.retire_refresh_token_keyspace().await.unwrap(), 2); + assert_eq!(store.retire_refresh_token_keyspace().await.unwrap(), 0); + for key in &refresh_keys { + assert!(db.get(key.as_slice()).await.unwrap().is_none()); + } + assert!( + db.get(auth_code_key.as_slice()).await.unwrap().is_some(), + "retiring refresh tokens must not touch the auth code prefix" + ); + } + async fn make_summary_store() -> (tempfile::TempDir, Arc) { let (directory, store) = test_util::sqlite_summary_store().await; (directory, Arc::new(store)) diff --git a/lib/components/fabro-store/src/test_util.rs b/lib/components/fabro-store/src/test_util.rs index bbc0b87174..e5bedb0d33 100644 --- a/lib/components/fabro-store/src/test_util.rs +++ b/lib/components/fabro-store/src/test_util.rs @@ -1,4 +1,14 @@ use crate::RunSummaryStore; +use crate::auth_session_store::AuthSessionStore; + +pub(crate) async fn sqlite_auth_session_store() -> (tempfile::TempDir, AuthSessionStore) { + let directory = tempfile::tempdir().unwrap(); + let database = fabro_db::Database::connect(directory.path().join("fabro.sqlite3")) + .await + .unwrap(); + database.migrate().await.unwrap(); + (directory, AuthSessionStore::new(database.clone_pool())) +} pub(crate) async fn sqlite_summary_store() -> (tempfile::TempDir, RunSummaryStore) { let directory = tempfile::tempdir().unwrap(); diff --git a/lib/foundation/fabro-db/migrations/2026072501_auth_sessions.sql b/lib/foundation/fabro-db/migrations/2026072501_auth_sessions.sql new file mode 100644 index 0000000000..170d92b127 --- /dev/null +++ b/lib/foundation/fabro-db/migrations/2026072501_auth_sessions.sql @@ -0,0 +1,49 @@ +-- A CLI auth session is a rotation chain: the identity and profile are facts +-- about the chain, not about any single token in it. Keeping them here means a +-- chain has exactly one owner, and `created_at_ms` is the real start of the +-- session rather than the newest token's issue time. +CREATE TABLE auth_sessions ( + id TEXT PRIMARY KEY NOT NULL, + identity_issuer TEXT NOT NULL, + identity_subject TEXT NOT NULL, + login TEXT NOT NULL, + name TEXT NOT NULL, + email TEXT NOT NULL, + avatar_url TEXT NOT NULL DEFAULT '', + user_agent TEXT NOT NULL DEFAULT '', + created_at_ms INTEGER NOT NULL, + last_used_at_ms INTEGER NOT NULL, + CHECK (length(id) = 36), + CHECK (length(identity_issuer) > 0), + CHECK (length(identity_subject) > 0) +); + +CREATE INDEX auth_sessions_by_identity + ON auth_sessions (identity_issuer, identity_subject, last_used_at_ms DESC); + +-- Rotated tokens are retained until they expire so a replayed token is still +-- recognisable as one that existed, rather than indistinguishable from a +-- forgery. `used_at_ms IS NULL` marks the one token that can still be spent. +CREATE TABLE refresh_tokens ( + token_hash BLOB PRIMARY KEY NOT NULL, + session_id TEXT NOT NULL REFERENCES auth_sessions(id) ON DELETE CASCADE, + issued_at_ms INTEGER NOT NULL, + expires_at_ms INTEGER NOT NULL, + used_at_ms INTEGER, + CHECK (length(token_hash) = 32), + CHECK (expires_at_ms > issued_at_ms) +); + +-- Deliberately no ordering CHECK between a session's timestamps and its +-- tokens'. Rotation stamps `now` from the process clock against rows written +-- by an earlier request, so an NTP step backwards would turn a harmless clock +-- anomaly into refresh failing outright for every affected session. + +-- Rotation marks the presented token used before issuing its successor, so a +-- chain can only ever hold one live token. Enforcing it here turns an implicit +-- code convention into a constraint, and lets the session listing find the +-- live token by index instead of grouping candidates in memory. +CREATE UNIQUE INDEX refresh_tokens_one_live_per_session + ON refresh_tokens (session_id) WHERE used_at_ms IS NULL; + +CREATE INDEX refresh_tokens_by_expiry ON refresh_tokens (expires_at_ms); diff --git a/lib/foundation/fabro-db/tests/sqlite.rs b/lib/foundation/fabro-db/tests/sqlite.rs index bc4981b062..d8bc7e6d5a 100644 --- a/lib/foundation/fabro-db/tests/sqlite.rs +++ b/lib/foundation/fabro-db/tests/sqlite.rs @@ -73,6 +73,16 @@ async fn connect_creates_parent_directory_and_migrate_is_idempotent() -> anyhow: .await?; assert_eq!(runs_table_count, 1); + for table in ["auth_sessions", "refresh_tokens"] { + let count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = ?", + ) + .bind(table) + .fetch_one(database.pool()) + .await?; + assert_eq!(count, 1, "{table} table should exist"); + } + let legacy_import_table_count: i64 = sqlx::query_scalar( "SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = 'legacy_imports'", ) @@ -390,6 +400,139 @@ INSERT INTO runs ( Ok(()) } +#[tokio::test] +async fn auth_sessions_schema_enforces_one_live_token_and_cascade() -> anyhow::Result<()> { + let dir = tempfile::tempdir()?; + let database = fabro_db::Database::connect(dir.path().join("fabro.sqlite3")).await?; + database.migrate().await?; + + let session = "11111111-1111-4111-8111-111111111111"; + insert_auth_session(database.pool(), session, "https://github.com", "12345").await?; + insert_refresh_token(database.pool(), &[1_u8; 32], session, 1_000, None).await?; + + // Rotation marks the old token used before issuing the new one, so a + // second live token in the same chain must be impossible. + assert!( + insert_refresh_token(database.pool(), &[2_u8; 32], session, 1_000, None) + .await + .is_err(), + "a session must not hold two live refresh tokens" + ); + // A used token alongside the live one is the normal post-rotation state. + insert_refresh_token(database.pool(), &[2_u8; 32], session, 1_000, Some(1_500)).await?; + + assert!( + insert_refresh_token( + database.pool(), + &[3_u8; 32], + "22222222-2222-4222-8222-222222222222", + 1_000, + None + ) + .await + .is_err(), + "a refresh token must reference an existing session" + ); + + sqlx::query("DELETE FROM auth_sessions WHERE id = ?") + .bind(session) + .execute(database.pool()) + .await?; + let orphaned: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM refresh_tokens") + .fetch_one(database.pool()) + .await?; + assert_eq!( + orphaned, 0, + "deleting a session should cascade to its tokens" + ); + + Ok(()) +} + +#[tokio::test] +async fn auth_sessions_schema_rejects_invalid_rows() -> anyhow::Result<()> { + let dir = tempfile::tempdir()?; + let database = fabro_db::Database::connect(dir.path().join("fabro.sqlite3")).await?; + database.migrate().await?; + + for (id, issuer, subject) in [ + ("too-short", "https://github.com", "12345"), + ("33333333-3333-4333-8333-333333333333", "", "12345"), + ( + "44444444-4444-4444-8444-444444444444", + "https://github.com", + "", + ), + ] { + assert!( + insert_auth_session(database.pool(), id, issuer, subject) + .await + .is_err(), + "auth session row should be rejected: id={id}, issuer={issuer}, subject={subject}" + ); + } + + let session = "55555555-5555-4555-8555-555555555555"; + insert_auth_session(database.pool(), session, "https://github.com", "12345").await?; + for (hash, expires_at_ms, used_at_ms) in + [(vec![9_u8; 31], 1_000, None), (vec![9_u8; 32], 0, None)] + { + assert!( + insert_refresh_token(database.pool(), &hash, session, expires_at_ms, used_at_ms) + .await + .is_err(), + "refresh token row should be rejected: len={}, expires_at_ms={expires_at_ms}", + hash.len() + ); + } + + Ok(()) +} + +async fn insert_auth_session( + pool: &fabro_db::DbPool, + id: &str, + identity_issuer: &str, + identity_subject: &str, +) -> Result<(), sqlx::Error> { + sqlx::query( + r" +INSERT INTO auth_sessions ( + id, identity_issuer, identity_subject, login, name, email, + created_at_ms, last_used_at_ms +) VALUES (?, ?, ?, 'octocat', 'The Octocat', 'octocat@example.com', 0, 0) +", + ) + .bind(id) + .bind(identity_issuer) + .bind(identity_subject) + .execute(pool) + .await?; + Ok(()) +} + +async fn insert_refresh_token( + pool: &fabro_db::DbPool, + token_hash: &[u8], + session_id: &str, + expires_at_ms: i64, + used_at_ms: Option, +) -> Result<(), sqlx::Error> { + sqlx::query( + r" +INSERT INTO refresh_tokens (token_hash, session_id, issued_at_ms, expires_at_ms, used_at_ms) +VALUES (?, ?, 0, ?, ?) +", + ) + .bind(token_hash) + .bind(session_id) + .bind(expires_at_ms) + .bind(used_at_ms) + .execute(pool) + .await?; + Ok(()) +} + #[tokio::test] async fn environments_schema_rejects_invalid_rows() -> anyhow::Result<()> { let dir = tempfile::tempdir()?;