diff --git a/contracts/events/agent.connected.json b/contracts/events/agent.connected.json new file mode 100644 index 0000000..522c615 --- /dev/null +++ b/contracts/events/agent.connected.json @@ -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" } + } +} diff --git a/contracts/events/agent.no_answer.json b/contracts/events/agent.no_answer.json new file mode 100644 index 0000000..78d1a7d --- /dev/null +++ b/contracts/events/agent.no_answer.json @@ -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" } + } +} diff --git a/contracts/events/agent.reserved.json b/contracts/events/agent.reserved.json new file mode 100644 index 0000000..0541026 --- /dev/null +++ b/contracts/events/agent.reserved.json @@ -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" } + } +} diff --git a/contracts/events/agent.ringing.json b/contracts/events/agent.ringing.json new file mode 100644 index 0000000..bc98216 --- /dev/null +++ b/contracts/events/agent.ringing.json @@ -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" } + } +} diff --git a/contracts/events/transfer.completed.json b/contracts/events/transfer.completed.json new file mode 100644 index 0000000..cb990fd --- /dev/null +++ b/contracts/events/transfer.completed.json @@ -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" } + } +} diff --git a/contracts/events/transfer.failed.json b/contracts/events/transfer.failed.json new file mode 100644 index 0000000..83032ec --- /dev/null +++ b/contracts/events/transfer.failed.json @@ -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" } + } +} diff --git a/migrations/sql/0033_escalation_retry_and_agent_state_postgres.sql b/migrations/sql/0033_escalation_retry_and_agent_state_postgres.sql new file mode 100644 index 0000000..6ba5bab --- /dev/null +++ b/migrations/sql/0033_escalation_retry_and_agent_state_postgres.sql @@ -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); diff --git a/migrations/sql/0033_escalation_retry_and_agent_state_sqlite.sql b/migrations/sql/0033_escalation_retry_and_agent_state_sqlite.sql new file mode 100644 index 0000000..bb5deb2 --- /dev/null +++ b/migrations/sql/0033_escalation_retry_and_agent_state_sqlite.sql @@ -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); diff --git a/services/asterisk_bridge_service/ami.py b/services/asterisk_bridge_service/ami.py index bf94d72..8177155 100644 --- a/services/asterisk_bridge_service/ami.py +++ b/services/asterisk_bridge_service/ami.py @@ -412,13 +412,18 @@ def ami_loop(stop_event=None) -> None: raise RuntimeError("AMI connection closed") if not frame: 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 - 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: bridge._STATE.set_connected(False) bridge._STATE.set_error(str(exc)) diff --git a/services/asterisk_bridge_service/app.py b/services/asterisk_bridge_service/app.py index 9b3fd67..9a91982 100644 --- a/services/asterisk_bridge_service/app.py +++ b/services/asterisk_bridge_service/app.py @@ -262,6 +262,7 @@ def _start_background_threads() -> None: extra_loops: list[tuple[str, Callable[..., Any]]] = [] if _ivr_fastagi_enabled(): 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( ami_loop=_ami_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 _create_escalation = bridge_voice_ai.create_escalation _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 _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_operator_connected = bridge_processing.process_operator_connected _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 _record_ami_payload = bridge_processing.record_ami_payload _retry_failed_events_once = bridge_processing.retry_failed_events_once diff --git a/services/asterisk_bridge_service/bridge_processing.py b/services/asterisk_bridge_service/bridge_processing.py index 1118f6e..14d4b19 100644 --- a/services/asterisk_bridge_service/bridge_processing.py +++ b/services/asterisk_bridge_service/bridge_processing.py @@ -4,6 +4,7 @@ from datetime import datetime, timezone import hashlib import json import logging +import os from pathlib import Path import threading from typing import Any @@ -884,18 +885,34 @@ def process_call_ended( open_escalation = session.execute( select(EscalationRow).where( EscalationRow.call_id == row.call_id, - EscalationRow.status.in_(["requested", "ringing"]), + EscalationRow.status.in_(["requested", "ringing", "connected"]), ) ).scalar_one_or_none() + was_talking = open_escalation is not None and open_escalation.status == "connected" if open_escalation is not None: open_escalation.status = "completed" if answered else "failed" open_escalation.completed_at = now if answered and not open_escalation.connected_at: open_escalation.connected_at = now - try: - bridge._release_routing_agent(row.call_id) - except Exception: - pass + if was_talking and open_escalation and open_escalation.real_agent_id: + try: + bridge._set_routing_agent_status(open_escalation.real_agent_id, "AFTER_CALL_WORK") + 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: bridge._notify_voice_ai_telephony_event( voice_session_id=link.voice_session_id, @@ -969,6 +986,48 @@ def process_operator_connected( link.telephony_status = "connected" link.connected_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( session, link, @@ -1173,6 +1232,55 @@ def process_recording_ready( 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: bridge = _bridge_app() 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) elif row.ami_event_name == f"{bridge._ami_prefix()}RecordingReady": bridge._process_recording_ready(session, row, payload) + elif row.ami_event_name in {"DialEnd", "Hangup"}: + bridge._process_agent_dial_outcome(session, row, payload) else: row.forward_status = "received" row.updated_at = utc_now_iso() @@ -1216,9 +1326,9 @@ def process_bridge_row(session, row: AsteriskEventLogRow) -> AsteriskEventLogRow 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() - 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) linked_id = bridge._extract_linked_id(payload, 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) +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: bridge = _bridge_app() active_stop_event = stop_event or bridge._background_stop_event() diff --git a/services/asterisk_bridge_service/voice_ai.py b/services/asterisk_bridge_service/voice_ai.py index 97f2958..4f4b4c9 100644 --- a/services/asterisk_bridge_service/voice_ai.py +++ b/services/asterisk_bridge_service/voice_ai.py @@ -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( *, call_id: str, @@ -188,7 +239,8 @@ def _reserve_routing_agent( level: str, tenant_id: str | None, required_skills: list[str] | None = None, -) -> str | None: + exclude_agent_ids: list[str] | None = None, +) -> dict[str, Any] | None: bridge = _bridge_app() try: response = bridge._post_json( @@ -198,7 +250,7 @@ def _reserve_routing_agent( "level": level, "tenant_id": tenant_id, "required_skills": required_skills or [], - "exclude_agent_ids": [], + "exclude_agent_ids": exclude_agent_ids or [], }, timeout_seconds=bridge._callcontrol_side_effect_timeout_seconds(), max_attempts=1, @@ -212,8 +264,43 @@ def _reserve_routing_agent( tenant_id, ) return None + agent_id = str((response or {}).get("agent_id") 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( @@ -233,14 +320,14 @@ def _resolve_handoff_extension( level = target_level or bridge._routing_level_for_queue_code(queue_code) 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) - reserved_extension = _reserve_routing_agent( + reserved_agent = _reserve_routing_agent( call_id=call_id, level=level, tenant_id=resolved_tenant_id, required_skills=required_skills, ) - if reserved_extension: - return queue_code, reserved_extension, level, resolved_tenant_id + if reserved_agent: + return queue_code, reserved_agent["extension"], level, resolved_tenant_id raise HTTPException(status_code=409, detail=f"No available {level} agent right now") extension = bridge._transfer_target_map().get(queue_code) @@ -666,6 +753,7 @@ def _escalation_to_out(row: EscalationRow) -> EscalationOut: summary=row.summary, status=row.status, assigned_agent_id=row.assigned_agent_id, + attempt_count=row.attempt_count or 0, requested_at=row.requested_at, connected_at=row.connected_at, completed_at=row.completed_at, @@ -709,28 +797,32 @@ def create_escalation(call_id: str, body: EscalationRequestIn, actor: dict) -> E session.flush() channel = _resolve_handoff_channel(session, link) - reserved_extension = _reserve_routing_agent( + agent = _reserve_routing_agent( call_id=call_id, level=body.target_level, tenant_id=link.tenant_id, required_skills=body.required_skills, ) - if not reserved_extension: + if not agent: escalation.status = "failed" escalation.completed_at = now 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") + 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: - bridge._ami_action( - "Redirect", - { - "Channel": channel, - "Context": bridge._transfer_context(), - "Exten": reserved_extension, - "Priority": 1, - }, - ) + _redirect_channel_to_agent(channel=channel, extension=agent["extension"]) except Exception as exc: try: 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.completed_at = utc_now_iso() 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 + set_routing_agent_status(agent["agent_id"], "RINGING") escalation.status = "ringing" - escalation.assigned_agent_id = reserved_extension + escalation.assigned_agent_id = agent["extension"] link.current_level = body.target_level link.required_skills_json = json.dumps(body.required_skills, ensure_ascii=False) link.priority = body.priority @@ -751,11 +846,114 @@ def create_escalation(call_id: str, body: EscalationRequestIn, actor: dict) -> E link.operator_extension = None link.updated_at = now 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) finally: 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: bridge = _bridge_app() assert_trusted_voice_runtime_actor(actor) diff --git a/services/routing_service/app.py b/services/routing_service/app.py index dc16848..1b1fab8 100644 --- a/services/routing_service/app.py +++ b/services/routing_service/app.py @@ -18,6 +18,7 @@ from services.shared.models import ( QueueOut, RoutingAgentReserveIn, RoutingAgentReserveOut, + RoutingAgentStatusIn, ) from services.shared.security import require_roles from services.shared.sql_init import init_sql_schema @@ -343,6 +344,7 @@ def _escalation_to_out(row: EscalationRow) -> EscalationOut: summary=row.summary, status=row.status, assigned_agent_id=row.assigned_agent_id, + attempt_count=row.attempt_count or 0, requested_at=row.requested_at, connected_at=row.connected_at, completed_at=row.completed_at, @@ -399,11 +401,32 @@ def release_agent_endpoint( _: dict = Depends(require_roles(Role.ADMIN)), ) -> dict: call_id = str(payload.get("call_id") or "").strip() - if not call_id: - raise HTTPException(status_code=400, detail="call_id is required") + agent_id = str(payload.get("agent_id") or "").strip() + 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() try: - agent = routing_engine.release_agent_by_call_id(session, call_id=call_id) - return {"call_id": call_id, "released_agent_id": agent.agent_id if agent else None} + if agent_id: + 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: session.close() diff --git a/services/routing_service/engine.py b/services/routing_service/engine.py index 09b5f96..370fbf3 100644 --- a/services/routing_service/engine.py +++ b/services/routing_service/engine.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +from datetime import datetime from sqlalchemy import select, text @@ -117,6 +118,41 @@ def release_agent_by_call_id(session, *, call_id: str) -> AgentRow | None: 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: agent = session.execute( select(AgentRow).where(AgentRow.agent_id == agent_id) diff --git a/services/shared/models.py b/services/shared/models.py index ff83c45..532c647 100644 --- a/services/shared/models.py +++ b/services/shared/models.py @@ -1311,11 +1311,16 @@ class EscalationOut(BaseModel): summary: str | None = None status: str assigned_agent_id: str | None = None + attempt_count: int = 0 requested_at: str connected_at: str | None = None completed_at: str | None = None +class RoutingAgentStatusIn(BaseModel): + status: str = Field(min_length=1) + + class RoutingAgentReserveIn(BaseModel): call_id: str = Field(min_length=1) level: AgentLevel diff --git a/services/shared/sql_models.py b/services/shared/sql_models.py index 1b9d6a1..b809545 100644 --- a/services/shared/sql_models.py +++ b/services/shared/sql_models.py @@ -846,6 +846,9 @@ class EscalationRow(Base): summary: Mapped[str | None] = mapped_column(Text, nullable=True) 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) + 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) connected_at: Mapped[str | None] = mapped_column(String(64), nullable=True) completed_at: Mapped[str | None] = mapped_column(String(64), nullable=True) diff --git a/tests/test_asterisk_bridge_service.py b/tests/test_asterisk_bridge_service.py index 069168c..4b35a7e 100644 --- a/tests/test_asterisk_bridge_service.py +++ b/tests/test_asterisk_bridge_service.py @@ -2195,6 +2195,7 @@ def test_disabled_bridge_startup_does_not_require_singleton_guard(monkeypatch): monkeypatch.delenv("ASTERISK_BRIDGE_ENABLED", raising=False) monkeypatch.setattr(bridge_module, "_ami_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( bridge_module, "_try_acquire_bridge_singleton_guard", @@ -2202,7 +2203,7 @@ def test_disabled_bridge_startup_does_not_require_singleton_guard(monkeypatch): ) 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): diff --git a/tests/test_escalation_retry.py b/tests/test_escalation_retry.py new file mode 100644 index 0000000..4de0120 --- /dev/null +++ b/tests/test_escalation_retry.py @@ -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() diff --git a/tests/test_routing_engine.py b/tests/test_routing_engine.py index d0aab69..6a149c0 100644 --- a/tests/test_routing_engine.py +++ b/tests/test_routing_engine.py @@ -1,3 +1,4 @@ +from datetime import datetime, timedelta, timezone import json from sqlalchemy import select @@ -134,3 +135,50 @@ def test_reserve_agent_excludes_disabled_and_excluded_ids(): assert reserved.agent_id != excluded.agent_id finally: 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()