feat: make vici configurable from UI

This commit is contained in:
Abhishek Kumar 2026-07-20 19:40:28 +05:30
parent cf80b20be1
commit 8628e576ca
46 changed files with 2277 additions and 618 deletions

View file

@ -30,8 +30,10 @@ from api.services.call_concurrency import (
CallConcurrencyLimitError,
call_concurrency,
)
from api.services.organization_preferences import external_pbx_integrations_enabled
from api.services.quota_service import authorize_workflow_run_start
from api.services.telephony.call_transfer_manager import get_call_transfer_manager
from api.services.telephony.providers.ari.external_pbx import create_adapter
from api.services.telephony.transfer_event_protocol import (
TransferEvent,
TransferEventType,
@ -56,6 +58,7 @@ class ARIConnection:
app_name: str,
app_password: str,
ws_client_name: str = "",
external_pbx_config: Optional[dict] = None,
):
self.organization_id = organization_id
self.telephony_configuration_id = telephony_configuration_id
@ -63,6 +66,8 @@ class ARIConnection:
self.app_name = app_name
self.app_password = app_password
self.ws_client_name = ws_client_name
self.external_pbx_config = external_pbx_config
self.external_pbx_adapter = create_adapter(external_pbx_config)
self._ws: Optional[websockets.ClientConnection] = None
self._task: Optional[asyncio.Task] = None
@ -287,7 +292,7 @@ class ARIConnection:
channel_state = channel.get("state", "unknown")
# Log all events for each channel
logger.debug(
logger.trace(
f"[ARI EVENT org={self.organization_id}] {event_type}: channel={channel_id}, state={channel_state}"
)
@ -413,7 +418,7 @@ class ARIConnection:
)
else:
logger.debug(
logger.trace(
f"[ARI org={self.organization_id}] Event: {event_type} "
f"channel={channel_id}"
)
@ -428,9 +433,9 @@ class ARIConnection:
async with session.request(method, url, auth=auth, **kwargs) as response:
response_text = await response.text()
if response.status not in (200, 201, 204):
logger.error(
logger.warning(
f"[ARI org={self.organization_id}] REST API error: "
f"{method} {path} -> {response.status}: {response_text}"
f"{method} {path} {kwargs} -> {response.status}: {response_text}"
)
return {}
if response_text:
@ -451,82 +456,33 @@ class ARIConnection:
)
return (result or {}).get("value", "") or ""
async def _capture_upstream_pbx(
async def _capture_external_pbx_call(
self, channel_id: str, channel_name: str = ""
) -> Optional[dict]:
"""Capture upstream-PBX identity from the inbound SIP headers.
The customer's real call leg lives on the upstream PBX, not on dograh, so
dograh drives hangup/transfer via the upstream's API using a captured
handle. Two providers are supported:
* FreeSWITCH bridges in with ``X-PBX-*`` headers; the handle is the
channel UUID (``X-PBX-UUID``), driven over ESL (uuid_kill/transfer).
* VICIdial patches in with ``X-VICIDIAL-*`` headers; the handle is the
callerid + remote-agent user, driven over ``ra_call_control``.
Returns None for non-upstream (direct) calls.
"""
"""Capture adapter-defined identity from inbound SIP headers."""
if self.external_pbx_adapter is None:
return None
# PJSIP_HEADER() only works on a PJSIP channel; on any other technology
# (Local, WebSocket, etc.) Asterisk returns a 500 ("This function
# requires a PJSIP channel"). Non-PJSIP legs carry no SIP headers to
# capture anyway, so skip the reads quietly instead of spamming errors.
if not channel_name.startswith("PJSIP/"):
logger.debug(
f"[ARI org={self.organization_id}] Skipping upstream_pbx capture "
f"[ARI org={self.organization_id}] Skipping external PBX capture "
f"for non-PJSIP channel {channel_id} ({channel_name or 'unknown'})"
)
return None
# FreeSWITCH: X-PBX-Provider marks the call; X-PBX-UUID is the ESL handle.
if (
await self._get_channel_var(channel_id, "PJSIP_HEADER(read,X-PBX-Provider)")
) == "freeswitch":
upstream = {
"provider": "freeswitch",
"uuid": await self._get_channel_var(
channel_id, "PJSIP_HEADER(read,X-PBX-UUID)"
),
"lead_id": await self._get_channel_var(
channel_id, "PJSIP_HEADER(read,X-PBX-Lead-ID)"
),
"campaign_id": await self._get_channel_var(
channel_id, "PJSIP_HEADER(read,X-PBX-Campaign)"
),
}
logger.info(
f"[ARI org={self.organization_id}] Captured upstream_pbx for channel "
f"{channel_id}: {upstream}"
)
return upstream
async def read_header(name: str) -> str:
return await self._get_channel_var(channel_id, f"PJSIP_HEADER(read,{name})")
# VICIdial: X-VICIDIAL-callerid + user is the ra_call_control handle.
callerid = await self._get_channel_var(
channel_id, "PJSIP_HEADER(read,X-VICIDIAL-callerid)"
)
if not callerid:
return None
upstream = {
"provider": "vicidial",
"callerid": callerid,
"agent_user": await self._get_channel_var(
channel_id, "PJSIP_HEADER(read,X-VICIDIAL-user)"
),
"lead_id": await self._get_channel_var(
channel_id, "PJSIP_HEADER(read,X-VICIDIAL-lead_id)"
),
"campaign_id": await self._get_channel_var(
channel_id, "PJSIP_HEADER(read,X-VICIDIAL-campaign_id)"
),
# The in-group the call arrived on, so a transfer can bounce it back
# to the same queue via INGROUPTRANSFER (destination "ingroup:source").
"ingroup_id": await self._get_channel_var(
channel_id, "PJSIP_HEADER(read,X-VICIDIAL-ingroup_id)"
),
}
logger.info(
f"[ARI org={self.organization_id}] Captured upstream_pbx for channel "
f"{channel_id}: {upstream}"
)
return upstream
identity = await self.external_pbx_adapter.capture_call_identity(read_header)
if identity:
logger.info(
f"[ARI org={self.organization_id}] Captured "
f"{self.external_pbx_adapter.type} call identity for channel {channel_id} identity: {identity}"
)
return identity
async def _create_external_media(
self,
@ -677,10 +633,8 @@ class ARIConnection:
# 3. Create workflow run
call_id = channel_id
# Capture the upstream-PBX (VICIdial) identity off the SIP headers so
# the hangup/transfer strategies can drive VICIdial's API. The
# customer's real leg lives on VICIdial; this is the handle for it.
upstream_pbx = await self._capture_upstream_pbx(
# Capture the configured external PBX identity from SIP headers.
external_pbx_call = await self._capture_external_pbx_call(
channel_id, channel.get("name", "")
)
workflow_run = await db_client.create_workflow_run(
@ -695,7 +649,7 @@ class ARIConnection:
"direction": "inbound",
"provider": "ari",
"telephony_configuration_id": self.telephony_configuration_id,
"upstream_pbx": upstream_pbx,
"external_pbx_call": external_pbx_call,
},
gathered_context={
"call_id": call_id,
@ -1242,6 +1196,7 @@ class ARIManager:
app_name = config["app_name"]
app_password = config["app_password"]
ws_client_name = config["ws_client_name"]
external_pbx_config = config.get("external_pbx")
conn = ARIConnection(
org_id,
@ -1250,6 +1205,7 @@ class ARIManager:
app_name,
app_password,
ws_client_name,
external_pbx_config,
)
key = conn.connection_key
@ -1274,6 +1230,7 @@ class ARIManager:
or existing.app_name != app_name
or existing.app_password != app_password
or existing.ws_client_name != ws_client_name
or existing.external_pbx_config != external_pbx_config
):
logger.info(
f"[ARI Manager] Config {telephony_configuration_id} "
@ -1311,6 +1268,11 @@ class ARIManager:
app_name = credentials.get("app_name")
app_password = credentials.get("app_password")
ws_client_name = credentials.get("ws_client_name", "")
external_pbx = credentials.get("external_pbx")
if external_pbx and not await external_pbx_integrations_enabled(
row.organization_id
):
external_pbx = None
if not all([ari_endpoint, app_name, app_password]):
logger.warning(
@ -1333,6 +1295,7 @@ class ARIManager:
"app_name": app_name,
"app_password": app_password,
"ws_client_name": ws_client_name,
"external_pbx": external_pbx,
}
)

View file

@ -446,3 +446,18 @@ class TelephonyProvider(ABC):
True if provider supports call transfers, False otherwise
"""
pass
async def transfer_external_pbx_call(
self,
*,
identity: Dict[str, Any],
destination: str,
field_updates: Optional[Dict[str, str]] = None,
) -> Optional[Dict[str, Any]]:
"""Handle an external-PBX-owned customer leg when one is present.
Providers without an external PBX return ``None`` so the ordinary
telephony transfer path continues unchanged.
"""
return None

View file

@ -0,0 +1,43 @@
"""Provider-neutral helpers for external-PBX workflow mappings."""
from __future__ import annotations
from typing import Any, Iterable, Mapping
def _read_path(context: Mapping[str, Any], path: str) -> Any:
normalized = path.strip()
if normalized.startswith("gathered_context."):
normalized = normalized.removeprefix("gathered_context.")
current: Any = context
for part in normalized.split("."):
if not isinstance(current, Mapping):
return None
current = current.get(part)
if current is None and "." not in normalized:
extracted = context.get("extracted_variables")
if isinstance(extracted, Mapping):
current = extracted.get(normalized)
return current
def resolve_external_pbx_field_mappings(
gathered_context: Mapping[str, Any] | None,
mappings: Iterable[Mapping[str, Any]] | None,
) -> dict[str, str]:
"""Return provider field -> non-empty gathered-context value."""
context = gathered_context or {}
resolved: dict[str, str] = {}
for mapping in mappings or []:
context_path = str(mapping.get("context_path", "")).strip()
destination_field = str(mapping.get("destination_field", "")).strip()
if not context_path or not destination_field:
continue
value = _read_path(context, context_path)
if value is None:
continue
text = str(value).strip()
if text:
resolved[destination_field] = text
return resolved

View file

@ -4,8 +4,10 @@ from typing import Any, Dict
from api.services.telephony.registry import (
ProviderSpec,
ProviderUICondition,
ProviderUIField,
ProviderUIMetadata,
ProviderUIOption,
register,
)
@ -20,6 +22,7 @@ def _config_loader(value: Dict[str, Any]) -> Dict[str, Any]:
"ari_endpoint": value.get("ari_endpoint"),
"app_name": value.get("app_name"),
"app_password": value.get("app_password"),
"external_pbx": value.get("external_pbx"),
"from_numbers": value.get("from_numbers", []),
}
@ -58,6 +61,92 @@ _UI_METADATA = ProviderUIMetadata(
type="string-array",
description="SIP extensions/numbers for outbound calls",
),
ProviderUIField(
name="external_pbx.type",
label="External PBX Type",
type="select",
required=False,
description=(
"Enable PBX-specific call control for calls patched into Dograh "
"through this Asterisk configuration."
),
options=[ProviderUIOption(value="vicidial", label="VICIdial")],
section="External PBX",
feature_gate="external_pbx_integrations",
),
ProviderUIField(
name="external_pbx.agent_api.url",
label="Agent API URL",
type="text",
description="Full VICIdial remote-agent API URL, ending in agc/api.php",
placeholder="https://vici.example.com/agc/api.php",
visible_when=ProviderUICondition(
field="external_pbx.type", equals="vicidial"
),
section="External PBX",
feature_gate="external_pbx_integrations",
),
ProviderUIField(
name="external_pbx.agent_api.username",
label="Agent API User",
type="text",
sensitive=True,
visible_when=ProviderUICondition(
field="external_pbx.type", equals="vicidial"
),
section="External PBX",
feature_gate="external_pbx_integrations",
),
ProviderUIField(
name="external_pbx.agent_api.password",
label="Agent API Password",
type="password",
sensitive=True,
visible_when=ProviderUICondition(
field="external_pbx.type", equals="vicidial"
),
section="External PBX",
feature_gate="external_pbx_integrations",
),
ProviderUIField(
name="external_pbx.non_agent_api.url",
label="Non-Agent API URL",
type="text",
required=False,
description=(
"Optional. Required only when a workflow updates VICIdial lead fields."
),
placeholder="https://vici.example.com/vicidial/non_agent_api.php",
visible_when=ProviderUICondition(
field="external_pbx.type", equals="vicidial"
),
section="External PBX",
feature_gate="external_pbx_integrations",
),
ProviderUIField(
name="external_pbx.non_agent_api.username",
label="Non-Agent API User",
type="text",
required=False,
sensitive=True,
visible_when=ProviderUICondition(
field="external_pbx.type", equals="vicidial"
),
section="External PBX",
feature_gate="external_pbx_integrations",
),
ProviderUIField(
name="external_pbx.non_agent_api.password",
label="Non-Agent API Password",
type="password",
required=False,
sensitive=True,
visible_when=ProviderUICondition(
field="external_pbx.type", equals="vicidial"
),
section="External PBX",
feature_gate="external_pbx_integrations",
),
],
)

View file

@ -1,8 +1,69 @@
"""ARI (Asterisk REST Interface) telephony configuration schemas."""
from typing import List, Literal
from typing import List, Literal, Optional
from pydantic import BaseModel, Field
from pydantic import BaseModel, Field, field_validator, model_validator
class VicidialAgentAPIConfiguration(BaseModel):
"""VICIdial remote-agent call-control API configuration."""
url: str = Field(..., min_length=1, description="Full URL to agc/api.php")
username: str = Field(..., min_length=1, description="VICIdial agent API user")
password: str = Field(..., min_length=1, description="VICIdial agent API password")
source: str = Field(default="dograh", description="VICIdial API source tag")
@field_validator("url")
@classmethod
def validate_http_url(cls, value: str) -> str:
stripped = value.strip()
if not stripped.startswith(("http://", "https://")):
raise ValueError("VICIdial agent API URL must use http:// or https://")
return stripped
class VicidialNonAgentAPIConfiguration(BaseModel):
"""Optional VICIdial non-agent API configuration for lead updates."""
url: Optional[str] = Field(default=None, description="Full non_agent_api.php URL")
username: Optional[str] = Field(default=None, description="Non-agent API user")
password: Optional[str] = Field(default=None, description="Non-agent API password")
source: str = Field(default="dograh", description="Non-agent API source tag")
@field_validator("url")
@classmethod
def validate_http_url(cls, value: Optional[str]) -> Optional[str]:
if value is None:
return None
stripped = value.strip()
if stripped and not stripped.startswith(("http://", "https://")):
raise ValueError("VICIdial non-agent API URL must use http:// or https://")
return stripped or None
@model_validator(mode="after")
def validate_complete_credentials(self):
supplied = [self.url, self.username, self.password]
if any(supplied) and not all(supplied):
raise ValueError(
"VICIdial non-agent API URL, username, and password must be "
"configured together"
)
return self
class VicidialExternalPBXConfiguration(BaseModel):
"""External-PBX configuration used by the VICIdial strategy adapter."""
type: Literal["vicidial"] = Field(default="vicidial")
agent_api: VicidialAgentAPIConfiguration
non_agent_api: Optional[VicidialNonAgentAPIConfiguration] = None
timeout_seconds: int = Field(default=8, ge=1, le=30)
@model_validator(mode="after")
def drop_empty_non_agent_configuration(self):
if self.non_agent_api is not None and not self.non_agent_api.url:
self.non_agent_api = None
return self
class ARIConfigurationRequest(BaseModel):
@ -20,6 +81,10 @@ class ARIConfigurationRequest(BaseModel):
default="",
description="websocket_client.conf connection name for externalMedia (e.g., dograh_staging)",
)
external_pbx: Optional[VicidialExternalPBXConfiguration] = Field(
default=None,
description="Optional external PBX connected through this Asterisk instance",
)
from_numbers: List[str] = Field(
default_factory=list,
description="List of SIP extensions/numbers for outbound calls (optional)",
@ -34,4 +99,5 @@ class ARIConfigurationResponse(BaseModel):
app_name: str
app_password: str # Masked
ws_client_name: str = ""
external_pbx: Optional[VicidialExternalPBXConfiguration] = None
from_numbers: List[str]

View file

@ -0,0 +1,14 @@
"""External-PBX adapter entrypoint for the ARI provider."""
from .base import ExternalPBXAdapter, ExternalPBXResult
from .registry import create_adapter, register_adapter, registered_adapter_types
from .vicidial import VicidialAdapter
register_adapter("vicidial", VicidialAdapter)
__all__ = [
"ExternalPBXAdapter",
"ExternalPBXResult",
"create_adapter",
"registered_adapter_types",
]

View file

@ -0,0 +1,44 @@
"""Contracts for PBXs that hand a customer leg to Dograh through Asterisk."""
from __future__ import annotations
from abc import ABC, abstractmethod
from dataclasses import dataclass
from typing import Awaitable, Callable, Mapping
HeaderReader = Callable[[str], Awaitable[str]]
@dataclass(frozen=True)
class ExternalPBXResult:
ok: bool
action: str
message: str
class ExternalPBXAdapter(ABC):
"""PBX-specific operations; ARI continues to own only Dograh's local leg."""
type: str
@abstractmethod
async def capture_call_identity(
self, read_header: HeaderReader
) -> dict[str, str] | None:
"""Read a stable upstream-call identity from inbound SIP headers."""
@abstractmethod
async def hangup(self, identity: Mapping[str, str]) -> ExternalPBXResult:
"""Hang up the customer leg owned by the external PBX."""
@abstractmethod
async def transfer(
self, identity: Mapping[str, str], destination: str
) -> ExternalPBXResult:
"""Transfer the customer leg to a PBX-native destination."""
@abstractmethod
async def update_fields(
self, identity: Mapping[str, str], fields: Mapping[str, str]
) -> ExternalPBXResult:
"""Update provider-native fields associated with the call."""

View file

@ -0,0 +1,34 @@
"""Registry for drop-in external-PBX adapters."""
from __future__ import annotations
from collections.abc import Callable
from typing import Any
from .base import ExternalPBXAdapter
AdapterFactory = Callable[[dict[str, Any]], ExternalPBXAdapter]
_FACTORIES: dict[str, AdapterFactory] = {}
def register_adapter(pbx_type: str, factory: AdapterFactory) -> None:
normalized = pbx_type.strip().lower()
if not normalized:
raise ValueError("External PBX type cannot be empty")
if normalized in _FACTORIES:
raise ValueError(f"External PBX adapter already registered: {normalized}")
_FACTORIES[normalized] = factory
def create_adapter(config: dict[str, Any] | None) -> ExternalPBXAdapter | None:
if not config:
return None
pbx_type = str(config.get("type", "")).strip().lower()
factory = _FACTORIES.get(pbx_type)
if factory is None:
raise ValueError(f"Unsupported external PBX type: {pbx_type or '<empty>'}")
return factory(config)
def registered_adapter_types() -> tuple[str, ...]:
return tuple(sorted(_FACTORIES))

View file

@ -0,0 +1,170 @@
"""VICIdial implementation of external-PBX call and lead operations."""
from __future__ import annotations
import re
from typing import Any, Mapping
import aiohttp
from loguru import logger
from .base import ExternalPBXAdapter, ExternalPBXResult, HeaderReader
_LEAD_FIELD_RE = re.compile(r"^[A-Za-z][A-Za-z0-9_]{0,63}$")
_RESERVED_LEAD_FIELDS = frozenset({"source", "user", "pass", "function", "lead_id"})
class VicidialAdapter(ExternalPBXAdapter):
type = "vicidial"
def __init__(self, config: dict[str, Any]):
agent_api = config.get("agent_api") or {}
non_agent_api = config.get("non_agent_api") or {}
self._agent_url = str(agent_api.get("url", "")).strip()
self._agent_user = str(agent_api.get("username", "")).strip()
self._agent_password = str(agent_api.get("password", ""))
self._agent_source = str(agent_api.get("source", "dograh")).strip()
self._non_agent_url = str(non_agent_api.get("url", "")).strip()
self._non_agent_user = str(non_agent_api.get("username", "")).strip()
self._non_agent_password = str(non_agent_api.get("password", ""))
self._non_agent_source = str(non_agent_api.get("source", "dograh")).strip()
self._timeout = aiohttp.ClientTimeout(
total=min(max(int(config.get("timeout_seconds", 8)), 1), 30)
)
async def capture_call_identity(
self, read_header: HeaderReader
) -> dict[str, str] | None:
callerid = (await read_header("X-VICIDIAL-callerid")).strip()
if not callerid:
return None
return {
"type": self.type,
"callerid": callerid,
"agent_user": (await read_header("X-VICIDIAL-user")).strip(),
"lead_id": (await read_header("X-VICIDIAL-lead_id")).strip(),
"campaign_id": (await read_header("X-VICIDIAL-campaign_id")).strip(),
"ingroup_id": (await read_header("X-VICIDIAL-ingroup_id")).strip(),
}
async def _agent_call_control(
self, identity: Mapping[str, str], stage: str, **extra: str
) -> ExternalPBXResult:
if not all([self._agent_url, self._agent_user, self._agent_password]):
return ExternalPBXResult(
False, stage.lower(), "Agent API is not configured"
)
if not identity.get("callerid") or not identity.get("agent_user"):
return ExternalPBXResult(
False, stage.lower(), "VICIdial call identity is incomplete"
)
params = {
"source": self._agent_source,
"user": self._agent_user,
"pass": self._agent_password,
"agent_user": identity["agent_user"],
"function": "ra_call_control",
"stage": stage,
"value": identity["callerid"],
**extra,
}
try:
async with aiohttp.ClientSession(timeout=self._timeout) as session:
async with session.get(self._agent_url, params=params) as response:
response_text = (await response.text()).strip()
ok = response.status == 200 and response_text.startswith("SUCCESS")
logger.info(
"[VICIdial] ra_call_control completed "
f"stage={stage} status={response.status} ok={ok}"
)
return ExternalPBXResult(
ok,
stage.lower(),
"VICIdial accepted the operation"
if ok
else "VICIdial rejected the operation",
)
except Exception as exc:
logger.error(f"[VICIdial] ra_call_control failed stage={stage}: {exc}")
return ExternalPBXResult(
False, stage.lower(), "VICIdial API request failed"
)
async def hangup(self, identity: Mapping[str, str]) -> ExternalPBXResult:
return await self._agent_call_control(identity, "HANGUP")
async def transfer(
self, identity: Mapping[str, str], destination: str
) -> ExternalPBXResult:
choice = destination.strip()
if choice.lower() == "source":
choice = str(identity.get("ingroup_id", "")).strip()
if not choice:
return ExternalPBXResult(
False, "ingrouptransfer", "No VICIdial in-group was resolved"
)
return await self._agent_call_control(
identity, "INGROUPTRANSFER", ingroup_choices=choice
)
async def update_fields(
self, identity: Mapping[str, str], fields: Mapping[str, str]
) -> ExternalPBXResult:
if not fields:
return ExternalPBXResult(True, "update_lead", "No lead fields configured")
if not all(
[self._non_agent_url, self._non_agent_user, self._non_agent_password]
):
return ExternalPBXResult(
False, "update_lead", "Non-agent API is not configured"
)
lead_id = str(identity.get("lead_id", "")).strip()
if not lead_id:
return ExternalPBXResult(
False, "update_lead", "No VICIdial lead ID captured"
)
safe_fields: dict[str, str] = {}
for key, value in fields.items():
normalized = str(key).strip()
if (
not _LEAD_FIELD_RE.fullmatch(normalized)
or normalized.lower() in _RESERVED_LEAD_FIELDS
):
logger.warning(
f"[VICIdial] Ignoring invalid lead field name: {normalized!r}"
)
continue
safe_fields[normalized] = str(value)
if not safe_fields:
return ExternalPBXResult(
False, "update_lead", "No valid lead fields resolved"
)
params = {
**safe_fields,
"source": self._non_agent_source,
"user": self._non_agent_user,
"pass": self._non_agent_password,
"function": "update_lead",
"lead_id": lead_id,
}
try:
async with aiohttp.ClientSession(timeout=self._timeout) as session:
async with session.get(self._non_agent_url, params=params) as response:
response_text = (await response.text()).strip()
ok = response.status == 200 and response_text.startswith("SUCCESS")
logger.info(
"[VICIdial] update_lead completed "
f"status={response.status} ok={ok} field_count={len(safe_fields)}"
)
return ExternalPBXResult(
ok,
"update_lead",
"VICIdial lead updated" if ok else "VICIdial rejected the lead update",
)
except Exception as exc:
logger.error(f"[VICIdial] update_lead failed: {exc}")
return ExternalPBXResult(
False, "update_lead", "VICIdial API request failed"
)

View file

@ -20,6 +20,7 @@ from api.services.telephony.base import (
NormalizedInboundData,
TelephonyProvider,
)
from api.services.telephony.providers.ari.external_pbx import create_adapter
if TYPE_CHECKING:
from fastapi import WebSocket
@ -51,6 +52,7 @@ class ARIProvider(TelephonyProvider):
self.app_name = config.get("app_name", "")
self.app_password = config.get("app_password", "")
self.from_numbers = config.get("from_numbers", [])
self.external_pbx_adapter = create_adapter(config.get("external_pbx"))
if isinstance(self.from_numbers, str):
self.from_numbers = [self.from_numbers]
@ -363,6 +365,53 @@ class ARIProvider(TelephonyProvider):
"""ARI supports call transfers via bridge manipulation."""
return True
async def transfer_external_pbx_call(
self,
*,
identity: Dict[str, Any],
destination: str,
field_updates: Optional[Dict[str, str]] = None,
) -> Optional[Dict[str, Any]]:
"""Delegate a PBX-owned customer leg to the configured adapter."""
adapter = self.external_pbx_adapter
if adapter is None or not identity:
return None
identity_type = identity.get("type") or identity.get("provider")
if identity_type != adapter.type:
logger.warning(
"[ARI External PBX] Captured identity does not match configured "
f"adapter: identity={identity_type!r} adapter={adapter.type!r}"
)
return {
"status": "failed",
"action": "external_pbx_transfer",
"message": "The external PBX call identity is invalid.",
"reason": "external_pbx_identity_mismatch",
}
update_result = None
if field_updates:
update_result = await adapter.update_fields(identity, field_updates)
if not update_result.ok:
logger.warning(
"[ARI External PBX] Field update failed; continuing transfer "
f"adapter={adapter.type} message={update_result.message}"
)
transfer_result = await adapter.transfer(identity, destination)
return {
"status": "success" if transfer_result.ok else "failed",
"action": "external_pbx_transfer",
"message": (
"Transferring your call now."
if transfer_result.ok
else "I'm sorry, I couldn't complete the transfer."
),
"reason": None if transfer_result.ok else "external_pbx_transfer_failed",
"field_update_ok": update_result.ok if update_result else None,
}
async def transfer_call(
self,
destination: str,

View file

@ -3,11 +3,14 @@
This module contains the business logic for Asterisk ARI call operations.
"""
from typing import Any, Dict
from typing import TYPE_CHECKING, Any, Dict
from loguru import logger
from pipecat.serializers.call_strategies import HangupStrategy, TransferStrategy
if TYPE_CHECKING:
from .external_pbx import ExternalPBXAdapter
class ARIBridgeSwapStrategy(TransferStrategy):
"""Implements bridge swap transfer for Asterisk ARI.
@ -200,6 +203,9 @@ class ARIBridgeSwapStrategy(TransferStrategy):
class ARIHangupStrategy(HangupStrategy):
"""Implements hangup for Asterisk ARI channels."""
def __init__(self, external_pbx_adapter: "ExternalPBXAdapter | None" = None):
self._external_pbx_adapter = external_pbx_adapter
async def execute_hangup(self, context: Dict[str, Any]) -> bool:
"""Hang up the Asterisk channel via ARI REST API."""
try:
@ -217,11 +223,9 @@ class ARIHangupStrategy(HangupStrategy):
)
return False
# If this call came from an upstream PBX (VICIdial), hang up its
# customer leg via its API FIRST -- so the upstream manages its own
# conference/remote-agent teardown cleanly instead of racing dograh's
# SIP BYE -- THEN drop dograh's own leg below.
await self._terminate_upstream_if_any(channel_id)
# The external PBX owns the real customer leg. End it before the
# local ARI leg so its conference teardown cannot race our SIP BYE.
await self._terminate_external_pbx_if_any(channel_id)
endpoint = f"{ari_endpoint}/ari/channels/{channel_id}"
auth = BasicAuth(app_name, app_password)
@ -250,8 +254,8 @@ class ARIHangupStrategy(HangupStrategy):
logger.exception(f"Failed to hang up Asterisk channel: {e}")
return False
async def _terminate_upstream_if_any(self, channel_id: str) -> None:
"""If this run came from an upstream PBX, hang up its customer leg via API.
async def _terminate_external_pbx_if_any(self, channel_id: str) -> None:
"""If configured, hang up the external PBX's customer leg first.
Reuses the same channel->run lookup the transfer strategy uses
(Redis ``ari:channel:{id}`` -> run_id -> ``initial_context``). Best-effort:
@ -262,7 +266,6 @@ class ARIHangupStrategy(HangupStrategy):
from api.constants import REDIS_URL
from api.db import db_client
from api.services.telephony.upstream_pbx import terminate_upstream_call
redis = aioredis.from_url(REDIS_URL, decode_responses=True)
run_id = await redis.get(f"ari:channel:{channel_id}")
@ -271,16 +274,48 @@ class ARIHangupStrategy(HangupStrategy):
run = await db_client.get_workflow_run_by_id(int(run_id))
if not run:
return
upstream = (run.initial_context or {}).get("upstream_pbx")
# If the call was already transferred to the upstream PBX, the customer
identity = (run.initial_context or {}).get("external_pbx_call")
# Read the legacy key for calls created before this refactor.
identity = identity or (run.initial_context or {}).get("upstream_pbx")
# If the call was already transferred to the external PBX, the customer
# leg has moved on -- do NOT hang it up (that would drop the transferred
# customer); just let dograh's own legs tear down below.
transferred = (run.gathered_context or {}).get("upstream_transferred")
if upstream and not transferred:
transferred = (run.gathered_context or {}).get("external_pbx_transferred")
transferred = transferred or (run.gathered_context or {}).get(
"upstream_transferred"
)
if identity and not transferred and self._external_pbx_adapter:
from api.services.telephony.external_pbx import (
resolve_external_pbx_field_mappings,
)
workflow_configurations = (
await db_client.get_workflow_run_configurations(
int(run_id), run.workflow.organization_id
)
)
field_updates = resolve_external_pbx_field_mappings(
run.gathered_context,
workflow_configurations.get("external_pbx_field_mappings", []),
)
if field_updates:
update_result = await self._external_pbx_adapter.update_fields(
identity, field_updates
)
if not update_result.ok:
logger.warning(
"[ARI Hangup] External PBX field update failed; "
f"continuing hangup: {update_result.message}"
)
logger.info(
f"[ARI Hangup] Upstream PBX call ({upstream.get('provider')}); "
"[ARI Hangup] External PBX call "
f"({self._external_pbx_adapter.type}); "
f"terminating customer leg via API before dropping dograh leg"
)
await terminate_upstream_call(upstream)
result = await self._external_pbx_adapter.hangup(identity)
if not result.ok:
logger.warning(
f"[ARI Hangup] External PBX rejected hangup: {result.message}"
)
except Exception as e:
logger.error(f"[ARI Hangup] upstream terminate check failed: {e}")
logger.error(f"[ARI Hangup] external PBX terminate check failed: {e}")

View file

@ -11,6 +11,7 @@ from api.services.pipecat.audio_mixer import build_audio_out_mixer
from api.services.pipecat.transport_params import realtime_param_overrides
from api.services.telephony.factory import load_credentials_for_transport
from .external_pbx import create_adapter
from .serializers import AsteriskFrameSerializer
from .strategies import ARIBridgeSwapStrategy, ARIHangupStrategy
@ -47,7 +48,9 @@ async def create_transport(
app_name=app_name,
app_password=app_password,
transfer_strategy=ARIBridgeSwapStrategy(),
hangup_strategy=ARIHangupStrategy(),
hangup_strategy=ARIHangupStrategy(
external_pbx_adapter=create_adapter(config.get("external_pbx"))
),
params=AsteriskFrameSerializer.InputParams(
asterisk_sample_rate=audio_config.transport_in_sample_rate,
sample_rate=audio_config.pipeline_sample_rate,

View file

@ -30,6 +30,22 @@ if TYPE_CHECKING:
from api.services.telephony.base import TelephonyProvider
@dataclass(frozen=True)
class ProviderUIOption:
"""One selectable value for a provider configuration field."""
value: str
label: str
@dataclass(frozen=True)
class ProviderUICondition:
"""Display a field only when another form value matches ``equals``."""
field: str
equals: Any
@dataclass(frozen=True)
class ProviderUIField:
"""One form field for the telephony configuration UI.
@ -46,6 +62,10 @@ class ProviderUIField:
sensitive: bool = False # If true, mask when displaying stored value
description: Optional[str] = None
placeholder: Optional[str] = None
options: Optional[List[ProviderUIOption]] = None
visible_when: Optional[ProviderUICondition] = None
section: Optional[str] = None
feature_gate: Optional[str] = None
@dataclass(frozen=True)

View file

@ -1,386 +0,0 @@
"""Upstream-PBX call control.
When a call originates on an upstream PBX (e.g. VICIdial) and is patched into
dograh over a SIP trunk, the *customer's* real call leg lives on the upstream
PBX, not on dograh's Asterisk -- dograh only owns its agent leg (the SIP leg
into Stasis + the externalMedia WebSocket). So when the AI decides to hang up or
transfer, dograh must tell the upstream PBX what to do via its API, and must do
so BEFORE tearing down its own SIP leg -- otherwise the SIP BYE races the
upstream PBX's own conference/remote-agent teardown and can leave it in an
inconsistent state.
The upstream identity (the handle for these API calls) is captured off the
inbound SIP headers in ari_manager and stored on the workflow run's
``initial_context["upstream_pbx"]``. The adapter is selected per call by
``upstream["provider"]``.
This deployment is VICIdial-focused; the FreeSWITCH adapter is retained but
inert unless an upstream tags itself ``freeswitch`` (via ``X-PBX-*`` headers).
Connection settings come from environment variables so the same image works
against a PBX on another server (the api container reaches it over normal egress
-- no shared ``pbx-net`` required):
VICIDIAL_API_URL e.g. http://vici.example.com/agc/api.php
VICIDIAL_API_USER VICIdial API user
VICIDIAL_API_PASS VICIdial API password
VICIDIAL_API_SOURCE source tag sent to the API (default: dograh)
VICIDIAL_NON_AGENT_API_URL e.g. http://vici.example.com/vicidial/non_agent_api.php
VICIDIAL_NON_AGENT_API_USER non-agent API user (distinct from the agent API)
VICIDIAL_NON_AGENT_API_PASS non-agent API password
VICIDIAL_NON_AGENT_API_SOURCE source tag sent to the non-agent API (default: dograh)
FREESWITCH_ESL_HOST FreeSWITCH Event Socket host (optional)
FREESWITCH_ESL_PORT FreeSWITCH Event Socket port (default: 8021)
FREESWITCH_ESL_PASSWORD FreeSWITCH ESL password (default: ClueCon)
"""
import asyncio
import os
import aiohttp
from loguru import logger
# --- VICIdial agent-API connection (from env; remote-server friendly) ---
_VICIDIAL_API_URL = os.getenv("VICIDIAL_API_URL", "")
_VICIDIAL_API_USER = os.getenv("VICIDIAL_API_USER", "")
_VICIDIAL_API_PASS = os.getenv("VICIDIAL_API_PASS", "")
_VICIDIAL_API_SOURCE = os.getenv("VICIDIAL_API_SOURCE", "dograh")
# --- VICIdial non-agent API (update_lead etc.; separate endpoint + creds) ---
_VICIDIAL_NON_AGENT_API_URL = os.getenv("VICIDIAL_NON_AGENT_API_URL", "")
_VICIDIAL_NON_AGENT_API_USER = os.getenv("VICIDIAL_NON_AGENT_API_USER", "")
_VICIDIAL_NON_AGENT_API_PASS = os.getenv("VICIDIAL_NON_AGENT_API_PASS", "")
_VICIDIAL_NON_AGENT_API_SOURCE = os.getenv("VICIDIAL_NON_AGENT_API_SOURCE", "dograh")
# Extraction variables whose name starts with this prefix are forwarded to the
# VICIdial non-agent ``update_lead`` API: the prefix is stripped to yield the
# raw lead column name and the extracted value is sent as that column's value.
# e.g. an extraction variable ``X-VICI-UPDATE-LEAD_address3`` with value ``Y``
# becomes ``address3=Y`` on the lead. This lets a workflow plumb arbitrary,
# conversation-derived fields into the VICIdial flow without code changes (see
# ``collect_update_lead_fields``).
UPDATE_LEAD_VAR_PREFIX = "X-VICI-UPDATE-LEAD_"
# API-control params that must never be overridden by a forwarded field -- a
# variable named e.g. ``X-VICI-UPDATE-LEAD_function`` would otherwise hijack the
# update_lead call. These are dropped (with a warning) from forwarded fields.
_UPDATE_LEAD_RESERVED_FIELDS = frozenset(
{"source", "user", "pass", "function", "lead_id"}
)
_HTTP_TIMEOUT = aiohttp.ClientTimeout(total=8)
# --- FreeSWITCH ESL connection (from env; retained, inert unless used) ---
# FreeSWITCH owns the customer leg; dograh drives hangup/transfer over the Event
# Socket Library by the channel UUID captured from the X-PBX-UUID header.
_FS_ESL_HOST = os.getenv("FREESWITCH_ESL_HOST", "")
_FS_ESL_PORT = int(os.getenv("FREESWITCH_ESL_PORT", "8021"))
_FS_ESL_PASSWORD = os.getenv("FREESWITCH_ESL_PASSWORD", "ClueCon")
_FS_ESL_TIMEOUT = 8
async def _ra_call_control(upstream: dict, stage: str, **extra) -> bool:
"""Invoke VICIdial's agent API ``ra_call_control`` for the captured RA call.
The call is identified by ``value`` (the VICIdial callerid captured from the
``X-VICIDIAL-callerid`` header) plus the remote-agent ``agent_user``.
"""
if not _VICIDIAL_API_URL:
logger.warning(
"[upstream_pbx] VICIDIAL_API_URL not configured — cannot drive "
f"VICIdial {stage}"
)
return False
params = {
"source": _VICIDIAL_API_SOURCE,
"user": _VICIDIAL_API_USER,
"pass": _VICIDIAL_API_PASS,
"agent_user": upstream.get("agent_user", ""),
"function": "ra_call_control",
"stage": stage,
"value": upstream.get("callerid", ""),
**extra,
}
try:
async with aiohttp.ClientSession() as session:
async with session.get(
_VICIDIAL_API_URL, params=params, timeout=_HTTP_TIMEOUT
) as resp:
text = (await resp.text()).strip()
ok = text.startswith("SUCCESS")
logger.info(
f"[upstream_pbx] VICIdial ra_call_control {stage} "
f"(agent_user={params['agent_user']}, value={params['value']}) -> {text}"
)
return ok
except Exception as e:
logger.error(f"[upstream_pbx] VICIdial ra_call_control {stage} failed: {e}")
return False
async def _non_agent_update_lead(lead_id: str, **fields) -> bool:
"""Invoke VICIdial's non-agent API ``update_lead`` for one lead.
Uses the dedicated non-agent API endpoint/credentials (distinct from the
agent API). ``fields`` are passed straight through as query params, e.g.
``address3="Y"``.
"""
if not _VICIDIAL_NON_AGENT_API_URL:
logger.warning(
"[upstream_pbx] VICIDIAL_NON_AGENT_API_URL not configured — cannot "
"update_lead"
)
return False
if not lead_id:
logger.warning(
"[upstream_pbx] update_lead requested but no lead_id captured — skipping"
)
return False
# ``fields`` is spread first so the API-control params below always win even
# if a forwarded field collides with one of them (defense in depth; the
# collector also drops reserved names).
params = {
**fields,
"source": _VICIDIAL_NON_AGENT_API_SOURCE,
"user": _VICIDIAL_NON_AGENT_API_USER,
"pass": _VICIDIAL_NON_AGENT_API_PASS,
"function": "update_lead",
"lead_id": lead_id,
}
try:
async with aiohttp.ClientSession() as session:
async with session.get(
_VICIDIAL_NON_AGENT_API_URL, params=params, timeout=_HTTP_TIMEOUT
) as resp:
text = (await resp.text()).strip()
ok = text.startswith("SUCCESS")
logger.info(
f"[upstream_pbx] VICIdial update_lead (lead_id={lead_id}, "
f"fields={fields}) -> {text}"
)
return ok
except Exception as e:
logger.error(f"[upstream_pbx] VICIdial update_lead failed: {e}")
return False
def collect_update_lead_fields(gathered_context: dict) -> dict:
"""Map ``X-VICI-UPDATE-LEAD_<field>`` extracted variables to update_lead fields.
Scans a workflow run's gathered context (its ``extracted_variables`` map) for
variables named with the :data:`UPDATE_LEAD_VAR_PREFIX` prefix and returns
``{<field>: <value>}`` for each one that has a non-empty value. The prefix is
stripped to yield the raw VICIdial lead column (e.g.
``X-VICI-UPDATE-LEAD_address3`` -> ``address3``).
Empty/None values are skipped so we never blank out an existing lead column,
and reserved API-control params are dropped so a stray variable name cannot
hijack the update_lead request.
"""
if not gathered_context:
return {}
extracted = gathered_context.get("extracted_variables")
if not isinstance(extracted, dict):
return {}
fields: dict[str, str] = {}
for key, value in extracted.items():
if not isinstance(key, str) or not key.startswith(UPDATE_LEAD_VAR_PREFIX):
continue
field = key[len(UPDATE_LEAD_VAR_PREFIX) :].strip()
if not field or value is None:
continue
text = str(value).strip()
if not text:
continue
if field in _UPDATE_LEAD_RESERVED_FIELDS:
logger.warning(
f"[upstream_pbx] Ignoring reserved update_lead field '{field}' "
f"from variable '{key}'"
)
continue
fields[field] = text
return fields
async def update_upstream_lead(upstream: dict, fields: dict) -> bool:
"""Update the upstream lead with ``fields`` before a transfer.
Dispatched by provider; currently only VICIdial (via the non-agent
``update_lead`` API). ``fields`` maps VICIdial lead column -> value, e.g.
``{"address3": "Y"}`` (typically built by
:func:`collect_update_lead_fields` from the run's extracted variables).
Best-effort: never blocks the transfer if it fails.
"""
if not upstream or not fields:
return False
if upstream.get("provider") == "vicidial":
return await _non_agent_update_lead(upstream.get("lead_id", ""), **fields)
return False
async def _fs_esl_api(command: str) -> tuple[bool, str]:
"""Run a FreeSWITCH ``api`` command over the Event Socket (inbound mode).
Connects, authenticates, issues ``api <command>`` and returns
``(ok, response_body)`` where ok is True when FreeSWITCH replied ``+OK``.
"""
if not _FS_ESL_HOST:
logger.warning("[upstream_pbx] FREESWITCH_ESL_HOST not configured")
return False, ""
async def _run() -> tuple[bool, str]:
reader, writer = await asyncio.open_connection(_FS_ESL_HOST, _FS_ESL_PORT)
try:
await reader.readuntil(b"\n\n") # "Content-Type: auth/request"
writer.write(f"auth {_FS_ESL_PASSWORD}\n\n".encode())
await writer.drain()
await reader.readuntil(b"\n\n") # auth command/reply
writer.write(f"api {command}\n\n".encode())
await writer.drain()
headers = (await reader.readuntil(b"\n\n")).decode(errors="replace")
length = 0
for line in headers.splitlines():
if line.lower().startswith("content-length:"):
length = int(line.split(":", 1)[1].strip())
body = (
(await reader.readexactly(length)).decode(errors="replace")
if length
else ""
)
return body.startswith("+OK"), body.strip()
finally:
writer.close()
try:
await writer.wait_closed()
except Exception:
pass
try:
return await asyncio.wait_for(_run(), timeout=_FS_ESL_TIMEOUT)
except Exception as e:
logger.error(f"[upstream_pbx] FreeSWITCH ESL '{command}' failed: {e}")
return False, ""
async def _fs_uuid_kill(upstream: dict) -> bool:
uuid = upstream.get("uuid", "")
if not uuid:
return False
ok, resp = await _fs_esl_api(f"uuid_kill {uuid}")
logger.info(f"[upstream_pbx] FreeSWITCH uuid_kill {uuid} -> {resp or ok}")
return ok
async def _fs_uuid_transfer(upstream: dict, destination: str) -> bool:
uuid = upstream.get("uuid", "")
if not uuid:
return False
number = destination.split("/")[-1]
# The customer leg is transferred into the FS dialplan extension
# dograh_xfer_<number>, which bridges to that registered user/agent.
ok, resp = await _fs_esl_api(
f"uuid_transfer {uuid} dograh_xfer_{number} XML dograh-customer"
)
logger.info(
f"[upstream_pbx] FreeSWITCH uuid_transfer {uuid} -> {number}: {resp or ok}"
)
return ok
async def terminate_upstream_call(upstream: dict) -> bool:
"""Hang up the upstream PBX's customer leg. Call this BEFORE dropping dograh's leg."""
if not upstream:
return False
provider = upstream.get("provider")
if provider == "vicidial":
return await _ra_call_control(upstream, "HANGUP")
if provider == "freeswitch":
return await _fs_uuid_kill(upstream)
return False
async def transfer_upstream_call(upstream: dict, destination: str) -> bool:
"""Transfer the upstream PBX's customer leg, dispatched by provider.
VICIdial: always INGROUPTRANSFER (these upstream customers are bounced back
to a queue/agent group, never a bare extension). An explicit ``ingroup:<id>``
destination picks that in-group; anything else (including a plain
extension/number) falls back to the in-group the call arrived on (captured
from the ``X-VICIDIAL-ingroup_id`` header).
FreeSWITCH: uuid_transfer the customer leg to the FS dialplan extension that
bridges to the target. (Both tolerate a leading ``PJSIP/`` in destination.)
"""
if not upstream:
return False
provider = upstream.get("provider")
if provider == "vicidial":
# An explicit "ingroup:<id>" destination names the in-group; everything
# else defaults to the in-group the call arrived on.
choice = ""
if destination.startswith("ingroup"):
_, _, choice = destination.partition(":")
choice = choice.strip()
if not choice or choice == "source":
choice = upstream.get("ingroup_id", "")
if not choice:
logger.warning(
"[upstream_pbx] VICIdial INGROUPTRANSFER requested but no in-group "
f"id available (destination={destination!r}, captured ingroup_id "
"is empty) -- not transferring"
)
return False
return await _ra_call_control(
upstream, "INGROUPTRANSFER", ingroup_choices=choice
)
if provider == "freeswitch":
return await _fs_uuid_transfer(upstream, destination)
return False
# --- Hardcoded post-conversation routing (VICIdial "address3" disposition) ---
# The workflow extracts ``X-VICI-UPDATE-LEAD_address3``; its final value decides
# where the customer is sent once the AI conversation ends:
# "Y" -> INGROUPTRANSFER into in-group "dograhtest1"
# "N" -> INGROUPTRANSFER into in-group "dograhtest2"
# anything else (including a missing/blank value) -> do NOT transfer; the
# customer leg is hung up instead.
# Matched case-insensitively on the stripped value.
ADDRESS3_INGROUP_ROUTES = {
"Y": "dograhtest1",
"N": "dograhtest2",
}
async def route_upstream_after_call(upstream: dict, fields: dict) -> tuple[str, bool]:
"""Dispatch the upstream customer leg from the extracted ``address3`` value.
Hardcoded business routing for VICIdial (see :data:`ADDRESS3_INGROUP_ROUTES`):
an ``address3`` of "Y"/"N" bounces the customer into in-group
``dograhtest1``/``dograhtest2`` respectively; any other value -- including a
missing one -- is treated as "no transfer" and the customer leg is hung up.
``fields`` is the ``{lead_column: value}`` map built by
:func:`collect_update_lead_fields` from the run's extracted variables.
Returns ``(action, ok)`` where ``action`` is ``"transfer"`` or ``"hangup"``
(so the caller can tear down dograh's own leg appropriately) and ``ok`` is the
upstream API result for that action.
"""
raw = (fields or {}).get("address3", "")
address3 = str(raw).strip().upper()
ingroup = ADDRESS3_INGROUP_ROUTES.get(address3)
if ingroup:
logger.info(
f"[upstream_pbx] address3={raw!r} -> INGROUPTRANSFER to in-group "
f"'{ingroup}'"
)
ok = await transfer_upstream_call(upstream, f"ingroup:{ingroup}")
return "transfer", ok
logger.info(
f"[upstream_pbx] address3={raw!r} is not a routable disposition "
"(expected Y or N) -- not transferring; hanging up the customer leg"
)
ok = await terminate_upstream_call(upstream)
return "hangup", ok