from __future__ import annotations from collections import defaultdict from datetime import datetime, timedelta, timezone import json import logging import math import os import re import time from typing import Any import httpx from fastapi import Depends, FastAPI, HTTPException, Query from sqlalchemy import select from services.shared.core import Role, new_id, utc_now_iso from services.shared.db import get_session from services.shared.ai_context_summary import ( dump_context_summary, render_context_summary_text, update_context_summary_from_assistant_turn, update_context_summary_from_user_turn, ) from services.shared.kb_localization import normalize_kb_language from services.shared.kb_search import search_kb_rows from services.shared.models import ( AIWhatsAppEnqueueIn, AIAnalyticsDrilldownFiltersOut, AIAnalyticsDrilldownItemOut, AIAnalyticsDrilldownOut, AIWhatsAppPauseIn, AIAnalyticsBreakdownsOut, AIAnalyticsChannelBreakdownOut, AIAnalyticsCoverageOut, AIAnalyticsFiltersOut, AIAnalyticsHandoffReasonBreakdownOut, AIAnalyticsMetricsOut, AIAnalyticsOutcomeBreakdownOut, AIAnalyticsOverviewOut, AIAnalyticsSessionDetailOut, AIAnalyticsSessionEventOut, AIAnalyticsSessionLinkedInteractionOut, AIAnalyticsTimeseriesOut, AIAnalyticsTimeseriesPointOut, AIAnalyticsTotalsOut, AIAnalyticsWindowOut, AIOperatorConfig, AIOperatorConfigOut, VoiceNameFlowAnalyticsBreakdownsOut, VoiceNameFlowAnalyticsCoverageOut, VoiceNameFlowAnalyticsFiltersOut, VoiceNameFlowAnalyticsFunnelStageOut, VoiceNameFlowAnalyticsHandoffBreakdownOut, VoiceNameFlowAnalyticsLanguageBreakdownOut, VoiceNameFlowAnalyticsMetricsOut, VoiceNameFlowAnalyticsOverviewOut, VoiceNameFlowAnalyticsQueueBreakdownOut, VoiceNameFlowAnalyticsTimeseriesOut, VoiceNameFlowAnalyticsTimeseriesPointOut, VoiceNameFlowAnalyticsTotalsOut, AITelegramEnqueueIn, AITelegramPauseIn, HealthResponse, VoiceNameCollectionConfig, VoiceNameCollectionConfigOut, VoiceTTSConfig, VoiceTTSConfigOut, VoiceAIStartIn, VoiceAIStartOut, VoiceAITurnIn, ) from services.shared.security import issue_app_token, require_roles from services.shared.sql_init import init_sql_schema from services.shared.sql_models import ( AIJobRow, AISessionRow, AITurnRow, AsteriskCallLinkRow, Customer, CustomerExternalIdentity, Interaction, InteractionTimeline, KBArticleRow, WhatsAppMessageRow, WhatsAppThreadRow, TelegramMessageRow, TelegramThreadRow, VoiceAISessionRow, ) from services.ai_orchestrator_service import voice as voice_flows from services.ai_orchestrator_service import operator_persona as persona from services.ai_orchestrator_service.voice_name_config import ( load_voice_name_collection_config, save_voice_name_collection_config, ) from services.shared.ai_operator_config import ( load_ai_operator_config, load_effective_ai_operator_config, save_ai_operator_config, ) from services.shared.voice_tts_config import load_voice_tts_config, save_voice_tts_config app = FastAPI(title="ai-orchestrator-service", version="1.0.0") init_sql_schema() logger = logging.getLogger(__name__) _AI_ANALYTICS_CHANNELS = {"telegram", "whatsapp"} _AI_ANALYTICS_TERMINAL_STATUSES = {"closed", "handoff_required", "human_owned", "error"} _AI_ANALYTICS_METRICS = { "containment_rate", "handoff_rate", "ai_latency_avg_ms", "closed_without_operator_rate", "human_touched_rate", } _AI_ANALYTICS_INTERVALS = {"hour", "day"} _AI_ANALYTICS_SLICES = {"all", "contained", "handoff", "human_touched", "closed_without_operator", "active", "error"} _AI_ANALYTICS_DRILLDOWN_SORT_FIELDS = {"created_at", "updated_at", "ai_latency_avg_ms", "status"} _VOICE_NAME_FLOW_STATUS_VALUES = {"name_obtained", "name_followup_required", "name_not_obtained"} _VOICE_NAME_FLOW_METRICS = { "scenario_calls", "start_capture_rate", "downstream_rescue_rate", "handoff_unconfirmed_rate", "manual_correction_rate", } _AI_ANALYTICS_REASON_LABELS = { "requested_human": "Запрос клиента на оператора", "knowledge_or_tool_gap": "Недостаточно знаний или tools", "policy_or_sensitive": "Policy или чувствительная тема", "delivery_or_runtime_error": "Ошибка доставки или runtime", "manual_claim": "Ручной takeover", "other": "Другая причина", } _AI_ANALYTICS_OUTCOME_LABELS = { "contained": "Containment", "handoff": "Handoff", "human_touched": "Human touched", "closed_without_operator": "Closed without operator", "active": "Active", "error": "Error", } _VOICE_NAME_FLOW_FUNNEL_LABELS = { "scenario_calls": "Звонки в сценарии", "start_obtained": "Имя взято сразу", "needed_downstream": "Потребовался downstream AI", "downstream_ai_obtained": "Имя добрал downstream AI", "handoff_confirmed_name": "Handoff с подтверждённым именем", "handoff_unconfirmed_name": "Handoff без подтверждённого имени", } def _parse_analytics_timestamp(raw: str, field_name: str) -> datetime: normalized = raw.strip() if normalized.endswith("Z"): normalized = f"{normalized[:-1]}+00:00" try: parsed = datetime.fromisoformat(normalized) except ValueError as exc: raise HTTPException(status_code=400, detail=f"Invalid {field_name}") from exc if parsed.tzinfo is None: parsed = parsed.replace(tzinfo=timezone.utc) return parsed.astimezone(timezone.utc) def _normalize_ai_analytics_channel(channel: str | None) -> tuple[str, list[str]]: normalized = (channel or "all").strip().lower() or "all" if normalized == "all": return normalized, ["telegram", "whatsapp"] if normalized in _AI_ANALYTICS_CHANNELS: return normalized, [normalized] return normalized, [] def _normalize_ai_analytics_metric(metric: str) -> str: normalized = (metric or "containment_rate").strip() if normalized not in _AI_ANALYTICS_METRICS: raise HTTPException(status_code=400, detail="Invalid metric") return normalized def _normalize_ai_analytics_interval(interval: str) -> str: normalized = (interval or "day").strip().lower() if normalized not in _AI_ANALYTICS_INTERVALS: raise HTTPException(status_code=400, detail="Invalid interval") return normalized def _normalize_ai_analytics_slice(slice_name: str | None) -> str: normalized = (slice_name or "all").strip().lower() or "all" if normalized not in _AI_ANALYTICS_SLICES: raise HTTPException(status_code=400, detail="Invalid slice") return normalized def _normalize_ai_analytics_sort(sort_by: str | None, sort_dir: str | None) -> tuple[str, str]: normalized_sort_by = (sort_by or "created_at").strip().lower() or "created_at" if normalized_sort_by not in _AI_ANALYTICS_DRILLDOWN_SORT_FIELDS: raise HTTPException(status_code=400, detail="Invalid sort_by") normalized_sort_dir = (sort_dir or "desc").strip().lower() or "desc" if normalized_sort_dir not in {"asc", "desc"}: raise HTTPException(status_code=400, detail="Invalid sort_dir") return normalized_sort_by, normalized_sort_dir def _normalize_ai_analytics_reason_key(reason_key: str | None) -> str | None: normalized = (reason_key or "").strip().lower() or None if normalized is None: return None if normalized not in _AI_ANALYTICS_REASON_LABELS: raise HTTPException(status_code=400, detail="Invalid reason_key") return normalized def _analytics_percent(numerator: int, denominator: int) -> float: if denominator <= 0: return 0.0 return round((numerator / denominator) * 100, 2) def _analytics_percentile(values: list[int], percentile: float) -> float | None: if not values: return None ordered = sorted(int(value) for value in values) index = max(0, math.ceil(percentile * len(ordered)) - 1) return float(ordered[index]) def _analytics_average(values: list[int]) -> float | None: if not values: return None return round(sum(values) / len(values), 2) def _analytics_window(range_from: datetime, range_to: datetime) -> AIAnalyticsWindowOut: return AIAnalyticsWindowOut(from_ts=range_from.isoformat(), to_ts=range_to.isoformat()) def _analytics_filters(range_from: datetime, range_to: datetime, queue_id: str | None, channel: str | None) -> AIAnalyticsFiltersOut: return AIAnalyticsFiltersOut( from_ts=range_from.isoformat(), to_ts=range_to.isoformat(), queue_id=queue_id, channel=channel, ) def _safe_json_loads(raw: str | None) -> dict[str, Any]: if not raw: return {} try: payload = json.loads(raw) except (TypeError, ValueError): return {} return payload if isinstance(payload, dict) else {} def _normalize_voice_name_language(value: str | None) -> str | None: normalized = str(value or "").strip().lower() if not normalized or normalized == "all": return None return normalized def _voice_name_filters(range_from: datetime, range_to: datetime, queue_id: str | None, language: str | None) -> VoiceNameFlowAnalyticsFiltersOut: return VoiceNameFlowAnalyticsFiltersOut( from_ts=range_from.isoformat(), to_ts=range_to.isoformat(), queue_id=queue_id, language=language, ) def _normalize_voice_name_metric(metric: str) -> str: normalized = str(metric or "").strip() if normalized not in _VOICE_NAME_FLOW_METRICS: raise HTTPException(status_code=400, detail="Unsupported metric") return normalized def _voice_name_interval_for_window(range_from: datetime, range_to: datetime) -> str: return "hour" if (range_to - range_from) <= timedelta(hours=36) else "day" def _normalize_ai_handoff_reason(reason: str | None, claimed_by_user: str | None = None) -> tuple[str | None, str | None, str | None]: raw = (reason or "").strip() or None if not raw and claimed_by_user: return "manual_claim", _AI_ANALYTICS_REASON_LABELS["manual_claim"], None if not raw: return None, None, None lowered = raw.lower() if any(token in lowered for token in ("operator", "жив", "человек", "customer requested", "requested by customer", "клиент запрос")): key = "requested_human" elif any(token in lowered for token in ("knowledge", "kb", "tools", "tool", "баз", "знани", "инструмент", "данных")): key = "knowledge_or_tool_gap" elif any(token in lowered for token in ("policy", "sensitive", "чувств", "политик", "комплаенс", "restricted")): key = "policy_or_sensitive" elif any(token in lowered for token in ("error", "timeout", "delivery", "runtime", "exception", "ошиб", "сбой", "failure")): key = "delivery_or_runtime_error" else: key = "other" return key, _AI_ANALYTICS_REASON_LABELS[key], raw def _build_ai_analytics_snapshot( row: AISessionRow, *, interaction: Interaction | None, thread: TelegramThreadRow | WhatsAppThreadRow | None, turns: list[AITurnRow], ) -> dict[str, Any]: resolved_queue_id = None if interaction and interaction.queue_id: resolved_queue_id = interaction.queue_id elif thread and getattr(thread, "queue_id", None): resolved_queue_id = getattr(thread, "queue_id") assigned_to = interaction.assigned_to if interaction else None claimed_by_user = getattr(thread, "claimed_by_user", None) if thread else None handoff_reason_source = (row.handoff_reason or getattr(thread, "ai_handoff_reason", None) or "").strip() or None reason_key, reason_label, raw_handoff_reason = _normalize_ai_handoff_reason(handoff_reason_source, claimed_by_user) assistant_turns = sum(1 for turn in turns if turn.role == "assistant") user_turns = sum(1 for turn in turns if turn.role == "user") tool_turns = sum(1 for turn in turns if turn.source_type == "tool") latency_values = [ int(turn.latency_ms) for turn in turns if turn.role == "assistant" and turn.source_type == "model" and turn.latency_ms is not None ] interaction_closed = bool(interaction and interaction.status == "closed") human_touched = bool(row.status == "human_owned" or assigned_to or claimed_by_user) handoff = bool(row.status in {"handoff_required", "human_owned"} or raw_handoff_reason) contained = bool(row.status == "closed" and not raw_handoff_reason and not human_touched and not assigned_to) closed_without_operator = bool(interaction_closed and interaction and not interaction.assigned_to) return { "session_id": row.session_id, "thread_id": row.thread_id, "interaction_id": row.interaction_id, "channel": row.channel, "queue_id": resolved_queue_id, "status": row.status, "created_at": row.created_at, "updated_at": row.updated_at, "closed_at": row.closed_at, "contained": contained, "handoff": handoff, "human_touched": human_touched, "closed_without_operator": closed_without_operator, "interaction_closed": interaction_closed, "assigned_to": assigned_to, "claimed_by_user": claimed_by_user, "assistant_turns": assistant_turns, "user_turns": user_turns, "tool_turns": tool_turns, "latencies": latency_values, "ai_latency_avg_ms": _analytics_average(latency_values), "ai_latency_p95_ms": _analytics_percentile(latency_values, 0.95), "reason_key": reason_key, "reason_label": reason_label, "raw_handoff_reason": raw_handoff_reason, "subject": interaction.subject if interaction else None, } def _ai_analytics_coverage_from_snapshots(snapshots: list[dict[str, Any]]) -> AIAnalyticsCoverageOut: return AIAnalyticsCoverageOut( sessions_with_interaction_id=sum(1 for item in snapshots if item.get("interaction_id")), sessions_with_queue_id=sum(1 for item in snapshots if item.get("queue_id")), sessions_with_latency_turns=sum(1 for item in snapshots if item.get("latencies")), sessions_with_terminal_state=sum(1 for item in snapshots if item.get("status") in _AI_ANALYTICS_TERMINAL_STATUSES), sessions_with_handoff_reason=sum(1 for item in snapshots if item.get("raw_handoff_reason")), ) def _ai_analytics_snapshot_matches_slice(snapshot: dict[str, Any], slice_name: str) -> bool: if slice_name == "all": return True if slice_name == "contained": return bool(snapshot.get("contained")) if slice_name == "handoff": return bool(snapshot.get("handoff")) if slice_name == "human_touched": return bool(snapshot.get("human_touched")) if slice_name == "closed_without_operator": return bool(snapshot.get("closed_without_operator")) if slice_name == "active": return str(snapshot.get("status") or "").lower() == "active" if slice_name == "error": return str(snapshot.get("status") or "").lower() == "error" return True def _ai_analytics_snapshot_matches_query(snapshot: dict[str, Any], query_text: str | None) -> bool: normalized = (query_text or "").strip().lower() if not normalized: return True haystack = " ".join( [ str(snapshot.get("session_id") or ""), str(snapshot.get("interaction_id") or ""), str(snapshot.get("thread_id") or ""), str(snapshot.get("raw_handoff_reason") or ""), ] ).lower() return normalized in haystack def _ai_analytics_snapshot_sort_value(snapshot: dict[str, Any], sort_by: str) -> Any: if sort_by == "updated_at": return snapshot.get("updated_at") or "" if sort_by == "ai_latency_avg_ms": latency = snapshot.get("ai_latency_avg_ms") if latency is None: return -1 return float(latency) if sort_by == "status": return str(snapshot.get("status") or "") return snapshot.get("created_at") or "" def _ai_analytics_snapshot_to_item(snapshot: dict[str, Any]) -> AIAnalyticsDrilldownItemOut: return AIAnalyticsDrilldownItemOut( session_id=str(snapshot.get("session_id") or ""), thread_id=snapshot.get("thread_id"), interaction_id=snapshot.get("interaction_id"), channel=str(snapshot.get("channel") or "unknown"), queue_id=snapshot.get("queue_id"), status=str(snapshot.get("status") or "unknown"), created_at=str(snapshot.get("created_at") or ""), updated_at=str(snapshot.get("updated_at") or ""), closed_at=snapshot.get("closed_at"), contained=bool(snapshot.get("contained")), handoff=bool(snapshot.get("handoff")), human_touched=bool(snapshot.get("human_touched")), closed_without_operator=bool(snapshot.get("closed_without_operator")), reason_key=snapshot.get("reason_key"), reason_label=snapshot.get("reason_label"), raw_handoff_reason=snapshot.get("raw_handoff_reason"), assigned_to=snapshot.get("assigned_to"), claimed_by_user=snapshot.get("claimed_by_user"), assistant_turns=int(snapshot.get("assistant_turns") or 0), user_turns=int(snapshot.get("user_turns") or 0), tool_turns=int(snapshot.get("tool_turns") or 0), ai_latency_avg_ms=snapshot.get("ai_latency_avg_ms"), ai_latency_p95_ms=snapshot.get("ai_latency_p95_ms"), ) def _empty_ai_analytics_overview( *, range_from: datetime, range_to: datetime, queue_id: str | None, channel: str | None, ) -> AIAnalyticsOverviewOut: return AIAnalyticsOverviewOut( window=_analytics_window(range_from, range_to), filters=_analytics_filters(range_from, range_to, queue_id, channel), totals=AIAnalyticsTotalsOut(), metrics=AIAnalyticsMetricsOut(), breakdowns=AIAnalyticsBreakdownsOut(), coverage=AIAnalyticsCoverageOut(), ) def _empty_ai_analytics_timeseries( *, range_from: datetime, range_to: datetime, metric: str, interval: str, queue_id: str | None, channel: str | None, ) -> AIAnalyticsTimeseriesOut: return AIAnalyticsTimeseriesOut( metric=metric, # type: ignore[arg-type] interval=interval, # type: ignore[arg-type] filters=_analytics_filters(range_from, range_to, queue_id, channel), points=[], ) def _empty_voice_name_flow_timeseries( *, range_from: datetime, range_to: datetime, metric: str, interval: str, queue_id: str | None, language: str | None, ) -> VoiceNameFlowAnalyticsTimeseriesOut: return VoiceNameFlowAnalyticsTimeseriesOut( metric=metric, # type: ignore[arg-type] interval=interval, # type: ignore[arg-type] filters=_voice_name_filters(range_from, range_to, queue_id, language), points=[], ) def _load_ai_analytics_snapshots( session, *, range_from: datetime, range_to: datetime, queue_id: str | None, channels: list[str], ) -> list[dict[str, Any]]: if not channels: return [] rows = session.execute( select(AISessionRow).where( AISessionRow.created_at >= range_from.isoformat(), AISessionRow.created_at < range_to.isoformat(), AISessionRow.channel.in_(channels), ) ).scalars().all() if not rows: return [] interaction_ids = {row.interaction_id for row in rows if row.interaction_id} telegram_thread_ids = {row.thread_id for row in rows if row.channel == "telegram" and row.thread_id} whatsapp_thread_ids = {row.thread_id for row in rows if row.channel == "whatsapp" and row.thread_id} session_ids = [row.session_id for row in rows] interactions = {} if interaction_ids: interactions = { row.interaction_id: row for row in session.execute( select(Interaction).where(Interaction.interaction_id.in_(interaction_ids)) ).scalars().all() } telegram_threads = {} if telegram_thread_ids: telegram_threads = { row.thread_id: row for row in session.execute( select(TelegramThreadRow).where(TelegramThreadRow.thread_id.in_(telegram_thread_ids)) ).scalars().all() } whatsapp_threads = {} if whatsapp_thread_ids: whatsapp_threads = { row.thread_id: row for row in session.execute( select(WhatsAppThreadRow).where(WhatsAppThreadRow.thread_id.in_(whatsapp_thread_ids)) ).scalars().all() } turns_by_session: dict[str, list[AITurnRow]] = defaultdict(list) for turn in session.execute( select(AITurnRow).where( AITurnRow.session_id.in_(session_ids), ) ).scalars().all(): turns_by_session[turn.session_id].append(turn) snapshots: list[dict[str, Any]] = [] for row in rows: interaction = interactions.get(row.interaction_id) if row.interaction_id else None if row.channel == "telegram": thread = telegram_threads.get(row.thread_id) if row.thread_id else None else: thread = whatsapp_threads.get(row.thread_id) if row.thread_id else None resolved_queue_id = None if interaction and interaction.queue_id: resolved_queue_id = interaction.queue_id elif thread and getattr(thread, "queue_id", None): resolved_queue_id = getattr(thread, "queue_id") if queue_id: if not resolved_queue_id or resolved_queue_id != queue_id: continue snapshots.append( _build_ai_analytics_snapshot( row, interaction=interaction, thread=thread, turns=list(turns_by_session.get(row.session_id, [])), ) ) return snapshots def _aggregate_ai_analytics_overview( snapshots: list[dict[str, Any]], *, range_from: datetime, range_to: datetime, queue_id: str | None, channel: str | None, ) -> AIAnalyticsOverviewOut: if not snapshots: return _empty_ai_analytics_overview( range_from=range_from, range_to=range_to, queue_id=queue_id, channel=channel, ) total_sessions = len(snapshots) contained_sessions = sum(1 for item in snapshots if item["contained"]) handoff_sessions = sum(1 for item in snapshots if item["handoff"]) closed_sessions = sum(1 for item in snapshots if item["status"] == "closed" or item["interaction_closed"]) closed_without_operator = sum(1 for item in snapshots if item["closed_without_operator"]) closed_candidates = sum(1 for item in snapshots if item["interaction_closed"]) human_touched_sessions = sum(1 for item in snapshots if item["human_touched"]) latency_values = [latency for item in snapshots for latency in item["latencies"]] totals = AIAnalyticsTotalsOut( sessions_started=total_sessions, sessions_contained=contained_sessions, sessions_handoff=handoff_sessions, sessions_closed=closed_sessions, sessions_closed_without_operator=closed_without_operator, assistant_turns=len(latency_values), ) metrics = AIAnalyticsMetricsOut( containment_rate=_analytics_percent(contained_sessions, total_sessions), handoff_rate=_analytics_percent(handoff_sessions, total_sessions), ai_latency_avg_ms=_analytics_average(latency_values), ai_latency_p95_ms=_analytics_percentile(latency_values, 0.95), closed_without_operator_rate=_analytics_percent(closed_without_operator, closed_candidates), human_touched_rate=_analytics_percent(human_touched_sessions, total_sessions), ) channel_breakdowns = [] grouped: dict[str, list[dict[str, Any]]] = defaultdict(list) for item in snapshots: grouped[item["channel"]].append(item) for resolved_channel, items in sorted(grouped.items(), key=lambda entry: (-len(entry[1]), entry[0])): channel_latencies = [latency for item in items for latency in item["latencies"]] channel_closed_candidates = sum(1 for item in items if item["interaction_closed"]) channel_contained = sum(1 for item in items if item["contained"]) channel_handoff = sum(1 for item in items if item["handoff"]) channel_closed_without_operator = sum(1 for item in items if item["closed_without_operator"]) human_touched_sessions = sum(1 for item in items if item["human_touched"]) channel_breakdowns.append( AIAnalyticsChannelBreakdownOut( channel=resolved_channel, sessions_started=len(items), sessions_contained=channel_contained, sessions_handoff=channel_handoff, sessions_closed_without_operator=channel_closed_without_operator, assistant_turns=len(channel_latencies), containment_rate=_analytics_percent(channel_contained, len(items)), handoff_rate=_analytics_percent(channel_handoff, len(items)), closed_without_operator_rate=_analytics_percent( channel_closed_without_operator, channel_closed_candidates, ), ai_latency_avg_ms=_analytics_average(channel_latencies), ai_only_sessions=len(items) - human_touched_sessions, human_touched_sessions=human_touched_sessions, ) ) outcome_breakdowns = [ AIAnalyticsOutcomeBreakdownOut( outcome="contained", label=_AI_ANALYTICS_OUTCOME_LABELS["contained"], sessions=contained_sessions, share=_analytics_percent(contained_sessions, total_sessions), ), AIAnalyticsOutcomeBreakdownOut( outcome="handoff", label=_AI_ANALYTICS_OUTCOME_LABELS["handoff"], sessions=handoff_sessions, share=_analytics_percent(handoff_sessions, total_sessions), ), AIAnalyticsOutcomeBreakdownOut( outcome="human_touched", label=_AI_ANALYTICS_OUTCOME_LABELS["human_touched"], sessions=human_touched_sessions, share=_analytics_percent(human_touched_sessions, total_sessions), ), AIAnalyticsOutcomeBreakdownOut( outcome="closed_without_operator", label=_AI_ANALYTICS_OUTCOME_LABELS["closed_without_operator"], sessions=closed_without_operator, share=_analytics_percent(closed_without_operator, total_sessions), ), AIAnalyticsOutcomeBreakdownOut( outcome="active", label=_AI_ANALYTICS_OUTCOME_LABELS["active"], sessions=sum(1 for item in snapshots if item["status"] == "active"), share=_analytics_percent(sum(1 for item in snapshots if item["status"] == "active"), total_sessions), ), AIAnalyticsOutcomeBreakdownOut( outcome="error", label=_AI_ANALYTICS_OUTCOME_LABELS["error"], sessions=sum(1 for item in snapshots if item["status"] == "error"), share=_analytics_percent(sum(1 for item in snapshots if item["status"] == "error"), total_sessions), ), ] handoff_reason_rows: list[AIAnalyticsHandoffReasonBreakdownOut] = [] handoff_reason_groups: dict[str, int] = defaultdict(int) for item in snapshots: if item["reason_key"]: handoff_reason_groups[str(item["reason_key"])] += 1 for reason_key, sessions_count in sorted(handoff_reason_groups.items(), key=lambda entry: (-entry[1], entry[0])): handoff_reason_rows.append( AIAnalyticsHandoffReasonBreakdownOut( reason_key=reason_key, label=_AI_ANALYTICS_REASON_LABELS.get(reason_key, reason_key), sessions=sessions_count, share=_analytics_percent(sessions_count, handoff_sessions or sessions_count), ) ) coverage = _ai_analytics_coverage_from_snapshots(snapshots) return AIAnalyticsOverviewOut( window=_analytics_window(range_from, range_to), filters=_analytics_filters(range_from, range_to, queue_id, channel), totals=totals, metrics=metrics, breakdowns=AIAnalyticsBreakdownsOut( by_channel=channel_breakdowns, by_outcome=outcome_breakdowns, by_handoff_reason=handoff_reason_rows, ), coverage=coverage, ) def _load_ai_analytics_overview( session, *, range_from: datetime, range_to: datetime, queue_id: str | None, channel: str | None, ) -> AIAnalyticsOverviewOut: normalized_channel, channels = _normalize_ai_analytics_channel(channel) snapshots = _load_ai_analytics_snapshots( session, range_from=range_from, range_to=range_to, queue_id=queue_id, channels=channels, ) return _aggregate_ai_analytics_overview( snapshots, range_from=range_from, range_to=range_to, queue_id=queue_id, channel=normalized_channel, ) def _timeseries_metric_value(metric: str, overview: AIAnalyticsOverviewOut) -> float | None: if metric == "containment_rate": return overview.metrics.containment_rate if metric == "handoff_rate": return overview.metrics.handoff_rate if metric == "ai_latency_avg_ms": return overview.metrics.ai_latency_avg_ms if metric == "human_touched_rate": return overview.metrics.human_touched_rate return overview.metrics.closed_without_operator_rate def _load_ai_analytics_drilldown( session, *, range_from: datetime, range_to: datetime, queue_id: str | None, channel: str | None, slice_name: str, reason_key: str | None, status: str | None, query_text: str | None, sort_by: str, sort_dir: str, limit: int, offset: int, ) -> AIAnalyticsDrilldownOut: normalized_channel, channels = _normalize_ai_analytics_channel(channel) if not channels: return AIAnalyticsDrilldownOut( items=[], total=0, limit=limit, offset=offset, filters=AIAnalyticsDrilldownFiltersOut( from_ts=range_from.isoformat(), to_ts=range_to.isoformat(), queue_id=queue_id, channel=normalized_channel, slice=slice_name, # type: ignore[arg-type] reason_key=reason_key, status=status, q=query_text, sort_by=sort_by, # type: ignore[arg-type] sort_dir=sort_dir, # type: ignore[arg-type] ), coverage=AIAnalyticsCoverageOut(), ) snapshots = _load_ai_analytics_snapshots( session, range_from=range_from, range_to=range_to, queue_id=queue_id, channels=channels, ) filtered = [ item for item in snapshots if _ai_analytics_snapshot_matches_slice(item, slice_name) and (not reason_key or item.get("reason_key") == reason_key) and (not status or item.get("status") == status) and _ai_analytics_snapshot_matches_query(item, query_text) ] reverse = sort_dir == "desc" filtered.sort( key=lambda item: ( _ai_analytics_snapshot_sort_value(item, sort_by), item.get("created_at") or "", item.get("session_id") or "", ), reverse=reverse, ) paged = filtered[offset : offset + limit] return AIAnalyticsDrilldownOut( items=[_ai_analytics_snapshot_to_item(item) for item in paged], total=len(filtered), limit=limit, offset=offset, filters=AIAnalyticsDrilldownFiltersOut( from_ts=range_from.isoformat(), to_ts=range_to.isoformat(), queue_id=queue_id, channel=normalized_channel, slice=slice_name, # type: ignore[arg-type] reason_key=reason_key, status=status, q=query_text, sort_by=sort_by, # type: ignore[arg-type] sort_dir=sort_dir, # type: ignore[arg-type] ), coverage=_ai_analytics_coverage_from_snapshots(filtered), ) def _load_ai_analytics_session_detail(session, session_id: str) -> AIAnalyticsSessionDetailOut: row = session.execute( select(AISessionRow).where(AISessionRow.session_id == session_id) ).scalar_one_or_none() if row is None or row.channel not in _AI_ANALYTICS_CHANNELS: raise HTTPException(status_code=404, detail="AI session not found") interaction = None if row.interaction_id: interaction = session.execute( select(Interaction).where(Interaction.interaction_id == row.interaction_id) ).scalar_one_or_none() thread = None if row.thread_id: if row.channel == "telegram": thread = session.execute( select(TelegramThreadRow).where(TelegramThreadRow.thread_id == row.thread_id) ).scalar_one_or_none() elif row.channel == "whatsapp": thread = session.execute( select(WhatsAppThreadRow).where(WhatsAppThreadRow.thread_id == row.thread_id) ).scalar_one_or_none() turns = session.execute( select(AITurnRow).where(AITurnRow.session_id == session_id).order_by(AITurnRow.created_at.asc(), AITurnRow.id.asc()) ).scalars().all() snapshot = _build_ai_analytics_snapshot(row, interaction=interaction, thread=thread, turns=turns) timeline: list[AIAnalyticsSessionEventOut] = [ AIAnalyticsSessionEventOut( ts=row.created_at, event_type="session_created", label="AI session created", status=row.status, metadata={ "channel": row.channel, "agent_profile": row.agent_profile, "language": row.language, }, ) ] for turn in turns: timeline.append( AIAnalyticsSessionEventOut( ts=turn.created_at, event_type="turn", label=f"{turn.role or 'turn'} via {turn.source_type or 'unknown'}", role=turn.role, source_type=turn.source_type, latency_ms=turn.latency_ms, finish_reason=turn.finish_reason, metadata={ "model": turn.model, "thread_id": turn.thread_id, "interaction_id": turn.interaction_id, }, ) ) if row.handoff_reason or getattr(thread, "ai_handoff_reason", None): timeline.append( AIAnalyticsSessionEventOut( ts=row.updated_at, event_type="handoff", label="Handoff recorded", status=row.status, metadata={ "reason_key": snapshot.get("reason_key"), "claimed_by_user": snapshot.get("claimed_by_user"), }, ) ) if row.closed_at: timeline.append( AIAnalyticsSessionEventOut( ts=row.closed_at, event_type="session_closed", label="AI session closed", status=row.status, metadata={ "contained": snapshot.get("contained"), "closed_without_operator": snapshot.get("closed_without_operator"), }, ) ) elif row.updated_at and row.updated_at != row.created_at: timeline.append( AIAnalyticsSessionEventOut( ts=row.updated_at, event_type="session_updated", label="AI session updated", status=row.status, metadata={ "human_touched": snapshot.get("human_touched"), "reason_key": snapshot.get("reason_key"), }, ) ) timeline.sort(key=lambda item: (item.ts, item.event_type)) interaction_snapshot = None if interaction: interaction_snapshot = AIAnalyticsSessionLinkedInteractionOut( interaction_id=interaction.interaction_id, channel=interaction.channel, queue_id=interaction.queue_id, status=interaction.status, assigned_to=interaction.assigned_to, subject=interaction.subject, created_at=interaction.created_at, updated_at=interaction.updated_at, ) return AIAnalyticsSessionDetailOut( session=_ai_analytics_snapshot_to_item(snapshot), interaction=interaction_snapshot, timeline=timeline, ) def _voice_name_policy_decision_from_turn(turn: AITurnRow) -> dict[str, Any] | None: if str(turn.source_type or "").strip() != "voice_policy": return None payload = _safe_json_loads(turn.payload_json) decision = payload.get("decision") if not isinstance(decision, dict): return None metadata = decision.get("metadata") source = metadata if isinstance(metadata, dict) else decision status = str(source.get("customer_name_status") or "").strip() if status not in _VOICE_NAME_FLOW_STATUS_VALUES: return None value = str(source.get("customer_name_value") or "").strip() or None name_source = str(source.get("customer_name_source") or "").strip() or None language = str(source.get("language") or decision.get("language") or "").strip().lower() or None return { "ts": turn.created_at, "status": status, "value": value, "source": name_source, "language": language, } def _voice_name_start_event_from_timeline(row: InteractionTimeline) -> dict[str, Any] | None: if str(row.action or "").strip() != "voice.start.completed": return None metadata = _safe_json_loads(row.metadata_json) status = str(metadata.get("customer_name_status") or "").strip() if status not in _VOICE_NAME_FLOW_STATUS_VALUES: return None value = str(metadata.get("customer_name_value") or "").strip() or None name_source = str(metadata.get("customer_name_source") or "").strip() or None language = str(metadata.get("language") or "").strip().lower() or None return { "ts": row.timestamp, "status": status, "value": value, "source": name_source, "language": language, } def _voice_name_primary_outcome(start_decision: dict[str, Any], final_state: dict[str, Any]) -> str: if start_decision.get("status") == "name_obtained": return "start_obtained" if final_state.get("status") == "name_obtained": return "downstream_ai_obtained" if final_state.get("status") == "name_followup_required": return "followup_required" return "name_not_obtained" def _voice_name_operator_handoff( row: VoiceAISessionRow, *, call_row: AsteriskCallLinkRow | None, interaction: Interaction | None, ) -> bool: handoff_reason = str( row.handoff_reason or (call_row.ai_handoff_reason if call_row else "") or "" ).strip() has_operator_owner = bool( (interaction and interaction.assigned_to) or (call_row and (call_row.claimed_by_user or call_row.operator_extension)) ) if has_operator_owner: return True if handoff_reason and handoff_reason != "voice_start_completed": return True call_ai_state = str(call_row.ai_state or "").strip() if call_row else "" row_status = str(row.status or "").strip() if call_ai_state == "human_owned" or row_status == "human_owned": return True if handoff_reason and handoff_reason != "voice_start_completed": return call_ai_state in {"handoff_requested", "handoff_required"} or row_status in {"handoff_requested", "handoff_required"} return False def _voice_name_analytics_coverage_from_snapshots(snapshots: list[dict[str, Any]]) -> VoiceNameFlowAnalyticsCoverageOut: note = ( "Статусы восстановления имени собраны из voice_policy turns и start-event metadata; ручное исправление считается отдельной overlay-метрикой." if snapshots else None ) return VoiceNameFlowAnalyticsCoverageOut( sessions_with_start_decision=sum(1 for item in snapshots if item.get("start_decision")), sessions_with_final_ai_state=sum(1 for item in snapshots if item.get("final_state")), sessions_with_manual_overlay=sum(1 for item in snapshots if item.get("manual_corrected")), note=note, ) def _voice_name_metric_value(metric: str, overview: VoiceNameFlowAnalyticsOverviewOut) -> float | None: if metric == "scenario_calls": return float(overview.totals.scenario_calls) if metric == "start_capture_rate": return overview.metrics.start_capture_rate if metric == "downstream_rescue_rate": return overview.metrics.downstream_rescue_rate if metric == "handoff_unconfirmed_rate": return overview.metrics.handoff_unconfirmed_rate return overview.metrics.manual_correction_rate def _voice_name_metric_denominator(metric: str, overview: VoiceNameFlowAnalyticsOverviewOut) -> int: if metric == "scenario_calls": return overview.totals.scenario_calls if metric == "start_capture_rate": return overview.totals.scenario_calls if metric == "downstream_rescue_rate": return overview.totals.needed_downstream return overview.totals.handoff_confirmed_name + overview.totals.handoff_unconfirmed_name def _build_voice_name_flow_snapshot( row: VoiceAISessionRow, *, call_row: AsteriskCallLinkRow | None, interaction: Interaction | None, turns: list[AITurnRow], start_event: InteractionTimeline | None, ) -> dict[str, Any] | None: decisions = [ item for item in ( _voice_name_policy_decision_from_turn(turn) for turn in sorted(turns, key=lambda candidate: (candidate.created_at, candidate.turn_id)) ) if item is not None ] fallback_start = _voice_name_start_event_from_timeline(start_event) if start_event else None start_decision = decisions[0] if decisions else fallback_start final_state = decisions[-1] if decisions else fallback_start if not start_decision and not final_state: return None resolved_queue_id = ( str(row.queue_id or "").strip() or (str(call_row.queue_id or "").strip() if call_row else "") or (str(interaction.queue_id or "").strip() if interaction else "") or "unknown" ) resolved_language = ( str(row.voice_start_language or "").strip().lower() or str(row.language or "").strip().lower() or (str(call_row.voice_start_language or "").strip().lower() if call_row else "") or str((start_decision or {}).get("language") or "").strip().lower() or "unknown" ) manual_corrected = ( str(row.customer_name_source or "").strip() == "manual" or (str(call_row.customer_name_source or "").strip() == "manual" if call_row else False) ) primary_outcome = _voice_name_primary_outcome(start_decision or {}, final_state or start_decision or {}) operator_handoff = _voice_name_operator_handoff(row, call_row=call_row, interaction=interaction) final_status = str((final_state or {}).get("status") or "name_not_obtained") return { "session_id": row.session_id, "call_id": row.call_id, "interaction_id": row.interaction_id or (call_row.interaction_id if call_row else None), "queue_id": resolved_queue_id, "language": resolved_language, "started_at": row.started_at, "status": row.status, "start_decision": start_decision, "final_state": final_state, "manual_corrected": manual_corrected, "primary_outcome": primary_outcome, "needed_downstream": primary_outcome != "start_obtained", "operator_handoff": operator_handoff, "handoff_confirmed_name": operator_handoff and final_status == "name_obtained", "handoff_unconfirmed_name": operator_handoff and final_status in {"name_followup_required", "name_not_obtained"}, } def _load_voice_name_flow_snapshots( session, *, range_from: datetime, range_to: datetime, queue_id: str | None, language: str | None, ) -> list[dict[str, Any]]: rows = session.execute( select(VoiceAISessionRow).where( VoiceAISessionRow.started_at >= range_from.isoformat(), VoiceAISessionRow.started_at < range_to.isoformat(), ) ).scalars().all() if not rows: return [] interaction_ids = {row.interaction_id for row in rows if row.interaction_id} call_ids = {row.call_id for row in rows if row.call_id} voice_session_ids = {row.session_id for row in rows if row.session_id} ai_session_ids = {row.ai_session_id for row in rows if row.ai_session_id} interactions = {} if interaction_ids: interactions = { item.interaction_id: item for item in session.execute( select(Interaction).where(Interaction.interaction_id.in_(interaction_ids)) ).scalars().all() } call_rows = {} for item in session.execute( select(AsteriskCallLinkRow).where( AsteriskCallLinkRow.call_id.in_(call_ids) if call_ids else False ) ).scalars().all(): call_rows[item.call_id] = item if voice_session_ids: for item in session.execute( select(AsteriskCallLinkRow).where(AsteriskCallLinkRow.voice_session_id.in_(voice_session_ids)) ).scalars().all(): if item.call_id not in call_rows: call_rows[item.call_id] = item turns_by_ai_session: dict[str, list[AITurnRow]] = defaultdict(list) if ai_session_ids: for turn in session.execute( select(AITurnRow).where(AITurnRow.session_id.in_(ai_session_ids)) ).scalars().all(): turns_by_ai_session[turn.session_id].append(turn) start_events: dict[str, InteractionTimeline] = {} if interaction_ids: timeline_rows = session.execute( select(InteractionTimeline).where( InteractionTimeline.interaction_id.in_(interaction_ids), InteractionTimeline.action == "voice.start.completed", ) ).scalars().all() for item in sorted(timeline_rows, key=lambda candidate: (candidate.timestamp, candidate.id)): if item.interaction_id not in start_events: start_events[item.interaction_id] = item normalized_language = _normalize_voice_name_language(language) snapshots: list[dict[str, Any]] = [] for row in rows: call_row = call_rows.get(row.call_id) interaction = interactions.get(row.interaction_id) if row.interaction_id else None snapshot = _build_voice_name_flow_snapshot( row, call_row=call_row, interaction=interaction, turns=list(turns_by_ai_session.get(str(row.ai_session_id or ""), [])), start_event=start_events.get(row.interaction_id) if row.interaction_id else None, ) if snapshot is None: continue if queue_id and snapshot["queue_id"] != queue_id: continue if normalized_language and snapshot["language"] != normalized_language: continue snapshots.append(snapshot) return snapshots def _voice_name_breakdown_metrics(snapshots: list[dict[str, Any]]) -> dict[str, float]: scenario_calls = len(snapshots) start_obtained = sum(1 for item in snapshots if item["primary_outcome"] == "start_obtained") downstream_ai_obtained = sum(1 for item in snapshots if item["primary_outcome"] == "downstream_ai_obtained") needed_downstream = sum(1 for item in snapshots if item.get("needed_downstream")) handoff_confirmed = sum(1 for item in snapshots if item.get("handoff_confirmed_name")) handoff_unconfirmed = sum(1 for item in snapshots if item.get("handoff_unconfirmed_name")) all_handoffs = handoff_confirmed + handoff_unconfirmed manual_corrected = sum(1 for item in snapshots if item.get("manual_corrected")) return { "start_capture_rate": _analytics_percent(start_obtained, scenario_calls), "downstream_rescue_rate": _analytics_percent(downstream_ai_obtained, needed_downstream), "handoff_unconfirmed_rate": _analytics_percent(handoff_unconfirmed, all_handoffs), "manual_correction_rate": _analytics_percent(manual_corrected, all_handoffs), } def _aggregate_voice_name_flow_overview( snapshots: list[dict[str, Any]], *, range_from: datetime, range_to: datetime, queue_id: str | None, language: str | None, ) -> VoiceNameFlowAnalyticsOverviewOut: if not snapshots: return VoiceNameFlowAnalyticsOverviewOut( window=_analytics_window(range_from, range_to), filters=_voice_name_filters(range_from, range_to, queue_id, language), totals=VoiceNameFlowAnalyticsTotalsOut(), metrics=VoiceNameFlowAnalyticsMetricsOut(), breakdowns=VoiceNameFlowAnalyticsBreakdownsOut(), coverage=VoiceNameFlowAnalyticsCoverageOut(), ) scenario_calls = len(snapshots) start_obtained = sum(1 for item in snapshots if item["primary_outcome"] == "start_obtained") downstream_ai_obtained = sum(1 for item in snapshots if item["primary_outcome"] == "downstream_ai_obtained") followup_required = sum(1 for item in snapshots if item["primary_outcome"] == "followup_required") name_not_obtained = sum(1 for item in snapshots if item["primary_outcome"] == "name_not_obtained") manual_corrected = sum(1 for item in snapshots if item.get("manual_corrected")) handoff_confirmed_name = sum(1 for item in snapshots if item.get("handoff_confirmed_name")) handoff_unconfirmed_name = sum(1 for item in snapshots if item.get("handoff_unconfirmed_name")) needed_downstream = sum(1 for item in snapshots if item.get("needed_downstream")) totals = VoiceNameFlowAnalyticsTotalsOut( scenario_calls=scenario_calls, start_obtained=start_obtained, downstream_ai_obtained=downstream_ai_obtained, followup_required=followup_required, name_not_obtained=name_not_obtained, manual_corrected=manual_corrected, handoff_confirmed_name=handoff_confirmed_name, handoff_unconfirmed_name=handoff_unconfirmed_name, needed_downstream=needed_downstream, ) metric_values = _voice_name_breakdown_metrics(snapshots) metrics = VoiceNameFlowAnalyticsMetricsOut(**metric_values) funnel = [ VoiceNameFlowAnalyticsFunnelStageOut( stage=stage, # type: ignore[arg-type] label=_VOICE_NAME_FLOW_FUNNEL_LABELS[stage], sessions=value, share=_analytics_percent(value, scenario_calls), ) for stage, value in ( ("scenario_calls", scenario_calls), ("start_obtained", start_obtained), ("needed_downstream", needed_downstream), ("downstream_ai_obtained", downstream_ai_obtained), ("handoff_confirmed_name", handoff_confirmed_name), ("handoff_unconfirmed_name", handoff_unconfirmed_name), ) ] by_language_rows: list[VoiceNameFlowAnalyticsLanguageBreakdownOut] = [] language_groups: dict[str, list[dict[str, Any]]] = defaultdict(list) for item in snapshots: language_groups[str(item["language"] or "unknown")].append(item) for resolved_language, items in sorted(language_groups.items(), key=lambda entry: (-len(entry[1]), entry[0])): breakdown_metrics = _voice_name_breakdown_metrics(items) by_language_rows.append( VoiceNameFlowAnalyticsLanguageBreakdownOut( language=resolved_language, scenario_calls=len(items), start_obtained=sum(1 for item in items if item["primary_outcome"] == "start_obtained"), downstream_ai_obtained=sum(1 for item in items if item["primary_outcome"] == "downstream_ai_obtained"), followup_required=sum(1 for item in items if item["primary_outcome"] == "followup_required"), name_not_obtained=sum(1 for item in items if item["primary_outcome"] == "name_not_obtained"), manual_corrected=sum(1 for item in items if item.get("manual_corrected")), handoff_confirmed_name=sum(1 for item in items if item.get("handoff_confirmed_name")), handoff_unconfirmed_name=sum(1 for item in items if item.get("handoff_unconfirmed_name")), **breakdown_metrics, ) ) by_queue_rows: list[VoiceNameFlowAnalyticsQueueBreakdownOut] = [] queue_groups: dict[str, list[dict[str, Any]]] = defaultdict(list) for item in snapshots: queue_groups[str(item["queue_id"] or "unknown")].append(item) for resolved_queue_id, items in sorted(queue_groups.items(), key=lambda entry: (-len(entry[1]), entry[0])): breakdown_metrics = _voice_name_breakdown_metrics(items) by_queue_rows.append( VoiceNameFlowAnalyticsQueueBreakdownOut( queue_id=resolved_queue_id, scenario_calls=len(items), start_obtained=sum(1 for item in items if item["primary_outcome"] == "start_obtained"), downstream_ai_obtained=sum(1 for item in items if item["primary_outcome"] == "downstream_ai_obtained"), followup_required=sum(1 for item in items if item["primary_outcome"] == "followup_required"), name_not_obtained=sum(1 for item in items if item["primary_outcome"] == "name_not_obtained"), manual_corrected=sum(1 for item in items if item.get("manual_corrected")), handoff_confirmed_name=sum(1 for item in items if item.get("handoff_confirmed_name")), handoff_unconfirmed_name=sum(1 for item in items if item.get("handoff_unconfirmed_name")), **breakdown_metrics, ) ) all_handoffs = handoff_confirmed_name + handoff_unconfirmed_name handoff_rows = [ VoiceNameFlowAnalyticsHandoffBreakdownOut( outcome="confirmed_name", label="Handoff с подтверждённым именем", sessions=handoff_confirmed_name, share=_analytics_percent(handoff_confirmed_name, all_handoffs), ), VoiceNameFlowAnalyticsHandoffBreakdownOut( outcome="unconfirmed_name", label="Handoff без подтверждённого имени", sessions=handoff_unconfirmed_name, share=_analytics_percent(handoff_unconfirmed_name, all_handoffs), ), ] return VoiceNameFlowAnalyticsOverviewOut( window=_analytics_window(range_from, range_to), filters=_voice_name_filters(range_from, range_to, queue_id, language), totals=totals, metrics=metrics, breakdowns=VoiceNameFlowAnalyticsBreakdownsOut( funnel=funnel, by_language=by_language_rows, by_queue=by_queue_rows, handoff=handoff_rows, ), coverage=_voice_name_analytics_coverage_from_snapshots(snapshots), ) def _load_voice_name_flow_overview( session, *, range_from: datetime, range_to: datetime, queue_id: str | None, language: str | None, ) -> VoiceNameFlowAnalyticsOverviewOut: snapshots = _load_voice_name_flow_snapshots( session, range_from=range_from, range_to=range_to, queue_id=queue_id, language=language, ) return _aggregate_voice_name_flow_overview( snapshots, range_from=range_from, range_to=range_to, queue_id=queue_id, language=_normalize_voice_name_language(language), ) def _bool_env(name: str, default: bool) -> bool: raw = os.getenv(name) if raw is None: return default return raw.strip().lower() in {"1", "true", "yes", "on"} def _int_env(name: str, default: int) -> int: raw = os.getenv(name) if raw is None: return default try: return int(raw.strip()) except ValueError: return default def _float_env(name: str, default: float) -> float: raw = os.getenv(name) if raw is None: return default try: return float(raw.strip()) except ValueError: return default def _ai_provider() -> str: return os.getenv("AI_PROVIDER", "stub").strip() or "stub" def _ai_api_base() -> str: return os.getenv("AI_API_BASE", "").rstrip("/") def _ai_api_key() -> str: return os.getenv("AI_API_KEY", "").strip() def _ai_model() -> str: return os.getenv("AI_MODEL", "stub-telegram-assistant").strip() or "stub-telegram-assistant" def _load_ai_operator_config(session) -> AIOperatorConfig: return load_effective_ai_operator_config(session) def _ai_timeout_seconds() -> float: return max(3.0, _float_env("AI_TIMEOUT_SECONDS", 20.0)) def _ai_telegram_enabled() -> bool: return _bool_env("AI_TELEGRAM_ENABLED", False) def _ai_telegram_always_reply() -> bool: return _bool_env("AI_TELEGRAM_ALWAYS_REPLY", False) def _ai_whatsapp_enabled() -> bool: return _bool_env("AI_WHATSAPP_ENABLED", False) def _ai_whatsapp_always_reply() -> bool: return _bool_env("AI_WHATSAPP_ALWAYS_REPLY", False) def _ai_max_context_messages() -> int: return max(4, _int_env("AI_TELEGRAM_MAX_CONTEXT_MESSAGES", 12)) def _ai_whatsapp_max_context_messages() -> int: return max(4, _int_env("AI_WHATSAPP_MAX_CONTEXT_MESSAGES", 12)) def _ai_max_kb_results() -> int: return max(1, _int_env("AI_TELEGRAM_MAX_KB_RESULTS", 3)) def _ai_whatsapp_max_kb_results() -> int: return max(1, _int_env("AI_WHATSAPP_MAX_KB_RESULTS", 3)) def _ai_handoff_threshold() -> float: return max(0.0, min(1.0, _float_env("AI_TELEGRAM_CONFIDENCE_HANDOFF_THRESHOLD", 0.65))) def _ai_whatsapp_handoff_threshold() -> float: return max(0.0, min(1.0, _float_env("AI_WHATSAPP_CONFIDENCE_HANDOFF_THRESHOLD", 0.65))) def _interaction_service_url() -> str: return os.getenv("INTERACTION_SERVICE_URL", "http://localhost:8004").rstrip("/") def _telegram_service_url() -> str: return os.getenv("TELEGRAM_ADAPTER_SERVICE_URL", "http://localhost:8007").rstrip("/") def _whatsapp_service_url() -> str: return os.getenv("WHATSAPP_ADAPTER_SERVICE_URL", "http://localhost:8019").rstrip("/") def _service_headers() -> dict[str, str]: token = issue_app_token( subject="svc:ai-orchestrator", username="ai-orchestrator", role="admin", auth_source="service", provider="ai-orchestrator", ttl_seconds=300, ) return {"Authorization": f"Bearer {token}"} def _interaction_request(method: str, path: str, *, payload: dict | None = None) -> dict: with httpx.Client(timeout=10.0) as client: response = client.request( method, f"{_interaction_service_url()}{path}", json=payload, headers=_service_headers(), ) response.raise_for_status() return response.json() def _telegram_request(method: str, path: str, *, payload: dict | None = None) -> dict: with httpx.Client(timeout=10.0) as client: response = client.request( method, f"{_telegram_service_url()}{path}", json=payload, headers=_service_headers(), ) response.raise_for_status() return response.json() def _whatsapp_request(method: str, path: str, *, payload: dict | None = None) -> dict: with httpx.Client(timeout=10.0) as client: response = client.request( method, f"{_whatsapp_service_url()}{path}", json=payload, headers=_service_headers(), ) response.raise_for_status() return response.json() def _response_error_detail(response: httpx.Response) -> str: try: payload = response.json() except Exception: # noqa: BLE001 payload = None if isinstance(payload, dict): detail = payload.get("detail") if isinstance(detail, str) and detail.strip(): return detail.strip() text = response.text.strip() if text: return text return f"HTTP {response.status_code}" def _conflict_result_from_telegram_error( session, *, job: AIJobRow, ai_session: AISessionRow, thread_id: str, exc: Exception, reply_message_id: str | None = None, ) -> dict[str, Any] | None: if not isinstance(exc, httpx.HTTPStatusError) or exc.response.status_code != 409: return None now = utc_now_iso() thread = _thread_or_404(session, thread_id) detail = _response_error_detail(exc.response).lower() if thread.status == "closed" or thread.ai_state == "closed" or "closed" in detail: ai_session.status = "closed" ai_session.closed_at = ai_session.closed_at or now ai_session.updated_at = now _mark_job_done(session, job) session.commit() result: dict[str, Any] = {"ok": True, "status": "closed", "job_id": job.job_id} if reply_message_id: result["reply_message_id"] = reply_message_id return result if thread.claimed_by_user or thread.ai_state == "human_owned" or "human operator" in detail: ai_session.status = "human_owned" ai_session.updated_at = now ai_session.handoff_reason = thread.ai_handoff_reason _mark_job_done(session, job) session.commit() result = {"ok": True, "status": "human_owned", "job_id": job.job_id} if reply_message_id: result["reply_message_id"] = reply_message_id return result return None def _conflict_result_from_whatsapp_error( session, *, job: AIJobRow, ai_session: AISessionRow, thread_id: str, exc: Exception, reply_message_id: str | None = None, ) -> dict[str, Any] | None: if not isinstance(exc, httpx.HTTPStatusError) or exc.response.status_code != 409: return None now = utc_now_iso() thread = _whatsapp_thread_or_404(session, thread_id) detail = _response_error_detail(exc.response).lower() if thread.status == "closed" or thread.ai_state == "closed" or "closed" in detail: ai_session.status = "closed" ai_session.closed_at = ai_session.closed_at or now ai_session.updated_at = now _mark_job_done(session, job) session.commit() result: dict[str, Any] = {"ok": True, "status": "closed", "job_id": job.job_id} if reply_message_id: result["reply_message_id"] = reply_message_id return result if thread.claimed_by_user or thread.ai_state == "human_owned" or "human operator" in detail: ai_session.status = "human_owned" ai_session.updated_at = now ai_session.handoff_reason = thread.ai_handoff_reason _mark_job_done(session, job) session.commit() result = {"ok": True, "status": "human_owned", "job_id": job.job_id} if reply_message_id: result["reply_message_id"] = reply_message_id return result return None def _push_timeline(session, interaction_id: str, action: str, metadata: dict | None = None) -> None: session.add( InteractionTimeline( interaction_id=interaction_id, timestamp=utc_now_iso(), action=action, metadata_json=json.dumps(metadata or {}, ensure_ascii=False), ) ) def _normalize_external_subject(value: str | None) -> str | None: raw = str(value or "").strip() if not raw: return None if raw.startswith("telegram:") or raw.startswith("whatsapp:"): return raw.split(":", 1)[1].strip() or None return raw def _customer_id_is_real(customer_id: str | None) -> bool: return str(customer_id or "").startswith("cus_") def _telegram_identity_subjects( telegram_user_id: str | None, chat_id: str, explicit: str | None = None, ) -> list[str]: values = [ _normalize_external_subject(telegram_user_id), _normalize_external_subject(explicit), _normalize_external_subject(chat_id), ] seen: set[str] = set() result: list[str] = [] for value in values: if not value or value in seen: continue seen.add(value) result.append(value) return result def _whatsapp_identity_subjects( whatsapp_user_id: str | None, chat_id: str, *, phone_number: str | None = None, explicit: str | None = None, ) -> list[str]: values = [ _normalize_external_subject(whatsapp_user_id), _normalize_external_subject(phone_number), _normalize_external_subject(explicit), _normalize_external_subject(chat_id), ] seen: set[str] = set() result: list[str] = [] for value in values: if not value or value in seen: continue seen.add(value) result.append(value) return result def _ensure_customer_external_identity( session, *, customer_id: str, channel: str, external_subject: str, display_name_snapshot: str | None, now: str, ) -> None: row = session.execute( select(CustomerExternalIdentity).where( CustomerExternalIdentity.channel == channel, CustomerExternalIdentity.external_subject == external_subject, ) ).scalar_one_or_none() if row: row.customer_id = customer_id if display_name_snapshot: row.display_name_snapshot = display_name_snapshot row.updated_at = now return session.add( CustomerExternalIdentity( identity_id=new_id("cei"), customer_id=customer_id, channel=channel, external_subject=external_subject, display_name_snapshot=display_name_snapshot, created_at=now, updated_at=now, ) ) def _resolve_or_create_customer_id(session, thread: TelegramThreadRow, interaction: Interaction) -> str: if _customer_id_is_real(interaction.customer_id): return str(interaction.customer_id) now = utc_now_iso() legacy_subject = _normalize_external_subject(interaction.customer_id) subjects = _telegram_identity_subjects(thread.telegram_user_id, thread.chat_id, legacy_subject) for subject in subjects: identity = session.execute( select(CustomerExternalIdentity).where( CustomerExternalIdentity.channel == "telegram", CustomerExternalIdentity.external_subject == subject, ) ).scalar_one_or_none() if identity: if thread.display_name and identity.display_name_snapshot != thread.display_name: identity.display_name_snapshot = thread.display_name identity.updated_at = now interaction.customer_id = identity.customer_id interaction.updated_at = now return identity.customer_id customer = Customer( customer_id=new_id("cus"), display_name=thread.display_name or thread.username or f"Telegram {thread.chat_id}", phones_json="[]", preferred_phone=None, tags_json=json.dumps(["telegram"], ensure_ascii=False), created_at=now, ) session.add(customer) interaction.customer_id = customer.customer_id interaction.updated_at = now for subject in subjects or [_normalize_external_subject(thread.chat_id) or thread.chat_id]: _ensure_customer_external_identity( session, customer_id=customer.customer_id, channel="telegram", external_subject=subject, display_name_snapshot=thread.display_name, now=now, ) return customer.customer_id def _resolve_or_create_whatsapp_customer_id(session, thread: WhatsAppThreadRow, interaction: Interaction) -> str: if _customer_id_is_real(interaction.customer_id): return str(interaction.customer_id) now = utc_now_iso() legacy_subject = _normalize_external_subject(interaction.customer_id) subjects = _whatsapp_identity_subjects( thread.whatsapp_user_id, thread.chat_id, phone_number=thread.phone_number, explicit=legacy_subject, ) for subject in subjects: identity = session.execute( select(CustomerExternalIdentity).where( CustomerExternalIdentity.channel == "whatsapp", CustomerExternalIdentity.external_subject == subject, ) ).scalar_one_or_none() if identity: if thread.display_name and identity.display_name_snapshot != thread.display_name: identity.display_name_snapshot = thread.display_name identity.updated_at = now interaction.customer_id = identity.customer_id interaction.updated_at = now return identity.customer_id phones = [thread.phone_number] if thread.phone_number else [] customer = Customer( customer_id=new_id("cus"), display_name=thread.display_name or thread.phone_number or f"WhatsApp {thread.chat_id}", phones_json=json.dumps(phones, ensure_ascii=False), preferred_phone=thread.phone_number, tags_json=json.dumps(["whatsapp"], ensure_ascii=False), created_at=now, ) session.add(customer) interaction.customer_id = customer.customer_id interaction.updated_at = now for subject in subjects or [_normalize_external_subject(thread.chat_id) or thread.chat_id]: _ensure_customer_external_identity( session, customer_id=customer.customer_id, channel="whatsapp", external_subject=subject, display_name_snapshot=thread.display_name, now=now, ) return customer.customer_id def _infer_language(text: str) -> str: source = str(text or "").lower() if re.search(r"[әіңғүұқөһ]", source): return "kz" kz_keywords = ("сәлем", "көмек", "рақмет", "өтінемін", "керек", "қалай", "қайырлы") ru_keywords = ("здравствуйте", "помощь", "оператор", "заявка", "тариф") kz_hits = sum(token in source for token in kz_keywords) ru_hits = sum(token in source for token in ru_keywords) if kz_hits > ru_hits and kz_hits > 0: return "kz" if ru_hits > kz_hits and ru_hits > 0: return "ru" if re.search(r"[әіңғүұқөһ]", source): return "kz" if any(token in source for token in ("сәлем", "көмек", "рақмет", "өтінемін", "баға", "тариф")): return "kz" return "ru" def _thread_or_404(session, thread_id: str) -> TelegramThreadRow: row = session.execute( select(TelegramThreadRow).where(TelegramThreadRow.thread_id == thread_id) ).scalar_one_or_none() if not row: raise HTTPException(status_code=404, detail="Telegram thread not found") return row def _whatsapp_thread_or_404(session, thread_id: str) -> WhatsAppThreadRow: row = session.execute( select(WhatsAppThreadRow).where(WhatsAppThreadRow.thread_id == thread_id) ).scalar_one_or_none() if not row: raise HTTPException(status_code=404, detail="WhatsApp thread not found") return row def _article_snippet(article: KBArticleRow, limit: int = 240) -> str: raw = f"{article.title}. {article.body}".strip() if len(raw) <= limit: return raw return f"{raw[: limit - 3]}..." def _kb_search(session, text: str, *, language: str | None = None) -> list[KBArticleRow]: if not str(text or "").strip(): return [] stmt = select(KBArticleRow).order_by(KBArticleRow.id.desc()) if language is not None: stmt = stmt.where(KBArticleRow.language == normalize_kb_language(language)) rows = session.execute(stmt).scalars().all() return search_kb_rows(rows, text, limit=_ai_max_kb_results()) def _last_messages(session, thread_id: str, limit: int) -> list[TelegramMessageRow]: rows = session.execute( select(TelegramMessageRow) .where(TelegramMessageRow.thread_id == thread_id) .order_by(TelegramMessageRow.id.desc()) .limit(max(limit, 1)) ).scalars().all() return list(reversed(rows)) def _last_whatsapp_messages(session, thread_id: str, limit: int) -> list[WhatsAppMessageRow]: rows = session.execute( select(WhatsAppMessageRow) .where(WhatsAppMessageRow.thread_id == thread_id) .order_by(WhatsAppMessageRow.id.desc()) .limit(max(limit, 1)) ).scalars().all() return list(reversed(rows)) def _select_trigger_message( session, *, thread_id: str, trigger_message_id: str | None, ) -> TelegramMessageRow | None: if trigger_message_id: direct = session.execute( select(TelegramMessageRow).where(TelegramMessageRow.message_id == trigger_message_id) ).scalar_one_or_none() if direct and direct.thread_id == thread_id and direct.author_type == "customer": return direct rows = session.execute( select(TelegramMessageRow) .where( TelegramMessageRow.thread_id == thread_id, TelegramMessageRow.author_type == "customer", ) .order_by(TelegramMessageRow.id.desc()) ).scalars().all() return rows[0] if rows else None def _select_whatsapp_trigger_message( session, *, thread_id: str, trigger_message_id: str | None, ) -> WhatsAppMessageRow | None: if trigger_message_id: direct = session.execute( select(WhatsAppMessageRow).where(WhatsAppMessageRow.message_id == trigger_message_id) ).scalar_one_or_none() if direct and direct.thread_id == thread_id and direct.author_type == "customer": return direct rows = session.execute( select(WhatsAppMessageRow) .where( WhatsAppMessageRow.thread_id == thread_id, WhatsAppMessageRow.author_type == "customer", ) .order_by(WhatsAppMessageRow.id.desc()) ).scalars().all() return rows[0] if rows else None def _existing_job_for_thread(session, thread_id: str) -> AIJobRow | None: rows = session.execute( select(AIJobRow) .where( AIJobRow.thread_id == thread_id, AIJobRow.status.in_(["pending", "running"]), ) .order_by(AIJobRow.id.desc()) ).scalars().all() return rows[0] if rows else None def _latest_job_for_trigger(session, thread_id: str, trigger_message_id: str | None) -> AIJobRow | None: if not trigger_message_id: return None rows = session.execute( select(AIJobRow) .where( AIJobRow.thread_id == thread_id, AIJobRow.trigger_message_id == trigger_message_id, ) .order_by(AIJobRow.id.desc()) ).scalars().all() return rows[0] if rows else None def _ensure_ai_session( session, *, thread: TelegramThreadRow, interaction: Interaction, customer_id: str, language: str, ) -> tuple[AISessionRow, bool]: existing = None if thread.ai_session_id: existing = session.execute( select(AISessionRow).where(AISessionRow.session_id == thread.ai_session_id) ).scalar_one_or_none() if existing and existing.status not in {"closed", "error"}: existing.customer_id = customer_id existing.interaction_id = interaction.interaction_id existing.language = language or existing.language existing.updated_at = utc_now_iso() return existing, False now = utc_now_iso() ai_session = AISessionRow( session_id=new_id("ais"), channel="telegram", thread_id=thread.thread_id, interaction_id=interaction.interaction_id, customer_id=customer_id, agent_profile="telegram_support", language=language, status="active", summary_text="", last_user_message_id=None, last_ai_message_id=None, handoff_reason=None, created_at=now, updated_at=now, closed_at=None, ) session.add(ai_session) thread.ai_session_id = ai_session.session_id thread.ai_state = "queued" thread.ai_handoff_reason = None thread.updated_at = now _push_timeline( session, interaction.interaction_id, "ai.session_started", {"thread_id": thread.thread_id, "session_id": ai_session.session_id, "language": language}, ) return ai_session, True def _ensure_job( session, *, thread: TelegramThreadRow, ai_session: AISessionRow, trigger_message_id: str | None, ) -> tuple[AIJobRow | None, bool]: if trigger_message_id and ai_session.last_user_message_id == trigger_message_id: existing = _latest_job_for_trigger(session, thread.thread_id, trigger_message_id) if existing: return existing, True existing = _existing_job_for_thread(session, thread.thread_id) if existing: return existing, True now = utc_now_iso() job = AIJobRow( job_id=new_id("aij"), session_id=ai_session.session_id, thread_id=thread.thread_id, trigger_message_id=trigger_message_id, status="pending", attempts=0, next_attempt_at=now, locked_until=None, last_error=None, created_at=now, updated_at=now, ) session.add(job) return job, False def _ensure_whatsapp_ai_session( session, *, thread: WhatsAppThreadRow, interaction: Interaction, customer_id: str, language: str, ) -> tuple[AISessionRow, bool]: existing = None if thread.ai_session_id: existing = session.execute( select(AISessionRow).where(AISessionRow.session_id == thread.ai_session_id) ).scalar_one_or_none() if existing and existing.status not in {"closed", "error"}: existing.customer_id = customer_id existing.interaction_id = interaction.interaction_id existing.language = language or existing.language existing.updated_at = utc_now_iso() return existing, False now = utc_now_iso() ai_session = AISessionRow( session_id=new_id("ais"), channel="whatsapp", thread_id=thread.thread_id, interaction_id=interaction.interaction_id, customer_id=customer_id, agent_profile="whatsapp_support", language=language, status="active", summary_text="", last_user_message_id=None, last_ai_message_id=None, handoff_reason=None, created_at=now, updated_at=now, closed_at=None, ) session.add(ai_session) thread.ai_session_id = ai_session.session_id thread.ai_state = "queued" thread.ai_handoff_reason = None thread.updated_at = now _push_timeline( session, interaction.interaction_id, "ai.session_started", {"thread_id": thread.thread_id, "session_id": ai_session.session_id, "language": language}, ) return ai_session, True def _looks_like_human_request(text: str) -> bool: source = str(text or "").lower() patterns = [ "оператор", "человек", "менеджер", "живой", "переведи", "transfer", "human", "сотрудник", "специалист", ] return any(token in source for token in patterns) def _is_sensitive_request(text: str) -> bool: source = str(text or "").lower() tokens = [ "жалоб", "претензи", "верните деньги", "деньги", "оплат", "договор", "суд", "юрист", "паспорт", "iin", "бин", "конфиден", ] return any(token in source for token in tokens) def _looks_like_resolution_confirmation(text: str) -> bool: source = str(text or "").lower() tokens = [ "спасибо", "решено", "все понятно", "всё понятно", "не нужно", "ок", "хорошо", "рахмет", "түсінікті", ] return any(token in source for token in tokens) def _extract_json_object(raw: str) -> dict[str, Any]: text = str(raw or "").strip() if not text: raise ValueError("Model returned empty response") try: parsed = json.loads(text) if isinstance(parsed, dict): return parsed except json.JSONDecodeError: pass match = re.search(r"\{.*\}", text, re.DOTALL) if not match: raise ValueError("Model response does not contain JSON object") parsed = json.loads(match.group(0)) if not isinstance(parsed, dict): raise ValueError("Model response JSON is not an object") return parsed def _openai_prompt( *, customer: Customer | None, interaction: Interaction, thread: Any, messages: list[Any], kb_results: list[KBArticleRow], language: str, conversation_summary_text: str = "", channel_label: str = "Telegram", channel_key: str = "telegram", operator_config: AIOperatorConfig | None = None, ) -> list[dict[str, str]]: customer_summary = { "customer_id": customer.customer_id if customer else interaction.customer_id, "display_name": customer.display_name if customer else thread.display_name, "tags": json.loads(customer.tags_json or "[]") if customer else [], "channel": channel_key, } history = [ { "author_type": item.author_type, "author_id": item.author_id, "direction": item.direction, "text": item.text, "created_at": item.created_at, } for item in messages ] kb_context = [ { "article_id": article.article_id, "title": article.title, "snippet": _article_snippet(article), } for article in kb_results ] system_prompt = persona.operator_system_prompt( language=language, channel_label=channel_label, is_voice=False, config=operator_config, ) user_prompt = { "customer": customer_summary, "interaction": { "interaction_id": interaction.interaction_id, "status": interaction.status, "queue_id": interaction.queue_id, "subject": interaction.subject, }, "thread": { "thread_id": thread.thread_id, "chat_id": thread.chat_id, "display_name": thread.display_name, }, "conversation_summary": conversation_summary_text, "kb_results": kb_context, "history": history, } return [ {"role": "system", "content": system_prompt}, {"role": "user", "content": json.dumps(user_prompt, ensure_ascii=False)}, ] def _request_structured_model_decision(messages: list[dict[str, str]]) -> dict[str, Any]: if not _ai_api_base() or not _ai_api_key(): raise RuntimeError("AI_API_BASE / AI_API_KEY are required for openai_compatible provider") started = time.perf_counter() payload = { "model": _ai_model(), "temperature": 0.2, "response_format": {"type": "json_object"}, "messages": messages, } with httpx.Client(timeout=_ai_timeout_seconds()) as client: response = client.post( f"{_ai_api_base()}/chat/completions", json=payload, headers={"Authorization": f"Bearer {_ai_api_key()}"}, ) response.raise_for_status() result = response.json() choice = ((result.get("choices") or [{}])[0] if isinstance(result.get("choices"), list) else {}) or {} message = choice.get("message") if isinstance(choice.get("message"), dict) else {} raw_content = message.get("content") or "" decision = _extract_json_object(str(raw_content)) decision["_model"] = result.get("model") or _ai_model() decision["_latency_ms"] = int((time.perf_counter() - started) * 1000) decision["_finish_reason"] = choice.get("finish_reason") return decision def _openai_compatible_decision( *, customer: Customer | None, interaction: Interaction, thread: Any, messages: list[Any], kb_results: list[KBArticleRow], language: str, conversation_summary_text: str = "", channel_label: str = "Telegram", channel_key: str = "telegram", operator_config: AIOperatorConfig | None = None, ) -> dict[str, Any]: return _request_structured_model_decision( _openai_prompt( customer=customer, interaction=interaction, thread=thread, messages=messages, kb_results=kb_results, language=language, conversation_summary_text=conversation_summary_text, channel_label=channel_label, channel_key=channel_key, operator_config=operator_config, ) ) def _sanitize_decision(raw: dict[str, Any], *, fallback_language: str) -> dict[str, Any]: decision = { "language": str(raw.get("language") or fallback_language or "ru"), "intent": str(raw.get("intent") or "unknown"), "reply_text": str(raw.get("reply_text") or "").strip(), "extracted_name": str(raw.get("extracted_name") or "").strip() or None, "confidence": float(raw.get("confidence") or 0.0), "needs_handoff": bool(raw.get("needs_handoff")), "handoff_reason": str(raw.get("handoff_reason") or "").strip() or None, "case_action": str(raw.get("case_action") or "keep_open"), "kb_refs": [str(item) for item in (raw.get("kb_refs") or []) if str(item).strip()], "_model": raw.get("_model") or _ai_model(), "_finish_reason": raw.get("_finish_reason"), "_latency_ms": int(raw.get("_latency_ms") or 0), } if decision["case_action"] not in {"none", "close", "escalate", "keep_open"}: decision["case_action"] = "keep_open" if decision["confidence"] < 0: decision["confidence"] = 0.0 if decision["confidence"] > 1: decision["confidence"] = 1.0 return decision def _stub_decision( *, customer: Customer | None, interaction: Interaction, last_user_message: TelegramMessageRow, kb_results: list[KBArticleRow], language: str, operator_config: AIOperatorConfig | None = None, ) -> dict[str, Any]: text = last_user_message.text if persona.is_identity_request(text): return { "language": language, "intent": "identity_question", "reply_text": persona.identity_reply(language, operator_config), "confidence": 0.98, "needs_handoff": False, "handoff_reason": None, "case_action": "keep_open", "kb_refs": [], } if _looks_like_human_request(text): return { "language": language, "intent": "handoff_request", "reply_text": "", "confidence": 0.2, "needs_handoff": True, "handoff_reason": "Клиент запросил живого оператора.", "case_action": "keep_open", "kb_refs": [], } if _is_sensitive_request(text): return { "language": language, "intent": "sensitive_request", "reply_text": "", "confidence": 0.25, "needs_handoff": True, "handoff_reason": "Нужен человек: запрос затрагивает чувствительную тему или действие вне доступного контекста.", "case_action": "keep_open", "kb_refs": [], } if _looks_like_resolution_confirmation(text): return { "language": language, "intent": "resolution_confirmed", "reply_text": persona.text_resolution_reply(language), "confidence": 0.9, "needs_handoff": False, "handoff_reason": None, "case_action": "close", "kb_refs": [], } if kb_results: best = kb_results[0] snippet = _article_snippet(best) return { "language": language, "intent": "kb_answer", "reply_text": ( f"Қысқаша айтайын: {snippet}" if language == "kz" else f"Коротко подскажу: {snippet}" ), "confidence": 0.84, "needs_handoff": False, "handoff_reason": None, "case_action": "keep_open", "kb_refs": [best.article_id], } return { "language": language, "intent": "clarification", "reply_text": ( "Нақтырақ айтыңызшы: нені тексеру немесе нақтылау керек?" if language == "kz" else "Уточните, пожалуйста, что именно нужно проверить или подсказать." ), "confidence": 0.45, "needs_handoff": False, "handoff_reason": None, "case_action": "keep_open", "kb_refs": [], } def _always_reply_fallback(last_user_text: str, language: str) -> str: normalized = str(last_user_text or "").strip().lower() greeting_tokens = ("привет", "здравствуйте", "добрый", "салем", "сә", "сәлем", "hello", "hi") return persona.human_fallback_reply( language, is_greeting=any(token in normalized for token in greeting_tokens), ) def _apply_always_reply_mode( decision: dict[str, Any], *, last_user_text: str, enabled: bool | None = None, handoff_threshold: float | None = None, ) -> dict[str, Any]: if enabled is None: enabled = _ai_telegram_always_reply() if handoff_threshold is None: handoff_threshold = _ai_handoff_threshold() if not enabled: return decision decision["needs_handoff"] = False decision["handoff_reason"] = None decision["case_action"] = "keep_open" if not str(decision.get("reply_text") or "").strip(): decision["reply_text"] = _always_reply_fallback(last_user_text, decision.get("language") or "ru") if float(decision.get("confidence") or 0.0) < handoff_threshold: decision["confidence"] = max(handoff_threshold, 0.7) return decision def _decide_reply( *, customer: Customer | None, interaction: Interaction, thread: Any, messages: list[Any], kb_results: list[KBArticleRow], language: str, conversation_summary_text: str = "", channel_label: str = "Telegram", channel_key: str = "telegram", operator_config: AIOperatorConfig | None = None, ) -> dict[str, Any]: last_user_message = next((item for item in reversed(messages) if item.author_type == "customer"), None) if last_user_message is None: raise RuntimeError("No customer message found for AI decision") if persona.is_identity_request(last_user_message.text): return _sanitize_decision( { "language": language, "intent": "identity_question", "reply_text": persona.identity_reply(language, operator_config), "confidence": 0.98, "needs_handoff": False, "handoff_reason": None, "case_action": "keep_open", "kb_refs": [], "_model": "operator_identity_policy", "_latency_ms": 1, "_finish_reason": "stop", }, fallback_language=language, ) if _ai_provider() == "openai_compatible": raw = _openai_compatible_decision( customer=customer, interaction=interaction, thread=thread, messages=messages, kb_results=kb_results, language=language, conversation_summary_text=conversation_summary_text, channel_label=channel_label, channel_key=channel_key, operator_config=operator_config, ) else: raw = _stub_decision( customer=customer, interaction=interaction, last_user_message=last_user_message, kb_results=kb_results, language=language, operator_config=operator_config, ) raw["_model"] = _ai_model() raw["_latency_ms"] = 1 raw["_finish_reason"] = "stop" return _sanitize_decision(raw, fallback_language=language) def _update_ai_session_context_summary_from_user_turn( ai_session: AISessionRow, *, channel: str, language: str, customer_name: str | None, text: str, now: str, ) -> str: summary = update_context_summary_from_user_turn( ai_session.context_summary_json, channel=channel, language=language, customer_name=customer_name, text=text, now=now, ) ai_session.context_summary_json = dump_context_summary(summary) ai_session.context_summary_updated_at = now return render_context_summary_text(summary) def _update_ai_session_context_summary_from_assistant_turn( ai_session: AISessionRow, *, language: str, customer_name: str | None, reply_text: str, decision: dict[str, Any], now: str, ) -> None: summary = update_context_summary_from_assistant_turn( ai_session.context_summary_json, language=language, customer_name=customer_name, reply_text=reply_text, decision_intent=str(decision.get("intent") or "").strip() or None, needs_handoff=bool(decision.get("needs_handoff")), handoff_reason=str(decision.get("handoff_reason") or "").strip() or None, now=now, ) ai_session.context_summary_json = dump_context_summary(summary) ai_session.context_summary_updated_at = now def _update_thread_after_close(session, thread: TelegramThreadRow, when: str) -> None: thread.status = "closed" thread.ai_state = "closed" thread.ai_handoff_reason = None thread.ai_last_model_at = when thread.updated_at = when def _update_whatsapp_thread_after_close(session, thread: WhatsAppThreadRow, when: str) -> None: thread.status = "closed" thread.ai_state = "closed" thread.ai_handoff_reason = None thread.ai_last_model_at = when thread.updated_at = when def _mark_job_done(session, job: AIJobRow, *, error: str | None = None) -> None: job.status = "failed" if error else "done" job.last_error = error job.updated_at = utc_now_iso() def _record_turn( session, *, session_id: str, thread_id: str, interaction_id: str, role: str, source_type: str, text: str, payload: dict[str, Any], model: str | None = None, finish_reason: str | None = None, latency_ms: int | None = None, ) -> None: session.add( AITurnRow( turn_id=new_id("ait"), session_id=session_id, thread_id=thread_id, interaction_id=interaction_id, role=role, source_type=source_type, text=text, payload_json=json.dumps(payload, ensure_ascii=False), model=model, finish_reason=finish_reason, latency_ms=latency_ms, created_at=utc_now_iso(), ) ) def _process_job(job_id: str) -> dict[str, Any]: session = get_session() try: job = session.execute(select(AIJobRow).where(AIJobRow.job_id == job_id)).scalar_one_or_none() if not job: raise HTTPException(status_code=404, detail="AI job not found") thread = _thread_or_404(session, job.thread_id) interaction = session.execute( select(Interaction).where(Interaction.interaction_id == thread.interaction_id) ).scalar_one() customer_id = _resolve_or_create_customer_id(session, thread, interaction) customer = session.execute( select(Customer).where(Customer.customer_id == customer_id) ).scalar_one_or_none() trigger_message = _select_trigger_message( session, thread_id=thread.thread_id, trigger_message_id=job.trigger_message_id, ) if not trigger_message: _mark_job_done(session, job, error="No trigger customer message found") session.commit() return {"ok": False, "status": job.status, "job_id": job.job_id} ai_session = session.execute( select(AISessionRow).where(AISessionRow.session_id == job.session_id) ).scalar_one() if thread.ai_state == "human_owned" or thread.claimed_by_user: ai_session.status = "human_owned" ai_session.updated_at = utc_now_iso() _mark_job_done(session, job) session.commit() return {"ok": True, "status": "human_owned", "job_id": job.job_id} job.status = "running" job.attempts = int(job.attempts or 0) + 1 job.updated_at = utc_now_iso() thread.ai_state = "thinking" thread.ai_handoff_reason = None thread.updated_at = utc_now_iso() ai_session.status = "active" ai_session.language = _infer_language(trigger_message.text) ai_session.customer_id = customer_id ai_session.last_user_message_id = trigger_message.message_id ai_session.updated_at = utc_now_iso() user_turn_now = ai_session.updated_at _record_turn( session, session_id=ai_session.session_id, thread_id=thread.thread_id, interaction_id=interaction.interaction_id, role="user", source_type="telegram", text=trigger_message.text, payload={"message_id": trigger_message.message_id, "author_type": trigger_message.author_type}, ) conversation_summary_text = _update_ai_session_context_summary_from_user_turn( ai_session, channel="telegram", language=ai_session.language or "ru", customer_name=customer.display_name if customer else thread.display_name, text=trigger_message.text, now=user_turn_now, ) session.commit() messages = _last_messages(session, thread.thread_id, _ai_max_context_messages()) kb_results = _kb_search(session, trigger_message.text, language=ai_session.language) operator_config = _load_ai_operator_config(session) decision = _decide_reply( customer=customer, interaction=interaction, thread=thread, messages=messages, kb_results=kb_results, language=ai_session.language or "ru", conversation_summary_text=conversation_summary_text, operator_config=operator_config, ) decision = _apply_always_reply_mode(decision, last_user_text=trigger_message.text) _record_turn( session, session_id=ai_session.session_id, thread_id=thread.thread_id, interaction_id=interaction.interaction_id, role="assistant", source_type="model", text=decision["reply_text"] or (decision["handoff_reason"] or decision["intent"]), payload=decision, model=decision["_model"], finish_reason=decision["_finish_reason"], latency_ms=decision["_latency_ms"], ) decision_now = utc_now_iso() ai_session.updated_at = decision_now needs_handoff = False if not _ai_telegram_always_reply(): needs_handoff = ( bool(decision["needs_handoff"]) or float(decision["confidence"]) < _ai_handoff_threshold() or _looks_like_human_request(trigger_message.text) or _is_sensitive_request(trigger_message.text) ) if needs_handoff: reason = decision["handoff_reason"] or "AI передаёт диалог оператору." try: _telegram_request( "POST", f"/integrations/telegram/threads/{thread.thread_id}/ai/handoff", payload={ "reason": reason, "agent_profile": ai_session.agent_profile, "trigger_message_id": trigger_message.message_id, "confidence": decision["confidence"], "payload": {"intent": decision["intent"], "kb_refs": decision["kb_refs"]}, }, ) except Exception as exc: # noqa: BLE001 resolved = _conflict_result_from_telegram_error( session, job=job, ai_session=ai_session, thread_id=thread.thread_id, exc=exc, ) if resolved is not None: return resolved raise ai_session.status = "handoff_required" ai_session.handoff_reason = reason ai_session.updated_at = utc_now_iso() ai_session.summary_text = decision["reply_text"] or (reason or ai_session.summary_text) _update_ai_session_context_summary_from_assistant_turn( ai_session, language=decision["language"], customer_name=customer.display_name if customer else thread.display_name, reply_text=decision["reply_text"] or reason, decision=decision, now=ai_session.updated_at, ) _mark_job_done(session, job) session.commit() return {"ok": True, "status": "handoff_required", "job_id": job.job_id} try: reply_payload = _telegram_request( "POST", f"/integrations/telegram/threads/{thread.thread_id}/ai/reply", payload={ "text": decision["reply_text"], "agent_profile": ai_session.agent_profile, "model": decision["_model"], "trigger_message_id": trigger_message.message_id, "language": decision["language"], "confidence": decision["confidence"], "kb_refs": decision["kb_refs"], "payload": {"intent": decision["intent"]}, }, ) except Exception as exc: # noqa: BLE001 resolved = _conflict_result_from_telegram_error( session, job=job, ai_session=ai_session, thread_id=thread.thread_id, exc=exc, ) if resolved is not None: return resolved raise ai_session.updated_at = utc_now_iso() ai_session.summary_text = decision["reply_text"] or (decision["handoff_reason"] or ai_session.summary_text) _update_ai_session_context_summary_from_assistant_turn( ai_session, language=decision["language"], customer_name=customer.display_name if customer else thread.display_name, reply_text=decision["reply_text"], decision=decision, now=ai_session.updated_at, ) if decision["case_action"] == "escalate": _interaction_request( "POST", f"/interactions/{interaction.interaction_id}/escalate", payload={"target_queue_id": interaction.queue_id or thread.queue_id or "q_telegram"}, ) try: _telegram_request( "POST", f"/integrations/telegram/threads/{thread.thread_id}/ai/handoff", payload={ "reason": decision["handoff_reason"] or "Нужна передача оператору по результатам AI-анализа.", "agent_profile": ai_session.agent_profile, "trigger_message_id": trigger_message.message_id, "confidence": decision["confidence"], "payload": {"intent": decision["intent"], "kb_refs": decision["kb_refs"]}, }, ) except Exception as exc: # noqa: BLE001 resolved = _conflict_result_from_telegram_error( session, job=job, ai_session=ai_session, thread_id=thread.thread_id, exc=exc, reply_message_id=reply_payload.get("message_id"), ) if resolved is not None: return resolved raise elif decision["case_action"] == "close" and _looks_like_resolution_confirmation(trigger_message.text): _interaction_request( "PATCH", f"/interactions/{interaction.interaction_id}/status", payload={"status": "closed"}, ) thread = _thread_or_404(session, thread.thread_id) _update_thread_after_close(session, thread, utc_now_iso()) ai_session.status = "closed" ai_session.closed_at = utc_now_iso() ai_session.updated_at = utc_now_iso() _push_timeline( session, interaction.interaction_id, "ai.case_closed", {"thread_id": thread.thread_id, "reply_message_id": reply_payload.get("message_id")}, ) _mark_job_done(session, job) session.commit() return {"ok": True, "status": "done", "job_id": job.job_id, "reply_message_id": reply_payload.get("message_id")} except HTTPException: raise except Exception as exc: # noqa: BLE001 logger.exception("AI Telegram job failed", extra={"job_id": job_id}) try: session.rollback() job = session.execute(select(AIJobRow).where(AIJobRow.job_id == job_id)).scalar_one_or_none() if job: job.status = "failed" job.last_error = str(exc)[:1000] job.updated_at = utc_now_iso() thread = _thread_or_404(session, job.thread_id) thread.ai_state = "error" thread.ai_handoff_reason = str(exc)[:240] thread.updated_at = utc_now_iso() interaction = session.execute( select(Interaction).where(Interaction.interaction_id == thread.interaction_id) ).scalar_one_or_none() if interaction: _push_timeline( session, interaction.interaction_id, "ai.error", {"thread_id": thread.thread_id, "job_id": job.job_id, "error": str(exc)[:500]}, ) session.commit() finally: pass raise HTTPException(status_code=502, detail=f"AI Telegram processing failed: {exc}") from exc finally: session.close() def _process_whatsapp_job(job_id: str) -> dict[str, Any]: session = get_session() try: job = session.execute(select(AIJobRow).where(AIJobRow.job_id == job_id)).scalar_one_or_none() if not job: raise HTTPException(status_code=404, detail="AI job not found") thread = _whatsapp_thread_or_404(session, job.thread_id) interaction = session.execute( select(Interaction).where(Interaction.interaction_id == thread.interaction_id) ).scalar_one() customer_id = _resolve_or_create_whatsapp_customer_id(session, thread, interaction) customer = session.execute( select(Customer).where(Customer.customer_id == customer_id) ).scalar_one_or_none() trigger_message = _select_whatsapp_trigger_message( session, thread_id=thread.thread_id, trigger_message_id=job.trigger_message_id, ) if not trigger_message: _mark_job_done(session, job, error="No trigger customer message found") session.commit() return {"ok": False, "status": job.status, "job_id": job.job_id} ai_session = session.execute( select(AISessionRow).where(AISessionRow.session_id == job.session_id) ).scalar_one() if thread.ai_state == "human_owned" or thread.claimed_by_user: ai_session.status = "human_owned" ai_session.updated_at = utc_now_iso() _mark_job_done(session, job) session.commit() return {"ok": True, "status": "human_owned", "job_id": job.job_id} job.status = "running" job.attempts = int(job.attempts or 0) + 1 job.updated_at = utc_now_iso() thread.ai_state = "thinking" thread.ai_handoff_reason = None thread.updated_at = utc_now_iso() ai_session.status = "active" ai_session.language = _infer_language(trigger_message.text) ai_session.customer_id = customer_id ai_session.last_user_message_id = trigger_message.message_id ai_session.updated_at = utc_now_iso() user_turn_now = ai_session.updated_at _record_turn( session, session_id=ai_session.session_id, thread_id=thread.thread_id, interaction_id=interaction.interaction_id, role="user", source_type="whatsapp", text=trigger_message.text, payload={"message_id": trigger_message.message_id, "author_type": trigger_message.author_type}, ) conversation_summary_text = _update_ai_session_context_summary_from_user_turn( ai_session, channel="whatsapp", language=ai_session.language or "ru", customer_name=customer.display_name if customer else thread.display_name, text=trigger_message.text, now=user_turn_now, ) session.commit() messages = _last_whatsapp_messages(session, thread.thread_id, _ai_whatsapp_max_context_messages()) kb_results = _kb_search( session, trigger_message.text, language=ai_session.language, )[: _ai_whatsapp_max_kb_results()] operator_config = _load_ai_operator_config(session) decision = _decide_reply( customer=customer, interaction=interaction, thread=thread, messages=messages, kb_results=kb_results, language=ai_session.language or "ru", conversation_summary_text=conversation_summary_text, channel_label="WhatsApp", channel_key="whatsapp", operator_config=operator_config, ) decision = _apply_always_reply_mode( decision, last_user_text=trigger_message.text, enabled=_ai_whatsapp_always_reply(), handoff_threshold=_ai_whatsapp_handoff_threshold(), ) _record_turn( session, session_id=ai_session.session_id, thread_id=thread.thread_id, interaction_id=interaction.interaction_id, role="assistant", source_type="model", text=decision["reply_text"] or (decision["handoff_reason"] or decision["intent"]), payload=decision, model=decision["_model"], finish_reason=decision["_finish_reason"], latency_ms=decision["_latency_ms"], ) ai_session.updated_at = utc_now_iso() needs_handoff = False if not _ai_whatsapp_always_reply(): needs_handoff = ( bool(decision["needs_handoff"]) or float(decision["confidence"]) < _ai_whatsapp_handoff_threshold() or _looks_like_human_request(trigger_message.text) or _is_sensitive_request(trigger_message.text) ) if needs_handoff: reason = decision["handoff_reason"] or "AI передаёт диалог оператору." try: _whatsapp_request( "POST", f"/integrations/whatsapp/threads/{thread.thread_id}/ai/handoff", payload={ "reason": reason, "agent_profile": ai_session.agent_profile, "trigger_message_id": trigger_message.message_id, "confidence": decision["confidence"], "payload": {"intent": decision["intent"], "kb_refs": decision["kb_refs"]}, }, ) except Exception as exc: # noqa: BLE001 resolved = _conflict_result_from_whatsapp_error( session, job=job, ai_session=ai_session, thread_id=thread.thread_id, exc=exc, ) if resolved is not None: return resolved raise ai_session.status = "handoff_required" ai_session.handoff_reason = reason ai_session.updated_at = utc_now_iso() ai_session.summary_text = decision["reply_text"] or (reason or ai_session.summary_text) _update_ai_session_context_summary_from_assistant_turn( ai_session, language=decision["language"], customer_name=customer.display_name if customer else thread.display_name, reply_text=decision["reply_text"] or reason, decision=decision, now=ai_session.updated_at, ) _mark_job_done(session, job) session.commit() return {"ok": True, "status": "handoff_required", "job_id": job.job_id} try: reply_payload = _whatsapp_request( "POST", f"/integrations/whatsapp/threads/{thread.thread_id}/ai/reply", payload={ "text": decision["reply_text"], "agent_profile": ai_session.agent_profile, "model": decision["_model"], "trigger_message_id": trigger_message.message_id, "language": decision["language"], "confidence": decision["confidence"], "kb_refs": decision["kb_refs"], "payload": {"intent": decision["intent"]}, }, ) except Exception as exc: # noqa: BLE001 resolved = _conflict_result_from_whatsapp_error( session, job=job, ai_session=ai_session, thread_id=thread.thread_id, exc=exc, ) if resolved is not None: return resolved raise ai_session.updated_at = utc_now_iso() ai_session.summary_text = decision["reply_text"] or (decision["handoff_reason"] or ai_session.summary_text) _update_ai_session_context_summary_from_assistant_turn( ai_session, language=decision["language"], customer_name=customer.display_name if customer else thread.display_name, reply_text=decision["reply_text"], decision=decision, now=ai_session.updated_at, ) if decision["case_action"] == "escalate": _interaction_request( "POST", f"/interactions/{interaction.interaction_id}/escalate", payload={"target_queue_id": interaction.queue_id or thread.queue_id or "q_whatsapp"}, ) try: _whatsapp_request( "POST", f"/integrations/whatsapp/threads/{thread.thread_id}/ai/handoff", payload={ "reason": decision["handoff_reason"] or "Нужна передача оператору по результатам AI-анализа.", "agent_profile": ai_session.agent_profile, "trigger_message_id": trigger_message.message_id, "confidence": decision["confidence"], "payload": {"intent": decision["intent"], "kb_refs": decision["kb_refs"]}, }, ) except Exception as exc: # noqa: BLE001 resolved = _conflict_result_from_whatsapp_error( session, job=job, ai_session=ai_session, thread_id=thread.thread_id, exc=exc, reply_message_id=reply_payload.get("message_id"), ) if resolved is not None: return resolved raise elif decision["case_action"] == "close" and _looks_like_resolution_confirmation(trigger_message.text): _interaction_request( "PATCH", f"/interactions/{interaction.interaction_id}/status", payload={"status": "closed"}, ) thread = _whatsapp_thread_or_404(session, thread.thread_id) _update_whatsapp_thread_after_close(session, thread, utc_now_iso()) ai_session.status = "closed" ai_session.closed_at = utc_now_iso() ai_session.updated_at = utc_now_iso() _push_timeline( session, interaction.interaction_id, "ai.case_closed", {"thread_id": thread.thread_id, "reply_message_id": reply_payload.get("message_id")}, ) _mark_job_done(session, job) session.commit() return {"ok": True, "status": "done", "job_id": job.job_id, "reply_message_id": reply_payload.get("message_id")} except HTTPException: raise except Exception as exc: # noqa: BLE001 logger.exception("AI WhatsApp job failed", extra={"job_id": job_id}) try: session.rollback() job = session.execute(select(AIJobRow).where(AIJobRow.job_id == job_id)).scalar_one_or_none() if job: job.status = "failed" job.last_error = str(exc)[:1000] job.updated_at = utc_now_iso() thread = _whatsapp_thread_or_404(session, job.thread_id) thread.ai_state = "error" thread.ai_handoff_reason = str(exc)[:240] thread.updated_at = utc_now_iso() interaction = session.execute( select(Interaction).where(Interaction.interaction_id == thread.interaction_id) ).scalar_one_or_none() if interaction: _push_timeline( session, interaction.interaction_id, "ai.error", {"thread_id": thread.thread_id, "job_id": job.job_id, "error": str(exc)[:500]}, ) session.commit() finally: pass raise HTTPException(status_code=502, detail=f"AI WhatsApp processing failed: {exc}") from exc finally: session.close() @app.get("/health", response_model=HealthResponse) def health() -> HealthResponse: return HealthResponse(status="ok", service="ai-orchestrator-service", version="v1") @app.get("/ai/analytics/overview", response_model=AIAnalyticsOverviewOut) def ai_analytics_overview( from_ts: str = Query(...), to_ts: str = Query(...), queue_id: str | None = None, channel: str | None = Query(default="all"), _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.ANALYST)), ) -> AIAnalyticsOverviewOut: range_from = _parse_analytics_timestamp(from_ts, "from_ts") range_to = _parse_analytics_timestamp(to_ts, "to_ts") if range_to <= range_from: raise HTTPException(status_code=400, detail="to_ts must be greater than from_ts") session = get_session() try: return _load_ai_analytics_overview( session, range_from=range_from, range_to=range_to, queue_id=queue_id, channel=channel, ) finally: session.close() @app.get("/ai/analytics/timeseries", response_model=AIAnalyticsTimeseriesOut) def ai_analytics_timeseries( from_ts: str = Query(...), to_ts: str = Query(...), metric: str = Query(default="containment_rate"), interval: str = Query(default="day"), queue_id: str | None = None, channel: str | None = Query(default="all"), _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.ANALYST)), ) -> AIAnalyticsTimeseriesOut: range_from = _parse_analytics_timestamp(from_ts, "from_ts") range_to = _parse_analytics_timestamp(to_ts, "to_ts") if range_to <= range_from: raise HTTPException(status_code=400, detail="to_ts must be greater than from_ts") normalized_metric = _normalize_ai_analytics_metric(metric) normalized_interval = _normalize_ai_analytics_interval(interval) normalized_channel, channels = _normalize_ai_analytics_channel(channel) if not channels: return _empty_ai_analytics_timeseries( range_from=range_from, range_to=range_to, metric=normalized_metric, interval=normalized_interval, queue_id=queue_id, channel=normalized_channel, ) step = timedelta(hours=1) if normalized_interval == "hour" else timedelta(days=1) points: list[AIAnalyticsTimeseriesPointOut] = [] session = get_session() try: cursor = range_from while cursor < range_to: bucket_from = cursor bucket_to = min(bucket_from + step, range_to) overview = _load_ai_analytics_overview( session, range_from=bucket_from, range_to=bucket_to, queue_id=queue_id, channel=normalized_channel, ) points.append( AIAnalyticsTimeseriesPointOut( ts=bucket_from.isoformat(), value=_timeseries_metric_value(normalized_metric, overview), sessions=overview.totals.sessions_started, assistant_turns=overview.totals.assistant_turns, ) ) cursor = bucket_to finally: session.close() return AIAnalyticsTimeseriesOut( metric=normalized_metric, # type: ignore[arg-type] interval=normalized_interval, # type: ignore[arg-type] filters=_analytics_filters(range_from, range_to, queue_id, normalized_channel), points=points, ) @app.get("/ai/analytics/voice-name-flow/overview", response_model=VoiceNameFlowAnalyticsOverviewOut) def voice_name_flow_overview( from_ts: str = Query(...), to_ts: str = Query(...), queue_id: str | None = None, language: str | None = None, _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.ANALYST)), ) -> VoiceNameFlowAnalyticsOverviewOut: range_from = _parse_analytics_timestamp(from_ts, "from_ts") range_to = _parse_analytics_timestamp(to_ts, "to_ts") if range_to <= range_from: raise HTTPException(status_code=400, detail="to_ts must be greater than from_ts") session = get_session() try: return _load_voice_name_flow_overview( session, range_from=range_from, range_to=range_to, queue_id=queue_id, language=language, ) finally: session.close() @app.get("/ai/analytics/voice-name-flow/timeseries", response_model=VoiceNameFlowAnalyticsTimeseriesOut) def voice_name_flow_timeseries( from_ts: str = Query(...), to_ts: str = Query(...), metric: str = Query(default="scenario_calls"), queue_id: str | None = None, language: str | None = None, _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.ANALYST)), ) -> VoiceNameFlowAnalyticsTimeseriesOut: range_from = _parse_analytics_timestamp(from_ts, "from_ts") range_to = _parse_analytics_timestamp(to_ts, "to_ts") if range_to <= range_from: raise HTTPException(status_code=400, detail="to_ts must be greater than from_ts") normalized_metric = _normalize_voice_name_metric(metric) normalized_language = _normalize_voice_name_language(language) normalized_interval = _voice_name_interval_for_window(range_from, range_to) step = timedelta(hours=1) if normalized_interval == "hour" else timedelta(days=1) points: list[VoiceNameFlowAnalyticsTimeseriesPointOut] = [] session = get_session() try: cursor = range_from while cursor < range_to: bucket_from = cursor bucket_to = min(bucket_from + step, range_to) overview = _load_voice_name_flow_overview( session, range_from=bucket_from, range_to=bucket_to, queue_id=queue_id, language=normalized_language, ) points.append( VoiceNameFlowAnalyticsTimeseriesPointOut( ts=bucket_from.isoformat(), value=_voice_name_metric_value(normalized_metric, overview), scenario_calls=overview.totals.scenario_calls, denominator=_voice_name_metric_denominator(normalized_metric, overview), ) ) cursor = bucket_to finally: session.close() return VoiceNameFlowAnalyticsTimeseriesOut( metric=normalized_metric, # type: ignore[arg-type] interval=normalized_interval, # type: ignore[arg-type] filters=_voice_name_filters(range_from, range_to, queue_id, normalized_language), points=points, ) @app.get("/ai/analytics/drilldown", response_model=AIAnalyticsDrilldownOut) def ai_analytics_drilldown( from_ts: str = Query(...), to_ts: str = Query(...), slice: str = Query(default="all"), queue_id: str | None = None, channel: str | None = Query(default="all"), reason_key: str | None = None, status: str | None = None, q: str | None = None, sort_by: str = Query(default="created_at"), sort_dir: str = Query(default="desc"), limit: int = Query(default=12, ge=1, le=200), offset: int = Query(default=0, ge=0), _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.ANALYST)), ) -> AIAnalyticsDrilldownOut: range_from = _parse_analytics_timestamp(from_ts, "from_ts") range_to = _parse_analytics_timestamp(to_ts, "to_ts") if range_to <= range_from: raise HTTPException(status_code=400, detail="to_ts must be greater than from_ts") normalized_slice = _normalize_ai_analytics_slice(slice) normalized_reason_key = _normalize_ai_analytics_reason_key(reason_key) normalized_sort_by, normalized_sort_dir = _normalize_ai_analytics_sort(sort_by, sort_dir) session = get_session() try: return _load_ai_analytics_drilldown( session, range_from=range_from, range_to=range_to, queue_id=queue_id, channel=channel, slice_name=normalized_slice, reason_key=normalized_reason_key, status=(status or "").strip() or None, query_text=q, sort_by=normalized_sort_by, sort_dir=normalized_sort_dir, limit=limit, offset=offset, ) finally: session.close() @app.get("/ai/analytics/sessions/{session_id}", response_model=AIAnalyticsSessionDetailOut) def ai_analytics_session_detail( session_id: str, _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR, Role.ANALYST)), ) -> AIAnalyticsSessionDetailOut: session = get_session() try: return _load_ai_analytics_session_detail(session, session_id) finally: session.close() @app.get("/ai/operator/config", response_model=AIOperatorConfigOut) def get_ai_operator_config( _: dict = Depends(require_roles(Role.ADMIN)), ) -> AIOperatorConfigOut: session = get_session() try: return load_ai_operator_config(session) finally: session.close() @app.put("/ai/operator/config", response_model=AIOperatorConfigOut) def put_ai_operator_config( payload: AIOperatorConfig, _: dict = Depends(require_roles(Role.ADMIN)), ) -> AIOperatorConfigOut: session = get_session() try: return save_ai_operator_config(session, payload) finally: session.close() @app.post("/ai/voice/sessions/{session_id}/start") def start_voice_ai_session( session_id: str, payload: VoiceAIStartIn, _: dict = Depends(require_roles(Role.ADMIN)), ) -> VoiceAIStartOut: return voice_flows.start_voice_session(session_id, payload) @app.get("/ai/voice/config/name-collection", response_model=VoiceNameCollectionConfigOut) def get_voice_name_collection_config( _: dict = Depends(require_roles(Role.ADMIN)), ) -> VoiceNameCollectionConfigOut: session = get_session() try: return load_voice_name_collection_config(session) finally: session.close() @app.put("/ai/voice/config/name-collection", response_model=VoiceNameCollectionConfigOut) def put_voice_name_collection_config( payload: VoiceNameCollectionConfig, _: dict = Depends(require_roles(Role.ADMIN)), ) -> VoiceNameCollectionConfigOut: session = get_session() try: return save_voice_name_collection_config(session, payload) finally: session.close() @app.get("/ai/voice/config/tts", response_model=VoiceTTSConfigOut) def get_voice_tts_config( _: dict = Depends(require_roles(Role.ADMIN)), ) -> VoiceTTSConfigOut: session = get_session() try: return load_voice_tts_config(session) finally: session.close() @app.put("/ai/voice/config/tts", response_model=VoiceTTSConfigOut) def put_voice_tts_config( payload: VoiceTTSConfig, _: dict = Depends(require_roles(Role.ADMIN)), ) -> VoiceTTSConfigOut: session = get_session() try: return save_voice_tts_config(session, payload) finally: session.close() @app.post("/ai/voice/sessions/{session_id}/turns") def turn_voice_ai_session( session_id: str, payload: VoiceAITurnIn, _: dict = Depends(require_roles(Role.ADMIN)), ) -> dict[str, Any]: return voice_flows.turn_voice_session(session_id, payload).model_dump() @app.post("/ai/voice/sessions/{session_id}/close") def close_voice_ai_session( session_id: str, _: dict = Depends(require_roles(Role.ADMIN)), ) -> dict[str, Any]: return voice_flows.close_voice_session(session_id) @app.post("/ai/telegram/threads/{thread_id}/enqueue") def enqueue_telegram_thread( thread_id: str, payload: AITelegramEnqueueIn, _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR)), ) -> dict[str, Any]: if not _ai_telegram_enabled(): return {"ok": True, "status": "disabled", "thread_id": thread_id} session = get_session() try: thread = _thread_or_404(session, thread_id) interaction = session.execute( select(Interaction).where(Interaction.interaction_id == thread.interaction_id) ).scalar_one() customer_id = _resolve_or_create_customer_id(session, thread, interaction) trigger_message = _select_trigger_message( session, thread_id=thread.thread_id, trigger_message_id=payload.trigger_message_id, ) if not trigger_message: return {"ok": False, "status": "skipped", "reason": "no customer message", "thread_id": thread_id} language = _infer_language(trigger_message.text) ai_session, _ = _ensure_ai_session( session, thread=thread, interaction=interaction, customer_id=customer_id, language=language, ) job, deduplicated = _ensure_job( session, thread=thread, ai_session=ai_session, trigger_message_id=trigger_message.message_id, ) session.commit() if deduplicated and job: return { "ok": True, "status": job.status, "thread_id": thread.thread_id, "session_id": ai_session.session_id, "job_id": job.job_id, "deduplicated": True, } if not job: return { "ok": True, "status": "skipped", "thread_id": thread.thread_id, "session_id": ai_session.session_id, "deduplicated": True, } job_id = job.job_id session_id = ai_session.session_id finally: session.close() result = _process_job(job_id) result["thread_id"] = thread_id result["session_id"] = session_id result["deduplicated"] = False return result @app.post("/ai/telegram/threads/{thread_id}/pause") def pause_telegram_thread_ai( thread_id: str, payload: AITelegramPauseIn, _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR)), ) -> dict[str, Any]: session = get_session() try: thread = _thread_or_404(session, thread_id) interaction = session.execute( select(Interaction).where(Interaction.interaction_id == thread.interaction_id) ).scalar_one_or_none() now = utc_now_iso() thread.ai_state = "human_owned" thread.ai_handoff_reason = payload.reason thread.updated_at = now if thread.ai_session_id: ai_session = session.execute( select(AISessionRow).where(AISessionRow.session_id == thread.ai_session_id) ).scalar_one_or_none() if ai_session: ai_session.status = "human_owned" ai_session.handoff_reason = payload.reason ai_session.updated_at = now if interaction: _push_timeline( session, interaction.interaction_id, "ai.handoff_requested", {"thread_id": thread.thread_id, "reason": payload.reason, "actor_user": payload.actor_user}, ) session.commit() return {"ok": True, "status": "human_owned", "thread_id": thread.thread_id} finally: session.close() @app.post("/ai/whatsapp/threads/{thread_id}/enqueue") def enqueue_whatsapp_thread( thread_id: str, payload: AIWhatsAppEnqueueIn, _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR)), ) -> dict[str, Any]: if not _ai_whatsapp_enabled(): return {"ok": True, "status": "disabled", "thread_id": thread_id} session = get_session() try: thread = _whatsapp_thread_or_404(session, thread_id) interaction = session.execute( select(Interaction).where(Interaction.interaction_id == thread.interaction_id) ).scalar_one() customer_id = _resolve_or_create_whatsapp_customer_id(session, thread, interaction) trigger_message = _select_whatsapp_trigger_message( session, thread_id=thread.thread_id, trigger_message_id=payload.trigger_message_id, ) if not trigger_message: return {"ok": False, "status": "skipped", "reason": "no customer message", "thread_id": thread_id} language = _infer_language(trigger_message.text) ai_session, _ = _ensure_whatsapp_ai_session( session, thread=thread, interaction=interaction, customer_id=customer_id, language=language, ) job, deduplicated = _ensure_job( session, thread=thread, ai_session=ai_session, trigger_message_id=trigger_message.message_id, ) session.commit() if deduplicated and job: return { "ok": True, "status": job.status, "thread_id": thread.thread_id, "session_id": ai_session.session_id, "job_id": job.job_id, "deduplicated": True, } if not job: return { "ok": True, "status": "skipped", "thread_id": thread.thread_id, "session_id": ai_session.session_id, "deduplicated": True, } job_id = job.job_id session_id = ai_session.session_id finally: session.close() result = _process_whatsapp_job(job_id) result["thread_id"] = thread_id result["session_id"] = session_id result["deduplicated"] = False return result @app.post("/ai/whatsapp/threads/{thread_id}/pause") def pause_whatsapp_thread_ai( thread_id: str, payload: AIWhatsAppPauseIn, _: dict = Depends(require_roles(Role.ADMIN, Role.SUPERVISOR)), ) -> dict[str, Any]: session = get_session() try: thread = _whatsapp_thread_or_404(session, thread_id) interaction = session.execute( select(Interaction).where(Interaction.interaction_id == thread.interaction_id) ).scalar_one_or_none() now = utc_now_iso() thread.ai_state = "human_owned" thread.ai_handoff_reason = payload.reason thread.updated_at = now if thread.ai_session_id: ai_session = session.execute( select(AISessionRow).where(AISessionRow.session_id == thread.ai_session_id) ).scalar_one_or_none() if ai_session: ai_session.status = "human_owned" ai_session.handoff_reason = payload.reason ai_session.updated_at = now if interaction: _push_timeline( session, interaction.interaction_id, "ai.handoff_requested", {"thread_id": thread.thread_id, "reason": payload.reason, "actor_user": payload.actor_user}, ) session.commit() return {"ok": True, "status": "human_owned", "thread_id": thread.thread_id} finally: session.close()