1use std::sync::Arc;
2use std::time::Duration;
3
4use async_trait::async_trait;
5use bytes::Bytes;
6use camel_api::datasource::DatasourceCatalog;
7use futures::TryStreamExt;
8use serde_json::Value as JsonValue;
9use sqlx::AnyPool;
10use sqlx::any::AnyPoolOptions;
11use sqlx::any::AnyRow;
12use tokio::sync::OnceCell;
13use tracing::{debug, error, info, warn};
14
15use camel_component_api::retry_async;
16use camel_component_api::{
17 Body, CamelError, Exchange, Message, RuntimeObservability, StreamBody, StreamMetadata,
18};
19use camel_component_api::{ConcurrencyModel, Consumer, ConsumerContext};
20
21use crate::config::{
22 PollStrategy, ProcessingStrategy, SqlEndpointConfig, SqlOutputType, TransactionMode,
23 enrich_db_url_with_ssl, redact_db_url,
24};
25use crate::headers;
26use crate::query::{QueryTemplate, parse_query_template, resolve_params};
27use crate::utils::{bind_json_values, is_retryable_sqlx_error, row_to_json};
28
29fn record_post_process_failure(
43 runtime: &dyn RuntimeObservability,
44 route_id: &str,
45 label: &str,
46 error: &CamelError,
47 message: &str,
48) {
49 runtime.metrics().increment_errors(route_id, label);
51 error!(error = %error, "{message}");
53}
54
55#[derive(Debug, Clone, Copy, Default)]
59struct PollOutcome {
60 was_empty: bool,
61}
62
63pub struct SqlConsumer {
64 pub(crate) config: SqlEndpointConfig,
65 pub(crate) pool: Arc<OnceCell<Arc<AnyPool>>>,
66 pub(crate) catalog: Option<Arc<dyn DatasourceCatalog>>,
67 stopped: bool,
68 runtime: Arc<dyn RuntimeObservability>,
71}
72
73impl SqlConsumer {
74 pub fn new(
75 config: SqlEndpointConfig,
76 pool: Arc<OnceCell<Arc<AnyPool>>>,
77 catalog: Option<Arc<dyn DatasourceCatalog>>,
78 runtime: Arc<dyn RuntimeObservability>,
79 ) -> Self {
80 Self {
81 config,
82 pool,
83 catalog,
84 stopped: false,
85 runtime,
86 }
87 }
88
89 async fn poll_database(
91 &self,
92 pool: &AnyPool,
93 context: &ConsumerContext,
94 template: &QueryTemplate,
95 ) -> Result<PollOutcome, CamelError> {
96 let route_id = context.route_id();
98
99 let empty_exchange = Exchange::new(Message::default());
101
102 let prepared = resolve_params(template, &empty_exchange, &self.config.in_separator)?;
104
105 debug!(query = %prepared.sql, "executing SQL consumer poll");
106
107 if self.config.output_type == SqlOutputType::StreamList {
108 return self.poll_database_stream(pool, context, &prepared).await;
109 }
110
111 let query = bind_json_values(sqlx::query(&prepared.sql), &prepared.bindings);
112 let rows: Vec<AnyRow> = query.fetch_all(pool).await.map_err(|e| {
113 warn!(error = %e, "SQL consumer poll query failed");
114 CamelError::ProcessorError(format!("Query execution failed: {}", e))
115 })?;
116
117 debug!(rows = rows.len(), "SQL consumer poll completed");
118
119 let was_empty = rows.is_empty();
120 if was_empty && !self.config.route_empty_result_set {
121 return Ok(PollOutcome { was_empty });
122 }
123
124 let rows_to_process: Vec<AnyRow> = if let Some(max) = self.config.max_messages_per_poll {
125 if max > 0 {
126 rows.into_iter().take(max as usize).collect()
127 } else {
128 rows
129 }
130 } else {
131 rows
132 };
133
134 if self.config.use_iterator {
135 for row in rows_to_process {
137 let row_json = row_to_json(&row)?;
138
139 let mut msg = Message::new(Body::Json(row_json.clone()));
141
142 if let Some(obj) = row_json.as_object() {
144 for (key, value) in obj {
145 msg.set_header(format!("CamelSql.{}", key), value.clone());
146 }
147 }
148
149 let exchange = Exchange::new(msg);
150
151 let result = context.send_and_wait(exchange).await;
153
154 if let Err(e) = self.handle_post_processing(pool, &result, &row_json).await {
156 record_post_process_failure(
157 self.runtime.as_ref(),
158 route_id,
159 "b-prime:sql:on-consume",
160 &e,
161 "Post-processing failed",
162 );
163 if self.config.break_batch_on_consume_fail {
164 return Err(e);
165 }
166 }
167
168 if let Err(ref consume_err) = result
170 && self.config.break_batch_on_consume_fail
171 {
172 return Err(consume_err.clone());
173 }
174 }
175 } else {
176 let rows_json: Vec<JsonValue> = rows_to_process
178 .iter()
179 .map(row_to_json)
180 .collect::<Result<Vec<_>, CamelError>>()?;
181
182 let row_count = rows_json.len();
183
184 let mut msg = Message::new(Body::Json(JsonValue::Array(rows_json.clone())));
186 msg.set_header(headers::ROW_COUNT, JsonValue::Number(row_count.into()));
187
188 let exchange = Exchange::new(msg);
189
190 let result = context.send_and_wait(exchange).await;
192
193 for row_json in rows_json.iter() {
197 if let Err(e) = self.handle_post_processing(pool, &result, row_json).await {
198 record_post_process_failure(
199 self.runtime.as_ref(),
200 route_id,
201 "b-prime:sql:on-consume-batch",
202 &e,
203 "Post-processing failed for batch row",
204 );
205 if self.config.break_batch_on_consume_fail {
206 return Err(e);
207 }
208 }
209 }
210
211 if let Err(ref consume_err) = result
213 && self.config.break_batch_on_consume_fail
214 {
215 return Err(consume_err.clone());
216 }
217 }
218
219 if let Some(ref batch_query) = self.config.on_consume_batch_complete {
221 let _ = self
222 .execute_post_query(pool, batch_query, &JsonValue::Null)
223 .await;
224 }
225
226 Ok(PollOutcome { was_empty })
227 }
228
229 async fn poll_database_stream(
230 &self,
231 pool: &AnyPool,
232 context: &ConsumerContext,
233 prepared: &crate::query::PreparedQuery,
234 ) -> Result<PollOutcome, CamelError> {
235 let pool_clone = pool.clone();
236 let sql_str = prepared.sql.clone();
237 let bindings = prepared.bindings.clone();
238
239 let byte_stream = async_stream::try_stream! {
240 let mut q = sqlx::query(&sql_str);
241 q = bind_json_values(q, &bindings);
242 let mut rows = q.fetch(&pool_clone);
243 while let Some(row) = rows.try_next().await.map_err(|e| {
244 CamelError::ProcessorError(format!("Query execution failed: {}", e))
245 })? {
246 let json_val = row_to_json(&row).map_err(|e| {
247 CamelError::ProcessorError(format!("JSON serialization failed: {}", e))
248 })?;
249 let mut bytes = serde_json::to_vec(&json_val)
250 .map_err(|e| CamelError::ProcessorError(format!("JSON serialization failed: {}", e)))?;
251 bytes.push(b'\n');
252 yield Bytes::from(bytes);
253 }
254 };
255
256 let msg = Message::new(Body::Stream(StreamBody {
257 stream: Arc::new(tokio::sync::Mutex::new(Some(Box::pin(byte_stream)))),
258 metadata: StreamMetadata {
259 content_type: Some("application/x-ndjson".to_string()),
260 size_hint: None,
261 origin: None,
262 },
263 }));
264
265 let exchange = Exchange::new(msg);
266 let result = context.send_and_wait(exchange).await;
267 if let Err(e) = result {
268 record_post_process_failure(
269 self.runtime.as_ref(),
270 context.route_id(),
271 "b-prime:sql:stream-list",
272 &e,
273 "StreamList consumer downstream processing failed",
274 );
275 return Err(e);
276 }
277
278 debug!("StreamList: consumer poll completed (lazy stream emitted)");
279 Ok(PollOutcome::default())
281 }
282
283 async fn handle_post_processing(
285 &self,
286 pool: &AnyPool,
287 result: &Result<Exchange, CamelError>,
288 row_json: &JsonValue,
289 ) -> Result<(), CamelError> {
290 match result {
291 Ok(_) => {
292 if let Some(ref on_consume) = self.config.on_consume {
294 self.execute_post_query(pool, on_consume, row_json).await?;
295 }
296 }
297 Err(_) => {
298 if let Some(ref on_consume_failed) = self.config.on_consume_failed {
300 self.execute_post_query(pool, on_consume_failed, row_json)
301 .await?;
302 }
303 }
304 }
305 Ok(())
306 }
307
308 async fn execute_post_query(
310 &self,
311 pool: &AnyPool,
312 query_str: &str,
313 row_json: &JsonValue,
314 ) -> Result<(), CamelError> {
315 let template = parse_query_template(query_str, self.config.placeholder)?;
317
318 let mut temp_msg = Message::new(Body::Json(row_json.clone()));
321 if let Some(obj) = row_json.as_object() {
322 for (key, value) in obj {
323 temp_msg.set_header(format!("CamelSql.{}", key), value.clone());
324 }
325 }
326 let temp_exchange = Exchange::new(temp_msg);
327
328 let prepared = resolve_params(&template, &temp_exchange, &self.config.in_separator)?;
330
331 let query = bind_json_values(sqlx::query(&prepared.sql), &prepared.bindings);
333 let result = query.execute(pool).await.map_err(|e| {
334 CamelError::ProcessorError(format!("Post-query execution failed: {}", e))
335 })?;
336
337 if result.rows_affected() == 0 {
339 warn!(
340 query = query_str,
341 "Post-processing query affected 0 rows — the row may not have been marked correctly"
342 );
343 }
344
345 Ok(())
346 }
347
348 async fn handle_poll_result(
352 &self,
353 pool: &AnyPool,
354 context: &ConsumerContext,
355 template: &QueryTemplate,
356 ) -> PollOutcome {
357 match self.poll_database(pool, context, template).await {
358 Ok(outcome) => outcome,
359 Err(e) => {
360 if self.config.bridge_error_handler {
362 warn!(error = %e, "SQL consumer poll failed (bridged)");
366 if let Err(route_err) = self.bridge_poll_error(context, e).await {
367 error!(error = %route_err, "Failed to bridge SQL consumer error to route");
370 }
371 } else {
372 record_post_process_failure(
373 self.runtime.as_ref(),
374 context.route_id(),
375 "b-prime:sql:poll-failed",
376 &e,
377 "SQL consumer poll failed",
378 );
379 }
380 PollOutcome::default()
382 }
383 }
384 }
385
386 async fn bridge_poll_error(
387 &self,
388 context: &ConsumerContext,
389 error: CamelError,
390 ) -> Result<(), CamelError> {
391 if !self.config.bridge_error_handler {
392 return Ok(());
393 }
394 let mut exchange = Exchange::new(Message::default());
395 exchange.set_error(error);
396 context.send_and_wait(exchange).await.map(|_| ())
397 }
398}
399
400#[async_trait]
401impl Consumer for SqlConsumer {
402 async fn start(&mut self, context: ConsumerContext) -> Result<(), CamelError> {
403 if self.stopped {
405 return Err(CamelError::Config(
406 "SQL consumer cannot be restarted after stop".into(),
407 ));
408 }
409
410 let route_id = context.route_id().to_string();
412 let catalog = self.catalog.clone();
413 let ds_name = self.config.datasource_name.clone();
414
415 self.config.resolve_defaults();
417 self.config.resolve_file_query().await?;
418
419 let pool = self
420 .pool
421 .get_or_try_init(|| async {
422 if let (Some(ref cat), Some(ref name)) = (catalog, ds_name) {
424 let handle = cat.get_pool(name).await?;
425 return handle.downcast::<AnyPool>();
426 }
427
428 sqlx::any::install_default_drivers();
431 let db_url = enrich_db_url_with_ssl(&self.config.db_url, &self.config)?;
432
433 let max_conn = self.config.max_connections.ok_or_else(|| {
434 CamelError::Config("max_connections not resolved for SQL consumer pool".into())
435 })?;
436 let min_conn = self.config.min_connections.ok_or_else(|| {
437 CamelError::Config("min_connections not resolved for SQL consumer pool".into())
438 })?;
439 let idle_timeout = self.config.idle_timeout_secs.ok_or_else(|| {
440 CamelError::Config(
441 "idle_timeout_secs not resolved for SQL consumer pool".into(),
442 )
443 })?;
444 let max_lifetime = self.config.max_lifetime_secs.ok_or_else(|| {
445 CamelError::Config(
446 "max_lifetime_secs not resolved for SQL consumer pool".into(),
447 )
448 })?;
449
450 info!(
451 db_url = %redact_db_url(&self.config.db_url),
452 "SQL consumer pool initializing"
453 );
454 let retry_policy = &self.config.retry;
455 let pool = retry_async::<_, _, _, _, sqlx::Error>(
456 retry_policy,
457 "sql",
458 "consumer-pool-init",
459 || {
460 async {
461 AnyPoolOptions::new()
462 .max_connections(max_conn)
463 .min_connections(min_conn)
464 .idle_timeout(Duration::from_secs(idle_timeout))
465 .max_lifetime(Duration::from_secs(max_lifetime))
466 .connect(&db_url)
467 .await
468 }
469 },
470 is_retryable_sqlx_error,
471 Some(self.runtime.metrics().as_ref()),
472 )
473 .await
474 .map_err(|e| {
475 self.runtime.health().force_unhealthy_for_route(
476 &route_id,
477 "g:sql:consumer-pool-init",
478 &e.to_string(),
479 );
480 error!(error = %e, db_url = %redact_db_url(&self.config.db_url), "SQL connect failed, giving up");
482 CamelError::EndpointCreationFailed(format!(
483 "Failed to connect to database: {}",
484 e
485 ))
486 })?;
487 Ok(Arc::new(pool))
488 })
489 .await?;
490
491 if self.config.transaction_mode == TransactionMode::Managed {
493 warn!("transactionManager not yet implemented; using Auto mode");
494 }
495
496 if self.config.processing_strategy == ProcessingStrategy::Scheduled {
498 debug!(
499 "Processing strategy: Scheduled (rows dispatched individually via send_and_wait)"
500 );
501 }
502 if self.config.poll_strategy == PollStrategy::Burst {
503 debug!("Poll strategy: Burst (rapid successive polls)");
504 }
505
506 if self.config.output_type == SqlOutputType::StreamList
507 && (self.config.on_consume.is_some()
508 || self.config.on_consume_failed.is_some()
509 || self.config.on_consume_batch_complete.is_some()
510 || self.config.break_on_empty)
511 {
512 warn!(
513 "onConsume/onConsumeFailed/onConsumeBatchComplete/breakOnEmpty are not executed in \
514 StreamList mode (rows are consumed lazily downstream)"
515 );
516 }
517
518 if self.config.on_consume.is_none() {
520 warn!(
521 "SQL consumer started without onConsume configured — consumed rows will not be marked/deleted"
522 );
523 }
524
525 info!(
526 db_url = %redact_db_url(&self.config.db_url),
527 query_len = self.config.query.len(),
528 "SQL consumer started"
529 );
530
531 let template = parse_query_template(&self.config.query, self.config.placeholder)
533 .map_err(|e| CamelError::Config(format!("Invalid query template: {}", e)))?;
534
535 if self.config.initial_delay_ms > 0 {
537 tokio::select! {
538 _ = context.cancelled() => {
539 info!("SQL consumer stopped during initial delay");
540 return Ok(());
541 }
542 _ = tokio::time::sleep(Duration::from_millis(self.config.initial_delay_ms)) => {}
543 }
544 }
545
546 let mut poll_count: u32 = 0;
562 loop {
563 if let Some(max_repeats) = self.config.repeat_count
565 && poll_count >= max_repeats
566 {
567 info!(
568 repeat_count = max_repeats,
569 "SQL consumer reached repeat_count limit, stopping"
570 );
571 break;
572 }
573
574 tokio::select! {
575 _ = context.cancelled() => {
576 info!("SQL consumer stopped");
577 break;
578 }
579 _ = tokio::time::sleep(Duration::from_millis(self.config.delay_ms)) => {
580 poll_count += 1;
581 let outcome = self.handle_poll_result(pool.as_ref(), &context, &template).await;
582 if self.config.break_on_empty && outcome.was_empty {
583 info!("SQL consumer stopping: break_on_empty triggered (poll returned 0 rows)");
584 break;
585 }
586 }
587 }
588 }
589
590 Ok(())
591 }
592
593 async fn stop(&mut self) -> Result<(), CamelError> {
594 if self.stopped {
596 debug!("SQL consumer stop called on already-stopped consumer");
597 return Ok(());
598 }
599
600 if let Some(pool) = self.pool.get() {
602 debug!("SQL consumer closing connection pool");
603 pool.close().await;
604 debug!("SQL consumer pool closed");
605 }
606
607 self.stopped = true;
608 info!("SQL consumer stopped");
609 Ok(())
610 }
611
612 fn concurrency_model(&self) -> ConcurrencyModel {
613 ConcurrencyModel::Sequential
617 }
618}
619
620#[cfg(test)]
621mod tests {
622 use super::*;
623 use camel_api::MetricsCollector;
624 use camel_component_api::HealthCheckRegistry;
625 use camel_component_api::test_support::PanicRuntimeObservability;
626 fn test_rt() -> std::sync::Arc<dyn camel_component_api::RuntimeObservability> {
627 std::sync::Arc::new(PanicRuntimeObservability)
628 }
629 use crate::config::SqlEndpointConfig;
630 use camel_component_api::ExchangeEnvelope;
631 use camel_component_api::UriConfig;
632 use sqlx::any::AnyPoolOptions;
633 use std::sync::Arc;
634 use std::sync::Mutex;
635 use std::time::Duration;
636 use tokio::sync::mpsc;
637 use tokio::time::timeout;
638 use tokio_util::sync::CancellationToken;
639
640 struct RecordingMetrics {
645 errors: Arc<Mutex<Vec<(String, String)>>>,
646 }
647
648 impl MetricsCollector for RecordingMetrics {
649 fn record_exchange_duration(&self, _: &str, _: Duration) {}
650 fn increment_errors(&self, route_id: &str, error_type: &str) {
651 self.errors
652 .lock()
653 .unwrap()
654 .push((route_id.to_string(), error_type.to_string()));
655 }
656 fn increment_exchanges(&self, _: &str) {}
657 fn set_queue_depth(&self, _: &str, _: usize) {}
658 fn record_circuit_breaker_change(&self, _: &str, _: &str, _: &str) {}
659 }
660
661 struct RecordingRuntime {
662 metrics_collector: Arc<RecordingMetrics>,
663 }
664
665 impl RecordingRuntime {
666 fn new(errors: Arc<Mutex<Vec<(String, String)>>>) -> Self {
667 Self {
668 metrics_collector: Arc::new(RecordingMetrics { errors }),
669 }
670 }
671 }
672
673 impl RuntimeObservability for RecordingRuntime {
674 fn metrics(&self) -> Arc<dyn MetricsCollector> {
675 self.metrics_collector.clone() as Arc<dyn MetricsCollector>
676 }
677 fn health(&self) -> Arc<dyn HealthCheckRegistry> {
678 panic!("RecordingRuntime::health not used in this test")
679 }
680 }
681
682 #[tracing_test::traced_test]
686 #[test]
687 fn record_post_process_failure_increments_errors_and_emits_error_log() {
688 let errors: Arc<Mutex<Vec<(String, String)>>> = Arc::new(Mutex::new(Vec::new()));
689 let runtime = Arc::new(RecordingRuntime::new(Arc::clone(&errors)));
690 let error = CamelError::ProcessorError("test failure".to_string());
691
692 record_post_process_failure(
694 runtime.as_ref(),
695 "test-route",
696 "b-prime:sql:on-consume",
697 &error,
698 "Post-processing failed",
699 );
700
701 let recorded = errors.lock().unwrap();
703 assert_eq!(recorded.len(), 1, "expected 1 increment_errors call");
704 assert_eq!(recorded[0].0, "test-route");
705 assert_eq!(recorded[0].1, "b-prime:sql:on-consume");
706 drop(recorded);
707
708 assert!(logs_contain("ERROR"), "helper must emit error! log");
710 assert!(
711 logs_contain("Post-processing failed"),
712 "helper must include the message in the log"
713 );
714 }
715
716 async fn sqlite_pool() -> AnyPool {
717 sqlx::any::install_default_drivers();
718 tokio::time::timeout(
721 std::time::Duration::from_secs(10),
722 AnyPoolOptions::new()
723 .max_connections(1)
724 .connect("sqlite::memory:"),
725 )
726 .await
727 .expect("sqlite pool connect timed out after 10s")
728 .expect("sqlite pool")
729 }
730
731 async fn seed_consumer_table(pool: &AnyPool) {
732 sqlx::query("CREATE TABLE jobs (id INTEGER PRIMARY KEY, processed INTEGER DEFAULT 0, failed INTEGER DEFAULT 0)")
733 .execute(pool)
734 .await
735 .expect("create table");
736 sqlx::query("INSERT INTO jobs (id, processed, failed) VALUES (1, 0, 0), (2, 0, 0)")
737 .execute(pool)
738 .await
739 .expect("seed rows");
740 }
741
742 fn config() -> SqlEndpointConfig {
743 let mut c =
744 SqlEndpointConfig::from_uri("sql:select * from t?db_url=postgres://localhost/test")
745 .unwrap();
746 c.resolve_defaults();
747 c
748 }
749
750 #[test]
751 fn consumer_concurrency_model() {
752 let c = SqlConsumer::new(config(), Arc::new(OnceCell::new()), None, test_rt());
753 assert_eq!(c.concurrency_model(), ConcurrencyModel::Sequential);
754 }
755
756 #[test]
757 fn consumer_stores_config() {
758 let mut config = SqlEndpointConfig::from_uri(
759 "sql:select * from t?db_url=postgres://localhost/test&delay=2000&onConsume=update t set done=true"
760 ).unwrap();
761 config.resolve_defaults();
762 let c = SqlConsumer::new(config.clone(), Arc::new(OnceCell::new()), None, test_rt());
763 assert_eq!(c.config.delay_ms, 2000);
764 assert!(c.config.on_consume.is_some());
765 }
766
767 #[tokio::test]
768 async fn poll_database_runs_on_consume_for_successful_rows() {
769 let pool = sqlite_pool().await;
770 seed_consumer_table(&pool).await;
771
772 let mut config = SqlEndpointConfig::from_uri(
773 "sql:select id, processed, failed from jobs where processed = 0 order by id?db_url=sqlite::memory:&onConsume=update jobs set processed=1 where id=:#id&initialDelay=0&delay=1",
774 )
775 .unwrap();
776 config.resolve_defaults();
777
778 let consumer = SqlConsumer::new(config.clone(), Arc::new(OnceCell::new()), None, test_rt());
779 let template = parse_query_template(&config.query, config.placeholder).unwrap();
780
781 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
782 tokio::spawn(async move {
783 loop {
784 match timeout(Duration::from_secs(2), rx.recv()).await {
785 Ok(Some(env)) => {
786 if let Some(reply_tx) = env.reply_tx {
787 let _ = reply_tx.send(Ok(env.exchange));
788 }
789 }
790 Ok(None) => break, Err(_) => break, }
793 }
794 });
795 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
796
797 consumer
798 .poll_database(&pool, &ctx, &template)
799 .await
800 .expect("poll must succeed");
801
802 let row = sqlx::query("select processed from jobs where id = 1")
803 .fetch_one(&pool)
804 .await
805 .expect("row 1");
806 let processed_1: i64 = sqlx::Row::try_get(&row, 0).expect("processed");
807
808 let row = sqlx::query("select processed from jobs where id = 2")
809 .fetch_one(&pool)
810 .await
811 .expect("row 2");
812 let processed_2: i64 = sqlx::Row::try_get(&row, 0).expect("processed");
813
814 assert_eq!(processed_1, 1);
815 assert_eq!(processed_2, 1);
816 }
817
818 #[tokio::test]
819 async fn poll_database_runs_on_consume_failed_when_downstream_fails() {
820 let pool = sqlite_pool().await;
821 seed_consumer_table(&pool).await;
822
823 let mut config = SqlEndpointConfig::from_uri(
824 "sql:select id, processed, failed from jobs where processed = 0 order by id?db_url=sqlite::memory:&onConsumeFailed=update jobs set failed=1 where id=:#id&initialDelay=0&delay=1",
825 )
826 .unwrap();
827 config.resolve_defaults();
828
829 let consumer = SqlConsumer::new(config.clone(), Arc::new(OnceCell::new()), None, test_rt());
830 let template = parse_query_template(&config.query, config.placeholder).unwrap();
831
832 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
833 tokio::spawn(async move {
834 loop {
835 match timeout(Duration::from_secs(2), rx.recv()).await {
836 Ok(Some(env)) => {
837 if let Some(reply_tx) = env.reply_tx {
838 let _ = reply_tx
839 .send(Err(CamelError::ProcessorError("downstream boom".into())));
840 }
841 }
842 Ok(None) => break, Err(_) => break, }
845 }
846 });
847 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
848
849 consumer
850 .poll_database(&pool, &ctx, &template)
851 .await
852 .expect("consumer should swallow downstream errors when breakBatchOnConsumeFail=false");
853
854 let row = sqlx::query("select failed from jobs where id = 1")
855 .fetch_one(&pool)
856 .await
857 .expect("row 1");
858 let failed_1: i64 = sqlx::Row::try_get(&row, 0).expect("failed");
859
860 let row = sqlx::query("select failed from jobs where id = 2")
861 .fetch_one(&pool)
862 .await
863 .expect("row 2");
864 let failed_2: i64 = sqlx::Row::try_get(&row, 0).expect("failed");
865
866 assert_eq!(failed_1, 1);
867 assert_eq!(failed_2, 1);
868 }
869
870 #[tokio::test]
871 async fn poll_database_breaks_batch_on_consume_fail() {
872 let pool = sqlite_pool().await;
873 seed_consumer_table(&pool).await;
874
875 let mut config = SqlEndpointConfig::from_uri(
876 "sql:select id, processed, failed from jobs where processed = 0 order by id?db_url=sqlite::memory:&onConsumeFailed=update jobs set failed=1 where id=:#id&breakBatchOnConsumeFail=true&initialDelay=0&delay=1",
877 )
878 .unwrap();
879 config.resolve_defaults();
880
881 let consumer = SqlConsumer::new(config.clone(), Arc::new(OnceCell::new()), None, test_rt());
882 let template = parse_query_template(&config.query, config.placeholder).unwrap();
883
884 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
885 tokio::spawn(async move {
886 loop {
887 match timeout(Duration::from_secs(2), rx.recv()).await {
888 Ok(Some(env)) => {
889 if let Some(reply_tx) = env.reply_tx {
890 let _ = reply_tx
891 .send(Err(CamelError::ProcessorError("downstream boom".into())));
892 }
893 }
894 Ok(None) => break, Err(_) => break, }
897 }
898 });
899 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
900
901 let err = consumer
902 .poll_database(&pool, &ctx, &template)
903 .await
904 .expect_err("must stop on first downstream failure");
905 assert!(err.to_string().contains("downstream boom"));
906
907 let row = sqlx::query("select failed from jobs where id = 1")
908 .fetch_one(&pool)
909 .await
910 .expect("row 1");
911 let failed_1: i64 = sqlx::Row::try_get(&row, 0).expect("failed");
912
913 let row = sqlx::query("select failed from jobs where id = 2")
914 .fetch_one(&pool)
915 .await
916 .expect("row 2");
917 let failed_2: i64 = sqlx::Row::try_get(&row, 0).expect("failed");
918
919 assert_eq!(failed_1, 1);
920 assert_eq!(failed_2, 0, "second row must not be processed");
921 }
922
923 #[tokio::test]
929 async fn consumer_no_panic_without_prior_resolve_defaults() {
930 let config = SqlEndpointConfig::from_uri(
931 "sql:select 1?db_url=sqlite::memory:&initialDelay=0&delay=1",
932 )
933 .unwrap();
934 assert!(config.max_connections.is_none());
936
937 let mut consumer = SqlConsumer::new(
938 config,
939 Arc::new(OnceCell::new()),
940 None,
941 std::sync::Arc::new(camel_component_api::test_support::NoopRuntimeObservability),
944 );
945 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
946 tokio::spawn(async move {
947 loop {
948 match timeout(Duration::from_secs(2), rx.recv()).await {
949 Ok(Some(env)) => {
950 if let Some(reply_tx) = env.reply_tx {
951 let _ = reply_tx.send(Ok(env.exchange));
952 }
953 }
954 Ok(None) => break, Err(_) => break, }
957 }
958 });
959 let token = CancellationToken::new();
960 let ctx = ConsumerContext::new(tx, token.clone(), "sql-test-route".to_string());
961
962 let consumer_handle = tokio::spawn(async move { consumer.start(ctx).await });
964
965 tokio::time::sleep(Duration::from_millis(50)).await;
967 token.cancel();
968
969 let result = consumer_handle.await.expect("task should not panic");
970 let _ = result;
972 }
973
974 #[tokio::test]
976 async fn stop_closes_pool() {
977 let pool = sqlite_pool().await;
978 seed_consumer_table(&pool).await;
979
980 let mut config = SqlEndpointConfig::from_uri(
981 "sql:select id from jobs?db_url=sqlite::memory:&onConsume=update jobs set processed=1 where id=:#id&initialDelay=0&delay=1",
982 )
983 .unwrap();
984 config.resolve_defaults();
985
986 let pool_cell = Arc::new(OnceCell::new());
987 pool_cell.set(Arc::new(pool.clone())).unwrap();
988
989 let mut consumer = SqlConsumer::new(config, pool_cell, None, test_rt());
990 consumer.stop().await.expect("stop should succeed");
991
992 assert!(
994 pool.is_closed(),
995 "Pool should be closed after consumer.stop()"
996 );
997 }
998
999 #[tokio::test]
1001 async fn double_stop_is_safe() {
1002 let pool = sqlite_pool().await;
1003 let mut config = SqlEndpointConfig::from_uri(
1004 "sql:select 1?db_url=sqlite::memory:&initialDelay=0&delay=1",
1005 )
1006 .unwrap();
1007 config.resolve_defaults();
1008
1009 let pool_cell = Arc::new(OnceCell::new());
1010 pool_cell.set(Arc::new(pool.clone())).unwrap();
1011
1012 let mut consumer = SqlConsumer::new(config, pool_cell, None, test_rt());
1013 consumer.stop().await.expect("first stop should succeed");
1014 consumer
1015 .stop()
1016 .await
1017 .expect("second stop should also succeed");
1018 }
1019
1020 #[tokio::test]
1022 async fn start_after_stop_rejected() {
1023 let pool = sqlite_pool().await;
1024 let mut config = SqlEndpointConfig::from_uri(
1025 "sql:select 1?db_url=sqlite::memory:&initialDelay=0&delay=1",
1026 )
1027 .unwrap();
1028 config.resolve_defaults();
1029
1030 let pool_cell = Arc::new(OnceCell::new());
1031 pool_cell.set(Arc::new(pool.clone())).unwrap();
1032
1033 let mut consumer = SqlConsumer::new(config, pool_cell, None, test_rt());
1034 consumer.stop().await.expect("stop should succeed");
1035
1036 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
1037 tokio::spawn(async move {
1038 loop {
1039 match timeout(Duration::from_secs(2), rx.recv()).await {
1040 Ok(Some(env)) => {
1041 if let Some(reply_tx) = env.reply_tx {
1042 let _ = reply_tx.send(Ok(env.exchange));
1043 }
1044 }
1045 Ok(None) => break, Err(_) => break, }
1048 }
1049 });
1050 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1051
1052 let result = consumer.start(ctx).await;
1053 assert!(result.is_err());
1054 let err_msg = result.unwrap_err().to_string();
1055 assert!(
1056 err_msg.contains("cannot be restarted") || err_msg.contains("after stop"),
1057 "Expected restart error, got: {}",
1058 err_msg
1059 );
1060 }
1061
1062 #[tokio::test]
1064 async fn batch_mode_per_row_post_processing() {
1065 let pool = sqlite_pool().await;
1066 seed_consumer_table(&pool).await;
1067
1068 let mut config = SqlEndpointConfig::from_uri(
1069 "sql:select id, processed, failed from jobs where processed = 0 order by id?db_url=sqlite::memory:&onConsume=update jobs set processed=1 where id=:#id&useIterator=false&initialDelay=0&delay=1",
1070 )
1071 .unwrap();
1072 config.resolve_defaults();
1073
1074 let consumer = SqlConsumer::new(config.clone(), Arc::new(OnceCell::new()), None, test_rt());
1075 let template = parse_query_template(&config.query, config.placeholder).unwrap();
1076
1077 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
1078 tokio::spawn(async move {
1079 loop {
1080 match timeout(Duration::from_secs(2), rx.recv()).await {
1081 Ok(Some(env)) => {
1082 if let Some(reply_tx) = env.reply_tx {
1083 let _ = reply_tx.send(Ok(env.exchange));
1084 }
1085 }
1086 Ok(None) => break, Err(_) => break, }
1089 }
1090 });
1091 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1092
1093 consumer
1094 .poll_database(&pool, &ctx, &template)
1095 .await
1096 .expect("poll must succeed");
1097
1098 let row = sqlx::query("select processed from jobs where id = 1")
1100 .fetch_one(&pool)
1101 .await
1102 .expect("row 1");
1103 let processed_1: i64 = sqlx::Row::try_get(&row, 0).expect("processed");
1104
1105 let row = sqlx::query("select processed from jobs where id = 2")
1106 .fetch_one(&pool)
1107 .await
1108 .expect("row 2");
1109 let processed_2: i64 = sqlx::Row::try_get(&row, 0).expect("processed");
1110
1111 assert_eq!(
1112 processed_1, 1,
1113 "row 1 should be marked processed via per-row onConsume"
1114 );
1115 assert_eq!(
1116 processed_2, 1,
1117 "row 2 should be marked processed via per-row onConsume"
1118 );
1119 }
1120
1121 #[tokio::test]
1123 async fn batch_mode_per_row_post_processing_on_failure() {
1124 let pool = sqlite_pool().await;
1125 seed_consumer_table(&pool).await;
1126
1127 let mut config = SqlEndpointConfig::from_uri(
1128 "sql:select id, processed, failed from jobs where processed = 0 order by id?db_url=sqlite::memory:&onConsumeFailed=update jobs set failed=1 where id=:#id&useIterator=false&initialDelay=0&delay=1",
1129 )
1130 .unwrap();
1131 config.resolve_defaults();
1132
1133 let consumer = SqlConsumer::new(config.clone(), Arc::new(OnceCell::new()), None, test_rt());
1134 let template = parse_query_template(&config.query, config.placeholder).unwrap();
1135
1136 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
1137 tokio::spawn(async move {
1138 loop {
1139 match timeout(Duration::from_secs(2), rx.recv()).await {
1140 Ok(Some(env)) => {
1141 if let Some(reply_tx) = env.reply_tx {
1142 let _ = reply_tx
1143 .send(Err(CamelError::ProcessorError("downstream boom".into())));
1144 }
1145 }
1146 Ok(None) => break, Err(_) => break, }
1149 }
1150 });
1151 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1152
1153 consumer
1154 .poll_database(&pool, &ctx, &template)
1155 .await
1156 .expect("consumer should swallow downstream errors when breakBatchOnConsumeFail=false");
1157
1158 let row = sqlx::query("select failed from jobs where id = 1")
1160 .fetch_one(&pool)
1161 .await
1162 .expect("row 1");
1163 let failed_1: i64 = sqlx::Row::try_get(&row, 0).expect("failed");
1164
1165 let row = sqlx::query("select failed from jobs where id = 2")
1166 .fetch_one(&pool)
1167 .await
1168 .expect("row 2");
1169 let failed_2: i64 = sqlx::Row::try_get(&row, 0).expect("failed");
1170
1171 assert_eq!(
1172 failed_1, 1,
1173 "row 1 should be marked failed via per-row onConsumeFailed"
1174 );
1175 assert_eq!(
1176 failed_2, 1,
1177 "row 2 should be marked failed via per-row onConsumeFailed"
1178 );
1179 }
1180
1181 #[tokio::test]
1182 async fn bridge_error_handler_routes_poll_errors_to_exchange_error() {
1183 let mut config = config();
1184 config.bridge_error_handler = true;
1185 let consumer = SqlConsumer::new(config, Arc::new(OnceCell::new()), None, test_rt());
1186
1187 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(4);
1188 tokio::spawn(async move {
1189 #[allow(clippy::never_loop)]
1190 loop {
1191 match timeout(Duration::from_secs(2), rx.recv()).await {
1192 Ok(Some(env)) => {
1193 assert!(env.exchange.error.is_some(), "exchange must carry error");
1194 if let Some(reply_tx) = env.reply_tx {
1195 let _ = reply_tx.send(Ok(env.exchange));
1196 }
1197 break;
1198 }
1199 Ok(None) => break, Err(_) => break, }
1202 }
1203 });
1204
1205 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1206 consumer
1207 .bridge_poll_error(&ctx, CamelError::ProcessorError("poll failed".into()))
1208 .await
1209 .expect("bridging should succeed");
1210 }
1211
1212 #[tracing_test::traced_test]
1216 #[tokio::test]
1217 async fn bridged_poll_failure_emits_warn_not_error() {
1218 let pool = sqlite_pool().await;
1219 let mut config = config();
1224 config.bridge_error_handler = true;
1225 config.query = "select * from nonexistent_table".to_string();
1227 config.resolve_defaults();
1228 let consumer = SqlConsumer::new(config.clone(), Arc::new(OnceCell::new()), None, test_rt());
1229 let template = parse_query_template(&config.query, config.placeholder).unwrap();
1230
1231 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(4);
1234 tokio::spawn(async move {
1235 loop {
1236 match timeout(Duration::from_secs(2), rx.recv()).await {
1237 Ok(Some(env)) => {
1238 if let Some(reply_tx) = env.reply_tx {
1239 let _ = reply_tx.send(Ok(env.exchange));
1240 }
1241 }
1242 Ok(None) => break, Err(_) => break, }
1245 }
1246 });
1247 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1248
1249 consumer.handle_poll_result(&pool, &ctx, &template).await;
1251
1252 assert!(
1254 !logs_contain("ERROR"),
1255 "bridged poll failure must not emit ERROR (handler owns it); check captured logs for stray ERROR lines"
1256 );
1257 assert!(
1259 logs_contain("WARN"),
1260 "bridged poll failure should emit warn! for operator visibility"
1261 );
1262 }
1263
1264 #[tracing_test::traced_test]
1274 #[tokio::test]
1275 async fn unbridged_send_and_wait_failure_emits_error_loud() {
1276 let pool = sqlite_pool().await;
1277 sqlx::query("CREATE TABLE items (id INTEGER PRIMARY KEY, name TEXT)")
1278 .execute(&pool)
1279 .await
1280 .expect("create table");
1281 sqlx::query("INSERT INTO items (id, name) VALUES (1, 'alpha')")
1282 .execute(&pool)
1283 .await
1284 .expect("seed rows");
1285
1286 let mut config = SqlEndpointConfig::from_uri(
1287 "sql:select id, name from items order by id?db_url=sqlite::memory:&outputType=StreamList&initialDelay=0&delay=1",
1288 )
1289 .unwrap();
1290 config.resolve_defaults();
1291 config.bridge_error_handler = false;
1293 let consumer = SqlConsumer::new(
1294 config.clone(),
1295 Arc::new(OnceCell::new()),
1296 None,
1297 Arc::new(RecordingRuntime::new(Arc::new(Mutex::new(Vec::new())))),
1298 );
1299 let template = parse_query_template(&config.query, config.placeholder).unwrap();
1300
1301 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
1303 tokio::spawn(async move {
1304 loop {
1305 match timeout(Duration::from_secs(2), rx.recv()).await {
1306 Ok(Some(env)) => {
1307 if let Some(reply_tx) = env.reply_tx {
1308 let _ = reply_tx.send(Err(CamelError::ProcessorError("boom".into())));
1309 }
1310 }
1311 Ok(None) => break, Err(_) => break, }
1314 }
1315 });
1316 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1317
1318 let _ = consumer.poll_database(&pool, &ctx, &template).await;
1319
1320 assert!(
1322 logs_contain("ERROR"),
1323 "unbridged send_and_wait failure MUST emit ERROR (consumer owns the signal)"
1324 );
1325 }
1326
1327 #[tracing_test::traced_test]
1330 #[tokio::test]
1331 async fn unbridged_handle_poll_result_emits_error_loud() {
1332 let pool = sqlite_pool().await;
1333 let mut config = config();
1336 config.bridge_error_handler = false;
1337 config.query = "select * from nonexistent_table".to_string();
1338 config.resolve_defaults();
1339 let consumer = SqlConsumer::new(
1340 config.clone(),
1341 Arc::new(OnceCell::new()),
1342 None,
1343 Arc::new(RecordingRuntime::new(Arc::new(Mutex::new(Vec::new())))),
1344 );
1345 let template = parse_query_template(&config.query, config.placeholder).unwrap();
1346
1347 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(4);
1349 tokio::spawn(async move {
1350 loop {
1351 match timeout(Duration::from_secs(2), rx.recv()).await {
1352 Ok(Some(env)) => {
1353 if let Some(reply_tx) = env.reply_tx {
1354 let _ = reply_tx.send(Ok(env.exchange));
1355 }
1356 }
1357 Ok(None) => break, Err(_) => break, }
1360 }
1361 });
1362 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1363
1364 consumer.handle_poll_result(&pool, &ctx, &template).await;
1365
1366 assert!(
1367 logs_contain("ERROR"),
1368 "unbridged handle_poll_result failure MUST emit ERROR (consumer owns signal)"
1369 );
1370 }
1371
1372 #[tokio::test]
1373 async fn stream_list_consumer_emits_ndjson_body() {
1374 let pool = sqlite_pool().await;
1375 sqlx::query("CREATE TABLE items (id INTEGER PRIMARY KEY, name TEXT)")
1376 .execute(&pool)
1377 .await
1378 .expect("create table");
1379 sqlx::query("INSERT INTO items (id, name) VALUES (1, 'alpha'), (2, 'beta'), (3, 'gamma')")
1380 .execute(&pool)
1381 .await
1382 .expect("seed rows");
1383
1384 let mut config = SqlEndpointConfig::from_uri(
1385 "sql:select id, name from items order by id?db_url=sqlite::memory:&outputType=StreamList&initialDelay=0&delay=1",
1386 )
1387 .unwrap();
1388 config.resolve_defaults();
1389
1390 let consumer = SqlConsumer::new(config.clone(), Arc::new(OnceCell::new()), None, test_rt());
1391 let template = parse_query_template(&config.query, config.placeholder).unwrap();
1392
1393 let (tx, rx) = mpsc::channel::<ExchangeEnvelope>(8);
1394 let (result_tx, result_rx) = tokio::sync::oneshot::channel::<Exchange>();
1395 tokio::spawn(async move {
1396 let mut rx = rx;
1397 match timeout(Duration::from_secs(2), rx.recv()).await {
1398 Ok(Some(env)) => {
1399 if let Some(reply_tx) = env.reply_tx {
1400 let _ = reply_tx.send(Ok(env.exchange.clone()));
1401 }
1402 let _ = result_tx.send(env.exchange);
1403 }
1404 Ok(None) => {}
1405 Err(_) => {}
1406 }
1407 });
1408 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1409
1410 consumer
1411 .poll_database(&pool, &ctx, &template)
1412 .await
1413 .expect("poll must succeed");
1414
1415 let exchange = result_rx.await.expect("should have received one exchange");
1416
1417 match exchange.input.body {
1418 Body::Stream(ref stream_body) => {
1419 let stream = stream_body.stream.clone();
1420 let mut guard = stream.lock().await;
1421 let stream_opt = guard.take();
1422 assert!(stream_opt.is_some(), "stream should be present");
1423
1424 use futures::StreamExt;
1425 let mut collected = Vec::new();
1426 let mut stream = stream_opt.unwrap();
1427 while let Some(chunk) = stream.next().await {
1428 let chunk = chunk.expect("stream chunk should not error");
1429 collected.extend_from_slice(&chunk);
1430 }
1431
1432 let ndjson = String::from_utf8(collected).expect("valid utf8");
1433 let lines: Vec<&str> = ndjson.trim().lines().collect();
1434 assert_eq!(lines.len(), 3, "should have 3 NDJSON lines");
1435
1436 let row0: serde_json::Value =
1437 serde_json::from_str(lines[0]).expect("valid json line 0");
1438 assert_eq!(row0["id"], 1);
1439 assert_eq!(row0["name"], "alpha");
1440
1441 let row1: serde_json::Value =
1442 serde_json::from_str(lines[1]).expect("valid json line 1");
1443 assert_eq!(row1["id"], 2);
1444 assert_eq!(row1["name"], "beta");
1445
1446 let row2: serde_json::Value =
1447 serde_json::from_str(lines[2]).expect("valid json line 2");
1448 assert_eq!(row2["id"], 3);
1449 assert_eq!(row2["name"], "gamma");
1450 }
1451 ref other => panic!("expected Body::Stream, got {:?}", other),
1452 }
1453 }
1454
1455 #[tokio::test]
1456 async fn stream_list_consumer_empty_result_set_emits_empty_stream() {
1457 let pool = sqlite_pool().await;
1458 sqlx::query("CREATE TABLE empty_items (id INTEGER PRIMARY KEY, name TEXT)")
1459 .execute(&pool)
1460 .await
1461 .expect("create table");
1462
1463 let mut config = SqlEndpointConfig::from_uri(
1464 "sql:select id, name from empty_items?db_url=sqlite::memory:&outputType=StreamList&initialDelay=0&delay=1",
1465 )
1466 .unwrap();
1467 config.resolve_defaults();
1468
1469 let consumer = SqlConsumer::new(config.clone(), Arc::new(OnceCell::new()), None, test_rt());
1470 let template = parse_query_template(&config.query, config.placeholder).unwrap();
1471
1472 let (tx, rx) = tokio::sync::oneshot::channel();
1473 let (mpsc_tx, mut mpsc_rx) = mpsc::channel::<ExchangeEnvelope>(8);
1474 tokio::spawn(async move {
1475 #[allow(clippy::never_loop)]
1476 loop {
1477 match timeout(Duration::from_secs(2), mpsc_rx.recv()).await {
1478 Ok(Some(env)) => {
1479 if let Some(reply_tx) = env.reply_tx {
1480 let _ = reply_tx.send(Ok(env.exchange.clone()));
1481 }
1482 let _ = tx.send(env.exchange);
1483 break;
1484 }
1485 Ok(None) => break, Err(_) => break, }
1488 }
1489 });
1490 let ctx = ConsumerContext::new(
1491 mpsc_tx,
1492 CancellationToken::new(),
1493 "sql-test-route".to_string(),
1494 );
1495
1496 consumer
1497 .poll_database(&pool, &ctx, &template)
1498 .await
1499 .expect("poll must succeed");
1500
1501 let exchange = rx
1502 .await
1503 .expect("StreamList should emit exchange even for empty results");
1504
1505 match exchange.input.body {
1506 Body::Stream(ref stream_body) => {
1507 let stream = stream_body.stream.clone();
1508 let mut guard = stream.lock().await;
1509 let stream_opt = guard.take();
1510
1511 use futures::StreamExt;
1512 let mut count = 0;
1513 if let Some(mut stream) = stream_opt {
1514 while let Some(chunk) = stream.next().await {
1515 let chunk = chunk.expect("stream chunk should not error");
1516 count += chunk.len();
1517 }
1518 }
1519 assert_eq!(count, 0, "empty table should produce zero stream bytes");
1520 }
1521 ref other => panic!("expected Body::Stream, got {:?}", other),
1522 }
1523 }
1524
1525 #[tokio::test]
1526 async fn break_on_empty_stops_after_drained_table() {
1527 let pool = sqlite_pool().await;
1528 seed_consumer_table(&pool).await;
1529
1530 let mut config = SqlEndpointConfig::from_uri(
1532 "sql:select id from jobs where processed = 0 order by id?db_url=sqlite::memory:&onConsume=update jobs set processed=1 where id=:#id&initialDelay=0&delay=50&breakOnEmpty=true&repeatCount=100",
1533 )
1534 .unwrap();
1535 config.resolve_defaults();
1536
1537 let pool_cell = Arc::new(OnceCell::new());
1540 pool_cell
1541 .set(Arc::new(pool.clone()))
1542 .expect("pool cell set");
1543 let mut consumer = SqlConsumer::new(config, pool_cell, None, test_rt());
1544
1545 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
1546 let route_cancel = CancellationToken::new();
1547 let received = Arc::new(std::sync::atomic::AtomicU32::new(0));
1551 let echo_cancel = route_cancel.clone();
1552 let counter = Arc::clone(&received);
1553 tokio::spawn(async move {
1555 loop {
1556 tokio::select! {
1557 _ = echo_cancel.cancelled() => break,
1558 Some(env) = rx.recv() => {
1559 counter.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1560 if let Some(reply_tx) = env.reply_tx {
1561 let _ = reply_tx.send(Ok(env.exchange));
1562 }
1563 }
1564 }
1565 }
1566 });
1567
1568 let ctx = ConsumerContext::new(tx, route_cancel.clone(), "sql-test-route".to_string());
1569
1570 let start = tokio::time::Instant::now();
1572 consumer.start(ctx).await.expect("start must succeed");
1573 let elapsed = start.elapsed();
1574
1575 let n = received.load(std::sync::atomic::Ordering::Relaxed);
1578 assert_eq!(
1579 n, 2,
1580 "productive poll should emit 2 envelopes (2 seeded rows), got {}",
1581 n
1582 );
1583
1584 let count_unprocessed: i64 =
1586 sqlx::query_scalar("select count(*) from jobs where processed = 0")
1587 .fetch_one(&pool)
1588 .await
1589 .expect("count");
1590 assert_eq!(count_unprocessed, 0);
1591
1592 assert!(
1595 elapsed < std::time::Duration::from_millis(500),
1596 "consumer should have stopped on empty poll, took {:?}",
1597 elapsed
1598 );
1599
1600 assert!(
1609 elapsed >= std::time::Duration::from_millis(100),
1610 "consumer must run a second (empty) poll before breaking on break_on_empty, \
1611 took {:?} — likely broke after the productive poll without seeing the empty one",
1612 elapsed
1613 );
1614 }
1615
1616 #[tokio::test]
1620 async fn handle_poll_result_error_does_not_signal_empty() {
1621 let pool = sqlite_pool().await;
1622
1623 let mut config = SqlEndpointConfig::from_uri(
1624 "sql:select * from this_table_does_not_exist?db_url=sqlite::memory:&breakOnEmpty=true&initialDelay=0&delay=1",
1625 )
1626 .unwrap();
1627 config.resolve_defaults();
1628
1629 let pool_cell = Arc::new(OnceCell::new());
1630 pool_cell.set(Arc::new(pool.clone())).unwrap();
1631 let consumer = SqlConsumer::new(
1634 config.clone(),
1635 pool_cell,
1636 None,
1637 Arc::new(RecordingRuntime::new(Arc::new(Mutex::new(Vec::new())))),
1638 );
1639
1640 let template = parse_query_template(&config.query, config.placeholder).unwrap();
1641
1642 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
1643 tokio::spawn(async move {
1644 loop {
1645 match timeout(Duration::from_secs(2), rx.recv()).await {
1646 Ok(Some(env)) => {
1647 if let Some(reply_tx) = env.reply_tx {
1648 let _ = reply_tx.send(Ok(env.exchange));
1649 }
1650 }
1651 Ok(None) => break, Err(_) => break, }
1654 }
1655 });
1656 let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1657
1658 let outcome = consumer.handle_poll_result(&pool, &ctx, &template).await;
1659 assert!(
1660 !outcome.was_empty,
1661 "poll error must NOT signal empty (was_empty should be false)"
1662 );
1663 }
1664
1665 #[derive(Debug, Default)]
1669 struct RecordingHealth {
1670 forced: Arc<Mutex<Vec<(String, String, String)>>>,
1671 }
1672
1673 impl HealthCheckRegistry for RecordingHealth {
1674 fn force_unhealthy_for_route(&self, route_id: &str, name: &str, reason: &str) {
1675 self.forced.lock().unwrap().push((
1676 route_id.to_string(),
1677 name.to_string(),
1678 reason.to_string(),
1679 ));
1680 }
1681 }
1682
1683 struct NoopMetricsForConsumer;
1684
1685 impl MetricsCollector for NoopMetricsForConsumer {
1686 fn record_exchange_duration(&self, _: &str, _: Duration) {}
1687 fn increment_errors(&self, _: &str, _: &str) {}
1688 fn increment_exchanges(&self, _: &str) {}
1689 fn set_queue_depth(&self, _: &str, _: usize) {}
1690 fn record_circuit_breaker_change(&self, _: &str, _: &str, _: &str) {}
1691 }
1692
1693 struct RecordingRuntimeWithHealth {
1694 health: Arc<RecordingHealth>,
1695 }
1696
1697 impl RuntimeObservability for RecordingRuntimeWithHealth {
1698 fn metrics(&self) -> Arc<dyn MetricsCollector> {
1699 Arc::new(NoopMetricsForConsumer)
1700 }
1701 fn health(&self) -> Arc<dyn HealthCheckRegistry> {
1702 self.health.clone()
1703 }
1704 }
1705
1706 #[tokio::test]
1709 async fn consumer_pool_init_failure_calls_force_unhealthy_for_route() {
1710 let health = Arc::new(RecordingHealth::default());
1711 let recorded_health = health.clone();
1712 let rt: Arc<dyn RuntimeObservability> = Arc::new(RecordingRuntimeWithHealth { health });
1713
1714 let mut config = SqlEndpointConfig::from_uri(
1715 "sql:select 1?db_url=postgres://nonexistent-host:5432/nonexistent_db&retryEnabled=false&initialDelay=0&delay=1",
1716 )
1717 .unwrap();
1718 config.max_connections = Some(1);
1719 config.min_connections = Some(0);
1720 config.idle_timeout_secs = Some(300);
1721 config.max_lifetime_secs = Some(1800);
1722
1723 let mut consumer = SqlConsumer::new(config, Arc::new(OnceCell::new()), None, rt);
1724
1725 let (tx, _rx) = mpsc::channel(8);
1726 let ctx = ConsumerContext::new(
1727 tx,
1728 CancellationToken::new(),
1729 "sql-consumer-test-route".to_string(),
1730 );
1731
1732 let result = consumer.start(ctx).await;
1733 assert!(result.is_err(), "pool init should fail with bad db_url");
1734
1735 let forced = recorded_health.forced.lock().unwrap();
1736 assert_eq!(
1737 forced.len(),
1738 1,
1739 "expected one force_unhealthy_for_route call"
1740 );
1741 assert_eq!(forced[0].0, "sql-consumer-test-route");
1742 assert_eq!(forced[0].1, "g:sql:consumer-pool-init");
1743 assert!(!forced[0].2.is_empty(), "reason should be non-empty");
1744 }
1745
1746 #[tokio::test]
1750 async fn retry_loop_invokes_operation_exactly_max_attempts_times() {
1751 use camel_component_api::NetworkRetryPolicy;
1752 use std::sync::Arc;
1753 use std::sync::atomic::{AtomicU32, Ordering};
1754
1755 let policy = NetworkRetryPolicy {
1756 max_attempts: 3,
1757 initial_delay: Duration::from_millis(1),
1758 max_delay: Duration::from_millis(1),
1759 multiplier: 1.0,
1760 ..NetworkRetryPolicy::default()
1761 };
1762
1763 let calls = Arc::new(AtomicU32::new(0));
1764 let calls_clone = Arc::clone(&calls);
1765
1766 let mut attempt: u32 = 0;
1767 let _result: Result<(), ()> = tokio::time::timeout(Duration::from_secs(30), async {
1768 loop {
1769 attempt += 1;
1770 calls_clone.fetch_add(1, Ordering::SeqCst);
1771 let op_result: Result<(), ()> = Err(());
1772 match op_result {
1773 Ok(v) => break Ok(v),
1774 Err(_) if policy.should_retry(attempt) => {
1775 let delay = policy.delay_for(attempt - 1);
1776 tokio::time::sleep(delay).await;
1777 continue;
1778 }
1779 Err(_) => break Err(()),
1780 }
1781 }
1782 })
1783 .await
1784 .expect("retry loop must finish within 30s");
1785
1786 assert_eq!(
1787 calls.load(Ordering::SeqCst),
1788 3,
1789 "max_attempts=3 must yield exactly 3 invocations"
1790 );
1791 }
1792
1793 #[tokio::test]
1796 #[tracing_test::traced_test]
1797 async fn break_on_empty_ignored_in_streamlist() {
1798 let pool = sqlite_pool().await;
1799 seed_consumer_table(&pool).await;
1800
1801 let mut config = SqlEndpointConfig::from_uri(
1805 "sql:select id from jobs?db_url=sqlite::memory:&outputType=StreamList&initialDelay=0&delay=1&breakOnEmpty=true&repeatCount=3",
1806 )
1807 .unwrap();
1808 config.resolve_defaults();
1809 assert_eq!(config.output_type, SqlOutputType::StreamList);
1810 assert!(config.break_on_empty);
1811
1812 let pool_cell = Arc::new(OnceCell::new());
1813 pool_cell
1814 .set(Arc::new(pool.clone()))
1815 .expect("pool cell set");
1816 let mut consumer = SqlConsumer::new(config, pool_cell, None, test_rt());
1817
1818 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
1819 let route_cancel = CancellationToken::new();
1820 let received = Arc::new(std::sync::atomic::AtomicU32::new(0));
1821 let echo_cancel = route_cancel.clone();
1822 let counter = Arc::clone(&received);
1823 tokio::spawn(async move {
1824 loop {
1825 tokio::select! {
1826 _ = echo_cancel.cancelled() => break,
1827 Some(env) = rx.recv() => {
1828 counter.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1829 if let Some(reply_tx) = env.reply_tx {
1830 let _ = reply_tx.send(Ok(env.exchange));
1831 }
1832 }
1833 }
1834 }
1835 });
1836
1837 let ctx = ConsumerContext::new(tx, route_cancel, "sql-test-route".to_string());
1838 consumer.start(ctx).await.expect("start must succeed");
1839
1840 assert!(
1843 logs_contain("breakOnEmpty"),
1844 "expected StreamList warn naming breakOnEmpty"
1845 );
1846
1847 assert!(
1851 received.load(std::sync::atomic::Ordering::Relaxed) >= 2,
1852 "StreamList must not break early with breakOnEmpty, got {} exchanges",
1853 received.load(std::sync::atomic::Ordering::Relaxed)
1854 );
1855 }
1856
1857 #[tokio::test]
1858 async fn break_on_empty_false_default_loops_on_empty() {
1859 let pool = sqlite_pool().await;
1860 seed_consumer_table(&pool).await;
1862 sqlx::query("delete from jobs")
1863 .execute(&pool)
1864 .await
1865 .expect("drain");
1866
1867 let mut config = SqlEndpointConfig::from_uri(
1870 "sql:select id from jobs where processed = 0?db_url=sqlite::memory:&initialDelay=0&delay=20&repeatCount=3",
1871 )
1872 .unwrap();
1873 config.resolve_defaults();
1874 assert!(!config.break_on_empty);
1875
1876 let pool_cell = Arc::new(OnceCell::new());
1877 pool_cell
1878 .set(Arc::new(pool.clone()))
1879 .expect("pool cell set");
1880 let mut consumer = SqlConsumer::new(config, pool_cell, None, test_rt());
1881
1882 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
1883 let route_cancel = CancellationToken::new();
1884 let echo_cancel = route_cancel.clone();
1885 tokio::spawn(async move {
1886 loop {
1887 tokio::select! {
1888 _ = echo_cancel.cancelled() => break,
1889 Some(env) = rx.recv() => {
1890 if let Some(reply_tx) = env.reply_tx {
1891 let _ = reply_tx.send(Ok(env.exchange));
1892 }
1893 }
1894 }
1895 }
1896 });
1897
1898 let ctx = ConsumerContext::new(tx, route_cancel, "sql-test-route".to_string());
1899 let start = tokio::time::Instant::now();
1900 consumer.start(ctx).await.expect("start must succeed");
1901 let elapsed = start.elapsed();
1902
1903 assert!(
1907 elapsed >= std::time::Duration::from_millis(55),
1908 "consumer should run all 3 polls (breakOnEmpty=false), took {:?}",
1909 elapsed
1910 );
1911 }
1912
1913 #[tokio::test]
1914 async fn break_on_empty_with_route_empty_result_set() {
1915 let pool = sqlite_pool().await;
1916 seed_consumer_table(&pool).await;
1917 sqlx::query("delete from jobs")
1918 .execute(&pool)
1919 .await
1920 .expect("drain");
1921 sqlx::query("CREATE TABLE batch_marks (n INTEGER NOT NULL DEFAULT 0)")
1927 .execute(&pool)
1928 .await
1929 .expect("create batch_marks");
1930 sqlx::query("INSERT INTO batch_marks (n) VALUES (0)")
1931 .execute(&pool)
1932 .await
1933 .expect("seed batch_marks");
1934
1935 let mut config = SqlEndpointConfig::from_uri(
1940 "sql:select id from jobs where processed = 0?db_url=sqlite::memory:&routeEmptyResultSet=true&breakOnEmpty=true&useIterator=false&onConsumeBatchComplete=update batch_marks set n = n + 1&initialDelay=0&delay=10&repeatCount=100",
1941 )
1942 .unwrap();
1943 config.resolve_defaults();
1944
1945 let pool_cell = Arc::new(OnceCell::new());
1946 pool_cell
1947 .set(Arc::new(pool.clone()))
1948 .expect("pool cell set");
1949 let mut consumer = SqlConsumer::new(config, pool_cell, None, test_rt());
1950
1951 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
1952 let route_cancel = CancellationToken::new();
1953 let received = Arc::new(std::sync::atomic::AtomicU32::new(0));
1954 let echo_cancel = route_cancel.clone();
1955 let counter = Arc::clone(&received);
1956 tokio::spawn(async move {
1957 loop {
1958 tokio::select! {
1959 _ = echo_cancel.cancelled() => break,
1960 Some(env) = rx.recv() => {
1961 counter.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1962 if let Some(reply_tx) = env.reply_tx {
1963 let _ = reply_tx.send(Ok(env.exchange));
1964 }
1965 }
1966 }
1967 }
1968 });
1969
1970 let ctx = ConsumerContext::new(tx, route_cancel, "sql-test-route".to_string());
1971 let start = tokio::time::Instant::now();
1972 consumer.start(ctx).await.expect("start must succeed");
1973
1974 let n = received.load(std::sync::atomic::Ordering::Relaxed);
1978 assert_eq!(n, 1, "expected exactly 1 empty-result exchange, got {}", n);
1979 assert!(
1980 start.elapsed() < std::time::Duration::from_millis(500),
1981 "consumer should have stopped after the first empty poll, took {:?}",
1982 start.elapsed()
1983 );
1984
1985 let batch_fires: i64 = sqlx::query_scalar("select n from batch_marks")
1990 .fetch_one(&pool)
1991 .await
1992 .expect("batch_marks n");
1993 assert_eq!(
1994 batch_fires, 1,
1995 "onConsumeBatchComplete must fire exactly once on the empty-poll \
1996 fall-through before break_on_empty, got {}",
1997 batch_fires
1998 );
1999 }
2000
2001 #[tokio::test]
2002 async fn repeat_count_zero_polls_never() {
2003 let pool = sqlite_pool().await;
2004 seed_consumer_table(&pool).await;
2005
2006 let mut config = SqlEndpointConfig::from_uri(
2008 "sql:select id from jobs?db_url=sqlite::memory:&initialDelay=0&delay=1&repeatCount=0&onConsume=update jobs set processed=1 where id=:#id",
2009 )
2010 .unwrap();
2011 config.resolve_defaults();
2012
2013 let pool_cell = Arc::new(OnceCell::new());
2016 pool_cell
2017 .set(Arc::new(pool.clone()))
2018 .expect("pool cell set");
2019 let mut consumer = SqlConsumer::new(config, pool_cell, None, test_rt());
2020
2021 let (tx, mut rx) = mpsc::channel::<ExchangeEnvelope>(8);
2022 let route_cancel = CancellationToken::new();
2023 let echo_cancel = route_cancel.clone();
2024 tokio::spawn(async move {
2025 loop {
2026 tokio::select! {
2027 _ = echo_cancel.cancelled() => break,
2028 Some(env) = rx.recv() => {
2029 if let Some(reply_tx) = env.reply_tx {
2030 let _ = reply_tx.send(Ok(env.exchange));
2031 }
2032 }
2033 }
2034 }
2035 });
2036
2037 let ctx = ConsumerContext::new(tx, route_cancel, "sql-test-route".to_string());
2038 consumer.start(ctx).await.expect("start must succeed");
2039
2040 let count_processed: i64 =
2042 sqlx::query_scalar("select count(*) from jobs where processed = 1")
2043 .fetch_one(&pool)
2044 .await
2045 .expect("count");
2046 assert_eq!(count_processed, 0, "repeatCount=0 must not poll");
2047 }
2048}