use std::fs;
use serde_json::json;
use crate::domain::ticket::TicketSnapshot;
use crate::flow::{Check, Confidence, Flow};
use crate::protocol::{ErrorBody, Request, RequestId, ResponseEnvelope, VerdictArgs, VerdictValue};
use crate::run_store::{PanelReportRecord, RunRecord};
use crate::runner::WorkerScope;
use crate::vendor::VendorErrorMatch;
use crate::worker::{WorkerRole, check_label, definition_of_done};
use crate::work_state::local::TicketRecord;
use super::commands::{local_lookup, mark_storage_full, run_lookup};
use super::dispatcher::{DispatcherState, conflict, internal, invalid_arguments, unauthorized};
pub(super) fn dispatch_worker(
state: &mut DispatcherState,
id: RequestId,
request: Request,
run_id: &str,
token: Option<&str>,
) -> ResponseEnvelope {
let scope = token.and_then(|presented| {
state
.worker_tokens
.get(run_id)
.filter(|issued| issued.token == presented)
.map(|issued| issued.scope.clone())
});
let Some(scope) = scope else {
return ResponseEnvelope::failure(
Some(id),
unauthorized("the presented token is not valid for this run"),
);
};
let data = match request {
Request::Brief(_) => handle_brief(state, run_id, &scope),
Request::Show(args) => match args.reference.as_deref() {
Some(reference) if args.limit.is_none() => handle_show(state, run_id, reference),
_ => Err(unauthorized(
"workers may only show their own run's ticket by exact id",
)),
},
Request::Note(args) => handle_note(state, run_id, &args.text),
Request::Verdict(args) => match &scope {
WorkerScope::Stage { .. } => handle_verdict(state, run_id, &scope, &args),
WorkerScope::PanelReviewer {
stage_index,
reviewer_index,
..
} => handle_panel_report(
state,
run_id,
&scope,
PanelSeat {
stage_index: *stage_index,
reviewer_index: *reviewer_index,
},
&args,
),
},
_ => Err(unauthorized(
"operator verbs are not available on a worker socket",
)),
};
match data {
Ok(data) => ResponseEnvelope::success(Some(id), data),
Err(error) => ResponseEnvelope::failure(Some(id), error),
}
}
struct ExecutingStage {
name: String,
attempt: u32,
check: Check,
}
fn executing_stage(run: &RunRecord, scope: &WorkerScope) -> Result<ExecutingStage, ErrorBody> {
if !matches!(run.state.as_str(), "running" | "driving") {
return Err(conflict("the run has no stage currently executing"));
}
let snapshot = run
.flow_json
.as_deref()
.ok_or_else(|| internal("the run has no flow snapshot"))?;
let flow: Flow = serde_json::from_str(snapshot)
.map_err(|error| internal(&format!("the run's flow snapshot is invalid: {error}")))?;
let (name, attempt) = match scope {
WorkerScope::Stage { stage, attempt } => (stage.clone(), *attempt),
WorkerScope::PanelReviewer { stage, attempt, .. } => (stage.clone(), *attempt),
};
let stage = flow
.stages
.iter()
.find(|stage| stage.name == name)
.ok_or_else(|| internal("the executing stage is not in the run's flow snapshot"))?;
Ok(ExecutingStage {
name,
attempt,
check: stage.result_check.clone(),
})
}
fn handle_brief(
state: &DispatcherState,
run_id: &str,
scope: &WorkerScope,
) -> Result<serde_json::Value, ErrorBody> {
let run = run_lookup(state, |run_store| run_store.run(run_id))?
.ok_or_else(|| internal("the run for this token no longer exists"))?;
let ticket = match run.ticket_json.as_deref() {
Some(snapshot) => serde_json::from_str::<TicketSnapshot>(snapshot)
.map_err(|error| internal(&format!("the run's ticket snapshot is invalid: {error}")))?,
None => {
let ticket = local_lookup(state, |work_state| work_state.ticket(&run.ticket_id))?
.ok_or_else(|| internal("the ticket for this run no longer exists"))?;
let body = ticket.body.unwrap_or_else(|| {
ticket
.file_path
.as_ref()
.and_then(|file_path| fs::read_to_string(state.root.join(file_path)).ok())
.unwrap_or_default()
});
TicketSnapshot {
id: ticket.id,
name: ticket.name,
blocked_by: ticket.blocked_by,
worktree: ticket.worktree,
target: ticket.target,
model: ticket.model,
effort: ticket.effort,
body,
}
}
};
let executing = executing_stage(&run, scope)?;
let role = match scope {
WorkerScope::Stage { .. } => WorkerRole::Stage,
WorkerScope::PanelReviewer { .. } => WorkerRole::PanelReviewer,
};
let definition_of_done = definition_of_done(role, &executing.check);
Ok(json!({
"run": run_id,
"ticket": {
"id": ticket.id,
"name": ticket.name,
"blocked_by": ticket.blocked_by,
"worktree": ticket.worktree,
"body": ticket.body,
"target": ticket.target,
"model": ticket.model,
"effort": ticket.effort,
},
"worktree": run.worktree_path,
"branch": run.branch,
"stage": {
"name": executing.name,
"attempt": executing.attempt,
"result_check": check_label(&executing.check),
},
"definition_of_done": definition_of_done,
}))
}
fn handle_show(
state: &DispatcherState,
run_id: &str,
reference: &str,
) -> Result<serde_json::Value, ErrorBody> {
let run = run_lookup(state, |run_store| run_store.run(run_id))?
.ok_or_else(|| internal("the run for this token no longer exists"))?;
if reference != run.ticket_id {
return Err(unauthorized("workers may only show their own run's ticket"));
}
let ticket = local_lookup(state, |work_state| work_state.ticket(&run.ticket_id))?
.ok_or_else(|| internal("the ticket for this run no longer exists"))?;
let vendor_error = current_ticket_vendor_error(state, &ticket)?;
Ok(ticket_show(reference, &ticket, vendor_error.as_ref()))
}
pub(super) fn ticket_show(
reference: &str,
ticket: &TicketRecord,
vendor_error: Option<&VendorErrorMatch>,
) -> serde_json::Value {
json!({
"ref": reference,
"kind": "ticket",
"value": {
"id": ticket.id,
"project": ticket.project_id,
"state": ticket.state,
"file": ticket.file_path,
"name": ticket.name,
"blocked_by": ticket.blocked_by,
"worktree": ticket.worktree,
"target": ticket.target,
"model": ticket.model,
"effort": ticket.effort,
"reason": vendor_error.map(|error| error.diagnostic.as_str()),
"classification": vendor_error,
},
})
}
pub(super) fn current_ticket_vendor_error(
state: &DispatcherState,
ticket: &TicketRecord,
) -> Result<Option<VendorErrorMatch>, ErrorBody> {
let vendor_error = run_lookup(state, |run_store| {
run_store.latest_vendor_error_for_ticket(&ticket.id)
})?;
if ticket.state != "ready" {
return Ok(vendor_error);
}
let cooldown_active = match ticket.target.as_deref() {
Some(target) => run_lookup(state, |run_store| {
run_store.active_cooldown_for_target(target, state.clock.now_ms())
})?
.is_some(),
None => false,
};
Ok(vendor_error.filter(|error| error.class.requires_cooldown() && cooldown_active))
}
fn handle_note(
state: &DispatcherState,
run_id: &str,
text: &str,
) -> Result<serde_json::Value, ErrorBody> {
let ordinal = run_lookup(state, |run_store| run_store.next_note_ordinal())?;
let note_id = format!("N{ordinal}");
state
.run_store
.insert_note(¬e_id, run_id, text, state.clock.now_ms())
.map_err(|error| {
mark_storage_full(state, &error);
internal(&format!("cannot record note: {error}"))
})?;
Ok(json!({"note": {"id": note_id, "run": run_id, "text": text}}))
}
struct PanelSeat {
stage_index: usize,
reviewer_index: usize,
}
fn handle_panel_report(
state: &DispatcherState,
run_id: &str,
scope: &WorkerScope,
seat: PanelSeat,
args: &VerdictArgs,
) -> Result<serde_json::Value, ErrorBody> {
let run = run_lookup(state, |run_store| run_store.run(run_id))?
.ok_or_else(|| internal("the run for this token no longer exists"))?;
let executing = executing_stage(&run, scope)?;
let reason = args
.reason
.as_deref()
.map(str::trim)
.filter(|reason| !reason.is_empty())
.ok_or_else(|| invalid_arguments("a panel reviewer must report a non-empty `--reason`"))?;
let verdict = match args.verdict {
VerdictValue::Pass => "pass",
VerdictValue::Fail => "fail",
};
let confidence = args
.confidence
.map_or(Confidence::default(), Confidence::from);
let record = PanelReportRecord {
stage: &executing.name,
stage_index: seat.stage_index,
attempt: executing.attempt,
reviewer_index: seat.reviewer_index,
verdict,
confidence: confidence.as_str(),
reason,
};
let inserted = state
.run_store
.record_panel_report(run_id, &record, state.clock.now_ms())
.map_err(|error| {
mark_storage_full(state, &error);
internal(&format!("cannot record panel report: {error}"))
})?;
if !inserted {
return Err(conflict(&format!(
"reviewer {} of stage `{}` has already reported",
seat.reviewer_index, executing.name
)));
}
Ok(json!({
"verdict": {
"run": run_id,
"stage": executing.name,
"reviewer": seat.reviewer_index,
"verdict": verdict,
"confidence": confidence.as_str(),
"reason": reason,
}
}))
}
fn handle_verdict(
state: &DispatcherState,
run_id: &str,
scope: &WorkerScope,
args: &VerdictArgs,
) -> Result<serde_json::Value, ErrorBody> {
let run = run_lookup(state, |run_store| run_store.run(run_id))?
.ok_or_else(|| internal("the run for this token no longer exists"))?;
let executing = executing_stage(&run, scope)?;
let stage_name = executing.name;
let attempt = executing.attempt;
if executing.check != Check::Reported {
return Err(unauthorized(&format!(
"stage `{stage_name}` does not use `result_check: reported`"
)));
}
let verdict = match args.verdict {
VerdictValue::Pass => "pass",
VerdictValue::Fail => "fail",
};
let confidence = args
.confidence
.map_or(Confidence::default(), Confidence::from);
let inserted = state
.run_store
.record_stage_verdict(
run_id,
&stage_name,
attempt,
verdict,
confidence.as_str(),
args.reason.as_deref(),
state.clock.now_ms(),
)
.map_err(|error| {
mark_storage_full(state, &error);
internal(&format!("cannot record stage verdict: {error}"))
})?;
if !inserted {
return Err(conflict(&format!(
"stage `{stage_name}` has already reported a verdict"
)));
}
Ok(json!({
"verdict": {
"run": run_id,
"stage": stage_name,
"attempt": attempt,
"verdict": verdict,
"confidence": confidence.as_str(),
"reason": args.reason,
}
}))
}