Last updated on

2. Adaptive Scheduling of Asynchronous Tasks in Rust


Introduction

In the previous article I described why I’ve chosen to use a public API for my web service to fetch blockchain data instead of maintaining a own node. In this follow-up I am going to describe the mechanics behind my module handling these async API-Calls written in Rust.

A simplified model is depicted in figure 1, illustrating a system where a Scheduler initiates RunningTasks, which can subsequently be halted and transitioned into StaleTasks. A RunningTask periodically calls a process function which is a part of the Schedulable Trait. This design aims not only to fulfil the immediate need for a dynamic polling mechanism but also to ensure the module’s versatility for potential integration into future projects.

fig 1. Scheduler Model
fig 1. Scheduler Model

Solution

fig 2. Scheduling Intervals Timeline
fig 2. Scheduling Intervals Timeline
Time SpanFrequency
ts0f0
ts1f1
ts2f2
after ts3Terminate

This whole setup of adjusting how often we check for updates is neatly packaged into a struct called ScheduleInterval:

#[derive(Clone, Debug)]
pub struct ScheduleInterval {
    pub duration: Duration, // This is how long the current checking spree lasts.
    pub frequency: Duration, // And this is how often we peek in to see if there's any new Bitcoin.
}

We slot these intervals into our config, the ScheduleConfig<Ctx> struct:

#[derive(Clone, Debug, Error)]
pub struct ScheduleConfig<Ctx> {
    pub ranges: Vec<ScheduleInterval>, // These are the different time frames we've set for checking in.
    pub clone_data: bool, // Tells us if we need to make copies of the data for each check.
    pub init_processing: bool, // Decides if we do a quick check right when we start up.
    pub ctx: Arc<Ctx>, // This is a shared context that all our tasks use.
}

In this ScheduleConfig, we’ve also got the option to clone our data for each task if needed. This flexibility is key for making our setup work smoothly with different kinds of data.

Clone or not to clone?

Rust’s borrow checker ensures that only one piece of code can mutate data at any time. This means if we send our data off to be worked on, we temporarily can’t use it for anything else. But, if we enable cloning, we can send a copy off for processing and still keep the original around for reading or writing (although, updating it might not be super useful since the processed results will overwrite it anyway). This cloning trick is especially slick for quick, light data that needs to be accessible all the time. Imagine someone’s trying to check the status of a transaction, but it’s locked up in processing – not ideal.

However, if we’re dealing with hefty data or it’s not critical to have it available every second, then skipping cloning and instead leaving a placeholder while the data’s being crunched is totally fine. This way, anyone needing the data will just have to wait a bit until the processing is finished. It’s all about balancing the need for speed and access against the weight of the data we’re juggling.

Kickstarting with Initial Processing

The init_processing flag is what tells our scheduler whether to fire up the fn process function right out of the gate. Without this flag, we’re leaving room for guesswork. Everyone might have their own take on whether it’s go-time for our tasks. By setting this flag, we clear up any ambiguity, ensuring everyone’s on the same page about when our tasks spring into action.

The Role of Context

In our scheduler, the generic Ctx is our Swiss Army knife for customizing how tasks interact with their environment. It’s like giving our tasks a backpack of tools before they head out on their journey. Need to keep clients in the loop with real-time updates? Pack in some ‘Server-Sent Events’. Want your tasks to chat with a database? Toss in a database connection.

This flexibility lets us tailor the context to suit the specific demands of each task, ensuring they have everything they need to succeed. Whether it’s dispatching updates to eager clients or accessing crucial data on the fly, the context makes it all possible, adapting seamlessly to the varied landscapes of our projects.

Schedulable Trait

#[async_trait]
pub trait Schedulable<K, V, Ctx>
where
    K: Clone + PartialEq + Eq + Hash + Send + Sync + 'static + Debug,
    V: Clone + Send + Sync + 'static + Debug,
    Ctx: Clone + Send + Sync + 'static + Debug,
{

    fn id(&self) -> K;
    async fn process(mut self, ctx: Arc<Ctx>) -> (V, TaskStatus)
    async fn time_out(self, ctx: Arc<Ctx>);
    async fn complete(self, ctx: Arc<Ctx>);
}

The Schedulable trait is designed as a contract for tasks to be managed by a Scheduler in Rust, requiring implementers to define behaviors which can occur during execution. It mandates that tasks be identifiable via a unique key (K) and that they carry task-specific data (V) along with a shared context (Ctx), all of which must adhere to common Rust traits for cloning, comparison, and asynchronous safety. Implementers must provide an id method for task identification, an asynchronous process method to define the task’s main logic (returning updated task data and a TaskStatus), and asynchronous time_out and complete methods for handling task expiry and successful completion, respectively.

The process method’s ability to return a TaskStatus allows for dynamic state management, facilitating tasks to be marked as running, canceled, completed, (or timed out), all depending on the conclusions drawn after fn process execution.

#[derive(Debug, Error, Clone, PartialEq)]
pub enum TaskStatus {
    Running,
    Canceled,
    Completed,
    Timeout
}

Unpacking the Scheduler Mechanics

Equally to our Schedulable trait we utilize the generics K, V , and Ctx.

Within the Scheduler’s structure, tasks are categorized into two main states: actively running and stale. Running tasks are securely managed within a thread-safe container (Arc<Mutex<HashMap<K, RunningTask<K, V, Ctx>>>>), ensuring concurrent operations are safely executed. Stale tasks, on the other hand, have completed their cycle and are stored in a regular HashMap, awaiting further instructions without the need for thread-safe protections.

A pivotal component of this system is the running_task_tx_sender, a messaging channel that allows tasks to communicate their id and status back to the Scheduler.

pub struct Scheduler<K, V, Ctx> {
    pub running_tasks: Arc<Mutex<HashMap<K, RunningTask<K, V, Ctx>>>>,
    pub stale_tasks: HashMap<K, StaleTask<K, V, Ctx>>,
    pub config: Arc<ScheduleConfig<Ctx>>,
    running_task_tx_sender: mpsc::Sender<(K, TaskStatus)>,
}

For the Scheduler to fulfill its purpose, it requires a suite of functions to manage the lifecycle of tasks:

  • new to instantiate the Scheduler with a given configuration.
  • add to register a new task within the Scheduler.
  • remove to extract and optionally return a completed or canceled task.
  • inspect to examine the current state of a task without modifying it.
  • stop to gracefully terminate an active task.
  • run to initiate or resume the execution of a task according to its schedule.
  • An internal listener function to catch incoming timeout signals from tasks.
pub fn new(config: ScheduleConfig<Ctx>) -> Self,
pub fn add(&mut self, item: V) -> K,
pub async fn remove(&mut self, id: K) -> Option<V>,
pub async fn inspect(&self, id: &K, f: impl Fn(RwLockReadGuard<'_, Option<V>>),) -> Result<(), SchedulerError>,
pub async fn stop(&mut self, id: K) -> Result<(), SchedulerError>,
pub async fn run(&mut self, id: K) -> Result<(), SchedulerError> ,
fn listener(mut receiver: mpsc::Receiver<(K, TaskStatus)>, running_tasks: Arc<Mutex<HashMap<K, RunningTask<K, V, Ctx>>>>,config: Arc<ScheduleConfig<Ctx>>)

These functions collectively empower the Scheduler to dynamically manage tasks, ensuring they are executed, monitored, and concluded in a controlled and orderly fashion.

Limitations

Schedulable Conditional Trait - V:Clone

#[async_trait]
pub trait Schedulable<K, V, Ctx>
where
    ...
    V: Clone + Send + Sync + 'static + Debug, //V Clone not necessary if Configuration is set to false
    ...
{

...
    async fn process(mut self, ctx: Arc<Ctx>) -> (V, TaskStatus)
...
}

The limitation around the Schedulable trait’s requirement for V to implement Clone stems from the requirement for type V to implement Clone, contingent on a flag within the scheduler’s configuration. When the clone flag is false, no cloning occurs, theoretically eliminating the need for V to implement Clone. However, the trait’s current design obliges all task data types to support cloning regardless, potentially burdening developers with unnecessary implementation for heavier data types where cloning isn’t utilized. This requirement can add unnecessary complexity.

Handling TaskStatus::Timeout in Fn Process

When a task exceeds the designated time span, it emits to the Scheduler and triggers the TaskStatus::Timeout state. This particular status arises as a direct consequence of surpassing the allocated time frame for task completion. As such, it’s not logically coherent for a extern entity to explicitly set its status to TaskStatus::Timeout through the process function. This function’s primary role is to perform the task’s core logic, while the management of timeout conditions is inherently handled by the scheduler’s timing mechanisms. Therefore, TaskStatus::Timeout serves as an indicator of a system-enforced state transition, rather than a status to be manually assigned by the task’s processing logic.

Cheers!

👉🏻 Check out the code here