""" Web UI demo for Gemini + VektoriAgent (single file). Why this version: - Single request per turn (no duplicate retrieval call) - Proof of memory: explicit retrieved fact/episode/sentence list with scores - Latency visibility: chat latency + evidence-render latency are shown in UI Run: export GOOGLE_API_KEY="your-gemini-api-key" uv pip install -e ".[litellm,sentence-transformers]" google-generativeai aiosqlite python3 examples/gemini_agent_ui.py Open: http://127.0.0.1:8765 """ from __future__ import annotations import asyncio import json import os import sys import threading import time from http import HTTPStatus from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path from uuid import uuid4 # Allow direct execution from source checkout. REPO_ROOT = Path(__file__).resolve().parents[1] if str(REPO_ROOT) not in sys.path: sys.path.insert(0, str(REPO_ROOT)) from vektori import AgentConfig, Vektori, VektoriAgent from vektori.models.factory import create_chat_model HOST = os.getenv("VEKTORI_UI_HOST", "127.0.0.1") PORT = int(os.getenv("VEKTORI_UI_PORT", "8765")) USER_ID = os.getenv("VEKTORI_DEMO_USER_ID", "demo-gemini-user") AGENT_ID = os.getenv("VEKTORI_DEMO_AGENT_ID", "gemini-web-ui-agent") SESSION_ID = os.getenv("VEKTORI_DEMO_SESSION_ID", f"gemini-web-ui-{uuid4()}") # Faster default than full flash model; can override via env. CHAT_MODEL = os.getenv("VEKTORI_CHAT_MODEL", "litellm:gemini/gemini-2.5-flash-lite") EXTRACTION_MODEL = os.getenv("VEKTORI_EXTRACT_MODEL", "gemini:gemini-2.5-flash-lite") EMBEDDING_MODEL = os.getenv("VEKTORI_EMBED_MODEL", "sentence-transformers:all-MiniLM-L6-v2") HTML = """ Vektori Gemini Agent UI

Vektori Gemini Agent UI

Fast chat + explicit retrieval evidence so memory claims are verifiable.

Agent Chat

Memory Proof + Latency

status: ready
retrieval_reason: -
facts / episodes / sentences: -
chat_latency: -
retrieval_latency: -
model_latency: -
other_overhead: -
evidence_render_latency: -
Proof = retrieved evidence below (not just generated text).
""" def _to_evidence_item(item: dict, prefer_score: bool = True) -> dict: score = None if prefer_score and "score" in item: score = item.get("score") elif "distance" in item: # distance lower is better; expose as score-like signal for readability try: score = 1 - float(item.get("distance", 1.0)) except Exception: score = None return { "text": item.get("text", ""), "session_id": item.get("session_id"), "score": score, } class Runtime: def __init__(self) -> None: self.loop = asyncio.new_event_loop() self.thread = threading.Thread(target=self._run_loop, daemon=True) self.thread.start() self.memory: Vektori | None = None self.agent: VektoriAgent | None = None self._last_retrieval_ms: int | None = None self._model_ms_accum: int = 0 self._model_calls: int = 0 self.run(self._init()) def _run_loop(self) -> None: asyncio.set_event_loop(self.loop) self.loop.run_forever() def run(self, coro): future = asyncio.run_coroutine_threadsafe(coro, self.loop) return future.result() async def _seed_if_needed(self) -> None: assert self.memory is not None existing = await self.memory.search( query="support teams handling many follow-ups", user_id=USER_ID, agent_id=AGENT_ID, depth="l0", top_k=1, ) if existing.get("facts"): return await self.memory.add( messages=[ {"role": "user", "content": "I prefer concise responses with bullet points."}, {"role": "assistant", "content": "Understood. I will keep responses concise."}, { "role": "user", "content": "Our demo target is support teams handling many follow-ups.", }, ], session_id="seed-web-ui-001", user_id=USER_ID, agent_id=AGENT_ID, ) await self.memory.add( messages=[ {"role": "user", "content": "Avoid long narratives unless I explicitly ask."}, {"role": "assistant", "content": "Got it. I will stay brief and action-oriented."}, ], session_id="seed-web-ui-002", user_id=USER_ID, agent_id=AGENT_ID, ) async def _init(self) -> None: self.memory = Vektori( embedding_model=EMBEDDING_MODEL, extraction_model=EXTRACTION_MODEL, async_extraction=False, ) chat_model = create_chat_model(CHAT_MODEL) self.agent = VektoriAgent( memory=self.memory, model=chat_model, user_id=USER_ID, agent_id=AGENT_ID, session_id=SESSION_ID, config=AgentConfig( retrieve_on_every_turn=True, background_add=True, # faster turn latency (writes async in background) retrieval_depth="l1", retrieval_top_k=3, reserve_response_tokens=220, max_context_tokens=6000, ), ) self._instrument_timings() await self._seed_if_needed() def _instrument_timings(self) -> None: assert self.memory is not None assert self.agent is not None original_search = self.memory.search async def timed_search(*args, **kwargs): started = time.perf_counter() out = await original_search(*args, **kwargs) self._last_retrieval_ms = int((time.perf_counter() - started) * 1000) return out self.memory.search = timed_search # type: ignore[assignment] original_complete = self.agent.model.complete async def timed_complete(*args, **kwargs): started = time.perf_counter() out = await original_complete(*args, **kwargs) self._model_ms_accum += int((time.perf_counter() - started) * 1000) self._model_calls += 1 return out self.agent.model.complete = timed_complete # type: ignore[assignment] async def chat(self, message: str) -> dict: assert self.agent is not None self._last_retrieval_ms = None self._model_ms_accum = 0 self._model_calls = 0 started = time.perf_counter() result = await self.agent.chat(message) chat_elapsed_ms = int((time.perf_counter() - started) * 1000) evidence_started = time.perf_counter() facts = [_to_evidence_item(x) for x in result.memories_used.get("facts", [])][:3] episodes = [_to_evidence_item(x) for x in result.memories_used.get("episodes", [])][:3] sentences = [ _to_evidence_item(x, prefer_score=False) for x in result.memories_used.get("sentences", []) ][:3] evidence_elapsed_ms = int((time.perf_counter() - evidence_started) * 1000) accounted = (self._last_retrieval_ms or 0) + self._model_ms_accum other_overhead_ms = max(0, chat_elapsed_ms - accounted) return { "content": result.content, "retrieval_debug": result.retrieval_debug, "chat_latency_ms": chat_elapsed_ms, "retrieval_latency_ms": self._last_retrieval_ms, "model_latency_ms": self._model_ms_accum, "model_calls": self._model_calls, "other_overhead_ms": other_overhead_ms, "facts": facts, "episodes": episodes, "sentences": sentences, "evidence_latency_ms": evidence_elapsed_ms, "status": "ok", } async def reset(self) -> dict: assert self.agent is not None self.agent.reset_window() return {"status": "conversation reset"} async def close(self) -> None: if self.agent is not None: await self.agent.close() if self.memory is not None: await self.memory.close() runtime: Runtime | None = None class Handler(BaseHTTPRequestHandler): def _json(self, payload: dict, status: HTTPStatus = HTTPStatus.OK) -> None: body = json.dumps(payload).encode("utf-8") self.send_response(status.value) self.send_header("Content-Type", "application/json; charset=utf-8") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) def _html(self, text: str) -> None: body = text.encode("utf-8") self.send_response(HTTPStatus.OK.value) self.send_header("Content-Type", "text/html; charset=utf-8") self.send_header("Content-Length", str(len(body))) self.end_headers() self.wfile.write(body) def _read_json(self) -> dict: length = int(self.headers.get("Content-Length", "0")) raw = self.rfile.read(length) if length > 0 else b"{}" return json.loads(raw.decode("utf-8")) def log_message(self, fmt: str, *args) -> None: # noqa: A003 return def do_GET(self) -> None: # noqa: N802 if self.path == "/": self._html(HTML) return if self.path == "/api/health": self._json({"status": "ok"}) return self._json({"error": "not found"}, HTTPStatus.NOT_FOUND) def do_POST(self) -> None: # noqa: N802 global runtime if runtime is None: self._json({"error": "runtime not initialized"}, HTTPStatus.INTERNAL_SERVER_ERROR) return try: if self.path == "/api/chat": data = self._read_json() message = str(data.get("message", "")).strip() if not message: self._json({"error": "message is required"}, HTTPStatus.BAD_REQUEST) return self._json(runtime.run(runtime.chat(message))) return if self.path == "/api/reset": self._json(runtime.run(runtime.reset())) return self._json({"error": "not found"}, HTTPStatus.NOT_FOUND) except Exception as e: self._json({"error": str(e)}, HTTPStatus.INTERNAL_SERVER_ERROR) def main() -> None: global runtime if not os.getenv("GOOGLE_API_KEY"): raise SystemExit("GOOGLE_API_KEY is not set. Export it and re-run.") runtime = Runtime() server = ThreadingHTTPServer((HOST, PORT), Handler) print(f"Vektori Gemini Web UI running at http://{HOST}:{PORT}") print("Press Ctrl+C to stop.") try: server.serve_forever() except KeyboardInterrupt: pass finally: server.server_close() if runtime is not None: runtime.run(runtime.close()) runtime.loop.call_soon_threadsafe(runtime.loop.stop) if __name__ == "__main__": main()