Merge pull request 'feat: no-answer retry, agent status machine, escalation events (ТЗ Phase 2)' (#8) from feature/escalation-no-answer-retry-and-agent-state into main
deploy / deploy (push) Successful in 32s

This commit was merged in pull request #8.
This commit is contained in:
2026-08-30 09:19:59 +00:00
19 changed files with 792 additions and 36 deletions
+13
View File
@@ -0,0 +1,13 @@
{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"title": "AgentConnected",
"type": "object",
"required": ["event", "call_id", "escalation_id", "agent_id", "created_at"],
"properties": {
"event": { "const": "AgentConnected" },
"call_id": { "type": "string" },
"escalation_id": { "type": "string" },
"agent_id": { "type": "string" },
"created_at": { "type": "string", "format": "date-time" }
}
}
+15
View File
@@ -0,0 +1,15 @@
{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"title": "AgentNoAnswer",
"type": "object",
"required": ["event", "call_id", "escalation_id", "agent_id", "dial_outcome", "created_at"],
"properties": {
"event": { "const": "AgentNoAnswer" },
"call_id": { "type": "string" },
"escalation_id": { "type": "string" },
"agent_id": { "type": ["string", "null"] },
"dial_outcome": { "type": "string" },
"attempt_count": { "type": "integer" },
"created_at": { "type": "string", "format": "date-time" }
}
}
+18
View File
@@ -0,0 +1,18 @@
{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"title": "AgentReserved",
"type": "object",
"required": ["event", "call_id", "escalation_id", "agent_id", "created_at"],
"properties": {
"event": { "const": "AgentReserved" },
"call_id": { "type": "string" },
"escalation_id": { "type": "string" },
"agent_id": { "type": "string" },
"tenant_id": { "type": ["string", "null"] },
"from_level": { "type": "string" },
"to_level": { "type": "string" },
"reason_code": { "type": "string" },
"attempt_count": { "type": "integer" },
"created_at": { "type": "string", "format": "date-time" }
}
}
+13
View File
@@ -0,0 +1,13 @@
{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"title": "AgentRinging",
"type": "object",
"required": ["event", "call_id", "escalation_id", "agent_id", "created_at"],
"properties": {
"event": { "const": "AgentRinging" },
"call_id": { "type": "string" },
"escalation_id": { "type": "string" },
"agent_id": { "type": "string" },
"created_at": { "type": "string", "format": "date-time" }
}
}
+13
View File
@@ -0,0 +1,13 @@
{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"title": "TransferCompleted",
"type": "object",
"required": ["event", "call_id", "escalation_id", "agent_id", "created_at"],
"properties": {
"event": { "const": "TransferCompleted" },
"call_id": { "type": "string" },
"escalation_id": { "type": "string" },
"agent_id": { "type": "string" },
"created_at": { "type": "string", "format": "date-time" }
}
}
+15
View File
@@ -0,0 +1,15 @@
{
"$schema": "https://json-schema.org/draft/2020-12/schema",
"title": "TransferFailed",
"type": "object",
"required": ["event", "call_id", "created_at"],
"properties": {
"event": { "const": "TransferFailed" },
"call_id": { "type": "string" },
"escalation_id": { "type": ["string", "null"] },
"agent_id": { "type": ["string", "null"] },
"reason": { "type": ["string", "null"] },
"error": { "type": ["string", "null"] },
"created_at": { "type": "string", "format": "date-time" }
}
}
@@ -0,0 +1,4 @@
ALTER TABLE escalations ADD COLUMN IF NOT EXISTS attempt_count INTEGER NOT NULL DEFAULT 0;
ALTER TABLE escalations ADD COLUMN IF NOT EXISTS real_agent_id TEXT;
ALTER TABLE escalations ADD COLUMN IF NOT EXISTS attempted_agent_ids_json TEXT NOT NULL DEFAULT '[]';
CREATE INDEX IF NOT EXISTS ix_escalations_real_agent_id ON escalations(real_agent_id);
@@ -0,0 +1,4 @@
ALTER TABLE escalations ADD COLUMN attempt_count INTEGER NOT NULL DEFAULT 0;
ALTER TABLE escalations ADD COLUMN real_agent_id TEXT;
ALTER TABLE escalations ADD COLUMN attempted_agent_ids_json TEXT NOT NULL DEFAULT '[]';
CREATE INDEX IF NOT EXISTS ix_escalations_real_agent_id ON escalations(real_agent_id);
+11 -6
View File
@@ -412,13 +412,18 @@ def ami_loop(stop_event=None) -> None:
raise RuntimeError("AMI connection closed") raise RuntimeError("AMI connection closed")
if not frame: if not frame:
continue continue
if frame.get("Event") != "UserEvent": event_name = frame.get("Event")
if event_name == "UserEvent":
user_event = str(frame.get("UserEvent") or "").strip()
if not user_event.startswith(bridge._ami_prefix()):
continue
bridge._STATE.set_last_event()
bridge._record_ami_payload(frame)
elif event_name in {"DialEnd", "Hangup"}:
bridge._STATE.set_last_event()
bridge._record_ami_payload(frame, event_name=event_name)
else:
continue continue
user_event = str(frame.get("UserEvent") or "").strip()
if not user_event.startswith(bridge._ami_prefix()):
continue
bridge._STATE.set_last_event()
bridge._record_ami_payload(frame)
except Exception as exc: except Exception as exc:
bridge._STATE.set_connected(False) bridge._STATE.set_connected(False)
bridge._STATE.set_error(str(exc)) bridge._STATE.set_error(str(exc))
+6
View File
@@ -262,6 +262,7 @@ def _start_background_threads() -> None:
extra_loops: list[tuple[str, Callable[..., Any]]] = [] extra_loops: list[tuple[str, Callable[..., Any]]] = []
if _ivr_fastagi_enabled(): if _ivr_fastagi_enabled():
extra_loops.append(("asterisk-ivr-fastagi-loop", _ivr_fastagi_loop)) extra_loops.append(("asterisk-ivr-fastagi-loop", _ivr_fastagi_loop))
extra_loops.append(("agent-acw-sweep-loop", _acw_sweep_loop))
bridge_runtime.start_background_threads( bridge_runtime.start_background_threads(
ami_loop=_ami_loop, ami_loop=_ami_loop,
failed_retry_loop=_failed_retry_loop, failed_retry_loop=_failed_retry_loop,
@@ -305,6 +306,9 @@ _update_voice_ai_call_state = bridge_voice_ai.update_call_ai_state
_voice_ai_summary_for_call = bridge_voice_ai.voice_ai_summary_for_call _voice_ai_summary_for_call = bridge_voice_ai.voice_ai_summary_for_call
_create_escalation = bridge_voice_ai.create_escalation _create_escalation = bridge_voice_ai.create_escalation
_release_routing_agent = bridge_voice_ai.release_routing_agent _release_routing_agent = bridge_voice_ai.release_routing_agent
_retry_escalation_no_answer = bridge_voice_ai.retry_escalation_no_answer
_routing_release_by_agent_id = bridge_voice_ai.routing_release_by_agent_id
_set_routing_agent_status = bridge_voice_ai.set_routing_agent_status
_first_non_empty = bridge_ami.first_non_empty _first_non_empty = bridge_ami.first_non_empty
_extract_call_id = bridge_ami.extract_call_id _extract_call_id = bridge_ami.extract_call_id
@@ -373,6 +377,8 @@ _process_audio_bridge_ended = bridge_processing.process_audio_bridge_ended
_process_call_ended = bridge_processing.process_call_ended _process_call_ended = bridge_processing.process_call_ended
_process_operator_connected = bridge_processing.process_operator_connected _process_operator_connected = bridge_processing.process_operator_connected
_process_recording_ready = bridge_processing.process_recording_ready _process_recording_ready = bridge_processing.process_recording_ready
_process_agent_dial_outcome = bridge_processing.process_agent_dial_outcome
_acw_sweep_loop = bridge_processing.acw_sweep_loop
_process_bridge_row = bridge_processing.process_bridge_row _process_bridge_row = bridge_processing.process_bridge_row
_record_ami_payload = bridge_processing.record_ami_payload _record_ami_payload = bridge_processing.record_ami_payload
_retry_failed_events_once = bridge_processing.retry_failed_events_once _retry_failed_events_once = bridge_processing.retry_failed_events_once
@@ -4,6 +4,7 @@ from datetime import datetime, timezone
import hashlib import hashlib
import json import json
import logging import logging
import os
from pathlib import Path from pathlib import Path
import threading import threading
from typing import Any from typing import Any
@@ -884,18 +885,34 @@ def process_call_ended(
open_escalation = session.execute( open_escalation = session.execute(
select(EscalationRow).where( select(EscalationRow).where(
EscalationRow.call_id == row.call_id, EscalationRow.call_id == row.call_id,
EscalationRow.status.in_(["requested", "ringing"]), EscalationRow.status.in_(["requested", "ringing", "connected"]),
) )
).scalar_one_or_none() ).scalar_one_or_none()
was_talking = open_escalation is not None and open_escalation.status == "connected"
if open_escalation is not None: if open_escalation is not None:
open_escalation.status = "completed" if answered else "failed" open_escalation.status = "completed" if answered else "failed"
open_escalation.completed_at = now open_escalation.completed_at = now
if answered and not open_escalation.connected_at: if answered and not open_escalation.connected_at:
open_escalation.connected_at = now open_escalation.connected_at = now
try: if was_talking and open_escalation and open_escalation.real_agent_id:
bridge._release_routing_agent(row.call_id) try:
except Exception: bridge._set_routing_agent_status(open_escalation.real_agent_id, "AFTER_CALL_WORK")
pass except Exception:
logger.warning("bridge.agent_acw_status_failed call_id=%s agent_id=%s", row.call_id, open_escalation.real_agent_id)
else:
try:
bridge._release_routing_agent(row.call_id)
except Exception:
pass
if open_escalation is not None:
try:
bridge._append_interaction_timeline(
interaction_id=link.interaction_id,
action="escalation.completed" if answered else "escalation.failed_client_disconnected",
metadata={"call_id": row.call_id, "escalation_id": open_escalation.escalation_id, "agent_id": open_escalation.real_agent_id},
)
except Exception:
pass
try: try:
bridge._notify_voice_ai_telephony_event( bridge._notify_voice_ai_telephony_event(
voice_session_id=link.voice_session_id, voice_session_id=link.voice_session_id,
@@ -969,6 +986,48 @@ def process_operator_connected(
link.telephony_status = "connected" link.telephony_status = "connected"
link.connected_at = now link.connected_at = now
link.updated_at = now link.updated_at = now
open_escalation = session.execute(
select(EscalationRow).where(
EscalationRow.call_id == row.call_id,
EscalationRow.status == "ringing",
)
).scalar_one_or_none()
if open_escalation is not None:
open_escalation.status = "connected"
open_escalation.connected_at = now
if open_escalation.real_agent_id:
try:
bridge._set_routing_agent_status(open_escalation.real_agent_id, "TALKING")
except Exception:
logger.warning("bridge.agent_talking_status_failed call_id=%s agent_id=%s", row.call_id, open_escalation.real_agent_id)
try:
bridge._append_interaction_timeline(
interaction_id=link.interaction_id,
action="escalation.agent_connected",
metadata={"call_id": row.call_id, "agent_id": open_escalation.real_agent_id, "escalation_id": open_escalation.escalation_id},
)
except Exception:
pass
try:
bridge._emit_voice_event(
event_type="AgentConnected",
call_id=row.call_id,
interaction_id=link.interaction_id,
payload={"escalation_id": open_escalation.escalation_id, "agent_id": open_escalation.real_agent_id},
)
except Exception:
pass
try:
bridge._emit_voice_event(
event_type="TransferCompleted",
call_id=row.call_id,
interaction_id=link.interaction_id,
payload={"escalation_id": open_escalation.escalation_id, "agent_id": open_escalation.real_agent_id},
)
except Exception:
pass
_upsert_voice_reporting_fact( _upsert_voice_reporting_fact(
session, session,
link, link,
@@ -1173,6 +1232,55 @@ def process_recording_ready(
local_path.unlink(missing_ok=True) local_path.unlink(missing_ok=True)
_NO_ANSWER_DIAL_STATUSES = {"NOANSWER", "BUSY", "CANCEL", "CHANUNAVAIL", "CONGESTION"}
_NO_ANSWER_HANGUP_CAUSES = {"17", "18", "19", "21", "34", "38"}
def process_agent_dial_outcome(session, row: AsteriskEventLogRow, payload: dict[str, Any]) -> None:
bridge = _bridge_app()
call_id = row.call_id
if not call_id:
bridge._mark_log_forwarded(session, row, interaction_id=None)
return
dial_status = str(payload.get("DialStatus") or "").strip().upper()
hangup_cause = str(payload.get("Cause") or "").strip()
is_no_answer_outcome = (
(row.ami_event_name == "DialEnd" and dial_status in _NO_ANSWER_DIAL_STATUSES)
or (row.ami_event_name == "Hangup" and hangup_cause in _NO_ANSWER_HANGUP_CAUSES)
)
if not is_no_answer_outcome:
bridge._mark_log_forwarded(session, row, interaction_id=None)
return
escalation = session.execute(
select(EscalationRow).where(
EscalationRow.call_id == call_id,
EscalationRow.status == "ringing",
)
).scalar_one_or_none()
if escalation is None:
bridge._mark_log_forwarded(session, row, interaction_id=None)
return
dialed_channel = payload.get("DestChannel") if row.ami_event_name == "DialEnd" else payload.get("Channel")
dialed_extension = bridge._extract_extension_from_channel(dialed_channel)
if dialed_extension and escalation.assigned_agent_id and dialed_extension != escalation.assigned_agent_id:
bridge._mark_log_forwarded(session, row, interaction_id=None)
return
try:
bridge._retry_escalation_no_answer(
session,
call_id=call_id,
dial_outcome=dial_status or f"hangup_cause_{hangup_cause}",
)
except Exception:
logger.exception("bridge.agent_dial_outcome_retry_failed call_id=%s", call_id)
bridge._mark_log_forwarded(session, row, interaction_id=None)
def process_bridge_row(session, row: AsteriskEventLogRow) -> AsteriskEventLogRow: def process_bridge_row(session, row: AsteriskEventLogRow) -> AsteriskEventLogRow:
bridge = _bridge_app() bridge = _bridge_app()
payload = json.loads(row.payload_json or "{}") payload = json.loads(row.payload_json or "{}")
@@ -1209,6 +1317,8 @@ def process_bridge_row(session, row: AsteriskEventLogRow) -> AsteriskEventLogRow
bridge._process_call_ended(session, row, payload) bridge._process_call_ended(session, row, payload)
elif row.ami_event_name == f"{bridge._ami_prefix()}RecordingReady": elif row.ami_event_name == f"{bridge._ami_prefix()}RecordingReady":
bridge._process_recording_ready(session, row, payload) bridge._process_recording_ready(session, row, payload)
elif row.ami_event_name in {"DialEnd", "Hangup"}:
bridge._process_agent_dial_outcome(session, row, payload)
else: else:
row.forward_status = "received" row.forward_status = "received"
row.updated_at = utc_now_iso() row.updated_at = utc_now_iso()
@@ -1216,9 +1326,9 @@ def process_bridge_row(session, row: AsteriskEventLogRow) -> AsteriskEventLogRow
return row return row
def record_ami_payload(payload: dict[str, Any]) -> None: def record_ami_payload(payload: dict[str, Any], *, event_name: str | None = None) -> None:
bridge = _bridge_app() bridge = _bridge_app()
event_name = str(payload.get("UserEvent") or "").strip() event_name = str(event_name or payload.get("UserEvent") or "").strip()
call_id = bridge._extract_call_id(payload) call_id = bridge._extract_call_id(payload)
linked_id = bridge._extract_linked_id(payload, call_id) linked_id = bridge._extract_linked_id(payload, call_id)
if not event_name or not call_id: if not event_name or not call_id:
@@ -1262,6 +1372,42 @@ def retry_failed_events_once() -> None:
bridge._process_claimed_bridge_event(bridge_event_id) bridge._process_claimed_bridge_event(bridge_event_id)
def acw_duration_seconds() -> int:
raw = str(os.environ.get("ACW_DURATION_SECONDS", "30")).strip()
try:
return max(int(raw), 1)
except ValueError:
return 30
def acw_sweep_interval_seconds() -> float:
return max(min(float(acw_duration_seconds()) / 2, 15.0), 5.0)
def acw_sweep_once() -> None:
from services.routing_service import engine as routing_engine
session = get_session()
try:
swept = routing_engine.sweep_after_call_work(session, older_than_seconds=acw_duration_seconds())
for agent in swept:
logger.warning("bridge.agent_acw_swept agent_id=%s", agent.agent_id)
finally:
session.close()
def acw_sweep_loop(stop_event: threading.Event | None = None) -> None:
bridge = _bridge_app()
active_stop_event = stop_event or bridge._background_stop_event()
while not active_stop_event.is_set():
try:
acw_sweep_once()
except Exception:
logger.exception("bridge.acw_sweep_failed")
if active_stop_event.wait(acw_sweep_interval_seconds()):
break
def failed_retry_loop(stop_event: threading.Event | None = None) -> None: def failed_retry_loop(stop_event: threading.Event | None = None) -> None:
bridge = _bridge_app() bridge = _bridge_app()
active_stop_event = stop_event or bridge._background_stop_event() active_stop_event = stop_event or bridge._background_stop_event()
+216 -18
View File
@@ -63,6 +63,57 @@ def append_interaction_timeline(
) )
def _append_escalation_timeline(
escalation: EscalationRow,
link: AsteriskCallLinkRow,
*,
action: str,
extra: dict[str, Any] | None = None,
) -> None:
try:
append_interaction_timeline(
interaction_id=link.interaction_id,
action=action,
metadata={
"call_id": escalation.call_id,
"escalation_id": escalation.escalation_id,
"from_level": escalation.from_level,
"to_level": escalation.to_level,
"attempt_count": escalation.attempt_count,
**(extra or {}),
},
)
except Exception:
pass
def _emit_escalation_event(
*,
event_type: str,
escalation: EscalationRow,
link: AsteriskCallLinkRow,
extra: dict[str, Any] | None = None,
) -> None:
bridge = _bridge_app()
try:
bridge._emit_voice_event(
event_type=event_type,
call_id=escalation.call_id,
interaction_id=link.interaction_id,
payload={
"escalation_id": escalation.escalation_id,
"tenant_id": escalation.tenant_id,
"from_level": escalation.from_level,
"to_level": escalation.to_level,
"reason_code": escalation.reason_code,
"attempt_count": escalation.attempt_count,
**(extra or {}),
},
)
except Exception:
LOGGER.warning("bridge.escalation_event_failed event_type=%s call_id=%s", event_type, escalation.call_id)
def start_voice_ai_session( def start_voice_ai_session(
*, *,
call_id: str, call_id: str,
@@ -188,7 +239,8 @@ def _reserve_routing_agent(
level: str, level: str,
tenant_id: str | None, tenant_id: str | None,
required_skills: list[str] | None = None, required_skills: list[str] | None = None,
) -> str | None: exclude_agent_ids: list[str] | None = None,
) -> dict[str, Any] | None:
bridge = _bridge_app() bridge = _bridge_app()
try: try:
response = bridge._post_json( response = bridge._post_json(
@@ -198,7 +250,7 @@ def _reserve_routing_agent(
"level": level, "level": level,
"tenant_id": tenant_id, "tenant_id": tenant_id,
"required_skills": required_skills or [], "required_skills": required_skills or [],
"exclude_agent_ids": [], "exclude_agent_ids": exclude_agent_ids or [],
}, },
timeout_seconds=bridge._callcontrol_side_effect_timeout_seconds(), timeout_seconds=bridge._callcontrol_side_effect_timeout_seconds(),
max_attempts=1, max_attempts=1,
@@ -212,8 +264,43 @@ def _reserve_routing_agent(
tenant_id, tenant_id,
) )
return None return None
agent_id = str((response or {}).get("agent_id") or "").strip()
extension = str((response or {}).get("extension") or "").strip() extension = str((response or {}).get("extension") or "").strip()
return extension or None if not agent_id or not extension:
return None
return {
"agent_id": agent_id,
"extension": extension,
"endpoint": (response or {}).get("endpoint"),
"display_name": (response or {}).get("display_name"),
}
def set_routing_agent_status(agent_id: str, status: str) -> None:
bridge = _bridge_app()
try:
bridge._patch_json(
f"{bridge._routing_service_url()}/internal/routing/agents/{agent_id}/status",
{"status": status},
timeout_seconds=bridge._callcontrol_side_effect_timeout_seconds(),
max_attempts=1,
retry_backoff_seconds=0.0,
)
except Exception:
LOGGER.warning("bridge.routing_set_status_failed agent_id=%s status=%s", agent_id, status)
def _redirect_channel_to_agent(*, channel: str, extension: str) -> None:
bridge = _bridge_app()
bridge._ami_action(
"Redirect",
{
"Channel": channel,
"Context": bridge._transfer_context(),
"Exten": extension,
"Priority": 1,
},
)
def _resolve_handoff_extension( def _resolve_handoff_extension(
@@ -233,14 +320,14 @@ def _resolve_handoff_extension(
level = target_level or bridge._routing_level_for_queue_code(queue_code) level = target_level or bridge._routing_level_for_queue_code(queue_code)
if level and call_id: if level and call_id:
resolved_tenant_id = tenant_id if tenant_id is not None else bridge._routing_tenant_for_queue_code(queue_code) resolved_tenant_id = tenant_id if tenant_id is not None else bridge._routing_tenant_for_queue_code(queue_code)
reserved_extension = _reserve_routing_agent( reserved_agent = _reserve_routing_agent(
call_id=call_id, call_id=call_id,
level=level, level=level,
tenant_id=resolved_tenant_id, tenant_id=resolved_tenant_id,
required_skills=required_skills, required_skills=required_skills,
) )
if reserved_extension: if reserved_agent:
return queue_code, reserved_extension, level, resolved_tenant_id return queue_code, reserved_agent["extension"], level, resolved_tenant_id
raise HTTPException(status_code=409, detail=f"No available {level} agent right now") raise HTTPException(status_code=409, detail=f"No available {level} agent right now")
extension = bridge._transfer_target_map().get(queue_code) extension = bridge._transfer_target_map().get(queue_code)
@@ -666,6 +753,7 @@ def _escalation_to_out(row: EscalationRow) -> EscalationOut:
summary=row.summary, summary=row.summary,
status=row.status, status=row.status,
assigned_agent_id=row.assigned_agent_id, assigned_agent_id=row.assigned_agent_id,
attempt_count=row.attempt_count or 0,
requested_at=row.requested_at, requested_at=row.requested_at,
connected_at=row.connected_at, connected_at=row.connected_at,
completed_at=row.completed_at, completed_at=row.completed_at,
@@ -709,28 +797,32 @@ def create_escalation(call_id: str, body: EscalationRequestIn, actor: dict) -> E
session.flush() session.flush()
channel = _resolve_handoff_channel(session, link) channel = _resolve_handoff_channel(session, link)
reserved_extension = _reserve_routing_agent( agent = _reserve_routing_agent(
call_id=call_id, call_id=call_id,
level=body.target_level, level=body.target_level,
tenant_id=link.tenant_id, tenant_id=link.tenant_id,
required_skills=body.required_skills, required_skills=body.required_skills,
) )
if not reserved_extension: if not agent:
escalation.status = "failed" escalation.status = "failed"
escalation.completed_at = now escalation.completed_at = now
session.commit() session.commit()
_emit_escalation_event(
event_type="TransferFailed",
escalation=escalation,
link=link,
extra={"reason": "no_available_agent"},
)
raise HTTPException(status_code=409, detail=f"No available {body.target_level} agent right now") raise HTTPException(status_code=409, detail=f"No available {body.target_level} agent right now")
escalation.real_agent_id = agent["agent_id"]
escalation.attempted_agent_ids_json = json.dumps([agent["agent_id"]], ensure_ascii=False)
session.flush()
_append_escalation_timeline(escalation, link, action="escalation.agent_reserved", extra={"agent_id": agent["agent_id"]})
_emit_escalation_event(event_type="AgentReserved", escalation=escalation, link=link, extra={"agent_id": agent["agent_id"]})
try: try:
bridge._ami_action( _redirect_channel_to_agent(channel=channel, extension=agent["extension"])
"Redirect",
{
"Channel": channel,
"Context": bridge._transfer_context(),
"Exten": reserved_extension,
"Priority": 1,
},
)
except Exception as exc: except Exception as exc:
try: try:
bridge._release_routing_agent(call_id) bridge._release_routing_agent(call_id)
@@ -739,10 +831,13 @@ def create_escalation(call_id: str, body: EscalationRequestIn, actor: dict) -> E
escalation.status = "failed" escalation.status = "failed"
escalation.completed_at = utc_now_iso() escalation.completed_at = utc_now_iso()
session.commit() session.commit()
_append_escalation_timeline(escalation, link, action="escalation.transfer_failed", extra={"agent_id": agent["agent_id"], "error": str(exc)})
_emit_escalation_event(event_type="TransferFailed", escalation=escalation, link=link, extra={"agent_id": agent["agent_id"], "error": str(exc)})
raise HTTPException(status_code=502, detail=f"Failed to redirect call to agent: {exc}") from exc raise HTTPException(status_code=502, detail=f"Failed to redirect call to agent: {exc}") from exc
set_routing_agent_status(agent["agent_id"], "RINGING")
escalation.status = "ringing" escalation.status = "ringing"
escalation.assigned_agent_id = reserved_extension escalation.assigned_agent_id = agent["extension"]
link.current_level = body.target_level link.current_level = body.target_level
link.required_skills_json = json.dumps(body.required_skills, ensure_ascii=False) link.required_skills_json = json.dumps(body.required_skills, ensure_ascii=False)
link.priority = body.priority link.priority = body.priority
@@ -751,11 +846,114 @@ def create_escalation(call_id: str, body: EscalationRequestIn, actor: dict) -> E
link.operator_extension = None link.operator_extension = None
link.updated_at = now link.updated_at = now
session.commit() session.commit()
_append_escalation_timeline(escalation, link, action="escalation.agent_ringing", extra={"agent_id": agent["agent_id"]})
_emit_escalation_event(event_type="AgentRinging", escalation=escalation, link=link, extra={"agent_id": agent["agent_id"]})
return _escalation_to_out(escalation) return _escalation_to_out(escalation)
finally: finally:
session.close() session.close()
def retry_escalation_no_answer(session, *, call_id: str, dial_outcome: str) -> None:
"""AC-08 / ТЗ §14: агент не ответил — освободить его и попробовать следующего.
Вызывается из bridge_processing при нативном AMI DialEnd/Hangup с исходом
NOANSWER/BUSY/CANCEL/CHANUNAVAIL/CONGESTION на канале агента, зарезервированного
под активную (status='ringing') эскалацию этого call_id.
"""
bridge = _bridge_app()
escalation = session.execute(
select(EscalationRow).where(
EscalationRow.call_id == call_id,
EscalationRow.status == "ringing",
)
).scalar_one_or_none()
if escalation is None:
return
link = bridge._find_call_link(session, call_id)
if link is None:
return
now = utc_now_iso()
no_answer_agent_id = escalation.real_agent_id
if no_answer_agent_id:
try:
routing_release_by_agent_id(no_answer_agent_id)
except Exception:
LOGGER.warning("bridge.escalation_retry_release_failed call_id=%s agent_id=%s", call_id, no_answer_agent_id)
escalation.attempt_count = (escalation.attempt_count or 0) + 1
session.flush()
_append_escalation_timeline(
escalation,
link,
action="escalation.agent_no_answer",
extra={"agent_id": no_answer_agent_id, "dial_outcome": dial_outcome},
)
_emit_escalation_event(
event_type="AgentNoAnswer",
escalation=escalation,
link=link,
extra={"agent_id": no_answer_agent_id, "dial_outcome": dial_outcome},
)
attempted_ids = json.loads(escalation.attempted_agent_ids_json or "[]")
channel = _resolve_handoff_channel(session, link)
required_skills = json.loads(escalation.required_skills_json or "[]")
next_agent = _reserve_routing_agent(
call_id=call_id,
level=escalation.to_level,
tenant_id=escalation.tenant_id,
required_skills=required_skills,
exclude_agent_ids=attempted_ids,
)
if not next_agent:
escalation.status = "failed"
escalation.completed_at = now
session.commit()
_append_escalation_timeline(escalation, link, action="escalation.failed_no_agents")
_emit_escalation_event(event_type="TransferFailed", escalation=escalation, link=link, extra={"reason": "no_more_agents"})
LOGGER.warning("bridge.escalation_retry_exhausted call_id=%s attempts=%s", call_id, escalation.attempt_count)
return
escalation.real_agent_id = next_agent["agent_id"]
escalation.attempted_agent_ids_json = json.dumps(attempted_ids + [next_agent["agent_id"]], ensure_ascii=False)
session.flush()
_append_escalation_timeline(escalation, link, action="escalation.agent_reserved", extra={"agent_id": next_agent["agent_id"]})
_emit_escalation_event(event_type="AgentReserved", escalation=escalation, link=link, extra={"agent_id": next_agent["agent_id"]})
try:
_redirect_channel_to_agent(channel=channel, extension=next_agent["extension"])
except Exception as exc:
try:
routing_release_by_agent_id(next_agent["agent_id"])
except Exception:
pass
escalation.status = "failed"
escalation.completed_at = utc_now_iso()
session.commit()
_append_escalation_timeline(escalation, link, action="escalation.transfer_failed", extra={"agent_id": next_agent["agent_id"], "error": str(exc)})
_emit_escalation_event(event_type="TransferFailed", escalation=escalation, link=link, extra={"agent_id": next_agent["agent_id"], "error": str(exc)})
return
set_routing_agent_status(next_agent["agent_id"], "RINGING")
escalation.assigned_agent_id = next_agent["extension"]
session.commit()
_append_escalation_timeline(escalation, link, action="escalation.agent_ringing", extra={"agent_id": next_agent["agent_id"]})
_emit_escalation_event(event_type="AgentRinging", escalation=escalation, link=link, extra={"agent_id": next_agent["agent_id"]})
def routing_release_by_agent_id(agent_id: str, *, next_status: str = "AVAILABLE") -> None:
bridge = _bridge_app()
bridge._post_json(
f"{bridge._routing_service_url()}/internal/routing/release-agent",
{"agent_id": agent_id, "next_status": next_status},
timeout_seconds=bridge._callcontrol_side_effect_timeout_seconds(),
max_attempts=1,
retry_backoff_seconds=0.0,
)
def update_call_ai_state(call_id: str, body: VoiceAICallStateUpdateIn, actor: dict) -> VoiceLiveCallOut: def update_call_ai_state(call_id: str, body: VoiceAICallStateUpdateIn, actor: dict) -> VoiceLiveCallOut:
bridge = _bridge_app() bridge = _bridge_app()
assert_trusted_voice_runtime_actor(actor) assert_trusted_voice_runtime_actor(actor)
+27 -4
View File
@@ -18,6 +18,7 @@ from services.shared.models import (
QueueOut, QueueOut,
RoutingAgentReserveIn, RoutingAgentReserveIn,
RoutingAgentReserveOut, RoutingAgentReserveOut,
RoutingAgentStatusIn,
) )
from services.shared.security import require_roles from services.shared.security import require_roles
from services.shared.sql_init import init_sql_schema from services.shared.sql_init import init_sql_schema
@@ -343,6 +344,7 @@ def _escalation_to_out(row: EscalationRow) -> EscalationOut:
summary=row.summary, summary=row.summary,
status=row.status, status=row.status,
assigned_agent_id=row.assigned_agent_id, assigned_agent_id=row.assigned_agent_id,
attempt_count=row.attempt_count or 0,
requested_at=row.requested_at, requested_at=row.requested_at,
connected_at=row.connected_at, connected_at=row.connected_at,
completed_at=row.completed_at, completed_at=row.completed_at,
@@ -399,11 +401,32 @@ def release_agent_endpoint(
_: dict = Depends(require_roles(Role.ADMIN)), _: dict = Depends(require_roles(Role.ADMIN)),
) -> dict: ) -> dict:
call_id = str(payload.get("call_id") or "").strip() call_id = str(payload.get("call_id") or "").strip()
if not call_id: agent_id = str(payload.get("agent_id") or "").strip()
raise HTTPException(status_code=400, detail="call_id is required") next_status = str(payload.get("next_status") or "AVAILABLE").strip() or "AVAILABLE"
if not call_id and not agent_id:
raise HTTPException(status_code=400, detail="call_id or agent_id is required")
session = get_session() session = get_session()
try: try:
agent = routing_engine.release_agent_by_call_id(session, call_id=call_id) if agent_id:
return {"call_id": call_id, "released_agent_id": agent.agent_id if agent else None} agent = routing_engine.release_agent_by_id(session, agent_id=agent_id, next_status=next_status)
else:
agent = routing_engine.release_agent_by_call_id(session, call_id=call_id)
return {"call_id": call_id or None, "released_agent_id": agent.agent_id if agent else None}
finally:
session.close()
@app.patch("/internal/routing/agents/{agent_id}/status")
def set_agent_status_endpoint(
agent_id: str,
payload: RoutingAgentStatusIn,
_: dict = Depends(require_roles(Role.ADMIN)),
) -> dict:
session = get_session()
try:
agent = routing_engine.set_agent_status(session, agent_id=agent_id, status=payload.status)
if agent is None:
raise HTTPException(status_code=404, detail="Agent not found")
return {"agent_id": agent.agent_id, "status": agent.status}
finally: finally:
session.close() session.close()
+36
View File
@@ -1,6 +1,7 @@
from __future__ import annotations from __future__ import annotations
import json import json
from datetime import datetime
from sqlalchemy import select, text from sqlalchemy import select, text
@@ -117,6 +118,41 @@ def release_agent_by_call_id(session, *, call_id: str) -> AgentRow | None:
return agent return agent
def set_agent_status(session, *, agent_id: str, status: str) -> AgentRow | None:
agent = session.execute(
select(AgentRow).where(AgentRow.agent_id == agent_id)
).scalar_one_or_none()
if agent is None:
return None
agent.status = status
agent.updated_at = utc_now_iso()
session.commit()
return agent
def sweep_after_call_work(session, *, older_than_seconds: int) -> list[AgentRow]:
cutoff = utc_now_iso()
rows = session.execute(
select(AgentRow).where(AgentRow.status == "AFTER_CALL_WORK")
).scalars().all()
swept: list[AgentRow] = []
for agent in rows:
try:
age_seconds = (
datetime.fromisoformat(cutoff) - datetime.fromisoformat(str(agent.updated_at))
).total_seconds()
except (TypeError, ValueError):
continue
if age_seconds >= older_than_seconds:
agent.status = "AVAILABLE"
agent.current_call_id = None
agent.updated_at = cutoff
swept.append(agent)
if swept:
session.commit()
return swept
def release_agent_by_id(session, *, agent_id: str, next_status: str = "AVAILABLE") -> AgentRow | None: def release_agent_by_id(session, *, agent_id: str, next_status: str = "AVAILABLE") -> AgentRow | None:
agent = session.execute( agent = session.execute(
select(AgentRow).where(AgentRow.agent_id == agent_id) select(AgentRow).where(AgentRow.agent_id == agent_id)
+5
View File
@@ -1311,11 +1311,16 @@ class EscalationOut(BaseModel):
summary: str | None = None summary: str | None = None
status: str status: str
assigned_agent_id: str | None = None assigned_agent_id: str | None = None
attempt_count: int = 0
requested_at: str requested_at: str
connected_at: str | None = None connected_at: str | None = None
completed_at: str | None = None completed_at: str | None = None
class RoutingAgentStatusIn(BaseModel):
status: str = Field(min_length=1)
class RoutingAgentReserveIn(BaseModel): class RoutingAgentReserveIn(BaseModel):
call_id: str = Field(min_length=1) call_id: str = Field(min_length=1)
level: AgentLevel level: AgentLevel
+3
View File
@@ -846,6 +846,9 @@ class EscalationRow(Base):
summary: Mapped[str | None] = mapped_column(Text, nullable=True) summary: Mapped[str | None] = mapped_column(Text, nullable=True)
status: Mapped[str] = mapped_column(String(32), index=True, default="requested") status: Mapped[str] = mapped_column(String(32), index=True, default="requested")
assigned_agent_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True) assigned_agent_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
real_agent_id: Mapped[str | None] = mapped_column(String(64), nullable=True, index=True)
attempted_agent_ids_json: Mapped[str] = mapped_column(Text, default="[]")
attempt_count: Mapped[int] = mapped_column(Integer, default=0)
requested_at: Mapped[str] = mapped_column(String(64), index=True) requested_at: Mapped[str] = mapped_column(String(64), index=True)
connected_at: Mapped[str | None] = mapped_column(String(64), nullable=True) connected_at: Mapped[str | None] = mapped_column(String(64), nullable=True)
completed_at: Mapped[str | None] = mapped_column(String(64), nullable=True) completed_at: Mapped[str | None] = mapped_column(String(64), nullable=True)
+2 -1
View File
@@ -2195,6 +2195,7 @@ def test_disabled_bridge_startup_does_not_require_singleton_guard(monkeypatch):
monkeypatch.delenv("ASTERISK_BRIDGE_ENABLED", raising=False) monkeypatch.delenv("ASTERISK_BRIDGE_ENABLED", raising=False)
monkeypatch.setattr(bridge_module, "_ami_loop", _wait_until_stopped) monkeypatch.setattr(bridge_module, "_ami_loop", _wait_until_stopped)
monkeypatch.setattr(bridge_module, "_failed_retry_loop", _wait_until_stopped) monkeypatch.setattr(bridge_module, "_failed_retry_loop", _wait_until_stopped)
monkeypatch.setattr(bridge_module, "_acw_sweep_loop", _wait_until_stopped)
monkeypatch.setattr( monkeypatch.setattr(
bridge_module, bridge_module,
"_try_acquire_bridge_singleton_guard", "_try_acquire_bridge_singleton_guard",
@@ -2202,7 +2203,7 @@ def test_disabled_bridge_startup_does_not_require_singleton_guard(monkeypatch):
) )
bridge_module._startup() bridge_module._startup()
assert len(bridge_module._background_threads_alive()) == 2 assert len(bridge_module._background_threads_alive()) == 3
def test_shutdown_releases_singleton_guard(monkeypatch): def test_shutdown_releases_singleton_guard(monkeypatch):
+190
View File
@@ -0,0 +1,190 @@
import json
import pytest
import services.asterisk_bridge_service.voice_ai as voice_ai
from services.routing_service import engine as routing_engine
from services.shared.core import new_id, utc_now_iso
from services.shared.db import get_session
from services.shared.sql_init import init_sql_schema
from services.shared.sql_models import AgentRow, AsteriskCallLinkRow, EscalationRow
def _make_agent(session, *, extension: str, status: str = "AVAILABLE"):
now = utc_now_iso()
row = AgentRow(
agent_id=new_id("agt"),
tenant_ids_json="[]",
extension=extension,
endpoint=None,
display_name=f"Agent {extension}",
level="L2",
skills_json="[]",
status=status,
max_concurrent_calls=1,
enabled=True,
created_at=now,
updated_at=now,
)
session.add(row)
session.commit()
session.refresh(row)
return row
def _make_call_link(session, *, call_id: str):
now = utc_now_iso()
link = AsteriskCallLinkRow(
call_id=call_id,
linked_id=call_id,
queue_code="voice_lab_ai",
queue_id="que_test",
interaction_id="int_test",
status="active",
telephony_status="ringing",
channel_name=f"PJSIP/tele2-kazgaz-{call_id}",
current_level="L2",
started_at=now,
updated_at=now,
)
session.add(link)
session.commit()
session.refresh(link)
return link
def _make_ringing_escalation(session, *, call_id: str, agent: AgentRow):
now = utc_now_iso()
escalation = EscalationRow(
escalation_id=new_id("esc"),
call_id=call_id,
tenant_id=None,
from_level="L1",
to_level="L2",
reason_code="AI_UNABLE_TO_RESOLVE",
required_skills_json="[]",
priority=3,
status="ringing",
assigned_agent_id=agent.extension,
real_agent_id=agent.agent_id,
attempted_agent_ids_json=json.dumps([agent.agent_id]),
requested_at=now,
)
session.add(escalation)
session.commit()
session.refresh(escalation)
return escalation
def _patch_routing_over_http(monkeypatch, session):
"""Make voice_ai's HTTP-facing routing helpers operate on the same test session directly."""
def fake_reserve(*, call_id, level, tenant_id, required_skills=None, exclude_agent_ids=None):
agent = routing_engine.reserve_agent(
session,
call_id=call_id,
level=level,
tenant_id=tenant_id,
required_skills=required_skills,
exclude_agent_ids=exclude_agent_ids,
)
if agent is None:
return None
return {
"agent_id": agent.agent_id,
"extension": agent.extension,
"endpoint": agent.endpoint,
"display_name": agent.display_name,
}
def fake_release_by_agent_id(agent_id, *, next_status="AVAILABLE"):
routing_engine.release_agent_by_id(session, agent_id=agent_id, next_status=next_status)
def fake_set_status(agent_id, status):
routing_engine.set_agent_status(session, agent_id=agent_id, status=status)
monkeypatch.setattr(voice_ai, "_reserve_routing_agent", fake_reserve)
monkeypatch.setattr(voice_ai, "routing_release_by_agent_id", fake_release_by_agent_id)
monkeypatch.setattr(voice_ai, "set_routing_agent_status", fake_set_status)
monkeypatch.setattr(voice_ai, "_append_escalation_timeline", lambda *a, **k: None)
monkeypatch.setattr(voice_ai, "_emit_escalation_event", lambda *a, **k: None)
monkeypatch.setattr(voice_ai, "_resolve_handoff_channel", lambda session, link: link.channel_name)
def test_retry_escalation_no_answer_moves_to_next_available_agent(monkeypatch):
init_sql_schema()
redirected_to: list[str] = []
session = get_session()
try:
_patch_routing_over_http(monkeypatch, session)
monkeypatch.setattr(
voice_ai,
"_redirect_channel_to_agent",
lambda *, channel, extension: redirected_to.append(extension),
)
agent_a = _make_agent(session, extension="2001")
agent_b = _make_agent(session, extension="2002")
routing_engine.reserve_agent(session, call_id="call-retry-1", level="L2", tenant_id=None, required_skills=[])
link = _make_call_link(session, call_id="call-retry-1")
escalation = _make_ringing_escalation(session, call_id="call-retry-1", agent=agent_a)
voice_ai.retry_escalation_no_answer(session, call_id="call-retry-1", dial_outcome="NOANSWER")
session.refresh(escalation)
session.refresh(agent_a)
session.refresh(agent_b)
assert escalation.attempt_count == 1
assert escalation.status == "ringing"
assert escalation.real_agent_id == agent_b.agent_id
assert escalation.assigned_agent_id == agent_b.extension
assert json.loads(escalation.attempted_agent_ids_json) == [agent_a.agent_id, agent_b.agent_id]
assert agent_a.status == "AVAILABLE"
assert agent_a.current_call_id is None
assert agent_b.status == "RINGING"
assert agent_b.current_call_id == "call-retry-1"
assert redirected_to == [agent_b.extension]
finally:
session.close()
def test_retry_escalation_no_answer_exhausts_pool_marks_failed(monkeypatch):
init_sql_schema()
session = get_session()
try:
_patch_routing_over_http(monkeypatch, session)
monkeypatch.setattr(voice_ai, "_redirect_channel_to_agent", lambda **kwargs: None)
agent_a = _make_agent(session, extension="3001")
routing_engine.reserve_agent(session, call_id="call-retry-2", level="L2", tenant_id=None, required_skills=[])
link = _make_call_link(session, call_id="call-retry-2")
escalation = _make_ringing_escalation(session, call_id="call-retry-2", agent=agent_a)
voice_ai.retry_escalation_no_answer(session, call_id="call-retry-2", dial_outcome="NOANSWER")
session.refresh(escalation)
session.refresh(agent_a)
assert escalation.attempt_count == 1
assert escalation.status == "failed"
assert agent_a.status == "AVAILABLE"
finally:
session.close()
def test_retry_escalation_no_answer_ignores_calls_without_ringing_escalation(monkeypatch):
init_sql_schema()
called = []
session = get_session()
try:
_patch_routing_over_http(monkeypatch, session)
monkeypatch.setattr(voice_ai, "_redirect_channel_to_agent", lambda **kwargs: called.append(kwargs))
_make_call_link(session, call_id="call-no-escalation")
voice_ai.retry_escalation_no_answer(session, call_id="call-no-escalation", dial_outcome="NOANSWER")
assert called == []
finally:
session.close()
+48
View File
@@ -1,3 +1,4 @@
from datetime import datetime, timedelta, timezone
import json import json
from sqlalchemy import select from sqlalchemy import select
@@ -134,3 +135,50 @@ def test_reserve_agent_excludes_disabled_and_excluded_ids():
assert reserved.agent_id != excluded.agent_id assert reserved.agent_id != excluded.agent_id
finally: finally:
session.close() session.close()
def test_set_agent_status_transitions_ringing_to_talking():
init_sql_schema()
session = get_session()
try:
agent = _make_agent(session, level="L2")
routing_engine.reserve_agent(session, call_id="call-status", level="L2", tenant_id=None, required_skills=[])
ringing = routing_engine.set_agent_status(session, agent_id=agent.agent_id, status="RINGING")
assert ringing is not None
assert ringing.status == "RINGING"
talking = routing_engine.set_agent_status(session, agent_id=agent.agent_id, status="TALKING")
assert talking.status == "TALKING"
missing = routing_engine.set_agent_status(session, agent_id="unknown-agent", status="AVAILABLE")
assert missing is None
finally:
session.close()
def test_sweep_after_call_work_releases_only_expired_agents():
init_sql_schema()
session = get_session()
try:
stale = _make_agent(session, level="L2", status="AFTER_CALL_WORK")
fresh = _make_agent(session, level="L2", status="AFTER_CALL_WORK")
stale_ts = (datetime.now(timezone.utc) - timedelta(seconds=120)).replace(microsecond=0).isoformat()
stale.updated_at = stale_ts
stale.current_call_id = "call-stale"
session.commit()
swept = routing_engine.sweep_after_call_work(session, older_than_seconds=30)
swept_ids = {a.agent_id for a in swept}
assert stale.agent_id in swept_ids
assert fresh.agent_id not in swept_ids
session.refresh(stale)
session.refresh(fresh)
assert stale.status == "AVAILABLE"
assert stale.current_call_id is None
assert fresh.status == "AFTER_CALL_WORK"
finally:
session.close()