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