This repository has no description
0

Configure Feed

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

core / knot2 / crates / knot-cob / src / lib.rs
66 kB 1999 lines
1mod backend; 2mod change; 3mod error; 4mod graph; 5#[cfg(feature = "instrument")] 6pub mod instrument; 7mod object; 8 9pub use change::{Change, ChangePayload, CobHome, Payload}; 10pub use error::{CobError, PayloadError}; 11pub use graph::{ChangeGraph, History}; 12pub use knot_types::{ActorId, ChangeId, CobId, TypeName}; 13pub use object::{Checkpoint, Evaluate, HistoryModel, Object, SnapshotStride, StateSize}; 14 15pub use backend::parse_cob_ref; 16 17use knot_git::{RefUpdate, Repo}; 18use knot_runtime::Signer; 19use knot_types::{Oid, UnixSeconds}; 20use serde::Serialize; 21use serde::de::DeserializeOwned; 22 23const MAX_CAS_RETRIES: usize = 16; 24const CHECKPOINT_GROWTH_DIVISOR: usize = 16; 25// don't ask 26 27fn checkpoint_stride(snapshot_stride: SnapshotStride, size: StateSize) -> usize { 28 snapshot_stride 29 .get() 30 .max(size.get() / CHECKPOINT_GROWTH_DIVISOR) 31 .min(backend::MAX_GRAPH_CHANGES) 32} 33 34#[derive(Debug, Clone, Copy, PartialEq, Eq)] 35pub struct Created { 36 pub object: CobId, 37 pub tip: ChangeId, 38} 39 40#[derive(Debug)] 41pub struct Delta { 42 pub changes: Vec<Change>, 43 pub tip: ChangeId, 44} 45 46const CHECKPOINT_FORMAT: u16 = 3; 47 48#[derive(Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] 49struct CheckpointDigest(Oid); 50 51fn checkpoint_digest<S: Serialize>( 52 object: CobId, 53 tip: ChangeId, 54 state: &S, 55) -> Result<CheckpointDigest, CobError> { 56 let bytes = serde_ipld_dagcbor::to_vec(&(object, tip, state)) 57 .map_err(|error| CobError::Write(error.to_string()))?; 58 let mut hasher = gix_hash::hasher(gix_hash::Kind::Sha1); 59 hasher.update(&bytes); 60 hasher 61 .try_finalize() 62 .map(|id| CheckpointDigest(Oid::from(id))) 63 .map_err(|error| CobError::Write(error.to_string())) 64} 65 66fn encode_checkpoint<S: Serialize>( 67 object: CobId, 68 tip: ChangeId, 69 state: &S, 70) -> Result<Vec<u8>, CobError> { 71 let digest = checkpoint_digest(object, tip, state)?; 72 serde_ipld_dagcbor::to_vec(&(CHECKPOINT_FORMAT, object, tip, digest, state)) 73 .map_err(|error| CobError::Write(error.to_string())) 74} 75 76fn decode_checkpoint<S: DeserializeOwned + Serialize>( 77 object: CobId, 78 bytes: &[u8], 79) -> Option<(ChangeId, S)> { 80 let (version, decoded_object, tip, digest, state): (u16, CobId, ChangeId, CheckpointDigest, S) = 81 serde_ipld_dagcbor::from_slice(bytes).ok()?; 82 (version == CHECKPOINT_FORMAT).then_some(())?; 83 (decoded_object == object).then_some(())?; 84 (checkpoint_digest(object, tip, &state).ok()? == digest).then_some(())?; 85 Some((tip, state)) 86} 87 88pub struct CobStore<'r> { 89 repo: &'r Repo, 90} 91 92impl<'r> CobStore<'r> { 93 pub fn new(repo: &'r Repo) -> Self { 94 Self { repo } 95 } 96 97 pub fn create<P: ChangePayload>( 98 &self, 99 home: &CobHome, 100 payload: &P, 101 signer: &dyn Signer, 102 timestamp: UnixSeconds, 103 ) -> Result<Created, CobError> { 104 let type_name = P::type_name(); 105 let bytes = payload.encode()?; 106 let tip = backend::write_change( 107 home, 108 self.repo, 109 &type_name, 110 &bytes, 111 &[], 112 None, 113 signer, 114 timestamp, 115 )?; 116 let object = CobId::new(tip.oid()); 117 let name = backend::cob_ref_name(&type_name, object)?; 118 self.repo.update_ref(&RefUpdate::Create { 119 name, 120 new: tip.oid(), 121 })?; 122 Ok(Created { object, tip }) 123 } 124 125 pub fn update<P: ChangePayload>( 126 &self, 127 home: &CobHome, 128 object: CobId, 129 payload: &P, 130 signer: &dyn Signer, 131 timestamp: UnixSeconds, 132 ) -> Result<ChangeId, CobError> { 133 let type_name = P::type_name(); 134 let expected = backend::resolve_tip(self.repo, &type_name, object)? 135 .map(ChangeId::new) 136 .ok_or(CobError::NoSuchObject(object))?; 137 self.append( 138 home, object, &type_name, expected, payload, signer, timestamp, 139 ) 140 } 141 142 pub fn extend<'p, P: ChangePayload + 'p>( 143 &self, 144 home: &CobHome, 145 object: CobId, 146 changes: impl IntoIterator<Item = (&'p P, UnixSeconds)>, 147 signer: &dyn Signer, 148 ) -> Result<Option<ChangeId>, CobError> { 149 let type_name = P::type_name(); 150 let expected = backend::resolve_tip(self.repo, &type_name, object)? 151 .map(ChangeId::new) 152 .ok_or(CobError::NoSuchObject(object))?; 153 self.chain(home, object, &type_name, expected, changes, signer) 154 } 155 156 #[allow(clippy::too_many_arguments)] 157 fn append<P: ChangePayload>( 158 &self, 159 home: &CobHome, 160 object: CobId, 161 type_name: &TypeName, 162 expected: ChangeId, 163 payload: &P, 164 signer: &dyn Signer, 165 timestamp: UnixSeconds, 166 ) -> Result<ChangeId, CobError> { 167 self.chain( 168 home, 169 object, 170 type_name, 171 expected, 172 std::iter::once((payload, timestamp)), 173 signer, 174 ) 175 .map(|tip| tip.expect("one change chains onto one tip")) 176 } 177 178 fn chain<'p, P: ChangePayload + 'p>( 179 &self, 180 home: &CobHome, 181 object: CobId, 182 type_name: &TypeName, 183 expected: ChangeId, 184 changes: impl IntoIterator<Item = (&'p P, UnixSeconds)>, 185 signer: &dyn Signer, 186 ) -> Result<Option<ChangeId>, CobError> { 187 let written = match changes.into_iter().try_fold( 188 Vec::<ChangeId>::new(), 189 |mut written, (payload, timestamp)| { 190 let parent = written.last().copied().unwrap_or(expected); 191 let write = || -> Result<ChangeId, CobError> { 192 let bytes = payload.encode()?; 193 backend::write_change( 194 home, 195 self.repo, 196 type_name, 197 &bytes, 198 &[parent], 199 Some(object), 200 signer, 201 timestamp, 202 ) 203 }; 204 match write() { 205 Ok(tip) => { 206 written.push(tip); 207 Ok(written) 208 } 209 Err(error) => Err((written, error)), 210 } 211 }, 212 ) { 213 Ok(written) => written, 214 Err((partial, error)) => { 215 partial 216 .iter() 217 .try_for_each(|change| self.repo.remove_loose_object(change.oid()))?; 218 return Err(error); 219 } 220 }; 221 written.last().copied().map_or(Ok(None), |tip| { 222 self.publish(type_name, object, expected, tip, &written) 223 .map(Some) 224 }) 225 } 226 227 fn publish( 228 &self, 229 type_name: &TypeName, 230 object: CobId, 231 expected: ChangeId, 232 tip: ChangeId, 233 written: &[ChangeId], 234 ) -> Result<ChangeId, CobError> { 235 let name = backend::cob_ref_name(type_name, object)?; 236 match self.repo.update_ref(&RefUpdate::Update { 237 name, 238 old: expected.oid(), 239 new: tip.oid(), 240 }) { 241 Ok(()) => Ok(tip), 242 Err(error) => { 243 let actual = backend::resolve_tip(self.repo, type_name, object)?; 244 if actual != Some(tip.oid()) { 245 written 246 .iter() 247 .try_for_each(|change| self.repo.remove_loose_object(change.oid()))?; 248 } 249 match actual { 250 actual if actual != Some(expected.oid()) => { 251 Err(CobError::StaleTip { object, expected }) 252 } 253 _ => Err(CobError::Git(error)), 254 } 255 } 256 } 257 } 258 259 pub fn update_with<E, D>( 260 &self, 261 home: &CobHome, 262 object: CobId, 263 signer: &dyn Signer, 264 timestamp: UnixSeconds, 265 decide: impl Fn(&E::State) -> Result<E::Change, D>, 266 ) -> Result<ChangeId, D> 267 where 268 E: Evaluate, 269 D: From<CobError>, 270 { 271 self.update_maybe::<E, D>(home, object, signer, timestamp, |state| { 272 decide(state).map(Some) 273 }) 274 .map(|tip| tip.expect("update_with always yields change to append")) 275 } 276 277 pub fn update_maybe<E, D>( 278 &self, 279 home: &CobHome, 280 object: CobId, 281 signer: &dyn Signer, 282 timestamp: UnixSeconds, 283 decide: impl Fn(&E::State) -> Result<Option<E::Change>, D>, 284 ) -> Result<Option<ChangeId>, D> 285 where 286 E: Evaluate, 287 D: From<CobError>, 288 { 289 let type_name = E::Change::type_name(); 290 let attempt = || -> Result<Option<Option<ChangeId>>, D> { 291 let (graph, expected) = backend::load_graph( 292 self.repo, 293 &type_name, 294 object, 295 E::HISTORY, 296 backend::MAX_GRAPH_CHANGES, 297 )?; 298 let (state, _history) = object::evaluate::<E>(graph, &type_name)?; 299 match decide(&state)? { 300 None => Ok(Some(None)), 301 Some(change) => { 302 match self.append( 303 home, object, &type_name, expected, &change, signer, timestamp, 304 ) { 305 Ok(id) => Ok(Some(Some(id))), 306 Err(CobError::StaleTip { .. }) => Ok(None), 307 Err(other) => Err(D::from(other)), 308 } 309 } 310 } 311 }; 312 (0..MAX_CAS_RETRIES) 313 .find_map(|_| attempt().transpose()) 314 .unwrap_or_else(|| Err(D::from(CobError::Contended(object)))) 315 } 316 317 pub fn graph<E: Evaluate>(&self, object: CobId) -> Result<ChangeGraph, CobError> { 318 Ok(backend::load_graph( 319 self.repo, 320 &E::Change::type_name(), 321 object, 322 E::HISTORY, 323 backend::MAX_GRAPH_CHANGES, 324 )? 325 .0) 326 } 327 328 pub fn get<E: Evaluate>(&self, object: CobId) -> Result<Object<E::State>, CobError> { 329 let type_name = E::Change::type_name(); 330 let (graph, _tip) = backend::load_graph( 331 self.repo, 332 &type_name, 333 object, 334 E::HISTORY, 335 backend::MAX_GRAPH_CHANGES, 336 )?; 337 let (state, history) = object::evaluate::<E>(graph, &type_name)?; 338 Ok(Object::new(object, type_name, state, history)) 339 } 340 341 pub fn verify<E: Evaluate>( 342 &self, 343 home: &CobHome, 344 object: CobId, 345 owner: &ActorId, 346 ) -> Result<(), CobError> { 347 let type_name = E::Change::type_name(); 348 let (graph, _tip) = backend::load_graph( 349 self.repo, 350 &type_name, 351 object, 352 E::HISTORY, 353 backend::MAX_GRAPH_CHANGES, 354 )?; 355 graph.into_ordered().into_iter().try_for_each(|change| { 356 if change.type_name != type_name { 357 return Err(CobError::UnexpectedChangeType { 358 change: change.id, 359 expected: type_name.clone(), 360 found: change.type_name, 361 }); 362 } 363 change 364 .verify(home, owner, Some(object)) 365 .then_some(()) 366 .ok_or(CobError::UnverifiedChange { change: change.id }) 367 }) 368 } 369 370 pub fn list<E: Evaluate>(&self) -> Result<Vec<CobId>, CobError> { 371 backend::list_objects(self.repo, &E::Change::type_name()) 372 } 373 374 pub fn changes_since<E: Evaluate>( 375 &self, 376 object: CobId, 377 since: Option<ChangeId>, 378 ) -> Result<Delta, CobError> { 379 let type_name = E::Change::type_name(); 380 let tip = backend::resolve_tip(self.repo, &type_name, object)? 381 .map(ChangeId::new) 382 .ok_or(CobError::NoSuchObject(object))?; 383 let collected = 384 backend::collect(self.repo, tip, object, backend::MAX_GRAPH_CHANGES, since)?; 385 match since { 386 None => backend::check_full_shape(&collected, object, E::HISTORY)?, 387 Some(since) => backend::check_delta_shape(&collected, object, since, E::HISTORY)?, 388 } 389 let changes = ChangeGraph::new(object, collected).into_ordered(); 390 Ok(Delta { changes, tip }) 391 } 392 393 pub fn update_with_checkpointed<E, D>( 394 &self, 395 home: &CobHome, 396 object: CobId, 397 signer: &dyn Signer, 398 timestamp: UnixSeconds, 399 decide: impl Fn(&E::State) -> Result<E::Change, D>, 400 ) -> Result<ChangeId, D> 401 where 402 E: Checkpoint, 403 E::State: Serialize + DeserializeOwned, 404 D: From<CobError>, 405 { 406 self.update_maybe_checkpointed::<E, D>(home, object, signer, timestamp, |state| { 407 decide(state).map(Some) 408 }) 409 .map(|tip| tip.expect("update_with_checkpointed always yields change to append")) 410 } 411 412 pub fn update_maybe_checkpointed<E, D>( 413 &self, 414 home: &CobHome, 415 object: CobId, 416 signer: &dyn Signer, 417 timestamp: UnixSeconds, 418 decide: impl Fn(&E::State) -> Result<Option<E::Change>, D>, 419 ) -> Result<Option<ChangeId>, D> 420 where 421 E: Checkpoint, 422 E::State: Serialize + DeserializeOwned, 423 D: From<CobError>, 424 { 425 let type_name = E::Change::type_name(); 426 let attempt = || -> Result<Option<Option<ChangeId>>, D> { 427 let (state, expected, suffix) = self.checkpointed_state::<E>(object)?; 428 match decide(&state)? { 429 None => Ok(Some(None)), 430 Some(change) => match self.append( 431 home, object, &type_name, expected, &change, signer, timestamp, 432 ) { 433 Ok(tip) => { 434 let stride = 435 checkpoint_stride(E::SNAPSHOT_STRIDE, E::checkpoint_size(&state)); 436 if suffix.saturating_add(1) >= stride { 437 let author = ActorId::from_secp256k1(signer.public_key().as_bytes()); 438 let folded = E::apply(state, change, &author); 439 if let Err(error) = self.write_checkpoint::<E>(object, tip, &folded) { 440 tracing::warn!( 441 cob = type_name.as_str(), 442 object = %object.oid().to_hex(), 443 %error, 444 "checkpoint write failed" 445 ); 446 } 447 } 448 Ok(Some(Some(tip))) 449 } 450 Err(CobError::StaleTip { .. }) => Ok(None), 451 Err(other) => Err(D::from(other)), 452 }, 453 } 454 }; 455 (0..MAX_CAS_RETRIES) 456 .find_map(|_| attempt().transpose()) 457 .unwrap_or_else(|| Err(D::from(CobError::Contended(object)))) 458 } 459 460 pub fn materialize<E>(&self, object: CobId) -> Result<(E::State, ChangeId), CobError> 461 where 462 E: Checkpoint, 463 E::State: Serialize + DeserializeOwned, 464 { 465 self.checkpointed_state::<E>(object) 466 .map(|(state, tip, _)| (state, tip)) 467 } 468 469 fn checkpointed_state<E>(&self, object: CobId) -> Result<(E::State, ChangeId, usize), CobError> 470 where 471 E: Checkpoint, 472 E::State: Serialize + DeserializeOwned, 473 { 474 let type_name = E::Change::type_name(); 475 let tip = backend::resolve_tip(self.repo, &type_name, object)? 476 .map(ChangeId::new) 477 .ok_or(CobError::NoSuchObject(object))?; 478 match self.load_checkpoint::<E>(object)? { 479 Some((checkpoint_tip, state)) if checkpoint_tip == tip => Ok((state, tip, 0)), 480 Some((checkpoint_tip, state)) => { 481 match self.changes_since::<E>(object, Some(checkpoint_tip)) { 482 Ok(delta) => { 483 let folded = object::fold_changes::<E>(state, &delta.changes, &type_name)?; 484 Ok((folded, tip, delta.changes.len())) 485 } 486 Err(_) => self 487 .full_state::<E>(object) 488 .map(|state| (state, tip, usize::MAX)), 489 } 490 } 491 None => self 492 .full_state::<E>(object) 493 .map(|state| (state, tip, usize::MAX)), 494 } 495 } 496 497 fn full_state<E: Evaluate>(&self, object: CobId) -> Result<E::State, CobError> { 498 let type_name = E::Change::type_name(); 499 let (graph, _tip) = backend::load_graph( 500 self.repo, 501 &type_name, 502 object, 503 E::HISTORY, 504 backend::rebuild_graph_limit(), 505 )?; 506 object::evaluate::<E>(graph, &type_name).map(|(state, _history)| state) 507 } 508 509 fn load_checkpoint<E>(&self, object: CobId) -> Result<Option<(ChangeId, E::State)>, CobError> 510 where 511 E: Checkpoint, 512 E::State: DeserializeOwned + Serialize, 513 { 514 let name = backend::checkpoint_ref_name(&E::Change::type_name(), object)?; 515 let Some(blob) = self.repo.find_ref(&name)? else { 516 return Ok(None); 517 }; 518 let Ok(bytes) = self.repo.read_blob(blob) else { 519 return Ok(None); 520 }; 521 Ok(decode_checkpoint::<E::State>(object, &bytes)) 522 } 523 524 fn write_checkpoint<E>( 525 &self, 526 object: CobId, 527 tip: ChangeId, 528 state: &E::State, 529 ) -> Result<(), CobError> 530 where 531 E: Checkpoint, 532 E::State: Serialize, 533 { 534 let bytes = encode_checkpoint(object, tip, state)?; 535 let blob = Oid::from( 536 self.repo 537 .git() 538 .write_blob(&bytes) 539 .map_err(|error| CobError::Write(error.to_string()))? 540 .detach(), 541 ); 542 let name = backend::checkpoint_ref_name(&E::Change::type_name(), object)?; 543 let update = match self.repo.find_ref(&name)? { 544 Some(old) => RefUpdate::Update { 545 name, 546 old, 547 new: blob, 548 }, 549 None => RefUpdate::Create { name, new: blob }, 550 }; 551 self.repo.update_ref(&update).map_err(CobError::from) 552 } 553} 554 555#[cfg(test)] 556mod tests { 557 use std::collections::BTreeSet; 558 559 use knot_git::{Layout, Repo}; 560 use knot_runtime::{K256Signer, SeededEntropy, Signature, Signer}; 561 use knot_types::{RepoDid, UnixSeconds}; 562 use proptest::prelude::*; 563 use serde::{Deserialize, Serialize}; 564 use tempfile::TempDir; 565 566 use super::*; 567 568 #[derive(Debug, Serialize, Deserialize)] 569 #[serde(tag = "op", content = "subject")] 570 enum Tag { 571 Add(String), 572 Remove(String), 573 } 574 575 impl ChangePayload for Tag { 576 const TYPE: &'static str = "sh.tangled.test.tag"; 577 } 578 579 #[test] 580 fn checkpoint_encoding_stays_byte_stable() { 581 let object = CobId::new(Oid::from_hex("aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa").unwrap()); 582 let tip = ChangeId::new(Oid::from_hex("bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb").unwrap()); 583 let mut state = BTreeSet::new(); 584 state.insert("kelp".to_string()); 585 state.insert("squid".to_string()); 586 let bytes = encode_checkpoint(object, tip, &state).unwrap(); 587 assert_eq!( 588 knot_types::lowercase_hex(&bytes), 589 "850378286161616161616161616161616161616161616161616161616161616161616161616161616161616178286262626262626262626262626262626262626262626262626262626262626262626262626262626278283861623531376164633234616530313138383962643634343930663836633562646333333632333782646b656c70657371756964" 590 ); 591 } 592 593 struct Tags; 594 595 impl Evaluate for Tags { 596 type State = BTreeSet<String>; 597 type Change = Tag; 598 599 const HISTORY: HistoryModel = HistoryModel::Convergent; 600 601 fn initial() -> Self::State { 602 BTreeSet::new() 603 } 604 605 fn apply(mut state: Self::State, change: Self::Change, _author: &ActorId) -> Self::State { 606 match change { 607 Tag::Add(subject) => { 608 state.insert(subject); 609 } 610 Tag::Remove(subject) => { 611 state.remove(&subject); 612 } 613 } 614 state 615 } 616 } 617 618 #[derive(Debug, Serialize, Deserialize)] 619 #[serde(tag = "op", content = "subject")] 620 enum NameChange { 621 Claim(String), 622 } 623 624 impl ChangePayload for NameChange { 625 const TYPE: &'static str = "sh.tangled.test.name"; 626 } 627 628 struct Names; 629 630 impl Evaluate for Names { 631 type State = BTreeSet<String>; 632 type Change = NameChange; 633 634 const HISTORY: HistoryModel = HistoryModel::Linear; 635 636 fn initial() -> Self::State { 637 BTreeSet::new() 638 } 639 640 fn apply(mut state: Self::State, change: Self::Change, _author: &ActorId) -> Self::State { 641 match change { 642 NameChange::Claim(subject) => { 643 state.insert(subject); 644 } 645 } 646 state 647 } 648 } 649 650 impl Checkpoint for Names { 651 const SNAPSHOT_STRIDE: SnapshotStride = SnapshotStride::new(4); 652 fn checkpoint_size(state: &Self::State) -> StateSize { 653 StateSize::new(state.len()) 654 } 655 } 656 657 #[derive(Debug)] 658 enum NameError { 659 Taken, 660 Cob(CobError), 661 } 662 663 impl From<CobError> for NameError { 664 fn from(error: CobError) -> Self { 665 NameError::Cob(error) 666 } 667 } 668 669 fn claim( 670 store: &CobStore, 671 object: CobId, 672 key: &K256Signer, 673 who: &str, 674 when: i64, 675 ) -> Result<Option<ChangeId>, NameError> { 676 store.update_maybe_checkpointed::<Names, NameError>( 677 &cob_home(), 678 object, 679 key, 680 at(when), 681 |state| match state.contains(who) { 682 true => Err(NameError::Taken), 683 false => Ok(Some(NameChange::Claim(who.to_string()))), 684 }, 685 ) 686 } 687 688 fn appended(result: Result<Option<ChangeId>, NameError>) -> ChangeId { 689 match result { 690 Ok(Some(id)) => id, 691 Ok(None) => panic!("distinct claim should append a change"), 692 Err(NameError::Taken) => panic!("name was unexpectedly already taken"), 693 Err(NameError::Cob(error)) => panic!("checkpointed write failed: {error:?}"), 694 } 695 } 696 697 fn fixture() -> (TempDir, Repo) { 698 let dir = tempfile::tempdir().unwrap(); 699 let layout = Layout::new(dir.path()); 700 let repo = layout 701 .create(&RepoDid::new("did:plc:squid").unwrap()) 702 .unwrap(); 703 (dir, repo) 704 } 705 706 fn signer(seed: u64) -> K256Signer { 707 K256Signer::generate(&SeededEntropy::new(seed)) 708 } 709 710 fn cob_home() -> CobHome { 711 CobHome::from(&RepoDid::new("did:plc:squid").unwrap()) 712 } 713 714 fn at(seconds: i64) -> UnixSeconds { 715 UnixSeconds::new(seconds) 716 } 717 718 fn tag_root(repo: &Repo, key: &K256Signer, subject: &str, when: i64) -> ChangeId { 719 backend::write_change( 720 &cob_home(), 721 repo, 722 &Tag::type_name(), 723 &Tag::Add(subject.into()).encode().unwrap(), 724 &[], 725 None, 726 key, 727 at(when), 728 ) 729 .unwrap() 730 } 731 732 fn tag_child( 733 repo: &Repo, 734 key: &K256Signer, 735 subject: &str, 736 parents: &[ChangeId], 737 object: CobId, 738 when: i64, 739 ) -> ChangeId { 740 backend::write_change( 741 &cob_home(), 742 repo, 743 &Tag::type_name(), 744 &Tag::Add(subject.into()).encode().unwrap(), 745 parents, 746 Some(object), 747 key, 748 at(when), 749 ) 750 .unwrap() 751 } 752 753 fn publish(repo: &Repo, object: CobId, tip: ChangeId) { 754 let name = backend::cob_ref_name(&Tag::type_name(), object).unwrap(); 755 repo.update_ref(&RefUpdate::Create { 756 name, 757 new: tip.oid(), 758 }) 759 .unwrap(); 760 } 761 762 fn forked_tag_object( 763 repo: &Repo, 764 key: &K256Signer, 765 left: &str, 766 right: &str, 767 merge: &str, 768 ) -> CobId { 769 let root = tag_root(repo, key, "base", 1); 770 let object = CobId::new(root.oid()); 771 let left = tag_child(repo, key, left, &[root], object, 2); 772 let right = tag_child(repo, key, right, &[root], object, 3); 773 let merge = tag_child(repo, key, merge, &[left, right], object, 4); 774 publish(repo, object, merge); 775 object 776 } 777 778 fn linear_chain(repo: &Repo, key: &K256Signer, subjects: &[&str]) -> CobId { 779 let (first, rest) = subjects.split_first().expect("chain needs a root subject"); 780 let root = tag_root(repo, key, first, 1); 781 let object = CobId::new(root.oid()); 782 let tip = rest 783 .iter() 784 .enumerate() 785 .fold(root, |parent, (index, subject)| { 786 tag_child(repo, key, subject, &[parent], object, index as i64 + 2) 787 }); 788 publish(repo, object, tip); 789 object 790 } 791 792 #[test] 793 fn extend_chains_the_same_tip_as_sequential_updates() { 794 let (_sequential_dir, sequential_repo) = fixture(); 795 let (_batched_dir, batched_repo) = fixture(); 796 let key = signer(3); 797 let changes = [ 798 NameChange::Claim("kelp".into()), 799 NameChange::Claim("squid".into()), 800 NameChange::Claim("whelk".into()), 801 ]; 802 let sequential = CobStore::new(&sequential_repo); 803 let batched = CobStore::new(&batched_repo); 804 let root = |store: &CobStore| { 805 store 806 .create(&cob_home(), &NameChange::Claim("uni".into()), &key, at(0)) 807 .unwrap() 808 .object 809 }; 810 let one_by_one = root(&sequential); 811 changes.iter().zip(1..).for_each(|(change, when)| { 812 sequential 813 .update(&cob_home(), one_by_one, change, &key, at(when)) 814 .unwrap(); 815 }); 816 let in_one_go = root(&batched); 817 let tip = batched 818 .extend( 819 &cob_home(), 820 in_one_go, 821 changes.iter().zip((1..).map(at)), 822 &key, 823 ) 824 .unwrap(); 825 826 let resolved = |repo: &Repo, object| { 827 backend::resolve_tip(repo, &NameChange::type_name(), object) 828 .unwrap() 829 .map(ChangeId::new) 830 }; 831 assert_eq!(one_by_one, in_one_go); 832 assert_eq!(tip, resolved(&batched_repo, in_one_go)); 833 assert_eq!(tip, resolved(&sequential_repo, one_by_one)); 834 assert_eq!( 835 sequential.materialize::<Names>(one_by_one).unwrap().0, 836 batched.materialize::<Names>(in_one_go).unwrap().0 837 ); 838 assert_eq!( 839 batched 840 .extend::<NameChange>(&cob_home(), in_one_go, std::iter::empty(), &key) 841 .unwrap(), 842 None 843 ); 844 assert_eq!(tip, resolved(&batched_repo, in_one_go)); 845 } 846 847 #[test] 848 fn checkpointed_writes_match_a_full_fold_and_leave_a_snapshot() { 849 let (_dir, repo) = fixture(); 850 let store = CobStore::new(&repo); 851 let key = signer(1); 852 let object = store 853 .create( 854 &cob_home(), 855 &NameChange::Claim("name0000".into()), 856 &key, 857 at(0), 858 ) 859 .unwrap() 860 .object; 861 862 let total = Names::SNAPSHOT_STRIDE.get() * 3; 863 (1..total).for_each(|index| { 864 appended(claim( 865 &store, 866 object, 867 &key, 868 &format!("name{index:04}"), 869 index as i64, 870 )); 871 }); 872 873 let folded = store.get::<Names>(object).unwrap(); 874 assert_eq!( 875 folded.state().len(), 876 total, 877 "every checkpointed write landed and full fold agrees" 878 ); 879 let snapshot = backend::checkpoint_ref_name(&NameChange::type_name(), object).unwrap(); 880 assert!( 881 repo.find_ref(&snapshot).unwrap().is_some(), 882 "snapshot ref was written once stride was crossed" 883 ); 884 } 885 886 #[test] 887 fn checkpoint_stride_grows_with_state_so_large_cobs_snapshot_less_often() { 888 assert_eq!( 889 checkpoint_stride(SnapshotStride::new(256), StateSize::new(0)), 890 256 891 ); 892 assert_eq!( 893 checkpoint_stride( 894 SnapshotStride::new(256), 895 StateSize::new(256 * CHECKPOINT_GROWTH_DIVISOR) 896 ), 897 256, 898 "up to stride*divisor entries the fixed floor governs, so small cobs snapshot exactly as before" 899 ); 900 assert_eq!( 901 checkpoint_stride(SnapshotStride::new(256), StateSize::new(1_600_000)), 902 100_000, 903 "a large cob snapshots at a fixed fraction of its size, so total snapshot bytes stay linear in the change count" 904 ); 905 assert_eq!( 906 checkpoint_stride(SnapshotStride::new(256), StateSize::new(4_000_000)), 907 backend::MAX_GRAPH_CHANGES, 908 "the stride stops at the serving limit so the boot tail always folds incrementally, never overflowing into a full rebuild" 909 ); 910 } 911 912 #[test] 913 fn a_large_checkpointed_cob_bounds_its_boot_suffix_to_a_fraction_of_its_size() { 914 let (_dir, repo) = fixture(); 915 let store = CobStore::new(&repo); 916 let key = signer(2); 917 let object = store 918 .create( 919 &cob_home(), 920 &NameChange::Claim("name0000".into()), 921 &key, 922 at(0), 923 ) 924 .unwrap() 925 .object; 926 927 let total = Names::SNAPSHOT_STRIDE.get() * CHECKPOINT_GROWTH_DIVISOR * 4; 928 (1..total).for_each(|index| { 929 claim( 930 &store, 931 object, 932 &key, 933 &format!("name{index:04}"), 934 index as i64, 935 ) 936 .unwrap(); 937 }); 938 939 assert_eq!( 940 store.get::<Names>(object).unwrap().state().len(), 941 total, 942 "full fold still agrees after adaptive snapshotting" 943 ); 944 945 let (_state, _tip, suffix) = store.checkpointed_state::<Names>(object).unwrap(); 946 let bound = checkpoint_stride(Names::SNAPSHOT_STRIDE, StateSize::new(total)); 947 assert!( 948 suffix <= bound && bound < total, 949 "boot replays only the {suffix}-change tail since the last snapshot, bounded by the adaptive stride {bound}, never the whole {total}-change history" 950 ); 951 } 952 953 #[test] 954 fn a_corrupt_checkpoint_never_serves_a_wrong_answer_and_heals_at_the_live_tip() { 955 type Forge = fn(CobId, Oid, &BTreeSet<String>) -> Vec<u8>; 956 let cases: &[Forge] = &[ 957 |_object, _tip, _state| b"not a valid checkpoint envelope".to_vec(), 958 |object, tip, state| { 959 serde_ipld_dagcbor::to_vec(&( 960 CHECKPOINT_FORMAT + 1, 961 object.oid().to_hex(), 962 tip.to_hex(), 963 checkpoint_digest(object, ChangeId::new(tip), state).unwrap(), 964 state, 965 )) 966 .unwrap() 967 }, 968 |object, tip, state| { 969 let mut tampered = state.clone(); 970 tampered.insert("intruder".to_string()); 971 serde_ipld_dagcbor::to_vec(&( 972 CHECKPOINT_FORMAT, 973 object.oid().to_hex(), 974 tip.to_hex(), 975 checkpoint_digest(object, ChangeId::new(tip), state).unwrap(), 976 tampered, 977 )) 978 .unwrap() 979 }, 980 ]; 981 982 cases.iter().enumerate().for_each(|(index, forge)| { 983 let (_dir, repo) = fixture(); 984 let store = CobStore::new(&repo); 985 let key = signer(200 + index as u64); 986 let object = store 987 .create(&cob_home(), &NameChange::Claim("squid".into()), &key, at(0)) 988 .unwrap() 989 .object; 990 (1..=Names::SNAPSHOT_STRIDE.get()).for_each(|step| { 991 claim(&store, object, &key, &format!("name{step:04}"), step as i64).unwrap(); 992 }); 993 994 let real_state = store.get::<Names>(object).unwrap().state().clone(); 995 let real_tip = backend::resolve_tip(&repo, &NameChange::type_name(), object) 996 .unwrap() 997 .unwrap(); 998 let forged = forge(object, real_tip, &real_state); 999 1000 let snapshot = backend::checkpoint_ref_name(&NameChange::type_name(), object).unwrap(); 1001 let live = repo.find_ref(&snapshot).unwrap().expect("snapshot exists"); 1002 let blob = Oid::from(repo.git().write_blob(&forged).unwrap().detach()); 1003 repo.update_ref(&RefUpdate::Update { 1004 name: snapshot.clone(), 1005 old: live, 1006 new: blob, 1007 }) 1008 .unwrap(); 1009 1010 let duplicate = claim(&store, object, &key, "squid", 100); 1011 assert!( 1012 matches!(duplicate, Err(NameError::Taken)), 1013 "corrupt checkpoint falls back to real state instead of waving a duplicate through" 1014 ); 1015 appended(claim(&store, object, &key, "intruder", 101)); 1016 1017 let healed = repo.find_ref(&snapshot).unwrap().unwrap(); 1018 let bytes = repo.read_blob(healed).unwrap(); 1019 let (tip, state) = decode_checkpoint::<BTreeSet<String>>(object, &bytes) 1020 .expect("healed snapshot decodes"); 1021 let live_tip = backend::resolve_tip(&repo, &NameChange::type_name(), object) 1022 .unwrap() 1023 .unwrap(); 1024 assert_eq!( 1025 tip, 1026 ChangeId::new(live_tip), 1027 "healed snapshot sits at the live tip" 1028 ); 1029 assert_eq!( 1030 &state, 1031 store.get::<Names>(object).unwrap().state(), 1032 "healed snapshot matches the full fold" 1033 ); 1034 }); 1035 } 1036 1037 #[test] 1038 fn a_checkpoint_envelope_is_bound_to_its_object_tip_and_format() { 1039 let (_dir, repo) = fixture(); 1040 let store = CobStore::new(&repo); 1041 let key = signer(6); 1042 let a = store 1043 .create(&cob_home(), &NameChange::Claim("squid".into()), &key, at(0)) 1044 .unwrap(); 1045 let b = store 1046 .create( 1047 &cob_home(), 1048 &NameChange::Claim("anemone".into()), 1049 &key, 1050 at(1), 1051 ) 1052 .unwrap(); 1053 let state = BTreeSet::from(["squid".to_string()]); 1054 1055 let bytes = encode_checkpoint(a.object, a.tip, &state).unwrap(); 1056 assert!( 1057 decode_checkpoint::<BTreeSet<String>>(b.object, &bytes).is_none(), 1058 "checkpoint minted for one object mustn't decode under another" 1059 ); 1060 assert_eq!( 1061 decode_checkpoint::<BTreeSet<String>>(a.object, &bytes), 1062 Some((a.tip, state.clone())), 1063 "checkpoint decodes to the tip and state its digest commits to under its own object" 1064 ); 1065 1066 let wrong_tip = ChangeId::new(Oid::from_hex(&"b".repeat(40)).unwrap()); 1067 let swapped_tip = serde_ipld_dagcbor::to_vec(&( 1068 CHECKPOINT_FORMAT, 1069 a.object.oid().to_hex(), 1070 wrong_tip.oid().to_hex(), 1071 checkpoint_digest(a.object, a.tip, &state).unwrap(), 1072 state.clone(), 1073 )) 1074 .unwrap(); 1075 assert!( 1076 decode_checkpoint::<BTreeSet<String>>(a.object, &swapped_tip).is_none(), 1077 "checkpoint whose tip is swapped away from the one its digest covers mustn't decode" 1078 ); 1079 1080 let future_version = serde_ipld_dagcbor::to_vec(&( 1081 CHECKPOINT_FORMAT + 1, 1082 a.object.oid().to_hex(), 1083 a.tip.oid().to_hex(), 1084 checkpoint_digest(a.object, a.tip, &state).unwrap(), 1085 state, 1086 )) 1087 .unwrap(); 1088 assert!( 1089 decode_checkpoint::<BTreeSet<String>>(a.object, &future_version).is_none(), 1090 "envelope whose format version isn't the current one mustn't decode" 1091 ); 1092 } 1093 1094 #[test] 1095 fn the_rebuild_walk_floors_at_the_serving_limit_and_scales_above_it_with_headroom() { 1096 assert_eq!( 1097 backend::rebuild_change_limit_for(None), 1098 2_000_000, 1099 "an unmeasurable host rebuilds up to the generous fixed ceiling" 1100 ); 1101 let available = |bytes| Some(knot_resource::AvailableBytes::new(bytes)); 1102 assert_eq!( 1103 backend::rebuild_change_limit_for(available(64 * 1024 * 1024)), 1104 backend::MAX_GRAPH_CHANGES, 1105 "a squeezed host floors the rebuild walk at the serving limit, never beneath it" 1106 ); 1107 assert_eq!( 1108 backend::rebuild_change_limit_for(available(0)), 1109 backend::MAX_GRAPH_CHANGES, 1110 "a host with no measured headroom still heals at least as deep as it serves" 1111 ); 1112 let one_gib = backend::rebuild_change_limit_for(available(1024 * 1024 * 1024)); 1113 assert!( 1114 (100_000..150_000).contains(&one_gib), 1115 "on a 1 GiB host the fold budget tracks the measured per-change cost near the serving limit, was {one_gib}" 1116 ); 1117 let roomy = backend::rebuild_change_limit_for(available(64 * 1024 * 1024 * 1024)); 1118 assert_eq!(roomy, 64 * 1024 * 1024 * 1024 / 4 / 2560); 1119 assert!( 1120 roomy > backend::MAX_GRAPH_CHANGES, 1121 "a roomy host rebuilds far past the serving limit, was {roomy}" 1122 ); 1123 } 1124 1125 #[test] 1126 fn tag_lifecycle_roundtrips_appends_reloads_and_keeps_signatures() { 1127 let (_dir, repo) = fixture(); 1128 let store = CobStore::new(&repo); 1129 let key = signer(1); 1130 1131 let created = store 1132 .create(&cob_home(), &Tag::Add("nel".into()), &key, at(1)) 1133 .unwrap(); 1134 let fresh = store.get::<Tags>(created.object).unwrap(); 1135 assert_eq!(fresh.id(), created.object); 1136 assert_eq!(fresh.state(), &BTreeSet::from(["nel".to_string()])); 1137 assert_eq!(fresh.history().len(), 1); 1138 assert_eq!(fresh.history().root(), created.tip); 1139 let root = backend::read_change(&repo, ChangeId::new(created.object.oid())).unwrap(); 1140 assert!(root.verify(&cob_home(), &root.author, None)); 1141 1142 store 1143 .update( 1144 &cob_home(), 1145 created.object, 1146 &Tag::Add("olaren".into()), 1147 &key, 1148 at(2), 1149 ) 1150 .unwrap(); 1151 let tip = store 1152 .update( 1153 &cob_home(), 1154 created.object, 1155 &Tag::Remove("nel".into()), 1156 &key, 1157 at(3), 1158 ) 1159 .unwrap(); 1160 1161 let object = store.get::<Tags>(created.object).unwrap(); 1162 assert_eq!(object.state(), &BTreeSet::from(["olaren".to_string()])); 1163 assert_eq!(object.history().len(), 3); 1164 assert_eq!(store.list::<Tags>().unwrap(), vec![created.object]); 1165 let order = object.history().traverse(Vec::new(), |mut acc, change| { 1166 acc.push(change.timestamp.get()); 1167 acc 1168 }); 1169 assert_eq!(order, vec![1, 2, 3]); 1170 assert_eq!(object.history().tips(), vec![tip]); 1171 1172 let reloaded = store.get::<Tags>(created.object).unwrap(); 1173 assert_eq!(object.state(), reloaded.state()); 1174 let graph = store.graph::<Tags>(created.object).unwrap(); 1175 let again = store.graph::<Tags>(created.object).unwrap(); 1176 assert_eq!(graph.len(), again.len()); 1177 assert_eq!(graph.causal_order(), again.causal_order()); 1178 1179 let head = backend::read_change(&repo, tip).unwrap(); 1180 assert!(head.verify(&cob_home(), &head.author, Some(created.object))); 1181 assert!(!head.verify(&cob_home(), &head.author, None)); 1182 assert!(!head.verify( 1183 &cob_home(), 1184 &head.author, 1185 Some(CobId::new(knot_types::Oid::null())) 1186 )); 1187 1188 let absent = CobId::new(knot_types::Oid::from_hex(&"0".repeat(40)).unwrap()); 1189 assert!(matches!( 1190 store.get::<Tags>(absent), 1191 Err(CobError::NoSuchObject(_)) 1192 )); 1193 } 1194 1195 #[test] 1196 fn undecodable_change_poisons_the_object() { 1197 let (_dir, repo) = fixture(); 1198 let key = signer(5); 1199 let nsid = Tag::type_name(); 1200 1201 let root = backend::write_change( 1202 &cob_home(), 1203 &repo, 1204 &nsid, 1205 &Tag::Add("nel".into()).encode().unwrap(), 1206 &[], 1207 None, 1208 &key, 1209 at(1), 1210 ) 1211 .unwrap(); 1212 let object = CobId::new(root.oid()); 1213 let garbage = backend::write_change( 1214 &cob_home(), 1215 &repo, 1216 &nsid, 1217 &[0xff, 0xff, 0xff], 1218 &[root], 1219 Some(object), 1220 &key, 1221 at(2), 1222 ) 1223 .unwrap(); 1224 let name = backend::cob_ref_name(&nsid, object).unwrap(); 1225 repo.update_ref(&RefUpdate::Create { 1226 name, 1227 new: garbage.oid(), 1228 }) 1229 .unwrap(); 1230 1231 assert!(matches!( 1232 CobStore::new(&repo).get::<Tags>(object), 1233 Err(CobError::UndecodableChange { .. }) 1234 )); 1235 } 1236 1237 #[test] 1238 fn merge_is_last_writer_wins_with_no_tombstone() { 1239 let key = signer(40); 1240 let nsid = Tag::type_name(); 1241 let resolve = |add_seconds: i64, remove_seconds: i64| -> bool { 1242 let (_dir, repo) = fixture(); 1243 let root = backend::write_change( 1244 &cob_home(), 1245 &repo, 1246 &nsid, 1247 &Tag::Add("seed".into()).encode().unwrap(), 1248 &[], 1249 None, 1250 &key, 1251 at(1), 1252 ) 1253 .unwrap(); 1254 let object = CobId::new(root.oid()); 1255 let add = backend::write_change( 1256 &cob_home(), 1257 &repo, 1258 &nsid, 1259 &Tag::Add("nel".into()).encode().unwrap(), 1260 &[root], 1261 Some(object), 1262 &key, 1263 at(add_seconds), 1264 ) 1265 .unwrap(); 1266 let remove = backend::write_change( 1267 &cob_home(), 1268 &repo, 1269 &nsid, 1270 &Tag::Remove("nel".into()).encode().unwrap(), 1271 &[root], 1272 Some(object), 1273 &key, 1274 at(remove_seconds), 1275 ) 1276 .unwrap(); 1277 let merge = backend::write_change( 1278 &cob_home(), 1279 &repo, 1280 &nsid, 1281 &Tag::Add("merged".into()).encode().unwrap(), 1282 &[add, remove], 1283 Some(object), 1284 &key, 1285 at(100), 1286 ) 1287 .unwrap(); 1288 let name = backend::cob_ref_name(&nsid, object).unwrap(); 1289 repo.update_ref(&RefUpdate::Create { 1290 name, 1291 new: merge.oid(), 1292 }) 1293 .unwrap(); 1294 CobStore::new(&repo) 1295 .get::<Tags>(object) 1296 .unwrap() 1297 .into_state() 1298 .contains("nel") 1299 }; 1300 1301 assert!(resolve(3, 2), "later add resurrects removed element"); 1302 assert!(!resolve(2, 3), "later remove wins under last-writer-wins"); 1303 } 1304 1305 #[test] 1306 fn tip_not_descending_from_root_is_detached() { 1307 let (_dir, repo) = fixture(); 1308 let key = signer(22); 1309 let nsid = Tag::type_name(); 1310 let root = backend::write_change( 1311 &cob_home(), 1312 &repo, 1313 &nsid, 1314 &Tag::Add("nel".into()).encode().unwrap(), 1315 &[], 1316 None, 1317 &key, 1318 at(1), 1319 ) 1320 .unwrap(); 1321 let object = CobId::new(root.oid()); 1322 let unrelated = backend::write_change( 1323 &cob_home(), 1324 &repo, 1325 &nsid, 1326 &Tag::Add("olaren".into()).encode().unwrap(), 1327 &[], 1328 None, 1329 &key, 1330 at(2), 1331 ) 1332 .unwrap(); 1333 let name = backend::cob_ref_name(&nsid, object).unwrap(); 1334 repo.update_ref(&RefUpdate::Create { 1335 name, 1336 new: unrelated.oid(), 1337 }) 1338 .unwrap(); 1339 1340 assert!(matches!( 1341 CobStore::new(&repo).get::<Tags>(object), 1342 Err(CobError::DetachedTip(_)) 1343 )); 1344 } 1345 1346 #[test] 1347 fn deep_history_loads_and_orders_without_overflow() { 1348 let count: i64 = 8_000; 1349 let (_dir, repo) = fixture(); 1350 let key = signer(23); 1351 let nsid = Tag::type_name(); 1352 let payload = Tag::Add("nel".into()).encode().unwrap(); 1353 let root = 1354 backend::write_change(&cob_home(), &repo, &nsid, &payload, &[], None, &key, at(0)) 1355 .unwrap(); 1356 let object = CobId::new(root.oid()); 1357 let tip = (1..count).fold(root, |parent, i| { 1358 backend::write_change( 1359 &cob_home(), 1360 &repo, 1361 &nsid, 1362 &payload, 1363 &[parent], 1364 Some(object), 1365 &key, 1366 at(i), 1367 ) 1368 .unwrap() 1369 }); 1370 publish(&repo, object, tip); 1371 let graph = CobStore::new(&repo).graph::<Tags>(object).unwrap(); 1372 assert_eq!(graph.len(), count as usize); 1373 assert_eq!(graph.causal_order().len(), count as usize); 1374 1375 let synthetic: u64 = 100_000; 1376 let oid = |index: u64| knot_types::Oid::from_hex(&format!("{index:040x}")).unwrap(); 1377 let actor = ActorId::from_secp256k1(&[0x02; 33]); 1378 let changes: std::collections::BTreeMap<ChangeId, Change> = (1..=synthetic) 1379 .map(|index| { 1380 let id = ChangeId::new(oid(index)); 1381 let parents = if index == 1 { 1382 Vec::new() 1383 } else { 1384 vec![ChangeId::new(oid(index - 1))] 1385 }; 1386 let change = Change { 1387 id, 1388 revision: oid(index), 1389 parents, 1390 type_name: Tag::type_name(), 1391 author: actor.clone(), 1392 signature: Signature::from_bytes(Vec::new()), 1393 payload: Payload::new(Vec::new()), 1394 timestamp: UnixSeconds::new(index as i64), 1395 }; 1396 (id, change) 1397 }) 1398 .collect(); 1399 let synthetic_order = ChangeGraph::new(CobId::new(oid(1)), changes).causal_order(); 1400 assert_eq!(synthetic_order.len(), synthetic as usize); 1401 assert_eq!(synthetic_order.first(), Some(&ChangeId::new(oid(1)))); 1402 assert_eq!(synthetic_order.last(), Some(&ChangeId::new(oid(synthetic)))); 1403 } 1404 1405 #[test] 1406 fn grafted_foreign_genesis_is_refused() { 1407 let (_dir, repo) = fixture(); 1408 let key = signer(31); 1409 let nsid = Tag::type_name(); 1410 let root = backend::write_change( 1411 &cob_home(), 1412 &repo, 1413 &nsid, 1414 &Tag::Add("nel".into()).encode().unwrap(), 1415 &[], 1416 None, 1417 &key, 1418 at(1), 1419 ) 1420 .unwrap(); 1421 let object = CobId::new(root.oid()); 1422 let child = backend::write_change( 1423 &cob_home(), 1424 &repo, 1425 &nsid, 1426 &Tag::Add("olaren".into()).encode().unwrap(), 1427 &[root], 1428 Some(object), 1429 &key, 1430 at(2), 1431 ) 1432 .unwrap(); 1433 let foreign = backend::write_change( 1434 &cob_home(), 1435 &repo, 1436 &nsid, 1437 &Tag::Add("evil".into()).encode().unwrap(), 1438 &[], 1439 None, 1440 &key, 1441 at(3), 1442 ) 1443 .unwrap(); 1444 let merge = backend::write_change( 1445 &cob_home(), 1446 &repo, 1447 &nsid, 1448 &Tag::Add("merge".into()).encode().unwrap(), 1449 &[child, foreign], 1450 Some(object), 1451 &key, 1452 at(4), 1453 ) 1454 .unwrap(); 1455 let name = backend::cob_ref_name(&nsid, object).unwrap(); 1456 repo.update_ref(&RefUpdate::Create { 1457 name, 1458 new: merge.oid(), 1459 }) 1460 .unwrap(); 1461 1462 assert!(matches!( 1463 CobStore::new(&repo).get::<Tags>(object), 1464 Err(CobError::MultipleRoots { .. }) 1465 )); 1466 } 1467 1468 #[test] 1469 fn cob_ref_name_and_parse_cob_ref_are_inverses() { 1470 let nsid = Tag::type_name(); 1471 let object = CobId::new(knot_types::Oid::from_hex(&"a".repeat(40)).unwrap()); 1472 let name = backend::cob_ref_name(&nsid, object).unwrap(); 1473 1474 assert_eq!(parse_cob_ref(name.as_str()), Some((nsid, object))); 1475 assert_eq!(parse_cob_ref("refs/heads/main"), None); 1476 assert_eq!( 1477 parse_cob_ref("refs/cobs/sh.tangled.test.tag/not-an-oid"), 1478 None 1479 ); 1480 } 1481 1482 fn loose_commits(repo: &Repo) -> BTreeSet<knot_types::Oid> { 1483 std::fs::read_dir(repo.path().join("objects")) 1484 .unwrap() 1485 .filter_map(Result::ok) 1486 .filter(|shard| { 1487 shard 1488 .file_name() 1489 .to_str() 1490 .map(|s| s.len() == 2) 1491 .unwrap_or(false) 1492 }) 1493 .flat_map(|shard| { 1494 let prefix = shard.file_name().to_str().unwrap().to_string(); 1495 std::fs::read_dir(shard.path()) 1496 .unwrap() 1497 .filter_map(Result::ok) 1498 .filter_map(move |entry| { 1499 let rest = entry.file_name().to_str()?.to_string(); 1500 knot_types::Oid::from_hex(&format!("{prefix}{rest}")).ok() 1501 }) 1502 .collect::<Vec<_>>() 1503 }) 1504 .filter(|oid| { 1505 repo.git() 1506 .find_object(oid.object_id()) 1507 .map(|object| object.kind == gix::object::Kind::Commit) 1508 .unwrap_or(false) 1509 }) 1510 .collect() 1511 } 1512 1513 #[test] 1514 fn a_contended_retry_keeps_both_writes_and_leaves_no_dangling_commit() { 1515 let (_dir, repo) = fixture(); 1516 let store = CobStore::new(&repo); 1517 let key = signer(124); 1518 let created = store 1519 .create(&cob_home(), &Tag::Add("base".into()), &key, at(1)) 1520 .unwrap(); 1521 1522 let injected = std::cell::Cell::new(false); 1523 let mine = store.update_with::<Tags, CobError>( 1524 &cob_home(), 1525 created.object, 1526 &key, 1527 at(3), 1528 |_state| { 1529 if !injected.replace(true) { 1530 store 1531 .update( 1532 &cob_home(), 1533 created.object, 1534 &Tag::Add("intruder".into()), 1535 &key, 1536 at(2), 1537 ) 1538 .unwrap(); 1539 } 1540 Ok(Tag::Add("mine".into())) 1541 }, 1542 ); 1543 assert!( 1544 mine.is_ok(), 1545 "handler retried instead of surfacing StaleTip" 1546 ); 1547 1548 let state = store.get::<Tags>(created.object).unwrap().into_state(); 1549 assert!( 1550 state.contains("intruder"), 1551 "concurrent append survives retry" 1552 ); 1553 assert!( 1554 state.contains("mine"), 1555 "retried append isn't lost to stale tip" 1556 ); 1557 1558 let reachable: BTreeSet<knot_types::Oid> = store 1559 .graph::<Tags>(created.object) 1560 .unwrap() 1561 .causal_order() 1562 .into_iter() 1563 .map(|change| change.oid()) 1564 .collect(); 1565 assert_eq!( 1566 reachable.len(), 1567 3, 1568 "live chain is base, intruder, then mine" 1569 ); 1570 assert_eq!( 1571 loose_commits(&repo), 1572 reachable, 1573 "abandoned compare-and-swap attempt left no dangling commit behind" 1574 ); 1575 } 1576 1577 #[test] 1578 fn exhausting_the_retry_budget_is_contended_and_leaves_no_dangling_commits() { 1579 let (_dir, repo) = fixture(); 1580 let store = CobStore::new(&repo); 1581 let key = signer(126); 1582 let created = store 1583 .create(&cob_home(), &Tag::Add("base".into()), &key, at(1)) 1584 .unwrap(); 1585 let calls = std::cell::Cell::new(0i64); 1586 1587 let result = store.update_with::<Tags, CobError>( 1588 &cob_home(), 1589 created.object, 1590 &key, 1591 at(1000), 1592 |_state| { 1593 let i = calls.get(); 1594 calls.set(i + 1); 1595 store 1596 .update( 1597 &cob_home(), 1598 created.object, 1599 &Tag::Add(format!("intruder{i}")), 1600 &key, 1601 at(100 + i), 1602 ) 1603 .unwrap(); 1604 Ok(Tag::Add("mine".into())) 1605 }, 1606 ); 1607 1608 assert!(matches!(result, Err(CobError::Contended(_)))); 1609 assert_eq!( 1610 calls.get(), 1611 MAX_CAS_RETRIES as i64, 1612 "decision ran exactly the retry budget before failing closed" 1613 ); 1614 1615 let reachable: BTreeSet<knot_types::Oid> = store 1616 .graph::<Tags>(created.object) 1617 .unwrap() 1618 .causal_order() 1619 .into_iter() 1620 .map(|change| change.oid()) 1621 .collect(); 1622 assert_eq!( 1623 reachable.len(), 1624 1 + MAX_CAS_RETRIES, 1625 "base plus one landed commit per injected intruder" 1626 ); 1627 assert_eq!( 1628 loose_commits(&repo), 1629 reachable, 1630 "every failed attempt across whole budget cleaned up its own orphan" 1631 ); 1632 } 1633 1634 #[test] 1635 fn an_identical_concurrent_change_is_not_cleaned_up_as_the_live_tip() { 1636 let (_dir, repo) = fixture(); 1637 let store = CobStore::new(&repo); 1638 let key = signer(125); 1639 let nsid = Tag::type_name(); 1640 let created = store 1641 .create(&cob_home(), &Tag::Add("base".into()), &key, at(1)) 1642 .unwrap(); 1643 let dup = Tag::Add("dup".into()); 1644 let bytes = dup.encode().unwrap(); 1645 1646 let first = backend::write_change( 1647 &cob_home(), 1648 &repo, 1649 &nsid, 1650 &bytes, 1651 &[created.tip], 1652 Some(created.object), 1653 &key, 1654 at(2), 1655 ) 1656 .unwrap(); 1657 let second = backend::write_change( 1658 &cob_home(), 1659 &repo, 1660 &nsid, 1661 &bytes, 1662 &[created.tip], 1663 Some(created.object), 1664 &key, 1665 at(2), 1666 ) 1667 .unwrap(); 1668 assert_eq!( 1669 first, second, 1670 "byte-identical changes are content-addressed to one commit" 1671 ); 1672 1673 let winner = store 1674 .append( 1675 &cob_home(), 1676 created.object, 1677 &nsid, 1678 created.tip, 1679 &dup, 1680 &key, 1681 at(2), 1682 ) 1683 .unwrap(); 1684 assert_eq!(winner, first); 1685 1686 let loser = store.append( 1687 &cob_home(), 1688 created.object, 1689 &nsid, 1690 created.tip, 1691 &dup, 1692 &key, 1693 at(2), 1694 ); 1695 assert!(matches!(loser, Err(CobError::StaleTip { .. }))); 1696 1697 assert_eq!( 1698 backend::resolve_tip(&repo, &nsid, created.object).unwrap(), 1699 Some(winner.oid()), 1700 "shared live tip survived identical-write loser's cleanup" 1701 ); 1702 assert!( 1703 store.get::<Tags>(created.object).is_ok(), 1704 "live tip wasn't deleted out from under the object" 1705 ); 1706 } 1707 1708 #[test] 1709 fn collect_refuses_a_graph_past_its_limit() { 1710 let (_dir, repo) = fixture(); 1711 let key = signer(34); 1712 let nsid = Tag::type_name(); 1713 let payload = Tag::Add("nel".into()).encode().unwrap(); 1714 let root = 1715 backend::write_change(&cob_home(), &repo, &nsid, &payload, &[], None, &key, at(0)) 1716 .unwrap(); 1717 let object = CobId::new(root.oid()); 1718 let tip = (1..4).fold(root, |parent, i| { 1719 backend::write_change( 1720 &cob_home(), 1721 &repo, 1722 &nsid, 1723 &payload, 1724 &[parent], 1725 Some(object), 1726 &key, 1727 at(i), 1728 ) 1729 .unwrap() 1730 }); 1731 1732 assert!(matches!( 1733 backend::collect(&repo, tip, object, 2, None), 1734 Err(CobError::HistoryTooLong(_)) 1735 )); 1736 assert_eq!( 1737 backend::collect(&repo, tip, object, 10, None) 1738 .unwrap() 1739 .len(), 1740 4 1741 ); 1742 } 1743 1744 #[test] 1745 fn changes_since_returns_a_suffix_enforces_shape_and_rejects_a_non_ancestor() { 1746 let (_dir, repo) = fixture(); 1747 let store = CobStore::new(&repo); 1748 let key = signer(60); 1749 1750 let created = store 1751 .create(&cob_home(), &Tag::Add("nel".into()), &key, at(1)) 1752 .unwrap(); 1753 let full = store.changes_since::<Tags>(created.object, None).unwrap(); 1754 assert_eq!(full.tip, created.tip); 1755 assert_eq!( 1756 full.changes.iter().map(|c| c.id).collect::<Vec<_>>(), 1757 vec![created.tip] 1758 ); 1759 1760 let second = store 1761 .update( 1762 &cob_home(), 1763 created.object, 1764 &Tag::Add("olaren".into()), 1765 &key, 1766 at(2), 1767 ) 1768 .unwrap(); 1769 let third = store 1770 .update( 1771 &cob_home(), 1772 created.object, 1773 &Tag::Add("teq".into()), 1774 &key, 1775 at(3), 1776 ) 1777 .unwrap(); 1778 1779 let delta = store 1780 .changes_since::<Tags>(created.object, Some(created.tip)) 1781 .unwrap(); 1782 assert_eq!(delta.tip, third); 1783 assert_eq!( 1784 delta.changes.iter().map(|c| c.id).collect::<Vec<_>>(), 1785 vec![second, third], 1786 "delta is the appended changes in causal order instead of whole graph" 1787 ); 1788 1789 let caught_up = store 1790 .changes_since::<Tags>(created.object, Some(third)) 1791 .unwrap(); 1792 assert!(caught_up.changes.is_empty()); 1793 assert_eq!(caught_up.tip, third); 1794 1795 let (_diverged_dir, diverged_repo) = fixture(); 1796 let diverged = linear_chain(&diverged_repo, &signer(71), &["nel", "olaren"]); 1797 let stranger = ChangeId::new(knot_types::Oid::from_hex(&"a".repeat(40)).unwrap()); 1798 assert!( 1799 matches!( 1800 CobStore::new(&diverged_repo).changes_since::<Tags>(diverged, Some(stranger)), 1801 Err(CobError::DivergedTip { .. }) 1802 ), 1803 "a since that doesn't descend from the indexed tip is refused, not re-folded onto stale state" 1804 ); 1805 1806 let (_forked_dir, forked_repo) = fixture(); 1807 let forked = forked_tag_object(&forked_repo, &signer(70), "nel", "olaren", "teq"); 1808 let forked_store = CobStore::new(&forked_repo); 1809 assert!( 1810 matches!( 1811 forked_store.get::<LinearTags>(forked), 1812 Err(CobError::ForkedHistory { .. }) 1813 ), 1814 "materialization fails closed on forked linear history" 1815 ); 1816 assert!( 1817 matches!( 1818 forked_store.changes_since::<LinearTags>(forked, None), 1819 Err(CobError::ForkedHistory { .. }) 1820 ), 1821 "changes_since refuses the same fork instead of folding it for the index" 1822 ); 1823 assert!( 1824 forked_store.changes_since::<Tags>(forked, None).is_ok(), 1825 "convergent class still folds the same graph" 1826 ); 1827 } 1828 1829 struct LinearTags; 1830 1831 impl Evaluate for LinearTags { 1832 type State = BTreeSet<String>; 1833 type Change = Tag; 1834 1835 const HISTORY: HistoryModel = HistoryModel::Linear; 1836 1837 fn initial() -> Self::State { 1838 BTreeSet::new() 1839 } 1840 1841 fn apply(state: Self::State, change: Self::Change, author: &ActorId) -> Self::State { 1842 Tags::apply(state, change, author) 1843 } 1844 } 1845 1846 #[test] 1847 fn verify_accepts_the_owner_and_rejects_a_foreign_signer_or_mismatched_type() { 1848 let (_dir, repo) = fixture(); 1849 let store = CobStore::new(&repo); 1850 let key = signer(51); 1851 let object = linear_chain(&repo, &key, &["nel", "olaren"]); 1852 let owner = ActorId::from_secp256k1(key.public_key().as_bytes()); 1853 assert!(store.verify::<Tags>(&cob_home(), object, &owner).is_ok()); 1854 1855 let stranger = ActorId::from_secp256k1(signer(52).public_key().as_bytes()); 1856 assert!(matches!( 1857 store.verify::<Tags>(&cob_home(), object, &stranger), 1858 Err(CobError::UnverifiedChange { .. }) 1859 )); 1860 1861 let (_other_dir, other_repo) = fixture(); 1862 let key = signer(53); 1863 let foreign = TypeName::new("sh.tangled.test.other").unwrap(); 1864 let root = tag_root(&other_repo, &key, "nel", 1); 1865 let object = CobId::new(root.oid()); 1866 let child = backend::write_change( 1867 &cob_home(), 1868 &other_repo, 1869 &foreign, 1870 &Tag::Add("olaren".into()).encode().unwrap(), 1871 &[root], 1872 Some(object), 1873 &key, 1874 at(2), 1875 ) 1876 .unwrap(); 1877 publish(&other_repo, object, child); 1878 let owner = ActorId::from_secp256k1(key.public_key().as_bytes()); 1879 assert!( 1880 matches!( 1881 CobStore::new(&other_repo).verify::<Tags>(&cob_home(), object, &owner), 1882 Err(CobError::UnexpectedChangeType { .. }) 1883 ), 1884 "owner-signed change whose type doesn't match namespace is refused at import" 1885 ); 1886 } 1887 1888 fn ops_and_perm() -> impl Strategy<Value = (Vec<(bool, u8)>, Vec<u32>)> { 1889 prop::collection::vec((any::<bool>(), 0u8..4u8), 1..8usize).prop_flat_map(|ops| { 1890 let len = ops.len(); 1891 (Just(ops), prop::collection::vec(any::<u32>(), len)) 1892 }) 1893 } 1894 1895 fn permutation(keys: &[u32]) -> Vec<usize> { 1896 let mut order: Vec<usize> = (0..keys.len()).collect(); 1897 order.sort_by_key(|&index| keys[index]); 1898 order 1899 } 1900 1901 fn model_state(ops: &[(bool, u8)]) -> BTreeSet<String> { 1902 ops.iter().fold( 1903 BTreeSet::from(["base".to_string()]), 1904 |mut state, (is_add, subject)| { 1905 let name = format!("s{subject}"); 1906 if *is_add { 1907 state.insert(name); 1908 } else { 1909 state.remove(&name); 1910 } 1911 state 1912 }, 1913 ) 1914 } 1915 1916 fn build_repo(ops: &[(bool, u8)], creation_order: &[usize]) -> (TempDir, Repo, CobId) { 1917 let (dir, repo) = fixture(); 1918 let key = signer(500); 1919 let nsid = Tag::type_name(); 1920 let root = tag_root(&repo, &key, "base", 1); 1921 let object = CobId::new(root.oid()); 1922 let child = |index: usize| { 1923 let (is_add, subject) = ops[index]; 1924 let name = format!("s{subject}"); 1925 let payload = if is_add { 1926 Tag::Add(name) 1927 } else { 1928 Tag::Remove(name) 1929 }; 1930 backend::write_change( 1931 &cob_home(), 1932 &repo, 1933 &nsid, 1934 &payload.encode().unwrap(), 1935 &[root], 1936 Some(object), 1937 &key, 1938 at(index as i64 + 2), 1939 ) 1940 .unwrap() 1941 }; 1942 creation_order.iter().for_each(|&index| { 1943 child(index); 1944 }); 1945 let parents: Vec<ChangeId> = (0..ops.len()).map(child).collect(); 1946 let merge = backend::write_change( 1947 &cob_home(), 1948 &repo, 1949 &nsid, 1950 &Tag::Remove("absent".into()).encode().unwrap(), 1951 &parents, 1952 Some(object), 1953 &key, 1954 at(ops.len() as i64 + 2), 1955 ) 1956 .unwrap(); 1957 publish(&repo, object, merge); 1958 (dir, repo, object) 1959 } 1960 1961 proptest! { 1962 #![proptest_config(ProptestConfig { cases: 32, ..ProptestConfig::default() })] 1963 1964 #[test] 1965 fn prop_concurrent_changes_converge_to_the_causal_fold((ops, keys) in ops_and_perm()) { 1966 let identity: Vec<usize> = (0..ops.len()).collect(); 1967 let model = model_state(&ops); 1968 1969 let (_canon_dir, canon_repo, canon_object) = build_repo(&ops, &identity); 1970 let canonical = CobStore::new(&canon_repo) 1971 .get::<Tags>(canon_object) 1972 .unwrap() 1973 .into_state(); 1974 1975 let (_perm_dir, perm_repo, perm_object) = build_repo(&ops, &permutation(&keys)); 1976 let permuted = CobStore::new(&perm_repo) 1977 .get::<Tags>(perm_object) 1978 .unwrap() 1979 .into_state(); 1980 1981 prop_assert_eq!(&canonical, &model); 1982 prop_assert_eq!(&permuted, &model); 1983 } 1984 1985 #[test] 1986 fn prop_rematerialization_is_idempotent((ops, keys) in ops_and_perm()) { 1987 let (_dir, repo, object) = build_repo(&ops, &permutation(&keys)); 1988 let store = CobStore::new(&repo); 1989 let first = store.get::<Tags>(object).unwrap(); 1990 let second = store.get::<Tags>(object).unwrap(); 1991 prop_assert_eq!(first.state(), second.state()); 1992 prop_assert_eq!(first.history().len(), second.history().len()); 1993 prop_assert_eq!( 1994 store.graph::<Tags>(object).unwrap().causal_order(), 1995 store.graph::<Tags>(object).unwrap().causal_order() 1996 ); 1997 } 1998 } 1999}