diff --git a/api/services/telephony/ari_manager.py b/api/services/telephony/ari_manager.py index cc0be6b6..34d5cbe0 100644 --- a/api/services/telephony/ari_manager.py +++ b/api/services/telephony/ari_manager.py @@ -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, diff --git a/api/services/telephony/upstream_pbx.py b/api/services/telephony/upstream_pbx.py index bb686927..284d546e 100644 --- a/api/services/telephony/upstream_pbx.py +++ b/api/services/telephony/upstream_pbx.py @@ -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_`` 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 + ``{: }`` 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 `` 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:`` -> 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:`` + 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:" 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 diff --git a/api/services/workflow/pipecat_engine.py b/api/services/workflow/pipecat_engine.py index a0d67947..2d41c1b7 100644 --- a/api/services/workflow/pipecat_engine.py +++ b/api/services/workflow/pipecat_engine.py @@ -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) diff --git a/api/services/workflow/pipecat_engine_custom_tools.py b/api/services/workflow/pipecat_engine_custom_tools.py index 6a75ac88..87ce5e76 100644 --- a/api/services/workflow/pipecat_engine_custom_tools.py +++ b/api/services/workflow/pipecat_engine_custom_tools.py @@ -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,