oxedyne/fe2o3/fe2o3_pearlite/src/collab/hub.rs
11.2 KiB, 21 runs
created by r1870400018:58598, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! The store a document's operations live in, behind a thin interface. |
| 2 | //! |
| 3 | //! A document's edit stream is a set of signed envelopes keyed by the document's [`DocId`]. The |
| 4 | //! [`Hub`] trait is the whole of what the fold and sign layers ask of a store: append one operation, |
| 5 | //! and load the document's whole log back. Keeping it this small is deliberate -- the o3db backing |
| 6 | //! here can be swapped for a networked one, or a browser's local one, without the layers above |
| 7 | //! noticing. |
| 8 | //! |
| 9 | //! # The o3db representation |
| 10 | //! |
| 11 | //! o3db is a keyed value store, so a document's log is held under one key -- `pearlite-collab:doc:<id>` |
| 12 | //! -- as a list of `[op_id, envelope]` pairs. An append reads that list, adds the pair if the |
| 13 | //! operation is not already there, and writes it back; a scan reads and decodes it. Reading the list |
| 14 | //! whole to append one entry is the simple, correct thing for a single writer, which is what this |
| 15 | //! increment has: the concurrent-append path, where two writers must not lose each other's operation, |
| 16 | //! belongs with the sync transport that carries operations between them, and is that increment's to |
| 17 | //! build. The interface does not change when it arrives -- only this one implementation behind it. |
| 18 | |
| 19 | use crate::collab::{ |
| 20 | op, |
| 21 | sign, |
| 22 | DocId, |
| 23 | MAX_DOC_BYTES, |
| 24 | MAX_OP_BYTES, |
| 25 | MAX_OPS_PER_DOC, |
| 26 | MAX_TIME_SKEW_SECS, |
| 27 | }; |
| 28 | |
| 29 | use oxedyne_fe2o3_ore::{ |
| 30 | envelope::Envelope, |
| 31 | id::OpId, |
| 32 | op::{ |
| 33 | Op, |
| 34 | Record, |
| 35 | }, |
| 36 | }; |
| 37 | |
| 38 | use oxedyne_fe2o3_core::prelude::*; |
| 39 | use oxedyne_fe2o3_jdat::{ |
| 40 | prelude::*, |
| 41 | id::NumIdDat, |
| 42 | }; |
| 43 | use oxedyne_fe2o3_iop_db::api::Database; |
| 44 | use oxedyne_fe2o3_o3db_sync::prelude::{ |
| 45 | Encrypter, |
| 46 | Hasher, |
| 47 | }; |
| 48 | |
| 49 | use std::time::{ |
| 50 | SystemTime, |
| 51 | UNIX_EPOCH, |
| 52 | }; |
| 53 | |
| 54 | /// The thin store the collaboration layer reaches a document's log through. |
| 55 | /// |
| 56 | /// It is deliberately just two operations. A store that can do more -- a networked hub, a local cache |
| 57 | /// -- offers it through this same pair, so nothing above has to know which it holds. |
| 58 | pub trait Hub { |
| 59 | /// Verifies a signed operation and, if it holds, appends it to a document's log, returning the |
| 60 | /// identifier it was stored under -- which is taken from the verified record, never from the |
| 61 | /// caller, so a peer cannot store an envelope under an identity it did not sign for. Appending an |
| 62 | /// operation the log already holds is a no-op that still returns its identifier, so a redelivery |
| 63 | /// does not duplicate it. An envelope that does not verify, is oversized, carries an operation of |
| 64 | /// the wrong voice, an annotation stamped for another document, a time far in the future, or that |
| 65 | /// would take the log past its size bound, is refused rather than stored. |
| 66 | fn put(&self, doc: &DocId, env: &Envelope) -> Outcome<OpId>; |
| 67 | |
| 68 | /// Loads a document's whole log: every operation stored under its identity, each with the envelope |
| 69 | /// that attests to it. The order is not significant -- the fold imposes its own -- and a document |
| 70 | /// with no operations yet is an empty log, not an error. |
| 71 | fn scan(&self, doc: &DocId) -> Outcome<Vec<(OpId, Envelope)>>; |
| 72 | } |
| 73 | |
| 74 | /// The receiving replica's wall clock in unix epoch seconds, for bounding an author-stated time. |
| 75 | fn now_secs() -> u64 { |
| 76 | match SystemTime::now().duration_since(UNIX_EPOCH) { |
| 77 | Ok(d) => d.as_secs(), |
| 78 | Err(_) => 0, // a clock before the epoch bounds nothing, which is safe: it only tightens. |
| 79 | } |
| 80 | } |
| 81 | |
| 82 | /// The author-stated time an operation carries, if it carries one, for the skew bound. |
| 83 | fn op_time(op: &Op) -> Option<u64> { |
| 84 | match op { |
| 85 | Op::Proposal { time, .. } => Some(*time), |
| 86 | Op::Said { time, .. } => Some(*time), |
| 87 | Op::Amended { time, .. } => Some(*time), |
| 88 | Op::Settled { time, .. } => Some(*time), |
| 89 | _ => None, |
| 90 | } |
| 91 | } |
| 92 | |
| 93 | /// Checks a verified record is one this log may hold: a Pearlite annotation operation, of the right |
| 94 | /// voice, whose body -- where it carries one -- decodes as an annotation stamped for this document, and |
| 95 | /// whose time is not implausibly far ahead of the receiving clock. |
| 96 | /// |
| 97 | /// This is ingest validation: the fold is poison-proof against a bad operation that is already stored, |
| 98 | /// and this is what keeps one from being stored in the first place. The two are deliberately |
| 99 | /// belt-and-braces, since an append-only log cannot take back what it once accepted. |
| 100 | fn validate_ingest(doc: &DocId, rec: &Record) -> Outcome<()> { |
| 101 | // The time bound applies to every operation that states one. |
| 102 | if let Some(time) = op_time(&rec.op) { |
| 103 | let ceiling = now_secs().saturating_add(MAX_TIME_SKEW_SECS); |
| 104 | if time > ceiling { |
| 105 | return Err(err!( |
| 106 | "The operation {} states a time of {}, more than {} seconds ahead of the receiving \ |
| 107 | clock; a time far in the future is a bid to win last-writer-wins ordering and is \ |
| 108 | refused.", rec.id(), time, MAX_TIME_SKEW_SECS; |
| 109 | Invalid, Input, Security, Excessive)); |
| 110 | } |
| 111 | } |
| 112 | match &rec.op { |
| 113 | Op::Proposal { voice, body, .. } | Op::Amended { voice, body, .. } => { |
| 114 | res!(check_voice(rec, voice)); |
| 115 | // The body must decode as an annotation, bounded, and be stamped for this document, or the |
| 116 | // operation does not belong in this log. |
| 117 | let ann = res!(op::annotation_from_body(body)); |
| 118 | match &ann.doc_id { |
| 119 | Some(id) if id == doc.as_str() => {}, |
| 120 | Some(id) => return Err(err!( |
| 121 | "The operation {} carries an annotation stamped for document {:?}, not {:?}; it \ |
| 122 | will not be stored here.", rec.id(), id, doc.as_str(); |
| 123 | Invalid, Input, Security, Mismatch)), |
| 124 | None => return Err(err!( |
| 125 | "The operation {} carries an annotation with no document identity, so it cannot be \ |
| 126 | confirmed to belong to {:?}.", rec.id(), doc.as_str(); |
| 127 | Invalid, Input, Missing)), |
| 128 | } |
| 129 | Ok(()) |
| 130 | }, |
| 131 | Op::Said { voice, .. } => check_voice(rec, voice), |
| 132 | // A settlement carries no voice or body of its own; its authority is checked at the fold, |
| 133 | // against the proposal it settles. Every other operation kind has no business in an annotation |
| 134 | // log at all. |
| 135 | Op::Settled { .. } => Ok(()), |
| 136 | other => Err(err!( |
| 137 | "The operation {} is a {}, which is not an annotation operation and does not belong in a \ |
| 138 | Pearlite collaboration log.", rec.id(), other.name(); |
| 139 | Invalid, Input, Mismatch)), |
| 140 | } |
| 141 | } |
| 142 | |
| 143 | /// Refuses an operation written under any voice but Pearlite's, so the log holds annotation threads and |
| 144 | /// not some other forge's proposals that happen to share the vocabulary. |
| 145 | fn check_voice(rec: &Record, voice: &str) -> Outcome<()> { |
| 146 | if voice != op::VOICE { |
| 147 | return Err(err!( |
| 148 | "The operation {} is written under the voice {:?}, not {:?}; only Pearlite annotation \ |
| 149 | operations belong in this log.", rec.id(), voice, op::VOICE; |
| 150 | Invalid, Input, Mismatch)); |
| 151 | } |
| 152 | Ok(()) |
| 153 | } |
| 154 | |
| 155 | /// A [`Hub`] backed by an o3db instance, reached through the blocking |
| 156 | /// [`Database`](oxedyne_fe2o3_iop_db::api::Database) trait so the store's own generic machinery stays |
| 157 | /// out of here. |
| 158 | pub struct O3dbHub<'a, const UIDL: usize, UID, ENC, KH, D> |
| 159 | where |
| 160 | UID: NumIdDat<UIDL> + Copy, |
| 161 | ENC: Encrypter, |
| 162 | KH: Hasher, |
| 163 | D: Database<UIDL, UID, ENC, KH>, |
| 164 | { |
| 165 | db: &'a D, |
| 166 | user: UID, |
| 167 | _p: std::marker::PhantomData<(ENC, KH)>, |
| 168 | } |
| 169 | |
| 170 | impl<'a, const UIDL: usize, UID, ENC, KH, D> O3dbHub<'a, UIDL, UID, ENC, KH, D> |
| 171 | where |
| 172 | UID: NumIdDat<UIDL> + Copy, |
| 173 | ENC: Encrypter, |
| 174 | KH: Hasher, |
| 175 | D: Database<UIDL, UID, ENC, KH>, |
| 176 | { |
| 177 | pub fn new(db: &'a D, user: UID) -> Self { |
| 178 | Self { db, user, _p: std::marker::PhantomData } |
| 179 | } |
| 180 | |
| 181 | /// The key a document's log is stored under. |
| 182 | fn key(doc: &DocId) -> Dat { |
| 183 | dat!(fmt!("pearlite-collab:doc:{}", doc.as_str())) |
| 184 | } |
| 185 | |
| 186 | /// Reads and decodes a document's log, or an empty one where nothing is stored yet. |
| 187 | fn load(&self, doc: &DocId) -> Outcome<Vec<(OpId, Envelope)>> { |
| 188 | let stored = res!(self.db.get(&Self::key(doc), None)); |
| 189 | let entries = match stored { |
| 190 | None => return Ok(Vec::new()), |
| 191 | Some((Dat::List(entries), _)) => entries, |
| 192 | Some((other, _)) => return Err(err!( |
| 193 | "A Pearlite hub log for {} is stored as a list of entries, found {:?}.", |
| 194 | doc, other; Input, Invalid, Mismatch)), |
| 195 | }; |
| 196 | let mut out = Vec::with_capacity(entries.len()); |
| 197 | for entry in entries { |
| 198 | let pair = try_extract_dat!(entry, List); |
| 199 | if pair.len() != 2 { |
| 200 | return Err(err!( |
| 201 | "A Pearlite hub log entry is an [op_id, envelope] pair, found {} elements.", |
| 202 | pair.len(); Input, Invalid, Mismatch)); |
| 203 | } |
| 204 | let op_id = res!(OpId::from_dat(&pair[0])); |
| 205 | let env = res!(Envelope::from_dat(&pair[1])); |
| 206 | out.push((op_id, env)); |
| 207 | } |
| 208 | Ok(out) |
| 209 | } |
| 210 | } |
| 211 | |
| 212 | impl<'a, const UIDL: usize, UID, ENC, KH, D> Hub for O3dbHub<'a, UIDL, UID, ENC, KH, D> |
| 213 | where |
| 214 | UID: NumIdDat<UIDL> + Copy, |
| 215 | ENC: Encrypter, |
| 216 | KH: Hasher, |
| 217 | D: Database<UIDL, UID, ENC, KH>, |
| 218 | { |
| 219 | fn put(&self, doc: &DocId, env: &Envelope) -> Outcome<OpId> { |
| 220 | // A peer does not get to write an unbounded blob into a log: the size of the sealed envelope is |
| 221 | // bounded before anything is decoded from it. |
| 222 | if env.payload().len() > MAX_OP_BYTES { |
| 223 | return Err(err!( |
| 224 | "A Pearlite operation of {} payload bytes exceeds the {}-byte maximum a single \ |
| 225 | operation may occupy.", env.payload().len(), MAX_OP_BYTES; |
| 226 | Invalid, Input, Excessive, Size)); |
| 227 | } |
| 228 | // Verify the signature, decode the record under bounds, and check the header's replica is the |
| 229 | // one the signer's key derives under this document -- so the identity the operation is stored |
| 230 | // under is one the signer actually signed for in this document, taken from the record and never |
| 231 | // from the caller, and an operation minted for another document is refused here. |
| 232 | let opened = res!(sign::open(env, doc)); |
| 233 | let id = opened.record().id(); |
| 234 | res!(validate_ingest(doc, opened.record())); |
| 235 | |
| 236 | let mut entries = res!(self.load(doc)); |
| 237 | if let Some((_, stored_env)) = entries.iter().find(|(stored, _)| *stored == id) { |
| 238 | // A redelivery of the very same envelope is a no-op; the same identity carrying a different |
| 239 | // envelope is a squat or a two-device-same-key race, and is surfaced rather than silently |
| 240 | // dropped -- the first-seen operation would otherwise vanish the second, or the reverse. |
| 241 | if stored_env == env { |
| 242 | return Ok(id); |
| 243 | } |
| 244 | return Err(err!( |
| 245 | "The document {} already holds an operation identified {} with a different envelope; a \ |
| 246 | second, differing operation under the same identity is a squat or a same-key race, and \ |
| 247 | is refused rather than either being dropped silently.", doc, id; |
| 248 | Invalid, Input, Security, Conflict)); |
| 249 | } |
| 250 | // The count is bounded so a peer cannot exhaust the store by appending without end. |
| 251 | if entries.len() >= MAX_OPS_PER_DOC { |
| 252 | return Err(err!( |
| 253 | "The document {} already holds {} operations, the maximum a single log may hold; a \ |
| 254 | further operation is refused.", doc, entries.len(); |
| 255 | Invalid, Input, Excessive, Size)); |
| 256 | } |
| 257 | // The total size is bounded too, so a log's real cost is capped and not only its length. |
| 258 | let total: usize = entries.iter().map(|(_, e)| e.payload().len()).sum(); |
| 259 | if total.saturating_add(env.payload().len()) > MAX_DOC_BYTES { |
| 260 | return Err(err!( |
| 261 | "The document {} holds {} operation payload bytes; a further {} would pass the {}-byte \ |
| 262 | maximum a single log may reach, and is refused.", |
| 263 | doc, total, env.payload().len(), MAX_DOC_BYTES; |
| 264 | Invalid, Input, Excessive, Size)); |
| 265 | } |
| 266 | entries.push((id, env.clone())); |
| 267 | let list = Dat::List(entries.iter() |
| 268 | .map(|(stored, e)| listdat![stored.to_dat(), e.to_dat()]) |
| 269 | .collect()); |
| 270 | res!(self.db.insert(Self::key(doc), list, self.user, None)); |
| 271 | Ok(id) |
| 272 | } |
| 273 | |
| 274 | fn scan(&self, doc: &DocId) -> Outcome<Vec<(OpId, Envelope)>> { |
| 275 | self.load(doc) |
| 276 | } |
| 277 | } |