diff --git a/src/backend/base/langflow/api/v1/deployments.py b/src/backend/base/langflow/api/v1/deployments.py index d6749481a7..26d1abfb86 100644 --- a/src/backend/base/langflow/api/v1/deployments.py +++ b/src/backend/base/langflow/api/v1/deployments.py @@ -2,7 +2,6 @@ from __future__ import annotations import time from collections.abc import AsyncIterator -from contextlib import asynccontextmanager from dataclasses import dataclass from typing import Annotated from uuid import UUID @@ -118,40 +117,48 @@ from langflow.services.telemetry.schema import DeploymentPayload @dataclass class DeploymentTelemetryCtx: - """Handler-local context for deployment telemetry; handlers set `provider` once known.""" + """Mutable context that routes write `provider` into; passed via Depends.""" provider: str = "unknown" -@asynccontextmanager -async def _track_deployment_telemetry(action: str) -> AsyncIterator[DeploymentTelemetryCtx]: - # Emits telemetry inline via ``await`` (like endpoints.py:282) rather than - # BackgroundTasks because BackgroundTasks are dropped when the handler - # raises, and we still need to record failure telemetry. The queue put - # is effectively instant so latency is not a concern. - ctx = DeploymentTelemetryCtx() - start = time.perf_counter() - success = True - error_message = "" - try: - yield ctx - except Exception as exc: - success = False - error_message = str(exc) - raise - finally: +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: - await get_telemetry_service().log_package_deployment( - DeploymentPayload( + 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=int(time.perf_counter() - start), + deployment_seconds=time.perf_counter() - started_at, deployment_success=success, - deployment_error_message=error_message, + deployment_error_type=type(error).__name__ if error else None, ) - ) - except Exception: # noqa: BLE001 - logger.debug("deployment telemetry emit failed", exc_info=True) + 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) @@ -311,31 +318,31 @@ async def create_provider_account( session: DbSession, payload: DeploymentProviderAccountCreateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(provider_create_telemetry)], ): - async with _track_deployment_telemetry("provider.create") as telemetry: - telemetry.provider = payload.provider_key - deployment_mapper = get_deployment_mapper(payload.provider_key) - deployment_adapter = resolve_deployment_adapter(payload.provider_key) + telemetry.provider = payload.provider_key + deployment_mapper = get_deployment_mapper(payload.provider_key) + deployment_adapter = resolve_deployment_adapter(payload.provider_key) - with handle_adapter_errors(mapper=deployment_mapper): - verify_input = deployment_mapper.resolve_verify_credentials_for_create(payload=payload) - await deployment_adapter.verify_credentials( - user_id=current_user.id, - payload=verify_input, - ) + with handle_adapter_errors(mapper=deployment_mapper): + verify_input = deployment_mapper.resolve_verify_credentials_for_create(payload=payload) + await deployment_adapter.verify_credentials( + user_id=current_user.id, + payload=verify_input, + ) - try: - provider_account_to_create = deployment_mapper.resolve_provider_account_create( - payload=payload, - user_id=current_user.id, - ) - provider_account = await create_provider_account_row( - session, - provider_account=provider_account_to_create, - ) - except ValueError as exc: - _raise_http_for_provider_account_value_error(exc) - return deployment_mapper.resolve_provider_account_response(provider_account) + try: + provider_account_to_create = deployment_mapper.resolve_provider_account_create( + payload=payload, + user_id=current_user.id, + ) + provider_account = await create_provider_account_row( + session, + provider_account=provider_account_to_create, + ) + except ValueError as exc: + _raise_http_for_provider_account_value_error(exc) + return deployment_mapper.resolve_provider_account_response(provider_account) @router.get("/providers", response_model=DeploymentProviderAccountListResponse, tags=["Deployment Providers"]) @@ -384,29 +391,29 @@ async def delete_provider_account( provider_id: DeploymentProviderAccountIdPath, session: DbSession, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(provider_delete_telemetry)], ): - async with _track_deployment_telemetry("provider.delete") as telemetry: - provider_account = await get_owned_provider_account_or_404( - provider_id=provider_id, - user_id=current_user.id, - db=session, + 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, + user_id=current_user.id, + ) + if deployment_count > 0: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail="Cannot delete provider account while deployments still exist.", ) - telemetry.provider = provider_account.provider_key - deployment_count = await _count_provider_deployments_after_reconciliation( - session=session, - provider_account=provider_account, - user_id=current_user.id, - ) - if deployment_count > 0: - raise HTTPException( - status_code=status.HTTP_409_CONFLICT, - detail="Cannot delete provider account while deployments still exist.", - ) - try: - await delete_provider_account_row(session, provider_account=provider_account) - except ValueError as exc: - _raise_http_for_provider_account_value_error(exc) - return Response(status_code=status.HTTP_204_NO_CONTENT) + try: + await delete_provider_account_row(session, provider_account=provider_account) + except ValueError as exc: + _raise_http_for_provider_account_value_error(exc) + return Response(status_code=status.HTTP_204_NO_CONTENT) @router.patch( @@ -419,51 +426,51 @@ async def update_provider_account( session: DbSession, payload: DeploymentProviderAccountUpdateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(provider_update_telemetry)], ): - async with _track_deployment_telemetry("provider.update") as 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 - if _field_was_explicitly_set(payload, "provider_data"): - deployment_adapter = resolve_deployment_adapter(provider_account.provider_key) - try: - verify_input = deployment_mapper.resolve_verify_credentials_for_update( - payload=payload, - existing_account=provider_account, - ) - except ValueError as exc: - _raise_http_for_provider_account_value_error(exc) - except NotImplementedError as exc: - raise HTTPException( - status_code=status.HTTP_501_NOT_IMPLEMENTED, - detail="This operation is not supported by the deployment provider.", - ) from exc - if verify_input is not None: - with handle_adapter_errors(mapper=deployment_mapper): - await deployment_adapter.verify_credentials( - user_id=current_user.id, - payload=verify_input, - ) + 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 + if _field_was_explicitly_set(payload, "provider_data"): + deployment_adapter = resolve_deployment_adapter(provider_account.provider_key) try: - update_kwargs = deployment_mapper.resolve_provider_account_update( + verify_input = deployment_mapper.resolve_verify_credentials_for_update( payload=payload, existing_account=provider_account, ) - updated = await update_provider_account_row( - session, - provider_account=provider_account, - **update_kwargs, - ) except ValueError as exc: _raise_http_for_provider_account_value_error(exc) - return deployment_mapper.resolve_provider_account_response(updated) + except NotImplementedError as exc: + raise HTTPException( + status_code=status.HTTP_501_NOT_IMPLEMENTED, + detail="This operation is not supported by the deployment provider.", + ) from exc + if verify_input is not None: + with handle_adapter_errors(mapper=deployment_mapper): + await deployment_adapter.verify_credentials( + user_id=current_user.id, + payload=verify_input, + ) + + try: + update_kwargs = deployment_mapper.resolve_provider_account_update( + payload=payload, + existing_account=provider_account, + ) + updated = await update_provider_account_row( + session, + provider_account=provider_account, + **update_kwargs, + ) + except ValueError as exc: + _raise_http_for_provider_account_value_error(exc) + return deployment_mapper.resolve_provider_account_response(updated) @router.post("", response_model=DeploymentCreateResponse, status_code=status.HTTP_201_CREATED) @@ -471,163 +478,161 @@ async def create_deployment( session: DbSession, payload: DeploymentCreateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(deployment_create_telemetry)], ): - async with _track_deployment_telemetry("deployment.create") as telemetry: - provider_id = payload.provider_id - provider_account = await get_owned_provider_account_or_404( - provider_id=provider_id, - user_id=current_user.id, - db=session, + provider_id = payload.provider_id + 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 + # fail fast if the deployment name already exists + # we could have races but that is more + # acceptable than provider-side rollback failure + if await deployment_name_exists( + session, + user_id=current_user.id, + deployment_provider_account_id=provider_id, + name=payload.name, + ): + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail=f"A deployment named '{payload.name}' already exists. " + "Please choose a different name or delete the existing deployment first.", ) - 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 - if await deployment_name_exists( + + deployment_adapter = resolve_deployment_adapter(provider_account.provider_key) + deployment_mapper = get_deployment_mapper(provider_account.provider_key) + existing_resource_key = deployment_mapper.util_existing_deployment_resource_key_for_create(payload) + if existing_resource_key is not None: + existing_deployment = await get_deployment_by_resource_key( session, user_id=current_user.id, deployment_provider_account_id=provider_id, - name=payload.name, - ): + resource_key=str(existing_resource_key), + ) + if existing_deployment is not None: raise HTTPException( status_code=status.HTTP_409_CONFLICT, - detail=f"A deployment named '{payload.name}' already exists. " - "Please choose a different name or delete the existing deployment first.", + detail=f"The agent '{existing_resource_key}' is already managed by Langflow. " + "Update it to make changes, or delete the existing deployment first.", ) - - deployment_adapter = resolve_deployment_adapter(provider_account.provider_key) - deployment_mapper = get_deployment_mapper(provider_account.provider_key) - existing_resource_key = deployment_mapper.util_existing_deployment_resource_key_for_create(payload) - if existing_resource_key is not None: - existing_deployment = await get_deployment_by_resource_key( - session, - user_id=current_user.id, - deployment_provider_account_id=provider_id, - resource_key=str(existing_resource_key), - ) - if existing_deployment is not None: - raise HTTPException( - status_code=status.HTTP_409_CONFLICT, - detail=f"The agent '{existing_resource_key}' is already managed by Langflow. " - "Update it to make changes, or delete the existing deployment first.", - ) - should_mutate_existing_resource = ( - existing_resource_key is not None - and deployment_mapper.util_should_mutate_provider_for_existing_deployment_create(payload) - ) - should_create_provider_resource = existing_resource_key is None - project_id = await resolve_project_id_for_deployment_create( - payload=payload, user_id=current_user.id, db=session - ) - flow_version_ids = deployment_mapper.util_create_flow_version_ids(payload) - await validate_project_scoped_flow_version_ids( - flow_version_ids=flow_version_ids, + should_mutate_existing_resource = ( + existing_resource_key is not None + and deployment_mapper.util_should_mutate_provider_for_existing_deployment_create(payload) + ) + should_create_provider_resource = existing_resource_key is None + project_id = await resolve_project_id_for_deployment_create(payload=payload, user_id=current_user.id, db=session) + flow_version_ids = deployment_mapper.util_create_flow_version_ids(payload) + await validate_project_scoped_flow_version_ids( + flow_version_ids=flow_version_ids, + user_id=current_user.id, + project_id=project_id, + db=session, + ) + if should_create_provider_resource: + adapter_payload = await deployment_mapper.resolve_deployment_create( user_id=current_user.id, project_id=project_id, db=session, + payload=payload, ) - if should_create_provider_resource: - adapter_payload = await deployment_mapper.resolve_deployment_create( + with handle_adapter_errors(mapper=deployment_mapper), deployment_provider_scope(provider_id): + provider_create_result = await deployment_adapter.create( + user_id=current_user.id, + payload=adapter_payload, + db=session, + ) + else: + # Existing-resource create starts as DB-only onboarding: no provider + # mutation is performed and created_* response fields stay empty. + provider_create_result = deployment_mapper.util_create_result_from_existing_resource( + existing_resource_key=str(existing_resource_key), + ) + if should_mutate_existing_resource: + # When create payload includes add_flows/upsert_tools, run provider + # update and normalize the update result into create-style created_*. + adapter_payload = await deployment_mapper.resolve_deployment_update_for_existing_create( user_id=current_user.id, project_id=project_id, db=session, payload=payload, ) with handle_adapter_errors(mapper=deployment_mapper), deployment_provider_scope(provider_id): - provider_create_result = await deployment_adapter.create( - user_id=current_user.id, + provider_update_result: DeploymentUpdateResult = await deployment_adapter.update( + deployment_id=existing_resource_key, payload=adapter_payload, - db=session, - ) - else: - # Existing-resource create starts as DB-only onboarding: no provider - # mutation is performed and created_* response fields stay empty. - provider_create_result = deployment_mapper.util_create_result_from_existing_resource( - existing_resource_key=str(existing_resource_key), - ) - if should_mutate_existing_resource: - # When create payload includes add_flows/upsert_tools, run provider - # update and normalize the update result into create-style created_*. - adapter_payload = await deployment_mapper.resolve_deployment_update_for_existing_create( user_id=current_user.id, - project_id=project_id, db=session, - payload=payload, ) - with handle_adapter_errors(mapper=deployment_mapper), deployment_provider_scope(provider_id): - provider_update_result: DeploymentUpdateResult = await deployment_adapter.update( - deployment_id=existing_resource_key, - payload=adapter_payload, - user_id=current_user.id, - db=session, - ) - provider_create_result = deployment_mapper.util_create_result_from_existing_update( - existing_resource_key=str(existing_resource_key), - result=provider_update_result, - ) - # if we get here, the deployment was created successfully in the provider - # so we need to create the deployment row and attach the flow versions - # in the DB - try: - deployment_row = await create_deployment_db( - session, - user_id=current_user.id, - project_id=project_id, - deployment_provider_account_id=provider_id, - resource_key=str(provider_create_result.id), - name=payload.name, - deployment_type=payload.type, - description=payload.description or None, + provider_create_result = deployment_mapper.util_create_result_from_existing_update( + existing_resource_key=str(existing_resource_key), + result=provider_update_result, ) + # if we get here, the deployment was created successfully in the provider + # so we need to create the deployment row and attach the flow versions + # in the DB + try: + deployment_row = await create_deployment_db( + session, + user_id=current_user.id, + project_id=project_id, + deployment_provider_account_id=provider_id, + resource_key=str(provider_create_result.id), + name=payload.name, + deployment_type=payload.type, + description=payload.description or None, + ) - snapshot_id_by_flow_version_id: dict[UUID, str] = {} - if flow_version_ids: - snapshot_id_by_flow_version_id = resolve_snapshot_map_for_create( - deployment_mapper=deployment_mapper, - result=provider_create_result, - flow_version_ids=flow_version_ids, - ) - await attach_flow_versions( + snapshot_id_by_flow_version_id: dict[UUID, str] = {} + if flow_version_ids: + snapshot_id_by_flow_version_id = resolve_snapshot_map_for_create( + deployment_mapper=deployment_mapper, + result=provider_create_result, flow_version_ids=flow_version_ids, + ) + await attach_flow_versions( + flow_version_ids=flow_version_ids, + user_id=current_user.id, + deployment_row_id=deployment_row.id, + snapshot_id_by_flow_version_id=snapshot_id_by_flow_version_id, + db=session, + ) + + await session.commit() + except Exception as exc: + # Compensate: delete the provider resource so it doesn't become orphaned. + # Only the deployment resource itself is deleted (e.g. the WXO agent). + # Secondary resources (snapshots/tools, configs) may remain orphaned -- + # this is intentional because snapshots/configs may be shared across deployments, + # making cascade-delete unsafe. + await session.rollback() + if should_create_provider_resource: + await rollback_provider_create( + deployment_adapter=deployment_adapter, + provider_id=provider_id, + resource_id=provider_create_result.id, + provider_result=provider_create_result.provider_result, user_id=current_user.id, - deployment_row_id=deployment_row.id, - snapshot_id_by_flow_version_id=snapshot_id_by_flow_version_id, db=session, ) - - await session.commit() - except Exception as exc: - # Compensate: delete the provider resource so it doesn't become orphaned. - # Only the deployment resource itself is deleted (e.g. the WXO agent). - # Secondary resources (snapshots/tools, configs) may remain orphaned -- - # this is intentional because snapshots/configs may be shared across deployments, - # making cascade-delete unsafe. - await session.rollback() - if should_create_provider_resource: - await rollback_provider_create( - deployment_adapter=deployment_adapter, - provider_id=provider_id, - resource_id=provider_create_result.id, - provider_result=provider_create_result.provider_result, - user_id=current_user.id, - db=session, - ) - elif should_mutate_existing_resource: - await rollback_provider_create( - deployment_adapter=deployment_adapter, - provider_id=provider_id, - resource_id=str(existing_resource_key), - provider_result=provider_create_result.provider_result, - allow_delete_fallback=False, - user_id=current_user.id, - db=session, - ) - if isinstance(exc, AttachmentConflictError): - raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(exc)) from exc - raise - return deployment_mapper.shape_deployment_create_result( - provider_create_result, deployment_row, provider_key=provider_account.provider_key - ) + elif should_mutate_existing_resource: + await rollback_provider_create( + deployment_adapter=deployment_adapter, + provider_id=provider_id, + resource_id=str(existing_resource_key), + provider_result=provider_create_result.provider_result, + allow_delete_fallback=False, + user_id=current_user.id, + db=session, + ) + if isinstance(exc, AttachmentConflictError): + raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(exc)) from exc + raise + return deployment_mapper.shape_deployment_create_result( + provider_create_result, deployment_row, provider_key=provider_account.provider_key + ) # exclude none as its associated with @@ -804,38 +809,33 @@ async def create_deployment_run( session: DbSession, payload: RunCreateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(deployment_run_telemetry)], ): - async with _track_deployment_telemetry("deployment.run") as telemetry: - ( - deployment_row, - deployment_adapter, - deployment_mapper, - _provider_key, - ) = await resolve_adapter_mapper_from_deployment( - deployment_id=deployment_id, + 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, + payload=payload, + ) + with ( + handle_adapter_errors(mapper=deployment_mapper), + deployment_provider_scope(deployment_row.deployment_provider_account_id), + ): + execution_result = await deployment_adapter.create_execution( + payload=adapter_execution_payload, 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, - payload=payload, - ) - with ( - handle_adapter_errors(mapper=deployment_mapper), - deployment_provider_scope(deployment_row.deployment_provider_account_id), - ): - execution_result = await deployment_adapter.create_execution( - payload=adapter_execution_payload, - user_id=current_user.id, - db=session, - ) - return deployment_mapper.shape_execution_create_result( - execution_result, - deployment_id=deployment_row.id, - ) + return deployment_mapper.shape_execution_create_result( + execution_result, + deployment_id=deployment_row.id, + ) @router.get("/{deployment_id}/runs/{run_id}", response_model=RunStatusResponse) @@ -1012,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. @@ -1022,151 +1023,150 @@ async def update_snapshot( from langflow.services.database.models.deployment.crud import get_deployment as get_deployment_row from langflow.services.database.models.flow_version.crud import get_flow_version_entry - async with _track_deployment_telemetry("snapshot.update") as telemetry: - snapshot_id = provider_snapshot_id.strip() + snapshot_id = provider_snapshot_id.strip() - attachment = await get_attachment_by_provider_snapshot_id( - session, - user_id=current_user.id, - provider_snapshot_id=snapshot_id, + attachment = await get_attachment_by_provider_snapshot_id( + session, + user_id=current_user.id, + provider_snapshot_id=snapshot_id, + ) + if attachment is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail=f"No attachment found for provider_snapshot_id '{snapshot_id}'.", ) - if attachment is None: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, - detail=f"No attachment found for provider_snapshot_id '{snapshot_id}'.", - ) - deployment = await get_deployment_row( - session, - user_id=current_user.id, - deployment_id=attachment.deployment_id, + deployment = await get_deployment_row( + session, + user_id=current_user.id, + deployment_id=attachment.deployment_id, + ) + if deployment is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail=f"Deployment for attachment (deployment_id={attachment.deployment_id}) not found.", ) - if deployment is None: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, - detail=f"Deployment for attachment (deployment_id={attachment.deployment_id}) not found.", - ) - flow_version = await get_flow_version_entry( - session, - version_id=body.flow_version_id, - user_id=current_user.id, + flow_version = await get_flow_version_entry( + session, + version_id=body.flow_version_id, + user_id=current_user.id, + ) + if flow_version is None: + raise HTTPException( + status_code=status.HTTP_404_NOT_FOUND, + detail=f"Flow version '{body.flow_version_id}' not found.", + ) + if flow_version.data is None: + raise HTTPException( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + detail=f"Flow version '{body.flow_version_id}' has no data.", ) - if flow_version is None: - raise HTTPException( - status_code=status.HTTP_404_NOT_FOUND, - detail=f"Flow version '{body.flow_version_id}' not found.", - ) - if flow_version.data is None: - raise HTTPException( - status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, - detail=f"Flow version '{body.flow_version_id}' has no data.", - ) - provider_account = await get_owned_provider_account_or_404( - provider_id=deployment.deployment_provider_account_id, + provider_account = await get_owned_provider_account_or_404( + provider_id=deployment.deployment_provider_account_id, + 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) + + from langflow.services.database.models.flow.model import Flow + + flow_row = await session.get(Flow, flow_version.flow_id) + + flow_artifact = deployment_mapper.resolve_snapshot_update_artifact( + flow_version=flow_version, + flow_row=flow_row, + deployment=deployment, + ) + + with ( + handle_adapter_errors(mapper=deployment_mapper), + deployment_provider_scope(deployment.deployment_provider_account_id), + ): + await deployment_adapter.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) - - from langflow.services.database.models.flow.model import Flow - - flow_row = await session.get(Flow, flow_version.flow_id) - - flow_artifact = deployment_mapper.resolve_snapshot_update_artifact( - flow_version=flow_version, - flow_row=flow_row, - deployment=deployment, + snapshot_id=snapshot_id, + flow_artifact=flow_artifact, ) - with ( - handle_adapter_errors(mapper=deployment_mapper), - deployment_provider_scope(deployment.deployment_provider_account_id), - ): - await deployment_adapter.update_snapshot( - user_id=current_user.id, - db=session, - snapshot_id=snapshot_id, - flow_artifact=flow_artifact, - ) - - # Provider mutation succeeded — update all local attachment rows that share - # this provider snapshot id. - # If the DB flush fails, attempt a best-effort compensating re-upload - # of the previous flow version's artifact. - # Concurrency note: rollback uses a single previously-read flow_version_id. - # This assumes the snapshot->flow_version invariant held before this call. - # Because the invariant is enforced in app logic (not a DB constraint), - # concurrent writers can still race and violate that assumption. - previous_flow_version_id = attachment.flow_version_id - try: - updated_rows = await update_flow_version_by_provider_snapshot_id( - session, - user_id=current_user.id, - provider_snapshot_id=snapshot_id, - flow_version_id=body.flow_version_id, - ) - if updated_rows == 0: - logger.warning( - "Snapshot '%s' update changed zero attachment rows after provider mutation " - "(user_id=%s, requested_flow_version_id=%s). Possible concurrent modification.", - snapshot_id, - current_user.id, - body.flow_version_id, - ) - await session.commit() - except Exception: - await session.rollback() + # Provider mutation succeeded — update all local attachment rows that share + # this provider snapshot id. + # If the DB flush fails, attempt a best-effort compensating re-upload + # of the previous flow version's artifact. + # Concurrency note: rollback uses a single previously-read flow_version_id. + # This assumes the snapshot->flow_version invariant held before this call. + # Because the invariant is enforced in app logic (not a DB constraint), + # concurrent writers can still race and violate that assumption. + previous_flow_version_id = attachment.flow_version_id + try: + updated_rows = await update_flow_version_by_provider_snapshot_id( + session, + user_id=current_user.id, + provider_snapshot_id=snapshot_id, + flow_version_id=body.flow_version_id, + ) + if updated_rows == 0: logger.warning( - "DB update/commit failed after provider snapshot update for snapshot '%s' " - "(requested_flow_version_id=%s). Attempting compensating provider rollback.", + "Snapshot '%s' update changed zero attachment rows after provider mutation " + "(user_id=%s, requested_flow_version_id=%s). Possible concurrent modification.", + snapshot_id, + current_user.id, + body.flow_version_id, + ) + await session.commit() + except Exception: + await session.rollback() + logger.warning( + "DB update/commit failed after provider snapshot update for snapshot '%s' " + "(requested_flow_version_id=%s). Attempting compensating provider rollback.", + snapshot_id, + body.flow_version_id, + exc_info=True, + ) + try: + prev_version = await get_flow_version_entry( + session, + version_id=previous_flow_version_id, + user_id=current_user.id, + ) + if prev_version and prev_version.data: + prev_artifact = deployment_mapper.resolve_snapshot_update_artifact( + flow_version=prev_version, + flow_row=flow_row, + deployment=deployment, + ) + with deployment_provider_scope(deployment.deployment_provider_account_id): + await deployment_adapter.update_snapshot( + user_id=current_user.id, + db=session, + snapshot_id=snapshot_id, + flow_artifact=prev_artifact, + ) + logger.info( + "Restored provider snapshot '%s' to previous flow_version_id=%s after DB commit failure.", + snapshot_id, + previous_flow_version_id, + ) + except Exception: # noqa: BLE001 + logger.warning( + "Best-effort rollback failed for snapshot '%s'. " + "Provider content reflects flow_version_id=%s but attachment " + "records point to flow_version_id=%s. Manual reconciliation may be needed.", snapshot_id, body.flow_version_id, + previous_flow_version_id, exc_info=True, ) - try: - prev_version = await get_flow_version_entry( - session, - version_id=previous_flow_version_id, - user_id=current_user.id, - ) - if prev_version and prev_version.data: - prev_artifact = deployment_mapper.resolve_snapshot_update_artifact( - flow_version=prev_version, - flow_row=flow_row, - deployment=deployment, - ) - with deployment_provider_scope(deployment.deployment_provider_account_id): - await deployment_adapter.update_snapshot( - user_id=current_user.id, - db=session, - snapshot_id=snapshot_id, - flow_artifact=prev_artifact, - ) - logger.info( - "Restored provider snapshot '%s' to previous flow_version_id=%s after DB commit failure.", - snapshot_id, - previous_flow_version_id, - ) - except Exception: # noqa: BLE001 - logger.warning( - "Best-effort rollback failed for snapshot '%s'. " - "Provider content reflects flow_version_id=%s but attachment " - "records point to flow_version_id=%s. Manual reconciliation may be needed.", - snapshot_id, - body.flow_version_id, - previous_flow_version_id, - exc_info=True, - ) - raise + raise - return SnapshotUpdateResponse( - flow_version_id=body.flow_version_id, - provider_snapshot_id=snapshot_id, - ) + return SnapshotUpdateResponse( + flow_version_id=body.flow_version_id, + provider_snapshot_id=snapshot_id, + ) # Internal note: keep exclude-none for lean responses; use explicit nulls only for intentional tri-state fields. @@ -1296,108 +1296,103 @@ async def update_deployment( session: DbSession, payload: DeploymentUpdateRequest, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(deployment_update_telemetry)], ): - async with _track_deployment_telemetry("deployment.update") as telemetry: - ( - deployment_row, - deployment_adapter, - deployment_mapper, - provider_key, - ) = await resolve_adapter_mapper_from_deployment( - deployment_id=deployment_id, + 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 + adapter_payload = await deployment_mapper.resolve_deployment_update( + user_id=current_user.id, + deployment_db_id=deployment_row_id, + db=session, + payload=payload, + ) + added_flow_version_ids, remove_flow_version_ids = resolve_flow_version_patch_for_update( + deployment_mapper=deployment_mapper, + payload=payload, + ) + await validate_project_scoped_flow_version_ids( + flow_version_ids=list(dict.fromkeys([*added_flow_version_ids, *remove_flow_version_ids])), + user_id=current_user.id, + project_id=deployment_row.project_id, + db=session, + ) + with handle_adapter_errors(mapper=deployment_mapper), deployment_provider_scope(deployment_provider_account_id): + update_result: DeploymentUpdateResult = await deployment_adapter.update( + deployment_id=deployment_resource_key, + payload=adapter_payload, 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 - adapter_payload = await deployment_mapper.resolve_deployment_update( + try: + existing_attachments = await list_deployment_attachments_for_flow_version_ids( + session, user_id=current_user.id, - deployment_db_id=deployment_row_id, - db=session, - payload=payload, + deployment_id=deployment_row_id, + flow_version_ids=added_flow_version_ids, ) - added_flow_version_ids, remove_flow_version_ids = resolve_flow_version_patch_for_update( + already_attached = {a.flow_version_id for a in existing_attachments} + newly_added_flow_version_ids = [fv for fv in added_flow_version_ids if fv not in already_attached] + added_snapshot_bindings = resolve_added_snapshot_bindings_for_update( deployment_mapper=deployment_mapper, - payload=payload, + added_flow_version_ids=newly_added_flow_version_ids, + result=update_result, ) - await validate_project_scoped_flow_version_ids( - flow_version_ids=list(dict.fromkeys([*added_flow_version_ids, *remove_flow_version_ids])), + await apply_flow_version_patch_attachments( user_id=current_user.id, - project_id=deployment_row.project_id, + deployment_row_id=deployment_row_id, + added_snapshot_bindings=added_snapshot_bindings, + remove_flow_version_ids=remove_flow_version_ids, db=session, ) - with handle_adapter_errors(mapper=deployment_mapper), deployment_provider_scope(deployment_provider_account_id): - update_result: DeploymentUpdateResult = await deployment_adapter.update( - deployment_id=deployment_resource_key, - payload=adapter_payload, - user_id=current_user.id, - db=session, - ) - try: - existing_attachments = await list_deployment_attachments_for_flow_version_ids( - session, - user_id=current_user.id, - deployment_id=deployment_row_id, - flow_version_ids=added_flow_version_ids, - ) - already_attached = {a.flow_version_id for a in existing_attachments} - newly_added_flow_version_ids = [fv for fv in added_flow_version_ids if fv not in already_attached] - added_snapshot_bindings = resolve_added_snapshot_bindings_for_update( - deployment_mapper=deployment_mapper, - added_flow_version_ids=newly_added_flow_version_ids, - result=update_result, - ) - await apply_flow_version_patch_attachments( - user_id=current_user.id, - deployment_row_id=deployment_row_id, - added_snapshot_bindings=added_snapshot_bindings, - remove_flow_version_ids=remove_flow_version_ids, - db=session, - ) - update_kwargs: dict = {} - if payload.name is not None and payload.name != deployment_row.name: - update_kwargs["name"] = payload.name - if _field_was_explicitly_set(payload, "description"): - if payload.description != deployment_row.description: - update_kwargs["description"] = payload.description - elif payload.description is not None and payload.description != deployment_row.description: + update_kwargs: dict = {} + if payload.name is not None and payload.name != deployment_row.name: + update_kwargs["name"] = payload.name + if _field_was_explicitly_set(payload, "description"): + if payload.description != deployment_row.description: update_kwargs["description"] = payload.description - if update_kwargs: - deployment_row = await update_deployment_db( - session, - deployment=deployment_row, - **update_kwargs, - ) - - await session.commit() - except Exception as exc: - # Provider was already mutated by deployment_adapter.update above. - # Roll back the session to discard any pending DB changes (or reset - # it from the "inactive" state after a failed commit) so the mapper - # can query the original attachment rows and build a compensating - # payload. - await session.rollback() - await rollback_provider_update( - deployment_adapter=deployment_adapter, - deployment_mapper=deployment_mapper, - deployment_db_id=deployment_row_id, - deployment_resource_key=deployment_resource_key, - deployment_provider_account_id=deployment_provider_account_id, - user_id=current_user.id, - db=session, + elif payload.description is not None and payload.description != deployment_row.description: + update_kwargs["description"] = payload.description + if update_kwargs: + deployment_row = await update_deployment_db( + session, + deployment=deployment_row, + **update_kwargs, ) - if isinstance(exc, AttachmentConflictError): - raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(exc)) from exc - raise - return deployment_mapper.shape_deployment_update_result( - update_result, - deployment_row, - provider_key=provider_key, + await session.commit() + except Exception as exc: + # Provider was already mutated by deployment_adapter.update above. + # Roll back the session to discard any pending DB changes (or reset + # it from the "inactive" state after a failed commit) so the mapper + # can query the original attachment rows and build a compensating + # payload. + await session.rollback() + await rollback_provider_update( + deployment_adapter=deployment_adapter, + deployment_mapper=deployment_mapper, + deployment_db_id=deployment_row_id, + deployment_resource_key=deployment_resource_key, + deployment_provider_account_id=deployment_provider_account_id, + user_id=current_user.id, + db=session, ) + if isinstance(exc, AttachmentConflictError): + raise HTTPException(status_code=status.HTTP_409_CONFLICT, detail=str(exc)) from exc + raise + + return deployment_mapper.shape_deployment_update_result( + update_result, + deployment_row, + provider_key=provider_key, + ) @router.delete("/{deployment_id}", status_code=status.HTTP_204_NO_CONTENT) @@ -1405,40 +1400,40 @@ async def delete_deployment( deployment_id: DeploymentIdPath, session: DbSession, current_user: CurrentActiveUser, + telemetry: Annotated[DeploymentTelemetryCtx, Depends(deployment_delete_telemetry)], *, include_provider: IncludeProviderDeleteQuery = True, ): - async with _track_deployment_telemetry("deployment.delete") as telemetry: - deployment_row, deployment_adapter, _provider_key = await resolve_adapter_from_deployment( - deployment_id=deployment_id, - 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): - await deployment_adapter.delete( - deployment_id=deployment_row.resource_key, - user_id=current_user.id, - db=session, - ) - except HTTPException as exc: - if exc.status_code != status.HTTP_404_NOT_FOUND: - raise - logger.warning( - "Deployment %s (resource_key=%s) already missing on provider %s during delete; deleting stale row.", - deployment_row.id, - deployment_row.resource_key, - deployment_row.deployment_provider_account_id, + deployment_row, deployment_adapter, _provider_key = await resolve_adapter_from_deployment( + deployment_id=deployment_id, + 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): + await deployment_adapter.delete( + deployment_id=deployment_row.resource_key, + user_id=current_user.id, + db=session, ) - await _delete_local_deployment_row_with_commit_retry( - session=session, - deployment_id=deployment_row.id, - user_id=current_user.id, - resource_key=deployment_row.resource_key, - ) - return Response(status_code=status.HTTP_204_NO_CONTENT) + except HTTPException as exc: + if exc.status_code != status.HTTP_404_NOT_FOUND: + raise + logger.warning( + "Deployment %s (resource_key=%s) already missing on provider %s during delete; deleting stale row.", + deployment_row.id, + deployment_row.resource_key, + deployment_row.deployment_provider_account_id, + ) + await _delete_local_deployment_row_with_commit_retry( + session=session, + deployment_id=deployment_row.id, + user_id=current_user.id, + resource_key=deployment_row.resource_key, + ) + return Response(status_code=status.HTTP_204_NO_CONTENT) @router.get( diff --git a/src/backend/base/langflow/services/telemetry/schema.py b/src/backend/base/langflow/services/telemetry/schema.py index ba55225f29..64fab76e85 100644 --- a/src/backend/base/langflow/services/telemetry/schema.py +++ b/src/backend/base/langflow/services/telemetry/schema.py @@ -22,9 +22,9 @@ class RunPayload(BasePayload): class DeploymentPayload(BasePayload): deployment_action: str = Field(serialization_alias="deploymentAction") deployment_provider: str = Field(serialization_alias="deploymentProvider") - deployment_seconds: int = Field(serialization_alias="deploymentSeconds") + deployment_seconds: float = Field(serialization_alias="deploymentSeconds") deployment_success: bool = Field(serialization_alias="deploymentSuccess") - deployment_error_message: str = Field("", serialization_alias="deploymentErrorMessage") + deployment_error_type: str | None = Field(default=None, serialization_alias="deploymentErrorType") class ShutdownPayload(BasePayload): diff --git a/src/backend/base/langflow/services/telemetry/service.py b/src/backend/base/langflow/services/telemetry/service.py index 97a7bd3298..4bafdfd42f 100644 --- a/src/backend/base/langflow/services/telemetry/service.py +++ b/src/backend/base/langflow/services/telemetry/service.py @@ -114,6 +114,12 @@ class TelemetryService(Service): async def log_package_deployment(self, payload: DeploymentPayload) -> None: await self._queue_event((self.send_telemetry_data, payload, "deployment")) + async def log_package_deployment_provider(self, payload: DeploymentPayload) -> None: + await self._queue_event((self.send_telemetry_data, payload, "deployment_provider")) + + async def log_package_deployment_run(self, payload: DeploymentPayload) -> None: + await self._queue_event((self.send_telemetry_data, payload, "deployment_run")) + async def log_package_shutdown(self) -> None: payload = ShutdownPayload(time_running=(datetime.now(timezone.utc) - self._start_time).seconds) await self._queue_event(payload) diff --git a/src/backend/tests/unit/api/v1/test_deployments_telemetry.py b/src/backend/tests/unit/api/v1/test_deployments_telemetry.py index 703162738f..134cd28ace 100644 --- a/src/backend/tests/unit/api/v1/test_deployments_telemetry.py +++ b/src/backend/tests/unit/api/v1/test_deployments_telemetry.py @@ -146,24 +146,6 @@ def mock_db_crud(mock_mapper): yield -def _get_deployment_call_payload(mock_telemetry_service): - """Inspect the single DeploymentPayload handed to background_tasks.add_task. - - Telemetry now goes through BackgroundTasks, so the mocked service method is - referenced rather than awaited. The payload is the positional arg passed to - ``add_task(log_package_deployment, payload)``. - """ - mock_method = mock_telemetry_service.log_package_deployment - # Background tasks may either await the coroutine function directly or - # queue it via BackgroundTasks; handle both shapes. - if mock_method.await_args_list: - return mock_method.await_args_list[-1].args[0] - if mock_method.call_args_list: - return mock_method.call_args_list[-1].args[0] - msg = "log_package_deployment was not invoked" - raise AssertionError(msg) - - @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 @@ -174,11 +156,11 @@ async def test_create_provider_account_telemetry( headers=logged_in_headers, ) assert response.status_code == status.HTTP_201_CREATED - payload = _get_deployment_call_payload(mock_telemetry_service) + 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.deployment_error_message == "" @pytest.mark.asyncio @@ -191,7 +173,8 @@ async def test_update_provider_account_telemetry( headers=logged_in_headers, ) assert response.status_code == status.HTTP_200_OK - payload = _get_deployment_call_payload(mock_telemetry_service) + 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 @@ -203,7 +186,8 @@ async def test_delete_provider_account_telemetry( ): response = await client.delete(f"api/v1/deployments/providers/{uuid4()}", headers=logged_in_headers) assert response.status_code == status.HTTP_204_NO_CONTENT - payload = _get_deployment_call_payload(mock_telemetry_service) + 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 @@ -219,7 +203,8 @@ async def test_create_deployment_telemetry( headers=logged_in_headers, ) assert response.status_code == status.HTTP_201_CREATED, response.json() - payload = _get_deployment_call_payload(mock_telemetry_service) + 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 @@ -231,7 +216,8 @@ async def test_update_deployment_telemetry( ): response = await client.patch(f"api/v1/deployments/{uuid4()}", json={"name": "Test"}, headers=logged_in_headers) assert response.status_code == status.HTTP_200_OK - payload = _get_deployment_call_payload(mock_telemetry_service) + 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 @@ -243,7 +229,8 @@ async def test_delete_deployment_telemetry( ): response = await client.delete(f"api/v1/deployments/{uuid4()}", headers=logged_in_headers) assert response.status_code == status.HTTP_204_NO_CONTENT - payload = _get_deployment_call_payload(mock_telemetry_service) + 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 @@ -257,7 +244,8 @@ async def test_create_deployment_run_telemetry( f"api/v1/deployments/{uuid4()}/runs", json={"provider_data": {}}, headers=logged_in_headers ) assert response.status_code == status.HTTP_201_CREATED, response.json() - payload = _get_deployment_call_payload(mock_telemetry_service) + 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 @@ -271,7 +259,8 @@ async def test_update_snapshot_telemetry( "api/v1/deployments/snapshots/snap-1", json={"flow_version_id": str(uuid4())}, headers=logged_in_headers ) assert response.status_code == status.HTTP_200_OK - payload = _get_deployment_call_payload(mock_telemetry_service) + 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 @@ -288,11 +277,12 @@ async def test_create_provider_account_telemetry_error( headers=logged_in_headers, ) assert response.status_code == status.HTTP_400_BAD_REQUEST - payload = _get_deployment_call_payload(mock_telemetry_service) + 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 # non-empty; HTTPException.str() includes the detail + assert payload.deployment_error_type == "HTTPException" @pytest.mark.asyncio @@ -307,8 +297,9 @@ async def test_cross_route_smoke_exception_after_provider_set( headers=logged_in_headers, ) assert response.status_code == status.HTTP_500_INTERNAL_SERVER_ERROR - payload = _get_deployment_call_payload(mock_telemetry_service) + 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 # non-empty error message captured + assert payload.deployment_error_type == "HTTPException" diff --git a/src/backend/tests/unit/services/telemetry/test_telemetry_schema.py b/src/backend/tests/unit/services/telemetry/test_telemetry_schema.py index e9f682776f..fbd3abf5e7 100644 --- a/src/backend/tests/unit/services/telemetry/test_telemetry_schema.py +++ b/src/backend/tests/unit/services/telemetry/test_telemetry_schema.py @@ -25,17 +25,17 @@ class TestDeploymentPayload: payload = DeploymentPayload( deployment_action="deployment.create", deployment_provider="test_provider", - deployment_seconds=2, + deployment_seconds=1.5, deployment_success=True, - deployment_error_message="", + deployment_error_type=None, client_type="oss", ) assert payload.deployment_action == "deployment.create" assert payload.deployment_provider == "test_provider" - assert payload.deployment_seconds == 2 + assert payload.deployment_seconds == 1.5 assert payload.deployment_success is True - assert payload.deployment_error_message == "" + assert payload.deployment_error_type is None assert payload.client_type == "oss" def test_deployment_payload_initialization_with_defaults(self): @@ -43,46 +43,36 @@ class TestDeploymentPayload: payload = DeploymentPayload( deployment_action="deployment.delete", deployment_provider="test_provider", - deployment_seconds=1, + deployment_seconds=0.5, deployment_success=False, - deployment_error_message="Invalid credentials", + deployment_error_type="ValueError", ) assert payload.deployment_action == "deployment.delete" assert payload.deployment_provider == "test_provider" - assert payload.deployment_seconds == 1 + assert payload.deployment_seconds == 0.5 assert payload.deployment_success is False - assert payload.deployment_error_message == "Invalid credentials" + assert payload.deployment_error_type == "ValueError" assert payload.client_type is None # Default value - def test_deployment_payload_error_message_defaults_to_empty(self): - """Omitting deployment_error_message yields an empty string (matches RunPayload).""" - payload = DeploymentPayload( - deployment_action="deployment.create", - deployment_provider="test_provider", - deployment_seconds=2, - deployment_success=True, - ) - assert payload.deployment_error_message == "" - def test_deployment_payload_serialization(self): """Test DeploymentPayload serialization to dictionary.""" payload = DeploymentPayload( deployment_action="deployment.update", deployment_provider="test_provider", - deployment_seconds=3, + deployment_seconds=2.0, deployment_success=True, - deployment_error_message="", + deployment_error_type=None, client_type="desktop", ) - data = payload.model_dump(by_alias=True) + data = payload.model_dump(by_alias=True, exclude_none=True) assert data["deploymentAction"] == "deployment.update" assert data["deploymentProvider"] == "test_provider" - assert data["deploymentSeconds"] == 3 + assert data["deploymentSeconds"] == 2.0 assert data["deploymentSuccess"] is True - assert data["deploymentErrorMessage"] == "" + assert "deploymentErrorType" not in data assert data["clientType"] == "desktop" def test_deployment_payload_roundtrip(self): @@ -90,9 +80,9 @@ class TestDeploymentPayload: payload = DeploymentPayload( deployment_action="provider.create", deployment_provider="test_provider", - deployment_seconds=1, + deployment_seconds=1.0, deployment_success=False, - deployment_error_message="boom", + deployment_error_type="Exception", client_type="oss", ) @@ -103,7 +93,7 @@ class TestDeploymentPayload: 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.deployment_error_type == payload.deployment_error_type assert new_payload.client_type == payload.client_type diff --git a/src/backend/tests/unit/test_telemetry.py b/src/backend/tests/unit/test_telemetry.py index 0ca63ebe18..9b0f197ff4 100644 --- a/src/backend/tests/unit/test_telemetry.py +++ b/src/backend/tests/unit/test_telemetry.py @@ -27,7 +27,7 @@ async def test_log_package_deployment(telemetry_service): payload = DeploymentPayload( deployment_action="deployment.create", deployment_provider="test_provider", - deployment_seconds=1, + deployment_seconds=1.0, deployment_success=True, ) await telemetry_service.log_package_deployment(payload) @@ -37,16 +37,48 @@ async def test_log_package_deployment(telemetry_service): 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, + 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()