mirror of
https://github.com/langflow-ai/langflow.git
synced 2026-07-24 07:02:37 +08:00
* perf(telemetry): batched off-pool writer for transactions + vertex_builds
Adds TelemetryWriterService that buffers transaction and vertex_build rows
in memory and drains them in batched INSERTs via a dedicated AsyncEngine
(pool_size=1 for SQLite, 2 for Postgres, max_overflow=0). Retention is
amortized in a 60s sweeper instead of running on every insert.
Producers (log_transaction, log_vertex_build) enqueue instead of opening
a DB session, so telemetry traffic no longer competes with the
request-handling pool. Falls back to the legacy direct-write path via
LANGFLOW_TELEMETRY_WRITER_ENABLED=false.
Durability: in-flight rows spill to a diskcache.Deque per PID on shutdown
and are restored on next startup; orphan PID directories left by crashed
workers are adopted.
* perf(telemetry): tighten error handling and add coverage
Review feedback from pr-review-toolkit + silent-failure-hunter:
- Retention sweep snapshots the dirty-flow sets before commit and only
clears them after it lands; a crashed sweep no longer drops the flows on
the floor, so per-flow caps cannot drift unboundedly.
- Producer fall-through is no longer silent. transaction_service and
log_vertex_build each log a one-shot WARNING when telemetry_writer_enabled
is True but the writer is not running.
- Lifespan startup failure now logs ERROR instead of WARNING — if the user
opted into the writer and it didn't come up, that's an error.
- Shutdown drain timeout no longer suppressed silently: logs a WARNING with
pending row count + a hint to raise telemetry_writer_shutdown_drain_s.
- Writer's flush loop catches asyncio.CancelledError separately and
re-prepends the in-flight batch to the buffer so teardown's disk spill
catches it.
- After 6 consecutive batch failures the writer emits a loud ERROR with
buffer depths so operators see sustained data-loss risk.
Added tests:
- test_sanitization_survives_writer_round_trip
- test_retention_failure_preserves_dirty_flows
- test_in_flight_batch_returned_on_cancel
* perf(telemetry): address copilot review
- _restore_from_disk + _adopt_orphan_outboxes now route through _enqueue so a
large disk-spilled or orphan outbox can't bypass telemetry_writer_max_queue
and OOM the process. Oldest rows are dropped and counted via the existing
dropped_transactions / dropped_vertex_builds counters.
- chmod 0o700 the outbox root + per-PID directory so sanitized-but-still-
sensitive payloads aren't exposed cross-user on multi-tenant hosts.
Suppressed on platforms where chmod is a no-op (Windows).
- Added test_adopt_orphan_outboxes_honors_max_queue.
The remaining two copilot notes (private API access to DatabaseService and
diskcache.Cache.close) were also flagged by the in-tree review; tracking
separately. The "diskcache not declared" note is a false positive — the
dependency is at src/backend/base/pyproject.toml:76.
* perf(telemetry): address coderabbit review
- Sweeper hands off dirty sets via capture-and-clear so concurrent
flushes during a retention pass aren't wiped by the post-commit
subtract; failure path restores the snapshot.
- Stress README uses a concrete Postgres DSN matching the docker
example instead of an unset env var.
- Test PID-probe loops bounded via a shared helper with pytest.fail
fallback.
* perf(telemetry): swap diskcache outbox for stdlib sqlite (CVE-2025-69872)
Replaces the diskcache.Deque-backed spill outbox with a stdlib sqlite3
outbox (WAL mode, JSON payloads in TEXT). diskcache 5.6.3 has
CVE-2025-69872 (pickle deserialization RCE for an attacker with write
access to the cache dir), no fixed version released, and was never
declared as a dependency — import failed on Python 3.10.
The replacement is encapsulated in a small _Outbox helper that owns the
connection, schema, JSON codec, and lifecycle. The codec uses tagged
wrappers for datetime and UUID so SQLAlchemy's typed columns accept
restored payloads on the way back out.
Hardening informed by a survey of OTel collector, Fluent Bit, Vector,
Prometheus remote_write, Datadog agent, Logstash, and Filebeat:
- Per-PID outbox dirs now stamp an owner.json (hostname + Linux boot_id
or time()-monotonic() proxy). Adoption only proceeds when host+boot
match; cross-host or pre-owner-file dirs are logged and skipped so a
recycled PID after a container restart cannot pull in a stranger's
spill data.
- PRAGMA synchronous=FULL on the spill connection so the shutdown
commit hits the platter (NORMAL only fsyncs on WAL checkpoint, which
may never run if the process exits immediately after commit).
- Spill honors telemetry_writer_max_queue with drop-oldest, matching
the producer-side overflow policy so a backlogged buffer at shutdown
can neither stall teardown nor fill disk.
- append_all encodes payloads up front and only clears the deque after
the transaction commits — no more partial-drain on mid-flight SQLite
failure. drain collapses to a single DELETE FROM outbox after the
SELECT.
- Exception handlers at the spill/restore/adopt boundaries narrowed to
(sqlite3.Error, OSError) so genuinely unexpected exceptions
propagate rather than being silently logged.
Tests cover: cross-host orphan skipped, missing-owner orphan skipped,
spill cap drops oldest, and a realistic UUID+datetime payload
surviving spill → restore → SQLAlchemy INSERT.
* [autofix.ci] apply automated fixes
* fix(telemetry): name wait_for inner tasks so pyleak can filter
Integration tests using @pyleak_marker were flagging an inner asyncio
task created by the telemetry writer's `wait_for(Event.wait(), ...)`
loops. The wrapper task is auto-named (Task-N) and lands mid-await
across the test boundary, so pyleak counts it as leaked. The writer
itself shuts down cleanly via teardown; the "leak" is a pattern
mismatch with pyleak's per-test snapshot model, not a real lifecycle
bug.
Wrap Event.wait() in an explicitly named task (`telemetry-writer-tick`,
`telemetry-writer-backoff`, `telemetry-sweeper-tick`) via a small
_wait_or_shutdown helper, and extend pyleak_marker with a task_name_filter
that excludes `telemetry-*` from leak detection. CI was green on earlier
commits in this branch only because the diskcache import error
prevented the writer from starting at all — fixing the import surfaced
this latent pyleak collision.
* [autofix.ci] apply automated fixes
* [autofix.ci] apply automated fixes (attempt 2/3)
* test(telemetry): disable writer in tests so reads see writes synchronously
* perf(telemetry): age out cross-host orphan outboxes on shared volumes
A pod on a shared PV (NFS/RWX) cannot adopt rows from a dead pod on a
different host because the dead pod's hostname doesn't match. Without a
janitor those directories accumulate forever. Add an owner-file mtime
heartbeat from the sweeper plus a prune pass that deletes cross-host
outboxes whose owner file hasn't been touched within
telemetry_writer_orphan_max_age_s (default 1h). Same-host orphans still
flow through the existing adoption path.
* perf(telemetry): add byte-aware flush + drop strategy
Today the writer bounds memory by row count only, so a worker logging fat
vertex_build artifacts can hold tens of MB per row and silently dwarf the
configured max_queue. Add a size_strategy switch ('count' | 'bytes' |
'either', default 'count' for parity) with two byte thresholds:
- batch_size_bytes (256KB) caps per-flush INSERT size when bytes apply
- max_queue_bytes (200MB) drops oldest by bytes when bytes apply
Parallel sizes deques mirror the payload buffers so accounting stays
consistent across drain, spill, restore, and the cancel/retry rollback.
Single rows above the byte budget are still emitted (no row is refused);
operators get dropped_*_bytes counters alongside dropped_*.
* test(telemetry): tighten byte-strategy assertions and cover end-to-end paths
Address gaps in the byte-aware strategy tests:
- pin exact dropped counts and remaining buffer state for the drop-by-bytes
case (was 'greater than zero')
- compute exact drain count from encoded size (was a 1-4 range)
- round-trip bytes through spill+restore to verify size deque is rebuilt
- round-trip bytes through orphan adoption for the same reason
- exercise the sweeper loop end-to-end to confirm heartbeat + prune are
wired in (previously only unit-tested in isolation)
Clarify in the settings docs that the byte caps measure encoded JSON size,
not Python in-memory footprint.
* ref: updates to telemetry writer PR (#13294)
* fix(telemetry-writer): harden error handling and add missing test coverage
Critical:
- teardown() now re-raises CancelledError after awaiting the writer task so the
asyncio cancellation chain propagates correctly on lifespan task kill
- suppress(OSError) on owner file write replaced with explicit error log so
operators know when disk-spilled rows will not be recoverable on restart
High / important:
- Escalation threshold check changed from == to >= so the error log fires on
every failure past the threshold, not just the 6th
- Dead `except OSError` on time.time() - time.monotonic() changed to
`except Exception` since time functions cannot raise OSError
- Removed redundant `from uuid import UUID` inside _run_retention_pass (already
imported at module top)
- Orphan directory cleanup: replaced blanket suppress(OSError) with per-child
suppress so ENOTEMPTY on an individual rmdir doesn't abort the whole loop and
leave the parent directory leaking silently; outer failure now logs at debug
Tests (3 new):
- test_retention_sweep_caps_vertex_builds_per_vertex: inserts 8 builds for a
single vertex with max_per_vertex=3 so the per-vertex DELETE subquery
actually executes (previous test used 8 distinct vertex IDs, bypassing it)
- test_either_strategy_trips_on_bytes_first: verifies bytes can be the first
trigger under 'either' strategy (previous test only covered the count-first path)
- test_writer_retries_on_batch_failure: injects 2 flush failures then success;
confirms failed_batches increments, rows are preserved in the buffer, and
flushed_rows reflects the final successful write
* ci: add stress-tests job to nightly build pipeline
Wires the stress-tests workflow into nightly so telemetry write stress
tests run automatically. Also gates release-nightly-build and
slack-notification on stress-tests result so a stress regression blocks
the nightly release and surfaces in Slack.
* ci: run stress tests in nightly without blocking release
Stress tests are informational for now — a failure is visible in the
workflow run and Slack but does not gate the nightly release or build.
* fix(telemetry-writer): address PR review on cancelled teardown and nightly stress tests
- Add the stress-tests.yml reusable workflow (was untracked, so the
nightly stress-tests job could never resolve)
- Notify Slack on stress-tests failure (add to slack-notification gate
and FAILED_JOB detection as non-blocking)
- Run sweeper cancel, disk spill, and engine dispose in a finally so a
cancelled teardown() still persists the in-memory buffer
- Add tests for the cancelled-teardown spill path and the >= escalation
threshold; drop the misleading wait-loop in the retry test
---------
Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com>
Co-authored-by: Jordan Frazier <122494242+jordanrfrazier@users.noreply.github.com>