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