use axum::Json;
use axum::extract::{Multipart, Path, State};
use axum::http::StatusCode;
use futures::StreamExt;
use kanade_shared::bin_platform::{
AgentPlatform, base_version_of_key, candidate_keys, platform_of_key,
};
use kanade_shared::exe_version::extract_pe_version;
use kanade_shared::kv::{
BUCKET_AGENT_CONFIG, KEY_AGENT_CONFIG_GLOBAL, OBJECT_AGENT_RELEASES, agent_config_group_key,
agent_config_pc_key, parse_agent_config_group_key, parse_agent_config_pc_key,
};
use kanade_shared::wire::ConfigScope;
use serde::{Deserialize, Serialize};
use tracing::{info, warn};
use super::AppState;
use crate::audit;
use crate::audit::Caller;
#[derive(Serialize)]
pub struct PublishResponse {
pub version: String,
pub key: String,
pub platform: String,
pub size: u64,
pub digest: Option<String>,
}
fn windows_label(pe: Option<String>, field: Option<&str>) -> Result<String, (StatusCode, String)> {
Ok(match (pe, field) {
(Some(pe), Some(v))
if v.trim().trim_start_matches('v') != pe.trim().trim_start_matches('v') =>
{
return Err((
StatusCode::BAD_REQUEST,
format!(
"version field '{v}' disagrees with the binary's embedded version '{pe}'; \
omit the field to use the embedded one, or pass the matching label"
),
));
}
(Some(pe), _) => pe,
(None, Some(v)) => v.to_owned(),
(None, None) => {
return Err((
StatusCode::BAD_REQUEST,
"couldn't extract VERSIONINFO from the uploaded binary — \
is it a Windows PE built with `winres`? Kanade ≥ v0.13.1 \
embeds the resource automatically; for an older binary \
include a 'version' form field with the label to use."
.to_owned(),
));
}
})
}
pub async fn publish(
State(state): State<AppState>,
caller: Caller,
mut multipart: Multipart,
) -> Result<Json<PublishResponse>, (StatusCode, String)> {
let mut bytes: Option<Vec<u8>> = None;
let mut version_field: Option<String> = None;
while let Some(field) = multipart
.next_field()
.await
.map_err(|e| (e.status(), format!("read multipart field: {e}")))?
{
match field.name().unwrap_or("") {
"file" => {
let buf = field
.bytes()
.await
.map_err(|e| (e.status(), format!("read file field: {e}")))?;
bytes = Some(buf.to_vec());
}
"version" => {
let text = field
.text()
.await
.map_err(|e| (e.status(), format!("read version field: {e}")))?;
let text = text.trim().to_string();
if !text.is_empty() {
version_field = Some(text);
}
}
other => {
warn!(field = other, "publish: ignoring unknown multipart field");
}
}
}
let bytes = bytes.ok_or((StatusCode::BAD_REQUEST, "missing 'file' field".into()))?;
if bytes.is_empty() {
return Err((StatusCode::BAD_REQUEST, "'file' field is empty".into()));
}
let platform = AgentPlatform::detect(&bytes).map_err(|e| (StatusCode::BAD_REQUEST, e))?;
let version = match platform {
AgentPlatform::WindowsX86_64 | AgentPlatform::WindowsAarch64 => {
windows_label(extract_pe_version(&bytes), version_field.as_deref())?
}
AgentPlatform::LinuxX86_64 | AgentPlatform::LinuxAarch64 | AgentPlatform::MacOSAarch64 => {
let format = match platform {
AgentPlatform::MacOSAarch64 => "macOS Mach-O",
_ => "Linux ELF",
};
version_field.ok_or((
StatusCode::BAD_REQUEST,
format!(
"no version: the uploaded binary is a {format} ({}), which carries no \
embedded VERSIONINFO — include a 'version' form field (e.g. X.Y.Z). A \
Windows PE built with `winres` (kanade ≥ v0.13.1) is auto-labelled.",
platform.as_str()
),
))?
}
};
let key = platform.release_key(&version);
kanade_shared::bin_platform::check_release_key(&key)
.map_err(|e| (StatusCode::BAD_REQUEST, e))?;
let size = bytes.len() as u64;
info!(
version,
key,
platform = platform.as_str(),
size,
"publish: uploading new agent binary"
);
let store = state
.jetstream
.get_object_store(OBJECT_AGENT_RELEASES)
.await
.map_err(|e| {
warn!(error = %e, "get_object_store agent_releases");
(
StatusCode::SERVICE_UNAVAILABLE,
format!(
"Object Store '{OBJECT_AGENT_RELEASES}' not found — the backend creates it at startup, so it was removed since or this broker is not the one expected; check `GET /api/jetstream/status` (restarting the backend recreates it)"
),
)
})?;
let mut cursor = std::io::Cursor::new(bytes);
let meta = store.put(key.as_str(), &mut cursor).await.map_err(|e| {
warn!(error = %e, "object_store.put");
(StatusCode::INTERNAL_SERVER_ERROR, e.to_string())
})?;
info!(version, key, digest = ?meta.digest, "publish: agent binary uploaded");
if let Err(e) =
crate::projector::object_meta::apply(&state.pool, OBJECT_AGENT_RELEASES, &meta).await
{
warn!(error = %e, %version, "object_meta write-through failed (watcher will heal)");
}
audit::record(
&state.nats,
"operator",
"agent_publish",
Some(&key),
Some(&caller),
serde_json::json!({
"version": version,
"platform": platform.as_str(),
"size": size,
"digest": meta.digest,
}),
)
.await;
Ok(Json(PublishResponse {
version,
key,
platform: platform.as_str().to_string(),
size,
digest: meta.digest,
}))
}
pub async fn delete_release(
State(state): State<AppState>,
Path(version): Path<String>,
caller: Caller,
) -> Result<StatusCode, (StatusCode, String)> {
let kv = state
.jetstream
.get_key_value(BUCKET_AGENT_CONFIG)
.await
.map_err(|e| (StatusCode::SERVICE_UNAVAILABLE, e.to_string()))?;
let mut keys = kv
.keys()
.await
.map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
while let Some(k) = keys.next().await {
let k = match k {
Ok(k) => k,
Err(_) => continue,
};
let entry = match kv.get(&k).await.map_err(|e| {
warn!(error = %e, %k, "kv.get scope");
(StatusCode::INTERNAL_SERVER_ERROR, e.to_string())
})? {
Some(b) => b,
None => continue,
};
let scope: ConfigScope = match serde_json::from_slice(&entry) {
Ok(s) => s,
Err(_) => continue,
};
if scope.target_version.as_deref() == Some(base_version_of_key(&version)) {
let base = base_version_of_key(&version);
let label = if k == KEY_AGENT_CONFIG_GLOBAL {
"global".to_string()
} else if let Some(g) = parse_agent_config_group_key(&k) {
format!("group:{g}")
} else if let Some(p) = parse_agent_config_pc_key(&k) {
format!("pc:{p}")
} else {
k.clone()
};
return Err((
StatusCode::CONFLICT,
format!(
"version '{base}' is the current target_version of scope '{label}' — \
clear or change that scope first (kanade config unset target_version --… )"
),
));
}
}
let store = state
.jetstream
.get_object_store(OBJECT_AGENT_RELEASES)
.await
.map_err(|e| (StatusCode::SERVICE_UNAVAILABLE, e.to_string()))?;
store.delete(&version).await.map_err(|e| {
warn!(error = %e, %version, "object_store.delete");
let msg = e.to_string();
if msg.contains("not found") || msg.contains("no objects") {
(
StatusCode::NOT_FOUND,
format!("version '{version}' not in Object Store"),
)
} else {
(StatusCode::INTERNAL_SERVER_ERROR, msg)
}
})?;
info!(%version, "publish: agent binary deleted");
if let Err(e) =
crate::projector::object_meta::delete_key(&state.pool, OBJECT_AGENT_RELEASES, &version)
.await
{
warn!(error = %e, %version, "object_meta write-through failed (watcher will heal)");
}
audit::record(
&state.nats,
"operator",
"agent_release_delete",
Some(&version),
Some(&caller),
serde_json::json!({}),
)
.await;
Ok(StatusCode::NO_CONTENT)
}
#[derive(Serialize)]
pub struct ReleaseRow {
pub version: String,
pub platform: String,
pub size: u64,
pub digest: Option<String>,
pub modified: Option<String>,
}
pub async fn list_releases(
State(state): State<AppState>,
) -> Result<Json<Vec<ReleaseRow>>, (StatusCode, String)> {
let metas = crate::projector::object_meta::list_bucket(&state.pool, OBJECT_AGENT_RELEASES)
.await
.map_err(|e| {
warn!(error = %e, "object_store_meta list agent_releases");
(StatusCode::INTERNAL_SERVER_ERROR, e.to_string())
})?;
let mut rows: Vec<ReleaseRow> = metas
.into_iter()
.map(|m| ReleaseRow {
platform: platform_of_key(&m.key).to_string(),
version: m.key,
size: m.size as u64,
digest: m.digest,
modified: m.modified,
})
.collect();
rows.sort_by(|a, b| b.modified.cmp(&a.modified));
Ok(Json(rows))
}
#[derive(Deserialize, Debug)]
#[serde(rename_all = "snake_case", tag = "type", content = "value")]
pub enum RolloutScope {
Global,
Group(String),
Pc(String),
}
#[derive(Deserialize, Debug)]
pub struct RolloutBody {
pub version: String,
pub scope: RolloutScope,
#[serde(default)]
pub jitter: Option<String>,
}
#[derive(Serialize)]
pub struct RolloutResponse {
pub version: String,
pub scope_key: String,
pub scope_label: String,
pub jitter: Option<String>,
}
pub async fn rollout(
State(state): State<AppState>,
caller: Caller,
Json(body): Json<RolloutBody>,
) -> Result<Json<RolloutResponse>, (StatusCode, String)> {
let (key, label) = match &body.scope {
RolloutScope::Global => (KEY_AGENT_CONFIG_GLOBAL.to_string(), "global".to_string()),
RolloutScope::Group(g) => (agent_config_group_key(g), format!("group:{g}")),
RolloutScope::Pc(p) => (agent_config_pc_key(p), format!("pc:{p}")),
};
let version = base_version_of_key(&body.version).to_string();
let store = state
.jetstream
.get_object_store(OBJECT_AGENT_RELEASES)
.await
.map_err(|e| (StatusCode::SERVICE_UNAVAILABLE, e.to_string()))?;
let mut any_key_exists = false;
for candidate in candidate_keys(&version) {
if store.info(&candidate).await.is_ok() {
any_key_exists = true;
break;
}
}
if !any_key_exists {
return Err((
StatusCode::NOT_FOUND,
format!(
"version '{version}' not found in {OBJECT_AGENT_RELEASES} — run `kanade agent publish` first"
),
));
}
let kv = state
.jetstream
.get_key_value(BUCKET_AGENT_CONFIG)
.await
.map_err(|e| (StatusCode::SERVICE_UNAVAILABLE, e.to_string()))?;
if let Some(j) = body.jitter.as_deref() {
humantime::parse_duration(j).map_err(|e| {
(
StatusCode::BAD_REQUEST,
format!("jitter: expected a humantime duration (e.g. 30s, 10m, 1h): {e}"),
)
})?;
}
kanade_shared::kv_cas::read_modify_write(&kv, &key, |scope: &mut ConfigScope| {
let before = scope.clone();
scope.target_version = Some(version.clone());
if let Some(j) = body.jitter.as_deref() {
scope.target_version_jitter = Some(j.to_owned());
}
*scope != before
})
.await
.map_err(|e| {
warn!(error = %e, %key, "rollout scope RMW");
(
StatusCode::INTERNAL_SERVER_ERROR,
format!("update {BUCKET_AGENT_CONFIG}.{key}: {e:#}"),
)
})?;
info!(
scope = %label,
version = %version,
jitter = ?body.jitter,
"rollout: target_version flipped via HTTP",
);
audit::record(
&state.nats,
"operator",
"agent_rollout",
Some(&key),
Some(&caller),
serde_json::json!({
"version": version,
"scope_label": label,
"jitter": body.jitter,
}),
)
.await;
Ok(Json(RolloutResponse {
version,
scope_key: key,
scope_label: label,
jitter: body.jitter,
}))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn pe_without_versioninfo_takes_the_explicit_label() {
assert_eq!(windows_label(None, Some("1.2.3")).unwrap(), "1.2.3");
let (code, msg) = windows_label(None, None).unwrap_err();
assert_eq!(code, StatusCode::BAD_REQUEST);
assert!(msg.contains("'version' form field"), "{msg}");
}
#[test]
fn embedded_version_wins_and_a_contradicting_field_is_rejected() {
assert_eq!(windows_label(Some("1.2.3".into()), None).unwrap(), "1.2.3");
assert_eq!(
windows_label(Some("1.2.3".into()), Some("v1.2.3")).unwrap(),
"1.2.3"
);
let (code, _) = windows_label(Some("1.2.3".into()), Some("9.9.9")).unwrap_err();
assert_eq!(code, StatusCode::BAD_REQUEST);
}
}