use super::compiled::CompiledTrigger;
use super::context::{self, TriggerEvent};
use super::{metrics, spec::PipelineRef};
use crate::serve::runner::{self, ConfigFormatWire, SubmitRequest};
use crate::serve::state::ServerState;
#[derive(Debug)]
pub enum FireOutcome {
Enqueued(String),
Coalesced,
Dropped(&'static str),
Error(String),
}
impl FireOutcome {
pub fn committed(&self) -> bool {
matches!(self, FireOutcome::Enqueued(_) | FireOutcome::Coalesced)
}
}
pub async fn resolve_config_text(
config: &PipelineRef,
event: &TriggerEvent,
name: &str,
fired_at: &str,
) -> Result<String, String> {
let raw = match config {
PipelineRef::Path(p) => tokio::fs::read_to_string(p)
.await
.map_err(|e| format!("reading pipeline config '{p}': {e}"))?,
PipelineRef::Inline(v) => {
serde_yaml::to_string(v).map_err(|e| format!("serializing inline pipeline: {e}"))?
}
};
context::substitute(&raw, event, name, fired_at)
}
pub fn build_submit_request(
compiled: &CompiledTrigger,
event: &TriggerEvent,
config_text: String,
fired_at: &str,
) -> SubmitRequest {
let name = compiled.name();
let mut labels = context::labels(name, event);
labels.extend(compiled.spec.run.labels.clone());
let run_name = compiled
.spec
.run
.name
.as_deref()
.map(|tpl| context::render_name(tpl, event, name, fired_at))
.unwrap_or_else(|| name.to_string());
SubmitRequest {
config: config_text,
config_format: ConfigFormatWire::Yaml,
name: Some(run_name),
labels,
timeout_secs: compiled.spec.run.timeout_secs,
doctor_first: false,
idempotency_key: Some(context::idempotency_key(name, event)),
clock: None,
concurrency: None,
callback: None,
require_approval: false,
reason: None,
budget: None,
approved_change: None,
selection: compiled.spec.run.selection.clone(),
}
}
pub async fn fire(
state: &ServerState,
compiled: &CompiledTrigger,
event: TriggerEvent,
fired_at: &str,
) -> FireOutcome {
let kind = compiled.kind_label();
metrics::fired(compiled.name(), kind);
let targets = match fire_targets(state, compiled).await {
Ok(t) => t,
Err(e) => {
metrics::error(compiled.name(), kind);
return FireOutcome::Error(e);
}
};
if targets.is_empty() {
tracing::info!(
trigger = compiled.name(),
"trigger fired with no tenants to run for"
);
return FireOutcome::Coalesced;
}
let mut outcomes = Vec::with_capacity(targets.len());
for tenant in targets {
let outcome = fire_one(state, compiled, &event, fired_at, tenant.as_deref()).await;
if let (Some(t), FireOutcome::Error(e)) = (&tenant, &outcome) {
tracing::warn!(trigger = compiled.name(), tenant = %t, error = %e, "tenant fire failed");
}
outcomes.push(outcome);
}
combine(outcomes)
}
pub fn combine(outcomes: Vec<FireOutcome>) -> FireOutcome {
if outcomes.len() == 1 {
return outcomes.into_iter().next().expect("one outcome");
}
if let Some(reason) = outcomes.iter().find_map(|o| match o {
FireOutcome::Dropped(r) => Some(*r),
_ => None,
}) {
return FireOutcome::Dropped(reason);
}
let errors: Vec<String> = outcomes
.iter()
.filter_map(|o| match o {
FireOutcome::Error(e) => Some(e.clone()),
_ => None,
})
.collect();
if !errors.is_empty() {
return FireOutcome::Error(errors.join("; "));
}
let ids: Vec<String> = outcomes
.into_iter()
.filter_map(|o| match o {
FireOutcome::Enqueued(id) => Some(id),
_ => None,
})
.collect();
if ids.is_empty() {
FireOutcome::Coalesced
} else {
FireOutcome::Enqueued(ids.join(","))
}
}
async fn fire_targets(
state: &ServerState,
compiled: &CompiledTrigger,
) -> Result<Vec<Option<String>>, String> {
let Some(selector) = &compiled.spec.tenants else {
return Ok(vec![None]);
};
#[cfg(feature = "tenants")]
{
let (run, skipped) = crate::serve::handlers::tenants::resolve_tenants(state, selector)
.await
.map_err(|e| e.api_error().error.message)?;
for s in skipped {
tracing::info!(
trigger = compiled.name(),
tenant = %s.tenant,
reason = s.reason.as_deref().unwrap_or(""),
"tenant skipped"
);
}
Ok(run.into_iter().map(Some).collect())
}
#[cfg(not(feature = "tenants"))]
{
let _ = (state, selector);
Err("`tenants:` needs a build with the `tenants` feature".into())
}
}
async fn fire_one(
state: &ServerState,
compiled: &CompiledTrigger,
event: &TriggerEvent,
fired_at: &str,
tenant: Option<&str>,
) -> FireOutcome {
let kind = compiled.kind_label();
let mut actor = crate::serve::rbac::AuthContext::trigger(compiled.name());
actor.tenant = tenant.map(str::to_string);
let idem = |key: String| match tenant {
Some(t) => format!("{key}:{t}"),
None => key,
};
let result = match (&compiled.spec.config, &compiled.spec.template) {
(Some(config), _) => {
let text = match resolve_config_text(config, event, compiled.name(), fired_at).await {
Ok(t) => t,
Err(e) => {
metrics::error(compiled.name(), kind);
return FireOutcome::Error(e);
}
};
let mut req = build_submit_request(compiled, event, text, fired_at);
req.idempotency_key = req.idempotency_key.map(idem);
runner::submit(state.clone(), req, actor)
.await
.map(|r| r.run_id)
}
(None, Some(tpl)) => {
submit_template(state, compiled, tpl, event, fired_at, actor, idem).await
}
(None, None) => Err(crate::serve::error::ServeError::BadConfig(
"trigger has neither config nor template".into(),
)),
};
match result {
Ok(run_id) => {
metrics::enqueued(compiled.name());
FireOutcome::Enqueued(run_id)
}
Err(crate::serve::error::ServeError::QueueFull { .. }) => {
metrics::dropped(compiled.name(), "queue_full");
FireOutcome::Dropped("queue_full")
}
Err(crate::serve::error::ServeError::TooManyRequests(_)) => {
metrics::dropped(compiled.name(), "tenant_limit");
FireOutcome::Dropped("tenant_limit")
}
Err(crate::serve::error::ServeError::Conflict(m)) => {
tracing::info!(trigger = compiled.name(), reason = %m, "trigger fire coalesced");
metrics::coalesced(compiled.name());
FireOutcome::Coalesced
}
Err(e) => {
metrics::error(compiled.name(), kind);
FireOutcome::Error(e.api_error().error.message)
}
}
}
#[cfg(feature = "templates")]
pub fn template_body(
compiled: &CompiledTrigger,
tpl: &super::spec::TemplateTrigger,
event: &TriggerEvent,
fired_at: &str,
) -> Result<crate::serve::handlers::templates::TriggerBody, String> {
fn selector(
field: &str,
v: &Option<String>,
) -> Result<Option<crate::serve::history::templates::VersionSelector>, String> {
v.as_ref()
.map(|s| {
serde_json::from_value(serde_json::Value::String(s.clone()))
.map_err(|e| format!("template {field} '{s}': {e}"))
})
.transpose()
}
let mut params = std::collections::BTreeMap::new();
for (k, v) in &tpl.params {
let v =
match v {
serde_json::Value::String(s) => serde_json::Value::String(
context::substitute_plain(s, event, compiled.name(), fired_at)?,
),
other => other.clone(),
};
params.insert(k.clone(), v);
}
let name = compiled.name();
let mut labels = context::labels(name, event);
labels.extend(compiled.spec.run.labels.clone());
Ok(crate::serve::handlers::templates::TriggerBody {
params,
version: selector("version", &tpl.version)?,
sink: tpl.sink.clone(),
sink_version: selector("sink_version", &tpl.sink_version)?,
overlay: tpl
.overlay
.as_ref()
.map(|o| crate::serve::handlers::templates::OverlayRef::Id(o.clone())),
overlay_version: selector("overlay_version", &tpl.overlay_version)?,
name: Some(
compiled
.spec
.run
.name
.as_deref()
.map(|t| context::render_name(t, event, name, fired_at))
.unwrap_or_else(|| name.to_string()),
),
labels,
timeout_secs: compiled.spec.run.timeout_secs,
idempotency_key: Some(context::idempotency_key(name, event)),
selection: compiled.spec.run.selection.clone(),
..Default::default()
})
}
#[cfg(feature = "templates")]
async fn submit_template(
state: &ServerState,
compiled: &CompiledTrigger,
tpl: &super::spec::TemplateTrigger,
event: &TriggerEvent,
fired_at: &str,
actor: crate::serve::rbac::AuthContext,
idem: impl Fn(String) -> String,
) -> Result<String, crate::serve::error::ServeError> {
use crate::serve::handlers::templates::{TriggerOutcome, trigger_template_outcome};
let mut body = template_body(compiled, tpl, event, fired_at)
.map_err(crate::serve::error::ServeError::BadConfig)?;
body.idempotency_key = body.idempotency_key.map(idem);
match trigger_template_outcome(state.clone(), actor, tpl.id.clone(), body).await? {
TriggerOutcome::Run(r) => Ok(r.run.run_id),
TriggerOutcome::PendingApproval(c) => Ok(format!("change:{}", c.id)),
}
}
#[cfg(not(feature = "templates"))]
async fn submit_template(
_state: &ServerState,
_compiled: &CompiledTrigger,
_tpl: &super::spec::TemplateTrigger,
_event: &TriggerEvent,
_fired_at: &str,
_actor: crate::serve::rbac::AuthContext,
_idem: impl Fn(String) -> String,
) -> Result<String, crate::serve::error::ServeError> {
Err(crate::serve::error::ServeError::BadConfig(
"a template trigger needs a build with the `templates` feature".into(),
))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn combine_prefers_drops_then_errors_then_ids() {
use FireOutcome::*;
assert!(matches!(combine(vec![Enqueued("a".into())]), Enqueued(ref id) if id == "a"));
assert!(matches!(
combine(vec![
Enqueued("a".into()),
Dropped("queue_full"),
Error("e".into())
]),
Dropped("queue_full")
));
assert!(matches!(
combine(vec![Enqueued("a".into()), Error("x".into()), Error("y".into())]),
Error(ref e) if e == "x; y"
));
assert!(matches!(
combine(vec![Enqueued("a".into()), Coalesced, Enqueued("b".into())]),
Enqueued(ref ids) if ids == "a,b"
));
assert!(matches!(combine(vec![Coalesced, Coalesced]), Coalesced));
}
#[cfg(feature = "templates")]
#[test]
fn template_bodies_carry_params_labels_and_the_tick_key() {
use crate::serve::triggers::spec::TemplateTrigger;
let mut c = compiled_webhook();
c.spec.run.name = Some("{name}@{tick}".into());
let tpl = TemplateTrigger {
id: "crm".into(),
version: Some("3".into()),
params: std::collections::BTreeMap::from([
("since".to_string(), serde_json::json!("${trigger.tick}")),
("n".to_string(), serde_json::json!(5)),
]),
sink: Some("wh".into()),
sink_version: Some("stable".into()),
overlay: Some("prod".into()),
overlay_version: None,
};
let e = TriggerEvent::Schedule {
tick: "2026-01-01T00:00:00Z".into(),
};
let b = template_body(&c, &tpl, &e, "now").unwrap();
assert_eq!(b.params["since"], "2026-01-01T00:00:00Z");
assert_eq!(b.params["n"], 5);
assert_eq!(b.name.as_deref(), Some("hook@2026-01-01T00:00:00Z"));
assert_eq!(
b.idempotency_key.as_deref(),
Some("trig:hook:2026-01-01T00:00:00Z")
);
assert_eq!(b.labels["faucet.trigger.type"], "schedule");
assert_eq!(b.sink.as_deref(), Some("wh"));
assert!(b.version.is_some() && b.sink_version.is_some() && b.overlay.is_some());
let bad = TemplateTrigger {
version: Some("latest".into()),
..tpl.clone()
};
assert!(
template_body(&c, &bad, &e, "now")
.unwrap_err()
.contains("version")
);
let bad_token = TemplateTrigger {
params: std::collections::BTreeMap::from([(
"x".to_string(),
serde_json::json!("${trigger.nope}"),
)]),
..tpl
};
assert!(template_body(&c, &bad_token, &e, "now").is_err());
}
#[cfg(feature = "tenants")]
#[tokio::test]
async fn a_tenant_fan_out_fires_once_per_tenant_and_skips_missing_ones() {
let state = crate::serve::test_support::test_state();
let now = chrono::Utc::now();
for id in ["acme", "globex"] {
state
.history()
.tenant_upsert(&crate::serve::history::tenants::TenantRecord {
id: id.into(),
name: None,
labels: Default::default(),
limits: Default::default(),
notifications: Vec::new(),
suspended: false,
created_at: now,
updated_at: now,
created_by: "t".into(),
})
.await
.unwrap();
}
let mut c = compiled_webhook();
c.spec.config = Some(PipelineRef::Path("/definitely/missing.yaml".into()));
c.spec.tenants = Some(crate::serve::history::tenants::TenantSelector::Named(vec![
"acme".into(),
"ghost".into(),
]));
assert_eq!(
fire_targets(&state, &c).await.unwrap(),
vec![Some("acme".into())]
);
c.spec.tenants = Some(crate::serve::history::tenants::TenantSelector::All(
"all".into(),
));
let event = TriggerEvent::Schedule { tick: "t".into() };
let out = fire(&state, &c, event.clone(), "now").await;
assert!(
matches!(out, FireOutcome::Error(ref e) if e.contains("missing.yaml")),
"{out:?}"
);
c.spec.tenants = Some(crate::serve::history::tenants::TenantSelector::Named(vec![
"ghost".into(),
]));
assert!(matches!(
fire(&state, &c, event, "now").await,
FireOutcome::Coalesced
));
}
use crate::serve::triggers::spec::{RunTemplate, TriggerKind, TriggerSpec};
fn compiled_webhook() -> CompiledTrigger {
CompiledTrigger {
spec: TriggerSpec {
name: "hook".into(),
enabled: true,
config: Some(PipelineRef::Path("/tmp/x.yaml".into())),
template: None,
tenants: None,
run: RunTemplate {
name: Some("{name}:{object_key}".into()),
labels: Default::default(),
timeout_secs: Some(60),
selection: None,
},
kind: TriggerKind::Webhook {
methods: vec!["POST".into()],
dedupe_header: None,
debounce_secs: 0,
},
},
webhook_path: Some("/v1/triggers/hook".into()),
}
}
#[test]
fn builds_request_with_labels_idem_and_timeout() {
let event = TriggerEvent::Object {
bucket: "b".into(),
key: "k".into(),
size: 1,
last_modified: "2026-06-12T00:00:00Z".into(),
};
let req = build_submit_request(&compiled_webhook(), &event, "version: 1".into(), "now");
assert_eq!(req.name.as_deref(), Some("hook:k"));
assert_eq!(req.timeout_secs, Some(60));
assert_eq!(
req.idempotency_key.as_deref(),
Some("trig:hook:b:k:2026-06-12T00:00:00Z")
);
assert_eq!(
req.labels.get("faucet.trigger.name").map(String::as_str),
Some("hook")
);
assert!(!req.doctor_first);
}
}