Files
app/backend/fluksio/core/config.py
T
stroblmeandClaude Opus 5 180da3d640 Stop paying five Redis round trips and a global lock per message
The engine was I/O-bound on its own state backend. `RedisState.lock()` is one
key — `pipeline:_lock` — for the whole process, taken five times a message at
two round trips each, and every cascade and every node read queued behind it.
Inside it, reading a node's inputs was three round trips per input (an EXISTS
for `in`, then EXISTS and GET for the value), writing was two updates that a
single transaction already gives, and the version counters went one INCR at a
time.

Replaced with the atomic command that was always available: `get_present` is
one MGET and tells a missing key from one holding null, so the lock it used to
be read under bought nothing; value and timestamp land in one `update`, which
is a MULTI/EXEC; `increment_multi` pipelines the counters. `values()` — what
every websocket snapshot calls — is two reads whatever the message count
instead of two per message.

Beside that: every webhook did its blocking XADD on the asyncio event loop
(MQTT already used `to_thread`); the per-execution `NodeOutcome` was built and
validated even with no run watching; `_minute` built a tz-aware datetime per
event on the loop thread to key a dict, and now keys on an int; `move_due`
promoted delayed items one round trip each, every second; `FLOW_MAX_CASCADES`
makes the in-flight ceiling a setting rather than a constant.

`orjson` replaces stdlib json where a message pays for it — state, the
journal, the engine side of the worker pipe. `fluksio-worker` stays
dependency-free, and the run-cache digest stays on stdlib so no stored key is
invalidated. A non-finite number now stores as `null` rather than the bare
`NaN` that was never JSON.

Measured with `scripts/bench_engine.py` against a real Redis, 200 messages:
a five-node chain went from 43.9 to 103.1 msg/s with p50 latency 2110ms →
782ms and p95 3913ms → 1439ms; one source into twenty consumers went from 5.4
to 33.7 msg/s. In memory, twenty consumers went from 187 to 448 msg/s.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01BpfSinyCBfjuieikyfMPbf
2026-08-26 10:12:25 +02:00

224 lines
8.8 KiB
Python

import os
import secrets
import warnings
from pathlib import Path
from typing import Annotated, Any, Literal, Self
from pydantic import (
AnyUrl,
BeforeValidator,
EmailStr,
HttpUrl,
computed_field,
model_validator,
)
from pydantic_settings import BaseSettings, SettingsConfigDict
#: Everything the engine keeps on disk, relative to :attr:`Settings.DATA_DIR`.
#: One setting to move the lot; each still overridable on its own, which is
#: what the container images do.
DERIVED_PATHS = {
"FLOWS_DIR": "flows",
"SECRETS_FILE": "secrets.enc",
"ALERTS_FILE": "alerts.json",
"PANELS_FILE": "panels.json",
"OAUTH_PRIVATE_KEY_FILE": "oauth-key.pem",
"CLOUD_CONFIG_FILE": "cloud.json",
}
def parse_cors(v: Any) -> list[str] | str:
if isinstance(v, str) and not v.startswith("["):
return [i.strip() for i in v.split(",") if i.strip()]
elif isinstance(v, list | str):
return v
raise ValueError(v)
class Settings(BaseSettings):
model_config = SettingsConfigDict(
# The stack's own file, one level above ./backend/. An installed
# `fluksio` has no such tree, so its CLI points this at the data
# directory instead — and at nothing it might find in the cwd.
env_file=os.environ.get("FLUKSIO_ENV_FILE", "../.env"),
env_ignore_empty=True,
extra="ignore",
)
API_V1_STR: str = "/api/v1"
SECRET_KEY: str = secrets.token_urlsafe(32)
# 60 minutes * 24 hours * 8 days = 8 days
ACCESS_TOKEN_EXPIRE_MINUTES: int = 60 * 24 * 8
FRONTEND_HOST: str = "http://localhost:5173"
ENVIRONMENT: Literal["local", "staging", "production"] = "local"
#: Everything this installation keeps: the database, the flow repository,
#: secrets, artifacts and the user venv. The paths below derive from it
#: unless they are set explicitly.
DATA_DIR: Path = Path("flow-data")
#: Any SQLAlchemy URL. The default puts SQLite in the data directory, which
#: is what makes `fluksio serve` need no infrastructure at all.
DATABASE_URL: str | None = None
# Flows live on disk as a git repository; secrets stay outside it.
FLOWS_DIR: Path = Path("flow-data/flows")
SECRETS_FILE: Path = Path("flow-data/secrets.enc")
# Which failures reach which channel. Beside the flows, not in them:
# alerting is the deployment's concern, not any one flow's.
ALERTS_FILE: Path = Path("flow-data/alerts.json")
# Which dashboards each device shows. Beside the flows for the same reason
# alerting is: where a screen hangs is the deployment's concern rather than
# any one dashboard's.
PANELS_FILE: Path = Path("flow-data/panels.json")
# Which interpreter node code runs on. "auto" adopts the venv the engine
# was installed into, when it was installed into one and there is no venv
# of its own to lose — which is the `pip install fluksio` beside your own
# packages case. "managed" always builds a separate one, which is what a
# container wants. A path names an interpreter outright.
NODE_VENV: str = "auto"
# The MCP endpoint, and the OAuth server agents authenticate against. Off
# until someone asks for it: it opens client registration to the network.
MCP_ENABLED: bool = False
# Unauthenticated test-only endpoints (user seeding). Requires an explicit
# opt-in on top of ENVIRONMENT=local, so a deployment that merely kept the
# default environment never exposes them.
PRIVATE_API_ENABLED: bool = False
DOMAIN: str = "localhost"
OAUTH_PRIVATE_KEY_FILE: Path = Path("flow-data/oauth-key.pem")
# Written only when someone enrols this installation with a portal.
# Its absence is what keeps remote access off.
CLOUD_CONFIG_FILE: Path = Path("flow-data/cloud.json")
OAUTH_CODE_EXPIRE_SECONDS: int = 60
# Short, because an agent's token is a bearer secret held by a program
# rather than a person, and it can refresh unattended.
MCP_TOKEN_EXPIRE_MINUTES: int = 60
MCP_REFRESH_EXPIRE_DAYS: int = 30
FLOW_MAX_WORKERS: int = 4
# How many cascades may be in flight at once. Sustained throughput is this
# over the mean cascade time, so an installation whose nodes wait on the
# network rather than on a CPU wants it higher than the core count.
FLOW_MAX_CASCADES: int = 4
# How long a python node may be silent before its worker is killed, unless
# the node sets its own. 0, the default, disables it: a dead worker still
# fails fast, and a slow one is left to finish. Set it where silence means
# stuck rather than working.
FLOW_NODE_TIMEOUT: float = 0.0
# How long the engine's own metrics, events and run records are kept.
OBS_RETENTION_DAYS: int = 30
# Without a Redis host the engine keeps its state in memory.
REDIS_HOST: str | None = None
REDIS_PORT: int = 6379
BACKEND_CORS_ORIGINS: Annotated[
list[AnyUrl] | str, BeforeValidator(parse_cors)
] = []
@model_validator(mode="before")
@classmethod
def _derive_data_paths(cls, data: Any) -> Any:
"""Put every stored thing under ``DATA_DIR`` unless it was named.
``setdefault``, so the container images keep their explicit ``/data``
paths and a checkout keeps ``flow-data/``.
"""
if not isinstance(data, dict):
return data
base = Path(str(data.get("DATA_DIR", "flow-data"))).expanduser()
data["DATA_DIR"] = base
for key, name in DERIVED_PATHS.items():
data.setdefault(key, base / name)
return data
@computed_field # type: ignore[prop-decorator]
@property
def oauth_issuer(self) -> str:
"""Who issues MCP tokens — this app, on its API host.
Kept separate from the app's own URL because a hosted deployment can
later point agents at a different issuer without the resource server
changing: it validates whatever issuer it is configured to trust.
"""
scheme = "http" if self.ENVIRONMENT == "local" else "https"
return f"{scheme}://api.{self.DOMAIN}"
@computed_field # type: ignore[prop-decorator]
@property
def mcp_resource(self) -> str:
"""The resource an MCP token is issued for (RFC 8707)."""
return f"{self.oauth_issuer}/mcp"
@computed_field # type: ignore[prop-decorator]
@property
def all_cors_origins(self) -> list[str]:
return [str(origin).rstrip("/") for origin in self.BACKEND_CORS_ORIGINS] + [
self.FRONTEND_HOST
]
PROJECT_NAME: str = "Fluksio"
SENTRY_DSN: HttpUrl | None = None
@computed_field # type: ignore[prop-decorator]
@property
def SQLALCHEMY_DATABASE_URI(self) -> str:
"""SQLite in the data directory, unless a URL says otherwise.
One engine process owns this database — the same reason the image runs
a single uvicorn worker — so a file beside the flows is the honest
shape for it, and needs nothing running to be one.
"""
if self.DATABASE_URL:
return self.DATABASE_URL
return f"sqlite:///{(self.DATA_DIR / 'fluksio.db').expanduser().resolve()}"
SMTP_TLS: bool = True
SMTP_SSL: bool = False
SMTP_PORT: int = 587
SMTP_HOST: str | None = None
SMTP_USER: str | None = None
SMTP_PASSWORD: str | None = None
EMAILS_FROM_EMAIL: EmailStr | None = None
EMAILS_FROM_NAME: str | None = None
@model_validator(mode="after")
def _set_default_emails_from(self) -> Self:
if not self.EMAILS_FROM_NAME:
self.EMAILS_FROM_NAME = self.PROJECT_NAME
return self
EMAIL_RESET_TOKEN_EXPIRE_HOURS: int = 48
@computed_field # type: ignore[prop-decorator]
@property
def emails_enabled(self) -> bool:
return bool(self.SMTP_HOST and self.EMAILS_FROM_EMAIL)
EMAIL_TEST_USER: EmailStr = "test@example.com"
# Absent means "the CLI will make one on first run" — a pip install is not
# asked for two environment variables before it can start.
FIRST_SUPERUSER: EmailStr | None = None
FIRST_SUPERUSER_PASSWORD: str | None = None
def _check_default_secret(self, var_name: str, value: str | None) -> None:
if value == "changethis":
message = (
f'The value of {var_name} is "changethis", '
"for security, please change it, at least for deployments."
)
if self.ENVIRONMENT == "local":
warnings.warn(message, stacklevel=1)
else:
raise ValueError(message)
@model_validator(mode="after")
def _enforce_non_default_secrets(self) -> Self:
self._check_default_secret("SECRET_KEY", self.SECRET_KEY)
self._check_default_secret(
"FIRST_SUPERUSER_PASSWORD", self.FIRST_SUPERUSER_PASSWORD
)
return self
# No arguments and no required environment: a fresh install boots on the
# defaults above, into a data directory of its own.
settings = Settings()