From f31007f9887bdac0796ed250e2213a370eb5f10b Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 19 Apr 2026 20:12:23 -0400 Subject: [PATCH] feat(server): move batch run selector resolution to the server Resolve rm/archive/unarchive selectors through the server-owned runs/resolve endpoint instead of CLI-side summary matching, and move active-run delete force semantics into DELETE /runs/{id}. --- docs/api-reference/fabro-api.yaml | 19 +- .../fabro-cli/src/commands/runs/archive.rs | 16 +- lib/crates/fabro-cli/src/commands/runs/rm.rs | 48 ++--- lib/crates/fabro-cli/src/server_client.rs | 20 ++- lib/crates/fabro-cli/src/server_runs.rs | 65 +------ lib/crates/fabro-cli/tests/it/cmd/archive.rs | 78 ++++++++ lib/crates/fabro-cli/tests/it/cmd/rm.rs | 166 +++++++++++++----- .../fabro-cli/tests/it/cmd/unarchive.rs | 78 ++++++++ lib/crates/fabro-server/src/server.rs | 147 +++++++++++++++- .../fabro-api-client/src/api/runs-api.ts | 30 ++-- 10 files changed, 491 insertions(+), 176 deletions(-) diff --git a/docs/api-reference/fabro-api.yaml b/docs/api-reference/fabro-api.yaml index aa7358b0e..1b9cb3a1f 100644 --- a/docs/api-reference/fabro-api.yaml +++ b/docs/api-reference/fabro-api.yaml @@ -544,9 +544,10 @@ paths: operationId: deleteRun tags: [Runs] summary: Delete Run - description: Deletes durable store state for a run. This does not remove any local run directory. + description: Deletes durable store state for a run. This does not remove any local run directory. Active runs require `force=true`. parameters: - $ref: "#/components/parameters/RunId" + - $ref: "#/components/parameters/ForceRunDelete" responses: "204": description: Run deleted or already absent @@ -556,6 +557,12 @@ paths: application/json: schema: $ref: "#/components/schemas/ErrorResponse" + "409": + description: Run is active and requires `force=true` + content: + application/json: + schema: + $ref: "#/components/schemas/ErrorResponse" /api/v1/runs/{id}/cancel: post: @@ -2099,6 +2106,16 @@ components: default: false example: false + ForceRunDelete: + name: force + in: query + required: false + description: Whether to force deletion of an active run. Defaults to `false`. + schema: + type: boolean + default: false + example: false + ModelProviderFilter: name: provider in: query diff --git a/lib/crates/fabro-cli/src/commands/runs/archive.rs b/lib/crates/fabro-cli/src/commands/runs/archive.rs index 4ba328ef1..1673b858a 100644 --- a/lib/crates/fabro-cli/src/commands/runs/archive.rs +++ b/lib/crates/fabro-cli/src/commands/runs/archive.rs @@ -7,9 +7,6 @@ use super::short_run_id; use crate::args::{RunsArchiveArgs, RunsUnarchiveArgs}; use crate::command_context::CommandContext; use crate::server_client; -use crate::server_runs::{ - ServerRunSummaryInfo, ServerSummaryLookup, resolve_server_run_from_summaries, -}; use crate::shared::print_json_pretty; pub(crate) async fn archive_command( @@ -19,12 +16,10 @@ pub(crate) async fn archive_command( printer: Printer, ) -> Result<()> { let ctx = CommandContext::for_target(&args.server, printer, cli.clone(), cli_layer)?; - let lookup = ServerSummaryLookup::from_client(ctx.server().await?).await?; run_bulk( Action::Archive, &args.runs, - lookup.client(), - lookup.runs(), + ctx.server().await?.as_ref(), cli, printer, ) @@ -38,12 +33,10 @@ pub(crate) async fn unarchive_command( printer: Printer, ) -> Result<()> { let ctx = CommandContext::for_target(&args.server, printer, cli.clone(), cli_layer)?; - let lookup = ServerSummaryLookup::from_client(ctx.server().await?).await?; run_bulk( Action::Unarchive, &args.runs, - lookup.client(), - lookup.runs(), + ctx.server().await?.as_ref(), cli, printer, ) @@ -73,7 +66,6 @@ async fn run_bulk( action: Action, identifiers: &[String], client: &server_client::ServerStoreClient, - runs: &[ServerRunSummaryInfo], cli: &CliSettings, printer: Printer, ) -> Result<()> { @@ -83,7 +75,7 @@ async fn run_bulk( let mut errors = Vec::new(); for identifier in identifiers { - let run = match resolve_server_run_from_summaries(runs, identifier) { + let run = match client.resolve_run(identifier).await { Ok(run) => run, Err(err) => { if !json { @@ -98,7 +90,7 @@ async fn run_bulk( } }; - let run_id = run.run_id(); + let run_id = run.run_id; let result = match action { Action::Archive => client.archive_run(&run_id).await, Action::Unarchive => client.unarchive_run(&run_id).await, diff --git a/lib/crates/fabro-cli/src/commands/runs/rm.rs b/lib/crates/fabro-cli/src/commands/runs/rm.rs index 590e5bca0..8ea688a5d 100644 --- a/lib/crates/fabro-cli/src/commands/runs/rm.rs +++ b/lib/crates/fabro-cli/src/commands/runs/rm.rs @@ -1,4 +1,4 @@ -use anyhow::{Context, Result, bail}; +use anyhow::{Result, bail}; use fabro_types::settings::CliSettings; use fabro_types::settings::cli::{CliLayer, OutputFormat}; use fabro_util::printer::Printer; @@ -7,9 +7,6 @@ use super::short_run_id; use crate::args::RunsRemoveArgs; use crate::command_context::CommandContext; use crate::server_client; -use crate::server_runs::{ - ServerRunSummaryInfo, ServerSummaryLookup, resolve_server_run_from_summaries, -}; use crate::shared::print_json_pretty; pub(crate) async fn remove_command( @@ -19,14 +16,12 @@ pub(crate) async fn remove_command( printer: Printer, ) -> Result<()> { let ctx = CommandContext::for_target(&args.server, printer, cli.clone(), cli_layer)?; - let lookup = ServerSummaryLookup::from_client(ctx.server().await?).await?; - remove_from(args, lookup.client(), lookup.runs(), cli, printer).await + remove_from(args, ctx.server().await?.as_ref(), cli, printer).await } async fn remove_from( args: &RunsRemoveArgs, client: &server_client::ServerStoreClient, - runs: &[ServerRunSummaryInfo], cli: &CliSettings, printer: Printer, ) -> Result<()> { @@ -36,7 +31,7 @@ async fn remove_from( let mut errors = Vec::new(); for identifier in &args.runs { - let run = match resolve_server_run_from_summaries(runs, identifier) { + let run = match client.resolve_run(identifier).await { Ok(run) => run, Err(err) => { if !json { @@ -51,15 +46,15 @@ async fn remove_from( } }; - if run.status().is_active() && !args.force { - let run_id = run.run_id().to_string(); - let error = format!( - "cannot remove active run {} (status: {}, use -f to force)", - short_run_id(&run_id), - run.status() - ); + let run_id = run.run_id.to_string(); + if let Err(err) = delete_server_run(client, &run.run_id, args.force).await { + let error = err.to_string(); if !json { - fabro_util::printerr!(printer, "{error}"); + if error.starts_with("cannot remove active run ") { + fabro_util::printerr!(printer, "{error}"); + } else { + fabro_util::printerr!(printer, "error: {identifier}: {error}"); + } } errors.push(serde_json::json!({ "identifier": identifier, @@ -68,19 +63,6 @@ async fn remove_from( had_errors = true; continue; } - - let run_id = run.run_id().to_string(); - if let Err(err) = delete_server_run(client, &run).await { - if !json { - fabro_util::printerr!(printer, "error: {identifier}: {err}"); - } - errors.push(serde_json::json!({ - "identifier": identifier, - "error": err.to_string(), - })); - had_errors = true; - continue; - } removed.push(run_id.clone()); if !json { fabro_util::printerr!(printer, "{}", short_run_id(&run_id)); @@ -102,10 +84,8 @@ async fn remove_from( async fn delete_server_run( client: &server_client::ServerStoreClient, - run: &ServerRunSummaryInfo, + run_id: &fabro_types::RunId, + force: bool, ) -> Result<()> { - client - .delete_store_run(&run.run_id()) - .await - .with_context(|| format!("failed to delete store state for {}", run.run_id())) + client.delete_store_run(run_id, force).await } diff --git a/lib/crates/fabro-cli/src/server_client.rs b/lib/crates/fabro-cli/src/server_client.rs index ad583f7c9..a3da368a8 100644 --- a/lib/crates/fabro-cli/src/server_client.rs +++ b/lib/crates/fabro-cli/src/server_client.rs @@ -831,14 +831,18 @@ impl ServerStoreClient { } } - pub(crate) async fn delete_store_run(&self, run_id: &RunId) -> Result<()> { - self.client - .delete_run() - .id(run_id.to_string()) - .send() - .await - .map_err(map_api_error)?; - Ok(()) + pub(crate) async fn delete_store_run(&self, run_id: &RunId, force: bool) -> Result<()> { + let mut url = fabro_http::Url::parse(&self.base_url) + .with_context(|| format!("invalid server base URL {}", self.base_url))?; + url.path_segments_mut() + .map_err(|()| anyhow!("server base URL cannot accept path segments"))? + .extend(["api", "v1", "runs", &run_id.to_string()]); + if force { + url.query_pairs_mut().append_pair("force", "true"); + } + + let response = self.http_client.delete(url).send().await?; + ensure_raw_response_success(response).await } pub(crate) async fn list_run_artifacts( diff --git a/lib/crates/fabro-cli/src/server_runs.rs b/lib/crates/fabro-cli/src/server_runs.rs index 86f3bcde3..7a8b9dffa 100644 --- a/lib/crates/fabro-cli/src/server_runs.rs +++ b/lib/crates/fabro-cli/src/server_runs.rs @@ -1,7 +1,7 @@ use std::collections::HashMap; use std::sync::Arc; -use anyhow::{Result, bail}; +use anyhow::Result; use chrono::{DateTime, Utc}; use fabro_store::RunSummary; use fabro_types::{RunId, RunStatus, StatusReason}; @@ -103,13 +103,6 @@ impl ServerSummaryLookup { } } -pub(crate) fn resolve_server_run_from_summaries( - runs: &[ServerRunSummaryInfo], - selector: &str, -) -> Result { - resolve_server_run_from_infos(runs, selector) -} - pub(crate) fn filter_server_runs( runs: &[ServerRunSummaryInfo], before: Option<&str>, @@ -136,59 +129,3 @@ pub(crate) fn filter_server_runs( .cloned() .collect() } - -fn resolve_server_run_from_infos( - runs: &[ServerRunSummaryInfo], - identifier: &str, -) -> Result { - let id_matches: Vec<_> = runs - .iter() - .filter(|run| run_id_matches(run.run_id(), identifier)) - .collect(); - - match id_matches.len() { - 1 => return Ok(id_matches[0].clone()), - count if count > 1 => { - let ids: Vec = id_matches - .iter() - .map(|run| run.run_id().to_string()) - .collect(); - bail!( - "Ambiguous prefix '{identifier}': {count} runs match: {}", - ids.join(", ") - ); - } - _ => {} - } - - let id_lower = identifier.to_lowercase(); - let id_collapsed = collapse_separators(&id_lower); - let workflow_match = runs - .iter() - .filter(|run| { - if let Some(slug) = run.workflow_slug() { - if slug.to_lowercase() == id_lower { - return true; - } - } - let name_lower = run.workflow_name().to_lowercase(); - name_lower.contains(&id_lower) - || collapse_separators(&name_lower).contains(&id_collapsed) - }) - .max_by_key(|run| run.run_id().created_at()); - - match workflow_match { - Some(run) => Ok(run.clone()), - None => { - bail!("No run found matching '{identifier}' (tried run ID prefix and workflow name)") - } - } -} - -fn collapse_separators(s: &str) -> String { - s.chars().filter(|c| *c != '-' && *c != '_').collect() -} - -fn run_id_matches(run_id: RunId, prefix: &str) -> bool { - run_id.to_string().starts_with(prefix) -} diff --git a/lib/crates/fabro-cli/tests/it/cmd/archive.rs b/lib/crates/fabro-cli/tests/it/cmd/archive.rs index fee73fc38..06f276557 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/archive.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/archive.rs @@ -1,4 +1,5 @@ use fabro_test::{fabro_snapshot, test_context}; +use httpmock::MockServer; use serde_json::Value; use super::support::{setup_completed_fast_dry_run, setup_created_fast_dry_run}; @@ -198,3 +199,80 @@ fn archive_mixed_batch_aggregates_errors() { assert_eq!(errors.len(), 1); assert_eq!(errors[0]["identifier"], bad); } + +#[test] +fn archive_resolves_selector_via_server_endpoint() { + let context = test_context!(); + let server = MockServer::start(); + let run_id = unique_run_id(); + let resolve_mock = server.mock(|when, then| { + when.method("GET") + .path("/api/v1/runs/resolve") + .query_param("selector", "nightly-build"); + then.status(200) + .header("Content-Type", "application/json") + .body( + serde_json::json!({ + "run_id": run_id, + "workflow_name": "Nightly Build", + "workflow_slug": "nightly-build", + "goal": "Nightly run", + "title": "Nightly run", + "labels": {}, + "host_repo_path": null, + "repository": { "name": "unknown" }, + "start_time": "2026-04-05T12:00:00Z", + "created_at": "2026-04-05T12:00:00Z", + "status": "succeeded", + "status_reason": null, + "blocked_reason": null, + "pending_control": null, + "duration_ms": 123, + "elapsed_secs": 0, + "total_usd_micros": null + }) + .to_string(), + ); + }); + let archive_mock = server.mock(|when, then| { + when.method("POST") + .path(format!("/api/v1/runs/{run_id}/archive")); + then.status(200) + .header("Content-Type", "application/json") + .body( + serde_json::json!({ + "id": run_id, + "status": "archived", + "blocked_reason": null, + "error": null, + "queue_position": null, + "status_reason": null, + "pending_control": null, + "created_at": "2026-04-05T12:00:00Z" + }) + .to_string(), + ); + }); + context.write_home( + ".fabro/settings.toml", + format!( + "_version = 1\n\n[cli.target]\ntype = \"http\"\nurl = \"{}/api/v1\"\n", + server.base_url() + ), + ); + + let mut filters = context.filters(); + filters.push(ulid_filter()); + let mut cmd = context.command(); + cmd.args(["archive", "nightly-build"]); + fabro_snapshot!(filters, cmd, @" + success: true + exit_code: 0 + ----- stdout ----- + ----- stderr ----- + [ULID] + "); + + resolve_mock.assert(); + archive_mock.assert(); +} diff --git a/lib/crates/fabro-cli/tests/it/cmd/rm.rs b/lib/crates/fabro-cli/tests/it/cmd/rm.rs index 69e2129d6..e5830dd06 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/rm.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/rm.rs @@ -84,7 +84,7 @@ fn rm_rejects_submitted_run_without_force() { exit_code: 1 ----- stdout ----- ----- stderr ----- - cannot remove active run [ULID] (status: submitted, use -f to force) + cannot remove active run [ULID] (status: submitted, use force=true or --force to force) error: some runs could not be removed "); } @@ -152,35 +152,39 @@ fn rm_force_removes_active_run() { let context = test_context!(); let run_id = unique_run_id(); let server = MockServer::start(); - let list_mock = server.mock(|when, then| { - when.method("GET").path("/api/v1/runs"); + let resolve_mock = server.mock(|when, then| { + when.method("GET") + .path("/api/v1/runs/resolve") + .query_param("selector", &run_id); then.status(200) .header("Content-Type", "application/json") .body( serde_json::json!({ - "data": [{ - "run_id": run_id, - "workflow_name": "Active Workflow", - "workflow_slug": "active-workflow", - "goal": "Active goal", - "title": "Active goal", - "labels": {}, - "host_repo_path": null, - "repository": { "name": "unknown" }, - "start_time": "2026-04-05T12:00:00Z", - "created_at": "2026-04-05T12:00:00Z", - "status": "running", - "status_reason": null, - "duration_ms": 123, - "total_usd_micros": null - }], - "meta": { "has_more": false } + "run_id": run_id, + "workflow_name": "Active Workflow", + "workflow_slug": "active-workflow", + "goal": "Active goal", + "title": "Active goal", + "labels": {}, + "host_repo_path": null, + "repository": { "name": "unknown" }, + "start_time": "2026-04-05T12:00:00Z", + "created_at": "2026-04-05T12:00:00Z", + "status": "running", + "status_reason": null, + "blocked_reason": null, + "pending_control": null, + "duration_ms": 123, + "elapsed_secs": 0, + "total_usd_micros": null }) .to_string(), ); }); let delete_mock = server.mock(|when, then| { - when.method("DELETE").path(format!("/api/v1/runs/{run_id}")); + when.method("DELETE") + .path(format!("/api/v1/runs/{run_id}")) + .query_param("force", "true"); then.status(204); }); context.write_home( @@ -205,7 +209,86 @@ fn rm_force_removes_active_run() { ----- stderr ----- [ULID] "); - list_mock.assert(); + resolve_mock.assert(); + delete_mock.assert(); +} + +#[test] +fn rm_without_force_uses_resolve_then_surfaces_server_conflict() { + let context = test_context!(); + let run_id = unique_run_id(); + let server = MockServer::start(); + let resolve_mock = server.mock(|when, then| { + when.method("GET") + .path("/api/v1/runs/resolve") + .query_param("selector", &run_id); + then.status(200) + .header("Content-Type", "application/json") + .body( + serde_json::json!({ + "run_id": run_id, + "workflow_name": "Active Workflow", + "workflow_slug": "active-workflow", + "goal": "Active goal", + "title": "Active goal", + "labels": {}, + "host_repo_path": null, + "repository": { "name": "unknown" }, + "start_time": "2026-04-05T12:00:00Z", + "created_at": "2026-04-05T12:00:00Z", + "status": "running", + "status_reason": null, + "blocked_reason": null, + "pending_control": null, + "duration_ms": 123, + "elapsed_secs": 0, + "total_usd_micros": null + }) + .to_string(), + ); + }); + let delete_mock = server.mock(|when, then| { + when.method("DELETE").path(format!("/api/v1/runs/{run_id}")); + then.status(409) + .header("Content-Type", "application/json") + .body( + serde_json::json!({ + "errors": [{ + "status": "409", + "title": "Conflict", + "detail": format!( + "cannot remove active run {} (status: running, use force=true or --force to force)", + &run_id[..12], + ), + }] + }) + .to_string(), + ); + }); + context.write_home( + ".fabro/settings.toml", + format!( + "_version = 1\n\n[cli.target]\ntype = \"http\"\nurl = \"{}/api/v1\"\n", + server.base_url() + ), + ); + + let mut filters = context.filters(); + filters.push(( + r"\b[0-9A-HJKMNP-TV-Z]{12}\b".to_string(), + "[ULID]".to_string(), + )); + let mut cmd = context.command(); + cmd.args(["rm", &run_id]); + fabro_snapshot!(filters, cmd, @" + success: false + exit_code: 1 + ----- stdout ----- + ----- stderr ----- + cannot remove active run [ULID] (status: running, use force=true or --force to force) + error: some runs could not be removed + "); + resolve_mock.assert(); delete_mock.assert(); } @@ -269,29 +352,28 @@ fn rm_uses_configured_server_target_without_local_run_dir() { let context = test_context!(); let run_id = unique_run_id(); let server = MockServer::start(); - let list_mock = server.mock(|when, then| { - when.method("GET").path("/api/v1/runs"); + let resolve_mock = server.mock(|when, then| { + when.method("GET") + .path("/api/v1/runs/resolve") + .query_param("selector", &run_id); then.status(200) .header("Content-Type", "application/json") .body( serde_json::json!({ - "data": [{ - "run_id": run_id, - "workflow_name": "Remote Workflow", - "workflow_slug": "remote-workflow", - "goal": "Remote goal", - "title": "Remote goal", - "labels": {}, - "host_repo_path": null, - "repository": { "name": "unknown" }, - "start_time": "2026-04-05T12:00:00Z", - "created_at": "2026-04-05T12:00:00Z", - "status": "succeeded", - "status_reason": null, - "duration_ms": 123, - "total_usd_micros": null - }], - "meta": { "has_more": false } + "run_id": run_id, + "workflow_name": "Remote Workflow", + "workflow_slug": "remote-workflow", + "goal": "Remote goal", + "title": "Remote goal", + "labels": {}, + "host_repo_path": null, + "repository": { "name": "unknown" }, + "start_time": "2026-04-05T12:00:00Z", + "created_at": "2026-04-05T12:00:00Z", + "status": "succeeded", + "status_reason": null, + "duration_ms": 123, + "total_usd_micros": null }) .to_string(), ); @@ -320,6 +402,6 @@ fn rm_uses_configured_server_target_without_local_run_dir() { String::from_utf8_lossy(&output.stdout), String::from_utf8_lossy(&output.stderr) ); - list_mock.assert(); + resolve_mock.assert(); delete_mock.assert(); } diff --git a/lib/crates/fabro-cli/tests/it/cmd/unarchive.rs b/lib/crates/fabro-cli/tests/it/cmd/unarchive.rs index d65438bdc..0ee460afa 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/unarchive.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/unarchive.rs @@ -1,4 +1,5 @@ use fabro_test::{fabro_snapshot, test_context}; +use httpmock::MockServer; use serde_json::Value; use super::support::{setup_completed_fast_dry_run, setup_created_fast_dry_run}; @@ -204,3 +205,80 @@ fn unarchive_mixed_batch_aggregates_errors() { assert_eq!(errors.len(), 1); assert_eq!(errors[0]["identifier"], active_run.run_id); } + +#[test] +fn unarchive_resolves_selector_via_server_endpoint() { + let context = test_context!(); + let server = MockServer::start(); + let run_id = unique_run_id(); + let resolve_mock = server.mock(|when, then| { + when.method("GET") + .path("/api/v1/runs/resolve") + .query_param("selector", "nightly-build"); + then.status(200) + .header("Content-Type", "application/json") + .body( + serde_json::json!({ + "run_id": run_id, + "workflow_name": "Nightly Build", + "workflow_slug": "nightly-build", + "goal": "Nightly run", + "title": "Nightly run", + "labels": {}, + "host_repo_path": null, + "repository": { "name": "unknown" }, + "start_time": "2026-04-05T12:00:00Z", + "created_at": "2026-04-05T12:00:00Z", + "status": "archived", + "status_reason": null, + "blocked_reason": null, + "pending_control": null, + "duration_ms": 123, + "elapsed_secs": 0, + "total_usd_micros": null + }) + .to_string(), + ); + }); + let unarchive_mock = server.mock(|when, then| { + when.method("POST") + .path(format!("/api/v1/runs/{run_id}/unarchive")); + then.status(200) + .header("Content-Type", "application/json") + .body( + serde_json::json!({ + "id": run_id, + "status": "succeeded", + "blocked_reason": null, + "error": null, + "queue_position": null, + "status_reason": null, + "pending_control": null, + "created_at": "2026-04-05T12:00:00Z" + }) + .to_string(), + ); + }); + context.write_home( + ".fabro/settings.toml", + format!( + "_version = 1\n\n[cli.target]\ntype = \"http\"\nurl = \"{}/api/v1\"\n", + server.base_url() + ), + ); + + let mut filters = context.filters(); + filters.push(ulid_filter()); + let mut cmd = context.command(); + cmd.args(["unarchive", "nightly-build"]); + fabro_snapshot!(filters, cmd, @" + success: true + exit_code: 0 + ----- stdout ----- + ----- stderr ----- + [ULID] + "); + + resolve_mock.assert(); + unarchive_mock.assert(); +} diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index f0ba62234..d01a1d76f 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -115,8 +115,8 @@ use crate::ip_allowlist::{IpAllowlistConfig, ip_allowlist_middleware}; use crate::jwt_auth::{ AuthMode, AuthenticatedService, AuthenticatedSubject, authenticate_service_parts, }; -use crate::run_selector::{ResolveRunError, resolve_run_by_selector}; use crate::run_files::{FilesInFlight, list_run_files, new_files_in_flight}; +use crate::run_selector::{ResolveRunError, resolve_run_by_selector}; use crate::server_secrets::{ LlmClientResult, ProviderCredentials, ServerSecrets, auth_issue_message, }; @@ -1410,7 +1410,7 @@ async fn prune_runs( } for run_id in &prune_plan.run_ids { - if let Err(response) = delete_run_internal(&state, *run_id).await { + if let Err(response) = delete_run_internal(&state, *run_id, true).await { return response; } } @@ -2858,6 +2858,12 @@ struct ResolveRunQuery { selector: String, } +#[derive(Debug, Default, serde::Deserialize)] +struct DeleteRunQuery { + #[serde(default)] + force: bool, +} + async fn resolve_run( _auth: AuthenticatedService, State(state): State>, @@ -2900,6 +2906,7 @@ async fn resolve_run( async fn delete_run( _auth: AuthenticatedService, State(state): State>, + Query(query): Query, Path(id): Path, ) -> Response { let id = match parse_run_id_path(&id) { @@ -2907,13 +2914,21 @@ async fn delete_run( Err(response) => return response, }; - match delete_run_internal(&state, id).await { + match delete_run_internal(&state, id, query.force).await { Ok(()) => StatusCode::NO_CONTENT.into_response(), Err(response) => response, } } -async fn delete_run_internal(state: &Arc, id: RunId) -> Result<(), Response> { +async fn delete_run_internal( + state: &Arc, + id: RunId, + force: bool, +) -> Result<(), Response> { + if !force { + reject_active_delete_without_force(state.as_ref(), &id).await?; + } + let managed_run = if let Ok(mut runs) = state.runs.lock() { runs.remove(&id) } else { @@ -2977,6 +2992,55 @@ async fn delete_run_internal(state: &Arc, id: RunId) -> Result<(), Res Ok(()) } +async fn reject_active_delete_without_force( + state: &AppState, + run_id: &RunId, +) -> Result<(), Response> { + let managed_status = state + .runs + .lock() + .ok() + .and_then(|runs| runs.get(run_id).map(|managed_run| managed_run.status)); + if let Some(status) = managed_status { + if matches!( + status, + RunStatus::Submitted + | RunStatus::Queued + | RunStatus::Starting + | RunStatus::Running + | RunStatus::Blocked + | RunStatus::Paused + ) { + return Err(ApiError::new( + StatusCode::CONFLICT, + active_run_delete_message(*run_id, status), + ) + .into_response()); + } + return Ok(()); + } + + match state.store.runs().find(run_id).await { + Ok(Some(summary)) if summary.status.is_active() => Err(ApiError::new( + StatusCode::CONFLICT, + active_run_delete_message(*run_id, summary.status), + ) + .into_response()), + Ok(_) => Ok(()), + Err(err) => { + Err(ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response()) + } + } +} + +fn active_run_delete_message(run_id: RunId, status: impl std::fmt::Display) -> String { + let run_id = run_id.to_string(); + let short_run_id = &run_id[..12.min(run_id.len())]; + format!( + "cannot remove active run {short_run_id} (status: {status}, use force=true or --force to force)" + ) +} + async fn terminate_worker_for_deletion( worker_pid: Option, worker_pgid: Option, @@ -9378,9 +9442,84 @@ slug = "fabro" let req = Request::builder() .method("DELETE") + .uri(api(&format!("/runs/{run_id}?force=true"))) + .body(Body::empty()) + .unwrap(); + let response = app.clone().oneshot(req).await.unwrap(); + assert_eq!(response.status(), StatusCode::NO_CONTENT); + + let req = Request::builder() + .method("GET") .uri(api(&format!("/runs/{run_id}"))) .body(Body::empty()) .unwrap(); + let response = app.oneshot(req).await.unwrap(); + assert_eq!(response.status(), StatusCode::NOT_FOUND); + } + + #[tokio::test] + async fn delete_active_run_requires_force() { + let state = create_app_state(); + let app = build_router(Arc::clone(&state), AuthMode::Disabled); + + let req = Request::builder() + .method("POST") + .uri(api("/runs")) + .header("content-type", "application/json") + .body(manifest_body(MINIMAL_DOT)) + .unwrap(); + + let response = app.clone().oneshot(req).await.unwrap(); + let body = body_json(response.into_body()).await; + let run_id = body["id"].as_str().unwrap(); + + let req = Request::builder() + .method("DELETE") + .uri(api(&format!("/runs/{run_id}"))) + .body(Body::empty()) + .unwrap(); + let response = app.clone().oneshot(req).await.unwrap(); + assert_eq!(response.status(), StatusCode::CONFLICT); + let body = body_json(response.into_body()).await; + let short_run_id = &run_id[..12.min(run_id.len())]; + let expected = format!( + "cannot remove active run {short_run_id} (status: submitted, use force=true or --force to force)" + ); + assert_eq!( + body["errors"][0]["detail"].as_str(), + Some(expected.as_str()) + ); + + let req = Request::builder() + .method("GET") + .uri(api(&format!("/runs/{run_id}"))) + .body(Body::empty()) + .unwrap(); + let response = app.oneshot(req).await.unwrap(); + assert_eq!(response.status(), StatusCode::OK); + } + + #[tokio::test] + async fn delete_active_run_force_succeeds() { + let state = create_app_state(); + let app = build_router(Arc::clone(&state), AuthMode::Disabled); + + let req = Request::builder() + .method("POST") + .uri(api("/runs")) + .header("content-type", "application/json") + .body(manifest_body(MINIMAL_DOT)) + .unwrap(); + + let response = app.clone().oneshot(req).await.unwrap(); + let body = body_json(response.into_body()).await; + let run_id = body["id"].as_str().unwrap(); + + let req = Request::builder() + .method("DELETE") + .uri(api(&format!("/runs/{run_id}?force=true"))) + .body(Body::empty()) + .unwrap(); let response = app.clone().oneshot(req).await.unwrap(); assert_eq!(response.status(), StatusCode::NO_CONTENT); diff --git a/lib/packages/fabro-api-client/src/api/runs-api.ts b/lib/packages/fabro-api-client/src/api/runs-api.ts index b0d1a87bd..046eff252 100644 --- a/lib/packages/fabro-api-client/src/api/runs-api.ts +++ b/lib/packages/fabro-api-client/src/api/runs-api.ts @@ -166,13 +166,14 @@ export const RunsApiAxiosParamCreator = function (configuration?: Configuration) }; }, /** - * Deletes durable store state for a run. This does not remove any local run directory. + * Deletes durable store state for a run. This does not remove any local run directory. Active runs require `force=true`. * @summary Delete Run * @param {string} id Unique run identifier (ULID). + * @param {boolean} [force] Whether to force deletion of an active run. Defaults to `false`. * @param {*} [options] Override http request option. * @throws {RequiredError} */ - deleteRun: async (id: string, options: RawAxiosRequestConfig = {}): Promise => { + deleteRun: async (id: string, force?: boolean, options: RawAxiosRequestConfig = {}): Promise => { // verify required parameter 'id' is not null or undefined assertParamExists('deleteRun', 'id', id) const localVarPath = `/api/v1/runs/{id}` @@ -194,6 +195,10 @@ export const RunsApiAxiosParamCreator = function (configuration?: Configuration) // http bearer authentication required await setBearerAuthToObject(localVarHeaderParameter, configuration) + if (force !== undefined) { + localVarQueryParameter['force'] = force; + } + localVarHeaderParameter['Accept'] = 'application/json'; setSearchParams(localVarUrlObj, localVarQueryParameter); @@ -719,14 +724,15 @@ export const RunsApiFp = function(configuration?: Configuration) { return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, /** - * Deletes durable store state for a run. This does not remove any local run directory. + * Deletes durable store state for a run. This does not remove any local run directory. Active runs require `force=true`. * @summary Delete Run * @param {string} id Unique run identifier (ULID). + * @param {boolean} [force] Whether to force deletion of an active run. Defaults to `false`. * @param {*} [options] Override http request option. * @throws {RequiredError} */ - async deleteRun(id: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { - const localVarAxiosArgs = await localVarAxiosParamCreator.deleteRun(id, options); + async deleteRun(id: string, force?: boolean, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { + const localVarAxiosArgs = await localVarAxiosParamCreator.deleteRun(id, force, options); const localVarOperationServerIndex = configuration?.serverIndex ?? 0; const localVarOperationServerBasePath = operationServerMap['RunsApi.deleteRun']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); @@ -918,14 +924,15 @@ export const RunsApiFactory = function (configuration?: Configuration, basePath? return localVarFp.createRun(runManifest, options).then((request) => request(axios, basePath)); }, /** - * Deletes durable store state for a run. This does not remove any local run directory. + * Deletes durable store state for a run. This does not remove any local run directory. Active runs require `force=true`. * @summary Delete Run * @param {string} id Unique run identifier (ULID). + * @param {boolean} [force] Whether to force deletion of an active run. Defaults to `false`. * @param {*} [options] Override http request option. * @throws {RequiredError} */ - deleteRun(id: string, options?: RawAxiosRequestConfig): AxiosPromise { - return localVarFp.deleteRun(id, options).then((request) => request(axios, basePath)); + deleteRun(id: string, force?: boolean, options?: RawAxiosRequestConfig): AxiosPromise { + return localVarFp.deleteRun(id, force, options).then((request) => request(axios, basePath)); }, /** * Temporary board-view list of managed runs. This endpoint is UI-oriented and may change as the app evolves. @@ -1082,14 +1089,15 @@ export class RunsApi extends BaseAPI { } /** - * Deletes durable store state for a run. This does not remove any local run directory. + * Deletes durable store state for a run. This does not remove any local run directory. Active runs require `force=true`. * @summary Delete Run * @param {string} id Unique run identifier (ULID). + * @param {boolean} [force] Whether to force deletion of an active run. Defaults to `false`. * @param {*} [options] Override http request option. * @throws {RequiredError} */ - public deleteRun(id: string, options?: RawAxiosRequestConfig) { - return RunsApiFp(this.configuration).deleteRun(id, options).then((request) => request(this.axios, this.basePath)); + public deleteRun(id: string, force?: boolean, options?: RawAxiosRequestConfig) { + return RunsApiFp(this.configuration).deleteRun(id, force, options).then((request) => request(this.axios, this.basePath)); } /**