1use super::backend::{DriverBackend, QueryResultKind, TransactionType};
2#[cfg(feature = "v2-query")]
3use super::real_driver::sanitize_prepared_connect_error;
4use super::real_driver::{RealTypeDBBackend, prepare_secure_connect_options};
5use crate::config::{SecureTypeDBSection, TypeDBSection};
6use crate::error::PipelineError;
7use crate::executor::QueryExecutor;
8use type_bridge_typedb_runtime::PreparedSecureConnectOptions;
9
10#[doc(hidden)]
16pub struct PreparedSecureTypeDBConnection {
17 address: String,
18 #[cfg(feature = "v2-query")]
19 database: String,
20 username: String,
21 password: String,
22 options: PreparedSecureConnectOptions,
23}
24
25impl PreparedSecureTypeDBConnection {
26 #[cfg(feature = "v2-query")]
28 #[doc(hidden)]
29 pub async fn connect_database(&self) -> Result<type_bridge_orm::Database, PipelineError> {
30 type_bridge_orm::Database::connect_prepared_secure_with_options(
31 &self.address,
32 &self.database,
33 &self.username,
34 &self.password,
35 self.options.clone(),
36 )
37 .await
38 .map_err(sanitize_prepared_connect_error)
39 }
40}
41
42pub struct TypeDBClient {
45 backend: Box<dyn DriverBackend>,
46}
47
48impl TypeDBClient {
49 #[cfg_attr(coverage_nightly, coverage(off))]
51 pub async fn connect(config: &TypeDBSection) -> Result<Self, PipelineError> {
52 let backend = RealTypeDBBackend::connect(config).await?;
53 Ok(Self {
54 backend: Box::new(backend),
55 })
56 }
57
58 #[cfg_attr(coverage_nightly, coverage(off))]
60 pub async fn connect_secure(config: &SecureTypeDBSection) -> Result<Self, PipelineError> {
61 let prepared = Self::prepare_secure_transport(config)?;
62 Self::connect_prepared_secure(&prepared).await
63 }
64
65 pub fn prepare_secure_transport(
67 config: &SecureTypeDBSection,
68 ) -> Result<PreparedSecureTypeDBConnection, PipelineError> {
69 let options = prepare_secure_connect_options(config)?;
72 let connection = &config.connection;
73 Ok(PreparedSecureTypeDBConnection {
74 address: connection.address.clone(),
75 #[cfg(feature = "v2-query")]
76 database: connection.database.clone(),
77 username: connection.username.clone(),
78 password: connection.password.clone(),
79 options,
80 })
81 }
82
83 pub async fn connect_prepared_secure(
85 prepared: &PreparedSecureTypeDBConnection,
86 ) -> Result<Self, PipelineError> {
87 let backend = RealTypeDBBackend::connect_prepared_secure(
88 &prepared.address,
89 &prepared.username,
90 &prepared.password,
91 prepared.options.clone(),
92 )
93 .await?;
94 Ok(Self {
95 backend: Box::new(backend),
96 })
97 }
98
99 #[cfg(test)]
101 pub(crate) fn with_backend(backend: Box<dyn DriverBackend>) -> Self {
102 Self { backend }
103 }
104
105 pub async fn execute(
110 &self,
111 database: &str,
112 typeql: &str,
113 tx_type: &str,
114 ) -> Result<serde_json::Value, PipelineError> {
115 let transaction_type = parse_transaction_type(tx_type)?;
116
117 let mut tx = self
118 .backend
119 .open_transaction(database, transaction_type)
120 .await?;
121
122 let answer = tx.query(typeql).await?;
123
124 let needs_commit = matches!(
125 transaction_type,
126 TransactionType::Write | TransactionType::Schema
127 );
128
129 let results = match answer {
130 QueryResultKind::Ok => {
131 if needs_commit {
132 tx.commit().await?;
133 }
134 serde_json::json!({ "ok": true })
135 }
136 QueryResultKind::Rows(rows) => {
137 if needs_commit {
138 tx.commit().await?;
139 }
140 serde_json::Value::Array(rows)
141 }
142 QueryResultKind::Documents(docs) => {
143 if needs_commit {
144 let _ = tx.commit().await;
145 }
146 serde_json::Value::Array(docs)
147 }
148 };
149
150 Ok(results)
151 }
152
153 pub async fn database_exists(&self, database: &str) -> Result<bool, PipelineError> {
155 self.backend.database_exists(database).await
156 }
157
158 pub async fn create_database(&self, database: &str) -> Result<(), PipelineError> {
160 self.backend.create_database(database).await
161 }
162
163 pub async fn delete_database(&self, database: &str) -> Result<(), PipelineError> {
165 self.backend.delete_database(database).await
166 }
167
168 pub async fn reset_database(&self, database: &str) -> Result<(), PipelineError> {
170 if self.database_exists(database).await? {
171 self.delete_database(database).await?;
172 }
173 self.create_database(database).await
174 }
175
176 pub fn is_connected(&self) -> bool {
178 self.backend.is_open()
179 }
180}
181
182impl QueryExecutor for TypeDBClient {
183 fn execute<'a>(
184 &'a self,
185 database: &'a str,
186 typeql: &'a str,
187 transaction_type: &'a str,
188 ) -> std::pin::Pin<
189 Box<dyn std::future::Future<Output = Result<serde_json::Value, PipelineError>> + Send + 'a>,
190 > {
191 Box::pin(async move { self.execute(database, typeql, transaction_type).await })
192 }
193
194 fn is_connected(&self) -> bool {
195 self.is_connected()
196 }
197}
198
199pub(crate) fn parse_transaction_type(tx_type: &str) -> Result<TransactionType, PipelineError> {
201 match tx_type {
202 "read" => Ok(TransactionType::Read),
203 "write" => Ok(TransactionType::Write),
204 "schema" => Ok(TransactionType::Schema),
205 other => Err(PipelineError::QueryExecution(format!(
206 "Unknown transaction type: {other}"
207 ))),
208 }
209}
210
211#[cfg(test)]
212#[cfg(feature = "band8")]
213#[cfg_attr(coverage_nightly, coverage(off))]
214mod tests {
215 use std::collections::VecDeque;
216 use std::error::Error as _;
217 use std::future::Future;
218 use std::pin::Pin;
219 use std::sync::Arc;
220 use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
221
222 use super::*;
223 use crate::error::PipelineError;
224 use crate::typedb::backend::TransactionOps;
225
226 struct MockTransaction {
231 query_result: Option<QueryResultKind>,
232 query_error: Option<String>,
233 commit_error: Option<String>,
234 committed: Arc<AtomicBool>,
235 query_called: Arc<AtomicBool>,
236 }
237
238 impl MockTransaction {
239 fn new(result: QueryResultKind) -> Self {
240 Self {
241 query_result: Some(result),
242 query_error: None,
243 commit_error: None,
244 committed: Arc::new(AtomicBool::new(false)),
245 query_called: Arc::new(AtomicBool::new(false)),
246 }
247 }
248
249 fn failing_query(msg: &str) -> Self {
250 Self {
251 query_result: None,
252 query_error: Some(msg.to_string()),
253 commit_error: None,
254 committed: Arc::new(AtomicBool::new(false)),
255 query_called: Arc::new(AtomicBool::new(false)),
256 }
257 }
258
259 fn with_commit_error(mut self, msg: &str) -> Self {
260 self.commit_error = Some(msg.to_string());
261 self
262 }
263 }
264
265 impl TransactionOps for MockTransaction {
266 fn query(
267 &mut self,
268 _typeql: &str,
269 ) -> Pin<Box<dyn Future<Output = Result<QueryResultKind, PipelineError>> + Send + '_>>
270 {
271 self.query_called.store(true, Ordering::SeqCst);
272 let result = self.query_result.take();
273 let error = self.query_error.take();
274 Box::pin(async move {
275 if let Some(msg) = error {
276 return Err(PipelineError::QueryExecution(msg));
277 }
278 Ok(result.expect("MockTransaction::query called more than once"))
279 })
280 }
281
282 fn commit(
283 &mut self,
284 ) -> Pin<Box<dyn Future<Output = Result<(), PipelineError>> + Send + '_>> {
285 self.committed.store(true, Ordering::SeqCst);
286 let error = self.commit_error.take();
287 Box::pin(async move {
288 if let Some(msg) = error {
289 return Err(PipelineError::QueryExecution(msg));
290 }
291 Ok(())
292 })
293 }
294 }
295
296 #[derive(Default)]
297 struct MockDatabaseAdmin {
298 exists_results: std::sync::Mutex<VecDeque<Result<bool, String>>>,
299 create_errors: std::sync::Mutex<VecDeque<String>>,
300 delete_errors: std::sync::Mutex<VecDeque<String>>,
301 exists_called: AtomicUsize,
302 create_called: AtomicUsize,
303 delete_called: AtomicUsize,
304 operations: std::sync::Mutex<Vec<String>>,
305 }
306
307 struct MockBackend {
308 transaction: std::sync::Mutex<Option<MockTransaction>>,
309 open_error: Option<String>,
310 is_open: bool,
311 open_called: Arc<AtomicUsize>,
312 database_admin: Arc<MockDatabaseAdmin>,
313 }
314
315 impl MockBackend {
316 fn new(tx: MockTransaction) -> Self {
317 Self {
318 transaction: std::sync::Mutex::new(Some(tx)),
319 open_error: None,
320 is_open: true,
321 open_called: Arc::new(AtomicUsize::new(0)),
322 database_admin: Arc::new(MockDatabaseAdmin::default()),
323 }
324 }
325
326 fn failing(msg: &str) -> Self {
327 Self {
328 transaction: std::sync::Mutex::new(None),
329 open_error: Some(msg.to_string()),
330 is_open: true,
331 open_called: Arc::new(AtomicUsize::new(0)),
332 database_admin: Arc::new(MockDatabaseAdmin::default()),
333 }
334 }
335
336 fn with_database_exists_results(self, results: Vec<Result<bool, String>>) -> Self {
337 *self.database_admin.exists_results.lock().unwrap() = results.into();
338 self
339 }
340
341 fn with_create_errors(self, errors: Vec<&str>) -> Self {
342 *self.database_admin.create_errors.lock().unwrap() =
343 errors.into_iter().map(str::to_string).collect();
344 self
345 }
346
347 fn with_delete_errors(self, errors: Vec<&str>) -> Self {
348 *self.database_admin.delete_errors.lock().unwrap() =
349 errors.into_iter().map(str::to_string).collect();
350 self
351 }
352
353 fn database_admin_state(&self) -> Arc<MockDatabaseAdmin> {
354 Arc::clone(&self.database_admin)
355 }
356 }
357
358 impl DriverBackend for MockBackend {
359 fn open_transaction(
360 &self,
361 _database: &str,
362 _tx_type: TransactionType,
363 ) -> Pin<Box<dyn Future<Output = Result<Box<dyn TransactionOps>, PipelineError>> + Send + '_>>
364 {
365 self.open_called.fetch_add(1, Ordering::SeqCst);
366 let tx = self.transaction.lock().unwrap().take();
367 let error = self.open_error.clone();
368 Box::pin(async move {
369 if let Some(msg) = error {
370 return Err(PipelineError::QueryExecution(msg));
371 }
372 Ok(
373 Box::new(tx.expect("MockBackend: no transaction configured"))
374 as Box<dyn TransactionOps>,
375 )
376 })
377 }
378
379 fn is_open(&self) -> bool {
380 self.is_open
381 }
382
383 fn database_exists(
384 &self,
385 database: &str,
386 ) -> Pin<Box<dyn Future<Output = Result<bool, PipelineError>> + Send + '_>> {
387 self.database_admin
388 .exists_called
389 .fetch_add(1, Ordering::SeqCst);
390 self.database_admin
391 .operations
392 .lock()
393 .unwrap()
394 .push(format!("exists:{database}"));
395 let result = self
396 .database_admin
397 .exists_results
398 .lock()
399 .unwrap()
400 .pop_front()
401 .unwrap_or(Ok(false));
402 Box::pin(async move { result.map_err(PipelineError::Connection) })
403 }
404
405 fn create_database(
406 &self,
407 database: &str,
408 ) -> Pin<Box<dyn Future<Output = Result<(), PipelineError>> + Send + '_>> {
409 self.database_admin
410 .create_called
411 .fetch_add(1, Ordering::SeqCst);
412 self.database_admin
413 .operations
414 .lock()
415 .unwrap()
416 .push(format!("create:{database}"));
417 let error = self
418 .database_admin
419 .create_errors
420 .lock()
421 .unwrap()
422 .pop_front();
423 Box::pin(async move {
424 if let Some(msg) = error {
425 return Err(PipelineError::Connection(msg));
426 }
427 Ok(())
428 })
429 }
430
431 fn delete_database(
432 &self,
433 database: &str,
434 ) -> Pin<Box<dyn Future<Output = Result<(), PipelineError>> + Send + '_>> {
435 self.database_admin
436 .delete_called
437 .fetch_add(1, Ordering::SeqCst);
438 self.database_admin
439 .operations
440 .lock()
441 .unwrap()
442 .push(format!("delete:{database}"));
443 let error = self
444 .database_admin
445 .delete_errors
446 .lock()
447 .unwrap()
448 .pop_front();
449 Box::pin(async move {
450 if let Some(msg) = error {
451 return Err(PipelineError::Connection(msg));
452 }
453 Ok(())
454 })
455 }
456 }
457
458 fn make_client(backend: MockBackend) -> TypeDBClient {
459 TypeDBClient::with_backend(Box::new(backend))
460 }
461
462 #[test]
467 fn parse_transaction_type_read() {
468 let result = parse_transaction_type("read").unwrap();
469 assert_eq!(result, TransactionType::Read);
470 }
471
472 #[test]
473 fn parse_transaction_type_write() {
474 let result = parse_transaction_type("write").unwrap();
475 assert_eq!(result, TransactionType::Write);
476 }
477
478 #[test]
479 fn parse_transaction_type_schema() {
480 let result = parse_transaction_type("schema").unwrap();
481 assert_eq!(result, TransactionType::Schema);
482 }
483
484 #[test]
485 fn parse_transaction_type_unknown() {
486 let result = parse_transaction_type("unknown");
487 let err = result.unwrap_err();
488 assert!(
489 matches!(&err, PipelineError::QueryExecution(msg) if msg.contains("Unknown transaction type: unknown"))
490 );
491 }
492
493 #[test]
494 fn parse_transaction_type_empty() {
495 let result = parse_transaction_type("");
496 assert!(result.is_err());
497 }
498
499 #[test]
500 fn parse_transaction_type_case_sensitive() {
501 let result = parse_transaction_type("Read");
502 assert!(result.is_err());
503 }
504
505 #[tokio::test]
510 async fn execute_ok_read_no_commit() {
511 let tx = MockTransaction::new(QueryResultKind::Ok);
512 let committed = tx.committed.clone();
513 let client = make_client(MockBackend::new(tx));
514
515 let result = client
516 .execute("db", "match $x isa thing;", "read")
517 .await
518 .unwrap();
519 assert_eq!(result, serde_json::json!({"ok": true}));
520 assert!(!committed.load(Ordering::SeqCst));
521 }
522
523 #[tokio::test]
524 async fn execute_ok_write_commits() {
525 let tx = MockTransaction::new(QueryResultKind::Ok);
526 let committed = tx.committed.clone();
527 let client = make_client(MockBackend::new(tx));
528
529 let result = client
530 .execute("db", "insert $x isa thing;", "write")
531 .await
532 .unwrap();
533 assert_eq!(result, serde_json::json!({"ok": true}));
534 assert!(committed.load(Ordering::SeqCst));
535 }
536
537 #[tokio::test]
538 async fn execute_ok_schema_commits() {
539 let tx = MockTransaction::new(QueryResultKind::Ok);
540 let committed = tx.committed.clone();
541 let client = make_client(MockBackend::new(tx));
542
543 let result = client
544 .execute("db", "define entity thing;", "schema")
545 .await
546 .unwrap();
547 assert_eq!(result, serde_json::json!({"ok": true}));
548 assert!(committed.load(Ordering::SeqCst));
549 }
550
551 #[tokio::test]
552 async fn execute_rows_read_no_commit() {
553 let rows = vec![
554 serde_json::json!({"name": "Alice"}),
555 serde_json::json!({"name": "Bob"}),
556 ];
557 let tx = MockTransaction::new(QueryResultKind::Rows(rows.clone()));
558 let committed = tx.committed.clone();
559 let client = make_client(MockBackend::new(tx));
560
561 let result = client
562 .execute("db", "match $p isa person;", "read")
563 .await
564 .unwrap();
565 assert_eq!(result, serde_json::Value::Array(rows));
566 assert!(!committed.load(Ordering::SeqCst));
567 }
568
569 #[tokio::test]
570 async fn execute_rows_write_commits() {
571 let rows = vec![serde_json::json!({"id": 1})];
572 let tx = MockTransaction::new(QueryResultKind::Rows(rows));
573 let committed = tx.committed.clone();
574 let client = make_client(MockBackend::new(tx));
575
576 let result = client
577 .execute("db", "insert $x isa thing;", "write")
578 .await
579 .unwrap();
580 assert!(result.is_array());
581 assert!(committed.load(Ordering::SeqCst));
582 }
583
584 #[tokio::test]
585 async fn execute_rows_data_preserved() {
586 let rows = vec![
587 serde_json::json!({"name": "Alice", "age": 30}),
588 serde_json::json!({"name": "Bob", "age": 25}),
589 ];
590 let tx = MockTransaction::new(QueryResultKind::Rows(rows.clone()));
591 let client = make_client(MockBackend::new(tx));
592
593 let result = client
594 .execute("db", "match $p isa person;", "read")
595 .await
596 .unwrap();
597 let arr = result.as_array().unwrap();
598 assert_eq!(arr.len(), 2);
599 assert_eq!(arr[0]["name"], "Alice");
600 assert_eq!(arr[1]["age"], 25);
601 }
602
603 #[tokio::test]
604 async fn execute_docs_read_no_commit() {
605 let docs = vec![serde_json::json!({"doc": "data"})];
606 let tx = MockTransaction::new(QueryResultKind::Documents(docs.clone()));
607 let committed = tx.committed.clone();
608 let client = make_client(MockBackend::new(tx));
609
610 let result = client
611 .execute("db", "match $p isa person; fetch {};", "read")
612 .await
613 .unwrap();
614 assert_eq!(result, serde_json::Value::Array(docs));
615 assert!(!committed.load(Ordering::SeqCst));
616 }
617
618 #[tokio::test]
619 async fn execute_docs_write_commits() {
620 let docs = vec![serde_json::json!({"doc": "data"})];
621 let tx = MockTransaction::new(QueryResultKind::Documents(docs));
622 let committed = tx.committed.clone();
623 let client = make_client(MockBackend::new(tx));
624
625 let result = client
626 .execute("db", "insert $x isa thing;", "write")
627 .await
628 .unwrap();
629 assert!(result.is_array());
630 assert!(committed.load(Ordering::SeqCst));
631 }
632
633 #[tokio::test]
634 async fn execute_docs_commit_error_ignored() {
635 let docs = vec![serde_json::json!({"doc": "data"})];
636 let tx = MockTransaction::new(QueryResultKind::Documents(docs.clone()))
637 .with_commit_error("commit failed");
638 let client = make_client(MockBackend::new(tx));
639
640 let result = client
642 .execute("db", "insert $x isa thing;", "write")
643 .await
644 .unwrap();
645 assert_eq!(result, serde_json::Value::Array(docs));
646 }
647
648 #[tokio::test]
649 async fn execute_transaction_open_failure() {
650 let client = make_client(MockBackend::failing("connection refused"));
651
652 let result = client.execute("db", "match $x isa thing;", "read").await;
653 let err = result.unwrap_err();
654 assert!(
655 matches!(&err, PipelineError::QueryExecution(msg) if msg.contains("connection refused"))
656 );
657 }
658
659 #[tokio::test]
660 async fn execute_query_failure() {
661 let tx = MockTransaction::failing_query("syntax error");
662 let client = make_client(MockBackend::new(tx));
663
664 let result = client.execute("db", "bad query", "read").await;
665 let err = result.unwrap_err();
666 assert!(matches!(&err, PipelineError::QueryExecution(msg) if msg.contains("syntax error")));
667 }
668
669 #[tokio::test]
670 async fn execute_commit_failure_ok_propagated() {
671 let tx = MockTransaction::new(QueryResultKind::Ok).with_commit_error("commit failed");
672 let client = make_client(MockBackend::new(tx));
673
674 let result = client.execute("db", "insert $x isa thing;", "write").await;
676 let err = result.unwrap_err();
677 assert!(
678 matches!(&err, PipelineError::QueryExecution(msg) if msg.contains("commit failed"))
679 );
680 }
681
682 #[tokio::test]
683 async fn execute_commit_failure_rows_propagated() {
684 let tx =
685 MockTransaction::new(QueryResultKind::Rows(vec![])).with_commit_error("commit failed");
686 let client = make_client(MockBackend::new(tx));
687
688 let result = client.execute("db", "insert $x isa thing;", "write").await;
690 let err = result.unwrap_err();
691 assert!(
692 matches!(&err, PipelineError::QueryExecution(msg) if msg.contains("commit failed"))
693 );
694 }
695
696 #[tokio::test]
697 async fn execute_invalid_transaction_type() {
698 let tx = MockTransaction::new(QueryResultKind::Ok);
699 let backend = MockBackend::new(tx);
700 let open_called = backend.open_called.clone();
701 let client = make_client(backend);
702
703 let result = client.execute("db", "match $x;", "invalid").await;
704 assert!(result.is_err());
705 assert_eq!(open_called.load(Ordering::SeqCst), 0);
707 }
708
709 #[test]
710 fn is_connected_delegates_to_backend() {
711 let mut backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok));
712 backend.is_open = true;
713 let client = make_client(backend);
714 assert!(client.is_connected());
715 }
716
717 #[test]
718 fn is_connected_false_when_backend_closed() {
719 let mut backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok));
720 backend.is_open = false;
721 let client = make_client(backend);
722 assert!(!client.is_connected());
723 }
724
725 #[tokio::test]
730 async fn database_exists_true_delegates_to_backend() {
731 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok))
732 .with_database_exists_results(vec![Ok(true)]);
733 let admin = backend.database_admin_state();
734 let client = make_client(backend);
735
736 assert!(client.database_exists("admin_db").await.unwrap());
737 assert_eq!(admin.exists_called.load(Ordering::SeqCst), 1);
738 assert_eq!(
739 admin.operations.lock().unwrap().clone(),
740 vec!["exists:admin_db".to_string()]
741 );
742 }
743
744 #[tokio::test]
745 async fn database_exists_false_delegates_to_backend() {
746 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok))
747 .with_database_exists_results(vec![Ok(false)]);
748 let admin = backend.database_admin_state();
749 let client = make_client(backend);
750
751 assert!(!client.database_exists("missing_db").await.unwrap());
752 assert_eq!(admin.exists_called.load(Ordering::SeqCst), 1);
753 assert_eq!(
754 admin.operations.lock().unwrap().clone(),
755 vec!["exists:missing_db".to_string()]
756 );
757 }
758
759 #[tokio::test]
760 async fn create_database_delegates_to_backend() {
761 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok));
762 let admin = backend.database_admin_state();
763 let client = make_client(backend);
764
765 client.create_database("new_db").await.unwrap();
766 assert_eq!(admin.create_called.load(Ordering::SeqCst), 1);
767 assert_eq!(
768 admin.operations.lock().unwrap().clone(),
769 vec!["create:new_db".to_string()]
770 );
771 }
772
773 #[tokio::test]
774 async fn delete_database_delegates_to_backend() {
775 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok));
776 let admin = backend.database_admin_state();
777 let client = make_client(backend);
778
779 client.delete_database("old_db").await.unwrap();
780 assert_eq!(admin.delete_called.load(Ordering::SeqCst), 1);
781 assert_eq!(
782 admin.operations.lock().unwrap().clone(),
783 vec!["delete:old_db".to_string()]
784 );
785 }
786
787 #[tokio::test]
788 async fn reset_database_deletes_existing_database_before_create() {
789 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok))
790 .with_database_exists_results(vec![Ok(true)]);
791 let admin = backend.database_admin_state();
792 let client = make_client(backend);
793
794 client.reset_database("reset_db").await.unwrap();
795 assert_eq!(admin.exists_called.load(Ordering::SeqCst), 1);
796 assert_eq!(admin.delete_called.load(Ordering::SeqCst), 1);
797 assert_eq!(admin.create_called.load(Ordering::SeqCst), 1);
798 assert_eq!(
799 admin.operations.lock().unwrap().clone(),
800 vec![
801 "exists:reset_db".to_string(),
802 "delete:reset_db".to_string(),
803 "create:reset_db".to_string(),
804 ]
805 );
806 }
807
808 #[tokio::test]
809 async fn reset_database_creates_when_database_is_absent() {
810 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok))
811 .with_database_exists_results(vec![Ok(false)]);
812 let admin = backend.database_admin_state();
813 let client = make_client(backend);
814
815 client.reset_database("reset_db").await.unwrap();
816 assert_eq!(admin.exists_called.load(Ordering::SeqCst), 1);
817 assert_eq!(admin.delete_called.load(Ordering::SeqCst), 0);
818 assert_eq!(admin.create_called.load(Ordering::SeqCst), 1);
819 assert_eq!(
820 admin.operations.lock().unwrap().clone(),
821 vec!["exists:reset_db".to_string(), "create:reset_db".to_string(),]
822 );
823 }
824
825 #[tokio::test]
826 async fn database_exists_propagates_backend_error() {
827 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok))
828 .with_database_exists_results(vec![Err("lookup failed".to_string())]);
829 let client = make_client(backend);
830
831 let err = client.database_exists("db").await.unwrap_err();
832 assert!(matches!(&err, PipelineError::Connection(msg) if msg.contains("lookup failed")));
833 }
834
835 #[tokio::test]
836 async fn create_database_propagates_backend_error() {
837 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok))
838 .with_create_errors(vec!["create failed"]);
839 let client = make_client(backend);
840
841 let err = client.create_database("db").await.unwrap_err();
842 assert!(matches!(&err, PipelineError::Connection(msg) if msg.contains("create failed")));
843 }
844
845 #[tokio::test]
846 async fn delete_database_propagates_backend_error() {
847 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok))
848 .with_delete_errors(vec!["delete failed"]);
849 let client = make_client(backend);
850
851 let err = client.delete_database("db").await.unwrap_err();
852 assert!(matches!(&err, PipelineError::Connection(msg) if msg.contains("delete failed")));
853 }
854
855 #[tokio::test]
856 async fn reset_database_propagates_lookup_error_without_mutating() {
857 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok))
858 .with_database_exists_results(vec![Err("lookup failed".to_string())]);
859 let admin = backend.database_admin_state();
860 let client = make_client(backend);
861
862 let err = client.reset_database("db").await.unwrap_err();
863 assert!(matches!(&err, PipelineError::Connection(msg) if msg.contains("lookup failed")));
864 assert_eq!(admin.delete_called.load(Ordering::SeqCst), 0);
865 assert_eq!(admin.create_called.load(Ordering::SeqCst), 0);
866 assert_eq!(
867 admin.operations.lock().unwrap().clone(),
868 vec!["exists:db".to_string()]
869 );
870 }
871
872 #[tokio::test]
873 async fn reset_database_propagates_delete_error_without_create() {
874 let backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok))
875 .with_database_exists_results(vec![Ok(true)])
876 .with_delete_errors(vec!["delete failed"]);
877 let admin = backend.database_admin_state();
878 let client = make_client(backend);
879
880 let err = client.reset_database("db").await.unwrap_err();
881 assert!(matches!(&err, PipelineError::Connection(msg) if msg.contains("delete failed")));
882 assert_eq!(admin.create_called.load(Ordering::SeqCst), 0);
883 assert_eq!(
884 admin.operations.lock().unwrap().clone(),
885 vec!["exists:db".to_string(), "delete:db".to_string()]
886 );
887 }
888
889 #[test]
894 fn type_db_client_implements_query_executor() {
895 fn assert_executor<T: QueryExecutor>() {}
896 assert_executor::<TypeDBClient>();
897 }
898
899 #[tokio::test]
900 async fn query_executor_execute_delegates_to_client() {
901 let tx = MockTransaction::new(QueryResultKind::Rows(vec![serde_json::json!({"x": 1})]));
902 let client = make_client(MockBackend::new(tx));
903 let executor: Box<dyn QueryExecutor> = Box::new(client);
904
905 let result = executor
906 .execute("db", "match $x isa thing;", "read")
907 .await
908 .unwrap();
909 assert!(result.is_array());
910 assert_eq!(result.as_array().unwrap().len(), 1);
911 }
912
913 #[test]
914 fn query_executor_is_connected_delegates_to_client() {
915 let mut backend = MockBackend::new(MockTransaction::new(QueryResultKind::Ok));
916 backend.is_open = true;
917 let client = make_client(backend);
918 let executor: Box<dyn QueryExecutor> = Box::new(client);
919 assert!(executor.is_connected());
920 }
921
922 #[tokio::test]
923 async fn prepared_first_and_v2_connections_never_expose_credentials_or_provider_identity() {
924 const ADDRESS: &str = "admin:TB_ADDRESS_SECRET@provider.invalid:1729";
925 const USERNAME: &str = "TB_USERNAME_SECRET";
926 const PASSWORD: &str = "TB_PASSWORD_SECRET";
927 let config = SecureTypeDBSection::new(
928 TypeDBSection {
929 address: ADDRESS.to_owned(),
930 database: "db".to_owned(),
931 username: USERNAME.to_owned(),
932 password: PASSWORD.to_owned(),
933 http_port: 8000,
934 server_version: Some("3.12.1".to_owned()),
935 },
936 crate::config::OutboundTlsMode::Disabled,
937 );
938 let prepared = TypeDBClient::prepare_secure_transport(&config).unwrap();
939
940 let first = match TypeDBClient::connect_prepared_secure(&prepared).await {
941 Ok(_) => panic!("malformed provider address unexpectedly connected"),
942 Err(error) => error,
943 };
944 let rendered = format!("{first}\n{first:?}");
945 for secret in [ADDRESS, USERNAME, PASSWORD] {
946 assert!(!rendered.contains(secret), "{secret}: {rendered}");
947 }
948 assert!(first.source().is_none());
949
950 #[cfg(feature = "v2-query")]
951 {
952 let second = match prepared.connect_database().await {
953 Ok(_) => panic!("malformed V2 provider address unexpectedly connected"),
954 Err(error) => error,
955 };
956 let rendered = format!("{second}\n{second:?}");
957 for secret in [ADDRESS, USERNAME, PASSWORD] {
958 assert!(!rendered.contains(secret), "{secret}: {rendered}");
959 }
960 assert!(second.source().is_none());
961 }
962 }
963
964 #[tokio::test]
969 #[ignore = "requires running TypeDB server"]
970 #[cfg_attr(coverage_nightly, coverage(off))]
971 async fn integration_connect_invalid_address() {
972 let config = TypeDBSection {
973 address: "localhost:99999".to_string(),
974 database: "test".to_string(),
975 username: "admin".to_string(),
976 password: "password".to_string(),
977 http_port: 8000,
978 server_version: None,
979 };
980 let result = TypeDBClient::connect(&config).await;
981 assert!(result.is_err());
982 }
983
984 fn live_config() -> TypeDBSection {
987 TypeDBSection {
988 address: std::env::var("TYPEDB_ADDRESS")
989 .unwrap_or_else(|_| "localhost:1729".to_string()),
990 database: std::env::var("TYPEDB_DATABASE").unwrap_or_else(|_| "test".to_string()),
991 username: "admin".to_string(),
992 password: "password".to_string(),
993 http_port: std::env::var("TYPEDB_HTTP_PORT")
994 .ok()
995 .and_then(|port| port.parse().ok())
996 .unwrap_or(8000),
997 server_version: std::env::var("TYPEDB_SERVER_VERSION").ok(),
998 }
999 }
1000
1001 #[tokio::test]
1002 #[ignore = "requires running TypeDB server"]
1003 #[cfg_attr(coverage_nightly, coverage(off))]
1004 async fn integration_connect_success() {
1005 let result = TypeDBClient::connect(&live_config()).await;
1006 assert!(result.is_ok());
1007 assert!(result.unwrap().is_connected());
1008 }
1009
1010 #[tokio::test]
1011 #[ignore = "requires running TypeDB server"]
1012 #[cfg_attr(coverage_nightly, coverage(off))]
1013 async fn integration_database_admin_roundtrip() {
1014 let config = live_config();
1015 let client = TypeDBClient::connect(&config)
1016 .await
1017 .expect("connect failed");
1018 let database = format!("type_bridge_server_admin_{}", uuid::Uuid::new_v4().simple());
1019
1020 if client.database_exists(&database).await.unwrap_or(false) {
1021 let _ = client.delete_database(&database).await;
1022 }
1023
1024 assert!(!client.database_exists(&database).await.unwrap());
1025 client
1026 .create_database(&database)
1027 .await
1028 .expect("create failed");
1029 assert!(client.database_exists(&database).await.unwrap());
1030
1031 client
1032 .reset_database(&database)
1033 .await
1034 .expect("reset existing database failed");
1035 assert!(client.database_exists(&database).await.unwrap());
1036
1037 client
1038 .delete_database(&database)
1039 .await
1040 .expect("delete failed");
1041 assert!(!client.database_exists(&database).await.unwrap());
1042
1043 client
1044 .reset_database(&database)
1045 .await
1046 .expect("reset absent database failed");
1047 assert!(client.database_exists(&database).await.unwrap());
1048
1049 client
1050 .delete_database(&database)
1051 .await
1052 .expect("cleanup delete failed");
1053 }
1054
1055 #[tokio::test]
1056 #[ignore = "requires running TypeDB server"]
1057 #[cfg_attr(coverage_nightly, coverage(off))]
1058 async fn integration_execute_roundtrip() {
1059 let config = live_config();
1060 let client = TypeDBClient::connect(&config)
1061 .await
1062 .expect("connect failed");
1063
1064 client
1065 .execute(&config.database, "define entity smoke_marker;", "schema")
1066 .await
1067 .expect("schema define failed");
1068
1069 let rows = client
1070 .execute(&config.database, "match entity $t;", "read")
1071 .await
1072 .expect("read query failed");
1073 let rows = rows.as_array().expect("read result must be a JSON array");
1074 assert!(
1075 !rows.is_empty(),
1076 "expected at least the smoke_marker entity type"
1077 );
1078 }
1079}