Revert "refactor: Align deployments telemetry with existing langflow patterns"

This reverts commit a9a09f491b.
This commit is contained in:
Jordan Frazier
2026-04-24 14:07:48 -04:00
parent 2fb841b8ca
commit 2199481642
6 changed files with 581 additions and 567 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: 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):

View File

@ -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)

View File

@ -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"

View File

@ -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

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,
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()