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