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}