feat(telephony): env-based VICIdial config + address3 post-call routing

Move VICIdial agent/non-agent API and FreeSWITCH ESL connection settings
out of hardcoded POC values into environment variables so the same image
works against a PBX on another server.

Add VICIdial update_lead forwarding: X-VICI-UPDATE-LEAD_* extracted
variables are mapped to lead columns and pushed via the non-agent API
before a transfer, with reserved API-control params dropped.

Add hardcoded address3 disposition routing (Y/N -> in-group transfer,
else hang up) and a synchronous final variable extraction so the
transfer path sees the freshest conversation state. Skip upstream_pbx
capture on non-PJSIP channels to avoid Asterisk 500s.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Abhishek 2026-07-20 08:22:54 +00:00
parent 688773d3b3
commit c8d7584f9f
4 changed files with 352 additions and 57 deletions

View file

@ -439,7 +439,9 @@ class ARIConnection:
)
return (result or {}).get("value", "") or ""
async def _capture_upstream_pbx(self, channel_id: str) -> Optional[dict]:
async def _capture_upstream_pbx(
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
@ -451,6 +453,17 @@ class ARIConnection:
callerid + remote-agent user, driven over ``ra_call_control``.
Returns None for non-upstream (direct) calls.
"""
# 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"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(
@ -493,6 +506,11 @@ class ARIConnection:
"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 "
@ -637,7 +655,9 @@ class ARIConnection:
# 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(channel_id)
upstream_pbx = await self._capture_upstream_pbx(
channel_id, channel.get("name", "")
)
workflow_run = await db_client.create_workflow_run(
name=f"ARI Inbound {caller_number}",
workflow_id=inbound_workflow_id,

View file

@ -1,4 +1,4 @@
"""Upstream-PBX call control (POC).
"""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
@ -11,32 +11,72 @@ 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"]``.
``initial_context["upstream_pbx"]``. The adapter is selected per call by
``upstream["provider"]``.
POC scope: the VICIdial agent-API connection is hardcoded below. The productized
version moves these into the ARI telephony-configuration credentials and selects
the adapter by ``upstream["provider"]`` (see the upstream-PBX seam design doc).
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
# --- POC hardcoded VICIdial agent-API connection (refine: telephony config creds) ---
_VICIDIAL_API_URL = "http://10.10.10.15/agc/api.php"
_VICIDIAL_API_USER = "6666"
_VICIDIAL_API_PASS = "1234"
_VICIDIAL_API_SOURCE = "dograh"
# --- 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)
# --- POC hardcoded FreeSWITCH ESL connection (refine: telephony config creds) ---
# --- 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 = "10.10.10.17"
_FS_ESL_PORT = 8021
_FS_ESL_PASSWORD = "ClueCon"
_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
@ -46,6 +86,12 @@ async def _ra_call_control(upstream: dict, stage: str, **extra) -> bool:
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,
@ -73,12 +119,116 @@ async def _ra_call_control(upstream: dict, stage: str, **extra) -> bool:
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)
@ -154,8 +304,11 @@ async def terminate_upstream_call(upstream: dict) -> bool:
async def transfer_upstream_call(upstream: dict, destination: str) -> bool:
"""Transfer the upstream PBX's customer leg, dispatched by provider.
VICIdial: ``ingroup:<id>`` -> INGROUPTRANSFER (to a queue/agent group);
anything else -> EXTENSIONTRANSFER to that number/extension.
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.)
"""
@ -163,14 +316,71 @@ async def transfer_upstream_call(upstream: dict, destination: str) -> bool:
return False
provider = upstream.get("provider")
if provider == "vicidial":
if destination.startswith("ingroup:"):
return await _ra_call_control(
upstream, "INGROUPTRANSFER", ingroup_choices=destination.split(":", 1)[1]
# 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"
)
number = destination.split("/")[-1]
return False
return await _ra_call_control(
upstream, "EXTENSIONTRANSFER", phone_number=number
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

View file

@ -99,6 +99,10 @@ class PipecatEngine:
self._gathered_context: dict = {}
self._user_response_timeout_task: Optional[asyncio.Task] = None
self._pending_extraction_tasks: set[asyncio.Task] = set()
# True once a final (synchronous) extraction has run, so the end-of-call
# and upstream-transfer paths don't redundantly re-extract the same
# terminal state.
self._final_extraction_done: bool = False
# Will be set later in initialize() when we have
# access to _context
@ -504,6 +508,29 @@ class PipecatEngine:
f"Incomplete: {incomplete}"
)
async def perform_final_variable_extraction(self) -> None:
"""Flush in-flight + current-node variable extraction synchronously.
Awaits any background extractions still running from previous nodes,
then runs the current node's extraction inline so callers that need the
freshest extracted variables before acting can rely on them -- e.g.
end_call_with_reason before disposing the call, or an upstream-PBX
transfer that maps extracted variables into the VICIdial update_lead
call before handing the customer off.
Idempotent: only the first call does work. The upstream-PBX transfer
runs this just before forwarding update_lead, so the subsequent
end_call_with_reason would otherwise re-extract the same terminal state.
"""
if self._final_extraction_done:
logger.debug("Final variable extraction already performed; skipping")
return
self._final_extraction_done = True
await self._await_pending_extractions()
await self._perform_variable_extraction_if_needed(
self._current_node, run_in_background=False
)
async def _setup_llm_context(self, node: Node) -> None:
"""Common method to set up LLM context"""
# Set OTel span name for tracing
@ -739,13 +766,8 @@ class PipecatEngine:
EndTaskReason.PIPELINE_ERROR.value,
EndTaskReason.VOICEMAIL_DETECTED.value,
):
# Await any in-flight background extractions from previous nodes
await self._await_pending_extractions()
# Perform final variable extraction synchronously before ending
await self._perform_variable_extraction_if_needed(
self._current_node, run_in_background=False
)
# Flush in-flight + current-node extractions synchronously before ending
await self.perform_final_variable_extraction()
frame_to_push = (
CancelFrame(reason=reason) if abort_immediately else EndFrame(reason=reason)

View file

@ -561,38 +561,81 @@ class CustomToolManager:
upstream = (workflow_run.initial_context or {}).get("upstream_pbx")
if upstream:
from api.services.telephony.upstream_pbx import (
transfer_upstream_call,
collect_update_lead_fields,
route_upstream_after_call,
update_upstream_lead,
)
# Play the transfer message first so its TTS overlaps the
# extraction latency below.
await self._play_config_message(config)
ok = await transfer_upstream_call(upstream, destination)
# Mark the run so the end-of-call hangup does NOT also hang up
# the now-transferred customer (see ARIHangupStrategy).
await db_client.update_workflow_run(
run_id=self._engine._workflow_run_id,
gathered_context={"upstream_transferred": True},
# Extract variables right before the transfer so any
# X-VICI-UPDATE-LEAD_* values reflect the final conversation
# state, then forward them to the upstream PBX's update_lead
# BEFORE handing the customer off. This lets a workflow plumb
# arbitrary, conversation-derived fields into the VICIdial
# lead. Best-effort: never blocks the transfer.
await self._engine.perform_final_variable_extraction()
update_fields = collect_update_lead_fields(
self._engine._gathered_context
)
await function_call_params.result_callback(
{
"status": "success" if ok else "failed",
"action": "upstream_transfer",
"message": (
"Transferring your call now."
if ok
else "I'm sorry, I couldn't complete the transfer."
),
},
properties=properties,
logger.info(
f"[transfer] update_lead fields from extracted "
f"variables: {update_fields}"
)
if ok:
# Give the upstream PBX time to redirect the customer out
# of the conference toward the destination BEFORE we drop
# dograh's own leg -- otherwise the teardown races the
# redirect and cancels the customer's call to the
# destination (it never rings).
await asyncio.sleep(4)
# Tear down dograh's own legs; the customer has moved to the
# upstream PBX side.
if update_fields:
await update_upstream_lead(upstream, update_fields)
# Hardcoded address3 routing (see route_upstream_after_call):
# "Y"/"N" transfer the customer into in-group
# dograhtest1/dograhtest2; anything else hangs the customer up
# instead of transferring. The destination configured on the
# transfer node is intentionally ignored for upstream calls.
action, ok = await route_upstream_after_call(
upstream, update_fields
)
if action == "transfer":
# Mark the run so the end-of-call hangup does NOT also hang
# up the now-transferred customer (see ARIHangupStrategy).
await db_client.update_workflow_run(
run_id=self._engine._workflow_run_id,
gathered_context={"upstream_transferred": True},
)
await function_call_params.result_callback(
{
"status": "success" if ok else "failed",
"action": "upstream_transfer",
"message": (
"Transferring your call now."
if ok
else "I'm sorry, I couldn't complete the transfer."
),
},
properties=properties,
)
if ok:
# Give the upstream PBX time to redirect the customer
# out of the conference toward the in-group BEFORE we
# drop dograh's own leg -- otherwise the teardown races
# the redirect and cancels the customer's call (it
# never rings).
await asyncio.sleep(4)
else:
# address3 wasn't Y/N: route_upstream_after_call already
# hung up the customer leg. Leave upstream_transferred
# unset so the end-of-call teardown re-confirms the hangup
# (a no-op if it already succeeded, a retry if it didn't).
await function_call_params.result_callback(
{
"status": "success",
"action": "upstream_hangup",
"message": "Ending the call now.",
},
properties=properties,
)
# Tear down dograh's own legs; the customer leg has already
# been handled on the upstream PBX side.
await self._engine.end_call_with_reason(
EndTaskReason.END_CALL_TOOL_REASON.value,
abort_immediately=True,