heddle_object_model/object/
thread_replication.rs1pub mod capture_visibility;
5pub mod delegated_import;
6pub mod git_import_converter;
7pub mod git_import_graph;
8pub mod initial_base;
9pub mod integration;
10pub mod local_integration;
11pub mod metadata;
12pub mod ownership_claim;
13pub mod ownership_resolution;
14pub mod source_author;
15use std::collections::BTreeSet;
16
17pub use capture_visibility::CaptureVisibility;
18use serde::{Deserialize, Serialize};
19pub use source_author::{AuthoredCapture, SOURCE_AUTHORIZATION_METHOD, SourceAuthor};
20
21use crate::{
22 error::{HeddleError, Result},
23 object::{CollaborationOperationEnvelope, ContentHash, State, StateId},
24};
25
26pub const GENESIS_FORMAT: &str = "heddle-thread-genesis-v1";
27pub const OPERATION_FORMAT: &str = "heddle-thread-operation-v1";
28pub const MAX_OPERATION_BYTES: usize = 256 * 1024;
29
30#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
33#[serde(rename_all = "snake_case")]
34pub enum GenesisOwner {
35 LocalKey([u8; 32]),
36 Account(uuid::Uuid),
37}
38
39impl GenesisOwner {
40 fn is_valid(&self) -> bool {
41 match self {
42 Self::LocalKey(key) => *key != [0; 32],
43 Self::Account(account) => !account.is_nil(),
44 }
45 }
46}
47
48#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
49#[serde(deny_unknown_fields)]
50pub struct ThreadGenesis {
51 pub version: u16,
52 pub spool: String,
55 pub parent: Option<ContentHash>,
56 pub base: StateId,
57 pub name: String,
58 pub intent: String,
59 pub owner: GenesisOwner,
61 pub creator: [u8; 32],
62 pub nonce: Vec<u8>,
63}
64
65impl ThreadGenesis {
66 pub fn encode(&self) -> Result<Vec<u8>> {
67 if self.version != 1
68 || !self.owner.is_valid()
69 || matches!(&self.owner, GenesisOwner::LocalKey(key) if *key != self.creator)
70 || !uuid::Uuid::parse_str(&self.spool)
71 .is_ok_and(|id| !id.is_nil() && id.to_string() == self.spool)
72 || self.name.is_empty()
73 || self.nonce.len() > 64
74 {
75 return Err(invalid("invalid Thread genesis"));
76 }
77 let bytes = rmp_serde::to_vec_named(self)?;
78 bounded(&bytes)?;
79 Ok(bytes)
80 }
81
82 pub fn id(&self) -> Result<ContentHash> {
83 Ok(ContentHash::compute_typed(GENESIS_FORMAT, &self.encode()?))
84 }
85
86 pub fn decode(bytes: &[u8]) -> Result<Self> {
87 bounded(bytes)?;
88 let genesis: Self = rmp_serde::from_slice(bytes)?;
89 if genesis.encode()? != bytes {
90 return Err(invalid("non-canonical Thread genesis"));
91 }
92 Ok(genesis)
93 }
94}
95
96#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize, Deserialize)]
97#[serde(rename_all = "snake_case")]
98pub enum ThreadFacet {
99 Source,
100 Discussion,
101 Metadata,
102}
103impl ThreadFacet {
104 pub const ALL: [Self; 3] = [Self::Source, Self::Discussion, Self::Metadata];
105}
106
107#[derive(Clone, Debug, PartialEq, Eq)]
109pub enum Admission {
110 Accepted,
111 Pending,
112 Rejected(String),
113}
114
115#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
116#[serde(tag = "kind", content = "canonical", rename_all = "snake_case")]
117#[allow(clippy::large_enum_variant)] pub enum ThreadOperationBody {
119 Capture(AuthoredCapture),
120 Integration(Vec<u8>),
121 LocalIntegration(Vec<u8>),
122 Discussion(Vec<u8>),
123 Context(Vec<u8>),
124 Metadata(Vec<u8>),
125}
126
127#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
130#[serde(deny_unknown_fields)]
131pub struct Capture {
132 pub state: Vec<u8>,
133 pub source_targets: Option<ContentHash>,
134 pub visibility: Option<CaptureVisibility>,
136}
137impl From<Vec<u8>> for Capture {
138 fn from(state: Vec<u8>) -> Self {
139 Self {
140 state,
141 source_targets: None,
142 visibility: None,
143 }
144 }
145}
146
147#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
148#[serde(deny_unknown_fields)]
149pub struct ThreadOperation {
150 pub version: u16,
151 pub thread: ContentHash,
152 pub parents: BTreeSet<ContentHash>,
153 pub publisher: [u8; 32],
156 pub body: ThreadOperationBody,
157}
158
159impl ThreadOperation {
160 pub fn reference_proof(
161 &self,
162 genesis: &ThreadGenesis,
163 ) -> Result<Option<crate::object::source_target::capture::ReferenceProof>> {
164 let Some(capture) = self.source_result()? else {
165 return Ok(None);
166 };
167 let Some(descriptor) = capture.source_targets else {
168 return Ok(None);
169 };
170 if self.thread != genesis.id()? {
171 return Err(invalid("capture reference Thread mismatch"));
172 }
173 Ok(Some(
174 crate::object::source_target::capture::ReferenceProof {
175 descriptor,
176 scope: crate::object::CollaborationScope {
177 spool: genesis.spool.parse().map_err(invalid)?,
178 thread: Some(self.thread),
179 },
180 state: State::decode_current_msgpack(&capture.state)?.id(),
181 },
182 ))
183 }
184 pub fn facet(&self) -> ThreadFacet {
185 match self.body {
186 ThreadOperationBody::Capture(_)
187 | ThreadOperationBody::Integration(_)
188 | ThreadOperationBody::LocalIntegration(_) => ThreadFacet::Source,
189 ThreadOperationBody::Metadata(_) => ThreadFacet::Metadata,
190 ThreadOperationBody::Discussion(_) | ThreadOperationBody::Context(_) => {
191 ThreadFacet::Discussion
192 }
193 }
194 }
195
196 pub fn source_author(&self) -> Result<Option<SourceAuthor>> {
199 if let ThreadOperationBody::Capture(capture) = &self.body {
200 return Ok(Some(capture.author.clone()));
201 }
202 Ok(self
203 .local_integration()?
204 .map(|integration| integration.author))
205 }
206
207 pub fn source_result(&self) -> Result<Option<Capture>> {
208 match &self.body {
209 ThreadOperationBody::Capture(capture) => Ok(Some(capture.result.clone())),
210 ThreadOperationBody::Integration(bytes) => {
211 Ok(Some(integration::HostedIntegration::decode(bytes)?.result))
212 }
213 ThreadOperationBody::LocalIntegration(bytes) => Ok(Some(
214 local_integration::LocalIntegration::decode(bytes)?.result,
215 )),
216 _ => Ok(None),
217 }
218 }
219
220 pub fn source_state(&self) -> Result<Option<State>> {
222 match &self.body {
223 ThreadOperationBody::Capture(bytes) => bytes.result.validated_state().map(Some),
224 ThreadOperationBody::Integration(bytes) => {
225 integration::HostedIntegration::decode(bytes)?
226 .resulting_state()
227 .map(Some)
228 }
229 ThreadOperationBody::LocalIntegration(bytes) => {
230 local_integration::LocalIntegration::decode(bytes)?
231 .resulting_state()
232 .map(Some)
233 }
234 _ => Ok(None),
235 }
236 }
237 pub fn local_integration(&self) -> Result<Option<local_integration::LocalIntegration>> {
238 match &self.body {
239 ThreadOperationBody::LocalIntegration(bytes) => {
240 local_integration::LocalIntegration::decode(bytes).map(Some)
241 }
242 _ => Ok(None),
243 }
244 }
245 pub fn integration(&self) -> Result<Option<integration::HostedIntegration>> {
246 match &self.body {
247 ThreadOperationBody::Integration(bytes) => {
248 integration::HostedIntegration::decode(bytes).map(Some)
249 }
250 _ => Ok(None),
251 }
252 }
253
254 pub fn context_revision(&self) -> Result<Option<crate::object::ContextRevision>> {
257 use crate::object::{CollaborationOperationBodyV1 as Body, CollaborationResolution};
258 match &self.body {
259 ThreadOperationBody::Context(bytes) => crate::object::ContextRevision::decode(bytes)
260 .map(Some)
261 .map_err(invalid),
262 ThreadOperationBody::Discussion(bytes) => {
263 let record = CollaborationOperationEnvelope::decode(bytes)
264 .map_err(invalid)?
265 .operation;
266 match record.body {
267 Body::Resolve {
268 resolution: CollaborationResolution::IntoContext { context },
269 } => Ok(Some(context)),
270 _ => Ok(None),
271 }
272 }
273 _ => Ok(None),
274 }
275 }
276
277 pub fn encode(&self) -> Result<Vec<u8>> {
278 if self.version != 1 {
279 return Err(invalid("unsupported Thread operation version"));
280 }
281 match &self.body {
282 ThreadOperationBody::Capture(bytes) => {
283 bytes.author.validate()?;
284 bytes.result.validated_state()?;
285 }
286 ThreadOperationBody::Integration(bytes) => {
287 integration::HostedIntegration::decode(bytes)?.validate_operation(self)?;
288 }
289 ThreadOperationBody::LocalIntegration(bytes) => {
290 local_integration::LocalIntegration::decode(bytes)?.validate_operation(self)?;
291 }
292 ThreadOperationBody::Metadata(bytes) => {
293 metadata::ThreadControl::decode(bytes)?.validate_operation(self)?;
294 }
295 ThreadOperationBody::Context(bytes) => {
296 let context = crate::object::ContextRevision::decode(bytes).map_err(invalid)?;
297 if context.metadata.scope.thread != Some(self.thread)
298 || context.parents.iter().copied().collect::<BTreeSet<_>>() != self.parents
299 {
300 return Err(invalid(
301 "context scope or causal parents differ from Thread operation",
302 ));
303 }
304 }
305 ThreadOperationBody::Discussion(bytes) => {
306 let decoded = CollaborationOperationEnvelope::decode(bytes).map_err(invalid)?;
307 if decoded
308 .operation
309 .metadata
310 .as_ref()
311 .is_some_and(|m| m.scope.thread != Some(self.thread))
312 {
313 return Err(invalid("collaboration metadata belongs to another Thread"));
314 }
315 if let Some(context) = self.context_revision()?
316 && (Some(&context.metadata) != decoded.operation.metadata.as_ref()
317 || context.extracted_from != Some(decoded.operation.discussion_id)
318 || context.parents.iter().copied().collect::<BTreeSet<_>>() != self.parents)
319 {
320 return Err(invalid(
321 "extracted context differs from signed discussion actor, scope or parents",
322 ));
323 }
324 if decoded.operation.encode().map_err(invalid)? != *bytes {
325 return Err(invalid("non-canonical discussion operation"));
326 }
327 }
328 }
329 let bytes = rmp_serde::to_vec_named(self)?;
330 bounded(&bytes)?;
331 Ok(bytes)
332 }
333
334 pub fn id(&self) -> Result<ContentHash> {
335 Ok(ContentHash::compute_typed(
336 OPERATION_FORMAT,
337 &self.encode()?,
338 ))
339 }
340
341 pub fn decode(bytes: &[u8]) -> Result<Self> {
342 bounded(bytes)?;
343 let operation: Self = rmp_serde::from_slice(bytes)?;
344 if operation.encode()? != bytes {
345 return Err(invalid("non-canonical Thread operation"));
346 }
347 Ok(operation)
348 }
349
350 pub fn validate_parents(&self, genesis: &ThreadGenesis, parents: &[Self]) -> Result<()> {
353 self.validate_parents_inner(genesis, parents, false)
354 }
355
356 fn validate_parents_inner(
358 &self,
359 genesis: &ThreadGenesis,
360 parents: &[Self],
361 imported: bool,
362 ) -> Result<()> {
363 if self.thread != genesis.id()? || parents.len() != self.parents.len() {
364 return Err(invalid("Thread or causal parent set mismatch"));
365 }
366 let ids = parents
367 .iter()
368 .map(Self::id)
369 .collect::<Result<BTreeSet<_>>>()?;
370 if ids != self.parents
371 || parents
372 .iter()
373 .any(|p| p.thread != self.thread || p.facet() != self.facet())
374 {
375 return Err(invalid("causal parents cross Thread or disclosure facet"));
376 }
377 if self
378 .source_result()?
379 .is_some_and(|result| result.source_targets.is_none())
380 {
381 for parent in parents {
382 if parent
383 .source_result()?
384 .is_some_and(|result| result.source_targets.is_some())
385 {
386 return Err(invalid(
387 "source evolution drops inherited reference closure",
388 ));
389 }
390 }
391 }
392 match &self.body {
393 ThreadOperationBody::Capture(bytes) => {
394 bytes.author.validate()?;
395 if let SourceAuthor::Account { spool, .. } = &bytes.author
396 && spool.to_string() != genesis.spool
397 {
398 return Err(invalid("original source author crosses Spool scope"));
399 }
400 let state = State::decode_current_msgpack(&bytes.result.state)?;
401 let mut source_parents = BTreeSet::new();
402 for parent in parents {
403 source_parents.insert(
404 parent
405 .source_state()?
406 .ok_or_else(|| invalid("capture parent is not source"))?
407 .id(),
408 );
409 }
410 let declared: BTreeSet<_> = state
411 .parents
412 .iter()
413 .copied()
414 .filter(|id| *id != genesis.base)
415 .collect();
416 if ((!imported || !self.parents.is_empty()) && declared != source_parents)
420 || (!imported && state.parents.is_empty())
421 || (imported
422 && (state.parents.contains(&genesis.base)
423 || genesis.base != initial_base::synthetic_initial_base()?.id()))
424 || state.parents.len() != state.parents.iter().collect::<BTreeSet<_>>().len()
425 {
426 return Err(invalid(
427 "capture source ancestry differs from causal parents",
428 ));
429 }
430 }
431 ThreadOperationBody::Integration(bytes) => {
432 let receipt = integration::HostedIntegration::decode(bytes)?;
433 receipt.validate_operation(self)?;
434 receipt.validate_parents(genesis, parents)?;
435 }
436 ThreadOperationBody::LocalIntegration(bytes) => {
437 let receipt = local_integration::LocalIntegration::decode(bytes)?;
438 receipt.validate_operation(self)?;
439 receipt.validate_parents(genesis, parents)?;
440 }
441 ThreadOperationBody::Metadata(bytes) => {
442 metadata::ThreadControl::decode(bytes)?.validate_parents(genesis, parents)?;
443 }
444 ThreadOperationBody::Context(bytes) => {
445 let context = crate::object::ContextRevision::decode(bytes).map_err(invalid)?;
446 if context.metadata.scope.spool.to_string() != genesis.spool {
447 return Err(invalid("context belongs to another spool"));
448 }
449 if parents.is_empty() && context.extracted_from.is_some() {
450 return Err(invalid(
451 "context extraction requires a signed discussion resolution",
452 ));
453 }
454 for parent in parents {
455 let parent = parent.context_revision()?.ok_or_else(|| {
456 invalid("context parent is not a context revision or extraction")
457 })?;
458 if parent.id != context.id
459 || parent.metadata.scope != context.metadata.scope
460 || parent.extracted_from != context.extracted_from
461 {
462 return Err(invalid("context parents belong to another record or scope"));
463 }
464 }
465 }
466 ThreadOperationBody::Discussion(bytes) => {
467 let operation = CollaborationOperationEnvelope::decode(bytes).map_err(invalid)?;
468 if operation
469 .operation
470 .metadata
471 .as_ref()
472 .is_some_and(|m| m.scope.spool.to_string() != genesis.spool)
473 {
474 return Err(invalid("collaboration metadata belongs to another spool"));
475 }
476 let mut discussion_parents = BTreeSet::new();
477 for parent in parents {
478 let ThreadOperationBody::Discussion(bytes) = &parent.body else {
479 return Err(invalid("discussion parent is not discussion"));
480 };
481 let parent = CollaborationOperationEnvelope::decode(bytes).map_err(invalid)?;
482 if parent.operation.discussion_id != operation.operation.discussion_id {
483 return Err(invalid("parents cross discussions"));
484 }
485 discussion_parents.insert(parent.operation_id);
486 }
487 if discussion_parents != operation.operation.parents.iter().copied().collect() {
488 return Err(invalid(
489 "discussion causality differs from its canonical operation",
490 ));
491 }
492 }
493 }
494 Ok(())
495 }
496}
497
498fn bounded(bytes: &[u8]) -> Result<()> {
499 if bytes.len() > MAX_OPERATION_BYTES {
500 return Err(invalid("Thread record exceeds the durable record bound"));
501 }
502 Ok(())
503}
504
505fn invalid(message: impl std::fmt::Display) -> HeddleError {
506 HeddleError::InvalidObject(message.to_string())
507}
508
509#[cfg(test)]
510mod capture_shape_tests {
511 use super::*;
512 use crate::object::{Attribution, Principal, Tree};
513 #[test]
514 fn byte_only_capture_shape_is_rejected_in_clean_cutover() {
515 #[derive(serde::Serialize)]
516 struct OldBody {
517 kind: &'static str,
518 canonical: Vec<u8>,
519 }
520 #[derive(serde::Serialize)]
521 struct OldOperation {
522 version: u16,
523 thread: ContentHash,
524 parents: BTreeSet<ContentHash>,
525 publisher: [u8; 32],
526 body: OldBody,
527 }
528 let state = State::new_snapshot(
529 Tree::new().hash(),
530 vec![],
531 Attribution::human(Principal::new("author", "")),
532 );
533 let bytes = rmp_serde::to_vec_named(&OldOperation {
534 version: 1,
535 thread: ContentHash::from_bytes([1; 32]),
536 parents: BTreeSet::new(),
537 publisher: [41; 32],
538 body: OldBody {
539 kind: "capture",
540 canonical: state.encode_current_msgpack().expect("state"),
541 },
542 })
543 .expect("old wire bytes");
544 assert!(
545 ThreadOperation::decode(&bytes).is_err(),
546 "capture has one typed shape, without a legacy byte decoder"
547 );
548 }
549}