Oregami
Repositories/oxedyne/fe2o3

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
19use 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
29use oxedyne_fe2o3_ore::{
30 envelope::Envelope,
31 id::OpId,
32 op::{
33 Op,
34 Record,
35 },
36};
37
38use oxedyne_fe2o3_core::prelude::*;
39use oxedyne_fe2o3_jdat::{
40 prelude::*,
41 id::NumIdDat,
42};
43use oxedyne_fe2o3_iop_db::api::Database;
44use oxedyne_fe2o3_o3db_sync::prelude::{
45 Encrypter,
46 Hasher,
47};
48
49use 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.
58pub 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.
75fn 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.
83fn 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.
100fn 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.
145fn 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.
158pub struct O3dbHub<'a, const UIDL: usize, UID, ENC, KH, D>
159where
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
170impl<'a, const UIDL: usize, UID, ENC, KH, D> O3dbHub<'a, UIDL, UID, ENC, KH, D>
171where
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
212impl<'a, const UIDL: usize, UID, ENC, KH, D> Hub for O3dbHub<'a, UIDL, UID, ENC, KH, D>
213where
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}