Files
app/backend/fluksio/flow/schemas.py
T
stroblmeandClaude Fable 5 a38e2745eb Add a Python SDK: flows declared in your own repository
A data scientist keeps their code where it is and decorates it: `@node`
declares a function's ports beside the function, `Flow(name, nodes=[...])`
says which of them make a flow, and `use(fn, wire=..., **settings)` rebinds
one for a single flow. `fluksio sync` uploads the document plus a generated
import shim per node, so the store still holds a complete, runnable,
git-versioned definition while the code it imports stays theirs.

`fluksio login|run|runs` and `flow.submit().wait()` are the client half, over
the run endpoints that already existed. Runs record the user repository's
commit beside the store's, so "what code produced this number" is answerable
on the side that now holds the code.

- `fluksio/sdk/`: ports, decorators, the flow builder and its checks, the shim
  generator, an HTTP client and sync. Standard library only at import, so
  `from fluksio import node` in a training script pulls in no engine.
- `FlowDef.origin` marks a flow code-defined; `Run.origin_commit` carries the
  repository's commit; `POST /modules/refresh` retires the workers without an
  install, which every sync calls — a worker holds the imported package in
  memory, so an edit to it is invisible until the process goes.
- The canvas shows a generated body read-only and names the repository to edit
  instead; a body edited there stops the next sync rather than being discarded.
- The worker's reporter carries inert `Port`, `node`, `use` and `Flow`, since
  the shim imports a module whose first line declares them.
- `examples/myresearch` is the worked example, `make sync-example` uploads it.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_012ue1tkFWB1bcGy3aWhCKpU
2026-08-23 20:16:08 +02:00

312 lines
9.5 KiB
Python

"""The persisted shape of a flow, shared by the store, the API and the editor.
A flow is structure plus code: this module is the structure. Node logic for
``python`` nodes lives beside it as a plain ``.py`` file.
"""
from __future__ import annotations
import re
from typing import Any, Literal
from pydantic import BaseModel, Field, field_validator
from fluksio.flow.messages import MessageSpec
NAME_PATTERN = re.compile(r"^[a-z][a-z0-9_]*$")
def _validate_name(value: str) -> str:
if not NAME_PATTERN.match(value):
raise ValueError(
"Use lowercase letters, digits and underscores, starting with a letter"
)
return value
class NodeDef(BaseModel):
"""A node as stored: identity, configuration and ports.
Deliberately no canvas position. The editor lays a flow out itself, so
where a node sits is a fact about the drawing rather than about the flow —
and a graph nobody can arrange is one worth keeping small.
"""
id: str
type: str = "python"
title: str = ""
#: This node's settings: constants of its function, stored with the flow.
#: A function node reads them as keyword arguments beside its ports, so a
#: setting cannot share a name with one.
params: dict[str, Any] = Field(default_factory=dict)
requires: list[MessageSpec] = Field(default_factory=list)
provides: list[MessageSpec] = Field(default_factory=list)
#: Name of a shared source in the library, instead of this node's own file.
#: Editing it edits the copy every flow using it runs.
source_ref: str | None = None
timeout: float | None = Field(
default=None,
gt=0,
description=(
"Seconds this node's code may run before it is stopped. This "
"covers the first call's imports, which can be much slower than "
"the body. Above 60 the engine may deliver its work again while "
"it is still running — in a batch run, which never redelivers, "
"it is an idle timeout instead: silence this long is a kill."
),
)
device: str | None = Field(
default=None,
description=(
"Label of the worker this node's code must run on, such as 'gpu'. "
"Empty means the engine's own workers. A run needing a label no "
"attached worker carries waits rather than failing."
),
)
device_policy: Literal["require", "prefer"] = Field(
default="require",
description=(
"What to do when no worker carries `device`: wait for one, or run "
"locally anyway."
),
)
@field_validator("id")
@classmethod
def _check_id(cls, value: str) -> str:
return _validate_name(value)
class FlowOrigin(BaseModel):
"""Where a flow was declared, when that was somewhere other than here.
A flow drawn on the canvas has no origin: the store is where it lives. One
stamped with this was declared with the decorators in somebody's own
repository and put here by ``fluksio sync``, so the node bodies below it
are generated imports and the code they run is versioned twice — once here
and once there. Its presence is what makes a flow code-defined.
Deliberately no timestamp. The store commits every change it is given, so
when a flow was last synced is a fact its own history already holds — and
one that would otherwise change on every sync, making an unchanged upload
look like a new version of the flow.
"""
kind: Literal["python"] = "python"
#: The repository root on the machine that ran ``sync``.
repo: str = ""
#: Its commit, and whether the tree had uncommitted changes at the time —
#: a run stamped with a dirty commit names code that was never stored.
commit: str = ""
dirty: bool = False
class FlowInput(BaseModel):
"""A message the flow starts with rather than computes."""
spec: MessageSpec
initial: Any | None = None
class FlowDef(BaseModel):
"""One atomic flow."""
name: str
title: str = ""
nodes: list[NodeDef] = Field(default_factory=list)
inputs: list[FlowInput] = Field(default_factory=list)
version: int = 1
mode: Literal["live", "batch"] = Field(
default="live",
description=(
"A live flow reacts to what arrives: its subscriptions, schedules "
"and webhooks run until it is stopped. A batch flow only runs when "
"a run asks it to, from its inputs to its outputs, and is never "
"activated."
),
)
outputs: list[str] = Field(
default_factory=list,
description=(
"Messages a batch run reports as its result, unqualified. Empty "
"means every message the flow ends up holding."
),
)
origin: FlowOrigin | None = Field(
default=None,
description=(
"Set when the flow was declared in code elsewhere and uploaded by "
"`fluksio sync`. Absent for a flow drawn on the canvas."
),
)
@field_validator("name")
@classmethod
def _check_name(cls, value: str) -> str:
return _validate_name(value)
class NodeSource(BaseModel):
"""The Python source of a node."""
code: str
Health = Literal["ok", "degraded", "down"]
class NodeStatusPublic(BaseModel):
"""Whether a node loaded, and how its connection is doing."""
id: str
status: str = "active"
error: str | None = None
health: Health = "ok"
health_detail: str | None = None
#: The node's last runtime failure, kept after it runs again: a failure
#: that fired an alert should leave a trace of what it was.
last_error: str = ""
last_error_ts: float | None = None
class MessageValue(BaseModel):
"""The last payload seen on a message."""
value: Any = None
ts: float | None = None
class HistoryPoint(BaseModel):
"""One numeric value a message carried, and when."""
ts: float
value: float
class MessageHistory(BaseModel):
"""A message's recent numeric values, oldest first.
Only numbers are recorded, so ``numeric`` tells the panel whether an empty
series means "nothing plottable here" or "nothing has arrived yet".
"""
message: str
numeric: bool = False
points: list[HistoryPoint] = Field(default_factory=list)
class FlowSummary(BaseModel):
name: str
title: str = ""
node_count: int = 0
error_count: int = 0
has_draft: bool = False
enabled: bool = True
paused: bool = False
# Its background tasks kept crashing, so the engine stopped restarting them.
quarantined: bool = False
#: Of the working copy, so publishing from a list needs no second read.
version: int = 1
class FlowsPublic(BaseModel):
data: list[FlowSummary]
count: int
class LibraryNode(BaseModel):
"""A node source shared across flows, and who is using it."""
name: str
used_by: list[str] = Field(default_factory=list)
class FlowStatePublic(BaseModel):
values: dict[str, MessageValue] = Field(default_factory=dict)
nodes: list[NodeStatusPublic] = Field(default_factory=list)
class ModulePackage(BaseModel):
"""One package installed in the venv node code runs on."""
name: str
version: str
class ModulesInfo(BaseModel):
"""The venv node code imports from, and the manifest that describes it."""
python_version: str = ""
venv_path: str = ""
requirements: str = ""
packages: list[ModulePackage] = Field(default_factory=list)
#: Whether what is installed matches the manifest.
applied: bool = False
class ApplyRequest(BaseModel):
"""A pip manifest, one requirement per line."""
requirements: str = ""
class ApplyResult(BaseModel):
ok: bool
output: str = ""
class BrainNode(BaseModel):
"""One neuron: a thing the engine talks to, or a node that only computes.
Nodes of the same type pointing at the same outside thing — one broker
topic, one URL, one bucket — are a single entry here, whichever flows they
sit in. ``members`` are the ``flow.node_id`` names behind it, which is also
what the live events are keyed by.
"""
id: str
label: str
kind: str
members: list[str] = Field(default_factory=list)
flows: list[str] = Field(default_factory=list)
issue: str | None = Field(
default=None,
description=(
"Why this neuron cannot run, if validation found something. A "
"failure the engine hits while running arrives over the socket "
"instead; this is the part that is already true before anything "
"fires, and so has to travel with the graph."
),
)
class BrainEdge(BaseModel):
"""Messages carrying values from one neuron to another."""
source: str
target: str
messages: list[str] = Field(default_factory=list)
class BrainGraph(BaseModel):
"""Every published flow at once, merged on what its nodes talk to."""
nodes: list[BrainNode] = Field(default_factory=list)
edges: list[BrainEdge] = Field(default_factory=list)
class NodeTypeInfo(BaseModel):
"""A node type the editor can offer, with its parameter schema."""
type: str
title: str
description: str
params_schema: dict[str, Any] = Field(default_factory=dict)
has_source: bool = False
#: Whether this type takes settings beyond the ones its schema declares.
#: A function node's settings are its author's to name, and reach `process`
#: as keyword arguments beside its ports.
free_params: bool = False
#: The package a connector came from; empty for the built-in types.
plugin: str | None = None