Files
app/backend/tests/test_metrics.py
T
stroblmeandClaude Fable 5 83c30aa1c7 Take the engine off the path node code imports from, and let it stop
The worker script is handed to the interpreter by path, so app/flow was
sys.path[0] for every node: `import queue` got the engine's. It now drops
its own directory before anything else imports, and runs with the
deployment's credentials scrubbed out of its environment.

Also: reload builds off the event loop, the pool wakes what is blocked on
it when it stops, a refused metrics flush is kept for the next one rather
than dropped, and the cascade events are paired through failures.

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

200 lines
5.9 KiB
Python

"""The collector writes down what the bus only ever broadcast.
Not under ``tests/flow`` with the rest of the engine: that package opts out of
the database, and writing to it is this module's whole job.
"""
import asyncio
import time
from datetime import datetime, timezone
from sqlmodel import Session, select
from app.flow.events import EventBus
from app.flow.metrics import MetricsCollector
from app.models import EngineEvent, FlowRun, MetricBucket
FLOW = "metrics-test"
NODE = "metrics-test.calc"
def _events(collector: MetricsCollector, ts: float, run: str) -> None:
collector.handle(
{
"type": "cascade_started",
"run": run,
"flow": FLOW,
"node": "metrics-test.in",
"cause": "external",
"deliveries": 1,
"ts": ts,
}
)
collector.handle(
{
"type": "node_executed",
"flow": FLOW,
"node": NODE,
"outputs": 2,
"duration_ms": 5.0,
"run": run,
"ts": ts,
}
)
collector.handle(
{"type": "work_latency", "flow": FLOW, "node": NODE, "lag_ms": 12.0, "ts": ts}
)
# The traceback arrives as its own event, just before the failure.
collector.handle(
{
"type": "node_log",
"flow": FLOW,
"node": NODE,
"level": "error",
"text": "Traceback: line 3, in run",
"ts": ts,
}
)
collector.handle(
{
"type": "node_error",
"flow": FLOW,
"node": NODE,
"error": "ValueError: bad input",
"run": run,
"ts": ts,
}
)
def test_events_become_rollups_failures_runs_and_audit(db: Session) -> None:
collector = MetricsCollector(EventBus())
# The current minute, pinned: the second flush has to land in the same
# bucket, and anything past the retention window is pruned on write.
ts = datetime.now(timezone.utc).replace(second=0, microsecond=0).timestamp()
_events(collector, ts, "1-0")
collector.handle(
{
"type": "audit",
"action": "published",
"flow": FLOW,
"user": "someone@example.com",
"ts": ts,
}
)
collector.handle(
{"type": "cascade_finished", "run": "1-0", "flow": FLOW, "ts": ts + 0.5}
)
asyncio.run(collector.flush())
bucket = db.exec(
select(MetricBucket).where(MetricBucket.flow == FLOW, MetricBucket.node == NODE)
).one()
assert (bucket.executions, bucket.errors, bucket.messages) == (1, 1, 2)
assert bucket.duration_max_ms == 5.0
assert (bucket.lag_max_ms, bucket.items) == (12.0, 1)
failure = db.exec(
select(EngineEvent).where(
EngineEvent.flow == FLOW, EngineEvent.type == "node_error"
)
).one()
assert "ValueError: bad input" in failure.detail
assert "line 3, in run" in failure.detail
audit = db.exec(
select(EngineEvent).where(EngineEvent.flow == FLOW, EngineEvent.type == "audit")
).one()
assert (audit.actor, audit.detail) == ("someone@example.com", "published")
run = db.exec(select(FlowRun).where(FlowRun.id == "1-0")).one()
assert (run.status, run.flow, run.source) == ("error", FLOW, "external")
assert (run.nodes, run.errors) == (1, 1)
assert run.duration_ms > 0
# The same minute, written again: the counters add rather than duplicate.
_events(collector, ts + 10, "2-0")
asyncio.run(collector.flush())
db.expire_all()
bucket = db.exec(
select(MetricBucket).where(MetricBucket.flow == FLOW, MetricBucket.node == NODE)
).one()
assert (bucket.executions, bucket.errors) == (2, 2)
# Never finished, so it is still open — the prune is what closes it.
open_run = db.exec(select(FlowRun).where(FlowRun.id == "2-0")).one()
assert open_run.status == "running"
def test_a_flush_the_database_refused_is_written_by_the_next_one(
db: Session, monkeypatch
) -> None:
collector = MetricsCollector(EventBus())
ts = datetime.now(timezone.utc).replace(second=0, microsecond=0).timestamp()
collector.handle(
{
"type": "audit",
"action": "held back",
"flow": FLOW,
"user": "held@example.com",
"ts": ts,
}
)
def refuse(*_args: object) -> None:
raise RuntimeError("the database is gone")
monkeypatch.setattr(collector, "_write", refuse)
asyncio.run(collector.flush())
monkeypatch.undo()
asyncio.run(collector.flush())
audit = db.exec(
select(EngineEvent).where(EngineEvent.actor == "held@example.com")
).one()
assert audit.detail == "held back"
def test_a_node_id_wider_than_the_column_still_records(db: Session) -> None:
collector = MetricsCollector(EventBus())
ts = datetime.now(timezone.utc).replace(second=0, microsecond=0).timestamp()
long_node = f"{FLOW}.{'w' * 400}"
collector.handle(
{
"type": "node_executed",
"flow": FLOW,
"node": long_node,
"duration_ms": 1.0,
"ts": ts,
}
)
asyncio.run(collector.flush())
bucket = db.exec(
select(MetricBucket).where(MetricBucket.node == long_node[:255])
).one()
assert bucket.executions == 1
def test_a_run_that_did_not_fail_reads_ok(db: Session) -> None:
collector = MetricsCollector(EventBus())
ts = time.time()
collector.handle(
{
"type": "cascade_started",
"run": "manual-abc",
"flow": FLOW,
"node": "metrics-test.in",
"cause": "manual",
"deliveries": 1,
"ts": ts,
}
)
collector.handle(
{"type": "cascade_finished", "run": "manual-abc", "flow": FLOW, "ts": ts + 0.1}
)
asyncio.run(collector.flush())
run = db.exec(select(FlowRun).where(FlowRun.id == "manual-abc")).one()
assert (run.status, run.source) == ("ok", "manual")