use crate::error::FaucetError;
use chrono::{DateTime, NaiveDate, NaiveDateTime, Utc};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::cmp::Ordering;
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(tag = "type")]
pub enum ReplicationMethod {
#[default]
FullTable,
Incremental,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum OnMissingKey {
#[default]
Keep,
Drop,
Fail,
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum KeyForm {
TopLevel,
DotPath(Vec<String>),
Pointer,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ReplicationKey {
raw: String,
form: KeyForm,
}
impl ReplicationKey {
pub fn parse(raw: &str) -> Result<Self, FaucetError> {
if raw.trim().is_empty() {
return Err(FaucetError::Config(
"replication_key must not be empty".to_owned(),
));
}
if raw.starts_with('/') {
return Ok(Self {
raw: raw.to_owned(),
form: KeyForm::Pointer,
});
}
if !raw.contains('.') {
return Ok(Self::top_level(raw));
}
let segments: Vec<String> = raw.split('.').map(str::to_owned).collect();
if segments.iter().any(String::is_empty) {
return Err(FaucetError::Config(format!(
"replication_key '{raw}': empty path segment (use the JSON Pointer form \
`/a/b` for field names that contain dots)"
)));
}
Ok(Self {
raw: raw.to_owned(),
form: KeyForm::DotPath(segments),
})
}
pub fn top_level(name: &str) -> Self {
Self {
raw: name.to_owned(),
form: KeyForm::TopLevel,
}
}
pub fn as_str(&self) -> &str {
&self.raw
}
pub fn is_pointer(&self) -> bool {
self.form == KeyForm::Pointer
}
pub fn is_nested(&self) -> bool {
self.form != KeyForm::TopLevel
}
pub fn resolve<'a>(&self, record: &'a Value) -> Option<&'a Value> {
match &self.form {
KeyForm::TopLevel => record.get(&self.raw),
KeyForm::Pointer => record.pointer(&self.raw),
KeyForm::DotPath(segments) => {
if let Some(v) = record.get(&self.raw) {
return Some(v);
}
segments.iter().try_fold(record, |cur, seg| match cur {
Value::Object(m) => m.get(seg),
Value::Array(a) => seg.parse::<usize>().ok().and_then(|i| a.get(i)),
_ => None,
})
}
}
}
fn resolve_present<'a>(&self, record: &'a Value) -> Option<&'a Value> {
self.resolve(record).filter(|v| !v.is_null())
}
}
#[derive(Debug, Clone, PartialEq, Default)]
pub struct IncrementalFilter {
pub records: Vec<Value>,
pub missing: usize,
}
pub fn filter_incremental_path(
records: Vec<Value>,
key: &ReplicationKey,
start: &Value,
on_missing: OnMissingKey,
) -> Result<IncrementalFilter, FaucetError> {
let mut missing = 0usize;
let mut kept = Vec::with_capacity(records.len());
for r in records {
let keep = match key.resolve_present(&r) {
None => {
missing += 1;
match on_missing {
OnMissingKey::Keep => true,
OnMissingKey::Drop => false,
OnMissingKey::Fail => {
return Err(FaucetError::Source(format!(
"incremental replication: a record lacks replication_key '{}' \
(on_missing_key: fail)",
key.as_str()
)));
}
}
}
Some(v) if type_rank(v) != type_rank(start) => {
tracing::warn!(
key = key.as_str(),
"incremental replication: record key type does not match the bookmark \
type; keeping the record to avoid silently dropping data"
);
true
}
Some(v) => json_gt(v, start),
};
if keep {
kept.push(r);
}
}
Ok(IncrementalFilter {
records: kept,
missing,
})
}
pub fn filter_incremental(records: Vec<Value>, key: &str, start: &Value) -> Vec<Value> {
let out = filter_incremental_path(
records,
&ReplicationKey::top_level(key),
start,
OnMissingKey::Keep,
)
.unwrap_or_default();
if out.missing > 0 {
tracing::warn!(
key,
missing = out.missing,
"incremental replication: {} record(s) lacked replication_key '{key}'; kept to \
avoid silent data loss",
out.missing
);
}
out.records
}
pub fn max_replication_value_path<'a>(
records: &'a [Value],
key: &ReplicationKey,
) -> Option<&'a Value> {
records
.iter()
.filter_map(|r| key.resolve_present(r))
.max_by(|a, b| json_compare(a, b))
}
pub fn max_replication_value<'a>(records: &'a [Value], key: &str) -> Option<&'a Value> {
records
.iter()
.filter_map(|r| r.get(key))
.max_by(|a, b| json_compare(a, b))
}
pub fn max_value(a: Value, b: Value) -> Value {
match json_compare(&a, &b) {
Ordering::Less => b,
_ => a,
}
}
fn type_rank(v: &Value) -> u8 {
match v {
Value::Null => 0,
Value::Bool(_) => 1,
Value::Number(_) => 2,
Value::String(_) => 3,
Value::Array(_) => 4,
Value::Object(_) => 5,
}
}
fn number_as_i128(n: &serde_json::Number) -> Option<i128> {
n.as_i64()
.map(i128::from)
.or_else(|| n.as_u64().map(i128::from))
}
pub(crate) fn json_compare(a: &Value, b: &Value) -> Ordering {
match (a, b) {
(Value::Number(an), Value::Number(bn)) => {
match (number_as_i128(an), number_as_i128(bn)) {
(Some(ai), Some(bi)) => ai.cmp(&bi),
_ => {
let af = an.as_f64().unwrap_or(f64::NAN);
let bf = bn.as_f64().unwrap_or(f64::NAN);
af.partial_cmp(&bf).unwrap_or_else(|| {
match (af.is_nan(), bf.is_nan()) {
(false, true) => Ordering::Less,
(true, false) => Ordering::Greater,
_ => Ordering::Equal,
}
})
}
}
}
(Value::String(x), Value::String(y)) => x.cmp(y),
(Value::Bool(x), Value::Bool(y)) => x.cmp(y),
(Value::Null, Value::Null) => Ordering::Equal,
(Value::Array(x), Value::Array(y)) => {
for (xi, yi) in x.iter().zip(y.iter()) {
let c = json_compare(xi, yi);
if c != Ordering::Equal {
return c;
}
}
x.len().cmp(&y.len())
}
(Value::Object(_), Value::Object(_)) => a.to_string().cmp(&b.to_string()),
_ => type_rank(a).cmp(&type_rank(b)),
}
}
pub fn json_gt(a: &Value, b: &Value) -> bool {
json_compare(a, b) == Ordering::Greater
}
pub const BIND_PLACEHOLDER: &str = "${bookmark}";
fn default_bind_template() -> String {
BIND_PLACEHOLDER.to_owned()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum BindTarget {
#[default]
Query,
Header,
Body,
Path,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum BindFormat {
#[default]
Raw,
Iso8601,
EpochS,
EpochMs,
Date,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum BindValueType {
#[default]
String,
Number,
}
impl BindValueType {
pub fn to_value(self, rendered: &str) -> Result<Value, FaucetError> {
match self {
Self::String => Ok(Value::String(rendered.to_owned())),
Self::Number => {
let n: Option<serde_json::Number> = rendered
.parse::<i64>()
.map(serde_json::Number::from)
.ok()
.or_else(|| rendered.parse::<u64>().ok().map(serde_json::Number::from))
.or_else(|| {
rendered
.parse::<f64>()
.ok()
.and_then(serde_json::Number::from_f64)
});
n.map(Value::Number).ok_or_else(|| {
FaucetError::Source(format!(
"bind: rendered value '{rendered}' is not a number (value_type: number)"
))
})
}
}
}
}
fn unescape_pointer_token(token: &str) -> String {
token.replace("~1", "/").replace("~0", "~")
}
pub fn set_body_pointer(body: &mut Value, pointer: &str, value: Value) -> Result<(), FaucetError> {
if !pointer.starts_with('/') {
return Err(FaucetError::Config(format!(
"JSON Pointer '{pointer}' must start with '/'"
)));
}
if let Some(slot) = body.pointer_mut(pointer) {
if slot.is_object() || slot.is_array() {
return Err(FaucetError::Source(format!(
"request body location '{pointer}' holds an object or array; a bind replaces \
a scalar value"
)));
}
*slot = value;
return Ok(());
}
let cut = pointer.rfind('/').unwrap_or(0);
let (parent, leaf) = (&pointer[..cut], &pointer[cut + 1..]);
let parent_value = if parent.is_empty() {
Some(body)
} else {
body.pointer_mut(parent)
};
match parent_value {
Some(Value::Object(map)) => {
map.insert(unescape_pointer_token(leaf), value);
Ok(())
}
_ => Err(FaucetError::Source(format!(
"request body has no location '{pointer}' (the pointer must resolve to an \
existing value, or to a new key of an existing object)"
))),
}
}
pub(crate) fn validate_bind_placement(
what: &str,
into: BindTarget,
name: &str,
path: Option<&str>,
) -> Result<(), FaucetError> {
let has_name = !name.trim().is_empty();
match (into, path) {
(BindTarget::Body, Some(p)) => {
if has_name {
return Err(FaucetError::Config(format!(
"{what}: set either `name` (top-level body field) or `path` (JSON Pointer), \
not both"
)));
}
if !p.starts_with('/') || p.len() < 2 {
return Err(FaucetError::Config(format!(
"{what}: `path` must be a JSON Pointer such as `/filters/0/value`, got '{p}'"
)));
}
Ok(())
}
(_, Some(_)) => Err(FaucetError::Config(format!(
"{what}: `path` applies only to `into: body`"
))),
(BindTarget::Body, None) if !has_name => Err(FaucetError::Config(format!(
"{what}: `into: body` needs `name` (top-level field) or `path` (JSON Pointer)"
))),
(_, None) if !has_name => Err(FaucetError::Config(format!(
"{what}: `name` must not be empty"
))),
_ => Ok(()),
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(deny_unknown_fields)]
pub struct ReplicationBind {
#[serde(default)]
pub into: BindTarget,
#[serde(default)]
pub name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub path: Option<String>,
#[serde(default)]
pub value_type: BindValueType,
#[serde(default = "default_bind_template")]
pub template: String,
#[serde(default)]
pub format: BindFormat,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub advance_from: Option<String>,
}
impl ReplicationBind {
pub fn validate(&self) -> Result<(), FaucetError> {
validate_bind_placement(
"replication bind",
self.into,
&self.name,
self.path.as_deref(),
)?;
if !self.template.contains(BIND_PLACEHOLDER) {
return Err(FaucetError::Config(format!(
"replication bind: `template` must contain the `{BIND_PLACEHOLDER}` placeholder"
)));
}
Ok(())
}
pub fn render(&self, bookmark: &Value) -> Result<String, FaucetError> {
let formatted = format_bookmark(bookmark, self.format)?;
Ok(self.template.replace(BIND_PLACEHOLDER, &formatted))
}
}
fn bookmark_instant(value: &Value) -> Result<DateTime<Utc>, FaucetError> {
match value {
Value::String(s) => {
let s = s.trim();
if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
return Ok(dt.with_timezone(&Utc));
}
if let Ok(d) = NaiveDate::parse_from_str(s, "%Y-%m-%d")
&& let Some(ndt) = d.and_hms_opt(0, 0, 0)
{
return Ok(DateTime::<Utc>::from_naive_utc_and_offset(ndt, Utc));
}
if let Ok(ndt) = NaiveDateTime::parse_from_str(s, "%Y-%m-%dT%H:%M:%S") {
return Ok(DateTime::<Utc>::from_naive_utc_and_offset(ndt, Utc));
}
Err(FaucetError::Config(format!(
"replication bind: cannot parse bookmark '{s}' as a timestamp \
(expected RFC 3339, YYYY-MM-DD, or YYYY-MM-DDTHH:MM:SS)"
)))
}
Value::Number(n) => {
let secs = n.as_i64().or_else(|| n.as_f64().map(|f| f as i64));
secs.and_then(|s| DateTime::<Utc>::from_timestamp(s, 0))
.ok_or_else(|| {
FaucetError::Config(format!(
"replication bind: numeric bookmark {n} is out of range for epoch seconds"
))
})
}
other => Err(FaucetError::Config(format!(
"replication bind: bookmark must be a string or number, got {other}"
))),
}
}
pub fn parse_instant(value: &Value) -> Result<DateTime<Utc>, FaucetError> {
bookmark_instant(value)
}
pub fn format_instant(dt: DateTime<Utc>, format: BindFormat) -> String {
match format {
BindFormat::Raw | BindFormat::Iso8601 => {
dt.to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
}
BindFormat::Date => dt.format("%Y-%m-%d").to_string(),
BindFormat::EpochS => dt.timestamp().to_string(),
BindFormat::EpochMs => dt.timestamp_millis().to_string(),
}
}
pub fn format_bookmark(value: &Value, format: BindFormat) -> Result<String, FaucetError> {
match format {
BindFormat::Raw => match value {
Value::String(s) => Ok(s.clone()),
Value::Number(n) => Ok(n.to_string()),
Value::Bool(b) => Ok(b.to_string()),
other => Err(FaucetError::Config(format!(
"replication bind: cannot render {other} as a raw scalar"
))),
},
BindFormat::Iso8601 => {
Ok(bookmark_instant(value)?.to_rfc3339_opts(chrono::SecondsFormat::Secs, true))
}
BindFormat::Date => Ok(bookmark_instant(value)?.format("%Y-%m-%d").to_string()),
BindFormat::EpochS => Ok(bookmark_instant(value)?.timestamp().to_string()),
BindFormat::EpochMs => Ok(bookmark_instant(value)?.timestamp_millis().to_string()),
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn test_filter_incremental_strings() {
let records = vec![
json!({"id": 1, "updated_at": "2024-01-01"}),
json!({"id": 2, "updated_at": "2024-06-01"}),
json!({"id": 3, "updated_at": "2024-12-01"}),
];
let start = json!("2024-06-01");
let filtered = filter_incremental(records, "updated_at", &start);
assert_eq!(filtered.len(), 1);
assert_eq!(filtered[0]["id"], 3);
}
#[test]
fn test_filter_incremental_numbers() {
let records = vec![
json!({"id": 1, "seq": 100}),
json!({"id": 2, "seq": 200}),
json!({"id": 3, "seq": 300}),
];
let start = json!(150);
let filtered = filter_incremental(records, "seq", &start);
assert_eq!(filtered.len(), 2);
assert_eq!(filtered[0]["id"], 2);
assert_eq!(filtered[1]["id"], 3);
}
#[test]
fn test_filter_incremental_missing_key_kept() {
let records = vec![
json!({"id": 1}),
json!({"id": 2, "updated_at": "2024-12-01"}),
json!({"id": 3, "updated_at": null}),
];
let start = json!("2024-01-01");
let filtered = filter_incremental(records, "updated_at", &start);
assert_eq!(filtered.len(), 3);
}
#[test]
fn replication_key_parse_forms() {
assert!(!ReplicationKey::parse("updated").unwrap().is_nested());
let dot = ReplicationKey::parse("fields.updated").unwrap();
assert!(dot.is_nested() && !dot.is_pointer());
assert_eq!(dot.as_str(), "fields.updated");
let ptr = ReplicationKey::parse("/a.b/c").unwrap();
assert!(ptr.is_pointer() && ptr.is_nested());
assert!(ReplicationKey::parse(" ").is_err());
assert!(ReplicationKey::parse("a..b").is_err());
assert!(ReplicationKey::parse(".a").is_err());
}
#[test]
fn replication_key_resolves_nested_array_and_pointer() {
let r = json!({
"fields": {"updated": "2024-06-01"},
"items": [{"date": 1}, {"date": 2}],
"a.b": {"c": 7},
"x": 5
});
let k = |s: &str| ReplicationKey::parse(s).unwrap();
assert_eq!(k("fields.updated").resolve(&r), Some(&json!("2024-06-01")));
assert_eq!(k("items.1.date").resolve(&r), Some(&json!(2)));
assert_eq!(k("items.x.date").resolve(&r), None);
assert_eq!(k("items.9.date").resolve(&r), None);
assert_eq!(k("x.y").resolve(&r), None);
assert_eq!(k("/a.b/c").resolve(&r), Some(&json!(7)));
assert_eq!(k("x").resolve(&r), Some(&json!(5)));
let flat = json!({"Account.LastModifiedDate": "2024"});
assert_eq!(
k("Account.LastModifiedDate").resolve(&flat),
Some(&json!("2024"))
);
assert_eq!(
ReplicationKey::top_level("a.b").resolve(&r),
Some(&json!({"c": 7}))
);
}
#[test]
fn filter_incremental_path_nested_and_policies() {
let records = || {
vec![
json!({"id": 1, "fields": {"updated": "2024-01-01"}}),
json!({"id": 2, "fields": {"updated": "2024-12-01"}}),
json!({"id": 3, "fields": {}}),
json!({"id": 4, "fields": {"updated": 5}}),
]
};
let key = ReplicationKey::parse("fields.updated").unwrap();
let start = json!("2024-06-01");
let keep = filter_incremental_path(records(), &key, &start, OnMissingKey::Keep).unwrap();
let ids: Vec<i64> = keep
.records
.iter()
.map(|r| r["id"].as_i64().unwrap())
.collect();
assert_eq!(ids, vec![2, 3, 4]);
assert_eq!(keep.missing, 1);
let drop = filter_incremental_path(records(), &key, &start, OnMissingKey::Drop).unwrap();
assert_eq!(drop.records.len(), 2);
assert_eq!(drop.missing, 1);
let err = filter_incremental_path(records(), &key, &start, OnMissingKey::Fail).unwrap_err();
assert!(err.to_string().contains("fields.updated"), "{err}");
}
#[test]
fn max_replication_value_path_skips_missing_and_null() {
let key = ReplicationKey::parse("fields.updated").unwrap();
let records = vec![
json!({"fields": {"updated": "2024-01-01"}}),
json!({"fields": {"updated": null}}),
json!({"fields": {"updated": "2024-12-01"}}),
json!({}),
];
assert_eq!(
max_replication_value_path(&records, &key),
Some(&json!("2024-12-01"))
);
assert!(max_replication_value_path(&records[1..2], &key).is_none());
}
#[test]
fn on_missing_key_serde() {
assert_eq!(OnMissingKey::default(), OnMissingKey::Keep);
let v: OnMissingKey = serde_json::from_value(json!("fail")).unwrap();
assert_eq!(v, OnMissingKey::Fail);
}
#[test]
fn test_filter_incremental_equal_excluded() {
let records = vec![
json!({"id": 1, "updated_at": "2024-06-01"}),
json!({"id": 2, "updated_at": "2024-06-02"}),
];
let start = json!("2024-06-01");
let filtered = filter_incremental(records, "updated_at", &start);
assert_eq!(filtered.len(), 1);
assert_eq!(filtered[0]["id"], 2);
}
#[test]
fn test_max_replication_value_strings() {
let records = vec![
json!({"updated_at": "2024-01-01"}),
json!({"updated_at": "2024-12-01"}),
json!({"updated_at": "2024-06-01"}),
];
let max = max_replication_value(&records, "updated_at").unwrap();
assert_eq!(max, &json!("2024-12-01"));
}
#[test]
fn test_max_replication_value_numbers() {
let records = vec![json!({"seq": 5}), json!({"seq": 10}), json!({"seq": 3})];
let max = max_replication_value(&records, "seq").unwrap();
assert_eq!(max, &json!(10));
}
#[test]
fn test_max_replication_value_empty() {
let records: Vec<Value> = vec![];
assert!(max_replication_value(&records, "updated_at").is_none());
}
#[test]
fn test_max_value_picks_larger_string() {
assert_eq!(
max_value(json!("2024-01-01"), json!("2024-06-01")),
json!("2024-06-01")
);
}
#[test]
fn test_max_value_picks_larger_number() {
assert_eq!(max_value(json!(5), json!(10)), json!(10));
}
#[test]
fn test_max_value_returns_a_on_type_mismatch() {
assert_eq!(max_value(json!("string"), json!(5)), json!("string"));
}
#[test]
fn filter_incremental_keeps_large_integer_beyond_f64_precision() {
let two_pow_53 = 9_007_199_254_740_992_i64; let records = vec![
json!({"id": 1, "seq": two_pow_53 + 1}),
json!({"id": 2, "seq": two_pow_53 + 2}),
];
let start = json!(two_pow_53);
let filtered = filter_incremental(records, "seq", &start);
assert_eq!(
filtered.len(),
2,
"both values are strictly greater than 2^53"
);
}
#[test]
fn json_compare_distinguishes_large_integers() {
let a = json!(9_007_199_254_740_993_i64); let b = json!(9_007_199_254_740_992_i64); assert_eq!(json_compare(&a, &b), Ordering::Greater);
}
#[test]
fn filter_incremental_keeps_records_on_type_mismatch() {
let records = vec![json!({"id": 1, "seq": 20_240_701})];
let start = json!("2024-06-01"); let filtered = filter_incremental(records, "seq", &start);
assert_eq!(filtered.len(), 1, "type mismatch must not silently drop");
}
fn bind(into: BindTarget, template: &str, format: BindFormat) -> ReplicationBind {
ReplicationBind {
into,
name: "updated_after".to_owned(),
template: template.to_owned(),
format,
advance_from: None,
path: None,
value_type: BindValueType::String,
}
}
#[test]
fn bind_defaults_template_to_bare_placeholder() {
let b: ReplicationBind =
serde_json::from_value(json!({ "name": "since" })).expect("deserializes");
assert_eq!(b.into, BindTarget::Query);
assert_eq!(b.template, "${bookmark}");
assert_eq!(b.format, BindFormat::Raw);
assert!(b.advance_from.is_none());
}
#[test]
fn bind_render_raw_string_and_number() {
let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Raw);
assert_eq!(b.render(&json!("2024-06-01")).unwrap(), "2024-06-01");
assert_eq!(b.render(&json!(150)).unwrap(), "150");
}
#[test]
fn bind_render_applies_operator_template() {
let b = bind(BindTarget::Query, "gte|${bookmark}", BindFormat::Raw);
assert_eq!(
b.render(&json!("2024-06-01T00:00:00Z")).unwrap(),
"gte|2024-06-01T00:00:00Z"
);
let l = bind(BindTarget::Query, "[${bookmark} TO *]", BindFormat::Raw);
assert_eq!(l.render(&json!("20240601")).unwrap(), "[20240601 TO *]");
}
#[test]
fn bind_format_iso8601_from_date_and_epoch() {
let b = bind(BindTarget::Header, "${bookmark}", BindFormat::Iso8601);
assert_eq!(
b.render(&json!("2024-06-01")).unwrap(),
"2024-06-01T00:00:00Z"
);
assert_eq!(
b.render(&json!(1_717_200_000)).unwrap(),
"2024-06-01T00:00:00Z"
);
}
#[test]
fn bind_format_epoch_s_and_ms_from_iso() {
let s = bind(BindTarget::Query, "${bookmark}", BindFormat::EpochS);
assert_eq!(
s.render(&json!("2024-06-01T00:00:00Z")).unwrap(),
"1717200000"
);
let ms = bind(BindTarget::Query, "${bookmark}", BindFormat::EpochMs);
assert_eq!(
ms.render(&json!("2024-06-01T00:00:00Z")).unwrap(),
"1717200000000"
);
}
#[test]
fn bind_format_date_truncates_datetime() {
let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Date);
assert_eq!(
b.render(&json!("2024-06-01T12:34:56Z")).unwrap(),
"2024-06-01"
);
}
#[test]
fn bind_format_naive_datetime_assumed_utc() {
let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Iso8601);
assert_eq!(
b.render(&json!("2024-06-01T08:00:00")).unwrap(),
"2024-06-01T08:00:00Z"
);
}
#[test]
fn bind_format_unparseable_string_errors() {
let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Iso8601);
assert!(b.render(&json!("not-a-date")).is_err());
}
#[test]
fn bind_format_raw_rejects_composite() {
let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Raw);
assert!(b.render(&json!({"a": 1})).is_err());
assert!(b.render(&json!(null)).is_err());
}
#[test]
fn bind_validate_rejects_empty_name_and_missing_placeholder() {
let mut b = bind(BindTarget::Query, "${bookmark}", BindFormat::Raw);
b.name = " ".to_owned();
assert!(b.validate().is_err());
let mut b2 = bind(BindTarget::Query, "no placeholder here", BindFormat::Raw);
b2.name = "since".to_owned();
assert!(b2.validate().is_err());
let ok = bind(BindTarget::Query, "gte|${bookmark}", BindFormat::Raw);
assert!(ok.validate().is_ok());
}
#[test]
fn bind_placement_validation() {
let mut b = bind(BindTarget::Body, "${bookmark}", BindFormat::Raw);
assert!(b.validate().is_ok());
b.path = Some("/a/0/b".into());
assert!(b.validate().unwrap_err().to_string().contains("not both"));
b.name.clear();
assert!(b.validate().is_ok());
b.path = Some("a".into());
assert!(b.validate().is_err());
b.path = Some("/".into());
assert!(b.validate().is_err());
b.path = None;
assert!(
b.validate()
.unwrap_err()
.to_string()
.contains("needs `name`")
);
let mut q = bind(BindTarget::Query, "${bookmark}", BindFormat::Raw);
q.path = Some("/a".into());
assert!(
q.validate()
.unwrap_err()
.to_string()
.contains("only to `into: body`")
);
let parsed: ReplicationBind = serde_json::from_value(json!({
"into": "body", "path": "/f/0/v", "value_type": "number"
}))
.unwrap();
assert!(parsed.validate().is_ok());
assert_eq!(parsed.value_type, BindValueType::Number);
}
#[test]
fn value_type_conversion() {
assert_eq!(BindValueType::String.to_value("5").unwrap(), json!("5"));
assert_eq!(BindValueType::Number.to_value("5").unwrap(), json!(5));
assert_eq!(
BindValueType::Number
.to_value("18446744073709551615")
.unwrap(),
json!(18_446_744_073_709_551_615_u64)
);
assert_eq!(BindValueType::Number.to_value("1.5").unwrap(), json!(1.5));
assert!(BindValueType::Number.to_value("x").is_err());
}
#[test]
fn set_body_pointer_rules() {
let mut body = json!({"filterGroups": [{"filters": [{"value": null}]}], "v": {}, "a/b": 1});
set_body_pointer(&mut body, "/filterGroups/0/filters/0/value", json!("x")).unwrap();
assert_eq!(body["filterGroups"][0]["filters"][0]["value"], json!("x"));
set_body_pointer(&mut body, "/v/after", json!("c")).unwrap();
assert_eq!(body["v"]["after"], json!("c"));
set_body_pointer(&mut body, "/top", json!(1)).unwrap();
assert_eq!(body["top"], json!(1));
set_body_pointer(&mut body, "/a~1b", json!(2)).unwrap();
assert_eq!(body["a/b"], json!(2));
set_body_pointer(&mut body, "/v/x~1y~0z", json!(3)).unwrap();
assert_eq!(body["v"]["x/y~z"], json!(3));
assert!(set_body_pointer(&mut body, "/filterGroups/1/filters", json!(1)).is_err());
assert!(set_body_pointer(&mut body, "/missing/leaf", json!(1)).is_err());
assert!(set_body_pointer(&mut body, "/v", json!(1)).is_err());
assert!(set_body_pointer(&mut body, "/filterGroups/0/filters/5", json!(1)).is_err());
assert!(set_body_pointer(&mut body, "nope", json!(1)).is_err());
}
#[test]
fn bind_format_bookmark_bool_raw() {
assert_eq!(
format_bookmark(&json!(true), BindFormat::Raw).unwrap(),
"true"
);
}
#[test]
fn bind_format_non_scalar_bookmark_errors() {
assert!(format_bookmark(&json!({"a": 1}), BindFormat::Iso8601).is_err());
assert!(format_bookmark(&json!(null), BindFormat::EpochS).is_err());
}
#[test]
fn bind_format_out_of_range_epoch_errors() {
assert!(format_bookmark(&json!(i64::MAX), BindFormat::Iso8601).is_err());
}
}