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
120fn sql_literal(v: &Value) -> String {
122 match v {
123 Value::String(s) => format!("'{}'", s.replace('\'', "''")),
124 Value::Number(n) => n.to_string(),
125 Value::Bool(b) => b.to_string(),
126 _ => "NULL".to_owned(),
129 }
130}
131
132#[derive(Debug, Clone, Default, Serialize, Deserialize, schemars::JsonSchema)]
136pub struct WriteSpec {
137 #[serde(default)]
139 pub write_mode: WriteMode,
140 #[serde(default)]
142 pub key: Vec<String>,
143 #[serde(default, skip_serializing_if = "Option::is_none")]
147 pub delete_marker: Option<DeleteMarker>,
148}
149
150impl WriteSpec {
151 pub fn validate(&self) -> Result<(), FaucetError> {
153 if matches!(self.write_mode, WriteMode::Upsert | WriteMode::Delete) && self.key.is_empty() {
154 return Err(FaucetError::Config(format!(
155 "write_mode: {} requires a non-empty `key`",
156 self.write_mode.as_str()
157 )));
158 }
159 Ok(())
160 }
161
162 pub fn dedups_by_key(&self) -> bool {
167 matches!(self.write_mode, WriteMode::Upsert | WriteMode::Delete) && !self.key.is_empty()
168 }
169
170 pub fn is_overwrite(&self) -> bool {
175 matches!(self.write_mode, WriteMode::Overwrite)
176 }
177}
178
179#[derive(Debug, Clone, PartialEq)]
181pub struct KeyTuple(pub Vec<(String, Value)>);
182
183#[derive(Debug, Default)]
187pub struct WritePlan {
188 pub upserts: Vec<Value>,
190 pub deletes: Vec<KeyTuple>,
192 pub failed: Vec<(usize, String)>,
194}
195
196#[derive(Clone)]
197enum Action {
198 Upsert(Value),
199 Delete(KeyTuple),
200}
201
202pub fn plan_writes(page: &[Value], spec: &WriteSpec) -> WritePlan {
206 debug_assert!(
207 matches!(spec.write_mode, WriteMode::Upsert | WriteMode::Delete),
208 "plan_writes is only for Upsert/Delete — Append and Overwrite are routed separately"
209 );
210 let mut plan = WritePlan::default();
211 let mut index: HashMap<String, usize> = HashMap::new();
212 let mut order: Vec<Action> = Vec::new();
213
214 for (i, rec) in page.iter().enumerate() {
215 let key_tuple = match extract_key(rec, &spec.key) {
216 Ok(k) => k,
217 Err(msg) => {
218 plan.failed.push((i, msg));
219 continue;
220 }
221 };
222 let canon = canonical(&key_tuple);
223
224 let is_delete = match spec.write_mode {
225 WriteMode::Delete => true,
226 WriteMode::Upsert => is_delete_marked(rec, spec.delete_marker.as_ref()),
227 WriteMode::Append | WriteMode::Overwrite => false,
228 };
229
230 let action = if is_delete {
231 Action::Delete(key_tuple)
232 } else {
233 Action::Upsert(strip_marker(rec.clone(), spec.delete_marker.as_ref()))
234 };
235
236 match index.get(&canon) {
237 Some(&slot) => order[slot] = action,
238 None => {
239 index.insert(canon, order.len());
240 order.push(action);
241 }
242 }
243 }
244
245 for action in order {
246 match action {
247 Action::Upsert(v) => plan.upserts.push(v),
248 Action::Delete(k) => plan.deletes.push(k),
249 }
250 }
251 plan
252}
253
254fn extract_key(rec: &Value, key: &[String]) -> Result<KeyTuple, String> {
257 let obj = rec
258 .as_object()
259 .ok_or_else(|| "record is not a JSON object".to_string())?;
260 let mut out = Vec::with_capacity(key.len());
261 for col in key {
262 match obj.get(col) {
263 None => return Err(format!("missing key column '{col}'")),
264 Some(Value::Null) => return Err(format!("null value for key column '{col}'")),
265 Some(v) => out.push((col.clone(), v.clone())),
266 }
267 }
268 Ok(KeyTuple(out))
269}
270
271fn is_delete_marked(rec: &Value, marker: Option<&DeleteMarker>) -> bool {
272 let Some(dm) = marker else { return false };
273 let Some(v) = rec.get(&dm.field) else {
274 return false;
275 };
276 let Some(s) = v.as_str() else { return false };
277 dm.values.iter().any(|m| m == s)
278}
279
280fn strip_marker(mut rec: Value, marker: Option<&DeleteMarker>) -> Value {
281 if let (Some(dm), Value::Object(map)) = (marker, &mut rec) {
282 map.remove(&dm.field);
283 }
284 rec
285}
286
287fn canonical(k: &KeyTuple) -> String {
289 let arr: Vec<&Value> = k.0.iter().map(|(_, v)| v).collect();
290 serde_json::to_string(&arr).expect("a Vec<&serde_json::Value> always serializes")
291}
292
293pub fn key_to_doc_id(k: &KeyTuple, separator: &str) -> String {
307 let _ = separator; if k.0.len() == 1 {
309 return match &k.0[0].1 {
310 Value::String(s) => s.clone(),
311 other => other.to_string(),
312 };
313 }
314 let values: Vec<&Value> = k.0.iter().map(|(_, v)| v).collect();
315 serde_json::to_string(&values).expect("a Vec<&serde_json::Value> always serializes")
316}
317
318pub fn key_to_filter(k: &KeyTuple) -> Map<String, Value> {
320 k.0.iter().map(|(c, v)| (c.clone(), v.clone())).collect()
321}
322
323#[cfg(test)]
324mod tests {
325 use super::*;
326 use serde_json::json;
327
328 fn upsert_spec(keys: &[&str]) -> WriteSpec {
329 WriteSpec {
330 write_mode: WriteMode::Upsert,
331 key: keys.iter().map(|s| s.to_string()).collect(),
332 delete_marker: None,
333 }
334 }
335
336 #[test]
337 fn upsert_extracts_key_and_keeps_row() {
338 let plan = plan_writes(&[json!({"id": 1, "name": "a"})], &upsert_spec(&["id"]));
339 assert_eq!(plan.upserts, vec![json!({"id": 1, "name": "a"})]);
340 assert!(plan.deletes.is_empty());
341 assert!(plan.failed.is_empty());
342 }
343
344 #[test]
345 fn key_to_doc_id_single_key_is_plain() {
346 let k = KeyTuple(vec![("id".into(), json!(7))]);
347 assert_eq!(key_to_doc_id(&k, "_"), "7");
348 let k = KeyTuple(vec![("name".into(), json!("alice"))]);
349 assert_eq!(key_to_doc_id(&k, "_"), "alice");
350 }
351
352 #[test]
353 fn key_to_doc_id_composite_is_injective() {
354 let k1 = KeyTuple(vec![("x".into(), json!("a_")), ("y".into(), json!("b"))]);
356 let k2 = KeyTuple(vec![("x".into(), json!("a")), ("y".into(), json!("_b"))]);
357 let id1 = key_to_doc_id(&k1, "_");
358 let id2 = key_to_doc_id(&k2, "_");
359 assert_ne!(id1, id2, "distinct composite keys must map to distinct ids");
360 let k3 = KeyTuple(vec![("x".into(), json!(1)), ("y".into(), json!("2"))]);
362 let k4 = KeyTuple(vec![("x".into(), json!("1")), ("y".into(), json!(2))]);
363 assert_ne!(key_to_doc_id(&k3, "_"), key_to_doc_id(&k4, "_"));
364 }
365
366 #[test]
367 fn missing_key_goes_to_failed_with_original_index() {
368 let plan = plan_writes(
369 &[json!({"id": 1}), json!({"name": "no-key"})],
370 &upsert_spec(&["id"]),
371 );
372 assert_eq!(plan.upserts.len(), 1);
373 assert_eq!(plan.failed.len(), 1);
374 assert_eq!(plan.failed[0].0, 1, "failed row keeps its page index");
375 }
376
377 #[test]
378 fn null_key_value_is_a_failure() {
379 let plan = plan_writes(&[json!({"id": null})], &upsert_spec(&["id"]));
380 assert!(plan.upserts.is_empty());
381 assert_eq!(plan.failed.len(), 1);
382 }
383
384 #[test]
385 fn delete_marker_routes_to_deletes_and_strips_marker() {
386 let spec = WriteSpec {
387 write_mode: WriteMode::Upsert,
388 key: vec!["id".into()],
389 delete_marker: Some(DeleteMarker {
390 field: "__op".into(),
391 values: vec!["d".into()],
392 }),
393 };
394 let plan = plan_writes(
395 &[
396 json!({"id": 1, "name": "a", "__op": "u"}),
397 json!({"id": 2, "__op": "d"}),
398 ],
399 &spec,
400 );
401 assert_eq!(plan.upserts, vec![json!({"id": 1, "name": "a"})]);
402 assert_eq!(plan.deletes.len(), 1);
403 assert_eq!(plan.deletes[0].0, vec![("id".to_string(), json!(2))]);
404 }
405
406 #[test]
407 fn last_write_wins_dedup_keeps_final_upsert() {
408 let plan = plan_writes(
409 &[json!({"id": 1, "v": "old"}), json!({"id": 1, "v": "new"})],
410 &upsert_spec(&["id"]),
411 );
412 assert_eq!(plan.upserts, vec![json!({"id": 1, "v": "new"})]);
413 }
414
415 #[test]
416 fn last_write_wins_delete_after_upsert_is_a_delete() {
417 let spec = WriteSpec {
418 write_mode: WriteMode::Upsert,
419 key: vec!["id".into()],
420 delete_marker: Some(DeleteMarker {
421 field: "__op".into(),
422 values: vec!["d".into()],
423 }),
424 };
425 let plan = plan_writes(
426 &[json!({"id": 1, "__op": "u"}), json!({"id": 1, "__op": "d"})],
427 &spec,
428 );
429 assert!(plan.upserts.is_empty());
430 assert_eq!(plan.deletes.len(), 1);
431 }
432
433 #[test]
434 fn delete_mode_routes_every_row_to_deletes() {
435 let spec = WriteSpec {
436 write_mode: WriteMode::Delete,
437 key: vec!["id".into()],
438 delete_marker: None,
439 };
440 let plan = plan_writes(&[json!({"id": 1}), json!({"id": 2})], &spec);
441 assert!(plan.upserts.is_empty());
442 assert_eq!(plan.deletes.len(), 2);
443 }
444
445 #[test]
446 fn composite_key_tuple_is_ordered() {
447 let plan = plan_writes(
448 &[json!({"a": 1, "b": 2, "v": 9})],
449 &upsert_spec(&["a", "b"]),
450 );
451 assert_eq!(plan.upserts.len(), 1);
452 let plan2 = plan_writes(
453 &[
454 json!({"a": 1, "b": 2, "v": "x"}),
455 json!({"a": 1, "b": 3, "v": "y"}),
456 ],
457 &upsert_spec(&["a", "b"]),
458 );
459 assert_eq!(plan2.upserts.len(), 2, "(1,2) and (1,3) are distinct keys");
460 }
461
462 #[test]
463 fn validate_rejects_upsert_without_key() {
464 let spec = WriteSpec {
465 write_mode: WriteMode::Upsert,
466 key: vec![],
467 delete_marker: None,
468 };
469 assert!(spec.validate().is_err());
470 }
471
472 #[test]
473 fn validate_allows_append_without_key() {
474 assert!(WriteSpec::default().validate().is_ok());
475 }
476
477 #[test]
478 fn dedups_by_key_requires_keyed_upsert_or_delete() {
479 assert!(!WriteSpec::default().dedups_by_key());
480 let upsert = WriteSpec {
481 write_mode: WriteMode::Upsert,
482 key: vec!["id".into()],
483 delete_marker: None,
484 };
485 assert!(upsert.dedups_by_key());
486 let delete = WriteSpec {
487 write_mode: WriteMode::Delete,
488 key: vec!["id".into()],
489 delete_marker: None,
490 };
491 assert!(delete.dedups_by_key());
492 let keyless = WriteSpec {
494 write_mode: WriteMode::Upsert,
495 key: vec![],
496 delete_marker: None,
497 };
498 assert!(!keyless.dedups_by_key());
499 }
500
501 #[test]
502 fn last_write_wins_upsert_after_delete_is_an_upsert() {
503 let spec = WriteSpec {
505 write_mode: WriteMode::Upsert,
506 key: vec!["id".into()],
507 delete_marker: Some(DeleteMarker {
508 field: "__op".into(),
509 values: vec!["d".into()],
510 }),
511 };
512 let plan = plan_writes(
513 &[
514 json!({"id": 1, "__op": "d"}),
515 json!({"id": 1, "v": 9, "__op": "u"}),
516 ],
517 &spec,
518 );
519 assert!(plan.deletes.is_empty());
520 assert_eq!(plan.upserts, vec![json!({"id": 1, "v": 9})]);
521 }
522
523 #[test]
524 fn overwrite_mode_flags_and_needs_no_key() {
525 let spec = WriteSpec {
526 write_mode: WriteMode::Overwrite,
527 ..Default::default()
528 };
529 assert!(spec.is_overwrite());
530 assert!(!spec.dedups_by_key());
531 assert!(spec.validate().is_ok());
533 assert_eq!(WriteMode::Overwrite.as_str(), "overwrite");
534 assert!(!WriteSpec::default().is_overwrite());
536 assert!(!upsert_spec(&["id"]).is_overwrite());
537 }
538
539 #[test]
540 fn overwrite_scope_window_validates_and_renders() {
541 let scope = OverwriteScope::Window {
542 column: "posting_date".into(),
543 from: json!("2024-06-01"),
544 to: json!("2024-07-01"),
545 };
546 assert_eq!(scope.column(), "posting_date");
547 assert!(scope.validate().is_ok());
548 let whr = scope.render_where_literal("\"posting_date\"");
549 assert_eq!(
550 whr,
551 "\"posting_date\" >= '2024-06-01' AND \"posting_date\" < '2024-07-01'"
552 );
553
554 assert!(
556 OverwriteScope::Window {
557 column: " ".into(),
558 from: json!(1),
559 to: json!(2)
560 }
561 .validate()
562 .is_err()
563 );
564 assert!(
565 OverwriteScope::Window {
566 column: "c".into(),
567 from: json!(null),
568 to: json!(2)
569 }
570 .validate()
571 .is_err()
572 );
573 }
574
575 #[test]
576 fn scope_window_number_bounds() {
577 let scope = OverwriteScope::Window {
578 column: "seq".into(),
579 from: json!(100),
580 to: json!(200),
581 };
582 assert!(scope.validate().is_ok());
583 assert_eq!(
584 scope.render_where_literal("`seq`"),
585 "`seq` >= 100 AND `seq` < 200"
586 );
587 }
588
589 #[test]
590 fn scope_literal_escapes_quotes() {
591 let scope = OverwriteScope::Window {
593 column: "c".into(),
594 from: json!("x' OR '1'='1"),
595 to: json!("z"),
596 };
597 let whr = scope.render_where_literal("\"c\"");
598 assert!(whr.contains("'x'' OR ''1''=''1'"), "{whr}");
599 }
600
601 #[test]
602 fn overwrite_deserializes_from_wire() {
603 let spec: WriteSpec = serde_json::from_value(json!({"write_mode": "overwrite"})).unwrap();
604 assert_eq!(spec.write_mode, WriteMode::Overwrite);
605 assert!(spec.is_overwrite());
606 }
607
608 #[test]
609 fn empty_page_produces_empty_plan() {
610 let plan = plan_writes(&[], &upsert_spec(&["id"]));
611 assert!(plan.upserts.is_empty());
612 assert!(plan.deletes.is_empty());
613 assert!(plan.failed.is_empty());
614 }
615
616 #[test]
617 fn delete_mode_dedups_repeated_key() {
618 let spec = WriteSpec {
620 write_mode: WriteMode::Delete,
621 key: vec!["id".into()],
622 delete_marker: None,
623 };
624 let plan = plan_writes(&[json!({"id": 1}), json!({"id": 1})], &spec);
625 assert_eq!(plan.deletes.len(), 1);
626 }
627}