1#![deny(missing_docs)]
2use 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
23pub trait RemoteQueryTransport: Send + Sync + 'static {
28 fn capabilities(&self) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + '_>>;
30
31 fn exchange<'a>(
33 &'a self,
34 request: &'a [u8],
35 ) -> Pin<Box<dyn Future<Output = Result<Vec<u8>>> + Send + 'a>>;
36}
37
38#[derive(Clone, Copy, Debug, Eq, PartialEq)]
40pub struct RemoteQueryLimits {
41 resources: QueryExecutionResourceLimits,
42}
43
44impl RemoteQueryLimits {
45 #[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 #[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
96pub 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 #[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 #[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
148pub 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 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 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 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 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}