210 lines
6.2 KiB
Python
210 lines
6.2 KiB
Python
from __future__ import annotations
|
|
|
|
import json
|
|
import sys
|
|
from datetime import datetime
|
|
from typing import Any
|
|
|
|
from sqlalchemy import select
|
|
|
|
from services.shared.models import AsteriskEventOut, VoiceCallActionOut, VoiceLiveCallOut
|
|
from services.shared.sql_models import (
|
|
AsteriskCallActionLogRow,
|
|
AsteriskCallLinkRow,
|
|
AsteriskEventLogRow,
|
|
VoiceEventRow,
|
|
)
|
|
|
|
|
|
def _bridge_app():
|
|
return sys.modules["services.asterisk_bridge_service.app"]
|
|
|
|
|
|
def to_out(row: AsteriskEventLogRow) -> AsteriskEventOut:
|
|
return AsteriskEventOut(
|
|
bridge_event_id=row.bridge_event_id,
|
|
ami_event_name=row.ami_event_name,
|
|
call_id=row.call_id,
|
|
linked_id=row.linked_id,
|
|
interaction_id=row.interaction_id,
|
|
recording_id=row.recording_id,
|
|
forward_status=row.forward_status,
|
|
payload=json.loads(row.payload_json or "{}"),
|
|
last_error=row.last_error,
|
|
created_at=row.created_at,
|
|
updated_at=row.updated_at,
|
|
)
|
|
|
|
|
|
def action_to_out(row: AsteriskCallActionLogRow) -> VoiceCallActionOut:
|
|
return VoiceCallActionOut(
|
|
action_id=row.action_id,
|
|
call_id=row.call_id,
|
|
interaction_id=row.interaction_id,
|
|
action_type=row.action_type,
|
|
actor_user=row.actor_user,
|
|
actor_role=row.actor_role,
|
|
request=json.loads(row.request_json or "{}"),
|
|
result_status=row.result_status, # type: ignore[arg-type]
|
|
ami_action_id=row.ami_action_id,
|
|
error=row.error,
|
|
created_at=row.created_at,
|
|
)
|
|
|
|
|
|
def safe_json_loads(raw: str | None) -> dict[str, Any]:
|
|
try:
|
|
payload = json.loads(raw or "{}")
|
|
except Exception: # noqa: BLE001
|
|
return {}
|
|
return payload if isinstance(payload, dict) else {}
|
|
|
|
|
|
def latest_voice_event_row(
|
|
session,
|
|
*,
|
|
call_id: str,
|
|
event_type: str,
|
|
include_reconciled: bool = True,
|
|
) -> VoiceEventRow | None:
|
|
rows = session.execute(
|
|
select(VoiceEventRow)
|
|
.where(VoiceEventRow.call_id == call_id)
|
|
.where(VoiceEventRow.event_type == event_type)
|
|
.order_by(VoiceEventRow.id.desc())
|
|
.limit(20)
|
|
).scalars().all()
|
|
if include_reconciled:
|
|
return rows[0] if rows else None
|
|
for row in rows:
|
|
payload = safe_json_loads(row.payload_json)
|
|
if payload.get("reconciled"):
|
|
continue
|
|
return row
|
|
return None
|
|
|
|
|
|
def latest_successful_action_row(
|
|
session,
|
|
*,
|
|
call_id: str,
|
|
action_types: tuple[str, ...],
|
|
) -> AsteriskCallActionLogRow | None:
|
|
rows = session.execute(
|
|
select(AsteriskCallActionLogRow)
|
|
.where(AsteriskCallActionLogRow.call_id == call_id)
|
|
.where(AsteriskCallActionLogRow.result_status == "ok")
|
|
.order_by(AsteriskCallActionLogRow.id.desc())
|
|
.limit(20)
|
|
).scalars().all()
|
|
for row in rows:
|
|
if row.action_type in action_types:
|
|
return row
|
|
return None
|
|
|
|
|
|
def latest_iso(*values: str | None) -> str | None:
|
|
bridge = _bridge_app()
|
|
parsed: list[tuple[datetime, str]] = []
|
|
for value in values:
|
|
iso = str(value or "").strip()
|
|
dt = bridge._parse_iso(iso)
|
|
if dt is None:
|
|
continue
|
|
parsed.append((dt, iso))
|
|
if not parsed:
|
|
return None
|
|
parsed.sort(key=lambda item: item[0], reverse=True)
|
|
return parsed[0][1]
|
|
|
|
|
|
def hangup_cause_for_call(session, *, call_id: str) -> str | None:
|
|
row = latest_voice_event_row(
|
|
session,
|
|
call_id=call_id,
|
|
event_type="call.ended",
|
|
include_reconciled=False,
|
|
) or latest_voice_event_row(
|
|
session,
|
|
call_id=call_id,
|
|
event_type="call.ended",
|
|
include_reconciled=True,
|
|
)
|
|
if row is None:
|
|
return None
|
|
payload = safe_json_loads(row.payload_json)
|
|
value = str(payload.get("hangup_cause") or "").strip()
|
|
return value or None
|
|
|
|
|
|
def terminal_action_for_call(
|
|
session,
|
|
*,
|
|
call_id: str,
|
|
row: AsteriskCallLinkRow,
|
|
) -> tuple[str | None, str | None, str | None]:
|
|
bridge = _bridge_app()
|
|
action = latest_successful_action_row(
|
|
session,
|
|
call_id=call_id,
|
|
action_types=("blind-transfer", "hangup"),
|
|
)
|
|
if action is not None:
|
|
request_payload = safe_json_loads(action.request_json)
|
|
target = None
|
|
if action.action_type == "blind-transfer":
|
|
target = str(
|
|
request_payload.get("target_value")
|
|
or request_payload.get("resolved_extension")
|
|
or ""
|
|
).strip() or None
|
|
return action.action_type, action.created_at, target
|
|
if bridge._call_is_ended(row):
|
|
return "ended", row.ended_at or row.updated_at, None
|
|
return None, None, None
|
|
|
|
|
|
def live_call_to_out(session, row: AsteriskCallLinkRow) -> VoiceLiveCallOut:
|
|
bridge = _bridge_app()
|
|
has_recording = bridge._has_recording(session, call_id=row.call_id)
|
|
terminal_action, terminal_at, terminal_target = terminal_action_for_call(
|
|
session,
|
|
call_id=row.call_id,
|
|
row=row,
|
|
)
|
|
last_transition_at = latest_iso(
|
|
terminal_at,
|
|
row.connected_at,
|
|
row.claimed_at,
|
|
row.started_at,
|
|
row.updated_at,
|
|
)
|
|
return VoiceLiveCallOut(
|
|
call_id=row.call_id,
|
|
interaction_id=row.interaction_id,
|
|
queue_id=row.queue_id,
|
|
queue_code=row.queue_code,
|
|
caller_number=row.caller_number,
|
|
caller_name=row.caller_name,
|
|
status=row.status,
|
|
telephony_status=row.telephony_status, # type: ignore[arg-type]
|
|
claimed_by_user=row.claimed_by_user,
|
|
claimed_at=row.claimed_at,
|
|
operator_extension=row.operator_extension,
|
|
channel_name=row.channel_name,
|
|
started_at=row.started_at,
|
|
connected_at=row.connected_at,
|
|
ended_at=row.ended_at,
|
|
updated_at=row.updated_at,
|
|
last_transition_at=last_transition_at,
|
|
hangup_cause=hangup_cause_for_call(session, call_id=row.call_id),
|
|
terminal_action=terminal_action,
|
|
terminal_target=terminal_target,
|
|
voice_session_id=row.voice_session_id,
|
|
ai_session_id=row.ai_session_id,
|
|
ai_state=row.ai_state, # type: ignore[arg-type]
|
|
ai_handoff_reason=row.ai_handoff_reason,
|
|
ai_last_model_at=row.ai_last_model_at,
|
|
has_recording=has_recording,
|
|
)
|