use std::{
collections::HashMap,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
time::Duration,
};
use arc_swap::ArcSwap;
use fraiseql_core::runtime::subscription::{ChangeSpineEnvelope, SubscriptionOperation};
use fraiseql_observers::{
ActionConfig as ObserverActionConfig, ActionExecutionDetail, ChangeLogListener,
ChangeLogListenerConfig, EntityEvent as ObserverEntityEvent, EventMatcher, FailurePolicy,
InMemoryTransport, ObserverDefinition, ObserverExecutor, RetryConfig as ObserverRetryConfig,
checkpoint::{
CheckpointState as ObserverCheckpointState, CheckpointStore, PostgresCheckpointStore,
},
config::{EmailSmtpConfig, TransportConfig, TransportKind},
transport::{EventFilter, EventTransport},
};
use futures::StreamExt;
use sqlx::PgPool;
use tokio::{
sync::{RwLock, mpsc, oneshot},
task::JoinHandle,
};
use tracing::{debug, error, info, warn};
use crate::{
ServerError,
observers::{Observer, ObserverRepository},
subscriptions::event_bridge::EntityEvent as BridgeEntityEvent,
};
#[cfg(test)]
mod tests;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ListenerSelection {
PostgresChangeLog,
TransportStream,
}
pub(crate) const fn listener_selection(kind: TransportKind) -> ListenerSelection {
match kind {
TransportKind::Postgres => ListenerSelection::PostgresChangeLog,
_ => ListenerSelection::TransportStream,
}
}
#[cfg(feature = "observers-nats")]
fn nats_config_from(
cfg: &fraiseql_observers::config::NatsTransportConfig,
) -> fraiseql_observers::transport::NatsConfig {
fraiseql_observers::transport::NatsConfig {
url: cfg.url.clone(),
stream_name: cfg.stream_name.clone(),
consumer_name: cfg.consumer_name.clone(),
subject_prefix: cfg.subject_prefix.clone(),
ack_wait_secs: cfg.jetstream.ack_wait_secs,
retention_max_messages: cfg.jetstream.max_msgs,
retention_max_bytes: cfg.jetstream.max_bytes,
..Default::default()
}
}
#[derive(Debug, Clone)]
pub struct ObserverRuntimeConfig {
pub pool: PgPool,
pub poll_interval_ms: u64,
pub batch_size: usize,
pub channel_capacity: usize,
pub auto_reload: bool,
pub reload_interval_secs: u64,
pub max_dlq_size: Option<usize>,
pub transport: TransportConfig,
pub email: Option<EmailSmtpConfig>,
pub log_payloads: bool,
pub listener_id: String,
pub redis: Option<fraiseql_observers::config::RedisConfig>,
}
impl ObserverRuntimeConfig {
#[must_use]
pub fn new(pool: PgPool) -> Self {
Self {
pool,
poll_interval_ms: 100,
batch_size: 100,
channel_capacity: 1000,
auto_reload: true,
reload_interval_secs: 60,
max_dlq_size: None,
transport: TransportConfig::default(),
email: None,
log_payloads: false,
listener_id: "change_log".to_string(),
redis: None,
}
}
#[must_use]
pub fn with_redis(mut self, redis: Option<fraiseql_observers::config::RedisConfig>) -> Self {
self.redis = redis;
self
}
#[must_use]
pub const fn with_poll_interval(mut self, ms: u64) -> Self {
self.poll_interval_ms = ms;
self
}
#[must_use]
pub const fn with_batch_size(mut self, size: usize) -> Self {
self.batch_size = size;
self
}
#[must_use]
pub const fn with_channel_capacity(mut self, capacity: usize) -> Self {
self.channel_capacity = capacity;
self
}
#[must_use]
pub const fn with_max_dlq_size(mut self, max: Option<usize>) -> Self {
self.max_dlq_size = max;
self
}
#[must_use]
pub fn with_transport(mut self, transport: TransportConfig) -> Self {
self.transport = transport;
self
}
#[must_use]
pub fn with_email(mut self, email: Option<EmailSmtpConfig>) -> Self {
self.email = email;
self
}
#[must_use]
pub const fn with_log_payloads(mut self, log_payloads: bool) -> Self {
self.log_payloads = log_payloads;
self
}
#[must_use]
pub fn with_listener_id(mut self, listener_id: impl Into<String>) -> Self {
self.listener_id = listener_id.into();
self
}
}
#[derive(Debug, Clone)]
pub struct RuntimeHealth {
pub running: bool,
pub observer_count: usize,
pub last_checkpoint: Option<i64>,
pub events_processed: u64,
pub errors: u64,
}
pub struct ObserverRuntime {
config: ObserverRuntimeConfig,
repository: ObserverRepository,
running: Arc<AtomicBool>,
task_handle: Option<JoinHandle<()>>,
shutdown_tx: Option<mpsc::Sender<()>>,
events_processed: Arc<std::sync::atomic::AtomicU64>,
errors: Arc<std::sync::atomic::AtomicU64>,
observer_count: Arc<std::sync::atomic::AtomicUsize>,
last_checkpoint: Arc<std::sync::atomic::AtomicI64>,
matcher: Arc<RwLock<Option<EventMatcher>>>,
executor: Arc<RwLock<Option<Arc<ObserverExecutor>>>>,
entity_type_index: Arc<ArcSwap<HashMap<(String, String), Vec<i64>>>>,
dlq: Arc<InMemoryDlq>,
event_bridge_sender: Option<mpsc::Sender<BridgeEntityEvent>>,
capture_dispatch: Option<CaptureDispatchFn>,
#[cfg(feature = "observers-cache")]
cache_invalidator: Option<Arc<fraiseql_observers::RedisCacheInvalidator>>,
}
pub type CaptureDispatchFn = Arc<dyn Fn(&ObserverEntityEvent) + Send + Sync>;
impl ObserverRuntime {
#[must_use]
pub fn new(config: ObserverRuntimeConfig) -> Self {
let repository = ObserverRepository::new(config.pool.clone());
let max_dlq_size = config.max_dlq_size;
Self {
config,
repository,
running: Arc::new(AtomicBool::new(false)),
task_handle: None,
shutdown_tx: None,
events_processed: Arc::new(std::sync::atomic::AtomicU64::new(0)),
errors: Arc::new(std::sync::atomic::AtomicU64::new(0)),
observer_count: Arc::new(std::sync::atomic::AtomicUsize::new(0)),
last_checkpoint: Arc::new(std::sync::atomic::AtomicI64::new(0)),
matcher: Arc::new(RwLock::new(None)),
executor: Arc::new(RwLock::new(None)),
entity_type_index: Arc::new(ArcSwap::from_pointee(HashMap::new())),
dlq: Arc::new(InMemoryDlq::new_with_max(max_dlq_size)),
event_bridge_sender: None,
capture_dispatch: None,
#[cfg(feature = "observers-cache")]
cache_invalidator: None,
}
}
pub fn set_event_bridge_sender(&mut self, sender: mpsc::Sender<BridgeEntityEvent>) {
self.event_bridge_sender = Some(sender);
}
pub fn set_capture_dispatch(&mut self, hook: CaptureDispatchFn) {
self.capture_dispatch = Some(hook);
}
async fn load_observers(
&self,
) -> Result<
(HashMap<String, ObserverDefinition>, HashMap<(String, String), Vec<i64>>),
ServerError,
> {
let query = crate::observers::ListObserversQuery {
page: 1,
page_size: 10000, entity_type: None,
event_type: None,
enabled: Some(true),
include_deleted: false,
};
let (observers, _total) = self.repository.list(&query, None).await?;
let mut definitions = HashMap::new();
let mut entity_type_index: HashMap<(String, String), Vec<i64>> = HashMap::new();
for observer in observers {
match Self::convert_observer(&observer) {
Ok(definition) => {
let entity_type =
observer.entity_type.clone().unwrap_or_else(|| "*".to_string());
let event_type =
observer.event_type.clone().unwrap_or_else(|| "INSERT".to_string());
entity_type_index
.entry((entity_type, event_type.to_uppercase()))
.or_default()
.push(observer.pk_observer);
definitions.insert(observer.name.clone(), definition);
},
Err(e) => {
warn!("Failed to convert observer {}: {}", observer.name, e);
},
}
}
info!("Loaded {} observers from database", definitions.len());
Ok((definitions, entity_type_index))
}
fn convert_observer(observer: &Observer) -> Result<ObserverDefinition, ServerError> {
let actions: Vec<ObserverActionConfig> = serde_json::from_value(observer.actions.clone())
.map_err(|e| {
ServerError::Validation(format!(
"Failed to parse actions for observer {}: {}",
observer.name, e
))
})?;
let retry_config: ObserverRetryConfig =
match serde_json::from_value(observer.retry_config.clone()) {
Ok(cfg) => cfg,
Err(e) => {
warn!(
observer = %observer.name,
error = %e,
"Observer retry_config could not be deserialized; using defaults. \
Check the stored JSON in tb_observer."
);
ObserverRetryConfig::default()
},
};
Ok(ObserverDefinition {
event_type: observer.event_type.clone().unwrap_or_else(|| "INSERT".to_string()),
entity: observer.entity_type.clone().unwrap_or_else(|| "*".to_string()),
condition: observer.condition_expression.clone(),
actions,
retry: retry_config,
on_failure: FailurePolicy::default(),
})
}
pub async fn start(&mut self) -> Result<(), ServerError> {
if self.running.load(Ordering::SeqCst) {
return Err(ServerError::ConfigError("Observer runtime already running".to_string()));
}
if listener_selection(self.config.transport.transport) == ListenerSelection::TransportStream
{
return Box::pin(self.start_transport_stream()).await;
}
info!("Starting observer runtime (PostgreSQL LISTEN/NOTIFY transport)...");
let (observers, entity_type_index) = self.load_observers().await?;
self.observer_count.store(observers.len(), Ordering::SeqCst);
let observers_for_mount = observers.clone();
let matcher = EventMatcher::build(observers).map_err(|e| {
ServerError::ConfigError(format!("Failed to build event matcher: {}", e))
})?;
let matcher_for_logging = matcher.clone();
let built = ObserverExecutor::new_with_email(
matcher.clone(),
self.dlq.clone(),
self.config.email.as_ref(),
)
.map_err(|e| ServerError::ConfigError(format!("invalid observer email config: {e}")))?
.with_database_pool(self.config.pool.clone());
#[cfg(feature = "observers-cache")]
let built = {
self.cache_invalidator =
connect_cache_invalidator(&observers_for_mount, self.config.redis.as_ref()).await?;
match self.cache_invalidator.clone() {
Some(invalidator) => built.with_cache_invalidation(invalidator),
None => built,
}
};
#[cfg(not(feature = "observers-cache"))]
reject_cache_actions_not_compiled_in(&observers_for_mount)?;
let executor = Arc::new(built);
{
let mut m = self.matcher.write().await;
*m = Some(matcher_for_logging.clone());
}
{
let mut ex = self.executor.write().await;
*ex = Some(executor.clone());
}
self.entity_type_index.store(Arc::new(entity_type_index));
let checkpoint_store = PostgresCheckpointStore::new(self.config.pool.clone());
let restored = if let Ok(state) = checkpoint_store.load(&self.config.listener_id).await {
state
} else {
{
crate::migration_lock::run_migration(
&self.config.pool,
fraiseql_observers::checkpoint::migration_sql(),
)
.await
.map_err(|e| {
ServerError::ConfigError(format!(
"observer checkpoint table is missing and could not be created \
(apply fraiseql-observers migration 02_create_observer_checkpoints.sql): {e}"
))
})?;
checkpoint_store.load(&self.config.listener_id).await.map_err(|e| {
ServerError::ConfigError(format!(
"failed to load observer checkpoint for listener '{}': {e}",
self.config.listener_id
))
})?
}
};
if let Some(state) = &restored {
info!(
listener_id = %self.config.listener_id,
resume_from = state.last_processed_id,
"Restored observer change-log checkpoint; resuming (no replay)"
);
self.last_checkpoint.store(state.last_processed_id, Ordering::Relaxed);
}
let mut listener_config = ChangeLogListenerConfig::new(self.config.pool.clone())
.with_poll_interval(self.config.poll_interval_ms)
.with_batch_size(self.config.batch_size)
.with_listener_id(self.config.listener_id.clone());
if let Some(state) = &restored {
listener_config = listener_config.with_resume_from(state.last_processed_id);
}
let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<()>(1);
self.shutdown_tx = Some(shutdown_tx);
let running = self.running.clone();
let events_processed = self.events_processed.clone();
let errors = self.errors.clone();
let last_checkpoint = self.last_checkpoint.clone();
let poll_interval = Duration::from_millis(self.config.poll_interval_ms);
let pool = self.config.pool.clone();
let log_payloads = self.config.log_payloads;
let listener_id = self.config.listener_id.clone();
let mut checkpoint_event_count = restored.as_ref().map_or(0, |s| s.event_count);
let capture_dispatch = self.capture_dispatch.clone();
let matcher_ref = Arc::clone(&self.matcher);
let executor_ref = Arc::clone(&self.executor);
let entity_type_index_ref = Arc::clone(&self.entity_type_index);
let bridge_sender = self.event_bridge_sender.clone();
let initial_matcher = {
let m = self.matcher.read().await;
m.clone().ok_or_else(|| {
ServerError::ConfigError(
"matcher not initialised before spawning background task".to_string(),
)
})?
};
let initial_executor = {
let ex = self.executor.read().await;
ex.clone().ok_or_else(|| {
ServerError::ConfigError(
"executor not initialised before spawning background task".to_string(),
)
})?
};
debug!("About to spawn background task");
running.store(true, Ordering::SeqCst);
let (ready_tx, ready_rx) = oneshot::channel::<()>();
debug!("Calling tokio::spawn()");
let handle = tokio::spawn(async move {
let mut listener = ChangeLogListener::new(listener_config);
let mut current_matcher = initial_matcher;
let mut current_executor = initial_executor;
debug!("Observer runtime background task spawned");
debug!("Poll interval: {:?}", poll_interval);
info!("Observer runtime started, beginning event processing loop");
let _ = ready_tx.send(());
loop {
tokio::select! {
_ = shutdown_rx.recv() => {
info!("Observer runtime received shutdown signal");
break;
}
result = listener.next_batch() => {
{
let m = matcher_ref.read().await;
if let Some(updated) = m.clone() {
current_matcher = updated;
}
}
{
let ex = executor_ref.read().await;
if let Some(updated) = ex.clone() {
current_executor = updated;
}
}
match result {
Ok(entries) => {
if entries.is_empty() {
tokio::time::sleep(poll_interval).await;
continue;
}
debug!("Processing batch of {} change log entries", entries.len());
for entry in &entries {
let event = match entry.to_entity_event() {
Ok(e) => e,
Err(e) => {
errors.fetch_add(1, Ordering::Relaxed);
warn!("Failed to convert change log entry to event: {}", e);
continue;
}
};
process_entity_event(
&event,
¤t_matcher,
¤t_executor,
&entity_type_index_ref,
&pool,
bridge_sender.as_ref(),
&events_processed,
&errors,
log_payloads,
)
.await;
if let Some(ref dispatch) = capture_dispatch {
dispatch(&event);
}
}
if let Err(e) = listener.record_dispatched(&entries).await {
errors.fetch_add(1, Ordering::Relaxed);
error!(
"Failed to record dispatched change-log rows; this batch \
may be re-delivered: {e}"
);
}
if let Some(last_entry) = entries.last() {
last_checkpoint.store(last_entry.id, Ordering::Relaxed);
checkpoint_event_count += entries.len();
let state = ObserverCheckpointState {
listener_id: listener_id.clone(),
last_processed_id: last_entry.id,
last_processed_at: chrono::Utc::now(),
batch_size: entries.len(),
event_count: checkpoint_event_count,
};
match checkpoint_store.save(&listener_id, &state).await {
Ok(()) => {
debug!(
"Checkpoint saved: listener_id={}, last_id={}",
listener_id, last_entry.id
);
}
Err(e) => {
error!("Failed to save checkpoint: {}", e);
}
}
}
}
Err(e) => {
errors.fetch_add(1, Ordering::Relaxed);
error!("Failed to fetch entries from change log: {}", e);
tokio::time::sleep(Duration::from_secs(1)).await;
}
}
}
}
if !running.load(Ordering::SeqCst) {
break;
}
}
info!("Observer runtime stopped");
});
debug!("tokio::spawn() returned, storing task handle");
self.task_handle = Some(handle);
ready_rx.await.map_err(|_| {
ServerError::ConfigError(
"observer background task exited before signalling readiness".to_string(),
)
})?;
info!("Runtime started successfully");
Ok(())
}
async fn start_transport_stream(&mut self) -> Result<(), ServerError> {
let kind = self.config.transport.transport;
info!("Starting observer runtime ({kind:?} transport)...");
let transport: Arc<dyn EventTransport> = match kind {
TransportKind::InMemory => Arc::new(InMemoryTransport::new()),
#[cfg(feature = "observers-nats")]
TransportKind::Nats => {
let nats_config = nats_config_from(&self.config.transport.nats);
let connected = fraiseql_observers::transport::NatsTransport::new(nats_config)
.await
.map_err(|e| {
ServerError::ConfigError(format!(
"observer NATS transport failed to connect to {}: {e}",
self.config.transport.nats.url
))
})?;
Arc::new(connected)
},
#[cfg(not(feature = "observers-nats"))]
TransportKind::Nats => {
return Err(ServerError::ConfigError(
"observer transport = \"nats\" but this binary lacks the observers-nats \
feature"
.to_string(),
));
},
other => {
return Err(ServerError::ConfigError(format!(
"unsupported observer transport for the stream path: {other:?}"
)));
},
};
let (observers, entity_type_index) = self.load_observers().await?;
self.observer_count.store(observers.len(), Ordering::SeqCst);
let matcher = EventMatcher::build(observers).map_err(|e| {
ServerError::ConfigError(format!("Failed to build event matcher: {}", e))
})?;
let executor = Arc::new(
ObserverExecutor::new_with_email(
matcher.clone(),
self.dlq.clone(),
self.config.email.as_ref(),
)
.map_err(|e| ServerError::ConfigError(format!("invalid observer email config: {e}")))?
.with_database_pool(self.config.pool.clone()),
);
{
let mut m = self.matcher.write().await;
*m = Some(matcher.clone());
}
{
let mut ex = self.executor.write().await;
*ex = Some(executor.clone());
}
self.entity_type_index.store(Arc::new(entity_type_index));
let mut stream = transport.subscribe(EventFilter::all_tenants()).await.map_err(|e| {
ServerError::ConfigError(format!("failed to subscribe to observer transport: {e}"))
})?;
let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<()>(1);
self.shutdown_tx = Some(shutdown_tx);
let (ready_tx, ready_rx) = oneshot::channel::<()>();
let running = self.running.clone();
let events_processed = self.events_processed.clone();
let errors = self.errors.clone();
let pool = self.config.pool.clone();
let log_payloads = self.config.log_payloads;
let matcher_ref = Arc::clone(&self.matcher);
let executor_ref = Arc::clone(&self.executor);
let entity_type_index_ref = Arc::clone(&self.entity_type_index);
let bridge_sender = self.event_bridge_sender.clone();
let mut current_matcher = matcher;
let mut current_executor = executor;
running.store(true, Ordering::SeqCst);
let handle = tokio::spawn(async move {
info!("Observer runtime stream loop started, beginning event processing");
let _ = ready_tx.send(());
loop {
tokio::select! {
_ = shutdown_rx.recv() => {
info!("Observer runtime received shutdown signal");
break;
}
maybe_event = stream.next() => {
{
let m = matcher_ref.read().await;
if let Some(updated) = m.clone() {
current_matcher = updated;
}
}
{
let ex = executor_ref.read().await;
if let Some(updated) = ex.clone() {
current_executor = updated;
}
}
match maybe_event {
Some(Ok(event)) => {
process_entity_event(
&event,
¤t_matcher,
¤t_executor,
&entity_type_index_ref,
&pool,
bridge_sender.as_ref(),
&events_processed,
&errors,
log_payloads,
)
.await;
},
Some(Err(e)) => {
errors.fetch_add(1, Ordering::Relaxed);
error!("Observer transport stream error: {}", e);
},
None => {
info!("Observer transport stream ended; stopping loop");
break;
},
}
}
}
if !running.load(Ordering::SeqCst) {
break;
}
}
info!("Observer runtime stopped");
});
self.task_handle = Some(handle);
ready_rx.await.map_err(|_| {
ServerError::ConfigError(
"observer background task exited before signalling readiness".to_string(),
)
})?;
info!("Observer runtime started ({kind:?} transport)");
Ok(())
}
pub async fn stop(&mut self) -> Result<(), ServerError> {
if !self.running.load(Ordering::SeqCst) {
return Ok(());
}
info!("Stopping observer runtime...");
self.running.store(false, Ordering::SeqCst);
if let Some(tx) = self.shutdown_tx.take() {
let _ = tx.send(()).await;
}
if let Some(handle) = self.task_handle.take() {
let _ = tokio::time::timeout(Duration::from_secs(10), handle).await;
}
info!("Observer runtime stopped");
Ok(())
}
#[must_use]
pub fn is_running(&self) -> bool {
self.running.load(Ordering::SeqCst)
}
#[must_use]
pub(crate) fn transport_requires_broker(&self) -> bool {
self.config.transport.transport != TransportKind::Postgres
}
#[must_use]
pub fn stream_replay_reader(
&self,
) -> Option<Arc<fraiseql_observers::listener::ChangeLogReplayReader>> {
match listener_selection(self.config.transport.transport) {
ListenerSelection::PostgresChangeLog => {
Some(Arc::new(fraiseql_observers::listener::ChangeLogReplayReader::new(
self.config.pool.clone(),
self.config.listener_id.clone(),
)))
},
ListenerSelection::TransportStream => None,
}
}
#[must_use]
pub(crate) const fn dlq(&self) -> &Arc<InMemoryDlq> {
&self.dlq
}
#[must_use]
pub(crate) const fn executor_ref(&self) -> &Arc<RwLock<Option<Arc<ObserverExecutor>>>> {
&self.executor
}
#[must_use]
pub fn health(&self) -> RuntimeHealth {
RuntimeHealth {
running: self.running.load(Ordering::SeqCst),
observer_count: self.observer_count.load(Ordering::SeqCst),
last_checkpoint: Some(self.last_checkpoint.load(Ordering::SeqCst)),
events_processed: self.events_processed.load(Ordering::SeqCst),
errors: self.errors.load(Ordering::SeqCst),
}
}
pub async fn reload_observers(&self) -> Result<usize, ServerError> {
debug!("Reloading observers from database");
let (observers, new_entity_type_index) = self.load_observers().await?;
let count = observers.len();
let new_matcher = EventMatcher::build(observers)
.map_err(|e| ServerError::ConfigError(format!("Failed to build matcher: {}", e)))?;
let rebuilt = ObserverExecutor::new_with_email(
new_matcher.clone(),
self.dlq.clone(),
self.config.email.as_ref(),
)
.map_err(|e| ServerError::ConfigError(format!("invalid observer email config: {e}")))?
.with_database_pool(self.config.pool.clone());
#[cfg(feature = "observers-cache")]
let rebuilt = match self.cache_invalidator.clone() {
Some(invalidator) => rebuilt.with_cache_invalidation(invalidator),
None => rebuilt,
};
let new_executor = Arc::new(rebuilt);
debug!("Swapping matcher, executor, and entity_type_index atomically");
{
let mut m = self.matcher.write().await;
*m = Some(new_matcher);
}
{
let mut ex = self.executor.write().await;
*ex = Some(new_executor);
}
self.entity_type_index.store(Arc::new(new_entity_type_index));
self.observer_count.store(count, Ordering::SeqCst);
info!("Reloaded {} observers successfully", count);
Ok(count)
}
}
#[allow(clippy::too_many_arguments)]
async fn process_entity_event(
event: &ObserverEntityEvent,
matcher: &EventMatcher,
executor: &Arc<ObserverExecutor>,
entity_type_index: &ArcSwap<HashMap<(String, String), Vec<i64>>>,
pool: &PgPool,
bridge_sender: Option<&mpsc::Sender<BridgeEntityEvent>>,
events_processed: &std::sync::atomic::AtomicU64,
errors: &std::sync::atomic::AtomicU64,
log_payloads: bool,
) {
let matching_observers = matcher.find_matches(event);
let process_result = executor.process_event(event).await;
match process_result {
Ok(summary) => {
events_processed.fetch_add(1, Ordering::Relaxed);
debug!(
"Event {} processed: {} actions succeeded, {} skipped",
event.id, summary.successful_actions, summary.conditions_skipped
);
let event_type_str = event.event_type.as_str().to_uppercase();
let observer_ids = entity_type_index
.load()
.get(&(event.entity_type.clone(), event_type_str.clone()))
.cloned();
if let Some(observer_ids) = observer_ids {
let status = if summary.successful_actions > 0 {
OBSERVER_LOG_STATUS_SUCCESS
} else {
OBSERVER_LOG_STATUS_FAILED
};
let representative = if status == "success" {
summary.action_details.iter().find(|d| d.success)
} else {
summary.action_details.iter().find(|d| !d.success)
}
.or_else(|| summary.action_details.first());
let fallback_duration_ms = if matching_observers.is_empty() {
None
} else {
#[allow(clippy::cast_precision_loss, clippy::cast_possible_truncation)]
Some((summary.total_duration_ms / matching_observers.len() as f64) as i32)
};
for observer_id in observer_ids {
write_observer_log(
pool,
observer_id,
event,
status,
representative,
fallback_duration_ms,
None,
log_payloads,
)
.await;
}
}
if let (Some(sender), Some(bridge_event)) = (bridge_sender, bridge_event_for(event)) {
forward_to_bridge(sender, bridge_event, &event.id.to_string()).await;
}
},
Err(e) => {
errors.fetch_add(1, Ordering::Relaxed);
error!("Failed to process event {}: {}", event.id, e);
let event_type_str = event.event_type.as_str().to_uppercase();
let observer_ids_err = entity_type_index
.load()
.get(&(event.entity_type.clone(), event_type_str))
.cloned();
if let Some(observer_ids) = observer_ids_err {
let error_message = e.to_string();
for observer_id in observer_ids {
write_observer_log(
pool,
observer_id,
event,
OBSERVER_LOG_STATUS_FAILED,
None,
None,
Some(&error_message),
log_payloads,
)
.await;
}
}
},
}
}
pub(crate) const fn subscription_operation_for(
kind: fraiseql_observers::EventKind,
) -> Option<SubscriptionOperation> {
use fraiseql_observers::EventKind;
match kind {
EventKind::Created => Some(SubscriptionOperation::Create),
EventKind::Updated => Some(SubscriptionOperation::Update),
EventKind::Deleted => Some(SubscriptionOperation::Delete),
EventKind::Custom => None,
}
}
pub(crate) fn bridge_event_for(event: &ObserverEntityEvent) -> Option<BridgeEntityEvent> {
let operation = subscription_operation_for(event.event_type)?;
let mut bridge_event = BridgeEntityEvent::new(
&event.entity_type,
event.entity_id.to_string(),
operation,
event.data.clone(),
);
if let Some(ref tid) = event.tenant_id {
bridge_event = bridge_event.with_tenant_id(tid);
}
let envelope = ChangeSpineEnvelope {
actor_type: event.actor_type.clone(),
acting_for: event.acting_for.clone(),
schema_version: event.schema_version.clone(),
tenant_id: event.tenant_id.clone(),
duration_ms: event.duration_ms,
seq: event.seq,
};
if !envelope.is_empty() {
bridge_event = bridge_event.with_change_spine(envelope);
}
Some(bridge_event)
}
async fn forward_to_bridge(
sender: &mpsc::Sender<BridgeEntityEvent>,
bridge_event: BridgeEntityEvent,
event_id: &str,
) {
if let Err(e) = sender.send(bridge_event).await {
error!("EventBridge task is gone; event {event_id} not forwarded to subscriptions: {e}");
}
}
const MAX_LOG_PAYLOAD_BYTES: usize = 64 * 1024;
fn truncate_log_payload(data: &serde_json::Value) -> serde_json::Value {
let size = serde_json::to_vec(data).map_or(0, |v| v.len());
if size > MAX_LOG_PAYLOAD_BYTES {
serde_json::json!({
"_truncated": true,
"_original_size_bytes": size,
})
} else {
data.clone()
}
}
#[cfg_attr(not(feature = "observers-cache"), allow(dead_code))]
fn declares_cache_action(observers: &HashMap<String, ObserverDefinition>) -> bool {
observers
.values()
.flat_map(|o| o.actions.iter())
.any(|a| matches!(a, ObserverActionConfig::Cache { .. }))
}
#[cfg(feature = "observers-cache")]
async fn connect_cache_invalidator(
observers: &HashMap<String, ObserverDefinition>,
redis: Option<&fraiseql_observers::config::RedisConfig>,
) -> crate::Result<Option<Arc<fraiseql_observers::RedisCacheInvalidator>>> {
if !declares_cache_action(observers) {
return Ok(None);
}
let Some(redis_config) = redis else {
return Err(ServerError::ConfigError(
"an observer declares a `cache` action but no Redis backend is configured; \
add [observers.runtime.redis] (url = \"redis://…\") or remove the action"
.to_string(),
));
};
let invalidator = fraiseql_observers::RedisCacheInvalidator::connect(redis_config)
.await
.map_err(|e| {
ServerError::ConfigError(format!(
"an observer declares a `cache` action but the configured Redis backend \
({}) could not be reached: {e}",
redis_config.url
))
})?;
info!(
redis_url = %redis_config.url,
"Mounted the Redis cache-invalidation transport for `cache` observer actions"
);
Ok(Some(Arc::new(invalidator)))
}
#[cfg(not(feature = "observers-cache"))]
fn reject_cache_actions_not_compiled_in(
observers: &HashMap<String, ObserverDefinition>,
) -> crate::Result<()> {
if declares_cache_action(observers) {
return Err(ServerError::ConfigError(
"an observer declares a `cache` action but this binary was built without the `observers-cache` feature, so the Redis cache-invalidation transport is not compiled in; rebuild with --features observers-cache or remove the action"
.to_string(),
));
}
Ok(())
}
const OBSERVER_LOG_STATUS_SUCCESS: &str = "success";
const OBSERVER_LOG_STATUS_FAILED: &str = "failed";
#[allow(clippy::too_many_arguments)]
async fn write_observer_log(
pool: &PgPool,
observer_id: i64,
event: &ObserverEntityEvent,
status: &str,
detail: Option<&ActionExecutionDetail>,
fallback_duration_ms: Option<i32>,
fallback_error: Option<&str>,
log_payloads: bool,
) {
let action_index = detail.map(|d| i32::try_from(d.action_index).unwrap_or(i32::MAX));
let action_type = detail.map(|d| d.action_type.clone());
let status_code = detail.and_then(|d| d.status_code).map(i32::from);
let duration_ms = detail
.map(|d| {
#[allow(clippy::cast_possible_truncation)]
{
d.duration_ms as i32
}
})
.or(fallback_duration_ms);
let response_payload =
detail.map(|d| serde_json::json!({ "message": d.message, "success": d.success }));
let request_payload = log_payloads.then(|| truncate_log_payload(&event.data));
let error_message = fallback_error
.map(ToString::to_string)
.or_else(|| detail.and_then(|d| d.error_message.clone()));
if let Err(e) = sqlx::query(
"INSERT INTO tb_observer_log
(fk_observer, event_id, entity_type, entity_id, event_type, status,
action_index, action_type, response_status_code, response_payload,
request_payload, duration_ms, error_message, attempt_number, max_attempts)
VALUES ($1, $2, $3, $4::uuid, $5, $6, $7, $8, $9, $10, $11, $12, $13, 1, 3)",
)
.bind(observer_id)
.bind(event.id)
.bind(&event.entity_type)
.bind(event.entity_id.to_string())
.bind(event.event_type.as_str())
.bind(status)
.bind(action_index)
.bind(action_type)
.bind(status_code)
.bind(response_payload)
.bind(request_payload)
.bind(duration_ms)
.bind(error_message)
.execute(pool)
.await
{
warn!(
"Failed to write observer {status} log (observer {observer_id}, event {}): {e}",
event.id
);
}
}
pub(crate) struct InMemoryDlq {
items: std::sync::Mutex<Vec<fraiseql_observers::DlqItem>>,
function_items: std::sync::Mutex<Vec<fraiseql_observers::FunctionDispatchRecord>>,
max_size: Option<usize>,
overflow_count: std::sync::atomic::AtomicUsize,
}
impl InMemoryDlq {
pub(crate) const fn new_with_max(max_size: Option<usize>) -> Self {
Self {
items: std::sync::Mutex::new(Vec::new()),
function_items: std::sync::Mutex::new(Vec::new()),
max_size,
overflow_count: std::sync::atomic::AtomicUsize::new(0),
}
}
pub(crate) fn function_count(&self) -> usize {
self.function_items.lock().expect("function_items mutex poisoned").len()
}
pub(crate) fn count(&self) -> usize {
self.items.lock().expect("items mutex poisoned").len()
}
pub(crate) fn overflow_count(&self) -> usize {
self.overflow_count.load(Ordering::Relaxed)
}
pub(crate) fn list_all(&self) -> Vec<fraiseql_observers::DlqItem> {
self.items.lock().expect("items mutex poisoned").clone()
}
pub(crate) fn get(&self, id: uuid::Uuid) -> Option<fraiseql_observers::DlqItem> {
self.items
.lock()
.expect("items mutex poisoned")
.iter()
.find(|item| item.id == id)
.cloned()
}
pub(crate) fn try_claim(&self, id: uuid::Uuid) -> Option<fraiseql_observers::DlqItem> {
let mut items = self.items.lock().expect("items mutex poisoned");
let pos = items.iter().position(|item| item.id == id)?;
Some(items.remove(pos))
}
pub(crate) fn reinsert(&self, item: fraiseql_observers::DlqItem) {
self.items.lock().expect("items mutex poisoned").push(item);
}
}
#[async_trait::async_trait]
impl fraiseql_observers::DeadLetterQueue for InMemoryDlq {
async fn push(
&self,
event: fraiseql_observers::EntityEvent,
action: fraiseql_observers::ActionConfig,
error: String,
) -> fraiseql_observers::Result<uuid::Uuid> {
let id = uuid::Uuid::new_v4();
let mut items = self.items.lock().expect("items mutex poisoned");
if let Some(max) = self.max_size {
if items.len() >= max {
warn!(
max_dlq_size = max,
action_type = action.action_type(),
event_id = %event.id,
"DLQ full; dropping failed action entry"
);
self.overflow_count.fetch_add(1, Ordering::Relaxed);
return Ok(id);
}
}
items.push(fraiseql_observers::DlqItem {
id,
event,
action,
error_message: error,
attempts: 0,
});
Ok(id)
}
async fn get_pending(
&self,
limit: i64,
) -> fraiseql_observers::Result<Vec<fraiseql_observers::DlqItem>> {
let items = self.items.lock().expect("items mutex poisoned");
#[allow(clippy::cast_sign_loss, clippy::cast_possible_truncation)]
let limit_usize = limit as usize;
Ok(items.iter().take(limit_usize).cloned().collect())
}
async fn mark_success(&self, id: uuid::Uuid) -> fraiseql_observers::Result<()> {
let mut items = self.items.lock().expect("items mutex poisoned");
items.retain(|i| i.id != id);
Ok(())
}
async fn mark_retry_failed(
&self,
id: uuid::Uuid,
error: &str,
) -> fraiseql_observers::Result<()> {
let mut items = self.items.lock().expect("items mutex poisoned");
if let Some(item) = items.iter_mut().find(|i| i.id == id) {
item.attempts = item.attempts.saturating_add(1);
item.error_message = error.to_string();
}
Ok(())
}
async fn push_function(
&self,
record: fraiseql_observers::FunctionDispatchRecord,
) -> fraiseql_observers::Result<uuid::Uuid> {
let id = record.id;
let mut items = self.function_items.lock().expect("function_items mutex poisoned");
if let Some(max) = self.max_size {
if items.len() >= max {
warn!(
max_dlq_size = max,
function = %record.function_name,
trigger = %record.trigger_type,
"DLQ full; dropping failed function dispatch entry"
);
self.overflow_count.fetch_add(1, Ordering::Relaxed);
crate::function_metrics::record_dlq_eviction();
crate::function_metrics::set_dlq_size(items.len());
return Ok(id);
}
}
items.push(record);
crate::function_metrics::set_dlq_size(items.len());
Ok(id)
}
async fn get_pending_functions(
&self,
limit: i64,
) -> fraiseql_observers::Result<Vec<fraiseql_observers::FunctionDispatchRecord>> {
let items = self.function_items.lock().expect("function_items mutex poisoned");
#[allow(clippy::cast_sign_loss, clippy::cast_possible_truncation)]
let limit_usize = limit as usize;
Ok(items.iter().take(limit_usize).cloned().collect())
}
}