use std::{
sync::Arc,
time::{Duration, Instant},
};
use chrono::{DateTime, Utc};
use fraiseql_functions::{
EventPayload, FunctionModule, FunctionObserver, FunctionResult, ResourceLimits,
host::{
dyn_context::DynHostContext,
live::{HostContextConfig, LiveHostContext, QueryExecutor},
},
triggers::{CronDecision, CronExecutionState, CronSchedule, CronTrigger},
};
use fraiseql_observers::{
DispatchSource, LeaseGuardedRunner, RunOutcome, derive_idempotency_token,
};
use tracing::{debug, info, warn};
use crate::{query_bridge::RunAsQueryExecutor, subsystems::BeforeMutationHooks};
#[derive(Clone)]
pub struct PgCronState {
pool: sqlx::PgPool,
}
impl PgCronState {
#[must_use]
pub const fn new(pool: sqlx::PgPool) -> Self {
Self { pool }
}
pub async fn init(&self) -> Result<(), sqlx::Error> {
crate::migration_lock::run_migration(
&self.pool,
fraiseql_functions::migrations::cron_migration_sql(),
)
.await
}
pub async fn load_last_fired(
&self,
function_name: &str,
cron_expr: &str,
) -> Result<Option<DateTime<Utc>>, sqlx::Error> {
sqlx::query_scalar(
"SELECT last_fired_at FROM _fraiseql_cron_state \
WHERE function_name = $1 AND cron_expr = $2",
)
.bind(function_name)
.bind(cron_expr)
.fetch_optional(&self.pool)
.await
}
pub async fn resume_state(
&self,
function_name: &str,
cron_expr: &str,
) -> Result<CronExecutionState, sqlx::Error> {
Ok(self
.load_last_fired(function_name, cron_expr)
.await?
.map_or_else(CronExecutionState::new, CronExecutionState::resuming_from))
}
pub async fn record_fire(
&self,
function_name: &str,
cron_expr: &str,
fired_at: DateTime<Utc>,
) -> Result<(), sqlx::Error> {
sqlx::query(
"INSERT INTO _fraiseql_cron_state \
(function_name, cron_expr, last_fired_at, fire_count) \
VALUES ($1, $2, $3, 1) \
ON CONFLICT (function_name, cron_expr) DO UPDATE SET \
last_fired_at = EXCLUDED.last_fired_at, \
fire_count = _fraiseql_cron_state.fire_count + 1, \
updated_at = now()",
)
.bind(function_name)
.bind(cron_expr)
.bind(fired_at)
.execute(&self.pool)
.await
.map(|_| ())
}
}
pub struct CronPoller {
function_name: String,
schedule: CronSchedule,
cron_expr: String,
module: FunctionModule,
observer: Arc<FunctionObserver>,
executor: Arc<dyn QueryExecutor>,
identity: fraiseql_core::security::SecurityContext,
runner: LeaseGuardedRunner,
cron_state: PgCronState,
host_config: HostContextConfig,
limits: ResourceLimits,
idempotency_key: Option<Arc<[u8]>>,
state: CronExecutionState,
}
impl CronPoller {
#[allow(clippy::too_many_arguments)]
#[must_use]
pub fn new(
trigger: &CronTrigger,
schedule: CronSchedule,
module: FunctionModule,
observer: Arc<FunctionObserver>,
executor: Arc<dyn QueryExecutor>,
identity: fraiseql_core::security::SecurityContext,
runner: LeaseGuardedRunner,
cron_state: PgCronState,
host_config: HostContextConfig,
limits: ResourceLimits,
idempotency_key: Option<Arc<[u8]>>,
initial_state: CronExecutionState,
) -> Self {
Self {
function_name: trigger.function_name.clone(),
schedule,
cron_expr: trigger.schedule.clone(),
module,
observer,
executor,
identity,
runner,
cron_state,
host_config,
limits,
idempotency_key,
state: initial_state,
}
}
fn idempotency_token(&self, payload: &EventPayload) -> String {
derive_idempotency_token(
self.idempotency_key.as_deref(),
DispatchSource::Source,
&self.module.name,
&payload.trigger_type,
&payload.data,
)
}
fn build_host(
&self,
payload: EventPayload,
idempotency_token: &str,
) -> Arc<dyn DynHostContext> {
Arc::new(
LiveHostContext::new(payload, self.host_config.clone())
.with_executor(Arc::clone(&self.executor))
.with_idempotency_token(idempotency_token.to_string())
.with_security_context(self.identity.clone()),
)
}
async fn fire_once(
&self,
now: DateTime<Utc>,
) -> RunOutcome<fraiseql_error::Result<FunctionResult>> {
let trigger = CronTrigger {
function_name: self.function_name.clone(),
schedule: self.cron_expr.clone(),
timezone: "UTC".to_string(),
};
let payload = trigger.build_payload(&now);
let token = self.idempotency_token(&payload);
let started = Instant::now();
let attempt = self
.runner
.run(|| async {
let host = self.build_host(payload.clone(), &token);
self.observer
.invoke_with_context(&self.module, payload.clone(), host, self.limits.clone())
.await
})
.await;
let elapsed = started.elapsed().as_secs_f64();
match attempt {
Ok(RunOutcome::Ran(result)) => {
let metric_result = match &result {
Ok(_) => crate::function_metrics::RESULT_OK,
Err(_) => crate::function_metrics::RESULT_ERROR,
};
crate::function_metrics::record_dispatch(
&self.function_name,
crate::function_metrics::KIND_CRON,
metric_result,
elapsed,
);
match &result {
Ok(_) => {
info!(
function = %self.function_name,
idempotency_token = %token,
duration_ms = elapsed * 1000.0,
"cron function fired"
);
if let Err(error) = self
.cron_state
.record_fire(&self.function_name, &self.cron_expr, now)
.await
{
warn!(
function = %self.function_name,
%error,
"cron state record failed — the function ran but the fire was not persisted"
);
}
},
Err(error) => warn!(
function = %self.function_name,
idempotency_token = %token,
%error,
"cron function invocation failed"
),
}
RunOutcome::Ran(result)
},
Ok(RunOutcome::SkippedNotLeader) => {
debug!(
function = %self.function_name,
"cron function skipped — another replica leads"
);
RunOutcome::SkippedNotLeader
},
Err(error) => {
warn!(
function = %self.function_name,
%error,
"cron lease acquire failed — skipping tick"
);
RunOutcome::SkippedNotLeader
},
}
}
pub async fn run_forever(mut self) {
let mut ticker = tokio::time::interval(Duration::from_mins(1));
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
ticker.tick().await;
info!(
function = %self.function_name,
schedule = %self.cron_expr,
"cron scheduler started"
);
loop {
ticker.tick().await;
let now = Utc::now();
match self.state.decide(&self.schedule, &now) {
CronDecision::Fire => {},
CronDecision::NotDue => continue,
CronDecision::AlreadyFired {
window_start,
last_executed,
} => {
warn!(
function = %self.function_name,
schedule = %self.cron_expr,
%window_start,
%last_executed,
"cron tick skipped — this window already fired"
);
continue;
},
}
self.state.record_execution(now);
let _ = self.fire_once(now).await;
}
}
}
pub async fn build_cron_pollers(
db_pool: &sqlx::PgPool,
executor: &Arc<arc_swap::ArcSwap<fraiseql_core::runtime::Executor>>,
hooks: &BeforeMutationHooks,
host_config: &HostContextConfig,
limits: &ResourceLimits,
) -> Result<Vec<CronPoller>, sqlx::Error> {
let cron_state = PgCronState::new(db_pool.clone());
let mut pollers = Vec::new();
for trigger in &hooks.trigger_registry.cron_triggers {
let schedule = match CronSchedule::parse(&trigger.schedule) {
Ok(schedule) => schedule,
Err(error) => {
warn!(
function = %trigger.function_name,
expression = %trigger.schedule,
%error,
"invalid cron schedule — skipping function"
);
continue;
},
};
let Some(module) = hooks.module_registry.get(&trigger.function_name).cloned() else {
continue;
};
let run_as = hooks.run_as.get(&trigger.function_name).cloned().unwrap_or_default();
let identity = run_as.identity(&trigger.function_name, &trigger.function_name);
let query_executor: Arc<dyn QueryExecutor> =
Arc::new(RunAsQueryExecutor::new(Arc::clone(executor), identity.clone()));
let initial_state =
cron_state.resume_state(&trigger.function_name, &trigger.schedule).await?;
pollers.push(CronPoller::new(
trigger,
schedule,
module,
Arc::clone(&hooks.observer),
query_executor,
identity,
LeaseGuardedRunner::postgres(
db_pool.clone(),
format!("cron:{}", trigger.function_name),
),
cron_state.clone(),
host_config.clone(),
limits.clone(),
hooks.idempotency_key.clone(),
initial_state,
));
}
Ok(pollers)
}
#[cfg(test)]
mod tests;