Skip to main content

sloop/
store.rs

1use std::collections::{BTreeMap, BTreeSet, HashMap};
2use std::fmt;
3use std::path::{Path, PathBuf};
4
5use rusqlite::{Connection, OptionalExtension, TransactionBehavior, params};
6
7use crate::domain::ticket::TicketState;
8
9pub const SCHEMA_VERSION: u32 = 13;
10
11const CONNECTION_PRAGMAS: &str = "
12PRAGMA foreign_keys = ON;
13PRAGMA journal_mode = WAL;
14PRAGMA busy_timeout = 5000;
15";
16
17const SCHEMA_V1: &str = "
18CREATE TABLE projects (
19    id              TEXT PRIMARY KEY,
20    file_path       TEXT UNIQUE,
21    source          TEXT NOT NULL DEFAULT 'local',
22    source_ref      TEXT,
23    title           TEXT NOT NULL,
24    created_at_ms   INTEGER NOT NULL,
25    updated_at_ms   INTEGER NOT NULL,
26    UNIQUE (source, source_ref),
27    CHECK (file_path IS NOT NULL OR source_ref IS NOT NULL)
28);
29
30CREATE TABLE tickets (
31    id              TEXT PRIMARY KEY,
32    project_id      TEXT NOT NULL REFERENCES projects(id),
33    file_path       TEXT UNIQUE,
34    source          TEXT NOT NULL DEFAULT 'local',
35    source_ref      TEXT,
36    state           TEXT NOT NULL,
37    attempts        INTEGER NOT NULL DEFAULT 0,
38    content_hash    TEXT,
39    name            TEXT NOT NULL DEFAULT '',
40    worktree        TEXT,
41    target          TEXT,
42    model           TEXT,
43    effort          TEXT,
44    flow            TEXT,
45    body            TEXT,
46    held_reason     TEXT,
47    missing_at_ms   INTEGER,
48    created_at_ms   INTEGER NOT NULL,
49    updated_at_ms   INTEGER NOT NULL,
50    UNIQUE (source, source_ref),
51    CHECK (file_path IS NOT NULL OR source_ref IS NOT NULL)
52);
53
54CREATE INDEX tickets_by_project_state
55ON tickets(project_id, state);
56
57-- Dependencies are normalized so references are foreign-key checked and
58-- graph reads do not require decoding serialized ticket data.
59CREATE TABLE ticket_blockers (
60    ticket_id       TEXT NOT NULL REFERENCES tickets(id) ON DELETE CASCADE,
61    blocker_id      TEXT NOT NULL REFERENCES tickets(id),
62    position        INTEGER NOT NULL,
63    PRIMARY KEY (ticket_id, blocker_id)
64);
65
66CREATE TABLE activations (
67    id              TEXT PRIMARY KEY,
68    kind            TEXT NOT NULL,
69    state           TEXT NOT NULL,
70    ticket_id       TEXT REFERENCES tickets(id),
71    project_id      TEXT REFERENCES projects(id),
72    eligible_at_ms  INTEGER,
73    interval_ms     INTEGER,
74    created_at_ms   INTEGER NOT NULL,
75    updated_at_ms   INTEGER NOT NULL,
76    CHECK (ticket_id IS NULL OR project_id IS NULL)
77);
78
79CREATE TABLE activation_filters (
80    activation_id   TEXT NOT NULL REFERENCES activations(id) ON DELETE CASCADE,
81    ticket_id       TEXT NOT NULL REFERENCES tickets(id),
82    PRIMARY KEY (activation_id, ticket_id)
83);
84
85CREATE TABLE runs (
86    id                    TEXT PRIMARY KEY,
87    activation_id         TEXT NOT NULL REFERENCES activations(id),
88    ticket_id             TEXT NOT NULL REFERENCES tickets(id),
89    state                 TEXT NOT NULL,
90    attempt               INTEGER NOT NULL,
91    branch                TEXT,
92    worktree_path         TEXT,
93    pid                   INTEGER,
94    pid_start_time        INTEGER,
95    process_group_id      INTEGER,
96    worker_token          TEXT,
97    worker_socket_path    TEXT,
98    started_at_ms         INTEGER,
99    exited_at_ms          INTEGER,
100    exit_code             INTEGER,
101    cleanup_eligible_at_ms INTEGER,
102    cleaned_at_ms         INTEGER,
103    flow_json             TEXT,
104    ticket_json           TEXT,
105    created_at_ms         INTEGER NOT NULL,
106    updated_at_ms         INTEGER NOT NULL
107);
108
109CREATE INDEX runs_by_ticket ON runs(ticket_id, created_at_ms);
110CREATE INDEX runs_by_activation ON runs(activation_id, created_at_ms);
111
112-- A lease is time-bounded ownership of a ticket by the daemon, taken
113-- atomically at claim time. `ticket_id` is the PRIMARY KEY and `run_id` is
114-- UNIQUE, so the engine itself enforces at most one lease per ticket and per
115-- run: the durable guard against double-spawn, backstopping the conditional
116-- `UPDATE ... WHERE state='ready'` in `claim_ticket`.
117--
118-- Leases are held only by the daemon; `owner_id` records which daemon process
119-- took the claim. Workers never hold, renew, or observe leases — a worker's
120-- only credential is a per-run capability token granting the worker verbs on
121-- its own run.
122--
123-- `expires_at_ms` gates renewal only: an expired lease cannot be renewed, so a
124-- revived process cannot resurrect a claim recovery has decided is lost.
125-- Liveness of a run is determined by process identity (pid + pid start time +
126-- process group id), never by lease expiry.
127--
128-- The daemon renews the lease of every run it supervises, so `expires_at_ms`
129-- stays in the future for as long as a run is alive and an expired row means
130-- nobody was there to renew it. Because renewal is strict, a daemon returning
131-- after longer than the TTL re-arms a readopted run's lapsed lease through
132-- `readopt_lease` rather than through renewal.
133--
134-- A lease is released by deleting its row: on settlement (`finish_run`) or on
135-- claim rollback (`abort_claim`). An expired-but-present row is evidence of an
136-- owner that died mid-work.
137CREATE TABLE leases (
138    ticket_id       TEXT PRIMARY KEY REFERENCES tickets(id),
139    run_id          TEXT NOT NULL UNIQUE REFERENCES runs(id),
140    owner_id        TEXT NOT NULL,
141    acquired_at_ms  INTEGER NOT NULL,
142    renewed_at_ms   INTEGER NOT NULL,
143    expires_at_ms   INTEGER NOT NULL
144);
145
146CREATE INDEX leases_by_expiry ON leases(expires_at_ms);
147
148CREATE TABLE run_evidence (
149    sequence        INTEGER PRIMARY KEY AUTOINCREMENT,
150    run_id          TEXT NOT NULL REFERENCES runs(id),
151    kind            TEXT NOT NULL,
152    observed_at_ms  INTEGER NOT NULL,
153    dedupe_key      TEXT UNIQUE,
154    data_json       TEXT NOT NULL
155);
156
157CREATE INDEX evidence_by_run ON run_evidence(run_id, sequence);
158
159CREATE TABLE aftercare_stages (
160    run_id          TEXT NOT NULL REFERENCES runs(id),
161    stage_index     INTEGER NOT NULL,
162    stage           TEXT NOT NULL,
163    state           TEXT NOT NULL,
164    attempt         INTEGER NOT NULL DEFAULT 1,
165    started_at_ms   INTEGER,
166    finished_at_ms  INTEGER,
167    exit_code       INTEGER,
168    evidence_json   TEXT,
169    PRIMARY KEY (run_id, stage_index, attempt)
170);
171
172CREATE TABLE cooldowns (
173    key             TEXT PRIMARY KEY,
174    until_ms        INTEGER NOT NULL,
175    reason          TEXT NOT NULL,
176    source_run_id   TEXT REFERENCES runs(id),
177    updated_at_ms   INTEGER NOT NULL
178);
179
180CREATE TABLE budget_reservations (
181    run_id              TEXT PRIMARY KEY REFERENCES runs(id),
182    reserved_tokens     INTEGER NOT NULL,
183    actual_tokens       INTEGER,
184    state               TEXT NOT NULL,
185    created_at_ms       INTEGER NOT NULL,
186    reconciled_at_ms    INTEGER
187);
188
189CREATE TABLE scheduler_state (
190    singleton       INTEGER PRIMARY KEY CHECK (singleton = 1),
191    paused          INTEGER NOT NULL CHECK (paused IN (0, 1)),
192    draining        INTEGER NOT NULL DEFAULT 0 CHECK (draining IN (0, 1)),
193    updated_at_ms   INTEGER NOT NULL
194);
195
196CREATE TABLE notes (
197    id              TEXT PRIMARY KEY,
198    run_id          TEXT NOT NULL REFERENCES runs(id),
199    text            TEXT NOT NULL,
200    recorded_at_ms  INTEGER NOT NULL
201);
202";
203
204const ID_COUNTER_SCHEMA: &str = "
205CREATE TABLE IF NOT EXISTS id_counters (
206    kind            TEXT PRIMARY KEY,
207    next_ordinal    INTEGER NOT NULL CHECK (next_ordinal > 0)
208);
209INSERT OR IGNORE INTO id_counters (kind, next_ordinal)
210SELECT 'activation', COALESCE(MAX(CAST(SUBSTR(id, 2) AS INTEGER)), 0) + 1 FROM activations;
211INSERT OR IGNORE INTO id_counters (kind, next_ordinal)
212SELECT 'note', COALESCE(MAX(CAST(SUBSTR(id, 2) AS INTEGER)), 0) + 1 FROM notes;
213";
214
215const RUN_SNAPSHOT_COLUMNS: &str = "
216ALTER TABLE runs ADD COLUMN flow_json TEXT;
217ALTER TABLE runs ADD COLUMN ticket_json TEXT;
218";
219
220// The activity feed read by `sloop watch`. Rows are written inside the same
221// transaction as the state transition they describe, so the feed can never
222// disagree with the tables it narrates.
223const EVENTS_SCHEMA: &str = "
224CREATE TABLE IF NOT EXISTS events (
225    sequence        INTEGER PRIMARY KEY AUTOINCREMENT,
226    occurred_at_ms  INTEGER NOT NULL,
227    kind            TEXT NOT NULL,
228    run_id          TEXT,
229    ticket_id       TEXT,
230    data_json       TEXT NOT NULL DEFAULT '{}'
231);
232";
233
234const TICKET_SOURCE_COLUMNS: &str = "
235ALTER TABLE tickets ADD COLUMN body TEXT;
236ALTER TABLE tickets ADD COLUMN held_reason TEXT;
237";
238
239const RESTART_DRAINING_COLUMN: &str = "
240ALTER TABLE scheduler_state ADD COLUMN draining INTEGER NOT NULL DEFAULT 0
241CHECK (draining IN (0, 1));
242";
243
244const WORKTREE_CLEANUP_COLUMNS: &str = "
245ALTER TABLE runs ADD COLUMN cleanup_eligible_at_ms INTEGER;
246ALTER TABLE runs ADD COLUMN cleaned_at_ms INTEGER;
247";
248
249/// Every value the `runs.state` column can hold. The ladder runs
250/// `claimed → running → aftercare` and then to one terminal state, either an
251/// outcome written by [`Store::finish_run`] or `aborted` from a rolled-back
252/// claim. Values are the exact strings already stored; there is no migration.
253#[derive(Debug, Clone, Copy, PartialEq, Eq)]
254pub enum RunState {
255    Claimed,
256    Running,
257    Aftercare,
258    /// A claim rolled back before a process existed.
259    Aborted,
260    Merged,
261    Failed,
262    NeedsReview,
263    Cancelled,
264    RateLimited,
265    Orphaned,
266}
267
268/// The nonterminal run states, in ladder order. A run in one of these still
269/// owns its lease and is a candidate for recovery.
270pub(crate) const NONTERMINAL_RUN_STATES: [RunState; 3] =
271    [RunState::Claimed, RunState::Running, RunState::Aftercare];
272
273impl RunState {
274    pub fn as_str(self) -> &'static str {
275        match self {
276            Self::Claimed => "claimed",
277            Self::Running => "running",
278            Self::Aftercare => "aftercare",
279            Self::Aborted => "aborted",
280            Self::Merged => "merged",
281            Self::Failed => "failed",
282            Self::NeedsReview => "needs_review",
283            Self::Cancelled => "cancelled",
284            Self::RateLimited => "rate_limited",
285            Self::Orphaned => "orphaned",
286        }
287    }
288
289    /// Reads a state written by an older or newer binary. An unrecognized
290    /// value is an error rather than a fallback: silently treating it as
291    /// nonterminal would let the daemon act on a run it cannot classify.
292    pub fn parse(value: &str) -> Result<Self, StoreError> {
293        match value {
294            "claimed" => Ok(Self::Claimed),
295            "running" => Ok(Self::Running),
296            "aftercare" => Ok(Self::Aftercare),
297            "aborted" => Ok(Self::Aborted),
298            "merged" => Ok(Self::Merged),
299            "failed" => Ok(Self::Failed),
300            "needs_review" => Ok(Self::NeedsReview),
301            "cancelled" => Ok(Self::Cancelled),
302            "rate_limited" => Ok(Self::RateLimited),
303            "orphaned" => Ok(Self::Orphaned),
304            other => Err(StoreError::UnknownRunState {
305                state: other.into(),
306            }),
307        }
308    }
309
310    /// Whether the run has stopped: no lease, no supervision, no renewal.
311    pub fn is_terminal(self) -> bool {
312        !NONTERMINAL_RUN_STATES.contains(&self)
313    }
314}
315
316/// Reads `runs.state` as a typed value. An unrecognized string fails the row
317/// rather than defaulting, so a state this binary does not understand can
318/// never be mistaken for a live or a settled run.
319impl rusqlite::types::FromSql for RunState {
320    fn column_result(value: rusqlite::types::ValueRef<'_>) -> rusqlite::types::FromSqlResult<Self> {
321        let text = value.as_str()?;
322        Self::parse(text).map_err(|error| rusqlite::types::FromSqlError::Other(Box::new(error)))
323    }
324}
325
326/// Binds the nonterminal states as `?1, ?2, ?3` for the `IN` clauses that
327/// select live runs.
328fn nonterminal_state_params() -> [&'static str; 3] {
329    [
330        NONTERMINAL_RUN_STATES[0].as_str(),
331        NONTERMINAL_RUN_STATES[1].as_str(),
332        NONTERMINAL_RUN_STATES[2].as_str(),
333    ]
334}
335
336impl From<crate::outcome::Outcome> for RunState {
337    fn from(outcome: crate::outcome::Outcome) -> Self {
338        use crate::outcome::Outcome;
339        match outcome {
340            Outcome::Merged => Self::Merged,
341            Outcome::Failed => Self::Failed,
342            Outcome::NeedsReview => Self::NeedsReview,
343            Outcome::Cancelled => Self::Cancelled,
344            Outcome::RateLimited => Self::RateLimited,
345            Outcome::Orphaned => Self::Orphaned,
346        }
347    }
348}
349
350#[derive(Debug, Clone, Copy, PartialEq, Eq)]
351pub enum ActivationKind {
352    Immediate,
353    Auto,
354    At,
355    Every,
356    Overnight,
357}
358
359impl ActivationKind {
360    pub fn as_str(self) -> &'static str {
361        match self {
362            Self::Immediate => "immediate",
363            Self::Auto => "auto",
364            Self::At => "at",
365            Self::Every => "every",
366            Self::Overnight => "overnight",
367        }
368    }
369}
370
371#[derive(Debug, Clone, Copy, PartialEq, Eq)]
372pub enum ActivationState {
373    Queued,
374    Completed,
375    Cancelled,
376}
377
378impl ActivationState {
379    pub fn as_str(self) -> &'static str {
380        match self {
381            Self::Queued => "queued",
382            Self::Completed => "completed",
383            Self::Cancelled => "cancelled",
384        }
385    }
386}
387
388#[derive(Debug, Clone, PartialEq, Eq)]
389pub struct NewActivation<'a> {
390    pub id: &'a str,
391    pub kind: ActivationKind,
392    pub ticket_id: Option<&'a str>,
393    pub project_id: Option<&'a str>,
394    pub eligible_at_ms: Option<i64>,
395    pub interval_ms: Option<i64>,
396}
397
398#[derive(Debug, Clone, PartialEq, Eq)]
399pub struct ClaimRequest<'a> {
400    pub ticket_id: &'a str,
401    pub run_id: &'a str,
402    pub activation_id: &'a str,
403    pub owner_id: &'a str,
404    pub lease_ms: i64,
405    pub next_activation_eligible_at_ms: Option<i64>,
406    pub flow_json: &'a str,
407    pub ticket_json: &'a str,
408}
409
410#[derive(Debug, Clone, PartialEq, Eq)]
411pub struct ClaimedRun {
412    pub run_id: String,
413    pub attempt: i64,
414    pub lease_expires_at_ms: i64,
415}
416
417#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
418pub struct TicketCounts {
419    pub ready: u64,
420    pub held: u64,
421    pub blocked: u64,
422    pub claimed: u64,
423    pub merged: u64,
424    pub failed: u64,
425    pub needs_review: u64,
426}
427
428/// One appended `run_evidence` row: a kind plus kind-specific JSON facts.
429#[derive(Debug, Clone, PartialEq, Eq)]
430pub struct EvidenceRecord {
431    pub kind: &'static str,
432    pub data_json: String,
433}
434
435/// One row of the activity feed, ordered by `sequence`.
436#[derive(Debug, Clone, PartialEq, Eq)]
437pub struct EventRecord {
438    pub sequence: i64,
439    pub occurred_at_ms: i64,
440    pub kind: String,
441    pub run_id: Option<String>,
442    pub ticket_id: Option<String>,
443    pub data_json: String,
444}
445
446/// One run's wall-clock boundaries, derived from the activity feed. Every
447/// field is optional because a run is observable at each stage of its life:
448/// claimed but not started, started but not finished.
449#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
450pub struct RunTimeline {
451    pub claimed_at_ms: Option<i64>,
452    pub started_at_ms: Option<i64>,
453    pub finished_at_ms: Option<i64>,
454}
455
456/// Appends one activity-feed row. Callers pass the transaction performing the
457/// transition so the event commits or rolls back with it.
458fn record_event(
459    connection: &Connection,
460    now_ms: i64,
461    kind: &str,
462    run_id: Option<&str>,
463    ticket_id: Option<&str>,
464    data_json: &str,
465) -> Result<(), rusqlite::Error> {
466    connection.execute(
467        "INSERT INTO events (occurred_at_ms, kind, run_id, ticket_id, data_json)
468         VALUES (?1, ?2, ?3, ?4, ?5)",
469        params![now_ms, kind, run_id, ticket_id, data_json],
470    )?;
471    Ok(())
472}
473
474#[derive(Debug, Clone, PartialEq, Eq)]
475pub struct CooldownUpdate<'a> {
476    pub target: &'a str,
477    pub until_ms: i64,
478    pub reason: &'a str,
479}
480
481#[derive(Debug, Clone, PartialEq, Eq)]
482pub struct CooldownRecord {
483    pub target: String,
484    pub until_ms: i64,
485    pub reason: String,
486}
487
488/// One executed aftercare stage, persisted alongside the run's outcome.
489#[derive(Debug, Clone, PartialEq, Eq)]
490pub struct StageRecord {
491    pub stage_index: usize,
492    pub stage: String,
493    pub state: String,
494    pub started_at_ms: i64,
495    pub finished_at_ms: i64,
496    pub exit_code: Option<i32>,
497    pub output_ref: String,
498    pub verdict_source: String,
499    pub reason: Option<String>,
500}
501
502#[derive(Debug, Clone, PartialEq, Eq)]
503pub struct QueuedActivation {
504    pub id: String,
505    pub kind: String,
506    pub ticket_id: Option<String>,
507    pub project_id: Option<String>,
508    pub eligible_at_ms: Option<i64>,
509    pub interval_ms: Option<i64>,
510}
511
512#[derive(Debug, Clone, PartialEq, Eq)]
513pub struct ActiveRun {
514    pub id: String,
515    pub ticket_id: String,
516    /// Carried so callers can render the run's alias without a second lookup.
517    pub attempt: i64,
518    pub ticket_name: String,
519    pub project_id: String,
520    pub state: String,
521}
522
523#[derive(Debug, Clone, PartialEq, Eq)]
524pub struct RunRecord {
525    pub id: String,
526    pub ticket_id: String,
527    /// The per-ticket attempt this run served. Frozen at claim, so it is the
528    /// second half of the run's alias.
529    pub attempt: i64,
530    pub state: String,
531    pub branch: Option<String>,
532    pub worktree_path: Option<String>,
533    pub pid: Option<i64>,
534    pub pid_start_time: Option<i64>,
535    pub process_group_id: Option<i64>,
536    pub exit_code: Option<i64>,
537    pub exited_at_ms: Option<i64>,
538    pub flow_json: Option<String>,
539    pub ticket_json: Option<String>,
540}
541
542/// Every `RunRecord` read uses this projection so the column order and the
543/// mapper below can never drift apart.
544const RUN_RECORD_SELECT: &str = "SELECT id, ticket_id, attempt, state, branch, worktree_path, pid,
545            pid_start_time, process_group_id, exit_code, exited_at_ms,
546            flow_json, ticket_json
547     FROM runs";
548
549const TICKET_RECORD_SELECT: &str =
550    "SELECT id, project_id, file_path, source, source_ref, state, name, worktree,
551            target, model, effort, flow, attempts, body, held_reason, created_at_ms
552     FROM tickets";
553
554fn run_record(row: &rusqlite::Row<'_>) -> rusqlite::Result<RunRecord> {
555    Ok(RunRecord {
556        id: row.get(0)?,
557        ticket_id: row.get(1)?,
558        attempt: row.get(2)?,
559        state: row.get(3)?,
560        branch: row.get(4)?,
561        worktree_path: row.get(5)?,
562        pid: row.get(6)?,
563        pid_start_time: row.get(7)?,
564        process_group_id: row.get(8)?,
565        exit_code: row.get(9)?,
566        exited_at_ms: row.get(10)?,
567        flow_json: row.get(11)?,
568        ticket_json: row.get(12)?,
569    })
570}
571
572/// A `needs_review` ticket paired with the preserved run branch whose tip the
573/// daemon can test for external integration against the default branch.
574#[derive(Debug, Clone, PartialEq, Eq)]
575pub(crate) struct NeedsReviewBranch {
576    pub(crate) ticket_id: String,
577    pub(crate) run_id: String,
578    pub(crate) branch: String,
579}
580
581#[derive(Debug, Clone, PartialEq, Eq)]
582pub(crate) struct WorktreeCleanupCandidate {
583    pub(crate) run_id: String,
584    pub(crate) ticket_id: String,
585    pub(crate) branch: String,
586    pub(crate) worktree_path: String,
587    pub(crate) cleanup_eligible_at_ms: i64,
588}
589
590/// One lease that must be classified when a daemon starts. Process identity
591/// and worker credentials are returned only to the daemon recovery path.
592#[derive(Debug, Clone, PartialEq, Eq)]
593pub(crate) struct RecoverableRun {
594    pub(crate) id: String,
595    pub(crate) ticket_id: String,
596    pub(crate) target: String,
597    pub(crate) state: RunState,
598    pub(crate) branch: Option<String>,
599    pub(crate) worktree_path: Option<String>,
600    pub(crate) pid: Option<i64>,
601    pub(crate) pid_start_time: Option<i64>,
602    pub(crate) process_group_id: Option<i64>,
603    pub(crate) worker_token: Option<String>,
604    pub(crate) worker_socket_path: Option<String>,
605    pub(crate) exit_code: Option<i64>,
606    pub(crate) lease_expires_at_ms: i64,
607    pub(crate) flow_json: Option<String>,
608    pub(crate) ticket_json: Option<String>,
609}
610
611#[derive(Debug, Clone, PartialEq, Eq)]
612pub struct ProjectRecord {
613    pub id: String,
614    pub file_path: Option<String>,
615    pub title: String,
616}
617
618#[derive(Debug, Clone, PartialEq, Eq)]
619pub struct ProjectNote {
620    pub id: String,
621    pub run_id: String,
622    pub ticket_id: String,
623    pub text: String,
624    pub recorded_at_ms: i64,
625}
626
627#[derive(Debug, Clone, PartialEq, Eq)]
628pub struct ProjectCommitEvidence {
629    pub run_id: String,
630    pub ticket_id: String,
631    pub data_json: String,
632}
633
634#[derive(Debug, Clone, PartialEq, Eq)]
635pub struct LocalTicketFile {
636    pub id: String,
637    pub file_path: String,
638    pub state: String,
639    pub missing_at_ms: Option<i64>,
640}
641
642#[derive(Debug, Clone, PartialEq, Eq)]
643pub struct TicketRecord {
644    pub id: String,
645    pub project_id: String,
646    pub file_path: Option<String>,
647    pub source: String,
648    pub source_ref: Option<String>,
649    pub state: String,
650    pub name: String,
651    pub blocked_by: Vec<String>,
652    pub worktree: Option<String>,
653    pub target: Option<String>,
654    pub model: Option<String>,
655    pub effort: Option<String>,
656    pub flow: Option<String>,
657    pub attempts: i64,
658    pub body: Option<String>,
659    pub held_reason: Option<String>,
660    /// When the ticket was registered. `sloop list` orders on this.
661    pub created_at_ms: i64,
662}
663
664#[derive(Debug, Clone, PartialEq, Eq)]
665pub struct ReindexTicket {
666    pub id: String,
667    pub project_id: String,
668    pub source: String,
669    pub source_ref: String,
670    pub file_path: Option<String>,
671    pub name: String,
672    pub blocked_by: Vec<String>,
673    pub worktree: String,
674    pub target: Option<String>,
675    pub model: Option<String>,
676    pub effort: Option<String>,
677    pub flow: String,
678    pub body: String,
679    pub held_reason: Option<String>,
680    pub derived_state: Option<TicketState>,
681}
682
683#[derive(Debug, Clone, PartialEq, Eq)]
684pub struct ReindexStateChange {
685    pub ticket_id: String,
686    pub previous_state: String,
687    pub state: String,
688}
689
690#[derive(Debug, Clone, Default, PartialEq, Eq)]
691pub struct ReindexResult {
692    pub state_changes: Vec<ReindexStateChange>,
693    pub rows_dropped: usize,
694}
695
696fn ticket_record(row: &rusqlite::Row<'_>) -> rusqlite::Result<TicketRecord> {
697    Ok(TicketRecord {
698        id: row.get(0)?,
699        project_id: row.get(1)?,
700        file_path: row.get(2)?,
701        source: row.get(3)?,
702        source_ref: row.get(4)?,
703        state: row.get(5)?,
704        name: row.get(6)?,
705        blocked_by: Vec::new(),
706        worktree: row.get(7)?,
707        target: row.get(8)?,
708        model: row.get(9)?,
709        effort: row.get(10)?,
710        flow: row.get(11)?,
711        attempts: row.get(12)?,
712        body: row.get(13)?,
713        held_reason: row.get(14)?,
714        created_at_ms: row.get(15)?,
715    })
716}
717
718fn replace_ticket_blockers(
719    transaction: &rusqlite::Transaction<'_>,
720    ticket_id: &str,
721    blocked_by: &[String],
722) -> rusqlite::Result<()> {
723    transaction.execute(
724        "DELETE FROM ticket_blockers WHERE ticket_id = ?1",
725        params![ticket_id],
726    )?;
727    for (position, blocker_id) in blocked_by.iter().enumerate() {
728        transaction.execute(
729            "INSERT OR IGNORE INTO ticket_blockers (ticket_id, blocker_id, position)
730             VALUES (?1, ?2, ?3)",
731            params![ticket_id, blocker_id, position as i64],
732        )?;
733    }
734    Ok(())
735}
736
737pub struct Store {
738    connection: Connection,
739}
740
741impl Store {
742    /// Opens (creating if needed) the database and migrates it to the current
743    /// schema version. The daemon is the only writer; `now_ms` is injected so
744    /// decision-adjacent timestamps never read the wall clock here.
745    pub fn open(path: &Path, now_ms: i64) -> Result<Self, StoreError> {
746        let connection = Connection::open(path).map_err(|source| StoreError::Open {
747            path: path.to_path_buf(),
748            source,
749        })?;
750        connection.execute_batch(CONNECTION_PRAGMAS)?;
751
752        let mut store = Self { connection };
753        store.migrate(now_ms)?;
754        Ok(store)
755    }
756
757    fn migrate(&mut self, now_ms: i64) -> Result<(), StoreError> {
758        let version: u32 = self
759            .connection
760            .query_row("PRAGMA user_version", [], |row| row.get(0))?;
761        match version {
762            0 => {
763                let transaction = self
764                    .connection
765                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
766                transaction.execute_batch(SCHEMA_V1)?;
767                transaction.execute_batch(ID_COUNTER_SCHEMA)?;
768                transaction.execute_batch(EVENTS_SCHEMA)?;
769                transaction.execute(
770                    "INSERT INTO scheduler_state (singleton, paused, draining, updated_at_ms)
771                     VALUES (1, 0, 0, ?1)",
772                    params![now_ms],
773                )?;
774                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
775                transaction.commit()?;
776                Ok(())
777            }
778            1 => {
779                let transaction = self
780                    .connection
781                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
782                transaction.execute_batch(
783                    "ALTER TABLE tickets ADD COLUMN model TEXT;
784                     ALTER TABLE tickets ADD COLUMN effort TEXT;
785                     ALTER TABLE tickets ADD COLUMN target TEXT;
786                     ALTER TABLE tickets ADD COLUMN name TEXT NOT NULL DEFAULT '';
787                     ALTER TABLE tickets ADD COLUMN worktree TEXT;
788                     ALTER TABLE tickets ADD COLUMN flow TEXT;
789                     ALTER TABLE tickets ADD COLUMN missing_at_ms INTEGER;
790                     ALTER TABLE runs ADD COLUMN worker_socket_path TEXT;
791                     CREATE TABLE ticket_blockers (
792                         ticket_id TEXT NOT NULL REFERENCES tickets(id) ON DELETE CASCADE,
793                         blocker_id TEXT NOT NULL REFERENCES tickets(id),
794                         position INTEGER NOT NULL,
795                         PRIMARY KEY (ticket_id, blocker_id)
796                     );",
797                )?;
798                transaction.execute_batch(RUN_SNAPSHOT_COLUMNS)?;
799                transaction.execute_batch(ID_COUNTER_SCHEMA)?;
800                transaction.execute_batch(EVENTS_SCHEMA)?;
801                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
802                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
803                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
804                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
805                transaction.commit()?;
806                Ok(())
807            }
808            2 => {
809                let transaction = self
810                    .connection
811                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
812                transaction.execute_batch(
813                    "ALTER TABLE tickets ADD COLUMN target TEXT;
814                     ALTER TABLE tickets ADD COLUMN name TEXT NOT NULL DEFAULT '';
815                     ALTER TABLE tickets ADD COLUMN worktree TEXT;
816                     ALTER TABLE tickets ADD COLUMN flow TEXT;
817                     ALTER TABLE tickets ADD COLUMN missing_at_ms INTEGER;
818                     ALTER TABLE runs ADD COLUMN worker_socket_path TEXT;
819                     CREATE TABLE ticket_blockers (
820                         ticket_id TEXT NOT NULL REFERENCES tickets(id) ON DELETE CASCADE,
821                         blocker_id TEXT NOT NULL REFERENCES tickets(id),
822                         position INTEGER NOT NULL,
823                         PRIMARY KEY (ticket_id, blocker_id)
824                     );",
825                )?;
826                transaction.execute_batch(RUN_SNAPSHOT_COLUMNS)?;
827                transaction.execute_batch(ID_COUNTER_SCHEMA)?;
828                transaction.execute_batch(EVENTS_SCHEMA)?;
829                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
830                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
831                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
832                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
833                transaction.commit()?;
834                Ok(())
835            }
836            3 => {
837                let transaction = self
838                    .connection
839                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
840                transaction.execute_batch(
841                    "ALTER TABLE tickets ADD COLUMN name TEXT NOT NULL DEFAULT '';
842                     ALTER TABLE tickets ADD COLUMN worktree TEXT;
843                     ALTER TABLE tickets ADD COLUMN flow TEXT;
844                     ALTER TABLE tickets ADD COLUMN missing_at_ms INTEGER;
845                     ALTER TABLE runs ADD COLUMN worker_socket_path TEXT;
846                     CREATE TABLE ticket_blockers (
847                         ticket_id TEXT NOT NULL REFERENCES tickets(id) ON DELETE CASCADE,
848                         blocker_id TEXT NOT NULL REFERENCES tickets(id),
849                         position INTEGER NOT NULL,
850                         PRIMARY KEY (ticket_id, blocker_id)
851                     );",
852                )?;
853                transaction.execute_batch(RUN_SNAPSHOT_COLUMNS)?;
854                transaction.execute_batch(ID_COUNTER_SCHEMA)?;
855                transaction.execute_batch(EVENTS_SCHEMA)?;
856                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
857                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
858                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
859                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
860                transaction.commit()?;
861                Ok(())
862            }
863            4 => {
864                let transaction = self
865                    .connection
866                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
867                transaction.execute_batch(
868                    "ALTER TABLE tickets ADD COLUMN flow TEXT;
869                     ALTER TABLE tickets ADD COLUMN missing_at_ms INTEGER;
870                     ALTER TABLE runs ADD COLUMN worker_socket_path TEXT;",
871                )?;
872                transaction.execute_batch(RUN_SNAPSHOT_COLUMNS)?;
873                transaction.execute_batch(ID_COUNTER_SCHEMA)?;
874                transaction.execute_batch(EVENTS_SCHEMA)?;
875                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
876                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
877                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
878                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
879                transaction.commit()?;
880                Ok(())
881            }
882            5 => {
883                let transaction = self
884                    .connection
885                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
886                transaction.execute_batch(
887                    "ALTER TABLE tickets ADD COLUMN missing_at_ms INTEGER;
888                         ALTER TABLE runs ADD COLUMN worker_socket_path TEXT;",
889                )?;
890                transaction.execute_batch(RUN_SNAPSHOT_COLUMNS)?;
891                transaction.execute_batch(ID_COUNTER_SCHEMA)?;
892                transaction.execute_batch(EVENTS_SCHEMA)?;
893                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
894                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
895                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
896                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
897                transaction.commit()?;
898                Ok(())
899            }
900            6 => {
901                let transaction = self
902                    .connection
903                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
904                transaction
905                    .execute_batch("ALTER TABLE runs ADD COLUMN worker_socket_path TEXT;")?;
906                transaction.execute_batch(RUN_SNAPSHOT_COLUMNS)?;
907                transaction.execute_batch(ID_COUNTER_SCHEMA)?;
908                transaction.execute_batch(EVENTS_SCHEMA)?;
909                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
910                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
911                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
912                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
913                transaction.commit()?;
914                Ok(())
915            }
916            7 => {
917                let transaction = self
918                    .connection
919                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
920                transaction.execute_batch(RUN_SNAPSHOT_COLUMNS)?;
921                transaction.execute_batch(ID_COUNTER_SCHEMA)?;
922                transaction.execute_batch(EVENTS_SCHEMA)?;
923                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
924                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
925                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
926                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
927                transaction.commit()?;
928                Ok(())
929            }
930            8 => {
931                let transaction = self
932                    .connection
933                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
934                transaction.execute_batch(RUN_SNAPSHOT_COLUMNS)?;
935                transaction.execute_batch(EVENTS_SCHEMA)?;
936                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
937                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
938                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
939                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
940                transaction.commit()?;
941                Ok(())
942            }
943            9 => {
944                let transaction = self
945                    .connection
946                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
947                transaction.execute_batch(EVENTS_SCHEMA)?;
948                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
949                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
950                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
951                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
952                transaction.commit()?;
953                Ok(())
954            }
955            10 => {
956                let transaction = self
957                    .connection
958                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
959                transaction.execute_batch(TICKET_SOURCE_COLUMNS)?;
960                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
961                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
962                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
963                transaction.commit()?;
964                Ok(())
965            }
966            11 => {
967                let transaction = self
968                    .connection
969                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
970                transaction.execute_batch(RESTART_DRAINING_COLUMN)?;
971                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
972                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
973                transaction.commit()?;
974                Ok(())
975            }
976            12 => {
977                let transaction = self
978                    .connection
979                    .transaction_with_behavior(TransactionBehavior::Immediate)?;
980                transaction.execute_batch(WORKTREE_CLEANUP_COLUMNS)?;
981                transaction.pragma_update(None, "user_version", SCHEMA_VERSION)?;
982                transaction.commit()?;
983                Ok(())
984            }
985            SCHEMA_VERSION => Ok(()),
986            newer => Err(StoreError::UnsupportedSchemaVersion(newer)),
987        }
988    }
989
990    pub fn insert_local_project(
991        &self,
992        id: &str,
993        file_path: &str,
994        title: &str,
995        now_ms: i64,
996    ) -> Result<(), StoreError> {
997        self.connection.execute(
998            "INSERT INTO projects (id, file_path, source, title, created_at_ms, updated_at_ms)
999             VALUES (?1, ?2, 'local', ?3, ?4, ?4)",
1000            params![id, file_path, title, now_ms],
1001        )?;
1002        Ok(())
1003    }
1004
1005    /// Inserts or refreshes a project indexed from a committed file. Startup
1006    /// and reindex call this for every configured project file, so it must tolerate
1007    /// rows that already exist.
1008    pub fn upsert_local_project(
1009        &self,
1010        id: &str,
1011        file_path: &str,
1012        title: &str,
1013        now_ms: i64,
1014    ) -> Result<(), StoreError> {
1015        self.connection.execute(
1016            "INSERT INTO projects (id, file_path, source, title, created_at_ms, updated_at_ms)
1017             VALUES (?1, ?2, 'local', ?3, ?4, ?4)
1018             ON CONFLICT(id) DO UPDATE SET
1019                 file_path = excluded.file_path,
1020                 title = excluded.title,
1021                 updated_at_ms = excluded.updated_at_ms",
1022            params![id, file_path, title, now_ms],
1023        )?;
1024        Ok(())
1025    }
1026
1027    pub fn project_exists(&self, id: &str) -> Result<bool, StoreError> {
1028        let found: Option<i64> = self
1029            .connection
1030            .query_row("SELECT 1 FROM projects WHERE id = ?1", params![id], |row| {
1031                row.get(0)
1032            })
1033            .optional()?;
1034        Ok(found.is_some())
1035    }
1036
1037    pub fn project(&self, id: &str) -> Result<Option<ProjectRecord>, StoreError> {
1038        self.connection
1039            .query_row(
1040                "SELECT id, file_path, title FROM projects WHERE id = ?1",
1041                params![id],
1042                |row| {
1043                    Ok(ProjectRecord {
1044                        id: row.get(0)?,
1045                        file_path: row.get(1)?,
1046                        title: row.get(2)?,
1047                    })
1048                },
1049            )
1050            .optional()
1051            .map_err(StoreError::from)
1052    }
1053
1054    #[allow(clippy::too_many_arguments)]
1055    pub fn insert_local_ticket(
1056        &self,
1057        id: &str,
1058        project_id: &str,
1059        file_path: &str,
1060        name: &str,
1061        blocked_by: &[String],
1062        worktree: &str,
1063        target: Option<&str>,
1064        model: Option<&str>,
1065        effort: Option<&str>,
1066        flow: &str,
1067        state: TicketState,
1068        now_ms: i64,
1069    ) -> Result<(), StoreError> {
1070        let transaction = self.immediate_transaction()?;
1071        transaction.execute(
1072            "INSERT INTO tickets
1073                 (id, project_id, file_path, source, state, name, worktree, target, model, effort,
1074                    flow, created_at_ms, updated_at_ms)
1075             VALUES (?1, ?2, ?3, 'local', ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?11)",
1076            params![
1077                id,
1078                project_id,
1079                file_path,
1080                state.as_str(),
1081                name,
1082                worktree,
1083                target,
1084                model,
1085                effort,
1086                flow,
1087                now_ms
1088            ],
1089        )?;
1090        replace_ticket_blockers(&transaction, id, blocked_by)?;
1091        transaction.commit()?;
1092        Ok(())
1093    }
1094
1095    #[allow(clippy::too_many_arguments)]
1096    pub fn update_local_ticket(
1097        &self,
1098        id: &str,
1099        name: &str,
1100        blocked_by: &[String],
1101        worktree: &str,
1102        target: Option<&str>,
1103        model: Option<&str>,
1104        effort: Option<&str>,
1105        flow: &str,
1106        now_ms: i64,
1107    ) -> Result<(), StoreError> {
1108        let transaction = self.immediate_transaction()?;
1109        transaction.execute(
1110            "UPDATE tickets
1111             SET name = ?2, worktree = ?3, target = ?4, model = ?5, effort = ?6, flow = ?7,
1112                  held_reason = NULL, missing_at_ms = NULL, updated_at_ms = ?8
1113             WHERE id = ?1",
1114            params![id, name, worktree, target, model, effort, flow, now_ms],
1115        )?;
1116        replace_ticket_blockers(&transaction, id, blocked_by)?;
1117        transaction.commit()?;
1118        Ok(())
1119    }
1120
1121    /// Applies a complete authored ticket snapshot without disturbing runtime
1122    /// history for IDs that remain present. Missing local rows and everything
1123    /// that depends on them are removed explicitly so the operation can report
1124    /// how much non-derivable state was discarded.
1125    pub fn apply_reindex(
1126        &self,
1127        project_ids: &[String],
1128        tickets: &[ReindexTicket],
1129        now_ms: i64,
1130    ) -> Result<ReindexResult, StoreError> {
1131        let existing: BTreeMap<String, TicketRecord> = self
1132            .tickets()?
1133            .into_iter()
1134            .map(|ticket| (ticket.id.clone(), ticket))
1135            .collect();
1136        let desired_ticket_ids: BTreeSet<&str> =
1137            tickets.iter().map(|ticket| ticket.id.as_str()).collect();
1138        let desired_project_ids: BTreeSet<&str> = project_ids.iter().map(String::as_str).collect();
1139        let transaction = self.immediate_transaction()?;
1140
1141        let stale_tickets = {
1142            let mut statement = transaction.prepare("SELECT id FROM tickets ORDER BY id")?;
1143            statement
1144                .query_map([], |row| row.get::<_, String>(0))?
1145                .collect::<Result<Vec<_>, _>>()?
1146                .into_iter()
1147                .filter(|id| !desired_ticket_ids.contains(id.as_str()))
1148                .collect::<Vec<_>>()
1149        };
1150        let stale_projects = {
1151            let mut statement = transaction
1152                .prepare("SELECT id FROM projects WHERE source = 'local' ORDER BY id")?;
1153            statement
1154                .query_map([], |row| row.get::<_, String>(0))?
1155                .collect::<Result<Vec<_>, _>>()?
1156                .into_iter()
1157                .filter(|id| !desired_project_ids.contains(id.as_str()))
1158                .collect::<Vec<_>>()
1159        };
1160
1161        let mut doomed_activations = BTreeSet::new();
1162        for ticket_id in &stale_tickets {
1163            let mut statement =
1164                transaction.prepare("SELECT id FROM activations WHERE ticket_id = ?1")?;
1165            doomed_activations.extend(
1166                statement
1167                    .query_map(params![ticket_id], |row| row.get::<_, String>(0))?
1168                    .collect::<Result<Vec<_>, _>>()?,
1169            );
1170        }
1171        for project_id in &stale_projects {
1172            let activations = {
1173                let mut statement = transaction
1174                    .prepare("SELECT id, state FROM activations WHERE project_id = ?1")?;
1175                statement
1176                    .query_map(params![project_id], |row| {
1177                        Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
1178                    })?
1179                    .collect::<Result<Vec<_>, _>>()?
1180            };
1181            for (activation_id, activation_state) in activations {
1182                if activation_state == "queued" {
1183                    doomed_activations.insert(activation_id);
1184                } else {
1185                    transaction.execute(
1186                        "UPDATE activations SET project_id = NULL WHERE id = ?1",
1187                        params![activation_id],
1188                    )?;
1189                }
1190            }
1191        }
1192
1193        let mut doomed_runs = BTreeSet::new();
1194        for ticket_id in &stale_tickets {
1195            let mut statement = transaction.prepare("SELECT id FROM runs WHERE ticket_id = ?1")?;
1196            doomed_runs.extend(
1197                statement
1198                    .query_map(params![ticket_id], |row| row.get::<_, String>(0))?
1199                    .collect::<Result<Vec<_>, _>>()?,
1200            );
1201        }
1202        for activation_id in &doomed_activations {
1203            let mut statement =
1204                transaction.prepare("SELECT id FROM runs WHERE activation_id = ?1")?;
1205            doomed_runs.extend(
1206                statement
1207                    .query_map(params![activation_id], |row| row.get::<_, String>(0))?
1208                    .collect::<Result<Vec<_>, _>>()?,
1209            );
1210        }
1211
1212        let mut rows_dropped = 0;
1213        for run_id in &doomed_runs {
1214            transaction.execute(
1215                "UPDATE cooldowns SET source_run_id = NULL WHERE source_run_id = ?1",
1216                params![run_id],
1217            )?;
1218            for table in [
1219                "leases",
1220                "run_evidence",
1221                "aftercare_stages",
1222                "budget_reservations",
1223                "notes",
1224            ] {
1225                rows_dropped += transaction.execute(
1226                    &format!("DELETE FROM {table} WHERE run_id = ?1"),
1227                    params![run_id],
1228                )?;
1229            }
1230            rows_dropped +=
1231                transaction.execute("DELETE FROM runs WHERE id = ?1", params![run_id])?;
1232        }
1233        for activation_id in &doomed_activations {
1234            rows_dropped += transaction.execute(
1235                "DELETE FROM activation_filters WHERE activation_id = ?1",
1236                params![activation_id],
1237            )?;
1238            rows_dropped += transaction.execute(
1239                "DELETE FROM activations WHERE id = ?1",
1240                params![activation_id],
1241            )?;
1242        }
1243        for ticket_id in &stale_tickets {
1244            rows_dropped += transaction.execute(
1245                "DELETE FROM activation_filters WHERE ticket_id = ?1",
1246                params![ticket_id],
1247            )?;
1248            rows_dropped += transaction.execute(
1249                "DELETE FROM leases WHERE ticket_id = ?1",
1250                params![ticket_id],
1251            )?;
1252            rows_dropped += transaction.execute(
1253                "DELETE FROM ticket_blockers WHERE ticket_id = ?1 OR blocker_id = ?1",
1254                params![ticket_id],
1255            )?;
1256            rows_dropped +=
1257                transaction.execute("DELETE FROM tickets WHERE id = ?1", params![ticket_id])?;
1258        }
1259        let mut state_changes = Vec::new();
1260        for ticket in tickets {
1261            let previous = existing.get(&ticket.id);
1262            let state = if ticket.held_reason.is_some() {
1263                TicketState::Held.as_str()
1264            } else {
1265                match (previous, ticket.derived_state) {
1266                    (Some(_), Some(derived)) => derived.as_str(),
1267                    (Some(existing), None) if existing.held_reason.is_some() => {
1268                        TicketState::Ready.as_str()
1269                    }
1270                    (Some(existing), None) => existing.state.as_str(),
1271                    (None, Some(derived)) => derived.as_str(),
1272                    (None, None) => TicketState::Ready.as_str(),
1273                }
1274            };
1275            if let Some(previous) = previous
1276                && previous.state != state
1277            {
1278                state_changes.push(ReindexStateChange {
1279                    ticket_id: ticket.id.clone(),
1280                    previous_state: previous.state.clone(),
1281                    state: state.to_owned(),
1282                });
1283                if state == TicketState::Merged.as_str()
1284                    && matches!(previous.state.as_str(), "failed" | "needs_review")
1285                {
1286                    transaction.execute(
1287                        "UPDATE runs SET cleanup_eligible_at_ms = ?2
1288                         WHERE ticket_id = ?1 AND state IN ('failed', 'needs_review')
1289                           AND cleanup_eligible_at_ms IS NULL AND cleaned_at_ms IS NULL",
1290                        params![ticket.id, now_ms],
1291                    )?;
1292                }
1293            }
1294            transaction.execute(
1295                "INSERT INTO tickets
1296                     (id, project_id, file_path, source, source_ref, state, name, worktree, target,
1297                      model, effort, flow, body, held_reason, created_at_ms, updated_at_ms)
1298                 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?15)
1299                 ON CONFLICT(id) DO UPDATE SET
1300                     project_id = excluded.project_id,
1301                     file_path = excluded.file_path,
1302                     source = excluded.source,
1303                     source_ref = excluded.source_ref,
1304                     state = excluded.state,
1305                     name = excluded.name,
1306                     worktree = excluded.worktree,
1307                     target = excluded.target,
1308                     model = excluded.model,
1309                     effort = excluded.effort,
1310                     flow = excluded.flow,
1311                     body = excluded.body,
1312                     held_reason = excluded.held_reason,
1313                     missing_at_ms = NULL,
1314                     updated_at_ms = excluded.updated_at_ms",
1315                params![
1316                    ticket.id,
1317                    ticket.project_id,
1318                    ticket.file_path,
1319                    ticket.source,
1320                    ticket.source_ref,
1321                    state,
1322                    ticket.name,
1323                    ticket.worktree,
1324                    ticket.target,
1325                    ticket.model,
1326                    ticket.effort,
1327                    ticket.flow,
1328                    ticket.body,
1329                    ticket.held_reason,
1330                    now_ms,
1331                ],
1332            )?;
1333        }
1334        for ticket in tickets {
1335            replace_ticket_blockers(&transaction, &ticket.id, &ticket.blocked_by)?;
1336        }
1337        for project_id in &stale_projects {
1338            rows_dropped +=
1339                transaction.execute("DELETE FROM projects WHERE id = ?1", params![project_id])?;
1340        }
1341
1342        state_changes.sort_by(|left, right| left.ticket_id.cmp(&right.ticket_id));
1343        transaction.commit()?;
1344        Ok(ReindexResult {
1345            state_changes,
1346            rows_dropped,
1347        })
1348    }
1349
1350    pub fn update_ticket_execution(
1351        &self,
1352        id: &str,
1353        target: Option<&str>,
1354        model: Option<&str>,
1355        effort: Option<&str>,
1356        now_ms: i64,
1357    ) -> Result<(), StoreError> {
1358        self.connection.execute(
1359            "UPDATE tickets SET target = ?2, model = ?3, effort = ?4, updated_at_ms = ?5 WHERE id = ?1",
1360            params![id, target, model, effort, now_ms],
1361        )?;
1362        Ok(())
1363    }
1364
1365    pub fn update_ticket_body(&self, id: &str, body: &str, now_ms: i64) -> Result<(), StoreError> {
1366        self.connection.execute(
1367            "UPDATE tickets SET body = ?2, updated_at_ms = ?3 WHERE id = ?1",
1368            params![id, body, now_ms],
1369        )?;
1370        Ok(())
1371    }
1372
1373    /// Version-two rows predate target snapshots. Once a repository has a
1374    /// target configuration, persist its default before dispatch can observe
1375    /// those rows.
1376    pub fn backfill_ticket_targets(
1377        &self,
1378        default_target: &str,
1379        now_ms: i64,
1380    ) -> Result<usize, StoreError> {
1381        self.connection
1382            .execute(
1383                "UPDATE tickets SET target = ?1, updated_at_ms = ?2 WHERE target IS NULL",
1384                params![default_target, now_ms],
1385            )
1386            .map_err(StoreError::from)
1387    }
1388
1389    pub fn ticket(&self, id: &str) -> Result<Option<TicketRecord>, StoreError> {
1390        let mut ticket = self
1391            .connection
1392            .query_row(
1393                &format!("{TICKET_RECORD_SELECT} WHERE id = ?1"),
1394                params![id],
1395                ticket_record,
1396            )
1397            .optional()?;
1398        if let Some(ticket) = ticket.as_mut() {
1399            ticket.blocked_by = self.ticket_blockers(&ticket.id)?;
1400        }
1401        Ok(ticket)
1402    }
1403
1404    /// Resolves a ticket by its human-facing name. Names are not guaranteed
1405    /// unique across projects, so the lowest id wins deterministically; `show`
1406    /// tries this only after an exact id match fails.
1407    pub fn ticket_by_name(&self, name: &str) -> Result<Option<TicketRecord>, StoreError> {
1408        let mut ticket = self
1409            .connection
1410            .query_row(
1411                &format!("{TICKET_RECORD_SELECT} WHERE name = ?1 ORDER BY id LIMIT 1"),
1412                params![name],
1413                ticket_record,
1414            )
1415            .optional()?;
1416        if let Some(ticket) = ticket.as_mut() {
1417            ticket.blocked_by = self.ticket_blockers(&ticket.id)?;
1418        }
1419        Ok(ticket)
1420    }
1421
1422    pub fn ticket_by_file(&self, file_path: &str) -> Result<Option<TicketRecord>, StoreError> {
1423        let mut ticket = self
1424            .connection
1425            .query_row(
1426                &format!("{TICKET_RECORD_SELECT} WHERE file_path = ?1"),
1427                params![file_path],
1428                ticket_record,
1429            )
1430            .optional()?;
1431        if let Some(ticket) = ticket.as_mut() {
1432            ticket.blocked_by = self.ticket_blockers(&ticket.id)?;
1433        }
1434        Ok(ticket)
1435    }
1436
1437    pub fn ticket_by_source_ref(
1438        &self,
1439        source: &str,
1440        source_ref: &str,
1441    ) -> Result<Option<TicketRecord>, StoreError> {
1442        let mut ticket = self
1443            .connection
1444            .query_row(
1445                &format!("{TICKET_RECORD_SELECT} WHERE source = ?1 AND source_ref = ?2"),
1446                params![source, source_ref],
1447                ticket_record,
1448            )
1449            .optional()?;
1450        if let Some(ticket) = ticket.as_mut() {
1451            ticket.blocked_by = self.ticket_blockers(&ticket.id)?;
1452        }
1453        Ok(ticket)
1454    }
1455
1456    /// Every ticket, newest registration first. `sloop list` answers "what is
1457    /// going on right now?", so recency leads; SQL settles the coarse order and
1458    /// a stable pass re-breaks ties on the id's numeric ordinal, which string
1459    /// comparison gets wrong (`TICK-9` sorts above `TICK-38`). Ids with no
1460    /// ordinal keep the deterministic `id DESC` order SQL gave them.
1461    pub fn tickets(&self) -> Result<Vec<TicketRecord>, StoreError> {
1462        let mut statement = self.connection.prepare(&format!(
1463            "{TICKET_RECORD_SELECT} ORDER BY created_at_ms DESC, id DESC"
1464        ))?;
1465        let mut tickets = statement
1466            .query_map([], ticket_record)?
1467            .collect::<Result<Vec<_>, _>>()?;
1468        tickets.sort_by_key(|ticket| {
1469            (
1470                std::cmp::Reverse(ticket.created_at_ms),
1471                std::cmp::Reverse(crate::ids::ordinal(&ticket.id).unwrap_or(0)),
1472            )
1473        });
1474        let mut blockers = self.all_ticket_blockers()?;
1475        for ticket in &mut tickets {
1476            ticket.blocked_by = blockers.remove(&ticket.id).unwrap_or_default();
1477        }
1478        Ok(tickets)
1479    }
1480
1481    pub fn tickets_for_project(&self, project_id: &str) -> Result<Vec<TicketRecord>, StoreError> {
1482        let mut statement = self.connection.prepare(&format!(
1483            "{TICKET_RECORD_SELECT} WHERE project_id = ?1 ORDER BY id"
1484        ))?;
1485        let mut tickets = statement
1486            .query_map(params![project_id], ticket_record)?
1487            .collect::<Result<Vec<_>, _>>()?;
1488        let mut blockers = self.all_ticket_blockers()?;
1489        for ticket in &mut tickets {
1490            ticket.blocked_by = blockers.remove(&ticket.id).unwrap_or_default();
1491        }
1492        Ok(tickets)
1493    }
1494
1495    pub fn ticket_dependencies(
1496        &self,
1497    ) -> Result<std::collections::BTreeMap<String, Vec<String>>, StoreError> {
1498        let mut dependencies = std::collections::BTreeMap::new();
1499        let mut statement = self.connection.prepare("SELECT id FROM tickets")?;
1500        let ids = statement.query_map([], |row| row.get::<_, String>(0))?;
1501        for id in ids {
1502            dependencies.insert(id?, Vec::new());
1503        }
1504        for (ticket_id, blockers) in self.all_ticket_blockers()? {
1505            if let Some(entry) = dependencies.get_mut(&ticket_id) {
1506                *entry = blockers;
1507            }
1508        }
1509        Ok(dependencies)
1510    }
1511
1512    /// Every ticket's blockers in one pass, keeping each list in declared
1513    /// order. Loading these per ticket turns any all-tickets read into a
1514    /// query per row, which is what the post path pays cycle checks against.
1515    fn all_ticket_blockers(&self) -> Result<BTreeMap<String, Vec<String>>, StoreError> {
1516        let mut statement = self.connection.prepare(
1517            "SELECT ticket_id, blocker_id FROM ticket_blockers
1518             ORDER BY ticket_id, position, blocker_id",
1519        )?;
1520        let rows = statement.query_map([], |row| {
1521            Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
1522        })?;
1523        let mut blockers: BTreeMap<String, Vec<String>> = BTreeMap::new();
1524        for row in rows {
1525            let (ticket_id, blocker_id) = row?;
1526            blockers.entry(ticket_id).or_default().push(blocker_id);
1527        }
1528        Ok(blockers)
1529    }
1530
1531    fn ticket_blockers(&self, id: &str) -> Result<Vec<String>, StoreError> {
1532        let mut statement = self.connection.prepare(
1533            "SELECT blocker_id FROM ticket_blockers
1534             WHERE ticket_id = ?1 ORDER BY position, blocker_id",
1535        )?;
1536        statement
1537            .query_map(params![id], |row| row.get(0))?
1538            .collect::<Result<Vec<_>, _>>()
1539            .map_err(StoreError::from)
1540    }
1541
1542    pub fn ticket_ids(&self) -> Result<Vec<String>, StoreError> {
1543        let mut statement = self.connection.prepare("SELECT id FROM tickets")?;
1544        let rows = statement.query_map([], |row| row.get(0))?;
1545        rows.collect::<Result<Vec<_>, _>>()
1546            .map_err(StoreError::from)
1547    }
1548
1549    pub fn local_ticket_files(&self) -> Result<Vec<LocalTicketFile>, StoreError> {
1550        let mut statement = self.connection.prepare(
1551            "SELECT id, file_path, state, missing_at_ms FROM tickets
1552             WHERE source = 'local' AND file_path IS NOT NULL
1553             ORDER BY id",
1554        )?;
1555        let rows = statement.query_map([], |row| {
1556            Ok(LocalTicketFile {
1557                id: row.get(0)?,
1558                file_path: row.get(1)?,
1559                state: row.get(2)?,
1560                missing_at_ms: row.get(3)?,
1561            })
1562        })?;
1563        rows.collect::<Result<Vec<_>, _>>()
1564            .map_err(StoreError::from)
1565    }
1566
1567    /// Whether run history, a lease, an activation, or another ticket's
1568    /// blocker list still points at this row; deleting it would then violate
1569    /// a foreign key or orphan run evidence.
1570    pub fn ticket_is_referenced(&self, id: &str) -> Result<bool, StoreError> {
1571        let referenced = self.connection.query_row(
1572            "SELECT EXISTS (SELECT 1 FROM runs WHERE ticket_id = ?1)
1573                 OR EXISTS (SELECT 1 FROM leases WHERE ticket_id = ?1)
1574                 OR EXISTS (SELECT 1 FROM activations WHERE ticket_id = ?1)
1575                 OR EXISTS (SELECT 1 FROM activation_filters WHERE ticket_id = ?1)
1576                 OR EXISTS (SELECT 1 FROM ticket_blockers WHERE blocker_id = ?1)",
1577            params![id],
1578            |row| row.get(0),
1579        )?;
1580        Ok(referenced)
1581    }
1582
1583    pub fn delete_ticket(&self, id: &str) -> Result<(), StoreError> {
1584        self.connection
1585            .execute("DELETE FROM tickets WHERE id = ?1", params![id])?;
1586        Ok(())
1587    }
1588
1589    /// Stamps a ticket whose committed file has disappeared. The stamp keeps
1590    /// the row out of selection without disturbing its state; an existing
1591    /// stamp is preserved so the deletion clock starts at the first pass.
1592    pub fn mark_ticket_missing(&self, id: &str, now_ms: i64) -> Result<(), StoreError> {
1593        self.connection.execute(
1594            "UPDATE tickets SET missing_at_ms = ?2, updated_at_ms = ?2
1595             WHERE id = ?1 AND missing_at_ms IS NULL",
1596            params![id, now_ms],
1597        )?;
1598        Ok(())
1599    }
1600
1601    pub fn clear_ticket_missing(&self, id: &str, now_ms: i64) -> Result<(), StoreError> {
1602        self.connection.execute(
1603            "UPDATE tickets SET missing_at_ms = NULL, updated_at_ms = ?2
1604             WHERE id = ?1 AND missing_at_ms IS NOT NULL",
1605            params![id, now_ms],
1606        )?;
1607        Ok(())
1608    }
1609
1610    pub fn ticket_state(&self, id: &str) -> Result<Option<String>, StoreError> {
1611        let state = self
1612            .connection
1613            .query_row(
1614                "SELECT state FROM tickets WHERE id = ?1",
1615                params![id],
1616                |row| row.get(0),
1617            )
1618            .optional()?;
1619        Ok(state)
1620    }
1621
1622    /// Applies the operator-controlled ready/held side-state transition. The
1623    /// conditional update prevents an override from stealing a live claim or
1624    /// rewriting an evidence-derived outcome.
1625    pub fn set_ticket_hold(
1626        &self,
1627        id: &str,
1628        state: TicketState,
1629        now_ms: i64,
1630    ) -> Result<String, StoreError> {
1631        debug_assert!(matches!(state, TicketState::Ready | TicketState::Held));
1632        let requested = state.as_str();
1633        let previous = self
1634            .ticket_state(id)?
1635            .ok_or_else(|| StoreError::TicketNotFound {
1636                ticket_id: id.into(),
1637            })?;
1638        if previous == requested {
1639            return Ok(previous);
1640        }
1641        let allowed_previous = match state {
1642            TicketState::Ready => TicketState::Held.as_str(),
1643            TicketState::Held => TicketState::Ready.as_str(),
1644            _ => unreachable!("hold transitions only use ready and held"),
1645        };
1646        let changed = self.connection.execute(
1647            "UPDATE tickets SET state = ?2, held_reason = NULL, updated_at_ms = ?3
1648             WHERE id = ?1 AND state = ?4",
1649            params![id, requested, now_ms, allowed_previous],
1650        )?;
1651        if changed != 1 {
1652            return Err(StoreError::TicketStateConflict {
1653                ticket_id: id.into(),
1654                state: previous,
1655                requested: requested.into(),
1656            });
1657        }
1658        Ok(previous)
1659    }
1660
1661    /// Returns a failed ticket to the ready queue and starts its attempt
1662    /// counter over. Other states remain evidence-derived and immutable here.
1663    pub fn retry_ticket(&self, id: &str, now_ms: i64) -> Result<String, StoreError> {
1664        let previous = self
1665            .ticket_state(id)?
1666            .ok_or_else(|| StoreError::TicketNotFound {
1667                ticket_id: id.into(),
1668            })?;
1669        let transaction = self.immediate_transaction()?;
1670        let changed = transaction.execute(
1671            "UPDATE tickets SET state = 'ready', held_reason = NULL, attempts = 0, updated_at_ms = ?2
1672             WHERE id = ?1 AND state = 'failed'",
1673            params![id, now_ms],
1674        )?;
1675        if changed != 1 {
1676            return Err(StoreError::TicketStateConflict {
1677                ticket_id: id.into(),
1678                state: previous,
1679                requested: TicketState::Ready.as_str().into(),
1680            });
1681        }
1682        transaction.execute(
1683            "UPDATE runs SET cleanup_eligible_at_ms = ?2
1684             WHERE ticket_id = ?1 AND state = 'failed'
1685               AND cleanup_eligible_at_ms IS NULL AND cleaned_at_ms IS NULL",
1686            params![id, now_ms],
1687        )?;
1688        transaction.commit()?;
1689        Ok(previous)
1690    }
1691
1692    pub fn insert_activation(
1693        &self,
1694        activation: &NewActivation<'_>,
1695        now_ms: i64,
1696    ) -> Result<(), StoreError> {
1697        self.connection.execute(
1698            "INSERT INTO activations
1699                 (id, kind, state, ticket_id, project_id, eligible_at_ms, interval_ms,
1700                  created_at_ms, updated_at_ms)
1701             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?8)",
1702            params![
1703                activation.id,
1704                activation.kind.as_str(),
1705                ActivationState::Queued.as_str(),
1706                activation.ticket_id,
1707                activation.project_id,
1708                activation.eligible_at_ms,
1709                activation.interval_ms,
1710                now_ms,
1711            ],
1712        )?;
1713        Ok(())
1714    }
1715
1716    pub fn insert_activation_filter(
1717        &self,
1718        activation_id: &str,
1719        ticket_id: &str,
1720    ) -> Result<(), StoreError> {
1721        self.connection.execute(
1722            "INSERT OR IGNORE INTO activation_filters (activation_id, ticket_id) VALUES (?1, ?2)",
1723            params![activation_id, ticket_id],
1724        )?;
1725        Ok(())
1726    }
1727
1728    pub fn queued_activations(&self) -> Result<Vec<QueuedActivation>, StoreError> {
1729        let mut statement = self.connection.prepare(
1730            "SELECT id, kind, ticket_id, project_id, eligible_at_ms, interval_ms
1731             FROM activations WHERE state = 'queued'
1732             ORDER BY created_at_ms, id",
1733        )?;
1734        let activations = statement
1735            .query_map([], |row| {
1736                Ok(QueuedActivation {
1737                    id: row.get(0)?,
1738                    kind: row.get(1)?,
1739                    ticket_id: row.get(2)?,
1740                    project_id: row.get(3)?,
1741                    eligible_at_ms: row.get(4)?,
1742                    interval_ms: row.get(5)?,
1743                })
1744            })?
1745            .collect::<Result<Vec<_>, _>>()?;
1746        Ok(activations)
1747    }
1748
1749    /// Queued activations whose time gate is open, oldest first.
1750    pub fn dispatchable_activations(
1751        &self,
1752        now_ms: i64,
1753    ) -> Result<Vec<QueuedActivation>, StoreError> {
1754        let mut statement = self.connection.prepare(
1755            "SELECT id, kind, ticket_id, project_id, eligible_at_ms, interval_ms
1756             FROM activations
1757             WHERE state = 'queued'
1758               AND (kind IN ('immediate', 'auto') OR eligible_at_ms <= ?1)
1759             ORDER BY created_at_ms, id",
1760        )?;
1761        let activations = statement
1762            .query_map(params![now_ms], |row| {
1763                Ok(QueuedActivation {
1764                    id: row.get(0)?,
1765                    kind: row.get(1)?,
1766                    ticket_id: row.get(2)?,
1767                    project_id: row.get(3)?,
1768                    eligible_at_ms: row.get(4)?,
1769                    interval_ms: row.get(5)?,
1770                })
1771            })?
1772            .collect::<Result<Vec<_>, _>>()?;
1773        Ok(activations)
1774    }
1775
1776    pub fn next_activation_eligible_at_ms(&self, now_ms: i64) -> Result<Option<i64>, StoreError> {
1777        self.connection
1778            .query_row(
1779                "SELECT MIN(eligible_at_ms) FROM activations
1780                 WHERE state = 'queued' AND eligible_at_ms > ?1",
1781                params![now_ms],
1782                |row| row.get(0),
1783            )
1784            .map_err(StoreError::from)
1785    }
1786
1787    /// Deterministic ready-work selection within an activation's scope:
1788    /// oldest registration first, ticket ID as the tiebreak. `--only` filters
1789    /// apply when the activation has filter rows.
1790    pub fn select_ready_ticket(
1791        &self,
1792        activation: &QueuedActivation,
1793        now_ms: i64,
1794    ) -> Result<Option<String>, StoreError> {
1795        let ticket = self
1796            .connection
1797            .query_row(
1798                "SELECT t.id FROM tickets t
1799                 WHERE t.state = 'ready'
1800                   AND t.missing_at_ms IS NULL
1801                   AND (?1 IS NULL OR t.project_id = ?1)
1802                   AND NOT EXISTS (SELECT 1 FROM ticket_blockers b
1803                                   JOIN tickets bt ON bt.id = b.blocker_id
1804                                   WHERE b.ticket_id = t.id
1805                                     AND bt.state != 'merged')
1806                    AND (NOT EXISTS (SELECT 1 FROM activation_filters f
1807                                    WHERE f.activation_id = ?2)
1808                        OR EXISTS (SELECT 1 FROM activation_filters f
1809                                    WHERE f.activation_id = ?2 AND f.ticket_id = t.id))
1810                   AND NOT EXISTS (SELECT 1 FROM cooldowns c
1811                                   WHERE c.key = 'agent_target:' || t.target
1812                                     AND c.until_ms > ?3)
1813                  ORDER BY t.created_at_ms, t.id
1814                  LIMIT 1",
1815                params![activation.project_id, activation.id, now_ms],
1816                |row| row.get(0),
1817            )
1818            .optional()?;
1819        Ok(ticket)
1820    }
1821
1822    pub fn ticket_is_dispatchable(&self, ticket_id: &str) -> Result<bool, StoreError> {
1823        self.connection
1824            .query_row(
1825                "SELECT EXISTS(
1826                     SELECT 1 FROM tickets t
1827                     WHERE t.id = ?1
1828                       AND t.state = 'ready'
1829                       AND t.missing_at_ms IS NULL
1830                       AND NOT EXISTS (SELECT 1 FROM ticket_blockers b
1831                                       JOIN tickets bt ON bt.id = b.blocker_id
1832                                       WHERE b.ticket_id = t.id
1833                                         AND bt.state != 'merged')
1834                 )",
1835                params![ticket_id],
1836                |row| row.get(0),
1837            )
1838            .map_err(StoreError::from)
1839    }
1840
1841    pub fn unmerged_blockers(&self, ticket_id: &str) -> Result<Vec<String>, StoreError> {
1842        let mut statement = self.connection.prepare(
1843            "SELECT b.blocker_id FROM ticket_blockers b
1844             JOIN tickets bt ON bt.id = b.blocker_id
1845             WHERE b.ticket_id = ?1 AND bt.state != 'merged'
1846             ORDER BY b.position, b.blocker_id",
1847        )?;
1848        statement
1849            .query_map(params![ticket_id], |row| row.get(0))?
1850            .collect::<Result<Vec<_>, _>>()
1851            .map_err(StoreError::from)
1852    }
1853
1854    /// Records a successful launch: the run turns `running` and carries the
1855    /// worktree, branch, and durable process identity.
1856    #[allow(clippy::too_many_arguments)]
1857    pub fn mark_run_running(
1858        &self,
1859        run_id: &str,
1860        branch: &str,
1861        worktree_path: &str,
1862        pid: u32,
1863        pid_start_time: Option<i64>,
1864        process_group_id: u32,
1865        worker_token: &str,
1866        worker_socket_path: &str,
1867        now_ms: i64,
1868    ) -> Result<(), StoreError> {
1869        let transaction = self.immediate_transaction()?;
1870        let changed = transaction.execute(
1871            "UPDATE runs
1872             SET state = ?10, branch = ?2, worktree_path = ?3, pid = ?4,
1873                 pid_start_time = ?5, process_group_id = ?6, worker_token = ?7,
1874                 worker_socket_path = ?8, started_at_ms = ?9, updated_at_ms = ?9
1875             WHERE id = ?1 AND state = ?11 AND exited_at_ms IS NULL",
1876            params![
1877                run_id,
1878                branch,
1879                worktree_path,
1880                i64::from(pid),
1881                pid_start_time,
1882                i64::from(process_group_id),
1883                worker_token,
1884                worker_socket_path,
1885                now_ms,
1886                RunState::Running.as_str(),
1887                RunState::Claimed.as_str(),
1888            ],
1889        )?;
1890        if changed != 1 {
1891            let state = transaction
1892                .query_row(
1893                    "SELECT state FROM runs WHERE id = ?1",
1894                    params![run_id],
1895                    |row| row.get(0),
1896                )
1897                .optional()?;
1898            return Err(StoreError::RunStateConflict {
1899                run_id: run_id.into(),
1900                state,
1901                requested: RunState::Running.as_str().into(),
1902            });
1903        }
1904        let ticket_id: String = transaction.query_row(
1905            "SELECT ticket_id FROM runs WHERE id = ?1",
1906            params![run_id],
1907            |row| row.get(0),
1908        )?;
1909        record_event(
1910            &transaction,
1911            now_ms,
1912            "run_started",
1913            Some(run_id),
1914            Some(&ticket_id),
1915            "{}",
1916        )?;
1917        transaction.commit()?;
1918        Ok(())
1919    }
1920
1921    /// Terminates a run in one transaction: the raw exit and derived outcome
1922    /// land on the run, evidence is appended, the lease is
1923    /// freed, and the ticket moves to its terminal state or back to `ready`
1924    /// when cancellation or recovery releases it.
1925    #[allow(clippy::too_many_arguments)]
1926    pub(crate) fn finish_run(
1927        &mut self,
1928        run_id: &str,
1929        ticket_id: &str,
1930        exit_code: Option<i32>,
1931        outcome: crate::outcome::Outcome,
1932        evidence: &[EvidenceRecord],
1933        cooldown: Option<&CooldownUpdate<'_>>,
1934        now_ms: i64,
1935    ) -> Result<bool, StoreError> {
1936        use crate::outcome::Outcome;
1937
1938        let transaction = self
1939            .connection
1940            .transaction_with_behavior(TransactionBehavior::Immediate)?;
1941        let run_state = RunState::from(outcome);
1942        let changed = transaction.execute(
1943            "UPDATE runs
1944             SET state = ?2, exited_at_ms = ?3, exit_code = ?4, updated_at_ms = ?3,
1945                 cleanup_eligible_at_ms = CASE WHEN ?2 = ?5 THEN ?3 ELSE NULL END
1946             WHERE id = ?1 AND exited_at_ms IS NULL",
1947            params![
1948                run_id,
1949                run_state.as_str(),
1950                now_ms,
1951                exit_code,
1952                RunState::Merged.as_str(),
1953            ],
1954        )?;
1955        if changed == 0 {
1956            let existing: Option<(String, Option<i64>)> = transaction
1957                .query_row(
1958                    "SELECT state, exited_at_ms FROM runs WHERE id = ?1",
1959                    params![run_id],
1960                    |row| Ok((row.get(0)?, row.get(1)?)),
1961                )
1962                .optional()?;
1963            match existing {
1964                Some((_, Some(_))) => {
1965                    transaction.commit()?;
1966                    return Ok(false);
1967                }
1968                Some((state, None)) => {
1969                    return Err(StoreError::RunStateConflict {
1970                        run_id: run_id.into(),
1971                        state: Some(state),
1972                        requested: run_state.as_str().into(),
1973                    });
1974                }
1975                None => {
1976                    return Err(StoreError::RunNotFound {
1977                        run_id: run_id.into(),
1978                    });
1979                }
1980            }
1981        }
1982        transaction.execute("DELETE FROM leases WHERE run_id = ?1", params![run_id])?;
1983
1984        let ticket_state = TicketState::after_outcome(outcome);
1985        transaction.execute(
1986            "UPDATE tickets SET state = ?2, held_reason = NULL, updated_at_ms = ?3
1987             WHERE id = ?1 AND state = 'claimed'",
1988            params![ticket_id, ticket_state.as_str(), now_ms],
1989        )?;
1990        if outcome == Outcome::RateLimited {
1991            transaction.execute(
1992                "UPDATE activations SET state = 'queued', updated_at_ms = ?2
1993                 WHERE id = (SELECT activation_id FROM runs WHERE id = ?1)",
1994                params![run_id, now_ms],
1995            )?;
1996        }
1997
1998        if let Some(cooldown) = cooldown {
1999            upsert_cooldown(&transaction, run_id, cooldown, now_ms)?;
2000        }
2001
2002        for record in evidence {
2003            transaction.execute(
2004                "INSERT OR IGNORE INTO run_evidence
2005                     (run_id, kind, observed_at_ms, dedupe_key, data_json)
2006                 VALUES (?1, ?2, ?3, 'settlement:' || ?1 || ':' || ?2, ?4)",
2007                params![run_id, record.kind, now_ms, record.data_json],
2008            )?;
2009        }
2010        record_event(
2011            &transaction,
2012            now_ms,
2013            "run_finished",
2014            Some(run_id),
2015            Some(ticket_id),
2016            &serde_json::json!({
2017                "outcome": outcome.as_str(),
2018                "exit_code": exit_code,
2019                "ticket_state": ticket_state.as_str(),
2020            })
2021            .to_string(),
2022        )?;
2023        transaction.commit()?;
2024        Ok(true)
2025    }
2026
2027    /// Records one completed flow stage. The flow index is the idempotency
2028    /// key, so recovery can re-derive the first stage still lacking a verdict.
2029    pub(crate) fn record_aftercare_stage(
2030        &self,
2031        run_id: &str,
2032        stage: &StageRecord,
2033    ) -> Result<(), StoreError> {
2034        let evidence_json = serde_json::json!({
2035            "output": stage.output_ref,
2036            "verdict_source": stage.verdict_source,
2037            "reason": stage.reason,
2038        })
2039        .to_string();
2040        self.connection.execute(
2041            "INSERT INTO aftercare_stages
2042                 (run_id, stage_index, stage, state, started_at_ms, finished_at_ms, exit_code,
2043                  evidence_json)
2044             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
2045             ON CONFLICT(run_id, stage_index, attempt) DO UPDATE SET
2046                 stage = excluded.stage,
2047                 state = excluded.state,
2048                 started_at_ms = excluded.started_at_ms,
2049                 finished_at_ms = excluded.finished_at_ms,
2050                 exit_code = excluded.exit_code,
2051                 evidence_json = excluded.evidence_json",
2052            params![
2053                run_id,
2054                stage.stage_index as i64,
2055                stage.stage,
2056                stage.state,
2057                stage.started_at_ms,
2058                stage.finished_at_ms,
2059                stage.exit_code,
2060                evidence_json,
2061            ],
2062        )?;
2063        Ok(())
2064    }
2065
2066    pub(crate) fn aftercare_stages(&self, run_id: &str) -> Result<Vec<StageRecord>, StoreError> {
2067        let mut statement = self.connection.prepare(
2068            "SELECT stage_index, stage, state, started_at_ms, finished_at_ms, exit_code,
2069                    evidence_json
2070             FROM aftercare_stages WHERE run_id = ?1 ORDER BY stage_index",
2071        )?;
2072        statement
2073            .query_map(params![run_id], |row| {
2074                let evidence_json: Option<String> = row.get(6)?;
2075                let output_ref = evidence_json
2076                    .as_deref()
2077                    .and_then(|value| serde_json::from_str::<serde_json::Value>(value).ok())
2078                    .and_then(|value| value["output"].as_str().map(str::to_owned))
2079                    .unwrap_or_default();
2080                let evidence = evidence_json
2081                    .as_deref()
2082                    .and_then(|value| serde_json::from_str::<serde_json::Value>(value).ok());
2083                Ok(StageRecord {
2084                    stage_index: row.get::<_, i64>(0)? as usize,
2085                    stage: row.get(1)?,
2086                    state: row.get(2)?,
2087                    started_at_ms: row.get(3)?,
2088                    finished_at_ms: row.get(4)?,
2089                    exit_code: row.get(5)?,
2090                    output_ref,
2091                    verdict_source: evidence
2092                        .as_ref()
2093                        .and_then(|value| value["verdict_source"].as_str())
2094                        .unwrap_or("exit_code")
2095                        .to_owned(),
2096                    reason: evidence
2097                        .as_ref()
2098                        .and_then(|value| value["reason"].as_str())
2099                        .map(str::to_owned),
2100                })
2101            })?
2102            .collect::<Result<Vec<_>, _>>()
2103            .map_err(StoreError::from)
2104    }
2105
2106    /// Checkpoints the agent's exit before aftercare starts. The lease and
2107    /// ticket remain claimed until final settlement, but recovery can now
2108    /// resume with the exact exit and branch-activity facts. Only the caller that
2109    /// wins this transition owns exit processing and aftercare for the run.
2110    #[allow(clippy::too_many_arguments)]
2111    pub(crate) fn record_agent_exit(
2112        &mut self,
2113        run_id: &str,
2114        exit_code: Option<i32>,
2115        capture_complete: bool,
2116        commits_json: &str,
2117        vendor_error: Option<&crate::vendor_error::VendorErrorMatch>,
2118        cooldown_until_ms: Option<i64>,
2119        now_ms: i64,
2120    ) -> Result<ExitClaim, StoreError> {
2121        let transaction = self
2122            .connection
2123            .transaction_with_behavior(TransactionBehavior::Immediate)?;
2124        let changed = transaction.execute(
2125            "UPDATE runs
2126             SET state = ?4, exit_code = ?2, updated_at_ms = ?3
2127             WHERE id = ?1 AND state = ?5 AND exited_at_ms IS NULL",
2128            params![
2129                run_id,
2130                exit_code,
2131                now_ms,
2132                RunState::Aftercare.as_str(),
2133                RunState::Running.as_str(),
2134            ],
2135        )?;
2136        if changed == 0 {
2137            let state: Option<String> = transaction
2138                .query_row(
2139                    "SELECT state FROM runs WHERE id = ?1",
2140                    params![run_id],
2141                    |row| row.get(0),
2142                )
2143                .optional()?;
2144            return match state {
2145                Some(state) => Ok(ExitClaim::AlreadyClaimed { state }),
2146                None => Err(StoreError::RunNotFound {
2147                    run_id: run_id.into(),
2148                }),
2149            };
2150        }
2151        for (kind, data_json) in [
2152            (
2153                "exit_classified",
2154                serde_json::json!({"exit_code": exit_code}).to_string(),
2155            ),
2156            ("commits_observed", commits_json.to_owned()),
2157        ] {
2158            transaction.execute(
2159                "INSERT OR IGNORE INTO run_evidence
2160                     (run_id, kind, observed_at_ms, dedupe_key, data_json)
2161                 VALUES (?1, ?2, ?3, 'settlement:' || ?1 || ':' || ?2, ?4)",
2162                params![run_id, kind, now_ms, data_json],
2163            )?;
2164        }
2165        if let Some(vendor_error) = vendor_error {
2166            transaction.execute(
2167                "INSERT OR IGNORE INTO run_evidence
2168                     (run_id, kind, observed_at_ms, dedupe_key, data_json)
2169                 VALUES (?1, 'vendor_error_classified', ?2,
2170                         'settlement:' || ?1 || ':vendor_error_classified', ?3)",
2171                params![
2172                    run_id,
2173                    now_ms,
2174                    vendor_error.evidence_json(cooldown_until_ms)
2175                ],
2176            )?;
2177        }
2178        if !capture_complete {
2179            transaction.execute(
2180                "INSERT OR IGNORE INTO run_evidence
2181                     (run_id, kind, observed_at_ms, dedupe_key, data_json)
2182                 VALUES (?1, 'capture_incomplete', ?2,
2183                         'settlement:' || ?1 || ':capture_incomplete', '{}')",
2184                params![run_id, now_ms],
2185            )?;
2186        }
2187        transaction.commit()?;
2188        Ok(ExitClaim::Claimed)
2189    }
2190
2191    pub(crate) fn record_aftercare_evidence(
2192        &self,
2193        run_id: &str,
2194        kind: &str,
2195        data_json: &str,
2196        now_ms: i64,
2197    ) -> Result<(), StoreError> {
2198        self.connection.execute(
2199            "INSERT INTO run_evidence
2200                 (run_id, kind, observed_at_ms, dedupe_key, data_json)
2201             VALUES (?1, ?2, ?3, 'settlement:' || ?1 || ':' || ?2, ?4)
2202             ON CONFLICT(dedupe_key) DO UPDATE SET
2203                 observed_at_ms = excluded.observed_at_ms,
2204                 data_json = excluded.data_json",
2205            params![run_id, kind, now_ms, data_json],
2206        )?;
2207        Ok(())
2208    }
2209
2210    pub(crate) fn clear_aftercare_process(&self, run_id: &str) -> Result<(), StoreError> {
2211        self.connection.execute(
2212            "DELETE FROM run_evidence
2213             WHERE run_id = ?1 AND dedupe_key = 'settlement:' || ?1 || ':aftercare_process'",
2214            params![run_id],
2215        )?;
2216        Ok(())
2217    }
2218
2219    /// Durably records an operator's cancellation intent, idempotently: the
2220    /// dedupe key makes a repeated `cancel` a no-op rather than new evidence.
2221    pub fn record_cancel_requested(&self, run_id: &str, now_ms: i64) -> Result<(), StoreError> {
2222        self.connection.execute(
2223            "INSERT OR IGNORE INTO run_evidence
2224                 (run_id, kind, observed_at_ms, dedupe_key, data_json)
2225             VALUES (?1, 'cancel_requested', ?2, 'cancel_requested:' || ?1, '{}')",
2226            params![run_id, now_ms],
2227        )?;
2228        Ok(())
2229    }
2230
2231    /// Whether cancellation intent was recorded for the run, so an exit event
2232    /// racing the cancel still resolves to `Cancelled`.
2233    pub fn cancellation_requested(&self, run_id: &str) -> Result<bool, StoreError> {
2234        let found: Option<i64> = self
2235            .connection
2236            .query_row(
2237                "SELECT 1 FROM run_evidence
2238                 WHERE run_id = ?1 AND kind = 'cancel_requested'",
2239                params![run_id],
2240                |row| row.get(0),
2241            )
2242            .optional()?;
2243        Ok(found.is_some())
2244    }
2245
2246    /// Appends a worker's advisory note. The agent's only write: it records
2247    /// text against the run and moves nothing.
2248    pub fn insert_note(
2249        &self,
2250        id: &str,
2251        run_id: &str,
2252        text: &str,
2253        now_ms: i64,
2254    ) -> Result<(), StoreError> {
2255        self.connection.execute(
2256            "INSERT INTO notes (id, run_id, text, recorded_at_ms)
2257             VALUES (?1, ?2, ?3, ?4)",
2258            params![id, run_id, text, now_ms],
2259        )?;
2260        Ok(())
2261    }
2262
2263    /// Notes recorded against one run, in the order they arrived.
2264    pub fn notes_for_run(&self, run_id: &str) -> Result<Vec<String>, StoreError> {
2265        let mut statement = self
2266            .connection
2267            .prepare("SELECT text FROM notes WHERE run_id = ?1 ORDER BY recorded_at_ms, id")?;
2268        let rows = statement
2269            .query_map(params![run_id], |row| row.get(0))?
2270            .collect::<Result<Vec<_>, _>>()?;
2271        Ok(rows)
2272    }
2273
2274    /// Records the first worker-reported verdict for one stage. The unique
2275    /// dedupe key is the at-most-once gate; later reports cannot overwrite it.
2276    pub(crate) fn record_stage_verdict(
2277        &self,
2278        run_id: &str,
2279        stage: &str,
2280        verdict: &str,
2281        reason: Option<&str>,
2282        now_ms: i64,
2283    ) -> Result<bool, StoreError> {
2284        let dedupe_key = format!("verdict:{run_id}:{stage}");
2285        let data_json =
2286            serde_json::json!({"stage": stage, "verdict": verdict, "reason": reason}).to_string();
2287        let inserted = self.connection.execute(
2288            "INSERT OR IGNORE INTO run_evidence
2289                 (run_id, kind, observed_at_ms, dedupe_key, data_json)
2290             VALUES (?1, 'stage_verdict', ?2, ?3, ?4)",
2291            params![run_id, now_ms, dedupe_key, data_json],
2292        )?;
2293        Ok(inserted == 1)
2294    }
2295
2296    pub fn notes_for_project(&self, project_id: &str) -> Result<Vec<ProjectNote>, StoreError> {
2297        let mut statement = self.connection.prepare(
2298            "SELECT n.id, n.run_id, r.ticket_id, n.text, n.recorded_at_ms
2299             FROM notes n
2300             JOIN runs r ON r.id = n.run_id
2301             JOIN tickets t ON t.id = r.ticket_id
2302             WHERE t.project_id = ?1
2303             ORDER BY r.ticket_id, n.recorded_at_ms, n.id",
2304        )?;
2305        statement
2306            .query_map(params![project_id], |row| {
2307                Ok(ProjectNote {
2308                    id: row.get(0)?,
2309                    run_id: row.get(1)?,
2310                    ticket_id: row.get(2)?,
2311                    text: row.get(3)?,
2312                    recorded_at_ms: row.get(4)?,
2313                })
2314            })?
2315            .collect::<Result<Vec<_>, _>>()
2316            .map_err(StoreError::from)
2317    }
2318
2319    pub fn commit_evidence_for_project(
2320        &self,
2321        project_id: &str,
2322    ) -> Result<Vec<ProjectCommitEvidence>, StoreError> {
2323        let mut statement = self.connection.prepare(
2324            "SELECT r.id, r.ticket_id, e.data_json
2325             FROM run_evidence e
2326             JOIN runs r ON r.id = e.run_id
2327             JOIN tickets t ON t.id = r.ticket_id
2328             WHERE t.project_id = ?1 AND e.kind = 'commits_observed'
2329             ORDER BY r.ticket_id, r.created_at_ms, r.id, e.sequence",
2330        )?;
2331        statement
2332            .query_map(params![project_id], |row| {
2333                Ok(ProjectCommitEvidence {
2334                    run_id: row.get(0)?,
2335                    ticket_id: row.get(1)?,
2336                    data_json: row.get(2)?,
2337                })
2338            })?
2339            .collect::<Result<Vec<_>, _>>()
2340            .map_err(StoreError::from)
2341    }
2342
2343    pub fn next_note_ordinal(&self) -> Result<i64, StoreError> {
2344        self.reserve_ordinal("note", "notes")
2345    }
2346
2347    /// Evidence rows for one run in observation order, as (kind, data_json).
2348    pub fn run_evidence(&self, run_id: &str) -> Result<Vec<(String, String)>, StoreError> {
2349        let mut statement = self.connection.prepare(
2350            "SELECT kind, data_json FROM run_evidence WHERE run_id = ?1 ORDER BY sequence",
2351        )?;
2352        let rows = statement
2353            .query_map(params![run_id], |row| Ok((row.get(0)?, row.get(1)?)))?
2354            .collect::<Result<Vec<_>, _>>()?;
2355        Ok(rows)
2356    }
2357
2358    pub fn vendor_error_for_run(
2359        &self,
2360        run_id: &str,
2361    ) -> Result<Option<crate::vendor_error::VendorErrorMatch>, StoreError> {
2362        let data: Option<String> = self
2363            .connection
2364            .query_row(
2365                "SELECT data_json FROM run_evidence
2366                 WHERE run_id = ?1 AND kind = 'vendor_error_classified'
2367                 ORDER BY sequence DESC LIMIT 1",
2368                params![run_id],
2369                |row| row.get(0),
2370            )
2371            .optional()?;
2372        Ok(data.and_then(|data| serde_json::from_str(&data).ok()))
2373    }
2374
2375    pub fn latest_vendor_error_for_ticket(
2376        &self,
2377        ticket_id: &str,
2378    ) -> Result<Option<crate::vendor_error::VendorErrorMatch>, StoreError> {
2379        let data: Option<String> = self
2380            .connection
2381            .query_row(
2382                "SELECT e.data_json FROM run_evidence e
2383                 JOIN runs r ON r.id = e.run_id
2384                 WHERE r.id = (SELECT latest.id FROM runs latest
2385                               WHERE latest.ticket_id = ?1
2386                               ORDER BY latest.created_at_ms DESC, latest.id DESC LIMIT 1)
2387                   AND e.kind = 'vendor_error_classified'
2388                 ORDER BY e.sequence DESC LIMIT 1",
2389                params![ticket_id],
2390                |row| row.get(0),
2391            )
2392            .optional()?;
2393        Ok(data.and_then(|data| serde_json::from_str(&data).ok()))
2394    }
2395
2396    /// Rolls back a claim whose launch failed before a process existed: the
2397    /// lease is released, the run is closed, and the ticket returns to
2398    /// `ready`. The consumed attempt is kept as evidence of the try.
2399    pub(crate) fn abort_claim(
2400        &mut self,
2401        run_id: &str,
2402        ticket_id: &str,
2403        now_ms: i64,
2404    ) -> Result<(), StoreError> {
2405        let transaction = self
2406            .connection
2407            .transaction_with_behavior(TransactionBehavior::Immediate)?;
2408        transaction.execute("DELETE FROM leases WHERE run_id = ?1", params![run_id])?;
2409        transaction.execute(
2410            "UPDATE runs
2411             SET state = ?3, exited_at_ms = ?2, updated_at_ms = ?2
2412             WHERE id = ?1 AND exited_at_ms IS NULL",
2413            params![run_id, now_ms, RunState::Aborted.as_str()],
2414        )?;
2415        transaction.execute(
2416            "UPDATE tickets SET state = 'ready', held_reason = NULL, updated_at_ms = ?2
2417             WHERE id = ?1 AND state = 'claimed'",
2418            params![ticket_id, now_ms],
2419        )?;
2420        record_event(
2421            &transaction,
2422            now_ms,
2423            "run_aborted",
2424            Some(run_id),
2425            Some(ticket_id),
2426            "{}",
2427        )?;
2428        transaction.commit()?;
2429        Ok(())
2430    }
2431
2432    /// Reads activity-feed rows with `sequence > after`, oldest first. The
2433    /// last row's sequence is the caller's next cursor.
2434    pub fn events_after(&self, after: i64, limit: usize) -> Result<Vec<EventRecord>, StoreError> {
2435        let mut statement = self.connection.prepare(
2436            "SELECT sequence, occurred_at_ms, kind, run_id, ticket_id, data_json
2437             FROM events WHERE sequence > ?1 ORDER BY sequence LIMIT ?2",
2438        )?;
2439        statement
2440            .query_map(params![after, limit as i64], |row| {
2441                Ok(EventRecord {
2442                    sequence: row.get(0)?,
2443                    occurred_at_ms: row.get(1)?,
2444                    kind: row.get(2)?,
2445                    run_id: row.get(3)?,
2446                    ticket_id: row.get(4)?,
2447                    data_json: row.get(5)?,
2448                })
2449            })?
2450            .collect::<Result<Vec<_>, _>>()
2451            .map_err(StoreError::from)
2452    }
2453
2454    /// The claim/start/finish instants of each named run, read back out of the
2455    /// activity feed. The feed is written in the same transaction as the
2456    /// transition it narrates, so these are the authoritative wall-clock
2457    /// boundaries of a run — nothing extra is stored to render them.
2458    ///
2459    /// A run with no finish row is still in flight; callers render it
2460    /// open-ended rather than inventing an end. `run_aborted` counts as a
2461    /// finish because it is the terminal row for runs that never settled.
2462    pub fn run_timelines(
2463        &self,
2464        run_ids: &[&str],
2465    ) -> Result<HashMap<String, RunTimeline>, StoreError> {
2466        if run_ids.is_empty() {
2467            return Ok(HashMap::new());
2468        }
2469        let placeholders = std::iter::repeat_n("?", run_ids.len())
2470            .collect::<Vec<_>>()
2471            .join(", ");
2472        let mut statement = self.connection.prepare(&format!(
2473            "SELECT run_id,
2474                    MIN(CASE WHEN kind = 'run_claimed' THEN occurred_at_ms END),
2475                    MIN(CASE WHEN kind = 'run_started' THEN occurred_at_ms END),
2476                    MAX(CASE WHEN kind IN ('run_finished', 'run_aborted')
2477                             THEN occurred_at_ms END)
2478             FROM events
2479             WHERE run_id IN ({placeholders})
2480             GROUP BY run_id"
2481        ))?;
2482        let rows = statement
2483            .query_map(rusqlite::params_from_iter(run_ids), |row| {
2484                Ok((
2485                    row.get::<_, String>(0)?,
2486                    RunTimeline {
2487                        claimed_at_ms: row.get(1)?,
2488                        started_at_ms: row.get(2)?,
2489                        finished_at_ms: row.get(3)?,
2490                    },
2491                ))
2492            })?
2493            .collect::<Result<HashMap<_, _>, _>>()?;
2494        Ok(rows)
2495    }
2496
2497    pub fn latest_event_sequence(&self) -> Result<i64, StoreError> {
2498        let latest = self.connection.query_row(
2499            "SELECT COALESCE(MAX(sequence), 0) FROM events",
2500            [],
2501            |row| row.get(0),
2502        )?;
2503        Ok(latest)
2504    }
2505
2506    /// Drops all but the newest `keep` activity-feed rows. Sequences are never
2507    /// reused after a trim, so cursors held by watchers stay valid.
2508    pub fn trim_events(&self, keep: i64) -> Result<(), StoreError> {
2509        self.connection.execute(
2510            "DELETE FROM events
2511             WHERE sequence <= (SELECT COALESCE(MAX(sequence), 0) FROM events) - ?1",
2512            params![keep],
2513        )?;
2514        Ok(())
2515    }
2516
2517    pub fn run(&self, id: &str) -> Result<Option<RunRecord>, StoreError> {
2518        let run = self
2519            .connection
2520            .query_row(
2521                &format!("{RUN_RECORD_SELECT} WHERE id = ?1"),
2522                params![id],
2523                run_record,
2524            )
2525            .optional()?;
2526        Ok(run)
2527    }
2528
2529    /// The run a `<ticket>-r<attempt>` alias names. The pair is unique because
2530    /// attempts are allocated once per ticket at claim time.
2531    pub fn run_for_ticket_attempt(
2532        &self,
2533        ticket_id: &str,
2534        attempt: i64,
2535    ) -> Result<Option<RunRecord>, StoreError> {
2536        let run = self
2537            .connection
2538            .query_row(
2539                &format!(
2540                    "{RUN_RECORD_SELECT} WHERE ticket_id = ?1 AND attempt = ?2
2541                     ORDER BY created_at_ms DESC LIMIT 1"
2542                ),
2543                params![ticket_id, attempt],
2544                run_record,
2545            )
2546            .optional()?;
2547        Ok(run)
2548    }
2549
2550    /// Every run a ticket has produced, newest attempt first, so a bare ticket
2551    /// reference can name the latest run and still report the earlier ones.
2552    pub fn runs_for_ticket(&self, ticket_id: &str) -> Result<Vec<RunRecord>, StoreError> {
2553        let mut statement = self.connection.prepare(&format!(
2554            "{RUN_RECORD_SELECT} WHERE ticket_id = ?1 ORDER BY attempt DESC, created_at_ms DESC"
2555        ))?;
2556        let runs = statement
2557            .query_map(params![ticket_id], run_record)?
2558            .collect::<Result<Vec<_>, _>>()?;
2559        Ok(runs)
2560    }
2561
2562    /// Runs whose internal id starts with `prefix`. More than one row means the
2563    /// reference is ambiguous, so the caller needs the candidates, not a pick.
2564    pub fn runs_with_id_prefix(&self, prefix: &str) -> Result<Vec<RunRecord>, StoreError> {
2565        // `LIKE` would treat `%` and `_` in a prefix as wildcards; run ids are
2566        // hexadecimal, but comparing on the substring keeps that beyond doubt.
2567        let mut statement = self.connection.prepare(&format!(
2568            "{RUN_RECORD_SELECT} WHERE SUBSTR(id, 1, ?2) = ?1 ORDER BY created_at_ms, id"
2569        ))?;
2570        let runs = statement
2571            .query_map(params![prefix, prefix.len() as i64], run_record)?
2572            .collect::<Result<Vec<_>, _>>()?;
2573        Ok(runs)
2574    }
2575
2576    /// Every `needs_review` ticket paired with the branch of the run that
2577    /// produced it, so the daemon can freshly test each branch tip for external
2578    /// integration. Only the newest `needs_review` run with a branch is
2579    /// returned per ticket; the tip itself is never cached here.
2580    pub(crate) fn needs_review_branches(&self) -> Result<Vec<NeedsReviewBranch>, StoreError> {
2581        let mut statement = self.connection.prepare(
2582            "SELECT t.id, r.id, r.branch
2583             FROM tickets t
2584             JOIN runs r ON r.id = (
2585                 SELECT r2.id FROM runs r2
2586                 WHERE r2.ticket_id = t.id
2587                   AND r2.state = 'needs_review'
2588                   AND r2.branch IS NOT NULL
2589                 ORDER BY r2.created_at_ms DESC, r2.id DESC
2590                 LIMIT 1
2591             )
2592             WHERE t.state = 'needs_review'",
2593        )?;
2594        statement
2595            .query_map([], |row| {
2596                Ok(NeedsReviewBranch {
2597                    ticket_id: row.get(0)?,
2598                    run_id: row.get(1)?,
2599                    branch: row.get(2)?,
2600                })
2601            })?
2602            .collect::<Result<Vec<_>, _>>()
2603            .map_err(StoreError::from)
2604    }
2605
2606    pub(crate) fn worktree_cleanup_candidates(
2607        &self,
2608    ) -> Result<Vec<WorktreeCleanupCandidate>, StoreError> {
2609        let mut statement = self.connection.prepare(
2610            "SELECT r.id, r.ticket_id, r.branch, r.worktree_path, r.cleanup_eligible_at_ms
2611             FROM runs r
2612             WHERE r.cleanup_eligible_at_ms IS NOT NULL
2613               AND r.cleaned_at_ms IS NULL
2614               AND r.branch IS NOT NULL
2615               AND r.worktree_path IS NOT NULL
2616               AND NOT EXISTS (SELECT 1 FROM leases l WHERE l.run_id = r.id)
2617             ORDER BY r.cleanup_eligible_at_ms, r.id",
2618        )?;
2619        statement
2620            .query_map([], |row| {
2621                Ok(WorktreeCleanupCandidate {
2622                    run_id: row.get(0)?,
2623                    ticket_id: row.get(1)?,
2624                    branch: row.get(2)?,
2625                    worktree_path: row.get(3)?,
2626                    cleanup_eligible_at_ms: row.get(4)?,
2627                })
2628            })?
2629            .collect::<Result<Vec<_>, _>>()
2630            .map_err(StoreError::from)
2631    }
2632
2633    pub(crate) fn next_worktree_cleanup_at_ms(
2634        &self,
2635        retention_ms: i64,
2636        now_ms: i64,
2637    ) -> Result<Option<i64>, StoreError> {
2638        let eligible_at: Option<i64> = self.connection.query_row(
2639            "SELECT MIN(r.cleanup_eligible_at_ms)
2640             FROM runs r
2641             WHERE r.cleanup_eligible_at_ms IS NOT NULL
2642               AND r.cleaned_at_ms IS NULL
2643               AND r.branch IS NOT NULL
2644               AND r.worktree_path IS NOT NULL
2645               AND NOT EXISTS (SELECT 1 FROM leases l WHERE l.run_id = r.id)",
2646            [],
2647            |row| row.get(0),
2648        )?;
2649        Ok(eligible_at.and_then(|value| {
2650            let deadline = value.saturating_add(retention_ms);
2651            (deadline > now_ms).then_some(deadline)
2652        }))
2653    }
2654
2655    pub(crate) fn mark_run_worktree_cleaned(
2656        &self,
2657        candidate: &WorktreeCleanupCandidate,
2658        now_ms: i64,
2659    ) -> Result<bool, StoreError> {
2660        let transaction = self.immediate_transaction()?;
2661        let changed = transaction.execute(
2662            "UPDATE runs SET cleaned_at_ms = ?2, updated_at_ms = ?2
2663             WHERE id = ?1 AND cleanup_eligible_at_ms IS NOT NULL
2664               AND cleaned_at_ms IS NULL
2665               AND NOT EXISTS (SELECT 1 FROM leases l WHERE l.run_id = runs.id)",
2666            params![candidate.run_id, now_ms],
2667        )?;
2668        if changed == 1 {
2669            record_event(
2670                &transaction,
2671                now_ms,
2672                "run_worktree_cleaned",
2673                Some(&candidate.run_id),
2674                Some(&candidate.ticket_id),
2675                &serde_json::json!({
2676                    "branch": candidate.branch,
2677                    "worktree": candidate.worktree_path,
2678                })
2679                .to_string(),
2680            )?;
2681        }
2682        transaction.commit()?;
2683        Ok(changed == 1)
2684    }
2685
2686    /// Settles a `needs_review` ticket whose run branch an operator merged by
2687    /// hand: the ticket becomes `merged`, releasing its `blocked_by` dependents
2688    /// exactly as a flow merge would, and the observation is recorded as
2689    /// evidence. The ticket-state gate makes a repeated pass a no-op and the
2690    /// `dedupe_key` UNIQUE gate keeps the evidence row unique across restarts.
2691    /// Returns whether this call performed the transition.
2692    pub(crate) fn settle_external_merge(
2693        &mut self,
2694        run_id: &str,
2695        ticket_id: &str,
2696        branch: &str,
2697        branch_tip: &str,
2698        observed_default_tip: &str,
2699        now_ms: i64,
2700    ) -> Result<bool, StoreError> {
2701        let transaction = self
2702            .connection
2703            .transaction_with_behavior(TransactionBehavior::Immediate)?;
2704        let changed = transaction.execute(
2705            "UPDATE tickets SET state = 'merged', held_reason = NULL, updated_at_ms = ?2
2706             WHERE id = ?1 AND state = 'needs_review'",
2707            params![ticket_id, now_ms],
2708        )?;
2709        if changed == 0 {
2710            transaction.commit()?;
2711            return Ok(false);
2712        }
2713        transaction.execute(
2714            "UPDATE runs SET cleanup_eligible_at_ms = ?2
2715             WHERE id = ?1 AND cleanup_eligible_at_ms IS NULL AND cleaned_at_ms IS NULL",
2716            params![run_id, now_ms],
2717        )?;
2718        let data_json = serde_json::json!({
2719            "branch": branch,
2720            "branch_tip": branch_tip,
2721            "observed_default_tip": observed_default_tip,
2722        })
2723        .to_string();
2724        transaction.execute(
2725            "INSERT OR IGNORE INTO run_evidence
2726                 (run_id, kind, observed_at_ms, dedupe_key, data_json)
2727             VALUES (?1, 'external_merge_observed', ?2, 'external_merge:' || ?1, ?3)",
2728            params![run_id, now_ms, data_json],
2729        )?;
2730        record_event(
2731            &transaction,
2732            now_ms,
2733            "external_merge_reconciled",
2734            Some(run_id),
2735            Some(ticket_id),
2736            &data_json,
2737        )?;
2738        transaction.commit()?;
2739        Ok(true)
2740    }
2741
2742    /// The ticket's live run as `(id, attempt)`. The attempt travels with the
2743    /// id so callers can name the run by alias without re-reading the row.
2744    pub fn active_run_for_ticket(
2745        &self,
2746        ticket_id: &str,
2747    ) -> Result<Option<(String, i64)>, StoreError> {
2748        let run = self
2749            .connection
2750            .query_row(
2751                "SELECT r.id, r.attempt FROM runs r
2752                 JOIN leases l ON l.run_id = r.id
2753                 WHERE r.ticket_id = ?1
2754                   AND r.state IN (?2, ?3, ?4)
2755                   AND r.exited_at_ms IS NULL
2756                 ORDER BY r.created_at_ms DESC, r.id DESC LIMIT 1",
2757                params![
2758                    ticket_id,
2759                    NONTERMINAL_RUN_STATES[0].as_str(),
2760                    NONTERMINAL_RUN_STATES[1].as_str(),
2761                    NONTERMINAL_RUN_STATES[2].as_str(),
2762                ],
2763                |row| Ok((row.get(0)?, row.get(1)?)),
2764            )
2765            .optional()?;
2766        Ok(run)
2767    }
2768
2769    /// Leased nonterminal runs that consume capacity, oldest first.
2770    pub fn active_runs(&self) -> Result<Vec<ActiveRun>, StoreError> {
2771        let mut statement = self.connection.prepare(
2772            "SELECT r.id, r.ticket_id, r.attempt, t.name, t.project_id, r.state FROM runs r
2773             JOIN leases l ON l.run_id = r.id
2774             JOIN tickets t ON t.id = r.ticket_id
2775             WHERE r.exited_at_ms IS NULL
2776               AND r.state IN (?1, ?2, ?3)
2777             ORDER BY r.created_at_ms, r.id",
2778        )?;
2779        let runs = statement
2780            .query_map(nonterminal_state_params(), |row| {
2781                Ok(ActiveRun {
2782                    id: row.get(0)?,
2783                    ticket_id: row.get(1)?,
2784                    attempt: row.get(2)?,
2785                    ticket_name: row.get(3)?,
2786                    project_id: row.get(4)?,
2787                    state: row.get(5)?,
2788                })
2789            })?
2790            .collect::<Result<Vec<_>, _>>()?;
2791        Ok(runs)
2792    }
2793
2794    /// Every nonterminal run that still owns a lease, oldest first. Startup
2795    /// must classify all of these before making another spawn decision.
2796    pub(crate) fn recoverable_runs(&self) -> Result<Vec<RecoverableRun>, StoreError> {
2797        let mut statement = self.connection.prepare(
2798            "SELECT r.id, r.ticket_id, t.target, r.state, r.branch, r.worktree_path,
2799                    r.pid, r.pid_start_time, r.process_group_id, r.worker_token,
2800                    r.worker_socket_path, r.exit_code, l.expires_at_ms, r.flow_json,
2801                    r.ticket_json
2802             FROM runs r
2803             JOIN leases l ON l.run_id = r.id
2804             JOIN tickets t ON t.id = r.ticket_id
2805             WHERE r.exited_at_ms IS NULL
2806               AND r.state IN (?1, ?2, ?3)
2807             ORDER BY r.created_at_ms, r.id",
2808        )?;
2809        statement
2810            .query_map(nonterminal_state_params(), |row| {
2811                Ok(RecoverableRun {
2812                    id: row.get(0)?,
2813                    ticket_id: row.get(1)?,
2814                    target: row.get(2)?,
2815                    state: row.get(3)?,
2816                    branch: row.get(4)?,
2817                    worktree_path: row.get(5)?,
2818                    pid: row.get(6)?,
2819                    pid_start_time: row.get(7)?,
2820                    process_group_id: row.get(8)?,
2821                    worker_token: row.get(9)?,
2822                    worker_socket_path: row.get(10)?,
2823                    exit_code: row.get(11)?,
2824                    lease_expires_at_ms: row.get(12)?,
2825                    flow_json: row.get(13)?,
2826                    ticket_json: row.get(14)?,
2827                })
2828            })?
2829            .collect::<Result<Vec<_>, _>>()
2830            .map_err(StoreError::from)
2831    }
2832
2833    /// Finds a still-queued activation of `kind` scoped to one ticket, used
2834    /// to keep reposting the same file idempotent.
2835    pub fn queued_ticket_activation(
2836        &self,
2837        ticket_id: &str,
2838        kind: ActivationKind,
2839    ) -> Result<Option<String>, StoreError> {
2840        let id = self
2841            .connection
2842            .query_row(
2843                "SELECT id FROM activations
2844                 WHERE ticket_id = ?1 AND kind = ?2 AND state = 'queued'
2845                 ORDER BY created_at_ms LIMIT 1",
2846                params![ticket_id, kind.as_str()],
2847                |row| row.get(0),
2848            )
2849            .optional()?;
2850        Ok(id)
2851    }
2852
2853    /// Moves a still-queued timed activation to a new eligibility instant,
2854    /// so reposting a ticket with a different `--at` time reschedules the
2855    /// existing activation instead of queueing a duplicate.
2856    pub fn reschedule_activation(
2857        &self,
2858        id: &str,
2859        eligible_at_ms: i64,
2860        now_ms: i64,
2861    ) -> Result<(), StoreError> {
2862        self.connection.execute(
2863            "UPDATE activations
2864             SET eligible_at_ms = ?2, updated_at_ms = ?3
2865             WHERE id = ?1 AND state = 'queued'",
2866            params![id, eligible_at_ms, now_ms],
2867        )?;
2868        Ok(())
2869    }
2870
2871    /// Reserves the next activation ordinal without reusing IDs removed by
2872    /// reindex.
2873    pub fn next_activation_ordinal(&self) -> Result<i64, StoreError> {
2874        self.reserve_ordinal("activation", "activations")
2875    }
2876
2877    /// Begins a write transaction that takes the write lock up front. A
2878    /// deferred transaction that reads before its first write can hit an
2879    /// immediate `SQLITE_BUSY` when a sibling connection commits in between:
2880    /// its snapshot is stale, so SQLite fails fast instead of honoring
2881    /// `busy_timeout`. Starting immediate makes contending writers queue.
2882    fn immediate_transaction(&self) -> Result<rusqlite::Transaction<'_>, rusqlite::Error> {
2883        rusqlite::Transaction::new_unchecked(&self.connection, TransactionBehavior::Immediate)
2884    }
2885
2886    fn reserve_ordinal(&self, kind: &str, table: &str) -> Result<i64, StoreError> {
2887        let transaction = self.immediate_transaction()?;
2888        let reserved: i64 = transaction.query_row(
2889            "SELECT next_ordinal FROM id_counters WHERE kind = ?1",
2890            params![kind],
2891            |row| row.get(0),
2892        )?;
2893        let existing: i64 = transaction.query_row(
2894            &format!("SELECT COALESCE(MAX(CAST(SUBSTR(id, 2) AS INTEGER)), 0) + 1 FROM {table}"),
2895            [],
2896            |row| row.get(0),
2897        )?;
2898        let ordinal = reserved.max(existing);
2899        transaction.execute(
2900            "UPDATE id_counters SET next_ordinal = ?2 WHERE kind = ?1",
2901            params![kind, ordinal + 1],
2902        )?;
2903        transaction.commit()?;
2904        Ok(ordinal)
2905    }
2906
2907    /// Claims a ready ticket for one run in a single transaction. The
2908    /// conditional update plus the primary key on `leases.ticket_id` are the
2909    /// durable guards against a double claim.
2910    pub(crate) fn claim_ticket(
2911        &mut self,
2912        claim: &ClaimRequest<'_>,
2913        now_ms: i64,
2914    ) -> Result<ClaimedRun, StoreError> {
2915        let transaction = self
2916            .connection
2917            .transaction_with_behavior(TransactionBehavior::Immediate)?;
2918
2919        let changed = transaction.execute(
2920            "UPDATE tickets
2921             SET state = 'claimed', held_reason = NULL, attempts = attempts + 1, updated_at_ms = ?2
2922             WHERE id = ?1 AND state = 'ready' AND missing_at_ms IS NULL
2923               AND NOT EXISTS (SELECT 1 FROM ticket_blockers b
2924                               JOIN tickets bt ON bt.id = b.blocker_id
2925                               WHERE b.ticket_id = tickets.id
2926                                 AND bt.state != 'merged')",
2927            params![claim.ticket_id, now_ms],
2928        )?;
2929        if changed != 1 {
2930            let state: Option<String> = transaction
2931                .query_row(
2932                    "SELECT CASE
2933                              WHEN missing_at_ms IS NOT NULL THEN 'missing'
2934                              WHEN state = 'ready' AND EXISTS (
2935                                  SELECT 1 FROM ticket_blockers b
2936                                  JOIN tickets bt ON bt.id = b.blocker_id
2937                                  WHERE b.ticket_id = tickets.id
2938                                    AND bt.state != 'merged'
2939                              ) THEN 'blocked'
2940                              ELSE state
2941                            END
2942                     FROM tickets WHERE id = ?1",
2943                    params![claim.ticket_id],
2944                    |row| row.get(0),
2945                )
2946                .optional()?;
2947            return Err(StoreError::TicketNotReady {
2948                ticket_id: claim.ticket_id.into(),
2949                state,
2950            });
2951        }
2952
2953        let activation_changed = match claim.next_activation_eligible_at_ms {
2954            Some(eligible_at_ms) => transaction.execute(
2955                "UPDATE activations
2956                 SET eligible_at_ms = ?2, updated_at_ms = ?3
2957                 WHERE id = ?1 AND state = 'queued' AND kind = 'every'",
2958                params![claim.activation_id, eligible_at_ms, now_ms],
2959            )?,
2960            None => transaction.execute(
2961                "UPDATE activations SET state = 'completed', updated_at_ms = ?2
2962                 WHERE id = ?1 AND state = 'queued' AND kind != 'every'",
2963                params![claim.activation_id, now_ms],
2964            )?,
2965        };
2966        if activation_changed != 1 {
2967            return Err(StoreError::ActivationNotQueued {
2968                activation_id: claim.activation_id.into(),
2969            });
2970        }
2971
2972        // The run's attempt counts runs, not the ticket's retry budget:
2973        // `retry` resets `tickets.attempts`, and a reused number would make two
2974        // runs answer to the same alias. Allocating inside the claim
2975        // transaction keeps the sequence gap-free under concurrent claims.
2976        let attempt: i64 = transaction.query_row(
2977            "SELECT COALESCE(MAX(attempt), 0) + 1 FROM runs WHERE ticket_id = ?1",
2978            params![claim.ticket_id],
2979            |row| row.get(0),
2980        )?;
2981
2982        transaction.execute(
2983            "INSERT INTO runs
2984                 (id, activation_id, ticket_id, state, attempt, flow_json, ticket_json,
2985                  created_at_ms, updated_at_ms)
2986             VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?8)",
2987            params![
2988                claim.run_id,
2989                claim.activation_id,
2990                claim.ticket_id,
2991                RunState::Claimed.as_str(),
2992                attempt,
2993                claim.flow_json,
2994                claim.ticket_json,
2995                now_ms,
2996            ],
2997        )?;
2998
2999        let expires_at_ms = now_ms + claim.lease_ms;
3000        transaction.execute(
3001            "INSERT INTO leases
3002                 (ticket_id, run_id, owner_id, acquired_at_ms, renewed_at_ms, expires_at_ms)
3003             VALUES (?1, ?2, ?3, ?4, ?4, ?5)",
3004            params![
3005                claim.ticket_id,
3006                claim.run_id,
3007                claim.owner_id,
3008                now_ms,
3009                expires_at_ms,
3010            ],
3011        )?;
3012        record_event(
3013            &transaction,
3014            now_ms,
3015            "run_claimed",
3016            Some(claim.run_id),
3017            Some(claim.ticket_id),
3018            &serde_json::json!({"attempt": attempt}).to_string(),
3019        )?;
3020
3021        transaction.commit()?;
3022        Ok(ClaimedRun {
3023            run_id: claim.run_id.into(),
3024            attempt,
3025            lease_expires_at_ms: expires_at_ms,
3026        })
3027    }
3028
3029    /// Re-arms the lease of a run this daemon has just adopted, returning the
3030    /// new expiry. Unlike [`Store::renew_lease`] this accepts an already
3031    /// expired lease: a daemon down longer than the TTL comes back to leases
3032    /// that lapsed while nobody was renewing them, and ordinary renewal could
3033    /// never lift them again. It is not a weaker renewal — the guard moves
3034    /// from the clock to the run itself, so only a run that has not settled
3035    /// can be re-armed, and a dead run's lease stays expired because recovery
3036    /// settles it instead of adopting it.
3037    pub(crate) fn readopt_lease(
3038        &mut self,
3039        ticket_id: &str,
3040        run_id: &str,
3041        lease_ms: i64,
3042        now_ms: i64,
3043    ) -> Result<i64, StoreError> {
3044        let expires_at_ms = now_ms + lease_ms;
3045        let changed = self.connection.execute(
3046            "UPDATE leases
3047             SET renewed_at_ms = ?3, expires_at_ms = ?4
3048             WHERE ticket_id = ?1 AND run_id = ?2
3049               AND EXISTS (SELECT 1 FROM runs
3050                           WHERE id = ?2 AND exited_at_ms IS NULL)",
3051            params![ticket_id, run_id, now_ms, expires_at_ms],
3052        )?;
3053        if changed != 1 {
3054            return Err(StoreError::LeaseNotHeld {
3055                ticket_id: ticket_id.into(),
3056                run_id: run_id.into(),
3057            });
3058        }
3059        Ok(expires_at_ms)
3060    }
3061
3062    /// Renews the lease that `run_id` holds on `ticket_id`, returning the new
3063    /// expiry. Renewal is strict: an expired lease cannot be renewed, so once
3064    /// recovery treats expiry as "run is lost" a revived run can never
3065    /// resurrect a lease that recovery may be reclaiming. Re-arming an expired
3066    /// lease is a separate, adoption-only verb ([`Store::readopt_lease`]).
3067    pub(crate) fn renew_lease(
3068        &mut self,
3069        ticket_id: &str,
3070        run_id: &str,
3071        lease_ms: i64,
3072        now_ms: i64,
3073    ) -> Result<i64, StoreError> {
3074        let expires_at_ms = now_ms + lease_ms;
3075        let changed = self.connection.execute(
3076            "UPDATE leases
3077             SET renewed_at_ms = ?3, expires_at_ms = ?4
3078             WHERE ticket_id = ?1 AND run_id = ?2 AND expires_at_ms > ?3",
3079            params![ticket_id, run_id, now_ms, expires_at_ms],
3080        )?;
3081        if changed != 1 {
3082            return Err(StoreError::LeaseNotHeld {
3083                ticket_id: ticket_id.into(),
3084                run_id: run_id.into(),
3085            });
3086        }
3087        Ok(expires_at_ms)
3088    }
3089
3090    pub fn paused(&self) -> Result<bool, StoreError> {
3091        let paused: i64 = self.connection.query_row(
3092            "SELECT paused FROM scheduler_state WHERE singleton = 1",
3093            [],
3094            |row| row.get(0),
3095        )?;
3096        Ok(paused != 0)
3097    }
3098
3099    pub fn clear_restart_draining(&self, now_ms: i64) -> Result<(), StoreError> {
3100        self.connection.execute(
3101            "UPDATE scheduler_state SET draining = 0, updated_at_ms = ?1 WHERE singleton = 1",
3102            params![now_ms],
3103        )?;
3104        Ok(())
3105    }
3106
3107    pub fn restart_draining(&self) -> Result<bool, StoreError> {
3108        let draining: i64 = self.connection.query_row(
3109            "SELECT draining FROM scheduler_state WHERE singleton = 1",
3110            [],
3111            |row| row.get(0),
3112        )?;
3113        Ok(draining != 0)
3114    }
3115
3116    pub fn active_cooldown_for_target(
3117        &self,
3118        target: &str,
3119        now_ms: i64,
3120    ) -> Result<Option<CooldownRecord>, StoreError> {
3121        self.connection
3122            .query_row(
3123                "SELECT ?1, until_ms, reason FROM cooldowns
3124                 WHERE key = 'agent_target:' || ?1 AND until_ms > ?2",
3125                params![target, now_ms],
3126                |row| {
3127                    Ok(CooldownRecord {
3128                        target: row.get(0)?,
3129                        until_ms: row.get(1)?,
3130                        reason: row.get(2)?,
3131                    })
3132                },
3133            )
3134            .optional()
3135            .map_err(StoreError::from)
3136    }
3137
3138    /// Number of runs currently holding a durable lease. Used as the capacity
3139    /// gate for repair spawns, which run inside an already-leased run.
3140    pub(crate) fn active_lease_count(&self) -> Result<usize, StoreError> {
3141        let count: i64 = self
3142            .connection
3143            .query_row("SELECT COUNT(*) FROM leases", [], |row| row.get(0))?;
3144        Ok(count as usize)
3145    }
3146
3147    /// Persists one repair attempt for a stage. The dedupe key is per
3148    /// (run, stage, attempt), so recovery counts consumed attempts without
3149    /// repeating or losing one, and the retry verdict can be filled in later
3150    /// by upserting the same key.
3151    pub(crate) fn record_repair_attempt(
3152        &self,
3153        run_id: &str,
3154        stage: &str,
3155        attempt: u32,
3156        data_json: &str,
3157        now_ms: i64,
3158    ) -> Result<(), StoreError> {
3159        self.connection.execute(
3160            "INSERT INTO run_evidence
3161                 (run_id, kind, observed_at_ms, dedupe_key, data_json)
3162             VALUES (?1, 'repair_attempt', ?2,
3163                     'repair:' || ?1 || ':' || ?3 || ':' || ?4, ?5)
3164             ON CONFLICT(dedupe_key) DO UPDATE SET
3165                 observed_at_ms = excluded.observed_at_ms,
3166                 data_json = excluded.data_json",
3167            params![run_id, now_ms, stage, attempt as i64, data_json],
3168        )?;
3169        Ok(())
3170    }
3171
3172    pub fn active_cooldowns(&self, now_ms: i64) -> Result<Vec<CooldownRecord>, StoreError> {
3173        let mut statement = self.connection.prepare(
3174            "SELECT SUBSTR(key, 14), until_ms, reason FROM cooldowns
3175             WHERE key LIKE 'agent_target:%' AND until_ms > ?1
3176             ORDER BY key",
3177        )?;
3178        statement
3179            .query_map(params![now_ms], |row| {
3180                Ok(CooldownRecord {
3181                    target: row.get(0)?,
3182                    until_ms: row.get(1)?,
3183                    reason: row.get(2)?,
3184                })
3185            })?
3186            .collect::<Result<Vec<_>, _>>()
3187            .map_err(StoreError::from)
3188    }
3189
3190    pub fn next_active_cooldown(&self, now_ms: i64) -> Result<Option<i64>, StoreError> {
3191        self.connection
3192            .query_row(
3193                "SELECT MIN(until_ms) FROM cooldowns WHERE until_ms > ?1",
3194                params![now_ms],
3195                |row| row.get(0),
3196            )
3197            .map_err(StoreError::from)
3198    }
3199
3200    pub fn set_paused(&self, paused: bool, now_ms: i64) -> Result<(), StoreError> {
3201        self.connection.execute(
3202            "UPDATE scheduler_state SET paused = ?1, updated_at_ms = ?2 WHERE singleton = 1",
3203            params![i64::from(paused), now_ms],
3204        )?;
3205        Ok(())
3206    }
3207
3208    pub fn begin_restart_draining(
3209        &mut self,
3210        active_runs: usize,
3211        now_ms: i64,
3212    ) -> Result<bool, StoreError> {
3213        let transaction = self
3214            .connection
3215            .transaction_with_behavior(TransactionBehavior::Immediate)?;
3216        let changed = transaction.execute(
3217            "UPDATE scheduler_state SET draining = 1, updated_at_ms = ?1
3218             WHERE singleton = 1 AND draining = 0",
3219            params![now_ms],
3220        )? != 0;
3221        if changed {
3222            record_event(
3223                &transaction,
3224                now_ms,
3225                "daemon_restart_requested",
3226                None,
3227                None,
3228                &serde_json::json!({"active_runs": active_runs}).to_string(),
3229            )?;
3230        }
3231        transaction.commit()?;
3232        Ok(changed)
3233    }
3234
3235    /// Resuming cancels both scheduler holds in one durable transition.
3236    pub fn resume_scheduler(&mut self, now_ms: i64) -> Result<bool, StoreError> {
3237        let transaction = self
3238            .connection
3239            .transaction_with_behavior(TransactionBehavior::Immediate)?;
3240        let was_draining: bool = transaction.query_row(
3241            "SELECT draining FROM scheduler_state WHERE singleton = 1",
3242            [],
3243            |row| row.get::<_, i64>(0).map(|value| value != 0),
3244        )?;
3245        transaction.execute(
3246            "UPDATE scheduler_state
3247             SET paused = 0, draining = 0, updated_at_ms = ?1
3248             WHERE singleton = 1",
3249            params![now_ms],
3250        )?;
3251        transaction.commit()?;
3252        Ok(was_draining)
3253    }
3254
3255    /// Performs a small committed write used to detect when SQLite can make
3256    /// progress again after returning `SQLITE_FULL`.
3257    pub(crate) fn probe_writable(&self, now_ms: i64) -> Result<(), StoreError> {
3258        self.connection.execute(
3259            "UPDATE scheduler_state SET updated_at_ms = ?1 WHERE singleton = 1",
3260            params![now_ms],
3261        )?;
3262        Ok(())
3263    }
3264
3265    pub fn ticket_counts(&self) -> Result<TicketCounts, StoreError> {
3266        let mut statement = self.connection.prepare(
3267            "SELECT CASE
3268                      WHEN t.state = 'ready' AND EXISTS (
3269                          SELECT 1 FROM ticket_blockers b
3270                          JOIN tickets bt ON bt.id = b.blocker_id
3271                          WHERE b.ticket_id = t.id AND bt.state != 'merged'
3272                      ) THEN 'blocked'
3273                      ELSE t.state
3274                    END AS display_state,
3275                    COUNT(*)
3276             FROM tickets t
3277             GROUP BY display_state",
3278        )?;
3279        let mut rows = statement.query([])?;
3280        let mut counts = TicketCounts::default();
3281        while let Some(row) = rows.next()? {
3282            let state: String = row.get(0)?;
3283            let count = row.get::<_, i64>(1)?.max(0) as u64;
3284            match state.as_str() {
3285                "ready" => counts.ready = count,
3286                "held" => counts.held = count,
3287                "blocked" => counts.blocked = count,
3288                "claimed" => counts.claimed = count,
3289                "merged" => counts.merged = count,
3290                "failed" => counts.failed = count,
3291                "needs_review" => counts.needs_review = count,
3292                _ => {}
3293            }
3294        }
3295        Ok(counts)
3296    }
3297}
3298
3299/// Whether the caller won the `running` → `aftercare` transition and with it
3300/// ownership of exit processing and aftercare for the run.
3301#[derive(Debug, PartialEq, Eq)]
3302pub(crate) enum ExitClaim {
3303    Claimed,
3304    AlreadyClaimed { state: String },
3305}
3306
3307fn upsert_cooldown(
3308    transaction: &rusqlite::Transaction<'_>,
3309    run_id: &str,
3310    cooldown: &CooldownUpdate<'_>,
3311    now_ms: i64,
3312) -> Result<(), rusqlite::Error> {
3313    transaction.execute(
3314        "INSERT INTO cooldowns (key, until_ms, reason, source_run_id, updated_at_ms)
3315         VALUES ('agent_target:' || ?1, ?2, ?3, ?4, ?5)
3316         ON CONFLICT(key) DO UPDATE SET
3317             until_ms = MAX(cooldowns.until_ms, excluded.until_ms),
3318             reason = CASE WHEN excluded.until_ms >= cooldowns.until_ms
3319                           THEN excluded.reason ELSE cooldowns.reason END,
3320             source_run_id = CASE WHEN excluded.until_ms >= cooldowns.until_ms
3321                                  THEN excluded.source_run_id ELSE cooldowns.source_run_id END,
3322             updated_at_ms = excluded.updated_at_ms",
3323        params![
3324            cooldown.target,
3325            cooldown.until_ms,
3326            cooldown.reason,
3327            run_id,
3328            now_ms
3329        ],
3330    )?;
3331    Ok(())
3332}
3333
3334#[derive(Debug)]
3335pub enum StoreError {
3336    Open {
3337        path: PathBuf,
3338        source: rusqlite::Error,
3339    },
3340    Sqlite(rusqlite::Error),
3341    UnsupportedSchemaVersion(u32),
3342    TicketNotReady {
3343        ticket_id: String,
3344        state: Option<String>,
3345    },
3346    TicketNotFound {
3347        ticket_id: String,
3348    },
3349    TicketStateConflict {
3350        ticket_id: String,
3351        state: String,
3352        requested: String,
3353    },
3354    ActivationNotQueued {
3355        activation_id: String,
3356    },
3357    LeaseNotHeld {
3358        ticket_id: String,
3359        run_id: String,
3360    },
3361    RunNotFound {
3362        run_id: String,
3363    },
3364    RunStateConflict {
3365        run_id: String,
3366        state: Option<String>,
3367        requested: String,
3368    },
3369    UnknownRunState {
3370        state: String,
3371    },
3372}
3373
3374impl StoreError {
3375    pub(crate) fn is_disk_full(&self) -> bool {
3376        let source = match self {
3377            Self::Open { source, .. } | Self::Sqlite(source) => source,
3378            _ => return false,
3379        };
3380        matches!(
3381            source,
3382            rusqlite::Error::SqliteFailure(error, _)
3383                if error.code == rusqlite::ffi::ErrorCode::DiskFull
3384        )
3385    }
3386}
3387
3388impl From<rusqlite::Error> for StoreError {
3389    fn from(source: rusqlite::Error) -> Self {
3390        Self::Sqlite(source)
3391    }
3392}
3393
3394impl fmt::Display for StoreError {
3395    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
3396        match self {
3397            Self::Open { path, source } => {
3398                write!(formatter, "cannot open {}: {source}", path.display())
3399            }
3400            Self::Sqlite(source) => write!(formatter, "database error: {source}"),
3401            Self::UnsupportedSchemaVersion(version) => {
3402                write!(formatter, "unsupported database schema version {version}")
3403            }
3404            Self::TicketNotReady { ticket_id, state } => match state {
3405                Some(state) => write!(formatter, "ticket `{ticket_id}` is `{state}`, not `ready`"),
3406                None => write!(formatter, "ticket `{ticket_id}` does not exist"),
3407            },
3408            Self::TicketNotFound { ticket_id } => {
3409                write!(formatter, "ticket `{ticket_id}` does not exist")
3410            }
3411            Self::TicketStateConflict {
3412                ticket_id,
3413                state,
3414                requested,
3415            } => write!(
3416                formatter,
3417                "ticket `{ticket_id}` is `{state}` and cannot be changed to `{requested}`"
3418            ),
3419            Self::ActivationNotQueued { activation_id } => write!(
3420                formatter,
3421                "activation `{activation_id}` is not queued for dispatch"
3422            ),
3423            Self::LeaseNotHeld { ticket_id, run_id } => write!(
3424                formatter,
3425                "run `{run_id}` does not hold the lease on ticket `{ticket_id}`"
3426            ),
3427            Self::RunNotFound { run_id } => write!(formatter, "run `{run_id}` does not exist"),
3428            Self::RunStateConflict {
3429                run_id,
3430                state,
3431                requested,
3432            } => match state {
3433                Some(state) => write!(
3434                    formatter,
3435                    "run `{run_id}` is `{state}` and cannot be changed to `{requested}`"
3436                ),
3437                None => write!(formatter, "run `{run_id}` does not exist"),
3438            },
3439            Self::UnknownRunState { state } => {
3440                write!(formatter, "unrecognized run state `{state}`")
3441            }
3442        }
3443    }
3444}
3445
3446impl std::error::Error for StoreError {}
3447
3448#[cfg(test)]
3449mod tests {
3450    use rusqlite::Connection;
3451    use tempfile::tempdir;
3452
3453    use super::{
3454        ActivationKind, ClaimRequest, ExitClaim, NewActivation, ReindexTicket, RunState,
3455        SCHEMA_VERSION, Store, StoreError,
3456    };
3457    use crate::domain::ticket::{TicketSnapshot, TicketState};
3458    use crate::flow::{Flow, Stage, StageKind, VerdictPolicy};
3459    use crate::outcome::Outcome;
3460
3461    fn open_seeded(path: &std::path::Path) -> Store {
3462        let store = Store::open(path, 1_000).unwrap();
3463        store
3464            .insert_local_project(
3465                "default",
3466                ".agents/sloop/projects/default.md",
3467                "Default",
3468                1_000,
3469            )
3470            .unwrap();
3471        store
3472            .insert_local_ticket(
3473                "T1",
3474                "default",
3475                ".agents/sloop/tickets/t1.md",
3476                "Ticket one",
3477                &[],
3478                "sloop/T1",
3479                Some("claude"),
3480                Some("sonnet"),
3481                Some("medium"),
3482                "default",
3483                TicketState::Ready,
3484                1_000,
3485            )
3486            .unwrap();
3487        store
3488            .insert_activation(
3489                &NewActivation {
3490                    id: "A1",
3491                    kind: ActivationKind::Immediate,
3492                    ticket_id: Some("T1"),
3493                    project_id: None,
3494                    eligible_at_ms: None,
3495                    interval_ms: None,
3496                },
3497                1_000,
3498            )
3499            .unwrap();
3500        store
3501    }
3502
3503    #[test]
3504    fn sqlite_full_errors_are_classified_for_backpressure() {
3505        let sqlite = rusqlite::Error::SqliteFailure(
3506            rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_FULL),
3507            None,
3508        );
3509        assert!(StoreError::from(sqlite).is_disk_full());
3510        assert!(
3511            !StoreError::TicketNotFound {
3512                ticket_id: "T1".into()
3513            }
3514            .is_disk_full()
3515        );
3516    }
3517
3518    #[test]
3519    fn writable_probe_commits_without_changing_pause_state() {
3520        let directory = tempdir().unwrap();
3521        let store = open_seeded(&directory.path().join("sloop.db"));
3522
3523        store.probe_writable(2_000).unwrap();
3524
3525        assert!(!store.paused().unwrap());
3526        let updated_at_ms: i64 = store
3527            .connection
3528            .query_row(
3529                "SELECT updated_at_ms FROM scheduler_state WHERE singleton = 1",
3530                [],
3531                |row| row.get(0),
3532            )
3533            .unwrap();
3534        assert_eq!(updated_at_ms, 2_000);
3535    }
3536
3537    fn claim_t1<'a>(run_id: &'a str) -> ClaimRequest<'a> {
3538        ClaimRequest {
3539            ticket_id: "T1",
3540            run_id,
3541            activation_id: "A1",
3542            owner_id: "daemon-1",
3543            lease_ms: 60_000,
3544            next_activation_eligible_at_ms: None,
3545            flow_json: "{}",
3546            ticket_json: "{}",
3547        }
3548    }
3549
3550    #[test]
3551    fn claims_persist_flow_and_ticket_snapshots() {
3552        let directory = tempdir().unwrap();
3553        let mut store = open_seeded(&directory.path().join("sloop.db"));
3554        let flow = Flow {
3555            name: "default".into(),
3556            stages: vec![
3557                Stage {
3558                    name: "build".into(),
3559                    kind: StageKind::Agent,
3560                    verdict: VerdictPolicy::Commits,
3561                    on_fail: None,
3562                },
3563                Stage {
3564                    name: "check".into(),
3565                    kind: StageKind::Exec {
3566                        cmd: vec!["cargo".into(), "test".into()],
3567                    },
3568                    verdict: VerdictPolicy::Exit,
3569                    on_fail: None,
3570                },
3571            ],
3572        };
3573        let ticket = TicketSnapshot {
3574            id: "T1".into(),
3575            name: "Ticket one".into(),
3576            blocked_by: vec![],
3577            worktree: Some("sloop/T1".into()),
3578            target: Some("claude".into()),
3579            model: Some("sonnet".into()),
3580            effort: Some("medium".into()),
3581            body: "# Original body\n".into(),
3582        };
3583        let flow_json = serde_json::to_string(&flow).unwrap();
3584        let ticket_json = serde_json::to_string(&ticket).unwrap();
3585
3586        store
3587            .claim_ticket(
3588                &ClaimRequest {
3589                    flow_json: &flow_json,
3590                    ticket_json: &ticket_json,
3591                    ..claim_t1("R1")
3592                },
3593                2_000,
3594            )
3595            .unwrap();
3596
3597        let run = store.run("R1").unwrap().unwrap();
3598        assert_eq!(
3599            serde_json::from_str::<Flow>(run.flow_json.as_deref().unwrap()).unwrap(),
3600            flow
3601        );
3602        assert_eq!(
3603            serde_json::from_str::<TicketSnapshot>(run.ticket_json.as_deref().unwrap()).unwrap(),
3604            ticket
3605        );
3606    }
3607
3608    #[test]
3609    fn missing_tickets_are_not_selected_and_cannot_be_claimed() {
3610        let directory = tempdir().unwrap();
3611        let mut store = open_seeded(&directory.path().join("sloop.db"));
3612        store.mark_ticket_missing("T1", 2_000).unwrap();
3613
3614        let activation = super::QueuedActivation {
3615            id: "A1".into(),
3616            kind: "immediate".into(),
3617            ticket_id: None,
3618            project_id: None,
3619            eligible_at_ms: None,
3620            interval_ms: None,
3621        };
3622        assert_eq!(store.select_ready_ticket(&activation, 2_000).unwrap(), None);
3623        match store.claim_ticket(&claim_t1("R1"), 2_000).unwrap_err() {
3624            StoreError::TicketNotReady { state, .. } => {
3625                assert_eq!(state.as_deref(), Some("missing"));
3626            }
3627            other => panic!("unexpected error: {other:?}"),
3628        }
3629
3630        // A second stamp must not restart the deletion clock.
3631        store.mark_ticket_missing("T1", 5_000).unwrap();
3632        assert_eq!(
3633            store.local_ticket_files().unwrap()[0].missing_at_ms,
3634            Some(2_000)
3635        );
3636
3637        store.clear_ticket_missing("T1", 6_000).unwrap();
3638        assert_eq!(
3639            store
3640                .select_ready_ticket(&activation, 6_000)
3641                .unwrap()
3642                .as_deref(),
3643            Some("T1")
3644        );
3645    }
3646
3647    #[test]
3648    fn blockers_gate_selection_claims_and_derived_counts_until_merged() {
3649        let directory = tempdir().unwrap();
3650        let mut store = open_seeded(&directory.path().join("sloop.db"));
3651        store
3652            .insert_local_ticket(
3653                "T2",
3654                "default",
3655                ".agents/sloop/tickets/t2.md",
3656                "Ticket two",
3657                &["T1".into()],
3658                "sloop/T2",
3659                Some("claude"),
3660                None,
3661                None,
3662                "default",
3663                TicketState::Ready,
3664                1_500,
3665            )
3666            .unwrap();
3667        let activation = super::QueuedActivation {
3668            id: "A1".into(),
3669            kind: "immediate".into(),
3670            ticket_id: None,
3671            project_id: None,
3672            eligible_at_ms: None,
3673            interval_ms: None,
3674        };
3675
3676        assert_eq!(store.unmerged_blockers("T2").unwrap(), ["T1"]);
3677        assert_eq!(
3678            store
3679                .select_ready_ticket(&activation, 2_000)
3680                .unwrap()
3681                .as_deref(),
3682            Some("T1")
3683        );
3684        assert_eq!(store.ticket_counts().unwrap().blocked, 1);
3685
3686        store
3687            .connection
3688            .execute("UPDATE tickets SET state = 'failed' WHERE id = 'T1'", [])
3689            .unwrap();
3690        assert_eq!(store.select_ready_ticket(&activation, 2_000).unwrap(), None);
3691        match store
3692            .claim_ticket(
3693                &ClaimRequest {
3694                    ticket_id: "T2",
3695                    run_id: "R2",
3696                    activation_id: "A1",
3697                    owner_id: "daemon-1",
3698                    lease_ms: 60_000,
3699                    next_activation_eligible_at_ms: None,
3700                    flow_json: "{}",
3701                    ticket_json: "{}",
3702                },
3703                2_000,
3704            )
3705            .unwrap_err()
3706        {
3707            StoreError::TicketNotReady { state, .. } => {
3708                assert_eq!(state.as_deref(), Some("blocked"));
3709            }
3710            other => panic!("unexpected error: {other:?}"),
3711        }
3712        assert_eq!(store.ticket("T2").unwrap().unwrap().attempts, 0);
3713
3714        store
3715            .connection
3716            .execute("UPDATE tickets SET state = 'merged' WHERE id = 'T1'", [])
3717            .unwrap();
3718        assert!(store.unmerged_blockers("T2").unwrap().is_empty());
3719        assert_eq!(
3720            store
3721                .select_ready_ticket(&activation, 2_000)
3722                .unwrap()
3723                .as_deref(),
3724            Some("T2")
3725        );
3726        let counts = store.ticket_counts().unwrap();
3727        assert_eq!(counts.ready, 1);
3728        assert_eq!(counts.blocked, 0);
3729    }
3730
3731    #[test]
3732    fn state_survives_reopening_the_database() {
3733        let directory = tempdir().unwrap();
3734        let path = directory.path().join("sloop.db");
3735
3736        let mut store = open_seeded(&path);
3737        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
3738        drop(store);
3739
3740        let store = Store::open(&path, 3_000).unwrap();
3741        assert_eq!(store.ticket_state("T1").unwrap().unwrap(), "claimed");
3742        assert_eq!(store.ticket_counts().unwrap().claimed, 1);
3743        let ticket = store.ticket("T1").unwrap().unwrap();
3744        assert_eq!(ticket.target.as_deref(), Some("claude"));
3745        assert_eq!(ticket.model.as_deref(), Some("sonnet"));
3746        assert_eq!(ticket.effort.as_deref(), Some("medium"));
3747        assert_eq!(ticket.name, "Ticket one");
3748        assert!(ticket.blocked_by.is_empty());
3749        assert_eq!(ticket.worktree.as_deref(), Some("sloop/T1"));
3750    }
3751
3752    #[test]
3753    fn blocked_by_and_worktree_round_trip() {
3754        let directory = tempdir().unwrap();
3755        let path = directory.path().join("sloop.db");
3756        let store = open_seeded(&path);
3757        store
3758            .insert_local_ticket(
3759                "T2",
3760                "default",
3761                ".agents/sloop/tickets/t2.md",
3762                "Ticket two",
3763                &["T1".to_owned()],
3764                "feature/t2",
3765                None,
3766                None,
3767                None,
3768                "default",
3769                TicketState::Ready,
3770                2_000,
3771            )
3772            .unwrap();
3773        drop(store);
3774
3775        let store = Store::open(&path, 3_000).unwrap();
3776        let ticket = store.ticket("T2").unwrap().unwrap();
3777        assert_eq!(ticket.name, "Ticket two");
3778        assert_eq!(ticket.blocked_by, ["T1"]);
3779        assert_eq!(ticket.worktree.as_deref(), Some("feature/t2"));
3780    }
3781
3782    #[test]
3783    fn a_claimed_ticket_cannot_be_claimed_again() {
3784        let directory = tempdir().unwrap();
3785        let mut store = open_seeded(&directory.path().join("sloop.db"));
3786
3787        let claimed = store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
3788        assert_eq!(claimed.attempt, 1);
3789        assert_eq!(claimed.lease_expires_at_ms, 62_000);
3790
3791        let error = store.claim_ticket(&claim_t1("R2"), 2_100).unwrap_err();
3792        assert!(matches!(
3793            error,
3794            StoreError::TicketNotReady { state: Some(ref state), .. } if state == "claimed"
3795        ));
3796    }
3797
3798    #[test]
3799    fn tickets_are_ordered_newest_first_and_include_attempts() {
3800        let directory = tempdir().unwrap();
3801        let mut store = open_seeded(&directory.path().join("sloop.db"));
3802        store
3803            .insert_local_project("alpha", ".agents/sloop/projects/alpha.md", "Alpha", 1_000)
3804            .unwrap();
3805        store
3806            .insert_local_ticket(
3807                "T0",
3808                "alpha",
3809                ".agents/sloop/tickets/t0.md",
3810                "Ticket zero",
3811                &[],
3812                "sloop/T0",
3813                None,
3814                None,
3815                None,
3816                "default",
3817                TicketState::Held,
3818                3_000,
3819            )
3820            .unwrap();
3821        store
3822            .insert_local_ticket(
3823                "T2",
3824                "default",
3825                ".agents/sloop/tickets/t2.md",
3826                "Ticket two",
3827                &[],
3828                "sloop/T2",
3829                None,
3830                None,
3831                None,
3832                "default",
3833                TicketState::Ready,
3834                1_000,
3835            )
3836            .unwrap();
3837        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
3838
3839        let tickets = store.tickets().unwrap();
3840        // T0 registered last, so it leads despite the lowest ordinal and a
3841        // different project; T2 and T1 tie on time and fall back to ordinal.
3842        assert_eq!(
3843            tickets
3844                .iter()
3845                .map(|ticket| ticket.id.as_str())
3846                .collect::<Vec<_>>(),
3847            ["T0", "T2", "T1"]
3848        );
3849        assert_eq!(tickets[0].attempts, 0);
3850        assert_eq!(tickets[1].attempts, 0);
3851        assert_eq!(tickets[2].attempts, 1);
3852    }
3853
3854    #[test]
3855    fn active_run_for_ticket_tracks_claimed_and_running_runs_only() {
3856        use crate::outcome::Outcome;
3857
3858        let directory = tempdir().unwrap();
3859        let mut store = open_seeded(&directory.path().join("sloop.db"));
3860        assert_eq!(store.active_run_for_ticket("T1").unwrap(), None);
3861
3862        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
3863        assert_eq!(
3864            store.active_run_for_ticket("T1").unwrap(),
3865            Some(("R1".into(), 1))
3866        );
3867        store
3868            .mark_run_running(
3869                "R1",
3870                "branch",
3871                "/tmp/worktree",
3872                1,
3873                Some(1),
3874                1,
3875                "token",
3876                "/runtime/R1.sock",
3877                2_100,
3878            )
3879            .unwrap();
3880        assert_eq!(
3881            store.active_run_for_ticket("T1").unwrap(),
3882            Some(("R1".into(), 1))
3883        );
3884
3885        store
3886            .finish_run("R1", "T1", Some(1), Outcome::Failed, &[], None, 2_200)
3887            .unwrap();
3888        assert_eq!(store.active_run_for_ticket("T1").unwrap(), None);
3889    }
3890
3891    #[test]
3892    fn aborted_claims_are_closed_and_no_longer_active() {
3893        let directory = tempdir().unwrap();
3894        let mut store = open_seeded(&directory.path().join("sloop.db"));
3895        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
3896
3897        store.abort_claim("R1", "T1", 2_100).unwrap();
3898
3899        assert_eq!(store.run("R1").unwrap().unwrap().state, "aborted");
3900        assert_eq!(store.active_run_for_ticket("T1").unwrap(), None);
3901        assert_eq!(store.ticket_state("T1").unwrap().as_deref(), Some("ready"));
3902    }
3903
3904    #[test]
3905    fn recoverable_runs_round_trip_process_identity_and_lease() {
3906        let directory = tempdir().unwrap();
3907        let path = directory.path().join("sloop.db");
3908        let mut store = open_seeded(&path);
3909        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
3910        store
3911            .mark_run_running(
3912                "R1",
3913                "sloop/T1-a1-R1",
3914                "/worktrees/R1",
3915                123,
3916                Some(456),
3917                123,
3918                "worker-token",
3919                "/runtime/R1.sock",
3920                2_100,
3921            )
3922            .unwrap();
3923        drop(store);
3924
3925        let store = Store::open(&path, 3_000).unwrap();
3926        let runs = store.recoverable_runs().unwrap();
3927        assert_eq!(runs.len(), 1);
3928        assert_eq!(runs[0].id, "R1");
3929        assert_eq!(runs[0].ticket_id, "T1");
3930        assert_eq!(runs[0].pid, Some(123));
3931        assert_eq!(runs[0].pid_start_time, Some(456));
3932        assert_eq!(runs[0].process_group_id, Some(123));
3933        assert_eq!(runs[0].worker_token.as_deref(), Some("worker-token"));
3934        assert_eq!(
3935            runs[0].worker_socket_path.as_deref(),
3936            Some("/runtime/R1.sock")
3937        );
3938        assert_eq!(runs[0].exit_code, None);
3939        assert_eq!(runs[0].lease_expires_at_ms, 62_000);
3940    }
3941
3942    #[test]
3943    fn agent_exit_and_aftercare_results_are_checkpointed_idempotently() {
3944        let directory = tempdir().unwrap();
3945        let mut store = open_seeded(&directory.path().join("sloop.db"));
3946        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
3947        store
3948            .mark_run_running(
3949                "R1",
3950                "branch",
3951                "/worktree",
3952                123,
3953                Some(456),
3954                123,
3955                "token",
3956                "/runtime/R1.sock",
3957                2_100,
3958            )
3959            .unwrap();
3960
3961        store
3962            .record_agent_exit(
3963                "R1",
3964                Some(0),
3965                true,
3966                r#"{"count":1,"oids":["abc"]}"#,
3967                None,
3968                None,
3969                2_200,
3970            )
3971            .unwrap();
3972        store
3973            .record_aftercare_evidence(
3974                "R1",
3975                "test_result",
3976                r#"{"passed":true,"exit_code":0}"#,
3977                2_300,
3978            )
3979            .unwrap();
3980        store
3981            .record_aftercare_evidence(
3982                "R1",
3983                "test_result",
3984                r#"{"passed":true,"exit_code":0}"#,
3985                2_400,
3986            )
3987            .unwrap();
3988
3989        let run = store.run("R1").unwrap().unwrap();
3990        assert_eq!(run.state, "aftercare");
3991        assert_eq!(run.exit_code, Some(0));
3992        assert_eq!(
3993            store.recoverable_runs().unwrap()[0].state,
3994            RunState::Aftercare
3995        );
3996        let evidence = store.run_evidence("R1").unwrap();
3997        assert_eq!(
3998            evidence
3999                .iter()
4000                .filter(|(kind, _)| kind == "test_result")
4001                .count(),
4002            1
4003        );
4004    }
4005
4006    fn running_r1(store: &mut Store) {
4007        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4008        store
4009            .mark_run_running(
4010                "R1",
4011                "branch",
4012                "/worktree",
4013                123,
4014                Some(456),
4015                123,
4016                "token",
4017                "/runtime/R1.sock",
4018                2_100,
4019            )
4020            .unwrap();
4021    }
4022
4023    #[test]
4024    fn agent_exit_checkpoint_is_an_exclusive_ownership_handoff() {
4025        let directory = tempdir().unwrap();
4026        let mut store = open_seeded(&directory.path().join("sloop.db"));
4027        running_r1(&mut store);
4028
4029        let first = store
4030            .record_agent_exit(
4031                "R1",
4032                Some(0),
4033                true,
4034                r#"{"oids":["abc"]}"#,
4035                None,
4036                None,
4037                2_200,
4038            )
4039            .unwrap();
4040        assert_eq!(first, ExitClaim::Claimed);
4041        assert_eq!(store.run("R1").unwrap().unwrap().state, "aftercare");
4042        let evidence = store.run_evidence("R1").unwrap();
4043        assert!(evidence.iter().any(|(kind, _)| kind == "exit_classified"));
4044        assert!(evidence.iter().any(|(kind, _)| kind == "commits_observed"));
4045
4046        let second = store
4047            .record_agent_exit("R1", Some(1), false, r#"{"oids":[]}"#, None, None, 2_300)
4048            .unwrap();
4049        assert_eq!(
4050            second,
4051            ExitClaim::AlreadyClaimed {
4052                state: "aftercare".into()
4053            }
4054        );
4055        let run = store.run("R1").unwrap().unwrap();
4056        assert_eq!(run.state, "aftercare");
4057        assert_eq!(run.exit_code, Some(0));
4058        assert_eq!(store.run_evidence("R1").unwrap(), evidence);
4059    }
4060
4061    #[test]
4062    fn agent_exit_checkpoint_reports_terminal_and_missing_runs() {
4063        let directory = tempdir().unwrap();
4064        let mut store = open_seeded(&directory.path().join("sloop.db"));
4065        running_r1(&mut store);
4066        store
4067            .finish_run("R1", "T1", Some(0), Outcome::Merged, &[], None, 2_200)
4068            .unwrap();
4069
4070        let claim = store
4071            .record_agent_exit(
4072                "R1",
4073                Some(0),
4074                true,
4075                r#"{"count":0,"oids":[]}"#,
4076                None,
4077                None,
4078                2_300,
4079            )
4080            .unwrap();
4081        assert_eq!(
4082            claim,
4083            ExitClaim::AlreadyClaimed {
4084                state: "merged".into()
4085            }
4086        );
4087        assert_eq!(store.run("R1").unwrap().unwrap().state, "merged");
4088
4089        let missing = store.record_agent_exit(
4090            "R9",
4091            Some(0),
4092            true,
4093            r#"{"count":0,"oids":[]}"#,
4094            None,
4095            None,
4096            2_300,
4097        );
4098        assert!(matches!(missing, Err(StoreError::RunNotFound { .. })));
4099    }
4100
4101    #[test]
4102    fn finish_run_settles_exactly_once() {
4103        let directory = tempdir().unwrap();
4104        let mut store = open_seeded(&directory.path().join("sloop.db"));
4105        running_r1(&mut store);
4106        store
4107            .record_agent_exit(
4108                "R1",
4109                Some(0),
4110                true,
4111                r#"{"count":1,"oids":["abc"]}"#,
4112                None,
4113                None,
4114                2_200,
4115            )
4116            .unwrap();
4117
4118        store
4119            .finish_run("R1", "T1", Some(0), Outcome::Merged, &[], None, 2_300)
4120            .unwrap();
4121        assert_eq!(store.ticket_state("T1").unwrap().as_deref(), Some("merged"));
4122        assert_eq!(store.active_run_for_ticket("T1").unwrap(), None);
4123        let evidence = store.run_evidence("R1").unwrap();
4124
4125        store
4126            .finish_run("R1", "T1", Some(1), Outcome::Failed, &[], None, 2_400)
4127            .unwrap();
4128        let run = store.run("R1").unwrap().unwrap();
4129        assert_eq!(run.state, "merged");
4130        assert_eq!(run.exit_code, Some(0));
4131        assert_eq!(store.ticket_state("T1").unwrap().as_deref(), Some("merged"));
4132        assert_eq!(store.run_evidence("R1").unwrap(), evidence);
4133    }
4134
4135    #[test]
4136    fn lifecycle_transitions_append_ordered_events() {
4137        let directory = tempdir().unwrap();
4138        let mut store = open_seeded(&directory.path().join("sloop.db"));
4139        running_r1(&mut store);
4140        store
4141            .finish_run("R1", "T1", Some(0), Outcome::Merged, &[], None, 2_300)
4142            .unwrap();
4143
4144        let events = store.events_after(0, 10).unwrap();
4145        let kinds: Vec<&str> = events.iter().map(|event| event.kind.as_str()).collect();
4146        assert_eq!(kinds, ["run_claimed", "run_started", "run_finished"]);
4147        assert!(events.iter().all(|event| {
4148            event.run_id.as_deref() == Some("R1") && event.ticket_id.as_deref() == Some("T1")
4149        }));
4150        let finished: serde_json::Value = serde_json::from_str(&events[2].data_json).unwrap();
4151        assert_eq!(finished["outcome"], "merged");
4152        assert_eq!(finished["ticket_state"], "merged");
4153
4154        // Settling twice is idempotent, so no duplicate event appears.
4155        store
4156            .finish_run("R1", "T1", Some(1), Outcome::Failed, &[], None, 2_400)
4157            .unwrap();
4158        assert_eq!(store.latest_event_sequence().unwrap(), events[2].sequence);
4159
4160        let rest = store.events_after(events[0].sequence, 10).unwrap();
4161        assert_eq!(rest.len(), 2);
4162        assert_eq!(rest[0].kind, "run_started");
4163
4164        store.trim_events(1).unwrap();
4165        let kept = store.events_after(0, 10).unwrap();
4166        assert_eq!(kept.len(), 1);
4167        assert_eq!(kept[0].sequence, events[2].sequence);
4168    }
4169
4170    #[test]
4171    fn abandoned_claims_append_an_abort_event() {
4172        let directory = tempdir().unwrap();
4173        let mut store = open_seeded(&directory.path().join("sloop.db"));
4174        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4175        store.abort_claim("R1", "T1", 2_100).unwrap();
4176
4177        let kinds: Vec<String> = store
4178            .events_after(0, 10)
4179            .unwrap()
4180            .into_iter()
4181            .map(|event| event.kind)
4182            .collect();
4183        assert_eq!(kinds, ["run_claimed", "run_aborted"]);
4184    }
4185
4186    #[test]
4187    fn operator_hold_transitions_are_narrow_and_idempotent() {
4188        let directory = tempdir().unwrap();
4189        let store = open_seeded(&directory.path().join("sloop.db"));
4190
4191        assert_eq!(
4192            store
4193                .set_ticket_hold("T1", TicketState::Held, 2_000)
4194                .unwrap(),
4195            "ready"
4196        );
4197        assert_eq!(store.ticket_counts().unwrap().held, 1);
4198        assert_eq!(
4199            store
4200                .set_ticket_hold("T1", TicketState::Held, 2_100)
4201                .unwrap(),
4202            "held"
4203        );
4204        assert_eq!(
4205            store
4206                .set_ticket_hold("T1", TicketState::Ready, 2_200)
4207                .unwrap(),
4208            "held"
4209        );
4210    }
4211
4212    #[test]
4213    fn validation_hold_reasons_set_and_clear_without_releasing_operator_holds() {
4214        let directory = tempdir().unwrap();
4215        let store = open_seeded(&directory.path().join("sloop.db"));
4216        let ticket = |held_reason: Option<&str>| ReindexTicket {
4217            id: "T1".into(),
4218            project_id: "default".into(),
4219            source: "markdown".into(),
4220            source_ref: ".agents/sloop/tickets/t1.md".into(),
4221            file_path: Some(".agents/sloop/tickets/t1.md".into()),
4222            name: "Ticket one".into(),
4223            blocked_by: Vec::new(),
4224            worktree: "sloop/T1".into(),
4225            target: Some("claude".into()),
4226            model: Some("sonnet".into()),
4227            effort: Some("medium".into()),
4228            flow: "default".into(),
4229            body: "work".into(),
4230            held_reason: held_reason.map(str::to_owned),
4231            derived_state: None,
4232        };
4233
4234        store
4235            .apply_reindex(
4236                &["default".into()],
4237                &[ticket(Some("flow `missing` is not defined"))],
4238                2_000,
4239            )
4240            .unwrap();
4241        let held = store.ticket("T1").unwrap().unwrap();
4242        assert_eq!(held.state, "held");
4243        assert_eq!(
4244            held.held_reason.as_deref(),
4245            Some("flow `missing` is not defined")
4246        );
4247
4248        store
4249            .apply_reindex(&["default".into()], &[ticket(None)], 2_100)
4250            .unwrap();
4251        let released = store.ticket("T1").unwrap().unwrap();
4252        assert_eq!(released.state, "ready");
4253        assert_eq!(released.held_reason, None);
4254
4255        store
4256            .set_ticket_hold("T1", TicketState::Held, 2_200)
4257            .unwrap();
4258        store
4259            .apply_reindex(&["default".into()], &[ticket(None)], 2_300)
4260            .unwrap();
4261        let operator_held = store.ticket("T1").unwrap().unwrap();
4262        assert_eq!(operator_held.state, "held");
4263        assert_eq!(operator_held.held_reason, None);
4264    }
4265
4266    #[test]
4267    fn operator_hold_cannot_steal_a_claim() {
4268        let directory = tempdir().unwrap();
4269        let mut store = open_seeded(&directory.path().join("sloop.db"));
4270        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4271
4272        assert!(matches!(
4273            store.set_ticket_hold("T1", TicketState::Held, 2_100),
4274            Err(StoreError::TicketStateConflict { state, .. }) if state == "claimed"
4275        ));
4276    }
4277
4278    #[test]
4279    fn retry_only_requeues_failed_tickets_and_resets_attempts() {
4280        use crate::outcome::Outcome;
4281
4282        let directory = tempdir().unwrap();
4283        let mut store = open_seeded(&directory.path().join("sloop.db"));
4284
4285        let first = store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4286        assert_eq!(first.attempt, 1);
4287        store
4288            .finish_run("R1", "T1", Some(0), Outcome::Failed, &[], None, 2_100)
4289            .unwrap();
4290
4291        assert_eq!(store.retry_ticket("T1", 2_200).unwrap(), "failed");
4292        assert_eq!(store.ticket_state("T1").unwrap().as_deref(), Some("ready"));
4293        store
4294            .insert_activation(
4295                &NewActivation {
4296                    id: "A2",
4297                    kind: ActivationKind::Immediate,
4298                    ticket_id: Some("T1"),
4299                    project_id: None,
4300                    eligible_at_ms: None,
4301                    interval_ms: None,
4302                },
4303                2_300,
4304            )
4305            .unwrap();
4306        let retried = store
4307            .claim_ticket(
4308                &ClaimRequest {
4309                    activation_id: "A2",
4310                    ..claim_t1("R2")
4311                },
4312                2_300,
4313            )
4314            .unwrap();
4315        // `retry` resets the ticket's attempt budget, but a run's attempt
4316        // counts runs of that ticket: it must keep climbing, or two runs would
4317        // answer to the same `T1-r1` alias.
4318        assert_eq!(retried.attempt, 2);
4319        assert_eq!(store.ticket("T1").unwrap().unwrap().attempts, 1);
4320
4321        assert!(matches!(
4322            store.retry_ticket("T1", 2_400),
4323            Err(StoreError::TicketStateConflict { state, .. }) if state == "claimed"
4324        ));
4325        assert!(matches!(
4326            store.retry_ticket("missing", 2_400),
4327            Err(StoreError::TicketNotFound { .. })
4328        ));
4329    }
4330
4331    #[test]
4332    fn claiming_an_unknown_ticket_reports_it_missing() {
4333        let directory = tempdir().unwrap();
4334        let mut store = open_seeded(&directory.path().join("sloop.db"));
4335
4336        let error = store
4337            .claim_ticket(
4338                &ClaimRequest {
4339                    ticket_id: "missing",
4340                    ..claim_t1("R1")
4341                },
4342                2_000,
4343            )
4344            .unwrap_err();
4345        assert!(matches!(
4346            error,
4347            StoreError::TicketNotReady { state: None, .. }
4348        ));
4349    }
4350
4351    #[test]
4352    fn concurrent_connections_cannot_both_claim_one_ticket() {
4353        let directory = tempdir().unwrap();
4354        let path = directory.path().join("sloop.db");
4355        open_seeded(&path);
4356
4357        let barrier = std::sync::Arc::new(std::sync::Barrier::new(2));
4358        let claims: Vec<_> = ["R1", "R2"]
4359            .into_iter()
4360            .map(|run_id| {
4361                let path = path.clone();
4362                let barrier = barrier.clone();
4363                std::thread::spawn(move || {
4364                    let mut store = Store::open(&path, 2_000).unwrap();
4365                    barrier.wait();
4366                    store.claim_ticket(&claim_t1(run_id), 2_000).is_ok()
4367                })
4368            })
4369            .collect();
4370
4371        let successes = claims
4372            .into_iter()
4373            .map(|handle| handle.join().unwrap())
4374            .filter(|claimed| *claimed)
4375            .count();
4376        assert_eq!(successes, 1);
4377    }
4378
4379    #[test]
4380    fn renewing_a_held_lease_extends_its_expiry() {
4381        let directory = tempdir().unwrap();
4382        let mut store = open_seeded(&directory.path().join("sloop.db"));
4383        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4384
4385        let expires = store.renew_lease("T1", "R1", 60_000, 10_000).unwrap();
4386        assert_eq!(expires, 70_000);
4387    }
4388
4389    #[test]
4390    fn a_run_cannot_renew_a_lease_it_does_not_hold() {
4391        let directory = tempdir().unwrap();
4392        let mut store = open_seeded(&directory.path().join("sloop.db"));
4393        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4394
4395        let error = store.renew_lease("T1", "R2", 60_000, 10_000).unwrap_err();
4396        assert!(matches!(error, StoreError::LeaseNotHeld { .. }));
4397    }
4398
4399    #[test]
4400    fn an_expired_lease_cannot_be_renewed() {
4401        let directory = tempdir().unwrap();
4402        let mut store = open_seeded(&directory.path().join("sloop.db"));
4403        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4404
4405        // The lease expires at 62_000; renewal at or after that must fail.
4406        let error = store.renew_lease("T1", "R1", 60_000, 62_000).unwrap_err();
4407        assert!(matches!(error, StoreError::LeaseNotHeld { .. }));
4408    }
4409
4410    #[test]
4411    fn every_run_state_round_trips_through_its_stored_string() {
4412        let states = [
4413            RunState::Claimed,
4414            RunState::Running,
4415            RunState::Aftercare,
4416            RunState::Aborted,
4417            RunState::Merged,
4418            RunState::Failed,
4419            RunState::NeedsReview,
4420            RunState::Cancelled,
4421            RunState::RateLimited,
4422            RunState::Orphaned,
4423        ];
4424        for state in states {
4425            assert_eq!(RunState::parse(state.as_str()).unwrap(), state);
4426        }
4427        // Every outcome `finish_run` can write is one of those variants.
4428        for outcome in [
4429            crate::outcome::Outcome::Merged,
4430            crate::outcome::Outcome::Failed,
4431            crate::outcome::Outcome::NeedsReview,
4432            crate::outcome::Outcome::Cancelled,
4433            crate::outcome::Outcome::RateLimited,
4434            crate::outcome::Outcome::Orphaned,
4435        ] {
4436            assert_eq!(RunState::from(outcome).as_str(), outcome.as_str());
4437            assert!(RunState::from(outcome).is_terminal());
4438        }
4439        for state in [RunState::Claimed, RunState::Running, RunState::Aftercare] {
4440            assert!(!state.is_terminal());
4441        }
4442        assert!(RunState::Aborted.is_terminal());
4443    }
4444
4445    #[test]
4446    fn an_unknown_stored_run_state_is_an_error_not_a_fallback() {
4447        let error = RunState::parse("half_running").unwrap_err();
4448        assert!(matches!(error, StoreError::UnknownRunState { state } if state == "half_running"));
4449    }
4450
4451    #[test]
4452    fn a_readopted_lease_is_re_armed_even_after_it_expired() {
4453        let directory = tempdir().unwrap();
4454        let mut store = open_seeded(&directory.path().join("sloop.db"));
4455        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4456
4457        // The lease expired at 62_000, so ordinary renewal is refused...
4458        assert!(store.renew_lease("T1", "R1", 60_000, 90_000).is_err());
4459        // ...while adoption re-arms it, and renewal works again afterwards.
4460        assert_eq!(
4461            store.readopt_lease("T1", "R1", 60_000, 90_000).unwrap(),
4462            150_000
4463        );
4464        assert_eq!(
4465            store.renew_lease("T1", "R1", 60_000, 100_000).unwrap(),
4466            160_000
4467        );
4468    }
4469
4470    #[test]
4471    fn a_settled_run_cannot_be_readopted() {
4472        let directory = tempdir().unwrap();
4473        let mut store = open_seeded(&directory.path().join("sloop.db"));
4474        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4475        store
4476            .finish_run(
4477                "R1",
4478                "T1",
4479                Some(0),
4480                crate::outcome::Outcome::Failed,
4481                &[],
4482                None,
4483                3_000,
4484            )
4485            .unwrap();
4486
4487        let error = store.readopt_lease("T1", "R1", 60_000, 4_000).unwrap_err();
4488        assert!(matches!(error, StoreError::LeaseNotHeld { .. }));
4489    }
4490
4491    #[test]
4492    fn ready_work_selection_is_deterministic_and_respects_filters() {
4493        let directory = tempdir().unwrap();
4494        let store = open_seeded(&directory.path().join("sloop.db"));
4495        store
4496            .insert_local_ticket(
4497                "T0",
4498                "default",
4499                ".agents/sloop/tickets/t0.md",
4500                "Ticket zero",
4501                &[],
4502                "sloop/T0",
4503                None,
4504                None,
4505                None,
4506                "default",
4507                TicketState::Ready,
4508                2_000,
4509            )
4510            .unwrap();
4511        store
4512            .insert_activation(
4513                &NewActivation {
4514                    id: "A2",
4515                    kind: ActivationKind::Immediate,
4516                    ticket_id: None,
4517                    project_id: None,
4518                    eligible_at_ms: None,
4519                    interval_ms: None,
4520                },
4521                2_000,
4522            )
4523            .unwrap();
4524        let activation = super::QueuedActivation {
4525            id: "A2".into(),
4526            kind: "immediate".into(),
4527            ticket_id: None,
4528            project_id: None,
4529            eligible_at_ms: None,
4530            interval_ms: None,
4531        };
4532
4533        // T1 was registered first, so it wins despite T0 sorting lower.
4534        assert_eq!(
4535            store
4536                .select_ready_ticket(&activation, 2_000)
4537                .unwrap()
4538                .as_deref(),
4539            Some("T1")
4540        );
4541
4542        store.insert_activation_filter("A2", "T0").unwrap();
4543        assert_eq!(
4544            store
4545                .select_ready_ticket(&activation, 2_000)
4546                .unwrap()
4547                .as_deref(),
4548            Some("T0")
4549        );
4550
4551        let scoped = super::QueuedActivation {
4552            project_id: Some("elsewhere".into()),
4553            ..activation
4554        };
4555        assert_eq!(store.select_ready_ticket(&scoped, 2_000).unwrap(), None);
4556    }
4557
4558    #[test]
4559    fn notes_round_trip_in_arrival_order() {
4560        let directory = tempdir().unwrap();
4561        let mut store = open_seeded(&directory.path().join("sloop.db"));
4562        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4563
4564        assert_eq!(store.next_note_ordinal().unwrap(), 1);
4565        store.insert_note("N1", "R1", "first", 3_000).unwrap();
4566        store.insert_note("N2", "R1", "second", 3_000).unwrap();
4567        assert_eq!(store.next_note_ordinal().unwrap(), 3);
4568
4569        assert_eq!(
4570            store.notes_for_run("R1").unwrap(),
4571            vec!["first".to_owned(), "second".to_owned()]
4572        );
4573        assert!(store.notes_for_run("R2").unwrap().is_empty());
4574    }
4575
4576    #[test]
4577    fn version_three_migrates_ticket_metadata_and_newer_schemas_are_rejected() {
4578        let directory = tempdir().unwrap();
4579        let path = directory.path().join("sloop.db");
4580        drop(Store::open(&path, 1_000).unwrap());
4581
4582        let connection = rusqlite::Connection::open(&path).unwrap();
4583        connection
4584            .execute_batch(
4585                "DROP TABLE ticket_blockers;
4586                 ALTER TABLE tickets DROP COLUMN name;
4587                 ALTER TABLE tickets DROP COLUMN worktree;
4588                 ALTER TABLE tickets DROP COLUMN flow;
4589                 ALTER TABLE tickets DROP COLUMN body;
4590                 ALTER TABLE tickets DROP COLUMN held_reason;
4591                 ALTER TABLE tickets DROP COLUMN missing_at_ms;
4592                 ALTER TABLE scheduler_state DROP COLUMN draining;
4593                 ALTER TABLE runs DROP COLUMN worker_socket_path;
4594                 ALTER TABLE runs DROP COLUMN flow_json;
4595                 ALTER TABLE runs DROP COLUMN ticket_json;
4596                 ALTER TABLE runs DROP COLUMN cleanup_eligible_at_ms;
4597                 ALTER TABLE runs DROP COLUMN cleaned_at_ms;",
4598            )
4599            .unwrap();
4600        connection.pragma_update(None, "user_version", 3).unwrap();
4601        drop(connection);
4602
4603        let store = Store::open(&path, 2_000).unwrap();
4604        assert!(!store.paused().unwrap());
4605        store
4606            .insert_local_project(
4607                "default",
4608                ".agents/sloop/projects/default.md",
4609                "Default",
4610                2_000,
4611            )
4612            .unwrap();
4613        store
4614            .insert_local_ticket(
4615                "T1",
4616                "default",
4617                ".agents/sloop/tickets/t1.md",
4618                "Ticket one",
4619                &[],
4620                "sloop/T1",
4621                Some("codex"),
4622                None,
4623                None,
4624                "default",
4625                TicketState::Ready,
4626                2_000,
4627            )
4628            .unwrap();
4629        assert_eq!(
4630            store.ticket("T1").unwrap().unwrap().target.as_deref(),
4631            Some("codex")
4632        );
4633        drop(store);
4634
4635        let connection = rusqlite::Connection::open(&path).unwrap();
4636        connection.pragma_update(None, "user_version", 99).unwrap();
4637        drop(connection);
4638
4639        assert!(matches!(
4640            Store::open(&path, 3_000),
4641            Err(StoreError::UnsupportedSchemaVersion(99))
4642        ));
4643    }
4644
4645    #[test]
4646    fn version_eight_migrates_existing_runs_with_null_snapshots() {
4647        let directory = tempdir().unwrap();
4648        let path = directory.path().join("sloop.db");
4649        let mut store = open_seeded(&path);
4650        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4651        drop(store);
4652
4653        let connection = rusqlite::Connection::open(&path).unwrap();
4654        connection
4655            .execute_batch(
4656                "ALTER TABLE runs DROP COLUMN flow_json;
4657                 ALTER TABLE runs DROP COLUMN ticket_json;
4658                 ALTER TABLE tickets DROP COLUMN body;
4659                 ALTER TABLE tickets DROP COLUMN held_reason;
4660                 ALTER TABLE scheduler_state DROP COLUMN draining;
4661                 ALTER TABLE runs DROP COLUMN cleanup_eligible_at_ms;
4662                 ALTER TABLE runs DROP COLUMN cleaned_at_ms;",
4663            )
4664            .unwrap();
4665        connection.pragma_update(None, "user_version", 8).unwrap();
4666        drop(connection);
4667
4668        let store = Store::open(&path, 3_000).unwrap();
4669        let run = store.run("R1").unwrap().unwrap();
4670        assert_eq!(run.flow_json, None);
4671        assert_eq!(run.ticket_json, None);
4672    }
4673
4674    #[test]
4675    fn version_ten_adds_source_metadata_without_disturbing_ticket_state() {
4676        let directory = tempdir().unwrap();
4677        let path = directory.path().join("sloop.db");
4678        let store = open_seeded(&path);
4679        store
4680            .connection
4681            .execute(
4682                "UPDATE tickets SET state = 'held', attempts = 3 WHERE id = 'T1'",
4683                [],
4684            )
4685            .unwrap();
4686        drop(store);
4687
4688        let connection = rusqlite::Connection::open(&path).unwrap();
4689        connection
4690            .execute_batch(
4691                "ALTER TABLE tickets DROP COLUMN body;
4692                 ALTER TABLE tickets DROP COLUMN held_reason;
4693                 ALTER TABLE scheduler_state DROP COLUMN draining;
4694                 ALTER TABLE runs DROP COLUMN cleanup_eligible_at_ms;
4695                 ALTER TABLE runs DROP COLUMN cleaned_at_ms;",
4696            )
4697            .unwrap();
4698        connection.pragma_update(None, "user_version", 10).unwrap();
4699        drop(connection);
4700
4701        let store = Store::open(&path, 3_000).unwrap();
4702        let ticket = store.ticket("T1").unwrap().unwrap();
4703        assert_eq!(ticket.state, "held");
4704        assert_eq!(ticket.attempts, 3);
4705        assert_eq!(ticket.body, None);
4706        assert_eq!(ticket.held_reason, None);
4707    }
4708
4709    #[test]
4710    fn configured_default_backfills_tickets_that_predate_target_snapshots() {
4711        let directory = tempdir().unwrap();
4712        let store = open_seeded(&directory.path().join("sloop.db"));
4713        store
4714            .update_ticket_execution("T1", None, Some("sonnet"), Some("medium"), 2_000)
4715            .unwrap();
4716
4717        assert_eq!(store.backfill_ticket_targets("codex", 3_000).unwrap(), 1);
4718        assert_eq!(
4719            store.ticket("T1").unwrap().unwrap().target.as_deref(),
4720            Some("codex")
4721        );
4722        assert_eq!(store.backfill_ticket_targets("claude", 4_000).unwrap(), 0);
4723        assert_eq!(
4724            store.ticket("T1").unwrap().unwrap().target.as_deref(),
4725            Some("codex")
4726        );
4727    }
4728
4729    #[test]
4730    fn tickets_with_unmerged_blockers_are_never_selected() {
4731        use crate::outcome::Outcome;
4732        let directory = tempdir().unwrap();
4733        let mut store = open_seeded(&directory.path().join("sloop.db"));
4734        store
4735            .insert_local_ticket(
4736                "T2",
4737                "default",
4738                ".agents/sloop/tickets/t2.md",
4739                "Ticket two",
4740                &["T1".into()],
4741                "sloop/T2",
4742                Some("claude"),
4743                Some("sonnet"),
4744                Some("medium"),
4745                "default",
4746                TicketState::Ready,
4747                1_500,
4748            )
4749            .unwrap();
4750        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4751
4752        let activation = super::QueuedActivation {
4753            id: "A1".into(),
4754            kind: "immediate".into(),
4755            ticket_id: None,
4756            project_id: None,
4757            eligible_at_ms: None,
4758            interval_ms: None,
4759        };
4760        // T1 is claimed and T2's blocker has not merged: nothing is ready.
4761        assert_eq!(store.select_ready_ticket(&activation, 2_000).unwrap(), None);
4762
4763        store
4764            .finish_run("R1", "T1", Some(0), Outcome::Merged, &[], None, 3_000)
4765            .unwrap();
4766        assert_eq!(
4767            store
4768                .select_ready_ticket(&activation, 3_000)
4769                .unwrap()
4770                .as_deref(),
4771            Some("T2")
4772        );
4773    }
4774
4775    #[test]
4776    fn finishing_a_run_settles_ticket_lease_and_evidence_atomically() {
4777        use crate::outcome::Outcome;
4778        let directory = tempdir().unwrap();
4779        let mut store = open_seeded(&directory.path().join("sloop.db"));
4780        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4781        store
4782            .record_aftercare_stage(
4783                "R1",
4784                &super::StageRecord {
4785                    stage_index: 0,
4786                    stage: "test".into(),
4787                    state: "passed".into(),
4788                    started_at_ms: 2_500,
4789                    finished_at_ms: 2_900,
4790                    exit_code: Some(0),
4791                    output_ref: "runs/R1/output.ndjson".into(),
4792                    verdict_source: "exit_code".into(),
4793                    reason: None,
4794                },
4795            )
4796            .unwrap();
4797
4798        store
4799            .finish_run(
4800                "R1",
4801                "T1",
4802                Some(0),
4803                Outcome::Merged,
4804                &[super::EvidenceRecord {
4805                    kind: "commits_observed",
4806                    data_json: "{\"oids\":[\"abc\",\"def\"]}".into(),
4807                }],
4808                None,
4809                3_000,
4810            )
4811            .unwrap();
4812
4813        assert_eq!(store.ticket_state("T1").unwrap().unwrap(), "merged");
4814        let run = store.run("R1").unwrap().unwrap();
4815        assert_eq!(run.state, "merged");
4816        assert_eq!(run.exit_code, Some(0));
4817        assert_eq!(run.exited_at_ms, Some(3_000));
4818        let evidence = store.run_evidence("R1").unwrap();
4819        assert_eq!(evidence[0].0, "commits_observed");
4820        assert_eq!(store.aftercare_stages("R1").unwrap()[0].stage, "test");
4821        // The lease is gone: the same run cannot renew it.
4822        assert!(store.renew_lease("T1", "R1", 60_000, 3_100).is_err());
4823    }
4824
4825    #[test]
4826    fn finishing_a_run_is_idempotent() {
4827        use crate::outcome::Outcome;
4828        let directory = tempdir().unwrap();
4829        let mut store = open_seeded(&directory.path().join("sloop.db"));
4830        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4831        let evidence = [super::EvidenceRecord {
4832            kind: "exit_classified",
4833            data_json: "{\"exit_code\":1}".into(),
4834        }];
4835
4836        store
4837            .finish_run("R1", "T1", Some(1), Outcome::Failed, &evidence, None, 3_000)
4838            .unwrap();
4839        store
4840            .finish_run("R1", "T1", Some(1), Outcome::Failed, &evidence, None, 3_100)
4841            .unwrap();
4842
4843        assert_eq!(store.run_evidence("R1").unwrap().len(), 1);
4844        assert_eq!(store.run("R1").unwrap().unwrap().exited_at_ms, Some(3_000));
4845    }
4846
4847    #[test]
4848    fn orphaning_a_run_releases_the_ticket_without_failing_it() {
4849        use crate::outcome::Outcome;
4850        let directory = tempdir().unwrap();
4851        let mut store = open_seeded(&directory.path().join("sloop.db"));
4852        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4853
4854        store
4855            .finish_run("R1", "T1", None, Outcome::Orphaned, &[], None, 3_000)
4856            .unwrap();
4857
4858        assert_eq!(store.run("R1").unwrap().unwrap().state, "orphaned");
4859        assert_eq!(store.ticket_state("T1").unwrap().as_deref(), Some("ready"));
4860    }
4861
4862    #[test]
4863    fn a_cancelled_outcome_returns_the_ticket_to_ready() {
4864        use crate::outcome::Outcome;
4865        let directory = tempdir().unwrap();
4866        let mut store = open_seeded(&directory.path().join("sloop.db"));
4867        store.claim_ticket(&claim_t1("R1"), 2_000).unwrap();
4868
4869        assert!(!store.cancellation_requested("R1").unwrap());
4870        store.record_cancel_requested("R1", 2_500).unwrap();
4871        store.record_cancel_requested("R1", 2_600).unwrap();
4872        assert!(store.cancellation_requested("R1").unwrap());
4873
4874        store
4875            .finish_run("R1", "T1", None, Outcome::Cancelled, &[], None, 3_000)
4876            .unwrap();
4877        assert_eq!(store.ticket_state("T1").unwrap().unwrap(), "ready");
4878        assert_eq!(store.ticket_counts().unwrap().ready, 1);
4879
4880        // Intent stayed deduplicated to one evidence row.
4881        let cancels = store
4882            .run_evidence("R1")
4883            .unwrap()
4884            .into_iter()
4885            .filter(|(kind, _)| kind == "cancel_requested")
4886            .count();
4887        assert_eq!(cancels, 1);
4888    }
4889
4890    #[test]
4891    fn paused_state_persists() {
4892        let directory = tempdir().unwrap();
4893        let path = directory.path().join("sloop.db");
4894
4895        let store = Store::open(&path, 1_000).unwrap();
4896        store.set_paused(true, 2_000).unwrap();
4897        drop(store);
4898
4899        assert!(Store::open(&path, 3_000).unwrap().paused().unwrap());
4900    }
4901
4902    #[test]
4903    fn restart_draining_is_durable_idempotent_and_cancelled_by_resume() {
4904        let directory = tempdir().unwrap();
4905        let path = directory.path().join("sloop.db");
4906        let mut store = Store::open(&path, 1_000).unwrap();
4907
4908        assert!(store.begin_restart_draining(2, 2_000).unwrap());
4909        assert!(!store.begin_restart_draining(2, 2_100).unwrap());
4910        assert!(store.restart_draining().unwrap());
4911        assert_eq!(
4912            store
4913                .events_after(0, 10)
4914                .unwrap()
4915                .iter()
4916                .filter(|event| event.kind == "daemon_restart_requested")
4917                .count(),
4918            1
4919        );
4920        drop(store);
4921
4922        let mut reopened = Store::open(&path, 3_000).unwrap();
4923        assert!(reopened.restart_draining().unwrap());
4924        assert!(reopened.resume_scheduler(4_000).unwrap());
4925        assert!(!reopened.restart_draining().unwrap());
4926    }
4927
4928    #[test]
4929    fn version_eleven_adds_restart_draining_state() {
4930        let directory = tempdir().unwrap();
4931        let path = directory.path().join("sloop.db");
4932        drop(Store::open(&path, 1_000).unwrap());
4933        let connection = Connection::open(&path).unwrap();
4934        connection
4935            .execute_batch(
4936                "ALTER TABLE scheduler_state DROP COLUMN draining;
4937                 ALTER TABLE runs DROP COLUMN cleanup_eligible_at_ms;
4938                 ALTER TABLE runs DROP COLUMN cleaned_at_ms;
4939                 PRAGMA user_version = 11;",
4940            )
4941            .unwrap();
4942        drop(connection);
4943
4944        let store = Store::open(&path, 2_000).unwrap();
4945        assert!(!store.restart_draining().unwrap());
4946        assert_eq!(
4947            store
4948                .connection
4949                .query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0))
4950                .unwrap(),
4951            SCHEMA_VERSION
4952        );
4953    }
4954}