Skip to main content

camel_component_sql/
consumer.rs

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
29/// Record a post-process (b′) failure for ADR-0012 outside-contract sites in this
30/// consumer. Increments the per-label error metric AND emits an `error!` log
31/// per ADR-0012 L57 + L70-72 (the metric is the operator signal; `error!`
32/// provides loud log visibility — b′ errors are NOT absorbed by route handlers).
33///
34/// Both the metric call and the `error!` live INSIDE this helper so that
35/// `lint-log-levels`'s `has_replacement_signal` (scripts/xtask/src/main.rs)
36/// sees both literals in the helper's function body. Call sites have NO
37/// `error!` of their own.
38///
39/// Regression-tested by:
40/// - `record_post_process_failure_increments_errors_and_emits_error_log` (helper unit)
41/// - `unbridged_send_and_wait_failure_emits_error_loud` (StreamList integration path)
42fn record_post_process_failure(
43    runtime: &dyn RuntimeObservability,
44    route_id: &str,
45    label: &str,
46    error: &CamelError,
47    message: &str,
48) {
49    // allow-open-label rc-otxh (label: caller-bounded b-prime literals at all four call sites; helper co-locates metric + error! for lint-log-levels)
50    runtime.metrics().increment_errors(route_id, label);
51    // log-policy: outside-contract
52    error!(error = %error, "{message}");
53}
54
55/// Outcome of a single poll cycle. Carries whether the poll returned zero rows,
56/// threaded up to the poll loop for `break_on_empty` (without propagating errors,
57/// which are swallowed/bridged by `handle_poll_result`).
58#[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 observability for metrics and health — used by the
69    /// `record_post_process_failure` helper for ADR-0012 (b′) metric calls.
70    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    /// Poll the database for new rows and process them.
90    async fn poll_database(
91        &self,
92        pool: &AnyPool,
93        context: &ConsumerContext,
94        template: &QueryTemplate,
95    ) -> Result<PollOutcome, CamelError> {
96        // Capture route_id from ConsumerContext for ADR-0012 metrics
97        let route_id = context.route_id();
98
99        // Create an empty exchange for parameter resolution (consumer has no input)
100        let empty_exchange = Exchange::new(Message::default());
101
102        // Resolve parameters
103        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            // Process each row individually
136            for row in rows_to_process {
137                let row_json = row_to_json(&row)?;
138
139                // Create exchange with the row as JSON body
140                let mut msg = Message::new(Body::Json(row_json.clone()));
141
142                // Set individual column headers with CamelSql. prefix per Apache Camel convention
143                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                // Send and wait for processing
152                let result = context.send_and_wait(exchange).await;
153
154                // Handle post-processing (onConsume/onConsumeFailed)
155                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 downstream processing itself failed, honour break_batch_on_consume_fail
169                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            // Process all rows as a single batch
177            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            // Create exchange with array of rows
185            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            // Send and wait for result
191            let result = context.send_and_wait(exchange).await;
192
193            // SQL-021: Run per-row post-processing even in batch mode so that
194            // onConsume/onConsumeFailed queries can reference row-specific parameters
195            // (e.g. `:#id`). Each row gets its own post-processing query execution.
196            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 downstream processing itself failed, honour break_batch_on_consume_fail
212            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        // Execute on_consume_batch_complete if configured
220        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        // StreamList ignores break_on_empty, so it always returns was_empty=false
280        Ok(PollOutcome::default())
281    }
282
283    /// Handle post-processing after a row is processed (onConsume/onConsumeFailed).
284    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                // Success - execute onConsume if configured
293                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                // Failure - execute onConsumeFailed if configured
299                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    /// Execute a post-processing query with the row data as parameters.
309    async fn execute_post_query(
310        &self,
311        pool: &AnyPool,
312        query_str: &str,
313        row_json: &JsonValue,
314    ) -> Result<(), CamelError> {
315        // Parse the query template
316        let template = parse_query_template(query_str, self.config.placeholder)?;
317
318        // Create a temporary exchange with the row as body for parameter resolution
319        // Populate CamelSql.* headers so named params can reference them
320        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        // Resolve parameters
329        let prepared = resolve_params(&template, &temp_exchange, &self.config.in_separator)?;
330
331        // Build and execute the query
332        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        // Warn if 0 rows affected (the row may not have been marked correctly)
338        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    /// Handle the result of a single poll cycle, including bridging if configured.
349    /// Extracted from `run()` so tests can exercise the error-handling branch directly.
350    /// Returns `PollOutcome` so the poll loop can decide whether to break on empty.
351    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                // Swallow the poll error (do NOT propagate via ? — the loop continues).
361                if self.config.bridge_error_handler {
362                    // log-policy: handler-owned
363                    // (category b-bridged: error will be wrapped as Exchange
364                    // and flow into the route's error handler)
365                    warn!(error = %e, "SQL consumer poll failed (bridged)");
366                    if let Err(route_err) = self.bridge_poll_error(context, e).await {
367                        // (the bridge channel itself broke — route will CrashNotification per ADR-0007)
368                        // log-policy: system-broken
369                        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                // An error is NOT an empty poll — the loop must continue.
381                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        // Reject double-start
404        if self.stopped {
405            return Err(CamelError::Config(
406                "SQL consumer cannot be restarted after stop".into(),
407            ));
408        }
409
410        // Step 1: Initialize the connection pool
411        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        // SQL-014: resolve file-based query before pool init, regardless of pool source
416        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                // Catalog path: resolve shared pool from the datasource catalog
423                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                // Install all compiled-in sqlx drivers so AnyPool can resolve them.
429                // This is idempotent; safe to call multiple times.
430                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                    // log-policy: outside-contract
481                    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        // SQL-002: warn if Managed transaction mode requested
492        if self.config.transaction_mode == TransactionMode::Managed {
493            warn!("transactionManager not yet implemented; using Auto mode");
494        }
495
496        // SQL-017/SQL-018: log processing and poll strategies
497        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        // Warn if no onConsume configured
519        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        // Step 2: Parse query template once (avoid re-parsing every poll)
532        let template = parse_query_template(&self.config.query, self.config.placeholder)
533            .map_err(|e| CamelError::Config(format!("Invalid query template: {}", e)))?;
534
535        // Step 3: Initial delay before starting polling
536        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        // Step 4: Polling loop
547        //
548        // This is a POLLING LOOP with fixed cadence (delay_ms), NOT a
549        // retry loop. It polls the database until cancelled or repeat_count
550        // is reached — there is no "transient error → retry with backoff"
551        // contract at this level. retry_async / retry_async_cancelable do
552        // not apply because they are designed for bounded retry, not
553        // repeated polling with uniform delay.
554        //
555        // The pool-connect retry at startup (Step 1) was migrated to
556        // retry_async in rc-d2r. The per-poll error handling (poll_database
557        // failures) is an error-bridge pattern, not a retry loop.
558        //
559        // See camel-redis/src/consumer.rs:325 for a similar polling-loop
560        // justification.
561        let mut poll_count: u32 = 0;
562        loop {
563            // SQL-015: check repeat_count limit
564            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        // Double-stop is safe — no-op after first stop
595        if self.stopped {
596            debug!("SQL consumer stop called on already-stopped consumer");
597            return Ok(());
598        }
599
600        // Close the connection pool if it was initialized
601        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        // Sequential is correct for SQL consumers: concurrent polls would fetch
614        // duplicate rows. The design doc mentioned SharedState (which doesn't exist
615        // in this runtime) — Sequential is the correct equivalent.
616        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    // -----------------------------------------------------------------------
641    // Recording metrics collector for testing increment_errors calls
642    // -----------------------------------------------------------------------
643
644    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    /// Regression test for ADR-0012: the record_post_process_failure helper
683    /// must increment the error metric with the correct route_id and label,
684    /// AND emit error! via tracing.
685    #[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        // Directly invoke the helper
693        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        // Verify MetricsCollector::increment_errors was called
702        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        // Verify error! was emitted
709        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        // Per-attempt deadline (lintwiden D4.3): a wedged pool connect
719        // must fail the helper instead of parking it.
720        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, // closed: drain complete
791                    Err(_) => break,   // stalled: drainer ends
792                }
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, // closed: drain complete
843                    Err(_) => break,   // stalled: drainer ends
844                }
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, // closed: drain complete
895                    Err(_) => break,   // stalled: drainer ends
896                }
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    // --- Phase B hardening tests ---
924
925    // SQL-001: Direct consumer construction without resolve_defaults does not panic.
926    // The consumer defensively calls resolve_defaults() during pool init, so the pool
927    // fields get resolved. This test verifies no panic occurs.
928    #[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        // Deliberately NOT calling resolve_defaults() — pool fields remain None
935        assert!(config.max_connections.is_none());
936
937        let mut consumer = SqlConsumer::new(
938            config,
939            Arc::new(OnceCell::new()),
940            None,
941            // Noop runtime: pool-init retries now record per-attempt
942            // telemetry by design, which the panic runtime would reject.
943            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, // closed: drain complete
955                    Err(_) => break,   // stalled: drainer ends
956                }
957            }
958        });
959        let token = CancellationToken::new();
960        let ctx = ConsumerContext::new(tx, token.clone(), "sql-test-route".to_string());
961
962        // Spawn the consumer and cancel it quickly — it should not panic
963        let consumer_handle = tokio::spawn(async move { consumer.start(ctx).await });
964
965        // Cancel after a short delay
966        tokio::time::sleep(Duration::from_millis(50)).await;
967        token.cancel();
968
969        let result = consumer_handle.await.expect("task should not panic");
970        // Should complete without panic (may be Ok or Err depending on timing)
971        let _ = result;
972    }
973
974    // SQL-008: stop() closes the pool
975    #[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        // After stop, the pool should be closed
993        assert!(
994            pool.is_closed(),
995            "Pool should be closed after consumer.stop()"
996        );
997    }
998
999    // SQL-008: double-stop is safe
1000    #[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    // SQL-008: start after stop is rejected
1021    #[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, // closed: drain complete
1046                    Err(_) => break,   // stalled: drainer ends
1047                }
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    // SQL-021: batch mode per-row post-processing
1063    #[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, // closed: drain complete
1087                    Err(_) => break,   // stalled: drainer ends
1088                }
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        // SQL-021: Each row should have been processed individually via onConsume
1099        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    // SQL-021: batch mode per-row onConsumeFailed when downstream fails
1122    #[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, // closed: drain complete
1147                    Err(_) => break,   // stalled: drainer ends
1148                }
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        // SQL-021: Each row should have onConsumeFailed executed individually
1159        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, // closed: drain complete
1200                    Err(_) => break,   // stalled: drainer ends
1201                }
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    /// Regression for ADR-0012: when bridge_error_handler=true, the poll
1213    /// failure must NOT emit error! (the route's error handler owns ERROR
1214    /// for bridged failures). Was previously duplicated at line 429 + 431.
1215    #[tracing_test::traced_test]
1216    #[tokio::test]
1217    async fn bridged_poll_failure_emits_warn_not_error() {
1218        let pool = sqlite_pool().await;
1219        // Do NOT create any table — the query against a non-existent
1220        // table will fail at fetch_all, returning Err BEFORE any
1221        // downstream send (so lines 103/205 are never reached).
1222
1223        let mut config = config();
1224        config.bridge_error_handler = true;
1225        // Query a non-existent table to trigger a query-failure poll error.
1226        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        // Healthy downstream — replies Ok so bridge_poll_error succeeds
1232        // and does NOT emit its own error!.
1233        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, // closed: drain complete
1243                    Err(_) => break,   // stalled: drainer ends
1244                }
1245            }
1246        });
1247        let ctx = ConsumerContext::new(tx, CancellationToken::new(), "sql-test-route".to_string());
1248
1249        // Drive poll — fetch_all will fail because the table is missing.
1250        consumer.handle_poll_result(&pool, &ctx, &template).await;
1251
1252        // The bridged path must NOT emit ERROR (handler owns it).
1253        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        // Sanity: warn! was emitted so the failure is still visible.
1258        assert!(
1259            logs_contain("WARN"),
1260            "bridged poll failure should emit warn! for operator visibility"
1261        );
1262    }
1263
1264    /// Regression for ADR-0012 "b-bridged discriminator": when
1265    /// send_and_wait returns Err on a NORMAL-DATA send (i.e., not a
1266    /// deliberate bridge_poll_error handoff), the route handler did NOT
1267    /// absorb the failure (consumer.rs:77-91 contract; error_handler.rs
1268    /// returns Ok in every branch). The consumer's error! is the only
1269    /// ERROR signal for the unhandled failure and MUST stay at error!.
1270    ///
1271    /// Protects consumer.rs:205 (StreamList downstream send) and any
1272    /// future site that uses send_and_wait on a non-bridge path.
1273    #[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        // Explicitly non-bridged: normal-data send path.
1292        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        // Downstream that returns Err — simulates unhandled route failure.
1302        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, // closed: drain complete
1312                    Err(_) => break,   // stalled: drainer ends
1313                }
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        // The unbridged path MUST emit ERROR — consumer owns the signal.
1321        assert!(
1322            logs_contain("ERROR"),
1323            "unbridged send_and_wait failure MUST emit ERROR (consumer owns the signal)"
1324        );
1325    }
1326
1327    /// Regression for ADR-0012: when bridge_error_handler=false, the unbridged
1328    /// branch of handle_poll_result MUST emit ERROR for unhandled poll failure.
1329    #[tracing_test::traced_test]
1330    #[tokio::test]
1331    async fn unbridged_handle_poll_result_emits_error_loud() {
1332        let pool = sqlite_pool().await;
1333        // Do NOT create any table — fetch_all will fail in poll_database.
1334
1335        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        // Healthy downstream task; should not be reached for this poll-failure path.
1348        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, // closed: drain complete
1358                    Err(_) => break,   // stalled: drainer ends
1359                }
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, // closed: drain complete
1486                    Err(_) => break,   // stalled: drainer ends
1487                }
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        // onConsume marks rows processed=1; query selects only processed=0 → drains in one poll.
1531        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        // Inject the seeded pool — start() otherwise self-initializes a disjoint
1538        // sqlite::memory: DB (per-connection private).
1539        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        // Count envelopes so we can pin the productive poll's row processing.
1548        // With use_iterator=true (default) and 2 seeded rows, the productive poll
1549        // must emit exactly 2 envelopes before the empty poll triggers the break.
1550        let received = Arc::new(std::sync::atomic::AtomicU32::new(0));
1551        let echo_cancel = route_cancel.clone();
1552        let counter = Arc::clone(&received);
1553        // Echo replies so the poll completes.
1554        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        // start() runs the poll loop to completion (returns when the loop breaks).
1571        let start = tokio::time::Instant::now();
1572        consumer.start(ctx).await.expect("start must succeed");
1573        let elapsed = start.elapsed();
1574
1575        // Productive poll must have emitted exactly 2 envelopes (one per row in
1576        // the seeded table). Locks "no rows skipped, no rows lost".
1577        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        // Both rows drained (processed=1) by the productive poll.
1585        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        // Upper bound: with delay=10ms and breakOnEmpty, the consumer should stop
1593        // well under the repeatCount=100 ceiling (which would take ~1s).
1594        assert!(
1595            elapsed < std::time::Duration::from_millis(500),
1596            "consumer should have stopped on empty poll, took {:?}",
1597            elapsed
1598        );
1599
1600        // Lower bound — LOCKS the [productive_poll, empty_poll] ordering.
1601        // With delay=50ms:
1602        //   - Correct: at least 2 full delays elapse (1st sleep + 2nd sleep) ≈ 100ms+,
1603        //     because the loop MUST run a second (empty) poll before breaking.
1604        //   - Buggy (break after the productive poll): only the 1st delay elapses
1605        //     ≈ 50ms+processing. 100ms threshold sits safely between the two and
1606        //     would FAIL if the consumer incorrectly set was_empty=true on the
1607        //     productive poll and broke early. Wide margin survives slow CI boxes.
1608        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    /// Regression: `handle_poll_result` with `break_on_empty=true` must NOT
1617    /// signal `was_empty: true` when `poll_database` returns an error — the
1618    /// loop should continue past the error.
1619    #[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        // Use RecordingRuntime (not test_rt/PanicRuntime) because
1632        // handle_poll_result calls record_post_process_failure → metrics().
1633        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, // closed: drain complete
1652                    Err(_) => break,   // stalled: drainer ends
1653                }
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    // ── ADR-0012 (g) regression tests ──────────────────────────────────
1666
1667    /// Fixture: captures `force_unhealthy_for_route` calls.
1668    #[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    /// Regression: consumer pool init failure calls force_unhealthy_for_route
1707    /// with correct route_id + name "g:sql:consumer-pool-init" + non-empty reason.
1708    #[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    /// Regression: max_attempts=N → exactly N invocations (caught OpenSearch off-by-one 1f5c4c2a).
1747    /// Replicates the exact retry loop from SqlConsumer::start() (consumer.rs:343-367):
1748    ///   attempt starts at 0, incremented at top, should_retry(attempt), delay_for(attempt-1)
1749    #[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    // ── break_on_empty edge-case tests ──────────────────────────────────
1794
1795    #[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        // StreamList + breakOnEmpty=true: warn must fire, breakOnEmpty ignored
1802        // (no break on empty — rows flow lazily). repeatCount=3 to distinguish
1803        // "ran 3 polls" from "broke on poll 1".
1804        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        // The startup warn must name breakOnEmpty (FAILS before Step 2 — current
1841        // warn message omits it).
1842        assert!(
1843            logs_contain("breakOnEmpty"),
1844            "expected StreamList warn naming breakOnEmpty"
1845        );
1846
1847        // Counter must prove the stream ran multiple polls (no early break).
1848        // repeatCount=3 with delay=1ms means 3 polls; even accounting for race
1849        // the counter must be >=2 if no early break.
1850        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        // Empty table (seed then drain) — every poll returns 0 rows.
1861        seed_consumer_table(&pool).await;
1862        sqlx::query("delete from jobs")
1863            .execute(&pool)
1864            .await
1865            .expect("drain");
1866
1867        // breakOnEmpty NOT set (default false); repeatCount=3 so the loop must run
1868        // all 3 polls (NOT break on the first empty poll). delay=20ms each.
1869        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        // Regression guard: with breakOnEmpty=false + repeatCount=3 + delay=20ms,
1904        // the loop runs all 3 polls (~60ms). If breakOnEmpty were mis-defaulted to
1905        // true, it would break on poll 1 (~20ms). Assert the full window ran.
1906        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        // Side-effect table for onConsumeBatchComplete: each fire of the
1922        // batch-complete callback increments `n` exactly once. With
1923        // breakOnEmpty=true on an empty table, the spec pins that the
1924        // empty-poll fall-through fires the callback exactly once before
1925        // termination.
1926        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        // Empty table + routeEmptyResultSet=true: empty polls fall through to the
1936        // batch path (emit empty result) instead of early-returning. breakOnEmpty=true
1937        // must break AFTER that batch processing → exactly one downstream exchange.
1938        // onConsumeBatchComplete is wired so we can observe the empty-poll fall-through.
1939        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        // Exactly one downstream exchange emitted (the first empty poll's empty result),
1975        // then breakOnEmpty stopped the loop. Must NOT be 0 (route_empty honored) and
1976        // must NOT be repeatCount=100 (break_on_empty honored).
1977        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        // onConsumeBatchComplete must fire exactly once on the empty-poll
1986        // fall-through before the loop breaks. A buggy version that broke
1987        // before invoking the batch callback would leave n=0; a version
1988        // that looped through repeatCount=100 would leave n=100.
1989        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        // repeatCount=0 → consumer exits before the first poll (guard at loop top).
2007        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        // Inject the seeded pool — start() otherwise self-initializes a disjoint
2014        // sqlite::memory: DB (per-connection private). Pattern: consumer.rs:967-970.
2015        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        // No poll ran → rows are untouched (processed=0).
2041        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}