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