1use 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
63pub const TRIGGERS_QUEUE: &str = "swale-triggers";
65
66#[derive(Debug, thiserror::Error)]
68pub enum Error {
69 #[error(transparent)]
71 Queue(#[from] taquba::Error),
72 #[error(transparent)]
74 Workflow(#[from] taquba_workflow::Error),
75 #[error(transparent)]
77 ObjectStore(#[from] taquba::object_store::Error),
78 #[error(transparent)]
80 Record(#[from] RecordError),
81 #[error(transparent)]
83 Definition(#[from] DefinitionError),
84 #[error("definition `{0}` is not in the definition store")]
86 UnknownDefinition(String),
87 #[error("graph `{0}` does not have an adopted definition")]
90 UnknownGraph(String),
91 #[error("the trigger of graph `{0}` does not determine a partition")]
94 NoPartition(String),
95 #[error("node `{node}`: pool `{pool}` does not have a runtime")]
97 UnknownPool {
98 node: String,
100 pool: String,
102 },
103 #[error("graph `{graph}` does not have a run for partition `{partition}`")]
105 UnknownGraphRun {
106 graph: String,
108 partition: Partition,
110 },
111 #[error("graph `{graph}` does not have a node `{node}`")]
113 UnknownNode {
114 graph: String,
116 node: String,
118 },
119 #[error("the run of graph `{graph}` for partition `{partition}` changed during the transition")]
122 Contended {
123 graph: String,
125 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 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#[derive(Debug, Clone)]
165pub struct SchedulerOptions {
166 pub concurrency: usize,
168 pub poll_interval: Duration,
170 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#[derive(Debug, Clone, PartialEq, Eq)]
186pub struct StartOutcome {
187 pub started: bool,
190 pub submitted: Vec<RunId>,
192}
193
194#[derive(Debug, Clone, PartialEq, Eq)]
196pub enum RerunOutcome {
197 Submitted(RunId),
199 Active(RunId),
202 NoRecord,
204 NotReady,
206}
207
208#[derive(Debug, Clone, Default, PartialEq, Eq)]
210pub struct ReconcileReport {
211 pub active_runs: usize,
213 pub submitted: usize,
216 pub settled: usize,
218 pub cancelled: usize,
220}
221
222pub struct Scheduler {
224 pub(crate) queue: Arc<Queue>,
225 definitions: Arc<DefinitionStore>,
226 pools: Arc<Pools>,
227 pub(crate) clock: Arc<dyn Clock>,
228 pub(crate) expiry: ExpiryIndex,
230}
231
232impl Scheduler {
233 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 pub fn definitions(&self) -> &Arc<DefinitionStore> {
249 &self.definitions
250 }
251
252 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 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 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 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, ¤t) {
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, ¤t), effects)
384 .await?;
385 Ok(if new {
386 RerunOutcome::Submitted(run_id)
387 } else {
388 RerunOutcome::Active(run_id)
389 })
390 }
391
392 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 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 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 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 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 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 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 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 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, ¤t);
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, ¤t),
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 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 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 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 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 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
835fn 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
846fn 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
869fn 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
887fn 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
897pub fn firing_headers(graph: &str) -> HashMap<String, String> {
900 HashMap::from([(HEADER_GRAPH.to_string(), graph.to_string())])
901}
902
903pub struct TriggerWorker {
905 scheduler: Arc<Scheduler>,
906}
907
908impl TriggerWorker {
909 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}