#[cfg(feature = "functions-runtime")]
use fraiseql_core::runtime::AfterMutationObserver;
use fraiseql_core::{runtime::CommittedMutation, schema::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,
mutation: &CommittedMutation<'_>,
) -> Vec<AfterMutationDispatch> {
let Some(event_kind) = event_kind_for(mutation.operation) else {
return Vec::new();
};
let entity_value = Some(mutation.entity).filter(|value| !value.is_null()).cloned();
let (old, new) = match event_kind {
EventKind::Delete => (entity_value, None),
_ => (None, entity_value),
};
let event = EntityEvent {
entity: mutation.entity_type.to_string(),
event_kind,
old,
new,
timestamp: chrono::Utc::now(),
};
hooks
.observer
.find_after_mutation_triggers(&hooks.trigger_registry, &event)
.into_iter()
.filter(|trigger| {
let held = trigger.predicates_hold(&event);
if !held {
crate::function_metrics::record_predicate_skip(&trigger.function_name);
}
held
})
.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 struct FunctionDispatchObserver {
hooks: std::sync::Arc<BeforeMutationHooks>,
executor:
std::sync::OnceLock<std::sync::Weak<arc_swap::ArcSwap<fraiseql_core::runtime::Executor>>>,
}
#[cfg(feature = "functions-runtime")]
impl FunctionDispatchObserver {
#[must_use]
pub const fn new(hooks: std::sync::Arc<BeforeMutationHooks>) -> Self {
Self {
hooks,
executor: std::sync::OnceLock::new(),
}
}
pub fn bind_executor(
&self,
executor: &std::sync::Arc<arc_swap::ArcSwap<fraiseql_core::runtime::Executor>>,
) {
let _ = self.executor.set(std::sync::Arc::downgrade(executor));
}
#[cfg(test)]
pub(crate) fn bound_executor(
&self,
) -> Option<std::sync::Arc<arc_swap::ArcSwap<fraiseql_core::runtime::Executor>>> {
self.executor.get().and_then(std::sync::Weak::upgrade)
}
fn plans_for(&self, mutation: &CommittedMutation<'_>) -> Vec<AfterMutationDispatch> {
if mutation.dispatch_depth > 0 {
return Vec::new();
}
plan_after_mutation_dispatch(&self.hooks, mutation)
}
}
#[cfg(feature = "functions-runtime")]
impl AfterMutationObserver for FunctionDispatchObserver {
fn on_committed(&self, mutation: &CommittedMutation<'_>) {
let plans = self.plans_for(mutation);
if plans.is_empty() {
return;
}
let query_executor_factory = self
.executor
.get()
.and_then(std::sync::Weak::upgrade)
.map(make_query_executor_factory);
spawn_after_mutation(
&self.hooks,
plans,
query_executor_factory,
mutation.security_ctx.cloned(),
);
}
}
#[cfg_attr(not(feature = "inbound"), allow(dead_code))]
pub fn plan_after_ingest_dispatch(
hooks: &BeforeMutationHooks,
message: &fraiseql_functions::InboundMessage,
) -> Vec<AfterMutationDispatch> {
hooks
.trigger_registry
.find_ingest_triggers(message)
.into_iter()
.filter_map(|trigger| {
let module = hooks.module_registry.get(&trigger.function_name)?.clone();
let payload = trigger.build_payload(message);
Some(AfterMutationDispatch { module, payload })
})
.collect()
}
pub const CAPTURED_WRITE_MARKER: &str = "fallback_trigger";
#[must_use]
pub fn plan_after_capture_dispatch(
hooks: &BeforeMutationHooks,
event: &EntityEvent,
cdc_source: Option<&str>,
) -> Vec<AfterMutationDispatch> {
if cdc_source != Some(CAPTURED_WRITE_MARKER) {
return Vec::new();
}
hooks
.observer
.find_after_capture_triggers(&hooks.trigger_registry, event)
.into_iter()
.filter(|trigger| {
let held = trigger.predicates_hold(event);
if !held {
crate::function_metrics::record_predicate_skip(&trigger.function_name);
}
held
})
.filter_map(|trigger| {
let module = hooks.module_registry.get(&trigger.function_name)?.clone();
let payload = EventPayload {
trigger_type: format!("after:capture:{}", trigger.function_name),
entity: event.entity.clone(),
event_kind: event.event_kind.to_string(),
data: serde_json::json!({
"event_kind": event.event_kind.as_str(),
"old": event.old,
"new": event.new,
}),
timestamp: event.timestamp,
};
Some(AfterMutationDispatch { module, payload })
})
.collect()
}
#[cfg(feature = "functions-runtime")]
#[must_use]
pub fn observer_event_to_capture(
event: &fraiseql_observers::EntityEvent,
) -> Option<(EntityEvent, Option<String>)> {
use fraiseql_observers::EventKind as ObserverEventKind;
let event_kind = match event.event_type {
ObserverEventKind::Created => EventKind::Insert,
ObserverEventKind::Updated => EventKind::Update,
ObserverEventKind::Deleted => EventKind::Delete,
ObserverEventKind::Custom => return None,
};
let (old, new) = match event_kind {
EventKind::Delete => (Some(event.data.clone()), None),
_ => (None, Some(event.data.clone())),
};
let fn_event = EntityEvent {
entity: event.entity_type.clone(),
event_kind,
old,
new,
timestamp: event.timestamp,
};
Some((fn_event, event.cdc_source.clone()))
}
#[cfg(feature = "functions-runtime")]
#[derive(Debug, Clone)]
pub struct FunctionDispatchSetting {
pub re_runnable: bool,
pub policy: fraiseql_observers::DispatchPolicy,
}
#[cfg(feature = "functions-runtime")]
impl Default for FunctionDispatchSetting {
fn default() -> Self {
Self {
re_runnable: false,
policy: fraiseql_observers::DispatchPolicy::new(
fraiseql_observers::RetryConfig::default(),
fraiseql_observers::FailurePolicy::Dlq,
),
}
}
}
#[cfg(feature = "functions-runtime")]
#[derive(Debug, Clone)]
pub struct DispatchDefaults {
pub retry: fraiseql_observers::RetryConfig,
pub dlq_max_size: Option<usize>,
}
#[cfg(feature = "functions-runtime")]
impl DispatchDefaults {
#[must_use]
pub fn from_env() -> Self {
Self::from_getter(|key| std::env::var(key).ok())
}
#[must_use]
pub fn from_getter(get: impl Fn(&str) -> Option<String>) -> Self {
let mut retry = fraiseql_observers::RetryConfig::default();
if let Some(value) =
get("FRAISEQL_FUNCTIONS_RETRY_MAX_ATTEMPTS").and_then(|s| s.parse().ok())
{
retry.max_attempts = value;
}
if let Some(value) =
get("FRAISEQL_FUNCTIONS_RETRY_INITIAL_DELAY_MS").and_then(|s| s.parse().ok())
{
retry.initial_delay_ms = value;
}
if let Some(value) =
get("FRAISEQL_FUNCTIONS_RETRY_MAX_DELAY_MS").and_then(|s| s.parse().ok())
{
retry.max_delay_ms = value;
}
let dlq_max_size = get("FRAISEQL_FUNCTIONS_DLQ_MAX_SIZE").and_then(|s| s.parse().ok());
Self {
retry,
dlq_max_size,
}
}
}
#[cfg(feature = "functions-runtime")]
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum DlqStoreKind {
#[default]
Memory,
Postgres,
}
#[cfg(feature = "functions-runtime")]
impl DlqStoreKind {
#[must_use]
pub fn resolve(compiled: Option<&str>, get_env: impl Fn(&str) -> Option<String>) -> Self {
let Some(raw) =
get_env("FRAISEQL_FUNCTIONS_DLQ_STORE").or_else(|| compiled.map(str::to_owned))
else {
return Self::Memory;
};
match raw.trim().to_ascii_lowercase().as_str() {
"postgres" | "pg" => Self::Postgres,
"memory" | "" => Self::Memory,
other => {
tracing::warn!(
value = other,
"unknown [functions] dlq_store — using the in-memory DLQ (\"memory\" | \"postgres\")"
);
Self::Memory
},
}
}
}
#[cfg(feature = "functions-runtime")]
#[must_use]
pub fn resolve_dispatch_settings(
definitions: &[fraiseql_functions::FunctionDefinition],
defaults: &DispatchDefaults,
) -> std::collections::HashMap<String, FunctionDispatchSetting> {
definitions
.iter()
.map(|definition| {
let retry = definition.retry.clone().unwrap_or_else(|| defaults.retry.clone());
let setting = FunctionDispatchSetting {
re_runnable: definition.re_runnable,
policy: fraiseql_observers::DispatchPolicy::new(
retry,
fraiseql_observers::FailurePolicy::Dlq,
),
};
(definition.name.clone(), setting)
})
.collect()
}
#[cfg(feature = "functions-runtime")]
pub type QueryExecutorFactory = std::sync::Arc<
dyn Fn(
fraiseql_core::security::SecurityContext,
) -> std::sync::Arc<dyn fraiseql_functions::host::live::QueryExecutor>
+ Send
+ Sync,
>;
#[cfg(feature = "functions-runtime")]
#[must_use]
pub fn make_query_executor_factory(
executor: std::sync::Arc<arc_swap::ArcSwap<fraiseql_core::runtime::Executor>>,
) -> QueryExecutorFactory {
std::sync::Arc::new(move |identity| {
std::sync::Arc::new(crate::query_bridge::RunAsQueryExecutor::new(
std::sync::Arc::clone(&executor),
identity,
)) as std::sync::Arc<dyn fraiseql_functions::host::live::QueryExecutor>
})
}
#[cfg(feature = "functions-runtime")]
#[must_use]
fn dispatch_idempotency_token(
key: Option<&[u8]>,
source: fraiseql_observers::DispatchSource,
function_name: &str,
payload: &EventPayload,
) -> String {
let trigger_identity =
format!("{}:{}:{}", payload.trigger_type, payload.entity, payload.event_kind);
fraiseql_observers::derive_idempotency_token(
key,
source,
function_name,
&trigger_identity,
&payload.data,
)
}
#[cfg(feature = "functions-runtime")]
#[derive(Clone)]
struct DurableDispatcher {
observer: std::sync::Arc<fraiseql_functions::FunctionObserver>,
host_config: fraiseql_functions::host::live::HostContextConfig,
limits: fraiseql_functions::ResourceLimits,
dlq: std::sync::Arc<dyn fraiseql_observers::DeadLetterQueue>,
source: fraiseql_observers::DispatchSource,
sender_resolver: Option<std::sync::Arc<dyn fraiseql_functions::SenderIdentityResolver>>,
email_transport: Option<std::sync::Arc<dyn fraiseql_functions::EmailTransport>>,
idempotency_key: Option<std::sync::Arc<[u8]>>,
query_executor_factory: Option<QueryExecutorFactory>,
run_as: Option<fraiseql_functions::RunAs>,
caller: Option<fraiseql_core::security::SecurityContext>,
}
#[cfg(feature = "functions-runtime")]
impl DurableDispatcher {
async fn invoke_once(
&self,
module: &FunctionModule,
payload: EventPayload,
idempotency_token: &str,
) -> fraiseql_error::Result<fraiseql_functions::FunctionResult> {
let live = self.build_host(&module.name, payload.clone(), idempotency_token);
let host: std::sync::Arc<dyn fraiseql_functions::host::dyn_context::DynHostContext> =
std::sync::Arc::new(live);
self.observer
.invoke_with_context(module, payload, host, self.limits.clone())
.await
}
fn build_host(
&self,
function_name: &str,
payload: EventPayload,
idempotency_token: &str,
) -> fraiseql_functions::host::live::LiveHostContext {
let run_as_identity = self
.run_as
.clone()
.unwrap_or_default()
.identity(function_name, idempotency_token);
let host_identity = self.caller.clone().unwrap_or_else(|| run_as_identity.clone());
let mut live =
fraiseql_functions::host::live::LiveHostContext::new(payload, self.host_config.clone())
.with_idempotency_token(idempotency_token)
.with_security_context(host_identity);
if let (Some(resolver), Some(transport)) =
(self.sender_resolver.as_ref(), self.email_transport.as_ref())
{
live =
live.with_email(std::sync::Arc::clone(resolver), std::sync::Arc::clone(transport));
}
if let Some(factory) = self.query_executor_factory.as_ref() {
live = live.with_executor(factory(run_as_identity));
}
live
}
async fn dispatch(
&self,
module: FunctionModule,
payload: EventPayload,
setting: &FunctionDispatchSetting,
) {
let function_name = module.name.clone();
let started = std::time::Instant::now();
let trigger_kind = crate::function_metrics::trigger_kind(self.source);
let idempotency_token = dispatch_idempotency_token(
self.idempotency_key.as_deref(),
self.source,
&function_name,
&payload,
);
if setting.re_runnable {
let (result, elapsed) = match self
.invoke_once(&module, payload, &idempotency_token)
.await
{
Ok(_) => {
tracing::debug!(function = %function_name, "re-runnable function dispatched");
(crate::function_metrics::RESULT_OK, started.elapsed().as_secs_f64())
},
Err(error) => {
tracing::warn!(
error = %error,
function = %function_name,
"re-runnable function failed (not retried)"
);
(crate::function_metrics::RESULT_ERROR, started.elapsed().as_secs_f64())
},
};
crate::function_metrics::record_dispatch(&function_name, trigger_kind, result, elapsed);
return;
}
let trigger_type = payload.trigger_type.clone();
let attempts = std::sync::atomic::AtomicU32::new(0);
let result =
fraiseql_observers::run_with_retry(
&setting.policy,
|error: &fraiseql_error::FraiseQLError| !error.is_client_error(),
|error: &fraiseql_error::FraiseQLError| match error {
fraiseql_error::FraiseQLError::ServiceUnavailable {
retry_after: Some(secs),
..
} => Some(std::time::Duration::from_secs(*secs)),
_ => None,
},
|n| {
attempts.store(n, std::sync::atomic::Ordering::Relaxed);
let attempt_module = module.clone();
let attempt_payload = payload.clone();
let attempt_token = idempotency_token.clone();
async move {
self.invoke_once(&attempt_module, attempt_payload, &attempt_token).await
}
},
)
.await;
let Err(error) = result else {
tracing::debug!(function = %function_name, "function dispatched");
crate::function_metrics::record_dispatch(
&function_name,
trigger_kind,
crate::function_metrics::RESULT_OK,
started.elapsed().as_secs_f64(),
);
return;
};
let attempts = attempts.load(std::sync::atomic::Ordering::Relaxed);
crate::function_metrics::record_dispatch(
&function_name,
trigger_kind,
crate::function_metrics::RESULT_DEAD_LETTERED,
started.elapsed().as_secs_f64(),
);
let dead_letter_token = idempotency_token.clone();
let record = fraiseql_observers::FunctionDispatchRecord::new(
self.source,
function_name.clone(),
trigger_type,
idempotency_token,
serde_json::to_value(&payload).unwrap_or(serde_json::Value::Null),
error.to_string(),
attempts,
);
match self.dlq.push_function(record).await {
Ok(_) => tracing::error!(
error = %error,
function = %function_name,
idempotency_token = %dead_letter_token,
attempts,
"function dead-lettered after exhausting retries"
),
Err(dlq_error) => tracing::error!(
error = %error,
dlq_error = %dlq_error,
function = %function_name,
idempotency_token = %dead_letter_token,
"function failed and could not be dead-lettered"
),
}
}
}
#[cfg(feature = "functions-runtime")]
pub fn spawn_after_mutation(
hooks: &BeforeMutationHooks,
plans: Vec<AfterMutationDispatch>,
query_executor_factory: Option<QueryExecutorFactory>,
caller: Option<fraiseql_core::security::SecurityContext>,
) {
spawn_dispatch(
hooks,
plans,
fraiseql_observers::DispatchSource::AfterMutation,
query_executor_factory,
caller,
);
}
#[cfg(feature = "inbound")]
pub fn spawn_after_ingest(
hooks: &BeforeMutationHooks,
plans: Vec<AfterMutationDispatch>,
query_executor_factory: Option<QueryExecutorFactory>,
) {
spawn_dispatch(
hooks,
plans,
fraiseql_observers::DispatchSource::AfterIngest,
query_executor_factory,
None,
);
}
#[cfg(feature = "functions-runtime")]
pub fn spawn_after_capture(
hooks: &BeforeMutationHooks,
plans: Vec<AfterMutationDispatch>,
query_executor_factory: Option<QueryExecutorFactory>,
) {
spawn_dispatch(
hooks,
plans,
fraiseql_observers::DispatchSource::AfterCapture,
query_executor_factory,
None,
);
}
#[cfg(feature = "functions-runtime")]
fn spawn_dispatch(
hooks: &BeforeMutationHooks,
plans: Vec<AfterMutationDispatch>,
source: fraiseql_observers::DispatchSource,
query_executor_factory: Option<QueryExecutorFactory>,
caller: Option<fraiseql_core::security::SecurityContext>,
) {
let dispatcher = DurableDispatcher {
observer: std::sync::Arc::clone(&hooks.observer),
host_config: host_context_config(),
limits: fraiseql_functions::ResourceLimits::default(),
dlq: std::sync::Arc::clone(&hooks.dlq),
source,
sender_resolver: hooks.sender_resolver.clone(),
email_transport: hooks.email_transport.clone(),
idempotency_key: hooks.idempotency_key.clone(),
query_executor_factory,
run_as: None,
caller,
};
for plan in plans {
let setting = hooks.dispatch_settings.get(&plan.module.name).cloned().unwrap_or_default();
let mut dispatcher = dispatcher.clone();
dispatcher.run_as = hooks.run_as.get(&plan.module.name).cloned();
tokio::spawn(async move {
dispatcher.dispatch(plan.module, plan.payload, &setting).await;
});
}
}
#[cfg(feature = "functions-runtime")]
pub fn host_context_config() -> fraiseql_functions::host::live::HostContextConfig {
host_context_config_from(|key| std::env::var(key).ok())
}
#[cfg(feature = "functions-runtime")]
fn host_context_config_from(
get: impl Fn(&str) -> Option<String>,
) -> fraiseql_functions::host::live::HostContextConfig {
let mut config = fraiseql_functions::host::live::HostContextConfig::default();
if let Some(domains) = get("FRAISEQL_FUNCTIONS_ALLOWED_DOMAINS") {
config.allowed_domains = domains
.split(',')
.map(str::trim)
.filter(|domain| !domain.is_empty())
.map(String::from)
.collect();
}
if let Some(names) = get("FRAISEQL_FUNCTIONS_ALLOWED_ENV_VARS") {
config.allowed_env_vars = names
.split(',')
.map(str::trim)
.filter(|name| !name.is_empty())
.map(String::from)
.collect();
}
config
}
#[cfg(test)]
mod tests;