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 sink_guarantee(&self) -> crate::idempotency::SinkGuarantee {
348 self.inner.sink_guarantee()
349 }
350 async fn current_schema(&self) -> Result<Option<Value>, FaucetError> {
351 self.inner.current_schema().await
352 }
353 fn supports_schema_evolution(&self) -> bool {
354 self.inner.supports_schema_evolution()
355 }
356 async fn evolve_schema(
357 &self,
358 evolution: &crate::drift::SchemaEvolution,
359 ) -> Result<(), FaucetError> {
360 self.inner.evolve_schema(evolution).await
361 }
362 fn config_schema(&self) -> Value {
363 self.inner.config_schema()
364 }
365 fn connector_name(&self) -> &'static str {
366 self.inner.connector_name()
367 }
368 fn dataset_uri(&self) -> String {
369 self.inner.dataset_uri()
370 }
371 fn is_overwrite(&self) -> bool {
372 self.inner.is_overwrite()
373 }
374 async fn begin_overwrite(&self) -> Result<(), FaucetError> {
375 self.inner.begin_overwrite().await
376 }
377 async fn commit_overwrite(&self) -> Result<(), FaucetError> {
378 self.inner.commit_overwrite().await
379 }
380 async fn abort_overwrite(&self) -> Result<(), FaucetError> {
381 self.inner.abort_overwrite().await
382 }
383}
384
385#[cfg(test)]
386mod tests {
387 use super::*;
388 use serde_json::json;
389
390 fn scope() -> BTreeMap<String, Value> {
391 BTreeMap::from([("contact_id".to_string(), json!(123))])
392 }
393
394 #[test]
395 fn policy_requires_a_non_empty_scope() {
396 let err = CleanupPolicy::new(BTreeMap::new(), vec!["id".into()], 10)
398 .expect_err("empty scope must be refused");
399 assert!(err.to_string().contains("at least one"), "{err}");
400 }
401
402 #[test]
403 fn policy_requires_a_key() {
404 let err = CleanupPolicy::new(scope(), vec![], 10).expect_err("no key must be refused");
405 assert!(err.to_string().contains("`key`"), "{err}");
406 }
407
408 #[test]
409 fn policy_refuses_a_null_scope_value() {
410 let s = BTreeMap::from([("contact_id".to_string(), Value::Null)]);
413 let err = CleanupPolicy::new(s, vec!["id".into()], 10).expect_err("null must be refused");
414 assert!(err.to_string().contains("null"), "{err}");
415 }
416
417 #[test]
418 fn policy_floors_max_keys_at_one() {
419 let p = CleanupPolicy::new(scope(), vec!["id".into()], 0).unwrap();
420 assert_eq!(p.max_keys, 1);
421 }
422
423 #[tokio::test]
424 async fn cleanup_tracker_forwards_overwrite_lifecycle() {
425 struct OvwSink {
428 log: std::sync::Mutex<Vec<&'static str>>,
429 }
430 #[async_trait::async_trait]
431 impl Sink for OvwSink {
432 async fn write_batch(&self, r: &[Value]) -> Result<usize, FaucetError> {
433 Ok(r.len())
434 }
435 fn is_overwrite(&self) -> bool {
436 true
437 }
438 async fn begin_overwrite(&self) -> Result<(), FaucetError> {
439 self.log.lock().unwrap().push("begin");
440 Ok(())
441 }
442 async fn commit_overwrite(&self) -> Result<(), FaucetError> {
443 self.log.lock().unwrap().push("commit");
444 Ok(())
445 }
446 async fn abort_overwrite(&self) -> Result<(), FaucetError> {
447 self.log.lock().unwrap().push("abort");
448 Ok(())
449 }
450 }
451 let inner = OvwSink {
452 log: std::sync::Mutex::new(Vec::new()),
453 };
454 let policy = CleanupPolicy::new(scope(), vec!["id".into()], 10).unwrap();
455 let tracker = CleanupTracker::new(&inner, &policy);
456 assert!(tracker.is_overwrite());
457 assert_eq!(tracker.write_batch(&[json!({"id": 1})]).await.unwrap(), 1);
459 tracker.begin_overwrite().await.unwrap();
460 tracker.commit_overwrite().await.unwrap();
461 tracker.abort_overwrite().await.unwrap();
462 assert_eq!(*inner.log.lock().unwrap(), vec!["begin", "commit", "abort"]);
463 }
464
465 #[test]
466 fn accumulates_keys_across_pages() {
467 let mut seen = SeenKeys::new();
468 let key = vec!["id".to_string()];
469 seen.record_page(&[json!({"id": 1}), json!({"id": 2})], &key, 100);
470 seen.record_page(&[json!({"id": 3})], &key, 100);
471 assert_eq!(seen.len(), 3);
472 assert!(!seen.overflowed());
473 }
474
475 #[test]
476 fn accumulates_composite_keys_in_declared_order() {
477 let mut seen = SeenKeys::new();
478 let key = vec!["a".to_string(), "b".to_string()];
479 seen.record_page(&[json!({"b": 2, "a": 1})], &key, 100);
480 assert_eq!(seen.len(), 1);
481 let t = &seen.keys()[0].0;
482 assert_eq!(
483 t[0].0, "a",
484 "key order follows the declared `key`, not the record"
485 );
486 assert_eq!(t[1].0, "b");
487 }
488
489 #[test]
490 fn skips_rows_with_a_missing_or_null_key() {
491 let mut seen = SeenKeys::new();
492 let key = vec!["id".to_string()];
493 seen.record_page(
494 &[
495 json!({"id": 1}),
496 json!({"other": 9}), json!({"id": null}), json!("not an object"),
499 ],
500 &key,
501 100,
502 );
503 assert_eq!(seen.len(), 1, "only the well-keyed row is tracked");
504 }
505
506 #[test]
507 fn overflow_is_sticky_and_frees_the_buffer() {
508 let mut seen = SeenKeys::new();
509 let key = vec!["id".to_string()];
510 let page: Vec<Value> = (0..5).map(|i| json!({"id": i})).collect();
511 seen.record_page(&page, &key, 3);
512 assert!(seen.overflowed(), "ceiling of 3 must trip on a 5-row page");
513 assert!(seen.is_empty(), "buffer is freed — the cleanup will refuse");
514 seen.record_page(&[json!({"id": 99})], &key, 3);
516 assert!(seen.overflowed());
517 assert!(seen.is_empty());
518 }
519
520 #[test]
521 fn overflow_error_explains_that_nothing_was_deleted() {
522 let seen = SeenKeys::new();
523 let msg = seen.overflow_error(50).to_string();
524 assert!(msg.contains("Nothing was deleted"), "{msg}");
525 assert!(msg.contains("50"), "{msg}");
526 }
527}