mirror of
https://github.com/dograh-hq/dograh.git
synced 2026-07-25 12:01:04 +02:00
chore: format and minor cleanups
This commit is contained in:
parent
8c9938bde1
commit
ce19d0cfdc
6 changed files with 22 additions and 18 deletions
|
|
@ -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}"
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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"])
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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})"
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue