Files
app/backend/tests/test_metrics.py
T
stroblmeandClaude Opus 5 2554488a73 Backend: real alert test results, state cleanup on delete/rename, node trigger errors, queue and collector fixes
`AlertManager.send` swallowed every delivery failure, so the alerts screen's
Test button answered 200 whatever happened — the one thing it exists for. It
takes `raise_on_error` now, which only the test route passes; the per-channel
loop keeps the swallow, because one dead channel must not stop the others
hearing about the same fault. A refused delivery answers 502 with whatever the
sender said.

Renaming a flow left its values under the old name for good: the delete path
already swept them, the rename path never did. It calls the same `forget_flow`,
which covers the messages and the `__ts__`/`__version__`/`__history__`
bookkeeping keyed by message name. Cleanup, not migration — they repopulate
under the new name on the next run.

Triggering a node by hand ran `Node.__call__` with nothing catching it, so a
node that raised produced a 500 and a stack trace in the server log, and
nothing at all on the canvas. `Pipeline.publish_error` is the reporting half of
`_execute_node` lifted out; both paths go through it, so a manual failure now
reads the same on the canvas and in the metrics as a queued one. The route
answers 400 with the node's error.

`MemoryWorkQueue.stats()` counts claimed-but-unacknowledged work rather than
reporting zero, so the health tile means something without Redis. The metrics
collector's held tracebacks are capped at `DETAIL_CAP` and swept on the same
`RUN_STALE_S` cutoff the open runs use, instead of one untruncated traceback
per node kept for the life of the process — a traceback still survives the
flush between the log and the failure it belongs to.

`GET /observability/events` takes `since`/`until`, the window `/runs` already
took, so a failures list can cover the span the charts beside it are drawn from.

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

227 lines
6.8 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 DETAIL_CAP, RUN_STALE_S, 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")
def test_a_traceback_no_failure_ever_claims_is_dropped() -> None:
"""One entry per node that ever errored would else last the process out."""
node = f"{FLOW}.stale"
collector = MetricsCollector(EventBus())
collector.handle(
{
"type": "node_log",
"flow": FLOW,
"node": node,
"level": "error",
"text": "x" * (DETAIL_CAP * 2),
"ts": time.time() - RUN_STALE_S - 1,
}
)
# Untruncated, one of these holds a whole run of a chatty node.
assert len(collector._tracebacks[(FLOW, node)][1]) == DETAIL_CAP
# Something to write, so the flush does not stop at its early return.
collector.handle(
{"type": "node_executed", "flow": FLOW, "node": node, "ts": time.time()}
)
asyncio.run(collector.flush())
assert collector._tracebacks == {}