use std::collections::BTreeMap;
use std::sync::OnceLock;
use serde::Deserialize;
use crate::sources::providers::open_connector::error::OpenConnectorError;
use crate::sources::providers::open_connector::json_to_arrow::RowConverter;
use crate::sources::providers::open_connector::json_to_arrow::{FieldMapping, FieldType};
use crate::sources::providers::open_connector::pagination::{
AbsentCursor, CursorContinuation, PaginationStrategy,
};
use crate::sources::providers::open_connector::row_path::RowPath;
use crate::sources::providers::open_connector::source_pack::{
FixedValue, RowShape, SourcePack, SourcePackTable,
};
use skardi_source_pack::filters::{Fidelity, FilterMapping, Operator, ValueFormat};
pub(crate) fn builtin(
asset: &'static str,
yaml: &'static str,
cell: &'static OnceLock<Result<SourcePack, String>>,
) -> Result<&'static SourcePack, OpenConnectorError> {
cell.get_or_init(|| parse_pack(yaml))
.as_ref()
.map_err(|reason| OpenConnectorError::SourcePackAssetInvalid {
asset: asset.to_string(),
reason: reason.clone(),
})
}
pub(crate) fn parse_pack(yaml: &str) -> Result<SourcePack, String> {
let doc: PackDoc = serde_yaml::from_str(yaml).map_err(|e| e.to_string())?;
if doc.tables.is_empty() {
return Err("a pack must declare at least one table".to_string());
}
let pack_name = leak_str(doc.pack);
let mut tables = Vec::with_capacity(doc.tables.len());
for (short_name, table) in doc.tables {
if short_name.contains('.') {
return Err(format!(
"table key '{short_name}' must be a bare name; the id is derived as \
<pack>.<table>"
));
}
tables.push(convert_table(pack_name, &short_name, table)?);
}
Ok(SourcePack {
name: pack_name,
version: doc.version,
tables: leak_slice(tables),
})
}
fn convert_table(
pack: &'static str,
short_name: &str,
doc: TableDoc,
) -> Result<SourcePackTable, String> {
let id = leak_str(format!("{pack}.{short_name}"));
let mut fields = Vec::with_capacity(doc.columns.len());
for column in doc.columns {
fields.push(convert_column(id, column)?);
}
let filters = doc
.filters
.into_iter()
.map(convert_filter)
.collect::<Vec<_>>();
let mut fixed_inputs = Vec::with_capacity(doc.fixed_inputs.len());
for (key, value) in doc.fixed_inputs {
let fixed = match value {
FixedValueDoc::Float(v) if !v.is_finite() => {
return Err(format!(
"{id}: fixed input '{key}' contains {v}, which has no JSON \
spelling; pin finite numbers only"
));
}
FixedValueDoc::Bool(v) => FixedValue::Bool(v),
FixedValueDoc::Int(v) => FixedValue::Int(v),
FixedValueDoc::Float(v) => FixedValue::Float(v),
FixedValueDoc::Str(v) => FixedValue::Str(leak_str(v)),
FixedValueDoc::StrList(v) => FixedValue::StrList(leak_str_slice(v)),
FixedValueDoc::Json(v) => {
let json = yaml_to_json(v)
.map_err(|reason| format!("{id}: fixed input '{key}' {reason}"))?;
FixedValue::Json(Box::leak(Box::new(json)))
}
};
fixed_inputs.push((leak_str(key), fixed));
}
let action_id = leak_str(doc.action);
let (pagination, continuation) = doc.pagination.into_parts(action_id);
let table = SourcePackTable {
id,
action_id,
row_path: leak_str(doc.row_path),
row_shape: doc.row_shape.into(),
fields: leak_slice(fields),
pagination,
required_resources: leak_str_slice(doc.resources.required),
optional_resources: leak_str_slice(doc.resources.optional),
exclusive_resources: leak_str_groups(doc.resources.exclusive),
fixed_inputs: leak_slice(fixed_inputs),
filters: leak_slice(filters),
error_path: doc.error_path.map(leak_str),
expected_fingerprint: doc.fingerprint.map(leak_str),
continuation,
};
validate_table(&table)?;
Ok(table)
}
fn validate_table(table: &SourcePackTable) -> Result<(), String> {
let id = table.id;
match table.row_shape {
RowShape::Array => {
RowPath::parse(table.row_path).map_err(|e| format!("{id}: {e}"))?;
}
RowShape::Object => {
RowPath::parse_object_root(table.row_path).map_err(|e| format!("{id}: {e}"))?;
if !matches!(table.pagination, PaginationStrategy::SinglePage { .. }) {
return Err(format!(
"{id}: row_shape 'object' requires pagination strategy 'single_page'"
));
}
}
}
if let Some(path) = table.error_path {
RowPath::parse(path).map_err(|e| format!("{id}: {e}"))?;
}
table
.pagination
.validate()
.map_err(|e| format!("{id}: {e}"))?;
RowConverter::new(table.fields).map_err(|e| format!("{id}: {e}"))?;
let mut columns = std::collections::HashSet::new();
for field in table.fields {
if !columns.insert(field.name) {
return Err(format!("{id}: duplicate column '{}'", field.name));
}
}
let mut mappings = std::collections::HashSet::new();
for filter in table.filters {
if !columns.contains(filter.column) {
return Err(format!(
"{id}: filter references undeclared column '{}'",
filter.column
));
}
if !mappings.insert((filter.column, format!("{:?}", filter.operator))) {
return Err(format!(
"{id}: duplicate filter mapping for column '{}' and operator {:?}",
filter.column, filter.operator
));
}
if matches!(
filter.value_format,
ValueFormat::EpochSeconds | ValueFormat::EpochSecondsString
) {
if !matches!(filter.operator, Operator::Gt | Operator::GtEq) {
return Err(format!(
"{id}: filter on '{}' renders as flooring epoch seconds, which is only \
sound for lower bounds — operator {:?} would drop rows; declare gt/gt_eq \
or use the rfc3339 format",
filter.column, filter.operator
));
}
if filter.fidelity != Fidelity::Inexact {
return Err(format!(
"{id}: filter on '{}' floors sub-second literals into a WIDER fetch, so it \
must be declared inexact (DataFusion re-filters locally); exact would \
surface rows the predicate excludes",
filter.column
));
}
}
}
for required in table.required_resources {
if table.optional_resources.contains(required) {
return Err(format!(
"{id}: resource '{required}' is declared both required and optional"
));
}
}
let mut grouped: Vec<&str> = Vec::new();
for group in table.exclusive_resources {
if group.len() < 2 {
return Err(format!(
"{id}: exclusive resource group {group:?} needs at least two members to be a \
choice"
));
}
for key in *group {
if !table.declares_resource(key) {
return Err(format!(
"{id}: exclusive resource group names '{key}', which the table does not \
declare as a resource"
));
}
if table.required_resources.contains(key) {
return Err(format!(
"{id}: resource '{key}' is required, so it cannot be one of a group of \
alternatives — every other member would be unreachable"
));
}
if grouped.contains(key) {
return Err(format!(
"{id}: resource '{key}' appears in more than one exclusive group"
));
}
grouped.push(key);
}
}
let pagination_params: Vec<&str> = table.pagination_input_keys();
if pagination_params.iter().any(|p| p.is_empty()) {
return Err(format!("{id}: pagination input names must be non-empty"));
}
match table.pagination {
PaginationStrategy::PageNumber {
page_param,
per_page_param,
per_page,
..
} => {
if page_param == per_page_param {
return Err(format!(
"{id}: pagination declares '{page_param}' as both the page and page-size input"
));
}
if per_page == 0 {
return Err(format!("{id}: pagination page size must be positive"));
}
}
PaginationStrategy::Cursor {
cursor_param,
page_size_param,
page_size,
..
} => {
if page_size_param == Some(cursor_param) {
return Err(format!(
"{id}: pagination declares '{cursor_param}' as both the cursor and page-size input"
));
}
if page_size_param.is_some() && page_size == 0 {
return Err(format!("{id}: pagination page size must be positive"));
}
}
PaginationStrategy::Keyset {
cursor_param,
page_size_param,
page_size,
..
} => {
if page_size_param == cursor_param {
return Err(format!(
"{id}: pagination declares '{cursor_param}' as both the cursor and page-size input"
));
}
if page_size == 0 {
return Err(format!("{id}: pagination page size must be positive"));
}
}
PaginationStrategy::SinglePage { .. } => {}
}
let mut filter_inputs = std::collections::HashSet::new();
for filter in table.filters {
if !filter_inputs.insert(filter.input_field) {
return Err(format!(
"{id}: two filter mappings target input '{}'; declare one mapping per input",
filter.input_field
));
}
if pagination_params.contains(&filter.input_field) {
return Err(format!(
"{id}: filter input '{}' collides with a pagination input, which is applied last and would overwrite the pushed predicate",
filter.input_field
));
}
if table.declares_resource(filter.input_field) {
return Err(format!(
"{id}: filter input '{}' collides with a declared resource",
filter.input_field
));
}
}
for (key, _) in table.fixed_inputs {
if table.declares_resource(key) {
return Err(format!(
"{id}: fixed input '{key}' collides with a declared resource — the \
request would carry an ambiguous value"
));
}
if pagination_params.contains(key) {
return Err(format!(
"{id}: fixed input '{key}' collides with a pagination input"
));
}
}
for param in &pagination_params {
if table.declares_resource(param) {
return Err(format!(
"{id}: pagination input '{param}' collides with a declared resource — \
pagination is applied last and would overwrite the binding's value \
from page 2 onward"
));
}
}
if let Some(continuation) = table.continuation {
if table.expected_fingerprint.is_none() {
return Err(format!(
"{id}: continuation declares a fingerprint but the table does not; pin \
the table's own action too or neither is gated"
));
}
if continuation.action_id == table.action_id && !continuation.cursor_only {
return Err(format!(
"{id}: continuation repeats the table's own action with the full input, \
which is the default; drop the block or declare `inputs: cursor_only`"
));
}
if continuation.action_id == table.action_id
&& table.expected_fingerprint != Some(continuation.expected_fingerprint)
{
return Err(format!(
"{id}: continuation names the table's own action '{}' but pins a different \
fingerprint ({} vs {}); one action has one contract, so no gateway can \
satisfy both",
continuation.action_id,
table.expected_fingerprint.unwrap_or("<none>"),
continuation.expected_fingerprint
));
}
if continuation.cursor_only
&& let Some(filter) = table
.filters
.iter()
.find(|f| matches!(f.fidelity, Fidelity::Exact))
{
return Err(format!(
"{id}: filter on '{}' is Exact, which removes the Filter node from the \
plan, but `inputs: cursor_only` cannot carry the predicate past page one; \
declare the filter Inexact or drop `cursor_only`",
filter.input_field
));
}
}
Ok(())
}
fn convert_column(table_id: &str, doc: ColumnDoc) -> Result<FieldMapping, String> {
let field_type = match (doc.column_type, doc.key) {
(ColumnType::Utf8ListFromObjectKey, Some(key)) => {
FieldType::Utf8ListFromObjectKey(leak_str(key))
}
(ColumnType::Utf8ListFromObjectKey, None) => {
return Err(format!(
"{table_id}: column '{}' has type utf8_list_from_object_key and needs `key`",
doc.name
));
}
(other, Some(_)) => {
return Err(format!(
"{table_id}: column '{}' declares `key`, which only \
utf8_list_from_object_key accepts (got {other:?})",
doc.name
));
}
(ColumnType::Boolean, None) => FieldType::Boolean,
(ColumnType::Int64, None) => FieldType::Int64,
(ColumnType::Uint64, None) => FieldType::UInt64,
(ColumnType::Float64, None) => FieldType::Float64,
(ColumnType::Utf8, None) => FieldType::Utf8,
(ColumnType::TimestampMsUtc, None) => FieldType::TimestampMillisUtc,
(ColumnType::TimestampSUtc, None) => FieldType::TimestampSecondsUtc,
(ColumnType::TimestampMsStringUtc, None) => FieldType::TimestampMillisStringUtc,
(ColumnType::TimestampSStringUtc, None) => FieldType::TimestampSecondsStringUtc,
(ColumnType::TimestampSFracStringUtc, None) => {
FieldType::TimestampSecondsFractionalStringUtc
}
(ColumnType::Utf8List, None) => FieldType::Utf8List,
(ColumnType::Json, None) => FieldType::Json,
};
Ok(FieldMapping {
name: leak_str(doc.name),
path: leak_str(doc.path),
field_type,
nullable: doc.nullable,
})
}
fn convert_filter(doc: FilterDoc) -> FilterMapping {
FilterMapping {
column: leak_str(doc.column),
operator: match doc.op {
OpDoc::Eq => Operator::Eq,
OpDoc::Gt => Operator::Gt,
OpDoc::GtEq => Operator::GtEq,
},
input_field: leak_str(doc.input),
fidelity: match doc.fidelity {
FidelityDoc::Exact => Fidelity::Exact,
FidelityDoc::Inexact => Fidelity::Inexact,
},
value_format: match doc.format {
FormatDoc::Verbatim => ValueFormat::Verbatim,
FormatDoc::Rfc3339 => ValueFormat::Rfc3339,
FormatDoc::EpochSeconds => ValueFormat::EpochSeconds,
FormatDoc::EpochSecondsString => ValueFormat::EpochSecondsString,
},
}
}
fn leak_str(s: String) -> &'static str {
Box::leak(s.into_boxed_str())
}
fn leak_slice<T>(v: Vec<T>) -> &'static [T] {
Box::leak(v.into_boxed_slice())
}
fn leak_str_slice(v: Vec<String>) -> &'static [&'static str] {
leak_slice(v.into_iter().map(leak_str).collect())
}
fn leak_str_groups(v: Vec<Vec<String>>) -> &'static [&'static [&'static str]] {
leak_slice(v.into_iter().map(leak_str_slice).collect())
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct PackDoc {
#[allow(dead_code, reason = "the tag's value is its validation")]
kind: KindTag,
pack: String,
version: u32,
tables: BTreeMap<String, TableDoc>,
}
#[derive(Deserialize)]
enum KindTag {
#[serde(rename = "pack")]
Pack,
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct TableDoc {
action: String,
row_path: String,
#[serde(default)]
row_shape: RowShapeDoc,
#[serde(default)]
fingerprint: Option<String>,
pagination: PaginationDoc,
#[serde(default)]
resources: ResourcesDoc,
#[serde(default)]
fixed_inputs: BTreeMap<String, FixedValueDoc>,
#[serde(default)]
error_path: Option<String>,
columns: Vec<ColumnDoc>,
#[serde(default)]
filters: Vec<FilterDoc>,
}
#[derive(Deserialize, Default, Clone, Copy)]
#[serde(rename_all = "snake_case", deny_unknown_fields)]
enum RowShapeDoc {
#[default]
Array,
Object,
}
impl From<RowShapeDoc> for RowShape {
fn from(doc: RowShapeDoc) -> Self {
match doc {
RowShapeDoc::Array => Self::Array,
RowShapeDoc::Object => Self::Object,
}
}
}
#[derive(Deserialize, Default)]
#[serde(deny_unknown_fields)]
struct ResourcesDoc {
#[serde(default)]
required: Vec<String>,
#[serde(default)]
optional: Vec<String>,
#[serde(default)]
exclusive: Vec<Vec<String>>,
}
#[derive(Deserialize)]
#[serde(tag = "strategy", rename_all = "snake_case", deny_unknown_fields)]
enum PaginationDoc {
PageNumber {
page_input: String,
page_size_input: String,
page_size: u32,
#[serde(default)]
total_pages_path: Option<String>,
#[serde(default)]
raw_page_size_path: Option<String>,
},
Cursor {
cursor_input: String,
next_cursor_path: String,
#[serde(default)]
page_size_input: Option<String>,
page_size: u32,
#[serde(default)]
has_more_path: Option<String>,
#[serde(default)]
continuation: Option<ContinuationDoc>,
},
Keyset {
cursor_input: String,
row_cursor_field: String,
page_size_input: String,
page_size: u32,
},
SinglePage {
#[serde(default)]
next_cursor_path: Option<String>,
},
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct ContinuationDoc {
#[serde(default)]
action: Option<String>,
fingerprint: String,
#[serde(default)]
inputs: ContinuationInputsDoc,
}
#[derive(Deserialize, Default)]
enum ContinuationInputsDoc {
#[serde(rename = "full")]
#[default]
Full,
#[serde(rename = "cursor_only")]
CursorOnly,
}
impl PaginationDoc {
fn into_parts(
self,
table_action: &'static str,
) -> (PaginationStrategy, Option<CursorContinuation>) {
match self {
Self::PageNumber {
page_input,
page_size_input,
page_size,
total_pages_path,
raw_page_size_path,
} => (
PaginationStrategy::PageNumber {
page_param: leak_str(page_input),
per_page_param: leak_str(page_size_input),
per_page: page_size,
total_pages_path: total_pages_path.map(leak_str),
raw_page_size_path: raw_page_size_path.map(leak_str),
},
None,
),
Self::Cursor {
cursor_input,
next_cursor_path,
page_size_input,
page_size,
has_more_path,
continuation,
} => (
PaginationStrategy::Cursor {
cursor_param: leak_str(cursor_input),
next_cursor_path: leak_str(next_cursor_path),
page_size_param: page_size_input.map(leak_str),
page_size,
has_more_path: has_more_path.map(leak_str),
absent_cursor: AbsentCursor::EndsTheScan,
},
continuation.map(|doc| CursorContinuation {
action_id: doc.action.map_or(table_action, leak_str),
expected_fingerprint: leak_str(doc.fingerprint),
cursor_only: matches!(doc.inputs, ContinuationInputsDoc::CursorOnly),
}),
),
Self::Keyset {
cursor_input,
row_cursor_field,
page_size_input,
page_size,
} => (
PaginationStrategy::Keyset {
cursor_param: leak_str(cursor_input),
row_cursor_field: leak_str(row_cursor_field),
page_size_param: leak_str(page_size_input),
page_size,
},
None,
),
Self::SinglePage { next_cursor_path } => (
PaginationStrategy::SinglePage {
next_cursor_path: next_cursor_path.map(leak_str),
},
None,
),
}
}
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct ColumnDoc {
name: String,
path: String,
#[serde(rename = "type")]
column_type: ColumnType,
#[serde(default)]
key: Option<String>,
nullable: bool,
}
#[derive(Debug, Deserialize, Clone, Copy)]
enum ColumnType {
#[serde(rename = "boolean")]
Boolean,
#[serde(rename = "int64")]
Int64,
#[serde(rename = "uint64")]
Uint64,
#[serde(rename = "float64")]
Float64,
#[serde(rename = "utf8")]
Utf8,
#[serde(rename = "timestamp_ms_utc")]
TimestampMsUtc,
#[serde(rename = "timestamp_s_utc")]
TimestampSUtc,
#[serde(rename = "timestamp_ms_string_utc")]
TimestampMsStringUtc,
#[serde(rename = "timestamp_s_string_utc")]
TimestampSStringUtc,
#[serde(rename = "timestamp_s_frac_string_utc")]
TimestampSFracStringUtc,
#[serde(rename = "utf8_list")]
Utf8List,
#[serde(rename = "utf8_list_from_object_key")]
Utf8ListFromObjectKey,
#[serde(rename = "json")]
Json,
}
#[derive(Deserialize)]
#[serde(untagged)]
enum FixedValueDoc {
Bool(bool),
Int(i64),
Float(f64),
Str(String),
StrList(Vec<String>),
Json(serde_yaml::Value),
}
fn yaml_to_json(value: serde_yaml::Value) -> Result<serde_json::Value, String> {
use serde_yaml::Value as Yaml;
Ok(match value {
Yaml::Null => serde_json::Value::Null,
Yaml::Bool(b) => serde_json::Value::from(b),
Yaml::Number(n) => {
if let Some(i) = n.as_i64() {
serde_json::Value::from(i)
} else if let Some(u) = n.as_u64() {
serde_json::Value::from(u)
} else {
let f = n.as_f64().unwrap_or(f64::NAN);
if !f.is_finite() {
return Err(format!(
"contains {f}, which has no JSON spelling; pin finite numbers only"
));
}
serde_json::Value::from(f)
}
}
Yaml::String(s) => serde_json::Value::from(s),
Yaml::Sequence(items) => serde_json::Value::Array(
items
.into_iter()
.map(yaml_to_json)
.collect::<Result<_, _>>()?,
),
Yaml::Mapping(map) => {
let mut out = serde_json::Map::with_capacity(map.len());
for (k, v) in map {
let Yaml::String(k) = k else {
return Err(
"contains a non-string mapping key, which JSON cannot represent"
.to_string(),
);
};
out.insert(k, yaml_to_json(v)?);
}
serde_json::Value::Object(out)
}
Yaml::Tagged(_) => {
return Err("contains a YAML tag, which JSON cannot represent".to_string());
}
})
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct FilterDoc {
column: String,
op: OpDoc,
input: String,
fidelity: FidelityDoc,
#[serde(default)]
format: FormatDoc,
}
#[derive(Deserialize)]
enum OpDoc {
#[serde(rename = "eq")]
Eq,
#[serde(rename = "gt")]
Gt,
#[serde(rename = "gt_eq")]
GtEq,
}
#[derive(Deserialize, Default)]
enum FidelityDoc {
#[serde(rename = "exact")]
#[default]
Exact,
#[serde(rename = "inexact")]
Inexact,
}
#[derive(Deserialize, Default)]
enum FormatDoc {
#[serde(rename = "verbatim")]
#[default]
Verbatim,
#[serde(rename = "rfc3339")]
Rfc3339,
#[serde(rename = "epoch_seconds")]
EpochSeconds,
#[serde(rename = "epoch_seconds_string")]
EpochSecondsString,
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn builtin_assets_parse_and_validate() {
let registry = crate::sources::providers::open_connector::builtin_pack_registry()
.expect("every embedded asset parses and validates");
for pack in registry.packs() {
assert!(!pack.tables.is_empty(), "{}: no tables", pack.name);
}
}
#[test]
fn nested_finite_values_in_a_json_pin_convert_faithfully() {
let pack = parse_pack(
r#"kind: pack
pack: demo
version: 1
tables:
things:
action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
fixed_inputs:
filter:
threshold: 1.5
flags: [true, 2, "three"]
inner: { level: null }
columns:
- { name: id, path: id, type: uint64, nullable: false }
"#,
)
.expect("nested finite values are legal");
let (key, value) = &pack.tables[0].fixed_inputs[0];
assert_eq!(*key, "filter");
assert_eq!(
value.to_json(),
serde_json::json!({
"threshold": 1.5,
"flags": [true, 2, "three"],
"inner": { "level": null }
})
);
}
#[test]
fn misspelled_keys_fail_loudly() {
let err = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10, total_page_path: "$.pages" }
columns:
- { name: id, path: id, type: uint64, nullable: false }
"#,
)
.unwrap_err();
assert!(err.contains("total_page_path"), "{err}");
}
#[test]
fn row_shape_defaults_to_array_for_every_existing_pack() {
let pack = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
action: demo.list
row_path: "$.items"
pagination: { strategy: single_page }
columns:
- { name: id, path: id, type: uint64, nullable: false }
"#,
)
.expect("a pack without row_shape still parses");
assert_eq!(pack.tables[0].row_shape, RowShape::Array);
}
#[test]
fn object_row_shape_parses_with_root_path_and_single_page() {
let pack = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
document_content:
action: demo.get_document_content
row_path: "$"
row_shape: object
pagination: { strategy: single_page }
columns:
- { name: content, path: content, type: utf8, nullable: true }
"#,
)
.expect("object rows at the root with single_page are valid");
assert_eq!(pack.tables[0].row_shape, RowShape::Object);
assert_eq!(pack.tables[0].row_path, "$");
}
#[test]
fn object_row_shape_requires_single_page() {
let err = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
document_content:
action: demo.get_document_content
row_path: "$"
row_shape: object
pagination:
strategy: cursor
cursor_input: pageToken
next_cursor_path: "$.pageToken"
page_size_input: pageSize
page_size: 100
columns:
- { name: content, path: content, type: utf8, nullable: true }
"#,
)
.unwrap_err();
assert!(err.contains("single_page"), "{err}");
}
#[test]
fn object_row_shape_rejects_a_non_root_path() {
let err = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
document_content:
action: demo.get_document_content
row_path: "$.document"
row_shape: object
pagination: { strategy: single_page }
columns:
- { name: content, path: content, type: utf8, nullable: true }
"#,
)
.unwrap_err();
assert!(err.contains("'$'"), "{err}");
}
#[test]
fn array_row_shape_still_rejects_the_root_path() {
let err = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
action: demo.list
row_path: "$"
pagination: { strategy: single_page }
columns:
- { name: id, path: id, type: uint64, nullable: false }
"#,
)
.unwrap_err();
assert!(err.contains("row path"), "{err}");
}
#[test]
fn single_page_strategy_parses_and_rejects_stray_keys() {
let pack = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
action: demo.list
row_path: "$.items"
pagination: { strategy: single_page }
columns:
- { name: id, path: id, type: uint64, nullable: false }
"#,
)
.expect("single_page is a valid strategy");
assert!(matches!(
pack.tables[0].pagination,
PaginationStrategy::SinglePage {
next_cursor_path: None
}
));
let pack = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
action: demo.list
row_path: "$.items"
pagination: { strategy: single_page, next_cursor_path: "$.nextPageToken" }
columns:
- { name: id, path: id, type: uint64, nullable: false }
"#,
)
.expect("single_page accepts an optional next_cursor_path");
assert!(matches!(
pack.tables[0].pagination,
PaginationStrategy::SinglePage {
next_cursor_path: Some("$.nextPageToken")
}
));
let err = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
action: demo.list
row_path: "$.items"
pagination: { strategy: single_page, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }
"#,
)
.unwrap_err();
assert!(err.contains("page_size"), "{err}");
}
#[test]
fn key_field_is_bound_to_the_plucking_type() {
let base = |columns: &str| {
format!(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
action: demo.list
row_path: "$.items"
pagination: {{ strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }}
columns:
{columns}
"#
)
};
let err = parse_pack(&base(
" - { name: labels, path: labels, type: utf8_list_from_object_key, nullable: true }",
))
.unwrap_err();
assert!(err.contains("needs `key`"), "{err}");
let err = parse_pack(&base(
" - { name: id, path: id, type: uint64, key: name, nullable: false }",
))
.unwrap_err();
assert!(err.contains("only"), "{err}");
}
#[test]
fn epoch_formats_are_lower_bound_inexact_only() {
let base = |filter: &str| {
format!(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
action: demo.list
row_path: "$.items"
pagination: {{ strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }}
columns:
- {{ name: created, path: created, type: timestamp_ms_utc, nullable: true }}
filters:
{filter}
"#
)
};
let err = parse_pack(&base(
" - { column: created, op: eq, input: at, fidelity: inexact, format: epoch_seconds }",
))
.unwrap_err();
assert!(
err.contains("only") && err.contains("lower bounds"),
"{err}"
);
let err = parse_pack(&base(
" - { column: created, op: gt_eq, input: since, fidelity: exact, format: epoch_seconds }",
))
.unwrap_err();
assert!(err.contains("inexact"), "{err}");
parse_pack(&base(
" - { column: created, op: gt_eq, input: since, fidelity: inexact, format: epoch_seconds_string }",
))
.expect("lower-bound inexact epoch mapping is legal");
}
#[test]
fn keyset_and_single_page_pagination_parse_and_validate() {
let base = |pagination: &str| {
format!(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
action: demo.list
row_path: "$.items"
pagination: {pagination}
columns:
- {{ name: id, path: id, type: utf8, nullable: false }}
"#
)
};
let pack = parse_pack(&base(
"{ strategy: keyset, cursor_input: after, row_cursor_field: id, page_size_input: limit, page_size: 200 }",
))
.expect("keyset parses");
assert!(matches!(
pack.tables[0].pagination,
PaginationStrategy::Keyset { page_size: 200, .. }
));
let pack = parse_pack(&base("{ strategy: single_page }")).expect("single_page parses");
assert!(matches!(
pack.tables[0].pagination,
PaginationStrategy::SinglePage { .. }
));
let err = parse_pack(&base(
"{ strategy: keyset, cursor_input: limit, row_cursor_field: id, page_size_input: limit, page_size: 200 }",
))
.unwrap_err();
assert!(err.contains("both the cursor and page-size input"), "{err}");
let err = parse_pack(&base(
"{ strategy: keyset, cursor_input: after, row_cursor_field: id, page_size_input: limit, page_size: 0 }",
))
.unwrap_err();
assert!(err.contains("page size must be positive"), "{err}");
let err = parse_pack(&base(
"{ strategy: keyset, cursor_input: after, row_cursor_field: \"$.id\", page_size_input: limit, page_size: 200 }",
))
.unwrap_err();
assert!(err.contains("relative to the row"), "{err}");
for pagination in [
"{ strategy: keyset, cursor_input: \"\", row_cursor_field: id, page_size_input: limit, page_size: 200 }",
"{ strategy: keyset, cursor_input: after, row_cursor_field: id, page_size_input: \"\", page_size: 200 }",
"{ strategy: cursor, cursor_input: \"\", next_cursor_path: \"$.next\", page_size: 100 }",
"{ strategy: cursor, cursor_input: cursor, next_cursor_path: \"$.next\", page_size_input: \"\", page_size: 100 }",
"{ strategy: page_number, page_input: \"\", page_size_input: perPage, page_size: 100 }",
"{ strategy: page_number, page_input: page, page_size_input: \"\", page_size: 100 }",
] {
let err = parse_pack(&base(pagination)).unwrap_err();
assert!(err.contains("must be non-empty"), "{err}");
}
}
#[test]
fn a_valid_exclusive_group_survives_the_round_trip() {
let pack = pack_with(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { optional: [byId, byPath], exclusive: [[byId, byPath]] }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
)
.expect("two declared optional resources are a valid group");
let table = &pack.tables[0];
let groups: Vec<Vec<&str>> = table
.exclusive_resources
.iter()
.map(|group| group.to_vec())
.collect();
assert_eq!(groups, vec![vec!["byId", "byPath"]]);
assert_eq!(
table.conflicting_resources(|key| ["byId", "byPath"].contains(&key)),
Some(("byId", "byPath")),
"both supplied is the ambiguity"
);
for one in ["byId", "byPath"] {
assert_eq!(
table.conflicting_resources(|key| key == one),
None,
"{one} alone names one scope"
);
}
assert_eq!(
table.conflicting_resources(|_| false),
None,
"neither supplied is the unscoped default"
);
}
#[test]
fn tables_without_an_exclusive_group_declare_none() {
let pack = pack_with(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { optional: [byId, byPath] }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
)
.expect("omitting the group stays valid");
assert!(pack.tables[0].exclusive_resources.is_empty());
assert_eq!(pack.tables[0].conflicting_resources(|_| true), None);
}
#[test]
fn dotted_table_keys_are_rejected() {
let err = parse_pack(
r#"
kind: pack
pack: demo
version: 1
tables:
demo.items:
action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }
"#,
)
.unwrap_err();
assert!(err.contains("bare name"), "{err}");
}
fn pack_with(table_body: &str) -> Result<SourcePack, String> {
parse_pack(&format!(
r#"
kind: pack
pack: demo
version: 1
tables:
items:
{table_body}
"#
))
}
#[test]
fn semantic_invariants_are_rejected_with_targeted_errors() {
for (body, expected) in [
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }
- { name: id, path: other, type: utf8, nullable: true }"#,
"duplicate column 'id'",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }
filters:
- { column: missing, op: eq, input: q, fidelity: inexact }"#,
"undeclared column 'missing'",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }
filters:
- { column: id, op: eq, input: a, fidelity: inexact }
- { column: id, op: eq, input: b, fidelity: inexact }"#,
"duplicate filter mapping",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { required: [owner], optional: [owner] }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"both required and optional",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { optional: [byId], exclusive: [[byId]] }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"at least two members",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { optional: [byId], exclusive: [[byId, byPath]] }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"does not declare as a resource",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { required: [byId], optional: [byPath], exclusive: [[byId, byPath]] }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"every other member would be unreachable",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { optional: [byId, byPath, byName], exclusive: [[byId, byPath], [byPath, byName]] }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"more than one exclusive group",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { required: [owner] }
fixed_inputs:
owner: acme
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"collides with a declared resource",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
fixed_inputs:
page: 1
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"collides with a pagination input",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: keyset, cursor_input: after, page_size_input: limit, page_size: 10, row_cursor_field: id }
resources: { required: [after] }
columns:
- { name: id, path: id, type: utf8, nullable: false }"#,
"pagination input 'after' collides with a declared resource",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { optional: [perPage] }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"pagination input 'perPage' collides with a declared resource",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }
filters:
- { column: id, op: eq, input: perPage, fidelity: exact }"#,
"collides with a pagination input",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
resources: { required: [owner] }
columns:
- { name: id, path: id, type: uint64, nullable: false }
filters:
- { column: id, op: eq, input: owner, fidelity: exact }"#,
"collides with a declared resource",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }
- { name: score, path: score, type: uint64, nullable: true }
filters:
- { column: id, op: eq, input: q, fidelity: inexact }
- { column: score, op: eq, input: q, fidelity: inexact }"#,
"two filter mappings target input 'q'",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
fixed_inputs:
threshold: .nan
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"no JSON spelling",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
fixed_inputs:
threshold: .inf
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"no JSON spelling",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
fixed_inputs:
threshold: -.inf
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"no JSON spelling",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
fixed_inputs:
filter:
threshold: .nan
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"no JSON spelling",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
fixed_inputs:
filter:
bounds: [1.5, .inf]
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"no JSON spelling",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
fixed_inputs:
filter:
1: numeric-key
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"non-string mapping key",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
fixed_inputs:
filter:
payload: !custom tagged
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"do not support enum input",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: page, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"both the page and page-size input",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 0 }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"page size must be positive",
),
(
r#" action: demo.list
row_path: "$.items"
pagination: { strategy: cursor, cursor_input: cursor, next_cursor_path: "$.next", page_size_input: cursor, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"both the cursor and page-size input",
),
(
r#" action: demo.list
row_path: "items"
pagination: { strategy: page_number, page_input: page, page_size_input: perPage, page_size: 10 }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"must start with '$.'",
),
] {
let err = pack_with(body).expect_err(expected);
assert!(err.contains(expected), "want {expected:?} in: {err}");
}
}
#[test]
fn split_action_continuation_parses_with_its_defaults() {
let pack = pack_with(
r#" action: demo.list
row_path: "$.entries"
fingerprint: aa
pagination:
strategy: cursor
cursor_input: cursor
next_cursor_path: "$.cursor"
page_size_input: limit
page_size: 2000
has_more_path: "$.hasMore"
continuation:
action: demo.list_continue
fingerprint: bb
inputs: cursor_only
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
)
.expect("continuation parses");
let continuation = pack.tables[0].continuation.expect("declared");
assert_eq!(continuation.action_id, "demo.list_continue");
assert_eq!(continuation.expected_fingerprint, "bb");
assert!(continuation.cursor_only);
let table = &pack.tables[0];
assert_eq!(
table.actions().collect::<Vec<_>>(),
vec!["demo.list", "demo.list_continue"]
);
assert_eq!(
table.gated_actions().collect::<Vec<_>>(),
vec![("demo.list", "aa"), ("demo.list_continue", "bb")]
);
}
#[test]
fn continuation_action_defaults_to_the_tables_own() {
let pack = pack_with(
r#" action: demo.list
row_path: "$.entries"
fingerprint: aa
pagination:
strategy: cursor
cursor_input: cursor
next_cursor_path: "$.cursor"
page_size: 100
continuation:
fingerprint: aa
inputs: cursor_only
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
)
.expect("continuation parses");
let continuation = pack.tables[0].continuation.expect("declared");
assert_eq!(continuation.action_id, "demo.list");
assert!(continuation.cursor_only);
}
#[test]
fn continuation_invariants_are_rejected_with_targeted_errors() {
for (body, expected) in [
(
r#" action: demo.list
row_path: "$.items"
pagination:
strategy: page_number
page_input: page
page_size_input: perPage
page_size: 10
continuation: { fingerprint: bb, inputs: cursor_only }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"unknown field",
),
(
r#" action: demo.list
row_path: "$.items"
pagination:
strategy: cursor
cursor_input: cursor
next_cursor_path: "$.cursor"
page_size: 10
continuation: { action: demo.list_continue, fingerprint: bb }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"the table does not",
),
(
r#" action: demo.list
row_path: "$.items"
fingerprint: aa
pagination:
strategy: cursor
cursor_input: cursor
next_cursor_path: "$.cursor"
page_size: 10
continuation: { fingerprint: aa }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"which is the default",
),
(
r#" action: demo.list
row_path: "$.items"
fingerprint: aa
pagination:
strategy: cursor
cursor_input: cursor
next_cursor_path: "$.cursor"
page_size: 10
continuation: { action: demo.list_continue, inputs: cursor_only }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"missing field `fingerprint`",
),
(
r#" action: demo.list
row_path: "$.items"
fingerprint: aa
pagination:
strategy: cursor
cursor_input: cursor
next_cursor_path: "$.cursor"
page_size: 10
continuation: { fingerprint: bb, inputs: cursor_only }
columns:
- { name: id, path: id, type: uint64, nullable: false }"#,
"pins a different fingerprint",
),
(
r#" action: demo.list
row_path: "$.items"
fingerprint: aa
pagination:
strategy: cursor
cursor_input: cursor
next_cursor_path: "$.cursor"
page_size: 10
continuation: { action: demo.list_continue, fingerprint: bb, inputs: cursor_only }
columns:
- { name: id, path: id, type: uint64, nullable: false }
- { name: state, path: state, type: utf8, nullable: true }
filters:
- { column: state, op: eq, input: state, fidelity: exact }"#,
"cannot carry the predicate past page one",
),
] {
let err = pack_with(body).expect_err(expected);
assert!(err.contains(expected), "want {expected:?} in: {err}");
}
}
#[test]
fn an_inexact_filter_stays_legal_alongside_a_cursor_only_continuation() {
let pack = pack_with(
r#" action: demo.list
row_path: "$.items"
fingerprint: aa
pagination:
strategy: cursor
cursor_input: cursor
next_cursor_path: "$.cursor"
page_size: 10
continuation: { action: demo.list_continue, fingerprint: bb, inputs: cursor_only }
columns:
- { name: id, path: id, type: uint64, nullable: false }
- { name: state, path: state, type: utf8, nullable: true }
filters:
- { column: state, op: eq, input: state, fidelity: inexact }"#,
)
.expect("an Inexact filter cannot return rows the query excluded");
assert_eq!(pack.tables[0].filters.len(), 1);
}
}