Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

8 changes: 8 additions & 0 deletions crates/rb-task/Cargo.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
[package]
name = "rb-task"
version.workspace = true
edition.workspace = true
license.workspace = true
description = "Dependency graphs and bounded task execution"

[dependencies]
45 changes: 45 additions & 0 deletions crates/rb-task/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
# rb-task

Dependency graphs, bounded workers, and task events using only the standard library.

```rust
use rb_task::{TaskEvents, Executor, TaskAction, TaskGraph};
use std::{collections::BTreeMap, sync::Arc};

let mut graph = TaskGraph::new();
let prepare = graph.add("prepare", []);
let finish = graph.add("finish", [prepare]);
let actions = BTreeMap::from([
(prepare, Arc::new(|context: rb_task::TaskContext| {
context.output("preparing input");
Ok(())
}) as TaskAction),
(finish, Arc::new(|_| Ok(())) as TaskAction),
]);
Executor::new(2, TaskEvents::default()).run(graph, actions)?;
# Ok::<(), rb_task::ExecutionError>(())
```

Subscribe a `TaskEventSink` to the `TaskEvents` to receive lifecycle, progress,
output, and duration events. Callbacks run synchronously on workers and may
run concurrently; slow callbacks delay their worker, and panics fail the run.
The crate does not format output.

Graphs are fixed per run. Task IDs are graph-local indexes; callers must not mix
IDs from different graphs. Ready tasks run within
the worker limit as dependencies finish. On an observed action failure or panic,
or an observer panic, execution stops dispatching and waits for running work.
Cancellation, retries, and graph expansion belong to the caller.

Run `cargo test -p rb-task` for the example and regression suite.

Runnable examples:

- [Dependencies](examples/dependencies.rs): two tasks run independently between
shared preparation and completion steps.
- [Reporting](examples/reporting.rs): observe lifecycle, output, and progress events.

```sh
cargo run -p rb-task --example dependencies
cargo run -p rb-task --example reporting
```
28 changes: 28 additions & 0 deletions crates/rb-task/examples/dependencies.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
use rb_task::{ExecutionError, Executor, TaskAction, TaskContext, TaskEvents, TaskGraph};
use std::{collections::BTreeMap, sync::Arc, thread, time::Duration};

fn main() -> Result<(), ExecutionError> {
let mut graph = TaskGraph::new();
let prepare = graph.add("prepare", []);
let left = graph.add("process left", [prepare]);
let right = graph.add("process right", [prepare]);
let finish = graph.add("finish", [left, right]);

let actions = BTreeMap::from([
(prepare, action("prepare")),
(left, action("process left")),
(right, action("process right")),
(finish, action("finish")),
]);

Executor::new(2, TaskEvents::default()).run(graph, actions)
}

fn action(label: &'static str) -> TaskAction {
Arc::new(move |context: TaskContext| {
println!("worker {}: starting {label}", context.worker);
thread::sleep(Duration::from_millis(50));
println!("worker {}: finished {label}", context.worker);
Ok(())
})
}
45 changes: 45 additions & 0 deletions crates/rb-task/examples/reporting.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
use rb_task::{
ExecutionError, Executor, TaskAction, TaskContext, TaskEvent, TaskEventSink, TaskEvents,
TaskGraph,
};
use std::{collections::BTreeMap, sync::Arc};

struct Console;

impl TaskEventSink for Console {
fn event(&self, event: TaskEvent) {
match event {
TaskEvent::Started { worker, label, .. } => {
println!("worker {worker}: {label}");
}
TaskEvent::Output { line, .. } => println!(" {line}"),
TaskEvent::Progress {
phase, done, total, ..
} => {
println!(" {phase}: {done}/{total}");
}
TaskEvent::Finished { elapsed, .. } => println!(" finished in {elapsed:?}"),
TaskEvent::Failed { error, .. } => println!(" failed: {error}"),
}
}
}

fn main() -> Result<(), ExecutionError> {
let events = TaskEvents::default();
events.subscribe(Arc::new(Console));

let mut graph = TaskGraph::new();
let process = graph.add("process documents", []);
let actions = BTreeMap::from([(
process,
Arc::new(|context: TaskContext| {
for done in 1..=3 {
context.output(format!("Processed document {done}"));
context.progress("processing", done, 3, None);
}
Ok(())
}) as TaskAction,
)]);

Executor::new(1, events).run(graph, actions)
}
38 changes: 38 additions & 0 deletions crates/rb-task/src/error.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,38 @@
use crate::{CompletionError, PlanError, TaskId};
use std::{error::Error, fmt};

/// Identifies a validation, task, or panic failure encountered during execution.
#[derive(Debug, Eq, PartialEq)]
pub enum ExecutionError {
InvalidActions,
InvalidGraph(PlanError),
InvalidCompletion(CompletionError),
TaskFailed { task: TaskId, message: String },
Panicked { task: TaskId },
}

impl fmt::Display for ExecutionError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::InvalidActions => formatter.write_str("missing or invalid task action"),
Self::InvalidGraph(error) => write!(formatter, "invalid task graph: {error}"),
Self::InvalidCompletion(error) => write!(formatter, "invalid task completion: {error}"),
Self::TaskFailed { task, message } => {
write!(formatter, "task {} failed: {message}", task.index())
}
Self::Panicked { task } => {
write!(formatter, "task {} or its observer panicked", task.index())
}
}
}
}

impl Error for ExecutionError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::InvalidGraph(error) => Some(error),
Self::InvalidCompletion(error) => Some(error),
_ => None,
}
}
}
67 changes: 67 additions & 0 deletions crates/rb-task/src/events.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
use crate::TaskId;
use std::sync::{Arc, Mutex};
use std::time::Duration;

/// Reports task lifecycle changes, output, and progress to observers.
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum TaskEvent {
Started {
task: TaskId,
worker: usize,
label: String,
elapsed: Duration,
},
Output {
task: TaskId,
worker: usize,
line: String,
elapsed: Duration,
},
Progress {
task: TaskId,
worker: usize,
phase: &'static str,
done: usize,
total: usize,
detail: Option<String>,
elapsed: Duration,
},
Finished {
task: TaskId,
worker: usize,
elapsed: Duration,
},
Failed {
task: TaskId,
worker: usize,
error: String,
elapsed: Duration,
},
}

/// Receives events synchronously on workers; callbacks may run concurrently.
pub trait TaskEventSink: Send + Sync {
fn event(&self, event: TaskEvent);
}

/// Shares subscriptions and delivers each event to registered observers.
#[derive(Default, Clone)]
pub struct TaskEvents {
sinks: Arc<Mutex<Vec<Arc<dyn TaskEventSink>>>>,
}

impl TaskEvents {
pub fn subscribe(&self, sink: Arc<dyn TaskEventSink>) {
self.sinks.lock().unwrap().push(sink);
}

pub fn emit(&self, event: TaskEvent) {
let sinks = self.sinks.lock().unwrap().clone();
for sink in sinks {
sink.event(event.clone());
}
}
}

#[cfg(test)]
mod tests;
59 changes: 59 additions & 0 deletions crates/rb-task/src/events/tests.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
use super::*;
use std::sync::{
atomic::{AtomicUsize, Ordering},
mpsc,
};
use std::thread;

struct CountingSink(AtomicUsize);

impl TaskEventSink for CountingSink {
fn event(&self, _event: TaskEvent) {
self.0.fetch_add(1, Ordering::Relaxed);
}
}

#[test]
fn event_bus_fans_out_without_owning_execution() {
let bus = TaskEvents::default();
let sink = Arc::new(CountingSink(AtomicUsize::new(0)));
let other = Arc::new(CountingSink(AtomicUsize::new(0)));
bus.subscribe(sink.clone());
bus.subscribe(other.clone());
bus.emit(TaskEvent::Finished {
task: TaskId(0),
worker: 0,
elapsed: Duration::ZERO,
});
assert_eq!(sink.0.load(Ordering::Relaxed), 1);
assert_eq!(other.0.load(Ordering::Relaxed), 1);
}

#[test]
fn event_callbacks_can_subscribe_without_locking_the_bus() {
struct Subscriber(TaskEvents);
struct Ignore;
impl TaskEventSink for Ignore {
fn event(&self, _: TaskEvent) {}
}
impl TaskEventSink for Subscriber {
fn event(&self, _: TaskEvent) {
self.0.subscribe(Arc::new(Ignore));
}
}
let bus = TaskEvents::default();
bus.subscribe(Arc::new(Subscriber(bus.clone())));
let emitting = bus.clone();
let (send, receive) = mpsc::channel();
thread::spawn(move || {
emitting.emit(TaskEvent::Finished {
task: TaskId(0),
worker: 0,
elapsed: Duration::ZERO,
});
send.send(()).unwrap();
});
receive.recv_timeout(Duration::from_secs(2)).unwrap();
assert_eq!(bus.sinks.lock().unwrap().len(), 2);
bus.sinks.lock().unwrap().clear();
}
Loading
Loading