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::{CronExecutionState, CronSchedule},
};
use fraiseql_observers::{
DispatchSource, LeaseGuardedRunner, PostgresSourceCursorStore, RunOutcome,
derive_idempotency_token,
};
use tracing::{debug, info, warn};
use super::metrics;
pub struct SourcePoller {
source_name: String,
schedule: CronSchedule,
module: FunctionModule,
observer: Arc<FunctionObserver>,
cursor_store: PostgresSourceCursorStore,
executor: Arc<dyn QueryExecutor>,
runner: LeaseGuardedRunner,
host_config: HostContextConfig,
limits: ResourceLimits,
idempotency_key: Option<Arc<[u8]>>,
log_payloads: bool,
state: CronExecutionState,
}
impl SourcePoller {
#[allow(clippy::too_many_arguments)]
#[must_use]
pub fn new(
source_name: impl Into<String>,
schedule: CronSchedule,
module: FunctionModule,
observer: Arc<FunctionObserver>,
cursor_store: PostgresSourceCursorStore,
executor: Arc<dyn QueryExecutor>,
runner: LeaseGuardedRunner,
host_config: HostContextConfig,
limits: ResourceLimits,
idempotency_key: Option<Arc<[u8]>>,
log_payloads: bool,
) -> Self {
Self {
source_name: source_name.into(),
schedule,
module,
observer,
cursor_store,
executor,
runner,
host_config,
limits,
idempotency_key,
log_payloads,
state: CronExecutionState::new(),
}
}
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_source_cursor(self.source_name.clone(), self.cursor_store.clone())
.with_executor(Arc::clone(&self.executor))
.with_idempotency_token(idempotency_token.to_string()),
)
}
async fn fire_once(
&self,
now: DateTime<Utc>,
) -> RunOutcome<fraiseql_error::Result<FunctionResult>> {
let payload = build_source_payload(&self.source_name, &self.schedule.expression, now);
let token = self.idempotency_token(&payload);
if self.log_payloads {
debug!(
source = %self.source_name,
idempotency_token = %token,
payload = %payload.data,
"source firing (payload logging enabled)"
);
}
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 label = match &result {
Ok(_) => metrics::RESULT_OK,
Err(_) => metrics::RESULT_ERROR,
};
metrics::record_fire(&self.source_name, label, elapsed);
match &result {
Ok(_) => info!(
source = %self.source_name,
idempotency_token = %token,
duration_ms = elapsed * 1000.0,
"source fired"
),
Err(error) => warn!(
source = %self.source_name,
idempotency_token = %token,
duration_ms = elapsed * 1000.0,
%error,
"source invocation failed — re-runs from the last cursor next tick"
),
}
RunOutcome::Ran(result)
},
Ok(RunOutcome::SkippedNotLeader) => {
metrics::record_skip_not_leader(&self.source_name);
debug!(
source = %self.source_name,
idempotency_token = %token,
"source skipped — another replica leads"
);
RunOutcome::SkippedNotLeader
},
Err(error) => {
warn!(
source = %self.source_name,
idempotency_token = %token,
%error,
"source lease acquire failed — skipping tick"
);
RunOutcome::SkippedNotLeader
},
}
}
pub async fn run_forever(mut self) {
let mut ticker = tokio::time::interval(Duration::from_secs(60));
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
ticker.tick().await;
info!(
source = %self.source_name,
schedule = %self.schedule.expression,
"source scheduler started"
);
loop {
ticker.tick().await;
let now = Utc::now();
if !self.state.should_execute(&self.schedule, &now) {
continue;
}
self.state.record_execution(now);
let _ = self.fire_once(now).await;
}
}
}
fn build_source_payload(source_name: &str, schedule: &str, now: DateTime<Utc>) -> EventPayload {
EventPayload {
trigger_type: format!("source:{source_name}"),
entity: "source".to_string(),
event_kind: "scheduled".to_string(),
data: serde_json::json!({
"source": source_name,
"schedule": schedule,
"scheduled_at": now.to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
}),
timestamp: now,
}
}
#[cfg(test)]
mod tests;