pub mod action_registry;
pub mod cache;
pub mod client;
pub mod config;
mod error;
pub mod exec;
pub mod filters;
pub mod json_to_arrow;
pub mod packs;
pub mod pagination;
mod raw_schema;
pub mod row_path;
pub mod source_pack;
pub mod table;
pub mod table_functions;
#[cfg(test)]
pub(crate) mod testutil;
pub use action_registry::{ActionMetadata, ActionRegistry};
pub use client::OpenConnectorClient;
pub use config::{OpenConnectorBinding, OpenConnectorConfig};
pub use error::OpenConnectorError;
pub use source_pack::{FixedValue, SourcePack, SourcePackRegistry, SourcePackTable};
pub use table::OpenConnectorTableProvider;
pub use table_functions::{GatewayHandle, OpenConnectorGateways, register_open_connector_udtfs};
use std::sync::Arc;
use std::time::Duration;
use datafusion::catalog::{
CatalogProvider, MemoryCatalogProvider, MemorySchemaProvider, SchemaProvider,
};
use datafusion::prelude::SessionContext;
use serde_json::Value;
use crate::sources::hierarchy::HierarchyLevel;
use anyhow::Result;
#[allow(clippy::too_many_arguments)]
pub async fn register_open_connector_tables(
session_ctx: &mut SessionContext,
name: &str,
connection_string: &str,
config: Option<&OpenConnectorConfig>,
read_write: bool,
hierarchy_level: HierarchyLevel,
udtf_gateways: Option<&OpenConnectorGateways>,
) -> Result<()> {
if hierarchy_level != HierarchyLevel::Catalog {
return Err(OpenConnectorError::CatalogHierarchyRequired {
name: name.to_string(),
}
.into());
}
if read_write {
return Err(OpenConnectorError::ReadWriteNotSupported {
name: name.to_string(),
}
.into());
}
let config = config.ok_or_else(|| OpenConnectorError::MissingConfig {
name: name.to_string(),
})?;
if connection_string.trim().is_empty() {
return Err(OpenConnectorError::EmptyGatewayUrl {
name: name.to_string(),
}
.into());
}
config.validate()?;
let client = Arc::new(OpenConnectorClient::from_config(connection_string, config)?);
client.health().await?;
let pack_registry = SourcePackRegistry::builtins()?;
let mut action_ids = config.raw_action_allowlist.clone();
for binding in &config.bindings {
let pack = pack_registry.require(&binding.source_pack)?;
SourcePackRegistry::check_version_pin(pack, binding.source_pack_version)?;
let mut tables = Vec::with_capacity(binding.tables.len());
for table_name in &binding.tables {
let table = pack_registry.table(pack, table_name)?;
action_ids.push(table.action_id.to_string());
for key in table.required_resources {
if !binding.resource.contains_key(*key) {
return Err(OpenConnectorError::MissingResourceInput {
binding: binding.name.clone(),
key: (*key).to_string(),
}
.into());
}
}
tables.push(table);
}
for key in binding.resource.keys() {
if !tables.iter().any(|table| table.declares_resource(key)) {
return Err(OpenConnectorError::UnknownResourceKey {
binding: binding.name.clone(),
key: key.clone(),
}
.into());
}
}
}
let registry = Arc::new(ActionRegistry::load(&client, &action_ids).await?);
let catalog = Arc::new(MemoryCatalogProvider::new());
let cache = Arc::new(cache::ScanCache::new(
Duration::from_secs(config.cache_ttl_seconds),
usize::try_from(config.cache_max_bytes).unwrap_or(usize::MAX),
));
let scan_timeout = Duration::from_secs(config.scan_timeout_seconds);
for binding in &config.bindings {
let pack = pack_registry.require(&binding.source_pack)?;
let schema_provider = Arc::new(MemorySchemaProvider::new());
for table_name in &binding.tables {
let table = pack_registry.table(pack, table_name)?;
if let Some(expected) = table.expected_fingerprint {
let actual = registry
.get(table.action_id)
.map(ActionMetadata::fingerprint);
if actual != Some(expected) {
return Err(OpenConnectorError::ActionContractMismatch {
table: table.id.to_string(),
reason: format!(
"action '{}' fingerprint mismatch (expected {expected}, discovered {})",
table.action_id,
actual.unwrap_or("<none>")
),
}
.into());
}
}
let provider = OpenConnectorTableProvider::new(
Arc::clone(&client),
Some(Arc::clone(&cache)),
name.to_string(),
Some(binding.name.clone()),
binding.connection_alias.clone(),
table,
pack.version,
Value::Object(binding.resource.clone().into_iter().collect()),
config.max_pages,
config.max_rows,
scan_timeout,
)?;
schema_provider
.register_table(table_name.clone(), Arc::new(provider))
.map_err(|e| OpenConnectorError::CatalogRegistrationFailed {
name: format!("{name}.{}.{table_name}", binding.name),
reason: format!("failed to register table into catalog schema: {e}"),
})?;
}
catalog
.register_schema(&binding.name, schema_provider)
.map_err(|e| OpenConnectorError::CatalogRegistrationFailed {
name: format!("{name}.{}", binding.name),
reason: format!("failed to register schema in catalog: {e}"),
})?;
}
session_ctx.register_catalog(name, catalog);
if let Some(gateways) = udtf_gateways {
gateways.write().unwrap_or_else(|p| p.into_inner()).insert(
name.to_string(),
Arc::new(GatewayHandle::new(
Arc::clone(&client),
Arc::clone(&cache),
Arc::clone(®istry),
config,
)),
);
}
tracing::info!(
gateway = %name,
bindings = config.bindings.len(),
actions = registry.len(),
"Open Connector catalog registered"
);
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::sources::providers::open_connector::testutil::{
MockGateway, MockResponse, RecordedRequest, discovery_ok, envelope_ok,
};
const TOKEN_ENV_HEALTH_FAIL: &str = "SKARDI_TEST_OC_REGISTER_TOKEN_HEALTH_FAIL";
fn valid_config(token_env: &str) -> OpenConnectorConfig {
serde_yaml::from_str(&format!(
r#"
runtime_token_env: {token_env}
raw_action_allowlist:
- github.list_repository_issues
bindings:
- name: github_skardi
source_pack: github
resource: {{ owner: SkardiLabs, repo: skardi }}
tables: [issues]
"#
))
.expect("parse config")
}
#[tokio::test]
async fn register_rejects_table_hierarchy_before_any_network() {
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
"http://127.0.0.1:1",
Some(&valid_config("UNUSED_ENV")),
false,
HierarchyLevel::Table,
None,
)
.await
.unwrap_err();
let err = err.downcast::<OpenConnectorError>().unwrap();
assert!(matches!(
err,
OpenConnectorError::CatalogHierarchyRequired { ref name } if name == "saas"
));
}
#[tokio::test]
async fn register_rejects_empty_gateway_url_before_any_network() {
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
" ",
Some(&valid_config("UNUSED_ENV")),
false,
HierarchyLevel::Catalog,
None,
)
.await
.unwrap_err();
let err = err.downcast::<OpenConnectorError>().unwrap();
assert!(matches!(
err,
OpenConnectorError::EmptyGatewayUrl { ref name } if name == "saas"
));
}
#[tokio::test]
async fn register_rejects_read_write_before_any_network() {
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
"http://127.0.0.1:1",
Some(&valid_config("UNUSED_ENV")),
true,
HierarchyLevel::Catalog,
None,
)
.await
.unwrap_err();
let err = err.downcast::<OpenConnectorError>().unwrap();
assert!(matches!(
err,
OpenConnectorError::ReadWriteNotSupported { ref name } if name == "saas"
));
}
#[tokio::test]
async fn register_rejects_missing_config_block_before_any_network() {
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
"http://127.0.0.1:1",
None,
false,
HierarchyLevel::Catalog,
None,
)
.await
.unwrap_err();
let err = err.downcast::<OpenConnectorError>().unwrap();
assert!(matches!(
err,
OpenConnectorError::MissingConfig { ref name } if name == "saas"
));
}
#[tokio::test]
async fn register_rejects_invalid_config_before_any_network() {
let mut ctx = SessionContext::new();
let invalid: OpenConnectorConfig = serde_yaml::from_str(
"runtime_token_env: ''\nbindings:\n - name: b\n source_pack: github\n tables: [issues]",
)
.expect("parse config");
let err = register_open_connector_tables(
&mut ctx,
"saas",
"http://127.0.0.1:1",
Some(&invalid),
false,
HierarchyLevel::Catalog,
None,
)
.await
.unwrap_err();
let err = err.downcast::<OpenConnectorError>().unwrap();
assert!(matches!(err, OpenConnectorError::EmptyRuntimeTokenEnv));
}
#[tokio::test]
async fn register_fails_missing_token_env_before_network() {
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
"http://127.0.0.1:1",
Some(&valid_config("SKARDI_TEST_OC_TOKEN_DEFINITELY_UNSET")),
false,
HierarchyLevel::Catalog,
None,
)
.await
.unwrap_err();
let err = err.downcast::<OpenConnectorError>().unwrap();
assert!(matches!(
err,
OpenConnectorError::MissingRuntimeToken { ref env }
if env == "SKARDI_TEST_OC_TOKEN_DEFINITELY_UNSET"
));
}
#[tokio::test]
async fn register_fails_health_check_before_discovery() {
let gateway = MockGateway::start(|_| MockResponse::new(503, "{}")).await;
unsafe {
std::env::set_var(TOKEN_ENV_HEALTH_FAIL, "test-token");
}
let mut ctx = SessionContext::new();
let result = register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&valid_config(TOKEN_ENV_HEALTH_FAIL)),
false,
HierarchyLevel::Catalog,
None,
)
.await;
unsafe {
std::env::remove_var(TOKEN_ENV_HEALTH_FAIL);
}
let err = result
.unwrap_err()
.downcast::<OpenConnectorError>()
.unwrap();
assert!(
matches!(err, OpenConnectorError::RetriesExhausted { .. }),
"got {err}"
);
let requests = gateway.requests();
assert!(
!requests.is_empty() && requests.iter().all(|r| r.path == "/v1/health"),
"only health calls were attempted: {:?}",
requests.iter().map(|r| &r.path).collect::<Vec<_>>()
);
}
#[tokio::test]
async fn register_builds_queryable_catalog_with_mock_pack() {
let gateway = MockGateway::start(|req| mock_gateway_handler(req, 5)).await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_BASIC, "test-token");
}
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&mock_config(TOKEN_ENV_CATALOG_BASIC, 0)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_BASIC);
}
let df = ctx
.sql("SELECT id, name FROM saas.ws.items ORDER BY id")
.await
.expect("plan");
let batches = df.collect().await.expect("collect");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 5);
let executes = execute_requests(&gateway);
assert_eq!(executes.len(), 3, "3 pages for 5 items at per_page=2");
assert!(executes[0].body.contains(r#""page":1"#));
assert!(executes[2].body.contains(r#""page":3"#));
}
#[tokio::test]
async fn scan_pushes_allowlisted_filter_and_stops_at_limit() {
let gateway = MockGateway::start(|req| mock_gateway_handler(req, 5)).await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_FILTER, "test-token");
}
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&mock_config(TOKEN_ENV_CATALOG_FILTER, 0)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_FILTER);
}
let df = ctx
.sql("SELECT id, value FROM saas.ws.items WHERE value > 3.0")
.await
.expect("plan");
let batches = df.collect().await.expect("collect");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 2, "values 4.0 and 5.0");
assert!(
execute_requests(&gateway)
.iter()
.all(|r| r.body.contains(r#""min_value":3"#)),
"min_value pushed on every page"
);
let gateway2 = MockGateway::start(|req| mock_gateway_handler(req, 5)).await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_FILTER, "test-token");
}
let mut ctx2 = SessionContext::new();
register_open_connector_tables(
&mut ctx2,
"saas",
&gateway2.url,
Some(&mock_config(TOKEN_ENV_CATALOG_FILTER, 0)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_FILTER);
}
let df = ctx2
.sql("SELECT id FROM saas.ws.items LIMIT 1")
.await
.expect("plan");
let batches = df.collect().await.expect("collect");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 1);
assert_eq!(
execute_requests(&gateway2).len(),
1,
"LIMIT 1 must stop after the first page"
);
}
#[tokio::test]
async fn in_band_error_key_fails_the_scan_as_a_provider_error() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path == "/v1/actions/mock.list_items" {
return MockResponse::ok(&discovery_ok("{}", r#"{"type": "object"}"#, true, None));
}
if req.method == "POST" && req.path == "/v1/actions/mock.list_items" {
return MockResponse::ok(&envelope_ok(r#"{"error": "missing_scope"}"#));
}
MockResponse::new(404, "{}")
})
.await;
let token_env = "SKARDI_TEST_OC_INBAND_ERROR";
unsafe {
std::env::set_var(token_env, "test-token");
}
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&mock_config(token_env, 0)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("registration succeeds");
unsafe {
std::env::remove_var(token_env);
}
let err = ctx
.sql("SELECT id FROM saas.ws.items")
.await
.expect("plan")
.collect()
.await
.expect_err("the in-band error must fail the scan");
let message = err.to_string();
assert!(
message.contains("missing_scope") && message.contains("mock.list_items"),
"the provider's own code and the action are named: {message}"
);
assert!(
!message.contains("row path"),
"never the misleading row-path error: {message}"
);
}
#[tokio::test]
async fn scan_deadline_bounds_retry_waits() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path == "/v1/actions/mock.list_items" {
return MockResponse::ok(&discovery_ok("{}", r#"{"type": "object"}"#, true, None));
}
if req.method == "POST" && req.path == "/v1/actions/mock.list_items" {
return MockResponse::new(429, "{}").with_header("retry-after", "2");
}
MockResponse::new(404, "{}")
})
.await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_TIMEOUT, "test-token");
}
let mut config = mock_config(TOKEN_ENV_CATALOG_TIMEOUT, 0);
config.scan_timeout_seconds = 1;
config.request_timeout_seconds = 30;
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&config),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_TIMEOUT);
}
let df = ctx.sql("SELECT id FROM saas.ws.items").await.expect("plan");
let err = df.collect().await.expect_err("scan must time out");
assert!(err.to_string().contains("timed out after 1s"), "got {err}");
assert_eq!(
execute_requests(&gateway).len(),
1,
"retry wait was cancelled"
);
}
#[tokio::test]
async fn gteq_is_not_pushed_to_strict_gt_input() {
let gateway = MockGateway::start(|req| mock_gateway_handler(req, 5)).await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_FILTER, "test-token");
}
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&mock_config(TOKEN_ENV_CATALOG_FILTER, 0)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_FILTER);
}
let before = execute_requests(&gateway).len();
let df = ctx
.sql("SELECT id FROM saas.ws.items WHERE value >= 3.0 ORDER BY id")
.await
.expect("plan");
let batches = df.collect().await.expect("collect");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 3, "boundary row id=3 must be present (ids 3,4,5)");
let new_requests = &execute_requests(&gateway)[before..];
assert!(
new_requests.iter().all(|r| !r.body.contains("min_value")),
"no min_value may be pushed for >=: {:?}",
new_requests.iter().map(|r| &r.body).collect::<Vec<_>>()
);
}
#[tokio::test]
async fn cached_scan_replays_without_new_requests() {
let gateway = MockGateway::start(|req| mock_gateway_handler(req, 3)).await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_CACHE, "test-token");
}
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&mock_config(TOKEN_ENV_CATALOG_CACHE, 60)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_CACHE);
}
for round in 1..=2 {
let df = ctx
.sql("SELECT id, name FROM saas.ws.items ORDER BY id")
.await
.expect("plan");
let batches = df.collect().await.expect("collect");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 3, "round {round}");
}
let executes = execute_requests(&gateway);
assert_eq!(
executes.len(),
2,
"second identical scan must be served from cache (3 items at per_page=2 → 2 live pages)"
);
}
#[tokio::test]
async fn limited_scan_is_cached_and_replayed() {
let gateway = MockGateway::start(|req| mock_gateway_handler(req, 5)).await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_CACHE, "test-token");
}
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&mock_config(TOKEN_ENV_CATALOG_CACHE, 60)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_CACHE);
}
for round in 1..=2 {
let df = ctx
.sql("SELECT id FROM saas.ws.items LIMIT 2")
.await
.expect("plan");
let batches = df.collect().await.expect("collect");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 2, "round {round}");
}
assert_eq!(
execute_requests(&gateway).len(),
1,
"first LIMIT scan fetches one page; the replay adds none"
);
}
#[tokio::test]
async fn full_scan_after_limited_scan_never_replays_the_truncated_entry() {
let gateway = MockGateway::start(|req| mock_gateway_handler(req, 5)).await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_LIMIT_FULL, "test-token");
}
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&mock_config(TOKEN_ENV_CATALOG_LIMIT_FULL, 60)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_LIMIT_FULL);
}
let df = ctx
.sql("SELECT id FROM saas.ws.items LIMIT 2")
.await
.expect("plan");
let rows: usize = df
.collect()
.await
.expect("collect")
.iter()
.map(|b| b.num_rows())
.sum();
assert_eq!(rows, 2);
let live_pages = execute_requests(&gateway).len();
let df = ctx.sql("SELECT id FROM saas.ws.items").await.expect("plan");
let rows: usize = df
.collect()
.await
.expect("collect")
.iter()
.map(|b| b.num_rows())
.sum();
assert_eq!(
rows, 5,
"the truncated entry must never serve a fuller query"
);
assert!(
execute_requests(&gateway).len() > live_pages,
"the full scan fetched live"
);
}
#[tokio::test]
async fn cached_empty_scan_replays_without_new_requests() {
let gateway = MockGateway::start(|req| mock_gateway_handler(req, 0)).await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_EMPTY_CACHE, "test-token");
}
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&mock_config(TOKEN_ENV_CATALOG_EMPTY_CACHE, 60)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_EMPTY_CACHE);
}
for round in 1..=2 {
let df = ctx.sql("SELECT id FROM saas.ws.items").await.expect("plan");
let batches = df.collect().await.expect("collect");
assert!(batches.is_empty(), "round {round} should be empty");
}
assert_eq!(
execute_requests(&gateway).len(),
1,
"second empty scan must be served from cache"
);
}
#[tokio::test]
async fn self_join_scans_compute_identical_keys_but_fetch_live() {
let gateway = MockGateway::start(|req| mock_gateway_handler(req, 3)).await;
unsafe {
std::env::set_var(TOKEN_ENV_CATALOG_SELFJOIN, "test-token");
}
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&mock_config(TOKEN_ENV_CATALOG_SELFJOIN, 60)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect("catalog registration succeeds");
unsafe {
std::env::remove_var(TOKEN_ENV_CATALOG_SELFJOIN);
}
let df = ctx
.sql(
"SELECT count(*) AS n FROM saas.ws.items i1 JOIN saas.ws.items i2 ON i1.id = i2.id",
)
.await
.expect("plan");
let batches = df.collect().await.expect("collect");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 1, "count(*) returns one row");
let before = execute_requests(&gateway).len();
assert_eq!(before, 4, "concurrent join sides both fetch live");
let df = ctx.sql("SELECT id FROM saas.ws.items").await.expect("plan");
let batches = df.collect().await.expect("collect");
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 3);
assert_eq!(
execute_requests(&gateway).len(),
before,
"no new live pages once earlier scans have completed and cached"
);
}
#[cfg(test)]
const TOKEN_ENV_CATALOG_BASIC: &str = "SKARDI_TEST_OC_REGISTER_CATALOG_BASIC";
#[cfg(test)]
const TOKEN_ENV_CATALOG_FILTER: &str = "SKARDI_TEST_OC_REGISTER_CATALOG_FILTER";
#[cfg(test)]
const TOKEN_ENV_CATALOG_CACHE: &str = "SKARDI_TEST_OC_REGISTER_CATALOG_CACHE";
#[cfg(test)]
const TOKEN_ENV_CATALOG_SELFJOIN: &str = "SKARDI_TEST_OC_REGISTER_CATALOG_SELFJOIN";
#[cfg(test)]
const TOKEN_ENV_CATALOG_EMPTY_CACHE: &str = "SKARDI_TEST_OC_REGISTER_CATALOG_EMPTY_CACHE";
#[cfg(test)]
const TOKEN_ENV_CATALOG_LIMIT_FULL: &str = "SKARDI_TEST_OC_REGISTER_CATALOG_LIMIT_FULL";
#[cfg(test)]
const TOKEN_ENV_CATALOG_TIMEOUT: &str = "SKARDI_TEST_OC_REGISTER_CATALOG_TIMEOUT";
#[cfg(test)]
fn mock_config(token_env: &str, cache_ttl_seconds: u64) -> OpenConnectorConfig {
serde_yaml::from_str(&format!(
r#"
runtime_token_env: {token_env}
cache_ttl_seconds: {cache_ttl_seconds}
bindings:
- name: ws
source_pack: mock
resource: {{ workspace: demo }}
tables: [items]
"#
))
.expect("parse config")
}
#[cfg(test)]
fn mock_items() -> Vec<serde_json::Value> {
(1..=5)
.map(|id| {
serde_json::json!({
"id": id,
"name": format!("item-{id}"),
"value": id as f64,
"tags": ["t1", "t2"],
"created_at": "2026-01-01T00:00:00Z"
})
})
.collect()
}
#[cfg(test)]
fn mock_gateway_handler(req: &RecordedRequest, total: usize) -> MockResponse {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path == "/v1/actions/mock.list_items" {
return MockResponse::ok(&discovery_ok("{}", r#"{"type": "object"}"#, true, None));
}
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 input = body.get("input").cloned().unwrap_or_default();
let page = input
.get("page")
.and_then(serde_json::Value::as_u64)
.unwrap_or(1) as usize;
let min_value = input.get("min_value").and_then(serde_json::Value::as_f64);
let items = mock_items();
let start = (page - 1) * 2;
let slice: Vec<_> = items
.into_iter()
.take(total)
.filter(|item| {
min_value.is_none_or(|min| {
item.get("value").and_then(serde_json::Value::as_f64) > Some(min)
})
})
.skip(start)
.take(2)
.collect();
return MockResponse::ok(&envelope_ok(
&serde_json::json!({ "items": slice }).to_string(),
));
}
MockResponse::new(404, "{}")
}
#[cfg(test)]
fn execute_requests(gateway: &MockGateway) -> Vec<RecordedRequest> {
gateway
.requests()
.into_iter()
.filter(|r| r.method == "POST")
.collect()
}
}