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 if crate::select::is_selection_error(&other) => {
ServeError::BadConfig(other.to_string())
}
other => ServeError::Unprocessable {
message: other.to_string(),
details: None,
},
}
}
fn store(state: &ServerState) -> TemplateStore {
state.history()
}
#[derive(Debug, Clone, Serialize, 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;
let mut summary = record.summary();
#[cfg(feature = "policy")]
if let Some(policy) = state.policy()
&& record.kind == crate::hub::TemplateKind::Pipeline
{
let parsed: Result<serde_json::Value, String> = match record.format {
crate::serve::load::ConfigFormat::Yaml => {
serde_yaml::from_str(&record.body).map_err(|e| e.to_string())
}
crate::serve::load::ConfigFormat::Json => {
serde_json::from_str(&record.body).map_err(|e| e.to_string())
}
};
if let Ok(doc) = parsed
&& let Ok(cfg) = crate::templates::store::validate_pipeline_body(&doc)
&& let Ok(nodes) = crate::expand::expand(&cfg)
&& let Ok(report) = crate::policy::evaluate_nodes(&policy, &nodes, &Default::default())
&& report.violated()
{
for v in report.all_violations() {
tracing::warn!(template = %record.id, "policy: {v}");
summary.warnings.push(format!("policy: {v}"));
}
}
}
Ok((StatusCode::CREATED, Json(summary)))
}
#[derive(Debug, Serialize)]
pub struct ListResponse {
pub templates: Vec<TemplateSummary>,
#[serde(skip_serializing_if = "Option::is_none")]
pub sync: Option<SyncInfo>,
}
#[derive(Debug, Clone, Serialize)]
pub struct SyncInfo {
pub origins: Vec<SyncOriginInfo>,
}
#[derive(Debug, Clone, Serialize)]
pub struct SyncOriginInfo {
pub name: String,
pub kind: &'static str,
pub prefix: String,
pub launch: String,
pub prune: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub interval_secs: Option<u64>,
}
#[cfg(feature = "templates-sync")]
fn sync_info(state: &ServerState) -> Option<SyncInfo> {
let file = state.templates_sync()?;
Some(SyncInfo {
origins: file
.origins
.iter()
.map(|o| SyncOriginInfo {
name: o.name.clone(),
kind: o.source.kind(),
prefix: o.prefix.clone(),
launch: serde_json::to_value(o.launch)
.ok()
.and_then(|v| v.as_str().map(str::to_string))
.unwrap_or_default(),
prune: serde_json::to_value(o.prune)
.ok()
.and_then(|v| v.as_str().map(str::to_string))
.unwrap_or_default(),
interval_secs: o.interval_secs,
})
.collect(),
})
}
#[cfg(not(feature = "templates-sync"))]
fn sync_info(_state: &ServerState) -> Option<SyncInfo> {
None
}
#[derive(Debug, Default, Deserialize)]
pub struct ListQuery {
#[serde(default)]
pub kind: Option<crate::hub::TemplateKind>,
}
pub async fn list_templates(
State(state): State<ServerState>,
Query(query): Query<ListQuery>,
) -> Result<Json<ListResponse>, ServeError> {
let mut templates = crate::templates::list_with_state(&store(&state))
.await
.map_err(map_err)?;
if let Some(kind) = query.kind {
templates.retain(|t| t.kind == kind);
}
Ok(Json(ListResponse {
templates,
sync: sync_info(&state),
}))
}
pub async fn template_matrix(State(state): State<ServerState>) -> Result<Json<Value>, ServeError> {
let store = store(&state);
let templates = crate::templates::list_with_state(&store)
.await
.map_err(map_err)?;
let mut sources = Vec::new();
let mut sinks = Vec::new();
for t in templates {
if !t.kind.is_hub() {
continue;
}
let version = t.state.as_ref().and_then(|s| s.stable).unwrap_or(t.version);
let Some(rec) = store
.template_get(&t.id, Some(version))
.await
.map_err(|e| ServeError::Internal(format!("template registry read: {e}")))?
else {
continue;
};
let doc = crate::templates::parse_body(&rec.body, rec.format).map_err(map_err)?;
let path = std::path::PathBuf::from(&t.id);
match t.kind {
crate::hub::TemplateKind::SourceTemplate => {
let s: crate::hub::SourceTemplate = serde_json::from_value(doc).map_err(|e| {
ServeError::Internal(format!("stored source-template '{}': {e}", t.id))
})?;
sources.push((path, s));
}
crate::hub::TemplateKind::SinkTemplate => {
let k: crate::hub::SinkTemplate = serde_json::from_value(doc).map_err(|e| {
ServeError::Internal(format!("stored sink-template '{}': {e}", t.id))
})?;
sinks.push((path, k));
}
crate::hub::TemplateKind::Pipeline | crate::hub::TemplateKind::Deployment => {}
}
}
sources.sort_by(|a, b| a.1.name.cmp(&b.1.name));
sinks.sort_by(|a, b| a.1.name.cmp(&b.1.name));
let cat = crate::hub::Catalog {
root: std::path::PathBuf::new(),
sources,
sinks,
};
Ok(Json(crate::hub::catalog::index_json_with(
&cat,
crate::hub::catalog::registry_run_command,
)))
}
#[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,
}))
}
#[derive(Debug, Default, Deserialize)]
pub struct RowsParams {
#[serde(default)]
pub version: Option<VersionSelector>,
#[serde(default)]
pub sink: Option<String>,
#[serde(default)]
pub sink_version: Option<VersionSelector>,
#[serde(default)]
pub overlay: Option<String>,
#[serde(default)]
pub overlay_version: Option<VersionSelector>,
#[serde(default)]
pub select: Option<String>,
#[serde(default)]
pub only: Option<String>,
#[serde(default)]
pub skip: Option<String>,
#[serde(default)]
pub tags: Option<String>,
#[serde(default)]
pub status: Option<String>,
#[serde(default)]
pub include_parents: Option<String>,
#[serde(default)]
pub state: Option<bool>,
}
fn csv(v: &Option<String>) -> Vec<String> {
v.as_deref()
.map(|s| {
s.split(',')
.map(str::trim)
.filter(|t| !t.is_empty())
.map(str::to_string)
.collect()
})
.unwrap_or_default()
}
impl RowsParams {
pub fn selection(&self) -> Result<Option<crate::select::SelectionRequest>, ServeError> {
let any = [
&self.select,
&self.only,
&self.skip,
&self.tags,
&self.status,
&self.include_parents,
]
.iter()
.any(|v| v.is_some());
if !any {
return Ok(None);
}
let sel = crate::select::RunSelection::resolve(
&csv(&self.select),
&csv(&self.only),
&csv(&self.skip),
&csv(&self.status),
&csv(&self.tags),
self.include_parents.as_deref(),
None,
)
.map_err(|e| ServeError::BadConfig(e.to_string()))?;
let mut req = crate::select::SelectionRequest::from_run_selection(&sel);
if self.include_parents.is_none() {
req.include_parents = None;
}
Ok(Some(req))
}
}
pub async fn template_rows(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Query(q): Query<RowsParams>,
) -> Result<Json<crate::hub::rows::RowsReport>, ServeError> {
let s = store(&state);
let selection = q.selection()?;
let version = crate::templates::resolve_version(&s, &id, q.version.unwrap_or_default())
.await
.map_err(map_err)?;
let sink_version = match &q.sink {
Some(sid) => Some(
crate::templates::resolve_version(&s, sid, q.sink_version.unwrap_or_default())
.await
.map_err(map_err)?,
),
None => None,
};
let report = crate::templates::rows::list_rows(
&s,
crate::templates::rows::RowsQuery {
id: &id,
version,
sink: q.sink.as_deref().zip(sink_version),
overlay: q
.overlay
.clone()
.map(|oid| crate::templates::OverlayChoice::Registered {
id: oid,
version: q.overlay_version.unwrap_or_default(),
}),
selection: selection.as_ref(),
state: q.state.unwrap_or(true),
},
)
.await
.map_err(map_err)?;
crate::serve::audit::write(&state, &actor, "template.rows", None, None, "ok").await;
Ok(Json(report))
}
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() }),
))
}
pub async fn deprecate_version(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path((id, version)): Path<(String, String)>,
Json(body): Json<DeprecateBody>,
) -> Result<Json<serde_json::Value>, ServeError> {
let s = store(&state);
let selector = VersionSelector::parse(&version).map_err(map_err)?;
let version = crate::templates::resolve_version(&s, &id, selector)
.await
.map_err(map_err)?;
crate::templates::set_version_deprecated(
&s,
&id,
version,
body.reason.clone(),
Some(&actor.principal),
!body.undo,
)
.await
.map_err(map_err)?;
tracing::info!(
principal = %actor.principal,
template = %id,
version,
deprecated = !body.undo,
"pipeline template version deprecation changed"
);
crate::serve::audit::write(
&state,
&actor,
"template.version_deprecate",
None,
None,
"ok",
)
.await;
Ok(Json(serde_json::json!({
"id": id,
"version": version,
"deprecated": !body.undo,
})))
}
#[derive(Debug, Clone, Default, Deserialize)]
pub struct TriggerBody {
#[serde(default)]
pub params: BTreeMap<String, Value>,
#[serde(default)]
pub env: BTreeMap<String, String>,
#[serde(default)]
pub sink: Option<String>,
#[serde(default)]
pub sink_version: Option<VersionSelector>,
#[serde(default)]
pub overlay: Option<OverlayRef>,
#[serde(default)]
pub overlay_version: Option<VersionSelector>,
#[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 concurrency: Option<usize>,
#[serde(default)]
pub clock: Option<String>,
#[serde(default)]
pub callback: Option<crate::serve::callback::CallbackSpec>,
#[serde(default)]
pub require_approval: bool,
#[serde(default)]
pub reason: Option<String>,
#[serde(default)]
pub budget: Option<faucet_core::BudgetSpec>,
#[serde(default)]
pub selection: Option<crate::select::SelectionRequest>,
}
#[derive(Debug, Clone, Deserialize, Serialize)]
#[serde(untagged)]
pub enum OverlayRef {
Id(String),
Inline(Value),
}
impl TriggerBody {
fn overlay_choice(&self) -> Option<crate::templates::OverlayChoice> {
self.overlay.as_ref().map(|o| match o {
OverlayRef::Id(id) => crate::templates::OverlayChoice::Registered {
id: id.clone(),
version: self.overlay_version.unwrap_or_default(),
},
OverlayRef::Inline(v) => crate::templates::OverlayChoice::Inline(v.clone()),
})
}
}
#[derive(Debug, Serialize)]
pub struct TriggerResponse {
#[serde(flatten)]
pub run: SubmitResponse,
pub template_id: String,
pub template_version: u32,
#[serde(skip_serializing_if = "Option::is_none")]
pub sink_template: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub sink_template_version: Option<u32>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub streams: Vec<crate::hub::compose::StreamPlan>,
#[serde(skip_serializing_if = "Option::is_none")]
pub overlay: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub overlay_version: Option<u32>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub overlay_contributes: Vec<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub warnings: Vec<String>,
pub params: BTreeMap<String, Value>,
#[serde(skip_serializing_if = "Option::is_none")]
pub deprecated: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub selection: Option<crate::select::SelectionRequest>,
}
const LABEL_TEMPLATE: &str = "template";
const LABEL_TEMPLATE_VERSION: &str = "template_version";
const LABEL_SINK_TEMPLATE: &str = "sink_template";
const LABEL_SINK_TEMPLATE_VERSION: &str = "sink_template_version";
const LABEL_OVERLAY: &str = "overlay";
const LABEL_OVERLAY_VERSION: &str = "overlay_version";
pub async fn trigger_template(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Json(body): Json<TriggerBody>,
) -> Result<axum::response::Response, ServeError> {
use axum::response::IntoResponse;
match trigger_template_outcome(state, actor, id, body).await? {
TriggerOutcome::Run(resp) => Ok((StatusCode::ACCEPTED, Json(resp)).into_response()),
TriggerOutcome::PendingApproval(change) => Ok((
StatusCode::ACCEPTED,
Json(runner::pending_approval_body(&change)),
)
.into_response()),
}
}
#[derive(Debug)]
#[allow(clippy::large_enum_variant)]
pub enum TriggerOutcome {
Run(TriggerResponse),
PendingApproval(Box<crate::serve::changes::ChangeRequest>),
}
pub async fn trigger_template_outcome(
state: ServerState,
actor: AuthContext,
id: String,
mut body: TriggerBody,
) -> Result<TriggerOutcome, ServeError> {
let overlay = body.overlay_choice();
let supplied: SuppliedParams = std::mem::take(&mut 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)?;
let deprecated = crate::templates::deprecation_warning(&tstate, want);
if let Some(reason) = &deprecated {
tracing::warn!(
template = %id,
version = want,
reason = %reason,
"triggering a DEPRECATED pipeline template"
);
}
let clustered = state.cluster().enabled();
let mode = if clustered {
crate::templates::Materialize::Persisted
} else {
crate::templates::Materialize::Local
};
let sink = crate::templates::SinkChoice {
id: body.sink.clone(),
version: body.sink_version.unwrap_or_default(),
overlay,
};
if let Some(sink_id) = &sink.id {
let sink_state = crate::templates::template_state(&s, sink_id)
.await
.map_err(map_err)?;
if sink_state.status == crate::serve::history::templates::TemplateStatus::Deprecated {
tracing::warn!(sink_template = %sink_id, "composing a DEPRECATED sink template");
}
}
let materialized = crate::templates::materialize_for_run_selected(
&s,
&id,
want,
&sink,
&supplied,
&body.env,
mode,
body.selection.as_ref(),
)
.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(),
);
if let (Some(sid), Some(sv)) = (&materialized.sink_id, materialized.sink_version) {
labels.insert(LABEL_SINK_TEMPLATE.into(), sid.clone());
labels.insert(LABEL_SINK_TEMPLATE_VERSION.into(), sv.to_string());
}
if let Some(oid) = &materialized.overlay_id {
labels.insert(LABEL_OVERLAY.into(), oid.clone());
if let Some(ov) = materialized.overlay_version {
labels.insert(LABEL_OVERLAY_VERSION.into(), ov.to_string());
}
}
if let Some(sel) = &body.selection {
labels.insert(crate::select::LABEL_SELECTION.into(), sel.canonical());
}
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,
concurrency: body.concurrency,
callback: body.callback,
require_approval: body.require_approval,
reason: body.reason,
budget: body.budget,
approved_change: None,
selection: materialized.selection.clone(),
};
let run = match runner::submit_gated(state.clone(), req, actor.clone()).await? {
runner::SubmitOutcome::Accepted(run) => run,
runner::SubmitOutcome::PendingApproval(change) => {
return Ok(TriggerOutcome::PendingApproval(change));
}
};
crate::serve::audit::write(
&state,
&actor,
"template.run",
Some(run.run_id.clone()),
None,
"ok",
)
.await;
Ok(TriggerOutcome::Run(TriggerResponse {
run,
template_id: materialized.template_id,
template_version: materialized.version,
sink_template: materialized.sink_id,
sink_template_version: materialized.sink_version,
streams: materialized.streams,
overlay: materialized.overlay_id,
overlay_version: materialized.overlay_version,
overlay_contributes: materialized.overlay_contributes,
warnings: materialized.warnings,
params: materialized.params_redacted,
deprecated,
selection: body.selection,
}))
}
#[cfg(test)]
mod tests {
async fn trigger_pair(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Json(body): Json<TriggerBody>,
) -> Result<(StatusCode, Json<TriggerResponse>), ServeError> {
match trigger_template_outcome(state, actor, id, body).await? {
TriggerOutcome::Run(r) => Ok((StatusCode::ACCEPTED, Json(r))),
TriggerOutcome::PendingApproval(c) => panic!("unexpected pending change {}", c.id),
}
}
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,
tenant: 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()), Query(ListQuery::default()))
.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:?}");
}
fn source_template(dir: &std::path::Path) -> String {
std::fs::write(dir.join("orders.csv"), "id,total\n1,10\n2,20\n").unwrap();
std::fs::write(dir.join("customers.csv"), "id,name\n1,alice\n").unwrap();
format!(
"kind: source-template
name: acme-exports
description: Acme — orders and customers exports
params:
data_dir: {{ type: string, default: {} }}
source:
type: csv
config:
path: \"${{param.data_dir}}/orders.csv\"
transforms:
- {{ type: keys_case, config: {{ mode: snake }} }}
streams:
- name: orders
primary_keys: [id]
write: [overwrite, upsert]
- name: customers
source: {{ config: {{ path: \"${{param.data_dir}}/customers.csv\" }} }}
primary_keys: [id]
write: [overwrite, upsert]
",
dir.display()
)
}
fn sink_template(dir: &std::path::Path) -> String {
format!(
"kind: sink-template
name: local-jsonl
description: Local JSON Lines files, one per stream
params:
out_dir: {{ type: string, default: {} }}
sink:
type: jsonl
config:
append: false
per_stream:
path: \"${{param.out_dir}}/${{source}}/${{stream}}.jsonl\"
write_mode_aliases:
overwrite: append
",
dir.display()
)
}
async fn register_body(state: &ServerState, config: String) -> TemplateSummary {
register_template(
State(state.clone()),
Extension(actor()),
Json(RegisterBody {
id: None,
config,
config_format: ConfigFormatWire::Yaml,
description: None,
tags: vec![],
launch: true,
}),
)
.await
.expect("register")
.1
.0
}
#[tokio::test]
async fn a_source_template_triggers_composed_with_a_sink_and_lists_by_kind() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
let src = register_body(&state, source_template(dir.path())).await;
assert_eq!(src.kind, crate::hub::TemplateKind::SourceTemplate);
assert_eq!(
src.description.as_deref(),
Some("Acme — orders and customers exports")
);
register_body(&state, sink_template(dir.path())).await;
register_demo(&state, &dir.path().join("o.jsonl")).await;
let all = list_templates(State(state.clone()), Query(ListQuery::default()))
.await
.unwrap()
.0;
assert_eq!(all.templates.len(), 3);
let sinks = list_templates(
State(state.clone()),
Query(ListQuery {
kind: Some(crate::hub::TemplateKind::SinkTemplate),
}),
)
.await
.unwrap()
.0;
assert_eq!(sinks.templates.len(), 1);
assert_eq!(sinks.templates[0].id, "local-jsonl");
let err = trigger_pair(
State(state.clone()),
Extension(actor()),
Path("acme-exports".into()),
Json(TriggerBody::default()),
)
.await
.unwrap_err();
match err {
ServeError::Unprocessable { message, .. } => {
assert!(message.contains("sink"), "{message}")
}
other => panic!("expected 422, got {other:?}"),
}
let err = trigger_pair(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(TriggerBody {
sink: Some("local-jsonl".into()),
params: [("tag".to_string(), json!("a"))].into(),
..Default::default()
}),
)
.await
.unwrap_err();
match err {
ServeError::Unprocessable { message, .. } => {
assert!(message.contains("takes no sink"), "{message}")
}
other => panic!("expected 422, got {other:?}"),
}
let (code, resp) = trigger_pair(
State(state.clone()),
Extension(actor()),
Path("acme-exports".into()),
Json(TriggerBody {
sink: Some("local-jsonl".into()),
..Default::default()
}),
)
.await
.expect("trigger");
assert_eq!(code, StatusCode::ACCEPTED);
assert_eq!(resp.0.template_id, "acme-exports");
assert_eq!(resp.0.sink_template.as_deref(), Some("local-jsonl"));
assert_eq!(resp.0.sink_template_version, Some(1));
assert_eq!(resp.0.streams.len(), 2);
let rec = state
.history()
.get(&resp.0.run.run_id)
.await
.unwrap()
.expect("run record");
assert_eq!(rec.labels[LABEL_TEMPLATE], "acme-exports");
assert_eq!(rec.labels[LABEL_SINK_TEMPLATE], "local-jsonl");
assert_eq!(rec.labels[LABEL_SINK_TEMPLATE_VERSION], "1");
assert_eq!(rec.name.as_deref(), Some("acme-exports"));
register_body(
&state,
"kind: deployment\nname: ops\nstate: { type: memory }\n".into(),
)
.await;
let (_, resp) = trigger_pair(
State(state.clone()),
Extension(actor()),
Path("acme-exports".into()),
Json(TriggerBody {
sink: Some("local-jsonl".into()),
overlay: Some(OverlayRef::Id("ops".into())),
..Default::default()
}),
)
.await
.expect("trigger with a registered overlay");
assert_eq!(resp.0.overlay.as_deref(), Some("ops"));
assert_eq!(resp.0.overlay_version, Some(1));
assert_eq!(resp.0.overlay_contributes, vec!["pipeline.state"]);
let rec = state
.history()
.get(&resp.0.run.run_id)
.await
.unwrap()
.expect("run record");
assert_eq!(rec.labels[LABEL_OVERLAY], "ops");
assert_eq!(rec.labels[LABEL_OVERLAY_VERSION], "1");
let body: TriggerBody = serde_json::from_value(json!({
"sink": "local-jsonl",
"overlay": { "execution": { "max_concurrent": 1 } }
}))
.unwrap();
let (_, resp) = trigger_pair(
State(state.clone()),
Extension(actor()),
Path("acme-exports".into()),
Json(body),
)
.await
.expect("trigger with an inline overlay");
assert_eq!(resp.0.overlay.as_deref(), Some("inline"));
assert_eq!(resp.0.overlay_contributes, vec!["execution"]);
let rec = state
.history()
.get(&resp.0.run.run_id)
.await
.unwrap()
.expect("run record");
assert!(!rec.labels.contains_key(LABEL_OVERLAY_VERSION));
let idx = template_matrix(State(state.clone())).await.unwrap().0;
assert!(
idx["sources"]
.as_array()
.unwrap()
.iter()
.all(|s| s["name"] != "ops")
);
assert!(
idx["sinks"]
.as_array()
.unwrap()
.iter()
.all(|s| s["name"] != "ops")
);
}
#[tokio::test]
async fn the_matrix_composes_registered_sources_with_registered_sinks() {
let dir = tempfile::tempdir().unwrap();
let state = test_state();
let idx = template_matrix(State(state.clone())).await.unwrap().0;
assert_eq!(idx["sources"].as_array().map(Vec::len), Some(0));
assert_eq!(idx["matrix"].as_array().map(Vec::len), Some(0));
register_body(&state, source_template(dir.path())).await;
register_body(&state, sink_template(dir.path())).await;
register_demo(&state, &dir.path().join("o.jsonl")).await;
let idx = template_matrix(State(state.clone())).await.unwrap().0;
assert_eq!(idx["sources"][0]["name"], json!("acme-exports"));
assert_eq!(
idx["sources"][0]["file"],
json!("acme-exports"),
"registry id where a catalog has a path"
);
assert_eq!(idx["sinks"][0]["name"], json!("local-jsonl"));
let cell = &idx["matrix"][0];
assert_eq!(cell["compatible"], json!(true));
assert_eq!(cell["streams"][0]["write_mode"], json!("append"));
assert_eq!(cell["streams"][0]["satisfies"], json!("overwrite"));
let cmd = cell["command"].as_str().unwrap();
assert!(
cmd.starts_with("faucet template run acme-exports --sink local-jsonl"),
"{cmd}"
);
}
#[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 a_trigger_under_require_approval_becomes_a_change_request() {
let dir = tempfile::tempdir().unwrap();
let mut cfg = crate::serve::test_support::test_config();
cfg.require_approval = vec![crate::serve::changes::ChangeKind::Run];
let state = crate::serve::test_support::state_from(&cfg);
register_demo(&state, &dir.path().join("o.jsonl")).await;
let 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!(resp.status(), StatusCode::ACCEPTED);
let body = axum::body::to_bytes(resp.into_body(), 1 << 20)
.await
.unwrap();
let v: serde_json::Value = serde_json::from_slice(&body).unwrap();
assert_eq!(v["status"], "pending_approval", "{v}");
assert!(v["change_id"].is_string());
let ok = trigger_template(
State(test_state_with(&dir).await),
Extension(actor()),
Path("tpl-demo".into()),
Json(TriggerBody {
params: [("tag".to_string(), json!("beta"))].into(),
..Default::default()
}),
)
.await
.expect("trigger");
assert_eq!(ok.status(), StatusCode::ACCEPTED);
}
async fn test_state_with(dir: &tempfile::TempDir) -> ServerState {
let state = test_state();
register_demo(&state, &dir.path().join("o2.jsonl")).await;
state
}
#[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_pair(
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_pair(
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_pair(
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_pair(
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_pair(
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 a_deprecated_version_warns_skips_newest_and_cannot_be_launched() {
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 body = deprecate_version(
State(state.clone()),
Extension(actor()),
Path(("tpl-demo".into(), "newest".into())),
Json(DeprecateBody {
reason: Some("bad build".into()),
undo: false,
}),
)
.await
.unwrap()
.0;
assert_eq!(
(body["version"].as_u64(), body["deprecated"].as_bool()),
(Some(2), Some(true))
);
let st = crate::templates::template_state(&store(&state), "tpl-demo")
.await
.unwrap();
assert_eq!(st.newest, Some(1), "`newest` skips the retired v2");
assert_eq!(
st.status,
crate::serve::history::templates::TemplateStatus::Launched
);
let resp = trigger_pair(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(TriggerBody {
version: Some(VersionSelector::Pinned(2)),
params: [("tag".to_string(), json!("x"))].into(),
..Default::default()
}),
)
.await
.expect("a pinned deprecated version still runs")
.1
.0;
assert_eq!(
resp.deprecated.as_deref(),
Some("v2 is deprecated: bad build")
);
let err = launch_template(
State(state.clone()),
Extension(actor()),
Path("tpl-demo".into()),
Json(LaunchBody {
version: Some(VersionSelector::Pinned(2)),
}),
)
.await
.unwrap_err();
assert!(matches!(err, ServeError::Unprocessable { .. }), "{err:?}");
let revived = deprecate_version(
State(state.clone()),
Extension(actor()),
Path(("tpl-demo".into(), "2".into())),
Json(DeprecateBody {
reason: None,
undo: true,
}),
)
.await
.unwrap()
.0;
assert_eq!(revived["deprecated"], false);
let st = crate::templates::template_state(&store(&state), "tpl-demo")
.await
.unwrap();
assert_eq!(st.newest, Some(2));
for bad in ["9", "latest"] {
assert!(
deprecate_version(
State(state.clone()),
Extension(actor()),
Path(("tpl-demo".into(), bad.into())),
Json(DeprecateBody::default()),
)
.await
.is_err(),
"{bad}"
);
}
}
#[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_pair(
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_pair(
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]
);
}
}
#[cfg(feature = "templates-sync")]
#[derive(Debug, Default, Deserialize)]
pub struct SyncBody {
#[serde(default)]
pub origin: Option<String>,
#[serde(default)]
pub dry_run: bool,
}
#[cfg(feature = "templates-sync")]
#[derive(Debug, Serialize)]
pub struct SyncResponse {
pub dry_run: bool,
pub reports: Vec<crate::templates::sync::SyncReport>,
pub origin_errors: Vec<OriginError>,
}
#[cfg(feature = "templates-sync")]
#[derive(Debug, Serialize)]
pub struct OriginError {
pub origin: String,
pub error: String,
}
#[cfg(feature = "templates-sync")]
fn require_sync(
state: &ServerState,
) -> Result<std::sync::Arc<crate::templates::sync::SyncFile>, ServeError> {
state
.templates_sync()
.ok_or_else(|| ServeError::Unprocessable {
message:
"this server has no template origins — start it with `--templates-sync <file>`"
.into(),
details: None,
})
}
#[cfg(feature = "templates-sync")]
pub async fn sync_templates(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
body: Option<Json<SyncBody>>,
) -> Result<Json<SyncResponse>, ServeError> {
let body = body.map(|Json(b)| b).unwrap_or_default();
let file = require_sync(&state)?;
let results = crate::templates::sync::sync_all(
&store(&state),
&file,
body.origin.as_deref(),
body.dry_run,
Some(&actor.principal),
)
.await
.map_err(map_err)?;
let mut reports = Vec::new();
let mut origin_errors = Vec::new();
for r in results {
match r {
Ok(rep) => reports.push(rep),
Err((origin, e)) => origin_errors.push(OriginError {
origin,
error: e.to_string(),
}),
}
}
let result = if origin_errors.is_empty() && reports.iter().all(|r| r.failed() == 0) {
"ok"
} else {
"partial"
};
tracing::info!(
principal = %actor.principal,
origins = reports.len(),
errors = origin_errors.len(),
dry_run = body.dry_run,
"template sync requested"
);
crate::serve::audit::write(&state, &actor, "template.sync", None, None, result).await;
Ok(Json(SyncResponse {
dry_run: body.dry_run,
reports,
origin_errors,
}))
}
#[cfg(feature = "templates-sync")]
#[derive(Debug, Deserialize)]
pub struct PublishBody {
pub origin: String,
#[serde(default)]
pub version: VersionSelector,
}
#[cfg(feature = "templates-sync")]
pub async fn publish_template(
State(state): State<ServerState>,
Extension(actor): Extension<AuthContext>,
Path(id): Path<String>,
Json(body): Json<PublishBody>,
) -> Result<Json<crate::templates::sync::PublishReport>, ServeError> {
let file = require_sync(&state)?;
let report =
crate::templates::sync::publish(&store(&state), &file, &id, &body.origin, body.version)
.await
.map_err(map_err)?;
tracing::info!(
principal = %actor.principal,
template = %id,
version = report.version,
origin = %report.origin,
"template published"
);
crate::serve::audit::write(&state, &actor, "template.publish", None, None, "ok").await;
Ok(Json(report))
}