Skip to main content

type_bridge/
remote.rs

1#![deny(missing_docs)]
2//! Authenticated one-exchange remote execution for generated queries.
3
4use std::future::Future;
5use std::marker::PhantomData;
6use std::pin::Pin;
7use std::sync::Arc;
8
9use type_bridge_contract::query_remote::RemoteCapabilities;
10use type_bridge_orm::_registry::DescriptorRegistry;
11use type_bridge_orm::query_v2_prepared::QueryAuthority;
12use type_bridge_orm::{
13    AnswerCancellation, InstalledRuntimeProjection, QueryExecutionDeadline,
14    QueryExecutionResourceLimits, RemoteModelQueryV2Error, ValidatedMatchRequest,
15    ValidatedMatchResult, lower_remote_query_diagnostic, prepare_remote_model_query_v2_with_budget,
16};
17
18use crate::Result;
19use crate::error::Error;
20use crate::query::QuerySession;
21use crate::schema::{Schema, SchemaPackage, Unbound};
22
23/// One caller-owned asynchronous transport for the authenticated V2 routes.
24///
25/// Implementations fetch the exact `/v2/capabilities` bytes once at connect
26/// time and perform exactly one `/v2/query` exchange per terminal.
27pub trait RemoteQueryTransport: Send + Sync + 'static {
28    /// Fetch the executor's exact signed capability advertisement.
29    fn capabilities(&self) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + '_>>;
30
31    /// Exchange one exact canonical request for one exact signed reply.
32    fn exchange<'a>(
33        &'a self,
34        request: &'a [u8],
35    ) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + 'a>>;
36}
37
38/// Explicit immutable budgets for one remote generated-query terminal.
39#[derive(Clone, Copy, Debug, Eq, PartialEq)]
40pub struct RemoteQueryLimits {
41    resources: QueryExecutionResourceLimits,
42}
43
44impl RemoteQueryLimits {
45    /// Construct one explicit remote response and hydration budget.
46    #[must_use]
47    pub const fn new(
48        max_items: u64,
49        max_bytes: u64,
50        max_collection_members: u64,
51        max_graph_nodes: u64,
52        max_attribute_values: u64,
53        max_role_players: u64,
54    ) -> Self {
55        Self {
56            resources: QueryExecutionResourceLimits::tightened(
57                type_bridge_orm::MAX_QUERY_TIMEOUT_MILLISECONDS,
58                max_items,
59                max_bytes,
60                max_graph_nodes,
61                max_attribute_values,
62                max_collection_members,
63                max_role_players,
64                type_bridge_orm::MAX_QUERY_STATEMENTS,
65            ),
66        }
67    }
68
69    /// Attach an optional executor deadline in milliseconds.
70    #[must_use]
71    pub const fn deadline_ms(mut self, deadline_ms: u64) -> Self {
72        self.resources.timeout_milliseconds =
73            if deadline_ms < type_bridge_orm::MAX_QUERY_TIMEOUT_MILLISECONDS {
74                deadline_ms
75            } else {
76                type_bridge_orm::MAX_QUERY_TIMEOUT_MILLISECONDS
77            };
78        self
79    }
80}
81
82impl From<RemoteQueryLimits> for QueryExecutionResourceLimits {
83    fn from(limits: RemoteQueryLimits) -> Self {
84        limits.resources
85    }
86}
87
88impl From<QueryExecutionResourceLimits> for RemoteQueryLimits {
89    fn from(resources: QueryExecutionResourceLimits) -> Self {
90        Self {
91            resources: resources.effective(),
92        }
93    }
94}
95
96/// Connection-time authority, transport, and limit configuration.
97pub struct RemoteConnectionOptions {
98    scope: Option<String>,
99    semantic_profile: Option<String>,
100    resources: QueryExecutionResourceLimits,
101    transport: Arc<dyn RemoteQueryTransport>,
102    advertisement: Option<Vec<u8>>,
103}
104
105impl RemoteConnectionOptions {
106    /// Construct remote options for one managed schema scope and semantic
107    /// profile.
108    #[must_use]
109    pub fn new(
110        scope: impl Into<String>,
111        semantic_profile: impl Into<String>,
112        limits: impl Into<QueryExecutionResourceLimits>,
113        transport: impl RemoteQueryTransport,
114    ) -> Self {
115        Self {
116            scope: Some(scope.into()),
117            semantic_profile: Some(semantic_profile.into()),
118            resources: limits.into().effective(),
119            transport: Arc::new(transport),
120            advertisement: None,
121        }
122    }
123
124    /// Construct normal generated-package options; schema scope and semantic
125    /// profile are derived from [`SchemaPackage`] during binding.
126    #[must_use]
127    pub fn generated(
128        limits: impl Into<QueryExecutionResourceLimits>,
129        transport: impl RemoteQueryTransport,
130    ) -> Self {
131        Self {
132            scope: None,
133            semantic_profile: None,
134            resources: limits.into().effective(),
135            transport: Arc::new(transport),
136            advertisement: None,
137        }
138    }
139}
140
141struct RemoteRuntime {
142    advertisement: Vec<u8>,
143    authority: Arc<QueryAuthority>,
144    resources: QueryExecutionResourceLimits,
145    transport: Arc<dyn RemoteQueryTransport>,
146}
147
148/// A client-owned remote generated-query database branded by schema `S`.
149pub struct RemoteDatabase<S: Schema = Unbound> {
150    options: Option<RemoteConnectionOptions>,
151    runtime: Option<Arc<RemoteRuntime>>,
152    installed: Option<Arc<InstalledRuntimeProjection>>,
153    registry: Option<Arc<DescriptorRegistry>>,
154    marker: PhantomData<fn() -> S>,
155}
156
157impl<S: Schema> std::fmt::Debug for RemoteDatabase<S> {
158    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
159        formatter
160            .debug_struct("RemoteDatabase")
161            .field("schema_bound", &self.installed.is_some())
162            .finish_non_exhaustive()
163    }
164}
165
166impl RemoteDatabase<Unbound> {
167    /// Fetch and validate one immutable executor advertisement.
168    pub async fn connect(mut options: RemoteConnectionOptions) -> Result<Self> {
169        let advertisement = options.transport.capabilities().await?;
170        RemoteCapabilities::decode(&advertisement).map_err(remote_diagnostic)?;
171        options.advertisement = Some(advertisement);
172        Ok(Self {
173            options: Some(options),
174            runtime: None,
175            installed: None,
176            registry: None,
177            marker: PhantomData,
178        })
179    }
180
181    /// Verify and bind one generated schema package and its remote authority.
182    pub fn with_schema<S: Schema>(mut self, schema: SchemaPackage<S>) -> Result<RemoteDatabase<S>> {
183        let (installed, embedded_authority) = schema.verify_and_install_with_authority()?;
184        let registry = Arc::new(installed.match_registry().map_err(Error::from_orm)?);
185        let options = self.options.take().ok_or_else(|| Error::Other {
186            message: "remote connection options are unavailable".into(),
187            source: None,
188        })?;
189        let authority = if let Some(embedded) = embedded_authority {
190            let embedded_scope = embedded.managed_scope().id().as_str();
191            let embedded_profile = embedded.semantic_profile().id().as_str();
192            if options
193                .scope
194                .as_deref()
195                .is_some_and(|scope| scope != embedded_scope)
196                || options
197                    .semantic_profile
198                    .as_deref()
199                    .is_some_and(|profile| profile != embedded_profile)
200            {
201                return Err(Error::SchemaVerification {
202                    message: "remote options disagree with generated schema authority".into(),
203                    source: None,
204                });
205            }
206            let declared =
207                type_bridge_contract::schema::encode_declared_schema(embedded.declared_schema())
208                    .map_err(|error| Error::SchemaVerification {
209                        message: "verified generated authority cannot reconstruct its declaration"
210                            .into(),
211                        source: Some(Box::new(error)),
212                    })?;
213            QueryAuthority::from_declared_bytes(&declared, embedded_scope, embedded_profile)
214                .map_err(remote_diagnostic)?
215        } else {
216            let declared =
217                schema
218                    .declared_schema_json()
219                    .ok_or_else(|| Error::SchemaVerification {
220                        message: "generated schema package omits remote declared-schema authority"
221                            .into(),
222                        source: None,
223                    })?;
224            let scope = options
225                .scope
226                .as_deref()
227                .ok_or_else(|| Error::SchemaVerification {
228                    message: "schema package has no embedded managed scope".into(),
229                    source: None,
230                })?;
231            let semantic_profile =
232                options
233                    .semantic_profile
234                    .as_deref()
235                    .ok_or_else(|| Error::SchemaVerification {
236                        message: "schema package has no embedded semantic profile".into(),
237                        source: None,
238                    })?;
239            QueryAuthority::from_declared_bytes(declared.as_bytes(), scope, semantic_profile)
240                .map_err(remote_diagnostic)?
241        };
242        if !authority.matches_semantic_fingerprint(installed.projection().semantic_fingerprint()) {
243            return Err(Error::SchemaVerification {
244                message: "remote declared-schema authority does not match the generated projection"
245                    .into(),
246                source: None,
247            });
248        }
249        let runtime = Arc::new(RemoteRuntime {
250            advertisement: options.advertisement.ok_or_else(|| Error::Other {
251                message: "remote capability advertisement is unavailable".into(),
252                source: None,
253            })?,
254            authority: Arc::new(authority),
255            resources: options.resources,
256            transport: options.transport,
257        });
258        Ok(RemoteDatabase {
259            options: None,
260            runtime: Some(runtime),
261            installed: Some(installed),
262            registry: Some(registry),
263            marker: PhantomData,
264        })
265    }
266}
267
268impl<S: Schema> RemoteDatabase<S> {
269    /// Start one owner-branded query session over this remote executor.
270    pub fn query(&self) -> Result<QuerySession<'_, S>> {
271        let runtime = self.runtime.as_ref().ok_or_else(remote_not_bound)?;
272        self.query_with_resources(runtime.resources, AnswerCancellation::default())
273    }
274
275    /// Start one remote query session with one common tighten-only resource
276    /// policy and caller-owned cooperative cancellation signal.
277    pub fn query_with_resources(
278        &self,
279        resources: QueryExecutionResourceLimits,
280        cancellation: AnswerCancellation,
281    ) -> Result<QuerySession<'_, S>> {
282        let installed = self.installed.as_deref().ok_or_else(remote_not_bound)?;
283        let registry = self.registry.as_ref().ok_or_else(remote_not_bound)?;
284        let runtime = self.runtime.as_ref().ok_or_else(remote_not_bound)?;
285        Ok(QuerySession::remote(
286            installed,
287            Arc::clone(registry),
288            self,
289            resources.constrained_by(runtime.resources),
290            cancellation,
291        ))
292    }
293
294    pub(crate) async fn execute_match(
295        &self,
296        registry: &DescriptorRegistry,
297        validated: ValidatedMatchRequest,
298        resources: QueryExecutionResourceLimits,
299        cancellation: AnswerCancellation,
300        deadline: QueryExecutionDeadline,
301    ) -> Result<(ValidatedMatchRequest, ValidatedMatchResult)> {
302        let runtime = self.runtime.as_ref().ok_or_else(remote_not_bound)?;
303        check_remote_execution_budget(&cancellation, deadline)?;
304        let pending = prepare_remote_model_query_v2_with_budget(
305            &runtime.authority,
306            registry,
307            validated,
308            &runtime.advertisement,
309            resources.remote(),
310            deadline,
311            &cancellation,
312        )
313        .map_err(remote_model_input_error)?;
314        check_remote_execution_budget(&cancellation, deadline)?;
315        let request = pending.request_bytes().to_vec();
316        let response = await_remote_exchange(
317            runtime.transport.as_ref(),
318            &request,
319            &cancellation,
320            deadline,
321        )
322        .await?;
323        check_remote_execution_budget(&cancellation, deadline)?;
324        let claimed = pending
325            .claim_reply_with_cancellation(&cancellation)
326            .map_err(remote_model_hydration_error)?;
327        if response.len() > claimed.response_snapshot_limit() {
328            return Err(Error::classified(
329                crate::ErrorCategory::ResourceLimit,
330                None,
331                "remote_response_limit",
332                Vec::new(),
333                "remote query reply exceeds the authenticated response ceiling",
334                None,
335            ));
336        }
337        let (request, result, _registry) = claimed
338            .decode_with_cancellation(&response, &cancellation)
339            .map_err(remote_model_hydration_error)?;
340        check_remote_execution_budget(&cancellation, deadline)?;
341        Ok((request, result))
342    }
343}
344
345async fn await_remote_exchange(
346    transport: &dyn RemoteQueryTransport,
347    request: &[u8],
348    cancellation: &AnswerCancellation,
349    deadline: QueryExecutionDeadline,
350) -> Result<Vec<u8>> {
351    check_remote_execution_budget(cancellation, deadline)?;
352    let exchange = transport.exchange(request);
353    tokio::pin!(exchange);
354    let cancellation_wait = cancellation.cancelled();
355    tokio::pin!(cancellation_wait);
356    let timeout = tokio::time::sleep_until(tokio::time::Instant::from_std(deadline.instant()));
357    tokio::pin!(timeout);
358    let result = tokio::select! {
359        biased;
360        result = &mut exchange => result,
361        () = &mut cancellation_wait => Err(remote_cancelled()),
362        () = &mut timeout => Err(remote_timeout()),
363    };
364    check_remote_execution_budget(cancellation, deadline)?;
365    result
366}
367
368fn check_remote_execution_budget(
369    cancellation: &AnswerCancellation,
370    deadline: QueryExecutionDeadline,
371) -> Result<()> {
372    deadline
373        .check(cancellation)
374        .map_err(|error| Error::from_sdk_execution(error, crate::ModelValidationPhase::Input))
375}
376
377fn remote_cancelled() -> Error {
378    Error::classified(
379        crate::ErrorCategory::Cancelled,
380        None,
381        "provider_cancelled",
382        Vec::new(),
383        "query execution was cancelled",
384        None,
385    )
386}
387
388fn remote_timeout() -> Error {
389    Error::classified(
390        crate::ErrorCategory::ResourceLimit,
391        None,
392        "transaction_deadline_exceeded",
393        Vec::new(),
394        "query execution exceeded its timeout",
395        None,
396    )
397}
398
399fn remote_not_bound() -> Error {
400    Error::ModelValidation {
401        phase: crate::ModelValidationPhase::Input,
402        code: "schema_not_bound".into(),
403        path: vec![],
404        message: "remote database is not schema-bound".into(),
405        source: None,
406    }
407}
408
409fn remote_diagnostic(error: type_bridge_contract::diagnostic::Diagnostic) -> Error {
410    Error::from_sdk_execution(
411        lower_remote_query_diagnostic(error),
412        crate::ModelValidationPhase::Input,
413    )
414}
415
416fn remote_model_input_error(error: RemoteModelQueryV2Error) -> Error {
417    match error {
418        RemoteModelQueryV2Error::Diagnostic(error) => Error::from_sdk_execution(
419            lower_remote_query_diagnostic(error),
420            crate::ModelValidationPhase::Input,
421        ),
422        RemoteModelQueryV2Error::Match(error) => {
423            Error::from_match(error, crate::ModelValidationPhase::Input)
424        }
425    }
426}
427
428fn remote_model_hydration_error(error: RemoteModelQueryV2Error) -> Error {
429    match error {
430        RemoteModelQueryV2Error::Diagnostic(error) => Error::from_sdk_execution(
431            lower_remote_query_diagnostic(error),
432            crate::ModelValidationPhase::Hydration,
433        ),
434        RemoteModelQueryV2Error::Match(error) => {
435            Error::from_match(error, crate::ModelValidationPhase::Hydration)
436        }
437    }
438}
439
440#[cfg(test)]
441mod tests {
442    use std::env;
443    use std::fs::{self, OpenOptions};
444    use std::io::Write as _;
445    use std::path::{Path, PathBuf};
446    use std::sync::Mutex;
447    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering as AtomicOrdering};
448
449    use sha2::{Digest as _, Sha256};
450    use tokio::sync::Notify;
451    use type_bridge_contract::capability::{CapabilityId, CapabilitySet};
452    use type_bridge_contract::codec::to_canonical_json;
453    use type_bridge_contract::diagnostic::{
454        Diagnostic, DiagnosticCategory, DiagnosticCode, DiagnosticPath, DiagnosticPathSegment,
455    };
456    use type_bridge_contract::fingerprint::SemanticProfileId;
457    use type_bridge_contract::managed_scope::ManagedScopeId;
458    use type_bridge_contract::migration_assertion::BindingId;
459    use type_bridge_contract::projection::{BindingTarget, ProjectionConfig};
460    use type_bridge_contract::query_plan::{ModelQueryV2, query_plan_v2_capability_vocabulary};
461    use type_bridge_contract::query_remote::RemoteExecutorBinding;
462    use type_bridge_contract::query_remote_v2::{
463        CAP_QUERY_REMOTE_STATEMENT_LIMIT, HydrationGraphV2, RemoteOutcomeV2, RemoteQueryFailureV2,
464        RemoteQueryRequestV2, RemoteQueryResponseV2, RemoteReducedValueV2, RemoteReductionRowV2,
465        RemoteResultKindV2, query_remote_v2_required_capabilities,
466    };
467    use type_bridge_contract::schema::{DocumentId, encode_declared_schema};
468    use type_bridge_orm::OrmError;
469    use type_bridge_orm::match_request::CapabilitySet as MatchCapabilitySet;
470    use type_bridge_orm::match_request::SessionHandle;
471    use type_bridge_orm::query_v2_remote::RemoteReplySigningKey;
472    use type_bridge_orm::session::backend::{
473        BoxFuture, DriverBackend, QueryResult, TransactionOps, TxType,
474    };
475    use type_bridge_schema::{
476        ManagedDeltaContext, SchemaDocumentSet, build_schema_authority, encode_schema_authority,
477        normalize_documents, project, resolve,
478    };
479    use type_bridge_schema_codegen::RustEmitter;
480
481    use super::*;
482    use crate::__codegen::{
483        self, CompleteModel, EncodedCreate, EntityModel, HydratedRow, HydrationCapability,
484        IntoEncodedCreate, MaterializeModel, Model, ThingModel, ValidationError,
485    };
486    use crate::schema::sealed;
487
488    struct TestSchema;
489    impl sealed::Sealed for TestSchema {}
490    impl Schema for TestSchema {}
491
492    #[derive(Debug)]
493    struct Person;
494    impl sealed::Sealed for Person {}
495    impl Model for Person {
496        type Schema = TestSchema;
497        const TYPE_ID_JSON: &'static str = r#"{"kind":"entity","label":"person"}"#;
498    }
499    impl ThingModel for Person {
500        fn thing_kind() -> __codegen::ThingKind {
501            __codegen::ThingKind::Entity
502        }
503    }
504    impl EntityModel for Person {}
505    impl CompleteModel for Person {
506        type Create = PersonCreate;
507
508        fn iid(&self) -> &str {
509            unreachable!()
510        }
511    }
512    impl MaterializeModel for Person {
513        fn materialize(
514            _: &HydratedRow,
515            _: &HydrationCapability,
516        ) -> std::result::Result<Self, ValidationError> {
517            Ok(Self)
518        }
519    }
520
521    #[derive(Clone)]
522    struct PersonCreate;
523    impl sealed::Sealed for PersonCreate {}
524    impl IntoEncodedCreate for PersonCreate {
525        fn into_encoded_create(self) -> std::result::Result<EncodedCreate, ValidationError> {
526            Ok(EncodedCreate::new(Person::TYPE_ID_JSON, vec![], vec![]))
527        }
528    }
529
530    #[derive(Default)]
531    struct DirectCancellationState {
532        entered: Notify,
533        opens: AtomicUsize,
534        statements: AtomicUsize,
535        closes: AtomicUsize,
536    }
537
538    struct DirectCancellationBackend {
539        state: Arc<DirectCancellationState>,
540    }
541
542    impl DriverBackend for DirectCancellationBackend {
543        fn match_capabilities(&self) -> MatchCapabilitySet {
544            MatchCapabilitySet::all()
545        }
546
547        fn open_transaction(
548            &self,
549            _database: &str,
550            _tx_type: TxType,
551        ) -> BoxFuture<'_, std::result::Result<Box<dyn TransactionOps>, OrmError>> {
552            let state = Arc::clone(&self.state);
553            Box::pin(async move {
554                state.opens.fetch_add(1, AtomicOrdering::SeqCst);
555                Ok(Box::new(DirectCancellationTransaction { state }) as Box<dyn TransactionOps>)
556            })
557        }
558
559        fn is_open(&self) -> bool {
560            true
561        }
562    }
563
564    struct DirectCancellationTransaction {
565        state: Arc<DirectCancellationState>,
566    }
567
568    impl TransactionOps for DirectCancellationTransaction {
569        fn query(
570            &mut self,
571            _typeql: &str,
572        ) -> BoxFuture<'_, std::result::Result<QueryResult, OrmError>> {
573            let state = Arc::clone(&self.state);
574            Box::pin(async move {
575                state.statements.fetch_add(1, AtomicOrdering::SeqCst);
576                state.entered.notify_one();
577                std::future::pending().await
578            })
579        }
580
581        fn commit(&mut self) -> BoxFuture<'_, std::result::Result<(), OrmError>> {
582            Box::pin(async { Ok(()) })
583        }
584
585        fn rollback(&mut self) -> BoxFuture<'_, std::result::Result<(), OrmError>> {
586            Box::pin(async { Ok(()) })
587        }
588
589        fn close(&mut self) -> BoxFuture<'_, std::result::Result<(), OrmError>> {
590            self.state.closes.fetch_add(1, AtomicOrdering::SeqCst);
591            Box::pin(async { Ok(()) })
592        }
593    }
594
595    fn direct_cancellation_database(
596        state: Arc<DirectCancellationState>,
597    ) -> crate::session::Database<TestSchema> {
598        let generated = package();
599        let installed = InstalledRuntimeProjection::from_verified_rust_json(
600            generated.runtime_projection_json().as_bytes(),
601            generated.semantic_fingerprint_json().as_bytes(),
602            generated.projection_fingerprint_json().as_bytes(),
603        )
604        .expect("Rust cancellation proof installs the generated projection");
605        crate::session::Database::from_test_parts(
606            type_bridge_orm::Database::with_backend(
607                Box::new(DirectCancellationBackend { state }),
608                "rust-sdk-v2-proof",
609            ),
610            installed,
611        )
612    }
613
614    async fn observe_sdk_v2_direct_cancellation() -> serde_json::Value {
615        let pre_state = Arc::new(DirectCancellationState::default());
616        let pre_database = direct_cancellation_database(Arc::clone(&pre_state));
617        let pre_signal = AnswerCancellation::default();
618        pre_signal.cancel();
619        let mut pre_session = pre_database
620            .query_with_resources(QueryExecutionResourceLimits::default(), pre_signal)
621            .expect("Rust direct cancellation proof opens a generated session");
622        let pre_person = pre_session
623            .exact::<Person>()
624            .expect("Rust direct cancellation proof binds a generated model");
625        let pre_result = pre_session
626            .query(pre_person)
627            .expect("Rust direct cancellation proof authors a generated query")
628            .count()
629            .await;
630        let pre_partial_result = pre_result.is_ok();
631        let pre_error = pre_result.expect_err("pre-cancellation rejects the generated terminal");
632        let pre_provider_calls = pre_state.opens.load(AtomicOrdering::SeqCst)
633            + pre_state.statements.load(AtomicOrdering::SeqCst);
634        assert_eq!(pre_error.category(), crate::ErrorCategory::Cancelled);
635        assert_eq!(pre_error.code(), Some("provider_cancelled"));
636        assert!(!pre_partial_result);
637        assert_eq!(pre_provider_calls, 0);
638
639        let in_flight_state = Arc::new(DirectCancellationState::default());
640        let in_flight_database = direct_cancellation_database(Arc::clone(&in_flight_state));
641        let in_flight_signal = AnswerCancellation::default();
642        let requester_state = Arc::clone(&in_flight_state);
643        let requester_signal = in_flight_signal.clone();
644        let requester = tokio::spawn(async move {
645            tokio::time::timeout(
646                std::time::Duration::from_secs(5),
647                requester_state.entered.notified(),
648            )
649            .await
650            .expect("direct provider await is reached before cancellation");
651            requester_signal.cancel();
652            true
653        });
654        let mut in_flight_session = in_flight_database
655            .query_with_resources(QueryExecutionResourceLimits::default(), in_flight_signal)
656            .expect("Rust in-flight cancellation proof opens a generated session");
657        let in_flight_person = in_flight_session
658            .exact::<Person>()
659            .expect("Rust in-flight cancellation proof binds a generated model");
660        let in_flight_result = tokio::time::timeout(
661            std::time::Duration::from_secs(5),
662            in_flight_session
663                .query(in_flight_person)
664                .expect("Rust in-flight cancellation proof authors a generated query")
665                .count(),
666        )
667        .await
668        .expect("cancellation wakes the pending provider await");
669        let in_flight_partial_result = in_flight_result.is_ok();
670        let in_flight_error =
671            in_flight_result.expect_err("in-flight cancellation rejects the generated terminal");
672        let provider_await_woken = requester
673            .await
674            .expect("direct cancellation requester completes")
675            && in_flight_state.statements.load(AtomicOrdering::SeqCst) == 1;
676        assert_eq!(in_flight_error.category(), crate::ErrorCategory::Cancelled);
677        assert_eq!(in_flight_error.code(), Some("provider_cancelled"));
678        assert!(!in_flight_partial_result);
679        assert!(provider_await_woken);
680        assert_eq!(in_flight_state.opens.load(AtomicOrdering::SeqCst), 1);
681        assert_eq!(in_flight_state.closes.load(AtomicOrdering::SeqCst), 1);
682
683        serde_json::json!({
684            "in_flight": {
685                "category": in_flight_error.category().as_str(),
686                "code": in_flight_error.code().expect("cancelled errors carry a stable code"),
687                "partial_result": in_flight_partial_result,
688                "provider_await_woken": provider_await_woken,
689            },
690            "pre_dispatch": {
691                "category": pre_error.category().as_str(),
692                "code": pre_error.code().expect("cancelled errors carry a stable code"),
693                "partial_result": pre_partial_result,
694                "provider_calls": pre_provider_calls,
695            },
696        })
697    }
698
699    #[test]
700    fn local_and_remote_failures_preserve_classification_codes_and_paths() {
701        let session = SessionHandle::new(Arc::new(DescriptorRegistry::new()));
702        let match_error = match session.exact("missing") {
703            Err(OrmError::Match(error)) => error,
704            Err(other) => panic!("unexpected ORM error: {other:?}"),
705            Ok(_) => panic!("missing descriptor unexpectedly resolved"),
706        };
707        let local = Error::from_orm(OrmError::Match(match_error.clone()));
708        let remote = remote_model_input_error(RemoteModelQueryV2Error::Match(match_error));
709
710        assert_eq!(local.category(), crate::ErrorCategory::QueryAuthoring);
711        assert_eq!(remote.category(), local.category());
712        assert_eq!(remote.code(), Some("unknown_descriptor"));
713        assert_eq!(remote.code(), local.code());
714        assert_eq!(remote.path(), local.path());
715        assert_eq!(
716            remote.model_validation_phase(),
717            local.model_validation_phase()
718        );
719
720        let diagnostic = Diagnostic::new(
721            DiagnosticCategory::UnsupportedCapability,
722            DiagnosticCode::new("missing_remote_capability").unwrap(),
723            "the remote executor does not advertise one required capability",
724        )
725        .at(DiagnosticPathSegment::Field("capabilities".into()))
726        .at(DiagnosticPathSegment::Index(2));
727        let classified = remote_diagnostic(diagnostic);
728
729        assert_eq!(classified.category(), crate::ErrorCategory::Capability);
730        assert_eq!(classified.code(), Some("missing_remote_capability"));
731        assert_eq!(
732            classified.path(),
733            Some(&["capabilities".to_owned(), "[2]".to_owned()][..])
734        );
735        assert_eq!(classified.model_validation_phase(), None);
736    }
737
738    #[test]
739    fn remote_model_resource_limits_use_released_generated_error_codes() {
740        let internal = Diagnostic::new(
741            DiagnosticCategory::ResourceLimit,
742            DiagnosticCode::new("query_v2_model_role_player_limit").unwrap(),
743            "internal model execution detail",
744        );
745        let error = remote_model_hydration_error(RemoteModelQueryV2Error::Diagnostic(internal));
746
747        assert_eq!(error.category(), crate::ErrorCategory::ResourceLimit);
748        assert_eq!(error.code(), Some("hydrated_role_player_limit"));
749        assert_eq!(error.path(), Some(&["provider_evidence".to_owned()][..]));
750        assert_eq!(
751            error.diagnostic_path(),
752            Some(
753                &[crate::ErrorPathSegment::Query(
754                    crate::QueryDiagnosticPathKind::ProviderEvidence,
755                )][..]
756            )
757        );
758        assert_eq!(error.model_validation_phase(), None);
759        assert_eq!(
760            error
761                .details()
762                .and_then(|details| details.get("query_category")),
763            Some(&crate::ErrorDetail::QueryCategory(
764                crate::QueryDiagnosticCategory::ResourceLimit,
765            ))
766        );
767    }
768
769    fn package() -> SchemaPackage<TestSchema> {
770        let documents = SchemaDocumentSet::parse([(
771            DocumentId::new("remote.yaml").unwrap(),
772            "format: typebridge.schema/v2\nattributes:\n  name: { value: string }\nentities:\n  person:\n    owns: { name: { key: true } }\n",
773        )])
774        .unwrap();
775        let declared = normalize_documents(&documents).unwrap();
776        let profile = SemanticProfileId::new("typedb-3.12.1/v1").unwrap();
777        let resolved = resolve(&declared, &profile).unwrap();
778        let authority = build_schema_authority(
779            &declared,
780            declared.required_capabilities(),
781            &ManagedDeltaContext::new(
782                ManagedScopeId::new("rust-client-test").unwrap(),
783                profile,
784                CapabilitySet::new(),
785            ),
786        )
787        .unwrap();
788        let emitter = RustEmitter::new();
789        let projection = project(
790            &resolved,
791            BindingTarget::Rust,
792            &ProjectionConfig::rust(),
793            &emitter.generator_handlers(),
794            &emitter.code_resources().unwrap(),
795        )
796        .unwrap();
797        let leak = |bytes: Vec<u8>| {
798            Box::leak(String::from_utf8(bytes).unwrap().into_boxed_str()) as &'static str
799        };
800        SchemaPackage::new_with_authority(
801            leak(to_canonical_json(projection.semantic_fingerprint()).unwrap()),
802            leak(to_canonical_json(projection.projection_fingerprint()).unwrap()),
803            leak(to_canonical_json(&projection).unwrap()),
804            leak(encode_schema_authority(&authority)),
805            leak(encode_declared_schema(&declared).unwrap()),
806            "rust-client-test",
807            "typedb-3.12.1/v1",
808        )
809    }
810
811    fn released_declared_package() -> SchemaPackage<TestSchema> {
812        let generated = package();
813        SchemaPackage::new_with_declared(
814            generated.semantic_fingerprint_json(),
815            generated.projection_fingerprint_json(),
816            generated.runtime_projection_json(),
817            generated
818                .declared_schema_json()
819                .expect("test package carries a declaration"),
820        )
821    }
822
823    struct UnusedTransport;
824
825    impl RemoteQueryTransport for UnusedTransport {
826        fn capabilities(&self) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + '_>> {
827            Box::pin(async { panic!("compatibility test performs no transport I/O") })
828        }
829
830        fn exchange<'a>(
831            &'a self,
832            _request: &'a [u8],
833        ) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + 'a>> {
834            Box::pin(async { panic!("compatibility test performs no transport I/O") })
835        }
836    }
837
838    fn unconnected_remote(mut options: RemoteConnectionOptions) -> RemoteDatabase<Unbound> {
839        options.advertisement = Some(Vec::new());
840        RemoteDatabase {
841            options: Some(options),
842            runtime: None,
843            installed: None,
844            registry: None,
845            marker: PhantomData,
846        }
847    }
848
849    #[test]
850    fn generated_authority_accepts_matching_legacy_options_and_rejects_overrides() {
851        let limits = || RemoteQueryLimits::new(10, 1 << 20, 10, 100, 100, 100);
852        let matching = RemoteConnectionOptions::new(
853            "rust-client-test",
854            "typedb-3.12.1/v1",
855            limits(),
856            UnusedTransport,
857        );
858        unconnected_remote(matching)
859            .with_schema(package())
860            .expect("matching 2.0.1-style options remain compatible");
861
862        for mismatched in [
863            RemoteConnectionOptions::new(
864                "other-scope",
865                "typedb-3.12.1/v1",
866                limits(),
867                UnusedTransport,
868            ),
869            RemoteConnectionOptions::new(
870                "rust-client-test",
871                "typedb-3.11.5/v1",
872                limits(),
873                UnusedTransport,
874            ),
875        ] {
876            let error = unconnected_remote(mismatched)
877                .with_schema(package())
878                .expect_err("caller strings cannot override generated authority");
879            assert!(
880                error
881                    .to_string()
882                    .contains("disagree with generated schema authority"),
883                "{error}"
884            );
885        }
886    }
887
888    #[test]
889    fn released_declared_package_retains_explicit_remote_options_compatibility() {
890        let limits = || RemoteQueryLimits::new(10, 1 << 20, 10, 100, 100, 100);
891        let options = RemoteConnectionOptions::new(
892            "rust-client-test",
893            "typedb-3.12.1/v1",
894            limits(),
895            UnusedTransport,
896        );
897        unconnected_remote(options)
898            .with_schema(released_declared_package())
899            .expect("2.0.1 generated package and explicit options remain compatible");
900
901        let generated_options = RemoteConnectionOptions::generated(limits(), UnusedTransport);
902        let error = unconnected_remote(generated_options)
903            .with_schema(released_declared_package())
904            .expect_err("detached 2.0.1 package cannot invent embedded deployment authority");
905        assert!(error.to_string().contains("no embedded managed scope"));
906    }
907
908    struct Transport {
909        advertisement_contract: RemoteCapabilities,
910        advertisement: Vec<u8>,
911        capabilities: Arc<Mutex<usize>>,
912        exchanges: Arc<Mutex<Vec<Vec<u8>>>>,
913        failure: Option<Diagnostic>,
914        signer: RemoteReplySigningKey,
915    }
916
917    impl RemoteQueryTransport for Transport {
918        fn capabilities(&self) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + '_>> {
919            *self.capabilities.lock().unwrap() += 1;
920            let bytes = self.advertisement.clone();
921            Box::pin(async move { Ok(bytes) })
922        }
923
924        fn exchange<'a>(
925            &'a self,
926            request: &'a [u8],
927        ) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + 'a>> {
928            self.exchanges.lock().unwrap().push(request.to_vec());
929            let response = (|| {
930                let request = RemoteQueryRequestV2::decode(request).map_err(remote_diagnostic)?;
931                request
932                    .validate_advertisement(&self.advertisement_contract)
933                    .map_err(remote_diagnostic)?;
934                if let Some(diagnostic) = &self.failure {
935                    return RemoteQueryFailureV2::bound(
936                        request.nonce(),
937                        &request.fingerprint().map_err(remote_diagnostic)?,
938                        diagnostic,
939                    )
940                    .and_then(|failure| {
941                        failure.encode_signed(
942                            &self.advertisement_contract.fingerprint()?,
943                            &self.signer,
944                        )
945                    })
946                    .map_err(remote_diagnostic);
947                }
948                let plan = request.plan().map_err(remote_diagnostic)?;
949                let root = BindingId::new(0).map_err(remote_diagnostic)?;
950                let outcome = match request.result_kind() {
951                    RemoteResultKindV2::DistinctCount => {
952                        RemoteOutcomeV2::DistinctCount { root, value: 7 }
953                    }
954                    RemoteResultKindV2::DistinctExists => {
955                        RemoteOutcomeV2::DistinctExists { root, value: true }
956                    }
957                    RemoteResultKindV2::HydratedRows => RemoteOutcomeV2::HydratedRows {
958                        graph: HydrationGraphV2::new(vec![]).map_err(remote_diagnostic)?,
959                        rows: vec![],
960                    },
961                    RemoteResultKindV2::HydratedPage => RemoteOutcomeV2::HydratedPage {
962                        entries: vec![],
963                        graph: HydrationGraphV2::new(vec![]).map_err(remote_diagnostic)?,
964                        limit: 2,
965                        offset: 0,
966                        root,
967                        total: Some(0),
968                    },
969                    RemoteResultKindV2::ModelReduction => {
970                        let Some(ModelQueryV2::Reduction {
971                            root,
972                            group,
973                            reducers,
974                            ..
975                        }) = plan
976                            .v2_compatibility()
977                            .and_then(|compatibility| compatibility.model_query())
978                        else {
979                            return Err(Error::Other {
980                                message: "test reduction lacks its model contract".into(),
981                                source: None,
982                            });
983                        };
984                        RemoteOutcomeV2::ModelReduction {
985                            graph: HydrationGraphV2::new(vec![]).map_err(remote_diagnostic)?,
986                            root: *root,
987                            group: group.clone(),
988                            reducers: reducers.clone(),
989                            rows: vec![RemoteReductionRowV2::new(
990                                None,
991                                vec![RemoteReducedValueV2::Count { value: 7 }],
992                            )],
993                        }
994                    }
995                    _ => {
996                        return Err(Error::Other {
997                            message: "test transport received an unexpected terminal".into(),
998                            source: None,
999                        });
1000                    }
1001                };
1002                RemoteQueryResponseV2::new(
1003                    request.nonce(),
1004                    &plan,
1005                    &request.fingerprint().map_err(remote_diagnostic)?,
1006                    request.result_kind(),
1007                    outcome,
1008                )
1009                .and_then(|response| {
1010                    response
1011                        .encode_signed(&self.advertisement_contract.fingerprint()?, &self.signer)
1012                })
1013                .map_err(remote_diagnostic)
1014            })();
1015            Box::pin(async move { response })
1016        }
1017    }
1018
1019    struct CancelBeforeDecodeTransport {
1020        inner: Transport,
1021        cancellation: AnswerCancellation,
1022        response_completed: Arc<AtomicBool>,
1023    }
1024
1025    impl RemoteQueryTransport for CancelBeforeDecodeTransport {
1026        fn capabilities(&self) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + '_>> {
1027            self.inner.capabilities()
1028        }
1029
1030        fn exchange<'a>(
1031            &'a self,
1032            request: &'a [u8],
1033        ) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + 'a>> {
1034            Box::pin(async move {
1035                let response = self.inner.exchange(request).await?;
1036                self.response_completed.store(true, AtomicOrdering::SeqCst);
1037                self.cancellation.cancel();
1038                Ok(response)
1039            })
1040        }
1041    }
1042
1043    struct ExchangeDropProbe(Arc<AtomicBool>);
1044
1045    impl Drop for ExchangeDropProbe {
1046        fn drop(&mut self) {
1047            self.0.store(true, AtomicOrdering::SeqCst);
1048        }
1049    }
1050
1051    struct CallerAbortTransport {
1052        advertisement: Vec<u8>,
1053        exchanges: Arc<AtomicUsize>,
1054        entered: Arc<Notify>,
1055        exchange_dropped: Arc<AtomicBool>,
1056    }
1057
1058    impl RemoteQueryTransport for CallerAbortTransport {
1059        fn capabilities(&self) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + '_>> {
1060            let advertisement = self.advertisement.clone();
1061            Box::pin(async move { Ok(advertisement) })
1062        }
1063
1064        fn exchange<'a>(
1065            &'a self,
1066            _request: &'a [u8],
1067        ) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + 'a>> {
1068            self.exchanges.fetch_add(1, AtomicOrdering::SeqCst);
1069            self.entered.notify_one();
1070            let probe = ExchangeDropProbe(Arc::clone(&self.exchange_dropped));
1071            Box::pin(async move {
1072                let _probe = probe;
1073                std::future::pending().await
1074            })
1075        }
1076    }
1077
1078    fn sdk_v2_transport(seed: u8) -> (Transport, Arc<Mutex<Vec<Vec<u8>>>>) {
1079        let signer = RemoteReplySigningKey::from_secret_bytes([seed; 32]);
1080        let mut capabilities = query_plan_v2_capability_vocabulary();
1081        for capability in query_remote_v2_required_capabilities(true) {
1082            capabilities.insert(capability);
1083        }
1084        capabilities.insert(CapabilityId::new(CAP_QUERY_REMOTE_STATEMENT_LIMIT).unwrap());
1085        let advertisement_contract = RemoteCapabilities::new(
1086            capabilities,
1087            RemoteExecutorBinding::new("rust-client-test", "epoch-00000000003").unwrap(),
1088            signer.public_key(),
1089        );
1090        let advertisement = advertisement_contract.encode().unwrap();
1091        let exchanges = Arc::new(Mutex::new(Vec::new()));
1092        (
1093            Transport {
1094                advertisement_contract,
1095                advertisement,
1096                capabilities: Arc::new(Mutex::new(0)),
1097                exchanges: Arc::clone(&exchanges),
1098                failure: None,
1099                signer,
1100            },
1101            exchanges,
1102        )
1103    }
1104
1105    async fn observe_sdk_v2_remote_cancellation() -> serde_json::Value {
1106        let (pre_transport, pre_exchanges) = sdk_v2_transport(0x71);
1107        let remote = RemoteDatabase::connect(RemoteConnectionOptions::generated(
1108            QueryExecutionResourceLimits::default(),
1109            pre_transport,
1110        ))
1111        .await
1112        .expect("Rust remote cancellation proof connects")
1113        .with_schema(package())
1114        .expect("Rust remote cancellation proof binds generated schema authority");
1115        let pre_signal = AnswerCancellation::default();
1116        pre_signal.cancel();
1117        let mut pre_session = remote
1118            .query_with_resources(QueryExecutionResourceLimits::default(), pre_signal)
1119            .expect("Rust pre-cancelled remote session opens");
1120        let pre_person = pre_session
1121            .exact::<Person>()
1122            .expect("Rust pre-cancelled remote model binds");
1123        let pre_result = pre_session
1124            .query(pre_person)
1125            .expect("Rust pre-cancelled remote query authors")
1126            .count()
1127            .await;
1128        let pre_partial_result = pre_result.is_ok();
1129        let pre_error = pre_result.expect_err("pre-cancellation rejects before remote exchange");
1130        let pre_exchange_count = pre_exchanges.lock().unwrap().len();
1131        assert_eq!(pre_error.category(), crate::ErrorCategory::Cancelled);
1132        assert_eq!(pre_error.code(), Some("provider_cancelled"));
1133        assert_eq!(pre_exchange_count, 0);
1134        assert!(!pre_partial_result);
1135
1136        let decode_signal = AnswerCancellation::default();
1137        let response_completed = Arc::new(AtomicBool::new(false));
1138        let (decode_inner, decode_exchanges) = sdk_v2_transport(0x72);
1139        let decode_remote = RemoteDatabase::connect(RemoteConnectionOptions::generated(
1140            QueryExecutionResourceLimits::default(),
1141            CancelBeforeDecodeTransport {
1142                inner: decode_inner,
1143                cancellation: decode_signal.clone(),
1144                response_completed: Arc::clone(&response_completed),
1145            },
1146        ))
1147        .await
1148        .expect("Rust decode-cancel remote connects")
1149        .with_schema(package())
1150        .expect("Rust decode-cancel remote binds generated schema authority");
1151        let mut decode_session = decode_remote
1152            .query_with_resources(QueryExecutionResourceLimits::default(), decode_signal)
1153            .expect("Rust decode-cancel session opens");
1154        let decode_person = decode_session
1155            .exact::<Person>()
1156            .expect("Rust decode-cancel model binds");
1157        let decode_result = decode_session
1158            .query(decode_person)
1159            .expect("Rust decode-cancel query authors")
1160            .count()
1161            .await;
1162        let decode_partial_result = decode_result.is_ok();
1163        let decode_error = decode_result.expect_err("cancellation rejects before reply decode");
1164        let decode_exchange_count = decode_exchanges.lock().unwrap().len();
1165        assert_eq!(decode_error.category(), crate::ErrorCategory::Cancelled);
1166        assert_eq!(decode_error.code(), Some("provider_cancelled"));
1167        assert_eq!(decode_exchange_count, 1);
1168        assert!(!decode_partial_result);
1169
1170        let (abort_contract_transport, _) = sdk_v2_transport(0x73);
1171        let abort_exchanges = Arc::new(AtomicUsize::new(0));
1172        let abort_entered = Arc::new(Notify::new());
1173        let abort_dropped = Arc::new(AtomicBool::new(false));
1174        let abort_transport = CallerAbortTransport {
1175            advertisement: abort_contract_transport.advertisement,
1176            exchanges: Arc::clone(&abort_exchanges),
1177            entered: Arc::clone(&abort_entered),
1178            exchange_dropped: Arc::clone(&abort_dropped),
1179        };
1180        let abort_remote = RemoteDatabase::connect(RemoteConnectionOptions::generated(
1181            QueryExecutionResourceLimits::default(),
1182            abort_transport,
1183        ))
1184        .await
1185        .expect("Rust caller-abort remote connects")
1186        .with_schema(package())
1187        .expect("Rust caller-abort remote binds generated schema authority");
1188        let abort_signal = AnswerCancellation::default();
1189        let requester_signal = abort_signal.clone();
1190        let requester_entered = Arc::clone(&abort_entered);
1191        let requester = tokio::spawn(async move {
1192            tokio::time::timeout(
1193                std::time::Duration::from_secs(5),
1194                requester_entered.notified(),
1195            )
1196            .await
1197            .expect("caller transport exchange starts before cancellation");
1198            requester_signal.cancel();
1199        });
1200        let mut abort_session = abort_remote
1201            .query_with_resources(QueryExecutionResourceLimits::default(), abort_signal)
1202            .expect("Rust caller-abort session opens");
1203        let abort_person = abort_session
1204            .exact::<Person>()
1205            .expect("Rust caller-abort model binds");
1206        let abort_error = tokio::time::timeout(
1207            std::time::Duration::from_secs(5),
1208            abort_session
1209                .query(abort_person)
1210                .expect("Rust caller-abort query authors")
1211                .count(),
1212        )
1213        .await
1214        .expect("caller cancellation wakes the remote exchange")
1215        .expect_err("caller cancellation rejects the remote terminal");
1216        requester
1217            .await
1218            .expect("caller transport cancellation requester completes");
1219        let caller_transport_abort_supported = abort_error.category()
1220            == crate::ErrorCategory::Cancelled
1221            && abort_error.code() == Some("provider_cancelled")
1222            && abort_exchanges.load(AtomicOrdering::SeqCst) == 1
1223            && abort_dropped.load(AtomicOrdering::SeqCst);
1224        assert!(caller_transport_abort_supported);
1225
1226        serde_json::json!({
1227            "before_exchange": {
1228                "category": pre_error.category().as_str(),
1229                "code": pre_error.code().expect("cancelled errors carry a stable code"),
1230                "exchange_count": pre_exchange_count,
1231                "partial_result": pre_partial_result,
1232            },
1233            "during_decode": {
1234                "category": decode_error.category().as_str(),
1235                "code": decode_error.code().expect("cancelled errors carry a stable code"),
1236                "exchange_count": decode_exchange_count,
1237                "partial_result": decode_partial_result,
1238            },
1239            "caller_transport_abort_supported": caller_transport_abort_supported,
1240            "server_exchange_cancelled_after_send": !response_completed.load(AtomicOrdering::SeqCst),
1241        })
1242    }
1243
1244    fn sdk_v2_query_category(error: &crate::Error) -> &'static str {
1245        error
1246            .details()
1247            .and_then(|details| {
1248                details.values().find_map(|detail| match detail {
1249                    crate::ErrorDetail::QueryCategory(
1250                        crate::QueryDiagnosticCategory::ResultDecode,
1251                    ) => Some("result_decode"),
1252                    _ => None,
1253                })
1254            })
1255            .expect("authenticated remote diagnostic carries result_decode")
1256    }
1257
1258    fn sdk_v2_remote_diagnostic_observation(
1259        error: &crate::Error,
1260        claim_consumed: bool,
1261    ) -> serde_json::Value {
1262        let path = error
1263            .diagnostic_path()
1264            .expect("authenticated remote diagnostic carries typed path")
1265            .iter()
1266            .map(|segment| match segment {
1267                crate::ErrorPathSegment::ContractField(value) => {
1268                    serde_json::json!({"kind": "contract_field", "value": value})
1269                }
1270                crate::ErrorPathSegment::Index(value) => {
1271                    serde_json::json!({"kind": "index", "value": value})
1272                }
1273                crate::ErrorPathSegment::ContractIdentity(value) => {
1274                    serde_json::json!({"kind": "contract_identity", "value": value})
1275                }
1276                other => panic!("unexpected authenticated proof path segment: {other:?}"),
1277            })
1278            .collect::<Vec<_>>();
1279        let mut details = serde_json::Map::new();
1280        for (name, detail) in error
1281            .details()
1282            .expect("authenticated remote diagnostic carries typed details")
1283        {
1284            let value = match detail {
1285                crate::ErrorDetail::Long(value) => {
1286                    serde_json::json!({"kind": "signed", "value": value.to_string()})
1287                }
1288                crate::ErrorDetail::Boolean(value) => {
1289                    serde_json::json!({"kind": "boolean", "value": value})
1290                }
1291                crate::ErrorDetail::QueryIdentity(value) => {
1292                    serde_json::json!({"kind": "query_identity", "value": value})
1293                }
1294                crate::ErrorDetail::QueryIdentityList(values) => {
1295                    serde_json::json!({"kind": "query_identity_list", "value": values})
1296                }
1297                crate::ErrorDetail::QueryCategory(_) => continue,
1298                other => panic!("unexpected authenticated proof detail: {other:?}"),
1299            };
1300            assert!(details.insert(name.clone(), value).is_none());
1301        }
1302        let visible = format!("{}{:?}{:?}", error.message(), path, details);
1303        let redacted = !visible.contains("provider-secret") && !visible.contains("must-not-cross");
1304        serde_json::json!({
1305            "category": match (error.category(), sdk_v2_query_category(error)) {
1306                (crate::ErrorCategory::ModelValidation, "result_decode") => "integrity",
1307                (category, _) => category.as_str(),
1308            },
1309            "query_category": sdk_v2_query_category(error),
1310            "code": error.code().expect("authenticated remote diagnostic carries a stable code"),
1311            "message": error.message(),
1312            "path": path,
1313            "details": details,
1314            "redacted": redacted,
1315            "claim_consumed": claim_consumed,
1316        })
1317    }
1318
1319    async fn observe_sdk_v2_remote_structured_diagnostic() -> serde_json::Value {
1320        let signer = RemoteReplySigningKey::from_secret_bytes([0x42; 32]);
1321        let mut capabilities = query_plan_v2_capability_vocabulary();
1322        for capability in query_remote_v2_required_capabilities(true) {
1323            capabilities.insert(capability);
1324        }
1325        let advertisement_contract = RemoteCapabilities::new(
1326            capabilities,
1327            RemoteExecutorBinding::new("rust-generated-acceptance", "epoch-00000000002").unwrap(),
1328            signer.public_key(),
1329        );
1330        let advertisement = advertisement_contract.encode().unwrap();
1331        let exchanges = Arc::new(Mutex::new(Vec::new()));
1332        let failure = Diagnostic::new(
1333            DiagnosticCategory::Integrity,
1334            DiagnosticCode::new("remote_application_failure").unwrap(),
1335            "provider-secret remote query literal and endpoint",
1336        )
1337        .with_path(DiagnosticPath::from_segments([
1338            DiagnosticPathSegment::Field("plan".into()),
1339            DiagnosticPathSegment::Index(2),
1340            DiagnosticPathSegment::Identifier("person".into()),
1341        ]))
1342        .with_detail("attempt", -7_i64)
1343        .with_detail("expected", vec!["person".to_owned(), "employee".to_owned()])
1344        .with_detail("retryable", false)
1345        .with_detail("subject", "person")
1346        .with_detail("provider_secret", "must-not-cross");
1347        let transport = Transport {
1348            advertisement_contract,
1349            advertisement,
1350            capabilities: Arc::new(Mutex::new(0)),
1351            exchanges: Arc::clone(&exchanges),
1352            failure: Some(failure),
1353            signer,
1354        };
1355        let remote = RemoteDatabase::connect(RemoteConnectionOptions::generated(
1356            RemoteQueryLimits::new(10, 1 << 20, 10, 100, 100, 100),
1357            transport,
1358        ))
1359        .await
1360        .expect("Rust structured-diagnostic proof connects")
1361        .with_schema(package())
1362        .expect("Rust structured-diagnostic proof binds generated schema authority");
1363        let mut session = remote
1364            .query()
1365            .expect("Rust structured-diagnostic proof opens a generated session");
1366        let person = session
1367            .exact::<Person>()
1368            .expect("Rust structured-diagnostic proof binds a generated model");
1369        let query = session
1370            .query(person)
1371            .expect("Rust structured-diagnostic proof authors a generated query");
1372        let error = query
1373            .one()
1374            .await
1375            .expect_err("signed application failure reaches the generated query facade");
1376        let exchange_count = exchanges.lock().unwrap().len();
1377        let claim_consumed = exchange_count == 1
1378            && error.code() == Some("remote_application_failure")
1379            && error.details().is_some();
1380        assert_eq!(exchange_count, 1);
1381        assert!(claim_consumed);
1382        let observation = sdk_v2_remote_diagnostic_observation(&error, claim_consumed);
1383        assert_eq!(observation["redacted"], true);
1384        observation
1385    }
1386
1387    fn sdk_v2_proof_source(root: &Path, relative: &str) -> serde_json::Value {
1388        let path = root.join(relative);
1389        let bytes = fs::read(&path)
1390            .unwrap_or_else(|error| panic!("proof source {} is readable: {error}", path.display()));
1391        serde_json::json!({
1392            "path": relative,
1393            "sha256": format!("{:x}", Sha256::digest(bytes)),
1394        })
1395    }
1396
1397    fn requested_sdk_v2_proof_fragment() -> Option<(PathBuf, String)> {
1398        let destination = env::var_os("TYPE_BRIDGE_SDK_V2_PROOF_FRAGMENT");
1399        let run_nonce = env::var_os("TYPE_BRIDGE_SDK_V2_PROOF_RUN_NONCE");
1400        assert_eq!(
1401            destination.is_some(),
1402            run_nonce.is_some(),
1403            "Rust proof fragment output and run nonce must be configured together",
1404        );
1405        let destination = PathBuf::from(destination?);
1406        assert!(
1407            destination.is_absolute(),
1408            "Rust proof destination must be absolute"
1409        );
1410        let parent = destination
1411            .parent()
1412            .expect("Rust proof destination has a parent");
1413        let parent_metadata =
1414            fs::symlink_metadata(parent).expect("Rust proof destination parent exists");
1415        assert!(
1416            parent_metadata.is_dir() && !parent_metadata.file_type().is_symlink(),
1417            "Rust proof destination parent must be a non-symlink directory"
1418        );
1419        assert!(
1420            matches!(fs::symlink_metadata(&destination), Err(error) if error.kind() == std::io::ErrorKind::NotFound),
1421            "Rust proof destination must not already exist"
1422        );
1423        let run_nonce = run_nonce
1424            .expect("Rust proof run nonce exists")
1425            .into_string()
1426            .expect("Rust proof run nonce must be UTF-8");
1427        assert!(
1428            run_nonce.len() == 64
1429                && run_nonce
1430                    .bytes()
1431                    .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)),
1432            "Rust proof run nonce must be 64 lowercase hexadecimal digits"
1433        );
1434        Some((destination, run_nonce))
1435    }
1436
1437    fn publish_sdk_v2_rust_proof_fragment(
1438        destination: &Path,
1439        run_nonce: &str,
1440        results: Vec<serde_json::Value>,
1441    ) {
1442        let core = Path::new(env!("CARGO_MANIFEST_DIR"))
1443            .parent()
1444            .and_then(Path::parent)
1445            .expect("Rust crate lives beneath type-bridge-core");
1446        let repository = core
1447            .parent()
1448            .expect("type-bridge-core has a repository parent");
1449        let sources = ["type-bridge-core/crates/rust/src/remote.rs"];
1450        let fragment = serde_json::json!({
1451            "format": "typebridge.sdk-v2-proof-fragment/v1",
1452            "binding": "rust",
1453            "semantic_profile": "typedb-3.12.1/v1",
1454            "run_nonce": run_nonce,
1455            "contract": {
1456                "allowlist": sdk_v2_proof_source(repository, "tests/contracts/sdk_conformance/sdk-v2/proof-fragment-allowlist-v1.json"),
1457                "journey": sdk_v2_proof_source(repository, "tests/contracts/sdk_conformance/sdk-v2/journey-v2.json"),
1458                "proof_schema": sdk_v2_proof_source(repository, "tests/contracts/sdk_conformance/sdk-v2/proof-fragment-schema-v1.json"),
1459            },
1460            "producer": {
1461                "id": "type-bridge-rust.generated-query-proof",
1462                "sources": sources
1463                    .iter()
1464                    .map(|source| sdk_v2_proof_source(repository, source))
1465                    .collect::<Vec<_>>(),
1466            },
1467            "results": results,
1468        });
1469        let mut bytes =
1470            to_canonical_json(&fragment).expect("Rust sdk-v2 proof fragment canonicalizes");
1471        bytes.push(b'\n');
1472        assert!(bytes.len() <= 64 * 1024);
1473        let mut output = OpenOptions::new()
1474            .write(true)
1475            .create_new(true)
1476            .open(destination)
1477            .expect("Rust sdk-v2 proof fragment is created once");
1478        output
1479            .write_all(&bytes)
1480            .expect("Rust sdk-v2 proof fragment is written completely");
1481        output
1482            .sync_all()
1483            .expect("Rust sdk-v2 proof fragment is durable");
1484    }
1485
1486    #[tokio::test]
1487    async fn remote_database_fetches_capabilities_once_and_exchanges_once_per_terminal() {
1488        let signer = RemoteReplySigningKey::from_secret_bytes([0x31; 32]);
1489        let mut capabilities = query_plan_v2_capability_vocabulary();
1490        for capability in query_remote_v2_required_capabilities(true) {
1491            capabilities.insert(capability);
1492        }
1493        capabilities.insert(CapabilityId::new(CAP_QUERY_REMOTE_STATEMENT_LIMIT).unwrap());
1494        let advertisement_contract = RemoteCapabilities::new(
1495            capabilities,
1496            RemoteExecutorBinding::new("rust-client-test", "epoch-00000000001").unwrap(),
1497            signer.public_key(),
1498        );
1499        let advertisement = advertisement_contract.encode().unwrap();
1500        let capability_calls = Arc::new(Mutex::new(0));
1501        let exchanges = Arc::new(Mutex::new(Vec::new()));
1502        let transport = Transport {
1503            advertisement_contract,
1504            advertisement,
1505            capabilities: Arc::clone(&capability_calls),
1506            exchanges: Arc::clone(&exchanges),
1507            failure: None,
1508            signer,
1509        };
1510        let connection_resources =
1511            QueryExecutionResourceLimits::tightened(20_000, 10, 1_000_000, 90, 80, 70, 60, 2);
1512        let options = RemoteConnectionOptions::generated(connection_resources, transport);
1513        let remote = RemoteDatabase::connect(options)
1514            .await
1515            .unwrap()
1516            .with_schema(package())
1517            .unwrap();
1518        let mut session = remote.query().unwrap();
1519        let person = session.exact::<Person>().unwrap();
1520        let query = session.query(person).unwrap();
1521
1522        assert_eq!(query.count().await.unwrap(), 7);
1523        assert!(query.exists().await.unwrap());
1524        assert!(
1525            query
1526                .rows(crate::RowsOptions::new(2))
1527                .await
1528                .unwrap()
1529                .is_empty()
1530        );
1531        let page = query
1532            .page_by(person, crate::PageOptions::new(2).include_total(true))
1533            .await
1534            .unwrap();
1535        assert!(page.items().is_empty());
1536        assert_eq!(page.total(), Some(0));
1537        let reduction = query
1538            .aggregate((crate::aggregate::count(),))
1539            .await
1540            .expect("typed reductions use the same authenticated exchange");
1541        assert_eq!(reduction, (7,));
1542        assert_eq!(*capability_calls.lock().unwrap(), 1);
1543
1544        let zero_timeout = QueryExecutionResourceLimits {
1545            timeout_milliseconds: 0,
1546            ..QueryExecutionResourceLimits::default()
1547        };
1548        let mut timed_session = remote
1549            .query_with_resources(zero_timeout, AnswerCancellation::default())
1550            .unwrap();
1551        let timed_person = timed_session.exact::<Person>().unwrap();
1552        let timed_error = timed_session
1553            .query(timed_person)
1554            .unwrap()
1555            .count()
1556            .await
1557            .expect_err("zero timeout must reject before caller transport");
1558        assert_eq!(timed_error.category(), crate::ErrorCategory::ResourceLimit);
1559        assert_eq!(timed_error.code(), Some("transaction_deadline_exceeded"));
1560
1561        let cancellation = AnswerCancellation::default();
1562        cancellation.cancel();
1563        let mut cancelled_session = remote
1564            .query_with_resources(QueryExecutionResourceLimits::default(), cancellation)
1565            .unwrap();
1566        let cancelled_person = cancelled_session.exact::<Person>().unwrap();
1567        let cancelled_error = cancelled_session
1568            .query(cancelled_person)
1569            .unwrap()
1570            .count()
1571            .await
1572            .expect_err("pre-cancellation must reject before caller transport");
1573        assert_eq!(cancelled_error.category(), crate::ErrorCategory::Cancelled);
1574        assert_eq!(cancelled_error.code(), Some("provider_cancelled"));
1575
1576        let terminal_resources =
1577            QueryExecutionResourceLimits::tightened(10_000, 20, 900_000, 100, 70, 80, 50, 1);
1578        let mut tightened_session = remote
1579            .query_with_resources(terminal_resources, AnswerCancellation::default())
1580            .unwrap();
1581        let tightened_person = tightened_session.exact::<Person>().unwrap();
1582        assert_eq!(
1583            tightened_session
1584                .query(tightened_person)
1585                .unwrap()
1586                .count()
1587                .await
1588                .unwrap(),
1589            7,
1590        );
1591
1592        let requests = exchanges.lock().unwrap();
1593        assert_eq!(requests.len(), 6);
1594        assert!(
1595            std::str::from_utf8(&requests[0])
1596                .unwrap()
1597                .contains("\"format\":\"typebridge.query-remote-request/v2\"")
1598        );
1599        let request = RemoteQueryRequestV2::decode(&requests[5]).unwrap();
1600        let limits = request.limits();
1601        assert!(
1602            limits
1603                .deadline_ms
1604                .is_some_and(|deadline| deadline <= 10_000),
1605            "the wire carries the remaining portion of the tighter terminal timeout",
1606        );
1607        assert_eq!(limits.max_items, 10);
1608        assert_eq!(limits.max_bytes, 900_000);
1609        assert_eq!(limits.max_graph_nodes, 90);
1610        assert_eq!(limits.max_attribute_values, 70);
1611        assert_eq!(limits.max_collection_members, 70);
1612        assert_eq!(limits.max_role_players, 50);
1613        assert_eq!(limits.max_statements, 1);
1614    }
1615
1616    #[tokio::test]
1617    async fn generated_remote_query_preserves_complete_authenticated_structured_diagnostic() {
1618        let signer = RemoteReplySigningKey::from_secret_bytes([0x42; 32]);
1619        let mut capabilities = query_plan_v2_capability_vocabulary();
1620        for capability in query_remote_v2_required_capabilities(true) {
1621            capabilities.insert(capability);
1622        }
1623        let advertisement_contract = RemoteCapabilities::new(
1624            capabilities,
1625            RemoteExecutorBinding::new("rust-generated-acceptance", "epoch-00000000002").unwrap(),
1626            signer.public_key(),
1627        );
1628        let advertisement = advertisement_contract.encode().unwrap();
1629        let capability_calls = Arc::new(Mutex::new(0));
1630        let exchanges = Arc::new(Mutex::new(Vec::new()));
1631        let diagnostic = Diagnostic::new(
1632            DiagnosticCategory::Integrity,
1633            DiagnosticCode::new("remote_application_failure").unwrap(),
1634            "the remote application rejected this query",
1635        )
1636        .with_path(DiagnosticPath::from_segments([
1637            DiagnosticPathSegment::Field("plan".into()),
1638            DiagnosticPathSegment::Index(2),
1639            DiagnosticPathSegment::Identifier("person".into()),
1640        ]))
1641        .with_detail("attempt", -7_i64)
1642        .with_detail("expected", vec!["person".to_owned(), "employee".to_owned()])
1643        .with_detail("retryable", false)
1644        .with_detail("subject", "person");
1645        let transport = Transport {
1646            advertisement_contract,
1647            advertisement,
1648            capabilities: Arc::clone(&capability_calls),
1649            exchanges: Arc::clone(&exchanges),
1650            failure: Some(diagnostic),
1651            signer,
1652        };
1653        let options = RemoteConnectionOptions::generated(
1654            RemoteQueryLimits::new(10, 1 << 20, 10, 100, 100, 100),
1655            transport,
1656        );
1657        let remote = RemoteDatabase::connect(options)
1658            .await
1659            .unwrap()
1660            .with_schema(package())
1661            .unwrap();
1662        let mut session = remote.query().unwrap();
1663        let person = session.exact::<Person>().unwrap();
1664        let error = session
1665            .query(person)
1666            .unwrap()
1667            .one()
1668            .await
1669            .expect_err("generated query must return the authenticated application failure");
1670
1671        assert_eq!(error.category(), crate::ErrorCategory::ModelValidation);
1672        assert_eq!(error.code(), Some("remote_application_failure"));
1673        assert_eq!(
1674            error.message(),
1675            "Typed query evidence does not match the validated request invocation"
1676        );
1677        assert_eq!(
1678            error.path(),
1679            Some(&["plan".to_owned(), "[2]".to_owned(), "person".to_owned()][..])
1680        );
1681        assert_eq!(
1682            error.diagnostic_path(),
1683            Some(
1684                &[
1685                    crate::ErrorPathSegment::ContractField("plan".into()),
1686                    crate::ErrorPathSegment::Index(2),
1687                    crate::ErrorPathSegment::ContractIdentity("person".into()),
1688                ][..]
1689            )
1690        );
1691        let details = error.details().expect("authenticated diagnostic details");
1692        assert_eq!(details.get("attempt"), Some(&crate::ErrorDetail::Long(-7)));
1693        assert_eq!(
1694            details.get("expected"),
1695            Some(&crate::ErrorDetail::QueryIdentityList(vec![
1696                "person".to_owned(),
1697                "employee".to_owned(),
1698            ]))
1699        );
1700        assert_eq!(
1701            details.get("retryable"),
1702            Some(&crate::ErrorDetail::Boolean(false))
1703        );
1704        assert_eq!(
1705            details.get("subject"),
1706            Some(&crate::ErrorDetail::QueryIdentity("person".to_owned()))
1707        );
1708        assert_eq!(
1709            details.get("query_category"),
1710            Some(&crate::ErrorDetail::QueryCategory(
1711                crate::QueryDiagnosticCategory::ResultDecode
1712            ))
1713        );
1714        assert_eq!(*capability_calls.lock().unwrap(), 1);
1715        assert_eq!(exchanges.lock().unwrap().len(), 1);
1716    }
1717
1718    #[tokio::test]
1719    async fn sdk_v2_rust_deterministic_proof_fragment() {
1720        let test_id = "remote::tests::sdk_v2_rust_deterministic_proof_fragment";
1721        let results = vec![
1722            serde_json::json!({
1723                "observation_ref": "cancellation_direct",
1724                "proof_kind": "direct_runtime",
1725                "test_id": test_id,
1726                "outcome": "passed",
1727                "observation": observe_sdk_v2_direct_cancellation().await,
1728            }),
1729            serde_json::json!({
1730                "observation_ref": "cancellation_remote",
1731                "proof_kind": "remote_runtime",
1732                "test_id": test_id,
1733                "outcome": "passed",
1734                "observation": observe_sdk_v2_remote_cancellation().await,
1735            }),
1736            serde_json::json!({
1737                "observation_ref": "remote_structured_diagnostic",
1738                "proof_kind": "diagnostic",
1739                "test_id": test_id,
1740                "outcome": "passed",
1741                "observation": observe_sdk_v2_remote_structured_diagnostic().await,
1742            }),
1743        ];
1744        if let Some((destination, run_nonce)) = requested_sdk_v2_proof_fragment() {
1745            publish_sdk_v2_rust_proof_fragment(&destination, &run_nonce, results);
1746        }
1747    }
1748}