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