Skip to main content

WorkQueue

Trait WorkQueue 

Source
pub trait WorkQueue: Send + Sync {
    // Required methods
    fn claim_due<'life0, 'life1, 'life2, 'async_trait>(
        &'life0 self,
        spec_ids: &'life1 [String],
        worker_id: &'life2 str,
        lease_secs: i64,
        batch: i64,
    ) -> Pin<Box<dyn Future<Output = Result<Vec<WorkItem>, SupervisorError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait,
             'life2: 'async_trait;
    fn renew<'life0, 'life1, 'life2, 'life3, 'async_trait>(
        &'life0 self,
        tenant_id: &'life1 str,
        id: &'life2 str,
        worker_id: &'life3 str,
        lease_version: i64,
        lease_secs: i64,
    ) -> Pin<Box<dyn Future<Output = Result<(), SupervisorError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait,
             'life2: 'async_trait,
             'life3: 'async_trait;
    fn commit_command<'life0, 'life1, 'life2, 'async_trait>(
        &'life0 self,
        item: &'life1 WorkItem,
        command: &'life2 WorkflowTransitionCommand,
    ) -> Pin<Box<dyn Future<Output = Result<(), SupervisorError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait,
             'life2: 'async_trait;

    // Provided methods
    fn load_revision<'life0, 'life1, 'async_trait>(
        &'life0 self,
        _item: &'life1 WorkItem,
    ) -> Pin<Box<dyn Future<Output = Result<WorkflowRevision, SupervisorError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait { ... }
    fn admit_evaluation<'life0, 'life1, 'async_trait>(
        &'life0 self,
        _item: &'life1 WorkItem,
    ) -> Pin<Box<dyn Future<Output = Result<(), SupervisorError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait { ... }
    fn prepare_evaluation<'life0, 'life1, 'async_trait>(
        &'life0 self,
        item: &'life1 WorkItem,
        load_revision: bool,
    ) -> Pin<Box<dyn Future<Output = Result<Option<WorkflowRevision>, SupervisorError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait { ... }
}
Expand description

Durable queue the supervisor claims from and commits to.

Required Methods§

Source

fn claim_due<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, spec_ids: &'life1 [String], worker_id: &'life2 str, lease_secs: i64, batch: i64, ) -> Pin<Box<dyn Future<Output = Result<Vec<WorkItem>, SupervisorError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Claim due instances for spec_ids; an empty slice selects all specs.

Source

fn renew<'life0, 'life1, 'life2, 'life3, 'async_trait>( &'life0 self, tenant_id: &'life1 str, id: &'life2 str, worker_id: &'life3 str, lease_version: i64, lease_secs: i64, ) -> Pin<Box<dyn Future<Output = Result<(), SupervisorError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait,

Renew a claim’s lease.

Source

fn commit_command<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, item: &'life1 WorkItem, command: &'life2 WorkflowTransitionCommand, ) -> Pin<Box<dyn Future<Output = Result<(), SupervisorError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Commit a driver command atomically.

Provided Methods§

Source

fn load_revision<'life0, 'life1, 'async_trait>( &'life0 self, _item: &'life1 WorkItem, ) -> Pin<Box<dyn Future<Output = Result<WorkflowRevision, SupervisorError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Read the immutable revision of an owned, currently leased instance. Stores without revision persistence reject dynamic spec execution.

Source

fn admit_evaluation<'life0, 'life1, 'async_trait>( &'life0 self, _item: &'life1 WorkItem, ) -> Pin<Box<dyn Future<Output = Result<(), SupervisorError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Recheck current owner admission before any graph/read or custom-driver evaluation. Queues without an authority source fail closed. Terminal system observations and cancellation still enter their existing convergence paths.

Source

fn prepare_evaluation<'life0, 'life1, 'async_trait>( &'life0 self, item: &'life1 WorkItem, load_revision: bool, ) -> Pin<Box<dyn Future<Output = Result<Option<WorkflowRevision>, SupervisorError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Recheck current admission and, when requested, load the exact immutable revision needed by a spec driver. Stores may override this combined boundary to avoid repeating the same claimed-revision and capability reads. The default preserves the original fail-closed sequence.

Dyn Compatibility§

This trait is dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§