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("gmail.yaml", include_str!("gmail.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::pagination::PaginationStrategy;
use crate::sources::providers::open_connector::row_path::RowPath;
use crate::sources::providers::open_connector::source_pack::{FixedValue, SourcePackTable};
use crate::sources::providers::open_connector::testutil::{
EnvVarGuard, MockGateway, MockResponse, collect, column_values, convert_first_page,
discovery_ok, envelope_err, envelope_ok, execute_inputs, fingerprint_uncovered_columns,
input_keys, utf8,
};
use crate::sources::providers::open_connector::{
OpenConnectorConfig, OpenConnectorGateways, register_open_connector_tables,
register_open_connector_udtfs,
};
use arrow::array::{Array, ListArray, 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 gmail_discovery(path: &str) -> MockResponse {
let output_schema = if path.ends_with("gmail.list_threads") {
include_str!("fixtures/gmail/contracts/list_threads.json")
} else if path.ends_with("gmail.fetch_emails") {
include_str!("fixtures/gmail/contracts/fetch_emails.json")
} else if path.ends_with("gmail.list_drafts") {
include_str!("fixtures/gmail/contracts/list_drafts.json")
} else if path.ends_with("gmail.list_labels") {
include_str!("fixtures/gmail/contracts/list_labels.json")
} else if path.ends_with("gmail.list_filters") {
include_str!("fixtures/gmail/contracts/list_filters.json")
} else {
r#"{"type": "object"}"#
};
MockResponse::ok(&discovery_ok("{}", output_schema, true, None))
}
fn convert_fixture(table: &SourcePackTable, fixture: &str) -> RecordBatch {
let page: Value = serde_json::from_str(fixture).expect("fixture parses");
convert_first_page(table, &page)
}
#[test]
fn threads_fixture_converts_the_live_page_shape() {
let batch = convert_fixture(
table("threads"),
include_str!("fixtures/gmail/threads.json"),
);
assert_eq!(batch.num_rows(), 3);
assert_eq!(utf8(&batch, "thread_id").value(0), "1900000000000a00");
assert_eq!(utf8(&batch, "history_id").value(0), "4200001");
assert!((0..3).all(|i| !utf8(&batch, "snippet").is_null(i)));
}
#[test]
fn messages_fixture_converts_the_live_summary_shape() {
let batch = convert_fixture(
table("messages"),
include_str!("fixtures/gmail/messages.json"),
);
assert_eq!(batch.num_rows(), 3);
assert_eq!(utf8(&batch, "message_id").value(0), "1900000000000b00");
let labels: &ListArray = batch
.column_by_name("label_ids")
.expect("column")
.as_any()
.downcast_ref()
.expect("List column");
let third = labels.value(2);
let third = third
.as_any()
.downcast_ref::<StringArray>()
.expect("Utf8 items");
assert_eq!(
(0..third.len()).map(|i| third.value(i)).collect::<Vec<_>>(),
vec!["IMPORTANT", "CATEGORY_UPDATES", "INBOX"]
);
let ts: &TimestampMillisecondArray = batch
.column_by_name("message_timestamp")
.expect("column")
.as_any()
.downcast_ref()
.expect("timestamp");
assert!((0..3).all(|i| !ts.is_null(i)));
}
#[test]
fn drafts_fixture_converts_through_nested_paths() {
let batch = convert_fixture(table("drafts"), include_str!("fixtures/gmail/drafts.json"));
assert_eq!(batch.num_rows(), 1);
assert_eq!(utf8(&batch, "id").value(0), "r-4000000000000000001");
assert_eq!(utf8(&batch, "message_id").value(0), "1900000000000c01");
assert_eq!(utf8(&batch, "thread_id").value(0), "1900000000000c01");
}
#[test]
fn executor_absent_spellings_convert_as_pinned() {
let batch = convert_first_page(
table("threads"),
&json!({"threads": [
{"threadId": "t-1", "snippet": "", "historyId": null, "extra": true},
]}),
);
assert!(utf8(&batch, "history_id").is_null(0));
assert_eq!(utf8(&batch, "snippet").value(0), "");
let batch = convert_first_page(
table("messages"),
&json!({"messages": [{
"messageId": "m-1", "threadId": "t-1", "labelIds": [],
"subject": "", "sender": "", "to": "",
"messageTimestamp": "2026-07-30T08:15:42.000Z",
}]}),
);
assert_eq!(utf8(&batch, "subject").value(0), "");
assert_eq!(utf8(&batch, "to_addresses").value(0), "");
let labels: &ListArray = batch
.column_by_name("label_ids")
.expect("column")
.as_any()
.downcast_ref()
.expect("List column");
assert!(!labels.is_null(0));
assert_eq!(labels.value(0).len(), 0);
let batch = convert_first_page(
table("drafts"),
&json!({"drafts": [{"id": "r-1", "message": {"messageId": "", "threadId": ""}}]}),
);
assert_eq!(utf8(&batch, "message_id").value(0), "");
}
#[test]
fn null_parent_on_a_nested_path_becomes_sql_null() {
let batch = convert_first_page(
table("drafts"),
&json!({"drafts": [{"id": "r-1", "message": null}]}),
);
assert!(utf8(&batch, "message_id").is_null(0));
assert!(utf8(&batch, "thread_id").is_null(0));
}
#[test]
fn labels_fixture_converts_with_absent_visibility_and_color() {
let batch = convert_fixture(table("labels"), include_str!("fixtures/gmail/labels.json"));
assert_eq!(batch.num_rows(), 15);
assert_eq!(utf8(&batch, "id").value(0), "CHAT");
assert_eq!(utf8(&batch, "message_list_visibility").value(0), "hide");
assert_eq!(utf8(&batch, "id").value(1), "SENT");
assert!(utf8(&batch, "message_list_visibility").is_null(1));
assert!(utf8(&batch, "label_list_visibility").is_null(1));
assert_eq!(utf8(&batch, "id").value(2), "INBOX");
assert!(utf8(&batch, "message_list_visibility").is_null(2));
assert!(utf8(&batch, "color").is_null(0));
assert_eq!(utf8(&batch, "type").value(14), "user");
let color: Value =
serde_json::from_str(utf8(&batch, "color").value(14)).expect("valid JSON");
assert_eq!(color["backgroundColor"], "#fb4c2f");
}
#[test]
fn filters_fixture_converts_with_opaque_json() {
let batch = convert_fixture(
table("filters"),
include_str!("fixtures/gmail/filters.json"),
);
assert_eq!(batch.num_rows(), 1);
let criteria: Value =
serde_json::from_str(utf8(&batch, "criteria").value(0)).expect("valid JSON");
assert_eq!(criteria["from"], "alerts@example.com");
let action: Value =
serde_json::from_str(utf8(&batch, "action").value(0)).expect("valid JSON");
assert_eq!(action["addLabelIds"][0], "Label_1");
}
#[test]
fn messages_mismatch_fixture_fails_with_the_targeted_error() {
let page: Value =
serde_json::from_str(include_str!("fixtures/gmail/messages_type_mismatch.json"))
.expect("fixture parses");
let t = table("messages");
let rows = RowPath::parse(t.row_path)
.expect("row path")
.rows(&page, 1)
.expect("row array");
let err = RowConverter::new(t.fields)
.expect("converter")
.convert(rows, 1)
.expect_err("a number where Utf8 is declared must fail conversion");
match err {
OpenConnectorError::ConversionFailed {
column,
page,
row,
found,
..
} => {
assert_eq!(column, "message_id");
assert_eq!(page, 1);
assert_eq!(row, 1, "the valid first row converts");
assert_eq!(found, "a number");
}
other => panic!("expected ConversionFailed, got {other}"),
}
}
#[test]
fn pinned_fingerprints_match_the_reconciled_contracts() {
let contracts = [
(
"threads",
include_str!("fixtures/gmail/contracts/list_threads.json"),
),
(
"messages",
include_str!("fixtures/gmail/contracts/fetch_emails.json"),
),
(
"drafts",
include_str!("fixtures/gmail/contracts/list_drafts.json"),
),
(
"labels",
include_str!("fixtures/gmail/contracts/list_labels.json"),
),
(
"filters",
include_str!("fixtures/gmail/contracts/list_filters.json"),
),
];
let mut mismatches = Vec::new();
for (short, contract) in contracts {
let schema: Value = serde_json::from_str(contract).expect("contract fixture parses");
let actual = fingerprint_schema(Some(&schema));
let t = table(short);
if t.expected_fingerprint != Some(actual.as_str()) {
mismatches.push(format!(
"{}: pinned {:?}, contract fixture hashes to {actual}",
t.id, t.expected_fingerprint
));
}
}
assert!(mismatches.is_empty(), "{}", mismatches.join("\n"));
}
#[test]
fn generated_inputs_are_accepted_by_the_captured_input_contracts() {
for (short, contract) in [
(
"threads",
include_str!("fixtures/gmail/contracts/inputs/list_threads.json"),
),
(
"messages",
include_str!("fixtures/gmail/contracts/inputs/fetch_emails.json"),
),
(
"drafts",
include_str!("fixtures/gmail/contracts/inputs/list_drafts.json"),
),
(
"labels",
include_str!("fixtures/gmail/contracts/inputs/list_labels.json"),
),
(
"filters",
include_str!("fixtures/gmail/contracts/inputs/list_filters.json"),
),
] {
let schema: Value =
serde_json::from_str(contract).expect("input contract fixture parses");
let properties = &schema["properties"];
let t = table(short);
assert_eq!(
schema["additionalProperties"],
json!(false),
"{short}: the action's input schema is strict"
);
let mut generated: Vec<&str> = t
.required_resources
.iter()
.chain(t.optional_resources)
.copied()
.collect();
generated.extend(t.fixed_inputs.iter().map(|(key, _)| *key));
match t.pagination {
PaginationStrategy::Cursor {
cursor_param,
page_size_param,
..
} => {
generated.push(cursor_param);
generated.extend(page_size_param);
}
PaginationStrategy::PageNumber {
page_param,
per_page_param,
..
} => {
generated.push(page_param);
generated.push(per_page_param);
}
PaginationStrategy::Keyset {
cursor_param,
page_size_param,
..
} => {
generated.push(cursor_param);
generated.push(page_size_param);
}
PaginationStrategy::SinglePage { .. } => {}
}
for key in &generated {
assert!(
!properties[*key].is_null(),
"{short}: `{key}` is not declared by the action's input schema"
);
}
if let Some(required) = schema["required"].as_array() {
for entry in required {
let entry = entry.as_str().expect("required entries are strings");
assert!(
generated.contains(&entry),
"{short}: the action requires `{entry}`, which this table never sends"
);
}
}
if let PaginationStrategy::Cursor {
page_size_param: Some(param),
page_size,
..
} = t.pagination
{
let declared = &properties[param];
if let Some(minimum) = declared["minimum"].as_u64() {
assert!(
u64::from(page_size) >= minimum,
"{short}: page size {page_size} is below `{param}`'s minimum {minimum}"
);
}
if let Some(maximum) = declared["maximum"].as_u64() {
assert!(
u64::from(page_size) <= maximum,
"{short}: page size {page_size} exceeds `{param}`'s maximum {maximum}"
);
}
}
for (key, value) in t.fixed_inputs {
if let (FixedValue::Str(value), Some(variants)) =
(value, properties[*key]["enum"].as_array())
{
assert!(
variants.iter().any(|v| v.as_str() == Some(*value)),
"{short}: fixed input `{key}: {value}` is outside the declared enum"
);
}
}
}
}
#[test]
fn columns_the_coverage_walker_cannot_resolve_are_pinned() {
for (short, contract, expected) in [
(
"threads",
include_str!("fixtures/gmail/contracts/list_threads.json"),
&[] as &[&str],
),
(
"messages",
include_str!("fixtures/gmail/contracts/fetch_emails.json"),
&[
"message_id",
"thread_id",
"label_ids",
"subject",
"sender",
"to_addresses",
"message_timestamp",
],
),
(
"drafts",
include_str!("fixtures/gmail/contracts/list_drafts.json"),
&[],
),
(
"labels",
include_str!("fixtures/gmail/contracts/list_labels.json"),
&[],
),
(
"filters",
include_str!("fixtures/gmail/contracts/list_filters.json"),
&[],
),
] {
let t = table(short);
assert_eq!(
fingerprint_uncovered_columns(contract, t.row_path, t.fields),
expected,
"fingerprint coverage changed for {short}"
);
}
}
fn gmail_config(token_env: &str, tables: &str, resource: &str) -> OpenConnectorConfig {
let resource_line = if resource.is_empty() {
String::new()
} else {
format!("resource: {resource}")
};
serde_yaml::from_str(&format!(
r#"
runtime_token_env: {token_env}
bindings:
- name: mail
source_pack: gmail
{resource_line}
tables: [{tables}]
"#
))
.expect("config parses")
}
async fn setup_with_gateway(
gateway: MockGateway,
token_env: &'static str,
tables: &str,
resource: &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(&gmail_config(token_env, tables, resource)),
false,
HierarchyLevel::Catalog,
Some(&gateways),
)
.await
.expect("gateway registration succeeds");
register_open_connector_udtfs(&ctx, gateways).expect("UDTF registration succeeds");
(gateway, ctx)
}
fn thread_row(id: &str) -> Value {
json!({"threadId": id, "snippet": "", "historyId": null})
}
fn message_row(id: &str) -> Value {
json!({
"messageId": id,
"threadId": format!("t-{id}"),
"labelIds": ["INBOX"],
"subject": "s",
"sender": "a@example.com",
"to": "b@example.com",
"messageTimestamp": "2026-07-30T08:15:42.000Z"
})
}
#[tokio::test]
async fn threads_cursor_scan_pages_with_its_own_declared_inputs() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.list_threads" {
let body: Value = serde_json::from_str(&req.body).unwrap_or_default();
let page = match body["input"].get("pageToken").and_then(Value::as_str) {
None => json!({"threads": [thread_row("t-1"), thread_row("t-2")],
"nextPageToken": "tok-2", "resultSizeEstimate": 3}),
Some("tok-2") => json!({"threads": [thread_row("t-3")],
"nextPageToken": null, "resultSizeEstimate": 3}),
Some(other) => return MockResponse::new(400, format!("bad token {other}")),
};
return MockResponse::ok(&envelope_ok(&page.to_string()));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_GMAIL_THREADS", "threads", "").await;
let batches = collect(
&ctx,
"SELECT thread_id FROM saas.mail.threads ORDER BY thread_id",
)
.await;
assert_eq!(
column_values(&batches, "thread_id"),
vec!["t-1", "t-2", "t-3"]
);
let inputs = execute_inputs(&gateway, "gmail.list_threads");
assert_eq!(inputs.len(), 2, "two cursor pages");
assert_eq!(inputs[1]["pageToken"], "tok-2");
for (page, (input, expected_keys)) in inputs
.iter()
.zip([vec!["maxResults"], vec!["maxResults", "pageToken"]])
.enumerate()
{
assert_eq!(input["maxResults"], 500, "page-size hint: {input}");
assert_eq!(input_keys(input), expected_keys, "page {} keys", page + 1);
}
}
#[tokio::test]
async fn messages_scan_pins_detail_summary_on_every_page() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.fetch_emails" {
let body: Value = serde_json::from_str(&req.body).unwrap_or_default();
let page = match body["input"].get("pageToken").and_then(Value::as_str) {
None => json!({"messages": [message_row("m-1")],
"nextPageToken": "tok-2", "resultSizeEstimate": 2}),
Some("tok-2") => json!({"messages": [message_row("m-2")],
"resultSizeEstimate": 2}),
Some(other) => return MockResponse::new(400, format!("bad token {other}")),
};
return MockResponse::ok(&envelope_ok(&page.to_string()));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_GMAIL_MESSAGES", "messages", "").await;
let batches = collect(
&ctx,
"SELECT message_id FROM saas.mail.messages ORDER BY message_id",
)
.await;
assert_eq!(column_values(&batches, "message_id"), vec!["m-1", "m-2"]);
let inputs = execute_inputs(&gateway, "gmail.fetch_emails");
assert_eq!(inputs.len(), 2, "two cursor pages");
for (page, (input, expected_keys)) in inputs
.iter()
.zip([
vec!["detail", "maxResults"],
vec!["detail", "maxResults", "pageToken"],
])
.enumerate()
{
assert_eq!(input["detail"], "summary", "fixed input pin: {input}");
assert_eq!(input["maxResults"], 100, "bounded page size: {input}");
assert_eq!(input_keys(input), expected_keys, "page {} keys", page + 1);
}
}
#[tokio::test]
async fn drafts_cursor_scan_pages_with_its_own_declared_inputs() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.list_drafts" {
let body: Value = serde_json::from_str(&req.body).unwrap_or_default();
let page = match body["input"].get("pageToken").and_then(Value::as_str) {
None => json!({"drafts": [
{"id": "d-1", "message": {"messageId": "m-1", "threadId": "t-1"}}],
"nextPageToken": "tok-2"}),
Some("tok-2") => json!({"drafts": [
{"id": "d-2", "message": {"messageId": "m-2", "threadId": "t-2"}}],
"nextPageToken": ""}),
Some(other) => return MockResponse::new(400, format!("bad token {other}")),
};
return MockResponse::ok(&envelope_ok(&page.to_string()));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_GMAIL_DRAFTS", "drafts", "").await;
let batches = collect(
&ctx,
"SELECT id, message_id FROM saas.mail.drafts ORDER BY id",
)
.await;
assert_eq!(column_values(&batches, "id"), vec!["d-1", "d-2"]);
assert_eq!(column_values(&batches, "message_id"), vec!["m-1", "m-2"]);
let inputs = execute_inputs(&gateway, "gmail.list_drafts");
assert_eq!(inputs.len(), 2, "empty-string token terminates");
assert_eq!(input_keys(&inputs[0]), vec!["maxResults"]);
assert_eq!(inputs[0]["maxResults"], 500);
assert_eq!(input_keys(&inputs[1]), vec!["maxResults", "pageToken"]);
}
#[tokio::test]
async fn single_page_tables_issue_exactly_one_request_with_no_inputs() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.list_labels" {
return MockResponse::ok(&envelope_ok(
&json!({"labels": [
{"id": "INBOX", "name": "INBOX", "type": "system"},
{"id": "Label_1", "name": "P/Redacted", "type": "user"}],
"nextPageToken": null})
.to_string(),
));
}
if req.method == "POST" && req.path == "/v1/actions/gmail.list_filters" {
return MockResponse::ok(&envelope_ok(
&json!({"filters": [{"id": "f-1", "criteria": {"from": "x@example.com"},
"action": {"addLabelIds": ["Label_1"]}}]})
.to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) = setup_with_gateway(
gateway,
"SKARDI_TEST_OC_GMAIL_SINGLE",
"labels, filters",
"",
)
.await;
let batches = collect(&ctx, "SELECT id FROM saas.mail.labels ORDER BY id").await;
assert_eq!(column_values(&batches, "id"), vec!["INBOX", "Label_1"]);
let batches = collect(&ctx, "SELECT id, criteria, action FROM saas.mail.filters").await;
assert_eq!(column_values(&batches, "id"), vec!["f-1"]);
for action in ["gmail.list_labels", "gmail.list_filters"] {
let inputs = execute_inputs(&gateway, action);
assert_eq!(inputs.len(), 1, "{action}: single page means one request");
assert_eq!(
input_keys(&inputs[0]),
Vec::<&str>::new(),
"{action}: empty input object"
);
}
}
#[tokio::test]
async fn a_single_page_table_refuses_a_live_continuation_token() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.list_labels" {
return MockResponse::ok(&envelope_ok(
&json!({"labels": [{"id": "INBOX", "name": "INBOX", "type": "system"}],
"nextPageToken": "tok-2"})
.to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_GMAIL_SHORT", "labels", "").await;
let err = ctx
.sql("SELECT id FROM saas.mail.labels")
.await
.expect("plan")
.collect()
.await
.expect_err("a continuation token must fail the single-page scan");
let message = err.to_string();
assert!(
message.contains("single-page") && message.contains("$.nextPageToken"),
"the broken premise and its path are named: {message}"
);
assert_eq!(
execute_inputs(&gateway, "gmail.list_labels").len(),
1,
"the refusal happens after one request, never by following the token"
);
}
#[tokio::test]
async fn single_page_tables_scan_an_empty_collection_cleanly() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" {
let empty = match req.path.as_str() {
"/v1/actions/gmail.list_labels" => json!({"labels": []}),
"/v1/actions/gmail.list_filters" => json!({"filters": []}),
_ => return MockResponse::new(404, "{}"),
};
return MockResponse::ok(&envelope_ok(&empty.to_string()));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) = setup_with_gateway(
gateway,
"SKARDI_TEST_OC_GMAIL_EMPTY_SINGLE",
"labels, filters",
"",
)
.await;
for table in ["labels", "filters"] {
let batches = collect(&ctx, &format!("SELECT id FROM saas.mail.{table}")).await;
assert_eq!(
batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
0,
"{table}: an empty collection is an empty result set"
);
assert_eq!(
execute_inputs(&gateway, &format!("gmail.list_{table}")).len(),
1,
"{table}: an empty page still means exactly one request"
);
}
}
#[tokio::test]
async fn optional_resources_forward_verbatim_and_only_where_declared() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" {
let empty = match req.path.as_str() {
"/v1/actions/gmail.fetch_emails" => {
json!({"messages": [], "resultSizeEstimate": 0, "nextPageToken": null})
}
"/v1/actions/gmail.list_threads" => {
json!({"threads": [], "resultSizeEstimate": 0, "nextPageToken": null})
}
"/v1/actions/gmail.list_drafts" => {
json!({"drafts": [], "nextPageToken": null})
}
_ => return MockResponse::new(404, "{}"),
};
return MockResponse::ok(&envelope_ok(&empty.to_string()));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) = setup_with_gateway(
gateway,
"SKARDI_TEST_OC_GMAIL_RESOURCES",
"threads, messages, drafts",
r#"{ query: "in:inbox -category:promotions", labelIds: [INBOX, IMPORTANT], includeSpamTrash: true }"#,
)
.await;
for sql in [
"SELECT * FROM saas.mail.messages",
"SELECT * FROM saas.mail.threads",
"SELECT * FROM saas.mail.drafts",
] {
let batches = collect(&ctx, sql).await;
assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 0);
}
let inputs = execute_inputs(&gateway, "gmail.fetch_emails");
assert_eq!(
input_keys(&inputs[0]),
vec![
"detail",
"includeSpamTrash",
"labelIds",
"maxResults",
"query"
]
);
assert_eq!(inputs[0]["query"], "in:inbox -category:promotions");
assert_eq!(inputs[0]["labelIds"], json!(["INBOX", "IMPORTANT"]));
assert_eq!(inputs[0]["includeSpamTrash"], json!(true));
let inputs = execute_inputs(&gateway, "gmail.list_threads");
assert_eq!(input_keys(&inputs[0]), vec!["maxResults", "query"]);
let inputs = execute_inputs(&gateway, "gmail.list_drafts");
assert_eq!(input_keys(&inputs[0]), vec!["maxResults"]);
}
#[tokio::test]
async fn predicates_stay_local_against_a_provider_that_cannot_narrow() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.fetch_emails" {
let mut wanted = message_row("m-1");
wanted["subject"] = json!("needle");
return MockResponse::ok(&envelope_ok(
&json!({"messages": [wanted, message_row("m-2"), message_row("m-3")],
"nextPageToken": null, "resultSizeEstimate": 3})
.to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_GMAIL_LOCAL", "messages", "").await;
let batches = collect(
&ctx,
"SELECT message_id FROM saas.mail.messages WHERE subject = 'needle'",
)
.await;
assert_eq!(column_values(&batches, "message_id"), vec!["m-1"]);
let inputs = execute_inputs(&gateway, "gmail.fetch_emails");
assert_eq!(
input_keys(&inputs[0]),
vec!["detail", "maxResults"],
"the predicate stayed local; no subject/query key was pushed"
);
}
#[tokio::test]
async fn limit_stops_cursor_pagination_early() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.list_threads" {
return MockResponse::ok(&envelope_ok(
&json!({"threads": [thread_row("t-1"), thread_row("t-2")],
"nextPageToken": "again", "resultSizeEstimate": 100})
.to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let (gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_GMAIL_LIMIT", "threads", "").await;
let batches = collect(&ctx, "SELECT thread_id FROM saas.mail.threads LIMIT 2").await;
assert_eq!(column_values(&batches, "thread_id").len(), 2);
assert_eq!(
execute_inputs(&gateway, "gmail.list_threads").len(),
1,
"one page satisfied LIMIT"
);
}
#[tokio::test]
async fn a_repeated_cursor_fails_as_a_pagination_loop() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.list_threads" {
return MockResponse::ok(&envelope_ok(
&json!({"threads": [thread_row("t-1")],
"nextPageToken": "stuck", "resultSizeEstimate": 2})
.to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let (_gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_GMAIL_LOOP", "threads", "").await;
let err = ctx
.sql("SELECT thread_id FROM saas.mail.threads")
.await
.expect("plan")
.collect()
.await
.expect_err("a non-advancing cursor must fail the scan");
let message = err.to_string();
assert!(
message.contains("pagination loop") && message.contains("stuck"),
"loop identity is named: {message}"
);
}
#[tokio::test]
async fn provider_errors_surface_through_the_gateway_failure_envelope() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.list_filters" {
return MockResponse::new(
403,
envelope_err(
"authorization_failed",
"Request had insufficient authentication scopes.",
),
);
}
MockResponse::new(404, "{}")
})
.await;
let (_gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_GMAIL_SCOPE", "filters", "").await;
let err = ctx
.sql("SELECT id FROM saas.mail.filters")
.await
.expect("plan")
.collect()
.await
.expect_err("a missing-scope failure must fail the scan");
let message = err.to_string();
assert!(
message.contains("authorization_failed") && message.contains("gmail.list_filters"),
"the gateway's error code and the action are named: {message}"
);
assert!(
!message.contains("row path"),
"never the misleading row-path error: {message}"
);
}
#[tokio::test]
async fn udtf_parity_for_labels() {
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 gmail_discovery(&req.path);
}
if req.method == "POST" && req.path == "/v1/actions/gmail.list_labels" {
return MockResponse::ok(&envelope_ok(
&json!({"labels": [{"id": "INBOX", "name": "INBOX", "type": "system"}]})
.to_string(),
));
}
MockResponse::new(404, "{}")
})
.await;
let (_gateway, ctx) =
setup_with_gateway(gateway, "SKARDI_TEST_OC_GMAIL_UDTF", "labels", "").await;
let from_table = collect(&ctx, "SELECT id, name, type FROM saas.mail.labels").await;
let from_udtf = collect(
&ctx,
"SELECT id, name, type FROM open_connector_query('saas', 'gmail.labels', '{}')",
)
.await;
assert_eq!(from_table[0].schema(), from_udtf[0].schema());
assert_eq!(
arrow::util::pretty::pretty_format_batches(&from_table)
.unwrap()
.to_string(),
arrow::util::pretty::pretty_format_batches(&from_udtf)
.unwrap()
.to_string()
);
}
#[tokio::test]
async fn drifted_contract_fails_registration_not_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 MockResponse::ok(&discovery_ok("{}", r#"{"type": "object"}"#, true, None));
}
MockResponse::new(404, "{}")
})
.await;
let _token = EnvVarGuard::set("SKARDI_TEST_OC_GMAIL_DRIFT", "test-token");
let mut ctx = SessionContext::new();
let gateways = OpenConnectorGateways::default();
let err = register_open_connector_tables(
&mut ctx,
"saas",
&gateway.url,
Some(&gmail_config("SKARDI_TEST_OC_GMAIL_DRIFT", "threads", "")),
false,
HierarchyLevel::Catalog,
Some(&gateways),
)
.await
.expect_err("a drifted contract must fail registration");
let message = err.to_string();
assert!(
message.contains("gmail.threads")
&& message.contains("gmail.list_threads")
&& message.contains("fingerprint mismatch"),
"table, action, and cause are named: {message}"
);
}
}