Skip to main content

rustpython_host_env/fs/
append_log.rs

1use parking_lot::Mutex;
2use std::{
3    fs::{File, OpenOptions},
4    io::{self, Write},
5    path::Path,
6};
7
8/// A replaceable append-only log destination, shared by its owner's consumers.
9/// File replacement and complete record writes are serialized per destination.
10#[derive(Debug, Default)]
11pub struct AppendLog {
12    file: Mutex<Option<File>>,
13}
14
15impl AppendLog {
16    /// Close the old destination before opening its replacement. On failure the
17    /// log stays disabled. The header is written only to empty regular files.
18    pub fn set_path(&self, path: Option<&Path>, header: &[u8]) -> io::Result<()> {
19        let mut slot = self.file.lock();
20        *slot = None;
21        let Some(path) = path else {
22            return Ok(());
23        };
24        let mut file = OpenOptions::new().create(true).append(true).open(path)?;
25        let metadata = file.metadata()?;
26        if metadata.is_file() && metadata.len() == 0 {
27            file.write_all(header)?;
28            file.flush()?;
29        }
30        *slot = Some(file);
31        Ok(())
32    }
33
34    pub fn enabled(&self) -> bool {
35        self.file.lock().is_some()
36    }
37
38    /// Append and flush a complete record, or do nothing if disabled.
39    pub fn write(&self, record: &[u8]) -> io::Result<()> {
40        if let Some(file) = self.file.lock().as_mut() {
41            file.write_all(record)?;
42            file.flush()?;
43        }
44        Ok(())
45    }
46}
47
48#[cfg(test)]
49mod tests {
50    use super::*;
51    use core::sync::atomic::{AtomicUsize, Ordering};
52    use std::path::PathBuf;
53
54    struct TestDir(PathBuf);
55
56    impl TestDir {
57        fn new() -> Self {
58            static NEXT: AtomicUsize = AtomicUsize::new(0);
59            loop {
60                let path = std::env::temp_dir().join(format!(
61                    "rustpython-append-log-{}-{}",
62                    std::process::id(),
63                    NEXT.fetch_add(1, Ordering::Relaxed)
64                ));
65                match std::fs::create_dir(&path) {
66                    Ok(()) => return Self(path),
67                    Err(e) if e.kind() == io::ErrorKind::AlreadyExists => continue,
68                    Err(e) => panic!("create test directory: {e}"),
69                }
70            }
71        }
72    }
73
74    impl Drop for TestDir {
75        fn drop(&mut self) {
76            std::fs::remove_dir_all(&self.0).unwrap();
77        }
78    }
79
80    #[test]
81    fn append_replace_and_disable() {
82        let dir = TestDir::new();
83        let first = dir.0.join("first");
84        let second = dir.0.join("second");
85        let log = AppendLog::default();
86        assert!(!log.enabled());
87        log.write(b"ignored\n").unwrap();
88        log.set_path(Some(&first), b"header\n").unwrap();
89        log.write(b"one\n").unwrap();
90        log.set_path(Some(&first), b"header\n").unwrap();
91        log.write(b"two\n").unwrap();
92        log.set_path(Some(&second), b"header\n").unwrap();
93        log.write(b"three\n").unwrap();
94        log.set_path(None, b"header\n").unwrap();
95        log.write(b"ignored\n").unwrap();
96        assert!(!log.enabled());
97        assert_eq!(std::fs::read(first).unwrap(), b"header\none\ntwo\n");
98        assert_eq!(std::fs::read(second).unwrap(), b"header\nthree\n");
99    }
100
101    #[test]
102    fn failed_replacement_disables_the_old_destination() {
103        let dir = TestDir::new();
104        let path = dir.0.join("log");
105        let log = AppendLog::default();
106        log.set_path(Some(&path), b"header\n").unwrap();
107        assert!(
108            log.set_path(Some(&dir.0.join("missing/log")), b"header\n")
109                .is_err()
110        );
111        assert!(!log.enabled());
112        log.write(b"ignored\n").unwrap();
113        assert_eq!(std::fs::read(path).unwrap(), b"header\n");
114    }
115
116    #[test]
117    fn concurrent_records_are_complete() {
118        let dir = TestDir::new();
119        let path = dir.0.join("log");
120        let log = AppendLog::default();
121        log.set_path(Some(&path), b"").unwrap();
122        std::thread::scope(|scope| {
123            for n in 0..4 {
124                let log = &log;
125                scope.spawn(move || {
126                    for _ in 0..100 {
127                        log.write(format!("{n}\n").as_bytes()).unwrap();
128                    }
129                });
130            }
131        });
132        let contents = std::fs::read_to_string(path).unwrap();
133        assert_eq!(contents.lines().count(), 400);
134        for n in 0..4 {
135            assert_eq!(
136                contents
137                    .lines()
138                    .filter(|line| *line == n.to_string())
139                    .count(),
140                100
141            );
142        }
143    }
144}