mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-01 02:04:24 +00:00
Every run executes on Petri, so the in-process legacy executor goes: `fabro-core` and, in `fabro-workflow`, the handlers, lifecycle, pipeline execution, routing, retry, conditions, node handlers, steering, agent memory, artifacts, checkpoints, command log, and the `start`, `resume`, `retry`, `fork`, `rewind` and `timeline` operations. The two are deleted together because the engine half of `fabro-workflow` was the only user of `fabro-core` and `fabro-core` the only runtime of that half; neither compiles without the other. Kept in `fabro-workflow`, narrowed: the parse/transform/validate/persist pipeline and `create`, `archive`, `validate` (workflow definitions still come from DOT and settings); the run tools (`run_tools`, moved from `handler/llm/fabro_tools.rs`) for Ask Fabro, `fabro exec` and Petri's host tools; the pull request pipeline (`pull_request`, moved from `pipeline/`, for the step 0 port); Run Files' diff helpers in `sandbox_git`; `git_identity`, `usage_rollup`, `run_status`, `run_materialization`, `web_search` and `workflow_bundle`. Server: `RegistryFactoryOverride` becomes `execute_in_process`; `RunAnswerTransport::InProcess` carries only the interviewer; the interrupt endpoint answers 501 `interrupt_unsupported` and every pair endpoint 501 `pair_unsupported` (status lists none); rewind, fork, retry and timeline handlers and routes are removed; the command log is served from the stage output blob; usage rollups accumulate from the settled projection after an in-process run as after a worker exit. Ported while here: - `materialize_admitted_run` materializes the goal and drops a disabled pull request block, as the legacy materializer did. - A run whose admitted graph has an agent or prompt node is refused at create when no LLM provider is ready (`fabro.model.no_ready_provider`); a workflow of commands and gates needs no model and is admitted. - The projection's question type falls back on the options, as the interview adapter does, so a gate with edge-label options answers as multiple choice. Tests: the server scenarios (lifecycle, run completion, SSE, helpers) run in process on Petri and assert Petri's stage labels and stream names; the reconcile tests assert Petri's relaunch semantics; legacy unit tests of the deleted executor are removed; three server unit tests the removal took with it are restored; the pair fixtures go with the pair feature. Petri test fixtures no longer name `[workflow] engine`. Still red after this commit, all legacy consumers the next steps delete or port: fabro-store's Slate/reducer fixtures and fabro-types legacy JSON tests (step 4); server unit tests over legacy run events (retry endpoints, list_run_events, artifacts, per-event pause/unpause, run history activation, legacy sandbox fixtures) (steps 3-4); CLI tests that parse legacy event envelopes, the legacy `events`/`attach`/`diff`/ `dump`/`inspect` snapshots, `run rewind`/`run fork`, the ACP and git-identity workflow tests, and the runner tests that drive the legacy worker by hand (steps 3-4); the web app's Petri fixtures still carry `engine` (regenerate with `FABRO_CAPTURE_PETRI_FIXTURES` in step 4). Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
1118 lines
38 KiB
Rust
1118 lines
38 KiB
Rust
use std::collections::HashSet;
|
|
use std::sync::Arc;
|
|
|
|
use chrono::Utc;
|
|
use tokio::time::{Instant, sleep_until};
|
|
|
|
use super::super::{
|
|
ApiError, AppState, AskFabroReadiness, BatchDeleteRunsRequest, BatchDeleteRunsResponse,
|
|
BatchDeleteRunsResult, BatchDeleteRunsResultOutcome, BatchDeleteRunsSummary,
|
|
BatchRunLifecycleRequest, BatchRunLifecycleResponse, BatchRunLifecycleResult,
|
|
BatchRunLifecycleResultOutcome, BatchRunLifecycleSummary, DeleteRunOutcome, DeleteRunSandbox,
|
|
DenyRunRequest, FailureReason, IntoResponse, Json, Path, PendingReason, Principal,
|
|
RequireRunManagementTarget, RequiredUser, Response, Router, RunAnswerTransport,
|
|
RunControlAction, RunExecutionMode, RunId, RunRunnableSource, RunStatus, StartRunRequest,
|
|
State, StatusCode, Storage, WORKER_CANCEL_GRACE, WorkflowError, append_control_request,
|
|
clear_live_run_state, delete_run_internal, durable_run_status, load_pending_control,
|
|
managed_run, operations, parse_run_id_path, persist_cancelled_run_status, post,
|
|
reject_if_archived, update_live_run_from_event, workflow_event,
|
|
};
|
|
use crate::worker_runtime::WorkerRef;
|
|
|
|
pub(super) fn routes() -> Router<Arc<AppState>> {
|
|
Router::new()
|
|
.route("/runs/{id}/cancel", post(cancel_run))
|
|
.route("/runs/{id}/start", post(start_run))
|
|
.route("/runs/{id}/approve", post(approve_run))
|
|
.route("/runs/{id}/deny", post(deny_run))
|
|
.route("/runs/{id}/pause", post(pause_run))
|
|
.route("/runs/{id}/unpause", post(unpause_run))
|
|
.route("/runs/archive", post(batch_archive_runs))
|
|
.route("/runs/delete", post(batch_delete_runs))
|
|
.route("/runs/unarchive", post(batch_unarchive_runs))
|
|
.route("/runs/{id}/archive", post(archive_run))
|
|
.route("/runs/{id}/unarchive", post(unarchive_run))
|
|
}
|
|
|
|
async fn run_response(state: &AppState, id: RunId, status: StatusCode) -> Response {
|
|
match state.stores.run_summaries.get(&id, Utc::now()).await {
|
|
Ok(Some(summary)) => {
|
|
(status, Json(state.decorate_run_summary(summary).await)).into_response()
|
|
}
|
|
Ok(None) => ApiError::not_found("Run not found.").into_response(),
|
|
Err(err) => {
|
|
ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response()
|
|
}
|
|
}
|
|
}
|
|
|
|
async fn start_run(
|
|
RequireRunManagementTarget(id, actor): RequireRunManagementTarget,
|
|
State(state): State<Arc<AppState>>,
|
|
body: Option<Json<StartRunRequest>>,
|
|
) -> Response {
|
|
if let Some(response) = reject_if_archived(state.as_ref(), &id).await {
|
|
return response;
|
|
}
|
|
let resume = body.is_some_and(|Json(req)| req.resume);
|
|
|
|
match queue_run_start(state.as_ref(), id, resume, actor).await {
|
|
Ok(()) => run_response(state.as_ref(), id, StatusCode::OK).await,
|
|
Err(err) => err.into_response(),
|
|
}
|
|
}
|
|
|
|
pub(in crate::server) async fn queue_run_start(
|
|
state: &AppState,
|
|
id: RunId,
|
|
resume: bool,
|
|
actor: Principal,
|
|
) -> Result<(), ApiError> {
|
|
{
|
|
let runs = state.runs.lock().expect("runs lock poisoned");
|
|
if let Some(managed_run) = runs.get(&id) {
|
|
if matches!(
|
|
managed_run.status,
|
|
RunStatus::Pending { .. }
|
|
| RunStatus::Runnable
|
|
| RunStatus::Starting
|
|
| RunStatus::Running
|
|
| RunStatus::Blocked { .. }
|
|
| RunStatus::Paused { .. }
|
|
) {
|
|
return Err(ApiError::new(
|
|
StatusCode::CONFLICT,
|
|
if resume {
|
|
"an engine process is still running for this run — cannot resume"
|
|
} else if matches!(
|
|
managed_run.status,
|
|
RunStatus::Pending { .. } | RunStatus::Runnable
|
|
) {
|
|
"start has already been requested for this run"
|
|
} else {
|
|
"an engine process is still running for this run — cannot start"
|
|
},
|
|
));
|
|
}
|
|
}
|
|
}
|
|
|
|
let Ok(run_store) = state.stores.runs.open_run(&id).await else {
|
|
return Err(ApiError::not_found("Run not found."));
|
|
};
|
|
let run_state = match run_store.state().await {
|
|
Ok(state) => state,
|
|
Err(err) => {
|
|
return Err(ApiError::new(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
format!("Failed to load run state: {err}"),
|
|
));
|
|
}
|
|
};
|
|
|
|
if resume {
|
|
if run_state.current_checkpoint().is_none() {
|
|
return Err(ApiError::new(
|
|
StatusCode::CONFLICT,
|
|
"no checkpoint to resume from",
|
|
));
|
|
}
|
|
} else {
|
|
let status = run_state.status;
|
|
if !matches!(status, RunStatus::Submitted) {
|
|
return Err(ApiError::new(
|
|
StatusCode::CONFLICT,
|
|
format!("cannot start run: status is {status}, expected submitted"),
|
|
));
|
|
}
|
|
}
|
|
|
|
let run_dir = Storage::new(state.server_storage_dir())
|
|
.run_scratch(&id)
|
|
.root()
|
|
.to_path_buf();
|
|
let dot_source = run_state.spec.graph_source.clone().unwrap_or_default();
|
|
let approval_required = !resume
|
|
&& matches!(
|
|
&actor,
|
|
Principal::Worker { run_id } if run_state.parent_id == Some(*run_id)
|
|
);
|
|
if let Err(err) =
|
|
workflow_event::append_event(&run_store, &id, &workflow_event::Event::RunStartRequested {
|
|
resume,
|
|
actor: Some(actor.clone()),
|
|
})
|
|
.await
|
|
{
|
|
return Err(ApiError::new(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
err.to_string(),
|
|
));
|
|
}
|
|
let (next_status, next_event) = if approval_required {
|
|
(
|
|
RunStatus::Pending {
|
|
reason: PendingReason::ApprovalRequired,
|
|
},
|
|
workflow_event::Event::RunPending {
|
|
reason: PendingReason::ApprovalRequired,
|
|
actor: Some(actor),
|
|
},
|
|
)
|
|
} else {
|
|
(RunStatus::Runnable, workflow_event::Event::RunRunnable {
|
|
source: RunRunnableSource::StartRequested,
|
|
actor: Some(actor),
|
|
})
|
|
};
|
|
if let Err(err) = workflow_event::append_event(&run_store, &id, &next_event).await {
|
|
return Err(ApiError::new(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
err.to_string(),
|
|
));
|
|
}
|
|
|
|
{
|
|
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
|
runs.insert(
|
|
id,
|
|
managed_run(
|
|
dot_source,
|
|
next_status,
|
|
id.created_at(),
|
|
run_dir,
|
|
if resume {
|
|
RunExecutionMode::Resume
|
|
} else {
|
|
RunExecutionMode::Start
|
|
},
|
|
),
|
|
);
|
|
}
|
|
|
|
if !approval_required {
|
|
state.scheduler_notify.notify_one();
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
async fn approve_run(
|
|
RequiredUser(user): RequiredUser,
|
|
Path(id): Path<String>,
|
|
State(state): State<Arc<AppState>>,
|
|
) -> Response {
|
|
let id = match parse_run_id_path(&id) {
|
|
Ok(id) => id,
|
|
Err(response) => return response,
|
|
};
|
|
if let Some(response) = reject_if_archived(state.as_ref(), &id).await {
|
|
return response;
|
|
}
|
|
let Ok(run_store) = state.stores.runs.open_run(&id).await else {
|
|
return ApiError::not_found("Run not found.").into_response();
|
|
};
|
|
let run_state = match run_store.state().await {
|
|
Ok(state) => state,
|
|
Err(err) => {
|
|
return ApiError::new(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
format!("Failed to load run state: {err}"),
|
|
)
|
|
.into_response();
|
|
}
|
|
};
|
|
if !matches!(run_state.status, RunStatus::Pending {
|
|
reason: PendingReason::ApprovalRequired,
|
|
}) {
|
|
return ApiError::new(StatusCode::CONFLICT, "Run is not pending approval.").into_response();
|
|
}
|
|
|
|
let actor = Some(Principal::User(user));
|
|
for event in [
|
|
workflow_event::Event::RunApproved {
|
|
actor: actor.clone(),
|
|
},
|
|
workflow_event::Event::RunRunnable {
|
|
source: RunRunnableSource::Approved,
|
|
actor,
|
|
},
|
|
] {
|
|
if let Err(err) = workflow_event::append_event(&run_store, &id, &event).await {
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
|
.into_response();
|
|
}
|
|
}
|
|
|
|
{
|
|
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
|
if let Some(managed_run) = runs.get_mut(&id) {
|
|
managed_run.status = RunStatus::Runnable;
|
|
} else {
|
|
let run_dir = Storage::new(state.server_storage_dir())
|
|
.run_scratch(&id)
|
|
.root()
|
|
.to_path_buf();
|
|
let dot_source = run_state.spec.graph_source.clone().unwrap_or_default();
|
|
runs.insert(
|
|
id,
|
|
managed_run(
|
|
dot_source,
|
|
RunStatus::Runnable,
|
|
id.created_at(),
|
|
run_dir,
|
|
RunExecutionMode::Start,
|
|
),
|
|
);
|
|
}
|
|
}
|
|
|
|
state.scheduler_notify.notify_one();
|
|
run_response(state.as_ref(), id, StatusCode::OK).await
|
|
}
|
|
|
|
async fn deny_run(
|
|
RequiredUser(user): RequiredUser,
|
|
Path(id): Path<String>,
|
|
State(state): State<Arc<AppState>>,
|
|
body: Option<Json<DenyRunRequest>>,
|
|
) -> Response {
|
|
let id = match parse_run_id_path(&id) {
|
|
Ok(id) => id,
|
|
Err(response) => return response,
|
|
};
|
|
if let Some(response) = reject_if_archived(state.as_ref(), &id).await {
|
|
return response;
|
|
}
|
|
let reason = body
|
|
.and_then(|Json(req)| req.reason)
|
|
.map(|reason| reason.trim().to_string())
|
|
.filter(|reason| !reason.is_empty());
|
|
let message = reason
|
|
.clone()
|
|
.unwrap_or_else(|| "Not approved for execution".to_string());
|
|
let Ok(run_store) = state.stores.runs.open_run(&id).await else {
|
|
return ApiError::not_found("Run not found.").into_response();
|
|
};
|
|
let run_state = match run_store.state().await {
|
|
Ok(state) => state,
|
|
Err(err) => {
|
|
return ApiError::new(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
format!("Failed to load run state: {err}"),
|
|
)
|
|
.into_response();
|
|
}
|
|
};
|
|
if !matches!(run_state.status, RunStatus::Pending {
|
|
reason: PendingReason::ApprovalRequired,
|
|
}) {
|
|
return ApiError::new(StatusCode::CONFLICT, "Run is not pending approval.").into_response();
|
|
}
|
|
|
|
let actor = Some(Principal::User(user));
|
|
let denied_event = workflow_event::Event::RunDenied {
|
|
reason: reason.clone(),
|
|
actor,
|
|
};
|
|
if let Err(err) = workflow_event::append_event(&run_store, &id, &denied_event).await {
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response();
|
|
}
|
|
let failure_event = workflow_event::Event::workflow_run_failed_from_error(
|
|
&WorkflowError::engine(message.clone()),
|
|
fabro_types::RunTiming::default(),
|
|
FailureReason::ApprovalDenied,
|
|
None,
|
|
None,
|
|
None,
|
|
None,
|
|
);
|
|
if let Err(err) = workflow_event::append_event(&run_store, &id, &failure_event).await {
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response();
|
|
}
|
|
|
|
{
|
|
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
|
if let Some(managed_run) = runs.get_mut(&id) {
|
|
managed_run.status = RunStatus::Failed {
|
|
reason: FailureReason::ApprovalDenied,
|
|
};
|
|
managed_run.error = Some(message);
|
|
clear_live_run_state(managed_run);
|
|
}
|
|
}
|
|
|
|
run_response(state.as_ref(), id, StatusCode::OK).await
|
|
}
|
|
|
|
fn schedule_worker_cancel_escalation(state: Arc<AppState>, run_id: RunId, worker_ref: WorkerRef) {
|
|
let requested_at = Instant::now();
|
|
let armed = {
|
|
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
|
let Some(run) = runs.get_mut(&run_id) else {
|
|
return;
|
|
};
|
|
if run.cancel_escalation_worker.as_ref() == Some(&worker_ref) {
|
|
false
|
|
} else {
|
|
run.cancel_escalation_worker = Some(worker_ref.clone());
|
|
true
|
|
}
|
|
};
|
|
if !armed {
|
|
tracing::debug!(
|
|
run_id = %run_id,
|
|
worker_kind = worker_ref.kind(),
|
|
worker_ref = ?worker_ref,
|
|
"Worker cancellation escalation is already armed"
|
|
);
|
|
return;
|
|
}
|
|
|
|
tokio::spawn(async move {
|
|
sleep_until(requested_at + WORKER_CANCEL_GRACE).await;
|
|
let should_escalate = {
|
|
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
|
let Some(run) = runs.get_mut(&run_id) else {
|
|
return;
|
|
};
|
|
run.escalation_still_current(&worker_ref)
|
|
};
|
|
if !should_escalate {
|
|
tracing::debug!(
|
|
run_id = %run_id,
|
|
worker_kind = worker_ref.kind(),
|
|
worker_ref = ?worker_ref,
|
|
"Skipping stale worker cancellation escalation"
|
|
);
|
|
return;
|
|
}
|
|
if !state.worker_runtime.is_alive(&worker_ref).await {
|
|
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
|
if let Some(run) = runs.get_mut(&run_id) {
|
|
run.clear_escalation_for(&worker_ref);
|
|
}
|
|
tracing::debug!(
|
|
run_id = %run_id,
|
|
worker_kind = worker_ref.kind(),
|
|
worker_ref = ?worker_ref,
|
|
"Skipping worker cancellation escalation because worker exited"
|
|
);
|
|
return;
|
|
}
|
|
let still_current = {
|
|
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
|
let Some(run) = runs.get_mut(&run_id) else {
|
|
return;
|
|
};
|
|
run.escalation_still_current(&worker_ref)
|
|
};
|
|
if !still_current {
|
|
tracing::debug!(
|
|
run_id = %run_id,
|
|
worker_kind = worker_ref.kind(),
|
|
worker_ref = ?worker_ref,
|
|
"Skipping worker cancellation escalation after liveness check"
|
|
);
|
|
return;
|
|
}
|
|
|
|
let elapsed_ms = u64::try_from(requested_at.elapsed().as_millis()).unwrap_or(u64::MAX);
|
|
tracing::warn!(
|
|
run_id = %run_id,
|
|
worker_kind = worker_ref.kind(),
|
|
worker_ref = ?worker_ref,
|
|
elapsed_ms,
|
|
"Force-stopping worker after cancellation grace period"
|
|
);
|
|
state.worker_runtime.force_stop(&worker_ref).await;
|
|
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
|
if let Some(run) = runs.get_mut(&run_id) {
|
|
run.clear_escalation_for(&worker_ref);
|
|
}
|
|
});
|
|
}
|
|
|
|
async fn cancel_run(
|
|
RequireRunManagementTarget(id, actor): RequireRunManagementTarget,
|
|
State(state): State<Arc<AppState>>,
|
|
) -> Response {
|
|
if let Some(response) = reject_if_archived(state.as_ref(), &id).await {
|
|
return response;
|
|
}
|
|
let durable_summary = match state.stores.run_summaries.get(&id, Utc::now()).await {
|
|
Ok(summary) => summary,
|
|
Err(err) => {
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
|
.into_response();
|
|
}
|
|
};
|
|
let pending_control = durable_summary
|
|
.as_ref()
|
|
.and_then(|summary| summary.lifecycle.pending_control);
|
|
let durable_status = durable_summary
|
|
.as_ref()
|
|
.map(|summary| summary.lifecycle.status);
|
|
let cancel_target = {
|
|
let mut runs = state.runs.lock().expect("runs lock poisoned");
|
|
match runs.get_mut(&id) {
|
|
Some(managed_run) => {
|
|
let managed_status = managed_run.status;
|
|
match managed_status {
|
|
RunStatus::Submitted
|
|
| RunStatus::Pending { .. }
|
|
| RunStatus::Runnable
|
|
| RunStatus::Starting
|
|
| RunStatus::Running
|
|
| RunStatus::Blocked { .. }
|
|
| RunStatus::Paused { .. } => {
|
|
let answer_transport = managed_run.answer_transport.clone();
|
|
let should_cancel_pending_interview =
|
|
matches!(
|
|
&answer_transport,
|
|
Some(RunAnswerTransport::InProcess { .. })
|
|
) && (matches!(managed_status, RunStatus::Blocked { .. })
|
|
|| matches!(durable_status, Some(RunStatus::Blocked { .. })));
|
|
let persist_cancelled_status = matches!(
|
|
managed_status,
|
|
RunStatus::Submitted | RunStatus::Pending { .. } | RunStatus::Runnable
|
|
) && !should_cancel_pending_interview;
|
|
if persist_cancelled_status {
|
|
managed_run.status = RunStatus::Failed {
|
|
reason: FailureReason::Cancelled,
|
|
};
|
|
}
|
|
let cancel_tx = if should_cancel_pending_interview {
|
|
None
|
|
} else {
|
|
managed_run.cancel_tx.take()
|
|
};
|
|
Some((
|
|
persist_cancelled_status,
|
|
answer_transport,
|
|
managed_run.cancel_token.clone(),
|
|
cancel_tx,
|
|
managed_run.worker_ref.clone(),
|
|
))
|
|
}
|
|
_ => {
|
|
return ApiError::new(StatusCode::CONFLICT, "Run is not cancellable.")
|
|
.into_response();
|
|
}
|
|
}
|
|
}
|
|
None => None,
|
|
}
|
|
};
|
|
let Some((persist_cancelled_status, answer_transport, cancel_token, cancel_tx, worker_ref)) =
|
|
cancel_target
|
|
else {
|
|
return unmanaged_cancel_response(state.as_ref(), id, actor, pending_control).await;
|
|
};
|
|
|
|
if pending_control != Some(RunControlAction::Cancel) {
|
|
if let Err(err) =
|
|
append_control_request(state.as_ref(), id, RunControlAction::Cancel, Some(actor)).await
|
|
{
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
|
.into_response();
|
|
}
|
|
}
|
|
|
|
if let Some(token) = &cancel_token {
|
|
token.cancel();
|
|
}
|
|
let sent_in_process_cancel = if let Some(cancel_tx) = cancel_tx {
|
|
let _ = cancel_tx.send(());
|
|
true
|
|
} else {
|
|
false
|
|
};
|
|
let delivered_control = if let Some(answer_transport) = answer_transport {
|
|
if sent_in_process_cancel
|
|
&& matches!(answer_transport, RunAnswerTransport::InProcess { .. })
|
|
{
|
|
true
|
|
} else {
|
|
answer_transport.cancel_run().await.is_ok()
|
|
}
|
|
} else {
|
|
false
|
|
};
|
|
tracing::debug!(
|
|
run_id = %id,
|
|
delivered_control,
|
|
"Processed cooperative run cancellation signal"
|
|
);
|
|
if let Some(worker_ref) = worker_ref {
|
|
if !delivered_control {
|
|
state.worker_runtime.request_stop(&worker_ref).await;
|
|
}
|
|
schedule_worker_cancel_escalation(Arc::clone(&state), id, worker_ref);
|
|
}
|
|
|
|
let response_status = if persist_cancelled_status {
|
|
if let Err(err) = persist_cancelled_run_status(state.as_ref(), id).await {
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
|
.into_response();
|
|
}
|
|
StatusCode::OK
|
|
} else {
|
|
StatusCode::ACCEPTED
|
|
};
|
|
|
|
run_response(state.as_ref(), id, response_status).await
|
|
}
|
|
|
|
async fn unmanaged_cancel_response(
|
|
state: &AppState,
|
|
id: RunId,
|
|
actor: Principal,
|
|
pending_control: Option<RunControlAction>,
|
|
) -> Response {
|
|
match durable_run_status(state, id).await {
|
|
Ok(Some(status)) if status.is_terminal() => ApiError::new(
|
|
StatusCode::CONFLICT,
|
|
"Run is already terminal and cannot be cancelled.",
|
|
)
|
|
.into_response(),
|
|
Ok(Some(RunStatus::Submitted | RunStatus::Pending { .. } | RunStatus::Runnable)) => {
|
|
if pending_control != Some(RunControlAction::Cancel) {
|
|
if let Err(err) =
|
|
append_control_request(state, id, RunControlAction::Cancel, Some(actor)).await
|
|
{
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
|
.into_response();
|
|
}
|
|
}
|
|
match persist_cancelled_run_status(state, id).await {
|
|
Ok(()) => run_response(state, id, StatusCode::OK).await,
|
|
Err(err) => ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
|
.into_response(),
|
|
}
|
|
}
|
|
Ok(Some(_)) => {
|
|
ApiError::new(StatusCode::CONFLICT, "Run is not cancellable.").into_response()
|
|
}
|
|
Ok(None) => ApiError::not_found("Run not found.").into_response(),
|
|
Err(err) => {
|
|
ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response()
|
|
}
|
|
}
|
|
}
|
|
|
|
/// How `pause_run` should enact the transition, chosen from the current run
|
|
/// status.
|
|
enum PauseMode {
|
|
/// Worker is running; ask it to pause via the worker control bus. Status
|
|
/// flips to `Paused` once the worker acknowledges.
|
|
Transport { transport: RunAnswerTransport },
|
|
/// Worker is blocked on a human gate; flip to `Paused` directly by
|
|
/// appending `RunPaused` ourselves.
|
|
AppendEvent,
|
|
}
|
|
|
|
/// How `unpause_run` should enact the transition.
|
|
enum UnpauseMode {
|
|
/// No outstanding block; ask the worker to resume via the worker control
|
|
/// bus.
|
|
Transport { transport: RunAnswerTransport },
|
|
/// Was paused while blocked; append `RunUnpaused` and let the reducer
|
|
/// restore the underlying blocked state from `Paused { prior_block }`.
|
|
AppendEvent,
|
|
}
|
|
|
|
async fn pause_run(
|
|
subject: RequiredUser,
|
|
State(state): State<Arc<AppState>>,
|
|
Path(id): Path<String>,
|
|
) -> Response {
|
|
let id = match parse_run_id_path(&id) {
|
|
Ok(id) => id,
|
|
Err(response) => return response,
|
|
};
|
|
if let Some(response) = reject_if_archived(state.as_ref(), &id).await {
|
|
return response;
|
|
}
|
|
let pending_control = match load_pending_control(state.as_ref(), id).await {
|
|
Ok(pending_control) => pending_control,
|
|
Err(err) => {
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
|
.into_response();
|
|
}
|
|
};
|
|
let mode = {
|
|
let runs = state.runs.lock().expect("runs lock poisoned");
|
|
match runs.get(&id) {
|
|
Some(managed_run) if managed_run.status == RunStatus::Running => {
|
|
let Some(transport) = managed_run.answer_transport.clone() else {
|
|
return ApiError::new(StatusCode::CONFLICT, "Run worker is not available.")
|
|
.into_response();
|
|
};
|
|
PauseMode::Transport { transport }
|
|
}
|
|
Some(managed_run) if matches!(managed_run.status, RunStatus::Blocked { .. }) => {
|
|
PauseMode::AppendEvent
|
|
}
|
|
Some(_) => {
|
|
return ApiError::new(StatusCode::CONFLICT, "Run is not pausable.").into_response();
|
|
}
|
|
None => return ApiError::not_found("Run not found.").into_response(),
|
|
}
|
|
};
|
|
|
|
if pending_control.is_some() {
|
|
return ApiError::new(
|
|
StatusCode::CONFLICT,
|
|
"Run control request is already pending.",
|
|
)
|
|
.into_response();
|
|
}
|
|
if let Err(err) = append_control_request(
|
|
state.as_ref(),
|
|
id,
|
|
RunControlAction::Pause,
|
|
Some(Principal::User(subject.0.clone())),
|
|
)
|
|
.await
|
|
{
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response();
|
|
}
|
|
match mode {
|
|
PauseMode::Transport { transport } => {
|
|
if transport.pause_run().await.is_err() {
|
|
return ApiError::new(
|
|
StatusCode::SERVICE_UNAVAILABLE,
|
|
"Failed to deliver pause request to the active run.",
|
|
)
|
|
.into_response();
|
|
}
|
|
}
|
|
PauseMode::AppendEvent => {
|
|
if let Some(response) = synchronous_transition(state.as_ref(), id, |events| {
|
|
events.push(workflow_event::Event::RunPaused);
|
|
})
|
|
.await
|
|
{
|
|
return response;
|
|
}
|
|
}
|
|
}
|
|
|
|
run_response(state.as_ref(), id, StatusCode::OK).await
|
|
}
|
|
|
|
async fn unpause_run(
|
|
subject: RequiredUser,
|
|
State(state): State<Arc<AppState>>,
|
|
Path(id): Path<String>,
|
|
) -> Response {
|
|
let id = match parse_run_id_path(&id) {
|
|
Ok(id) => id,
|
|
Err(response) => return response,
|
|
};
|
|
if let Some(response) = reject_if_archived(state.as_ref(), &id).await {
|
|
return response;
|
|
}
|
|
let pending_control = match load_pending_control(state.as_ref(), id).await {
|
|
Ok(pending_control) => pending_control,
|
|
Err(err) => {
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
|
.into_response();
|
|
}
|
|
};
|
|
let mode = {
|
|
let runs = state.runs.lock().expect("runs lock poisoned");
|
|
match runs.get(&id) {
|
|
Some(managed_run) => match managed_run.status {
|
|
RunStatus::Paused {
|
|
prior_block: Some(_),
|
|
} => UnpauseMode::AppendEvent,
|
|
RunStatus::Paused { prior_block: None } => {
|
|
let Some(transport) = managed_run.answer_transport.clone() else {
|
|
return ApiError::new(StatusCode::CONFLICT, "Run worker is not available.")
|
|
.into_response();
|
|
};
|
|
UnpauseMode::Transport { transport }
|
|
}
|
|
_ => {
|
|
return ApiError::new(StatusCode::CONFLICT, "Run is not paused.")
|
|
.into_response();
|
|
}
|
|
},
|
|
None => return ApiError::not_found("Run not found.").into_response(),
|
|
}
|
|
};
|
|
|
|
if pending_control.is_some() {
|
|
return ApiError::new(
|
|
StatusCode::CONFLICT,
|
|
"Run control request is already pending.",
|
|
)
|
|
.into_response();
|
|
}
|
|
if let Err(err) = append_control_request(
|
|
state.as_ref(),
|
|
id,
|
|
RunControlAction::Unpause,
|
|
Some(Principal::User(subject.0.clone())),
|
|
)
|
|
.await
|
|
{
|
|
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response();
|
|
}
|
|
match mode {
|
|
UnpauseMode::Transport { transport } => {
|
|
if transport.unpause_run().await.is_err() {
|
|
return ApiError::new(
|
|
StatusCode::SERVICE_UNAVAILABLE,
|
|
"Failed to deliver unpause request to the active run.",
|
|
)
|
|
.into_response();
|
|
}
|
|
}
|
|
UnpauseMode::AppendEvent => {
|
|
if let Some(response) = synchronous_transition(state.as_ref(), id, |events| {
|
|
events.push(workflow_event::Event::RunUnpaused);
|
|
})
|
|
.await
|
|
{
|
|
return response;
|
|
}
|
|
}
|
|
}
|
|
|
|
run_response(state.as_ref(), id, StatusCode::OK).await
|
|
}
|
|
|
|
async fn archive_run(
|
|
RequireRunManagementTarget(id, actor): RequireRunManagementTarget,
|
|
State(state): State<Arc<AppState>>,
|
|
) -> Response {
|
|
run_archive_action(state, actor, id, ArchiveAction::Archive).await
|
|
}
|
|
|
|
async fn unarchive_run(
|
|
RequireRunManagementTarget(id, actor): RequireRunManagementTarget,
|
|
State(state): State<Arc<AppState>>,
|
|
) -> Response {
|
|
run_archive_action(state, actor, id, ArchiveAction::Unarchive).await
|
|
}
|
|
|
|
async fn batch_archive_runs(
|
|
RequiredUser(user): RequiredUser,
|
|
State(state): State<Arc<AppState>>,
|
|
Json(request): Json<BatchRunLifecycleRequest>,
|
|
) -> Response {
|
|
Box::pin(batch_run_archive_action(
|
|
state,
|
|
Principal::User(user),
|
|
request,
|
|
ArchiveAction::Archive,
|
|
))
|
|
.await
|
|
}
|
|
|
|
async fn batch_unarchive_runs(
|
|
RequiredUser(user): RequiredUser,
|
|
State(state): State<Arc<AppState>>,
|
|
Json(request): Json<BatchRunLifecycleRequest>,
|
|
) -> Response {
|
|
Box::pin(batch_run_archive_action(
|
|
state,
|
|
Principal::User(user),
|
|
request,
|
|
ArchiveAction::Unarchive,
|
|
))
|
|
.await
|
|
}
|
|
|
|
async fn batch_delete_runs(
|
|
_auth: RequiredUser,
|
|
State(state): State<Arc<AppState>>,
|
|
Json(request): Json<BatchDeleteRunsRequest>,
|
|
) -> Response {
|
|
let force = request.force;
|
|
let ids = match validate_batch_run_ids(request.run_ids) {
|
|
Ok(ids) => ids,
|
|
Err(err) => return err.into_response(),
|
|
};
|
|
|
|
let mut results = Vec::with_capacity(ids.len());
|
|
for id in ids {
|
|
results.push(batch_delete_run_item(state.as_ref(), id, force).await);
|
|
}
|
|
|
|
let requested = results.len() as u64;
|
|
let succeeded = results.iter().filter(|result| result.ok).count() as u64;
|
|
(
|
|
StatusCode::OK,
|
|
Json(BatchDeleteRunsResponse {
|
|
results,
|
|
summary: BatchDeleteRunsSummary {
|
|
requested,
|
|
succeeded,
|
|
failed: requested - succeeded,
|
|
},
|
|
}),
|
|
)
|
|
.into_response()
|
|
}
|
|
|
|
#[derive(Clone, Copy)]
|
|
enum ArchiveAction {
|
|
Archive,
|
|
Unarchive,
|
|
}
|
|
|
|
const MAX_BATCH_RUN_IDS: usize = 250;
|
|
|
|
async fn batch_run_archive_action(
|
|
state: Arc<AppState>,
|
|
actor: Principal,
|
|
request: BatchRunLifecycleRequest,
|
|
action: ArchiveAction,
|
|
) -> Response {
|
|
let ids = match validate_batch_run_ids(request.run_ids) {
|
|
Ok(ids) => ids,
|
|
Err(err) => return err.into_response(),
|
|
};
|
|
|
|
// Resolve Ask Fabro readiness once per batch instead of inside each
|
|
// per-item summary lookup; readiness is identical for every run in the
|
|
// request and resolving it performs LLM credential work.
|
|
let readiness = state.ask_fabro_readiness().await;
|
|
let mut results = Vec::with_capacity(ids.len());
|
|
for id in ids {
|
|
results.push(
|
|
batch_run_archive_item(state.as_ref(), &readiness, actor.clone(), id, action).await,
|
|
);
|
|
}
|
|
|
|
let requested = results.len() as u64;
|
|
let succeeded = results.iter().filter(|result| result.ok).count() as u64;
|
|
(
|
|
StatusCode::OK,
|
|
Json(BatchRunLifecycleResponse {
|
|
results,
|
|
summary: BatchRunLifecycleSummary {
|
|
requested,
|
|
succeeded,
|
|
failed: requested - succeeded,
|
|
},
|
|
}),
|
|
)
|
|
.into_response()
|
|
}
|
|
|
|
fn validate_batch_run_ids(run_ids: Vec<String>) -> Result<Vec<RunId>, ApiError> {
|
|
if run_ids.is_empty() {
|
|
return Err(ApiError::bad_request(
|
|
"run_ids must contain at least one run ID.",
|
|
));
|
|
}
|
|
if run_ids.len() > MAX_BATCH_RUN_IDS {
|
|
return Err(ApiError::bad_request(format!(
|
|
"run_ids must contain no more than {MAX_BATCH_RUN_IDS} run IDs.",
|
|
)));
|
|
}
|
|
|
|
let mut seen = HashSet::with_capacity(run_ids.len());
|
|
let mut ids = Vec::with_capacity(run_ids.len());
|
|
for raw in run_ids {
|
|
let id = raw.parse::<RunId>().map_err(|_| {
|
|
ApiError::bad_request(format!("run_ids contains invalid run ID: {raw}"))
|
|
})?;
|
|
if !seen.insert(id) {
|
|
return Err(ApiError::bad_request(
|
|
"run_ids must not contain duplicate IDs.",
|
|
));
|
|
}
|
|
ids.push(id);
|
|
}
|
|
Ok(ids)
|
|
}
|
|
|
|
async fn batch_delete_run_item(state: &AppState, id: RunId, force: bool) -> BatchDeleteRunsResult {
|
|
match delete_run_internal(state, id, force).await {
|
|
Ok(DeleteRunOutcome::Deleted) => {
|
|
batch_delete_success(id, BatchDeleteRunsResultOutcome::Deleted, None)
|
|
}
|
|
Ok(DeleteRunOutcome::AlreadyAbsent) => {
|
|
batch_delete_success(id, BatchDeleteRunsResultOutcome::AlreadyAbsent, None)
|
|
}
|
|
Ok(DeleteRunOutcome::Preserved(response)) => batch_delete_success(
|
|
id,
|
|
BatchDeleteRunsResultOutcome::SandboxPreserved,
|
|
Some(response.sandbox),
|
|
),
|
|
Err(error) => {
|
|
let outcome = match error.status() {
|
|
StatusCode::CONFLICT => BatchDeleteRunsResultOutcome::Conflict,
|
|
_ => BatchDeleteRunsResultOutcome::Error,
|
|
};
|
|
BatchDeleteRunsResult {
|
|
run_id: id.to_string(),
|
|
ok: false,
|
|
outcome,
|
|
sandbox: None,
|
|
error: Some(error.into_response_entry()),
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
fn batch_delete_success(
|
|
id: RunId,
|
|
outcome: BatchDeleteRunsResultOutcome,
|
|
sandbox: Option<DeleteRunSandbox>,
|
|
) -> BatchDeleteRunsResult {
|
|
BatchDeleteRunsResult {
|
|
run_id: id.to_string(),
|
|
ok: true,
|
|
outcome,
|
|
sandbox,
|
|
error: None,
|
|
}
|
|
}
|
|
|
|
async fn batch_run_archive_item(
|
|
state: &AppState,
|
|
readiness: &AskFabroReadiness,
|
|
actor: Principal,
|
|
id: RunId,
|
|
action: ArchiveAction,
|
|
) -> BatchRunLifecycleResult {
|
|
let outcome = match run_archive_operation(state, &id, Some(actor), action).await {
|
|
Ok(outcome) => outcome,
|
|
Err(err) => {
|
|
let api_error = archive_workflow_error_to_api_error(err);
|
|
let result_outcome = match api_error.status() {
|
|
StatusCode::NOT_FOUND => BatchRunLifecycleResultOutcome::NotFound,
|
|
StatusCode::CONFLICT => BatchRunLifecycleResultOutcome::Conflict,
|
|
_ => BatchRunLifecycleResultOutcome::Error,
|
|
};
|
|
return batch_result_failure(id, result_outcome, api_error);
|
|
}
|
|
};
|
|
|
|
match state.stores.run_summaries.get(&id, Utc::now()).await {
|
|
Ok(Some(summary)) => BatchRunLifecycleResult {
|
|
run_id: id.to_string(),
|
|
ok: true,
|
|
outcome,
|
|
run: Some(readiness.decorate(summary)),
|
|
error: None,
|
|
},
|
|
Ok(None) => batch_result_failure(
|
|
id,
|
|
BatchRunLifecycleResultOutcome::Error,
|
|
ApiError::new(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
"Failed to load run summary after lifecycle action.",
|
|
),
|
|
),
|
|
Err(err) => batch_result_failure(
|
|
id,
|
|
BatchRunLifecycleResultOutcome::Error,
|
|
ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()),
|
|
),
|
|
}
|
|
}
|
|
|
|
fn batch_result_failure(
|
|
id: RunId,
|
|
outcome: BatchRunLifecycleResultOutcome,
|
|
error: ApiError,
|
|
) -> BatchRunLifecycleResult {
|
|
BatchRunLifecycleResult {
|
|
run_id: id.to_string(),
|
|
ok: false,
|
|
outcome,
|
|
run: None,
|
|
error: Some(error.into_response_entry()),
|
|
}
|
|
}
|
|
|
|
async fn run_archive_operation(
|
|
state: &AppState,
|
|
id: &RunId,
|
|
actor: Option<Principal>,
|
|
action: ArchiveAction,
|
|
) -> Result<BatchRunLifecycleResultOutcome, WorkflowError> {
|
|
match action {
|
|
ArchiveAction::Archive => operations::archive(&state.stores.runs, id, actor)
|
|
.await
|
|
.map(|outcome| match outcome {
|
|
operations::ArchiveOutcome::Archived { .. } => {
|
|
BatchRunLifecycleResultOutcome::Archived
|
|
}
|
|
operations::ArchiveOutcome::AlreadyArchived => {
|
|
BatchRunLifecycleResultOutcome::AlreadyArchived
|
|
}
|
|
}),
|
|
ArchiveAction::Unarchive => operations::unarchive(&state.stores.runs, id, actor)
|
|
.await
|
|
.map(|outcome| match outcome {
|
|
operations::UnarchiveOutcome::Unarchived { .. } => {
|
|
BatchRunLifecycleResultOutcome::Unarchived
|
|
}
|
|
operations::UnarchiveOutcome::NotArchived { .. } => {
|
|
BatchRunLifecycleResultOutcome::NotArchived
|
|
}
|
|
}),
|
|
}
|
|
}
|
|
|
|
fn archive_workflow_error_to_api_error(err: WorkflowError) -> ApiError {
|
|
match err {
|
|
WorkflowError::Precondition(message) => ApiError::new(StatusCode::CONFLICT, message),
|
|
WorkflowError::RunNotFound(_) => ApiError::not_found("Run not found."),
|
|
err => ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()),
|
|
}
|
|
}
|
|
|
|
async fn run_archive_action(
|
|
state: Arc<AppState>,
|
|
actor: Principal,
|
|
id: RunId,
|
|
action: ArchiveAction,
|
|
) -> Response {
|
|
match run_archive_operation(state.as_ref(), &id, Some(actor), action).await {
|
|
Ok(_) => archive_status_response(state.as_ref(), id).await,
|
|
Err(err) => archive_workflow_error_to_api_error(err).into_response(),
|
|
}
|
|
}
|
|
|
|
async fn archive_status_response(state: &AppState, id: RunId) -> Response {
|
|
run_response(state, id, StatusCode::OK).await
|
|
}
|
|
|
|
/// Persist a synchronous pause/unpause transition: append the caller-supplied
|
|
/// events to the run store and mirror the new status in the in-memory run map.
|
|
/// Returns `Some(Response)` on error, `None` on success.
|
|
async fn synchronous_transition(
|
|
state: &AppState,
|
|
id: RunId,
|
|
append_events: impl FnOnce(&mut Vec<workflow_event::Event>),
|
|
) -> Option<Response> {
|
|
let run_store = match state.stores.runs.open_run(&id).await {
|
|
Ok(run_store) => run_store,
|
|
Err(err) => {
|
|
return Some(
|
|
ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response(),
|
|
);
|
|
}
|
|
};
|
|
let mut events = Vec::new();
|
|
append_events(&mut events);
|
|
for event in events {
|
|
if let Err(err) = workflow_event::append_event(&run_store, &id, &event).await {
|
|
return Some(
|
|
ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response(),
|
|
);
|
|
}
|
|
let stored = workflow_event::to_run_event(&id, &event);
|
|
update_live_run_from_event(state, id, &stored);
|
|
}
|
|
None
|
|
}
|