kanade-backend 0.54.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
1058
//! Job endpoints:
//!   * `POST /api/jobs/{job_id}/kill` — runtime control. Looks up every
//!     in-flight execution of `{job_id}` from the `executions` table
//!     (status pending / running) and publishes `kill.{exec_id}` per
//!     deployment so agents actually receive the signal (spec §2.6
//!     Layer 3). Pre-v0.29 this published `kill.{cmd_id}`, which no
//!     agent subscribes to — the kill button on the SPA was effectively
//!     a no-op since v0.27.
//!   * `GET / POST /api/jobs` + `DELETE /api/jobs/{id}` — catalog CRUD
//!     (v0.15). Schedules reference catalog rows by `job_id`.

use std::collections::HashMap;

use async_nats::jetstream::kv::Config as KvConfig;
use axum::Json;
use axum::extract::{Path, State};
use axum::http::StatusCode;
use axum::http::header::HeaderMap;
use futures::TryStreamExt;
use kanade_shared::kv::{
    BUCKET_JOBS, BUCKET_JOBS_YAML, BUCKET_SCHEDULES, BUCKET_SCRIPT_CURRENT, BUCKET_SCRIPT_STATUS,
    SCRIPT_STATUS_REVOKED,
};
use kanade_shared::manifest::{ExecuteShell, InventoryHint, Manifest, RepoOrigin, Schedule};
use kanade_shared::subject;
use kanade_shared::wire::RunAs;
use serde::Serialize;
use sqlx::{Row, SqlitePool};
use tracing::{info, warn};

use super::AppState;
use super::yaml_body::{YamlOrJson, mirror_yaml, yaml_headers};
use crate::audit;
use crate::audit::Caller;

pub async fn kill(
    State(state): State<AppState>,
    Path(job_id): Path<String>,
) -> Result<StatusCode, (StatusCode, String)> {
    // v0.29 / Issue #19: the agent listens on `kill.{exec_id}`, never
    // on `kill.{cmd_id}`. The path param here is the cmd / manifest
    // id, so we have to expand it to every still-running exec_id and
    // publish per-exec. status IN ('pending', 'running') skips
    // already-completed deployments — there's nothing to kill on those.
    let rows = sqlx::query(
        "SELECT exec_id FROM executions \
         WHERE job_id = ? AND status IN ('pending', 'running')",
    )
    .bind(&job_id)
    .fetch_all(&state.pool)
    .await
    .map_err(|e| {
        warn!(error = %e, job_id, "kill: lookup running execs");
        (
            StatusCode::INTERNAL_SERVER_ERROR,
            format!("lookup running execs: {e}"),
        )
    })?;

    let exec_ids: Vec<String> = rows
        .into_iter()
        .map(|r| r.try_get::<String, _>("exec_id").unwrap_or_default())
        .filter(|s| !s.is_empty())
        .collect();

    if exec_ids.is_empty() {
        // No running deployments → there's nothing to kill. Return
        // 204 anyway: the operator's mental model is "I clicked kill,
        // it's not running", which is what 204 + zero-published
        // conveys. A 404 here would just confuse the SPA.
        info!(
            %job_id,
            "kill: no running executions for this job (no-op)",
        );
        return Ok(StatusCode::NO_CONTENT);
    }

    for exec_id in &exec_ids {
        if let Err(e) = state
            .nats
            .publish(subject::kill(exec_id), bytes::Bytes::new())
            .await
        {
            warn!(error = %e, %job_id, %exec_id, "publish kill failed");
            return Err((
                StatusCode::INTERNAL_SERVER_ERROR,
                format!("publish kill.{exec_id}: {e}"),
            ));
        }
    }
    // flush so the subjects are on the wire before we ack the
    // operator — without it, a fast operator-then-shutdown could
    // theoretically drop the kill on the floor.
    let _ = state.nats.flush().await;
    info!(
        %job_id,
        kill_count = exec_ids.len(),
        "kill signal fanned out to running execs",
    );
    Ok(StatusCode::NO_CONTENT)
}

#[derive(Serialize)]
pub struct JobSummary {
    pub id: String,
    pub version: String,
    pub description: Option<String>,
    pub inventory: bool,
    /// #1492: human-readable notes about any `inventory.explode`
    /// derived-table schema migration this `create` triggered — e.g.
    /// `"example_items: added column 'kind'"` or `"example_items:
    /// rebuilt (primary_key/columns changed); copied 12 row(s), lost
    /// 1 row(s) to the new primary key"`. Empty when the manifest has
    /// no explode specs, or every explode table already matched the
    /// manifest. Surfaced by `kanade job create` so a schema drift
    /// that used to silently freeze a derived table now shows up in
    /// the CLI's own success output.
    #[serde(default, skip_serializing_if = "Vec::is_empty")]
    pub schema_changes: Vec<String>,
}

/// v0.30 / PR γ: in-flight counters joined onto each `/api/jobs`
/// row so the Jobs page can show "is anything running right now"
/// at a glance — the operator's decision input for kill / revoke
/// without having to drill into Activity. Sourced from
/// `executions.status`, which the v0.29 results projector now
/// maintains correctly.
#[derive(Serialize, Default, Debug, Clone, PartialEq, Eq)]
pub struct JobLiveCounts {
    /// `executions.status = 'running'` — at least one result has
    /// landed but more are still in flight.
    pub running: i64,
    /// `executions.status = 'pending'` — fan-out published but no
    /// result has landed yet. Distinguished from `running` so the
    /// operator can tell "nothing reported back" from "partially
    /// reported".
    pub pending: i64,
}

/// The `execute:` fields the Jobs page's catalog table renders
/// (`web/src/pages/Jobs.tsx` `JobRow`). The full `Execute` is NOT
/// serialised here — `execute.script` bodies were ~85 % of the list
/// payload (#1214①) and neither list consumer (Jobs page, Exec
/// picker) reads them. The YAML editor fetches the full source via
/// `GET /api/jobs/{id}/yaml` instead.
#[derive(Serialize)]
pub struct ExecuteSummary {
    pub shell: ExecuteShell,
    pub timeout: String,
    pub run_as: RunAs,
    #[serde(skip_serializing_if = "Option::is_none")]
    pub cwd: Option<String>,
}

/// `GET /api/jobs` row shape — a projection of the registered
/// Manifest carrying only the fields the SPA's two list consumers
/// read, plus the `live` counters aggregated from `executions`.
/// Field-for-field this keeps the old flattened shape (`job.id`,
/// `job.execute.shell`, …), so the SPA needs no change; what changed
/// is what's absent: `execute.script` / `script_file` /
/// `script_object` and the hint blocks the table never renders
/// (#1214①, ~500 KB → ~75 KB on the measured catalog).
///
/// `kanade job list` pretty-prints this verbatim — dropping the
/// script body from its output is the deliberate compatibility call
/// from #1214: `list` is a catalog overview and
/// `GET /api/jobs/{id}/yaml` remains the path to a full manifest.
#[derive(Serialize)]
pub struct JobListRow {
    pub id: String,
    pub version: String,
    pub description: Option<String>,
    pub execute: ExecuteSummary,
    pub inventory: Option<InventoryHint>,
    #[serde(skip_serializing_if = "Vec::is_empty")]
    pub tags: Vec<String>,
    #[serde(skip_serializing_if = "Option::is_none")]
    pub origin: Option<RepoOrigin>,
    pub live: JobLiveCounts,
}

impl JobListRow {
    fn from_manifest(m: Manifest, live: JobLiveCounts) -> Self {
        Self {
            id: m.id,
            version: m.version,
            description: m.description,
            execute: ExecuteSummary {
                shell: m.execute.shell,
                timeout: m.execute.timeout,
                run_as: m.execute.run_as,
                cwd: m.execute.cwd,
            },
            inventory: m.inventory,
            tags: m.tags,
            origin: m.origin,
            live,
        }
    }
}

/// GET /api/jobs — list every registered job + live in-flight
/// counters from the executions table.
pub async fn list(
    State(s): State<AppState>,
) -> Result<Json<Vec<JobListRow>>, (StatusCode, String)> {
    let kv = match s.jetstream.get_key_value(BUCKET_JOBS).await {
        Ok(k) => k,
        Err(_) => return Ok(Json(Vec::new())),
    };
    let keys_stream = kv
        .keys()
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("kv keys: {e}")))?;
    let keys: Vec<String> = keys_stream
        .try_collect()
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("kv keys: {e}")))?;
    // #1214②: fan the per-key GETs out in parallel — the previous
    // sequential loop serialised one NATS round-trip per catalog row
    // on the Jobs page's critical path. Mirrors scripts::list_status
    // ("Gemini #47 review"), bounded by the operator-facing catalog
    // size (~10s-100s of jobs).
    let fetches = keys.into_iter().map(|k| {
        let kv = kv.clone();
        async move {
            match kv.get(&k).await {
                Ok(Some(bytes)) => serde_json::from_slice::<Manifest>(&bytes).ok(),
                _ => None,
            }
        }
    });
    let mut manifests: Vec<Manifest> = futures::future::join_all(fetches)
        .await
        .into_iter()
        .flatten()
        .collect();
    manifests.sort_by(|a, b| a.id.cmp(&b.id));

    // v0.30 / PR γ: one GROUP BY query for the whole list instead of
    // N round-trips. `executions` lives in SQLite so this is local.
    // A missing `executions` row for a job (= never fired since the
    // backend started) yields the default zeros via the HashMap
    // lookup fallback below.
    let live_counts = fetch_live_counts(&s.pool).await.unwrap_or_else(|e| {
        warn!(error = %e, "jobs list: live count aggregation failed; returning zeros");
        HashMap::new()
    });

    let out: Vec<JobListRow> = manifests
        .into_iter()
        .map(|m| {
            let live = live_counts.get(&m.id).cloned().unwrap_or_default();
            JobListRow::from_manifest(m, live)
        })
        .collect();
    Ok(Json(out))
}

/// Aggregate `executions.status` counts by `job_id` so the
/// `/api/jobs` list can attach per-row live counters in one round
/// trip. Returned map omits jobs with no executions entirely; the
/// caller falls back to `JobLiveCounts::default()`.
async fn fetch_live_counts(
    pool: &SqlitePool,
) -> Result<HashMap<String, JobLiveCounts>, sqlx::Error> {
    // Gemini #71 perf fix: filter on status BEFORE the aggregation
    // so SQLite skips completed rows entirely instead of summing
    // CASE-zeros for the (growing forever) historical tail. The
    // resulting empty groups disappear from the output map; the
    // caller already falls back to `JobLiveCounts::default()` for
    // jobs not present in the map, so semantics are preserved.
    let rows = sqlx::query(
        "SELECT job_id,
                SUM(CASE WHEN status = 'running' THEN 1 ELSE 0 END) AS running,
                SUM(CASE WHEN status = 'pending' THEN 1 ELSE 0 END) AS pending
           FROM executions
          WHERE status IN ('running', 'pending')
          GROUP BY job_id",
    )
    .fetch_all(pool)
    .await?;
    let mut out: HashMap<String, JobLiveCounts> = HashMap::with_capacity(rows.len());
    for r in rows {
        let job_id: String = r.try_get("job_id").unwrap_or_default();
        if job_id.is_empty() {
            continue;
        }
        out.insert(
            job_id,
            JobLiveCounts {
                running: r.try_get("running").unwrap_or(0),
                pending: r.try_get("pending").unwrap_or(0),
            },
        );
    }
    Ok(out)
}

/// #1492 review (R2-1-1): build the error response for a catalog
/// write (KV bucket lookup / serialize / put) that failed AFTER
/// `ensure_tables_atomic` already committed an explode-table schema
/// migration for this manifest. There is no way to undo that SQLite
/// commit from here — it's a different transactional system than the
/// NATS KV write that just failed, so this handler can't offer
/// all-or-nothing across both. What it CAN do is make the gap loud
/// instead of silent: the `note` strings already computed (rebuilt?
/// which columns? rows lost to the new primary key?) go into both the
/// HTTP error body — so the CLI prints them instead of just "KV put:
/// ..." — and an audit record, so the drift is still discoverable
/// later even if the operator doesn't read this response closely.
/// `schema_changes` is empty on the common path (no explode specs, or
/// every explode table already matched), in which case this is just
/// the plain error.
async fn catalog_write_failed(
    s: &AppState,
    job: &Manifest,
    caller: &Caller,
    schema_changes: &[String],
    reason: String,
) -> (StatusCode, String) {
    if !schema_changes.is_empty() {
        warn!(
            job_id = %job.id,
            ?schema_changes,
            error = %reason,
            "job_upsert: catalog write failed after explode schema migration already committed",
        );
        audit::record(
            &s.nats,
            "operator",
            "job_upsert_failed_after_schema_migration",
            Some(&job.id),
            Some(caller),
            serde_json::json!({
                "error": reason,
                "schema_changes": schema_changes,
            }),
        )
        .await;
    }
    (
        StatusCode::INTERNAL_SERVER_ERROR,
        catalog_write_failure_message(&reason, schema_changes),
    )
}

/// Pure message-building half of [`catalog_write_failed`], split out
/// so it's testable without a live NATS client (audit::record needs
/// one; this doesn't). `schema_changes` empty ⇒ the plain error,
/// unchanged from before this fix.
fn catalog_write_failure_message(reason: &str, schema_changes: &[String]) -> String {
    if schema_changes.is_empty() {
        return reason.to_string();
    }
    format!(
        "{reason}; NOTE: this job's inventory.explode derived table schema was already \
         migrated before this failure and was NOT rolled back ({}) — the manifest was NOT \
         saved, so the catalog and the derived table now disagree on that table's shape \
         until this job is created successfully or the table receives another exec result",
        schema_changes.join("; ")
    )
}

/// POST /api/jobs — upsert a Manifest into the job catalog. The KV
/// key is `manifest.id`.
///
/// Accepts JSON (`application/json`, default) or YAML
/// (`application/yaml`, `text/yaml`). When the body is YAML, the raw
/// source is mirrored verbatim into `BUCKET_JOBS_YAML` so the SPA's
/// YAML editor preserves operator comments + script block-scalar
/// indentation across edits. JSON callers get a `serde_yaml::to_string`
/// fallback so the YAML bucket stays in lockstep — the operator just
/// loses comment fidelity on that path (it has nothing to preserve).
pub async fn create(
    State(s): State<AppState>,
    caller: Caller,
    body: YamlOrJson<Manifest>,
) -> Result<Json<JobSummary>, (StatusCode, String)> {
    let YamlOrJson {
        value: job,
        raw_yaml,
    } = body;
    // SPEC §2.4.1: exactly one of script / script_file /
    // script_object must be set. Enforce at the write boundary
    // so the JOBS KV never stores ambiguous manifests.
    job.validate().map_err(|e| (StatusCode::BAD_REQUEST, e))?;
    // #918: `script_file` is CLI-side sugar — `kanade job create`
    // reads the repo-local file and inlines it into `script` before
    // POSTing, so a manifest that reaches this handler still carrying
    // `script_file` came in through the raw API / SPA editor, where no
    // filesystem resolution ever happens. Storing it would 400 at
    // every exec (`resolve_script_source`), so reject it here with an
    // actionable message rather than accepting a job that can never
    // run.
    if job.execute.script_file.is_some() {
        return Err((
            StatusCode::BAD_REQUEST,
            "execute.script_file is resolved CLI-side by `kanade job create` and cannot be \
             stored via the API — inline the script into `execute.script`, or publish it as a \
             `script_object`"
                .to_string(),
        ));
    }
    // #1492: reconcile every `inventory.explode` derived table's SQL
    // schema against this manifest BEFORE the KV write below. Without
    // this, a manifest edit that changed `primary_key` / `columns`
    // stored cleanly (this handler never touched the derived table at
    // all — only the startup pass and the per-result hot path did, and
    // both used a plain `CREATE TABLE IF NOT EXISTS` that's a no-op
    // once the table exists), so the derived table kept its original
    // schema forever while `job create` / `job validate` / `exec` all
    // kept reporting success. Running the reconcile here, and failing
    // the request (not writing the manifest) if it errors, means a job
    // catalog entry is never left pointing at a derived table we know
    // is out of sync with it.
    //
    // `ensure_tables_atomic` runs every spec inside ONE transaction
    // (review R1-2-1): a manifest with several explode specs where an
    // earlier spec's migration succeeds but a later one fails must not
    // leave the earlier table already mutated while this handler
    // reports failure — that would both lose the earlier migration's
    // own report (discarded along with the 500 response) and leave the
    // DB and the (unwritten) catalog disagreeing about that table's
    // shape. A single transaction means a failure anywhere rolls every
    // spec in this call back together, so schema_changes below is only
    // ever populated with changes that actually landed.
    let mut schema_changes = Vec::new();
    if let Some(specs) = job.inventory.as_ref().and_then(|inv| inv.explode.as_ref()) {
        let changes = crate::projector::explode::ensure_tables_atomic(&s.pool, specs)
            .await
            .map_err(|e| {
                (
                    StatusCode::INTERNAL_SERVER_ERROR,
                    format!("inventory.explode schema migration failed: {e:#}"),
                )
            })?;
        for (spec, change) in specs.iter().zip(changes) {
            if !change.is_notable() {
                continue;
            }
            let mut note = if change.rebuilt {
                format!(
                    "{}: rebuilt (primary_key/columns changed); copied {} row(s)",
                    spec.table, change.rows_copied
                )
            } else {
                format!(
                    "{}: added column(s) {}",
                    spec.table,
                    change.added_columns.join(", ")
                )
            };
            if change.rows_lost > 0 {
                note.push_str(&format!(
                    ", lost {} row(s) to the new primary key (will reappear on next exec)",
                    change.rows_lost
                ));
            }
            warn!(job_id = %job.id, table = %spec.table, %note, "explode: derived table schema migrated");
            schema_changes.push(note);
        }
    }

    // #1492 review (R2-1-1): everything from here on writes to the KV
    // catalog, which is a different transactional system than the
    // SQLite migration `ensure_tables_atomic` already committed above
    // — there's no 2-phase commit spanning NATS KV and SQLite, so a
    // failure in any of these steps can't be made to un-migrate the
    // derived table. `catalog_write_failed` makes that unavoidable gap
    // visible instead of silent: it folds the already-committed
    // `schema_changes` into both the error response and an audit
    // record, mirroring `delete`'s `job_delete_failed_post_revoke`
    // event for the same "committed side effect, primary write still
    // failed" shape.
    let kv = match s
        .jetstream
        .create_key_value(KvConfig {
            bucket: BUCKET_JOBS.into(),
            history: 5,
            ..Default::default()
        })
        .await
    {
        Ok(kv) => kv,
        Err(e) => {
            return Err(catalog_write_failed(
                &s,
                &job,
                &caller,
                &schema_changes,
                format!("ensure KV: {e}"),
            )
            .await);
        }
    };
    let body_bytes = match serde_json::to_vec(&job) {
        Ok(b) => b,
        Err(e) => {
            return Err(catalog_write_failed(
                &s,
                &job,
                &caller,
                &schema_changes,
                format!("serialize: {e}"),
            )
            .await);
        }
    };
    if let Err(e) = kv.put(&job.id, body_bytes.into()).await {
        return Err(catalog_write_failed(
            &s,
            &job,
            &caller,
            &schema_changes,
            format!("KV put: {e}"),
        )
        .await);
    }

    // Keep `script_current` in lockstep with the catalog on every upsert,
    // NOT just on `kanade exec` (exec.rs #258). `script_current.<id>` is
    // the canonical "current version" the agent's Layer 2 version-pin gate
    // (§2.6.2) checks against. Previously only the exec / backend-fire path
    // wrote it, so a `runs_on: agent` job — which fires from its own
    // BUCKET_JOBS cache and never routes through exec_manifest — left the
    // KV frozen at whatever version was last exec'd. After a version bump
    // the KV lagged the catalog until someone ran a manual exec, and the
    // agent's own fresh fires got self-rejected (exit 124, version-pin
    // mismatch). Writing it here means the pin baseline tracks the catalog
    // regardless of `runs_on`. Best-effort: a failure must not roll back
    // the catalog write above (the agent-local path no longer depends on
    // this KV — see CommandSource::LocalScheduler — and the backend-fire
    // path re-pins it in exec_manifest anyway), so warn and continue.
    match s.jetstream.get_key_value(BUCKET_SCRIPT_CURRENT).await {
        Ok(cur) => {
            if let Err(e) = cur
                .put(
                    &job.id,
                    bytes::Bytes::from(job.version.clone().into_bytes()),
                )
                .await
            {
                warn!(error = %e, job_id = %job.id, "jobs: script_current put failed; catalog is current");
            }
        }
        Err(e) => {
            warn!(error = %e, "jobs: script_current KV missing; skipping version-pin refresh")
        }
    }

    // Operator-facing YAML mirror — best-effort. A failure here doesn't
    // roll back the JSON catalog (which the scheduler / agents read)
    // because the user's update is functionally already in place; we
    // just may lose a round of comment preservation. Warn-log so the
    // gap is observable.
    let yaml_source = raw_yaml.unwrap_or_else(|| {
        serde_yaml::to_string(&job)
            .unwrap_or_else(|_| String::from("# YAML mirror unavailable for this entry"))
    });
    if let Err(e) = mirror_yaml(&s, BUCKET_JOBS_YAML, &job.id, &yaml_source).await {
        warn!(error = %e, job_id = %job.id, "jobs: YAML mirror put failed; JSON catalog is current");
    }

    let summary = JobSummary {
        id: job.id.clone(),
        version: job.version.clone(),
        description: job.description.clone(),
        inventory: job.inventory.is_some(),
        schema_changes,
    };
    info!(job_id = %job.id, version = %job.version, "job upserted");
    audit::record(
        &s.nats,
        "operator",
        "job_upsert",
        Some(&job.id),
        Some(&caller),
        serde_json::json!({
            "version": job.version,
            "inventory": job.inventory.is_some(),
        }),
    )
    .await;
    Ok(Json(summary))
}

/// `GET /api/jobs/{id}/yaml` — fetch the operator's YAML source for
/// the job, used by the SPA editor to populate Edit modals without
/// losing comments / block-scalar formatting. Falls back to a
/// `serde_yaml::to_string` dump of the JSON catalog row when the
/// YAML mirror is missing (legacy entries from before this endpoint).
pub async fn get_yaml(
    State(s): State<AppState>,
    Path(id): Path<String>,
) -> Result<(StatusCode, HeaderMap, String), (StatusCode, String)> {
    if let Ok(kv) = s.jetstream.get_key_value(BUCKET_JOBS_YAML).await
        && let Ok(Some(bytes)) = kv.get(&id).await
        && let Ok(yaml_str) = String::from_utf8(bytes.to_vec())
    {
        return Ok((StatusCode::OK, yaml_headers(), yaml_str));
    }

    let kv = s
        .jetstream
        .get_key_value(BUCKET_JOBS)
        .await
        .map_err(|_| (StatusCode::NOT_FOUND, format!("job '{id}' not found")))?;
    let bytes = kv
        .get(&id)
        .await
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("KV get: {e}")))?
        .ok_or_else(|| (StatusCode::NOT_FOUND, format!("job '{id}' not found")))?;
    let manifest: Manifest = serde_json::from_slice(&bytes)
        .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, format!("decode: {e}")))?;
    let yaml = serde_yaml::to_string(&manifest).map_err(|e| {
        (
            StatusCode::INTERNAL_SERVER_ERROR,
            format!("encode YAML: {e}"),
        )
    })?;
    Ok((StatusCode::OK, yaml_headers(), yaml))
}

/// DELETE /api/jobs/{id} — 409 if any Schedule references it.
///
/// v0.27 (SPEC §2.6.4 (b)) cascades a Layer 2 revoke: before the
/// Manifest is removed from `BUCKET_JOBS`, the handler writes
/// `script_status.{id} = REVOKED` so any Command already in flight
/// (live core sub delivery in progress, or stored in `STREAM_EXEC`
/// awaiting a reconnecting agent) gets skipped by the agent's
/// `handle_command` KV check. Without this, deleting a Manifest only
/// stops *future* exec calls — already-published Commands would
/// still run.
///
/// To undo the cascade: re-create the Manifest with
/// `kanade job create`, then `kanade unrevoke <id>` to flip
/// `script_status` back to `ACTIVE`.
const CASCADE_FACTS_SQL: &str = "DELETE FROM inventory_facts WHERE job_id = ?";
const CASCADE_HISTORY_SQL: &str = "DELETE FROM inventory_history WHERE job_id = ?";

/// Drop a job's inventory rows (facts + history), keyed by `job_id`.
/// Both deletes run in one transaction so the two tables can't end up
/// half-cleared if the second fails. Returns `(facts_deleted,
/// history_deleted)`. The `delete` handler runs this AND the
/// `delete_cascades_inventory_for_that_job_only` test calls it directly,
/// so the test exercises the real SQL (a table/column typo fails the
/// test rather than silently passing reimplemented SQL).
async fn cascade_inventory(pool: &SqlitePool, job_id: &str) -> sqlx::Result<(u64, u64)> {
    let mut tx = pool.begin().await?;
    let facts = sqlx::query(CASCADE_FACTS_SQL)
        .bind(job_id)
        .execute(&mut *tx)
        .await?
        .rows_affected();
    let history = sqlx::query(CASCADE_HISTORY_SQL)
        .bind(job_id)
        .execute(&mut *tx)
        .await?
        .rows_affected();
    tx.commit().await?;
    Ok((facts, history))
}

pub async fn delete(
    State(s): State<AppState>,
    Path(id): Path<String>,
    caller: Caller,
) -> Result<StatusCode, (StatusCode, String)> {
    if let Ok(kv) = s.jetstream.get_key_value(BUCKET_SCHEDULES).await
        && let Ok(keys_stream) = kv.keys().await
    {
        let keys: Vec<String> = keys_stream.try_collect().await.unwrap_or_default();
        for k in keys {
            if let Ok(Some(bytes)) = kv.get(&k).await
                && let Ok(sched) = serde_json::from_slice::<Schedule>(&bytes)
                && sched.job_id == id
            {
                return Err((
                    StatusCode::CONFLICT,
                    format!(
                        "job '{id}' is referenced by schedule '{}'; remove the schedule first",
                        sched.id
                    ),
                ));
            }
        }
    }

    // v0.27 — SPEC §2.6.4 (b) cascade revoke: every job delete also
    // writes `script_status.{cmd_id} = REVOKED` so any in-flight
    // Command for this manifest (publish-already-emitted but the
    // agent hasn't run yet, or about to be replayed from STREAM_EXEC
    // on reconnect) gets caught by the Layer 2 KV check and skipped.
    // Without this, removing a Manifest only stops *future* exec
    // calls — Commands already in the broker would still execute on
    // any agent that reads them. We revoke FIRST, then delete the
    // Manifest, so that if delete somehow fails we're still in a safe
    // (revoked) state. Idempotent — re-revoking an already-REVOKED
    // entry is a no-op put. v0.27 round-2 review (gemini #36
    // line 208): resolve BOTH KV handles upfront before any write —
    // that way a missing / unreachable BUCKET_JOBS surfaces as a
    // clean 404 with zero side effects, instead of leaking a revoke
    // that has no matching delete.
    let status_kv = s
        .jetstream
        .get_key_value(BUCKET_SCRIPT_STATUS)
        .await
        .map_err(|e| {
            warn!(
                error = %e,
                bucket = BUCKET_SCRIPT_STATUS,
                "job_delete cascade revoke: status KV unavailable",
            );
            (
                StatusCode::SERVICE_UNAVAILABLE,
                format!("script_status bucket missing: {e}"),
            )
        })?;
    let kv = s.jetstream.get_key_value(BUCKET_JOBS).await.map_err(|e| {
        warn!(error = %e, "jobs KV missing on delete");
        (StatusCode::NOT_FOUND, "jobs bucket missing".to_string())
    })?;

    status_kv
        .put(&id, bytes::Bytes::from(SCRIPT_STATUS_REVOKED))
        .await
        .map_err(|e| {
            warn!(
                error = %e,
                job_id = %id,
                bucket = BUCKET_SCRIPT_STATUS,
                "job_delete cascade revoke: status KV put failed",
            );
            (
                StatusCode::INTERNAL_SERVER_ERROR,
                format!("script_status put: {e}"),
            )
        })?;
    // If the manifest delete fails *after* we successfully cascaded
    // the revoke, the operator needs to know that `script_status.{id}`
    // is now REVOKED so they can `kanade unrevoke <id>` as part of the
    // recovery. We audit + surface the revoke state in the error body
    // rather than silently dropping it (CodeRabbit #36 review).
    if let Err(e) = kv.delete(&id).await {
        warn!(
            error = %e,
            job_id = %id,
            cascade_revoke = true,
            "job_delete failed after cascade revoke",
        );
        audit::record(
            &s.nats,
            "operator",
            "job_delete_failed_post_revoke",
            Some(&id),
            Some(&caller),
            serde_json::json!({ "cascade_revoke": true, "error": e.to_string() }),
        )
        .await;
        return Err((
            StatusCode::INTERNAL_SERVER_ERROR,
            format!(
                "kv delete: {e}; script_status.{id} is already REVOKED — `kanade unrevoke {id}` to recover"
            ),
        ));
    }
    // Cascade inventory cleanup. inventory_facts / inventory_history are
    // keyed by job_id, so this job's rows are unambiguously its own —
    // unlike check_status (keyed by check slug, deliberately NOT cascaded
    // so a same-named replacement job keeps its rows). Without this a
    // deleted inventory job leaves orphaned facts stuck under a dead
    // job_id on the SPA — there is no other cleanup path for
    // inventory_facts (inventory_history additionally has a 90d TTL).
    //
    // Best-effort: the manifest is already gone from KV, so a failed
    // sweep must not fail the request — that would only mean orphans
    // remain (the pre-cascade status quo). Surface it via warn + audit
    // rather than a misleading 500 for an already-deleted job.
    let (facts_deleted, history_deleted, cascade_error) =
        match cascade_inventory(&s.pool, &id).await {
            Ok((f, h)) => (f, h, None),
            Err(e) => {
                warn!(error = %e, job_id = %id, "job_delete inventory cascade failed");
                (0, 0, Some(e.to_string()))
            }
        };

    info!(
        job_id = %id,
        cascade_revoke = true,
        inventory_facts_deleted = facts_deleted,
        inventory_history_deleted = history_deleted,
        cascade_ok = cascade_error.is_none(),
        "job deleted",
    );
    audit::record(
        &s.nats,
        "operator",
        "job_delete",
        Some(&id),
        Some(&caller),
        serde_json::json!({
            "cascade_revoke": true,
            "inventory_facts_deleted": facts_deleted,
            "inventory_history_deleted": history_deleted,
            // null on success; an error string distinguishes a failed
            // sweep (orphans may remain) from a genuine "0 rows matched".
            "inventory_cascade_error": cascade_error,
        }),
    )
    .await;
    Ok(StatusCode::NO_CONTENT)
}

/// Lookup helper for scheduler + projector. Returns `Ok(None)` when
/// the key is absent so callers can warn-and-skip without unwrapping
/// a fatal error.
pub async fn fetch(
    js: &async_nats::jetstream::Context,
    job_id: &str,
) -> anyhow::Result<Option<Manifest>> {
    let kv = match js.get_key_value(BUCKET_JOBS).await {
        Ok(k) => k,
        Err(_) => return Ok(None),
    };
    let Some(bytes) = kv.get(job_id).await? else {
        return Ok(None);
    };
    let job: Manifest = serde_json::from_slice(&bytes)?;
    Ok(Some(job))
}

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

    async fn fresh_pool() -> SqlitePool {
        let pool = SqlitePoolOptions::new()
            .max_connections(1)
            .connect("sqlite::memory:")
            .await
            .expect("open sqlite memory");
        sqlx::migrate!("./migrations")
            .run(&pool)
            .await
            .expect("run migrations");
        pool
    }

    async fn insert_exec(
        pool: &SqlitePool,
        exec_id: &str,
        job_id: &str,
        status: &str,
        target: i64,
    ) {
        sqlx::query(
            "INSERT INTO executions
                (exec_id, job_id, version, initiated_by, target_count, status)
             VALUES (?, ?, '1.0.0', 'tester', ?, ?)",
        )
        .bind(exec_id)
        .bind(job_id)
        .bind(target)
        .bind(status)
        .execute(pool)
        .await
        .unwrap();
    }

    #[tokio::test]
    async fn fetch_live_counts_groups_by_job_id() {
        // Three execs for "inv-hw" (2 running, 1 pending), one exec for
        // "patch-x" (just running). The aggregation should produce a
        // map keyed by job_id with the right partition.
        let pool = fresh_pool().await;
        insert_exec(&pool, "e1", "inv-hw", "running", 10).await;
        insert_exec(&pool, "e2", "inv-hw", "running", 10).await;
        insert_exec(&pool, "e3", "inv-hw", "pending", 10).await;
        insert_exec(&pool, "e4", "patch-x", "running", 5).await;

        let counts = fetch_live_counts(&pool).await.unwrap();
        assert_eq!(
            counts.get("inv-hw"),
            Some(&JobLiveCounts {
                running: 2,
                pending: 1,
            }),
        );
        assert_eq!(
            counts.get("patch-x"),
            Some(&JobLiveCounts {
                running: 1,
                pending: 0,
            }),
        );
    }

    #[tokio::test]
    async fn fetch_live_counts_excludes_completed() {
        // 'completed' status (= projector saw all target_count
        // results) shouldn't count toward live. Only 'running' and
        // 'pending' are operationally "in flight".
        let pool = fresh_pool().await;
        insert_exec(&pool, "e1", "j", "completed", 5).await;
        insert_exec(&pool, "e2", "j", "running", 5).await;

        let counts = fetch_live_counts(&pool).await.unwrap();
        let live = counts.get("j").expect("j has at least one exec");
        assert_eq!(live.running, 1);
        assert_eq!(live.pending, 0);
    }

    #[tokio::test]
    async fn fetch_live_counts_empty_when_no_executions() {
        let pool = fresh_pool().await;
        let counts = fetch_live_counts(&pool).await.unwrap();
        assert!(counts.is_empty());
    }

    async fn insert_inv_fact(pool: &SqlitePool, pc: &str, job: &str) {
        sqlx::query("INSERT INTO inventory_facts (pc_id, job_id, facts_json) VALUES (?, ?, '{}')")
            .bind(pc)
            .bind(job)
            .execute(pool)
            .await
            .unwrap();
    }

    async fn insert_inv_hist(pool: &SqlitePool, pc: &str, job: &str) {
        sqlx::query(
            "INSERT INTO inventory_history (pc_id, job_id, field_path, change_kind)
             VALUES (?, ?, 'x', 'added')",
        )
        .bind(pc)
        .bind(job)
        .execute(pool)
        .await
        .unwrap();
    }

    // Calls the real `cascade_inventory` the handler uses (not
    // reimplemented SQL), so a table/column typo fails here: dropping a
    // job's inventory must hit only that job_id, leaving other jobs' rows.
    #[tokio::test]
    async fn delete_cascades_inventory_for_that_job_only() {
        let pool = fresh_pool().await;
        insert_inv_fact(&pool, "pc-1", "inv-hw").await;
        insert_inv_fact(&pool, "pc-2", "inv-hw").await;
        insert_inv_fact(&pool, "pc-1", "inv-sw").await; // a different job
        insert_inv_hist(&pool, "pc-1", "inv-hw").await;
        insert_inv_hist(&pool, "pc-1", "inv-sw").await;

        let (facts, hist) = cascade_inventory(&pool, "inv-hw").await.unwrap();
        assert_eq!(facts, 2, "both inv-hw fact rows cleared");
        assert_eq!(hist, 1, "inv-hw history row cleared");

        // inv-sw is untouched.
        let facts_left: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM inventory_facts")
            .fetch_one(&pool)
            .await
            .unwrap();
        let hist_left: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM inventory_history")
            .fetch_one(&pool)
            .await
            .unwrap();
        assert_eq!(facts_left, 1, "the other job's fact survives");
        assert_eq!(hist_left, 1, "the other job's history survives");
    }

    /// #1214①: the list response must NOT carry the script body (it
    /// was ~85 % of the payload) while keeping the flattened shape
    /// the SPA's `JobRow` reads.
    #[test]
    fn job_list_row_omits_script_but_keeps_catalog_shape() {
        let m: Manifest = serde_yaml::from_str(
            "id: inv-hw\n\
             version: 1.2.3\n\
             description: hardware inventory\n\
             execute:\n\
             \x20 shell: powershell\n\
             \x20 timeout: 30s\n\
             \x20 run_as: system\n\
             \x20 cwd: C:\\Temp\n\
             \x20 script: |\n\
             \x20   Get-ComputerInfo | Out-Null\n\
             tags: [inventory, windows]\n",
        )
        .unwrap();
        let row = JobListRow::from_manifest(
            m,
            JobLiveCounts {
                running: 2,
                pending: 1,
            },
        );
        let v = serde_json::to_value(&row).unwrap();

        assert_eq!(v["id"], "inv-hw");
        assert_eq!(v["version"], "1.2.3");
        assert_eq!(v["description"], "hardware inventory");
        assert_eq!(v["execute"]["shell"], "powershell");
        assert_eq!(v["execute"]["timeout"], "30s");
        assert_eq!(v["execute"]["run_as"], "system");
        assert_eq!(v["execute"]["cwd"], "C:\\Temp");
        assert_eq!(v["tags"], serde_json::json!(["inventory", "windows"]));
        assert_eq!(v["live"], serde_json::json!({"running": 2, "pending": 1}));
        // Serialized as null (not omitted) — the SPA types it as
        // `unknown | null`.
        assert!(v.get("inventory").is_some());
        assert!(v["inventory"].is_null());

        // The whole point of #1214①: no script source in the list row.
        assert!(v["execute"].get("script").is_none());
        assert!(v["execute"].get("script_file").is_none());
        assert!(v["execute"].get("script_object").is_none());
    }

    /// Review R2-1-1: when the KV catalog write fails with no explode
    /// migration in play (the common case), the response is exactly
    /// the underlying error — no change from before this fix.
    #[test]
    fn catalog_write_failure_message_is_plain_reason_when_no_schema_changes() {
        let msg = catalog_write_failure_message("KV put: connection refused", &[]);
        assert_eq!(msg, "KV put: connection refused");
    }

    /// Review R2-1-1: when `ensure_tables_atomic` already committed a
    /// migration before the KV write failed, that fact — including any
    /// row loss from a rebuild — must be readable in the error body
    /// the CLI prints, not just in a server-side log line. Before this
    /// fix the handler returned bare `format!("KV put: {e}")`,
    /// discarding the already-computed `schema_changes` entirely.
    #[test]
    fn catalog_write_failure_message_surfaces_committed_schema_changes() {
        let changes = vec![
            "example_items: rebuilt (primary_key/columns changed); copied 3 row(s), lost 1 \
             row(s) to the new primary key (will reappear on next exec)"
                .to_string(),
        ];
        let msg = catalog_write_failure_message("KV put: timeout", &changes);
        assert!(msg.starts_with("KV put: timeout"), "{msg}");
        assert!(
            msg.contains("example_items: rebuilt"),
            "migration note must appear in the response body: {msg}"
        );
        assert!(
            msg.contains("lost 1 row(s)"),
            "row loss must be visible in the response body, not just logs: {msg}"
        );
        assert!(
            msg.contains("was NOT saved"),
            "must make clear the catalog write itself failed despite the migration landing: {msg}"
        );
    }
}