# -*- coding: utf-8 -*- """Astraea — the world's best Washington divorce-law attorney. Multi-agent RAG + no-KYC profiles (About Me) + document prep + comms missions + TTS. """ import json import os import re import threading import time import urllib.parse import urllib.request import uuid from flask import Flask, jsonify, render_template, request, send_file, Response import rag import store import comms import documents import forms import tts import userdocs from agents import AGENTS, AGENT_BY_ID, OLLAMA_URL, RAG_MODEL, GENERAL_MODEL app = Flask(__name__) # ── Rate limiting (SQLite-backed sliding window; shared across gunicorn workers) ── def _rate_limit(key, max_events, window_seconds): """True if allowed, False if the caller exceeded max_events in window. Counters live in the astraea.db `rate_limits` table so all workers share them.""" now = time.time() try: c = store._conn() c.execute("CREATE TABLE IF NOT EXISTS rate_limits (key TEXT, ts REAL)") c.execute("DELETE FROM rate_limits WHERE ts < ?", (now - 3600,)) row = c.execute("SELECT COUNT(*) FROM rate_limits WHERE key=? AND ts >= ?", (key, now - window_seconds)).fetchone() if row[0] >= max_events: c.commit(); c.close(); return False c.execute("INSERT INTO rate_limits (key, ts) VALUES (?,?)", (key, now)) c.commit(); c.close() return True except Exception: return True # fail open — availability over strictness def _client_ip(): return request.headers.get("X-Real-IP") or request.remote_addr or "unknown" # ── Professional avatars (gradient + glyph, rendered client-side) ── AVATARS = [ {"id": 0, "glyph": "\u2696", "c1": "#6ea8ff", "c2": "#a78bfa", "label": "Scales"}, {"id": 1, "glyph": "\U0001F3DB", "c1": "#7dd3fc", "c2": "#0ea5e9", "label": "Court"}, {"id": 2, "glyph": "\U0001F6E1", "c1": "#4ade80", "c2": "#059669", "label": "Shield"}, {"id": 3, "glyph": "\U0001F4DC", "c1": "#f5d06f", "c2": "#d97706", "label": "Scroll"}, {"id": 4, "glyph": "\U0001F54A", "c1": "#f0abfc", "c2": "#c026d3", "label": "Dove"}, {"id": 5, "glyph": "\U0001F4BC", "c1": "#fdba74", "c2": "#ea580c", "label": "Briefcase"}, {"id": 6, "glyph": "\U0001F512", "c1": "#67e8f9", "c2": "#0891b2", "label": "Key"}, {"id": 7, "glyph": "\u2731", "c1": "#fda4af", "c2": "#e11d48", "label": "Star"}, {"id": 8, "glyph": "\U0001F3AF", "c1": "#a3e635", "c2": "#65a30d", "label": "Target"}, {"id": 9, "glyph": "\U0001F525", "c1": "#f87171", "c2": "#b91c1c", "label": "Flame"}, {"id": 10, "glyph": "\U0001F31F", "c1": "#fbbf24", "c2": "#b45309", "label": "Glow"}, {"id": 11, "glyph": "\U0001F9ED", "c1": "#5eead4", "c2": "#0d9488", "label": "Compass"}, ] UPLOAD_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "uploads") FILE_CATEGORIES = ["Court filings", "Correspondence", "Financial", "Agreements", "Evidence", "Medical & safety", "Generated forms", "Other"] def _docx_text(path): import zipfile as _z try: with _z.ZipFile(path) as z: xml = z.read("word/document.xml").decode("utf-8", "ignore") out = [] for p in re.split(r"", xml): ts = re.findall(r"]*>(.*?)", p) if ts: out.append("".join(ts)) return "\n".join(out) except Exception: return "" def _reindex_user(uid): """Rebuild a user's document vector index in a background thread.""" try: userdocs.build(uid, store.list_files(uid), UPLOAD_DIR) except Exception: pass def _autofill_values(uid, form_id): """Map the user's About Me + display name onto a form's fields via the LLM.""" spec = forms.FORM_NAMES.get(form_id) fm = forms.FIELDMAP.get(form_id, []) if not spec or not fm: return {} profile = store.get_profile(uid) if uid else {} about = profile.get("about_me", "") display = profile.get("display_name", "") # filter to a manageable set: all text fields + checkboxes with a real label (capped) text_fields = [f for f in fm if f.get("type") != "checkbox"] cb_fields = [f for f in fm if f.get("type") == "checkbox" and f.get("label", "").strip()] cb_fields = cb_fields[:40] sel = text_fields + cb_fields field_desc = ", ".join( f"{f['key']} ({f.get('label', '')})" + (" [CB]" if f.get('type') == 'checkbox' else "") for f in sel ) prompt = ( "You are filling a Washington State court form. Map the user's facts to the form fields.\n" f"Form: {spec['name']}\n" f"Fields: {field_desc}\n\n" f"User's display name: {display or '(not set)'}\n" f"User's information (About Me):\n{about or '(not provided)'}\n\n" "Output ONLY valid JSON. For text fields, output the value (or empty string if unknown). " "For CHECKBOX fields, output 'X' ONLY when the user's facts clearly and unambiguously " "match that option; otherwise output ''. Be conservative — do not guess. " "Petitioner is the user themself (use the display name); respondent is their spouse." ) try: resp = _ollama_chat(RAG_MODEL, [{"role": "user", "content": prompt}], temperature=0.1, num_predict=1600) content = resp.get("message", {}).get("content", "") m = re.search(r"\{.*\}", content, re.DOTALL) vals = json.loads(m.group(0)) if m else {} except Exception: vals = {} return {f["key"]: str(vals.get(f["key"], "") or "") for f in sel} def _save_pdf_to_vault(uid, path, name): data = open(path, "rb").read() fn = f"gen_{uid}_{uuid.uuid4().hex[:8]}.pdf" with open(os.path.join(UPLOAD_DIR, fn), "wb") as f: f.write(data) return store.add_file(uid, "Generated forms", fn, name + ".pdf", "application/pdf", len(data)) # ── Ollama helpers ── def _ollama_chat(model, messages, tools=None, num_ctx=24000, num_predict=1400, temperature=0.2, timeout=180): payload = {"model": model, "messages": messages, "stream": False, "think": False, "options": {"temperature": temperature, "num_predict": num_predict, "num_ctx": num_ctx}} if tools: payload["tools"] = tools req = urllib.request.Request(f"{OLLAMA_URL}/api/chat", data=json.dumps(payload).encode(), headers={"Content-Type": "application/json"}) op = urllib.request.build_opener(urllib.request.ProxyHandler({})) with op.open(req, timeout=timeout) as r: return json.loads(r.read().decode("utf-8")) def _ollama_chat_stream(model, messages, num_ctx=24000, num_predict=1400, temperature=0.2, timeout=180): """Yield content deltas from Ollama's streaming /api/chat (no tools — the streaming path is for final answers only; tool-calling stays non-streaming).""" payload = {"model": model, "messages": messages, "stream": True, "think": False, "options": {"temperature": temperature, "num_predict": num_predict, "num_ctx": num_ctx}} req = urllib.request.Request(f"{OLLAMA_URL}/api/chat", data=json.dumps(payload).encode(), headers={"Content-Type": "application/json"}) op = urllib.request.build_opener(urllib.request.ProxyHandler({})) with op.open(req, timeout=timeout) as r: for line in r: line = line.strip() if not line: continue try: chunk = json.loads(line) except Exception: continue if chunk.get("done"): break msg = chunk.get("message", {}) delta = msg.get("content") or "" if delta: yield delta def _web_search(query, limit=5): """Search the live web via local SearXNG (separate service, not the LLM).""" try: url = f"http://10.30.20.35:6969/search?q={urllib.parse.quote(query)}&format=json" req = urllib.request.Request(url, headers={"User-Agent": "Mozilla/5.0 (Astraea)"}) op = urllib.request.build_opener(urllib.request.ProxyHandler({})) with op.open(req, timeout=15) as r: data = json.loads(r.read().decode("utf-8")) out = [] for it in data.get("results", [])[:limit]: out.append({"title": it.get("title", ""), "url": it.get("url", ""), "snippet": (it.get("content") or "")[:400]}) return out except Exception as e: return [{"error": f"web search failed: {e}"}] def _strip_sources(text): text = re.split(r"\n\s*(?:Sources|References|Citations)\s*:\s*\n", text, flags=re.I)[0] return text.strip() # ── About Me injection (designated plainly, used silently) ── def about_me_block(about_me): if not about_me or not about_me.strip(): return "" return ( "\n\n=== ABOUT THE USER (PRIVATE CONTEXT) ===\n" f"{about_me.strip()}\n" "=== END ABOUT THE USER ===\n" "This is private background about the user's situation. Use it silently to tailor your " "answer to their facts. NEVER reference it, quote it, summarize it back, or say anything " "like 'based on what you shared'. Answer ONLY what the user asks." ) # ── Comms tools (agent can invoke to send SMS/email/call) ── TOOLS = [ {"type": "function", "function": { "name": "send_sms", "description": "Send a professionally-worded text message on the user's behalf.", "parameters": {"type": "object", "properties": { "recipient": {"type": "string", "description": "recipient phone number in E.164 (e.g. +14255550123)"}, "message": {"type": "string", "description": "the message text"}}, "required": ["recipient", "message"]}}}, {"type": "function", "function": { "name": "send_email", "description": "Send a formal email on the user's behalf.", "parameters": {"type": "object", "properties": { "recipient": {"type": "string", "description": "recipient email address"}, "subject": {"type": "string"}, "body": {"type": "string"}}, "required": ["recipient", "subject", "body"]}}}, {"type": "function", "function": { "name": "make_call", "description": "Place a phone call that speaks a message.", "parameters": {"type": "object", "properties": { "recipient": {"type": "string", "description": "recipient phone number in E.164"}, "message": {"type": "string", "description": "the message to speak"}}, "required": ["recipient", "message"]}}}, {"type": "function", "function": { "name": "web_search", "description": "Search the live web for current facts, statutes, or answers not in your reference documents. Use when you don't know something or need up-to-date information.", "parameters": {"type": "object", "properties": { "query": {"type": "string", "description": "the search query"}}, "required": ["query"]}}}, {"type": "function", "function": { "name": "fill_form", "description": "Fill an official Washington court form with the user's info and save the filled PDF to their case documents. Use when the user asks to fill out / prepare / complete paperwork. Form ids: fl200=Summons, fl201=Petition, fl001=Confidential Info, fl140=Parenting Plan, wscss=Child Support Worksheets, fl131=Financial Declaration, fl231=Findings & Conclusions, fl241=Final Divorce Order, fl211=Response, fl223=Motion for Temporary Order.", "parameters": {"type": "object", "properties": { "form_id": {"type": "string", "description": "the form id, e.g. fl201"}, "values": {"type": "object", "description": "optional explicit field values; omitted fields auto-fill from the user's profile"}}, "required": ["form_id"]}}}, ] def _execute_tool(name, args, settings, uid, display_name): recipient = args.get("recipient", "") if comms.no_contact_blocked(settings, recipient): return ("BLOCKED: a no-contact order is in effect for this person. No message was sent. " "Inform the user you cannot contact this person due to the no-contact order.") try: if name == "send_sms": comms.send_sms(settings, recipient, comms.sms_template(args.get("message", ""))) return f"SMS sent to {recipient}." if name == "send_email": body = comms.email_template(display_name, recipient, args.get("subject", ""), args.get("body", "")) comms.send_email(settings, recipient, args.get("subject", ""), body) return f"Email sent to {recipient}." if name == "make_call": twiml_url = request.host_url.rstrip("/") + "/twilio/voice" comms.make_call(settings, recipient, twiml_url) return f"Call initiated to {recipient}." if name == "web_search": res = _web_search(args.get("query", ""), limit=5) if not res: return "No results found." return json.dumps(res, ensure_ascii=False) if name == "fill_form": form_id = args.get("form_id", "") if not forms.FORM_NAMES.get(form_id): return "Unknown form. Use one of: " + ", ".join(forms.FORM_NAMES.keys()) values = args.get("values") or {} auto = _autofill_values(uid, form_id) merged = {k: str(values.get(k) or auto.get(k) or "") for k in auto} path = forms.fill_form(form_id, merged) if not path: return "Could not generate the form." _save_pdf_to_vault(uid, path, forms.FORM_NAMES[form_id]["name"]) return (f"Filled '{forms.FORM_NAMES[form_id]['name']}' and saved it to the user's " f"Case Documents. Tell the user to open their profile → Case documents to download it.") except Exception as e: return f"ERROR sending {name}: {e}" return "Unknown tool." # ── RAG chat with tool-calling ── def safety_block(settings): if not (settings.get("domestic_violence") or settings.get("protection_order")): return "" lines = ["\n\n=== SAFETY CONTEXT (PRIVATE — CRITICAL) ==="] if settings.get("domestic_violence"): lines.append("DOMESTIC VIOLENCE IS A FACTOR IN THIS CASE.") if settings.get("protection_order"): lines.append("An active protection order is in effect (RCW 7.105).") notes = (settings.get("safety_notes") or "").strip() if notes: lines.append(f"Safety notes: {notes}") lines.append("=== END SAFETY CONTEXT ===") lines.append( "Prioritize the user's safety above all else. Apply RCW 26.09.191 restrictions on " "residential time and decision-making where domestic violence or abuse is a factor. " "Recommend protection orders, supervised visitation, and safety planning where " "appropriate. NEVER advise or facilitate direct contact with the abusive or restrained party." ) return "\n".join(lines) def _build_answer(agent, message, history, about_me, settings, uid, display_name): top = rag.retrieve(message, agent["id"], top_k=5) context_blocks = [] citations = [] seen = set() for i, (score, c) in enumerate(top, 1): key = (c["source"], c["title"]) if key in seen: continue seen.add(key) context_blocks.append(f"[{i}] ({c['source']} — {c['title']})\n{c['text']}") citations.append({"source": c["source"], "title": c["title"], "score": round(score, 3), "snippet": c["text"][:260]}) # user's own uploaded documents (per-user vector index) — retrieve alongside case law if uid: try: for d in userdocs.retrieve(uid, message, top_k=4): context_blocks.append(f"[USER DOC — {d['source']}]\n{d['text']}") citations.append({"source": d["source"], "title": "your document", "score": d["score"], "snippet": d["text"][:260]}) except Exception: pass system = ( f"{agent['system']}\n\n" "Rules:\n" "- Give REAL, practical legal advice: state the rule, apply it to the user's facts, " "tell them what to do next.\n" "- Answer DIRECTLY. Do not restate the question or narrate your reasoning.\n" "- Ground every legal claim in the reference documents. Cite the RCW section AND the " "controlling case law by name.\n" "- Be plain-English, specific to Washington State.\n" "- Items marked '[USER DOC]' are the user's own uploaded files — reference them by name " "to ground your answer in their actual case.\n" "- You may use the available tools (send_sms / send_email / make_call) ONLY when the " "user explicitly asks you to contact someone on their behalf.\n" "- If you don't know the answer, or need current/up-to-date facts not in the reference " "documents, use the web_search tool and answer from those results.\n" + about_me_block(about_me) + safety_block(settings) ) user = ("REFERENCE DOCUMENTS:\n\n" + "\n\n".join(context_blocks) + f"\n\nUSER QUESTION: {message}") messages = [{"role": "system", "content": system}] if history: for turn in history[-6:]: if turn.get("role") in ("user", "assistant"): messages.append({"role": turn["role"], "content": turn["content"]}) messages.append({"role": "user", "content": user}) try: resp = _ollama_chat(RAG_MODEL, messages, tools=TOOLS) except Exception: try: resp = _ollama_chat(GENERAL_MODEL, messages) except Exception: return {"answer": "The legal engine is temporarily unavailable. Try again shortly.", "citations": [], "grounded": False} # tool-calling loop for _ in range(3): msg = resp.get("message", {}) tool_calls = msg.get("tool_calls") or [] if not tool_calls: break messages.append(msg) for tc in tool_calls: fn = tc.get("function", {}) name = fn.get("name", "") args = fn.get("arguments", {}) if isinstance(args, str): try: args = json.loads(args) except Exception: args = {} result = _execute_tool(name, args, settings, uid, display_name) messages.append({"role": "tool", "content": result}) resp = _ollama_chat(RAG_MODEL, messages, tools=TOOLS) final = resp.get("message", {}) answer = (final.get("content") or final.get("thinking") or "").strip() if not answer: # transient empty response — fall back to the reliable model (no tools) try: resp = _ollama_chat(GENERAL_MODEL, messages) answer = (resp.get("message", {}).get("content") or "").strip() except Exception: answer = "" answer = _strip_sources(answer) return {"answer": answer, "citations": citations, "grounded": bool(context_blocks)} # ── Auth helpers ── def _uid(): token = request.headers.get("Authorization", "").replace("Bearer ", "") return store.user_from_token(token) def _profile_or_none(): uid = _uid() return uid, store.get_profile(uid) if uid else None, store.get_settings(uid) if uid else {} # ── Pages ── @app.route("/") def landing(): return render_template("landing.html", agents=AGENTS) @app.route("/app") def app_page(): return render_template("index.html", agents=AGENTS) # ── Auth API ── @app.route("/api/register", methods=["POST"]) def api_register(): ip = _client_ip() if not _rate_limit(f"register:{ip}", 5, 3600): return jsonify({"error": "Too many attempts. Try again later."}), 429 d = request.get_json(force=True, silent=True) or {} u = (d.get("username") or "").strip().lower() p = (d.get("password") or "").strip() if len(u) < 2 or len(u) > 30 or not u.isalnum(): return jsonify({"error": "Username must be 2-30 alphanumeric chars"}), 400 if len(p) < 8: return jsonify({"error": "Password must be 8+ chars"}), 400 c = store._conn() if c.execute("SELECT id FROM users WHERE username=?", (u,)).fetchone(): c.close() return jsonify({"error": "Username taken"}), 409 c.close() uid = store.create_user(u, p) return jsonify({"ok": True, "token": store.issue_token(uid), "username": u}) @app.route("/api/login", methods=["POST"]) def api_login(): ip = _client_ip() # 10 attempts / 15 min / IP, and 10 / 15 min per username (brute-force guard) d = request.get_json(force=True, silent=True) or {} u = (d.get("username") or "").strip().lower() if not _rate_limit(f"login-ip:{ip}", 10, 900) or \ (u and not _rate_limit(f"login-user:{u}", 10, 900)): return jsonify({"error": "Too many attempts. Try again later."}), 429 p = (d.get("password") or "").strip() uid = store.authenticate(u, p) if not uid: return jsonify({"error": "Invalid username or password"}), 401 return jsonify({"ok": True, "token": store.issue_token(uid), "username": u}) @app.route("/api/me") def api_me(): uid = _uid() if not uid: return jsonify({"authenticated": False}), 401 return jsonify({"authenticated": True, "profile": store.get_profile(uid), "settings": store.get_settings(uid)}) @app.route("/api/profile", methods=["GET", "POST"]) def api_profile(): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 if request.method == "GET": return jsonify(store.get_profile(uid)) d = request.get_json(force=True, silent=True) or {} prof = store.set_profile(uid, display_name=d.get("display_name"), avatar=d.get("avatar"), about_me=d.get("about_me")) return jsonify(prof) @app.route("/api/profile/pic", methods=["POST"]) def api_profile_pic(): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 f = request.files.get("file") if not f or not f.filename: return jsonify({"error": "no file"}), 400 data = f.read() if len(data) > 5 * 1024 * 1024: return jsonify({"error": "image too large (max 5 MB)"}), 400 ext = os.path.splitext(f.filename)[1].lower() if ext not in (".png", ".jpg", ".jpeg", ".gif", ".webp"): return jsonify({"error": "unsupported type — use PNG/JPG/GIF/WebP"}), 400 up = os.path.join(os.path.dirname(os.path.abspath(__file__)), "static", "uploads") os.makedirs(up, exist_ok=True) fn = f"u{uid}{ext}" with open(os.path.join(up, fn), "wb") as out: out.write(data) prof = store.set_profile(uid, profile_pic=f"/static/uploads/{fn}") return jsonify(prof) # ── Document vault ── @app.route("/api/files", methods=["GET", "POST"]) def api_files(): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 if request.method == "GET": return jsonify({"categories": FILE_CATEGORIES, "files": store.list_files(uid)}) cat = request.form.get("category", "Other") if cat not in FILE_CATEGORIES: cat = "Other" f = request.files.get("file") if not f or not f.filename: return jsonify({"error": "no file"}), 400 data = f.read() if len(data) > 25 * 1024 * 1024: return jsonify({"error": "file too large (max 25 MB)"}), 400 os.makedirs(UPLOAD_DIR, exist_ok=True) ext = os.path.splitext(f.filename)[1].lower()[:12] fn = f"{uid}_{uuid.uuid4().hex[:12]}{ext}" with open(os.path.join(UPLOAD_DIR, fn), "wb") as out: out.write(data) fid = store.add_file(uid, cat, fn, f.filename, f.mimetype or "", len(data)) threading.Thread(target=_reindex_user, args=(uid,), daemon=True).start() return jsonify({"ok": True, "id": fid, "category": cat}) @app.route("/api/files/reindex", methods=["POST"]) def api_files_reindex(): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 n = userdocs.build(uid, store.list_files(uid), UPLOAD_DIR) return jsonify({"ok": True, "chunks": n}) @app.route("/api/files/status") def api_files_status(): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 return jsonify(userdocs.status(uid)) @app.route("/api/files//raw") def api_file_raw(fid): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 rec = store.get_file(uid, fid) if not rec: return jsonify({"error": "not found"}), 404 path = os.path.join(UPLOAD_DIR, rec["filename"]) if not os.path.exists(path): return jsonify({"error": "missing"}), 404 return send_file(path, mimetype=rec["mime"] or "application/octet-stream", as_attachment=False, download_name=rec["orig_name"]) @app.route("/api/files//text") def api_file_text(fid): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 rec = store.get_file(uid, fid) if not rec: return jsonify({"error": "not found"}), 404 path = os.path.join(UPLOAD_DIR, rec["filename"]) ext = os.path.splitext(rec["orig_name"])[1].lower() if ext == ".docx": return jsonify({"text": _docx_text(path), "mime": "text/plain"}) if ext in (".txt", ".md", ".csv", ".log", ".json", ".html", ".xml"): try: with open(path, "r", encoding="utf-8", errors="replace") as fh: return jsonify({"text": fh.read()[:200000], "mime": "text/plain"}) except Exception: return jsonify({"text": "", "mime": "text/plain"}) return jsonify({"text": "", "mime": rec["mime"] or ""}) @app.route("/api/files/", methods=["DELETE"]) def api_file_delete(fid): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 rec = store.get_file(uid, fid) if rec: p = os.path.join(UPLOAD_DIR, rec["filename"]) if os.path.exists(p): try: os.remove(p) except Exception: pass store.delete_file(uid, fid) return jsonify({"ok": True}) # ── Form filling (official WA court forms via pymupdf overlay) ── @app.route("/api/forms") def api_forms(): return jsonify(forms.list_forms()) @app.route("/api/forms/autofill", methods=["POST"]) def api_forms_autofill(): uid = _uid() d = request.get_json(force=True, silent=True) or {} form_id = d.get("form_id", "") if not forms.FORM_NAMES.get(form_id): return jsonify({"error": "unknown form"}), 400 values = _autofill_values(uid, form_id) return jsonify({"form_id": form_id, "values": values}) @app.route("/api/forms/fill", methods=["POST"]) def api_forms_fill(): uid = _uid() d = request.get_json(force=True, silent=True) or {} form_id = d.get("form_id", "") values = d.get("values", {}) path = forms.fill_form(form_id, values) if not path: return jsonify({"error": "could not fill form"}), 400 return send_file(path, mimetype="application/pdf", as_attachment=True, download_name=f"{form_id}_filled.pdf") @app.route("/api/settings", methods=["GET", "POST"]) def api_settings(): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 if request.method == "GET": return jsonify(store.get_settings(uid)) d = request.get_json(force=True, silent=True) or {} s = store.set_settings(uid, **d) # never echo the auth token back in full if "twilio_auth_token" in s: s["twilio_auth_token"] = "••••" if s["twilio_auth_token"] else "" if "smtp_pass" in s: s["smtp_pass"] = "••••" if s["smtp_pass"] else "" return jsonify(s) @app.route("/api/avatars") def api_avatars(): return jsonify(AVATARS) # ── Chat ── @app.route("/api/chat", methods=["POST"]) def api_chat(): uid = _uid() d = request.get_json(force=True, silent=True) or {} agent_id = (d.get("agent_id") or "navigator").strip() message = (d.get("message") or "").strip() client_history = d.get("history") or [] agent = AGENT_BY_ID.get(agent_id) if not agent: return jsonify({"error": "Unknown agent"}), 400 if not message: return jsonify({"error": "Message required"}), 400 message = message[:4000] profile = store.get_profile(uid) if uid else {"about_me": "", "display_name": ""} settings = store.get_settings(uid) if uid else {} # persisted history for authed users; client history as fallback for anonymous history = store.get_conversation(uid, agent_id) if uid else client_history t0 = time.time() result = _build_answer(agent, message, history, profile.get("about_me", ""), settings, uid, profile.get("display_name", "")) if uid and result.get("answer"): store.save_message(uid, agent_id, "user", message) store.save_message(uid, agent_id, "assistant", result["answer"], result.get("citations")) result["latency_ms"] = round((time.time() - t0) * 1000) result["agent"] = agent_id return jsonify(result) # ── Streaming chat (SSE) — same retrieval + safety context, streamed answer ── def _chat_context(agent_id, message): """Shared retrieval/context assembly for /api/chat and /api/chat/stream.""" agent = AGENT_BY_ID.get(agent_id) if not agent: return None, None, None, None message = message[:4000] uid = _uid() profile = store.get_profile(uid) if uid else {"about_me": "", "display_name": ""} settings = store.get_settings(uid) if uid else {} return agent, uid, profile, settings @app.route("/api/chat/stream", methods=["POST"]) def api_chat_stream(): """SSE stream: 'meta' event first (citations), then 'delta' events with the answer text, then 'done'. Falls back to a single 'delta' with an error if the engine fails. Persists the conversation like /api/chat does.""" d = request.get_json(force=True, silent=True) or {} agent_id = (d.get("agent_id") or "navigator").strip() message = (d.get("message") or "").strip() client_history = d.get("history") or [] ip = _client_ip() if not _rate_limit(f"chat:{ip}", 30, 3600): return jsonify({"error": "Too many requests. Try again later."}), 429 agent, uid, profile, settings = _chat_context(agent_id, message) if not agent: return jsonify({"error": "Unknown agent"}), 400 if not message: return jsonify({"error": "Message required"}), 400 history = store.get_conversation(uid, agent_id) if uid else client_history # reuse the exact same context assembly as /api/chat by refactoring _build_answer's front half top = rag.retrieve(message, agent["id"], top_k=5) context_blocks, citations, seen = [], [], set() for i, (score, c) in enumerate(top, 1): key = (c["source"], c["title"]) if key in seen: continue seen.add(key) context_blocks.append(f"[{i}] ({c['source']} — {c['title']})\n{c['text']}") citations.append({"source": c["source"], "title": c["title"], "score": round(score, 3), "snippet": c["text"][:260]}) if uid: try: for dd in userdocs.retrieve(uid, message, top_k=4): context_blocks.append(f"[USER DOC — {dd['source']}]\n{dd['text']}") citations.append({"source": dd["source"], "title": "your document", "score": dd["score"], "snippet": dd["text"][:260]}) except Exception: pass system = ( f"{agent['system']}\n\n" "Rules:\n" "- Give REAL, practical legal advice: state the rule, apply it to the user's facts, " "tell them what to do next.\n" "- Answer DIRECTLY. Do not restate the question or narrate your reasoning.\n" "- Ground every legal claim in the reference documents. Cite the RCW section AND the " "controlling case law by name.\n" "- Be plain-English, specific to Washington State.\n" "- Items marked '[USER DOC]' are the user's own uploaded files — reference them by name " "to ground your answer in their actual case.\n" + about_me_block(profile.get("about_me", "")) + safety_block(settings) ) user = ("REFERENCE DOCUMENTS:\n\n" + "\n\n".join(context_blocks) + f"\n\nUSER QUESTION: {message}") messages = [{"role": "system", "content": system}] if history: for turn in history[-6:]: if turn.get("role") in ("user", "assistant"): messages.append({"role": turn["role"], "content": turn["content"]}) messages.append({"role": "user", "content": user}) def generate(): full = [] try: yield f"event: meta\ndata: {json.dumps({'citations': citations, 'agent': agent_id})}\n\n" try: for delta in _ollama_chat_stream(RAG_MODEL, messages): full.append(delta) yield f"event: delta\ndata: {json.dumps({'t': delta})}\n\n" except Exception: # model failure — fall back to the reliable general model try: for delta in _ollama_chat_stream(GENERAL_MODEL, messages): full.append(delta) yield f"event: delta\ndata: {json.dumps({'t': delta})}\n\n" except Exception: yield "event: delta\ndata: {\"t\": \"The legal engine is temporarily unavailable. Try again shortly.\"}\n\n" answer = _strip_sources("".join(full).strip()) if uid and answer: store.save_message(uid, agent_id, "user", message) store.save_message(uid, agent_id, "assistant", answer, citations) yield f"event: done\ndata: {json.dumps({'citations': citations, 'grounded': bool(context_blocks)})}\n\n" except GeneratorExit: # client disconnected mid-stream try: answer = _strip_sources("".join(full).strip()) if uid and answer: store.save_message(uid, agent_id, "user", message) store.save_message(uid, agent_id, "assistant", answer, citations) except Exception: pass raise return Response(generate(), mimetype="text/event-stream", headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"}) @app.route("/api/conversation/", methods=["GET", "DELETE"]) def api_conversation(agent_id): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 if agent_id not in AGENT_BY_ID: return jsonify({"error": "unknown agent"}), 400 if request.method == "DELETE": store.clear_conversation(uid, agent_id) return jsonify({"ok": True}) return jsonify({"history": store.get_conversation(uid, agent_id)}) # ── Documents ── @app.route("/api/documents/types") def api_doc_types(): return jsonify([{"id": k, "name": v["name"], "form": v["form"], "description": v["description"]} for k, v in documents.DOC_TYPES.items()]) @app.route("/api/documents", methods=["POST"]) def api_documents(): uid = _uid() d = request.get_json(force=True, silent=True) or {} doc_type = (d.get("doc_type") or "").strip() instructions = (d.get("instructions") or "").strip() if doc_type not in documents.DOC_TYPES: return jsonify({"error": "Unknown document type"}), 400 profile = store.get_profile(uid) if uid else {"about_me": ""} # retrieve relevant context for grounding top = rag.retrieve(documents.DOC_TYPES[doc_type]["name"], "divorce", top_k=4) ctx = "\n\n".join(f"{c['text']}" for s, c in top) try: text = documents.generate(doc_type, profile.get("about_me", ""), instructions, ctx) except Exception as e: return jsonify({"error": f"generation failed: {e}"}), 500 return jsonify({"ok": True, "document": text, "name": documents.DOC_TYPES[doc_type]["name"]}) # ── TTS ── @app.route("/api/voices") def api_voices(): return jsonify({"voices": tts.list_voices(), "default": "aria"}) @app.route("/api/speak", methods=["POST"]) def api_speak(): d = request.get_json(force=True, silent=True) or {} text = (d.get("text") or "").strip() voice = d.get("voice") or "aria" if not text: return jsonify({"error": "no text"}), 400 try: path = tts.synthesize(text, voice) except Exception as e: return jsonify({"error": str(e)}), 500 return send_file(path, mimetype="audio/mpeg") # ── Missions (manual comms) ── @app.route("/api/mission", methods=["POST"]) def api_mission(): uid = _uid() if not uid: return jsonify({"error": "auth required"}), 401 d = request.get_json(force=True, silent=True) or {} mtype = d.get("type") recipient = d.get("recipient") or "" settings = store.get_settings(uid) profile = store.get_profile(uid) if comms.no_contact_blocked(settings, recipient): return jsonify({"ok": False, "error": "Blocked: a no-contact order is in effect for this person."}), 403 try: if mtype == "sms": comms.send_sms(settings, recipient, comms.sms_template(d.get("message", ""))) elif mtype == "email": body = comms.email_template(profile.get("display_name", ""), recipient, d.get("subject", ""), d.get("body", "")) comms.send_email(settings, recipient, d.get("subject", ""), body) elif mtype == "call": comms.make_call(settings, recipient, request.host_url.rstrip("/") + "/twilio/voice") else: return jsonify({"ok": False, "error": "unknown type"}), 400 except Exception as e: return jsonify({"ok": False, "error": str(e)}), 500 return jsonify({"ok": True}) # ── Twilio voice webhook (TwiML ) ── @app.route("/twilio/voice", methods=["GET", "POST"]) def twilio_voice(): msg = request.values.get("message", "This is an automated message from a legal assistant.") twiml = (f'' f'{msg}') return Response(twiml, mimetype="text/xml") @app.route("/api/health") def health(): return jsonify({"ok": True, "agent_count": len(AGENTS), "index": rag.index_status()}) @app.route("/api/index/status") def index_status(): return jsonify(rag.index_status()) @app.route("/api/reindex", methods=["POST"]) def reindex(): try: idx = rag.build_index() return jsonify({"ok": True, "chunks": len(idx)}) except Exception as e: return jsonify({"ok": False, "error": str(e)}), 500 store.init_db() rag.start_background_index() if __name__ == "__main__": app.run(host="0.0.0.0", port=5000, threaded=True)