Skip to main content

heddle_api/
import_authority.rs

1//! HYBRID canonical import authority. The existing heddle capability verifier
2//! supplies independently selected owner history/current permission. This
3//! module verifies the NEW typed parent grant and its bounded job child;
4//! it neither enrolls incoming roots nor substitutes for owner verification.
5use 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
125/// Conservative canonical URL grammar for v1: lowercase DNS HTTPS host, no
126/// userinfo/query/fragment/port/escapes, ASCII unreserved path segments. No
127/// implicit URL normalization is performed by a signing implementation.
128pub 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}
185// Structural validation never substitutes a synthetic historical clock.
186fn 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
243/// Check independently observed OID knowledge before signing/preparation.
244/// Empty/unknown is None, never an implicit observe-mode selection or consent.
245pub 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
259/// Converter option octets are selected verbatim from authenticated discovery.
260pub 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
273/// Resolve explicit custody using a CURRENT authenticated connection provider.
274/// Public URLs never select an adapter by domain. Network/grant checks remain host gates.
275pub 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
322/// Advisory Git size in KiB. UNKNOWN is explicit and must carry zero.
323pub 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
331/// Preserve selected identity exactly. Redirect/SSRF and current grants are host gates.
332pub 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
359/// Structural binding to signed provider/URL; durable custody association is a host invariant.
360fn 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
382/// Discovery may report unknown. Preparing/signing requires known=true; no SHA-1 fallback.
383pub 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
420/// Compare independent repository discovery before signing, including known-OID pinning.
421pub 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        // An authenticated retained pin selects its original commit, not today's head.
449        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
461/// Host resolves the selector anew and checks current grants and selected-commit
462/// availability before issuing a reservation, including renewals. Retained state
463/// MUST come from an authenticated job-state read or receiver-owned durable state;
464/// the opaque predecessor binds its independently verified authority to that CAS.
465pub 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        // An activated job retains its Prepare byte maximum; it can only narrow.
530        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
612/// Host-side negotiation with a CURRENT authenticated configuration and CAS.
613/// Only an empty destination token is filled; all other choices survive exactly.
614pub 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
663/// Independently retained refs/slots for one non-terminal logical job. Include
664/// Prepared jobs and all original refs until logical termination, not just the
665/// current certificate's remaining scope. Hosts provide the complete inventory
666/// under the same transaction that reserves/activates; this helper stores nothing.
667pub struct ImportSpoolReservation<'a> {
668    pub spool_uuid: &'a [u8],
669    pub logical_job_id: &'a [u8],
670    pub branches: &'a [ImportBranchLimitV1],
671}
672
673/// Per-spool full refs and (full ref, slot_id) keys are exclusive across jobs,
674/// regardless of provider/source. Same-job renewals retain exact branch identity.
675/// Run after scope negotiation and before atomically reserving every ref/slot.
676pub 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
720/// Current destination configuration/ownership/policy CAS at activation. Other
721/// imports do not advance this token. Exact stored replay is resolved first.
722pub 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
734/// Browser-side comparison before signing. A host cannot silently negotiate.
735pub 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
784/// Enforce inclusive decoded protobuf sizes before semantic admission.
785/// Transport must also bound the whole uncompressed payload before decoding.
786pub 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
801/// Source association/provider comes from the host's authenticated resolver,
802/// never a projection hint. Native base decoding/identity and current grants
803/// remain host gates; this validates carrier, bounds and exact originals.
804pub 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    // A frozen pin names the selected commit, even if the branch head moves.
892    // OBSERVE still cannot hide a currently known selected OID by clearing hints.
893    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    // Current converter/options/budget support is also rechecked at activation.
907    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    // Arrays are ordered exactly like the signed branch manifest, no second
916    // association-by-name payload and no incoming regenerated branch originals.
917    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
940/// This closed route has no initial-submission or authority-attachment role.
941pub fn validate_import_source(_: &ImportSourceRequest) -> Result<(), Reject> {
942    Err(Reject::ImportSourceRequiresCommit)
943}
944
945/// Complete initial validation, excluding native model/owner history, live
946/// source access and transaction checks owned by the hosted implementation.
947pub 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
1009/// Use a durable caller-scoped idempotency row BEFORE rechecking expired job
1010/// authority. Host stores the original request/receipt atomically with activation.
1011pub 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
1072/// Public context from the existing owner/keyring verifier, not from fields in
1073/// the incoming bundle. now is receiver/host time for new work or independently
1074/// verified witness observation time for retained history, NEVER author time.
1075pub struct ImportOwnerExpectation<'a> {
1076    pub identity: &'a ImportIdentityV1,
1077    pub owner_public_key: &'a [u8],
1078    pub owner_chain_digest: &'a [u8],
1079    /// From the independently verified EFFECTIVE owner state at `now`, not
1080    /// the immutable root. Accepted claim/deferral clearing is unbounded.
1081    pub authority_expires_at_seconds: i64,
1082    pub now_unix_seconds: i64,
1083    pub forbidden_job_keys: &'a [Vec<u8>], // Every user/root/witness key, including tombstones.
1084    pub known_job_associations: &'a [(Vec<u8>, Vec<u8>)], // key -> logical job.
1085}
1086/// Independently verified owner facts for retained bundle verification. Historical
1087/// check times are selected internally from authenticated receipts, never callers.
1088#[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    /// Effective interval of these independently authenticated owner facts.
1095    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}
1113/// Inputs MUST come from the selected, independently verified effective state.
1114/// Historical verification selects the state at its authenticated observation.
1115pub 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}
1297/// Frozen canonical projection. The browser completes only the fields absent
1298/// here; it cannot normalize or reduce even an otherwise authorized scope.
1299pub 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
1316/// Validate Commit against the HOST-STORED preparation and independently
1317/// selected current authority. Native genesis/envelope verification, online
1318/// revocation, custody uniqueness and transactional activation remain host gates.
1319/// Retained renewal genesis bindings require their original accepted context.
1320pub 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}
1329/// Browser signing preflight only: no execution or admission token is returned.
1330pub 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    // Use i128 for host arithmetic so extreme advertised uint64 bounds cannot
1363    // wrap. Skew permits a future not-before, never grace after expiry.
1364    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    // Future not-before within skew can be committed, but verify_new_operation
1401    // still refuses execution until that exact signed time. Parent/owner expiry
1402    // and containment are checked without extending them by skew.
1403    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}
1484/// Structural signature + scoped operation only. This does not establish
1485/// publication, current policy, cancellation, leases or conversion correctness.
1486pub 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}
1585/// CAS state MUST be receiver-owned and held under the publication/renewal
1586/// transaction fence. Returns replacement only after all remaining-slot checks.
1587pub 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        // Old certificates may already omit previously committed slots. Only
1655        // charge all committed slots against the original logical-job budgets;
1656        // the host's initial manifest/limits remain durable across renewals.
1657        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}
1733/// Persistent unique slot identity: logical job/ref/slot. Exact replay returns
1734/// the old receipt/manifest, never another publication or fresh witness.
1735pub 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}
1762/// Verify exact committed publication in addition to job signature/scope. The
1763/// owner/keyring verifier must resolve the statement's accepted state/order;
1764/// witness signature alone cannot establish that user authority or disclosure.
1765pub 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
1838/// Producer-owned logical-job/lease fence for RetryImportSource and final
1839/// publication. The physical retry row never supplies a new logical identity.
1840pub 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
1887/// Typed adapter boundary: an authentic unrelated capability/online role must
1888/// never be selected as the parent of an import certificate.
1889pub 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}
1904/// Hybrid dispatch has no legacy execution arm, even for an authentic witness.
1905pub 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}
1961/// Recompute transport commitments from exact native evidence. The caller's
1962/// native verifier additionally authenticates manifest membership, receipt
1963/// subjects/bases, accepting authority and canonical native encoding.
1964pub 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// These readers extract only the native commitment selectors. Full native
2049// canonicality, model validity, membership and authority remain the native gate.
2050#[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;
2220/// Caller constructs this from independently verified native originals and
2221/// accepted owner/policy/landing context. This matching layer verifies original
2222/// signatures and exact payload commitments separately from witness trust; it
2223/// does not replace native causal, authority or landing-model verification.
2224pub 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}
2464/// Completeness and digest addressing only. Trust/signature verification still
2465/// uses independently selected owner contexts at each witnessed historical time.
2466fn 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            // Every branch needs its original admission, not merely its
2552            // binding. Proof-only retirement lookup cannot recover a payload.
2553            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}
2650/// Reference completeness only. Native verification must authenticate every
2651/// selected policy, its owner context and the receipt before using its time.
2652pub(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}
2692/// Prepare response is a coherent proposal, never authority. Caller must compare
2693/// the IDs, predecessor and manifest before asking its device to sign renewal.
2694pub 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
2716/// Caller-generated non-nil UUID, reserved as the first physical operation ID.
2717/// Occupancy is checked under the host's reservation/activation transaction.
2718pub 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
2737/// Evaluate inside the cancellation transaction, after caller-scoped replay lookup.
2738/// The selector names only the ACTIVE delegation. Parent revocation is independent.
2739pub 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}
2768/// Exact replay acknowledges the stored cancellation without advancing the epoch.
2769pub 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}
2778/// Check both independent revocation selectors, regardless of Cancel's selector.
2779pub 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
2799/// Compute the exact remaining total and slots from an authenticated manifest.
2800/// An empty result is completed work and cannot pass new-scope validation.
2801pub 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/// Signature/scope-verified recovery evidence. This is neither currently
2879/// executable authority nor evidence of historical admission. No receipt needed.
2880///
2881/// ```compile_fail
2882/// use heddle_api::import_authority::{VerifiedImportRenewalPredecessor, verify_new_operation};
2883/// use heddle_api::heddle::api::v1alpha2::{SignedDelegatedImportOperationV1, ImportResultManifestV1};
2884/// fn cannot_execute(op: &SignedDelegatedImportOperationV1, old: &VerifiedImportRenewalPredecessor, manifest: &ImportResultManifestV1) {
2885///     verify_new_operation(op, old, 0, manifest);
2886/// }
2887/// ```
2888#[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}
2909/// State MUST come from authenticated job-state read / Prepare or receiver-owned durable state.
2910/// Expected owner context independently verifies the predecessor's selected history.
2911pub 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}
2928/// Current replacement authority remains mandatory. Activation MUST recheck the
2929/// durable active predecessor/epoch/manifest and terminal state atomically.
2930pub 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
2951/// Validate discovery metadata only; a well-shaped selector grants no authority.
2952pub 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
2962/// Project a visible operation's durable HYBRID association into the existing
2963/// writer-only state read. Missing/unknown subject or selector is unavailable.
2964/// Malformed present metadata is rejected; never substitute an attempt ID.
2965/// This validates shape, not operation visibility, writer access or signatures.
2966pub 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
2991/// Finite destination-writer read; transport authentication/authorization belongs
2992/// to the generated RPC contract. Validate before any storage lookup.
2993pub 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    // An empty authenticated snapshot is recovery evidence, not proof of past
3040    // admission. Nonempty snapshots retain the full public export closure.
3041    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
3083/// Supply the authenticated read, never an incoming untrusted state assertion.
3084pub 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
3124/// A publication/renewal race requires recomputation and a new exact Prepare
3125/// before signing. Compare the complete snapshot, including signed predecessor.
3126pub 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
3149/// Composition/reference validation only. Independently verify owner histories,
3150/// policies, accepted admissions/publications and the current owner head before
3151/// using them. The authenticated read supplies exact retained evidence.
3152/// logical_job_terminal is independently selected under the host transaction.
3153pub 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    // Normalize only the four permitted additions/current selectors. Every
3205    // other retained byte (including original parents and accepted history)
3206    // must stay exact. A pending candidate can never enter accepted arrays.
3207    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
3282/// Full typed renewal checks with separate independently authenticated historical
3283/// predecessor and current replacement contexts. Host transaction gates remain.
3284pub 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
3330/// Frozen protobuf request replay is distinct from HYBRID signed-body encoding.
3331/// Hosts compare retained raw bytes before CAS/expiry checks, in caller scope.
3332pub 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/// Independently installed descriptor pin; never learned from a response.
3343#[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/// Receiver-owned durable data, committed atomically with accepted history.
3351/// Do not deserialize this from the incoming proof or lower its clock floor.
3352#[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/// Witnessed requires an authenticated observation for every accepted delegation.
3361/// Recovery can retain witnessed history but grants no admission for the active tail.
3362#[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}
3375/// Verify retained evidence after authenticating the state RPC. Owner contexts
3376/// are resolved independently at each authenticated time. The verifier derives times
3377/// from the first authenticated publication, else initial admission.
3378/// Unwitnessed certificates are time-free recovery only.
3379/// The mandatory hook verifies the selected policy chain and owner/native context
3380/// at each authenticated statement time (heddle capability-verifier / WASM).
3381/// With no statements it receives None and verifies time-free policy closure.
3382/// No result is returned on failure. Recovery without a witnessed prefix or new
3383/// authenticated set returns the input snapshot unchanged. A witnessed prefix
3384/// extends rollback protection; an unwitnessed tail is never persisted. A new set
3385/// can additionally advance set trust and clock. None remains None when no
3386/// set is carried and no input exists. Persist under the trust lock; inspect
3387/// snapshot_advanced for actual durable change, independently of evidence.
3388pub 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    // Authenticate statements before using their observation times or policies.
3478    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        // Time-free policy/owner/native closure; this asserts no admission event.
3502        verify_policy(bundle, None)?;
3503    }
3504    // Derive owner-selection times only from authenticated receipts. Prefer the
3505    // first publication for each delegation; initial admissions stand alone too.
3506    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        // First means accepted publication order, independent of array order.
3535        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    // Resolve only after authentication. Facts must cover the selected time;
3550    // current claimed state cannot stand in for an earlier deferred state.
3551    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        // Also keep newly supplied known job keys disjoint from witness roles.
3582        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                // Every original admission must independently satisfy [N,E).
3664                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    // Operations are accepted publication order, while manifest slots are sorted.
3738    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
3974/// These facts MUST come from current authenticated host state, never request
3975/// claims. Retry/Renew fetches use this exact authorized source association.
3976/// Classify independently retained original authority and actual native admission.
3977/// A later certificate cannot reopen the originals' exclusive admission window.
3978pub 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
3996/// Independently authenticated current caller facts for source custody checks.
3997pub 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
4013/// Destination control is independent of source custody. Cancel never fetches.
4014/// Run at Retry admission / Renew activation after checking signed authority.
4015pub 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
4041/// Current host facts read under the same fence as the writer-only job response.
4042pub 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}
4048/// Compute caller-relative advice; NEVER use this result as admission authority.
4049/// Cancel needs no source relationship. Only the caller's relationship is exposed.
4050pub 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}
4130/// Validate current read disclosures. Historical frozen carriers are validated
4131/// separately by validate_job_state_response / validate_retry_state_response.
4132pub 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
4188/// Validate the additive writer disclosure separately from historical recovery
4189/// carriers, whose frozen alpha.30 bytes have no retry_availability field.
4190pub 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
4198/// Receiver-owned facts locked together with job/authority/source state. The
4199/// operation-lineage association is durable host state, not an ID inference.
4200pub 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
4209/// Run after caller-scoped frozen replay lookup, under ONE admission transaction.
4210/// Full owner/policy/revocation/lease checks and allocation remain host gates.
4211pub 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
4305/// Check a host-generated UUID against the COMPLETE durable attempt set, including
4306/// the first lineage UUID. Persist the allocation, both direct links and receipt
4307/// atomically; do not generate a UUID from the request idempotency key.
4308pub 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
4336/// Authority-only acknowledgement. Empty Applied.resulting_versions is valid;
4337/// scheduling is forbidden and the activated epoch/certificate comes from state.
4338pub 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
4352/// Exact frozen Retry replay returns the stored receipt without allocating.
4353pub fn check_retry_replay(request_bytes: &[u8], stored_bytes: &[u8]) -> Result<(), Reject> {
4354    check_renew_replay(request_bytes, stored_bytes)
4355}
4356
4357/// Under the publication transaction, use the complete cumulative manifest of
4358/// this logical job. Renewals MUST already be narrowed to predecessor remainder.
4359pub 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
4380/// Independently retained original certificate and native-admission fact.
4381pub struct ImportOriginalAdmission<'a> {
4382    pub original: &'a SignedImportJobDelegationV1,
4383    pub admitted: bool,
4384    pub now_unix_seconds: i64,
4385}
4386/// A missed original window requires a new job, never another renewal or retry.
4387pub 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}
4395/// Hosts run this under the same inventory lock before reserving a new job.
4396/// Only this job's reservations in this spool are released.
4397pub 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}