pub struct DistributedDagExecutor<Q: WorkQueuePort> { /* private fields */ }Expand description
Executes a DAG wave using a WorkQueuePort to distribute node-level tasks
across workers.
Workers are spawned as Tokio tasks that pull from the queue, call the
appropriate service, and acknowledge results. For local development the
LocalWorkQueue is used; in production any queue backend can be plugged
in without changing this executor.
§Example
use stygian_graph::adapters::distributed::{DistributedDagExecutor, LocalWorkQueue};
use stygian_graph::ports::work_queue::WorkTask;
use stygian_graph::adapters::noop::NoopService;
use serde_json::json;
use std::sync::Arc;
use std::collections::HashMap;
let queue = Arc::new(LocalWorkQueue::new());
let executor = DistributedDagExecutor::new(queue, 4);
let mut services: HashMap<String, Arc<dyn stygian_graph::ports::ScrapingService>> =
HashMap::new();
services.insert("noop".to_string(), Arc::new(NoopService));
let tasks = vec![WorkTask {
id: "p1::fetch::01".to_string(),
pipeline_id: "p1".to_string(),
node_name: "fetch".to_string(),
input: json!({"url": "https://example.com"}),
wave: 0,
attempt: 0,
idempotency_key: "ik-01".to_string(),
}];
let results = executor.execute_wave("p1", tasks, &services).await.unwrap();
assert!(!results.is_empty() || results.is_empty()); // noop returns emptyImplementations§
Source§impl<Q: WorkQueuePort + 'static> DistributedDagExecutor<Q>
impl<Q: WorkQueuePort + 'static> DistributedDagExecutor<Q>
Sourcepub fn new(queue: Arc<Q>, worker_concurrency: usize) -> Self
pub fn new(queue: Arc<Q>, worker_concurrency: usize) -> Self
Create a new executor with the given work queue and worker concurrency.
worker_concurrency controls how many parallel worker tasks drain the
queue.
Sourcepub async fn execute_wave(
&self,
pipeline_id: &str,
tasks: Vec<WorkTask>,
services: &HashMap<String, Arc<dyn ScrapingService>>,
) -> Result<Vec<(String, Value)>>
pub async fn execute_wave( &self, pipeline_id: &str, tasks: Vec<WorkTask>, services: &HashMap<String, Arc<dyn ScrapingService>>, ) -> Result<Vec<(String, Value)>>
Execute a single wave of tasks, distributing them across workers.
Returns (node_name, output) pairs for all tasks in the wave.
§Panics
Panics if an internal Mutex is poisoned (i.e. another thread panicked
while holding the lock). Treat this as unrecoverable.
§Errors
Returns StygianError when a service reports a failure, the executor
is shut down, or a worker task cannot be enqueued.
Auto Trait Implementations§
impl<Q> Freeze for DistributedDagExecutor<Q>
impl<Q> RefUnwindSafe for DistributedDagExecutor<Q>where
Q: RefUnwindSafe,
impl<Q> Send for DistributedDagExecutor<Q>
impl<Q> Sync for DistributedDagExecutor<Q>
impl<Q> Unpin for DistributedDagExecutor<Q>
impl<Q> UnsafeUnpin for DistributedDagExecutor<Q>
impl<Q> UnwindSafe for DistributedDagExecutor<Q>where
Q: RefUnwindSafe,
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self>
fn instrument(self, span: Span) -> Instrumented<Self>
Source§fn in_current_span(self) -> Instrumented<Self>
fn in_current_span(self) -> Instrumented<Self>
Source§impl<T> Paint for Twhere
T: ?Sized,
impl<T> Paint for Twhere
T: ?Sized,
Source§fn fg(&self, value: Color) -> Painted<&T>
fn fg(&self, value: Color) -> Painted<&T>
Returns a styled value derived from self with the foreground set to
value.
This method should be used rarely. Instead, prefer to use color-specific
builder methods like red() and
green(), which have the same functionality but are
pithier.
§Example
Set foreground color to white using fg():
use yansi::{Paint, Color};
painted.fg(Color::White);Set foreground color to white using white().
use yansi::Paint;
painted.white();Source§fn bright_black(&self) -> Painted<&T>
fn bright_black(&self) -> Painted<&T>
Source§fn bright_red(&self) -> Painted<&T>
fn bright_red(&self) -> Painted<&T>
Source§fn bright_green(&self) -> Painted<&T>
fn bright_green(&self) -> Painted<&T>
Source§fn bright_yellow(&self) -> Painted<&T>
fn bright_yellow(&self) -> Painted<&T>
Source§fn bright_blue(&self) -> Painted<&T>
fn bright_blue(&self) -> Painted<&T>
Source§fn bright_magenta(&self) -> Painted<&T>
fn bright_magenta(&self) -> Painted<&T>
Source§fn bright_cyan(&self) -> Painted<&T>
fn bright_cyan(&self) -> Painted<&T>
Source§fn bright_white(&self) -> Painted<&T>
fn bright_white(&self) -> Painted<&T>
Source§fn bg(&self, value: Color) -> Painted<&T>
fn bg(&self, value: Color) -> Painted<&T>
Returns a styled value derived from self with the background set to
value.
This method should be used rarely. Instead, prefer to use color-specific
builder methods like on_red() and
on_green(), which have the same functionality but
are pithier.
§Example
Set background color to red using fg():
use yansi::{Paint, Color};
painted.bg(Color::Red);Set background color to red using on_red().
use yansi::Paint;
painted.on_red();Source§fn on_primary(&self) -> Painted<&T>
fn on_primary(&self) -> Painted<&T>
Source§fn on_magenta(&self) -> Painted<&T>
fn on_magenta(&self) -> Painted<&T>
Source§fn on_bright_black(&self) -> Painted<&T>
fn on_bright_black(&self) -> Painted<&T>
Source§fn on_bright_red(&self) -> Painted<&T>
fn on_bright_red(&self) -> Painted<&T>
Source§fn on_bright_green(&self) -> Painted<&T>
fn on_bright_green(&self) -> Painted<&T>
Source§fn on_bright_yellow(&self) -> Painted<&T>
fn on_bright_yellow(&self) -> Painted<&T>
Source§fn on_bright_blue(&self) -> Painted<&T>
fn on_bright_blue(&self) -> Painted<&T>
Source§fn on_bright_magenta(&self) -> Painted<&T>
fn on_bright_magenta(&self) -> Painted<&T>
Source§fn on_bright_cyan(&self) -> Painted<&T>
fn on_bright_cyan(&self) -> Painted<&T>
Source§fn on_bright_white(&self) -> Painted<&T>
fn on_bright_white(&self) -> Painted<&T>
Source§fn attr(&self, value: Attribute) -> Painted<&T>
fn attr(&self, value: Attribute) -> Painted<&T>
Enables the styling Attribute value.
This method should be used rarely. Instead, prefer to use
attribute-specific builder methods like bold() and
underline(), which have the same functionality
but are pithier.
§Example
Make text bold using attr():
use yansi::{Paint, Attribute};
painted.attr(Attribute::Bold);Make text bold using using bold().
use yansi::Paint;
painted.bold();Source§fn rapid_blink(&self) -> Painted<&T>
fn rapid_blink(&self) -> Painted<&T>
Source§fn quirk(&self, value: Quirk) -> Painted<&T>
fn quirk(&self, value: Quirk) -> Painted<&T>
Enables the yansi Quirk value.
This method should be used rarely. Instead, prefer to use quirk-specific
builder methods like mask() and
wrap(), which have the same functionality but are
pithier.
§Example
Enable wrapping using .quirk():
use yansi::{Paint, Quirk};
painted.quirk(Quirk::Wrap);Enable wrapping using wrap().
use yansi::Paint;
painted.wrap();Source§fn clear(&self) -> Painted<&T>
👎Deprecated since 1.0.1: renamed to resetting() due to conflicts with Vec::clear().
The clear() method will be removed in a future release.
fn clear(&self) -> Painted<&T>
renamed to resetting() due to conflicts with Vec::clear().
The clear() method will be removed in a future release.
Source§fn whenever(&self, value: Condition) -> Painted<&T>
fn whenever(&self, value: Condition) -> Painted<&T>
Conditionally enable styling based on whether the Condition value
applies. Replaces any previous condition.
See the crate level docs for more details.
§Example
Enable styling painted only when both stdout and stderr are TTYs:
use yansi::{Paint, Condition};
painted.red().on_yellow().whenever(Condition::STDOUTERR_ARE_TTY);