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 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/// Ownership fixed by the original signed genesis. Hosting a key-owned Thread
31/// requires a separately verified, explicit ownership claim.
32#[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    /// Canonical non-nil UUID shared by local and hosted replicas. Mutable
53    /// namespace/name addresses and native spool-link locators are not identity.
54    pub spool: String,
55    pub parent: Option<ContentHash>,
56    pub base: StateId,
57    pub name: String,
58    pub intent: String,
59    /// Immutable original ownership; uploading does not transfer it.
60    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/// Durable causal admission. Receiving bytes alone does not accept an operation.
108#[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)] // wire codec; boxing would change MessagePack layout
118pub 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/// 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::LocalIntegration(_) => ThreadFacet::Source,
189            ThreadOperationBody::Metadata(_) => ThreadFacet::Metadata,
190            ThreadOperationBody::Discussion(_) | ThreadOperationBody::Context(_) => {
191                ThreadFacet::Discussion
192            }
193        }
194    }
195
196    /// Uniform source result for capture and both integration authorities.
197    /// Original authored source identity. Hosted execution uses its executor proof.
198    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    /// The exact source revision represented by a capture or hosted integration.
221    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    /// A context root may be extracted atomically by a discussion resolution.
255    /// Later context revisions retain that original signed resolution as parent.
256    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    /// Called after every parent is present. A peer cannot relabel a private
351    /// dependency as public or attach an operation to a different Thread.
352    pub fn validate_parents(&self, genesis: &ThreadGenesis, parents: &[Self]) -> Result<()> {
353        self.validate_parents_inner(genesis, parents, false)
354    }
355
356    // Only an authenticated DelegatedImport can select imported ancestry.
357    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                // A first import has one native tip operation. Its converted
417                // Git ancestors travel as State closure, attested by the
418                // converter witness, rather than causal native operations.
419                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}