1use 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#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
25#[serde(deny_unknown_fields)]
26pub struct HostedIntegration {
27 pub version: u16,
28 pub spool: Uuid,
29 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 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 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 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#[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#[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}