fix: preserve dependency queue atomicity

This commit is contained in:
2026-08-25 12:05:43 +09:00
parent a41147916b
commit 5857e6121c
+161 -3
View File
@@ -4702,6 +4702,24 @@ impl TicketBackend for LocalTicketBackend {
planned.push((dir, target)); planned.push((dir, target));
} }
let snapshots = planned
.iter()
.map(|(dir, _)| {
let item_path = dir.join("item.md");
let thread_path = dir.join("thread.md");
Ok((
item_path.clone(),
fs::read(&item_path).map_err(|error| io_err(&item_path, error))?,
thread_path.clone(),
if thread_path.exists() {
Some(fs::read(&thread_path).map_err(|error| io_err(&thread_path, error))?)
} else {
None
},
))
})
.collect::<Result<Vec<_>>>()?;
for (dir, target) in planned { for (dir, target) in planned {
let at = now_utc(); let at = now_utc();
let mut change = TicketStateChange::new( let mut change = TicketStateChange::new(
@@ -4711,7 +4729,7 @@ impl TicketBackend for LocalTicketBackend {
self.queued_ready_body(queued_by), self.queued_ready_body(queued_by),
); );
change.author = Some(queued_by.to_string()); change.author = Some(queued_by.to_string());
self.apply_workflow_state_change( if let Err(error) = self.apply_workflow_state_change(
&dir, &dir,
TicketWorkflowState::Ready, TicketWorkflowState::Ready,
TicketWorkflowState::Queued, TicketWorkflowState::Queued,
@@ -4723,7 +4741,29 @@ impl TicketBackend for LocalTicketBackend {
("ref_selector", target.ref_selector.as_str()), ("ref_selector", target.ref_selector.as_str()),
("queue_root_ticket", requested_ticket.as_str()), ("queue_root_ticket", requested_ticket.as_str()),
], ],
)?; ) {
for (item_path, item_before, thread_path, thread_before) in &snapshots {
if fs::read(item_path).ok().as_deref() != Some(item_before.as_slice()) {
fs::write(item_path, item_before)
.map_err(|rollback| io_err(item_path, rollback))?;
}
match thread_before {
Some(thread_before)
if fs::read(thread_path).ok().as_deref()
!= Some(thread_before.as_slice()) =>
{
fs::write(thread_path, thread_before)
.map_err(|rollback| io_err(thread_path, rollback))?;
}
None if thread_path.exists() => {
fs::remove_file(thread_path)
.map_err(|rollback| io_err(thread_path, rollback))?;
}
_ => {}
}
}
return Err(error);
}
} }
Ok(TicketQueueOutcome { Ok(TicketQueueOutcome {
@@ -5690,11 +5730,19 @@ fn dependency_queue_plan(
fn visit( fn visit(
ticket: &str, ticket: &str,
states: &HashMap<String, TicketWorkflowState>,
prerequisites: &BTreeMap<String, BTreeSet<String>>, prerequisites: &BTreeMap<String, BTreeSet<String>>,
marks: &mut BTreeMap<String, u8>, marks: &mut BTreeMap<String, u8>,
stack: &mut Vec<String>, stack: &mut Vec<String>,
ordered: &mut Vec<String>, ordered: &mut Vec<String>,
) -> Result<()> { ) -> Result<()> {
let state = states
.get(ticket)
.copied()
.ok_or_else(|| TicketError::NotFound(ticket.to_owned()))?;
if ticket_state_resolved(state) {
return Ok(());
}
match marks.get(ticket).copied() { match marks.get(ticket).copied() {
Some(2) => return Ok(()), Some(2) => return Ok(()),
Some(1) => { Some(1) => {
@@ -5713,7 +5761,7 @@ fn dependency_queue_plan(
stack.push(ticket.to_owned()); stack.push(ticket.to_owned());
if let Some(dependencies) = prerequisites.get(ticket) { if let Some(dependencies) = prerequisites.get(ticket) {
for dependency in dependencies { for dependency in dependencies {
visit(dependency, prerequisites, marks, stack, ordered)?; visit(dependency, states, prerequisites, marks, stack, ordered)?;
} }
} }
stack.pop(); stack.pop();
@@ -5725,6 +5773,7 @@ fn dependency_queue_plan(
let mut ordered = Vec::new(); let mut ordered = Vec::new();
visit( visit(
requested_ticket, requested_ticket,
states,
&prerequisites, &prerequisites,
&mut BTreeMap::new(), &mut BTreeMap::new(),
&mut Vec::new(), &mut Vec::new(),
@@ -7323,6 +7372,7 @@ mod tests {
let relations = [ let relations = [
dependency_relation("root", "done"), dependency_relation("root", "done"),
dependency_relation("done", "planning"), dependency_relation("done", "planning"),
dependency_relation("done", "root"),
]; ];
assert_eq!( assert_eq!(
@@ -7740,6 +7790,59 @@ state: planning
); );
} }
#[test]
fn sqlite_queue_ignores_cycle_behind_done_dependency() {
let temp = TempDir::new().unwrap();
let backend = SqliteTicketBackend::open(temp.path().join("tickets.db"), "workspace-test")
.unwrap()
.with_target_authority(Arc::new(TestTargetAuthority));
let mut root_input = NewTicket::new("Ready root");
root_input.workflow_state = Some(TicketWorkflowState::Ready);
root_input.repository_id = Some("main".to_string());
let root = backend.create(root_input).unwrap();
let mut done_input = NewTicket::new("Done dependency");
done_input.workflow_state = Some(TicketWorkflowState::Done);
done_input.repository_id = Some("main".to_string());
let done = backend.create(done_input).unwrap();
for (ticket_id, target) in [
(root.id.clone(), done.id.clone()),
(done.id.clone(), root.id.clone()),
] {
backend
.add_ticket_relation(
TicketIdOrSlug::Id(ticket_id),
NewTicketRelation {
kind: TicketRelationKind::DependsOn,
target,
note: None,
author: None,
},
)
.unwrap();
}
let outcome = backend
.queue_ready(TicketIdOrSlug::Id(root.id.clone()), "orchestrator")
.unwrap();
assert_eq!(outcome.queued_tickets, vec![root.id.clone()]);
assert_eq!(
backend
.show(TicketIdOrSlug::Id(root.id))
.unwrap()
.meta
.workflow_state,
TicketWorkflowState::Queued
);
assert_eq!(
backend
.show(TicketIdOrSlug::Id(done.id))
.unwrap()
.meta
.workflow_state,
TicketWorkflowState::Done
);
}
#[test] #[test]
fn sqlite_queue_cycle_diagnostic_leaves_all_tickets_ready() { fn sqlite_queue_cycle_diagnostic_leaves_all_tickets_ready() {
let temp = TempDir::new().unwrap(); let temp = TempDir::new().unwrap();
@@ -8680,6 +8783,61 @@ state: planning
); );
} }
#[cfg(unix)]
#[test]
fn local_queue_rolls_back_ready_dependency_when_later_write_fails() {
use std::os::unix::fs::PermissionsExt;
let tmp = TempDir::new().unwrap();
let backend = backend(&tmp);
let mut dependency_input = NewTicket::new("Ready dependency");
dependency_input.workflow_state = Some(TicketWorkflowState::Ready);
dependency_input.repository_id = Some("main".to_string());
dependency_input.ref_selector = Some("develop".to_string());
let dependency = backend.create(dependency_input).unwrap();
let mut root_input = NewTicket::new("Ready root");
root_input.workflow_state = Some(TicketWorkflowState::Ready);
root_input.repository_id = Some("main".to_string());
root_input.ref_selector = Some("develop".to_string());
let root = backend.create(root_input).unwrap();
backend
.add_ticket_relation(
TicketIdOrSlug::Id(root.id.clone()),
NewTicketRelation {
kind: TicketRelationKind::DependsOn,
target: dependency.id.clone(),
note: None,
author: None,
},
)
.unwrap();
let root_dir = backend
.find_ticket_dir(&TicketIdOrSlug::Id(root.id.clone()))
.unwrap();
let root_item = root_dir.join("item.md");
fs::set_permissions(&root_dir, fs::Permissions::from_mode(0o555)).unwrap();
fs::set_permissions(&root_item, fs::Permissions::from_mode(0o444)).unwrap();
let result = backend.queue_ready(TicketIdOrSlug::Id(root.id.clone()), "test");
fs::set_permissions(&root_dir, fs::Permissions::from_mode(0o755)).unwrap();
fs::set_permissions(&root_item, fs::Permissions::from_mode(0o644)).unwrap();
assert!(result.is_err());
for ticket_id in [dependency.id, root.id] {
let ticket = backend.show(TicketIdOrSlug::Id(ticket_id)).unwrap();
assert_eq!(ticket.meta.workflow_state, TicketWorkflowState::Ready);
assert!(
!ticket
.events
.iter()
.any(|event| event.to.as_deref() == Some("queued")),
"Ticket {} retained queued events: {:?}",
ticket.meta.id,
ticket.events
);
}
}
#[test] #[test]
fn state_cannot_be_changed_through_generic_field_api() { fn state_cannot_be_changed_through_generic_field_api() {
let tmp = TempDir::new().unwrap(); let tmp = TempDir::new().unwrap();