Skip to main content

type_bridge_contract/
query_remote.rs

1//! Versioned fail-closed wire envelopes for remote query execution.
2//!
3//! One validated plan/result contract serves direct and server execution:
4//! the request carries the exact canonical plan bytes plus the invocation
5//! (operation, input rows, caller budgets), the exact executor advertisement,
6//! an absolute bounded expiry, and a caller nonce; the response binds that
7//! nonce and the whole request fingerprint so replayed or foreign evidence is
8//! rejected before any host object is constructed. Envelope formats are
9//! versioned independently of the plan format.
10
11use serde::de::{DeserializeSeed, IgnoredAny, MapAccess, SeqAccess, Visitor};
12use serde::{Deserialize, Deserializer, Serialize};
13use sha2::{Digest, Sha256};
14
15use crate::codec::{
16    from_canonical_json_with_limits, to_canonical_json, to_canonical_json_with_limits,
17};
18use crate::diagnostic::{Diagnostic, DiagnosticCategory, DiagnosticCode};
19use crate::fingerprint::{
20    CanonicalizationVersion, Fingerprint, FingerprintDigest, FingerprintDomain,
21};
22use crate::id::TypeId;
23use crate::limits::{
24    MAX_CANONICAL_COLLECTION_LEN, MAX_REMOTE_ENVELOPE_BYTES, REMOTE_ENVELOPE_CODEC_LIMITS,
25    REMOTE_REQUEST_CODEC_LIMITS,
26};
27use crate::query_plan::{
28    InputRow, QueryInvocation, QueryOperation, QueryPlan, QueryPlanFingerprint, decode_query_plan,
29};
30use crate::value::CanonicalValue;
31
32/// The exact wire discriminator for first-format remote requests.
33pub const QUERY_REMOTE_REQUEST_FORMAT_V1: &str = "typebridge.query-remote-request/v1";
34/// The exact wire discriminator for first-format remote responses.
35pub const QUERY_REMOTE_RESPONSE_FORMAT_V1: &str = "typebridge.query-remote-response/v1";
36/// The exact wire discriminator for first-format remote failures.
37pub const QUERY_REMOTE_FAILURE_FORMAT_V1: &str = "typebridge.query-remote-failure/v1";
38/// The authenticated outer reply format used for both successes and failures.
39pub const QUERY_REMOTE_SIGNED_REPLY_FORMAT_V1: &str = "typebridge.query-remote-signed-reply/v1";
40/// Domain separating Ed25519 reply signatures from every other signed value.
41pub const QUERY_REMOTE_REPLY_SIGNATURE_DOMAIN: &str = "typebridge.query.remote-reply-signature/v1";
42/// Domain separating deterministic reply-signing key identifiers.
43pub const QUERY_REMOTE_REPLY_KEY_ID_DOMAIN: &str = "typebridge.query.remote-reply-key-id/v1";
44
45const NONCE_MIN_BYTES: usize = 16;
46const NONCE_MAX_BYTES: usize = 128;
47const EXECUTOR_COMPONENT_MIN_BYTES: usize = 16;
48const EXECUTOR_COMPONENT_MAX_BYTES: usize = 128;
49/// Default lifetime for a remote request with no explicit caller deadline.
50///
51/// This matches the standalone executor's mandatory execution timeout and
52/// keeps one omitted field from occupying replay capacity for hours.
53pub const DEFAULT_REMOTE_DEADLINE_MS: u64 = 30 * 1_000;
54/// Longest caller deadline admitted by the first remote format: five minutes.
55pub const MAX_REMOTE_DEADLINE_MS: u64 = 5 * 60 * 1_000;
56/// Maximum positive client/server wall-clock skew admitted by remote preflight.
57///
58/// The request lifetime is still bounded to [`MAX_REMOTE_DEADLINE_MS`]
59/// between its fingerprint-bound preparation and expiry timestamps. This allowance only
60/// prevents a client clock up to one minute ahead from being rejected as a
61/// forged future request.
62pub const MAX_REMOTE_CLOCK_SKEW_MS: u64 = 60 * 1_000;
63
64/// Caller execution budgets carried with one remote invocation.
65///
66/// Budgets tighten provider and session ceilings; they never raise them.
67#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
68#[serde(deny_unknown_fields)]
69pub struct RemoteLimits {
70    /// Optional wall-clock deadline in milliseconds.
71    pub deadline_ms: Option<u64>,
72    /// Maximum successful typed-response bytes the caller accepts.
73    ///
74    /// Authenticated structured failures are control-plane evidence bounded by
75    /// [`MAX_REMOTE_ENVELOPE_BYTES`], not by this success-data budget. This lets
76    /// an executor report that even the smallest success cannot fit when the
77    /// caller deliberately supplies a zero or otherwise tiny budget.
78    pub max_bytes: u64,
79    /// Maximum answer items the caller accepts.
80    pub max_items: u64,
81    /// Maximum aggregate list members the caller accepts across documents.
82    pub max_collection_members: u64,
83}
84
85#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
86#[serde(rename_all = "snake_case")]
87enum RemoteOperation {
88    Rows,
89    Count,
90    Exists,
91}
92
93impl RemoteOperation {
94    const fn from_operation(operation: QueryOperation) -> Self {
95        match operation {
96            QueryOperation::Rows => Self::Rows,
97            QueryOperation::Count => Self::Count,
98            QueryOperation::Exists => Self::Exists,
99        }
100    }
101
102    const fn operation(self) -> QueryOperation {
103        match self {
104            Self::Rows => QueryOperation::Rows,
105            Self::Count => QueryOperation::Count,
106            Self::Exists => QueryOperation::Exists,
107        }
108    }
109}
110
111/// One complete remote invocation of a reusable validated plan.
112#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
113#[serde(deny_unknown_fields)]
114pub struct RemoteQueryRequest {
115    // Pre-release wire ledger: before any 2.0.0 artifact shipped, /v1 gained
116    // an exact capability-advertisement fingerprint and absolute bounded
117    // preparation/expiry timestamps. Relative deadline_ms remains the caller
118    // API, but it is resolved once during preparation; executors never restart
119    // or replay the full relative window. Before any 2.0.0 artifact shipped,
120    // the omitted-deadline lifetime became 30 seconds, the explicit ceiling
121    // became five minutes, and decode began requiring the absolute timestamps
122    // to equal that fingerprint-bound declared/default lifetime exactly.
123    advertisement: String,
124    expires_at_unix_ms: u64,
125    format: String,
126    limits: RemoteLimits,
127    nonce: String,
128    operation: RemoteOperation,
129    // Pre-release wire ledger: the /v1 request embedded the plan as one
130    // JSON string until 2.0.0 shipped, which capped remote plans at the
131    // 1 MiB per-string ceiling instead of the 16 MiB document limit the
132    // plan contract states. The plan now embeds as the canonical JSON
133    // object itself, so local and remote share one plan size limit.
134    plan: serde_json::Value,
135    prepared_at_unix_ms: u64,
136    rows: Vec<Vec<Option<CanonicalValue>>>,
137}
138
139impl RemoteQueryRequest {
140    /// Bind one plan invocation and caller budgets into a request envelope.
141    pub fn new(
142        plan: &QueryPlan,
143        invocation: &QueryInvocation,
144        advertisement: &RemoteCapabilities,
145        limits: RemoteLimits,
146        nonce: impl Into<String>,
147        prepared_at_unix_ms: u64,
148    ) -> Result<Self, Diagnostic> {
149        if !invocation.binds(plan)? {
150            return Err(envelope_failure(
151                DiagnosticCategory::Integrity,
152                "query_remote_invocation_plan_mismatch",
153                "invocation does not bind the exact plan fingerprint",
154            ));
155        }
156        let nonce = nonce.into();
157        check_nonce(&nonce)?;
158        validate_remote_limits(limits)?;
159        let expires_at_unix_ms = prepared_at_unix_ms
160            .checked_add(limits.deadline_ms.unwrap_or(DEFAULT_REMOTE_DEADLINE_MS))
161            .ok_or_else(remote_time_invalid)?;
162        validate_remote_time_shape(prepared_at_unix_ms, expires_at_unix_ms, limits)?;
163        let advertisement = advertisement.fingerprint()?.digest_hex();
164        // Canonical bytes prove the plan encodes within contract limits;
165        // the parsed tree of those exact bytes is what the envelope embeds.
166        let plan_value = serde_json::from_slice::<serde_json::Value>(&plan.canonical_bytes()?)
167            .map_err(|_| {
168                envelope_failure(
169                    DiagnosticCategory::Integrity,
170                    "query_remote_plan_unencodable",
171                    "the plan cannot be embedded as canonical JSON",
172                )
173            })?;
174        Ok(Self {
175            advertisement,
176            expires_at_unix_ms,
177            format: QUERY_REMOTE_REQUEST_FORMAT_V1.to_owned(),
178            limits,
179            nonce,
180            operation: RemoteOperation::from_operation(invocation.operation()),
181            plan: plan_value,
182            prepared_at_unix_ms,
183            rows: invocation
184                .inputs()
185                .iter()
186                .map(|row| row.values().to_vec())
187                .collect(),
188        })
189    }
190
191    /// Encode exact canonical envelope bytes.
192    pub fn encode(&self) -> Result<Vec<u8>, Diagnostic> {
193        to_canonical_json_with_limits(self, REMOTE_REQUEST_CODEC_LIMITS)
194    }
195
196    /// Decode one request envelope, rejecting unknown fields and formats.
197    pub fn decode(bytes: &[u8]) -> Result<Self, Diagnostic> {
198        let request = from_canonical_json_with_limits::<Self>(bytes, REMOTE_REQUEST_CODEC_LIMITS)?;
199        if request.format != QUERY_REMOTE_REQUEST_FORMAT_V1 {
200            return Err(envelope_failure(
201                DiagnosticCategory::InvalidContract,
202                "query_remote_format_unsupported",
203                "remote request wire format is unsupported",
204            ));
205        }
206        check_nonce(&request.nonce)?;
207        validate_remote_limits(request.limits)?;
208        FingerprintDigest::from_hex(&request.advertisement).map_err(|_| {
209            envelope_failure(
210                DiagnosticCategory::InvalidContract,
211                "query_remote_advertisement_invalid",
212                "remote request carries a malformed advertisement fingerprint",
213            )
214        })?;
215        validate_remote_time_shape(
216            request.prepared_at_unix_ms,
217            request.expires_at_unix_ms,
218            request.limits,
219        )?;
220        if request.encode()? != bytes {
221            return Err(envelope_failure(
222                DiagnosticCategory::Integrity,
223                "query_remote_request_wire_mismatch",
224                "remote request bytes normalize after trusted reconstruction",
225            ));
226        }
227        Ok(request)
228    }
229
230    /// Rebuild the trusted plan from the embedded canonical document.
231    ///
232    /// The embedded tree re-encodes to exact canonical bytes and runs the
233    /// full plan wire decoder, so every structural plan check applies to
234    /// remote plans exactly as it does to local ones.
235    pub fn plan(&self) -> Result<QueryPlan, Diagnostic> {
236        decode_query_plan(&to_canonical_json(&self.plan)?)
237    }
238
239    /// Rebuild the validated invocation against the carried plan.
240    pub fn invocation(&self, plan: &QueryPlan) -> Result<QueryInvocation, Diagnostic> {
241        QueryInvocation::new(
242            plan,
243            self.operation.operation(),
244            self.rows
245                .iter()
246                .map(|row| InputRow::new(row.clone()))
247                .collect(),
248        )
249    }
250
251    /// Return the caller budgets.
252    #[must_use]
253    pub const fn limits(&self) -> RemoteLimits {
254        self.limits
255    }
256
257    /// Return the caller nonce echoed by the response.
258    #[must_use]
259    pub fn nonce(&self) -> &str {
260        &self.nonce
261    }
262
263    /// Return whether this request binds the executor's exact advertisement.
264    pub fn binds_advertisement(
265        &self,
266        advertisement: &RemoteCapabilities,
267    ) -> Result<bool, Diagnostic> {
268        let expected = advertisement.fingerprint()?;
269        let actual = FingerprintDigest::from_hex(&self.advertisement).map_err(|_| {
270            envelope_failure(
271                DiagnosticCategory::InvalidContract,
272                "query_remote_advertisement_invalid",
273                "remote request carries a malformed advertisement fingerprint",
274            )
275        })?;
276        Ok(actual == expected.as_fingerprint().digest())
277    }
278
279    /// Validate this request's absolute lifetime at one executor clock sample.
280    ///
281    /// Clients at most [`MAX_REMOTE_CLOCK_SKEW_MS`] ahead are accepted. The
282    /// expiry itself is exclusive: a request at or beyond it is rejected.
283    /// The returned duration reaches the absolute expiry and is the replay
284    /// retention horizon. It can exceed the declared execution duration by
285    /// the admitted positive clock skew; use [`Self::remaining_execution_ms`]
286    /// to bound execution itself.
287    pub fn remaining_lifetime_ms(&self, now_unix_ms: u64) -> Result<u64, Diagnostic> {
288        validate_remote_time_shape(
289            self.prepared_at_unix_ms,
290            self.expires_at_unix_ms,
291            self.limits,
292        )?;
293        if self.prepared_at_unix_ms > now_unix_ms.saturating_add(MAX_REMOTE_CLOCK_SKEW_MS) {
294            return Err(envelope_failure(
295                DiagnosticCategory::Integrity,
296                "query_remote_time_future",
297                "remote request preparation time exceeds the allowed clock skew",
298            ));
299        }
300        self.expires_at_unix_ms
301            .checked_sub(now_unix_ms)
302            .filter(|remaining| *remaining > 0)
303            .ok_or_else(|| {
304                envelope_failure(
305                    DiagnosticCategory::ResourceLimit,
306                    "query_remote_request_expired",
307                    "remote request absolute expiry has elapsed",
308                )
309            })
310    }
311
312    /// Return the remaining execution duration without granting clock skew.
313    ///
314    /// Positive skew can place absolute expiry later than `now + declared
315    /// duration`. Execution therefore uses the smaller of the absolute
316    /// remaining horizon and the fingerprint-bound declared lifetime, while
317    /// replay retention continues through the full absolute horizon.
318    pub fn remaining_execution_ms(&self, now_unix_ms: u64) -> Result<u64, Diagnostic> {
319        let absolute_remaining = self.remaining_lifetime_ms(now_unix_ms)?;
320        let declared_lifetime = self
321            .expires_at_unix_ms
322            .checked_sub(self.prepared_at_unix_ms)
323            .ok_or_else(remote_time_invalid)?;
324        Ok(absolute_remaining.min(declared_lifetime))
325    }
326
327    /// Return the request's exclusive absolute expiry timestamp.
328    #[must_use]
329    pub const fn expires_at_unix_ms(&self) -> u64 {
330        self.expires_at_unix_ms
331    }
332}
333
334/// One typed value of a remote result row.
335#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
336#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
337pub enum RemoteValue {
338    /// An entity or relation reference.
339    Thing {
340        /// The provider instance identity.
341        iid: String,
342        /// The validated runtime type.
343        type_id: TypeId,
344    },
345    /// An attribute instance with its parsed canonical value.
346    Attribute {
347        /// The validated runtime attribute type.
348        type_id: TypeId,
349        /// The exact typed scalar value.
350        value: CanonicalValue,
351    },
352    /// A pure typed value.
353    Value {
354        /// The exact typed scalar value.
355        value: CanonicalValue,
356    },
357    /// An explicit absence in an optional column.
358    Absent,
359}
360
361/// One typed field value of a remote fetched document.
362#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
363#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
364pub enum RemoteFieldValue {
365    /// One exact typed scalar.
366    Scalar {
367        /// The exact typed scalar value.
368        value: CanonicalValue,
369    },
370    /// An explicit absence in an optional scalar field.
371    Absent,
372    /// A typed list of attribute values.
373    List {
374        /// The exact typed list elements.
375        values: Vec<CanonicalValue>,
376    },
377}
378
379/// The typed terminal outcome of one remote invocation.
380#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
381#[serde(tag = "kind", rename_all = "snake_case", deny_unknown_fields)]
382pub enum RemoteOutcome {
383    /// Evidence-validated projected rows in provider order.
384    Rows {
385        /// Positional row values per validated output column.
386        rows: Vec<Vec<RemoteValue>>,
387    },
388    /// Evidence-validated fetched documents in provider order.
389    Documents {
390        /// Positional field values per validated document column.
391        documents: Vec<Vec<RemoteFieldValue>>,
392    },
393    /// The exact number of returned answers.
394    Count {
395        /// The counted answers.
396        value: u64,
397    },
398    /// Whether at least one answer exists.
399    Exists {
400        /// The existence verdict.
401        value: bool,
402    },
403}
404
405/// Fingerprint domain for whole remote request envelopes.
406pub const QUERY_REMOTE_REQUEST_FINGERPRINT_DOMAIN: &str = "typebridge.query.remote-request";
407/// Canonicalization identifier for whole remote request envelopes.
408pub const QUERY_REMOTE_REQUEST_CANONICALIZATION: &str = "typebridge.query-remote-request/v1";
409
410/// The canonical fingerprint of one complete request envelope.
411///
412/// Covers every request field — plan bytes, operation, input rows, limits,
413/// and nonce — so evidence carrying it is bound to exactly one invocation,
414/// never merely to a plan/nonce pair whose rows may differ.
415#[derive(Clone, Debug, Eq, PartialEq)]
416pub struct RemoteRequestFingerprint(Fingerprint);
417
418impl RemoteRequestFingerprint {
419    /// Compute the fingerprint of exact request envelope bytes.
420    pub fn compute(request_bytes: &[u8]) -> Result<Self, Diagnostic> {
421        Ok(Self(Fingerprint::compute(
422            FingerprintDomain::new(QUERY_REMOTE_REQUEST_FINGERPRINT_DOMAIN)?,
423            CanonicalizationVersion::new(QUERY_REMOTE_REQUEST_CANONICALIZATION)?,
424            None,
425            request_bytes,
426        )))
427    }
428
429    /// Return the generic fingerprint.
430    #[must_use]
431    pub const fn as_fingerprint(&self) -> &Fingerprint {
432        &self.0
433    }
434
435    fn digest_hex(&self) -> String {
436        self.0.digest().to_hex()
437    }
438}
439
440/// Exact Ed25519 public key trusted to authenticate replies from one executor.
441#[derive(Clone, Copy, Eq, Hash, PartialEq)]
442pub struct RemoteSigningPublicKey([u8; 32]);
443
444impl RemoteSigningPublicKey {
445    /// Construct a public key from its exact Ed25519 bytes.
446    #[must_use]
447    pub const fn from_bytes(bytes: [u8; 32]) -> Self {
448        Self(bytes)
449    }
450
451    /// Borrow the exact Ed25519 public-key bytes.
452    #[must_use]
453    pub const fn as_bytes(&self) -> &[u8; 32] {
454        &self.0
455    }
456
457    /// Encode the key as fixed-width lowercase hexadecimal wire text.
458    #[must_use]
459    pub fn to_hex(self) -> String {
460        encode_hex(&self.0)
461    }
462
463    fn from_hex(value: &str) -> Result<Self, ()> {
464        decode_fixed_hex(value).map(Self)
465    }
466}
467
468impl std::fmt::Debug for RemoteSigningPublicKey {
469    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
470        formatter
471            .debug_tuple("RemoteSigningPublicKey")
472            .field(&self.to_hex())
473            .finish()
474    }
475}
476
477impl Serialize for RemoteSigningPublicKey {
478    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
479    where
480        S: serde::Serializer,
481    {
482        serializer.serialize_str(&self.to_hex())
483    }
484}
485
486impl<'de> Deserialize<'de> for RemoteSigningPublicKey {
487    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
488    where
489        D: Deserializer<'de>,
490    {
491        let value = String::deserialize(deserializer)?;
492        Self::from_hex(&value).map_err(|()| serde::de::Error::custom("invalid remote signing key"))
493    }
494}
495
496/// Deterministic identity of one exact remote reply-signing public key.
497///
498/// The identifier is domain-separated from every other digest and is carried
499/// alongside the key in capability advertisements and signed outer replies.
500#[derive(Clone, Copy, Eq, Hash, PartialEq)]
501pub struct RemoteSigningKeyId([u8; 32]);
502
503impl RemoteSigningKeyId {
504    /// Derive the identifier for one exact Ed25519 public key.
505    #[must_use]
506    pub fn for_public_key(key: RemoteSigningPublicKey) -> Self {
507        let mut hasher = Sha256::new();
508        hasher.update(QUERY_REMOTE_REPLY_KEY_ID_DOMAIN.as_bytes());
509        hasher.update([0]);
510        hasher.update(key.as_bytes());
511        Self(hasher.finalize().into())
512    }
513
514    /// Borrow the exact key-identifier bytes.
515    #[must_use]
516    pub const fn as_bytes(&self) -> &[u8; 32] {
517        &self.0
518    }
519
520    /// Encode the identifier as fixed-width lowercase hexadecimal wire text.
521    #[must_use]
522    pub fn to_hex(self) -> String {
523        encode_hex(&self.0)
524    }
525
526    fn from_hex(value: &str) -> Result<Self, ()> {
527        decode_fixed_hex(value).map(Self)
528    }
529}
530
531impl std::fmt::Debug for RemoteSigningKeyId {
532    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
533        formatter
534            .debug_tuple("RemoteSigningKeyId")
535            .field(&self.to_hex())
536            .finish()
537    }
538}
539
540impl Serialize for RemoteSigningKeyId {
541    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
542    where
543        S: serde::Serializer,
544    {
545        serializer.serialize_str(&self.to_hex())
546    }
547}
548
549impl<'de> Deserialize<'de> for RemoteSigningKeyId {
550    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
551    where
552        D: Deserializer<'de>,
553    {
554        let value = String::deserialize(deserializer)?;
555        Self::from_hex(&value)
556            .map_err(|()| serde::de::Error::custom("invalid remote signing key identifier"))
557    }
558}
559
560/// Exact Ed25519 signature carried by the authenticated outer reply.
561#[derive(Clone, Copy, Eq, PartialEq)]
562pub struct RemoteReplySignature([u8; 64]);
563
564impl RemoteReplySignature {
565    /// Construct a signature from its exact Ed25519 bytes.
566    #[must_use]
567    pub const fn from_bytes(bytes: [u8; 64]) -> Self {
568        Self(bytes)
569    }
570
571    /// Borrow the exact Ed25519 signature bytes.
572    #[must_use]
573    pub const fn as_bytes(&self) -> &[u8; 64] {
574        &self.0
575    }
576
577    fn to_hex(self) -> String {
578        encode_hex(&self.0)
579    }
580
581    fn from_hex(value: &str) -> Result<Self, ()> {
582        decode_fixed_hex(value).map(Self)
583    }
584}
585
586impl std::fmt::Debug for RemoteReplySignature {
587    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
588        formatter.write_str("RemoteReplySignature(..)")
589    }
590}
591
592/// Domain-separated digest signed for one canonical unsigned outer reply.
593#[derive(Clone, Copy, Debug, Eq, PartialEq)]
594pub struct RemoteReplySigningDigest([u8; 32]);
595
596impl RemoteReplySigningDigest {
597    /// Borrow the digest bytes passed to the Ed25519 implementation.
598    #[must_use]
599    pub const fn as_bytes(&self) -> &[u8; 32] {
600        &self.0
601    }
602}
603
604/// Binding-neutral signing operation used by the contract wire encoder.
605pub trait RemoteReplySigner {
606    /// Return the public key paired with this signer.
607    fn public_key(&self) -> RemoteSigningPublicKey;
608
609    /// Sign one domain-separated canonical reply digest.
610    fn sign(&self, digest: &RemoteReplySigningDigest) -> RemoteReplySignature;
611}
612
613/// Binding-neutral signature verifier used before any reply payload is decoded.
614pub trait RemoteReplyVerifier {
615    /// Verify one digest against the exact trusted public key.
616    fn verify(
617        &self,
618        key: RemoteSigningPublicKey,
619        digest: &RemoteReplySigningDigest,
620        signature: &RemoteReplySignature,
621    ) -> bool;
622}
623
624/// Expected authenticated success shape used for allocation-free budget scans.
625#[derive(Clone, Copy, Debug, Eq, PartialEq)]
626pub enum RemoteOutcomeShape {
627    /// Selected rows with this exact output width.
628    Rows {
629        /// Exact number of values in every selected row.
630        width: usize,
631    },
632    /// Fetched documents with this exact output width.
633    Documents {
634        /// Exact number of fields in every fetched document.
635        width: usize,
636    },
637    /// A scalar count.
638    Count,
639    /// A scalar existence verdict.
640    Exists,
641}
642
643/// Caller limits checked before a typed remote outcome is allocated.
644#[derive(Clone, Copy, Debug, Eq, PartialEq)]
645pub struct RemoteReplyDecodeLimits {
646    /// Expected success outcome shape.
647    pub shape: RemoteOutcomeShape,
648    /// Maximum authenticated successful outer-reply bytes accepted by the caller.
649    ///
650    /// Request-bound failure envelopes remain subject to the protocol hard
651    /// ceiling so their typed diagnostic can be surfaced at any success budget.
652    pub max_bytes: u64,
653    /// Maximum rows, documents, count value, or positive existence item accepted by the caller.
654    pub max_items: u64,
655    /// Maximum aggregate list members accepted across fetched documents.
656    pub max_collection_members: u64,
657}
658
659/// One successful remote execution bound to its request and plan.
660#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
661#[serde(deny_unknown_fields)]
662pub struct RemoteQueryResponse {
663    format: String,
664    nonce: String,
665    outcome: RemoteOutcome,
666    plan: String,
667    // Pre-release wire ledger: the /v1 response gained the whole-request
668    // binding before any 2.0.0 artifact shipped; no released bytes change.
669    request: String,
670}
671
672impl RemoteQueryResponse {
673    /// Bind one outcome to the request nonce, plan, and whole request.
674    pub fn new(
675        nonce: impl Into<String>,
676        plan: &QueryPlanFingerprint,
677        request: &RemoteRequestFingerprint,
678        outcome: RemoteOutcome,
679    ) -> Result<Self, Diagnostic> {
680        let nonce = nonce.into();
681        check_nonce(&nonce)?;
682        Ok(Self {
683            format: QUERY_REMOTE_RESPONSE_FORMAT_V1.to_owned(),
684            nonce,
685            outcome,
686            plan: plan.as_fingerprint().digest().to_hex(),
687            request: request.digest_hex(),
688        })
689    }
690
691    /// Encode one authenticated outer reply using the exact trusted advertisement.
692    pub fn encode_signed(
693        &self,
694        advertisement: &RemoteCapabilitiesFingerprint,
695        signer: &impl RemoteReplySigner,
696    ) -> Result<Vec<u8>, Diagnostic> {
697        encode_signed_reply(&self.encode_payload()?, advertisement, signer)
698    }
699
700    /// Return the exact authenticated wire length without invoking a signer.
701    pub fn signed_encoded_len(
702        &self,
703        advertisement: &RemoteCapabilitiesFingerprint,
704        key: RemoteSigningPublicKey,
705    ) -> Result<usize, Diagnostic> {
706        Ok(signed_reply_encoded_len(
707            &self.encode_payload()?,
708            advertisement,
709            key,
710        ))
711    }
712
713    fn encode_payload(&self) -> Result<Vec<u8>, Diagnostic> {
714        to_canonical_json_with_limits(self, REMOTE_ENVELOPE_CODEC_LIMITS)
715    }
716
717    fn decode_bound(bytes: &[u8]) -> Result<Self, Diagnostic> {
718        let response =
719            from_canonical_json_with_limits::<Self>(bytes, REMOTE_ENVELOPE_CODEC_LIMITS)?;
720        if response.format != QUERY_REMOTE_RESPONSE_FORMAT_V1 {
721            return Err(envelope_failure(
722                DiagnosticCategory::InvalidContract,
723                "query_remote_format_unsupported",
724                "remote response wire format is unsupported",
725            ));
726        }
727        if response.encode_payload()? != bytes {
728            return Err(envelope_failure(
729                DiagnosticCategory::Integrity,
730                "query_remote_response_wire_mismatch",
731                "remote response bytes normalize after trusted reconstruction",
732            ));
733        }
734        Ok(response)
735    }
736
737    /// Return the typed outcome.
738    #[must_use]
739    pub const fn outcome(&self) -> &RemoteOutcome {
740        &self.outcome
741    }
742
743    /// Consume the envelope and return its typed outcome.
744    #[must_use]
745    pub fn into_outcome(self) -> RemoteOutcome {
746        self.outcome
747    }
748}
749
750/// One structured remote failure bound to its request.
751#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
752#[serde(deny_unknown_fields)]
753pub struct RemoteQueryFailure {
754    category: DiagnosticCategory,
755    code: String,
756    format: String,
757    message: String,
758    nonce: Option<String>,
759    // Pre-release wire ledger: the /v1 failure gained the optional
760    // whole-request binding before any 2.0.0 artifact shipped.
761    request: Option<String>,
762}
763
764impl RemoteQueryFailure {
765    /// Bind one structured diagnostic to the request nonce.
766    #[must_use]
767    pub fn new(nonce: Option<String>, diagnostic: &Diagnostic) -> Self {
768        Self {
769            category: diagnostic.category(),
770            code: diagnostic.code().as_str().to_owned(),
771            format: QUERY_REMOTE_FAILURE_FORMAT_V1.to_owned(),
772            message: diagnostic.message().to_owned(),
773            nonce,
774            request: None,
775        }
776    }
777
778    /// Bind one structured diagnostic to the nonce and the exact request.
779    #[must_use]
780    pub fn bound(
781        nonce: impl Into<String>,
782        request: &RemoteRequestFingerprint,
783        diagnostic: &Diagnostic,
784    ) -> Self {
785        Self {
786            category: diagnostic.category(),
787            code: diagnostic.code().as_str().to_owned(),
788            format: QUERY_REMOTE_FAILURE_FORMAT_V1.to_owned(),
789            message: diagnostic.message().to_owned(),
790            nonce: Some(nonce.into()),
791            request: Some(request.digest_hex()),
792        }
793    }
794
795    /// Verify this failure's binding to the request the caller sent.
796    ///
797    /// Both the nonce and whole-request digest are mandatory on a reply to a
798    /// valid request. Pre-decode transport failures use a separate channel;
799    /// an unbound envelope is never treated as request-correlated evidence.
800    pub fn verify_binding(
801        &self,
802        expected_nonce: &str,
803        expected_request: &RemoteRequestFingerprint,
804    ) -> Result<(), Diagnostic> {
805        verify_failure_binding(
806            self.nonce.as_deref(),
807            self.request.as_deref(),
808            expected_nonce,
809            expected_request,
810        )
811    }
812
813    /// Encode one authenticated outer reply using the exact trusted advertisement.
814    pub fn encode_signed(
815        &self,
816        advertisement: &RemoteCapabilitiesFingerprint,
817        signer: &impl RemoteReplySigner,
818    ) -> Result<Vec<u8>, Diagnostic> {
819        encode_signed_reply(&self.encode_payload()?, advertisement, signer)
820    }
821
822    /// Encode an authenticated failure, replacing an unencodable diagnostic
823    /// with a fixed bounded internal failure while preserving safe bindings.
824    ///
825    /// Server transports use this path so production failures can never
826    /// collapse to empty or unsigned bytes merely because an upstream error
827    /// message exceeded a codec ceiling.
828    #[must_use]
829    pub fn encode_signed_or_fallback(
830        &self,
831        advertisement: &RemoteCapabilitiesFingerprint,
832        signer: &impl RemoteReplySigner,
833    ) -> Vec<u8> {
834        match self.encode_signed(advertisement, signer) {
835            Ok(encoded) => encoded,
836            Err(_) => encode_minimal_signed_failure(self, advertisement, signer),
837        }
838    }
839
840    fn encode_payload(&self) -> Result<Vec<u8>, Diagnostic> {
841        to_canonical_json_with_limits(self, REMOTE_ENVELOPE_CODEC_LIMITS)
842    }
843
844    /// Decode an already authenticated failure payload.
845    pub fn decode_payload(bytes: &[u8]) -> Result<Self, Diagnostic> {
846        let failure = from_canonical_json_with_limits::<Self>(bytes, REMOTE_ENVELOPE_CODEC_LIMITS)?;
847        if failure.format != QUERY_REMOTE_FAILURE_FORMAT_V1 {
848            return Err(envelope_failure(
849                DiagnosticCategory::InvalidContract,
850                "query_remote_format_unsupported",
851                "remote failure wire format is unsupported",
852            ));
853        }
854        Ok(failure)
855    }
856
857    /// Rebuild the structured diagnostic.
858    pub fn diagnostic(&self) -> Result<Diagnostic, Diagnostic> {
859        let code = DiagnosticCode::new(self.code.clone()).map_err(|_| {
860            envelope_failure(
861                DiagnosticCategory::InvalidContract,
862                "query_remote_code_invalid",
863                "remote failure carries a malformed diagnostic code",
864            )
865        })?;
866        Ok(Diagnostic::new(self.category, code, self.message.clone()))
867    }
868
869    /// Return the echoed request nonce, when the request decoded far enough.
870    #[must_use]
871    pub fn nonce(&self) -> Option<&str> {
872        self.nonce.as_deref()
873    }
874}
875
876/// The exact wire discriminator for first-format capability advertisements.
877pub const QUERY_REMOTE_CAPABILITIES_FORMAT_V1: &str = "typebridge.query-remote-capabilities/v1";
878/// Fingerprint domain for exact executor capability advertisements.
879pub const QUERY_REMOTE_CAPABILITIES_FINGERPRINT_DOMAIN: &str =
880    "typebridge.query.remote-capabilities";
881/// Canonicalization identifier for capability-advertisement fingerprints.
882pub const QUERY_REMOTE_CAPABILITIES_CANONICALIZATION: &str =
883    "typebridge.query-remote-capabilities/v1";
884
885/// One logical executor identity and one concrete process/shared-store epoch.
886///
887/// Standalone executors generate a fresh pair at startup. A multi-instance
888/// deployment may share a pair only together with a globally atomic replay
889/// store; otherwise advertisements differ and cross-instance requests fail
890/// closed at preflight.
891#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
892#[serde(deny_unknown_fields)]
893pub struct RemoteExecutorBinding {
894    epoch: String,
895    identity: String,
896}
897
898impl RemoteExecutorBinding {
899    /// Construct one validated executor identity/epoch pair.
900    pub fn new(identity: impl Into<String>, epoch: impl Into<String>) -> Result<Self, Diagnostic> {
901        let binding = Self {
902            epoch: epoch.into(),
903            identity: identity.into(),
904        };
905        validate_executor_component(&binding.identity)?;
906        validate_executor_component(&binding.epoch)?;
907        Ok(binding)
908    }
909
910    /// Return the logical executor identity.
911    #[must_use]
912    pub fn identity(&self) -> &str {
913        &self.identity
914    }
915
916    /// Return the concrete process/shared-store epoch.
917    #[must_use]
918    pub fn epoch(&self) -> &str {
919        &self.epoch
920    }
921}
922
923/// The canonical fingerprint of one exact capability advertisement.
924#[derive(Clone, Debug, Eq, PartialEq)]
925pub struct RemoteCapabilitiesFingerprint(Fingerprint);
926
927impl RemoteCapabilitiesFingerprint {
928    /// Return the generic fingerprint.
929    #[must_use]
930    pub const fn as_fingerprint(&self) -> &Fingerprint {
931        &self.0
932    }
933
934    /// Return the fixed-width lowercase advertisement digest.
935    #[must_use]
936    pub fn digest_hex(&self) -> String {
937        self.0.digest().to_hex()
938    }
939}
940
941/// One executor capability advertisement for pre-flight negotiation.
942///
943/// A client checks its plan's required capabilities against this set and
944/// refuses to send unsupported plans; the executor re-checks on receipt.
945#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
946#[serde(deny_unknown_fields)]
947pub struct RemoteCapabilities {
948    capabilities: crate::capability::CapabilitySet,
949    // Pre-release wire ledger: /v1 gained an executor identity and epoch
950    // before 2.0.0 shipped. Requests fingerprint these exact bytes so a
951    // restart or a different standalone instance cannot accept captured
952    // requests prepared for an earlier executor incarnation.
953    executor: RemoteExecutorBinding,
954    format: String,
955    // This key is an explicit caller trust input. Its inclusion in the exact
956    // advertisement fingerprint binds every prepared request to one signer.
957    reply_key: RemoteSigningPublicKey,
958    // The domain-separated identity is redundant by design: decoding verifies
959    // it against `reply_key`, while signed replies bind both exact values.
960    reply_key_id: RemoteSigningKeyId,
961}
962
963impl RemoteCapabilities {
964    /// Advertise one executor capability set.
965    #[must_use]
966    pub fn new(
967        capabilities: crate::capability::CapabilitySet,
968        executor: RemoteExecutorBinding,
969        reply_key: RemoteSigningPublicKey,
970    ) -> Self {
971        let reply_key_id = RemoteSigningKeyId::for_public_key(reply_key);
972        Self {
973            capabilities,
974            executor,
975            format: QUERY_REMOTE_CAPABILITIES_FORMAT_V1.to_owned(),
976            reply_key,
977            reply_key_id,
978        }
979    }
980
981    /// Encode exact canonical envelope bytes.
982    pub fn encode(&self) -> Result<Vec<u8>, Diagnostic> {
983        to_canonical_json_with_limits(self, REMOTE_ENVELOPE_CODEC_LIMITS)
984    }
985
986    /// Decode one advertisement, rejecting unknown fields and formats.
987    pub fn decode(bytes: &[u8]) -> Result<Self, Diagnostic> {
988        let advertisement =
989            from_canonical_json_with_limits::<Self>(bytes, REMOTE_ENVELOPE_CODEC_LIMITS)?;
990        if advertisement.format != QUERY_REMOTE_CAPABILITIES_FORMAT_V1 {
991            return Err(envelope_failure(
992                DiagnosticCategory::InvalidContract,
993                "query_remote_format_unsupported",
994                "remote capability wire format is unsupported",
995            ));
996        }
997        validate_executor_component(advertisement.executor.identity())?;
998        validate_executor_component(advertisement.executor.epoch())?;
999        if advertisement.reply_key_id != RemoteSigningKeyId::for_public_key(advertisement.reply_key)
1000        {
1001            return Err(remote_signature_invalid());
1002        }
1003        if advertisement.encode()? != bytes {
1004            return Err(envelope_failure(
1005                DiagnosticCategory::Integrity,
1006                "query_remote_capabilities_wire_mismatch",
1007                "remote capability bytes normalize after trusted reconstruction",
1008            ));
1009        }
1010        Ok(advertisement)
1011    }
1012
1013    /// Return the advertised capability set.
1014    #[must_use]
1015    pub const fn capabilities(&self) -> &crate::capability::CapabilitySet {
1016        &self.capabilities
1017    }
1018
1019    /// Return this advertisement's executor identity and epoch.
1020    #[must_use]
1021    pub const fn executor(&self) -> &RemoteExecutorBinding {
1022        &self.executor
1023    }
1024
1025    /// Return the exact public key trusted to authenticate executor replies.
1026    #[must_use]
1027    pub const fn reply_key(&self) -> RemoteSigningPublicKey {
1028        self.reply_key
1029    }
1030
1031    /// Return the deterministic identity of the advertised reply key.
1032    #[must_use]
1033    pub const fn reply_key_id(&self) -> RemoteSigningKeyId {
1034        self.reply_key_id
1035    }
1036
1037    /// Fingerprint the exact canonical advertisement, including its epoch.
1038    pub fn fingerprint(&self) -> Result<RemoteCapabilitiesFingerprint, Diagnostic> {
1039        Ok(RemoteCapabilitiesFingerprint(Fingerprint::compute(
1040            FingerprintDomain::new(QUERY_REMOTE_CAPABILITIES_FINGERPRINT_DOMAIN)?,
1041            CanonicalizationVersion::new(QUERY_REMOTE_CAPABILITIES_CANONICALIZATION)?,
1042            None,
1043            &self.encode()?,
1044        )))
1045    }
1046}
1047
1048/// One decoded remote reply: a typed response or a request-bound failure.
1049#[derive(Clone, Debug, PartialEq)]
1050pub enum RemoteReply {
1051    /// A typed successful response, fully bound to the request.
1052    Response(RemoteQueryResponse),
1053    /// A structured failure whose request bindings were verified.
1054    Failure(RemoteQueryFailure),
1055}
1056
1057#[derive(Deserialize)]
1058#[serde(deny_unknown_fields)]
1059struct SignedRemoteReplyPeek<'wire> {
1060    #[serde(borrow)]
1061    advertisement: &'wire str,
1062    #[serde(borrow)]
1063    format: &'wire str,
1064    #[serde(borrow)]
1065    key: &'wire str,
1066    #[serde(borrow)]
1067    key_id: &'wire str,
1068    #[serde(borrow)]
1069    payload: &'wire serde_json::value::RawValue,
1070    #[serde(borrow)]
1071    signature: &'wire str,
1072}
1073
1074/// Borrowed correlation fields parsed before any reply outcome is materialized.
1075///
1076/// Canonical replies borrow all four strings directly from the input. Escaped
1077/// spellings are non-canonical and fail this precheck without allocating;
1078/// ignored outcome values are traversed without constructing a
1079/// `serde_json::Value` or a [`RemoteOutcome`].
1080#[derive(Deserialize)]
1081struct RemoteReplyBindingPeek<'wire> {
1082    #[serde(borrow)]
1083    format: &'wire str,
1084    #[serde(borrow)]
1085    nonce: Option<&'wire str>,
1086    #[serde(borrow)]
1087    plan: Option<&'wire str>,
1088    #[serde(borrow)]
1089    request: Option<&'wire str>,
1090}
1091
1092/// Decode one reply envelope of either kind and verify its request binding.
1093///
1094/// Success and failure envelopes share one entry point so every caller —
1095/// including the Python and Node bindings — decodes both outcomes with the
1096/// same nonce and whole-request correlation checks.
1097#[expect(
1098    clippy::too_many_arguments,
1099    reason = "the trust-boundary API keeps every expected binding and verifier explicit"
1100)]
1101pub fn decode_remote_reply(
1102    bytes: &[u8],
1103    expected_nonce: &str,
1104    expected_plan: &QueryPlanFingerprint,
1105    expected_request: &RemoteRequestFingerprint,
1106    expected_advertisement: &RemoteCapabilitiesFingerprint,
1107    trusted_key: RemoteSigningPublicKey,
1108    limits: RemoteReplyDecodeLimits,
1109    verifier: &impl RemoteReplyVerifier,
1110) -> Result<RemoteReply, Diagnostic> {
1111    // Every reply is capped before JSON parsing or signature work. The caller
1112    // budget applies only to successful data: applying it before authenticating
1113    // the reply kind would make a valid request-bound failure undecodable at a
1114    // tiny budget, including the failure explaining that no success can fit.
1115    preflight_remote_reply_size(bytes, u64::MAX)?;
1116    let payload = verify_signed_reply(bytes, expected_advertisement, trusted_key, verifier)?;
1117    let peek = peek_remote_reply_binding(payload)?;
1118    if peek.format == QUERY_REMOTE_RESPONSE_FORMAT_V1 {
1119        verify_response_binding(
1120            peek.nonce,
1121            peek.plan,
1122            peek.request,
1123            expected_nonce,
1124            expected_plan,
1125            expected_request,
1126        )?;
1127        preflight_remote_reply_size(bytes, limits.max_bytes)?;
1128        preflight_remote_response_shape(payload, limits)?;
1129        return Ok(RemoteReply::Response(RemoteQueryResponse::decode_bound(
1130            payload,
1131        )?));
1132    }
1133    if peek.format == QUERY_REMOTE_FAILURE_FORMAT_V1 {
1134        verify_failure_binding(peek.nonce, peek.request, expected_nonce, expected_request)?;
1135        let failure = RemoteQueryFailure::decode_payload(payload)?;
1136        failure.verify_binding(expected_nonce, expected_request)?;
1137        return Ok(RemoteReply::Failure(failure));
1138    }
1139    Err(envelope_failure(
1140        DiagnosticCategory::InvalidContract,
1141        "query_remote_format_unsupported",
1142        "remote reply wire format is unsupported",
1143    ))
1144}
1145
1146/// Authenticate and decode an uncorrelated remote failure.
1147///
1148/// This is reserved for transport failures that occur before a request can be
1149/// decoded and fingerprinted. Request-correlated clients must use
1150/// [`decode_remote_reply`] so nonce, plan, and whole-request bindings are
1151/// mandatory.
1152pub fn decode_signed_remote_failure(
1153    bytes: &[u8],
1154    expected_advertisement: &RemoteCapabilitiesFingerprint,
1155    trusted_key: RemoteSigningPublicKey,
1156    max_bytes: u64,
1157    verifier: &impl RemoteReplyVerifier,
1158) -> Result<RemoteQueryFailure, Diagnostic> {
1159    preflight_remote_reply_size(bytes, max_bytes)?;
1160    let payload = verify_signed_reply(bytes, expected_advertisement, trusted_key, verifier)?;
1161    RemoteQueryFailure::decode_payload(payload)
1162}
1163
1164fn preflight_remote_reply_size(bytes: &[u8], caller_max_bytes: u64) -> Result<(), Diagnostic> {
1165    let wire_max_bytes = u64::try_from(MAX_REMOTE_ENVELOPE_BYTES).unwrap_or(u64::MAX);
1166    let effective_max_bytes = caller_max_bytes.min(wire_max_bytes);
1167    if u64::try_from(bytes.len()).unwrap_or(u64::MAX) <= effective_max_bytes {
1168        return Ok(());
1169    }
1170    if caller_max_bytes < wire_max_bytes {
1171        return Err(remote_response_oversized());
1172    }
1173    Err(remote_envelope_too_large())
1174}
1175
1176/// Shared authenticated-outer-envelope byte preflight for additive payload
1177/// versions. This does not inspect or reconstruct a version-specific payload.
1178pub(crate) fn preflight_signed_reply_size(
1179    bytes: &[u8],
1180    caller_max_bytes: u64,
1181) -> Result<(), Diagnostic> {
1182    preflight_remote_reply_size(bytes, caller_max_bytes)
1183}
1184
1185fn verify_signed_reply<'wire>(
1186    bytes: &'wire [u8],
1187    expected_advertisement: &RemoteCapabilitiesFingerprint,
1188    trusted_key: RemoteSigningPublicKey,
1189    verifier: &impl RemoteReplyVerifier,
1190) -> Result<&'wire [u8], Diagnostic> {
1191    let outer = serde_json::from_slice::<SignedRemoteReplyPeek<'wire>>(bytes).map_err(|_| {
1192        envelope_failure(
1193            DiagnosticCategory::InvalidContract,
1194            "query_remote_reply_malformed",
1195            "remote reply is not a signed JSON envelope",
1196        )
1197    })?;
1198    let expected_advertisement = expected_advertisement.digest_hex();
1199    let expected_key = trusted_key.to_hex();
1200    let expected_key_id = RemoteSigningKeyId::for_public_key(trusted_key).to_hex();
1201    let signature =
1202        RemoteReplySignature::from_hex(outer.signature).map_err(|()| remote_signature_invalid())?;
1203    if outer.advertisement != expected_advertisement
1204        || outer.key != expected_key
1205        || outer.key_id != expected_key_id
1206        || outer.signature != signature.to_hex()
1207    {
1208        return Err(remote_signature_invalid());
1209    }
1210    let digest = remote_reply_signing_digest(
1211        outer.advertisement,
1212        outer.format,
1213        outer.key,
1214        outer.key_id,
1215        outer.payload.get().as_bytes(),
1216    );
1217    if !verifier.verify(trusted_key, &digest, &signature) {
1218        return Err(remote_signature_invalid());
1219    }
1220    if outer.format != QUERY_REMOTE_SIGNED_REPLY_FORMAT_V1 {
1221        return Err(envelope_failure(
1222            DiagnosticCategory::InvalidContract,
1223            "query_remote_format_unsupported",
1224            "signed remote reply wire format is unsupported",
1225        ));
1226    }
1227    let canonical = canonical_signed_reply(
1228        outer.advertisement,
1229        outer.format,
1230        outer.key,
1231        outer.key_id,
1232        outer.payload.get().as_bytes(),
1233        outer.signature,
1234    );
1235    if canonical != bytes {
1236        return Err(envelope_failure(
1237            DiagnosticCategory::InvalidContract,
1238            "non_canonical_json",
1239            "input is valid JSON but not the canonical encoding",
1240        ));
1241    }
1242    Ok(outer.payload.get().as_bytes())
1243}
1244
1245/// Authenticate the unchanged signed outer envelope for an additive payload
1246/// version. Payload dispatch remains the caller's version-specific boundary.
1247pub(crate) fn verify_signed_reply_payload<'wire>(
1248    bytes: &'wire [u8],
1249    expected_advertisement: &RemoteCapabilitiesFingerprint,
1250    trusted_key: RemoteSigningPublicKey,
1251    verifier: &impl RemoteReplyVerifier,
1252) -> Result<&'wire [u8], Diagnostic> {
1253    verify_signed_reply(bytes, expected_advertisement, trusted_key, verifier)
1254}
1255
1256fn encode_signed_reply(
1257    payload: &[u8],
1258    advertisement: &RemoteCapabilitiesFingerprint,
1259    signer: &impl RemoteReplySigner,
1260) -> Result<Vec<u8>, Diagnostic> {
1261    let encoded = encode_signed_reply_unchecked(payload, advertisement, signer);
1262    if encoded.len() > MAX_REMOTE_ENVELOPE_BYTES {
1263        return Err(envelope_failure(
1264            DiagnosticCategory::ResourceLimit,
1265            "query_remote_envelope_too_large",
1266            "remote reply exceeds the envelope byte ceiling",
1267        ));
1268    }
1269    Ok(encoded)
1270}
1271
1272/// Encode additive versioned payload bytes in the unchanged signed outer
1273/// envelope.
1274pub(crate) fn encode_signed_reply_payload(
1275    payload: &[u8],
1276    advertisement: &RemoteCapabilitiesFingerprint,
1277    signer: &impl RemoteReplySigner,
1278) -> Result<Vec<u8>, Diagnostic> {
1279    encode_signed_reply(payload, advertisement, signer)
1280}
1281
1282/// Encode already bounded, canonical additive payload bytes.
1283///
1284/// This is used only by the fixed internal-failure fallback after the normal
1285/// checked encoder has failed. Version-specific modules must construct the
1286/// payload from static text plus previously validated ASCII bindings.
1287pub(crate) fn encode_signed_reply_payload_unchecked(
1288    payload: &[u8],
1289    advertisement: &RemoteCapabilitiesFingerprint,
1290    signer: &impl RemoteReplySigner,
1291) -> Vec<u8> {
1292    encode_signed_reply_unchecked(payload, advertisement, signer)
1293}
1294
1295fn encode_signed_reply_unchecked(
1296    payload: &[u8],
1297    advertisement: &RemoteCapabilitiesFingerprint,
1298    signer: &impl RemoteReplySigner,
1299) -> Vec<u8> {
1300    let advertisement = advertisement.digest_hex();
1301    let key = signer.public_key().to_hex();
1302    let key_id = RemoteSigningKeyId::for_public_key(signer.public_key()).to_hex();
1303    let digest = remote_reply_signing_digest(
1304        &advertisement,
1305        QUERY_REMOTE_SIGNED_REPLY_FORMAT_V1,
1306        &key,
1307        &key_id,
1308        payload,
1309    );
1310    let signature = signer.sign(&digest).to_hex();
1311    canonical_signed_reply(
1312        &advertisement,
1313        QUERY_REMOTE_SIGNED_REPLY_FORMAT_V1,
1314        &key,
1315        &key_id,
1316        payload,
1317        &signature,
1318    )
1319}
1320
1321fn encode_minimal_signed_failure(
1322    original: &RemoteQueryFailure,
1323    advertisement: &RemoteCapabilitiesFingerprint,
1324    signer: &impl RemoteReplySigner,
1325) -> Vec<u8> {
1326    const PREFIX: &[u8] = b"{\"category\":\"integrity\",\"code\":\"query_remote_internal_failure\",\"format\":\"typebridge.query-remote-failure/v1\",\"message\":\"executor could not encode the original remote failure\",\"nonce\":";
1327    let nonce = original
1328        .nonce
1329        .as_deref()
1330        .filter(|nonce| check_nonce(nonce).is_ok());
1331    let request = original.request.as_deref().filter(|request| {
1332        request.len() == 64
1333            && request
1334                .bytes()
1335                .all(|byte| byte.is_ascii_digit() || matches!(byte, b'a'..=b'f'))
1336    });
1337    let mut payload = Vec::with_capacity(PREFIX.len() + 256);
1338    payload.extend_from_slice(PREFIX);
1339    append_optional_safe_ascii(&mut payload, nonce);
1340    payload.extend_from_slice(b",\"request\":");
1341    append_optional_safe_ascii(&mut payload, request);
1342    payload.push(b'}');
1343    encode_signed_reply_unchecked(&payload, advertisement, signer)
1344}
1345
1346fn append_optional_safe_ascii(encoded: &mut Vec<u8>, value: Option<&str>) {
1347    if let Some(value) = value {
1348        encoded.push(b'"');
1349        encoded.extend_from_slice(value.as_bytes());
1350        encoded.push(b'"');
1351    } else {
1352        encoded.extend_from_slice(b"null");
1353    }
1354}
1355
1356fn signed_reply_encoded_len(
1357    payload: &[u8],
1358    advertisement: &RemoteCapabilitiesFingerprint,
1359    key: RemoteSigningPublicKey,
1360) -> usize {
1361    canonical_signed_reply(
1362        &advertisement.digest_hex(),
1363        QUERY_REMOTE_SIGNED_REPLY_FORMAT_V1,
1364        &key.to_hex(),
1365        &RemoteSigningKeyId::for_public_key(key).to_hex(),
1366        payload,
1367        &encode_hex(&[0_u8; 64]),
1368    )
1369    .len()
1370}
1371
1372/// Return the exact unchanged signed-envelope length for additive payload
1373/// bytes without constructing a signature.
1374pub(crate) fn signed_reply_payload_encoded_len(
1375    payload: &[u8],
1376    advertisement: &RemoteCapabilitiesFingerprint,
1377    key: RemoteSigningPublicKey,
1378) -> usize {
1379    signed_reply_encoded_len(payload, advertisement, key)
1380}
1381
1382fn remote_reply_signing_digest(
1383    advertisement: &str,
1384    format: &str,
1385    key: &str,
1386    key_id: &str,
1387    payload: &[u8],
1388) -> RemoteReplySigningDigest {
1389    let prefix = canonical_signed_reply_prefix(advertisement, format, key, key_id);
1390    let mut hasher = Sha256::new();
1391    hasher.update(QUERY_REMOTE_REPLY_SIGNATURE_DOMAIN.as_bytes());
1392    hasher.update([0]);
1393    hasher.update(&prefix);
1394    hasher.update(payload);
1395    hasher.update(b"}");
1396    RemoteReplySigningDigest(hasher.finalize().into())
1397}
1398
1399fn canonical_signed_reply_prefix(
1400    advertisement: &str,
1401    format: &str,
1402    key: &str,
1403    key_id: &str,
1404) -> Vec<u8> {
1405    format!(
1406        "{{\"advertisement\":\"{advertisement}\",\"format\":\"{format}\",\"key\":\"{key}\",\"key_id\":\"{key_id}\",\"payload\":"
1407    )
1408    .into_bytes()
1409}
1410
1411fn canonical_signed_reply(
1412    advertisement: &str,
1413    format: &str,
1414    key: &str,
1415    key_id: &str,
1416    payload: &[u8],
1417    signature: &str,
1418) -> Vec<u8> {
1419    let prefix = canonical_signed_reply_prefix(advertisement, format, key, key_id);
1420    let suffix = format!(",\"signature\":\"{signature}\"}}");
1421    let mut encoded = Vec::with_capacity(prefix.len() + payload.len() + suffix.len());
1422    encoded.extend_from_slice(&prefix);
1423    encoded.extend_from_slice(payload);
1424    encoded.extend_from_slice(suffix.as_bytes());
1425    encoded
1426}
1427
1428/// Stable rejection for an unauthenticated or foreign remote reply.
1429#[must_use]
1430pub fn remote_signature_invalid() -> Diagnostic {
1431    envelope_failure(
1432        DiagnosticCategory::Integrity,
1433        "query_remote_signature_invalid",
1434        "remote reply signature or trusted executor binding is invalid",
1435    )
1436}
1437
1438fn remote_response_oversized() -> Diagnostic {
1439    envelope_failure(
1440        DiagnosticCategory::ResourceLimit,
1441        "query_remote_response_oversized",
1442        "response envelope exceeds the caller byte budget",
1443    )
1444}
1445
1446fn remote_envelope_too_large() -> Diagnostic {
1447    envelope_failure(
1448        DiagnosticCategory::ResourceLimit,
1449        "query_remote_envelope_too_large",
1450        "remote reply exceeds the envelope byte ceiling",
1451    )
1452}
1453
1454#[derive(Deserialize)]
1455struct RemoteResponseOutcomePeek<'wire> {
1456    #[serde(borrow)]
1457    outcome: &'wire serde_json::value::RawValue,
1458}
1459
1460#[derive(Clone, Copy)]
1461enum SequenceScanKind {
1462    Rows,
1463    Documents,
1464}
1465
1466#[derive(Clone, Copy)]
1467enum ShapeLimitExceeded {
1468    Items,
1469    Width,
1470    Members,
1471    Outcome,
1472    Evidence,
1473}
1474
1475#[derive(Default)]
1476struct ShapeScanState {
1477    items: u64,
1478    members: u64,
1479    exceeded: Option<ShapeLimitExceeded>,
1480}
1481
1482fn preflight_remote_response_shape(
1483    payload: &[u8],
1484    limits: RemoteReplyDecodeLimits,
1485) -> Result<(), Diagnostic> {
1486    let response =
1487        serde_json::from_slice::<RemoteResponseOutcomePeek<'_>>(payload).map_err(|_| {
1488            envelope_failure(
1489                DiagnosticCategory::InvalidContract,
1490                "query_remote_reply_malformed",
1491                "remote response payload is malformed",
1492            )
1493        })?;
1494    let mut state = ShapeScanState::default();
1495    let mut deserializer = serde_json::Deserializer::from_slice(response.outcome.get().as_bytes());
1496    let result = match limits.shape {
1497        RemoteOutcomeShape::Rows { width } => OutcomeShapeSeed {
1498            field: "rows",
1499            kind: SequenceScanKind::Rows,
1500            width,
1501            max_items: limits
1502                .max_items
1503                .min(u64::try_from(MAX_CANONICAL_COLLECTION_LEN).unwrap_or(u64::MAX)),
1504            max_members: limits
1505                .max_collection_members
1506                .min(u64::try_from(MAX_CANONICAL_COLLECTION_LEN).unwrap_or(u64::MAX)),
1507            state: &mut state,
1508        }
1509        .deserialize(&mut deserializer),
1510        RemoteOutcomeShape::Documents { width } => OutcomeShapeSeed {
1511            field: "documents",
1512            kind: SequenceScanKind::Documents,
1513            width,
1514            max_items: limits
1515                .max_items
1516                .min(u64::try_from(MAX_CANONICAL_COLLECTION_LEN).unwrap_or(u64::MAX)),
1517            max_members: limits
1518                .max_collection_members
1519                .min(u64::try_from(MAX_CANONICAL_COLLECTION_LEN).unwrap_or(u64::MAX)),
1520            state: &mut state,
1521        }
1522        .deserialize(&mut deserializer),
1523        RemoteOutcomeShape::Count => ScalarOutcomeSeed {
1524            kind: ScalarScanKind::Count,
1525            max_items: limits.max_items,
1526            state: &mut state,
1527        }
1528        .deserialize(&mut deserializer),
1529        RemoteOutcomeShape::Exists => ScalarOutcomeSeed {
1530            kind: ScalarScanKind::Exists,
1531            max_items: limits.max_items,
1532            state: &mut state,
1533        }
1534        .deserialize(&mut deserializer),
1535    };
1536    if let Some(exceeded) = state.exceeded {
1537        return Err(match exceeded {
1538            ShapeLimitExceeded::Items => envelope_failure(
1539                DiagnosticCategory::ResourceLimit,
1540                "query_remote_response_oversized",
1541                "response rows, documents, or scalar evidence exceed the caller item budget",
1542            ),
1543            ShapeLimitExceeded::Width => envelope_failure(
1544                DiagnosticCategory::Integrity,
1545                "query_remote_evidence_mismatch",
1546                "response evidence does not conform to the validated output schema",
1547            ),
1548            ShapeLimitExceeded::Members => envelope_failure(
1549                DiagnosticCategory::ResourceLimit,
1550                "query_v2_document_member_limit",
1551                "document lists exceed the aggregate member ceiling",
1552            ),
1553            ShapeLimitExceeded::Outcome => envelope_failure(
1554                DiagnosticCategory::Integrity,
1555                "query_remote_outcome_mismatch",
1556                "response outcome kind does not match the invoked operation",
1557            ),
1558            ShapeLimitExceeded::Evidence => envelope_failure(
1559                DiagnosticCategory::Integrity,
1560                "query_remote_evidence_mismatch",
1561                "response evidence does not conform to the validated output schema",
1562            ),
1563        });
1564    }
1565    result.map_err(|_| {
1566        envelope_failure(
1567            DiagnosticCategory::InvalidContract,
1568            "query_remote_reply_malformed",
1569            "remote response outcome is malformed",
1570        )
1571    })
1572}
1573
1574#[derive(Clone, Copy)]
1575enum ScalarScanKind {
1576    Count,
1577    Exists,
1578}
1579
1580impl ScalarScanKind {
1581    const fn wire_name(self) -> &'static str {
1582        match self {
1583            Self::Count => "count",
1584            Self::Exists => "exists",
1585        }
1586    }
1587}
1588
1589struct ScalarOutcomeSeed<'scan> {
1590    kind: ScalarScanKind,
1591    max_items: u64,
1592    state: &'scan mut ShapeScanState,
1593}
1594
1595impl<'de> DeserializeSeed<'de> for ScalarOutcomeSeed<'_> {
1596    type Value = ();
1597
1598    fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
1599    where
1600        D: Deserializer<'de>,
1601    {
1602        deserializer.deserialize_map(ScalarOutcomeVisitor { seed: self })
1603    }
1604}
1605
1606struct ScalarOutcomeVisitor<'scan> {
1607    seed: ScalarOutcomeSeed<'scan>,
1608}
1609
1610impl<'de> Visitor<'de> for ScalarOutcomeVisitor<'_> {
1611    type Value = ();
1612
1613    fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1614        formatter.write_str("an exact count or exists remote outcome object")
1615    }
1616
1617    fn visit_map<M>(self, mut map: M) -> Result<Self::Value, M::Error>
1618    where
1619        M: MapAccess<'de>,
1620    {
1621        let mut found_kind = false;
1622        let mut found_value = false;
1623        while let Some(key) = map.next_key::<&str>()? {
1624            if key == "kind" {
1625                if found_kind {
1626                    self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1627                    return Err(serde::de::Error::custom(
1628                        "remote scalar outcome kind is duplicated",
1629                    ));
1630                }
1631                found_kind = true;
1632                if map.next_value::<&str>()? != self.seed.kind.wire_name() {
1633                    self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1634                    return Err(serde::de::Error::custom(
1635                        "remote scalar outcome kind does not match the expected operation",
1636                    ));
1637                }
1638            } else if key == "value" {
1639                if found_value {
1640                    self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1641                    return Err(serde::de::Error::custom(
1642                        "remote scalar outcome value is duplicated",
1643                    ));
1644                }
1645                found_value = true;
1646                let value = match self.seed.kind {
1647                    ScalarScanKind::Count => map.next_value::<u64>().map(|value| {
1648                        if value > self.seed.max_items {
1649                            self.seed.state.exceeded = Some(ShapeLimitExceeded::Items);
1650                        }
1651                    }),
1652                    ScalarScanKind::Exists => map.next_value::<bool>().map(|value| {
1653                        if value && self.seed.max_items == 0 {
1654                            self.seed.state.exceeded = Some(ShapeLimitExceeded::Items);
1655                        }
1656                    }),
1657                };
1658                if value.is_err() {
1659                    self.seed.state.exceeded = Some(ShapeLimitExceeded::Evidence);
1660                }
1661                value?;
1662            } else {
1663                self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1664                return Err(serde::de::Error::custom(
1665                    "remote scalar outcome carries an unexpected field",
1666                ));
1667            }
1668        }
1669        if found_kind && found_value {
1670            Ok(())
1671        } else {
1672            self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1673            Err(serde::de::Error::custom(
1674                "remote scalar outcome does not carry its exact fields",
1675            ))
1676        }
1677    }
1678}
1679
1680struct OutcomeShapeSeed<'scan> {
1681    field: &'static str,
1682    kind: SequenceScanKind,
1683    width: usize,
1684    max_items: u64,
1685    max_members: u64,
1686    state: &'scan mut ShapeScanState,
1687}
1688
1689impl<'de> DeserializeSeed<'de> for OutcomeShapeSeed<'_> {
1690    type Value = ();
1691
1692    fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
1693    where
1694        D: Deserializer<'de>,
1695    {
1696        deserializer.deserialize_map(OutcomeShapeVisitor { seed: self })
1697    }
1698}
1699
1700struct OutcomeShapeVisitor<'scan> {
1701    seed: OutcomeShapeSeed<'scan>,
1702}
1703
1704impl<'de> Visitor<'de> for OutcomeShapeVisitor<'_> {
1705    type Value = ();
1706
1707    fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1708        formatter.write_str("a remote outcome object")
1709    }
1710
1711    fn visit_map<M>(self, mut map: M) -> Result<Self::Value, M::Error>
1712    where
1713        M: MapAccess<'de>,
1714    {
1715        let mut found = false;
1716        let mut found_kind = false;
1717        while let Some(key) = map.next_key::<&str>()? {
1718            if key == self.seed.field {
1719                if found {
1720                    self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1721                    return Err(serde::de::Error::custom(
1722                        "remote outcome field is duplicated",
1723                    ));
1724                }
1725                found = true;
1726                map.next_value_seed(OutcomeSequenceSeed {
1727                    kind: self.seed.kind,
1728                    width: self.seed.width,
1729                    max_items: self.seed.max_items,
1730                    max_members: self.seed.max_members,
1731                    state: self.seed.state,
1732                })?;
1733            } else if key == "kind" {
1734                if found_kind {
1735                    self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1736                    return Err(serde::de::Error::custom(
1737                        "remote outcome kind is duplicated",
1738                    ));
1739                }
1740                found_kind = true;
1741                if map.next_value::<&str>()? != self.seed.field {
1742                    self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1743                    return Err(serde::de::Error::custom(
1744                        "remote outcome kind does not match the expected shape",
1745                    ));
1746                }
1747            } else {
1748                self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1749                return Err(serde::de::Error::custom(
1750                    "remote outcome carries an unexpected field",
1751                ));
1752            }
1753        }
1754        if found && found_kind {
1755            Ok(())
1756        } else {
1757            self.seed.state.exceeded = Some(ShapeLimitExceeded::Outcome);
1758            Err(serde::de::Error::custom(
1759                "remote outcome does not carry the expected field",
1760            ))
1761        }
1762    }
1763}
1764
1765struct OutcomeSequenceSeed<'scan> {
1766    kind: SequenceScanKind,
1767    width: usize,
1768    max_items: u64,
1769    max_members: u64,
1770    state: &'scan mut ShapeScanState,
1771}
1772
1773impl<'de> DeserializeSeed<'de> for OutcomeSequenceSeed<'_> {
1774    type Value = ();
1775
1776    fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
1777    where
1778        D: Deserializer<'de>,
1779    {
1780        deserializer.deserialize_seq(OutcomeSequenceVisitor { seed: self })
1781    }
1782}
1783
1784struct OutcomeSequenceVisitor<'scan> {
1785    seed: OutcomeSequenceSeed<'scan>,
1786}
1787
1788impl<'de> Visitor<'de> for OutcomeSequenceVisitor<'_> {
1789    type Value = ();
1790
1791    fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1792        formatter.write_str("a remote row or document sequence")
1793    }
1794
1795    fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
1796    where
1797        A: SeqAccess<'de>,
1798    {
1799        loop {
1800            let next = sequence.next_element_seed(OutputItemSeed {
1801                kind: self.seed.kind,
1802                width: self.seed.width,
1803                max_items: self.seed.max_items,
1804                max_members: self.seed.max_members,
1805                state: self.seed.state,
1806            })?;
1807            if next.is_none() {
1808                return Ok(());
1809            }
1810        }
1811    }
1812}
1813
1814struct OutputItemSeed<'scan> {
1815    kind: SequenceScanKind,
1816    width: usize,
1817    max_items: u64,
1818    max_members: u64,
1819    state: &'scan mut ShapeScanState,
1820}
1821
1822impl<'de> DeserializeSeed<'de> for OutputItemSeed<'_> {
1823    type Value = ();
1824
1825    fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
1826    where
1827        D: Deserializer<'de>,
1828    {
1829        self.state.items = self.state.items.saturating_add(1);
1830        if self.state.items > self.max_items {
1831            self.state.exceeded = Some(ShapeLimitExceeded::Items);
1832            return Err(serde::de::Error::custom("remote item budget exceeded"));
1833        }
1834        deserializer.deserialize_seq(OutputItemVisitor { seed: self })
1835    }
1836}
1837
1838struct OutputItemVisitor<'scan> {
1839    seed: OutputItemSeed<'scan>,
1840}
1841
1842impl<'de> Visitor<'de> for OutputItemVisitor<'_> {
1843    type Value = ();
1844
1845    fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1846        formatter.write_str("one positional remote output item")
1847    }
1848
1849    fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
1850    where
1851        A: SeqAccess<'de>,
1852    {
1853        let mut width = 0_usize;
1854        loop {
1855            let next = sequence.next_element_seed(OutputFieldSeed {
1856                kind: self.seed.kind,
1857                position: &mut width,
1858                max_width: self.seed.width,
1859                max_members: self.seed.max_members,
1860                state: self.seed.state,
1861            })?;
1862            if next.is_none() {
1863                if width == self.seed.width {
1864                    return Ok(());
1865                }
1866                self.seed.state.exceeded = Some(ShapeLimitExceeded::Width);
1867                return Err(serde::de::Error::custom(
1868                    "remote output width does not match the validated shape",
1869                ));
1870            }
1871        }
1872    }
1873}
1874
1875struct OutputFieldSeed<'scan> {
1876    kind: SequenceScanKind,
1877    position: &'scan mut usize,
1878    max_width: usize,
1879    max_members: u64,
1880    state: &'scan mut ShapeScanState,
1881}
1882
1883impl<'de> DeserializeSeed<'de> for OutputFieldSeed<'_> {
1884    type Value = ();
1885
1886    fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
1887    where
1888        D: Deserializer<'de>,
1889    {
1890        *self.position = self.position.saturating_add(1);
1891        if *self.position > self.max_width {
1892            self.state.exceeded = Some(ShapeLimitExceeded::Width);
1893            return Err(serde::de::Error::custom("remote output width exceeded"));
1894        }
1895        match self.kind {
1896            SequenceScanKind::Rows => IgnoredAny::deserialize(deserializer).map(|_| ()),
1897            SequenceScanKind::Documents => deserializer.deserialize_map(DocumentFieldVisitor {
1898                max_members: self.max_members,
1899                state: self.state,
1900            }),
1901        }
1902    }
1903}
1904
1905struct DocumentFieldVisitor<'scan> {
1906    max_members: u64,
1907    state: &'scan mut ShapeScanState,
1908}
1909
1910impl<'de> Visitor<'de> for DocumentFieldVisitor<'_> {
1911    type Value = ();
1912
1913    fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1914        formatter.write_str("one remote document field object")
1915    }
1916
1917    fn visit_map<M>(self, mut map: M) -> Result<Self::Value, M::Error>
1918    where
1919        M: MapAccess<'de>,
1920    {
1921        let mut kind = None;
1922        let mut value = false;
1923        let mut values = false;
1924        while let Some(key) = map.next_key::<&str>()? {
1925            if key == "kind" {
1926                if kind.is_some() {
1927                    self.state.exceeded = Some(ShapeLimitExceeded::Evidence);
1928                    return Err(serde::de::Error::custom(
1929                        "remote document field kind is duplicated",
1930                    ));
1931                }
1932                kind = Some(map.next_value::<&str>()?);
1933            } else if key == "value" {
1934                if value {
1935                    self.state.exceeded = Some(ShapeLimitExceeded::Evidence);
1936                    return Err(serde::de::Error::custom(
1937                        "remote document scalar value is duplicated",
1938                    ));
1939                }
1940                value = true;
1941                map.next_value::<IgnoredAny>()?;
1942            } else if key == "values" {
1943                if values {
1944                    self.state.exceeded = Some(ShapeLimitExceeded::Evidence);
1945                    return Err(serde::de::Error::custom(
1946                        "remote document list values are duplicated",
1947                    ));
1948                }
1949                values = true;
1950                map.next_value_seed(MemberSequenceSeed {
1951                    max_members: self.max_members,
1952                    state: self.state,
1953                })?;
1954            } else {
1955                self.state.exceeded = Some(ShapeLimitExceeded::Evidence);
1956                return Err(serde::de::Error::custom(
1957                    "remote document field carries an unexpected member",
1958                ));
1959            }
1960        }
1961        let exact = match kind {
1962            Some("absent") => !value && !values,
1963            Some("scalar") => value && !values,
1964            Some("list") => !value && values,
1965            Some(_) | None => false,
1966        };
1967        if exact {
1968            Ok(())
1969        } else {
1970            self.state.exceeded = Some(ShapeLimitExceeded::Evidence);
1971            Err(serde::de::Error::custom(
1972                "remote document field does not match its declared kind",
1973            ))
1974        }
1975    }
1976}
1977
1978struct MemberSequenceSeed<'scan> {
1979    max_members: u64,
1980    state: &'scan mut ShapeScanState,
1981}
1982
1983impl<'de> DeserializeSeed<'de> for MemberSequenceSeed<'_> {
1984    type Value = ();
1985
1986    fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
1987    where
1988        D: Deserializer<'de>,
1989    {
1990        deserializer.deserialize_seq(MemberSequenceVisitor { seed: self })
1991    }
1992}
1993
1994struct MemberSequenceVisitor<'scan> {
1995    seed: MemberSequenceSeed<'scan>,
1996}
1997
1998impl<'de> Visitor<'de> for MemberSequenceVisitor<'_> {
1999    type Value = ();
2000
2001    fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
2002        formatter.write_str("a remote document member sequence")
2003    }
2004
2005    fn visit_seq<A>(self, mut sequence: A) -> Result<Self::Value, A::Error>
2006    where
2007        A: SeqAccess<'de>,
2008    {
2009        loop {
2010            let next = sequence.next_element_seed(MemberSeed {
2011                max_members: self.seed.max_members,
2012                state: self.seed.state,
2013            })?;
2014            if next.is_none() {
2015                return Ok(());
2016            }
2017        }
2018    }
2019}
2020
2021struct MemberSeed<'scan> {
2022    max_members: u64,
2023    state: &'scan mut ShapeScanState,
2024}
2025
2026impl<'de> DeserializeSeed<'de> for MemberSeed<'_> {
2027    type Value = ();
2028
2029    fn deserialize<D>(self, deserializer: D) -> Result<Self::Value, D::Error>
2030    where
2031        D: Deserializer<'de>,
2032    {
2033        self.state.members = self.state.members.saturating_add(1);
2034        if self.state.members > self.max_members {
2035            self.state.exceeded = Some(ShapeLimitExceeded::Members);
2036            return Err(serde::de::Error::custom(
2037                "remote document member budget exceeded",
2038            ));
2039        }
2040        IgnoredAny::deserialize(deserializer).map(|_| ())
2041    }
2042}
2043
2044fn peek_remote_reply_binding(bytes: &[u8]) -> Result<RemoteReplyBindingPeek<'_>, Diagnostic> {
2045    if bytes.len() > MAX_REMOTE_ENVELOPE_BYTES {
2046        return Err(envelope_failure(
2047            DiagnosticCategory::ResourceLimit,
2048            "query_remote_envelope_too_large",
2049            "remote reply exceeds the envelope byte ceiling",
2050        ));
2051    }
2052    serde_json::from_slice(bytes).map_err(|_| {
2053        envelope_failure(
2054            DiagnosticCategory::InvalidContract,
2055            "query_remote_reply_malformed",
2056            "remote reply is not a JSON envelope with a format discriminator",
2057        )
2058    })
2059}
2060
2061fn verify_response_binding(
2062    nonce: Option<&str>,
2063    plan: Option<&str>,
2064    request: Option<&str>,
2065    expected_nonce: &str,
2066    expected_plan: &QueryPlanFingerprint,
2067    expected_request: &RemoteRequestFingerprint,
2068) -> Result<(), Diagnostic> {
2069    if nonce != Some(expected_nonce) {
2070        return Err(envelope_failure(
2071            DiagnosticCategory::Integrity,
2072            "query_remote_nonce_mismatch",
2073            "response evidence does not echo the request nonce",
2074        ));
2075    }
2076    let Some(plan) = plan else {
2077        return Err(envelope_failure(
2078            DiagnosticCategory::Integrity,
2079            "query_remote_plan_mismatch",
2080            "response evidence does not bind the invoked plan",
2081        ));
2082    };
2083    let echoed = FingerprintDigest::from_hex(plan)?;
2084    if echoed != expected_plan.as_fingerprint().digest() {
2085        return Err(envelope_failure(
2086            DiagnosticCategory::Integrity,
2087            "query_remote_plan_mismatch",
2088            "response evidence does not bind the invoked plan",
2089        ));
2090    }
2091    let Some(request) = request else {
2092        return Err(envelope_failure(
2093            DiagnosticCategory::Integrity,
2094            "query_remote_request_mismatch",
2095            "response evidence does not bind the exact request envelope",
2096        ));
2097    };
2098    let echoed_request = FingerprintDigest::from_hex(request)?;
2099    if echoed_request != expected_request.as_fingerprint().digest() {
2100        return Err(envelope_failure(
2101            DiagnosticCategory::Integrity,
2102            "query_remote_request_mismatch",
2103            "response evidence does not bind the exact request envelope",
2104        ));
2105    }
2106    Ok(())
2107}
2108
2109fn verify_failure_binding(
2110    nonce: Option<&str>,
2111    request: Option<&str>,
2112    expected_nonce: &str,
2113    expected_request: &RemoteRequestFingerprint,
2114) -> Result<(), Diagnostic> {
2115    let nonce = nonce.ok_or_else(|| {
2116        envelope_failure(
2117            DiagnosticCategory::Integrity,
2118            "query_remote_failure_unbound",
2119            "failure evidence is not bound to the request nonce",
2120        )
2121    })?;
2122    if nonce != expected_nonce {
2123        return Err(envelope_failure(
2124            DiagnosticCategory::Integrity,
2125            "query_remote_nonce_mismatch",
2126            "failure evidence does not echo the request nonce",
2127        ));
2128    }
2129    let request = request.ok_or_else(|| {
2130        envelope_failure(
2131            DiagnosticCategory::Integrity,
2132            "query_remote_failure_unbound",
2133            "failure evidence is not bound to the request envelope",
2134        )
2135    })?;
2136    let echoed = FingerprintDigest::from_hex(request)?;
2137    if echoed != expected_request.as_fingerprint().digest() {
2138        return Err(envelope_failure(
2139            DiagnosticCategory::Integrity,
2140            "query_remote_request_mismatch",
2141            "failure evidence does not bind the exact request envelope",
2142        ));
2143    }
2144    Ok(())
2145}
2146
2147/// Convert one caller-supplied limit into the unsigned wire range.
2148///
2149/// Every binding funnels its limit arguments through this exact
2150/// conversion, so a negative or out-of-range budget fails with one
2151/// stable diagnostic in every language instead of silently saturating
2152/// to zero or dropping the budget entirely.
2153pub fn checked_remote_limit(value: i128) -> Result<u64, Diagnostic> {
2154    u64::try_from(value).map_err(|_| remote_limit_invalid())
2155}
2156
2157/// The stable rejection every out-of-range limit argument maps to.
2158///
2159/// Exposed so bindings whose integer representation exceeds `i128`
2160/// (JavaScript `BigInt`) reject unrepresentable values with exactly
2161/// this diagnostic instead of inventing their own.
2162#[must_use]
2163pub fn remote_limit_invalid() -> Diagnostic {
2164    envelope_failure(
2165        DiagnosticCategory::InvalidContract,
2166        "query_remote_limit_invalid",
2167        "remote limits are unsigned 64-bit integers",
2168    )
2169}
2170
2171/// Convert one optional caller-supplied deadline into the wire range.
2172///
2173/// A negative deadline is rejected — never silently mapped to "no
2174/// deadline", which would remove the bound instead of enforcing it.
2175pub fn checked_remote_deadline(value: Option<i128>) -> Result<Option<u64>, Diagnostic> {
2176    value
2177        .map(|value| {
2178            let value = checked_remote_limit(value)?;
2179            if value > MAX_REMOTE_DEADLINE_MS {
2180                return Err(remote_deadline_limit());
2181            }
2182            Ok(value)
2183        })
2184        .transpose()
2185}
2186
2187fn validate_remote_limits(limits: RemoteLimits) -> Result<(), Diagnostic> {
2188    if limits
2189        .deadline_ms
2190        .is_some_and(|value| value > MAX_REMOTE_DEADLINE_MS)
2191    {
2192        return Err(remote_deadline_limit());
2193    }
2194    Ok(())
2195}
2196
2197fn validate_remote_time_shape(
2198    prepared_at_unix_ms: u64,
2199    expires_at_unix_ms: u64,
2200    limits: RemoteLimits,
2201) -> Result<(), Diagnostic> {
2202    let lifetime = expires_at_unix_ms
2203        .checked_sub(prepared_at_unix_ms)
2204        .ok_or_else(remote_time_invalid)?;
2205    let declared_lifetime = limits.deadline_ms.unwrap_or(DEFAULT_REMOTE_DEADLINE_MS);
2206    if lifetime != declared_lifetime {
2207        return Err(remote_time_invalid());
2208    }
2209    Ok(())
2210}
2211
2212fn remote_time_invalid() -> Diagnostic {
2213    envelope_failure(
2214        DiagnosticCategory::InvalidContract,
2215        "query_remote_time_invalid",
2216        "remote request timestamps do not form a bounded absolute lifetime",
2217    )
2218}
2219
2220/// Stable rejection for a deadline outside the supported monotonic range.
2221#[must_use]
2222pub fn remote_deadline_limit() -> Diagnostic {
2223    envelope_failure(
2224        DiagnosticCategory::ResourceLimit,
2225        "query_remote_deadline_limit",
2226        "remote deadline exceeds the maximum supported duration",
2227    )
2228}
2229
2230fn check_nonce(nonce: &str) -> Result<(), Diagnostic> {
2231    let valid = (NONCE_MIN_BYTES..=NONCE_MAX_BYTES).contains(&nonce.len())
2232        && nonce
2233            .bytes()
2234            .all(|byte| byte.is_ascii_alphanumeric() || byte == b'-');
2235    if valid {
2236        Ok(())
2237    } else {
2238        Err(envelope_failure(
2239            DiagnosticCategory::InvalidContract,
2240            "query_remote_nonce_invalid",
2241            "request nonces are 16-128 ASCII alphanumeric or dash bytes",
2242        ))
2243    }
2244}
2245
2246fn validate_executor_component(value: &str) -> Result<(), Diagnostic> {
2247    let valid = (EXECUTOR_COMPONENT_MIN_BYTES..=EXECUTOR_COMPONENT_MAX_BYTES)
2248        .contains(&value.len())
2249        && value
2250            .bytes()
2251            .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.'));
2252    if valid {
2253        Ok(())
2254    } else {
2255        Err(envelope_failure(
2256            DiagnosticCategory::InvalidContract,
2257            "query_remote_executor_invalid",
2258            "executor identity and epoch are 16-128 safe ASCII bytes",
2259        ))
2260    }
2261}
2262
2263fn encode_hex(bytes: &[u8]) -> String {
2264    const DIGITS: &[u8; 16] = b"0123456789abcdef";
2265    let mut encoded = String::with_capacity(bytes.len() * 2);
2266    for byte in bytes {
2267        encoded.push(char::from(DIGITS[usize::from(byte >> 4)]));
2268        encoded.push(char::from(DIGITS[usize::from(byte & 0x0f)]));
2269    }
2270    encoded
2271}
2272
2273fn decode_fixed_hex<const N: usize>(value: &str) -> Result<[u8; N], ()> {
2274    if value.len() != N * 2 || !value.bytes().all(|byte| byte.is_ascii_hexdigit()) {
2275        return Err(());
2276    }
2277    let mut decoded = [0_u8; N];
2278    for (index, pair) in value.as_bytes().as_chunks::<2>().0.iter().enumerate() {
2279        decoded[index] = (hex_nibble(pair[0])? << 4) | hex_nibble(pair[1])?;
2280    }
2281    Ok(decoded)
2282}
2283
2284fn hex_nibble(byte: u8) -> Result<u8, ()> {
2285    match byte {
2286        b'0'..=b'9' => Ok(byte - b'0'),
2287        b'a'..=b'f' => Ok(byte - b'a' + 10),
2288        b'A'..=b'F' => Ok(byte - b'A' + 10),
2289        _ => Err(()),
2290    }
2291}
2292
2293fn envelope_failure(
2294    category: DiagnosticCategory,
2295    code: &'static str,
2296    message: &'static str,
2297) -> Diagnostic {
2298    Diagnostic::new(
2299        category,
2300        DiagnosticCode::new(code).expect("static remote envelope code"),
2301        message,
2302    )
2303}
2304
2305#[cfg(test)]
2306mod tests {
2307    use super::*;
2308
2309    #[test]
2310    fn canonical_reply_binding_peek_borrows_correlation_strings() {
2311        let bytes = br#"{"format":"typebridge.query-remote-response/v1","nonce":"remote-nonce-0123456789abcdef","outcome":{"kind":"exists","value":true},"plan":"aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa","request":"bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"}"#;
2312        let peek: RemoteReplyBindingPeek<'_> = serde_json::from_slice(bytes).expect("binding peek");
2313
2314        assert!(bytes.as_ptr_range().contains(&peek.format.as_ptr()));
2315        assert!(
2316            peek.nonce
2317                .is_some_and(|value| bytes.as_ptr_range().contains(&value.as_ptr()))
2318        );
2319        assert!(
2320            peek.plan
2321                .is_some_and(|value| bytes.as_ptr_range().contains(&value.as_ptr()))
2322        );
2323        assert!(
2324            peek.request
2325                .is_some_and(|value| bytes.as_ptr_range().contains(&value.as_ptr()))
2326        );
2327    }
2328}