Merge commit '0496cd907bc7bb96e9aa1c6d385bedb616bf3233' into work/00001M10FJVA2-orchestrator-queue-notice

This commit is contained in:
2026-08-27 12:46:27 +09:00
8 changed files with 187 additions and 65 deletions
+52 -20
View File
@@ -576,6 +576,7 @@ pub(crate) enum IntakeRegistryUpdate {
pub(crate) struct ReadyTicketPlanningReturnRequest {
workspace_root: PathBuf,
ticket_id: String,
ticket_key: String,
user_instruction: String,
followup: ReadyTicketPlanningReturnFollowup,
}
@@ -2031,11 +2032,18 @@ impl DashboardApp {
return None;
};
let ticket_id = ticket.id.clone();
let ticket_key = match required_ticket_handoff_key(ticket.resource_key.as_deref()) {
Ok(ticket_key) => ticket_key.to_string(),
Err(error) => {
self.notice = Some(error);
return None;
}
};
let mut context =
TicketRoleLaunchContext::new(current_workspace_root(), TicketRole::Intake);
context.ticket = Some(TicketRef::id(ticket_id.clone()));
context.user_instruction = Some(format!(
"Continue Intake for existing Ticket {ticket_id}. Do not create a duplicate Ticket unless the user explicitly requests one. Read ShowTicket body/thread/artifacts before making routing or requirements decisions."
"Continue Intake for existing Ticket {ticket_key}. Do not create a duplicate Ticket unless the user explicitly requests one. Read ShowTicket body/thread/artifacts before making routing or requirements decisions."
));
let store = match PanelRegistryStore::default_for_workspace(&context.workspace_root) {
Ok(store) => store,
@@ -2048,7 +2056,7 @@ impl DashboardApp {
Ok(Some(claim)) => {
let status = local_claim_status_for_pod(&claim.worker_name, &self.list);
self.notice = Some(existing_ticket_claim_notice(
&ticket_id,
&ticket_key,
&claim.worker_name,
status,
));
@@ -2076,7 +2084,7 @@ impl DashboardApp {
self.sending = true;
self.notice = Some(format!(
"Launching Ticket Intake for {} as {}",
ticket_id, planned.worker_name
ticket_key, planned.worker_name
));
Some(IntakeLaunchRequest {
context,
@@ -2147,10 +2155,17 @@ impl DashboardApp {
return None;
};
let ticket_id = ticket.id.clone();
let ticket_key = match required_ticket_handoff_key(ticket.resource_key.as_deref()) {
Ok(ticket_key) => ticket_key.to_string(),
Err(error) => {
self.notice = Some(error);
return None;
}
};
if ticket.workflow_state != TicketWorkflowState::Ready {
self.notice = Some(format!(
"Ticket {} is {}; expected ready before returning to planning.",
ticket_id,
ticket_key,
ticket.workflow_state.as_str()
));
return None;
@@ -2202,7 +2217,7 @@ impl DashboardApp {
TicketRoleLaunchContext::new(workspace_root.clone(), TicketRole::Intake);
context.ticket = Some(TicketRef::id(ticket_id.clone()));
context.user_instruction = Some(build_ready_ticket_refinement_launch_instruction(
&ticket_id,
&ticket_key,
&user_instruction,
));
let peer_registration = self.prepare_intake_peer_registration(&mut context);
@@ -2226,11 +2241,12 @@ impl DashboardApp {
self.sending = true;
self.notice = Some(format!(
"Returning ready Ticket {} to planning for refinement…",
ticket_id
ticket_key
));
Some(ReadyTicketPlanningReturnRequest {
workspace_root,
ticket_id,
ticket_key,
user_instruction,
followup,
})
@@ -3882,21 +3898,35 @@ fn bounded_refinement_instruction(input: &str) -> String {
.to_string()
}
fn build_ready_ticket_refinement_thread_body(ticket_id: &str, instruction: &str) -> String {
fn required_ticket_handoff_key(resource_key: Option<&str>) -> Result<&str, String> {
let resource_key = resource_key.ok_or_else(|| {
"Ticket handoff is unavailable because the canonical T-* resource key is missing. Refresh the panel and retry."
.to_string()
})?;
let sequence = resource_key.strip_prefix("T-").filter(|sequence| {
!sequence.is_empty() && sequence.bytes().all(|byte| byte.is_ascii_digit())
});
sequence.map(|_| resource_key).ok_or_else(|| {
"Ticket handoff is unavailable because the canonical T-* resource key is invalid. Refresh the panel and retry."
.to_string()
})
}
fn build_ready_ticket_refinement_thread_body(ticket_key: &str, instruction: &str) -> String {
format!(
"Panel returned ready Ticket {ticket_id} to planning for requirements sync. This is not Queue routing and must not start implementation.\n\n## User refinement instruction\n\n{instruction}\n"
"Panel returned ready Ticket {ticket_key} to planning for requirements sync. This is not Queue routing and must not start implementation.\n\n## User refinement instruction\n\n{instruction}\n"
)
}
fn build_ready_ticket_refinement_launch_instruction(ticket_id: &str, instruction: &str) -> String {
fn build_ready_ticket_refinement_launch_instruction(ticket_key: &str, instruction: &str) -> String {
format!(
"Continue Ticket Intake / requirements sync for existing Ticket {ticket_id}. The Panel has returned the Ticket from ready to planning; do not queue the Ticket, do not route implementation, and do not create a duplicate unless the user explicitly asks for one. Read ShowTicket body/thread/artifacts before making requirements or readiness decisions.\n\nUser refinement instruction:\n\n{instruction}"
"Continue Ticket Intake / requirements sync for existing Ticket {ticket_key}. The Panel has returned the Ticket from ready to planning; do not queue the Ticket, do not route implementation, and do not create a duplicate unless the user explicitly asks for one. Read ShowTicket body/thread/artifacts before making requirements or readiness decisions.\n\nUser refinement instruction:\n\n{instruction}"
)
}
fn build_ready_ticket_refinement_notify(ticket_id: &str, instruction: &str) -> String {
fn build_ready_ticket_refinement_notify(ticket_key: &str, instruction: &str) -> String {
format!(
"Ticket {ticket_id} was returned from ready to planning from the Panel for requirements sync. Continue Intake/refinement only; do not Queue or route implementation. Read the Ticket thread for the recorded state change and user instruction.\n\nUser refinement instruction:\n\n{instruction}"
"Ticket {ticket_key} was returned from ready to planning from the Panel for requirements sync. Continue Intake/refinement only; do not Queue or route implementation. Read the Ticket thread for the recorded state change and user instruction.\n\nUser refinement instruction:\n\n{instruction}"
)
}
@@ -3925,10 +3955,12 @@ async fn dispatch_ready_ticket_planning_return(
let ticket = backend
.show(id.clone())
.map_err(|error| TicketActionError::Ticket(error.to_string()))?;
let ticket_key =
required_ticket_handoff_key(Some(&request.ticket_key)).map_err(TicketActionError::Stale)?;
if ticket.meta.workflow_state != TicketWorkflowState::Ready {
return Err(TicketActionError::Stale(format!(
"Ticket {} is {}; expected ready before returning it to planning. Refresh the panel and retry if appropriate.",
ticket.meta.id,
ticket_key,
ticket.meta.workflow_state.as_str()
)));
}
@@ -3937,7 +3969,7 @@ async fn dispatch_ready_ticket_planning_return(
TicketWorkflowState::Planning.as_str(),
"panel_return_to_planning",
MarkdownText::from(build_ready_ticket_refinement_thread_body(
&ticket.meta.id,
ticket_key,
&request.user_instruction,
)),
);
@@ -3951,7 +3983,7 @@ async fn dispatch_ready_ticket_planning_return(
ReadyTicketPlanningReturnOutcome {
notice: format!(
"Ticket {} returned to planning for refinement; launching Ticket Intake…",
ticket.meta.id
ticket_key
),
followup: ReadyTicketPlanningReturnAfterMutation::LaunchIntake(request),
}
@@ -3961,19 +3993,19 @@ async fn dispatch_ready_ticket_planning_return(
socket_path,
} => {
let message =
build_ready_ticket_refinement_notify(&ticket.meta.id, &request.user_instruction);
build_ready_ticket_refinement_notify(ticket_key, &request.user_instruction);
match send_notify_only(&socket_path, message, true).await {
Ok(()) => ReadyTicketPlanningReturnOutcome {
notice: format!(
"Ticket {} returned to planning for refinement; notified live Intake Worker {}.",
ticket.meta.id, worker_name
ticket_key, worker_name
),
followup: ReadyTicketPlanningReturnAfterMutation::None,
},
Err(error) => ReadyTicketPlanningReturnOutcome {
notice: bounded_panel_diagnostic(format!(
"Ticket {} returned to planning and instruction was recorded, but notifying Intake Worker {} failed: {}",
ticket.meta.id, worker_name, error
ticket_key, worker_name, error
)),
followup: ReadyTicketPlanningReturnAfterMutation::None,
},
@@ -3984,7 +4016,7 @@ async fn dispatch_ready_ticket_planning_return(
ReadyTicketPlanningReturnOutcome {
notice: format!(
"Ticket {} returned to planning for refinement; opening/restoring claimed Intake Worker {}…",
ticket.meta.id, worker_name
ticket_key, worker_name
),
followup: ReadyTicketPlanningReturnAfterMutation::OpenClaim(request),
}
@@ -3993,7 +4025,7 @@ async fn dispatch_ready_ticket_planning_return(
ReadyTicketPlanningReturnOutcome {
notice: bounded_panel_diagnostic(format!(
"Ticket {} returned to planning and instruction was recorded, but Intake launch was not attempted because existing Intake claim {} is stale; inspect or clear the local claim before launching another Intake Worker.",
ticket.meta.id, worker_name
ticket_key, worker_name
)),
followup: ReadyTicketPlanningReturnAfterMutation::None,
}
+24
View File
@@ -390,6 +390,7 @@ fn planning_return_request(
ReadyTicketPlanningReturnRequest {
workspace_root: temp.path().to_path_buf(),
ticket_id,
ticket_key: "T-482".to_string(),
user_instruction: instruction.to_string(),
followup: ReadyTicketPlanningReturnFollowup::BlockedByStaleClaim {
worker_name: "stale-intake".to_string(),
@@ -494,6 +495,7 @@ fn ready_ticket_intake_enter_prepares_planning_return_not_queue_or_generic_launc
};
assert_eq!(request.ticket_id, "20260608-000123-ready");
assert_eq!(request.ticket_key, "T-1");
assert_eq!(request.user_instruction, "clarify expected behavior");
assert!(matches!(
request.followup,
@@ -515,6 +517,7 @@ async fn planning_return_with_launch_followup_changes_state_before_launch_follow
let request = ReadyTicketPlanningReturnRequest {
workspace_root: temp.path().to_path_buf(),
ticket_id: ticket_id.clone(),
ticket_key: "T-482".to_string(),
user_instruction: "launch intake after state change".to_string(),
followup: ReadyTicketPlanningReturnFollowup::LaunchIntake(IntakeLaunchRequest {
context: TicketRoleLaunchContext::new(temp.path().to_path_buf(), TicketRole::Intake),
@@ -3504,6 +3507,27 @@ fn ticket_action_error_records_f2_diagnostic_details() {
assert!(!app.panel_diagnostic_open);
}
#[test]
fn ready_ticket_refinement_projection_uses_only_canonical_resource_key() {
const INTERNAL_ID: &str = "00001KZVNXFNK";
let thread = build_ready_ticket_refinement_thread_body("T-482", "Clarify rollback.");
let launch = build_ready_ticket_refinement_launch_instruction("T-482", "Clarify rollback.");
let notify = build_ready_ticket_refinement_notify("T-482", "Clarify rollback.");
for projection in [&thread, &launch, &notify] {
assert!(projection.contains("T-482"));
assert!(!projection.contains(INTERNAL_ID));
}
}
#[test]
fn ticket_handoff_fails_closed_without_canonical_resource_key() {
assert_eq!(required_ticket_handoff_key(Some("T-482")), Ok("T-482"));
for invalid in [None, Some(""), Some("00001KZVNXFNK"), Some("T-key")] {
assert!(required_ticket_handoff_key(invalid).is_err());
}
}
fn plain_line(line: &Line<'_>) -> String {
line.spans
.iter()
@@ -89,18 +89,19 @@ impl Tool for SpawnTicketCoderTool {
let input: SpawnTicketCoderInput = serde_json::from_str(input_json).map_err(|error| {
ToolError::InvalidArgument(format!("invalid {TOOL_NAME} input: {error}"))
})?;
let ticket_id = authority_id(input.ticket_id, "ticket_id")?;
let workflow_state = self
let ticket_ref = authority_id(input.ticket_id, "ticket_id")?;
let ticket = self
.ticket_service
.workflow_state(&ticket_id)
.ticket_handoff(&ticket_ref)
.map_err(|error| ToolError::ExecutionFailed(error.to_string()))?;
if !matches!(
workflow_state,
ticket.workflow_state,
ticket::TicketWorkflowState::Queued | ticket::TicketWorkflowState::InProgress
) {
return Err(ToolError::ExecutionFailed(format!(
"Ticket {ticket_id} must be queued or inprogress before spawning its Coder; current state is {}",
workflow_state.as_str()
"Ticket {} must be queued or inprogress before spawning its Coder; current state is {}",
ticket.resource_key,
ticket.workflow_state.as_str()
)));
}
let call_id = non_empty(ctx.call_id, "tool call_id")?;
@@ -115,14 +116,14 @@ impl Tool for SpawnTicketCoderTool {
)?,
relative_cwd,
profile: CODER_PROFILE.to_string(),
ticket_id: Some(ticket_id.clone()),
operation_id: Some(format!("spawn-ticket-coder:{ticket_id}:{call_id}")),
display_name: format!("Coder · {ticket_id}"),
ticket_id: Some(ticket.id.clone()),
operation_id: Some(format!("spawn-ticket-coder:{}:{call_id}", ticket.id)),
display_name: format!("Coder · {}", ticket.resource_key),
initial_submit: vec![
Segment::Flow {
selector: CODER_FLOW.to_string(),
},
Segment::text(format!("Implement Ticket {ticket_id}.")),
Segment::text(format!("Implement Ticket {}.", ticket.resource_key)),
],
})
.await
@@ -134,7 +135,7 @@ impl Tool for SpawnTicketCoderTool {
)));
}
Ok(ToolOutput {
summary: format!("Spawned Coder for Ticket {ticket_id}"),
summary: format!("Spawned Coder for Ticket {}", ticket.resource_key),
content: Some(response.body),
attachments: Vec::new(),
})
@@ -201,21 +202,31 @@ mod tests {
use crate::worker::{WorkspaceClientError, WorkspaceResponse};
use super::*;
use crate::feature::builtin::ticket::TicketHandoff;
#[derive(Default)]
struct RecordingTicketService;
impl TicketService for RecordingTicketService {
fn workflow_state(&self, _ticket_id: &str) -> Result<TicketWorkflowState, TicketError> {
Ok(TicketWorkflowState::Queued)
fn ticket_handoff(&self, ticket_ref: &str) -> Result<TicketHandoff, TicketError> {
assert_eq!(ticket_ref, "T-482");
Ok(TicketHandoff {
id: "00001KZXN51C7".to_string(),
resource_key: "T-482".to_string(),
workflow_state: TicketWorkflowState::Queued,
})
}
}
struct FixedTicketService(TicketWorkflowState);
impl TicketService for FixedTicketService {
fn workflow_state(&self, _ticket_id: &str) -> Result<TicketWorkflowState, TicketError> {
Ok(self.0)
fn ticket_handoff(&self, _ticket_ref: &str) -> Result<TicketHandoff, TicketError> {
Ok(TicketHandoff {
id: "00001KZXN51C7".to_string(),
resource_key: "T-482".to_string(),
workflow_state: self.0,
})
}
}
@@ -247,7 +258,7 @@ mod tests {
};
tool.execute(
&serde_json::json!({
"ticket_id": "00001KZXN51C7",
"ticket_id": "T-482",
"runtime_id": "runtime-1",
"working_directory_id": "workdir-1"
})
@@ -265,16 +276,20 @@ mod tests {
request.operation_id.as_deref(),
Some("spawn-ticket-coder:00001KZXN51C7:call-7")
);
assert_eq!(request.display_name, "Coder · 00001KZXN51C7");
assert_eq!(request.display_name, "Coder · T-482");
assert_eq!(
request.initial_submit,
vec![
Segment::Flow {
selector: CODER_FLOW.to_string()
},
Segment::text("Implement Ticket 00001KZXN51C7.")
Segment::text("Implement Ticket T-482.")
]
);
assert!(!request.display_name.contains("00001KZXN51C7"));
assert!(request.initial_submit.iter().all(|segment| {
!Segment::flatten_to_text(std::slice::from_ref(segment)).contains("00001KZXN51C7")
}));
}
#[tokio::test]
+34 -5
View File
@@ -267,7 +267,20 @@ pub const TICKET_SERVICE_ID: &str = "ticket.authority";
const TICKET_SERVICE_VERSION: &str = "1";
pub trait TicketService: Send + Sync {
fn workflow_state(&self, ticket_id: &str) -> Result<TicketWorkflowState, TicketError>;
fn ticket_handoff(&self, ticket_ref: &str) -> Result<TicketHandoff, TicketError>;
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TicketHandoff {
pub id: String,
pub resource_key: String,
pub workflow_state: TicketWorkflowState,
}
fn is_canonical_ticket_resource_key(resource_key: &str) -> bool {
resource_key.strip_prefix("T-").is_some_and(|sequence| {
!sequence.is_empty() && sequence.bytes().all(|byte| byte.is_ascii_digit())
})
}
struct BackendTicketService {
@@ -275,10 +288,18 @@ struct BackendTicketService {
}
impl TicketService for BackendTicketService {
fn workflow_state(&self, ticket_id: &str) -> Result<TicketWorkflowState, TicketError> {
self.backend
.show(ticket_id.into())
.map(|ticket| ticket.meta.workflow_state)
fn ticket_handoff(&self, ticket_ref: &str) -> Result<TicketHandoff, TicketError> {
let ticket = self.backend.show(ticket_ref.into())?;
let resource_key = ticket
.meta
.resource_key
.filter(|key| is_canonical_ticket_resource_key(key))
.ok_or_else(|| TicketError::Conflict("ticket resource key is unavailable".into()))?;
Ok(TicketHandoff {
id: ticket.meta.id,
resource_key,
workflow_state: ticket.meta.workflow_state,
})
}
}
@@ -1770,6 +1791,14 @@ provider = "github"
assert_eq!(removed.target, "01TARGET");
}
#[test]
fn ticket_handoff_accepts_only_canonical_ticket_resource_keys() {
assert!(is_canonical_ticket_resource_key("T-482"));
for invalid in ["", "00001KZVNXFNK", "T-", "T-key", "O-482"] {
assert!(!is_canonical_ticket_resource_key(invalid));
}
}
#[test]
fn workspace_http_backend_executes_ticket_create_operation() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
+6 -3
View File
@@ -102,7 +102,6 @@ pub enum WorkerPrompt {
AgentsMdSection,
ResidentMemorySummarySection,
WorkerOrchestrationGuidanceSection,
TicketEventCompanionNotice,
SubWorkerSpawnToolDescription,
}
@@ -122,7 +121,6 @@ impl WorkerPrompt {
Self::WorkerOrchestrationGuidanceSection => {
"internal.worker_orchestration_guidance_section"
}
Self::TicketEventCompanionNotice => "worker.ticket_event_companion_notice",
Self::SubWorkerSpawnToolDescription => "internal.sub_worker_spawn_tool_description",
}
}
@@ -139,7 +137,6 @@ impl WorkerPrompt {
WorkerPrompt::AgentsMdSection,
WorkerPrompt::ResidentMemorySummarySection,
WorkerPrompt::WorkerOrchestrationGuidanceSection,
WorkerPrompt::TicketEventCompanionNotice,
WorkerPrompt::SubWorkerSpawnToolDescription,
];
}
@@ -686,6 +683,12 @@ mod tests {
fn builtin_dcdl_catalog_loads() {
let catalog = PromptCatalog::builtins_only().unwrap();
assert!(!catalog.projection.templates.is_empty());
assert!(
!catalog
.projection
.templates
.contains_key("worker.ticket_event_companion_notice")
);
}
#[test]
+38 -8
View File
@@ -6412,9 +6412,15 @@ fn worker_ticket_source_context(
}
}
fn ticket_notification_content(ticket_id: &str, current_state: &str) -> String {
fn canonical_ticket_resource_key(resource_key: &str) -> Option<&str> {
let sequence = resource_key.strip_prefix("T-")?;
(!sequence.is_empty() && sequence.bytes().all(|byte| byte.is_ascii_digit()))
.then_some(resource_key)
}
fn ticket_notification_content(resource_key: &str, current_state: &str) -> String {
format!(
"Ticket notification: ticket_id={ticket_id} current_state={current_state}. Reread the Ticket before acting."
"Ticket {resource_key} changed to {current_state}. Reread the current Ticket before acting."
)
}
@@ -6496,6 +6502,16 @@ fn notify_ticket_recipients(
current_state: &str,
source: Option<RuntimeWorkerRef>,
) {
let Ok(Some(resource_key)) =
api.store
.resource_key(workspace_id, WorkspaceResourceKind::Ticket, ticket_id)
else {
return;
};
let Some(resource_key) = canonical_ticket_resource_key(&resource_key) else {
return;
};
let mut recipients = Vec::new();
if let Some(assignment) = api
.store
@@ -6516,7 +6532,7 @@ fn notify_ticket_recipients(
recipients.sort();
recipients.dedup();
let content = ticket_notification_content(ticket_id, current_state);
let content = ticket_notification_content(resource_key, current_state);
for recipient in recipients {
if source.as_ref().is_some_and(|source| source == &recipient) {
continue;
@@ -18219,15 +18235,25 @@ mod tests {
}
#[test]
fn ticket_notification_projection_exposes_only_ticket_and_current_state() {
fn ticket_notification_requires_canonical_ticket_resource_key() {
assert_eq!(canonical_ticket_resource_key("T-429"), Some("T-429"));
for invalid in ["", "00001KZ9SR97B", "T-", "T-key", "O-429"] {
assert_eq!(canonical_ticket_resource_key(invalid), None);
}
}
#[test]
fn ticket_notification_projection_exposes_only_resource_key_and_current_state() {
const INTERNAL_ID: &str = "00001KZ9SR97B";
for current_state in ["queued", "inprogress"] {
let content = ticket_notification_content("00001KZ9SR97B", current_state);
let content = ticket_notification_content("T-429", current_state);
assert_eq!(
content,
format!(
"Ticket notification: ticket_id=00001KZ9SR97B current_state={current_state}. Reread the Ticket before acting."
"Ticket T-429 changed to {current_state}. Reread the current Ticket before acting."
)
);
assert!(!content.contains(INTERNAL_ID));
for forbidden in [
"workspace_id",
"event_sequence",
@@ -18405,9 +18431,13 @@ mod tests {
assert_eq!(inputs.len(), expected_states.len());
for ((recipient, content), current_state) in inputs.iter().zip(expected_states) {
assert_eq!(recipient.worker_id.to_string(), orchestrator.worker_id);
assert!(!content.contains(&ticket.id));
assert_eq!(
content,
&ticket_notification_content(&ticket.id, current_state)
&ticket_notification_content(
ticket.resource_key.as_deref().unwrap(),
current_state,
)
);
}
}
@@ -19523,7 +19553,7 @@ mod tests {
assert_eq!(
notifications[0].1,
ticket_notification_content(
ticket_ref.id.as_str(),
ticket_ref.resource_key.as_deref().unwrap(),
TicketWorkflowState::Queued.as_str()
)
);