work_dispatcher 0.1.1

A simple concurrent data processing framework
Documentation
use anyhow::Result;
use async_trait::async_trait;
use flume::Sender;
use std::time::Duration;
use tracing::{Level, info};
use work_dispatcher::{Processor, Producer, WorkDispatcher};

#[derive(Debug)]
struct CsvRow {
    id: i32,
    data: String,
}

#[derive(Clone)]
struct DbConnection(String);

struct CsvProducer {
    count: usize,
}

#[async_trait]
impl Producer for CsvProducer {
    type Item = CsvRow;
    async fn run(self, sender: Sender<Self::Item>) {
        info!("[Producer] Starting to generate {} records.", self.count);
        for i in 0..self.count {
            let row = CsvRow {
                id: i as i32,
                data: format!("Data for row {}", i),
            };
            sender.send_async(row).await.unwrap();
        }
    }
}

#[derive(Clone)]
struct DbProcessor;

#[async_trait]
impl Processor for DbProcessor {
    type Item = CsvRow;
    type Context = DbConnection;
    async fn process(&self, item: Self::Item, context: &Self::Context) {
        info!(
            "[Processor] Writing item #{} data:{} to DB '{}'",
            item.id, item.data, context.0
        );
        // Simulate an async database call
        tokio::time::sleep(Duration::from_secs(1)).await;
    }
}

#[tokio::main]
async fn main() -> Result<()> {
    tracing_subscriber::fmt()
        .with_max_level(Level::DEBUG)
        .init();

    // 1. Setup producer, processor, and context
    let producer = CsvProducer { count: 100 };
    let processor = DbProcessor;
    let context = DbConnection("mongodb://localhost:27017".to_string());

    // 2. Create and run the dispatcher
    WorkDispatcher::new(producer, processor, context)
        .workers(8)
        .buffer(500)
        .run()
        .await?;

    Ok(())
}