1use crate::error::FaucetError;
4use serde::{Deserialize, Serialize};
5use serde_json::{Map, Value};
6use std::collections::HashMap;
7
8#[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 #[default]
22 Append,
23 Upsert,
25 Delete,
27 Overwrite,
32}
33
34impl WriteMode {
35 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#[derive(Debug, Clone, Default, Serialize, Deserialize, schemars::JsonSchema, PartialEq, Eq)]
49pub struct DeleteMarker {
50 pub field: String,
52 pub values: Vec<String>,
54}
55
56#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, schemars::JsonSchema)]
62#[serde(rename_all = "snake_case")]
63#[non_exhaustive]
64pub enum OverwriteScope {
65 Window {
67 column: String,
69 from: Value,
71 to: Value,
73 },
74}
75
76impl OverwriteScope {
77 pub fn column(&self) -> &str {
79 match self {
80 OverwriteScope::Window { column, .. } => column,
81 }
82 }
83
84 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 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
120pub 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 _ => "NULL".to_owned(),
146 }
147}
148
149#[derive(Debug, Clone, Default, Serialize, Deserialize, schemars::JsonSchema)]
153pub struct WriteSpec {
154 #[serde(default)]
156 pub write_mode: WriteMode,
157 #[serde(default)]
159 pub key: Vec<String>,
160 #[serde(default, skip_serializing_if = "Option::is_none")]
164 pub delete_marker: Option<DeleteMarker>,
165 #[serde(default, skip_serializing_if = "Option::is_none")]
170 pub rollback: Option<crate::rollback::RollbackWriteSpec>,
171}
172
173impl WriteSpec {
174 pub fn rollback_run_id(&self) -> Option<&str> {
176 self.rollback.as_ref().map(|r| r.run_id.as_str())
177 }
178
179 pub fn journals(&self) -> bool {
181 self.rollback.as_ref().is_some_and(|r| r.journal)
182 }
183
184 pub fn keeps_previous(&self) -> bool {
186 self.rollback.as_ref().is_some_and(|r| r.keep_previous)
187 }
188
189 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 pub fn dedups_by_key(&self) -> bool {
205 matches!(self.write_mode, WriteMode::Upsert | WriteMode::Delete) && !self.key.is_empty()
206 }
207
208 pub fn is_overwrite(&self) -> bool {
213 matches!(self.write_mode, WriteMode::Overwrite)
214 }
215}
216
217#[derive(Debug, Clone, PartialEq)]
219pub struct KeyTuple(pub Vec<(String, Value)>);
220
221#[derive(Debug, Default)]
225pub struct WritePlan {
226 pub upserts: Vec<Value>,
228 pub deletes: Vec<KeyTuple>,
230 pub failed: Vec<(usize, String)>,
232}
233
234#[derive(Clone)]
235enum Action {
236 Upsert(Value),
237 Delete(KeyTuple),
238}
239
240pub 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
292pub fn record_key(rec: &Value, key: &[String]) -> Option<KeyTuple> {
295 extract_key(rec, key).ok()
296}
297
298fn 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
331fn 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
337pub fn key_to_doc_id(k: &KeyTuple, separator: &str) -> String {
351 let _ = separator; 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
362pub 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 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 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 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 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 assert!(spec.validate().is_ok());
586 assert_eq!(WriteMode::Overwrite.as_str(), "overwrite");
587 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 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 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 assert_eq!(sql_literal(&json!("'' ' ''")), "''''' '' '''''");
651 assert_eq!(sql_literal(&json!("")), "''");
653 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 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 let out = sql_literal(&json!("x' OR 1=1 --"));
668 assert_eq!(out, "'x'' OR 1=1 --'");
669 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 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 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}