Skip to main content

Workflow

Struct Workflow 

Source
#[non_exhaustive]
pub struct Workflow {
Show 15 fields pub id: String, pub name: String, pub priority: u32, pub description: Option<String>, pub condition: Value, pub tasks: Vec<Task>, pub continue_on_error: bool, pub channel: String, pub version: u32, pub status: WorkflowStatus, pub rollout: Option<Rollout>, pub loop_config: Option<LoopConfig>, pub tags: Vec<String>, pub created_at: Option<DateTime<Utc>>, pub updated_at: Option<DateTime<Utc>>, /* private fields */
}
Expand description

Workflow represents a collection of tasks that execute sequentially (also known as a Rule in rules-engine terminology).

Conditions are evaluated against the full message context, including data, metadata, and temp_data fields. A collection of tasks with a condition, executed against a message.

#[non_exhaustive]: construct through Workflow::new, Workflow::rule or Workflow::from_json and assign the public fields you need. Field reads and writes are unaffected.

Same reason as Task: several fields are engine internals marked not part of the stable API, and struct-literal construction forced callers to name them.

Fields (Non-exhaustive)§

This struct is marked as non-exhaustive
Non-exhaustive structs could have additional fields added in future. Therefore, non-exhaustive structs cannot be constructed in external crates using the traditional Struct { .. } syntax; cannot be matched against without a wildcard ..; and struct update syntax will not work.
§id: String§name: String§priority: u32§description: Option<String>§condition: Value§tasks: Vec<Task>

The workflow’s steps, flattened.

The JSON tasks array holds steps: an element carrying a tasks key is a TaskGroup, anything else is a Task. The parser flattens the tree in document order and records each group’s span on the task that opens it (Task::group_starts), so this stays a flat list and the executor keeps walking &[Task] slices.

§continue_on_error: bool§channel: String

Channel for routing (default: “default”)

§version: u32

Version number for rule versioning (default: 1)

§status: WorkflowStatus

Workflow status — Active, Paused, or Archived (default: Active)

§rollout: Option<Rollout>

Traffic split for this workflow. None (the default) means the workflow is not part of a split and runs for every message.

A workflow with a rollout is skipped when the message’s crate::Message::routing_bucket falls outside the range. A message with no bucket is admitted — see Rollout.

§loop_config: Option<LoopConfig>

Engine-managed loop over this workflow’s task list. None (the default) runs the task list exactly once — the historical behaviour, on a code path that carries no loop overhead.

See LoopConfig for the per-sweep contract.

§tags: Vec<String>

Tags for categorization and filtering

§created_at: Option<DateTime<Utc>>

Creation timestamp

§updated_at: Option<DateTime<Utc>>

Last update timestamp

Implementations§

Source§

impl Workflow

Source

pub fn validate_authored(json: &Value) -> Vec<WorkflowIssue>

Check authored workflow JSON without building an engine.

Returns empty if and only if the JSON parses into a Workflow and that workflow validates.

That is the shape question, and it is the whole of it. It is not the same as “this engine can run it”: Engine::build also resolves every task to a handler and parses custom inputs, so a structurally perfect definition naming an unregistered function still aborts a build. Engine::check_workflow answers that half; run both.

§How the guarantee holds

Three stages. A structural walk collects every semantic violation with the coordinate the author typed; if it finds none, the document is then actually parsed and validated, and either failure is reported as one further issue. The second and third stages are what make the promise true by construction: the crate’s serde schema is far larger than any rule list — a "priority": "high" or a map task missing mappings breaks no semantic rule and still cannot load — and mirroring it here would recreate the very drift this API exists to remove.

§Example
use dataflow_rs::{IssueCode, Workflow};
use serde_json::json;

let broken = json!({
    "id": "w", "name": "w", "priority": 0,
    "tasks": [
        {"id": "dup", "name": "a", "function": {"name": "map", "input": {"mappings": []}}},
        {"id": "dup", "name": "b", "function": {"name": "map", "input": {"mappings": []}}}
    ]
});

let issues = Workflow::validate_authored(&broken);
assert_eq!(issues[0].code, IssueCode::DuplicateStepId);
assert_eq!(issues[0].path.as_deref(), Some("tasks[1].id"));
assert_eq!(issues[0].task_id.as_deref(), Some("dup"));

Every problem is reported, not just the first:

let issues = Workflow::validate_authored(&json!({
    "id": "", "name": "w", "tasks": [{"id": "t", "name": "t"}]
}));

let codes: Vec<IssueCode> = issues.iter().map(|i| i.code).collect();
assert!(codes.contains(&IssueCode::EmptyWorkflowId));
assert!(codes.contains(&IssueCode::MissingFunction));
Source§

impl Workflow

Source

pub fn new() -> Self

Source

pub fn rule(id: &str, name: &str, condition: Value, tasks: Vec<Task>) -> Self

Create a workflow (rule) with a condition and tasks.

This is a convenience constructor for the IFTTT-style rules engine pattern: IF condition THEN execute tasks.

§Arguments
  • id - Unique identifier for the rule
  • name - Human-readable name
  • condition - JSONLogic condition evaluated against the full context (data, metadata, temp_data)
  • tasks - Actions to execute when the condition is met
Source

pub fn from_json(json_str: &str) -> Result<Self>

Load workflow from JSON string

Source

pub fn from_file<P: AsRef<Path>>(path: P) -> Result<Self>

Load workflow from JSON file

Source

pub fn validate(&self) -> Result<()>

Validate the workflow structure

Source§

impl Workflow

Source

pub fn connector_refs(&self) -> impl Iterator<Item = ConnectorRef<'_>>

Every connector reference in this workflow, in task order.

Tasks whose function names no connector are skipped. One item is yielded per task, not per distinct connector: two tasks on the same connector yield two items. Callers wanting a distinct set collect one themselves.

Does not require a compiled workflow — this reads only deserialized fields, so it works on the output of Workflow::from_json before the engine has compiled it.

Which configs carry a connector is this crate’s fact; deriving it here rather than reimplementing the set downstream is the point.

Trait Implementations§

Source§

impl Clone for Workflow

Source§

fn clone(&self) -> Workflow

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for Workflow

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more
Source§

impl Default for Workflow

Source§

fn default() -> Self

Returns the “default value” for a type. Read more
Source§

impl<'de> Deserialize<'de> for Workflow

Source§

fn deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>
where __D: Deserializer<'de>,

Deserialize this value from the given Serde deserializer. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> DeserializeOwned for T
where T: for<'de> Deserialize<'de>,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.