Add modular data structures, parallel pipelines, and system profiling
Establishes foundational infrastructure across multiple crates by introducing unified persistent bit matrix storage with columnar, packed, and implicit variants, alongside De Bruijn graph node encoding and unitig iteration logic. Adds a macro-driven parallel pipeline scheduler featuring NUMA-aware runners, bounded channels, and memory budgets to enforce concurrency limits. Implements streaming nucleotide parsers with pooled page buffers for FASTA, FASTQ, and Genbank formats, complemented by system resource monitoring, progress tracking, and stage profiling utilities. Collectively, these changes provide the core data models, execution frameworks, and I/O pipelines required for downstream k-mer indexing and analysis workloads.
This commit is contained in:
@@ -1,878 +0,0 @@
|
||||
use crossbeam_channel::{Receiver, Select, Sender, bounded};
|
||||
use std::error::Error;
|
||||
use std::fmt;
|
||||
use std::marker::PhantomData;
|
||||
use std::sync::Arc;
|
||||
use std::thread;
|
||||
|
||||
/// Error type for pipeline operations.
|
||||
#[derive(Debug)]
|
||||
pub enum PipelineError {
|
||||
/// A stage received a `PipelineData` variant it did not expect.
|
||||
TypeMismatch,
|
||||
/// The step kind is not compatible with the data type.
|
||||
StepKindMismatch(&'static str),
|
||||
/// The source has no more data to produce.
|
||||
EndOfStream,
|
||||
/// An error occurred inside a stage (e.g., I/O, parsing, custom logic).
|
||||
StepError(Box<dyn Error + Send + Sync>),
|
||||
}
|
||||
|
||||
impl fmt::Display for PipelineError {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
match self {
|
||||
PipelineError::TypeMismatch => write!(f, "data type mismatch in pipeline stage"),
|
||||
PipelineError::StepKindMismatch(s) => write!(f, "step kind mismatch: {}", s),
|
||||
PipelineError::EndOfStream => write!(f, "end of input stream"),
|
||||
PipelineError::StepError(e) => write!(f, "stage error: {}", e),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Error for PipelineError {
|
||||
fn source(&self) -> Option<&(dyn Error + 'static)> {
|
||||
match self {
|
||||
PipelineError::StepError(e) => Some(e.as_ref()),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ── Function types ────────────────────────────────────────────────────────────
|
||||
|
||||
/// Fonction source : appelée répétitivement, retourne le prochain item ou EndOfStream.
|
||||
/// `FnMut` car elle maintient un état interne (position dans l'itérateur).
|
||||
pub type SourceFn<D> = Box<dyn FnMut() -> Result<D, PipelineError> + Send>;
|
||||
|
||||
/// Fonction sink : consomme un item final, peut échouer (erreur d'I/O, etc.).
|
||||
pub type SinkFn<D> = Box<dyn Fn(D) -> Result<(), PipelineError> + Send>;
|
||||
|
||||
/// Fonction de transformation partagée entre workers via Arc.
|
||||
pub type SharedFn<D> = Arc<dyn Fn(D) -> Result<D, PipelineError> + Send + Sync>;
|
||||
|
||||
/// Fonction de transformation 1→N (flat map) partagée entre workers via Arc.
|
||||
///
|
||||
/// La fonction reçoit l'item d'entrée, un canal `push` pour envoyer chaque item
|
||||
/// produit, et un canal `delta` pour signaler au scheduler combien d'items
|
||||
/// supplémentaires sont entrés dans le pipeline (N-1 si N items produits).
|
||||
/// Elle doit appeler `delta.send(N - 1)` **après** avoir poussé tous les items.
|
||||
pub type SharedFlatFn<D> =
|
||||
Arc<dyn Fn(D, &Sender<Result<D, PipelineError>>, &Sender<isize>) + Send + Sync>;
|
||||
|
||||
// ── Stage enum ────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Une étape du pipeline : transform classique (1→1) ou flat transform (1→N).
|
||||
pub enum Stage<D> {
|
||||
Transform(SharedFn<D>),
|
||||
Flat(SharedFlatFn<D>),
|
||||
}
|
||||
|
||||
impl<D> Clone for Stage<D> {
|
||||
fn clone(&self) -> Self {
|
||||
match self {
|
||||
Stage::Transform(f) => Stage::Transform(Arc::clone(f)),
|
||||
Stage::Flat(f) => Stage::Flat(Arc::clone(f)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ── Worker task ───────────────────────────────────────────────────────────────
|
||||
|
||||
enum WorkerTask<D> {
|
||||
Transform(D, usize),
|
||||
Flat(D, usize),
|
||||
}
|
||||
|
||||
// ── Thread runners ────────────────────────────────────────────────────────────
|
||||
|
||||
fn source_runner<DATA>(
|
||||
mut source: SourceFn<DATA>,
|
||||
capacity: usize,
|
||||
) -> (
|
||||
Receiver<Result<DATA, PipelineError>>,
|
||||
thread::JoinHandle<()>,
|
||||
)
|
||||
where
|
||||
DATA: Send + Sync + 'static,
|
||||
{
|
||||
let (tx, rx) = bounded(capacity);
|
||||
let handle = thread::spawn(move || {
|
||||
loop {
|
||||
match source() {
|
||||
Ok(data) => {
|
||||
if tx.send(Ok(data)).is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Err(PipelineError::EndOfStream) => break,
|
||||
Err(e) => {
|
||||
eprintln!("Source error: {:?}", e);
|
||||
let _ = tx.send(Err(e));
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
(rx, handle)
|
||||
}
|
||||
|
||||
/// Lance un thread worker du pool.
|
||||
///
|
||||
/// Gère deux types de tâches :
|
||||
/// - `Transform` : applique `f(data)` et envoie le résultat dans `result_tx`.
|
||||
/// - `Flat` : appelle `f(data, &push_tx, &delta_tx)` ; la fonction elle-même
|
||||
/// pousse ses items dans `push_tx` et envoie `N-1` dans `delta_tx`.
|
||||
fn transform_runner<DATA>(
|
||||
task_rx: Receiver<WorkerTask<DATA>>,
|
||||
stages: Vec<Stage<DATA>>,
|
||||
stage_txs: Vec<Sender<Result<DATA, PipelineError>>>,
|
||||
flat_delta_tx: Sender<isize>,
|
||||
) -> thread::JoinHandle<()>
|
||||
where
|
||||
DATA: Send + Sync + 'static,
|
||||
{
|
||||
thread::spawn(move || {
|
||||
while let Ok(task) = task_rx.recv() {
|
||||
match task {
|
||||
WorkerTask::Transform(data, idx) => {
|
||||
if let Stage::Transform(f) = &stages[idx] {
|
||||
let _ = stage_txs[idx].send(f(data));
|
||||
}
|
||||
}
|
||||
WorkerTask::Flat(data, idx) => {
|
||||
if let Stage::Flat(f) = &stages[idx] {
|
||||
f(data, &stage_txs[idx], &flat_delta_tx);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/// Lance le thread sink.
|
||||
fn sink_runner<DATA>(
|
||||
sink: SinkFn<DATA>,
|
||||
capacity: usize,
|
||||
) -> (
|
||||
Sender<DATA>,
|
||||
Receiver<PipelineError>,
|
||||
thread::JoinHandle<()>,
|
||||
)
|
||||
where
|
||||
DATA: Send + Sync + 'static,
|
||||
{
|
||||
let (data_tx, data_rx) = bounded(capacity);
|
||||
let (err_tx, err_rx) = bounded(capacity);
|
||||
let handle = thread::spawn(move || {
|
||||
for data in data_rx {
|
||||
if let Err(e) = sink(data) {
|
||||
let _ = err_tx.send(e);
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
(data_tx, err_rx, handle)
|
||||
}
|
||||
|
||||
// ── Pipeline ──────────────────────────────────────────────────────────────────
|
||||
|
||||
pub struct Pipeline<DATA> {
|
||||
source: SourceFn<DATA>,
|
||||
stages: Vec<Stage<DATA>>,
|
||||
sink: SinkFn<DATA>,
|
||||
}
|
||||
|
||||
impl<DATA> Pipeline<DATA> {
|
||||
pub fn new(
|
||||
source: SourceFn<DATA>,
|
||||
stages: Vec<Stage<DATA>>,
|
||||
sink: SinkFn<DATA>,
|
||||
) -> Self {
|
||||
Self { source, stages, sink }
|
||||
}
|
||||
}
|
||||
|
||||
// ── WorkerPool ────────────────────────────────────────────────────────────────
|
||||
|
||||
pub struct WorkerPool<DATA> {
|
||||
pipeline: Pipeline<DATA>,
|
||||
handles: Vec<std::thread::JoinHandle<()>>,
|
||||
n_workers: usize,
|
||||
capacity: usize,
|
||||
}
|
||||
|
||||
impl<DATA> WorkerPool<DATA>
|
||||
where
|
||||
DATA: Send + Sync + 'static,
|
||||
{
|
||||
pub fn new(pipeline: Pipeline<DATA>, n_workers: usize, capacity: usize) -> Self {
|
||||
Self {
|
||||
pipeline,
|
||||
handles: Vec::new(),
|
||||
n_workers,
|
||||
capacity,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn run(mut self) {
|
||||
let n = self.pipeline.stages.len();
|
||||
|
||||
// ── Canaux inter-stages ────────────────────────────────────────────
|
||||
// stage_txs[i] / stage_rxs[i] : sortie du stage i
|
||||
let mut stage_txs: Vec<Sender<Result<DATA, PipelineError>>> = Vec::new();
|
||||
let mut stage_rxs: Vec<Receiver<Result<DATA, PipelineError>>> = Vec::new();
|
||||
for _ in 0..n {
|
||||
let (tx, rx) = bounded(self.capacity);
|
||||
stage_txs.push(tx);
|
||||
stage_rxs.push(rx);
|
||||
}
|
||||
|
||||
// ── Source thread ──────────────────────────────────────────────────
|
||||
let (source_rx, src_handle) = source_runner(self.pipeline.source, self.capacity);
|
||||
self.handles.push(src_handle);
|
||||
|
||||
let stages = self.pipeline.stages;
|
||||
|
||||
// ── Canal delta pour les flat stages ───────────────────────────────
|
||||
// Chaque flat worker envoie `N-1` ici après avoir poussé N items.
|
||||
// Le scheduler ajuste `in_flight` en conséquence.
|
||||
let (flat_delta_tx, flat_delta_rx) = bounded::<isize>(self.capacity);
|
||||
|
||||
// ── Worker pool ────────────────────────────────────────────────────
|
||||
let (worker_tx, worker_rx): (Sender<WorkerTask<DATA>>, Receiver<WorkerTask<DATA>>) =
|
||||
bounded(self.capacity);
|
||||
|
||||
for _ in 0..self.n_workers {
|
||||
self.handles.push(transform_runner(
|
||||
worker_rx.clone(),
|
||||
stages.iter().map(Stage::clone).collect(),
|
||||
stage_txs.clone(),
|
||||
flat_delta_tx.clone(),
|
||||
));
|
||||
}
|
||||
// Le scheduler ne tient plus flat_delta_tx : les workers le détiennent.
|
||||
// On le drop ici pour que le canal se ferme quand les workers terminent.
|
||||
drop(flat_delta_tx);
|
||||
|
||||
// ── Sink thread ────────────────────────────────────────────────────
|
||||
let (sink_tx, sink_err_rx, sink_handle) = sink_runner(self.pipeline.sink, self.capacity);
|
||||
self.handles.push(sink_handle);
|
||||
|
||||
// ── Boucle principale ──────────────────────────────────────────────
|
||||
//
|
||||
// `in_flight` (isize) = nb d'items qui doivent encore atteindre le sink.
|
||||
// Peut temporairement être négatif si un flat worker a poussé ses items
|
||||
// avant que le scheduler ait reçu le delta correspondant.
|
||||
//
|
||||
// `flat_workers_active` = nb de flat workers en cours d'exécution.
|
||||
// Empêche la terminaison prématurée quand in_flight vaut 0 mais qu'un
|
||||
// flat worker n'a pas encore envoyé son delta.
|
||||
//
|
||||
// Priorités du Select biaisé (index le plus bas = priorité la plus haute) :
|
||||
// 0 → sink_err_rx (arrêt immédiat sur erreur sink)
|
||||
// 1 → flat_delta_rx (mettre à jour in_flight avant de dispatcher)
|
||||
// 2..=n+1 → stage_rxs[n-1..0] (vider le pipeline en priorité)
|
||||
// n+2 → source_rx (dernier recours : nouvelles données)
|
||||
//
|
||||
// Quand k = 0 : erreur du sink
|
||||
// Quand k = 1 : delta d'un flat worker
|
||||
// Quand 2 ≤ k ≤ n+1 : résultat du stage n+1-k
|
||||
// Quand k = n+2 : item source
|
||||
//
|
||||
// Terminaison : source tarie ET in_flight == 0 ET aucun flat worker actif.
|
||||
{
|
||||
let mut source_done = false;
|
||||
let mut in_flight: isize = 0;
|
||||
let mut flat_workers_active: usize = 0;
|
||||
|
||||
loop {
|
||||
if source_done && in_flight == 0 && flat_workers_active == 0 {
|
||||
break;
|
||||
}
|
||||
|
||||
let mut sel = Select::new_biased();
|
||||
sel.recv(&sink_err_rx); // index 0
|
||||
sel.recv(&flat_delta_rx); // index 1
|
||||
for rx in stage_rxs.iter().rev() {
|
||||
sel.recv(rx); // indices 2..=n+1
|
||||
}
|
||||
let src_idx = if !source_done {
|
||||
Some(sel.recv(&source_rx)) // index n+2
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
let oper = sel.select();
|
||||
let k = oper.index();
|
||||
|
||||
if k == 0 {
|
||||
// ── Erreur du sink ────────────────────────────────────
|
||||
match oper.recv(&sink_err_rx) {
|
||||
Ok(e) => { eprintln!("Sink error: {:?}", e); break; }
|
||||
Err(_) => break,
|
||||
}
|
||||
} else if k == 1 {
|
||||
// ── Delta d'un flat worker ────────────────────────────
|
||||
// delta = N - 1 (N items poussés, 1 item consommé)
|
||||
match oper.recv(&flat_delta_rx) {
|
||||
Ok(delta) => {
|
||||
in_flight += delta;
|
||||
flat_workers_active -= 1;
|
||||
}
|
||||
Err(_) => {}
|
||||
}
|
||||
} else if src_idx == Some(k) {
|
||||
// ── Nouvel item depuis la source ──────────────────────
|
||||
match oper.recv(&source_rx) {
|
||||
Ok(Ok(data)) => {
|
||||
if n == 0 {
|
||||
let _ = sink_tx.send(data);
|
||||
} else {
|
||||
in_flight += 1;
|
||||
dispatch(
|
||||
data, 0,
|
||||
&stages, &worker_tx,
|
||||
&mut flat_workers_active,
|
||||
);
|
||||
}
|
||||
}
|
||||
Ok(Err(e)) => eprintln!("Source error: {:?}", e),
|
||||
Err(_) => source_done = true,
|
||||
}
|
||||
} else {
|
||||
// ── Résultat d'un stage intermédiaire ─────────────────
|
||||
// k ∈ [2, n+1] → stage = n+1 - k
|
||||
let stage = n + 1 - k;
|
||||
match oper.recv(&stage_rxs[stage]) {
|
||||
Ok(Ok(data)) => {
|
||||
if stage == n - 1 {
|
||||
in_flight -= 1;
|
||||
let _ = sink_tx.send(data);
|
||||
} else {
|
||||
dispatch(
|
||||
data, stage + 1,
|
||||
&stages, &worker_tx,
|
||||
&mut flat_workers_active,
|
||||
);
|
||||
}
|
||||
}
|
||||
Ok(Err(e)) => eprintln!("Stage {} error: {:?}", stage, e),
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
drop(worker_tx);
|
||||
drop(sink_tx);
|
||||
|
||||
for h in self.handles {
|
||||
let _ = h.join();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ── Pipe ──────────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Typed, composable iterator transformer.
|
||||
///
|
||||
/// A `Pipe<D, In, Out>` is a pure description of pipeline stages — no threads,
|
||||
/// no channels, no scheduler. Call `.apply(iter, n_workers, capacity)` to start
|
||||
/// execution and get back a `PipeIter<Out>`.
|
||||
///
|
||||
/// Compose two pipes with `.then()`: the resulting `Pipe` holds the concatenated
|
||||
/// stage list, so a single scheduler is created when `.apply()` is eventually called.
|
||||
pub struct Pipe<D, In, Out> {
|
||||
stages: Vec<Stage<D>>,
|
||||
wrap: Arc<dyn Fn(In) -> D + Send + Sync>,
|
||||
unwrap: Arc<dyn Fn(D) -> Out + Send + Sync>,
|
||||
_phantom: PhantomData<(In, Out)>,
|
||||
}
|
||||
|
||||
impl<D, In, Out> Pipe<D, In, Out> {
|
||||
/// Build a `Pipe` from stages and wrap/unwrap converters.
|
||||
/// Prefer the `make_pipe!` macro.
|
||||
pub fn new(
|
||||
stages: Vec<Stage<D>>,
|
||||
wrap: Arc<dyn Fn(In) -> D + Send + Sync>,
|
||||
unwrap: Arc<dyn Fn(D) -> Out + Send + Sync>,
|
||||
) -> Self {
|
||||
Self { stages, wrap, unwrap, _phantom: PhantomData }
|
||||
}
|
||||
|
||||
/// Concatenate stages from two pipes into one.
|
||||
///
|
||||
/// Requires `Out` of `self` == `In` of `other`. The single scheduler
|
||||
/// created at `.apply()` time sees the full combined stage list.
|
||||
pub fn then<Next>(self, other: Pipe<D, Out, Next>) -> Pipe<D, In, Next> {
|
||||
Pipe {
|
||||
stages: self.stages.into_iter().chain(other.stages).collect(),
|
||||
wrap: self.wrap,
|
||||
unwrap: other.unwrap,
|
||||
_phantom: PhantomData,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<D, In, Out> Pipe<D, In, Out>
|
||||
where
|
||||
D: Send + Sync + 'static,
|
||||
In: Send + 'static,
|
||||
Out: Send + 'static,
|
||||
{
|
||||
/// Run the pipeline in a background thread; returns an iterator over the output.
|
||||
pub fn apply(
|
||||
self,
|
||||
input: impl Iterator<Item = In> + Send + 'static,
|
||||
n_workers: usize,
|
||||
capacity: usize,
|
||||
) -> PipeIter<Out> {
|
||||
let wrap = Arc::clone(&self.wrap);
|
||||
let unwrap = Arc::clone(&self.unwrap);
|
||||
|
||||
let mut iter = input;
|
||||
let source: SourceFn<D> = Box::new(move || match iter.next() {
|
||||
Some(x) => Ok(wrap(x)),
|
||||
None => Err(PipelineError::EndOfStream),
|
||||
});
|
||||
|
||||
let (out_tx, out_rx) = bounded::<Out>(capacity);
|
||||
let sink: SinkFn<D> = Box::new(move |data: D| {
|
||||
out_tx.send(unwrap(data)).map_err(|_| {
|
||||
PipelineError::StepError(Box::new(std::io::Error::new(
|
||||
std::io::ErrorKind::BrokenPipe,
|
||||
"output channel closed",
|
||||
)))
|
||||
})
|
||||
});
|
||||
|
||||
let pipeline = Pipeline::new(source, self.stages, sink);
|
||||
let handle = thread::spawn(move || {
|
||||
WorkerPool::new(pipeline, n_workers, capacity).run();
|
||||
});
|
||||
|
||||
PipeIter { rx: out_rx, handle: Some(handle) }
|
||||
}
|
||||
}
|
||||
|
||||
// ── PipeIter ──────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Iterator over the output of `Pipe::apply()`.
|
||||
pub struct PipeIter<Out> {
|
||||
rx: Receiver<Out>,
|
||||
handle: Option<thread::JoinHandle<()>>,
|
||||
}
|
||||
|
||||
impl<Out> Iterator for PipeIter<Out> {
|
||||
type Item = Out;
|
||||
|
||||
fn next(&mut self) -> Option<Out> {
|
||||
self.rx.recv().ok()
|
||||
}
|
||||
}
|
||||
|
||||
impl<Out> Drop for PipeIter<Out> {
|
||||
fn drop(&mut self) {
|
||||
// Drain buffered items so the scheduler can unblock if the channel is full.
|
||||
while self.rx.try_recv().is_ok() {}
|
||||
if let Some(h) = self.handle.take() {
|
||||
let _ = h.join();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Envoie `data` au stage `stage_idx`.
|
||||
/// Pour un `Transform`, empile une `WorkerTask::Transform`.
|
||||
/// Pour un `Flat`, incrémente `flat_workers_active` et empile une `WorkerTask::Flat`.
|
||||
#[inline]
|
||||
fn dispatch<DATA>(
|
||||
data: DATA,
|
||||
stage_idx: usize,
|
||||
stages: &[Stage<DATA>],
|
||||
worker_tx: &Sender<WorkerTask<DATA>>,
|
||||
flat_workers_active: &mut usize,
|
||||
) {
|
||||
match &stages[stage_idx] {
|
||||
Stage::Transform(_) => {
|
||||
let _ = worker_tx.send(WorkerTask::Transform(data, stage_idx));
|
||||
}
|
||||
Stage::Flat(_) => {
|
||||
*flat_workers_active += 1;
|
||||
let _ = worker_tx.send(WorkerTask::Flat(data, stage_idx));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ── Macros ────────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Creates a `SourceFn` from an iterator of plain values.
|
||||
#[macro_export]
|
||||
macro_rules! make_source {
|
||||
($enum:ident, $iterator:expr, $output:ident) => {{
|
||||
let mut iter = $iterator.into_iter();
|
||||
Box::new(
|
||||
move || -> ::std::result::Result<$enum, $crate::PipelineError> {
|
||||
match iter.next() {
|
||||
Some(x) => Ok($enum::$output(x)),
|
||||
None => Err($crate::PipelineError::EndOfStream),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<dyn FnMut() -> ::std::result::Result<$enum, $crate::PipelineError> + Send>
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `SourceFn` from an iterator of `Result<T, E>`.
|
||||
#[macro_export]
|
||||
macro_rules! make_source_fallible {
|
||||
($enum:ident, $iterator:expr, $output:ident) => {{
|
||||
let mut iter = $iterator.into_iter();
|
||||
Box::new(
|
||||
move || -> ::std::result::Result<$enum, $crate::PipelineError> {
|
||||
match iter.next() {
|
||||
Some(Ok(x)) => Ok($enum::$output(x)),
|
||||
Some(Err(e)) => Err($crate::PipelineError::StepError(Box::new(e))),
|
||||
None => Err($crate::PipelineError::EndOfStream),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<dyn FnMut() -> ::std::result::Result<$enum, $crate::PipelineError> + Send>
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `Stage::Transform` from a pure (non-fallible) function `Fn(T) -> U`.
|
||||
#[macro_export]
|
||||
macro_rules! make_transform {
|
||||
($enum:ident, $func:tt, $input:ident, $output:ident) => {{
|
||||
let __f = $func;
|
||||
$crate::Stage::Transform(
|
||||
::std::sync::Arc::from(Box::new(
|
||||
move |data: $enum| -> ::std::result::Result<$enum, $crate::PipelineError> {
|
||||
match data {
|
||||
$enum::$input(x) => Ok($enum::$output(__f(x))),
|
||||
_ => Err($crate::PipelineError::TypeMismatch),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<
|
||||
dyn Fn($enum) -> ::std::result::Result<$enum, $crate::PipelineError>
|
||||
+ Send
|
||||
+ Sync,
|
||||
>)
|
||||
)
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `Stage::Transform` from a fallible function `Fn(T) -> Result<U, E>`.
|
||||
#[macro_export]
|
||||
macro_rules! make_transform_fallible {
|
||||
($enum:ident, $func:tt, $input:ident, $output:ident) => {{
|
||||
let __f = $func;
|
||||
$crate::Stage::Transform(
|
||||
::std::sync::Arc::from(Box::new(
|
||||
move |data: $enum| -> ::std::result::Result<$enum, $crate::PipelineError> {
|
||||
match data {
|
||||
$enum::$input(inner) => {
|
||||
let result = __f(inner)
|
||||
.map_err(|e| $crate::PipelineError::StepError(Box::new(e)))?;
|
||||
Ok($enum::$output(result))
|
||||
}
|
||||
_ => Err($crate::PipelineError::TypeMismatch),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<
|
||||
dyn Fn($enum) -> ::std::result::Result<$enum, $crate::PipelineError>
|
||||
+ Send
|
||||
+ Sync,
|
||||
>)
|
||||
)
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `Stage::Flat` from a function `Fn(T) -> impl IntoIterator<Item = U>`.
|
||||
///
|
||||
/// Pour chaque item produit par l'itérateur, il est poussé individuellement dans
|
||||
/// le canal de sortie, permettant au scheduler de dispatcher les items en parallèle
|
||||
/// dès qu'un worker est disponible.
|
||||
#[macro_export]
|
||||
macro_rules! make_flat_transform {
|
||||
($enum:ident, $func:tt, $input:ident, $output:ident) => {{
|
||||
let __f = $func;
|
||||
$crate::Stage::Flat(
|
||||
::std::sync::Arc::new(
|
||||
move |data: $enum,
|
||||
push: &$crate::PipelineSender<
|
||||
::std::result::Result<$enum, $crate::PipelineError>,
|
||||
>,
|
||||
delta: &$crate::PipelineSender<isize>| {
|
||||
match data {
|
||||
$enum::$input(inner) => {
|
||||
let mut count: isize = 0;
|
||||
for item in __f(inner) {
|
||||
push.send(Ok($enum::$output(item))).ok();
|
||||
count += 1;
|
||||
}
|
||||
delta.send(count - 1).ok();
|
||||
}
|
||||
_ => {
|
||||
push.send(Err($crate::PipelineError::TypeMismatch)).ok();
|
||||
delta.send(0).ok();
|
||||
}
|
||||
}
|
||||
},
|
||||
) as $crate::SharedFlatFn<$enum>
|
||||
)
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `Stage::Flat` from a fallible function
|
||||
/// `Fn(T) -> Result<impl IntoIterator<Item = U>, E>`.
|
||||
///
|
||||
/// Si la fonction retourne `Err`, une erreur est poussée dans le canal et aucun
|
||||
/// item normal n'est produit.
|
||||
#[macro_export]
|
||||
macro_rules! make_flat_transform_fallible {
|
||||
($enum:ident, $func:tt, $input:ident, $output:ident) => {{
|
||||
let __f = $func;
|
||||
$crate::Stage::Flat(
|
||||
::std::sync::Arc::new(
|
||||
move |data: $enum,
|
||||
push: &$crate::PipelineSender<
|
||||
::std::result::Result<$enum, $crate::PipelineError>,
|
||||
>,
|
||||
delta: &$crate::PipelineSender<isize>| {
|
||||
match data {
|
||||
$enum::$input(inner) => match __f(inner) {
|
||||
Ok(iter) => {
|
||||
let mut count: isize = 0;
|
||||
for item in iter {
|
||||
push.send(Ok($enum::$output(item))).ok();
|
||||
count += 1;
|
||||
}
|
||||
delta.send(count - 1).ok();
|
||||
}
|
||||
Err(e) => {
|
||||
push.send(Err($crate::PipelineError::StepError(Box::new(e))))
|
||||
.ok();
|
||||
delta.send(0).ok();
|
||||
}
|
||||
},
|
||||
_ => {
|
||||
push.send(Err($crate::PipelineError::TypeMismatch)).ok();
|
||||
delta.send(0).ok();
|
||||
}
|
||||
}
|
||||
},
|
||||
) as $crate::SharedFlatFn<$enum>
|
||||
)
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `SinkFn` from a function that consumes a concrete value and returns `()`.
|
||||
#[macro_export]
|
||||
macro_rules! make_sink {
|
||||
($enum:ident, $func:tt, $input:ident) => {{
|
||||
let __f = $func;
|
||||
Box::new(
|
||||
move |data: $enum| -> ::std::result::Result<(), $crate::PipelineError> {
|
||||
match data {
|
||||
$enum::$input(x) => {
|
||||
__f(x);
|
||||
Ok(())
|
||||
}
|
||||
_ => Err($crate::PipelineError::TypeMismatch),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<dyn Fn($enum) -> ::std::result::Result<(), $crate::PipelineError> + Send>
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `SinkFn` from a fallible function that returns `Result<(), E>`.
|
||||
#[macro_export]
|
||||
macro_rules! make_sink_fallible {
|
||||
($enum:ident, $func:tt, $input:ident) => {{
|
||||
let __f = $func;
|
||||
Box::new(
|
||||
move |data: $enum| -> ::std::result::Result<(), $crate::PipelineError> {
|
||||
match data {
|
||||
$enum::$input(inner) => {
|
||||
__f(inner).map_err(|e| $crate::PipelineError::StepError(Box::new(e)))
|
||||
}
|
||||
_ => Err($crate::PipelineError::TypeMismatch),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<dyn Fn($enum) -> ::std::result::Result<(), $crate::PipelineError> + Send>
|
||||
}};
|
||||
}
|
||||
|
||||
/// Construit un `Pipeline` à partir d'une source, d'une liste de stages et d'un sink.
|
||||
///
|
||||
/// Syntaxe :
|
||||
/// ```ignore
|
||||
/// make_pipeline! {
|
||||
/// MyData,
|
||||
/// source my_iter => Variant, // source non-fallible
|
||||
/// source? my_iter => Variant, // source fallible (Result<T, E>)
|
||||
/// | func: In => Out, // transform 1→1 non-fallible
|
||||
/// |? func: In => Out, // transform 1→1 fallible
|
||||
/// || func: In => Out, // flat transform 1→N non-fallible
|
||||
/// ||? func: In => Out, // flat transform 1→N fallible
|
||||
/// sink my_func @ Variant, // sink non-fallible
|
||||
/// sink? my_func @ Variant, // sink fallible
|
||||
/// }
|
||||
/// ```
|
||||
#[macro_export]
|
||||
macro_rules! make_pipeline {
|
||||
// ── Points d'entrée ──────────────────────────────────────────────────
|
||||
|
||||
($enum:ident, source $src:expr => $src_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum,
|
||||
{ $crate::make_source!($enum, $src, $src_out) },
|
||||
[],
|
||||
$($rest)*)
|
||||
};
|
||||
($enum:ident, source? $src:expr => $src_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum,
|
||||
{ $crate::make_source_fallible!($enum, $src, $src_out) },
|
||||
[],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// ── Accumulation des stages ──────────────────────────────────────────
|
||||
|
||||
// transform 1→1 non-fallible
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
| $tf:tt : $t_in:ident => $t_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum, $source,
|
||||
[$($acc)* $crate::make_transform!($enum, $tf, $t_in, $t_out),],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// transform 1→1 fallible
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
|? $tf:tt : $t_in:ident => $t_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum, $source,
|
||||
[$($acc)* $crate::make_transform_fallible!($enum, $tf, $t_in, $t_out),],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// flat transform 1→N non-fallible
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
|| $tf:tt : $t_in:ident => $t_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum, $source,
|
||||
[$($acc)* $crate::make_flat_transform!($enum, $tf, $t_in, $t_out),],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// flat transform 1→N fallible
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
||? $tf:tt : $t_in:ident => $t_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum, $source,
|
||||
[$($acc)* $crate::make_flat_transform_fallible!($enum, $tf, $t_in, $t_out),],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// ── Terminaison : sink ───────────────────────────────────────────────
|
||||
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
sink $sink_fn:tt @ $sink_in:ident $(,)?) => {
|
||||
$crate::Pipeline::new(
|
||||
$source,
|
||||
vec![$($acc)*],
|
||||
$crate::make_sink!($enum, $sink_fn, $sink_in),
|
||||
)
|
||||
};
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
sink? $sink_fn:tt @ $sink_in:ident $(,)?) => {
|
||||
$crate::Pipeline::new(
|
||||
$source,
|
||||
vec![$($acc)*],
|
||||
$crate::make_sink_fallible!($enum, $sink_fn, $sink_in),
|
||||
)
|
||||
};
|
||||
}
|
||||
|
||||
/// Builds a typed `Pipe<D, In, Out>` — sourceless and sinkless.
|
||||
///
|
||||
/// Syntax:
|
||||
/// ```ignore
|
||||
/// make_pipe! {
|
||||
/// MyData : InType => OutType,
|
||||
/// | func : InVariant => OutVariant, // transform 1→1
|
||||
/// |? func : InVariant => OutVariant, // transform 1→1 fallible
|
||||
/// || func : InVariant => OutVariant, // flat transform 1→N
|
||||
/// ||? func : InVariant => OutVariant, // flat transform 1→N fallible
|
||||
/// }
|
||||
/// ```
|
||||
#[macro_export]
|
||||
macro_rules! make_pipe {
|
||||
// ── Entry: first stage | ─────────────────────────────────────────────
|
||||
($enum:ident : $in_ty:ty => $out_ty:ty,
|
||||
| $tf:tt : $fi:ident => $fo:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$crate::make_transform!($enum, $tf, $fi, $fo),], $fo, $($rest)*)
|
||||
};
|
||||
// ── Entry: first stage |? ────────────────────────────────────────────
|
||||
($enum:ident : $in_ty:ty => $out_ty:ty,
|
||||
|? $tf:tt : $fi:ident => $fo:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$crate::make_transform_fallible!($enum, $tf, $fi, $fo),], $fo, $($rest)*)
|
||||
};
|
||||
// ── Entry: first stage || ────────────────────────────────────────────
|
||||
($enum:ident : $in_ty:ty => $out_ty:ty,
|
||||
|| $tf:tt : $fi:ident => $fo:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$crate::make_flat_transform!($enum, $tf, $fi, $fo),], $fo, $($rest)*)
|
||||
};
|
||||
// ── Entry: first stage ||? ───────────────────────────────────────────
|
||||
($enum:ident : $in_ty:ty => $out_ty:ty,
|
||||
||? $tf:tt : $fi:ident => $fo:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$crate::make_flat_transform_fallible!($enum, $tf, $fi, $fo),], $fo, $($rest)*)
|
||||
};
|
||||
|
||||
// ── Accumulation: | ──────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident,
|
||||
| $tf:tt : $ti:ident => $to:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$($acc)* $crate::make_transform!($enum, $tf, $ti, $to),], $to, $($rest)*)
|
||||
};
|
||||
// ── Accumulation: |? ─────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident,
|
||||
|? $tf:tt : $ti:ident => $to:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$($acc)* $crate::make_transform_fallible!($enum, $tf, $ti, $to),], $to, $($rest)*)
|
||||
};
|
||||
// ── Accumulation: || ─────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident,
|
||||
|| $tf:tt : $ti:ident => $to:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$($acc)* $crate::make_flat_transform!($enum, $tf, $ti, $to),], $to, $($rest)*)
|
||||
};
|
||||
// ── Accumulation: ||? ────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident,
|
||||
||? $tf:tt : $ti:ident => $to:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$($acc)* $crate::make_flat_transform_fallible!($enum, $tf, $ti, $to),], $to, $($rest)*)
|
||||
};
|
||||
|
||||
// ── Termination ───────────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident $(,)?) => {
|
||||
$crate::Pipe::new(
|
||||
vec![$($acc)*],
|
||||
::std::sync::Arc::new(|x: $in_ty| $enum::$fi(x)),
|
||||
::std::sync::Arc::new(|d: $enum| -> $out_ty {
|
||||
if let $enum::$lo(x) = d { x }
|
||||
else { ::std::unreachable!("unexpected pipeline data variant in make_pipe!") }
|
||||
}),
|
||||
)
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
use std::error::Error;
|
||||
use std::fmt;
|
||||
|
||||
/// Error type for pipeline operations.
|
||||
#[derive(Debug)]
|
||||
pub enum PipelineError {
|
||||
/// A stage received a `PipelineData` variant it did not expect.
|
||||
TypeMismatch,
|
||||
/// The step kind is not compatible with the data type.
|
||||
StepKindMismatch(&'static str),
|
||||
/// The source has no more data to produce.
|
||||
EndOfStream,
|
||||
/// An error occurred inside a stage (e.g., I/O, parsing, custom logic).
|
||||
StepError(Box<dyn Error + Send + Sync>),
|
||||
}
|
||||
|
||||
impl fmt::Display for PipelineError {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
match self {
|
||||
PipelineError::TypeMismatch => write!(f, "data type mismatch in pipeline stage"),
|
||||
PipelineError::StepKindMismatch(s) => write!(f, "step kind mismatch: {}", s),
|
||||
PipelineError::EndOfStream => write!(f, "end of input stream"),
|
||||
PipelineError::StepError(e) => write!(f, "stage error: {}", e),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Error for PipelineError {
|
||||
fn source(&self) -> Option<&(dyn Error + 'static)> {
|
||||
match self {
|
||||
PipelineError::StepError(e) => Some(e.as_ref()),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,373 @@
|
||||
// ── Macros ────────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Creates a `SourceFn` from an iterator of plain values.
|
||||
#[macro_export]
|
||||
macro_rules! make_source {
|
||||
($enum:ident, $iterator:expr, $output:ident) => {{
|
||||
let mut iter = $iterator.into_iter();
|
||||
Box::new(
|
||||
move || -> ::std::result::Result<$enum, $crate::PipelineError> {
|
||||
match iter.next() {
|
||||
Some(x) => Ok($enum::$output(x)),
|
||||
None => Err($crate::PipelineError::EndOfStream),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<dyn FnMut() -> ::std::result::Result<$enum, $crate::PipelineError> + Send>
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `SourceFn` from an iterator of `Result<T, E>`.
|
||||
#[macro_export]
|
||||
macro_rules! make_source_fallible {
|
||||
($enum:ident, $iterator:expr, $output:ident) => {{
|
||||
let mut iter = $iterator.into_iter();
|
||||
Box::new(
|
||||
move || -> ::std::result::Result<$enum, $crate::PipelineError> {
|
||||
match iter.next() {
|
||||
Some(Ok(x)) => Ok($enum::$output(x)),
|
||||
Some(Err(e)) => Err($crate::PipelineError::StepError(Box::new(e))),
|
||||
None => Err($crate::PipelineError::EndOfStream),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<dyn FnMut() -> ::std::result::Result<$enum, $crate::PipelineError> + Send>
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `Stage::Transform` from a pure (non-fallible) function `Fn(T) -> U`.
|
||||
#[macro_export]
|
||||
macro_rules! make_transform {
|
||||
($enum:ident, $func:tt, $input:ident, $output:ident) => {{
|
||||
let __f = $func;
|
||||
$crate::Stage::Transform(
|
||||
::std::sync::Arc::from(Box::new(
|
||||
move |data: $enum| -> ::std::result::Result<$enum, $crate::PipelineError> {
|
||||
match data {
|
||||
$enum::$input(x) => Ok($enum::$output(__f(x))),
|
||||
_ => Err($crate::PipelineError::TypeMismatch),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<
|
||||
dyn Fn($enum) -> ::std::result::Result<$enum, $crate::PipelineError>
|
||||
+ Send
|
||||
+ Sync,
|
||||
>)
|
||||
)
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `Stage::Transform` from a fallible function `Fn(T) -> Result<U, E>`.
|
||||
#[macro_export]
|
||||
macro_rules! make_transform_fallible {
|
||||
($enum:ident, $func:tt, $input:ident, $output:ident) => {{
|
||||
let __f = $func;
|
||||
$crate::Stage::Transform(
|
||||
::std::sync::Arc::from(Box::new(
|
||||
move |data: $enum| -> ::std::result::Result<$enum, $crate::PipelineError> {
|
||||
match data {
|
||||
$enum::$input(inner) => {
|
||||
let result = __f(inner)
|
||||
.map_err(|e| $crate::PipelineError::StepError(Box::new(e)))?;
|
||||
Ok($enum::$output(result))
|
||||
}
|
||||
_ => Err($crate::PipelineError::TypeMismatch),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<
|
||||
dyn Fn($enum) -> ::std::result::Result<$enum, $crate::PipelineError>
|
||||
+ Send
|
||||
+ Sync,
|
||||
>)
|
||||
)
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `Stage::Flat` from a function `Fn(T) -> impl IntoIterator<Item = U>`.
|
||||
///
|
||||
/// Pour chaque item produit par l'itérateur, il est poussé individuellement dans
|
||||
/// le canal de sortie, permettant au scheduler de dispatcher les items en parallèle
|
||||
/// dès qu'un worker est disponible.
|
||||
#[macro_export]
|
||||
macro_rules! make_flat_transform {
|
||||
($enum:ident, $func:tt, $input:ident, $output:ident) => {{
|
||||
let __f = $func;
|
||||
$crate::Stage::Flat(
|
||||
::std::sync::Arc::new(
|
||||
move |data: $enum,
|
||||
push: &$crate::PipelineSender<
|
||||
::std::result::Result<$enum, $crate::PipelineError>,
|
||||
>,
|
||||
delta: &$crate::PipelineSender<isize>| {
|
||||
match data {
|
||||
$enum::$input(inner) => {
|
||||
let mut count: isize = 0;
|
||||
for item in __f(inner) {
|
||||
push.send(Ok($enum::$output(item))).ok();
|
||||
count += 1;
|
||||
}
|
||||
delta.send(count - 1).ok();
|
||||
}
|
||||
_ => {
|
||||
push.send(Err($crate::PipelineError::TypeMismatch)).ok();
|
||||
delta.send(0).ok();
|
||||
}
|
||||
}
|
||||
},
|
||||
) as $crate::SharedFlatFn<$enum>
|
||||
)
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `Stage::Flat` from a fallible function
|
||||
/// `Fn(T) -> Result<impl IntoIterator<Item = U>, E>`.
|
||||
///
|
||||
/// Si la fonction retourne `Err`, une erreur est poussée dans le canal et aucun
|
||||
/// item normal n'est produit.
|
||||
#[macro_export]
|
||||
macro_rules! make_flat_transform_fallible {
|
||||
($enum:ident, $func:tt, $input:ident, $output:ident) => {{
|
||||
let __f = $func;
|
||||
$crate::Stage::Flat(
|
||||
::std::sync::Arc::new(
|
||||
move |data: $enum,
|
||||
push: &$crate::PipelineSender<
|
||||
::std::result::Result<$enum, $crate::PipelineError>,
|
||||
>,
|
||||
delta: &$crate::PipelineSender<isize>| {
|
||||
match data {
|
||||
$enum::$input(inner) => match __f(inner) {
|
||||
Ok(iter) => {
|
||||
let mut count: isize = 0;
|
||||
for item in iter {
|
||||
push.send(Ok($enum::$output(item))).ok();
|
||||
count += 1;
|
||||
}
|
||||
delta.send(count - 1).ok();
|
||||
}
|
||||
Err(e) => {
|
||||
push.send(Err($crate::PipelineError::StepError(Box::new(e))))
|
||||
.ok();
|
||||
delta.send(0).ok();
|
||||
}
|
||||
},
|
||||
_ => {
|
||||
push.send(Err($crate::PipelineError::TypeMismatch)).ok();
|
||||
delta.send(0).ok();
|
||||
}
|
||||
}
|
||||
},
|
||||
) as $crate::SharedFlatFn<$enum>
|
||||
)
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `SinkFn` from a function that consumes a concrete value and returns `()`.
|
||||
#[macro_export]
|
||||
macro_rules! make_sink {
|
||||
($enum:ident, $func:tt, $input:ident) => {{
|
||||
let __f = $func;
|
||||
Box::new(
|
||||
move |data: $enum| -> ::std::result::Result<(), $crate::PipelineError> {
|
||||
match data {
|
||||
$enum::$input(x) => {
|
||||
__f(x);
|
||||
Ok(())
|
||||
}
|
||||
_ => Err($crate::PipelineError::TypeMismatch),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<dyn Fn($enum) -> ::std::result::Result<(), $crate::PipelineError> + Send>
|
||||
}};
|
||||
}
|
||||
|
||||
/// Creates a `SinkFn` from a fallible function that returns `Result<(), E>`.
|
||||
#[macro_export]
|
||||
macro_rules! make_sink_fallible {
|
||||
($enum:ident, $func:tt, $input:ident) => {{
|
||||
let __f = $func;
|
||||
Box::new(
|
||||
move |data: $enum| -> ::std::result::Result<(), $crate::PipelineError> {
|
||||
match data {
|
||||
$enum::$input(inner) => {
|
||||
__f(inner).map_err(|e| $crate::PipelineError::StepError(Box::new(e)))
|
||||
}
|
||||
_ => Err($crate::PipelineError::TypeMismatch),
|
||||
}
|
||||
},
|
||||
)
|
||||
as Box<dyn Fn($enum) -> ::std::result::Result<(), $crate::PipelineError> + Send>
|
||||
}};
|
||||
}
|
||||
|
||||
/// Construit un `Pipeline` à partir d'une source, d'une liste de stages et d'un sink.
|
||||
///
|
||||
/// Syntaxe :
|
||||
/// ```ignore
|
||||
/// make_pipeline! {
|
||||
/// MyData,
|
||||
/// source my_iter => Variant, // source non-fallible
|
||||
/// source? my_iter => Variant, // source fallible (Result<T, E>)
|
||||
/// | func: In => Out, // transform 1→1 non-fallible
|
||||
/// |? func: In => Out, // transform 1→1 fallible
|
||||
/// || func: In => Out, // flat transform 1→N non-fallible
|
||||
/// ||? func: In => Out, // flat transform 1→N fallible
|
||||
/// sink my_func @ Variant, // sink non-fallible
|
||||
/// sink? my_func @ Variant, // sink fallible
|
||||
/// }
|
||||
/// ```
|
||||
#[macro_export]
|
||||
macro_rules! make_pipeline {
|
||||
// ── Points d'entrée ──────────────────────────────────────────────────
|
||||
|
||||
($enum:ident, source $src:expr => $src_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum,
|
||||
{ $crate::make_source!($enum, $src, $src_out) },
|
||||
[],
|
||||
$($rest)*)
|
||||
};
|
||||
($enum:ident, source? $src:expr => $src_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum,
|
||||
{ $crate::make_source_fallible!($enum, $src, $src_out) },
|
||||
[],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// ── Accumulation des stages ──────────────────────────────────────────
|
||||
|
||||
// transform 1→1 non-fallible
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
| $tf:tt : $t_in:ident => $t_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum, $source,
|
||||
[$($acc)* $crate::make_transform!($enum, $tf, $t_in, $t_out),],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// transform 1→1 fallible
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
|? $tf:tt : $t_in:ident => $t_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum, $source,
|
||||
[$($acc)* $crate::make_transform_fallible!($enum, $tf, $t_in, $t_out),],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// flat transform 1→N non-fallible
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
|| $tf:tt : $t_in:ident => $t_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum, $source,
|
||||
[$($acc)* $crate::make_flat_transform!($enum, $tf, $t_in, $t_out),],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// flat transform 1→N fallible
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
||? $tf:tt : $t_in:ident => $t_out:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipeline!(@build $enum, $source,
|
||||
[$($acc)* $crate::make_flat_transform_fallible!($enum, $tf, $t_in, $t_out),],
|
||||
$($rest)*)
|
||||
};
|
||||
|
||||
// ── Terminaison : sink ───────────────────────────────────────────────
|
||||
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
sink $sink_fn:tt @ $sink_in:ident $(,)?) => {
|
||||
$crate::Pipeline::new(
|
||||
$source,
|
||||
vec![$($acc)*],
|
||||
$crate::make_sink!($enum, $sink_fn, $sink_in),
|
||||
)
|
||||
};
|
||||
(@build $enum:ident, $source:tt, [$($acc:tt)*],
|
||||
sink? $sink_fn:tt @ $sink_in:ident $(,)?) => {
|
||||
$crate::Pipeline::new(
|
||||
$source,
|
||||
vec![$($acc)*],
|
||||
$crate::make_sink_fallible!($enum, $sink_fn, $sink_in),
|
||||
)
|
||||
};
|
||||
}
|
||||
|
||||
/// Builds a typed `Pipe<D, In, Out>` — sourceless and sinkless.
|
||||
///
|
||||
/// Syntax:
|
||||
/// ```ignore
|
||||
/// make_pipe! {
|
||||
/// MyData : InType => OutType,
|
||||
/// | func : InVariant => OutVariant, // transform 1→1
|
||||
/// |? func : InVariant => OutVariant, // transform 1→1 fallible
|
||||
/// || func : InVariant => OutVariant, // flat transform 1→N
|
||||
/// ||? func : InVariant => OutVariant, // flat transform 1→N fallible
|
||||
/// }
|
||||
/// ```
|
||||
#[macro_export]
|
||||
macro_rules! make_pipe {
|
||||
// ── Entry: first stage | ─────────────────────────────────────────────
|
||||
($enum:ident : $in_ty:ty => $out_ty:ty,
|
||||
| $tf:tt : $fi:ident => $fo:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$crate::make_transform!($enum, $tf, $fi, $fo),], $fo, $($rest)*)
|
||||
};
|
||||
// ── Entry: first stage |? ────────────────────────────────────────────
|
||||
($enum:ident : $in_ty:ty => $out_ty:ty,
|
||||
|? $tf:tt : $fi:ident => $fo:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$crate::make_transform_fallible!($enum, $tf, $fi, $fo),], $fo, $($rest)*)
|
||||
};
|
||||
// ── Entry: first stage || ────────────────────────────────────────────
|
||||
($enum:ident : $in_ty:ty => $out_ty:ty,
|
||||
|| $tf:tt : $fi:ident => $fo:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$crate::make_flat_transform!($enum, $tf, $fi, $fo),], $fo, $($rest)*)
|
||||
};
|
||||
// ── Entry: first stage ||? ───────────────────────────────────────────
|
||||
($enum:ident : $in_ty:ty => $out_ty:ty,
|
||||
||? $tf:tt : $fi:ident => $fo:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$crate::make_flat_transform_fallible!($enum, $tf, $fi, $fo),], $fo, $($rest)*)
|
||||
};
|
||||
|
||||
// ── Accumulation: | ──────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident,
|
||||
| $tf:tt : $ti:ident => $to:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$($acc)* $crate::make_transform!($enum, $tf, $ti, $to),], $to, $($rest)*)
|
||||
};
|
||||
// ── Accumulation: |? ─────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident,
|
||||
|? $tf:tt : $ti:ident => $to:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$($acc)* $crate::make_transform_fallible!($enum, $tf, $ti, $to),], $to, $($rest)*)
|
||||
};
|
||||
// ── Accumulation: || ─────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident,
|
||||
|| $tf:tt : $ti:ident => $to:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$($acc)* $crate::make_flat_transform!($enum, $tf, $ti, $to),], $to, $($rest)*)
|
||||
};
|
||||
// ── Accumulation: ||? ────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident,
|
||||
||? $tf:tt : $ti:ident => $to:ident, $($rest:tt)*) => {
|
||||
$crate::make_pipe!(@build $enum : $in_ty => $out_ty, $fi,
|
||||
[$($acc)* $crate::make_flat_transform_fallible!($enum, $tf, $ti, $to),], $to, $($rest)*)
|
||||
};
|
||||
|
||||
// ── Termination ───────────────────────────────────────────────────────
|
||||
(@build $enum:ident : $in_ty:ty => $out_ty:ty, $fi:ident,
|
||||
[$($acc:tt)*], $lo:ident $(,)?) => {
|
||||
$crate::Pipe::new(
|
||||
vec![$($acc)*],
|
||||
::std::sync::Arc::new(|x: $in_ty| $enum::$fi(x)),
|
||||
::std::sync::Arc::new(|d: $enum| -> $out_ty {
|
||||
if let $enum::$lo(x) = d { x }
|
||||
else { ::std::unreachable!("unexpected pipeline data variant in make_pipe!") }
|
||||
}),
|
||||
)
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
//! Scheduler: a channel/thread-based pipeline runtime with typed, composable
|
||||
//! stages (`Pipe`), plus the macro-based `PipelineData`-enum runtime it grew
|
||||
//! out of (`Pipeline`/`WorkerPool`).
|
||||
//!
|
||||
//! Submodules: [`error`] (`PipelineError`), [`types`] (function types, `Stage`),
|
||||
//! [`runner`] (thread bodies), [`pool`] (`Pipeline`/`WorkerPool` scheduler
|
||||
//! loop), [`pipe`] (`Pipe`/`PipeIter`), [`macros`] (`make_pipe!` and friends).
|
||||
|
||||
mod error;
|
||||
mod macros;
|
||||
mod pipe;
|
||||
mod pool;
|
||||
mod runner;
|
||||
mod types;
|
||||
|
||||
pub use error::PipelineError;
|
||||
pub use pipe::{Pipe, PipeIter};
|
||||
pub use pool::{Pipeline, WorkerPool};
|
||||
pub use types::{SharedFlatFn, SharedFn, SinkFn, SourceFn, Stage};
|
||||
@@ -0,0 +1,117 @@
|
||||
use crossbeam_channel::{Receiver, bounded};
|
||||
use std::marker::PhantomData;
|
||||
use std::sync::Arc;
|
||||
use std::thread;
|
||||
|
||||
use super::error::PipelineError;
|
||||
use super::pool::{Pipeline, WorkerPool};
|
||||
use super::types::{SinkFn, SourceFn, Stage};
|
||||
|
||||
// ── Pipe ──────────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Typed, composable iterator transformer.
|
||||
///
|
||||
/// A `Pipe<D, In, Out>` is a pure description of pipeline stages — no threads,
|
||||
/// no channels, no scheduler. Call `.apply(iter, n_workers, capacity)` to start
|
||||
/// execution and get back a `PipeIter<Out>`.
|
||||
///
|
||||
/// Compose two pipes with `.then()`: the resulting `Pipe` holds the concatenated
|
||||
/// stage list, so a single scheduler is created when `.apply()` is eventually called.
|
||||
pub struct Pipe<D, In, Out> {
|
||||
stages: Vec<Stage<D>>,
|
||||
wrap: Arc<dyn Fn(In) -> D + Send + Sync>,
|
||||
unwrap: Arc<dyn Fn(D) -> Out + Send + Sync>,
|
||||
_phantom: PhantomData<(In, Out)>,
|
||||
}
|
||||
|
||||
impl<D, In, Out> Pipe<D, In, Out> {
|
||||
/// Build a `Pipe` from stages and wrap/unwrap converters.
|
||||
/// Prefer the `make_pipe!` macro.
|
||||
pub fn new(
|
||||
stages: Vec<Stage<D>>,
|
||||
wrap: Arc<dyn Fn(In) -> D + Send + Sync>,
|
||||
unwrap: Arc<dyn Fn(D) -> Out + Send + Sync>,
|
||||
) -> Self {
|
||||
Self { stages, wrap, unwrap, _phantom: PhantomData }
|
||||
}
|
||||
|
||||
/// Concatenate stages from two pipes into one.
|
||||
///
|
||||
/// Requires `Out` of `self` == `In` of `other`. The single scheduler
|
||||
/// created at `.apply()` time sees the full combined stage list.
|
||||
pub fn then<Next>(self, other: Pipe<D, Out, Next>) -> Pipe<D, In, Next> {
|
||||
Pipe {
|
||||
stages: self.stages.into_iter().chain(other.stages).collect(),
|
||||
wrap: self.wrap,
|
||||
unwrap: other.unwrap,
|
||||
_phantom: PhantomData,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<D, In, Out> Pipe<D, In, Out>
|
||||
where
|
||||
D: Send + Sync + 'static,
|
||||
In: Send + 'static,
|
||||
Out: Send + 'static,
|
||||
{
|
||||
/// Run the pipeline in a background thread; returns an iterator over the output.
|
||||
pub fn apply(
|
||||
self,
|
||||
input: impl Iterator<Item = In> + Send + 'static,
|
||||
n_workers: usize,
|
||||
capacity: usize,
|
||||
) -> PipeIter<Out> {
|
||||
let wrap = Arc::clone(&self.wrap);
|
||||
let unwrap = Arc::clone(&self.unwrap);
|
||||
|
||||
let mut iter = input;
|
||||
let source: SourceFn<D> = Box::new(move || match iter.next() {
|
||||
Some(x) => Ok(wrap(x)),
|
||||
None => Err(PipelineError::EndOfStream),
|
||||
});
|
||||
|
||||
let (out_tx, out_rx) = bounded::<Out>(capacity);
|
||||
let sink: SinkFn<D> = Box::new(move |data: D| {
|
||||
out_tx.send(unwrap(data)).map_err(|_| {
|
||||
PipelineError::StepError(Box::new(std::io::Error::new(
|
||||
std::io::ErrorKind::BrokenPipe,
|
||||
"output channel closed",
|
||||
)))
|
||||
})
|
||||
});
|
||||
|
||||
let pipeline = Pipeline::new(source, self.stages, sink);
|
||||
let handle = thread::spawn(move || {
|
||||
WorkerPool::new(pipeline, n_workers, capacity).run();
|
||||
});
|
||||
|
||||
PipeIter { rx: out_rx, handle: Some(handle) }
|
||||
}
|
||||
}
|
||||
|
||||
// ── PipeIter ──────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Iterator over the output of `Pipe::apply()`.
|
||||
pub struct PipeIter<Out> {
|
||||
rx: Receiver<Out>,
|
||||
handle: Option<thread::JoinHandle<()>>,
|
||||
}
|
||||
|
||||
impl<Out> Iterator for PipeIter<Out> {
|
||||
type Item = Out;
|
||||
|
||||
fn next(&mut self) -> Option<Out> {
|
||||
self.rx.recv().ok()
|
||||
}
|
||||
}
|
||||
|
||||
impl<Out> Drop for PipeIter<Out> {
|
||||
fn drop(&mut self) {
|
||||
// Drain buffered items so the scheduler can unblock if the channel is full.
|
||||
while self.rx.try_recv().is_ok() {}
|
||||
if let Some(h) = self.handle.take() {
|
||||
let _ = h.join();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,203 @@
|
||||
use crossbeam_channel::{Receiver, Select, Sender, bounded};
|
||||
|
||||
use super::error::PipelineError;
|
||||
use super::runner::{dispatch, sink_runner, source_runner, transform_runner};
|
||||
use super::types::{SinkFn, SourceFn, Stage, WorkerTask};
|
||||
|
||||
// ── Pipeline ──────────────────────────────────────────────────────────────────
|
||||
|
||||
pub struct Pipeline<DATA> {
|
||||
source: SourceFn<DATA>,
|
||||
stages: Vec<Stage<DATA>>,
|
||||
sink: SinkFn<DATA>,
|
||||
}
|
||||
|
||||
impl<DATA> Pipeline<DATA> {
|
||||
pub fn new(
|
||||
source: SourceFn<DATA>,
|
||||
stages: Vec<Stage<DATA>>,
|
||||
sink: SinkFn<DATA>,
|
||||
) -> Self {
|
||||
Self { source, stages, sink }
|
||||
}
|
||||
}
|
||||
|
||||
// ── WorkerPool ────────────────────────────────────────────────────────────────
|
||||
|
||||
pub struct WorkerPool<DATA> {
|
||||
pipeline: Pipeline<DATA>,
|
||||
handles: Vec<std::thread::JoinHandle<()>>,
|
||||
n_workers: usize,
|
||||
capacity: usize,
|
||||
}
|
||||
|
||||
impl<DATA> WorkerPool<DATA>
|
||||
where
|
||||
DATA: Send + Sync + 'static,
|
||||
{
|
||||
pub fn new(pipeline: Pipeline<DATA>, n_workers: usize, capacity: usize) -> Self {
|
||||
Self {
|
||||
pipeline,
|
||||
handles: Vec::new(),
|
||||
n_workers,
|
||||
capacity,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn run(mut self) {
|
||||
let n = self.pipeline.stages.len();
|
||||
|
||||
// ── Canaux inter-stages ────────────────────────────────────────────
|
||||
// stage_txs[i] / stage_rxs[i] : sortie du stage i
|
||||
let mut stage_txs: Vec<Sender<Result<DATA, PipelineError>>> = Vec::new();
|
||||
let mut stage_rxs: Vec<Receiver<Result<DATA, PipelineError>>> = Vec::new();
|
||||
for _ in 0..n {
|
||||
let (tx, rx) = bounded(self.capacity);
|
||||
stage_txs.push(tx);
|
||||
stage_rxs.push(rx);
|
||||
}
|
||||
|
||||
// ── Source thread ──────────────────────────────────────────────────
|
||||
let (source_rx, src_handle) = source_runner(self.pipeline.source, self.capacity);
|
||||
self.handles.push(src_handle);
|
||||
|
||||
let stages = self.pipeline.stages;
|
||||
|
||||
// ── Canal delta pour les flat stages ───────────────────────────────
|
||||
// Chaque flat worker envoie `N-1` ici après avoir poussé N items.
|
||||
// Le scheduler ajuste `in_flight` en conséquence.
|
||||
let (flat_delta_tx, flat_delta_rx) = bounded::<isize>(self.capacity);
|
||||
|
||||
// ── Worker pool ────────────────────────────────────────────────────
|
||||
let (worker_tx, worker_rx): (Sender<WorkerTask<DATA>>, Receiver<WorkerTask<DATA>>) =
|
||||
bounded(self.capacity);
|
||||
|
||||
for _ in 0..self.n_workers {
|
||||
self.handles.push(transform_runner(
|
||||
worker_rx.clone(),
|
||||
stages.iter().map(Stage::clone).collect(),
|
||||
stage_txs.clone(),
|
||||
flat_delta_tx.clone(),
|
||||
));
|
||||
}
|
||||
// Le scheduler ne tient plus flat_delta_tx : les workers le détiennent.
|
||||
// On le drop ici pour que le canal se ferme quand les workers terminent.
|
||||
drop(flat_delta_tx);
|
||||
|
||||
// ── Sink thread ────────────────────────────────────────────────────
|
||||
let (sink_tx, sink_err_rx, sink_handle) = sink_runner(self.pipeline.sink, self.capacity);
|
||||
self.handles.push(sink_handle);
|
||||
|
||||
// ── Boucle principale ──────────────────────────────────────────────
|
||||
//
|
||||
// `in_flight` (isize) = nb d'items qui doivent encore atteindre le sink.
|
||||
// Peut temporairement être négatif si un flat worker a poussé ses items
|
||||
// avant que le scheduler ait reçu le delta correspondant.
|
||||
//
|
||||
// `flat_workers_active` = nb de flat workers en cours d'exécution.
|
||||
// Empêche la terminaison prématurée quand in_flight vaut 0 mais qu'un
|
||||
// flat worker n'a pas encore envoyé son delta.
|
||||
//
|
||||
// Priorités du Select biaisé (index le plus bas = priorité la plus haute) :
|
||||
// 0 → sink_err_rx (arrêt immédiat sur erreur sink)
|
||||
// 1 → flat_delta_rx (mettre à jour in_flight avant de dispatcher)
|
||||
// 2..=n+1 → stage_rxs[n-1..0] (vider le pipeline en priorité)
|
||||
// n+2 → source_rx (dernier recours : nouvelles données)
|
||||
//
|
||||
// Quand k = 0 : erreur du sink
|
||||
// Quand k = 1 : delta d'un flat worker
|
||||
// Quand 2 ≤ k ≤ n+1 : résultat du stage n+1-k
|
||||
// Quand k = n+2 : item source
|
||||
//
|
||||
// Terminaison : source tarie ET in_flight == 0 ET aucun flat worker actif.
|
||||
{
|
||||
let mut source_done = false;
|
||||
let mut in_flight: isize = 0;
|
||||
let mut flat_workers_active: usize = 0;
|
||||
|
||||
loop {
|
||||
if source_done && in_flight == 0 && flat_workers_active == 0 {
|
||||
break;
|
||||
}
|
||||
|
||||
let mut sel = Select::new_biased();
|
||||
sel.recv(&sink_err_rx); // index 0
|
||||
sel.recv(&flat_delta_rx); // index 1
|
||||
for rx in stage_rxs.iter().rev() {
|
||||
sel.recv(rx); // indices 2..=n+1
|
||||
}
|
||||
let src_idx = if !source_done {
|
||||
Some(sel.recv(&source_rx)) // index n+2
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
let oper = sel.select();
|
||||
let k = oper.index();
|
||||
|
||||
if k == 0 {
|
||||
// ── Erreur du sink ────────────────────────────────────
|
||||
match oper.recv(&sink_err_rx) {
|
||||
Ok(e) => { eprintln!("Sink error: {:?}", e); break; }
|
||||
Err(_) => break,
|
||||
}
|
||||
} else if k == 1 {
|
||||
// ── Delta d'un flat worker ────────────────────────────
|
||||
// delta = N - 1 (N items poussés, 1 item consommé)
|
||||
match oper.recv(&flat_delta_rx) {
|
||||
Ok(delta) => {
|
||||
in_flight += delta;
|
||||
flat_workers_active -= 1;
|
||||
}
|
||||
Err(_) => {}
|
||||
}
|
||||
} else if src_idx == Some(k) {
|
||||
// ── Nouvel item depuis la source ──────────────────────
|
||||
match oper.recv(&source_rx) {
|
||||
Ok(Ok(data)) => {
|
||||
if n == 0 {
|
||||
let _ = sink_tx.send(data);
|
||||
} else {
|
||||
in_flight += 1;
|
||||
dispatch(
|
||||
data, 0,
|
||||
&stages, &worker_tx,
|
||||
&mut flat_workers_active,
|
||||
);
|
||||
}
|
||||
}
|
||||
Ok(Err(e)) => eprintln!("Source error: {:?}", e),
|
||||
Err(_) => source_done = true,
|
||||
}
|
||||
} else {
|
||||
// ── Résultat d'un stage intermédiaire ─────────────────
|
||||
// k ∈ [2, n+1] → stage = n+1 - k
|
||||
let stage = n + 1 - k;
|
||||
match oper.recv(&stage_rxs[stage]) {
|
||||
Ok(Ok(data)) => {
|
||||
if stage == n - 1 {
|
||||
in_flight -= 1;
|
||||
let _ = sink_tx.send(data);
|
||||
} else {
|
||||
dispatch(
|
||||
data, stage + 1,
|
||||
&stages, &worker_tx,
|
||||
&mut flat_workers_active,
|
||||
);
|
||||
}
|
||||
}
|
||||
Ok(Err(e)) => eprintln!("Stage {} error: {:?}", stage, e),
|
||||
Err(_) => break,
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
drop(worker_tx);
|
||||
drop(sink_tx);
|
||||
|
||||
for h in self.handles {
|
||||
let _ = h.join();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,118 @@
|
||||
use crossbeam_channel::{Receiver, Sender, bounded};
|
||||
use std::thread;
|
||||
|
||||
use super::error::PipelineError;
|
||||
use super::types::{SinkFn, SourceFn, Stage, WorkerTask};
|
||||
|
||||
// ── Thread runners ────────────────────────────────────────────────────────────
|
||||
|
||||
pub(super) fn source_runner<DATA>(
|
||||
mut source: SourceFn<DATA>,
|
||||
capacity: usize,
|
||||
) -> (
|
||||
Receiver<Result<DATA, PipelineError>>,
|
||||
thread::JoinHandle<()>,
|
||||
)
|
||||
where
|
||||
DATA: Send + Sync + 'static,
|
||||
{
|
||||
let (tx, rx) = bounded(capacity);
|
||||
let handle = thread::spawn(move || {
|
||||
loop {
|
||||
match source() {
|
||||
Ok(data) => {
|
||||
if tx.send(Ok(data)).is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
Err(PipelineError::EndOfStream) => break,
|
||||
Err(e) => {
|
||||
eprintln!("Source error: {:?}", e);
|
||||
let _ = tx.send(Err(e));
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
(rx, handle)
|
||||
}
|
||||
|
||||
/// Lance un thread worker du pool.
|
||||
///
|
||||
/// Gère deux types de tâches :
|
||||
/// - `Transform` : applique `f(data)` et envoie le résultat dans `result_tx`.
|
||||
/// - `Flat` : appelle `f(data, &push_tx, &delta_tx)` ; la fonction elle-même
|
||||
/// pousse ses items dans `push_tx` et envoie `N-1` dans `delta_tx`.
|
||||
pub(super) fn transform_runner<DATA>(
|
||||
task_rx: Receiver<WorkerTask<DATA>>,
|
||||
stages: Vec<Stage<DATA>>,
|
||||
stage_txs: Vec<Sender<Result<DATA, PipelineError>>>,
|
||||
flat_delta_tx: Sender<isize>,
|
||||
) -> thread::JoinHandle<()>
|
||||
where
|
||||
DATA: Send + Sync + 'static,
|
||||
{
|
||||
thread::spawn(move || {
|
||||
while let Ok(task) = task_rx.recv() {
|
||||
match task {
|
||||
WorkerTask::Transform(data, idx) => {
|
||||
if let Stage::Transform(f) = &stages[idx] {
|
||||
let _ = stage_txs[idx].send(f(data));
|
||||
}
|
||||
}
|
||||
WorkerTask::Flat(data, idx) => {
|
||||
if let Stage::Flat(f) = &stages[idx] {
|
||||
f(data, &stage_txs[idx], &flat_delta_tx);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/// Lance le thread sink.
|
||||
pub(super) fn sink_runner<DATA>(
|
||||
sink: SinkFn<DATA>,
|
||||
capacity: usize,
|
||||
) -> (
|
||||
Sender<DATA>,
|
||||
Receiver<PipelineError>,
|
||||
thread::JoinHandle<()>,
|
||||
)
|
||||
where
|
||||
DATA: Send + Sync + 'static,
|
||||
{
|
||||
let (data_tx, data_rx) = bounded(capacity);
|
||||
let (err_tx, err_rx) = bounded(capacity);
|
||||
let handle = thread::spawn(move || {
|
||||
for data in data_rx {
|
||||
if let Err(e) = sink(data) {
|
||||
let _ = err_tx.send(e);
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
(data_tx, err_rx, handle)
|
||||
}
|
||||
|
||||
/// Envoie `data` au stage `stage_idx`.
|
||||
/// Pour un `Transform`, empile une `WorkerTask::Transform`.
|
||||
/// Pour un `Flat`, incrémente `flat_workers_active` et empile une `WorkerTask::Flat`.
|
||||
#[inline]
|
||||
pub(super) fn dispatch<DATA>(
|
||||
data: DATA,
|
||||
stage_idx: usize,
|
||||
stages: &[Stage<DATA>],
|
||||
worker_tx: &Sender<WorkerTask<DATA>>,
|
||||
flat_workers_active: &mut usize,
|
||||
) {
|
||||
match &stages[stage_idx] {
|
||||
Stage::Transform(_) => {
|
||||
let _ = worker_tx.send(WorkerTask::Transform(data, stage_idx));
|
||||
}
|
||||
Stage::Flat(_) => {
|
||||
*flat_workers_active += 1;
|
||||
let _ = worker_tx.send(WorkerTask::Flat(data, stage_idx));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,49 @@
|
||||
use crossbeam_channel::Sender;
|
||||
use std::sync::Arc;
|
||||
|
||||
use super::error::PipelineError;
|
||||
|
||||
// ── Function types ────────────────────────────────────────────────────────────
|
||||
|
||||
/// Fonction source : appelée répétitivement, retourne le prochain item ou EndOfStream.
|
||||
/// `FnMut` car elle maintient un état interne (position dans l'itérateur).
|
||||
pub type SourceFn<D> = Box<dyn FnMut() -> Result<D, PipelineError> + Send>;
|
||||
|
||||
/// Fonction sink : consomme un item final, peut échouer (erreur d'I/O, etc.).
|
||||
pub type SinkFn<D> = Box<dyn Fn(D) -> Result<(), PipelineError> + Send>;
|
||||
|
||||
/// Fonction de transformation partagée entre workers via Arc.
|
||||
pub type SharedFn<D> = Arc<dyn Fn(D) -> Result<D, PipelineError> + Send + Sync>;
|
||||
|
||||
/// Fonction de transformation 1→N (flat map) partagée entre workers via Arc.
|
||||
///
|
||||
/// La fonction reçoit l'item d'entrée, un canal `push` pour envoyer chaque item
|
||||
/// produit, et un canal `delta` pour signaler au scheduler combien d'items
|
||||
/// supplémentaires sont entrés dans le pipeline (N-1 si N items produits).
|
||||
/// Elle doit appeler `delta.send(N - 1)` **après** avoir poussé tous les items.
|
||||
pub type SharedFlatFn<D> =
|
||||
Arc<dyn Fn(D, &Sender<Result<D, PipelineError>>, &Sender<isize>) + Send + Sync>;
|
||||
|
||||
// ── Stage enum ────────────────────────────────────────────────────────────────
|
||||
|
||||
/// Une étape du pipeline : transform classique (1→1) ou flat transform (1→N).
|
||||
pub enum Stage<D> {
|
||||
Transform(SharedFn<D>),
|
||||
Flat(SharedFlatFn<D>),
|
||||
}
|
||||
|
||||
impl<D> Clone for Stage<D> {
|
||||
fn clone(&self) -> Self {
|
||||
match self {
|
||||
Stage::Transform(f) => Stage::Transform(Arc::clone(f)),
|
||||
Stage::Flat(f) => Stage::Flat(Arc::clone(f)),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// ── Worker task ───────────────────────────────────────────────────────────────
|
||||
|
||||
pub(super) enum WorkerTask<D> {
|
||||
Transform(D, usize),
|
||||
Flat(D, usize),
|
||||
}
|
||||
Reference in New Issue
Block a user