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 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 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 #[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 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#[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 #[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}