Skip to main content

type_bridge_server/typedb/
client.rs

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/// One immutable server connection identity bound to one prepared transport.
11///
12/// The fields are intentionally private: a caller cannot pair trust material,
13/// HTTP discovery settings, or a selected driver band with a different host or
14/// credential set after preparation.
15#[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    /// Connect the V2 ORM authority through this exact prepared identity.
27    #[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
42/// Wrapper around the TypeDB Rust driver providing a clean async API
43/// for query execution and schema retrieval.
44pub struct TypeDBClient {
45    backend: Box<dyn DriverBackend>,
46}
47
48impl TypeDBClient {
49    /// Connect to a TypeDB server using the provided configuration.
50    #[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    /// Connect using the standalone server's validated secure configuration.
59    #[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    /// Validate and snapshot one secure transport policy for reuse.
66    pub fn prepare_secure_transport(
67        config: &SecureTypeDBSection,
68    ) -> Result<PreparedSecureTypeDBConnection, PipelineError> {
69        // Resolve trust before retaining any credential-bearing connection
70        // identity, then bind both into one opaque value.
71        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    /// Connect through an already prepared immutable transport snapshot.
84    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    /// Create a TypeDBClient with a custom backend (for testing).
100    #[cfg(test)]
101    pub(crate) fn with_backend(backend: Box<dyn DriverBackend>) -> Self {
102        Self { backend }
103    }
104
105    /// Execute a TypeQL query and return results as JSON.
106    ///
107    /// For read transactions, the transaction is used directly.
108    /// For write and schema transactions, the transaction is committed after execution.
109    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    /// Return whether a TypeDB database exists on the connected server.
154    pub async fn database_exists(&self, database: &str) -> Result<bool, PipelineError> {
155        self.backend.database_exists(database).await
156    }
157
158    /// Create a TypeDB database.
159    pub async fn create_database(&self, database: &str) -> Result<(), PipelineError> {
160        self.backend.create_database(database).await
161    }
162
163    /// Delete a TypeDB database.
164    pub async fn delete_database(&self, database: &str) -> Result<(), PipelineError> {
165        self.backend.delete_database(database).await
166    }
167
168    /// Delete a TypeDB database if it exists, then create it.
169    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    /// Check if the driver connection is open.
177    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
199/// Parse a transaction type string into a TypeDB TransactionType.
200pub(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    // =============================================
227    // Mock infrastructure
228    // =============================================
229
230    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    // =============================================
463    // parse_transaction_type tests
464    // =============================================
465
466    #[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    // =============================================
506    // execute tests (via MockBackend)
507    // =============================================
508
509    #[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        // Documents + write: commit error is intentionally ignored (let _ = ...)
641        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        // Ok + write: commit error IS propagated
675        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        // Rows + write: commit error IS propagated
689        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        // Backend should never be called if transaction type is invalid
706        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    // =============================================
726    // database admin tests (via MockBackend)
727    // =============================================
728
729    #[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    // =============================================
890    // QueryExecutor trait impl tests
891    // =============================================
892
893    #[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    // =============================================
965    // Integration tests (require running TypeDB)
966    // =============================================
967
968    #[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    /// Live-test target resolved from the environment so the suite can point
985    /// at any disposable TypeDB container instead of a fixed local install.
986    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}