dedup/transform/
stateful.rs1
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
11pub struct StatefulDedupTransform {
16 inner: DedupTransformer,
18 finished: bool,
20}
21
22impl StatefulDedupTransform {
23 #[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 #[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 #[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 #[must_use]
52 pub fn clusters(&mut self) -> &[DuplicateCluster] {
53 self.inner.clusters()
54 }
55
56 #[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() }
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