#[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
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: StringChannel for routing (default: “default”)
version: u32Version number for rule versioning (default: 1)
status: WorkflowStatusWorkflow 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 for categorization and filtering
created_at: Option<DateTime<Utc>>Creation timestamp
updated_at: Option<DateTime<Utc>>Last update timestamp
Implementations§
Source§impl Workflow
impl Workflow
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
impl Workflow
pub fn new() -> Self
Sourcepub fn rule(id: &str, name: &str, condition: Value, tasks: Vec<Task>) -> Self
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 rulename- Human-readable namecondition- JSONLogic condition evaluated against the full context (data, metadata, temp_data)tasks- Actions to execute when the condition is met
Source§impl Workflow
impl Workflow
Sourcepub fn connector_refs(&self) -> impl Iterator<Item = ConnectorRef<'_>>
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.