forked from
tranquil.farm/tranquil-pds
Our Personal Data Server from scratch!
10 kB
307 lines
1use std::collections::HashSet;
2
3use serde::{Deserialize, Serialize};
4
5use super::backlink_ops::remove_backlinks_for_record;
6use super::backlinks::{BacklinkValue, backlink_by_user_key, backlink_key, discriminant_to_path};
7use super::encoding::KeyReader;
8use super::keys::{KeyTag, UserHash};
9use super::records::{RecordValue, record_key};
10use super::repo_meta::{RepoMetaValue, repo_meta_key};
11use super::user_blocks::{user_block_key, user_block_user_prefix};
12use crate::metastore::MetastoreError;
13
14const MUTATION_SET_VERSION: u8 = 1;
15
16#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
17pub struct CommitMutationSet {
18 pub new_root_cid: Vec<u8>,
19 pub new_rev: String,
20 pub record_upserts: Vec<RecordMutationUpsert>,
21 pub record_deletes: Vec<RecordMutationDelete>,
22 pub block_inserts: Vec<Vec<u8>>,
23 pub block_deletes: Vec<Vec<u8>>,
24 pub backlink_adds: Vec<BacklinkMutation>,
25 pub backlink_remove_uris: Vec<String>,
26}
27
28#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
29pub struct RecordMutationUpsert {
30 pub collection: String,
31 pub rkey: String,
32 pub cid_bytes: Vec<u8>,
33}
34
35#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
36pub struct RecordMutationDelete {
37 pub collection: String,
38 pub rkey: String,
39}
40
41#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
42pub struct BacklinkMutation {
43 pub uri: String,
44 pub path: u8,
45 pub link_to: String,
46}
47
48const MAX_MUTATION_SET_ENTRIES: usize = 50_000;
49
50impl CommitMutationSet {
51 pub fn serialize(&self) -> Result<Vec<u8>, MetastoreError> {
52 self.validate_size()?;
53 let payload = postcard::to_allocvec(self)
54 .map_err(|_| MetastoreError::CorruptData("CommitMutationSet serialization failed"))?;
55 let mut buf = Vec::with_capacity(1 + payload.len());
56 buf.push(MUTATION_SET_VERSION);
57 buf.extend_from_slice(&payload);
58 Ok(buf)
59 }
60
61 fn validate_size(&self) -> Result<(), MetastoreError> {
62 let total = self.record_upserts.len()
63 + self.record_deletes.len()
64 + self.block_inserts.len()
65 + self.block_deletes.len()
66 + self.backlink_adds.len()
67 + self.backlink_remove_uris.len();
68 match total <= MAX_MUTATION_SET_ENTRIES {
69 true => Ok(()),
70 false => {
71 tracing::warn!(
72 total_entries = total,
73 max = MAX_MUTATION_SET_ENTRIES,
74 "CommitMutationSet exceeds entry limit"
75 );
76 Err(MetastoreError::InvalidInput(
77 "CommitMutationSet exceeds maximum entry count",
78 ))
79 }
80 }
81 }
82
83 pub fn deserialize(bytes: &[u8]) -> Option<Self> {
84 let (&version, payload) = bytes.split_first()?;
85 match version {
86 MUTATION_SET_VERSION => match postcard::from_bytes(payload) {
87 Ok(v) => Some(v),
88 Err(e) => {
89 tracing::warn!(%e, "failed to deserialize CommitMutationSet payload");
90 None
91 }
92 },
93 _ => {
94 tracing::warn!(version, "unknown CommitMutationSet version");
95 None
96 }
97 }
98 }
99}
100
101pub fn replay_mutation_set(
102 batch: &mut fjall::OwnedWriteBatch,
103 repo_data: &fjall::Keyspace,
104 indexes: &fjall::Keyspace,
105 user_hash: UserHash,
106 current_meta: &RepoMetaValue,
107 mutation_set: &CommitMutationSet,
108) -> Result<(), MetastoreError> {
109 mutation_set.validate_size()?;
110
111 let updated_meta = RepoMetaValue {
112 repo_root_cid: mutation_set.new_root_cid.clone(),
113 repo_rev: mutation_set.new_rev.clone(),
114 ..current_meta.clone()
115 };
116 let meta_key = repo_meta_key(user_hash);
117 batch.insert(repo_data, meta_key.as_slice(), updated_meta.serialize());
118
119 mutation_set.record_upserts.iter().for_each(|u| {
120 let key = record_key(user_hash, &u.collection, &u.rkey);
121 let value = RecordValue {
122 record_cid: u.cid_bytes.clone(),
123 takedown_ref: None,
124 };
125 batch.insert(repo_data, key.as_slice(), value.serialize());
126 });
127
128 mutation_set.record_deletes.iter().for_each(|d| {
129 let key = record_key(user_hash, &d.collection, &d.rkey);
130 batch.remove(repo_data, key.as_slice());
131 });
132
133 mutation_set.block_inserts.iter().for_each(|cid_bytes| {
134 let key = user_block_key(user_hash, &mutation_set.new_rev, cid_bytes);
135 batch.insert(repo_data, key.as_slice(), []);
136 });
137
138 delete_user_blocks_by_cid_scan(batch, repo_data, user_hash, &mutation_set.block_deletes)?;
139
140 mutation_set
141 .backlink_remove_uris
142 .iter()
143 .try_for_each(|uri_str| {
144 let uri = tranquil_types::AtUri::from(uri_str.clone());
145 let collection = uri.collection().ok_or(MetastoreError::CorruptData(
146 "backlink URI missing collection",
147 ))?;
148 let rkey = uri
149 .rkey()
150 .ok_or(MetastoreError::CorruptData("backlink URI missing rkey"))?;
151
152 remove_backlinks_for_record(batch, indexes, user_hash, collection, rkey)
153 })?;
154
155 mutation_set.backlink_adds.iter().try_for_each(|bl| {
156 let uri = tranquil_types::AtUri::from(bl.uri.clone());
157 let collection = uri.collection().ok_or(MetastoreError::CorruptData(
158 "backlink URI missing collection",
159 ))?;
160 let rkey = uri
161 .rkey()
162 .ok_or(MetastoreError::CorruptData("backlink URI missing rkey"))?;
163
164 match discriminant_to_path(bl.path) {
165 None => {
166 tracing::warn!(
167 path = bl.path,
168 uri = %bl.uri,
169 "skipping backlink with unknown path discriminant during recovery"
170 );
171 }
172 Some(_) => {
173 let primary = backlink_key(&bl.link_to, user_hash, collection, rkey);
174 let value = BacklinkValue {
175 source_uri: bl.uri.clone(),
176 path: bl.path,
177 };
178 batch.insert(indexes, primary.as_slice(), value.serialize());
179
180 let reverse = backlink_by_user_key(user_hash, collection, rkey, &bl.link_to);
181 batch.insert(indexes, reverse.as_slice(), []);
182 }
183 }
184 Ok::<_, MetastoreError>(())
185 })
186}
187
188fn delete_user_blocks_by_cid_scan(
189 batch: &mut fjall::OwnedWriteBatch,
190 repo_data: &fjall::Keyspace,
191 user_hash: UserHash,
192 block_cids: &[Vec<u8>],
193) -> Result<(), MetastoreError> {
194 match block_cids.is_empty() {
195 true => Ok(()),
196 false => {
197 let cid_set: HashSet<&[u8]> = block_cids.iter().map(|c| c.as_slice()).collect();
198 let prefix = user_block_user_prefix(user_hash);
199 repo_data.prefix(prefix.as_slice()).try_for_each(|guard| {
200 let (key_bytes, _) = guard.into_inner().map_err(MetastoreError::Fjall)?;
201 match extract_cid_from_user_block_key(&key_bytes) {
202 Some(cid) if cid_set.contains(cid) => {
203 batch.remove(repo_data, key_bytes.as_ref());
204 Ok(())
205 }
206 _ => Ok(()),
207 }
208 })
209 }
210 }
211}
212
213fn extract_cid_from_user_block_key(key_bytes: &[u8]) -> Option<&[u8]> {
214 let mut reader = KeyReader::new(key_bytes);
215 let tag = reader.tag()?;
216
217 if tag != KeyTag::USER_BLOCKS.raw() {
218 tracing::warn!(
219 tag,
220 "unexpected key tag in user_block prefix scan during recovery"
221 );
222 return None;
223 }
224
225 if reader.u64().and_then(|_| reader.string()).is_none() {
226 tracing::warn!("user_block key has corrupt user_hash or rev during recovery");
227 return None;
228 }
229
230 let remaining = reader.remaining();
231 match remaining.is_empty() {
232 true => {
233 tracing::warn!("user_block key has no CID suffix during recovery");
234 None
235 }
236 false => Some(remaining),
237 }
238}
239
240#[cfg(test)]
241mod tests {
242 use super::*;
243
244 #[test]
245 fn mutation_set_roundtrip() {
246 let ms = CommitMutationSet {
247 new_root_cid: vec![0x01, 0x71, 0x12, 0x20],
248 new_rev: "rev1".to_owned(),
249 record_upserts: vec![RecordMutationUpsert {
250 collection: "app.bsky.feed.post".to_owned(),
251 rkey: "3k2abc".to_owned(),
252 cid_bytes: vec![0xDE, 0xAD],
253 }],
254 record_deletes: vec![RecordMutationDelete {
255 collection: "app.bsky.feed.like".to_owned(),
256 rkey: "3k2del".to_owned(),
257 }],
258 block_inserts: vec![vec![0x01, 0x02]],
259 block_deletes: vec![vec![0x03, 0x04]],
260 backlink_adds: vec![BacklinkMutation {
261 uri: "at://did:plc:olaren/app.bsky.feed.like/3k2abc".to_owned(),
262 path: 1,
263 link_to: "at://did:plc:teq/app.bsky.feed.post/3k2xyz".to_owned(),
264 }],
265 backlink_remove_uris: vec!["at://did:plc:olaren/app.bsky.feed.like/3k2old".to_owned()],
266 };
267
268 let bytes = ms.serialize().unwrap();
269 assert_eq!(bytes[0], MUTATION_SET_VERSION);
270 let recovered = CommitMutationSet::deserialize(&bytes).unwrap();
271 assert_eq!(recovered, ms);
272 }
273
274 #[test]
275 fn mutation_set_empty_roundtrip() {
276 let ms = CommitMutationSet {
277 new_root_cid: vec![],
278 new_rev: String::new(),
279 record_upserts: vec![],
280 record_deletes: vec![],
281 block_inserts: vec![],
282 block_deletes: vec![],
283 backlink_adds: vec![],
284 backlink_remove_uris: vec![],
285 };
286
287 let recovered = CommitMutationSet::deserialize(&ms.serialize().unwrap()).unwrap();
288 assert_eq!(recovered, ms);
289 }
290
291 #[test]
292 fn unknown_version_returns_none() {
293 let ms = CommitMutationSet {
294 new_root_cid: vec![],
295 new_rev: String::new(),
296 record_upserts: vec![],
297 record_deletes: vec![],
298 block_inserts: vec![],
299 block_deletes: vec![],
300 backlink_adds: vec![],
301 backlink_remove_uris: vec![],
302 };
303 let mut bytes = ms.serialize().unwrap();
304 bytes[0] = 99;
305 assert!(CommitMutationSet::deserialize(&bytes).is_none());
306 }
307}