mirror of
https://github.com/langflow-ai/langflow.git
synced 2026-07-25 17:33:09 +08:00
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.
This commit is contained in:
@ -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):
|
||||
|
||||
@ -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")
|
||||
|
||||
|
||||
@ -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)
|
||||
|
||||
305
src/backend/tests/unit/api/v1/test_deployments_telemetry.py
Normal file
305
src/backend/tests/unit/api/v1/test_deployments_telemetry.py
Normal file
@ -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"
|
||||
@ -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."""
|
||||
|
||||
|
||||
@ -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"}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user