"""
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.
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()