refactor: Align deployments telemetry with existing langflow patterns

- DeploymentPayload: replace deployment_error_type (class name) with
  deployment_error_message: str (captures str(exc)) to match
  RunPayload/PlaygroundPayload shape; switch deployment_seconds from
  float to int for consistency with other payloads.
- TelemetryService: collapse log_package_deployment_provider and
  log_package_deployment_run into a single log_package_deployment so
  payload type maps to one Scarf path (the action is already carried in
  DeploymentPayload.deployment_action).
- deployments.py: replace the yield-based FastAPI Depends factory with
  an @asynccontextmanager helper (_track_deployment_telemetry) used
  inline via `async with` in each of the 8 instrumented routes. Emits
  telemetry via `await telemetry_service.log_package_deployment(...)`
  inside the CM's finally, matching the endpoints.py:282 inline-await
  pattern. background_tasks.add_task was rejected because FastAPI drops
  background tasks when the handler raises, which would silently lose
  all failure telemetry.
- Tests updated for the new payload fields, single service method,
  and inline-await emission.
This commit is contained in:
Jordan Frazier
2026-04-24 11:59:13 -04:00
parent 7bcb27021a
commit a9a09f491b
6 changed files with 573 additions and 587 deletions

File diff suppressed because it is too large Load Diff

View File

@ -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: float = Field(serialization_alias="deploymentSeconds")
deployment_seconds: int = Field(serialization_alias="deploymentSeconds")
deployment_success: bool = Field(serialization_alias="deploymentSuccess")
deployment_error_type: str | None = Field(default=None, serialization_alias="deploymentErrorType")
deployment_error_message: str = Field("", serialization_alias="deploymentErrorMessage")
class ShutdownPayload(BasePayload):

View File

@ -114,12 +114,6 @@ 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)

View File

@ -146,6 +146,24 @@ 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
@ -156,11 +174,11 @@ async def test_create_provider_account_telemetry(
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]
payload = _get_deployment_call_payload(mock_telemetry_service)
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
@ -173,8 +191,7 @@ async def test_update_provider_account_telemetry(
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]
payload = _get_deployment_call_payload(mock_telemetry_service)
assert payload.deployment_action == "provider.update"
assert payload.deployment_provider == "watsonx-orchestrate"
assert payload.deployment_success is True
@ -186,8 +203,7 @@ 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
mock_telemetry_service.log_package_deployment_provider.assert_awaited_once()
payload = mock_telemetry_service.log_package_deployment_provider.call_args[0][0]
payload = _get_deployment_call_payload(mock_telemetry_service)
assert payload.deployment_action == "provider.delete"
assert payload.deployment_provider == "watsonx-orchestrate"
assert payload.deployment_success is True
@ -203,8 +219,7 @@ async def test_create_deployment_telemetry(
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]
payload = _get_deployment_call_payload(mock_telemetry_service)
assert payload.deployment_action == "deployment.create"
assert payload.deployment_provider == "watsonx-orchestrate"
assert payload.deployment_success is True
@ -216,8 +231,7 @@ 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
mock_telemetry_service.log_package_deployment.assert_awaited_once()
payload = mock_telemetry_service.log_package_deployment.call_args[0][0]
payload = _get_deployment_call_payload(mock_telemetry_service)
assert payload.deployment_action == "deployment.update"
assert payload.deployment_provider == "watsonx-orchestrate"
assert payload.deployment_success is True
@ -229,8 +243,7 @@ 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
mock_telemetry_service.log_package_deployment.assert_awaited_once()
payload = mock_telemetry_service.log_package_deployment.call_args[0][0]
payload = _get_deployment_call_payload(mock_telemetry_service)
assert payload.deployment_action == "deployment.delete"
assert payload.deployment_provider == "watsonx-orchestrate"
assert payload.deployment_success is True
@ -244,8 +257,7 @@ 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()
mock_telemetry_service.log_package_deployment_run.assert_awaited_once()
payload = mock_telemetry_service.log_package_deployment_run.call_args[0][0]
payload = _get_deployment_call_payload(mock_telemetry_service)
assert payload.deployment_action == "deployment.run"
assert payload.deployment_provider == "watsonx-orchestrate"
assert payload.deployment_success is True
@ -259,8 +271,7 @@ 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
mock_telemetry_service.log_package_deployment.assert_awaited_once()
payload = mock_telemetry_service.log_package_deployment.call_args[0][0]
payload = _get_deployment_call_payload(mock_telemetry_service)
assert payload.deployment_action == "snapshot.update"
assert payload.deployment_provider == "watsonx-orchestrate"
assert payload.deployment_success is True
@ -277,12 +288,11 @@ async def test_create_provider_account_telemetry_error(
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]
payload = _get_deployment_call_payload(mock_telemetry_service)
assert payload.deployment_action == "provider.create"
assert payload.deployment_provider == "watsonx-orchestrate"
assert payload.deployment_success is False
assert payload.deployment_error_type == "HTTPException"
assert payload.deployment_error_message # non-empty; HTTPException.str() includes the detail
@pytest.mark.asyncio
@ -297,9 +307,8 @@ async def test_cross_route_smoke_exception_after_provider_set(
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]
payload = _get_deployment_call_payload(mock_telemetry_service)
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"
assert payload.deployment_error_message # non-empty error message captured

View File

@ -25,17 +25,17 @@ class TestDeploymentPayload:
payload = DeploymentPayload(
deployment_action="deployment.create",
deployment_provider="test_provider",
deployment_seconds=1.5,
deployment_seconds=2,
deployment_success=True,
deployment_error_type=None,
deployment_error_message="",
client_type="oss",
)
assert payload.deployment_action == "deployment.create"
assert payload.deployment_provider == "test_provider"
assert payload.deployment_seconds == 1.5
assert payload.deployment_seconds == 2
assert payload.deployment_success is True
assert payload.deployment_error_type is None
assert payload.deployment_error_message == ""
assert payload.client_type == "oss"
def test_deployment_payload_initialization_with_defaults(self):
@ -43,36 +43,46 @@ class TestDeploymentPayload:
payload = DeploymentPayload(
deployment_action="deployment.delete",
deployment_provider="test_provider",
deployment_seconds=0.5,
deployment_seconds=1,
deployment_success=False,
deployment_error_type="ValueError",
deployment_error_message="Invalid credentials",
)
assert payload.deployment_action == "deployment.delete"
assert payload.deployment_provider == "test_provider"
assert payload.deployment_seconds == 0.5
assert payload.deployment_seconds == 1
assert payload.deployment_success is False
assert payload.deployment_error_type == "ValueError"
assert payload.deployment_error_message == "Invalid credentials"
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=2.0,
deployment_seconds=3,
deployment_success=True,
deployment_error_type=None,
deployment_error_message="",
client_type="desktop",
)
data = payload.model_dump(by_alias=True, exclude_none=True)
data = payload.model_dump(by_alias=True)
assert data["deploymentAction"] == "deployment.update"
assert data["deploymentProvider"] == "test_provider"
assert data["deploymentSeconds"] == 2.0
assert data["deploymentSeconds"] == 3
assert data["deploymentSuccess"] is True
assert "deploymentErrorType" not in data
assert data["deploymentErrorMessage"] == ""
assert data["clientType"] == "desktop"
def test_deployment_payload_roundtrip(self):
@ -80,9 +90,9 @@ class TestDeploymentPayload:
payload = DeploymentPayload(
deployment_action="provider.create",
deployment_provider="test_provider",
deployment_seconds=1.0,
deployment_seconds=1,
deployment_success=False,
deployment_error_type="Exception",
deployment_error_message="boom",
client_type="oss",
)
@ -93,7 +103,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_type == payload.deployment_error_type
assert new_payload.deployment_error_message == payload.deployment_error_message
assert new_payload.client_type == payload.client_type

View File

@ -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.0,
deployment_seconds=1,
deployment_success=True,
)
await telemetry_service.log_package_deployment(payload)
@ -37,48 +37,16 @@ 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.0,
deployment_seconds=1,
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()