Cómo construir un CLI concurrente en Rust con Tokio y Clap

rust Cómo construir un CLI concurrente en Rust con Tokio y Clap

Cómo construir un CLI concurrente en Rust con Tokio y Clap

Objetivo: crear un CLI que procese varios archivos en paralelo con control de concurrencia, buen manejo de errores y logging estructurado.

Requisitos

  • Rust (stable, edición 2021)
  • Conocimientos básicos de async/await
  • Cargo instalado

Estructura del proyecto

my-cli/
├─ Cargo.toml
└─ src/
   └─ main.rs

Cargo.toml

[package]
name = "my-cli"
version = "0.1.0"
edition = "2021"

[dependencies]
tokio = { version = "1", features = ["full"] }
clap = { version = "4", features = ["derive"] }
tracing = "0.1"
tracing-subscriber = "0.3"
anyhow = "1.0"
futures = "0.3"

src/main.rs

use clap::Parser;
use tokio::sync::Semaphore;
use anyhow::{Context, Result};
use tracing::{info, error};
use std::sync::Arc;
use futures::stream::{FuturesUnordered, StreamExt};

#[derive(Parser, Debug)]
#[command(about = "CLI concurrente para procesar archivos")]
struct Args {
    /// Paths de archivos a procesar
    #[arg(required = true)]
    files: Vec<:path::pathbuf>,

    /// Límite de concurrencia
    #[arg(short, long, default_value_t = 4)]
    concurrency: usize,
}

#[tokio::main]
async fn main() -> Result<()> {
    tracing_subscriber::fmt().init();
    let args = Args::parse();

    let sem = Arc::new(Semaphore::new(args.concurrency));
    let mut tasks = FuturesUnordered::new();

    for path in args.files {
        let permit = sem.clone().acquire_owned().await.unwrap();
        let path_clone = path.clone();
        tasks.push(tokio::spawn(async move {
            let _permit = permit; // mantiene el permiso hasta que la tarea termine
            match process_file(path_clone).await {
                Ok(len) => info!(file = ?path_clone, bytes = len, "Procesado"),
                Err(e) => error!(file = ?path_clone, error = ?e, "Error"),
            }
        }));
    }

    while let Some(res) = tasks.next().await {
        if let Err(e) = res {
            error!("task panicked: {:?}", e);
        }
    }

    Ok(())
}

async fn process_file(path: std::path::PathBuf) -> Result {
    let data = tokio::fs::read(&path).await.context("lectura")?;
    // Simular trabajo CPU-bound lanzándolo al pool de bloqueo
    let len = tokio::task::spawn_blocking(move || heavy_work(&data))
        .await
        .context("spawn_blocking")??;
    Ok(len)
}

fn heavy_work(data: &[u8]) -> usize {
    // Ejemplo: contar ocurrencias de la letra 'e'
    data.iter().filter(|&&b| b == b'e').count()
}

Explicación de decisiones

  • Tokio: runtime async maduro para IO concurrente y utilidades como spawn_blocking y semáforos.
  • Clap (derive): parseo de argumentos y documentación automática del CLI.
  • Semaphore (tokio::sync::Semaphore): limita el número de tareas activas para evitar saturación de I/O o memoria.
  • spawn_blocking: cualquier trabajo CPU-bound debe ejecutarse fuera del hilo async principal para no bloquear el reactor.
  • tracing + tracing-subscriber: logging estructurado, útil en concurrencia para filtrar y correlacionar eventos.
  • anyhow: manejo de errores con contexto sin boilerplate excesivo.
  • FuturesUnordered: recolectar tareas concurrentes y esperar sus resultados mientras llegan.

Uso

cargo run -- file1.txt file2.txt -c 8

Puntos importantes

  • No uses tokio::fs::read_to_string para archivos muy grandes si solo necesitas bytes: read evita conversiones innecesarias.
  • Usa spawn_blocking para operaciones CPU-bound; de lo contrario bloquearás el runtime y degradarás otras tareas async.
  • Semaphore controla la concurrencia efectiva y evita OOM o saturación del disco/FS.
  • Añade contexto a errores con anyhow::Context para diagnósticos más útiles en producción.

Consejo avanzado: si prefieres un enfoque más funcional, reemplaza el bucle y FuturesUnordered por

use futures::stream::{self, StreamExt};

stream::iter(args.files)
    .map(|path| {
        let sem = sem.clone();
        async move {
            let _permit = sem.acquire_owned().await.unwrap();
            process_file(path).await
        }
    })
    .buffer_unordered(args.concurrency)
    .for_each(|res| async move { /* manejar resultado */ })
    .await;

Puedes ampliar esta base añadiendo métricas (prometheus), retries exponenciales (tokio-retry) o un modo streaming para procesar archivos grandes por chunks.

Advertencia: evita mezclar operaciones blocking intensivas sin offloading; aunque el ejemplo usa spawn_blocking, operaciones muy pesadas en gran número pueden necesitar un pool dedicado o arquitectura diferente.

Comentarios
¿Quieres comentar?

Inicia sesión con Telegram para participar en la conversación


Comentarios (0)

Aún no hay comentarios. ¡Sé el primero en comentar!

Iniciar Sesión