use fraiseql_core::schema::{CompiledSchema, MutationOperation};
use fraiseql_functions::{EntityEvent, EventKind, EventPayload, FunctionModule};
use crate::subsystems::BeforeMutationHooks;
pub struct AfterMutationDispatch {
pub module: FunctionModule,
pub payload: EventPayload,
}
pub const fn event_kind_for(operation: &MutationOperation) -> Option<EventKind> {
match operation {
MutationOperation::Insert { .. } => Some(EventKind::Insert),
MutationOperation::Update { .. } => Some(EventKind::Update),
MutationOperation::Delete { .. } => Some(EventKind::Delete),
_ => None,
}
}
pub fn plan_after_mutation_dispatch(
hooks: &BeforeMutationHooks,
schema: &CompiledSchema,
mutation_name: &str,
response_data: &serde_json::Value,
) -> Vec<AfterMutationDispatch> {
let Some(definition) = schema.find_mutation(mutation_name) else {
return Vec::new();
};
let Some(event_kind) = event_kind_for(&definition.operation) else {
return Vec::new();
};
let entity_value = response_data
.get("data")
.and_then(|data| data.get(mutation_name))
.filter(|value| !value.is_null())
.cloned();
let (old, new) = match event_kind {
EventKind::Delete => (entity_value, None),
_ => (None, entity_value),
};
let event = EntityEvent {
entity: definition.return_type.clone(),
event_kind,
old,
new,
timestamp: chrono::Utc::now(),
};
hooks
.observer
.find_after_mutation_triggers(&hooks.trigger_registry, &event)
.into_iter()
.filter_map(|trigger| {
let module = hooks.module_registry.get(&trigger.function_name)?.clone();
let payload = trigger.build_payload(&event);
Some(AfterMutationDispatch { module, payload })
})
.collect()
}
#[cfg(feature = "functions-runtime")]
pub fn spawn_after_mutation(hooks: &BeforeMutationHooks, plans: Vec<AfterMutationDispatch>) {
let config = host_context_config();
let limits = fraiseql_functions::ResourceLimits::default();
for plan in plans {
let observer = std::sync::Arc::clone(&hooks.observer);
let config = config.clone();
let limits = limits.clone();
tokio::spawn(async move {
let function_name = plan.module.name.clone();
let host: std::sync::Arc<
dyn fraiseql_functions::runtime::wasm::host_bridge::DynHostContext,
> = std::sync::Arc::new(fraiseql_functions::host::live::LiveHostContext::new(
plan.payload.clone(),
config,
));
match observer.invoke_with_context(&plan.module, plan.payload, host, limits).await {
Ok(_) => {
tracing::debug!(function = %function_name, "after:mutation function dispatched");
},
Err(error) => {
tracing::error!(
error = %error,
function = %function_name,
"after:mutation function failed",
);
},
}
});
}
}
#[cfg(feature = "functions-runtime")]
fn host_context_config() -> fraiseql_functions::host::live::HostContextConfig {
let mut config = fraiseql_functions::host::live::HostContextConfig::default();
if let Ok(domains) = std::env::var("FRAISEQL_FUNCTIONS_ALLOWED_DOMAINS") {
config.allowed_domains = domains
.split(',')
.map(str::trim)
.filter(|domain| !domain.is_empty())
.map(String::from)
.collect();
}
config
}
#[cfg(test)]
mod tests;