use std::sync::OnceLock;
use crate::sources::providers::open_connector::error::OpenConnectorError;
use crate::sources::providers::open_connector::source_pack::SourcePack;
use super::loader;
static PACK: OnceLock<Result<SourcePack, String>> = OnceLock::new();
pub fn pack() -> Result<&'static SourcePack, OpenConnectorError> {
loader::builtin("dropbox.yaml", include_str!("dropbox.yaml"), &PACK)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::sources::hierarchy::HierarchyLevel;
use crate::sources::providers::open_connector::action_registry::fingerprint_schema;
use crate::sources::providers::open_connector::json_to_arrow::RowConverter;
use crate::sources::providers::open_connector::row_path::RowPath;
use crate::sources::providers::open_connector::source_pack::SourcePackTable;
use crate::sources::providers::open_connector::testutil::{
EnvVarGuard, MockGateway, MockResponse, discovery_ok, envelope_ok,
fingerprint_uncovered_columns,
};
use crate::sources::providers::open_connector::{
OpenConnectorConfig, OpenConnectorGateways, register_open_connector_tables,
register_open_connector_udtfs,
};
use arrow::array::{Array, BooleanArray, Int64Array, StringArray, TimestampMillisecondArray};
use arrow::record_batch::RecordBatch;
use datafusion::prelude::SessionContext;
use serde_json::{Value, json};
fn table(short: &str) -> &'static SourcePackTable {
pack()
.expect("embedded asset is test-pinned to parse")
.tables
.iter()
.find(|t| t.id.rsplit('.').next() == Some(short))
.unwrap_or_else(|| panic!("table {short}"))
}
fn contracts() -> Vec<(&'static str, &'static str)> {
vec![
(
"dropbox.list_folder",
include_str!("fixtures/dropbox/contracts/list_folder.json"),
),
(
"dropbox.list_folder_continue",
include_str!("fixtures/dropbox/contracts/list_folder_continue.json"),
),
(
"dropbox.list_shared_links",
include_str!("fixtures/dropbox/contracts/list_shared_links.json"),
),
(
"dropbox.search_files",
include_str!("fixtures/dropbox/contracts/search_files.json"),
),
(
"dropbox.search_files_continue",
include_str!("fixtures/dropbox/contracts/search_files_continue.json"),
),
]
}
#[test]
fn pinned_fingerprints_match_the_committed_contracts() {
let mut mismatches = Vec::new();
for (action, contract) in contracts() {
let schema: Value = serde_json::from_str(contract).expect("contract parses");
let actual = fingerprint_schema(Some(&schema));
let expected: Vec<&str> = pack()
.expect("parses")
.tables
.iter()
.flat_map(SourcePackTable::gated_actions)
.filter(|(id, _)| *id == action)
.map(|(_, fingerprint)| fingerprint)
.collect();
assert!(
!expected.is_empty(),
"{action} has a captured contract but no table pins it"
);
for pin in expected {
if pin != actual {
mismatches.push(format!("{action}: pinned {pin}, actual {actual}"));
}
}
}
assert!(mismatches.is_empty(), "fingerprint drift:\n{mismatches:#?}");
}
#[test]
fn the_cursor_only_claim_holds_against_the_committed_input_contracts() {
for short in ["files", "file_search"] {
let table = table(short);
let continuation = table
.continuation
.unwrap_or_else(|| panic!("{short} is a split-action table"));
assert!(
continuation.cursor_only,
"{short} declares `inputs: cursor_only`"
);
let captured =
continuation_input_contract(continuation.action_id).unwrap_or_else(|| {
panic!("{} has a committed input contract", continuation.action_id)
});
let schema: Value = serde_json::from_str(captured).expect("input contract parses");
table
.check_continuation_inputs(Some(&schema))
.unwrap_or_else(|e| {
panic!(
"{short}: the committed input contract for {} no longer satisfies \
`inputs: cursor_only`: {e}",
continuation.action_id
)
});
}
}
#[test]
fn every_mapped_column_is_inside_the_fingerprint_gate() {
for (short, contract) in [
(
"files",
include_str!("fixtures/dropbox/contracts/list_folder.json"),
),
(
"shared_links",
include_str!("fixtures/dropbox/contracts/list_shared_links.json"),
),
(
"file_search",
include_str!("fixtures/dropbox/contracts/search_files.json"),
),
] {
let table = table(short);
let uncovered = fingerprint_uncovered_columns(contract, table.row_path, table.fields);
assert!(
uncovered.is_empty(),
"{short}: columns outside the fingerprint gate: {uncovered:?}"
);
}
}
#[test]
fn split_action_tables_declare_a_cursor_only_continuation() {
for (short, opener, continues) in [
(
"files",
"dropbox.list_folder",
"dropbox.list_folder_continue",
),
(
"file_search",
"dropbox.search_files",
"dropbox.search_files_continue",
),
] {
let table = table(short);
assert_eq!(table.action_id, opener);
let continuation = table
.continuation
.unwrap_or_else(|| panic!("{short} must declare a continuation"));
assert_eq!(continuation.action_id, continues);
assert!(
continuation.cursor_only,
"{short}: the continue action accepts the cursor and nothing else"
);
assert_eq!(
table.gated_actions().count(),
2,
"{short}: both actions must be fingerprint-gated"
);
}
let shared = table("shared_links");
assert!(shared.continuation.is_none());
assert_eq!(shared.gated_actions().count(), 1);
}
#[test]
fn no_table_pushes_a_filter_or_maps_shared_link_only_columns_onto_files() {
for short in ["files", "shared_links", "file_search"] {
assert!(
table(short).filters.is_empty(),
"{short}: Dropbox exposes no faithful column predicate; \
pushing one needs a documented rationale first"
);
}
let files: Vec<&str> = table("files").fields.iter().map(|f| f.name).collect();
for absent in ["url", "expires_at", "link_permissions"] {
assert!(
!files.contains(&absent),
"files must not map '{absent}': files/list_folder never returns it, \
so the column would be structurally always-NULL"
);
}
let shared: Vec<&str> = table("shared_links")
.fields
.iter()
.map(|f| f.name)
.collect();
for present in ["url", "expires_at", "link_permissions"] {
assert!(shared.contains(&present), "shared_links must map {present}");
}
for absent in [
"path_display",
"is_downloadable",
"content_hash",
"sharing_info",
] {
assert!(
!shared.contains(&absent),
"shared_links must not map '{absent}': SharedLinkMetadata never \
carries it, so the column would be structurally always-NULL"
);
}
let search: Vec<&str> = table("file_search").fields.iter().map(|f| f.name).collect();
assert!(!search.contains(&"highlight_spans"));
for (key, _) in table("file_search").fixed_inputs {
assert_ne!(*key, "includeHighlights");
}
}
fn convert_fixture(table: &SourcePackTable, fixture: &str) -> RecordBatch {
let page: Value = serde_json::from_str(fixture).expect("fixture parses");
let rows = RowPath::parse(table.row_path)
.expect("row path")
.rows(&page, 1)
.expect("row array");
RowConverter::new(table.fields)
.expect("converter")
.convert(rows, 1)
.expect("fixture converts")
}
fn utf8<'a>(batch: &'a RecordBatch, name: &str) -> &'a StringArray {
batch
.column_by_name(name)
.unwrap_or_else(|| panic!("column {name}"))
.as_any()
.downcast_ref()
.expect("Utf8 column")
}
#[test]
fn files_fixture_converts_nulls_nested_and_extra_fields() {
let batch = convert_fixture(table("files"), include_str!("fixtures/dropbox/files.json"));
assert_eq!(batch.num_rows(), 3);
assert_eq!(
utf8(&batch, "tag").iter().collect::<Vec<_>>(),
vec![Some("folder"), Some("file"), Some("file")]
);
let sizes: &Int64Array = batch
.column_by_name("size_bytes")
.expect("size_bytes")
.as_any()
.downcast_ref()
.expect("Int64 column");
assert!(sizes.is_null(0), "a folder has no size");
assert_eq!(sizes.value(1), 284_913);
assert_eq!(sizes.value(2), 0, "an empty file is 0, not NULL");
let downloadable: &BooleanArray = batch
.column_by_name("is_downloadable")
.expect("is_downloadable")
.as_any()
.downcast_ref()
.expect("Boolean column");
assert!(downloadable.is_null(0));
assert!(downloadable.value(1));
let modified: &TimestampMillisecondArray = batch
.column_by_name("server_modified")
.expect("server_modified")
.as_any()
.downcast_ref()
.expect("Timestamp column");
assert!(modified.is_null(0), "folders carry no serverModified");
assert!(modified.value(1) > 0, "ISO 8601 parses through RFC 3339");
let sharing = utf8(&batch, "sharing_info");
assert!(sharing.is_null(0));
assert!(
sharing.value(1).contains("parent_shared_folder_id"),
"nested object preserved: {}",
sharing.value(1)
);
}
#[test]
fn empty_page_converts_to_zero_rows() {
let batch = convert_fixture(
table("files"),
include_str!("fixtures/dropbox/files_empty.json"),
);
assert_eq!(batch.num_rows(), 0);
assert_eq!(batch.schema().fields().len(), table("files").fields.len());
}
#[test]
fn every_table_converts_an_empty_page_and_keeps_its_schema() {
for table in pack().expect("embedded asset parses").tables {
let batch = RowConverter::new(table.fields)
.expect("converter")
.convert(&[], 1)
.expect("empty page");
assert_eq!(batch.num_rows(), 0, "{}", table.id);
assert_eq!(
batch.schema().fields().len(),
table.fields.len(),
"{} keeps its stable schema on empty results",
table.id
);
}
}
#[test]
fn schema_mismatch_fails_with_the_full_error_identity() {
let table = table("files");
let page: Value =
serde_json::from_str(include_str!("fixtures/dropbox/files_type_mismatch.json"))
.expect("fixture parses");
let rows = RowPath::parse(table.row_path)
.expect("row path")
.rows(&page, 7)
.expect("row array");
let err = RowConverter::new(table.fields)
.expect("converter")
.convert(rows, 7)
.expect_err("a string where int64 is declared must fail");
let rendered = err.to_string();
for fragment in [
"size_bytes",
"sizeBytes",
"page 7",
"row 1",
"expected integer",
"found a string",
] {
assert!(
rendered.contains(fragment),
"error must carry {fragment:?}: {rendered}"
);
}
assert!(
!rendered.contains("not-a-number"),
"the offending VALUE must never appear: {rendered}"
);
}
#[test]
fn shared_links_fixture_populates_the_link_only_columns() {
let batch = convert_fixture(
table("shared_links"),
include_str!("fixtures/dropbox/shared_links.json"),
);
assert_eq!(batch.num_rows(), 2);
let urls = utf8(&batch, "url");
assert!(urls.value(0).starts_with("https://www.dropbox.com/scl/fi/"));
assert!(urls.value(1).starts_with("https://www.dropbox.com/scl/fo/"));
let expires: &TimestampMillisecondArray = batch
.column_by_name("expires_at")
.expect("expires_at")
.as_any()
.downcast_ref()
.expect("Timestamp column");
assert!(expires.value(0) > 0, "a link with an expiry");
assert!(expires.is_null(1), "a link without one");
let permissions = utf8(&batch, "link_permissions");
assert!(permissions.value(0).contains("resolved_visibility"));
}
#[test]
fn file_search_fixture_reads_through_the_nested_metadata_block() {
let batch = convert_fixture(
table("file_search"),
include_str!("fixtures/dropbox/file_search.json"),
);
assert_eq!(batch.num_rows(), 2);
assert_eq!(
utf8(&batch, "match_type").iter().collect::<Vec<_>>(),
vec![Some("filename"), Some("content")]
);
assert_eq!(
utf8(&batch, "name").iter().collect::<Vec<_>>(),
vec![Some("redacted-report.pdf"), Some("notes.txt")]
);
}
#[test]
fn a_null_metadata_parent_fails_the_scan_naming_the_non_nullable_column() {
let table = table("file_search");
let page: Value = serde_json::from_str(include_str!(
"fixtures/dropbox/file_search_null_parent.json"
))
.expect("fixture parses");
let rows = RowPath::parse(table.row_path)
.expect("row path")
.rows(&page, 1)
.expect("row array");
let err = RowConverter::new(table.fields)
.expect("converter")
.convert(rows, 1)
.expect_err("a null parent under a non-nullable column must fail");
assert!(
err.to_string().contains("tag"),
"the non-nullable column names itself: {err}"
);
}
#[test]
fn complete_collection_pins_ride_every_files_request() {
let files = table("files");
let pinned: Vec<(&str, Value)> = files
.fixed_inputs
.iter()
.map(|(key, value)| (*key, value.to_json()))
.collect();
assert_eq!(
pinned,
vec![
("includeDeleted", Value::Bool(false)),
("includeMountedFolders", Value::Bool(true)),
("recursive", Value::Bool(true)),
],
"files means every file under `path`, tombstones excluded"
);
let search: Vec<(&str, Value)> = table("file_search")
.fixed_inputs
.iter()
.map(|(key, value)| (*key, value.to_json()))
.collect();
assert_eq!(search, vec![("fileStatus", Value::from("active"))]);
}
fn continuation_input_contract(action: &str) -> Option<&'static str> {
match action {
"dropbox.list_folder_continue" => Some(include_str!(
"fixtures/dropbox/contracts/inputs/list_folder_continue.json"
)),
"dropbox.search_files_continue" => Some(include_str!(
"fixtures/dropbox/contracts/inputs/search_files_continue.json"
)),
_ => None,
}
}
fn dropbox_discovery(path: &str) -> MockResponse {
let action = path.rsplit('/').next().unwrap_or_default();
let input_schema = continuation_input_contract(action).unwrap_or("{}");
let output_schema = match action {
"dropbox.list_folder" => include_str!("fixtures/dropbox/contracts/list_folder.json"),
"dropbox.list_folder_continue" => {
include_str!("fixtures/dropbox/contracts/list_folder_continue.json")
}
"dropbox.list_shared_links" => {
include_str!("fixtures/dropbox/contracts/list_shared_links.json")
}
"dropbox.search_files" => include_str!("fixtures/dropbox/contracts/search_files.json"),
"dropbox.search_files_continue" => {
include_str!("fixtures/dropbox/contracts/search_files_continue.json")
}
_ => r#"{"type": "object"}"#,
};
MockResponse::ok(&discovery_ok(input_schema, output_schema, true, None))
}
fn dropbox_config(token_env: &str, tables: &str) -> OpenConnectorConfig {
let resource = if tables.contains("file_search") {
"resource: { query: redacted }"
} else {
""
};
serde_yaml::from_str(&format!(
r#"
runtime_token_env: {token_env}
bindings:
- name: ws
source_pack: dropbox
{resource}
tables: [{tables}]
"#
))
.expect("config parses")
}
async fn setup_with_gateway(
gateway: MockGateway,
token_env: &'static str,
tables: &str,
) -> (MockGateway, SessionContext) {
let _token = EnvVarGuard::set(token_env, "test-token");
let gateways = OpenConnectorGateways::default();
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&dropbox_config(token_env, tables)),
false,
HierarchyLevel::Catalog,
Some(&gateways),
)
.await
.expect("gateway registration succeeds");
register_open_connector_udtfs(&ctx, gateways).expect("UDTF registration succeeds");
(gateway, ctx)
}
async fn collect(ctx: &SessionContext, sql: &str) -> Vec<RecordBatch> {
ctx.sql(sql)
.await
.expect("plan")
.collect()
.await
.expect("collect")
}
fn execute_calls(gateway: &MockGateway) -> Vec<(String, Value)> {
gateway
.requests()
.into_iter()
.filter(|r| r.method == "POST")
.map(|r| {
let body: Value = serde_json::from_str(&r.body).expect("request body is JSON");
(r.path.clone(), body["input"].clone())
})
.collect()
}
fn entry(name: &str) -> Value {
json!({
"tag": "file", "name": name, "id": format!("id:{name}"),
"pathDisplay": format!("/{name}"), "pathLower": format!("/{name}"),
"clientModified": null, "serverModified": null, "rev": null,
"sizeBytes": null, "isDownloadable": null, "contentHash": null,
"url": null, "expiresAt": null, "sharingInfo": null,
"linkPermissions": null
})
}
fn names_of(batches: &[RecordBatch]) -> Vec<String> {
batches
.iter()
.flat_map(|b| {
utf8(b, "name")
.iter()
.map(|v| v.expect("name is non-null").to_string())
.collect::<Vec<_>>()
})
.collect()
}
#[tokio::test]
async fn files_pages_through_the_continue_action_with_only_a_cursor() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
match req.path.as_str() {
"/v1/actions/dropbox.list_folder" => MockResponse::ok(&envelope_ok(
&json!({"entries": [entry("a"), entry("b")],
"cursor": "cur-2", "hasMore": true})
.to_string(),
)),
"/v1/actions/dropbox.list_folder_continue" => {
let body: Value = serde_json::from_str(&req.body).unwrap_or_default();
match body["input"].get("cursor").and_then(Value::as_str) {
Some("cur-2") => MockResponse::ok(&envelope_ok(
&json!({"entries": [entry("c")],
"cursor": "cur-3", "hasMore": false})
.to_string(),
)),
other => MockResponse::new(400, format!("bad cursor {other:?}")),
}
}
_ => MockResponse::new(404, "{}"),
}
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_DROPBOX_FILES", "files").await;
let batches = collect(&ctx, "SELECT name FROM saas.ws.files ORDER BY name").await;
assert_eq!(names_of(&batches), vec!["a", "b", "c"]);
let calls = execute_calls(&gateway);
assert_eq!(calls.len(), 2, "two pages");
let (path, input) = &calls[0];
assert_eq!(path, "/v1/actions/dropbox.list_folder");
assert_eq!(input["recursive"], json!(true));
assert_eq!(input["includeMountedFolders"], json!(true));
assert_eq!(input["includeDeleted"], json!(false));
assert_eq!(input["limit"], json!(2000));
assert!(
input.get("cursor").is_none(),
"no cursor exists on page one"
);
let (path, input) = &calls[1];
assert_eq!(
path, "/v1/actions/dropbox.list_folder_continue",
"page two targets the CONTINUE action"
);
assert_eq!(
input,
&json!({"cursor": "cur-2"}),
"the continue action declares `cursor` as its only property, so \
anything else here is a 400 on the real wire"
);
}
#[tokio::test]
async fn shared_links_pages_through_its_own_action_with_the_full_input() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
if req.path == "/v1/actions/dropbox.list_shared_links" {
let body: Value = serde_json::from_str(&req.body).unwrap_or_default();
let page = match body["input"].get("cursor").and_then(Value::as_str) {
None => json!({"links": [entry("l1")], "cursor": "sl-2", "hasMore": true}),
Some("sl-2") => {
json!({"links": [entry("l2")], "cursor": null, "hasMore": false})
}
Some(other) => return MockResponse::new(400, format!("bad cursor {other}")),
};
return MockResponse::ok(&envelope_ok(&page.to_string()));
}
MockResponse::new(404, "{}")
})
.await;
let _token = EnvVarGuard::set("SKARDI_TEST_OC_DROPBOX_SHARED_LINKS", "test-token");
let config: OpenConnectorConfig = serde_yaml::from_str(
r#"
runtime_token_env: SKARDI_TEST_OC_DROPBOX_SHARED_LINKS
bindings:
- name: ws
source_pack: dropbox
resource:
path: /Redacted Folder/redacted-report.pdf
directOnly: true
tables: [shared_links]
"#,
)
.expect("config parses");
let mut ctx = SessionContext::new();
let gateways = OpenConnectorGateways::default();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&config),
false,
HierarchyLevel::Catalog,
Some(&gateways),
)
.await
.expect("registration succeeds");
let batches = collect(&ctx, "SELECT name FROM saas.ws.shared_links ORDER BY name").await;
assert_eq!(names_of(&batches), vec!["l1", "l2"]);
let calls = execute_calls(&gateway);
assert_eq!(calls.len(), 2);
for (path, _) in &calls {
assert_eq!(
path, "/v1/actions/dropbox.list_shared_links",
"both pages use the same action"
);
}
assert_eq!(
calls[0].1,
json!({"path": "/Redacted Folder/redacted-report.pdf", "directOnly": true}),
"page one carries both resources and no cursor"
);
assert_eq!(
calls[1].1,
json!({"path": "/Redacted Folder/redacted-report.pdf", "directOnly": true,
"cursor": "sl-2"}),
"page two REPEATS both resources alongside the cursor — this is what \
`inputs: cursor_only` would have removed"
);
assert!(
calls[0].1.get("limit").is_none() && calls[1].1.get("limit").is_none(),
"this action declares no page-size input"
);
}
#[tokio::test]
async fn file_search_forwards_its_required_query_and_pinned_status() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
if req.path == "/v1/actions/dropbox.search_files" {
return MockResponse::ok(&envelope_ok(
&json!({"matches": [{"matchType": "filename", "metadata": entry("hit"),
"highlightSpans": []}],
"cursor": null, "hasMore": false})
.to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_DROPBOX_SEARCH", "file_search").await;
let batches = collect(&ctx, "SELECT name FROM saas.ws.file_search").await;
assert_eq!(names_of(&batches), vec!["hit"]);
let calls = execute_calls(&gateway);
assert_eq!(calls.len(), 1, "a null cursor ends the scan");
let input = &calls[0].1;
assert_eq!(input["query"], json!("redacted"), "required resource");
assert_eq!(input["fileStatus"], json!("active"), "pinned");
assert_eq!(input["maxResults"], json!(1000));
assert!(
input.get("includeHighlights").is_none(),
"negative space: highlights are never requested"
);
}
#[tokio::test]
async fn a_missing_required_query_fails_before_any_action_call() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
MockResponse::new(404, "{}")
})
.await;
let _token = EnvVarGuard::set("SKARDI_TEST_OC_DROPBOX_NO_QUERY", "test-token");
let config: OpenConnectorConfig = serde_yaml::from_str(
r#"
runtime_token_env: SKARDI_TEST_OC_DROPBOX_NO_QUERY
bindings:
- name: ws
source_pack: dropbox
tables: [file_search]
"#,
)
.expect("config parses");
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&config),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect_err("a search table without a query must not register");
assert!(err.to_string().contains("query"), "{err}");
assert!(
gateway
.requests()
.iter()
.all(|r| r.path == "/v1/health" && r.method == "GET"),
"the resource guard runs before any action is discovered or executed"
);
}
#[tokio::test]
async fn a_drifted_continuation_contract_fails_registration() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
if req.path.ends_with("dropbox.list_folder_continue") {
return MockResponse::ok(&discovery_ok(
"{}",
r#"{"type":"object","properties":{"entries":{"type":"array"}}}"#,
true,
None,
));
}
return dropbox_discovery(&req.path);
}
MockResponse::new(404, "{}")
})
.await;
let _token = EnvVarGuard::set("SKARDI_TEST_OC_DROPBOX_DRIFT", "test-token");
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&dropbox_config("SKARDI_TEST_OC_DROPBOX_DRIFT", "files")),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect_err("a drifted continuation contract must fail registration");
let rendered = err.to_string();
assert!(
rendered.contains("dropbox.files"),
"names the table: {rendered}"
);
assert!(
rendered.contains("dropbox.list_folder_continue"),
"names the CONTINUATION action, not the opener: {rendered}"
);
}
#[tokio::test]
async fn limit_pushdown_stops_the_scan_before_the_continue_action() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
if req.path == "/v1/actions/dropbox.list_folder" {
return MockResponse::ok(&envelope_ok(
&json!({"entries": [entry("a"), entry("b")],
"cursor": "cur-2", "hasMore": true})
.to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_DROPBOX_LIMIT", "files").await;
let batches = collect(&ctx, "SELECT name FROM saas.ws.files LIMIT 1").await;
assert_eq!(names_of(&batches).len(), 1);
let calls = execute_calls(&gateway);
assert_eq!(
calls.len(),
1,
"the LIMIT was satisfied on page one; no continuation fetch"
);
}
#[tokio::test]
async fn a_gateway_failure_surfaces_the_providers_code() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
if req.path == "/v1/actions/dropbox.list_folder" {
return MockResponse::new(
409,
crate::sources::providers::open_connector::testutil::envelope_err(
"provider_error",
"path/not_found",
),
);
}
MockResponse::new(404, "{}")
})
.await;
let (_gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_DROPBOX_ERR", "files").await;
let err = ctx
.sql("SELECT name FROM saas.ws.files")
.await
.expect("plan")
.collect()
.await
.expect_err("a provider failure must fail the scan");
let rendered = err.to_string();
assert!(
rendered.contains("path/not_found") || rendered.contains("provider_error"),
"the provider's own code must surface: {rendered}"
);
assert!(
table("files").error_path.is_none(),
"no in-band error path is declared for Dropbox"
);
}
#[tokio::test]
async fn file_search_pages_through_its_own_continue_action() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
match req.path.as_str() {
"/v1/actions/dropbox.search_files" => MockResponse::ok(&envelope_ok(
&json!({"matches": [
{"matchType": "filename", "metadata": entry("s1"),
"highlightSpans": []}],
"cursor": "s-2", "hasMore": true})
.to_string(),
)),
"/v1/actions/dropbox.search_files_continue" => {
let body: Value = serde_json::from_str(&req.body).unwrap_or_default();
match body["input"].get("cursor").and_then(Value::as_str) {
Some("s-2") => MockResponse::ok(&envelope_ok(
&json!({"matches": [
{"matchType": "content", "metadata": entry("s2"),
"highlightSpans": []}],
"cursor": null, "hasMore": false})
.to_string(),
)),
other => MockResponse::new(400, format!("bad cursor {other:?}")),
}
}
_ => MockResponse::new(404, "{}"),
}
})
.await;
let (gateway, ctx) = setup_with_gateway(
gateway,
"SKARDI_TEST_OC_DROPBOX_SEARCH_PAGES",
"file_search",
)
.await;
let batches = collect(&ctx, "SELECT name FROM saas.ws.file_search ORDER BY name").await;
assert_eq!(names_of(&batches), vec!["s1", "s2"]);
let calls = execute_calls(&gateway);
assert_eq!(calls.len(), 2, "two pages");
let (path, input) = &calls[0];
assert_eq!(path, "/v1/actions/dropbox.search_files");
let mut keys: Vec<&str> = input
.as_object()
.expect("object")
.keys()
.map(String::as_str)
.collect();
keys.sort_unstable();
assert_eq!(
keys,
vec!["fileStatus", "maxResults", "query"],
"page one carries EXACTLY the required resource, the pin and the page size"
);
assert_eq!(
input["maxResults"],
json!(1000),
"this table's own page size"
);
let (path, input) = &calls[1];
assert_eq!(
path, "/v1/actions/dropbox.search_files_continue",
"page two targets THIS table's continue action"
);
assert_eq!(
input,
&json!({"cursor": "s-2"}),
"no query, no fileStatus, no maxResults — the continue action declares none of them"
);
}
#[tokio::test]
async fn shared_links_forwards_both_of_its_optional_resources() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
if req.path == "/v1/actions/dropbox.list_shared_links" {
return MockResponse::ok(&envelope_ok(
&json!({"links": [entry("l1")], "cursor": null, "hasMore": false}).to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let _token = EnvVarGuard::set("SKARDI_TEST_OC_DROPBOX_RESOURCES", "test-token");
let config: OpenConnectorConfig = serde_yaml::from_str(
r#"
runtime_token_env: SKARDI_TEST_OC_DROPBOX_RESOURCES
bindings:
- name: ws
source_pack: dropbox
resource:
path: /Redacted Folder/redacted-report.pdf
directOnly: true
tables: [shared_links]
"#,
)
.expect("config parses");
let mut ctx = SessionContext::new();
let gateways = OpenConnectorGateways::default();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&config),
false,
HierarchyLevel::Catalog,
Some(&gateways),
)
.await
.expect("registration succeeds");
let batches = collect(&ctx, "SELECT name FROM saas.ws.shared_links").await;
assert_eq!(names_of(&batches), vec!["l1"]);
let calls = execute_calls(&gateway);
assert_eq!(calls.len(), 1);
assert_eq!(
calls[0].1,
json!({"path": "/Redacted Folder/redacted-report.pdf", "directOnly": true}),
"both resources forward verbatim, and directOnly is a JSON boolean"
);
}
#[tokio::test]
async fn a_files_page_without_the_has_more_signal_fails_the_scan() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
if req.path == "/v1/actions/dropbox.list_folder" {
return MockResponse::ok(&envelope_ok(
&json!({"entries": [entry("a")], "cursor": "cur-2"}).to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_DROPBOX_NO_HAS_MORE", "files").await;
let err = ctx
.sql("SELECT name FROM saas.ws.files")
.await
.expect("plan")
.collect()
.await
.expect_err("a page missing the declared has-more signal must fail");
let rendered = err.to_string();
assert!(
rendered.contains("$.hasMore") && rendered.contains("absent"),
"the error names the declared path and what was found: {rendered}"
);
assert_eq!(
execute_calls(&gateway).len(),
1,
"the scan stops at the drifted page instead of continuing"
);
}
#[tokio::test]
async fn a_drifted_opening_contract_fails_registration_for_every_table() {
for (short, action) in [
("files", "dropbox.list_folder"),
("shared_links", "dropbox.list_shared_links"),
("file_search", "dropbox.search_files"),
] {
let drifted = action.to_string();
let gateway = MockGateway::start(move |req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
if req.path.ends_with(&drifted) {
return MockResponse::ok(&discovery_ok(
"{}",
r#"{"type":"object","properties":{"drifted":{"type":"string"}}}"#,
true,
None,
));
}
return dropbox_discovery(&req.path);
}
MockResponse::new(404, "{}")
})
.await;
let env = format!("SKARDI_TEST_OC_DROPBOX_DRIFT_{}", short.to_uppercase());
let _token = EnvVarGuard::set(&env, "test-token");
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&dropbox_config(&env, short)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect_err("a drifted opening contract must fail registration");
let rendered = err.to_string();
assert!(
rendered.contains(&format!("dropbox.{short}")) && rendered.contains(action),
"{short}: error must name the table and the drifted action: {rendered}"
);
}
}
#[tokio::test]
async fn a_continue_action_demanding_more_than_a_cursor_fails_registration() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
if req.path.ends_with("dropbox.list_folder_continue") {
return MockResponse::ok(&discovery_ok(
r#"{"type":"object",
"properties":{"cursor":{"type":"string"},"path":{"type":"string"}},
"required":["cursor","path"],
"additionalProperties":false}"#,
include_str!("fixtures/dropbox/contracts/list_folder_continue.json"),
true,
None,
));
}
return dropbox_discovery(&req.path);
}
MockResponse::new(404, "{}")
})
.await;
let _token = EnvVarGuard::set("SKARDI_TEST_OC_DROPBOX_CURSOR_ONLY", "test-token");
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&dropbox_config(
"SKARDI_TEST_OC_DROPBOX_CURSOR_ONLY",
"files",
)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect_err("a continue action requiring more than the cursor must fail registration");
let rendered = err.to_string();
assert!(
rendered.contains("dropbox.list_folder_continue") && rendered.contains("path"),
"the error names the continuation action and the unsatisfiable input: {rendered}"
);
}
#[tokio::test]
async fn a_continue_action_without_an_input_schema_is_refused() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
if req.path.ends_with("dropbox.list_folder_continue") {
return MockResponse::ok(&discovery_ok(
"null",
include_str!("fixtures/dropbox/contracts/list_folder_continue.json"),
true,
None,
));
}
return dropbox_discovery(&req.path);
}
MockResponse::new(404, "{}")
})
.await;
let _token = EnvVarGuard::set("SKARDI_TEST_OC_DROPBOX_NO_INPUT_SCHEMA", "test-token");
let mut ctx = SessionContext::new();
let err = register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&dropbox_config(
"SKARDI_TEST_OC_DROPBOX_NO_INPUT_SCHEMA",
"files",
)),
false,
HierarchyLevel::Catalog,
None,
)
.await
.expect_err("an unverifiable cursor_only claim must be refused");
assert!(
err.to_string().contains("no input schema"),
"the error says why it cannot be verified: {err}"
);
}
#[tokio::test]
async fn udtf_parity_for_files_crosses_the_split_action_boundary() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
match req.path.as_str() {
"/v1/actions/dropbox.list_folder" => MockResponse::ok(&envelope_ok(
&json!({"entries": [entry("via-udtf")], "cursor": "c-2", "hasMore": true})
.to_string(),
)),
"/v1/actions/dropbox.list_folder_continue" => {
let body: Value = serde_json::from_str(&req.body).unwrap_or_default();
match body["input"].get("cursor").and_then(Value::as_str) {
Some("c-2") => MockResponse::ok(&envelope_ok(
&json!({"entries": [entry("via-udtf-page-2")],
"cursor": null, "hasMore": false})
.to_string(),
)),
other => MockResponse::new(400, format!("bad cursor {other:?}")),
}
}
_ => MockResponse::new(404, "{}"),
}
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_DROPBOX_UDTF", "files").await;
let batches = collect(
&ctx,
"SELECT name FROM open_connector_query('saas', 'dropbox.files', '{}')",
)
.await;
assert_eq!(names_of(&batches), vec!["via-udtf", "via-udtf-page-2"]);
let calls = execute_calls(&gateway);
let (path, input) = calls
.iter()
.find(|(path, _)| path.ends_with("dropbox.list_folder_continue"))
.expect("the UDTF path issued a continuation request");
assert!(path.ends_with("dropbox.list_folder_continue"));
assert_eq!(*input, json!({"cursor": "c-2"}), "cursor only, on page two");
}
fn dropbox_allowlist_config(token_env: &str, allowlist: &[&str]) -> OpenConnectorConfig {
let allowlist = allowlist.join(", ");
serde_yaml::from_str(&format!(
r#"
runtime_token_env: {token_env}
raw_action_allowlist: [{allowlist}]
bindings:
- name: ws
source_pack: dropbox
tables: [shared_links]
"#
))
.expect("config parses")
}
async fn allowlist_only_ctx(
gateway: &MockGateway,
token_env: &'static str,
allowlist: &[&str],
) -> SessionContext {
let _token = EnvVarGuard::set(token_env, "test-token");
let gateways = OpenConnectorGateways::default();
let mut ctx = SessionContext::new();
register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&dropbox_allowlist_config(token_env, allowlist)),
false,
HierarchyLevel::Catalog,
Some(&gateways),
)
.await
.expect("registration gates only the BOUND table, so it succeeds");
register_open_connector_udtfs(&ctx, gateways).expect("UDTF registration succeeds");
ctx
}
async fn udtf_error(ctx: &SessionContext, sql: &str) -> String {
match ctx.sql(sql).await {
Err(e) => e.to_string(),
Ok(df) => df
.collect()
.await
.expect_err("the query must not succeed")
.to_string(),
}
}
#[tokio::test]
async fn allowlisting_only_the_opening_action_fails_udtf_planning() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
return dropbox_discovery(&req.path);
}
MockResponse::new(404, "{}")
})
.await;
let ctx = allowlist_only_ctx(
&gateway,
"SKARDI_TEST_OC_DROPBOX_UDTF_ALLOWLIST",
&["dropbox.list_folder"],
)
.await;
let rendered = udtf_error(
&ctx,
"SELECT name FROM open_connector_query('saas', 'dropbox.files', '{}')",
)
.await;
assert!(
rendered.contains("was not discovered"),
"the undiscovered-action diagnostic: {rendered}"
);
assert!(
rendered.contains("dropbox.list_folder_continue"),
"names the CONTINUATION action, not the opener: {rendered}"
);
}
#[tokio::test]
async fn a_drifted_continuation_contract_fails_udtf_planning() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
if req.path.ends_with("dropbox.list_folder_continue") {
return MockResponse::ok(&discovery_ok(
"{}",
r#"{"type":"object","properties":{"entries":{"type":"array"}}}"#,
true,
None,
));
}
return dropbox_discovery(&req.path);
}
MockResponse::new(404, "{}")
})
.await;
let ctx = allowlist_only_ctx(
&gateway,
"SKARDI_TEST_OC_DROPBOX_UDTF_DRIFT",
&["dropbox.list_folder", "dropbox.list_folder_continue"],
)
.await;
let rendered = udtf_error(
&ctx,
"SELECT name FROM open_connector_query('saas', 'dropbox.files', '{}')",
)
.await;
assert!(
rendered.contains("fingerprint mismatch"),
"the contract gate refused it: {rendered}"
);
assert!(
rendered.contains("dropbox.list_folder_continue"),
"names the CONTINUATION action, not the opener: {rendered}"
);
}
#[tokio::test]
async fn an_unverifiable_cursor_only_claim_fails_udtf_planning() {
let gateway = MockGateway::start(|req| {
if req.method == "GET" && req.path == "/v1/health" {
return MockResponse::ok("{}");
}
if req.method == "GET" && req.path.starts_with("/v1/actions/") {
if req.path.ends_with("dropbox.list_folder_continue") {
return MockResponse::ok(&discovery_ok(
"{}",
include_str!("fixtures/dropbox/contracts/list_folder_continue.json"),
true,
None,
));
}
return dropbox_discovery(&req.path);
}
MockResponse::new(404, "{}")
})
.await;
let ctx = allowlist_only_ctx(
&gateway,
"SKARDI_TEST_OC_DROPBOX_UDTF_INPUTS",
&["dropbox.list_folder", "dropbox.list_folder_continue"],
)
.await;
let rendered = udtf_error(
&ctx,
"SELECT name FROM open_connector_query('saas', 'dropbox.files', '{}')",
)
.await;
assert!(
rendered.contains("no input properties"),
"the input gate said why it cannot be verified: {rendered}"
);
assert!(
rendered.contains("dropbox.list_folder_continue"),
"names the continuation action: {rendered}"
);
}
}