Skip to main content

faucet_core/
replication.rs

1//! Incremental replication support.
2
3use crate::error::FaucetError;
4use chrono::{DateTime, NaiveDate, NaiveDateTime, Utc};
5use schemars::JsonSchema;
6use serde::{Deserialize, Serialize};
7use serde_json::Value;
8use std::cmp::Ordering;
9
10/// Determines how records are replicated from the source.
11#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize, JsonSchema)]
12#[serde(tag = "type")]
13pub enum ReplicationMethod {
14    /// All records are fetched on every run (default).
15    #[default]
16    FullTable,
17    /// Only records where the `replication_key` field is strictly greater than
18    /// the stored bookmark (`start_replication_value`) are kept.
19    Incremental,
20}
21
22/// What to do with a record whose replication key is missing or `null` (#747).
23#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
24#[serde(rename_all = "snake_case")]
25pub enum OnMissingKey {
26    /// Keep the record, count it, and warn (default): dropping it would be
27    /// silent data loss.
28    #[default]
29    Keep,
30    /// Drop the record (counted and warned, never silent).
31    Drop,
32    /// Fail the run.
33    Fail,
34}
35
36#[derive(Debug, Clone, PartialEq, Eq)]
37enum KeyForm {
38    TopLevel,
39    DotPath(Vec<String>),
40    Pointer,
41}
42
43/// A compiled replication key: a top-level field name, a dot path
44/// (`fields.updated`, numeric segments index arrays), or an RFC 6901 JSON
45/// Pointer (`/fields/updated`, for field names that contain a dot) (#747).
46///
47/// A dot path first tries the whole string as a literal top-level field, so a
48/// flat column literally named `Account.LastModifiedDate` (a CSV header) still
49/// resolves.
50#[derive(Debug, Clone, PartialEq, Eq)]
51pub struct ReplicationKey {
52    raw: String,
53    form: KeyForm,
54}
55
56impl ReplicationKey {
57    /// Parse a user-facing key (see the type docs for the accepted forms).
58    pub fn parse(raw: &str) -> Result<Self, FaucetError> {
59        if raw.trim().is_empty() {
60            return Err(FaucetError::Config(
61                "replication_key must not be empty".to_owned(),
62            ));
63        }
64        if raw.starts_with('/') {
65            return Ok(Self {
66                raw: raw.to_owned(),
67                form: KeyForm::Pointer,
68            });
69        }
70        if !raw.contains('.') {
71            return Ok(Self::top_level(raw));
72        }
73        let segments: Vec<String> = raw.split('.').map(str::to_owned).collect();
74        if segments.iter().any(String::is_empty) {
75            return Err(FaucetError::Config(format!(
76                "replication_key '{raw}': empty path segment (use the JSON Pointer form \
77                 `/a/b` for field names that contain dots)"
78            )));
79        }
80        Ok(Self {
81            raw: raw.to_owned(),
82            form: KeyForm::DotPath(segments),
83        })
84    }
85
86    /// A literal top-level field name, never interpreted as a path.
87    pub fn top_level(name: &str) -> Self {
88        Self {
89            raw: name.to_owned(),
90            form: KeyForm::TopLevel,
91        }
92    }
93
94    /// The key as configured.
95    pub fn as_str(&self) -> &str {
96        &self.raw
97    }
98
99    /// Whether the key is a JSON Pointer (`/a/b`).
100    pub fn is_pointer(&self) -> bool {
101        self.form == KeyForm::Pointer
102    }
103
104    /// Whether the key addresses a nested value (dot path or pointer).
105    pub fn is_nested(&self) -> bool {
106        self.form != KeyForm::TopLevel
107    }
108
109    /// Resolve the key against one record.
110    pub fn resolve<'a>(&self, record: &'a Value) -> Option<&'a Value> {
111        match &self.form {
112            KeyForm::TopLevel => record.get(&self.raw),
113            KeyForm::Pointer => record.pointer(&self.raw),
114            KeyForm::DotPath(segments) => {
115                if let Some(v) = record.get(&self.raw) {
116                    return Some(v);
117                }
118                segments.iter().try_fold(record, |cur, seg| match cur {
119                    Value::Object(m) => m.get(seg),
120                    Value::Array(a) => seg.parse::<usize>().ok().and_then(|i| a.get(i)),
121                    _ => None,
122                })
123            }
124        }
125    }
126
127    fn resolve_present<'a>(&self, record: &'a Value) -> Option<&'a Value> {
128        self.resolve(record).filter(|v| !v.is_null())
129    }
130}
131
132/// The result of [`filter_incremental_path`].
133#[derive(Debug, Clone, PartialEq, Default)]
134pub struct IncrementalFilter {
135    /// Records to write.
136    pub records: Vec<Value>,
137    /// Records whose key was missing or `null` (kept or dropped per policy).
138    pub missing: usize,
139}
140
141/// Filter `records` to those where `key > start`, using a compiled key.
142///
143/// A record whose key is missing or `null` is handled per `on_missing`
144/// (counted in [`IncrementalFilter::missing`] either way); a record whose key
145/// has a different JSON type than `start` is kept with a warning. Returns
146/// `Err` only for [`OnMissingKey::Fail`].
147pub fn filter_incremental_path(
148    records: Vec<Value>,
149    key: &ReplicationKey,
150    start: &Value,
151    on_missing: OnMissingKey,
152) -> Result<IncrementalFilter, FaucetError> {
153    let mut missing = 0usize;
154    let mut kept = Vec::with_capacity(records.len());
155    for r in records {
156        let keep = match key.resolve_present(&r) {
157            None => {
158                missing += 1;
159                match on_missing {
160                    OnMissingKey::Keep => true,
161                    OnMissingKey::Drop => false,
162                    OnMissingKey::Fail => {
163                        return Err(FaucetError::Source(format!(
164                            "incremental replication: a record lacks replication_key '{}' \
165                             (on_missing_key: fail)",
166                            key.as_str()
167                        )));
168                    }
169                }
170            }
171            Some(v) if type_rank(v) != type_rank(start) => {
172                tracing::warn!(
173                    key = key.as_str(),
174                    "incremental replication: record key type does not match the bookmark \
175                     type; keeping the record to avoid silently dropping data"
176                );
177                true
178            }
179            Some(v) => json_gt(v, start),
180        };
181        if keep {
182            kept.push(r);
183        }
184    }
185    Ok(IncrementalFilter {
186        records: kept,
187        missing,
188    })
189}
190
191/// Filter `records` to only those where `record[key] > start`.
192///
193/// `key` is a literal top-level field. Strings compare lexicographically
194/// (ISO-8601 dates compare correctly this way); integers compare exactly
195/// (no `f64` precision loss); floats compare as `f64`.
196///
197/// Records missing the key (or holding `null`) are **kept** and a warning is
198/// logged (#747): dropping them silently is data loss. Likewise a record whose
199/// key value is a *different JSON type* than `start` is kept (#78/#27).
200pub fn filter_incremental(records: Vec<Value>, key: &str, start: &Value) -> Vec<Value> {
201    let out = filter_incremental_path(
202        records,
203        &ReplicationKey::top_level(key),
204        start,
205        OnMissingKey::Keep,
206    )
207    .unwrap_or_default();
208    if out.missing > 0 {
209        tracing::warn!(
210            key,
211            missing = out.missing,
212            "incremental replication: {} record(s) lacked replication_key '{key}'; kept to \
213             avoid silent data loss",
214            out.missing
215        );
216    }
217    out.records
218}
219
220/// Return the maximum non-null value of `key` across all records, if any.
221pub fn max_replication_value_path<'a>(
222    records: &'a [Value],
223    key: &ReplicationKey,
224) -> Option<&'a Value> {
225    records
226        .iter()
227        .filter_map(|r| key.resolve_present(r))
228        .max_by(|a, b| json_compare(a, b))
229}
230
231/// Return the maximum value of `record[key]` across all records, if any.
232pub fn max_replication_value<'a>(records: &'a [Value], key: &str) -> Option<&'a Value> {
233    records
234        .iter()
235        .filter_map(|r| r.get(key))
236        .max_by(|a, b| json_compare(a, b))
237}
238
239/// Return the larger of two replication values using the same ordering as
240/// [`max_replication_value`] (string lexicographic, numeric for numbers,
241/// falling back to `a` on type mismatch).
242pub fn max_value(a: Value, b: Value) -> Value {
243    match json_compare(&a, &b) {
244        Ordering::Less => b,
245        _ => a,
246    }
247}
248
249/// Type-rank for a total ordering across JSON value kinds, so comparisons of
250/// differing types are deterministic instead of collapsing to `Equal`.
251fn type_rank(v: &Value) -> u8 {
252    match v {
253        Value::Null => 0,
254        Value::Bool(_) => 1,
255        Value::Number(_) => 2,
256        Value::String(_) => 3,
257        Value::Array(_) => 4,
258        Value::Object(_) => 5,
259    }
260}
261
262/// Exact integer view of a JSON number (`i64` or `u64`), widened to `i128` so
263/// both halves of the range compare without `f64` precision loss. `None` for
264/// non-integral (floating) numbers.
265fn number_as_i128(n: &serde_json::Number) -> Option<i128> {
266    n.as_i64()
267        .map(i128::from)
268        .or_else(|| n.as_u64().map(i128::from))
269}
270
271/// Total ordering over JSON values used for replication bookmarks.
272///
273/// - Numbers: compared exactly as `i128` when both are integral (so cursors
274///   above 2^53 don't lose precision); otherwise as `f64`, with NaN ordered
275///   last.
276/// - Same-type scalars/containers: natural ordering (strings lexicographic,
277///   bools `false < true`, arrays element-wise, objects by serialized form).
278/// - Different types: ordered by [`type_rank`] so the result is always total.
279pub(crate) fn json_compare(a: &Value, b: &Value) -> Ordering {
280    match (a, b) {
281        (Value::Number(an), Value::Number(bn)) => {
282            match (number_as_i128(an), number_as_i128(bn)) {
283                (Some(ai), Some(bi)) => ai.cmp(&bi),
284                _ => {
285                    let af = an.as_f64().unwrap_or(f64::NAN);
286                    let bf = bn.as_f64().unwrap_or(f64::NAN);
287                    af.partial_cmp(&bf).unwrap_or_else(|| {
288                        // At least one NaN — order NaN last, deterministically.
289                        match (af.is_nan(), bf.is_nan()) {
290                            (false, true) => Ordering::Less,
291                            (true, false) => Ordering::Greater,
292                            _ => Ordering::Equal,
293                        }
294                    })
295                }
296            }
297        }
298        (Value::String(x), Value::String(y)) => x.cmp(y),
299        (Value::Bool(x), Value::Bool(y)) => x.cmp(y),
300        (Value::Null, Value::Null) => Ordering::Equal,
301        (Value::Array(x), Value::Array(y)) => {
302            for (xi, yi) in x.iter().zip(y.iter()) {
303                let c = json_compare(xi, yi);
304                if c != Ordering::Equal {
305                    return c;
306                }
307            }
308            x.len().cmp(&y.len())
309        }
310        // Objects have no natural order; use the serialized form for a stable
311        // total order (objects as replication keys are pathological).
312        (Value::Object(_), Value::Object(_)) => a.to_string().cmp(&b.to_string()),
313        // Different JSON types — order by type rank so comparison is total.
314        _ => type_rank(a).cmp(&type_rank(b)),
315    }
316}
317
318/// Total-order "greater than" over JSON values, using the same comparison
319/// [`filter_incremental`] applies to replication keys (numbers numerically,
320/// strings lexicographically — so RFC3339 timestamps order correctly). Public
321/// so callers bounding a replay window (e.g. `faucet backfill --to-bookmark`)
322/// compare exactly like the incremental filter does.
323pub fn json_gt(a: &Value, b: &Value) -> bool {
324    json_compare(a, b) == Ordering::Greater
325}
326
327// ── Server-side incremental push-down (#513) ─────────────────────────────────
328
329/// The placeholder replaced by the formatted bookmark inside a
330/// [`ReplicationBind::template`].
331pub const BIND_PLACEHOLDER: &str = "${bookmark}";
332
333fn default_bind_template() -> String {
334    BIND_PLACEHOLDER.to_owned()
335}
336
337/// Where a rendered bookmark is injected into the outgoing request.
338#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
339#[serde(rename_all = "snake_case")]
340pub enum BindTarget {
341    /// A query-string parameter (default) — e.g. `?updated_after=…`.
342    #[default]
343    Query,
344    /// A request header — e.g. `If-Modified-Since: …`.
345    Header,
346    /// A field of the JSON request body (POST-search APIs): a top-level field
347    /// named by `name`, or any existing location addressed by a JSON Pointer
348    /// `path` (#748).
349    Body,
350    /// A `{name}` placeholder in the request path.
351    Path,
352}
353
354/// How the bookmark value is formatted before it is substituted into the
355/// [`ReplicationBind::template`].
356///
357/// For every non-[`Raw`](BindFormat::Raw) format the bookmark is first parsed
358/// into an instant: a string is read as RFC 3339, a bare `YYYY-MM-DD` date
359/// (midnight UTC), or a naive `YYYY-MM-DDTHH:MM:SS` (assumed UTC); a JSON
360/// number is read as **epoch seconds**. It is then re-emitted in the target
361/// representation, so `epoch_ms` ← ISO string and `iso8601` ← epoch number both
362/// work.
363#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
364#[serde(rename_all = "snake_case")]
365pub enum BindFormat {
366    /// Emit the scalar verbatim (string as-is, number as its decimal form).
367    /// The default; no timestamp parsing.
368    #[default]
369    Raw,
370    /// RFC 3339 / ISO-8601 UTC timestamp, e.g. `2024-06-01T00:00:00Z`.
371    Iso8601,
372    /// Unix epoch **seconds** (integer).
373    EpochS,
374    /// Unix epoch **milliseconds** (integer).
375    EpochMs,
376    /// Calendar date `YYYY-MM-DD` (UTC).
377    Date,
378}
379
380/// The JSON type a body-target bind writes (#748).
381#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize, JsonSchema)]
382#[serde(rename_all = "snake_case")]
383pub enum BindValueType {
384    /// A JSON string (default).
385    #[default]
386    String,
387    /// A JSON number (e.g. an `epoch_ms` bookmark an API requires unquoted).
388    Number,
389}
390
391impl BindValueType {
392    /// Convert a rendered bind value into the JSON value to write.
393    pub fn to_value(self, rendered: &str) -> Result<Value, FaucetError> {
394        match self {
395            Self::String => Ok(Value::String(rendered.to_owned())),
396            Self::Number => {
397                let n: Option<serde_json::Number> = rendered
398                    .parse::<i64>()
399                    .map(serde_json::Number::from)
400                    .ok()
401                    .or_else(|| rendered.parse::<u64>().ok().map(serde_json::Number::from))
402                    .or_else(|| {
403                        rendered
404                            .parse::<f64>()
405                            .ok()
406                            .and_then(serde_json::Number::from_f64)
407                    });
408                n.map(Value::Number).ok_or_else(|| {
409                    FaucetError::Source(format!(
410                        "bind: rendered value '{rendered}' is not a number (value_type: number)"
411                    ))
412                })
413            }
414        }
415    }
416}
417
418fn unescape_pointer_token(token: &str) -> String {
419    token.replace("~1", "/").replace("~0", "~")
420}
421
422/// Write `value` into `body` at the RFC 6901 JSON Pointer `pointer` (#748).
423///
424/// The pointer must address an existing scalar (or `null`) — or a missing
425/// final key whose parent is an existing object. Intermediate objects and
426/// array elements are never created, and an object/array target is refused
427/// (a bind replaces a value; it does not merge).
428pub fn set_body_pointer(body: &mut Value, pointer: &str, value: Value) -> Result<(), FaucetError> {
429    if !pointer.starts_with('/') {
430        return Err(FaucetError::Config(format!(
431            "JSON Pointer '{pointer}' must start with '/'"
432        )));
433    }
434    if let Some(slot) = body.pointer_mut(pointer) {
435        if slot.is_object() || slot.is_array() {
436            return Err(FaucetError::Source(format!(
437                "request body location '{pointer}' holds an object or array; a bind replaces \
438                 a scalar value"
439            )));
440        }
441        *slot = value;
442        return Ok(());
443    }
444    let cut = pointer.rfind('/').unwrap_or(0);
445    let (parent, leaf) = (&pointer[..cut], &pointer[cut + 1..]);
446    let parent_value = if parent.is_empty() {
447        Some(body)
448    } else {
449        body.pointer_mut(parent)
450    };
451    match parent_value {
452        Some(Value::Object(map)) => {
453            map.insert(unescape_pointer_token(leaf), value);
454            Ok(())
455        }
456        _ => Err(FaucetError::Source(format!(
457            "request body has no location '{pointer}' (the pointer must resolve to an \
458             existing value, or to a new key of an existing object)"
459        ))),
460    }
461}
462
463/// Load-time check of a bind's placement: a body bind needs exactly one of
464/// `name` / `path`; every other target needs `name` and refuses `path`.
465pub(crate) fn validate_bind_placement(
466    what: &str,
467    into: BindTarget,
468    name: &str,
469    path: Option<&str>,
470) -> Result<(), FaucetError> {
471    let has_name = !name.trim().is_empty();
472    match (into, path) {
473        (BindTarget::Body, Some(p)) => {
474            if has_name {
475                return Err(FaucetError::Config(format!(
476                    "{what}: set either `name` (top-level body field) or `path` (JSON Pointer), \
477                     not both"
478                )));
479            }
480            if !p.starts_with('/') || p.len() < 2 {
481                return Err(FaucetError::Config(format!(
482                    "{what}: `path` must be a JSON Pointer such as `/filters/0/value`, got '{p}'"
483                )));
484            }
485            Ok(())
486        }
487        (_, Some(_)) => Err(FaucetError::Config(format!(
488            "{what}: `path` applies only to `into: body`"
489        ))),
490        (BindTarget::Body, None) if !has_name => Err(FaucetError::Config(format!(
491            "{what}: `into: body` needs `name` (top-level field) or `path` (JSON Pointer)"
492        ))),
493        (_, None) if !has_name => Err(FaucetError::Config(format!(
494            "{what}: `name` must not be empty"
495        ))),
496        _ => Ok(()),
497    }
498}
499
500/// Declarative binding of the stored bookmark into the **outgoing request** —
501/// "server-side incremental push-down" (#513).
502///
503/// Today faucet tracks bookmarks and filters incrementally *client-side* (after
504/// download). A bind lets a source instead push the bookmark into the request
505/// (query param / header / body field / path) so the server returns only the
506/// new rows. The existing client-side [`filter_incremental`] stays active as a
507/// safety net for servers that don't honour the filter exactly.
508#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
509#[serde(deny_unknown_fields)]
510pub struct ReplicationBind {
511    /// Where to place the rendered value.
512    #[serde(default)]
513    pub into: BindTarget,
514    /// The parameter / header / body-field / path-placeholder name. Optional
515    /// only for `into: body` with a `path`.
516    #[serde(default)]
517    pub name: String,
518    /// `into: body` only: an RFC 6901 JSON Pointer into the configured `body`
519    /// (`/filterGroups/0/filters/0/value`) instead of a top-level `name` (#748).
520    /// It must resolve to an existing scalar or to a new key of an existing
521    /// object; array elements are never created.
522    #[serde(default, skip_serializing_if = "Option::is_none")]
523    pub path: Option<String>,
524    /// JSON type written by a body bind: `string` (default) or `number`.
525    #[serde(default)]
526    pub value_type: BindValueType,
527    /// Template rendered with [`BIND_PLACEHOLDER`] (`${bookmark}`) replaced by
528    /// the formatted bookmark. Defaults to the bare `${bookmark}`; set e.g.
529    /// `"gte|${bookmark}"` (an operator-prefixed filter) or `"[${bookmark} TO *]"` (Lucene).
530    #[serde(default = "default_bind_template")]
531    pub template: String,
532    /// How to format the bookmark before substitution.
533    #[serde(default)]
534    pub format: BindFormat,
535    /// Optional JSONPath into the response body to advance the bookmark from,
536    /// instead of `max(record[replication_key])`.
537    #[serde(default, skip_serializing_if = "Option::is_none")]
538    pub advance_from: Option<String>,
539}
540
541impl ReplicationBind {
542    /// Validate the binding at config-load time.
543    pub fn validate(&self) -> Result<(), FaucetError> {
544        validate_bind_placement(
545            "replication bind",
546            self.into,
547            &self.name,
548            self.path.as_deref(),
549        )?;
550        if !self.template.contains(BIND_PLACEHOLDER) {
551            return Err(FaucetError::Config(format!(
552                "replication bind: `template` must contain the `{BIND_PLACEHOLDER}` placeholder"
553            )));
554        }
555        Ok(())
556    }
557
558    /// Render the binding for a concrete bookmark: format the value, then
559    /// substitute it into the template.
560    pub fn render(&self, bookmark: &Value) -> Result<String, FaucetError> {
561        let formatted = format_bookmark(bookmark, self.format)?;
562        Ok(self.template.replace(BIND_PLACEHOLDER, &formatted))
563    }
564}
565
566/// Parse a bookmark value into a UTC instant (see [`BindFormat`] for the rules).
567fn bookmark_instant(value: &Value) -> Result<DateTime<Utc>, FaucetError> {
568    match value {
569        Value::String(s) => {
570            let s = s.trim();
571            if let Ok(dt) = DateTime::parse_from_rfc3339(s) {
572                return Ok(dt.with_timezone(&Utc));
573            }
574            if let Ok(d) = NaiveDate::parse_from_str(s, "%Y-%m-%d")
575                && let Some(ndt) = d.and_hms_opt(0, 0, 0)
576            {
577                return Ok(DateTime::<Utc>::from_naive_utc_and_offset(ndt, Utc));
578            }
579            if let Ok(ndt) = NaiveDateTime::parse_from_str(s, "%Y-%m-%dT%H:%M:%S") {
580                return Ok(DateTime::<Utc>::from_naive_utc_and_offset(ndt, Utc));
581            }
582            Err(FaucetError::Config(format!(
583                "replication bind: cannot parse bookmark '{s}' as a timestamp \
584                 (expected RFC 3339, YYYY-MM-DD, or YYYY-MM-DDTHH:MM:SS)"
585            )))
586        }
587        Value::Number(n) => {
588            let secs = n.as_i64().or_else(|| n.as_f64().map(|f| f as i64));
589            secs.and_then(|s| DateTime::<Utc>::from_timestamp(s, 0))
590                .ok_or_else(|| {
591                    FaucetError::Config(format!(
592                        "replication bind: numeric bookmark {n} is out of range for epoch seconds"
593                    ))
594                })
595        }
596        other => Err(FaucetError::Config(format!(
597            "replication bind: bookmark must be a string or number, got {other}"
598        ))),
599    }
600}
601
602/// Parse a bookmark value into a UTC instant (public wrapper over the internal
603/// parser; see [`BindFormat`] for the accepted forms). Used by datetime window
604/// slicing (#527) to resolve the sweep's start bound from the stored bookmark.
605pub fn parse_instant(value: &Value) -> Result<DateTime<Utc>, FaucetError> {
606    bookmark_instant(value)
607}
608
609/// Format an already-resolved UTC instant per [`BindFormat`]. Unlike
610/// [`format_bookmark`] (which takes an arbitrary scalar and, for [`BindFormat::Raw`],
611/// echoes it verbatim), this always has a real instant, so `Raw` and `Iso8601`
612/// both emit an RFC 3339 UTC timestamp. Used to render window boundaries (#527).
613pub fn format_instant(dt: DateTime<Utc>, format: BindFormat) -> String {
614    match format {
615        BindFormat::Raw | BindFormat::Iso8601 => {
616            dt.to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
617        }
618        BindFormat::Date => dt.format("%Y-%m-%d").to_string(),
619        BindFormat::EpochS => dt.timestamp().to_string(),
620        BindFormat::EpochMs => dt.timestamp_millis().to_string(),
621    }
622}
623
624/// Format a bookmark value per [`BindFormat`].
625pub fn format_bookmark(value: &Value, format: BindFormat) -> Result<String, FaucetError> {
626    match format {
627        BindFormat::Raw => match value {
628            Value::String(s) => Ok(s.clone()),
629            Value::Number(n) => Ok(n.to_string()),
630            Value::Bool(b) => Ok(b.to_string()),
631            other => Err(FaucetError::Config(format!(
632                "replication bind: cannot render {other} as a raw scalar"
633            ))),
634        },
635        BindFormat::Iso8601 => {
636            Ok(bookmark_instant(value)?.to_rfc3339_opts(chrono::SecondsFormat::Secs, true))
637        }
638        BindFormat::Date => Ok(bookmark_instant(value)?.format("%Y-%m-%d").to_string()),
639        BindFormat::EpochS => Ok(bookmark_instant(value)?.timestamp().to_string()),
640        BindFormat::EpochMs => Ok(bookmark_instant(value)?.timestamp_millis().to_string()),
641    }
642}
643
644#[cfg(test)]
645mod tests {
646    use super::*;
647    use serde_json::json;
648
649    #[test]
650    fn test_filter_incremental_strings() {
651        let records = vec![
652            json!({"id": 1, "updated_at": "2024-01-01"}),
653            json!({"id": 2, "updated_at": "2024-06-01"}),
654            json!({"id": 3, "updated_at": "2024-12-01"}),
655        ];
656        let start = json!("2024-06-01");
657        let filtered = filter_incremental(records, "updated_at", &start);
658        assert_eq!(filtered.len(), 1);
659        assert_eq!(filtered[0]["id"], 3);
660    }
661
662    #[test]
663    fn test_filter_incremental_numbers() {
664        let records = vec![
665            json!({"id": 1, "seq": 100}),
666            json!({"id": 2, "seq": 200}),
667            json!({"id": 3, "seq": 300}),
668        ];
669        let start = json!(150);
670        let filtered = filter_incremental(records, "seq", &start);
671        assert_eq!(filtered.len(), 2);
672        assert_eq!(filtered[0]["id"], 2);
673        assert_eq!(filtered[1]["id"], 3);
674    }
675
676    #[test]
677    fn test_filter_incremental_missing_key_kept() {
678        // #747: a record without the key used to be dropped silently.
679        let records = vec![
680            json!({"id": 1}),
681            json!({"id": 2, "updated_at": "2024-12-01"}),
682            json!({"id": 3, "updated_at": null}),
683        ];
684        let start = json!("2024-01-01");
685        let filtered = filter_incremental(records, "updated_at", &start);
686        assert_eq!(filtered.len(), 3);
687    }
688
689    #[test]
690    fn replication_key_parse_forms() {
691        assert!(!ReplicationKey::parse("updated").unwrap().is_nested());
692        let dot = ReplicationKey::parse("fields.updated").unwrap();
693        assert!(dot.is_nested() && !dot.is_pointer());
694        assert_eq!(dot.as_str(), "fields.updated");
695        let ptr = ReplicationKey::parse("/a.b/c").unwrap();
696        assert!(ptr.is_pointer() && ptr.is_nested());
697        assert!(ReplicationKey::parse(" ").is_err());
698        assert!(ReplicationKey::parse("a..b").is_err());
699        assert!(ReplicationKey::parse(".a").is_err());
700    }
701
702    #[test]
703    fn replication_key_resolves_nested_array_and_pointer() {
704        let r = json!({
705            "fields": {"updated": "2024-06-01"},
706            "items": [{"date": 1}, {"date": 2}],
707            "a.b": {"c": 7},
708            "x": 5
709        });
710        let k = |s: &str| ReplicationKey::parse(s).unwrap();
711        assert_eq!(k("fields.updated").resolve(&r), Some(&json!("2024-06-01")));
712        assert_eq!(k("items.1.date").resolve(&r), Some(&json!(2)));
713        assert_eq!(k("items.x.date").resolve(&r), None);
714        assert_eq!(k("items.9.date").resolve(&r), None);
715        assert_eq!(k("x.y").resolve(&r), None);
716        assert_eq!(k("/a.b/c").resolve(&r), Some(&json!(7)));
717        assert_eq!(k("x").resolve(&r), Some(&json!(5)));
718        let flat = json!({"Account.LastModifiedDate": "2024"});
719        assert_eq!(
720            k("Account.LastModifiedDate").resolve(&flat),
721            Some(&json!("2024"))
722        );
723        assert_eq!(
724            ReplicationKey::top_level("a.b").resolve(&r),
725            Some(&json!({"c": 7}))
726        );
727    }
728
729    #[test]
730    fn filter_incremental_path_nested_and_policies() {
731        let records = || {
732            vec![
733                json!({"id": 1, "fields": {"updated": "2024-01-01"}}),
734                json!({"id": 2, "fields": {"updated": "2024-12-01"}}),
735                json!({"id": 3, "fields": {}}),
736                json!({"id": 4, "fields": {"updated": 5}}),
737            ]
738        };
739        let key = ReplicationKey::parse("fields.updated").unwrap();
740        let start = json!("2024-06-01");
741        let keep = filter_incremental_path(records(), &key, &start, OnMissingKey::Keep).unwrap();
742        let ids: Vec<i64> = keep
743            .records
744            .iter()
745            .map(|r| r["id"].as_i64().unwrap())
746            .collect();
747        assert_eq!(ids, vec![2, 3, 4]);
748        assert_eq!(keep.missing, 1);
749        let drop = filter_incremental_path(records(), &key, &start, OnMissingKey::Drop).unwrap();
750        assert_eq!(drop.records.len(), 2);
751        assert_eq!(drop.missing, 1);
752        let err = filter_incremental_path(records(), &key, &start, OnMissingKey::Fail).unwrap_err();
753        assert!(err.to_string().contains("fields.updated"), "{err}");
754    }
755
756    #[test]
757    fn max_replication_value_path_skips_missing_and_null() {
758        let key = ReplicationKey::parse("fields.updated").unwrap();
759        let records = vec![
760            json!({"fields": {"updated": "2024-01-01"}}),
761            json!({"fields": {"updated": null}}),
762            json!({"fields": {"updated": "2024-12-01"}}),
763            json!({}),
764        ];
765        assert_eq!(
766            max_replication_value_path(&records, &key),
767            Some(&json!("2024-12-01"))
768        );
769        assert!(max_replication_value_path(&records[1..2], &key).is_none());
770    }
771
772    #[test]
773    fn on_missing_key_serde() {
774        assert_eq!(OnMissingKey::default(), OnMissingKey::Keep);
775        let v: OnMissingKey = serde_json::from_value(json!("fail")).unwrap();
776        assert_eq!(v, OnMissingKey::Fail);
777    }
778
779    #[test]
780    fn test_filter_incremental_equal_excluded() {
781        let records = vec![
782            json!({"id": 1, "updated_at": "2024-06-01"}),
783            json!({"id": 2, "updated_at": "2024-06-02"}),
784        ];
785        let start = json!("2024-06-01");
786        let filtered = filter_incremental(records, "updated_at", &start);
787        assert_eq!(filtered.len(), 1);
788        assert_eq!(filtered[0]["id"], 2);
789    }
790
791    #[test]
792    fn test_max_replication_value_strings() {
793        let records = vec![
794            json!({"updated_at": "2024-01-01"}),
795            json!({"updated_at": "2024-12-01"}),
796            json!({"updated_at": "2024-06-01"}),
797        ];
798        let max = max_replication_value(&records, "updated_at").unwrap();
799        assert_eq!(max, &json!("2024-12-01"));
800    }
801
802    #[test]
803    fn test_max_replication_value_numbers() {
804        let records = vec![json!({"seq": 5}), json!({"seq": 10}), json!({"seq": 3})];
805        let max = max_replication_value(&records, "seq").unwrap();
806        assert_eq!(max, &json!(10));
807    }
808
809    #[test]
810    fn test_max_replication_value_empty() {
811        let records: Vec<Value> = vec![];
812        assert!(max_replication_value(&records, "updated_at").is_none());
813    }
814
815    #[test]
816    fn test_max_value_picks_larger_string() {
817        assert_eq!(
818            max_value(json!("2024-01-01"), json!("2024-06-01")),
819            json!("2024-06-01")
820        );
821    }
822
823    #[test]
824    fn test_max_value_picks_larger_number() {
825        assert_eq!(max_value(json!(5), json!(10)), json!(10));
826    }
827
828    #[test]
829    fn test_max_value_returns_a_on_type_mismatch() {
830        // String outranks Number in the total type-rank ordering, so the
831        // larger (a) is returned.
832        assert_eq!(max_value(json!("string"), json!(5)), json!("string"));
833    }
834
835    #[test]
836    fn filter_incremental_keeps_large_integer_beyond_f64_precision() {
837        // Regression for #78/#27: integer cursors above 2^53 lose precision
838        // when compared as f64, so a genuinely-greater value compared Equal
839        // and was silently dropped.
840        let two_pow_53 = 9_007_199_254_740_992_i64; // 2^53
841        let records = vec![
842            json!({"id": 1, "seq": two_pow_53 + 1}),
843            json!({"id": 2, "seq": two_pow_53 + 2}),
844        ];
845        let start = json!(two_pow_53);
846        let filtered = filter_incremental(records, "seq", &start);
847        assert_eq!(
848            filtered.len(),
849            2,
850            "both values are strictly greater than 2^53"
851        );
852    }
853
854    #[test]
855    fn json_compare_distinguishes_large_integers() {
856        let a = json!(9_007_199_254_740_993_i64); // 2^53 + 1
857        let b = json!(9_007_199_254_740_992_i64); // 2^53
858        assert_eq!(json_compare(&a, &b), Ordering::Greater);
859    }
860
861    #[test]
862    fn filter_incremental_keeps_records_on_type_mismatch() {
863        // Regression for #78/#27: a bookmark/key type mismatch must not be
864        // silently treated as "not greater" and the record dropped — that is
865        // data loss. Keep the record instead.
866        let records = vec![json!({"id": 1, "seq": 20_240_701})];
867        let start = json!("2024-06-01"); // string bookmark vs numeric key
868        let filtered = filter_incremental(records, "seq", &start);
869        assert_eq!(filtered.len(), 1, "type mismatch must not silently drop");
870    }
871
872    // ── ReplicationBind (#513) ──────────────────────────────────────────────
873
874    fn bind(into: BindTarget, template: &str, format: BindFormat) -> ReplicationBind {
875        ReplicationBind {
876            into,
877            name: "updated_after".to_owned(),
878            template: template.to_owned(),
879            format,
880            advance_from: None,
881            path: None,
882            value_type: BindValueType::String,
883        }
884    }
885
886    #[test]
887    fn bind_defaults_template_to_bare_placeholder() {
888        let b: ReplicationBind =
889            serde_json::from_value(json!({ "name": "since" })).expect("deserializes");
890        assert_eq!(b.into, BindTarget::Query);
891        assert_eq!(b.template, "${bookmark}");
892        assert_eq!(b.format, BindFormat::Raw);
893        assert!(b.advance_from.is_none());
894    }
895
896    #[test]
897    fn bind_render_raw_string_and_number() {
898        let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Raw);
899        assert_eq!(b.render(&json!("2024-06-01")).unwrap(), "2024-06-01");
900        assert_eq!(b.render(&json!(150)).unwrap(), "150");
901    }
902
903    #[test]
904    fn bind_render_applies_operator_template() {
905        let b = bind(BindTarget::Query, "gte|${bookmark}", BindFormat::Raw);
906        assert_eq!(
907            b.render(&json!("2024-06-01T00:00:00Z")).unwrap(),
908            "gte|2024-06-01T00:00:00Z"
909        );
910        // Lucene range form.
911        let l = bind(BindTarget::Query, "[${bookmark} TO *]", BindFormat::Raw);
912        assert_eq!(l.render(&json!("20240601")).unwrap(), "[20240601 TO *]");
913    }
914
915    #[test]
916    fn bind_format_iso8601_from_date_and_epoch() {
917        let b = bind(BindTarget::Header, "${bookmark}", BindFormat::Iso8601);
918        assert_eq!(
919            b.render(&json!("2024-06-01")).unwrap(),
920            "2024-06-01T00:00:00Z"
921        );
922        // Epoch seconds → ISO.
923        assert_eq!(
924            b.render(&json!(1_717_200_000)).unwrap(),
925            "2024-06-01T00:00:00Z"
926        );
927    }
928
929    #[test]
930    fn bind_format_epoch_s_and_ms_from_iso() {
931        let s = bind(BindTarget::Query, "${bookmark}", BindFormat::EpochS);
932        assert_eq!(
933            s.render(&json!("2024-06-01T00:00:00Z")).unwrap(),
934            "1717200000"
935        );
936        let ms = bind(BindTarget::Query, "${bookmark}", BindFormat::EpochMs);
937        assert_eq!(
938            ms.render(&json!("2024-06-01T00:00:00Z")).unwrap(),
939            "1717200000000"
940        );
941    }
942
943    #[test]
944    fn bind_format_date_truncates_datetime() {
945        let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Date);
946        assert_eq!(
947            b.render(&json!("2024-06-01T12:34:56Z")).unwrap(),
948            "2024-06-01"
949        );
950    }
951
952    #[test]
953    fn bind_format_naive_datetime_assumed_utc() {
954        let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Iso8601);
955        assert_eq!(
956            b.render(&json!("2024-06-01T08:00:00")).unwrap(),
957            "2024-06-01T08:00:00Z"
958        );
959    }
960
961    #[test]
962    fn bind_format_unparseable_string_errors() {
963        let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Iso8601);
964        assert!(b.render(&json!("not-a-date")).is_err());
965    }
966
967    #[test]
968    fn bind_format_raw_rejects_composite() {
969        let b = bind(BindTarget::Query, "${bookmark}", BindFormat::Raw);
970        assert!(b.render(&json!({"a": 1})).is_err());
971        assert!(b.render(&json!(null)).is_err());
972    }
973
974    #[test]
975    fn bind_validate_rejects_empty_name_and_missing_placeholder() {
976        let mut b = bind(BindTarget::Query, "${bookmark}", BindFormat::Raw);
977        b.name = "  ".to_owned();
978        assert!(b.validate().is_err());
979
980        let mut b2 = bind(BindTarget::Query, "no placeholder here", BindFormat::Raw);
981        b2.name = "since".to_owned();
982        assert!(b2.validate().is_err());
983
984        let ok = bind(BindTarget::Query, "gte|${bookmark}", BindFormat::Raw);
985        assert!(ok.validate().is_ok());
986    }
987
988    #[test]
989    fn bind_placement_validation() {
990        let mut b = bind(BindTarget::Body, "${bookmark}", BindFormat::Raw);
991        assert!(b.validate().is_ok());
992        b.path = Some("/a/0/b".into());
993        assert!(b.validate().unwrap_err().to_string().contains("not both"));
994        b.name.clear();
995        assert!(b.validate().is_ok());
996        b.path = Some("a".into());
997        assert!(b.validate().is_err());
998        b.path = Some("/".into());
999        assert!(b.validate().is_err());
1000        b.path = None;
1001        assert!(
1002            b.validate()
1003                .unwrap_err()
1004                .to_string()
1005                .contains("needs `name`")
1006        );
1007        let mut q = bind(BindTarget::Query, "${bookmark}", BindFormat::Raw);
1008        q.path = Some("/a".into());
1009        assert!(
1010            q.validate()
1011                .unwrap_err()
1012                .to_string()
1013                .contains("only to `into: body`")
1014        );
1015        let parsed: ReplicationBind = serde_json::from_value(json!({
1016            "into": "body", "path": "/f/0/v", "value_type": "number"
1017        }))
1018        .unwrap();
1019        assert!(parsed.validate().is_ok());
1020        assert_eq!(parsed.value_type, BindValueType::Number);
1021    }
1022
1023    #[test]
1024    fn value_type_conversion() {
1025        assert_eq!(BindValueType::String.to_value("5").unwrap(), json!("5"));
1026        assert_eq!(BindValueType::Number.to_value("5").unwrap(), json!(5));
1027        assert_eq!(
1028            BindValueType::Number
1029                .to_value("18446744073709551615")
1030                .unwrap(),
1031            json!(18_446_744_073_709_551_615_u64)
1032        );
1033        assert_eq!(BindValueType::Number.to_value("1.5").unwrap(), json!(1.5));
1034        assert!(BindValueType::Number.to_value("x").is_err());
1035    }
1036
1037    #[test]
1038    fn set_body_pointer_rules() {
1039        let mut body = json!({"filterGroups": [{"filters": [{"value": null}]}], "v": {}, "a/b": 1});
1040        set_body_pointer(&mut body, "/filterGroups/0/filters/0/value", json!("x")).unwrap();
1041        assert_eq!(body["filterGroups"][0]["filters"][0]["value"], json!("x"));
1042        set_body_pointer(&mut body, "/v/after", json!("c")).unwrap();
1043        assert_eq!(body["v"]["after"], json!("c"));
1044        set_body_pointer(&mut body, "/top", json!(1)).unwrap();
1045        assert_eq!(body["top"], json!(1));
1046        set_body_pointer(&mut body, "/a~1b", json!(2)).unwrap();
1047        assert_eq!(body["a/b"], json!(2));
1048        set_body_pointer(&mut body, "/v/x~1y~0z", json!(3)).unwrap();
1049        assert_eq!(body["v"]["x/y~z"], json!(3));
1050        assert!(set_body_pointer(&mut body, "/filterGroups/1/filters", json!(1)).is_err());
1051        assert!(set_body_pointer(&mut body, "/missing/leaf", json!(1)).is_err());
1052        assert!(set_body_pointer(&mut body, "/v", json!(1)).is_err());
1053        assert!(set_body_pointer(&mut body, "/filterGroups/0/filters/5", json!(1)).is_err());
1054        assert!(set_body_pointer(&mut body, "nope", json!(1)).is_err());
1055    }
1056
1057    #[test]
1058    fn bind_format_bookmark_bool_raw() {
1059        assert_eq!(
1060            format_bookmark(&json!(true), BindFormat::Raw).unwrap(),
1061            "true"
1062        );
1063    }
1064
1065    #[test]
1066    fn bind_format_non_scalar_bookmark_errors() {
1067        // A composite / null bookmark cannot be parsed into an instant.
1068        assert!(format_bookmark(&json!({"a": 1}), BindFormat::Iso8601).is_err());
1069        assert!(format_bookmark(&json!(null), BindFormat::EpochS).is_err());
1070    }
1071
1072    #[test]
1073    fn bind_format_out_of_range_epoch_errors() {
1074        // i64::MAX seconds is far outside chrono's representable range.
1075        assert!(format_bookmark(&json!(i64::MAX), BindFormat::Iso8601).is_err());
1076    }
1077}