Metadata-Version: 2.5
Name: pyflowred
Version: 0.8.0
Summary: A Python port of Node-RED: the same editor, the same flow semantics, a Python runtime.
Project-URL: Homepage, https://github.com/DataVoyage/pyflowred
Project-URL: Documentation, https://github.com/DataVoyage/pyflowred/blob/main/docs/operations.md
Project-URL: Node reference, https://github.com/DataVoyage/pyflowred/blob/main/docs/nodes/README.md
Project-URL: Source, https://github.com/DataVoyage/pyflowred
Project-URL: Changelog, https://github.com/DataVoyage/pyflowred/blob/main/CHANGELOG.md
Author: Melvin Neffle
License: Apache-2.0
Keywords: automation,flow,iot,low-code,node-red,streaming
Classifier: Development Status :: 4 - Beta
Classifier: Environment :: Web Environment
Classifier: Framework :: AsyncIO
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: Apache Software License
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Python :: 3.13
Classifier: Topic :: Software Development :: Libraries :: Application Frameworks
Classifier: Topic :: System :: Monitoring
Requires-Python: >=3.13
Requires-Dist: bcrypt>=4.2
Requires-Dist: beautifulsoup4>=4.12
Requires-Dist: chevron>=0.14
Requires-Dist: confluent-kafka>=2.5
Requires-Dist: croniter>=6.0
Requires-Dist: cryptography>=44.0
Requires-Dist: duckdb>=1.1
Requires-Dist: fastmcp>=4.0
Requires-Dist: google-auth>=2.35
Requires-Dist: httpx>=0.28
Requires-Dist: jsonata-python>=0.7.0
Requires-Dist: jsonschema>=4.23
Requires-Dist: lxml>=5.3
Requires-Dist: obstore>=0.11
Requires-Dist: packaging>=24.0
Requires-Dist: paho-mqtt>=2.1
Requires-Dist: prometheus-client>=0.21
Requires-Dist: pyarrow>=18.0
Requires-Dist: python-dateutil>=2.9
Requires-Dist: python-lsp-server[pyflakes]>=1.12
Requires-Dist: python-multipart>=0.0.20
Requires-Dist: pyyaml>=6.0
Requires-Dist: rocksdict>=0.3.29
Requires-Dist: starlette>=0.41
Requires-Dist: uvicorn[standard]>=0.32
Requires-Dist: watchfiles>=1.0
Requires-Dist: websockets>=14.0
Provides-Extra: dev
Requires-Dist: playwright>=1.49; extra == 'dev'
Requires-Dist: pytest-asyncio>=0.24; extra == 'dev'
Requires-Dist: pytest>=8.3; extra == 'dev'
Requires-Dist: ruff>=0.8; extra == 'dev'
Description-Content-Type: text/markdown

# PyFlowRED

A port of [Node-RED](https://nodered.org) 5.0.7 to Python.

The same editor, the same flow file, the same node set, the same HTTP admin
API — with a Python runtime and nodes written in Python.

Flows can also be kept as **Python files** instead of JSON, and the node set
adds what an industrial ingest needs: stream operators on RocksDB, buffering,
Parquet, Kafka, Cloud Storage and a DuckDB store you can run SQL across.

![The PyFlowRED editor](https://raw.githubusercontent.com/DataVoyage/pyflowred/main/docs/editor.png)

## What "port" means here

| Layer | Origin |
|---|---|
| Editor client | Node-RED source, vendored and built with the same pipeline |
| Node edit dialogs and message catalogues | Node-RED source, adapted only where the runtime language shows through |
| HTTP admin API | Re-implemented in Python, wire-compatible with the editor |
| Flow runtime, node registry, context, credentials, storage | Re-implemented in Python |
| Projects (git-backed flows), multiplayer, palette catalogue | Re-implemented in Python |
| Core nodes | Re-implemented in Python — all 50 types |
| Nineteen further node types | New — [see below](#the-nodes-it-adds) |

The editor is not a reimplementation. It is Node-RED's own editor, talking to
a Python implementation of the API it expects. Flows exported from Node-RED
import into PyFlowRED and back again: the flow file is the same JSON, and the
credentials file uses the same encryption.

See [docs/porting.md](https://github.com/DataVoyage/pyflowred/blob/main/docs/porting.md) for the full correspondence.

## The nodes it adds

Beyond the ported core set, **twenty-four node types in seven packages** come with
the distribution — no palette install, no dependency to add. They are aimed at
the industrial case: readings arriving from a plant that have to become
measurements and events before anything can be done with them.

| Package | Types | |
|---|---|---|
| **[Stream processing](https://github.com/DataVoyage/pyflowred/blob/main/docs/nodes/stream-processing.md)** | `stream-state` `aggregate` `derive` `threshold` `transition` `assemble` `compress` `pattern` | Seven window shapes, eighteen aggregation functions, hysteresis and dwell, state durations, swinging-door compression, multi-condition patterns — on keyed state in RocksDB |
| **[Google Cloud Storage](https://github.com/DataVoyage/pyflowred/blob/main/docs/nodes/cloud-storage.md)** | `gcs-config` `gcs in` `gcs out` `gcs list` `gcs delete` | Read, write, list and delete on obstore, with every value settable from the message |
| **[Kafka](https://github.com/DataVoyage/pyflowred/blob/main/docs/nodes/kafka.md)** | `kafka-broker` `kafka in` `kafka out` | Consume and produce on librdkafka, with every librdkafka property reachable |
| **[Buffering and Parquet](https://github.com/DataVoyage/pyflowred/blob/main/docs/nodes/buffering.md)** | `buffer` `parquet` | Hold messages unchanged in RocksDB until an interval, a count or a control message closes the block, then write it as one Parquet file |
| **[DuckDB](https://github.com/DataVoyage/pyflowred/blob/main/docs/nodes/duckdb.md)** | `duckdb-store` `duckdb out` `duckdb` | A body of messages on disk with SQL across all of it — twenty million rows in 171 MB, a grouped aggregate in 34 ms, and the query in a thread so it never holds the loop |
| **[Secrets](https://github.com/DataVoyage/pyflowred/blob/main/docs/nodes/secrets.md)** | `secret` | Fetch a secret at message time; `${secret:NAME}` for connection-time credentials |
| **[Prometheus](https://github.com/DataVoyage/pyflowred/blob/main/docs/nodes/prometheus.md)** | `prometheus-metric-config` `prometheus out` | Declare a metric in a flow and feed it, into the same registry the runtime's own metrics use |

The state a stream operator keeps lives in RocksDB rather than in the process,
which is the difference between a hundred thousand keys costing tens of
megabytes resident and costing a gigabyte of Python objects. A window closed by
a control message — a shift, a batch, a production order, a cleaning cycle — is
one mechanism rather than four nodes, and it is what an availability figure is
actually computed from.

Every one of these carries its help text in the editor's info sidebar, and a
**worked example flow** — *Import > Examples > node-red* in the editor. Sixteen
of them, one per node: inject, the node, debug, and a comment block beside each
row saying what to look for.
[docs/nodes/](https://github.com/DataVoyage/pyflowred/blob/main/docs/nodes/README.md) is the long form: every setting, what each
node puts on the message, and the reasoning where the behaviour is not the
obvious one.

## An agent can work on a running instance

An **MCP server on the instance's own port**, so a coding agent asks the live
runtime what is true, changes the flows through the path the editor uses, and
is told by the same runtime whether the change held. Off unless switched on,
in the editor under *Settings → Agent*.

It knows what the node set is because it reads the same edit dialogs the
editor does, it lints a change before any of it is applied, and its failures
say what to try instead rather than where our code gave up. Nothing it returns
is an unbounded stream: a tap answers with the **shape** of what went past —
every field, its types with their split, its range — because twenty thousand
messages would fill an agent's context and teach it nothing that twenty would
not.

[docs/agent.md](https://github.com/DataVoyage/pyflowred/blob/main/docs/agent.md)
is the reference.

## Where it deliberately differs

Two things cannot be carried over one-to-one, because they *are* the
JavaScript runtime:

1. **The Function node runs Python.** The same three tabs (*On Start*, *On
   Message*, *On Stop*), the same `node` / `flow` / `env` objects, the same
   send and done semantics — but the body is Python:

   ```python
   msg["payload"] = msg["payload"] * 2
   return msg
   ```

   The body is compiled into an `async def`, so `await` works directly.
   `global` is a Python keyword, so the global context object is `global_`.

   The editor knows it is Python. `pylsp` runs alongside the runtime, so the
   Function node has completion, hover, signature help and problems as you
   type — for the standard library and for anything the Setup tab installed,
   and for `msg`, `node`, `flow`, `global_` and `env`, whose stubs are
   generated from the runtime's own classes rather than kept by hand.

2. **Packages come from PyPI.** The Function node's *Setup* tab and the
   palette manager install Python distributions. A PyFlowRED node package is an
   ordinary distribution that advertises a `pyflowred.nodes` entry point.

And one thing it adds, which follows from the first:

3. **The flow file can be Python too.** Point `flowFile` at a `.py` and the
   flows are kept as code — one file, every tab, no JSON. The editor is
   unchanged and still where flows are drawn; a deploy writes the file back.

   ```python
   readings = f.mqtt_in("readings", topic="plant/+/temp", broker=plant)
   lathe, mill, other = by_machine       # a Switch with three outputs

   readings >> by_machine
   lathe >> archive
   attempt[1] >> waited                  # a cycle needs no special syntax
   ```

   It runs on its own — `python flow.py` starts the engine — and a change
   reads in a diff. See
   [docs/porting.md](https://github.com/DataVoyage/pyflowred/blob/main/docs/porting.md#flows-as-python)
   and
   [examples/python-flows](https://github.com/DataVoyage/pyflowred/blob/main/examples/python-flows/README.md).

## Installing

```bash
pip install pyflowred     # or: uv pip install pyflowred
pyflowred
```

One name throughout: the distribution on PyPI, the command, the import
package, the `PYFLOWRED_*` settings and the editor are all `pyflowred`.
(`pyred` was already taken on PyPI, which is where the name came from; up to
0.1.0 the import package and command were still `pyred`, and 0.2.0 renamed
them to match.)

The wheel carries the built editor, so this needs Python and nothing else —
no Node, no build step. The editor is then at <http://127.0.0.1:1880/>.

## In a container

Two images, same application. `Dockerfile` installs the released package with
`uv` and carries no build tooling; `Dockerfile.source` builds this checkout in
three stages, Node included.

```bash
docker build -t pyflowred:0.8.0 .                        # released
docker build -f Dockerfile.source -t pyflowred:source .  # this checkout
docker run -d -p 1880:1880 -v pyflowred-data:/data pyflowred:0.8.0
```

Your own CA is a build argument or a mount, and reaches every library in the
runtime. See [docker/README.md](https://github.com/DataVoyage/pyflowred/blob/main/docker/README.md).

## How fast, measured against Node-RED

`tools/iotload.py` drives the same IoT ingest — three MQTT topic groups, two
Kafka topics, each split three ways — into this runtime and into Node-RED
5.0.7, steady and in bursts. What it found:

* **Throughput is one over the CPU a message costs, and nothing else.** Three
  settings of injected work predicted the observed ceiling to within a few
  percent, which makes sizing arithmetic: sustainable rate ≈ 1000 / (CPU ms
  per message), and `pyflowred health` reports that number.
* **A fan-out costs what its branches cost.** Both runtimes drive their flows
  from one event loop, so a split schedules its branches — but nothing is
  serialised that would otherwise have run in parallel.
* **The gap to Node-RED is raw speed, not scheduling**: about 3× with no work
  in the flow, 1.06× with half a millisecond of it, and the other way round
  with two. The more a flow actually does, the less the language matters.
* **Under a burst this stays answerable and Node-RED does not** — 479 ms
  against 10 517 ms at the 95th percentile while draining five thousand
  messages.

Numbers, method and caveats in
[docs/benchmarks.md](https://github.com/DataVoyage/pyflowred/blob/main/docs/benchmarks.md).

## Is it keeping up?

```bash
pyflowred health                         # the local instance
pyflowred health https://host/pyflowred  # any instance
```

```
  [ok  ] flows         running
  [ok  ] event_loop    0.0% of 60 samples over 100 ms, worst 1 ms
  [ok  ] cpu           2% of one core, about 49x this load fits
  [ok  ] throughput    312.3 messages/s received, 65 us of CPU each
  [ok  ] node_errors   0.00 node errors/s
  [ok  ] backpressure  0 in flight, deepest node queue 3
```

Exit code 0 healthy, 1 degraded, 2 failing, 3 unreachable. The same report is
at `<httpAdminRoot>health` as JSON — 200 while the instance is doing its job,
503 when it is not — for a Kubernetes probe or an uptime monitor. The runtime
is one asyncio loop, so what matters is not the core count but whether that
loop keeps up; see
[docs/operations.md](https://github.com/DataVoyage/pyflowred/blob/main/docs/operations.md#health-and-whether-one-core-is-enough).

## Building it

Node.js is needed once, to build the editor. It is never needed at runtime.

```bash
npm install          # editor build dependencies
npm run build        # builds src/pyflowred/editor_client/public
uv build --wheel     # ~10 MB, py3-none-any
```

That build output is a git-ignored artifact, so the wheel has to pull it back
in deliberately (`tool.hatch.build.targets.wheel.artifacts`). A build hook
refuses to produce a wheel when the editor has not been built, because the
failure mode otherwise is a package that installs, starts, and serves a blank
page. The editor *source* is excluded from the wheel and kept in the sdist,
where it can be rebuilt.

## Working on it

```bash
npm install && npm run build
uv venv --python 3.13
uv pip install -e ".[dev]"
.venv/bin/pyflowred
```

> Without `adminAuth`, anyone who can reach the editor can run arbitrary
> Python as the user running the process. See
> [docs/operations.md](https://github.com/DataVoyage/pyflowred/blob/main/docs/operations.md#protecting-the-editor).

## Writing a node

```python
def register(RED):
    class LowerCaseNode(RED.Node):
        def __init__(self, config):
            super().__init__(config)
            self.on("input", self._on_input)

        def _on_input(self, msg, send, done):
            msg["payload"] = msg["payload"].lower()
            send(msg)
            done()

    RED.nodes.register_type("lower-case", LowerCaseNode)
```

with the matching `.html` next to it, exactly as in Node-RED. A complete
example package is in `examples/pyflowred-node-lowercase`; the details are in
[docs/writing-nodes.md](https://github.com/DataVoyage/pyflowred/blob/main/docs/writing-nodes.md).

## Layout

```
src/pyflowred/
  util/            property expressions, JSONata, mustache, cloning, logging, i18n
  runtime/         flow engine, node registry, context, credentials, storage
  registry/        node discovery, loading, palette installation
  api/             HTTP admin API, editor routes, comms WebSocket (Starlette)
  metrics/         the Prometheus registry, runtime instrumentation, endpoint
  lsp/             Python language support for the Function node's editor
  flow/            flows declared in Python: the library, the generator, the storage
  streaming/       keyed state on RocksDB, windows, aggregate functions
  nodes/core/      the core node set: Python implementations + Node-RED dialogs
  nodes/core/examples/  a worked example flow per added node type
  editor_client/   the Node-RED editor, vendored and built
tools/
  editor-build/    the Node-RED editor build pipeline, re-pointed at this repo
  devserver.sh     start/stop a throwaway instance
  uicheck.py       load the editor in a real browser and report on it
  uiflow.py        deploy and run a flow through the editor, in a real browser
  uiprojects.py    commit a deploy through the project UI, in a real browser
  uigcs.py         check the Cloud Storage dialogs are typed inputs, in a browser
  uistream.py      check the stream processing dialogs, in a browser
  uilsp.py         check the Function node's completion and problems, in a browser
  uiflows.py       check a deploy from the editor reaches the Python flow file
  exampleflows.py  the per-node example flows the editor offers on Import
  e2eflows.py      the end-to-end integration test: 18 tabs, 157 assertions
  coverage.py      what a run reached: node types and dialog choices
  e2e/             the MQTT broker and Kafka the test talks to
  benchflows.py    seven benchmark flows, 82 nodes across 33 node types
  iotflows.py      an IoT ingest that splits three ways, for either runtime
  iotload.py       drive it over MQTT and Kafka, steady and in bursts
  loadtest.py      drive them and sample the runtime's CPU, RSS and fds
  hopprofile.py    take one node hop apart: what each part of the path costs
  rawbaseline.py   a bare Starlette endpoint, to measure the runtime against
Dockerfile         the released package, installed with uv
Dockerfile.source  this checkout, editor and all
docker/            entrypoint, compose, CA drop-in
examples/
  e2e-stack/       the integration test as a compose stack you can bring up
  python-flows/    a flow file written as Python: a cycle, a subflow, a Switch
  kafka-iot/       a whole plant on one Kafka topic: discovery, KPIs, no configuration
  pyflowred-node-lowercase/  a node package, complete and installable
docs/              porting notes, node authoring, operations, benchmarks
docs/nodes/        reference for the twenty-four node types this adds
tests/             the test suite, including a browser test of the editor
```

## Tests

```bash
.venv/bin/python -m pytest tests
```

[docs/test-coverage.md](https://github.com/DataVoyage/pyflowred/blob/main/docs/test-coverage.md) is the standing record of what
is tested, how, and what is not - including the 24% of end-to-end assertions
that so far only check that something arrived rather than what it was.

The suite covers every node group against real sockets (including an
in-process MQTT broker), the runtime (subflows, groups, context, credentials,
deploy types), the admin API, projects against a real `git` binary, and — when
Playwright is available — the actual editor in a headless browser: import a
flow, press Deploy, open the Function node, click Inject, read the debug
sidebar.

On top of that there is an **end-to-end integration test** that runs against a
real service rather than in-process: eighteen tabs, 803 nodes, 157 assertions.
One command brings the whole thing up in containers — application, MQTT,
Kafka, and a one-shot that deploys the flows and runs the checks:

```bash
docker compose -f examples/e2e-stack/compose.yml up -d
# the editor, with all eighteen tabs in it:
open http://127.0.0.1:1881/pyflowred/
```

It is also the quickest way to *read* this project: every node type it adds is
wired into a working flow with the assertion that proves it sitting next to
it. See
[examples/e2e-stack](https://github.com/DataVoyage/pyflowred/blob/main/examples/e2e-stack/README.md).

## Licence

Apache-2.0, as Node-RED is. The vendored editor and node dialogs keep their
original copyright headers.
