Skip to main content

heddle_object_model/object/thread_replication/
integration.rs

1//! Hosted landing is an executor attestation, never a human signature or a
2//! transferable grant. Trust comes from the receiver's selected remote and
3//! verified immutable Spool genesis, independently of incoming operation bytes.
4//! HYBRID receivers resolve a purpose-4 witness through crypto/repo's trust
5//! layer before using these unchanged source, executor and ancestry bindings.
6//! Delegated import content has its own typed model and supplies no executor.
7use std::collections::BTreeSet;
8
9use serde::{Deserialize, Serialize};
10use uuid::Uuid;
11
12use super::{ThreadGenesis, ThreadOperation, bounded, invalid};
13use crate::{
14    error::Result,
15    object::{ContentHash, State, StateId},
16};
17
18pub const SPOOL_GENESIS_TRUST_FORMAT: &str = "heddle-spool-owner-genesis-executor-trust-v1";
19pub const EXECUTION_FORMAT: &str = "heddle-hosted-integration-v1";
20pub const MAX_EXECUTION_EVIDENCE: usize = 128;
21
22/// This receipt is signed by the enclosing Thread operation's publisher. Source
23/// originals stay in their own Thread; this fact evolves only the target graph.
24#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
25#[serde(deny_unknown_fields)]
26pub struct HostedIntegration {
27    pub version: u16,
28    pub spool: Uuid,
29    /// Digest of verified immutable Spool owner genesis, not its mutable owner.
30    pub spool_genesis: ContentHash,
31    pub executor: [u8; 32],
32    pub source_thread: ContentHash,
33    pub source_operation: ContentHash,
34    pub source_revision: StateId,
35    pub target_thread: ContentHash,
36    pub expected_target_frontier: BTreeSet<ContentHash>,
37    /// Exact canonical resulting State. Content installation is a separate,
38    /// authorized closure transfer and never writes a checkout on receive.
39    pub result: super::Capture,
40    pub initiating_request_proof: ContentHash,
41    pub review_policy_version: ContentHash,
42    pub review_evidence: BTreeSet<ContentHash>,
43    pub executed_at_ms: i64,
44}
45impl HostedIntegration {
46    pub fn encode(&self) -> Result<Vec<u8>> {
47        if self.version != 1
48            || self.spool.is_nil()
49            || self.source_thread == self.target_thread
50            || self.expected_target_frontier.len() > 128
51            || self.review_evidence.len() > MAX_EXECUTION_EVIDENCE
52            || self.executed_at_ms < 0
53        {
54            return Err(invalid("invalid or unbounded hosted integration receipt"));
55        }
56        let state = self.result.validated_state()?;
57        if state.encode_current_msgpack()? != self.result.state {
58            return Err(invalid("non-canonical integration result"));
59        }
60        let bytes = rmp_serde::to_vec_named(self)?;
61        bounded(&bytes)?;
62        Ok(bytes)
63    }
64    pub fn decode(bytes: &[u8]) -> Result<Self> {
65        bounded(bytes)?;
66        let receipt: Self = rmp_serde::from_slice(bytes)?;
67        if receipt.encode()? != bytes {
68            return Err(invalid("non-canonical hosted integration receipt"));
69        }
70        Ok(receipt)
71    }
72    pub fn id(&self) -> Result<ContentHash> {
73        Ok(ContentHash::compute_typed(
74            EXECUTION_FORMAT,
75            &self.encode()?,
76        ))
77    }
78    pub fn resulting_state(&self) -> Result<State> {
79        self.result.validated_state()
80    }
81    /// Supply the independently authenticated original source operation.
82    pub fn validate_source(&self, source: &ThreadOperation) -> Result<()> {
83        if source.thread != self.source_thread
84            || source.id()? != self.source_operation
85            || source
86                .source_state()?
87                .is_none_or(|state| state.id() != self.source_revision)
88        {
89            return Err(invalid(
90                "integration differs from original source operation",
91            ));
92        }
93        if self.result.source_targets.is_none()
94            && source
95                .source_result()?
96                .is_some_and(|result| result.source_targets.is_some())
97        {
98            return Err(invalid("integration drops source reference closure"));
99        }
100        Ok(())
101    }
102    pub(super) fn validate_operation(&self, operation: &ThreadOperation) -> Result<()> {
103        if self.target_thread != operation.thread
104            || self.executor != operation.publisher
105            || self.expected_target_frontier != operation.parents
106        {
107            return Err(invalid(
108                "integration receipt differs from signed executor, target or frontier",
109            ));
110        }
111        Ok(())
112    }
113    pub(super) fn validate_parents(
114        &self,
115        genesis: &ThreadGenesis,
116        parents: &[ThreadOperation],
117    ) -> Result<()> {
118        if self.spool.to_string() != genesis.spool {
119            return Err(invalid("integration receipt belongs to another Spool"));
120        }
121        let state = self.resulting_state()?;
122        // A fast-forward retains the original State identity; the trusted
123        // executor attests its ancestry. A merge must explicitly retain both
124        // the selected source revision and every target source parent.
125        if state.id() != self.source_revision {
126            let mut expected = BTreeSet::from([self.source_revision]);
127            if parents.is_empty() {
128                expected.insert(genesis.base);
129            }
130            for parent in parents {
131                expected.insert(
132                    parent
133                        .source_state()?
134                        .ok_or_else(|| invalid("integration parent is not source"))?
135                        .id(),
136                );
137            }
138            if state.parents.iter().copied().collect::<BTreeSet<_>>() != expected
139                || state.parents.len() != expected.len()
140            {
141                return Err(invalid(
142                    "integration result drops source or target ancestry",
143                ));
144            }
145        }
146        Ok(())
147    }
148}
149
150/// Construct only from receiver-owned remote configuration and verified Spool
151/// genesis. Copying these fields from an incoming receipt establishes no trust.
152/// A hosted server also requires its persisted original execution receipt.
153#[derive(Clone, Debug, PartialEq, Eq)]
154pub struct TrustedHostedExecutor {
155    pub spool: Uuid,
156    pub spool_genesis: ContentHash,
157    pub executor: [u8; 32],
158}
159impl TrustedHostedExecutor {
160    pub fn authorize(&self, operation: &ThreadOperation) -> Result<()> {
161        let receipt = operation
162            .hosted_execution_binding()?
163            .ok_or_else(|| invalid("executor trust requires a hosted execution"))?;
164        if receipt.spool != self.spool
165            || receipt.spool_genesis != self.spool_genesis
166            || receipt.executor != self.executor
167        {
168            return Err(invalid(
169                "hosted execution has no independently trusted executor for this Spool genesis",
170            ));
171        }
172        Ok(())
173    }
174}
175
176/// Extracted after canonical receipt and enclosing operation binding checks.
177/// It identifies an execution; receiving these bytes never establishes trust.
178#[derive(Clone, Copy, Debug, PartialEq, Eq)]
179pub struct HostedExecutionBinding {
180    pub kind: HostedExecutionKind,
181    pub spool: Uuid,
182    pub spool_genesis: ContentHash,
183    pub executor: [u8; 32],
184    pub initiating_request_proof: ContentHash,
185    pub review_policy_version: Option<ContentHash>,
186    pub executed_at_ms: i64,
187}
188#[derive(Clone, Copy, Debug, PartialEq, Eq)]
189pub enum HostedExecutionKind {
190    Integration,
191}
192impl ThreadOperation {
193    pub fn hosted_execution_binding(&self) -> Result<Option<HostedExecutionBinding>> {
194        if let Some(value) = self.integration()? {
195            value.validate_operation(self)?;
196            return Ok(Some(HostedExecutionBinding {
197                kind: HostedExecutionKind::Integration,
198                spool: value.spool,
199                spool_genesis: value.spool_genesis,
200                executor: value.executor,
201                initiating_request_proof: value.initiating_request_proof,
202                review_policy_version: Some(value.review_policy_version),
203                executed_at_ms: value.executed_at_ms,
204            }));
205        }
206        Ok(None)
207    }
208}
209
210#[cfg(test)]
211mod tests {
212    use super::*;
213    use crate::object::{Attribution, Principal, Tree, thread_replication::ThreadOperationBody};
214
215    fn fixture() -> (
216        ThreadGenesis,
217        ThreadOperation,
218        HostedIntegration,
219        ThreadOperation,
220    ) {
221        let genesis = ThreadGenesis {
222            version: 1,
223            spool: Uuid::from_u128(7).to_string(),
224            parent: None,
225            base: StateId::from_bytes([1; 32]),
226            name: "target".into(),
227            intent: "hosted landing".into(),
228            owner: crate::object::thread_replication::GenesisOwner::Account(Uuid::from_u128(12)),
229            creator: [2; 32],
230            nonce: vec![3; 16],
231        };
232        let target = State::new_snapshot(
233            Tree::new().hash(),
234            vec![genesis.base],
235            Attribution::human(Principal::new("target", "target@example.test")),
236        );
237        let target_operation = ThreadOperation {
238            version: 1,
239            thread: genesis.id().expect("Thread"),
240            parents: BTreeSet::new(),
241            publisher: [2; 32],
242            body: ThreadOperationBody::Capture(
243                crate::object::thread_replication::AuthoredCapture::local(
244                    target.encode_current_msgpack().expect("target").into(),
245                ),
246            ),
247        };
248        let source = State::new_snapshot(
249            Tree::new().hash(),
250            vec![target.id()],
251            Attribution::human(Principal::new("source", "source@example.test")),
252        );
253        let receipt = HostedIntegration {
254            version: 1,
255            spool: Uuid::from_u128(7),
256            spool_genesis: ContentHash::from_bytes([4; 32]),
257            executor: [5; 32],
258            source_thread: ContentHash::from_bytes([6; 32]),
259            source_operation: ContentHash::from_bytes([7; 32]),
260            source_revision: source.id(),
261            target_thread: genesis.id().expect("Thread"),
262            expected_target_frontier: BTreeSet::from([target_operation.id().expect("parent")]),
263            result: source.encode_current_msgpack().expect("source").into(),
264            initiating_request_proof: ContentHash::from_bytes([8; 32]),
265            review_policy_version: ContentHash::from_bytes([9; 32]),
266            review_evidence: BTreeSet::from([ContentHash::from_bytes([10; 32])]),
267            executed_at_ms: 100,
268        };
269        let operation = ThreadOperation {
270            version: 1,
271            thread: receipt.target_thread,
272            parents: receipt.expected_target_frontier.clone(),
273            publisher: receipt.executor,
274            body: ThreadOperationBody::Integration(receipt.encode().expect("receipt")),
275        };
276        (genesis, target_operation, receipt, operation)
277    }
278    #[test]
279    fn integration_is_bound_to_executor_target_frontier_and_scope() {
280        let (genesis, parent, receipt, operation) = fixture();
281        operation
282            .validate_parents(&genesis, std::slice::from_ref(&parent))
283            .expect("valid integration");
284        assert_eq!(
285            ThreadOperation::decode(&operation.encode().expect("encode")).expect("canonical"),
286            operation
287        );
288        let mut changed = operation.clone();
289        changed.publisher = [99; 32];
290        assert!(changed.encode().is_err(), "receipt cannot change executor");
291        changed = operation.clone();
292        changed.parents.clear();
293        assert!(
294            changed.encode().is_err(),
295            "receipt cannot change target frontier"
296        );
297        changed = operation.clone();
298        changed.thread = ContentHash::from_bytes([99; 32]);
299        assert!(
300            changed.encode().is_err(),
301            "receipt cannot change target Thread"
302        );
303        let mut wrong_scope = receipt;
304        wrong_scope.spool = Uuid::from_u128(99);
305        changed = operation;
306        changed.body =
307            ThreadOperationBody::Integration(wrong_scope.encode().expect("structural receipt"));
308        assert!(
309            changed.validate_parents(&genesis, &[parent]).is_err(),
310            "receipt cannot change immutable Spool scope"
311        );
312    }
313    #[test]
314    fn integration_source_evolution_preserves_merge_parents_and_accepts_later_captures() {
315        let (genesis, parent, mut receipt, mut operation) = fixture();
316        let target = parent.source_state().expect("source").expect("target");
317        let merge = State::new_snapshot(
318            Tree::new().hash(),
319            vec![target.id(), receipt.source_revision],
320            Attribution::human(Principal::new("executor", "weft@example.test")),
321        );
322        receipt.result = merge.encode_current_msgpack().expect("merge").into();
323        operation.body = ThreadOperationBody::Integration(receipt.encode().expect("receipt"));
324        operation
325            .validate_parents(&genesis, std::slice::from_ref(&parent))
326            .expect("explicit merge ancestry");
327        let after = State::new_snapshot(
328            Tree::new().hash(),
329            vec![merge.id()],
330            Attribution::human(Principal::new("agent", "agent@example.test")),
331        );
332        let capture = ThreadOperation {
333            version: 1,
334            thread: operation.thread,
335            parents: BTreeSet::from([operation.id().expect("integration")]),
336            publisher: [3; 32],
337            body: ThreadOperationBody::Capture(
338                crate::object::thread_replication::AuthoredCapture::local(
339                    after.encode_current_msgpack().expect("capture").into(),
340                ),
341            ),
342        };
343        capture
344            .validate_parents(&genesis, std::slice::from_ref(&operation))
345            .expect("capture after hosted integration");
346        let dropped = State::new_snapshot(
347            Tree::new().hash(),
348            vec![target.id()],
349            Attribution::human(Principal::new("executor", "weft@example.test")),
350        );
351        receipt.result = dropped
352            .encode_current_msgpack()
353            .expect("dropped ancestry")
354            .into();
355        operation.body = ThreadOperationBody::Integration(receipt.encode().expect("receipt"));
356        assert!(
357            operation.validate_parents(&genesis, &[parent]).is_err(),
358            "merge cannot drop source ancestry"
359        );
360    }
361    #[test]
362    fn integration_signature_identity_does_not_create_executor_trust() {
363        let (_, _, _, operation) = fixture();
364        let trust = TrustedHostedExecutor {
365            spool: Uuid::from_u128(7),
366            spool_genesis: ContentHash::from_bytes([4; 32]),
367            executor: [5; 32],
368        };
369        trust
370            .authorize(&operation)
371            .expect("independent selected remote pin");
372        let mut different = trust.clone();
373        different.executor = [9; 32];
374        assert!(
375            different.authorize(&operation).is_err(),
376            "different endpoint is not user authority"
377        );
378        different = trust.clone();
379        different.spool_genesis = ContentHash::from_bytes([9; 32]);
380        assert!(
381            different.authorize(&operation).is_err(),
382            "same UUID does not replace immutable genesis"
383        );
384        different = trust;
385        different.spool = Uuid::from_u128(9);
386        assert!(
387            different.authorize(&operation).is_err(),
388            "executor trust is scoped to one Spool"
389        );
390    }
391}