Skip to main content

pitchfork_cli/cli/
log_sink.rs

1use crate::Result;
2use crate::daemon_id::DaemonId;
3use crate::log_parse::ParsedLog;
4use crate::log_store::LogStore;
5use crate::log_store::sqlite::LOG_STORE;
6use tokio::io::AsyncReadExt;
7
8/// Number of parsed lines to accumulate before writing them as one batch.
9const BATCH_SIZE: usize = 100;
10
11/// Longest a parsed line waits in the batch before being written.
12const FLUSH_INTERVAL: std::time::Duration = std::time::Duration::from_millis(100);
13
14/// How many parsed lines may be queued for writing before reading slows down.
15///
16/// Reading and writing run separately so a slow write cannot stop the pipe being
17/// drained, but the queue is bounded: output that sustainably outpaces the store
18/// has to push back on the daemon eventually, which is preferable to growing
19/// without limit.
20const QUEUE_DEPTH: usize = 8192;
21
22/// Bytes read from the pipe at a time.
23const READ_CHUNK: usize = 8192;
24
25/// Longest run of bytes treated as a single line.
26///
27/// Output containing no newline must not accumulate indefinitely: a daemon
28/// emitting an endless stream, or binary data, would otherwise grow the buffer
29/// until the sink was killed for using too much memory — whereupon the
30/// supervisor would start another sink and repeat it.
31const MAX_LINE_BYTES: usize = 64 * 1024;
32
33/// Reads a daemon's output on stdin and writes it to the log store
34///
35/// Spawned by the supervisor as a sibling of the daemon, holding the read end
36/// of the daemon's output pipe. Keeping the reader in its own process is what
37/// makes logging survive a supervisor crash: the pipe still has a reader, so
38/// the daemon is neither killed by SIGPIPE nor blocked, and no output is lost.
39/// Exits when the pipe reaches end of file, which happens once the daemon and
40/// every descendant holding the write end have gone.
41#[derive(Debug, clap::Args)]
42#[clap(hide = true, verbatim_doc_comment)]
43pub struct LogSink {
44    /// Qualified id of the daemon whose output this is
45    #[clap(long)]
46    daemon_id: String,
47
48    /// Log format to parse lines with (`json`, `logfmt`, `auto`, or `text`)
49    #[clap(long, default_value = "text")]
50    log_format: String,
51}
52
53impl LogSink {
54    pub async fn run(&self) -> Result<()> {
55        let id = DaemonId::parse(&self.daemon_id)?;
56
57        // Reading and writing are separate tasks. A write to SQLite can block —
58        // for as long as the store's busy timeout, if another writer holds the
59        // lock — and this process is the only reader of the daemon's pipe, so a
60        // write must never stop it being drained.
61        let (tx, rx) = tokio::sync::mpsc::channel::<ParsedLog>(QUEUE_DEPTH);
62        let writer = tokio::spawn(write_batches(id.clone(), rx));
63
64        let read_result = read_lines(tx, &self.log_format).await;
65
66        // The sender has been dropped by now, so the writer drains its queue and
67        // returns; wait for it so nothing queued is lost on exit.
68        let _ = writer.await;
69
70        read_result.map_err(|e| {
71            miette::miette!("log sink for {id} could not read the daemon's output: {e}")
72        })
73    }
74}
75
76/// Split the daemon's output into lines and queue them for writing.
77///
78/// Returns once the pipe reaches end of file. A read error is propagated so the
79/// process can exit non-zero: exiting cleanly would tell the supervisor the
80/// stream had finished and it would stop replacing this sink.
81async fn read_lines(
82    tx: tokio::sync::mpsc::Sender<ParsedLog>,
83    log_format: &str,
84) -> std::io::Result<()> {
85    let mut stdin = tokio::io::stdin();
86    let mut chunk = vec![0u8; READ_CHUNK];
87    let mut line: Vec<u8> = Vec::with_capacity(256);
88    // Whether the last line was emitted because it reached the cap rather than
89    // because it ended. A newline arriving straight afterwards terminates the
90    // line already written, so it must not produce an empty one.
91    let mut split_at_cap = false;
92
93    loop {
94        let read = stdin.read(&mut chunk).await?;
95        if read == 0 {
96            break;
97        }
98        for &byte in &chunk[..read] {
99            if byte == b'\n' {
100                if split_at_cap && line.is_empty() {
101                    split_at_cap = false;
102                    continue;
103                }
104                split_at_cap = false;
105                queue(&tx, &mut line, log_format).await?;
106            } else {
107                line.push(byte);
108                split_at_cap = false;
109                // Emit an over-long run as its own line rather than letting the
110                // buffer grow without bound.
111                if line.len() >= MAX_LINE_BYTES {
112                    queue_capped(&tx, &mut line, log_format).await?;
113                    split_at_cap = true;
114                }
115            }
116        }
117    }
118
119    // Anything written without a trailing newline is still output.
120    if !line.is_empty() {
121        queue(&tx, &mut line, log_format).await?;
122    }
123    Ok(())
124}
125
126/// Emit a line that has reached the length cap, keeping any trailing bytes that
127/// form an incomplete character.
128///
129/// Splitting purely by byte count would cut a multi-byte character in half, and
130/// converting each half on its own turns one valid character into two
131/// replacement characters.
132async fn queue_capped(
133    tx: &tokio::sync::mpsc::Sender<ParsedLog>,
134    line: &mut Vec<u8>,
135    log_format: &str,
136) -> std::io::Result<()> {
137    let split = split_before_incomplete_char(line);
138    let tail = line.split_off(split);
139    let result = queue(tx, line, log_format).await;
140    *line = tail;
141    result
142}
143
144/// Length to cut `bytes` at so no character is left half-written.
145///
146/// Decided by inspecting the final bytes rather than by asking `from_utf8` where
147/// the string stops being valid: that reports the *first* problem, so a single
148/// invalid byte earlier in the line would hide an unfinished character at the
149/// end, and the character would be split after all.
150fn split_before_incomplete_char(bytes: &[u8]) -> usize {
151    let len = bytes.len();
152    // A character is at most four bytes, so only the last few can be unfinished.
153    for i in (len.saturating_sub(4)..len).rev() {
154        let byte = bytes[i];
155        if byte & 0b1100_0000 == 0b1000_0000 {
156            continue; // a continuation byte; keep looking back for its lead
157        }
158        let expected = match byte {
159            0x00..=0x7f => 1,
160            b if b >> 5 == 0b110 => 2,
161            b if b >> 4 == 0b1110 => 3,
162            b if b >> 3 == 0b11110 => 4,
163            // Not a valid lead byte at all, so nothing is pending; the lossy
164            // conversion will render it.
165            _ => 1,
166        };
167        return if i + expected > len && i > 0 { i } else { len };
168    }
169    len
170}
171
172/// Parse `line` and hand it to the writer, clearing it either way.
173///
174/// A closed queue means the writer task is gone, which is a failure rather than
175/// the end of the stream: reporting it as success would tell the supervisor this
176/// sink had reached end of file, and it would stop replacing it while the daemon
177/// was still writing.
178async fn queue(
179    tx: &tokio::sync::mpsc::Sender<ParsedLog>,
180    line: &mut Vec<u8>,
181    log_format: &str,
182) -> std::io::Result<()> {
183    // Convert lossily: a daemon emitting a stray non-UTF-8 byte must not be able
184    // to stop its own logging.
185    let text = String::from_utf8_lossy(line);
186    let parsed = crate::log_parse::parse(text.trim_end_matches('\r'), log_format);
187    line.clear();
188    tx.send(parsed)
189        .await
190        .map_err(|_| std::io::Error::other("log writer stopped"))
191}
192
193/// Write queued lines in batches until the queue closes.
194async fn write_batches(id: DaemonId, mut rx: tokio::sync::mpsc::Receiver<ParsedLog>) {
195    let mut batch: Vec<ParsedLog> = Vec::with_capacity(BATCH_SIZE);
196    let mut flush_interval = tokio::time::interval(FLUSH_INTERVAL);
197    flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
198
199    loop {
200        let closed = tokio::select! {
201            received = rx.recv_many(&mut batch, BATCH_SIZE) => received == 0,
202            _ = flush_interval.tick() => false,
203        };
204        flush(&id, &mut batch).await;
205        if closed {
206            break;
207        }
208    }
209}
210
211/// Write one batch, off the runtime so the SQLite call cannot stall other tasks.
212async fn flush(id: &DaemonId, batch: &mut Vec<ParsedLog>) {
213    if batch.is_empty() {
214        return;
215    }
216    let daemon_id = id.clone();
217    let entries = std::mem::take(batch);
218    let written = tokio::task::spawn_blocking(move || {
219        LOG_STORE.append_structured_batch(&daemon_id, &entries)
220    })
221    .await;
222    if let Ok(Err(e)) = written {
223        // Nothing useful to do but report it: the supervisor is not necessarily
224        // alive to be told, and dropping a batch is preferable to stalling the
225        // daemon behind a pipe nobody is draining.
226        error!("log sink failed to write batch for {id}: {e}");
227    }
228}