lightflow
Health Gecti
- License — License: Apache-2.0
- Description — Repository has a description
- Active repo — Last push 0 days ago
- Community trust — 12 GitHub stars
Code Gecti
- Code scan — Scanned 12 files during light audit, no dangerous patterns found
Permissions Gecti
- Permissions — No dangerous permissions requested
Bu listing icin henuz AI raporu yok.
Lightweight, local-first Python DAG workflow compiler and execution engine with human-in-the-loop checkpoints for CLI and AI-agent workflows.
Lightflow: Deterministic DAG Workflow Engine for AI Agents
Lightflow is a lightweight, local-first workflow compiler and execution
engine that gives AI agents a structural backbone for multi-step
workflows—turning steps an agent would otherwise follow from prose instructions
into an explicit DAG of python_action stages and operator_action checkpoints
with dependencies, retries, and rollbacks enforced by the engine.
Every step runs across process boundaries as a plain CLI command that an agent
or operator can start, pause, and resume against a single append-only JSON
ledger (passport.json). Rather than burying "ask before step 7" inside a
long prompt blob where context decay and momentum cause agents to skip ahead,
Lightflow halts the process (exit 2) and delivers each checkpoint as a fresh,
dedicated turn at the exact moment of decision. It provides zero-daemon,
append-only state guardrails in pure Python—depending only on PyYAML andjsonschema, with built-in standard-library fallbacks when running directly
from a checkout.
flowchart LR
S1["1. fetch_top_stories<br/>(python_action + retry)<br/>Fetch live API data"] --> S2{"2. editorial_gate<br/>(operator_action · exit 2)<br/>Human/Agent approval + JSON Schema"}
S2 -->|payload.outputs.editorial_gate.approved == true| S3["3. publish_digest<br/>(python_action)<br/>Idempotent side effect"]
Key Features
- Append-Only Passport Ledger (
passport.json): Every stage transition
appends an immutableStamp(PENDING,COMPLETED,PAUSED,SKIPPED,FAILED) guarded by cross-platform file locks and atomicos.replace,
recording operator and agent session provenance (LIGHTFLOW_OPERATOR/ANTIGRAVITY_CONVERSATION_ID). - Selective Subgraph Re-Arming: When a stage fails, fix the underlying
code or environment and runresume. Lightflow re-arms only the failed
stage and its downstream collateral subgraph—never re-executing completed
upstream side effects. - Human-in-the-Loop Checkpoints (
operator_action): Suspends execution
with exit code2so the gate arrives as its own dedicated turn rather than
a buried prompt instruction, validates operator payloads against JSON Schema,
and auto-generates drift-freeresumecommands. - Safe, CEL-like Expressions:
run_if, polling conditions, and dynamic
gate instructions use a small subset of CEL syntax
(string(),int(),size(),has()), evaluated by a pure-Python AST
walker without native C/Go dependencies or unsafeeval(). - YAML & JSON Manifests: Author workflows declaratively in YAML
(.yaml) or JSON (.json) (with built-in.textproto
compatibility). - Zero-Dependency Interactive Visualizer (
visualizer.html): Generates a
standalone, single-file HTML5 DAG canvas, execution trace log, and
payload/manifest inspector with zero external CDN requests. - First-Class Agent Integrations:
- Model Context Protocol (MCP): Built-in stdio MCP server
(lightflow-mcp) for Claude Code, Claude Desktop, Gemini CLI, and
Cursor, plusCLAUDE.md. - Agent Skills & Streaming Notifications: Includes
skills/run_lightflows/SKILL.md
(executing and resuming workflows),skills/create_lightflows/SKILL.md
(authoring new workflows), andbenchmarks/SKILL.md, plus live Mermaid
progress notifications when running insideagentapi-compatible hosts.
- Model Context Protocol (MCP): Built-in stdio MCP server
- Consistency & Token Efficiency (Benchmarks):
Evaluated across 7 matched domain workflow pairs (3to22stages) and15live zero-context LLM subagent sessions (N=5per cohort), Lightflow
achieves38xlower runtime token variance (±4 tokvs.±144 tok),3.1x-3.8xless LLM-generated code,6.4xfewer session skill
tokens (1universal runner skill vs.7workflow skills), and100%
vs.50%first-try compliance under prose-to-code drift.
When to Use Lightflow (and When Not To)
Reach for Lightflow when all of these are true:
- The task has several ordered steps with real side effects (modifying files,
calling external APIs, publishing artifacts) that must not be repeated on a
retry. - You want a human or an agent to inspect intermediate output and approve a
checkpoint partway through (operator_action). - You want the run to be inspectable on disk (
passport.json) and resumable
across separate CLI invocations or compacted agent sessions.
Look elsewhere when:
- You need concurrent execution. Execution is strictly single-threaded and
evaluates stages sequentially in deterministic topological order. - You need unattended scheduling. Lightflow does not include a cron
daemon; triggerlightflow startfromcron,systemd, or CI if you need
a recurring run. - You need shared multi-user state. State is a local file guarded by an
advisory lock (flock/msvcrt), not a networked service. - It is a single step. Just call the function directly.
| Category | Representative Tools | Infrastructure | Execution Model | Where Lightflow Differs |
|---|---|---|---|---|
| Distributed Orchestrators | Temporal, Prefect, Airflow | Server daemon + database | Distributed workers | Zero daemon or DB; state is a single atomic local JSON ledger (passport.json). |
| Cyclic Agent Graphs | LangGraph, Burr | In-process Python library | Cyclic state machines | Process-boundary suspension (exit 2) with auto-generated CLI resume commands and AST-restricted CEL-like conditions. |
| Function / Task DAGs | Hamilton, doit, Kedro | Local Python library | Single-run in-memory DAG | Built-in operator_action JSON Schema gates, rollback_action hooks, and append-only Stamp history. |
Lightflow DAG vs. Traditional SKILL.md + actions.py (Full A/B Benchmark)
| Dimension | Lightflow DAG (Arm A) |
Steelmanned Traditional Skill (Arm B: SKILL.md + actions.py) |
Empirical Delta (benchmarks/) |
|---|---|---|---|
| Skill Context Loaded | 1 universal runner skill (~1,097 tok once per session; 0 tok on workflows 2–7) |
1 bespoke SKILL.md per workflow (~7,075 tok across 7 workflows) |
6.4x fewer skill tokens (O(1) vs. O(N)) |
Agent-Generated Code (cmd) |
Declarative CLI (lightflow start / resume, ~795 tok across 7 workflows) |
Inline python3 -c polling loops, branches & rollbacks (~2,992 tok) |
3.8x less generated code (6.5x on 22-stage DAG) |
Multi-Run Consistency (N=5 Live Subagent Sessions, 8 & 9 Stages) |
~1,847 ± 4 tok warm runtime (~502 ± 4 tok cmd, 10/10 pass) |
~2,433 ± 144 tok warm runtime (~1,565 ± 97 tok cmd, 10/10 pass) |
38x lower runtime token SD (±4 vs. ±144 tok) & 3.1x less code |
| Prose-to-Code Drift Resilience | Engine unpacks returns & enforces DAG schema (10/10 first-try pass) |
1 omitted return-type detail crashes Stage 2 after Stage 1 mutates disk (5/10 first-try pass) |
100% vs. 50% first-try compliance under drift |
Session Total (Skill + cmd + stdout) |
~2,944 ± 4 tok (live 2-wf) • ~5,570 tok (7-wf ladder) |
~5,904 ± 144 tok (live 2-wf) • ~10,992 tok (7-wf ladder) |
49%–50% lower total session tokens |
Crossover point: For a single 3-stage linear workflow in a one-off session, a bespoke
SKILL.md+/tmp/state.jsonis ~760 tokens lighter (~732vs.~1,495 tok); Lightflow breaks even by the 2nd–3rd workflow in a session or on any single workflow with $\ge 8$ stages, conditional branches, polling, or rollbacks.
Installation
git clone https://github.com/google/lightflow.git
cd lightflow
pip install -e .
(Or run directly from a checkout without installing dependencies usingPYTHONPATH=. python3 -m lightflow ...)
Where state lives
Each run's passport.json is written to~/.lightflow/lightflow_state_<log_id>/passport.json by default. SetLIGHTFLOW_STATE_DIR=/some/dir to relocate the lightflow_state_<log_id>/
directories (for example into a project-local or CI-scoped folder). Gate and
re-arm stamps attribute the action to --operator if passed, LIGHTFLOW_OPERATOR
if set, or <user> (agent:<ANTIGRAVITY_CONVERSATION_ID>) when invoked inside an
agent session so you can trace any approval back to its conversation transcript.
Security model
python_import executes arbitrary Python by design — treat manifests as
code and review them like you would a script. For manifests authored by an
agent, set LIGHTFLOW_ALLOWED_IMPORT_PREFIXES to a comma-separated list of
module prefixes (e.g.LIGHTFLOW_ALLOWED_IMPORT_PREFIXES=my_pkg.actions,examples) and the engine will
refuse to import any action outside those prefixes, while disabling
manifest-local sys.path injection and module eviction so an untrusted manifest
directory cannot shadow allowlisted packages.
60-Second Runnable Quickstart
Run the included live Hacker News curation workflow right out of the box:
# 1. Dry-run simulation (validates DAG topology, CEL-like expressions & JSON Schemas)
lightflow dry_run --lightflow=examples/hn_digest
# 2. Start execution (fetches live HN top stories, then suspends at 'editorial_gate' with exit code 2)
lightflow start --lightflow=examples/hn_digest --log_id=hn_demo
# 3. Approve the human gate and finish the workflow
lightflow resume --lightflow=examples/hn_digest --log_id=hn_demo \
--stage=editorial_gate --resolution=APPROVE \
--payload='{"editor_note": "Top engineering reads today.", "output_path": "/tmp/hn_digest.md"}'
# 4. Generate a standalone interactive HTML5 DAG & timeline visualizer
lightflow visualize --lightflow=examples/hn_digest --log_id=hn_demo --output=/tmp/viz.html
--lightflow accepts a workflow directory (containing lightflow.yaml,lightflow.yml, lightflow.json, or lightflow.textproto) or a direct file
path. python3 -m lightflow ... is equivalent to lightflow ... whenever the
console script is not on your PATH. On shells that strip inner double quotes
(such as Windows PowerShell 5.1), you can also pass JSON payloads from a file or
stdin via [email protected] (or --payload=@-).
Included Runnable Examples (examples/)
| Tier | Example | Stages & Control Flow | Pattern Demonstrated | External API / Sandbox |
|---|---|---|---|---|
| Starter | examples/hn_digest/ |
3 stages · 1 gate · retry |
Live API → Editorial Gate → Idempotent Publish: Fetches top stories, previews #1 in the gate prompt, and writes a Markdown digest. | Hacker News Firebase API (zero auth) |
| Starter | examples/usgs_seismic_alert/ |
3 stages · 1 gate · retry |
Live Telemetry Triage: Polls 24h M4.5+ global seismic events, pauses for severity classification (INFO / ADVISORY / WATCH), and emits a bulletin. |
USGS GeoJSON Feed (zero auth) |
| Starter | examples/pypi_upgrade_guard/ |
3 stages · 1 gate · 1 rollback |
Live Audit → Approval Gate → Automatic rollback_action: Audits pinned packages against PyPI, applies upgrades, and auto-restores .bak on smoke failure. |
PyPI JSON API (zero auth) |
| Starter | examples/async_job_watcher/ |
4 stages · 1 poll · ALL_DONE |
Async Background Polling & Cleanup: Spawns a detached background worker, polls its status file via polling_policy until READY, and cleans up temp files. |
Detached Worker + Local JSON Status |
| Production Recovery | examples/blue_green_release/ |
8 stages · 1 gate · 1 poll · 1 rollback |
Blue/Green Canary Release & Traffic Rollback: Parallel security/integration diamond, green warmup poll, cutover gate, automatic revert_canary_traffic on 5xx SLO breach, and surgical recovery. |
Local Deployment Slot & Canary State Machine |
| Production Recovery | examples/incident_db_failover/ |
9 stages · 1 gate · 1 poll · 1 rollback |
Regional DB Failover & Pooler Rollback: Mutually exclusive run_if WAL replay vs. zero-loss branch, ALL_DONE STONITH fence, IC gate, one-way standby promotion, and PgBouncer rollback. |
Local HA PostgreSQL & PgBouncer Sandbox |
| Enterprise Scale | examples/tenant_gitops_onboarding/ |
22 stages · 2 gates · 6 polls |
Multi-Branch Kubernetes, Cloud IAM & GitOps Onboarding: 4-phase pipeline with 3 branching patterns (early handoff, 3-way IAM fan-out + ALL_DONE join, S3/Vault + KMS diamond) and ArgoCD sync loops. |
Local GitOps / Cloud IAM / SCIM / KMS / ArgoCD Sandbox |
| Authoring Meta-DAG | examples/create_lightflow/ |
7 stages · 3 gates |
Staggered Agent Skill Delivery ("A Lightflow to Create Lightflows"): 3 design & coding gates interleaved with deterministic AST, unittest, and dry_run verification stages. |
Python ast + unittest + LightflowEngine |
Authoring Your Own Workflow
1. Write Python Actions (my_pkg/actions.py)
Each python_action is a regular Python function that receives a copy of the
current payload dictionary (plus any static_kwargs from the manifest) and
returns either (payload_delta_dict, message_str) or just payload_delta_dict
(which Lightflow merges into payload and records underpayload.outputs.<stage_name>):
from typing import Any
def fetch_rows(
payload: dict[str, Any], limit: int = 500, dry_run: bool = False
) -> tuple[dict[str, Any], str]:
if dry_run or payload.get("dry_run"):
return {"count": 0, "summary": "0 rows (dry run)"}, "[DRY RUN] Skipped fetch"
return {"count": limit, "summary": f"{limit} rows"}, f"Fetched {limit} rows"
def publish_report(payload: dict[str, Any]) -> tuple[dict[str, Any], str]:
summary = payload["outputs"]["fetch"]["summary"]
return {"published": True}, f"Published report ({summary})"
2. Define the Workflow Manifest (lightflow.yaml)
name: daily_revenue_report
description: Fetch rows, wait for human approval, and publish report.
actions:
- id: fetch_rows
python_import: my_pkg.actions.fetch_rows
- id: publish_report
python_import: my_pkg.actions.publish_report
stages:
- name: fetch
description: Fetch daily rows
python_action:
action_id: fetch_rows
static_kwargs:
limit: 500
retry_policy:
max_attempts: 3
- name: approve
description: Human review gate
run_after:
- fetch
operator_action:
instructions: "'Ready to publish: ' + payload.outputs.fetch.summary"
json_schema:
type: object
required: [approved]
properties:
approved:
type: boolean
- name: publish
description: Publish verified report
run_after:
- approve
run_if: "payload.outputs.approve.approved == true"
python_action:
action_id: publish_report
3. Stage & Expression Reference
| Stage Field | Purpose |
|---|---|
run_after |
Upstream stage names that must finish before this stage evaluates. |
run_if |
CEL-like expression evaluated against payload; stage is stamped SKIPPED if false. |
trigger_rule |
ALL_SUCCESS (default) or ALL_DONE (runs even if an upstream dependency failed or skipped, useful for cleanup stages). |
timeout_seconds |
Bounds how long the engine waits for each attempt (0–86400s). A timeout does not cancel the action: it keeps running in a background thread. Pass explicit timeouts to the action's own I/O clients, and make timed actions idempotent. |
retry_policy |
Opt-in retry (max_attempts, initial_backoff_seconds, backoff_multiplier) with ±10% jitter for exceptions and timeouts. |
polling_policy |
Polls python_action until condition evaluates to true (interval_seconds, timeout_seconds, max_attempts). |
rollback_action |
Compensating python_action invoked automatically after a stage exhausts its retries (skipped if a timed-out attempt's background thread is still running); its output delta is merged into payload. |
Expression Evaluator (run_if, polling_policy.condition,operator_action.instructions):
- Operators:
&&,||,!,==,!=,<,<=,>,>=,+,-,*,/ - Built-ins:
size(x)/len(x),has(x.field),string(x),int(x),double(x),bool(x)(all null-safe:string(null) == "",int(null) == 0) - String & Collection Methods:
contains(),startsWith(),endsWith(),replace(),split(),strip(),lower(),upper() - Per-Stage Output Namespacing: Every stage's delta is merged into
payloadand also recorded underpayload.outputs.<stage_name>to prevent
key collisions across stages.
Execution & Failure Semantics
Stage Lifecycle (passport.json)
Each stage carries zero or more stamps in the passport. A stage's current
state is its most recent stamp; earlier stamps are retained as an audit log.
stateDiagram-v2
[*] --> Unstamped: never reached
Unstamped --> PENDING: engine starts stage
PENDING --> COMPLETED: action returned
PENDING --> FAILED: action raised (retries exhausted)
PENDING --> PAUSED: operator_action reached (exit 2)
Unstamped --> SKIPPED: run_if false / upstream skipped
Unstamped --> FAILED: upstream failed (ALL_SUCCESS)
PAUSED --> COMPLETED: resume --resolution=APPROVE
FAILED --> PENDING: start/resume re-arms subgraph
COMPLETED --> [*]
When a stage fails:
- The failing stage is stamped
FAILEDwith its traceback, and itsrollback_actionruns if configured. - Downstream dependents under
ALL_SUCCESSare stampedFAILEDas collateral
(Upstream dependency failed), while independent parallel branches in a
diamond continue to completion andALL_DONEcleanup stages still execute. - Running
lightflow resumeappends a freshPENDINGstamp to the failed
stage and its collateral dependents—re-running only the broken subgraph
while preserving completed upstream work and keeping the earlierFAILED
stamp in history. To selectively re-run an alreadyCOMPLETEDorSKIPPED
stage (and its downstream dependents) without wiping upstream completed
work, pass--rerun=<stage_name>(lightflow resume --lightflow=<dir> --log_id=<id> --rerun=<stage_name>, or add--no-cascadeto re-run only
that stage).
CLI Exit Codes
| Exit Code | Meaning | Typical Next Step |
|---|---|---|
0 |
Workflow completed (or was already completed). | Done. |
1 |
Stage failure, compile/schema error, or state lock contention. | Inspect lightflow status, fix the cause, and run lightflow resume. |
2 |
Suspended at an operator_action gate awaiting human/agent input. |
Run lightflow resume --stage=<name> --resolution=APPROVE --payload='{...}'. |
3 |
status only: no saved state for the log ID. |
Check the log ID, or start. |
4 |
status only: another process is running the run. |
Wait, then check status again. |
lightflow status exits with the code of the run it inspects (0 completed,1 failed, rejected, or interrupted, 2 paused), so a script or agent can poll
it and branch exactly as it would on start or resume. It prints one line per
stage plus any paused gate's Instructions, Payload Schema, and resume command;--verbose adds the Passport payload, full stamp log, and Mermaid diagram.
Agent Ecosystem Setup
For Claude Code, Claude Desktop, Gemini CLI & Cursor (MCP)
Add lightflow-mcp to your claude_desktop_config.json or .mcp.json:
{
"mcpServers": {
"lightflow": {
"command": "lightflow-mcp",
"args": []
}
}
}
Copy CLAUDE.md into your repository root (or installskills/run_lightflows/SKILL.md andskills/create_lightflows/SKILL.md) to
instruct your agent on executing, resuming, and authoring Lightflow workflows.
The create_lightflow scaffolding meta-workflow those files reference lives in
this repository under examples/create_lightflow/; run it from a clone ofgithub.com/google/lightflow, or pass the absolute path to that directory in
your clone (--lightflow=/path/to/lightflow/examples/create_lightflow).
License & Disclaimer
Licensed under the Apache 2.0 License.
This is not an officially supported Google product. Eligibility for the
Google Open Source Software Vulnerability Rewards Program
is determined by the
Google Open Source Software Vulnerability Reward Program Rules.
Yorumlar (0)
Yorum birakmak icin giris yap.
Yorum birakSonuc bulunamadi