feat: queue ready dependency closures atomically
This commit is contained in:
@@ -3047,10 +3047,7 @@ VALUES ('workspace-test', 'ticket', 4);
|
||||
assert_eq!(tickets.items[0].record_source, "sqlite_yoi_ticket");
|
||||
assert_eq!(tickets.items[0].id, "00000000001J2");
|
||||
assert_eq!(tickets.items[0].state, "ready");
|
||||
assert_eq!(
|
||||
tickets.items[0].workspace_action_priority,
|
||||
"ready_for_queue"
|
||||
);
|
||||
assert_eq!(tickets.items[0].workspace_action_priority, "background");
|
||||
let ticket_by_key = authority.ticket(&tickets.items[0].resource_key).unwrap();
|
||||
assert_eq!(ticket_by_key.id, tickets.items[0].id);
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
use std::collections::{BTreeMap, HashMap, HashSet};
|
||||
use std::collections::{BTreeMap, BTreeSet, HashMap, HashSet};
|
||||
use std::path::{Component, Path, PathBuf};
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::sync::{Arc, Mutex, Weak};
|
||||
@@ -3951,6 +3951,39 @@ fn reject_unguarded_ticket_completion(operation: &TicketBackendOperation) -> Res
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn queue_assignment_candidates(
|
||||
backend: &dyn TicketBackend,
|
||||
requested_ticket_id: &str,
|
||||
) -> ticket::Result<Vec<String>> {
|
||||
fn visit(
|
||||
backend: &dyn TicketBackend,
|
||||
ticket_id: &str,
|
||||
visited: &mut BTreeSet<String>,
|
||||
ready: &mut BTreeSet<String>,
|
||||
) -> ticket::Result<()> {
|
||||
if !visited.insert(ticket_id.to_owned()) {
|
||||
return Ok(());
|
||||
}
|
||||
let ticket = backend.show(TicketIdOrSlug::Id(ticket_id.to_owned()))?;
|
||||
if ticket.meta.workflow_state == TicketWorkflowState::Ready {
|
||||
ready.insert(ticket.meta.id.clone());
|
||||
}
|
||||
for blocker in ticket.relations.blockers {
|
||||
visit(backend, &blocker.blocking_ticket, visited, ready)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
let mut ready = BTreeSet::new();
|
||||
visit(
|
||||
backend,
|
||||
requested_ticket_id,
|
||||
&mut BTreeSet::new(),
|
||||
&mut ready,
|
||||
)?;
|
||||
Ok(ready.into_iter().collect())
|
||||
}
|
||||
|
||||
async fn execute_ticket_rest_operation(
|
||||
api: &WorkspaceApi,
|
||||
workspace_id: &str,
|
||||
@@ -3988,27 +4021,34 @@ async fn execute_ticket_rest_operation(
|
||||
.to_string(),
|
||||
)
|
||||
})?;
|
||||
let assignment = active_orchestrator_assignment(api, workspace_id, &ticket.meta.id)?
|
||||
.ok_or_else(|| {
|
||||
Error::TicketAssignmentConflict(
|
||||
"Queue requires role=orchestrator assignment to workspace-orchestrator"
|
||||
.to_string(),
|
||||
)
|
||||
})?;
|
||||
let candidates =
|
||||
queue_assignment_candidates(&backend, &ticket.meta.id).map_err(Error::from)?;
|
||||
let mut assignment_ids = BTreeMap::new();
|
||||
for ticket_id in candidates {
|
||||
let assignment = active_orchestrator_assignment(api, workspace_id, &ticket_id)?
|
||||
.ok_or_else(|| {
|
||||
Error::TicketAssignmentConflict(format!(
|
||||
"Queue requires role=orchestrator assignment for Ticket {ticket_id}"
|
||||
))
|
||||
})?;
|
||||
assignment_ids.insert(ticket_id, assignment.assignment_id);
|
||||
}
|
||||
let assignment_json = serde_json::to_string(&assignment_ids).map_err(|error| {
|
||||
Error::Config(format!("failed to encode Queue assignment fence: {error}"))
|
||||
})?;
|
||||
let operation_id = new_id("tqueue");
|
||||
let fingerprint = Sha256::digest(format!(
|
||||
"ticket-queue:v1\0{workspace_id}\0{}\0{}\0{}",
|
||||
"ticket-queue:v2\0{workspace_id}\0{}\0{}\0{assignment_json}",
|
||||
ticket.meta.id,
|
||||
ticket.meta.workflow_state.as_str(),
|
||||
assignment.assignment_id
|
||||
))
|
||||
.iter()
|
||||
.map(|byte| format!("{byte:02x}"))
|
||||
.collect::<String>();
|
||||
event_attributes.extend([
|
||||
(
|
||||
"orchestrator_assignment_id".to_string(),
|
||||
assignment.assignment_id,
|
||||
"queue_orchestrator_assignments".to_string(),
|
||||
assignment_json,
|
||||
),
|
||||
(
|
||||
"routing_principal".to_string(),
|
||||
@@ -4029,18 +4069,30 @@ async fn execute_ticket_rest_operation(
|
||||
}
|
||||
|
||||
let result = execute_ticket_backend_operation(&backend, operation).map_err(Error::from)?;
|
||||
if is_mutation
|
||||
&& let Some(target) = target
|
||||
&& let Ok(ticket) = backend.show(target)
|
||||
{
|
||||
notify_ticket_recipients(
|
||||
api,
|
||||
workspace_id,
|
||||
&ticket.meta.id,
|
||||
&previous_state,
|
||||
ticket.meta.workflow_state.as_str(),
|
||||
source,
|
||||
);
|
||||
if is_mutation {
|
||||
if let TicketBackendOperationResult::QueueOutcome(outcome) = &result {
|
||||
for ticket_id in &outcome.queued_tickets {
|
||||
notify_ticket_recipients(
|
||||
api,
|
||||
workspace_id,
|
||||
ticket_id,
|
||||
TicketWorkflowState::Ready.as_str(),
|
||||
TicketWorkflowState::Queued.as_str(),
|
||||
source.clone(),
|
||||
);
|
||||
}
|
||||
} else if let Some(target) = target
|
||||
&& let Ok(ticket) = backend.show(target)
|
||||
{
|
||||
notify_ticket_recipients(
|
||||
api,
|
||||
workspace_id,
|
||||
&ticket.meta.id,
|
||||
&previous_state,
|
||||
ticket.meta.workflow_state.as_str(),
|
||||
source,
|
||||
);
|
||||
}
|
||||
}
|
||||
Ok(result)
|
||||
}
|
||||
@@ -4364,7 +4416,7 @@ async fn scoped_queue_ticket_record(
|
||||
State(api): State<WorkspaceApi>,
|
||||
AxumPath((workspace_id, id)): AxumPath<(String, String)>,
|
||||
headers: HeaderMap,
|
||||
) -> ApiResult<StatusCode> {
|
||||
) -> ApiResult<Json<ticket::TicketQueueOutcome>> {
|
||||
let result = execute_ticket_rest_operation(
|
||||
&api,
|
||||
&workspace_id,
|
||||
@@ -4375,7 +4427,10 @@ async fn scoped_queue_ticket_record(
|
||||
},
|
||||
)
|
||||
.await?;
|
||||
ticket_rest_unit(result)
|
||||
ticket_rest_result(result, |result| match result {
|
||||
TicketBackendOperationResult::QueueOutcome(outcome) => Some(outcome),
|
||||
_ => None,
|
||||
})
|
||||
}
|
||||
|
||||
#[derive(Debug, serde::Deserialize)]
|
||||
@@ -16404,7 +16459,7 @@ mod tests {
|
||||
);
|
||||
|
||||
assign_test_orchestrator(&api, &ticket.id);
|
||||
scoped_queue_ticket_record(State(api.clone()), AxumPath(path), HeaderMap::new())
|
||||
let _ = scoped_queue_ticket_record(State(api.clone()), AxumPath(path), HeaderMap::new())
|
||||
.await
|
||||
.unwrap();
|
||||
let queued = backend.show(ticket.id.into()).unwrap();
|
||||
@@ -16422,6 +16477,88 @@ mod tests {
|
||||
assert!(event.attributes.contains_key("routing_request_fingerprint"));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn queue_requires_assignments_for_every_ready_dependency() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
init_clean_git_workspace(dir.path());
|
||||
let api = test_api(dir.path()).await;
|
||||
let backend = browser_ticket_backend(&api).unwrap();
|
||||
let mut dependency_input = ticket::NewTicket::new("Ready dependency");
|
||||
dependency_input.workflow_state = Some(TicketWorkflowState::Ready);
|
||||
dependency_input.repository_id = Some(TEST_REPOSITORY_ID.to_string());
|
||||
dependency_input.ref_selector = Some("develop".to_string());
|
||||
let dependency = backend.create(dependency_input).unwrap();
|
||||
let mut root_input = ticket::NewTicket::new("Queue root");
|
||||
root_input.workflow_state = Some(TicketWorkflowState::Ready);
|
||||
root_input.repository_id = Some(TEST_REPOSITORY_ID.to_string());
|
||||
root_input.ref_selector = Some("develop".to_string());
|
||||
let root = backend.create(root_input).unwrap();
|
||||
backend
|
||||
.add_ticket_relation(
|
||||
root.id.clone().into(),
|
||||
ticket::NewTicketRelation {
|
||||
kind: ticket::TicketRelationKind::DependsOn,
|
||||
target: dependency.id.clone(),
|
||||
note: Some("must queue first".to_string()),
|
||||
author: Some("test".to_string()),
|
||||
},
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assign_test_orchestrator(&api, &root.id);
|
||||
let error = scoped_queue_ticket_record(
|
||||
State(api.clone()),
|
||||
AxumPath((TEST_WORKSPACE_ID.to_string(), root.id.clone())),
|
||||
HeaderMap::new(),
|
||||
)
|
||||
.await
|
||||
.unwrap_err()
|
||||
.into_response();
|
||||
assert_eq!(error.status(), StatusCode::CONFLICT);
|
||||
assert_eq!(
|
||||
backend
|
||||
.show(root.id.clone().into())
|
||||
.unwrap()
|
||||
.meta
|
||||
.workflow_state,
|
||||
TicketWorkflowState::Ready
|
||||
);
|
||||
assert_eq!(
|
||||
backend
|
||||
.show(dependency.id.clone().into())
|
||||
.unwrap()
|
||||
.meta
|
||||
.workflow_state,
|
||||
TicketWorkflowState::Ready
|
||||
);
|
||||
|
||||
assign_test_orchestrator(&api, &dependency.id);
|
||||
let Json(outcome) = scoped_queue_ticket_record(
|
||||
State(api.clone()),
|
||||
AxumPath((TEST_WORKSPACE_ID.to_string(), root.id.clone())),
|
||||
HeaderMap::new(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(outcome.requested_ticket, root.id);
|
||||
assert_eq!(
|
||||
outcome.queued_tickets,
|
||||
vec![dependency.id.clone(), root.id.clone()]
|
||||
);
|
||||
for ticket_id in [dependency.id, root.id] {
|
||||
let queued = backend.show(ticket_id.clone().into()).unwrap();
|
||||
assert_eq!(queued.meta.workflow_state, TicketWorkflowState::Queued);
|
||||
assert_eq!(
|
||||
queued
|
||||
.events
|
||||
.last()
|
||||
.and_then(|event| event.attributes.get("orchestrator_assignment_id"))
|
||||
.cloned(),
|
||||
Some(format!("orchestrator-{ticket_id}"))
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn queued_ticket_mutation_succeeds_without_orchestrator() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
@@ -17200,9 +17337,13 @@ mod tests {
|
||||
workspace_id: TEST_WORKSPACE_ID.to_string(),
|
||||
id: ticket_id.clone(),
|
||||
};
|
||||
let mut related_input = ticket::NewTicket::new("Related Browser Ticket");
|
||||
related_input.workflow_state = Some(TicketWorkflowState::Ready);
|
||||
related_input.repository_id = Some(TEST_REPOSITORY_ID.to_string());
|
||||
related_input.ref_selector = Some("develop".to_string());
|
||||
let related_ticket_id = browser_ticket_backend(&api)
|
||||
.unwrap()
|
||||
.create(ticket::NewTicket::new("Related Browser Ticket"))
|
||||
.create(related_input)
|
||||
.unwrap()
|
||||
.id;
|
||||
browser_ticket_backend(&api)
|
||||
@@ -17274,6 +17415,7 @@ mod tests {
|
||||
.unwrap();
|
||||
assert_eq!(ready.state, "ready");
|
||||
assign_test_orchestrator(&api, &ticket_id);
|
||||
assign_test_orchestrator(&api, &related_ticket_id);
|
||||
let Json(ready_detail) = scoped_get_ticket(State(api.clone()), AxumPath(path()))
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
Reference in New Issue
Block a user