pitchfork-cli 2.24.1

Daemons with DX
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
use crate::Result;
use crate::daemon_id::DaemonId;
use crate::log_parse::ParsedLog;
use crate::log_store::LogStore;
use crate::log_store::sqlite::LOG_STORE;
use tokio::io::AsyncReadExt;

/// Number of parsed lines to accumulate before writing them as one batch.
const BATCH_SIZE: usize = 100;

/// Longest a parsed line waits in the batch before being written.
const FLUSH_INTERVAL: std::time::Duration = std::time::Duration::from_millis(100);

/// How many parsed lines may be queued for writing before reading slows down.
///
/// Reading and writing run separately so a slow write cannot stop the pipe being
/// drained, but the queue is bounded: output that sustainably outpaces the store
/// has to push back on the daemon eventually, which is preferable to growing
/// without limit.
const QUEUE_DEPTH: usize = 8192;

/// Bytes read from the pipe at a time.
const READ_CHUNK: usize = 8192;

/// Longest run of bytes treated as a single line.
///
/// Output containing no newline must not accumulate indefinitely: a daemon
/// emitting an endless stream, or binary data, would otherwise grow the buffer
/// until the sink was killed for using too much memory — whereupon the
/// supervisor would start another sink and repeat it.
const MAX_LINE_BYTES: usize = 64 * 1024;

/// Reads a daemon's output on stdin and writes it to the log store
///
/// Spawned by the supervisor as a sibling of the daemon, holding the read end
/// of the daemon's output pipe. Keeping the reader in its own process is what
/// makes logging survive a supervisor crash: the pipe still has a reader, so
/// the daemon is neither killed by SIGPIPE nor blocked, and no output is lost.
/// Exits when the pipe reaches end of file, which happens once the daemon and
/// every descendant holding the write end have gone.
#[derive(Debug, usage_rs::Args)]
#[usage(verbatim_doc_comment)]
pub struct LogSink {
    /// Qualified id of the daemon whose output this is
    #[usage(long)]
    daemon_id: String,

    /// Log format to parse lines with (`json`, `logfmt`, `auto`, or `text`)
    #[usage(long, default = "text")]
    log_format: String,

    /// Regex whose first match means the daemon is ready
    ///
    /// Set for a daemon configured with `ready_output`. The supervisor cannot
    /// match it itself — this process holds the output — so the match is
    /// reported back over IPC.
    #[usage(long)]
    ready_pattern: Option<String>,

    /// Token identifying the start attempt this sink belongs to
    ///
    /// Quoted back when reporting a match. The supervisor drops reports whose
    /// token is no longer current, so a sink still draining a failed attempt
    /// cannot mark that daemon's retry ready.
    #[usage(long, default_value_t = 0, default = "0")]
    relay_token: u64,

    /// Report lines so the supervisor can fire the daemon's `on_output` hook
    ///
    /// Without `--output-filter` or `--output-regex` every line qualifies,
    /// which is what a hook with no pattern asks for.
    #[usage(long)]
    report_output: bool,

    /// Only report lines containing this substring
    #[usage(long)]
    output_filter: Option<String>,

    /// Only report lines matching this regex
    #[usage(long)]
    output_regex: Option<String>,

    /// Shortest gap between reported lines, in milliseconds
    #[usage(long, default_value_t = 1000, default = "1000")]
    output_debounce_ms: u64,
}

impl LogSink {
    pub async fn run(&self) -> Result<()> {
        let id = DaemonId::parse(&self.daemon_id)?;

        // A pattern that does not compile is reported and then ignored, rather
        // than failing the sink: refusing to start would leave the daemon's
        // output unread, which is far worse than a readiness check that never
        // fires. The supervisor validates patterns too, so this is a backstop.
        let compile = |what: &str, pattern: &str| {
            regex::Regex::new(pattern)
                .map_err(|e| error!("log sink for {id} ignoring unparsable {what}: {e}"))
                .ok()
        };
        let ready_pattern = self
            .ready_pattern
            .as_deref()
            .and_then(|p| compile("ready pattern", p));
        let hook = self.report_output.then(|| HookMatcher {
            filter: self.output_filter.clone(),
            regex: self
                .output_regex
                .as_deref()
                .and_then(|p| compile("output pattern", p)),
            debounce: std::time::Duration::from_millis(self.output_debounce_ms),
            last_reported: None,
        });

        // Reading and writing are separate tasks. A write to SQLite can block —
        // for as long as the store's busy timeout, if another writer holds the
        // lock — and this process is the only reader of the daemon's pipe, so a
        // write must never stop it being drained.
        let (tx, rx) = tokio::sync::mpsc::channel::<SinkEvent>(QUEUE_DEPTH);
        let writer = tokio::spawn(write_batches(id.clone(), self.relay_token, rx));

        let read_result =
            read_lines(tx, &self.log_format, ReadyMatcher::new(ready_pattern, hook)).await;

        // The sender has been dropped by now, so the writer drains its queue and
        // returns; wait for it so nothing queued is lost on exit.
        let _ = writer.await;

        read_result.map_err(|e| {
            miette::miette!("log sink for {id} could not read the daemon's output: {e}")
        })
    }
}

/// Something for the writer task to do, in the order the reader saw it.
///
/// Reporting a readiness match travels the same queue as the lines rather than
/// jumping ahead of them, so the line that triggered the match is always in the
/// log store by the time the supervisor hears about it — `collect_startup_logs`
/// and `pitchfork logs` would otherwise be able to miss it.
enum SinkEvent {
    Line(ParsedLog),
    Report(Report),
}

/// A line the supervisor needs to see, and why.
#[derive(Debug, PartialEq, Eq)]
struct Report {
    text: String,
    /// Whether it passed the `on_output` hook's filter and debounce. False for
    /// a line reported only because it matched the readiness pattern — firing a
    /// hook that filters for something else would be wrong.
    fires_hook: bool,
}

/// How much of a capped line is carried forward for matching.
///
/// A line longer than [`MAX_LINE_BYTES`] is emitted in pieces, and a readiness
/// pattern straddling a split would match none of them — the daemon would then
/// be killed at its readiness timeout despite having announced itself. Keeping
/// the tail of the previous piece closes that for any pattern shorter than this
/// while still bounding what is held.
const MATCH_CARRY_BYTES: usize = 4 * 1024;

/// Watches a daemon's output for the things the supervisor would look for if it
/// could still read the stream: the readiness pattern, and whatever fires the
/// `on_output` hook.
///
/// What the supervisor does about a reported line — mark the daemon ready, run
/// the hook — remains its own business.
struct ReadyMatcher {
    /// Readiness pattern, cleared once it has matched: readiness happens once.
    pattern: Option<regex::Regex>,
    /// The `on_output` hook's filter and rate limit, if the daemon has one.
    hook: Option<HookMatcher>,
    /// Tail of the previous piece of a line split at the length cap. Empty
    /// whenever the last piece ended at a real newline.
    carried: String,
}

/// The `on_output` hook's line filter and its rate limit.
struct HookMatcher {
    filter: Option<String>,
    regex: Option<regex::Regex>,
    debounce: std::time::Duration,
    last_reported: Option<std::time::Instant>,
}

impl HookMatcher {
    /// Whether this line should fire the hook, consuming the debounce window.
    ///
    /// The debounce is applied here rather than by the supervisor because this
    /// process is the one that sees every line: enforcing it here means one IPC
    /// message per window instead of one per line, which matters for a hook
    /// with no filter, where every line qualifies.
    fn matches(&mut self, clean: &str) -> bool {
        let matched = match (&self.filter, &self.regex) {
            // Mutually exclusive, and a hook setting both is rejected before it
            // ever reaches this process.
            (Some(substr), _) => clean.contains(substr.as_str()),
            (None, Some(re)) => re.is_match(clean),
            // A hook with neither fires on every line.
            (None, None) => true,
        };
        if !matched {
            return false;
        }
        let now = std::time::Instant::now();
        if self
            .last_reported
            .is_some_and(|last| now.duration_since(last) < self.debounce)
        {
            return false;
        }
        self.last_reported = Some(now);
        true
    }
}

impl ReadyMatcher {
    fn new(pattern: Option<regex::Regex>, hook: Option<HookMatcher>) -> Self {
        Self {
            pattern,
            hook,
            carried: String::new(),
        }
    }

    /// Whether anything is still being watched for. Once nothing is, matching is
    /// skipped: stripping ANSI from every line of a chatty daemon is not free.
    fn is_watching(&self) -> bool {
        self.pattern.is_some() || self.hook.is_some()
    }

    /// Test `text` — one whole line, or one piece of an over-long one — and
    /// return what matched, which is what the supervisor is told about.
    ///
    /// `split_at_cap` says another piece of the same logical line follows.
    fn consider(&mut self, text: &str, split_at_cap: bool) -> Option<Report> {
        if !self.is_watching() {
            self.carried.clear();
            return None;
        }
        let clean = console::strip_ansi_codes(text);
        let candidate = if self.carried.is_empty() {
            clean.into_owned()
        } else {
            format!("{}{clean}", self.carried)
        };

        self.carried = if split_at_cap {
            let start = candidate.len().saturating_sub(MATCH_CARRY_BYTES);
            // Never split a character in half; walk forward to a boundary.
            let start = (start..candidate.len())
                .find(|i| candidate.is_char_boundary(*i))
                .unwrap_or(candidate.len());
            candidate[start..].to_string()
        } else {
            String::new()
        };

        // A line can qualify for either reason, or both. The supervisor
        // re-checks the readiness pattern on what it is sent, but it cannot
        // re-check the hook without redoing the debounce, so that answer is
        // carried explicitly.
        let mut ready_matched = false;
        if self
            .pattern
            .as_ref()
            .is_some_and(|re| re.is_match(&candidate))
        {
            self.pattern = None;
            ready_matched = true;
        }
        let fires_hook = self
            .hook
            .as_mut()
            .is_some_and(|hook| hook.matches(&candidate));

        // Send the text the patterns actually matched against, not just this
        // piece of it: the supervisor re-matches it, and hands it to the hook.
        (ready_matched || fires_hook).then_some(Report {
            text: candidate,
            fires_hook,
        })
    }
}

/// Split the daemon's output into lines and queue them for writing.
///
/// Returns once the pipe reaches end of file. A read error is propagated so the
/// process can exit non-zero: exiting cleanly would tell the supervisor the
/// stream had finished and it would stop replacing this sink.
async fn read_lines(
    tx: tokio::sync::mpsc::Sender<SinkEvent>,
    log_format: &str,
    mut matcher: ReadyMatcher,
) -> std::io::Result<()> {
    let mut stdin = tokio::io::stdin();
    let mut chunk = vec![0u8; READ_CHUNK];
    let mut line: Vec<u8> = Vec::with_capacity(256);
    // Whether the last line was emitted because it reached the cap rather than
    // because it ended. A newline arriving straight afterwards terminates the
    // line already written, so it must not produce an empty one.
    let mut split_at_cap = false;

    loop {
        let read = stdin.read(&mut chunk).await?;
        if read == 0 {
            break;
        }
        for &byte in &chunk[..read] {
            if byte == b'\n' {
                if split_at_cap && line.is_empty() {
                    split_at_cap = false;
                    continue;
                }
                split_at_cap = false;
                queue(&tx, &mut line, log_format, &mut matcher).await?;
            } else {
                line.push(byte);
                split_at_cap = false;
                // Emit an over-long run as its own line rather than letting the
                // buffer grow without bound.
                if line.len() >= MAX_LINE_BYTES {
                    queue_capped(&tx, &mut line, log_format, &mut matcher).await?;
                    split_at_cap = true;
                }
            }
        }
    }

    // Anything written without a trailing newline is still output.
    if !line.is_empty() {
        queue(&tx, &mut line, log_format, &mut matcher).await?;
    }
    Ok(())
}

/// Emit a line that has reached the length cap, keeping any trailing bytes that
/// form an incomplete character.
///
/// Splitting purely by byte count would cut a multi-byte character in half, and
/// converting each half on its own turns one valid character into two
/// replacement characters.
async fn queue_capped(
    tx: &tokio::sync::mpsc::Sender<SinkEvent>,
    line: &mut Vec<u8>,
    log_format: &str,
    matcher: &mut ReadyMatcher,
) -> std::io::Result<()> {
    let split = split_before_incomplete_char(line);
    let tail = line.split_off(split);
    let result = queue_piece(tx, line, log_format, matcher).await;
    *line = tail;
    result
}

/// Queue one piece of a line that hit the length cap, telling the matcher that
/// the rest of the logical line is still to come.
async fn queue_piece(
    tx: &tokio::sync::mpsc::Sender<SinkEvent>,
    line: &mut Vec<u8>,
    log_format: &str,
    matcher: &mut ReadyMatcher,
) -> std::io::Result<()> {
    let text = String::from_utf8_lossy(line);
    let text = text.trim_end_matches('\r');
    let parsed = crate::log_parse::parse(text, log_format);
    let report = matcher.consider(text, true);
    line.clear();

    tx.send(SinkEvent::Line(parsed))
        .await
        .map_err(|_| std::io::Error::other("log writer stopped"))?;
    if let Some(report) = report {
        tx.send(SinkEvent::Report(report))
            .await
            .map_err(|_| std::io::Error::other("log writer stopped"))?;
    }
    Ok(())
}

/// Length to cut `bytes` at so no character is left half-written.
///
/// Decided by inspecting the final bytes rather than by asking `from_utf8` where
/// the string stops being valid: that reports the *first* problem, so a single
/// invalid byte earlier in the line would hide an unfinished character at the
/// end, and the character would be split after all.
fn split_before_incomplete_char(bytes: &[u8]) -> usize {
    let len = bytes.len();
    // A character is at most four bytes, so only the last few can be unfinished.
    for i in (len.saturating_sub(4)..len).rev() {
        let byte = bytes[i];
        if byte & 0b1100_0000 == 0b1000_0000 {
            continue; // a continuation byte; keep looking back for its lead
        }
        let expected = match byte {
            0x00..=0x7f => 1,
            b if b >> 5 == 0b110 => 2,
            b if b >> 4 == 0b1110 => 3,
            b if b >> 3 == 0b11110 => 4,
            // Not a valid lead byte at all, so nothing is pending; the lossy
            // conversion will render it.
            _ => 1,
        };
        return if i + expected > len && i > 0 { i } else { len };
    }
    len
}

/// Parse `line` and hand it to the writer, clearing it either way.
///
/// A closed queue means the writer task is gone, which is a failure rather than
/// the end of the stream: reporting it as success would tell the supervisor this
/// sink had reached end of file, and it would stop replacing it while the daemon
/// was still writing.
async fn queue(
    tx: &tokio::sync::mpsc::Sender<SinkEvent>,
    line: &mut Vec<u8>,
    log_format: &str,
    matcher: &mut ReadyMatcher,
) -> std::io::Result<()> {
    // Convert lossily: a daemon emitting a stray non-UTF-8 byte must not be able
    // to stop its own logging.
    let text = String::from_utf8_lossy(line);
    let text = text.trim_end_matches('\r');
    let parsed = crate::log_parse::parse(text, log_format);
    // Strip ANSI before matching so a pattern works whether or not the daemon
    // colours its output, matching what in-process capture did.
    let report = matcher.consider(text, false);
    line.clear();

    tx.send(SinkEvent::Line(parsed))
        .await
        .map_err(|_| std::io::Error::other("log writer stopped"))?;
    if let Some(report) = report {
        tx.send(SinkEvent::Report(report))
            .await
            .map_err(|_| std::io::Error::other("log writer stopped"))?;
    }
    Ok(())
}

/// Write queued lines in batches until the queue closes.
async fn write_batches(
    id: DaemonId,
    relay_token: u64,
    mut rx: tokio::sync::mpsc::Receiver<SinkEvent>,
) {
    let mut events: Vec<SinkEvent> = Vec::with_capacity(BATCH_SIZE);
    let mut batch: Vec<ParsedLog> = Vec::with_capacity(BATCH_SIZE);
    let mut flush_interval = tokio::time::interval(FLUSH_INTERVAL);
    flush_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);

    loop {
        let closed = tokio::select! {
            received = rx.recv_many(&mut events, BATCH_SIZE) => received == 0,
            _ = flush_interval.tick() => false,
        };
        for event in events.drain(..) {
            match event {
                SinkEvent::Line(parsed) => batch.push(parsed),
                SinkEvent::Report(report) => {
                    // Everything up to and including the reported line goes to
                    // the store before the supervisor is told, so whatever it
                    // does next can already read it.
                    flush(&id, &mut batch).await;
                    report_line(&id, relay_token, report).await;
                }
            }
        }
        flush(&id, &mut batch).await;
        if closed {
            break;
        }
    }
}

/// Hand the supervisor a line it needs to act on.
///
/// Failure is logged and dropped. There is no supervisor to retry against if it
/// has crashed — and if it has, nothing is waiting on this daemon's readiness or
/// hooks — while the daemon's output keeps being captured either way.
async fn report_line(id: &DaemonId, relay_token: u64, report: Report) {
    // `autostart: false` — a sink must never bring a supervisor into being.
    match crate::ipc::client::IpcClient::connect(false).await {
        Ok(client) => {
            if let Err(e) = client
                .sink_output_line(id.clone(), relay_token, report.fires_hook, report.text)
                .await
            {
                warn!("log sink for {id} could not report a line of output: {e}");
            }
        }
        Err(e) => {
            warn!("log sink for {id} could not reach the supervisor to report output: {e}");
        }
    }
}

/// Write one batch, off the runtime so the SQLite call cannot stall other tasks.
async fn flush(id: &DaemonId, batch: &mut Vec<ParsedLog>) {
    if batch.is_empty() {
        return;
    }
    let daemon_id = id.clone();
    let entries = std::mem::take(batch);
    let written = tokio::task::spawn_blocking(move || {
        LOG_STORE.append_structured_batch(&daemon_id, &entries)
    })
    .await;
    if let Ok(Err(e)) = written {
        // Nothing useful to do but report it: the supervisor is not necessarily
        // alive to be told, and dropping a batch is preferable to stalling the
        // daemon behind a pipe nobody is draining.
        error!("log sink failed to write batch for {id}: {e}");
    }
}

#[cfg(test)]
mod tests {
    use super::{HookMatcher, MATCH_CARRY_BYTES, ReadyMatcher};

    fn matcher(pattern: &str) -> ReadyMatcher {
        ReadyMatcher::new(Some(regex::Regex::new(pattern).unwrap()), None)
    }

    fn hook_matcher(filter: Option<&str>, debounce_ms: u64) -> ReadyMatcher {
        ReadyMatcher::new(
            None,
            Some(HookMatcher {
                filter: filter.map(str::to_string),
                regex: None,
                debounce: std::time::Duration::from_millis(debounce_ms),
                last_reported: None,
            }),
        )
    }

    #[test]
    fn reports_the_first_match_and_then_stops_looking() {
        let mut m = matcher("READY");
        assert_eq!(m.consider("starting up", false), None);
        assert_eq!(
            m.consider("READY to serve", false)
                .map(|r| r.text)
                .as_deref(),
            Some("READY to serve")
        );
        // Readiness happens once; later matches are somebody else's business.
        assert_eq!(m.consider("READY again", false), None);
    }

    #[test]
    fn matches_a_pattern_split_across_the_line_cap() {
        // The daemon emitted one enormous line whose announcement straddles the
        // point where the sink had to cut it.
        let mut m = matcher("SERVER READY");
        assert_eq!(m.consider("....SERVER ", true), None);
        let report = m
            .consider("READY....", false)
            .expect("should match across the split");
        assert!(
            report.text.contains("SERVER READY"),
            "reported {:?}",
            report.text
        );
    }

    #[test]
    fn does_not_match_across_a_completed_line() {
        // Two separate lines are not one line: a pattern spanning them must not
        // match, or "SERVER" at the end of one line plus "READY" at the start of
        // the next would look like an announcement.
        let mut m = matcher("SERVER READY");
        assert_eq!(m.consider("SERVER ", false), None);
        assert_eq!(m.consider("READY", false), None);
    }

    #[test]
    fn carries_a_bounded_amount_of_a_capped_line() {
        let mut m = matcher("nothing-matches-this");
        m.consider(&"x".repeat(MATCH_CARRY_BYTES * 3), true);
        assert!(m.carried.len() <= MATCH_CARRY_BYTES);
    }

    #[test]
    fn carrying_never_splits_a_character() {
        // A carry boundary landing mid-character would panic on the slice, or
        // corrupt the text a pattern is matched against.
        let mut m = matcher("nothing-matches-this");
        m.consider(&"é".repeat(MATCH_CARRY_BYTES), true);
        assert!(m.carried.chars().all(|c| c == 'é'));
    }

    #[test]
    fn hook_reports_matching_lines_within_the_debounce_window() {
        let mut m = hook_matcher(Some("ALERT"), 0);
        assert_eq!(m.consider("nothing here", false), None);
        let report = m.consider("ALERT disk full", false).expect("should report");
        assert_eq!(report.text, "ALERT disk full");
        assert!(report.fires_hook);
        // Unlike readiness, a hook keeps firing.
        assert!(m.consider("ALERT again", false).is_some());
    }

    #[test]
    fn hook_debounce_suppresses_a_second_line_in_the_same_window() {
        let mut m = hook_matcher(None, 60_000);
        assert!(m.consider("first", false).is_some());
        // Every line matches a hook with no filter, so only the window stops it.
        assert_eq!(m.consider("second", false), None);
    }

    #[test]
    fn a_readiness_only_line_does_not_fire_the_hook() {
        // A daemon can have both, filtered differently. The line that announces
        // readiness is nothing to do with the hook, and reporting it must not
        // be taken as permission to run the hook's command.
        let mut m = ReadyMatcher::new(
            Some(regex::Regex::new("READY").unwrap()),
            Some(HookMatcher {
                filter: Some("ALERT".to_string()),
                regex: None,
                debounce: std::time::Duration::from_millis(0),
                last_reported: None,
            }),
        );
        let report = m.consider("READY to serve", false).expect("should report");
        assert!(!report.fires_hook, "readiness line must not fire the hook");
        // ...and a line matching both says so.
        let mut m = ReadyMatcher::new(
            Some(regex::Regex::new("READY").unwrap()),
            Some(HookMatcher {
                filter: Some("READY".to_string()),
                regex: None,
                debounce: std::time::Duration::from_millis(0),
                last_reported: None,
            }),
        );
        assert!(
            m.consider("READY", false)
                .expect("should report")
                .fires_hook
        );
    }

    #[test]
    fn strips_ansi_before_matching() {
        let mut m = matcher("^READY$");
        assert!(m.consider("\x1b[32mREADY\x1b[0m", false).is_some());
    }
}