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,
callback: None,
}
}
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 text =
match resolve_config_text(&compiled.spec.config, &event, compiled.name(), fired_at).await {
Ok(t) => t,
Err(e) => {
metrics::error(compiled.name(), kind);
return FireOutcome::Error(e);
}
};
let req = build_submit_request(compiled, &event, text, fired_at);
let actor = crate::serve::rbac::AuthContext::trigger(compiled.name());
match runner::submit(state.clone(), req, actor).await {
Ok(resp) => {
metrics::enqueued(compiled.name());
FireOutcome::Enqueued(resp.run_id)
}
Err(crate::serve::error::ServeError::QueueFull { .. }) => {
metrics::dropped(compiled.name(), "queue_full");
FireOutcome::Dropped("queue_full")
}
Err(crate::serve::error::ServeError::Conflict(_)) => {
metrics::coalesced(compiled.name());
FireOutcome::Coalesced
}
Err(e) => {
metrics::error(compiled.name(), kind);
FireOutcome::Error(e.api_error().error.message)
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::serve::triggers::spec::{RunTemplate, TriggerKind, TriggerSpec};
fn compiled_webhook() -> CompiledTrigger {
CompiledTrigger {
spec: TriggerSpec {
name: "hook".into(),
enabled: true,
config: PipelineRef::Path("/tmp/x.yaml".into()),
run: RunTemplate {
name: Some("{name}:{object_key}".into()),
labels: Default::default(),
timeout_secs: Some(60),
},
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);
}
}