A second bus subscriber folds executions, errors, timings and queue lag into per-minute rollups, keeps failures with their traceback and an audit trail of who published what, and records one row per cascade — manual runs and previews included, under an id of their own that writes no idempotency markers. Read back through /observability/*, which always answers 200 so a degraded engine still renders its own health screen. Also fixes two things found on the way: node-health alerts read `status` where the engine publishes `health`, so a device dropping never alerted anyone, and the Redis queue reported `parked: 0` whatever was held. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017MeiWk3Yq12n2pTvnQWYvt
150 lines
4.5 KiB
Python
150 lines
4.5 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_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")
|