Skip to main content

type_bridge_server/
pipeline.rs

1use std::collections::HashMap;
2use std::time::Instant;
3
4use type_bridge_core_lib::_schema::TypeSchema;
5use type_bridge_core_lib::ast::Clause;
6use type_bridge_core_lib::compiler::QueryCompiler;
7use type_bridge_core_lib::validation::ValidationEngine;
8
9use crate::error::PipelineError;
10use crate::executor::QueryExecutor;
11use crate::interceptor::crud_interceptor::{CrudInterceptor, CrudInterceptorAdapter};
12use crate::interceptor::{Interceptor, InterceptorChain, RequestContext};
13#[cfg(feature = "v2-query")]
14use crate::interceptor::{V2PolicyOutcome, V2PolicyRequest};
15use crate::schema_source::SchemaSource;
16
17/// Input for a structured (AST-based) query.
18pub struct QueryInput {
19    /// Database override, or `None` to use the pipeline default.
20    pub database: Option<String>,
21    /// Provider transaction mode requested by the caller.
22    pub transaction_type: String,
23    /// Structured query clauses to validate and execute.
24    pub clauses: Vec<Clause>,
25    /// Transport and application metadata exposed to interceptors.
26    pub metadata: HashMap<String, serde_json::Value>,
27}
28
29/// Input for a validation-only request.
30pub struct ValidateInput {
31    /// Structured query clauses to validate without execution.
32    pub clauses: Vec<Clause>,
33}
34
35/// Output from a successful pipeline execution.
36#[derive(Debug)]
37pub struct QueryOutput {
38    /// Provider result encoded as JSON.
39    pub results: serde_json::Value,
40    /// Unique identifier assigned to the request.
41    pub request_id: String,
42    /// Whole-pipeline execution duration in milliseconds.
43    pub execution_time_ms: u64,
44    /// Interceptor names applied to the request, in request order.
45    pub interceptors_applied: Vec<String>,
46}
47
48/// Output from a validation-only request.
49#[derive(Debug)]
50pub struct ValidateOutput {
51    /// Whether every supplied clause passed static validation.
52    pub is_valid: bool,
53    /// Deterministically ordered validation findings.
54    pub errors: Vec<ValidationErrorDetail>,
55}
56
57/// A single validation error.
58#[derive(Debug)]
59pub struct ValidationErrorDetail {
60    /// Stable machine-readable validation code.
61    pub code: String,
62    /// Human-readable validation explanation.
63    pub message: String,
64    /// Logical location of the invalid query element.
65    pub path: String,
66}
67
68#[cfg_attr(coverage_nightly, coverage(off))]
69fn log_query_execution(database: &str, transaction_type: &str, typeql: &str) {
70    tracing::info!(database, transaction_type, "Executing query");
71    tracing::debug!(typeql, "Compiled TypeQL");
72}
73
74/// Transport-agnostic query pipeline.
75///
76/// Encapsulates the full query lifecycle: validate → intercept → compile → execute → intercept.
77/// Use [`PipelineBuilder`] to construct an instance.
78///
79/// # Example
80///
81/// This example is ignored because it relies on application-defined executor
82/// and schema-source implementations.
83///
84/// ```rust,ignore
85/// use type_bridge_server::{PipelineBuilder, QueryInput};
86///
87/// let pipeline = PipelineBuilder::new(my_executor)
88///     .with_schema_source(my_schema_source)
89///     .with_default_database("my_db")
90///     .build()?;
91///
92/// let output = pipeline.execute_query(QueryInput { ... }).await?;
93/// ```
94pub struct QueryPipeline {
95    schema: Option<TypeSchema>,
96    validation_engine: ValidationEngine,
97    interceptor_chain: InterceptorChain,
98    default_database: String,
99    executor: Box<dyn QueryExecutor>,
100    skip_validation: bool,
101}
102
103impl QueryPipeline {
104    /// Prove every configured policy explicitly covers typed V2 requests.
105    ///
106    /// Server startup calls this before constructing the V2 router/listener;
107    /// request admission rechecks it defensively through
108    /// [`Self::v2_request_context`].
109    #[cfg(feature = "v2-query")]
110    pub fn validate_v2_coverage(&self) -> Result<(), PipelineError> {
111        self.interceptor_chain
112            .ensure_v2_coverage()
113            .map_err(|error| PipelineError::Interceptor(error.to_string()))
114    }
115
116    /// Construct one V2 policy context after proving every configured
117    /// interceptor explicitly supports typed plans.
118    #[cfg(feature = "v2-query")]
119    pub fn v2_request_context(
120        &self,
121        metadata: HashMap<String, serde_json::Value>,
122    ) -> Result<RequestContext, PipelineError> {
123        self.validate_v2_coverage()?;
124        Ok(RequestContext {
125            request_id: uuid::Uuid::new_v4().to_string(),
126            client_id: "unknown".to_owned(),
127            database: self.default_database.clone(),
128            transaction_type: "read".to_owned(),
129            metadata,
130            timestamp: chrono::Utc::now(),
131            crud_info: None,
132        })
133    }
134
135    /// Authenticate and rate-limit V2 transport metadata before envelope
136    /// decoding or other attacker-controlled contract work.
137    #[cfg(feature = "v2-query")]
138    pub async fn begin_v2_request(
139        &self,
140        context: &mut RequestContext,
141    ) -> Result<(), PipelineError> {
142        self.interceptor_chain
143            .execute_v2_transport(context)
144            .await
145            .map_err(|error| PipelineError::Interceptor(error.to_string()))
146    }
147
148    /// Authorize one exact validated V2 plan before replay/provider admission.
149    #[cfg(feature = "v2-query")]
150    pub async fn authorize_v2_request(
151        &self,
152        request: &V2PolicyRequest<'_>,
153        context: &mut RequestContext,
154    ) -> Result<(), PipelineError> {
155        self.interceptor_chain
156            .execute_v2_request(request, context)
157            .await
158            .map_err(|error| PipelineError::Interceptor(error.to_string()))
159    }
160
161    /// Run configured response and audit policies for every V2 outcome.
162    #[cfg(feature = "v2-query")]
163    pub async fn finish_v2_request(
164        &self,
165        outcome: &V2PolicyOutcome<'_>,
166        context: &RequestContext,
167    ) -> Result<(), PipelineError> {
168        self.interceptor_chain
169            .execute_v2_response(outcome, context)
170            .await
171            .map_err(|error| PipelineError::Interceptor(error.to_string()))
172    }
173
174    /// Execute a structured (AST-based) query through the full pipeline.
175    pub async fn execute_query(&self, input: QueryInput) -> Result<QueryOutput, PipelineError> {
176        let start = Instant::now();
177        let request_id = uuid::Uuid::new_v4().to_string();
178        let database = input
179            .database
180            .unwrap_or_else(|| self.default_database.clone());
181
182        let mut ctx = RequestContext {
183            request_id: request_id.clone(),
184            client_id: "unknown".to_string(),
185            database: database.clone(),
186            transaction_type: input.transaction_type.clone(),
187            metadata: input.metadata,
188            timestamp: chrono::Utc::now(),
189            crud_info: None,
190        };
191
192        // Validate against schema
193        if !self.skip_validation
194            && let Some(schema) = &self.schema
195        {
196            let result = self
197                .validation_engine
198                .validate_query(&input.clauses, schema);
199            if !result.is_valid {
200                return Err(PipelineError::Validation(format!(
201                    "{} validation error(s)",
202                    result.errors.len()
203                )));
204            }
205        }
206
207        // Run request interceptors
208        let clauses = self
209            .interceptor_chain
210            .execute_request(input.clauses, &mut ctx)
211            .await
212            .map_err(|e| PipelineError::Interceptor(e.to_string()))?;
213
214        // Compile to TypeQL
215        let compiler = QueryCompiler::new();
216        let typeql = compiler.compile(&clauses);
217        ctx.metadata.insert(
218            "compiled_typeql".to_string(),
219            serde_json::Value::String(typeql.clone()),
220        );
221
222        // Execute
223        log_query_execution(&database, &input.transaction_type, &typeql);
224
225        let results = self
226            .executor
227            .execute(&database, &typeql, &input.transaction_type)
228            .await?;
229
230        // Run response interceptors
231        self.interceptor_chain
232            .execute_response(&results, &ctx)
233            .await
234            .map_err(|e| PipelineError::Interceptor(e.to_string()))?;
235
236        let elapsed = start.elapsed().as_millis() as u64;
237
238        Ok(QueryOutput {
239            results,
240            request_id,
241            execution_time_ms: elapsed,
242            interceptors_applied: self
243                .interceptor_chain
244                .interceptor_names()
245                .into_iter()
246                .map(String::from)
247                .collect(),
248        })
249    }
250
251    /// Validate clauses against the loaded schema without executing.
252    pub fn validate(&self, input: &ValidateInput) -> Result<ValidateOutput, PipelineError> {
253        let schema = self
254            .schema
255            .as_ref()
256            .ok_or_else(|| PipelineError::Schema("No schema loaded".to_string()))?;
257
258        let result = self
259            .validation_engine
260            .validate_query(&input.clauses, schema);
261
262        let errors = result
263            .errors
264            .iter()
265            .map(|e| ValidationErrorDetail {
266                code: e.code.clone(),
267                message: e.message.clone(),
268                path: e.path.clone(),
269            })
270            .collect();
271
272        Ok(ValidateOutput {
273            is_valid: result.is_valid,
274            errors,
275        })
276    }
277
278    /// Get the loaded schema, if any.
279    pub fn schema(&self) -> Option<&TypeSchema> {
280        self.schema.as_ref()
281    }
282
283    /// Check if the backend executor is connected.
284    pub fn is_connected(&self) -> bool {
285        self.executor.is_connected()
286    }
287
288    /// Get the default database name.
289    pub fn default_database(&self) -> &str {
290        &self.default_database
291    }
292}
293
294/// Builder for constructing a [`QueryPipeline`].
295///
296/// # Example
297///
298/// This example is ignored because it relies on application-defined executor,
299/// schema-source, and interceptor implementations.
300///
301/// ```rust,ignore
302/// use type_bridge_server::PipelineBuilder;
303///
304/// let pipeline = PipelineBuilder::new(my_executor)
305///     .with_schema_source(FileSchemaSource::new("schema.tql"))
306///     .with_interceptor(AuditLogInterceptor::new(&config)?)
307///     .with_default_database("my_db")
308///     .build()?;
309/// ```
310pub struct PipelineBuilder {
311    executor: Box<dyn QueryExecutor>,
312    schema_source: Option<Box<dyn SchemaSource>>,
313    interceptors: Vec<Box<dyn Interceptor>>,
314    default_database: String,
315    skip_validation: bool,
316}
317
318impl PipelineBuilder {
319    /// Create a new builder with the given query executor.
320    pub fn new(executor: impl QueryExecutor + 'static) -> Self {
321        Self {
322            executor: Box::new(executor),
323            schema_source: None,
324            interceptors: Vec::new(),
325            default_database: String::new(),
326            skip_validation: false,
327        }
328    }
329
330    /// Set the schema source. The schema will be loaded during [`build()`](Self::build).
331    pub fn with_schema_source(mut self, source: impl SchemaSource + 'static) -> Self {
332        self.schema_source = Some(Box::new(source));
333        self
334    }
335
336    /// Add an interceptor to the pipeline chain.
337    pub fn with_interceptor(mut self, interceptor: impl Interceptor + 'static) -> Self {
338        self.interceptors.push(Box::new(interceptor));
339        self
340    }
341
342    /// Set the default database name used when requests don't specify one.
343    pub fn with_default_database(mut self, database: impl Into<String>) -> Self {
344        self.default_database = database.into();
345        self
346    }
347
348    /// Add a CRUD-aware interceptor to the pipeline chain.
349    ///
350    /// The interceptor is automatically wrapped in a [`CrudInterceptorAdapter`]
351    /// that extracts [`CrudInfo`](crate::interceptor::CrudInfo) and delegates
352    /// to the CRUD-specific hooks.
353    pub fn with_crud_interceptor(self, interceptor: impl CrudInterceptor + 'static) -> Self {
354        self.with_interceptor(CrudInterceptorAdapter::new(interceptor))
355    }
356
357    /// Skip schema validation during query execution.
358    ///
359    /// The schema is still loaded (and accessible via [`QueryPipeline::schema`]),
360    /// but queries are not validated against it before execution.
361    pub fn with_skip_validation(mut self) -> Self {
362        self.skip_validation = true;
363        self
364    }
365
366    /// Build the pipeline, loading the schema if a source was provided.
367    pub fn build(self) -> Result<QueryPipeline, PipelineError> {
368        let schema = match self.schema_source {
369            Some(source) => Some(source.load()?),
370            None => None,
371        };
372
373        Ok(QueryPipeline {
374            schema,
375            validation_engine: ValidationEngine::new(),
376            interceptor_chain: InterceptorChain::new(self.interceptors),
377            default_database: self.default_database,
378            executor: self.executor,
379            skip_validation: self.skip_validation,
380        })
381    }
382}
383
384#[cfg(test)]
385#[cfg_attr(coverage_nightly, coverage(off))]
386mod tests {
387    use std::future::Future;
388    use std::pin::Pin;
389    use std::sync::Arc;
390    use std::sync::atomic::{AtomicUsize, Ordering};
391
392    use type_bridge_core_lib::ast::{Constraint, Pattern, Value};
393
394    use super::*;
395    use crate::interceptor::traits::InterceptError;
396    use crate::test_helpers::{MockExecutor, make_pipeline, make_simple_clauses};
397
398    fn init_tracing() -> tracing::subscriber::DefaultGuard {
399        let subscriber = tracing_subscriber::fmt()
400            .with_max_level(tracing::Level::DEBUG)
401            .with_test_writer()
402            .finish();
403        tracing::subscriber::set_default(subscriber)
404    }
405
406    // --- Helper interceptors ---
407
408    struct PassthroughInterceptor {
409        name: String,
410    }
411
412    impl Interceptor for PassthroughInterceptor {
413        fn name(&self) -> &str {
414            &self.name
415        }
416        fn on_request<'a>(
417            &'a self,
418            clauses: Vec<Clause>,
419            _ctx: &'a mut RequestContext,
420        ) -> Pin<Box<dyn Future<Output = Result<Vec<Clause>, InterceptError>> + Send + 'a>>
421        {
422            Box::pin(async move { Ok(clauses) })
423        }
424    }
425
426    struct RejectingRequestInterceptor;
427
428    impl Interceptor for RejectingRequestInterceptor {
429        fn name(&self) -> &str {
430            "rejector"
431        }
432        fn on_request<'a>(
433            &'a self,
434            _clauses: Vec<Clause>,
435            _ctx: &'a mut RequestContext,
436        ) -> Pin<Box<dyn Future<Output = Result<Vec<Clause>, InterceptError>> + Send + 'a>>
437        {
438            Box::pin(async {
439                Err(InterceptError::AccessDenied {
440                    reason: "test rejection".into(),
441                })
442            })
443        }
444
445        #[cfg(feature = "v2-query")]
446        fn supports_v2(&self) -> bool {
447            true
448        }
449
450        #[cfg(feature = "v2-query")]
451        fn on_v2_transport<'a>(
452            &'a self,
453            _ctx: &'a mut RequestContext,
454        ) -> Pin<Box<dyn Future<Output = Result<(), InterceptError>> + Send + 'a>> {
455            Box::pin(async {
456                Err(InterceptError::AccessDenied {
457                    reason: "test rejection".into(),
458                })
459            })
460        }
461    }
462
463    struct RejectingResponseInterceptor;
464
465    impl Interceptor for RejectingResponseInterceptor {
466        fn name(&self) -> &str {
467            "resp-rejector"
468        }
469        fn on_request<'a>(
470            &'a self,
471            clauses: Vec<Clause>,
472            _ctx: &'a mut RequestContext,
473        ) -> Pin<Box<dyn Future<Output = Result<Vec<Clause>, InterceptError>> + Send + 'a>>
474        {
475            Box::pin(async move { Ok(clauses) })
476        }
477        fn on_response<'a>(
478            &'a self,
479            _result: &'a serde_json::Value,
480            _ctx: &'a RequestContext,
481        ) -> Pin<Box<dyn Future<Output = Result<(), InterceptError>> + Send + 'a>> {
482            Box::pin(async { Err(InterceptError::Internal("response rejected".into())) })
483        }
484    }
485
486    struct CountingInterceptor {
487        name: String,
488        count: Arc<AtomicUsize>,
489    }
490
491    impl Interceptor for CountingInterceptor {
492        fn name(&self) -> &str {
493            &self.name
494        }
495        fn on_request<'a>(
496            &'a self,
497            clauses: Vec<Clause>,
498            _ctx: &'a mut RequestContext,
499        ) -> Pin<Box<dyn Future<Output = Result<Vec<Clause>, InterceptError>> + Send + 'a>>
500        {
501            Box::pin(async move {
502                self.count.fetch_add(1, Ordering::SeqCst);
503                Ok(clauses)
504            })
505        }
506    }
507
508    /// SchemaSource that always fails.
509    struct FailingSchemaSource;
510
511    impl crate::schema_source::SchemaSource for FailingSchemaSource {
512        fn load(&self) -> Result<TypeSchema, PipelineError> {
513            Err(PipelineError::Schema("source failed".into()))
514        }
515    }
516
517    fn make_query_input(clauses: Vec<Clause>) -> QueryInput {
518        QueryInput {
519            database: None,
520            transaction_type: "read".to_string(),
521            clauses,
522            metadata: HashMap::new(),
523        }
524    }
525
526    fn make_query_input_with_db(clauses: Vec<Clause>, db: &str) -> QueryInput {
527        QueryInput {
528            database: Some(db.to_string()),
529            transaction_type: "read".to_string(),
530            clauses,
531            metadata: HashMap::new(),
532        }
533    }
534
535    // =============================================
536    // PipelineBuilder tests
537    // =============================================
538
539    #[test]
540    fn builder_without_schema_source() {
541        let pipeline = PipelineBuilder::new(MockExecutor::new()).build().unwrap();
542        assert!(pipeline.schema().is_none());
543    }
544
545    #[test]
546    fn builder_with_valid_schema_source() {
547        let pipeline = make_pipeline(MockExecutor::new(), true);
548        assert!(pipeline.schema().is_some());
549        let schema = pipeline.schema().unwrap();
550        assert!(schema.entities.contains_key("person"));
551    }
552
553    #[test]
554    fn builder_with_failing_schema_source() {
555        let result = PipelineBuilder::new(MockExecutor::new())
556            .with_schema_source(FailingSchemaSource)
557            .build();
558        let err = result.err().expect("Expected build error");
559        assert!(matches!(&err, PipelineError::Schema(msg) if msg.contains("source failed")));
560    }
561
562    #[test]
563    fn builder_with_default_database() {
564        let pipeline = PipelineBuilder::new(MockExecutor::new())
565            .with_default_database("mydb")
566            .build()
567            .unwrap();
568        assert_eq!(pipeline.default_database(), "mydb");
569    }
570
571    #[test]
572    fn builder_default_empty_database() {
573        let pipeline = PipelineBuilder::new(MockExecutor::new()).build().unwrap();
574        assert_eq!(pipeline.default_database(), "");
575    }
576
577    #[tokio::test]
578    async fn builder_with_interceptors() {
579        let pipeline = PipelineBuilder::new(MockExecutor::new())
580            .with_interceptor(PassthroughInterceptor {
581                name: "first".into(),
582            })
583            .with_interceptor(PassthroughInterceptor {
584                name: "second".into(),
585            })
586            .build()
587            .unwrap();
588        assert!(pipeline.schema().is_none());
589
590        let input = make_query_input(vec![]);
591        let output = pipeline.execute_query(input).await.unwrap();
592        assert_eq!(output.interceptors_applied, vec!["first", "second"]);
593    }
594
595    #[tokio::test]
596    #[cfg(feature = "v2-query")]
597    async fn v2_policy_rejection_is_fail_closed() {
598        let pipeline = PipelineBuilder::new(MockExecutor::new())
599            .with_interceptor(RejectingRequestInterceptor)
600            .build()
601            .unwrap();
602        let mut context = pipeline
603            .v2_request_context(HashMap::from([(
604                "transport".to_owned(),
605                serde_json::json!("v2"),
606            )]))
607            .expect("V2-aware policy coverage");
608        let error = pipeline
609            .begin_v2_request(&mut context)
610            .await
611            .expect_err("the same request policy must gate V2 envelopes");
612        assert!(
613            matches!(error, PipelineError::Interceptor(message) if message.contains("test rejection"))
614        );
615    }
616
617    #[cfg(feature = "v2-query")]
618    struct RewritingInterceptor;
619
620    #[cfg(feature = "v2-query")]
621    impl Interceptor for RewritingInterceptor {
622        fn name(&self) -> &str {
623            "rewriter"
624        }
625
626        fn on_request<'a>(
627            &'a self,
628            _clauses: Vec<Clause>,
629            _ctx: &'a mut RequestContext,
630        ) -> Pin<Box<dyn Future<Output = Result<Vec<Clause>, InterceptError>> + Send + 'a>>
631        {
632            Box::pin(async { Ok(make_simple_clauses()) })
633        }
634    }
635
636    #[tokio::test]
637    #[cfg(feature = "v2-query")]
638    async fn v2_rejects_legacy_ast_policies_without_typed_coverage() {
639        let pipeline = PipelineBuilder::new(MockExecutor::new())
640            .with_interceptor(RewritingInterceptor)
641            .build()
642            .unwrap();
643        let startup_error = pipeline
644            .validate_v2_coverage()
645            .expect_err("startup must reject a legacy-only policy before routing");
646        assert!(
647            matches!(startup_error, PipelineError::Interceptor(message) if message.contains("does not declare typed V2 coverage"))
648        );
649        let error = pipeline
650            .v2_request_context(HashMap::new())
651            .expect_err("request admission must retain the defensive recheck");
652        assert!(
653            matches!(error, PipelineError::Interceptor(message) if message.contains("does not declare typed V2 coverage"))
654        );
655    }
656
657    // =============================================
658    // execute_query tests
659    // =============================================
660
661    #[tokio::test]
662    async fn execute_query_uses_input_database() {
663        let executor = MockExecutor::new();
664        let calls = executor.calls.clone();
665        let pipeline = make_pipeline(executor, false);
666
667        let input = make_query_input_with_db(vec![], "custom_db");
668        pipeline.execute_query(input).await.unwrap();
669
670        let recorded = calls.lock().unwrap();
671        assert_eq!(recorded[0].0, "custom_db");
672    }
673
674    #[tokio::test]
675    async fn execute_query_uses_default_database_when_none() {
676        let executor = MockExecutor::new();
677        let calls = executor.calls.clone();
678        let pipeline = make_pipeline(executor, false);
679
680        let input = make_query_input(vec![]);
681        pipeline.execute_query(input).await.unwrap();
682
683        let recorded = calls.lock().unwrap();
684        assert_eq!(recorded[0].0, "test_db"); // from make_pipeline
685    }
686
687    #[tokio::test]
688    async fn execute_query_skips_validation_when_no_schema() {
689        let pipeline = make_pipeline(MockExecutor::new(), false);
690        let clauses = vec![Clause::Match(vec![Pattern::Entity {
691            variable: "x".to_string(),
692            type_name: "nonexistent_type".to_string(),
693            constraints: vec![],
694            is_strict: false,
695        }])];
696        let input = make_query_input(clauses);
697        let result = pipeline.execute_query(input).await;
698        assert!(result.is_ok());
699    }
700
701    #[tokio::test]
702    async fn execute_query_validates_when_schema_present_valid() {
703        let pipeline = make_pipeline(MockExecutor::new(), true);
704        let input = make_query_input(make_simple_clauses());
705        let result = pipeline.execute_query(input).await;
706        assert!(result.is_ok());
707    }
708
709    #[tokio::test]
710    async fn execute_query_validates_when_schema_present_invalid() {
711        let pipeline = make_pipeline(MockExecutor::new(), true);
712        let clauses = vec![Clause::Match(vec![Pattern::Entity {
713            variable: "p".to_string(),
714            type_name: "person".to_string(),
715            constraints: vec![Constraint::Has {
716                attr_name: "nonexistent_attr".to_string(),
717                value: Value::Literal(type_bridge_core_lib::ast::LiteralValue {
718                    value: serde_json::json!("val"),
719                    value_type: "string".to_string(),
720                }),
721            }],
722            is_strict: false,
723        }])];
724        let input = make_query_input(clauses);
725        let result = pipeline.execute_query(input).await;
726        let err = result.unwrap_err();
727        assert!(matches!(&err, PipelineError::Validation(msg) if msg.contains("validation error")));
728    }
729
730    #[tokio::test]
731    async fn execute_query_request_interceptor_failure() {
732        assert_eq!(RejectingRequestInterceptor.name(), "rejector");
733        let pipeline = PipelineBuilder::new(MockExecutor::new())
734            .with_interceptor(RejectingRequestInterceptor)
735            .build()
736            .unwrap();
737        let input = make_query_input(vec![]);
738        let result = pipeline.execute_query(input).await;
739        let err = result.unwrap_err();
740        assert!(matches!(&err, PipelineError::Interceptor(msg) if msg.contains("test rejection")));
741    }
742
743    #[tokio::test]
744    async fn execute_query_executor_failure() {
745        let pipeline = make_pipeline(MockExecutor::failing("db crash"), false);
746        let input = make_query_input(vec![]);
747        let result = pipeline.execute_query(input).await;
748        let err = result.unwrap_err();
749        assert!(matches!(&err, PipelineError::QueryExecution(msg) if msg.contains("db crash")));
750    }
751
752    #[tokio::test]
753    async fn execute_query_response_interceptor_failure() {
754        assert_eq!(RejectingResponseInterceptor.name(), "resp-rejector");
755        let pipeline = PipelineBuilder::new(MockExecutor::new())
756            .with_interceptor(RejectingResponseInterceptor)
757            .build()
758            .unwrap();
759        let input = make_query_input(vec![]);
760        let result = pipeline.execute_query(input).await;
761        let err = result.unwrap_err();
762        assert!(
763            matches!(&err, PipelineError::Interceptor(msg) if msg.contains("response rejected"))
764        );
765    }
766
767    #[tokio::test]
768    async fn execute_query_success_output_fields() {
769        let _guard = init_tracing();
770        let count = Arc::new(AtomicUsize::new(0));
771        let pipeline =
772            PipelineBuilder::new(MockExecutor::with_result(serde_json::json!({"ok": true})))
773                .with_default_database("test_db")
774                .with_interceptor(CountingInterceptor {
775                    name: "counter".into(),
776                    count: count.clone(),
777                })
778                .build()
779                .unwrap();
780
781        let input = make_query_input(vec![]);
782        let output = pipeline.execute_query(input).await.unwrap();
783
784        assert!(!output.request_id.is_empty());
785        assert_eq!(output.results, serde_json::json!({"ok": true}));
786        assert_eq!(output.interceptors_applied, vec!["counter"]);
787        assert_eq!(count.load(Ordering::SeqCst), 1);
788    }
789
790    #[tokio::test]
791    async fn execute_query_empty_clauses_success() {
792        let pipeline = make_pipeline(MockExecutor::new(), false);
793        let input = make_query_input(vec![]);
794        let result = pipeline.execute_query(input).await;
795        assert!(result.is_ok());
796    }
797
798    #[tokio::test]
799    async fn execute_query_compiled_typeql_in_metadata() {
800        let executor = MockExecutor::new();
801        let calls = executor.calls.clone();
802        let pipeline = make_pipeline(executor, false);
803
804        let clauses = make_simple_clauses();
805        let input = make_query_input(clauses);
806        pipeline.execute_query(input).await.unwrap();
807
808        let recorded = calls.lock().unwrap();
809        assert!(!recorded[0].1.is_empty());
810    }
811
812    #[tokio::test]
813    async fn execute_query_passes_transaction_type() {
814        let executor = MockExecutor::new();
815        let calls = executor.calls.clone();
816        let pipeline = make_pipeline(executor, false);
817
818        let input = QueryInput {
819            database: None,
820            transaction_type: "write".to_string(),
821            clauses: vec![],
822            metadata: HashMap::new(),
823        };
824        pipeline.execute_query(input).await.unwrap();
825
826        let recorded = calls.lock().unwrap();
827        assert_eq!(recorded[0].2, "write");
828    }
829
830    // =============================================
831    // validate tests
832    // =============================================
833
834    #[test]
835    fn validate_no_schema_returns_error() {
836        let pipeline = make_pipeline(MockExecutor::new(), false);
837        let input = ValidateInput { clauses: vec![] };
838        let result = pipeline.validate(&input);
839        let err = result.unwrap_err();
840        assert!(matches!(&err, PipelineError::Schema(msg) if msg.contains("No schema loaded")));
841    }
842
843    #[test]
844    fn validate_valid_clauses() {
845        let pipeline = make_pipeline(MockExecutor::new(), true);
846        let input = ValidateInput {
847            clauses: make_simple_clauses(),
848        };
849        let result = pipeline.validate(&input).unwrap();
850        assert!(result.is_valid);
851        assert!(result.errors.is_empty());
852    }
853
854    #[test]
855    fn validate_invalid_clauses() {
856        let pipeline = make_pipeline(MockExecutor::new(), true);
857        let input = ValidateInput {
858            clauses: vec![Clause::Match(vec![Pattern::Entity {
859                variable: "p".to_string(),
860                type_name: "person".to_string(),
861                constraints: vec![Constraint::Has {
862                    attr_name: "nonexistent_attr".to_string(),
863                    value: Value::Literal(type_bridge_core_lib::ast::LiteralValue {
864                        value: serde_json::json!("val"),
865                        value_type: "string".to_string(),
866                    }),
867                }],
868                is_strict: false,
869            }])],
870        };
871        let result = pipeline.validate(&input).unwrap();
872        assert!(!result.is_valid);
873        assert!(!result.errors.is_empty());
874    }
875
876    #[test]
877    fn validate_error_detail_fields() {
878        let pipeline = make_pipeline(MockExecutor::new(), true);
879        let input = ValidateInput {
880            clauses: vec![Clause::Match(vec![Pattern::Entity {
881                variable: "x".to_string(),
882                type_name: "person".to_string(),
883                constraints: vec![Constraint::Has {
884                    attr_name: "nonexistent_attr".to_string(),
885                    value: Value::Literal(type_bridge_core_lib::ast::LiteralValue {
886                        value: serde_json::json!("val"),
887                        value_type: "string".to_string(),
888                    }),
889                }],
890                is_strict: false,
891            }])],
892        };
893        let result = pipeline.validate(&input).unwrap();
894        assert!(!result.is_valid);
895        let error = &result.errors[0];
896        assert!(!error.code.is_empty());
897        assert!(!error.message.is_empty());
898    }
899
900    #[test]
901    fn validate_empty_clauses_with_schema() {
902        let pipeline = make_pipeline(MockExecutor::new(), true);
903        let input = ValidateInput { clauses: vec![] };
904        let result = pipeline.validate(&input).unwrap();
905        assert!(result.is_valid);
906    }
907
908    // =============================================
909    // Accessor tests
910    // =============================================
911
912    #[test]
913    fn schema_returns_some_when_loaded() {
914        let pipeline = make_pipeline(MockExecutor::new(), true);
915        assert!(pipeline.schema().is_some());
916    }
917
918    #[test]
919    fn schema_returns_none_when_not_loaded() {
920        let pipeline = make_pipeline(MockExecutor::new(), false);
921        assert!(pipeline.schema().is_none());
922    }
923
924    #[test]
925    fn is_connected_delegates_to_executor() {
926        let executor = MockExecutor::new();
927        *executor.connected.lock().unwrap() = true;
928        let pipeline = make_pipeline(executor, false);
929        assert!(pipeline.is_connected());
930    }
931
932    #[test]
933    fn is_connected_false_when_executor_disconnected() {
934        let executor = MockExecutor::new();
935        *executor.connected.lock().unwrap() = false;
936        let pipeline = make_pipeline(executor, false);
937        assert!(!pipeline.is_connected());
938    }
939
940    #[test]
941    fn default_database_returns_configured_value() {
942        let pipeline = PipelineBuilder::new(MockExecutor::new())
943            .with_default_database("my_database")
944            .build()
945            .unwrap();
946        assert_eq!(pipeline.default_database(), "my_database");
947    }
948
949    #[test]
950    fn default_database_empty_when_not_set() {
951        let pipeline = PipelineBuilder::new(MockExecutor::new()).build().unwrap();
952        assert_eq!(pipeline.default_database(), "");
953    }
954}