use crate::params::SuppliedParams;
use crate::serve::error::ServeError;
use crate::serve::history::templates::{
TemplateRecord, TemplateState, TemplateSummary, VersionChannel, VersionSelector,
};
use crate::serve::rbac::AuthContext;
use crate::serve::runner::{self, ConfigFormatWire, SubmitRequest, SubmitResponse};
use crate::serve::state::ServerState;
use crate::templates::{RegisterRequest, TemplateStore};
use axum::Json;
use axum::extract::{Extension, Path, Query, State};
use axum::http::StatusCode;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::BTreeMap;
fn map_err(e: crate::error::CliError) -> ServeError {
use crate::error::CliError;
match e {
CliError::UnknownPipelineTemplate { .. } => ServeError::NotFound,
CliError::Internal(m) => ServeError::Internal(m),
other => ServeError::Unprocessable {
message: other.to_string(),
details: None,
},
}
}
fn store(state: &ServerState) -> TemplateStore {
state.history()
}
#[derive(Debug, Deserialize)]
pub struct RegisterBody {
#[serde(default)]
pub id: Option<String>,
pub config: String,
#[serde(default)]
pub config_format: ConfigFormatWire,
#[serde(default)]
pub description: Option<String>,
#[serde(default)]
pub tags: Vec<VersionChannel>,
#[serde(default)]
pub launch: bool,
}
pub async fn register_template(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Json(body): Json<RegisterBody>,
) -> Result<(StatusCode, Json<TemplateSummary>), ServeError> {
let record = crate::templates::register(
&store(&state),
RegisterRequest {
id: body.id,
body: body.config,
format: body.config_format.into(),
description: body.description,
tags: body.tags,
launch: body.launch,
created_by: Some(actor.principal.clone()),
},
)
.await
.map_err(map_err)?;
let fingerprint = crate::serve::idempotency::fingerprint(
&serde_json::Value::String(record.body.clone()),
record.name.as_deref(),
);
tracing::info!(
principal = %actor.principal,
template = %record.id,
version = record.version,
"registered pipeline template"
);
crate::serve::audit::write(
&state,
&actor,
"template.register",
None,
Some(fingerprint),
"ok",
)
.await;
Ok((StatusCode::CREATED, Json(record.summary())))
}
#[derive(Debug, Serialize)]
pub struct ListResponse {
pub templates: Vec<TemplateSummary>,
}
pub async fn list_templates(
State(state): State<ServerState>,
) -> Result<Json<ListResponse>, ServeError> {
let templates = crate::templates::list_with_state(&store(&state))
.await
.map_err(map_err)?;
Ok(Json(ListResponse { templates }))
}
#[derive(Debug, Default, Deserialize)]
pub struct VersionQuery {
#[serde(default)]
pub version: Option<VersionSelector>,
#[serde(default)]
pub clean: bool,
}
impl VersionQuery {
fn selector(&self) -> VersionSelector {
self.version.unwrap_or_default()
}
}
#[derive(Debug, Serialize)]
pub struct GetResponse {
#[serde(flatten)]
pub template: TemplateRecord,
#[serde(flatten)]
pub state: TemplateState,
pub is_stable: bool,
pub launches: Vec<crate::serve::history::templates::LaunchRecord>,
}
pub async fn get_template(
State(state): State<ServerState>,
Path(id): Path<String>,
Query(q): Query<VersionQuery>,
) -> Result<Json<GetResponse>, ServeError> {
let s = store(&state);
let want = crate::templates::resolve_version(&s, &id, q.selector())
.await
.map_err(map_err)?;
let mut template = s
.template_get(&id, Some(want))
.await
.map_err(|e| ServeError::Internal(e.to_string()))?
.ok_or(ServeError::NotFound)?;
if q.clean {
template.body = crate::templates::clean_config_yaml(&template.body).map_err(map_err)?;
}
let state = crate::templates::template_state(&s, &id)
.await
.map_err(map_err)?;
let launches = s
.template_launches(&id)
.await
.map_err(|e| ServeError::Internal(e.to_string()))?;
let is_stable = state.stable == Some(template.version);
Ok(Json(GetResponse {
template,
state,
is_stable,
launches,
}))
}
pub async fn delete_template(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Query(q): Query<VersionQuery>,
) -> Result<StatusCode, ServeError> {
let s = store(&state);
let target = match q.version {
None => None,
Some(selector) => Some(
crate::templates::resolve_version(&s, &id, selector)
.await
.map_err(map_err)?,
),
};
let removed = s
.template_delete(&id, target)
.await
.map_err(|e| ServeError::Internal(e.to_string()))?;
if removed == 0 {
return Err(ServeError::NotFound);
}
tracing::info!(
principal = %actor.principal,
template = %id,
version = ?target,
removed,
"deleted pipeline template version(s)"
);
crate::serve::audit::write(&state, &actor, "template.delete", None, None, "ok").await;
Ok(StatusCode::NO_CONTENT)
}
#[derive(Debug, Deserialize)]
pub struct PromoteBody {
pub tag: VersionChannel,
#[serde(default)]
pub version: Option<VersionSelector>,
}
#[derive(Debug, Serialize)]
pub struct PromoteResponse {
pub id: String,
pub tag: String,
pub version: u32,
}
pub async fn promote_template(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Json(body): Json<PromoteBody>,
) -> Result<Json<PromoteResponse>, ServeError> {
let version = crate::templates::promote(
&store(&state),
&id,
body.tag,
body.version.unwrap_or_default(),
)
.await
.map_err(map_err)?;
tracing::info!(
principal = %actor.principal,
template = %id,
tag = body.tag.as_str(),
version,
"promoted pipeline template channel"
);
crate::serve::audit::write(&state, &actor, "template.promote", None, None, "ok").await;
Ok(Json(PromoteResponse {
id,
tag: body.tag.as_str().to_string(),
version,
}))
}
#[derive(Debug, Default, Deserialize)]
pub struct LaunchBody {
#[serde(default)]
pub version: Option<VersionSelector>,
}
#[derive(Debug, Serialize)]
pub struct LaunchResponse {
pub id: String,
pub version: u32,
#[serde(skip_serializing_if = "Option::is_none")]
pub replaced: Option<u32>,
pub already_launched: bool,
pub status: String,
}
pub async fn launch_template(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Json(body): Json<LaunchBody>,
) -> Result<Json<LaunchResponse>, ServeError> {
let target = body.version.unwrap_or_else(VersionSelector::newest);
let outcome = crate::templates::launch(&store(&state), &id, target, Some(&actor.principal))
.await
.map_err(map_err)?;
finish_launch(&state, &actor, &id, outcome, "template.launch").await
}
pub async fn rollback_template(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
) -> Result<Json<LaunchResponse>, ServeError> {
let outcome = crate::templates::rollback(&store(&state), &id, Some(&actor.principal))
.await
.map_err(map_err)?;
finish_launch(&state, &actor, &id, outcome, "template.rollback").await
}
async fn finish_launch(
state: &ServerState,
actor: &AuthContext,
id: &str,
outcome: crate::templates::LaunchOutcome,
action: &str,
) -> Result<Json<LaunchResponse>, ServeError> {
let status = crate::templates::template_state(&store(state), id)
.await
.map_err(map_err)?
.status;
tracing::info!(
principal = %actor.principal,
template = %id,
version = outcome.version,
replaced = ?outcome.replaced,
already_launched = outcome.already_launched,
action,
"pipeline template launch"
);
crate::serve::audit::write(state, actor, action, None, None, "ok").await;
Ok(Json(LaunchResponse {
id: id.to_string(),
version: outcome.version,
replaced: outcome.replaced,
already_launched: outcome.already_launched,
status: status.as_str().to_string(),
}))
}
#[derive(Debug, Default, Deserialize)]
pub struct DeprecateBody {
#[serde(default)]
pub reason: Option<String>,
#[serde(default)]
pub undo: bool,
}
pub async fn deprecate_template(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Json(body): Json<DeprecateBody>,
) -> Result<Json<serde_json::Value>, ServeError> {
let status = crate::templates::set_deprecated(
&store(&state),
&id,
body.reason.clone(),
Some(&actor.principal),
!body.undo,
)
.await
.map_err(map_err)?;
let action = if body.undo {
"template.undeprecate"
} else {
"template.deprecate"
};
tracing::info!(
principal = %actor.principal,
template = %id,
status = status.as_str(),
"pipeline template deprecation changed"
);
crate::serve::audit::write(&state, &actor, action, None, None, "ok").await;
Ok(Json(
serde_json::json!({ "id": id, "status": status.as_str() }),
))
}
#[derive(Debug, Default, Deserialize)]
pub struct TriggerBody {
#[serde(default)]
pub params: BTreeMap<String, Value>,
#[serde(default)]
pub env: BTreeMap<String, String>,
#[serde(default)]
pub version: Option<VersionSelector>,
#[serde(default)]
pub name: Option<String>,
#[serde(default)]
pub labels: BTreeMap<String, String>,
#[serde(default)]
pub timeout_secs: Option<u64>,
#[serde(default)]
pub doctor_first: bool,
#[serde(default)]
pub idempotency_key: Option<String>,
#[serde(default)]
pub clock: Option<String>,
#[serde(default)]
pub callback: Option<crate::serve::callback::CallbackSpec>,
}
#[derive(Debug, Serialize)]
pub struct TriggerResponse {
#[serde(flatten)]
pub run: SubmitResponse,
pub template_id: String,
pub template_version: u32,
pub params: BTreeMap<String, Value>,
#[serde(skip_serializing_if = "Option::is_none")]
pub deprecated: Option<String>,
}
const LABEL_TEMPLATE: &str = "template";
const LABEL_TEMPLATE_VERSION: &str = "template_version";
pub async fn trigger_template(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Json(body): Json<TriggerBody>,
) -> Result<(StatusCode, Json<TriggerResponse>), ServeError> {
let supplied: SuppliedParams = body.params.into_iter().collect();
let s = store(&state);
let want = crate::templates::resolve_version(&s, &id, body.version.unwrap_or_default())
.await
.map_err(map_err)?;
let tstate = crate::templates::template_state(&s, &id)
.await
.map_err(map_err)?;
if tstate.status == crate::serve::history::templates::TemplateStatus::Deprecated {
tracing::warn!(
template = %id,
version = want,
reason = tstate
.deprecation
.as_ref()
.and_then(|d| d.reason.as_deref())
.unwrap_or("(none given)"),
"triggering a DEPRECATED pipeline template"
);
}
let clustered = state.cluster().enabled();
let mode = if clustered {
crate::templates::Materialize::Persisted
} else {
crate::templates::Materialize::Local
};
let materialized = crate::templates::materialize(&s, &id, want, &supplied, &body.env, mode)
.await
.map_err(map_err)?;
if clustered && (materialized.used_secret_params || !body.env.is_empty()) {
let what = if materialized.used_secret_params {
"declares `secret: true` param(s)"
} else {
"was triggered with `env` overrides"
};
return Err(ServeError::Unprocessable {
message: format!(
"this template {what}, and a clustered server persists the materialized config \
so a peer can execute it — which would store the value in the shared \
run-history database. Reference the secret from the template body instead \
(`${{env:VAR}}`, `${{vault:…}}`, `${{aws-sm:…}}`, … — all resolved on the \
executing instance, never persisted), or trigger it on a non-clustered server"
),
details: None,
});
}
let mut labels = body.labels;
labels.insert(LABEL_TEMPLATE.into(), materialized.template_id.clone());
labels.insert(
LABEL_TEMPLATE_VERSION.into(),
materialized.version.to_string(),
);
let req = SubmitRequest {
config: materialized.body.clone(),
config_format: ConfigFormatWire::Json,
name: body.name.or_else(|| materialized.name.clone()),
labels,
timeout_secs: body.timeout_secs,
doctor_first: body.doctor_first,
idempotency_key: body.idempotency_key,
clock: body.clock,
callback: body.callback,
};
let run = runner::submit(state.clone(), req, actor.clone()).await?;
crate::serve::audit::write(
&state,
&actor,
"template.run",
Some(run.run_id.clone()),
None,
"ok",
)
.await;
Ok((
StatusCode::ACCEPTED,
Json(TriggerResponse {
run,
template_id: materialized.template_id,
template_version: materialized.version,
params: materialized.params_redacted,
deprecated: (tstate.status
== crate::serve::history::templates::TemplateStatus::Deprecated)
.then(|| {
tstate
.deprecation
.as_ref()
.and_then(|d| d.reason.clone())
.unwrap_or_else(|| "this template is deprecated".to_string())
}),
}),
))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::serve::history::AuditFilter;
use crate::serve::rbac::Role;
use crate::serve::test_support::test_state;
use serde_json::json;
fn actor() -> AuthContext {
AuthContext {
principal: "tester".into(),
role: Role::Admin,
source_ip: None,
}
}
fn template_yaml(out: &std::path::Path) -> String {
format!(
"version: 1\nname: tpl-demo\nparams:\n tag: {{ required: true }}\n page: {{ type: int, default: 7 }}\npipeline:\n source:\n type: csv\n config:\n path: ./missing-${{param.tag}}.csv\n sink:\n type: jsonl\n config:\n path: {}\n",
out.display()
)
}
async fn register_demo_opts(
state: &ServerState,
out: &std::path::Path,
launch: bool,
) -> TemplateSummary {
register_template(
State(state.clone()),
Extension(actor()),
Json(RegisterBody {
id: None,
config: template_yaml(out),
config_format: ConfigFormatWire::Yaml,
description: Some("demo".into()),
tags: vec![],
launch,
}),
)
.await
.expect("register")
.1
.0
}
async fn register_demo(state: &ServerState, out: &std::path::Path) -> TemplateSummary {
register_demo_opts(state, out, true).await
}
async fn get(
state: &ServerState,
id: &str,
q: VersionQuery,
) -> Result<GetResponse, ServeError> {
get_template(State(state.clone()), Path(id.into()), Query(q))
.await
.map(|j| j.0)
}
#[tokio::test]
async fn register_list_get_delete_round_trip() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
let summary = register_demo(&state, &dir.path().join("o.jsonl")).await;
assert_eq!(summary.id, "tpl-demo");
assert_eq!(summary.version, 1);
assert_eq!(summary.created_by.as_deref(), Some("tester"));
assert!(summary.params["tag"].required);
let listed = list_templates(State(state.clone())).await.unwrap().0;
assert_eq!(listed.templates.len(), 1);
let st = listed.templates[0]
.state
.as_ref()
.expect("state on list rows");
assert_eq!(st.status.as_str(), "launched");
assert_eq!(st.stable, Some(1));
let got = get(&state, "tpl-demo", VersionQuery::default())
.await
.unwrap();
assert_eq!(got.template.version, 1);
assert_eq!(got.state.versions, vec![1]);
assert!(got.is_stable);
assert_eq!(got.launches.len(), 1, "the launch is recorded");
assert!(got.template.body.contains("${param.tag}"));
assert!(matches!(
get(&state, "nope", VersionQuery::default()).await,
Err(ServeError::NotFound)
));
assert!(matches!(
get(
&state,
"tpl-demo",
VersionQuery {
version: Some(VersionSelector::Pinned(9)),
clean: false,
},
)
.await,
Err(ServeError::NotFound)
));
let code = delete_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Query(VersionQuery::default()),
)
.await
.unwrap();
assert_eq!(code, StatusCode::NO_CONTENT);
assert!(matches!(
delete_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Query(VersionQuery::default())
)
.await,
Err(ServeError::NotFound)
));
let entries = state
.history()
.list_audit(&AuditFilter {
limit: 20,
..Default::default()
})
.await
.unwrap();
let actions: Vec<&str> = entries.iter().map(|e| e.action.as_str()).collect();
assert!(actions.contains(&"template.register"), "{actions:?}");
assert!(actions.contains(&"template.delete"), "{actions:?}");
}
#[tokio::test]
async fn register_rejects_an_invalid_config() {
let state = test_state();
let err = register_template(
State(state),
Extension(actor()),
Json(RegisterBody {
id: None,
config: "version: 1\nname: x\nbogus_key: 1\npipeline: {}\n".into(),
config_format: ConfigFormatWire::Yaml,
description: None,
tags: vec![],
launch: false,
}),
)
.await
.unwrap_err();
assert!(matches!(err, ServeError::Unprocessable { .. }), "{err:?}");
}
#[tokio::test]
async fn trigger_binds_params_and_stamps_provenance() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
register_demo(&state, &dir.path().join("o.jsonl")).await;
let (code, resp) = trigger_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(TriggerBody {
params: [("tag".to_string(), json!("alpha"))].into(),
..Default::default()
}),
)
.await
.expect("trigger");
assert_eq!(code, StatusCode::ACCEPTED);
assert_eq!(resp.0.template_id, "tpl-demo");
assert_eq!(resp.0.template_version, 1);
assert_eq!(resp.0.params["tag"], json!("alpha"));
assert_eq!(resp.0.params["page"], json!(7));
assert!(resp.0.deprecated.is_none());
let rec = state
.history()
.get(&resp.0.run.run_id)
.await
.unwrap()
.expect("run record");
assert_eq!(rec.labels[LABEL_TEMPLATE], "tpl-demo");
assert_eq!(rec.labels[LABEL_TEMPLATE_VERSION], "1");
assert_eq!(rec.name.as_deref(), Some("tpl-demo"));
let entries = state
.history()
.list_audit(&AuditFilter {
action: Some("template.run".into()),
limit: 10,
..Default::default()
})
.await
.unwrap();
assert_eq!(entries.len(), 1);
assert_eq!(
entries[0].run_id.as_deref(),
Some(resp.0.run.run_id.as_str())
);
}
#[tokio::test]
async fn a_draft_template_cannot_be_triggered_unpinned() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
register_demo_opts(&state, &dir.path().join("o.jsonl"), false).await;
let err = trigger_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(TriggerBody {
params: [("tag".to_string(), json!("x"))].into(),
..Default::default()
}),
)
.await
.unwrap_err();
match err {
ServeError::Unprocessable { message, .. } => {
assert!(message.contains("no launched version"), "{message}");
assert!(message.contains("launch"), "{message}");
}
other => panic!("expected 422, got {other:?}"),
}
let resp = trigger_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(TriggerBody {
params: [("tag".to_string(), json!("x"))].into(),
version: Some(VersionSelector::newest()),
..Default::default()
}),
)
.await
.expect("explicit newest runs a draft")
.1
.0;
assert_eq!(resp.template_version, 1);
}
#[tokio::test]
async fn launch_moves_callers_and_registering_does_not() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
register_demo(&state, &dir.path().join("v1.jsonl")).await; register_demo_opts(&state, &dir.path().join("v2.jsonl"), false).await;
let trigger = |version: Option<VersionSelector>| {
let state = state.clone();
async move {
trigger_template(
State(state),
Extension(actor()),
Path("tpl-demo".into()),
Json(TriggerBody {
params: [("tag".to_string(), json!("x"))].into(),
version,
..Default::default()
}),
)
.await
.expect("trigger")
.1
.0
}
};
assert_eq!(trigger(None).await.template_version, 1);
assert_eq!(
trigger(Some(VersionSelector::newest()))
.await
.template_version,
2
);
let resp = launch_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(LaunchBody::default()),
)
.await
.unwrap()
.0;
assert_eq!((resp.version, resp.replaced), (2, Some(1)));
assert_eq!(resp.status, "launched");
assert!(!resp.already_launched);
assert_eq!(trigger(None).await.template_version, 2);
let resp = rollback_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
)
.await
.unwrap()
.0;
assert_eq!((resp.version, resp.replaced), (1, Some(2)));
assert_eq!(trigger(None).await.template_version, 1);
for action in ["template.launch", "template.rollback"] {
let entries = state
.history()
.list_audit(&AuditFilter {
action: Some(action.into()),
limit: 10,
..Default::default()
})
.await
.unwrap();
assert_eq!(entries.len(), 1, "{action}");
}
}
#[tokio::test]
async fn deprecation_warns_but_keeps_serving() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
register_demo(&state, &dir.path().join("o.jsonl")).await;
let body = deprecate_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(DeprecateBody {
reason: Some("superseded by tenant-sync-v2".into()),
undo: false,
}),
)
.await
.unwrap()
.0;
assert_eq!(body["status"], "deprecated");
let resp = trigger_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(TriggerBody {
params: [("tag".to_string(), json!("x"))].into(),
..Default::default()
}),
)
.await
.expect("a deprecated template still runs")
.1
.0;
assert_eq!(resp.template_version, 1);
assert_eq!(
resp.deprecated.as_deref(),
Some("superseded by tenant-sync-v2")
);
let err = launch_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(LaunchBody::default()),
)
.await
.unwrap_err();
assert!(matches!(err, ServeError::Unprocessable { .. }), "{err:?}");
let body = deprecate_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(DeprecateBody {
reason: None,
undo: true,
}),
)
.await
.unwrap()
.0;
assert_eq!(body["status"], "launched");
}
#[tokio::test]
async fn promote_moves_a_channel_without_touching_what_is_live() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
register_demo(&state, &dir.path().join("v1.jsonl")).await; register_demo_opts(&state, &dir.path().join("v2.jsonl"), false).await;
let resp = promote_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(PromoteBody {
tag: VersionChannel::PreProd,
version: Some(VersionSelector::newest()),
}),
)
.await
.unwrap()
.0;
assert_eq!((resp.tag.as_str(), resp.version), ("pre-prod", 2));
let got = get(&state, "tpl-demo", VersionQuery::default())
.await
.unwrap();
assert_eq!(got.state.stable, Some(1), "promote must not move `stable`");
assert_eq!(got.state.tags["pre-prod"], 2);
assert!(
!got.state.tags.contains_key("stable"),
"derived, never stored"
);
let err = promote_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(PromoteBody {
tag: VersionChannel::Stable,
version: Some(VersionSelector::Pinned(1)),
}),
)
.await
.unwrap_err();
match err {
ServeError::Unprocessable { message, .. } => {
assert!(message.contains("derived"), "{message}")
}
other => panic!("expected 422, got {other:?}"),
}
assert!(
serde_json::from_value::<PromoteBody>(json!({"tag": "latest", "version": 1})).is_err()
);
}
#[tokio::test]
async fn selecting_an_unset_channel_is_unprocessable() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
register_demo(&state, &dir.path().join("v1.jsonl")).await;
let err = get(
&state,
"tpl-demo",
VersionQuery {
version: Some(VersionSelector::Channel(VersionChannel::Canary)),
clean: false,
},
)
.await
.unwrap_err();
match err {
ServeError::Unprocessable { message, .. } => {
assert!(message.contains("no `canary` version"), "{message}")
}
other => panic!("expected 422, got {other:?}"),
}
}
#[tokio::test]
async fn delete_by_selector_removes_one_version_but_omitted_removes_all() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
register_demo(&state, &dir.path().join("v1.jsonl")).await; register_demo_opts(&state, &dir.path().join("v2.jsonl"), false).await;
register_demo_opts(&state, &dir.path().join("v3.jsonl"), false).await;
delete_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Query(VersionQuery {
version: Some(VersionSelector::newest()),
clean: false,
}),
)
.await
.unwrap();
assert_eq!(
state.history().template_versions("tpl-demo").await.unwrap(),
vec![2, 1]
);
delete_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Query(VersionQuery::default()),
)
.await
.unwrap();
assert!(state.history().template_list().await.unwrap().is_empty());
}
#[tokio::test]
async fn secret_params_are_refused_on_a_clustered_server() {
let dir = tempfile::tempdir().unwrap();
let state = crate::serve::test_support::test_state_clustered();
let body = format!(
"version: 1\nname: tpl-secret\nparams:\n token: {{ required: true, secret: true }}\npipeline:\n source:\n type: csv\n config:\n path: ./x-${{param.token}}.csv\n sink:\n type: jsonl\n config:\n path: {}\n",
dir.path().join("o.jsonl").display()
);
let _registered = register_template(
State(state.clone()),
Extension(actor()),
Json(RegisterBody {
id: None,
config: body,
config_format: ConfigFormatWire::Yaml,
description: None,
tags: vec![],
launch: true,
}),
)
.await
.expect("register");
let err = trigger_template(
State(state),
Extension(actor()),
Path("tpl-secret".into()),
Json(TriggerBody {
params: [("token".to_string(), json!("super-secret-value"))].into(),
..Default::default()
}),
)
.await
.unwrap_err();
match err {
ServeError::Unprocessable { message, .. } => {
assert!(message.contains("clustered"), "{message}");
assert!(!message.contains("super-secret-value"), "leaked: {message}");
}
other => panic!("expected 422, got {other:?}"),
}
}
#[tokio::test]
async fn env_overrides_are_refused_on_a_clustered_server() {
let dir = tempfile::tempdir().unwrap();
let state = crate::serve::test_support::test_state_clustered();
let body = format!(
"version: 1\nname: tpl-env\npipeline:\n source:\n type: csv\n config:\n path: \"${{env:SRC_PATH}}\"\n sink:\n type: jsonl\n config:\n path: {}\n",
dir.path().join("o.jsonl").display()
);
let _registered = register_template(
State(state.clone()),
Extension(actor()),
Json(RegisterBody {
id: None,
config: body,
config_format: ConfigFormatWire::Yaml,
description: None,
tags: vec![],
launch: true,
}),
)
.await
.expect("register");
let err = trigger_template(
State(state),
Extension(actor()),
Path("tpl-env".into()),
Json(TriggerBody {
env: [("SRC_PATH".to_string(), "s3cret-path".to_string())].into(),
..Default::default()
}),
)
.await
.unwrap_err();
match err {
ServeError::Unprocessable { message, .. } => {
assert!(message.contains("env"), "{message}");
assert!(!message.contains("s3cret-path"), "leaked: {message}");
}
other => panic!("expected 422, got {other:?}"),
}
}
#[tokio::test]
async fn a_clustered_trigger_persists_tokens_not_resolved_values() {
let dir = tempfile::tempdir().unwrap();
unsafe { std::env::set_var("FAUCET_TEST_C5_SECRET", "hunter2-should-not-persist") };
let s = store(&crate::serve::test_support::test_state_clustered());
let body = format!(
"version: 1\nname: tpl-c5\npipeline:\n source:\n type: csv\n config:\n path: \"${{env:FAUCET_TEST_C5_SECRET}}\"\n sink:\n type: jsonl\n config:\n path: {}\n",
dir.path().join("o.jsonl").display()
);
crate::templates::register(
&s,
crate::templates::RegisterRequest {
id: None,
body,
format: crate::serve::load::ConfigFormat::Yaml,
description: None,
tags: vec![],
launch: true,
created_by: None,
},
)
.await
.expect("register");
let persisted = crate::templates::materialize(
&s,
"tpl-c5",
1,
&Default::default(),
&Default::default(),
crate::templates::Materialize::Persisted,
)
.await
.expect("materialize");
assert!(
!persisted.body.contains("hunter2-should-not-persist"),
"a resolved secret must never reach a persisted body: {}",
persisted.body
);
assert!(
persisted.body.contains("${env:FAUCET_TEST_C5_SECRET}"),
"the directive must survive as a token for the executor: {}",
persisted.body
);
let local = crate::templates::materialize(
&s,
"tpl-c5",
1,
&Default::default(),
&Default::default(),
crate::templates::Materialize::Local,
)
.await
.expect("materialize");
assert!(
local.body.contains("hunter2-should-not-persist"),
"{}",
local.body
);
unsafe { std::env::remove_var("FAUCET_TEST_C5_SECRET") };
}
#[test]
fn version_query_deserializes_channels_and_numbers() {
assert!(
serde_json::from_value::<VersionQuery>(json!({}))
.unwrap()
.selector()
.is_stable(),
"an omitted selector means `stable`"
);
assert!(
serde_json::from_value::<VersionQuery>(json!({ "version": "stable" }))
.unwrap()
.selector()
.is_stable()
);
for wire in [json!({ "version": "2" }), json!({ "version": 2 })] {
let q: VersionQuery = serde_json::from_value(wire.clone()).unwrap();
assert_eq!(q.selector().pinned(), Some(2), "{wire}");
}
for bad in [
json!({ "version": "nope" }),
json!({ "version": 0 }),
json!({ "version": "latest" }),
] {
assert!(
serde_json::from_value::<VersionQuery>(bad.clone()).is_err(),
"{bad} should be rejected"
);
}
}
#[tokio::test]
async fn version_pinning_selects_an_older_body() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
register_demo(&state, &dir.path().join("v1.jsonl")).await;
let v2 = register_demo(&state, &dir.path().join("v2.jsonl")).await;
assert_eq!(v2.version, 2);
let got = get_template(
State(state.clone()),
Path("tpl-demo".into()),
Query(VersionQuery {
version: Some(VersionSelector::Pinned(1)),
clean: false,
}),
)
.await
.unwrap()
.0;
assert!(got.template.body.contains("v1.jsonl"));
assert_eq!(got.state.versions, vec![2, 1]);
assert!(!got.is_stable, "v2 is live, so a pinned v1 is not");
delete_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Query(VersionQuery {
version: Some(VersionSelector::Pinned(1)),
clean: false,
}),
)
.await
.unwrap();
assert_eq!(
state.history().template_versions("tpl-demo").await.unwrap(),
vec![2]
);
}
}