Skip to main content

sloop/
store.rs

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