diff --git a/Cargo.lock b/Cargo.lock index 4bba62a53..751a7d93f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2744,7 +2744,6 @@ dependencies = [ "jsonwebtoken", "lithos-llm", "mime_guess", - "multer", "object_store", "pebble-agent", "pebble-coding-agent", @@ -4622,23 +4621,6 @@ dependencies = [ "windows-sys 0.61.2", ] -[[package]] -name = "multer" -version = "3.1.0" -source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "83e87776546dc87511aa5ee218730c92b666d7264ab6ed41f9d215af9cd5224b" -dependencies = [ - "bytes", - "encoding_rs", - "futures-util", - "http 1.4.0", - "httparse", - "memchr", - "mime", - "spin", - "version_check", -] - [[package]] name = "native-tls" version = "0.2.18" diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 02c227ee0..646c1fad7 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -3106,24 +3106,6 @@ paths: schema: type: string format: binary - multipart/form-data: - schema: - type: object - required: - - manifest - properties: - manifest: - $ref: "#/components/schemas/ArtifactBatchUploadManifest" - additionalProperties: - type: string - format: binary - description: | - Strict multipart upload format. The `manifest` part must arrive first with JSON - matching `ArtifactBatchUploadManifest`. Each subsequent file part name must match - a manifest entry `part` value. - encoding: - manifest: - contentType: application/json responses: "200": description: Blob written @@ -3845,58 +3827,6 @@ paths: application/json: schema: $ref: "#/components/schemas/ErrorResponse" - post: - operationId: putStageArtifact - tags: [Run Internals] - summary: Put Stage Artifact - description: | - Uploads one or more artifacts for a stage. Intended for trusted internal callers. - - The server accepts both: - - `application/octet-stream` for single-file uploads with the `filename` query parameter - - strict manifest-first `multipart/form-data` uploads documented by `ArtifactBatchUploadManifest` - - The generated Rust client currently exposes the octet-stream variant because the OpenAPI - code generator in this repo does not support multiple request media types on one operation. - parameters: - - $ref: "#/components/parameters/RunId" - - $ref: "#/components/parameters/StageId" - - $ref: "#/components/parameters/ArtifactRetry" - - name: filename - in: query - required: false - description: Relative artifact path for `application/octet-stream` uploads. Ignored for multipart uploads. - schema: - type: string - requestBody: - required: true - content: - application/octet-stream: - schema: - type: string - format: binary - responses: - "204": - description: Artifact written - "400": - description: Invalid filename, multipart manifest, checksum, or upload body - headers: - x-request-id: - $ref: "#/components/headers/XRequestId" - content: - application/json: - schema: - $ref: "#/components/schemas/ErrorResponse" - "404": - description: Run not found - headers: - x-request-id: - $ref: "#/components/headers/XRequestId" - content: - application/json: - schema: - $ref: "#/components/schemas/ErrorResponse" - /api/v1/runs/{id}/stages/{stageId}/artifacts/download: get: operationId: getStageArtifact @@ -11200,48 +11130,6 @@ components: items: $ref: "#/components/schemas/ArtifactEntry" - ArtifactBatchUploadEntry: - description: One file entry in a strict multipart artifact upload manifest. - type: object - required: - - part - - path - properties: - part: - type: string - description: Multipart field name for the file part. - example: file1 - path: - type: string - description: Relative artifact path to store. - example: src/lib.rs - sha256: - type: ["string", "null"] - description: Optional SHA-256 checksum for the file contents; hex input is case-insensitive. - example: 3f785df4c5b7d3f1f4c1f0ecb0f55f1d9f6f6a3d9f0a8a98f7a74f29d1f81a2c - expected_bytes: - type: ["integer", "null"] - format: int64 - minimum: 0 - description: Optional exact byte length expected for the file part. - example: 1234 - content_type: - type: ["string", "null"] - description: Optional client-supplied content type for the file part. - example: text/plain - - ArtifactBatchUploadManifest: - description: Manifest for strict multipart artifact uploads. - type: object - required: - - entries - properties: - entries: - type: array - minItems: 1 - items: - $ref: "#/components/schemas/ArtifactBatchUploadEntry" - RunArtifactEntry: description: A captured artifact file for a run. type: object diff --git a/lib/apps/fabro-server/Cargo.toml b/lib/apps/fabro-server/Cargo.toml index f3b1d5a90..847ac14a2 100644 --- a/lib/apps/fabro-server/Cargo.toml +++ b/lib/apps/fabro-server/Cargo.toml @@ -100,7 +100,6 @@ mime_guess.workspace = true regex.workspace = true semver.workspace = true walkdir.workspace = true -multer = "3" thiserror.workspace = true percent-encoding.workspace = true url = "2" diff --git a/lib/apps/fabro-server/src/principal_middleware.rs b/lib/apps/fabro-server/src/principal_middleware.rs index ca3b16bc2..190583e22 100644 --- a/lib/apps/fabro-server/src/principal_middleware.rs +++ b/lib/apps/fabro-server/src/principal_middleware.rs @@ -66,7 +66,6 @@ pub(crate) struct RequireWorkerRunSegment(pub(crate) RunId, pub(crate) String); pub(crate) struct RequireRunManagementTarget(pub(crate) RunId, pub(crate) Principal); pub(crate) struct RequireRunBlob(pub(crate) RunId, pub(crate) BlobHash); pub(crate) struct RequireRunStageScoped(pub(crate) RunId, pub(crate) String); -pub(crate) struct RequireStageArtifact(pub(crate) RunId, pub(crate) StageId); pub(crate) struct RequireCommandLog(pub(crate) RunId, pub(crate) StageId); #[derive(Clone, Debug)] @@ -343,24 +342,6 @@ impl FromRequestParts> for RequireRunStageScoped { } } -impl FromRequestParts> for RequireStageArtifact { - type Rejection = Response; - - async fn from_request_parts( - parts: &mut Parts, - state: &Arc, - ) -> Result { - let Path((id, stage_id)): Path<(String, String)> = Path::from_request_parts(parts, state) - .await - .map_err(IntoResponse::into_response)?; - let run_id = parse_run_id_path(&id)?; - let stage_id = parse_stage_id_path(&stage_id)?; - require_worker_or_user_for_run(&auth_slot_from_parts(parts), &run_id) - .map_err(IntoResponse::into_response)?; - Ok(Self(run_id, stage_id)) - } -} - impl FromRequestParts> for RequireCommandLog { type Rejection = Response; diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 830eaad43..b85d6a34f 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -10,7 +10,7 @@ use anyhow::Context as _; use axum::body::Body; #[cfg(test)] use axum::body::to_bytes; -use axum::extract::{self as axum_extract, DefaultBodyLimit, Path, Query, State}; +use axum::extract::{self as axum_extract, Path, Query, State}; use axum::http::{HeaderMap, Method, StatusCode, header}; use axum::middleware::{self, Next}; use axum::response::sse::{Event, KeepAlive, Sse}; @@ -113,7 +113,6 @@ use fabro_workflow::{Error as WorkflowError, operations, pull_request}; use futures_util::future::join_all; use lithos_llm::catalog::ProviderId; use lithos_llm::types::Usage; -use sha2::{Digest, Sha256}; use tempfile::NamedTempFile; use tokio::fs; use tokio::io::{AsyncBufReadExt, AsyncRead, AsyncWriteExt, BufReader}; @@ -146,8 +145,8 @@ use crate::jwt_auth::{self, AuthMode}; use crate::petri_runs::PetriRuns; use crate::principal_middleware::{ AuthContextSlot, RequestAuth, RequestAuthContext, RequireRunBlob, RequireRunManagementTarget, - RequireRunScoped, RequireStageArtifact, RequireWorkerRunScoped, RequireWorkerRunSegment, - RequiredUser, principal_middleware, + RequireRunScoped, RequireWorkerRunScoped, RequireWorkerRunSegment, RequiredUser, + principal_middleware, }; use crate::request_id::{self, RequestId}; use crate::run_files::{FilesInFlight, new_files_in_flight}; @@ -3073,14 +3072,6 @@ fn validate_relative_artifact_path(kind: &str, value: &str) -> Result) -> Response { - ApiError::bad_request(detail.into()).into_response() -} - -fn payload_too_large_response(detail: impl Into) -> Response { - ApiError::new(StatusCode::PAYLOAD_TOO_LARGE, detail.into()).into_response() -} - fn octet_stream_response(bytes: Bytes) -> Response { ( StatusCode::OK, diff --git a/lib/apps/fabro-server/src/server/handler/artifacts.rs b/lib/apps/fabro-server/src/server/handler/artifacts.rs index 9cae243ab..3ccc0a8ba 100644 --- a/lib/apps/fabro-server/src/server/handler/artifacts.rs +++ b/lib/apps/fabro-server/src/server/handler/artifacts.rs @@ -20,13 +20,11 @@ use tracing::warn; use super::super::{ ApiError, AppState, ArtifactEntry, ArtifactKey, ArtifactListResponse, AsyncWriteExt, Body, - Bytes, DefaultBodyLimit, Digest, HashMap, HashSet, HeaderMap, IntoResponse, Json, NodeArtifact, - Path, Query, RequireRunBlob, RequireRunScoped, RequireStageArtifact, RequiredUser, Response, - Router, RunArtifactEntry, RunArtifactListResponse, RunId, Sha256, StageArtifactEntry, StageId, - State, StatusCode, StreamExt, WriteBlobResponse, axum_extract, bad_request_response, get, - header, octet_stream_response, parse_run_id_path, parse_stage_id_path, - payload_too_large_response, post, reject_if_archived, required_query_param, - validate_relative_artifact_path, + Bytes, HashMap, IntoResponse, Json, NodeArtifact, Path, Query, RequireRunBlob, + RequireRunScoped, RequiredUser, Response, Router, RunArtifactEntry, RunArtifactListResponse, + RunId, StageArtifactEntry, State, StatusCode, WriteBlobResponse, get, header, + octet_stream_response, parse_run_id_path, parse_stage_id_path, post, reject_if_archived, + required_query_param, validate_relative_artifact_path, }; pub(super) fn routes() -> Router> { @@ -38,9 +36,7 @@ pub(super) fn routes() -> Router> { .route("/runs/{id}/artifacts/download", get(download_run_artifacts)) .route( "/runs/{id}/stages/{stageId}/artifacts", - get(list_stage_artifacts) - .post(put_stage_artifact) - .layer(DefaultBodyLimit::disable()), + get(list_stage_artifacts), ) .route( "/runs/{id}/stages/{stageId}/artifacts/download", @@ -56,28 +52,6 @@ struct ArtifactFilenameParams { retry: Option, } -const MAX_SINGLE_ARTIFACT_BYTES: u64 = 10 * 1024 * 1024; -const MAX_MULTIPART_ARTIFACTS: usize = 100; -const MAX_MULTIPART_REQUEST_BYTES: u64 = 50 * 1024 * 1024; -const MAX_MULTIPART_MANIFEST_BYTES: usize = 256 * 1024; - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -struct ArtifactBatchUploadManifest { - entries: Vec, -} - -#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -struct ArtifactBatchUploadEntry { - part: String, - path: String, - #[serde(default, skip_serializing_if = "Option::is_none")] - sha256: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - expected_bytes: Option, - #[serde(default, skip_serializing_if = "Option::is_none")] - content_type: Option, -} - async fn get_checkpoint( _auth: RequiredUser, State(state): State>, @@ -131,24 +105,8 @@ async fn read_run_blob( } } -async fn ensure_run_exists(state: &AppState, run_id: &RunId) -> Result<(), Response> { - match state - .stores - .run_summaries - .get(run_id, chrono::Utc::now()) - .await - { - Ok(Some(_)) => Ok(()), - Ok(None) => Err(ApiError::not_found("Run not found.").into_response()), - Err(err) => { - Err(ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response()) - } - } -} - /// Where an artifact's bytes are: the blob table, for one the run's -/// hooks collected, or the artifact store, for one uploaded to the stage -/// artifact endpoint. +/// hooks collected, or the artifact store, for one written there directly. #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum ArtifactBytes { Blob(BlobHash), @@ -156,9 +114,9 @@ enum ArtifactBytes { } /// Every artifact of the run, each once: the ones the run's projection -/// records, with their bytes in the blob table, and the ones uploaded to -/// the artifact store. A path uploaded for a stage and retry the projection -/// also collected is the projection's. +/// records, with their bytes in the blob table, and the ones written to +/// the artifact store. A path in the store for a stage and retry the +/// projection also collected is the projection's. async fn run_artifacts( state: &AppState, run_id: &RunId, @@ -505,415 +463,6 @@ async fn list_stage_artifacts( } } -enum ArtifactUploadContentType { - OctetStream, - Multipart { boundary: String }, -} - -struct ValidatedArtifactBatchEntry { - path: String, - sha256: Option, - expected_bytes: Option, -} - -#[allow( - clippy::result_large_err, - reason = "Upload content-type parsing returns HTTP client errors directly." -)] -fn artifact_upload_content_type( - headers: &HeaderMap, -) -> Result { - let value = headers - .get(header::CONTENT_TYPE) - .and_then(|value| value.to_str().ok()) - .ok_or_else(|| { - ApiError::new( - StatusCode::UNSUPPORTED_MEDIA_TYPE, - "artifact uploads require a supported Content-Type", - ) - .into_response() - })?; - - let mime = value.split(';').next().unwrap_or(value).trim(); - match mime { - "application/octet-stream" => Ok(ArtifactUploadContentType::OctetStream), - "multipart/form-data" => multer::parse_boundary(value) - .map(|boundary| ArtifactUploadContentType::Multipart { boundary }) - .map_err(|err| bad_request_response(format!("invalid multipart boundary: {err}"))), - _ => Err(ApiError::new( - StatusCode::UNSUPPORTED_MEDIA_TYPE, - "artifact uploads only support application/octet-stream or multipart/form-data", - ) - .into_response()), - } -} - -#[allow( - clippy::result_large_err, - reason = "Content-Length parsing returns HTTP client errors directly." -)] -fn content_length_from_headers(headers: &HeaderMap) -> Result, Response> { - headers - .get(header::CONTENT_LENGTH) - .map(|value| { - value - .to_str() - .map_err(|err| { - bad_request_response(format!("invalid content-length header: {err}")) - }) - .and_then(|value| { - value.parse::().map_err(|err| { - bad_request_response(format!("invalid content-length header: {err}")) - }) - }) - }) - .transpose() -} - -#[allow( - clippy::result_large_err, - reason = "Multipart manifest parsing returns HTTP client errors directly." -)] -async fn read_multipart_manifest( - field: &mut multer::Field<'_>, -) -> Result { - let mut manifest_bytes = Vec::new(); - while let Some(chunk) = field - .chunk() - .await - .map_err(|err| bad_request_response(format!("invalid multipart body: {err}")))? - { - manifest_bytes.extend_from_slice(&chunk); - if manifest_bytes.len() > MAX_MULTIPART_MANIFEST_BYTES { - return Err(payload_too_large_response( - "multipart manifest exceeds the server limit", - )); - } - } - - serde_json::from_slice(&manifest_bytes) - .map_err(|err| bad_request_response(format!("invalid multipart manifest: {err}"))) -} - -#[allow( - clippy::result_large_err, - reason = "Artifact batch validation returns HTTP client errors directly." -)] -fn validate_artifact_batch_manifest( - manifest: ArtifactBatchUploadManifest, -) -> Result, Response> { - if manifest.entries.is_empty() { - return Err(bad_request_response( - "multipart manifest must include at least one artifact entry", - )); - } - if manifest.entries.len() > MAX_MULTIPART_ARTIFACTS { - return Err(payload_too_large_response(format!( - "multipart upload exceeds the {MAX_MULTIPART_ARTIFACTS} artifact limit" - ))); - } - - let mut entries = HashMap::with_capacity(manifest.entries.len()); - let mut seen_paths = HashSet::new(); - let mut expected_total_bytes = 0_u64; - - for entry in manifest.entries { - if entry.part.is_empty() { - return Err(bad_request_response( - "multipart manifest part names must not be empty", - )); - } - if entry.part == "manifest" { - return Err(bad_request_response( - "multipart manifest part name 'manifest' is reserved", - )); - } - let path = validate_relative_artifact_path("manifest path", &entry.path)?; - if !seen_paths.insert(path.clone()) { - return Err(bad_request_response(format!( - "duplicate artifact path in multipart manifest: {path}" - ))); - } - if let Some(sha256) = entry.sha256.as_ref() { - if sha256.len() != 64 || !sha256.bytes().all(|byte| byte.is_ascii_hexdigit()) { - return Err(bad_request_response(format!( - "invalid sha256 for multipart part {}", - entry.part - ))); - } - } - if let Some(expected_bytes) = entry.expected_bytes { - if expected_bytes > MAX_SINGLE_ARTIFACT_BYTES { - return Err(payload_too_large_response(format!( - "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 {MAX_MULTIPART_REQUEST_BYTES} byte limit" - ))); - } - } - if entries - .insert(entry.part.clone(), ValidatedArtifactBatchEntry { - path, - sha256: entry.sha256.map(|value| value.to_ascii_lowercase()), - expected_bytes: entry.expected_bytes, - }) - .is_some() - { - return Err(bad_request_response(format!( - "duplicate multipart part name in manifest: {}", - entry.part - ))); - } - } - - Ok(entries) -} - -async fn upload_stage_artifact_octet_stream( - state: &AppState, - run_id: &RunId, - stage_id: &StageId, - retry: u32, - filename: String, - body: Body, - content_length: Option, -) -> Response { - let relative_path = match validate_relative_artifact_path("filename", &filename) { - Ok(path) => path, - Err(response) => return response, - }; - - if content_length.is_some_and(|length| length > MAX_SINGLE_ARTIFACT_BYTES) { - return payload_too_large_response(format!( - "artifact exceeds the {MAX_SINGLE_ARTIFACT_BYTES} byte limit" - )); - } - - let mut writer = match state.artifact_store.writer( - run_id, - &ArtifactKey::new(stage_id.clone(), retry, relative_path), - ) { - Ok(writer) => writer, - Err(err) => { - return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) - .into_response(); - } - }; - - let mut bytes_written = 0_u64; - let mut data_stream = body.into_data_stream(); - while let Some(chunk) = data_stream.next().await { - let chunk = match chunk - .map_err(|err| bad_request_response(format!("invalid request body: {err}"))) - { - Ok(chunk) => chunk, - Err(response) => return response, - }; - bytes_written = - 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 {MAX_SINGLE_ARTIFACT_BYTES} byte limit" - )); - } - if let Err(err) = writer.write_all(&chunk).await { - return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) - .into_response(); - } - } - - match writer.shutdown().await { - Ok(()) => StatusCode::NO_CONTENT.into_response(), - Err(err) => { - ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()).into_response() - } - } -} - -async fn upload_stage_artifact_multipart( - state: &AppState, - run_id: &RunId, - stage_id: &StageId, - retry: u32, - boundary: String, - body: Body, -) -> Response { - let mut multipart = multer::Multipart::new(body.into_data_stream(), boundary); - let Some(mut manifest_field) = (match multipart - .next_field() - .await - .map_err(|err| bad_request_response(format!("invalid multipart body: {err}"))) - { - Ok(field) => field, - Err(response) => return response, - }) else { - return bad_request_response("multipart upload must begin with a manifest part"); - }; - - if manifest_field.name() != Some("manifest") { - return bad_request_response("multipart upload must begin with a manifest part"); - } - - let manifest = match read_multipart_manifest(&mut manifest_field).await { - Ok(manifest) => manifest, - Err(response) => return response, - }; - drop(manifest_field); - let mut expected_parts = match validate_artifact_batch_manifest(manifest) { - Ok(entries) => entries, - Err(response) => return response, - }; - let mut total_bytes = 0_u64; - - while let Some(mut field) = match multipart - .next_field() - .await - .map_err(|err| bad_request_response(format!("invalid multipart body: {err}"))) - { - Ok(field) => field, - Err(response) => return response, - } { - let Some(part_name) = field.name().map(ToOwned::to_owned) else { - return bad_request_response("multipart file parts must be named"); - }; - let Some(entry) = expected_parts.remove(&part_name) else { - return bad_request_response(format!("unexpected multipart part: {part_name}")); - }; - - let mut writer = match state.artifact_store.writer( - run_id, - &ArtifactKey::new(stage_id.clone(), retry, entry.path.clone()), - ) { - Ok(writer) => writer, - Err(err) => { - return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) - .into_response(); - } - }; - let mut bytes_written = 0_u64; - let mut sha256 = Sha256::new(); - - while let Some(chunk) = match field - .chunk() - .await - .map_err(|err| bad_request_response(format!("invalid multipart body: {err}"))) - { - Ok(chunk) => chunk, - Err(response) => return response, - } { - let chunk_len = u64::try_from(chunk.len()).unwrap_or(u64::MAX); - bytes_written = bytes_written.saturating_add(chunk_len); - total_bytes = total_bytes.saturating_add(chunk_len); - - if bytes_written > MAX_SINGLE_ARTIFACT_BYTES { - return payload_too_large_response(format!( - "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 {MAX_MULTIPART_REQUEST_BYTES} byte limit" - )); - } - - sha256.update(&chunk); - if let Err(err) = writer.write_all(&chunk).await { - return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) - .into_response(); - } - } - - if let Some(expected_bytes) = entry.expected_bytes { - if bytes_written != expected_bytes { - return bad_request_response(format!( - "multipart part {part_name} expected {expected_bytes} bytes but received {bytes_written}" - )); - } - } - if let Some(expected_sha256) = entry.sha256.as_ref() { - let actual_sha256 = hex::encode(sha256.finalize()); - if actual_sha256 != *expected_sha256 { - return bad_request_response(format!( - "multipart part {part_name} sha256 did not match manifest" - )); - } - } - - if let Err(err) = writer.shutdown().await { - return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) - .into_response(); - } - } - - if !expected_parts.is_empty() { - let mut missing = expected_parts.into_keys().collect::>(); - missing.sort(); - return bad_request_response(format!( - "multipart upload is missing part(s): {}", - missing.join(", ") - )); - } - - StatusCode::NO_CONTENT.into_response() -} - -async fn put_stage_artifact( - State(state): State>, - RequireStageArtifact(id, stage_id): RequireStageArtifact, - Query(params): Query, - request: axum_extract::Request, -) -> Response { - let (parts, body) = request.into_parts(); - if let Some(response) = reject_if_archived(state.as_ref(), &id).await { - return response; - } - if let Err(response) = ensure_run_exists(state.as_ref(), &id).await { - return response; - } - let retry = match required_query_param(params.retry.as_ref(), "retry") { - Ok(retry) => retry, - Err(response) => return response, - }; - - let content_length = match content_length_from_headers(&parts.headers) { - Ok(length) => length, - Err(response) => return response, - }; - match artifact_upload_content_type(&parts.headers) { - Ok(ArtifactUploadContentType::OctetStream) => { - let filename = match required_query_param(params.filename.as_ref(), "filename") { - Ok(filename) => filename, - Err(response) => return response, - }; - upload_stage_artifact_octet_stream( - state.as_ref(), - &id, - &stage_id, - retry, - filename, - body, - content_length, - ) - .await - } - 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 {MAX_MULTIPART_REQUEST_BYTES} byte limit" - )); - } - upload_stage_artifact_multipart(state.as_ref(), &id, &stage_id, retry, boundary, body) - .await - } - Err(response) => response, - } -} - async fn get_stage_artifact( _auth: RequiredUser, State(state): State>, diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 56dc76b85..3aaba9f85 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -4913,31 +4913,21 @@ fn named_workflow_dot(name: &str, goal: &str) -> String { ) } -fn multipart_body( - boundary: &str, - manifest: &serde_json::Value, - files: &[(&str, &str, &[u8])], -) -> Body { - let mut body = Vec::new(); - body.extend_from_slice(format!("--{boundary}\r\n").as_bytes()); - body.extend_from_slice(b"Content-Disposition: form-data; name=\"manifest\"\r\n"); - body.extend_from_slice(b"Content-Type: application/json\r\n\r\n"); - body.extend_from_slice(serde_json::to_string(manifest).unwrap().as_bytes()); - body.extend_from_slice(b"\r\n"); - - for (part, filename, bytes) in files { - body.extend_from_slice(format!("--{boundary}\r\n").as_bytes()); - body.extend_from_slice( - format!("Content-Disposition: form-data; name=\"{part}\"; filename=\"{filename}\"\r\n") - .as_bytes(), - ); - body.extend_from_slice(b"Content-Type: application/octet-stream\r\n\r\n"); - body.extend_from_slice(bytes); - body.extend_from_slice(b"\r\n"); - } - - body.extend_from_slice(format!("--{boundary}--\r\n").as_bytes()); - Body::from(body) +/// Write one artifact for a stage the way the hooks do: straight into the +/// artifact store. +async fn seed_stage_artifact( + state: &AppState, + run_id: &str, + stage_id: &str, + retry: u32, + relative_path: &str, + bytes: &[u8], +) { + let run_id = run_id.parse::().unwrap(); + let key = ArtifactKey::new(stage_id.parse::().unwrap(), retry, relative_path); + let mut writer = state.artifact_store.writer(&run_id, &key).unwrap(); + writer.write_all(bytes).await.unwrap(); + writer.shutdown().await.unwrap(); } /// Create a run via POST /runs, then start it via POST /runs/{id}/start. @@ -7233,17 +7223,7 @@ async fn stage_artifacts_round_trip() { let run_id = create_run(&app, MINIMAL_DOT).await; let stage_id = "code@2"; - - let req = Request::builder() - .method("POST") - .uri(api(&format!( - "/runs/{run_id}/stages/{stage_id}/artifacts?filename=src/lib.rs&retry=1" - ))) - .header("content-type", "application/octet-stream") - .body(Body::from("fn main() {}")) - .unwrap(); - let response = app.clone().oneshot(req).await.unwrap(); - assert_status!(response, StatusCode::NO_CONTENT).await; + seed_stage_artifact(&state, &run_id, stage_id, 1, "src/lib.rs", b"fn main() {}").await; let req = Request::builder() .method("GET") @@ -7287,16 +7267,15 @@ async fn stage_artifacts_keep_same_filename_per_retry() { let stage_id = "code@2"; for (retry, body) in [(1, "first"), (2, "second")] { - let req = Request::builder() - .method("POST") - .uri(api(&format!( - "/runs/{run_id}/stages/{stage_id}/artifacts?filename=logs/output.txt&retry={retry}" - ))) - .header("content-type", "application/octet-stream") - .body(Body::from(body)) - .unwrap(); - let response = app.clone().oneshot(req).await.unwrap(); - assert_status!(response, StatusCode::NO_CONTENT).await; + seed_stage_artifact( + &state, + &run_id, + stage_id, + retry, + "logs/output.txt", + body.as_bytes(), + ) + .await; } let req = Request::builder() @@ -7391,25 +7370,6 @@ async fn create_run_keeps_missing_project_and_workflow_names_absent() { assert_eq!(run_state.spec.graph_name(), Some("Demo")); } -#[tokio::test] -async fn stage_artifact_upload_rejects_invalid_filename() { - let state = test_app_state(); - let app = crate::test_support::build_test_router(Arc::clone(&state)); - - let run_id = create_run(&app, MINIMAL_DOT).await; - - let req = Request::builder() - .method("POST") - .uri(api(&format!( - "/runs/{run_id}/stages/code@2/artifacts?filename=../escape.txt&retry=1" - ))) - .header("content-type", "application/octet-stream") - .body(Body::from("nope")) - .unwrap(); - let response = app.oneshot(req).await.unwrap(); - assert_status!(response, StatusCode::BAD_REQUEST).await; -} - #[tokio::test] async fn worker_token_accepts_run_scoped_routes_and_falls_back_to_user_jwt() { let (state, app) = jwt_auth_app(); @@ -7759,85 +7719,6 @@ async fn base_worker_token_is_rejected_by_run_tool_only_routes() { } } -#[tokio::test] -async fn worker_token_controls_stage_artifact_route() { - let (_state, app) = jwt_auth_app(); - let user_jwt = issue_test_user_jwt(); - let run_id = create_run_with_bearer(&app, &user_jwt).await; - let worker_token = issue_test_worker_token(&run_id); - let other_run_id = create_run_with_bearer(&app, &user_jwt).await; - let mismatched_worker_token = issue_test_worker_token(&other_run_id); - - let response = app - .clone() - .oneshot( - Request::builder() - .method(Method::POST) - .uri(api(&format!( - "/runs/{run_id}/stages/code@2/artifacts?filename=artifact.txt&retry=1" - ))) - .header(header::AUTHORIZATION, format!("Bearer {worker_token}")) - .header(header::CONTENT_TYPE, "application/octet-stream") - .body(Body::from("artifact")) - .unwrap(), - ) - .await - .unwrap(); - assert_status!(response, StatusCode::NO_CONTENT).await; - - let response = app - .clone() - .oneshot( - Request::builder() - .method(Method::POST) - .uri(api(&format!( - "/runs/{run_id}/stages/code@2/artifacts?filename=artifact.txt&retry=1" - ))) - .header(header::AUTHORIZATION, format!("Bearer {user_jwt}")) - .header(header::CONTENT_TYPE, "application/octet-stream") - .body(Body::from("artifact")) - .unwrap(), - ) - .await - .unwrap(); - assert_status!(response, StatusCode::NO_CONTENT).await; - - let response = app - .clone() - .oneshot( - Request::builder() - .method(Method::POST) - .uri(api(&format!( - "/runs/{run_id}/stages/code@2/artifacts?filename=artifact.txt&retry=1" - ))) - .header( - header::AUTHORIZATION, - format!("Bearer {mismatched_worker_token}"), - ) - .header(header::CONTENT_TYPE, "application/octet-stream") - .body(Body::from("artifact")) - .unwrap(), - ) - .await - .unwrap(); - assert_status!(response, StatusCode::FORBIDDEN).await; - - let response = app - .oneshot( - Request::builder() - .method(Method::POST) - .uri(api(&format!( - "/runs/{run_id}/stages/code@2/artifacts?filename=artifact.txt&retry=1" - ))) - .header(header::CONTENT_TYPE, "application/octet-stream") - .body(Body::from("artifact")) - .unwrap(), - ) - .await - .unwrap(); - assert_status!(response, StatusCode::UNAUTHORIZED).await; -} - #[tokio::test] async fn worker_token_is_rejected_on_user_only_routes() { let (_state, app) = jwt_auth_app(); @@ -7916,104 +7797,6 @@ async fn worker_token_is_rejected_on_user_only_routes() { assert_ne!(response.status(), StatusCode::UNAUTHORIZED); } -#[tokio::test] -async fn stage_artifacts_multipart_round_trip() { - let state = test_app_state(); - let app = crate::test_support::build_test_router(Arc::clone(&state)); - - let run_id = create_run(&app, MINIMAL_DOT).await; - let stage_id = "code@2"; - let source_bytes = b"fn main() {}\n"; - let log_bytes = b"build ok\n"; - let manifest = serde_json::json!({ - "entries": [ - { - "part": "file1", - "path": "src/lib.rs", - "sha256": hex::encode(Sha256::digest(source_bytes)), - "expected_bytes": source_bytes.len(), - "content_type": "text/plain" - }, - { - "part": "file2", - "path": "logs/output.txt", - "sha256": hex::encode(Sha256::digest(log_bytes)), - "expected_bytes": log_bytes.len(), - "content_type": "text/plain" - } - ] - }); - let boundary = "fabro-test-boundary"; - - let req = Request::builder() - .method("POST") - .uri(api(&format!( - "/runs/{run_id}/stages/{stage_id}/artifacts?retry=1" - ))) - .header( - "content-type", - format!("multipart/form-data; boundary={boundary}"), - ) - .body(multipart_body(boundary, &manifest, &[ - ("file1", "src/lib.rs", source_bytes), - ("file2", "logs/output.txt", log_bytes), - ])) - .unwrap(); - let response = app.clone().oneshot(req).await.unwrap(); - assert_status!(response, StatusCode::NO_CONTENT).await; - - let req = Request::builder() - .method("GET") - .uri(api(&format!("/runs/{run_id}/stages/{stage_id}/artifacts"))) - .body(Body::empty()) - .unwrap(); - let response = app.clone().oneshot(req).await.unwrap(); - let body = response_json!(response, StatusCode::OK).await; - assert_eq!(body["data"][0]["filename"], "logs/output.txt"); - assert_eq!(body["data"][0]["retry"], 1); - assert_eq!(body["data"][0]["size"], log_bytes.len()); - assert_eq!(body["data"][1]["filename"], "src/lib.rs"); - assert_eq!(body["data"][1]["retry"], 1); - assert_eq!(body["data"][1]["size"], source_bytes.len()); - - let req = Request::builder() - .method("GET") - .uri(api(&format!( - "/runs/{run_id}/stages/{stage_id}/artifacts/download?filename=logs/output.txt&retry=1" - ))) - .body(Body::empty()) - .unwrap(); - let response = app.oneshot(req).await.unwrap(); - let bytes = response_bytes!(response, StatusCode::OK).await; - assert_eq!(&bytes[..], log_bytes); -} - -#[tokio::test] -async fn stage_artifacts_multipart_requires_manifest_first() { - let state = test_app_state(); - let app = crate::test_support::build_test_router(Arc::clone(&state)); - - let run_id = create_run(&app, MINIMAL_DOT).await; - let boundary = "fabro-test-boundary"; - let body = format!( - "--{boundary}\r\nContent-Disposition: form-data; name=\"file1\"; filename=\"src/lib.rs\"\r\n\r\nfn main() {{}}\r\n--{boundary}\r\nContent-Disposition: form-data; name=\"manifest\"\r\nContent-Type: application/json\r\n\r\n{{\"entries\":[{{\"part\":\"file1\",\"path\":\"src/lib.rs\"}}]}}\r\n--{boundary}--\r\n" - ); - - let req = Request::builder() - .method("POST") - .uri(api(&format!( - "/runs/{run_id}/stages/code@2/artifacts?retry=1" - ))) - .header( - "content-type", - format!("multipart/form-data; boundary={boundary}"), - ) - .body(Body::from(body)) - .unwrap(); - let response = app.oneshot(req).await.unwrap(); - assert_status!(response, StatusCode::BAD_REQUEST).await; -} - #[tokio::test] async fn create_run_accepts_explicit_title() { let state = test_app_state(); diff --git a/lib/foundation/fabro-client/src/client.rs b/lib/foundation/fabro-client/src/client.rs index 3e28a0f34..6f661d017 100644 --- a/lib/foundation/fabro-client/src/client.rs +++ b/lib/foundation/fabro-client/src/client.rs @@ -1,7 +1,6 @@ use std::collections::VecDeque; use std::future::Future; use std::num::NonZeroU64; -use std::path::Path; use std::pin::Pin; use std::sync::{Arc, RwLock}; @@ -9,25 +8,22 @@ use anyhow::{Context as _, Result, anyhow, bail}; use bytes::Bytes; use fabro_api::types; use fabro_api::types::RunControlAcknowledgement; -use fabro_http::header::{ACCEPT, AUTHORIZATION, CONTENT_LENGTH, CONTENT_TYPE}; -use fabro_http::multipart::{Form, Part}; +use fabro_http::header::{ACCEPT, AUTHORIZATION}; use fabro_types::settings::run::MergeStrategy; use fabro_types::{ - ArtifactUpload, BlobHash, Model, ModelTestMode, PairId, PairMessageRecord, PairMessageRequest, - PairRecord, PairStartRequest, PairTranscriptResponse, Run, RunId, RunPairStatusResponse, - RunProjection, RunSessionMetadata, RunStreamItem, SessionEvent, SessionId, StageId, - WorkflowVersion, WorkflowVersionId, + BlobHash, Model, ModelTestMode, PairId, PairMessageRecord, PairMessageRequest, PairRecord, + PairStartRequest, PairTranscriptResponse, Run, RunId, RunPairStatusResponse, RunProjection, + RunSessionMetadata, RunStreamItem, SessionEvent, SessionId, StageId, WorkflowVersion, + WorkflowVersionId, }; use fabro_util::exit::{ErrorExt, ExitClass}; use futures::future::BoxFuture; use futures::{Stream, StreamExt}; use lithos_llm::catalog::ProviderId; use lithos_llm::types::ReasoningEffort; -use serde::{Deserialize, Serialize}; -use tokio::fs::File; +use serde::Deserialize; use tokio::sync::Mutex; use tokio::time; -use tokio_util::io::ReaderStream; use crate::credential::Credential; use crate::error::{ @@ -150,23 +146,6 @@ struct OAuthErrorBody { error_description: Option, } -#[derive(Debug, Serialize)] -struct ArtifactBatchUploadManifest { - entries: Vec, -} - -#[derive(Debug, Serialize)] -struct ArtifactBatchUploadEntry { - part: String, - path: String, - #[serde(skip_serializing_if = "Option::is_none")] - sha256: Option, - #[serde(skip_serializing_if = "Option::is_none")] - expected_bytes: Option, - #[serde(skip_serializing_if = "Option::is_none")] - content_type: Option, -} - impl RunStreamItemStream { #[must_use] pub fn new(stream: progenitor_client::ByteStream) -> Self { @@ -2148,142 +2127,6 @@ impl Client { Ok(bytes) } - #[expect( - clippy::disallowed_types, - reason = "Client builds raw server API request URLs for wire transit; logging redaction is handled at log boundaries." - )] - fn stage_artifacts_url( - &self, - run_id: &RunId, - stage_id: &StageId, - retry: u32, - ) -> Result { - let base_url = self.base_url(); - let mut url = fabro_http::Url::parse(&base_url) - .with_context(|| format!("invalid server base URL {base_url}"))?; - url.path_segments_mut() - .map_err(|()| anyhow!("server base URL cannot accept path segments"))? - .extend([ - "api", - "v1", - "runs", - &run_id.to_string(), - "stages", - &stage_id.to_string(), - "artifacts", - ]); - url.query_pairs_mut() - .append_pair("retry", &retry.to_string()); - Ok(url) - } - - pub async fn upload_stage_artifact_file( - &self, - run_id: &RunId, - stage_id: &StageId, - retry: u32, - filename: &str, - path: &Path, - bearer_token: &str, - ) -> Result<()> { - let mut url = self.stage_artifacts_url(run_id, stage_id, retry)?; - url.query_pairs_mut().append_pair("filename", filename); - - let file = File::open(path) - .await - .with_context(|| format!("failed to open artifact {}", path.display()))?; - let content_length = file - .metadata() - .await - .with_context(|| format!("failed to stat artifact {}", path.display()))? - .len(); - let body = fabro_http::Body::wrap_stream(ReaderStream::new(file)); - - let response = self - .current_state() - .http_client - .post(url) - .bearer_auth(bearer_token) - .header(CONTENT_TYPE, "application/octet-stream") - .header(CONTENT_LENGTH, content_length.to_string()) - .body(body) - .send() - .await - .with_context(|| format!("failed to upload artifact {}", path.display()))?; - classify_http_response(response) - .await? - .map(|_| ()) - .map_err(|failure| raw_response_failure_error(&failure)) - } - - pub async fn upload_stage_artifact_batch( - &self, - run_id: &RunId, - stage_id: &StageId, - retry: u32, - artifact_capture_dir: &Path, - artifacts: &[ArtifactUpload], - bearer_token: &str, - ) -> Result<()> { - let url = self.stage_artifacts_url(run_id, stage_id, retry)?; - let mut manifest_entries = Vec::with_capacity(artifacts.len()); - let mut file_parts = Vec::with_capacity(artifacts.len()); - - for (index, artifact) in artifacts.iter().enumerate() { - let part_name = format!("file{}", index + 1); - let path = artifact_capture_dir.join(&artifact.path); - let file = File::open(&path) - .await - .with_context(|| format!("failed to open artifact {}", path.display()))?; - let content_length = file - .metadata() - .await - .with_context(|| format!("failed to stat artifact {}", path.display()))? - .len(); - - manifest_entries.push(ArtifactBatchUploadEntry { - part: part_name.clone(), - path: artifact.path.clone(), - sha256: Some(artifact.content_sha256.clone()), - expected_bytes: Some(artifact.bytes), - content_type: Some(artifact.mime.clone()), - }); - - file_parts.push(( - part_name, - Part::stream_with_length( - fabro_http::Body::wrap_stream(ReaderStream::new(file)), - content_length, - ) - .file_name(artifact.path.clone()), - )); - } - - let manifest = ArtifactBatchUploadManifest { - entries: manifest_entries, - }; - 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); - } - - let response = self - .current_state() - .http_client - .post(url) - .bearer_auth(bearer_token) - .multipart(form) - .send() - .await - .context("failed to upload artifact batch")?; - classify_http_response(response) - .await? - .map(|_| ()) - .map_err(|failure| raw_response_failure_error(&failure)) - } - pub async fn generate_preview_url( &self, run_id: &RunId, diff --git a/lib/foundation/fabro-types/src/artifact.rs b/lib/foundation/fabro-types/src/artifact.rs deleted file mode 100644 index 0cd6af42a..000000000 --- a/lib/foundation/fabro-types/src/artifact.rs +++ /dev/null @@ -1,30 +0,0 @@ -use serde::{Deserialize, Serialize}; - -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -pub struct ArtifactUpload { - pub path: String, - pub mime: String, - pub content_md5: String, - pub content_sha256: String, - pub bytes: u64, -} - -#[cfg(test)] -mod tests { - use super::ArtifactUpload; - - #[test] - fn round_trips_through_serde_json() { - let artifact = ArtifactUpload { - path: "artifacts/log.txt".to_string(), - mime: "text/plain".to_string(), - content_md5: "md5".to_string(), - content_sha256: "sha256".to_string(), - bytes: 42, - }; - - let value = serde_json::to_value(&artifact).unwrap(); - let parsed: ArtifactUpload = serde_json::from_value(value).unwrap(); - assert_eq!(parsed, artifact); - } -} diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index 8a88e3af3..a7f703d44 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -1,7 +1,6 @@ extern crate self as fabro_types; pub mod agent_props; -pub mod artifact; pub mod auth; pub mod blob_hash; pub mod blob_ref; @@ -69,7 +68,6 @@ pub use agent_props::{ AgentEventProps, AgentSessionActivatedProps, AgentToolsAvailableProps, CODING_EVENT_NAMES, SessionCapability, StagePromptProps, coding_event_name, is_coding_event_name, }; -pub use artifact::ArtifactUpload; pub use auth::{IdpIdentity, IdpIdentityError}; pub use blob_hash::BlobHash; pub use blob_ref::{ diff --git a/lib/packages/fabro-api-client/src/.openapi-generator/FILES b/lib/packages/fabro-api-client/src/.openapi-generator/FILES index 85099e35b..3546fa7d0 100644 --- a/lib/packages/fabro-api-client/src/.openapi-generator/FILES +++ b/lib/packages/fabro-api-client/src/.openapi-generator/FILES @@ -55,8 +55,6 @@ models/aggregate-usage-totals.ts models/aggregate-usage.ts models/api-question.ts models/approval-mode.ts -models/artifact-batch-upload-entry.ts -models/artifact-batch-upload-manifest.ts models/artifact-entry.ts models/artifact-list-response.ts models/artifacts-settings.ts diff --git a/lib/packages/fabro-api-client/src/api/run-internals-api.ts b/lib/packages/fabro-api-client/src/api/run-internals-api.ts index e99a16e04..09e391eaf 100644 --- a/lib/packages/fabro-api-client/src/api/run-internals-api.ts +++ b/lib/packages/fabro-api-client/src/api/run-internals-api.ts @@ -59,8 +59,6 @@ import type { StageContextWindow } from '../models'; import type { WorkflowSettings } from '../models'; // @ts-ignore import type { WriteBlobResponse } from '../models'; -// @ts-ignore -import type { WriteRunBlobRequest } from '../models'; /** * RunInternalsApi - axios parameter creator */ @@ -799,67 +797,6 @@ export const RunInternalsApiAxiosParamCreator = function (configuration?: Config options: localVarRequestOptions, }; }, - /** - * Uploads one or more artifacts for a stage. Intended for trusted internal callers. The server accepts both: - `application/octet-stream` for single-file uploads with the `filename` query parameter - strict manifest-first `multipart/form-data` uploads documented by `ArtifactBatchUploadManifest` The generated Rust client currently exposes the octet-stream variant because the OpenAPI code generator in this repo does not support multiple request media types on one operation. - * @summary Put Stage Artifact - * @param {string} id Unique run identifier (ULID). - * @param {string} stageId Identifier of a stage within a run\'s workflow graph, serialized as `node_id@visit`. - * @param {number} retry Retry attempt number for the artifact. - * @param {File} body - * @param {string} [filename] Relative artifact path for `application/octet-stream` uploads. Ignored for multipart uploads. - * @param {*} [options] Override http request option. - * @throws {RequiredError} - */ - putStageArtifact: async (id: string, stageId: string, retry: number, body: File, filename?: string, options: RawAxiosRequestConfig = {}): Promise => { - // verify required parameter 'id' is not null or undefined - assertParamExists('putStageArtifact', 'id', id) - // verify required parameter 'stageId' is not null or undefined - assertParamExists('putStageArtifact', 'stageId', stageId) - // verify required parameter 'retry' is not null or undefined - assertParamExists('putStageArtifact', 'retry', retry) - // verify required parameter 'body' is not null or undefined - assertParamExists('putStageArtifact', 'body', body) - const localVarPath = `/api/v1/runs/{id}/stages/{stageId}/artifacts` - .replace(`{${"id"}}`, encodeURIComponent(String(id))) - .replace(`{${"stageId"}}`, encodeURIComponent(String(stageId))); - // use dummy base URL string because the URL constructor only accepts absolute URLs. - const localVarUrlObj = new URL(localVarPath, DUMMY_BASE_URL); - let baseOptions; - if (configuration) { - baseOptions = configuration.baseOptions; - } - - const localVarRequestOptions = { method: 'POST', ...baseOptions, ...options}; - const localVarHeaderParameter = {} as any; - const localVarQueryParameter = {} as any; - - // authentication SessionCookie required - - // authentication BearerAuth required - // http bearer authentication required - await setBearerAuthToObject(localVarHeaderParameter, configuration) - - if (retry !== undefined) { - localVarQueryParameter['retry'] = retry; - } - - if (filename !== undefined) { - localVarQueryParameter['filename'] = filename; - } - - localVarHeaderParameter['Content-Type'] = 'application/octet-stream'; - localVarHeaderParameter['Accept'] = 'application/json'; - - setSearchParams(localVarUrlObj, localVarQueryParameter); - let headersFromBaseOptions = baseOptions && baseOptions.headers ? baseOptions.headers : {}; - localVarRequestOptions.headers = {...localVarHeaderParameter, ...headersFromBaseOptions, ...options.headers}; - localVarRequestOptions.data = serializeDataIfNeeded(body, localVarRequestOptions, configuration) - - return { - url: toPathString(localVarUrlObj), - options: localVarRequestOptions, - }; - }, /** * The blob with this digest, if the store holds one. * @summary Read Petri Blob @@ -1405,23 +1342,6 @@ export const RunInternalsApiFp = function(configuration?: Configuration) { const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.openPetriRun']?.[localVarOperationServerIndex]?.url; return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); }, - /** - * Uploads one or more artifacts for a stage. Intended for trusted internal callers. The server accepts both: - `application/octet-stream` for single-file uploads with the `filename` query parameter - strict manifest-first `multipart/form-data` uploads documented by `ArtifactBatchUploadManifest` The generated Rust client currently exposes the octet-stream variant because the OpenAPI code generator in this repo does not support multiple request media types on one operation. - * @summary Put Stage Artifact - * @param {string} id Unique run identifier (ULID). - * @param {string} stageId Identifier of a stage within a run\'s workflow graph, serialized as `node_id@visit`. - * @param {number} retry Retry attempt number for the artifact. - * @param {File} body - * @param {string} [filename] Relative artifact path for `application/octet-stream` uploads. Ignored for multipart uploads. - * @param {*} [options] Override http request option. - * @throws {RequiredError} - */ - async putStageArtifact(id: string, stageId: string, retry: number, body: File, filename?: string, options?: RawAxiosRequestConfig): Promise<(axios?: AxiosInstance, basePath?: string) => AxiosPromise> { - const localVarAxiosArgs = await localVarAxiosParamCreator.putStageArtifact(id, stageId, retry, body, filename, options); - const localVarOperationServerIndex = configuration?.serverIndex ?? 0; - const localVarOperationServerBasePath = operationServerMap['RunInternalsApi.putStageArtifact']?.[localVarOperationServerIndex]?.url; - return (axios, basePath) => createRequestFunction(localVarAxiosArgs, globalAxios, BASE_PATH, configuration)(axios, localVarOperationServerBasePath || basePath); - }, /** * The blob with this digest, if the store holds one. * @summary Read Petri Blob @@ -1707,20 +1627,6 @@ export const RunInternalsApiFactory = function (configuration?: Configuration, b openPetriRun(id: string, petriOpenRequest: PetriOpenRequest, options?: RawAxiosRequestConfig): AxiosPromise { return localVarFp.openPetriRun(id, petriOpenRequest, options).then((request) => request(axios, basePath)); }, - /** - * Uploads one or more artifacts for a stage. Intended for trusted internal callers. The server accepts both: - `application/octet-stream` for single-file uploads with the `filename` query parameter - strict manifest-first `multipart/form-data` uploads documented by `ArtifactBatchUploadManifest` The generated Rust client currently exposes the octet-stream variant because the OpenAPI code generator in this repo does not support multiple request media types on one operation. - * @summary Put Stage Artifact - * @param {string} id Unique run identifier (ULID). - * @param {string} stageId Identifier of a stage within a run\'s workflow graph, serialized as `node_id@visit`. - * @param {number} retry Retry attempt number for the artifact. - * @param {File} body - * @param {string} [filename] Relative artifact path for `application/octet-stream` uploads. Ignored for multipart uploads. - * @param {*} [options] Override http request option. - * @throws {RequiredError} - */ - putStageArtifact(id: string, stageId: string, retry: number, body: File, filename?: string, options?: RawAxiosRequestConfig): AxiosPromise { - return localVarFp.putStageArtifact(id, stageId, retry, body, filename, options).then((request) => request(axios, basePath)); - }, /** * The blob with this digest, if the store holds one. * @summary Read Petri Blob @@ -1999,21 +1905,6 @@ export class RunInternalsApi extends BaseAPI { return RunInternalsApiFp(this.configuration).openPetriRun(id, petriOpenRequest, options).then((request) => request(this.axios, this.basePath)); } - /** - * Uploads one or more artifacts for a stage. Intended for trusted internal callers. The server accepts both: - `application/octet-stream` for single-file uploads with the `filename` query parameter - strict manifest-first `multipart/form-data` uploads documented by `ArtifactBatchUploadManifest` The generated Rust client currently exposes the octet-stream variant because the OpenAPI code generator in this repo does not support multiple request media types on one operation. - * @summary Put Stage Artifact - * @param {string} id Unique run identifier (ULID). - * @param {string} stageId Identifier of a stage within a run\'s workflow graph, serialized as `node_id@visit`. - * @param {number} retry Retry attempt number for the artifact. - * @param {File} body - * @param {string} [filename] Relative artifact path for `application/octet-stream` uploads. Ignored for multipart uploads. - * @param {*} [options] Override http request option. - * @throws {RequiredError} - */ - public putStageArtifact(id: string, stageId: string, retry: number, body: File, filename?: string, options?: RawAxiosRequestConfig) { - return RunInternalsApiFp(this.configuration).putStageArtifact(id, stageId, retry, body, filename, options).then((request) => request(this.axios, this.basePath)); - } - /** * The blob with this digest, if the store holds one. * @summary Read Petri Blob diff --git a/lib/packages/fabro-api-client/src/models/artifact-batch-upload-entry.ts b/lib/packages/fabro-api-client/src/models/artifact-batch-upload-entry.ts deleted file mode 100644 index 160e12f90..000000000 --- a/lib/packages/fabro-api-client/src/models/artifact-batch-upload-entry.ts +++ /dev/null @@ -1,41 +0,0 @@ -/* tslint:disable */ -/* eslint-disable */ -/** - * Fabro Run API - * HTTP API for managing Fabro workflow run executions. - * - * The version of the OpenAPI document: 0.2.0 - * - * - * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). - * https://openapi-generator.tech - * Do not edit the class manually. - */ - - - -/** - * One file entry in a strict multipart artifact upload manifest. - */ -export interface ArtifactBatchUploadEntry { - /** - * Multipart field name for the file part. - */ - 'part': string; - /** - * Relative artifact path to store. - */ - 'path': string; - /** - * Optional SHA-256 checksum for the file contents; hex input is case-insensitive. - */ - 'sha256'?: string | null; - /** - * Optional exact byte length expected for the file part. - */ - 'expected_bytes'?: number | null; - /** - * Optional client-supplied content type for the file part. - */ - 'content_type'?: string | null; -} diff --git a/lib/packages/fabro-api-client/src/models/artifact-batch-upload-manifest.ts b/lib/packages/fabro-api-client/src/models/artifact-batch-upload-manifest.ts deleted file mode 100644 index 080a3e575..000000000 --- a/lib/packages/fabro-api-client/src/models/artifact-batch-upload-manifest.ts +++ /dev/null @@ -1,25 +0,0 @@ -/* tslint:disable */ -/* eslint-disable */ -/** - * Fabro Run API - * HTTP API for managing Fabro workflow run executions. - * - * The version of the OpenAPI document: 0.2.0 - * - * - * NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech). - * https://openapi-generator.tech - * Do not edit the class manually. - */ - - -// May contain unused imports in some cases -// @ts-ignore -import type { ArtifactBatchUploadEntry } from './artifact-batch-upload-entry'; - -/** - * Manifest for strict multipart artifact uploads. - */ -export interface ArtifactBatchUploadManifest { - 'entries': Array; -} diff --git a/lib/packages/fabro-api-client/src/models/index.ts b/lib/packages/fabro-api-client/src/models/index.ts index 6156dbf68..53baa3c13 100644 --- a/lib/packages/fabro-api-client/src/models/index.ts +++ b/lib/packages/fabro-api-client/src/models/index.ts @@ -26,8 +26,6 @@ export * from './aggregate-usage'; export * from './aggregate-usage-totals'; export * from './api-question'; export * from './approval-mode'; -export * from './artifact-batch-upload-entry'; -export * from './artifact-batch-upload-manifest'; export * from './artifact-entry'; export * from './artifact-list-response'; export * from './artifacts-settings';