diff --git a/lib/crates/fabro-checkpoint/src/metadata.rs b/lib/crates/fabro-checkpoint/src/metadata.rs index cdbe86711..d6d634ae1 100644 --- a/lib/crates/fabro-checkpoint/src/metadata.rs +++ b/lib/crates/fabro-checkpoint/src/metadata.rs @@ -214,6 +214,7 @@ mod tests { repo_origin_url: None, base_branch: None, labels: HashMap::new(), + artifact_storage: None, provenance: None, } } diff --git a/lib/crates/fabro-cli/src/commands/run/attach.rs b/lib/crates/fabro-cli/src/commands/run/attach.rs index 9c0585489..51a7e02b8 100644 --- a/lib/crates/fabro-cli/src/commands/run/attach.rs +++ b/lib/crates/fabro-cli/src/commands/run/attach.rs @@ -80,8 +80,7 @@ pub(crate) async fn attach_run_with_client( .or(state_exit_code) .unwrap_or(ExitCode::from(1)), json_output, - ) - .await; + ); } let stream = client.attach_run_events(run_id, Some(next_seq)).await?; @@ -98,7 +97,7 @@ pub(crate) async fn attach_run_with_client( .await } -async fn replay_run_with_client( +fn replay_run_with_client( verbose: bool, events: Vec, exit_code: ExitCode, diff --git a/lib/crates/fabro-cli/src/commands/run/output.rs b/lib/crates/fabro-cli/src/commands/run/output.rs index 017a0b846..703a3860c 100644 --- a/lib/crates/fabro-cli/src/commands/run/output.rs +++ b/lib/crates/fabro-cli/src/commands/run/output.rs @@ -180,7 +180,7 @@ pub(crate) fn print_run_conclusion( ); if let Some(billing) = conclusion.billing.as_ref() { - let total_tokens = i64::try_from(billing.total_tokens).unwrap_or(i64::MAX); + let total_tokens = billing.total_tokens; if total_tokens > 0 { if let Some(total_usd_micros) = billing.total_usd_micros { if total_usd_micros > 0 { @@ -206,12 +206,8 @@ pub(crate) fn print_run_conclusion( "{}", styles.dim.apply_to(format!( "Cache: {} read, {} write", - format_tokens_human( - i64::try_from(billing.cache_read_tokens).unwrap_or(i64::MAX) - ), - format_tokens_human( - i64::try_from(billing.cache_write_tokens).unwrap_or(i64::MAX) - ), + format_tokens_human(billing.cache_read_tokens), + format_tokens_human(billing.cache_write_tokens), )), ); } @@ -220,9 +216,7 @@ pub(crate) fn print_run_conclusion( "{}", styles.dim.apply_to(format!( "Reasoning: {} tokens", - format_tokens_human( - i64::try_from(billing.reasoning_tokens).unwrap_or(i64::MAX) - ), + format_tokens_human(billing.reasoning_tokens), )), ); } diff --git a/lib/crates/fabro-cli/src/commands/run/runner.rs b/lib/crates/fabro-cli/src/commands/run/runner.rs index 2a9c7dbb6..04e8a3e5f 100644 --- a/lib/crates/fabro-cli/src/commands/run/runner.rs +++ b/lib/crates/fabro-cli/src/commands/run/runner.rs @@ -12,6 +12,7 @@ use fabro_types::{EventBody, RunBlobId, RunEvent, RunId, Settings, StatusReason} use fabro_workflow::artifact_snapshot::CapturedArtifactInfo; use fabro_workflow::artifact_upload::StageArtifactUploader; use fabro_workflow::event::{Emitter, RunEventSink}; +use fabro_workflow::operations::{self, StartServices}; use fabro_workflow::run_control::RunControlState; use fabro_workflow::runtime_store::{RunStoreBackend, RunStoreHandle}; #[cfg(unix)] @@ -67,7 +68,7 @@ pub(crate) async fn execute( run_record, client.clone_for_reuse(), artifact_upload_token, - )?; + ); let scratch = RunScratch::new(&run_dir); let interviewer = Arc::new(FileInterviewer::new( scratch.interview_request_path(), @@ -78,7 +79,7 @@ pub(crate) async fn execute( let cancel_token = Arc::new(AtomicBool::new(false)); install_signal_handlers(Arc::clone(&run_control), Arc::clone(&cancel_token))?; let github_app = maybe_build_github_app_credentials(&run_record.settings)?; - let services = fabro_workflow::operations::StartServices { + let services = StartServices { run_id, cancel_token: Some(Arc::clone(&cancel_token)), emitter: Arc::new(Emitter::new(run_id)), @@ -100,10 +101,10 @@ pub(crate) async fn execute( match mode { RunWorkerMode::Start => { - fabro_workflow::operations::start(&run_dir, services).await?; + operations::start(&run_dir, services).await?; } RunWorkerMode::Resume => { - fabro_workflow::operations::resume(&run_dir, services).await?; + operations::resume(&run_dir, services).await?; } } @@ -115,19 +116,21 @@ fn build_artifact_uploader( run_record: &fabro_types::RunRecord, client: server_client::ServerStoreClient, artifact_upload_token: Option, -) -> Result>> { +) -> Option> { if !run_record.uses_object_backed_artifacts() { - return Ok(None); + return None; } - let token = artifact_upload_token - .ok_or_else(|| anyhow!("run {run_id} is configured for object-backed artifacts but the worker did not receive an artifact upload token"))?; + let uploader: Arc = match artifact_upload_token { + Some(token) => Arc::new(HttpArtifactUploader { + run_id, + client, + bearer_token: token, + }), + None => Arc::new(MissingArtifactUploadTokenUploader { run_id }), + }; - Ok(Some(Arc::new(HttpArtifactUploader { - run_id, - client, - bearer_token: token, - }))) + Some(uploader) } struct HttpArtifactUploader { @@ -174,6 +177,25 @@ impl StageArtifactUploader for HttpArtifactUploader { } } +struct MissingArtifactUploadTokenUploader { + run_id: RunId, +} + +#[async_trait] +impl StageArtifactUploader for MissingArtifactUploadTokenUploader { + async fn upload_stage_artifacts( + &self, + _stage_id: &fabro_types::StageId, + _artifact_capture_dir: &Path, + _artifacts: &[CapturedArtifactInfo], + ) -> Result<()> { + Err(anyhow!( + "run {} is configured for object-backed artifacts but the worker did not receive an artifact upload token", + self.run_id + )) + } +} + #[derive(Clone)] struct HttpRunStore { run_id: RunId, diff --git a/lib/crates/fabro-cli/src/server_client.rs b/lib/crates/fabro-cli/src/server_client.rs index 6bcac2133..5ae32e89f 100644 --- a/lib/crates/fabro-cli/src/server_client.rs +++ b/lib/crates/fabro-cli/src/server_client.rs @@ -11,8 +11,11 @@ use fabro_store::{EventEnvelope, RunSummary, StageId}; use fabro_types::{RunBlobId, RunEvent, RunId, Settings}; use fabro_workflow::artifact_snapshot::CapturedArtifactInfo; use futures::StreamExt; +use reqwest::header::{CONTENT_LENGTH, CONTENT_TYPE}; +use reqwest::multipart::{Form, Part}; use serde::Serialize; use serde::de::DeserializeOwned; +use tokio::fs::File; use tokio::time::sleep; use tokio_util::io::ReaderStream; @@ -20,6 +23,7 @@ use crate::args::ServerTargetArgs; use crate::commands::server::start; use crate::sse; use crate::user_config; +use crate::user_config::cli_http_client_builder; #[derive(Clone)] pub(crate) struct ServerStoreClient { @@ -55,16 +59,13 @@ impl RunAttachEventStream { return Ok(Some(event)); } - match self.stream.next().await { - Some(chunk) => { - let chunk = chunk.map_err(|err| anyhow!("{err}"))?; - self.pending_bytes.extend_from_slice(&chunk); - self.buffer_sse_events(false)?; - } - None => { - self.buffer_sse_events(true)?; - return Ok(self.buffered_events.pop_front()); - } + if let Some(chunk) = self.stream.next().await { + let chunk = chunk.map_err(|err| anyhow!("{err}"))?; + self.pending_bytes.extend_from_slice(&chunk); + self.buffer_sse_events(false)?; + } else { + self.buffer_sse_events(true)?; + return Ok(self.buffered_events.pop_front()); } } } @@ -199,7 +200,7 @@ fn normalize_remote_server_target(api_url: &str) -> String { } async fn connect_unix_socket_api_client_bundle(path: &Path) -> Result { - let http_client = crate::user_config::cli_http_client_builder() + let http_client = cli_http_client_builder() .unix_socket(path) .no_proxy() .build() @@ -552,7 +553,7 @@ impl ServerStoreClient { let mut url = reqwest::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"))? + .map_err(|()| anyhow!("server base URL cannot accept path segments"))? .extend([ "api", "v1", @@ -576,7 +577,7 @@ impl ServerStoreClient { let mut url = self.stage_artifacts_url(run_id, stage_id)?; url.query_pairs_mut().append_pair("filename", filename); - let file = tokio::fs::File::open(path) + let file = File::open(path) .await .with_context(|| format!("failed to open artifact {}", path.display()))?; let content_length = file @@ -590,8 +591,8 @@ impl ServerStoreClient { .http_client .post(url) .bearer_auth(bearer_token) - .header(reqwest::header::CONTENT_TYPE, "application/octet-stream") - .header(reqwest::header::CONTENT_LENGTH, content_length.to_string()) + .header(CONTENT_TYPE, "application/octet-stream") + .header(CONTENT_LENGTH, content_length.to_string()) .body(body) .send() .await @@ -614,7 +615,7 @@ impl ServerStoreClient { for (index, artifact) in artifacts.iter().enumerate() { let part_name = format!("file{}", index + 1); let path = artifact_capture_dir.join(&artifact.path); - let file = tokio::fs::File::open(&path) + let file = File::open(&path) .await .with_context(|| format!("failed to open artifact {}", path.display()))?; let content_length = file @@ -633,7 +634,7 @@ impl ServerStoreClient { file_parts.push(( part_name, - reqwest::multipart::Part::stream_with_length( + Part::stream_with_length( reqwest::Body::wrap_stream(ReaderStream::new(file)), content_length, ) @@ -644,9 +645,9 @@ impl ServerStoreClient { let manifest = ArtifactBatchUploadManifest { entries: manifest_entries, }; - let manifest_part = reqwest::multipart::Part::text(serde_json::to_string(&manifest)?) - .mime_str("application/json")?; - let mut form = reqwest::multipart::Form::new().part("manifest", manifest_part); + let manifest_part = + Part::text(serde_json::to_string(&manifest)?).mime_str("application/json")?; + let mut form = Form::new().part("manifest", manifest_part); for (part_name, part) in file_parts { form = form.part(part_name, part); } diff --git a/lib/crates/fabro-cli/src/user_config.rs b/lib/crates/fabro-cli/src/user_config.rs index a6b929d9f..672f4ba11 100644 --- a/lib/crates/fabro-cli/src/user_config.rs +++ b/lib/crates/fabro-cli/src/user_config.rs @@ -8,6 +8,7 @@ use fabro_types::Settings; use tracing::debug; use crate::args::ServerTargetArgs; +use fabro_util::version::FABRO_VERSION; pub(crate) fn load_settings() -> anyhow::Result { load_settings_with_config_and_storage_dir(None, None) @@ -133,8 +134,7 @@ pub(crate) fn exec_server_target( } pub(crate) fn cli_http_client_builder() -> reqwest::ClientBuilder { - reqwest::Client::builder() - .user_agent(format!("fabro-cli/{}", fabro_util::version::FABRO_VERSION)) + reqwest::Client::builder().user_agent(format!("fabro-cli/{FABRO_VERSION}")) } pub(crate) fn build_server_client( diff --git a/lib/crates/fabro-cli/tests/it/cmd/attach.rs b/lib/crates/fabro-cli/tests/it/cmd/attach.rs index d556b6643..a3b983227 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/attach.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/attach.rs @@ -571,6 +571,7 @@ fn attach_json_errors_without_prompting_for_human_input() { "event": "run.created", "id": "[EVENT_ID]", "properties": { + "artifact_storage": "object_store_v1", "graph": { "attrs": { "goal": { diff --git a/lib/crates/fabro-cli/tests/it/cmd/run.rs b/lib/crates/fabro-cli/tests/it/cmd/run.rs index 90e31d4fa..dad77e9fa 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/run.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/run.rs @@ -646,6 +646,7 @@ fn json_run_implies_auto_approve_for_human_gates() { "event": "run.created", "id": "[EVENT_ID]", "properties": { + "artifact_storage": "object_store_v1", "graph": { "attrs": { "goal": { diff --git a/lib/crates/fabro-cli/tests/it/cmd/runner.rs b/lib/crates/fabro-cli/tests/it/cmd/runner.rs index 8482c7c21..f4c11d718 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/runner.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/runner.rs @@ -44,11 +44,11 @@ fn help() { --json Output as JSON [env: FABRO_JSON=] --server Fabro server target: http(s) URL or absolute Unix socket path --debug Enable DEBUG-level logging (default is INFO) [env: FABRO_DEBUG=] - --run-dir Run scratch directory --no-upgrade-check Disable automatic upgrade check [env: FABRO_NO_UPGRADE_CHECK=true] + --run-dir Run scratch directory + --quiet Suppress non-essential output [env: FABRO_QUIET=] --run-id Run ID --mode Worker mode [possible values: start, resume] - --quiet Suppress non-essential output [env: FABRO_QUIET=] --verbose Enable verbose output [env: FABRO_VERBOSE=] -h, --help Print help ----- stderr ----- diff --git a/lib/crates/fabro-llm/src/generate.rs b/lib/crates/fabro-llm/src/generate.rs index 4f06156ce..4f50a8722 100644 --- a/lib/crates/fabro-llm/src/generate.rs +++ b/lib/crates/fabro-llm/src/generate.rs @@ -221,7 +221,7 @@ pub async fn generate(params: GenerateParams) -> Result i64 { + i64::try_from(value).unwrap_or_else(|_| { + if value.is_negative() { + i64::MIN + } else { + i64::MAX + } + }) +} + +#[allow(clippy::cast_possible_truncation, clippy::cast_precision_loss)] +fn saturating_rounded_f64_to_i64(value: f64) -> i64 { + if !value.is_finite() { + return if value.is_sign_negative() { + i64::MIN + } else { + i64::MAX + }; + } + + if value <= i64::MIN as f64 { + i64::MIN + } else if value >= i64::MAX as f64 { + i64::MAX + } else { + value as i64 + } +} #[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default, Serialize, Deserialize)] pub struct UsdMicros(pub i64); @@ -45,7 +74,7 @@ impl PricePerMTok { #[must_use] pub fn from_usd(usd: f64) -> Self { Self { - usd_micros: (usd * USD_MICROS_PER_USD as f64).round() as i64, + usd_micros: saturating_rounded_f64_to_i64((usd * USD_MICROS_PER_USD_F64).round()), } } @@ -59,7 +88,7 @@ impl PricePerMTok { #[must_use] pub fn bill(self, tokens: i64) -> UsdMicros { let total = i128::from(tokens) * i128::from(self.usd_micros); - UsdMicros((total / TOKENS_PER_MTOK) as i64) + UsdMicros(saturating_i128_to_i64(total / TOKENS_PER_MTOK)) } } @@ -218,6 +247,7 @@ pub struct ModelPricing { pub policy: ModelPricingPolicy, } +#[allow(clippy::empty_structs_with_brackets)] #[derive(Debug, Clone, PartialEq, Eq, Default, Serialize, Deserialize)] pub struct OpenAiBillingFacts {} @@ -533,10 +563,10 @@ fn bill_gemini( .map(|segment| { let token_seconds = i128::from(segment.cached_tokens) * i128::from(segment.ttl_seconds); - UsdMicros( - (token_seconds * i128::from(storage.usd_micros_per_mtok_second) - / TOKENS_PER_MTOK) as i64, - ) + UsdMicros(saturating_i128_to_i64( + token_seconds * i128::from(storage.usd_micros_per_mtok_second) + / TOKENS_PER_MTOK, + )) }) .sum::(); total += storage_cost; @@ -689,4 +719,20 @@ mod tests { assert_eq!(pricing.bill(&input), None); } + + #[test] + fn price_per_mtok_bill_saturates_large_totals() { + let price = PricePerMTok { + usd_micros: i64::MAX, + }; + + assert_eq!(price.bill(i64::MAX), UsdMicros(i64::MAX)); + } + + #[test] + fn price_per_mtok_from_usd_saturates_large_inputs() { + let price = PricePerMTok::from_usd(f64::MAX); + + assert_eq!(price.usd_micros, i64::MAX); + } } diff --git a/lib/crates/fabro-server/src/serve.rs b/lib/crates/fabro-server/src/serve.rs index 202d82012..02a308a56 100644 --- a/lib/crates/fabro-server/src/serve.rs +++ b/lib/crates/fabro-server/src/serve.rs @@ -110,7 +110,7 @@ fn apply_runtime_settings( fn use_in_memory_store() -> bool { !matches!( std::env::var(TEST_IN_MEMORY_STORE_ENV).ok().as_deref(), - None | Some("") | Some("0") | Some("false") | Some("no") + None | Some("" | "0" | "false" | "no") ) } @@ -472,6 +472,7 @@ async fn wait_for_shutdown(mut shutdown_rx: watch::Receiver) { let _ = shutdown_rx.changed().await; } +#[allow(clippy::print_stderr)] fn announce_server_ready(bind_addr: &Bind, styles: &'static Styles, dry_run_mode: bool) { set_server_title(ServerTitlePhase::Listening, Some(bind_addr)); info!(bind = %bind_addr, dry_run = dry_run_mode, "API server started"); diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 4c0633d1a..7c7a1249d 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -7,6 +7,7 @@ use std::sync::{Arc, Mutex, RwLock}; use std::time::{Duration, Instant}; use crate::bind::Bind; +use axum::body::Body; #[cfg(test)] use axum::body::to_bytes; use axum::extract::{self as axum_extract, DefaultBodyLimit, Path, Query, State}; @@ -42,12 +43,12 @@ use fabro_workflow::handler::HandlerRegistry; use futures_util::stream; use jsonwebtoken::{Algorithm, DecodingKey, EncodingKey, Header, Validation}; use object_store::memory::InMemory as MemoryObjectStore; -use rand::RngCore; +use rand::{RngCore, rngs::OsRng}; use sha2::{Digest, Sha256}; use tempfile::NamedTempFile; use tokio::fs; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; -use tokio::process::Command; +use tokio::process::{ChildStderr, Command}; use tokio::sync::Notify; use tokio::sync::RwLock as AsyncRwLock; use tokio::sync::broadcast; @@ -109,6 +110,8 @@ pub fn default_page_limit() -> u32 { 20 } +const ATTACH_REPLAY_BATCH_LIMIT: usize = 256; + #[derive(serde::Deserialize)] pub struct PaginationParams { #[serde(rename = "page[limit]", default = "default_page_limit")] @@ -457,7 +460,7 @@ impl AppState { fn artifact_upload_token_keys() -> ArtifactUploadTokenKeys { let mut secret = [0_u8; 32]; - rand::rngs::OsRng.fill_bytes(&mut secret); + OsRng.fill_bytes(&mut secret); let mut validation = Validation::new(Algorithm::HS256); validation.set_required_spec_claims(&["iss", "iat", "exp"]); @@ -475,17 +478,15 @@ fn maybe_authorize_artifact_upload_token( run_id: &RunId, keys: &ArtifactUploadTokenKeys, ) -> Result { - let header = match parts + let Some(header) = parts .headers .get(header::AUTHORIZATION) .and_then(|value| value.to_str().ok()) - { - Some(header) => header, - None => return Ok(false), + else { + return Ok(false); }; - let token = match header.strip_prefix("Bearer ") { - Some(token) => token, - None => return Ok(false), + let Some(token) = header.strip_prefix("Bearer ") else { + return Ok(false); }; let claims = @@ -1743,10 +1744,10 @@ async fn list_board_runs( .map(|(id, managed_run)| { ( *id, - managed_run.status.clone(), + managed_run.status, managed_run.error.clone(), queue_positions.get(id).copied(), - managed_run.created_at.clone(), + managed_run.created_at, ) }) .collect::>() @@ -1773,7 +1774,7 @@ async fn list_board_runs( let summary = summaries.get(id); RunStatusResponse { id: id.to_string(), - status: status.clone(), + status: *status, error: error.as_ref().map(|msg| RunError { message: msg.clone(), }), @@ -1782,7 +1783,7 @@ async fn list_board_runs( .and_then(|summary| summary.status_reason.map(api_status_reason)), pending_control: summary .and_then(|summary| summary.pending_control.map(api_pending_control)), - created_at: created_at.clone(), + created_at: *created_at, } }) .collect(); @@ -1892,8 +1893,6 @@ async fn terminate_worker_for_deletion(worker_pid: Option, worker_pgid: Opt sleep(Duration::from_millis(50)).await; } } - - return; } #[cfg(not(unix))] @@ -2387,7 +2386,7 @@ fn update_live_run_from_event(state: &Arc, run_id: RunId, event: &RunE match &event.body { EventBody::RunStarting(_) => managed_run.status = RunStatus::Starting, EventBody::RunRunning(_) | EventBody::RunUnpaused(_) => { - managed_run.status = RunStatus::Running + managed_run.status = RunStatus::Running; } EventBody::RunPaused(_) => managed_run.status = RunStatus::Paused, EventBody::RunCompleted(_) => { @@ -2409,7 +2408,7 @@ fn update_live_run_from_event(state: &Arc, run_id: RunId, event: &RunE async fn drain_worker_stderr( run_id: RunId, run_dir: PathBuf, - stderr: tokio::process::ChildStderr, + stderr: ChildStderr, ) -> anyhow::Result<()> { let log_path = run_dir.join("runtime").join(WORKER_STDERR_LOG); if let Some(parent) = log_path.parent() { @@ -2502,9 +2501,8 @@ fn worker_command( mode: RunExecutionMode, run_dir: &std::path::Path, ) -> anyhow::Result { - let exe = std::env::var_os("CARGO_BIN_EXE_fabro") - .map(PathBuf::from) - .unwrap_or(std::env::current_exe()?); + let exe = + std::env::var_os("CARGO_BIN_EXE_fabro").map_or(std::env::current_exe()?, PathBuf::from); let storage_dir = state .settings .read() @@ -2561,6 +2559,7 @@ fn api_question_from_interview_question(id: &str, question: &Question) -> ApiQue } } +#[allow(clippy::result_large_err)] fn answer_from_request(req: SubmitAnswerRequest, question: &Question) -> Result { if let Some(key) = req.selected_option_key { let option = question @@ -3664,19 +3663,16 @@ async fn attach_run_events( } }, }; - const ATTACH_REPLAY_BATCH_LIMIT: usize = 256; - let (sender, receiver) = mpsc::unbounded_channel(); tokio::spawn(async move { let mut next_seq = start_seq; loop { - let replay_batch = match run_store + let Ok(replay_batch) = run_store .list_events_from_with_limit(next_seq, ATTACH_REPLAY_BATCH_LIMIT) .await - { - Ok(events) => events, - Err(_) => return, + else { + return; }; let replay_has_more = replay_batch.len() > ATTACH_REPLAY_BATCH_LIMIT; @@ -3700,21 +3696,19 @@ async fn attach_run_events( continue; } - let state = match run_store.state().await { - Ok(state) => state, - Err(_) => return, + let Ok(state) = run_store.state().await else { + return; }; if run_projection_is_active(&state) { break; } - let tail_batch = match run_store + let Ok(tail_batch) = run_store .list_events_from_with_limit(next_seq, ATTACH_REPLAY_BATCH_LIMIT) .await - { - Ok(events) => events, - Err(_) => return, + else { + return; }; let tail_has_more = tail_batch.len() > ATTACH_REPLAY_BATCH_LIMIT; @@ -3741,9 +3735,8 @@ async fn attach_run_events( return; } - let mut live_stream = match run_store.watch_events_from(next_seq) { - Ok(stream) => stream, - Err(_) => return, + let Ok(mut live_stream) = run_store.watch_events_from(next_seq) else { + return; }; while let Some(result) = live_stream.next().await { @@ -4071,8 +4064,7 @@ fn validate_artifact_batch_manifest( } if manifest.entries.len() > MAX_MULTIPART_ARTIFACTS { return Err(payload_too_large_response(format!( - "multipart upload exceeds the {} artifact limit", - MAX_MULTIPART_ARTIFACTS + "multipart upload exceeds the {MAX_MULTIPART_ARTIFACTS} artifact limit" ))); } @@ -4108,15 +4100,13 @@ fn validate_artifact_batch_manifest( if let Some(expected_bytes) = entry.expected_bytes { if expected_bytes > MAX_SINGLE_ARTIFACT_BYTES { return Err(payload_too_large_response(format!( - "artifact {} exceeds the {} byte limit", - path, MAX_SINGLE_ARTIFACT_BYTES + "artifact {path} exceeds the {MAX_SINGLE_ARTIFACT_BYTES} byte limit" ))); } expected_total_bytes = expected_total_bytes.saturating_add(expected_bytes); if expected_total_bytes > MAX_MULTIPART_REQUEST_BYTES { return Err(payload_too_large_response(format!( - "multipart upload exceeds the {} byte limit", - MAX_MULTIPART_REQUEST_BYTES + "multipart upload exceeds the {MAX_MULTIPART_REQUEST_BYTES} byte limit" ))); } } @@ -4146,7 +4136,7 @@ async fn upload_stage_artifact_octet_stream( run_id: &RunId, stage_id: &StageId, filename: String, - body: axum::body::Body, + body: Body, content_length: Option, ) -> Response { let relative_path = match validate_relative_artifact_path("filename", &filename) { @@ -4156,8 +4146,7 @@ async fn upload_stage_artifact_octet_stream( if content_length.is_some_and(|length| length > MAX_SINGLE_ARTIFACT_BYTES) { return payload_too_large_response(format!( - "artifact exceeds the {} byte limit", - MAX_SINGLE_ARTIFACT_BYTES + "artifact exceeds the {MAX_SINGLE_ARTIFACT_BYTES} byte limit" )); } @@ -4185,8 +4174,7 @@ async fn upload_stage_artifact_octet_stream( bytes_written.saturating_add(u64::try_from(chunk.len()).unwrap_or(u64::MAX)); if bytes_written > MAX_SINGLE_ARTIFACT_BYTES { return payload_too_large_response(format!( - "artifact exceeds the {} byte limit", - MAX_SINGLE_ARTIFACT_BYTES + "artifact exceeds the {MAX_SINGLE_ARTIFACT_BYTES} byte limit" )); } if let Err(err) = writer.write_all(&chunk).await { @@ -4208,7 +4196,7 @@ async fn upload_stage_artifact_multipart( run_id: &RunId, stage_id: &StageId, boundary: String, - body: axum::body::Body, + body: Body, ) -> Response { let mut multipart = multer::Multipart::new(body.into_data_stream(), boundary); let Some(mut manifest_field) = (match multipart @@ -4276,14 +4264,13 @@ async fn upload_stage_artifact_multipart( if bytes_written > MAX_SINGLE_ARTIFACT_BYTES { return payload_too_large_response(format!( - "artifact {} exceeds the {} byte limit", - entry.path, MAX_SINGLE_ARTIFACT_BYTES + "artifact {} exceeds the {MAX_SINGLE_ARTIFACT_BYTES} byte limit", + entry.path )); } if total_bytes > MAX_MULTIPART_REQUEST_BYTES { return payload_too_large_response(format!( - "multipart upload exceeds the {} byte limit", - MAX_MULTIPART_REQUEST_BYTES + "multipart upload exceeds the {MAX_MULTIPART_REQUEST_BYTES} byte limit" )); } @@ -4374,8 +4361,7 @@ async fn put_stage_artifact( Ok(ArtifactUploadContentType::Multipart { boundary }) => { if content_length.is_some_and(|length| length > MAX_MULTIPART_REQUEST_BYTES) { return payload_too_large_response(format!( - "multipart upload exceeds the {} byte limit", - MAX_MULTIPART_REQUEST_BYTES + "multipart upload exceeds the {MAX_MULTIPART_REQUEST_BYTES} byte limit" )); } upload_stage_artifact_multipart(state.as_ref(), &id, &stage_id, boundary, body).await diff --git a/lib/crates/fabro-store/src/artifact_store.rs b/lib/crates/fabro-store/src/artifact_store.rs index f277a1c3d..6389d1d1a 100644 --- a/lib/crates/fabro-store/src/artifact_store.rs +++ b/lib/crates/fabro-store/src/artifact_store.rs @@ -2,6 +2,7 @@ use std::sync::Arc; use bytes::Bytes; use futures::StreamExt; +use object_store::buffered::BufWriter; use object_store::{ObjectStore, path::Path as ObjectPath}; use percent_encoding::{AsciiSet, NON_ALPHANUMERIC, percent_decode_str, utf8_percent_encode}; use tokio::io::AsyncWriteExt; @@ -57,14 +58,9 @@ impl ArtifactStore { Ok(()) } - pub fn writer( - &self, - run_id: &RunId, - node: &StageId, - filename: &str, - ) -> Result { + pub fn writer(&self, run_id: &RunId, node: &StageId, filename: &str) -> Result { let path = self.artifact_path(run_id, node, filename)?; - Ok(object_store::buffered::BufWriter::with_capacity( + Ok(BufWriter::with_capacity( Arc::clone(&self.object_store), path, STREAM_BUFFER_BYTES, diff --git a/lib/crates/fabro-store/src/slate/run_store.rs b/lib/crates/fabro-store/src/slate/run_store.rs index 8b0025713..14db0cf70 100644 --- a/lib/crates/fabro-store/src/slate/run_store.rs +++ b/lib/crates/fabro-store/src/slate/run_store.rs @@ -270,7 +270,7 @@ impl RunDatabase { } } Err(broadcast::error::TryRecvError::Empty) => break, - Err(broadcast::error::TryRecvError::Lagged(_)) => continue, + Err(broadcast::error::TryRecvError::Lagged(_)) => {} Err(broadcast::error::TryRecvError::Closed) => return, } } diff --git a/lib/crates/fabro-test/src/lib.rs b/lib/crates/fabro-test/src/lib.rs index 7830fa27f..c9ee10d3c 100644 --- a/lib/crates/fabro-test/src/lib.rs +++ b/lib/crates/fabro-test/src/lib.rs @@ -11,6 +11,7 @@ use fabro_types::RunId; use regex::Regex; use serde::Serialize; use serde_json::{Map, Value, json}; +use toml::{Value as TomlValue, map::Map as TomlMap}; /// Walk up from `start` to find the repo-level `test/` fixtures directory. pub fn find_test_fixtures_dir(start: &Path) -> Option { @@ -373,9 +374,9 @@ fn write_settings_file(path: &Path, storage_dir: &Path, rest: &str) { .unwrap_or_else(|err| panic!("failed to write {}: {err}", path.display())); } -fn parse_settings_table(contents: &str, source: &Path) -> toml::map::Map { +fn parse_settings_table(contents: &str, source: &Path) -> TomlMap { let stripped = strip_managed_storage_settings(contents); - let value = toml::from_str::(stripped) + let value = toml::from_str::(stripped) .unwrap_or_else(|err| panic!("failed to parse {}: {err}", source.display())); let Some(table) = value.as_table() else { panic!("expected {} to contain a TOML table", source.display()); @@ -383,7 +384,7 @@ fn parse_settings_table(contents: &str, source: &Path) -> toml::map::Map) { +fn write_settings_table(path: &Path, table: &TomlMap) { ensure_parent_dir(path); let mut contents = toml::to_string(table) .unwrap_or_else(|err| panic!("failed to serialize {}: {err}", path.display())); @@ -394,29 +395,29 @@ fn write_settings_table(path: &Path, table: &toml::map::Map .unwrap_or_else(|err| panic!("failed to write {}: {err}", path.display())); } -fn server_target_from_table(table: &toml::map::Map) -> Option { +fn server_target_from_table(table: &TomlMap) -> Option { table .get("server") - .and_then(toml::Value::as_table) + .and_then(TomlValue::as_table) .and_then(|server| server.get("target")) - .and_then(toml::Value::as_str) + .and_then(TomlValue::as_str) .map(ToOwned::to_owned) } -fn set_server_target(table: &mut toml::map::Map, socket_path: &Path) { +fn set_server_target(table: &mut TomlMap, socket_path: &Path) { let server_entry = table .entry("server".to_string()) - .or_insert_with(|| toml::Value::Table(toml::map::Map::new())); + .or_insert_with(|| TomlValue::Table(TomlMap::new())); let Some(server_table) = server_entry.as_table_mut() else { panic!("expected [server] to be a TOML table"); }; server_table.insert( "target".to_string(), - toml::Value::String(socket_path.display().to_string()), + TomlValue::String(socket_path.display().to_string()), ); } -fn clear_server_target(table: &mut toml::map::Map) { +fn clear_server_target(table: &mut TomlMap) { let Some(server_entry) = table.get_mut("server") else { return; }; @@ -446,7 +447,7 @@ fn sync_home_settings( (table, had_explicit_storage, had_explicit_target) } Err(err) if err.kind() == std::io::ErrorKind::NotFound => { - (toml::map::Map::new(), false, false) + (TomlMap::new(), false, false) } Err(err) => panic!("failed to read {}: {err}", settings_path.display()), }; @@ -454,7 +455,7 @@ fn sync_home_settings( if !had_explicit_storage { table.insert( "storage_dir".to_string(), - toml::Value::String(storage_dir.display().to_string()), + TomlValue::String(storage_dir.display().to_string()), ); table.remove("data_dir"); } @@ -509,7 +510,7 @@ fn server_record_pid(storage_dir: &Path) -> Option { } fn server_running(server: &ServerPaths) -> bool { - server_record_pid(&server.storage_dir).is_some_and(|pid| fabro_proc::process_alive(pid)) + server_record_pid(&server.storage_dir).is_some_and(fabro_proc::process_alive) } fn wait_for_server_running(server: &ServerPaths) { @@ -558,13 +559,12 @@ fn ensure_server_running(fabro_bin: &Path, server: &ServerPaths, config_path: &P .unwrap_or_else(|err| panic!("failed to execute {}: {err}", fabro_bin.display())); let stderr = String::from_utf8_lossy(&output.stderr); - if !output.status.success() && !stderr.contains("Server already running") { - panic!( - "failed to start test server:\nstdout:\n{}\nstderr:\n{}", - String::from_utf8_lossy(&output.stdout), - stderr - ); - } + assert!( + output.status.success() || stderr.contains("Server already running"), + "failed to start test server:\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&output.stdout), + stderr + ); wait_for_server_running(server); } @@ -1073,8 +1073,8 @@ impl TestContext { if fabro_bin_exists(&self.fabro_bin) { ensure_server_running(&self.fabro_bin, &server, &settings_path); } - self.storage_dir = server.storage_dir.clone(); - self.active_socket_path = server.socket_path.clone(); + self.storage_dir.clone_from(&server.storage_dir); + self.active_socket_path.clone_from(&server.socket_path); self.isolated_server = Some(server); self } diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs index 89d46ca63..b91317427 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/api.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs @@ -567,7 +567,7 @@ impl CodergenBackend for AgentApiBackend { let mut total_usage = TokenCounts::default(); for turn in &session.history().turns()[turns_before..] { if let Turn::Assistant { usage, .. } = turn { - total_usage = total_usage + *usage.clone(); + total_usage += *usage.clone(); } } diff --git a/lib/crates/fabro-workflow/src/lifecycle/artifact.rs b/lib/crates/fabro-workflow/src/lifecycle/artifact.rs index 945f705bb..9856198fb 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/artifact.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/artifact.rs @@ -5,6 +5,7 @@ use std::time::Duration; use async_trait::async_trait; use fabro_types::StageId; +use tokio::time::sleep; use fabro_core::graph::NodeSpec; use fabro_core::lifecycle::{AttemptContext, AttemptResultContext, RunLifecycle}; @@ -12,7 +13,7 @@ use fabro_core::outcome::NodeResult; use fabro_core::state::ExecutionState; use crate::artifact::{offload_large_values, sync_artifacts_to_env}; -use crate::artifact_snapshot::collect_artifacts; +use crate::artifact_snapshot::{CapturedArtifactInfo, collect_artifacts}; use crate::artifact_upload::StageArtifactUploader; use crate::event::{Emitter, Event, RunNoticeLevel}; use crate::graph::WorkflowGraph; @@ -207,7 +208,7 @@ impl ArtifactLifecycle { &self, stage_id: &StageId, artifact_capture_dir: &std::path::Path, - artifacts: &[crate::artifact_snapshot::CapturedArtifactInfo], + artifacts: &[CapturedArtifactInfo], ) -> Result<(), String> { let Some(uploader) = self.artifact_uploader.as_ref() else { return Ok(()); @@ -224,7 +225,7 @@ impl ArtifactLifecycle { } if let Some(delay) = ARTIFACT_UPLOAD_RETRY_DELAYS.get(attempt) { - tokio::time::sleep(*delay).await; + sleep(*delay).await; } } diff --git a/lib/crates/fabro-workflow/src/operations/create.rs b/lib/crates/fabro-workflow/src/operations/create.rs index 0edbf44fe..18957d5f1 100644 --- a/lib/crates/fabro-workflow/src/operations/create.rs +++ b/lib/crates/fabro-workflow/src/operations/create.rs @@ -703,6 +703,7 @@ mod tests { host_repo_path: None, repo_origin_url: None, base_branch: None, + artifact_storage: None, provenance: None, }, ) @@ -751,6 +752,7 @@ mod tests { host_repo_path: Some(dir.path().display().to_string()), repo_origin_url: None, base_branch: Some("main".to_string()), + artifact_storage: None, provenance: None, }, ) @@ -831,6 +833,7 @@ mod tests { host_repo_path: None, repo_origin_url: None, base_branch: None, + artifact_storage: None, provenance: None, }, ) @@ -874,6 +877,7 @@ mod tests { host_repo_path: None, repo_origin_url: Some("https://github.com/acme/widgets".to_string()), base_branch: None, + artifact_storage: None, provenance: None, }, ) @@ -914,6 +918,7 @@ mod tests { host_repo_path: None, repo_origin_url: None, base_branch: None, + artifact_storage: None, provenance: None, }, ) @@ -956,6 +961,7 @@ mod tests { host_repo_path: None, repo_origin_url: None, base_branch: None, + artifact_storage: None, provenance: Some(fabro_types::RunProvenance { server: Some(fabro_types::RunServerProvenance { version: "0.9.0".to_string(), diff --git a/lib/crates/fabro-workflow/src/operations/start.rs b/lib/crates/fabro-workflow/src/operations/start.rs index cb3850699..7bd4a6ffa 100644 --- a/lib/crates/fabro-workflow/src/operations/start.rs +++ b/lib/crates/fabro-workflow/src/operations/start.rs @@ -846,6 +846,7 @@ mod tests { host_repo_path: None, repo_origin_url: None, base_branch: None, + artifact_storage: None, provenance: None, }, ) diff --git a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs index 16925eef5..a31765a83 100644 --- a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs +++ b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs @@ -144,6 +144,7 @@ fn persisted_workflow(graph: Graph, source: String, run_dir: &Path, run_id: RunI repo_origin_url: None, base_branch: Some("main".to_string()), labels: HashMap::new(), + artifact_storage: None, provenance: None, }, ) diff --git a/lib/crates/fabro-workflow/src/pipeline/initialize.rs b/lib/crates/fabro-workflow/src/pipeline/initialize.rs index 4af604ffb..891b724d6 100644 --- a/lib/crates/fabro-workflow/src/pipeline/initialize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/initialize.rs @@ -749,6 +749,7 @@ mod tests { repo_origin_url: None, base_branch: Some("main".to_string()), labels: HashMap::new(), + artifact_storage: None, provenance: None, }, ) diff --git a/lib/crates/fabro-workflow/src/pipeline/retro.rs b/lib/crates/fabro-workflow/src/pipeline/retro.rs index bf1648260..9fdb40dc5 100644 --- a/lib/crates/fabro-workflow/src/pipeline/retro.rs +++ b/lib/crates/fabro-workflow/src/pipeline/retro.rs @@ -34,7 +34,7 @@ pub async fn run_retro(options: &RetroOptions, dry_run: bool) -> Option { return None; }; - let completed_stages = crate::build_completed_stages(&cp, options.failed); + let completed_stages = crate::build_completed_stages(cp, options.failed); let stage_durations = match options.run_store.list_events().await { Ok(events) => crate::extract_stage_durations_from_events(&events), Err(err) => {