forked from
tranquil.farm/tranquil-pds
Our Personal Data Server from scratch!
11 kB
346 lines
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(¤t_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(¤t_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}