Skip to main content

faucet_core/
write_mode.rs

1//! Unified write-mode types + planner shared by every upsert-capable sink.
2
3use crate::error::FaucetError;
4use serde::{Deserialize, Serialize};
5use serde_json::{Map, Value};
6use std::collections::HashMap;
7
8/// Write semantics for a sink. Serialized snake_case. Default `Append`.
9// `#[non_exhaustive]`: this is a deliberate extension point — adding a write
10// mode (as `Overwrite` was, #492) is an additive change that ships as a minor
11// release. Downstream connectors that `match` on it must carry a wildcard arm;
12// the built-in sinks already gate on the specific modes they implement. Kept a
13// plain comment (not rustdoc) so the schema/rustdoc description is unchanged.
14#[derive(
15    Debug, Clone, Copy, Default, Serialize, Deserialize, schemars::JsonSchema, PartialEq, Eq,
16)]
17#[serde(rename_all = "snake_case")]
18#[non_exhaustive]
19pub enum WriteMode {
20    /// Insert every record (today's behaviour).
21    #[default]
22    Append,
23    /// Insert-or-update by `key`; optionally route delete-marked rows to deletes.
24    Upsert,
25    /// Delete by `key` for every record.
26    Delete,
27    /// Replace the entire destination with this run's records (truncate-load /
28    /// full refresh). The old contents are swapped out atomically only after the
29    /// run completes successfully, so a mid-run failure leaves them intact. No
30    /// `key` is required — it is a whole-dataset operation, not a keyed one.
31    Overwrite,
32}
33
34impl WriteMode {
35    /// Lowercase wire name, for error messages.
36    pub fn as_str(&self) -> &'static str {
37        match self {
38            WriteMode::Append => "append",
39            WriteMode::Upsert => "upsert",
40            WriteMode::Delete => "delete",
41            WriteMode::Overwrite => "overwrite",
42        }
43    }
44}
45
46/// Identifies a record as a delete (vs. an upsert) by a marker field's value.
47/// e.g. `{ field: "__op", values: ["d", "delete"] }`.
48#[derive(Debug, Clone, Default, Serialize, Deserialize, schemars::JsonSchema, PartialEq, Eq)]
49pub struct DeleteMarker {
50    /// Field name whose value flags a delete.
51    pub field: String,
52    /// Values of `field` that mean "this row is a delete".
53    pub values: Vec<String>,
54}
55
56/// Scope for a **scoped/windowed overwrite** (#518): with `write_mode:
57/// overwrite`, replace only the destination rows matching this scope instead of
58/// truncating the whole table. The sink-side sibling of scoped cleanup (#478).
59///
60/// v1 supports a half-open date/number **window** on a single column.
61#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, schemars::JsonSchema)]
62#[serde(rename_all = "snake_case")]
63#[non_exhaustive]
64pub enum OverwriteScope {
65    /// Replace rows where `column` is in the half-open range `[from, to)`.
66    Window {
67        /// The destination column to window on.
68        column: String,
69        /// Inclusive lower bound.
70        from: Value,
71        /// Exclusive upper bound.
72        to: Value,
73    },
74}
75
76impl OverwriteScope {
77    /// The scoped column name.
78    pub fn column(&self) -> &str {
79        match self {
80            OverwriteScope::Window { column, .. } => column,
81        }
82    }
83
84    /// Validate the scope at config-load time.
85    pub fn validate(&self) -> Result<(), FaucetError> {
86        match self {
87            OverwriteScope::Window { column, from, to } => {
88                if column.trim().is_empty() {
89                    return Err(FaucetError::Config(
90                        "overwrite scope: window `column` must not be empty".into(),
91                    ));
92                }
93                if from.is_null() || to.is_null() {
94                    return Err(FaucetError::Config(
95                        "overwrite scope: window `from`/`to` must not be null".into(),
96                    ));
97                }
98                Ok(())
99            }
100        }
101    }
102
103    /// Render a SQL `WHERE` predicate over the (already-quoted) column using
104    /// escaped SQL **literals** for the bounds. Literals (not bind params) are
105    /// used so the engine coerces a string bound to the column's real type
106    /// (`date >= '2024-06-01'` works; a text-typed bind param would not). String
107    /// literals have their single quotes doubled, so a crafted bound cannot
108    /// break out of the quotes.
109    pub fn render_where_literal(&self, quoted_col: &str) -> String {
110        match self {
111            OverwriteScope::Window { from, to, .. } => format!(
112                "{quoted_col} >= {} AND {quoted_col} < {}",
113                sql_literal(from),
114                sql_literal(to)
115            ),
116        }
117    }
118}
119
120/// Render a JSON scalar as an **ANSI-standard** SQL literal: strings are
121/// single-quoted with embedded quotes doubled (`''`), numbers and booleans pass
122/// through in their canonical form, and anything else becomes `NULL`.
123///
124/// This is the house escaper for every dialect that follows the ANSI rule —
125/// Postgres (with `standard_conforming_strings`, the default since 9.1), BigQuery
126/// string literals built by the overwrite/cleanup planners, MySQL under
127/// `NO_BACKSLASH_ESCAPES`, SQL Server, and SQLite. **Public so connectors escape
128/// through one audited implementation instead of hand-rolling a third
129/// convention** (#654 M14): the copies that existed before it was exported each
130/// invented their own rules, so a fix to one reached none of the others.
131///
132/// A dialect whose string literals treat some *other* character as a
133/// meta-character cannot use this function — doubling quotes alone does not
134/// neutralise a trailing escape character. Such a dialect must keep its own
135/// escaper, and that escaper's doc comment must state which extra characters it
136/// escapes and why (see `faucet_common_clickhouse::sql_literal`, where a
137/// backslash is a literal meta-character).
138pub fn sql_literal(v: &Value) -> String {
139    match v {
140        Value::String(s) => format!("'{}'", s.replace('\'', "''")),
141        Value::Number(n) => n.to_string(),
142        Value::Bool(b) => b.to_string(),
143        // Validated non-null upstream; any residual maps to NULL (never matches
144        // a range comparison, so the delete is a safe no-op rather than wrong).
145        _ => "NULL".to_owned(),
146    }
147}
148
149/// Shared write-mode config, embedded in each upsert-capable sink config via
150/// `#[serde(flatten)]` so `write_mode` / `key` / `delete_marker` appear at the
151/// sink-config top level.
152#[derive(Debug, Clone, Default, Serialize, Deserialize, schemars::JsonSchema)]
153pub struct WriteSpec {
154    /// Append (default), upsert, or delete.
155    #[serde(default)]
156    pub write_mode: WriteMode,
157    /// Key columns. Required and non-empty for upsert/delete; ignored for append.
158    #[serde(default)]
159    pub key: Vec<String>,
160    /// Optional. Upsert only: rows whose `field` matches one of `values` are
161    /// deletes; all others are upserts. The marker field is stripped from
162    /// upsert rows before writing.
163    #[serde(default, skip_serializing_if = "Option::is_none")]
164    pub delete_marker: Option<DeleteMarker>,
165    /// Runtime rollback settings (#706): the run id to journal under, the
166    /// run-id column, and whether to journal before-images / keep a replaced
167    /// table. Injected per invocation by the CLI when a config has a
168    /// top-level `rollback:` block — not meant to be written by hand.
169    #[serde(default, skip_serializing_if = "Option::is_none")]
170    pub rollback: Option<crate::rollback::RollbackWriteSpec>,
171}
172
173impl WriteSpec {
174    /// The run id this invocation writes under, when rollback is on.
175    pub fn rollback_run_id(&self) -> Option<&str> {
176        self.rollback.as_ref().map(|r| r.run_id.as_str())
177    }
178
179    /// Whether upserts/deletes must journal before-images (#706).
180    pub fn journals(&self) -> bool {
181        self.rollback.as_ref().is_some_and(|r| r.journal)
182    }
183
184    /// Whether an overwrite keeps the replaced table (#706).
185    pub fn keeps_previous(&self) -> bool {
186        self.rollback.as_ref().is_some_and(|r| r.keep_previous)
187    }
188
189    /// Validate internal consistency at config-load time.
190    pub fn validate(&self) -> Result<(), FaucetError> {
191        if matches!(self.write_mode, WriteMode::Upsert | WriteMode::Delete) && self.key.is_empty() {
192            return Err(FaucetError::Config(format!(
193                "write_mode: {} requires a non-empty `key`",
194                self.write_mode.as_str()
195            )));
196        }
197        Ok(())
198    }
199
200    /// Whether this spec makes writes converge by key — `write_mode: upsert`
201    /// or `delete` with a non-empty `key`. The canonical implementation of
202    /// [`Sink::dedups_by_key`](crate::Sink::dedups_by_key) for sinks that
203    /// flatten a `WriteSpec` into their config.
204    pub fn dedups_by_key(&self) -> bool {
205        matches!(self.write_mode, WriteMode::Upsert | WriteMode::Delete) && !self.key.is_empty()
206    }
207
208    /// Whether this spec requests full-destination replacement
209    /// ([`WriteMode::Overwrite`]). The canonical implementation of
210    /// [`Sink::is_overwrite`](crate::Sink::is_overwrite) for sinks that flatten
211    /// a `WriteSpec` into their config.
212    pub fn is_overwrite(&self) -> bool {
213        matches!(self.write_mode, WriteMode::Overwrite)
214    }
215}
216
217/// Ordered key column → value pairs, in `key` declaration order.
218#[derive(Debug, Clone, PartialEq)]
219pub struct KeyTuple(pub Vec<(String, Value)>);
220
221/// The partition of a page by write mode. Infallible to build — per-row
222/// failures (missing/null key) land in `failed` with their original page index
223/// so the caller can route them to a DLQ or abort.
224#[derive(Debug, Default)]
225pub struct WritePlan {
226    /// Rows to insert-or-update, deduped (last-write-wins), marker stripped.
227    pub upserts: Vec<Value>,
228    /// Key tuples to delete, deduped.
229    pub deletes: Vec<KeyTuple>,
230    /// `(page_index, message)` for rows whose key could not be extracted.
231    pub failed: Vec<(usize, String)>,
232}
233
234#[derive(Clone)]
235enum Action {
236    Upsert(Value),
237    Delete(KeyTuple),
238}
239
240/// Partition `page` into upserts + deletes per `spec`. The single place all six
241/// sinks share. `WriteMode::Append` should never reach here (callers route
242/// append separately); if it does, every row is treated as an upsert.
243pub fn plan_writes(page: &[Value], spec: &WriteSpec) -> WritePlan {
244    debug_assert!(
245        matches!(spec.write_mode, WriteMode::Upsert | WriteMode::Delete),
246        "plan_writes is only for Upsert/Delete — Append and Overwrite are routed separately"
247    );
248    let mut plan = WritePlan::default();
249    let mut index: HashMap<String, usize> = HashMap::new();
250    let mut order: Vec<Action> = Vec::new();
251
252    for (i, rec) in page.iter().enumerate() {
253        let key_tuple = match extract_key(rec, &spec.key) {
254            Ok(k) => k,
255            Err(msg) => {
256                plan.failed.push((i, msg));
257                continue;
258            }
259        };
260        let canon = canonical(&key_tuple);
261
262        let is_delete = match spec.write_mode {
263            WriteMode::Delete => true,
264            WriteMode::Upsert => is_delete_marked(rec, spec.delete_marker.as_ref()),
265            WriteMode::Append | WriteMode::Overwrite => false,
266        };
267
268        let action = if is_delete {
269            Action::Delete(key_tuple)
270        } else {
271            Action::Upsert(strip_marker(rec.clone(), spec.delete_marker.as_ref()))
272        };
273
274        match index.get(&canon) {
275            Some(&slot) => order[slot] = action,
276            None => {
277                index.insert(canon, order.len());
278                order.push(action);
279            }
280        }
281    }
282
283    for action in order {
284        match action {
285            Action::Upsert(v) => plan.upserts.push(v),
286            Action::Delete(k) => plan.deletes.push(k),
287        }
288    }
289    plan
290}
291
292/// The key tuple of a record in `key` order, or `None` when a key column is
293/// missing or null — the same rule [`plan_writes`] applies.
294pub fn record_key(rec: &Value, key: &[String]) -> Option<KeyTuple> {
295    extract_key(rec, key).ok()
296}
297
298/// Pull the key columns out of a record in `key` order. Missing key or null
299/// key value is an error.
300fn extract_key(rec: &Value, key: &[String]) -> Result<KeyTuple, String> {
301    let obj = rec
302        .as_object()
303        .ok_or_else(|| "record is not a JSON object".to_string())?;
304    let mut out = Vec::with_capacity(key.len());
305    for col in key {
306        match obj.get(col) {
307            None => return Err(format!("missing key column '{col}'")),
308            Some(Value::Null) => return Err(format!("null value for key column '{col}'")),
309            Some(v) => out.push((col.clone(), v.clone())),
310        }
311    }
312    Ok(KeyTuple(out))
313}
314
315fn is_delete_marked(rec: &Value, marker: Option<&DeleteMarker>) -> bool {
316    let Some(dm) = marker else { return false };
317    let Some(v) = rec.get(&dm.field) else {
318        return false;
319    };
320    let Some(s) = v.as_str() else { return false };
321    dm.values.iter().any(|m| m == s)
322}
323
324fn strip_marker(mut rec: Value, marker: Option<&DeleteMarker>) -> Value {
325    if let (Some(dm), Value::Object(map)) = (marker, &mut rec) {
326        map.remove(&dm.field);
327    }
328    rec
329}
330
331/// Stable canonical string for a key tuple, for dedup.
332fn canonical(k: &KeyTuple) -> String {
333    let arr: Vec<&Value> = k.0.iter().map(|(_, v)| v).collect();
334    serde_json::to_string(&arr).expect("a Vec<&serde_json::Value> always serializes")
335}
336
337/// Render a key tuple into a single document id (Elasticsearch `_id`).
338///
339/// A single-column key is rendered as its plain string / JSON form (no
340/// separator can collide). A **composite** key is rendered as a canonical JSON
341/// array of its values rather than a separator-join: a plain join is not
342/// injective — e.g. `["a_", "b"]` and `["a", "_b"]` both collapse to `"a__b"`
343/// under separator `"_"`, silently overwriting two distinct rows with one. JSON
344/// encoding escapes any separator-like characters in the values, so distinct key
345/// tuples always map to distinct ids.
346///
347/// Assumes each key column has a consistent JSON type across records (the
348/// normal case for SQL and CDC sources); it does not disambiguate, e.g., the
349/// integer `7` from the string `"7"` in the same column.
350pub fn key_to_doc_id(k: &KeyTuple, separator: &str) -> String {
351    let _ = separator; // retained for API stability; no separator can collide now
352    if k.0.len() == 1 {
353        return match &k.0[0].1 {
354            Value::String(s) => s.clone(),
355            other => other.to_string(),
356        };
357    }
358    let values: Vec<&Value> = k.0.iter().map(|(_, v)| v).collect();
359    serde_json::to_string(&values).expect("a Vec<&serde_json::Value> always serializes")
360}
361
362/// Build a Mongo/ES filter document `{ col: value, … }` from a key tuple.
363pub fn key_to_filter(k: &KeyTuple) -> Map<String, Value> {
364    k.0.iter().map(|(c, v)| (c.clone(), v.clone())).collect()
365}
366
367#[cfg(test)]
368mod tests {
369    use super::*;
370    use serde_json::json;
371
372    fn upsert_spec(keys: &[&str]) -> WriteSpec {
373        WriteSpec {
374            write_mode: WriteMode::Upsert,
375            key: keys.iter().map(|s| s.to_string()).collect(),
376            delete_marker: None,
377            rollback: None,
378        }
379    }
380
381    #[test]
382    fn upsert_extracts_key_and_keeps_row() {
383        let plan = plan_writes(&[json!({"id": 1, "name": "a"})], &upsert_spec(&["id"]));
384        assert_eq!(plan.upserts, vec![json!({"id": 1, "name": "a"})]);
385        assert!(plan.deletes.is_empty());
386        assert!(plan.failed.is_empty());
387    }
388
389    #[test]
390    fn key_to_doc_id_single_key_is_plain() {
391        let k = KeyTuple(vec![("id".into(), json!(7))]);
392        assert_eq!(key_to_doc_id(&k, "_"), "7");
393        let k = KeyTuple(vec![("name".into(), json!("alice"))]);
394        assert_eq!(key_to_doc_id(&k, "_"), "alice");
395    }
396
397    #[test]
398    fn key_to_doc_id_composite_is_injective() {
399        // ["a_", "b"] and ["a", "_b"] must NOT collide (the F13 separator bug).
400        let k1 = KeyTuple(vec![("x".into(), json!("a_")), ("y".into(), json!("b"))]);
401        let k2 = KeyTuple(vec![("x".into(), json!("a")), ("y".into(), json!("_b"))]);
402        let id1 = key_to_doc_id(&k1, "_");
403        let id2 = key_to_doc_id(&k2, "_");
404        assert_ne!(id1, id2, "distinct composite keys must map to distinct ids");
405        // Mixed types also stay distinct.
406        let k3 = KeyTuple(vec![("x".into(), json!(1)), ("y".into(), json!("2"))]);
407        let k4 = KeyTuple(vec![("x".into(), json!("1")), ("y".into(), json!(2))]);
408        assert_ne!(key_to_doc_id(&k3, "_"), key_to_doc_id(&k4, "_"));
409    }
410
411    #[test]
412    fn missing_key_goes_to_failed_with_original_index() {
413        let plan = plan_writes(
414            &[json!({"id": 1}), json!({"name": "no-key"})],
415            &upsert_spec(&["id"]),
416        );
417        assert_eq!(plan.upserts.len(), 1);
418        assert_eq!(plan.failed.len(), 1);
419        assert_eq!(plan.failed[0].0, 1, "failed row keeps its page index");
420    }
421
422    #[test]
423    fn null_key_value_is_a_failure() {
424        let plan = plan_writes(&[json!({"id": null})], &upsert_spec(&["id"]));
425        assert!(plan.upserts.is_empty());
426        assert_eq!(plan.failed.len(), 1);
427    }
428
429    #[test]
430    fn delete_marker_routes_to_deletes_and_strips_marker() {
431        let spec = WriteSpec {
432            write_mode: WriteMode::Upsert,
433            key: vec!["id".into()],
434            delete_marker: Some(DeleteMarker {
435                field: "__op".into(),
436                values: vec!["d".into()],
437            }),
438            rollback: None,
439        };
440        let plan = plan_writes(
441            &[
442                json!({"id": 1, "name": "a", "__op": "u"}),
443                json!({"id": 2, "__op": "d"}),
444            ],
445            &spec,
446        );
447        assert_eq!(plan.upserts, vec![json!({"id": 1, "name": "a"})]);
448        assert_eq!(plan.deletes.len(), 1);
449        assert_eq!(plan.deletes[0].0, vec![("id".to_string(), json!(2))]);
450    }
451
452    #[test]
453    fn last_write_wins_dedup_keeps_final_upsert() {
454        let plan = plan_writes(
455            &[json!({"id": 1, "v": "old"}), json!({"id": 1, "v": "new"})],
456            &upsert_spec(&["id"]),
457        );
458        assert_eq!(plan.upserts, vec![json!({"id": 1, "v": "new"})]);
459    }
460
461    #[test]
462    fn last_write_wins_delete_after_upsert_is_a_delete() {
463        let spec = WriteSpec {
464            write_mode: WriteMode::Upsert,
465            key: vec!["id".into()],
466            delete_marker: Some(DeleteMarker {
467                field: "__op".into(),
468                values: vec!["d".into()],
469            }),
470            rollback: None,
471        };
472        let plan = plan_writes(
473            &[json!({"id": 1, "__op": "u"}), json!({"id": 1, "__op": "d"})],
474            &spec,
475        );
476        assert!(plan.upserts.is_empty());
477        assert_eq!(plan.deletes.len(), 1);
478    }
479
480    #[test]
481    fn delete_mode_routes_every_row_to_deletes() {
482        let spec = WriteSpec {
483            write_mode: WriteMode::Delete,
484            key: vec!["id".into()],
485            delete_marker: None,
486            rollback: None,
487        };
488        let plan = plan_writes(&[json!({"id": 1}), json!({"id": 2})], &spec);
489        assert!(plan.upserts.is_empty());
490        assert_eq!(plan.deletes.len(), 2);
491    }
492
493    #[test]
494    fn composite_key_tuple_is_ordered() {
495        let plan = plan_writes(
496            &[json!({"a": 1, "b": 2, "v": 9})],
497            &upsert_spec(&["a", "b"]),
498        );
499        assert_eq!(plan.upserts.len(), 1);
500        let plan2 = plan_writes(
501            &[
502                json!({"a": 1, "b": 2, "v": "x"}),
503                json!({"a": 1, "b": 3, "v": "y"}),
504            ],
505            &upsert_spec(&["a", "b"]),
506        );
507        assert_eq!(plan2.upserts.len(), 2, "(1,2) and (1,3) are distinct keys");
508    }
509
510    #[test]
511    fn validate_rejects_upsert_without_key() {
512        let spec = WriteSpec {
513            write_mode: WriteMode::Upsert,
514            key: vec![],
515            delete_marker: None,
516            rollback: None,
517        };
518        assert!(spec.validate().is_err());
519    }
520
521    #[test]
522    fn validate_allows_append_without_key() {
523        assert!(WriteSpec::default().validate().is_ok());
524    }
525
526    #[test]
527    fn dedups_by_key_requires_keyed_upsert_or_delete() {
528        assert!(!WriteSpec::default().dedups_by_key());
529        let upsert = WriteSpec {
530            write_mode: WriteMode::Upsert,
531            key: vec!["id".into()],
532            delete_marker: None,
533            rollback: None,
534        };
535        assert!(upsert.dedups_by_key());
536        let delete = WriteSpec {
537            write_mode: WriteMode::Delete,
538            key: vec!["id".into()],
539            delete_marker: None,
540            rollback: None,
541        };
542        assert!(delete.dedups_by_key());
543        // An (invalid) keyless upsert never claims keyed dedup.
544        let keyless = WriteSpec {
545            write_mode: WriteMode::Upsert,
546            key: vec![],
547            delete_marker: None,
548            rollback: None,
549        };
550        assert!(!keyless.dedups_by_key());
551    }
552
553    #[test]
554    fn last_write_wins_upsert_after_delete_is_an_upsert() {
555        // Inverse of the delete-after-upsert case: [delete, upsert] → upsert wins.
556        let spec = WriteSpec {
557            write_mode: WriteMode::Upsert,
558            key: vec!["id".into()],
559            delete_marker: Some(DeleteMarker {
560                field: "__op".into(),
561                values: vec!["d".into()],
562            }),
563            rollback: None,
564        };
565        let plan = plan_writes(
566            &[
567                json!({"id": 1, "__op": "d"}),
568                json!({"id": 1, "v": 9, "__op": "u"}),
569            ],
570            &spec,
571        );
572        assert!(plan.deletes.is_empty());
573        assert_eq!(plan.upserts, vec![json!({"id": 1, "v": 9})]);
574    }
575
576    #[test]
577    fn overwrite_mode_flags_and_needs_no_key() {
578        let spec = WriteSpec {
579            write_mode: WriteMode::Overwrite,
580            ..Default::default()
581        };
582        assert!(spec.is_overwrite());
583        assert!(!spec.dedups_by_key());
584        // Overwrite is a whole-dataset op — no key required, so validate passes.
585        assert!(spec.validate().is_ok());
586        assert_eq!(WriteMode::Overwrite.as_str(), "overwrite");
587        // Non-overwrite specs report false.
588        assert!(!WriteSpec::default().is_overwrite());
589        assert!(!upsert_spec(&["id"]).is_overwrite());
590    }
591
592    #[test]
593    fn overwrite_scope_window_validates_and_renders() {
594        let scope = OverwriteScope::Window {
595            column: "posting_date".into(),
596            from: json!("2024-06-01"),
597            to: json!("2024-07-01"),
598        };
599        assert_eq!(scope.column(), "posting_date");
600        assert!(scope.validate().is_ok());
601        let whr = scope.render_where_literal("\"posting_date\"");
602        assert_eq!(
603            whr,
604            "\"posting_date\" >= '2024-06-01' AND \"posting_date\" < '2024-07-01'"
605        );
606
607        // Empty column / null bounds are rejected.
608        assert!(
609            OverwriteScope::Window {
610                column: " ".into(),
611                from: json!(1),
612                to: json!(2)
613            }
614            .validate()
615            .is_err()
616        );
617        assert!(
618            OverwriteScope::Window {
619                column: "c".into(),
620                from: json!(null),
621                to: json!(2)
622            }
623            .validate()
624            .is_err()
625        );
626    }
627
628    #[test]
629    fn scope_window_number_bounds() {
630        let scope = OverwriteScope::Window {
631            column: "seq".into(),
632            from: json!(100),
633            to: json!(200),
634        };
635        assert!(scope.validate().is_ok());
636        assert_eq!(
637            scope.render_where_literal("`seq`"),
638            "`seq` >= 100 AND `seq` < 200"
639        );
640    }
641
642    #[test]
643    fn sql_literal_renders_each_scalar_kind_by_the_ansi_rule() {
644        // Strings: single-quoted, embedded quotes doubled — and *only* quotes,
645        // since a backslash is an ordinary character under the ANSI rule.
646        assert_eq!(sql_literal(&json!("plain")), "'plain'");
647        assert_eq!(sql_literal(&json!("O'Brien")), "'O''Brien'");
648        assert_eq!(sql_literal(&json!("a\\b")), "'a\\b'");
649        // Every quote is doubled, not just the first: 5 interior quotes → 10.
650        assert_eq!(sql_literal(&json!("'' ' ''")), "''''' '' '''''");
651        // An empty string stays a valid (empty) literal, not bare quotes-less text.
652        assert_eq!(sql_literal(&json!("")), "''");
653        // Numbers and booleans pass through in canonical form.
654        assert_eq!(sql_literal(&json!(42)), "42");
655        assert_eq!(sql_literal(&json!(-1.5)), "-1.5");
656        assert_eq!(sql_literal(&json!(true)), "true");
657        assert_eq!(sql_literal(&json!(false)), "false");
658        // Null and non-scalars degrade to NULL rather than injecting structure.
659        assert_eq!(sql_literal(&Value::Null), "NULL");
660        assert_eq!(sql_literal(&json!({"a": 1})), "NULL");
661        assert_eq!(sql_literal(&json!([1, 2])), "NULL");
662    }
663
664    #[test]
665    fn sql_literal_neutralises_a_quote_breakout_attempt() {
666        // The whole point of the escaper: a crafted value must stay one literal.
667        let out = sql_literal(&json!("x' OR 1=1 --"));
668        assert_eq!(out, "'x'' OR 1=1 --'");
669        // Exactly one opening and one closing quote delimit the literal: every
670        // interior quote is part of a doubled pair.
671        let interior = &out[1..out.len() - 1];
672        assert!(
673            interior.matches('\'').count().is_multiple_of(2),
674            "unpaired quote inside literal: {out}"
675        );
676    }
677
678    #[test]
679    fn scope_literal_escapes_quotes() {
680        // A crafted bound cannot break out of the string literal.
681        let scope = OverwriteScope::Window {
682            column: "c".into(),
683            from: json!("x' OR '1'='1"),
684            to: json!("z"),
685        };
686        let whr = scope.render_where_literal("\"c\"");
687        assert!(whr.contains("'x'' OR ''1''=''1'"), "{whr}");
688    }
689
690    #[test]
691    fn overwrite_deserializes_from_wire() {
692        let spec: WriteSpec = serde_json::from_value(json!({"write_mode": "overwrite"})).unwrap();
693        assert_eq!(spec.write_mode, WriteMode::Overwrite);
694        assert!(spec.is_overwrite());
695    }
696
697    #[test]
698    fn empty_page_produces_empty_plan() {
699        let plan = plan_writes(&[], &upsert_spec(&["id"]));
700        assert!(plan.upserts.is_empty());
701        assert!(plan.deletes.is_empty());
702        assert!(plan.failed.is_empty());
703    }
704
705    #[test]
706    fn delete_mode_dedups_repeated_key() {
707        // Same key deleted twice in one page collapses to a single delete.
708        let spec = WriteSpec {
709            write_mode: WriteMode::Delete,
710            key: vec!["id".into()],
711            delete_marker: None,
712            rollback: None,
713        };
714        let plan = plan_writes(&[json!({"id": 1}), json!({"id": 1})], &spec);
715        assert_eq!(plan.deletes.len(), 1);
716    }
717}