From 8e907a6a9513fedf2b3d65f81c52892c621a9bb7 Mon Sep 17 00:00:00 2001 From: arys Date: Sat, 9 May 2026 09:25:29 +0500 Subject: [PATCH] sales push --- .env.example | 11 +++ .idea/.gitignore | 10 +++ .idea/checkstyle-idea.xml | 16 ++++ .idea/git_toolbox_prj.xml | 15 ++++ .idea/material_theme_project_new.xml | 12 +++ .idea/misc.xml | 6 ++ .idea/modules.xml | 8 ++ .idea/realtime_voice_service.iml | 9 ++ .idea/vcs.xml | 6 ++ audioop_compat.py | 110 +++++++++++++++++++++++ core/filler_audio.py | 2 +- core/session.py | 8 ++ main.py | 106 +++++++++++++++++++--- providers/stt.py | 2 +- providers/tts.py | 2 +- sales_sync_client.py | 126 +++++++++++++++++++++++++++ transports/audiosocket.py | 2 +- 17 files changed, 436 insertions(+), 15 deletions(-) create mode 100644 .idea/.gitignore create mode 100644 .idea/checkstyle-idea.xml create mode 100644 .idea/git_toolbox_prj.xml create mode 100644 .idea/material_theme_project_new.xml create mode 100644 .idea/misc.xml create mode 100644 .idea/modules.xml create mode 100644 .idea/realtime_voice_service.iml create mode 100644 .idea/vcs.xml create mode 100644 audioop_compat.py create mode 100644 sales_sync_client.py diff --git a/.env.example b/.env.example index 86e4afe..ff73675 100644 --- a/.env.example +++ b/.env.example @@ -80,3 +80,14 @@ SEMANTIC_ENDPOINTING_HOLD_MS=600 REALTIME_VOICE_FILLER_DELAY_MS=300 REALTIME_VOICE_TTS_CHUNK_SOFT_MIN_CHARS=10 REALTIME_VOICE_TTS_CHUNK_SOFT_MIN_WORDS=2 + +CRM_INTERACTION_ENABLED=false +CRM_INTERACTION_SERVICE_URL=http://interaction-service:8000 +CRM_APP_TOKEN_SECRET= +CRM_QUEUE_ID= + +SALES_SYNC_ENABLED=true +SALES_SERVICE_URL=http://sales-service:8000 +SALES_SYNC_TIMEOUT_SECONDS=8 +REALTIME_VOICE_QUEUE_ID= +REALTIME_VOICE_QUEUE_CODE= diff --git a/.idea/.gitignore b/.idea/.gitignore new file mode 100644 index 0000000..30cf57e --- /dev/null +++ b/.idea/.gitignore @@ -0,0 +1,10 @@ +# Default ignored files +/shelf/ +/workspace.xml +# Editor-based HTTP Client requests +/httpRequests/ +# Ignored default folder with query files +/queries/ +# Datasource local storage ignored files +/dataSources/ +/dataSources.local.xml diff --git a/.idea/checkstyle-idea.xml b/.idea/checkstyle-idea.xml new file mode 100644 index 0000000..ae8e08c --- /dev/null +++ b/.idea/checkstyle-idea.xml @@ -0,0 +1,16 @@ + + + + 13.4.0 + JavaOnly + true + + + \ No newline at end of file diff --git a/.idea/git_toolbox_prj.xml b/.idea/git_toolbox_prj.xml new file mode 100644 index 0000000..02b915b --- /dev/null +++ b/.idea/git_toolbox_prj.xml @@ -0,0 +1,15 @@ + + + + + + + \ No newline at end of file diff --git a/.idea/material_theme_project_new.xml b/.idea/material_theme_project_new.xml new file mode 100644 index 0000000..bfc5c67 --- /dev/null +++ b/.idea/material_theme_project_new.xml @@ -0,0 +1,12 @@ + + + + + + \ No newline at end of file diff --git a/.idea/misc.xml b/.idea/misc.xml new file mode 100644 index 0000000..31e1ebc --- /dev/null +++ b/.idea/misc.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/.idea/modules.xml b/.idea/modules.xml new file mode 100644 index 0000000..2f9d213 --- /dev/null +++ b/.idea/modules.xml @@ -0,0 +1,8 @@ + + + + + + + + \ No newline at end of file diff --git a/.idea/realtime_voice_service.iml b/.idea/realtime_voice_service.iml new file mode 100644 index 0000000..d6ebd48 --- /dev/null +++ b/.idea/realtime_voice_service.iml @@ -0,0 +1,9 @@ + + + + + + + + + \ No newline at end of file diff --git a/.idea/vcs.xml b/.idea/vcs.xml new file mode 100644 index 0000000..35eb1dd --- /dev/null +++ b/.idea/vcs.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/audioop_compat.py b/audioop_compat.py new file mode 100644 index 0000000..ef99ee9 --- /dev/null +++ b/audioop_compat.py @@ -0,0 +1,110 @@ +from __future__ import annotations + +import math +import struct +from array import array + +try: + import audioop as _audioop +except ModuleNotFoundError: + try: + import audioop_lts as _audioop # type: ignore[import-not-found] + except ModuleNotFoundError: + _audioop = None + + +def _native_is_little_endian() -> bool: + return struct.pack("=h", 1) == struct.pack(" array: + if width != 2: + raise NotImplementedError("audioop_compat fallback currently supports only 16-bit PCM") + samples = array("h") + samples.frombytes(fragment) + if not _native_is_little_endian(): + samples.byteswap() + return samples + + +def _write_int16_samples(samples: array) -> bytes: + output = array("h", samples) + if not _native_is_little_endian(): + output.byteswap() + return output.tobytes() + + +def _clip_int16(value: float) -> int: + return max(-32768, min(32767, int(round(value)))) + + +class _AudioopFallback: + @staticmethod + def rms(fragment: bytes, width: int) -> int: + samples = _read_int16_samples(fragment, width) + if not samples: + return 0 + mean_square = sum(sample * sample for sample in samples) / len(samples) + return int(math.sqrt(mean_square)) + + @staticmethod + def max(fragment: bytes, width: int) -> int: + samples = _read_int16_samples(fragment, width) + if not samples: + return 0 + return max(abs(sample) for sample in samples) + + @staticmethod + def tomono(fragment: bytes, width: int, lfactor: float, rfactor: float) -> bytes: + samples = _read_int16_samples(fragment, width) + if len(samples) % 2 != 0: + raise ValueError("Stereo PCM must contain an even number of samples") + mono = array("h") + for index in range(0, len(samples), 2): + left = samples[index] + right = samples[index + 1] + mono.append(_clip_int16((left * lfactor) + (right * rfactor))) + return _write_int16_samples(mono) + + @staticmethod + def ratecv( + fragment: bytes, + width: int, + nchannels: int, + inrate: int, + outrate: int, + state, + weightA: int = 1, + weightB: int = 0, + ) -> tuple[bytes, None]: + del state, weightA, weightB + if nchannels <= 0: + raise ValueError("nchannels must be positive") + if inrate <= 0 or outrate <= 0: + raise ValueError("Sample rates must be positive") + samples = _read_int16_samples(fragment, width) + if not samples or inrate == outrate: + return fragment, None + if len(samples) % nchannels != 0: + raise ValueError("PCM fragment size does not match channel count") + + frame_count = len(samples) // nchannels + output_frame_count = max(1, int(round(frame_count * outrate / inrate))) + output = array("h") + + for out_index in range(output_frame_count): + position = out_index * inrate / outrate + left_index = min(int(position), frame_count - 1) + right_index = min(left_index + 1, frame_count - 1) + fraction = max(0.0, min(1.0, position - left_index)) + for channel_index in range(nchannels): + left_sample = samples[(left_index * nchannels) + channel_index] + right_sample = samples[(right_index * nchannels) + channel_index] + interpolated = left_sample + ((right_sample - left_sample) * fraction) + output.append(_clip_int16(interpolated)) + return _write_int16_samples(output), None + + +audioop = _audioop or _AudioopFallback() + +__all__ = ["audioop"] diff --git a/core/filler_audio.py b/core/filler_audio.py index ff91e16..bc287dd 100644 --- a/core/filler_audio.py +++ b/core/filler_audio.py @@ -1,7 +1,6 @@ from __future__ import annotations import asyncio -import audioop import io import logging import os @@ -10,6 +9,7 @@ import wave from collections.abc import AsyncIterable from collections.abc import Iterable +from realtime_voice_service.audioop_compat import audioop from realtime_voice_service.providers.base import BaseTTS diff --git a/core/session.py b/core/session.py index f3452e4..dc809b9 100644 --- a/core/session.py +++ b/core/session.py @@ -1214,6 +1214,14 @@ class CallSession: def conversation(self) -> tuple[tuple[str, str], ...]: return tuple(self._conversation) + @property + def customer_name(self) -> str | None: + return self._customer_name + + @property + def session_language(self) -> str | None: + return self._session_language + async def run(self) -> None: if self._assistant_task is not None: raise RuntimeError("CallSession.run() can only be called once per session") diff --git a/main.py b/main.py index dcc08bc..6328309 100644 --- a/main.py +++ b/main.py @@ -6,10 +6,12 @@ import logging import os import uuid from contextlib import asynccontextmanager +from datetime import datetime, timezone from fastapi import FastAPI, WebSocket from realtime_voice_service import crm_client +from realtime_voice_service import sales_sync_client from realtime_voice_service.core.filler_audio import FillerAudioLibrary from realtime_voice_service.core.session import CallSession from realtime_voice_service.core.session import _normalize_voice_pronunciation @@ -148,6 +150,52 @@ def _llm_model_name() -> str: return str(os.getenv("OPENAI_LLM_MODEL", "gpt-4o-mini")).strip() or "gpt-4o-mini" +def _utc_now_iso() -> str: + return datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z") + + +def _conversation_transcript(session: CallSession) -> str | None: + lines = [f"{speaker}: {text.strip()}" for speaker, text in session.conversation if str(text).strip()] + if not lines: + return None + return "\n".join(lines) + + +def _conversation_summary(session: CallSession) -> str | None: + user_turns = [text.strip() for speaker, text in session.conversation if speaker == "user" and str(text).strip()] + assistant_turns = [text.strip() for speaker, text in session.conversation if speaker == "assistant" and str(text).strip()] + if not user_turns and not assistant_turns: + return None + parts: list[str] = [] + if user_turns: + parts.append(f"Customer asked: {user_turns[0][:220]}") + if assistant_turns: + parts.append(f"Assistant response: {assistant_turns[-1][:220]}") + if session.customer_name: + parts.append(f"Customer name: {session.customer_name}") + return " | ".join(parts)[:700] + + +def _session_sync_metadata( + session: CallSession, + transport: BaseMediaTransport, + *, + failure_reason: str | None = None, +) -> dict[str, object]: + return { + "protocol": transport.protocol, + "sample_rate_hz": transport.sample_rate_hz, + "frame_duration_ms": transport.frame_duration_ms, + "frame_bytes": transport.frame_bytes, + "session_language": session.session_language, + "customer_name": session.customer_name, + "interruptions": list(session.interruptions), + "latency_ms": dict(session.last_latency_ms), + "conversation_entries": len(session.conversation), + "failure_reason": failure_reason, + } + + class RealtimeVoiceService: def __init__(self) -> None: @@ -276,19 +324,53 @@ class RealtimeVoiceService: async def _run_transport_session(self, transport: BaseMediaTransport) -> None: session = self._build_session(transport) + started_at = _utc_now_iso() interaction_id = await crm_client.create_interaction(call_id=session.session_id) + await sales_sync_client.sync_voice_session( + call_id=session.session_id, + interaction_id=interaction_id, + voice_session_id=session.session_id, + ai_state=str(session.state.value).lower(), + telephony_status="connected", + call_status="started", + started_at=started_at, + summary="Realtime voice session started", + metadata=_session_sync_metadata(session, transport), + ) + failure_reason: str | None = None + call_status = "completed" async with self._track_session(session): - LOGGER.info( - "starting realtime session %s via %s sample_rate=%s frame_ms=%s frame_bytes=%s", - session.session_id, - transport.protocol, - transport.sample_rate_hz, - transport.frame_duration_ms, - transport.frame_bytes, - ) - await session.run() - if interaction_id: - await crm_client.close_interaction(interaction_id) + try: + LOGGER.info( + "starting realtime session %s via %s sample_rate=%s frame_ms=%s frame_bytes=%s", + session.session_id, + transport.protocol, + transport.sample_rate_hz, + transport.frame_duration_ms, + transport.frame_bytes, + ) + await session.run() + except Exception as exc: + call_status = "failed" + failure_reason = str(exc)[:300] + raise + finally: + await sales_sync_client.sync_voice_session( + call_id=session.session_id, + interaction_id=interaction_id, + caller_name=session.customer_name, + voice_session_id=session.session_id, + ai_state=str(session.state.value).lower(), + telephony_status="ended", + call_status=call_status, + started_at=started_at, + ended_at=_utc_now_iso(), + summary=_conversation_summary(session), + transcript_text=_conversation_transcript(session), + metadata=_session_sync_metadata(session, transport, failure_reason=failure_reason), + ) + if interaction_id: + await crm_client.close_interaction(interaction_id) def _build_session(self, transport: BaseMediaTransport) -> CallSession: return CallSession( @@ -364,6 +446,8 @@ async def health() -> dict[str, object]: "sample_rate_hz": service.sample_rate_hz, "active_sessions": service.active_session_count, "audiosocket_port": _audiosocket_port(), + "sales_sync_enabled": os.getenv("SALES_SYNC_ENABLED", "0"), + "crm_interaction_enabled": os.getenv("CRM_INTERACTION_ENABLED", "0"), } diff --git a/providers/stt.py b/providers/stt.py index 8e455c0..754b17f 100644 --- a/providers/stt.py +++ b/providers/stt.py @@ -1,7 +1,6 @@ from __future__ import annotations import asyncio -import audioop import base64 import contextlib import io @@ -13,6 +12,7 @@ import wave from collections.abc import Sequence from urllib.parse import urlencode +from realtime_voice_service.audioop_compat import audioop from realtime_voice_service.providers.base import BaseSTT from realtime_voice_service.providers.base import BaseSTTStream from realtime_voice_service.providers.base import PartialTranscriptCallback diff --git a/providers/tts.py b/providers/tts.py index 001e283..f5311d1 100644 --- a/providers/tts.py +++ b/providers/tts.py @@ -1,7 +1,6 @@ from __future__ import annotations import asyncio -import audioop import base64 import json import logging @@ -12,6 +11,7 @@ from collections.abc import AsyncIterable from urllib.parse import quote from urllib.parse import urlencode +from realtime_voice_service.audioop_compat import audioop from realtime_voice_service.providers.base import BaseTTS diff --git a/sales_sync_client.py b/sales_sync_client.py new file mode 100644 index 0000000..aeb7fcf --- /dev/null +++ b/sales_sync_client.py @@ -0,0 +1,126 @@ +from __future__ import annotations + +import logging +import os +from typing import Any + +import httpx + +from realtime_voice_service.crm_client import _issue_service_token + + +LOGGER = logging.getLogger("uvicorn.error") + + +def _enabled() -> bool: + return os.getenv("SALES_SYNC_ENABLED", "0").strip().lower() in {"1", "true", "yes", "on"} + + +def _base_url() -> str: + return str(os.getenv("SALES_SERVICE_URL", "http://sales-service:8000")).rstrip("/") + + +def _timeout_seconds() -> float: + raw = str(os.getenv("SALES_SYNC_TIMEOUT_SECONDS", "8")).strip() + try: + return max(float(raw), 1.0) + except ValueError: + return 8.0 + + +def _secret() -> str: + return str(os.getenv("CRM_APP_TOKEN_SECRET", "")).strip() + + +def _service_headers() -> dict[str, str] | None: + secret = _secret() + if not secret: + LOGGER.warning("sales sync: CRM_APP_TOKEN_SECRET not configured") + return None + token = _issue_service_token(secret) + return {"Authorization": f"Bearer {token}"} + + +def _queue_id() -> str | None: + raw = str( + os.getenv("REALTIME_VOICE_QUEUE_ID") + or os.getenv("CRM_QUEUE_ID") + or "" + ).strip() + return raw or None + + +def _queue_code() -> str | None: + raw = str(os.getenv("REALTIME_VOICE_QUEUE_CODE", "")).strip() + return raw or None + + +async def sync_voice_session( + *, + call_id: str, + interaction_id: str | None = None, + caller_number: str | None = None, + caller_name: str | None = None, + voice_session_id: str | None = None, + ai_session_id: str | None = None, + ai_state: str | None = None, + handoff_reason: str | None = None, + telephony_status: str | None = None, + call_status: str | None = None, + started_at: str | None = None, + ended_at: str | None = None, + summary: str | None = None, + transcript_text: str | None = None, + metadata: dict[str, Any] | None = None, +) -> dict[str, Any] | None: + if not _enabled(): + return None + headers = _service_headers() + if headers is None: + return None + + payload: dict[str, Any] = { + "call_id": call_id, + "interaction_id": interaction_id, + "queue_id": _queue_id(), + "queue_code": _queue_code(), + "caller_number": caller_number, + "caller_name": caller_name, + "voice_session_id": voice_session_id or call_id, + "ai_session_id": ai_session_id, + "ai_state": ai_state, + "handoff_reason": handoff_reason, + "telephony_status": telephony_status, + "call_status": call_status, + "started_at": started_at, + "ended_at": ended_at, + "summary": summary, + "transcript_text": transcript_text, + "metadata": { + "source_service": "realtime_voice_service", + **(metadata or {}), + }, + } + try: + async with httpx.AsyncClient(timeout=_timeout_seconds()) as client: + response = await client.post( + f"{_base_url()}/internal/sales-sync/voice", + json=payload, + headers=headers, + ) + response.raise_for_status() + data = response.json() + LOGGER.info( + "sales sync: voice session synced call_id=%s voice_session_id=%s status=%s", + call_id, + voice_session_id or call_id, + call_status, + ) + return data if isinstance(data, dict) else None + except Exception: + LOGGER.exception( + "sales sync: failed to sync voice session call_id=%s voice_session_id=%s", + call_id, + voice_session_id or call_id, + ) + return None diff --git a/transports/audiosocket.py b/transports/audiosocket.py index 1f0f7c2..183e96a 100644 --- a/transports/audiosocket.py +++ b/transports/audiosocket.py @@ -1,7 +1,6 @@ from __future__ import annotations import asyncio -import audioop import logging import os import struct @@ -9,6 +8,7 @@ import time import uuid from typing import Awaitable, Callable +from realtime_voice_service.audioop_compat import audioop from realtime_voice_service.transports.base import BaseMediaTransport