kanade-backend 0.47.0

axum + SQLite projection backend for the kanade endpoint-management system. Hosts /api/* and the embedded SPA dashboard, projects JetStream streams into SQLite, drives the cron scheduler
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
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
//! Turns heartbeat gaps into durable `agent_offline` / `agent_online`
//! events.
//!
//! Heartbeats are core NATS with no history: a single `agents.last_heartbeat`
//! column, overwritten every 30 s. That is enough to say "this agent is not
//! reporting *right now*" and structurally incapable of saying "this agent
//! was not reporting last Tuesday" — so an outage, once it becomes history,
//! reverts to looking exactly like a machine that was quietly idle (#1089).
//!
//! obs_events, by contrast, are durable and carry their own `at`. Recording
//! the two transitions as events costs two rows per outage and makes the
//! outage permanently visible, instead of storing 2,880 "still alive"
//! messages per agent per day to reconstruct the same handful of edges.
//!
//! # What this deliberately will not claim
//!
//! Only outages observed *within a single backend run* are recorded. While
//! the backend is down nobody is watching heartbeats, so on restart every
//! agent looks stale — a naive sweep would announce a fleet-wide outage that
//! never happened. `observed_since` gates that: an agent whose last heartbeat
//! predates this process is one we have not seen report, and we say nothing
//! about it rather than guessing.
//!
//! The cost is real and accepted: an agent that drops while the backend is
//! restarting is never recorded. Reporting an outage that did not happen is
//! worse than missing one that did, and the whole point of these events is to
//! be trustworthy about what was observed.
use std::collections::HashMap;

use anyhow::Result;
use chrono::{DateTime, Utc};
use kanade_shared::{subject, wire::ObsEvent};
use sqlx::{Row, SqlitePool};
use tracing::{debug, info, warn};

use crate::api::agents::ALIVE_THRESHOLD;

/// `source` on the events this writes. Distinct from any `agent:*` scheme:
/// these are the backend's own inferences, not something a host reported, and
/// an operator reading the Events table should be able to tell the difference.
const SOURCE: &str = "backend:heartbeat-watchdog";

pub const KIND_OFFLINE: &str = "agent_offline";
pub const KIND_ONLINE: &str = "agent_online";

/// Per-agent state across sweeps.
#[derive(Debug, Clone, Copy)]
struct Outage {
    /// Last heartbeat before the gap — the `at` of the emitted offline event
    /// and the key that makes re-emission idempotent.
    since: DateTime<Utc>,
}

/// What a sweep concluded, before anything is published.
///
/// Split from the publishing so the rules — which agents are eligible, when
/// an outage opens, when it closes, when to stay quiet — can be tested
/// without NATS. Verifying them against a live broker would mean a second
/// backend sharing the production durable consumer, which would take messages
/// away from the real projector.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Action {
    /// Stopped reporting. `at` is the last heartbeat, the newest instant with
    /// evidence behind it.
    Offline { pc_id: String, at: DateTime<Utc> },
    /// Reporting again. `at` is the beat that proved it.
    Online {
        pc_id: String,
        at: DateTime<Utc>,
        since: DateTime<Utc>,
    },
}

impl Action {
    fn pc_id(&self) -> &str {
        match self {
            Action::Offline { pc_id, .. } | Action::Online { pc_id, .. } => pc_id,
        }
    }
}

/// Apply one sweep's outcomes to the open-outage map.
///
/// `outcomes` pairs each action with whether it was published. Pure, because
/// the interesting failure is not in the broker call but in what the state is
/// left as afterwards.
///
/// A sweep can emit two actions for one host — the close-then-reopen pair for
/// an agent that recovered and died again inside one interval. Handling those
/// independently loses data: if the `Online` publish fails and the `Offline`
/// then succeeds, inserting the new outage overwrites the entry that was
/// still waiting to be closed. The next sweep sees `outage.since == last`,
/// concludes it has already recorded this outage, and stays silent forever —
/// so the close is not retried, it is *dropped*, and the log ends up with the
/// gap-swallowing shape this feature exists to prevent.
///
/// So a host with any failed action this sweep is left entirely untouched.
/// `decide` regenerates the identical pair next time from the unchanged
/// state, and the two retry together.
fn apply_outcomes(
    open: &mut HashMap<String, Outage>,
    outcomes: &[(Action, bool)],
) -> (usize, usize) {
    let failed: std::collections::HashSet<&str> = outcomes
        .iter()
        .filter(|(_, ok)| !ok)
        .map(|(a, _)| a.pc_id())
        .collect();

    let mut went_offline = 0usize;
    let mut came_online = 0usize;
    for (action, ok) in outcomes {
        if !ok || failed.contains(action.pc_id()) {
            continue;
        }
        match action {
            Action::Offline { pc_id, at } => {
                open.insert(pc_id.clone(), Outage { since: *at });
                went_offline += 1;
            }
            Action::Online { pc_id, .. } => {
                open.remove(pc_id);
                came_online += 1;
            }
        }
    }
    (went_offline, came_online)
}

/// Pure decision step. `agents` is `(pc_id, last_heartbeat)`.
fn decide(
    agents: &[(String, DateTime<Utc>)],
    observed_since: DateTime<Utc>,
    open: &HashMap<String, Outage>,
    now: DateTime<Utc>,
) -> Vec<Action> {
    let cutoff = now - ALIVE_THRESHOLD;
    let mut out = Vec::new();
    for (pc_id, last) in agents {
        // Never observed reporting during this run — say nothing. Covers both
        // a backend restart (every agent looks stale until its next beat) and
        // an agent that has been dead since before we started.
        if *last < observed_since {
            continue;
        }
        if *last > cutoff {
            if let Some(outage) = open.get(pc_id) {
                out.push(Action::Online {
                    pc_id: pc_id.clone(),
                    at: *last,
                    since: outage.since,
                });
            }
            continue;
        }
        // Stale.
        if let Some(outage) = open.get(pc_id) {
            // Same outage as last sweep — already recorded, nothing to say.
            if outage.since == *last {
                continue;
            }
            // A NEWER heartbeat than the one that opened the outage, yet
            // already stale again: the agent recovered and died again inside
            // one sweep interval. We observed that beat, so the recovery is
            // evidence we hold and must not discard — closing the first
            // outage here is the only chance to record it. Skipping it would
            // leave two `agent_offline` events with no recovery between them,
            // and the strip would read the pair as one continuous outage,
            // swallowing an interval the host was demonstrably up.
            out.push(Action::Online {
                pc_id: pc_id.clone(),
                at: *last,
                since: outage.since,
            });
        }
        out.push(Action::Offline {
            pc_id: pc_id.clone(),
            at: *last,
        });
    }
    out
}

/// Which hosts are still mid-outage, given each one's NEWEST watchdog record.
///
/// Split out for the same reason `decide` is: the rule is one line and the
/// query around it needs a database. `latest` is `(pc_id, kind, at)`, one row
/// per host.
fn open_from_latest(latest: &[(String, String, DateTime<Utc>)]) -> HashMap<String, Outage> {
    latest
        .iter()
        .filter(|(_, kind, _)| kind == KIND_OFFLINE)
        .map(|(pc_id, _, at)| (pc_id.clone(), Outage { since: *at }))
        .collect()
}

pub struct AgentWatchdog {
    /// When this process started watching. Nothing before it is asserted.
    ///
    /// Compared against `agents.last_heartbeat`, which is the agent's own
    /// `Utc::now()` at send time — so this gate straddles two clocks that
    /// nothing here synchronises. An agent running far enough fast can push a
    /// genuinely pre-restart beat past the gate, at which point a later
    /// silence produces a durable `agent_offline` for a window the backend
    /// never actually watched.
    ///
    /// The same mixed comparison already underlies `ALIVE_THRESHOLD` in
    /// `api/agents.rs` and the scheduler, but those are live flags that
    /// self-correct on the next beat; this is the first place it feeds a
    /// permanent record, so the assumption is worth naming: the fleet's
    /// clocks are expected to be roughly in step (NTP). Skew beyond a couple
    /// of minutes degrades what these events mean.
    observed_since: DateTime<Utc>,
    /// Agents currently believed offline, keyed by pc_id.
    open: HashMap<String, Outage>,
}

impl AgentWatchdog {
    pub fn new(observed_since: DateTime<Utc>) -> Self {
        Self {
            observed_since,
            open: HashMap::new(),
        }
    }

    /// Re-open the outages a previous process left unclosed.
    ///
    /// `open` lives in this struct and nowhere else, so a backend restart
    /// forgets every outage in flight. The `agent_offline` events are already
    /// durable in `obs_events`; what is lost is the knowledge that they are
    /// still waiting for a close. `decide` then sees a recovered agent with
    /// no open entry and emits nothing, so the pair is never completed —
    /// permanently. The strip has no right edge to draw and hatches from the
    /// outage to the live edge, which is the reported "the events came back
    /// but the stretch in between stays unknown". A backend restart is a
    /// deploy, so this is routine rather than exotic.
    ///
    /// Closing a stretch that spans the restart is not a claim about what
    /// happened while nobody was watching — that claim is already in the log,
    /// written by the process that observed the host go quiet. This only
    /// bounds it. Leaving it unbounded is the stronger and more wrong of the
    /// two, and it is what "say nothing about what we did not observe" turns
    /// into if the close is never written.
    ///
    /// Newest watchdog record per host decides: an `agent_online` means the
    /// pair completed and there is nothing to reopen. Both kinds are
    /// considered whatever wrote them — an agent's own `agent:startup` closes
    /// an outage just as the watchdog's does, and the strip reads them the
    /// same way. Hosts with no `agents` row are skipped: `sweep` only looks at
    /// registered agents, so an entry for anything else could never be
    /// closed.
    ///
    /// Failure is not fatal. Without the seed the watchdog behaves exactly as
    /// it did before, so the caller logs and carries on.
    pub async fn restore(&mut self, pool: &SqlitePool) -> Result<usize> {
        let rows = sqlx::query(
            "SELECT pc_id, kind, at FROM ( \
               SELECT pc_id, kind, at, \
                      ROW_NUMBER() OVER (PARTITION BY pc_id ORDER BY at DESC, id DESC) AS rn \
               FROM obs_events \
               WHERE kind IN (?1, ?2) \
                 AND pc_id IN (SELECT pc_id FROM agents) \
             ) WHERE rn = 1",
        )
        .bind(KIND_OFFLINE)
        .bind(KIND_ONLINE)
        .fetch_all(pool)
        .await?;

        let latest: Vec<(String, String, DateTime<Utc>)> = rows
            .iter()
            .filter_map(|r| {
                Some((
                    r.try_get("pc_id").ok()?,
                    r.try_get("kind").ok()?,
                    r.try_get("at").ok()?,
                ))
            })
            .collect();

        self.open = open_from_latest(&latest);
        if !self.open.is_empty() {
            info!(
                outages = self.open.len(),
                "restored open agent outages recorded before this process started",
            );
        }
        Ok(self.open.len())
    }

    /// One pass. Returns `(offline_emitted, online_emitted)`.
    ///
    /// Publishes to `obs.<pc_id>` so the events take the same durable path as
    /// an agent's own — they land in the same stream, are projected by the
    /// same consumer, and are subject to the same retention. The backend does
    /// not write `obs_events` rows directly; doing so here would create a
    /// second ingest path that the stream never sees.
    pub async fn sweep(
        &mut self,
        pool: &SqlitePool,
        js: &async_nats::jetstream::Context,
        now: DateTime<Utc>,
    ) -> Result<(usize, usize)> {
        let rows = sqlx::query(
            "SELECT pc_id, last_heartbeat FROM agents WHERE last_heartbeat IS NOT NULL",
        )
        .fetch_all(pool)
        .await?;

        let agents: Vec<(String, DateTime<Utc>)> = rows
            .iter()
            .filter_map(|r| {
                let pc_id: String = r.try_get("pc_id").ok()?;
                let last: DateTime<Utc> = r.try_get("last_heartbeat").ok()?;
                Some((pc_id, last))
            })
            .collect();

        // Publish first, collect what actually landed, then update state in
        // one pass. A host whose first action fails has its remaining actions
        // skipped rather than half-applied — see `apply_outcomes`.
        let mut outcomes: Vec<(Action, bool)> = Vec::new();
        let mut failed: std::collections::HashSet<String> = std::collections::HashSet::new();

        for action in decide(&agents, self.observed_since, &self.open, now) {
            if failed.contains(action.pc_id()) {
                // An earlier action for this host failed. Attempting this one
                // would record half a pair; leave the whole thing for the
                // next sweep, which regenerates it from unchanged state.
                outcomes.push((action, false));
                continue;
            }
            let published = match &action {
                Action::Offline { pc_id, at } => {
                    let r = self.publish(js, pc_id, KIND_OFFLINE, *at, *at).await;
                    match &r {
                        Ok(()) => info!(
                            pc_id = %pc_id,
                            last_heartbeat = %at,
                            threshold_secs = ALIVE_THRESHOLD.num_seconds(),
                            "agent stopped reporting; recorded agent_offline",
                        ),
                        Err(e) => {
                            warn!(pc_id = %pc_id, error = %e, "failed to publish agent_offline")
                        }
                    }
                    r.is_ok()
                }
                Action::Online { pc_id, at, since } => {
                    let r = self.publish(js, pc_id, KIND_ONLINE, *at, *since).await;
                    match &r {
                        Ok(()) => info!(
                            pc_id = %pc_id,
                            offline_since = %since,
                            back_at = %at,
                            "agent recovered; recorded agent_online",
                        ),
                        Err(e) => {
                            warn!(pc_id = %pc_id, error = %e, "failed to publish agent_online")
                        }
                    }
                    r.is_ok()
                }
            };
            if !published {
                failed.insert(action.pc_id().to_string());
            }
            outcomes.push((action, published));
        }

        let (went_offline, came_online) = apply_outcomes(&mut self.open, &outcomes);

        debug!(
            agents = agents.len(),
            went_offline, came_online, "agent watchdog sweep"
        );
        Ok((went_offline, came_online))
    }

    /// `at` is the instant the event describes; `key` distinguishes one
    /// outage from the next.
    ///
    /// `at` for an offline event is the **last heartbeat**, not that instant
    /// plus a heartbeat interval. A beat at T proves the agent was alive at
    /// T and nothing after it, so the unknown stretch opens at T. Padding
    /// forward would assert liveness across an interval nobody observed —
    /// and would need the agent's configured `heartbeat_interval`, which is
    /// per-PC and live-reloadable, to even compute.
    async fn publish(
        &self,
        js: &async_nats::jetstream::Context,
        pc_id: &str,
        kind: &str,
        at: DateTime<Utc>,
        key: DateTime<Utc>,
    ) -> Result<()> {
        // `event_record_id` must be Some and deterministic. The projector
        // dedups on UNIQUE(pc_id, source, event_record_id), and SQL NULL
        // never equals NULL — a None here would let every redelivery or
        // repeated sweep insert another row.
        let event = ObsEvent {
            pc_id: pc_id.to_string(),
            at,
            kind: kind.to_string(),
            source: SOURCE.to_string(),
            event_record_id: Some(format!("{kind}:{}", key.timestamp_millis())),
            payload: serde_json::json!({
                "last_heartbeat": key,
                "threshold_secs": ALIVE_THRESHOLD.num_seconds(),
            }),
        };
        let bytes = serde_json::to_vec(&event)?;
        let ack = js.publish(subject::obs(pc_id), bytes.into()).await?;
        ack.await?;
        Ok(())
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    fn ts(secs: i64) -> DateTime<Utc> {
        DateTime::from_timestamp(1_800_000_000 + secs, 0).unwrap()
    }

    #[test]
    fn offline_event_ids_are_deterministic_per_outage() {
        // Same outage → same id → the projector's UNIQUE drops the repeat.
        let a = format!("{KIND_OFFLINE}:{}", ts(0).timestamp_millis());
        let b = format!("{KIND_OFFLINE}:{}", ts(0).timestamp_millis());
        assert_eq!(a, b);
        // A later outage on the same host gets its own id.
        let c = format!("{KIND_OFFLINE}:{}", ts(600).timestamp_millis());
        assert_ne!(a, c);
    }

    #[test]
    fn online_and_offline_ids_do_not_collide() {
        // Both are keyed on the same instant when an outage is a single
        // sweep long; only the kind prefix separates them, so it must be
        // part of the id.
        let off = format!("{KIND_OFFLINE}:{}", ts(0).timestamp_millis());
        let on = format!("{KIND_ONLINE}:{}", ts(0).timestamp_millis());
        assert_ne!(off, on);
    }

    /// `observed_since` well before every fixture, so it is out of the way
    /// unless a test is specifically about it.
    const WATCH_START: i64 = 0;

    fn open_with(pc: &str, since: i64) -> HashMap<String, Outage> {
        HashMap::from([(pc.to_string(), Outage { since: ts(since) })])
    }

    /// `now` far enough past the heartbeat to be stale.
    fn stale_now(hb: i64) -> DateTime<Utc> {
        ts(hb) + ALIVE_THRESHOLD + chrono::Duration::seconds(1)
    }

    #[test]
    fn a_fresh_agent_produces_nothing() {
        let agents = vec![("pc1".into(), ts(100))];
        let out = decide(&agents, ts(WATCH_START), &HashMap::new(), ts(110));
        assert!(out.is_empty());
    }

    #[test]
    fn a_stale_agent_goes_offline_at_its_last_heartbeat() {
        let agents = vec![("pc1".into(), ts(100))];
        let out = decide(&agents, ts(WATCH_START), &HashMap::new(), stale_now(100));
        assert_eq!(
            out,
            vec![Action::Offline {
                pc_id: "pc1".into(),
                at: ts(100),
            }]
        );
    }

    #[test]
    fn an_already_recorded_outage_is_not_repeated() {
        // The sweep runs every 5 minutes; without this an agent down for an
        // hour would emit a dozen identical events. (The projector would
        // dedup them on the deterministic id, but re-publishing every tick
        // is still wasted work and noise in the log.)
        let agents = vec![("pc1".into(), ts(100))];
        let out = decide(
            &agents,
            ts(WATCH_START),
            &open_with("pc1", 100),
            stale_now(100),
        );
        assert!(out.is_empty());
    }

    #[test]
    fn recovery_closes_the_outage() {
        let agents = vec![("pc1".into(), ts(900))];
        let out = decide(&agents, ts(WATCH_START), &open_with("pc1", 100), ts(905));
        assert_eq!(
            out,
            vec![Action::Online {
                pc_id: "pc1".into(),
                at: ts(900),
                since: ts(100),
            }]
        );
    }

    // Recovered and died again inside one sweep interval. The beat at 900 was
    // observed, so it has to be recorded — otherwise the log holds two
    // `agent_offline` events with no recovery between them and the strip
    // reads them as one continuous outage, swallowing an interval the host
    // was demonstrably up.
    #[test]
    fn a_recovery_missed_between_sweeps_is_still_recorded() {
        let agents = vec![("pc1".into(), ts(900))];
        let out = decide(
            &agents,
            ts(WATCH_START),
            &open_with("pc1", 100),
            stale_now(900),
        );
        assert_eq!(
            out,
            vec![
                Action::Online {
                    pc_id: "pc1".into(),
                    at: ts(900),
                    since: ts(100),
                },
                Action::Offline {
                    pc_id: "pc1".into(),
                    at: ts(900),
                },
            ],
            "the close must come first, and both must be emitted",
        );
    }

    // The property this whole module hinges on: while the backend was down
    // nobody watched heartbeats, so on restart every agent looks stale. A
    // sweep that ignored `observed_since` would announce a fleet-wide outage
    // that never happened.
    #[test]
    fn a_restart_does_not_invent_a_fleet_wide_outage() {
        let agents: Vec<(String, DateTime<Utc>)> =
            (0..50).map(|i| (format!("pc{i}"), ts(100))).collect();
        // Watching began after every one of those heartbeats.
        let out = decide(&agents, ts(500), &HashMap::new(), stale_now(500));
        assert!(
            out.is_empty(),
            "claimed {} outages for agents never observed reporting",
            out.len()
        );
    }

    #[test]
    fn an_agent_seen_after_the_watch_started_is_still_reported() {
        // The gate must not be so blunt that it silences everything after a
        // restart — an agent that reported and *then* went quiet is exactly
        // what we want to catch.
        let agents = vec![("pc1".into(), ts(600))];
        let out = decide(&agents, ts(500), &HashMap::new(), stale_now(600));
        assert_eq!(out.len(), 1);
    }

    #[test]
    fn the_threshold_boundary_is_not_stale() {
        // Exactly at the cutoff counts as reporting: `last > cutoff` is the
        // comparison, so a beat landing precisely on it must not flip.
        let agents = vec![("pc1".into(), ts(100))];
        let just_inside = ts(100) + ALIVE_THRESHOLD - chrono::Duration::seconds(1);
        assert!(decide(&agents, ts(WATCH_START), &HashMap::new(), just_inside).is_empty());
    }

    // ---- apply_outcomes: what a partial publish failure leaves behind ----

    fn off(pc: &str, at: i64) -> Action {
        Action::Offline {
            pc_id: pc.into(),
            at: ts(at),
        }
    }
    fn on(pc: &str, at: i64, since: i64) -> Action {
        Action::Online {
            pc_id: pc.into(),
            at: ts(at),
            since: ts(since),
        }
    }

    #[test]
    fn a_successful_pair_leaves_the_new_outage_open() {
        let mut open = open_with("pc1", 100);
        let out = apply_outcomes(
            &mut open,
            &[(on("pc1", 900, 100), true), (off("pc1", 900), true)],
        );
        assert_eq!(out, (1, 1));
        assert_eq!(open["pc1"].since, ts(900));
    }

    // The bug this guard exists for. A failed close followed by a successful
    // open used to overwrite the entry, so the next sweep saw
    // `outage.since == last`, decided the outage was already recorded, and
    // never retried the close — dropping it permanently and leaving the log
    // with the gap-swallowing shape the close was added to prevent.
    #[test]
    fn a_failed_close_does_not_let_the_reopen_bury_it() {
        let mut open = open_with("pc1", 100);
        let out = apply_outcomes(
            &mut open,
            &[(on("pc1", 900, 100), false), (off("pc1", 900), true)],
        );
        assert_eq!(out, (0, 0), "nothing may be counted as recorded");
        assert_eq!(
            open["pc1"].since,
            ts(100),
            "the outage still awaiting its close must survive untouched",
        );
        // And the next sweep must therefore regenerate the whole pair.
        let agents = vec![("pc1".into(), ts(900))];
        assert_eq!(
            decide(&agents, ts(WATCH_START), &open, stale_now(900)),
            vec![on("pc1", 900, 100), off("pc1", 900)],
        );
    }

    // Failure in the other order. The rule is deliberately blunt — any
    // failure for a host leaves that host's state entirely alone — rather
    // than "apply the ones that succeeded". Partial application is safe in
    // this direction and unsafe in the other, and depending on which is which
    // is precisely the subtlety that produced the bug this guard fixes.
    //
    // The cost is one redundant re-publish next sweep, which the deterministic
    // `event_record_id` makes a no-op at the projector.
    #[test]
    fn a_failed_open_leaves_the_host_untouched_and_replays_the_pair() {
        let mut open = open_with("pc1", 100);
        let out = apply_outcomes(
            &mut open,
            &[(on("pc1", 900, 100), true), (off("pc1", 900), false)],
        );
        assert_eq!(out, (0, 0));
        assert_eq!(open["pc1"].since, ts(100));

        let agents = vec![("pc1".into(), ts(900))];
        assert_eq!(
            decide(&agents, ts(WATCH_START), &open, stale_now(900)),
            vec![on("pc1", 900, 100), off("pc1", 900)],
            "the pair replays; the already-landed close dedups on its id",
        );
    }

    #[test]
    fn one_hosts_failure_does_not_hold_back_another() {
        let mut open = HashMap::new();
        let out = apply_outcomes(
            &mut open,
            &[(off("pc1", 100), false), (off("pc2", 100), true)],
        );
        assert_eq!(out, (1, 0));
        assert!(!open.contains_key("pc1"));
        assert_eq!(open["pc2"].since, ts(100));
    }

    #[test]
    fn hosts_are_judged_independently() {
        let agents = vec![
            ("fresh".into(), ts(1000)),
            ("stale".into(), ts(100)),
            ("unseen".into(), ts(-100)),
        ];
        // `now` chosen so `stale` is past the threshold while `fresh` is not.
        let out = decide(&agents, ts(0), &HashMap::new(), stale_now(100));
        assert_eq!(
            out,
            vec![Action::Offline {
                pc_id: "stale".into(),
                at: ts(100),
            }]
        );
    }

    // -----------------------------------------------------------------
    // Surviving a restart. `open` is in-memory, so without this an outage
    // that spanned a deploy never received its closing `agent_online` and
    // the strip hatched from the outage to the live edge forever.
    // -----------------------------------------------------------------

    fn latest(rows: &[(&str, &str, i64)]) -> Vec<(String, String, DateTime<Utc>)> {
        rows.iter()
            .map(|(pc, kind, at)| (pc.to_string(), kind.to_string(), ts(*at)))
            .collect()
    }

    #[test]
    fn a_host_whose_newest_record_is_offline_is_still_down() {
        let open = open_from_latest(&latest(&[("pc1", KIND_OFFLINE, 100)]));
        assert_eq!(open["pc1"].since, ts(100));
    }

    #[test]
    fn a_host_whose_outage_already_closed_is_not_reopened() {
        let open = open_from_latest(&latest(&[("pc1", KIND_ONLINE, 900)]));
        assert!(open.is_empty());
    }

    #[test]
    fn a_restored_outage_closes_on_the_next_beat() {
        // The whole point: the recovery the previous process could not record
        // is emitted by this one, keyed on the ORIGINAL outage instant so the
        // projector's UNIQUE(pc_id, source, event_record_id) still applies.
        let open = open_from_latest(&latest(&[("pc1", KIND_OFFLINE, 100)]));
        let agents = vec![("pc1".into(), ts(900))];
        assert_eq!(
            decide(&agents, ts(WATCH_START), &open, ts(905)),
            vec![Action::Online {
                pc_id: "pc1".into(),
                at: ts(900),
                since: ts(100),
            }],
        );
    }

    #[tokio::test]
    async fn restore_reads_the_newest_watchdog_record_per_host() {
        let pool = SqlitePool::connect("sqlite::memory:").await.unwrap();
        sqlx::migrate!("./migrations").run(&pool).await.unwrap();

        // down: dropped and never came back. back: dropped and recovered.
        // stray: an outage for a host with no `agents` row — nothing in
        // `sweep` could ever close it, so it must not be carried.
        for pc in ["down", "back"] {
            sqlx::query("INSERT INTO agents (pc_id, last_heartbeat) VALUES (?1, ?2)")
                .bind(pc)
                .bind(ts(900))
                .execute(&pool)
                .await
                .unwrap();
        }
        let rows = [
            ("down", KIND_OFFLINE, 100),
            ("back", KIND_OFFLINE, 100),
            ("back", KIND_ONLINE, 500),
            ("stray", KIND_OFFLINE, 100),
        ];
        for (i, (pc, kind, at)) in rows.iter().enumerate() {
            sqlx::query(
                "INSERT INTO obs_events (pc_id, at, kind, source, event_record_id, payload) \
                 VALUES (?1, ?2, ?3, ?4, ?5, '{}')",
            )
            .bind(pc)
            .bind(ts(*at))
            .bind(*kind)
            .bind(SOURCE)
            .bind(format!("{kind}:{i}"))
            .execute(&pool)
            .await
            .unwrap();
        }

        let mut wd = AgentWatchdog::new(ts(1000));
        assert_eq!(wd.restore(&pool).await.unwrap(), 1);
        assert_eq!(wd.open["down"].since, ts(100));
        assert!(!wd.open.contains_key("back"));
        assert!(!wd.open.contains_key("stray"));
    }
}