Skip to main content

heddle_object_model/object/
thread_replication.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Portable Thread identity and immutable replication operations. Source and
3//! discussion causality have separate graphs so selective sharing is closed.
4pub 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/// Ownership fixed by the original signed genesis. Hosting a key-owned Thread
30/// requires a separately verified, explicit ownership claim.
31#[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    /// Canonical non-nil UUID shared by local and hosted replicas. Mutable
52    /// namespace/name addresses and native spool-link locators are not identity.
53    pub spool: String,
54    pub parent: Option<ContentHash>,
55    pub base: StateId,
56    pub name: String,
57    pub intent: String,
58    /// Immutable original ownership; uploading does not transfer it.
59    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/// Durable causal admission. Receiving bytes alone does not accept an operation.
107#[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)] // wire codec; boxing would change MessagePack layout
117pub 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/// Source result shared by authored captures and executor-derived integrations.
128/// Per-Thread reference metadata never changes State identity.
129#[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    /// Original author-signed privacy declarations; independent of the courier.
135    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    /// Signing key of the publisher of this immutable operation. Attribution
154    /// inside a source/discussion object is retained, not promoted to authority.
155    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    /// Uniform source result for capture and both integration authorities.
198    /// Original authored source identity. Hosted execution uses its executor proof.
199    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    /// The exact source revision represented by a capture or hosted integration.
225    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    /// A context root may be extracted atomically by a discussion resolution.
271    /// Later context revisions retain that original signed resolution as parent.
272    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    /// Called after every parent is present. A peer cannot relabel a private
370    /// dependency as public or attach an operation to a different Thread.
371    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}