diff --git a/.sqlx/query-6b51995c40519a63f85c70f29ca8bd6ec1963c8562d78215d980785dc46a6384.json b/.sqlx/query-18bbda5582db1b32d02ab8a3eee970c9508b9bd67239c2f936639a9f863b30ff.json similarity index 79% rename from .sqlx/query-6b51995c40519a63f85c70f29ca8bd6ec1963c8562d78215d980785dc46a6384.json rename to .sqlx/query-18bbda5582db1b32d02ab8a3eee970c9508b9bd67239c2f936639a9f863b30ff.json index cabd11f..86c323f 100644 --- a/.sqlx/query-6b51995c40519a63f85c70f29ca8bd6ec1963c8562d78215d980785dc46a6384.json +++ b/.sqlx/query-18bbda5582db1b32d02ab8a3eee970c9508b9bd67239c2f936639a9f863b30ff.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT id, did, handle, email, created_at, deactivated_at, takedown_ref, is_admin\n FROM users WHERE did = $1", + "query": "SELECT id, did, handle, email, created_at, deactivated_at, takedown_ref, is_admin, inbound_migration\n FROM users WHERE handle = $1", "describe": { "columns": [ { @@ -42,6 +42,11 @@ "ordinal": 7, "name": "is_admin", "type_info": "Bool" + }, + { + "ordinal": 8, + "name": "inbound_migration", + "type_info": "Bool" } ], "parameters": { @@ -57,8 +62,9 @@ false, true, true, + false, false ] }, - "hash": "6b51995c40519a63f85c70f29ca8bd6ec1963c8562d78215d980785dc46a6384" + "hash": "18bbda5582db1b32d02ab8a3eee970c9508b9bd67239c2f936639a9f863b30ff" } diff --git a/.sqlx/query-23201d4e26bc650939e30f69fb0bca00d351d057098afebc1017f70a84b4bd22.json b/.sqlx/query-23201d4e26bc650939e30f69fb0bca00d351d057098afebc1017f70a84b4bd22.json deleted file mode 100644 index 45565d5..0000000 --- a/.sqlx/query-23201d4e26bc650939e30f69fb0bca00d351d057098afebc1017f70a84b4bd22.json +++ /dev/null @@ -1,14 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "UPDATE users SET deactivated_at = NULL WHERE did = $1", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Text" - ] - }, - "nullable": [] - }, - "hash": "23201d4e26bc650939e30f69fb0bca00d351d057098afebc1017f70a84b4bd22" -} diff --git a/.sqlx/query-51e029e572777e6a103fd7fd5550494de9d4cac7e3ff84e27ddec1a6aaefc047.json b/.sqlx/query-51e029e572777e6a103fd7fd5550494de9d4cac7e3ff84e27ddec1a6aaefc047.json new file mode 100644 index 0000000..c07be51 --- /dev/null +++ b/.sqlx/query-51e029e572777e6a103fd7fd5550494de9d4cac7e3ff84e27ddec1a6aaefc047.json @@ -0,0 +1,14 @@ +{ + "db_name": "PostgreSQL", + "query": "UPDATE users SET deactivated_at = NULL, inbound_migration = FALSE WHERE did = $1", + "describe": { + "columns": [], + "parameters": { + "Left": [ + "Text" + ] + }, + "nullable": [] + }, + "hash": "51e029e572777e6a103fd7fd5550494de9d4cac7e3ff84e27ddec1a6aaefc047" +} diff --git a/.sqlx/query-584bceda60d8b6a02e7dc44d833e3fba13151f36ba9f75c64e33d6cb534cc939.json b/.sqlx/query-f1c4ec28b02d09ffce35aa8249c1747a70c12a3ecfc1ff6ca1847840f770db2f.json similarity index 79% rename from .sqlx/query-584bceda60d8b6a02e7dc44d833e3fba13151f36ba9f75c64e33d6cb534cc939.json rename to .sqlx/query-f1c4ec28b02d09ffce35aa8249c1747a70c12a3ecfc1ff6ca1847840f770db2f.json index db0d912..ba8c015 100644 --- a/.sqlx/query-584bceda60d8b6a02e7dc44d833e3fba13151f36ba9f75c64e33d6cb534cc939.json +++ b/.sqlx/query-f1c4ec28b02d09ffce35aa8249c1747a70c12a3ecfc1ff6ca1847840f770db2f.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT id, did, handle, email, created_at, deactivated_at, takedown_ref, is_admin\n FROM users WHERE handle = $1", + "query": "SELECT id, did, handle, email, created_at, deactivated_at, takedown_ref, is_admin, inbound_migration\n FROM users WHERE did = $1", "describe": { "columns": [ { @@ -42,6 +42,11 @@ "ordinal": 7, "name": "is_admin", "type_info": "Bool" + }, + { + "ordinal": 8, + "name": "inbound_migration", + "type_info": "Bool" } ], "parameters": { @@ -57,8 +62,9 @@ false, true, true, + false, false ] }, - "hash": "584bceda60d8b6a02e7dc44d833e3fba13151f36ba9f75c64e33d6cb534cc939" + "hash": "f1c4ec28b02d09ffce35aa8249c1747a70c12a3ecfc1ff6ca1847840f770db2f" } diff --git a/crates/tranquil-api/src/identity/account.rs b/crates/tranquil-api/src/identity/account.rs index 83ea72a..463157b 100644 --- a/crates/tranquil-api/src/identity/account.rs +++ b/crates/tranquil-api/src/identity/account.rs @@ -510,6 +510,7 @@ pub async fn create_account( telegram_username: comms.telegram, signal_username: comms.signal, deactivated_at, + inbound_migration: is_migration || is_did_web_byod, encrypted_key_bytes: repo.encrypted_key_bytes, encryption_version: tranquil_pds::config::ENCRYPTION_VERSION, reserved_key_id, diff --git a/crates/tranquil-api/src/repo/import.rs b/crates/tranquil-api/src/repo/import.rs index 87973c0..f060938 100644 --- a/crates/tranquil-api/src/repo/import.rs +++ b/crates/tranquil-api/src/repo/import.rs @@ -106,7 +106,7 @@ pub async fn import_repo( .map(|c| c.import.skip_verification) .unwrap_or(false) }); - let is_migration = user.deactivated_at.is_some(); + let is_migration = user.inbound_migration && user.deactivated_at.is_some(); if skip_verification { warn!("Skipping all CAR verification for import (SKIP_IMPORT_VERIFICATION=true)"); } else if is_migration { diff --git a/crates/tranquil-api/src/server/mod.rs b/crates/tranquil-api/src/server/mod.rs index 015f51e..4c01aae 100644 --- a/crates/tranquil-api/src/server/mod.rs +++ b/crates/tranquil-api/src/server/mod.rs @@ -51,6 +51,7 @@ pub use session::{ auto_resend_verification, confirm_signup, create_session, delete_session, get_legacy_login_preference, get_session, list_sessions, refresh_session, resend_verification, revoke_all_sessions, revoke_session, update_legacy_login_preference, update_locale, + verification_blocks_login, }; pub use signing_key::reserve_signing_key; pub use totp::{ diff --git a/crates/tranquil-api/src/server/session.rs b/crates/tranquil-api/src/server/session.rs index a8b186b..ac12d0f 100644 --- a/crates/tranquil-api/src/server/session.rs +++ b/crates/tranquil-api/src/server/session.rs @@ -8,7 +8,7 @@ use bcrypt::verify; use serde::{Deserialize, Serialize}; use serde_json::json; use tracing::{error, info, warn}; -use tranquil_db_traits::{SessionId, TokenFamilyId}; +use tranquil_db_traits::{ChannelVerificationStatus, SessionId, TokenFamilyId}; use tranquil_pds::api::error::{ApiError, DbResultExt}; use tranquil_pds::api::{EmptyResponse, PreferredLocaleOutput, SuccessResponse}; use tranquil_pds::auth::{ @@ -20,6 +20,13 @@ use tranquil_pds::state::AppState; use tranquil_pds::types::{AccountState, Did, Handle, PlainPassword}; use tranquil_types::TokenId; +pub fn verification_blocks_login(channel_verification: &ChannelVerificationStatus) -> bool { + !tranquil_config::get() + .server + .disable_account_verification_gate + && !channel_verification.has_any_verified() +} + #[derive(Deserialize)] #[serde(rename_all = "camelCase")] pub struct CreateSessionInput { @@ -129,14 +136,13 @@ pub async fn create_session( warn!("Login attempt for takendown account: {}", row.did); return Err(ApiError::AccountTakedown); } - let is_verified = row.channel_verification.has_any_verified(); let is_delegated = state .repos .delegation .is_delegated_account(&row.did) .await .unwrap_or(false); - if !is_verified && !is_delegated { + if verification_blocks_login(&row.channel_verification) && !is_delegated { warn!("Login attempt for unverified account: {}", row.did); let resend_info = auto_resend_verification(&state, &row.did).await; let handle = resend_info diff --git a/crates/tranquil-config/src/lib.rs b/crates/tranquil-config/src/lib.rs index 9897fc7..f742bd2 100644 --- a/crates/tranquil-config/src/lib.rs +++ b/crates/tranquil-config/src/lib.rs @@ -456,6 +456,11 @@ pub struct ServerConfig { #[config(env = "DISABLE_RATE_LIMITING", default = false)] pub disable_rate_limiting: bool, + /// Skip the verified-comms-channel gate for login and record writes. + /// Please keep this off unless you're an invite-only PDS! + #[config(env = "DISABLE_ACCOUNT_VERIFICATION_GATE", default = false)] + pub disable_account_verification_gate: bool, + /// List of additional banned words for handle validation. #[config(env = "PDS_BANNED_WORDS", parse_env = split_comma_list)] pub banned_words: Option>, diff --git a/crates/tranquil-db-traits/src/user.rs b/crates/tranquil-db-traits/src/user.rs index 6e8ccd4..6d7dff2 100644 --- a/crates/tranquil-db-traits/src/user.rs +++ b/crates/tranquil-db-traits/src/user.rs @@ -44,6 +44,7 @@ pub struct UserRow { pub deactivated_at: Option>, pub takedown_ref: Option, pub is_admin: bool, + pub inbound_migration: bool, } #[derive(Debug, Clone)] @@ -995,6 +996,7 @@ pub struct CreatePasswordAccountInput { pub telegram_username: Option, pub signal_username: Option, pub deactivated_at: Option>, + pub inbound_migration: bool, pub encrypted_key_bytes: Vec, pub encryption_version: i32, pub reserved_key_id: Option, diff --git a/crates/tranquil-db/src/postgres/user.rs b/crates/tranquil-db/src/postgres/user.rs index 22d94b6..ac92e52 100644 --- a/crates/tranquil-db/src/postgres/user.rs +++ b/crates/tranquil-db/src/postgres/user.rs @@ -47,7 +47,7 @@ pub(crate) fn map_sqlx_error(e: sqlx::Error) -> DbError { impl UserRepository for PostgresUserRepository { async fn get_by_did(&self, did: &Did) -> Result, DbError> { let row = sqlx::query!( - r#"SELECT id, did, handle, email, created_at, deactivated_at, takedown_ref, is_admin + r#"SELECT id, did, handle, email, created_at, deactivated_at, takedown_ref, is_admin, inbound_migration FROM users WHERE did = $1"#, did.as_str() ) @@ -64,12 +64,13 @@ impl UserRepository for PostgresUserRepository { deactivated_at: r.deactivated_at, takedown_ref: r.takedown_ref, is_admin: r.is_admin, + inbound_migration: r.inbound_migration, })) } async fn get_by_handle(&self, handle: &Handle) -> Result, DbError> { let row = sqlx::query!( - r#"SELECT id, did, handle, email, created_at, deactivated_at, takedown_ref, is_admin + r#"SELECT id, did, handle, email, created_at, deactivated_at, takedown_ref, is_admin, inbound_migration FROM users WHERE handle = $1"#, handle.as_str() ) @@ -86,6 +87,7 @@ impl UserRepository for PostgresUserRepository { deactivated_at: r.deactivated_at, takedown_ref: r.takedown_ref, is_admin: r.is_admin, + inbound_migration: r.inbound_migration, })) } @@ -1863,7 +1865,7 @@ impl UserRepository for PostgresUserRepository { async fn activate_account(&self, did: &Did) -> Result { let result = sqlx::query!( - "UPDATE users SET deactivated_at = NULL WHERE did = $1", + "UPDATE users SET deactivated_at = NULL, inbound_migration = FALSE WHERE did = $1", did.as_str() ) .execute(&self.pool) @@ -2426,8 +2428,8 @@ impl UserRepository for PostgresUserRepository { handle, email, did, password_hash, preferred_comms_channel, discord_username, telegram_username, signal_username, - is_admin, deactivated_at, email_verified - ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, FALSE) RETURNING id"#, + is_admin, deactivated_at, inbound_migration, email_verified + ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, FALSE) RETURNING id"#, ) .bind(input.handle.as_str()) .bind(&input.email) @@ -2439,6 +2441,7 @@ impl UserRepository for PostgresUserRepository { .bind(&input.signal_username) .bind(is_first_user) .bind(input.deactivated_at) + .bind(input.inbound_migration) .fetch_one(&mut *tx) .await; diff --git a/crates/tranquil-oauth-server/src/endpoints/authorize/login.rs b/crates/tranquil-oauth-server/src/endpoints/authorize/login.rs index c6c5ae6..19e1ade 100644 --- a/crates/tranquil-oauth-server/src/endpoints/authorize/login.rs +++ b/crates/tranquil-oauth-server/src/endpoints/authorize/login.rs @@ -488,8 +488,7 @@ pub async fn authorize_post( if !password_valid { return show_login_error("Invalid identifier or password.", json_response); } - let is_verified = user.channel_verification.has_any_verified(); - if !is_verified { + if tranquil_api::server::verification_blocks_login(&user.channel_verification) { let resend_info = tranquil_api::server::auto_resend_verification(&state, &user.did).await; let handle = resend_info .as_ref() @@ -854,8 +853,7 @@ pub async fn authorize_select( ); } }; - let is_verified = user.channel_verification.has_any_verified(); - if !is_verified { + if tranquil_api::server::verification_blocks_login(&user.channel_verification) { let resend_info = tranquil_api::server::auto_resend_verification(&state, &did).await; return ( StatusCode::FORBIDDEN, diff --git a/crates/tranquil-oauth-server/src/endpoints/authorize/passkey.rs b/crates/tranquil-oauth-server/src/endpoints/authorize/passkey.rs index f7684ce..6771b5f 100644 --- a/crates/tranquil-oauth-server/src/endpoints/authorize/passkey.rs +++ b/crates/tranquil-oauth-server/src/endpoints/authorize/passkey.rs @@ -289,9 +289,7 @@ async fn passkey_start_named( .into_response(); } - let is_verified = user.channel_verification.has_any_verified(); - - if !is_verified { + if tranquil_api::server::verification_blocks_login(&user.channel_verification) { let resend_info = tranquil_api::server::auto_resend_verification(&state, &user.did).await; return ( StatusCode::FORBIDDEN, diff --git a/crates/tranquil-oauth-server/src/endpoints/authorize/registration.rs b/crates/tranquil-oauth-server/src/endpoints/authorize/registration.rs index 5ccc9ae..9e0ca73 100644 --- a/crates/tranquil-oauth-server/src/endpoints/authorize/registration.rs +++ b/crates/tranquil-oauth-server/src/endpoints/authorize/registration.rs @@ -160,8 +160,10 @@ pub async fn register_complete( .into_response(); } - let is_verified = match state.repos.user.get_session_info_by_did(&did).await { - Ok(Some(info)) => info.channel_verification.has_any_verified(), + let login_blocked = match state.repos.user.get_session_info_by_did(&did).await { + Ok(Some(info)) => { + tranquil_api::server::verification_blocks_login(&info.channel_verification) + } Ok(None) => { return ( StatusCode::FORBIDDEN, @@ -189,7 +191,7 @@ pub async fn register_complete( } }; - if !is_verified { + if login_blocked { let resend_info = tranquil_api::server::auto_resend_verification(&state, &did).await; return ( StatusCode::FORBIDDEN, diff --git a/crates/tranquil-oauth-server/src/sso_endpoints.rs b/crates/tranquil-oauth-server/src/sso_endpoints.rs index 988ff26..4b1b354 100644 --- a/crates/tranquil-oauth-server/src/sso_endpoints.rs +++ b/crates/tranquil-oauth-server/src/sso_endpoints.rs @@ -402,13 +402,15 @@ async fn handle_sso_login( } }; - let is_verified = match state + let login_blocked = match state .repos .user .get_session_info_by_did(&identity.did) .await { - Ok(Some(info)) => info.channel_verification.has_any_verified(), + Ok(Some(info)) => { + tranquil_api::server::verification_blocks_login(&info.channel_verification) + } Ok(None) => { tracing::error!("User not found for SSO login: {}", identity.did); return redirect_to_error("Account not found"); @@ -419,7 +421,7 @@ async fn handle_sso_login( } }; - if !is_verified { + if login_blocked { tracing::warn!( did = %identity.did, provider = %provider.as_str(), diff --git a/crates/tranquil-pds/src/auth/account_verified.rs b/crates/tranquil-pds/src/auth/account_verified.rs index 29f2ec3..0871761 100644 --- a/crates/tranquil-pds/src/auth/account_verified.rs +++ b/crates/tranquil-pds/src/auth/account_verified.rs @@ -21,6 +21,13 @@ pub async fn require_verified_or_delegated<'a>( state: &AppState, user: &'a AuthenticatedUser, ) -> Result, ApiError> { + if tranquil_config::get() + .server + .disable_account_verification_gate + { + return Ok(AccountVerified { user }); + } + let is_verified = state .repos .user diff --git a/crates/tranquil-pds/src/state.rs b/crates/tranquil-pds/src/state.rs index 6a66128..186ffec 100644 --- a/crates/tranquil-pds/src/state.rs +++ b/crates/tranquil-pds/src/state.rs @@ -610,6 +610,17 @@ fn wire_tranquil_store( } } + if std::env::var("TRANQUIL_PURGE_ORPHAN_REPOS").is_ok_and(|v| v == "1") { + match metastore + .repo_ops() + .purge_orphan_repos(metastore.database()) + { + Ok(0) => tracing::info!("orphan repo purge: no orphans found"), + Ok(n) => tracing::info!(purged = n, "orphan repo purge: removed orphan repo_meta"), + Err(e) => tracing::error!(error = %e, "orphan repo purge failed"), + } + } + let notifier = bridge.notifier(); let signal_db = metastore.database().clone(); let signal_ks = metastore.signal_keyspace(); diff --git a/crates/tranquil-pds/tests/store_parity.rs b/crates/tranquil-pds/tests/store_parity.rs index 93a1365..7c01639 100644 --- a/crates/tranquil-pds/tests/store_parity.rs +++ b/crates/tranquil-pds/tests/store_parity.rs @@ -150,6 +150,7 @@ async fn seed_user(repos: &PostgresRepositories, did: &Did, handle: &Handle) -> telegram_username: None, signal_username: None, deactivated_at: None, + inbound_migration: false, encrypted_key_bytes: vec![0u8; 32], encryption_version: 0, reserved_key_id: None, @@ -1465,6 +1466,37 @@ async fn parity_delete_all_records() { assert_eq!(store_colls.len(), 0); } +#[tokio::test] +async fn parity_account_deletion_clears_records_on_reregister() { + let f = ParityFixture::new().await; + let did = test_did("cuttle"); + let handle = test_handle("cuttle"); + let collection = test_nsid("post"); + + let (pg_uid, store_uid) = seed_repos(&f, &did, &handle).await; + + let records: Vec<(Rkey, CidLink)> = (0u8..3) + .map(|i| (test_rkey(&format!("3l{:02}aaaaaaaaa", i)), test_cid(i + 1))) + .collect(); + seed_records(&f.pg, pg_uid, &collection, &records).await; + seed_records(&f.store, store_uid, &collection, &records).await; + + f.pg.user + .delete_account_complete(pg_uid, &did) + .await + .unwrap(); + f.store + .user + .delete_account_complete(store_uid, &did) + .await + .unwrap(); + + let (pg_uid2, store_uid2) = seed_repos(&f, &did, &handle).await; + + assert_eq!(f.pg.repo.count_records(pg_uid2).await.unwrap(), 0); + assert_eq!(f.store.repo.count_records(store_uid2).await.unwrap(), 0); +} + #[tokio::test] async fn parity_plc_tokens() { let f = ParityFixture::new().await; diff --git a/crates/tranquil-store/src/gauntlet/runner.rs b/crates/tranquil-store/src/gauntlet/runner.rs index 0aa1c1a..00b0437 100644 --- a/crates/tranquil-store/src/gauntlet/runner.rs +++ b/crates/tranquil-store/src/gauntlet/runner.rs @@ -1216,9 +1216,10 @@ pub(super) async fn apply_op( Op::ExternalDeleteDataFile { choice } => { let s = harness.store.clone(); let pick = choice.0; - let lost_cids = tokio::task::spawn_blocking(move || externally_delete_data_file(&s, pick)) - .await - .map_err(|e| OpError::Join(e.to_string()))??; + let lost_cids = + tokio::task::spawn_blocking(move || externally_delete_data_file(&s, pick)) + .await + .map_err(|e| OpError::Join(e.to_string()))??; if !lost_cids.is_empty() { oracle.mark_blocks_lost(lost_cids); } @@ -1561,9 +1562,10 @@ async fn apply_op_concurrent( let mut guard = shared.write.lock().await; let s = shared.store.clone(); let pick = choice.0; - let lost_cids = tokio::task::spawn_blocking(move || externally_delete_data_file(&s, pick)) - .await - .map_err(|e| OpError::Join(e.to_string()))??; + let lost_cids = + tokio::task::spawn_blocking(move || externally_delete_data_file(&s, pick)) + .await + .map_err(|e| OpError::Join(e.to_string()))??; if !lost_cids.is_empty() { guard.oracle.mark_blocks_lost(lost_cids); } diff --git a/crates/tranquil-store/src/metastore/handler.rs b/crates/tranquil-store/src/metastore/handler.rs index 118cdd3..7ff0c2b 100644 --- a/crates/tranquil-store/src/metastore/handler.rs +++ b/crates/tranquil-store/src/metastore/handler.rs @@ -5026,6 +5026,20 @@ fn dispatch(state: &HandlerState, request: MetastoreR } } +fn purge_repo_side_data( + state: &HandlerState, + user_id: Uuid, + did: &Did, +) -> Result<(), MetastoreError> { + let _ = state.metastore.blob_ops().delete_blobs_by_user(user_id)?; + let mut batch = state.metastore.database().batch(); + state + .metastore + .backlink_ops() + .remove_backlinks_by_repo(&mut batch, UserHash::from_did(did.as_str()))?; + batch.commit().map_err(MetastoreError::Fjall) +} + fn dispatch_user(state: &HandlerState, req: UserRequest) { let user = state.metastore.user_ops(); match req { @@ -5698,10 +5712,10 @@ fn dispatch_user(state: &HandlerState, req: UserReque let _ = tx.send(user.get_user_key_by_did(&did).map_err(metastore_to_db)); } UserRequest::DeleteAccountComplete { user_id, did, tx } => { - let _ = tx.send( - user.delete_account_complete(user_id, &did) - .map_err(metastore_to_db), - ); + let result = purge_repo_side_data(state, user_id, &did) + .and_then(|()| user.delete_account_complete(user_id, &did)) + .map_err(metastore_to_db); + let _ = tx.send(result); } UserRequest::SetUserTakedown { did, @@ -5714,10 +5728,10 @@ fn dispatch_user(state: &HandlerState, req: UserReque ); } UserRequest::AdminDeleteAccountComplete { user_id, did, tx } => { - let _ = tx.send( - user.admin_delete_account_complete(user_id, &did) - .map_err(metastore_to_db), - ); + let result = purge_repo_side_data(state, user_id, &did) + .and_then(|()| user.admin_delete_account_complete(user_id, &did)) + .map_err(metastore_to_db); + let _ = tx.send(result); } UserRequest::GetUserForDidDoc { did, tx } => { let _ = tx.send(user.get_user_for_did_doc(&did).map_err(metastore_to_db)); diff --git a/crates/tranquil-store/src/metastore/repo_meta.rs b/crates/tranquil-store/src/metastore/repo_meta.rs index 65a100f..27843fe 100644 --- a/crates/tranquil-store/src/metastore/repo_meta.rs +++ b/crates/tranquil-store/src/metastore/repo_meta.rs @@ -70,6 +70,18 @@ pub fn handle_key(handle_lower: &str) -> SmallVec<[u8; 128]> { .build() } +pub fn stage_repo_meta_removal( + batch: &mut fjall::OwnedWriteBatch, + repo_data: &fjall::Keyspace, + user_hash: UserHash, + handle: &str, +) { + batch.remove(repo_data, repo_meta_key(user_hash).as_slice()); + if !handle.is_empty() { + batch.remove(repo_data, handle_key(handle).as_slice()); + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/crates/tranquil-store/src/metastore/repo_ops.rs b/crates/tranquil-store/src/metastore/repo_ops.rs index 318b19d..a3094f1 100644 --- a/crates/tranquil-store/src/metastore/repo_ops.rs +++ b/crates/tranquil-store/src/metastore/repo_ops.rs @@ -6,8 +6,12 @@ use uuid::Uuid; use super::MetastoreError; use super::encoding::KeyReader; use super::keys::{KeyTag, UserHash}; -use super::repo_meta::{RepoMetaValue, RepoStatus, handle_key, repo_meta_key, repo_meta_prefix}; -use super::scan::{count_prefix, point_lookup}; +use super::records::record_user_prefix; +use super::repo_meta::{ + RepoMetaValue, RepoStatus, handle_key, repo_meta_key, repo_meta_prefix, stage_repo_meta_removal, +}; +use super::scan::{count_prefix, delete_all_by_prefix, point_lookup}; +use super::user_blocks::user_block_user_prefix; use super::user_hash::UserHashMap; use tranquil_types::{CidLink, Did, Handle}; @@ -241,13 +245,7 @@ impl RepoOps { let meta = self.get_meta_value(key.as_slice())?; let mut batch = db.batch(); - batch.remove(&self.repo_data, key.as_slice()); - - match meta.handle.is_empty() { - true => {} - false => batch.remove(&self.repo_data, handle_key(&meta.handle).as_slice()), - } - + stage_repo_meta_removal(&mut batch, &self.repo_data, user_hash, &meta.handle); self.user_hashes.stage_remove(&mut batch, &user_id); match batch.commit() { @@ -259,6 +257,48 @@ impl RepoOps { } } + pub fn purge_orphan_repos(&self, db: &fjall::Database) -> Result { + let prefix = repo_meta_prefix(); + let orphans: Vec<(UserHash, String)> = self + .repo_data + .prefix(prefix.as_slice()) + .map(|guard| -> Result, MetastoreError> { + let (k, v) = guard.into_inner().map_err(MetastoreError::Fjall)?; + let user_hash = parse_repo_meta_key_hash(&k) + .ok_or(MetastoreError::CorruptData("invalid repo_meta key"))?; + match self.user_hashes.get_uuid(&user_hash) { + Some(_) => Ok(None), + None => { + let handle = match RepoMetaValue::deserialize(&v) { + Some(meta) => meta.handle, + None => { + tracing::warn!( + user_hash = user_hash.raw(), + "could not deserialize orphan repo_meta to recover handle for cleanup" + ); + String::new() + } + }; + Ok(Some((user_hash, handle))) + } + } + }) + .filter_map(Result::transpose) + .collect::, _>>()?; + + match orphans.is_empty() { + true => Ok(0), + false => { + let mut batch = db.batch(); + orphans.iter().try_for_each(|(user_hash, handle)| { + stage_full_repo_data_removal(&mut batch, &self.repo_data, *user_hash, handle) + })?; + batch.commit().map_err(MetastoreError::Fjall)?; + Ok(orphans.len()) + } + } + } + pub fn get_repo(&self, user_id: Uuid) -> Result, MetastoreError> { let user_hash = match self.user_hashes.get(&user_id) { Some(h) => h, @@ -569,9 +609,26 @@ fn parse_repo_meta_key_hash(key_bytes: &[u8]) -> Option { Some(UserHash::from_raw(hash)) } +pub(super) fn stage_full_repo_data_removal( + batch: &mut fjall::OwnedWriteBatch, + repo_data: &Keyspace, + user_hash: UserHash, + handle: &str, +) -> Result<(), MetastoreError> { + stage_repo_meta_removal(batch, repo_data, user_hash, handle); + delete_all_by_prefix(repo_data, batch, record_user_prefix(user_hash).as_slice())?; + delete_all_by_prefix( + repo_data, + batch, + user_block_user_prefix(user_hash).as_slice(), + )?; + Ok(()) +} + #[cfg(test)] mod tests { use super::*; + use crate::metastore::partitions::Partition; use crate::metastore::{Metastore, MetastoreConfig}; fn test_config() -> MetastoreConfig { @@ -725,6 +782,121 @@ mod tests { assert!(ops.get_repo(uid_a).unwrap().is_none()); } + #[test] + fn purge_orphan_repos_removes_entries_with_missing_reverse_mapping() { + let (_dir, ms) = open_fresh(); + let ops = ms.repo_ops(); + let orphan_id = uuid::Uuid::new_v4(); + let orphan_did = test_did("limpet"); + let orphan_handle = test_handle("limpet"); + let live_id = uuid::Uuid::new_v4(); + let live_did = test_did("whelk"); + let live_handle = test_handle("whelk"); + let cid = test_cid_link(9); + + ops.create_repo( + ms.database(), + orphan_id, + &orphan_did, + &orphan_handle, + &cid, + "rev1", + ) + .unwrap(); + ops.create_repo( + ms.database(), + live_id, + &live_did, + &live_handle, + &cid, + "rev1", + ) + .unwrap(); + + let mut batch = ms.database().batch(); + ms.user_hashes().stage_remove(&mut batch, &orphan_id); + batch.commit().unwrap(); + + assert!(matches!( + ops.list_repos_paginated(None, 100), + Err(MetastoreError::CorruptData( + "user_hash has no reverse mapping" + )) + )); + + assert_eq!(ops.purge_orphan_repos(ms.database()).unwrap(), 1); + + let repos = ops.list_repos_paginated(None, 100).unwrap(); + assert_eq!(repos.len(), 1); + assert_eq!(repos[0].user_id, live_id); + assert!(ops.lookup_handle(&orphan_handle).unwrap().is_none()); + assert!(ops.lookup_handle(&live_handle).unwrap().is_some()); + + assert_eq!(ops.purge_orphan_repos(ms.database()).unwrap(), 0); + } + + #[test] + fn purge_orphan_repos_removes_records_and_blocks_for_orphan_only() { + let (_dir, ms) = open_fresh(); + let ops = ms.repo_ops(); + let orphan_id = uuid::Uuid::new_v4(); + let orphan_did = test_did("scallop"); + let orphan_handle = test_handle("scallop"); + let live_id = uuid::Uuid::new_v4(); + let live_did = test_did("mussel"); + let live_handle = test_handle("mussel"); + let cid = test_cid_link(3); + + ops.create_repo( + ms.database(), + orphan_id, + &orphan_did, + &orphan_handle, + &cid, + "rev1", + ) + .unwrap(); + ops.create_repo( + ms.database(), + live_id, + &live_did, + &live_handle, + &cid, + "rev1", + ) + .unwrap(); + + let orphan_hash = ms.user_hashes().get(&orphan_id).unwrap(); + let live_hash = ms.user_hashes().get(&live_id).unwrap(); + + let seed = |hash: UserHash| { + let mut batch = ms.database().batch(); + let repo_data = ms.partition(Partition::RepoData); + let mut rec_key = record_user_prefix(hash); + rec_key.extend_from_slice(b"app.bsky.feed.post/seed"); + batch.insert(repo_data, rec_key.as_slice(), b"r"); + let mut blk_key = user_block_user_prefix(hash); + blk_key.extend_from_slice(b"seed-cid"); + batch.insert(repo_data, blk_key.as_slice(), b"b"); + batch.commit().unwrap(); + }; + seed(orphan_hash); + seed(live_hash); + + let mut batch = ms.database().batch(); + ms.user_hashes().stage_remove(&mut batch, &orphan_id); + batch.commit().unwrap(); + + let count = |prefix: &[u8]| ms.partition(Partition::RepoData).prefix(prefix).count(); + + assert_eq!(ops.purge_orphan_repos(ms.database()).unwrap(), 1); + + assert_eq!(count(record_user_prefix(orphan_hash).as_slice()), 0); + assert_eq!(count(user_block_user_prefix(orphan_hash).as_slice()), 0); + assert_eq!(count(record_user_prefix(live_hash).as_slice()), 1); + assert_eq!(count(user_block_user_prefix(live_hash).as_slice()), 1); + } + #[test] fn delete_nonexistent_user_returns_error() { let (_dir, ms) = open_fresh(); diff --git a/crates/tranquil-store/src/metastore/user_ops.rs b/crates/tranquil-store/src/metastore/user_ops.rs index 73d62b6..5730633 100644 --- a/crates/tranquil-store/src/metastore/user_ops.rs +++ b/crates/tranquil-store/src/metastore/user_ops.rs @@ -8,7 +8,7 @@ use super::MetastoreError; use super::infra_schema::{channel_to_u8, u8_to_channel}; use super::keys::UserHash; use super::repo_meta::{RepoMetaValue, RepoStatus, handle_key, repo_meta_key}; -use super::repo_ops::cid_link_to_bytes; +use super::repo_ops::{cid_link_to_bytes, stage_full_repo_data_removal}; use super::scan::{count_prefix, delete_all_by_prefix, point_lookup}; use super::sessions::{SessionIndexValue, session_by_access_key}; use super::user_hash::UserHashMap; @@ -182,6 +182,7 @@ impl UserOps { .and_then(DateTime::from_timestamp_millis), takedown_ref: val.takedown_ref.clone(), is_admin: val.is_admin, + inbound_migration: val.inbound_migration, }) } @@ -2244,6 +2245,7 @@ impl UserOps { self.mutate_user(user_hash, |u| { u.deactivated_at_ms = None; u.delete_after_ms = None; + u.inbound_migration = false; }) } @@ -2421,10 +2423,27 @@ impl UserOps { let mut batch = self.db.batch(); self.delete_user_data(&mut batch, user_hash, &user)?; + self.stage_repo_data_removal(&mut batch, user_hash)?; self.user_hashes.stage_remove(&mut batch, &user_id); batch.commit().map_err(MetastoreError::Fjall) } + fn stage_repo_data_removal( + &self, + batch: &mut fjall::OwnedWriteBatch, + user_hash: UserHash, + ) -> Result<(), MetastoreError> { + let handle = point_lookup( + &self.repo_data, + repo_meta_key(user_hash).as_slice(), + RepoMetaValue::deserialize, + "invalid repo_meta value", + )? + .map(|m| m.handle) + .unwrap_or_default(); + stage_full_repo_data_removal(batch, &self.repo_data, user_hash, &handle) + } + pub fn set_user_takedown( &self, did: &Did, @@ -2657,6 +2676,7 @@ impl UserOps { account_type: AccountType, password_required: bool, is_admin: bool, + inbound_migration: bool, ) -> UserValue { let now_ms = Utc::now().timestamp_millis(); UserValue { @@ -2692,6 +2712,7 @@ impl UserOps { signal_username: signal_username.map(str::to_owned), signal_verified: false, delete_after_ms: None, + inbound_migration, } } @@ -2841,6 +2862,7 @@ impl UserOps { AccountType::Personal, true, is_admin, + input.inbound_migration, ); self.write_new_account(&user_value, &input.commit_cid, &input.repo_rev) @@ -2867,6 +2889,7 @@ impl UserOps { AccountType::Delegated, false, false, + false, ); let result = self.write_new_account(&user_value, &input.commit_cid, &input.repo_rev)?; @@ -2898,6 +2921,7 @@ impl UserOps { AccountType::Personal, false, is_admin, + false, ); let result = self.write_new_account(&user_value, &input.commit_cid, &input.repo_rev)?; @@ -2942,6 +2966,7 @@ impl UserOps { AccountType::Personal, false, is_admin, + false, ); self.write_new_account(&user_value, &input.commit_cid, &input.repo_rev) diff --git a/crates/tranquil-store/src/metastore/users.rs b/crates/tranquil-store/src/metastore/users.rs index e0e3c95..a2d0607 100644 --- a/crates/tranquil-store/src/metastore/users.rs +++ b/crates/tranquil-store/src/metastore/users.rs @@ -4,7 +4,7 @@ use smallvec::SmallVec; use super::encoding::KeyBuilder; use super::keys::{KeyTag, UserHash}; -const USER_SCHEMA_VERSION: u8 = 1; +const USER_SCHEMA_VERSION: u8 = 2; const PASSKEY_SCHEMA_VERSION: u8 = 1; const TOTP_SCHEMA_VERSION: u8 = 1; const BACKUP_CODE_SCHEMA_VERSION: u8 = 1; @@ -48,6 +48,7 @@ pub struct UserValue { pub signal_username: Option, pub signal_verified: bool, pub delete_after_ms: Option, + pub inbound_migration: bool, } impl UserValue { @@ -63,6 +64,11 @@ impl UserValue { let (&version, payload) = bytes.split_first()?; match version { USER_SCHEMA_VERSION => postcard::from_bytes(payload).ok(), + 1 => { + let mut extended = payload.to_vec(); + extended.push(0); + postcard::from_bytes(&extended).ok() + } _ => None, } } @@ -534,6 +540,7 @@ mod tests { signal_username: None, signal_verified: false, delete_after_ms: None, + inbound_migration: false, }; let bytes = val.serialize(); assert_eq!(bytes[0], USER_SCHEMA_VERSION); @@ -541,6 +548,54 @@ mod tests { assert_eq!(val, decoded); } + #[test] + fn deserialize_v1_defaults_inbound_migration_false() { + let val = UserValue { + id: uuid::Uuid::new_v4(), + did: "did:plc:squid".to_owned(), + handle: "witchcraft.systems".to_owned(), + email: Some("nel@oyster.cafe".to_owned()), + email_verified: true, + password_hash: Some("hashed".to_owned()), + created_at_ms: 1700000000000, + deactivated_at_ms: Some(1700000000000), + takedown_ref: None, + is_admin: false, + preferred_comms_channel: None, + key_bytes: vec![1, 2, 3], + encryption_version: 1, + account_type: 0, + password_required: true, + two_factor_enabled: false, + email_2fa_enabled: false, + totp_enabled: false, + allow_legacy_login: false, + preferred_locale: None, + invites_disabled: false, + migrated_to_pds: None, + migrated_at_ms: None, + discord_username: None, + discord_id: None, + discord_verified: false, + telegram_username: None, + telegram_chat_id: None, + telegram_verified: false, + signal_username: None, + signal_verified: false, + delete_after_ms: None, + inbound_migration: true, + }; + let v2 = val.serialize(); + let mut v1 = Vec::with_capacity(v2.len() - 1); + v1.push(1); + v1.extend_from_slice(&v2[1..v2.len() - 1]); + let decoded = UserValue::deserialize(&v1).expect("v1 user record must still decode"); + assert!(!decoded.inbound_migration); + assert_eq!(decoded.did, val.did); + assert_eq!(decoded.handle, val.handle); + assert_eq!(decoded.deactivated_at_ms, val.deactivated_at_ms); + } + #[test] fn passkey_value_roundtrip() { let val = PasskeyValue { @@ -774,6 +829,7 @@ mod tests { signal_username: None, signal_verified: false, delete_after_ms: None, + inbound_migration: false, }; assert_eq!(user.channel_verification(), 0); user.email_verified = true; @@ -821,6 +877,7 @@ mod tests { signal_username: None, signal_verified: false, delete_after_ms: None, + inbound_migration: false, }; let mut bytes = val.serialize(); bytes[0] = 99; diff --git a/crates/tranquil-store/tests/gauntlet_index_backed.rs b/crates/tranquil-store/tests/gauntlet_index_backed.rs index cedf448..00d9747 100644 --- a/crates/tranquil-store/tests/gauntlet_index_backed.rs +++ b/crates/tranquil-store/tests/gauntlet_index_backed.rs @@ -76,27 +76,28 @@ fn index_backed_by_disk_invariant_catches_phantom_after_external_delete() { #[tokio::test] async fn external_corruption_scenario_survives_many_seeds() { - let failures: Vec = futures::future::join_all((0..5).map(Seed).map(|seed| async move { - let cfg = config_for(Scenario::ExternalCorruption, seed); - let report = Gauntlet::new(cfg).expect("build gauntlet").run().await; - (seed, report) - })) - .await - .into_iter() - .filter(|(_, r)| !r.is_clean()) - .map(|(seed, r)| { - format!( - "seed {}: {} violations\n {}", - seed.0, - r.violations.len(), - r.violations - .iter() - .map(|v| format!("{}: {}", v.invariant, v.detail)) - .collect::>() - .join("\n ") - ) - }) - .collect(); + let failures: Vec = + futures::future::join_all((0..5).map(Seed).map(|seed| async move { + let cfg = config_for(Scenario::ExternalCorruption, seed); + let report = Gauntlet::new(cfg).expect("build gauntlet").run().await; + (seed, report) + })) + .await + .into_iter() + .filter(|(_, r)| !r.is_clean()) + .map(|(seed, r)| { + format!( + "seed {}: {} violations\n {}", + seed.0, + r.violations.len(), + r.violations + .iter() + .map(|v| format!("{}: {}", v.invariant, v.detail)) + .collect::>() + .join("\n ") + ) + }) + .collect(); assert!(failures.is_empty(), "{}", failures.join("\n---\n")); } diff --git a/example.toml b/example.toml index 1cc1012..d3ad2bb 100644 --- a/example.toml +++ b/example.toml @@ -62,6 +62,14 @@ # Default value: false #disable_rate_limiting = false +# Skip the verified-comms-channel gate for login and record writes. +# Please keep this off unless you're an invite-only PDS! +# +# Can also be specified via environment variable `DISABLE_ACCOUNT_VERIFICATION_GATE`. +# +# Default value: false +#disable_account_verification_gate = false + # List of additional banned words for handle validation. # # Can also be specified via environment variable `PDS_BANNED_WORDS`. diff --git a/justfile b/justfile index 1c57ee5..87c398e 100644 --- a/justfile +++ b/justfile @@ -7,6 +7,8 @@ run-dev: docker compose --profile dev up run-release: cargo run -p tranquil-server --release +gen-config: + cargo run -p tranquil-server -- config-template > example.toml build: cargo build build-release: diff --git a/migrations/20260523_add_inbound_migration_flag.sql b/migrations/20260523_add_inbound_migration_flag.sql new file mode 100644 index 0000000..d8d0b6f --- /dev/null +++ b/migrations/20260523_add_inbound_migration_flag.sql @@ -0,0 +1 @@ +ALTER TABLE users ADD COLUMN inbound_migration BOOLEAN NOT NULL DEFAULT FALSE;