diff --git a/src/backend/base/langflow/api/v1/deployments.py b/src/backend/base/langflow/api/v1/deployments.py index ad056d6233..785c676d28 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,57 @@ 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 into; passed via Depends.""" + + provider: str = "unknown" + wxo_tenant_id: str | None = None + + +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_message: str = "" + try: + yield ctx + except Exception as exc: + success = False + error_message = str(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_message=error_message, + wxo_tenant_id=ctx.wxo_tenant_id, + ) + 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) @@ -266,7 +320,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) @@ -288,6 +344,7 @@ async def create_provider_account( ) except ValueError as exc: _raise_http_for_provider_account_value_error(exc) + telemetry.wxo_tenant_id = provider_account.provider_tenant_id return deployment_mapper.resolve_provider_account_response(provider_account) @@ -337,12 +394,15 @@ 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 + telemetry.wxo_tenant_id = provider_account.provider_tenant_id deployment_count = await _count_provider_deployments_after_reconciliation( session=session, provider_account=provider_account, @@ -370,12 +430,15 @@ 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 + telemetry.wxo_tenant_id = provider_account.provider_tenant_id deployment_mapper = get_deployment_mapper(provider_account.provider_key) verify_input = None @@ -420,6 +483,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 +491,8 @@ async def create_deployment( user_id=current_user.id, db=session, ) + telemetry.provider = provider_account.provider_key + telemetry.wxo_tenant_id = provider_account.provider_tenant_id # fail fast if the deployment name already exists # we could have races but that is more # acceptable than provider-side rollback failure @@ -749,12 +815,21 @@ 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_row, + deployment_adapter, + deployment_mapper, + _provider_key, + provider_tenant_id, + ) = await resolve_adapter_mapper_from_deployment( deployment_id=deployment_id, user_id=current_user.id, db=session, ) + telemetry.provider = _provider_key + telemetry.wxo_tenant_id = provider_tenant_id adapter_execution_payload = await deployment_mapper.resolve_execution_create( deployment_resource_key=deployment_row.resource_key, db=session, @@ -783,7 +858,13 @@ async def get_deployment_run( session: DbSessionReadOnly, current_user: CurrentActiveUser, ): - deployment_row, deployment_adapter, deployment_mapper, _provider_key = await resolve_adapter_mapper_from_deployment( + ( + deployment_row, + deployment_adapter, + deployment_mapper, + _provider_key, + _provider_tenant_id, + ) = await resolve_adapter_mapper_from_deployment( deployment_id=deployment_id, user_id=current_user.id, db=session, @@ -950,6 +1031,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 +1087,8 @@ async def update_snapshot( user_id=current_user.id, db=session, ) + telemetry.provider = provider_account.provider_key + telemetry.wxo_tenant_id = provider_account.provider_tenant_id deployment_adapter = resolve_deployment_adapter(provider_account.provider_key) deployment_mapper = get_deployment_mapper(provider_account.provider_key) @@ -1116,7 +1200,13 @@ async def get_deployment( session: DbSession, current_user: CurrentActiveUser, ): - deployment_row, deployment_adapter, deployment_mapper, provider_key = await resolve_adapter_mapper_from_deployment( + ( + deployment_row, + deployment_adapter, + deployment_mapper, + provider_key, + _provider_tenant_id, + ) = await resolve_adapter_mapper_from_deployment( deployment_id=deployment_id, user_id=current_user.id, db=session, @@ -1232,12 +1322,21 @@ 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_row, + deployment_adapter, + deployment_mapper, + provider_key, + provider_tenant_id, + ) = await resolve_adapter_mapper_from_deployment( deployment_id=deployment_id, user_id=current_user.id, db=session, ) + telemetry.provider = provider_key + telemetry.wxo_tenant_id = provider_tenant_id deployment_row_id = deployment_row.id deployment_resource_key = deployment_row.resource_key deployment_provider_account_id = deployment_row.deployment_provider_account_id @@ -1334,14 +1433,17 @@ async def delete_deployment( deployment_id: DeploymentIdPath, session: DbSession, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(deployment_delete_telemetry)], *, include_provider: IncludeProviderDeleteQuery = True, ): - deployment_row, deployment_adapter, _provider_key = await resolve_adapter_from_deployment( + deployment_row, deployment_adapter, _provider_key, provider_tenant_id = await resolve_adapter_from_deployment( deployment_id=deployment_id, user_id=current_user.id, db=session, ) + telemetry.provider = _provider_key + telemetry.wxo_tenant_id = provider_tenant_id if include_provider: try: with handle_adapter_errors(), deployment_provider_scope(deployment_row.deployment_provider_account_id): @@ -1377,7 +1479,7 @@ async def get_deployment_status( session: DbSessionReadOnly, current_user: CurrentActiveUser, ): - deployment_row, deployment_adapter, provider_key = await resolve_adapter_from_deployment( + deployment_row, deployment_adapter, provider_key, _provider_tenant_id = await resolve_adapter_from_deployment( deployment_id=deployment_id, user_id=current_user.id, db=session, @@ -1422,7 +1524,13 @@ async def list_deployment_flow_versions( ), ] = None, ): - deployment_row, deployment_adapter, deployment_mapper, _provider_key = await resolve_adapter_mapper_from_deployment( + ( + deployment_row, + deployment_adapter, + deployment_mapper, + _provider_key, + _provider_tenant_id, + ) = await resolve_adapter_mapper_from_deployment( deployment_id=deployment_id, user_id=current_user.id, db=session, diff --git a/src/backend/base/langflow/api/v1/mappers/deployments/helpers.py b/src/backend/base/langflow/api/v1/mappers/deployments/helpers.py index 673ac99f5e..77f5fd176e 100644 --- a/src/backend/base/langflow/api/v1/mappers/deployments/helpers.py +++ b/src/backend/base/langflow/api/v1/mappers/deployments/helpers.py @@ -437,8 +437,8 @@ async def resolve_adapter_from_deployment( deployment_id: UUID, user_id: UUID, db: DbSession, -) -> tuple[Deployment, DeploymentServiceProtocol, str]: - """Returns ``(deployment_row, adapter, provider_key)``.""" +) -> tuple[Deployment, DeploymentServiceProtocol, str, str | None]: + """Returns ``(deployment_row, adapter, provider_key, provider_tenant_id)``.""" deployment_row = await get_deployment_row_or_404(deployment_id=deployment_id, user_id=user_id, db=db) provider_account = await get_owned_provider_account_or_404( provider_id=deployment_row.deployment_provider_account_id, @@ -446,7 +446,7 @@ async def resolve_adapter_from_deployment( db=db, ) deployment_adapter = resolve_deployment_adapter(provider_account.provider_key) - return deployment_row, deployment_adapter, provider_account.provider_key + return deployment_row, deployment_adapter, provider_account.provider_key, provider_account.provider_tenant_id async def resolve_adapter_mapper_from_deployment( @@ -454,8 +454,8 @@ async def resolve_adapter_mapper_from_deployment( deployment_id: UUID, user_id: UUID, db: DbSession, -) -> tuple[Deployment, DeploymentServiceProtocol, BaseDeploymentMapper, str]: - """Returns ``(deployment_row, adapter, mapper, provider_key)``.""" +) -> tuple[Deployment, DeploymentServiceProtocol, BaseDeploymentMapper, str, str | None]: + """Returns ``(deployment_row, adapter, mapper, provider_key, provider_tenant_id)``.""" from langflow.api.v1.mappers.deployments.registry import get_deployment_mapper deployment_row = await get_deployment_row_or_404(deployment_id=deployment_id, user_id=user_id, db=db) @@ -466,7 +466,13 @@ async def resolve_adapter_mapper_from_deployment( ) deployment_adapter = resolve_deployment_adapter(provider_account.provider_key) deployment_mapper = get_deployment_mapper(provider_account.provider_key) - return deployment_row, deployment_adapter, deployment_mapper, provider_account.provider_key + return ( + deployment_row, + deployment_adapter, + deployment_mapper, + provider_account.provider_key, + provider_account.provider_tenant_id, + ) async def resolve_project_id_for_deployment_create( diff --git a/src/backend/base/langflow/services/telemetry/schema.py b/src/backend/base/langflow/services/telemetry/schema.py index eecb38694d..b6937138fe 100644 --- a/src/backend/base/langflow/services/telemetry/schema.py +++ b/src/backend/base/langflow/services/telemetry/schema.py @@ -19,6 +19,15 @@ 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_message: str = Field(default="", serialization_alias="deploymentErrorMessage") + wxo_tenant_id: str | None = Field(default=None, serialization_alias="wxoTenantId") + + 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_deployment_route_handlers.py b/src/backend/tests/unit/api/v1/test_deployment_route_handlers.py index 58da55bf33..e46b9ee4a3 100644 --- a/src/backend/tests/unit/api/v1/test_deployment_route_handlers.py +++ b/src/backend/tests/unit/api/v1/test_deployment_route_handlers.py @@ -16,6 +16,7 @@ from uuid import uuid4 import pytest from fastapi import HTTPException +from langflow.api.v1.deployments import DeploymentTelemetryCtx from langflow.api.v1.mappers.deployments.contracts import ProviderSnapshotBinding from langflow.api.v1.schemas.deployments import ( DeploymentLlmListResponse, @@ -58,12 +59,14 @@ def _fake_provider_account( provider_key: str = DeploymentProviderKey.WATSONX_ORCHESTRATE, provider_url: str = "https://api.us-south.wxo.cloud.ibm.com/instances/tenant-1", api_key: str = "encrypted-key", + provider_tenant_id: str | None = "tenant-1", ) -> SimpleNamespace: return SimpleNamespace( id=uuid4(), provider_key=provider_key, provider_url=provider_url, api_key=api_key, + provider_tenant_id=provider_tenant_id, ) @@ -85,6 +88,10 @@ def _fake_user() -> SimpleNamespace: return SimpleNamespace(id=uuid4()) +def _fake_telemetry() -> DeploymentTelemetryCtx: + return DeploymentTelemetryCtx() + + def _fake_attachment(*, provider_snapshot_id: str | None = None) -> SimpleNamespace: return SimpleNamespace( flow_version_id=uuid4(), @@ -182,7 +189,9 @@ class TestCreateDeploymentRollback: payload.description = None with pytest.raises(RuntimeError, match="DB commit failed"): - await create_deployment(session=session, payload=payload, current_user=_fake_user()) + await create_deployment( + session=session, payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) mock_rollback.assert_awaited_once() assert mock_rollback.call_args.kwargs["resource_id"] == create_result.id @@ -240,7 +249,9 @@ class TestCreateDeploymentRollback: payload.description = None mapper.shape_deployment_create_result.return_value = MagicMock() - await create_deployment(session=session, payload=payload, current_user=_fake_user()) + await create_deployment( + session=session, payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) mock_rollback.assert_not_awaited() @@ -301,7 +312,9 @@ class TestCreateDeploymentRollback: payload.description = "desc" with pytest.raises(RuntimeError, match="DB commit failed"): - await create_deployment(session=session, payload=payload, current_user=_fake_user()) + await create_deployment( + session=session, payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) mock_rollback.assert_awaited_once() assert mock_rollback.call_args.kwargs["resource_id"] == "existing-agent-1" @@ -362,7 +375,9 @@ class TestCreateDeploymentExistingAgent: payload.description = None mapper.shape_deployment_create_result.return_value = MagicMock() - await create_deployment(session=session, payload=payload, current_user=_fake_user()) + await create_deployment( + session=session, payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) _ = (mock_name_exists, mock_get_by_resource_key, mock_validate_fv, mock_attach) adapter.create.assert_not_awaited() @@ -423,7 +438,9 @@ class TestCreateDeploymentExistingAgent: payload.description = "desc" mapper.shape_deployment_create_result.return_value = MagicMock() - await create_deployment(session=session, payload=payload, current_user=_fake_user()) + await create_deployment( + session=session, payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) _ = (mock_name_exists, mock_get_by_resource_key, mock_validate_fv, mock_attach) adapter.create.assert_not_awaited() @@ -487,7 +504,9 @@ class TestCreateDeploymentExistingAgent: payload.description = "desc" mapper.shape_deployment_create_result.return_value = MagicMock() - await create_deployment(session=session, payload=payload, current_user=_fake_user()) + await create_deployment( + session=session, payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) _ = (mock_name_exists, mock_get_by_resource_key, mock_validate_fv, mock_attach) adapter.create.assert_not_awaited() @@ -538,7 +557,9 @@ class TestCreateDeploymentExistingAgent: payload.description = None with pytest.raises(HTTPException) as exc_info: - await create_deployment(session=AsyncMock(), payload=payload, current_user=_fake_user()) + await create_deployment( + session=AsyncMock(), payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) assert exc_info.value.status_code == 409 _ = mock_name_exists @@ -1062,6 +1083,7 @@ class TestUpdateSnapshotRoute: body=SnapshotUpdateRequest(flow_version_id=target_flow_version_id), session=session, current_user=user, + telemetry=_fake_telemetry(), ) assert response.flow_version_id == target_flow_version_id @@ -1140,6 +1162,7 @@ class TestUpdateSnapshotRoute: body=SnapshotUpdateRequest(flow_version_id=target_flow_version_id), session=session, current_user=user, + telemetry=_fake_telemetry(), ) session.commit.assert_awaited_once() @@ -1176,7 +1199,7 @@ class TestListDeploymentFlowVersionsRoute: ) ] ) - mock_resolve.return_value = (deployment_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve.return_value = (deployment_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") mock_list_flow_versions_synced.return_value = (rows, 7, snapshot_result) mapper.shape_flow_version_list_result.return_value = SimpleNamespace( page=2, @@ -1257,6 +1280,7 @@ class TestProviderAccountRoutes: session=session, payload=DeploymentProviderAccountUpdateRequest(name="renamed"), current_user=_fake_user(), + telemetry=_fake_telemetry(), ) mapper.resolve_verify_credentials_for_update.assert_not_called() @@ -1298,6 +1322,7 @@ class TestProviderAccountRoutes: session=AsyncMock(), payload=DeploymentProviderAccountUpdateRequest(provider_data={"api_key": "new-api-key"}), current_user=_fake_user(), + telemetry=_fake_telemetry(), ) assert exc_info.value.status_code == 401 @@ -1339,6 +1364,7 @@ class TestProviderAccountRoutes: session=AsyncMock(), payload=payload, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) assert exc_info.value.status_code == 409 @@ -1377,6 +1403,7 @@ class TestProviderAccountRoutes: session=AsyncMock(), payload=DeploymentProviderAccountUpdateRequest(name="prod"), current_user=_fake_user(), + telemetry=_fake_telemetry(), ) assert exc_info.value.status_code == 409 @@ -1421,6 +1448,7 @@ class TestProviderAccountRoutes: session=AsyncMock(), payload=DeploymentProviderAccountUpdateRequest(provider_data={"tenant_id": "tenant-renamed"}), current_user=_fake_user(), + telemetry=_fake_telemetry(), ) @pytest.mark.asyncio @@ -1461,6 +1489,7 @@ class TestProviderAccountRoutes: session=AsyncMock(), payload=DeploymentProviderAccountUpdateRequest(provider_data={"tenant_id": "tenant-renamed"}), current_user=_fake_user(), + telemetry=_fake_telemetry(), ) @pytest.mark.asyncio @@ -1492,6 +1521,7 @@ class TestProviderAccountRoutes: provider_id=existing_account.id, session=AsyncMock(), current_user=_fake_user(), + telemetry=_fake_telemetry(), ) assert exc_info.value.status_code == 409 @@ -1528,6 +1558,7 @@ class TestProviderAccountRoutes: provider_id=existing_account.id, session=session, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) assert response.status_code == 204 @@ -1582,6 +1613,7 @@ class TestProviderAccountRoutes: provider_id=existing_account.id, session=AsyncMock(), current_user=_fake_user(), + telemetry=_fake_telemetry(), ) assert response.status_code == 204 @@ -1629,6 +1661,7 @@ class TestProviderAccountRoutes: provider_id=existing_account.id, session=AsyncMock(), current_user=_fake_user(), + telemetry=_fake_telemetry(), ) assert exc_info.value.status_code == 409 @@ -1663,6 +1696,7 @@ class TestProviderAccountRoutes: provider_id=existing_account.id, session=session, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) assert response.status_code == 204 @@ -1708,7 +1742,7 @@ class TestUpdateDeploymentRollback: adapter.update.return_value = update_result mapper.resolve_deployment_update = AsyncMock(return_value=MagicMock()) mapper.shape_deployment_update_result.return_value = MagicMock() - mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") session = AsyncMock() session.commit.side_effect = RuntimeError("DB commit failed") @@ -1723,6 +1757,7 @@ class TestUpdateDeploymentRollback: session=session, payload=payload, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) mock_rollback.assert_awaited_once() @@ -1760,7 +1795,7 @@ class TestUpdateDeploymentRollback: adapter.update.return_value = update_result mapper.resolve_deployment_update = AsyncMock(return_value=MagicMock()) mapper.shape_deployment_update_result.return_value = MagicMock() - mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") session = AsyncMock() session.commit.return_value = None @@ -1774,6 +1809,7 @@ class TestUpdateDeploymentRollback: session=session, payload=payload, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) mock_rollback.assert_not_awaited() @@ -1816,7 +1852,7 @@ class TestUpdateDeploymentAlreadyAttachedFiltering: adapter.update.return_value = update_result mapper.resolve_deployment_update = AsyncMock(return_value=MagicMock()) mapper.shape_deployment_update_result.return_value = MagicMock() - mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") reused_fv_id = uuid4() new_fv_id = uuid4() @@ -1838,6 +1874,7 @@ class TestUpdateDeploymentAlreadyAttachedFiltering: session=session, payload=payload, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) mock_resolve_snap.assert_called_once() @@ -1874,7 +1911,7 @@ class TestUpdateDeploymentAlreadyAttachedFiltering: adapter.update.return_value = update_result mapper.resolve_deployment_update = AsyncMock(return_value=MagicMock()) mapper.shape_deployment_update_result.return_value = MagicMock() - mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") fv_id_1 = uuid4() fv_id_2 = uuid4() @@ -1897,6 +1934,7 @@ class TestUpdateDeploymentAlreadyAttachedFiltering: session=session, payload=payload, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) resolved_fv_ids = mock_resolve_snap.call_args.kwargs["added_flow_version_ids"] @@ -1928,7 +1966,7 @@ class TestUpdateDeploymentAlreadyAttachedFiltering: adapter.update.return_value = update_result mapper.resolve_deployment_update = AsyncMock(return_value=MagicMock()) mapper.shape_deployment_update_result.return_value = MagicMock() - mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") fv_id_1 = uuid4() fv_id_2 = uuid4() @@ -1948,6 +1986,7 @@ class TestUpdateDeploymentAlreadyAttachedFiltering: session=session, payload=payload, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) resolved_fv_ids = mock_resolve_snap.call_args.kwargs["added_flow_version_ids"] @@ -1989,7 +2028,7 @@ class TestUpdateDeploymentMetadataPersistence: adapter.update.return_value = update_result mapper.resolve_deployment_update = AsyncMock(return_value=MagicMock()) mapper.shape_deployment_update_result.return_value = MagicMock() - mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") mock_update_db.return_value = updated_row session = AsyncMock() @@ -2007,6 +2046,7 @@ class TestUpdateDeploymentMetadataPersistence: session=session, payload=payload, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) mock_update_db.assert_awaited_once() @@ -2039,7 +2079,7 @@ class TestGetDeploymentSync: dep_row = _fake_deployment_row() adapter = AsyncMock() adapter.get.side_effect = DeploymentNotFoundError(message="gone") - mock_resolve.return_value = (dep_row, adapter, BaseDeploymentMapper(), "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, BaseDeploymentMapper(), "watsonx-orchestrate", "tenant-1") user = _fake_user() session = AsyncMock() @@ -2066,7 +2106,7 @@ class TestGetDeploymentSync: dep_row = _fake_deployment_row() adapter = AsyncMock() adapter.get.side_effect = AuthenticationError(message="bad creds", error_code="authentication_error") - mock_resolve.return_value = (dep_row, adapter, BaseDeploymentMapper(), "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, BaseDeploymentMapper(), "watsonx-orchestrate", "tenant-1") session = AsyncMock() @@ -2091,7 +2131,7 @@ class TestGetDeploymentSync: dep_row = _fake_deployment_row() adapter = AsyncMock() adapter.get.side_effect = ServiceUnavailableError(message="provider down") - mock_resolve.return_value = (dep_row, adapter, BaseDeploymentMapper(), "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, BaseDeploymentMapper(), "watsonx-orchestrate", "tenant-1") session = AsyncMock() @@ -2144,7 +2184,7 @@ class TestGetDeploymentSync: } } adapter.get.return_value = provider_deployment - mock_resolve.return_value = (dep_row, adapter, _MapperForGet(), "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, _MapperForGet(), "watsonx-orchestrate", "tenant-1") session = AsyncMock() result = await get_deployment(deployment_id=dep_row.id, session=session, current_user=_fake_user()) @@ -2186,7 +2226,7 @@ class TestGetDeploymentSync: ProviderSnapshotBinding(resource_key=dep_row.resource_key, snapshot_id="snap-1") ] mapper.shape_deployment_get_data.return_value = None - mock_resolve.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") att_good = _fake_attachment(provider_snapshot_id="snap-1") mock_list_att.return_value = [att_good] @@ -2224,7 +2264,7 @@ class TestGetDeploymentSync: adapter.get.return_value = provider_deployment mapper.extract_snapshot_bindings_for_get.side_effect = NotImplementedError("not supported") mapper.shape_deployment_get_data.return_value = None - mock_resolve.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") mock_list_att.return_value = [_fake_attachment(provider_snapshot_id=None)] session = AsyncMock() @@ -2262,7 +2302,7 @@ class TestGetDeploymentSync: ProviderSnapshotBinding(resource_key=dep_row.resource_key, snapshot_id="snap-1") ] mapper.shape_deployment_get_data.return_value = None - mock_resolve.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") att1 = _fake_attachment(provider_snapshot_id="snap-1") att2 = _fake_attachment(provider_snapshot_id="snap-2") @@ -2299,7 +2339,7 @@ class TestGetDeploymentSync: ProviderSnapshotBinding(resource_key="agent-rk-1", snapshot_id="snap-1") ] mapper.shape_deployment_get_data.return_value = None - mock_resolve.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") mock_list_att.return_value = [_fake_attachment(provider_snapshot_id="snap-1")] session = AsyncMock() @@ -2336,12 +2376,14 @@ class TestDeleteDeployment: dep_row = _fake_deployment_row() adapter = AsyncMock() adapter.delete.side_effect = DeploymentNotFoundError(message="gone") - mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate", "tenant-1") user = _fake_user() session = AsyncMock() - response = await delete_deployment(deployment_id=dep_row.id, session=session, current_user=user) + response = await delete_deployment( + deployment_id=dep_row.id, session=session, current_user=user, telemetry=_fake_telemetry() + ) assert response.status_code == 204 mock_delete_row.assert_awaited_once_with(session, user_id=user.id, deployment_id=dep_row.id) @@ -2362,12 +2404,14 @@ class TestDeleteDeployment: dep_row = _fake_deployment_row() adapter = AsyncMock() adapter.delete.side_effect = AuthenticationError(message="bad creds", error_code="authentication_error") - mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate", "tenant-1") session = AsyncMock() with pytest.raises(HTTPException) as exc_info: - await delete_deployment(deployment_id=dep_row.id, session=session, current_user=_fake_user()) + await delete_deployment( + deployment_id=dep_row.id, session=session, current_user=_fake_user(), telemetry=_fake_telemetry() + ) assert exc_info.value.status_code == 401 mock_delete_row.assert_not_awaited() @@ -2386,13 +2430,15 @@ class TestDeleteDeployment: dep_row = _fake_deployment_row() adapter = AsyncMock() - mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate", "tenant-1") user = _fake_user() session = AsyncMock() session.commit.side_effect = [RuntimeError("commit failed"), None] - response = await delete_deployment(deployment_id=dep_row.id, session=session, current_user=user) + response = await delete_deployment( + deployment_id=dep_row.id, session=session, current_user=user, telemetry=_fake_telemetry() + ) assert response.status_code == 204 assert mock_delete_row.await_count == 2 @@ -2412,13 +2458,15 @@ class TestDeleteDeployment: dep_row = _fake_deployment_row() adapter = AsyncMock() - mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate", "tenant-1") session = AsyncMock() session.commit.side_effect = [RuntimeError("commit failed"), RuntimeError("still failing")] with pytest.raises(HTTPException) as exc_info: - await delete_deployment(deployment_id=dep_row.id, session=session, current_user=_fake_user()) + await delete_deployment( + deployment_id=dep_row.id, session=session, current_user=_fake_user(), telemetry=_fake_telemetry() + ) assert exc_info.value.status_code == 500 assert mock_delete_row.await_count == 2 @@ -2438,13 +2486,17 @@ class TestDeleteDeployment: dep_row = _fake_deployment_row() adapter = AsyncMock() - mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate", "tenant-1") user = _fake_user() session = AsyncMock() response = await delete_deployment( - deployment_id=dep_row.id, session=session, current_user=user, include_provider=True + deployment_id=dep_row.id, + session=session, + current_user=user, + include_provider=True, + telemetry=_fake_telemetry(), ) assert response.status_code == 204 @@ -2464,13 +2516,17 @@ class TestDeleteDeployment: dep_row = _fake_deployment_row() adapter = AsyncMock() - mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate") + mock_resolve.return_value = (dep_row, adapter, "watsonx-orchestrate", "tenant-1") user = _fake_user() session = AsyncMock() response = await delete_deployment( - deployment_id=dep_row.id, session=session, current_user=user, include_provider=False + deployment_id=dep_row.id, + session=session, + current_user=user, + include_provider=False, + telemetry=_fake_telemetry(), ) assert response.status_code == 204 @@ -2503,7 +2559,9 @@ class TestCreateDeploymentDuplicateName: payload.name = "duplicate-name" with pytest.raises(HTTPException) as exc_info: - await create_deployment(session=AsyncMock(), payload=payload, current_user=_fake_user()) + await create_deployment( + session=AsyncMock(), payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) assert exc_info.value.status_code == 409 assert "duplicate-name" in exc_info.value.detail @@ -2529,7 +2587,9 @@ class TestCreateDeploymentDuplicateName: payload.name = "taken" with pytest.raises(HTTPException): - await create_deployment(session=AsyncMock(), payload=payload, current_user=_fake_user()) + await create_deployment( + session=AsyncMock(), payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) mock_resolve_adapter.assert_not_called() @@ -2577,7 +2637,9 @@ class TestCreateDeploymentProjectValidation: payload.provider_id = pa.id with pytest.raises(HTTPException) as exc_info: - await create_deployment(session=AsyncMock(), payload=payload, current_user=_fake_user()) + await create_deployment( + session=AsyncMock(), payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) assert exc_info.value.status_code == 404 mock_validate_fv.assert_awaited_once() @@ -2632,7 +2694,9 @@ class TestCreateDeploymentProjectValidation: patch(f"{ROUTES_MODULE}.attach_flow_versions", new_callable=AsyncMock), ): mapper.shape_deployment_create_result.return_value = MagicMock() - await create_deployment(session=session, payload=payload, current_user=_fake_user()) + await create_deployment( + session=session, payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) mock_validate_fv.assert_awaited_once() assert mock_validate_fv.call_args.kwargs["flow_version_ids"] == [] @@ -2678,7 +2742,9 @@ class TestCreateDeploymentSchemaValidation: payload.provider_id = pa.id with pytest.raises(HTTPException) as exc_info: - await create_deployment(session=AsyncMock(), payload=payload, current_user=_fake_user()) + await create_deployment( + session=AsyncMock(), payload=payload, current_user=_fake_user(), telemetry=_fake_telemetry() + ) assert exc_info.value.status_code == 422 mock_validate_fv.assert_awaited_once() @@ -2708,7 +2774,7 @@ class TestUpdateDeploymentProjectValidation: adapter = AsyncMock() mapper = MagicMock() mapper.resolve_deployment_update = AsyncMock(return_value=MagicMock()) - mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate") + mock_resolve_amm.return_value = (dep_row, adapter, mapper, "watsonx-orchestrate", "tenant-1") add_ids = [uuid4()] remove_ids = [uuid4()] @@ -2724,6 +2790,7 @@ class TestUpdateDeploymentProjectValidation: session=AsyncMock(), payload=payload, current_user=_fake_user(), + telemetry=_fake_telemetry(), ) assert exc_info.value.status_code == 404 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..b74f4e168b --- /dev/null +++ b/src/backend/tests/unit/api/v1/test_deployments_telemetry.py @@ -0,0 +1,325 @@ +# 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", provider_tenant_id="tenant-test" + ) + mock_get_owned.return_value = AsyncMock( + id=uuid4(), provider_key="watsonx-orchestrate", provider_tenant_id="tenant-test" + ) + 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", + "tenant-test", + ) + 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", + "tenant-test", + ) + 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 + assert payload.wxo_tenant_id == "tenant-test" + + +@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 + assert payload.wxo_tenant_id == "tenant-test" + + +@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 + assert payload.wxo_tenant_id == "tenant-test" + + +@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 + assert payload.wxo_tenant_id == "tenant-test" + + +@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 + assert payload.wxo_tenant_id == "tenant-test" + + +@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 + assert payload.wxo_tenant_id == "tenant-test" + + +@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 + assert payload.wxo_tenant_id == "tenant-test" + + +@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 + assert payload.wxo_tenant_id == "tenant-test" + + +@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_message == "400: Invalid credentials" + # verify_credentials raises before the row is created, so the tenant_id is not yet set. + assert payload.wxo_tenant_id is None + + +@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_message == ( + "500: An unexpected error occurred while communicating with the deployment provider." + ) + # tenant_id is captured before the adapter call that raises. + assert payload.wxo_tenant_id == "tenant-test" diff --git a/src/backend/tests/unit/api/v1/test_endpoints.py b/src/backend/tests/unit/api/v1/test_endpoints.py index c3fb0fedbf..e2b6591dbd 100644 --- a/src/backend/tests/unit/api/v1/test_endpoints.py +++ b/src/backend/tests/unit/api/v1/test_endpoints.py @@ -385,6 +385,6 @@ async def test_deprecated_upload_enforces_max_file_size( headers=logged_in_headers, ) - assert response.status_code == status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, ( + assert response.status_code == status.HTTP_413_CONTENT_TOO_LARGE, ( f"Expected 413 for oversized upload, got {response.status_code}: {response.text}" ) 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..9aec2a000b 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,93 @@ 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_message="", + wxo_tenant_id="tenant-abc", + 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_message == "" + assert payload.wxo_tenant_id == "tenant-abc" + 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_message="404: Provider not found.", + ) + + 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_message == "404: Provider not found." + assert payload.wxo_tenant_id is None # Default value + 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_message="", + wxo_tenant_id="tenant-xyz", + 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 data["deploymentErrorMessage"] == "" + assert data["wxoTenantId"] == "tenant-xyz" + 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_message="500: Adapter exploded.", + wxo_tenant_id="tenant-abc", + 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_message == payload.deployment_error_message + assert new_payload.wxo_tenant_id == payload.wxo_tenant_id + 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"} diff --git a/src/frontend/package-lock.json b/src/frontend/package-lock.json index 69eef6fbe2..4089ca3f79 100644 --- a/src/frontend/package-lock.json +++ b/src/frontend/package-lock.json @@ -8281,6 +8281,7 @@ "integrity": "sha512-riJjyv1/mHLIPX4RwiK+oW9/4c3TEUeORHKefKAKnZ5kyslbN+HXowtbaVEqt4IMUB7OXlfixcs6gsFeo/jhiQ==", "dev": true, "license": "Apache-2.0", + "peer": true, "peerDependencies": { "bare-abort-controller": "*" }, diff --git a/src/lfx/src/lfx/services/adapters/deployment/exceptions.py b/src/lfx/src/lfx/services/adapters/deployment/exceptions.py index 36f19c2c9c..3cb090671b 100644 --- a/src/lfx/src/lfx/services/adapters/deployment/exceptions.py +++ b/src/lfx/src/lfx/services/adapters/deployment/exceptions.py @@ -307,7 +307,7 @@ def raise_as_deployment_error( raise InvalidDeploymentOperationError(message=message, cause=cause) from cause if status_code == status.HTTP_405_METHOD_NOT_ALLOWED: raise InvalidDeploymentOperationError(message=message, cause=cause) from cause - if status_code == status.HTTP_413_REQUEST_ENTITY_TOO_LARGE: + if status_code == status.HTTP_413_CONTENT_TOO_LARGE: raise InvalidContentError(message=message, cause=cause) from cause if status_code == status.HTTP_415_UNSUPPORTED_MEDIA_TYPE: raise InvalidContentError(message=message, cause=cause) from cause