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