This repository has no description
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}