Skip to main content

sloop/work_state/
local.rs

1use std::collections::{BTreeMap, BTreeSet};
2use std::fs;
3use std::path::Path;
4use std::sync::Arc;
5use std::sync::atomic::{AtomicI64, Ordering};
6use std::time::Duration;
7
8use async_trait::async_trait;
9use rusqlite::{Connection, OptionalExtension, Transaction, TransactionBehavior, params};
10use serde_json::{Value, json};
11
12use crate::clock::{Clock, SystemClock};
13use crate::config::{AgentConfig, expand_agent_cmd};
14use crate::db::{Db, StoreError};
15use crate::domain::ticket::TicketState;
16use crate::domain::trigger as domain_trigger;
17use crate::domain::work::{
18    Disposition, ExecutionHints, OwnerId, SourceVersion, TicketRef, WorkOutcome, WorkTicket,
19    WorkTicketState,
20};
21use crate::flow::Flow;
22use crate::frontmatter;
23use crate::ids::next_id;
24use crate::work_state::exec::ExecTicketSource;
25use crate::work_state::trigger;
26use crate::work_state::{
27    ActiveClaim, ClaimResult, ClaimStrength, ReindexError, SourceError, TicketFeeder, WorkState,
28};
29
30const TICKET_RECORD_SELECT: &str =
31    "SELECT id, project_id, file_path, source, source_ref, state, name, worktree,
32            target, model, effort, flow, attempts, body, held_reason, created_at_ms, updated_at_ms
33     FROM tickets";
34
35#[derive(Debug, Clone, PartialEq, Eq)]
36pub(crate) struct ClaimTransaction<'a> {
37    pub(crate) ticket_id: &'a str,
38    pub(crate) run_id: &'a str,
39    pub(crate) trigger_id: &'a str,
40    pub(crate) owner_id: &'a str,
41    pub(crate) lease_ms: i64,
42}
43
44#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
45pub struct TicketCounts {
46    pub ready: u64,
47    pub held: u64,
48    pub blocked: u64,
49    pub claimed: u64,
50    pub merged: u64,
51    pub failed: u64,
52    pub needs_review: u64,
53}
54
55#[derive(Debug, Clone, PartialEq, Eq)]
56pub struct ProjectRecord {
57    pub id: String,
58    pub file_path: Option<String>,
59    pub title: String,
60}
61
62#[derive(Debug, Clone, PartialEq, Eq)]
63pub struct LocalTicketFile {
64    pub id: String,
65    pub file_path: String,
66    pub state: String,
67    pub missing_at_ms: Option<i64>,
68}
69
70#[derive(Debug, Clone, PartialEq, Eq)]
71pub struct TicketRecord {
72    pub id: String,
73    pub project_id: String,
74    pub file_path: Option<String>,
75    pub source: String,
76    pub source_ref: Option<String>,
77    pub state: String,
78    pub name: String,
79    pub blocked_by: Vec<String>,
80    pub worktree: Option<String>,
81    pub target: Option<String>,
82    pub model: Option<String>,
83    pub effort: Option<String>,
84    pub flow: Option<String>,
85    pub attempts: i64,
86    pub body: Option<String>,
87    pub held_reason: Option<String>,
88    /// When the ticket was registered. `sloop list` orders on this.
89    pub created_at_ms: i64,
90    pub updated_at_ms: i64,
91}
92
93#[derive(Debug, Clone, PartialEq, Eq)]
94pub struct ReindexTicket {
95    pub id: String,
96    pub project_id: String,
97    pub source: String,
98    pub source_ref: String,
99    pub file_path: Option<String>,
100    pub name: String,
101    pub blocked_by: Vec<String>,
102    pub worktree: String,
103    pub target: Option<String>,
104    pub model: Option<String>,
105    pub effort: Option<String>,
106    pub flow: String,
107    pub body: String,
108    pub held_reason: Option<String>,
109    pub derived_state: Option<TicketState>,
110}
111
112#[derive(Debug, Clone, PartialEq, Eq)]
113pub struct ReindexStateChange {
114    pub ticket_id: String,
115    pub previous_state: String,
116    pub state: String,
117}
118
119#[derive(Debug, Clone, Default, PartialEq, Eq)]
120pub struct ReindexResult {
121    pub state_changes: Vec<ReindexStateChange>,
122    pub rows_dropped: usize,
123}
124
125pub(crate) struct LocalTicketWrite<'a> {
126    pub id: &'a str,
127    pub project_id: &'a str,
128    pub file_path: &'a str,
129    pub name: &'a str,
130    pub blocked_by: &'a [String],
131    pub worktree: &'a str,
132    pub target: Option<&'a str>,
133    pub model: Option<&'a str>,
134    pub effort: Option<&'a str>,
135    pub flow: &'a str,
136    pub state: TicketState,
137    pub body: &'a str,
138    pub content_hash: &'a str,
139    pub now_ms: i64,
140}
141
142fn ticket_record(row: &rusqlite::Row<'_>) -> rusqlite::Result<TicketRecord> {
143    Ok(TicketRecord {
144        id: row.get(0)?,
145        project_id: row.get(1)?,
146        file_path: row.get(2)?,
147        source: row.get(3)?,
148        source_ref: row.get(4)?,
149        state: row.get(5)?,
150        name: row.get(6)?,
151        blocked_by: Vec::new(),
152        worktree: row.get(7)?,
153        target: row.get(8)?,
154        model: row.get(9)?,
155        effort: row.get(10)?,
156        flow: row.get(11)?,
157        attempts: row.get(12)?,
158        body: row.get(13)?,
159        held_reason: row.get(14)?,
160        created_at_ms: row.get(15)?,
161        updated_at_ms: row.get(16)?,
162    })
163}
164
165pub(crate) mod tx {
166    use super::*;
167
168    pub(crate) fn insert_authored_ticket(
169        transaction: &Transaction<'_>,
170        ticket: &LocalTicketWrite<'_>,
171    ) -> Result<(), StoreError> {
172        transaction.execute(
173            "INSERT INTO tickets
174                 (id, project_id, file_path, source, state, name, worktree, target, model, effort,
175                    flow, body, content_hash, created_at_ms, updated_at_ms)
176             VALUES (?1, ?2, ?3, 'local', ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?13)",
177            params![
178                ticket.id,
179                ticket.project_id,
180                ticket.file_path,
181                ticket.state.as_str(),
182                ticket.name,
183                ticket.worktree,
184                ticket.target,
185                ticket.model,
186                ticket.effort,
187                ticket.flow,
188                ticket.body,
189                ticket.content_hash,
190                ticket.now_ms,
191            ],
192        )?;
193        replace_ticket_blockers(transaction, ticket.id, ticket.blocked_by)?;
194        Ok(())
195    }
196
197    pub(crate) fn update_authored_ticket(
198        transaction: &Transaction<'_>,
199        ticket: &LocalTicketWrite<'_>,
200    ) -> Result<(), StoreError> {
201        transaction.execute(
202            "UPDATE tickets
203             SET name = ?2, worktree = ?3, target = ?4, model = ?5, effort = ?6, flow = ?7,
204                 body = ?8, content_hash = ?9, held_reason = NULL, missing_at_ms = NULL,
205                 updated_at_ms = ?10
206             WHERE id = ?1",
207            params![
208                ticket.id,
209                ticket.name,
210                ticket.worktree,
211                ticket.target,
212                ticket.model,
213                ticket.effort,
214                ticket.flow,
215                ticket.body,
216                ticket.content_hash,
217                ticket.now_ms,
218            ],
219        )?;
220        replace_ticket_blockers(transaction, ticket.id, ticket.blocked_by)?;
221        Ok(())
222    }
223
224    pub(crate) fn insert_lease(
225        transaction: &Transaction<'_>,
226        claim: &ClaimTransaction<'_>,
227        now_ms: i64,
228    ) -> Result<i64, StoreError> {
229        let expires_at_ms = now_ms + claim.lease_ms;
230        transaction.execute(
231            "INSERT INTO leases
232                 (ticket_id, run_id, owner_id, acquired_at_ms, renewed_at_ms, expires_at_ms)
233             VALUES (?1, ?2, ?3, ?4, ?4, ?5)",
234            params![
235                claim.ticket_id,
236                claim.run_id,
237                claim.owner_id,
238                now_ms,
239                expires_at_ms,
240            ],
241        )?;
242        Ok(expires_at_ms)
243    }
244
245    pub(crate) fn delete_lease(
246        transaction: &Transaction<'_>,
247        run_id: &str,
248    ) -> rusqlite::Result<usize> {
249        transaction.execute("DELETE FROM leases WHERE run_id = ?1", params![run_id])
250    }
251
252    pub(crate) fn replace_ticket_blockers(
253        transaction: &Transaction<'_>,
254        ticket_id: &str,
255        blocked_by: &[String],
256    ) -> rusqlite::Result<()> {
257        transaction.execute(
258            "DELETE FROM ticket_blockers WHERE ticket_id = ?1",
259            params![ticket_id],
260        )?;
261        for (position, blocker_id) in blocked_by.iter().enumerate() {
262            transaction.execute(
263                "INSERT OR IGNORE INTO ticket_blockers (ticket_id, blocker_id, position)
264                 VALUES (?1, ?2, ?3)",
265                params![ticket_id, blocker_id, position as i64],
266            )?;
267        }
268        Ok(())
269    }
270
271    pub(crate) fn claim_ticket(
272        transaction: &Transaction<'_>,
273        claim: &ClaimTransaction<'_>,
274        now_ms: i64,
275    ) -> Result<(), StoreError> {
276        let changed = transaction.execute(
277            "UPDATE tickets
278             SET state = 'claimed', held_reason = NULL, attempts = attempts + 1, updated_at_ms = ?2
279             WHERE id = ?1 AND state = 'ready' AND missing_at_ms IS NULL
280               AND NOT EXISTS (SELECT 1 FROM ticket_blockers b
281                               JOIN tickets bt ON bt.id = b.blocker_id
282                               WHERE b.ticket_id = tickets.id
283                                 AND bt.state != 'merged')",
284            params![claim.ticket_id, now_ms],
285        )?;
286        if changed != 1 {
287            let state: Option<String> = transaction
288                .query_row(
289                    "SELECT CASE
290                              WHEN missing_at_ms IS NOT NULL THEN 'missing'
291                              WHEN state = 'ready' AND EXISTS (
292                                  SELECT 1 FROM ticket_blockers b
293                                  JOIN tickets bt ON bt.id = b.blocker_id
294                                  WHERE b.ticket_id = tickets.id
295                                    AND bt.state != 'merged'
296                              ) THEN 'blocked'
297                              ELSE state
298                            END
299                     FROM tickets WHERE id = ?1",
300                    params![claim.ticket_id],
301                    |row| row.get(0),
302                )
303                .optional()?;
304            return Err(StoreError::TicketNotReady {
305                ticket_id: claim.ticket_id.into(),
306                state,
307            });
308        }
309        Ok(())
310    }
311
312    pub(crate) fn settle_ticket(
313        transaction: &Transaction<'_>,
314        ticket_id: &str,
315        ticket_state: TicketState,
316        now_ms: i64,
317    ) -> rusqlite::Result<usize> {
318        transaction.execute(
319            "UPDATE tickets SET state = ?2, held_reason = NULL, updated_at_ms = ?3
320             WHERE id = ?1 AND state = 'claimed'",
321            params![ticket_id, ticket_state.as_str(), now_ms],
322        )
323    }
324
325    pub(crate) fn abort_ticket(
326        transaction: &Transaction<'_>,
327        ticket_id: &str,
328        now_ms: i64,
329    ) -> rusqlite::Result<usize> {
330        transaction.execute(
331            "UPDATE tickets SET state = 'ready', held_reason = NULL, updated_at_ms = ?2
332             WHERE id = ?1 AND state = 'claimed'",
333            params![ticket_id, now_ms],
334        )
335    }
336
337    pub(crate) fn settle_external_merge(
338        transaction: &Transaction<'_>,
339        ticket_id: &str,
340        now_ms: i64,
341    ) -> rusqlite::Result<usize> {
342        transaction.execute(
343            "UPDATE tickets SET state = 'merged', held_reason = NULL, updated_at_ms = ?2
344             WHERE id = ?1 AND state = 'needs_review'",
345            params![ticket_id, now_ms],
346        )
347    }
348
349    pub(crate) fn retry_ticket(
350        transaction: &Transaction<'_>,
351        id: &str,
352        now_ms: i64,
353    ) -> rusqlite::Result<usize> {
354        transaction.execute(
355            "UPDATE tickets SET state = 'ready', held_reason = NULL, attempts = 0, updated_at_ms = ?2
356             WHERE id = ?1 AND state = 'failed'",
357            params![id, now_ms],
358        )
359    }
360}
361
362pub struct LocalSqlite {
363    db: Db,
364    clock: Arc<dyn Clock>,
365    last_sync_ms: AtomicI64,
366    outcome_reporter: Option<ExecTicketSource>,
367}
368
369impl LocalSqlite {
370    pub fn from_db(db: Db) -> Self {
371        Self::from_db_with_clock(db, Arc::new(SystemClock))
372    }
373
374    pub fn from_db_with_clock(db: Db, clock: Arc<dyn Clock>) -> Self {
375        Self::from_db_with_clock_and_reporter(db, clock, None)
376    }
377
378    pub(crate) fn from_db_with_clock_and_reporter(
379        db: Db,
380        clock: Arc<dyn Clock>,
381        outcome_reporter: Option<ExecTicketSource>,
382    ) -> Self {
383        Self {
384            db,
385            clock,
386            last_sync_ms: AtomicI64::new(i64::MIN),
387            outcome_reporter,
388        }
389    }
390
391    pub(crate) fn db(&self) -> Db {
392        self.db.clone()
393    }
394
395    pub fn last_sync_ms(&self) -> Option<i64> {
396        match self.last_sync_ms.load(Ordering::Acquire) {
397            i64::MIN => None,
398            timestamp => Some(timestamp),
399        }
400    }
401
402    pub fn insert_local_project(
403        &self,
404        id: &str,
405        file_path: &str,
406        title: &str,
407        now_ms: i64,
408    ) -> Result<(), StoreError> {
409        self.db.lock().execute(
410            "INSERT INTO projects (id, file_path, source, title, created_at_ms, updated_at_ms)
411             VALUES (?1, ?2, 'local', ?3, ?4, ?4)",
412            params![id, file_path, title, now_ms],
413        )?;
414        Ok(())
415    }
416
417    /// Inserts or refreshes a project indexed from a committed file. Startup
418    /// and reindex call this for every configured project file, so it must tolerate
419    /// rows that already exist.
420    pub fn upsert_local_project(
421        &self,
422        id: &str,
423        file_path: &str,
424        title: &str,
425        now_ms: i64,
426    ) -> Result<(), StoreError> {
427        self.db.lock().execute(
428            "INSERT INTO projects (id, file_path, source, title, created_at_ms, updated_at_ms)
429             VALUES (?1, ?2, 'local', ?3, ?4, ?4)
430             ON CONFLICT(id) DO UPDATE SET
431                 file_path = excluded.file_path,
432                 title = excluded.title,
433                 updated_at_ms = excluded.updated_at_ms",
434            params![id, file_path, title, now_ms],
435        )?;
436        Ok(())
437    }
438
439    pub fn project_exists(&self, id: &str) -> Result<bool, StoreError> {
440        let found: Option<i64> = self
441            .db
442            .lock()
443            .query_row("SELECT 1 FROM projects WHERE id = ?1", params![id], |row| {
444                row.get(0)
445            })
446            .optional()?;
447        Ok(found.is_some())
448    }
449
450    pub fn project(&self, id: &str) -> Result<Option<ProjectRecord>, StoreError> {
451        self.db
452            .lock()
453            .query_row(
454                "SELECT id, file_path, title FROM projects WHERE id = ?1",
455                params![id],
456                |row| {
457                    Ok(ProjectRecord {
458                        id: row.get(0)?,
459                        file_path: row.get(1)?,
460                        title: row.get(2)?,
461                    })
462                },
463            )
464            .optional()
465            .map_err(StoreError::from)
466    }
467
468    #[allow(clippy::too_many_arguments)]
469    pub fn insert_local_ticket(
470        &self,
471        id: &str,
472        project_id: &str,
473        file_path: &str,
474        name: &str,
475        blocked_by: &[String],
476        worktree: &str,
477        target: Option<&str>,
478        model: Option<&str>,
479        effort: Option<&str>,
480        flow: &str,
481        state: TicketState,
482        now_ms: i64,
483    ) -> Result<(), StoreError> {
484        let mut connection = self.db.lock();
485        let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
486        transaction.execute(
487            "INSERT INTO tickets
488                 (id, project_id, file_path, source, state, name, worktree, target, model, effort,
489                    flow, created_at_ms, updated_at_ms)
490             VALUES (?1, ?2, ?3, 'local', ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?11)",
491            params![
492                id,
493                project_id,
494                file_path,
495                state.as_str(),
496                name,
497                worktree,
498                target,
499                model,
500                effort,
501                flow,
502                now_ms
503            ],
504        )?;
505        tx::replace_ticket_blockers(&transaction, id, blocked_by)?;
506        transaction.commit()?;
507        Ok(())
508    }
509
510    #[allow(clippy::too_many_arguments)]
511    pub fn update_local_ticket(
512        &self,
513        id: &str,
514        name: &str,
515        blocked_by: &[String],
516        worktree: &str,
517        target: Option<&str>,
518        model: Option<&str>,
519        effort: Option<&str>,
520        flow: &str,
521        now_ms: i64,
522    ) -> Result<(), StoreError> {
523        let mut connection = self.db.lock();
524        let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
525        transaction.execute(
526            "UPDATE tickets
527             SET name = ?2, worktree = ?3, target = ?4, model = ?5, effort = ?6, flow = ?7,
528                  held_reason = NULL, missing_at_ms = NULL, updated_at_ms = ?8
529             WHERE id = ?1",
530            params![id, name, worktree, target, model, effort, flow, now_ms],
531        )?;
532        tx::replace_ticket_blockers(&transaction, id, blocked_by)?;
533        transaction.commit()?;
534        Ok(())
535    }
536
537    /// Pulls and applies a complete authored ticket snapshot. Invalid authored
538    /// tickets remain visible on hold; only source-wide identity conflicts and
539    /// dependency cycles reject the pull.
540    #[allow(clippy::too_many_arguments)]
541    pub(crate) fn sync_from_source<DropRuns, MarkRuns>(
542        &self,
543        root: &Path,
544        ticket_source: &TicketFeeder,
545        worktree_dir: &Path,
546        now_ms: i64,
547        ticket_prefix: &str,
548        project_ids: &[String],
549        agent: Option<&AgentConfig>,
550        flows: &BTreeMap<String, Flow>,
551        default_flow: &str,
552        drop_runs: DropRuns,
553        mark_runs_cleanup_eligible: MarkRuns,
554    ) -> Result<Value, ReindexError>
555    where
556        DropRuns:
557            FnMut(&Transaction<'_>, &[String], &BTreeSet<String>) -> Result<usize, StoreError>,
558        MarkRuns: FnMut(&Transaction<'_>, &str, i64) -> Result<(), StoreError>,
559    {
560        let authored = ticket_source
561            .pull()
562            .map_err(|error| ReindexError(error.to_string()))?;
563        let mut known_ids: Vec<String> = authored
564            .iter()
565            .filter_map(|ticket| ticket.frontmatter.id.clone())
566            .collect();
567        let mut unique_ids = BTreeSet::new();
568        for id in &known_ids {
569            if !unique_ids.insert(id.clone()) {
570                return Err(ReindexError(format!(
571                    "duplicate ticket ID `{id}` in the configured ticket directory"
572                )));
573            }
574        }
575        let mut unique_refs = BTreeSet::new();
576        for ticket in &authored {
577            if !unique_refs.insert((&ticket.source, &ticket.source_ref)) {
578                return Err(ReindexError(format!(
579                    "duplicate source reference `{}` from `{}`",
580                    ticket.source_ref, ticket.source
581                )));
582            }
583        }
584        let known_projects: BTreeSet<&str> = project_ids.iter().map(String::as_str).collect();
585        let fallback_project = if known_projects.contains("default") {
586            "default".to_owned()
587        } else {
588            project_ids.first().cloned().ok_or_else(|| {
589                ReindexError("cannot index tickets without an indexed project".into())
590            })?
591        };
592
593        let mut tickets = Vec::with_capacity(authored.len());
594        let mut assigned_ids = BTreeSet::new();
595        for authored_ticket in authored {
596            let prior_id = self
597                .ticket_by_source_ref(&authored_ticket.source, &authored_ticket.source_ref)
598                .map_err(|error| ReindexError(error.to_string()))?
599                .map(|ticket| ticket.id);
600            let id = match authored_ticket.frontmatter.id.clone().or(prior_id) {
601                Some(id) => id,
602                None => {
603                    let id = next_id(ticket_prefix, known_ids.iter().map(String::as_str))
604                        .map_err(|error| ReindexError(error.to_string()))?;
605                    known_ids.push(id.clone());
606                    id
607                }
608            };
609            if !assigned_ids.insert(id.clone()) {
610                return Err(ReindexError(format!(
611                    "duplicate ticket ID `{id}` in the pulled ticket source"
612                )));
613            }
614            if !known_ids.contains(&id) {
615                known_ids.push(id.clone());
616            }
617            let authored_project = authored_ticket
618                .frontmatter
619                .project
620                .clone()
621                .unwrap_or_else(|| "default".to_owned());
622            let mut held_reason = authored_ticket.validation_error.clone();
623            if authored_ticket.frontmatter.name.trim().is_empty() {
624                held_reason.get_or_insert_with(|| {
625                    format!(
626                        "{}: frontmatter field `name` is required and must be non-empty",
627                        authored_ticket.source_ref
628                    )
629                });
630            }
631            if authored_ticket.body.trim().is_empty() {
632                held_reason.get_or_insert_with(|| {
633                    format!(
634                        "{}: ticket body must be non-empty",
635                        authored_ticket.source_ref
636                    )
637                });
638            }
639            let project = if known_projects.contains(authored_project.as_str()) {
640                authored_project.clone()
641            } else {
642                held_reason.get_or_insert_with(|| {
643                    format!(
644                        "{}: project `{authored_project}` is not indexed",
645                        authored_ticket.source_ref
646                    )
647                });
648                fallback_project.clone()
649            };
650            let flow = authored_ticket
651                .frontmatter
652                .flow
653                .clone()
654                .unwrap_or_else(|| default_flow.to_owned());
655            if !flows.contains_key(&flow) {
656                held_reason.get_or_insert_with(|| {
657                    format!(
658                        "{}: flow `{flow}` is not defined",
659                        authored_ticket.source_ref
660                    )
661                });
662            }
663            let target = match authored_ticket.frontmatter.target.as_deref() {
664                Some(target) if agent.is_some_and(|agent| agent.targets.contains_key(target)) => {
665                    Some(target.to_owned())
666                }
667                Some(target) => {
668                    held_reason.get_or_insert_with(|| {
669                        format!(
670                            "{}: agent target `{target}` is not configured",
671                            authored_ticket.source_ref
672                        )
673                    });
674                    Some(target.to_owned())
675                }
676                None => agent.map(|agent| agent.default_target.clone()),
677            };
678            if let (Some(agent), Some(target)) = (agent, target.as_deref()) {
679                if let Some(command) = agent.targets.get(target)
680                    && let Err(message) = expand_agent_cmd(
681                        command,
682                        authored_ticket.frontmatter.model.as_deref(),
683                        authored_ticket.frontmatter.effort.as_deref(),
684                        "",
685                    )
686                {
687                    held_reason.get_or_insert_with(|| {
688                        format!(
689                            "{}: ticket using agent target `{target}` {message}",
690                            authored_ticket.source_ref
691                        )
692                    });
693                }
694            }
695            let worktree = match authored_ticket.frontmatter.worktree.clone() {
696                Some(worktree) => worktree,
697                None => {
698                    let stem = authored_ticket
699                        .file_path
700                        .as_deref()
701                        .and_then(Path::file_stem)
702                        .and_then(|stem| stem.to_str());
703                    match crate::ids::default_worktree(stem, &id) {
704                        Ok(branch) => branch,
705                        Err(reason) => {
706                            held_reason.get_or_insert_with(|| {
707                                format!("{}: {reason}", authored_ticket.source_ref)
708                            });
709                            format!("sloop/{id}")
710                        }
711                    }
712                }
713            };
714            if held_reason.is_none()
715                && let (Some(path), Some(content)) = (
716                    authored_ticket.file_path.as_ref(),
717                    authored_ticket.original_content.as_ref(),
718                )
719                && let Some(updated) = frontmatter::stamp(content, &id, &project, &worktree, &flow)
720                    .map_err(|error| {
721                        ReindexError(format!("{}: {error}", authored_ticket.source_ref))
722                    })?
723            {
724                let absolute = root.join(path);
725                fs::write(&absolute, updated)
726                    .map_err(|source| ReindexError::io(&absolute, source))?;
727            }
728
729            tickets.push(ReindexTicket {
730                id,
731                project_id: project,
732                source: authored_ticket.source,
733                source_ref: authored_ticket.source_ref,
734                file_path: authored_ticket
735                    .file_path
736                    .map(|path| path.to_string_lossy().into_owned()),
737                name: authored_ticket.frontmatter.name,
738                blocked_by: authored_ticket.frontmatter.blocked_by,
739                worktree,
740                target,
741                model: authored_ticket.frontmatter.model,
742                effort: authored_ticket.frontmatter.effort,
743                flow,
744                body: authored_ticket.body,
745                held_reason,
746                derived_state: None,
747            });
748        }
749
750        let ticket_ids: BTreeSet<String> = tickets.iter().map(|ticket| ticket.id.clone()).collect();
751        let mut dependencies = BTreeMap::new();
752        for ticket in &mut tickets {
753            let unknown_blocker = ticket
754                .blocked_by
755                .iter()
756                .find(|blocker| !ticket_ids.contains(*blocker))
757                .cloned();
758            if let Some(blocker) = unknown_blocker {
759                ticket.held_reason.get_or_insert_with(|| {
760                    format!(
761                        "ticket `{}` field `blocked_by` references unknown ticket `{blocker}`; edit `{}` to drop or correct the reference",
762                        ticket.id, ticket.source_ref
763                    )
764                });
765                ticket.blocked_by.clear();
766            } else {
767                dependencies.insert(ticket.id.clone(), ticket.blocked_by.clone());
768            }
769        }
770        if let Some(chain) = crate::domain::graph::find_cycle(&dependencies) {
771            return Err(ReindexError(format!(
772                "field `blocked_by` creates a dependency cycle: {}",
773                chain.join(" -> ")
774            )));
775        }
776
777        crate::work_state::reindex_evidence::derive_states(root, worktree_dir, &mut tickets)?;
778        let result = self
779            .apply_reindex(
780                project_ids,
781                &tickets,
782                now_ms,
783                drop_runs,
784                mark_runs_cleanup_eligible,
785            )
786            .map_err(|error| ReindexError(error.to_string()))?;
787        self.last_sync_ms.store(now_ms, Ordering::Release);
788        let state_changes: Vec<Value> = result
789            .state_changes
790            .into_iter()
791            .map(|change| {
792                json!({
793                    "ticket": change.ticket_id,
794                    "previous_state": change.previous_state,
795                    "state": change.state,
796                })
797            })
798            .collect();
799        Ok(json!({
800            "projects_indexed": project_ids.len(),
801            "tickets_indexed": tickets.len(),
802            "tickets_state_changed": state_changes.len(),
803            "state_changes": state_changes,
804            "rows_dropped": result.rows_dropped,
805        }))
806    }
807
808    /// Applies a complete authored ticket snapshot without disturbing runtime
809    /// history for IDs that remain present. Cross-store cleanup is supplied by
810    /// coordination and runs in this transaction.
811    pub(crate) fn apply_reindex<DropRuns, MarkRuns>(
812        &self,
813        project_ids: &[String],
814        tickets: &[ReindexTicket],
815        now_ms: i64,
816        mut drop_runs: DropRuns,
817        mut mark_runs_cleanup_eligible: MarkRuns,
818    ) -> Result<ReindexResult, StoreError>
819    where
820        DropRuns:
821            FnMut(&Transaction<'_>, &[String], &BTreeSet<String>) -> Result<usize, StoreError>,
822        MarkRuns: FnMut(&Transaction<'_>, &str, i64) -> Result<(), StoreError>,
823    {
824        let mut connection = self.db.lock();
825        let existing: BTreeMap<String, TicketRecord> = Self::tickets_on(&connection)?
826            .into_iter()
827            .map(|ticket| (ticket.id.clone(), ticket))
828            .collect();
829        let desired_ticket_ids: BTreeSet<&str> =
830            tickets.iter().map(|ticket| ticket.id.as_str()).collect();
831        let desired_project_ids: BTreeSet<&str> = project_ids.iter().map(String::as_str).collect();
832        let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
833
834        let stale_tickets = {
835            let mut statement = transaction.prepare("SELECT id FROM tickets ORDER BY id")?;
836            statement
837                .query_map([], |row| row.get::<_, String>(0))?
838                .collect::<Result<Vec<_>, _>>()?
839                .into_iter()
840                .filter(|id| !desired_ticket_ids.contains(id.as_str()))
841                .collect::<Vec<_>>()
842        };
843        let stale_projects = {
844            let mut statement = transaction
845                .prepare("SELECT id FROM projects WHERE source = 'local' ORDER BY id")?;
846            statement
847                .query_map([], |row| row.get::<_, String>(0))?
848                .collect::<Result<Vec<_>, _>>()?
849                .into_iter()
850                .filter(|id| !desired_project_ids.contains(id.as_str()))
851                .collect::<Vec<_>>()
852        };
853
854        let plan = trigger::deletion_plan(&transaction, &stale_tickets, &stale_projects)?;
855
856        let mut rows_dropped = drop_runs(&transaction, &stale_tickets, &plan.doomed)?;
857        rows_dropped += trigger::apply_deletion(&transaction, &plan)?;
858        for ticket_id in &stale_tickets {
859            rows_dropped += trigger::delete_filters_for_ticket(&transaction, ticket_id)?;
860            rows_dropped += transaction.execute(
861                "DELETE FROM leases WHERE ticket_id = ?1",
862                params![ticket_id],
863            )?;
864            rows_dropped += transaction.execute(
865                "DELETE FROM ticket_blockers WHERE ticket_id = ?1 OR blocker_id = ?1",
866                params![ticket_id],
867            )?;
868            rows_dropped +=
869                transaction.execute("DELETE FROM tickets WHERE id = ?1", params![ticket_id])?;
870        }
871        let mut state_changes = Vec::new();
872        for ticket in tickets {
873            let previous = existing.get(&ticket.id);
874            let state = if ticket.held_reason.is_some() {
875                TicketState::Held.as_str()
876            } else {
877                match (previous, ticket.derived_state) {
878                    (Some(_), Some(derived)) => derived.as_str(),
879                    (Some(existing), None) if existing.held_reason.is_some() => {
880                        TicketState::Ready.as_str()
881                    }
882                    (Some(existing), None) => existing.state.as_str(),
883                    (None, Some(derived)) => derived.as_str(),
884                    (None, None) => TicketState::Ready.as_str(),
885                }
886            };
887            if let Some(previous) = previous
888                && previous.state != state
889            {
890                state_changes.push(ReindexStateChange {
891                    ticket_id: ticket.id.clone(),
892                    previous_state: previous.state.clone(),
893                    state: state.to_owned(),
894                });
895                if state == TicketState::Merged.as_str()
896                    && matches!(previous.state.as_str(), "failed" | "needs_review")
897                {
898                    mark_runs_cleanup_eligible(&transaction, &ticket.id, now_ms)?;
899                }
900            }
901            transaction.execute(
902                "INSERT INTO tickets
903                     (id, project_id, file_path, source, source_ref, state, name, worktree, target,
904                      model, effort, flow, body, held_reason, created_at_ms, updated_at_ms)
905                 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?15)
906                 ON CONFLICT(id) DO UPDATE SET
907                     project_id = excluded.project_id,
908                     file_path = excluded.file_path,
909                     source = excluded.source,
910                     source_ref = excluded.source_ref,
911                     state = excluded.state,
912                     name = excluded.name,
913                     worktree = excluded.worktree,
914                     target = excluded.target,
915                     model = excluded.model,
916                     effort = excluded.effort,
917                     flow = excluded.flow,
918                     body = excluded.body,
919                     held_reason = excluded.held_reason,
920                     missing_at_ms = NULL,
921                     updated_at_ms = excluded.updated_at_ms",
922                params![
923                    ticket.id,
924                    ticket.project_id,
925                    ticket.file_path,
926                    ticket.source,
927                    ticket.source_ref,
928                    state,
929                    ticket.name,
930                    ticket.worktree,
931                    ticket.target,
932                    ticket.model,
933                    ticket.effort,
934                    ticket.flow,
935                    ticket.body,
936                    ticket.held_reason,
937                    now_ms,
938                ],
939            )?;
940        }
941        for ticket in tickets {
942            tx::replace_ticket_blockers(&transaction, &ticket.id, &ticket.blocked_by)?;
943        }
944        for project_id in &stale_projects {
945            rows_dropped +=
946                transaction.execute("DELETE FROM projects WHERE id = ?1", params![project_id])?;
947        }
948
949        state_changes.sort_by(|left, right| left.ticket_id.cmp(&right.ticket_id));
950        transaction.commit()?;
951        Ok(ReindexResult {
952            state_changes,
953            rows_dropped,
954        })
955    }
956
957    pub fn update_ticket_execution(
958        &self,
959        id: &str,
960        target: Option<&str>,
961        model: Option<&str>,
962        effort: Option<&str>,
963        now_ms: i64,
964    ) -> Result<(), StoreError> {
965        self.db.lock().execute(
966            "UPDATE tickets SET target = ?2, model = ?3, effort = ?4, updated_at_ms = ?5 WHERE id = ?1",
967            params![id, target, model, effort, now_ms],
968        )?;
969        Ok(())
970    }
971
972    pub fn update_ticket_body(&self, id: &str, body: &str, now_ms: i64) -> Result<(), StoreError> {
973        self.db.lock().execute(
974            "UPDATE tickets SET body = ?2, updated_at_ms = ?3 WHERE id = ?1",
975            params![id, body, now_ms],
976        )?;
977        Ok(())
978    }
979
980    pub fn backfill_ticket_targets(
981        &self,
982        default_target: &str,
983        now_ms: i64,
984    ) -> Result<usize, StoreError> {
985        self.db
986            .lock()
987            .execute(
988                "UPDATE tickets SET target = ?1, updated_at_ms = ?2 WHERE target IS NULL",
989                params![default_target, now_ms],
990            )
991            .map_err(StoreError::from)
992    }
993
994    pub fn ticket(&self, id: &str) -> Result<Option<TicketRecord>, StoreError> {
995        let connection = self.db.lock();
996        let mut ticket = connection
997            .query_row(
998                &format!("{TICKET_RECORD_SELECT} WHERE id = ?1"),
999                params![id],
1000                ticket_record,
1001            )
1002            .optional()?;
1003        if let Some(ticket) = ticket.as_mut() {
1004            ticket.blocked_by = Self::ticket_blockers(&connection, &ticket.id)?;
1005        }
1006        Ok(ticket)
1007    }
1008
1009    pub fn ticket_by_name(&self, name: &str) -> Result<Option<TicketRecord>, StoreError> {
1010        let connection = self.db.lock();
1011        let mut ticket = connection
1012            .query_row(
1013                &format!("{TICKET_RECORD_SELECT} WHERE name = ?1 ORDER BY id LIMIT 1"),
1014                params![name],
1015                ticket_record,
1016            )
1017            .optional()?;
1018        if let Some(ticket) = ticket.as_mut() {
1019            ticket.blocked_by = Self::ticket_blockers(&connection, &ticket.id)?;
1020        }
1021        Ok(ticket)
1022    }
1023
1024    pub fn ticket_by_file(&self, file_path: &str) -> Result<Option<TicketRecord>, StoreError> {
1025        let connection = self.db.lock();
1026        let mut ticket = connection
1027            .query_row(
1028                &format!("{TICKET_RECORD_SELECT} WHERE file_path = ?1"),
1029                params![file_path],
1030                ticket_record,
1031            )
1032            .optional()?;
1033        if let Some(ticket) = ticket.as_mut() {
1034            ticket.blocked_by = Self::ticket_blockers(&connection, &ticket.id)?;
1035        }
1036        Ok(ticket)
1037    }
1038
1039    pub fn ticket_by_source_ref(
1040        &self,
1041        source: &str,
1042        source_ref: &str,
1043    ) -> Result<Option<TicketRecord>, StoreError> {
1044        let connection = self.db.lock();
1045        let mut ticket = connection
1046            .query_row(
1047                &format!("{TICKET_RECORD_SELECT} WHERE source = ?1 AND source_ref = ?2"),
1048                params![source, source_ref],
1049                ticket_record,
1050            )
1051            .optional()?;
1052        if let Some(ticket) = ticket.as_mut() {
1053            ticket.blocked_by = Self::ticket_blockers(&connection, &ticket.id)?;
1054        }
1055        Ok(ticket)
1056    }
1057
1058    pub fn tickets(&self) -> Result<Vec<TicketRecord>, StoreError> {
1059        Self::tickets_on(&self.db.lock())
1060    }
1061
1062    fn tickets_on(connection: &Connection) -> Result<Vec<TicketRecord>, StoreError> {
1063        let mut statement = connection.prepare(&format!(
1064            "{TICKET_RECORD_SELECT} ORDER BY created_at_ms DESC, id DESC"
1065        ))?;
1066        let mut tickets = statement
1067            .query_map([], ticket_record)?
1068            .collect::<Result<Vec<_>, _>>()?;
1069        tickets.sort_by_key(|ticket| {
1070            (
1071                std::cmp::Reverse(ticket.created_at_ms),
1072                std::cmp::Reverse(crate::ids::ordinal(&ticket.id).unwrap_or(0)),
1073            )
1074        });
1075        let mut blockers = Self::all_ticket_blockers(connection)?;
1076        for ticket in &mut tickets {
1077            ticket.blocked_by = blockers.remove(&ticket.id).unwrap_or_default();
1078        }
1079        Ok(tickets)
1080    }
1081
1082    pub fn tickets_for_project(&self, project_id: &str) -> Result<Vec<TicketRecord>, StoreError> {
1083        let connection = self.db.lock();
1084        let mut statement = connection.prepare(&format!(
1085            "{TICKET_RECORD_SELECT} WHERE project_id = ?1 ORDER BY id"
1086        ))?;
1087        let mut tickets = statement
1088            .query_map(params![project_id], ticket_record)?
1089            .collect::<Result<Vec<_>, _>>()?;
1090        let mut blockers = Self::all_ticket_blockers(&connection)?;
1091        for ticket in &mut tickets {
1092            ticket.blocked_by = blockers.remove(&ticket.id).unwrap_or_default();
1093        }
1094        Ok(tickets)
1095    }
1096
1097    pub fn ticket_dependencies(&self) -> Result<BTreeMap<String, Vec<String>>, StoreError> {
1098        let mut dependencies = BTreeMap::new();
1099        let connection = self.db.lock();
1100        let mut statement = connection.prepare("SELECT id FROM tickets")?;
1101        let ids = statement.query_map([], |row| row.get::<_, String>(0))?;
1102        for id in ids {
1103            dependencies.insert(id?, Vec::new());
1104        }
1105        for (ticket_id, blockers) in Self::all_ticket_blockers(&connection)? {
1106            if let Some(entry) = dependencies.get_mut(&ticket_id) {
1107                *entry = blockers;
1108            }
1109        }
1110        Ok(dependencies)
1111    }
1112
1113    pub fn select_ready_ticket(
1114        &self,
1115        project_id: Option<&str>,
1116        trigger_id: &str,
1117        now_ms: i64,
1118    ) -> Result<Option<String>, StoreError> {
1119        self.db
1120            .lock()
1121            .query_row(
1122                &format!(
1123                    "SELECT t.id FROM tickets t
1124                     WHERE t.state = 'ready'
1125                       AND t.missing_at_ms IS NULL
1126                       AND (?1 IS NULL OR t.project_id = ?1)
1127                       AND NOT EXISTS (SELECT 1 FROM ticket_blockers b
1128                                       JOIN tickets bt ON bt.id = b.blocker_id
1129                                       WHERE b.ticket_id = t.id
1130                                         AND bt.state != 'merged')
1131                       AND {passes_filters}
1132                       AND NOT EXISTS (SELECT 1 FROM cooldowns c
1133                                       WHERE c.key = 'agent_target:' || t.target
1134                                         AND c.until_ms > ?3)
1135                      ORDER BY t.created_at_ms, t.id
1136                      LIMIT 1",
1137                    passes_filters = trigger::passes_filters("?2"),
1138                ),
1139                params![project_id, trigger_id, now_ms],
1140                |row| row.get(0),
1141            )
1142            .optional()
1143            .map_err(StoreError::from)
1144    }
1145
1146    pub fn ticket_is_dispatchable(&self, ticket_id: &str) -> Result<bool, StoreError> {
1147        self.db
1148            .lock()
1149            .query_row(
1150                "SELECT EXISTS(
1151                     SELECT 1 FROM tickets t
1152                     WHERE t.id = ?1
1153                       AND t.state = 'ready'
1154                       AND t.missing_at_ms IS NULL
1155                       AND NOT EXISTS (SELECT 1 FROM ticket_blockers b
1156                                       JOIN tickets bt ON bt.id = b.blocker_id
1157                                       WHERE b.ticket_id = t.id
1158                                         AND bt.state != 'merged')
1159                 )",
1160                params![ticket_id],
1161                |row| row.get(0),
1162            )
1163            .map_err(StoreError::from)
1164    }
1165
1166    pub fn unmerged_blockers(&self, ticket_id: &str) -> Result<Vec<String>, StoreError> {
1167        let connection = self.db.lock();
1168        let mut statement = connection.prepare(
1169            "SELECT b.blocker_id FROM ticket_blockers b
1170             JOIN tickets bt ON bt.id = b.blocker_id
1171             WHERE b.ticket_id = ?1 AND bt.state != 'merged'
1172             ORDER BY b.position, b.blocker_id",
1173        )?;
1174        statement
1175            .query_map(params![ticket_id], |row| row.get(0))?
1176            .collect::<Result<Vec<_>, _>>()
1177            .map_err(StoreError::from)
1178    }
1179
1180    pub fn readopt_lease(
1181        &self,
1182        ticket_id: &str,
1183        run_id: &str,
1184        lease_ms: i64,
1185        now_ms: i64,
1186    ) -> Result<i64, StoreError> {
1187        let expires_at_ms = now_ms + lease_ms;
1188        let changed = self.db.lock().execute(
1189            "UPDATE leases
1190             SET renewed_at_ms = ?3, expires_at_ms = ?4
1191             WHERE ticket_id = ?1 AND run_id = ?2
1192               AND EXISTS (SELECT 1 FROM runs
1193                           WHERE id = ?2 AND exited_at_ms IS NULL)",
1194            params![ticket_id, run_id, now_ms, expires_at_ms],
1195        )?;
1196        if changed != 1 {
1197            return Err(StoreError::LeaseNotHeld {
1198                ticket_id: ticket_id.into(),
1199                run_id: run_id.into(),
1200            });
1201        }
1202        Ok(expires_at_ms)
1203    }
1204
1205    pub fn renew_lease(
1206        &self,
1207        ticket_id: &str,
1208        run_id: &str,
1209        lease_ms: i64,
1210        now_ms: i64,
1211    ) -> Result<i64, StoreError> {
1212        let expires_at_ms = now_ms + lease_ms;
1213        let changed = self.db.lock().execute(
1214            "UPDATE leases
1215             SET renewed_at_ms = ?3, expires_at_ms = ?4
1216             WHERE ticket_id = ?1 AND run_id = ?2 AND expires_at_ms > ?3",
1217            params![ticket_id, run_id, now_ms, expires_at_ms],
1218        )?;
1219        if changed != 1 {
1220            return Err(StoreError::LeaseNotHeld {
1221                ticket_id: ticket_id.into(),
1222                run_id: run_id.into(),
1223            });
1224        }
1225        Ok(expires_at_ms)
1226    }
1227
1228    fn all_ticket_blockers(
1229        connection: &Connection,
1230    ) -> Result<BTreeMap<String, Vec<String>>, StoreError> {
1231        let mut statement = connection.prepare(
1232            "SELECT ticket_id, blocker_id FROM ticket_blockers
1233             ORDER BY ticket_id, position, blocker_id",
1234        )?;
1235        let rows = statement.query_map([], |row| {
1236            Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
1237        })?;
1238        let mut blockers: BTreeMap<String, Vec<String>> = BTreeMap::new();
1239        for row in rows {
1240            let (ticket_id, blocker_id) = row?;
1241            blockers.entry(ticket_id).or_default().push(blocker_id);
1242        }
1243        Ok(blockers)
1244    }
1245
1246    fn ticket_blockers(connection: &Connection, id: &str) -> Result<Vec<String>, StoreError> {
1247        let mut statement = connection.prepare(
1248            "SELECT blocker_id FROM ticket_blockers
1249             WHERE ticket_id = ?1 ORDER BY position, blocker_id",
1250        )?;
1251        statement
1252            .query_map(params![id], |row| row.get(0))?
1253            .collect::<Result<Vec<_>, _>>()
1254            .map_err(StoreError::from)
1255    }
1256
1257    pub fn ticket_ids(&self) -> Result<Vec<String>, StoreError> {
1258        let connection = self.db.lock();
1259        let mut statement = connection.prepare("SELECT id FROM tickets")?;
1260        let rows = statement.query_map([], |row| row.get(0))?;
1261        rows.collect::<Result<Vec<_>, _>>()
1262            .map_err(StoreError::from)
1263    }
1264
1265    pub fn local_ticket_files(&self) -> Result<Vec<LocalTicketFile>, StoreError> {
1266        let connection = self.db.lock();
1267        let mut statement = connection.prepare(
1268            "SELECT id, file_path, state, missing_at_ms FROM tickets
1269             WHERE source = 'local' AND file_path IS NOT NULL
1270             ORDER BY id",
1271        )?;
1272        let rows = statement.query_map([], |row| {
1273            Ok(LocalTicketFile {
1274                id: row.get(0)?,
1275                file_path: row.get(1)?,
1276                state: row.get(2)?,
1277                missing_at_ms: row.get(3)?,
1278            })
1279        })?;
1280        rows.collect::<Result<Vec<_>, _>>()
1281            .map_err(StoreError::from)
1282    }
1283
1284    pub(crate) fn ticket_has_work_references_on(
1285        connection: &Connection,
1286        id: &str,
1287    ) -> rusqlite::Result<bool> {
1288        connection.query_row(
1289            "SELECT EXISTS (SELECT 1 FROM leases WHERE ticket_id = ?1)
1290                 OR EXISTS (SELECT 1 FROM triggers WHERE ticket_id = ?1)
1291                 OR EXISTS (SELECT 1 FROM trigger_filters WHERE ticket_id = ?1)
1292                 OR EXISTS (SELECT 1 FROM ticket_blockers WHERE blocker_id = ?1)",
1293            params![id],
1294            |row| row.get(0),
1295        )
1296    }
1297
1298    pub(crate) fn ticket_has_work_references(&self, id: &str) -> Result<bool, StoreError> {
1299        Self::ticket_has_work_references_on(&self.db.lock(), id).map_err(StoreError::from)
1300    }
1301
1302    pub fn delete_ticket(&self, id: &str) -> Result<(), StoreError> {
1303        self.db
1304            .lock()
1305            .execute("DELETE FROM tickets WHERE id = ?1", params![id])?;
1306        Ok(())
1307    }
1308
1309    pub fn mark_ticket_missing(&self, id: &str, now_ms: i64) -> Result<(), StoreError> {
1310        self.db.lock().execute(
1311            "UPDATE tickets SET missing_at_ms = ?2, updated_at_ms = ?2
1312             WHERE id = ?1 AND missing_at_ms IS NULL",
1313            params![id, now_ms],
1314        )?;
1315        Ok(())
1316    }
1317
1318    pub fn clear_ticket_missing(&self, id: &str, now_ms: i64) -> Result<(), StoreError> {
1319        self.db.lock().execute(
1320            "UPDATE tickets SET missing_at_ms = NULL, updated_at_ms = ?2
1321             WHERE id = ?1 AND missing_at_ms IS NOT NULL",
1322            params![id, now_ms],
1323        )?;
1324        Ok(())
1325    }
1326
1327    pub fn ticket_state(&self, id: &str) -> Result<Option<String>, StoreError> {
1328        Self::ticket_state_on(&self.db.lock(), id)
1329    }
1330
1331    pub(crate) fn ticket_state_on(
1332        connection: &Connection,
1333        id: &str,
1334    ) -> Result<Option<String>, StoreError> {
1335        connection
1336            .query_row(
1337                "SELECT state FROM tickets WHERE id = ?1",
1338                params![id],
1339                |row| row.get(0),
1340            )
1341            .optional()
1342            .map_err(StoreError::from)
1343    }
1344
1345    pub fn set_ticket_hold(
1346        &self,
1347        id: &str,
1348        state: TicketState,
1349        now_ms: i64,
1350    ) -> Result<String, StoreError> {
1351        debug_assert!(matches!(state, TicketState::Ready | TicketState::Held));
1352        let connection = self.db.lock();
1353        let requested = state.as_str();
1354        let previous =
1355            Self::ticket_state_on(&connection, id)?.ok_or_else(|| StoreError::TicketNotFound {
1356                ticket_id: id.into(),
1357            })?;
1358        if previous == requested {
1359            return Ok(previous);
1360        }
1361        let allowed_previous = match state {
1362            TicketState::Ready => TicketState::Held.as_str(),
1363            TicketState::Held => TicketState::Ready.as_str(),
1364            _ => unreachable!("hold transitions only use ready and held"),
1365        };
1366        let changed = connection.execute(
1367            "UPDATE tickets SET state = ?2, held_reason = NULL, updated_at_ms = ?3
1368             WHERE id = ?1 AND state = ?4",
1369            params![id, requested, now_ms, allowed_previous],
1370        )?;
1371        if changed != 1 {
1372            return Err(StoreError::TicketStateConflict {
1373                ticket_id: id.into(),
1374                state: previous,
1375                requested: requested.into(),
1376            });
1377        }
1378        Ok(previous)
1379    }
1380
1381    pub(crate) fn retry_ticket<Cleanup>(
1382        &self,
1383        id: &str,
1384        now_ms: i64,
1385        cleanup: Cleanup,
1386    ) -> Result<String, StoreError>
1387    where
1388        Cleanup: FnOnce(&Transaction<'_>, &str, i64) -> Result<(), StoreError>,
1389    {
1390        let mut connection = self.db.lock();
1391        let previous =
1392            Self::ticket_state_on(&connection, id)?.ok_or_else(|| StoreError::TicketNotFound {
1393                ticket_id: id.into(),
1394            })?;
1395        let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
1396        let changed = tx::retry_ticket(&transaction, id, now_ms)?;
1397        if changed != 1 {
1398            return Err(StoreError::TicketStateConflict {
1399                ticket_id: id.into(),
1400                state: previous,
1401                requested: TicketState::Ready.as_str().into(),
1402            });
1403        }
1404        cleanup(&transaction, id, now_ms)?;
1405        transaction.commit()?;
1406        Ok(previous)
1407    }
1408
1409    pub fn ticket_counts(&self) -> Result<TicketCounts, StoreError> {
1410        let connection = self.db.lock();
1411        let mut statement = connection.prepare(
1412            "SELECT CASE
1413                      WHEN t.state = 'ready' AND EXISTS (
1414                          SELECT 1 FROM ticket_blockers b
1415                          JOIN tickets bt ON bt.id = b.blocker_id
1416                          WHERE b.ticket_id = t.id AND bt.state != 'merged'
1417                      ) THEN 'blocked'
1418                      ELSE t.state
1419                    END AS display_state,
1420                    COUNT(*)
1421             FROM tickets t
1422             GROUP BY display_state",
1423        )?;
1424        let mut rows = statement.query([])?;
1425        let mut counts = TicketCounts::default();
1426        while let Some(row) = rows.next()? {
1427            let state: String = row.get(0)?;
1428            let count = row.get::<_, i64>(1)?.max(0) as u64;
1429            match state.as_str() {
1430                "ready" => counts.ready = count,
1431                "held" => counts.held = count,
1432                "blocked" => counts.blocked = count,
1433                "claimed" => counts.claimed = count,
1434                "merged" => counts.merged = count,
1435                "failed" => counts.failed = count,
1436                "needs_review" => counts.needs_review = count,
1437                _ => {}
1438            }
1439        }
1440        Ok(counts)
1441    }
1442
1443    fn ticket_on(connection: &Connection, id: &str) -> Result<Option<TicketRecord>, StoreError> {
1444        let mut ticket = connection
1445            .query_row(
1446                &format!("{TICKET_RECORD_SELECT} WHERE id = ?1"),
1447                params![id],
1448                ticket_record,
1449            )
1450            .optional()?;
1451        if let Some(ticket) = ticket.as_mut() {
1452            ticket.blocked_by = Self::ticket_blockers(connection, &ticket.id)?;
1453        }
1454        Ok(ticket)
1455    }
1456}
1457
1458fn source_error(error: StoreError) -> SourceError {
1459    match error {
1460        StoreError::TicketNotFound { .. }
1461        | StoreError::TicketStateConflict { .. }
1462        | StoreError::TriggerNotQueued { .. }
1463        | StoreError::LeaseNotHeld { .. }
1464        | StoreError::TicketNotReady { .. } => SourceError::Rejected {
1465            message: error.to_string(),
1466        },
1467        _ => SourceError::Unavailable { retry_after: None },
1468    }
1469}
1470
1471fn lease_owner(owner: &OwnerId, trigger_id: &str) -> String {
1472    json!({"owner": owner.0, "trigger": trigger_id}).to_string()
1473}
1474
1475fn decode_lease_owner(stored: &str) -> (OwnerId, Option<String>) {
1476    let parsed = serde_json::from_str::<Value>(stored).ok();
1477    let owner = parsed
1478        .as_ref()
1479        .and_then(|value| value["owner"].as_str())
1480        .unwrap_or(stored);
1481    let trigger = parsed
1482        .as_ref()
1483        .and_then(|value| value["trigger"].as_str())
1484        .map(str::to_owned);
1485    (OwnerId(owner.into()), trigger)
1486}
1487
1488fn lease_ms(ttl: Duration) -> Result<i64, SourceError> {
1489    i64::try_from(ttl.as_millis()).map_err(|_| SourceError::Rejected {
1490        message: "lease duration is outside the supported range".into(),
1491    })
1492}
1493
1494fn work_ticket(
1495    record: TicketRecord,
1496    blocked: bool,
1497    owner: OwnerId,
1498) -> Result<WorkTicket, SourceError> {
1499    let state = match record.state.as_str() {
1500        "ready" => TicketState::Ready,
1501        "held" => TicketState::Held,
1502        "claimed" => TicketState::Claimed,
1503        "merged" => TicketState::Merged,
1504        "failed" => TicketState::Failed,
1505        "needs_review" => TicketState::NeedsReview,
1506        state => {
1507            return Err(SourceError::Corrupt {
1508                message: format!("ticket `{}` has unknown state `{state}`", record.id),
1509            });
1510        }
1511    };
1512    let attempts = u32::try_from(record.attempts).map_err(|_| SourceError::Corrupt {
1513        message: format!(
1514            "ticket `{}` has invalid attempt count {}",
1515            record.id, record.attempts
1516        ),
1517    })?;
1518
1519    Ok(WorkTicket {
1520        id: record.id,
1521        project_id: record.project_id,
1522        name: record.name,
1523        body: record.body.unwrap_or_default(),
1524        state: WorkTicketState::from_ticket_state(
1525            state,
1526            blocked,
1527            record.held_reason.unwrap_or_default(),
1528            owner,
1529        ),
1530        blocked_by: record.blocked_by,
1531        attempts,
1532        hints: ExecutionHints {
1533            worktree: record.worktree,
1534            trigger_id: None,
1535            target: record.target,
1536            model: record.model,
1537            effort: record.effort,
1538            flow: record.flow,
1539        },
1540        version: SourceVersion(record.updated_at_ms.to_string()),
1541    })
1542}
1543
1544/// SQLite satisfies atomic claims with an IMMEDIATE transaction and a
1545/// conditional ticket update. The ticket row is authoritative for attempts:
1546/// claim consumes one attempt, while release preserves that count as evidence
1547/// of the completed try. This backend never records runs, and its local-source
1548/// outcome push is idempotent because release already applied the durable
1549/// ticket state.
1550#[async_trait]
1551impl WorkState for LocalSqlite {
1552    fn claim_strength(&self) -> ClaimStrength {
1553        ClaimStrength::Atomic
1554    }
1555
1556    async fn pull_ready(&self) -> Result<Vec<WorkTicket>, SourceError> {
1557        let now_ms = self.clock.now_ms();
1558        let triggers = self.dispatchable_triggers(now_ms).map_err(source_error)?;
1559        let mut seen = BTreeSet::new();
1560        let mut selected = Vec::new();
1561        for trigger in triggers {
1562            let ticket_id = match trigger.ticket_id {
1563                Some(ticket_id) => self
1564                    .ticket_is_dispatchable(&ticket_id)
1565                    .map_err(source_error)?
1566                    .then_some(ticket_id),
1567                None => self
1568                    .select_ready_ticket(trigger.project_id.as_deref(), &trigger.id, now_ms)
1569                    .map_err(source_error)?,
1570            };
1571            if let Some(ticket_id) = ticket_id
1572                && seen.insert(ticket_id.clone())
1573            {
1574                selected.push(ticket_id);
1575            }
1576        }
1577
1578        selected
1579            .into_iter()
1580            .map(|id| {
1581                let record = self.ticket(&id).map_err(source_error)?.ok_or_else(|| {
1582                    SourceError::Corrupt {
1583                        message: format!("selected ticket `{id}` no longer exists"),
1584                    }
1585                })?;
1586                work_ticket(record, false, OwnerId(String::new()))
1587            })
1588            .collect()
1589    }
1590
1591    async fn active_claims(&self) -> Result<Vec<ActiveClaim>, SourceError> {
1592        let connection = self.db.lock();
1593        let mut statement = connection
1594            .prepare(
1595                "SELECT t.id, t.source, t.source_ref, l.run_id
1596                 FROM leases l
1597                 JOIN tickets t ON t.id = l.ticket_id
1598                 ORDER BY l.acquired_at_ms, l.ticket_id",
1599            )
1600            .map_err(StoreError::from)
1601            .map_err(source_error)?;
1602        statement
1603            .query_map([], |row| {
1604                Ok(ActiveClaim {
1605                    ticket: TicketRef {
1606                        id: row.get(0)?,
1607                        source: row.get(1)?,
1608                        source_ref: row.get(2)?,
1609                    },
1610                    owner: OwnerId(row.get(3)?),
1611                })
1612            })
1613            .map_err(StoreError::from)
1614            .map_err(source_error)?
1615            .collect::<Result<Vec<_>, _>>()
1616            .map_err(StoreError::from)
1617            .map_err(source_error)
1618    }
1619
1620    async fn claim(
1621        &self,
1622        ticket: &TicketRef,
1623        owner: &OwnerId,
1624        ttl: Duration,
1625    ) -> Result<ClaimResult, SourceError> {
1626        let now_ms = self.clock.now_ms();
1627        let lease_ms = lease_ms(ttl)?;
1628        let mut connection = self.db.lock();
1629        connection
1630            .pragma_update(None, "foreign_keys", false)
1631            .map_err(StoreError::from)
1632            .map_err(source_error)?;
1633        let result = (|| {
1634            let transaction = connection
1635                .transaction_with_behavior(TransactionBehavior::Immediate)
1636                .map_err(StoreError::from)
1637                .map_err(source_error)?;
1638            let Some(trigger) =
1639                trigger::claimable_on(&transaction, &ticket.id, now_ms).map_err(source_error)?
1640            else {
1641                return Ok(ClaimResult::Lost { held_by: None });
1642            };
1643            let effects = domain_trigger::step(
1644                &mut domain_trigger::Trigger::from(&trigger),
1645                domain_trigger::Event::Fired,
1646                now_ms,
1647            );
1648            if effects
1649                .iter()
1650                .any(|effect| matches!(effect, domain_trigger::Effect::Fault(_)))
1651            {
1652                return Err(SourceError::Corrupt {
1653                    message: format!("recurring trigger `{}` has an invalid cadence", trigger.id),
1654                });
1655            }
1656            let stored_owner = lease_owner(owner, &trigger.id);
1657            let claim = ClaimTransaction {
1658                ticket_id: &ticket.id,
1659                run_id: &owner.0,
1660                trigger_id: &trigger.id,
1661                owner_id: &stored_owner,
1662                lease_ms,
1663            };
1664
1665            match tx::claim_ticket(&transaction, &claim, now_ms) {
1666                Ok(()) => {}
1667                Err(StoreError::TicketNotReady { .. }) => {
1668                    let held_by = transaction
1669                        .query_row(
1670                            "SELECT owner_id FROM leases WHERE ticket_id = ?1",
1671                            params![ticket.id],
1672                            |row| row.get::<_, String>(0),
1673                        )
1674                        .optional()
1675                        .map_err(StoreError::from)
1676                        .map_err(source_error)?
1677                        .map(|stored| decode_lease_owner(&stored).0);
1678                    return Ok(ClaimResult::Lost { held_by });
1679                }
1680                Err(error) => return Err(source_error(error)),
1681            }
1682            trigger::consume(&transaction, &trigger.id, &effects, now_ms).map_err(source_error)?;
1683            tx::insert_lease(&transaction, &claim, now_ms).map_err(source_error)?;
1684            let record = Self::ticket_on(&transaction, &ticket.id)
1685                .map_err(source_error)?
1686                .ok_or_else(|| SourceError::Corrupt {
1687                    message: format!("claimed ticket `{}` no longer exists", ticket.id),
1688                })?;
1689            let mut ticket = work_ticket(record, false, owner.clone())?;
1690            ticket.hints.trigger_id = Some(trigger.id.clone());
1691            transaction
1692                .commit()
1693                .map_err(StoreError::from)
1694                .map_err(source_error)?;
1695            Ok(ClaimResult::Claimed { ticket })
1696        })();
1697        let restored = connection
1698            .pragma_update(None, "foreign_keys", true)
1699            .map_err(StoreError::from)
1700            .map_err(source_error);
1701        match (result, restored) {
1702            (_, Err(error)) => Err(error),
1703            (result, Ok(())) => result,
1704        }
1705    }
1706
1707    async fn renew(&self, ticket: &TicketRef, owner: &OwnerId) -> Result<ClaimResult, SourceError> {
1708        let now_ms = self.clock.now_ms();
1709        let lease_ms = self
1710            .db
1711            .lock()
1712            .query_row(
1713                "SELECT expires_at_ms - renewed_at_ms FROM leases
1714                 WHERE ticket_id = ?1 AND run_id = ?2",
1715                params![ticket.id, owner.0],
1716                |row| row.get::<_, i64>(0),
1717            )
1718            .optional()
1719            .map_err(StoreError::from)
1720            .map_err(source_error)?;
1721        let Some(lease_ms) = lease_ms else {
1722            return Ok(ClaimResult::Lost { held_by: None });
1723        };
1724        match self.renew_lease(&ticket.id, &owner.0, lease_ms, now_ms) {
1725            Ok(_) => {}
1726            Err(StoreError::LeaseNotHeld { .. }) => {
1727                return Ok(ClaimResult::Lost { held_by: None });
1728            }
1729            Err(error) => return Err(source_error(error)),
1730        }
1731        let record = self
1732            .ticket(&ticket.id)
1733            .map_err(source_error)?
1734            .ok_or_else(|| SourceError::Corrupt {
1735                message: format!("leased ticket `{}` no longer exists", ticket.id),
1736            })?;
1737        Ok(ClaimResult::Claimed {
1738            ticket: work_ticket(record, false, owner.clone())?,
1739        })
1740    }
1741
1742    async fn release(
1743        &self,
1744        ticket: &TicketRef,
1745        owner: &OwnerId,
1746        disposition: Disposition,
1747    ) -> Result<(), SourceError> {
1748        let now_ms = self.clock.now_ms();
1749        let mut connection = self.db.lock();
1750        let transaction = connection
1751            .transaction_with_behavior(TransactionBehavior::Immediate)
1752            .map_err(StoreError::from)
1753            .map_err(source_error)?;
1754        let lease = transaction
1755            .query_row(
1756                "SELECT run_id, owner_id FROM leases WHERE ticket_id = ?1",
1757                params![ticket.id],
1758                |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
1759            )
1760            .optional()
1761            .map_err(StoreError::from)
1762            .map_err(source_error)?;
1763        let claimed_trigger_id = match lease {
1764            Some((run_id, stored_owner)) if run_id == owner.0 => {
1765                decode_lease_owner(&stored_owner).1
1766            }
1767            Some(_) => {
1768                transaction
1769                    .commit()
1770                    .map_err(StoreError::from)
1771                    .map_err(source_error)?;
1772                return Ok(());
1773            }
1774            None => {
1775                let (state, held_reason) = transaction
1776                    .query_row(
1777                        "SELECT state, held_reason FROM tickets WHERE id = ?1",
1778                        params![ticket.id],
1779                        |row| Ok((row.get::<_, String>(0)?, row.get::<_, Option<String>>(1)?)),
1780                    )
1781                    .optional()
1782                    .map_err(StoreError::from)
1783                    .map_err(source_error)?
1784                    .ok_or_else(|| SourceError::Rejected {
1785                        message: format!("ticket `{}` does not exist", ticket.id),
1786                    })?;
1787                let converged = match &disposition {
1788                    Disposition::Complete => state == TicketState::Merged.as_str(),
1789                    Disposition::Retry { .. } => state == TicketState::Ready.as_str(),
1790                    Disposition::Park { reason } if reason == "needs-review" => {
1791                        state == TicketState::NeedsReview.as_str()
1792                    }
1793                    Disposition::Park { reason } => {
1794                        state == TicketState::Held.as_str()
1795                            && held_reason.as_deref() == Some(reason.as_str())
1796                    }
1797                    Disposition::Abandon => state == TicketState::Failed.as_str(),
1798                };
1799                if converged {
1800                    transaction
1801                        .commit()
1802                        .map_err(StoreError::from)
1803                        .map_err(source_error)?;
1804                    return Ok(());
1805                }
1806                if matches!(disposition, Disposition::Complete)
1807                    && state == TicketState::NeedsReview.as_str()
1808                {
1809                    let changed = tx::settle_external_merge(&transaction, &ticket.id, now_ms)
1810                        .map_err(StoreError::from)
1811                        .map_err(source_error)?;
1812                    if changed == 1 {
1813                        trigger::complete_for_ticket(&transaction, &ticket.id, now_ms)
1814                            .map_err(source_error)?;
1815                    }
1816                    transaction
1817                        .commit()
1818                        .map_err(StoreError::from)
1819                        .map_err(source_error)?;
1820                    return Ok(());
1821                }
1822                if state != TicketState::Claimed.as_str() {
1823                    transaction
1824                        .commit()
1825                        .map_err(StoreError::from)
1826                        .map_err(source_error)?;
1827                    return Ok(());
1828                }
1829                return Err(SourceError::Rejected {
1830                    message: format!("ticket `{}` is no longer claimed", ticket.id),
1831                });
1832            }
1833        };
1834
1835        let changed = match disposition {
1836            Disposition::Complete => {
1837                let changed =
1838                    tx::settle_ticket(&transaction, &ticket.id, TicketState::Merged, now_ms)
1839                        .map_err(StoreError::from)
1840                        .map_err(source_error)?;
1841                if changed == 1 {
1842                    trigger::complete_for_ticket(&transaction, &ticket.id, now_ms)
1843                        .map_err(source_error)?;
1844                }
1845                changed
1846            }
1847            Disposition::Retry { not_before_ms } => {
1848                let changed = tx::abort_ticket(&transaction, &ticket.id, now_ms)
1849                    .map_err(StoreError::from)
1850                    .map_err(source_error)?;
1851                if let Some(eligible_at_ms) = not_before_ms {
1852                    let trigger_id = match claimed_trigger_id {
1853                        Some(trigger_id) => trigger_id,
1854                        None => trigger::for_release(&transaction, &ticket.id)
1855                            .map_err(source_error)?
1856                            .ok_or_else(|| SourceError::Corrupt {
1857                                message: format!("ticket `{}` has no trigger to retry", ticket.id),
1858                            })?,
1859                    };
1860                    let mut facts = trigger::facts(&transaction, &trigger_id)
1861                        .map_err(source_error)?
1862                        .ok_or_else(|| SourceError::Corrupt {
1863                            message: format!("trigger `{trigger_id}` no longer exists"),
1864                        })?;
1865                    let effects = domain_trigger::step(
1866                        &mut facts,
1867                        domain_trigger::Event::Requeued { eligible_at_ms },
1868                        now_ms,
1869                    );
1870                    for effect in effects {
1871                        if let domain_trigger::Effect::Requeue { eligible_at_ms } = effect {
1872                            trigger::requeue(&transaction, &trigger_id, eligible_at_ms, now_ms)
1873                                .map_err(source_error)?;
1874                        }
1875                    }
1876                }
1877                changed
1878            }
1879            Disposition::Park { reason } => {
1880                let ticket_state = if reason == "needs-review" {
1881                    TicketState::NeedsReview
1882                } else {
1883                    TicketState::Held
1884                };
1885                let changed = tx::settle_ticket(&transaction, &ticket.id, ticket_state, now_ms)
1886                    .map_err(StoreError::from)
1887                    .map_err(source_error)?;
1888                if changed == 1 && ticket_state == TicketState::Held {
1889                    transaction
1890                        .execute(
1891                            "UPDATE tickets SET held_reason = ?2 WHERE id = ?1 AND state = 'held'",
1892                            params![ticket.id, reason],
1893                        )
1894                        .map_err(StoreError::from)
1895                        .map_err(source_error)?;
1896                }
1897                changed
1898            }
1899            Disposition::Abandon => {
1900                tx::settle_ticket(&transaction, &ticket.id, TicketState::Failed, now_ms)
1901                    .map_err(StoreError::from)
1902                    .map_err(source_error)?
1903            }
1904        };
1905        if changed != 1 {
1906            return Err(SourceError::Rejected {
1907                message: format!("ticket `{}` is no longer claimed", ticket.id),
1908            });
1909        }
1910        tx::delete_lease(&transaction, &owner.0)
1911            .map_err(StoreError::from)
1912            .map_err(source_error)?;
1913        transaction
1914            .commit()
1915            .map_err(StoreError::from)
1916            .map_err(source_error)?;
1917        Ok(())
1918    }
1919
1920    async fn push_outcome(&self, outcome: &WorkOutcome) -> Result<(), SourceError> {
1921        let Some(reporter) = self.outcome_reporter.clone() else {
1922            return Ok(());
1923        };
1924        let ticket_id = outcome.ticket_id.clone();
1925        let verdict = outcome.verdict;
1926        reporter
1927            .report(&ticket_id, &verdict)
1928            .map_err(|error| SourceError::Rejected {
1929                message: error.to_string(),
1930            })
1931    }
1932}
1933
1934#[cfg(test)]
1935mod tests {
1936    use std::collections::BTreeMap;
1937    use std::process::Command;
1938    use std::sync::Arc;
1939    use std::time::Duration;
1940
1941    use tempfile::tempdir;
1942
1943    use crate::db::{Db, REVERT_TRIGGER_RENAME, StoreError};
1944    use crate::domain::ticket::TicketState;
1945    use crate::domain::work::{Disposition, OwnerId, TicketRef, WorkOutcome, WorkTicket};
1946    use crate::outcome::Outcome;
1947    use crate::work_state::{ClaimResult, TicketFeeder, WorkState};
1948
1949    use crate::domain::trigger::{Effect, TriggerKind};
1950    use crate::work_state::trigger::{self, NewTrigger, QueuedTrigger};
1951
1952    use super::{ClaimTransaction, LocalSqlite, ReindexTicket};
1953
1954    fn open_seeded(path: &std::path::Path) -> LocalSqlite {
1955        let store = LocalSqlite::from_db(Db::open(path, 1_000).unwrap());
1956        store
1957            .insert_local_project(
1958                "default",
1959                ".agents/sloop/projects/default.md",
1960                "Default",
1961                1_000,
1962            )
1963            .unwrap();
1964        store
1965            .insert_local_ticket(
1966                "T1",
1967                "default",
1968                ".agents/sloop/tickets/t1.md",
1969                "Ticket one",
1970                &[],
1971                "sloop/T1",
1972                Some("claude"),
1973                Some("sonnet"),
1974                Some("medium"),
1975                "default",
1976                TicketState::Ready,
1977                1_000,
1978            )
1979            .unwrap();
1980        store
1981            .insert_trigger(
1982                &NewTrigger {
1983                    id: "TR1",
1984                    kind: TriggerKind::Immediate,
1985                    ticket_id: Some("T1"),
1986                    project_id: None,
1987                    eligible_at_ms: None,
1988                    interval_ms: None,
1989                },
1990                1_000,
1991            )
1992            .unwrap();
1993        store
1994    }
1995
1996    fn insert_claimed_run(store: &LocalSqlite, claim: &ClaimTransaction<'_>, now_ms: i64) -> i64 {
1997        let connection = store.db.lock();
1998        let attempt = connection
1999            .query_row(
2000                "SELECT COALESCE(MAX(attempt), 0) + 1 FROM runs WHERE ticket_id = ?1",
2001                rusqlite::params![claim.ticket_id],
2002                |row| row.get(0),
2003            )
2004            .unwrap();
2005        connection
2006            .execute(
2007                "INSERT INTO runs
2008                     (id, trigger_id, ticket_id, state, attempt, flow_json, ticket_json,
2009                      created_at_ms, updated_at_ms)
2010                 VALUES (?1, ?2, ?3, 'claimed', ?4, '{}', '{}', ?5, ?5)",
2011                rusqlite::params![
2012                    claim.run_id,
2013                    claim.trigger_id,
2014                    claim.ticket_id,
2015                    attempt,
2016                    now_ms
2017                ],
2018            )
2019            .unwrap();
2020        attempt
2021    }
2022
2023    fn granted_claim(store: &LocalSqlite, claim: &ClaimTransaction<'_>, now_ms: i64) -> i64 {
2024        let mut connection = store.db.lock();
2025        connection
2026            .pragma_update(None, "foreign_keys", false)
2027            .unwrap();
2028        {
2029            let transaction = connection
2030                .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
2031                .unwrap();
2032            super::tx::claim_ticket(&transaction, claim, now_ms).unwrap();
2033            trigger::consume(&transaction, claim.trigger_id, &[Effect::Complete], now_ms).unwrap();
2034            super::tx::insert_lease(&transaction, claim, now_ms).unwrap();
2035            transaction.commit().unwrap();
2036        }
2037        connection
2038            .pragma_update(None, "foreign_keys", true)
2039            .unwrap();
2040        drop(connection);
2041        insert_claimed_run(store, claim, now_ms)
2042    }
2043
2044    fn settle_for_test(store: &LocalSqlite, run_id: &str, outcome: Outcome, now_ms: i64) {
2045        let state = outcome.as_str();
2046        let ticket_state = match outcome {
2047            Outcome::Merged => "merged",
2048            Outcome::Failed => "failed",
2049            Outcome::NeedsReview => "needs_review",
2050            Outcome::Cancelled | Outcome::RateLimited | Outcome::Orphaned => "ready",
2051        };
2052        let mut connection = store.db.lock();
2053        let transaction = connection.transaction().unwrap();
2054        transaction
2055            .execute(
2056                "UPDATE runs SET state = ?2, exited_at_ms = ?3, updated_at_ms = ?3 WHERE id = ?1",
2057                rusqlite::params![run_id, state, now_ms],
2058            )
2059            .unwrap();
2060        transaction
2061            .execute(
2062                "UPDATE tickets SET state = ?2, updated_at_ms = ?3
2063                 WHERE id = (SELECT ticket_id FROM runs WHERE id = ?1)",
2064                rusqlite::params![run_id, ticket_state, now_ms],
2065            )
2066            .unwrap();
2067        transaction
2068            .execute("DELETE FROM leases WHERE run_id = ?1", [run_id])
2069            .unwrap();
2070        transaction.commit().unwrap();
2071    }
2072
2073    fn select_ready_ticket(
2074        store: &LocalSqlite,
2075        trigger: &QueuedTrigger,
2076        now_ms: i64,
2077    ) -> Option<String> {
2078        store
2079            .select_ready_ticket(trigger.project_id.as_deref(), &trigger.id, now_ms)
2080            .unwrap()
2081    }
2082
2083    fn apply_reindex(store: &LocalSqlite, tickets: &[ReindexTicket], now_ms: i64) {
2084        store
2085            .apply_reindex(
2086                &["default".into()],
2087                tickets,
2088                now_ms,
2089                |_, _, _| Ok(0),
2090                |_, _, _| Ok(()),
2091            )
2092            .unwrap();
2093    }
2094
2095    fn open_seeded_local() -> (tempfile::TempDir, LocalSqlite) {
2096        let directory = tempdir().unwrap();
2097        let local =
2098            LocalSqlite::from_db(Db::open(&directory.path().join("sloop.db"), 1_000).unwrap());
2099        local
2100            .insert_local_project(
2101                "default",
2102                ".agents/sloop/projects/default.md",
2103                "Default",
2104                1_000,
2105            )
2106            .unwrap();
2107        local
2108            .insert_local_ticket(
2109                "T1",
2110                "default",
2111                ".agents/sloop/tickets/t1.md",
2112                "Ticket one",
2113                &[],
2114                "sloop/T1",
2115                Some("claude"),
2116                Some("sonnet"),
2117                Some("medium"),
2118                "default",
2119                TicketState::Ready,
2120                1_000,
2121            )
2122            .unwrap();
2123        local
2124            .insert_trigger(
2125                &NewTrigger {
2126                    id: "TR1",
2127                    kind: TriggerKind::Immediate,
2128                    ticket_id: Some("T1"),
2129                    project_id: None,
2130                    eligible_at_ms: None,
2131                    interval_ms: None,
2132                },
2133                1_000,
2134            )
2135            .unwrap();
2136        (directory, local)
2137    }
2138
2139    fn ticket_ref() -> TicketRef {
2140        TicketRef {
2141            id: "T1".into(),
2142            source: "local".into(),
2143            source_ref: None,
2144        }
2145    }
2146
2147    async fn claim_local(local: &LocalSqlite, run_id: &str) -> ClaimResult {
2148        local
2149            .claim(
2150                &ticket_ref(),
2151                &OwnerId(run_id.into()),
2152                Duration::from_secs(60),
2153            )
2154            .await
2155            .unwrap()
2156    }
2157
2158    #[tokio::test]
2159    async fn local_work_state_claim_is_atomic() {
2160        let (_directory, local) = open_seeded_local();
2161
2162        assert!(matches!(
2163            claim_local(&local, "R1").await,
2164            ClaimResult::Claimed {
2165                ticket: WorkTicket { attempts: 1, .. }
2166            }
2167        ));
2168        {
2169            let connection = local.db.lock();
2170            assert_eq!(
2171                connection
2172                    .query_row("PRAGMA foreign_keys", [], |row| row.get::<_, i64>(0))
2173                    .unwrap(),
2174                1
2175            );
2176            assert_eq!(
2177                connection
2178                    .query_row("SELECT COUNT(*) FROM runs", [], |row| row.get::<_, i64>(0))
2179                    .unwrap(),
2180                0
2181            );
2182            assert_eq!(
2183                connection
2184                    .prepare("PRAGMA foreign_key_check")
2185                    .unwrap()
2186                    .query_map([], |_| Ok(()))
2187                    .unwrap()
2188                    .count(),
2189                1
2190            );
2191        }
2192        assert_eq!(local.active_claims().await.unwrap().len(), 1);
2193        insert_claimed_run(
2194            &local,
2195            &ClaimTransaction {
2196                ticket_id: "T1",
2197                run_id: "R1",
2198                trigger_id: "TR1",
2199                owner_id: "R1",
2200                lease_ms: 60_000,
2201            },
2202            2_000,
2203        );
2204        assert_eq!(
2205            local
2206                .db
2207                .lock()
2208                .prepare("PRAGMA foreign_key_check")
2209                .unwrap()
2210                .query_map([], |_| Ok(()))
2211                .unwrap()
2212                .count(),
2213            0
2214        );
2215        assert_eq!(
2216            claim_local(&local, "R1").await,
2217            ClaimResult::Lost { held_by: None }
2218        );
2219    }
2220
2221    #[tokio::test]
2222    async fn local_work_state_retry_preserves_attempts_until_the_next_claim() {
2223        let (_directory, local) = open_seeded_local();
2224        claim_local(&local, "R1").await;
2225
2226        local
2227            .release(
2228                &ticket_ref(),
2229                &OwnerId("R1".into()),
2230                Disposition::Retry {
2231                    not_before_ms: Some(2_000),
2232                },
2233            )
2234            .await
2235            .unwrap();
2236        let retried = local.ticket("T1").unwrap().unwrap();
2237        assert_eq!(retried.state, "ready");
2238        assert_eq!(retried.attempts, 1);
2239        assert_eq!(
2240            local
2241                .queued_triggers()
2242                .unwrap()
2243                .first()
2244                .and_then(|trigger| trigger.eligible_at_ms),
2245            Some(2_000)
2246        );
2247
2248        assert!(matches!(
2249            claim_local(&local, "R2").await,
2250            ClaimResult::Claimed {
2251                ticket: WorkTicket { attempts: 2, .. }
2252            }
2253        ));
2254    }
2255
2256    #[tokio::test]
2257    async fn local_work_state_park_and_abandon_apply_the_disposition() {
2258        let (_park_directory, parked) = open_seeded_local();
2259        claim_local(&parked, "R1").await;
2260        parked
2261            .release(
2262                &ticket_ref(),
2263                &OwnerId("R1".into()),
2264                Disposition::Park {
2265                    reason: "operator review".into(),
2266                },
2267            )
2268            .await
2269            .unwrap();
2270        let ticket = parked.ticket("T1").unwrap().unwrap();
2271        assert_eq!(ticket.state, "held");
2272        assert_eq!(ticket.held_reason.as_deref(), Some("operator review"));
2273
2274        let (_abandon_directory, abandoned) = open_seeded_local();
2275        claim_local(&abandoned, "R1").await;
2276        abandoned
2277            .release(&ticket_ref(), &OwnerId("R1".into()), Disposition::Abandon)
2278            .await
2279            .unwrap();
2280        assert_eq!(abandoned.ticket("T1").unwrap().unwrap().state, "failed");
2281    }
2282
2283    #[tokio::test]
2284    async fn local_work_state_release_is_idempotent_and_preserves_outcome_states() {
2285        let (_complete_directory, completed) = open_seeded_local();
2286        claim_local(&completed, "R1").await;
2287        completed
2288            .release(&ticket_ref(), &OwnerId("R1".into()), Disposition::Complete)
2289            .await
2290            .unwrap();
2291        completed
2292            .release(&ticket_ref(), &OwnerId("R1".into()), Disposition::Complete)
2293            .await
2294            .unwrap();
2295        assert_eq!(completed.ticket("T1").unwrap().unwrap().state, "merged");
2296
2297        let (_review_directory, review) = open_seeded_local();
2298        claim_local(&review, "R1").await;
2299        review
2300            .release(
2301                &ticket_ref(),
2302                &OwnerId("R1".into()),
2303                Disposition::Park {
2304                    reason: "needs-review".into(),
2305                },
2306            )
2307            .await
2308            .unwrap();
2309        assert_eq!(review.ticket("T1").unwrap().unwrap().state, "needs_review");
2310    }
2311
2312    fn insert_ready_ticket(local: &LocalSqlite, id: &str, now_ms: i64) {
2313        local
2314            .insert_local_ticket(
2315                id,
2316                "default",
2317                &format!(".agents/sloop/tickets/{}.md", id.to_lowercase()),
2318                "Another ticket",
2319                &[],
2320                &format!("sloop/{id}"),
2321                Some("claude"),
2322                Some("sonnet"),
2323                Some("medium"),
2324                "default",
2325                TicketState::Ready,
2326                now_ms,
2327            )
2328            .unwrap();
2329    }
2330
2331    fn insert_queued_trigger(
2332        local: &LocalSqlite,
2333        id: &str,
2334        kind: TriggerKind,
2335        ticket_id: Option<&str>,
2336        project_id: Option<&str>,
2337    ) {
2338        local
2339            .insert_trigger(
2340                &NewTrigger {
2341                    id,
2342                    kind,
2343                    ticket_id,
2344                    project_id,
2345                    eligible_at_ms: None,
2346                    interval_ms: None,
2347                },
2348                1_000,
2349            )
2350            .unwrap();
2351    }
2352
2353    fn queued_trigger_ids(local: &LocalSqlite) -> Vec<String> {
2354        local
2355            .queued_triggers()
2356            .unwrap()
2357            .into_iter()
2358            .map(|trigger| trigger.id)
2359            .collect()
2360    }
2361
2362    #[tokio::test]
2363    async fn local_work_state_complete_retires_triggers_pinned_to_the_merged_ticket() {
2364        let (_directory, local) = open_seeded_local();
2365        insert_ready_ticket(&local, "T2", 1_000);
2366        local
2367            .insert_trigger(
2368                &NewTrigger {
2369                    id: "TR2",
2370                    kind: TriggerKind::Every,
2371                    ticket_id: Some("T1"),
2372                    project_id: None,
2373                    eligible_at_ms: Some(50_000),
2374                    interval_ms: Some(60_000),
2375                },
2376                1_000,
2377            )
2378            .unwrap();
2379        insert_queued_trigger(&local, "TR3", TriggerKind::Immediate, Some("T2"), None);
2380        insert_queued_trigger(&local, "TR4", TriggerKind::Auto, None, Some("default"));
2381
2382        claim_local(&local, "R1").await;
2383        local
2384            .release(&ticket_ref(), &OwnerId("R1".into()), Disposition::Complete)
2385            .await
2386            .unwrap();
2387
2388        assert_eq!(local.ticket("T1").unwrap().unwrap().state, "merged");
2389        assert_eq!(queued_trigger_ids(&local), vec!["TR3", "TR4"]);
2390        assert_eq!(
2391            local
2392                .select_ready_ticket(Some("default"), "TR4", 2_000)
2393                .unwrap(),
2394            Some("T2".to_owned())
2395        );
2396        assert_eq!(
2397            local
2398                .pull_ready()
2399                .await
2400                .unwrap()
2401                .into_iter()
2402                .map(|ticket| ticket.id)
2403                .collect::<Vec<_>>(),
2404            vec!["T2".to_owned()]
2405        );
2406    }
2407
2408    #[tokio::test]
2409    async fn local_work_state_external_merge_retires_pinned_triggers() {
2410        let (_directory, local) = open_seeded_local();
2411        insert_queued_trigger(&local, "TR2", TriggerKind::Immediate, Some("T1"), None);
2412        claim_local(&local, "R1").await;
2413        local
2414            .release(
2415                &ticket_ref(),
2416                &OwnerId("R1".into()),
2417                Disposition::Park {
2418                    reason: "needs-review".into(),
2419                },
2420            )
2421            .await
2422            .unwrap();
2423        assert_eq!(queued_trigger_ids(&local), vec!["TR2"]);
2424
2425        local
2426            .release(&ticket_ref(), &OwnerId("R1".into()), Disposition::Complete)
2427            .await
2428            .unwrap();
2429
2430        assert_eq!(local.ticket("T1").unwrap().unwrap().state, "merged");
2431        assert!(queued_trigger_ids(&local).is_empty());
2432    }
2433
2434    #[tokio::test]
2435    async fn local_work_state_non_merge_dispositions_leave_pinned_triggers_queued() {
2436        for disposition in [
2437            Disposition::Abandon,
2438            Disposition::Park {
2439                reason: "operator review".into(),
2440            },
2441            Disposition::Park {
2442                reason: "needs-review".into(),
2443            },
2444            Disposition::Retry {
2445                not_before_ms: None,
2446            },
2447        ] {
2448            let (_directory, local) = open_seeded_local();
2449            insert_queued_trigger(&local, "TR2", TriggerKind::Immediate, Some("T1"), None);
2450            claim_local(&local, "R1").await;
2451            local
2452                .release(&ticket_ref(), &OwnerId("R1".into()), disposition)
2453                .await
2454                .unwrap();
2455            assert_eq!(queued_trigger_ids(&local), vec!["TR2"]);
2456        }
2457    }
2458
2459    #[test]
2460    fn complete_merged_ticket_triggers_sweeps_only_stranded_pinned_rows() {
2461        let (_directory, local) = open_seeded_local();
2462        insert_ready_ticket(&local, "T2", 1_000);
2463        insert_queued_trigger(&local, "TR2", TriggerKind::Immediate, Some("T2"), None);
2464        insert_queued_trigger(&local, "TR3", TriggerKind::Auto, None, Some("default"));
2465        local
2466            .db
2467            .lock()
2468            .execute("UPDATE tickets SET state = 'merged' WHERE id = 'T1'", [])
2469            .unwrap();
2470
2471        assert_eq!(
2472            local.complete_merged_ticket_triggers(2_000).unwrap(),
2473            vec![("TR1".to_owned(), "T1".to_owned())]
2474        );
2475        assert_eq!(queued_trigger_ids(&local), vec!["TR2", "TR3"]);
2476
2477        assert!(
2478            local
2479                .complete_merged_ticket_triggers(3_000)
2480                .unwrap()
2481                .is_empty()
2482        );
2483        assert_eq!(queued_trigger_ids(&local), vec!["TR2", "TR3"]);
2484    }
2485
2486    #[tokio::test]
2487    async fn local_work_state_denies_renewal_of_an_expired_lease() {
2488        let (_directory, local) = open_seeded_local();
2489        claim_local(&local, "R1").await;
2490        local
2491            .db
2492            .lock()
2493            .execute(
2494                "UPDATE leases SET expires_at_ms = renewed_at_ms WHERE run_id = 'R1'",
2495                [],
2496            )
2497            .unwrap();
2498
2499        assert_eq!(
2500            local
2501                .renew(&ticket_ref(), &OwnerId("R1".into()))
2502                .await
2503                .unwrap(),
2504            ClaimResult::Lost { held_by: None }
2505        );
2506    }
2507
2508    #[tokio::test]
2509    async fn local_work_state_reports_exec_outcomes_with_the_existing_wire_shape() {
2510        let directory = tempdir().unwrap();
2511        let request_path = directory.path().join("report.json");
2512        let feeder = TicketFeeder::exec(
2513            directory.path(),
2514            vec![
2515                "sh".into(),
2516                "-c".into(),
2517                "cat > \"$1\"".into(),
2518                "ticket-source".into(),
2519                request_path.to_string_lossy().into_owned(),
2520            ],
2521        );
2522        let local = LocalSqlite::from_db_with_clock_and_reporter(
2523            Db::open(&directory.path().join("sloop.db"), 1_000).unwrap(),
2524            Arc::new(crate::clock::SystemClock),
2525            feeder.exec_reporter(),
2526        );
2527
2528        local
2529            .push_outcome(&WorkOutcome {
2530                ticket_id: "T1".into(),
2531                owner: OwnerId("R1".into()),
2532                verdict: Outcome::Merged,
2533                branch: Some("sloop/T1".into()),
2534                commit_count: 1,
2535                attempt: 1,
2536                finished_at_ms: 2_000,
2537            })
2538            .await
2539            .unwrap();
2540
2541        let request: serde_json::Value =
2542            serde_json::from_str(&std::fs::read_to_string(request_path).unwrap()).unwrap();
2543        assert_eq!(
2544            request,
2545            serde_json::json!({
2546                "verb": "report",
2547                "ticket": "T1",
2548                "outcome": "merged",
2549            })
2550        );
2551    }
2552
2553    #[tokio::test]
2554    async fn local_work_state_without_an_exec_reporter_keeps_outcome_push_a_no_op() {
2555        let directory = tempdir().unwrap();
2556        let local =
2557            LocalSqlite::from_db(Db::open(&directory.path().join("sloop.db"), 1_000).unwrap());
2558
2559        local
2560            .push_outcome(&WorkOutcome {
2561                ticket_id: "T1".into(),
2562                owner: OwnerId("R1".into()),
2563                verdict: Outcome::Merged,
2564                branch: None,
2565                commit_count: 0,
2566                attempt: 1,
2567                finished_at_ms: 2_000,
2568            })
2569            .await
2570            .unwrap();
2571    }
2572
2573    #[test]
2574    fn successful_sync_records_the_last_sync_timestamp_in_memory() {
2575        let directory = tempdir().unwrap();
2576        assert!(
2577            Command::new("git")
2578                .args(["init", "--quiet"])
2579                .current_dir(directory.path())
2580                .status()
2581                .unwrap()
2582                .success()
2583        );
2584        let work_state =
2585            LocalSqlite::from_db(Db::open(&directory.path().join("sloop.db"), 1_000).unwrap());
2586        work_state
2587            .insert_local_project("default", "projects/default.md", "Default", 1_000)
2588            .unwrap();
2589        let empty_source = TicketFeeder::markdown(directory.path(), "tickets");
2590        let failing_source = TicketFeeder::exec(
2591            directory.path(),
2592            vec![
2593                "sh".into(),
2594                "-c".into(),
2595                "printf 'pull failed' >&2; exit 1".into(),
2596            ],
2597        );
2598        let project_ids = vec!["default".to_owned()];
2599        let drop_runs = |_: &rusqlite::Transaction<'_>, _: &[String], _: &_| Ok(0);
2600        let mark_runs = |_: &rusqlite::Transaction<'_>, _: &str, _: i64| Ok(());
2601
2602        assert_eq!(work_state.last_sync_ms(), None);
2603        work_state
2604            .sync_from_source(
2605                directory.path(),
2606                &empty_source,
2607                &directory.path().join("worktrees"),
2608                2_000,
2609                "T",
2610                &project_ids,
2611                None,
2612                &BTreeMap::new(),
2613                "default",
2614                drop_runs,
2615                mark_runs,
2616            )
2617            .unwrap();
2618        assert_eq!(work_state.last_sync_ms(), Some(2_000));
2619
2620        let error = work_state
2621            .sync_from_source(
2622                directory.path(),
2623                &failing_source,
2624                &directory.path().join("worktrees"),
2625                3_000,
2626                "T",
2627                &project_ids,
2628                None,
2629                &BTreeMap::new(),
2630                "default",
2631                drop_runs,
2632                mark_runs,
2633            )
2634            .unwrap_err();
2635        assert!(error.to_string().contains("pull failed"));
2636        assert_eq!(work_state.last_sync_ms(), Some(2_000));
2637    }
2638
2639    fn claim_t1<'a>(run_id: &'a str) -> ClaimTransaction<'a> {
2640        ClaimTransaction {
2641            ticket_id: "T1",
2642            run_id,
2643            trigger_id: "TR1",
2644            owner_id: "daemon-1",
2645            lease_ms: 60_000,
2646        }
2647    }
2648
2649    #[test]
2650    fn renewing_a_held_lease_extends_its_expiry() {
2651        let directory = tempdir().unwrap();
2652        let store = open_seeded(&directory.path().join("sloop.db"));
2653        granted_claim(&store, &claim_t1("R1"), 2_000);
2654
2655        let expires = store.renew_lease("T1", "R1", 60_000, 10_000).unwrap();
2656        assert_eq!(expires, 70_000);
2657    }
2658
2659    #[test]
2660    fn a_run_cannot_renew_a_lease_it_does_not_hold() {
2661        let directory = tempdir().unwrap();
2662        let store = open_seeded(&directory.path().join("sloop.db"));
2663        granted_claim(&store, &claim_t1("R1"), 2_000);
2664
2665        let error = store.renew_lease("T1", "R2", 60_000, 10_000).unwrap_err();
2666        assert!(matches!(error, StoreError::LeaseNotHeld { .. }));
2667    }
2668
2669    #[test]
2670    fn an_expired_lease_cannot_be_renewed() {
2671        let directory = tempdir().unwrap();
2672        let store = open_seeded(&directory.path().join("sloop.db"));
2673        granted_claim(&store, &claim_t1("R1"), 2_000);
2674
2675        let error = store.renew_lease("T1", "R1", 60_000, 62_000).unwrap_err();
2676        assert!(matches!(error, StoreError::LeaseNotHeld { .. }));
2677    }
2678
2679    #[test]
2680    fn a_readopted_lease_is_re_armed_even_after_it_expired() {
2681        let directory = tempdir().unwrap();
2682        let store = open_seeded(&directory.path().join("sloop.db"));
2683        granted_claim(&store, &claim_t1("R1"), 2_000);
2684
2685        assert!(store.renew_lease("T1", "R1", 60_000, 90_000).is_err());
2686        assert_eq!(
2687            store.readopt_lease("T1", "R1", 60_000, 90_000).unwrap(),
2688            150_000
2689        );
2690        assert_eq!(
2691            store.renew_lease("T1", "R1", 60_000, 100_000).unwrap(),
2692            160_000
2693        );
2694    }
2695
2696    #[test]
2697    fn a_settled_run_cannot_be_readopted() {
2698        let directory = tempdir().unwrap();
2699        let store = open_seeded(&directory.path().join("sloop.db"));
2700        granted_claim(&store, &claim_t1("R1"), 2_000);
2701        settle_for_test(&store, "R1", Outcome::Failed, 3_000);
2702
2703        let error = store.readopt_lease("T1", "R1", 60_000, 4_000).unwrap_err();
2704        assert!(matches!(error, StoreError::LeaseNotHeld { .. }));
2705    }
2706
2707    fn insert_ready_t0(store: &LocalSqlite, now_ms: i64) {
2708        store
2709            .insert_local_ticket(
2710                "T0",
2711                "default",
2712                ".agents/sloop/tickets/t0.md",
2713                "Ticket zero",
2714                &[],
2715                "sloop/T0",
2716                None,
2717                None,
2718                None,
2719                "default",
2720                TicketState::Ready,
2721                now_ms,
2722            )
2723            .unwrap();
2724    }
2725
2726    #[test]
2727    fn an_trigger_pinned_to_another_ticket_does_not_answer_for_this_one() {
2728        let directory = tempdir().unwrap();
2729        let store = open_seeded(&directory.path().join("sloop.db"));
2730        insert_ready_t0(&store, 2_000);
2731
2732        assert!(store.has_claimable_trigger("T1", 2_000).unwrap());
2733        assert!(!store.has_claimable_trigger("T0", 2_000).unwrap());
2734
2735        granted_claim(&store, &claim_t1("R1"), 2_000);
2736        settle_for_test(&store, "R1", Outcome::Merged, 3_000);
2737        store
2738            .insert_trigger(
2739                &NewTrigger {
2740                    id: "TR2",
2741                    kind: TriggerKind::Immediate,
2742                    ticket_id: Some("T1"),
2743                    project_id: None,
2744                    eligible_at_ms: None,
2745                    interval_ms: None,
2746                },
2747                3_000,
2748            )
2749            .unwrap();
2750
2751        assert!(!store.queued_triggers().unwrap().is_empty());
2752        assert!(!store.has_claimable_trigger("T0", 3_000).unwrap());
2753    }
2754
2755    #[test]
2756    fn an_unpinned_trigger_answers_for_every_ticket_it_could_select() {
2757        let directory = tempdir().unwrap();
2758        let store = open_seeded(&directory.path().join("sloop.db"));
2759        insert_ready_t0(&store, 2_000);
2760        store
2761            .insert_trigger(
2762                &NewTrigger {
2763                    id: "TR2",
2764                    kind: TriggerKind::Immediate,
2765                    ticket_id: None,
2766                    project_id: None,
2767                    eligible_at_ms: None,
2768                    interval_ms: None,
2769                },
2770                2_000,
2771            )
2772            .unwrap();
2773        let unpinned = QueuedTrigger {
2774            id: "TR2".into(),
2775            kind: TriggerKind::Immediate,
2776            ticket_id: None,
2777            project_id: None,
2778            eligible_at_ms: None,
2779            interval_ms: None,
2780        };
2781
2782        assert_eq!(
2783            select_ready_ticket(&store, &unpinned, 2_000).as_deref(),
2784            Some("T1")
2785        );
2786        assert!(store.has_claimable_trigger("T1", 2_000).unwrap());
2787        assert!(store.has_claimable_trigger("T0", 2_000).unwrap());
2788
2789        store.insert_trigger_filter("TR2", "T1").unwrap();
2790        assert!(!store.has_claimable_trigger("T0", 2_000).unwrap());
2791    }
2792
2793    #[test]
2794    fn ready_work_selection_is_deterministic_and_respects_filters() {
2795        let directory = tempdir().unwrap();
2796        let store = open_seeded(&directory.path().join("sloop.db"));
2797        store
2798            .insert_local_ticket(
2799                "T0",
2800                "default",
2801                ".agents/sloop/tickets/t0.md",
2802                "Ticket zero",
2803                &[],
2804                "sloop/T0",
2805                None,
2806                None,
2807                None,
2808                "default",
2809                TicketState::Ready,
2810                2_000,
2811            )
2812            .unwrap();
2813        store
2814            .insert_trigger(
2815                &NewTrigger {
2816                    id: "TR2",
2817                    kind: TriggerKind::Immediate,
2818                    ticket_id: None,
2819                    project_id: None,
2820                    eligible_at_ms: None,
2821                    interval_ms: None,
2822                },
2823                2_000,
2824            )
2825            .unwrap();
2826        let trigger = QueuedTrigger {
2827            id: "TR2".into(),
2828            kind: TriggerKind::Immediate,
2829            ticket_id: None,
2830            project_id: None,
2831            eligible_at_ms: None,
2832            interval_ms: None,
2833        };
2834
2835        assert_eq!(
2836            select_ready_ticket(&store, &trigger, 2_000).as_deref(),
2837            Some("T1")
2838        );
2839
2840        store.insert_trigger_filter("TR2", "T0").unwrap();
2841        assert_eq!(
2842            select_ready_ticket(&store, &trigger, 2_000).as_deref(),
2843            Some("T0")
2844        );
2845
2846        let scoped = QueuedTrigger {
2847            project_id: Some("elsewhere".into()),
2848            ..trigger
2849        };
2850        assert_eq!(select_ready_ticket(&store, &scoped, 2_000), None);
2851    }
2852
2853    #[test]
2854    fn tickets_with_unmerged_blockers_are_never_selected() {
2855        let directory = tempdir().unwrap();
2856        let store = open_seeded(&directory.path().join("sloop.db"));
2857        store
2858            .insert_local_ticket(
2859                "T2",
2860                "default",
2861                ".agents/sloop/tickets/t2.md",
2862                "Ticket two",
2863                &["T1".into()],
2864                "sloop/T2",
2865                Some("claude"),
2866                Some("sonnet"),
2867                Some("medium"),
2868                "default",
2869                TicketState::Ready,
2870                1_500,
2871            )
2872            .unwrap();
2873        granted_claim(&store, &claim_t1("R1"), 2_000);
2874
2875        let trigger = QueuedTrigger {
2876            id: "TR1".into(),
2877            kind: TriggerKind::Immediate,
2878            ticket_id: None,
2879            project_id: None,
2880            eligible_at_ms: None,
2881            interval_ms: None,
2882        };
2883        assert_eq!(select_ready_ticket(&store, &trigger, 2_000), None);
2884
2885        settle_for_test(&store, "R1", Outcome::Merged, 3_000);
2886        assert_eq!(
2887            select_ready_ticket(&store, &trigger, 3_000).as_deref(),
2888            Some("T2")
2889        );
2890    }
2891
2892    #[test]
2893    fn missing_tickets_are_not_selected_and_cannot_be_claimed() {
2894        let directory = tempdir().unwrap();
2895        let store = open_seeded(&directory.path().join("sloop.db"));
2896        store.mark_ticket_missing("T1", 2_000).unwrap();
2897        let trigger = QueuedTrigger {
2898            id: "TR1".into(),
2899            kind: TriggerKind::Immediate,
2900            ticket_id: None,
2901            project_id: None,
2902            eligible_at_ms: None,
2903            interval_ms: None,
2904        };
2905        assert_eq!(select_ready_ticket(&store, &trigger, 2_000), None);
2906        assert_eq!(store.ticket("T1").unwrap().unwrap().attempts, 0);
2907
2908        store.mark_ticket_missing("T1", 5_000).unwrap();
2909        assert_eq!(
2910            store.local_ticket_files().unwrap()[0].missing_at_ms,
2911            Some(2_000)
2912        );
2913        store.clear_ticket_missing("T1", 6_000).unwrap();
2914        assert_eq!(
2915            select_ready_ticket(&store, &trigger, 6_000).as_deref(),
2916            Some("T1")
2917        );
2918    }
2919
2920    #[test]
2921    fn blockers_gate_selection_claims_and_derived_counts_until_merged() {
2922        let directory = tempdir().unwrap();
2923        let store = open_seeded(&directory.path().join("sloop.db"));
2924        store
2925            .insert_local_ticket(
2926                "T2",
2927                "default",
2928                ".agents/sloop/tickets/t2.md",
2929                "Ticket two",
2930                &["T1".into()],
2931                "sloop/T2",
2932                Some("claude"),
2933                None,
2934                None,
2935                "default",
2936                TicketState::Ready,
2937                1_500,
2938            )
2939            .unwrap();
2940        let trigger = QueuedTrigger {
2941            id: "TR1".into(),
2942            kind: TriggerKind::Immediate,
2943            ticket_id: None,
2944            project_id: None,
2945            eligible_at_ms: None,
2946            interval_ms: None,
2947        };
2948
2949        assert_eq!(store.unmerged_blockers("T2").unwrap(), ["T1"]);
2950        assert_eq!(
2951            select_ready_ticket(&store, &trigger, 2_000).as_deref(),
2952            Some("T1")
2953        );
2954        assert_eq!(store.ticket_counts().unwrap().blocked, 1);
2955
2956        store
2957            .db
2958            .lock()
2959            .execute("UPDATE tickets SET state = 'failed' WHERE id = 'T1'", [])
2960            .unwrap();
2961        assert_eq!(select_ready_ticket(&store, &trigger, 2_000), None);
2962        assert_eq!(store.ticket("T2").unwrap().unwrap().attempts, 0);
2963
2964        store
2965            .db
2966            .lock()
2967            .execute("UPDATE tickets SET state = 'merged' WHERE id = 'T1'", [])
2968            .unwrap();
2969        assert!(store.unmerged_blockers("T2").unwrap().is_empty());
2970        assert_eq!(
2971            select_ready_ticket(&store, &trigger, 2_000).as_deref(),
2972            Some("T2")
2973        );
2974        let counts = store.ticket_counts().unwrap();
2975        assert_eq!(counts.ready, 1);
2976        assert_eq!(counts.blocked, 0);
2977    }
2978
2979    #[test]
2980    fn state_survives_reopening_the_database() {
2981        let directory = tempdir().unwrap();
2982        let path = directory.path().join("sloop.db");
2983        let store = open_seeded(&path);
2984        granted_claim(&store, &claim_t1("R1"), 2_000);
2985        drop(store);
2986
2987        let store = LocalSqlite::from_db(Db::open(&path, 3_000).unwrap());
2988        assert_eq!(
2989            store.ticket_state("T1").unwrap().as_deref(),
2990            Some("claimed")
2991        );
2992        assert_eq!(store.ticket_counts().unwrap().claimed, 1);
2993        let ticket = store.ticket("T1").unwrap().unwrap();
2994        assert_eq!(ticket.target.as_deref(), Some("claude"));
2995        assert_eq!(ticket.model.as_deref(), Some("sonnet"));
2996        assert_eq!(ticket.effort.as_deref(), Some("medium"));
2997        assert_eq!(ticket.name, "Ticket one");
2998        assert!(ticket.blocked_by.is_empty());
2999        assert_eq!(ticket.worktree.as_deref(), Some("sloop/T1"));
3000    }
3001
3002    #[test]
3003    fn blocked_by_and_worktree_round_trip() {
3004        let directory = tempdir().unwrap();
3005        let path = directory.path().join("sloop.db");
3006        let store = open_seeded(&path);
3007        store
3008            .insert_local_ticket(
3009                "T2",
3010                "default",
3011                ".agents/sloop/tickets/t2.md",
3012                "Ticket two",
3013                &["T1".to_owned()],
3014                "feature/t2",
3015                None,
3016                None,
3017                None,
3018                "default",
3019                TicketState::Ready,
3020                2_000,
3021            )
3022            .unwrap();
3023        drop(store);
3024
3025        let ticket = LocalSqlite::from_db(Db::open(&path, 3_000).unwrap())
3026            .ticket("T2")
3027            .unwrap()
3028            .unwrap();
3029        assert_eq!(ticket.name, "Ticket two");
3030        assert_eq!(ticket.blocked_by, ["T1"]);
3031        assert_eq!(ticket.worktree.as_deref(), Some("feature/t2"));
3032    }
3033
3034    #[test]
3035    fn tickets_are_ordered_newest_first_and_include_attempts() {
3036        let directory = tempdir().unwrap();
3037        let store = open_seeded(&directory.path().join("sloop.db"));
3038        store
3039            .insert_local_project("alpha", ".agents/sloop/projects/alpha.md", "Alpha", 1_000)
3040            .unwrap();
3041        for (id, project, state, created_at_ms) in [
3042            ("T0", "alpha", TicketState::Held, 3_000),
3043            ("T2", "default", TicketState::Ready, 1_000),
3044        ] {
3045            store
3046                .insert_local_ticket(
3047                    id,
3048                    project,
3049                    &format!(".agents/sloop/tickets/{}.md", id.to_lowercase()),
3050                    id,
3051                    &[],
3052                    &format!("sloop/{id}"),
3053                    None,
3054                    None,
3055                    None,
3056                    "default",
3057                    state,
3058                    created_at_ms,
3059                )
3060                .unwrap();
3061        }
3062        granted_claim(&store, &claim_t1("R1"), 2_000);
3063
3064        let tickets = store.tickets().unwrap();
3065        assert_eq!(
3066            tickets
3067                .iter()
3068                .map(|ticket| ticket.id.as_str())
3069                .collect::<Vec<_>>(),
3070            ["T0", "T2", "T1"]
3071        );
3072        assert_eq!(tickets[2].attempts, 1);
3073    }
3074
3075    #[test]
3076    fn operator_hold_transitions_are_narrow_and_idempotent() {
3077        let directory = tempdir().unwrap();
3078        let store = open_seeded(&directory.path().join("sloop.db"));
3079        assert_eq!(
3080            store
3081                .set_ticket_hold("T1", TicketState::Held, 2_000)
3082                .unwrap(),
3083            "ready"
3084        );
3085        assert_eq!(store.ticket_counts().unwrap().held, 1);
3086        assert_eq!(
3087            store
3088                .set_ticket_hold("T1", TicketState::Held, 2_100)
3089                .unwrap(),
3090            "held"
3091        );
3092        assert_eq!(
3093            store
3094                .set_ticket_hold("T1", TicketState::Ready, 2_200)
3095                .unwrap(),
3096            "held"
3097        );
3098    }
3099
3100    #[test]
3101    fn validation_hold_reasons_set_and_clear_without_releasing_operator_holds() {
3102        let directory = tempdir().unwrap();
3103        let store = open_seeded(&directory.path().join("sloop.db"));
3104        let ticket = |held_reason: Option<&str>| ReindexTicket {
3105            id: "T1".into(),
3106            project_id: "default".into(),
3107            source: "markdown".into(),
3108            source_ref: ".agents/sloop/tickets/t1.md".into(),
3109            file_path: Some(".agents/sloop/tickets/t1.md".into()),
3110            name: "Ticket one".into(),
3111            blocked_by: Vec::new(),
3112            worktree: "sloop/T1".into(),
3113            target: Some("claude".into()),
3114            model: Some("sonnet".into()),
3115            effort: Some("medium".into()),
3116            flow: "default".into(),
3117            body: "work".into(),
3118            held_reason: held_reason.map(str::to_owned),
3119            derived_state: None,
3120        };
3121
3122        apply_reindex(
3123            &store,
3124            &[ticket(Some("flow `missing` is not defined"))],
3125            2_000,
3126        );
3127        assert_eq!(
3128            store.ticket("T1").unwrap().unwrap().held_reason.as_deref(),
3129            Some("flow `missing` is not defined")
3130        );
3131        apply_reindex(&store, &[ticket(None)], 2_100);
3132        assert_eq!(store.ticket_state("T1").unwrap().as_deref(), Some("ready"));
3133
3134        store
3135            .set_ticket_hold("T1", TicketState::Held, 2_200)
3136            .unwrap();
3137        apply_reindex(&store, &[ticket(None)], 2_300);
3138        let operator_held = store.ticket("T1").unwrap().unwrap();
3139        assert_eq!(operator_held.state, "held");
3140        assert_eq!(operator_held.held_reason, None);
3141    }
3142
3143    #[test]
3144    fn operator_hold_cannot_steal_a_claim() {
3145        let directory = tempdir().unwrap();
3146        let store = open_seeded(&directory.path().join("sloop.db"));
3147        granted_claim(&store, &claim_t1("R1"), 2_000);
3148        assert!(matches!(
3149            store.set_ticket_hold("T1", TicketState::Held, 2_100),
3150            Err(StoreError::TicketStateConflict { state, .. }) if state == "claimed"
3151        ));
3152    }
3153
3154    #[test]
3155    fn retry_only_requeues_failed_tickets_and_resets_attempts() {
3156        let directory = tempdir().unwrap();
3157        let store = open_seeded(&directory.path().join("sloop.db"));
3158        assert_eq!(granted_claim(&store, &claim_t1("R1"), 2_000), 1);
3159        settle_for_test(&store, "R1", Outcome::Failed, 2_100);
3160
3161        assert_eq!(
3162            store.retry_ticket("T1", 2_200, |_, _, _| Ok(())).unwrap(),
3163            "failed"
3164        );
3165        store
3166            .insert_trigger(
3167                &NewTrigger {
3168                    id: "TR2",
3169                    kind: TriggerKind::Immediate,
3170                    ticket_id: Some("T1"),
3171                    project_id: None,
3172                    eligible_at_ms: None,
3173                    interval_ms: None,
3174                },
3175                2_300,
3176            )
3177            .unwrap();
3178        assert_eq!(
3179            granted_claim(
3180                &store,
3181                &ClaimTransaction {
3182                    trigger_id: "TR2",
3183                    ..claim_t1("R2")
3184                },
3185                2_300,
3186            ),
3187            2
3188        );
3189        assert_eq!(store.ticket("T1").unwrap().unwrap().attempts, 1);
3190        assert!(matches!(
3191            store.retry_ticket("T1", 2_400, |_, _, _| Ok(())),
3192            Err(StoreError::TicketStateConflict { state, .. }) if state == "claimed"
3193        ));
3194        assert!(matches!(
3195            store.retry_ticket("missing", 2_400, |_, _, _| Ok(())),
3196            Err(StoreError::TicketNotFound { .. })
3197        ));
3198    }
3199
3200    #[test]
3201    fn configured_default_backfills_tickets_that_predate_target_snapshots() {
3202        let directory = tempdir().unwrap();
3203        let store = open_seeded(&directory.path().join("sloop.db"));
3204        store
3205            .update_ticket_execution("T1", None, Some("sonnet"), Some("medium"), 2_000)
3206            .unwrap();
3207        assert_eq!(store.backfill_ticket_targets("codex", 3_000).unwrap(), 1);
3208        assert_eq!(
3209            store.ticket("T1").unwrap().unwrap().target.as_deref(),
3210            Some("codex")
3211        );
3212        assert_eq!(store.backfill_ticket_targets("claude", 4_000).unwrap(), 0);
3213    }
3214
3215    /// The rename migration must be lossless. A queued trigger is the durable
3216    /// record that someone asked for work; it is in no committed file and in no
3217    /// commit, so `reindex` cannot rebuild one and a migration that dropped it
3218    /// would silently cancel a scheduled run. The ids move too, which means
3219    /// every column pointing at one has to move with them — SQLite propagates a
3220    /// table rename into a `REFERENCES` clause but never propagates an `UPDATE`
3221    /// to the rows that reference it.
3222    #[test]
3223    fn the_rename_migration_carries_every_trigger_id_and_reference_across() {
3224        let directory = tempdir().unwrap();
3225        let path = directory.path().join("sloop.db");
3226
3227        {
3228            let local = open_seeded(&path);
3229            insert_ready_ticket(&local, "T2", 1_000);
3230            local
3231                .insert_trigger(
3232                    &NewTrigger {
3233                        id: "TR2",
3234                        kind: TriggerKind::Every,
3235                        ticket_id: None,
3236                        project_id: Some("default"),
3237                        eligible_at_ms: Some(5_000),
3238                        interval_ms: Some(60_000),
3239                    },
3240                    1_100,
3241                )
3242                .unwrap();
3243            local.insert_trigger_filter("TR2", "T2").unwrap();
3244            insert_queued_trigger(&local, "TR3", TriggerKind::Auto, Some("T2"), None);
3245            trigger::complete_for_ticket(&local.db.lock(), "T2", 1_150).unwrap();
3246            local
3247                .db
3248                .lock()
3249                .execute_batch(
3250                    "INSERT INTO runs
3251                         (id, trigger_id, ticket_id, state, attempt, flow_json, ticket_json,
3252                          created_at_ms, updated_at_ms)
3253                     VALUES ('R1', 'TR1', 'T1', 'claimed', 1, '{}', '{}', 1200, 1200);
3254                     INSERT INTO leases
3255                         (ticket_id, run_id, owner_id, acquired_at_ms, renewed_at_ms,
3256                          expires_at_ms)
3257                     VALUES ('T1', 'R1', '{\"owner\":\"R1\",\"trigger\":\"TR1\"}',
3258                             1200, 1200, 601200);",
3259                )
3260                .unwrap();
3261        }
3262
3263        {
3264            let connection = rusqlite::Connection::open(&path).unwrap();
3265            connection.execute_batch(REVERT_TRIGGER_RENAME).unwrap();
3266            connection.pragma_update(None, "user_version", 15).unwrap();
3267        }
3268        let local = LocalSqlite::from_db(Db::open(&path, 2_000).unwrap());
3269
3270        let queued = local.queued_triggers().unwrap();
3271        assert_eq!(
3272            queued.iter().map(|t| t.id.as_str()).collect::<Vec<_>>(),
3273            ["TR1", "TR2"]
3274        );
3275        let recurring = &queued[1];
3276        assert_eq!(recurring.kind, TriggerKind::Every);
3277        assert_eq!(recurring.project_id.as_deref(), Some("default"));
3278        assert_eq!(recurring.eligible_at_ms, Some(5_000));
3279        assert_eq!(recurring.interval_ms, Some(60_000));
3280
3281        let connection = local.db.lock();
3282        let states: Vec<(String, String)> = connection
3283            .prepare("SELECT id, state FROM triggers ORDER BY id")
3284            .unwrap()
3285            .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
3286            .unwrap()
3287            .map(Result::unwrap)
3288            .collect();
3289        assert_eq!(
3290            states,
3291            [
3292                ("TR1".to_owned(), "queued".to_owned()),
3293                ("TR2".to_owned(), "queued".to_owned()),
3294                ("TR3".to_owned(), "completed".to_owned()),
3295            ]
3296        );
3297
3298        let filter: (String, String) = connection
3299            .query_row(
3300                "SELECT trigger_id, ticket_id FROM trigger_filters",
3301                [],
3302                |row| Ok((row.get(0)?, row.get(1)?)),
3303            )
3304            .unwrap();
3305        assert_eq!(filter, ("TR2".to_owned(), "T2".to_owned()));
3306        let run_trigger: String = connection
3307            .query_row("SELECT trigger_id FROM runs WHERE id = 'R1'", [], |row| {
3308                row.get(0)
3309            })
3310            .unwrap();
3311        assert_eq!(run_trigger, "TR1");
3312        let stored_owner: String = connection
3313            .query_row(
3314                "SELECT owner_id FROM leases WHERE ticket_id = 'T1'",
3315                [],
3316                |row| row.get(0),
3317            )
3318            .unwrap();
3319        let (owner, claimed_trigger) = super::decode_lease_owner(&stored_owner);
3320        assert_eq!(owner.0, "R1");
3321        assert_eq!(claimed_trigger.as_deref(), Some("TR1"));
3322        assert_eq!(stored_owner, super::lease_owner(&owner, "TR1"));
3323
3324        assert_eq!(
3325            connection
3326                .prepare("PRAGMA foreign_key_check")
3327                .unwrap()
3328                .query_map([], |_| Ok(()))
3329                .unwrap()
3330                .count(),
3331            0
3332        );
3333        drop(connection);
3334
3335        assert_eq!(
3336            local
3337                .enqueue_trigger(
3338                    &trigger::EnqueueRequest {
3339                        kind: TriggerKind::Immediate,
3340                        ticket_id: None,
3341                        project_id: None,
3342                        eligible_at_ms: None,
3343                        interval_ms: None,
3344                        filters: &[],
3345                        duplicates: trigger::Duplicates::Allow,
3346                    },
3347                    2_000,
3348                )
3349                .unwrap()
3350                .id,
3351            "TR4"
3352        );
3353    }
3354}