1use crate::heddle::api::v1alpha2::*;
6pub use crate::hybrid_codec::{Reject, strict_decode};
7use crate::hybrid_codec::{canonical, field, hash, key_id, record, signing_digest, verify, width};
8
9pub const PERMISSION_DOMAIN: &str = "heddle-import-member-permission-v1";
10pub const GENESIS_DOMAIN: &str = "heddle-import-genesis-authority-v1";
11pub const DELEGATION_DOMAIN: &str = "heddle-import-job-delegation-v1";
12pub const OPERATION_DOMAIN: &str = "heddle-delegated-import-operation-v1";
13pub const MANIFEST_DOMAIN: &str = "heddle-import-result-manifest-v1";
14pub const PUBLICATION_DOMAIN: &str = "heddle-import-publication-payload-v1";
15pub const MAX_BRANCHES: usize = 256;
16pub const MAX_REF_BYTES: usize = 1024;
17pub const MAX_RECORD_BYTES: usize = 64 * 1024;
18pub const MAX_BUNDLE_BYTES: usize = 1024 * 1024;
19pub const MAX_DELEGATION_WINDOW_SECONDS: u64 = 7 * 24 * 60 * 60;
20pub const MAX_COMMIT_REQUEST_BYTES: usize = 2 * MAX_BUNDLE_BYTES;
21pub const CANCELLATION_NAMESPACE: &str = "heddle-import-cancel-v1";
22
23record!(AuthorizationSignature, signer_key_id:b, signature:b);
24record!(RecordSignature, public_key:b, signature:b);
25record!(SignedRecord, format:s, canonical_record:b, signatures:l);
26record!(ImportFrontierV1, format_version:u, thread_id:b, operation_ids:h);
27record!(ImportContentV1, format_version:u, canonical_capture:b);
28record!(ImportBoundaryAcceptanceV1, binding:m, signed_acceptance:m, originals_manifest:b, publication_intent:b, original_receipts:l);
29record!(ImportGenesisWitnessV1, format_version:u, binding:m, original_genesis:m, creator_authority_envelope:b, boundary_acceptance:o);
30record!(ImportAuthorityWitnessV1, format_version:u, kind:e, original:m, dependencies:l, authority_envelope:b, boundary_acceptances:l);
31record!(HostedLandingRequestProofV1, format_version:u, signing_identity:s, method_path:s, timestamp_millis:u, nonce:b, request_body:b, signature:m);
32record!(HostedLandingWitnessV1, format_version:u, execution:m, request:m, source_operation:m, review_evidence:l, authority_envelope:b);
33record!(ImportIdentityV1, spool_uuid:b, spool_genesis_digest:b, owner_id:b,
34 owner_account_uuid:b, owner_state_hash:b, ownership_transfer_sequence:u);
35record!(ImportOwnerChainV1, spool_genesis_digest:b, owner_state_hashes:q, transfer_audit_hashes:q);
36record!(ImportBranchLimitV1, ref_name:s, hash_algorithm:e, ref_mode:e, pinned_commit_oid:b,
37 genesis_digest:b, target_thread_id:b, expected_frontier_digest:b, slot_id:u, ref_disclosure:e);
38record!(ImportPermissionScopeV1, provider:s, source_url:s, branches:l, destination_version:b,
39 options_digest:b, converter_version:s, max_operations:u, max_result_bytes:u);
40record!(ImportMemberPermissionV1, format_version:u, identity:m, logical_job_id:b,
41 retry_lineage_id:b, subject_public_key:b, purpose:e, scope:m, not_before_unix_seconds:u,
42 expires_at_unix_seconds:u, cancellation_id:b, owner_chain_digest:b, nonce:b);
43record!(SignedImportMemberPermissionV1, body:m, owner_signature:m);
44record!(ImportGenesisAuthorityV1, format_version:u, identity:m, genesis_digest:b,
45 original_creator_signature:b, creator_public_key:b, creator_authority_envelope_digest:b,
46 parent_permission_digest:b, owner_chain_digest:b);
47record!(SignedImportGenesisAuthorityV1, body:m, creator_signature:m);
48record!(ImportBranchManifestV1, limit:m, genesis_authority_digest:b);
49record!(ImportJobDelegationV1, format_version:u, identity:m, delegation_id:b, logical_job_id:b,
50 retry_lineage_id:b, job_public_key:b, job_key_id:b, delegating_public_key:b,
51 parent_permission_digest:b, owner_chain_digest:b, purpose:e, scope:m, branch_manifest:l,
52 not_before_unix_seconds:u, expires_at_unix_seconds:u, cancellation_id:b);
53record!(ImportJobPreparationV1, format_version:u, identity:m, delegation_id:b, logical_job_id:b,
54 retry_lineage_id:b, job_public_key:b, job_key_id:b, owner_chain_digest:b, purpose:e,
55 scope:m, cancellation_id:b);
56record!(SignedImportJobDelegationV1, body:m, delegating_signature:m);
57record!(ImportCommittedSlotV1, ref_name:s, slot_id:u, signed_operation_digest:b,
58 resulting_frontier_digest:b, result_bytes:u);
59record!(ImportResultManifestV1, format_version:u, logical_job_id:b, retry_lineage_id:b, slots:l);
60record!(DelegatedImportOperationV1, format_version:u, spool_uuid:b, spool_genesis_digest:b,
61 logical_job_id:b, retry_lineage_id:b, physical_operation_id:b, delegation_digest:b,
62 ref_name:s, slot_id:u, hash_algorithm:e, observed_commit_oid:b, genesis_digest:b,
63 target_thread_id:b, expected_frontier_digest:b, resulting_frontier_digest:b,
64 resulting_content_digest:b, result_bytes:u, options_digest:b, converter_version:s);
65record!(SignedDelegatedImportOperationV1, body:m, job_signature:m);
66record!(ImportPublicationWitnessV1, format_version:u, signed_operation_digest:b, delegation_digest:b,
67 logical_job_id:b, retry_lineage_id:b, physical_operation_id:b, ref_name:s, slot_id:u,
68 hash_algorithm:e, observed_commit_oid:b, expected_frontier_digest:b, resulting_frontier_digest:b,
69 terminal_manifest_digest:b);
70
71pub fn signed_permission_digest(v: &SignedImportMemberPermissionV1) -> Result<Vec<u8>, Reject> {
72 signing_digest("heddle-signed-import-member-permission-v1", v)
73}
74pub fn owner_chain_digest(v: &ImportOwnerChainV1) -> Result<Vec<u8>, Reject> {
75 width(&v.spool_genesis_digest, 32)?;
76 if v.owner_state_hashes.is_empty()
77 || v.owner_state_hashes.len() > 64
78 || v.transfer_audit_hashes.len() > 64
79 {
80 return Err(Reject::Bounds);
81 }
82 for h in v.owner_state_hashes.iter().chain(&v.transfer_audit_hashes) {
83 width(h, 32)?;
84 }
85 if v.owner_state_hashes.windows(2).any(|w| w[0] >= w[1]) {
86 return Err(Reject::Canonical);
87 }
88 signing_digest("heddle-import-owner-chain-v1", v)
89}
90pub fn signed_genesis_digest(v: &SignedImportGenesisAuthorityV1) -> Result<Vec<u8>, Reject> {
91 signing_digest("heddle-signed-import-genesis-authority-v1", v)
92}
93pub fn signed_delegation_digest(v: &SignedImportJobDelegationV1) -> Result<Vec<u8>, Reject> {
94 signing_digest("heddle-signed-import-job-delegation-v1", v)
95}
96pub fn signed_operation_digest(v: &SignedDelegatedImportOperationV1) -> Result<Vec<u8>, Reject> {
97 signing_digest("heddle-signed-delegated-import-operation-v1", v)
98}
99pub fn manifest_digest(v: &ImportResultManifestV1) -> Result<Vec<u8>, Reject> {
100 signing_digest(MANIFEST_DOMAIN, v)
101}
102pub fn verify_authorization_signature(
103 key: &[u8],
104 domain: &str,
105 body: &impl crate::hybrid_codec::Canonical,
106 signature: &AuthorizationSignature,
107) -> Result<(), Reject> {
108 if signature.signer_key_id != key_id(key) {
109 return Err(Reject::Signature);
110 }
111 let bytes = canonical(body)?;
112 if bytes.len() > MAX_RECORD_BYTES {
113 return Err(Reject::Bounds);
114 }
115 verify(
116 key,
117 &hash(&[domain.as_bytes(), &bytes]),
118 &signature.signature,
119 )
120}
121
122pub fn canonical_https(value: &str, origin: bool) -> Result<(), Reject> {
126 if value.len() > 2048 {
127 return Err(Reject::Bounds);
128 }
129 let rest = value.strip_prefix("https://").ok_or(Reject::Canonical)?;
130 let (host, path) = rest.split_once('/').unwrap_or((rest, ""));
131 if host.is_empty()
132 || host.len() > 253
133 || host.split('.').any(|part| {
134 part.is_empty()
135 || part.len() > 63
136 || part.starts_with('-')
137 || part.ends_with('-')
138 || !part
139 .bytes()
140 .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-')
141 })
142 || (origin && rest != host)
143 || (!origin && path.is_empty())
144 || path.split('/').any(|p| {
145 p.is_empty() && !origin
146 || p == "."
147 || p == ".."
148 || !p
149 .bytes()
150 .all(|b| b.is_ascii_alphanumeric() || b"-._~".contains(&b))
151 })
152 {
153 return Err(Reject::Canonical);
154 }
155 Ok(())
156}
157fn identity(value: &ImportIdentityV1) -> Result<(), Reject> {
158 for v in [&value.spool_uuid, &value.owner_account_uuid] {
159 width(v, 16)?;
160 if v.iter().all(|b| *b == 0) {
161 return Err(Reject::Canonical);
162 }
163 }
164 for v in [
165 &value.spool_genesis_digest,
166 &value.owner_id,
167 &value.owner_state_hash,
168 ] {
169 width(v, 32)?;
170 }
171 Ok(())
172}
173fn interval(start: i64, end: i64, now: i64) -> Result<(), Reject> {
174 if start < 0 || end <= start {
175 return Err(Reject::Semantic);
176 }
177 if now < start || now >= end {
178 return Err(Reject::Expired);
179 }
180 Ok(())
181}
182fn validity(start: i64, end: i64, now: i64, current: bool) -> Result<(), Reject> {
184 if start < 0 || end <= start {
185 return Err(Reject::Semantic);
186 }
187 if current {
188 interval(start, end, now)?;
189 }
190 Ok(())
191}
192fn branch_ref(name: &str) -> Result<(), Reject> {
197 let short = name.strip_prefix("refs/heads/").ok_or(Reject::Canonical)?;
198 if short.is_empty()
199 || short.starts_with('-')
200 || short == "HEAD"
201 || name.ends_with('/')
202 || name.ends_with('.')
203 || name.contains("..")
204 || name.contains("//")
205 || name.contains("@{")
206 || name
207 .split('/')
208 .any(|p| p.starts_with('.') || p.ends_with(".lock"))
209 || name.bytes().any(|b| {
210 b <= b' ' || b == 0x7f || matches!(b, b'~' | b'^' | b':' | b'?' | b'*' | b'[' | b'\\')
211 })
212 {
213 return Err(Reject::Canonical);
214 }
215 if name.len() > MAX_REF_BYTES {
216 return Err(Reject::Bounds);
217 }
218 Ok(())
219}
220
221fn branch(value: &ImportBranchLimitV1) -> Result<(), Reject> {
222 branch_ref(&value.ref_name)?;
223 let size = match value.hash_algorithm {
224 1 => 20,
225 2 => 32,
226 _ => return Err(Reject::Version),
227 };
228 match value.ref_mode {
229 1 => {
230 width(&value.pinned_commit_oid, size)?;
231 if value.ref_disclosure != 0 {
232 return Err(Reject::RefDisclosure);
233 }
234 }
235 2 if value.pinned_commit_oid.is_empty() => {
236 if value.ref_disclosure != 1 {
237 return Err(Reject::RefDisclosure);
238 }
239 }
240 _ => return Err(Reject::Semantic),
241 }
242 for v in [
243 &value.genesis_digest,
244 &value.target_thread_id,
245 &value.expected_frontier_digest,
246 ] {
247 width(v, 32)?;
248 }
249 Ok(())
250}
251
252pub fn validate_ref_selection(
255 value: &ImportBranchLimitV1,
256 known_commit_oid: Option<&[u8]>,
257) -> Result<(), Reject> {
258 branch(value)?;
259 if let Some(oid) = known_commit_oid {
260 width(oid, if value.hash_algorithm == 1 { 20 } else { 32 })?;
261 if value.ref_mode != 1 || value.pinned_commit_oid != oid {
262 return Err(Reject::RefPinning);
263 }
264 }
265 Ok(())
266}
267
268pub fn conversion_options_digest(version: &str, options: &[u8]) -> Result<Vec<u8>, Reject> {
270 if version.is_empty() || version.len() > 128 || !version.is_ascii() {
271 return Err(Reject::Canonical);
272 }
273 if options.len() > 4096 {
274 return Err(Reject::Bounds);
275 }
276 let mut bytes = Vec::new();
277 crate::hybrid_codec::counted(&mut bytes, version.as_bytes())?;
278 crate::hybrid_codec::counted(&mut bytes, options)?;
279 Ok(hash(&[b"heddle-import-conversion-options-v1", &bytes]))
280}
281
282pub fn resolve_import_provider(
285 source: &ProviderRepository,
286 connection_provider: Option<&str>,
287) -> Result<&'static str, Reject> {
288 canonical_https(&source.clone_url, false)?;
289 if source.provider_repository_id.len() > 4096 || source.name.len() > 4096 {
290 return Err(Reject::Bounds);
291 }
292 if let Some(connection) = &source.connection {
293 let positive = |v: &str| {
294 !v.is_empty()
295 && v.bytes().all(|b| b.is_ascii_digit())
296 && v.parse::<u64>().is_ok_and(|id| id > 0)
297 };
298 if connection_provider != Some("github")
299 || connection.spool.is_some()
300 || connection.id.is_empty()
301 || !positive(&source.provider_repository_id)
302 || !positive(&source.installation_id)
303 {
304 return Err(Reject::SourceSelection);
305 }
306 let path = source
307 .clone_url
308 .strip_prefix("https://github.com/")
309 .ok_or(Reject::SourceSelection)?;
310 let parts: Vec<_> = path.split('/').collect();
311 if parts.len() != 2
312 || parts[0].is_empty()
313 || parts[1].strip_suffix(".git").is_none_or(str::is_empty)
314 {
315 return Err(Reject::SourceSelection);
316 }
317 Ok("github")
318 } else {
319 if connection_provider.is_some()
320 || source.private
321 || !source.installation_id.is_empty()
322 || (!source.provider_repository_id.is_empty()
323 && source.provider_repository_id != source.clone_url)
324 {
325 return Err(Reject::SourceSelection);
326 }
327 Ok("public-git")
328 }
329}
330
331pub fn validate_repository_size_estimate(source: &ProviderRepository) -> Result<(), Reject> {
333 match source.size_estimate_state {
334 0 if source.git_size_kib == 0 => Ok(()),
335 1 if source.connection.is_some() => Ok(()),
336 _ => Err(Reject::Canonical),
337 }
338}
339
340pub fn validate_resolve_import_source_response(
342 request: &ResolveImportSourceRequest,
343 response: &ResolveImportSourceResponse,
344 connection_provider: Option<&str>,
345) -> Result<(), Reject> {
346 use prost::Message;
347 if response.encoded_len() > MAX_BUNDLE_BYTES {
348 return Err(Reject::Bounds);
349 }
350 let selected = request.source.as_ref().ok_or(Reject::SourceSelection)?;
351 let resolved = response.source.as_ref().ok_or(Reject::SourceSelection)?;
352 let provider = resolve_import_provider(selected, connection_provider)?;
353 if resolve_import_provider(resolved, connection_provider)? != provider
354 || selected.clone_url != resolved.clone_url
355 || selected.connection != resolved.connection
356 || selected.installation_id != resolved.installation_id
357 || selected.private != resolved.private
358 || (selected.provider_repository_id != resolved.provider_repository_id
359 && (selected.connection.is_some() || !selected.provider_repository_id.is_empty()))
360 || (resolved.connection.is_none() && resolved.provider_repository_id != selected.clone_url)
361 {
362 return Err(Reject::SourceSelection);
363 }
364 validate_repository_size_estimate(resolved)?;
365 validate_repository_hash_algorithm(resolved, false)
366}
367
368fn validate_retained_import_source(
370 selector: &ImportSourceSelectionV1,
371 scope: &ImportPermissionScopeV1,
372) -> Result<(), Reject> {
373 let source = ProviderRepository {
374 connection: selector.connection.clone(),
375 provider_repository_id: selector.provider_repository_id.clone(),
376 installation_id: selector.installation_id.clone(),
377 private: selector.private,
378 clone_url: scope.source_url.clone(),
379 ..Default::default()
380 };
381 let provider =
382 resolve_import_provider(&source, selector.connection.as_ref().map(|_| "github"))?;
383 if scope.provider != provider
384 || (selector.connection.is_none() && selector.provider_repository_id != scope.source_url)
385 {
386 return Err(Reject::SourceSelection);
387 }
388 Ok(())
389}
390
391pub fn validate_repository_hash_algorithm(
393 source: &ProviderRepository,
394 known: bool,
395) -> Result<(), Reject> {
396 let size = match source.hash_algorithm {
397 1 => Some(40),
398 2 => Some(64),
399 0 if !known => None,
400 _ => return Err(Reject::Version),
401 };
402 if source.refs.len() > 512 {
403 return Err(Reject::Bounds);
404 }
405 for (i, r) in source.refs.iter().enumerate() {
406 if r.hash_algorithm != source.hash_algorithm {
407 return Err(Reject::SourceSelection);
408 }
409 if i > 0 && source.refs[i - 1].name >= r.name {
410 return Err(Reject::Canonical);
411 }
412 if !r.head_oid.is_empty() {
413 let size = size.ok_or(Reject::Version)?;
414 if r.head_oid.len() != size {
415 return Err(Reject::SourceSelection);
416 }
417 if !r
418 .head_oid
419 .bytes()
420 .all(|b| b.is_ascii_digit() || (b'a'..=b'f').contains(&b))
421 {
422 return Err(Reject::Canonical);
423 }
424 }
425 }
426 Ok(())
427}
428
429pub fn validate_discovered_import_scope(
431 scope: &ImportPermissionScopeV1,
432 source: &ProviderRepository,
433) -> Result<(), Reject> {
434 validate_discovered_import_scope_inner(scope, source)
435}
436
437fn validate_discovered_import_scope_inner(
438 scope: &ImportPermissionScopeV1,
439 source: &ProviderRepository,
440) -> Result<(), Reject> {
441 validate_repository_hash_algorithm(source, true)?;
442 if scope.source_url != source.clone_url {
443 return Err(Reject::SourceSelection);
444 }
445 for b in &scope.branches {
446 if b.hash_algorithm != source.hash_algorithm {
447 return Err(Reject::SourceSelection);
448 }
449 let oid = source
450 .refs
451 .iter()
452 .find(|r| r.name == b.ref_name)
453 .filter(|r| !r.head_oid.is_empty())
454 .map(|r| hex::decode(&r.head_oid).map_err(|_| Reject::Canonical))
455 .transpose()?;
456 validate_ref_selection(b, oid.as_deref())?;
457 }
458 Ok(())
459}
460
461pub fn prepare_import_source_scope(
466 request: &PrepareImportJobRequest,
467 current_source: &ProviderRepository,
468 connection_provider: Option<&str>,
469 configuration: &GetImportConfigurationResponse,
470 current_destination_version: &[u8],
471) -> Result<ImportPermissionScopeV1, Reject> {
472 let selector = request.source.as_ref().ok_or(Reject::SourceSelection)?;
473 if selector.connection != current_source.connection
474 || (selector.provider_repository_id != current_source.provider_repository_id
475 && (selector.connection.is_some() || !selector.provider_repository_id.is_empty()))
476 || selector.installation_id != current_source.installation_id
477 || selector.private != current_source.private
478 {
479 return Err(Reject::SourceSelection);
480 }
481 let scope = request.proposed_scope.as_ref().ok_or(Reject::Canonical)?;
482 let provider = resolve_import_provider(current_source, connection_provider)?;
483 if scope.provider != provider {
484 return Err(Reject::SourceSelection);
485 }
486 validate_discovered_import_scope(scope, current_source)?;
487 validate_import_configuration(configuration)?;
488 prepare_scope(scope, configuration, current_destination_version)
489}
490
491fn validate_provider_support(
492 provider: &str,
493 configuration: &GetImportConfigurationResponse,
494) -> Result<(), Reject> {
495 if !configuration
496 .providers
497 .iter()
498 .any(|p| p.provider == provider)
499 {
500 return Err(Reject::SourceSelection);
501 }
502 Ok(())
503}
504
505pub fn validate_import_configuration(v: &GetImportConfigurationResponse) -> Result<(), Reject> {
506 use prost::Message;
507 if v.encoded_len() > MAX_BUNDLE_BYTES || v.converters.is_empty() || v.converters.len() > 32 {
508 return Err(Reject::Bounds);
509 }
510 for (i, c) in v.converters.iter().enumerate() {
511 for text in [&c.converter_version, &c.options_encoding] {
512 if text.is_empty() || text.len() > 128 || !text.is_ascii() {
513 return Err(Reject::Canonical);
514 }
515 }
516 if i > 0 && v.converters[i - 1].converter_version >= c.converter_version {
517 return Err(Reject::Canonical);
518 }
519 if c.canonical_options.is_empty()
520 || c.canonical_options.len() > 64
521 || c.canonical_options.iter().any(|o| o.len() > 4096)
522 {
523 return Err(Reject::Bounds);
524 }
525 if c.canonical_options.windows(2).any(|w| w[0] >= w[1])
526 || !c.canonical_options.contains(&c.default_options)
527 {
528 return Err(Reject::Canonical);
529 }
530 }
531 if v.providers.is_empty() || v.providers.len() > 2 {
532 return Err(Reject::Bounds);
533 }
534 for (i, p) in v.providers.iter().enumerate() {
535 if i > 0 && v.providers[i - 1].provider >= p.provider {
536 return Err(Reject::Canonical);
537 }
538 let mode = match p.provider.as_str() {
539 "github" => 1,
540 "public-git" => 2,
541 _ => return Err(Reject::SourceSelection),
542 };
543 if p.source_modes != [mode] {
544 return Err(Reject::SourceSelection);
545 }
546 }
547 if let Some(default) = &v.default_converter_version
548 && !v.converters.iter().any(|c| &c.converter_version == default)
549 {
550 return Err(Reject::Canonical);
551 }
552 let l = v.limits.as_ref().ok_or(Reject::Canonical)?;
553 if l.max_branches == 0
554 || l.max_branches as usize > MAX_BRANCHES
555 || l.max_operations == 0
556 || l.max_operations as usize > MAX_BRANCHES
557 || l.max_result_bytes == 0
558 {
559 return Err(Reject::Bounds);
560 }
561 Ok(())
562}
563
564pub fn prepare_scope(
567 proposed: &ImportPermissionScopeV1,
568 configuration: &GetImportConfigurationResponse,
569 current_destination_version: &[u8],
570) -> Result<ImportPermissionScopeV1, Reject> {
571 use ImportPreparationRefusalReason as Reason;
572 validate_import_configuration(configuration)?;
573 width(current_destination_version, 32)?;
574 if !proposed.destination_version.is_empty() && proposed.destination_version.len() != 32 {
575 return Err(Reject::PreparationRefused(Reason::InvalidScope));
576 }
577 if !proposed.destination_version.is_empty()
578 && proposed.destination_version != current_destination_version
579 {
580 return Err(Reject::PreparationRefused(Reason::DestinationConflict));
581 }
582 let mut selected = proposed.clone();
583 selected.destination_version = current_destination_version.to_vec();
584 validate_scope(&selected).map_err(|_| Reject::PreparationRefused(Reason::InvalidScope))?;
585 validate_provider_support(&selected.provider, configuration)
586 .map_err(|_| Reject::PreparationRefused(Reason::InvalidScope))?;
587 let converter = configuration
588 .converters
589 .iter()
590 .find(|c| c.converter_version == selected.converter_version)
591 .ok_or(Reject::PreparationRefused(Reason::UnsupportedConverter))?;
592 let supported = converter
593 .canonical_options
594 .iter()
595 .try_fold(false, |found, o| {
596 Ok::<_, Reject>(
597 found
598 || conversion_options_digest(&converter.converter_version, o)?
599 == selected.options_digest,
600 )
601 })?;
602 if !supported {
603 return Err(Reject::PreparationRefused(Reason::UnsupportedOptions));
604 }
605 let limits = configuration.limits.as_ref().ok_or(Reject::Canonical)?;
606 if selected.branches.len() > limits.max_branches as usize
607 || selected.max_operations > limits.max_operations
608 || selected.max_result_bytes > limits.max_result_bytes
609 {
610 return Err(Reject::PreparationRefused(Reason::BudgetExceeded));
611 }
612 Ok(selected)
613}
614
615pub struct ImportSpoolReservation<'a> {
620 pub spool_uuid: &'a [u8],
621 pub logical_job_id: &'a [u8],
622 pub branches: &'a [ImportBranchLimitV1],
623}
624
625pub fn check_import_spool_reservations(
629 spool_uuid: &[u8],
630 logical_job_id: &[u8],
631 scope: &ImportPermissionScopeV1,
632 reservations: &[ImportSpoolReservation<'_>],
633) -> Result<(), Reject> {
634 initial_operation_id(spool_uuid, false)?;
635 initial_operation_id(logical_job_id, false)?;
636 validate_scope(scope)?;
637 for reservation in reservations {
638 initial_operation_id(reservation.spool_uuid, false)?;
639 initial_operation_id(reservation.logical_job_id, false)?;
640 if reservation.spool_uuid != spool_uuid {
641 continue;
642 }
643 let same_job = reservation.logical_job_id == logical_job_id;
644 for selected in &scope.branches {
645 let conflict = if same_job {
646 !reservation.branches.contains(selected)
647 } else {
648 reservation.branches.iter().any(|held| {
649 held.ref_name == selected.ref_name
650 || held.target_thread_id == selected.target_thread_id
651 || held.genesis_digest == selected.genesis_digest
652 })
653 };
654 if conflict {
655 return Err(Reject::PreparationRefused(
656 ImportPreparationRefusalReason::DestinationConflict,
657 ));
658 }
659 }
660 }
661 Ok(())
662}
663
664pub fn import_commit_conflict_failure(
667 rejection: &Reject,
668) -> Option<crate::heddle::api::common::CallFailure> {
669 use crate::heddle::api::common::{CallFailure, CallFailureCode, ErrorDetail, ErrorReason};
670 let (code, reason) = match rejection {
671 Reject::PreparationRefused(ImportPreparationRefusalReason::DestinationConflict) => (
672 CallFailureCode::AlreadyExists,
673 ErrorReason::ImportDestinationConflict,
674 ),
675 Reject::StaleContext => (CallFailureCode::Aborted, ErrorReason::VersionConflict),
676 _ => return None,
677 };
678 Some(CallFailure {
679 code: code as i32,
680 message: String::new(),
681 error: Some(ErrorDetail {
682 reason: reason as i32,
683 ..Default::default()
684 }),
685 })
686}
687
688pub fn check_import_destination_version(
691 scope: &ImportPermissionScopeV1,
692 current_destination_version: &[u8],
693) -> Result<(), Reject> {
694 width(&scope.destination_version, 32)?;
695 width(current_destination_version, 32)?;
696 if scope.destination_version != current_destination_version {
697 return Err(Reject::StaleContext);
698 }
699 Ok(())
700}
701
702pub fn validate_preparation_response(
704 request: &PrepareImportJobRequest,
705 response: &PrepareImportJobResponse,
706) -> Result<(), Reject> {
707 if let Some(refusal) = &response.refusal {
708 let reason = ImportPreparationRefusalReason::try_from(refusal.reason)
709 .map_err(|_| Reject::Version)?;
710 if reason == ImportPreparationRefusalReason::Unspecified
711 || refusal.field.len() > 256
712 || !refusal.field.is_ascii()
713 || response.proposal.is_some()
714 || response.reservation_expires_at_unix_seconds != 0
715 || response.prepared_at_unix_seconds != 0
716 || response.max_validity_duration_seconds != 0
717 || response.clock_skew_allowance_seconds != 0
718 {
719 return Err(Reject::Canonical);
720 }
721 return Err(Reject::PreparationRefused(reason));
722 }
723 let p = response.proposal.as_ref().ok_or(Reject::Canonical)?;
724 let returned = p.scope.as_ref().ok_or(Reject::Canonical)?;
725 let mut requested = request.proposed_scope.clone().ok_or(Reject::Canonical)?;
726 if requested.destination_version.is_empty() {
727 requested.destination_version = returned.destination_version.clone();
728 }
729 if canonical(&requested)? != canonical(returned)?
730 || p.identity != request.identity
731 || p.retry_lineage_id != request.retry_lineage_id
732 {
733 return Err(Reject::PreparedFields);
734 }
735 validate_scope(returned)?;
736 Ok(())
737}
738
739pub fn validate_commit_request_bounds(request: &CommitImportJobRequest) -> Result<(), Reject> {
742 use prost::Message;
743 if request.encoded_len() > MAX_COMMIT_REQUEST_BYTES
744 || request
745 .proof
746 .as_ref()
747 .ok_or(Reject::Canonical)?
748 .encoded_len()
749 > MAX_BUNDLE_BYTES
750 {
751 return Err(Reject::Bounds);
752 }
753 Ok(())
754}
755
756pub fn validate_commit_request(
760 request: &CommitImportJobRequest,
761 resolved_provider: &str,
762 current_source: &ProviderRepository,
763 configuration: &GetImportConfigurationResponse,
764) -> Result<(), Reject> {
765 use prost::Message;
766 validate_commit_request_bounds(request)?;
767 let source = request.source.as_ref().ok_or(Reject::SourceSelection)?;
768 let proof = request.proof.as_ref().ok_or(Reject::Canonical)?;
769 if request.client_operation_id.is_empty()
770 || request.client_operation_id.len() > 128
771 || request.destination.as_ref().is_none_or(|d| d.id.is_empty())
772 {
773 return Err(Reject::Canonical);
774 }
775 if request.initial_base_state.len() > 4096 || proof.encoded_len() > MAX_BUNDLE_BYTES {
776 return Err(Reject::Bounds);
777 }
778 if proof.format_version != 1
779 || proof.delegations.len() != 1
780 || !proof.operations.is_empty()
781 || proof.terminal_manifest.is_some()
782 || !proof.manifests.is_empty()
783 {
784 return Err(Reject::Canonical);
785 }
786 let d = proof.delegations[0]
787 .body
788 .as_ref()
789 .ok_or(Reject::Canonical)?;
790 let scope = d.scope.as_ref().ok_or(Reject::Canonical)?;
791 validate_scope(scope)?;
792 let id = d.identity.as_ref().ok_or(Reject::Canonical)?;
793 identity(id)?;
794 let uuid = hex::encode(&id.spool_uuid);
795 let destination_id = format!(
796 "{}-{}-{}-{}-{}",
797 &uuid[..8],
798 &uuid[8..12],
799 &uuid[12..16],
800 &uuid[16..20],
801 &uuid[20..]
802 );
803 if request
804 .destination
805 .as_ref()
806 .is_none_or(|s| s.id != destination_id)
807 {
808 return Err(Reject::Scope);
809 }
810 if source.clone_url != scope.source_url
811 || resolved_provider != scope.provider
812 || source.provider_repository_id.len() > 4096
813 || source.name.len() > 4096
814 {
815 return Err(Reject::SourceSelection);
816 }
817 if source.connection != current_source.connection
818 || (source.provider_repository_id != current_source.provider_repository_id
819 && (source.connection.is_some() || !source.provider_repository_id.is_empty()))
820 || source.clone_url != current_source.clone_url
821 || source.installation_id != current_source.installation_id
822 || source.private != current_source.private
823 {
824 return Err(Reject::SourceSelection);
825 }
826 let provider = resolve_import_provider(
827 current_source,
828 current_source
829 .connection
830 .as_ref()
831 .map(|_| resolved_provider),
832 )?;
833 if provider != resolved_provider {
834 return Err(Reject::SourceSelection);
835 }
836 validate_import_configuration(configuration)?;
837 validate_provider_support(provider, configuration)?;
838 validate_repository_hash_algorithm(current_source, true)?;
839 if source.hash_algorithm != current_source.hash_algorithm {
840 return Err(Reject::SourceSelection);
841 }
842 for b in &scope.branches {
845 if b.hash_algorithm != current_source.hash_algorithm {
846 return Err(Reject::SourceSelection);
847 }
848 if b.ref_mode == 2
849 && current_source
850 .refs
851 .iter()
852 .any(|r| r.name == b.ref_name && !r.head_oid.is_empty())
853 {
854 return Err(Reject::RefPinning);
855 }
856 }
857 prepare_scope(scope, configuration, &scope.destination_version)?;
859 if proof.original_geneses.len() != scope.branches.len()
860 || proof.creator_authority_envelopes.len() != scope.branches.len()
861 || proof.genesis_authorities.len() != scope.branches.len()
862 || d.branch_manifest.len() != scope.branches.len()
863 {
864 return Err(Reject::GenesisBinding);
865 }
866 for (i, b) in scope.branches.iter().enumerate() {
869 let original = &proof.original_geneses[i];
870 let binding = &proof.genesis_authorities[i];
871 let g = binding.body.as_ref().ok_or(Reject::GenesisBinding)?;
872 let m = &d.branch_manifest[i];
873 if native_id(original) != b.genesis_digest
874 || g.genesis_digest != b.genesis_digest
875 || m.limit.as_ref() != Some(b)
876 || m.genesis_authority_digest != signed_genesis_digest(binding)?
877 || g.creator_authority_envelope_digest != hash(&[&proof.creator_authority_envelopes[i]])
878 || proof.creator_authority_envelopes[i].is_empty()
879 || proof.creator_authority_envelopes[i].len() > MAX_RECORD_BYTES
880 || !original.signatures.iter().any(|s| {
881 s.public_key == g.creator_public_key && s.signature == g.original_creator_signature
882 })
883 {
884 return Err(Reject::GenesisBinding);
885 }
886 verify_native(original, "heddle-thread-genesis-v1")?;
887 }
888 Ok(())
889}
890
891pub fn validate_import_source(_: &ImportSourceRequest) -> Result<(), Reject> {
893 Err(Reject::ImportSourceRequiresCommit)
894}
895
896pub fn verify_commit_submission(
899 request: &CommitImportJobRequest,
900 prepared: &PrepareImportJobResponse,
901 resolved_provider: &str,
902 current_source: &ProviderRepository,
903 configuration: &GetImportConfigurationResponse,
904 expected: &ImportOwnerExpectation<'_>,
905) -> Result<VerifiedImportDelegation, Reject> {
906 validate_commit_request(request, resolved_provider, current_source, configuration)?;
907 let proof = request.proof.as_ref().ok_or(Reject::Canonical)?;
908 let member = proof.member_permission.as_ref();
909 verify_prepared_delegation(
910 prepared,
911 &proof.delegations[0],
912 member,
913 &proof.genesis_authorities,
914 expected,
915 )
916}
917
918pub fn validate_commit_response(
919 request: &CommitImportJobRequest,
920 response: &MutationResponse,
921) -> Result<(), Reject> {
922 let receipt = response.receipt.as_ref().ok_or(Reject::PendingOperation)?;
923 let Some(mutation_receipt::Outcome::PendingOperation(operation)) = &receipt.outcome else {
924 return Err(Reject::PendingOperation);
925 };
926 if request.client_operation_id.is_empty()
927 || receipt.client_operation_id != request.client_operation_id
928 || request.destination.is_none()
929 || operation.spool != request.destination
930 || operation.id
931 != initial_operation_id(
932 &request
933 .proof
934 .as_ref()
935 .ok_or(Reject::PendingOperation)?
936 .delegations
937 .first()
938 .and_then(|d| d.body.as_ref())
939 .ok_or(Reject::PendingOperation)?
940 .retry_lineage_id,
941 false,
942 )?
943 {
944 return Err(Reject::PendingOperation);
945 }
946 Ok(())
947}
948
949pub fn check_commit_replay(
953 request: &CommitImportJobRequest,
954 stored: &CommitImportJobRequest,
955 response: &MutationResponse,
956) -> Result<(), Reject> {
957 if request != stored {
958 return Err(Reject::OperationIdReused);
959 }
960 validate_commit_response(request, response)
961}
962pub fn validate_scope(value: &ImportPermissionScopeV1) -> Result<(), Reject> {
963 canonical_https(&value.source_url, false)?;
964 if value.provider.is_empty()
965 || value.provider.len() > 64
966 || !value
967 .provider
968 .bytes()
969 .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-')
970 || value.converter_version.is_empty()
971 || value.converter_version.len() > 128
972 || !value.converter_version.is_ascii()
973 {
974 return Err(Reject::Canonical);
975 }
976 width(&value.destination_version, 32)?;
977 width(&value.options_digest, 32)?;
978 if value.branches.is_empty()
979 || value.branches.len() > MAX_BRANCHES
980 || value.max_operations == 0
981 || value.max_operations as usize > MAX_BRANCHES
982 || value.max_result_bytes == 0
983 {
984 return Err(Reject::Bounds);
985 }
986 if (value.max_operations as usize) < value.branches.len() {
987 return Err(Reject::Scope);
988 }
989 for (i, b) in value.branches.iter().enumerate() {
990 branch(b)?;
991 if value.branches[..i].iter().any(|other| {
992 other.target_thread_id == b.target_thread_id || other.genesis_digest == b.genesis_digest
993 }) {
994 return Err(Reject::Scope);
995 }
996 if i > 0 && value.branches[i - 1].ref_name.as_bytes() >= b.ref_name.as_bytes() {
997 return Err(Reject::Canonical);
998 }
999 }
1000 Ok(())
1001}
1002fn scope_subset(child: &ImportPermissionScopeV1, parent: &ImportPermissionScopeV1) -> bool {
1003 child.provider == parent.provider
1004 && child.source_url == parent.source_url
1005 && child.destination_version == parent.destination_version
1006 && child.options_digest == parent.options_digest
1007 && child.converter_version == parent.converter_version
1008 && child.max_operations <= parent.max_operations
1009 && child.max_result_bytes <= parent.max_result_bytes
1010 && child.branches.iter().all(|c| parent.branches.contains(c))
1011}
1012
1013pub struct ImportOwnerExpectation<'a> {
1017 pub identity: &'a ImportIdentityV1,
1018 pub owner_public_key: &'a [u8],
1019 pub owner_chain_digest: &'a [u8],
1020 pub authority_expires_at_seconds: i64,
1023 pub now_unix_seconds: i64,
1024 pub forbidden_job_keys: &'a [Vec<u8>], pub known_job_associations: &'a [(Vec<u8>, Vec<u8>)], }
1027#[derive(Clone, Copy)]
1030pub struct ImportBundleOwnerExpectation<'a> {
1031 pub identity: &'a ImportIdentityV1,
1032 pub owner_public_key: &'a [u8],
1033 pub owner_chain_digest: &'a [u8],
1034 pub authority_expires_at_seconds: i64,
1035 pub effective_from_unix_seconds: i64,
1037 pub effective_until_unix_seconds: Option<i64>,
1038 pub forbidden_job_keys: &'a [Vec<u8>],
1039 pub forbidden_landing_keys: &'a [Vec<u8>],
1040 pub known_job_associations: &'a [(Vec<u8>, Vec<u8>)],
1041}
1042impl<'a> ImportBundleOwnerExpectation<'a> {
1043 fn at(self, time: i64, associations: &'a [(Vec<u8>, Vec<u8>)]) -> ImportOwnerExpectation<'a> {
1044 ImportOwnerExpectation {
1045 identity: self.identity,
1046 owner_public_key: self.owner_public_key,
1047 owner_chain_digest: self.owner_chain_digest,
1048 authority_expires_at_seconds: self.authority_expires_at_seconds,
1049 now_unix_seconds: time,
1050 forbidden_job_keys: self.forbidden_job_keys,
1051 known_job_associations: associations,
1052 }
1053 }
1054}
1055pub fn effective_owner_authority_expiry(
1058 deferred_human: bool,
1059 claimable_until: i64,
1060) -> Result<i64, Reject> {
1061 if deferred_human {
1062 if claimable_until <= 0 {
1063 return Err(Reject::Scope);
1064 }
1065 Ok(claimable_until)
1066 } else {
1067 Ok(i64::MAX)
1068 }
1069}
1070pub fn verify_member_permission(
1071 signed: &SignedImportMemberPermissionV1,
1072 expected: &ImportOwnerExpectation<'_>,
1073) -> Result<(), Reject> {
1074 verify_member_permission_inner(signed, expected, true)
1075}
1076fn verify_member_permission_inner(
1077 signed: &SignedImportMemberPermissionV1,
1078 expected: &ImportOwnerExpectation<'_>,
1079 current: bool,
1080) -> Result<(), Reject> {
1081 let p = signed.body.as_ref().ok_or(Reject::ImportPermission)?;
1082 if p.format_version != 1 || p.purpose != 1 {
1083 return Err(Reject::ImportPermission);
1084 }
1085 identity(p.identity.as_ref().ok_or(Reject::Canonical)?)?;
1086 if p.identity.as_ref() != Some(expected.identity)
1087 || p.owner_chain_digest != expected.owner_chain_digest
1088 {
1089 return Err(Reject::Root);
1090 }
1091 width(&p.logical_job_id, 16)?;
1092 width(&p.retry_lineage_id, 16)?;
1093 width(&p.subject_public_key, 32)?;
1094 width(&p.cancellation_id, 32)?;
1095 width(&p.nonce, 32)?;
1096 width(&p.owner_chain_digest, 32)?;
1097 validate_scope(p.scope.as_ref().ok_or(Reject::Canonical)?)?;
1098 validity(
1099 p.not_before_unix_seconds,
1100 p.expires_at_unix_seconds,
1101 expected.now_unix_seconds,
1102 current,
1103 )?;
1104 if p.expires_at_unix_seconds > expected.authority_expires_at_seconds {
1105 return Err(Reject::Scope);
1106 }
1107 verify_authorization_signature(
1108 expected.owner_public_key,
1109 PERMISSION_DOMAIN,
1110 p,
1111 signed.owner_signature.as_ref().ok_or(Reject::Signature)?,
1112 )
1113}
1114
1115#[derive(Debug, Clone, PartialEq)]
1116pub struct VerifiedImportDelegation {
1117 body: ImportJobDelegationV1,
1118 digest: Vec<u8>,
1119 member: Option<SignedImportMemberPermissionV1>,
1120}
1121impl VerifiedImportDelegation {
1122 pub fn body(&self) -> &ImportJobDelegationV1 {
1123 &self.body
1124 }
1125 pub fn digest(&self) -> &[u8] {
1126 &self.digest
1127 }
1128}
1129pub fn verify_delegation(
1130 signed: &SignedImportJobDelegationV1,
1131 member: Option<&SignedImportMemberPermissionV1>,
1132 expected: &ImportOwnerExpectation<'_>,
1133) -> Result<VerifiedImportDelegation, Reject> {
1134 verify_delegation_inner(signed, member, expected, true)
1135}
1136fn verify_delegation_inner(
1137 signed: &SignedImportJobDelegationV1,
1138 member: Option<&SignedImportMemberPermissionV1>,
1139 expected: &ImportOwnerExpectation<'_>,
1140 current: bool,
1141) -> Result<VerifiedImportDelegation, Reject> {
1142 let d = signed.body.as_ref().ok_or(Reject::Canonical)?;
1143 if d.format_version != 1 || d.purpose != 1 {
1144 return Err(Reject::Version);
1145 }
1146 identity(d.identity.as_ref().ok_or(Reject::Canonical)?)?;
1147 if d.identity.as_ref() != Some(expected.identity)
1148 || d.owner_chain_digest != expected.owner_chain_digest
1149 {
1150 return Err(Reject::Root);
1151 }
1152 for v in [&d.delegation_id, &d.logical_job_id, &d.retry_lineage_id] {
1153 width(v, 16)?;
1154 if v.iter().all(|b| *b == 0) {
1155 return Err(Reject::Canonical);
1156 }
1157 }
1158 for v in [
1159 &d.job_public_key,
1160 &d.job_key_id,
1161 &d.delegating_public_key,
1162 &d.parent_permission_digest,
1163 &d.owner_chain_digest,
1164 &d.cancellation_id,
1165 ] {
1166 width(v, 32)?;
1167 }
1168 if d.job_key_id != key_id(&d.job_public_key) {
1169 return Err(Reject::Canonical);
1170 }
1171 if d.job_public_key == d.delegating_public_key
1172 || d.job_public_key == expected.owner_public_key
1173 || expected.forbidden_job_keys.contains(&d.job_public_key)
1174 {
1175 return Err(Reject::KeyRole);
1176 }
1177 if expected
1178 .known_job_associations
1179 .iter()
1180 .any(|(k, j)| k == &d.job_public_key && j != &d.logical_job_id)
1181 {
1182 return Err(Reject::Scope);
1183 }
1184 let scope = d.scope.as_ref().ok_or(Reject::Canonical)?;
1185 validate_scope(scope)?;
1186 if d.branch_manifest.len() != scope.branches.len() {
1187 return Err(Reject::Scope);
1188 }
1189 for (m, b) in d.branch_manifest.iter().zip(&scope.branches) {
1190 width(&m.genesis_authority_digest, 32)?;
1191 if m.limit.as_ref() != Some(b) {
1192 return Err(Reject::Scope);
1193 }
1194 }
1195 if i128::from(d.expires_at_unix_seconds) - i128::from(d.not_before_unix_seconds)
1196 > i128::from(MAX_DELEGATION_WINDOW_SECONDS)
1197 {
1198 return Err(Reject::ValidityBounds);
1199 }
1200 validity(
1201 d.not_before_unix_seconds,
1202 d.expires_at_unix_seconds,
1203 expected.now_unix_seconds,
1204 current,
1205 )?;
1206 if d.expires_at_unix_seconds > expected.authority_expires_at_seconds {
1207 return Err(Reject::Scope);
1208 }
1209 if d.delegating_public_key == expected.owner_public_key {
1210 if member.is_some() || d.parent_permission_digest != vec![0; 32] {
1211 return Err(Reject::ImportPermission);
1212 }
1213 } else {
1214 let member = member.ok_or(Reject::ImportPermission)?;
1215 verify_member_permission_inner(member, expected, current)?;
1216 let p = member.body.as_ref().ok_or(Reject::ImportPermission)?;
1217 if d.parent_permission_digest != signed_permission_digest(member)?
1218 || d.delegating_public_key != p.subject_public_key
1219 || d.logical_job_id != p.logical_job_id
1220 || d.retry_lineage_id != p.retry_lineage_id
1221 || d.not_before_unix_seconds < p.not_before_unix_seconds
1222 || d.expires_at_unix_seconds > p.expires_at_unix_seconds
1223 || !scope_subset(scope, p.scope.as_ref().ok_or(Reject::Canonical)?)
1224 {
1225 return Err(Reject::Scope);
1226 }
1227 }
1228 verify_authorization_signature(
1229 &d.delegating_public_key,
1230 DELEGATION_DOMAIN,
1231 d,
1232 signed
1233 .delegating_signature
1234 .as_ref()
1235 .ok_or(Reject::Signature)?,
1236 )?;
1237 Ok(VerifiedImportDelegation {
1238 body: d.clone(),
1239 digest: signed_delegation_digest(signed)?,
1240 member: member.cloned(),
1241 })
1242}
1243pub fn delegation_preparation(d: &ImportJobDelegationV1) -> ImportJobPreparationV1 {
1246 ImportJobPreparationV1 {
1247 format_version: d.format_version,
1248 identity: d.identity.clone(),
1249 delegation_id: d.delegation_id.clone(),
1250 logical_job_id: d.logical_job_id.clone(),
1251 retry_lineage_id: d.retry_lineage_id.clone(),
1252 job_public_key: d.job_public_key.clone(),
1253 job_key_id: d.job_key_id.clone(),
1254 owner_chain_digest: d.owner_chain_digest.clone(),
1255 purpose: d.purpose,
1256 scope: d.scope.clone(),
1257 cancellation_id: d.cancellation_id.clone(),
1258 }
1259}
1260
1261pub fn verify_prepared_delegation(
1265 prepared: &PrepareImportJobResponse,
1266 signed: &SignedImportJobDelegationV1,
1267 member: Option<&SignedImportMemberPermissionV1>,
1268 geneses: &[SignedImportGenesisAuthorityV1],
1269 expected: &ImportOwnerExpectation<'_>,
1270) -> Result<VerifiedImportDelegation, Reject> {
1271 verify_prepared_inner(prepared, signed, member, geneses, expected, false)
1272}
1273pub fn preflight_prepared_delegation(
1275 prepared: &PrepareImportJobResponse,
1276 signed: &SignedImportJobDelegationV1,
1277 member: Option<&SignedImportMemberPermissionV1>,
1278 geneses: &[SignedImportGenesisAuthorityV1],
1279 expected: &ImportOwnerExpectation<'_>,
1280) -> Result<(), Reject> {
1281 verify_prepared_inner(prepared, signed, member, geneses, expected, true).map(|_| ())
1282}
1283fn verify_prepared_inner(
1284 prepared: &PrepareImportJobResponse,
1285 signed: &SignedImportJobDelegationV1,
1286 member: Option<&SignedImportMemberPermissionV1>,
1287 geneses: &[SignedImportGenesisAuthorityV1],
1288 expected: &ImportOwnerExpectation<'_>,
1289 browser: bool,
1290) -> Result<VerifiedImportDelegation, Reject> {
1291 let proposal = prepared.proposal.as_ref().ok_or(Reject::Canonical)?;
1292 let d = signed.body.as_ref().ok_or(Reject::Canonical)?;
1293 if canonical(proposal)? != canonical(&delegation_preparation(d))? {
1294 return Err(Reject::PreparedFields);
1295 }
1296 let scope = proposal.scope.as_ref().ok_or(Reject::Canonical)?;
1297 if d.branch_manifest.len() != scope.branches.len() {
1298 return Err(Reject::PreparedFields);
1299 }
1300 for (m, b) in d.branch_manifest.iter().zip(&scope.branches) {
1301 let limit = m.limit.as_ref().ok_or(Reject::PreparedFields)?;
1302 if canonical(limit)? != canonical(b)? {
1303 return Err(Reject::PreparedFields);
1304 }
1305 }
1306 let now = i128::from(expected.now_unix_seconds);
1309 let start = i128::from(d.not_before_unix_seconds);
1310 let end = i128::from(d.expires_at_unix_seconds);
1311 let at = i128::from(prepared.prepared_at_unix_seconds);
1312 let skew = i128::from(prepared.clock_skew_allowance_seconds);
1313 if at < 0
1314 || now < 0
1315 || if browser { now + skew < at } else { now < at }
1316 || i128::from(prepared.reservation_expires_at_unix_seconds) != at + 3600
1317 || now >= i128::from(prepared.reservation_expires_at_unix_seconds)
1318 {
1319 return Err(Reject::Expired);
1320 }
1321 if prepared.max_validity_duration_seconds == 0
1322 || prepared.max_validity_duration_seconds > MAX_DELEGATION_WINDOW_SECONDS
1323 || start < 0
1324 || start < at - skew
1325 || start > now + skew
1326 || end <= start
1327 || end <= now
1328 || end - start > i128::from(prepared.max_validity_duration_seconds)
1329 {
1330 return Err(Reject::ValidityBounds);
1331 }
1332 if let Some(parent) = member {
1333 if browser {
1334 let p = parent.body.as_ref().ok_or(Reject::ImportPermission)?;
1335 verify_member_permission_inner(parent, expected, false)?;
1336 if i128::from(p.not_before_unix_seconds) > now + skew
1337 || i128::from(p.expires_at_unix_seconds) <= now
1338 {
1339 return Err(Reject::Expired);
1340 }
1341 } else {
1342 verify_member_permission(parent, expected)?;
1343 }
1344 }
1345 let at_start = ImportOwnerExpectation {
1349 now_unix_seconds: expected.now_unix_seconds.max(d.not_before_unix_seconds),
1350 ..*expected
1351 };
1352 let verified = verify_delegation_inner(
1353 signed,
1354 member,
1355 if browser { expected } else { &at_start },
1356 !browser,
1357 )?;
1358 if geneses.len() != d.branch_manifest.len() {
1359 return Err(Reject::GenesisBinding);
1360 }
1361 for m in &d.branch_manifest {
1362 let branch = m.limit.as_ref().ok_or(Reject::Canonical)?;
1363 let g = geneses
1364 .iter()
1365 .find(|g| signed_genesis_digest(g).is_ok_and(|h| h == m.genesis_authority_digest))
1366 .ok_or(Reject::GenesisBinding)?;
1367 let body = g.body.as_ref().ok_or(Reject::GenesisBinding)?;
1368 if body.genesis_digest != branch.genesis_digest {
1369 return Err(Reject::GenesisBinding);
1370 }
1371 {
1372 verify_genesis_authority(
1373 g,
1374 &verified,
1375 &branch.genesis_digest,
1376 &body.original_creator_signature,
1377 &body.creator_authority_envelope_digest,
1378 )?;
1379 }
1380 }
1381 Ok(verified)
1382}
1383
1384pub fn verify_genesis_authority(
1385 signed: &SignedImportGenesisAuthorityV1,
1386 delegation: &VerifiedImportDelegation,
1387 original_genesis_digest: &[u8],
1388 original_signature: &[u8],
1389 envelope_digest: &[u8],
1390) -> Result<(), Reject> {
1391 let g = signed.body.as_ref().ok_or(Reject::Canonical)?;
1392 let d = &delegation.body;
1393 if g.format_version != 1 {
1394 return Err(Reject::Version);
1395 }
1396 width(&g.original_creator_signature, 64)?;
1397 for v in [
1398 &g.genesis_digest,
1399 &g.creator_public_key,
1400 &g.creator_authority_envelope_digest,
1401 &g.parent_permission_digest,
1402 &g.owner_chain_digest,
1403 ] {
1404 width(v, 32)?;
1405 }
1406 if g.identity != d.identity
1407 || g.creator_public_key != d.delegating_public_key
1408 || g.parent_permission_digest != d.parent_permission_digest
1409 || g.owner_chain_digest != d.owner_chain_digest
1410 || g.genesis_digest != original_genesis_digest
1411 || g.original_creator_signature != original_signature
1412 || g.creator_authority_envelope_digest != envelope_digest
1413 || !d.branch_manifest.iter().any(|m| {
1414 m.limit
1415 .as_ref()
1416 .is_some_and(|b| b.genesis_digest == g.genesis_digest)
1417 && signed_genesis_digest(signed).is_ok_and(|h| h == m.genesis_authority_digest)
1418 })
1419 {
1420 return Err(Reject::Scope);
1421 }
1422 verify_authorization_signature(
1423 &g.creator_public_key,
1424 GENESIS_DOMAIN,
1425 g,
1426 signed.creator_signature.as_ref().ok_or(Reject::Signature)?,
1427 )
1428}
1429pub fn verify_operation(
1432 signed: &SignedDelegatedImportOperationV1,
1433 delegation: &VerifiedImportDelegation,
1434) -> Result<(), Reject> {
1435 let o = signed.body.as_ref().ok_or(Reject::Canonical)?;
1436 let d = &delegation.body;
1437 if o.format_version != 1 {
1438 return Err(Reject::Version);
1439 }
1440 let id = d.identity.as_ref().ok_or(Reject::Canonical)?;
1441 width(&o.physical_operation_id, 16)?;
1442 for v in [
1443 &o.spool_genesis_digest,
1444 &o.delegation_digest,
1445 &o.genesis_digest,
1446 &o.target_thread_id,
1447 &o.expected_frontier_digest,
1448 &o.resulting_frontier_digest,
1449 &o.resulting_content_digest,
1450 &o.options_digest,
1451 ] {
1452 width(v, 32)?;
1453 }
1454 width(&o.spool_uuid, 16)?;
1455 width(&o.logical_job_id, 16)?;
1456 width(&o.retry_lineage_id, 16)?;
1457 let scope = d.scope.as_ref().ok_or(Reject::Canonical)?;
1458 let b = scope
1459 .branches
1460 .iter()
1461 .find(|b| b.ref_name == o.ref_name && b.slot_id == o.slot_id)
1462 .ok_or(Reject::Scope)?;
1463 let oid_len = match o.hash_algorithm {
1464 1 => 20,
1465 2 => 32,
1466 _ => return Err(Reject::Version),
1467 };
1468 width(&o.observed_commit_oid, oid_len)?;
1469 if o.spool_uuid != id.spool_uuid
1470 || o.spool_genesis_digest != id.spool_genesis_digest
1471 || o.logical_job_id != d.logical_job_id
1472 || o.retry_lineage_id != d.retry_lineage_id
1473 || o.delegation_digest != delegation.digest
1474 || o.hash_algorithm != b.hash_algorithm
1475 || (b.ref_mode == 1 && o.observed_commit_oid != b.pinned_commit_oid)
1476 || o.genesis_digest != b.genesis_digest
1477 || o.target_thread_id != b.target_thread_id
1478 || o.expected_frontier_digest != b.expected_frontier_digest
1479 || o.result_bytes > scope.max_result_bytes
1480 || o.result_bytes == 0
1481 || o.options_digest != scope.options_digest
1482 || o.converter_version != scope.converter_version
1483 {
1484 return Err(Reject::Scope);
1485 }
1486 verify_authorization_signature(
1487 &d.job_public_key,
1488 OPERATION_DOMAIN,
1489 o,
1490 signed.job_signature.as_ref().ok_or(Reject::Signature)?,
1491 )
1492}
1493pub fn verify_new_operation(
1494 signed: &SignedDelegatedImportOperationV1,
1495 delegation: &VerifiedImportDelegation,
1496 now_seconds: i64,
1497 committed_before: &ImportResultManifestV1,
1498) -> Result<(), Reject> {
1499 interval(
1500 delegation.body.not_before_unix_seconds,
1501 delegation.body.expires_at_unix_seconds,
1502 now_seconds,
1503 )?;
1504 check_import_publication_budget(signed, delegation, committed_before)
1505}
1506pub fn validate_manifest(m: &ImportResultManifestV1) -> Result<(), Reject> {
1507 if m.format_version != 1 {
1508 return Err(Reject::Version);
1509 }
1510 width(&m.logical_job_id, 16)?;
1511 width(&m.retry_lineage_id, 16)?;
1512 if m.slots.len() > MAX_BRANCHES {
1513 return Err(Reject::Bounds);
1514 }
1515 for (i, s) in m.slots.iter().enumerate() {
1516 width(&s.signed_operation_digest, 32)?;
1517 width(&s.resulting_frontier_digest, 32)?;
1518 branch_ref(&s.ref_name)?;
1519 if s.result_bytes == 0 {
1520 return Err(Reject::Bounds);
1521 }
1522 if i > 0 && (&m.slots[i - 1].ref_name, m.slots[i - 1].slot_id) >= (&s.ref_name, s.slot_id) {
1523 return Err(Reject::Canonical);
1524 }
1525 }
1526 Ok(())
1527}
1528
1529pub fn check_slot_replay(
1532 committed: &ImportResultManifestV1,
1533 signed: &SignedDelegatedImportOperationV1,
1534) -> Result<bool, Reject> {
1535 validate_manifest(committed)?;
1536 let o = signed.body.as_ref().ok_or(Reject::Canonical)?;
1537 if committed.logical_job_id != o.logical_job_id
1538 || committed.retry_lineage_id != o.retry_lineage_id
1539 {
1540 return Err(Reject::Scope);
1541 }
1542 match committed
1543 .slots
1544 .iter()
1545 .find(|s| s.ref_name == o.ref_name && s.slot_id == o.slot_id)
1546 {
1547 Some(s)
1548 if s.signed_operation_digest == signed_operation_digest(signed)?
1549 && s.resulting_frontier_digest == o.resulting_frontier_digest
1550 && s.result_bytes == o.result_bytes =>
1551 {
1552 Ok(true)
1553 }
1554 Some(_) => Err(Reject::SlotConflict),
1555 None => Ok(false),
1556 }
1557}
1558pub fn verify_publication(
1562 operation: &SignedDelegatedImportOperationV1,
1563 delegation: &VerifiedImportDelegation,
1564 manifest: &ImportResultManifestV1,
1565 statement: &crate::heddle::api::common::SignedHostedWitnessStatementV1,
1566 set: &crate::witness_trust::VerifiedWitnessSet,
1567 proof: Option<&crate::heddle::api::common::HostedWitnessHistoryProofV1>,
1568 now_ms: i64,
1569) -> Result<crate::witness_trust::ResolvedWitnessStatement, Reject> {
1570 validate_statement_boundary(statement.body.as_ref().ok_or(Reject::Canonical)?)?;
1571 verify_operation(operation, delegation)?;
1572 if !check_slot_replay(manifest, operation)? {
1573 return Err(Reject::Scope);
1574 }
1575 let mut before = manifest.clone();
1576 let o = operation.body.as_ref().ok_or(Reject::Canonical)?;
1577 before
1578 .slots
1579 .retain(|slot| slot.ref_name != o.ref_name || slot.slot_id != o.slot_id);
1580 check_import_publication_budget(operation, delegation, &before)?;
1581 let d = &delegation.body;
1582 let id = d.identity.as_ref().ok_or(Reject::Canonical)?;
1583 let s = statement.body.as_ref().ok_or(Reject::Canonical)?;
1584 let payload = ImportPublicationWitnessV1 {
1585 format_version: 1,
1586 signed_operation_digest: signed_operation_digest(operation)?,
1587 delegation_digest: delegation.digest.clone(),
1588 logical_job_id: o.logical_job_id.clone(),
1589 retry_lineage_id: o.retry_lineage_id.clone(),
1590 physical_operation_id: o.physical_operation_id.clone(),
1591 ref_name: o.ref_name.clone(),
1592 slot_id: o.slot_id,
1593 hash_algorithm: o.hash_algorithm,
1594 observed_commit_oid: o.observed_commit_oid.clone(),
1595 expected_frontier_digest: o.expected_frontier_digest.clone(),
1596 resulting_frontier_digest: o.resulting_frontier_digest.clone(),
1597 terminal_manifest_digest: manifest_digest(manifest)?,
1598 };
1599 if s.purpose != 3
1600 || s.spool_uuid != id.spool_uuid
1601 || s.spool_genesis_digest != id.spool_genesis_digest
1602 || s.owner_id != id.owner_id
1603 || s.owner_state_hash != id.owner_state_hash
1604 || s.ownership_transfer_sequence != id.ownership_transfer_sequence
1605 || s.authority_digest != delegation.digest
1606 || s.original_signatures_digest
1607 != hash(&[&operation
1608 .job_signature
1609 .as_ref()
1610 .ok_or(Reject::Signature)?
1611 .signature])
1612 || s.canonical_payload != canonical(&payload)?
1613 || s.basis != 1
1614 {
1615 return Err(Reject::Scope);
1616 }
1617 interval(
1618 d.not_before_unix_seconds,
1619 d.expires_at_unix_seconds,
1620 s.observed_at_unix_millis / 1000,
1621 )?;
1622 crate::witness_trust::resolve_statement(set, statement, proof, false, now_ms)
1623}
1624pub fn require_hybrid_peer(
1625 protocol: Option<&crate::heddle::api::common::ProtocolCompatibility>,
1626) -> Result<(), Reject> {
1627 let protocol = protocol.ok_or(Reject::Protocol)?;
1628 if protocol.protocol_version != 2 || protocol.mandatory_features != [1] {
1629 return Err(Reject::Protocol);
1630 }
1631 Ok(())
1632}
1633
1634pub fn check_job_fence(
1637 logical_job_id: &[u8],
1638 active_delegation_digest: &[u8],
1639 expected_epoch: u64,
1640 active: &VerifiedImportDelegation,
1641 durable_epoch: u64,
1642) -> Result<(), Reject> {
1643 if logical_job_id != active.body.logical_job_id {
1644 return Err(Reject::Scope);
1645 }
1646 if expected_epoch != durable_epoch || active_delegation_digest != active.digest {
1647 return Err(Reject::StaleContext);
1648 }
1649 Ok(())
1650}
1651pub fn validate_public_bundle(bundle: &ImportPublicProofBundleV1) -> Result<(), Reject> {
1652 validate_bundle_bounds(bundle)?;
1653 validate_bundle_history(bundle, true)
1654}
1655fn validate_bundle_bounds(bundle: &ImportPublicProofBundleV1) -> Result<(), Reject> {
1656 crate::writer_authority::validate_owner_histories(&bundle.owner_histories)?;
1657 use prost::Message;
1658 if bundle.format_version != 1 {
1659 return Err(Reject::Version);
1660 }
1661 if bundle.encoded_len() > MAX_BUNDLE_BYTES
1662 || bundle.owner_histories.len() > 64
1663 || bundle.ownership_transfers.len() > 64
1664 || bundle.genesis_authorities.len() > MAX_BRANCHES
1665 || bundle.delegations.len() != 1
1666 || bundle.operations.len() > MAX_BRANCHES
1667 || bundle.statements.len() > 1024
1668 || bundle.history_proofs.len() > 1024
1669 || bundle.policies.len() > 256
1670 || bundle.original_geneses.len() > MAX_BRANCHES
1671 || bundle.creator_authority_envelopes.len() > MAX_BRANCHES
1672 || bundle.manifests.len() > MAX_BRANCHES
1673 || bundle.genesis_witnesses.len() > 256
1674 || bundle.authority_witnesses.len() > 256
1675 || bundle.landing_witnesses.len() > 256
1676 {
1677 return Err(Reject::Bounds);
1678 }
1679 Ok(())
1680}
1681
1682pub enum ImportPermissionEvidence<'a> {
1685 Import(&'a SignedImportMemberPermissionV1),
1686 OwnerCapability(&'a SignedOwnerCapability),
1687 OnlineRole(&'a str),
1688}
1689pub fn select_import_permission(
1690 evidence: ImportPermissionEvidence<'_>,
1691) -> Result<&SignedImportMemberPermissionV1, Reject> {
1692 match evidence {
1693 ImportPermissionEvidence::Import(p) => Ok(p),
1694 ImportPermissionEvidence::OwnerCapability(_) | ImportPermissionEvidence::OnlineRole(_) => {
1695 Err(Reject::ImportPermission)
1696 }
1697 }
1698}
1699pub fn require_import_operation_format(format: &str) -> Result<(), Reject> {
1701 if format != OPERATION_DOMAIN {
1702 return Err(Reject::Protocol);
1703 }
1704 Ok(())
1705}
1706pub fn frontier_digest(frontier: &ImportFrontierV1) -> Result<Vec<u8>, Reject> {
1707 if frontier.format_version != 1 {
1708 return Err(Reject::Version);
1709 }
1710 width(&frontier.thread_id, 32)?;
1711 if frontier.operation_ids.len() > 128 {
1712 return Err(Reject::Bounds);
1713 }
1714 for id in &frontier.operation_ids {
1715 width(id, 32)?;
1716 }
1717 if frontier.operation_ids.windows(2).any(|w| w[0] >= w[1]) {
1718 return Err(Reject::Canonical);
1719 }
1720 signing_digest("heddle-import-frontier-v1", frontier)
1721}
1722pub fn content_digest(content: &ImportContentV1) -> Result<Vec<u8>, Reject> {
1723 if content.format_version != 1 {
1724 return Err(Reject::Version);
1725 }
1726 if content.canonical_capture.is_empty() {
1727 return Err(Reject::Bounds);
1728 }
1729 signing_digest("heddle-import-content-v1", content)
1730}
1731pub fn signed_native_digest(record: &SignedRecord) -> Result<Vec<u8>, Reject> {
1732 signing_digest("heddle-signed-native-record-v1", record)
1733}
1734pub(crate) fn verify_native(record: &SignedRecord, format: &str) -> Result<(), Reject> {
1735 if record.format != format {
1736 return Err(Reject::Version);
1737 }
1738 if record.canonical_record.is_empty()
1739 || record.canonical_record.len() > MAX_RECORD_BYTES
1740 || record.signatures.is_empty()
1741 || record.signatures.len() > 16
1742 {
1743 return Err(Reject::Bounds);
1744 }
1745 let mut previous: Option<&[u8]> = None;
1746 let input = [format.as_bytes(), b"\0", &record.canonical_record].concat();
1747 for s in &record.signatures {
1748 if previous.is_some_and(|p| p >= s.public_key.as_slice()) {
1749 return Err(Reject::Canonical);
1750 }
1751 verify(&s.public_key, &input, &s.signature)?;
1752 previous = Some(&s.public_key);
1753 }
1754 Ok(())
1755}
1756pub fn verify_boundary_acceptance(e: &ImportBoundaryAcceptanceV1) -> Result<(), Reject> {
1760 let binding = e.binding.as_ref().ok_or(Reject::BoundaryAcceptance)?;
1761 validate_boundary_binding(binding)?;
1762 let acceptance = e
1763 .signed_acceptance
1764 .as_ref()
1765 .ok_or(Reject::BoundaryAcceptance)?;
1766 verify_native(acceptance, "heddle-original-boundary-acceptance-v1")?;
1767 if acceptance.signatures.len() != 1 {
1768 return Err(Reject::Signature);
1769 }
1770 if e.originals_manifest.is_empty()
1771 || e.publication_intent.is_empty()
1772 || e.originals_manifest.len() > MAX_RECORD_BYTES
1773 || e.publication_intent.len() > MAX_RECORD_BYTES
1774 || e.original_receipts.is_empty()
1775 || e.original_receipts.len() > 128
1776 {
1777 return Err(Reject::Bounds);
1778 }
1779 if binding.acceptance_id != native_id(acceptance)
1780 || binding.signed_acceptance_digest != signed_native_digest(acceptance)?
1781 || binding.originals_manifest_digest
1782 != boundary_octets_digest(
1783 "heddle-boundary-originals-manifest-v1",
1784 &e.originals_manifest,
1785 )
1786 || binding.publication_intent_digest
1787 != boundary_octets_digest(
1788 "heddle-boundary-publication-intent-v1",
1789 &e.publication_intent,
1790 )
1791 {
1792 return Err(Reject::BoundaryAcceptance);
1793 }
1794 let native: NativeBoundarySelection =
1795 rmp_serde::from_slice(&acceptance.canonical_record).map_err(|_| Reject::Canonical)?;
1796 if native.originals_manifest.as_slice()
1797 != native_octets_id(
1798 "heddle-original-publication-manifest-v1",
1799 &e.originals_manifest,
1800 )
1801 || native.publication_intent.as_slice()
1802 != native_octets_id(
1803 "heddle-original-publication-intent-v1",
1804 &e.publication_intent,
1805 )
1806 {
1807 return Err(Reject::BoundaryAcceptance);
1808 }
1809 let mut digests = Vec::new();
1810 for receipt in &e.original_receipts {
1811 if ![
1812 "heddle-thread-genesis-admission-v2",
1813 "heddle-thread-authority-admission-v3",
1814 ]
1815 .contains(&receipt.format.as_str())
1816 {
1817 return Err(Reject::Version);
1818 }
1819 verify_native(receipt, &receipt.format)?;
1820 if receipt.signatures.len() != 1 {
1821 return Err(Reject::Signature);
1822 }
1823 let native: NativeBoundaryReceipt =
1824 rmp_serde::from_slice(&receipt.canonical_record).map_err(|_| Reject::Canonical)?;
1825 if native.basis
1826 != (NativeBoundaryBasis::BoundaryAcceptance {
1827 acceptance: binding
1828 .acceptance_id
1829 .as_slice()
1830 .try_into()
1831 .map_err(|_| Reject::Canonical)?,
1832 })
1833 {
1834 return Err(Reject::BoundaryAcceptance);
1835 }
1836 digests.push(signed_native_digest(receipt)?);
1837 }
1838 if digests != binding.original_receipt_digests {
1839 return Err(Reject::BoundaryAcceptance);
1840 }
1841 Ok(())
1842}
1843#[derive(serde::Deserialize)]
1846struct NativeBoundarySelection {
1847 originals_manifest: [u8; 32],
1848 publication_intent: [u8; 32],
1849}
1850#[derive(serde::Deserialize, PartialEq)]
1851enum NativeBoundaryBasis {
1852 OriginalAuthority,
1853 BoundaryAcceptance { acceptance: [u8; 32] },
1854}
1855#[derive(serde::Deserialize)]
1856struct NativeBoundaryReceipt {
1857 basis: NativeBoundaryBasis,
1858 thread: [u8; 32],
1859 subject: Option<NativeBoundarySubject>,
1860}
1861#[derive(serde::Deserialize)]
1862enum NativeBoundarySubject {
1863 Operation([u8; 32]),
1864 OwnershipClaim([u8; 32]),
1865 OwnershipResolution([u8; 32]),
1866}
1867pub(crate) fn native_octets_id(format: &str, bytes: &[u8]) -> Vec<u8> {
1868 let mut h = blake3::Hasher::new();
1869 h.update(format.as_bytes());
1870 h.update(&(bytes.len() as u64).to_le_bytes());
1871 h.update(b"\0");
1872 h.update(bytes);
1873 h.finalize().as_bytes().to_vec()
1874}
1875pub(crate) fn boundary_original(
1876 e: &ImportBoundaryAcceptanceV1,
1877 original: &SignedRecord,
1878) -> Result<(), Reject> {
1879 let id = native_id(original);
1880 for receipt in &e.original_receipts {
1881 let value: NativeBoundaryReceipt =
1882 rmp_serde::from_slice(&receipt.canonical_record).map_err(|_| Reject::Canonical)?;
1883 let (format, subject) = match value.subject {
1884 None if receipt.format == "heddle-thread-genesis-admission-v2" => {
1885 ("heddle-thread-genesis-v1", value.thread)
1886 }
1887 Some(NativeBoundarySubject::Operation(id)) => ("heddle-thread-operation-v1", id),
1888 Some(NativeBoundarySubject::OwnershipClaim(id)) => {
1889 ("heddle-thread-ownership-claim-v1", id)
1890 }
1891 Some(NativeBoundarySubject::OwnershipResolution(id)) => {
1892 ("heddle-thread-ownership-resolution-v1", id)
1893 }
1894 _ => return Err(Reject::BoundaryAcceptance),
1895 };
1896 if original.format == format && id == subject {
1897 return Ok(());
1898 }
1899 }
1900 Err(Reject::BoundaryAcceptance)
1901}
1902pub fn boundary_octets_digest(domain: &str, bytes: &[u8]) -> Vec<u8> {
1903 hash(&[
1904 domain.as_bytes(),
1905 &(bytes.len() as u32).to_be_bytes(),
1906 bytes,
1907 ])
1908}
1909pub fn validate_boundary_binding(
1910 b: &crate::heddle::api::common::HostedWitnessBoundaryAcceptanceV1,
1911) -> Result<(), Reject> {
1912 if b.format_version != 1 {
1913 return Err(Reject::Version);
1914 }
1915 for digest in [
1916 &b.acceptance_id,
1917 &b.signed_acceptance_digest,
1918 &b.originals_manifest_digest,
1919 &b.publication_intent_digest,
1920 ] {
1921 width(digest, 32)?;
1922 }
1923 if b.original_receipt_digests.is_empty() || b.original_receipt_digests.len() > 128 {
1924 return Err(Reject::Bounds);
1925 }
1926 for digest in &b.original_receipt_digests {
1927 width(digest, 32)?;
1928 }
1929 if b.original_receipt_digests.windows(2).any(|w| w[0] >= w[1]) {
1930 return Err(Reject::Canonical);
1931 }
1932 Ok(())
1933}
1934pub fn validate_statement_boundary(
1935 s: &crate::heddle::api::common::HostedWitnessStatementV1,
1936) -> Result<(), Reject> {
1937 match (s.basis, s.boundary_acceptance.as_ref()) {
1938 (1, None) => Ok(()),
1939 (2, Some(b)) if s.purpose == 1 || s.purpose == 2 => validate_boundary_binding(b),
1940 _ => Err(Reject::BoundaryAcceptance),
1941 }
1942}
1943pub(crate) fn match_boundary(
1944 s: &crate::heddle::api::common::HostedWitnessStatementV1,
1945 evidence: &[ImportBoundaryAcceptanceV1],
1946) -> Result<(), Reject> {
1947 validate_statement_boundary(s)?;
1948 let mut previous = None;
1949 for e in evidence {
1950 verify_boundary_acceptance(e)?;
1951 let b = e.binding.as_ref().ok_or(Reject::BoundaryAcceptance)?;
1952 if previous.is_some_and(|p: &[u8]| p >= b.acceptance_id.as_slice()) {
1953 return Err(Reject::Canonical);
1954 }
1955 previous = Some(b.acceptance_id.as_slice());
1956 }
1957 if let Some(binding) = &s.boundary_acceptance
1958 && !evidence.iter().any(|e| e.binding.as_ref() == Some(binding))
1959 {
1960 return Err(Reject::BoundaryAcceptance);
1961 }
1962 Ok(())
1963}
1964fn native_dependencies(
1965 records: &[SignedRecord],
1966 evidence: &[ImportBoundaryAcceptanceV1],
1967 reviews_only: bool,
1968) -> Result<(), Reject> {
1969 if records.len() > 128 {
1970 return Err(Reject::Bounds);
1971 }
1972 let mut previous = None;
1973 for record in records {
1974 if reviews_only {
1975 require_review_operation(record)?;
1976 }
1977 match record.format.as_str() {
1978 "heddle-thread-genesis-v1"
1979 | "heddle-thread-operation-v1"
1980 | "heddle-thread-ownership-claim-v1"
1981 | "heddle-thread-ownership-resolution-v1" => (),
1982 "heddle-original-boundary-acceptance-v1"
1983 | "heddle-thread-genesis-admission-v2"
1984 | "heddle-thread-authority-admission-v3" => {
1985 if !evidence.iter().any(|e| {
1986 e.signed_acceptance.as_ref() == Some(record)
1987 || e.original_receipts.contains(record)
1988 }) {
1989 return Err(Reject::BoundaryAcceptance);
1990 }
1991 }
1992 _ => return Err(Reject::Version),
1993 }
1994 verify_native(record, &record.format)?;
1995 let digest = signed_native_digest(record)?;
1996 if previous.as_ref().is_some_and(|p| p >= &digest) {
1997 return Err(Reject::Canonical);
1998 }
1999 previous = Some(digest);
2000 }
2001 Ok(())
2002}
2003#[derive(serde::Deserialize)]
2004struct ReviewOperation {
2005 body: ReviewBody,
2006}
2007#[derive(serde::Deserialize)]
2008struct ReviewBody {
2009 kind: String,
2010 canonical: Vec<u8>,
2011}
2012#[derive(serde::Deserialize)]
2013struct ReviewControl {
2014 control: ReviewKind,
2015}
2016#[derive(serde::Deserialize)]
2017struct ReviewKind {
2018 kind: String,
2019}
2020fn require_review_operation(record: &SignedRecord) -> Result<(), Reject> {
2021 if record.format != "heddle-thread-operation-v1" {
2022 return Err(Reject::Semantic);
2023 }
2024 let operation: ReviewOperation =
2025 rmp_serde::from_slice(&record.canonical_record).map_err(|_| Reject::Semantic)?;
2026 if operation.body.kind != "metadata" {
2027 return Err(Reject::Semantic);
2028 }
2029 let control: ReviewControl =
2030 rmp_serde::from_slice(&operation.body.canonical).map_err(|_| Reject::Semantic)?;
2031 if control.control.kind != "review" {
2032 return Err(Reject::Semantic);
2033 }
2034 Ok(())
2035}
2036pub(crate) fn original_signatures(
2037 records: &[&SignedRecord],
2038 extra: &[RecordSignature],
2039) -> Result<Vec<u8>, Reject> {
2040 let signatures = records
2041 .iter()
2042 .flat_map(|r| r.signatures.iter())
2043 .chain(extra.iter())
2044 .collect::<Vec<_>>();
2045 let mut out = (signatures.len() as u32).to_be_bytes().to_vec();
2046 for s in signatures {
2047 s.write(&mut out)?;
2048 }
2049 Ok(hash(&[b"heddle-hosted-original-signatures-v1", &out]))
2050}
2051use crate::hybrid_codec::Canonical;
2052pub fn verify_landing_key_roles(
2054 payload: &HostedLandingWitnessV1,
2055 known_job_keys: &[Vec<u8>],
2056 forbidden_keys: &[Vec<u8>],
2057) -> Result<(), Reject> {
2058 let key = &payload
2059 .request
2060 .as_ref()
2061 .and_then(|r| r.signature.as_ref())
2062 .ok_or(Reject::Signature)?
2063 .public_key;
2064 if known_job_keys.contains(key) || forbidden_keys.contains(key) {
2065 return Err(Reject::KeyRole);
2066 }
2067 Ok(())
2068}
2069
2070#[derive(Clone, Copy)]
2072pub enum WitnessPayload<'a> {
2073 Genesis(&'a ImportGenesisWitnessV1),
2074 Authority(&'a ImportAuthorityWitnessV1),
2075 Landing(&'a HostedLandingWitnessV1),
2076}
2077pub fn verify_witness_payload(
2078 statement: &crate::heddle::api::common::HostedWitnessStatementV1,
2079 payload: WitnessPayload<'_>,
2080) -> Result<(), Reject> {
2081 let (purpose, bytes, authority, signatures, publisher) = match payload {
2082 WitnessPayload::Genesis(p) => {
2083 if p.format_version != 1 {
2084 return Err(Reject::Version);
2085 }
2086 let original = p.original_genesis.as_ref().ok_or(Reject::Canonical)?;
2087 let binding = p.binding.as_ref().ok_or(Reject::Canonical)?;
2088 let b = binding.body.as_ref().ok_or(Reject::Canonical)?;
2089 match_boundary(
2090 statement,
2091 &p.boundary_acceptance.iter().cloned().collect::<Vec<_>>(),
2092 )?;
2093 if let Some(e) = &p.boundary_acceptance {
2094 boundary_original(e, original)?;
2095 }
2096 if (statement.basis == 2) != p.boundary_acceptance.is_some() {
2097 return Err(Reject::BoundaryAcceptance);
2098 }
2099 verify_native(original, "heddle-thread-genesis-v1")?;
2100 if native_id(original) != b.genesis_digest {
2101 return Err(Reject::Scope);
2102 }
2103 if let Some(id) = &b.identity {
2104 if statement.spool_uuid != id.spool_uuid
2105 || statement.spool_genesis_digest != id.spool_genesis_digest
2106 || statement.owner_id != id.owner_id
2107 || statement.owner_state_hash != id.owner_state_hash
2108 || statement.ownership_transfer_sequence != id.ownership_transfer_sequence
2109 {
2110 return Err(Reject::Scope);
2111 }
2112 } else {
2113 return Err(Reject::Canonical);
2114 }
2115 let creator = original
2116 .signatures
2117 .iter()
2118 .find(|s| s.public_key == b.creator_public_key)
2119 .ok_or(Reject::Signature)?;
2120 if b.original_creator_signature != creator.signature
2121 || b.creator_authority_envelope_digest != hash(&[&p.creator_authority_envelope])
2122 {
2123 return Err(Reject::Scope);
2124 }
2125 verify_authorization_signature(
2126 &b.creator_public_key,
2127 GENESIS_DOMAIN,
2128 b,
2129 binding
2130 .creator_signature
2131 .as_ref()
2132 .ok_or(Reject::Signature)?,
2133 )?;
2134 (
2135 1,
2136 canonical(p)?,
2137 signed_genesis_digest(binding)?,
2138 original_signatures(&[original], &[])?,
2139 key_id(&b.creator_public_key),
2140 )
2141 }
2142 WitnessPayload::Authority(p) => {
2143 if p.format_version != 1 {
2144 return Err(Reject::Version);
2145 }
2146 let original = p.original.as_ref().ok_or(Reject::Canonical)?;
2147 let format = match p.kind {
2148 1 => "heddle-thread-operation-v1",
2149 2 => "heddle-thread-ownership-claim-v1",
2150 3 => "heddle-thread-ownership-resolution-v1",
2151 _ => return Err(Reject::Version),
2152 };
2153 verify_native(original, format)?;
2154 if (p.kind == 2 || p.kind == 3) && original.signatures.len() != 2 {
2155 return Err(Reject::Signature);
2156 }
2157 if !original
2158 .signatures
2159 .iter()
2160 .any(|s| key_id(&s.public_key) == statement.publisher_key_id)
2161 {
2162 return Err(Reject::Signature);
2163 }
2164 if p.authority_envelope.is_empty() || p.authority_envelope.len() > MAX_RECORD_BYTES {
2165 return Err(Reject::Bounds);
2166 }
2167 match_boundary(statement, &p.boundary_acceptances)?;
2168 if let Some(b) = &statement.boundary_acceptance {
2169 let e = p
2170 .boundary_acceptances
2171 .iter()
2172 .find(|e| e.binding.as_ref() == Some(b))
2173 .ok_or(Reject::BoundaryAcceptance)?;
2174 boundary_original(e, original)?;
2175 }
2176 native_dependencies(&p.dependencies, &p.boundary_acceptances, false)?;
2177 let records = std::iter::once(original)
2178 .chain(p.dependencies.iter())
2179 .collect::<Vec<_>>();
2180 (
2181 2,
2182 canonical(p)?,
2183 hash(&[
2184 b"heddle-hosted-authority-envelope-v1",
2185 &(p.authority_envelope.len() as u32).to_be_bytes(),
2186 &p.authority_envelope,
2187 ]),
2188 original_signatures(&records, &[])?,
2189 statement.publisher_key_id.clone(),
2190 )
2191 }
2192 WitnessPayload::Landing(p) => {
2193 if p.format_version != 1 {
2194 return Err(Reject::Version);
2195 }
2196 let execution = p.execution.as_ref().ok_or(Reject::Canonical)?;
2197 let source = p.source_operation.as_ref().ok_or(Reject::Canonical)?;
2198 let request = p.request.as_ref().ok_or(Reject::Canonical)?;
2199 if request.format_version != 1
2200 || request.method_path != "/heddle.api.v1alpha2.ThreadService/LandThread"
2201 {
2202 return Err(Reject::Version);
2203 }
2204 verify_native(execution, "heddle-thread-operation-v1")?;
2205 verify_native(source, "heddle-thread-operation-v1")?;
2206 match_boundary(statement, &[])?;
2207 native_dependencies(&p.review_evidence, &[], true)?;
2208 let signature = request.signature.as_ref().ok_or(Reject::Signature)?;
2209 if request.signing_identity
2210 != format!(
2211 "principal:device-key:{}",
2212 hex::encode(&signature.public_key)
2213 )
2214 {
2215 return Err(Reject::Signature);
2216 }
2217 width(&request.nonce, 16)?;
2218 if request.timestamp_millis <= 0
2219 || request.request_body.is_empty()
2220 || request.request_body.len() > MAX_RECORD_BYTES
2221 || p.authority_envelope.is_empty()
2222 || p.authority_envelope.len() > MAX_RECORD_BYTES
2223 {
2224 return Err(Reject::Bounds);
2225 }
2226 let input = crate::signing::unary_bytes(
2227 &request.signing_identity,
2228 &request.method_path,
2229 request.timestamp_millis,
2230 &request.nonce,
2231 &request.request_body,
2232 );
2233 verify(&signature.public_key, &input, &signature.signature)?;
2234 let records = [execution, source]
2235 .into_iter()
2236 .chain(p.review_evidence.iter())
2237 .collect::<Vec<_>>();
2238 (
2239 4,
2240 canonical(p)?,
2241 hash(&[
2242 b"heddle-hosted-authority-envelope-v1",
2243 &(p.authority_envelope.len() as u32).to_be_bytes(),
2244 &p.authority_envelope,
2245 ]),
2246 original_signatures(&records, std::slice::from_ref(signature))?,
2247 key_id(&signature.public_key),
2248 )
2249 }
2250 };
2251 if bytes.len() > MAX_RECORD_BYTES {
2252 return Err(Reject::Bounds);
2253 }
2254 if statement.purpose != purpose
2255 || statement.canonical_payload != bytes
2256 || statement.authority_digest != authority
2257 || statement.original_signatures_digest != signatures
2258 || statement.publisher_key_id != publisher
2259 {
2260 return Err(Reject::Scope);
2261 }
2262 Ok(())
2263}
2264
2265pub fn resolve_bundle_permission<'a>(
2266 bundle: &'a ImportPublicProofBundleV1,
2267 digest: &[u8],
2268) -> Result<Option<&'a SignedImportMemberPermissionV1>, Reject> {
2269 width(digest, 32)?;
2270 if digest == [0; 32] {
2271 return Ok(None);
2272 }
2273 bundle
2274 .member_permission
2275 .iter()
2276 .find(|p| signed_permission_digest(p).is_ok_and(|d| d == digest))
2277 .map(Some)
2278 .ok_or(Reject::ImportPermission)
2279}
2280pub fn resolve_bundle_manifest<'a>(
2281 bundle: &'a ImportPublicProofBundleV1,
2282 digest: &[u8],
2283) -> Result<&'a ImportResultManifestV1, Reject> {
2284 width(digest, 32)?;
2285 bundle
2286 .manifests
2287 .iter()
2288 .find(|m| manifest_digest(m).is_ok_and(|d| d == digest))
2289 .ok_or(Reject::StaleManifest)
2290}
2291pub fn publication_payload(
2292 operation: &SignedDelegatedImportOperationV1,
2293 manifest: &ImportResultManifestV1,
2294) -> Result<ImportPublicationWitnessV1, Reject> {
2295 let o = operation.body.as_ref().ok_or(Reject::Canonical)?;
2296 Ok(ImportPublicationWitnessV1 {
2297 format_version: 1,
2298 signed_operation_digest: signed_operation_digest(operation)?,
2299 delegation_digest: o.delegation_digest.clone(),
2300 logical_job_id: o.logical_job_id.clone(),
2301 retry_lineage_id: o.retry_lineage_id.clone(),
2302 physical_operation_id: o.physical_operation_id.clone(),
2303 ref_name: o.ref_name.clone(),
2304 slot_id: o.slot_id,
2305 hash_algorithm: o.hash_algorithm,
2306 observed_commit_oid: o.observed_commit_oid.clone(),
2307 expected_frontier_digest: o.expected_frontier_digest.clone(),
2308 resulting_frontier_digest: o.resulting_frontier_digest.clone(),
2309 terminal_manifest_digest: manifest_digest(manifest)?,
2310 })
2311}
2312fn validate_bundle_history(
2315 bundle: &ImportPublicProofBundleV1,
2316 require_admissions: bool,
2317) -> Result<(), Reject> {
2318 fn sorted<T>(
2319 values: &[T],
2320 digest: impl Fn(&T) -> Result<Vec<u8>, Reject>,
2321 ) -> Result<(), Reject> {
2322 let mut previous = None;
2323 for value in values {
2324 let d = digest(value)?;
2325 if previous.as_ref().is_some_and(|p| p >= &d) {
2326 return Err(Reject::Canonical);
2327 }
2328 previous = Some(d);
2329 }
2330 Ok(())
2331 }
2332 let mut foreign = crate::foreign_dependencies::References::new(
2333 &bundle.foreign_dependencies,
2334 ForeignDependencyOrigin::Import,
2335 )?;
2336 for p in &bundle.authority_witnesses {
2337 let subject_thread =
2338 crate::foreign_dependencies::thread(p.original.as_ref().ok_or(Reject::Canonical)?)?;
2339 if !bundle.genesis_witnesses.iter().any(|g| {
2340 g.original_genesis
2341 .as_ref()
2342 .is_some_and(|g| native_id(g) == subject_thread)
2343 }) {
2344 return Err(Reject::Scope);
2345 }
2346 }
2347 let job_keys = bundle
2348 .delegations
2349 .iter()
2350 .filter_map(|d| d.body.as_ref().map(|d| d.job_public_key.clone()))
2351 .collect::<Vec<_>>();
2352 for p in &bundle.landing_witnesses {
2353 verify_landing_key_roles(p, &job_keys, &[])?;
2354 let execution = p.execution.as_ref().ok_or(Reject::Canonical)?;
2355 let thread = crate::foreign_dependencies::thread(execution)?;
2356 if !bundle.genesis_witnesses.iter().any(|p| {
2357 p.original_genesis
2358 .as_ref()
2359 .is_some_and(|g| native_id(g) == thread)
2360 }) {
2361 return Err(Reject::Scope);
2362 }
2363 }
2364 for original in bundle
2366 .authority_witnesses
2367 .iter()
2368 .flat_map(|p| p.original.iter().chain(p.dependencies.iter()))
2369 .chain(bundle.landing_witnesses.iter().flat_map(|p| {
2370 p.execution
2371 .iter()
2372 .chain(p.source_operation.iter())
2373 .chain(p.review_evidence.iter())
2374 }))
2375 {
2376 if [
2377 "heddle-original-boundary-acceptance-v1",
2378 "heddle-thread-genesis-admission-v2",
2379 "heddle-thread-authority-admission-v3",
2380 ]
2381 .contains(&original.format.as_str())
2382 {
2383 continue;
2384 }
2385 let thread = crate::foreign_dependencies::thread(original)?;
2386 if !bundle.genesis_witnesses.iter().any(|p| {
2387 p.original_genesis
2388 .as_ref()
2389 .is_some_and(|g| native_id(g) == thread)
2390 }) {
2391 foreign.require(original)?;
2392 } else if original.format == "heddle-thread-genesis-v1"
2393 && !bundle
2394 .genesis_witnesses
2395 .iter()
2396 .any(|p| p.original_genesis.as_ref() == Some(original))
2397 {
2398 return Err(Reject::Scope);
2399 }
2400 }
2401 foreign.finish()?;
2402 for operation in &bundle.operations {
2403 if operation.body.is_none() {
2404 return Err(Reject::Canonical);
2405 }
2406 }
2407 for statement in &bundle.statements {
2408 let s = statement.body.as_ref().ok_or(Reject::Canonical)?;
2409 validate_statement_boundary(s)?;
2410 require_policy_history(
2411 &bundle.policies,
2412 &s.spool_uuid,
2413 s.policy_sequence,
2414 &s.policy_state_hash,
2415 )?;
2416 let envelope = match s.purpose {
2417 1 if s.basis == 2 => Some(
2418 bundle
2419 .genesis_witnesses
2420 .iter()
2421 .find(|p| canonical(*p).is_ok_and(|v| v == s.canonical_payload))
2422 .ok_or(Reject::Scope)?
2423 .creator_authority_envelope
2424 .as_slice(),
2425 ),
2426 2 => Some(
2427 bundle
2428 .authority_witnesses
2429 .iter()
2430 .find(|p| canonical(*p).is_ok_and(|bytes| bytes == s.canonical_payload))
2431 .ok_or(Reject::Scope)?
2432 .authority_envelope
2433 .as_slice(),
2434 ),
2435 4 => Some(
2436 bundle
2437 .landing_witnesses
2438 .iter()
2439 .find(|p| canonical(*p).is_ok_and(|bytes| bytes == s.canonical_payload))
2440 .ok_or(Reject::Scope)?
2441 .authority_envelope
2442 .as_slice(),
2443 ),
2444 _ => None,
2445 };
2446 if let Some(envelope) = envelope {
2447 crate::writer_authority::check_witness_writer(
2448 s,
2449 envelope,
2450 &bundle.owner_histories,
2451 &bundle.policies,
2452 crate::writer_authority::spool_account_for_statement(
2453 s,
2454 bundle
2455 .delegations
2456 .iter()
2457 .filter_map(|d| d.body.as_ref()?.identity.as_ref())
2458 .chain(
2459 bundle
2460 .genesis_authorities
2461 .iter()
2462 .filter_map(|g| g.body.as_ref()?.identity.as_ref()),
2463 ),
2464 &bundle.ownership_transfers,
2465 )?,
2466 bundle
2467 .authority_witnesses
2468 .iter()
2469 .find(|p| {
2470 (p.kind == 2 || p.kind == 3)
2471 && canonical(*p).is_ok_and(|v| v == s.canonical_payload)
2472 })
2473 .and_then(|p| p.original.as_ref())
2474 .map(|r| r.signatures.as_slice())
2475 .unwrap_or(&[]),
2476 if s.purpose == 1 {
2477 bundle
2478 .genesis_witnesses
2479 .iter()
2480 .find(|p| canonical(*p).is_ok_and(|v| v == s.canonical_payload))
2481 .and_then(|p| p.boundary_acceptance.as_ref())
2482 } else {
2483 bundle
2484 .authority_witnesses
2485 .iter()
2486 .find(|p| canonical(*p).is_ok_and(|v| v == s.canonical_payload))
2487 .and_then(|p| {
2488 p.boundary_acceptances
2489 .iter()
2490 .find(|e| e.binding == s.boundary_acceptance)
2491 })
2492 },
2493 )?;
2494 }
2495 }
2496 sorted(&bundle.manifests, manifest_digest)?;
2497 if let Some(p) = &bundle.member_permission
2498 && resolve_bundle_permission(bundle, &signed_permission_digest(p)?)? != Some(p)
2499 {
2500 return Err(Reject::ImportPermission);
2501 }
2502 let terminal = bundle.terminal_manifest.as_ref().ok_or(Reject::Canonical)?;
2503 if resolve_bundle_manifest(bundle, &manifest_digest(terminal)?)? != terminal {
2504 return Err(Reject::Canonical);
2505 }
2506 if bundle.delegations.len() != 1 {
2507 return Err(Reject::Canonical);
2508 }
2509 for d in &bundle.delegations {
2510 let body = d.body.as_ref().ok_or(Reject::Canonical)?;
2511 resolve_bundle_permission(bundle, &body.parent_permission_digest)?;
2512 for branch in &body.branch_manifest {
2513 let g = bundle
2514 .genesis_authorities
2515 .iter()
2516 .find(|g| {
2517 signed_genesis_digest(g).is_ok_and(|h| h == branch.genesis_authority_digest)
2518 })
2519 .ok_or(Reject::Scope)?;
2520 let b = g.body.as_ref().ok_or(Reject::Canonical)?;
2521 resolve_bundle_permission(bundle, &b.parent_permission_digest)?;
2522 if !bundle
2523 .original_geneses
2524 .iter()
2525 .any(|o| native_id(o) == b.genesis_digest)
2526 || !bundle
2527 .creator_authority_envelopes
2528 .iter()
2529 .any(|e| hash(&[e]) == b.creator_authority_envelope_digest)
2530 {
2531 return Err(Reject::Scope);
2532 }
2533 if require_admissions
2536 && bundle.operations.iter().any(|o| {
2537 o.body
2538 .as_ref()
2539 .is_some_and(|o| o.genesis_digest == b.genesis_digest)
2540 })
2541 && !bundle.genesis_witnesses.iter().any(|payload| {
2542 payload.binding.as_ref() == Some(g)
2543 && payload.original_genesis.as_ref().is_some_and(|o| {
2544 native_id(o) == b.genesis_digest && bundle.original_geneses.contains(o)
2545 })
2546 && hash(&[&payload.creator_authority_envelope])
2547 == b.creator_authority_envelope_digest
2548 && canonical(payload).is_ok_and(|bytes| {
2549 bundle.statements.iter().any(|s| {
2550 s.body
2551 .as_ref()
2552 .is_some_and(|s| s.purpose == 1 && s.canonical_payload == bytes)
2553 })
2554 })
2555 })
2556 {
2557 return Err(Reject::Scope);
2558 }
2559 }
2560 }
2561 for manifest in &bundle.manifests {
2562 validate_manifest(manifest)?;
2563 if manifest.logical_job_id != terminal.logical_job_id
2564 || manifest.retry_lineage_id != terminal.retry_lineage_id
2565 {
2566 return Err(Reject::Scope);
2567 }
2568 for slot in &manifest.slots {
2569 let operation = bundle
2570 .operations
2571 .iter()
2572 .find(|o| {
2573 signed_operation_digest(o).is_ok_and(|d| d == slot.signed_operation_digest)
2574 })
2575 .ok_or(Reject::Scope)?;
2576 if !check_slot_replay(manifest, operation)? || !check_slot_replay(terminal, operation)?
2577 {
2578 return Err(Reject::Scope);
2579 }
2580 }
2581 }
2582 for operation in &bundle.operations {
2583 let o = operation.body.as_ref().ok_or(Reject::Canonical)?;
2584 if !bundle
2585 .delegations
2586 .iter()
2587 .any(|d| signed_delegation_digest(d).is_ok_and(|h| h == o.delegation_digest))
2588 || !check_slot_replay(terminal, operation)?
2589 {
2590 return Err(Reject::Scope);
2591 }
2592 if !bundle.manifests.iter().any(|m| {
2593 check_slot_replay(m, operation) == Ok(true)
2594 && publication_payload(operation, m)
2595 .and_then(|p| canonical(&p))
2596 .is_ok_and(|p| {
2597 bundle.statements.iter().any(|s| {
2598 s.body
2599 .as_ref()
2600 .is_some_and(|s| s.purpose == 3 && s.canonical_payload == p)
2601 })
2602 })
2603 }) {
2604 return Err(Reject::Scope);
2605 }
2606 }
2607 for statement in &bundle.statements {
2608 let s = statement.body.as_ref().ok_or(Reject::Canonical)?;
2609 let found = match s.purpose {
2610 1 => bundle
2611 .genesis_witnesses
2612 .iter()
2613 .any(|p| canonical(p).is_ok_and(|p| p == s.canonical_payload)),
2614 2 => bundle
2615 .authority_witnesses
2616 .iter()
2617 .any(|p| canonical(p).is_ok_and(|p| p == s.canonical_payload)),
2618 3 => bundle.operations.iter().any(|o| {
2619 bundle.manifests.iter().any(|m| {
2620 publication_payload(o, m)
2621 .and_then(|p| canonical(&p))
2622 .is_ok_and(|p| p == s.canonical_payload)
2623 })
2624 }),
2625 4 => bundle
2626 .landing_witnesses
2627 .iter()
2628 .any(|p| canonical(p).is_ok_and(|p| p == s.canonical_payload)),
2629 _ => return Err(Reject::Version),
2630 };
2631 if !found {
2632 return Err(Reject::Scope);
2633 }
2634 }
2635 Ok(())
2636}
2637pub(crate) fn require_policy_history(
2643 policies: &[SignedSpoolPolicyRecord],
2644 spool: &[u8],
2645 mut sequence: u64,
2646 state_hash: &[u8],
2647) -> Result<(), Reject> {
2648 let mut previous: Option<(&[u8], u64)> = None;
2649 for signed in policies {
2650 let policy = signed.body.as_ref().ok_or(Reject::Canonical)?;
2651 let digest = policy_state_digest(policy)?;
2652 if digest != policy.policy_state_hash {
2653 return Err(Reject::Canonical);
2654 }
2655 let current = (policy.spool_uuid.as_slice(), policy.sequence);
2656 if previous.is_some_and(|prev| prev >= current) {
2657 return Err(Reject::Canonical);
2658 }
2659 previous = Some(current);
2660 }
2661 let mut state_hash = state_hash.to_vec();
2662 let mut successor_revoked: Option<&[Vec<u8>]> = None;
2663 for _ in 0..=policies.len() {
2664 width(&state_hash, 32)?;
2665 if sequence == 0 {
2666 return if state_hash == [0; 32] {
2667 Ok(())
2668 } else {
2669 Err(Reject::Scope)
2670 };
2671 }
2672 let mut matches = policies.iter().filter_map(|p| p.body.as_ref()).filter(|p| {
2673 p.spool_uuid == spool && p.sequence == sequence && p.policy_state_hash == state_hash
2674 });
2675 let policy = matches.next().ok_or(Reject::Scope)?;
2676 let revoked = policy
2678 .policy
2679 .as_ref()
2680 .map(|p| p.revoked_key_ids.as_slice())
2681 .unwrap_or(&[]);
2682 if successor_revoked.is_some_and(|next| revoked.iter().any(|id| !next.contains(id))) {
2683 return Err(Reject::Scope);
2684 }
2685 successor_revoked = Some(revoked);
2686 let head = policy.expected_head.as_ref().ok_or(Reject::Canonical)?;
2687 if head.sequence.checked_add(1) != Some(sequence) {
2688 return Err(Reject::Scope);
2689 }
2690 sequence = head.sequence;
2691 state_hash = head.state_hash.clone();
2692 }
2693 Err(Reject::Scope)
2694}
2695pub fn policy_state_digest(p: &SignedPolicyBody) -> Result<Vec<u8>, Reject> {
2697 use crate::hybrid_codec::counted;
2698 if p.format_version != 1 || !p.merge_parent_state_hashes.is_empty() {
2699 return Err(Reject::Version);
2700 }
2701 width(&p.spool_uuid, 16)?;
2702 width(&p.owner_id, 32)?;
2703 width(&p.owner_state_hash, 32)?;
2704 let head = p.expected_head.as_ref().ok_or(Reject::Canonical)?;
2705 width(&head.state_hash, 32)?;
2706 let policy = p.policy.as_ref().ok_or(Reject::Canonical)?;
2707 if p.merge_policies.len() != 2
2708 || p.merge_policies[0].setting_key != "max_audience"
2709 || p.merge_policies[0].semantics != 1
2710 || p.merge_policies[1].setting_key != "revoked_key_ids"
2711 || p.merge_policies[1].semantics != 2
2712 {
2713 return Err(Reject::Canonical);
2714 }
2715 let mut out = p.format_version.to_be_bytes().to_vec();
2716 counted(&mut out, &p.spool_uuid)?;
2717 counted(&mut out, &head.state_hash)?;
2718 out.extend_from_slice(&head.sequence.to_be_bytes());
2719 out.extend_from_slice(&p.sequence.to_be_bytes());
2720 out.extend_from_slice(&0u32.to_be_bytes());
2721 out.extend_from_slice(&(policy.revoked_key_ids.len() as u32).to_be_bytes());
2722 for (i, id) in policy.revoked_key_ids.iter().enumerate() {
2723 width(id, 32)?;
2724 if i > 0 && policy.revoked_key_ids[i - 1] >= *id {
2725 return Err(Reject::Canonical);
2726 }
2727 counted(&mut out, id)?;
2728 }
2729 out.push(u8::from(policy.max_audience.is_some()));
2730 if let Some(audience) = policy.max_audience {
2731 if !(1..=3).contains(&audience) {
2732 return Err(Reject::Canonical);
2733 }
2734 out.extend_from_slice(&(audience as u32).to_be_bytes());
2735 }
2736 out.extend_from_slice(&2u32.to_be_bytes());
2737 for rule in &p.merge_policies {
2738 counted(&mut out, rule.setting_key.as_bytes())?;
2739 out.extend_from_slice(&(rule.semantics as u32).to_be_bytes());
2740 }
2741 counted(&mut out, &p.owner_id)?;
2742 counted(&mut out, &p.owner_state_hash)?;
2743 out.extend_from_slice(&p.ownership_transfer_sequence.to_be_bytes());
2744 Ok(hash(&[b"heddle-spool-signed-policy-v2", &out]))
2745}
2746pub(crate) fn native_id(record: &SignedRecord) -> Vec<u8> {
2747 let mut h = blake3::Hasher::new();
2748 h.update(record.format.as_bytes());
2749 h.update(&(record.canonical_record.len() as u64).to_le_bytes());
2750 h.update(b"\0");
2751 h.update(&record.canonical_record);
2752 h.finalize().as_bytes().to_vec()
2753}
2754
2755pub fn initial_operation_id(lineage: &[u8], occupied: bool) -> Result<String, Reject> {
2758 width(lineage, 16)?;
2759 if lineage.iter().all(|b| *b == 0) {
2760 return Err(Reject::Canonical);
2761 }
2762 if occupied {
2763 return Err(Reject::OperationIdReused);
2764 }
2765 let h = hex::encode(lineage);
2766 Ok(format!(
2767 "{}-{}-{}-{}-{}",
2768 &h[..8],
2769 &h[8..12],
2770 &h[12..16],
2771 &h[16..20],
2772 &h[20..]
2773 ))
2774}
2775
2776pub fn check_cancel_request(
2779 request: &CancelImportJobRequest,
2780 active: &SignedImportJobDelegationV1,
2781 durable_epoch: u64,
2782 cancelled: bool,
2783) -> Result<(), Reject> {
2784 let d = active.body.as_ref().ok_or(Reject::Canonical)?;
2785 width(&request.logical_job_id, 16)?;
2786 width(&request.cancellation_id, 32)?;
2787 let id = d.identity.as_ref().ok_or(Reject::Canonical)?;
2788 if request.client_operation_id.is_empty()
2789 || request.destination.as_ref().is_none_or(|s| {
2790 initial_operation_id(&id.spool_uuid, false).map_or(true, |uuid| s.id != uuid)
2791 })
2792 || request.logical_job_id != d.logical_job_id
2793 {
2794 return Err(Reject::Scope);
2795 }
2796 if durable_epoch == 0 || request.expected_authority_epoch != durable_epoch {
2797 return Err(Reject::StaleContext);
2798 }
2799 if request.cancellation_id != d.cancellation_id {
2800 return Err(Reject::Scope);
2801 }
2802 if cancelled {
2803 return Err(Reject::Revoked);
2804 }
2805 Ok(())
2806}
2807pub fn check_cancel_replay(
2809 request: &CancelImportJobRequest,
2810 stored: &CancelImportJobRequest,
2811) -> Result<(), Reject> {
2812 if request != stored {
2813 return Err(Reject::OperationIdReused);
2814 }
2815 Ok(())
2816}
2817pub fn check_import_revocations(
2819 delegation: &SignedImportJobDelegationV1,
2820 member: Option<&SignedImportMemberPermissionV1>,
2821 revoked: &[Vec<u8>],
2822) -> Result<(), Reject> {
2823 let d = delegation.body.as_ref().ok_or(Reject::Canonical)?;
2824 width(&d.cancellation_id, 32)?;
2825 if revoked.contains(&d.cancellation_id) {
2826 return Err(Reject::Revoked);
2827 }
2828 if let Some(parent) = member {
2829 let p = parent.body.as_ref().ok_or(Reject::ImportPermission)?;
2830 width(&p.cancellation_id, 32)?;
2831 if revoked.contains(&p.cancellation_id) {
2832 return Err(Reject::Revoked);
2833 }
2834 }
2835 Ok(())
2836}
2837
2838pub fn remaining_import_scope(
2841 scope: &ImportPermissionScopeV1,
2842 committed: &ImportResultManifestV1,
2843) -> Result<ImportPermissionScopeV1, Reject> {
2844 validate_scope(scope)?;
2845 validate_manifest(committed)?;
2846 let mut remaining = scope.clone();
2847 for slot in &committed.slots {
2848 if !scope
2849 .branches
2850 .iter()
2851 .any(|b| b.ref_name == slot.ref_name && b.slot_id == slot.slot_id)
2852 {
2853 return Err(Reject::Scope);
2854 }
2855 {
2856 remaining.max_result_bytes = remaining
2857 .max_result_bytes
2858 .checked_sub(slot.result_bytes)
2859 .ok_or(Reject::Scope)?;
2860 remaining.max_operations = remaining
2861 .max_operations
2862 .checked_sub(1)
2863 .ok_or(Reject::Scope)?;
2864 }
2865 }
2866 remaining.branches.retain(|b| {
2867 !committed
2868 .slots
2869 .iter()
2870 .any(|s| s.ref_name == b.ref_name && s.slot_id == b.slot_id)
2871 });
2872 Ok(remaining)
2873}
2874
2875pub fn validate_hybrid_import_job_selector(
2877 selector: &HybridImportJobSelector,
2878) -> Result<(), Reject> {
2879 width(&selector.logical_job_id, 16)?;
2880 if selector.logical_job_id.iter().all(|byte| *byte == 0) {
2881 return Err(Reject::Canonical);
2882 }
2883 Ok(())
2884}
2885
2886pub fn import_job_state_request_from_operation(
2891 operation: &OperationRecord,
2892) -> Result<Option<GetImportJobStateRequest>, Reject> {
2893 let Some(operation_subject::Subject::Import(subject)) = operation
2894 .subject
2895 .as_ref()
2896 .and_then(|subject| subject.subject.as_ref())
2897 else {
2898 return Ok(None);
2899 };
2900 let Some(selector) = subject.hybrid_job.as_ref() else {
2901 return Ok(None);
2902 };
2903 validate_hybrid_import_job_selector(selector)?;
2904 let request = GetImportJobStateRequest {
2905 destination: operation
2906 .r#ref
2907 .as_ref()
2908 .and_then(|record| record.spool.clone()),
2909 logical_job_id: selector.logical_job_id.clone(),
2910 };
2911 validate_job_state_request(&request)?;
2912 Ok(Some(request))
2913}
2914
2915pub fn validate_job_state_request(request: &GetImportJobStateRequest) -> Result<(), Reject> {
2918 use prost::Message;
2919 if request.encoded_len() > 4096 {
2920 return Err(Reject::Bounds);
2921 }
2922 initial_operation_id(&request.logical_job_id, false)?;
2923 let destination = request.destination.as_ref().ok_or(Reject::Scope)?;
2924 let compact = destination.id.replace('-', "");
2925 let raw = hex::decode(&compact).map_err(|_| Reject::Scope)?;
2926 if initial_operation_id(&raw, false).map_or(true, |id| id != destination.id) {
2927 return Err(Reject::Scope);
2928 }
2929 Ok(())
2930}
2931
2932pub fn validate_job_state_response(
2934 request: &GetImportJobStateRequest,
2935 response: &GetImportJobStateResponse,
2936) -> Result<(), Reject> {
2937 use prost::Message;
2938 validate_job_state_request(request)?;
2939 if response.encoded_len() > 4096 {
2940 return Err(Reject::Bounds);
2941 }
2942 if !(1..=5).contains(&response.status) || response.authority_epoch == 0 {
2943 return Err(Reject::Canonical);
2944 }
2945 width(&response.active_cancellation_id, 32)?;
2946 width(&response.active_delegation_digest, 32)?;
2947 validate_retry_availability(response, request.destination.as_ref().ok_or(Reject::Scope)?)
2948}
2949
2950#[derive(Debug, Clone, PartialEq)]
2952pub struct ImportWitnessRootPin {
2953 pub authority: String,
2954 pub root_id: String,
2955 pub public_key: Vec<u8>,
2956 pub epoch: u64,
2957}
2958#[derive(Debug, Clone, PartialEq)]
2961pub struct ImportWitnessSnapshot {
2962 pub root: ImportWitnessRootPin,
2963 pub witness_set: crate::heddle::api::common::SignedHostedWitnessSetV1,
2964 pub clock_floor_unix_millis: i64,
2965 pub job_associations: Vec<(Vec<u8>, Vec<u8>)>,
2966 pub accepted_history: Vec<ImportPublicProofBundleV1>,
2967}
2968#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2971pub enum ImportBundleEvidence {
2972 Recovery,
2973 Witnessed,
2974}
2975#[derive(Debug, Clone, PartialEq)]
2976pub struct VerifiedImportBundleWitnesses {
2977 pub evidence: ImportBundleEvidence,
2978 pub accepted_history: ImportResultManifestV1,
2979 pub snapshot: Option<ImportWitnessSnapshot>,
2980 pub snapshot_advanced: bool,
2981 pub owner_check_time_unix_seconds: Option<i64>,
2982}
2983pub fn check_import_genesis_publication_pair(
2986 delegation: &VerifiedImportDelegation,
2987 admission: &crate::heddle::api::common::HostedWitnessStatementV1,
2988 publication: &crate::heddle::api::common::HostedWitnessStatementV1,
2989) -> Result<(), Reject> {
2990 validity(
2991 delegation.body.not_before_unix_seconds,
2992 delegation.body.expires_at_unix_seconds,
2993 admission.observed_at_unix_millis / 1000,
2994 true,
2995 )?;
2996 if admission.purpose != 1
2997 || publication.purpose != 3
2998 || admission.observed_at_unix_millis != publication.observed_at_unix_millis
2999 || admission.host_transaction_id != publication.host_transaction_id
3000 || admission.executor_id != publication.executor_id
3001 || admission.admission_order >= publication.admission_order
3002 {
3003 return Err(Reject::Transition);
3004 }
3005 Ok(())
3006}
3007
3008pub fn verify_import_bundle_witnesses<'a>(
3022 bundle: &ImportPublicProofBundleV1,
3023 pin: &ImportWitnessRootPin,
3024 snapshot: Option<&ImportWitnessSnapshot>,
3025 now_ms: i64,
3026 mut owner_at: impl FnMut(Option<i64>) -> Result<ImportBundleOwnerExpectation<'a>, Reject>,
3027 mut verify_policy: impl FnMut(
3028 &ImportPublicProofBundleV1,
3029 Option<&crate::heddle::api::common::HostedWitnessStatementV1>,
3030 ) -> Result<(), Reject>,
3031) -> Result<VerifiedImportBundleWitnesses, Reject> {
3032 use crate::witness_trust as trust;
3033 width(&pin.public_key, 32)?;
3034 if pin.epoch == 0 {
3035 return Err(Reject::StaleContext);
3036 }
3037 validate_bundle_bounds(bundle)?;
3038 let terminal = bundle.terminal_manifest.as_ref().ok_or(Reject::Canonical)?;
3039 validate_bundle_history(bundle, !terminal.slots.is_empty())?;
3040 let mut associations = snapshot.map_or_else(Vec::new, |s| s.job_associations.clone());
3041 for d in &bundle.delegations {
3042 let b = d.body.as_ref().ok_or(Reject::Canonical)?;
3043 if associations
3044 .iter()
3045 .any(|(key, job)| *key == b.job_public_key && *job != b.logical_job_id)
3046 {
3047 return Err(Reject::KeyRole);
3048 }
3049 if !associations.iter().any(|(key, _)| *key == b.job_public_key) {
3050 associations.push((b.job_public_key.clone(), b.logical_job_id.clone()));
3051 }
3052 }
3053 let job_keys = associations
3054 .iter()
3055 .map(|(key, _)| key.clone())
3056 .collect::<Vec<_>>();
3057 fn expectation<'a>(
3058 root: &'a ImportWitnessRootPin,
3059 floor: i64,
3060 now: i64,
3061 keys: &'a [Vec<u8>],
3062 ) -> trust::SetExpectation<'a> {
3063 trust::SetExpectation {
3064 authority: &root.authority,
3065 root_id: &root.root_id,
3066 root_public_key: &root.public_key,
3067 root_epoch: root.epoch,
3068 now_unix_millis: now,
3069 clock_floor_unix_millis: floor,
3070 known_job_keys: keys,
3071 }
3072 }
3073 let previous = if let Some(s) = snapshot {
3074 if s.root.authority != pin.authority {
3075 return Err(Reject::Root);
3076 }
3077 let restored = trust::restore_history_snapshot(
3078 &s.witness_set,
3079 &expectation(&s.root, 0, now_ms, &job_keys),
3080 )?;
3081 Some(restored)
3082 } else {
3083 None
3084 };
3085 let carried = bundle.witness_set.as_ref();
3086 if carried.is_none() && !bundle.statements.is_empty() {
3087 return Err(Reject::Canonical);
3088 }
3089 let new_set =
3090 carried.is_some_and(|c| snapshot.is_none_or(|s| s.root != *pin || s.witness_set != *c));
3091 let selected = expectation(
3092 pin,
3093 snapshot.map_or(0, |s| s.clock_floor_unix_millis),
3094 now_ms,
3095 &job_keys,
3096 );
3097 let set = carried
3098 .map(|carried| {
3099 if snapshot.is_some_and(|s| s.root != *pin) {
3100 trust::verify_set_after_root_replacement(
3101 carried,
3102 &selected,
3103 previous.as_ref().ok_or(Reject::StaleContext)?,
3104 )
3105 } else {
3106 trust::verify_set(carried, &selected, previous.as_ref())
3107 }
3108 })
3109 .transpose()?;
3110 let mut resolved = Vec::new();
3112 for signed in &bundle.statements {
3113 let s = signed.body.as_ref().ok_or(Reject::Canonical)?;
3114 let set = set.as_ref().ok_or(Reject::Canonical)?;
3115 let entry = set
3116 .body()
3117 .entries
3118 .iter()
3119 .find(|e| e.executor_id == s.executor_id)
3120 .ok_or(Reject::Root)?;
3121 let proof = if entry.state == 2 {
3122 let leaf = trust::leaf_digest(s.purpose, &canonical(s)?, &signed.signature)?;
3123 bundle.history_proofs.iter().find(|p| {
3124 p.purpose == s.purpose && trust::verify_inclusion(&leaf, p, entry).is_ok()
3125 })
3126 } else {
3127 None
3128 };
3129 let context = trust::resolve_statement(set, signed, proof, false, now_ms)?;
3130 verify_policy(bundle, Some(s))?;
3131 resolved.push((context, proof));
3132 }
3133 if bundle.statements.is_empty() {
3134 verify_policy(bundle, None)?;
3136 }
3137 let mut prefix = ImportResultManifestV1 {
3142 slots: vec![],
3143 ..terminal.clone()
3144 };
3145 let mut selected_publications = Vec::new();
3146 for operation in &bundle.operations {
3147 let digest = signed_operation_digest(operation)?;
3148 let slot = terminal
3149 .slots
3150 .iter()
3151 .find(|slot| slot.signed_operation_digest == digest)
3152 .ok_or(Reject::Scope)?;
3153 prefix.slots.push(slot.clone());
3154 prefix
3155 .slots
3156 .sort_by(|a, b| (&a.ref_name, a.slot_id).cmp(&(&b.ref_name, b.slot_id)));
3157 let payload = canonical(&publication_payload(operation, &prefix)?)?;
3158 let matches = bundle
3159 .statements
3160 .iter()
3161 .filter(|signed| {
3162 signed
3163 .body
3164 .as_ref()
3165 .is_some_and(|s| s.purpose == 3 && s.canonical_payload == payload)
3166 })
3167 .collect::<Vec<_>>();
3168 if matches.len() != 1 {
3169 return Err(Reject::Transition);
3170 }
3171 selected_publications.push(matches[0]);
3172 }
3173 if bundle.statements.iter().any(|signed| {
3174 signed
3175 .body
3176 .as_ref()
3177 .is_some_and(|s| s.purpose == 3 && !selected_publications.contains(&signed))
3178 }) {
3179 return Err(Reject::Transition);
3180 }
3181 let mut admitted_geneses = std::collections::BTreeSet::<Vec<u8>>::new();
3182 let mut times = vec![None; bundle.delegations.len()];
3183 let mut publications = Vec::new();
3184 for signed in &bundle.statements {
3185 let statement = signed.body.as_ref().ok_or(Reject::Canonical)?;
3186 if statement.purpose != 3 {
3187 continue;
3188 }
3189 let (operation, manifest) = bundle
3190 .operations
3191 .iter()
3192 .flat_map(|o| bundle.manifests.iter().map(move |m| (o, m)))
3193 .find(|(o, m)| {
3194 publication_payload(o, m)
3195 .and_then(|p| canonical(&p))
3196 .is_ok_and(|b| b == statement.canonical_payload)
3197 })
3198 .ok_or(Reject::Scope)?;
3199 let digest = &operation
3200 .body
3201 .as_ref()
3202 .ok_or(Reject::Canonical)?
3203 .delegation_digest;
3204 let index = bundle
3205 .delegations
3206 .iter()
3207 .position(|d| signed_delegation_digest(d).is_ok_and(|h| h == *digest))
3208 .ok_or(Reject::Scope)?;
3209 let time = statement.observed_at_unix_millis / 1000;
3210 if times[index].is_none_or(|(order, _)| statement.admission_order < order) {
3212 times[index] = Some((statement.admission_order, time));
3213 }
3214 publications.push((signed, operation, manifest, index));
3215 }
3216 let mut resolve_owner = |time: Option<i64>| {
3219 let facts = owner_at(time)?;
3220 if facts.effective_from_unix_seconds < 0
3221 || facts
3222 .effective_until_unix_seconds
3223 .is_some_and(|end| end <= facts.effective_from_unix_seconds)
3224 || time.is_some_and(|t| {
3225 t < facts.effective_from_unix_seconds
3226 || facts
3227 .effective_until_unix_seconds
3228 .is_some_and(|end| t >= end)
3229 })
3230 {
3231 return Err(Reject::Scope);
3232 }
3233 for (key, job) in facts.known_job_associations {
3234 if associations.iter().any(|(k, j)| k == key && j != job) {
3235 return Err(Reject::KeyRole);
3236 }
3237 if !associations.iter().any(|(k, _)| k == key) {
3238 associations.push((key.clone(), job.clone()));
3239 }
3240 }
3241 if facts
3242 .forbidden_job_keys
3243 .iter()
3244 .any(|key| associations.iter().any(|(known, _)| known == key))
3245 {
3246 return Err(Reject::KeyRole);
3247 }
3248 if set.as_ref().is_some_and(|set| {
3250 facts.known_job_associations.iter().any(|(key, _)| {
3251 *key == pin.public_key || set.body().entries.iter().any(|e| e.public_key == *key)
3252 })
3253 }) {
3254 return Err(Reject::KeyRole);
3255 }
3256 Ok((facts, associations.clone()))
3257 };
3258 let mut verified = Vec::new();
3259 for (i, d) in bundle.delegations.iter().enumerate() {
3260 let body = d.body.as_ref().ok_or(Reject::Canonical)?;
3261 let parent = resolve_bundle_permission(bundle, &body.parent_permission_digest)?;
3262 let time = times[i].map(|(_, time)| time);
3263 let (facts, selected_associations) = resolve_owner(time)?;
3264 let owner = facts.at(time.unwrap_or(0), &selected_associations);
3265 let token = verify_delegation_inner(d, parent, &owner, times[i].is_some())?;
3266 verified.push(token);
3267 }
3268 let initial = verified.first().ok_or(Reject::Canonical)?;
3269 for branch in &initial.body.branch_manifest {
3270 let limit = branch.limit.as_ref().ok_or(Reject::Canonical)?;
3271 let binding = bundle
3272 .genesis_authorities
3273 .iter()
3274 .find(|g| signed_genesis_digest(g).is_ok_and(|h| h == branch.genesis_authority_digest))
3275 .ok_or(Reject::Scope)?;
3276 let original = bundle
3277 .original_geneses
3278 .iter()
3279 .find(|o| native_id(o) == limit.genesis_digest)
3280 .ok_or(Reject::Scope)?;
3281 let g = binding.body.as_ref().ok_or(Reject::Canonical)?;
3282 let envelope = bundle
3283 .creator_authority_envelopes
3284 .iter()
3285 .find(|e| hash(&[e]) == g.creator_authority_envelope_digest)
3286 .ok_or(Reject::Scope)?;
3287 verify_native(original, "heddle-thread-genesis-v1")?;
3288 let signature = original
3289 .signatures
3290 .iter()
3291 .find(|s| s.public_key == g.creator_public_key)
3292 .ok_or(Reject::Scope)?;
3293 verify_genesis_authority(
3294 binding,
3295 initial,
3296 &limit.genesis_digest,
3297 &signature.signature,
3298 &hash(&[envelope]),
3299 )?;
3300 }
3301 for ((context, proof), signed) in resolved.iter().zip(&bundle.statements) {
3302 let s = signed.body.as_ref().ok_or(Reject::Canonical)?;
3303 if !verified.iter().any(|d| {
3304 d.body.identity.as_ref().is_some_and(|id| {
3305 s.spool_uuid == id.spool_uuid
3306 && s.spool_genesis_digest == id.spool_genesis_digest
3307 && s.owner_id == id.owner_id
3308 && s.owner_state_hash == id.owner_state_hash
3309 && s.ownership_transfer_sequence == id.ownership_transfer_sequence
3310 })
3311 }) {
3312 return Err(Reject::Scope);
3313 }
3314 match s.purpose {
3315 1 => {
3316 let time = s.observed_at_unix_millis / 1000;
3318 let (facts, selected_associations) = resolve_owner(Some(time))?;
3319 let owner = facts.at(time, &selected_associations);
3320 verify_delegation(
3321 &bundle.delegations[0],
3322 resolve_bundle_permission(bundle, &initial.body.parent_permission_digest)?,
3323 &owner,
3324 )?;
3325 let payload = bundle
3326 .genesis_witnesses
3327 .iter()
3328 .find(|p| canonical(*p).is_ok_and(|b| b == s.canonical_payload))
3329 .ok_or(Reject::Scope)?;
3330 let genesis = &payload
3331 .binding
3332 .as_ref()
3333 .and_then(|g| g.body.as_ref())
3334 .ok_or(Reject::Canonical)?
3335 .genesis_digest;
3336 if !admitted_geneses.insert(genesis.clone()) {
3337 return Err(Reject::Transition);
3338 }
3339 let publication = bundle
3340 .operations
3341 .iter()
3342 .zip(&selected_publications)
3343 .find(|(o, _)| {
3344 o.body
3345 .as_ref()
3346 .is_some_and(|o| &o.genesis_digest == genesis)
3347 })
3348 .and_then(|(_, s)| s.body.as_ref())
3349 .ok_or(Reject::Transition)?;
3350 check_import_genesis_publication_pair(initial, s, publication)?;
3351 verify_witness_payload(
3352 s,
3353 WitnessPayload::Genesis(
3354 bundle
3355 .genesis_witnesses
3356 .iter()
3357 .find(|p| canonical(*p).is_ok_and(|b| b == s.canonical_payload))
3358 .ok_or(Reject::Scope)?,
3359 ),
3360 )?;
3361 }
3362 2 => verify_witness_payload(
3363 s,
3364 WitnessPayload::Authority(
3365 bundle
3366 .authority_witnesses
3367 .iter()
3368 .find(|p| canonical(*p).is_ok_and(|b| b == s.canonical_payload))
3369 .ok_or(Reject::Scope)?,
3370 ),
3371 )?,
3372 4 => {
3373 let payload = bundle
3374 .landing_witnesses
3375 .iter()
3376 .find(|p| canonical(*p).is_ok_and(|b| b == s.canonical_payload))
3377 .ok_or(Reject::Scope)?;
3378 let (facts, selected_associations) =
3379 resolve_owner(Some(s.observed_at_unix_millis / 1000))?;
3380 let known = selected_associations
3381 .iter()
3382 .map(|(key, _)| key.clone())
3383 .collect::<Vec<_>>();
3384 verify_landing_key_roles(payload, &known, facts.forbidden_landing_keys)?;
3385 verify_witness_payload(s, WitnessPayload::Landing(payload))?;
3386 }
3387 3 => {
3388 let (_, operation, manifest, index) = publications
3389 .iter()
3390 .find(|(statement, _, _, _)| *statement == signed)
3391 .ok_or(Reject::Scope)?;
3392 let d = &bundle.delegations[*index];
3393 let body = d.body.as_ref().ok_or(Reject::Canonical)?;
3394 let time = s.observed_at_unix_millis / 1000;
3395 let (facts, selected_associations) = resolve_owner(Some(time))?;
3396 let owner = facts.at(time, &selected_associations);
3397 let delegation = verify_delegation(
3398 d,
3399 resolve_bundle_permission(bundle, &body.parent_permission_digest)?,
3400 &owner,
3401 )?;
3402 verify_publication(
3403 operation,
3404 &delegation,
3405 manifest,
3406 signed,
3407 set.as_ref().ok_or(Reject::Canonical)?,
3408 *proof,
3409 now_ms,
3410 )?;
3411 }
3412 _ => return Err(Reject::Version),
3413 }
3414 trust::recheck_context(
3415 context,
3416 set.as_ref().ok_or(Reject::Canonical)?,
3417 signed,
3418 now_ms,
3419 )?;
3420 }
3421 let mut progressive = ImportResultManifestV1 {
3423 slots: vec![],
3424 ..terminal.clone()
3425 };
3426 let mut order = 0;
3427 let mut observed = 0;
3428 let mut consumed = std::collections::BTreeMap::<(Vec<u8>, Vec<u8>), (u32, u64)>::new();
3429 let mut total_bytes = 0_u64;
3430 for o in &bundle.operations {
3431 let digest = signed_operation_digest(o)?;
3432 let slot = terminal
3433 .slots
3434 .iter()
3435 .find(|s| s.signed_operation_digest == digest)
3436 .ok_or(Reject::Scope)?;
3437 progressive.slots.push(slot.clone());
3438 progressive
3439 .slots
3440 .sort_by(|a, b| (&a.ref_name, a.slot_id).cmp(&(&b.ref_name, b.slot_id)));
3441 let payload = canonical(&publication_payload(o, &progressive)?)?;
3442 let s = bundle
3443 .statements
3444 .iter()
3445 .filter_map(|s| s.body.as_ref())
3446 .find(|s| s.purpose == 3 && s.canonical_payload == payload)
3447 .ok_or(Reject::Transition)?;
3448 if s.admission_order <= order || s.observed_at_unix_millis < observed {
3449 return Err(Reject::Transition);
3450 }
3451 order = s.admission_order;
3452 observed = s.observed_at_unix_millis;
3453 let digest = &o.body.as_ref().ok_or(Reject::Canonical)?.delegation_digest;
3454 let index = verified
3455 .iter()
3456 .position(|d| d.digest == *digest)
3457 .ok_or(Reject::Scope)?;
3458 let result_bytes = o.body.as_ref().ok_or(Reject::Canonical)?.result_bytes;
3459 total_bytes = total_bytes
3460 .checked_add(result_bytes)
3461 .ok_or(Reject::Bounds)?;
3462 let original_scope = initial.body.scope.as_ref().ok_or(Reject::Canonical)?;
3463 if total_bytes > original_scope.max_result_bytes
3464 || progressive.slots.len() > original_scope.max_operations as usize
3465 {
3466 return Err(Reject::Scope);
3467 }
3468 let delegation = &verified[index];
3469 let mut committed_before = progressive.clone();
3470 let operation_digest = signed_operation_digest(o)?;
3471 committed_before
3472 .slots
3473 .retain(|s| s.signed_operation_digest != operation_digest);
3474 check_import_publication_budget(o, delegation, &committed_before)?;
3475 let scope = delegation.body.scope.as_ref().ok_or(Reject::Canonical)?;
3476 let mut budgets = vec![(&delegation.digest, scope)];
3477 if let Some(parent) = &delegation.member {
3478 budgets.push((
3479 &delegation.body.parent_permission_digest,
3480 parent
3481 .body
3482 .as_ref()
3483 .and_then(|p| p.scope.as_ref())
3484 .ok_or(Reject::Canonical)?,
3485 ));
3486 }
3487 for (digest, scope) in budgets {
3488 let (operations, bytes) = consumed
3489 .entry((digest.clone(), delegation.body.logical_job_id.clone()))
3490 .or_default();
3491 *operations = operations.checked_add(1).ok_or(Reject::Bounds)?;
3492 *bytes = bytes.checked_add(result_bytes).ok_or(Reject::Bounds)?;
3493 if *operations > scope.max_operations || *bytes > scope.max_result_bytes {
3494 return Err(Reject::Scope);
3495 }
3496 }
3497 }
3498 if progressive != *terminal {
3499 return Err(Reject::Scope);
3500 }
3501 let mut history = snapshot.map_or_else(Vec::new, |s| s.accepted_history.clone());
3502 if let Some(old) = history.iter_mut().find(|b| {
3503 b.terminal_manifest
3504 .as_ref()
3505 .is_some_and(|m| m.logical_job_id == terminal.logical_job_id)
3506 && b.delegations
3507 .first()
3508 .and_then(|d| d.body.as_ref())
3509 .and_then(|d| d.identity.as_ref())
3510 .map(|id| &id.spool_uuid)
3511 == bundle
3512 .delegations
3513 .first()
3514 .and_then(|d| d.body.as_ref())
3515 .and_then(|d| d.identity.as_ref())
3516 .map(|id| &id.spool_uuid)
3517 }) {
3518 if !bundle.delegations.starts_with(&old.delegations)
3519 || !bundle.operations.starts_with(&old.operations)
3520 || bundle.owner_genesis != old.owner_genesis
3521 || !bundle
3522 .ownership_transfers
3523 .starts_with(&old.ownership_transfers)
3524 || !old
3525 .owner_histories
3526 .iter()
3527 .all(|v| bundle.owner_histories.contains(v))
3528 || bundle.member_permission != old.member_permission
3529 || !old
3530 .genesis_authorities
3531 .iter()
3532 .all(|v| bundle.genesis_authorities.contains(v))
3533 || !old
3534 .original_geneses
3535 .iter()
3536 .all(|v| bundle.original_geneses.contains(v))
3537 || !old
3538 .creator_authority_envelopes
3539 .iter()
3540 .all(|v| bundle.creator_authority_envelopes.contains(v))
3541 || !old.manifests.iter().all(|v| bundle.manifests.contains(v))
3542 || !old.statements.iter().all(|v| bundle.statements.contains(v))
3543 || !old.policies.iter().all(|v| bundle.policies.contains(v))
3544 || !old
3545 .genesis_witnesses
3546 .iter()
3547 .all(|v| bundle.genesis_witnesses.contains(v))
3548 || !old
3549 .authority_witnesses
3550 .iter()
3551 .all(|v| bundle.authority_witnesses.contains(v))
3552 || !old
3553 .landing_witnesses
3554 .iter()
3555 .all(|v| bundle.landing_witnesses.contains(v))
3556 || !old
3557 .foreign_dependencies
3558 .iter()
3559 .all(|v| bundle.foreign_dependencies.contains(v))
3560 {
3561 return Err(Reject::HighWater);
3562 }
3563 *old = bundle.clone();
3564 } else {
3565 history.push(bundle.clone());
3566 }
3567 let witnessed = times[0].is_some();
3568 let witnessed_prefix = usize::from(witnessed);
3569 let mut persisted = snapshot.cloned();
3570 if witnessed_prefix > 0 || new_set {
3571 let mut persisted_associations =
3572 snapshot.map_or_else(Vec::new, |s| s.job_associations.clone());
3573 if witnessed_prefix > 0 {
3574 for d in &bundle.delegations[..witnessed_prefix] {
3575 let b = d.body.as_ref().ok_or(Reject::Canonical)?;
3576 if !persisted_associations
3577 .iter()
3578 .any(|(key, _)| *key == b.job_public_key)
3579 {
3580 persisted_associations
3581 .push((b.job_public_key.clone(), b.logical_job_id.clone()));
3582 }
3583 }
3584 }
3585 persisted = Some(ImportWitnessSnapshot {
3586 root: pin.clone(),
3587 witness_set: carried.ok_or(Reject::Canonical)?.clone(),
3588 clock_floor_unix_millis: now_ms,
3589 job_associations: persisted_associations,
3590 accepted_history: if witnessed_prefix > 0 {
3591 history
3592 } else {
3593 snapshot.map_or_else(Vec::new, |s| s.accepted_history.clone())
3594 },
3595 });
3596 }
3597 Ok(VerifiedImportBundleWitnesses {
3598 evidence: if witnessed {
3599 ImportBundleEvidence::Witnessed
3600 } else {
3601 ImportBundleEvidence::Recovery
3602 },
3603 snapshot_advanced: persisted.as_ref() != snapshot,
3604 owner_check_time_unix_seconds: times[0].map(|(_, time)| time),
3605 accepted_history: terminal.clone(),
3606 snapshot: persisted,
3607 })
3608}
3609
3610pub struct ImportControlCaller<'a> {
3612 pub authenticated_pop: bool,
3613 pub destination_writer: bool,
3614 pub caller_account: &'a str,
3615 pub connection_owner_account: Option<&'a str>,
3616 pub authorized_source: Option<&'a ImportSourceSelectionV1>,
3617 pub exact_grants_current: bool,
3618 pub selected_commits_available: bool,
3619}
3620#[derive(Clone, Copy)]
3621pub enum ImportControlAction {
3622 Cancel,
3623 Retry,
3624}
3625
3626pub fn check_import_control_caller(
3629 action: ImportControlAction,
3630 retained: &ImportSourceSelectionV1,
3631 scope: &ImportPermissionScopeV1,
3632 caller: &ImportControlCaller<'_>,
3633) -> Result<(), Reject> {
3634 if !caller.authenticated_pop || !caller.destination_writer || caller.caller_account.is_empty() {
3635 return Err(Reject::Scope);
3636 }
3637 if matches!(action, ImportControlAction::Cancel) {
3638 return Ok(());
3639 }
3640 validate_retained_import_source(retained, scope)?;
3641 if retained.connection.is_some()
3642 && (caller.connection_owner_account != Some(caller.caller_account)
3643 || caller.authorized_source != Some(retained)
3644 || !caller.exact_grants_current)
3645 {
3646 return Err(Reject::SourceSelection);
3647 }
3648 if !caller.selected_commits_available {
3649 return Err(Reject::SourceSelection);
3650 }
3651 Ok(())
3652}
3653
3654fn validate_retry_availability(
3655 response: &GetImportJobStateResponse,
3656 destination: &SpoolRef,
3657) -> Result<(), Reject> {
3658 use get_import_job_state_response::RetryAvailability;
3659 match response
3660 .retry_availability
3661 .as_ref()
3662 .ok_or(Reject::Canonical)?
3663 {
3664 RetryAvailability::EligibleRetryTarget(target) => {
3665 if response.status != 1 {
3666 return Err(Reject::Canonical);
3667 }
3668 let operation = target.operation_ref.as_ref().ok_or(Reject::Canonical)?;
3669 let raw = hex::decode(operation.id.replace('-', "")).map_err(|_| Reject::Canonical)?;
3670 if initial_operation_id(&raw, false)? != operation.id {
3671 return Err(Reject::Canonical);
3672 }
3673 if operation.spool.as_ref() != Some(destination) {
3674 return Err(Reject::Scope);
3675 }
3676 if target.operation_version.is_empty() || target.operation_version.len() > 256 {
3677 return Err(Reject::Bounds);
3678 }
3679 }
3680 RetryAvailability::RetryUnavailable(reason) => {
3681 if !matches!(
3682 (response.status, *reason),
3683 (1, 1 | 2 | 6) | (2, 3) | (3, 4) | (4, 5) | (5, 7)
3684 ) {
3685 return Err(Reject::Canonical);
3686 }
3687 }
3688 }
3689 Ok(())
3690}
3691
3692pub fn validate_retry_state_response(
3694 request: &GetImportJobStateRequest,
3695 response: &GetImportJobStateResponse,
3696) -> Result<(), Reject> {
3697 validate_job_state_response(request, response)?;
3698 validate_retry_availability(response, request.destination.as_ref().ok_or(Reject::Scope)?)
3699}
3700
3701pub struct ImportRetryAdmission<'a> {
3704 pub read: &'a GetImportJobStateResponse,
3705 pub original: &'a OperationRecord,
3706 pub retry_lineage_id: &'a [u8],
3707 pub logical_job_terminal: bool,
3708 pub retained_source: &'a ImportSourceSelectionV1,
3709 pub committed_manifest: &'a ImportResultManifestV1,
3710 pub now_unix_seconds: i64,
3711}
3712
3713pub fn check_retry_admission(
3716 request: &RetryImportSourceRequest,
3717 context: &ImportRetryAdmission<'_>,
3718 active: &VerifiedImportDelegation,
3719 caller: &ImportControlCaller<'_>,
3720) -> Result<(), Reject> {
3721 use get_import_job_state_response::RetryAvailability;
3722 let destination = request
3723 .original_operation
3724 .as_ref()
3725 .and_then(|r| r.spool.clone());
3726 let read_request = GetImportJobStateRequest {
3727 destination,
3728 logical_job_id: request.logical_job_id.clone(),
3729 };
3730 validate_retry_state_response(&read_request, context.read)?;
3731 if request.client_operation_id.is_empty() || request.client_operation_id.len() > 128 {
3732 return Err(Reject::Canonical);
3733 }
3734 if context.logical_job_terminal || context.read.status != 1 {
3735 return Err(Reject::Revoked);
3736 }
3737 let Some(RetryAvailability::EligibleRetryTarget(target)) = &context.read.retry_availability
3738 else {
3739 return Err(Reject::StaleContext);
3740 };
3741 let target_ref = target.operation_ref.as_ref().ok_or(Reject::Canonical)?;
3742 let original_ref = context.original.r#ref.as_ref().ok_or(Reject::Scope)?;
3743 if request.original_operation.as_ref() != Some(original_ref)
3744 || original_ref.spool != target_ref.spool
3745 || original_ref.id != target_ref.id
3746 || context.retry_lineage_id != active.body.retry_lineage_id
3747 || import_job_state_request_from_operation(context.original)?.as_ref()
3748 != Some(&read_request)
3749 {
3750 return Err(Reject::Scope);
3751 }
3752 if request.expected_operation_version != target.operation_version
3753 || context.original.version != target.operation_version
3754 || context.original.superseded_by.is_some()
3755 || !matches!(context.original.state, 4 | 5)
3756 {
3757 return Err(Reject::StaleContext);
3758 }
3759 let state = context.read;
3760 if state.active_delegation_digest != active.digest
3761 || state.active_cancellation_id != active.body.cancellation_id
3762 {
3763 return Err(Reject::StaleContext);
3764 }
3765 check_job_fence(
3766 &request.logical_job_id,
3767 &request.active_delegation_digest,
3768 request.expected_authority_epoch,
3769 active,
3770 state.authority_epoch,
3771 )?;
3772 interval(
3773 active.body.not_before_unix_seconds,
3774 active.body.expires_at_unix_seconds,
3775 context.now_unix_seconds,
3776 )?;
3777 let manifest = context.committed_manifest;
3778 validate_manifest(manifest)?;
3779 if manifest.logical_job_id != active.body.logical_job_id
3780 || manifest.retry_lineage_id != active.body.retry_lineage_id
3781 {
3782 return Err(Reject::Scope);
3783 }
3784 let scope = active.body.scope.as_ref().ok_or(Reject::Canonical)?;
3785 if !scope.branches.iter().any(|b| {
3786 !manifest
3787 .slots
3788 .iter()
3789 .any(|s| s.ref_name == b.ref_name && s.slot_id == b.slot_id)
3790 }) {
3791 return Err(Reject::StaleContext);
3792 }
3793 check_import_control_caller(
3794 ImportControlAction::Retry,
3795 context.retained_source,
3796 scope,
3797 caller,
3798 )
3799}
3800
3801pub fn validate_retry_response(
3805 request: &RetryImportSourceRequest,
3806 response: &MutationResponse,
3807 prior_attempt_ids: &[String],
3808) -> Result<(), Reject> {
3809 let receipt = response.receipt.as_ref().ok_or(Reject::PendingOperation)?;
3810 let Some(mutation_receipt::Outcome::PendingOperation(operation)) = &receipt.outcome else {
3811 return Err(Reject::PendingOperation);
3812 };
3813 let original = request
3814 .original_operation
3815 .as_ref()
3816 .ok_or(Reject::PendingOperation)?;
3817 let raw = hex::decode(operation.id.replace('-', "")).map_err(|_| Reject::PendingOperation)?;
3818 if initial_operation_id(&raw, false).map_or(true, |id| id != operation.id)
3819 || request.client_operation_id.is_empty()
3820 || receipt.client_operation_id != request.client_operation_id
3821 || original.spool.is_none()
3822 || operation.spool != original.spool
3823 || operation.id == original.id
3824 || operation.id == request.client_operation_id
3825 || prior_attempt_ids.contains(&operation.id)
3826 {
3827 return Err(Reject::PendingOperation);
3828 }
3829 Ok(())
3830}
3831
3832pub fn check_retry_replay(request_bytes: &[u8], stored_bytes: &[u8]) -> Result<(), Reject> {
3834 if request_bytes.len() > 2 * MAX_BUNDLE_BYTES {
3835 return Err(Reject::Bounds);
3836 }
3837 if request_bytes != stored_bytes {
3838 return Err(Reject::OperationIdReused);
3839 }
3840 Ok(())
3841}
3842
3843pub fn check_import_publication_budget(
3846 signed: &SignedDelegatedImportOperationV1,
3847 active: &VerifiedImportDelegation,
3848 committed_before: &ImportResultManifestV1,
3849) -> Result<(), Reject> {
3850 verify_operation(signed, active)?;
3851 let scope = active.body.scope.as_ref().ok_or(Reject::Canonical)?;
3852 if committed_before.logical_job_id != active.body.logical_job_id
3853 || committed_before.retry_lineage_id != active.body.retry_lineage_id
3854 {
3855 return Err(Reject::Scope);
3856 }
3857 let remaining = remaining_import_scope(scope, committed_before)?;
3858 if check_slot_replay(committed_before, signed)? {
3859 return Ok(());
3860 }
3861 let o = signed.body.as_ref().ok_or(Reject::Canonical)?;
3862 if remaining.max_operations == 0 || o.result_bytes > remaining.max_result_bytes {
3863 return Err(Reject::Scope);
3864 }
3865 Ok(())
3866}