From ce19d0cfdc2471aeb12661aca279827ea4cefcfa Mon Sep 17 00:00:00 2001 From: Abhishek Kumar Date: Tue, 21 Jul 2026 13:14:35 +0530 Subject: [PATCH] chore: format and minor cleanups --- api/db/workflow_run_client.py | 9 +++++++-- api/routes/workflow.py | 2 +- api/routes/workflow_text_chat.py | 2 +- api/services/campaign/campaign_call_dispatcher.py | 12 +++++------- .../telephony/providers/cloudonix/provider.py | 8 ++------ api/tests/test_from_number_pool_isolation.py | 7 ++++++- 6 files changed, 22 insertions(+), 18 deletions(-) diff --git a/api/db/workflow_run_client.py b/api/db/workflow_run_client.py index 95e05f3a..7b84d826 100644 --- a/api/db/workflow_run_client.py +++ b/api/db/workflow_run_client.py @@ -55,8 +55,13 @@ class WorkflowRunClient(BaseDBClient): raise ValueError(f"Workflow with ID {workflow_id} not found") if definition_id is not None: - definition = await session.get(WorkflowDefinitionModel, definition_id) - if not definition or definition.workflow_id != workflow.id: + definition_result = await session.execute( + select(WorkflowDefinitionModel.id).where( + WorkflowDefinitionModel.id == definition_id, + WorkflowDefinitionModel.workflow_id == workflow.id, + ) + ) + if definition_result.scalar_one_or_none() is None: raise ValueError( f"Workflow definition {definition_id} does not belong to " f"workflow {workflow.id}" diff --git a/api/routes/workflow.py b/api/routes/workflow.py index 2020468d..daa0c184 100644 --- a/api/routes/workflow.py +++ b/api/routes/workflow.py @@ -47,11 +47,11 @@ from api.services.storage import storage_fs from api.services.workflow.dto import ReactFlowDTO, sanitize_workflow_definition from api.services.workflow.duplicate import duplicate_workflow from api.services.workflow.errors import ItemKind, WorkflowError +from api.services.workflow.run_creation import prepare_workflow_run_inputs from api.services.workflow.run_usage_response import ( format_public_cost_info, format_public_usage_info, ) -from api.services.workflow.run_creation import prepare_workflow_run_inputs from api.services.workflow.trigger_paths import ( TriggerPathIssue, ensure_trigger_paths, diff --git a/api/routes/workflow_text_chat.py b/api/routes/workflow_text_chat.py index dd7844be..d6db5529 100644 --- a/api/routes/workflow_text_chat.py +++ b/api/routes/workflow_text_chat.py @@ -11,6 +11,7 @@ from api.db.models import UserModel, WorkflowRunTextSessionModel from api.enums import WorkflowRunMode from api.services.auth.depends import get_user_with_selected_organization from api.services.quota_service import authorize_workflow_run_start +from api.services.workflow.run_creation import prepare_workflow_run_inputs from api.services.workflow.text_chat_session_service import ( TextChatPendingTurnLostError, TextChatSessionExecutionError, @@ -25,7 +26,6 @@ from api.services.workflow.text_chat_session_service import ( normalize_text_chat_session_data, rewind_text_chat_session_state, ) -from api.services.workflow.run_creation import prepare_workflow_run_inputs router = APIRouter(prefix="/workflow", tags=["workflow-text-chat"]) diff --git a/api/services/campaign/campaign_call_dispatcher.py b/api/services/campaign/campaign_call_dispatcher.py index d50a639a..39688473 100644 --- a/api/services/campaign/campaign_call_dispatcher.py +++ b/api/services/campaign/campaign_call_dispatcher.py @@ -238,14 +238,12 @@ class CampaignCallDispatcher: try: # Get workflow details - workflow = await db_client.get_workflow_by_id(campaign.workflow_id) + workflow = await db_client.get_workflow( + campaign.workflow_id, + organization_id=campaign.organization_id, + ) if not workflow: raise ValueError(f"Workflow {campaign.workflow_id} not found") - if workflow.organization_id != campaign.organization_id: - raise ValueError( - f"Workflow {campaign.workflow_id} does not belong to " - f"organization {campaign.organization_id}" - ) # Extract phone number phone_number = queued_run.context_variables.get("phone_number") @@ -308,7 +306,7 @@ class CampaignCallDispatcher: from_number, telephony_configuration_id=campaign.telephony_configuration_id, ) - except Exception as e: + except Exception: # Release slot and from_number on error if slot_bound and workflow_run: await call_concurrency.release_workflow_run_slot(workflow_run.id) diff --git a/api/services/telephony/providers/cloudonix/provider.py b/api/services/telephony/providers/cloudonix/provider.py index 75d80e06..ff428589 100644 --- a/api/services/telephony/providers/cloudonix/provider.py +++ b/api/services/telephony/providers/cloudonix/provider.py @@ -1111,9 +1111,7 @@ class CloudonixProvider(TelephonyProvider): from_number = random.choice(self.from_numbers) backend_endpoint, _ = await get_backend_endpoints() - callback_url = ( - f"{backend_endpoint}/api/v1/telephony/cloudonix/transfer-result/{transfer_id}" - ) + callback_url = f"{backend_endpoint}/api/v1/telephony/cloudonix/transfer-result/{transfer_id}" endpoint = f"{self.base_url}/calls/{self.domain_id}/application" data: Dict[str, Any] = { @@ -1126,9 +1124,7 @@ class CloudonixProvider(TelephonyProvider): data.update(kwargs) headers = self._get_auth_headers() - masked_destination = ( - f"***{destination[-4:]}" if len(destination) > 4 else "***" - ) + masked_destination = f"***{destination[-4:]}" if len(destination) > 4 else "***" logger.info( f"[Cloudonix Transfer] Dialing {masked_destination} into conference " f"{conference_name} (transfer_id={transfer_id})" diff --git a/api/tests/test_from_number_pool_isolation.py b/api/tests/test_from_number_pool_isolation.py index 5e393111..4ce498ee 100644 --- a/api/tests/test_from_number_pool_isolation.py +++ b/api/tests/test_from_number_pool_isolation.py @@ -281,7 +281,7 @@ class TestDispatcherThreadsTelephonyConfig: ), ), ): - mock_db.get_workflow_by_id = AsyncMock(return_value=SimpleNamespace(id=1)) + mock_db.get_workflow = AsyncMock(return_value=SimpleNamespace(id=1)) mock_db.create_workflow_run = AsyncMock(return_value=workflow_run) mock_db.update_workflow_run = AsyncMock() mock_concurrency.bind_workflow_run = AsyncMock() @@ -300,6 +300,11 @@ class TestDispatcherThreadsTelephonyConfig: ) await dispatcher.dispatch_call(queued_run, campaign, slot) + mock_db.get_workflow.assert_awaited_once_with( + campaign.workflow_id, + organization_id=org_id, + ) + # acquire_from_number on rate_limiter must be called with the # campaign's telephony_configuration_id. assert mock_rl.acquire_from_number.await_count == 1