Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
7a45e72
perf(telemetry): batched off-pool writer for transactions + vertex_bu…
ogabrielluiz May 14, 2026
a5f1677
perf(telemetry): tighten error handling and add coverage
ogabrielluiz May 14, 2026
9358bdc
perf(telemetry): address copilot review
ogabrielluiz May 14, 2026
1c61e7d
perf(telemetry): address coderabbit review
ogabrielluiz May 14, 2026
651ecf7
perf(telemetry): swap diskcache outbox for stdlib sqlite (CVE-2025-69…
ogabrielluiz May 15, 2026
4ca231c
[autofix.ci] apply automated fixes
autofix-ci[bot] May 15, 2026
e38f058
fix(telemetry): name wait_for inner tasks so pyleak can filter
ogabrielluiz May 15, 2026
ac48873
Merge branch 'release-1.10.0' into perf/telemetry-writer
ogabrielluiz May 15, 2026
eca0f2e
[autofix.ci] apply automated fixes
autofix-ci[bot] May 15, 2026
1120375
[autofix.ci] apply automated fixes (attempt 2/3)
autofix-ci[bot] May 15, 2026
2e6629b
test(telemetry): disable writer in tests so reads see writes synchron…
ogabrielluiz May 15, 2026
b08167d
perf(telemetry): age out cross-host orphan outboxes on shared volumes
ogabrielluiz May 15, 2026
6b329ea
perf(telemetry): add byte-aware flush + drop strategy
ogabrielluiz May 15, 2026
9c22aa7
test(telemetry): tighten byte-strategy assertions and cover end-to-en…
ogabrielluiz May 15, 2026
88e32a1
Merge branch 'release-1.10.0' into perf/telemetry-writer
ogabrielluiz May 22, 2026
2edfa33
ref: updates to telemetry writer PR (#13294)
jordanrfrazier Jun 2, 2026
0f2f790
Merge branch 'release-1.10.0' into perf/telemetry-writer
ogabrielluiz Jun 5, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 14 additions & 4 deletions .github/workflows/nightly_build.yml
Original file line number Diff line number Diff line change
Expand Up @@ -283,11 +283,19 @@ jobs:
# python-versions: '["3.10", "3.11", "3.12", "3.13"]'
# ref: ${{ needs.create-nightly-tag.outputs.tag }}

stress-tests:
if: github.repository == 'langflow-ai/langflow' && !inputs.skip_backend_tests
name: Run Stress Tests
needs: [resolve-release-branch, create-nightly-tag]
uses: ./.github/workflows/stress-tests.yml
with:
ref: ${{ needs.resolve-release-branch.outputs.branch }}

release-nightly-build:
if: github.repository == 'langflow-ai/langflow' && always() && needs.create-nightly-tag.result == 'success' && needs.frontend-tests-linux.result != 'failure' && needs.backend-unit-tests.result != 'failure'
name: Run Nightly Langflow Build
needs:
[validate-inputs, create-nightly-tag, frontend-tests-linux, frontend-tests-windows, backend-unit-tests]
[validate-inputs, create-nightly-tag, frontend-tests-linux, frontend-tests-windows, backend-unit-tests, stress-tests]
uses: ./.github/workflows/release_nightly.yml
with:
build_docker_base: true
Expand Down Expand Up @@ -316,12 +324,12 @@ jobs:

slack-notification:
name: Send Slack Notification
needs: [frontend-tests-linux, frontend-tests-windows, backend-unit-tests, release-nightly-build, db-migration-validation]
if: ${{ github.repository == 'langflow-ai/langflow' && !inputs.skip_slack && always() && (needs.release-nightly-build.result == 'failure' || needs.frontend-tests-linux.result == 'failure' || needs.frontend-tests-windows.result == 'failure' || needs.backend-unit-tests.result == 'failure' || needs.db-migration-validation.result == 'failure' || needs.release-nightly-build.result == 'success') }}
needs: [frontend-tests-linux, frontend-tests-windows, backend-unit-tests, stress-tests, release-nightly-build, db-migration-validation]
if: ${{ github.repository == 'langflow-ai/langflow' && !inputs.skip_slack && always() && (needs.release-nightly-build.result == 'failure' || needs.frontend-tests-linux.result == 'failure' || needs.frontend-tests-windows.result == 'failure' || needs.backend-unit-tests.result == 'failure' || needs.stress-tests.result == 'failure' || needs.db-migration-validation.result == 'failure' || needs.release-nightly-build.result == 'success') }}
runs-on: ubuntu-latest
steps:
- name: Send failure notification to Slack
if: ${{ needs.release-nightly-build.result == 'failure' || needs.frontend-tests-linux.result == 'failure' || needs.frontend-tests-windows.result == 'failure' || needs.backend-unit-tests.result == 'failure' || needs.db-migration-validation.result == 'failure' }}
if: ${{ needs.release-nightly-build.result == 'failure' || needs.frontend-tests-linux.result == 'failure' || needs.frontend-tests-windows.result == 'failure' || needs.backend-unit-tests.result == 'failure' || needs.stress-tests.result == 'failure' || needs.db-migration-validation.result == 'failure' }}
run: |
# Determine which job failed
FAILED_JOB="unknown"
Expand All @@ -333,6 +341,8 @@ jobs:
FAILED_JOB="frontend-tests-windows (non-blocking)"
elif [ "${{ needs.backend-unit-tests.result }}" == "failure" ]; then
FAILED_JOB="backend-unit-tests"
elif [ "${{ needs.stress-tests.result }}" == "failure" ]; then
FAILED_JOB="stress-tests (non-blocking)"
elif [ "${{ needs.db-migration-validation.result }}" == "failure" ]; then
FAILED_JOB="db-migration-validation"
fi
Expand Down
56 changes: 56 additions & 0 deletions .github/workflows/stress-tests.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
name: Stress Tests

on:
workflow_dispatch: {}
workflow_call:
inputs:
ref:
description: "Branch or tag to test"
required: false
type: string
default: ""

jobs:
telemetry-writes:
name: Telemetry Writes (Postgres)
runs-on: ubuntu-latest

services:
postgres:
image: postgres:16
env:
POSTGRES_USER: langflow
POSTGRES_PASSWORD: langflow
POSTGRES_DB: langflow
ports:
- 5432:5432
options: >-
--health-cmd="pg_isready -U langflow"
--health-interval=10s
--health-timeout=5s
--health-retries=5

steps:
- name: Checkout code
uses: actions/checkout@v6
with:
ref: ${{ inputs.ref || github.sha }}

- name: Install uv
uses: astral-sh/setup-uv@v6
with:
enable-cache: true
cache-dependency-glob: "uv.lock"
python-version: "3.12"
prune-cache: false

- name: Install dependencies
run: uv sync --extra postgresql

- name: Run stress test
env:
DB_URL: "postgresql+psycopg://langflow:langflow@localhost:5432/langflow" # pragma: allowlist secret
LANGFLOW_TELEMETRY_WRITER_ENABLED: "true"
run: |
uv run python src/backend/tests/stress/stress_telemetry_writes.py \
--concurrency 50 --seconds 15
23 changes: 23 additions & 0 deletions .secrets.baseline
Original file line number Diff line number Diff line change
Expand Up @@ -2352,6 +2352,29 @@
"line_number": 774
}
],
"src/backend/base/langflow/initial_setup/starter_projects/Pok\u00e9dex Agent.json": [
{
"type": "Hex High Entropy String",
"filename": "src/backend/base/langflow/initial_setup/starter_projects/Pok\u00e9dex Agent.json",
"hashed_secret": "54ed260e3bc31bc77ee06754dff850981d39a66c",
"is_verified": false,
"line_number": 115
},
{
"type": "Hex High Entropy String",
"filename": "src/backend/base/langflow/initial_setup/starter_projects/Pok\u00e9dex Agent.json",
"hashed_secret": "d6e6d7b4b115cd3b9d172623199f8c403055fecc",
"is_verified": false,
"line_number": 394
},
{
"type": "Hex High Entropy String",
"filename": "src/backend/base/langflow/initial_setup/starter_projects/Pok\u00e9dex Agent.json",
"hashed_secret": "ecdbe30d19e36761df3620b37270f23c8e3eaa4c",
"is_verified": false,
"line_number": 774
}
],
"src/backend/base/langflow/initial_setup/starter_projects/Portfolio Website Code Generator.json": [
{
"type": "Hex High Entropy String",
Expand Down
7 changes: 7 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -285,6 +285,13 @@ external = ["RUF027"]

[tool.ruff.lint.per-file-ignores]
"scripts/*" = ["D1", "INP", "T201"]
"src/backend/tests/stress/*" = [
"D1", # Docstrings: this is a manual CLI tool, not a test target
"INP001", # Stress tests live outside the unit test package
"S101", # Asserts ok in stress harness
"S311", # Standard PRNG is fine for stress payload jitter
"T201", # Print statements: this is a CLI script
]
"scripts/gp/tests/*" = [
"D1", # Missing docstrings
"PLR2004", # Magic value comparisons
Expand Down
17 changes: 17 additions & 0 deletions src/backend/base/langflow/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -198,6 +198,23 @@ async def lifespan(_app: FastAPI):
await initialize_services(fix_migration=fix_migration)
await logger.adebug(f"Services initialized in {asyncio.get_event_loop().time() - start_time:.2f}s")

# Start the telemetry writer (no-op when telemetry_writer_enabled is False).
try:
from langflow.services.deps import get_telemetry_writer_service

telemetry_writer = get_telemetry_writer_service()
if telemetry_writer is not None and telemetry_writer.is_enabled():
await telemetry_writer.start()
except Exception as exc: # noqa: BLE001
# If the user explicitly opted in (telemetry_writer_enabled=True)
# but startup failed, this is an error not a warning — every
# subsequent write will silently fall back to the legacy direct-
# write path that this feature was built to replace.
await logger.aerror(
f"Failed to start telemetry writer; transactions and vertex_build "
f"writes will use the legacy direct-write path: {exc}"
)

current_time = asyncio.get_event_loop().time()
await logger.adebug("Setting up LLM caching")
setup_llm_caching()
Expand Down
7 changes: 7 additions & 0 deletions src/backend/base/langflow/services/deps.py
Original file line number Diff line number Diff line change
Expand Up @@ -284,3 +284,10 @@ def get_memory_base_service():
from langflow.services.memory_base.factory import MemoryBaseServiceFactory

return get_service(ServiceType.MEMORY_BASE_SERVICE, MemoryBaseServiceFactory())


def get_telemetry_writer_service():
"""Return the TelemetryWriterService instance (always registered when langflow is installed)."""
from langflow.services.telemetry_writer.factory import TelemetryWriterServiceFactory

return get_service(ServiceType.TELEMETRY_WRITER_SERVICE, TelemetryWriterServiceFactory())
1 change: 1 addition & 0 deletions src/backend/base/langflow/services/schema.py
Original file line number Diff line number Diff line change
Expand Up @@ -24,3 +24,4 @@ class ServiceType(str, Enum):
JOB_SERVICE = "jobs_service"
FLOW_EVENTS_SERVICE = "flow_events_service"
MEMORY_BASE_SERVICE = "memory_base_service"
TELEMETRY_WRITER_SERVICE = "telemetry_writer_service"
12 changes: 12 additions & 0 deletions src/backend/base/langflow/services/telemetry_writer/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
"""Telemetry writer service.

Async batched writer for transaction and vertex_build rows backed by a
disk-persisted SQLite outbox and a dedicated database connection. Removes
write-path contention from the request-handling connection pool so heavy load
no longer triggers SQLite "database is locked" or PostgreSQL pool timeouts.
"""

from langflow.services.telemetry_writer.factory import TelemetryWriterServiceFactory
from langflow.services.telemetry_writer.service import TelemetryWriterService

__all__ = ["TelemetryWriterService", "TelemetryWriterServiceFactory"]
20 changes: 20 additions & 0 deletions src/backend/base/langflow/services/telemetry_writer/factory.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
from __future__ import annotations

from typing import TYPE_CHECKING

from typing_extensions import override

from langflow.services.factory import ServiceFactory
from langflow.services.telemetry_writer.service import TelemetryWriterService

if TYPE_CHECKING:
from langflow.services.settings.service import SettingsService


class TelemetryWriterServiceFactory(ServiceFactory):
def __init__(self) -> None:
super().__init__(TelemetryWriterService)

@override
def create(self, settings_service: SettingsService):
return TelemetryWriterService(settings_service)
Loading
Loading