"""Customer Service Agent: multi-tenant config, knowledge, public chat helpers. Domain layer for website customer agents. Does not call the model backend; the API layer builds prompts and forwards inference. """ from __future__ import annotations import base64 import hashlib import hmac import ipaddress import json import logging import re import secrets import smtplib import socket import sqlite3 import urllib.error import urllib.request from dataclasses import dataclass from datetime import UTC, datetime, timedelta from email.message import EmailMessage from html import unescape from typing import Any from urllib.parse import urlparse LOGGER = logging.getLogger(__name__) SITE_KEY_PREFIX = "pk_live_" SITE_KEY_SECRET_BYTES = 24 MAX_AGENT_NAME_CHARS = 120 MAX_AGENT_GREETING_CHARS = 800 MAX_AGENT_VOICE_CHARS = 4_000 MAX_AGENT_ESCALATE_EMAIL_CHARS = 254 MAX_AGENT_ESCALATE_URL_CHARS = 500 MAX_AGENT_WEBHOOK_URL_CHARS = 500 MAX_BRAND_COLOR_CHARS = 32 MAX_BRAND_LOGO_URL_CHARS = 500 MAX_LAUNCHER_LABEL_CHARS = 40 MAX_ALLOWED_ORIGINS = 20 MAX_ORIGIN_CHARS = 200 HEX_COLOR_RE = re.compile(r"^#(?:[0-9a-fA-F]{3}|[0-9a-fA-F]{6})$") MAX_KNOWLEDGE_TITLE_CHARS = 200 MAX_KNOWLEDGE_BODY_CHARS = 80_000 MAX_KNOWLEDGE_DOCS_PER_AGENT = 40 MAX_KNOWLEDGE_FILE_BYTES = 400_000 MAX_CHUNK_CHARS = 900 MAX_RETRIEVED_CHUNKS = 6 MAX_PUBLIC_QUESTION_CHARS = 2_000 MAX_PUBLIC_HISTORY_MESSAGES = 8 MAX_PUBLIC_HISTORY_ITEM_CHARS = 1_200 MAX_VISITOR_NAME_CHARS = 120 MAX_VISITOR_EMAIL_CHARS = 254 MAX_ESCALATE_MESSAGE_CHARS = 4_000 MAX_CONVERSATION_ID_CHARS = 64 DEFAULT_ESCALATION_RETENTION_DAYS = 90 MAX_URL_IMPORT_BYTES = 500_000 MAX_URL_IMPORT_CHARS = 80_000 URL_IMPORT_TIMEOUT_SECONDS = 12 MIN_RETRIEVAL_HIT_SCORE = 0.35 TOKEN_RE = re.compile(r"[a-z0-9']+", re.IGNORECASE) HTML_TAG_RE = re.compile(r"<[^>]+>") SCRIPT_STYLE_RE = re.compile(r"(?is)<(script|style|noscript)[^>]*>.*?") HARD_REFUSE_MARKERS = ( "diagnose my", "what medication", "prescribe", "am i going to die", "file a lawsuit", "draft a will for me", "how do i hack", "bypass payment", "steal card", ) SENSITIVE_MARKERS = ( "abuse", "assault", "harassment", "suicide", "self-harm", "legal action", "lawyer", "attorney", "refund dispute", "chargeback", "speak to a human", "speak to a person", "talk to a human", "talk to a person", "real person", "manager", "complaint", ) def utc_now() -> str: return datetime.now(UTC).replace(microsecond=0).isoformat().replace("+00:00", "Z") def hash_site_key(site_key: str) -> str: return hashlib.sha256(f"synderesis-site-key:{site_key}".encode("utf-8")).hexdigest() def generate_site_key() -> str: return SITE_KEY_PREFIX + secrets.token_urlsafe(SITE_KEY_SECRET_BYTES) def uuid_str() -> str: return secrets.token_hex(16) def normalize_origin(value: str) -> str: raw = value.strip() if not raw: raise ValueError("origin must be non-empty") if len(raw) > MAX_ORIGIN_CHARS: raise ValueError("origin is too long") if raw == "*": return "*" parsed = urlparse(raw if "://" in raw else f"https://{raw}") if parsed.scheme not in {"http", "https"}: raise ValueError("origin must be http or https") if not parsed.netloc: raise ValueError("origin must include a host") return f"{parsed.scheme}://{parsed.netloc}".lower() def parse_allowed_origins(raw: str | list[str] | None) -> list[str]: if raw is None: return [] if isinstance(raw, str): items = [part.strip() for part in raw.replace("\n", ",").split(",") if part.strip()] else: items = [str(part).strip() for part in raw if str(part).strip()] if len(items) > MAX_ALLOWED_ORIGINS: raise ValueError(f"at most {MAX_ALLOWED_ORIGINS} allowed origins") return [normalize_origin(item) for item in items] def origins_json(origins: list[str]) -> str: return json.dumps(origins, ensure_ascii=False) def load_origins(raw: str | None) -> list[str]: if not raw: return [] try: data = json.loads(raw) except json.JSONDecodeError: return [] if not isinstance(data, list): return [] return [str(item) for item in data if isinstance(item, str)] def origin_allowed(allowed: list[str], request_origin: str | None) -> bool: if not allowed: return request_origin in (None, "") if "*" in allowed: return True if not request_origin: return False try: normalized = normalize_origin(request_origin) except ValueError: return False return normalized in {item.lower() for item in allowed if item != "*"} def chunk_text(text: str, max_chars: int = MAX_CHUNK_CHARS) -> list[str]: cleaned = re.sub(r"\r\n?", "\n", text).strip() if not cleaned: return [] paragraphs = [part.strip() for part in re.split(r"\n\s*\n", cleaned) if part.strip()] chunks: list[str] = [] for paragraph in paragraphs: if len(paragraph) <= max_chars: chunks.append(paragraph) continue start = 0 while start < len(paragraph): end = min(len(paragraph), start + max_chars) if end < len(paragraph): split_at = paragraph.rfind(" ", start, end) if split_at > start + max_chars // 2: end = split_at piece = paragraph[start:end].strip() if piece: chunks.append(piece) start = end return chunks def tokenize(text: str) -> list[str]: return [token.casefold() for token in TOKEN_RE.findall(text) if len(token) > 1] def score_chunk(query_tokens: list[str], chunk: str, *, title: str = "") -> float: if not query_tokens: return 0.0 blob = f"{title}\n{chunk}".strip() chunk_tokens = tokenize(blob) if not chunk_tokens: return 0.0 chunk_set = set(chunk_tokens) unique_query = list(dict.fromkeys(query_tokens)) hits = sum(1 for token in unique_query if token in chunk_set) if hits == 0: return 0.0 density = hits / max(len(chunk_tokens), 1) coverage = hits / max(len(unique_query), 1) # Light phrase bonus: consecutive query tokens found as a span in text. phrase_bonus = 0.0 lowered = blob.casefold() if len(unique_query) >= 2: for i in range(len(unique_query) - 1): pair = f"{unique_query[i]} {unique_query[i + 1]}" if pair in lowered: phrase_bonus += 0.25 title_tokens = set(tokenize(title)) title_hits = sum(1 for token in unique_query if token in title_tokens) title_bonus = 0.35 * (title_hits / max(len(unique_query), 1)) return coverage * 2.0 + density + phrase_bonus + title_bonus def retrieve_chunks( chunks: list[dict[str, Any]], query: str, limit: int = MAX_RETRIEVED_CHUNKS, ) -> list[dict[str, Any]]: query_tokens = tokenize(query) scored: list[tuple[float, dict[str, Any]]] = [] for chunk in chunks: text = str(chunk.get("text") or "") title = str(chunk.get("title") or "") score = score_chunk(query_tokens, text, title=title) if score > 0: enriched = {**chunk, "score": score} scored.append((score, enriched)) scored.sort(key=lambda item: item[0], reverse=True) if scored: return [item[1] for item in scored[:limit]] # Weak fallback only when something exists; mark low score so knowledge_hit stays false. fallback = [] for chunk in chunks[: min(limit, len(chunks))]: fallback.append({**chunk, "score": 0.0}) return fallback def knowledge_is_hit(retrieved: list[dict[str, Any]]) -> bool: if not retrieved: return False top = float(retrieved[0].get("score") or 0.0) return top >= MIN_RETRIEVAL_HIT_SCORE def looks_sensitive(question: str) -> bool: lowered = question.casefold() return any(marker in lowered for marker in SENSITIVE_MARKERS) def looks_hard_refuse(question: str) -> bool: lowered = question.casefold() return any(marker in lowered for marker in HARD_REFUSE_MARKERS) def build_customer_agent_system_prompt( *, business_name: str, voice_instructions: str, escalate_email: str, knowledge_blocks: list[str], ) -> str: knowledge = "\n\n".join(f"[Source {idx + 1}]\n{block}" for idx, block in enumerate(knowledge_blocks)) if not knowledge.strip(): knowledge = ( "(No business knowledge documents were uploaded yet. " "Say you do not have that information and offer escalation.)" ) voice = voice_instructions.strip() or ( "Warm, concise, professional, and consistent with a Catholic-aligned organisation. " "Do not preach. Do not invent doctrine or policy." ) contact = escalate_email.strip() or "the team" return ( f"You are the customer service agent for {business_name}.\n" f"Voice and style: {voice}\n\n" "Rules:\n" "1. Answer only using the Approved business knowledge below and ordinary conversational courtesy.\n" "2. If the knowledge does not contain the answer, say you do not know and offer to connect the customer to a human.\n" "3. Never invent prices, shipping times, stock, legal advice, medical advice, or policies.\n" "4. Refuse off-brand, deceptive, or dignity-violating requests; offer a polite alternative.\n" "5. Keep replies short and friendly (usually under 120 words).\n" f"6. When the issue is sensitive, angry, unresolved, or needs a person, tell the customer you will escalate to {contact}.\n" "7. Do not claim to be an official organ of the Catholic Church.\n" "8. Respond in the language the customer writes in when you can.\n" "9. For medical diagnosis, legal representation, or criminal requests, refuse briefly and offer human handoff.\n\n" "Approved business knowledge:\n" f"{knowledge}\n\n" "Return plain text for the customer. Do not wrap the answer in JSON unless the customer asked for JSON." ) def should_recommend_escalation(question: str, answer: str, knowledge_hit: bool) -> bool: if looks_hard_refuse(question) or looks_sensitive(question): return True lowered = answer.casefold() uncertainty = ( "i don't know" in lowered or "i do not know" in lowered or "don't have that information" in lowered or "do not have that information" in lowered or "not in our" in lowered or "pass you to" in lowered or "connect you with" in lowered or "leave a message" in lowered or "escalate" in lowered ) if uncertainty: return True if not knowledge_hit and len(question.strip()) > 12: return True return False def hard_refuse_answer(escalate_email: str = "") -> str: contact = escalate_email.strip() or "our team" return ( "I cannot help with that request. For anything medical, legal, or sensitive, " f"please leave a message for {contact} and a person will follow up." ) def demo_answer_from_knowledge( *, question: str, knowledge_blocks: list[str], knowledge_hit: bool, escalate_email: str = "", business_name: str = "our shop", ) -> tuple[str, bool]: """Offline/demo reply from retrieved FAQ text (no LLM). Returns (answer, escalate).""" contact = escalate_email.strip() or "our team" if looks_hard_refuse(question): return hard_refuse_answer(contact), True if looks_sensitive(question) and not knowledge_hit: return ( f"That sounds like something a person should handle. " f"Please leave a message for {contact} and the {business_name} team will follow up.", True, ) if knowledge_hit and knowledge_blocks: # Keep it short: top block(s) as the grounded answer. primary = knowledge_blocks[0] # Strip leading "Title: " style if present if ":" in primary[:80]: _, _, rest = primary.partition(":") if rest.strip(): primary = rest.strip() answer = ( f"{primary}\n\n" f"If you need something else, just ask — or leave a message for {contact}." ) return answer[:1200], False return ( f"I don't have that information in our published FAQ. " f"Please leave a message for {contact} and someone from {business_name} will help you.", True, ) def html_to_text(html: str) -> str: cleaned = SCRIPT_STYLE_RE.sub(" ", html or "") cleaned = HTML_TAG_RE.sub(" ", cleaned) cleaned = unescape(cleaned) cleaned = re.sub(r"[ \t]+", " ", cleaned) cleaned = re.sub(r"\n{3,}", "\n\n", cleaned) return cleaned.strip() def _is_blocked_ip(ip: str) -> bool: """Return True if IP must not be contacted (SSRF guard, IPv4 + IPv6).""" try: addr = ipaddress.ip_address(ip) except ValueError: return True # Unwrap IPv4-mapped IPv6 (:ffff:x.x.x.x) if isinstance(addr, ipaddress.IPv6Address) and addr.ipv4_mapped is not None: addr = addr.ipv4_mapped if bool( addr.is_private or addr.is_loopback or addr.is_link_local or addr.is_reserved or addr.is_multicast or addr.is_unspecified ): return True # Carrier-grade NAT / shared address space if isinstance(addr, ipaddress.IPv4Address) and addr in ipaddress.ip_network("100.64.0.0/10"): return True # IPv6 unique local (fc00::/7) — covered by is_private on modern Python; keep explicit if isinstance(addr, ipaddress.IPv6Address): if addr in ipaddress.ip_network("fc00::/7"): return True # Site-local deprecated fec0::/10 if addr in ipaddress.ip_network("fec0::/10"): return True return False def validate_public_https_url(url: str, *, purpose: str = "url") -> str: """Validate an https URL that will be fetched server-side (import or webhook).""" raw = (url or "").strip() if not raw: raise ValueError(f"{purpose} is required") parsed = urlparse(raw) if parsed.scheme != "https": raise ValueError(f"only https {purpose}s are allowed") if not parsed.hostname: raise ValueError(f"{purpose} must include a host") if parsed.username or parsed.password: raise ValueError(f"{purpose} must not include credentials") host = parsed.hostname if host in {"localhost"} or host.endswith(".local") or host.endswith(".internal"): raise ValueError("local hosts are not allowed") # Literal IP in hostname try: literal = ipaddress.ip_address(host.strip("[]")) if _is_blocked_ip(str(literal)): raise ValueError(f"{purpose} points to a private or blocked address") except ValueError as exc: if "private or blocked" in str(exc): raise # not a literal IP — resolve try: infos = socket.getaddrinfo(host, 443, type=socket.SOCK_STREAM) except socket.gaierror as gai_exc: raise ValueError(f"could not resolve host: {host}") from gai_exc if not infos: raise ValueError(f"could not resolve host: {host}") for info in infos: ip = info[4][0] if _is_blocked_ip(ip): raise ValueError(f"{purpose} resolves to a private or blocked address") return raw def fetch_url_text(url: str) -> tuple[str, str]: """Fetch a public HTTPS page and return (title, plain_text).""" safe_url = validate_public_https_url(url) req = urllib.request.Request( safe_url, headers={ "User-Agent": "SynderesisCustomerAgent/1.3 (+knowledge-import)", "Accept": "text/html,application/xhtml+xml,text/plain;q=0.9,*/*;q=0.1", }, method="GET", ) try: with urllib.request.urlopen(req, timeout=URL_IMPORT_TIMEOUT_SECONDS) as resp: content_type = (resp.headers.get("Content-Type") or "").split(";")[0].strip().lower() data = resp.read(MAX_URL_IMPORT_BYTES + 1) except Exception as exc: # noqa: BLE001 raise ValueError(f"could not fetch url: {exc}") from exc if len(data) > MAX_URL_IMPORT_BYTES: raise ValueError("page is too large to import") raw = data.decode("utf-8", errors="replace") title = urlparse(safe_url).path.rsplit("/", 1)[-1] or urlparse(safe_url).hostname or "Imported page" if "html" in content_type or "]*>(.*?)", raw) if title_match: title = html_to_text(title_match.group(1))[:MAX_KNOWLEDGE_TITLE_CHARS] or title body = html_to_text(raw) else: body = raw.strip() body = body[:MAX_URL_IMPORT_CHARS].strip() if len(body) < 40: raise ValueError("imported page had too little text") return title[:MAX_KNOWLEDGE_TITLE_CHARS], body def ensure_customer_agent_schema(connection: sqlite3.Connection) -> None: """Create customer-agent tables if missing (idempotent).""" connection.executescript( """ CREATE TABLE IF NOT EXISTS customer_agents ( id TEXT PRIMARY KEY, customer_id TEXT NOT NULL, name TEXT NOT NULL, greeting TEXT NOT NULL DEFAULT '', voice_instructions TEXT NOT NULL DEFAULT '', escalate_email TEXT NOT NULL DEFAULT '', escalate_url TEXT NOT NULL DEFAULT '', escalate_webhook_url TEXT NOT NULL DEFAULT '', allowed_origins_json TEXT NOT NULL DEFAULT '[]', site_key_prefix TEXT NOT NULL, site_key_hash TEXT NOT NULL UNIQUE, status TEXT NOT NULL DEFAULT 'active', created_at TEXT NOT NULL, updated_at TEXT NOT NULL, FOREIGN KEY (customer_id) REFERENCES registered_users(customer_id) ON DELETE CASCADE ); CREATE INDEX IF NOT EXISTS idx_customer_agents_owner ON customer_agents(customer_id, updated_at); CREATE INDEX IF NOT EXISTS idx_customer_agents_prefix ON customer_agents(site_key_prefix); CREATE TABLE IF NOT EXISTS customer_agent_docs ( id TEXT PRIMARY KEY, agent_id TEXT NOT NULL, title TEXT NOT NULL, body TEXT NOT NULL, source_type TEXT NOT NULL DEFAULT 'manual', created_at TEXT NOT NULL, updated_at TEXT NOT NULL, FOREIGN KEY (agent_id) REFERENCES customer_agents(id) ON DELETE CASCADE ); CREATE INDEX IF NOT EXISTS idx_customer_agent_docs_agent ON customer_agent_docs(agent_id, updated_at); CREATE TABLE IF NOT EXISTS customer_agent_chunks ( id INTEGER PRIMARY KEY AUTOINCREMENT, agent_id TEXT NOT NULL, doc_id TEXT NOT NULL, chunk_index INTEGER NOT NULL, text TEXT NOT NULL, FOREIGN KEY (agent_id) REFERENCES customer_agents(id) ON DELETE CASCADE, FOREIGN KEY (doc_id) REFERENCES customer_agent_docs(id) ON DELETE CASCADE ); CREATE INDEX IF NOT EXISTS idx_customer_agent_chunks_agent ON customer_agent_chunks(agent_id); CREATE TABLE IF NOT EXISTS customer_agent_escalations ( id TEXT PRIMARY KEY, agent_id TEXT NOT NULL, conversation_id TEXT NOT NULL DEFAULT '', visitor_name TEXT NOT NULL DEFAULT '', visitor_email TEXT NOT NULL DEFAULT '', message TEXT NOT NULL DEFAULT '', transcript_json TEXT NOT NULL DEFAULT '[]', reason TEXT NOT NULL DEFAULT '', status TEXT NOT NULL DEFAULT 'open', email_status TEXT NOT NULL DEFAULT '', webhook_status TEXT NOT NULL DEFAULT '', delivery_error TEXT NOT NULL DEFAULT '', created_at TEXT NOT NULL, FOREIGN KEY (agent_id) REFERENCES customer_agents(id) ON DELETE CASCADE ); CREATE INDEX IF NOT EXISTS idx_customer_agent_escalations_agent ON customer_agent_escalations(agent_id, created_at); CREATE TABLE IF NOT EXISTS customer_agent_chat_events ( id INTEGER PRIMARY KEY AUTOINCREMENT, agent_id TEXT NOT NULL, created_at TEXT NOT NULL, question TEXT NOT NULL DEFAULT '', answer_preview TEXT NOT NULL DEFAULT '', knowledge_hit INTEGER NOT NULL DEFAULT 0, escalate_recommended INTEGER NOT NULL DEFAULT 0, hard_refuse INTEGER NOT NULL DEFAULT 0, retrieved_count INTEGER NOT NULL DEFAULT 0, top_score REAL NOT NULL DEFAULT 0, FOREIGN KEY (agent_id) REFERENCES customer_agents(id) ON DELETE CASCADE ); CREATE INDEX IF NOT EXISTS idx_customer_agent_chat_events_agent ON customer_agent_chat_events(agent_id, created_at); """ ) agent_cols = {row["name"] for row in connection.execute("PRAGMA table_info(customer_agents)").fetchall()} if "escalate_webhook_url" not in agent_cols: connection.execute( "ALTER TABLE customer_agents ADD COLUMN escalate_webhook_url TEXT NOT NULL DEFAULT ''" ) for col, default in ( ("brand_primary_color", ""), ("brand_position", "right"), ("brand_logo_url", ""), ("launcher_label", ""), ): if col not in agent_cols: connection.execute( f"ALTER TABLE customer_agents ADD COLUMN {col} TEXT NOT NULL DEFAULT '{default}'" ) esc_cols = {row["name"] for row in connection.execute("PRAGMA table_info(customer_agent_escalations)").fetchall()} if "email_status" not in esc_cols: connection.execute( "ALTER TABLE customer_agent_escalations ADD COLUMN email_status TEXT NOT NULL DEFAULT ''" ) if "webhook_status" not in esc_cols: connection.execute( "ALTER TABLE customer_agent_escalations ADD COLUMN webhook_status TEXT NOT NULL DEFAULT ''" ) if "delivery_error" not in esc_cols: connection.execute( "ALTER TABLE customer_agent_escalations ADD COLUMN delivery_error TEXT NOT NULL DEFAULT ''" ) def normalize_brand_color(value: str) -> str: raw = (value or "").strip() if not raw: return "" if not HEX_COLOR_RE.match(raw): raise ValueError("brand_primary_color must be a hex color like #2e4a86") return raw.lower() def normalize_brand_position(value: str) -> str: raw = (value or "right").strip().lower() if raw not in {"left", "right"}: raise ValueError("brand_position must be left or right") return raw def normalize_brand_logo_url(value: str) -> str: raw = (value or "").strip() if not raw: return "" if not (raw.startswith("https://") or raw.startswith("http://")): raise ValueError("brand_logo_url must be an http(s) URL") if len(raw) > MAX_BRAND_LOGO_URL_CHARS: raise ValueError("brand_logo_url is too long") return raw @dataclass(frozen=True) class CustomerAgentRecord: id: str customer_id: str name: str greeting: str voice_instructions: str escalate_email: str escalate_url: str escalate_webhook_url: str allowed_origins: list[str] site_key_prefix: str status: str created_at: str updated_at: str brand_primary_color: str = "" brand_position: str = "right" brand_logo_url: str = "" launcher_label: str = "" def public_dict(self) -> dict[str, Any]: return { "object": "synderesis.customer_agent.public", "id": self.id, "name": self.name, "greeting": self.greeting or default_greeting(self.name), "escalate_email": self.escalate_email, "escalate_url": self.escalate_url, "brand": { "primary_color": self.brand_primary_color or "#2e4a86", "position": self.brand_position or "right", "logo_url": self.brand_logo_url, "launcher_label": self.launcher_label, }, "powered_by": "Synderesis", } def owner_dict( self, *, include_site_key: str | None = None, public_base_url: str = "https://www.synderesis.eu", knowledge_count: int | None = None, ) -> dict[str, Any]: payload: dict[str, Any] = { "object": "synderesis.customer_agent", "id": self.id, "name": self.name, "greeting": self.greeting, "voice_instructions": self.voice_instructions, "escalate_email": self.escalate_email, "escalate_url": self.escalate_url, "escalate_webhook_url": self.escalate_webhook_url, "allowed_origins": self.allowed_origins, "site_key_prefix": self.site_key_prefix, "status": self.status, "created_at": self.created_at, "updated_at": self.updated_at, "brand_primary_color": self.brand_primary_color, "brand_position": self.brand_position or "right", "brand_logo_url": self.brand_logo_url, "launcher_label": self.launcher_label, "embed_snippet": embed_snippet_placeholder(self.site_key_prefix, public_base_url), "readiness": readiness_checklist(self, knowledge_count=knowledge_count), } if include_site_key: payload["site_key"] = include_site_key payload["embed_snippet"] = embed_snippet(include_site_key, public_base_url) payload["site_key_shown_once"] = True return payload def readiness_checklist(agent: CustomerAgentRecord, *, knowledge_count: int | None = None) -> dict[str, Any]: """Go-live checklist for admin UI.""" has_name = bool(agent.name.strip()) has_handoff = bool(agent.escalate_email.strip() or agent.escalate_webhook_url.strip()) origins = agent.allowed_origins has_origins = bool(origins) origins_locked = has_origins and origins != ["*"] and "*" not in origins has_knowledge = knowledge_count is None or knowledge_count > 0 checks = [ {"id": "name", "label": "Business name set", "ok": has_name}, {"id": "handoff", "label": "Escalation email or webhook configured", "ok": has_handoff}, {"id": "knowledge", "label": "FAQ / knowledge uploaded", "ok": bool(has_knowledge)}, {"id": "origins", "label": "Allowed websites set", "ok": has_origins}, { "id": "origins_locked", "label": "Production origins locked (not only *)", "ok": origins_locked, "optional": True, }, ] required_ok = all(item["ok"] for item in checks if not item.get("optional")) return { "ready_for_test": required_ok or (has_name and has_knowledge and has_origins), "ready_for_production": required_ok and origins_locked and has_handoff and has_knowledge, "checks": checks, } def default_greeting(name: str) -> str: label = name.strip() or "us" return ( f"Hello — welcome to {label}. Ask me about products, hours, shipping, or policies. " "I will hand sensitive questions to a person." ) def embed_snippet(site_key: str, base_url: str = "https://www.synderesis.eu") -> str: base = base_url.rstrip("/") return ( f'' ) def embed_snippet_placeholder(site_key_prefix: str, base_url: str = "https://www.synderesis.eu") -> str: base = base_url.rstrip("/") return ( f'' ) def _row_get(row: sqlite3.Row, key: str, default: str = "") -> str: keys = row.keys() if key not in keys: return default return str(row[key] or default) def row_to_agent(row: sqlite3.Row) -> CustomerAgentRecord: return CustomerAgentRecord( id=str(row["id"]), customer_id=str(row["customer_id"]), name=str(row["name"]), greeting=str(row["greeting"] or ""), voice_instructions=str(row["voice_instructions"] or ""), escalate_email=str(row["escalate_email"] or ""), escalate_url=str(row["escalate_url"] or ""), escalate_webhook_url=_row_get(row, "escalate_webhook_url"), allowed_origins=load_origins(str(row["allowed_origins_json"] or "[]")), site_key_prefix=str(row["site_key_prefix"]), status=str(row["status"]), created_at=str(row["created_at"]), updated_at=str(row["updated_at"]), brand_primary_color=_row_get(row, "brand_primary_color"), brand_position=_row_get(row, "brand_position", "right") or "right", brand_logo_url=_row_get(row, "brand_logo_url"), launcher_label=_row_get(row, "launcher_label"), ) def create_agent( connection: sqlite3.Connection, *, customer_id: str, name: str, greeting: str = "", voice_instructions: str = "", escalate_email: str = "", escalate_url: str = "", escalate_webhook_url: str = "", allowed_origins: list[str] | None = None, brand_primary_color: str = "", brand_position: str = "right", brand_logo_url: str = "", launcher_label: str = "", ) -> tuple[CustomerAgentRecord, str]: agent_id = uuid_str() site_key = generate_site_key() now = utc_now() origins = allowed_origins or [] color = normalize_brand_color(brand_primary_color) position = normalize_brand_position(brand_position) logo = normalize_brand_logo_url(brand_logo_url) label = launcher_label.strip()[:MAX_LAUNCHER_LABEL_CHARS] webhook = escalate_webhook_url.strip() if webhook: webhook = validate_public_https_url(webhook, purpose="webhook") connection.execute( """ INSERT INTO customer_agents ( id, customer_id, name, greeting, voice_instructions, escalate_email, escalate_url, escalate_webhook_url, allowed_origins_json, site_key_prefix, site_key_hash, status, brand_primary_color, brand_position, brand_logo_url, launcher_label, created_at, updated_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'active', ?, ?, ?, ?, ?, ?) """, ( agent_id, customer_id, name.strip(), greeting.strip(), voice_instructions.strip(), escalate_email.strip(), escalate_url.strip(), webhook, origins_json(origins), site_key[:16], hash_site_key(site_key), color, position, logo, label, now, now, ), ) row = connection.execute("SELECT * FROM customer_agents WHERE id = ?", (agent_id,)).fetchone() assert row is not None return row_to_agent(row), site_key def list_agents_for_customer(connection: sqlite3.Connection, customer_id: str) -> list[CustomerAgentRecord]: rows = connection.execute( """ SELECT * FROM customer_agents WHERE customer_id = ? ORDER BY updated_at DESC """, (customer_id,), ).fetchall() return [row_to_agent(row) for row in rows] def get_agent_for_customer( connection: sqlite3.Connection, customer_id: str, agent_id: str, ) -> CustomerAgentRecord | None: row = connection.execute( "SELECT * FROM customer_agents WHERE id = ? AND customer_id = ?", (agent_id, customer_id), ).fetchone() return row_to_agent(row) if row else None def get_agent_by_site_key(connection: sqlite3.Connection, site_key: str) -> CustomerAgentRecord | None: key = site_key.strip() if not key.startswith(SITE_KEY_PREFIX): return None row = connection.execute( "SELECT * FROM customer_agents WHERE site_key_hash = ? AND status = 'active'", (hash_site_key(key),), ).fetchone() return row_to_agent(row) if row else None def update_agent( connection: sqlite3.Connection, *, customer_id: str, agent_id: str, name: str | None = None, greeting: str | None = None, voice_instructions: str | None = None, escalate_email: str | None = None, escalate_url: str | None = None, escalate_webhook_url: str | None = None, allowed_origins: list[str] | None = None, status: str | None = None, brand_primary_color: str | None = None, brand_position: str | None = None, brand_logo_url: str | None = None, launcher_label: str | None = None, ) -> CustomerAgentRecord | None: agent = get_agent_for_customer(connection, customer_id, agent_id) if agent is None: return None next_name = agent.name if name is None else name.strip() next_greeting = agent.greeting if greeting is None else greeting.strip() next_voice = agent.voice_instructions if voice_instructions is None else voice_instructions.strip() next_email = agent.escalate_email if escalate_email is None else escalate_email.strip() next_url = agent.escalate_url if escalate_url is None else escalate_url.strip() if escalate_webhook_url is None: next_webhook = agent.escalate_webhook_url else: raw_hook = escalate_webhook_url.strip() next_webhook = validate_public_https_url(raw_hook, purpose="webhook") if raw_hook else "" next_origins = agent.allowed_origins if allowed_origins is None else allowed_origins next_status = agent.status if status is None else status.strip() next_color = agent.brand_primary_color if brand_primary_color is None else normalize_brand_color(brand_primary_color) next_position = ( agent.brand_position if brand_position is None else normalize_brand_position(brand_position) ) next_logo = agent.brand_logo_url if brand_logo_url is None else normalize_brand_logo_url(brand_logo_url) next_label = ( agent.launcher_label if launcher_label is None else launcher_label.strip()[:MAX_LAUNCHER_LABEL_CHARS] ) if next_status not in {"active", "disabled"}: raise ValueError("status must be active or disabled") now = utc_now() connection.execute( """ UPDATE customer_agents SET name = ?, greeting = ?, voice_instructions = ?, escalate_email = ?, escalate_url = ?, escalate_webhook_url = ?, allowed_origins_json = ?, status = ?, brand_primary_color = ?, brand_position = ?, brand_logo_url = ?, launcher_label = ?, updated_at = ? WHERE id = ? AND customer_id = ? """, ( next_name, next_greeting, next_voice, next_email, next_url, next_webhook, origins_json(next_origins), next_status, next_color, next_position, next_logo, next_label, now, agent_id, customer_id, ), ) return get_agent_for_customer(connection, customer_id, agent_id) def rotate_site_key( connection: sqlite3.Connection, *, customer_id: str, agent_id: str, ) -> tuple[CustomerAgentRecord, str] | None: agent = get_agent_for_customer(connection, customer_id, agent_id) if agent is None: return None site_key = generate_site_key() now = utc_now() connection.execute( """ UPDATE customer_agents SET site_key_prefix = ?, site_key_hash = ?, updated_at = ? WHERE id = ? AND customer_id = ? """, (site_key[:16], hash_site_key(site_key), now, agent_id, customer_id), ) updated = get_agent_for_customer(connection, customer_id, agent_id) assert updated is not None return updated, site_key def delete_agent(connection: sqlite3.Connection, *, customer_id: str, agent_id: str) -> bool: cursor = connection.execute( "DELETE FROM customer_agents WHERE id = ? AND customer_id = ?", (agent_id, customer_id), ) return cursor.rowcount > 0 def count_docs(connection: sqlite3.Connection, agent_id: str) -> int: row = connection.execute( "SELECT COUNT(*) AS n FROM customer_agent_docs WHERE agent_id = ?", (agent_id,), ).fetchone() return int(row["n"] if row else 0) def reindex_doc(connection: sqlite3.Connection, agent_id: str, doc_id: str, body: str) -> None: connection.execute("DELETE FROM customer_agent_chunks WHERE doc_id = ?", (doc_id,)) for index, chunk in enumerate(chunk_text(body)): connection.execute( """ INSERT INTO customer_agent_chunks (agent_id, doc_id, chunk_index, text) VALUES (?, ?, ?, ?) """, (agent_id, doc_id, index, chunk), ) def add_knowledge_doc( connection: sqlite3.Connection, *, agent_id: str, title: str, body: str, source_type: str = "manual", ) -> dict[str, Any]: if count_docs(connection, agent_id) >= MAX_KNOWLEDGE_DOCS_PER_AGENT: raise ValueError(f"at most {MAX_KNOWLEDGE_DOCS_PER_AGENT} knowledge documents per agent") doc_id = uuid_str() now = utc_now() clean_title = title.strip()[:MAX_KNOWLEDGE_TITLE_CHARS] or "Untitled" clean_body = body.strip() if not clean_body: raise ValueError("knowledge body must be non-empty") if len(clean_body) > MAX_KNOWLEDGE_BODY_CHARS: raise ValueError(f"knowledge body exceeds {MAX_KNOWLEDGE_BODY_CHARS} characters") connection.execute( """ INSERT INTO customer_agent_docs (id, agent_id, title, body, source_type, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?) """, (doc_id, agent_id, clean_title, clean_body, source_type.strip() or "manual", now, now), ) reindex_doc(connection, agent_id, doc_id, clean_body) connection.execute( "UPDATE customer_agents SET updated_at = ? WHERE id = ?", (now, agent_id), ) return { "object": "synderesis.customer_agent.doc", "id": doc_id, "agent_id": agent_id, "title": clean_title, "source_type": source_type.strip() or "manual", "chars": len(clean_body), "created_at": now, "updated_at": now, } def list_knowledge_docs(connection: sqlite3.Connection, agent_id: str) -> list[dict[str, Any]]: rows = connection.execute( """ SELECT id, agent_id, title, source_type, created_at, updated_at, LENGTH(body) AS chars FROM customer_agent_docs WHERE agent_id = ? ORDER BY updated_at DESC """, (agent_id,), ).fetchall() return [ { "object": "synderesis.customer_agent.doc", "id": str(row["id"]), "agent_id": str(row["agent_id"]), "title": str(row["title"]), "source_type": str(row["source_type"]), "chars": int(row["chars"] or 0), "created_at": str(row["created_at"]), "updated_at": str(row["updated_at"]), } for row in rows ] def delete_knowledge_doc(connection: sqlite3.Connection, *, agent_id: str, doc_id: str) -> bool: cursor = connection.execute( "DELETE FROM customer_agent_docs WHERE id = ? AND agent_id = ?", (doc_id, agent_id), ) if cursor.rowcount: connection.execute( "UPDATE customer_agents SET updated_at = ? WHERE id = ?", (utc_now(), agent_id), ) return cursor.rowcount > 0 def load_agent_chunks(connection: sqlite3.Connection, agent_id: str) -> list[dict[str, Any]]: rows = connection.execute( """ SELECT c.id, c.doc_id, c.chunk_index, c.text, d.title FROM customer_agent_chunks c JOIN customer_agent_docs d ON d.id = c.doc_id WHERE c.agent_id = ? ORDER BY d.updated_at DESC, c.chunk_index ASC """, (agent_id,), ).fetchall() return [ { "id": int(row["id"]), "doc_id": str(row["doc_id"]), "title": str(row["title"]), "chunk_index": int(row["chunk_index"]), "text": str(row["text"]), } for row in rows ] def create_escalation( connection: sqlite3.Connection, *, agent_id: str, conversation_id: str = "", visitor_name: str = "", visitor_email: str = "", message: str = "", transcript: list[dict[str, str]] | None = None, reason: str = "", ) -> dict[str, Any]: escalation_id = uuid_str() now = utc_now() connection.execute( """ INSERT INTO customer_agent_escalations ( id, agent_id, conversation_id, visitor_name, visitor_email, message, transcript_json, reason, status, email_status, webhook_status, delivery_error, created_at ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, 'open', '', '', '', ?) """, ( escalation_id, agent_id, conversation_id.strip()[:MAX_CONVERSATION_ID_CHARS], visitor_name.strip()[:MAX_VISITOR_NAME_CHARS], visitor_email.strip()[:MAX_VISITOR_EMAIL_CHARS], message.strip()[:MAX_ESCALATE_MESSAGE_CHARS], json.dumps(transcript or [], ensure_ascii=False), reason.strip()[:200], now, ), ) return { "object": "synderesis.customer_agent.escalation", "id": escalation_id, "agent_id": agent_id, "status": "open", "created_at": now, "email_status": "", "webhook_status": "", } def update_escalation_delivery( connection: sqlite3.Connection, escalation_id: str, *, email_status: str = "", webhook_status: str = "", delivery_error: str = "", ) -> None: connection.execute( """ UPDATE customer_agent_escalations SET email_status = CASE WHEN ? != '' THEN ? ELSE email_status END, webhook_status = CASE WHEN ? != '' THEN ? ELSE webhook_status END, delivery_error = ? WHERE id = ? """, (email_status, email_status, webhook_status, webhook_status, delivery_error[:2000], escalation_id), ) def list_escalations(connection: sqlite3.Connection, agent_id: str, limit: int = 50) -> list[dict[str, Any]]: rows = connection.execute( """ SELECT id, agent_id, conversation_id, visitor_name, visitor_email, message, reason, status, COALESCE(email_status, '') AS email_status, COALESCE(webhook_status, '') AS webhook_status, COALESCE(delivery_error, '') AS delivery_error, created_at FROM customer_agent_escalations WHERE agent_id = ? ORDER BY created_at DESC LIMIT ? """, (agent_id, limit), ).fetchall() return [ { "object": "synderesis.customer_agent.escalation", "id": str(row["id"]), "agent_id": str(row["agent_id"]), "conversation_id": str(row["conversation_id"] or ""), "visitor_name": str(row["visitor_name"] or ""), "visitor_email": str(row["visitor_email"] or ""), "message": str(row["message"] or ""), "reason": str(row["reason"] or ""), "status": str(row["status"]), "email_status": str(row["email_status"] or ""), "webhook_status": str(row["webhook_status"] or ""), "delivery_error": str(row["delivery_error"] or ""), "created_at": str(row["created_at"]), } for row in rows ] def extract_knowledge_text_from_upload(*, filename: str, data: bytes) -> tuple[str, str]: """Return (title, body) from an uploaded knowledge file.""" if len(data) > MAX_KNOWLEDGE_FILE_BYTES: raise ValueError(f"file exceeds {MAX_KNOWLEDGE_FILE_BYTES} bytes") name = (filename or "FAQ").strip() or "FAQ" lower = name.casefold() if lower.endswith((".txt", ".md", ".markdown", ".csv")): body = data.decode("utf-8", errors="replace").strip() title = name.rsplit(".", 1)[0][:MAX_KNOWLEDGE_TITLE_CHARS] if not body: raise ValueError("file is empty") return title, body if lower.endswith(".pdf"): body = _extract_pdf_text(data) if not body: raise ValueError( "PDF contained no extractable text (scanned images need OCR; try a .txt/.md export)" ) return name.rsplit(".", 1)[0][:MAX_KNOWLEDGE_TITLE_CHARS], body[:MAX_KNOWLEDGE_BODY_CHARS] raise ValueError("supported knowledge files: .txt, .md, .csv, .pdf") def _extract_pdf_text(data: bytes) -> str: """Extract text from PDF using pypdf (declared in pyproject.toml).""" import io try: from pypdf import PdfReader # type: ignore except ImportError as exc: raise ValueError( "PDF support requires the pypdf package. Run `uv sync` in the Synderesis repo " "(pypdf is listed in pyproject.toml), or upload a .txt / .md file instead." ) from exc try: reader = PdfReader(io.BytesIO(data)) parts: list[str] = [] for page in reader.pages[:40]: parts.append(page.extract_text() or "") return "\n\n".join(part.strip() for part in parts if part and part.strip()).strip() except ValueError: raise except Exception as exc: # noqa: BLE001 raise ValueError(f"could not read PDF text: {exc}") from exc def decode_base64_file(data_base64: str) -> bytes: raw = data_base64.strip() if "," in raw and raw.lower().startswith("data:"): raw = raw.split(",", 1)[1] try: return base64.b64decode(raw, validate=False) except Exception as exc: # noqa: BLE001 raise ValueError("invalid base64 file data") from exc def purge_old_escalations( connection: sqlite3.Connection, *, retention_days: int = DEFAULT_ESCALATION_RETENTION_DAYS, ) -> int: if retention_days <= 0: return 0 cutoff = (datetime.now(UTC) - timedelta(days=retention_days)).replace(microsecond=0).isoformat().replace("+00:00", "Z") cursor = connection.execute( "DELETE FROM customer_agent_escalations WHERE created_at < ?", (cutoff,), ) deleted = int(cursor.rowcount or 0) connection.execute( "DELETE FROM customer_agent_chat_events WHERE created_at < ?", (cutoff,), ) return deleted def log_chat_event( connection: sqlite3.Connection, *, agent_id: str, question: str, answer: str, knowledge_hit: bool, escalate_recommended: bool, hard_refuse: bool, retrieved_count: int, top_score: float, ) -> None: connection.execute( """ INSERT INTO customer_agent_chat_events ( agent_id, created_at, question, answer_preview, knowledge_hit, escalate_recommended, hard_refuse, retrieved_count, top_score ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) """, ( agent_id, utc_now(), question.strip()[:500], (answer or "").strip()[:300], 1 if knowledge_hit else 0, 1 if escalate_recommended else 0, 1 if hard_refuse else 0, int(retrieved_count), float(top_score), ), ) def coverage_report(connection: sqlite3.Connection, agent_id: str, *, limit: int = 25) -> dict[str, Any]: totals = connection.execute( """ SELECT COUNT(*) AS n, COALESCE(SUM(knowledge_hit), 0) AS hits, COALESCE(SUM(escalate_recommended), 0) AS escalates, COALESCE(SUM(hard_refuse), 0) AS refuses FROM customer_agent_chat_events WHERE agent_id = ? """, (agent_id,), ).fetchone() n = int(totals["n"] or 0) hits = int(totals["hits"] or 0) escalates = int(totals["escalates"] or 0) refuses = int(totals["refuses"] or 0) unknown_rows = connection.execute( """ SELECT question, COUNT(*) AS times, MAX(created_at) AS last_seen FROM customer_agent_chat_events WHERE agent_id = ? AND knowledge_hit = 0 GROUP BY question ORDER BY times DESC, last_seen DESC LIMIT ? """, (agent_id, limit), ).fetchall() unknowns = [ { "question": str(row["question"]), "times": int(row["times"] or 0), "last_seen": str(row["last_seen"] or ""), } for row in unknown_rows ] return { "object": "synderesis.customer_agent.coverage", "agent_id": agent_id, "total_chats": n, "knowledge_hit_rate": (hits / n) if n else 0.0, "escalate_rate": (escalates / n) if n else 0.0, "hard_refuse_count": refuses, "unknown_questions": unknowns, "suggestion": ( "Add FAQ entries covering the unknown questions below." if unknowns else "No coverage gaps logged yet — keep testing real customer questions." ), } def export_agent_bundle(connection: sqlite3.Connection, agent: CustomerAgentRecord) -> dict[str, Any]: """Owner export of agent config + knowledge metadata + escalations (no site key secret).""" docs = list_knowledge_docs(connection, agent.id) escalations = list_escalations(connection, agent.id, limit=500) return { "object": "synderesis.customer_agent.export", "exported_at": utc_now(), "agent": agent.owner_dict(knowledge_count=len(docs)), "knowledge": docs, "escalations": escalations, } @dataclass(frozen=True) class EscalationMailConfig: host: str = "" port: int = 587 username: str = "" password: str = "" from_addr: str = "" use_tls: bool = True enabled: bool = False def send_escalation_email( mail: EscalationMailConfig, *, to_addr: str, agent_name: str, escalation: dict[str, Any], ) -> str: """Send escalation email. Returns status: sent|skipped|failed:...""" if not mail.enabled or not mail.host or not mail.from_addr or not to_addr: return "skipped" subject = f"[Synderesis] Message for {agent_name}" transcript = escalation.get("transcript") or [] if isinstance(transcript, str): transcript_text = transcript else: transcript_text = "\n".join( f"{item.get('role', '?')}: {item.get('content', '')}" for item in transcript ) body = ( f"Agent: {agent_name}\n" f"Escalation ID: {escalation.get('id', '')}\n" f"Visitor: {escalation.get('visitor_name', '')} <{escalation.get('visitor_email', '')}>\n" f"Reason: {escalation.get('reason', '')}\n\n" f"Message:\n{escalation.get('message', '')}\n\n" f"Transcript:\n{transcript_text}\n" ) msg = EmailMessage() msg["Subject"] = subject msg["From"] = mail.from_addr msg["To"] = to_addr msg.set_content(body) try: with smtplib.SMTP(mail.host, mail.port, timeout=20) as smtp: if mail.use_tls: smtp.starttls() if mail.username: smtp.login(mail.username, mail.password) smtp.send_message(msg) return "sent" except Exception as exc: # noqa: BLE001 LOGGER.warning("escalation email failed: %s", exc) return f"failed:{exc}" def post_escalation_webhook( *, webhook_url: str, signing_secret: str, agent: CustomerAgentRecord, escalation: dict[str, Any], ) -> str: """POST signed JSON to customer webhook. Returns status: sent|skipped|failed:...""" url = (webhook_url or "").strip() if not url: return "skipped" try: url = validate_public_https_url(url, purpose="webhook") except ValueError as exc: return f"failed:{exc}" payload = { "object": "synderesis.customer_agent.escalation.event", "event": "escalation.created", "agent": {"id": agent.id, "name": agent.name, "customer_id": agent.customer_id}, "escalation": escalation, } body = json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode("utf-8") headers = { "Content-Type": "application/json", "User-Agent": "SynderesisCustomerAgent/1.1", } if signing_secret: digest = hmac.new(signing_secret.encode("utf-8"), body, hashlib.sha256).hexdigest() headers["X-Synderesis-Signature"] = f"sha256={digest}" req = urllib.request.Request(url, data=body, headers=headers, method="POST") try: with urllib.request.urlopen(req, timeout=15) as resp: if 200 <= int(resp.status) < 300: return "sent" return f"failed:http_{resp.status}" except urllib.error.HTTPError as exc: return f"failed:http_{exc.code}" except Exception as exc: # noqa: BLE001 LOGGER.warning("escalation webhook failed: %s", exc) return f"failed:{exc}" def deliver_escalation_notifications( *, agent: CustomerAgentRecord, escalation: dict[str, Any], mail: EscalationMailConfig, webhook_signing_secret: str = "", ) -> dict[str, str]: """Send email + webhook; returns status map.""" # Attach transcript for email body if only stored as JSON on create payload. email_status = send_escalation_email( mail, to_addr=agent.escalate_email, agent_name=agent.name, escalation=escalation, ) webhook_status = post_escalation_webhook( webhook_url=agent.escalate_webhook_url, signing_secret=webhook_signing_secret, agent=agent, escalation=escalation, ) return {"email_status": email_status, "webhook_status": webhook_status} def seed_demo_bethlehem_agent( connection: sqlite3.Connection, customer_id: str, ) -> tuple[CustomerAgentRecord, str]: """Create a sample agent used by the marketing demo when none exists.""" agent, site_key = create_agent( connection, customer_id=customer_id, name="Bethlehem Books & Gifts", greeting=( "Hello, and welcome to Bethlehem Books & Gifts! " "Ask me about books and gifts, shipping, opening hours, or anything else." ), voice_instructions=( "You are warm, helpful, concise and professional, and you reflect the shop's " "Catholic values without preaching." ), escalate_email="help@bethlehembooks.example", allowed_origins=["*"], ) faq = ( "Opening hours: Monday to Saturday 9:00-18:00. Closed Sundays and holy days of obligation.\n\n" "Shipping: We ship across the EU in 3-5 working days. Free shipping over EUR 60.\n\n" "Range: Catholic books and Bibles, rosaries, icons, candles, and gifts for Baptism, " "First Communion, and Confirmation.\n\n" "Returns: Accepted within 30 days on unused items in original condition.\n\n" "Location: Family-run Catholic bookshop based in Ireland." ) add_knowledge_doc( connection, agent_id=agent.id, title="Shop FAQ", body=faq, source_type="manual", ) return agent, site_key