use std::collections::{BTreeMap, HashSet};
use std::future::Future;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use taquba::{Clock, Queue};
use taquba_cron::{Backfill, BackfillStart, CronScheduler, Schedule, ScheduleHandle};
use tokio_util::sync::CancellationToken;
use crate::graph::Graph;
use crate::pools::Pools;
use crate::records::JsonBytes;
use crate::records::{self, GraphRecord, RequestOutcome, RequestRecord};
use crate::request::{Request, RequestId, RequestStore};
use crate::scheduler::{Error, Scheduler, SchedulerOptions, TRIGGERS_QUEUE, firing_headers};
#[derive(Debug, Clone)]
pub struct DaemonOptions {
pub scheduler: SchedulerOptions,
pub sync_interval: Duration,
pub retention: Duration,
}
impl Default for DaemonOptions {
fn default() -> Self {
DaemonOptions {
scheduler: SchedulerOptions::default(),
sync_interval: Duration::from_secs(30),
retention: Duration::from_secs(90 * 86_400),
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct SyncReport {
pub adopted: Vec<(String, String)>,
pub refused: Vec<(String, String)>,
pub schedules: Vec<Schedule>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct RequestReport {
pub applied: Vec<(RequestId, RequestOutcome)>,
}
pub struct Daemon {
queue: Arc<Queue>,
scheduler: Arc<Scheduler>,
pools: Arc<Pools>,
requests: RequestStore,
clock: Arc<dyn Clock>,
logged: Mutex<HashSet<String>>,
}
impl Daemon {
pub fn new(
queue: Arc<Queue>,
scheduler: Arc<Scheduler>,
pools: Arc<Pools>,
requests: RequestStore,
) -> Self {
let clock = queue.clock();
Daemon {
queue,
scheduler,
pools,
requests,
clock,
logged: Mutex::new(HashSet::new()),
}
}
pub async fn sync(&self) -> Result<SyncReport, Error> {
let definitions = self.scheduler.definitions();
let current = definitions.current().await?;
let mut report = SyncReport::default();
let mut graphs: BTreeMap<String, Arc<Graph>> = BTreeMap::new();
let mut pending = Vec::new();
for (name, hash) in ¤t {
let record = self.scheduler.graph_record(name).await?;
if let Some(record) = &record
&& let Some(graph) = definitions.get(&record.definition).await?
{
graphs.insert(name.clone(), graph);
}
if record.as_ref().map(|r| &r.definition) != Some(hash) {
pending.push((name, hash));
}
}
for (name, hash) in pending {
let graph = match definitions.get(hash).await {
Ok(Some(graph)) => graph,
Ok(None) => {
self.refuse(
&mut report,
name,
hash,
"the definition object is absent".into(),
);
continue;
}
Err(e) => {
self.refuse(&mut report, name, hash, e.to_string());
continue;
}
};
if let Err(reason) = self.check(name, &graph, &graphs) {
self.refuse(&mut report, name, hash, reason);
continue;
}
self.adopt(hash, name).await?;
tracing::info!(graph = %name, definition = %hash, "definition adopted");
report.adopted.push((name.clone(), hash.clone()));
graphs.insert(name.clone(), graph);
}
report.schedules = graphs.values().filter_map(|g| schedule_of(g)).collect();
Ok(report)
}
pub async fn apply_requests(&self) -> Result<RequestReport, Error> {
let mut report = RequestReport::default();
for (id, bytes) in self.requests.list().await? {
let key = records::request_key(&id);
if self.queue.view().kv_get(&key).await?.is_none() {
match Request::from_bytes(&bytes) {
Ok(request) => match self.scheduler.handle_request(&id, &request).await {
Ok(record) => {
log_outcome(&id, &record);
report.applied.push((id.clone(), record.outcome));
}
Err(e) => {
tracing::warn!(request = %id, error = %e, "request failed");
continue;
}
},
Err(e) => {
tracing::warn!(request = %id, error = %e, "the object is not a request");
}
}
}
self.requests.remove(&id).await?;
}
Ok(report)
}
pub async fn run<F: Future<Output = ()>>(
&self,
options: DaemonOptions,
shutdown: F,
) -> Result<(), Error> {
let stop = CancellationToken::new();
let pool_handles = self.pools.spawn(&stop);
let scheduler_handle = self
.scheduler
.clone()
.spawn(options.scheduler.clone(), stop.clone().cancelled_owned());
let cron = CronScheduler::new(self.queue.clone());
let handle = cron.handle();
let cron = cron.spawn(std::future::pending());
let mut shutdown = std::pin::pin!(shutdown);
loop {
match self.sync().await {
Ok(report) => self.register(&handle, report.schedules),
Err(e) => tracing::warn!(error = %e, "sync pass failed"),
}
match self.apply_requests().await {
Ok(_) => {
if let Err(e) = self.scheduler.expire(options.retention).await {
tracing::warn!(error = %e, "retention pass failed");
}
}
Err(e) => tracing::warn!(error = %e, "request pass failed"),
}
tokio::select! {
() = tokio::time::sleep(options.sync_interval) => {}
() = &mut shutdown => break,
}
}
let _ = cron.shutdown().await;
stop.cancel();
for handle in pool_handles {
let _ = handle.wait().await;
}
scheduler_handle.wait().await
}
fn register(&self, handle: &ScheduleHandle, mut schedules: Vec<Schedule>) {
loop {
match handle.replace_all(schedules.clone()) {
Ok(()) => return,
Err(taquba_cron::Error::UnboundedStart(graph)) => {
self.log_once(&format!("cron/{graph}"), || {
tracing::error!(%graph, "the catch-up window is unbounded, and the schedule does not fire");
});
schedules.retain(|schedule| schedule.name != graph);
}
Err(e) => {
tracing::warn!(error = %e, "cron schedules refused");
return;
}
}
}
}
fn check(
&self,
name: &str,
graph: &Graph,
adopted: &BTreeMap<String, Arc<Graph>>,
) -> Result<(), String> {
if graph.name() != name {
return Err(format!("the definition is of graph `{}`", graph.name()));
}
for node in graph.nodes() {
if self.pools.runtime(node.pool()).is_none() {
return Err(format!(
"node `{}`: pool `{}` does not have a runtime",
node.name(),
node.pool()
));
}
}
for (other, other_graph) in adopted {
if other == name {
continue;
}
if let Some(node) = graph.conflicting_asset(other_graph) {
return Err(format!(
"node `{}`: asset `{}` is produced by graph `{other}`",
node.name(),
node.asset().unwrap_or_default()
));
}
}
Ok(())
}
fn refuse(&self, report: &mut SyncReport, name: &str, hash: &str, reason: String) {
self.log_once(hash, || {
tracing::error!(graph = %name, definition = %hash, %reason, "definition refused");
});
report.refused.push((name.to_string(), reason));
}
fn log_once(&self, key: &str, log: impl FnOnce()) {
let first = self
.logged
.lock()
.expect("the log set is not poisoned")
.insert(key.to_string());
if first {
log();
}
}
async fn adopt(&self, hash: &str, name: &str) -> Result<(), Error> {
let record = GraphRecord {
definition: hash.to_string(),
adopted_at_ms: self.clock.now_ms(),
};
self.queue
.kv_put(&records::graph_key(name), &record.to_bytes())
.await?;
Ok(())
}
}
fn log_outcome(id: &RequestId, record: &RequestRecord) {
match &record.outcome {
RequestOutcome::Refused { reason } => {
tracing::warn!(request = %id, %reason, "request refused");
}
outcome => {
tracing::info!(request = %id, ?outcome, "request applied");
}
}
}
fn schedule_of(graph: &Graph) -> Option<Schedule> {
let backfill = graph.catchup().map(|lookback| Backfill {
lookback,
start: BackfillStart::Lookback,
});
Some(
Schedule::new(
graph.name(),
graph.schedule()?.clone(),
TRIGGERS_QUEUE,
Vec::new(),
)
.headers(firing_headers(graph.name()))
.backfill(backfill),
)
}