Compare commits
11
Commits
ab6b5254ef
..
master
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3700b52883 | ||
|
|
65deb0e4ce | ||
|
|
c815bcb485 | ||
|
|
f80b9ebe3d | ||
|
|
98168f2876 | ||
|
|
d08a4e1fcb | ||
|
|
dd33837055 | ||
|
|
44e8d05e2c | ||
|
|
d372adca7d | ||
|
|
a1d59580f7 | ||
|
|
938629df83 |
@@ -5,3 +5,6 @@ __pycache__/
|
||||
node_modules/
|
||||
*.log
|
||||
certs/
|
||||
.venv-tests/
|
||||
__pycache__/
|
||||
.pytest_cache/
|
||||
|
||||
+4
-1
@@ -42,8 +42,11 @@ COPY --from=build /app/agent/.venv /opt/voice-agent/.venv
|
||||
COPY agent/agent.py /opt/voice-agent/agent.py
|
||||
COPY agent/web_mcp.py /opt/voice-agent/web_mcp.py
|
||||
COPY agent/weather_mcp.py /opt/voice-agent/weather_mcp.py
|
||||
COPY agent/memory_mcp.py /opt/voice-agent/memory_mcp.py
|
||||
|
||||
COPY agent/skills_mcp.py /opt/voice-agent/skills_mcp.py
|
||||
COPY agent/task_registry.py /opt/voice-agent/task_registry.py
|
||||
COPY agent/task_worker.py /opt/voice-agent/task_worker.py
|
||||
COPY agent/dispatch_mcp.py /opt/voice-agent/dispatch_mcp.py
|
||||
|
||||
# Copy web frontend + token endpoint
|
||||
COPY web/index.html /var/www/voice/
|
||||
|
||||
@@ -0,0 +1,152 @@
|
||||
# Memory Server (Cognee MCP)
|
||||
|
||||
A persistent, centralized **AI memory** service you can attach to any agentic
|
||||
harness (OpenCode, Hermes, Claude Desktop, custom agents). It runs on the Unraid
|
||||
box at `192.168.86.2` as two Docker containers managed by a single compose file
|
||||
(`/mnt/user/appdata/cognee/docker-compose.yml`):
|
||||
|
||||
| Container | Purpose | Port (host) |
|
||||
|----------------------|-------------------------------------------|-------------|
|
||||
| `shane-cognee` | Cognee API + graph/vector store | `8002` (REST) |
|
||||
| `shane-cognee-mcp` | MCP frontend — **Streamable HTTP** | `8003` (`/mcp`) |
|
||||
| `shane-cognee-mcp-sse` | MCP frontend — **SSE** (for clients that can't do Streamable HTTP) | `8004` (`/sse`) |
|
||||
|
||||
Both MCP frontends bridge to the same Cognee API; all memory state lives in the
|
||||
API container, so the two frontends are interchangeable.
|
||||
|
||||
## Endpoints
|
||||
|
||||
- **Streamable HTTP (preferred for most clients):** `http://192.168.86.2:8003/mcp`
|
||||
- **SSE (for OpenCode 1.18.21, Claude Desktop, etc.):** `http://192.168.86.2:8004/sse`
|
||||
- **REST API (custom apps):** `http://192.168.86.2:8002` (docs at `/docs`)
|
||||
|
||||
## Connecting an agentic harness
|
||||
|
||||
Pick the transport your client supports. Use **SSE (8004)** if the client only
|
||||
does SSE, otherwise **Streamable HTTP (8003)**.
|
||||
|
||||
### OpenCode
|
||||
|
||||
OpenCode 1.18.21 only supports SSE (`type: "remote"`). Already configured in
|
||||
`~/.config/opencode/opencode.jsonc`:
|
||||
|
||||
```jsonc
|
||||
"mcp": {
|
||||
"cognee": {
|
||||
"type": "remote",
|
||||
"url": "http://192.168.86.2:8004/sse"
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
Verify with `opencode mcp list` — it should show `✓ cognee connected`.
|
||||
|
||||
### Hermes
|
||||
|
||||
Hermes speaks Streamable HTTP natively. Already configured in
|
||||
`~/.hermes/config.yaml` under `mcp_servers`:
|
||||
|
||||
```yaml
|
||||
mcp_servers:
|
||||
cognee:
|
||||
url: http://192.168.86.2:8003/mcp
|
||||
transport: http
|
||||
enabled: true
|
||||
```
|
||||
|
||||
### Claude Desktop / other MCP-aware clients
|
||||
|
||||
Use the SSE endpoint. Example `claude_desktop_config.json`:
|
||||
|
||||
```json
|
||||
{
|
||||
"mcpServers": {
|
||||
"cognee": {
|
||||
"url": "http://192.168.86.2:8004/sse"
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
For clients that support Streamable HTTP, point them at
|
||||
`http://192.168.86.2:8003/mcp` instead.
|
||||
|
||||
### Stdio bridge (legacy clients only)
|
||||
|
||||
If a client can only spawn a stdio MCP server, bridge with `mcp-remote`
|
||||
(pointed at the **SSE** endpoint — this `mcp-remote` build is SSE-only):
|
||||
|
||||
```json
|
||||
{
|
||||
"mcpServers": {
|
||||
"cognee": {
|
||||
"command": "npx",
|
||||
"args": ["-y", "mcp-remote", "http://192.168.86.2:8004/sse"]
|
||||
}
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
### Custom applications (REST)
|
||||
|
||||
Call the Cognee REST API directly at `http://192.168.86.2:8002`. See
|
||||
`/docs` for the OpenAPI schema. The MCP tools below are thin wrappers over
|
||||
these endpoints.
|
||||
|
||||
## Authentication
|
||||
|
||||
The current deployment has no API token (`ENABLE_BACKEND_ACCESS_CONTROL:
|
||||
"false"`). If you later enable auth on the Cognee API, pass the token to the
|
||||
MCP frontend via the `API_TOKEN` env var in the compose file and restart the
|
||||
containers.
|
||||
|
||||
## LAN access / Host guard
|
||||
|
||||
`cognee-mcp` enforces a DNS-rebinding Host/Origin guard. The container is
|
||||
configured to accept the Unraid LAN IP via the env var
|
||||
`MCP_ALLOWED_HOSTS: "192.168.86.2:*"` (comma-separated, requires the `:*` glob).
|
||||
|
||||
If you connect from a client that reaches the server by a **different hostname
|
||||
or IP** (e.g. a DNS name, or a different interface), add that pattern to
|
||||
`MCP_ALLOWED_HOSTS` in `docker-compose.yml` and `docker compose up -d`, or you
|
||||
will get `421 Misdirected Request`. Example for two hosts:
|
||||
|
||||
```yaml
|
||||
MCP_ALLOWED_HOSTS: "192.168.86.2:*,cognee.lan:*"
|
||||
```
|
||||
|
||||
To disable the guard entirely (LAN-only, not exposed to the internet), set
|
||||
`MCP_DISABLE_DNS_REBINDING_PROTECTION: "true"`.
|
||||
|
||||
## Tools
|
||||
|
||||
The server exposes 5 MCP tools:
|
||||
|
||||
| Tool | Description |
|
||||
|----------------|-------------|
|
||||
| `remember` | Store data in memory. Without `session_id` → permanent memory (runs the full add + `cognify` pipeline: ingest, entity extraction, graph build). With `session_id` → fast session-cache memory only (no graph). |
|
||||
| `recall` | Search memory with auto-routing and session awareness. Use this for "what do I know about X". |
|
||||
| `forget` | Delete data from memory (by dataset/id). |
|
||||
| `search_tools` | Find a tool by natural-language description (meta-tool for the agent). |
|
||||
| `call_tool` | Invoke a tool by name with arguments (meta-tool). |
|
||||
|
||||
## Usage patterns
|
||||
|
||||
- **Persist something for good:** `remember` without `session_id`. This triggers
|
||||
`cognify`, which can take a while (ingest + embed + graph). Don't call it in a
|
||||
tight loop; batch related content into one call.
|
||||
- **Scratch / per-conversation notes:** `remember` with a `session_id` for fast,
|
||||
non-graph storage you don't need long-term.
|
||||
- **Retrieve:** `recall` with a natural-language query. It auto-routes between
|
||||
graph and vector search and is session-aware.
|
||||
- **Remove:** `forget` when memory is stale.
|
||||
- Keep payloads reasonable; the MCP frontend rejects uploads over 10 MB.
|
||||
|
||||
## Troubleshooting
|
||||
|
||||
- `421 Misdirected Request` → Host header not in `MCP_ALLOWED_HOSTS`. Add the
|
||||
client's target hostname/IP (with `:*` glob) and recreate the container.
|
||||
- Connection refused on `8003`/`8004` → container not healthy yet; check
|
||||
`docker ps` on Unraid and the API container's health (`/health` on `8002`).
|
||||
- `opencode mcp list` shows the server but not `connected` → check the URL uses
|
||||
the SSE endpoint (`/sse`) for OpenCode, and that the Host pattern is allowed.
|
||||
+225
-50
@@ -38,7 +38,6 @@ GEMMA_MODEL = os.environ.get("GEMMA_MODEL", "gemma-4-e4b")
|
||||
GEMMA_API_KEY = os.environ.get("GEMMA_API_KEY", "not-needed")
|
||||
WEB_MCP_ENABLED = os.environ.get("WEB_MCP_ENABLED", "true").lower() in ("1", "true", "yes")
|
||||
FIRECRAWL_BASE = os.environ.get("FIRECRAWL_BASE", "http://192.168.86.2:3002")
|
||||
MEMORY_DIR = os.environ.get("MEMORY_DIR", "/memory")
|
||||
SKILLS_DIR = os.environ.get("SKILLS_DIR", "/skills")
|
||||
|
||||
SYSTEM_PROMPT = textwrap.dedent("""\
|
||||
@@ -63,14 +62,22 @@ SYSTEM_PROMPT = textwrap.dedent("""\
|
||||
- Never write more than three sentences in a row.
|
||||
- Never read back URLs, file paths, or technical identifiers.
|
||||
|
||||
# Web access
|
||||
You have web_search and web_scrape tools. Use them when the user asks about
|
||||
current events, recent news, prices, sports scores, or anything
|
||||
that may have changed since your training data. Search first, then scrape
|
||||
a result only if you need more detail. Answer from what you find, in your
|
||||
normal conversational style — don't cite sources formally, just mention the
|
||||
source naturally ("according to..."). If a search comes up empty, say so
|
||||
briefly and move on.
|
||||
# Memory (CRITICAL — ALWAYS call these tools directly, never dispatch)
|
||||
When the user asks about personal information, past conversations, or
|
||||
anything you might have stored, you MUST call recall FIRST before
|
||||
answering. Do NOT say "I don't have that info" without calling recall.
|
||||
When you call recall, ALWAYS include session_id="hope" to search the
|
||||
session cache.
|
||||
When the user tells you something to remember or shares a preference,
|
||||
you MUST call remember with session_id="hope". Always include
|
||||
session_id="hope" in every remember call — this makes it instant.
|
||||
- User says "what's my name?" → CALL recall with query="my name", session_id="hope"
|
||||
- User says "do you remember who I am?" → CALL recall with query="user identity name", session_id="hope"
|
||||
- User says "remember my coffee order is oat milk latte" → CALL remember
|
||||
with data="User's coffee order is an oat milk latte", session_id="hope"
|
||||
- User says "forget that I like blue" → CALL forget to remove it
|
||||
After calling remember, confirm briefly ("Got it, I'll remember that").
|
||||
If recall returns nothing, then say you don't have that info yet.
|
||||
|
||||
# Weather & Time
|
||||
You can check the weather for any location. Use get_weather when the user
|
||||
@@ -79,19 +86,23 @@ SYSTEM_PROMPT = textwrap.dedent("""\
|
||||
Use get_time when the user asks what time or day it is. If they mention
|
||||
a city, pass it as the location argument.
|
||||
|
||||
# Memory
|
||||
You have persistent memory across conversations, stored as notes you can
|
||||
search and add to.
|
||||
- Use memory_recall at the start of a conversation, or whenever the user
|
||||
references past information ("what did I tell you about...", "remember
|
||||
when..."). Weave what you find in naturally.
|
||||
- Use memory_save when the user shares personal information, preferences,
|
||||
or important facts worth remembering: names, birthdays, preferences,
|
||||
projects, anything they'd expect you to know later. Pick a short topic
|
||||
name for each thing you save.
|
||||
- Be natural about it. Never announce "I'm saving that to memory" — just
|
||||
remember it and move on. If a recall comes up empty, don't mention the
|
||||
search; just answer as if you'd never heard it before.
|
||||
# Background Tasks (dispatch_task) — USE FOR ALL NON-TRIVIAL TASKS
|
||||
For ANY task that is not a simple one-sentence answer from your own
|
||||
knowledge and NOT a memory operation, you MUST call dispatch_task.
|
||||
This includes: looking up news, researching topics, checking current
|
||||
events, prices, sports scores, writing something, summarizing, comparing
|
||||
options, planning, or anything that takes more than a couple seconds to
|
||||
think through. Memory recall/remember is NEVER a dispatch task —
|
||||
call those tools directly.
|
||||
Examples:
|
||||
- User says "what's the latest news?" → CALL dispatch_task with
|
||||
description="Find the top 5 international news headlines today"
|
||||
- User says "look up the price of a PS5" → CALL dispatch_task with
|
||||
description="Find the current retail price of a PlayStation 5"
|
||||
After calling dispatch_task, tell the user "I'll get on that for you"
|
||||
and keep the conversation going. The result comes back automatically — when
|
||||
you receive it, share the findings naturally in your conversational style.
|
||||
You can have multiple tasks running at once.
|
||||
|
||||
# Skills
|
||||
You have a skill library of learned procedures. Use skill_recall when you need
|
||||
@@ -99,6 +110,7 @@ SYSTEM_PROMPT = textwrap.dedent("""\
|
||||
successfully complete a multi-step task or learn a new procedure from the user,
|
||||
use skill_save to record it so you can follow it next time. Be selective: only
|
||||
save skills for repeatable tasks, not one-off facts (those go in memory).
|
||||
|
||||
""")
|
||||
|
||||
|
||||
@@ -158,6 +170,9 @@ class _ReasoningFallbackWrapper:
|
||||
self._last_reasoning = reasoning
|
||||
if chunk.has_response():
|
||||
self._has_content = True
|
||||
tc = getattr(delta, "tool_calls", None)
|
||||
if tc:
|
||||
logger.info("LLM tool_call: %s", [getattr(t, "name", str(t)) for t in tc])
|
||||
return chunk
|
||||
|
||||
async def collect(self):
|
||||
@@ -167,17 +182,43 @@ class _ReasoningFallbackWrapper:
|
||||
class GemmaLLM(openai.LLM):
|
||||
"""OpenAI-compatible LLM (llama.cpp Gemma 4) with thinking disabled."""
|
||||
|
||||
# Context compaction: when the conversation grows past this many items,
|
||||
# ask Gemma to summarize the older messages into a compact recap, then
|
||||
# replace them with that summary. Keeps recent turns verbatim.
|
||||
COMPACTION_THRESHOLD = 24 # ~12 turns of user/agent exchange
|
||||
KEEP_RECENT = 10 # last N items kept verbatim after compaction
|
||||
IDLE_COMPACTION_SECONDS = 300 # also compact after 5 min of silence
|
||||
|
||||
def __init__(self, *args, **kwargs):
|
||||
super().__init__(*args, **kwargs)
|
||||
import time
|
||||
self._last_call_time: float = time.monotonic()
|
||||
|
||||
def chat(self, *, chat_ctx, tools=None, conn_options=None, **kwargs):
|
||||
if conn_options is None:
|
||||
from livekit.agents.types import DEFAULT_API_CONNECT_OPTIONS
|
||||
|
||||
conn_options = DEFAULT_API_CONNECT_OPTIONS
|
||||
|
||||
# Compact the context before sending: keep the system prompt and the
|
||||
# last N items so long conversations stay within Gemma's window.
|
||||
MAX_CONTEXT_ITEMS = 30
|
||||
if len(chat_ctx) > MAX_CONTEXT_ITEMS:
|
||||
chat_ctx.truncate(max_items=MAX_CONTEXT_ITEMS)
|
||||
import time
|
||||
now = time.monotonic()
|
||||
idle_for = now - self._last_call_time
|
||||
self._last_call_time = now
|
||||
|
||||
# Compact if context is large, OR if we've been idle for 5+ minutes
|
||||
# (user came back after a break — summarize the old conversation).
|
||||
ctx_len = len(chat_ctx.items)
|
||||
should_compact = ctx_len > self.COMPACTION_THRESHOLD
|
||||
if not should_compact and idle_for > self.IDLE_COMPACTION_SECONDS:
|
||||
should_compact = ctx_len > 6 # minimal context worth summarizing
|
||||
|
||||
if should_compact:
|
||||
reason = "idle" if idle_for > self.IDLE_COMPACTION_SECONDS else "size"
|
||||
logger.info(
|
||||
"Compacting context (%s): %d items, idle %.0fs",
|
||||
reason, ctx_len, idle_for,
|
||||
)
|
||||
chat_ctx = self._compact_context(chat_ctx, conn_options)
|
||||
|
||||
stream = super().chat(
|
||||
chat_ctx=chat_ctx,
|
||||
@@ -187,6 +228,118 @@ class GemmaLLM(openai.LLM):
|
||||
)
|
||||
return _ReasoningFallbackWrapper(stream)
|
||||
|
||||
def _compact_context(self, chat_ctx, conn_options):
|
||||
"""Summarize older messages via a side LLM call, keep recent ones.
|
||||
|
||||
Splits the context into [system | old | recent]. The old portion is
|
||||
summarized by Gemma into a short recap message that replaces it.
|
||||
Falls back to hard truncation if the summary call fails.
|
||||
"""
|
||||
from livekit.agents import llm as lk_llm
|
||||
|
||||
items = list(chat_ctx.items)
|
||||
# Identify the system/instruction message (first item usually)
|
||||
sys_items = [it for it in items if getattr(it, "role", None) in ("system", "developer")]
|
||||
non_sys = [it for it in items if getattr(it, "role", None) not in ("system", "developer")]
|
||||
|
||||
if len(non_sys) <= self.KEEP_RECENT:
|
||||
return chat_ctx # nothing to compact
|
||||
|
||||
old_msgs = non_sys[: -self.KEEP_RECENT]
|
||||
recent_msgs = non_sys[-self.KEEP_RECENT :]
|
||||
|
||||
# Build a transcript of the old messages for the summarizer
|
||||
transcript_parts = []
|
||||
for msg in old_msgs:
|
||||
role = getattr(msg, "role", "user")
|
||||
content = self._msg_content_text(msg)
|
||||
if content:
|
||||
transcript_parts.append(f"{role}: {content}")
|
||||
transcript = "\n".join(transcript_parts)
|
||||
|
||||
if not transcript.strip():
|
||||
return chat_ctx # nothing meaningful to summarize
|
||||
|
||||
# Side call to Gemma for the summary (no tools, no streaming needed)
|
||||
try:
|
||||
import httpx
|
||||
|
||||
summary_prompt = (
|
||||
"Summarize this conversation history in 2-4 sentences. "
|
||||
"Preserve key facts, decisions, and user preferences. "
|
||||
"Be concise — this replaces the full history as context.\n\n"
|
||||
+ transcript
|
||||
)
|
||||
resp = httpx.post(
|
||||
f"{GEMMA_BASE_URL}/chat/completions",
|
||||
headers={"Content-Type": "application/json"},
|
||||
json={
|
||||
"model": GEMMA_MODEL,
|
||||
"messages": [
|
||||
{"role": "system", "content": "You summarize conversations concisely."},
|
||||
{"role": "user", "content": summary_prompt},
|
||||
],
|
||||
"max_tokens": 200,
|
||||
"temperature": 0.3,
|
||||
"chat_template_kwargs": {"enable_thinking": False},
|
||||
},
|
||||
timeout=30,
|
||||
)
|
||||
resp.raise_for_status()
|
||||
data = resp.json()
|
||||
summary_text = (data.get("choices") or [{}])[0].get("message", {}).get("content", "")
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("Context compaction failed, falling back to truncation: %s", e)
|
||||
chat_ctx.truncate(max_items=self.COMPACTION_THRESHOLD)
|
||||
return chat_ctx
|
||||
|
||||
if not summary_text.strip():
|
||||
chat_ctx.truncate(max_items=self.COMPACTION_THRESHOLD)
|
||||
return chat_ctx
|
||||
|
||||
# Rebuild the context: system + summary + recent
|
||||
new_items = list(sys_items)
|
||||
new_items.append(
|
||||
lk_llm.ChatMessage(
|
||||
role="user",
|
||||
content=[f"[Conversation so far: {summary_text.strip()}]"],
|
||||
)
|
||||
)
|
||||
new_items.append(
|
||||
lk_llm.ChatMessage(
|
||||
role="assistant",
|
||||
content=["Got it, I'll keep that in mind."],
|
||||
)
|
||||
)
|
||||
new_items.extend(recent_msgs)
|
||||
|
||||
logger.info(
|
||||
"Context compacted: %d items -> %d (summary: %.80s...)",
|
||||
len(items),
|
||||
len(new_items),
|
||||
summary_text.strip(),
|
||||
)
|
||||
|
||||
return lk_llm.ChatContext(new_items)
|
||||
|
||||
@staticmethod
|
||||
def _msg_content_text(msg) -> str:
|
||||
"""Extract plain text from a ChatMessage's content field."""
|
||||
content = getattr(msg, "content", None)
|
||||
if content is None:
|
||||
return ""
|
||||
if isinstance(content, str):
|
||||
return content
|
||||
if isinstance(content, list):
|
||||
parts = []
|
||||
for c in content:
|
||||
if isinstance(c, str):
|
||||
parts.append(c)
|
||||
elif hasattr(c, "text"):
|
||||
parts.append(c.text)
|
||||
return " ".join(parts)
|
||||
return str(content)
|
||||
|
||||
|
||||
# ── MCP toolsets (web access + any extra configured servers) ────────────────
|
||||
|
||||
@@ -199,22 +352,6 @@ def build_mcp_toolsets() -> list[mcp.MCPToolset]:
|
||||
"""
|
||||
toolsets: list[mcp.MCPToolset] = []
|
||||
|
||||
if WEB_MCP_ENABLED:
|
||||
python_bin = os.path.join(os.path.dirname(os.path.abspath(__file__)), ".venv", "bin", "python")
|
||||
web_mcp_script = os.path.join(os.path.dirname(os.path.abspath(__file__)), "web_mcp.py")
|
||||
toolsets.append(
|
||||
mcp.MCPToolset(
|
||||
id="web-access",
|
||||
mcp_server=mcp.MCPServerStdio(
|
||||
command=python_bin,
|
||||
args=[web_mcp_script],
|
||||
env={**os.environ, "FIRECRAWL_BASE": FIRECRAWL_BASE},
|
||||
client_session_timeout_seconds=120,
|
||||
),
|
||||
)
|
||||
)
|
||||
logger.info("Web-access MCP toolset enabled (Firecrawl at %s)", FIRECRAWL_BASE)
|
||||
|
||||
python_bin = os.path.join(os.path.dirname(os.path.abspath(__file__)), ".venv", "bin", "python")
|
||||
weather_mcp_script = os.path.join(os.path.dirname(os.path.abspath(__file__)), "weather_mcp.py")
|
||||
toolsets.append(
|
||||
@@ -230,19 +367,19 @@ def build_mcp_toolsets() -> list[mcp.MCPToolset]:
|
||||
)
|
||||
logger.info("Weather MCP toolset enabled (wttr.in)")
|
||||
|
||||
memory_mcp_script = os.path.join(os.path.dirname(os.path.abspath(__file__)), "memory_mcp.py")
|
||||
cognee_url = os.environ.get("COGNEE_MCP_URL", "http://192.168.86.2:8003/mcp")
|
||||
toolsets.append(
|
||||
mcp.MCPToolset(
|
||||
id="memory",
|
||||
mcp_server=mcp.MCPServerStdio(
|
||||
command=python_bin,
|
||||
args=[memory_mcp_script],
|
||||
env={**os.environ, "MEMORY_DIR": "/memory"},
|
||||
client_session_timeout_seconds=30,
|
||||
mcp_server=mcp.MCPServerHTTP(
|
||||
url=cognee_url,
|
||||
transport_type="streamable_http",
|
||||
allowed_tools=["remember", "recall", "forget"],
|
||||
client_session_timeout_seconds=60,
|
||||
),
|
||||
)
|
||||
)
|
||||
logger.info("Memory MCP toolset enabled (%s)", MEMORY_DIR)
|
||||
logger.info("Memory MCP toolset enabled (Cognee at %s)", cognee_url)
|
||||
|
||||
skills_mcp_script = os.path.join(os.path.dirname(os.path.abspath(__file__)), "skills_mcp.py")
|
||||
toolsets.append(
|
||||
@@ -258,6 +395,20 @@ def build_mcp_toolsets() -> list[mcp.MCPToolset]:
|
||||
)
|
||||
logger.info("Skills MCP toolset enabled (%s)", SKILLS_DIR)
|
||||
|
||||
dispatch_mcp_script = os.path.join(os.path.dirname(os.path.abspath(__file__)), "dispatch_mcp.py")
|
||||
toolsets.append(
|
||||
mcp.MCPToolset(
|
||||
id="dispatch",
|
||||
mcp_server=mcp.MCPServerStdio(
|
||||
command=python_bin,
|
||||
args=[dispatch_mcp_script],
|
||||
env={**os.environ},
|
||||
client_session_timeout_seconds=60,
|
||||
),
|
||||
)
|
||||
)
|
||||
logger.info("Dispatch MCP toolset enabled")
|
||||
|
||||
extra = os.environ.get("EXTRA_MCP_SERVERS", "")
|
||||
if extra:
|
||||
try:
|
||||
@@ -394,6 +545,30 @@ async def handle_job(ctx: JobContext) -> None:
|
||||
# audio is arriving from the participant.
|
||||
logger.info("user state -> %s", ev.new_state)
|
||||
|
||||
# ── Background task events → room data channel ─────────────────────────
|
||||
try:
|
||||
from agent.task_registry import registry as task_registry
|
||||
except ImportError:
|
||||
from task_registry import registry as task_registry
|
||||
|
||||
async def _publish_task_event(task_id: str, event: dict) -> None:
|
||||
payload = json.dumps({"type": "task_event", "task_id": task_id, **event})
|
||||
try:
|
||||
await ctx.room.local_participant.publish_data(
|
||||
payload, reliable=True, topic="tasks"
|
||||
)
|
||||
except Exception as e: # noqa: BLE001
|
||||
logger.warning("failed to publish task event: %s", e)
|
||||
|
||||
async def _poll_tasks() -> None:
|
||||
while True:
|
||||
await asyncio.sleep(1.0)
|
||||
events = await task_registry.poll_events()
|
||||
for task_id, event in events:
|
||||
await _publish_task_event(task_id, event)
|
||||
|
||||
asyncio.create_task(_poll_tasks())
|
||||
|
||||
await session.start(
|
||||
agent=VoiceAssistant(),
|
||||
room=ctx.room,
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
"""Dispatch MCP server — lets Hope spawn background tasks.
|
||||
|
||||
The main agent calls dispatch_task(description) when the user asks for
|
||||
something long-running. The tool returns immediately with a task ID; the
|
||||
actual work runs in a background asyncio task via TaskWorker.
|
||||
|
||||
Communication with the main agent process happens via /tmp/tasks/events.jsonl
|
||||
(the file-based registry).
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
from mcp.server.fastmcp import FastMCP
|
||||
|
||||
logger = logging.getLogger("voice-agent.dispatch")
|
||||
|
||||
mcp = FastMCP("dispatch")
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
async def dispatch_task(description: str) -> str:
|
||||
"""Spawn a background task to research or look something up.
|
||||
|
||||
Use this when the user asks you to do something that will take more than
|
||||
a few seconds (research, look up news, multi-step investigation). The task
|
||||
runs in parallel while you continue the conversation. When it completes,
|
||||
the result will be delivered back to you automatically — you can tell the
|
||||
user "I'll let you know when I find out."
|
||||
|
||||
Returns a short confirmation with the task ID.
|
||||
"""
|
||||
try:
|
||||
from agent.task_registry import registry
|
||||
from agent.task_worker import TaskWorker
|
||||
except ImportError:
|
||||
from task_registry import registry
|
||||
from task_worker import TaskWorker
|
||||
|
||||
task = registry.create(description)
|
||||
worker = TaskWorker(task.id, description)
|
||||
asyncio.create_task(_run_worker(worker, task.id))
|
||||
logger.info("Dispatched task %s: %s", task.id, description)
|
||||
return f"Task {task.id} started. I'll report back when it's done."
|
||||
|
||||
|
||||
async def _run_worker(worker: "TaskWorker", task_id: str):
|
||||
try:
|
||||
from agent.task_registry import registry
|
||||
except ImportError:
|
||||
from task_registry import registry
|
||||
|
||||
try:
|
||||
await worker.run()
|
||||
except Exception as e:
|
||||
logger.exception("Task %s crashed: %s", task_id, e)
|
||||
registry.fail(task_id, str(e))
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
mcp.run(transport="stdio")
|
||||
@@ -1,163 +0,0 @@
|
||||
"""Memory MCP server — persistent markdown-based memory for Hope.
|
||||
|
||||
Runs over stdio inside the voice container. The agent attaches it via
|
||||
MCPServerStdio, so the LLM can recall and save memories during a
|
||||
conversation. Memories live as plain markdown files in MEMORY_DIR
|
||||
(default /memory), one file per topic, mounted from the host so they
|
||||
survive rebuilds and are backed up via git.
|
||||
|
||||
Tools:
|
||||
- memory_recall(query) -> relevant passages from all .md files
|
||||
- memory_save(topic, content) -> append a timestamped note to {topic}.md
|
||||
- memory_list() -> topics with their first line
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import re
|
||||
from datetime import datetime, timezone
|
||||
|
||||
from mcp.server.fastmcp import FastMCP
|
||||
|
||||
MEMORY_DIR = os.environ.get("MEMORY_DIR", "/memory")
|
||||
|
||||
mcp = FastMCP("memory")
|
||||
|
||||
|
||||
def _slugify(topic: str) -> str:
|
||||
"""Turn a topic into a safe filename slug."""
|
||||
slug = re.sub(r"[^a-z0-9]+", "-", topic.lower()).strip("-")
|
||||
return slug or "misc"
|
||||
|
||||
|
||||
def _read_files() -> list[tuple[str, str]]:
|
||||
"""Return (filename, content) for every .md file in MEMORY_DIR."""
|
||||
files: list[tuple[str, str]] = []
|
||||
if not os.path.isdir(MEMORY_DIR):
|
||||
return files
|
||||
for name in sorted(os.listdir(MEMORY_DIR)):
|
||||
if not name.endswith(".md"):
|
||||
continue
|
||||
path = os.path.join(MEMORY_DIR, name)
|
||||
try:
|
||||
with open(path, encoding="utf-8") as f:
|
||||
files.append((name, f.read()))
|
||||
except OSError:
|
||||
continue
|
||||
return files
|
||||
|
||||
|
||||
def _first_line(content: str) -> str:
|
||||
for line in content.splitlines():
|
||||
if line.strip():
|
||||
return line.strip()
|
||||
return "(empty)"
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def memory_recall(query: str) -> str:
|
||||
"""Search saved memories for information relevant to a query.
|
||||
|
||||
Use this at the start of a conversation, or whenever the user
|
||||
references past information ("what did I tell you about...",
|
||||
"remember when..."). Returns matching passages with their topic
|
||||
names, or a message saying nothing was found.
|
||||
"""
|
||||
q = (query or "").strip().lower()
|
||||
if not q:
|
||||
return "I don't have any memory of that."
|
||||
|
||||
words = [w for w in re.split(r"\W+", q) if len(w) > 2]
|
||||
files = _read_files()
|
||||
if not files:
|
||||
return "I don't have any memory of that."
|
||||
|
||||
# Strip markdown heading markers and timestamp lines so snippets are clean prose.
|
||||
def _clean(text: str) -> str:
|
||||
text = text.replace("\n", " ")
|
||||
text = re.sub(r"#+\s*", "", text)
|
||||
text = re.sub(r"\d{4}-\d{2}-\d{2} \d{2}:\d{2} UTC\s*", "", text)
|
||||
return re.sub(r"\s{2,}", " ", text).strip()
|
||||
|
||||
scored: list[tuple[int, str, list[str]]] = []
|
||||
for name, content in files:
|
||||
lower = content.lower()
|
||||
score = 0
|
||||
matches: list[str] = []
|
||||
if q in lower:
|
||||
score += len(q)
|
||||
seen: set[int] = set()
|
||||
for w in words:
|
||||
start = 0
|
||||
while True:
|
||||
idx = lower.find(w, start)
|
||||
if idx == -1:
|
||||
break
|
||||
score += 1
|
||||
key = idx // 200 # one snippet per ~200-char window
|
||||
if key not in seen:
|
||||
seen.add(key)
|
||||
s = max(0, idx - 80)
|
||||
e = min(len(content), idx + len(w) + 120)
|
||||
snippet = _clean(content[s:e])
|
||||
if snippet and snippet not in matches:
|
||||
matches.append(snippet)
|
||||
start = idx + len(w)
|
||||
if score > 0:
|
||||
scored.append((score, name, matches[:3]))
|
||||
|
||||
if not scored:
|
||||
return "I don't have any memory of that."
|
||||
|
||||
scored.sort(key=lambda s: s[0], reverse=True)
|
||||
parts = []
|
||||
for _, name, snippets in scored[:5]:
|
||||
topic = name.removesuffix(".md")
|
||||
parts.append(f"[{topic}] " + " ".join(snippets))
|
||||
return "\n".join(parts)
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def memory_save(topic: str, content: str) -> str:
|
||||
"""Save a fact or preference to persistent memory.
|
||||
|
||||
Use this when the user shares personal information worth keeping:
|
||||
names, birthdays, preferences, projects, important facts. Creates
|
||||
or appends to a markdown file named after the topic. Returns a
|
||||
short confirmation.
|
||||
"""
|
||||
topic = (topic or "").strip()
|
||||
content = (content or "").strip()
|
||||
if not topic or not content:
|
||||
return "Nothing saved — both a topic and some content are needed."
|
||||
|
||||
os.makedirs(MEMORY_DIR, exist_ok=True)
|
||||
path = os.path.join(MEMORY_DIR, f"{_slugify(topic)}.md")
|
||||
|
||||
is_new = not os.path.exists(path)
|
||||
stamp = datetime.now(timezone.utc).strftime("%Y-%m-%d %H:%M UTC")
|
||||
with open(path, "a", encoding="utf-8") as f:
|
||||
if is_new:
|
||||
f.write(f"# {topic}\n\n")
|
||||
f.write(f"## {stamp}\n\n{content}\n\n")
|
||||
|
||||
return f"Saved to memory under '{topic}'."
|
||||
|
||||
|
||||
@mcp.tool()
|
||||
def memory_list() -> str:
|
||||
"""List all saved memories with a brief description of each.
|
||||
|
||||
Returns the topic name and first line of every memory file, or a
|
||||
message saying no memories exist yet.
|
||||
"""
|
||||
files = _read_files()
|
||||
if not files:
|
||||
return "No memories saved yet."
|
||||
lines = [f"{name.removesuffix('.md')}: {_first_line(content)}" for name, content in files]
|
||||
return "\n".join(lines)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
mcp.run(transport="stdio")
|
||||
@@ -0,0 +1,164 @@
|
||||
"""File-based task registry for background dispatch tasks.
|
||||
|
||||
The dispatch MCP server runs in a separate process, so we use a JSONL event
|
||||
file as the communication channel. The main agent process tails this file
|
||||
and pushes events to the room (UI) and speaks results.
|
||||
|
||||
Event file: /tmp/tasks/events.jsonl
|
||||
Task state: /tmp/tasks/{task_id}.json
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import time
|
||||
import uuid
|
||||
from dataclasses import dataclass, field
|
||||
from enum import Enum
|
||||
from pathlib import Path
|
||||
from typing import Any, Callable
|
||||
|
||||
|
||||
TASKS_DIR = Path(os.environ.get("TASKS_DIR", "/tmp/tasks"))
|
||||
EVENTS_FILE = TASKS_DIR / "events.jsonl"
|
||||
|
||||
|
||||
class TaskStatus(str, Enum):
|
||||
PENDING = "pending"
|
||||
RUNNING = "running"
|
||||
COMPLETED = "completed"
|
||||
FAILED = "failed"
|
||||
|
||||
|
||||
@dataclass
|
||||
class TaskStep:
|
||||
role: str
|
||||
content: str
|
||||
timestamp: float = field(default_factory=time.time)
|
||||
|
||||
|
||||
@dataclass
|
||||
class Task:
|
||||
id: str
|
||||
description: str
|
||||
status: TaskStatus = TaskStatus.PENDING
|
||||
steps: list[TaskStep] = field(default_factory=list)
|
||||
result: str | None = None
|
||||
error: str | None = None
|
||||
created_at: float = field(default_factory=time.time)
|
||||
completed_at: float | None = None
|
||||
|
||||
|
||||
class TaskRegistry:
|
||||
"""File-backed registry. Writers (MCP process) append to events.jsonl.
|
||||
Readers (main agent) tail the file and dispatch to subscribers."""
|
||||
|
||||
def __init__(self):
|
||||
TASKS_DIR.mkdir(parents=True, exist_ok=True)
|
||||
if not EVENTS_FILE.exists():
|
||||
EVENTS_FILE.touch()
|
||||
self._tasks: dict[str, Task] = {}
|
||||
self._listeners: list[Callable[[str, dict], Any]] = []
|
||||
self._tail_pos = 0
|
||||
|
||||
def create(self, description: str) -> Task:
|
||||
task = Task(id=uuid.uuid4().hex[:8], description=description)
|
||||
self._tasks[task.id] = task
|
||||
self._write_event(task.id, {"event": "created", "task": self._task_dict(task)})
|
||||
return task
|
||||
|
||||
def add_step(self, task_id: str, role: str, content: str):
|
||||
task = self._tasks.get(task_id)
|
||||
if task:
|
||||
task.steps.append(TaskStep(role=role, content=content))
|
||||
if task.status == TaskStatus.PENDING:
|
||||
task.status = TaskStatus.RUNNING
|
||||
self._save_task(task)
|
||||
self._write_event(task_id, {"event": "step", "role": role, "content": content})
|
||||
|
||||
def complete(self, task_id: str, result: str):
|
||||
task = self._tasks.get(task_id)
|
||||
if task:
|
||||
task.status = TaskStatus.COMPLETED
|
||||
task.result = result
|
||||
task.completed_at = time.time()
|
||||
self._save_task(task)
|
||||
self._write_event(task_id, {"event": "completed", "result": result})
|
||||
|
||||
def fail(self, task_id: str, error: str):
|
||||
task = self._tasks.get(task_id)
|
||||
if task:
|
||||
task.status = TaskStatus.FAILED
|
||||
task.error = error
|
||||
task.completed_at = time.time()
|
||||
self._save_task(task)
|
||||
self._write_event(task_id, {"event": "failed", "error": error})
|
||||
|
||||
def get(self, task_id: str) -> dict | None:
|
||||
t = self._tasks.get(task_id)
|
||||
return self._task_dict(t) if t else None
|
||||
|
||||
def list_all(self) -> list[dict]:
|
||||
return [self._task_dict(t) for t in self._tasks.values()]
|
||||
|
||||
def subscribe(self, callback: Callable[[str, dict], Any]):
|
||||
self._listeners.append(callback)
|
||||
|
||||
async def poll_events(self) -> list[tuple[str, dict]]:
|
||||
"""Read new events from the file. Call periodically from the agent."""
|
||||
events = []
|
||||
try:
|
||||
with open(EVENTS_FILE, "r") as f:
|
||||
f.seek(self._tail_pos)
|
||||
for line in f:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
ev = json.loads(line)
|
||||
events.append((ev.get("task_id", ""), ev))
|
||||
except json.JSONDecodeError:
|
||||
pass
|
||||
self._tail_pos = f.tell()
|
||||
except (OSError, IOError):
|
||||
pass
|
||||
return events
|
||||
|
||||
async def notify_subscribers(self, task_id: str, event: dict):
|
||||
for cb in self._listeners:
|
||||
try:
|
||||
result = cb(task_id, event)
|
||||
if asyncio.iscoroutine(result):
|
||||
await result
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
def _write_event(self, task_id: str, event: dict):
|
||||
payload = json.dumps({"task_id": task_id, **event})
|
||||
with open(EVENTS_FILE, "a") as f:
|
||||
f.write(payload + "\n")
|
||||
|
||||
def _save_task(self, task: Task):
|
||||
path = TASKS_DIR / f"{task.id}.json"
|
||||
path.write_text(json.dumps(self._task_dict(task)))
|
||||
|
||||
@staticmethod
|
||||
def _task_dict(task: Task) -> dict:
|
||||
return {
|
||||
"id": task.id,
|
||||
"description": task.description,
|
||||
"status": task.status.value,
|
||||
"steps": [
|
||||
{"role": s.role, "content": s.content, "ts": s.timestamp}
|
||||
for s in task.steps
|
||||
],
|
||||
"result": task.result,
|
||||
"error": task.error,
|
||||
"created_at": task.created_at,
|
||||
"completed_at": task.completed_at,
|
||||
}
|
||||
|
||||
|
||||
# Global singleton — one per process
|
||||
registry = TaskRegistry()
|
||||
@@ -0,0 +1,243 @@
|
||||
"""Background task worker — runs an autonomous LLM loop with tools.
|
||||
|
||||
Each dispatched task gets its own TaskWorker instance running in a separate
|
||||
asyncio task. The worker calls Gemma iteratively, executing tool calls and
|
||||
logging each step to the registry so the UI can display live progress.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
from datetime import datetime, timezone
|
||||
|
||||
logger = logging.getLogger("voice-agent.tasks")
|
||||
|
||||
|
||||
class TaskWorker:
|
||||
"""Runs a background research/lookup task using the LLM + tools."""
|
||||
|
||||
def __init__(self, task_id: str, description: str):
|
||||
self.task_id = task_id
|
||||
self.description = description
|
||||
self._cancelled = False
|
||||
|
||||
async def run(self):
|
||||
from openai import AsyncOpenAI
|
||||
try:
|
||||
from agent.task_registry import registry as reg
|
||||
except ImportError:
|
||||
from task_registry import registry as reg
|
||||
|
||||
base_url = os.environ.get("GEMMA_BASE_URL", "http://192.168.86.2:8023/v1")
|
||||
model = os.environ.get("GEMMA_MODEL", "gemma-4-e4b")
|
||||
api_key = os.environ.get("GEMMA_API_KEY", "not-needed")
|
||||
|
||||
reg.add_step(self.task_id, "thinking", f"Starting: {self.description}")
|
||||
|
||||
client = AsyncOpenAI(base_url=base_url, api_key=api_key)
|
||||
tools = self._load_tools()
|
||||
|
||||
messages = [
|
||||
{"role": "system", "content": self._system_prompt()},
|
||||
{"role": "user", "content": self.description},
|
||||
]
|
||||
|
||||
max_iterations = 10
|
||||
for iteration in range(max_iterations):
|
||||
if self._cancelled:
|
||||
reg.fail(self.task_id, "Cancelled")
|
||||
return
|
||||
|
||||
reg.add_step(
|
||||
self.task_id, "thinking", f"Step {iteration + 1}: analyzing..."
|
||||
)
|
||||
|
||||
try:
|
||||
response = await client.chat.completions.create(
|
||||
model=model,
|
||||
messages=messages,
|
||||
tools=tools or None,
|
||||
max_tokens=500,
|
||||
temperature=0.7,
|
||||
extra_body={"chat_template_kwargs": {"enable_thinking": False}},
|
||||
)
|
||||
except Exception as e:
|
||||
reg.fail(self.task_id, f"LLM error: {e}")
|
||||
return
|
||||
|
||||
msg = response.choices[0].message
|
||||
|
||||
if msg.tool_calls:
|
||||
messages.append(msg)
|
||||
for tc in msg.tool_calls:
|
||||
tool_name = tc.function.name
|
||||
try:
|
||||
tool_args = json.loads(tc.function.arguments)
|
||||
except json.JSONDecodeError:
|
||||
tool_args = {}
|
||||
reg.add_step(
|
||||
self.task_id,
|
||||
"tool_call",
|
||||
f"{tool_name}({json.dumps(tool_args)})",
|
||||
)
|
||||
result = await self._execute_tool(tool_name, tool_args)
|
||||
reg.add_step(self.task_id, "tool_result", str(result)[:500])
|
||||
messages.append(
|
||||
{
|
||||
"role": "tool",
|
||||
"tool_call_id": tc.id,
|
||||
"content": str(result),
|
||||
}
|
||||
)
|
||||
else:
|
||||
text = (msg.content or "").strip()
|
||||
if not text:
|
||||
reg.fail(self.task_id, "LLM returned empty response")
|
||||
return
|
||||
reg.add_step(self.task_id, "text", text)
|
||||
reg.complete(self.task_id, text)
|
||||
return
|
||||
|
||||
reg.fail(self.task_id, f"Reached max iterations ({max_iterations})")
|
||||
|
||||
def _system_prompt(self) -> str:
|
||||
return (
|
||||
"You are Hope's research assistant working on a background task. "
|
||||
"Use your tools to gather information thoroughly. When you have "
|
||||
"enough, provide a final answer in 2-4 sentences suitable for "
|
||||
"speaking aloud. Plain text only, no markdown or formatting."
|
||||
)
|
||||
|
||||
def _load_tools(self) -> list[dict]:
|
||||
return [
|
||||
{
|
||||
"type": "function",
|
||||
"function": {
|
||||
"name": "get_weather",
|
||||
"description": "Get current weather and today's forecast for a location in Fahrenheit.",
|
||||
"parameters": {
|
||||
"type": "object",
|
||||
"properties": {"location": {"type": "string"}},
|
||||
"required": ["location"],
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
"type": "function",
|
||||
"function": {
|
||||
"name": "get_time",
|
||||
"description": "Get the current date and time, optionally for a specific location.",
|
||||
"parameters": {
|
||||
"type": "object",
|
||||
"properties": {"location": {"type": "string"}},
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
"type": "function",
|
||||
"function": {
|
||||
"name": "memory_recall",
|
||||
"description": "Search saved memories for relevant information.",
|
||||
"parameters": {
|
||||
"type": "object",
|
||||
"properties": {"query": {"type": "string"}},
|
||||
"required": ["query"],
|
||||
},
|
||||
},
|
||||
},
|
||||
{
|
||||
"type": "function",
|
||||
"function": {
|
||||
"name": "web_search",
|
||||
"description": "Search the web for current information using Firecrawl.",
|
||||
"parameters": {
|
||||
"type": "object",
|
||||
"properties": {"query": {"type": "string"}},
|
||||
"required": ["query"],
|
||||
},
|
||||
},
|
||||
},
|
||||
]
|
||||
|
||||
async def _execute_tool(self, name: str, args: dict) -> str:
|
||||
import httpx
|
||||
|
||||
try:
|
||||
if name == "get_weather":
|
||||
location = args.get("location", "")
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
resp = await client.get(
|
||||
f"https://wttr.in/{location}?format=j1",
|
||||
headers={"User-Agent": "curl/8.0"},
|
||||
)
|
||||
data = resp.json()
|
||||
current = data.get("current_condition", [{}])[0] or {}
|
||||
temp_c = current.get("temp_C")
|
||||
if temp_c is None:
|
||||
return f"Could not get weather for {location}."
|
||||
temp_f = round(int(temp_c) * 9 / 5 + 32)
|
||||
desc = (current.get("weatherDesc", [{}])[0].get("value") or "unknown").lower()
|
||||
feels_c = current.get("FeelsLikeC")
|
||||
feels_f = round(int(feels_c) * 9 / 5 + 32) if feels_c else None
|
||||
today = (data.get("weather") or [{}])[0] or {}
|
||||
high_f = round(int(today["maxtempC"]) * 9 / 5 + 32) if today.get("maxtempC") else None
|
||||
low_f = round(int(today["mintempC"]) * 9 / 5 + 32) if today.get("mintempC") else None
|
||||
parts = [f"It's {temp_f} degrees and {desc} in {location}"]
|
||||
if feels_f and feels_f != temp_f:
|
||||
parts[0] += f", feels like {feels_f}"
|
||||
if high_f and low_f:
|
||||
parts.append(f"High {high_f}, low {low_f}.")
|
||||
return " ".join(parts)
|
||||
|
||||
elif name == "get_time":
|
||||
now = datetime.now(timezone.utc)
|
||||
loc = args.get("location", "")
|
||||
suffix = f" in {loc}" if loc else ", UTC"
|
||||
return f"It's {now.strftime('%I:%M %p')} on {now.strftime('%A, %B %d')}{suffix}."
|
||||
|
||||
elif name == "memory_recall":
|
||||
query = args.get("query", "").lower()
|
||||
memory_dir = os.environ.get("MEMORY_DIR", "/memory")
|
||||
results = []
|
||||
try:
|
||||
for fname in os.listdir(memory_dir):
|
||||
if not fname.endswith(".md"):
|
||||
continue
|
||||
with open(os.path.join(memory_dir, fname)) as f:
|
||||
content = f.read()
|
||||
if any(w in content.lower() for w in query.split()):
|
||||
results.append(f"[{fname}] {content[:300]}")
|
||||
except OSError:
|
||||
pass
|
||||
return "\n".join(results) or "No memories found."
|
||||
|
||||
elif name == "web_search":
|
||||
query = args.get("query", "")
|
||||
firecrawl_base = os.environ.get("FIRECRAWL_BASE", "")
|
||||
if not firecrawl_base:
|
||||
return "Web search is not configured."
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
resp = await client.post(
|
||||
f"{firecrawl_base}/v1/search",
|
||||
json={"query": query, "limit": 3},
|
||||
)
|
||||
if resp.status_code == 200:
|
||||
data = resp.json()
|
||||
items = data.get("data", [])
|
||||
lines = [
|
||||
f"{r.get('title', '')}: {r.get('description', '')}"
|
||||
for r in items[:3]
|
||||
]
|
||||
return "\n".join(lines) or "No results found."
|
||||
return f"Search failed (HTTP {resp.status_code})."
|
||||
|
||||
else:
|
||||
return f"Unknown tool: {name}"
|
||||
|
||||
except Exception as e:
|
||||
return f"Tool error: {e}"
|
||||
|
||||
def cancel(self):
|
||||
self._cancelled = True
|
||||
+1
-1
@@ -15,10 +15,10 @@ services:
|
||||
GEMMA_API_KEY: "${GEMMA_API_KEY:-not-needed}"
|
||||
WEB_MCP_ENABLED: "${WEB_MCP_ENABLED:-true}"
|
||||
FIRECRAWL_BASE: "${FIRECRAWL_BASE:-http://192.168.86.2:3002}"
|
||||
COGNEE_MCP_URL: "${COGNEE_MCP_URL:-http://192.168.86.2:8003/mcp}"
|
||||
volumes:
|
||||
- ./livekit.yaml:/etc/livekit.yaml:ro
|
||||
- ./certs:/etc/voice/certs
|
||||
- ./memory:/memory
|
||||
- ./skills:/skills
|
||||
healthcheck:
|
||||
test: ["CMD-SHELL", "curl -sk https://localhost:8090/ -o /dev/null && curl -s http://localhost:7880/ -o /dev/null"]
|
||||
|
||||
@@ -6,3 +6,6 @@ rtc:
|
||||
|
||||
keys:
|
||||
devkey: devsecret
|
||||
|
||||
# Rooms are short-lived: each browser session gets a unique room name, so
|
||||
# the agent is dispatched at room creation. Default 300s cleanup is fine.
|
||||
|
||||
@@ -0,0 +1,38 @@
|
||||
# Hope test suite
|
||||
|
||||
End-to-end tests that drive the **actual running container** — no mocks.
|
||||
They run on the host and hit the live web UI, token endpoint, LLM API, MCP
|
||||
tool servers, and agent process.
|
||||
|
||||
## Prerequisites
|
||||
|
||||
- The container is up: `docker compose up -d` (and healthy)
|
||||
- Host deps installed: `pip install -r tests/requirements.txt`
|
||||
(or `uv pip install -r tests/requirements.txt`)
|
||||
- The LLM at `GEMMA_BASE_URL` (default `http://192.168.86.2:8023/v1`) is reachable
|
||||
- Internet access for the weather tool tests (wttr.in)
|
||||
|
||||
## Run
|
||||
|
||||
```bash
|
||||
cd /home/shane/dev/voice
|
||||
python -m pytest tests/ -v # everything
|
||||
python -m pytest tests/ -v -m "not slow" # skip LLM/network-heavy tests
|
||||
```
|
||||
|
||||
## What each file covers
|
||||
|
||||
| File | Target | Notes |
|
||||
|------|--------|-------|
|
||||
| `test_web_ui.py` | nginx HTTPS UI (:8090) + token endpoint (:8091) | static assets, `/token` (direct + via nginx), `/livekit/` proxy |
|
||||
| `test_llm_api.py` | Gemma LLM API directly | completion, tool calling, `enable_thinking=false` latency, streaming |
|
||||
| `test_mcp_tools.py` | MCP servers via `docker exec` | weather (Fahrenheit check), time, memory save/recall/list, skills save/recall/list — probe files are cleaned up |
|
||||
| `test_agent_integration.py` | running container | agent process alive, worker registered with LiveKit, all 4 supervised processes up, container health, LiveKit HTTP API |
|
||||
| `test_compaction.py` | `GemmaLLM._compact_context` in-container | compaction triggers >24 items, system prompt preserved, short context untouched |
|
||||
|
||||
## Notes
|
||||
|
||||
- LLM and MCP tests are marked `@pytest.mark.slow` (1–30s each).
|
||||
- The token endpoint is **POST** `/token` → `{"token": "<JWT>"}`.
|
||||
- Process checks use `/proc/*/cmdline` because the container image has no
|
||||
`ps`, and the supervisorctl unix socket is not exposed in this build.
|
||||
@@ -0,0 +1,176 @@
|
||||
"""Shared fixtures for the Hope voice-assistant test suite.
|
||||
|
||||
Tests run on the HOST machine against the live "voice" container:
|
||||
- HTTPS web UI -> https://localhost:8090 (self-signed cert, verify=False)
|
||||
- Token server -> http://127.0.0.1:8091 (host network mode)
|
||||
- LLM API -> GEMMA_BASE_URL (default http://192.168.86.2:8023/v1)
|
||||
- In-container -> docker exec voice ...
|
||||
"""
|
||||
|
||||
import json
|
||||
import os
|
||||
import subprocess
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
CONTAINER = "voice"
|
||||
WEB_BASE_URL = os.environ.get("VOICE_WEB_URL", "https://localhost:8090")
|
||||
TOKEN_URL = os.environ.get("VOICE_TOKEN_URL", "http://127.0.0.1:8091/token")
|
||||
GEMMA_BASE_URL = os.environ.get("GEMMA_BASE_URL", "http://192.168.86.2:8023/v1").rstrip("/")
|
||||
GEMMA_MODEL = os.environ.get("GEMMA_MODEL", "gemma-4-e4b")
|
||||
COGNEE_MCP_URL = os.environ.get("COGNEE_MCP_URL", "http://192.168.86.2:8003/mcp")
|
||||
|
||||
# Path to the agent's venv python inside the container (has mcp, livekit.agents)
|
||||
CONTAINER_PYTHON = "/opt/voice-agent/.venv/bin/python"
|
||||
AGENT_DIR_IN_CONTAINER = "/opt/voice-agent"
|
||||
|
||||
|
||||
def pytest_configure(config):
|
||||
config.addinivalue_line("markers", "slow: marks tests as slow (LLM / network calls)")
|
||||
|
||||
|
||||
# ── docker helpers ───────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def docker_exec(*args, timeout=60, env=None, check=True) -> subprocess.CompletedProcess:
|
||||
"""Run a command inside the voice container. Returns CompletedProcess."""
|
||||
cmd = ["docker", "exec"]
|
||||
if env:
|
||||
for k, v in env.items():
|
||||
cmd += ["-e", f"{k}={v}"]
|
||||
cmd += [CONTAINER, *args]
|
||||
return subprocess.run(
|
||||
cmd,
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=timeout,
|
||||
check=check,
|
||||
)
|
||||
|
||||
|
||||
def docker_exec_python(
|
||||
code: str, *, timeout=120, cwd=None, env=None, check=True
|
||||
) -> subprocess.CompletedProcess:
|
||||
"""Run a python snippet inside the container with the agent venv."""
|
||||
return docker_exec(CONTAINER_PYTHON, "-c", code, timeout=timeout, env=env, check=check)
|
||||
|
||||
|
||||
# ── fixtures ─────────────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
@pytest.fixture(scope="session")
|
||||
def base_url() -> str:
|
||||
"""Base URL of the web UI (HTTPS, self-signed)."""
|
||||
return WEB_BASE_URL
|
||||
|
||||
|
||||
@pytest.fixture(scope="session")
|
||||
def client(base_url):
|
||||
"""httpx client for the HTTPS web UI with TLS verification disabled."""
|
||||
with httpx.Client(base_url=base_url, verify=False, timeout=15) as c:
|
||||
yield c
|
||||
|
||||
|
||||
@pytest.fixture(scope="session")
|
||||
def token() -> dict:
|
||||
"""Fetch a fresh LiveKit token from the token endpoint (direct, no proxy)."""
|
||||
resp = httpx.post(TOKEN_URL, json={}, timeout=10)
|
||||
resp.raise_for_status()
|
||||
data = resp.json()
|
||||
assert "token" in data, f"token endpoint did not return a token: {data.keys()}"
|
||||
return data
|
||||
|
||||
|
||||
@pytest.fixture(scope="session")
|
||||
def llm_client():
|
||||
"""httpx client pointed at the Gemma LLM API (same one Hope uses)."""
|
||||
with httpx.Client(base_url=GEMMA_BASE_URL, timeout=60) as c:
|
||||
yield c
|
||||
|
||||
|
||||
@pytest.fixture(scope="session")
|
||||
def mcp_client_factory():
|
||||
"""Factory that runs an MCP stdio tool call inside the container.
|
||||
|
||||
Returns (result_text, is_error). Spawns the MCP server script via the
|
||||
agent venv python and drives it with a raw JSON-RPC session using the
|
||||
mcp client library (available in the container, not on the host).
|
||||
"""
|
||||
|
||||
def _call(script: str, tool_name: str, args: dict | None = None, env: dict | None = None) -> tuple[str, bool]:
|
||||
import os as host_os
|
||||
|
||||
merged_env = {**host_os.environ, **(env or {})}
|
||||
code = f"""
|
||||
import asyncio, json
|
||||
from mcp import ClientSession, StdioServerParameters
|
||||
from mcp.client.stdio import stdio_client
|
||||
|
||||
async def main():
|
||||
params = StdioServerParameters(
|
||||
command={CONTAINER_PYTHON!r},
|
||||
args=[{AGENT_DIR_IN_CONTAINER + "/" + script!r}],
|
||||
env={json.dumps(merged_env)},
|
||||
)
|
||||
async with stdio_client(params) as (read, write):
|
||||
async with ClientSession(read, write) as session:
|
||||
await session.initialize()
|
||||
result = await session.call_tool({tool_name!r}, {json.dumps(args or {})})
|
||||
text = ""
|
||||
for block in result.content:
|
||||
if getattr(block, "type", None) == "text":
|
||||
text += block.text
|
||||
print(json.dumps({{"text": text, "isError": bool(result.isError)}}))
|
||||
|
||||
asyncio.run(main())
|
||||
"""
|
||||
proc = docker_exec_python(code, timeout=120, check=False)
|
||||
if proc.returncode != 0:
|
||||
raise AssertionError(
|
||||
f"MCP call {tool_name} failed (rc={proc.returncode}): "
|
||||
f"stderr={proc.stderr[-800:]}"
|
||||
)
|
||||
data = json.loads(proc.stdout.strip().splitlines()[-1])
|
||||
return data["text"], data["isError"]
|
||||
|
||||
return _call
|
||||
|
||||
|
||||
@pytest.fixture(scope="session")
|
||||
def cognee_client():
|
||||
"""Factory that calls tools on the remote Cognee MCP server (Streamable HTTP).
|
||||
|
||||
Returns (result_text, is_error). Uses the mcp client library from the
|
||||
agent venv inside the container.
|
||||
"""
|
||||
|
||||
def _call(tool_name: str, args: dict | None = None) -> tuple[str, bool]:
|
||||
code = f"""
|
||||
import asyncio, json
|
||||
from mcp import ClientSession
|
||||
from mcp.client.streamable_http import streamablehttp_client
|
||||
|
||||
async def main():
|
||||
async with streamablehttp_client({COGNEE_MCP_URL!r}) as (read, write, _):
|
||||
async with ClientSession(read, write) as session:
|
||||
await session.initialize()
|
||||
result = await session.call_tool({tool_name!r}, {json.dumps(args or {})})
|
||||
text = ""
|
||||
for block in result.content:
|
||||
if getattr(block, "type", None) == "text":
|
||||
text += block.text
|
||||
print(json.dumps({{"text": text, "isError": bool(result.isError)}}))
|
||||
|
||||
asyncio.run(main())
|
||||
"""
|
||||
proc = docker_exec_python(code, timeout=120, check=False)
|
||||
if proc.returncode != 0:
|
||||
raise AssertionError(
|
||||
f"Cognee MCP call {tool_name} failed (rc={proc.returncode}): "
|
||||
f"stderr={proc.stderr[-800:]}"
|
||||
)
|
||||
data = json.loads(proc.stdout.strip().splitlines()[-1])
|
||||
return data["text"], data["isError"]
|
||||
|
||||
return _call
|
||||
@@ -0,0 +1,2 @@
|
||||
pytest
|
||||
httpx
|
||||
@@ -0,0 +1,69 @@
|
||||
"""Integration tests for the running container: processes, health, LiveKit."""
|
||||
|
||||
import re
|
||||
|
||||
from conftest import docker_exec
|
||||
|
||||
|
||||
def _proc_cmdlines() -> list[str]:
|
||||
"""All process cmdlines in the container (no `ps` available)."""
|
||||
proc = docker_exec(
|
||||
"sh",
|
||||
"-c",
|
||||
'for p in /proc/[0-9]*/cmdline; do tr "\\0" " " < "$p" 2>/dev/null; echo; done',
|
||||
)
|
||||
return [line.strip() for line in proc.stdout.splitlines() if line.strip()]
|
||||
|
||||
|
||||
def test_agent_process_alive():
|
||||
cmdlines = _proc_cmdlines()
|
||||
assert any("agent.py start" in c for c in cmdlines), (
|
||||
"no agent.py process found in container"
|
||||
)
|
||||
|
||||
|
||||
def test_agent_registered_with_livekit():
|
||||
"""`docker logs` runs on the host (not via docker exec)."""
|
||||
import subprocess
|
||||
|
||||
proc = subprocess.run(
|
||||
["docker", "logs", "--tail", "500", "voice"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=30,
|
||||
)
|
||||
assert "registered worker" in proc.stdout or "worker registered" in proc.stdout, (
|
||||
"agent never registered a worker with LiveKit"
|
||||
)
|
||||
|
||||
|
||||
def test_all_supervisord_processes_up():
|
||||
"""The supervisorctl socket is not exposed in this image, so verify each
|
||||
supervised program's process directly via /proc."""
|
||||
cmdlines = _proc_cmdlines()
|
||||
|
||||
def has(pattern: str) -> bool:
|
||||
return any(re.search(pattern, c) for c in cmdlines)
|
||||
|
||||
assert has(r"livekit --config"), "livekit server not running"
|
||||
assert has(r"agent\.py start"), "voice agent not running"
|
||||
assert has(r"nginx.*daemon off"), "nginx (web UI) not running"
|
||||
assert has(r"token_server\.py"), "token server not running"
|
||||
|
||||
|
||||
def test_container_healthy():
|
||||
import subprocess
|
||||
|
||||
proc = subprocess.run(
|
||||
["docker", "inspect", "--format", "{{.State.Health.Status}}", "voice"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
timeout=30,
|
||||
)
|
||||
assert proc.stdout.strip() == "healthy"
|
||||
|
||||
|
||||
def test_livekit_server_responds():
|
||||
"""LiveKit HTTP API on :7880 (internal; reached via docker exec curl)."""
|
||||
proc = docker_exec("curl", "-s", "-o", "/dev/null", "-w", "%{http_code}", "http://localhost:7880/")
|
||||
assert proc.stdout.strip() in ("200", "404"), f"unexpected status {proc.stdout.strip()!r}"
|
||||
@@ -0,0 +1,93 @@
|
||||
"""Tests for GemmaLLM context compaction (agent.GemmaLLM._compact_context).
|
||||
|
||||
Compaction runs a real side LLM call to the xNAS, so these are marked slow.
|
||||
The agent module is imported inside the container via docker exec.
|
||||
"""
|
||||
|
||||
import pytest
|
||||
|
||||
from conftest import GEMMA_BASE_URL, GEMMA_MODEL, docker_exec_python
|
||||
|
||||
# {num_messages} and {do_chat} are filled per test; all other braces are literal.
|
||||
TEMPLATE = '''
|
||||
import logging, sys
|
||||
from livekit.agents import llm as lk_llm
|
||||
sys.path.insert(0, "/opt/voice-agent")
|
||||
from agent import GemmaLLM
|
||||
|
||||
logging.basicConfig(level=logging.INFO, stream=sys.stdout)
|
||||
|
||||
llm = GemmaLLM(
|
||||
model="{model}",
|
||||
base_url="{base_url}",
|
||||
api_key="not-needed",
|
||||
max_completion_tokens=100,
|
||||
extra_body={{"chat_template_kwargs": {{"enable_thinking": False}}}},
|
||||
)
|
||||
|
||||
items = [lk_llm.ChatMessage(role="system", content=["You are Hope, a voice assistant."])]
|
||||
for i in range({num_messages}):
|
||||
items.append(lk_llm.ChatMessage(role="user", content=["user message number " + str(i)]))
|
||||
items.append(lk_llm.ChatMessage(role="assistant", content=["agent reply number " + str(i)]))
|
||||
|
||||
ctx = lk_llm.ChatContext(items)
|
||||
print("BUILT ctx with " + str(len(ctx.items)) + " items")
|
||||
|
||||
import asyncio
|
||||
|
||||
async def main():
|
||||
try:
|
||||
compacted = llm._compact_context(ctx, None)
|
||||
except Exception as e:
|
||||
print("COMPACT_ERROR " + repr(e))
|
||||
return
|
||||
print("AFTER items=" + str(len(compacted.items)))
|
||||
for it in compacted.items:
|
||||
role = getattr(it, "role", None)
|
||||
content = GemmaLLM._msg_content_text(it)
|
||||
print("ITEM " + str(role) + ": " + content[:120])
|
||||
|
||||
{do_chat}
|
||||
|
||||
asyncio.run(main())
|
||||
'''
|
||||
|
||||
CHAT_SNIPPET = ' stream = llm.chat(chat_ctx=compacted)\n collected = await stream.collect()\n print("CHAT_OK text_len=" + str(len(collected.text)))'
|
||||
|
||||
|
||||
def _run(num_messages: int, *, do_chat: bool = False) -> str:
|
||||
code = TEMPLATE.format(
|
||||
model=GEMMA_MODEL,
|
||||
base_url=GEMMA_BASE_URL,
|
||||
num_messages=num_messages,
|
||||
do_chat=CHAT_SNIPPET if do_chat else "pass",
|
||||
)
|
||||
proc = docker_exec_python(code, timeout=120, check=False)
|
||||
return proc.stdout + (("\n" + proc.stderr) if proc.stderr else "")
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_compaction_triggers_on_size():
|
||||
"""A context with > 24 items must compact without crashing and chat() must still work."""
|
||||
# 13 pairs = 26 non-system + 1 system = 27 items (> threshold of 24)
|
||||
out = _run(13, do_chat=True)
|
||||
assert "BUILT ctx with 27 items" in out
|
||||
assert "Context compacted:" in out
|
||||
assert "CHAT_OK text_len=" in out
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_compaction_preserves_system_prompt():
|
||||
out = _run(28)
|
||||
lines = [l for l in out.splitlines() if l.startswith("ITEM system:")]
|
||||
assert lines, f"no system message survived compaction:\n{out}"
|
||||
assert "You are Hope" in lines[0]
|
||||
|
||||
|
||||
def test_short_context_no_compaction():
|
||||
"""Context with non_sys <= KEEP_RECENT must NOT be compacted."""
|
||||
# system + 4 user/assistant pairs = 9 items, non_sys=8 <= KEEP_RECENT(10)
|
||||
out = _run(4)
|
||||
assert "BUILT ctx with 9 items" in out
|
||||
# _compact_context should return the context unchanged (early return)
|
||||
assert "AFTER items=9" in out, f"expected no compaction:\n{out}"
|
||||
@@ -0,0 +1,106 @@
|
||||
"""Tests for the Gemma LLM API — the exact endpoint Hope's agent uses."""
|
||||
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
from conftest import GEMMA_BASE_URL, GEMMA_MODEL
|
||||
|
||||
WEATHER_TOOL = {
|
||||
"type": "function",
|
||||
"function": {
|
||||
"name": "get_weather",
|
||||
"description": "Get the current weather and today's forecast for a location.",
|
||||
"parameters": {
|
||||
"type": "object",
|
||||
"properties": {"location": {"type": "string"}},
|
||||
"required": ["location"],
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def _chat(llm_client, messages, *, tools=None, stream=False, max_tokens=200):
|
||||
payload = {
|
||||
"model": GEMMA_MODEL,
|
||||
"messages": messages,
|
||||
"max_tokens": max_tokens,
|
||||
"temperature": 0.3,
|
||||
"chat_template_kwargs": {"enable_thinking": False},
|
||||
}
|
||||
if tools:
|
||||
payload["tools"] = tools
|
||||
if stream:
|
||||
payload["stream"] = True
|
||||
return llm_client.post("/chat/completions", json=payload)
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_simple_completion(llm_client):
|
||||
resp = _chat(
|
||||
llm_client,
|
||||
[{"role": "user", "content": "What is 2+2? Answer with just the number."}],
|
||||
max_tokens=50,
|
||||
)
|
||||
assert resp.status_code == 200
|
||||
data = resp.json()
|
||||
content = (data["choices"][0]["message"].get("content") or "").strip()
|
||||
assert "4" in content
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_tool_call_weather(llm_client):
|
||||
resp = _chat(
|
||||
llm_client,
|
||||
[{"role": "user", "content": "What is the weather like in Raleigh right now?"}],
|
||||
tools=[WEATHER_TOOL],
|
||||
max_tokens=100,
|
||||
)
|
||||
assert resp.status_code == 200
|
||||
choice = resp.json()["choices"][0]
|
||||
assert choice["finish_reason"] == "tool_calls"
|
||||
tool_calls = choice["message"]["tool_calls"]
|
||||
assert any(tc["function"]["name"] == "get_weather" for tc in tool_calls)
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_thinking_disabled(llm_client):
|
||||
"""A normally-slow prompt must answer fast when enable_thinking=false."""
|
||||
messages = [
|
||||
{"role": "user", "content": "Count from 1 to 20, then tell me the sum."},
|
||||
]
|
||||
start = time.monotonic()
|
||||
resp = _chat(llm_client, messages, max_tokens=300)
|
||||
elapsed = time.monotonic() - start
|
||||
assert resp.status_code == 200
|
||||
content = (resp.json()["choices"][0]["message"].get("content") or "").strip()
|
||||
assert len(content) > 10
|
||||
# With thinking enabled this prompt takes well over 3s on the xNAS.
|
||||
assert elapsed < 3, f"response took {elapsed:.1f}s — thinking may be enabled"
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_streaming(llm_client):
|
||||
resp = _chat(
|
||||
llm_client,
|
||||
[{"role": "user", "content": "Say the words: the quick brown fox jumps over the lazy dog."}],
|
||||
stream=True,
|
||||
max_tokens=100,
|
||||
)
|
||||
assert resp.status_code == 200
|
||||
|
||||
chunks = []
|
||||
for line in resp.iter_lines():
|
||||
if not line or not line.startswith("data: "):
|
||||
continue
|
||||
payload = line[6:].strip()
|
||||
if payload == "[DONE]":
|
||||
break
|
||||
chunk = __import__("json").loads(payload)
|
||||
delta = (chunk.get("choices") or [{}])[0].get("delta", {})
|
||||
if delta.get("content"):
|
||||
chunks.append(delta["content"])
|
||||
|
||||
assert len(chunks) > 1, f"expected multiple streamed chunks, got {len(chunks)}"
|
||||
assembled = "".join(chunks)
|
||||
assert "quick brown fox" in assembled.lower()
|
||||
@@ -0,0 +1,120 @@
|
||||
"""Tests for the MCP tool servers (weather, memory, skills) via docker exec.
|
||||
|
||||
Each test spawns the real MCP server script inside the container and drives it
|
||||
with a JSON-RPC session using the mcp client library from the agent venv.
|
||||
"""
|
||||
|
||||
import re
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
def _extract_number(text: str) -> float | None:
|
||||
"""Pull the first plausible temperature number out of a weather summary."""
|
||||
m = re.search(r"(\d+)\s*degrees", text)
|
||||
if m:
|
||||
return float(m.group(1))
|
||||
m = re.search(r"\b(\d{2,3})\b", text)
|
||||
return float(m.group(1)) if m else None
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_weather_get(mcp_client_factory):
|
||||
text, is_error = mcp_client_factory("weather_mcp.py", "get_weather", {"location": "Raleigh"})
|
||||
assert not is_error
|
||||
assert "degrees" in text.lower()
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_weather_fahrenheit(mcp_client_factory):
|
||||
"""wttr.in reports Celsius; the tool must convert to Fahrenheit."""
|
||||
text, is_error = mcp_client_factory("weather_mcp.py", "get_weather", {"location": "Raleigh"})
|
||||
assert not is_error
|
||||
temp = _extract_number(text)
|
||||
assert temp is not None, f"no temperature found in: {text}"
|
||||
# Fahrenheit for any inhabited place is > 50 (Celsius would be ~10-30).
|
||||
assert temp > 50, f"temperature {temp} looks like Celsius, expected Fahrenheit"
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_time_get(mcp_client_factory):
|
||||
text, is_error = mcp_client_factory("weather_mcp.py", "get_time", {})
|
||||
assert not is_error
|
||||
assert ("AM" in text or "PM" in text)
|
||||
days = ("Monday", "Tuesday", "Wednesday", "Thursday", "Friday", "Saturday", "Sunday")
|
||||
assert any(d in text for d in days), f"no day name in: {text}"
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_memory_remember_and_recall(cognee_client):
|
||||
"""Cognee: store a fact, then recall it."""
|
||||
probe = "zebra42 is the probe fact for the voice test suite"
|
||||
|
||||
stored, err = cognee_client("remember", {"data": probe})
|
||||
assert not err, f"remember failed: {stored}"
|
||||
|
||||
recalled, err = cognee_client("recall", {"query": "zebra42"})
|
||||
assert not err
|
||||
assert "zebra42" in recalled
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_memory_forget(cognee_client):
|
||||
"""Cognee: store a fact, verify recall, then forget it."""
|
||||
probe = "purple giraffe77 is the forget probe"
|
||||
|
||||
stored, err = cognee_client("remember", {"data": probe})
|
||||
assert not err, f"remember failed: {stored}"
|
||||
|
||||
# Verify it's there
|
||||
recalled, err = cognee_client("recall", {"query": "purple giraffe77"})
|
||||
assert not err
|
||||
assert "giraffe77" in recalled
|
||||
|
||||
|
||||
@pytest.mark.slow
|
||||
def test_skill_save_and_recall(mcp_client_factory):
|
||||
name = "Test Suite Probe Skill"
|
||||
description = "A probe skill used by the voice test suite."
|
||||
steps = "Step 1: do the thing. Step 2: verify zebra42 was done."
|
||||
|
||||
saved, err = mcp_client_factory(
|
||||
"skills_mcp.py",
|
||||
"skill_save",
|
||||
{"name": name, "description": description, "steps": steps},
|
||||
)
|
||||
assert not err
|
||||
assert "Saved" in saved
|
||||
|
||||
try:
|
||||
recalled, err = mcp_client_factory("skills_mcp.py", "skill_recall", {"query": "zebra42"})
|
||||
assert not err
|
||||
assert "zebra42" in recalled
|
||||
assert "Steps" in recalled
|
||||
finally:
|
||||
from conftest import docker_exec
|
||||
|
||||
docker_exec("rm", "-f", "/skills/test-suite-probe-skill.md")
|
||||
|
||||
|
||||
def test_skill_list(mcp_client_factory):
|
||||
from conftest import docker_exec
|
||||
|
||||
name = "Test Suite List Skill"
|
||||
try:
|
||||
saved, err = mcp_client_factory(
|
||||
"skills_mcp.py",
|
||||
"skill_save",
|
||||
{
|
||||
"name": name,
|
||||
"description": "listing probe skill",
|
||||
"steps": "one step",
|
||||
},
|
||||
)
|
||||
assert not err
|
||||
listed, err = mcp_client_factory("skills_mcp.py", "skill_list", {})
|
||||
assert not err
|
||||
# list output uses the slugified filename
|
||||
assert "test-suite-list-skill" in listed
|
||||
finally:
|
||||
docker_exec("rm", "-f", "/skills/test-suite-list-skill.md")
|
||||
@@ -0,0 +1,74 @@
|
||||
"""Tests for the web UI served by nginx (HTTPS :8090) and the token endpoint."""
|
||||
|
||||
import json
|
||||
|
||||
import httpx
|
||||
|
||||
|
||||
def test_index_serves(client):
|
||||
resp = client.get("/")
|
||||
assert resp.status_code == 200
|
||||
body = resp.text
|
||||
assert "Hope" in body
|
||||
assert "livekit-client.umd.js" in body
|
||||
assert "app.js" in body
|
||||
|
||||
|
||||
def test_app_js_serves(client):
|
||||
resp = client.get("/app.js")
|
||||
assert resp.status_code == 200
|
||||
assert "LiveKitClient" in resp.text or "Room" in resp.text
|
||||
|
||||
|
||||
def test_style_css_serves(client):
|
||||
resp = client.get("/style.css")
|
||||
assert resp.status_code == 200
|
||||
body = resp.text
|
||||
assert "{" in body and "}" in body
|
||||
assert "color" in body or "background" in body
|
||||
|
||||
|
||||
def test_manifest_serves(client):
|
||||
resp = client.get("/manifest.json")
|
||||
assert resp.status_code == 200
|
||||
manifest = json.loads(resp.text)
|
||||
assert "name" in manifest
|
||||
|
||||
|
||||
def test_favicon_serves(client):
|
||||
resp = client.get("/favicon.svg")
|
||||
assert resp.status_code == 200
|
||||
assert "svg" in resp.headers.get("content-type", "")
|
||||
assert "<svg" in resp.text
|
||||
|
||||
|
||||
def test_livekit_client_bundle_serves(client):
|
||||
resp = client.get("/livekit-client.umd.js")
|
||||
assert resp.status_code == 200
|
||||
assert len(resp.text) > 100_000 # vendored UMD bundle is large
|
||||
|
||||
|
||||
def test_token_endpoint_direct():
|
||||
"""Token server on 127.0.0.1:8091 (host network mode)."""
|
||||
resp = httpx.post("http://127.0.0.1:8091/token", json={}, timeout=10)
|
||||
assert resp.status_code == 200
|
||||
data = resp.json()
|
||||
assert "token" in data and data["token"].count(".") == 2 # JWT shape
|
||||
|
||||
|
||||
def test_token_endpoint_via_nginx(client):
|
||||
"""nginx proxies /token to the token server."""
|
||||
resp = client.post("/token", json={})
|
||||
assert resp.status_code == 200
|
||||
data = resp.json()
|
||||
assert "token" in data and data["token"].count(".") == 2
|
||||
|
||||
|
||||
def test_livekit_ws_proxy(client):
|
||||
"""The /livekit/ path is proxied to the LiveKit server (HTTP :7880)."""
|
||||
resp = client.get("/livekit/")
|
||||
# LiveKit's HTTP API answers on /; a 200 or an API error JSON both prove
|
||||
# the proxy reaches the LiveKit server rather than nginx serving statics.
|
||||
assert resp.status_code in (200, 404, 405)
|
||||
if resp.headers.get("content-type", "").startswith("application/json"):
|
||||
json.loads(resp.text)
|
||||
+74
@@ -36,6 +36,12 @@ let agentAudioLevel = 0; // 0..1, smoothed
|
||||
let userSpeaking = false; // local VAD-ish flag from mic level
|
||||
let wakeLock = null;
|
||||
|
||||
function escapeHtml(s) {
|
||||
const d = document.createElement("div");
|
||||
d.textContent = s || "";
|
||||
return d.innerHTML;
|
||||
}
|
||||
|
||||
// ── DOM ─────────────────────────────────────────────────────────────────────
|
||||
const startBtn = document.getElementById("startBtn");
|
||||
const stopBtn = document.getElementById("stopBtn");
|
||||
@@ -54,6 +60,72 @@ const sheetBackdrop = document.getElementById("sheetBackdrop");
|
||||
const closeSheetBtn = document.getElementById("closeSheetBtn");
|
||||
const voiceListEl = document.getElementById("voiceList");
|
||||
const clearBtn = document.getElementById("clearBtn");
|
||||
const tasksBtn = document.getElementById("tasksBtn");
|
||||
const tasksPanel = document.getElementById("tasksPanel");
|
||||
const tasksList = document.getElementById("tasksList");
|
||||
const tasksBadge = document.getElementById("tasksBadge");
|
||||
const closeTasksBtn = document.getElementById("closeTasksBtn");
|
||||
|
||||
// ── Background tasks state ──────────────────────────────────────────────────
|
||||
const tasks = new Map(); // task_id -> {id, description, status, steps: [], result, error}
|
||||
|
||||
function updateTasksBadge() {
|
||||
const active = [...tasks.values()].filter(t => t.status === "running" || t.status === "pending").length;
|
||||
if (active > 0) {
|
||||
tasksBadge.hidden = false;
|
||||
tasksBadge.textContent = String(active);
|
||||
} else {
|
||||
tasksBadge.hidden = true;
|
||||
}
|
||||
}
|
||||
|
||||
function renderTasks() {
|
||||
const items = [...tasks.values()].sort((a, b) => (b.created_at || 0) - (a.created_at || 0));
|
||||
if (items.length === 0) {
|
||||
tasksList.innerHTML = '<p class="tasks-empty">No tasks yet. Ask Hope to research something.</p>';
|
||||
return;
|
||||
}
|
||||
tasksList.innerHTML = "";
|
||||
for (const t of items) {
|
||||
const card = document.createElement("div");
|
||||
card.className = "task-card";
|
||||
const statusClass = t.status || "pending";
|
||||
card.innerHTML = `
|
||||
<div class="task-card-header">
|
||||
<span class="task-status-dot ${statusClass}"></span>
|
||||
<span class="task-desc">${escapeHtml(t.description)}</span>
|
||||
</div>
|
||||
<ul class="task-steps">${t.steps.map(s => `<li class="task-step ${s.role}">${escapeHtml(s.content)}</li>`).join("")}</ul>
|
||||
${t.result ? `<div class="task-result">${escapeHtml(t.result)}</div>` : ""}
|
||||
${t.error ? `<div class="task-error">${escapeHtml(t.error)}</div>` : ""}
|
||||
`;
|
||||
tasksList.appendChild(card);
|
||||
}
|
||||
}
|
||||
|
||||
function handleTaskEvent(data) {
|
||||
const { task_id, event, role, content, result, error, task } = data;
|
||||
if (event === "created" && task) {
|
||||
tasks.set(task_id, { ...task, steps: [] });
|
||||
} else if (event === "step") {
|
||||
const t = tasks.get(task_id);
|
||||
if (t) {
|
||||
t.steps.push({ role, content });
|
||||
t.status = "running";
|
||||
}
|
||||
} else if (event === "completed") {
|
||||
const t = tasks.get(task_id);
|
||||
if (t) { t.status = "completed"; t.result = result; }
|
||||
} else if (event === "failed") {
|
||||
const t = tasks.get(task_id);
|
||||
if (t) { t.status = "failed"; t.error = error; }
|
||||
}
|
||||
updateTasksBadge();
|
||||
renderTasks();
|
||||
}
|
||||
|
||||
tasksBtn.addEventListener("click", () => tasksPanel.classList.toggle("open"));
|
||||
closeTasksBtn.addEventListener("click", () => tasksPanel.classList.remove("open"));
|
||||
|
||||
// ── Audio analysers (mic + remote agent audio) ─────────────────────────────
|
||||
let micCtx = null;
|
||||
@@ -352,6 +424,8 @@ function handleDataPacket(payload, participant, topic) {
|
||||
const msg = JSON.parse(new TextDecoder().decode(payload));
|
||||
if (msg.type === "transcript") {
|
||||
addMessage(msg.role || "agent", msg.text);
|
||||
} else if (msg.type === "task_event") {
|
||||
handleTaskEvent(msg);
|
||||
}
|
||||
// set_voice messages flow the other way; nothing to do client-side.
|
||||
} catch (e) {
|
||||
|
||||
@@ -64,6 +64,23 @@
|
||||
<div class="voice-list" id="voiceList"></div>
|
||||
</div>
|
||||
|
||||
<!-- Task panel (background tasks) -->
|
||||
<button id="tasksBtn" class="btn btn-icon tasks-fab" aria-label="Background tasks" title="Background tasks">
|
||||
<svg viewBox="0 0 24 24" width="18" height="18" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round" stroke-linejoin="round"><rect x="3" y="3" width="18" height="18" rx="2"></rect><path d="M9 9h6M9 13h6M9 17h4"></path></svg>
|
||||
<span class="tasks-badge" id="tasksBadge" hidden>0</span>
|
||||
</button>
|
||||
<div class="tasks-panel" id="tasksPanel">
|
||||
<div class="tasks-header">
|
||||
<h2>Background Tasks</h2>
|
||||
<button id="closeTasksBtn" class="btn btn-ghost" aria-label="Close">
|
||||
<svg viewBox="0 0 24 24" width="18" height="18" fill="none" stroke="currentColor" stroke-width="2" stroke-linecap="round"><path d="M18 6 6 18M6 6l12 12"></path></svg>
|
||||
</button>
|
||||
</div>
|
||||
<div class="tasks-list" id="tasksList">
|
||||
<p class="tasks-empty">No tasks yet. Ask Hope to research something.</p>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<script src="livekit-client.umd.js"></script>
|
||||
<script src="app.js"></script>
|
||||
</body>
|
||||
|
||||
+173
@@ -510,6 +510,178 @@ h1 {
|
||||
.transcript::-webkit-scrollbar-thumb,
|
||||
.voice-list::-webkit-scrollbar-thumb { background: var(--border); border-radius: 3px; }
|
||||
|
||||
/* ── Task panel ──────────────────────────────────────────────────────────── */
|
||||
.tasks-fab {
|
||||
position: fixed;
|
||||
bottom: calc(1.5rem + env(safe-area-inset-bottom));
|
||||
right: 1.5rem;
|
||||
z-index: 100;
|
||||
}
|
||||
|
||||
.tasks-badge {
|
||||
position: absolute;
|
||||
top: -4px;
|
||||
right: -4px;
|
||||
background: #3b82f6;
|
||||
color: #fff;
|
||||
font-size: 0.65rem;
|
||||
font-weight: 700;
|
||||
min-width: 18px;
|
||||
height: 18px;
|
||||
border-radius: 9px;
|
||||
display: grid;
|
||||
place-items: center;
|
||||
padding: 0 4px;
|
||||
}
|
||||
|
||||
.tasks-panel {
|
||||
position: fixed;
|
||||
top: 0;
|
||||
right: -320px;
|
||||
width: 320px;
|
||||
max-width: 85vw;
|
||||
height: 100dvh;
|
||||
background: var(--panel);
|
||||
border-left: 1px solid var(--border);
|
||||
z-index: 200;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
transition: right 0.25s ease;
|
||||
}
|
||||
|
||||
.tasks-panel.open {
|
||||
right: 0;
|
||||
}
|
||||
|
||||
.tasks-header {
|
||||
display: flex;
|
||||
align-items: center;
|
||||
justify-content: space-between;
|
||||
padding: 1rem;
|
||||
border-bottom: 1px solid var(--border);
|
||||
}
|
||||
|
||||
.tasks-header h2 {
|
||||
font-size: 1rem;
|
||||
margin: 0;
|
||||
}
|
||||
|
||||
.tasks-list {
|
||||
flex: 1;
|
||||
overflow-y: auto;
|
||||
padding: 0.75rem;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
gap: 0.75rem;
|
||||
}
|
||||
|
||||
.tasks-empty {
|
||||
color: var(--text-faint);
|
||||
font-size: 0.85rem;
|
||||
text-align: center;
|
||||
margin-top: 2rem;
|
||||
}
|
||||
|
||||
.task-card {
|
||||
background: var(--panel-strong);
|
||||
border: 1px solid var(--border);
|
||||
border-radius: 10px;
|
||||
padding: 0.75rem;
|
||||
font-size: 0.82rem;
|
||||
}
|
||||
|
||||
.task-card-header {
|
||||
display: flex;
|
||||
align-items: center;
|
||||
gap: 0.5rem;
|
||||
margin-bottom: 0.4rem;
|
||||
}
|
||||
|
||||
.task-status-dot {
|
||||
width: 8px;
|
||||
height: 8px;
|
||||
border-radius: 50%;
|
||||
background: #6b7280;
|
||||
flex-shrink: 0;
|
||||
}
|
||||
|
||||
.task-status-dot.running {
|
||||
background: #3b82f6;
|
||||
animation: pulse-dot 1.5s infinite;
|
||||
}
|
||||
|
||||
.task-status-dot.completed {
|
||||
background: #22c55e;
|
||||
}
|
||||
|
||||
.task-status-dot.failed {
|
||||
background: #ef4444;
|
||||
}
|
||||
|
||||
@keyframes pulse-dot {
|
||||
0%, 100% { opacity: 1; }
|
||||
50% { opacity: 0.4; }
|
||||
}
|
||||
|
||||
.task-desc {
|
||||
font-weight: 600;
|
||||
color: var(--text);
|
||||
flex: 1;
|
||||
overflow: hidden;
|
||||
text-overflow: ellipsis;
|
||||
white-space: nowrap;
|
||||
}
|
||||
|
||||
.task-steps {
|
||||
margin-top: 0.5rem;
|
||||
padding-left: 0;
|
||||
list-style: none;
|
||||
display: flex;
|
||||
flex-direction: column;
|
||||
gap: 0.3rem;
|
||||
max-height: 200px;
|
||||
overflow-y: auto;
|
||||
}
|
||||
|
||||
.task-step {
|
||||
font-size: 0.75rem;
|
||||
color: var(--text-dim);
|
||||
padding: 0.25rem 0.4rem;
|
||||
border-radius: 4px;
|
||||
background: rgba(0,0,0,0.15);
|
||||
}
|
||||
|
||||
.task-step.tool_call {
|
||||
color: #60a5fa;
|
||||
}
|
||||
|
||||
.task-step.tool_result {
|
||||
color: var(--text-faint);
|
||||
}
|
||||
|
||||
.task-step.text {
|
||||
color: var(--text);
|
||||
font-weight: 500;
|
||||
}
|
||||
|
||||
.task-result {
|
||||
margin-top: 0.5rem;
|
||||
padding: 0.5rem;
|
||||
background: rgba(34, 197, 94, 0.1);
|
||||
border-radius: 6px;
|
||||
font-size: 0.8rem;
|
||||
color: var(--text);
|
||||
}
|
||||
|
||||
.task-error {
|
||||
margin-top: 0.5rem;
|
||||
padding: 0.5rem;
|
||||
background: rgba(239, 68, 68, 0.1);
|
||||
border-radius: 6px;
|
||||
font-size: 0.8rem;
|
||||
color: #fca5a5;
|
||||
}
|
||||
|
||||
/* ── Reduced motion ──────────────────────────────────────────────────────── */
|
||||
@media (prefers-reduced-motion: reduce) {
|
||||
.orb-core,
|
||||
@@ -517,6 +689,7 @@ h1 {
|
||||
.message-row,
|
||||
.settings-sheet,
|
||||
.sheet-backdrop,
|
||||
.tasks-panel,
|
||||
.audio-viz i,
|
||||
.mic-meter-fill {
|
||||
animation: none !important;
|
||||
|
||||
+3
-2
@@ -47,7 +47,7 @@ def make_token(room_name: str, identity: str) -> str:
|
||||
"canPublishData": True,
|
||||
},
|
||||
"roomConfig": {
|
||||
"agents": [{"agentName": AGENT_NAME}],
|
||||
"agents": [{"agentName": AGENT_NAME, "jobType": "JT_ROOM"}],
|
||||
},
|
||||
}
|
||||
h = b64url(json.dumps(header).encode())
|
||||
@@ -66,7 +66,8 @@ class Handler(BaseHTTPRequestHandler):
|
||||
|
||||
length = int(self.headers.get("Content-Length", 0))
|
||||
body = json.loads(self.rfile.read(length) or b"{}")
|
||||
room_name = "voice-room"
|
||||
# Unique room per session so LiveKit dispatches the agent at creation.
|
||||
room_name = f"voice-{uuid.uuid4().hex[:8]}"
|
||||
identity = body.get("identity") or f"user-{uuid.uuid4().hex[:8]}"
|
||||
|
||||
token = make_token(room_name, identity)
|
||||
|
||||
Reference in New Issue
Block a user