//! Record writes, as an observer sees them. //! //! [`LifecycleEvent`](crate::LifecycleEvent) says what happened to an //! *account*. This says what happened to a *repository*: one or more records //! changed, here is which ones and what they now say. The two are separate //! traits rather than one enum because they are separate streams with separate //! costs — a lifecycle event is four short strings and happens a few times per //! account, and a commit carries whole records and happens every time an agent //! speaks. //! //! A [`CommitEvent`] names the commit its operations landed in, which is what //! makes it reconcilable: a consumer holding every event up to some point //! holds exactly the records that commit covers, and can ask this server, or //! read an export, and compare. Since the server keeps its commits, the //! comparison runs the other way too: the revision an event carries is one a //! consumer can hand back as `com.atproto.sync.getRepo`'s `since`. //! //! # One event per commit, not one per record //! //! A single write and a batch of several both make exactly one commit — see //! [`crate::provision`] — and this is the type that says so on the stream: one //! [`CommitEvent`] per commit, carrying every [`CommitOp`] the commit made. //! `applyWrites` is why the list exists; an ordinary single write is the //! one-element case, made with [`CommitEvent::one`]. //! //! What is still absent is the commit's *blocks* — the //! signed commit itself and the tree nodes proving a record against it. Those //! are a request away at `com.atproto.sync.getRecord` and are deliberately not //! inlined here. use std::sync::Mutex; use serde_json::Value; /// Where a single-record write left the repository, before it is folded into /// a [`CommitOp`]. /// /// Three content identifiers' worth of fact: what the record is called, what /// the commit that now covers it is called, and which revision that commit is. /// A consumer that has these can go and check them — against /// `com.atproto.sync.getLatestCommit`, against an export, or against the proof /// at `com.atproto.sync.getRecord` — which is the whole difference between a /// stream that can be reconciled and one that can only be believed. #[derive(Debug, Clone, PartialEq, Eq)] pub struct Position { /// The CID of the record, over its canonical DAG-CBOR encoding. pub cid: String, /// The CID of the signed commit that covers it. pub commit: String, /// The revision that commit is at. /// /// Not a sequence number and not comparable across repositories, but it /// does move for every commit to one: a revision is minted per commit /// rather than read off the record keys, so a write to a fixed key such /// as `self` moves it like any other. See [`crate::history`]. pub rev: String, } /// What a commit did to one record. /// /// A deletion carries no bytes because there is nothing left to hand over: the /// record it names is gone, and [`CommitOp::deleted`] is the only way to say /// that. It does still name what left, which is what a consumer needs to tell /// a removal it already applied from one it disagrees about. #[derive(Debug, Clone, PartialEq)] pub enum RecordChange { /// The record now holds this value, named by this CID. Written { /// The CID naming the record's contents. cid: String, /// The record itself, exactly as it was stored. record: Value, }, /// The record is gone, and this is the CID it was named by. /// /// `None` only when the record could not be named, which a stored record /// always can be — see `provision::previous_cid`. Carried because a /// consumer of this stream can compare it against the record it holds and /// tell "the record I have is the one that left" from "we disagree about /// what was there", which is the difference between reconciling and /// guessing. Deleted { /// The CID the record was named by. prev: Option, }, } /// One record, and what a commit did to it. #[derive(Debug, Clone, PartialEq)] pub struct CommitOp { /// The collection the record is filed under. pub collection: String, /// The record key. pub rkey: String, /// What happened to it. pub change: RecordChange, } impl CommitOp { /// A record that was written, whether created or replaced. pub fn written( collection: impl Into, rkey: impl Into, cid: impl Into, record: Value, ) -> Self { Self { collection: collection.into(), rkey: rkey.into(), change: RecordChange::Written { cid: cid.into(), record, }, } } /// A record that was deleted, and the CID it held. pub fn deleted( collection: impl Into, rkey: impl Into, prev: Option, ) -> Self { Self { collection: collection.into(), rkey: rkey.into(), change: RecordChange::Deleted { prev }, } } /// The CID the key held before this operation, or `None` for a create. /// /// A write does not carry it — this stream's writes say what a record now /// is, not what it was — so this is `Some` only on a deletion. pub fn prev(&self) -> Option<&str> { match &self.change { RecordChange::Written { .. } => None, RecordChange::Deleted { prev } => prev.as_deref(), } } /// The AT-URI naming this record, under `did`. pub fn uri(&self, did: &str) -> String { format!("at://{did}/{}/{}", self.collection, self.rkey) } /// The CID a write named, or `None` for a deletion. pub fn cid(&self) -> Option<&str> { match &self.change { RecordChange::Written { cid, .. } => Some(cid), RecordChange::Deleted { .. } => None, } } } /// One or more records changed, in one commit. /// /// Owned values, for the same reason [`LifecycleEvent`](crate::LifecycleEvent) /// uses them: a sink may queue this, put it on a socket, or hold it long after /// the write that produced it returned. #[derive(Debug, Clone, PartialEq)] pub struct CommitEvent { /// The repository the commit belongs to, as a DID. pub did: String, /// The CID of the signed commit that covers every operation below. pub commit: String, /// The revision that commit is at. pub rev: String, /// Every record the commit touched, in the order the writes were asked /// for. pub ops: Vec, } impl CommitEvent { /// Assembles a commit event covering several operations. pub fn new( did: impl Into, commit: impl Into, rev: impl Into, ops: Vec, ) -> Self { Self { did: did.into(), commit: commit.into(), rev: rev.into(), ops, } } /// A commit covering exactly one written record. /// /// What an ordinary, unbatched write makes: one operation, in a commit /// that exists only for it. pub fn one( did: impl Into, collection: impl Into, rkey: impl Into, record: Value, at: Position, ) -> Self { Self { did: did.into(), commit: at.commit, rev: at.rev, ops: vec![CommitOp::written(collection, rkey, at.cid, record)], } } } /// A destination for commit events. /// /// Emitted synchronously inside the write that produced it and, like /// [`LifecycleSink`](crate::LifecycleSink), with no error channel: the record /// is already stored by the time this is called, and letting a subscriber turn /// a completed write into a failed one would be the tail wagging the dog. pub trait CommitSink: Send + Sync { /// Delivers one commit, with every operation it made. fn emit(&self, event: &CommitEvent); } /// A sink that keeps every commit it was handed, in order. /// /// Shipped rather than left in a test module for the same reason /// [`RecordingSink`](crate::RecordingSink) is: ordering is part of the /// contract, and anything embedding a [`Registry`](crate::Registry) needs the /// same tool to assert it. #[derive(Debug, Default)] pub struct RecordingCommitSink { commits: Mutex>, } impl RecordingCommitSink { /// A sink with nothing recorded. pub fn new() -> Self { Self::default() } /// Every commit so far, oldest first. pub fn commits(&self) -> Vec { self.commits .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) .clone() } } impl CommitSink for RecordingCommitSink { fn emit(&self, event: &CommitEvent) { self.commits .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()) .push(event.clone()); } } #[cfg(test)] mod tests { use super::*; /// A position with the right shape. These tests are about the event's /// fields and not about what the CIDs in them name. fn somewhere() -> Position { Position { cid: "bafyreie5737gdxlw5i64vzichcalba3z2v5n6icifvx5xytvske7mr3hpm".to_owned(), commit: "bafyreiaqmnopqrstuvwxyz234567abcdefghijklmnopqrstuvwxyz2345".to_owned(), rev: "3lqmnopqrstuv".to_owned(), } } #[test] fn a_commit_names_its_record_as_an_at_uri() { let commit = CommitEvent::one( "did:web:kestrel.agents.localhost%3A3000", "com.example.thing", "3lqmnop", serde_json::json!({"text": "reading the cart"}), somewhere(), ); assert_eq!( commit.ops[0].uri(&commit.did), "at://did:web:kestrel.agents.localhost%3A3000/com.example.thing/3lqmnop" ); } #[test] fn a_batch_commit_carries_every_operation_it_made() { let commit = CommitEvent::new( "did:web:a", "bafyreiacommit", "3lqmnopqrstuv", vec![ CommitOp::written("c", "1", "bafyreiarecord1", Value::Null), CommitOp::deleted("c", "2", None), CommitOp::written("c", "3", "bafyreiarecord3", Value::Null), ], ); assert_eq!(commit.ops.len(), 3); assert_eq!(commit.ops[0].cid(), Some("bafyreiarecord1")); assert_eq!(commit.ops[1].cid(), None); assert_eq!(commit.ops[1].change, RecordChange::Deleted { prev: None }); } #[test] fn the_recording_sink_keeps_order() { let sink = RecordingCommitSink::new(); sink.emit(&CommitEvent::one( "did:web:a", "c", "1", Value::Null, somewhere(), )); sink.emit(&CommitEvent::one( "did:web:a", "c", "2", Value::Null, somewhere(), )); let commits = sink.commits(); assert_eq!(commits.len(), 2); assert_eq!(commits[0].ops[0].rkey, "1"); assert_eq!(commits[1].ops[0].rkey, "2"); } }