apalis-postgres 1.0.0-rc.9

Background task processing for rust using apalis and postgres
Documentation
use std::time::Duration;

use apalis::prelude::*;
use apalis_postgres::{Config, *};
use apalis_workflow::*;

async fn get_name(user_id: u32) -> Result<String, BoxDynError> {
    Ok(user_id.to_string())
}

async fn get_age(user_id: u32) -> Result<usize, BoxDynError> {
    Ok(user_id as usize + 20)
}

async fn get_address(user_id: u32) -> Result<usize, BoxDynError> {
    Ok(user_id as usize + 100)
}

async fn collector(
    (name, age, address): (String, usize, usize),
    wrk: WorkerContext,
) -> Result<usize, BoxDynError> {
    let result = name.parse::<usize>()? + age + address;
    wrk.stop().unwrap();
    Ok(result)
}

#[tokio::main]
async fn main() {
    let dag_flow = GraphFlow::new("user-etl-workflow");
    let get_name = dag_flow.add_task(get_name);
    let get_age = dag_flow.add_task(get_age);
    let get_address = dag_flow.add_task(get_address);
    dag_flow
        .add_task(collector)
        .depends_on((&get_name, &get_age, &get_address)); // Order and types matters here

    dag_flow.validate().unwrap();

    println!("Executing workflow:\n{}", dag_flow);

    let pool = PgPool::connect(&std::env::var("DATABASE_URL").unwrap())
        .await
        .unwrap();
    PostgresStorage::setup(&pool).await.unwrap();

    let config = Config::default().queue("test-workflow");
    let mut backend = PostgresStorage::new(&pool)
        .with_config(config)
        .with_pubsub()
        .poll_with_interval(Duration::from_secs(1)); // Wake the worker once a second if its sleeping

    backend.start_fan_out(vec![42u32, 43, 44]).await.unwrap();

    let worker = WorkerBuilder::new("rango-tango")
        .backend(backend)
        .on_event(|ctx, ev| {
            println!("On Event = {:?}", ev);
            if matches!(ev, Event::Error(_)) {
                ctx.stop().unwrap();
            }
        })
        .build(dag_flow);

    worker.run().await.unwrap();
}