feat(onboarding): PLC/blob file upload (multipart) + single-file ingest fix (#163)
CI / Check (push) Has been skipped
CI / Detect Changes (push) Successful in 3s
CI / Deploy Agent (push) Successful in 5m39s
CI / Deploy Dashboard (push) Successful in 4m7s
CI / Deploy Docs (push) Has been skipped
CI / Deploy MCP (push) Successful in 2m38s
CI / Check (push) Has been skipped
CI / Detect Changes (push) Successful in 3s
CI / Deploy Agent (push) Successful in 5m39s
CI / Deploy Dashboard (push) Successful in 4m7s
CI / Deploy Docs (push) Has been skipped
CI / Deploy MCP (push) Successful in 2m38s
This commit was merged in pull request #163.
This commit is contained in:
@@ -34,7 +34,7 @@ hex = { workspace = true }
|
||||
uuid = { workspace = true }
|
||||
secrecy = { workspace = true }
|
||||
regex = { workspace = true }
|
||||
axum = "0.8"
|
||||
axum = { version = "0.8", features = ["multipart"] }
|
||||
tower-http = { version = "0.6", features = ["cors", "trace", "set-header"] }
|
||||
git2 = "0.20"
|
||||
octocrab = "0.44"
|
||||
@@ -65,5 +65,5 @@ tokio = { workspace = true }
|
||||
mongodb = { workspace = true }
|
||||
uuid = { workspace = true }
|
||||
secrecy = { workspace = true }
|
||||
axum = "0.8"
|
||||
axum = { version = "0.8", features = ["multipart"] }
|
||||
tower-http = { version = "0.6", features = ["cors"] }
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
|
||||
use axum::extract::{Extension, Path, Query};
|
||||
use axum::extract::{Extension, Multipart, Path, Query};
|
||||
use axum::http::StatusCode;
|
||||
use axum::Json;
|
||||
use mongodb::bson::{doc, oid::ObjectId, to_bson};
|
||||
@@ -377,6 +377,116 @@ pub async fn add_artifact(
|
||||
get_target(Extension(agent), tenant, Path(id)).await
|
||||
}
|
||||
|
||||
/// POST /api/v1/targets/{id}/artifacts/upload — attach an artifact by uploading
|
||||
/// its file (PLC project, firmware image, source archive, mobile package). The
|
||||
/// bytes are written to the artifact blob store and referenced by `stored_path`,
|
||||
/// so ingest resolves them locally (no URL fetch).
|
||||
///
|
||||
/// Multipart fields: `file` (required), `kind` (required, snake_case
|
||||
/// `ArtifactKind`), `plc_format` (optional, for PLC projects).
|
||||
#[tracing::instrument(skip_all, fields(target_id = %id))]
|
||||
pub async fn upload_artifact(
|
||||
Extension(agent): AgentExt,
|
||||
tenant: TenantCtx,
|
||||
Path(id): Path<String>,
|
||||
mut multipart: Multipart,
|
||||
) -> Result<Json<ApiResponse<OnboardedTarget>>, StatusCode> {
|
||||
let oid = parse_oid(&id)?;
|
||||
let db = tenant_db(&agent, &tenant).await?;
|
||||
if db
|
||||
.onboarded_targets()
|
||||
.find_one(doc! { "_id": oid })
|
||||
.await
|
||||
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?
|
||||
.is_none()
|
||||
{
|
||||
return Err(StatusCode::NOT_FOUND);
|
||||
}
|
||||
|
||||
let mut kind: Option<ArtifactKind> = None;
|
||||
let mut plc_format: Option<PlcFormat> = None;
|
||||
let mut filename = String::from("upload.bin");
|
||||
let mut bytes: Option<axum::body::Bytes> = None;
|
||||
|
||||
while let Some(field) = multipart
|
||||
.next_field()
|
||||
.await
|
||||
.map_err(|_| StatusCode::BAD_REQUEST)?
|
||||
{
|
||||
match field.name().unwrap_or("") {
|
||||
"kind" => {
|
||||
let v = field.text().await.map_err(|_| StatusCode::BAD_REQUEST)?;
|
||||
kind = parse_enum(&v);
|
||||
}
|
||||
"plc_format" => {
|
||||
let v = field.text().await.map_err(|_| StatusCode::BAD_REQUEST)?;
|
||||
plc_format = parse_enum(&v);
|
||||
}
|
||||
"file" => {
|
||||
if let Some(fname) = field.file_name() {
|
||||
filename = fname.to_string();
|
||||
}
|
||||
bytes = Some(field.bytes().await.map_err(|_| StatusCode::BAD_REQUEST)?);
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
let (Some(kind), Some(bytes)) = (kind, bytes) else {
|
||||
return Err(StatusCode::BAD_REQUEST);
|
||||
};
|
||||
|
||||
// Store the uploaded bytes under the artifact blob store.
|
||||
let safe_name: String = filename
|
||||
.chars()
|
||||
.map(|c| {
|
||||
if c.is_ascii_alphanumeric() || matches!(c, '.' | '-' | '_') {
|
||||
c
|
||||
} else {
|
||||
'_'
|
||||
}
|
||||
})
|
||||
.collect();
|
||||
let dir = std::path::Path::new(&agent.config.artifact_store_base_path)
|
||||
.join("uploads")
|
||||
.join(&id);
|
||||
std::fs::create_dir_all(&dir).map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;
|
||||
let dest = dir.join(format!("{}_{safe_name}", uuid::Uuid::new_v4()));
|
||||
std::fs::write(&dest, bytes.as_ref()).map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;
|
||||
|
||||
// Build the artifact for this kind, referencing the stored file.
|
||||
let mut artifact = match kind {
|
||||
ArtifactKind::PlcProject => Artifact::plc_project(
|
||||
filename.clone(),
|
||||
plc_format.unwrap_or(PlcFormat::PlcopenXml),
|
||||
),
|
||||
ArtifactKind::FirmwareImage => Artifact::firmware_image(filename.clone()),
|
||||
ArtifactKind::SourceArchive => Artifact::source_archive(filename.clone()),
|
||||
ArtifactKind::MobilePackage => Artifact::mobile_package(filename.clone()),
|
||||
// Non-file kinds (git repo, live URL, container ref, text) use the JSON
|
||||
// add-artifact endpoint, not upload.
|
||||
_ => return Err(StatusCode::BAD_REQUEST),
|
||||
};
|
||||
artifact.stored_path = Some(dest.to_string_lossy().to_string());
|
||||
artifact.size_bytes = Some(bytes.len() as u64);
|
||||
|
||||
let artifact_bson = to_bson(&artifact).map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;
|
||||
db.onboarded_targets()
|
||||
.update_one(
|
||||
doc! { "_id": oid },
|
||||
doc! { "$push": { "artifacts": artifact_bson }, "$set": { "updated_at": mongodb::bson::DateTime::now() } },
|
||||
)
|
||||
.await
|
||||
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;
|
||||
|
||||
get_target(Extension(agent), tenant, Path(id)).await
|
||||
}
|
||||
|
||||
/// Deserialize a snake_case enum value from a plain string.
|
||||
fn parse_enum<T: for<'de> Deserialize<'de>>(s: &str) -> Option<T> {
|
||||
serde_json::from_value(serde_json::Value::String(s.to_string())).ok()
|
||||
}
|
||||
|
||||
/// GET /api/v1/targets/{id}/applicable-scans — the scan-applicability matrix.
|
||||
#[tracing::instrument(skip_all, fields(target_id = %id))]
|
||||
pub async fn applicable_scans_for_target(
|
||||
|
||||
@@ -26,6 +26,10 @@ pub fn build_router() -> Router {
|
||||
"/api/v1/targets/{id}/artifacts",
|
||||
post(handlers::onboarding::add_artifact),
|
||||
)
|
||||
.route(
|
||||
"/api/v1/targets/{id}/artifacts/upload",
|
||||
post(handlers::onboarding::upload_artifact),
|
||||
)
|
||||
.route(
|
||||
"/api/v1/targets/{id}/applicable-scans",
|
||||
get(handlers::onboarding::applicable_scans_for_target),
|
||||
|
||||
@@ -162,14 +162,27 @@ fn ingest_blob(
|
||||
match blob::extract_zip(&stored, &dest) {
|
||||
Ok(()) => dest,
|
||||
Err(e) => {
|
||||
// Not a zip (e.g. a tar.gz source archive) — keep the blob and
|
||||
// note it so later stages can decide what to do.
|
||||
// Not a zip container — this is a single uploaded file (e.g. a
|
||||
// `.st`/`.xml` PLC project or a `.tar.gz`). The content-addressed
|
||||
// blob has no extension, so materialize it into a working dir
|
||||
// under its original name; extension-based scanners (PLC) can then
|
||||
// discover it and report a readable path.
|
||||
facts.push(DetectedFact::new(
|
||||
"archive_unextracted",
|
||||
e.to_string(),
|
||||
"ingest",
|
||||
));
|
||||
stored.clone()
|
||||
match materialize_single(&stored, &dest, &blob_file_name(artifact)) {
|
||||
Ok(dir) => dir,
|
||||
Err(copy_err) => {
|
||||
facts.push(DetectedFact::new(
|
||||
"materialize_failed",
|
||||
copy_err.to_string(),
|
||||
"ingest",
|
||||
));
|
||||
stored.clone()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
} else {
|
||||
@@ -186,6 +199,27 @@ fn ingest_blob(
|
||||
})
|
||||
}
|
||||
|
||||
/// Copy a stored blob into `dest`/`name`, returning `dest`. Used when an
|
||||
/// "extractable" artifact turns out to be a single file rather than an archive.
|
||||
fn materialize_single(stored: &Path, dest: &Path, name: &str) -> Result<PathBuf, AgentError> {
|
||||
std::fs::create_dir_all(dest)?;
|
||||
std::fs::copy(stored, dest.join(name))?;
|
||||
Ok(dest.to_path_buf())
|
||||
}
|
||||
|
||||
/// A safe, single-segment file name for an artifact, preserving the original
|
||||
/// extension so scanners can identify it. Derives from `source_ref` (the
|
||||
/// uploaded/original file name); `file_name` strips any directory components,
|
||||
/// so this is traversal-safe. Falls back to the artifact id.
|
||||
fn blob_file_name(artifact: &Artifact) -> String {
|
||||
Path::new(&artifact.source_ref)
|
||||
.file_name()
|
||||
.and_then(|n| n.to_str())
|
||||
.map(str::to_string)
|
||||
.filter(|s| !s.is_empty())
|
||||
.unwrap_or_else(|| format!("artifact-{}", artifact.id))
|
||||
}
|
||||
|
||||
/// An artifact with no on-disk form: record a single fact, no hash/path.
|
||||
fn metadata_only(artifact: &Artifact, fact: DetectedFact) -> IngestedArtifact {
|
||||
IngestedArtifact {
|
||||
@@ -331,4 +365,52 @@ mod tests {
|
||||
assert_eq!(creds.ssh_key_path.as_deref(), Some("/default/ssh/key"));
|
||||
assert!(creds.auth_token.is_none());
|
||||
}
|
||||
|
||||
/// A single uploaded PLC file (not an archive) must land in a working dir
|
||||
/// under its original name so the PLC scanner can discover it by extension
|
||||
/// and report a readable path — the demo's upload → scan path.
|
||||
#[test]
|
||||
fn single_uploaded_plc_file_is_materialized_and_scannable() {
|
||||
use compliance_core::models::PlcFormat;
|
||||
|
||||
let scratch = Scratch::new();
|
||||
let store = scratch.0.join("store");
|
||||
// Simulate the upload handler: bytes written to an `uploads/` path,
|
||||
// `source_ref` carrying the original (clean) file name.
|
||||
let uploads = scratch.0.join("uploads");
|
||||
std::fs::create_dir_all(&uploads).expect("mkdir uploads");
|
||||
let uploaded = uploads.join("a1b2c3_pump_station.st");
|
||||
std::fs::write(
|
||||
&uploaded,
|
||||
"PROGRAM P\nVAR\n ApiKey : STRING := 'sk-live-1234';\nEND_VAR\nEND_PROGRAM\n",
|
||||
)
|
||||
.expect("write st");
|
||||
|
||||
let mut artifact = Artifact::plc_project("pump_station.st", PlcFormat::StructuredText);
|
||||
artifact.stored_path = Some(uploaded.to_string_lossy().to_string());
|
||||
|
||||
let ctx = ctx_for(&store, "t-plc");
|
||||
let out = ingest_artifact(&artifact, &ctx).expect("ingest");
|
||||
|
||||
// Working path is a directory (not the extensionless blob) holding the
|
||||
// file under its original name.
|
||||
let wp = out.working_path.expect("working path");
|
||||
assert!(wp.is_dir(), "expected a working dir, got {wp:?}");
|
||||
assert!(wp.join("pump_station.st").is_file());
|
||||
|
||||
// The PLC scanner finds the hardcoded credential and reports a clean path.
|
||||
let findings = crate::pipeline::plc::analyze_tree(&wp, "t-plc");
|
||||
assert!(
|
||||
!findings.is_empty(),
|
||||
"scanner should flag the uploaded file"
|
||||
);
|
||||
assert!(findings
|
||||
.iter()
|
||||
.any(|f| f.rule_id.as_deref() == Some("plc-hardcoded-credential")));
|
||||
assert_eq!(
|
||||
findings[0].file_path.as_deref(),
|
||||
Some("pump_station.st"),
|
||||
"finding should reference the original file name"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user