pub async fn run_steps<C, F, Fut>(
budget: &Budget,
cost: u32,
cursor: C,
step: F,
) -> Result<Option<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(())
}
}