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