diff --git a/.sqlx/query-8003624cedbac8b094c83933578517abfb2eaf8e59d1d52c7ea59bf5d11cfcfe.json b/.sqlx/query-8003624cedbac8b094c83933578517abfb2eaf8e59d1d52c7ea59bf5d11cfcfe.json new file mode 100644 index 0000000..5cd6618 --- /dev/null +++ b/.sqlx/query-8003624cedbac8b094c83933578517abfb2eaf8e59d1d52c7ea59bf5d11cfcfe.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM session_tokens WHERE id = $1 AND did = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Int4", + "Text" + ] + }, + "nullable": [] + }, + "hash": "8003624cedbac8b094c83933578517abfb2eaf8e59d1d52c7ea59bf5d11cfcfe" +} diff --git a/.sqlx/query-847ce3c34985d0957526c87e0a20c6b4e5daae08a338f7635def682ac0689cf6.json b/.sqlx/query-847ce3c34985d0957526c87e0a20c6b4e5daae08a338f7635def682ac0689cf6.json deleted file mode 100644 index 0c0e68f..0000000 --- a/.sqlx/query-847ce3c34985d0957526c87e0a20c6b4e5daae08a338f7635def682ac0689cf6.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "DELETE FROM session_tokens WHERE access_jti = $1", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Text" - ] - }, - "nullable": [] - }, - "hash": "847ce3c34985d0957526c87e0a20c6b4e5daae08a338f7635def682ac0689cf6" -} diff --git a/.sqlx/query-a27e93bc594babbada10afe5c3e33a65909ec69c579329916833e4b0fe2332d3.json b/.sqlx/query-a27e93bc594babbada10afe5c3e33a65909ec69c579329916833e4b0fe2332d3.json new file mode 100644 index 0000000..edec496 --- /dev/null +++ b/.sqlx/query-a27e93bc594babbada10afe5c3e33a65909ec69c579329916833e4b0fe2332d3.json @@ -0,0 +1,15 @@ +{ + "db_name": "PostgreSQL", + "query": "DELETE FROM session_tokens WHERE access_jti = $1 AND did = $2", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [] + }, + "hash": "a27e93bc594babbada10afe5c3e33a65909ec69c579329916833e4b0fe2332d3" +} diff --git a/.sqlx/query-cf874abcb72017e775fe699a0b77ae9341355f30e4af84968ffeb9135dba745f.json b/.sqlx/query-cf874abcb72017e775fe699a0b77ae9341355f30e4af84968ffeb9135dba745f.json deleted file mode 100644 index 7586fd3..0000000 --- a/.sqlx/query-cf874abcb72017e775fe699a0b77ae9341355f30e4af84968ffeb9135dba745f.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "DELETE FROM session_tokens WHERE id = $1", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Int4" - ] - }, - "nullable": [] - }, - "hash": "cf874abcb72017e775fe699a0b77ae9341355f30e4af84968ffeb9135dba745f" -} diff --git a/crates/tranquil-api/src/server/session.rs b/crates/tranquil-api/src/server/session.rs index 74ab9cd..38dc556 100644 --- a/crates/tranquil-api/src/server/session.rs +++ b/crates/tranquil-api/src/server/session.rs @@ -452,7 +452,12 @@ pub async fn delete_session( ) -> Result, ApiError> { let jti = tranquil_pds::auth::extract_jti_from_headers(&headers) .ok_or(ApiError::AuthenticationRequired)?; - match state.repos.session.delete_session_by_access_jti(&jti).await { + match state + .repos + .session + .delete_session_by_access_jti(&jti, &auth.did) + .await + { Ok(rows) if rows > 0 => { let session_cache_key = tranquil_pds::cache_keys::session_key(&auth.did, &jti); let _ = state.cache.delete(&session_cache_key).await; @@ -505,72 +510,15 @@ pub async fn refresh_session( ))); } }; - match state.repos.session.lookup_refresh_grace(&refresh_jti).await { - Ok(tranquil_db_traits::RefreshGraceLookup::NotUsed) => {} - Ok(tranquil_db_traits::RefreshGraceLookup::Replay(replay)) => { - // Verify the presented token's signature before issuing anything: a - // rotated jti is a public value, so we must never mint for a forged - // unsigned token carrying it. - let key = match tranquil_pds::config::decrypt_key( - &replay.key_bytes, - Some(replay.encryption_version), - ) { - Ok(k) => k, - Err(e) => { - error!("Failed to decrypt user key for grace replay: {:?}", e); - return Err(ApiError::InternalError(None)); - } - }; - if tranquil_pds::auth::verify_refresh_token(&refresh_token, &key).is_err() { - return Err(ApiError::AuthenticationFailed(Some( - "Invalid refresh token".into(), - ))); - } - info!( - "Refresh token reuse within grace window for jti: {refresh_jti}; replaying tokens" - ); - let (access_jwt, refresh_jwt) = remint_grace_tokens(&replay, &key)?; - return build_refresh_session_output(&state, replay.did, access_jwt, refresh_jwt).await; - } - Ok(tranquil_db_traits::RefreshGraceLookup::Compromised { - session_id, - key_bytes, - encryption_version, - }) => { - // Never revoke a session for an unverified token: verify the - // signature first, so a forged unsigned token cannot force a logout. - let key = match tranquil_pds::config::decrypt_key(&key_bytes, Some(encryption_version)) - { - Ok(k) => k, - Err(e) => { - error!("Failed to decrypt user key for grace check: {:?}", e); - return Err(ApiError::InternalError(None)); - } - }; - if tranquil_pds::auth::verify_refresh_token(&refresh_token, &key).is_err() { - return Err(ApiError::AuthenticationFailed(Some( - "Invalid refresh token".into(), - ))); - } - warn!( - "Refresh token reuse outside grace window for jti: {refresh_jti}; revoking session" - ); - if let Err(e) = state.repos.session.delete_session_by_id(session_id).await { - error!( - "Failed to revoke session {} for refresh token reuse: {:?}", - session_id.as_i32(), - e - ); - return Err(ApiError::InternalError(None)); - } - return Err(ApiError::AuthenticationFailed(Some( - "Refresh token has been revoked due to suspected compromise".into(), - ))); - } - Err(e) => { - error!("Database error checking refresh token grace: {:?}", e); - return Err(ApiError::InternalError(None)); - } + if let Some(result) = dispatch_refresh_grace( + &state, + &refresh_token, + &refresh_jti, + state.repos.session.lookup_refresh_grace(&refresh_jti).await, + ) + .await + { + return result; } let session_row = match state .repos @@ -580,9 +528,18 @@ pub async fn refresh_session( { Ok(Some(row)) => row, Ok(None) => { - return Err(ApiError::AuthenticationFailed(Some( - "Invalid refresh token".into(), - ))); + return dispatch_refresh_grace( + &state, + &refresh_token, + &refresh_jti, + state.repos.session.lookup_refresh_grace(&refresh_jti).await, + ) + .await + .unwrap_or_else(|| { + Err(ApiError::AuthenticationFailed(Some( + "Invalid refresh token".into(), + ))) + }); } Err(e) => { error!("Database error fetching session: {:?}", e); @@ -628,6 +585,7 @@ pub async fn refresh_session( } }; let refresh_data = tranquil_db_traits::SessionRefreshData { + did: session_row.did.clone(), old_refresh_jti: refresh_jti.clone(), session_id: session_row.id, new_access_jti: new_access_meta.jti.clone(), @@ -670,6 +628,102 @@ pub async fn refresh_session( build_refresh_session_output(&state, session_row.did, access_jwt, refresh_jwt).await } +async fn dispatch_refresh_grace( + state: &AppState, + refresh_token: &str, + presented_jti: &str, + lookup: Result, +) -> Option, ApiError>> { + match lookup { + Ok(tranquil_db_traits::RefreshGraceLookup::NotUsed) => None, + Ok(tranquil_db_traits::RefreshGraceLookup::Replay(replay)) => { + Some(serve_refresh_grace_replay(state, refresh_token, presented_jti, replay).await) + } + Ok(tranquil_db_traits::RefreshGraceLookup::Compromised { + did, + session_id, + key_bytes, + encryption_version, + }) => Some(Err(revoke_compromised_session( + state, + refresh_token, + presented_jti, + did, + session_id, + key_bytes, + encryption_version, + ) + .await)), + Err(e) => { + error!("Database error checking refresh token grace: {:?}", e); + Some(Err(ApiError::InternalError(None))) + } + } +} + +async fn serve_refresh_grace_replay( + state: &AppState, + refresh_token: &str, + presented_jti: &str, + replay: tranquil_db_traits::RefreshGraceReplay, +) -> Result, ApiError> { + let key = + match tranquil_pds::config::decrypt_key(&replay.key_bytes, Some(replay.encryption_version)) + { + Ok(k) => k, + Err(e) => { + error!("Failed to decrypt user key for grace replay: {:?}", e); + return Err(ApiError::InternalError(None)); + } + }; + if tranquil_pds::auth::verify_refresh_token(refresh_token, &key).is_err() { + return Err(ApiError::AuthenticationFailed(Some( + "Invalid refresh token".into(), + ))); + } + info!("Refresh token reuse within grace window for jti: {presented_jti}; replaying tokens"); + let (access_jwt, refresh_jwt) = remint_grace_tokens(&replay, &key)?; + build_refresh_session_output(state, replay.did, access_jwt, refresh_jwt).await +} + +async fn revoke_compromised_session( + state: &AppState, + refresh_token: &str, + presented_jti: &str, + did: Did, + session_id: SessionId, + key_bytes: Vec, + encryption_version: i32, +) -> ApiError { + let key = match tranquil_pds::config::decrypt_key(&key_bytes, Some(encryption_version)) { + Ok(k) => k, + Err(e) => { + error!("Failed to decrypt user key for grace check: {:?}", e); + return ApiError::InternalError(None); + } + }; + if tranquil_pds::auth::verify_refresh_token(refresh_token, &key).is_err() { + return ApiError::AuthenticationFailed(Some("Invalid refresh token".into())); + } + warn!("Refresh token reuse outside grace window for jti: {presented_jti}; revoking session"); + if let Err(e) = state + .repos + .session + .delete_session_by_id(session_id, &did) + .await + { + error!( + "Failed to revoke session {} for refresh token reuse: {:?}", + session_id.as_i32(), + e + ); + return ApiError::InternalError(None); + } + ApiError::AuthenticationFailed(Some( + "Refresh token has been revoked due to suspected compromise".into(), + )) +} + /// Re-mint the access/refresh JWTs for a grace-window replay from the session's /// current jtis and signing key. We never persist the signed JWTs; they are /// reconstructed on demand so a benignly-racing client converges on the same @@ -1112,7 +1166,7 @@ pub async fn revoke_session( state .repos .session - .delete_session_by_id(session_id) + .delete_session_by_id(session_id, &auth.did) .await .log_db_err("deleting session")?; let cache_key = tranquil_pds::cache_keys::session_key(&auth.did, &access_jti); diff --git a/crates/tranquil-db-traits/src/session.rs b/crates/tranquil-db-traits/src/session.rs index ab589b6..1b5eb88 100644 --- a/crates/tranquil-db-traits/src/session.rs +++ b/crates/tranquil-db-traits/src/session.rs @@ -195,6 +195,7 @@ pub enum RefreshGraceLookup { NotUsed, Replay(RefreshGraceReplay), Compromised { + did: Did, session_id: SessionId, key_bytes: Vec, encryption_version: i32, @@ -203,6 +204,7 @@ pub enum RefreshGraceLookup { #[derive(Debug, Clone)] pub struct SessionRefreshData { + pub did: Did, pub old_refresh_jti: String, pub session_id: SessionId, pub new_access_jti: String, @@ -225,18 +227,17 @@ pub trait SessionRepository: Send + Sync { refresh_jti: &str, ) -> Result, DbError>; - async fn update_session_tokens( + async fn delete_session_by_access_jti( &self, - session_id: SessionId, - new_access_jti: &str, - new_refresh_jti: &str, - new_access_expires_at: DateTime, - new_refresh_expires_at: DateTime, - ) -> Result<(), DbError>; - - async fn delete_session_by_access_jti(&self, access_jti: &str) -> Result; + access_jti: &str, + did: &Did, + ) -> Result; - async fn delete_session_by_id(&self, session_id: SessionId) -> Result; + async fn delete_session_by_id( + &self, + session_id: SessionId, + did: &Did, + ) -> Result; async fn delete_sessions_by_did(&self, did: &Did) -> Result; diff --git a/crates/tranquil-db/src/postgres/session.rs b/crates/tranquil-db/src/postgres/session.rs index bbb24f2..fe31087 100644 --- a/crates/tranquil-db/src/postgres/session.rs +++ b/crates/tranquil-db/src/postgres/session.rs @@ -114,38 +114,15 @@ impl SessionRepository for PostgresSessionRepository { })) } - async fn update_session_tokens( + async fn delete_session_by_access_jti( &self, - session_id: SessionId, - new_access_jti: &str, - new_refresh_jti: &str, - new_access_expires_at: DateTime, - new_refresh_expires_at: DateTime, - ) -> Result<(), DbError> { - sqlx::query!( - r#" - UPDATE session_tokens - SET access_jti = $1, refresh_jti = $2, access_expires_at = $3, - refresh_expires_at = $4, updated_at = NOW() - WHERE id = $5 - "#, - new_access_jti, - new_refresh_jti, - new_access_expires_at, - new_refresh_expires_at, - session_id.as_i32() - ) - .execute(&self.pool) - .await - .map_err(map_sqlx_error)?; - - Ok(()) - } - - async fn delete_session_by_access_jti(&self, access_jti: &str) -> Result { + access_jti: &str, + did: &Did, + ) -> Result { let result = sqlx::query!( - "DELETE FROM session_tokens WHERE access_jti = $1", - access_jti + "DELETE FROM session_tokens WHERE access_jti = $1 AND did = $2", + access_jti, + did.as_str() ) .execute(&self.pool) .await @@ -154,10 +131,15 @@ impl SessionRepository for PostgresSessionRepository { Ok(result.rows_affected()) } - async fn delete_session_by_id(&self, session_id: SessionId) -> Result { + async fn delete_session_by_id( + &self, + session_id: SessionId, + did: &Did, + ) -> Result { let result = sqlx::query!( - "DELETE FROM session_tokens WHERE id = $1", - session_id.as_i32() + "DELETE FROM session_tokens WHERE id = $1 AND did = $2", + session_id.as_i32(), + did.as_str() ) .execute(&self.pool) .await @@ -308,6 +290,7 @@ impl SessionRepository for PostgresSessionRepository { })) } else { Ok(RefreshGraceLookup::Compromised { + did: Did::from(r.did), session_id: SessionId::new(r.session_id), key_bytes: r.key_bytes, encryption_version: r.encryption_version.unwrap_or(0), diff --git a/crates/tranquil-pds/src/delegation/scopes.rs b/crates/tranquil-pds/src/delegation/scopes.rs index f8bf636..07b0196 100644 --- a/crates/tranquil-pds/src/delegation/scopes.rs +++ b/crates/tranquil-pds/src/delegation/scopes.rs @@ -12,8 +12,7 @@ pub struct ScopePreset { pub scopes: &'static str, } -pub const OWNER_FULL_SCOPES: &str = - "atproto repo:* blob:*/* identity:* account:*?action=manage"; +pub const OWNER_FULL_SCOPES: &str = "atproto repo:* blob:*/* identity:* account:*?action=manage"; pub const EDITOR_FULL_SCOPES: &str = "atproto repo:*?action=create repo:*?action=update repo:*?action=delete blob:*/*"; @@ -210,8 +209,7 @@ mod tests { .find(|p| p.name == "editor") .expect("editor preset") .scopes; - let requested = - "atproto repo:*?action=create identity:* account:*?action=manage blob:*/*"; + let requested = "atproto repo:*?action=create identity:* account:*?action=manage blob:*/*"; let result = intersect_scopes(requested, editor); assert!(result.split_whitespace().any(|s| s == "atproto")); assert!(result.contains("repo:*?action=create")); diff --git a/crates/tranquil-pds/src/scheduled.rs b/crates/tranquil-pds/src/scheduled.rs index 08a1fb6..93f8da7 100644 --- a/crates/tranquil-pds/src/scheduled.rs +++ b/crates/tranquil-pds/src/scheduled.rs @@ -180,9 +180,7 @@ pub async fn collect_current_repo_blocks( let block = match block_store.get(&cid).await { Ok(Some(b)) => b, Ok(None) => continue, - Err(e) - if crate::api::error::ApiError::detail_is_repo_corruption(&format!("{e:#}")) => - { + Err(e) if crate::api::error::ApiError::detail_is_repo_corruption(&format!("{e:#}")) => { warn!(cid = %cid, error = %format!("{e:#}"), "skipping corrupt block during repo walk"); continue; } diff --git a/crates/tranquil-pds/src/state.rs b/crates/tranquil-pds/src/state.rs index 5d6c3f3..d00da7c 100644 --- a/crates/tranquil-pds/src/state.rs +++ b/crates/tranquil-pds/src/state.rs @@ -488,12 +488,7 @@ fn migrate_delegation_preset_scopes(metastore: &tranquil_store::metastore::Metas "repo:*?action=create repo:*?action=update repo:*?action=delete blob:*/*"; let infra = metastore.infra_ops(); - if infra - .get_server_config(MARKER_KEY) - .ok() - .flatten() - .is_some() - { + if infra.get_server_config(MARKER_KEY).ok().flatten().is_some() { return; } @@ -505,16 +500,21 @@ fn migrate_delegation_preset_scopes(metastore: &tranquil_store::metastore::Metas return; } }; - let editors = - match ops.remap_grant_scopes(LEGACY_EDITOR_SCOPES, crate::delegation::EDITOR_FULL_SCOPES) { - Ok(n) => n, - Err(e) => { - tracing::error!(error = ?e, "delegation editor-scope migration failed, will retry on next start"); - return; - } - }; + let editors = match ops + .remap_grant_scopes(LEGACY_EDITOR_SCOPES, crate::delegation::EDITOR_FULL_SCOPES) + { + Ok(n) => n, + Err(e) => { + tracing::error!(error = ?e, "delegation editor-scope migration failed, will retry on next start"); + return; + } + }; if owners + editors > 0 { - tracing::info!(owners, editors, "upgraded legacy delegation grants to preset scopes"); + tracing::info!( + owners, + editors, + "upgraded legacy delegation grants to preset scopes" + ); } if let Err(e) = infra.upsert_server_config(MARKER_KEY, "done") { diff --git a/crates/tranquil-pds/tests/scope_edge_cases.rs b/crates/tranquil-pds/tests/scope_edge_cases.rs index 51744b8..59f3d3c 100644 --- a/crates/tranquil-pds/tests/scope_edge_cases.rs +++ b/crates/tranquil-pds/tests/scope_edge_cases.rs @@ -225,8 +225,8 @@ fn test_delegation_validate_multiple() { } #[test] -fn test_delegation_intersect_empty_granted_returns_empty() { - assert_eq!(intersect_scopes("atproto", ""), ""); +fn test_delegation_intersect_empty_grant_keeps_only_atproto() { + assert_eq!(intersect_scopes("atproto", ""), "atproto"); assert_eq!(intersect_scopes("repo:*", ""), ""); } diff --git a/crates/tranquil-store/src/metastore/client.rs b/crates/tranquil-store/src/metastore/client.rs index 7e70c89..4136489 100644 --- a/crates/tranquil-store/src/metastore/client.rs +++ b/crates/tranquil-store/src/metastore/client.rs @@ -1463,33 +1463,16 @@ impl tranquil_db_traits::SessionRepository for Metastore recv(rx).await } - async fn update_session_tokens( + async fn delete_session_by_access_jti( &self, - session_id: tranquil_db_traits::SessionId, - new_access_jti: &str, - new_refresh_jti: &str, - new_access_expires_at: DateTime, - new_refresh_expires_at: DateTime, - ) -> Result<(), DbError> { - let (tx, rx) = oneshot::channel(); - self.pool.send(MetastoreRequest::Session( - SessionRequest::UpdateSessionTokens { - session_id, - new_access_jti: new_access_jti.to_owned(), - new_refresh_jti: new_refresh_jti.to_owned(), - new_access_expires_at, - new_refresh_expires_at, - tx, - }, - ))?; - recv(rx).await - } - - async fn delete_session_by_access_jti(&self, access_jti: &str) -> Result { + access_jti: &str, + did: &Did, + ) -> Result { let (tx, rx) = oneshot::channel(); self.pool.send(MetastoreRequest::Session( SessionRequest::DeleteSessionByAccessJti { access_jti: access_jti.to_owned(), + did: did.clone(), tx, }, ))?; @@ -1499,10 +1482,15 @@ impl tranquil_db_traits::SessionRepository for Metastore async fn delete_session_by_id( &self, session_id: tranquil_db_traits::SessionId, + did: &Did, ) -> Result { let (tx, rx) = oneshot::channel(); self.pool.send(MetastoreRequest::Session( - SessionRequest::DeleteSessionById { session_id, tx }, + SessionRequest::DeleteSessionById { + session_id, + did: did.clone(), + tx, + }, ))?; recv(rx).await } diff --git a/crates/tranquil-store/src/metastore/delegation_ops.rs b/crates/tranquil-store/src/metastore/delegation_ops.rs index 19b4c81..5a5ef96 100644 --- a/crates/tranquil-store/src/metastore/delegation_ops.rs +++ b/crates/tranquil-store/src/metastore/delegation_ops.rs @@ -471,8 +471,13 @@ mod tests { let ctrl_owner = did("did:plc:olaren"); let ctrl_editor = did("did:plc:teq"); - ops.create_delegation(&owner, &ctrl_owner, &DbScope::new("atproto").unwrap(), &owner) - .unwrap(); + ops.create_delegation( + &owner, + &ctrl_owner, + &DbScope::new("atproto").unwrap(), + &owner, + ) + .unwrap(); ops.create_delegation( &owner, &ctrl_editor, @@ -519,7 +524,11 @@ mod tests { UserHash::from_did("did:plc:scallop"), ); let mut batch = ops.db.batch(); - batch.insert(&ops.indexes, corrupt_key.as_slice(), b"not a grant".as_slice()); + batch.insert( + &ops.indexes, + corrupt_key.as_slice(), + b"not a grant".as_slice(), + ); batch.commit().unwrap(); assert_eq!(ops.remap_grant_scopes("atproto", OWNER_FULL).unwrap(), 1); diff --git a/crates/tranquil-store/src/metastore/handler.rs b/crates/tranquil-store/src/metastore/handler.rs index 2edd004..c22c32d 100644 --- a/crates/tranquil-store/src/metastore/handler.rs +++ b/crates/tranquil-store/src/metastore/handler.rs @@ -125,6 +125,7 @@ fn metastore_to_db(e: MetastoreError) -> DbError { } } +#[derive(Debug, Clone, Copy, PartialEq, Eq)] enum Routing { Sharded(u64), Global, @@ -848,20 +849,14 @@ pub enum SessionRequest { refresh_jti: String, tx: Tx>, }, - UpdateSessionTokens { - session_id: SessionId, - new_access_jti: String, - new_refresh_jti: String, - new_access_expires_at: DateTime, - new_refresh_expires_at: DateTime, - tx: Tx<()>, - }, DeleteSessionByAccessJti { access_jti: String, + did: Did, tx: Tx, }, DeleteSessionById { session_id: SessionId, + did: Did, tx: Tx, }, DeleteSessionsByDid { @@ -955,10 +950,10 @@ impl SessionRequest { Self::CreateSession { .. } | Self::GetSessionByAccessJti { .. } | Self::GetSessionForRefresh { .. } - | Self::LookupRefreshGrace { .. } - | Self::DeleteSessionByAccessJti { .. } - | Self::DeleteSessionById { .. } => Routing::Global, - Self::DeleteSessionsByDid { did, .. } + | Self::LookupRefreshGrace { .. } => Routing::Global, + Self::DeleteSessionByAccessJti { did, .. } + | Self::DeleteSessionById { did, .. } + | Self::DeleteSessionsByDid { did, .. } | Self::DeleteSessionsByDidExceptJti { did, .. } | Self::ListSessionsByDid { did, .. } | Self::GetSessionAccessJtiById { did, .. } @@ -970,7 +965,7 @@ impl SessionRequest { | Self::GetSessionMfaStatus { did, .. } | Self::UpdateMfaVerified { did, .. } | Self::GetAppPasswordHashesByDid { did, .. } => did_to_routing(did.as_str()), - Self::UpdateSessionTokens { .. } | Self::RefreshSessionAtomic { .. } => Routing::Global, + Self::RefreshSessionAtomic { data, .. } => did_to_routing(data.did.as_str()), Self::ListAppPasswords { user_id, .. } | Self::GetAppPasswordsForLogin { user_id, .. } | Self::GetAppPasswordByName { user_id, .. } @@ -3637,40 +3632,27 @@ fn dispatch_session(state: &HandlerState, req: SessionRequest) .map_err(metastore_to_db); let _ = tx.send(result); } - SessionRequest::UpdateSessionTokens { - session_id, - new_access_jti, - new_refresh_jti, - new_access_expires_at, - new_refresh_expires_at, + SessionRequest::DeleteSessionByAccessJti { + access_jti, + did, tx, } => { let result = state .metastore .session_ops() - .update_session_tokens( - session_id, - &new_access_jti, - &new_refresh_jti, - new_access_expires_at, - new_refresh_expires_at, - ) + .delete_session_by_access_jti(&access_jti, &did) .map_err(metastore_to_db); let _ = tx.send(result); } - SessionRequest::DeleteSessionByAccessJti { access_jti, tx } => { - let result = state - .metastore - .session_ops() - .delete_session_by_access_jti(&access_jti) - .map_err(metastore_to_db); - let _ = tx.send(result); - } - SessionRequest::DeleteSessionById { session_id, tx } => { + SessionRequest::DeleteSessionById { + session_id, + did, + tx, + } => { let result = state .metastore .session_ops() - .delete_session_by_id(session_id) + .delete_session_by_id(session_id, &did) .map_err(metastore_to_db); let _ = tx.send(result); } @@ -5799,14 +5781,19 @@ fn dispatch_user(state: &HandlerState, req: UserReque let infra = state.metastore.infra_ops(); let code = input.invite_code.as_deref(); let result = reserve_invite(&infra, code).and_then(|()| { - finalize_account(&infra, code, user.create_password_account(&input), |result| { - if let Some(key_id) = input.reserved_key_id { - infra - .mark_signing_key_used(key_id) - .map_err(|e| CreateAccountError::Database(e.to_string()))?; - } - Ok(result.user_id) - }) + finalize_account( + &infra, + code, + user.create_password_account(&input), + |result| { + if let Some(key_id) = input.reserved_key_id { + infra + .mark_signing_key_used(key_id) + .map_err(|e| CreateAccountError::Database(e.to_string()))?; + } + Ok(result.user_id) + }, + ) }); let _ = tx.send(result); } @@ -5840,14 +5827,19 @@ fn dispatch_user(state: &HandlerState, req: UserReque let infra = state.metastore.infra_ops(); let code = input.invite_code.as_deref(); let result = reserve_invite(&infra, code).and_then(|()| { - finalize_account(&infra, code, user.create_passkey_account(&input), |result| { - if let Some(key_id) = input.reserved_key_id { - infra - .mark_signing_key_used(key_id) - .map_err(|e| CreateAccountError::Database(e.to_string()))?; - } - Ok(result.user_id) - }) + finalize_account( + &infra, + code, + user.create_passkey_account(&input), + |result| { + if let Some(key_id) = input.reserved_key_id { + infra + .mark_signing_key_used(key_id) + .map_err(|e| CreateAccountError::Database(e.to_string()))?; + } + Ok(result.user_id) + }, + ) }); let _ = tx.send(result); } @@ -6280,6 +6272,58 @@ mod tests { assert_eq!(indices, vec![0, 1, 2, 3, 0, 1, 2, 3]); } + #[test] + fn session_mutations_route_by_did() { + let dir = tempfile::TempDir::new().unwrap(); + let ms = Metastore::open( + dir.path(), + MetastoreConfig { + cache_size_bytes: 1024 * 1024, + }, + ) + .unwrap(); + let user_hashes = ms.user_hashes().as_ref(); + let did = Did::from("did:plc:limpet".to_string()); + let expected = did_to_routing(did.as_str()); + let sid = SessionId::new(7); + + let (tx, _rx) = oneshot::channel(); + let delete_by_id = SessionRequest::DeleteSessionById { + session_id: sid, + did: did.clone(), + tx, + }; + let (tx, _rx) = oneshot::channel(); + let delete_by_jti = SessionRequest::DeleteSessionByAccessJti { + access_jti: "a".to_string(), + did: did.clone(), + tx, + }; + let (tx, _rx) = oneshot::channel(); + let delete_by_did = SessionRequest::DeleteSessionsByDid { + did: did.clone(), + tx, + }; + let (tx, _rx) = oneshot::channel(); + let refresh = SessionRequest::RefreshSessionAtomic { + data: tranquil_db_traits::SessionRefreshData { + did: did.clone(), + old_refresh_jti: "r0".to_string(), + session_id: sid, + new_access_jti: "a1".to_string(), + new_refresh_jti: "r1".to_string(), + new_access_expires_at: Utc::now(), + new_refresh_expires_at: Utc::now(), + }, + tx, + }; + + assert_eq!(delete_by_id.routing(user_hashes), expected); + assert_eq!(delete_by_jti.routing(user_hashes), expected); + assert_eq!(delete_by_did.routing(user_hashes), expected); + assert_eq!(refresh.routing(user_hashes), expected); + } + #[tokio::test] async fn shutdown_completes_inflight() { let h = setup(); diff --git a/crates/tranquil-store/src/metastore/mod.rs b/crates/tranquil-store/src/metastore/mod.rs index 511261a..a6a8e3f 100644 --- a/crates/tranquil-store/src/metastore/mod.rs +++ b/crates/tranquil-store/src/metastore/mod.rs @@ -396,6 +396,7 @@ impl Metastore { #[cfg(test)] mod tests { use super::*; + use tranquil_types::Did; fn open_fresh() -> (tempfile::TempDir, Metastore) { let dir = tempfile::TempDir::new().unwrap(); @@ -474,6 +475,7 @@ mod tests { } fn legacy_refresh_data( + did: &Did, session_id: tranquil_db_traits::SessionId, old_refresh_jti: &str, new_access_jti: &str, @@ -481,6 +483,7 @@ mod tests { ) -> tranquil_db_traits::SessionRefreshData { let now = chrono::Utc::now(); tranquil_db_traits::SessionRefreshData { + did: did.clone(), old_refresh_jti: old_refresh_jti.to_string(), session_id, new_access_jti: new_access_jti.to_string(), @@ -494,15 +497,16 @@ mod tests { fn legacy_refresh_grace_replays_within_window() { use tranquil_db_traits::{RefreshGraceLookup, RefreshSessionResult}; let (_dir, ms) = open_fresh(); - create_test_user(&ms, "did:plc:grace", "grace.test"); + let did = Did::new("did:plc:grace".to_string()).unwrap(); + create_test_user(&ms, did.as_str(), "grace.test"); let ops = ms.session_ops(); let session_id = ops - .create_session(&legacy_session_create("did:plc:grace", "acc0", "ref0")) + .create_session(&legacy_session_create(did.as_str(), "acc0", "ref0")) .unwrap(); // The winning request rotates ref0 -> ref1. - let win = legacy_refresh_data(session_id, "ref0", "acc1", "ref1"); + let win = legacy_refresh_data(&did, session_id, "ref0", "acc1", "ref1"); assert!(matches!( ops.refresh_session_atomic(&win).unwrap(), RefreshSessionResult::Success @@ -523,7 +527,7 @@ mod tests { // The atomic path (two requests both past the used-check) also yields // the winner's current tokens rather than revoking. - let lose = legacy_refresh_data(session_id, "ref0", "accX", "refX"); + let lose = legacy_refresh_data(&did, session_id, "ref0", "accX", "refX"); match ops.refresh_session_atomic(&lose).unwrap() { RefreshSessionResult::GraceReplay(replay) => { assert_eq!(replay.access_jti, "acc1"); @@ -550,16 +554,21 @@ mod tests { fn legacy_refresh_superseded_token_within_grace_replays() { use tranquil_db_traits::{RefreshGraceLookup, RefreshSessionResult}; let (_dir, ms) = open_fresh(); - create_test_user(&ms, "did:plc:reuse", "reuse.test"); + let did = Did::new("did:plc:reuse".to_string()).unwrap(); + create_test_user(&ms, did.as_str(), "reuse.test"); let ops = ms.session_ops(); let session_id = ops - .create_session(&legacy_session_create("did:plc:reuse", "acc0", "ref0")) - .unwrap(); - ops.refresh_session_atomic(&legacy_refresh_data(session_id, "ref0", "acc1", "ref1")) - .unwrap(); - ops.refresh_session_atomic(&legacy_refresh_data(session_id, "ref1", "acc2", "ref2")) + .create_session(&legacy_session_create(did.as_str(), "acc0", "ref0")) .unwrap(); + ops.refresh_session_atomic(&legacy_refresh_data( + &did, session_id, "ref0", "acc1", "ref1", + )) + .unwrap(); + ops.refresh_session_atomic(&legacy_refresh_data( + &did, session_id, "ref1", "acc2", "ref2", + )) + .unwrap(); // ref0 is two rotations back but was rotated just now: still in window. match ops.lookup_refresh_grace("ref0").unwrap() { @@ -570,7 +579,7 @@ mod tests { other => panic!("expected Replay, got {other:?}"), } match ops - .refresh_session_atomic(&legacy_refresh_data(session_id, "ref0", "z", "z")) + .refresh_session_atomic(&legacy_refresh_data(&did, session_id, "ref0", "z", "z")) .unwrap() { RefreshSessionResult::GraceReplay(replay) => { @@ -583,6 +592,49 @@ mod tests { assert!(ops.get_session_by_access_jti("acc2").unwrap().is_some()); } + #[test] + fn rotated_refresh_jti_is_none_from_fetch_but_replay_from_grace() { + use tranquil_db_traits::RefreshGraceLookup; + let (_dir, ms) = open_fresh(); + let did = Did::new("did:plc:whelk".to_string()).unwrap(); + create_test_user(&ms, did.as_str(), "whelk.test"); + let ops = ms.session_ops(); + + let session_id = ops + .create_session(&legacy_session_create(did.as_str(), "acc0", "ref0")) + .unwrap(); + ops.refresh_session_atomic(&legacy_refresh_data( + &did, session_id, "ref0", "acc1", "ref1", + )) + .unwrap(); + + assert!(ops.get_session_for_refresh("ref0").unwrap().is_none()); + match ops.lookup_refresh_grace("ref0").unwrap() { + RefreshGraceLookup::Replay(replay) => assert_eq!(replay.refresh_jti, "ref1"), + other => panic!("expected Replay, got {other:?}"), + } + } + + #[test] + fn delete_session_is_scoped_to_did() { + let (_dir, ms) = open_fresh(); + let owner = Did::new("did:plc:limpet".to_string()).unwrap(); + let other = Did::new("did:plc:scallop".to_string()).unwrap(); + create_test_user(&ms, owner.as_str(), "limpet.test"); + let ops = ms.session_ops(); + + let session_id = ops + .create_session(&legacy_session_create(owner.as_str(), "acc0", "ref0")) + .unwrap(); + + assert_eq!(ops.delete_session_by_id(session_id, &other).unwrap(), 0); + assert_eq!(ops.delete_session_by_access_jti("acc0", &other).unwrap(), 0); + assert!(ops.get_session_by_access_jti("acc0").unwrap().is_some()); + + assert_eq!(ops.delete_session_by_access_jti("acc0", &owner).unwrap(), 1); + assert!(ops.get_session_by_access_jti("acc0").unwrap().is_none()); + } + // An old-format (12-byte, no rotated_at) used marker has unknown rotation // time and must classify as Compromised, never Replay. #[test] @@ -622,11 +674,12 @@ mod tests { REFRESH_GRACE_PERIOD_SECS, RefreshGraceLookup, RefreshSessionResult, }; let (_dir, ms) = open_fresh(); - create_test_user(&ms, "did:plc:stale", "stale.test"); + let did = Did::new("did:plc:stale".to_string()).unwrap(); + create_test_user(&ms, did.as_str(), "stale.test"); let ops = ms.session_ops(); let session_id = ops - .create_session(&legacy_session_create("did:plc:stale", "acc0", "ref0")) + .create_session(&legacy_session_create(did.as_str(), "acc0", "ref0")) .unwrap(); // Overwrite ref0's used marker with a current 20-byte marker rotated 3h ago, @@ -659,7 +712,7 @@ mod tests { } assert!(matches!( - ops.refresh_session_atomic(&legacy_refresh_data(session_id, "ref0", "z", "z")) + ops.refresh_session_atomic(&legacy_refresh_data(&did, session_id, "ref0", "z", "z")) .unwrap(), RefreshSessionResult::Compromise )); diff --git a/crates/tranquil-store/src/metastore/session_ops.rs b/crates/tranquil-store/src/metastore/session_ops.rs index 60cd655..6d8dd92 100644 --- a/crates/tranquil-store/src/metastore/session_ops.rs +++ b/crates/tranquil-store/src/metastore/session_ops.rs @@ -211,6 +211,7 @@ impl SessionOps { } Ok(RefreshGraceLookup::Compromised { + did: Did::from(session.did.clone()), session_id: SessionId::new(session_id), key_bytes: user.key_bytes, encryption_version: user.encryption_version, @@ -397,72 +398,11 @@ impl SessionOps { } } - pub fn update_session_tokens( + pub fn delete_session_by_access_jti( &self, - session_id: SessionId, - new_access_jti: &str, - new_refresh_jti: &str, - new_access_expires_at: DateTime, - new_refresh_expires_at: DateTime, - ) -> Result<(), MetastoreError> { - let mut session = self - .load_session_by_id(session_id.as_i32())? - .ok_or(MetastoreError::InvalidInput("session not found"))?; - - let user_hash = self.resolve_user_hash_from_did(&session.did); - let old_access_jti = session.access_jti.clone(); - let old_refresh_jti = session.refresh_jti.clone(); - - session.access_jti = new_access_jti.to_owned(); - session.refresh_jti = new_refresh_jti.to_owned(); - session.access_expires_at_ms = new_access_expires_at.timestamp_millis(); - session.refresh_expires_at_ms = new_refresh_expires_at.timestamp_millis(); - session.updated_at_ms = Utc::now().timestamp_millis(); - - let new_access_index = SessionIndexValue { - user_hash: user_hash.raw(), - session_id: session.id, - }; - let new_refresh_index = SessionIndexValue { - user_hash: user_hash.raw(), - session_id: session.id, - }; - - let mut batch = self.db.batch(); - batch.remove( - &self.auth, - session_by_access_key(&old_access_jti).as_slice(), - ); - batch.remove( - &self.auth, - session_by_refresh_key(&old_refresh_jti).as_slice(), - ); - batch.insert( - &self.auth, - session_primary_key(session.id).as_slice(), - session.serialize(), - ); - batch.insert( - &self.auth, - session_by_access_key(new_access_jti).as_slice(), - new_access_index.serialize(session.refresh_expires_at_ms), - ); - batch.insert( - &self.auth, - session_by_refresh_key(new_refresh_jti).as_slice(), - new_refresh_index.serialize(session.refresh_expires_at_ms), - ); - batch.insert( - &self.auth, - session_by_did_key(user_hash, session.id).as_slice(), - serialize_by_did_value(session.refresh_expires_at_ms), - ); - batch.commit().map_err(MetastoreError::Fjall)?; - - Ok(()) - } - - pub fn delete_session_by_access_jti(&self, access_jti: &str) -> Result { + access_jti: &str, + did: &Did, + ) -> Result { let index_key = session_by_access_key(access_jti); let index_val: Option = point_lookup( &self.auth, @@ -477,8 +417,8 @@ impl SessionOps { }; let session = match self.load_session_by_id(idx.session_id)? { - Some(s) => s, - None => return Ok(0), + Some(s) if s.did == did.as_str() => s, + _ => return Ok(0), }; let mut batch = self.db.batch(); @@ -488,10 +428,14 @@ impl SessionOps { Ok(1) } - pub fn delete_session_by_id(&self, session_id: SessionId) -> Result { + pub fn delete_session_by_id( + &self, + session_id: SessionId, + did: &Did, + ) -> Result { let session = match self.load_session_by_id(session_id.as_i32())? { - Some(s) => s, - None => return Ok(0), + Some(s) if s.did == did.as_str() => s, + _ => return Ok(0), }; let mut batch = self.db.batch();