Skip to main content

faucet_core/
local_outputs.rs

1//! Local sink output tracking — the provenance record a retention GC deletes
2//! from (#587).
3//!
4//! A local run writes real files: `out.jsonl`, `rows.csv`, a directory of
5//! UUID-named parquet parts. Repeated local iteration (or a long-running
6//! `faucet serve` used for local testing) piles them up, so faucet grows a
7//! retention garbage-collector for them. That GC has one hard requirement:
8//!
9//! > **It may only delete files faucet itself created.** Never a glob, never a
10//! > directory wipe — not even for "clean all".
11//!
12//! Which means the deletion has to work from a *recorded list of concrete
13//! paths*, and the only component that knows those paths is the sink that
14//! opened them. A path cannot be re-derived after the fact:
15//!
16//! - The catalog stores a **canonical** dataset URI with `${now.*}` segments
17//!   folded back to their tokens, so a dated path there is deliberately not the
18//!   file that exists on disk.
19//! - A rolling parquet sink names each part with a fresh UUID. Nothing outside
20//!   the sink can enumerate them without globbing the directory — which is
21//!   exactly what the guardrail forbids.
22//!
23//! So each local-file sink accumulates a [`LocalOutputLog`] as it opens files
24//! and reports it through [`Sink::local_outputs`](crate::Sink::local_outputs);
25//! the CLI records that list after a successful run and the GC deletes only
26//! from it.
27//!
28//! ## Why `pre_existing` is tracked, and why the GC must honour it
29//!
30//! "faucet wrote this path" is not the same claim as "faucet created this file".
31//! A sink pointed at an existing file — `append: true` onto a colleague's
32//! export, or a mistyped path landing on a real file — writes to a file it did
33//! not create. Deleting that on a retention sweep is data loss of somebody
34//! else's data, and no retention window makes it acceptable.
35//!
36//! So the flag is captured **at the first open**, before the file is created
37//! ([`LocalOutputLog::record_open`] takes it from a caller-supplied probe), and
38//! it is sticky: a second run that truncates a path faucet created earlier is
39//! still faucet's own output, and re-recording must not reclassify it. The GC
40//! refuses to delete a `pre_existing` entry — including on an explicit,
41//! single-path "delete now" — and says why instead.
42
43use std::collections::BTreeMap;
44use std::path::{Path, PathBuf};
45use std::sync::Mutex;
46
47/// One concrete local file a sink opened during a run.
48#[derive(Debug, Clone, PartialEq, Eq)]
49pub struct LocalOutput {
50    /// The file the sink wrote, as the sink addressed it.
51    pub path: PathBuf,
52    /// The file already existed the first time this sink opened it — faucet
53    /// appended to (or truncated) a file it did not create. A retention GC must
54    /// never delete such a file; see the module docs.
55    pub pre_existing: bool,
56    /// The sink truncated the file when it first opened it, so every byte in
57    /// it is faucet's own output even when [`pre_existing`](Self::pre_existing)
58    /// is set. It makes a pre-existing file safe to *read back* (a preview can
59    /// only show what faucet wrote); it never makes it collectable.
60    pub replaced: bool,
61}
62
63impl LocalOutput {
64    /// A file faucet created.
65    pub fn created(path: impl Into<PathBuf>) -> Self {
66        Self {
67            path: path.into(),
68            pre_existing: false,
69            replaced: false,
70        }
71    }
72
73    /// A file that already existed when faucet first opened it.
74    pub fn pre_existing(path: impl Into<PathBuf>) -> Self {
75        Self {
76            path: path.into(),
77            pre_existing: true,
78            replaced: false,
79        }
80    }
81
82    /// A file that already existed and that faucet truncated on its first open.
83    pub fn replaced(path: impl Into<PathBuf>) -> Self {
84        Self {
85            path: path.into(),
86            pre_existing: true,
87            replaced: true,
88        }
89    }
90}
91
92/// Whether `path` already exists on disk — the probe a sink runs **before**
93/// opening a file for the first time, to fill
94/// [`LocalOutput::pre_existing`].
95///
96/// A path that cannot be stat-ed (a permission error on the parent, say) is
97/// reported as `true`: the conservative answer, since it makes the GC leave the
98/// file alone rather than delete something faucet may not have created.
99pub fn probe_pre_existing(path: &Path) -> bool {
100    match std::fs::symlink_metadata(path) {
101        Ok(_) => true,
102        Err(e) if e.kind() == std::io::ErrorKind::NotFound => false,
103        // Anything else (EACCES on the parent directory, EIO, a broken mount)
104        // is not a "definitely absent" answer, so do not claim faucet created it.
105        Err(_) => true,
106    }
107}
108
109/// A sink's accumulating list of the local files it has opened.
110///
111/// Shared shape for every local-file sink (and for third-party ones): dedup by
112/// path, `pre_existing` fixed by the **first** open of each path, insertion
113/// order preserved for stable reporting. Cheap enough to call on every open —
114/// a file open already costs a syscall.
115///
116/// Poisoned-lock safety: the accumulator is provenance metadata, never on the
117/// data path, so a poisoned mutex degrades to "record nothing" rather than
118/// panicking a sink mid-write. A lost record means the GC does not know about
119/// the file and leaves it on disk — the safe direction.
120#[derive(Debug, Default)]
121pub struct LocalOutputLog {
122    /// Path → (`pre_existing`, `replaced`, insertion index). `BTreeMap` for the
123    /// dedup; the index restores first-seen order on read.
124    seen: Mutex<BTreeMap<PathBuf, (bool, bool, usize)>>,
125}
126
127impl LocalOutputLog {
128    pub fn new() -> Self {
129        Self::default()
130    }
131
132    /// Record that the sink opened `path`, taking `pre_existing` from a probe
133    /// run before the open (see [`probe_pre_existing`]).
134    ///
135    /// Idempotent per path: a re-open (the per-page `flush()` → reopen cycle the
136    /// file sinks perform, or a later run truncating the same path) keeps the
137    /// classification captured the first time.
138    pub fn record_open(&self, path: impl Into<PathBuf>, pre_existing: bool) {
139        self.record_open_with(path, pre_existing, false);
140    }
141
142    /// [`record_open`](Self::record_open), also saying whether this first open
143    /// truncates the file (see [`LocalOutput::replaced`]). First open wins for
144    /// both flags.
145    pub fn record_open_with(&self, path: impl Into<PathBuf>, pre_existing: bool, truncates: bool) {
146        let path = path.into();
147        if let Ok(mut seen) = self.seen.lock() {
148            let next = seen.len();
149            seen.entry(path).or_insert((pre_existing, truncates, next));
150        }
151    }
152
153    /// Record a first open of `path`, probing the filesystem for
154    /// `pre_existing` **only if** the path has not been recorded yet.
155    ///
156    /// This is the entry point for sinks whose open path is async or runs inside
157    /// `spawn_blocking`: it keeps the stat off the hot path once the file is
158    /// known, and it cannot reclassify an already-recorded path.
159    pub fn record_open_probing(&self, path: impl Into<PathBuf>) {
160        self.record_open_probing_with(path, false);
161    }
162
163    /// [`record_open_probing`](Self::record_open_probing) for a first open that
164    /// truncates the file when `truncates` is set (see [`LocalOutput::replaced`]).
165    pub fn record_open_probing_with(&self, path: impl Into<PathBuf>, truncates: bool) {
166        let path = path.into();
167        let known = self
168            .seen
169            .lock()
170            .map(|seen| seen.contains_key(&path))
171            .unwrap_or(true);
172        if !known {
173            let pre_existing = probe_pre_existing(&path);
174            self.record_open_with(path, pre_existing, truncates);
175        }
176    }
177
178    /// The files recorded so far, in first-seen order.
179    pub fn snapshot(&self) -> Vec<LocalOutput> {
180        let Ok(seen) = self.seen.lock() else {
181            return Vec::new();
182        };
183        let mut rows: Vec<(usize, LocalOutput)> = seen
184            .iter()
185            .map(|(path, (pre_existing, replaced, idx))| {
186                (
187                    *idx,
188                    LocalOutput {
189                        path: path.clone(),
190                        pre_existing: *pre_existing,
191                        replaced: *pre_existing && *replaced,
192                    },
193                )
194            })
195            .collect();
196        rows.sort_by_key(|(idx, _)| *idx);
197        rows.into_iter().map(|(_, out)| out).collect()
198    }
199
200    /// Whether anything has been recorded.
201    pub fn is_empty(&self) -> bool {
202        self.seen.lock().map(|s| s.is_empty()).unwrap_or(true)
203    }
204}
205
206#[cfg(test)]
207mod tests {
208    use super::*;
209
210    #[test]
211    fn records_in_first_seen_order() {
212        let log = LocalOutputLog::new();
213        log.record_open("/tmp/b.jsonl", false);
214        log.record_open("/tmp/a.jsonl", false);
215        let snap = log.snapshot();
216        assert_eq!(
217            snap.iter().map(|o| o.path.clone()).collect::<Vec<_>>(),
218            vec![PathBuf::from("/tmp/b.jsonl"), PathBuf::from("/tmp/a.jsonl")],
219            "insertion order, not the BTreeMap's sort order"
220        );
221    }
222
223    #[test]
224    fn dedupes_by_path() {
225        let log = LocalOutputLog::new();
226        log.record_open("/tmp/a.jsonl", false);
227        log.record_open("/tmp/a.jsonl", false);
228        assert_eq!(log.snapshot().len(), 1);
229    }
230
231    #[test]
232    fn first_open_classification_is_sticky() {
233        // The flush→reopen cycle re-opens an existing file; that second open must
234        // not flip a faucet-created file into "pre-existing" (which would make it
235        // permanently un-collectable).
236        let log = LocalOutputLog::new();
237        log.record_open("/tmp/a.jsonl", false);
238        log.record_open("/tmp/a.jsonl", true);
239        assert!(!log.snapshot()[0].pre_existing);
240
241        // And the reverse: a file faucet did not create stays that way.
242        let log = LocalOutputLog::new();
243        log.record_open("/tmp/theirs.csv", true);
244        log.record_open("/tmp/theirs.csv", false);
245        assert!(log.snapshot()[0].pre_existing);
246    }
247
248    #[test]
249    fn replaced_is_first_open_wins_and_only_meaningful_for_pre_existing() {
250        let log = LocalOutputLog::new();
251        log.record_open_with("/tmp/theirs.jsonl", true, true);
252        log.record_open_with("/tmp/theirs.jsonl", true, false);
253        log.record_open_with("/tmp/ours.jsonl", false, true);
254        log.record_open("/tmp/appended.jsonl", true);
255        let snap = log.snapshot();
256        assert_eq!(snap[0], LocalOutput::replaced("/tmp/theirs.jsonl"));
257        assert_eq!(snap[1], LocalOutput::created("/tmp/ours.jsonl"));
258        assert_eq!(snap[2], LocalOutput::pre_existing("/tmp/appended.jsonl"));
259    }
260
261    #[test]
262    fn probing_with_truncation_flags_an_existing_file_as_replaced() {
263        let dir = std::env::temp_dir().join(format!("faucet-lo-{}", std::process::id()));
264        std::fs::create_dir_all(&dir).unwrap();
265        let existing = dir.join("existing.jsonl");
266        std::fs::write(&existing, b"x").unwrap();
267        let log = LocalOutputLog::new();
268        log.record_open_probing_with(&existing, true);
269        log.record_open_probing_with(dir.join("new.jsonl"), true);
270        let snap = log.snapshot();
271        assert!(snap[0].pre_existing && snap[0].replaced);
272        assert!(!snap[1].pre_existing && !snap[1].replaced);
273        std::fs::remove_dir_all(&dir).unwrap();
274    }
275
276    #[test]
277    fn empty_log_reports_empty() {
278        let log = LocalOutputLog::new();
279        assert!(log.is_empty());
280        assert!(log.snapshot().is_empty());
281        log.record_open("/tmp/a.jsonl", false);
282        assert!(!log.is_empty());
283    }
284
285    #[test]
286    fn probe_reports_absent_and_present_paths() {
287        let dir = tempfile::tempdir().unwrap();
288        let missing = dir.path().join("nope.jsonl");
289        assert!(!probe_pre_existing(&missing));
290        std::fs::write(&missing, b"").unwrap();
291        assert!(probe_pre_existing(&missing));
292    }
293
294    #[test]
295    fn probe_reports_a_dangling_symlink_as_pre_existing() {
296        // `symlink_metadata` (not `metadata`) so a link whose target is gone is
297        // still "something is already at this path" — creating through it would
298        // write the target, and deleting it later is not faucet's call.
299        let dir = tempfile::tempdir().unwrap();
300        let link = dir.path().join("link.jsonl");
301        #[cfg(unix)]
302        std::os::unix::fs::symlink(dir.path().join("absent-target"), &link).unwrap();
303        #[cfg(not(unix))]
304        std::fs::write(&link, b"").unwrap();
305        assert!(probe_pre_existing(&link));
306    }
307
308    #[test]
309    fn record_open_probing_probes_once_then_reuses() {
310        let dir = tempfile::tempdir().unwrap();
311        let path = dir.path().join("out.jsonl");
312        let log = LocalOutputLog::new();
313        // Absent at first open → faucet's own file.
314        log.record_open_probing(&path);
315        assert!(!log.snapshot()[0].pre_existing);
316        // Now it exists (the sink created it), but the second open must not
317        // re-probe and reclassify.
318        std::fs::write(&path, b"{}\n").unwrap();
319        log.record_open_probing(&path);
320        assert_eq!(log.snapshot().len(), 1);
321        assert!(!log.snapshot()[0].pre_existing);
322    }
323
324    #[test]
325    fn record_open_probing_marks_a_file_faucet_did_not_create() {
326        let dir = tempfile::tempdir().unwrap();
327        let path = dir.path().join("theirs.jsonl");
328        std::fs::write(&path, b"existing\n").unwrap();
329        let log = LocalOutputLog::new();
330        log.record_open_probing(&path);
331        assert!(log.snapshot()[0].pre_existing);
332    }
333
334    #[test]
335    fn constructors_set_the_flag() {
336        assert!(!LocalOutput::created("/tmp/a").pre_existing);
337        assert!(LocalOutput::pre_existing("/tmp/a").pre_existing);
338    }
339}