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
17pub struct QueryInput {
19 pub database: Option<String>,
21 pub transaction_type: String,
23 pub clauses: Vec<Clause>,
25 pub metadata: HashMap<String, serde_json::Value>,
27}
28
29pub struct ValidateInput {
31 pub clauses: Vec<Clause>,
33}
34
35#[derive(Debug)]
37pub struct QueryOutput {
38 pub results: serde_json::Value,
40 pub request_id: String,
42 pub execution_time_ms: u64,
44 pub interceptors_applied: Vec<String>,
46}
47
48#[derive(Debug)]
50pub struct ValidateOutput {
51 pub is_valid: bool,
53 pub errors: Vec<ValidationErrorDetail>,
55}
56
57#[derive(Debug)]
59pub struct ValidationErrorDetail {
60 pub code: String,
62 pub message: String,
64 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
74pub 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 #[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 #[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 #[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 #[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 #[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 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 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 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 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 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 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 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 pub fn schema(&self) -> Option<&TypeSchema> {
280 self.schema.as_ref()
281 }
282
283 pub fn is_connected(&self) -> bool {
285 self.executor.is_connected()
286 }
287
288 pub fn default_database(&self) -> &str {
290 &self.default_database
291 }
292}
293
294pub 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 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 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 pub fn with_interceptor(mut self, interceptor: impl Interceptor + 'static) -> Self {
338 self.interceptors.push(Box::new(interceptor));
339 self
340 }
341
342 pub fn with_default_database(mut self, database: impl Into<String>) -> Self {
344 self.default_database = database.into();
345 self
346 }
347
348 pub fn with_crud_interceptor(self, interceptor: impl CrudInterceptor + 'static) -> Self {
354 self.with_interceptor(CrudInterceptorAdapter::new(interceptor))
355 }
356
357 pub fn with_skip_validation(mut self) -> Self {
362 self.skip_validation = true;
363 self
364 }
365
366 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 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 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 #[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 #[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"); }
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 #[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 #[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}