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::utils::{ 2 CommitError, CommitParams, RecordOp, commit_and_log, get_current_root_cid, 3}; 4use crate::repo::record::write::{CommitInfo, prepare_repo_write}; 5use axum::{ 6 Json, 7 extract::State, 8 http::StatusCode, 9 response::{IntoResponse, Response}, 10}; 11use cid::Cid; 12use jacquard_repo::{commit::Commit, mst::Mst, storage::BlockStore}; 13use serde::{Deserialize, Serialize}; 14use serde_json::json; 15use std::str::FromStr; 16use std::sync::Arc; 17use tracing::error; 18use tranquil_pds::api::error::ApiError; 19use tranquil_pds::auth::{Active, Auth, VerifyScope}; 20use tranquil_pds::cid_types::CommitCid; 21use tranquil_pds::delegation::DelegationActionType; 22use tranquil_pds::repo::tracking::TrackingBlockStore; 23use tranquil_pds::state::AppState; 24use tranquil_pds::types::{AtIdentifier, AtUri, Nsid, Rkey}; 25 26#[derive(Deserialize)] 27pub struct DeleteRecordInput { 28 pub repo: AtIdentifier, 29 pub collection: Nsid, 30 pub rkey: Rkey, 31 #[serde(rename = "swapRecord")] 32 pub swap_record: Option<String>, 33 #[serde(rename = "swapCommit")] 34 pub swap_commit: Option<String>, 35} 36 37#[derive(Serialize)] 38#[serde(rename_all = "camelCase")] 39pub struct DeleteRecordOutput { 40 #[serde(skip_serializing_if = "Option::is_none")] 41 pub commit: Option<CommitInfo>, 42} 43 44pub async fn delete_record( 45 State(state): State<AppState>, 46 auth: Auth<Active>, 47 Json(input): Json<DeleteRecordInput>, 48) -> Result<Response, tranquil_pds::api::error::ApiError> { 49 let scope_proof = match auth.verify_repo_delete(&input.collection) { 50 Ok(proof) => proof, 51 Err(e) => return Ok(e.into_response()), 52 }; 53 54 let repo_auth = match prepare_repo_write(&state, &scope_proof, &input.repo).await { 55 Ok(res) => res, 56 Err(err_res) => return Ok(err_res), 57 }; 58 59 let did = repo_auth.did; 60 let user_id = repo_auth.user_id; 61 let controller_did = repo_auth.controller_did; 62 63 let _write_lock = state.repo_write_locks.lock(user_id).await; 64 let current_root_cid = get_current_root_cid(&state, user_id).await?; 65 66 if let Some(swap_commit) = &input.swap_commit 67 && CommitCid::from_str(swap_commit).ok().as_ref() != Some(&current_root_cid) 68 { 69 return Ok(ApiError::InvalidSwap(Some("Repo has been modified".into())).into_response()); 70 } 71 let tracking_store = TrackingBlockStore::new(state.block_store.clone()); 72 let commit_bytes = match tracking_store.get(current_root_cid.as_cid()).await { 73 Ok(Some(b)) => b, 74 _ => { 75 return Ok( 76 ApiError::InternalError(Some("Commit block not found".into())).into_response(), 77 ); 78 } 79 }; 80 let commit = match Commit::from_cbor(&commit_bytes) { 81 Ok(c) => c, 82 _ => { 83 return Ok( 84 ApiError::InternalError(Some("Failed to parse commit".into())).into_response(), 85 ); 86 } 87 }; 88 let mst = Mst::load(Arc::new(tracking_store.clone()), commit.data, None); 89 let key = format!("{}/{}", input.collection, input.rkey); 90 if let Some(swap_record_str) = &input.swap_record { 91 let expected_cid = Cid::from_str(swap_record_str).ok(); 92 let actual_cid = mst.get(&key).await.ok().flatten(); 93 if expected_cid != actual_cid { 94 return Ok(ApiError::InvalidSwap(Some( 95 "Record has been modified or does not exist".into(), 96 )) 97 .into_response()); 98 } 99 } 100 let prev_record_cid = mst.get(&key).await.ok().flatten(); 101 if prev_record_cid.is_none() { 102 return Ok((StatusCode::OK, Json(DeleteRecordOutput { commit: None })).into_response()); 103 } 104 let new_mst = match mst.delete(&key).await { 105 Ok(m) => m, 106 Err(e) => { 107 error!("Failed to delete from MST: {:?}", e); 108 return Ok(ApiError::InternalError(Some(format!( 109 "Failed to delete from MST: {:?}", 110 e 111 ))) 112 .into_response()); 113 } 114 }; 115 let new_mst_root = match new_mst.persist().await { 116 Ok(c) => c, 117 Err(e) => { 118 error!("Failed to persist MST: {:?}", e); 119 return Ok( 120 ApiError::InternalError(Some("Failed to persist MST".into())).into_response(), 121 ); 122 } 123 }; 124 let collection_for_audit = input.collection.to_string(); 125 let rkey_for_audit = input.rkey.to_string(); 126 let op = RecordOp::Delete { 127 collection: input.collection.clone(), 128 rkey: input.rkey.clone(), 129 prev: prev_record_cid, 130 }; 131 let mut new_mst_blocks = std::collections::BTreeMap::new(); 132 let mut old_mst_blocks = std::collections::BTreeMap::new(); 133 if new_mst 134 .blocks_for_path(&key, &mut new_mst_blocks) 135 .await 136 .is_err() 137 { 138 return Ok( 139 ApiError::InternalError(Some("Failed to get new MST blocks for path".into())) 140 .into_response(), 141 ); 142 } 143 if mst 144 .blocks_for_path(&key, &mut old_mst_blocks) 145 .await 146 .is_err() 147 { 148 return Ok( 149 ApiError::InternalError(Some("Failed to get old MST blocks for path".into())) 150 .into_response(), 151 ); 152 } 153 let mut relevant_blocks = new_mst_blocks.clone(); 154 relevant_blocks.extend(old_mst_blocks.iter().map(|(k, v)| (*k, v.clone()))); 155 let written_cids: Vec<Cid> = tracking_store 156 .get_all_relevant_cids() 157 .into_iter() 158 .chain(relevant_blocks.keys().copied()) 159 .collect::<std::collections::HashSet<_>>() 160 .into_iter() 161 .collect(); 162 let written_cids_str: Vec<String> = written_cids.iter().map(|c| c.to_string()).collect(); 163 let obsolete_cids: Vec<Cid> = std::iter::once(current_root_cid.into_cid()) 164 .chain( 165 old_mst_blocks 166 .keys() 167 .filter(|cid| !new_mst_blocks.contains_key(*cid)) 168 .copied(), 169 ) 170 .chain(prev_record_cid) 171 .collect(); 172 let commit_result = match commit_and_log( 173 &state, 174 CommitParams { 175 did: &did, 176 user_id, 177 current_root_cid: Some(current_root_cid.into_cid()), 178 prev_data_cid: Some(commit.data), 179 new_mst_root, 180 ops: vec![op], 181 blocks_cids: &written_cids_str, 182 blobs: &[], 183 obsolete_cids, 184 }, 185 ) 186 .await 187 { 188 Ok(res) => res, 189 Err(e) => return Ok(ApiError::from(e).into_response()), 190 }; 191 192 if let Some(ref controller) = controller_did { 193 let _ = state 194 .delegation_repo 195 .log_delegation_action( 196 &did, 197 controller, 198 Some(controller), 199 DelegationActionType::RepoWrite, 200 Some(json!({ 201 "action": "delete", 202 "collection": collection_for_audit, 203 "rkey": rkey_for_audit 204 })), 205 None, 206 None, 207 ) 208 .await; 209 } 210 211 let deleted_uri = AtUri::from_parts(&did, &input.collection, &input.rkey); 212 if let Err(e) = state 213 .backlink_repo 214 .remove_backlinks_by_uri(&deleted_uri) 215 .await 216 { 217 error!("Failed to remove backlinks for {}: {}", deleted_uri, e); 218 } 219 220 Ok(( 221 StatusCode::OK, 222 Json(DeleteRecordOutput { 223 commit: Some(CommitInfo { 224 cid: commit_result.commit_cid.to_string(), 225 rev: commit_result.rev, 226 }), 227 }), 228 ) 229 .into_response()) 230} 231 232use tranquil_pds::types::Did; 233use uuid::Uuid; 234 235pub async fn delete_record_internal( 236 state: &AppState, 237 did: &Did, 238 user_id: Uuid, 239 collection: &Nsid, 240 rkey: &Rkey, 241) -> Result<(), CommitError> { 242 let _write_lock = state.repo_write_locks.lock(user_id).await; 243 244 let root_cid_str = state 245 .repo_repo 246 .get_repo_root_cid_by_user_id(user_id) 247 .await 248 .map_err(|e| CommitError::DatabaseError(e.to_string()))? 249 .ok_or(CommitError::RepoNotFound)?; 250 251 let current_root_cid = 252 Cid::from_str(root_cid_str.as_str()).map_err(|e| CommitError::InvalidCid(e.to_string()))?; 253 254 let tracking_store = TrackingBlockStore::new(state.block_store.clone()); 255 let commit_bytes = tracking_store 256 .get(&current_root_cid) 257 .await 258 .map_err(|e| CommitError::BlockStoreFailed(format!("{:?}", e)))? 259 .ok_or(CommitError::BlockStoreFailed( 260 "Commit block not found".into(), 261 ))?; 262 263 let commit = Commit::from_cbor(&commit_bytes) 264 .map_err(|e| CommitError::CommitParseFailed(format!("{:?}", e)))?; 265 266 let mst = Mst::load(Arc::new(tracking_store.clone()), commit.data, None); 267 let key = format!("{}/{}", collection, rkey); 268 269 let prev_record_cid = mst 270 .get(&key) 271 .await 272 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 273 274 let Some(prev_cid) = prev_record_cid else { 275 return Ok(()); 276 }; 277 278 let new_mst = mst 279 .delete(&key) 280 .await 281 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 282 283 let new_mst_root = new_mst 284 .persist() 285 .await 286 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 287 288 let op = RecordOp::Delete { 289 collection: collection.clone(), 290 rkey: rkey.clone(), 291 prev: Some(prev_cid), 292 }; 293 294 let mut new_mst_blocks = std::collections::BTreeMap::new(); 295 let mut old_mst_blocks = std::collections::BTreeMap::new(); 296 297 new_mst 298 .blocks_for_path(&key, &mut new_mst_blocks) 299 .await 300 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 301 302 mst.blocks_for_path(&key, &mut old_mst_blocks) 303 .await 304 .map_err(|e| CommitError::MstOperationFailed(format!("{:?}", e)))?; 305 306 let mut relevant_blocks = new_mst_blocks.clone(); 307 relevant_blocks.extend(old_mst_blocks.iter().map(|(k, v)| (*k, v.clone()))); 308 309 let written_cids: Vec<Cid> = tracking_store 310 .get_all_relevant_cids() 311 .into_iter() 312 .chain(relevant_blocks.keys().copied()) 313 .collect::<std::collections::HashSet<_>>() 314 .into_iter() 315 .collect(); 316 317 let written_cids_str: Vec<String> = written_cids.iter().map(|c| c.to_string()).collect(); 318 319 let obsolete_cids: Vec<Cid> = std::iter::once(current_root_cid) 320 .chain( 321 old_mst_blocks 322 .keys() 323 .filter(|cid| !new_mst_blocks.contains_key(*cid)) 324 .copied(), 325 ) 326 .chain(std::iter::once(prev_cid)) 327 .collect(); 328 329 commit_and_log( 330 state, 331 CommitParams { 332 did, 333 user_id, 334 current_root_cid: Some(current_root_cid), 335 prev_data_cid: Some(commit.data), 336 new_mst_root, 337 ops: vec![op], 338 blocks_cids: &written_cids_str, 339 blobs: &[], 340 obsolete_cids, 341 }, 342 ) 343 .await?; 344 345 Ok(()) 346}