Skip to content

Commit d525188

Browse files
authored
Merge branch 'release-1.11.0' into fix/message-history-order-direction
2 parents fcdb1cb + 89a1643 commit d525188

9 files changed

Lines changed: 281 additions & 28 deletions

File tree

src/backend/base/langflow/services/job_queue/service.py

Lines changed: 70 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -129,6 +129,9 @@ class JobQueueService(Service):
129129

130130
name = "job_queue_service"
131131

132+
# Backend label used on OpenTelemetry metrics. Subclasses override.
133+
_backend_label: str = "memory"
134+
132135
def __init__(self) -> None:
133136
"""Initialize the JobQueueService.
134137
@@ -142,6 +145,44 @@ def __init__(self) -> None:
142145
self._closed = False
143146
self.ready = False
144147
self.CLEANUP_GRACE_PERIOD = 300 # 5 minutes before cleaning up marked tasks
148+
# Cached OpenTelemetry handle (resolved lazily on first emit).
149+
# _otel_resolved is True after the first lookup attempt, including failures,
150+
# so we don't retry on every bump. _otel remains None when lookup fails.
151+
self._otel: Any = None
152+
self._otel_resolved: bool = False
153+
154+
def _get_otel(self) -> Any:
155+
"""Best-effort lookup of the OpenTelemetry singleton. Cached after first call.
156+
157+
Returns the ``ot`` handle, or ``None`` if telemetry is unavailable. Never raises —
158+
instrumentation must not break the queue.
159+
"""
160+
if self._otel_resolved:
161+
return self._otel
162+
self._otel_resolved = True
163+
try:
164+
from langflow.services.deps import get_telemetry_service
165+
166+
self._otel = get_telemetry_service().ot
167+
except Exception: # noqa: BLE001 - telemetry must not crash the queue
168+
self._otel = None
169+
return self._otel
170+
171+
def _emit_otel_counter(self, metric_name: str, labels: dict[str, str], value: float = 1.0) -> None:
172+
"""Best-effort OpenTelemetry counter bump. Silent on failure."""
173+
ot = self._get_otel()
174+
if ot is None:
175+
return
176+
with contextlib.suppress(Exception):
177+
ot.increment_counter(metric_name, labels, value)
178+
179+
def _emit_otel_up_down(self, metric_name: str, value: float, labels: dict[str, str]) -> None:
180+
"""Best-effort OpenTelemetry up-down counter delta. Silent on failure."""
181+
ot = self._get_otel()
182+
if ot is None:
183+
return
184+
with contextlib.suppress(Exception):
185+
ot.up_down_counter(metric_name, value, labels)
145186

146187
def is_started(self) -> bool:
147188
"""Check if the JobQueueService has started.
@@ -241,6 +282,7 @@ def create_queue(self, job_id: str) -> tuple[asyncio.Queue, EventManager]:
241282

242283
# Register the queue without an active task.
243284
self._queues[job_id] = (main_queue, event_manager, None, None)
285+
self._emit_otel_up_down("langflow_job_queue_active_jobs", 1, {"backend": self._backend_label})
244286
logger.debug(f"Queue and event manager successfully created for job_id {job_id}")
245287
return main_queue, event_manager
246288

@@ -436,7 +478,8 @@ async def cleanup_job(self, job_id: str) -> None:
436478

437479
await logger.adebug(f"Removed {items_cleared} items from queue for job_id {job_id}")
438480
# Remove the job entry from the registry
439-
self._queues.pop(job_id, None)
481+
if self._queues.pop(job_id, None) is not None:
482+
self._emit_otel_up_down("langflow_job_queue_active_jobs", -1, {"backend": self._backend_label})
440483
self._job_owners.pop(job_id, None)
441484
self._public_jobs.discard(job_id)
442485
await logger.adebug(f"Cleanup successful for job_id {job_id}: resources have been released.")
@@ -810,6 +853,8 @@ class RedisJobQueueService(JobQueueService):
810853
ACTIVITY_PREFIX = _ACTIVITY_PREFIX
811854
PUBLIC_JOB_PREFIX = _PUBLIC_JOB_PREFIX
812855

856+
_backend_label: str = "redis"
857+
813858
def __init__(
814859
self,
815860
host: str = "localhost",
@@ -880,6 +925,16 @@ def cross_worker_cancel_enabled(self) -> bool:
880925
"""
881926
return self._cancel_channel_enabled
882927

928+
def _bump_cancel_stat(self, key: str, value: int = 1) -> None:
929+
"""Increment a cancel_stat counter and emit the matching OTel counter event.
930+
931+
Single point of mutation for ``self._cancel_stats`` so ``/monitor/job_queue``
932+
and Prometheus stay in lockstep. The OTel emit is best-effort and never
933+
raises (see :meth:`_emit_otel_counter` on the base class).
934+
"""
935+
self._cancel_stats[key] += value
936+
self._emit_otel_counter("langflow_job_queue_cancel_events_total", {"event_type": key}, float(value))
937+
883938
def _stream_key(self, job_id: str) -> str:
884939
return f"{self.STREAM_PREFIX}{job_id}"
885940

@@ -1208,7 +1263,7 @@ async def touch_activity(self, job_id: str) -> None:
12081263
except asyncio.CancelledError:
12091264
raise
12101265
except Exception as exc: # noqa: BLE001
1211-
self._cancel_stats["activity_touch_errors"] += 1
1266+
self._bump_cancel_stat("activity_touch_errors")
12121267
await logger.adebug(f"touch_activity SET failed for {job_id}: {exc}")
12131268

12141269
def _owner_refresh_interval_s(self) -> float:
@@ -1248,7 +1303,7 @@ async def _check_pending_cancel_marker(self, job_id: str) -> None:
12481303
try:
12491304
if await self._client.exists(marker_key):
12501305
await self._client.delete(marker_key)
1251-
self._cancel_stats["marker_hit"] += 1
1306+
self._bump_cancel_stat("marker_hit")
12521307
await self._handle_cancel(job_id, source="marker")
12531308
except Exception as exc: # noqa: BLE001
12541309
await logger.awarning(f"Pending cancel marker check failed for {job_id}: {exc}")
@@ -1306,7 +1361,7 @@ async def _run_polling_watchdog(self) -> None:
13061361
try:
13071362
raw = await self._client.get(self._activity_key(job_id))
13081363
except Exception as exc: # noqa: BLE001
1309-
self._cancel_stats["activity_get_errors"] += 1
1364+
self._bump_cancel_stat("activity_get_errors")
13101365
await logger.adebug(f"polling watchdog: GET failed for {job_id}: {exc}")
13111366
continue
13121367
# Default to "infinitely stale" so a missed elif below cannot
@@ -1324,7 +1379,7 @@ async def _run_polling_watchdog(self) -> None:
13241379
try:
13251380
last = float(raw.decode() if isinstance(raw, bytes) else raw)
13261381
except (ValueError, TypeError) as exc:
1327-
self._cancel_stats["activity_parse_errors"] += 1
1382+
self._bump_cancel_stat("activity_parse_errors")
13281383
await logger.awarning(
13291384
f"polling watchdog: ignoring malformed activity value for {job_id}: {exc}"
13301385
)
@@ -1335,7 +1390,7 @@ async def _run_polling_watchdog(self) -> None:
13351390
# Stale → cancel this job. Local cancel on owned jobs skips the
13361391
# pubsub round-trip and stays correct even during a dispatcher
13371392
# reconnect window.
1338-
self._cancel_stats["polling_watchdog_kills"] += 1
1393+
self._bump_cancel_stat("polling_watchdog_kills")
13391394
await logger.ainfo(f"polling watchdog: reclaiming abandoned job {job_id} (age={age:.1f}s)")
13401395
await self._handle_cancel(job_id, source="watchdog")
13411396
with contextlib.suppress(Exception):
@@ -1383,7 +1438,7 @@ async def _run_cancel_dispatcher(self) -> None:
13831438
job_id = channel_str[len(self.CANCEL_CHANNEL_PREFIX) :]
13841439
await self._handle_cancel(job_id, source="pubsub")
13851440
# listen() returned cleanly — treat as a soft disconnect and reconnect.
1386-
self._cancel_stats["dispatcher_reconnects"] += 1
1441+
self._bump_cancel_stat("dispatcher_reconnects")
13871442
await logger.awarning(f"Cancel dispatcher pubsub.listen() ended; reconnecting in {backoff:.1f}s")
13881443
except asyncio.CancelledError:
13891444
with contextlib.suppress(Exception):
@@ -1393,15 +1448,15 @@ async def _run_cancel_dispatcher(self) -> None:
13931448
except (ConnectionError, TimeoutError, OSError) as exc:
13941449
# Expected transient failure: Redis dropped the pubsub, network
13951450
# blip, broker restart. Reconnect quietly via the backoff loop.
1396-
self._cancel_stats["dispatcher_reconnects"] += 1
1451+
self._bump_cancel_stat("dispatcher_reconnects")
13971452
await logger.awarning(f"Cancel dispatcher disconnect (retrying in {backoff:.1f}s): {exc!r}")
13981453
except Exception as exc: # noqa: BLE001
13991454
# Unexpected exception — likely a bug in dispatch logic, NOT a
14001455
# Redis problem. Surface at error level with traceback so it
14011456
# reaches Sentry / log aggregation, then still reconnect so a
14021457
# one-off bug doesn't kill cross-worker cancel permanently.
1403-
self._cancel_stats["dispatcher_reconnects"] += 1
1404-
self._cancel_stats["dispatcher_internal_errors"] += 1
1458+
self._bump_cancel_stat("dispatcher_reconnects")
1459+
self._bump_cancel_stat("dispatcher_internal_errors")
14051460
await logger.aerror(
14061461
f"Cancel dispatcher internal error (retrying in {backoff:.1f}s): {exc!r}",
14071462
exc_info=True,
@@ -1418,7 +1473,7 @@ async def _run_cancel_dispatcher(self) -> None:
14181473

14191474
async def _on_cancel_dispatcher_connection_reconnect(self, _connection: Any) -> None:
14201475
"""Record redis-py reconnects that happen inside an active PubSub."""
1421-
self._cancel_stats["dispatcher_reconnects"] += 1
1476+
self._bump_cancel_stat("dispatcher_reconnects")
14221477
with contextlib.suppress(Exception):
14231478
await logger.awarning("Cancel dispatcher pubsub connection reconnected transparently")
14241479

@@ -1458,10 +1513,10 @@ async def _apply_cancel(self, job_id: str, *, source: str, wait_for_cleanup: boo
14581513
"""
14591514
entry = self._queues.get(job_id)
14601515
if entry is None:
1461-
self._cancel_stats["dispatched_foreign"] += 1
1516+
self._bump_cancel_stat("dispatched_foreign")
14621517
await logger.adebug(f"Cancel for {job_id} ignored on this worker (not owner); source={source}")
14631518
return
1464-
self._cancel_stats["dispatched_owned"] += 1
1519+
self._bump_cancel_stat("dispatched_owned")
14651520
await logger.ainfo(f"Cancel applied to {job_id} (source={source})")
14661521
main_queue, _, task, _ = entry
14671522
if task is not None and not task.done():
@@ -1569,9 +1624,9 @@ async def signal_cancel(self, job_id: str) -> int:
15691624
await self._client.set(self._cancel_marker_key(job_id), "1", ex=self._cancel_marker_ttl)
15701625
receivers = int(await self._client.publish(self._cancel_channel(job_id), "1"))
15711626
except Exception:
1572-
self._cancel_stats["publish_errors"] += 1
1627+
self._bump_cancel_stat("publish_errors")
15731628
raise
1574-
self._cancel_stats["published"] += 1
1629+
self._bump_cancel_stat("published")
15751630
await logger.ainfo(f"signal_cancel: job_id={job_id} receivers={receivers}")
15761631
return receivers
15771632

src/backend/base/langflow/services/telemetry/opentelemetry.py

Lines changed: 23 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -140,6 +140,25 @@ def _register_metric(self) -> None:
140140
metric_type=MetricType.COUNTER,
141141
labels={"flow_id": mandatory_label},
142142
)
143+
self._add_metric(
144+
name="langflow_job_queue_cancel_events_total",
145+
description=(
146+
"Job queue cancel-channel and watchdog events. event_type is one of: "
147+
"published, marker_hit, dispatched_owned, dispatched_foreign, publish_errors, "
148+
"dispatcher_reconnects, dispatcher_internal_errors, polling_watchdog_kills, "
149+
"activity_touch_errors, activity_get_errors, activity_parse_errors."
150+
),
151+
unit="",
152+
metric_type=MetricType.COUNTER,
153+
labels={"event_type": mandatory_label},
154+
)
155+
self._add_metric(
156+
name="langflow_job_queue_active_jobs",
157+
description="Active jobs tracked by the job queue on this worker.",
158+
unit="",
159+
metric_type=MetricType.UP_DOWN_COUNTER,
160+
labels={"backend": mandatory_label},
161+
)
143162

144163
def __init__(self, *, prometheus_enabled: bool = True):
145164
# Only initialize once
@@ -154,8 +173,10 @@ def __init__(self, *, prometheus_enabled: bool = True):
154173
# Get existing meter provider if any
155174
existing_provider = metrics.get_meter_provider()
156175

157-
# Check if FastAPI instrumentation is already set up
158-
if hasattr(existing_provider, "get_meter") and existing_provider.get_meter("http.server"):
176+
# Reuse a concrete SDK provider installed by another integration. The
177+
# default API proxy also returns meters, but it has no readers and must
178+
# be replaced so Prometheus can collect Langflow metrics.
179+
if isinstance(existing_provider, MeterProvider):
159180
self._meter_provider = existing_provider
160181
else:
161182
resource = Resource.create({"service.name": "langflow"})

src/backend/tests/unit/components/processing/test_operations_component.py

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,13 +9,16 @@
99

1010
import pandas as pd
1111
import pytest
12+
from lfx.components.processing.data_operations import DataOperationsComponent
13+
from lfx.components.processing.dataframe_operations import DataFrameOperationsComponent
1214
from lfx.components.processing.operations import (
1315
JSON_OPERATIONS,
1416
OPERATIONS_BY_TYPE,
1517
TABLE_OPERATIONS,
1618
TEXT_OPERATIONS,
1719
OperationsComponent,
1820
)
21+
from lfx.components.processing.text_operations import TextOperations
1922
from lfx.schema import Data
2023
from lfx.schema.dataframe import DataFrame
2124
from lfx.schema.dotdict import dotdict
@@ -44,6 +47,14 @@ def file_names_mapping(self):
4447
return []
4548

4649

50+
@pytest.mark.parametrize(
51+
"legacy_component",
52+
[DataOperationsComponent, DataFrameOperationsComponent, TextOperations],
53+
)
54+
def test_legacy_operations_components_point_to_unified_replacement(legacy_component):
55+
assert legacy_component.replacement == ["processing.Operations"]
56+
57+
4758
class TestOperationCatalog:
4859
"""The unified picker exposes every operation across all three data types.
4960

src/backend/tests/unit/test_redis_job_queue_service.py

Lines changed: 116 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2904,3 +2904,119 @@ async def test_initialize_services_fails_fast_when_redis_queue_unreachable(monke
29042904
finally:
29052905
manager.factories.clear()
29062906
manager.services.clear()
2907+
2908+
2909+
# ---------------------------------------------------------------------------
2910+
# OpenTelemetry instrumentation
2911+
# ---------------------------------------------------------------------------
2912+
2913+
2914+
class _OtelRecorder:
2915+
"""Capture OTel emissions so tests can assert what the queue exported.
2916+
2917+
Substituted for the real OT singleton via ``service._otel``. Mirrors the
2918+
surface area used by ``_emit_otel_counter`` and ``_emit_otel_up_down``.
2919+
"""
2920+
2921+
def __init__(self) -> None:
2922+
self.counters: list[tuple[str, dict[str, str], float]] = []
2923+
self.up_downs: list[tuple[str, float, dict[str, str]]] = []
2924+
2925+
def increment_counter(self, name: str, labels: dict[str, str], value: float = 1.0) -> None:
2926+
self.counters.append((name, dict(labels), value))
2927+
2928+
def up_down_counter(self, name: str, value: float, labels: dict[str, str]) -> None:
2929+
self.up_downs.append((name, value, dict(labels)))
2930+
2931+
2932+
def _attach_recorder(service: Any) -> _OtelRecorder:
2933+
"""Bypass the lazy OTel resolver and inject a recorder."""
2934+
recorder = _OtelRecorder()
2935+
service._otel = recorder
2936+
service._otel_resolved = True
2937+
return recorder
2938+
2939+
2940+
@pytest.mark.asyncio
2941+
async def test_bump_cancel_stat_updates_dict_and_otel_counter():
2942+
service, _client = await _make_service(cancel_channel_enabled=False)
2943+
recorder = _attach_recorder(service)
2944+
try:
2945+
service._bump_cancel_stat("published")
2946+
service._bump_cancel_stat("marker_hit", value=2)
2947+
2948+
# Dict mirror still works for /monitor/job_queue.
2949+
assert service._cancel_stats["published"] == 1
2950+
assert service._cancel_stats["marker_hit"] == 2
2951+
2952+
# OTel counter received both bumps with event_type label.
2953+
assert (
2954+
"langflow_job_queue_cancel_events_total",
2955+
{"event_type": "published"},
2956+
1.0,
2957+
) in recorder.counters
2958+
assert (
2959+
"langflow_job_queue_cancel_events_total",
2960+
{"event_type": "marker_hit"},
2961+
2.0,
2962+
) in recorder.counters
2963+
finally:
2964+
await _stop_service(service)
2965+
2966+
2967+
@pytest.mark.asyncio
2968+
async def test_create_and_cleanup_move_active_jobs_up_down_counter():
2969+
service, _client = await _make_service(cancel_channel_enabled=False)
2970+
recorder = _attach_recorder(service)
2971+
try:
2972+
job_id = uuid.uuid4().hex
2973+
service.create_queue(job_id)
2974+
await service.cleanup_job(job_id)
2975+
2976+
deltas = [
2977+
(value, labels) for name, value, labels in recorder.up_downs if name == "langflow_job_queue_active_jobs"
2978+
]
2979+
assert (1, {"backend": "redis"}) in deltas
2980+
assert (-1, {"backend": "redis"}) in deltas
2981+
finally:
2982+
await _stop_service(service)
2983+
2984+
2985+
@pytest.mark.asyncio
2986+
async def test_otel_emit_is_silent_when_telemetry_unavailable():
2987+
"""A broken OT handle must never propagate out of the emit helpers."""
2988+
service, _client = await _make_service(cancel_channel_enabled=False)
2989+
try:
2990+
explosion = "telemetry exploded"
2991+
2992+
class _BrokenOt:
2993+
def increment_counter(self, *_a, **_kw):
2994+
raise RuntimeError(explosion)
2995+
2996+
def up_down_counter(self, *_a, **_kw):
2997+
raise RuntimeError(explosion)
2998+
2999+
service._otel = _BrokenOt()
3000+
service._otel_resolved = True
3001+
3002+
# These must not raise even though OT itself does.
3003+
service._bump_cancel_stat("published")
3004+
service._emit_otel_up_down("langflow_job_queue_active_jobs", 1, {"backend": "redis"})
3005+
# Dict mirror still updated despite OT failure.
3006+
assert service._cancel_stats["published"] == 1
3007+
finally:
3008+
await _stop_service(service)
3009+
3010+
3011+
@pytest.mark.asyncio
3012+
async def test_all_cancel_stat_keys_route_through_helper():
3013+
"""Every key initialized in _cancel_stats must be a valid argument to _bump_cancel_stat."""
3014+
service, _client = await _make_service(cancel_channel_enabled=False)
3015+
recorder = _attach_recorder(service)
3016+
try:
3017+
for key in list(service._cancel_stats):
3018+
service._bump_cancel_stat(key)
3019+
emitted_event_types = {labels["event_type"] for _, labels, _ in recorder.counters}
3020+
assert emitted_event_types == set(service._cancel_stats.keys())
3021+
finally:
3022+
await _stop_service(service)

0 commit comments

Comments
 (0)