Skip to main content

dedup/transform/
stateful.rs

1
2
3use crate::cluster::DuplicateCluster;
4use crate::config::Config;
5use crate::error::Result;
6use tenshift_core::sample::Sample;
7use tracing::instrument;
8
9use super::DedupTransformer;
10
11/// A stateful deduplication transform that buffers samples.
12///
13/// This is the recommended transform for deduplication as it properly
14/// handles the stateful nature of duplicate detection.
15pub struct StatefulDedupTransform {
16    /// Inner transformer.
17    inner: DedupTransformer,
18    /// Whether we've finished processing.
19    finished: bool,
20}
21
22impl StatefulDedupTransform {
23    /// Create a new stateful deduplication transform.
24    ///
25    /// # Errors
26    ///
27    /// Returns an error if the configuration is invalid.
28    #[instrument(skip(config), level = "debug")]
29    pub fn new(config: Config) -> Result<Self> {
30        Ok(Self {
31            inner: DedupTransformer::new(config)?,
32            finished: false,
33        })
34    }
35
36    /// Set the text field name.
37    #[must_use]
38    pub fn with_text_field(mut self, field: impl Into<String>) -> Self {
39        self.inner = self.inner.with_text_field(field);
40        self
41    }
42
43    /// Enable marking duplicates.
44    #[must_use]
45    pub fn with_mark_duplicates(mut self, enabled: bool) -> Self {
46        self.inner = self.inner.with_mark_duplicates(enabled);
47        self
48    }
49
50    /// Get duplicate clusters.
51    #[must_use]
52    pub fn clusters(&mut self) -> &[DuplicateCluster] {
53        self.inner.clusters()
54    }
55
56    /// Get statistics.
57    #[must_use]
58    pub fn stats(&self) -> crate::lsh::LshStats {
59        self.inner.stats()
60    }
61}
62
63impl tenshift_core::transform::StatefulTransform for StatefulDedupTransform {
64    fn push(&mut self, sample: Sample) -> Vec<Sample> {
65        self.inner.push(sample);
66        Vec::new() // Buffer all samples, output on finish
67    }
68
69    fn finish(&mut self) -> Vec<Sample> {
70        if self.finished {
71            return Vec::new();
72        }
73        self.finished = true;
74        self.inner.finish_batch()
75    }
76
77    fn name(&self) -> &str {
78        "stateful_dedup"
79    }
80}
81