mukul-chauhan-methdai's picture
Direct send_email block to re-fire register_guest with the same name
9845bad
Raw
History Blame Contribute Delete
14.2 kB
"""Session manager for the receptionist state machine.
Single source of truth for the active visitor's session. The face worker
(background thread), conversation controller, and tools (asyncio) all mutate
state through this manager. All mutations take a lock; reads return snapshots.
State changes are dispatched to a subscriber callback so the realtime handler
can push them to the LLM as context events (mirroring how face events are
pushed today via FaceRecognitionWorker.set_face_event_callback).
"""
from __future__ import annotations
import logging
import threading
import time
from dataclasses import dataclass, field, replace
from datetime import datetime
from typing import Any, Callable, Optional
from reachy_mini_receptionist.receptionist_state import ReceptionState
logger = logging.getLogger(__name__)
@dataclass
class VisitorSession:
"""Snapshot of the current visitor's session state."""
current_state: ReceptionState = ReceptionState.IDLE
visitor_name: Optional[str] = None
recognized_face_name: Optional[str] = None
employee_name: Optional[str] = None
purpose: Optional[str] = None
matched_appointment: Optional[dict[str, Any]] = None
email_sent_to: Optional[str] = None
error_message: Optional[str] = None
started_at: float = field(default_factory=time.time)
last_event_at: float = field(default_factory=time.time)
# Last completed visitor transcript — used by register_guest's
# confirmation enforcement so the backend can verify the visitor
# actually said "yes" before saving a face under a name. The realtime
# handler populates this on every input_audio_transcription.completed.
last_user_transcript: Optional[str] = None
# How many times register_guest has refused this session (no
# confirmation, mistranscribed name, etc.). After 2 refusals the
# bot should hand off to an operator instead of looping forever.
register_guest_rejections: int = 0
# Last name proposed via a rejected register_guest(confirmed=false)
# call. Used by send_email's "no visitor identity" block message to
# give a directive error ("call register_guest(name=<this>,
# confirmed=true) now") instead of telling the model to ask the
# visitor again — which it already did 5 seconds ago.
last_register_guest_name: Optional[str] = None
def to_dict(self) -> dict[str, Any]:
"""Serialize for /api/session and LLM context messages."""
return {
"current_state": self.current_state.value,
"visitor_name": self.visitor_name,
"recognized_face_name": self.recognized_face_name,
"employee_name": self.employee_name,
"purpose": self.purpose,
"matched_appointment": self.matched_appointment,
"email_sent_to": self.email_sent_to,
"error_message": self.error_message,
"started_at_iso": datetime.fromtimestamp(self.started_at).isoformat(timespec="seconds"),
"last_event_at_iso": datetime.fromtimestamp(self.last_event_at).isoformat(timespec="seconds"),
}
StateChangeCallback = Callable[[ReceptionState, ReceptionState, VisitorSession], None]
"""Called as ``callback(previous_state, new_state, snapshot)`` on every transition."""
# How long a non-IDLE session can sit without any transition before we
# consider it stale and reset. Visitors who walk up but never speak should
# not hold state forever.
DEFAULT_IDLE_TIMEOUT_SECONDS: float = 60.0
class SessionManager:
"""Thread-safe holder for the active visitor session."""
def __init__(
self,
idle_timeout_seconds: float = DEFAULT_IDLE_TIMEOUT_SECONDS,
visitor_log: Any | None = None,
) -> None:
self._lock = threading.Lock()
self._session = VisitorSession()
self._callback: Optional[StateChangeCallback] = None
self._idle_timeout_seconds = float(idle_timeout_seconds)
# Optional sink: when reset() is called, persist the pre-reset
# session here so both the face-loss path AND the timeout path
# produce log entries.
self._visitor_log = visitor_log
@property
def session(self) -> VisitorSession:
"""Return a snapshot of the current session. Safe from any thread."""
with self._lock:
return replace(self._session)
@property
def current_state(self) -> ReceptionState:
with self._lock:
return self._session.current_state
def subscribe(self, callback: Optional[StateChangeCallback]) -> None:
"""Register (or clear, with ``None``) the state-change callback.
Only one subscriber, matching the face worker's pattern.
"""
with self._lock:
self._callback = callback
def transition(self, new_state: ReceptionState, **updates: Any) -> VisitorSession:
"""Move to a new state and optionally update session fields atomically.
Returns a snapshot. The callback is dispatched OUTSIDE the lock so
subscribers can safely call back in without deadlocking.
"""
with self._lock:
self._validate_updates(updates)
previous_state = self._session.current_state
self._session = replace(
self._session,
current_state=new_state,
last_event_at=time.time(),
**updates,
)
snapshot = replace(self._session)
callback = self._callback
logger.info(
"Session: %s -> %s (visitor=%r employee=%r)",
previous_state.value,
new_state.value,
snapshot.visitor_name,
snapshot.employee_name,
)
if previous_state != new_state and callback is not None:
try:
callback(previous_state, new_state, snapshot)
except Exception as e:
logger.warning(
"SessionManager: callback raised %s: %s",
type(e).__name__,
e,
)
return snapshot
def update(self, **updates: Any) -> VisitorSession:
"""Update session fields without changing state.
For stashing data (e.g. the recognized face name) where no transition
is appropriate yet.
"""
with self._lock:
self._validate_updates(updates)
self._session = replace(
self._session, **updates, last_event_at=time.time()
)
return replace(self._session)
# ------------------------------------------------------------------
# Atomic claim helpers (used by send_email to prevent races where two
# concurrent tool calls both pass a read-then-check dedupe and end up
# delivering the same notification twice).
# ------------------------------------------------------------------
def try_claim_email(self, to: str) -> bool:
"""Atomically claim the right to send an email to ``to``.
Returns True iff the caller wins the race; False if another
concurrent caller already claimed (or completed) a send to the
same recipient within this session. Callers that win MUST call
``release_email_claim(to)`` on failure so a retry can succeed.
The mutation here intentionally does NOT fire the state-change
callback — this is an in-flight reservation, not a transition.
The real NOTIFIED transition (which pushes context to the LLM)
still happens via the controller's on_tool_completed path after
Resend confirms delivery.
"""
norm = (to or "").strip().lower()
if not norm:
return False
with self._lock:
existing = (self._session.email_sent_to or "").strip().lower()
if existing == norm:
return False
self._session = replace(
self._session,
email_sent_to=to,
last_event_at=time.time(),
)
return True
def release_email_claim(self, to: str) -> None:
"""Undo a previous ``try_claim_email`` after a delivery failure."""
norm = (to or "").strip().lower()
if not norm:
return
with self._lock:
existing = (self._session.email_sent_to or "").strip().lower()
if existing == norm:
self._session = replace(
self._session,
email_sent_to=None,
last_event_at=time.time(),
)
def reset(self) -> VisitorSession:
"""Drop the current visitor and return to IDLE.
If a ``visitor_log`` was provided at construction, the pre-reset
snapshot is persisted there before the fields are cleared. This
means BOTH the controller's face-loss path AND the timeout path
produce a log entry — without each call site needing to remember
to capture the snapshot.
"""
# Capture pre-reset state for logging. We do this outside the lock
# (the manager handles its own concurrency on read).
pre_snapshot = self.session
pre_state = pre_snapshot.current_state
if (
self._visitor_log is not None
and pre_state != ReceptionState.IDLE
):
try:
self._visitor_log.record_visit(pre_snapshot, pre_state)
except Exception as e:
logger.warning(
"SessionManager: visitor_log.record_visit raised %s: %s",
type(e).__name__, e,
)
return self.transition(
ReceptionState.IDLE,
visitor_name=None,
recognized_face_name=None,
employee_name=None,
purpose=None,
matched_appointment=None,
email_sent_to=None,
error_message=None,
started_at=time.time(),
last_user_transcript=None,
register_guest_rejections=0,
)
# ------------------------------------------------------------------
# Idle handling
# ------------------------------------------------------------------
def bump_register_guest_rejection(self) -> int:
"""Increment the per-session register_guest rejection counter.
Returns the new count. Used by register_guest to detect when the
visitor's name cannot be captured reliably (after 2+ rejections,
the bot should hand off to an operator rather than looping).
"""
with self._lock:
self._session = replace(
self._session,
register_guest_rejections=self._session.register_guest_rejections + 1,
last_event_at=time.time(),
)
return self._session.register_guest_rejections
def record_user_transcript(self, transcript: str) -> None:
"""Stash the most recent visitor utterance for backend introspection.
Used by register_guest's confirmation guard so we can verify that
the visitor actually said an affirmation before the LLM tries to
save a face under that name. Does NOT fire the state-change
callback — purely a side-channel record.
"""
with self._lock:
self._session = replace(
self._session,
last_user_transcript=transcript,
last_event_at=time.time(),
)
def record_register_guest_attempt(self, name: str) -> None:
"""Stash the name proposed via a rejected register_guest(confirmed=false).
Used by send_email's no-visitor-identity block path so we can tell
the LLM exactly which name to re-call register_guest with, instead
of asking the visitor again. Does NOT fire the state-change callback.
"""
with self._lock:
self._session = replace(
self._session,
last_register_guest_name=(name or "").strip() or None,
last_event_at=time.time(),
)
def touch(self) -> None:
"""Bump ``last_event_at`` without changing any other session field.
Used by the realtime handler when the visitor starts speaking or
a transcript completes — those events prove the session is still
active even though no state transition fires. Without this, a
flow like 'I heard Henry — is that right?' followed by the
visitor thinking for 65s would get auto-reset out from under the
bot, and the visitor's eventual 'yes' would arrive into an empty
IDLE session.
Does NOT fire the state-change callback — this is a pure
liveness ping, not a transition.
"""
with self._lock:
self._session = replace(
self._session, last_event_at=time.time()
)
def is_stale(self, now: Optional[float] = None) -> bool:
"""Return True if a non-IDLE session has had no transition for the timeout.
Does NOT mutate. Callers can decide whether to reset.
"""
ref = now if now is not None else time.time()
with self._lock:
if self._session.current_state == ReceptionState.IDLE:
return False
return (ref - self._session.last_event_at) > self._idle_timeout_seconds
def maybe_reset_if_stale(self) -> bool:
"""Reset the session if it has gone stale. Returns True if reset happened.
Safe to call from any thread on any tick. Designed to be invoked
from the realtime handler's existing periodic ``emit()`` loop so
we don't need another timer thread.
"""
if not self.is_stale():
return False
logger.info(
"Session idle for > %.0fs — auto-resetting",
self._idle_timeout_seconds,
)
self.reset()
return True
@staticmethod
def _validate_updates(updates: dict[str, Any]) -> None:
unknown = set(updates) - set(VisitorSession.__dataclass_fields__)
if unknown:
raise ValueError(f"Unknown session fields: {sorted(unknown)!r}")