Claude Skill

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,

LLM Mart · 0 points · 0 views 0 listing impressions 0 install-command copies
Virus-scanned Reviewed automatically before listing.

Full trust report

Download oaustegard-claude-skills-plugins_environment-and-config_skills_flowing-e39c726.zip · 24 KB
Part of oaustegard/claude-skills — 39 skills

Install

skills CLI npx skills add https://github.com/oaustegard/claude-skills/tree/main/plugins/environment-and-config/skills/flowing
Claude Code claude plugin marketplace add https://llmmart.ai/marketplace.json && claude plugin install oaustegard-claude-skills@llmmart
Git 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=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 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. flowing orders 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; flowing enforces the execution order once they exist.
  • 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 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 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 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.

No comments yet.

Reviews (0)

No reviews yet.

Related