Spaces:
Sleeping
Sleeping
File size: 23,925 Bytes
6d23228 d402333 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 | """
Orchestrator β Task Queue Dialogue Management
===============================================
Replaces the one-intent-one-slot FSM. Fixes every issue in the feedback:
P0 Compound requests β NLU decomposes; ALL tasks enter a queue;
each is acknowledged and executed in order.
Nothing is silently dropped.
P0 Wrong destructive β destructive intents require BOTH
intent confidence β₯ 0.75 AND explicit confirmation.
Below threshold β clarifying question, never action.
P1 Named entity ignored β slots from NLU prefill the task; only
genuinely missing slots are asked.
P1 Dead-end fallback β "unknown" triggers ONE clarifying question with
candidate intents; second failure β human handoff
WITH full context attached to the CRM ticket.
P2 No balance guard β transfers are checked against balance before
confirmation; insufficient funds is surfaced.
P2 Number formatting β fmt_amount() used everywhere: "β¦350,000".
"""
import re
import uuid
import logging
from dataclasses import dataclass, field
from datetime import datetime
from typing import Optional
from nlu import NLU, INTENT_SCHEMA
logger = logging.getLogger(__name__)
CONFIDENCE_GATE_DESTRUCTIVE = 0.75 # below this, never act on money/card intents
CONFIDENCE_GATE_NORMAL = 0.50
MAX_CLARIFY_ATTEMPTS = 2 # then human handoff with context
MAX_TURNS = 20
# ββ Formatting (P2) βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
def fmt_amount(raw) -> str:
"""'350000' β 'β¦350,000' β single source of truth for money display."""
try:
n = int(str(raw).replace(",", "").replace(".", "").replace("β¦", "").strip())
return f"β¦{n:,}"
except (ValueError, TypeError):
return f"β¦{raw}"
def amount_int(raw) -> Optional[int]:
try:
return int(str(raw).replace(",", "").replace(".", "").replace("β¦", "").strip())
except (ValueError, TypeError):
return None
# ββ Task model ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
@dataclass
class Task:
task_id: str
intent: str
slots: dict
confidence: float
status: str = "pending" # pending | collecting | confirming | done | failed
result: Optional[str] = None
@property
def required_slots(self) -> list:
return INTENT_SCHEMA.get(self.intent, {}).get("slots", [])
@property
def missing_slots(self) -> list:
# account_id is optional if we already verified identity this session
return [s for s in self.required_slots if s not in self.slots]
@property
def is_destructive(self) -> bool:
return INTENT_SCHEMA.get(self.intent, {}).get("destructive", False)
@dataclass
class Session:
session_id: str = ""
tasks: list = field(default_factory=list) # task queue
active_task: Optional[Task] = None
awaiting_slot: Optional[str] = None
awaiting_confirmation: bool = False
verified_account: Optional[str] = None # identity, session-scoped
balance: Optional[int] = None # cached after lookup
clarify_attempts: int = 0
turn: int = 0
escalated: bool = False
history: list = field(default_factory=list)
started_at: str = field(default_factory=lambda: datetime.utcnow().isoformat())
# ββ Slot questions ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
SLOT_QUESTIONS = {
"account_id": "Could you give me your account number, please?",
"recipient": "Who would you like to send the money to?",
"amount": "How much would you like to send?",
"location": "Which city or area are you in?",
"issue_desc": "Could you briefly describe the problem?",
"order_id": "What is your order number?",
"return_reason": "What is the reason for the return?",
}
class Orchestrator:
def __init__(self, crm=None, nlu: Optional[NLU] = None):
self.crm = crm
self.nlu = nlu or NLU()
def new_session(self) -> Session:
return Session(session_id=str(uuid.uuid4())[:8])
# ββ Main entry ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
def respond(self, english_text: str, hausa_text: str,
session: Session) -> tuple[str, Session, bool]:
session.turn += 1
session.history.append({"role": "user", "en": english_text,
"ha": hausa_text, "turn": session.turn})
if session.turn >= MAX_TURNS:
return self._escalate(session,
"We've been talking a while β let me connect you to a colleague "
"who can wrap this up quickly.")
# 1. NLU with pending-slot context
pending_intent = session.active_task.intent if session.active_task else None
nlu_result = self.nlu.parse(english_text,
pending_intent=pending_intent,
pending_slot=session.awaiting_slot)
tasks_found = nlu_result["tasks"]
# 2. Route the parse
response = self._handle_parse(tasks_found, english_text, session)
# 3. Drive the queue until we need user input
followup = self._drive_queue(session)
if followup:
response = (response + " " + followup).strip() if response else followup
session.history.append({"role": "agent", "en": response,
"ha": None, "turn": session.turn})
return response, session, session.escalated
# ββ Parse handling ββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
def _handle_parse(self, tasks_found: list, raw_text: str,
session: Session) -> str:
parts = []
for parsed in tasks_found:
intent = parsed["intent"]
conf = parsed["confidence"]
slots = parsed["slots"]
# ββ Universal intents ββββββββββββββββββββββββββββββββββββββββββββ
if intent == "human_agent":
return self._escalate(session,
"Of course β connecting you to a human agent now. "
f"Your reference is REF-{session.session_id}. "
"They will see everything we discussed.")[0]
if intent in ("cancel",):
session.active_task = None
session.awaiting_slot = None
session.awaiting_confirmation = False
session.tasks = [t for t in session.tasks if t.status == "done"]
parts.append("Alright, I've cancelled that.")
continue
if intent == "goodbye":
pending = [t for t in session.tasks
if t.status in ("pending", "collecting", "confirming")]
if pending:
parts.append(f"Before you go β we still have "
f"{len(pending)} pending request(s). "
f"Say 'cancel' to drop them, or continue.")
else:
parts.append("Thank you for calling. Have a wonderful day!")
continue
if intent == "greeting" and session.turn <= 2:
parts.append("Hello! I can help you check balances, send money, "
"pay bills, track orders, or report an issue. "
"What can I do for you?")
continue
# ββ Confirmation answers for the active task βββββββββββββββββββββ
if session.awaiting_confirmation and intent in (
"confirmation_yes", "confirmation_no"):
parts.append(self._handle_confirmation(
intent == "confirmation_yes", session))
continue
# ββ Slot answer for the active task ββββββββββββββββββββββββββββββ
if (session.awaiting_slot and session.active_task
and intent in ("unknown", session.active_task.intent)):
filled = self._try_fill_pending_slot(raw_text, slots, session)
if filled:
continue # queue driver will move it forward
# ββ Unknown β clarify, not dead-end (P1) βββββββββββββββββββββββββ
if intent == "unknown" or conf < CONFIDENCE_GATE_NORMAL:
parts.append(self._clarify_or_handoff(raw_text, session))
continue
# ββ New task β enqueue with prefilled slots (P0 + P1) ββββββββββββ
task = Task(task_id=str(uuid.uuid4())[:6],
intent=intent, slots=dict(slots), confidence=conf)
# identity carries over within the session
if session.verified_account and "account_id" in task.required_slots:
task.slots.setdefault("account_id", session.verified_account)
session.tasks.append(task)
session.clarify_attempts = 0
# Acknowledge compound requests explicitly (P0 #1)
new_tasks = [t for t in session.tasks if t.status == "pending"]
if len(new_tasks) > 1:
names = ", then ".join(self._intent_label(t.intent) for t in new_tasks)
parts.insert(0, f"Got it β I'll {names}, one at a time.")
return " ".join(p for p in parts if p)
# ββ Queue driver ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
def _drive_queue(self, session: Session) -> str:
"""
Advance the active task; when it completes, pull the next from the
queue. Stops (returns a question) whenever user input is needed.
"""
out = []
while True:
# promote next task if none active
if session.active_task is None:
nxt = next((t for t in session.tasks if t.status == "pending"), None)
if nxt is None:
break
session.active_task = nxt
nxt.status = "collecting"
# identity verified earlier in this session carries forward β
# never re-ask for the account number (P1)
if session.verified_account and "account_id" in nxt.required_slots:
nxt.slots.setdefault("account_id", session.verified_account)
task = session.active_task
# 1. Missing slots? Ask for exactly ONE (but never one we have)
missing = task.missing_slots
if missing:
slot = missing[0]
session.awaiting_slot = slot
out.append(SLOT_QUESTIONS.get(
slot, f"Could you provide the {slot.replace('_', ' ')}?"))
return " ".join(out)
# 2. Destructive β confidence gate + confirmation (P0 #2)
if task.is_destructive and not session.awaiting_confirmation:
if task.confidence < CONFIDENCE_GATE_DESTRUCTIVE:
session.awaiting_confirmation = True
task.status = "confirming"
out.append(self._describe_action(task) +
" β did I understand that correctly? "
"Please say yes or no.")
return " ".join(out)
# balance guard before confirming a transfer (P2)
guard_msg = self._balance_guard(task, session)
if guard_msg:
task.status = "failed"
session.active_task = None
session.awaiting_slot = None
out.append(guard_msg)
continue
session.awaiting_confirmation = True
task.status = "confirming"
out.append(self._describe_action(task) +
" Say yes to confirm or no to cancel.")
return " ".join(out)
# 3. Execute non-destructive task
result = self._execute(task, session)
task.status = "done"
task.result = result
out.append(result)
session.active_task = None
session.awaiting_slot = None
# queue empty
if out:
remaining = [t for t in session.tasks if t.status == "pending"]
if not remaining and not session.awaiting_confirmation:
out.append("Is there anything else I can help you with?")
return " ".join(out)
# ββ Confirmation handling βββββββββββββββββββββββββββββββββββββββββββββββββ
def _handle_confirmation(self, confirmed: bool, session: Session) -> str:
task = session.active_task
session.awaiting_confirmation = False
if task is None:
return ""
if confirmed:
result = self._execute(task, session)
task.status = "done"
task.result = result
session.active_task = None
return result
task.status = "failed"
session.active_task = None
return ("No problem, I've cancelled that. "
"Just tell me if you'd like to do something else.")
# ββ Slot filling ββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
def _try_fill_pending_slot(self, raw_text: str, nlu_slots: dict,
session: Session) -> bool:
task = session.active_task
slot = session.awaiting_slot
if not task or not slot:
return False
value = nlu_slots.get(slot)
# heuristic extraction for short answers
if not value:
t = raw_text.strip()
if slot == "account_id":
m = re.search(r'\b(\d{6,12})\b', t)
value = m.group(1) if m else None
elif slot == "amount":
m = re.search(r'\b(\d{1,3}(?:[,\.]\d{3})*|\d+)\b', t)
value = m.group(1) if m else None
elif slot == "recipient":
m = re.match(r'^(?:to\s+)?([a-zA-Z]{2,20})$', t)
value = m.group(1) if m else None
elif slot in ("issue_desc", "return_reason", "location"):
# ANY non-trivial answer counts β "too small" is a valid
# return reason (fixes P1 dead-end)
value = t if len(t) >= 2 else None
elif slot == "order_id":
m = re.search(r'\b([a-zA-Z0-9\-]{4,20})\b', t)
value = m.group(1) if m else None
if value:
task.slots[slot] = str(value)
session.awaiting_slot = None
session.clarify_attempts = 0
if slot == "account_id":
session.verified_account = str(value)
return True
return False
# ββ Clarify β handoff (P1 dead-end fix) ββββββββββββββββββββββββββββββββββ
def _clarify_or_handoff(self, raw_text: str, session: Session) -> str:
session.clarify_attempts += 1
if session.clarify_attempts >= MAX_CLARIFY_ATTEMPTS:
return self._escalate(session,
"I want to make sure you get proper help β connecting you "
"to a human agent now. They'll see our whole conversation, "
"so you won't have to repeat anything.")[0]
return ("Sorry β just to be sure I get this right: are you asking "
"about your balance, a transfer, a payment, an order, "
"or something else? You can also say it in your own words.")
# ββ Execution βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
def _execute(self, task: Task, session: Session) -> str:
i = task.intent
s = task.slots
if i == "balance_inquiry":
bal = self._mock_balance(s.get("account_id", session.verified_account or "0"))
session.balance = bal
return (f"Your account balance is {fmt_amount(bal)} "
f"as of {datetime.utcnow().strftime('%d %b %Y')}.")
if i == "send_money":
tx = f"TXN-{session.session_id}-{task.task_id.upper()}"
amt = fmt_amount(s['amount'])
if session.balance is not None:
session.balance -= amount_int(s["amount"]) or 0
return (f"Done β {amt} sent to {s['recipient'].title()}. "
f"Transaction reference: {tx}.")
if i == "bill_payment":
tx = f"TXN-{session.session_id}-{task.task_id.upper()}"
return (f"Your payment of {fmt_amount(s['amount'])} has been "
f"processed. Reference: {tx}.")
if i == "block_card":
return (f"Your card linked to account ending "
f"β¦{s.get('account_id', '')[-4:]} is now blocked. "
f"A replacement can be requested at any branch.")
if i == "branch_info":
loc = s.get("location", "your area")
return (f"Our closest branch to {loc} is at 12 Ahmadu Bello Way β "
f"open Monday to Friday, 8am to 4pm. You can pick up an "
f"ATM card there with a valid ID.")
if i == "card_request":
return ("You can collect a new ATM card at any branch with a "
"valid ID, or I can order one to be delivered β "
"just say 'deliver my card'.")
if i == "track_order":
oid = s.get("order_id", "")
return (f"Order {oid} is out for delivery and should arrive "
f"within 2 business days.")
if i == "return_item":
ticket = self._crm_ticket(session,
f"Return request β order {s.get('order_id','?')} β "
f"reason: {s.get('return_reason','?')}")
return (f"I've registered your return for order "
f"{s.get('order_id','')} (reason: {s.get('return_reason','')}). "
f"Ticket {ticket}. You'll receive a pickup label by SMS.")
if i == "report_issue":
ticket = self._crm_ticket(session, s.get("issue_desc", raw := ""))
return (f"I've created support ticket {ticket}. "
f"Our team will contact you within 24 hours.")
return "Done."
# ββ Balance guard (P2) ββββββββββββββββββββββββββββββββββββββββββββββββββββ
def _balance_guard(self, task: Task, session: Session) -> Optional[str]:
if task.intent != "send_money":
return None
amt = amount_int(task.slots.get("amount"))
if amt is None:
return None
# look up balance if we haven't yet this session
if session.balance is None and task.slots.get("account_id"):
session.balance = self._mock_balance(task.slots["account_id"])
if session.balance is not None and amt > session.balance:
return (f"I can't process that transfer: you asked to send "
f"{fmt_amount(amt)} but your balance is "
f"{fmt_amount(session.balance)}. "
f"Would you like to send a smaller amount?")
return None
# ββ Escalation with context (P1) βββββββββββββββββββββββββββββββββββββββββ
def _escalate(self, session: Session, message: str) -> tuple[str, Session, bool]:
session.escalated = True
transcript = "\n".join(
f"[{h['turn']}] {h['role']}: {h.get('en','')}" for h in session.history)
self._crm_ticket(session,
f"ESCALATION β full context attached:\n{transcript}",
subject=f"Voice escalation {session.session_id}")
pending = [t for t in session.tasks
if t.status in ("pending", "collecting", "confirming")]
if pending:
message += (f" Note for the agent: {len(pending)} request(s) "
f"still open ({', '.join(t.intent for t in pending)}).")
return message, session, True
# ββ Helpers βββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββββ
def _describe_action(self, task: Task) -> str:
s = task.slots
if task.intent == "send_money":
return (f"You want to send {fmt_amount(s['amount'])} "
f"to {s['recipient'].title()}.")
if task.intent == "bill_payment":
return f"You want to pay {fmt_amount(s['amount'])}."
if task.intent == "block_card":
return (f"You want to BLOCK the card on account "
f"ending β¦{s.get('account_id','')[-4:]}.")
return f"You want to {self._intent_label(task.intent)}."
@staticmethod
def _intent_label(intent: str) -> str:
return {
"balance_inquiry": "check your balance",
"send_money": "make a transfer",
"bill_payment": "pay a bill",
"block_card": "block your card",
"branch_info": "find a branch",
"card_request": "get a card",
"track_order": "track your order",
"return_item": "process a return",
"report_issue": "log your issue",
}.get(intent, intent.replace("_", " "))
@staticmethod
def _mock_balance(account_id: str) -> int:
import hashlib
seed = int(hashlib.md5(str(account_id).encode()).hexdigest()[:6], 16)
return (seed % 400_000) + 50_000
def _crm_ticket(self, session: Session, description: str,
subject: str = "") -> str:
if self.crm:
try:
res = self.crm.create_ticket(
subject=subject or f"Voice session {session.session_id}",
description=description)
return res["ticket_id"]
except Exception as e:
logger.warning(f"CRM ticket failed: {e}")
return f"TKT-{session.session_id}-{str(uuid.uuid4())[:4].upper()}"
|