session: deletes scope to did, route muts by did

Lewis: May this revision serve well! <lu5a@proton.me>
This commit is contained in:
Lewis
2026-06-28 09:55:11 +03:00
parent 9dc184ee33
commit aab1a945c2
16 changed files with 402 additions and 328 deletions
@@ -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"
}
@@ -1,14 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM session_tokens WHERE access_jti = $1",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Text"
]
},
"nullable": []
},
"hash": "847ce3c34985d0957526c87e0a20c6b4e5daae08a338f7635def682ac0689cf6"
}
@@ -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"
}
@@ -1,14 +0,0 @@
{
"db_name": "PostgreSQL",
"query": "DELETE FROM session_tokens WHERE id = $1",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Int4"
]
},
"nullable": []
},
"hash": "cf874abcb72017e775fe699a0b77ae9341355f30e4af84968ffeb9135dba745f"
}
+125 -71
View File
@@ -452,7 +452,12 @@ pub async fn delete_session(
) -> Result<Json<EmptyResponse>, 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<tranquil_db_traits::RefreshGraceLookup, tranquil_db_traits::DbError>,
) -> Option<Result<Json<RefreshSessionOutput>, 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<Json<RefreshSessionOutput>, 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<u8>,
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);
+11 -10
View File
@@ -195,6 +195,7 @@ pub enum RefreshGraceLookup {
NotUsed,
Replay(RefreshGraceReplay),
Compromised {
did: Did,
session_id: SessionId,
key_bytes: Vec<u8>,
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<Option<SessionForRefresh>, DbError>;
async fn update_session_tokens(
async fn delete_session_by_access_jti(
&self,
access_jti: &str,
did: &Did,
) -> Result<u64, DbError>;
async fn delete_session_by_id(
&self,
session_id: SessionId,
new_access_jti: &str,
new_refresh_jti: &str,
new_access_expires_at: DateTime<Utc>,
new_refresh_expires_at: DateTime<Utc>,
) -> Result<(), DbError>;
async fn delete_session_by_access_jti(&self, access_jti: &str) -> Result<u64, DbError>;
async fn delete_session_by_id(&self, session_id: SessionId) -> Result<u64, DbError>;
did: &Did,
) -> Result<u64, DbError>;
async fn delete_sessions_by_did(&self, did: &Did) -> Result<u64, DbError>;
+16 -33
View File
@@ -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<Utc>,
new_refresh_expires_at: DateTime<Utc>,
) -> 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<u64, DbError> {
access_jti: &str,
did: &Did,
) -> Result<u64, DbError> {
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<u64, DbError> {
async fn delete_session_by_id(
&self,
session_id: SessionId,
did: &Did,
) -> Result<u64, DbError> {
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),
+2 -4
View File
@@ -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"));
+1 -3
View File
@@ -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;
}
+15 -15
View File
@@ -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") {
@@ -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:*", ""), "");
}
+11 -23
View File
@@ -1463,33 +1463,16 @@ impl<S: StorageIO + 'static> 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<Utc>,
new_refresh_expires_at: DateTime<Utc>,
) -> 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<u64, DbError> {
access_jti: &str,
did: &Did,
) -> Result<u64, DbError> {
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<S: StorageIO + 'static> tranquil_db_traits::SessionRepository for Metastore
async fn delete_session_by_id(
&self,
session_id: tranquil_db_traits::SessionId,
did: &Did,
) -> Result<u64, DbError> {
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
}
@@ -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);
+96 -52
View File
@@ -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<Option<tranquil_db_traits::SessionForRefresh>>,
},
UpdateSessionTokens {
session_id: SessionId,
new_access_jti: String,
new_refresh_jti: String,
new_access_expires_at: DateTime<Utc>,
new_refresh_expires_at: DateTime<Utc>,
tx: Tx<()>,
},
DeleteSessionByAccessJti {
access_jti: String,
did: Did,
tx: Tx<u64>,
},
DeleteSessionById {
session_id: SessionId,
did: Did,
tx: Tx<u64>,
},
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<S: StorageIO>(state: &HandlerState<S>, 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 } => {
SessionRequest::DeleteSessionById {
session_id,
did,
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 } => {
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<S: StorageIO + 'static>(state: &HandlerState<S>, 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<S: StorageIO + 'static>(state: &HandlerState<S>, 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();
+67 -14
View File
@@ -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
));
@@ -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<Utc>,
new_refresh_expires_at: DateTime<Utc>,
) -> 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<u64, MetastoreError> {
access_jti: &str,
did: &Did,
) -> Result<u64, MetastoreError> {
let index_key = session_by_access_key(access_jti);
let index_val: Option<SessionIndexValue> = 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<u64, MetastoreError> {
pub fn delete_session_by_id(
&self,
session_id: SessionId,
did: &Did,
) -> Result<u64, MetastoreError> {
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();