Commit Graph

20 Commits

Author SHA1 Message Date
2ddac97c77 feat(job_queue): cross-worker cancel + polling watchdog for multi-worker Redis-backed builds (#13084)
* feat: implement Redis job queue and fakeredis support
- Refactored the job queue service to support Redis-backed management for cross-worker scaling.
- Added environment variables for configuration:
    - `LANGFLOW_JOB_QUEUE_TYPE=redis`
    - `LANGFLOW_REDIS_QUEUE_DB=1`
- Updated job ownership methods to be asynchronous for improved concurrency handling.
- Enhanced Redis cache service with namespacing via key prefixes.
- Introduced `fakeredis` for in-memory Redis simulation in testin>
- Added comprehensive unit tests for Redis job queue components.

* fix: run  experimental warning for RedisCache usage only once

- Introduced a mechanism to emit a one-time warning for the RedisCache experimental feature during server runtime.
- The warning is logged only if no other worker has already emitted it, ensuring clarity for users regarding the experimental status of RedisCache.
- The implementation includes a temporary file check to prevent multiple warnings across different processes.

* docs: document environment variables for worker management
- Added documentation for LANGFLOW_GUNICORN_PRELOAD to explain preloading for better performance.
- Detailed the use of LANGFLOW_JOB_QUEUE_TYPE for specifying backends (e.g., Redis).
- Included LANGFLOW_REDIS_QUEUE_DB to define the database index for job queues.
- Updated the "High-Load Environments" guide with these optimal configurations.

* docs: updated 'High-load and multi-worker environments' section

* feat: enhance RedisJobQueueService with consumer wrapper management
- Introduced a caching mechanism for Redis stream consumers to optimize job data retrieval.
- Added methods to manage consumer wrappers, ensuring they are reused across sequential polls.
- Implemented cleanup logic to cancel and clear consumer wrappers during job cleanup and service stop.
- Expanded unit tests to verify consumer wrapper reuse and cleanup behavior.

* fix: ensure Redis keys are deleted during job cleanup even on cancellation

- Updated the cleanup_job method in RedisJobQueueService to guarantee Redis keys are removed even if the job cleanup is interrupted by a CancelledError.
- Added a new unit test to verify that Redis keys are deleted correctly when cleanup is called during task cancellation.

* fix: manage connection check task in RedisJobQueueService

- Added handling for the connection check task in the stop method to ensure it is properly cancelled and awaited if still running.
- This change improves resource management and prevents potential issues during service shutdown.

* fix: handle unpublished sentinel requeue on cancellation in RedisJobQueueService

- Updated the job processing logic to ensure that if a job is cancelled during the xadd operation, the unpublished sentinel is requeued instead of being dropped.
- Introduced a new unit test to verify this behavior, ensuring robustness in job handling during cancellations.

* fix: improve Redis job queue service with enhanced configuration and cleanup

- Added atexit cleanup to remove stale temporary files for RedisCache.
- Refactored Redis job queue service to use shared constants for stream prefixes, improving maintainability.
- Updated type hints for better clarity and consistency in RedisQueueWrapper and RedisJobQueueService.
- Enhanced error handling with configurable backoff for transient read failures.

* fix: enhance Redis job queue service with maxlen configuration for xadd

- Updated the xadd method in RedisJobQueueService to include maxlen and approximate parameters, improving stream management and preventing excessive memory usage.

* fix: enhance job ownership retrieval in RedisJobQueueService

- Updated the get_job_owner method to refresh the Redis key TTL on successful lookups, ensuring long-running jobs maintain their ownership anchor.
- Improved code clarity by extracting the owner key into a variable and adding detailed docstring explanations for better understanding of the TTL management.

* fix: improve job ownership handling in cleanup_job method

- Enhanced the cleanup_job method in RedisJobQueueService to accurately capture job ownership before deleting Redis keys, preventing potential data corruption in multi-worker scenarios.
- Added comments for clarity on ownership logic and its implications during job cleanup.

* fix: optimize TTL management in RedisJobQueueService

- Introduced periodic TTL refresh logic in the _bridge_to_redis method to enhance Redis stream management, reducing round-trips and improving throughput.
- Added constants for TTL refresh events and seconds to maintain clarity and configurability.
- Updated event handling to ensure TTL is refreshed appropriately based on event count and time elapsed.

* fix: improve event handling in flow response management

- Removed unnecessary error handling for missing event tasks in get_flow_events_response, allowing for smoother operation when no task exists.
- Updated create_flow_response to handle optional event_task parameter, ensuring proper cleanup during disconnections.
- Added unit tests to verify behavior when event tasks are missing, enhancing robustness in streaming scenarios.

* fix: enhance RedisQueueWrapper with startup grace period and stream observation

- Added a startup grace period to prevent premature end-of-stream signals when the producer has not yet issued its first XADD.
- Introduced a flag to track whether the stream has been observed, improving the handling of early polling scenarios.
- Updated logic to ensure proper handling of stream existence checks and logging for better debugging during job processing.

* fix: implement cleanup for old cross-worker job queues in RedisJobQueueService

- Added a new method to clean up done cross-worker consumer wrappers that are not owned by the current worker, ensuring proper resource management.
- Enhanced the existing cleanup logic to prevent memory leaks by explicitly pruning stale entries from the consumer wrappers dictionary.
- Improved logging to provide better visibility into the cleanup process for cross-worker jobs.

* fix: enhance job handling in JobQueueService with guarded task execution

- Updated the start_job method to accept a Coroutine type for task_coro, ensuring type safety.
- Introduced a new _guarded_task method to wrap job coroutines, guaranteeing that unhandled exceptions emit an error event and write a sentinel to the Redis Stream, improving reliability in job processing.
- Enhanced documentation to clarify the behavior of the new task handling mechanism and its implications for cross-worker consumers.

* fix: improve client disconnection handling in create_flow_response

- Added logging for scenarios where a client disconnects without an associated event_task, clarifying that the producer will continue running until the build completes.
- Documented the limitation regarding cross-worker passive disconnects and the need for a Redis side-channel for proper cancellation, enhancing observability in the logs.

* fix: enhance RedisQueueWrapper to manage initial read state and buffer behavior

- Introduced a flag to track the completion of the first XREAD call, ensuring proper buffer management during the initial read phase.
- Updated the empty method to reflect the state of the buffer accurately, preventing premature exits from the drain loop until the first read is complete.
- Improved documentation to clarify the behavior of the new flag and its impact on job processing.

* fix: clarify RedisQueueWrapper behavior in tests for buffer state and no-op operations

- Updated test documentation to explain the behavior of the empty() method in relation to the first XREAD completion, ensuring accurate understanding of buffer state during job processing.
- Adjusted assertions in tests to reflect the intended behavior of the RedisQueueWrapper, specifically regarding the no-op nature of put_nowait and its impact on the internal buffer.

* fix: update dependencies in pyproject.toml and uv.lock

- Added fakeredis dependency with a minimum version of 2.0.0 to both pyproject.toml and uv.lock files.
- Ensured proper formatting and comments in pyproject.toml for clarity on onnxruntime version constraints.

* feat(job_queue): cross-worker cancel via Redis pub/sub + configurable startup grace

- Add LANGFLOW_REDIS_QUEUE_STARTUP_GRACE_S (default 30.0) so deployments
  with slow producer-worker cold-start can bump the wrapper's early-poll
  grace period without editing code.

- Add LANGFLOW_REDIS_QUEUE_CANCEL_CHANNEL_ENABLED (default True) wiring a
  per-job Redis pub/sub channel (langflow:cancel:<job_id>) so
  POST /build/{job_id}/cancel works cross-worker. The producer subscribes
  when the job starts; any worker can publish via
  RedisJobQueueService.signal_cancel and the subscriber cancels the local
  build task on receipt.

- cancel_flow_build now falls back to signal_cancel when event_task is
  None on this worker, replacing the previous no-op for cross-worker
  cancellation.

- Subscriber tasks are tracked in _cancel_subscribers and torn down in
  cleanup_job/stop alongside the bridge tasks.

- Tests: 3 new unit tests covering cross-worker signal propagation, the
  disabled-flag no-op, and the wrapper's per-instance startup_grace_s
  override.

* fix(job_queue): close three cross-worker cancel corner cases

Discovered by a deeper local stress pass over the prototype:

1. **Signal-before-subscribe race**: publish from worker A would be lost if
   worker B's subscriber had not finished SUBSCRIBE yet. signal_cancel now
   also sets a short-lived langflow:cancel-marker:<job_id> key, and the
   subscriber checks it immediately after subscribing — so a cancel that
   raced the SUBSCRIBE still fires.

2. **Restart subscriber leak**: start_job called twice for the same job_id
   silently overwrote the dict entry and orphaned the previous subscriber
   (and its Redis pubsub connection). Now cancels the previous subscriber
   before storing the new one.

3. **Slow end-of-stream after cross-worker cancel**: cancelling the local
   task did not unblock cross-worker consumers — the bridge sat waiting on
   local_queue.get() until periodic cleanup ran ~5 min later. The
   subscriber now puts a sentinel on the local queue after task.cancel()
   so the bridge flushes a clean end-of-stream marker to Redis. Measured
   8ms from signal_cancel to consumer end-of-stream against real Redis.

Tests: 3 new unit tests covering each corner case; 9-scenario local
stress harness against real Redis confirms all green.

* refactor(job_queue): address self-review of cross-worker cancel

Addresses the eight points from the PR review of the previous commits:

1. Configurable marker TTL — new LANGFLOW_REDIS_QUEUE_CANCEL_MARKER_TTL
   setting (default 60s) replaces the hardcoded constant.

2. Connection-pool footprint is O(1) in active jobs — replaced the per-job
   pub/sub subscriber with a single PSUBSCRIBE dispatcher per worker
   listening on langflow:cancel:*.

3. Prompt stream cleanup after cancel — the dispatcher now waits for the
   bridge to drain the sentinel before triggering cleanup_job, so the
   Redis stream + owner keys are deleted within milliseconds instead of
   waiting for the 5-minute periodic cleanup grace.

4. signal_cancel raises on Redis failure — publish errors propagate to the
   caller instead of silently returning 0. The cancel HTTP endpoint
   catches and returns False so the client can retry.

5. Auth note in signal_cancel docstring — explicit note that callers must
   verify authorization; the HTTP cancel endpoint already does via
   _verify_job_ownership before calling through.

6. Structured cancel logging — INFO-level logs on publish, marker hit,
   owned dispatch, foreign dispatch. _cancel_stats counters expose
   {published, marker_hit, dispatched_owned, dispatched_foreign,
   publish_errors} for ops/metrics.

7. redis-py version compatibility — _close_pubsub falls back from
   aclose() to close(), handles both sync and async return values.

8. Fire-and-forget tasks hold strong references — _background_tasks set
   keeps marker checks and post-cancel cleanups from being GC'd before
   they run; each task self-removes via add_done_callback. stop() drains
   the set on shutdown.

Tests: 24 passing (1 skipped due to timing); deep stress harness verified
6 scenarios against real Redis.

* feat(job_queue): propagate client disconnect via signal_cancel for cross-worker builds

Closes the remaining cross-worker passive-disconnect gap. Previously, when a
client closed its streaming connection on a non-owner worker, the producer
worker kept emitting events until the build completed naturally. Now the
disconnect handler in create_flow_response publishes a cross-worker cancel
via signal_cancel so the owning worker stops promptly.

- create_flow_response accepts optional queue_service + job_id kwargs and
  uses them in on_disconnect for the cross-worker case (event_task is None).
  Both are keyword-only and default to None, preserving the single-worker
  contract.
- get_flow_events_response wires both through.
- New unit test test_cross_worker_disconnect_publishes_signal_cancel covers
  the full pubsub propagation path with two services sharing a fake Redis.
- Tested end-to-end against real Redis with a two-worker harness: 10/10
  scenarios pass, including the new disconnect propagation case.

* feat(job_queue): production-harden cross-worker cancel for multi-worker setups

Five complementary improvements that make the Redis-backed cancel/lifecycle
path safe to leave running unattended under real ops pressure.

A. Dispatcher reconnect loop with exponential backoff. _run_cancel_dispatcher
   used to exit silently on any Redis-side error (broker restart, network
   blip, listen() hiccup), leaving the worker permanently blind to
   cross-worker cancels until process restart. The loop now reconnects with
   capped backoff (max 30s), and a dispatcher_reconnects counter exposes the
   event for monitoring.

B. Outer timeout on _post_cancel_cleanup. The background cleanup task could
   hang indefinitely if cleanup_job stalled (Redis pathology during DELETE);
   wrap in asyncio.wait_for and let periodic cleanup retry instead.

C. Public metrics_snapshot() on JobQueueService (memory + Redis backends)
   plus GET /monitor/job_queue endpoint behind get_current_active_user.
   Surfaces backend, active_jobs, bridge_count, consumer_wrapper_count,
   background_task_count, cancel_dispatcher_running, and the full
   cancel_stats counter set.

D. Deterministic rewrite of test_redis_service_signal_cancel_flushes_sentinel_to_consumer.
   The previous version was skipped roughly half the time due to a fakeredis
   timing race; the new version no-ops cleanup_job for the duration of the
   assertion so the bridge XADD ordering can be observed reliably. No more
   skipped tests in the suite (33 pass, 0 skip).

E. Polling-mode stale-client watchdog. Polling has no persistent connection
   for on_disconnect, so abandoned polling builds previously ran to natural
   completion even after the client gave up. New flow:
     * touch_activity(job_id) writes langflow:activity:<job_id> on every poll
       and on streaming-response open; consume_and_yield refreshes it every
       10s while the connection is held.
     * _run_polling_watchdog scans owned jobs every
       LANGFLOW_REDIS_QUEUE_POLLING_WATCHDOG_INTERVAL_S (default 15s) and
       publishes signal_cancel for jobs whose activity is older than
       LANGFLOW_REDIS_QUEUE_POLLING_STALE_THRESHOLD_S (default 90s).
     * Streaming clients are protected by the heartbeat refresh; threshold <=
       0 disables the watchdog entirely.
     * cleanup_job folds the activity key into the existing DEL round-trip
       so successful builds clean up after themselves.

End-to-end coverage: 33 unit tests pass (up from 25, with the previous skip
now passing deterministically); a two-worker real-Redis harness exercises 12
scenarios including dispatcher reconnect, post-cancel timeout, watchdog
reclamation, and metrics_snapshot schema completeness.

* fix(job_queue): address PR review + Codex adversarial review findings

Acts on every must-fix finding from a multi-agent review pass (code-reviewer,
pr-test-analyzer, silent-failure-hunter, comment-analyzer) plus the Codex
adversarial review. No behavior is left to chance under realistic operational
conditions.

Critical fixes:

- Streaming heartbeat now runs as an INDEPENDENT task for the lifetime of the
  response, not coupled to event yield cadence. A quiet build (long graph
  step, slow LLM, no tokens for a while) previously had no heartbeat path, so
  the polling watchdog could reclaim a live streaming client after the
  threshold. The new heartbeat fires every streaming_activity_refresh_s
  regardless of queue activity, and is cancelled cleanly on both
  on_disconnect AND natural stream end.

- Watchdog start-grace window. A new in-memory _job_start_times[job_id]
  timestamp is set in start_job() BEFORE super().start_job() so a job can
  never appear in self._queues without a corresponding start time. When the
  watchdog scans and the Redis activity key is missing, it skips the kill
  until (now - start_ts) >= threshold. Without this, a slow background
  touch_activity (event-loop pressure, Redis blip) could nuke a brand-new
  build on the very first watchdog tick — the prior code's comment promised
  the guard but didn't implement it.

- touch_activity errors are now observable. Replaces a blanket
  contextlib.suppress(Exception) with an explicit try/except that bumps
  activity_touch_errors and emits a debug log. CancelledError is no longer
  swallowed. Operators get a signal when the heartbeat is silently failing.

- Watchdog UnboundLocalError closed. Initializes last = 0.0 before the
  raw-is-None branch so a malformed activity value (parse failure) can't
  crash the entire watchdog task. Adds activity_parse_errors and
  activity_get_errors counters.

- Dispatcher except is split: ConnectionError/TimeoutError/OSError stay at
  awarning (transient Redis), every other Exception goes to aerror with
  exc_info=True AND increments dispatcher_internal_errors. Code bugs in
  _handle_cancel no longer look identical to Redis disconnects.

- _post_cancel_cleanup narrows the bridge-wait suppress to TimeoutError
  only, with a dedicated awarning naming the consequence ("cross-worker
  consumers may see late end-of-stream"). Real bridge failures bubble up
  via a separate Exception arm with awarning instead of being silent.

- Polling watchdog uses local _handle_cancel for owned jobs instead of
  round-tripping through pubsub. Faster (no Redis RTT), dispatcher-
  independent (works during a reconnect backoff), and keeps the
  cancel_stats["published"] counter honest as a count of *external* cancels.

- /monitor/job_queue gated to get_current_active_superuser instead of
  get_current_active_user — the snapshot exposes process-wide tenant
  workload + cancel activity, which matters in multi-tenant deployments.

Docstring corrections:

- Removed the "cross-worker cancel is a known limitation" remnant from
  RedisJobQueueService.get_queue_data — the limitation was closed in the
  prior commit, and the new docstring tells callers exactly how cross-worker
  cancel travels (signal_cancel → dispatcher).

- touch_activity docstring rewritten. The previous text described a
  "TTL-based reasoning" fallback that didn't exist; the new text correctly
  explains the 4x-threshold TTL and how it interacts with the watchdog's
  start-time grace window.

New tests (5 added, total 38 pass):

- test_polling_watchdog_grants_start_grace_window: brand-new job with
  deleted activity key survives until the threshold passes, then dies.
- test_polling_watchdog_skips_malformed_activity_value: malformed value
  bumps activity_parse_errors and does NOT kill the job.
- test_streaming_heartbeat_runs_independent_of_event_yield: quiet streaming
  build is not reclaimed; heartbeat task is cancelled on disconnect.
- test_concurrent_cancels_from_multiple_workers_are_idempotent: two workers
  publish cancel simultaneously, no crashes, single observable cancel.
- test_dispatcher_internal_error_logged_at_error_level: bug in
  _handle_cancel bumps dispatcher_internal_errors and dispatcher_reconnects;
  the dispatcher task is still running afterwards.

Plus: test_metrics_snapshot_exposes_cancel_stats_and_counters now pins the
full cancel_stats key set (using set equality), so adding an increment
without registering the key — or vice versa — fails the test instead of
producing a silent KeyError in production.

End-to-end coverage: 38 unit tests pass; the two-worker real-Redis harness
passes all 12 scenarios on commit HEAD.

* fix: harden redis job queue cancellation

* feat(job_queue): seamless redis setup + actionable errors for unsupported delivery

- RedisQueueWrapper: restore _BUFFER_MAXSIZE bounded buffer and the
  _on_fill_done done-callback safety net that release-1.10.0 added in #12588,
  so a slow consumer cannot grow the buffer without bound and a crashing or
  cancelled _fill_task cannot leave consumers stuck on await get().
- build.get_flow_events_response: explicit exhaustiveness guard. Unknown
  EventDeliveryType values now return HTTP 400 with the supported set and a
  remediation hint instead of silently falling through to the polling path.
- lfx settings.set_event_delivery: when workers > 1 without a redis queue,
  upgrade the warning to name the requested mode, the forced fallback, and
  the LANGFLOW_JOB_QUEUE_TYPE env var that would preserve the original mode.
- Tests: port the three RedisQueueWrapper safety tests from release-1.10.0
  and add coverage for the new event_delivery guard.

* fix(job_queue): polling watchdog only watches jobs with a registered owner

Surfaced by locust load testing on a Redis-queue, 2-worker setup:

  2016 requests / 2016 polling_watchdog_kills (1:1 ratio)

TaskService.launch_task uses JobQueueService.start_job for server-internal
tasks (telemetry, background work) without calling register_job_owner.
These tasks never refresh the activity heartbeat because no polling
client is involved, so the watchdog's start-time fallback was killing
every single one once it crossed polling_stale_threshold_s.

Under fast-completing flows the user-visible result was unaffected (the
build finished before the watchdog struck) — but a long-running internal
task (slow LLM call, retrieval, etc.) would be cancelled mid-flight even
though no one was waiting on it.

Scope the watchdog to jobs in self._job_owners: user-facing flow builds
register an owner; TaskService tasks don't.  After this change the same
locust run reports polling_watchdog_kills=0 and dispatched_owned=0.

Existing watchdog tests register an owner explicitly to mirror the real
flow path.  Adds a regression test (TaskService-style start_job with no
owner must not be reclaimed).

* test(job_queue): add coverage for fixes surfaced by load testing

Three regression tests for the recent fixes:

- TaskService integration: launch a task through
  TaskService.fire_and_forget_task and assert the polling watchdog
  leaves it alone (no owner registered, no kill).  Catches a future
  regression where someone adds register_job_owner inside TaskService.

- Settings multi-worker warning: verify Settings.set_event_delivery
  warns with the requested mode AND names LANGFLOW_JOB_QUEUE_TYPE in
  the fallback log when workers > 1 and the job queue is in-memory.
  An operator seeing missing events needs that env var name to fix it.

- Bounded buffer backpressure: floods the fill task past the maxsize
  ceiling and verifies qsize() never exceeds _BUFFER_MAXSIZE.  The
  previous test only asserted the constant, not the behavior.

Skipped two suggested tests:
- TestClient version of the unknown-event_delivery guard: FastAPI's
  enum validation returns 422 before our handler runs, so HTTP cannot
  reach the defensive 400.  The existing unit test is the right test.
- Exact watchdog kill-count assertion (== 1): the watchdog can tick
  again before cleanup_job removes the job from _queues, making any
  exact-count check flaky on slow CI.  The existing >= 1 is correct.

* fix(job_queue): use asyncio.create_task for trace cleanup on cancel

background_tasks.add_task() is silently dropped after FastAPI drains
the POST /build response queue — by the time any cancel fires (local,
cross-worker pub/sub, or polling watchdog), the task list from the
original HTTP request is already consumed.

Replace the background_tasks call in _run_vertex_build's CancelledError
handler with asyncio.create_task(graph.end_all_traces_in_context()())
so the trace cleanup runs as an independent task regardless of the
background_tasks lifecycle.

Adds a regression test that exercises generate_flow_events end-to-end:
a blocking vertex is cancelled mid-flight and the test asserts that
end_all_traces is eventually called via the spawned task.

* [autofix.ci] apply automated fixes

* fix: addresses a few pr review findings (#13179)

* fix(job_queue): address PR review findings from cross-worker cancel implementation

- Restore _last_id cursor update to after await buffer.put() in _fill_from_redis;
  advancing before the await loses the event if the task is cancelled while put()
  is blocked on a full buffer
- Restore XREAD persistent-error grace period: track elapsed error time and deliver
  an end-of-stream sentinel after _STARTUP_GRACE_S of continuous failures instead
  of looping forever, which left consumers stuck on await get()
- Fix cancel_flow_build returning True when no task and no signal_cancel (in-memory
  backend); return False to correctly signal that cancellation did not take effect
- Decouple polling watchdog from cancel_channel_enabled; the watchdog uses only
  local state and _handle_cancel with no pub/sub dependency, so disabling the
  cancel channel must not silently disable stale-client reclamation
- Restore _BUFFER_MAXSIZE to 1000 (reverts unexplained 10x increase to 10_000)

* test(job_queue): add regression coverage for PR review fixes

- test_fill_from_redis_last_id_not_advanced_before_put_completes: verifies
  the XREAD cursor (_last_id) stays at the prior message's position when the
  fill task is cancelled while put() is blocked on a full buffer
- test_fill_from_redis_persistent_xread_error_delivers_sentinel: verifies a
  persistent XREAD failure eventually delivers the end-of-stream sentinel
  after _STARTUP_GRACE_S of continuous errors (not an infinite retry loop)
- test_cancel_flow_build_returns_false_for_in_memory_backend_without_task:
  verifies cancel_flow_build returns False when no local task exists and the
  queue service has no cross-worker cancel path
- test_polling_watchdog_runs_when_cancel_channel_disabled: verifies the
  polling watchdog reclaims stale jobs even when cancel_channel_enabled=False

* format and comment updates

* comments

* [autofix.ci] apply automated fixes

---------

Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>

* fix(job_queue): address CodeRabbit review findings

- cancel_flow_build: gate signal_cancel on a new cross_worker_cancel_enabled
  property so Redis with cancel_channel_enabled=False is treated the same as
  the in-memory backend (signal_cancel short-circuits to 0 without setting
  the marker, so the previous unconditional return True was misleading).
- _run_vertex_build: hold a strong reference to the trace cleanup task on
  CancelledError (Ruff RUF006). Tasks self-discard via add_done_callback.
- RedisQueueWrapper.finish_with_sentinel: stop awaiting on a bounded buffer
  during teardown. Evict + put_nowait mirrors _on_fill_done so shutdown can
  no longer hang on a slow or abandoned consumer.
- RedisJobQueueService.cleanup_job: wrap the Redis DEL in try/except so
  Redis failures during teardown surface as warnings instead of escaping
  stop() and explicit cancel paths.
- Settings (lfx): add Field bounds on the new Redis timing knobs so a
  negative startup_grace_s or non-positive cancel_marker_ttl /
  watchdog_interval_s fails fast at config load instead of silently
  reopening races.
- Streaming heartbeat interval: lift to module-level STREAMING_ACTIVITY_REFRESH_S
  so the heartbeat test can monkeypatch it.

Tests:
- _make_service tracks shared_client ownership so the first _stop_service
  doesn't aclose the FakeRedis under a sibling worker in cross-worker tests.
- test_polling_watchdog_skips_fresh_activity now registers a job owner;
  without it the watchdog skips the job entirely and the test would pass
  even with a broken touch_activity.
- test_streaming_heartbeat_runs_independent_of_event_yield patches the
  interval down and waits for the spawned task to refresh the activity
  timestamp, instead of calling touch_activity manually.
- Fix Ruff ARG001, EM101, TRY003, PT018 in test stubs.

* test(job_queue): cover behavior changes from CodeRabbit review fixes

- test_cancel_flow_build_returns_false_for_redis_with_disabled_cancel_channel:
  Redis with cancel_channel_enabled=False has signal_cancel but no marker
  delivery; cancel_flow_build must report failure, not false success.
- test_finish_with_sentinel_does_not_hang_on_full_buffer: regression cover for
  the await-on-bounded-buffer hang. Fills the buffer to capacity and asserts
  finish_with_sentinel returns within a tight budget.
- test_redis_service_cleanup_swallows_redis_delete_error: proxies a FakeRedis
  whose delete() raises; cleanup_job must absorb it and complete local
  teardown.
- test_redis_queue_bounds (lfx): pydantic ValidationError on negative
  startup_grace_s / non-positive cancel_marker_ttl /
  polling_watchdog_interval_s; zero remains valid for startup_grace_s and
  polling_stale_threshold_s (documented disable switch).

* log updates

* [autofix.ci] apply automated fixes

* test(cors): accept both '*' and ['*'] for Python 3.14 compatibility

pydantic-settings parses LANGFLOW_CORS_ORIGINS='*' as the raw string '*'
on Python 3.13 but as ['*'] (a single-element list) on Python 3.14+. The
`str | list[str]` union on Settings.cors_origins resolves differently
between versions, but both shapes carry the same 'all origins' semantic.

Updates two assertions in test_security_cors.py that previously required
the exact string form and were failing on Group 4 / Python 3.14 in CI:
- test_default_cors_settings_current_behavior
- test_wildcard_with_credentials_allowed_current_behavior

The PR itself (cross-worker cancel + polling watchdog) is unrelated to
CORS handling -- the failure is a Python-3.14-introduced incompatibility
in this pre-existing test that happens to surface against this branch's
CI matrix.

* [autofix.ci] apply automated fixes

---------

Co-authored-by: Arek Mateusiak <arek@a4bits.com>
Co-authored-by: Jordan Frazier <122494242+jordanrfrazier@users.noreply.github.com>
Co-authored-by: Jordan Frazier <jordan.frazier@datastax.com>
Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
2026-05-25 19:15:05 +00:00
09bd68d349 docs: text operations component (#13152)
docs-add-text-operations-component
2026-05-19 18:25:50 +00:00
c78e10e85a docs: add directory depth limitation warning to custom components (#13104)
* docs: add directory depth limitation warning to custom components

Add a warning admonition to the custom components documentation
explaining the MAX_DEPTH=2 limitation for component discovery.

The warning clarifies:
- Components must be at most 2 levels deep (category/component.py)
- Each first-level directory becomes an independent category in the UI
- How to modify MAX_DEPTH for advanced use cases

This addresses confusion around subdirectory support mentioned in issue #13102.

* docs: simplify directory depth warning per review feedback

Co-Authored-By: Antonio <antonio@oriontech.me>

---------

Co-authored-by: Antonio <antonio@oriontech.me>
2026-05-18 18:26:49 +00:00
5e9f0c37bc Merge branch 'release-1.9.3' into release-1.10.0
# Conflicts:
#	.secrets.baseline
#	docs/package.json
#	pyproject.toml
#	src/backend/base/langflow/initial_setup/starter_projects/Vector Store RAG.json
#	src/backend/base/pyproject.toml
#	src/backend/tests/unit/components/files_and_knowledge/test_retrieval.py
#	src/frontend/package-lock.json
#	src/frontend/package.json
#	src/frontend/src/components/core/parameterRenderComponent/components/modelInputComponent/__tests__/ModelInputComponent.test.tsx
#	src/frontend/src/components/core/parameterRenderComponent/components/modelInputComponent/index.tsx
#	src/frontend/src/modals/modelProviderModal/hooks/useProviderConfiguration.ts
#	src/lfx/pyproject.toml
#	src/lfx/src/lfx/_assets/component_index.json
#	src/sdk/pyproject.toml
#	uv.lock
2026-05-14 21:20:42 -07:00
ea29bcc6b9 feat: Redis-backed job queue for multi-worker deployments (#12588)
* feat: implement Redis job queue and fakeredis support
- Refactored the job queue service to support Redis-backed management for cross-worker scaling.
- Added environment variables for configuration:
    - `LANGFLOW_JOB_QUEUE_TYPE=redis`
    - `LANGFLOW_REDIS_QUEUE_DB=1`
- Updated job ownership methods to be asynchronous for improved concurrency handling.
- Enhanced Redis cache service with namespacing via key prefixes.
- Introduced `fakeredis` for in-memory Redis simulation in testin>
- Added comprehensive unit tests for Redis job queue components.

* fix: run  experimental warning for RedisCache usage only once

- Introduced a mechanism to emit a one-time warning for the RedisCache experimental feature during server runtime.
- The warning is logged only if no other worker has already emitted it, ensuring clarity for users regarding the experimental status of RedisCache.
- The implementation includes a temporary file check to prevent multiple warnings across different processes.

* docs: document environment variables for worker management
- Added documentation for LANGFLOW_GUNICORN_PRELOAD to explain preloading for better performance.
- Detailed the use of LANGFLOW_JOB_QUEUE_TYPE for specifying backends (e.g., Redis).
- Included LANGFLOW_REDIS_QUEUE_DB to define the database index for job queues.
- Updated the "High-Load Environments" guide with these optimal configurations.

* docs: updated 'High-load and multi-worker environments' section

* feat: enhance RedisJobQueueService with consumer wrapper management
- Introduced a caching mechanism for Redis stream consumers to optimize job data retrieval.
- Added methods to manage consumer wrappers, ensuring they are reused across sequential polls.
- Implemented cleanup logic to cancel and clear consumer wrappers during job cleanup and service stop.
- Expanded unit tests to verify consumer wrapper reuse and cleanup behavior.

* fix: ensure Redis keys are deleted during job cleanup even on cancellation

- Updated the cleanup_job method in RedisJobQueueService to guarantee Redis keys are removed even if the job cleanup is interrupted by a CancelledError.
- Added a new unit test to verify that Redis keys are deleted correctly when cleanup is called during task cancellation.

* fix: manage connection check task in RedisJobQueueService

- Added handling for the connection check task in the stop method to ensure it is properly cancelled and awaited if still running.
- This change improves resource management and prevents potential issues during service shutdown.

* fix: handle unpublished sentinel requeue on cancellation in RedisJobQueueService

- Updated the job processing logic to ensure that if a job is cancelled during the xadd operation, the unpublished sentinel is requeued instead of being dropped.
- Introduced a new unit test to verify this behavior, ensuring robustness in job handling during cancellations.

* fix: improve Redis job queue service with enhanced configuration and cleanup

- Added atexit cleanup to remove stale temporary files for RedisCache.
- Refactored Redis job queue service to use shared constants for stream prefixes, improving maintainability.
- Updated type hints for better clarity and consistency in RedisQueueWrapper and RedisJobQueueService.
- Enhanced error handling with configurable backoff for transient read failures.

* fix: enhance Redis job queue service with maxlen configuration for xadd

- Updated the xadd method in RedisJobQueueService to include maxlen and approximate parameters, improving stream management and preventing excessive memory usage.

* fix: enhance job ownership retrieval in RedisJobQueueService

- Updated the get_job_owner method to refresh the Redis key TTL on successful lookups, ensuring long-running jobs maintain their ownership anchor.
- Improved code clarity by extracting the owner key into a variable and adding detailed docstring explanations for better understanding of the TTL management.

* fix: improve job ownership handling in cleanup_job method

- Enhanced the cleanup_job method in RedisJobQueueService to accurately capture job ownership before deleting Redis keys, preventing potential data corruption in multi-worker scenarios.
- Added comments for clarity on ownership logic and its implications during job cleanup.

* fix: optimize TTL management in RedisJobQueueService

- Introduced periodic TTL refresh logic in the _bridge_to_redis method to enhance Redis stream management, reducing round-trips and improving throughput.
- Added constants for TTL refresh events and seconds to maintain clarity and configurability.
- Updated event handling to ensure TTL is refreshed appropriately based on event count and time elapsed.

* fix: improve event handling in flow response management

- Removed unnecessary error handling for missing event tasks in get_flow_events_response, allowing for smoother operation when no task exists.
- Updated create_flow_response to handle optional event_task parameter, ensuring proper cleanup during disconnections.
- Added unit tests to verify behavior when event tasks are missing, enhancing robustness in streaming scenarios.

* fix: enhance RedisQueueWrapper with startup grace period and stream observation

- Added a startup grace period to prevent premature end-of-stream signals when the producer has not yet issued its first XADD.
- Introduced a flag to track whether the stream has been observed, improving the handling of early polling scenarios.
- Updated logic to ensure proper handling of stream existence checks and logging for better debugging during job processing.

* fix: implement cleanup for old cross-worker job queues in RedisJobQueueService

- Added a new method to clean up done cross-worker consumer wrappers that are not owned by the current worker, ensuring proper resource management.
- Enhanced the existing cleanup logic to prevent memory leaks by explicitly pruning stale entries from the consumer wrappers dictionary.
- Improved logging to provide better visibility into the cleanup process for cross-worker jobs.

* fix: enhance job handling in JobQueueService with guarded task execution

- Updated the start_job method to accept a Coroutine type for task_coro, ensuring type safety.
- Introduced a new _guarded_task method to wrap job coroutines, guaranteeing that unhandled exceptions emit an error event and write a sentinel to the Redis Stream, improving reliability in job processing.
- Enhanced documentation to clarify the behavior of the new task handling mechanism and its implications for cross-worker consumers.

* fix: improve client disconnection handling in create_flow_response

- Added logging for scenarios where a client disconnects without an associated event_task, clarifying that the producer will continue running until the build completes.
- Documented the limitation regarding cross-worker passive disconnects and the need for a Redis side-channel for proper cancellation, enhancing observability in the logs.

* fix: enhance RedisQueueWrapper to manage initial read state and buffer behavior

- Introduced a flag to track the completion of the first XREAD call, ensuring proper buffer management during the initial read phase.
- Updated the empty method to reflect the state of the buffer accurately, preventing premature exits from the drain loop until the first read is complete.
- Improved documentation to clarify the behavior of the new flag and its impact on job processing.

* fix: clarify RedisQueueWrapper behavior in tests for buffer state and no-op operations

- Updated test documentation to explain the behavior of the empty() method in relation to the first XREAD completion, ensuring accurate understanding of buffer state during job processing.
- Adjusted assertions in tests to reflect the intended behavior of the RedisQueueWrapper, specifically regarding the no-op nature of put_nowait and its impact on the internal buffer.

* fix: update dependencies in pyproject.toml and uv.lock

- Added fakeredis dependency with a minimum version of 2.0.0 to both pyproject.toml and uv.lock files.
- Ensured proper formatting and comments in pyproject.toml for clarity on onnxruntime version constraints.

* fix: enhance RedisQueueWrapper and job queue handling for improved cancellation and buffer management

- Introduced a structural protocol `_CancellableQueue` to ensure queues can handle cancellation properly during client disconnects.
- Updated `RedisQueueWrapper` to implement this protocol, allowing for graceful cancellation of background tasks.
- Added a maximum size limit to the internal buffer to prevent unbounded memory usage and ensure backpressure on slow consumers.
- Implemented a done callback to handle unexpected fill task crashes, ensuring consumers are not left hanging indefinitely.
- Enhanced unit tests to verify compliance with the new protocol and the behavior of the buffer under various conditions.

* fix: RedisQueueWrapper sentinel delivery, error time-bound, and cancel response

- Deliver end-of-stream sentinel on fill-task cancellation (_on_fill_done now
  handles both cancelled and exception paths so consumers are never left hanging)
- Add _error_start time-bound to xread and exists() error loops: after
  _STARTUP_GRACE_S seconds of continuous Redis errors the sentinel is delivered
  instead of retrying forever
- Advance _last_id cursor only after buffer.put() succeeds so cancellation mid-put
  does not silently skip that message in the Redis cursor
- Return False from cancel_flow_build when event_task is None (cross-worker path)
  so the HTTP response correctly reports success=False instead of false success

* [autofix.ci] apply automated fixes

---------

Co-authored-by: ogabrielluiz <gabriel@langflow.org>
Co-authored-by: Jordan Frazier <jordan.frazier@datastax.com>
Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
2026-05-13 21:09:51 +00:00
00bc383f86 docs: SSRF protection enabled by default (#13106)
docs-ssrf-protection-enabled-by-default
2026-05-13 11:30:37 -04:00
7c347a56ea docs: release note refactor (#12987)
* docs-1.10-release-notes-and-link-out-to-previous-versions

* docs-remove-outdated-windows-upgrade-note
2026-05-06 15:44:07 +00:00
6d2e255f5d Merge remote-tracking branch 'origin/release-1.9.2' into release-1.10.0
# Conflicts:
#	.secrets.baseline
#	pyproject.toml
#	src/backend/base/langflow/initial_setup/starter_projects/Knowledge Retrieval.json
#	src/backend/base/langflow/initial_setup/starter_projects/Vector Store RAG.json
#	src/backend/base/pyproject.toml
#	src/frontend/package-lock.json
#	src/frontend/package.json
#	src/lfx/pyproject.toml
#	src/lfx/src/lfx/_assets/component_index.json
#	src/sdk/pyproject.toml
#	uv.lock
2026-05-01 14:22:58 -07:00
bc16a400d6 docs: add telemetry for deployments endpoints (#12922)
* docs-add-telemetry-for-deployments-endpoints

* peer-review
2026-04-30 13:47:32 +00:00
d965a3ae6c docs: wxo deployment is behind feature flag (#12937)
* docs-add-release-note-and-tip-for-wxo-feature-flag

* remove-unnecessary-information-from-1.10

* Apply suggestions from code review

Co-authored-by: April I. Murphy <36110273+aimurphy@users.noreply.github.com>

---------

Co-authored-by: April I. Murphy <36110273+aimurphy@users.noreply.github.com>
2026-04-29 20:52:40 +00:00
9dcdc704ec docs: change upgrade reinstall order (#12911)
* docs-change-upgrade-reinstall-order

* docs-make-langflow-oss-version-section
2026-04-28 20:23:02 +00:00
1e9a96f5b9 docs: include host field for external db in Helm chart (#12899)
docs-include-host-value-for-external-db
2026-04-27 17:16:56 +00:00
849158ed51 docs: tip for running local ollama models (#12840)
add-tip-about-running-large-models-locally

(cherry picked from commit bc2bd31c05)
2026-04-23 17:49:52 -07:00
8524d49db4 docs: concatenate and merge options for dataframe operations component (#12835)
concat-and-merge-table-component

(cherry picked from commit 919054fc17)
2026-04-23 17:49:52 -07:00
0ab5f040f8 docs: clarify trace storage when multiple providers are enabled (#12816)
clarify-trace-storage-with-multiple-providers

(cherry picked from commit f1863c512d)
2026-04-23 17:49:52 -07:00
f78830730d docs: wxo feature (#12677)
* deployment-sidebars-and-release-note

* fix-sidebars-name

* manage-deployments

* add-requests

* clarify-endpoint-difference

* improve-anchors

* Apply suggestions from code review

Co-authored-by: April I. Murphy <36110273+aimurphy@users.noreply.github.com>

* Apply suggestions from code review

Co-authored-by: Mendon Kissling <59585235+mendonk@users.noreply.github.com>

* clarify-menu

* apply-changes-to-next-and-latest

---------

Co-authored-by: April I. Murphy <36110273+aimurphy@users.noreply.github.com>
(cherry picked from commit 4898aa8621)
2026-04-23 17:49:52 -07:00
4beab56b6b docs: add multi-model opensearch component (#12799)
add-opensearch-multi-model

(cherry picked from commit 15edfe1fec)
2026-04-23 17:49:52 -07:00
aced13ef79 docs: clarify autosaving and versioned flow saving (#12794)
* add-clarification-to-1.9-and-next

* Apply suggestions from code review

Co-authored-by: April I. Murphy <36110273+aimurphy@users.noreply.github.com>

---------

Co-authored-by: April I. Murphy <36110273+aimurphy@users.noreply.github.com>
(cherry picked from commit 7fd8469917)
2026-04-23 17:49:52 -07:00
71fbf5ff00 docs: cut version 1.9.0 (#12681)
* version-1.9.0-docs

* [autofix.ci] apply automated fixes

* [autofix.ci] apply automated fixes (attempt 2/3)

---------

Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
2026-04-14 02:07:01 +00:00
b36444f5d9 docs: add versioning (#12218)
* fix: nightly now properly gets 1.9.0 branch (#12215)

before it was attempting to pull release-notes as letters are alphanumerically after numbers when we sort -V then grab tail
now we only look at branch names that follow the pattern '^release-[0-9]+\.[0-9]+\.[0-9]+$'

* docs: add search icon (#12216)

add-back-svg

* initial-content

* cut-1.8-release-and-include-next-version

* stage-1.8.0-and-next

---------

Co-authored-by: Adam-Aghili <149833988+Adam-Aghili@users.noreply.github.com>
2026-03-18 20:03:49 +00:00