Skip to main content

swale/
scheduler.rs

1//! The scheduler: it starts graph runs, submits a node when its upstreams
2//! satisfy its trigger rule, and settles the state of a graph run.
3//!
4//! Every task instance runs on the runtime of its pool ([`Pools`]). The
5//! terminal hook of each pool writes the node's record and enqueues an
6//! [`Event`], and the [`Scheduler`] is the [`Worker`] of the events queue.
7//! Every submit is idempotent on the deterministic run id, so a redelivered
8//! event and a repeated reconciler pass do not submit a second task
9//! instance.
10//!
11//! A job on the triggers queue ([`TRIGGERS_QUEUE`]) is a cron firing without
12//! a payload: the `swale.graph` header identifies the graph, and the
13//! `cron.previous_fire_ms` header, the start of the schedule interval that
14//! the firing ends, determines the partition. The [`TriggerWorker`] reads
15//! both, and [`Scheduler::handle_trigger`] starts the graph run of the
16//! graph's adopted definition for that partition. A start request starts the
17//! run of every partition it lists ([`Scheduler::start_runs`]).
18//!
19//! The events worker submits the ready downstreams of a terminated task
20//! instance, and the reconciler ([`Scheduler::reconcile`]) submits every
21//! ready node of every active graph run. Each reads the node records of the
22//! run once, applies [`crate::readiness`] and writes the final state when
23//! the run is finished. The reconciler also cancels the active runs of a
24//! cancelled graph run. A lost event delays a graph run by one reconciler
25//! interval at most.
26//!
27//! A request from another process ([`Scheduler::handle_request`], in the
28//! [`crate::request`] module) starts graph runs, reruns a node or cancels a
29//! graph run. A rerun of a succeeded node writes the expected rerun count of
30//! the node and its downstreams to the graph run record before the submit.
31//! A crash between the two writes leaves a run that the reconciler
32//! completes.
33
34use std::collections::{BTreeMap, HashMap};
35use std::future::Future;
36use std::sync::Arc;
37use std::time::Duration;
38
39use taquba::{
40    Clock, ExpiryIndex, JobRecord, LeaseHandle, PermanentFailure, Queue, SettlementEffects, Worker,
41    WorkerError, WorkerHandle,
42};
43use taquba_cron::PREVIOUS_FIRE_MS_HEADER;
44use taquba_workflow::{RunId, RunOptions, RunSpec, RunState};
45use tokio_util::sync::CancellationToken;
46
47use crate::definition_store::{DefinitionError, DefinitionStore};
48use crate::graph::{Graph, Node};
49use crate::hook::{EVENTS_QUEUE, Event};
50use crate::input::TaskInput;
51use crate::partition::Partition;
52use crate::pools::Pools;
53use crate::readiness::{
54    NodeState, current_records, is_ready, node_states, rerun_scope, settled_state,
55};
56use crate::records::JsonBytes;
57use crate::records::{
58    self, EXPIRY_PREFIX, Entry, Expiring, GraphRecord, GraphRunRecord, GraphRunState, NodeRecord,
59    ReadError, RecordError, RecordStatus,
60};
61use crate::task::{self, HEADER_GRAPH, TaskIdentity};
62
63/// The queue of the triggers.
64pub const TRIGGERS_QUEUE: &str = "swale-triggers";
65
66/// A failure of the scheduler.
67#[derive(Debug, thiserror::Error)]
68pub enum Error {
69    /// The queue failed.
70    #[error(transparent)]
71    Queue(#[from] taquba::Error),
72    /// The workflow runtime failed.
73    #[error(transparent)]
74    Workflow(#[from] taquba_workflow::Error),
75    /// The object store failed.
76    #[error(transparent)]
77    ObjectStore(#[from] taquba::object_store::Error),
78    /// A record is not valid JSON.
79    #[error(transparent)]
80    Record(#[from] RecordError),
81    /// The definition store failed.
82    #[error(transparent)]
83    Definition(#[from] DefinitionError),
84    /// The graph run records a definition the store does not have.
85    #[error("definition `{0}` is not in the definition store")]
86    UnknownDefinition(String),
87    /// The graph does not have a graph record, so the process did not adopt
88    /// a definition of the graph.
89    #[error("graph `{0}` does not have an adopted definition")]
90    UnknownGraph(String),
91    /// The trigger does not list a partition, and it is not a cron firing
92    /// with an earlier occurrence of the schedule.
93    #[error("the trigger of graph `{0}` does not determine a partition")]
94    NoPartition(String),
95    /// The pool of a node does not have a runtime.
96    #[error("node `{node}`: pool `{pool}` does not have a runtime")]
97    UnknownPool {
98        /// The node.
99        node: String,
100        /// The pool.
101        pool: String,
102    },
103    /// The graph run does not exist.
104    #[error("graph `{graph}` does not have a run for partition `{partition}`")]
105    UnknownGraphRun {
106        /// The graph.
107        graph: String,
108        /// The partition.
109        partition: Partition,
110    },
111    /// The node is not in the graph.
112    #[error("graph `{graph}` does not have a node `{node}`")]
113    UnknownNode {
114        /// The graph.
115        graph: String,
116        /// The node.
117        node: String,
118    },
119    /// The graph run record changed during the transition, which a retry
120    /// applies to the new record.
121    #[error("the run of graph `{graph}` for partition `{partition}` changed during the transition")]
122    Contended {
123        /// The graph.
124        graph: String,
125        /// The partition.
126        partition: Partition,
127    },
128}
129
130impl From<ReadError> for Error {
131    fn from(e: ReadError) -> Self {
132        match e {
133            ReadError::Queue(e) => Error::Queue(e),
134            ReadError::Record(e) => Error::Record(e),
135        }
136    }
137}
138
139impl Error {
140    /// Whether a retry cannot change the outcome: a graph, a graph run or a
141    /// node that does not exist, a trigger without a partition, a malformed
142    /// record, or a permanent error of the queue or the runtime. A request
143    /// with such an error is refused, and a worker with it dead-letters its
144    /// job.
145    pub fn is_permanent(&self) -> bool {
146        match self {
147            Error::Queue(e) => e.is_permanent(),
148            Error::Workflow(e) => e.is_permanent(),
149            Error::Record(_)
150            | Error::UnknownGraph(_)
151            | Error::UnknownGraphRun { .. }
152            | Error::UnknownNode { .. }
153            | Error::NoPartition(_) => true,
154            Error::ObjectStore(_)
155            | Error::Definition(_)
156            | Error::UnknownDefinition(_)
157            | Error::UnknownPool { .. }
158            | Error::Contended { .. } => false,
159        }
160    }
161}
162
163/// The settings of [`Scheduler::run`].
164#[derive(Debug, Clone)]
165pub struct SchedulerOptions {
166    /// Events handled at a time, and triggers handled at a time.
167    pub concurrency: usize,
168    /// The poll interval of the events worker and of the triggers worker.
169    pub poll_interval: Duration,
170    /// The time between reconciler passes.
171    pub reconcile_interval: Duration,
172}
173
174impl Default for SchedulerOptions {
175    fn default() -> Self {
176        SchedulerOptions {
177            concurrency: 4,
178            poll_interval: Duration::from_millis(250),
179            reconcile_interval: Duration::from_secs(60),
180        }
181    }
182}
183
184/// The result of [`Scheduler::start_run`].
185#[derive(Debug, Clone, PartialEq, Eq)]
186pub struct StartOutcome {
187    /// Whether this call created the graph run. `false` when a run for the
188    /// partition existed, in which case nothing was submitted.
189    pub started: bool,
190    /// The run ids of the root nodes submitted.
191    pub submitted: Vec<RunId>,
192}
193
194/// The result of [`Scheduler::rerun`].
195#[derive(Debug, Clone, PartialEq, Eq)]
196pub enum RerunOutcome {
197    /// The task instance at the next rerun count was submitted.
198    Submitted(RunId),
199    /// The task instance at the next rerun count is active from an earlier
200    /// rerun, so nothing was submitted.
201    Active(RunId),
202    /// The node does not have a record.
203    NoRecord,
204    /// The records of the node's upstreams do not satisfy its trigger rule.
205    NotReady,
206}
207
208/// Counts of one reconciler pass.
209#[derive(Debug, Clone, Default, PartialEq, Eq)]
210pub struct ReconcileReport {
211    /// Active graph runs visited.
212    pub active_runs: usize,
213    /// Nodes newly submitted. A submit of a run that was active is not
214    /// counted.
215    pub submitted: usize,
216    /// Graph runs whose state was written.
217    pub settled: usize,
218    /// Task instance runs cancelled for cancelled graph runs.
219    pub cancelled: usize,
220}
221
222/// The scheduler.
223pub struct Scheduler {
224    pub(crate) queue: Arc<Queue>,
225    definitions: Arc<DefinitionStore>,
226    pools: Arc<Pools>,
227    pub(crate) clock: Arc<dyn Clock>,
228    /// The expiry index of the records, at [`EXPIRY_PREFIX`].
229    pub(crate) expiry: ExpiryIndex,
230}
231
232impl Scheduler {
233    /// A scheduler over `queue`, running the graphs of `definitions` on
234    /// `pools`, with [`EVENTS_QUEUE`] as its events queue and
235    /// [`TRIGGERS_QUEUE`] as its triggers queue.
236    pub fn new(queue: Arc<Queue>, definitions: Arc<DefinitionStore>, pools: Arc<Pools>) -> Self {
237        let clock = queue.clock();
238        Scheduler {
239            queue,
240            definitions,
241            pools,
242            clock,
243            expiry: ExpiryIndex::new(EXPIRY_PREFIX),
244        }
245    }
246
247    /// The definition store.
248    pub fn definitions(&self) -> &Arc<DefinitionStore> {
249        &self.definitions
250    }
251
252    /// The graph record of `graph`, or `None` when the process did not adopt
253    /// a definition of the graph.
254    pub async fn graph_record(&self, graph: &str) -> Result<Option<GraphRecord>, Error> {
255        Ok(records::read(self.queue.view(), &records::graph_key(graph)).await?)
256    }
257
258    /// Starts the graph run of the definition `hash` for `partition`: writes
259    /// the graph run record and submits the root nodes. A second call for
260    /// the same graph and partition does not change anything.
261    pub async fn start_run(
262        &self,
263        hash: &str,
264        partition: &Partition,
265    ) -> Result<StartOutcome, Error> {
266        let graph = self.graph(hash).await?;
267        let record = GraphRunRecord {
268            definition: hash.to_string(),
269            requested_at_ms: self.clock.now_ms(),
270            state: GraphRunState::Active,
271            settled_at_ms: None,
272            expected_reruns: BTreeMap::new(),
273        };
274        let key = records::graph_run_key(graph.name(), partition);
275        if !self
276            .queue
277            .kv_compare_put(&key, None, &record.to_bytes())
278            .await?
279        {
280            return Ok(StartOutcome {
281                started: false,
282                submitted: Vec::new(),
283            });
284        }
285        let mut submitted = Vec::new();
286        for node in graph.roots() {
287            let (run_id, _) = self
288                .submit_node(
289                    node,
290                    identity(&graph, hash, partition, node, 0),
291                    &BTreeMap::new(),
292                    |_| SettlementEffects::default(),
293                )
294                .await?;
295            submitted.push(run_id);
296        }
297        Ok(StartOutcome {
298            started: true,
299            submitted,
300        })
301    }
302
303    /// Runs `node` again for `partition` at the next rerun count, when its
304    /// upstreams satisfy its trigger rule. The graph run returns to the
305    /// active state. After a succeeded record, the graph run record expects
306    /// the next count of the node and of every node in its
307    /// [`rerun_scope`], so each runs again with the new outputs. The outcome
308    /// states the run id submitted, or why the node was not rerun.
309    pub async fn rerun(
310        &self,
311        graph_name: &str,
312        partition: &Partition,
313        node: &str,
314    ) -> Result<RerunOutcome, Error> {
315        self.rerun_with(graph_name, partition, node, |_| {
316            SettlementEffects::default()
317        })
318        .await
319    }
320
321    /// [`Self::rerun`] with the effects of `effects`, given the run id,
322    /// committed with the submit.
323    pub(crate) async fn rerun_with(
324        &self,
325        graph_name: &str,
326        partition: &Partition,
327        node: &str,
328        effects: impl FnOnce(&RunId) -> SettlementEffects,
329    ) -> Result<RerunOutcome, Error> {
330        let key = records::graph_run_key(graph_name, partition);
331        let Some((run, bytes)) = self.graph_run(&key).await? else {
332            return Err(Error::UnknownGraphRun {
333                graph: graph_name.to_string(),
334                partition: partition.clone(),
335            });
336        };
337        let graph = self.graph(&run.definition).await?;
338        let node = graph.node(node).ok_or_else(|| Error::UnknownNode {
339            graph: graph_name.to_string(),
340            node: node.to_string(),
341        })?;
342        let records = self.node_records(&graph, partition).await?;
343        let Some(record) = records.get(node.name()) else {
344            return Ok(RerunOutcome::NoRecord);
345        };
346        let mut active = GraphRunRecord {
347            state: GraphRunState::Active,
348            settled_at_ms: None,
349            ..run.clone()
350        };
351        if record.status == RecordStatus::Succeeded && run.is_current(node.name(), record) {
352            for scoped in rerun_scope(&graph, node) {
353                if let Some(scoped_record) = records.get(scoped.name()) {
354                    active
355                        .expected_reruns
356                        .insert(scoped.name().to_string(), scoped_record.rerun + 1);
357                }
358            }
359        }
360        let current = current_records(&active, &records);
361        if !is_ready(node, &current) {
362            return Ok(RerunOutcome::NotReady);
363        }
364        if active != run
365            && !self
366                .queue
367                .kv_compare_put(&key, Some(&bytes), &active.to_bytes())
368                .await?
369        {
370            return Err(Error::Contended {
371                graph: graph_name.to_string(),
372                partition: partition.clone(),
373            });
374        }
375        let identity = identity(
376            &graph,
377            &record.definition,
378            partition,
379            node,
380            record.rerun + 1,
381        );
382        let (run_id, new) = self
383            .submit_node(node, identity, &upstream_records(node, &current), effects)
384            .await?;
385        Ok(if new {
386            RerunOutcome::Submitted(run_id)
387        } else {
388            RerunOutcome::Active(run_id)
389        })
390    }
391
392    /// Cancels the graph run: writes the cancelled state with the settle time
393    /// and cancels every active task instance. `false` when the run does not
394    /// exist or is not active.
395    pub async fn cancel_run(&self, graph_name: &str, partition: &Partition) -> Result<bool, Error> {
396        let key = records::graph_run_key(graph_name, partition);
397        let Some((run, bytes)) = self.graph_run(&key).await? else {
398            return Ok(false);
399        };
400        if run.state != GraphRunState::Active {
401            return Ok(false);
402        }
403        let cancelled = GraphRunRecord {
404            state: GraphRunState::Cancelled,
405            ..run.clone()
406        };
407        if !self
408            .commit_settled(graph_name, partition, &bytes, cancelled)
409            .await?
410        {
411            return Ok(false);
412        }
413        let graph = self.graph(&run.definition).await?;
414        self.cancel_active_runs(&graph, partition, &run).await?;
415        Ok(true)
416    }
417
418    /// Handles a cron firing of `graph_name`: starts the graph run of the
419    /// graph's adopted definition for the partition that contains
420    /// `interval_start_ms`, the occurrence of the schedule before the firing
421    /// time. Returns the partition, or `None` when its graph run existed.
422    pub async fn handle_trigger(
423        &self,
424        graph_name: &str,
425        interval_start_ms: Option<u64>,
426    ) -> Result<Option<Partition>, Error> {
427        let record = self.adopted(graph_name).await?;
428        let graph = self.graph(&record.definition).await?;
429        let partition = interval_start_ms
430            .and_then(|ms| Partition::of_time(graph.partitioning(), ms))
431            .ok_or_else(|| Error::NoPartition(graph_name.to_string()))?;
432        let started = self
433            .start_adopted(graph_name, &record, std::slice::from_ref(&partition))
434            .await?;
435        Ok(started.into_iter().next())
436    }
437
438    /// Starts the graph run of the adopted definition of `graph_name` for
439    /// every partition of `partitions`. A partition with a graph run is
440    /// unchanged. Returns the partitions whose graph run the call started.
441    pub async fn start_runs(
442        &self,
443        graph_name: &str,
444        partitions: &[Partition],
445    ) -> Result<Vec<Partition>, Error> {
446        let record = self.adopted(graph_name).await?;
447        self.start_adopted(graph_name, &record, partitions).await
448    }
449
450    async fn adopted(&self, graph_name: &str) -> Result<GraphRecord, Error> {
451        self.graph_record(graph_name)
452            .await?
453            .ok_or_else(|| Error::UnknownGraph(graph_name.to_string()))
454    }
455
456    async fn start_adopted(
457        &self,
458        graph_name: &str,
459        record: &GraphRecord,
460        partitions: &[Partition],
461    ) -> Result<Vec<Partition>, Error> {
462        let mut started = Vec::new();
463        for partition in partitions {
464            if self.start_run(&record.definition, partition).await?.started {
465                tracing::info!(graph = %graph_name, %partition, "graph run started");
466                started.push(partition.clone());
467            }
468        }
469        Ok(started)
470    }
471
472    /// Handles the termination of one task instance: submits every ready
473    /// downstream of its node, then settles the graph run.
474    pub async fn handle_event(&self, event: &Event) -> Result<(), Error> {
475        let key = records::graph_run_key(&event.graph, &event.partition);
476        let Some((run, bytes)) = self.graph_run(&key).await? else {
477            return Ok(());
478        };
479        if run.state != GraphRunState::Active {
480            return Ok(());
481        }
482        let graph = self.graph(&run.definition).await?;
483        let Some(node) = graph.node(&event.node) else {
484            return Ok(());
485        };
486        let downstreams: Vec<&Node> = node
487            .downstreams()
488            .iter()
489            .map(|name| {
490                graph
491                    .node(name)
492                    .expect("a downstream name is a node of the graph")
493            })
494            .collect();
495        self.advance(&graph, &event.partition, &run, &bytes, downstreams)
496            .await?;
497        Ok(())
498    }
499
500    /// One reconciler pass over every graph run record.
501    pub async fn reconcile(&self) -> Result<ReconcileReport, Error> {
502        let mut report = ReconcileReport::default();
503        let runs: Vec<Entry<GraphRunRecord>> =
504            records::scan(self.queue.view(), records::GRAPH_RUNS_PREFIX.as_bytes()).await?;
505        for Entry {
506            key,
507            bytes,
508            record: run,
509        } in runs
510        {
511            let Some((graph_name, partition)) = records::parse_graph_run_key(&key) else {
512                continue;
513            };
514            let graph = match self.definitions.get(&run.definition).await {
515                Ok(Some(graph)) => graph,
516                Ok(None) => {
517                    tracing::warn!(graph = %graph_name, %partition, definition = %run.definition, "graph run records an unknown definition");
518                    continue;
519                }
520                Err(e) => {
521                    tracing::warn!(graph = %graph_name, %partition, definition = %run.definition, error = %e, "the definition of a graph run does not load");
522                    continue;
523                }
524            };
525            match run.state {
526                GraphRunState::Active => {
527                    report.active_runs += 1;
528                    let (submitted, settled) = self
529                        .advance(&graph, &partition, &run, &bytes, graph.nodes())
530                        .await?;
531                    report.submitted += submitted;
532                    if settled {
533                        report.settled += 1;
534                    }
535                }
536                GraphRunState::Cancelled => {
537                    report.cancelled += self.cancel_active_runs(&graph, &partition, &run).await?;
538                }
539                GraphRunState::Complete | GraphRunState::Failed => {}
540            }
541        }
542        Ok(report)
543    }
544
545    /// Runs the events worker, the triggers worker and the reconciler until
546    /// `shutdown` resolves.
547    pub async fn run<F: Future<Output = ()>>(
548        self: Arc<Self>,
549        options: SchedulerOptions,
550        shutdown: F,
551    ) -> Result<(), Error> {
552        let stop = CancellationToken::new();
553        let worker = taquba::run_worker_concurrent(
554            &self.queue,
555            EVENTS_QUEUE,
556            self.clone(),
557            options.concurrency,
558            options.poll_interval,
559            stop.clone().cancelled_owned(),
560        );
561        let triggers = taquba::run_worker_concurrent(
562            &self.queue,
563            TRIGGERS_QUEUE,
564            Arc::new(TriggerWorker::new(self.clone())),
565            options.concurrency,
566            options.poll_interval,
567            stop.clone().cancelled_owned(),
568        );
569        let reconciler = async {
570            let mut interval = tokio::time::interval(options.reconcile_interval);
571            interval.tick().await;
572            loop {
573                tokio::select! {
574                    _ = interval.tick() => {
575                        if let Err(e) = self.reconcile().await {
576                            tracing::warn!(error = %e, "reconciler pass failed");
577                        }
578                    }
579                    () = stop.cancelled() => return,
580                }
581            }
582        };
583        let mut all = std::pin::pin!(async {
584            let (worker, triggers, ()) = tokio::join!(worker, triggers, reconciler);
585            worker.and(triggers)
586        });
587        tokio::select! {
588            result = &mut all => Ok(result?),
589            () = shutdown => {
590                stop.cancel();
591                Ok(all.await?)
592            }
593        }
594    }
595
596    /// Spawns [`Self::run`] as a task.
597    pub fn spawn<F>(
598        self: Arc<Self>,
599        options: SchedulerOptions,
600        shutdown: F,
601    ) -> WorkerHandle<Result<(), Error>>
602    where
603        F: Future<Output = ()> + Send + 'static,
604    {
605        WorkerHandle::spawn(shutdown, move |stop| async move {
606            self.run(options, stop.cancelled_owned()).await
607        })
608    }
609
610    async fn graph(&self, hash: &str) -> Result<Arc<Graph>, Error> {
611        self.definitions
612            .get(hash)
613            .await?
614            .ok_or_else(|| Error::UnknownDefinition(hash.to_string()))
615    }
616
617    async fn graph_run(&self, key: &[u8]) -> Result<Option<(GraphRunRecord, Vec<u8>)>, Error> {
618        let Some(bytes) = self.queue.view().kv_get(key).await? else {
619            return Ok(None);
620        };
621        let record = records::parse::<GraphRunRecord>(key, &bytes)?;
622        Ok(Some((record, bytes.to_vec())))
623    }
624
625    async fn node_record(
626        &self,
627        graph: &Graph,
628        partition: &Partition,
629        node: &Node,
630    ) -> Result<Option<NodeRecord>, Error> {
631        let key = records::node_record_key(graph.name(), partition, node);
632        Ok(records::read(self.queue.view(), &key).await?)
633    }
634
635    /// The node records of the graph run of `graph` for `partition`, by
636    /// node name.
637    async fn node_records(
638        &self,
639        graph: &Graph,
640        partition: &Partition,
641    ) -> Result<BTreeMap<String, NodeRecord>, Error> {
642        let mut records = BTreeMap::new();
643        for node in graph.nodes() {
644            if let Some(record) = self.node_record(graph, partition, node).await? {
645                records.insert(node.name().to_string(), record);
646            }
647        }
648        Ok(records)
649    }
650
651    /// Submits every ready node among `candidates` at its next rerun count,
652    /// then settles the active graph run. Returns the count of new submits
653    /// and whether the final state was written.
654    async fn advance<'a>(
655        &self,
656        graph: &Graph,
657        partition: &Partition,
658        run: &GraphRunRecord,
659        bytes: &[u8],
660        candidates: impl IntoIterator<Item = &'a Node>,
661    ) -> Result<(usize, bool), Error> {
662        let records = self.node_records(graph, partition).await?;
663        let current = current_records(run, &records);
664        let states = node_states(graph, &current);
665        let mut submitted = 0;
666        for node in candidates {
667            if states[node.name()] != NodeState::Ready {
668                continue;
669            }
670            let rerun = records.get(node.name()).map_or(0, |r| r.rerun + 1);
671            let (_, new) = self
672                .submit_node(
673                    node,
674                    identity(graph, &run.definition, partition, node, rerun),
675                    &upstream_records(node, &current),
676                    |_| SettlementEffects::default(),
677                )
678                .await?;
679            if new {
680                submitted += 1;
681            }
682        }
683        let settled = self
684            .settle(graph, partition, run, bytes, &records, &states)
685            .await?;
686        Ok((submitted, settled))
687    }
688
689    /// Submits the task instance `identity` of `node`, with the effects of
690    /// `effects` committed with a new submit. Returns the run id and whether
691    /// the submit was new.
692    async fn submit_node(
693        &self,
694        node: &Node,
695        identity: TaskIdentity,
696        upstreams: &BTreeMap<String, NodeRecord>,
697        effects: impl FnOnce(&RunId) -> SettlementEffects,
698    ) -> Result<(RunId, bool), Error> {
699        let runtime = self
700            .pools
701            .runtime(node.pool())
702            .ok_or_else(|| Error::UnknownPool {
703                node: node.name().to_string(),
704                pool: node.pool().to_string(),
705            })?;
706        let run_id = identity.run_id();
707        let outcome = runtime
708            .submit(RunSpec {
709                run_id: Some(run_id.clone()),
710                input: TaskInput::new(node, upstreams).to_bytes(),
711                options: RunOptions {
712                    headers: identity.headers(),
713                    max_attempts_per_step: Some(node.retries() + 1),
714                    ..RunOptions::default()
715                },
716                effects: effects(&run_id),
717            })
718            .await?;
719        if outcome.newly_submitted {
720            tracing::info!(run_id = %run_id, pool = node.pool(), "task instance submitted");
721        }
722        Ok((run_id, outcome.newly_submitted))
723    }
724
725    /// Writes the final state of the graph run when it is reached: the state of
726    /// [`settled_state`] with the settle time, once no unrecorded task instance
727    /// is active. The write removes an expected rerun count that a record
728    /// reached, and the expiry index entry of the run commits with it. `true`
729    /// when the state was written.
730    async fn settle(
731        &self,
732        graph: &Graph,
733        partition: &Partition,
734        run: &GraphRunRecord,
735        bytes: &[u8],
736        records: &BTreeMap<String, NodeRecord>,
737        states: &BTreeMap<String, NodeState>,
738    ) -> Result<bool, Error> {
739        let Some(state) = settled_state(states) else {
740            return Ok(false);
741        };
742        // A blocked node can have an active run at count 0 from the time it was
743        // ready, and a failed or cancelled node can have an active rerun.
744        for node in graph.nodes() {
745            if let Some(run_id) = unrecorded_run_id(graph, partition, run, node, records)
746                && self.run_is_active(node, &run_id).await?
747            {
748                return Ok(false);
749            }
750        }
751        let mut settled = GraphRunRecord {
752            state,
753            ..run.clone()
754        };
755        settled.expected_reruns.retain(|name, expected| {
756            records
757                .get(name)
758                .is_none_or(|record| record.rerun < *expected)
759        });
760        let written = self
761            .commit_settled(graph.name(), partition, bytes, settled)
762            .await?;
763        if written {
764            tracing::info!(graph = graph.name(), %partition, state = ?state, "graph run settled");
765        }
766        Ok(written)
767    }
768
769    /// Commits the graph run record `settled` with the clock's time as its
770    /// settle time and the expiry index entry of the run at that time, against
771    /// the stored bytes `expected` of the active record. `false` when the
772    /// stored record differs from `expected`.
773    async fn commit_settled(
774        &self,
775        graph_name: &str,
776        partition: &Partition,
777        expected: &[u8],
778        settled: GraphRunRecord,
779    ) -> Result<bool, Error> {
780        let key = records::graph_run_key(graph_name, partition);
781        let settled_at_ms = self.clock.now_ms();
782        let settled = GraphRunRecord {
783            settled_at_ms: Some(settled_at_ms),
784            ..settled
785        };
786        let expiring = Expiring::Run {
787            graph: graph_name.to_string(),
788            partition: partition.clone(),
789        };
790        let effects = SettlementEffects::default()
791            .kv_put(key.clone(), settled.to_bytes())
792            .expiry_entry(&self.expiry, settled_at_ms, &expiring.suffix());
793        Ok(self
794            .queue
795            .kv_compare_commit(&key, Some(expected), effects)
796            .await?
797            .is_some())
798    }
799
800    async fn run_is_active(&self, node: &Node, run_id: &RunId) -> Result<bool, Error> {
801        let Some(runtime) = self.pools.runtime(node.pool()) else {
802            return Ok(false);
803        };
804        Ok(matches!(
805            runtime.status(run_id).await?,
806            Some(status) if !matches!(status.state, RunState::Terminated(_))
807        ))
808    }
809
810    /// Cancels the unrecorded task instance run of every node. Returns the
811    /// count of runs cancelled.
812    async fn cancel_active_runs(
813        &self,
814        graph: &Graph,
815        partition: &Partition,
816        run: &GraphRunRecord,
817    ) -> Result<usize, Error> {
818        let records = self.node_records(graph, partition).await?;
819        let mut cancelled = 0;
820        for node in graph.nodes() {
821            let Some(run_id) = unrecorded_run_id(graph, partition, run, node, &records) else {
822                continue;
823            };
824            let Some(runtime) = self.pools.runtime(node.pool()) else {
825                continue;
826            };
827            if runtime.cancel(&run_id).await? {
828                cancelled += 1;
829            }
830        }
831        Ok(cancelled)
832    }
833}
834
835/// The records of the upstreams of `node` among `records`, by node name.
836fn upstream_records(
837    node: &Node,
838    records: &BTreeMap<String, NodeRecord>,
839) -> BTreeMap<String, NodeRecord> {
840    node.upstreams()
841        .iter()
842        .filter_map(|name| Some((name.clone(), records.get(name)?.clone())))
843        .collect()
844}
845
846/// The run id of the task instance of `node` that does not have a current
847/// record. It is the run at count 0 of a node without a record, and the run
848/// at the next count after a failed, cancelled or superseded record. `None`
849/// after a current succeeded record.
850fn unrecorded_run_id(
851    graph: &Graph,
852    partition: &Partition,
853    run: &GraphRunRecord,
854    node: &Node,
855    records: &BTreeMap<String, NodeRecord>,
856) -> Option<RunId> {
857    let rerun = match records.get(node.name()) {
858        Some(record)
859            if record.status == RecordStatus::Succeeded && run.is_current(node.name(), record) =>
860        {
861            return None;
862        }
863        Some(record) => record.rerun + 1,
864        None => 0,
865    };
866    Some(task::run_id(graph.name(), partition, node.name(), rerun))
867}
868
869/// The identity of the task instance of `node` at the rerun count `rerun`.
870fn identity(
871    graph: &Graph,
872    hash: &str,
873    partition: &Partition,
874    node: &Node,
875    rerun: u32,
876) -> TaskIdentity {
877    TaskIdentity {
878        graph: graph.name().to_string(),
879        partition: partition.clone(),
880        node: node.name().to_string(),
881        asset: node.asset().map(str::to_string),
882        definition: hash.to_string(),
883        rerun,
884    }
885}
886
887/// The failure of a worker for `error`: permanent when the error is, and
888/// retried otherwise.
889fn worker_error(error: Error) -> WorkerError {
890    if error.is_permanent() {
891        PermanentFailure::new(error.to_string()).into()
892    } else {
893        Box::new(error)
894    }
895}
896
897/// The headers of the cron schedule of `graph`, which every firing of the
898/// schedule includes.
899pub fn firing_headers(graph: &str) -> HashMap<String, String> {
900    HashMap::from([(HEADER_GRAPH.to_string(), graph.to_string())])
901}
902
903/// The [`Worker`] of the triggers queue.
904pub struct TriggerWorker {
905    scheduler: Arc<Scheduler>,
906}
907
908impl TriggerWorker {
909    /// A worker that starts graph runs on `scheduler`.
910    pub fn new(scheduler: Arc<Scheduler>) -> Self {
911        TriggerWorker { scheduler }
912    }
913}
914
915impl Worker for TriggerWorker {
916    async fn process(&self, job: &JobRecord, _lease: &LeaseHandle) -> Result<(), WorkerError> {
917        let graph = job.headers.get(HEADER_GRAPH).ok_or_else(|| {
918            PermanentFailure::new(format!("the job does not have the `{HEADER_GRAPH}` header"))
919        })?;
920        let interval_start_ms = job
921            .headers
922            .get(PREVIOUS_FIRE_MS_HEADER)
923            .and_then(|value| value.parse().ok());
924        self.scheduler
925            .handle_trigger(graph, interval_start_ms)
926            .await
927            .map(|_| ())
928            .map_err(worker_error)
929    }
930}
931
932impl Worker for Scheduler {
933    async fn process(&self, job: &JobRecord, _lease: &LeaseHandle) -> Result<(), WorkerError> {
934        let event = Event::from_bytes(&job.payload)
935            .map_err(|e| PermanentFailure::new(format!("the payload is not an event: {e}")))?;
936        self.handle_event(&event).await.map_err(worker_error)
937    }
938}