use std::collections::VecDeque;
use std::fmt;
use std::sync::Arc;
use std::time::{Duration, Instant};
use arrow::datatypes::SchemaRef;
use arrow::record_batch::RecordBatch;
use datafusion::common::Statistics;
use datafusion::error::{DataFusionError, Result as DFResult};
use datafusion::execution::TaskContext;
use datafusion::physical_expr::EquivalenceProperties;
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
use datafusion::physical_plan::{
DisplayAs, DisplayFormatType, Distribution, ExecutionPlan, Partitioning, PlanProperties,
SendableRecordBatchStream,
};
use futures::stream;
use serde_json::Value;
use super::cache::{ScanCache, ScanKeyParts, scan_cache_key, schema_fingerprint};
use super::client::OpenConnectorClient;
use super::error::OpenConnectorError;
use super::json_to_arrow::RowConverter;
use super::row_path::RowPath;
pub use skardi_source_pack::scan::{ActionScan, ScanBounds, ScanTarget};
pub struct OpenConnectorExec {
client: Arc<OpenConnectorClient>,
cache: Option<Arc<ScanCache>>,
gateway: String,
binding: Option<String>,
connection_alias: Option<String>,
target: ScanTarget,
converter: Arc<RowConverter>,
row_path: RowPath,
resource: Value,
filter_inputs: Vec<(String, Value)>,
projection: Option<Vec<usize>>,
limit: Option<usize>,
max_pages: u32,
max_rows: u64,
scan_timeout: Duration,
schema: SchemaRef,
properties: PlanProperties,
}
impl fmt::Debug for OpenConnectorExec {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("OpenConnectorExec")
.field("table", &self.target.table_id)
.field("action", &self.target.action_id)
.field("limit", &self.limit)
.field("projection", &self.projection)
.finish()
}
}
impl OpenConnectorExec {
#[allow(clippy::too_many_arguments)]
pub fn new(
client: Arc<OpenConnectorClient>,
cache: Option<Arc<ScanCache>>,
gateway: String,
binding: Option<String>,
connection_alias: Option<String>,
target: ScanTarget,
converter: Arc<RowConverter>,
row_path: RowPath,
resource: Value,
filter_inputs: Vec<(String, Value)>,
projection: Option<Vec<usize>>,
limit: Option<usize>,
max_pages: u32,
max_rows: u64,
scan_timeout: Duration,
) -> DFResult<Self> {
let full_schema = Arc::clone(converter.schema());
let schema = match &projection {
Some(indices) => Arc::new(full_schema.project(indices)?),
None => full_schema,
};
let properties = PlanProperties::new(
EquivalenceProperties::new(Arc::clone(&schema)),
Partitioning::UnknownPartitioning(1),
EmissionType::Incremental,
Boundedness::Bounded,
);
Ok(Self {
client,
cache,
gateway,
binding,
connection_alias,
target,
converter,
row_path,
resource,
filter_inputs,
projection,
limit,
max_pages,
max_rows,
scan_timeout,
schema,
properties,
})
}
}
impl DisplayAs for OpenConnectorExec {
fn fmt_as(&self, _t: DisplayFormatType, f: &mut fmt::Formatter) -> fmt::Result {
write!(
f,
"OpenConnectorExec: table={} action={} limit={:?}",
self.target.table_id, self.target.action_id, self.limit
)
}
}
impl ExecutionPlan for OpenConnectorExec {
fn name(&self) -> &str {
"OpenConnectorExec"
}
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn properties(&self) -> &PlanProperties {
&self.properties
}
fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
vec![]
}
fn with_new_children(
self: Arc<Self>,
children: Vec<Arc<dyn ExecutionPlan>>,
) -> DFResult<Arc<dyn ExecutionPlan>> {
if children.is_empty() {
Ok(self)
} else {
Err(DataFusionError::Internal(
"OpenConnectorExec is a leaf plan and takes no children".to_string(),
))
}
}
fn required_input_distribution(&self) -> Vec<Distribution> {
vec![]
}
fn statistics(&self) -> DFResult<Statistics> {
Ok(Statistics::new_unknown(&self.schema))
}
fn execute(
&self,
partition: usize,
_context: Arc<TaskContext>,
) -> DFResult<SendableRecordBatchStream> {
if partition != 0 {
return Err(DataFusionError::Internal(format!(
"OpenConnectorExec has a single partition, got partition {partition}"
)));
}
let state = ScanState::new(self).map_err(|e| DataFusionError::External(Box::new(e)))?;
let stream = stream::try_unfold(state, |mut state| async move {
match state.next_page().await {
Ok(Some(batch)) => Ok(Some((batch, state))),
Ok(None) => Ok(None),
Err(e) => {
tracing::warn!(
gateway = %state.gateway,
binding = state.binding.as_deref().unwrap_or("<udtf>"),
table = %state.scan.target().table_id,
action = %state.scan.target().action_id,
pages_fetched = state.scan.pages_fetched(),
error = %e,
"Open Connector scan failed"
);
Err(DataFusionError::External(Box::new(e)))
}
}
});
Ok(Box::pin(RecordBatchStreamAdapter::new(
Arc::clone(&self.schema),
stream,
)))
}
}
struct ScanState {
scan: ActionScan,
cache: Option<Arc<ScanCache>>,
cache_key: String,
gateway: String,
binding: Option<String>,
converter: Arc<RowConverter>,
projection: Option<Vec<usize>>,
limit_remaining: Option<usize>,
max_rows: u64,
rows_emitted: u64,
replay: VecDeque<RecordBatch>,
fetched: Vec<RecordBatch>,
done: bool,
cache_stored: bool,
started: Instant,
cache_hit: bool,
rows_returned: u64,
completion_logged: bool,
}
impl ScanState {
fn new(exec: &OpenConnectorExec) -> Result<Self, OpenConnectorError> {
let converter_schema = exec.converter.schema();
let projection_names: Vec<String> = match &exec.projection {
Some(indices) => indices
.iter()
.map(|&i| converter_schema.field(i).name().clone())
.collect(),
None => converter_schema
.fields()
.iter()
.map(|f| f.name().clone())
.collect(),
};
let cache_key = scan_cache_key(&ScanKeyParts {
gateway: &exec.gateway,
connection_alias: exec.connection_alias.as_deref(),
action_id: &exec.target.action_id,
source_pack_version: exec.target.source_pack_version,
resource: &exec.resource,
filter_inputs: &exec.filter_inputs,
projection: &projection_names,
limit: exec.limit,
schema_fingerprint: &schema_fingerprint(&exec.schema),
row_shape: exec.target.row_shape,
});
let cached = exec
.cache
.as_ref()
.and_then(|cache| cache.get(&cache_key))
.map(VecDeque::from);
let done = cached.is_some();
let replay = cached.unwrap_or_default();
let scan = ActionScan::new(
exec.client.clone(),
exec.target.clone(),
exec.connection_alias.clone(),
exec.resource.clone(),
exec.filter_inputs.clone(),
exec.row_path.as_str(),
ScanBounds {
max_pages: exec.max_pages,
timeout: exec.scan_timeout,
},
)?;
Ok(Self {
scan,
cache: exec.cache.clone(),
cache_key,
gateway: exec.gateway.clone(),
binding: exec.binding.clone(),
converter: exec.converter.clone(),
projection: exec.projection.clone(),
limit_remaining: exec.limit,
max_rows: exec.max_rows,
rows_emitted: 0,
replay,
fetched: Vec::new(),
done,
cache_stored: false,
started: Instant::now(),
cache_hit: done,
rows_returned: 0,
completion_logged: false,
})
}
fn store_cache(&mut self) {
if self.cache_stored {
return;
}
if let Some(cache) = &self.cache {
cache.put(self.cache_key.clone(), self.fetched.clone());
}
self.cache_stored = true;
}
fn timeout_error(&self) -> OpenConnectorError {
self.scan.timeout_error()
}
fn log_completion(&mut self) {
if self.completion_logged {
return;
}
self.completion_logged = true;
tracing::info!(
gateway = %self.gateway,
binding = self.binding.as_deref().unwrap_or("<udtf>"),
table = %self.scan.target().table_id,
action = %self.scan.target().action_id,
cache_hit = self.cache_hit,
pages = self.scan.pages_fetched(),
rows = self.rows_returned,
duration_ms = self.started.elapsed().as_millis() as u64,
"Open Connector scan completed"
);
}
async fn next_page(&mut self) -> Result<Option<RecordBatch>, OpenConnectorError> {
if let Some(batch) = self.replay.pop_front() {
self.rows_returned += batch.num_rows() as u64;
if self.replay.is_empty() {
self.log_completion();
}
return Ok(Some(batch));
}
if self.done || self.limit_remaining == Some(0) {
self.done = true;
self.log_completion();
return Ok(None);
}
let page = self.scan.page();
let envelope = self.scan.fetch_page().await?;
let rows = self.scan.rows(&envelope)?;
let batch = self.converter.convert(rows, page)?;
if Instant::now() >= self.scan.deadline() {
return Err(self.timeout_error());
}
let batch = match &self.projection {
Some(indices) => {
batch
.project(indices)
.map_err(|e| OpenConnectorError::ConversionFailed {
path: self.scan.row_path().as_str().to_string(),
column: "<projection>".to_string(),
page,
row: 0,
expected: "a valid projection".to_string(),
found: e.to_string(),
})?
}
None => batch,
};
let batch = match &mut self.limit_remaining {
Some(remaining) => {
let take = (*remaining).min(batch.num_rows());
let sliced = batch.slice(0, take);
*remaining -= take;
sliced
}
None => batch,
};
self.rows_emitted += batch.num_rows() as u64;
if self.rows_emitted > self.max_rows {
return Err(OpenConnectorError::ScanBoundsExceeded {
table: self.scan.target().table_id.to_string(),
bound: "max_rows",
limit: self.max_rows,
});
}
if batch.num_rows() > 0 {
self.fetched.push(batch.clone());
}
if self.limit_remaining == Some(0) {
self.done = true;
self.store_cache();
}
if !self.done {
let more = self.scan.advance(&envelope, rows)?;
if !more {
self.done = true;
self.store_cache();
}
}
if batch.num_rows() == 0 && self.done {
self.log_completion();
return Ok(None);
}
self.rows_returned += batch.num_rows() as u64;
if self.done {
self.log_completion();
}
Ok(Some(batch))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::sources::providers::open_connector::json_to_arrow::{FieldMapping, FieldType};
use crate::sources::providers::open_connector::packs::mock;
use crate::sources::providers::open_connector::pagination::{
AbsentCursor, CursorContinuation, PaginationStrategy,
};
use crate::sources::providers::open_connector::source_pack::{
FixedValue, RowShape, SourcePackTable,
};
use crate::sources::providers::open_connector::testutil::{
CapturedEvent, MockGateway, MockResponse, RecordedRequest, capture_events, envelope_ok,
};
use futures::StreamExt;
use serde_json::json;
use std::sync::Mutex;
fn build_exec(
client: Arc<OpenConnectorClient>,
cache: Option<Arc<ScanCache>>,
limit: Option<usize>,
source_pack_version: u32,
) -> OpenConnectorExec {
let table = &mock::pack().expect("embedded asset parses").tables[0];
OpenConnectorExec::new(
client,
cache,
"saas".to_string(),
Some("ws".to_string()),
None,
ScanTarget::from_pack_table(table, source_pack_version),
Arc::new(RowConverter::new(table.fields).expect("converter")),
RowPath::parse(table.row_path).expect("row path"),
json!({}),
vec![],
None,
limit,
10,
1000,
Duration::from_secs(30),
)
.expect("build exec")
}
fn exec_with_version(source_pack_version: u32) -> OpenConnectorExec {
let client = Arc::new(
OpenConnectorClient::new("http://127.0.0.1:1", "t", Duration::from_secs(1))
.expect("build client"),
);
build_exec(client, None, None, source_pack_version)
}
fn items_handler(req: &RecordedRequest, total: usize) -> MockResponse {
if req.method == "POST" && req.path == "/v1/actions/mock.list_items" {
let body: serde_json::Value = serde_json::from_str(&req.body).unwrap_or_default();
let page = body
.get("input")
.and_then(|input| input.get("page"))
.and_then(serde_json::Value::as_u64)
.unwrap_or(1) as usize;
let items: Vec<_> = (1..=total)
.map(|id| json!({"id": id, "name": format!("item-{id}")}))
.skip((page - 1) * 2)
.take(2)
.collect();
return MockResponse::ok(&envelope_ok(&json!({"items": items}).to_string()));
}
MockResponse::new(404, "{}")
}
async fn online_client(gateway: &MockGateway) -> Arc<OpenConnectorClient> {
Arc::new(
OpenConnectorClient::new(&gateway.url, "t", Duration::from_secs(5))
.expect("build client"),
)
}
fn split_action_table() -> &'static SourcePackTable {
Box::leak(Box::new(SourcePackTable {
id: "split.entries",
action_id: "split.list",
row_path: "$.entries",
row_shape: RowShape::Array,
fields: &[FieldMapping {
name: "id",
path: "id",
field_type: FieldType::UInt64,
nullable: false,
}],
pagination: PaginationStrategy::Cursor {
cursor_param: "cursor",
next_cursor_path: "$.cursor",
page_size_param: Some("limit"),
page_size: 2,
has_more_path: Some("$.hasMore"),
absent_cursor: AbsentCursor::EndsTheScan,
},
required_resources: &[],
optional_resources: &["path"],
exclusive_resources: &[],
fixed_inputs: &[("recursive", FixedValue::Bool(true))],
filters: &[],
error_path: None,
expected_fingerprint: None,
continuation: Some(CursorContinuation {
action_id: "split.list_continue",
expected_fingerprint: "unused-in-exec",
cursor_only: true,
}),
}))
}
fn object_row_table() -> &'static SourcePackTable {
Box::leak(Box::new(SourcePackTable {
id: "docs.document_content",
action_id: "docs.get_document_content",
row_path: "$",
row_shape: RowShape::Object,
fields: &[
FieldMapping {
name: "document_id",
path: "documentId",
field_type: FieldType::Utf8,
nullable: false,
},
FieldMapping {
name: "content",
path: "content",
field_type: FieldType::Utf8,
nullable: true,
},
],
pagination: PaginationStrategy::SinglePage {
next_cursor_path: None,
},
required_resources: &["documentId"],
optional_resources: &[],
exclusive_resources: &[],
fixed_inputs: &[],
filters: &[],
error_path: None,
expected_fingerprint: None,
continuation: None,
}))
}
#[tokio::test]
async fn object_row_response_becomes_exactly_one_row() {
let gateway = MockGateway::start(|req: &RecordedRequest| {
if req.path == "/v1/actions/docs.get_document_content" {
return MockResponse::ok(&envelope_ok(
&json!({"documentId": "doc-1", "content": "hello world"}).to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let table = object_row_table();
let exec = OpenConnectorExec::new(
online_client(&gateway).await,
None,
"saas".to_string(),
Some("ws".to_string()),
None,
ScanTarget::from_pack_table(table, 1),
Arc::new(RowConverter::new(table.fields).expect("converter")),
RowPath::parse_object_root(table.row_path).expect("object root"),
json!({"documentId": "doc-1"}),
vec![],
None,
None,
10,
1000,
Duration::from_secs(30),
)
.expect("build exec");
let mut state = ScanState::new(&exec).expect("state");
let first = state
.next_page()
.await
.expect("page")
.expect("one batch is emitted");
assert_eq!(
first.num_rows(),
1,
"the response object is exactly one row"
);
let content = first
.column(1)
.as_any()
.downcast_ref::<datafusion::arrow::array::StringArray>()
.expect("utf8 column");
assert_eq!(content.value(0), "hello world", "values survive conversion");
assert!(
state.next_page().await.expect("second page").is_none(),
"a single-page object table stops after one call"
);
}
#[tokio::test]
async fn object_row_table_fails_loudly_when_the_response_is_not_an_object() {
let gateway = MockGateway::start(|req: &RecordedRequest| {
if req.path == "/v1/actions/docs.get_document_content" {
return MockResponse::ok(&envelope_ok(&json!(null).to_string()));
}
MockResponse::new(404, "{}")
})
.await;
let table = object_row_table();
let exec = OpenConnectorExec::new(
online_client(&gateway).await,
None,
"saas".to_string(),
Some("ws".to_string()),
None,
ScanTarget::from_pack_table(table, 1),
Arc::new(RowConverter::new(table.fields).expect("converter")),
RowPath::parse_object_root(table.row_path).expect("object root"),
json!({"documentId": "doc-1"}),
vec![],
None,
None,
10,
1000,
Duration::from_secs(30),
)
.expect("build exec");
let mut state = ScanState::new(&exec).expect("state");
let err = state.next_page().await.expect_err("null must fail loudly");
assert!(
err.to_string().contains("expected an object"),
"unexpected error: {err}"
);
}
#[tokio::test]
async fn continuation_pages_target_the_continue_action_with_only_a_cursor() {
let requests: Arc<Mutex<Vec<(String, serde_json::Value)>>> =
Arc::new(Mutex::new(Vec::new()));
let recorder = Arc::clone(&requests);
let gateway = MockGateway::start(move |req: &RecordedRequest| {
let body: serde_json::Value = serde_json::from_str(&req.body).unwrap_or_default();
let input = body.get("input").cloned().unwrap_or(json!({}));
recorder
.lock()
.unwrap_or_else(|p| p.into_inner())
.push((req.path.clone(), input));
match req.path.as_str() {
"/v1/actions/split.list" => MockResponse::ok(&envelope_ok(
&json!({"entries": [{"id": 1}, {"id": 2}], "cursor": "c2", "hasMore": true})
.to_string(),
)),
"/v1/actions/split.list_continue" => MockResponse::ok(&envelope_ok(
&json!({"entries": [{"id": 3}], "cursor": "c3", "hasMore": false}).to_string(),
)),
_ => MockResponse::new(404, "{}"),
}
})
.await;
let table = split_action_table();
let exec = OpenConnectorExec::new(
online_client(&gateway).await,
None,
"saas".to_string(),
Some("ws".to_string()),
None,
ScanTarget::from_pack_table(table, 1),
Arc::new(RowConverter::new(table.fields).expect("converter")),
RowPath::parse(table.row_path).expect("row path"),
json!({"path": "/docs"}),
vec![],
None,
None,
10,
1000,
Duration::from_secs(30),
)
.expect("build exec");
let mut state = ScanState::new(&exec).expect("state");
let mut rows = 0;
while let Some(batch) = state.next_page().await.expect("page") {
rows += batch.num_rows();
}
assert_eq!(rows, 3, "both pages are scanned");
let requests = requests.lock().unwrap_or_else(|p| p.into_inner()).clone();
assert_eq!(requests.len(), 2, "exactly two gateway calls");
let (path, input) = &requests[0];
assert_eq!(path, "/v1/actions/split.list", "page one opens the listing");
assert_eq!(input.get("path"), Some(&json!("/docs")), "resource");
assert_eq!(input.get("recursive"), Some(&json!(true)), "fixed input");
assert_eq!(input.get("limit"), Some(&json!(2)), "page size");
assert!(
input.get("cursor").is_none(),
"no cursor exists on the first page"
);
let (path, input) = &requests[1];
assert_eq!(
path, "/v1/actions/split.list_continue",
"page two targets the continue action, not the opening one"
);
assert_eq!(
input,
&json!({"cursor": "c2"}),
"the continue action's schema accepts the cursor and nothing else"
);
}
#[tokio::test]
async fn limit_satisfied_scan_logs_completion_with_the_final_batch() {
let gateway = MockGateway::start(|req| items_handler(req, 5)).await;
let exec = build_exec(online_client(&gateway).await, None, Some(1), 1);
let mut state = ScanState::new(&exec).expect("state");
let batch = state.next_page().await.expect("page").expect("batch");
assert_eq!(batch.num_rows(), 1);
assert!(
state.completion_logged,
"LIMIT-satisfied scan must log completion on its final batch, \
not on a poll that never comes"
);
assert_eq!(
state.rows_returned, 1,
"the event must count the final batch"
);
}
#[tokio::test]
async fn exhausted_scan_logs_completion_with_a_nonempty_final_page() {
let gateway = MockGateway::start(|req| items_handler(req, 3)).await;
let exec = build_exec(online_client(&gateway).await, None, None, 1);
let mut state = ScanState::new(&exec).expect("state");
state.next_page().await.expect("page 1").expect("full page");
assert!(
!state.completion_logged,
"scan is still running after page 1"
);
state
.next_page()
.await
.expect("page 2")
.expect("short page");
assert!(
state.completion_logged,
"exhaustion must log with the final batch"
);
assert_eq!(state.rows_returned, 3);
}
#[tokio::test]
async fn cached_replay_logs_completion_with_the_last_replayed_batch() {
let gateway = MockGateway::start(|req| items_handler(req, 3)).await;
let cache = Arc::new(ScanCache::new(Duration::from_secs(60), usize::MAX));
let exec = build_exec(online_client(&gateway).await, Some(cache), None, 1);
let mut state = ScanState::new(&exec).expect("live state");
while state.next_page().await.expect("live page").is_some() {}
let mut state = ScanState::new(&exec).expect("replay state");
assert!(state.cache_hit, "second identical scan replays from cache");
let batches = state.replay.len();
assert!(batches > 0);
for index in 1..=batches {
state
.next_page()
.await
.expect("replayed page")
.expect("batch");
assert_eq!(
state.completion_logged,
index == batches,
"completion logs exactly when the replay queue drains"
);
}
assert_eq!(state.rows_returned, 3);
}
const COMPLETED: &str = "Open Connector scan completed";
const FAILED: &str = "Open Connector scan failed";
fn events_with_message(
events: &Arc<Mutex<Vec<CapturedEvent>>>,
message: &str,
) -> Vec<CapturedEvent> {
events
.lock()
.unwrap_or_else(|p| p.into_inner())
.iter()
.filter(|event| event.message == message)
.cloned()
.collect()
}
#[tokio::test]
async fn limit_satisfied_stream_emits_exactly_one_completion_event() {
let gateway = MockGateway::start(|req| items_handler(req, 5)).await;
let exec = build_exec(online_client(&gateway).await, None, Some(1), 1);
let (_guard, events) = capture_events();
let mut stream = exec
.execute(0, Arc::new(TaskContext::default()))
.expect("stream");
let batch = stream.next().await.expect("a batch").expect("no error");
assert_eq!(batch.num_rows(), 1);
drop(stream);
let completed = events_with_message(&events, COMPLETED);
assert_eq!(completed.len(), 1, "exactly one completion event");
let event = &completed[0];
assert_eq!(event.level, tracing::Level::INFO);
assert_eq!(event.field("gateway"), Some("saas"));
assert_eq!(event.field("binding"), Some("ws"));
assert_eq!(event.field("table"), Some("mock.items"));
assert_eq!(event.field("action"), Some("mock.list_items"));
assert_eq!(event.field("cache_hit"), Some("false"));
assert_eq!(event.field("pages"), Some("1"));
assert_eq!(event.field("rows"), Some("1"));
assert!(event.fields.contains_key("duration_ms"));
assert!(events_with_message(&events, FAILED).is_empty());
}
#[tokio::test]
async fn empty_scan_emits_exactly_one_completion_event() {
let gateway = MockGateway::start(|req| items_handler(req, 0)).await;
let exec = build_exec(online_client(&gateway).await, None, None, 1);
let (_guard, events) = capture_events();
let mut stream = exec
.execute(0, Arc::new(TaskContext::default()))
.expect("stream");
assert!(
stream.next().await.is_none(),
"empty scan yields no batches"
);
assert!(stream.next().await.is_none(), "stream stays terminated");
let completed = events_with_message(&events, COMPLETED);
assert_eq!(completed.len(), 1, "exactly one completion event");
assert_eq!(completed[0].field("rows"), Some("0"));
assert_eq!(completed[0].field("pages"), Some("1"));
assert_eq!(completed[0].field("cache_hit"), Some("false"));
}
#[tokio::test]
async fn cache_replay_emits_exactly_one_completion_event_marked_cache_hit() {
let gateway = MockGateway::start(|req| items_handler(req, 3)).await;
let cache = Arc::new(ScanCache::new(Duration::from_secs(60), usize::MAX));
let exec = build_exec(online_client(&gateway).await, Some(cache), None, 1);
let (_guard, events) = capture_events();
for round in 1..=2 {
let mut stream = exec
.execute(0, Arc::new(TaskContext::default()))
.expect("stream");
let mut rows = 0;
while let Some(batch) = stream.next().await {
rows += batch.expect("no error").num_rows();
}
assert_eq!(rows, 3, "round {round}");
}
let completed = events_with_message(&events, COMPLETED);
assert_eq!(
completed.len(),
2,
"one completion event per scan, replay included"
);
assert_eq!(completed[0].field("cache_hit"), Some("false"));
assert_eq!(completed[0].field("pages"), Some("2"));
assert_eq!(completed[1].field("cache_hit"), Some("true"));
assert_eq!(
completed[1].field("pages"),
Some("0"),
"a replay fetches no live pages"
);
assert_eq!(completed[1].field("rows"), Some("3"));
}
#[tokio::test]
async fn failed_scan_emits_one_failure_event_and_no_completion() {
let gateway = MockGateway::start(|req| {
if req.method == "POST" {
MockResponse::new(500, r#"{"message":"boom"}"#)
} else {
MockResponse::new(404, "{}")
}
})
.await;
let exec = build_exec(online_client(&gateway).await, None, None, 1);
let (_guard, events) = capture_events();
let mut stream = exec
.execute(0, Arc::new(TaskContext::default()))
.expect("stream");
let result = stream.next().await.expect("one item");
assert!(result.is_err(), "scan must fail");
let failed = events_with_message(&events, FAILED);
assert_eq!(failed.len(), 1, "exactly one failure event");
let event = &failed[0];
assert_eq!(event.level, tracing::Level::WARN);
assert_eq!(event.field("gateway"), Some("saas"));
assert_eq!(event.field("binding"), Some("ws"));
assert_eq!(event.field("table"), Some("mock.items"));
assert_eq!(event.field("action"), Some("mock.list_items"));
assert_eq!(event.field("pages_fetched"), Some("0"));
let error = event.field("error").expect("error field");
assert!(
error.contains("HTTP 500"),
"error carries the terminal status: {error}"
);
assert!(events_with_message(&events, COMPLETED).is_empty());
}
#[test]
fn cache_key_uses_the_bound_pack_version() {
let state_v1 = ScanState::new(&exec_with_version(1)).expect("state v1");
let state_v2 = ScanState::new(&exec_with_version(2)).expect("state v2");
assert!(
state_v1.cache_key.contains(r#""source_pack_version":1"#),
"v1 key: {}",
state_v1.cache_key
);
assert!(
state_v2.cache_key.contains(r#""source_pack_version":2"#),
"v2 key: {}",
state_v2.cache_key
);
assert_ne!(state_v1.cache_key, state_v2.cache_key);
}
}