Skip to main content

run_steps

Function run_steps 

Source
pub async fn run_steps<C, F, Fut>(
    budget: &Budget,
    cost: u32,
    cursor: C,
    step: F,
) -> Result<Option<C>>
where F: FnMut(C) -> Fut, Fut: Future<Output = Result<Step<C>>>,
Expand description

Runs step from cursor while budget covers its cost (the calls one step makes), for work too large for one invocation: a page of rows per step, the next page’s cursor between them (Active Job’s continuations).

Returns Some(cursor) when the budget ran out first: enqueue the job again with it, and the next run continues there. None when a step answered Step::Done. A failing step stops the run with its error; a retried job starts again from the cursor it was enqueued with, so make steps safe to repeat.

§Errors

The first error a step returns.

§Examples

use ocre::jobs::{Budget, Step, run_steps};
use ocre::{Ctx, Query, Result};
use serde::{Deserialize, Serialize};

#[derive(Deserialize)]
struct Row {
    id: i64,
}

#[derive(Serialize, Deserialize)]
pub struct Reindex {
    after_id: i64,
}

impl Reindex {
    pub async fn perform(self, ctx: &Ctx) -> Result<()> {
        let budget = Budget::new(Budget::FREE_D1_QUERIES - 1); // 1 left to enqueue the rest
        // Each step reads a page of 100 rows and writes them back: 2 queries.
        let rest = run_steps(&budget, 2, self.after_id, |after_id| async move {
            let db = ctx.db()?;
            let page: Vec<Row> = Query::table("tracks").gt("id", after_id).order_asc("id").limit(100).all(&db).await?;
            let Some(last) = page.last() else { return Ok(Step::Done) };
            // ... one bulk::update of the page ...
            Ok(Step::Next(last.id))
        })
        .await?;
        if let Some(after_id) = rest {
            ocre::jobs::enqueue(ctx, &Reindex { after_id }).await?;
        }
        Ok(())
    }
}