Metadata-Version: 2.5
Name: dagflows
Version: 0.2.1
Summary: SDK for authoring and running Dagflows workflows
Project-URL: Homepage, https://github.com/dagflows/sdk-python
Project-URL: Repository, https://github.com/dagflows/sdk-python
Project-URL: Issues, https://github.com/dagflows/sdk-python/issues
Author: dagflows
License-Expression: Apache-2.0
License-File: LICENSE
Keywords: dag,orchestration,workflow
Classifier: Development Status :: 4 - Beta
Classifier: Intended Audience :: Developers
Classifier: Programming Language :: Python :: 3.11
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: System :: Distributed Computing
Classifier: Typing :: Typed
Requires-Python: >=3.11
Provides-Extra: httpx
Requires-Dist: httpx>=0.27; extra == 'httpx'
Description-Content-Type: text/markdown

# dagflows

Python SDK for authoring and running Dagflows workflows.

## Installation

<!-- not-tested: installs from the package index -->
```bash
pip install dagflows
```

Requires Python 3.11+.

To enable the optional async HTTP client backend:

<!-- not-tested: installs from the package index -->
```bash
pip install dagflows[httpx]
```

## Quickstart

Define a workflow and its node functions in a module:

```python
# app/workflow.py
from dagflows import Workflow

wf = Workflow("order_pipeline", max_concurrent_nodes=4)


@wf.node()
def fetch_orders(ctx, inputs):
    return {"orders": [{"id": 1, "amount": 100}, {"id": 2, "amount": 250}]}


@wf.node(depends=[fetch_orders])
def calculate_totals(ctx, inputs):
    orders = inputs["fetch_orders"].value()["orders"]
    total = sum(order["amount"] for order in orders)
    return {"total": total}
```

### Handler Signatures

A handler declares what it needs, so a node that only reads its parents does
not carry an unused `ctx`. Ask for arguments **by name**, in any order:

```python
# Both
@wf.node()
def step_a(ctx, inputs):
    ...

# Inputs only
@wf.node(depends=[step_a])
def step_b(inputs):
    ...

# Context only
@wf.node()
def step_c(ctx):
    ...

# Neither
@wf.node()
def step_d():
    return {"status": "ready"}

# Async works identically
@wf.node(depends=[step_a])
async def step_async(ctx, inputs):
    async for item in inputs["step_a"]:
        ...
    return {"done": True}
```

### Annotating

Annotate for editor completion, which the SDK ships `py.typed` to support:

```python
from dagflows import Ctx, Inputs


@wf.node()
def step(ctx: Ctx, inputs: Inputs):
    inputs.one().value()      # completion on both
```

An annotation also frees the name, so call the parameters whatever suits you:

```python
@wf.node()
def step(payload: Inputs, context: Ctx):
    ...
```

A parameter that is neither named nor annotated is refused, naming the fix:

```
transform: cannot tell what to pass to 'data'.
Name it 'ctx' or 'inputs', or annotate it:

    def transform(data: Inputs): ...
```

That refusal happens while the manifest is emitted, so a bad signature fails
the build rather than the first run.

## Node Configuration

### Resource Limits and Retries

Configure execution limits and retry policies on `@wf.node()`:

```python
from dagflows import ExecutionConfig, RetryConfig, Workflow

wf = Workflow("data_pipeline")


@wf.node(
    execution=ExecutionConfig(timeout=300, memory_limit_mb=512),
    retry=RetryConfig(max_attempts=3, initial_backoff_ms=1000),
)
def extract_data(ctx, inputs):
    ...
```

### Cross-Project Dependencies

To depend on a node from another project or workflow, use `wf.external_node`:

```python
@wf.node(depends=[fetch_orders, wf.external_node("inventory_check")])
def fulfill_order(ctx, inputs):
    ...
```

## Standalone Node Scripts (Optional)

In addition to defining nodes within a workflow module, you can also write standalone script files using `dagflows.run`:

```python
# app/nodes/custom_task.py
import dagflows


def handler(ctx, inputs):
    data = inputs.one().value()
    return {"processed": True, "count": len(data)}


if __name__ == "__main__":
    dagflows.run(handler)
```

## Working with Inputs

Inputs are accessed through input handles:

```python
upstream = inputs["fetch_orders"]

# Retrieve complete value, refused if it exceeds the node's memory limit
data = upstream.value()

# Iterate without holding it all (sync and async iteration both work)
for record in upstream:
    process_record(record)

# Raw chunks for opaque content: a generator, not a bytes object
for chunk in upstream.bytes():
    sink.write(chunk)
```

Iteration yields one **record** at a time for row oriented payloads, which is
what lets a node read more than it can hold. A JSON input has no records, so it
yields the whole document as a single item.

| content type | `for x in handle` yields |
| --- | --- |
| `application/x-ndjson`, `text/csv` | one row at a time |
| `application/json` | the whole document, once |

Iteration re-opens the source each pass, so looping twice costs a second read
rather than silently yielding nothing the second time.

If a node has exactly one parent dependency, use `inputs.one()`:

```python
single_input = inputs.one().value()
```

## Producing Outputs

Return a Python dictionary/value directly, or use `Result` for explicit routing and formatting:

```python
from dagflows import ContentType, Result

# Direct value return
return {"count": 42}

# Route downstream execution to a specific branch
return Result(output={"status": "approved"}, next="process_payment")

# Halt execution of this branch
return Result(output={}, stop=True)

# Explicit content type specification
return Result(output=stream_generator(), content_type=ContentType.NDJSON)
```

## Error Handling

Raise `Fail` to signal structured execution errors with retry instructions:

```python
from dagflows import EXECUTION, Fail

raise Fail("Payment gateway unavailable", category=EXECUTION, retry_after=30)
```

Available error categories:
- `EXECUTION`
- `INFRASTRUCTURE`
- `TIMEOUT`
- `PERMANENT`

## Command Line Interface

The SDK provides CLI commands invoked via `python -m dagflows`:

### Manifest Management

```bash
# Generate manifest file
python -m dagflows build manifest app.workflow -o dagflows-manifest.json

# Check if an existing manifest is up to date
python -m dagflows build manifest app.workflow --check

# Validate workflow definition without writing output
python -m dagflows build validate app.workflow
```

### Local Development & Testing

`dev run` executes a node with no VM, no platform and no network. It takes a
**script path**, so it applies to standalone node scripts:

```bash
python -m dagflows dev run app/nodes/custom_task.py --input fetch_orders=orders.json
```

`dev fixture` writes a starting input envelope, so nobody has to invent one. It
takes either form:

```bash
# a standalone script
python -m dagflows dev fixture app/nodes/custom_task.py -o fixture.json

# a node declared with @wf.node
python -m dagflows dev fixture app.workflow:calculate_totals --input fetch_orders=orders.json -o m.json
```

It prints the command that runs the node against what it wrote, which for a
decorated node goes through the interpreter because a `module:function`
entrypoint is not a file:

```bash
DAGFLOWS_INPUT=m.json DAGFLOWS_OUTPUT=out.json python -m dagflows invoke --node calculate_totals
```

That is the local loop for decorator-defined nodes, since `dev run` needs a
file to execute.

`dev run` accepts several options worth knowing:

```text
--input users=rows.ndjson     a file, content type inferred from the suffix
--input users='{"n": 1}'      inline json
--memory-limit-mb 512         what the node believes it has
--inline-max-bytes 262144     when an output would offload
--keep-fixture <path>         write the envelope instead of discarding it
```

All CLI commands support `--json` for machine-readable output. Exit codes are
`0` success, `1` the operation failed, `2` the command was wrong.

## License

Apache-2.0
