Our Personal Data Server from scratch!
0

Configure Feed

Select the types of activity you want to include in your feed.

1use crate::repo::record::write::{CommitInfo, prepare_repo_write}; 2use axum::{Json, extract::State}; 3use cid::Cid; 4use jacquard_repo::{commit::Commit, mst::Mst, storage::BlockStore}; 5use serde::{Deserialize, Serialize}; 6use serde_json::json; 7use std::str::FromStr; 8use std::sync::Arc; 9use tracing::error; 10use tranquil_pds::api::error::ApiError; 11use tranquil_pds::auth::{Active, Auth, VerifyScope}; 12use tranquil_pds::repo::TrackingBlockStore; 13use tranquil_pds::repo_ops::{ 14 CommitError, FinalizeParams, RecordOp, begin_repo_write, finalize_repo_write, 15}; 16use tranquil_pds::state::AppState; 17use tranquil_pds::types::{AtIdentifier, AtUri, Did, Nsid, Rkey}; 18 19#[derive(Deserialize)] 20pub struct DeleteRecordInput { 21 pub repo: AtIdentifier, 22 pub collection: Nsid, 23 pub rkey: Rkey, 24 #[serde(rename = "swapRecord")] 25 pub swap_record: Option<String>, 26 #[serde(rename = "swapCommit")] 27 pub swap_commit: Option<String>, 28} 29 30#[derive(Serialize)] 31#[serde(rename_all = "camelCase")] 32pub struct DeleteRecordOutput { 33 #[serde(skip_serializing_if = "Option::is_none")] 34 pub commit: Option<CommitInfo>, 35} 36 37pub async fn delete_record( 38 State(state): State<AppState>, 39 auth: Auth<Active>, 40 Json(input): Json<DeleteRecordInput>, 41) -> Result<Json<DeleteRecordOutput>, ApiError> { 42 let scope_proof = auth.verify_repo_delete(&input.collection)?; 43 let repo_auth = prepare_repo_write(&state, &scope_proof, &input.repo).await?; 44 let did = repo_auth.did; 45 let user_id = repo_auth.user_id; 46 let controller_did = repo_auth.controller_did; 47 48 let (ctx, mst) = begin_repo_write(&state, user_id, input.swap_commit.as_deref()).await?; 49 50 let key = format!("{}/{}", input.collection, input.rkey); 51 52 if let Some(swap_record_str) = &input.swap_record { 53 let expected_cid = Cid::from_str(swap_record_str).ok(); 54 let actual_cid = mst.get(&key).await.ok().flatten(); 55 if expected_cid != actual_cid { 56 return Err(ApiError::InvalidSwap(Some( 57 "Record has been modified or does not exist".into(), 58 ))); 59 } 60 } 61 62 let prev_record_cid = mst.get(&key).await.ok().flatten(); 63 if prev_record_cid.is_none() { 64 return Ok(Json(DeleteRecordOutput { commit: None })); 65 } 66 67 let new_mst = mst.delete(&key).await.map_err(|e| { 68 error!("Failed to delete from MST: {:?}", e); 69 ApiError::InternalError(Some("Failed to delete from MST".into())) 70 })?; 71 72 let op = RecordOp::Delete { 73 collection: input.collection.clone(), 74 rkey: input.rkey.clone(), 75 prev: prev_record_cid, 76 }; 77 78 let modified_keys = [key]; 79 let deleted_uri = AtUri::from_parts(&did, &input.collection, &input.rkey); 80 81 let commit_result = finalize_repo_write( 82 &state, 83 ctx, 84 new_mst, 85 FinalizeParams { 86 did: &did, 87 user_id, 88 controller_did: controller_did.as_ref(), 89 delegation_detail: controller_did.as_ref().map(|_| { 90 json!({ 91 "action": "delete", 92 "collection": input.collection, 93 "rkey": input.rkey 94 }) 95 }), 96 ops: vec![op], 97 modified_keys: &modified_keys, 98 blob_cids: &[], 99 backlinks_to_add: vec![], 100 backlinks_to_remove: vec![deleted_uri], 101 }, 102 ) 103 .await?; 104 105 Ok(Json(DeleteRecordOutput { 106 commit: Some(CommitInfo { 107 cid: commit_result.commit_cid.to_string(), 108 rev: commit_result.rev, 109 }), 110 })) 111} 112 113use uuid::Uuid; 114 115pub async fn delete_record_internal( 116 state: &AppState, 117 did: &Did, 118 user_id: Uuid, 119 collection: &Nsid, 120 rkey: &Rkey, 121) -> Result<(), CommitError> { 122 use tranquil_pds::repo_ops::{CommitParams, RecordOp, commit_and_log}; 123 124 let _write_lock = state.repo_write_locks.lock(user_id).await; 125 126 let root_cid_str = state 127 .repos 128 .repo 129 .get_repo_root_cid_by_user_id(user_id) 130 .await 131 .map_err(|e| CommitError::DatabaseError(e.to_string()))? 132 .ok_or(CommitError::RepoNotFound)?; 133 134 let current_root_cid = 135 Cid::from_str(root_cid_str.as_str()).map_err(|e| CommitError::InvalidCid(e.to_string()))?; 136 137 let tracking_store = TrackingBlockStore::new(state.block_store.clone()); 138 let commit_bytes = tracking_store 139 .get(&current_root_cid) 140 .await 141 .map_err(|e| CommitError::BlockStoreFailed(format!("{:?}", e)))? 142 .ok_or(CommitError::BlockStoreFailed( 143 "Commit block not found".into(), 144 ))?; 145 146 let commit = Commit::from_cbor(&commit_bytes) 147 .map_err(|e| CommitError::CommitParseFailed(format!("{:?}", e)))?; 148 149 let mst = Mst::load(Arc::new(tracking_store.clone()), commit.data, None); 150 let key = format!("{}/{}", collection, rkey); 151 152 let prev_record_cid = mst 153 .get(&key) 154 .await 155 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 156 157 let Some(prev_cid) = prev_record_cid else { 158 return Ok(()); 159 }; 160 161 let new_mst = mst 162 .delete(&key) 163 .await 164 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 165 166 let new_mst_root = new_mst 167 .persist() 168 .await 169 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 170 171 let op = RecordOp::Delete { 172 collection: collection.clone(), 173 rkey: rkey.clone(), 174 prev: Some(prev_cid), 175 }; 176 177 let mut new_mst_blocks = std::collections::BTreeMap::new(); 178 let mut old_mst_blocks = std::collections::BTreeMap::new(); 179 180 new_mst 181 .blocks_for_path(&key, &mut new_mst_blocks) 182 .await 183 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 184 185 mst.blocks_for_path(&key, &mut old_mst_blocks) 186 .await 187 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 188 189 let obsolete_cids: Vec<Cid> = std::iter::once(current_root_cid) 190 .chain( 191 old_mst_blocks 192 .keys() 193 .filter(|cid| !new_mst_blocks.contains_key(*cid)) 194 .copied(), 195 ) 196 .chain(std::iter::once(prev_cid)) 197 .collect(); 198 199 let mut relevant_blocks = new_mst_blocks; 200 relevant_blocks.extend(old_mst_blocks); 201 202 let written_cids: Vec<Cid> = tracking_store 203 .get_all_relevant_cids() 204 .into_iter() 205 .chain(relevant_blocks.keys().copied()) 206 .collect::<std::collections::HashSet<_>>() 207 .into_iter() 208 .collect(); 209 210 let written_cids_str: Vec<String> = written_cids.iter().map(ToString::to_string).collect(); 211 212 let deleted_uri = AtUri::from_parts(did.as_str(), collection.as_str(), rkey.as_str()); 213 commit_and_log( 214 state, 215 CommitParams { 216 did, 217 user_id, 218 current_root_cid: Some(current_root_cid), 219 prev_data_cid: Some(commit.data), 220 new_mst_root, 221 ops: vec![op], 222 blocks_cids: &written_cids_str, 223 blobs: &[], 224 obsolete_cids, 225 backlinks_to_add: vec![], 226 backlinks_to_remove: vec![deleted_uri], 227 }, 228 ) 229 .await?; 230 231 Ok(()) 232}