""" Agent pending intent state — MongoDB with in-memory fallback. """ import os from datetime import datetime, timezone from typing import Any, Dict, Optional try: from pymongo import MongoClient from pymongo.errors import PyMongoError except ImportError: # pragma: no cover MongoClient = None PyMongoError = Exception MONGO_HOST = os.getenv("MONGO_HOST", "localhost") MONGO_PORT = int(os.getenv("MONGO_PORT", "27017")) MONGO_USER = os.getenv("MONGO_USER", "") MONGO_PASSWORD = os.getenv("MONGO_PASSWORD", "") MONGO_DATABASE = os.getenv("MONGO_DATABASE", "chat_history") PENDING_COLLECTION = os.getenv("AGENT_PENDING_COLLECTION", "agent_pending") PENDING_TTL_SECONDS = int(os.getenv("AGENT_PENDING_TTL_SECONDS", str(30 * 60))) class AgentPendingStore: """Stores clarify/pending context per botId for multi-turn slot filling.""" def __init__(self): self._memory: Dict[str, Dict[str, Any]] = {} self._collection = None self._init_mongo() def _init_mongo(self) -> None: if MongoClient is None: print("[AgentPending] pymongo unavailable — using in-memory store") return try: if MONGO_USER and MONGO_PASSWORD: uri = ( f"mongodb://{MONGO_USER}:{MONGO_PASSWORD}@{MONGO_HOST}:{MONGO_PORT}/" f"{MONGO_DATABASE}?authSource={MONGO_DATABASE}" ) else: uri = f"mongodb://{MONGO_HOST}:{MONGO_PORT}/" client = MongoClient(uri, serverSelectionTimeoutMS=3000) client.admin.command("ping") db = client[MONGO_DATABASE] self._collection = db[PENDING_COLLECTION] self._collection.create_index("bot_id", unique=True, background=True) self._collection.create_index( [("updated_at", 1)], expireAfterSeconds=PENDING_TTL_SECONDS, background=True, name="idx_pending_ttl", ) print(f"[AgentPending] Mongo connected: {MONGO_HOST}:{MONGO_PORT}/{PENDING_COLLECTION}") except PyMongoError as exc: print(f"[AgentPending] Mongo unavailable ({exc}) — using in-memory store") self._collection = None def get(self, bot_id: Optional[str]) -> Optional[Dict[str, Any]]: key = self._normalize_bot_id(bot_id) if self._collection is not None: doc = self._collection.find_one({"bot_id": key}, {"_id": 0}) if doc: return { "pendingIntentType": doc.get("pending_intent_type"), "pendingParams": doc.get("pending_params") or {}, "missingParams": doc.get("missing_params") or [], } return None return self._memory.get(key) def save( self, bot_id: Optional[str], *, pending_intent_type: str, pending_params: Optional[Dict[str, Any]] = None, missing_params: Optional[list] = None, ) -> None: key = self._normalize_bot_id(bot_id) payload = { "pendingIntentType": pending_intent_type, "pendingParams": pending_params or {}, "missingParams": missing_params or [], } if self._collection is not None: now = datetime.now(timezone.utc) self._collection.update_one( {"bot_id": key}, { "$set": { "bot_id": key, "pending_intent_type": pending_intent_type, "pending_params": pending_params or {}, "missing_params": missing_params or [], "updated_at": now, } }, upsert=True, ) return self._memory[key] = payload def clear(self, bot_id: Optional[str]) -> None: key = self._normalize_bot_id(bot_id) if self._collection is not None: self._collection.delete_one({"bot_id": key}) return self._memory.pop(key, None) @staticmethod def _normalize_bot_id(bot_id: Optional[str]) -> str: if bot_id and str(bot_id).strip(): return str(bot_id).strip() return "anonymous"