diff --git a/src/backend/base/langflow/services/background_execution/metrics.py b/src/backend/base/langflow/services/background_execution/metrics.py new file mode 100644 index 0000000000..f24eef93a6 --- /dev/null +++ b/src/backend/base/langflow/services/background_execution/metrics.py @@ -0,0 +1,44 @@ +"""Best-effort metric emission for background execution. Never raises into the runner.""" + +from __future__ import annotations + +from lfx.log.logger import logger + +from langflow.services.deps import get_telemetry_service + + +def _emit(fn_name: str, name: str, labels: dict, value: float | None = None) -> None: + try: + ot = get_telemetry_service().ot + fn = getattr(ot, fn_name) + if value is None: + fn(name, labels) + else: + fn(name, value, labels) + except Exception as exc: # noqa: BLE001 - observability must never break execution + logger.debug(f"bg metric emit skipped ({name}): {exc}") + + +def emit_job_started(*, backend: str) -> None: + _emit("increment_counter", "langflow_bg_jobs_started_total", {"backend": backend}) + + +def emit_job_completed(*, backend: str) -> None: + _emit("increment_counter", "langflow_bg_jobs_completed_total", {"backend": backend}) + + +def emit_job_failed(*, reason: str, backend: str) -> None: + _emit("increment_counter", "langflow_bg_jobs_failed_total", {"reason": reason, "backend": backend}) + + +def emit_orphan_reconciled(*, backend: str) -> None: + _emit("increment_counter", "langflow_bg_orphans_reconciled_total", {"backend": backend}) + + +def emit_job_duration(*, seconds: float, outcome: str, backend: str) -> None: + _emit( + "observe_histogram", + "langflow_bg_job_duration_seconds", + {"outcome": outcome, "backend": backend}, + value=seconds, + ) diff --git a/src/backend/tests/unit/services/background_execution/test_metrics_emit.py b/src/backend/tests/unit/services/background_execution/test_metrics_emit.py new file mode 100644 index 0000000000..76cdd48fd2 --- /dev/null +++ b/src/backend/tests/unit/services/background_execution/test_metrics_emit.py @@ -0,0 +1,49 @@ +"""Tests for the best-effort background-execution metric emit helpers. + +Two guarantees are exercised: + +1. The swallow path: when the telemetry service is unavailable (a real runtime + condition before the service manager is initialized), every ``emit_*`` helper + must return without raising. The ``monkeypatch`` here only simulates the + service being absent; it does not mock telemetry behavior under test. +2. The happy path: with the real telemetry service available, emitting reaches + the real OpenTelemetry instrument without raising, and the registry accepts + the labels we send (``validate_labels`` does not raise). +""" + +from langflow.services.background_execution import metrics as bgm +from langflow.services.deps import get_telemetry_service + + +def test_emit_swallows_when_telemetry_missing(monkeypatch): + """Every emit_* must swallow when the telemetry service cannot be resolved.""" + + def _raise(): + msg = "no service" + raise RuntimeError(msg) + + monkeypatch.setattr(bgm, "get_telemetry_service", _raise) + + # None of these may propagate the RuntimeError. + bgm.emit_job_started(backend="default") + bgm.emit_job_completed(backend="default") + bgm.emit_job_failed(reason="error", backend="default") + bgm.emit_orphan_reconciled(backend="default") + bgm.emit_job_duration(seconds=1.0, outcome="completed", backend="default") + + +def test_emit_reaches_real_instrument(): + """With the real telemetry service, emitting does not raise and labels validate.""" + bgm.emit_job_started(backend="default") + bgm.emit_job_completed(backend="default") + bgm.emit_job_failed(reason="error", backend="default") + bgm.emit_orphan_reconciled(backend="default") + bgm.emit_job_duration(seconds=1.0, outcome="completed", backend="default") + + # The labels each helper sends must be accepted by the real registry. + ot = get_telemetry_service().ot + ot.validate_labels("langflow_bg_jobs_started_total", {"backend": "default"}) + ot.validate_labels("langflow_bg_jobs_completed_total", {"backend": "default"}) + ot.validate_labels("langflow_bg_jobs_failed_total", {"reason": "error", "backend": "default"}) + ot.validate_labels("langflow_bg_orphans_reconciled_total", {"backend": "default"}) + ot.validate_labels("langflow_bg_job_duration_seconds", {"outcome": "completed", "backend": "default"})