From 87f2d98b52c6cfa245185ade037bb2583fa235e0 Mon Sep 17 00:00:00 2001 From: Hamza Rashid Date: Fri, 24 Apr 2026 03:13:27 +0000 Subject: [PATCH] feat: Add telemetry to deployments API Instruments 8 CUD-shaped deployment routes to emit telemetry events for tracking usage, duration, and error rates. Changes: - `schema.py`: Added `DeploymentPayload` Pydantic model to define the structure of telemetry data for deployment events (action, provider, seconds, success, error_type). - `service.py`: Added three async methods (`log_package_deployment`, `log_package_deployment_provider`, `log_package_deployment_run`) to enqueue deployment events to the telemetry queue. - `deployments.py`: - Created `DeploymentTelemetryCtx` dataclass and a non-invasive yield-based FastAPI dependency (`_make_telemetry_dep`) to handle timing, success/failure detection, and event emission. - Instrumented 8 routes (create/update/delete for deployments and providers, create run, update snapshot) by injecting the telemetry dependency and setting the provider context. - `test_telemetry_schema.py`: Added unit tests for `DeploymentPayload` initialization, defaults, serialization, and roundtrip. - `test_telemetry.py`: Added unit tests for the new `log_package_deployment*` service methods, including a check for the `do_not_track` setting. - `test_deployments_telemetry.py`: Created a new integration test file with 10 tests covering happy paths, error paths (e.g., `HTTPException` mapping), and cross-route smoke tests for all instrumented endpoints. --- .../base/langflow/api/v1/deployments.py | 80 ++++- .../langflow/services/telemetry/schema.py | 8 + .../langflow/services/telemetry/service.py | 10 + .../unit/api/v1/test_deployments_telemetry.py | 305 ++++++++++++++++++ .../telemetry/test_telemetry_schema.py | 81 +++++ src/backend/tests/unit/test_telemetry.py | 77 +++++ 6 files changed, 555 insertions(+), 6 deletions(-) create mode 100644 src/backend/tests/unit/api/v1/test_deployments_telemetry.py diff --git a/src/backend/base/langflow/api/v1/deployments.py b/src/backend/base/langflow/api/v1/deployments.py index ad056d6233..d24875e30f 100644 --- a/src/backend/base/langflow/api/v1/deployments.py +++ b/src/backend/base/langflow/api/v1/deployments.py @@ -1,5 +1,8 @@ from __future__ import annotations +import time +from collections.abc import AsyncIterator +from dataclasses import dataclass from typing import Annotated from uuid import UUID @@ -108,6 +111,55 @@ from langflow.services.database.models.flow_version_deployment_attachment.crud i list_deployment_attachments_for_flow_version_ids, update_flow_version_by_provider_snapshot_id, ) +from langflow.services.deps import get_telemetry_service +from langflow.services.telemetry.schema import DeploymentPayload + + +@dataclass +class DeploymentTelemetryCtx: + """Mutable context that routes write `provider` into; passed via Depends.""" + + provider: str = "unknown" + + +def _make_telemetry_dep(action: str, log_method_name: str): + async def _dep() -> AsyncIterator[DeploymentTelemetryCtx]: + ctx = DeploymentTelemetryCtx() + started_at = time.perf_counter() + success: bool = True + error: Exception | None = None + try: + yield ctx + except Exception as exc: + success = False + error = exc + raise + finally: + try: + ts = get_telemetry_service() + payload = DeploymentPayload( + deployment_action=action, + deployment_provider=ctx.provider, + deployment_seconds=time.perf_counter() - started_at, + deployment_success=success, + deployment_error_type=type(error).__name__ if error else None, + ) + await getattr(ts, log_method_name)(payload) + except Exception: # noqa: BLE001 + logger.debug("deployment telemetry emit failed", exc_info=True) + + return _dep + + +deployment_create_telemetry = _make_telemetry_dep("deployment.create", "log_package_deployment") +deployment_update_telemetry = _make_telemetry_dep("deployment.update", "log_package_deployment") +deployment_delete_telemetry = _make_telemetry_dep("deployment.delete", "log_package_deployment") +deployment_run_telemetry = _make_telemetry_dep("deployment.run", "log_package_deployment_run") +provider_create_telemetry = _make_telemetry_dep("provider.create", "log_package_deployment_provider") +provider_update_telemetry = _make_telemetry_dep("provider.update", "log_package_deployment_provider") +provider_delete_telemetry = _make_telemetry_dep("provider.delete", "log_package_deployment_provider") +snapshot_update_telemetry = _make_telemetry_dep("snapshot.update", "log_package_deployment") + router = APIRouter(prefix="/deployments", tags=["Deployments"], include_in_schema=False) @@ -204,7 +256,7 @@ async def _count_provider_deployments_after_reconciliation( size=deployment_count, deployment_type=None, ) - except Exception: # noqa: BLE001 + except Exception: # noqa: BLE001 # noqa: BLE001 logger.warning( "Failed to reconcile deployments before deleting provider account %s; falling back to local count.", provider_account.id, @@ -232,7 +284,7 @@ async def _delete_local_deployment_row_with_commit_retry( try: await delete_deployment_by_id(session, user_id=user_id, deployment_id=deployment_id) await session.commit() - except Exception: # noqa: BLE001 + except Exception: # noqa: BLE001 # noqa: BLE001 await session.rollback() logger.warning( "Local deployment cleanup failed for deployment %s (resource_key=%s) after provider delete; retrying.", @@ -266,7 +318,9 @@ async def create_provider_account( session: DbSession, payload: DeploymentProviderAccountCreateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(provider_create_telemetry)], ): + telemetry.provider = payload.provider_key deployment_mapper = get_deployment_mapper(payload.provider_key) deployment_adapter = resolve_deployment_adapter(payload.provider_key) @@ -337,12 +391,14 @@ async def delete_provider_account( provider_id: DeploymentProviderAccountIdPath, session: DbSession, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(provider_delete_telemetry)], ): provider_account = await get_owned_provider_account_or_404( provider_id=provider_id, user_id=current_user.id, db=session, ) + telemetry.provider = provider_account.provider_key deployment_count = await _count_provider_deployments_after_reconciliation( session=session, provider_account=provider_account, @@ -370,12 +426,14 @@ async def update_provider_account( session: DbSession, payload: DeploymentProviderAccountUpdateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(provider_update_telemetry)], ): provider_account = await get_owned_provider_account_or_404( provider_id=provider_id, user_id=current_user.id, db=session, ) + telemetry.provider = provider_account.provider_key deployment_mapper = get_deployment_mapper(provider_account.provider_key) verify_input = None @@ -420,6 +478,7 @@ async def create_deployment( session: DbSession, payload: DeploymentCreateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(deployment_create_telemetry)], ): provider_id = payload.provider_id provider_account = await get_owned_provider_account_or_404( @@ -427,6 +486,7 @@ async def create_deployment( user_id=current_user.id, db=session, ) + telemetry.provider = provider_account.provider_key # fail fast if the deployment name already exists # we could have races but that is more # acceptable than provider-side rollback failure @@ -749,12 +809,14 @@ async def create_deployment_run( session: DbSession, payload: RunCreateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(deployment_run_telemetry)], ): deployment_row, deployment_adapter, deployment_mapper, _provider_key = await resolve_adapter_mapper_from_deployment( deployment_id=deployment_id, user_id=current_user.id, db=session, ) + telemetry.provider = _provider_key adapter_execution_payload = await deployment_mapper.resolve_execution_create( deployment_resource_key=deployment_row.resource_key, db=session, @@ -950,6 +1012,7 @@ async def update_snapshot( body: SnapshotUpdateRequest, session: DbSession, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(snapshot_update_telemetry)], ): """Replace an existing provider snapshot's content with a new flow version. @@ -1005,6 +1068,7 @@ async def update_snapshot( user_id=current_user.id, db=session, ) + telemetry.provider = provider_account.provider_key deployment_adapter = resolve_deployment_adapter(provider_account.provider_key) deployment_mapper = get_deployment_mapper(provider_account.provider_key) @@ -1087,7 +1151,7 @@ async def update_snapshot( snapshot_id, previous_flow_version_id, ) - except Exception: # noqa: BLE001 + except Exception: # noqa: BLE001 # noqa: BLE001 logger.warning( "Best-effort rollback failed for snapshot '%s'. " "Provider content reflects flow_version_id=%s but attachment " @@ -1140,7 +1204,7 @@ async def get_deployment( try: await delete_deployment_by_id(session, user_id=current_user.id, deployment_id=deployment_row.id) await session.commit() - except Exception: # noqa: BLE001 + except Exception: # noqa: BLE001 # noqa: BLE001 logger.warning( "Failed to delete stale deployment row %s; returning 404 anyway", deployment_row.id, @@ -1184,7 +1248,7 @@ async def get_deployment( session, user_id=current_user.id, deployment_id=deployment_row.id ) attached_count = len(attachments) - except Exception: # noqa: BLE001 + except Exception: # noqa: BLE001 # noqa: BLE001 logger.warning( "Binding-aware sync failed for deployment %s; returning unverified attachment count", deployment_row.id, @@ -1196,7 +1260,7 @@ async def get_deployment( session, user_id=current_user.id, deployment_id=deployment_row.id ) attached_count = len(attachments) - except Exception: # noqa: BLE001 + except Exception: # noqa: BLE001 # noqa: BLE001 logger.warning( "Fallback attachment count query also failed for deployment %s; defaulting to 0", deployment_row.id, @@ -1232,12 +1296,14 @@ async def update_deployment( session: DbSession, payload: DeploymentUpdateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(deployment_update_telemetry)], ): deployment_row, deployment_adapter, deployment_mapper, provider_key = await resolve_adapter_mapper_from_deployment( deployment_id=deployment_id, user_id=current_user.id, db=session, ) + telemetry.provider = provider_key deployment_row_id = deployment_row.id deployment_resource_key = deployment_row.resource_key deployment_provider_account_id = deployment_row.deployment_provider_account_id @@ -1334,6 +1400,7 @@ async def delete_deployment( deployment_id: DeploymentIdPath, session: DbSession, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(deployment_delete_telemetry)], *, include_provider: IncludeProviderDeleteQuery = True, ): @@ -1342,6 +1409,7 @@ async def delete_deployment( user_id=current_user.id, db=session, ) + telemetry.provider = _provider_key if include_provider: try: with handle_adapter_errors(), deployment_provider_scope(deployment_row.deployment_provider_account_id): diff --git a/src/backend/base/langflow/services/telemetry/schema.py b/src/backend/base/langflow/services/telemetry/schema.py index eecb38694d..64fab76e85 100644 --- a/src/backend/base/langflow/services/telemetry/schema.py +++ b/src/backend/base/langflow/services/telemetry/schema.py @@ -19,6 +19,14 @@ class RunPayload(BasePayload): run_id: str | None = Field(None, serialization_alias="runId") +class DeploymentPayload(BasePayload): + deployment_action: str = Field(serialization_alias="deploymentAction") + deployment_provider: str = Field(serialization_alias="deploymentProvider") + deployment_seconds: float = Field(serialization_alias="deploymentSeconds") + deployment_success: bool = Field(serialization_alias="deploymentSuccess") + deployment_error_type: str | None = Field(default=None, serialization_alias="deploymentErrorType") + + class ShutdownPayload(BasePayload): time_running: int = Field(serialization_alias="timeRunning") diff --git a/src/backend/base/langflow/services/telemetry/service.py b/src/backend/base/langflow/services/telemetry/service.py index 6bdf34e3f8..4bafdfd42f 100644 --- a/src/backend/base/langflow/services/telemetry/service.py +++ b/src/backend/base/langflow/services/telemetry/service.py @@ -18,6 +18,7 @@ from langflow.services.telemetry.schema import ( ComponentIndexPayload, ComponentInputsPayload, ComponentPayload, + DeploymentPayload, EmailPayload, ExceptionPayload, PlaygroundPayload, @@ -110,6 +111,15 @@ class TelemetryService(Service): async def log_package_run(self, payload: RunPayload) -> None: await self._queue_event((self.send_telemetry_data, payload, "run")) + async def log_package_deployment(self, payload: DeploymentPayload) -> None: + await self._queue_event((self.send_telemetry_data, payload, "deployment")) + + async def log_package_deployment_provider(self, payload: DeploymentPayload) -> None: + await self._queue_event((self.send_telemetry_data, payload, "deployment_provider")) + + async def log_package_deployment_run(self, payload: DeploymentPayload) -> None: + await self._queue_event((self.send_telemetry_data, payload, "deployment_run")) + async def log_package_shutdown(self) -> None: payload = ShutdownPayload(time_running=(datetime.now(timezone.utc) - self._start_time).seconds) await self._queue_event(payload) diff --git a/src/backend/tests/unit/api/v1/test_deployments_telemetry.py b/src/backend/tests/unit/api/v1/test_deployments_telemetry.py new file mode 100644 index 0000000000..134cd28ace --- /dev/null +++ b/src/backend/tests/unit/api/v1/test_deployments_telemetry.py @@ -0,0 +1,305 @@ +# ruff: noqa: ARG001 +from __future__ import annotations + +from contextlib import ExitStack +from typing import TYPE_CHECKING +from unittest.mock import AsyncMock, MagicMock, patch +from uuid import uuid4 + +import pytest +from fastapi import status + +if TYPE_CHECKING: + from httpx import AsyncClient + +# We'll use a mocked adapter so we don't need real credentials. +# We need to mock the adapter resolution and the telemetry service. + + +@pytest.fixture +def mock_telemetry_service(): + with patch("langflow.api.v1.deployments.get_telemetry_service") as mock_get: + mock_ts = AsyncMock() + mock_get.return_value = mock_ts + yield mock_ts + + +@pytest.fixture +def mock_adapter(): + with patch("langflow.api.v1.deployments.resolve_deployment_adapter") as mock_resolve: + mock_ad = AsyncMock() + mock_resolve.return_value = mock_ad + yield mock_ad + + +@pytest.fixture +def mock_mapper(): + with patch("langflow.api.v1.deployments.get_deployment_mapper") as mock_get: + mock_map = MagicMock() + # Ensure it returns something valid for verify_credentials + mock_map.resolve_verify_credentials_for_create.return_value = {} + mock_map.resolve_verify_credentials_for_update.return_value = {} + mock_map.resolve_provider_account_create.return_value = AsyncMock( + id=uuid4(), provider_key="watsonx-orchestrate" + ) + mock_map.resolve_provider_account_response.return_value = { + "id": str(uuid4()), + "provider_key": "watsonx-orchestrate", + "name": "Test", + } + mock_map.util_existing_deployment_resource_key_for_create.return_value = None + mock_map.resolve_deployment_create = AsyncMock(return_value={}) + mock_map.resolve_deployment_update = AsyncMock(return_value={}) + mock_map.util_create_flow_version_ids.return_value = [] + mock_map.shape_deployment_create_result.return_value = { + "id": str(uuid4()), + "provider_id": str(uuid4()), + "provider_key": "watsonx-orchestrate", + "name": "Test", + "type": "agent", + "resource_key": "res-1", + } + mock_map.shape_deployment_update_result.return_value = { + "id": str(uuid4()), + "provider_id": str(uuid4()), + "provider_key": "watsonx-orchestrate", + "name": "Test", + "type": "agent", + "resource_key": "res-1", + } + mock_map.resolve_execution_create = AsyncMock(return_value={}) + mock_map.shape_execution_create_result.return_value = {"id": "run-1", "deployment_id": str(uuid4())} + mock_map.resolve_snapshot_update_artifact.return_value = {} + mock_get.return_value = mock_map + yield mock_map + + +@pytest.fixture +def mock_db_crud(mock_mapper): + with ExitStack() as stack: + mock_create = stack.enter_context(patch("langflow.api.v1.deployments.create_provider_account_row")) + mock_get_owned = stack.enter_context(patch("langflow.api.v1.deployments.get_owned_provider_account_or_404")) + _mock_del_prov = stack.enter_context(patch("langflow.api.v1.deployments.delete_provider_account_row")) + _mock_upd_prov = stack.enter_context(patch("langflow.api.v1.deployments.update_provider_account_row")) + mock_name_exists = stack.enter_context(patch("langflow.api.v1.deployments.deployment_name_exists")) + mock_proj_id = stack.enter_context( + patch("langflow.api.v1.deployments.resolve_project_id_for_deployment_create") + ) + mock_create_dep = stack.enter_context(patch("langflow.api.v1.deployments.create_deployment_db")) + _mock_attach = stack.enter_context(patch("langflow.api.v1.deployments.attach_flow_versions")) + mock_res_am = stack.enter_context(patch("langflow.api.v1.deployments.resolve_adapter_mapper_from_deployment")) + mock_res_patch = stack.enter_context(patch("langflow.api.v1.deployments.resolve_flow_version_patch_for_update")) + _mock_val_proj = stack.enter_context( + patch("langflow.api.v1.deployments.validate_project_scoped_flow_version_ids") + ) + mock_list_att = stack.enter_context( + patch("langflow.api.v1.deployments.list_deployment_attachments_for_flow_version_ids") + ) + _mock_apply_patch = stack.enter_context( + patch("langflow.api.v1.deployments.apply_flow_version_patch_attachments") + ) + mock_upd_dep = stack.enter_context(patch("langflow.api.v1.deployments.update_deployment_db")) + mock_res_ad = stack.enter_context(patch("langflow.api.v1.deployments.resolve_adapter_from_deployment")) + _mock_del_dep = stack.enter_context( + patch("langflow.api.v1.deployments._delete_local_deployment_row_with_commit_retry") + ) + mock_get_att = stack.enter_context(patch("langflow.api.v1.deployments.get_attachment_by_provider_snapshot_id")) + mock_get_dep_row = stack.enter_context( + patch("langflow.services.database.models.deployment.crud.get_deployment") + ) + mock_get_fv = stack.enter_context( + patch("langflow.services.database.models.flow_version.crud.get_flow_version_entry") + ) + mock_upd_fv = stack.enter_context( + patch("langflow.api.v1.deployments.update_flow_version_by_provider_snapshot_id") + ) + mock_count_deps = stack.enter_context( + patch("langflow.api.v1.deployments._count_provider_deployments_after_reconciliation") + ) + + mock_create.return_value = AsyncMock(id=uuid4(), provider_key="watsonx-orchestrate") + mock_get_owned.return_value = AsyncMock(id=uuid4(), provider_key="watsonx-orchestrate") + mock_name_exists.return_value = False + mock_proj_id.return_value = uuid4() + mock_create_dep.return_value = AsyncMock(id=uuid4()) + + mock_res_am.return_value = ( + AsyncMock(id=uuid4(), resource_key="res-1", deployment_provider_account_id=uuid4(), project_id=uuid4()), + AsyncMock(), + mock_mapper, + "watsonx-orchestrate", + ) + mock_res_patch.return_value = ([], []) + mock_list_att.return_value = [] + mock_upd_dep.return_value = AsyncMock() + mock_res_ad.return_value = ( + AsyncMock(id=uuid4(), resource_key="res-1", deployment_provider_account_id=uuid4()), + AsyncMock(), + "watsonx-orchestrate", + ) + mock_get_att.return_value = AsyncMock(deployment_id=uuid4(), flow_version_id=uuid4()) + mock_get_dep_row.return_value = AsyncMock(deployment_provider_account_id=uuid4()) + mock_get_fv.return_value = AsyncMock(flow_id=uuid4(), data={}) + mock_upd_fv.return_value = 1 + mock_count_deps.return_value = 0 + + yield + + +@pytest.mark.asyncio +async def test_create_provider_account_telemetry( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + response = await client.post( + "api/v1/deployments/providers", + json={"provider_key": "watsonx-orchestrate", "name": "Test", "provider_data": {"foo": "bar"}}, + headers=logged_in_headers, + ) + assert response.status_code == status.HTTP_201_CREATED + mock_telemetry_service.log_package_deployment_provider.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment_provider.call_args[0][0] + assert payload.deployment_action == "provider.create" + assert payload.deployment_provider == "watsonx-orchestrate" + assert payload.deployment_success is True + + +@pytest.mark.asyncio +async def test_update_provider_account_telemetry( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + response = await client.patch( + f"api/v1/deployments/providers/{uuid4()}", + json={"name": "Test", "provider_data": {"foo": "bar"}}, + headers=logged_in_headers, + ) + assert response.status_code == status.HTTP_200_OK + mock_telemetry_service.log_package_deployment_provider.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment_provider.call_args[0][0] + assert payload.deployment_action == "provider.update" + assert payload.deployment_provider == "watsonx-orchestrate" + assert payload.deployment_success is True + + +@pytest.mark.asyncio +async def test_delete_provider_account_telemetry( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + response = await client.delete(f"api/v1/deployments/providers/{uuid4()}", headers=logged_in_headers) + assert response.status_code == status.HTTP_204_NO_CONTENT + mock_telemetry_service.log_package_deployment_provider.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment_provider.call_args[0][0] + assert payload.deployment_action == "provider.delete" + assert payload.deployment_provider == "watsonx-orchestrate" + assert payload.deployment_success is True + + +@pytest.mark.asyncio +async def test_create_deployment_telemetry( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + response = await client.post( + "api/v1/deployments", + json={"provider_id": str(uuid4()), "name": "Test", "type": "agent", "provider_data": {}}, + headers=logged_in_headers, + ) + assert response.status_code == status.HTTP_201_CREATED, response.json() + mock_telemetry_service.log_package_deployment.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment.call_args[0][0] + assert payload.deployment_action == "deployment.create" + assert payload.deployment_provider == "watsonx-orchestrate" + assert payload.deployment_success is True + + +@pytest.mark.asyncio +async def test_update_deployment_telemetry( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + response = await client.patch(f"api/v1/deployments/{uuid4()}", json={"name": "Test"}, headers=logged_in_headers) + assert response.status_code == status.HTTP_200_OK + mock_telemetry_service.log_package_deployment.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment.call_args[0][0] + assert payload.deployment_action == "deployment.update" + assert payload.deployment_provider == "watsonx-orchestrate" + assert payload.deployment_success is True + + +@pytest.mark.asyncio +async def test_delete_deployment_telemetry( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + response = await client.delete(f"api/v1/deployments/{uuid4()}", headers=logged_in_headers) + assert response.status_code == status.HTTP_204_NO_CONTENT + mock_telemetry_service.log_package_deployment.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment.call_args[0][0] + assert payload.deployment_action == "deployment.delete" + assert payload.deployment_provider == "watsonx-orchestrate" + assert payload.deployment_success is True + + +@pytest.mark.asyncio +async def test_create_deployment_run_telemetry( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + response = await client.post( + f"api/v1/deployments/{uuid4()}/runs", json={"provider_data": {}}, headers=logged_in_headers + ) + assert response.status_code == status.HTTP_201_CREATED, response.json() + mock_telemetry_service.log_package_deployment_run.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment_run.call_args[0][0] + assert payload.deployment_action == "deployment.run" + assert payload.deployment_provider == "watsonx-orchestrate" + assert payload.deployment_success is True + + +@pytest.mark.asyncio +async def test_update_snapshot_telemetry( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + response = await client.patch( + "api/v1/deployments/snapshots/snap-1", json={"flow_version_id": str(uuid4())}, headers=logged_in_headers + ) + assert response.status_code == status.HTTP_200_OK + mock_telemetry_service.log_package_deployment.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment.call_args[0][0] + assert payload.deployment_action == "snapshot.update" + assert payload.deployment_provider == "watsonx-orchestrate" + assert payload.deployment_success is True + + +@pytest.mark.asyncio +async def test_create_provider_account_telemetry_error( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + mock_adapter.verify_credentials.side_effect = ValueError("Invalid credentials") + response = await client.post( + "api/v1/deployments/providers", + json={"provider_key": "watsonx-orchestrate", "name": "Test", "provider_data": {"foo": "bar"}}, + headers=logged_in_headers, + ) + assert response.status_code == status.HTTP_400_BAD_REQUEST + mock_telemetry_service.log_package_deployment_provider.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment_provider.call_args[0][0] + assert payload.deployment_action == "provider.create" + assert payload.deployment_provider == "watsonx-orchestrate" + assert payload.deployment_success is False + assert payload.deployment_error_type == "HTTPException" + + +@pytest.mark.asyncio +async def test_cross_route_smoke_exception_after_provider_set( + client: AsyncClient, mock_telemetry_service, mock_adapter, mock_mapper, mock_db_crud, logged_in_headers +): + # Simulate an error during adapter.create after the provider has been set in the route + mock_adapter.create.side_effect = RuntimeError("Something went wrong") + response = await client.post( + "api/v1/deployments", + json={"provider_id": str(uuid4()), "name": "Test", "type": "agent", "provider_data": {}}, + headers=logged_in_headers, + ) + assert response.status_code == status.HTTP_500_INTERNAL_SERVER_ERROR + mock_telemetry_service.log_package_deployment.assert_awaited_once() + payload = mock_telemetry_service.log_package_deployment.call_args[0][0] + assert payload.deployment_action == "deployment.create" + assert payload.deployment_provider == "watsonx-orchestrate" # Provider should be captured! + assert payload.deployment_success is False + assert payload.deployment_error_type == "HTTPException" diff --git a/src/backend/tests/unit/services/telemetry/test_telemetry_schema.py b/src/backend/tests/unit/services/telemetry/test_telemetry_schema.py index 7203a45228..fbd3abf5e7 100644 --- a/src/backend/tests/unit/services/telemetry/test_telemetry_schema.py +++ b/src/backend/tests/unit/services/telemetry/test_telemetry_schema.py @@ -8,6 +8,7 @@ import re import pytest from langflow.services.telemetry.schema import ( ComponentPayload, + DeploymentPayload, EmailPayload, PlaygroundPayload, RunPayload, @@ -16,6 +17,86 @@ from langflow.services.telemetry.schema import ( ) +class TestDeploymentPayload: + """Test cases for DeploymentPayload.""" + + def test_deployment_payload_initialization_with_valid_data(self): + """Test DeploymentPayload initialization with valid parameters.""" + payload = DeploymentPayload( + deployment_action="deployment.create", + deployment_provider="test_provider", + deployment_seconds=1.5, + deployment_success=True, + deployment_error_type=None, + client_type="oss", + ) + + assert payload.deployment_action == "deployment.create" + assert payload.deployment_provider == "test_provider" + assert payload.deployment_seconds == 1.5 + assert payload.deployment_success is True + assert payload.deployment_error_type is None + assert payload.client_type == "oss" + + def test_deployment_payload_initialization_with_defaults(self): + """Test DeploymentPayload initialization with default values.""" + payload = DeploymentPayload( + deployment_action="deployment.delete", + deployment_provider="test_provider", + deployment_seconds=0.5, + deployment_success=False, + deployment_error_type="ValueError", + ) + + assert payload.deployment_action == "deployment.delete" + assert payload.deployment_provider == "test_provider" + assert payload.deployment_seconds == 0.5 + assert payload.deployment_success is False + assert payload.deployment_error_type == "ValueError" + assert payload.client_type is None # Default value + + def test_deployment_payload_serialization(self): + """Test DeploymentPayload serialization to dictionary.""" + payload = DeploymentPayload( + deployment_action="deployment.update", + deployment_provider="test_provider", + deployment_seconds=2.0, + deployment_success=True, + deployment_error_type=None, + client_type="desktop", + ) + + data = payload.model_dump(by_alias=True, exclude_none=True) + + assert data["deploymentAction"] == "deployment.update" + assert data["deploymentProvider"] == "test_provider" + assert data["deploymentSeconds"] == 2.0 + assert data["deploymentSuccess"] is True + assert "deploymentErrorType" not in data + assert data["clientType"] == "desktop" + + def test_deployment_payload_roundtrip(self): + """Test DeploymentPayload roundtrip serialization.""" + payload = DeploymentPayload( + deployment_action="provider.create", + deployment_provider="test_provider", + deployment_seconds=1.0, + deployment_success=False, + deployment_error_type="Exception", + client_type="oss", + ) + + data = payload.model_dump() + new_payload = DeploymentPayload(**data) + + assert new_payload.deployment_action == payload.deployment_action + assert new_payload.deployment_provider == payload.deployment_provider + assert new_payload.deployment_seconds == payload.deployment_seconds + assert new_payload.deployment_success == payload.deployment_success + assert new_payload.deployment_error_type == payload.deployment_error_type + assert new_payload.client_type == payload.client_type + + class TestRunPayload: """Test cases for RunPayload.""" diff --git a/src/backend/tests/unit/test_telemetry.py b/src/backend/tests/unit/test_telemetry.py index b43ed64ce9..9b0f197ff4 100644 --- a/src/backend/tests/unit/test_telemetry.py +++ b/src/backend/tests/unit/test_telemetry.py @@ -4,6 +4,83 @@ from concurrent.futures import ThreadPoolExecutor, as_completed import pytest from langflow.services.telemetry.opentelemetry import OpenTelemetry +from langflow.services.telemetry.schema import DeploymentPayload +from langflow.services.telemetry.service import TelemetryService + + +@pytest.fixture +def mock_settings_service(mocker): + settings = mocker.MagicMock() + settings.settings.telemetry_base_url = "http://test.telemetry" + settings.settings.prometheus_enabled = False + settings.settings.do_not_track = False + return settings + + +@pytest.fixture +def telemetry_service(mock_settings_service): + return TelemetryService(mock_settings_service) + + +@pytest.mark.asyncio +async def test_log_package_deployment(telemetry_service): + payload = DeploymentPayload( + deployment_action="deployment.create", + deployment_provider="test_provider", + deployment_seconds=1.0, + deployment_success=True, + ) + await telemetry_service.log_package_deployment(payload) + func, queued_payload, path = await telemetry_service.telemetry_queue.get() + assert func == telemetry_service.send_telemetry_data + assert queued_payload == payload + assert path == "deployment" + + +@pytest.mark.asyncio +async def test_log_package_deployment_provider(telemetry_service): + payload = DeploymentPayload( + deployment_action="provider.create", + deployment_provider="test_provider", + deployment_seconds=1.0, + deployment_success=True, + ) + await telemetry_service.log_package_deployment_provider(payload) + func, queued_payload, path = await telemetry_service.telemetry_queue.get() + assert func == telemetry_service.send_telemetry_data + assert queued_payload == payload + assert path == "deployment_provider" + + +@pytest.mark.asyncio +async def test_log_package_deployment_run(telemetry_service): + payload = DeploymentPayload( + deployment_action="deployment.run", + deployment_provider="test_provider", + deployment_seconds=1.0, + deployment_success=True, + ) + await telemetry_service.log_package_deployment_run(payload) + func, queued_payload, path = await telemetry_service.telemetry_queue.get() + assert func == telemetry_service.send_telemetry_data + assert queued_payload == payload + assert path == "deployment_run" + + +@pytest.mark.asyncio +async def test_log_package_deployment_do_not_track(telemetry_service): + telemetry_service.do_not_track = True + payload = DeploymentPayload( + deployment_action="deployment.create", + deployment_provider="test_provider", + deployment_seconds=1.0, + deployment_success=True, + ) + await telemetry_service.log_package_deployment(payload) + await telemetry_service.log_package_deployment_provider(payload) + await telemetry_service.log_package_deployment_run(payload) + assert telemetry_service.telemetry_queue.empty() + fixed_labels = {"flow_id": "this_flow_id", "service": "this", "user": "that"}