1use crate::error::FaucetError;
56use crate::traits::Sink;
57use crate::write_mode::KeyTuple;
58use serde::{Deserialize, Serialize};
59use serde_json::Value;
60use std::collections::BTreeMap;
61
62pub const DEFAULT_MAX_KEYS: usize = 100_000;
67
68#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, schemars::JsonSchema)]
75#[serde(rename_all = "snake_case")]
76pub enum CleanupMode {
77 DeleteMissing,
79}
80
81#[derive(Debug, Clone)]
83pub struct CleanupPolicy {
84 pub scope: BTreeMap<String, Value>,
92 pub key: Vec<String>,
94 pub max_keys: usize,
96}
97
98impl CleanupPolicy {
99 pub fn new(
101 scope: BTreeMap<String, Value>,
102 key: Vec<String>,
103 max_keys: usize,
104 ) -> Result<Self, FaucetError> {
105 if scope.is_empty() {
106 return Err(FaucetError::Config(
110 "cleanup: the completeness claim (`complete_for`) must name at least one \
111 column — an empty scope would match every row in the destination"
112 .into(),
113 ));
114 }
115 if key.is_empty() {
116 return Err(FaucetError::Config(
117 "cleanup: requires a non-empty `key` so a written row can be told apart \
118 from a stale one"
119 .into(),
120 ));
121 }
122 if scope.values().any(Value::is_null) {
123 return Err(FaucetError::Config(
124 "cleanup: the completeness claim contains a null value — an unresolved \
125 scope token would delete the wrong rows"
126 .into(),
127 ));
128 }
129 Ok(Self {
130 scope,
131 key,
132 max_keys: max_keys.max(1),
133 })
134 }
135}
136
137#[derive(Debug, Default)]
142pub struct SeenKeys {
143 keys: Vec<KeyTuple>,
144 overflowed: bool,
147}
148
149impl SeenKeys {
150 pub fn new() -> Self {
151 Self::default()
152 }
153
154 pub fn record_page(&mut self, page: &[Value], key: &[String], max_keys: usize) {
159 if self.overflowed {
160 return;
161 }
162 for rec in page {
163 let Some(obj) = rec.as_object() else { continue };
164 let mut tuple = Vec::with_capacity(key.len());
165 let mut complete = true;
166 for k in key {
167 match obj.get(k) {
168 Some(v) if !v.is_null() => tuple.push((k.clone(), v.clone())),
169 _ => {
170 complete = false;
171 break;
172 }
173 }
174 }
175 if !complete {
176 continue;
177 }
178 if self.keys.len() >= max_keys {
179 self.overflowed = true;
180 self.keys.clear(); return;
182 }
183 self.keys.push(KeyTuple(tuple));
184 }
185 }
186
187 pub fn overflowed(&self) -> bool {
189 self.overflowed
190 }
191
192 pub fn len(&self) -> usize {
193 self.keys.len()
194 }
195
196 pub fn is_empty(&self) -> bool {
197 self.keys.is_empty()
198 }
199
200 pub fn keys(&self) -> &[KeyTuple] {
201 &self.keys
202 }
203
204 pub fn overflow_error(&self, max_keys: usize) -> FaucetError {
208 FaucetError::Config(format!(
209 "cleanup: this invocation wrote more than {max_keys} rows in the claimed scope, \
210 so the set of written keys could not be tracked. Nothing was deleted — a \
211 partial delete would remove rows the run actually wrote. Narrow the scope \
212 (a smaller `complete_for`), or raise the ceiling if the destination can take \
213 a delete of this size"
214 ))
215 }
216}
217
218pub struct CleanupTracker<'a, S: Sink + ?Sized> {
238 inner: &'a S,
239 key: Vec<String>,
240 max_keys: usize,
241 seen: std::sync::Mutex<SeenKeys>,
242}
243
244impl<'a, S: Sink + ?Sized> CleanupTracker<'a, S> {
245 pub fn new(inner: &'a S, policy: &CleanupPolicy) -> Self {
246 Self {
247 inner,
248 key: policy.key.clone(),
249 max_keys: policy.max_keys,
250 seen: std::sync::Mutex::new(SeenKeys::new()),
251 }
252 }
253
254 fn record(&self, records: &[Value]) {
255 if let Ok(mut seen) = self.seen.lock() {
256 seen.record_page(records, &self.key, self.max_keys);
257 }
258 }
259
260 pub async fn finish(&self, policy: &CleanupPolicy) -> Result<u64, FaucetError> {
263 let seen = {
267 let mut guard = self
268 .seen
269 .lock()
270 .map_err(|_| FaucetError::Sink("cleanup: key tracker poisoned".into()))?;
271 if guard.overflowed() {
272 return Err(guard.overflow_error(policy.max_keys));
275 }
276 std::mem::take(&mut *guard)
277 };
278 self.inner.cleanup_scope(&policy.scope, &seen).await
279 }
280
281 pub fn tracked(&self) -> usize {
283 self.seen.lock().map(|g| g.len()).unwrap_or(0)
284 }
285}
286
287#[async_trait::async_trait]
288impl<S: Sink + ?Sized> Sink for CleanupTracker<'_, S> {
289 async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
290 let n = self.inner.write_batch(records).await?;
291 self.record(records);
292 Ok(n)
293 }
294
295 async fn write_batch_partial(
296 &self,
297 records: &[Value],
298 ) -> Result<Vec<crate::traits::RowOutcome>, FaucetError> {
299 let out = self.inner.write_batch_partial(records).await?;
300 self.record(records);
303 Ok(out)
304 }
305
306 async fn write_batch_idempotent(
307 &self,
308 records: &[Value],
309 scope: &str,
310 token: &str,
311 ) -> Result<usize, FaucetError> {
312 let n = self
313 .inner
314 .write_batch_idempotent(records, scope, token)
315 .await?;
316 self.record(records);
317 Ok(n)
318 }
319
320 async fn flush(&self) -> Result<(), FaucetError> {
321 self.inner.flush().await
322 }
323
324 fn supports_cleanup(&self) -> bool {
326 self.inner.supports_cleanup()
327 }
328 async fn cleanup_scope(
329 &self,
330 scope: &BTreeMap<String, Value>,
331 seen: &SeenKeys,
332 ) -> Result<u64, FaucetError> {
333 self.inner.cleanup_scope(scope, seen).await
334 }
335 fn supports_idempotent_writes(&self) -> bool {
336 self.inner.supports_idempotent_writes()
337 }
338 async fn last_committed_token(&self, scope: &str) -> Result<Option<String>, FaucetError> {
339 self.inner.last_committed_token(scope).await
340 }
341 fn supported_write_modes(&self) -> &'static [crate::write_mode::WriteMode] {
342 self.inner.supported_write_modes()
343 }
344 fn dedups_by_key(&self) -> bool {
345 self.inner.dedups_by_key()
346 }
347 fn batch_atomicity(&self) -> crate::dlq::BatchAtomicity {
348 self.inner.batch_atomicity()
349 }
350 fn sink_guarantee(&self) -> crate::idempotency::SinkGuarantee {
351 self.inner.sink_guarantee()
352 }
353 async fn current_schema(&self) -> Result<Option<Value>, FaucetError> {
354 self.inner.current_schema().await
355 }
356 fn supports_schema_evolution(&self) -> bool {
357 self.inner.supports_schema_evolution()
358 }
359 async fn evolve_schema(
360 &self,
361 evolution: &crate::drift::SchemaEvolution,
362 ) -> Result<(), FaucetError> {
363 self.inner.evolve_schema(evolution).await
364 }
365 fn config_schema(&self) -> Value {
366 self.inner.config_schema()
367 }
368 fn connector_name(&self) -> &'static str {
369 self.inner.connector_name()
370 }
371 fn dataset_uri(&self) -> String {
372 self.inner.dataset_uri()
373 }
374 async fn local_outputs(&self) -> Vec<crate::local_outputs::LocalOutput> {
375 self.inner.local_outputs().await
376 }
377 fn is_overwrite(&self) -> bool {
378 self.inner.is_overwrite()
379 }
380 async fn begin_overwrite(&self) -> Result<(), FaucetError> {
381 self.inner.begin_overwrite().await
382 }
383 async fn commit_overwrite(&self) -> Result<(), FaucetError> {
384 self.inner.commit_overwrite().await
385 }
386 async fn abort_overwrite(&self) -> Result<(), FaucetError> {
387 self.inner.abort_overwrite().await
388 }
389 async fn complete_run(&self) -> Result<(), FaucetError> {
390 self.inner.complete_run().await
391 }
392}
393
394#[cfg(test)]
395mod tests {
396 use super::*;
397 use serde_json::json;
398
399 fn scope() -> BTreeMap<String, Value> {
400 BTreeMap::from([("contact_id".to_string(), json!(123))])
401 }
402
403 #[test]
404 fn policy_requires_a_non_empty_scope() {
405 let err = CleanupPolicy::new(BTreeMap::new(), vec!["id".into()], 10)
407 .expect_err("empty scope must be refused");
408 assert!(err.to_string().contains("at least one"), "{err}");
409 }
410
411 #[test]
412 fn policy_requires_a_key() {
413 let err = CleanupPolicy::new(scope(), vec![], 10).expect_err("no key must be refused");
414 assert!(err.to_string().contains("`key`"), "{err}");
415 }
416
417 #[test]
418 fn policy_refuses_a_null_scope_value() {
419 let s = BTreeMap::from([("contact_id".to_string(), Value::Null)]);
422 let err = CleanupPolicy::new(s, vec!["id".into()], 10).expect_err("null must be refused");
423 assert!(err.to_string().contains("null"), "{err}");
424 }
425
426 #[test]
427 fn policy_floors_max_keys_at_one() {
428 let p = CleanupPolicy::new(scope(), vec!["id".into()], 0).unwrap();
429 assert_eq!(p.max_keys, 1);
430 }
431
432 #[tokio::test]
433 async fn cleanup_tracker_forwards_overwrite_lifecycle() {
434 struct OvwSink {
437 log: std::sync::Mutex<Vec<&'static str>>,
438 }
439 #[async_trait::async_trait]
440 impl Sink for OvwSink {
441 async fn write_batch(&self, r: &[Value]) -> Result<usize, FaucetError> {
442 Ok(r.len())
443 }
444 fn is_overwrite(&self) -> bool {
445 true
446 }
447 async fn begin_overwrite(&self) -> Result<(), FaucetError> {
448 self.log.lock().unwrap().push("begin");
449 Ok(())
450 }
451 async fn commit_overwrite(&self) -> Result<(), FaucetError> {
452 self.log.lock().unwrap().push("commit");
453 Ok(())
454 }
455 async fn abort_overwrite(&self) -> Result<(), FaucetError> {
456 self.log.lock().unwrap().push("abort");
457 Ok(())
458 }
459 }
460 let inner = OvwSink {
461 log: std::sync::Mutex::new(Vec::new()),
462 };
463 let policy = CleanupPolicy::new(scope(), vec!["id".into()], 10).unwrap();
464 let tracker = CleanupTracker::new(&inner, &policy);
465 assert!(tracker.is_overwrite());
466 assert_eq!(
467 tracker.batch_atomicity(),
468 crate::dlq::BatchAtomicity::BestEffort
469 );
470 assert_eq!(tracker.write_batch(&[json!({"id": 1})]).await.unwrap(), 1);
472 tracker.begin_overwrite().await.unwrap();
473 tracker.commit_overwrite().await.unwrap();
474 tracker.abort_overwrite().await.unwrap();
475 tracker.complete_run().await.unwrap();
476 assert_eq!(*inner.log.lock().unwrap(), vec!["begin", "commit", "abort"]);
477 }
478
479 #[test]
480 fn accumulates_keys_across_pages() {
481 let mut seen = SeenKeys::new();
482 let key = vec!["id".to_string()];
483 seen.record_page(&[json!({"id": 1}), json!({"id": 2})], &key, 100);
484 seen.record_page(&[json!({"id": 3})], &key, 100);
485 assert_eq!(seen.len(), 3);
486 assert!(!seen.overflowed());
487 }
488
489 #[test]
490 fn accumulates_composite_keys_in_declared_order() {
491 let mut seen = SeenKeys::new();
492 let key = vec!["a".to_string(), "b".to_string()];
493 seen.record_page(&[json!({"b": 2, "a": 1})], &key, 100);
494 assert_eq!(seen.len(), 1);
495 let t = &seen.keys()[0].0;
496 assert_eq!(
497 t[0].0, "a",
498 "key order follows the declared `key`, not the record"
499 );
500 assert_eq!(t[1].0, "b");
501 }
502
503 #[test]
504 fn skips_rows_with_a_missing_or_null_key() {
505 let mut seen = SeenKeys::new();
506 let key = vec!["id".to_string()];
507 seen.record_page(
508 &[
509 json!({"id": 1}),
510 json!({"other": 9}), json!({"id": null}), json!("not an object"),
513 ],
514 &key,
515 100,
516 );
517 assert_eq!(seen.len(), 1, "only the well-keyed row is tracked");
518 }
519
520 #[test]
521 fn overflow_is_sticky_and_frees_the_buffer() {
522 let mut seen = SeenKeys::new();
523 let key = vec!["id".to_string()];
524 let page: Vec<Value> = (0..5).map(|i| json!({"id": i})).collect();
525 seen.record_page(&page, &key, 3);
526 assert!(seen.overflowed(), "ceiling of 3 must trip on a 5-row page");
527 assert!(seen.is_empty(), "buffer is freed — the cleanup will refuse");
528 seen.record_page(&[json!({"id": 99})], &key, 3);
530 assert!(seen.overflowed());
531 assert!(seen.is_empty());
532 }
533
534 #[test]
535 fn overflow_error_explains_that_nothing_was_deleted() {
536 let seen = SeenKeys::new();
537 let msg = seen.overflow_error(50).to_string();
538 assert!(msg.contains("Nothing was deleted"), "{msg}");
539 assert!(msg.contains("50"), "{msg}");
540 }
541}