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