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}