diff --git a/migrations/postgres/20260721000000_drop_repo_state_sig.sql b/migrations/postgres/20260721000000_drop_repo_state_sig.sql new file mode 100644 index 0000000..bfe7b1f --- /dev/null +++ b/migrations/postgres/20260721000000_drop_repo_state_sig.sql @@ -0,0 +1,11 @@ +-- Drop the asymmetric commit signature from space repo state. +-- +-- Permissioned Data Diary 7 establishes that asymmetric signatures on data are +-- an anti-pattern for permissioned data: a signature is an irrevocable, +-- redistributable proof of authenticity that survives escaping the space's +-- access boundary. Space commits are now authenticated solely by the deniable +-- HMAC (`mac`, keyed from the per-commit `ikm`). +-- +-- The column was never populated in production — the commit machinery had no +-- production callers until the write path was wired up alongside this change. +ALTER TABLE happyview_space_repo_state DROP COLUMN sig; diff --git a/migrations/sqlite/20260721000000_drop_repo_state_sig.sql b/migrations/sqlite/20260721000000_drop_repo_state_sig.sql new file mode 100644 index 0000000..bfe7b1f --- /dev/null +++ b/migrations/sqlite/20260721000000_drop_repo_state_sig.sql @@ -0,0 +1,11 @@ +-- Drop the asymmetric commit signature from space repo state. +-- +-- Permissioned Data Diary 7 establishes that asymmetric signatures on data are +-- an anti-pattern for permissioned data: a signature is an irrevocable, +-- redistributable proof of authenticity that survives escaping the space's +-- access boundary. Space commits are now authenticated solely by the deniable +-- HMAC (`mac`, keyed from the per-commit `ikm`). +-- +-- The column was never populated in production — the commit machinery had no +-- production callers until the write path was wired up alongside this change. +ALTER TABLE happyview_space_repo_state DROP COLUMN sig; diff --git a/src/jobs/worker.rs b/src/jobs/worker.rs index 273c69e..5fb921f 100644 --- a/src/jobs/worker.rs +++ b/src/jobs/worker.rs @@ -364,23 +364,59 @@ async fn execute_job(state: &AppState, job: &super::Job) { } Err(e) => { let error = format!("{e}"); - tracing::error!(job_id = %job.id, %error, "job script failed"); - let _ = db::set_error(state, &job.id, &error).await; - log_event( - &state.db, - EventLog { - event_type: "job.failed".to_string(), - severity: Severity::Error, - actor_did: Some(job.created_by.clone()), - subject: Some(job.job_type.clone()), - detail: serde_json::json!({ - "job_id": job.id, - "error": error, - }), - }, - backend, - ) - .await; + match db::should_stop(state, &job.id).await { + Some("pausing") => { + let _ = db::set_status(state, &job.id, "paused").await; + tracing::info!(job_id = %job.id, "job paused (script error during stop: {error})"); + log_event( + &state.db, + EventLog { + event_type: "job.paused".to_string(), + severity: Severity::Info, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ "job_id": job.id }), + }, + backend, + ) + .await; + } + Some("cancelling") => { + let _ = db::set_status(state, &job.id, "cancelled").await; + tracing::info!(job_id = %job.id, "job cancelled (script error during stop: {error})"); + log_event( + &state.db, + EventLog { + event_type: "job.cancelled".to_string(), + severity: Severity::Info, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ "job_id": job.id }), + }, + backend, + ) + .await; + } + _ => { + tracing::error!(job_id = %job.id, %error, "job script failed"); + let _ = db::set_error(state, &job.id, &error).await; + log_event( + &state.db, + EventLog { + event_type: "job.failed".to_string(), + severity: Severity::Error, + actor_did: Some(job.created_by.clone()), + subject: Some(job.job_type.clone()), + detail: serde_json::json!({ + "job_id": job.id, + "error": error, + }), + }, + backend, + ) + .await; + } + } } } } diff --git a/src/spaces/car.rs b/src/spaces/car.rs index e4156ba..09b2b7f 100644 --- a/src/spaces/car.rs +++ b/src/spaces/car.rs @@ -100,10 +100,6 @@ pub fn serialize_repo(commit: &SignedCommit, records: &[SpaceRecord]) -> Result< ciborium::Value::Text("ikm".into()), ciborium::Value::Bytes(commit.ikm.to_vec()), ), - ( - ciborium::Value::Text("sig".into()), - ciborium::Value::Bytes(commit.sig.clone()), - ), ( ciborium::Value::Text("mac".into()), ciborium::Value::Bytes(commit.mac.to_vec()), @@ -159,7 +155,6 @@ mod tests { ver: 1, hash: [0u8; 32], ikm: [0u8; 32], - sig: vec![0u8; 64], mac: [0u8; 32], rev: "3k2rev1".to_string(), } @@ -180,7 +175,6 @@ mod tests { ver: 1, hash: [0xAA; 32], ikm: [0xBB; 32], - sig: vec![0xCC; 64], mac: [0xDD; 32], rev: "3k2rev1".to_string(), }; diff --git a/src/spaces/commit.rs b/src/spaces/commit.rs index e7e90b1..5fec04b 100644 --- a/src/spaces/commit.rs +++ b/src/spaces/commit.rs @@ -1,6 +1,5 @@ use hkdf::Hkdf; use hmac::{Hmac, KeyInit, Mac}; -use k256::ecdsa::{Signature, SigningKey, VerifyingKey, signature::Signer, signature::Verifier}; use rand::Rng; use sha2::Sha256; @@ -10,7 +9,6 @@ pub struct SignedCommit { pub ver: u32, pub hash: [u8; 32], pub ikm: [u8; 32], - pub sig: Vec, pub mac: [u8; 32], pub rev: String, } @@ -48,16 +46,12 @@ pub fn sign_commit( space_uri: &str, author_did: &str, rev: &str, - signing_key: &SigningKey, ) -> Result { let mut ikm = [0u8; 32]; rand::rng().fill_bytes(&mut ikm); let ctx = build_context(space_uri, author_did, rev, &ikm); - // sig covers space + author + rev + ikm, NOT the hash — prevents rebroadcast proof - let sig: Signature = signing_key.sign(&ctx); - // mac = HMAC-SHA256(HKDF-SHA256(ikm, ctx), hash) let hk = Hkdf::::new(None, &ikm); let mut derived_key = [0u8; 32]; @@ -73,7 +67,6 @@ pub fn sign_commit( ver: 1, hash: *hash, ikm, - sig: sig.to_bytes().to_vec(), mac, rev: rev.to_string(), }) @@ -83,7 +76,6 @@ pub fn verify_commit( commit: &SignedCommit, space_uri: &str, author_did: &str, - verifying_key: &VerifyingKey, ) -> Result<(), AppError> { if commit.ver != 1 { return Err(AppError::BadRequest(format!( @@ -94,12 +86,6 @@ pub fn verify_commit( let ctx = build_context(space_uri, author_did, &commit.rev, &commit.ikm); - let sig = Signature::from_slice(commit.sig.as_slice()) - .map_err(|_| AppError::Auth("invalid commit signature format".into()))?; - verifying_key - .verify(&ctx, &sig) - .map_err(|_| AppError::Auth("commit signature verification failed".into()))?; - // Recompute and verify MAC let hk = Hkdf::::new(None, &commit.ikm); let mut derived_key = [0u8; 32]; @@ -120,13 +106,6 @@ pub fn verify_commit( #[cfg(test)] mod tests { use super::*; - use k256::ecdsa::SigningKey; - - fn test_signing_key() -> SigningKey { - let mut bytes = [0u8; 32]; - bytes[31] = 1; // valid non-zero scalar - SigningKey::from_slice(&bytes[..]).unwrap() - } #[test] fn context_string_format() { @@ -185,74 +164,69 @@ mod tests { #[test] fn sign_and_verify_roundtrip() { - let sk = test_signing_key(); - let vk = *sk.verifying_key(); let hash = [0xCC; 32]; let space = "at://did:plc:abc/space/com.example.forum/main"; - let commit = sign_commit(&hash, space, "did:plc:testuser", "3k2rev1", &sk).unwrap(); + let commit = sign_commit(&hash, space, "did:plc:testuser", "3k2rev1").unwrap(); assert_eq!(commit.hash, hash); assert_eq!(commit.rev, "3k2rev1"); assert_eq!(commit.mac.len(), 32); - assert!(!commit.sig.is_empty()); assert_eq!(commit.ver, 1); - assert!(verify_commit(&commit, space, "did:plc:testuser", &vk).is_ok()); + assert!(verify_commit(&commit, space, "did:plc:testuser").is_ok()); } #[test] fn commit_has_version() { - let sk = test_signing_key(); let hash = [0xCC; 32]; let space = "at://did:plc:abc/space/com.example.forum/main"; - let commit = sign_commit(&hash, space, "did:plc:testuser", "rev1", &sk).unwrap(); + let commit = sign_commit(&hash, space, "did:plc:testuser", "rev1").unwrap(); assert_eq!(commit.ver, 1); } #[test] fn verify_rejects_wrong_author() { - let sk = test_signing_key(); - let vk = *sk.verifying_key(); let hash = [0xAA; 32]; let space = "at://did:plc:abc/space/com.example.forum/main"; - let commit = sign_commit(&hash, space, "did:plc:user1", "rev1", &sk).unwrap(); - assert!(verify_commit(&commit, space, "did:plc:user1", &vk).is_ok()); - assert!(verify_commit(&commit, space, "did:plc:user2", &vk).is_err()); + let commit = sign_commit(&hash, space, "did:plc:user1", "rev1").unwrap(); + assert!(verify_commit(&commit, space, "did:plc:user1").is_ok()); + assert!(verify_commit(&commit, space, "did:plc:user2").is_err()); } #[test] - fn verify_rejects_wrong_key() { - let sk1 = test_signing_key(); - let mut bytes2 = [0u8; 32]; - bytes2[31] = 2; - let sk2 = SigningKey::from_slice(&bytes2[..]).unwrap(); - let vk2 = *sk2.verifying_key(); - - let hash = [0xDD; 32]; + fn verify_rejects_tampered_hash() { + let hash = [0xEE; 32]; let space = "at://did:plc:abc/space/com.example.forum/main"; - let commit = sign_commit(&hash, space, "did:plc:testuser", "rev1", &sk1).unwrap(); - assert!(verify_commit(&commit, space, "did:plc:testuser", &vk2).is_err()); + let mut commit = sign_commit(&hash, space, "did:plc:testuser", "rev1").unwrap(); + commit.hash[0] ^= 0xFF; // tamper + assert!(verify_commit(&commit, space, "did:plc:testuser").is_err()); } #[test] - fn verify_rejects_tampered_hash() { - let sk = test_signing_key(); - let vk = *sk.verifying_key(); + fn verify_rejects_tampered_ikm() { let hash = [0xEE; 32]; let space = "at://did:plc:abc/space/com.example.forum/main"; - let mut commit = sign_commit(&hash, space, "did:plc:testuser", "rev1", &sk).unwrap(); - commit.hash[0] ^= 0xFF; // tamper - assert!(verify_commit(&commit, space, "did:plc:testuser", &vk).is_err()); + let mut commit = sign_commit(&hash, space, "did:plc:testuser", "rev1").unwrap(); + commit.ikm[0] ^= 0xFF; // tamper — changes both the HKDF salt and the context + assert!(verify_commit(&commit, space, "did:plc:testuser").is_err()); + } + + #[test] + fn verify_rejects_tampered_rev() { + let hash = [0xEE; 32]; + let space = "at://did:plc:abc/space/com.example.forum/main"; + + let mut commit = sign_commit(&hash, space, "did:plc:testuser", "rev1").unwrap(); + commit.rev = "rev2".into(); + assert!(verify_commit(&commit, space, "did:plc:testuser").is_err()); } #[test] fn verify_rejects_wrong_space() { - let sk = test_signing_key(); - let vk = *sk.verifying_key(); let hash = [0xFF; 32]; let commit = sign_commit( @@ -260,7 +234,6 @@ mod tests { "at://did:plc:abc/space/com.example.forum/main", "did:plc:user", "rev1", - &sk, ) .unwrap(); assert!( @@ -268,7 +241,6 @@ mod tests { &commit, "at://did:plc:xyz/space/com.example.forum/other", "did:plc:user", - &vk ) .is_err() ); @@ -276,43 +248,38 @@ mod tests { #[test] fn different_ikm_per_commit() { - let sk = test_signing_key(); let hash = [0xAA; 32]; let space = "at://did:plc:abc/space/com.example.forum/main"; - let c1 = sign_commit(&hash, space, "did:plc:testuser", "rev1", &sk).unwrap(); - let c2 = sign_commit(&hash, space, "did:plc:testuser", "rev1", &sk).unwrap(); + let c1 = sign_commit(&hash, space, "did:plc:testuser", "rev1").unwrap(); + let c2 = sign_commit(&hash, space, "did:plc:testuser", "rev1").unwrap(); // Each call generates fresh ikm assert_ne!(c1.ikm, c2.ikm); + assert_ne!(c1.mac, c2.mac); // But both verify - let vk = *sk.verifying_key(); - assert!(verify_commit(&c1, space, "did:plc:testuser", &vk).is_ok()); - assert!(verify_commit(&c2, space, "did:plc:testuser", &vk).is_ok()); + assert!(verify_commit(&c1, space, "did:plc:testuser").is_ok()); + assert!(verify_commit(&c2, space, "did:plc:testuser").is_ok()); } #[test] fn verify_rejects_unknown_version() { - let sk = test_signing_key(); - let vk = *sk.verifying_key(); let hash = [0xCC; 32]; let space = "at://did:plc:abc/space/com.example.forum/main"; - let mut commit = sign_commit(&hash, space, "did:plc:testuser", "rev1", &sk).unwrap(); + let mut commit = sign_commit(&hash, space, "did:plc:testuser", "rev1").unwrap(); commit.ver = 2; - assert!(verify_commit(&commit, space, "did:plc:testuser", &vk).is_err()); + assert!(verify_commit(&commit, space, "did:plc:testuser").is_err()); } #[test] fn verify_rejects_tampered_mac() { - let sk = test_signing_key(); - let vk = *sk.verifying_key(); let hash = [0xCC; 32]; let space = "at://did:plc:abc/space/com.example.forum/main"; - let mut commit = sign_commit(&hash, space, "did:plc:testuser", "rev1", &sk).unwrap(); - assert!(verify_commit(&commit, space, "did:plc:testuser", &vk).is_ok()); + let mut commit = sign_commit(&hash, space, "did:plc:testuser", "rev1").unwrap(); + assert!(verify_commit(&commit, space, "did:plc:testuser").is_ok()); commit.mac[0] ^= 0xFF; - assert!(verify_commit(&commit, space, "did:plc:testuser", &vk).is_err()); + assert!(verify_commit(&commit, space, "did:plc:testuser").is_err()); } } diff --git a/src/spaces/db.rs b/src/spaces/db.rs index 77ef7de..4a22729 100644 --- a/src/spaces/db.rs +++ b/src/spaces/db.rs @@ -428,7 +428,7 @@ fn parse_member_row(r: MemberRow) -> Result { // --------------------------------------------------------------------------- pub async fn upsert_space_record( - pool: &sqlx::AnyPool, + executor: impl sqlx::Executor<'_, Database = sqlx::Any>, backend: DatabaseBackend, record: &SpaceRecord, ) -> Result<(), AppError> { @@ -455,7 +455,7 @@ pub async fn upsert_space_record( .bind(&record_json) .bind(&record.cid) .bind(&now) - .execute(pool) + .execute(executor) .await .map_err(|e| AppError::Internal(format!("failed to upsert space record: {e}")))?; @@ -463,7 +463,7 @@ pub async fn upsert_space_record( } pub async fn get_space_record( - pool: &sqlx::AnyPool, + executor: impl sqlx::Executor<'_, Database = sqlx::Any>, backend: DatabaseBackend, uri: &str, ) -> Result, AppError> { @@ -474,7 +474,7 @@ pub async fn get_space_record( let row: Option = crate::db::query_as(&sql) .bind(uri) - .fetch_optional(pool) + .fetch_optional(executor) .await .map_err(|e| AppError::Internal(format!("failed to get space record: {e}")))?; @@ -591,7 +591,7 @@ pub async fn list_all_space_records( } pub async fn insert_space_record( - pool: &sqlx::AnyPool, + executor: impl sqlx::Executor<'_, Database = sqlx::Any>, backend: DatabaseBackend, record: &SpaceRecord, ) -> Result<(), AppError> { @@ -613,7 +613,7 @@ pub async fn insert_space_record( .bind(&record_json) .bind(&record.cid) .bind(&now) - .execute(pool) + .execute(executor) .await .map_err(|e| { let msg = e.to_string(); @@ -628,7 +628,7 @@ pub async fn insert_space_record( } pub async fn upsert_space_record_with_swap( - pool: &sqlx::AnyPool, + conn: &mut sqlx::AnyConnection, backend: DatabaseBackend, record: &SpaceRecord, swap_cid: &str, @@ -648,12 +648,12 @@ pub async fn upsert_space_record_with_swap( .bind(&now) .bind(&record.uri) .bind(swap_cid) - .execute(pool) + .execute(&mut *conn) .await .map_err(|e| AppError::Internal(format!("failed to update space record: {e}")))?; if result.rows_affected() == 0 { - let existing = get_space_record(pool, backend, &record.uri).await?; + let existing = get_space_record(&mut *conn, backend, &record.uri).await?; if existing.is_some() { return Err(AppError::Conflict("Record CID mismatch".into())); } @@ -664,7 +664,7 @@ pub async fn upsert_space_record_with_swap( } pub async fn delete_space_record( - pool: &sqlx::AnyPool, + executor: impl sqlx::Executor<'_, Database = sqlx::Any>, backend: DatabaseBackend, uri: &str, ) -> Result { @@ -672,7 +672,7 @@ pub async fn delete_space_record( let result = crate::db::query(&sql) .bind(uri) - .execute(pool) + .execute(executor) .await .map_err(|e| AppError::Internal(format!("failed to delete space record: {e}")))?; @@ -680,7 +680,7 @@ pub async fn delete_space_record( } pub async fn delete_space_record_with_swap( - pool: &sqlx::AnyPool, + conn: &mut sqlx::AnyConnection, backend: DatabaseBackend, uri: &str, swap_cid: &str, @@ -693,12 +693,12 @@ pub async fn delete_space_record_with_swap( let result = crate::db::query(&sql) .bind(uri) .bind(swap_cid) - .execute(pool) + .execute(&mut *conn) .await .map_err(|e| AppError::Internal(format!("failed to delete space record: {e}")))?; if result.rows_affected() == 0 { - let existing = get_space_record(pool, backend, uri).await?; + let existing = get_space_record(&mut *conn, backend, uri).await?; if existing.is_some() { return Err(AppError::Conflict("Record CID mismatch".into())); } @@ -709,7 +709,7 @@ pub async fn delete_space_record_with_swap( } pub async fn update_space_revision( - pool: &sqlx::AnyPool, + executor: impl sqlx::Executor<'_, Database = sqlx::Any>, backend: DatabaseBackend, space_id: &str, revision: &str, @@ -724,7 +724,7 @@ pub async fn update_space_revision( .bind(revision) .bind(&now) .bind(space_id) - .execute(pool) + .execute(executor) .await .map_err(|e| AppError::Internal(format!("failed to update space revision: {e}")))?; @@ -736,20 +736,20 @@ pub async fn update_space_revision( // --------------------------------------------------------------------------- pub async fn get_or_create_repo_state( - pool: &sqlx::AnyPool, + conn: &mut sqlx::AnyConnection, backend: DatabaseBackend, space_id: &str, author_did: &str, ) -> Result { let sql = adapt_sql( - "SELECT id, space_id, author_did, lthash_state, rev, hash, ikm, sig, mac, updated_at FROM happyview_space_repo_state WHERE space_id = ? AND author_did = ?", + "SELECT id, space_id, author_did, lthash_state, rev, hash, ikm, mac, updated_at FROM happyview_space_repo_state WHERE space_id = ? AND author_did = ?", backend, ); let row: Option = crate::db::query_as(&sql) .bind(space_id) .bind(author_did) - .fetch_optional(pool) + .fetch_optional(&mut *conn) .await .map_err(|e| AppError::Internal(format!("failed to get repo state: {e}")))?; @@ -761,7 +761,7 @@ pub async fn get_or_create_repo_state( let now = now_rfc3339(); let default_lthash = vec![0u8; 2048]; let insert_sql = adapt_sql( - "INSERT INTO happyview_space_repo_state (id, space_id, author_did, lthash_state, rev, hash, ikm, sig, mac, updated_at) VALUES (?, ?, ?, ?, NULL, NULL, NULL, NULL, NULL, ?)", + "INSERT INTO happyview_space_repo_state (id, space_id, author_did, lthash_state, rev, hash, ikm, mac, updated_at) VALUES (?, ?, ?, ?, NULL, NULL, NULL, NULL, ?)", backend, ); crate::db::query(&insert_sql) @@ -770,7 +770,7 @@ pub async fn get_or_create_repo_state( .bind(author_did) .bind(&default_lthash) .bind(&now) - .execute(pool) + .execute(&mut *conn) .await .map_err(|e| AppError::Internal(format!("failed to create repo state: {e}")))?; @@ -782,20 +782,19 @@ pub async fn get_or_create_repo_state( rev: None, hash: None, ikm: None, - sig: None, mac: None, updated_at: now, }) } pub async fn update_repo_state( - pool: &sqlx::AnyPool, + executor: impl sqlx::Executor<'_, Database = sqlx::Any>, backend: DatabaseBackend, state: &RepoState, ) -> Result<(), AppError> { let now = now_rfc3339(); let sql = adapt_sql( - "UPDATE happyview_space_repo_state SET lthash_state = ?, rev = ?, hash = ?, ikm = ?, sig = ?, mac = ?, updated_at = ? WHERE id = ?", + "UPDATE happyview_space_repo_state SET lthash_state = ?, rev = ?, hash = ?, ikm = ?, mac = ?, updated_at = ? WHERE id = ?", backend, ); @@ -804,11 +803,10 @@ pub async fn update_repo_state( .bind(&state.rev) .bind(&state.hash) .bind(&state.ikm) - .bind(&state.sig) .bind(&state.mac) .bind(&now) .bind(&state.id) - .execute(pool) + .execute(executor) .await .map_err(|e| AppError::Internal(format!("failed to update repo state: {e}")))?; @@ -824,7 +822,6 @@ type RepoStateRow = ( Option>, Option>, Option>, - Option>, String, ); @@ -837,9 +834,8 @@ fn parse_repo_state_row(r: RepoStateRow) -> Result { rev: r.4, hash: r.5, ikm: r.6, - sig: r.7, - mac: r.8, - updated_at: r.9, + mac: r.7, + updated_at: r.8, }) } diff --git a/src/spaces/integration_tests.rs b/src/spaces/integration_tests.rs index 84914e6..2eb6efc 100644 --- a/src/spaces/integration_tests.rs +++ b/src/spaces/integration_tests.rs @@ -9,13 +9,6 @@ mod tests { use crate::spaces::commit::{sign_commit, verify_commit}; use crate::spaces::lthash::{LtHashState, record_element}; - use k256::ecdsa::SigningKey; - - fn test_signing_key() -> SigningKey { - let mut bytes = [0u8; 32]; - bytes[31] = 1; - SigningKey::from_slice(&bytes[..]).unwrap() - } /// Add two records, generate a commit over the hash, verify it. #[test] @@ -39,15 +32,13 @@ mod tests { "hash must change after second add" ); - let sk = test_signing_key(); - let vk = *sk.verifying_key(); let space_uri = "at://did:plc:abc/space/com.example.forum/main"; let rev = "3k2rev1"; - let commit = sign_commit(&hash_after_ab, space_uri, "did:plc:testuser", rev, &sk).unwrap(); + let commit = sign_commit(&hash_after_ab, space_uri, "did:plc:testuser", rev).unwrap(); assert_eq!(commit.hash, hash_after_ab); assert_eq!(commit.rev, rev); - assert!(verify_commit(&commit, space_uri, "did:plc:testuser", &vk).is_ok()); + assert!(verify_commit(&commit, space_uri, "did:plc:testuser").is_ok()); } /// Remove a record — hash must change back toward the previous state. @@ -74,11 +65,9 @@ mod tests { "hash after delete must match single-record state" ); - let sk = test_signing_key(); - let vk = *sk.verifying_key(); let space_uri = "at://did:plc:abc/space/com.example.forum/main"; - let commit = sign_commit(&hash_one, space_uri, "did:plc:testuser", "3k2rev2", &sk).unwrap(); - assert!(verify_commit(&commit, space_uri, "did:plc:testuser", &vk).is_ok()); + let commit = sign_commit(&hash_one, space_uri, "did:plc:testuser", "3k2rev2").unwrap(); + assert!(verify_commit(&commit, space_uri, "did:plc:testuser").is_ok()); } /// Commit signed for one hash must not verify against a different hash. @@ -92,15 +81,13 @@ mod tests { state_b.add(&record_element("col", "key2", "cid2")); let hash_b = state_b.hash(); - let sk = test_signing_key(); - let vk = *sk.verifying_key(); let space_uri = "at://did:plc:abc/space/com.example.forum/main"; - let commit_a = sign_commit(&hash_a, space_uri, "did:plc:testuser", "rev1", &sk).unwrap(); + let commit_a = sign_commit(&hash_a, space_uri, "did:plc:testuser", "rev1").unwrap(); // Tamper: swap in hash_b let mut tampered = commit_a; tampered.hash = hash_b; - assert!(verify_commit(&tampered, space_uri, "did:plc:testuser", &vk).is_err()); + assert!(verify_commit(&tampered, space_uri, "did:plc:testuser").is_err()); } // ----------------------------------------------------------------------- diff --git a/src/spaces/oplog.rs b/src/spaces/oplog.rs index 3c9299d..7099003 100644 --- a/src/spaces/oplog.rs +++ b/src/spaces/oplog.rs @@ -3,7 +3,7 @@ use crate::error::AppError; use crate::spaces::types::{OplogAction, OplogEntry}; pub async fn append_op( - pool: &sqlx::AnyPool, + executor: impl sqlx::Executor<'_, Database = sqlx::Any>, backend: DatabaseBackend, entry: &OplogEntry, ) -> Result<(), AppError> { @@ -23,23 +23,69 @@ pub async fn append_op( .bind(&entry.cid) .bind(&entry.prev) .bind(&entry.created_at) - .execute(pool) + .execute(executor) .await .map_err(|e| AppError::Internal(format!("failed to append oplog entry: {e}")))?; Ok(()) } +pub fn encode_cursor(rev: &str, idx: i32) -> String { + format!("{rev}:{idx}") +} + +fn parse_cursor(cursor: &str) -> (&str, i32) { + match cursor.rsplit_once(':') { + Some((rev, idx)) => match idx.parse::() { + Ok(idx) => (rev, idx), + Err(_) => (cursor, i32::MAX), + }, + None => (cursor, i32::MAX), + } +} + +type OplogRow = ( + String, + String, + String, + String, + i32, + String, + String, + String, + Option, + Option, + String, +); + +fn paginate( + mut rows: Vec, + limit: i64, + key: impl Fn(&T) -> (String, i32), +) -> (Vec, Option) { + let has_more = rows.len() as i64 > limit; + rows.truncate(limit.max(0) as usize); + let next = if has_more { + rows.last().map(|r| { + let (rev, idx) = key(r); + encode_cursor(&rev, idx) + }) + } else { + None + }; + (rows, next) +} + pub async fn list_ops( - pool: &sqlx::AnyPool, + executor: impl sqlx::Executor<'_, Database = sqlx::Any>, backend: DatabaseBackend, space_id: &str, author_did: &str, - since_rev: Option<&str>, + cursor: Option<&str>, limit: i64, -) -> Result, AppError> { - let sql = if since_rev.is_some() { +) -> Result<(Vec, Option), AppError> { + let sql = if cursor.is_some() { adapt_sql( - "SELECT id, space_id, author_did, rev, idx, action, collection, rkey, cid, prev, created_at FROM happyview_space_record_oplog WHERE space_id = ? AND author_did = ? AND rev > ? ORDER BY rev, idx LIMIT ?", + "SELECT id, space_id, author_did, rev, idx, action, collection, rkey, cid, prev, created_at FROM happyview_space_record_oplog WHERE space_id = ? AND author_did = ? AND (rev > ? OR (rev = ? AND idx > ?)) ORDER BY rev, idx LIMIT ?", backend, ) } else { @@ -49,34 +95,25 @@ pub async fn list_ops( ) }; - type OplogRow = ( - String, - String, - String, - String, - i32, - String, - String, - String, - Option, - Option, - String, - ); - let mut query = crate::db::query_as::(&sql) .bind(space_id) .bind(author_did); - if let Some(rev) = since_rev { - query = query.bind(rev); + if let Some(c) = cursor { + let (rev, idx) = parse_cursor(c); + query = query.bind(rev.to_string()).bind(rev.to_string()).bind(idx); } - query = query.bind(limit); + // Over-fetch by one to detect whether another page exists. + query = query.bind(limit + 1); let rows = query - .fetch_all(pool) + .fetch_all(executor) .await .map_err(|e| AppError::Internal(format!("failed to list oplog entries: {e}")))?; - rows.into_iter() + let (rows, next_cursor) = paginate(rows, limit, |r| (r.3.clone(), r.4)); + + let entries = rows + .into_iter() .map(|r| { let action = OplogAction::parse(&r.5) .ok_or_else(|| AppError::Internal(format!("invalid oplog action: {}", r.5)))?; @@ -95,20 +132,22 @@ pub async fn list_ops( created_at: r.10, }) }) - .collect() + .collect::, AppError>>()?; + + Ok((entries, next_cursor)) } pub async fn list_ops_with_values( - pool: &sqlx::AnyPool, + executor: impl sqlx::Executor<'_, Database = sqlx::Any>, backend: DatabaseBackend, space_id: &str, author_did: &str, - since_rev: Option<&str>, + cursor: Option<&str>, limit: i64, -) -> Result, AppError> { - let sql = if since_rev.is_some() { +) -> Result<(Vec, Option), AppError> { + let sql = if cursor.is_some() { adapt_sql( - "SELECT o.id, o.space_id, o.author_did, o.rev, o.idx, o.action, o.collection, o.rkey, o.cid, o.prev, o.created_at, r.record FROM happyview_space_record_oplog o LEFT JOIN happyview_space_records r ON r.space_id = o.space_id AND r.author_did = o.author_did AND r.collection = o.collection AND r.rkey = o.rkey AND r.cid = o.cid WHERE o.space_id = ? AND o.author_did = ? AND o.rev > ? ORDER BY o.rev, o.idx LIMIT ?", + "SELECT o.id, o.space_id, o.author_did, o.rev, o.idx, o.action, o.collection, o.rkey, o.cid, o.prev, o.created_at, r.record FROM happyview_space_record_oplog o LEFT JOIN happyview_space_records r ON r.space_id = o.space_id AND r.author_did = o.author_did AND r.collection = o.collection AND r.rkey = o.rkey AND r.cid = o.cid WHERE o.space_id = ? AND o.author_did = ? AND (o.rev > ? OR (o.rev = ? AND o.idx > ?)) ORDER BY o.rev, o.idx LIMIT ?", backend, ) } else { @@ -136,16 +175,21 @@ pub async fn list_ops_with_values( let mut query = crate::db::query_as::(&sql) .bind(space_id) .bind(author_did); - if let Some(rev) = since_rev { - query = query.bind(rev); + if let Some(c) = cursor { + let (rev, idx) = parse_cursor(c); + query = query.bind(rev.to_string()).bind(rev.to_string()).bind(idx); } - query = query.bind(limit); + // Over-fetch by one to detect whether another page exists. + query = query.bind(limit + 1); - let rows = query.fetch_all(pool).await.map_err(|e| { + let rows = query.fetch_all(executor).await.map_err(|e| { AppError::Internal(format!("failed to list oplog entries with values: {e}")) })?; - rows.into_iter() + let (rows, next_cursor) = paginate(rows, limit, |r| (r.3.clone(), r.4)); + + let entries = rows + .into_iter() .map(|r| { let action = OplogAction::parse(&r.5) .ok_or_else(|| AppError::Internal(format!("invalid oplog action: {}", r.5)))?; @@ -170,5 +214,54 @@ pub async fn list_ops_with_values( created_at: r.10, }) }) - .collect() + .collect::, AppError>>()?; + + Ok((entries, next_cursor)) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn cursor_roundtrip() { + let c = encode_cursor("3k2rev1", 4); + assert_eq!(c, "3k2rev1:4"); + assert_eq!(parse_cursor(&c), ("3k2rev1", 4)); + } + + #[test] + fn bare_rev_cursor_skips_whole_rev() { + // Backwards-compat: a rev-only cursor must exclude every op of that rev, so + // the `rev = ? AND idx > ?` arm can never match. + assert_eq!(parse_cursor("3k2rev1"), ("3k2rev1", i32::MAX)); + } + + #[test] + fn malformed_idx_falls_back_to_whole_rev() { + assert_eq!( + parse_cursor("3k2rev1:notanint"), + ("3k2rev1:notanint", i32::MAX) + ); + } + + #[test] + fn paginate_returns_no_cursor_when_page_not_full() { + let rows = vec![("a".to_string(), 0), ("b".to_string(), 1)]; + let (out, cursor) = paginate(rows, 5, |r| (r.0.clone(), r.1)); + assert_eq!(out.len(), 2); + assert!(cursor.is_none()); + } + + #[test] + fn paginate_truncates_and_emits_cursor_from_last_kept_row() { + let rows = vec![ + ("rev1".to_string(), 0), + ("rev1".to_string(), 1), + ("rev1".to_string(), 2), + ]; + let (out, cursor) = paginate(rows, 2, |r| (r.0.clone(), r.1)); + assert_eq!(out.len(), 2); + assert_eq!(cursor.as_deref(), Some("rev1:1")); + } } diff --git a/src/spaces/routes.rs b/src/spaces/routes.rs index 0b7b02d..385b84c 100644 --- a/src/spaces/routes.rs +++ b/src/spaces/routes.rs @@ -593,6 +593,14 @@ async fn apply_writes( } let mut results = Vec::with_capacity(input.writes.len()); + let rev = generate_tid(); + let mut ops: Vec = Vec::with_capacity(input.writes.len()); + + let mut tx = state + .db + .begin() + .await + .map_err(|e| AppError::Internal(format!("failed to begin transaction: {e}")))?; for op in input.writes { match op { @@ -612,13 +620,20 @@ async fn apply_writes( uri: record_uri.clone(), space_id: space.id.clone(), author_did: did.clone(), - collection, - rkey, + collection: collection.clone(), + rkey: rkey.clone(), record: value, cid: cid.clone(), indexed_at: now_rfc3339(), }; - db::insert_space_record(&state.db, state.db_backend, &record).await?; + db::insert_space_record(&mut *tx, state.db_backend, &record).await?; + ops.push(service::AppliedOp { + action: OplogAction::Create, + collection, + rkey, + new_cid: Some(cid.clone()), + old_cid: None, + }); results.push(serde_json::json!({ "uri": record_uri, "cid": cid, @@ -640,23 +655,40 @@ async fn apply_writes( uri: record_uri.clone(), space_id: space.id.clone(), author_did: did.clone(), - collection, - rkey, + collection: collection.clone(), + rkey: rkey.clone(), record: value, cid: cid.clone(), indexed_at: now_rfc3339(), }; + let old_cid = match swap_record.as_deref() { + Some(swap) => Some(swap.to_string()), + None => db::get_space_record(&mut *tx, state.db_backend, &record_uri) + .await? + .map(|r| r.cid), + }; if let Some(swap_cid) = swap_record { db::upsert_space_record_with_swap( - &state.db, + &mut tx, state.db_backend, &record, &swap_cid, ) .await?; } else { - db::upsert_space_record(&state.db, state.db_backend, &record).await?; + db::upsert_space_record(&mut *tx, state.db_backend, &record).await?; } + ops.push(service::AppliedOp { + action: if old_cid.is_some() { + OplogAction::Update + } else { + OplogAction::Create + }, + collection, + rkey, + new_cid: Some(cid.clone()), + old_cid, + }); results.push(serde_json::json!({ "uri": record_uri, "cid": cid, @@ -671,37 +703,50 @@ async fn apply_writes( "at://{}/space/{}/{}/{}/{}/{}", space.did, space.type_nsid, space.skey, did, collection, rkey ); - if let Some(swap_cid) = swap_record { + let old_cid = if let Some(swap_cid) = swap_record { db::delete_space_record_with_swap( - &state.db, + &mut tx, state.db_backend, &record_uri, &swap_cid, ) .await?; + swap_cid } else { // Mirrors `service::delete_record`'s non-swap ownership // check: only the record's own author may delete it. let existing = - db::get_space_record(&state.db, state.db_backend, &record_uri).await?; - match existing { + db::get_space_record(&mut *tx, state.db_backend, &record_uri).await?; + let cid = match existing { Some(r) if r.author_did != did => { return Err(AppError::Forbidden( "You can only delete your own records".into(), )); } None => return Err(AppError::NotFound("Record not found".into())), - _ => {} - } - db::delete_space_record(&state.db, state.db_backend, &record_uri).await?; - } + Some(r) => r.cid, + }; + db::delete_space_record(&mut *tx, state.db_backend, &record_uri).await?; + cid + }; + ops.push(service::AppliedOp { + action: OplogAction::Delete, + collection, + rkey, + new_cid: None, + old_cid: Some(old_cid), + }); results.push(serde_json::json!({})); } } } - let rev = generate_tid(); - db::update_space_revision(&state.db, state.db_backend, &space.id, &rev).await?; + service::commit_write(&mut tx, state.db_backend, &space, &did, &rev, &ops).await?; + tx.commit() + .await + .map_err(|e| AppError::Internal(format!("failed to commit transaction: {e}")))?; + + service::notify_ops(&state, &space, &did, &ops).await; Ok(Json(serde_json::json!({ "results": results, @@ -962,6 +1007,25 @@ async fn get_delegation_token( // Protocol endpoint implementations // --------------------------------------------------------------------------- +fn repo_state_commit_json(repo_state: &RepoState) -> Result, AppError> { + let Some(h) = repo_state.hash.as_deref() else { + return Ok(None); + }; + let ikm = repo_state.ikm.as_deref().ok_or_else(|| { + AppError::Internal("corrupt repo state: hash present but ikm missing".into()) + })?; + let mac = repo_state.mac.as_deref().ok_or_else(|| { + AppError::Internal("corrupt repo state: hash present but mac missing".into()) + })?; + Ok(Some(serde_json::json!({ + "ver": 1, + "hash": base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(h), + "ikm": base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(ikm), + "mac": base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(mac), + "rev": repo_state.rev, + }))) +} + async fn get_latest_commit( State(state): State, claims: XrpcClaims, @@ -982,30 +1046,15 @@ async fn get_latest_commit( let read_access = SpaceReadAccess::from_space_access(membership); check_read_access(&did, ¶ms.did, read_access, has_credential)?; + let mut conn = state + .db + .acquire() + .await + .map_err(|e| AppError::Internal(format!("failed to acquire connection: {e}")))?; let repo_state = - db::get_or_create_repo_state(&state.db, state.db_backend, &space.id, ¶ms.did).await?; + db::get_or_create_repo_state(&mut conn, state.db_backend, &space.id, ¶ms.did).await?; - let commit = if let Some(h) = repo_state.hash.as_ref() { - let ikm = repo_state.ikm.as_deref().ok_or_else(|| { - AppError::Internal("corrupt repo state: hash present but ikm missing".into()) - })?; - let sig = repo_state.sig.as_deref().ok_or_else(|| { - AppError::Internal("corrupt repo state: hash present but sig missing".into()) - })?; - let mac = repo_state.mac.as_deref().ok_or_else(|| { - AppError::Internal("corrupt repo state: hash present but mac missing".into()) - })?; - Some(serde_json::json!({ - "ver": 1, - "hash": base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(h), - "ikm": base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(ikm), - "sig": base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(sig), - "mac": base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(mac), - "rev": repo_state.rev, - })) - } else { - None - }; + let commit = repo_state_commit_json(&repo_state)?; Ok(Json(serde_json::json!({ "rev": repo_state.rev, @@ -1033,8 +1082,14 @@ async fn get_repo( let read_access = SpaceReadAccess::from_space_access(membership); check_read_access(&did, ¶ms.did, read_access, has_credential)?; + let mut conn = state + .db + .acquire() + .await + .map_err(|e| AppError::Internal(format!("failed to acquire connection: {e}")))?; let repo_state = - db::get_or_create_repo_state(&state.db, state.db_backend, &space.id, ¶ms.did).await?; + db::get_or_create_repo_state(&mut conn, state.db_backend, &space.id, ¶ms.did).await?; + drop(conn); let records = db::list_all_space_records(&state.db, state.db_backend, &space.id, ¶ms.did).await?; @@ -1057,9 +1112,6 @@ async fn get_repo( .ok_or_else(|| AppError::Internal("corrupt repo state: missing mac".into()))? .try_into() .map_err(|_| AppError::Internal("corrupt repo state: mac is not 32 bytes".into()))?; - let sig = repo_state - .sig - .ok_or_else(|| AppError::Internal("corrupt repo state: missing sig".into()))?; let rev = repo_state .rev .ok_or_else(|| AppError::Internal("corrupt repo state: missing rev".into()))?; @@ -1067,7 +1119,6 @@ async fn get_repo( ver: 1, hash, ikm, - sig, mac, rev, }; @@ -1102,10 +1153,10 @@ async fn list_repo_ops( let read_access = SpaceReadAccess::from_space_access(membership); check_read_access(&did, ¶ms.did, read_access, has_credential)?; - let limit = params.limit.unwrap_or(100).min(1000); + let limit = params.limit.unwrap_or(100).clamp(1, 1000); let exclude_values = params.exclude_values.unwrap_or(false); - let ops = if exclude_values { + let (ops, cursor) = if exclude_values { oplog::list_ops( &state.db, state.db_backend, @@ -1127,7 +1178,20 @@ async fn list_repo_ops( .await? }; - Ok(Json(serde_json::json!({ "ops": ops }))) + let mut conn = state + .db + .acquire() + .await + .map_err(|e| AppError::Internal(format!("failed to acquire connection: {e}")))?; + let repo_state = + db::get_or_create_repo_state(&mut conn, state.db_backend, &space.id, ¶ms.did).await?; + let commit = repo_state_commit_json(&repo_state)?; + + Ok(Json(serde_json::json!({ + "ops": ops, + "cursor": cursor, + "commit": commit, + }))) } async fn list_repos( @@ -1137,11 +1201,6 @@ async fn list_repos( ) -> Result { let space = service::resolve_space(&state, ¶ms.space).await?; - // The repo list is the space's participant list. Its visibility follows the - // `membershipPublic` config (like `get_space` / `listMembers`): public when - // set, otherwise the caller must be an authenticated member (or authority / - // space-credential holder). Previously this required *some* auth but never - // checked membership, leaking any private space's participants (M2). if !space.config.membership_public { let did = require_auth_or_credential(&state, &claims).await?; service::require_membership( diff --git a/src/spaces/service.rs b/src/spaces/service.rs index b799302..175972b 100644 --- a/src/spaces/service.rs +++ b/src/spaces/service.rs @@ -3,11 +3,12 @@ //! the exact same membership/validation/write logic. use crate::AppState; -use crate::db::{adapt_sql, now_rfc3339}; +use crate::db::{DatabaseBackend, adapt_sql, now_rfc3339}; use crate::error::AppError; use crate::lua::tid::generate_tid; +use crate::spaces::lthash::LtHashState; use crate::spaces::types::*; -use crate::spaces::{SpaceUri, db, members}; +use crate::spaces::{SpaceUri, commit, db, lthash, members, notifications, oplog}; use sha2::{Digest, Sha256}; pub(crate) async fn resolve_space(state: &AppState, space_ref: &str) -> Result { @@ -117,6 +118,100 @@ pub(crate) async fn require_membership( Ok(access) } +pub(crate) struct AppliedOp { + pub action: OplogAction, + pub collection: String, + pub rkey: String, + /// CID written by this op — `None` for a delete. + pub new_cid: Option, + /// CID this op superseded — `None` for a create. + pub old_cid: Option, +} + +pub(crate) async fn commit_write( + conn: &mut sqlx::AnyConnection, + backend: DatabaseBackend, + space: &Space, + author_did: &str, + rev: &str, + ops: &[AppliedOp], +) -> Result<(), AppError> { + let mut repo_state = + db::get_or_create_repo_state(&mut *conn, backend, &space.id, author_did).await?; + + let lthash_bytes: [u8; 2048] = repo_state.lthash_state.as_slice().try_into().map_err(|_| { + AppError::Internal(format!( + "corrupt repo state: lthash is {} bytes, expected 2048", + repo_state.lthash_state.len() + )) + })?; + let mut set_hash = LtHashState::from_bytes(lthash_bytes); + + let now = now_rfc3339(); + for (idx, op) in ops.iter().enumerate() { + if let Some(old) = op.old_cid.as_deref() { + set_hash.remove(<hash::record_element(&op.collection, &op.rkey, old)); + } + if let Some(new) = op.new_cid.as_deref() { + set_hash.add(<hash::record_element(&op.collection, &op.rkey, new)); + } + + let entry = OplogEntry { + id: uuid::Uuid::new_v4().to_string(), + space_id: space.id.clone(), + author_did: author_did.to_string(), + rev: rev.to_string(), + idx: idx as i32, + action: op.action, + collection: op.collection.clone(), + rkey: op.rkey.clone(), + cid: op.new_cid.clone(), + prev: op.old_cid.clone(), + value: None, + created_at: now.clone(), + }; + oplog::append_op(&mut *conn, backend, &entry).await?; + } + + let space_uri = format!( + "at://{}/space/{}/{}", + space.did, space.type_nsid, space.skey + ); + let signed = commit::sign_commit(&set_hash.hash(), &space_uri, author_did, rev)?; + + repo_state.lthash_state = set_hash.as_bytes().to_vec(); + repo_state.rev = Some(signed.rev); + repo_state.hash = Some(signed.hash.to_vec()); + repo_state.ikm = Some(signed.ikm.to_vec()); + repo_state.mac = Some(signed.mac.to_vec()); + db::update_repo_state(&mut *conn, backend, &repo_state).await?; + + db::update_space_revision(&mut *conn, backend, &space.id, rev).await?; + + Ok(()) +} + +pub(crate) async fn notify_ops( + state: &AppState, + space: &Space, + author_did: &str, + ops: &[AppliedOp], +) { + for op in ops { + let _ = notifications::dispatch_write_notification( + &state.db, + state.db_backend, + &state.http, + &space.id, + author_did, + &op.collection, + &op.rkey, + op.new_cid.as_deref(), + ) + .await; + } +} + pub(crate) async fn create_record( state: &AppState, did: &str, @@ -145,9 +240,28 @@ pub(crate) async fn create_record( cid: cid.clone(), indexed_at: now_rfc3339(), }; - db::insert_space_record(&state.db, state.db_backend, &rec).await?; let rev = generate_tid(); - db::update_space_revision(&state.db, state.db_backend, &space.id, &rev).await?; + let ops = vec![AppliedOp { + action: OplogAction::Create, + collection: collection.to_string(), + rkey: rec.rkey.clone(), + new_cid: Some(cid.clone()), + old_cid: None, + }]; + + let mut tx = state + .db + .begin() + .await + .map_err(|e| AppError::Internal(format!("failed to begin transaction: {e}")))?; + db::insert_space_record(&mut *tx, state.db_backend, &rec).await?; + commit_write(&mut tx, state.db_backend, &space, did, &rev, &ops).await?; + tx.commit() + .await + .map_err(|e| AppError::Internal(format!("failed to commit transaction: {e}")))?; + + notify_ops(state, &space, did, &ops).await; + Ok((record_uri, cid)) } @@ -181,13 +295,45 @@ pub(crate) async fn put_record( cid: cid.clone(), indexed_at: now_rfc3339(), }; + let rev = generate_tid(); + let mut tx = state + .db + .begin() + .await + .map_err(|e| AppError::Internal(format!("failed to begin transaction: {e}")))?; + + let old_cid = match swap_cid.as_deref() { + Some(swap) => Some(swap.to_string()), + None => db::get_space_record(&mut *tx, state.db_backend, &record_uri) + .await? + .map(|r| r.cid), + }; + let action = if old_cid.is_some() { + OplogAction::Update + } else { + OplogAction::Create + }; + if let Some(swap) = swap_cid { - db::upsert_space_record_with_swap(&state.db, state.db_backend, &rec, &swap).await?; + db::upsert_space_record_with_swap(&mut tx, state.db_backend, &rec, &swap).await?; } else { - db::upsert_space_record(&state.db, state.db_backend, &rec).await?; + db::upsert_space_record(&mut *tx, state.db_backend, &rec).await?; } - let rev = generate_tid(); - db::update_space_revision(&state.db, state.db_backend, &space.id, &rev).await?; + + let ops = vec![AppliedOp { + action, + collection: collection.to_string(), + rkey: rkey.to_string(), + new_cid: Some(cid.clone()), + old_cid, + }]; + commit_write(&mut tx, state.db_backend, &space, did, &rev, &ops).await?; + tx.commit() + .await + .map_err(|e| AppError::Internal(format!("failed to commit transaction: {e}")))?; + + notify_ops(state, &space, did, &ops).await; + Ok((record_uri, cid)) } @@ -206,23 +352,45 @@ pub(crate) async fn delete_record( "at://{}/space/{}/{}/{}/{}/{}", space.did, space.type_nsid, space.skey, did, collection, rkey ); - if let Some(swap) = swap_cid { - db::delete_space_record_with_swap(&state.db, state.db_backend, &record_uri, &swap).await?; + let rev = generate_tid(); + let mut tx = state + .db + .begin() + .await + .map_err(|e| AppError::Internal(format!("failed to begin transaction: {e}")))?; + + let old_cid = if let Some(swap) = swap_cid { + db::delete_space_record_with_swap(&mut tx, state.db_backend, &record_uri, &swap).await?; + swap } else { - let record = db::get_space_record(&state.db, state.db_backend, &record_uri).await?; - match record { + let record = db::get_space_record(&mut *tx, state.db_backend, &record_uri).await?; + let cid = match record { Some(r) if r.author_did != did => { return Err(AppError::Forbidden( "You can only delete your own records".into(), )); } None => return Err(AppError::NotFound("Record not found".into())), - _ => {} - } - db::delete_space_record(&state.db, state.db_backend, &record_uri).await?; - } - let rev = generate_tid(); - db::update_space_revision(&state.db, state.db_backend, &space.id, &rev).await?; + Some(r) => r.cid, + }; + db::delete_space_record(&mut *tx, state.db_backend, &record_uri).await?; + cid + }; + + let ops = vec![AppliedOp { + action: OplogAction::Delete, + collection: collection.to_string(), + rkey: rkey.to_string(), + new_cid: None, + old_cid: Some(old_cid), + }]; + commit_write(&mut tx, state.db_backend, &space, did, &rev, &ops).await?; + tx.commit() + .await + .map_err(|e| AppError::Internal(format!("failed to commit transaction: {e}")))?; + + notify_ops(state, &space, did, &ops).await; + Ok(()) } diff --git a/src/spaces/types.rs b/src/spaces/types.rs index a2c7d06..2fba8b2 100644 --- a/src/spaces/types.rs +++ b/src/spaces/types.rs @@ -154,7 +154,6 @@ pub struct RepoState { pub rev: Option, pub hash: Option>, pub ikm: Option>, - pub sig: Option>, pub mac: Option>, pub updated_at: String, } diff --git a/tests/spaces_db.rs b/tests/spaces_db.rs index 6ddb678..d47d961 100644 --- a/tests/spaces_db.rs +++ b/tests/spaces_db.rs @@ -179,7 +179,8 @@ async fn get_or_create_repo_state_creates_default() { .expect("create_space failed"); let author_did = "did:plc:repo-author"; - let state = spaces_db::get_or_create_repo_state(&pool, backend, &space_id, author_did) + let mut conn = pool.acquire().await.expect("acquire failed"); + let state = spaces_db::get_or_create_repo_state(&mut conn, backend, &space_id, author_did) .await .expect("get_or_create_repo_state failed"); @@ -210,10 +211,11 @@ async fn get_or_create_repo_state_is_idempotent() { .expect("create_space failed"); let author_did = "did:plc:idem-author"; - let first = spaces_db::get_or_create_repo_state(&pool, backend, &space_id, author_did) + let mut conn = pool.acquire().await.expect("acquire failed"); + let first = spaces_db::get_or_create_repo_state(&mut conn, backend, &space_id, author_did) .await .expect("first call failed"); - let second = spaces_db::get_or_create_repo_state(&pool, backend, &space_id, author_did) + let second = spaces_db::get_or_create_repo_state(&mut conn, backend, &space_id, author_did) .await .expect("second call failed"); @@ -240,7 +242,8 @@ async fn update_repo_state_persists_fields() { .expect("create_space failed"); let author_did = "did:plc:update-author"; - let mut state = spaces_db::get_or_create_repo_state(&pool, backend, &space_id, author_did) + let mut conn = pool.acquire().await.expect("acquire failed"); + let mut state = spaces_db::get_or_create_repo_state(&mut conn, backend, &space_id, author_did) .await .expect("get_or_create failed"); @@ -251,7 +254,7 @@ async fn update_repo_state_persists_fields() { .await .expect("update_repo_state failed"); - let reloaded = spaces_db::get_or_create_repo_state(&pool, backend, &space_id, author_did) + let reloaded = spaces_db::get_or_create_repo_state(&mut conn, backend, &space_id, author_did) .await .expect("reload failed"); @@ -315,9 +318,13 @@ async fn oplog_append_and_list() { .expect("append_op failed"); } - let all_ops = oplog::list_ops(&pool, backend, &space_id, author_did, None, 10) + let (all_ops, cursor) = oplog::list_ops(&pool, backend, &space_id, author_did, None, 10) .await .expect("list_ops failed"); + assert!( + cursor.is_none(), + "page is not full, so there is no next page" + ); assert_eq!(all_ops.len(), 3); assert!(matches!(all_ops[0].action, OplogAction::Create)); assert!(matches!(all_ops[1].action, OplogAction::Update)); @@ -365,12 +372,98 @@ async fn oplog_list_with_since_rev_cursor() { .expect("append_op failed"); } - let after_rev2 = oplog::list_ops(&pool, backend, &space_id, author_did, Some("rev-0002"), 10) - .await - .expect("list_ops with cursor failed"); + let (after_rev2, _) = + oplog::list_ops(&pool, backend, &space_id, author_did, Some("rev-0002"), 10) + .await + .expect("list_ops with cursor failed"); assert_eq!(after_rev2.len(), 3); assert_eq!(after_rev2[0].rev, "rev-0003"); + + let (after_rev2_idx, _) = oplog::list_ops( + &pool, + backend, + &space_id, + author_did, + Some(&oplog::encode_cursor("rev-0002", -1)), + 10, + ) + .await + .expect("list_ops with composite cursor failed"); + + assert_eq!(after_rev2_idx.len(), 4); + assert_eq!(after_rev2_idx[0].rev, "rev-0002"); +} + +#[tokio::test] +#[serial] +async fn oplog_paginates_within_a_single_rev() { + common::require_db!(); + let pool = test_db::test_pool().await; + let backend = test_db::test_backend(); + test_db::truncate_all(&pool).await; + + let space_id = new_id(); + let space = make_space( + &space_id, + "did:plc:batch-owner", + "com.example.batch", + "batch-skey", + ); + spaces_db::create_space(&pool, backend, &space) + .await + .expect("create_space failed"); + + let author_did = "did:plc:batch-author"; + + for i in 0..5 { + let entry = OplogEntry { + id: new_id(), + space_id: space_id.clone(), + author_did: author_did.to_string(), + rev: "rev-batch".to_string(), + idx: i, + action: OplogAction::Create, + collection: "com.example.item".to_string(), + rkey: format!("item-{i}"), + cid: Some(format!("bafy{i}")), + prev: None, + value: None, + created_at: now_rfc3339(), + }; + oplog::append_op(&pool, backend, &entry) + .await + .expect("append_op failed"); + } + + let (page1, cursor) = oplog::list_ops(&pool, backend, &space_id, author_did, None, 2) + .await + .expect("page 1 failed"); + assert_eq!(page1.len(), 2); + assert_eq!(page1[0].idx, 0); + assert_eq!(page1[1].idx, 1); + let cursor = cursor.expect("expected a next-page cursor"); + + let (page2, cursor2) = oplog::list_ops(&pool, backend, &space_id, author_did, Some(&cursor), 2) + .await + .expect("page 2 failed"); + assert_eq!(page2.len(), 2); + assert_eq!(page2[0].idx, 2); + assert_eq!(page2[1].idx, 3); + + let (page3, cursor3) = oplog::list_ops( + &pool, + backend, + &space_id, + author_did, + Some(&cursor2.expect("expected a second cursor")), + 2, + ) + .await + .expect("page 3 failed"); + assert_eq!(page3.len(), 1, "the tail op must not be dropped"); + assert_eq!(page3[0].idx, 4); + assert!(cursor3.is_none()); } // --------------------------------------------------------------------------- diff --git a/tests/spaces_repo_state.rs b/tests/spaces_repo_state.rs new file mode 100644 index 0000000..b44e2fe --- /dev/null +++ b/tests/spaces_repo_state.rs @@ -0,0 +1,434 @@ +mod common; + +use axum::body::Body; +use axum::http::{HeaderName, HeaderValue, Request, StatusCode}; +use happyview::db::now_rfc3339; +use happyview::spaces::db as spaces_db; +use happyview::spaces::types::*; +use http_body_util::BodyExt; +use serde_json::{Value, json}; +use serial_test::serial; +use tower::ServiceExt; +use uuid::Uuid; + +use common::app::TestApp; + +const TYPE_NSID: &str = "com.example.repostate"; +const COLLECTION: &str = "com.example.item"; + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +fn rand_did(label: &str) -> String { + format!("did:plc:{label}{}", Uuid::new_v4().simple()) +} + +fn rand_skey(label: &str) -> String { + format!("{label}{}", Uuid::new_v4().simple()) +} + +async fn enable_spaces(app: &TestApp) { + let (name, value) = app.admin_cookie(); + let req = Request::builder() + .method("PUT") + .uri("/admin/settings/feature.spaces_enabled") + .header(name, value) + .header("content-type", "application/json") + .body(Body::from(json!({ "value": "true" }).to_string())) + .unwrap(); + assert!( + app.router + .clone() + .oneshot(req) + .await + .unwrap() + .status() + .is_success(), + "failed to enable spaces" + ); +} + +fn cookie_for(app: &TestApp, did: &str) -> (HeaderName, HeaderValue) { + common::auth::admin_cookie_header(did, &app.state.cookie_key) +} + +async fn json_of(resp: axum::http::Response) -> Value { + let body = resp.into_body().collect().await.unwrap().to_bytes(); + serde_json::from_slice(&body).unwrap_or(json!(null)) +} + +async fn create_space(app: &TestApp, authority: &str, skey: &str) -> (String, String) { + let now = now_rfc3339(); + let id = Uuid::new_v4().to_string(); + let space = Space { + id: id.clone(), + did: authority.to_string(), + authority_did: authority.to_string(), + creator_did: authority.to_string(), + type_nsid: TYPE_NSID.to_string(), + skey: skey.to_string(), + display_name: None, + description: None, + mint_policy: MintPolicy::MemberList, + app_access: AppAccess::Open, + managing_app_did: None, + config: SpaceConfig::default(), + revision: None, + created_at: now.clone(), + updated_at: now, + }; + spaces_db::create_space(&app.state.db, app.state.db_backend, &space) + .await + .expect("create_space failed"); + let uri = format!("at://{authority}/space/{TYPE_NSID}/{skey}"); + (id, uri) +} + +async fn add_member(app: &TestApp, space_id: &str, did: &str, access: SpaceAccess) { + spaces_db::add_member( + &app.state.db, + app.state.db_backend, + &SpaceMember { + id: Uuid::new_v4().to_string(), + space_id: space_id.to_string(), + did: did.to_string(), + access, + is_delegation: false, + granted_by: None, + created_at: now_rfc3339(), + }, + ) + .await + .expect("add_member failed"); +} + +async fn setup(label: &str) -> (TestApp, String, String, String) { + let app = TestApp::new().await; + enable_spaces(&app).await; + let authority = rand_did(label); + let (space_id, space_uri) = create_space(&app, &authority, &rand_skey(label)).await; + add_member(&app, &space_id, &authority, SpaceAccess::Write).await; + (app, space_id, space_uri, authority) +} + +async fn create_record(app: &TestApp, space_uri: &str, did: &str, body: Value) -> Value { + let (name, value) = cookie_for(app, did); + let req = Request::builder() + .method("POST") + .uri("/xrpc/com.atproto.space.createRecord") + .header("content-type", "application/json") + .header(name, value) + .body(Body::from( + json!({ "space": space_uri, "collection": COLLECTION, "record": body }).to_string(), + )) + .unwrap(); + let resp = app.router.clone().oneshot(req).await.unwrap(); + assert!( + resp.status().is_success(), + "createRecord failed: {}", + resp.status() + ); + json_of(resp).await +} + +async fn get_repo_state(app: &TestApp, space_uri: &str, did: &str) -> Value { + let (name, value) = cookie_for(app, did); + let req = Request::builder() + .method("GET") + .uri(format!( + "/xrpc/com.atproto.space.getRepoState?space={}&did={}", + urlencoding::encode(space_uri), + urlencoding::encode(did) + )) + .header(name, value) + .body(Body::empty()) + .unwrap(); + let resp = app.router.clone().oneshot(req).await.unwrap(); + assert_eq!(resp.status(), StatusCode::OK, "getRepoState failed"); + json_of(resp).await +} + +async fn list_repo_ops(app: &TestApp, space_uri: &str, did: &str, extra: &str) -> Value { + let (name, value) = cookie_for(app, did); + let req = Request::builder() + .method("GET") + .uri(format!( + "/xrpc/com.atproto.space.listRepoOps?space={}&did={}{extra}", + urlencoding::encode(space_uri), + urlencoding::encode(did) + )) + .header(name, value) + .body(Body::empty()) + .unwrap(); + let resp = app.router.clone().oneshot(req).await.unwrap(); + assert_eq!(resp.status(), StatusCode::OK, "listRepoOps failed"); + json_of(resp).await +} + +// --------------------------------------------------------------------------- +// Tests +// --------------------------------------------------------------------------- + +#[tokio::test] +#[serial] +async fn write_populates_commit() { + common::require_db!(); + let (app, _space_id, space_uri, did) = setup("commit").await; + + // No writes yet: no commit exists. + let before = get_repo_state(&app, &space_uri, &did).await; + assert!( + before["commit"].is_null(), + "a repo with no writes must have no commit, got {before}" + ); + + create_record(&app, &space_uri, &did, json!({ "text": "hello" })).await; + + let after = get_repo_state(&app, &space_uri, &did).await; + let commit = &after["commit"]; + assert!(!commit.is_null(), "commit must exist after a write"); + assert_eq!(commit["ver"], 1); + assert!(commit["hash"].as_str().is_some(), "hash must be set"); + assert!(commit["ikm"].as_str().is_some(), "ikm must be set"); + assert!(commit["mac"].as_str().is_some(), "mac must be set"); + assert!(commit["rev"].as_str().is_some(), "rev must be set"); + assert!(after["rev"].as_str().is_some()); +} + +#[tokio::test] +#[serial] +async fn commit_has_no_asymmetric_signature() { + common::require_db!(); + let (app, _space_id, space_uri, did) = setup("nosig").await; + create_record(&app, &space_uri, &did, json!({ "text": "hi" })).await; + + let state = get_repo_state(&app, &space_uri, &did).await; + let commit = state["commit"].as_object().expect("commit must exist"); + assert!( + !commit.contains_key("sig"), + "commits must not expose an asymmetric signature, got {commit:?}" + ); +} + +#[tokio::test] +#[serial] +async fn lthash_remove_undoes_add() { + common::require_db!(); + let (app, _space_id, space_uri, did) = setup("lthash").await; + + async fn delete(app: &TestApp, space_uri: &str, did: &str, rkey: &str) { + let (name, value) = cookie_for(app, did); + let req = Request::builder() + .method("POST") + .uri("/xrpc/com.atproto.space.deleteRecord") + .header("content-type", "application/json") + .header(name, value) + .body(Body::from( + json!({ "space": space_uri, "collection": COLLECTION, "rkey": rkey }).to_string(), + )) + .unwrap(); + let resp = app.router.clone().oneshot(req).await.unwrap(); + assert!( + resp.status().is_success(), + "deleteRecord failed: {}", + resp.status() + ); + } + + fn rkey_of(created: &Value) -> String { + created["uri"] + .as_str() + .expect("uri") + .rsplit('/') + .next() + .expect("rkey") + .to_string() + } + + // Record A stays put; record B is the one we add then remove. + create_record(&app, &space_uri, &did, json!({ "text": "keeper" })).await; + let hash_a = get_repo_state(&app, &space_uri, &did).await["commit"]["hash"] + .as_str() + .expect("hash") + .to_string(); + + let b = create_record(&app, &space_uri, &did, json!({ "text": "transient" })).await; + let hash_ab = get_repo_state(&app, &space_uri, &did).await["commit"]["hash"] + .as_str() + .expect("hash") + .to_string(); + assert_ne!( + hash_a, hash_ab, + "adding a record must change the multiset hash" + ); + + delete(&app, &space_uri, &did, &rkey_of(&b)).await; + let hash_after_delete = get_repo_state(&app, &space_uri, &did).await["commit"]["hash"] + .as_str() + .expect("hash") + .to_string(); + + assert_eq!( + hash_a, hash_after_delete, + "removing B must return the hash to its exact pre-B value" + ); +} + +#[tokio::test] +#[serial] +async fn list_repo_ops_returns_ops_and_bundled_commit() { + common::require_db!(); + let (app, _space_id, space_uri, did) = setup("ops").await; + create_record(&app, &space_uri, &did, json!({ "text": "one" })).await; + + let body = list_repo_ops(&app, &space_uri, &did, "").await; + let ops = body["ops"].as_array().expect("ops array"); + assert_eq!(ops.len(), 1, "expected exactly one op, got {body}"); + assert_eq!(ops[0]["action"], "create"); + assert_eq!(ops[0]["collection"], COLLECTION); + assert_eq!(ops[0]["idx"], 0); + assert!(ops[0]["cid"].as_str().is_some()); + assert!( + ops[0]["prev"].is_null(), + "a create has nothing to link back to" + ); + + assert!( + body["cursor"].is_null(), + "page is not full, so there is no next page" + ); + assert!( + !body["commit"].is_null(), + "listRepoOps must bundle the current MAC'd hash" + ); + assert_eq!(body["commit"]["hash"], { + let s = get_repo_state(&app, &space_uri, &did).await; + s["commit"]["hash"].clone() + }); +} + +#[tokio::test] +#[serial] +async fn get_repo_returns_car_after_write() { + common::require_db!(); + let (app, _space_id, space_uri, did) = setup("car").await; + create_record(&app, &space_uri, &did, json!({ "text": "in a car" })).await; + + let (name, value) = cookie_for(&app, &did); + let req = Request::builder() + .method("GET") + .uri(format!( + "/xrpc/com.atproto.space.getRepo?space={}&did={}", + urlencoding::encode(&space_uri), + urlencoding::encode(&did) + )) + .header(name, value) + .body(Body::empty()) + .unwrap(); + let resp = app.router.clone().oneshot(req).await.unwrap(); + assert_eq!(resp.status(), StatusCode::OK, "getRepo must succeed"); + assert_eq!( + resp.headers().get("content-type").unwrap(), + "application/vnd.ipld.car" + ); + let bytes = resp.into_body().collect().await.unwrap().to_bytes(); + assert!(bytes.len() > 10, "CAR body looks empty"); +} + +#[tokio::test] +#[serial] +async fn apply_writes_shares_one_rev_with_ordered_idx() { + common::require_db!(); + let (app, _space_id, space_uri, did) = setup("batch").await; + + let (name, value) = cookie_for(&app, &did); + let req = Request::builder() + .method("POST") + .uri("/xrpc/com.atproto.space.applyWrites") + .header("content-type", "application/json") + .header(name, value) + .body(Body::from( + json!({ + "space": space_uri, + "writes": [ + { "action": "create", "collection": COLLECTION, "value": { "n": 1 } }, + { "action": "create", "collection": COLLECTION, "value": { "n": 2 } }, + { "action": "create", "collection": COLLECTION, "value": { "n": 3 } }, + ] + }) + .to_string(), + )) + .unwrap(); + let resp = app.router.clone().oneshot(req).await.unwrap(); + assert!( + resp.status().is_success(), + "applyWrites failed: {}", + resp.status() + ); + + let body = list_repo_ops(&app, &space_uri, &did, "").await; + let ops = body["ops"].as_array().expect("ops"); + assert_eq!(ops.len(), 3); + + let rev = ops[0]["rev"].as_str().expect("rev"); + for (i, op) in ops.iter().enumerate() { + assert_eq!(op["rev"].as_str().unwrap(), rev, "batch shares one rev"); + assert_eq!(op["idx"], i as i64, "ops are ordered by idx within the rev"); + } + + // One commit for the whole batch, stamped with the batch rev. + let state = get_repo_state(&app, &space_uri, &did).await; + assert_eq!(state["commit"]["rev"].as_str().unwrap(), rev); +} + +#[tokio::test] +#[serial] +async fn list_repo_ops_paginates_within_a_batch_rev() { + common::require_db!(); + let (app, _space_id, space_uri, did) = setup("page").await; + + let writes: Vec = (0..5) + .map(|n| json!({ "action": "create", "collection": COLLECTION, "value": { "n": n } })) + .collect(); + let (name, value) = cookie_for(&app, &did); + let req = Request::builder() + .method("POST") + .uri("/xrpc/com.atproto.space.applyWrites") + .header("content-type", "application/json") + .header(name, value) + .body(Body::from( + json!({ "space": space_uri, "writes": writes }).to_string(), + )) + .unwrap(); + assert!( + app.router + .clone() + .oneshot(req) + .await + .unwrap() + .status() + .is_success() + ); + + let mut seen = Vec::new(); + let mut extra = "&limit=2".to_string(); + loop { + let body = list_repo_ops(&app, &space_uri, &did, &extra).await; + let ops = body["ops"].as_array().unwrap().clone(); + assert!(ops.len() <= 2, "page must respect the limit"); + for op in &ops { + seen.push(op["idx"].as_i64().unwrap()); + } + match body["cursor"].as_str() { + Some(c) => extra = format!("&limit=2&cursor={}", urlencoding::encode(c)), + None => break, + } + } + + assert_eq!( + seen, + vec![0, 1, 2, 3, 4], + "every op in the batch must be visited exactly once, in order" + ); +}