use crate::HttpError;
use ares_store::schedules as db_schedules;
use ares_agent::trigger;
use ares_types::types::AppError;
use axum::{
extract::State,
http::{HeaderMap, StatusCode},
Json,
};
use serde::Deserialize;
use std::sync::Arc;
use cordis::Context;
#[derive(Debug, Deserialize)]
pub struct FieldChangeEvent {
pub tenant_id: String,
pub table: String,
pub column: String,
pub record_id: String,
pub old_value: serde_json::Value,
pub new_value: serde_json::Value,
}
pub async fn handle_field_change(
State(ctx): State<Arc<Context>>,
headers: HeaderMap,
Json(payload): Json<FieldChangeEvent>,
) -> crate::Result<StatusCode> {
verify_webhook_secret(&headers)?;
if let Some(svc) = ctx.get::<ares_agent::trigger::TriggerService>() {
svc.dispatch_field_change(
&payload.tenant_id,
&payload.table,
&payload.column,
&payload.record_id,
payload.old_value.clone(),
payload.new_value.clone(),
&ctx,
)
.await
.map_err(|e| AppError::Internal(e.to_string()))?;
return Ok(StatusCode::OK);
}
let __pool_1 = ctx.get::<ares_store::TenantDb>().expect("not provided").pool().clone();
let store = db_schedules::EventTriggerStore::new(&__pool_1);
let triggers = store
.list_by_event_type(&payload.tenant_id, "field_change")
.await?;
let matching: Vec<_> = triggers
.into_iter()
.filter(|t| t.enabled)
.filter(|t| {
let table_match = t
.event_config
.get("table")
.and_then(|v| v.as_str())
.map(|tbl| tbl == payload.table)
.unwrap_or(false);
let column_match = t
.event_config
.get("column")
.and_then(|v| v.as_str())
.map(|col| col == payload.column)
.unwrap_or(false);
table_match && column_match
})
.collect();
let app_state = ctx.clone();
for trigger in matching {
let context = serde_json::json!({
"event": "field_change",
"table": payload.table,
"column": payload.column,
"record_id": payload.record_id,
"old_value": payload.old_value,
"new_value": payload.new_value,
});
let message = serde_json::to_string(&context).unwrap_or_default();
if let Err(e) =
trigger::execute_triggered_agent(&trigger, &message, &app_state).await
{
tracing::warn!(
trigger_id = %trigger.id,
agent = %trigger.target_agent,
error = %e,
"Field-change trigger execution failed"
);
}
}
Ok(StatusCode::OK)
}
fn verify_webhook_secret(headers: &HeaderMap) -> crate::Result<()> {
let expected = std::env::var("WEBHOOK_SECRET").unwrap_or_default();
if expected.is_empty() {
return Ok(());
}
let provided = headers
.get("X-Webhook-Secret")
.and_then(|h| h.to_str().ok())
.unwrap_or("");
if provided == expected {
Ok(())
} else {
Err(HttpError::from(AppError::Auth("Invalid webhook secret".to_string())))
}
}
#[cfg(test)]
mod tests {
use super::*;
use axum::http::HeaderValue;
use std::sync::Mutex;
static WEBHOOK_SECRET_ENV_LOCK: Mutex<()> = Mutex::new(());
#[test]
fn verify_webhook_secret_empty_env_allows_all() {
let _guard = WEBHOOK_SECRET_ENV_LOCK.lock().expect("env lock poisoned");
std::env::remove_var("WEBHOOK_SECRET");
let mut headers = HeaderMap::new();
headers.insert("X-Webhook-Secret", HeaderValue::from_static("anything"));
assert!(verify_webhook_secret(&headers).is_ok());
}
#[test]
fn verify_webhook_secret_rejects_mismatch() {
let _guard = WEBHOOK_SECRET_ENV_LOCK.lock().expect("env lock poisoned");
std::env::set_var("WEBHOOK_SECRET", "secret456");
let mut headers = HeaderMap::new();
headers.insert("X-Webhook-Secret", HeaderValue::from_static("wrong"));
assert!(verify_webhook_secret(&headers).is_err());
std::env::remove_var("WEBHOOK_SECRET");
}
#[test]
fn verify_webhook_secret_accepts_match() {
let _guard = WEBHOOK_SECRET_ENV_LOCK.lock().expect("env lock poisoned");
std::env::set_var("WEBHOOK_SECRET", "secret456");
let mut headers = HeaderMap::new();
headers.insert("X-Webhook-Secret", HeaderValue::from_static("secret456"));
assert!(verify_webhook_secret(&headers).is_ok());
std::env::remove_var("WEBHOOK_SECRET");
}
}