kanade-backend 0.51.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
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
1001
1002
1003
1004
1005
1006
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
//! Issue #246 — `/api/obs_events*` routes.
//!
//! Drives the SPA Events page (per-PC timeline) and the planned
//! fleet-wide observability dashboard. Three endpoints:
//!
//! - `GET /api/obs_events?pc_id=&from=&to=&kind=&source=&limit=`
//!   Filtered list, ordered by `at DESC`. Filters are optional;
//!   omitting them all returns the most recent events fleet-wide.
//! - `GET /api/obs_events/kinds` — distinct `kind` strings the
//!   SPA's filter chip needs to populate without a separate query.
//! - `GET /api/obs_events/sources` — distinct `source` strings for
//!   the include/exclude chips (Issue #391).
//! - `GET /api/obs_events/recent?limit=` — convenience alias for
//!   "newest N events fleet-wide" (same as `obs_events` with no
//!   `pc_id` and `limit` default 50).
//!
//! Pagination is keyset (`before_id`) rather than offset so a long-
//! tail timeline view doesn't drift when new events arrive between
//! pages — matches the inventory and audit endpoints' shape.
//! Pagination is deferred to the SPA PR; the first cut returns
//! up to `limit` rows in a single call.

use axum::Json;
use axum::extract::{Query, State};
use axum::http::StatusCode;
use chrono::{DateTime, Utc};

use super::sql_like::contains_like;
use super::time_bounds::bounds_in_range;
use serde::{Deserialize, Serialize};
use sqlx::sqlite::SqliteRow;
use sqlx::{Row, SqlitePool};
use tracing::warn;

/// Default page size when the caller doesn't specify `limit`.
/// Generous enough to render a "today on this PC" table without
/// pagination chrome; cheap enough server-side that an accidental
/// no-filter call doesn't pull millions of rows.
const DEFAULT_LIMIT: i64 = 200;
/// Hard ceiling. A misbehaving caller asking for `limit=1_000_000`
/// would otherwise exhaust SQLite's working memory on a busy fleet.
const MAX_LIMIT: i64 = 5_000;

#[derive(Deserialize)]
pub struct ListQuery {
    pub pc_id: Option<String>,
    /// RFC3339 lower bound (inclusive).
    pub from: Option<DateTime<Utc>>,
    /// RFC3339 upper bound (exclusive).
    pub to: Option<DateTime<Utc>>,
    /// Exact-match filter on `kind` (e.g. `logon`, `boot`).
    pub kind: Option<String>,
    /// Exact-match filter on `source` (e.g. `winlog:Security`).
    pub source: Option<String>,
    /// Exact-match filter on `payload.logon_type` (Issue #366).
    /// Windows LogonType numbers: 2 interactive, 3 network,
    /// 4 batch, 5 service, 7 unlock, 10 RDP, 11 cached. Only
    /// logon/logoff events carry the field, so combining this
    /// with a non-logon `kind` filter returns nothing — which is
    /// the honest answer. Kept for URL compatibility; the SPA now
    /// sends the generic `payload_key`/`payload_value` pair
    /// (Issue #391) instead.
    pub logon_type: Option<i64>,
    /// Issue #391: comma-separated include / exclude lists for
    /// `kind` and `source`. Both compose with the single-value
    /// `kind` / `source` gates above (which stay for URL
    /// compatibility); an empty include list (after splitting)
    /// means "no constraint", same as absent.
    pub kinds: Option<String>,
    pub kinds_ex: Option<String>,
    pub sources: Option<String>,
    pub sources_ex: Option<String>,
    /// Issue #391: generic payload filter — `payload_key=user` +
    /// `payload_value=yukimemi` matches rows whose
    /// `payload.<key> == <value>`. The key is restricted to
    /// `[A-Za-z0-9_]+` (400 otherwise) so the `'$.' || ?` path
    /// concat below can't be steered into other JSONPath syntax.
    /// The value is matched as text AND, when it parses as a
    /// number, numerically — so `payload_value=2` matches the
    /// JSON number 2 that the collectors emit.
    pub payload_key: Option<String>,
    pub payload_value: Option<String>,
    /// Issue #1343: match PCs whose operator-managed `agent_meta`
    /// carries this text in ANY value — the same "global attribute
    /// search" the Agents page offers (#1061), so an operator can ask
    /// "events from the sales department's machines" without knowing
    /// which key holds that.
    ///
    /// Filters the PC set, not the events: `agent_meta` describes
    /// machines, and an event's own payload has nothing to do with it.
    /// A PC with no metadata projected is therefore excluded outright —
    /// see the empty-state note in the SPA, which has to distinguish
    /// "no machine matches" from "metadata was never populated".
    pub meta_any: Option<String>,
    pub limit: Option<i64>,
}

/// Turn a comma-separated filter list into the JSON-array string
/// the `json_each(?)` binds expect. Blank segments are dropped;
/// an empty result collapses to `None` ("no constraint") so a
/// trailing comma can't accidentally filter everything out.
fn csv_to_json_array(csv: &Option<String>) -> Option<String> {
    let vals: Vec<&str> = csv
        .as_deref()?
        .split(',')
        .map(str::trim)
        .filter(|s| !s.is_empty())
        .collect();
    if vals.is_empty() {
        return None;
    }
    // Values are data, not SQL — serde_json escaping keeps the
    // array well-formed whatever the kind/source strings contain.
    serde_json::to_string(&vals).ok()
}

#[derive(Serialize)]
pub struct EventRow {
    pub id: i64,
    pub pc_id: String,
    pub at: DateTime<Utc>,
    pub kind: String,
    pub source: String,
    pub event_record_id: Option<String>,
    /// Parsed back into a JSON Value so the SPA receives structured
    /// data instead of a stringified JSON column.
    pub payload: serde_json::Value,
}

#[derive(Serialize)]
pub struct ListResponse {
    pub events: Vec<EventRow>,
}

/// The filter gates that are the same whichever statement runs.
///
/// One copy, spliced into both by `concat!`, because two hand-maintained
/// copies of a fifteen-bind WHERE is precisely the shape that drifts — and a
/// drift here is a filter that silently stops applying on one code path.
/// `concat!` runs at compile time, so each statement is still a single
/// `&'static str` and the crate's no-dynamic-SQL lint is satisfied by
/// construction rather than by review.
macro_rules! list_sql {
    ($pc_gate:literal) => {
        list_sql!("", $pc_gate)
    };
    ($prefix:literal, $pc_gate:literal) => {
        concat!(
            $prefix,
            "SELECT id, pc_id, at, kind, source, event_record_id, payload
             FROM obs_events
             WHERE ",
            $pc_gate,
            "
           AND (?2 IS NULL OR at    >= ?2)
           AND (?3 IS NULL OR at    <  ?3)
           AND (?4 IS NULL OR kind   = ?4)
           AND (?5 IS NULL OR source = ?5)
           AND (?6 IS NULL OR json_extract(payload, '$.logon_type') = ?6)
           AND (?7 IS NULL OR kind   IN     (SELECT value FROM json_each(?7)))
           AND (?8 IS NULL OR kind   NOT IN (SELECT value FROM json_each(?8)))
           AND (?9 IS NULL OR source IN     (SELECT value FROM json_each(?9)))
           AND (?10 IS NULL OR source NOT IN (SELECT value FROM json_each(?10)))
           AND (?11 IS NULL
                OR json_extract(payload, '$.' || ?11) = ?12
                OR (?13 IS NOT NULL AND json_extract(payload, '$.' || ?11) = ?13))
           -- `ESCAPE '\\'` in the Rust source is a single backslash in the
           -- SQL: `contains_like` escapes the operator's `%` / `_` with
           -- it. Writing `'\'` here would emit an EMPTY escape clause,
           -- which SQLite rejects only when the LIKE is actually
           -- evaluated — so the bug hides behind the `?14 IS NULL`
           -- short-circuit and every unfiltered query still passes.
           AND (?14 IS NULL
                OR pc_id IN (SELECT pc_id FROM agent_meta
                             WHERE value LIKE ?14 ESCAPE '\\'))
         ORDER BY at DESC, id DESC
         LIMIT ?15"
        )
    };
}

/// No `pc_id` filter: the gate is inert and the planner walks `at` newest
/// first, which is what the default view wants.
const LIST_ANY_PC: &str = list_sql!("(?1 IS NULL OR pc_id  = ?1)");

/// `pc_id` filter present, stated as a plain equality so SQLite can use
/// `idx_obs_events_pc_at`.
///
/// The NULL-gated form `(?1 IS NULL OR pc_id = ?1)` cannot be an index
/// constraint: whether the bind is NULL is a run-time fact, so the planner has
/// to keep the scan general. `EXPLAIN QUERY PLAN` showed `SCAN obs_events
/// USING INDEX idx_obs_events_at` for every combination of filters — neither
/// `idx_obs_events_pc_at` nor `idx_obs_events_kind_at` was ever reachable
/// through this query.
///
/// That is fast while the filter is loose, because `LIMIT` fills quickly, and
/// degrades as it tightens, because the scan cannot stop until it has found
/// enough rows or exhausted the table. Narrowing therefore made the page
/// SLOWER, which is the opposite of what an operator expects and why the
/// report was hard to place. Measured on 504,000 rows (400 agents x 30 days),
/// one PC plus the `boot`/`shutdown` chips — the reported combination, and the
/// worst case because it matches ~60 rows in a month:
///
/// | statement | rows | time |
/// |---|---:|---:|
/// | NULL-gated | 60 | 542 ms |
/// | this one | 60 | **2.3 ms** |
///
/// One PC without kind filters goes 455 ms -> 10 ms. The default view is
/// unchanged; it never depended on these indexes.
///
/// Binds are identical in both statements — `?1` is still bound, it just no
/// longer has to be tested for NULL — so the call site keeps one bind
/// sequence and the two cannot disagree about parameter order.
const LIST_ONE_PC: &str = list_sql!("pc_id = ?1");

/// `EXPLAIN QUERY PLAN` for each statement, from the SAME template.
///
/// The plan test used to run a hand-written, cut-down query — so a regression
/// in the real statement (a gate added ahead of the `pc_id` equality, a
/// reordering that defeats the index) would have left it green, because it
/// was planning a different query. The planner reads the whole `WHERE` and
/// `ORDER BY`, so the only thing worth planning is the text that ships.
#[cfg(test)]
const EXPLAIN_ONE_PC: &str = list_sql!("EXPLAIN QUERY PLAN ", "pc_id = ?1");
#[cfg(test)]
const EXPLAIN_ANY_PC: &str = list_sql!("EXPLAIN QUERY PLAN ", "(?1 IS NULL OR pc_id  = ?1)");

/// Which statement this request runs.
///
/// Extracted rather than inlined at the call site so it can be pinned: a test
/// that only compares the two constants passes just as happily when the
/// handler stops choosing between them.
fn list_statement(pc_id: Option<&str>) -> &'static str {
    match pc_id {
        Some(_) => LIST_ONE_PC,
        None => LIST_ANY_PC,
    }
}

/// `GET /api/obs_events`.
pub async fn list(
    State(pool): State<SqlitePool>,
    Query(q): Query<ListQuery>,
) -> Result<Json<ListResponse>, StatusCode> {
    let limit = match q.limit {
        None => DEFAULT_LIMIT,
        Some(n) if n > 0 && n <= MAX_LIMIT => n,
        _ => return Err(StatusCode::BAD_REQUEST),
    };

    // Issue #1076/#1126: reject any lexically-uncomparable date bound
    // before it reaches the string-compared `at` gates below (see
    // `time_bounds`).
    if !bounds_in_range([q.from, q.to]) {
        return Err(StatusCode::BAD_REQUEST);
    }

    // Issue #391: payload key allow-list — the key is interpolated
    // into a JSONPath via `'$.' || ?`, so anything beyond
    // identifier characters (quotes, brackets, dots) is rejected
    // up front rather than left to SQLite's path parser.
    let payload_key = match q.payload_key.as_deref().map(str::trim) {
        None | Some("") => None,
        Some(k) if k.chars().all(|c| c.is_ascii_alphanumeric() || c == '_') => Some(k),
        Some(_) => return Err(StatusCode::BAD_REQUEST),
    };
    // The pair only constrains when BOTH halves are present —
    // otherwise neutralise it entirely. A key without a value
    // would otherwise bind `?12` to NULL and the
    // `json_extract(...) = NULL` comparison blanks the whole
    // result set for direct API callers (Gemini #394 medium; the
    // SPA always sends both).
    let payload_value = payload_key.and(q.payload_value.as_deref());
    let payload_key = payload_value.and(payload_key);
    // Numeric twin: collectors emit numbers as JSON numbers
    // (logon_type: 2), and SQLite's `json_extract` returns them
    // typed — a text-only bind would never match. f64 covers i64
    // payload values under SQLite's numeric comparison rules.
    let payload_value_num: Option<f64> = payload_value.and_then(|v| v.parse().ok());

    // Issue #1343: blank (or whitespace-only) is "no filter", not "match
    // every PC whose metadata contains the empty string" — the latter
    // would silently drop every host with no metadata the moment the box
    // was cleared.
    let meta_any_like = q
        .meta_any
        .as_deref()
        .map(str::trim)
        .filter(|s| !s.is_empty())
        .map(contains_like);

    let kinds_json = csv_to_json_array(&q.kinds);
    let kinds_ex_json = csv_to_json_array(&q.kinds_ex);
    let sources_json = csv_to_json_array(&q.sources);
    let sources_ex_json = csv_to_json_array(&q.sources_ex);

    // Static SQL with "param IS NULL OR column = param" gates per
    // optional filter, so the SQL string stays a `&'static str` —
    // what `kanade-backend`'s lint config requires (dynamic SQL is
    // blocked at the lint level to prevent accidental injection
    // surfaces).
    //
    // The gate is NOT a planner no-op, which is what #1350 was: a
    // run-time NULL check cannot be an index constraint, so every
    // filter combination scanned `idx_obs_events_at`. The lint's
    // requirement is that the SQL TEXT not depend on run-time
    // values — not that there be only one of them — so `pc_id`,
    // the filter that unlocks a usable index, now picks between two
    // fixed statements. See `LIST_ONE_PC`.
    // `json_extract` on the `?6` gate: `payload` is stored as JSON
    // text, so the logon_type filter digs into it at query time.
    // No index on the expression — acceptable because the filter
    // composes with the indexed gates above and the table is
    // cleanup-bounded (see cleanup.rs).
    // Issue #391 additions keep the same static-SQL discipline:
    // the include/exclude lists arrive as JSON-array strings and
    // unpack inside SQLite via `json_each(?)` — one bind per list,
    // no dynamic IN-clause assembly. The generic payload gate
    // builds its JSONPath from a validated identifier (`'$.' || ?`)
    // and compares against the text bind plus, when the value is
    // numeric, the f64 twin.
    let rows = sqlx::query(list_statement(q.pc_id.as_deref()))
        .bind(q.pc_id.as_deref())
        .bind(q.from)
        .bind(q.to)
        .bind(q.kind.as_deref())
        .bind(q.source.as_deref())
        .bind(q.logon_type)
        .bind(kinds_json)
        .bind(kinds_ex_json)
        .bind(sources_json)
        .bind(sources_ex_json)
        .bind(payload_key)
        .bind(payload_value)
        .bind(payload_value_num)
        .bind(meta_any_like)
        .bind(limit)
        .fetch_all(&pool)
        .await
        .map_err(|e| {
            warn!(error = %e, "obs_events list query");
            StatusCode::INTERNAL_SERVER_ERROR
        })?;

    let events = rows
        .into_iter()
        .filter_map(|r| match row_to_event(&r) {
            Ok(e) => Some(e),
            Err(e) => {
                // Gemini #248 HIGH: surface schema mismatches /
                // type errors instead of returning blank fields.
                // We can't propagate the error here without
                // changing the response shape; warn-log + drop
                // the row keeps the API consistent while making
                // any bug operator-visible via agent.log.
                warn!(error = %e, "obs_events: drop row that failed to decode");
                None
            }
        })
        .collect();

    Ok(Json(ListResponse { events }))
}

/// Decode one `obs_events` row into an `EventRow`. Errors propagate
/// (vs the previous `unwrap_or_default()` shape which silently
/// returned empty strings on a column-rename / type mismatch). The
/// caller drops the row + logs; an alternative would be to 500 the
/// whole response, but a single bad row in a 200-row page
/// shouldn't take the whole timeline down.
fn row_to_event(r: &SqliteRow) -> sqlx::Result<EventRow> {
    let raw: String = r.try_get("payload")?;
    // `payload` is JSON text we stored ourselves, so a parse
    // failure means data corruption (someone hand-edited the
    // table) rather than a schema mismatch — bubble the same
    // error type out so the warn-log captures both cases.
    let payload = serde_json::from_str(&raw).map_err(|e| {
        sqlx::Error::Decode(format!("obs_events.payload not valid JSON: {e}").into())
    })?;
    Ok(EventRow {
        id: r.try_get("id")?,
        pc_id: r.try_get("pc_id")?,
        at: r.try_get("at")?,
        kind: r.try_get("kind")?,
        source: r.try_get("source")?,
        event_record_id: r.try_get("event_record_id")?,
        payload,
    })
}

/// How many PCs one `lane_seeds` call may ask about. The swimlane draws at
/// most `CHART_MAX_PCS` (40) hosts, so this is that with headroom rather
/// than an arbitrary cap — a caller asking for more is asking for something
/// the strip cannot render.
const MAX_SEED_PCS: usize = 64;

#[derive(Deserialize)]
pub struct SeedsQuery {
    /// Comma-separated `pc_id`s — the hosts the strip is about to draw.
    pub pcs: String,
    /// The window start. Seeds are the newest event strictly before it.
    pub before: DateTime<Utc>,
}

/// `GET /api/obs_events/lane_seeds`.
///
/// The newest event before `before`, per PC and per swimlane lane.
///
/// The Events page fetches a WINDOW of events, so a host that did not reboot
/// inside it reports no power event at all, and the strip cannot tell "this
/// host has no winlog collector" from "this host simply stayed up" (#1256).
/// The Analytics `op_timeline` query has always seeded itself this way; this
/// gives the Events page the same footing so the two surfaces stop
/// disagreeing about the same host and window.
///
/// Deliberately a separate call rather than a wider `list`: seeding the whole
/// fleet would mean walking back per kind until every PC is covered, which a
/// host that last rebooted months ago makes unbounded. Scoped to the ≤40 PCs
/// actually drawn, each lookup is an index seek on `(pc_id, at DESC)` — the
/// same shape `op_timeline` already runs per PC.
pub async fn lane_seeds(
    State(pool): State<SqlitePool>,
    Query(q): Query<SeedsQuery>,
) -> Result<Json<ListResponse>, StatusCode> {
    let pcs: Vec<&str> = q
        .pcs
        .split(',')
        .map(str::trim)
        .filter(|s| !s.is_empty())
        .collect();
    if pcs.is_empty() || pcs.len() > MAX_SEED_PCS {
        return Err(StatusCode::BAD_REQUEST);
    }

    // One seek per PC rather than a dynamic `IN (…)`: keeps the SQL static
    // (the same reason `op_timeline` spells its kind lists out) and bounded
    // by MAX_SEED_PCS.
    //
    // The lane CASE must stay aligned with `OP_LANES` in the SPA's
    // OperationalTimeline.tsx and with the `op_timeline` query in
    // analytics.rs — see `tests::lane_seed_kinds_match_op_timeline`.
    let mut events = Vec::new();
    for pc in pcs {
        let rows = sqlx::query(
            "SELECT id, pc_id, at, kind, source, event_record_id, payload FROM (                SELECT id, pc_id, at, kind, source, event_record_id, payload,                       ROW_NUMBER() OVER (                         PARTITION BY CASE                           WHEN kind IN ('boot', 'shutdown', 'unexpected_shutdown',                                         'log_service_started', 'log_service_stopped') THEN 'power'                           WHEN kind IN ('logon', 'logoff') THEN 'session'                           WHEN kind IN ('sleep', 'resume') THEN 'sleep'                           WHEN kind IN ('active', 'idle') THEN 'active'                         END                         ORDER BY at DESC                       ) AS rn                FROM obs_events                WHERE pc_id = ?1 AND at < ?2                  AND kind IN ('boot', 'shutdown', 'unexpected_shutdown',                               'log_service_started', 'log_service_stopped',                               'logon', 'logoff', 'sleep', 'resume',                               'active', 'idle')              ) WHERE rn = 1 ORDER BY at",
        )
        .bind(pc)
        .bind(q.before)
        .fetch_all(&pool)
        .await
        .map_err(|e| {
            warn!(error = %e, pc_id = %pc, "lane_seeds query failed");
            StatusCode::INTERNAL_SERVER_ERROR
        })?;
        for r in &rows {
            match row_to_event(r) {
                Ok(ev) => events.push(ev),
                // Same posture as `list`: one undecodable row must not cost
                // the whole strip its seeds.
                Err(e) => warn!(error = %e, "skipping undecodable lane seed"),
            }
        }
    }
    Ok(Json(ListResponse { events }))
}

#[derive(Serialize)]
pub struct KindsResponse {
    pub kinds: Vec<String>,
}

/// `GET /api/obs_events/kinds`.
pub async fn kinds(State(pool): State<SqlitePool>) -> Result<Json<KindsResponse>, StatusCode> {
    let rows = sqlx::query("SELECT DISTINCT kind FROM obs_events ORDER BY kind")
        .fetch_all(&pool)
        .await
        .map_err(|e| {
            warn!(error = %e, "obs_events kinds query");
            StatusCode::INTERNAL_SERVER_ERROR
        })?;
    // Drop rows that fail to decode (same handling rationale as
    // `list` above — operator sees the warn, the API stays useful).
    let kinds = rows
        .into_iter()
        .filter_map(|r| match r.try_get::<String, _>("kind") {
            Ok(k) => Some(k),
            Err(e) => {
                warn!(error = %e, "obs_events kinds: drop row that failed to decode kind");
                None
            }
        })
        .collect();
    Ok(Json(KindsResponse { kinds }))
}

#[derive(Serialize)]
pub struct SourcesResponse {
    pub sources: Vec<String>,
}

/// `GET /api/obs_events/sources` (Issue #391) — distinct `source`
/// strings for the SPA's include/exclude chips, mirroring `kinds`.
pub async fn sources(State(pool): State<SqlitePool>) -> Result<Json<SourcesResponse>, StatusCode> {
    let rows = sqlx::query("SELECT DISTINCT source FROM obs_events ORDER BY source")
        .fetch_all(&pool)
        .await
        .map_err(|e| {
            warn!(error = %e, "obs_events sources query");
            StatusCode::INTERNAL_SERVER_ERROR
        })?;
    let sources = rows
        .into_iter()
        .filter_map(|r| match r.try_get::<String, _>("source") {
            Ok(s) => Some(s),
            Err(e) => {
                warn!(error = %e, "obs_events sources: drop row that failed to decode source");
                None
            }
        })
        .collect();
    Ok(Json(SourcesResponse { sources }))
}

#[derive(Deserialize)]
pub struct RecentQuery {
    pub limit: Option<i64>,
}

/// `GET /api/obs_events/recent?limit=`. Convenience alias for
/// `/api/obs_events` with no `pc_id`. Lower default `limit` (50)
/// suited to a dashboard "latest activity" card.
pub async fn recent(
    State(pool): State<SqlitePool>,
    Query(q): Query<RecentQuery>,
) -> Result<Json<ListResponse>, StatusCode> {
    let limit = match q.limit {
        None => 50,
        Some(n) if n > 0 && n <= MAX_LIMIT => n,
        _ => return Err(StatusCode::BAD_REQUEST),
    };
    // Same shape as `list` with no filters, just a different
    // default `limit`. Kept as a sibling handler (vs delegating to
    // `list`) so the response model + log labels stay specific to
    // "recent" — small clarity win over saving a few lines.
    let rows = sqlx::query(
        "SELECT id, pc_id, at, kind, source, event_record_id, payload
         FROM obs_events
         ORDER BY at DESC, id DESC
         LIMIT ?",
    )
    .bind(limit)
    .fetch_all(&pool)
    .await
    .map_err(|e| {
        warn!(error = %e, "obs_events recent query");
        StatusCode::INTERNAL_SERVER_ERROR
    })?;

    let events = rows
        .into_iter()
        .filter_map(|r| match row_to_event(&r) {
            Ok(e) => Some(e),
            Err(e) => {
                warn!(error = %e, "obs_events recent: drop row that failed to decode");
                None
            }
        })
        .collect();
    Ok(Json(ListResponse { events }))
}

#[cfg(test)]
mod tests {
    use super::*;
    use axum::extract::{Query, State};
    use axum::http::StatusCode;
    use chrono::{TimeZone, Utc};
    use sqlx::sqlite::SqlitePoolOptions;

    async fn fresh_pool() -> SqlitePool {
        let pool = SqlitePoolOptions::new()
            .max_connections(1)
            .connect("sqlite::memory:")
            .await
            .unwrap();
        sqlx::migrate!("./migrations").run(&pool).await.unwrap();
        // One ordinary row so a query that ISN'T rejected returns
        // something — that's what makes the "inverted filter dumps the
        // whole table" bug observable in the `Ok`-path assertions.
        sqlx::query(
            "INSERT INTO obs_events (pc_id, at, kind, source, event_record_id, payload)
             VALUES ('pc-01', ?, 'logon', 'winlog:Security', '1', '{}')",
        )
        .bind(Utc.with_ymd_and_hms(2026, 5, 28, 10, 41, 0).unwrap())
        .execute(&pool)
        .await
        .unwrap();
        pool
    }

    /// The two statements must be the SAME QUERY, differing only in whether
    /// the `pc_id` gate is stated as an index-usable equality.
    ///
    /// They are spliced from one `list_sql!` body so the shared gates cannot
    /// drift, but the splice itself could still be wrong — a bind renumbered
    /// in one arm, a gate accidentally landing inside the `$pc_gate` literal.
    /// This compares the rendered text rather than trusting that.
    #[test]
    fn the_two_statements_differ_only_in_the_pc_gate() {
        assert_eq!(
            LIST_ANY_PC.replace("(?1 IS NULL OR pc_id  = ?1)", "<GATE>"),
            LIST_ONE_PC.replace("pc_id = ?1", "<GATE>"),
        );
        // And the numbering is untouched: every bind the handler supplies
        // appears in both.
        for n in 1..=15 {
            let tok = format!("?{n}");
            assert!(LIST_ANY_PC.contains(&tok), "any-pc statement lost {tok}");
            assert!(LIST_ONE_PC.contains(&tok), "one-pc statement lost {tok}");
        }
    }

    /// `EXPLAIN QUERY PLAN`'s human-readable text lives in the `detail`
    /// column; column 0 is an integer node id, so `query_scalar` decodes the
    /// wrong thing and fails at run time rather than at compile time.
    ///
    /// Binds all fifteen parameters, because it plans the production
    /// statement rather than a reduced stand-in.
    async fn plan_details(pool: &SqlitePool, sql: &'static str) -> Vec<String> {
        sqlx::query(sql)
            .bind(Some("pc-01"))
            .bind(Option::<DateTime<Utc>>::None)
            .bind(Option::<DateTime<Utc>>::None)
            .bind(Option::<String>::None)
            .bind(Option::<String>::None)
            .bind(Option::<i64>::None)
            .bind(Option::<String>::None)
            .bind(Option::<String>::None)
            .bind(Option::<String>::None)
            .bind(Option::<String>::None)
            .bind(Option::<String>::None)
            .bind(Option::<String>::None)
            .bind(Option::<f64>::None)
            .bind(Option::<String>::None)
            .bind(50i64)
            .fetch_all(pool)
            .await
            .unwrap()
            .iter()
            .map(|r| r.get::<String, _>("detail"))
            .collect()
    }

    /// The handler must actually pick between them.
    ///
    /// Added after a perturbation found the gap: making the call site always
    /// use `LIST_ANY_PC` left every test in this module green, because they
    /// all exercised the constants and a hand-written `EXPLAIN` rather than
    /// the selection. The statements being right is worth nothing if nothing
    /// chooses the fast one.
    #[test]
    fn a_pc_filter_selects_the_sargable_statement() {
        assert_eq!(list_statement(Some("pc-01")), LIST_ONE_PC);
        assert_eq!(list_statement(None), LIST_ANY_PC);
        assert_ne!(LIST_ONE_PC, LIST_ANY_PC);
    }

    /// The point of the split: with a `pc_id` the planner must reach
    /// `idx_obs_events_pc_at` instead of scanning `idx_obs_events_at`.
    ///
    /// Plans `LIST_ONE_PC` itself, via an `EXPLAIN` variant spliced from the
    /// same `list_sql!` template. A hand-written stand-in would keep passing
    /// while the shipped statement regressed — the planner's answer depends
    /// on the whole `WHERE` and `ORDER BY`, so a cut-down query is a
    /// different question.
    ///
    /// Asserted on the PLAN, not on a duration, so it holds on a slow CI box
    /// and fails for the actual reason if someone re-gates `pc_id`.
    #[tokio::test]
    async fn the_pc_filter_reaches_its_index() {
        let pool = fresh_pool().await;
        let detail = plan_details(&pool, EXPLAIN_ONE_PC).await.join(" | ");
        assert!(
            detail.contains("idx_obs_events_pc_at"),
            "expected the pc index, got: {detail}"
        );
        assert!(
            detail.contains("SEARCH"),
            "expected a SEARCH (seek), got: {detail}"
        );

        let old = plan_details(&pool, EXPLAIN_ANY_PC).await.join(" | ");
        assert!(
            !old.contains("idx_obs_events_pc_at"),
            "the NULL-gated form should NOT reach the pc index; if it now does,              SQLite got smarter and this split may be unnecessary: {old}"
        );
    }

    /// Same rows out, whichever statement ran. The split is a planner
    /// concern; it must not change what the API returns.
    #[tokio::test]
    async fn both_statements_return_the_same_rows() {
        let pool = fresh_pool().await;
        let t = |h: u32| Utc.with_ymd_and_hms(2026, 5, 28, h, 0, 0).unwrap();
        seed_row(&pool, "pc-01", t(11), "boot", "10").await;
        seed_row(&pool, "pc-02", t(12), "boot", "11").await;
        seed_row(&pool, "pc-01", t(13), "shutdown", "12").await;

        let ids = |rows: Vec<sqlx::sqlite::SqliteRow>| -> Vec<String> {
            rows.iter()
                .map(|r| {
                    let pc: String = r.get("pc_id");
                    let kind: String = r.get("kind");
                    format!("{pc}/{kind}")
                })
                .collect()
        };
        let run = async |sql: &'static str| {
            sqlx::query(sql)
                .bind(Some("pc-01"))
                .bind(Option::<DateTime<Utc>>::None)
                .bind(Option::<DateTime<Utc>>::None)
                .bind(Option::<String>::None)
                .bind(Option::<String>::None)
                .bind(Option::<i64>::None)
                .bind(Option::<String>::None)
                .bind(Option::<String>::None)
                .bind(Option::<String>::None)
                .bind(Option::<String>::None)
                .bind(Option::<String>::None)
                .bind(Option::<String>::None)
                .bind(Option::<f64>::None)
                .bind(Option::<String>::None)
                .bind(50i64)
                .fetch_all(&pool)
                .await
                .unwrap()
        };
        let a = ids(run(LIST_ANY_PC).await);
        let b = ids(run(LIST_ONE_PC).await);
        assert_eq!(a, b, "the split changed the result set");
        // 13:00 shutdown, 11:00 boot, then `fresh_pool`'s 10:41 logon.
        assert_eq!(b, vec!["pc-01/shutdown", "pc-01/boot", "pc-01/logon"]);
    }

    /// A `ListQuery` with only the two date bounds set — every other
    /// filter absent. Keeps the handler tests focused on #1076.
    fn bounds_query(from: Option<DateTime<Utc>>, to: Option<DateTime<Utc>>) -> ListQuery {
        ListQuery {
            pc_id: None,
            from,
            to,
            kind: None,
            source: None,
            logon_type: None,
            kinds: None,
            kinds_ex: None,
            sources: None,
            sources_ex: None,
            payload_key: None,
            payload_value: None,
            meta_any: None,
            limit: None,
        }
    }

    // The `bound_in_range` unit test moved to `time_bounds`; the
    // handler-level guard tests below stay to prove `list` still 400s.

    #[tokio::test]
    async fn list_rejects_expanded_year_from_bound() {
        // The exact failure the issue reports: a year-10000 lower bound
        // used to sort below every row and return the whole table. It
        // must now 400 instead of silently dumping everything.
        let pool = fresh_pool().await;
        let q = bounds_query(
            Some(Utc.with_ymd_and_hms(10000, 1, 1, 0, 0, 0).unwrap()),
            None,
        );
        let res = list(State(pool), Query(q)).await;
        assert!(matches!(res, Err(StatusCode::BAD_REQUEST)));
    }

    #[tokio::test]
    async fn list_rejects_expanded_year_to_bound() {
        // The `to` side inverts the other way (answered 0 rows), equally
        // wrong — reject it too.
        let pool = fresh_pool().await;
        let q = bounds_query(
            None,
            Some(Utc.with_ymd_and_hms(10000, 1, 1, 0, 0, 0).unwrap()),
        );
        let res = list(State(pool), Query(q)).await;
        assert!(matches!(res, Err(StatusCode::BAD_REQUEST)));
    }

    #[tokio::test]
    async fn list_accepts_in_range_bounds() {
        // Control: an ordinary window straddling the seeded row is
        // accepted and returns it — proves the guard doesn't reject
        // legitimate bounds.
        let pool = fresh_pool().await;
        let q = bounds_query(
            Some(Utc.with_ymd_and_hms(2026, 1, 1, 0, 0, 0).unwrap()),
            Some(Utc.with_ymd_and_hms(2027, 1, 1, 0, 0, 0).unwrap()),
        );
        let res = list(State(pool), Query(q))
            .await
            .expect("in-range bounds must be accepted");
        assert_eq!(res.0.events.len(), 1);
    }
    /// Insert one obs_event, with a unique `event_record_id` so the table's
    /// `UNIQUE(pc_id, source, event_record_id)` doesn't collapse the rows.
    async fn seed_row(pool: &SqlitePool, pc: &str, at: DateTime<Utc>, kind: &str, rec: &str) {
        sqlx::query(
            "INSERT INTO obs_events (pc_id, at, kind, source, event_record_id, payload)
             VALUES (?, ?, ?, 'winlog:System', ?, '{}')",
        )
        .bind(pc)
        .bind(at)
        .bind(kind)
        .bind(rec)
        .execute(pool)
        .await
        .unwrap();
    }

    /// The lane CASE in `lane_seeds` must group kinds exactly the way
    /// `op_timeline` (analytics.rs) and `OP_LANES` (OperationalTimeline.tsx)
    /// do. There are three copies of that mapping; this pins the one this
    /// module owns, the same way `analytics::tests::op_timeline_kind_set_is_stable`
    /// pins the query next door.
    ///
    /// Behavioural rather than string-matching on the SQL: it asserts what
    /// the grouping DOES — one seed per lane, the newest of each — so a CASE
    /// edited into a different shape fails even if it still parses.
    #[tokio::test]
    async fn lane_seed_kinds_match_op_timeline() {
        let pool = fresh_pool().await;
        let at = |h: u32| Utc.with_ymd_and_hms(2026, 6, 17, h, 0, 0).unwrap();

        // Two events per lane before the window. The SECOND of each pair is
        // what a correct lane grouping returns.
        for (i, (older, newer)) in [
            ("boot", "shutdown"), // power
            ("logon", "logoff"),  // session
            ("sleep", "resume"),  // sleep
            ("active", "idle"),   // active
        ]
        .iter()
        .enumerate()
        {
            seed_row(&pool, "seedpc", at(i as u32 * 2), older, &format!("o{i}")).await;
            seed_row(
                &pool,
                "seedpc",
                at(i as u32 * 2 + 1),
                newer,
                &format!("n{i}"),
            )
            .await;
        }
        // Neither a lane kind nor an observation kind: must not seed anything.
        seed_row(&pool, "seedpc", at(15), "app_sample", "x1").await;
        // Observation kinds drive no lane, so they must not seed either — a
        // seeded `agent_online` would assert liveness the window never saw.
        seed_row(&pool, "seedpc", at(16), "agent_online", "x2").await;

        let res = lane_seeds(
            State(pool),
            Query(SeedsQuery {
                pcs: "seedpc".into(),
                before: Utc.with_ymd_and_hms(2026, 6, 18, 0, 0, 0).unwrap(),
            }),
        )
        .await
        .expect("lane_seeds must succeed");

        let mut got: Vec<&str> = res.0.events.iter().map(|e| e.kind.as_str()).collect();
        got.sort_unstable();
        assert_eq!(
            got,
            ["idle", "logoff", "resume", "shutdown"],
            "one seed per lane, the newest of each — and nothing from a              non-lane kind"
        );
    }

    #[tokio::test]
    async fn lane_seeds_rejects_an_empty_or_oversized_pc_list() {
        let pool = fresh_pool().await;
        let before = Utc.with_ymd_and_hms(2026, 6, 18, 0, 0, 0).unwrap();
        for pcs in [
            "".to_string(),
            ",, ,".to_string(),
            (0..MAX_SEED_PCS + 1)
                .map(|i| format!("pc{i}"))
                .collect::<Vec<_>>()
                .join(","),
        ] {
            let res = lane_seeds(State(pool.clone()), Query(SeedsQuery { pcs, before })).await;
            assert!(matches!(res, Err(StatusCode::BAD_REQUEST)));
        }
    }

    // ---- #1343: filter events by the PC's agent_meta ----

    /// `fresh_pool` plus a second host, and metadata for both. `pc-01`
    /// keeps the seeded `logon`; `pc-02` gets its own event so a filter
    /// that selects the wrong host is visible as a wrong row rather than
    /// as an empty result.
    async fn pool_with_meta() -> SqlitePool {
        let pool = fresh_pool().await;
        sqlx::query(
            "INSERT INTO obs_events (pc_id, at, kind, source, event_record_id, payload)
             VALUES ('pc-02', ?, 'boot', 'winlog:System', '2', '{}')",
        )
        .bind(Utc.with_ymd_and_hms(2026, 5, 28, 10, 42, 0).unwrap())
        .execute(&pool)
        .await
        .unwrap();
        for (pc, key, value) in [
            ("pc-01", "department", "Sales"),
            ("pc-01", "owner", "Ann"),
            ("pc-02", "department", "Engineering"),
        ] {
            sqlx::query("INSERT INTO agent_meta (pc_id, key, value) VALUES (?, ?, ?)")
                .bind(pc)
                .bind(key)
                .bind(value)
                .execute(&pool)
                .await
                .unwrap();
        }
        pool
    }

    async fn list_with_meta_any(pool: &SqlitePool, meta_any: Option<&str>) -> Vec<EventRow> {
        let q = ListQuery {
            meta_any: meta_any.map(str::to_string),
            ..bounds_query(None, None)
        };
        list(State(pool.clone()), Query(q)).await.unwrap().0.events
    }

    #[tokio::test]
    async fn meta_any_matches_a_value_under_any_key() {
        let pool = pool_with_meta().await;
        // `department` on one host, `owner` on the same host — the point
        // of the global search is that the operator needn't know which
        // key carries the text.
        let rows = list_with_meta_any(&pool, Some("Sales")).await;
        assert_eq!(rows.len(), 1);
        assert_eq!(rows[0].pc_id, "pc-01");

        let rows = list_with_meta_any(&pool, Some("Ann")).await;
        assert_eq!(rows.len(), 1);
        assert_eq!(rows[0].pc_id, "pc-01");
    }

    #[tokio::test]
    async fn meta_any_is_a_contains_match() {
        let pool = pool_with_meta().await;
        let rows = list_with_meta_any(&pool, Some("ngineer")).await;
        assert_eq!(rows.len(), 1);
        assert_eq!(rows[0].pc_id, "pc-02");
    }

    #[tokio::test]
    async fn absent_or_blank_meta_any_filters_nothing() {
        let pool = pool_with_meta().await;
        // Blank must mean "no filter", NOT "value LIKE '%%'". The two
        // differ exactly where it matters: a host with no metadata row
        // survives the first and is dropped by the second, so a cleared
        // search box would silently shrink the fleet.
        sqlx::query(
            "INSERT INTO obs_events (pc_id, at, kind, source, event_record_id, payload)
             VALUES ('pc-03-no-meta', ?, 'boot', 'winlog:System', '3', '{}')",
        )
        .bind(Utc.with_ymd_and_hms(2026, 5, 28, 10, 43, 0).unwrap())
        .execute(&pool)
        .await
        .unwrap();

        for probe in [None, Some(""), Some("   ")] {
            let rows = list_with_meta_any(&pool, probe).await;
            assert_eq!(rows.len(), 3, "probe {probe:?} should not filter");
            assert!(rows.iter().any(|r| r.pc_id == "pc-03-no-meta"));
        }
    }

    #[tokio::test]
    async fn meta_any_excludes_pcs_with_no_metadata() {
        // The documented consequence, asserted so it stays intentional:
        // the filter runs through `agent_meta`, so a host with nothing
        // projected cannot match. The SPA's empty state exists because
        // of this, and must not be dropped as redundant.
        let pool = pool_with_meta().await;
        sqlx::query(
            "INSERT INTO obs_events (pc_id, at, kind, source, event_record_id, payload)
             VALUES ('pc-03-no-meta', ?, 'boot', 'winlog:System', '3', '{}')",
        )
        .bind(Utc.with_ymd_and_hms(2026, 5, 28, 10, 43, 0).unwrap())
        .execute(&pool)
        .await
        .unwrap();
        let rows = list_with_meta_any(&pool, Some("Sales")).await;
        assert!(rows.iter().all(|r| r.pc_id != "pc-03-no-meta"));
    }

    #[tokio::test]
    async fn meta_any_treats_like_metacharacters_literally() {
        // `%` and `_` are LIKE wildcards. Unescaped, a search for `_`
        // matches every single-character value and `%` matches
        // everything — so an operator typing a literal underscore would
        // get the whole fleet back and read it as a match.
        let pool = fresh_pool().await;
        for (pc, value) in [("pc-01", "a_b"), ("pc-02", "axb")] {
            sqlx::query("INSERT INTO agent_meta (pc_id, key, value) VALUES (?, 'tag', ?)")
                .bind(pc)
                .bind(value)
                .execute(&pool)
                .await
                .unwrap();
        }
        sqlx::query(
            "INSERT INTO obs_events (pc_id, at, kind, source, event_record_id, payload)
             VALUES ('pc-02', ?, 'boot', 'winlog:System', '2', '{}')",
        )
        .bind(Utc.with_ymd_and_hms(2026, 5, 28, 10, 42, 0).unwrap())
        .execute(&pool)
        .await
        .unwrap();

        let rows = list_with_meta_any(&pool, Some("a_b")).await;
        assert_eq!(rows.len(), 1, "`_` must not act as a wildcard");
        assert_eq!(rows[0].pc_id, "pc-01");

        let rows = list_with_meta_any(&pool, Some("%")).await;
        assert!(rows.is_empty(), "`%` must not match every value");
    }

    #[tokio::test]
    async fn meta_any_composes_with_the_other_filters() {
        // It narrows the PC set; it must not widen or replace anything.
        let pool = pool_with_meta().await;
        let q = ListQuery {
            meta_any: Some("Sales".into()),
            kinds: Some("boot".into()),
            ..bounds_query(None, None)
        };
        let rows = list(State(pool.clone()), Query(q)).await.unwrap().0.events;
        // pc-01 matches the metadata but has only a `logon`; pc-02 has
        // the `boot` but the wrong department. The honest answer is none.
        assert!(rows.is_empty());
    }
}