flowing
Runs a multi-step procedure as a Python DAG, so ordering, branching and retries are enforced by the runner rather than described in prose a model can generate past. Use for "run these steps in order and retry the flaky one until the check passes", "build a pipeline that fetches,
Install
npx skills add https://github.com/oaustegard/claude-skills/tree/main/plugins/environment-and-config/skills/flowing
claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install oaustegard-claude-skills@llmmart
git clone https://github.com/oaustegard/claude-skills.git
The skills CLI installs just this skill, for any of its supported agents. Claude Code installs the whole oaustegard/claude-skills collection as a plugin from our marketplace. Git is the plain clone.
README
flowing
A lightweight DAG workflow runner for Claude's ephemeral containers. Declare steps, wire dependencies, run once — control flow lives in code, not in prose imperatives.
The problem it solves
Multi-step procedures are usually written as prose: "first fetch X, then validate Y, then if Z retry up to 3 times, otherwise skip ahead." An LLM reads and generates past prose like that — the gate is a suggestion, not a wall. Skipped validation, retries that never happen, branches taken on stale state.
A @task graph is structural instead. A step physically cannot run until its inputs are bound to its parameters. A gate that fires on missing or bad input can't be stepped over. The runner owns branching, retrying, validating, failure propagation, and parallelism; the LLM only supplies judgment at the leaves.
from flowing import task, Flow
@task
def fetch_data():
return {"items": [1, 2, 3]}
@task(depends_on=[fetch_data])
def process(fetch_data): # param name matches the dep's name
return sum(fetch_data["items"])
@task(depends_on=[process])
def store(process):
print(f"Result: {process}")
Flow(store).run() # topo-sorts into layers, parallel within a layer
Control-flow primitives
The distinctive part — branches and contracts as graph structure, not if statements buried in task bodies.
| Primitive | What it does | Use for |
|---|---|---|
when= |
Predicate over dep values; falsy → task SKIPPED, skip cascades to dependents | Branch selection in the topology |
validate= |
Checks dep values before the body runs; raise → FAILED with no retry | Enforceable input contracts between steps |
retry_until= |
Predicate over the return value; falsy → retry, consuming the retry= budget |
Self-correcting LLM steps (generate → check → regenerate) |
retry_until= is distinct from retry= alone: retry= only retries on a raised exception, retry_until= retries on output shape.
Also handles
- Parallel execution — independent tasks in a layer run on a thread pool.
- Resume —
run()→ fix →resume()re-runs from the failure point, keeping succeeded tasks cached.override()injects a corrected value for a step resolved out-of-band. - Detached side-effects —
detached=Truetasks (memory writes, notifications) run after the main DAG and never block it on failure. timeout_s=,retry=with exponential backoff,fail_fast=.
Layout
| File | Audience | Contents |
|---|---|---|
SKILL.md |
Claude | Trigger, mental model, quick start, the three primitives, when / when-not-to-use |
references/reference.md |
Claude | Full API — every @task parameter, Flow methods, resume/override, detached auto-discovery, signature gotchas |
scripts/flowing.py |
— | The runner itself (no third-party dependencies) |
tests/test_flowing.py |
— | 28 tests — python3 -m unittest tests.test_flowing |
CHANGELOG.md |
— | Version history |
When to reach for it
Use it when a procedure has branches that matter, steps with input contracts, an LLM step that needs to converge, 3+ operations that can parallelize, or a pipeline where late failures shouldn't waste early work.
Skip it for a single sequential operation (just call the function), for a next step that needs open-ended reasoning about the prior result rather than a predicate (use a think loop), or for async / distributed workflows (this is single-container, thread-pool based).
Complements
- orchestrating-agents — parallel API instances and delegated sub-tasks.
flowingorders and gates work within one container; orchestrating-agents fans work out across many. - tiling-tree — MECE partitioning of a problem space. Tiling-tree decides what the branches are;
flowingenforces the execution order once they exist. - tracking-todos — a human-legible checklist for loose, evolving work.
flowingis for procedures whose shape is known up front and worth encoding as a graph.
Skill manifest
NOT SUPERSEDED BY DYNAMIC WORKFLOWS — read first
Claude Code's dynamic workflows orchestrate subagents (separate contexts, fan-out to 16-concurrent / 1000-agent). This skill is a different primitive: single-context control flow over YOUR OWN tool calls, with durable side-effects and checkpoint resume. The workflows runtime explicitly cannot touch the filesystem or shell directly — its agents do the work and the script only coordinates them. Flowing is the inverse: the script does the work.
Use flowing for an in-context pipeline (3+ steps, branches, retries, validation, detached side-effects). Use a workflow when you need many subagents. They compose; they do not compete. Do not abandon flowing for a workflow — you would lose the durable side-effects and the cross-session checkpoint that hub-spoke depends on.
Flowing — Control Flow in Code, Not Prose
When a procedure needs 3+ steps with branches, retries, or contracts, encode it as a DAG of Python tasks instead of prose imperatives. Prose like "first X, then Y, then if Z retry 3×" is read and generated past. A @task graph is structural: a step physically cannot run until its inputs are bound, and gates that fire on bad inputs can't be skipped.
The runner owns control flow — branching, retrying, validating, propagating failures, parallelizing. You provide judgment at the leaves. Runner: scripts/flowing.py.
Quick Start
from flowing import task, Flow
@task
def fetch_data():
return {"items": [1, 2, 3]}
@task(depends_on=[fetch_data])
def process(fetch_data): # param name must match the dep's name
return sum(fetch_data["items"])
@task(depends_on=[process])
def store(process):
print(f"Result: {process}")
Flow(store).run() # topo-sorts, runs each layer, parallel within a layer
Each task receives its dependencies as kwargs named after them. Independent tasks in the same layer run in parallel.
Control-Flow Primitives
Encode branches and contracts as graph structure, not if statements inside task bodies.
when= — conditional gate
Run the task only if the predicate (over gathered dep values) is truthy. Falsy → SKIPPED, and the skip propagates to dependents.
@task(depends_on=[fetch], when=lambda fetch: fetch["needs_processing"])
def process(fetch):
return transform(fetch["payload"])
validate= — edge contract
Check gathered dep values before the body runs. Raise → FAILED with no retry (bad inputs don't fix themselves). Pass → proceed.
def must_have_items(fetch):
if not fetch.get("items"):
raise ValueError("fetch returned empty payload")
@task(depends_on=[fetch], validate=must_have_items)
def process(fetch):
return sum(fetch["items"])
retry_until= — predicate-driven loop
Run the body, then call retry_until(value). True → done. False → retry, consuming the retry= budget. Use for self-correcting LLM steps: generate, check, regenerate.
@task(retry=4, retry_until=lambda r: r["valid"])
def generate_until_valid():
candidate = llm_call(...)
return {"valid": passes_schema(candidate), "candidate": candidate}
Distinct from retry= alone, which only retries on a raised exception.
Other capabilities
- Parallel execution — independent tasks in a layer run on a thread pool (
max_workers=). detached=True— side-effect tasks (memory writes, notifications) that run after the main DAG and never block it on failure.- In-process resume —
flow.run()→ fix →flow.resume()re-runs from the failure point, keeping succeeded tasks cached in memory (same process only).flow.override(task, value)injects a corrected result. - Durable journal (
journal_path=) — opt-in content-addressed replay that survives container death.Flow(term, journal_path="/path/run.jsonl").run()appends each succeeded task's result to an append-only JSONL keyed by astep_key= SHA-256 over the task's bytecode + itswhen/validate/retry_untilbodies + its dependencies' keys (chained, so an upstream change propagates downstream). A laterrun()— even in a fresh container — replays the unchanged prefix from the journal and only executes tasks whose key is absent; editing a task body busts its key and re-runs it and its dependents, while cosmetic knobs (retry=,timeout_s=,name) do not. This is the cross-session checkpoint hub-spoke work relies on. Caveat: results are pickled, so non-picklable return values simply re-run; closure-captured values are not part of the key (only the task body's own code is). timeout_s=,retry=with exponential backoff,fail_fast=.
Read references/reference.md before using anything beyond the quick start and the three primitives above — it covers every @task parameter, the Flow methods, resume/override, detached auto-discovery, and the validate=/when= signature-matching gotcha.
When to use
- A procedure has branches that matter →
when=makes them structural. - Steps have input contracts →
validate=makes them enforceable. - An LLM step needs to converge →
retry_until=puts the check in the loop. - 3+ independent operations that can parallelize.
- Multi-step pipelines where late failures shouldn't waste early work.
- Side-effects that shouldn't block the critical path →
detached=True.
When NOT to use
- A single sequential operation — just call the function.
- The next step needs reasoning about the prior result that can't be a predicate — use a think loop.
- Async or distributed workflows — this is single-container, thread-pool based.
Authoring discipline
If you find yourself writing prose like "first call X, validate Y, then if Z retry up to 3 times" — that is a flowing graph. Refactor before shipping. Prose imperatives don't enforce; @task graphs do.
Files (claude-skills)
-
references
-
reference.md 5 KB
# Flowing — API Reference Full surface of `scripts/flowing.py`. The core workflow and the three control-flow primitives are in [SKILL.md](../SKILL.md); read this when you need anything beyond the quick start. ## `@task` decorator ```python @task( depends_on=[other_task], # TaskDefs this task consumes retry=2, # extra attempts on raised exception (total = 1 + retry) retry_backoff_base_ms=1000, # exponential backoff base retry_max_backoff_ms=30_000, # backoff ceiling timeout_s=60.0, # abort a hung body (see below) detached=True, # side-effect task; see "Detached tasks" name="custom_name", # override the TaskDef name (default: function name) when=lambda **deps: bool, # conditional gate — falsy skips the task validate=lambda **deps: None, # edge contract — raise to FAIL with no retry retry_until=lambda result: bool, # predicate loop — falsy retries the body ) def my_step(other_task): return result ``` ### Dependency naming A task body receives each dependency as a kwarg named after the dependency's TaskDef name — its function name, or the `name=` override. If the body has no parameter of that name and no `**kwargs`, `flow.run()` raises a clear `ValueError` at graph-build time rather than failing mid-run with a `TypeError`. ### `timeout_s` If set, the body runs in a one-shot worker thread; a call that overruns is aborted as a retryable `TimeoutError` and consumes the `retry=` budget like any other failure. Python can't kill the orphaned thread, so it keeps running until the container exits — fine for run-once ephemeral use, but not a hard cancel. ## `Flow` class ```python flow = Flow(terminal_task, max_workers=5, fail_fast=True) results = flow.run() # dict[str, StepResult] flow.summary() # human-readable status table flow.value(some_task) # return value of a succeeded task ``` `fail_fast` stops the *next* layer from starting once a task in the current layer fails. Siblings already running in parallel can't be killed and run to completion on their pool threads. ## Resume from failure ```python flow = Flow(terminal) results = flow.run() # step_3 fails flow.override(step_3, corrected_value) # inject a fix obtained out-of-band results = flow.resume() # step_1, step_2 stay cached; step_4+ runs ``` `flow.resume()` clears FAILED/SKIPPED tasks from results and re-runs them, keeping SUCCEEDED tasks cached. `flow.override(td, value)` manually injects a succeeded result. ## Detached tasks (non-blocking side-effects) ```python @task(depends_on=[create_issue], detached=True) def store_memory(create_issue): remember(create_issue["url"], ...) ``` Detached tasks run in topologically-sorted layers after the main DAG. Failures land in `flow.detached_failures` and never trigger `fail_fast`. All dependencies must be SUCCEEDED for a detached task to run. ### Auto-discovery A detached task whose `depends_on` are all reachable from the declared terminals is auto-discovered — you don't pass it as a terminal: ```python @task def main_step(): ... @task(depends_on=[main_step], detached=True) def store(main_step): ... Flow(main_step).run() # `store` runs automatically after main_step succeeds ``` Detached tasks may depend on other detached tasks; the chain is layered and runs in dependency order. Detached tasks whose deps are NOT reachable from any terminal are ignored — they belong to a different graph and should be terminals of their own `Flow`. ## `clear_registry()` Module-level function that empties the registry `@task` appends to. For run-once container use you never need it; call it between independent flows in the same process (tests, REPLs) so detached auto-discovery can't pull a stale task into an unrelated graph. ## `validate=` and `when=` signatures `validate=` and `when=` callables receive gathered dep values as kwargs *by dep name*, the same way task bodies do. A validator written for a specific dep: ```python def must_have_title(fetch_url_meta): if not fetch_url_meta.get("title"): raise ValueError("missing title") ``` works only on tasks whose dep is named `fetch_url_meta`. Reuse it on a task whose dep is named `fetch_bad_meta` and you get `TypeError: must_have_title() got an unexpected keyword argument 'fetch_bad_meta'` at validate time, surfacing as a confusing FAIL with the wrong reason. Two patterns to avoid the trap: ```python # A) Reusable: take **kwargs, look up by expected key def must_have_title(**kwargs): meta = next(iter(kwargs.values())) if not meta.get("title"): raise ValueError("missing title") # B) Factory: bind the dep name explicitly at task definition def must_have_title_of(dep_name): def v(**kwargs): if not kwargs[dep_name].get("title"): raise ValueError(f"{dep_name}: missing title") return v @task(depends_on=[fetch], validate=must_have_title_of("fetch")) def process(fetch): ... ```
-
-
scripts
-
flowing.py 30.9 KB
""" flowing — lightweight DAG workflow runner for Claude.ai containers. Declare steps. Wire dependencies. Run once. No think loops. Usage: from flowing import task, Flow @task def fetch_data(): return {"items": [1, 2, 3]} @task(depends_on=[fetch_data]) def process(fetch_data): return sum(fetch_data["items"]) @task(depends_on=[process]) def store(process): print(f"Result: {process}") Flow(store).run() Resume from failure: flow = Flow(terminal) results = flow.run() # step 3 fails flow.override(step_3, corrected_value) results = flow.resume() # runs step 4+ with corrected step 3 Detached side-effects: @task(depends_on=[create_issue], detached=True) def store_memory(create_issue): remember(create_issue["url"], ...) # Failure here does NOT block the main pipeline """ from __future__ import annotations import base64 import hashlib import inspect import json import pickle import sys import time from collections.abc import Callable from concurrent.futures import ( ThreadPoolExecutor, as_completed, ) from concurrent.futures import ( TimeoutError as FuturesTimeoutError, ) from dataclasses import dataclass, field from enum import Enum from typing import Any class StepState(Enum): # _run_step is synchronous and only ever returns a terminal state, so # PENDING/RUNNING/RETRYING would never be observable — they are omitted # rather than defined-but-never-assigned. SUCCEEDED = "succeeded" FAILED = "failed" SKIPPED = "skipped" @dataclass class StepResult: name: str state: StepState value: Any = None error: Exception | None = None duration_ms: float = 0 attempts: int = 0 step_key: str | None = None # content-addressed key (set only when journaling) cached: bool = False # True if this result was replayed from the journal @dataclass class TaskDef: """A declared workflow step.""" name: str fn: Callable depends_on: list[TaskDef] = field(default_factory=list) retry: int = 0 retry_backoff_base_ms: int = 1000 retry_max_backoff_ms: int = 30_000 timeout_s: float | None = None detached: bool = False # v1.1 control-flow primitives when: Callable | None = None # gate: receives gathered kwargs, returns bool. False -> SKIPPED validate: Callable | None = None # edge contract: receives gathered kwargs, raises on bad inputs. Raise -> FAILED, no retry retry_until: Callable | None = None # predicate loop: receives task return value, returns bool. False -> retry (uses retry= budget) def __call__(self, *args, **kwargs): return self.fn(*args, **kwargs) def __hash__(self): return id(self) def __eq__(self, other): return self is other # Module-level registry of all TaskDefs created via the @task decorator. # Used by Flow._collect_tasks to auto-discover detached tasks whose dependencies # are reachable from declared terminals (v1.2.0). Without this, a detached task # downstream of a terminal would silently never run, since the dep walk only # traverses depends_on backward. _TASK_REGISTRY: list[TaskDef] = [] def clear_registry() -> None: """Empty the module-level task registry. The registry accumulates every TaskDef created via @task for the life of the process. For run-once container use that is fine. Call this between independent flows in the same process (tests, REPLs) so detached auto-discovery can't pull a stale task into an unrelated graph. """ _TASK_REGISTRY.clear() # @lat: [[orchestration#DAG Workflow Runner]] def task( fn: Callable | None = None, *, depends_on: list[TaskDef] | None = None, retry: int = 0, retry_backoff_base_ms: int = 1000, retry_max_backoff_ms: int = 30_000, timeout_s: float | None = None, name: str | None = None, detached: bool = False, when: Callable | None = None, validate: Callable | None = None, retry_until: Callable | None = None, ) -> TaskDef: def wrap(f: Callable) -> TaskDef: td = TaskDef( name=name or f.__name__, fn=f, depends_on=depends_on or [], retry=retry, retry_backoff_base_ms=retry_backoff_base_ms, retry_max_backoff_ms=retry_max_backoff_ms, timeout_s=timeout_s, detached=detached, when=when, validate=validate, retry_until=retry_until, ) _TASK_REGISTRY.append(td) return td if fn is not None: return wrap(fn) return wrap def _log(msg: str, **kw): parts = [f"[flow] {msg}"] for k, v in kw.items(): parts.append(f"{k}={v}") print(" ".join(parts), file=sys.stderr, flush=True) def _fn_fingerprint(fn: Callable | None) -> str: """Stable-ish fingerprint of a callable's behaviour. Prefers bytecode (co_code + stringified co_consts), which changes when the body changes but is stable across re-imports of unchanged source. Falls back to source text, then to a qualified-name repr. Never raises. """ if fn is None: return "" try: code = fn.__code__ return "bc:" + hashlib.sha256( code.co_code + repr(code.co_consts).encode("utf-8", "replace") ).hexdigest() except Exception: pass try: return "src:" + hashlib.sha256( inspect.getsource(fn).encode("utf-8", "replace") ).hexdigest() except Exception: return "id:" + getattr(fn, "__qualname__", repr(fn)) STEP_KEY_VERSION = 1 def compute_step_key(td: TaskDef, dep_keys: list[str]) -> str: """Content-addressed, chained key for one task. Hashes, in order: the version, the sorted dependency keys (so an upstream change propagates through the chain — pi-flow's prefix sensitivity), and the fingerprints of the task body plus its when/validate/retry_until control callables. The result is `v<version>:<sha256>`. What is deliberately NOT in the key: name, retry/backoff/timeout knobs, detached flag. Tuning a reliability knob or renaming must not bust the cache; changing behaviour must. """ h = hashlib.sha256() h.update(f"v{STEP_KEY_VERSION}\n".encode()) for k in sorted(dep_keys): h.update(b"dep:") h.update(k.encode()) h.update(b"\n") for label, f in (("fn", td.fn), ("when", td.when), ("validate", td.validate), ("retry_until", td.retry_until)): h.update(label.encode()) h.update(b"=") h.update(_fn_fingerprint(f).encode()) h.update(b"\n") return f"v{STEP_KEY_VERSION}:{h.hexdigest()}" def _run_step(td: TaskDef, results: dict[str, StepResult]) -> StepResult: kwargs = {} for dep in td.depends_on: r = results[dep.name] if r.state != StepState.SUCCEEDED: _log(f"SKIP {td.name}", reason=f"dependency {dep.name} did not succeed (state={r.state.value})") return StepResult(name=td.name, state=StepState.SKIPPED, attempts=0) kwargs[dep.name] = r.value # GATE 1: when() — conditional skip. Truthy -> proceed; falsy -> SKIPPED (downstream cascades). if td.when is not None: try: if not td.when(**kwargs): _log(f"SKIP {td.name}", reason="when() returned False") return StepResult(name=td.name, state=StepState.SKIPPED, attempts=0) except Exception as e: _log(f"FAIL {td.name}", reason=f"when() raised: {str(e)[:120]}") return StepResult(name=td.name, state=StepState.FAILED, error=e, attempts=0) # GATE 2: validate() — edge contract. Raise -> FAILED with NO retry (bad inputs won't fix themselves). if td.validate is not None: try: td.validate(**kwargs) except Exception as e: _log(f"FAIL {td.name}", reason=f"validate() raised: {str(e)[:120]}") return StepResult(name=td.name, state=StepState.FAILED, error=e, attempts=0) max_attempts = 1 + td.retry last_error = None last_value = None last_dur = 0.0 for attempt in range(1, max_attempts + 1): _log(f"{'RUN' if attempt == 1 else 'RETRY'} {td.name}", attempt=f"{attempt}/{max_attempts}") t0 = time.monotonic() try: if td.timeout_s is not None: # Run the body in a one-shot worker so a hung call can't stall # the whole flow. Python can't kill a running thread, so on # timeout the orphaned worker keeps going until the container # exits — acceptable for run-once ephemeral use. shutdown is # wait=False precisely so we don't re-block on that orphan. one_shot = ThreadPoolExecutor(max_workers=1) fut = one_shot.submit(td.fn, **kwargs) try: value = fut.result(timeout=td.timeout_s) except FuturesTimeoutError: raise TimeoutError( f"{td.name} exceeded timeout_s={td.timeout_s}" ) from None finally: one_shot.shutdown(wait=False) else: value = td.fn(**kwargs) last_dur = (time.monotonic() - t0) * 1000 # GATE 3: retry_until() — predicate-driven loop. True -> done; False -> retry (consumes retry budget). if td.retry_until is not None: try: ok = td.retry_until(value) except Exception as e: _log(f"FAIL {td.name}", reason=f"retry_until() raised: {str(e)[:120]}") return StepResult( name=td.name, state=StepState.FAILED, error=e, value=value, duration_ms=last_dur, attempts=attempt, ) if not ok: last_value = value last_error = ValueError( f"retry_until predicate returned False (attempt {attempt}/{max_attempts})" ) _log(f"PREDICATE_FAIL {td.name}", ms=f"{last_dur:.0f}", attempt=f"{attempt}/{max_attempts}") if attempt < max_attempts: delay_ms = min( td.retry_backoff_base_ms * (2 ** (attempt - 1)), td.retry_max_backoff_ms, ) time.sleep(delay_ms / 1000) continue # Out of attempts; predicate still False -> FAILED with last value preserved return StepResult( name=td.name, state=StepState.FAILED, error=last_error, value=last_value, duration_ms=last_dur, attempts=attempt, ) _log(f"OK {td.name}", ms=f"{last_dur:.0f}") return StepResult( name=td.name, state=StepState.SUCCEEDED, value=value, duration_ms=last_dur, attempts=attempt, ) except Exception as e: last_dur = (time.monotonic() - t0) * 1000 last_error = e _log(f"FAIL {td.name}", ms=f"{last_dur:.0f}", error=str(e)[:120], attempt=f"{attempt}/{max_attempts}") if attempt < max_attempts: delay_ms = min( td.retry_backoff_base_ms * (2 ** (attempt - 1)), td.retry_max_backoff_ms, ) time.sleep(delay_ms / 1000) return StepResult( name=td.name, state=StepState.FAILED, error=last_error, duration_ms=last_dur, attempts=max_attempts, ) # @lat: [[orchestration#DAG Workflow Runner]] class Flow: def __init__(self, *terminals: TaskDef, max_workers: int = 5, fail_fast: bool = True, journal_path: str | None = None): if not terminals: raise ValueError("Flow requires at least one terminal task") self.terminals = list(terminals) self.max_workers = max_workers self.fail_fast = fail_fast self.results: dict[str, StepResult] = {} self.detached_failures: list[StepResult] = [] # Content-addressed durable journal (opt-in). When set, run() replays # the unchanged prefix from disk and only executes tasks whose # step_key is absent — surviving container death and re-running any # task whose body changed (and, via key chaining, its dependents). self.journal_path = journal_path self._journal_completed: dict[str, Any] = {} # step_key -> value self._step_keys: dict[TaskDef, str] = {} def _load_journal(self) -> None: """Load completed step_keys -> values from the JSONL journal. Append-only, tolerant of a truncated trailing line (crash mid-write). A value that fails to unpickle is dropped, so its task simply re-runs. """ self._journal_completed = {} if not self.journal_path: return try: with open(self.journal_path, "r", encoding="utf-8") as fh: lines = fh.readlines() except FileNotFoundError: return for line in lines: line = line.strip() if not line: continue try: entry = json.loads(line) except json.JSONDecodeError: continue # truncated final line — discard if entry.get("kind") != "Completed": continue try: val = pickle.loads(base64.b64decode(entry["value"])) except Exception: continue # unpicklable/corrupt -> treat as not-cached self._journal_completed[entry["key"]] = val def _compute_keys(self, tasks: set[TaskDef]) -> None: """Compute chained step_keys for all tasks in dependency order.""" self._step_keys = {} remaining = set(tasks) # Iterate to fixed point: a task is keyable once all its in-graph deps # have keys. Detached-or-not, deps may sit in another partition. while remaining: progressed = False for t in list(remaining): dep_keys = [] ready = True for dep in t.depends_on: if dep in tasks: if dep not in self._step_keys: ready = False break dep_keys.append(self._step_keys[dep]) if not ready: continue self._step_keys[t] = compute_step_key(t, dep_keys) remaining.discard(t) progressed = True if not progressed: # Dependency outside `tasks` (shouldn't happen) — key on body only. for t in list(remaining): self._step_keys[t] = compute_step_key(t, []) remaining.discard(t) def _cached_result(self, td: TaskDef) -> StepResult | None: """Return a replayed SUCCEEDED result if td's step_key is journaled.""" if not self.journal_path: return None key = self._step_keys.get(td) if key is not None and key in self._journal_completed: return StepResult( name=td.name, state=StepState.SUCCEEDED, value=self._journal_completed[key], step_key=key, cached=True, ) return None def _journal_append(self, result: StepResult) -> None: """Append a Completed entry for a freshly-succeeded, non-cached task.""" if not self.journal_path or result.cached: return if result.state != StepState.SUCCEEDED or result.step_key is None: return try: blob = base64.b64encode(pickle.dumps(result.value)).decode("ascii") except Exception as e: _log(f"JOURNAL_SKIP {result.name}", reason=f"unpicklable value: {str(e)[:80]}") return line = json.dumps({"kind": "Completed", "key": result.step_key, "name": result.name, "value": blob}) with open(self.journal_path, "a", encoding="utf-8") as fh: fh.write(line + "\n") def _collect_tasks(self) -> tuple[set[TaskDef], set[TaskDef]]: """Collect all tasks, separating main DAG from detached tasks. v1.2.0: Auto-discovers detached tasks whose dependencies are all reachable from declared terminals. Iterates to fixed point so that chains of detached tasks (detachB -> detachA -> main) are picked up. Detached tasks whose deps are NOT reachable are ignored — they belong to a different graph. """ all_tasks: set[TaskDef] = set() stack = list(self.terminals) while stack: t = stack.pop() if t not in all_tasks: all_tasks.add(t) stack.extend(t.depends_on) # Fixed-point pull-in of detached tasks from the module registry whose # deps are all already reachable. Repeat until no new task is added so # detached-on-detached chains resolve correctly. while True: added = False for t in _TASK_REGISTRY: if t.detached and t not in all_tasks: if t.depends_on and all(dep in all_tasks for dep in t.depends_on): all_tasks.add(t) added = True if not added: break main_tasks = {t for t in all_tasks if not t.detached} detached_tasks = {t for t in all_tasks if t.detached} return main_tasks, detached_tasks def _build_layers(self, tasks: set[TaskDef]) -> list[list[TaskDef]]: """Topological sort into parallel execution layers.""" in_degree: dict[TaskDef, int] = {t: 0 for t in tasks} dependents: dict[TaskDef, list[TaskDef]] = {t: [] for t in tasks} for t in tasks: for dep in t.depends_on: if dep in tasks: in_degree[t] += 1 dependents[dep].append(t) layers: list[list[TaskDef]] = [] current = [t for t, d in in_degree.items() if d == 0] visited = 0 while current: layers.append(current) visited += len(current) next_layer = [] for t in current: for dep in dependents[t]: in_degree[dep] -= 1 if in_degree[dep] == 0: next_layer.append(dep) current = next_layer if visited != len(tasks): raise ValueError("Cycle detected in task graph") return layers def _validate_signatures(self, tasks: set[TaskDef]) -> None: """Fail at graph-build time if a task body can't receive a dep by name. Deps are passed to the body as kwargs keyed by the producer's TaskDef name. If the consumer's signature has no matching parameter (and no **kwargs), the run would otherwise die mid-flight with a confusing TypeError. Catch it here with a message that names both ends. """ for t in tasks: if not t.depends_on: continue try: params = inspect.signature(t.fn).parameters except (ValueError, TypeError): continue # builtins / C functions — can't introspect, skip if any(p.kind == inspect.Parameter.VAR_KEYWORD for p in params.values()): continue # **kwargs absorbs anything accepted = { name for name, p in params.items() if p.kind in (inspect.Parameter.POSITIONAL_OR_KEYWORD, inspect.Parameter.KEYWORD_ONLY) } for dep in t.depends_on: if dep.name not in accepted: raise ValueError( f"Task '{t.name}' depends on '{dep.name}', but its " f"function has no parameter named '{dep.name}'. " f"Rename the parameter to match, add **kwargs, or set " f"@task(name=...) on the dependency." ) def _execute(self, layers: list[list[TaskDef]], skip_succeeded: bool = False) -> None: """Execute layers. If skip_succeeded=True, skip tasks already SUCCEEDED in self.results.""" for layer_idx, layer in enumerate(layers): # Filter out already-succeeded tasks when resuming if skip_succeeded: layer = [t for t in layer if not ( t.name in self.results and self.results[t.name].state == StepState.SUCCEEDED )] if not layer: continue # Journal replay: tasks whose step_key is already Completed return # their cached value without re-running. A miss falls through to a # live run; key chaining means an edited task's dependents miss too. if self.journal_path: cached_hits = [] live = [] for t in layer: cr = self._cached_result(t) if cr is not None: self.results[t.name] = cr cached_hits.append(t.name) else: live.append(t) if cached_hits: _log("CACHED", tasks=",".join(cached_hits)) layer = live if not layer: continue parallel = len(layer) > 1 _log(f"LAYER {layer_idx}", tasks=",".join(t.name for t in layer), parallel=parallel) if parallel: with ThreadPoolExecutor(max_workers=self.max_workers) as pool: futures = { pool.submit(_run_step, td, self.results): td for td in layer } for future in as_completed(futures): td = futures[future] result = future.result() if self.journal_path: result.step_key = self._step_keys.get(td) self.results[result.name] = result self._journal_append(result) if self.fail_fast and result.state == StepState.FAILED: _log(f"FAIL_FAST triggered by {result.name}") # Cancels queued-but-unstarted siblings. Siblings # already running can't be killed — they finish on # their pool threads — but fail_fast's real job is # done: the next layer won't start. pool.shutdown(wait=False, cancel_futures=True) break else: for td in layer: result = _run_step(td, self.results) if self.journal_path: result.step_key = self._step_keys.get(td) self.results[result.name] = result self._journal_append(result) if self.fail_fast and result.state == StepState.FAILED: _log(f"FAIL_FAST triggered by {result.name}") break if self.fail_fast: failed = [r for r in self.results.values() if r.state == StepState.FAILED] if failed: break def _execute_detached(self, detached_tasks: set[TaskDef]) -> None: """Run detached tasks in topologically-sorted layers after the main DAG. Failures collected in self.detached_failures, not propagated. v1.2.0: Layered execution (was single-layer). A detached task may depend on another detached task; layering ensures the dependency runs first and its result is available. """ if not detached_tasks: return detached_layers = self._build_layers(detached_tasks) for layer in detached_layers: # Filter to tasks whose deps (main + earlier detached layers) all succeeded runnable = [] for t in layer: deps_ok = all( t_dep.name in self.results and self.results[t_dep.name].state == StepState.SUCCEEDED for t_dep in t.depends_on ) if deps_ok: runnable.append(t) else: skip_result = StepResult( name=t.name, state=StepState.SKIPPED, attempts=0) self.results[t.name] = skip_result # Journal replay for detached side-effects: a cached hit means the # side-effect already fired on a prior run — don't re-fire it. if self.journal_path and runnable: still = [] cached_hits = [] for t in runnable: cr = self._cached_result(t) if cr is not None: self.results[t.name] = cr cached_hits.append(t.name) else: still.append(t) if cached_hits: _log("CACHED_DETACHED", tasks=",".join(cached_hits)) runnable = still if not runnable: continue _log("DETACHED", tasks=",".join(t.name for t in runnable)) if len(runnable) > 1: with ThreadPoolExecutor(max_workers=self.max_workers) as pool: futures = { pool.submit(_run_step, td, self.results): td for td in runnable } for future in as_completed(futures): td = futures[future] result = future.result() if self.journal_path: result.step_key = self._step_keys.get(td) self.results[result.name] = result self._journal_append(result) if result.state == StepState.FAILED: self.detached_failures.append(result) else: for td in runnable: result = _run_step(td, self.results) if self.journal_path: result.step_key = self._step_keys.get(td) self.results[result.name] = result self._journal_append(result) if result.state == StepState.FAILED: self.detached_failures.append(result) def run(self) -> dict[str, StepResult]: """Execute the full DAG from scratch.""" self.results = {} self.detached_failures = [] main_tasks, detached_tasks = self._collect_tasks() self._validate_signatures(main_tasks | detached_tasks) if self.journal_path: self._compute_keys(main_tasks | detached_tasks) self._load_journal() main_layers = self._build_layers(main_tasks) total_tasks = sum(len(l) for l in main_layers) + len(detached_tasks) _log("START", tasks=total_tasks, layers=len(main_layers), terminals=",".join(t.name for t in self.terminals), detached=len(detached_tasks)) flow_t0 = time.monotonic() # Execute main DAG self._execute(main_layers) # Execute detached tasks (only if main DAG didn't fail, or their deps succeeded) self._execute_detached(detached_tasks) flow_dur = (time.monotonic() - flow_t0) * 1000 succeeded = sum(1 for r in self.results.values() if r.state == StepState.SUCCEEDED) failed = sum(1 for r in self.results.values() if r.state == StepState.FAILED) skipped = sum(1 for r in self.results.values() if r.state == StepState.SKIPPED) _log("DONE", ms=f"{flow_dur:.0f}", succeeded=succeeded, failed=failed, skipped=skipped) return self.results def resume(self) -> dict[str, StepResult]: """Re-run from failure point. SUCCEEDED tasks keep their cached values. FAILED and SKIPPED tasks are cleared from results and re-execute.""" self.detached_failures = [] # Reset FAILED and SKIPPED tasks to_reset = [name for name, r in self.results.items() if r.state in (StepState.FAILED, StepState.SKIPPED)] for name in to_reset: del self.results[name] main_tasks, detached_tasks = self._collect_tasks() if self.journal_path: self._compute_keys(main_tasks | detached_tasks) self._load_journal() main_layers = self._build_layers(main_tasks) cached = sum(1 for r in self.results.values() if r.state == StepState.SUCCEEDED) _log("RESUME", cached=cached, reset=len(to_reset)) flow_t0 = time.monotonic() # Execute with skip_succeeded=True self._execute(main_layers, skip_succeeded=True) # Re-run detached tasks (they may have been skipped/failed before) self._execute_detached(detached_tasks) flow_dur = (time.monotonic() - flow_t0) * 1000 succeeded = sum(1 for r in self.results.values() if r.state == StepState.SUCCEEDED) failed = sum(1 for r in self.results.values() if r.state == StepState.FAILED) skipped = sum(1 for r in self.results.values() if r.state == StepState.SKIPPED) _log("DONE", ms=f"{flow_dur:.0f}", succeeded=succeeded, failed=failed, skipped=skipped) return self.results def override(self, td: TaskDef, value: Any) -> None: """Manually set a result without re-running the task. Use when you have fixed the problem externally and have the correct value.""" self.results[td.name] = StepResult( name=td.name, state=StepState.SUCCEEDED, value=value, duration_ms=0, attempts=0, ) def value(self, td: TaskDef) -> Any: r = self.results.get(td.name) if r is None: raise KeyError(f"Task {td.name} not found in results") if r.state != StepState.SUCCEEDED: raise RuntimeError(f"Task {td.name} did not succeed (state={r.state})") return r.value def summary(self) -> str: # Build the detached-name set once — including auto-discovered # detached tasks, which a terminal-only walk would miss. _, detached_tasks = self._collect_tasks() detached_names = {t.name for t in detached_tasks} lines = [] for name, r in self.results.items(): status = r.state.value.upper() dur = f"{r.duration_ms:.0f}ms" att = f"x{r.attempts}" if r.attempts > 1 else "" err = f" err={r.error}" if r.error else "" det = " [detached]" if name in detached_names else "" lines.append(f" {status:9s} {name} ({dur}{att}){det}{err}") return "\n".join(lines)
-
-
tests
-
test_flowing.py 19.6 KB
"""Tests for flowing control-flow primitives. Coverage: - v1.0 backward compat (existing DAG, retry, override, resume, detached) - v1.1: when= conditional gate - v1.1: validate= edge contract - v1.1: retry_until= predicate loop - Composition of new primitives - v1.3: timeout_s enforcement, signature validation, clear_registry """ import os import sys import time import unittest sys.path.insert(0, os.path.join(os.path.dirname(__file__), "..", "scripts")) from flowing import Flow, StepState, clear_registry, task class TestBackwardCompat(unittest.TestCase): """v1.0 features must keep working unchanged.""" def test_simple_chain(self): @task def a(): return 1 @task(depends_on=[a]) def b(a): return a + 1 @task(depends_on=[b]) def c(b): return b * 10 flow = Flow(c) flow.run() self.assertEqual(flow.value(c), 20) def test_retry_on_exception(self): calls = {"n": 0} @task(retry=2, retry_backoff_base_ms=1) def flaky(): calls["n"] += 1 if calls["n"] < 3: raise RuntimeError("nope") return "ok" flow = Flow(flaky) flow.run() self.assertEqual(flow.value(flaky), "ok") self.assertEqual(calls["n"], 3) def test_failure_propagates_skip(self): @task def fails(): raise RuntimeError("boom") @task(depends_on=[fails]) def downstream(fails): return fails + 1 flow = Flow(downstream, fail_fast=False) flow.run() self.assertEqual(flow.results["fails"].state, StepState.FAILED) self.assertEqual(flow.results["downstream"].state, StepState.SKIPPED) def test_override_and_resume(self): calls = {"a": 0, "b": 0} @task def a(): calls["a"] += 1 return 5 @task(depends_on=[a]) def b(a): calls["b"] += 1 if calls["b"] == 1: raise RuntimeError("first call fails") return a * 2 flow = Flow(b) flow.run() self.assertEqual(flow.results["b"].state, StepState.FAILED) flow.resume() self.assertEqual(flow.value(b), 10) self.assertEqual(calls["a"], 1, "succeeded task should not re-run") class TestWhenGate(unittest.TestCase): """when= conditional skip — falsy returns mark task SKIPPED.""" def test_when_true_runs(self): @task def upstream(): return {"ready": True} @task(depends_on=[upstream], when=lambda upstream: upstream["ready"]) def gated(upstream): return "ran" flow = Flow(gated) flow.run() self.assertEqual(flow.results["gated"].state, StepState.SUCCEEDED) self.assertEqual(flow.value(gated), "ran") def test_when_false_skips(self): ran = {"flag": False} @task def upstream(): return {"ready": False} @task(depends_on=[upstream], when=lambda upstream: upstream["ready"]) def gated(upstream): ran["flag"] = True return "should not run" flow = Flow(gated) flow.run() self.assertEqual(flow.results["gated"].state, StepState.SKIPPED) self.assertFalse(ran["flag"], "task body must not execute when when() is False") def test_when_skip_propagates_to_dependents(self): @task def upstream(): return {"ready": False} @task(depends_on=[upstream], when=lambda upstream: upstream["ready"]) def middle(upstream): return "middle" @task(depends_on=[middle]) def downstream(middle): return middle + " + downstream" flow = Flow(downstream) flow.run() self.assertEqual(flow.results["middle"].state, StepState.SKIPPED) self.assertEqual(flow.results["downstream"].state, StepState.SKIPPED) def test_when_raises_fails(self): @task def upstream(): return {"ready": True} @task(depends_on=[upstream], when=lambda upstream: 1 / 0) def gated(upstream): return "unreachable" flow = Flow(gated, fail_fast=False) flow.run() self.assertEqual(flow.results["gated"].state, StepState.FAILED) self.assertIsInstance(flow.results["gated"].error, ZeroDivisionError) class TestValidateGate(unittest.TestCase): """validate= edge contract — raise marks task FAILED with no retry.""" def test_validate_passes(self): def must_be_dict(upstream): assert isinstance(upstream, dict), "upstream must be dict" @task def upstream(): return {"k": "v"} @task(depends_on=[upstream], validate=must_be_dict) def consumer(upstream): return upstream["k"] flow = Flow(consumer) flow.run() self.assertEqual(flow.value(consumer), "v") def test_validate_fails_no_retry(self): body_calls = {"n": 0} def reject(upstream): raise ValueError(f"bad input: {upstream}") @task def upstream(): return "wrong shape" @task(depends_on=[upstream], validate=reject, retry=5, retry_backoff_base_ms=1) def consumer(upstream): body_calls["n"] += 1 return "should not run" flow = Flow(consumer, fail_fast=False) flow.run() self.assertEqual(flow.results["consumer"].state, StepState.FAILED) self.assertEqual(body_calls["n"], 0, "validate failure must not run task body") self.assertEqual(flow.results["consumer"].attempts, 0, "validate failure must not consume retry budget") self.assertIsInstance(flow.results["consumer"].error, ValueError) class TestRetryUntil(unittest.TestCase): """retry_until= predicate-driven loop — re-runs body until predicate(value) is True.""" def test_retry_until_succeeds_after_n(self): calls = {"n": 0} @task(retry=5, retry_backoff_base_ms=1, retry_until=lambda v: v["valid"]) def converging(): calls["n"] += 1 return {"valid": calls["n"] >= 3, "attempt": calls["n"]} flow = Flow(converging) flow.run() self.assertEqual(flow.results["converging"].state, StepState.SUCCEEDED) self.assertEqual(flow.value(converging)["attempt"], 3) self.assertEqual(flow.results["converging"].attempts, 3) def test_retry_until_exhausts(self): calls = {"n": 0} @task(retry=2, retry_backoff_base_ms=1, retry_until=lambda v: False) def never_satisfies(): calls["n"] += 1 return {"attempt": calls["n"]} flow = Flow(never_satisfies, fail_fast=False) flow.run() r = flow.results["never_satisfies"] self.assertEqual(r.state, StepState.FAILED) self.assertEqual(r.attempts, 3, "should consume full retry budget (1 + retry)") self.assertEqual(calls["n"], 3) # last value preserved on FAILED for diagnostics self.assertEqual(r.value, {"attempt": 3}) def test_retry_until_first_attempt_pass(self): calls = {"n": 0} @task(retry=5, retry_backoff_base_ms=1, retry_until=lambda v: v == "ok") def immediate(): calls["n"] += 1 return "ok" flow = Flow(immediate) flow.run() self.assertEqual(flow.value(immediate), "ok") self.assertEqual(calls["n"], 1, "should not retry when predicate passes first time") def test_retry_until_predicate_raises(self): @task(retry=5, retry_backoff_base_ms=1, retry_until=lambda v: 1 / 0) def victim(): return "value" flow = Flow(victim, fail_fast=False) flow.run() r = flow.results["victim"] self.assertEqual(r.state, StepState.FAILED) self.assertIsInstance(r.error, ZeroDivisionError) self.assertEqual(r.value, "value", "value preserved when predicate itself raises") class TestComposition(unittest.TestCase): """The three primitives compose.""" def test_when_then_validate_then_retry_until(self): # when=True -> proceed; validate passes; retry_until succeeds 2nd attempt body_calls = {"n": 0} @task def upstream(): return {"go": True, "input": [1, 2, 3]} @task( depends_on=[upstream], when=lambda upstream: upstream["go"], validate=lambda upstream: ( None if isinstance(upstream["input"], list) else (_ for _ in ()).throw(ValueError("bad input")) ), retry=3, retry_backoff_base_ms=1, retry_until=lambda v: v["good"], ) def converging(upstream): body_calls["n"] += 1 return {"good": body_calls["n"] >= 2, "n": body_calls["n"]} flow = Flow(converging) flow.run() self.assertEqual(flow.results["converging"].state, StepState.SUCCEEDED) self.assertEqual(flow.value(converging)["n"], 2) class TestDetachedAutoDiscovery(unittest.TestCase): """v1.2: detached tasks whose deps are reachable from declared terminals should be auto-discovered and run, without needing to be passed as terminals. Regression: in v1.1.1, `Flow(assemble)` with a `store_memory(detached=True, depends_on=[assemble])` defined elsewhere would silently NOT run store_memory. The user had to know to pass it: `Flow(assemble, store_memory)`. The SKILL.md "Run in a final layer after the main DAG" implied auto-discovery. """ def test_detached_downstream_of_terminal_auto_discovered(self): """Detached task whose dep IS the terminal should run automatically.""" side_effects = [] @task def main_step(): return "result" @task(depends_on=[main_step], detached=True) def store(main_step): side_effects.append(("stored", main_step)) return "stored" flow = Flow(main_step) # NOTE: store NOT passed as terminal results = flow.run() self.assertEqual(results["main_step"].state, StepState.SUCCEEDED) self.assertIn("store", results) self.assertEqual(results["store"].state, StepState.SUCCEEDED) self.assertEqual(side_effects, [("stored", "result")]) def test_detached_passed_as_terminal_still_works(self): """Backward compat: explicitly passing detached as terminal still works.""" side_effects = [] @task def main_step(): return "result" @task(depends_on=[main_step], detached=True) def store(main_step): side_effects.append(main_step) flow = Flow(main_step, store) flow.run() self.assertEqual(side_effects, ["result"]) def test_detached_chain_auto_discovered(self): """Detached-on-detached: if detachB depends on detachA depends on main, both should be auto-discovered when only main is the terminal.""" order = [] @task def main_step(): order.append("main") return 1 @task(depends_on=[main_step], detached=True) def detach_a(main_step): order.append("a") return main_step + 1 @task(depends_on=[detach_a], detached=True) def detach_b(detach_a): order.append("b") return detach_a + 1 flow = Flow(main_step) results = flow.run() self.assertEqual(results["main_step"].state, StepState.SUCCEEDED) self.assertEqual(results["detach_a"].state, StepState.SUCCEEDED) self.assertEqual(results["detach_b"].state, StepState.SUCCEEDED) self.assertEqual(order, ["main", "a", "b"]) def test_unrelated_detached_not_picked_up(self): """Detached tasks whose deps are NOT reachable from terminals should NOT run — they belong to a different graph.""" side_effects = [] @task def graph_a_step(): return "a" @task def graph_b_step(): return "b" @task(depends_on=[graph_b_step], detached=True) def graph_b_side(graph_b_step): side_effects.append(graph_b_step) flow = Flow(graph_a_step) # only graph A results = flow.run() self.assertEqual(results["graph_a_step"].state, StepState.SUCCEEDED) self.assertNotIn("graph_b_step", results) self.assertNotIn("graph_b_side", results) self.assertEqual(side_effects, []) def test_detached_failure_still_isolated(self): """Auto-discovered detached failure must still NOT trigger fail_fast and must land in flow.detached_failures (existing detached semantics).""" @task def main_step(): return "ok" @task(depends_on=[main_step], detached=True) def flaky_side(main_step): raise RuntimeError("side-effect blew up") flow = Flow(main_step) results = flow.run() self.assertEqual(results["main_step"].state, StepState.SUCCEEDED) self.assertEqual(results["flaky_side"].state, StepState.FAILED) self.assertEqual(len(flow.detached_failures), 1) self.assertEqual(flow.detached_failures[0].name, "flaky_side") class TestTimeout(unittest.TestCase): """timeout_s aborts a hung body and is retryable (v1.3).""" def test_timeout_fails_slow_task(self): @task(timeout_s=0.05) def slow(): time.sleep(0.5) return "done" flow = Flow(slow, fail_fast=False) flow.run() r = flow.results["slow"] self.assertEqual(r.state, StepState.FAILED) self.assertIsInstance(r.error, TimeoutError) def test_timeout_consumes_retry_then_succeeds(self): calls = {"n": 0} @task(timeout_s=0.05, retry=2, retry_backoff_base_ms=1) def sometimes_slow(): calls["n"] += 1 if calls["n"] == 1: time.sleep(0.5) # first attempt times out return calls["n"] flow = Flow(sometimes_slow) flow.run() self.assertEqual(flow.results["sometimes_slow"].state, StepState.SUCCEEDED) self.assertEqual(flow.value(sometimes_slow), 2) def test_no_timeout_runs_normally(self): @task def quick(): return "ok" flow = Flow(quick) flow.run() self.assertEqual(flow.value(quick), "ok") class TestSignatureValidation(unittest.TestCase): """Mismatched dep/parameter names are caught at graph-build time (v1.3).""" def test_mismatched_param_raises(self): @task def producer(): return 1 @task(depends_on=[producer]) def consumer(wrong_name): return wrong_name + 1 flow = Flow(consumer) with self.assertRaises(ValueError) as ctx: flow.run() msg = str(ctx.exception) self.assertIn("producer", msg) self.assertIn("consumer", msg) def test_kwargs_absorbs_any_dep(self): @task def producer(): return 7 @task(depends_on=[producer]) def consumer(**kwargs): return kwargs["producer"] flow = Flow(consumer) flow.run() self.assertEqual(flow.value(consumer), 7) def test_name_override_matches_param(self): @task(name="aliased") def producer(): return 3 @task(depends_on=[producer]) def consumer(aliased): return aliased * 2 flow = Flow(consumer) flow.run() self.assertEqual(flow.value(consumer), 6) class TestClearRegistry(unittest.TestCase): """clear_registry empties the module-level task registry (v1.3).""" def test_clear_registry_drops_detached_candidates(self): side_effects = [] @task def shared(): return 1 @task(depends_on=[shared], detached=True) def stale_side(shared): side_effects.append("ran") # Without clearing, stale_side would auto-discover onto any later # flow built on `shared`. Clear it first. clear_registry() flow = Flow(shared) results = flow.run() self.assertEqual(results["shared"].state, StepState.SUCCEEDED) self.assertNotIn("stale_side", results) self.assertEqual(side_effects, []) def test_registry_empty_after_clear(self): @task def throwaway(): return None clear_registry() from flowing import _TASK_REGISTRY self.assertEqual(_TASK_REGISTRY, []) class TestContentAddressedJournal(unittest.TestCase): """Opt-in durable journal: cross-process replay + edit divergence.""" def setUp(self): clear_registry() import tempfile self.jp = tempfile.mktemp(suffix=".jsonl") def tearDown(self): if os.path.exists(self.jp): os.unlink(self.jp) def _build(self, mult): clear_registry() calls = {"a": 0, "b": 0, "c": 0} @task def a(): calls["a"] += 1 return 10 if mult == 2: @task(depends_on=[a], name="b") def b(a): calls["b"] += 1 return a * 2 else: @task(depends_on=[a], name="b") def b(a): calls["b"] += 1 return a * 3 @task(depends_on=[b], name="c") def c(b): calls["c"] += 1 return b + 1 return c, calls def test_no_journal_path_is_unchanged(self): c, calls = self._build(2) r = Flow(c).run() self.assertEqual(r["c"].value, 21) self.assertIsNone(r["c"].step_key) self.assertFalse(r["c"].cached) def test_replay_across_fresh_flow(self): c, calls = self._build(2) Flow(c, journal_path=self.jp).run() self.assertEqual((calls["a"], calls["b"], calls["c"]), (1, 1, 1)) c2, calls2 = self._build(2) r2 = Flow(c2, journal_path=self.jp).run() self.assertEqual((calls2["a"], calls2["b"], calls2["c"]), (0, 0, 0)) self.assertEqual(r2["c"].value, 21) self.assertTrue(r2["c"].cached and r2["a"].cached) def test_edit_diverges_and_cascades(self): c, _ = self._build(2) Flow(c, journal_path=self.jp).run() c3, calls3 = self._build(3) # b's body edited r3 = Flow(c3, journal_path=self.jp).run() self.assertEqual(calls3["a"], 0) # unchanged upstream cached self.assertEqual(calls3["b"], 1) # edited task re-runs self.assertEqual(calls3["c"], 1) # dependent cascades (chained key) self.assertEqual(r3["c"].value, 31) self.assertTrue(r3["a"].cached) self.assertFalse(r3["b"].cached or r3["c"].cached) def test_cosmetic_knob_does_not_bust_key(self): clear_registry() calls = {"x": 0} @task(retry=5, name="x") def x(): calls["x"] += 1 return 99 Flow(x, journal_path=self.jp).run() clear_registry() @task(retry=0, name="x") # knob changed, body identical def x2(): calls["x"] += 1 return 99 Flow(x2, journal_path=self.jp).run() self.assertEqual(calls["x"], 1) def test_truncated_trailing_line_tolerated(self): c, _ = self._build(2) Flow(c, journal_path=self.jp).run() with open(self.jp, "a") as fh: fh.write('{"kind":"Completed","key":"v1:dead","nam') # torn write c2, calls2 = self._build(2) r = Flow(c2, journal_path=self.jp).run() self.assertEqual((calls2["a"], calls2["b"], calls2["c"]), (0, 0, 0)) self.assertEqual(r["c"].value, 21) if __name__ == "__main__": unittest.main()
-
-
CHANGELOG.md 9.5 KB
# flowing - Changelog All notable changes to the `flowing` skill are documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/). ## [1.5.0] - 2026-08-25 ### Other - Applicability boundaries, real failure signals, findable descriptions (#774) - Deprecate mapping-codebases; adopt ruff 0.16.0 baseline (#747) ## [1.4.0] - 2026-06-03 ### Other - flowing v1.4.0: content-addressed durable journal (cross-session resume) (#682) - Add surface routing to orchestration skills (native CC workflows vs custom) (#675) ## [1.4.0] - 2026-06-02 ### Added - **Content-addressed durable journal (`Flow(journal_path=...)`)** — opt-in cross-process replay, adapted from the StepKey/divergence design in `antstanley/pi-flow` specs 07. Each succeeded task appends a `Completed` entry to an append-only JSONL, keyed by a chained `step_key` = SHA-256 over the task body's bytecode + its `when`/`validate`/`retry_until` fingerprints + sorted dependency keys. A subsequent `run()` — including in a fresh container after the in-memory `self.results` is gone — replays the unchanged prefix from the journal and executes only tasks whose key is absent. This delivers the cross-session checkpoint `SKILL.md` already advertised but `resume()` only provided in-memory. - **Divergence + cascade** — editing a task body changes its `step_key`; key chaining means its dependents' keys change too, so the edited task and everything downstream re-run while the unchanged upstream prefix stays cached. Cosmetic knobs (`retry`, `retry_backoff_*`, `timeout_s`, `detached`, `name`) are excluded from the key, so tuning them does not bust the cache. - **Crash tolerance** — journal load discards a truncated trailing line (partial write on crash) rather than failing; non-picklable return values are skipped with a log line and simply re-run on the next pass. - `StepResult` gains `step_key` and `cached` fields (both `None`/`False` on the no-journal path, which is byte-for-byte unchanged). - New `compute_step_key()` / `STEP_KEY_VERSION` module surface; 5 new journal tests (33 total). ### Notes - The fingerprint hashes the task's own code, not values it closes over — a task parameterized by a captured variable will not see that variable change reflected in its key. Pass inputs through `depends_on` (the flowing analog of pi-flow's explicit `args`) for them to participate in divergence. ## [1.3.2] - 2026-05-15 ### Other - flowing v1.3.2: add human-facing README.md (#648) ## [1.3.2] - 2026-05-14 ### Added - **`README.md`** — human-facing overview for browsing the skill on GitHub, missing until now. Frames the problem (prose imperatives get generated past; a `@task` graph is structural), summarizes the three control-flow primitives, maps the file layout by audience, and points to sibling orchestration skills. Complements `SKILL.md` (agent instructions) and `references/reference.md` (agent reference). ## [1.3.1] - 2026-05-14 ### Other - flowing v1.3.1: move API reference out of SKILL.md into references/ (#647) ## [1.3.1] - 2026-05-14 ### Changed - **SKILL.md trimmed to imperative instructions.** The exhaustive API reference — full `@task` signature, `Flow` methods, resume/override, detached auto-discovery, and the `validate=`/`when=` signature gotcha — was reference material *about* the skill, not instructions *to* the agent. Moved it to `references/reference.md` (progressive disclosure); SKILL.md now keeps the trigger, mental model, quick start, the three control-flow primitives, when/when-not-to-use, and a pointer to the reference. Dropped inline version-annotation cruft (`v1.1`, `v1.2.0`, `v1.3`) that belongs in this changelog. ## [1.3.0] - 2026-05-14 ### Other - flowing v1.3.0: enforce timeout_s, validate signatures, prune dead code (#646) ## [1.3.0] - 2026-05-14 ### Added - **`timeout_s` is now enforced.** It was declared on `TaskDef` and accepted by `@task` but never read — the body ran unbounded. A task with `timeout_s` set now runs in a one-shot worker; overrunning the limit aborts the attempt as a retryable `TimeoutError` that consumes the `retry=` budget. Python can't kill the orphaned thread, so it runs until the container exits (acceptable for run-once use). - **Graph-build-time signature validation.** `flow.run()` now checks that every task body has a parameter for each of its dependencies (or `**kwargs`). A mismatch raises a clear `ValueError` naming both tasks, instead of failing mid-run with a confusing `TypeError`. - **`clear_registry()`** — module-level helper to empty `_TASK_REGISTRY`. Lets tests and long-lived REPLs run independent flows without stale detached tasks leaking across them via auto-discovery. - 8 new tests (timeout enforcement + retry interaction, signature validation, registry clearing). Suite now 28 tests. ### Fixed - **`fail_fast` parallel cancellation** now uses `pool.shutdown(wait=False, cancel_futures=True)`, which cancels queued-but-unstarted siblings. Siblings already running still can't be killed; the guarantee is "the next layer won't start," now stated explicitly in code and SKILL.md. - **`summary()` was O(n²)** — it rebuilt the task-def set on every row. Built once now, and it also picks up auto-discovered detached tasks that the old terminal-only walk missed. ### Removed - Dead `traceback` import and unused `functools.wraps` import. - `StepState.PENDING`, `RUNNING`, `RETRYING` — defined but never assignable, since `_run_step` is synchronous and only ever returns a terminal state. - Dead `_topo_sort()` function — superseded by `Flow._build_layers` and never called. - `Flow._all_task_defs()` — folded into `summary()`'s single call to `_collect_tasks()`. ## [1.2.1] - 2026-05-08 ### Other - Add flowing/SKILL.md (#627) ## [1.2.0] - 2026-05-08 ### Added - add mapping-features skill for behavioral web app documentation (#432) - add deep_read sub-agent for context-lean page processing ### Fixed - auto-discover detached tasks downstream of terminals (v1.2.0) (#613) ### Other - flowing v1.1: add when=, validate=, retry_until= control-flow primitives (#611) - Remove _MAP.md files, direct agents to tree-sitting for code navigation (#545) - Regenerate _MAP.md files after @lat: backlink insertion (#504) - Lattice v2: bidirectional source-anchored knowledge graph (#503) ## [1.2.0] - 2026-05-07 ### Fixed - **Detached tasks downstream of terminals are now auto-discovered.** In v1.1.1, `Flow(main)` would silently skip a `@task(detached=True, depends_on=[main])` defined elsewhere; the task had to be passed as an additional terminal (`Flow(main, side_effect)`). The SKILL.md said "Run in a final layer after the main DAG" which implied auto-discovery. Now matches the docs: any detached task in the module registry whose `depends_on` are all reachable from declared terminals joins the run automatically. Detached tasks with unreachable deps are still ignored (they belong to a different graph). - **Detached-on-detached chains now layer correctly.** Previously `_execute_detached` ran all detached tasks in one parallel layer, so `detachB(depends_on=[detachA], detached=True)` would be SKIPPED because `detachA` hadn't completed yet. Detached execution now uses topological layering inside the detached subset. ### Added - Module-level `_TASK_REGISTRY` populated by the `@task` decorator. Used by `Flow._collect_tasks` to find detached candidates for auto-discovery. - Test class `TestDetachedAutoDiscovery` (5 tests): direct downstream auto-discovery, backward-compat with explicit terminal, detached-on-detached chains, isolation of unrelated detached tasks, failure-isolation preservation. Total suite now 20 tests. ## [1.1.1] - 2026-05-07 ### Documentation - SKILL.md: added "Validator and predicate signatures" subsection clarifying that `validate=` and `when=` callables receive gathered dep values as kwargs by dep name. Reusing a validator across tasks with differently-named deps raises `TypeError` at validate time, surfacing as a confusing FAIL. Documents two patterns to handle reuse: `**kwargs` lookup and a name-binding factory. ## [1.1.0] - 2026-05-07 ### Added - **`when=` — conditional gate.** Receives gathered dep values as kwargs; falsy return marks the task SKIPPED and propagates to dependents. Use for branch selection in DAG topology rather than in-body `if` statements that no-op downstream tasks. - **`validate=` — edge contract.** Receives gathered dep values as kwargs; raise marks the task FAILED with **no retry** (bad inputs don't fix themselves). Validator runs before the task body; on failure the body never executes and the retry budget is preserved (`attempts=0`). - **`retry_until=` — predicate-driven loop.** Receives the task's return value; falsy return triggers a retry that consumes the existing `retry=` budget (with the same exponential backoff). On exhaustion, the last value is preserved on the FAILED result for diagnostics. Distinct from `retry=` alone, which only retries on raised exception — this retries on output shape. - Test suite at `tests/test_flowing.py` covering backward compat (chains, retry, fail propagation, override+resume), the three new primitives, and their composition. 15 tests, all green. ### Changed - SKILL.md reframed: control-first rather than throughput-first. Original motivation (cut serial tool calls to fit the 20/turn budget) is no longer the primary lever — the budget is now 50/turn and tool calls are faster. Control flow that doesn't bluff past gates is the durable value. ## [1.0.0] - 2026-03-20 ### Added - Add flowing skill — standalone DAG runner with resume, override, and detached tasks -
README.md 4.1 KB
# flowing A lightweight DAG workflow runner for Claude's ephemeral containers. Declare steps, wire dependencies, run once — control flow lives in code, not in prose imperatives. ## The problem it solves Multi-step procedures are usually written as prose: *"first fetch X, then validate Y, then if Z retry up to 3 times, otherwise skip ahead."* An LLM reads and **generates past** prose like that — the gate is a suggestion, not a wall. Skipped validation, retries that never happen, branches taken on stale state. A `@task` graph is **structural** instead. A step physically cannot run until its inputs are bound to its parameters. A gate that fires on missing or bad input can't be stepped over. The runner owns branching, retrying, validating, failure propagation, and parallelism; the LLM only supplies judgment at the leaves. ```python from flowing import task, Flow @task def fetch_data(): return {"items": [1, 2, 3]} @task(depends_on=[fetch_data]) def process(fetch_data): # param name matches the dep's name return sum(fetch_data["items"]) @task(depends_on=[process]) def store(process): print(f"Result: {process}") Flow(store).run() # topo-sorts into layers, parallel within a layer ``` ## Control-flow primitives The distinctive part — branches and contracts as graph structure, not `if` statements buried in task bodies. | Primitive | What it does | Use for | |---|---|---| | `when=` | Predicate over dep values; falsy → task SKIPPED, skip cascades to dependents | Branch selection in the topology | | `validate=` | Checks dep values before the body runs; raise → FAILED with **no retry** | Enforceable input contracts between steps | | `retry_until=` | Predicate over the return value; falsy → retry, consuming the `retry=` budget | Self-correcting LLM steps (generate → check → regenerate) | `retry_until=` is distinct from `retry=` alone: `retry=` only retries on a raised exception, `retry_until=` retries on *output shape*. ## Also handles - **Parallel execution** — independent tasks in a layer run on a thread pool. - **Resume** — `run()` → fix → `resume()` re-runs from the failure point, keeping succeeded tasks cached. `override()` injects a corrected value for a step resolved out-of-band. - **Detached side-effects** — `detached=True` tasks (memory writes, notifications) run after the main DAG and never block it on failure. - **`timeout_s=`**, **`retry=`** with exponential backoff, **`fail_fast=`**. ## Layout | File | Audience | Contents | |---|---|---| | [`SKILL.md`](SKILL.md) | Claude | Trigger, mental model, quick start, the three primitives, when / when-not-to-use | | [`references/reference.md`](references/reference.md) | Claude | Full API — every `@task` parameter, `Flow` methods, resume/override, detached auto-discovery, signature gotchas | | [`scripts/flowing.py`](scripts/flowing.py) | — | The runner itself (no third-party dependencies) | | [`tests/test_flowing.py`](tests/test_flowing.py) | — | 28 tests — `python3 -m unittest tests.test_flowing` | | [`CHANGELOG.md`](CHANGELOG.md) | — | Version history | ## When to reach for it Use it when a procedure has branches that matter, steps with input contracts, an LLM step that needs to converge, 3+ operations that can parallelize, or a pipeline where late failures shouldn't waste early work. Skip it for a single sequential operation (just call the function), for a next step that needs open-ended *reasoning* about the prior result rather than a predicate (use a think loop), or for async / distributed workflows (this is single-container, thread-pool based). ## Complements - **[orchestrating-agents](../orchestrating-agents)** — parallel API instances and delegated sub-tasks. `flowing` orders and gates work *within* one container; orchestrating-agents fans work *out* across many. - **[tiling-tree](../tiling-tree)** — MECE partitioning of a problem space. Tiling-tree decides *what* the branches are; `flowing` enforces the execution order once they exist. - **[tracking-todos](../tracking-todos)** — a human-legible checklist for loose, evolving work. `flowing` is for procedures whose shape is known up front and worth encoding as a graph. -
SKILL.md 6.9 KB
--- name: flowing description: Runs a multi-step procedure as a Python DAG, so ordering, branching and retries are enforced by the runner rather than described in prose a model can generate past. Use for "run these steps in order and retry the flaky one until the check passes", "build a pipeline that fetches, validates, then skips the upload when nothing changed", "make sure these steps cannot be skipped", "resume from where it broke instead of redoing the expensive early stages", "run these independent calls at once and merge the results", or any procedure of 3+ steps with branches, input contracts, or side effects that must not block the critical path. Primitives are depends_on, when=, validate=, retry_until=, detached= and journal_path=. Not for a single sequential call, for steps needing reasoning between them that no predicate captures, or for async and distributed work. To audit whether one verification check can actually go red, use gating. To fan work out across many subagents, use a dynamic workflow. metadata: version: 1.5.0 --- ## NOT SUPERSEDED BY DYNAMIC WORKFLOWS — read first Claude Code's dynamic workflows orchestrate **subagents** (separate contexts, fan-out to 16-concurrent / 1000-agent). This skill is a **different primitive**: single-context control flow over YOUR OWN tool calls, with durable side-effects and checkpoint resume. The workflows runtime explicitly cannot touch the filesystem or shell directly — its agents do the work and the script only coordinates them. Flowing is the inverse: the script does the work. Use flowing for an in-context pipeline (3+ steps, branches, retries, validation, detached side-effects). Use a workflow when you need many subagents. They compose; they do not compete. Do not abandon flowing for a workflow — you would lose the durable side-effects and the cross-session checkpoint that hub-spoke depends on. # Flowing — Control Flow in Code, Not Prose When a procedure needs 3+ steps with branches, retries, or contracts, encode it as a DAG of Python tasks instead of prose imperatives. Prose like "first X, then Y, then if Z retry 3×" is read and generated past. A `@task` graph is structural: a step physically cannot run until its inputs are bound, and gates that fire on bad inputs can't be skipped. The runner owns control flow — branching, retrying, validating, propagating failures, parallelizing. You provide judgment at the leaves. Runner: `scripts/flowing.py`. ## Quick Start ```python from flowing import task, Flow @task def fetch_data(): return {"items": [1, 2, 3]} @task(depends_on=[fetch_data]) def process(fetch_data): # param name must match the dep's name return sum(fetch_data["items"]) @task(depends_on=[process]) def store(process): print(f"Result: {process}") Flow(store).run() # topo-sorts, runs each layer, parallel within a layer ``` Each task receives its dependencies as kwargs named after them. Independent tasks in the same layer run in parallel. ## Control-Flow Primitives Encode branches and contracts as graph structure, not `if` statements inside task bodies. ### `when=` — conditional gate Run the task only if the predicate (over gathered dep values) is truthy. Falsy → SKIPPED, and the skip propagates to dependents. ```python @task(depends_on=[fetch], when=lambda fetch: fetch["needs_processing"]) def process(fetch): return transform(fetch["payload"]) ``` ### `validate=` — edge contract Check gathered dep values before the body runs. Raise → FAILED with **no retry** (bad inputs don't fix themselves). Pass → proceed. ```python def must_have_items(fetch): if not fetch.get("items"): raise ValueError("fetch returned empty payload") @task(depends_on=[fetch], validate=must_have_items) def process(fetch): return sum(fetch["items"]) ``` ### `retry_until=` — predicate-driven loop Run the body, then call `retry_until(value)`. True → done. False → retry, consuming the `retry=` budget. Use for self-correcting LLM steps: generate, check, regenerate. ```python @task(retry=4, retry_until=lambda r: r["valid"]) def generate_until_valid(): candidate = llm_call(...) return {"valid": passes_schema(candidate), "candidate": candidate} ``` Distinct from `retry=` alone, which only retries on a raised exception. ## Other capabilities - **Parallel execution** — independent tasks in a layer run on a thread pool (`max_workers=`). - **`detached=True`** — side-effect tasks (memory writes, notifications) that run after the main DAG and never block it on failure. - **In-process resume** — `flow.run()` → fix → `flow.resume()` re-runs from the failure point, keeping succeeded tasks cached **in memory** (same process only). `flow.override(task, value)` injects a corrected result. - **Durable journal (`journal_path=`)** — opt-in content-addressed replay that survives container death. `Flow(term, journal_path="/path/run.jsonl").run()` appends each succeeded task's result to an append-only JSONL keyed by a `step_key` = SHA-256 over the task's bytecode + its `when`/`validate`/`retry_until` bodies + its dependencies' keys (chained, so an upstream change propagates downstream). A later `run()` — even in a fresh container — replays the unchanged prefix from the journal and only executes tasks whose key is absent; editing a task body busts its key and re-runs it and its dependents, while cosmetic knobs (`retry=`, `timeout_s=`, `name`) do not. This is the cross-session checkpoint hub-spoke work relies on. Caveat: results are pickled, so non-picklable return values simply re-run; closure-captured values are not part of the key (only the task body's own code is). - **`timeout_s=`**, **`retry=`** with exponential backoff, **`fail_fast=`**. Read [references/reference.md](references/reference.md) before using anything beyond the quick start and the three primitives above — it covers every `@task` parameter, the `Flow` methods, resume/override, detached auto-discovery, and the `validate=`/`when=` signature-matching gotcha. ## When to use - A procedure has branches that matter → `when=` makes them structural. - Steps have input contracts → `validate=` makes them enforceable. - An LLM step needs to converge → `retry_until=` puts the check in the loop. - 3+ independent operations that can parallelize. - Multi-step pipelines where late failures shouldn't waste early work. - Side-effects that shouldn't block the critical path → `detached=True`. ## When NOT to use - A single sequential operation — just call the function. - The next step needs *reasoning* about the prior result that can't be a predicate — use a think loop. - Async or distributed workflows — this is single-container, thread-pool based. ## Authoring discipline If you find yourself writing prose like *"first call X, validate Y, then if Z retry up to 3 times"* — that is a flowing graph. Refactor before shipping. Prose imperatives don't enforce; `@task` graphs do.
Comments (0)
Sign in to join the conversation.
Reviews (0)
No reviews yet.
No comments yet.