Our Personal Data Server from scratch!
0

Configure Feed

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

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}