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