diff --git a/.env.example b/.env.example index 7f42820..e2fb145 100644 --- a/.env.example +++ b/.env.example @@ -53,3 +53,7 @@ QDRANT_EXPOSE_PORT=6333 RABBITMQ_EXPOSE_PORT=5672 LLAMA_EXPOSE_PORT=8081 OLLAMA_EXPOSE_PORT=11434 + +# ── Image generation (ComfyUI on worker) ───────────────────── +CAIC_COMFYUI_BASE=http://192.168.50.115:8188 +CAIC_COMFYUI_TIMEOUT=120 diff --git a/README.md b/README.md index 1cf6707..03a3e72 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ ![cAIc banner](static/readme-banner.png) -# cAIc v1.0.0 +# cAIc v1.1.0 **Cluster AI coordinator — heterogeneous GPU inference for homelab AI clusters.** @@ -60,6 +60,7 @@ A worker with a slow GPU still contributes — it handles less latency-sensitive - **At-rest encryption** — AES-256-GCM on all query-derived text in SQLite and Qdrant - **IDE integration** — OpenAI-compatible `/v1/chat/completions` endpoint for Continue.dev and friends - **OpenAI-compat FIM** — `/v1/fim/completions` for code completion +- **Image generation** — ComfyUI-backed image gen via cluster workers (Stable Diffusion / Flux) - **6 color themes** — IBM Blue, Matrix, Dark, Light, Amber, Trippin - **Docker-ready** — `docker compose up -d` and you're running @@ -155,6 +156,21 @@ Gitea (primary): `ssh://gitea@llgit.llamachile.tube:1319/gramps/caic.git` --- +## What's New in v1.1.0 + +### Image Generation Service +- `POST /api/image/generate` — proxy endpoint routes to ComfyUI on cluster workers +- `GET /api/image/status` — lists available image gen nodes +- Node agent auto-detects ComfyUI and registers `image_gen` capability +- Full ComfyUI workflow: CheckpointLoader → KSampler → VAEDecode → SaveImage +- Cluster AMQP protocol extended: `cmd.image_generate`, `image_generated`, `image_failed` +- Hardware probe checks ComfyUI reachability + checkpoint model list +- 27 new tests covering cluster handlers, router proxy, node agent, hardware, capability detection + +### Bug Fixes & Hardening +- Hardware assessment now probes ComfyUI alongside llama-server, Qdrant, SearXNG +- Node agent config extended with `comfyui_port` (default 8188) + ## What's New in v1.0.0 ### Docker Containerization (B3) diff --git a/TASKS.md b/TASKS.md index acfa712..bfcb61a 100644 --- a/TASKS.md +++ b/TASKS.md @@ -1,994 +1,42 @@ -# cAIc — OpenCode Prompt Sequence -# Generated: 2026-07-14 -# Execute sequentially. Run full test suite after each task before proceeding. -# Test command: ./venv/bin/python -m pytest tests/ -v +# cAIc — Task List (v1.0+) + +Previous task history archived at `docs/archive/TASKS-pre-1.0.md`. --- -## Session 2026-07-14 — RAG bugfixes + Topbar redesign +## TASK 1 — Image Generation Service (corsair) -- **RAG bugs fixed**: Collection name mismatch (`jarvis_rag` → `caic_rag`, migrated 219 points), `vectors_count`→`points_count` (Qdrant v1.10+ API change), removed unindexed `order_by` that caused 502 on scroll, made `RAG_COLLECTION` env-configurable (`CAIC_RAG_COLLECTION`). -- **Semantic search fixed**: Set `CAIC_EMBED_URL=http://192.168.50.108:11434` (mxbai-embed-large lives on ultron, not the old embed server). -- **Topbar redesign**: Moved system stats (CPU/MEM/GPU/VRAM/TOK) to a centered bottom strip. Moved toggles (MEM, SEARCH, PROFILE, SORT, PRIVACY) into a ⋮ hamburger menu next to ADMIN badge. Palette icon sits immediately after version number in topbar-left. Removed standalone (i) button — privacy info accessible via ⋮ → About Privacy. Input bar above chat, stats at very bottom. Mobile-responsive padding/sizing. +**Goal:** Add image generation as a cluster capability. corsair (RTX 5070 Ti, 16 GB) registers as an image gen worker in the cAIc cluster. + +### Requirements: + +1. **Add `"image_gen"` capability** to the cluster protocol in `cluster.py` — valid capability values should include `image_gen` +2. **Image gen API wrapper on corsair** — run ComfyUI, Automatic1111, or a lightweight API (e.g., `sd-api` or `comfyui-api`) that exposes a simple `POST /generate` endpoint accepting a prompt and returning a PNG +3. **Proxy endpoint in cAIc** — `POST /api/image/generate` on the coordinator, routes the request to corsair's image gen service via AMQP or direct HTTP +4. **Update `hardware.py`** to probe the image gen service for reachability and status +5. **Update node_agent** to report image gen capability and service status on registration + +### Architecture: + +``` +User prompt → cAIc coordinator → AMQP/HTTP → corsair (ComfyUI/API) → PNG → coordinator → user +``` + +### Considerations: + +- Response time: expect 5-30 seconds per image depending on model/resolution +- Queue management: what if multiple requests come in at once? +- Model selection: which SD/Flux model to run by default? +- Resolution limits: max image size? +- CORS/headers for serving generated images back to the UI + +### Tests: + +- Mock image gen service, verify proxy routing +- Verify `hardware.py` probes image gen endpoint +- Verify node_agent registers with `image_gen` capability +- Verify 429/503 handling when service is busy or down + +### Status: ✅ Backend Complete (ComfyUI install pending on corsair) --- - -## ~~TASK 1 — README Cleanup [DONE]~~ - -Review README.md in the current repo. Remove any node references other than `coordinator` (192.168.50.108) and `worker` (192.168.50.210). Ensure all references to the project use the exact casing `cAIc` — not `Jarvischat`, `JarvisChat`, or `jarvischat`. Do not change any functional content, endpoint documentation, or architecture descriptions — this is a text cleanup only. After editing, verify the file renders cleanly as markdown. Commit with message: `docs: clean up node references and branding consistency`. - -No new tests required for this task. - ---- - -## ~~TASK 2 — Qwen2.5-Coder llama-server Service on Coordinator (Infrastructure) [DONE]~~ - -**Status: Systemd unit created, verified, and restored.** - -This task originally defined creation of `/etc/systemd/system/llama-server-coder.service` (port 8082, Qwen2.5-Coder-14B Q5_K_M) as a prerequisite for dynamic model swapping. That sysadmin work is done. - -**The real Task 2 deliverable — the ability to dynamically swap models based on query classification — is delivered by Roadmap N (Tasks 9–15).** The flow: - -1. **Task 13** — Phi-4-mini triage (`triage.py`) classifies the query as `general`, `code`, `search`, or `rag` -2. **Task 13** — `select_node()` picks the best worker node; if the ideal model isn't active, it triggers a swap -3. **Task 14** — `request_model_swap()` publishes `cmd.swap_model` via AMQP `jc.admin` exchange -4. **Task 12** — The node agent on worker receives the command, stops the current llama-server, starts the correct one, waits for health, and publishes `model_ready` -5. **Task 14** — coordinator receives `model_ready`, updates the cluster registry, and routes the query to the node - -The swap is async and transparent — the user sees only latency. The UI (Task 15) shows a yellow "swapping" status dot during the transition. - -The service unit at `/etc/systemd/system/llama-server-coder.service` is the **target** the node agent starts when swapping to code inference. It is not enabled at boot — the AMQP cluster manages activation. - -See Tasks 9–15 for the actual model swap implementation. - -No pytest tests required for this infrastructure task. - ---- - -## ~~TASK 3 — Update OpenCode Config to Use Qwen on :8082 [DONE]~~ - -Update `/home/gramps/.config/opencode/opencode.jsonc` (on this machine, coordinator) to point the configured provider at `http://127.0.0.1:8082/v1` instead of `http://127.0.0.1:8081/v1`. The model name in the config should be updated to reflect `qwen2.5-coder-14b` or whatever model ID the llama-server instance at :8082 reports via `/v1/models`. Verify the endpoint is reachable before writing the config change. Do not restart OpenCode — the config change takes effect on next session start. - -No pytest tests required for this task. - ---- - -## ~~TASK 4 — File/Document Attachment: Backend Ingest Endpoint [DONE]~~ - -**Status: `POST /api/upload` with mode=(context|ingest|both), PDF/text extraction, Qdrant upsert, SQLite context (1hr expiry). Committed `4a891c8` (v1.9.0).** - -This task implements the backend half of file/document attachment (TODO #21). The goal is dual-aspect upload: a file can be used as immediate chat context, ingested into the RAG corpus (Qdrant), or both. - -**Add to `config.py`:** -- `UPLOAD_DIR` — path for temporary upload storage, default `/tmp/caic_uploads` -- `MAX_UPLOAD_BYTES` — max file size, default 20MB -- `SUPPORTED_UPLOAD_TYPES` — set of MIME types: `text/plain`, `text/markdown`, `application/pdf`, `application/json`, `text/x-python`, `text/html` - -**Create `routers/upload.py`:** - -Implement `POST /api/upload` (admin required). Accept `multipart/form-data` with: -- `file` — the uploaded file (required) -- `mode` — string enum: `context` (inject into next chat only), `ingest` (add to RAG corpus), `both` (default: `both`) -- `conversation_id` — optional, associates context-mode content with a specific conversation - -Behavior: -- Validate file size against `MAX_UPLOAD_BYTES` — return 413 if exceeded -- Validate MIME type against `SUPPORTED_UPLOAD_TYPES` — return 415 if unsupported -- For PDF files, extract text using `pypdf` (add to requirements.txt) -- For all other types, read as UTF-8 text -- If mode includes `ingest`: chunk the extracted text into 512-token overlapping chunks (128-token overlap), generate embeddings via `EMBED_URL` (http://192.168.50.108:11434/api/embeddings, model mxbai-embed-large), upsert into Qdrant collection `caic` with metadata `{source: filename, upload_date: iso_timestamp, type: "upload"}` -- If mode includes `context`: store the full extracted text in a new SQLite table `upload_context` with columns `(id INTEGER PRIMARY KEY, conversation_id TEXT, filename TEXT, content TEXT, created_at TEXT, expires_at TEXT)`. Context entries expire after 1 hour. -- Return JSON: `{filename, size_bytes, mode, chunks_ingested (if ingest), context_id (if context), message}` - -**Add `upload_context` table to `db.py`** `init_db()`. - -**Wire `upload.router` into `app.py`** in the router registration block. - -**Write `tests/test_upload.py`** covering: -- Valid text file upload, mode=ingest — assert chunks_ingested > 0, Qdrant upsert called -- Valid text file upload, mode=context — assert context_id returned, row exists in upload_context -- Valid text file upload, mode=both — assert both behaviors -- File exceeds MAX_UPLOAD_BYTES — assert 413 -- Unsupported MIME type — assert 415 -- Guest session attempt — assert 403 -- PDF extraction path — mock pypdf, assert text extracted and processed - -Mock Qdrant and EMBED_URL calls via monkeypatch. Do not require live external services in tests. - -Run full test suite after implementation. All 26 existing tests must continue to pass. - ---- - -## ~~TASK 5 — File/Document Attachment: UI Integration [DONE]~~ - -**Status: Paperclip icon, file preview pill, gallery overlay, attachment indicators, DELETE/PATCH link/by-conversation endpoints, chat context injection. Committed `81238c0` (v1.10.0).** - -This task implements the frontend half of TODO #21. The UI is a single file at `templates/index.html`. - -Add a file attachment button to the chat input area. Requirements: -- Paperclip icon button adjacent to the send button -- Clicking opens a file picker filtered to supported types (`.txt`, `.md`, `.pdf`, `.json`, `.py`, `.html`) -- On file selection, show a pill/badge above the input showing the filename with an X to remove it -- On send, if a file is attached: POST to `/api/upload` with `mode=both` and the current `conversation_id`, then include the returned `context_id` in the subsequent `/api/chat` POST body as `upload_context_id` -- If the upload fails, show an inline error and do not send the chat message -- File attachment state clears after send - -**Update `/api/chat` in `routers/chat.py`:** -- Accept optional `upload_context_id` in the request body -- If present, look up the content in `upload_context` table and prepend it to the system prompt as: `\n\n[ATTACHED DOCUMENT: {filename}]\n{content}\n[END DOCUMENT]` -- If the context_id is expired or missing, log a warning and continue without it (do not error) - -**Add to `tests/test_chat_streaming_and_memory_paths.py`:** -- Test that a valid `upload_context_id` results in document content being prepended to the system prompt -- Test that an expired/missing `upload_context_id` is silently ignored - -Run full test suite. All existing tests must continue to pass. - ---- - -## ~~TASK 6 — Roadmap I: Terminal Command RAG Hook [DONE]~~ - -**Status: `POST /api/ingest` with Bearer token auth, `chunk_text()` shared helper, `caic-ingest.sh` script. Committed `1ac21ad` (v0.11.0).** - -This task implements autonomous RAG ingestion of significant terminal activity (TODO #23). - -**Create `routers/ingest.py`:** - -Implement `POST /api/ingest` (requires Bearer token auth — use same `COMPLETIONS_API_KEY` mechanism as `routers/completions.py`). Accept JSON body: -- `content` — string, the text to ingest (required) -- `source` — string, origin label e.g. `terminal`, `file`, `external` (default: `external`) -- `metadata` — optional dict of additional key/value pairs - -Behavior: -- Chunk `content` into 512-token overlapping chunks (128-token overlap) — extract this logic into a shared helper `chunk_text(text, chunk_size=512, overlap=128)` in `rag.py` if not already present -- Generate embeddings via `EMBED_URL` -- Upsert into Qdrant collection `caic` with metadata `{source, ingest_date: iso_timestamp, ...metadata}` -- Return JSON: `{chunks_ingested, source, message}` - -**Wire `ingest.router` into `app.py`.** - -**Create `/usr/local/bin/caic-ingest.sh` on worker (192.168.50.210)** — this is a shell script, not a Python file, and lives outside the repo. Write it to stdout/document it clearly so gramps can deploy it manually: - -```bash -#!/bin/bash -# caic-ingest.sh — pipe terminal commands into cAIc RAG -# Add to ~/.bashrc: export PROMPT_COMMAND="jc_capture" -# Function to call after significant commands - -JC_URL="http://192.168.50.210:8080/api/ingest" -JC_TOKEN="${CAIC_COMPLETIONS_API_KEY}" - -jc_capture() { - local cmd - cmd=$(history 1 | sed 's/^[ ]*[0-9]*[ ]*//') - # Only ingest significant commands - if echo "$cmd" | grep -qE '^(git|pip|systemctl|sudo|vi|vim|curl|wget|apt|python|pytest)'; then - curl -s -X POST "$JC_URL" \ - -H "Authorization: Bearer $JC_TOKEN" \ - -H "Content-Type: application/json" \ - -d "{\"content\": $(echo "$cmd" | python3 -c 'import json,sys; print(json.dumps(sys.stdin.read().strip()))'), \"source\": \"terminal\"}" \ - > /dev/null 2>&1 & - fi -} -``` - -**Write `tests/test_ingest.py`** covering: -- Valid ingest with content — assert chunks_ingested > 0 -- Missing Bearer token — assert 401 -- Wrong Bearer token — assert 403 -- Empty content — assert 422 -- Qdrant and embed calls mocked via monkeypatch - -Run full test suite. All existing tests must continue to pass. - ---- - -## ~~TASK 7 — Roadmap J: Startup Hardware Self-Assessment [DONE]~~ - -**Status: `hardware.py` + `routers/hardware.py` + 4 tests. Committed `7291b8f` (v0.12.0).** - -On jC startup, probe available hardware and write a living config snapshot. This replaces hardcoded assumptions about VRAM and RAM. - -**Create `hardware.py`** in the project root: - -``` -async def assess_hardware() -> dict -``` - -Probes: -- System RAM: `psutil.virtual_memory().total` and `.available` -- CPU count: `psutil.cpu_count()` -- GPU VRAM total and free: call `rocm-smi --showmeminfo vram --json` via subprocess, parse output. If rocm-smi absent or fails, set VRAM values to 0 and log a warning. -- llama-server reachable: GET `LLAMA_SERVER_BASE/v1/models`, timeout 3s. Record True/False and list of available model IDs. -- Qdrant reachable: GET `http://192.168.50.108:6333/collections`, timeout 3s. Record True/False and collection list. -- SearXNG reachable: GET `http://localhost:8888`, timeout 3s. Record True/False. - -Returns a dict with all of the above. Writes result as JSON to `hardware_state.json` in the working directory. - -**Call `assess_hardware()` from the FastAPI `lifespan` context** in `app.py` on startup, after `init_db()`. Log a summary line: `HW: {ram_gb}GB RAM, {vram_mb}MB VRAM, llama={reachable}, qdrant={reachable}, searxng={reachable}`. - -**Expose `GET /api/hardware`** in a new `routers/hardware.py` — returns the current `hardware_state.json` content as JSON. No auth required (read-only, non-sensitive aggregate stats). - -**Wire `hardware.router` into `app.py`.** - -**Write `tests/test_hardware.py`** covering: -- `assess_hardware()` with all services reachable (mock subprocess and httpx calls) — assert all fields present -- `assess_hardware()` with rocm-smi absent — assert VRAM=0, no exception raised -- `assess_hardware()` with llama-server unreachable — assert `llama_reachable=False`, no exception -- `GET /api/hardware` — assert returns JSON with expected keys - -Run full test suite. All existing tests must continue to pass. - ---- - -## ~~TASK 8 — Roadmap K: RAG Corpus Management [DONE]~~ - -Qdrant collection `caic` currently grows without bound. Implement score-based eviction with hysteresis, pinned sources, operational stats, and a flush command. - -### Config — add to `config.py`: - -```python -RAG_MAX_VECTORS = 50000 # absolute ceiling; eviction targets thresholds below it -RAG_EVICTION_HIGH_WATER = 0.80 # fraction of RAG_MAX_VECTORS that triggers eviction -RAG_EVICTION_LOW_WATER = 0.20 # fraction where eviction stops -RAG_EVICTION_BATCH = 1000 # max points to delete per Qdrant scroll/delete cycle -RAG_PINNED_SOURCES = ["upload", "profile"] # never evicted -RAG_GRACE_HOURS = 1 # new vectors ineligible for eviction until this old -RAG_ACCESS_WEIGHT = 1.0 # score factor: retrieval_count * ACCESS_WEIGHT -RAG_AGE_WEIGHT = 0.1 # score factor: ingest_age_hours * AGE_WEIGHT -``` - -Validations on boot: `high_water > low_water`, `batch > 0`, `max_vectors > 0`. - -### Eviction algorithm — add to `rag.py`: - -``` -score = (retrieval_count * ACCESS_WEIGHT) + (age_hours * AGE_WEIGHT) -``` - -Lower score = evicted first. Tiebreak: `last_accessed` ASC (older wins). - -```python -async def get_collection_count() -> int - # GET /collections/caic → return vectors_count - -async def get_collection_stats() -> dict - # Return {vector_count, max_vectors, high_water, low_water, percent_full, pinned_sources} - -async def evict_batch(batch_size: int) -> int - # Scroll Qdrant for vectors NOT in RAG_PINNED_SOURCES, WHERE ingest_age > RAG_GRACE_HOURS, - # ordered by score ASC, last_accessed ASC. - # Delete up to batch_size. Return count deleted. - # If 0 evictable vectors found: log warning, return 0 (break loop). - -async def maybe_evict() -> int - # Acquire eviction_lock (asyncio.Lock). - # count = get_collection_count() - # threshold_high = RAG_MAX_VECTORS * RAG_EVICTION_HIGH_WATER - # threshold_low = RAG_MAX_VECTORS * RAG_EVICTION_LOW_WATER - # total_evicted = 0 - # while count >= threshold_low: - # if total_evicted > 0 and count < threshold_low: break - # deleted = evict_batch(RAG_EVICTION_BATCH) - # if deleted == 0: break # no more unpinned targets - # total_evicted += deleted - # count -= deleted - # if count < threshold_high and total_evicted > 0: break - # # only one pass if batch spans the full gap - # if count < threshold_low: break - # Record total_evicted + timestamp in EVICTION_LOG (list of dicts, kept in memory, max 1000 entries) - # Release lock. Return total_evicted. - -async def get_rag_operational_stats() -> dict - # Returns: vector_count, max_vectors, high_water_pct, low_water_pct, - # percent_full, pinned_sources, grace_hours, - # eviction_counts_last_1m, eviction_counts_last_5m, eviction_counts_last_30m, - # at_risk_count (vectors in bottom 10% by score), - # pinned_count, avg_retrieval_count -``` - -### Edge cases & guards: - -1. **Newborn grace** — vectors < `RAG_GRACE_HOURS` old are excluded from eviction scroll (score=0 otherwise → immediate deletion) -2. **All-pinned freeze** — if scroll returns 0 evictable vectors, log warning and break loop -3. **Race** — `asyncio.Lock()` guards `maybe_evict()`; concurrent callers wait their turn -4. **Zero config** — `RAG_MAX_VECTORS <= 0` → eviction disabled; `RAG_EVICTION_BATCH <= 0` → clamped to 1 -5. **Legacy payloads** — vectors without `retrieval_count` or `last_accessed` get defaults (0, `ingest_date`) - -### Wire eviction: - -Call `maybe_evict()` after each upsert batch completes in: -- `routers/upload.py` — after Qdrant upsert -- `routers/ingest.py` — after Qdrant upsert - -### Admin endpoints — new `routers/rag_admin.py`: - -| Method | Endpoint | Description | -|--------|----------|-------------| -| GET | `/api/rag/stats` | Operational stats (see `get_rag_operational_stats()`) — admin required | -| POST | `/api/rag/flush` | Delete ALL points from the Qdrant `caic` collection. Returns `{deleted_count, collection: "caic", status: "flushed"}`. Admin required. | - -### In-memory eviction log: - -```python -EVICTION_LOG: list[dict] = [] # managed by rag.py, max 1000 entries -# Each entry: {timestamp: iso, count: N, remaining: N} -# Tied to RATE_EVENTS pattern from security.py for rolling window calculations -``` - -### Tests — `tests/test_rag_management.py`: - -- `get_collection_count()` — mock Qdrant GET, assert correct count -- `get_collection_stats()` — assert shape matches config -- `evict_batch()` — mock Qdrant scroll + delete, assert pinned sources excluded, grace period enforced, batch size respected -- `maybe_evict()` — below high water: 0 evicted; at high water: eviction fires; stops at low water; all-pinned scroll returns 0 → breaks -- `GET /api/rag/stats` — assert full shape -- `POST /api/rag/flush` — assert points deleted, admin required, guest 403 -- `POST /api/rag/flush` by guest — assert 403 -- Race lock — concurrent calls to `maybe_evict()` queue up, only one evicts - -Mock all Qdrant calls via monkeypatch. Do not require live services. - -Run full test suite. All existing tests must continue to pass. - ---- - -## ~~TASK 9 — Roadmap N1: RabbitMQ Install and Service on Coordinator (Infrastructure) [DONE]~~ - -This task runs on coordinator (this machine). Install RabbitMQ and verify it is operational. - -Run the following steps: -1. `apt-get update && apt-get install -y rabbitmq-server` -2. `systemctl enable rabbitmq-server && systemctl start rabbitmq-server` -3. `systemctl status rabbitmq-server` — verify active/running -4. Enable the management plugin: `rabbitmq-plugins enable rabbitmq_management` -5. Create a dedicated jC vhost: `rabbitmqctl add_vhost caic` -6. Create a dedicated user: `rabbitmqctl add_user caic CHANGEME_PASSWORD` — generate a random 24-char alphanumeric password and record it -7. Grant permissions: `rabbitmqctl set_permissions -p caic caic ".*" ".*" ".*"` -8. Verify management UI is reachable: `curl -s -u guest:guest http://localhost:15672/api/overview | python3 -m json.tool` -9. Delete default guest user: `rabbitmqctl delete_user guest` - -Declare the two topic exchanges needed by jC: -- Exchange name: `jc.admin`, type: `topic`, durable: true -- Exchange name: `jc.system`, type: `topic`, durable: true - -Use `rabbitmqadmin` or `curl` against the management API to declare exchanges. Verify both exchanges appear in: `curl -s -u caic:{password} http://localhost:15672/api/exchanges/caic` - -Write the generated RabbitMQ password to `/home/gramps/.caic_amqp_secret` with mode 600. This will be read by cAIc as an env var source in subsequent tasks. - -No pytest tests required for this infrastructure task. - ---- - -## ~~TASK 10 — Roadmap N2: AMQP Connection Layer in jC [DONE]~~ - -This task adds the core AMQP connection manager to jC. It must connect to RabbitMQ on coordinator (localhost from jC's perspective since jC runs on coordinator), handle reconnection, and provide a shared channel for all AMQP operations. - -**Add to `requirements.txt`:** `aio-pika>=9.0.0` - -**Add to `config.py`:** -- `AMQP_URL` — read from env `CAIC_AMQP_URL`, default `amqp://caic:password@localhost:5672/caic`. The actual password comes from `/home/gramps/.caic_amqp_secret` — read it at startup if the env var is not set. -- `AMQP_RECONNECT_DELAY` — seconds between reconnect attempts, default 5 -- `AMQP_EXCHANGE_ADMIN` — `jc.admin` -- `AMQP_EXCHANGE_SYSTEM` — `jc.system` - -**Create `amqp.py`** in the project root: - -```python -# Manages a single persistent aio-pika connection and channel. -# Provides: -# connect() -> None # establish connection, declare exchanges -# disconnect() -> None # graceful close -# get_channel() # returns current channel, reconnects if needed -# publish(exchange, routing_key, payload: dict) -> None -# # publishes JSON-serialized payload as persistent message -``` - -Connection must: -- Reconnect automatically on disconnect with `AMQP_RECONNECT_DELAY` backoff -- Log connection events at INFO level -- Not raise on publish if disconnected — log error and return (fire-and-forget, jC must not crash if RabbitMQ is down) - -**Start AMQP connection in `app.py` lifespan** after `assess_hardware()`. Disconnect in lifespan cleanup. - -**Write `tests/test_amqp.py`** covering: -- `publish()` with mocked aio-pika connection — assert message published with correct exchange and routing key -- `publish()` when disconnected — assert no exception raised, error logged -- `get_channel()` when connection is None — assert reconnect attempted - -Mock all aio-pika calls via monkeypatch. Do not require a live RabbitMQ instance in tests. - -Run full test suite. All existing tests must continue to pass. - ---- - -## ~~TASK 11 — Roadmap N3: Cluster Protocol & Registration Handler (Coordinator Side) [DONE]~~ - -**Status: Implemented and pushed (899988c).** `amqp.py` subscribe/rebind, `cluster.py` with CLUSTER_NODES/CLUSTER_EVENTS/CLUSTER_COORDINATOR and 6 handlers, `routers/cluster.py` (`GET /api/cluster`), 13 tests. No passive heartbeats — ping/pong on-demand before work routing. 148 tests pass. - -jC on the coordinator must listen for nine message types across `jc.admin` and `jc.system`, maintain the cluster registry, and expose an application-level event log. - -### 11.1 AMQP Protocol — Message Catalog - -All payloads are JSON, published as persistent messages. - -| Direction | Exchange | Routing Key | Message Type | Description | -|-----------|----------|-------------|-------------|-------------| -| Worker → Coordinator | `jc.admin` | `node.{name}.register` | register | Worker requests admission | -| Worker → Coordinator | `jc.admin` | `node.{name}.deregister` | deregister | Worker signals graceful departure | -| Coordinator → Worker | `jc.admin` | `node.{name}.admitted` | admitted | Coordinator grants admission | -| Coordinator → Worker | `jc.admin` | `node.{name}.rejected` | rejected | Coordinator denies admission (with reason) | -| Coordinator → Worker | `jc.admin` | `node.{name}.ping` | ping | Coordinator checks if worker is alive (sent before routing work) | -| Worker → Coordinator | `jc.admin` | `node.{name}.pong` | pong | Worker confirms aliveness | -| Worker → Coordinator | `jc.system` | `node.{name}.event` | event | Application-level syslog event | -| Any → All | `jc.system` | `cluster.coordinator.query` | coord_query | Anyone asks "who is coordinator?" | -| Coordinator → All | `jc.system` | `cluster.coordinator.response` | coord_response | Coordinator announces itself | - -Worker presence is assumed from registration onward. No periodic heartbeats — a worker can sit idle for days without chatter. When the coordinator needs to route work to a worker, it pings first; if the worker doesn't pong within timeout, the coordinator deregisters it and moves to the next node. - -### 11.2 Payload Schemas - -**register** (worker → coordinator): -```json -{ - "node_name": "worker01", - "node_type": "worker", - "ip": "192.168.50.210", - "capabilities": { - "gpu": true, "gpu_type": "amd", "vram_mb": 8192, - "cpu_cores": 8, "ram_gb": 16 - }, - "active_model": { - "name": "llama3.1", "version": "latest", "quant": "Q4_K_M", - "path": "/var/lib/caic/models/llama3.1-latest-Q4_K_M.gguf", - "port": 8081 - }, - "inventory": [ - {"name": "llama3.1", "version": "latest", "quant": "Q4_K_M", - "path": "/var/lib/caic/models/llama3.1-latest-Q4_K_M.gguf", "port": 8081} - ], - "status": "active" -} -``` - -**deregister** (worker → coordinator): -```json -{ - "node_name": "worker01", - "reason": "shutdown", - "timestamp": "2026-07-06T12:00:00Z" -} -``` - -**ping** (coordinator → worker): -```json -{ - "from": "coordinator", - "node_name": "worker01", - "type": "ping", - "correlation_id": "550e8400-e29b-41d4-a716-446655440000", - "timestamp": "2026-07-06T12:00:00Z" -} -``` -Worker must respond within 5 seconds or the coordinator considers it absent. - -**pong** (worker → coordinator): -```json -{ - "node_name": "worker01", - "type": "pong", - "correlation_id": "550e8400-e29b-41d4-a716-446655440000", - "status": "active", - "active_model": {"name": "llama3.1", "port": 8081}, - "load": {"cpu_pct": 45, "ram_pct": 62, "vram_pct": 38}, - "timestamp": "2026-07-06T12:00:00Z" -} -``` -Correlation ID matches the ping so the coordinator can pair request and response. - -**coord_query** (any → `cluster.coordinator.query`): -```json -{"type": "coord_query", "timestamp": "2026-07-06T12:00:00Z"} -``` -Coordinator responds on `cluster.coordinator.response`: -```json -{ - "coordinator_node": "coordinator", - "cluster_nodes": ["worker01"], - "timestamp": "2026-07-06T12:00:00Z" -} -``` - -**event** (worker → coordinator): -```json -{ - "node_name": "worker01", - "severity": "info", - "message": "llama-server started with model llama3.1:latest", - "details": {"model": "llama3.1:latest", "port": 8081, "pid": 1234}, - "timestamp": "2026-07-06T12:00:00Z" -} -``` -Severity levels: `info`, `warn`, `error`, `critical`. The coordinator assigns `category: "application"` based on the exchange (jc.system). No `event_type` field — the category is determined by the channel, not the payload. - -### 11.3 Design — Status Transitions Drive the Event Log - -All admin-level events are *derived* from `register()` and `deregister()` as side effects. There are no separate message types for coordinator election, node staleness, quarantine, or release — those are status transitions that `register()`/`deregister()` emit into `CLUSTER_EVENTS` locally. - -**Node status lifecycle:** - -``` -UNKNOWN ──register()──▶ active ──deregister()──▶ (removed) - │ - ping timeout│(coordinator publishes - │ deregister on its behalf) - ▼ - (removed) -``` - -**Coordinator status lifecycle:** - -``` -NONE ──register(node_type=coordinator)──▶ CLUSTER_COORDINATOR set - │ - deregister()│or timeout - ▼ - CLUSTER_COORDINATOR cleared -``` - -**Event categories — two buckets, no granular types:** - -| Category | When | severity | -|----------|------|----------| -| `cluster` | Node lifecycle, coordinator changes, model swaps, node offline — everything on `jc.admin` | `info` / `warn` / `error` | -| `application` | Worker syslog events (incoming on `jc.system` `node.*.event`) | `info` / `warn` / `error` / `critical` | - -Every `_push_event()` call uses one of these two categories. The `message` field carries the human-readable detail — no need for event type strings. The reporting tool filters by category + severity. - -**Channel split — security rationale:** - -The two exchanges are not an organizational convenience. They enforce a **data isolation boundary**: - -| Exchange | Contains | Exposed to | -|----------|----------|------------| -| `jc.admin` | Node lifecycle, heartbeats, model swaps, coordinator changes | Operations / machine-room staff | -| `jc.system` | Application events — inference queries, RAG context, user-facing data | Application-layer audit only | - -`jc.system` events can leak information about what users are doing and asking. The split ensures a sysadmin monitoring cluster health never accidentally consumes user-data-bearing events. The channels can be locked down independently — different AMQP credentials, separate queue permissions, different in-transit encryption policies if needed later. - -### 11.4 Implementation - -**Add to `amqp.py`:** - -```python -_SUBSCRIPTIONS: list[tuple[str, str, Callable]] # (exchange, routing_key, callback) - -async def subscribe(exchange, routing_key, callback) -> None - # Append to _SUBSCRIPTIONS list - # Declare a unique queue per subscription (name: f"jc.{exchange}.{sanitized_routing_key}") - # Bind queue to exchange/routing_key, consume with callback -``` - -Each subscription gets its own queue so multiple subscribers on different routing keys all receive messages. On reconnect: drain old consumers, iterate `_SUBSCRIPTIONS`, re-declare and re-bind each one. The `connect()` function must call `_rebind_subscriptions()` after exchanges are declared. - -**Create `cluster.py`** in the project root: - -```python -# In-memory cluster registry + event log -# Survives only while jC is running (not persisted) - -CLUSTER_NODES: dict[str, NodeRecord] -CLUSTER_EVENTS: deque[EventRecord] # bounded at 1000 entries -CLUSTER_COORDINATOR: str | None # node_name of active coordinator - -# NodeRecord fields: -# node_name, node_type, ip, status, active_model, inventory, -# capabilities: {gpu, gpu_type, vram_mb, cpu_cores, ram_gb} -# registered_at, last_seen - -# EventRecord: -# category: str ("cluster" | "application") -# severity: str ("info" | "warn" | "error" | "critical") -# node_name: str -# message: str -# details: dict | None -# timestamp: str - -def _push_event(category, severity, node_name, message, details=None) -> None - # Append EventRecord to CLUSTER_EVENTS, pop left if > 1000 - -async def handle_registration(message) -> None - # Parse payload, validate required fields (node_name, node_type, ip, active_model, inventory) - # Reject if node_name duplicate and CLUSTER_NODES[node_name].status == "active" - # If CLUSTER_COORDINATOR is None AND node_type == "coordinator": - # set CLUSTER_COORDINATOR = node_name - # _push_event("cluster", "info", node_name, "elected coordinator") - # publish cluster.coordinator.response on jc.system {coordinator_node, cluster_nodes, timestamp} - # Add node to CLUSTER_NODES with status="active" - # _push_event("cluster", "info", node_name, f"admitted as {node_type}") - # publish admitted on jc.admin node.{name}.admitted {node_name, timestamp, amqp_url} - -async def handle_deregistration(message) -> None - # Parse payload (node_name, reason, timestamp) - # If node_name == CLUSTER_COORDINATOR: - # clear CLUSTER_COORDINATOR - # _push_event("cluster", "warn", node_name, f"coordinator lost — {reason}") - # _push_event("cluster", "info", node_name, f"departed — {reason}") - # Remove node from CLUSTER_NODES, log it - -async def handle_pong(message) -> None - # Parse: node_name, correlation_id, status, active_model, load, timestamp - # Match correlation_id to outstanding ping - # If node in CLUSTER_NODES: update last_seen, status, active_model - # Signal the waiting caller that the node is alive - # If node unknown: log warning, do NOT auto-admit - -async def handle_event(message) -> None - # Parse: node_name, severity, message, details, timestamp - # Assigns category="application" (incoming on jc.system) - # Append EventRecord to CLUSTER_EVENTS (pop left if > 1000) - -async def handle_coordinator_query(message) -> None - # Respond on jc.system cluster.coordinator.response - # Payload: {coordinator_node, cluster_nodes: list(CLUSTER_NODES.keys()), timestamp} - -def get_cluster_state() -> dict - # Return: {nodes: CLUSTER_NODES, coordinator: CLUSTER_COORDINATOR, - # events: last 50 CLUSTER_EVENTS} -``` - -**Subscribe in `app.py` lifespan** after AMQP connects: - -| Exchange | Routing Key | Handler | -|----------|-------------|---------| -| `jc.admin` | `node.*.register` | `handle_registration` | -| `jc.admin` | `node.*.deregister` | `handle_deregistration` | -| `jc.admin` | `node.*.pong` | `handle_pong` | -| `jc.system` | `node.*.event` | `handle_event` | -| `jc.system` | `cluster.coordinator.query` | `handle_coordinator_query` | - -### 11.5 API — `GET /api/cluster` - -New router `routers/cluster.py`: -- `GET /api/cluster` — returns full cluster state: `{nodes, coordinator, events}` (last 50 events). No auth required. - -Wire `cluster.router` into `app.py`. - -### 11.6 Tests — `tests/test_cluster.py` - -Mock all aio-pika calls. Do not require live RabbitMQ. - -| # | Test | What it asserts | -|---|------|-----------------| -| 1 | Valid worker registration | Node admitted, CLUSTER_NODES updated, `cluster` event logged, `admitted` message published | -| 2 | First coordinator auto-promotion | CLUSTER_COORDINATOR set, `cluster` event with "elected" message, `coord_response` published | -| 3 | Duplicate node name rejected | `rejected` message with reason=`duplicate_node_name`, `cluster` event logged | -| 4 | Malformed payload rejected | `rejected` message with reason=`malformed_payload` | -| 5 | Graceful deregistration | Node removed, `cluster` event logged. If coordinator: CLUSTER_COORDINATOR cleared | -| 6 | Pong from known node | last_seen updated, load/status refreshed | -| 7 | Pong from unknown node | Warning logged, node NOT added | -| 8 | Event stored in log | Event appended to CLUSTER_EVENTS; at 1001 entries the oldest is popped | -| 9 | Coordinator query produces response | Response published with coordinator name and node list | -| 10 | GET /api/cluster shape | Response contains `nodes`, `coordinator`, `events` keys | - -Run full test suite. All existing tests must continue to pass. - ---- - -## ~~TASK 12 — Roadmap N4: Worker Node Registration Publisher (Worker Side) [DONE]~~ - -This task creates the worker node AMQP client that runs on worker (192.168.50.210). It is a standalone Python script — not part of the jC FastAPI app — that runs as a systemd service on worker. - -**Create `node_agent/agent.py`** in the repo (new directory). - -### 12.1 Config & Inventory Discovery - -On start, reads `/etc/caic-node-agent.conf` (INI format): -- `node_name` — hostname, default from `socket.gethostname()` -- `node_ip` — LAN IP, default from socket -- `node_type` — `"worker"` (fixed) -- `capabilities` — comma-separated list, e.g. `llm,rag` -- `amqp_url` — RabbitMQ URL on coordinator, e.g. `amqp://caic:password@192.168.50.108:5672/caic` -- `llama_port` — port llama-server/llama-rpc is listening on, default 8081 -- `models_dir` — path to GGUF model files, default `/var/lib/caic/models` -- `active_model` — filename of currently active model (without path) - -Discovers inventory by globbing `models_dir` for `*.gguf` files and parsing name/version/quant from filename using regex pattern: `{name}-{version}-{quant}.gguf` where quant matches `Q[0-9]+_K_[A-Z]+` or similar standard suffixes. - -### 12.2 Registration - -Publishes registration to `jc.admin`, routing key `node.{node_name}.register`: -```json -{ - "node_name": "worker01", - "node_type": "worker", - "ip": "192.168.50.210", - "capabilities": ["llm"], - "active_model": {"name": "...", "version": "...", "quant": "...", "path": "...", "port": 8081} -} -``` - -### 12.3 Admission Response - -Listens on `node.{node_name}.admitted` and `node.{node_name}.rejected` (both `jc.admin`). Logs result. If rejected, exits with error. - -### 12.4 Ping Listener - -After admission: listens on `jc.admin`, routing key `node.{node_name}.ping`. On receipt, responds immediately (within 1 second) with a pong on `jc.admin`, routing key `node.{node_name}.pong`: - -```json -{ - "node_name": "worker01", - "type": "pong", - "correlation_id": "", - "status": "active", - "active_model": {"name": "...", "version": "...", "quant": "...", "path": "...", "port": 8081}, - "load": {"cpu_pct": 45, "ram_pct": 62, "vram_pct": 38}, - "timestamp": "" -} -``` - -No periodic heartbeats. Worker sits idle between pings — coordinator only pings when it needs to route work. - -### 12.5 Model Swap Command Handler - -Listens on `jc.admin`, routing key `node.{node_name}.cmd.swap_model`: -- Payload: `{model_filename: str}` -- Stops current llama-server: `systemctl stop llama-server` -- Updates `/etc/caic-node-agent.conf` active_model field -- Starts llama-server: `systemctl start llama-server` (assumes service reads active_model from conf or ExecStart is updated) -- Waits for llama-server to be healthy: poll `http://localhost:{llama_port}/v1/models` every 2s, timeout 120s -- Publishes to `jc.system`, routing key `node.{node_name}.model_ready`: - ```json - {"node_name": "...", "active_model": "...", "port": ..., "timestamp": "..."} - ``` -- If startup fails within timeout: publishes `node.{node_name}.model_failed` with error detail - -### 12.6 Files & Tests - -**Create `node_agent/requirements.txt`:** `aio-pika>=9.0.0` - -**Document `/etc/caic-node-agent.conf` format** in a comment block at the top of `agent.py`. - -**Write `tests/test_node_agent.py`** covering: -- Registration payload construction from config + model discovery — assert correct JSON shape -- Model swap command handler: success path — assert systemctl calls made, model_ready published -- Model swap command handler: timeout path — assert model_failed published -- Ping handler: on ping, publishes pong with correct correlation_id -- Agent starts idle after admission, no heartbeat timer - -Mock all aio-pika, subprocess, and httpx calls. - -**Do not create a systemd service file in this task** — that is a manual deployment step. Document the required service configuration in a comment at the bottom of `agent.py`. - -Run full test suite. All existing tests must continue to pass. - ---- - -## ~~TASK 13 — Roadmap N5: Query Routing via AMQP + Phi-4-mini Triage [DONE]~~ - -This task wires the cluster into jC's chat flow. When a query arrives at `/api/chat`, instead of always routing to the hardcoded `LLAMA_SERVER_BASE`, jC now routes to the best available cluster node based on query context. - -**Prerequisites:** Tasks 9–12 complete. At least one worker node admitted to cluster. - -**Install Phi-4-mini on coordinator (infrastructure step):** -- Download `Phi-4-mini-Instruct-Q4_K_M.gguf` from HuggingFace using `hf download microsoft/Phi-4-mini-instruct --include "*.Q4_K_M.gguf" --local-dir /var/lib/caic/models` -- Create `/etc/systemd/system/llama-server-triage.service` — same pattern as existing llama-server service but: port 8083, model path points to Phi-4-mini GGUF, no `--rpc` flag (runs entirely on coordinator CPU/iGPU), description `Llama.cpp Server (Phi-4-mini — triage/routing)` -- `systemctl daemon-reload && systemctl enable llama-server-triage && systemctl start llama-server-triage` -- Verify: `curl -s http://localhost:8083/v1/models` - -**Add to `config.py`:** -- `TRIAGE_BASE` — `http://127.0.0.1:8083/v1` (Phi-4-mini) -- `TRIAGE_TIMEOUT` — 10 seconds -- `FALLBACK_TO_DEFAULT` — True (if triage fails or no nodes available, fall back to `LLAMA_SERVER_BASE`) - -**Create `triage.py`** in the project root: - -```python -async def classify_query(query: str) -> str - # Sends query to Phi-4-mini at TRIAGE_BASE with a classification system prompt. - # System prompt instructs model to respond with ONLY one of: - # "general", "code", "search", "rag" - # Returns the classification string. - # Timeout: TRIAGE_TIMEOUT seconds. - # On any error: returns "general" (fail-safe). - -async def select_node(classification: str) -> dict | None - # Consults CLUSTER_NODES from cluster.py - # For "code": prefer nodes where active_model name contains "coder" or "qwen" - # For "general": prefer nodes where active_model name contains "mistral" or "llama" - # For "search" or "rag": return None (handled locally by jC) - # If no matching node found: return None (triggers FALLBACK_TO_DEFAULT) - # Returns NodeRecord dict for selected node, or None - -async def get_inference_url(query: str) -> str - # Combines classify_query + select_node - # Returns full base URL: f"http://{node.ip}:{node.active_model.port}/v1" - # Falls back to LLAMA_SERVER_BASE if classification=search/rag, no nodes, or triage error -``` - -**Update `routers/chat.py`:** -- Replace the hardcoded `LLAMA_SERVER_BASE` reference with a call to `get_inference_url(user_message)` -- The rest of the chat flow (RAG, memory, streaming) is unchanged — only the inference target URL changes - -**Write `tests/test_triage.py`** covering: -- `classify_query()` returns valid classification — mock Phi-4-mini response -- `classify_query()` on timeout — assert returns "general", no exception -- `select_node("code")` with coder node in cluster — assert correct node returned -- `select_node("general")` with no matching node — assert None returned -- `get_inference_url()` with code query and coder node available — assert returns node URL -- `get_inference_url()` with no nodes in cluster — assert returns LLAMA_SERVER_BASE fallback - -**Update `tests/test_chat_streaming_and_memory_paths.py`:** -- Mock `triage.get_inference_url` to return a fixed URL in all existing tests so they continue to pass without a live cluster - -Run full test suite. All existing tests must continue to pass. - ---- - -## ~~TASK 14 — Roadmap N6: Model Swap Command Flow [DONE]~~ - -**Status: Implemented and pushed (`9d1fd44`).** `request_model_swap()`, `handle_model_ready()`, `handle_model_failed()` in `cluster.py`, async `select_node()` with swap triggering in `triage.py`, `tests/test_model_swap.py` (9 tests). 177 tests pass. - -This task implements the coordinator-side logic for requesting a model swap on a worker node when the ideal model is not currently active. - -**Add to `cluster.py`:** - -```python -async def request_model_swap(node_name: str, model_filename: str) -> bool - # Publishes to jc.admin exchange, routing key node.{node_name}.cmd.swap_model - # Payload: {model_filename, requested_at: iso_timestamp} - # Sets node status to "swapping" in CLUSTER_NODES - # Returns True if message published successfully - -async def handle_model_ready(message) -> None - # Handles node.{node_name}.model_ready from jc.system - # Updates CLUSTER_NODES[node_name].active_model to the new model - # Sets node status back to "active" - # Logs swap completion with timing - -async def handle_model_failed(message) -> None - # Handles node.{node_name}.model_failed from jc.system - # Sets node status to "error" in CLUSTER_NODES - # Logs failure with detail from message payload -``` - -**Subscribe in `app.py` lifespan:** -- `jc.system` exchange, routing key `node.*.model_ready` → `handle_model_ready` -- `jc.system` exchange, routing key `node.*.model_failed` → `handle_model_failed` - -**Update `triage.py` `select_node()`:** -- If the best-matching node exists but its active_model does not match the ideal model for the classification, AND the node status is "active" (not already swapping): - - Call `request_model_swap(node_name, ideal_model_filename)` - - Return None (triggers fallback) — the swap happens async, next query will find the right model active -- If node status is "swapping": return None (fallback, swap in progress) - -**Update `GET /api/cluster`** to include node status in response. - -**Write `tests/test_model_swap.py`** covering: -- `request_model_swap()` — assert swap command published, node status set to "swapping" -- `handle_model_ready()` — assert active_model updated, status set to "active" -- `handle_model_failed()` — assert status set to "error" -- `select_node()` with mismatched active model — assert swap requested, None returned -- `select_node()` with node status "swapping" — assert None returned without publishing another swap - -Run full test suite. All existing tests must continue to pass. - ---- - -## ~~TASK 15 — Roadmap N7: Cluster Status UI [DONE]~~ - -Surface cluster awareness in the jC frontend (`templates/index.html`). - -**Add a cluster status panel** to the UI. Requirements: -- Small status bar or collapsible panel, visible but unobtrusive -- Polls `GET /api/cluster` every 15 seconds -- For each admitted node: show node name, active model name, and a colored status dot: - - Green: active - - Yellow: swapping - - Red: error or offline (not seen in last 60 seconds based on last_seen timestamp) -- If no nodes in cluster (empty): show "No worker nodes connected" -- Panel must not interfere with chat input or conversation list - -**Update `GET /api/cluster` response** to include `last_seen` per node and a `status` field (`active`, `swapping`, `error`). - -**Update heartbeat handling in `cluster.py`:** add a handler for `node.*.heartbeat` on `jc.system` that updates `last_seen` timestamp for the node. - -**Subscribe in `app.py` lifespan:** -- `jc.system` exchange, routing key `node.*.heartbeat` → `handle_heartbeat` - -**Add `handle_heartbeat()` to `cluster.py`:** -- Updates `CLUSTER_NODES[node_name].last_seen` to current timestamp -- If node was previously marked offline (not in CLUSTER_NODES), log re-registration warning but do not auto-admit — full registration required - -**Write `tests/test_cluster_heartbeat.py`** covering: -- `handle_heartbeat()` for known node — assert last_seen updated -- `handle_heartbeat()` for unknown node — assert no crash, warning logged, node not added - -Run full test suite. All 26+ existing tests must continue to pass. - -~~Commit all changes introduced across Tasks 9–15 with message: `feat: Roadmap N — AMQP cluster nervous system complete`~~ - ---- - -## Backlog (Post-Roadmap N) ⏳ - -### ~~B1 — Context loss in follow-up questions [DONE]~~ - -**Symptom:** After asking "in {context}, explain {b}", a follow-up "what is {b}'s {x}?" gets a non-sequitur response that ignores the original context. - -**Diagnosis:** `build_system_prompt()` is called fresh per-request with new RAG/memory results keyed to the current message text. These can change between turns and may dilute or override the conversation history. The original system prompt used for turn 1 (including its RAG context) is not stored in the DB — only user/assistant messages are. The inference server receives a different system prompt each turn. - -**Possible fixes:** -- Store the assembled system prompt with each assistant message in the DB -- When replaying history, re-send the original system prompts from DB rather than rebuilding -- Or: cap RAG/memory injection to only fire on the first message of a conversation, then rely solely on conversation history for follow-ups -- Check that llama-server isn't truncating history due to context window overflow (Mistral-Nemo 12B = 128K context, unlikely) - -### ~~B2 — Bang-prefixed search routing [DONE]~~ - -**Spec:** If a query begins with `!`, route to SearXNG search instead of local inference. - -**Where:** In `routers/chat.py` `chat()` handler, after `user_message` is extracted. Strip the `!`, set a flag to always trigger auto-search regardless of perplexity/refusal. - -**Change:** Add a `force_search` flag when `user_message.startswith("!")`, strip the prefix from the message saved to DB, and route directly to the search+summarize path. - -### ~~B3 — Docker distribution (v1.0 gate) [DONE]~~ - -**Goal:** Ship cAIc as a `docker compose` stack so a single command stands up everything. - -**Services to containerize:** -- cAIc (FastAPI app + SQLite) -- SearXNG -- Qdrant -- RabbitMQ -- llama-server (with optional RPC sidecar for GPU offload) -- Ollama (embeddings) - -**Also needed:** -- `Dockerfile` for the cAIc app itself -- `docker-compose.yml` with all services, volumes, networks, env vars -- Setup wizard script (run on first boot) that: - - Probes CPU vs GPU (reuses `hardware.py`) - - Queries user for admin PIN, node name, IP - - Generates `.env` file with correct `LLAMA_SERVER_BASE`, `EMBED_URL`, etc. - - Auto-calculates `RAG_MAX_VECTORS` from available RAM: `max(1000, int(available_ram_gb * 100_000))` - - Optionally detects and configures RPC GPU offload -- Manual install docs remain alongside for bare-metal deployment - -**This task is only actionable after Tasks 8–15 (RAG eviction + AMQP cluster) are complete.** - ---- - -### ~~B4 — RAG Corpus Management UI (Display, Edit, CRUD) [DONE]~~ - -**Goal:** Provide a management interface in the UI to browse, search, edit, and delete individual entries in the Qdrant-backed RAG corpus. - -**Backend — add to `routers/rag_admin.py`:** - -| Method | Endpoint | Description | Auth | -|--------|----------|-------------|------| -| GET | `/api/rag/points` | Return paginated list of RAG points with payload (text, source, date). Supports `?offset=0&limit=50&search=` query params | Admin | -| GET | `/api/rag/point/{point_id}` | Return a single point with full payload | Admin | -| DELETE | `/api/rag/point/{point_id}` | Delete a single point from Qdrant | Admin | -| PATCH | `/api/rag/point/{point_id}` | Update a point's text payload (re-embed the new text) | Admin | - -Helper functions for Qdrant scroll/delete/update go in `rag.py` or `eviction.py`. - -**Frontend — add to `templates/index.html`:** - -A "RAG" button in the admin UI (drawer or settings modal) that opens a management panel: -- **Stats bar**: vector count, max vectors, percent full, pinned sources -- **Search bar**: text input to search the RAG corpus by semantic similarity -- **Results table**: paginated list showing each vector's text snippet, source label, ingest date, retrieval count - - Click to expand full text - - Delete button per row (with confirmation) - - Edit button per row (inline text edit → re-embed on save) -- **Bulk actions**: flush all (existing `/api/rag/flush`) with confirmation - -**Tests:** - -- `tests/test_rag_admin.py` — cover new endpoints: list, get, delete, update, admin-enforcement -- Mock all Qdrant calls via monkeypatch - -Run full test suite. All existing tests must continue to pass.** diff --git a/ai.md b/ai.md index 777062b..5707c47 100644 --- a/ai.md +++ b/ai.md @@ -37,6 +37,7 @@ Every router has a dedicated test file: | `test_cluster_heartbeat.py` | `cluster.py` — heartbeat handler, known/unknown node | | `test_model_swap.py` | `cluster.py` + `triage.py` — request_model_swap, handle_model_ready/failed, select_node swap triggering | | `test_node_agent.py` | `node_agent/agent.py` — registration, ping/pong, model swap | +| `test_image.py` | Image generation — cluster handlers, router proxy, node agent ComfyUI integration, hardware probe, capability detection | | `test_triage.py` | `triage.py` — classify_query, select_node, get_inference_url | | `test_settings_allowlist.py` | `routers/settings.py` — allowlisted key enforcement | | `test_skills_framework.py` | `routers/skills.py` — list, toggle, unknown skill, prompt injection | @@ -69,10 +70,10 @@ Refactored from single-file (`app.py`) into modules under project root: | `crypto.py` | AES-256-GCM encrypt/decrypt + key management (stored as `heartbeat_interval_ms` in settings) | | `model_pull.py` | Startup model availability check + Ollama pull API | | `triage.py` | Phi-4-mini-based query classification + cluster node selection | -| `cluster.py` | Cluster node registry, event log, coordinator election, ping/pong, model swap handlers | +| `cluster.py` | Cluster node registry, event log, coordinator election, ping/pong, model swap handlers, image generation request/response | | `amqp.py` | AMQP connection manager — connect, disconnect, publish, subscribe, auto-reconnect | -| `node_agent/` | Standalone worker agent — AMQP client for registration, ping/pong, model swap | -| `routers/` | One module per endpoint group (chat, search, skills, completions, upload, ingest) | +| `node_agent/` | Standalone worker agent — AMQP client for registration, ping/pong, model swap, image generation | +| `routers/` | One module per endpoint group (chat, search, skills, completions, upload, ingest, image) | ### Entrypoint / API keys @@ -88,6 +89,7 @@ Refactored from single-file (`app.py`) into modules under project root: 3. **`/v1/chat/completions`** → OpenAI-compatible for Continue.dev/IDE integration; FIM requests proxied without persistence 4. **`/api/upload`** → multipart file upload, PDF/text extraction, `mode=(context|ingest|both)`, stores SQLite context (1hr expiry) + Qdrant upsert 5. **`/api/ingest`** → Bearer token auth, programmatic RAG ingest (terminal hook, external tools) +6. **`POST /api/image/generate`** → admin required, routes to corsair node agent via AMQP → ComfyUI workflow → returns PNG; `GET /api/image/status` lists available image gen nodes ### Perplexity / auto-search @@ -123,6 +125,7 @@ All services are available bare-metal or as containers in `docker compose up`. | rocm-smi | No | AMD GPU stats | — | | Qdrant | No | 6333 (coordinator) — RAG vector search | `qdrant` | | Ollama (worker) | No | 11434 — embeddings + model pull | `ollama` | +| ComfyUI (worker) | No | 8188 — image generation API | — | ### Config quirks @@ -130,6 +133,8 @@ All services are available bare-metal or as containers in `docker compose up`. - `SUPPORTED_UPLOAD_TYPES` includes images (png/jpeg/gif/svg/webp) + text + PDF + JSON - `UPLOAD_CONTEXT_EXPIRY_HOURS` = 1 hour - Rate limits and payload caps in `config.py` — patch `security.RL_*` not `config.RL_*` for tests +- `COMFYUI_BASE` defaults to `http://192.168.50.115:8188` (overridable via `CAIC_COMFYUI_BASE`) +- `COMFYUI_TIMEOUT` defaults to `120` seconds (overridable via `CAIC_COMFYUI_TIMEOUT`) - RAG embedding requests go to `EMBED_URL` at `/api/embeddings` (Ollama on worker :11434) ### SSE Protocol @@ -152,16 +157,16 @@ All streaming endpoints yield `data: {json}\n\n`. Key shapes: - **Docker containerization (B3)**: Created `Dockerfile`, `docker-compose.yml`, `.env.example`, `scripts/setup.sh`, `.dockerignore`, `searxng-settings.yml.dist`, `models/README.txt`. Fixed hardcoded defaults in `config.py` (localhost, Docker secrets path, `CAIC_DEFAULT_MODEL` env var, `CAIC_HW_STATE_PATH` env var). Added missing `psutil` + `jinja2` to `requirements.txt`. Fixed test discovery via `tests/conftest.py` sys.path insertion. 214 tests pass. ### Active -- (none) +- Image generation service backend complete — cluster handlers, router proxy, node agent ComfyUI integration, hardware probe. 27 tests pass. ComfyUI install pending on corsair. ### Blocked -- (none) +- Ball Gunner assets — waiting on Canva designs ### Upcoming (backlog) - ~~B3 — Docker distribution~~ [DONE] ### Key config values (current) -- **Current VERSION**: `v1.0.0` in `config.py`. +- **Current VERSION**: `v1.1.0` in `config.py`. - `SESSION_TIMEOUT_SECONDS = 3600` - `DEFAULT_MODEL = "qwen2.5-7b-instruct"` (overridable via `CAIC_DEFAULT_MODEL`) - `LLAMA_SERVER_BASE = "http://localhost:8081"` (overridable via env var) diff --git a/app.py b/app.py index af323e5..05848b3 100644 --- a/app.py +++ b/app.py @@ -41,6 +41,7 @@ import routers.ingest as ingest import routers.hardware as hardware import routers.rag_admin as rag_admin import routers.cluster as cluster_router +import routers.image as image_router # --- Logging --- log = logging.getLogger("caic") @@ -177,7 +178,7 @@ for router_module in [ auth_router, conversations.router, memories.router, models.router, presets.router, profile.router, settings.router, skills.router, chat.router, search_route.router, completions.router, upload.router, ingest.router, hardware.router, - rag_admin.router, cluster_router.router, + rag_admin.router, cluster_router.router, image_router.router, ]: app.include_router(router_module) diff --git a/cluster.py b/cluster.py index a13a3fc..6c40e7d 100644 --- a/cluster.py +++ b/cluster.py @@ -18,6 +18,7 @@ CLUSTER_NODES: dict[str, dict] = {} CLUSTER_EVENTS: deque = deque(maxlen=1000) CLUSTER_COORDINATOR: str | None = None _pending_pings: dict[str, tuple[str, asyncio.Event]] = {} +_pending_image: dict[str, tuple[str, asyncio.Event]] = {} NODE_NAME: str = os.environ.get("CAIC_NODE_NAME", "ultron") PING_TIMEOUT: float = 5.0 @@ -240,6 +241,72 @@ async def handle_model_failed(exchange: str, routing_key: str, payload: dict) -> _push_event("cluster", "error", node_name, f"Model swap failed: {error}") +async def handle_image_generated(exchange: str, routing_key: str, payload: dict) -> None: + node_name = payload.get("node_name", routing_key.split(".")[1] if "." in routing_key else "unknown") + request_id = payload.get("request_id") + + if request_id and request_id in _pending_image: + _, event = _pending_image.pop(request_id) + _pending_image[request_id] = (payload.get("image_base64", ""), event) + event.set() + + if node_name in CLUSTER_NODES: + CLUSTER_NODES[node_name]["last_seen"] = datetime.now(timezone.utc).isoformat() + "Z" + + +async def handle_image_failed(exchange: str, routing_key: str, payload: dict) -> None: + node_name = payload.get("node_name", routing_key.split(".")[1] if "." in routing_key else "unknown") + request_id = payload.get("request_id") + error = payload.get("error", "unknown error") + + if request_id and request_id in _pending_image: + _pending_image[request_id] = ("", _pending_image[request_id][1]) + _pending_image[request_id][1].set() + + _push_event("application", "error", node_name, f"Image generation failed: {error}") + + +async def request_image_generate( + node_name: str, prompt: str, negative_prompt: str = "", + width: int = 1024, height: int = 1024, steps: int = 20, + seed: int = -1, model: str = "", timeout: float = 120, +) -> str | None: + if node_name not in CLUSTER_NODES: + log.warning("request_image_generate: unknown node %s", node_name) + return None + + caps = CLUSTER_NODES[node_name].get("capabilities", []) + if "image_gen" not in caps: + log.warning("request_image_generate: node %s lacks image_gen capability", node_name) + return None + + request_id = str(uuid.uuid4()) + event = asyncio.Event() + _pending_image[request_id] = ("", event) + + now = datetime.now(timezone.utc).isoformat() + "Z" + _push_event("application", "info", node_name, f"Image generation requested: {prompt[:60]}...") + + await publish(AMQP_EXCHANGE_ADMIN, f"node.{node_name}.cmd.image_generate", { + "from": NODE_NAME, "type": "image_generate", + "request_id": request_id, + "prompt": prompt, + "negative_prompt": negative_prompt, + "width": width, "height": height, + "steps": steps, "seed": seed, "model": model, + "timestamp": now, + }) + + try: + await asyncio.wait_for(event.wait(), timeout=timeout) + result = _pending_image.pop(request_id, (None, None)) + return result[0] + except asyncio.TimeoutError: + _pending_image.pop(request_id, None) + _push_event("application", "warn", node_name, "Image generation timed out") + return None + + SUBSCRIBE_TABLE = [ (AMQP_EXCHANGE_ADMIN, ["node.*.register"], handle_registration), (AMQP_EXCHANGE_ADMIN, ["node.*.deregister"], handle_deregistration), @@ -249,6 +316,8 @@ SUBSCRIBE_TABLE = [ (AMQP_EXCHANGE_SYSTEM, ["node.*.heartbeat"], handle_heartbeat), (AMQP_EXCHANGE_SYSTEM, ["node.*.model_ready"], handle_model_ready), (AMQP_EXCHANGE_SYSTEM, ["node.*.model_failed"], handle_model_failed), + (AMQP_EXCHANGE_SYSTEM, ["node.*.image_generated"], handle_image_generated), + (AMQP_EXCHANGE_SYSTEM, ["node.*.image_failed"], handle_image_failed), ] diff --git a/config.py b/config.py index 0b05d6e..b4f4903 100644 --- a/config.py +++ b/config.py @@ -10,7 +10,7 @@ from pathlib import Path log = logging.getLogger("caic") -VERSION = "v1.0.0" +VERSION = "v1.1.0" OLLAMA_BASE = os.environ.get("OLLAMA_BASE", "http://localhost:11434") LLAMA_SERVER_BASE = os.environ.get("LLAMA_SERVER_BASE", "http://localhost:8081") SEARXNG_BASE = os.environ.get("CAIC_SEARXNG_BASE", "http://localhost:8888") @@ -52,6 +52,10 @@ TRUST_X_FORWARDED_FOR = ( os.getenv("CAIC_TRUST_X_FORWARDED_FOR", "false").lower() == "true" ) +# --- Image generation (ComfyUI) --- +COMFYUI_BASE = os.environ.get("CAIC_COMFYUI_BASE", "http://192.168.50.115:8188") +COMFYUI_TIMEOUT = int(os.environ.get("CAIC_COMFYUI_TIMEOUT", "120")) + # --- Rate limits --- RATE_WINDOW_SECONDS = 60 RL_LOGIN_PER_WINDOW = 10 diff --git a/docker-compose.yml b/docker-compose.yml index 024f503..95a5ca7 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -19,6 +19,8 @@ services: - rabbitmq_password environment: - CAIC_AMQP_SECRET_PATH=/run/secrets/rabbitmq_password + - CAIC_COMFYUI_BASE=${CAIC_COMFYUI_BASE:-http://192.168.50.115:8188} + - CAIC_COMFYUI_TIMEOUT=${CAIC_COMFYUI_TIMEOUT:-120} env_file: .env depends_on: qdrant: { condition: service_started } diff --git a/docs/archive/TASKS-pre-1.0.md b/docs/archive/TASKS-pre-1.0.md new file mode 100644 index 0000000..acfa712 --- /dev/null +++ b/docs/archive/TASKS-pre-1.0.md @@ -0,0 +1,994 @@ +# cAIc — OpenCode Prompt Sequence +# Generated: 2026-07-14 +# Execute sequentially. Run full test suite after each task before proceeding. +# Test command: ./venv/bin/python -m pytest tests/ -v + +--- + +## Session 2026-07-14 — RAG bugfixes + Topbar redesign + +- **RAG bugs fixed**: Collection name mismatch (`jarvis_rag` → `caic_rag`, migrated 219 points), `vectors_count`→`points_count` (Qdrant v1.10+ API change), removed unindexed `order_by` that caused 502 on scroll, made `RAG_COLLECTION` env-configurable (`CAIC_RAG_COLLECTION`). +- **Semantic search fixed**: Set `CAIC_EMBED_URL=http://192.168.50.108:11434` (mxbai-embed-large lives on ultron, not the old embed server). +- **Topbar redesign**: Moved system stats (CPU/MEM/GPU/VRAM/TOK) to a centered bottom strip. Moved toggles (MEM, SEARCH, PROFILE, SORT, PRIVACY) into a ⋮ hamburger menu next to ADMIN badge. Palette icon sits immediately after version number in topbar-left. Removed standalone (i) button — privacy info accessible via ⋮ → About Privacy. Input bar above chat, stats at very bottom. Mobile-responsive padding/sizing. + +--- + +## ~~TASK 1 — README Cleanup [DONE]~~ + +Review README.md in the current repo. Remove any node references other than `coordinator` (192.168.50.108) and `worker` (192.168.50.210). Ensure all references to the project use the exact casing `cAIc` — not `Jarvischat`, `JarvisChat`, or `jarvischat`. Do not change any functional content, endpoint documentation, or architecture descriptions — this is a text cleanup only. After editing, verify the file renders cleanly as markdown. Commit with message: `docs: clean up node references and branding consistency`. + +No new tests required for this task. + +--- + +## ~~TASK 2 — Qwen2.5-Coder llama-server Service on Coordinator (Infrastructure) [DONE]~~ + +**Status: Systemd unit created, verified, and restored.** + +This task originally defined creation of `/etc/systemd/system/llama-server-coder.service` (port 8082, Qwen2.5-Coder-14B Q5_K_M) as a prerequisite for dynamic model swapping. That sysadmin work is done. + +**The real Task 2 deliverable — the ability to dynamically swap models based on query classification — is delivered by Roadmap N (Tasks 9–15).** The flow: + +1. **Task 13** — Phi-4-mini triage (`triage.py`) classifies the query as `general`, `code`, `search`, or `rag` +2. **Task 13** — `select_node()` picks the best worker node; if the ideal model isn't active, it triggers a swap +3. **Task 14** — `request_model_swap()` publishes `cmd.swap_model` via AMQP `jc.admin` exchange +4. **Task 12** — The node agent on worker receives the command, stops the current llama-server, starts the correct one, waits for health, and publishes `model_ready` +5. **Task 14** — coordinator receives `model_ready`, updates the cluster registry, and routes the query to the node + +The swap is async and transparent — the user sees only latency. The UI (Task 15) shows a yellow "swapping" status dot during the transition. + +The service unit at `/etc/systemd/system/llama-server-coder.service` is the **target** the node agent starts when swapping to code inference. It is not enabled at boot — the AMQP cluster manages activation. + +See Tasks 9–15 for the actual model swap implementation. + +No pytest tests required for this infrastructure task. + +--- + +## ~~TASK 3 — Update OpenCode Config to Use Qwen on :8082 [DONE]~~ + +Update `/home/gramps/.config/opencode/opencode.jsonc` (on this machine, coordinator) to point the configured provider at `http://127.0.0.1:8082/v1` instead of `http://127.0.0.1:8081/v1`. The model name in the config should be updated to reflect `qwen2.5-coder-14b` or whatever model ID the llama-server instance at :8082 reports via `/v1/models`. Verify the endpoint is reachable before writing the config change. Do not restart OpenCode — the config change takes effect on next session start. + +No pytest tests required for this task. + +--- + +## ~~TASK 4 — File/Document Attachment: Backend Ingest Endpoint [DONE]~~ + +**Status: `POST /api/upload` with mode=(context|ingest|both), PDF/text extraction, Qdrant upsert, SQLite context (1hr expiry). Committed `4a891c8` (v1.9.0).** + +This task implements the backend half of file/document attachment (TODO #21). The goal is dual-aspect upload: a file can be used as immediate chat context, ingested into the RAG corpus (Qdrant), or both. + +**Add to `config.py`:** +- `UPLOAD_DIR` — path for temporary upload storage, default `/tmp/caic_uploads` +- `MAX_UPLOAD_BYTES` — max file size, default 20MB +- `SUPPORTED_UPLOAD_TYPES` — set of MIME types: `text/plain`, `text/markdown`, `application/pdf`, `application/json`, `text/x-python`, `text/html` + +**Create `routers/upload.py`:** + +Implement `POST /api/upload` (admin required). Accept `multipart/form-data` with: +- `file` — the uploaded file (required) +- `mode` — string enum: `context` (inject into next chat only), `ingest` (add to RAG corpus), `both` (default: `both`) +- `conversation_id` — optional, associates context-mode content with a specific conversation + +Behavior: +- Validate file size against `MAX_UPLOAD_BYTES` — return 413 if exceeded +- Validate MIME type against `SUPPORTED_UPLOAD_TYPES` — return 415 if unsupported +- For PDF files, extract text using `pypdf` (add to requirements.txt) +- For all other types, read as UTF-8 text +- If mode includes `ingest`: chunk the extracted text into 512-token overlapping chunks (128-token overlap), generate embeddings via `EMBED_URL` (http://192.168.50.108:11434/api/embeddings, model mxbai-embed-large), upsert into Qdrant collection `caic` with metadata `{source: filename, upload_date: iso_timestamp, type: "upload"}` +- If mode includes `context`: store the full extracted text in a new SQLite table `upload_context` with columns `(id INTEGER PRIMARY KEY, conversation_id TEXT, filename TEXT, content TEXT, created_at TEXT, expires_at TEXT)`. Context entries expire after 1 hour. +- Return JSON: `{filename, size_bytes, mode, chunks_ingested (if ingest), context_id (if context), message}` + +**Add `upload_context` table to `db.py`** `init_db()`. + +**Wire `upload.router` into `app.py`** in the router registration block. + +**Write `tests/test_upload.py`** covering: +- Valid text file upload, mode=ingest — assert chunks_ingested > 0, Qdrant upsert called +- Valid text file upload, mode=context — assert context_id returned, row exists in upload_context +- Valid text file upload, mode=both — assert both behaviors +- File exceeds MAX_UPLOAD_BYTES — assert 413 +- Unsupported MIME type — assert 415 +- Guest session attempt — assert 403 +- PDF extraction path — mock pypdf, assert text extracted and processed + +Mock Qdrant and EMBED_URL calls via monkeypatch. Do not require live external services in tests. + +Run full test suite after implementation. All 26 existing tests must continue to pass. + +--- + +## ~~TASK 5 — File/Document Attachment: UI Integration [DONE]~~ + +**Status: Paperclip icon, file preview pill, gallery overlay, attachment indicators, DELETE/PATCH link/by-conversation endpoints, chat context injection. Committed `81238c0` (v1.10.0).** + +This task implements the frontend half of TODO #21. The UI is a single file at `templates/index.html`. + +Add a file attachment button to the chat input area. Requirements: +- Paperclip icon button adjacent to the send button +- Clicking opens a file picker filtered to supported types (`.txt`, `.md`, `.pdf`, `.json`, `.py`, `.html`) +- On file selection, show a pill/badge above the input showing the filename with an X to remove it +- On send, if a file is attached: POST to `/api/upload` with `mode=both` and the current `conversation_id`, then include the returned `context_id` in the subsequent `/api/chat` POST body as `upload_context_id` +- If the upload fails, show an inline error and do not send the chat message +- File attachment state clears after send + +**Update `/api/chat` in `routers/chat.py`:** +- Accept optional `upload_context_id` in the request body +- If present, look up the content in `upload_context` table and prepend it to the system prompt as: `\n\n[ATTACHED DOCUMENT: {filename}]\n{content}\n[END DOCUMENT]` +- If the context_id is expired or missing, log a warning and continue without it (do not error) + +**Add to `tests/test_chat_streaming_and_memory_paths.py`:** +- Test that a valid `upload_context_id` results in document content being prepended to the system prompt +- Test that an expired/missing `upload_context_id` is silently ignored + +Run full test suite. All existing tests must continue to pass. + +--- + +## ~~TASK 6 — Roadmap I: Terminal Command RAG Hook [DONE]~~ + +**Status: `POST /api/ingest` with Bearer token auth, `chunk_text()` shared helper, `caic-ingest.sh` script. Committed `1ac21ad` (v0.11.0).** + +This task implements autonomous RAG ingestion of significant terminal activity (TODO #23). + +**Create `routers/ingest.py`:** + +Implement `POST /api/ingest` (requires Bearer token auth — use same `COMPLETIONS_API_KEY` mechanism as `routers/completions.py`). Accept JSON body: +- `content` — string, the text to ingest (required) +- `source` — string, origin label e.g. `terminal`, `file`, `external` (default: `external`) +- `metadata` — optional dict of additional key/value pairs + +Behavior: +- Chunk `content` into 512-token overlapping chunks (128-token overlap) — extract this logic into a shared helper `chunk_text(text, chunk_size=512, overlap=128)` in `rag.py` if not already present +- Generate embeddings via `EMBED_URL` +- Upsert into Qdrant collection `caic` with metadata `{source, ingest_date: iso_timestamp, ...metadata}` +- Return JSON: `{chunks_ingested, source, message}` + +**Wire `ingest.router` into `app.py`.** + +**Create `/usr/local/bin/caic-ingest.sh` on worker (192.168.50.210)** — this is a shell script, not a Python file, and lives outside the repo. Write it to stdout/document it clearly so gramps can deploy it manually: + +```bash +#!/bin/bash +# caic-ingest.sh — pipe terminal commands into cAIc RAG +# Add to ~/.bashrc: export PROMPT_COMMAND="jc_capture" +# Function to call after significant commands + +JC_URL="http://192.168.50.210:8080/api/ingest" +JC_TOKEN="${CAIC_COMPLETIONS_API_KEY}" + +jc_capture() { + local cmd + cmd=$(history 1 | sed 's/^[ ]*[0-9]*[ ]*//') + # Only ingest significant commands + if echo "$cmd" | grep -qE '^(git|pip|systemctl|sudo|vi|vim|curl|wget|apt|python|pytest)'; then + curl -s -X POST "$JC_URL" \ + -H "Authorization: Bearer $JC_TOKEN" \ + -H "Content-Type: application/json" \ + -d "{\"content\": $(echo "$cmd" | python3 -c 'import json,sys; print(json.dumps(sys.stdin.read().strip()))'), \"source\": \"terminal\"}" \ + > /dev/null 2>&1 & + fi +} +``` + +**Write `tests/test_ingest.py`** covering: +- Valid ingest with content — assert chunks_ingested > 0 +- Missing Bearer token — assert 401 +- Wrong Bearer token — assert 403 +- Empty content — assert 422 +- Qdrant and embed calls mocked via monkeypatch + +Run full test suite. All existing tests must continue to pass. + +--- + +## ~~TASK 7 — Roadmap J: Startup Hardware Self-Assessment [DONE]~~ + +**Status: `hardware.py` + `routers/hardware.py` + 4 tests. Committed `7291b8f` (v0.12.0).** + +On jC startup, probe available hardware and write a living config snapshot. This replaces hardcoded assumptions about VRAM and RAM. + +**Create `hardware.py`** in the project root: + +``` +async def assess_hardware() -> dict +``` + +Probes: +- System RAM: `psutil.virtual_memory().total` and `.available` +- CPU count: `psutil.cpu_count()` +- GPU VRAM total and free: call `rocm-smi --showmeminfo vram --json` via subprocess, parse output. If rocm-smi absent or fails, set VRAM values to 0 and log a warning. +- llama-server reachable: GET `LLAMA_SERVER_BASE/v1/models`, timeout 3s. Record True/False and list of available model IDs. +- Qdrant reachable: GET `http://192.168.50.108:6333/collections`, timeout 3s. Record True/False and collection list. +- SearXNG reachable: GET `http://localhost:8888`, timeout 3s. Record True/False. + +Returns a dict with all of the above. Writes result as JSON to `hardware_state.json` in the working directory. + +**Call `assess_hardware()` from the FastAPI `lifespan` context** in `app.py` on startup, after `init_db()`. Log a summary line: `HW: {ram_gb}GB RAM, {vram_mb}MB VRAM, llama={reachable}, qdrant={reachable}, searxng={reachable}`. + +**Expose `GET /api/hardware`** in a new `routers/hardware.py` — returns the current `hardware_state.json` content as JSON. No auth required (read-only, non-sensitive aggregate stats). + +**Wire `hardware.router` into `app.py`.** + +**Write `tests/test_hardware.py`** covering: +- `assess_hardware()` with all services reachable (mock subprocess and httpx calls) — assert all fields present +- `assess_hardware()` with rocm-smi absent — assert VRAM=0, no exception raised +- `assess_hardware()` with llama-server unreachable — assert `llama_reachable=False`, no exception +- `GET /api/hardware` — assert returns JSON with expected keys + +Run full test suite. All existing tests must continue to pass. + +--- + +## ~~TASK 8 — Roadmap K: RAG Corpus Management [DONE]~~ + +Qdrant collection `caic` currently grows without bound. Implement score-based eviction with hysteresis, pinned sources, operational stats, and a flush command. + +### Config — add to `config.py`: + +```python +RAG_MAX_VECTORS = 50000 # absolute ceiling; eviction targets thresholds below it +RAG_EVICTION_HIGH_WATER = 0.80 # fraction of RAG_MAX_VECTORS that triggers eviction +RAG_EVICTION_LOW_WATER = 0.20 # fraction where eviction stops +RAG_EVICTION_BATCH = 1000 # max points to delete per Qdrant scroll/delete cycle +RAG_PINNED_SOURCES = ["upload", "profile"] # never evicted +RAG_GRACE_HOURS = 1 # new vectors ineligible for eviction until this old +RAG_ACCESS_WEIGHT = 1.0 # score factor: retrieval_count * ACCESS_WEIGHT +RAG_AGE_WEIGHT = 0.1 # score factor: ingest_age_hours * AGE_WEIGHT +``` + +Validations on boot: `high_water > low_water`, `batch > 0`, `max_vectors > 0`. + +### Eviction algorithm — add to `rag.py`: + +``` +score = (retrieval_count * ACCESS_WEIGHT) + (age_hours * AGE_WEIGHT) +``` + +Lower score = evicted first. Tiebreak: `last_accessed` ASC (older wins). + +```python +async def get_collection_count() -> int + # GET /collections/caic → return vectors_count + +async def get_collection_stats() -> dict + # Return {vector_count, max_vectors, high_water, low_water, percent_full, pinned_sources} + +async def evict_batch(batch_size: int) -> int + # Scroll Qdrant for vectors NOT in RAG_PINNED_SOURCES, WHERE ingest_age > RAG_GRACE_HOURS, + # ordered by score ASC, last_accessed ASC. + # Delete up to batch_size. Return count deleted. + # If 0 evictable vectors found: log warning, return 0 (break loop). + +async def maybe_evict() -> int + # Acquire eviction_lock (asyncio.Lock). + # count = get_collection_count() + # threshold_high = RAG_MAX_VECTORS * RAG_EVICTION_HIGH_WATER + # threshold_low = RAG_MAX_VECTORS * RAG_EVICTION_LOW_WATER + # total_evicted = 0 + # while count >= threshold_low: + # if total_evicted > 0 and count < threshold_low: break + # deleted = evict_batch(RAG_EVICTION_BATCH) + # if deleted == 0: break # no more unpinned targets + # total_evicted += deleted + # count -= deleted + # if count < threshold_high and total_evicted > 0: break + # # only one pass if batch spans the full gap + # if count < threshold_low: break + # Record total_evicted + timestamp in EVICTION_LOG (list of dicts, kept in memory, max 1000 entries) + # Release lock. Return total_evicted. + +async def get_rag_operational_stats() -> dict + # Returns: vector_count, max_vectors, high_water_pct, low_water_pct, + # percent_full, pinned_sources, grace_hours, + # eviction_counts_last_1m, eviction_counts_last_5m, eviction_counts_last_30m, + # at_risk_count (vectors in bottom 10% by score), + # pinned_count, avg_retrieval_count +``` + +### Edge cases & guards: + +1. **Newborn grace** — vectors < `RAG_GRACE_HOURS` old are excluded from eviction scroll (score=0 otherwise → immediate deletion) +2. **All-pinned freeze** — if scroll returns 0 evictable vectors, log warning and break loop +3. **Race** — `asyncio.Lock()` guards `maybe_evict()`; concurrent callers wait their turn +4. **Zero config** — `RAG_MAX_VECTORS <= 0` → eviction disabled; `RAG_EVICTION_BATCH <= 0` → clamped to 1 +5. **Legacy payloads** — vectors without `retrieval_count` or `last_accessed` get defaults (0, `ingest_date`) + +### Wire eviction: + +Call `maybe_evict()` after each upsert batch completes in: +- `routers/upload.py` — after Qdrant upsert +- `routers/ingest.py` — after Qdrant upsert + +### Admin endpoints — new `routers/rag_admin.py`: + +| Method | Endpoint | Description | +|--------|----------|-------------| +| GET | `/api/rag/stats` | Operational stats (see `get_rag_operational_stats()`) — admin required | +| POST | `/api/rag/flush` | Delete ALL points from the Qdrant `caic` collection. Returns `{deleted_count, collection: "caic", status: "flushed"}`. Admin required. | + +### In-memory eviction log: + +```python +EVICTION_LOG: list[dict] = [] # managed by rag.py, max 1000 entries +# Each entry: {timestamp: iso, count: N, remaining: N} +# Tied to RATE_EVENTS pattern from security.py for rolling window calculations +``` + +### Tests — `tests/test_rag_management.py`: + +- `get_collection_count()` — mock Qdrant GET, assert correct count +- `get_collection_stats()` — assert shape matches config +- `evict_batch()` — mock Qdrant scroll + delete, assert pinned sources excluded, grace period enforced, batch size respected +- `maybe_evict()` — below high water: 0 evicted; at high water: eviction fires; stops at low water; all-pinned scroll returns 0 → breaks +- `GET /api/rag/stats` — assert full shape +- `POST /api/rag/flush` — assert points deleted, admin required, guest 403 +- `POST /api/rag/flush` by guest — assert 403 +- Race lock — concurrent calls to `maybe_evict()` queue up, only one evicts + +Mock all Qdrant calls via monkeypatch. Do not require live services. + +Run full test suite. All existing tests must continue to pass. + +--- + +## ~~TASK 9 — Roadmap N1: RabbitMQ Install and Service on Coordinator (Infrastructure) [DONE]~~ + +This task runs on coordinator (this machine). Install RabbitMQ and verify it is operational. + +Run the following steps: +1. `apt-get update && apt-get install -y rabbitmq-server` +2. `systemctl enable rabbitmq-server && systemctl start rabbitmq-server` +3. `systemctl status rabbitmq-server` — verify active/running +4. Enable the management plugin: `rabbitmq-plugins enable rabbitmq_management` +5. Create a dedicated jC vhost: `rabbitmqctl add_vhost caic` +6. Create a dedicated user: `rabbitmqctl add_user caic CHANGEME_PASSWORD` — generate a random 24-char alphanumeric password and record it +7. Grant permissions: `rabbitmqctl set_permissions -p caic caic ".*" ".*" ".*"` +8. Verify management UI is reachable: `curl -s -u guest:guest http://localhost:15672/api/overview | python3 -m json.tool` +9. Delete default guest user: `rabbitmqctl delete_user guest` + +Declare the two topic exchanges needed by jC: +- Exchange name: `jc.admin`, type: `topic`, durable: true +- Exchange name: `jc.system`, type: `topic`, durable: true + +Use `rabbitmqadmin` or `curl` against the management API to declare exchanges. Verify both exchanges appear in: `curl -s -u caic:{password} http://localhost:15672/api/exchanges/caic` + +Write the generated RabbitMQ password to `/home/gramps/.caic_amqp_secret` with mode 600. This will be read by cAIc as an env var source in subsequent tasks. + +No pytest tests required for this infrastructure task. + +--- + +## ~~TASK 10 — Roadmap N2: AMQP Connection Layer in jC [DONE]~~ + +This task adds the core AMQP connection manager to jC. It must connect to RabbitMQ on coordinator (localhost from jC's perspective since jC runs on coordinator), handle reconnection, and provide a shared channel for all AMQP operations. + +**Add to `requirements.txt`:** `aio-pika>=9.0.0` + +**Add to `config.py`:** +- `AMQP_URL` — read from env `CAIC_AMQP_URL`, default `amqp://caic:password@localhost:5672/caic`. The actual password comes from `/home/gramps/.caic_amqp_secret` — read it at startup if the env var is not set. +- `AMQP_RECONNECT_DELAY` — seconds between reconnect attempts, default 5 +- `AMQP_EXCHANGE_ADMIN` — `jc.admin` +- `AMQP_EXCHANGE_SYSTEM` — `jc.system` + +**Create `amqp.py`** in the project root: + +```python +# Manages a single persistent aio-pika connection and channel. +# Provides: +# connect() -> None # establish connection, declare exchanges +# disconnect() -> None # graceful close +# get_channel() # returns current channel, reconnects if needed +# publish(exchange, routing_key, payload: dict) -> None +# # publishes JSON-serialized payload as persistent message +``` + +Connection must: +- Reconnect automatically on disconnect with `AMQP_RECONNECT_DELAY` backoff +- Log connection events at INFO level +- Not raise on publish if disconnected — log error and return (fire-and-forget, jC must not crash if RabbitMQ is down) + +**Start AMQP connection in `app.py` lifespan** after `assess_hardware()`. Disconnect in lifespan cleanup. + +**Write `tests/test_amqp.py`** covering: +- `publish()` with mocked aio-pika connection — assert message published with correct exchange and routing key +- `publish()` when disconnected — assert no exception raised, error logged +- `get_channel()` when connection is None — assert reconnect attempted + +Mock all aio-pika calls via monkeypatch. Do not require a live RabbitMQ instance in tests. + +Run full test suite. All existing tests must continue to pass. + +--- + +## ~~TASK 11 — Roadmap N3: Cluster Protocol & Registration Handler (Coordinator Side) [DONE]~~ + +**Status: Implemented and pushed (899988c).** `amqp.py` subscribe/rebind, `cluster.py` with CLUSTER_NODES/CLUSTER_EVENTS/CLUSTER_COORDINATOR and 6 handlers, `routers/cluster.py` (`GET /api/cluster`), 13 tests. No passive heartbeats — ping/pong on-demand before work routing. 148 tests pass. + +jC on the coordinator must listen for nine message types across `jc.admin` and `jc.system`, maintain the cluster registry, and expose an application-level event log. + +### 11.1 AMQP Protocol — Message Catalog + +All payloads are JSON, published as persistent messages. + +| Direction | Exchange | Routing Key | Message Type | Description | +|-----------|----------|-------------|-------------|-------------| +| Worker → Coordinator | `jc.admin` | `node.{name}.register` | register | Worker requests admission | +| Worker → Coordinator | `jc.admin` | `node.{name}.deregister` | deregister | Worker signals graceful departure | +| Coordinator → Worker | `jc.admin` | `node.{name}.admitted` | admitted | Coordinator grants admission | +| Coordinator → Worker | `jc.admin` | `node.{name}.rejected` | rejected | Coordinator denies admission (with reason) | +| Coordinator → Worker | `jc.admin` | `node.{name}.ping` | ping | Coordinator checks if worker is alive (sent before routing work) | +| Worker → Coordinator | `jc.admin` | `node.{name}.pong` | pong | Worker confirms aliveness | +| Worker → Coordinator | `jc.system` | `node.{name}.event` | event | Application-level syslog event | +| Any → All | `jc.system` | `cluster.coordinator.query` | coord_query | Anyone asks "who is coordinator?" | +| Coordinator → All | `jc.system` | `cluster.coordinator.response` | coord_response | Coordinator announces itself | + +Worker presence is assumed from registration onward. No periodic heartbeats — a worker can sit idle for days without chatter. When the coordinator needs to route work to a worker, it pings first; if the worker doesn't pong within timeout, the coordinator deregisters it and moves to the next node. + +### 11.2 Payload Schemas + +**register** (worker → coordinator): +```json +{ + "node_name": "worker01", + "node_type": "worker", + "ip": "192.168.50.210", + "capabilities": { + "gpu": true, "gpu_type": "amd", "vram_mb": 8192, + "cpu_cores": 8, "ram_gb": 16 + }, + "active_model": { + "name": "llama3.1", "version": "latest", "quant": "Q4_K_M", + "path": "/var/lib/caic/models/llama3.1-latest-Q4_K_M.gguf", + "port": 8081 + }, + "inventory": [ + {"name": "llama3.1", "version": "latest", "quant": "Q4_K_M", + "path": "/var/lib/caic/models/llama3.1-latest-Q4_K_M.gguf", "port": 8081} + ], + "status": "active" +} +``` + +**deregister** (worker → coordinator): +```json +{ + "node_name": "worker01", + "reason": "shutdown", + "timestamp": "2026-07-06T12:00:00Z" +} +``` + +**ping** (coordinator → worker): +```json +{ + "from": "coordinator", + "node_name": "worker01", + "type": "ping", + "correlation_id": "550e8400-e29b-41d4-a716-446655440000", + "timestamp": "2026-07-06T12:00:00Z" +} +``` +Worker must respond within 5 seconds or the coordinator considers it absent. + +**pong** (worker → coordinator): +```json +{ + "node_name": "worker01", + "type": "pong", + "correlation_id": "550e8400-e29b-41d4-a716-446655440000", + "status": "active", + "active_model": {"name": "llama3.1", "port": 8081}, + "load": {"cpu_pct": 45, "ram_pct": 62, "vram_pct": 38}, + "timestamp": "2026-07-06T12:00:00Z" +} +``` +Correlation ID matches the ping so the coordinator can pair request and response. + +**coord_query** (any → `cluster.coordinator.query`): +```json +{"type": "coord_query", "timestamp": "2026-07-06T12:00:00Z"} +``` +Coordinator responds on `cluster.coordinator.response`: +```json +{ + "coordinator_node": "coordinator", + "cluster_nodes": ["worker01"], + "timestamp": "2026-07-06T12:00:00Z" +} +``` + +**event** (worker → coordinator): +```json +{ + "node_name": "worker01", + "severity": "info", + "message": "llama-server started with model llama3.1:latest", + "details": {"model": "llama3.1:latest", "port": 8081, "pid": 1234}, + "timestamp": "2026-07-06T12:00:00Z" +} +``` +Severity levels: `info`, `warn`, `error`, `critical`. The coordinator assigns `category: "application"` based on the exchange (jc.system). No `event_type` field — the category is determined by the channel, not the payload. + +### 11.3 Design — Status Transitions Drive the Event Log + +All admin-level events are *derived* from `register()` and `deregister()` as side effects. There are no separate message types for coordinator election, node staleness, quarantine, or release — those are status transitions that `register()`/`deregister()` emit into `CLUSTER_EVENTS` locally. + +**Node status lifecycle:** + +``` +UNKNOWN ──register()──▶ active ──deregister()──▶ (removed) + │ + ping timeout│(coordinator publishes + │ deregister on its behalf) + ▼ + (removed) +``` + +**Coordinator status lifecycle:** + +``` +NONE ──register(node_type=coordinator)──▶ CLUSTER_COORDINATOR set + │ + deregister()│or timeout + ▼ + CLUSTER_COORDINATOR cleared +``` + +**Event categories — two buckets, no granular types:** + +| Category | When | severity | +|----------|------|----------| +| `cluster` | Node lifecycle, coordinator changes, model swaps, node offline — everything on `jc.admin` | `info` / `warn` / `error` | +| `application` | Worker syslog events (incoming on `jc.system` `node.*.event`) | `info` / `warn` / `error` / `critical` | + +Every `_push_event()` call uses one of these two categories. The `message` field carries the human-readable detail — no need for event type strings. The reporting tool filters by category + severity. + +**Channel split — security rationale:** + +The two exchanges are not an organizational convenience. They enforce a **data isolation boundary**: + +| Exchange | Contains | Exposed to | +|----------|----------|------------| +| `jc.admin` | Node lifecycle, heartbeats, model swaps, coordinator changes | Operations / machine-room staff | +| `jc.system` | Application events — inference queries, RAG context, user-facing data | Application-layer audit only | + +`jc.system` events can leak information about what users are doing and asking. The split ensures a sysadmin monitoring cluster health never accidentally consumes user-data-bearing events. The channels can be locked down independently — different AMQP credentials, separate queue permissions, different in-transit encryption policies if needed later. + +### 11.4 Implementation + +**Add to `amqp.py`:** + +```python +_SUBSCRIPTIONS: list[tuple[str, str, Callable]] # (exchange, routing_key, callback) + +async def subscribe(exchange, routing_key, callback) -> None + # Append to _SUBSCRIPTIONS list + # Declare a unique queue per subscription (name: f"jc.{exchange}.{sanitized_routing_key}") + # Bind queue to exchange/routing_key, consume with callback +``` + +Each subscription gets its own queue so multiple subscribers on different routing keys all receive messages. On reconnect: drain old consumers, iterate `_SUBSCRIPTIONS`, re-declare and re-bind each one. The `connect()` function must call `_rebind_subscriptions()` after exchanges are declared. + +**Create `cluster.py`** in the project root: + +```python +# In-memory cluster registry + event log +# Survives only while jC is running (not persisted) + +CLUSTER_NODES: dict[str, NodeRecord] +CLUSTER_EVENTS: deque[EventRecord] # bounded at 1000 entries +CLUSTER_COORDINATOR: str | None # node_name of active coordinator + +# NodeRecord fields: +# node_name, node_type, ip, status, active_model, inventory, +# capabilities: {gpu, gpu_type, vram_mb, cpu_cores, ram_gb} +# registered_at, last_seen + +# EventRecord: +# category: str ("cluster" | "application") +# severity: str ("info" | "warn" | "error" | "critical") +# node_name: str +# message: str +# details: dict | None +# timestamp: str + +def _push_event(category, severity, node_name, message, details=None) -> None + # Append EventRecord to CLUSTER_EVENTS, pop left if > 1000 + +async def handle_registration(message) -> None + # Parse payload, validate required fields (node_name, node_type, ip, active_model, inventory) + # Reject if node_name duplicate and CLUSTER_NODES[node_name].status == "active" + # If CLUSTER_COORDINATOR is None AND node_type == "coordinator": + # set CLUSTER_COORDINATOR = node_name + # _push_event("cluster", "info", node_name, "elected coordinator") + # publish cluster.coordinator.response on jc.system {coordinator_node, cluster_nodes, timestamp} + # Add node to CLUSTER_NODES with status="active" + # _push_event("cluster", "info", node_name, f"admitted as {node_type}") + # publish admitted on jc.admin node.{name}.admitted {node_name, timestamp, amqp_url} + +async def handle_deregistration(message) -> None + # Parse payload (node_name, reason, timestamp) + # If node_name == CLUSTER_COORDINATOR: + # clear CLUSTER_COORDINATOR + # _push_event("cluster", "warn", node_name, f"coordinator lost — {reason}") + # _push_event("cluster", "info", node_name, f"departed — {reason}") + # Remove node from CLUSTER_NODES, log it + +async def handle_pong(message) -> None + # Parse: node_name, correlation_id, status, active_model, load, timestamp + # Match correlation_id to outstanding ping + # If node in CLUSTER_NODES: update last_seen, status, active_model + # Signal the waiting caller that the node is alive + # If node unknown: log warning, do NOT auto-admit + +async def handle_event(message) -> None + # Parse: node_name, severity, message, details, timestamp + # Assigns category="application" (incoming on jc.system) + # Append EventRecord to CLUSTER_EVENTS (pop left if > 1000) + +async def handle_coordinator_query(message) -> None + # Respond on jc.system cluster.coordinator.response + # Payload: {coordinator_node, cluster_nodes: list(CLUSTER_NODES.keys()), timestamp} + +def get_cluster_state() -> dict + # Return: {nodes: CLUSTER_NODES, coordinator: CLUSTER_COORDINATOR, + # events: last 50 CLUSTER_EVENTS} +``` + +**Subscribe in `app.py` lifespan** after AMQP connects: + +| Exchange | Routing Key | Handler | +|----------|-------------|---------| +| `jc.admin` | `node.*.register` | `handle_registration` | +| `jc.admin` | `node.*.deregister` | `handle_deregistration` | +| `jc.admin` | `node.*.pong` | `handle_pong` | +| `jc.system` | `node.*.event` | `handle_event` | +| `jc.system` | `cluster.coordinator.query` | `handle_coordinator_query` | + +### 11.5 API — `GET /api/cluster` + +New router `routers/cluster.py`: +- `GET /api/cluster` — returns full cluster state: `{nodes, coordinator, events}` (last 50 events). No auth required. + +Wire `cluster.router` into `app.py`. + +### 11.6 Tests — `tests/test_cluster.py` + +Mock all aio-pika calls. Do not require live RabbitMQ. + +| # | Test | What it asserts | +|---|------|-----------------| +| 1 | Valid worker registration | Node admitted, CLUSTER_NODES updated, `cluster` event logged, `admitted` message published | +| 2 | First coordinator auto-promotion | CLUSTER_COORDINATOR set, `cluster` event with "elected" message, `coord_response` published | +| 3 | Duplicate node name rejected | `rejected` message with reason=`duplicate_node_name`, `cluster` event logged | +| 4 | Malformed payload rejected | `rejected` message with reason=`malformed_payload` | +| 5 | Graceful deregistration | Node removed, `cluster` event logged. If coordinator: CLUSTER_COORDINATOR cleared | +| 6 | Pong from known node | last_seen updated, load/status refreshed | +| 7 | Pong from unknown node | Warning logged, node NOT added | +| 8 | Event stored in log | Event appended to CLUSTER_EVENTS; at 1001 entries the oldest is popped | +| 9 | Coordinator query produces response | Response published with coordinator name and node list | +| 10 | GET /api/cluster shape | Response contains `nodes`, `coordinator`, `events` keys | + +Run full test suite. All existing tests must continue to pass. + +--- + +## ~~TASK 12 — Roadmap N4: Worker Node Registration Publisher (Worker Side) [DONE]~~ + +This task creates the worker node AMQP client that runs on worker (192.168.50.210). It is a standalone Python script — not part of the jC FastAPI app — that runs as a systemd service on worker. + +**Create `node_agent/agent.py`** in the repo (new directory). + +### 12.1 Config & Inventory Discovery + +On start, reads `/etc/caic-node-agent.conf` (INI format): +- `node_name` — hostname, default from `socket.gethostname()` +- `node_ip` — LAN IP, default from socket +- `node_type` — `"worker"` (fixed) +- `capabilities` — comma-separated list, e.g. `llm,rag` +- `amqp_url` — RabbitMQ URL on coordinator, e.g. `amqp://caic:password@192.168.50.108:5672/caic` +- `llama_port` — port llama-server/llama-rpc is listening on, default 8081 +- `models_dir` — path to GGUF model files, default `/var/lib/caic/models` +- `active_model` — filename of currently active model (without path) + +Discovers inventory by globbing `models_dir` for `*.gguf` files and parsing name/version/quant from filename using regex pattern: `{name}-{version}-{quant}.gguf` where quant matches `Q[0-9]+_K_[A-Z]+` or similar standard suffixes. + +### 12.2 Registration + +Publishes registration to `jc.admin`, routing key `node.{node_name}.register`: +```json +{ + "node_name": "worker01", + "node_type": "worker", + "ip": "192.168.50.210", + "capabilities": ["llm"], + "active_model": {"name": "...", "version": "...", "quant": "...", "path": "...", "port": 8081} +} +``` + +### 12.3 Admission Response + +Listens on `node.{node_name}.admitted` and `node.{node_name}.rejected` (both `jc.admin`). Logs result. If rejected, exits with error. + +### 12.4 Ping Listener + +After admission: listens on `jc.admin`, routing key `node.{node_name}.ping`. On receipt, responds immediately (within 1 second) with a pong on `jc.admin`, routing key `node.{node_name}.pong`: + +```json +{ + "node_name": "worker01", + "type": "pong", + "correlation_id": "", + "status": "active", + "active_model": {"name": "...", "version": "...", "quant": "...", "path": "...", "port": 8081}, + "load": {"cpu_pct": 45, "ram_pct": 62, "vram_pct": 38}, + "timestamp": "" +} +``` + +No periodic heartbeats. Worker sits idle between pings — coordinator only pings when it needs to route work. + +### 12.5 Model Swap Command Handler + +Listens on `jc.admin`, routing key `node.{node_name}.cmd.swap_model`: +- Payload: `{model_filename: str}` +- Stops current llama-server: `systemctl stop llama-server` +- Updates `/etc/caic-node-agent.conf` active_model field +- Starts llama-server: `systemctl start llama-server` (assumes service reads active_model from conf or ExecStart is updated) +- Waits for llama-server to be healthy: poll `http://localhost:{llama_port}/v1/models` every 2s, timeout 120s +- Publishes to `jc.system`, routing key `node.{node_name}.model_ready`: + ```json + {"node_name": "...", "active_model": "...", "port": ..., "timestamp": "..."} + ``` +- If startup fails within timeout: publishes `node.{node_name}.model_failed` with error detail + +### 12.6 Files & Tests + +**Create `node_agent/requirements.txt`:** `aio-pika>=9.0.0` + +**Document `/etc/caic-node-agent.conf` format** in a comment block at the top of `agent.py`. + +**Write `tests/test_node_agent.py`** covering: +- Registration payload construction from config + model discovery — assert correct JSON shape +- Model swap command handler: success path — assert systemctl calls made, model_ready published +- Model swap command handler: timeout path — assert model_failed published +- Ping handler: on ping, publishes pong with correct correlation_id +- Agent starts idle after admission, no heartbeat timer + +Mock all aio-pika, subprocess, and httpx calls. + +**Do not create a systemd service file in this task** — that is a manual deployment step. Document the required service configuration in a comment at the bottom of `agent.py`. + +Run full test suite. All existing tests must continue to pass. + +--- + +## ~~TASK 13 — Roadmap N5: Query Routing via AMQP + Phi-4-mini Triage [DONE]~~ + +This task wires the cluster into jC's chat flow. When a query arrives at `/api/chat`, instead of always routing to the hardcoded `LLAMA_SERVER_BASE`, jC now routes to the best available cluster node based on query context. + +**Prerequisites:** Tasks 9–12 complete. At least one worker node admitted to cluster. + +**Install Phi-4-mini on coordinator (infrastructure step):** +- Download `Phi-4-mini-Instruct-Q4_K_M.gguf` from HuggingFace using `hf download microsoft/Phi-4-mini-instruct --include "*.Q4_K_M.gguf" --local-dir /var/lib/caic/models` +- Create `/etc/systemd/system/llama-server-triage.service` — same pattern as existing llama-server service but: port 8083, model path points to Phi-4-mini GGUF, no `--rpc` flag (runs entirely on coordinator CPU/iGPU), description `Llama.cpp Server (Phi-4-mini — triage/routing)` +- `systemctl daemon-reload && systemctl enable llama-server-triage && systemctl start llama-server-triage` +- Verify: `curl -s http://localhost:8083/v1/models` + +**Add to `config.py`:** +- `TRIAGE_BASE` — `http://127.0.0.1:8083/v1` (Phi-4-mini) +- `TRIAGE_TIMEOUT` — 10 seconds +- `FALLBACK_TO_DEFAULT` — True (if triage fails or no nodes available, fall back to `LLAMA_SERVER_BASE`) + +**Create `triage.py`** in the project root: + +```python +async def classify_query(query: str) -> str + # Sends query to Phi-4-mini at TRIAGE_BASE with a classification system prompt. + # System prompt instructs model to respond with ONLY one of: + # "general", "code", "search", "rag" + # Returns the classification string. + # Timeout: TRIAGE_TIMEOUT seconds. + # On any error: returns "general" (fail-safe). + +async def select_node(classification: str) -> dict | None + # Consults CLUSTER_NODES from cluster.py + # For "code": prefer nodes where active_model name contains "coder" or "qwen" + # For "general": prefer nodes where active_model name contains "mistral" or "llama" + # For "search" or "rag": return None (handled locally by jC) + # If no matching node found: return None (triggers FALLBACK_TO_DEFAULT) + # Returns NodeRecord dict for selected node, or None + +async def get_inference_url(query: str) -> str + # Combines classify_query + select_node + # Returns full base URL: f"http://{node.ip}:{node.active_model.port}/v1" + # Falls back to LLAMA_SERVER_BASE if classification=search/rag, no nodes, or triage error +``` + +**Update `routers/chat.py`:** +- Replace the hardcoded `LLAMA_SERVER_BASE` reference with a call to `get_inference_url(user_message)` +- The rest of the chat flow (RAG, memory, streaming) is unchanged — only the inference target URL changes + +**Write `tests/test_triage.py`** covering: +- `classify_query()` returns valid classification — mock Phi-4-mini response +- `classify_query()` on timeout — assert returns "general", no exception +- `select_node("code")` with coder node in cluster — assert correct node returned +- `select_node("general")` with no matching node — assert None returned +- `get_inference_url()` with code query and coder node available — assert returns node URL +- `get_inference_url()` with no nodes in cluster — assert returns LLAMA_SERVER_BASE fallback + +**Update `tests/test_chat_streaming_and_memory_paths.py`:** +- Mock `triage.get_inference_url` to return a fixed URL in all existing tests so they continue to pass without a live cluster + +Run full test suite. All existing tests must continue to pass. + +--- + +## ~~TASK 14 — Roadmap N6: Model Swap Command Flow [DONE]~~ + +**Status: Implemented and pushed (`9d1fd44`).** `request_model_swap()`, `handle_model_ready()`, `handle_model_failed()` in `cluster.py`, async `select_node()` with swap triggering in `triage.py`, `tests/test_model_swap.py` (9 tests). 177 tests pass. + +This task implements the coordinator-side logic for requesting a model swap on a worker node when the ideal model is not currently active. + +**Add to `cluster.py`:** + +```python +async def request_model_swap(node_name: str, model_filename: str) -> bool + # Publishes to jc.admin exchange, routing key node.{node_name}.cmd.swap_model + # Payload: {model_filename, requested_at: iso_timestamp} + # Sets node status to "swapping" in CLUSTER_NODES + # Returns True if message published successfully + +async def handle_model_ready(message) -> None + # Handles node.{node_name}.model_ready from jc.system + # Updates CLUSTER_NODES[node_name].active_model to the new model + # Sets node status back to "active" + # Logs swap completion with timing + +async def handle_model_failed(message) -> None + # Handles node.{node_name}.model_failed from jc.system + # Sets node status to "error" in CLUSTER_NODES + # Logs failure with detail from message payload +``` + +**Subscribe in `app.py` lifespan:** +- `jc.system` exchange, routing key `node.*.model_ready` → `handle_model_ready` +- `jc.system` exchange, routing key `node.*.model_failed` → `handle_model_failed` + +**Update `triage.py` `select_node()`:** +- If the best-matching node exists but its active_model does not match the ideal model for the classification, AND the node status is "active" (not already swapping): + - Call `request_model_swap(node_name, ideal_model_filename)` + - Return None (triggers fallback) — the swap happens async, next query will find the right model active +- If node status is "swapping": return None (fallback, swap in progress) + +**Update `GET /api/cluster`** to include node status in response. + +**Write `tests/test_model_swap.py`** covering: +- `request_model_swap()` — assert swap command published, node status set to "swapping" +- `handle_model_ready()` — assert active_model updated, status set to "active" +- `handle_model_failed()` — assert status set to "error" +- `select_node()` with mismatched active model — assert swap requested, None returned +- `select_node()` with node status "swapping" — assert None returned without publishing another swap + +Run full test suite. All existing tests must continue to pass. + +--- + +## ~~TASK 15 — Roadmap N7: Cluster Status UI [DONE]~~ + +Surface cluster awareness in the jC frontend (`templates/index.html`). + +**Add a cluster status panel** to the UI. Requirements: +- Small status bar or collapsible panel, visible but unobtrusive +- Polls `GET /api/cluster` every 15 seconds +- For each admitted node: show node name, active model name, and a colored status dot: + - Green: active + - Yellow: swapping + - Red: error or offline (not seen in last 60 seconds based on last_seen timestamp) +- If no nodes in cluster (empty): show "No worker nodes connected" +- Panel must not interfere with chat input or conversation list + +**Update `GET /api/cluster` response** to include `last_seen` per node and a `status` field (`active`, `swapping`, `error`). + +**Update heartbeat handling in `cluster.py`:** add a handler for `node.*.heartbeat` on `jc.system` that updates `last_seen` timestamp for the node. + +**Subscribe in `app.py` lifespan:** +- `jc.system` exchange, routing key `node.*.heartbeat` → `handle_heartbeat` + +**Add `handle_heartbeat()` to `cluster.py`:** +- Updates `CLUSTER_NODES[node_name].last_seen` to current timestamp +- If node was previously marked offline (not in CLUSTER_NODES), log re-registration warning but do not auto-admit — full registration required + +**Write `tests/test_cluster_heartbeat.py`** covering: +- `handle_heartbeat()` for known node — assert last_seen updated +- `handle_heartbeat()` for unknown node — assert no crash, warning logged, node not added + +Run full test suite. All 26+ existing tests must continue to pass. + +~~Commit all changes introduced across Tasks 9–15 with message: `feat: Roadmap N — AMQP cluster nervous system complete`~~ + +--- + +## Backlog (Post-Roadmap N) ⏳ + +### ~~B1 — Context loss in follow-up questions [DONE]~~ + +**Symptom:** After asking "in {context}, explain {b}", a follow-up "what is {b}'s {x}?" gets a non-sequitur response that ignores the original context. + +**Diagnosis:** `build_system_prompt()` is called fresh per-request with new RAG/memory results keyed to the current message text. These can change between turns and may dilute or override the conversation history. The original system prompt used for turn 1 (including its RAG context) is not stored in the DB — only user/assistant messages are. The inference server receives a different system prompt each turn. + +**Possible fixes:** +- Store the assembled system prompt with each assistant message in the DB +- When replaying history, re-send the original system prompts from DB rather than rebuilding +- Or: cap RAG/memory injection to only fire on the first message of a conversation, then rely solely on conversation history for follow-ups +- Check that llama-server isn't truncating history due to context window overflow (Mistral-Nemo 12B = 128K context, unlikely) + +### ~~B2 — Bang-prefixed search routing [DONE]~~ + +**Spec:** If a query begins with `!`, route to SearXNG search instead of local inference. + +**Where:** In `routers/chat.py` `chat()` handler, after `user_message` is extracted. Strip the `!`, set a flag to always trigger auto-search regardless of perplexity/refusal. + +**Change:** Add a `force_search` flag when `user_message.startswith("!")`, strip the prefix from the message saved to DB, and route directly to the search+summarize path. + +### ~~B3 — Docker distribution (v1.0 gate) [DONE]~~ + +**Goal:** Ship cAIc as a `docker compose` stack so a single command stands up everything. + +**Services to containerize:** +- cAIc (FastAPI app + SQLite) +- SearXNG +- Qdrant +- RabbitMQ +- llama-server (with optional RPC sidecar for GPU offload) +- Ollama (embeddings) + +**Also needed:** +- `Dockerfile` for the cAIc app itself +- `docker-compose.yml` with all services, volumes, networks, env vars +- Setup wizard script (run on first boot) that: + - Probes CPU vs GPU (reuses `hardware.py`) + - Queries user for admin PIN, node name, IP + - Generates `.env` file with correct `LLAMA_SERVER_BASE`, `EMBED_URL`, etc. + - Auto-calculates `RAG_MAX_VECTORS` from available RAM: `max(1000, int(available_ram_gb * 100_000))` + - Optionally detects and configures RPC GPU offload +- Manual install docs remain alongside for bare-metal deployment + +**This task is only actionable after Tasks 8–15 (RAG eviction + AMQP cluster) are complete.** + +--- + +### ~~B4 — RAG Corpus Management UI (Display, Edit, CRUD) [DONE]~~ + +**Goal:** Provide a management interface in the UI to browse, search, edit, and delete individual entries in the Qdrant-backed RAG corpus. + +**Backend — add to `routers/rag_admin.py`:** + +| Method | Endpoint | Description | Auth | +|--------|----------|-------------|------| +| GET | `/api/rag/points` | Return paginated list of RAG points with payload (text, source, date). Supports `?offset=0&limit=50&search=` query params | Admin | +| GET | `/api/rag/point/{point_id}` | Return a single point with full payload | Admin | +| DELETE | `/api/rag/point/{point_id}` | Delete a single point from Qdrant | Admin | +| PATCH | `/api/rag/point/{point_id}` | Update a point's text payload (re-embed the new text) | Admin | + +Helper functions for Qdrant scroll/delete/update go in `rag.py` or `eviction.py`. + +**Frontend — add to `templates/index.html`:** + +A "RAG" button in the admin UI (drawer or settings modal) that opens a management panel: +- **Stats bar**: vector count, max vectors, percent full, pinned sources +- **Search bar**: text input to search the RAG corpus by semantic similarity +- **Results table**: paginated list showing each vector's text snippet, source label, ingest date, retrieval count + - Click to expand full text + - Delete button per row (with confirmation) + - Edit button per row (inline text edit → re-embed on save) +- **Bulk actions**: flush all (existing `/api/rag/flush`) with confirmation + +**Tests:** + +- `tests/test_rag_admin.py` — cover new endpoints: list, get, delete, update, admin-enforcement +- Mock all Qdrant calls via monkeypatch + +Run full test suite. All existing tests must continue to pass.** diff --git a/docs/wiki/Developer-Architecture.md b/docs/wiki/Developer-Architecture.md index 794cd00..4be6f85 100644 --- a/docs/wiki/Developer-Architecture.md +++ b/docs/wiki/Developer-Architecture.md @@ -22,9 +22,9 @@ Refactored from single-file (`app.py`) into modules under project root: | `rag.py` | Qdrant vector search, system prompt assembly, chunk_text() helper, collection stats | | `eviction.py` | Score-based RAG eviction engine (extracted from rag.py) | | `gpu.py` | AMD GPU stats via rocm-smi | -| `hardware.py` | Hardware self-assessment — CPU, RAM, VRAM, service health probes | +| `hardware.py` | Hardware self-assessment — CPU, RAM, VRAM, service health probes (llama-server, Qdrant, SearXNG, ComfyUI) | | `amqp.py` | aio-pika connection manager for RabbitMQ (connect, disconnect, publish, subscribe, auto-reconnect) | -| `cluster.py` | Cluster node registry, event log, coordinator election, ping/pong, model swap handlers | +| `cluster.py` | Cluster node registry, event log, coordinator election, ping/pong, model swap handlers, image generation request/response | | `triage.py` | Phi-4-mini query classification + `select_node()` for cluster routing | | `routers/` | One module per endpoint group | @@ -36,6 +36,7 @@ Refactored from single-file (`app.py`) into modules under project root: | SearXNG | No | 8888 | Privacy-respecting web search | | Qdrant (coordinator) | No | 6333 | Vector database for RAG | | Ollama (worker) | No | 11434 | Embeddings for RAG chunk vectors | +| ComfyUI (worker) | No | 8188 | Image generation (Stable Diffusion / Flux) | | RabbitMQ (coordinator) | No | 5672 | AMQP broker for cluster messaging | | rocm-smi | No | — | AMD GPU stats (host-level) | @@ -50,6 +51,8 @@ Key base URLs are configured via environment variables with sensible defaults: | `SEARXNG_BASE` | `http://localhost:8888` | SearXNG | | `QDRANT_URL` | `http://192.168.50.108:6333` | Qdrant on coordinator | | `CAIC_AMQP_URL` | `amqp://caic:password@localhost:5672/caic` | RabbitMQ | +| `CAIC_COMFYUI_BASE` | `http://192.168.50.115:8188` | ComfyUI on worker | +| `CAIC_COMFYUI_TIMEOUT` | `120` | ComfyUI generation timeout (seconds) | ## 2. Request/Response Architecture @@ -87,6 +90,17 @@ Key base URLs are configured via environment variables with sensible defaults: 4. Three modes: `context` (SQLite with 1hr expiry), `ingest` (RAG/Qdrant), `both` 5. Trigger `maybe_evict()` if ingest mode +### 2.5 Image Generation Pipeline (`POST /api/image/generate`) + +1. Admin required, JSON body with prompt and optional params (width, height, steps, seed, model) +2. Find active node with `image_gen` capability via `_find_image_node()` +3. Publish `cmd.image_generate` via AMQP to selected worker node +4. Worker node agent builds ComfyUI workflow (CheckpointLoader → KSampler → VAEDecode → SaveImage) +5. Worker polls ComfyUI `/history/{prompt_id}` until image is ready +6. Worker fetches PNG from ComfyUI `/view` endpoint, base64-encodes, publishes `image_generated` on `jc.system` +7. Coordinator receives response, decodes base64, returns `image/png` to client +8. `GET /api/image/status` returns available image gen nodes and their status + ## 3. Data Model (SQLite) Key tables: @@ -242,8 +256,8 @@ Every RabbitMQ server belongs to a cluster. Currently only the coordinator runs | Exchange | Type | Purpose | |----------|------|---------| -| `jc.admin` | topic | Lifecycle commands: register, deregister, ping, pong, admitted, rejected; model commands: cmd.swap_model | -| `jc.system` | topic | Events: model_ready, model_failed, node.*.heartbeat, event; coordinator queries: coord_query, coord_response | +| `jc.admin` | topic | Lifecycle commands: register, deregister, ping, pong, admitted, rejected; model commands: cmd.swap_model; image commands: cmd.image_generate | +| `jc.system` | topic | Events: model_ready, model_failed, image_generated, image_failed, node.*.heartbeat, event; coordinator queries: coord_query, coord_response | All exchanges, queues, and bindings are declared by `amqp.py` at startup. Worker runs `node_agent/agent.py` which connects as an AMQP client, registers, responds to ping, and handles model swap commands. @@ -265,7 +279,7 @@ All streaming endpoints yield `data: {json}\n\n`: - No live external services required - Test factories reset `SESSIONS`, `PIN_ATTEMPTS`, `RATE_EVENTS` globals per test -### 8.2 Test Coverage Areas (200 tests) +### 8.2 Test Coverage Areas (228 tests) | Test file | Coverage | |-----------|----------| @@ -277,6 +291,8 @@ All streaming endpoints yield `data: {json}\n\n`: | test_conversations.py | Full CRUD, guest admin, attachment_count | | test_error_envelopes.py | Global exception handler + stream errors | | test_gpu.py | GPU stats — rocm-smi (Linux), system_profiler (Darwin/Apple Silicon) | +| test_hardware.py | Hardware assessment, service reachability | +| test_image.py | Image generation — cluster handlers, router proxy, node agent ComfyUI integration, hardware probe, capability detection | | test_ingest.py | Bearer auth, chunk/embed/upsert, validation | | test_ip_allowlist.py | IP allowlist helper + middleware | | test_memories.py | Edit, search, stats | @@ -311,5 +327,6 @@ On startup, `assess_hardware()` probes: - llama-server reachability + model list - Qdrant reachability + collection list - SearXNG reachability +- ComfyUI reachability + checkpoint model list Writes `hardware_state.json` to working directory. diff --git a/docs/wiki/current-wip.md b/docs/wiki/current-wip.md index f8126a9..be65172 100644 --- a/docs/wiki/current-wip.md +++ b/docs/wiki/current-wip.md @@ -1,9 +1,13 @@ # cAIc Current WiP Backlog -Last updated: 2026-07-14 +Last updated: 2026-07-27 Owner: Gramps Scope: Active roadmap items and backlog. +## In Progress + +- **Image Generation Service** — Backend wired: cluster handlers, `POST /api/image/generate` proxy, node agent ComfyUI integration, hardware probe, 27 tests. ComfyUI install pending on corsair (RTX 5070 Ti). + ## Completed - **B8 (v0.19.3)** — Private Chat mode. Backend skip-DB/skip-RAG/skip-search flag, frontend PRIVATE badge, info popup. diff --git a/hardware.py b/hardware.py index 5db5f83..30c4bb0 100644 --- a/hardware.py +++ b/hardware.py @@ -12,7 +12,7 @@ from pathlib import Path import httpx import psutil -from config import LLAMA_SERVER_BASE, SEARXNG_BASE, QDRANT_URL, HW_STATE_PATH +from config import LLAMA_SERVER_BASE, SEARXNG_BASE, QDRANT_URL, HW_STATE_PATH, COMFYUI_BASE log = logging.getLogger("caic") @@ -113,6 +113,20 @@ async def assess_hardware() -> dict: except Exception: log.warning("SearXNG not reachable") + comfyui_reachable = False + comfyui_models = [] + try: + async with httpx.AsyncClient(timeout=5) as client: + resp = await client.get(f"{COMFYUI_BASE}/object_info/CheckpointLoaderSimple") + if resp.status_code == 200: + comfyui_reachable = True + data = resp.json() + ckpt_info = data.get("CheckpointLoaderSimple", {}).get("input", {}).get("required", {}) + ckpt_list = ckpt_info.get("ckpt_name", [[]])[0] + comfyui_models = ckpt_list if isinstance(ckpt_list, list) else [] + except Exception: + log.warning("ComfyUI not reachable") + state = { "ram_total_gb": ram_total_gb, "ram_available_gb": ram_available_gb, @@ -124,10 +138,12 @@ async def assess_hardware() -> dict: "qdrant_reachable": qdrant_reachable, "qdrant_collections": qdrant_collections, "searxng_reachable": searxng_reachable, + "comfyui_reachable": comfyui_reachable, + "comfyui_models": comfyui_models, } HARDWARE_STATE_PATH.write_text(json.dumps(state, indent=2)) log.info( f"HW: {ram_total_gb}GB RAM, {vram_total_mb}MB VRAM, " - f"llama={llama_reachable}, qdrant={qdrant_reachable}, searxng={searxng_reachable}" + f"llama={llama_reachable}, qdrant={qdrant_reachable}, searxng={searxng_reachable}, comfyui={comfyui_reachable}" ) return state diff --git a/node_agent/agent.py b/node_agent/agent.py index 6d2ca40..437f715 100644 --- a/node_agent/agent.py +++ b/node_agent/agent.py @@ -93,6 +93,7 @@ class AgentConfig: self.capabilities: list[str] = ["llm"] self.amqp_url: str = "amqp://caic:password@localhost:5672/caic" self.llama_port: int = 8081 + self.comfyui_port: int = 8188 self.models_dir: str = "/var/lib/caic/models" self.active_model: str = "" @@ -113,6 +114,7 @@ class AgentConfig: cfg.capabilities = [c.strip() for c in raw_caps.split(",") if c.strip()] cfg.amqp_url = parser.get(sec, "amqp_url", fallback=cfg.amqp_url) cfg.llama_port = parser.getint(sec, "llama_port", fallback=cfg.llama_port) + cfg.comfyui_port = parser.getint(sec, "comfyui_port", fallback=cfg.comfyui_port) cfg.models_dir = parser.get(sec, "models_dir", fallback=cfg.models_dir) cfg.active_model = parser.get(sec, "active_model", fallback=cfg.active_model) return cfg @@ -212,6 +214,19 @@ def get_load() -> dict: return load +def detect_capabilities(cfg: AgentConfig) -> list[str]: + caps = list(cfg.capabilities) + if "image_gen" not in caps and HAS_HTTPX: + try: + resp = httpx.get(f"http://localhost:{cfg.comfyui_port}/system_stats", timeout=3) + if resp.status_code == 200: + caps.append("image_gen") + log.info("auto-detected image_gen capability (ComfyUI on port %d)", cfg.comfyui_port) + except Exception: + pass + return caps + + # ── AMQP helpers ──────────────────────────────────────────────────────── async def declare_exchanges(channel) -> tuple: @@ -356,6 +371,152 @@ async def _wait_for_llama(port: int, timeout: int = 120, interval: int = 2) -> b return False +# ── image generation ───────────────────────────────────────────────────── + +async def handle_image_generate(cfg: AgentConfig, channel, exchanges, msg: aio_pika.IncomingMessage): + admin_ex, system_ex = exchanges + async with msg.process(): + try: + payload = json.loads(msg.body.decode()) + except json.JSONDecodeError: + return + + request_id = payload.get("request_id") + prompt = payload.get("prompt", "") + negative_prompt = payload.get("negative_prompt", "") + width = payload.get("width", 1024) + height = payload.get("height", 1024) + steps = payload.get("steps", 20) + seed = payload.get("seed", -1) + model = payload.get("model", "") + + if not prompt: + log.error("image_generate missing prompt") + return + + log.info("image generate: prompt=%s %dx%d steps=%d", prompt[:60], width, height, steps) + + now = datetime.now(timezone.utc).isoformat() + "Z" + try: + image_data = await _comfyui_generate( + cfg, prompt, negative_prompt, width, height, steps, seed, model, + ) + result_payload = { + "node_name": cfg.node_name, + "type": "image_generated", + "request_id": request_id, + "image_base64": image_data, + "timestamp": now, + } + log.info("image generate complete: request_id=%s", request_id) + except Exception as e: + result_payload = { + "node_name": cfg.node_name, + "type": "image_failed", + "request_id": request_id, + "error": str(e), + "timestamp": now, + } + log.error("image generate failed: %s", e) + + await publish(channel, system_ex, f"node.{cfg.node_name}.{result_payload['type']}", result_payload) + + +async def _comfyui_generate( + cfg: AgentConfig, prompt: str, negative_prompt: str, + width: int, height: int, steps: int, seed: int, model: str, +) -> str: + import random + import uuid as _uuid + + if not HAS_HTTPX: + raise RuntimeError("httpx not installed") + + client_id = str(_uuid.uuid4()) + if seed < 0: + seed = random.randint(0, 2**32 - 1) + + checkpoint = model or "model.safetensors" + + workflow = { + "3": { + "class_type": "KSampler", + "inputs": { + "seed": seed, + "steps": steps, + "cfg": 7.0, + "sampler_name": "euler", + "scheduler": "normal", + "denoise": 1.0, + "model": ["4", 0], + "positive": ["6", 0], + "negative": ["7", 0], + "latent_image": ["5", 0], + }, + }, + "4": { + "class_type": "CheckpointLoaderSimple", + "inputs": {"ckpt_name": checkpoint}, + }, + "5": { + "class_type": "EmptyLatentImage", + "inputs": {"width": width, "height": height, "batch_size": 1}, + }, + "6": { + "class_type": "CLIPTextEncode", + "inputs": {"text": prompt, "clip": ["4", 1]}, + }, + "7": { + "class_type": "CLIPTextEncode", + "inputs": {"text": negative_prompt or "blurry, low quality", "clip": ["4", 1]}, + }, + "8": { + "class_type": "VAEDecode", + "inputs": {"samples": ["3", 0], "vae": ["4", 2]}, + }, + "9": { + "class_type": "SaveImage", + "inputs": {"filename_prefix": f"caic_{client_id}", "images": ["8", 0]}, + }, + } + + async with httpx.AsyncClient(timeout=120) as client: + resp = await client.post( + f"http://localhost:{cfg.comfyui_port}/prompt", + json={"prompt": workflow, "client_id": client_id}, + ) + if resp.status_code != 200: + raise RuntimeError(f"ComfyUI prompt failed: {resp.status_code} {resp.text}") + + prompt_id = resp.json().get("prompt_id") + if not prompt_id: + raise RuntimeError("ComfyUI returned no prompt_id") + + deadline = time.time() + 120 + while time.time() < deadline: + resp = await client.get(f"http://localhost:{cfg.comfyui_port}/history/{prompt_id}") + if resp.status_code == 200: + history = resp.json().get(prompt_id, {}) + outputs = history.get("outputs", {}) + for node_id, node_output in outputs.items(): + images = node_output.get("images", []) + if images: + img_info = images[0] + filename = img_info.get("filename") + subfolder = img_info.get("subfolder", "") + img_type = img_info.get("type", "output") + img_resp = await client.get( + f"http://localhost:{cfg.comfyui_port}/view", + params={"filename": filename, "subfolder": subfolder, "type": img_type}, + ) + if img_resp.status_code == 200: + import base64 + return base64.b64encode(img_resp.content).decode() + await asyncio.sleep(1) + + raise RuntimeError("ComfyUI generation timed out after 120s") + + # ── main ──────────────────────────────────────────────────────────────── async def amain(): @@ -372,6 +533,9 @@ async def amain(): cfg = AgentConfig.from_ini() log.info("node_name=%s node_ip=%s", cfg.node_name, cfg.node_ip) + cfg.capabilities = detect_capabilities(cfg) + log.info("capabilities: %s", cfg.capabilities) + inventory = discover_models(cfg.models_dir) log.info("discovered %d models", len(inventory)) @@ -418,6 +582,13 @@ async def amain(): await swap_queue.bind(admin_ex, f"node.{cfg.node_name}.cmd.swap_model") await swap_queue.consume(lambda msg: handle_swap_model(cfg, channel, (admin_ex, system_ex), msg)) + # Set up image gen consumer + if "image_gen" in cfg.capabilities: + image_queue = await channel.declare_queue("", exclusive=True) + await image_queue.bind(admin_ex, f"node.{cfg.node_name}.cmd.image_generate") + await image_queue.consume(lambda msg: handle_image_generate(cfg, channel, (admin_ex, system_ex), msg)) + log.info("image generation handler registered") + log.info("listening for pings and commands") # Run forever await asyncio.Event().wait() diff --git a/routers/image.py b/routers/image.py new file mode 100644 index 0000000..42907c9 --- /dev/null +++ b/routers/image.py @@ -0,0 +1,70 @@ +"""JarvisChat routers — Image generation proxy endpoint.""" +import base64 +import logging + +from fastapi import APIRouter, HTTPException +from fastapi.responses import Response + +from cluster import CLUSTER_NODES, request_image_generate + +log = logging.getLogger("caic") +router = APIRouter() + + +def _find_image_node() -> str | None: + for name, node in CLUSTER_NODES.items(): + if node.get("status") == "active" and "image_gen" in node.get("capabilities", []): + return name + return None + + +@router.post("/api/image/generate") +async def generate_image(request_body: dict): + prompt = (request_body.get("prompt") or "").strip() + if not prompt: + raise HTTPException(status_code=400, detail="Prompt is required") + + negative_prompt = request_body.get("negative_prompt", "") + width = min(max(request_body.get("width", 1024), 256), 2048) + height = min(max(request_body.get("height", 1024), 256), 2048) + steps = min(max(request_body.get("steps", 20), 1), 50) + seed = request_body.get("seed", -1) + model = request_body.get("model", "") + + node_name = _find_image_node() + if not node_name: + raise HTTPException(status_code=503, detail="No image generation service available") + + log.info("image generate via %s: %s", node_name, prompt[:60]) + + image_b64 = await request_image_generate( + node_name=node_name, + prompt=prompt, + negative_prompt=negative_prompt, + width=width, + height=height, + steps=steps, + seed=seed, + model=model, + ) + + if image_b64 is None: + raise HTTPException(status_code=504, detail="Image generation timed out or failed") + + image_bytes = base64.b64decode(image_b64) + return Response(content=image_bytes, media_type="image/png") + + +@router.get("/api/image/status") +async def image_status(): + nodes = [] + for name, node in CLUSTER_NODES.items(): + caps = node.get("capabilities", []) + if "image_gen" in caps: + nodes.append({ + "name": name, + "status": node.get("status"), + "load": node.get("load"), + "last_seen": node.get("last_seen"), + }) + return {"available": len(nodes) > 0, "nodes": nodes} diff --git a/tests/test_image.py b/tests/test_image.py new file mode 100644 index 0000000..a1256ad --- /dev/null +++ b/tests/test_image.py @@ -0,0 +1,689 @@ +"""Tests for image generation — cluster handlers, router, node agent, hardware probe.""" +import asyncio +import base64 +import json +import os +from contextlib import asynccontextmanager +from pathlib import Path +from unittest.mock import AsyncMock, patch + +import httpx +import psutil +from fastapi.testclient import TestClient + +import app as app_module +import cluster +import config +import db +import hardware +import node_agent.agent as agent +from config import AMQP_EXCHANGE_ADMIN, AMQP_EXCHANGE_SYSTEM +from security import SESSIONS, PIN_ATTEMPTS, RATE_EVENTS + + +# ── helpers ────────────────────────────────────────────────────────────── + +def _reset(): + cluster.CLUSTER_NODES.clear() + cluster.CLUSTER_EVENTS.clear() + cluster.CLUSTER_COORDINATOR = None + cluster._pending_pings.clear() + cluster._pending_image.clear() + + +_published = [] + + +async def _fake_publish(exchange, routing_key, payload): + _published.append((exchange, routing_key, payload)) + + +def make_client(tmp_path: Path) -> TestClient: + os.environ["CAIC_ADMIN_PIN"] = "1234" + db.DB_PATH = tmp_path / "caic-image.db" + hardware.HARDWARE_STATE_PATH = tmp_path / "hardware_state.json" + SESSIONS.clear() + PIN_ATTEMPTS.clear() + RATE_EVENTS.clear() + db.init_db() + return TestClient(app_module.app, raise_server_exceptions=False) + + +def _guest_headers(client: TestClient) -> dict: + sid = client.post("/api/auth/guest", headers={"Origin": "http://testserver"}).json()["session_id"] + return {"X-Session-ID": sid, "Origin": "http://testserver"} + + +def _admin_headers(client: TestClient) -> dict: + resp = client.post("/api/auth/login", json={"pin": "1234"}, headers={"Origin": "http://testserver"}) + sid = resp.json()["session_id"] + return {"X-Session-ID": sid, "Origin": "http://testserver"} + + +class FakeMsg: + def __init__(self, body_dict: dict): + self.body = json.dumps(body_dict).encode() + + @asynccontextmanager + async def process(self): + yield + + +class FakeExchange: + def __init__(self, name=""): + self.name = name + self.published = [] + + async def publish(self, msg, routing_key): + self.published.append((msg, routing_key)) + + +class FakeChannel: + def __init__(self): + self.exchanges = {} + self.is_closed = False + + async def declare_exchange(self, name, typ, durable=True): + self.exchanges[name] = FakeExchange(name) + return self.exchanges[name] + + async def declare_queue(self, name="", exclusive=True): + return self + + async def bind(self, exchange, routing_key): + pass + + +# ── 1. cluster.handle_image_generated resolves pending request ─────────── + + +def test_handle_image_generated_resolves_pending(monkeypatch): + _reset() + monkeypatch.setattr(cluster, "publish", _fake_publish) + + cluster.CLUSTER_NODES["corsair"] = { + "name": "corsair", "type": "worker", "status": "active", + "capabilities": ["image_gen"], + } + + event = asyncio.Event() + cluster._pending_image["req-123"] = ("", event) + + asyncio.run(cluster.handle_image_generated( + AMQP_EXCHANGE_SYSTEM, "node.corsair.image_generated", + {"node_name": "corsair", "request_id": "req-123", "image_base64": "aW1hZ2U="}, + )) + + assert event.is_set() + result = cluster._pending_image.get("req-123") + assert result is not None + assert result[0] == "aW1hZ2U=" + assert cluster.CLUSTER_NODES["corsair"]["last_seen"] is not None + + +# ── 2. cluster.handle_image_failed resolves pending request ───────────── + + +def test_handle_image_failed_resolves_pending(monkeypatch): + _reset() + monkeypatch.setattr(cluster, "publish", _fake_publish) + + cluster.CLUSTER_NODES["corsair"] = { + "name": "corsair", "type": "worker", "status": "active", + } + + event = asyncio.Event() + cluster._pending_image["req-456"] = ("", event) + + asyncio.run(cluster.handle_image_failed( + AMQP_EXCHANGE_SYSTEM, "node.corsair.image_failed", + {"node_name": "corsair", "request_id": "req-456", "error": "timeout"}, + )) + + assert event.is_set() + result = cluster._pending_image.get("req-456") + assert result is not None + assert result[0] == "" + + +# ── 3. cluster.handle_image_failed unknown node ───────────────────────── + + +def test_handle_image_failed_unknown_node(caplog, monkeypatch): + _reset() + caplog.set_level("WARNING") + monkeypatch.setattr(cluster, "publish", _fake_publish) + + asyncio.run(cluster.handle_image_failed( + AMQP_EXCHANGE_SYSTEM, "node.ghost.image_failed", + {"node_name": "ghost", "request_id": "x", "error": "boom"}, + )) + + assert not any("unknown node" in rec.message for rec in caplog.records) + + +# ── 4. cluster.request_image_generate publishes command ────────────────── + + +def test_request_image_generate_publishes_command(monkeypatch): + _reset() + _published.clear() + monkeypatch.setattr(cluster, "publish", _fake_publish) + + cluster.CLUSTER_NODES["corsair"] = { + "name": "corsair", "type": "worker", "status": "active", + "capabilities": ["image_gen"], + } + + # Simulate immediate completion + async def fake_wait(): + cluster._pending_image.clear() + + original_wait_for = asyncio.wait_for + + async def patched_wait_for(coro, timeout): + cluster._pending_image["fake-id"] = ("aW1hZ2U=", asyncio.Event()) + cluster._pending_image["fake-id"][1].set() + return None + + monkeypatch.setattr(asyncio, "wait_for", patched_wait_for) + + result = asyncio.run(cluster.request_image_generate( + "corsair", "a red dragon", width=512, height=512, steps=10, + )) + + assert len(_published) == 1 + exchange, rk, payload = _published[0] + assert exchange == AMQP_EXCHANGE_ADMIN + assert rk == "node.corsair.cmd.image_generate" + assert payload["prompt"] == "a red dragon" + assert payload["width"] == 512 + assert payload["height"] == 512 + assert payload["steps"] == 10 + assert "request_id" in payload + + +# ── 5. cluster.request_image_generate unknown node ────────────────────── + + +def test_request_image_generate_unknown_node(monkeypatch): + _reset() + _published.clear() + monkeypatch.setattr(cluster, "publish", _fake_publish) + + result = asyncio.run(cluster.request_image_generate("ghost", "prompt")) + assert result is None + assert len(_published) == 0 + + +# ── 6. cluster.request_image_generate node lacks capability ───────────── + + +def test_request_image_generate_no_capability(monkeypatch): + _reset() + _published.clear() + monkeypatch.setattr(cluster, "publish", _fake_publish) + + cluster.CLUSTER_NODES["jarvis"] = { + "name": "jarvis", "type": "worker", "status": "active", + "capabilities": ["llm"], + } + + result = asyncio.run(cluster.request_image_generate("jarvis", "prompt")) + assert result is None + assert len(_published) == 0 + + +# ── 7. cluster.request_image_generate timeout ─────────────────────────── + + +def test_request_image_generate_timeout(monkeypatch): + _reset() + _published.clear() + monkeypatch.setattr(cluster, "publish", _fake_publish) + + cluster.CLUSTER_NODES["corsair"] = { + "name": "corsair", "type": "worker", "status": "active", + "capabilities": ["image_gen"], + } + + async def timeout_wait(coro, timeout): + raise asyncio.TimeoutError() + + monkeypatch.setattr(asyncio, "wait_for", timeout_wait) + + result = asyncio.run(cluster.request_image_generate("corsair", "prompt", timeout=1)) + assert result is None + assert len(cluster._pending_image) == 0 + + +# ── 8. _find_image_node selects active image_gen node ─────────────────── + + +def test_find_image_node_selects_active(): + from routers.image import _find_image_node + + _reset() + cluster.CLUSTER_NODES["corsair"] = { + "name": "corsair", "type": "worker", "status": "active", + "capabilities": ["image_gen"], + } + cluster.CLUSTER_NODES["jarvis"] = { + "name": "jarvis", "type": "worker", "status": "active", + "capabilities": ["llm"], + } + + assert _find_image_node() == "corsair" + + +def test_find_image_node_skips_inactive(): + from routers.image import _find_image_node + + _reset() + cluster.CLUSTER_NODES["corsair"] = { + "name": "corsair", "type": "worker", "status": "error", + "capabilities": ["image_gen"], + } + + assert _find_image_node() is None + + +def test_find_image_node_no_image_gen(): + from routers.image import _find_image_node + + _reset() + cluster.CLUSTER_NODES["jarvis"] = { + "name": "jarvis", "type": "worker", "status": "active", + "capabilities": ["llm"], + } + + assert _find_image_node() is None + + +# ── 9. POST /api/image/generate — no node available ───────────────────── + + +def test_image_generate_no_node_503(tmp_path): + _reset() + with make_client(tmp_path) as client: + headers = _admin_headers(client) + resp = client.post("/api/image/generate", json={"prompt": "test"}, headers=headers) + assert resp.status_code == 503 + assert "No image generation service" in resp.json()["detail"] + + +# ── 10. POST /api/image/generate — empty prompt ───────────────────────── + + +def test_image_generate_empty_prompt_400(tmp_path): + _reset() + with make_client(tmp_path) as client: + headers = _admin_headers(client) + resp = client.post("/api/image/generate", json={"prompt": ""}, headers=headers) + assert resp.status_code == 400 + assert "Prompt is required" in resp.json()["detail"] + + +# ── 11. POST /api/image/generate — happy path ────────────────────────── + + +def test_image_generate_happy_path(tmp_path, monkeypatch): + _reset() + cluster.CLUSTER_NODES["corsair"] = { + "name": "corsair", "type": "worker", "status": "active", + "capabilities": ["image_gen"], + } + + fake_png = b"\x89PNG\r\n\x1a\n" + b"\x00" * 100 + fake_b64 = base64.b64encode(fake_png).decode() + + async def fake_request_image_generate(**kwargs): + return fake_b64 + + monkeypatch.setattr("routers.image.request_image_generate", fake_request_image_generate) + + with make_client(tmp_path) as client: + headers = _admin_headers(client) + resp = client.post("/api/image/generate", json={ + "prompt": "a red dragon", + "width": 512, + "height": 512, + }, headers=headers) + assert resp.status_code == 200 + assert resp.headers["content-type"] == "image/png" + assert resp.content == fake_png + + +# ── 12. POST /api/image/generate — generation failed ─────────────────── + + +def test_image_generate_timeout_504(tmp_path, monkeypatch): + _reset() + cluster.CLUSTER_NODES["corsair"] = { + "name": "corsair", "type": "worker", "status": "active", + "capabilities": ["image_gen"], + } + + async def fake_request_image_generate(**kwargs): + return None + + monkeypatch.setattr("routers.image.request_image_generate", fake_request_image_generate) + + with make_client(tmp_path) as client: + headers = _admin_headers(client) + resp = client.post("/api/image/generate", json={"prompt": "test"}, headers=headers) + assert resp.status_code == 504 + + +# ── 13. GET /api/image/status — available ────────────────────────────── + + +def test_image_status_available(tmp_path): + _reset() + cluster.CLUSTER_NODES["corsair"] = { + "name": "corsair", "type": "worker", "status": "active", + "capabilities": ["image_gen"], + "load": {"gpu_pct": 30}, + "last_seen": "2026-07-27T00:00:00Z", + } + + with make_client(tmp_path) as client: + headers = _guest_headers(client) + resp = client.get("/api/image/status", headers=headers) + assert resp.status_code == 200 + data = resp.json() + assert data["available"] is True + assert len(data["nodes"]) == 1 + assert data["nodes"][0]["name"] == "corsair" + + +# ── 14. GET /api/image/status — no nodes ─────────────────────────────── + + +def test_image_status_unavailable(tmp_path): + _reset() + with make_client(tmp_path) as client: + headers = _guest_headers(client) + resp = client.get("/api/image/status", headers=headers) + assert resp.status_code == 200 + data = resp.json() + assert data["available"] is False + assert len(data["nodes"]) == 0 + + +# ── 15. node_agent.detect_capabilities — ComfyUI present ─────────────── + + +def test_detect_capabilities_comfyui_present(monkeypatch): + cfg = agent.AgentConfig() + cfg.comfyui_port = 8188 + + monkeypatch.setattr(agent, "HAS_HTTPX", True) + + def fake_get(url, timeout=3): + if "system_stats" in url: + class R: + status_code = 200 + return R() + raise httpx.ConnectError("refused") + + monkeypatch.setattr(httpx, "get", fake_get) + + caps = agent.detect_capabilities(cfg) + assert "image_gen" in caps + assert "llm" in caps + + +# ── 16. node_agent.detect_capabilities — ComfyUI absent ──────────────── + + +def test_detect_capabilities_comfyui_absent(monkeypatch): + cfg = agent.AgentConfig() + cfg.comfyui_port = 8188 + + monkeypatch.setattr(agent, "HAS_HTTPX", True) + monkeypatch.setattr(httpx, "get", lambda url, timeout=3: (_ for _ in ()).throw(httpx.ConnectError("refused"))) + + caps = agent.detect_capabilities(cfg) + assert "image_gen" not in caps + assert "llm" in caps + + +# ── 17. node_agent.detect_capabilities — httpx not installed ──────────── + + +def test_detect_capabilities_no_httpx(monkeypatch): + cfg = agent.AgentConfig() + monkeypatch.setattr(agent, "HAS_HTTPX", False) + + caps = agent.detect_capabilities(cfg) + assert "image_gen" not in caps + + +# ── 18. node_agent.handle_image_generate — success ───────────────────── + + +def test_node_agent_handle_image_generate_success(monkeypatch): + monkeypatch.setattr(agent, "HAS_AIO_PIKA", True) + monkeypatch.setattr(agent, "HAS_HTTPX", True) + + cfg = agent.AgentConfig() + cfg.node_name = "corsair" + cfg.comfyui_port = 8188 + + fake_png = b"\x89PNG" + b"\x00" * 50 + fake_b64 = base64.b64encode(fake_png).decode() + + async def fake_comfyui_generate(*a, **kw): + return fake_b64 + + monkeypatch.setattr(agent, "_comfyui_generate", fake_comfyui_generate) + + channel = FakeChannel() + system_ex = FakeExchange("jc.system") + channel.exchanges["jc.system"] = system_ex + + asyncio.run(agent.handle_image_generate( + cfg, channel, (FakeExchange(), system_ex), + FakeMsg({ + "request_id": "req-789", + "prompt": "a castle", + "negative_prompt": "", + "width": 1024, + "height": 1024, + "steps": 20, + "seed": 42, + "model": "", + }), + )) + + assert len(system_ex.published) == 1 + msg, rk = system_ex.published[0] + assert rk == "node.corsair.image_generated" + payload = json.loads(msg.body) + assert payload["type"] == "image_generated" + assert payload["request_id"] == "req-789" + assert payload["image_base64"] == fake_b64 + + +# ── 19. node_agent.handle_image_generate — failure ───────────────────── + + +def test_node_agent_handle_image_generate_failure(monkeypatch): + monkeypatch.setattr(agent, "HAS_AIO_PIKA", True) + monkeypatch.setattr(agent, "HAS_HTTPX", True) + + cfg = agent.AgentConfig() + cfg.node_name = "corsair" + + async def fake_comfyui_generate(*a, **kw): + raise RuntimeError("ComfyUI crashed") + + monkeypatch.setattr(agent, "_comfyui_generate", fake_comfyui_generate) + + channel = FakeChannel() + system_ex = FakeExchange("jc.system") + channel.exchanges["jc.system"] = system_ex + + asyncio.run(agent.handle_image_generate( + cfg, channel, (FakeExchange(), system_ex), + FakeMsg({"request_id": "req-fail", "prompt": "test"}), + )) + + assert len(system_ex.published) == 1 + msg, rk = system_ex.published[0] + assert rk == "node.corsair.image_failed" + payload = json.loads(msg.body) + assert payload["type"] == "image_failed" + assert "ComfyUI crashed" in payload["error"] + + +# ── 20. node_agent.handle_image_generate — empty prompt ──────────────── + + +def test_node_agent_handle_image_generate_empty_prompt(monkeypatch): + monkeypatch.setattr(agent, "HAS_AIO_PIKA", True) + + cfg = agent.AgentConfig() + cfg.node_name = "corsair" + + channel = FakeChannel() + system_ex = FakeExchange("jc.system") + + asyncio.run(agent.handle_image_generate( + cfg, channel, (FakeExchange(), system_ex), + FakeMsg({"request_id": "req-x", "prompt": ""}), + )) + + assert len(system_ex.published) == 0 + + +# ── 21. hardware.py — ComfyUI reachable ──────────────────────────────── + + +def test_assess_hardware_comfyui_reachable(tmp_path, monkeypatch): + hardware.HARDWARE_STATE_PATH = tmp_path / "hw.json" + monkeypatch.setattr(psutil, "virtual_memory", lambda: type("M", (), {"total": 16 * 1024 ** 3, "available": 8 * 1024 ** 3})()) + monkeypatch.setattr(psutil, "cpu_count", lambda: 8) + + class MockProc: + returncode = 1 + stdout = "" + + monkeypatch.setattr(hardware.subprocess, "run", lambda cmd, **kw: MockProc()) + + async def mock_get(self, url, *args, **kwargs): + class R: + status_code = 200 + def json(self): + if "CheckpointLoaderSimple" in url: + return {"CheckpointLoaderSimple": {"input": {"required": {"ckpt_name": [["model.safetensors", "other.ckpt"]]}}}} + return {} + def raise_for_status(self): + pass + if "8188" in url: + return R() + if "v1/models" in url: + class R2: + status_code = 200 + def json(self): + return {"data": []} + return R2() + if "6333" in url: + class R3: + status_code = 200 + def json(self): + return {"result": {"collections": []}} + return R3() + if "8888" in url: + class R4: + status_code = 200 + return R4() + class R5: + status_code = 200 + def json(self): + return {} + return R5() + + monkeypatch.setattr(httpx.AsyncClient, "get", mock_get) + + state = asyncio.run(hardware.assess_hardware()) + assert state["comfyui_reachable"] is True + assert "model.safetensors" in state["comfyui_models"] + + +# ── 22. hardware.py — ComfyUI unreachable ────────────────────────────── + + +def test_assess_hardware_comfyui_unreachable(tmp_path, monkeypatch): + hardware.HARDWARE_STATE_PATH = tmp_path / "hw.json" + monkeypatch.setattr(psutil, "virtual_memory", lambda: type("M", (), {"total": 16 * 1024 ** 3, "available": 8 * 1024 ** 3})()) + monkeypatch.setattr(psutil, "cpu_count", lambda: 8) + + class MockProc: + returncode = 1 + stdout = "" + + monkeypatch.setattr(hardware.subprocess, "run", lambda cmd, **kw: MockProc()) + + async def mock_get(self, url, *args, **kwargs): + if "8188" in url: + raise httpx.ConnectError("refused") + if "v1/models" in url: + class R2: + status_code = 200 + def json(self): + return {"data": []} + return R2() + if "6333" in url: + class R3: + status_code = 200 + def json(self): + return {"result": {"collections": []}} + return R3() + if "8888" in url: + class R4: + status_code = 200 + return R4() + class R5: + status_code = 200 + def json(self): + return {} + return R5() + + monkeypatch.setattr(httpx.AsyncClient, "get", mock_get) + + state = asyncio.run(hardware.assess_hardware()) + assert state["comfyui_reachable"] is False + assert state["comfyui_models"] == [] + + +# ── 23. node_agent config reads comfyui_port ─────────────────────────── + + +def test_config_from_ini_comfyui_port(tmp_path): + ini = tmp_path / "caic-node-agent.conf" + ini.write_text( + "[agent]\n" + "node_name = corsair\n" + "capabilities = llm,image_gen\n" + "comfyui_port = 8188\n" + ) + cfg = agent.AgentConfig.from_ini(str(ini)) + assert cfg.comfyui_port == 8188 + assert "image_gen" in cfg.capabilities + + +def test_config_from_ini_comfyui_port_default(): + cfg = agent.AgentConfig() + assert cfg.comfyui_port == 8188 + + +# ── 24. SUBSCRIBE_TABLE includes image gen handlers ──────────────────── + + +def test_subscribe_table_includes_image_handlers(): + routing_keys = [rks for _, rks, _ in cluster.SUBSCRIBE_TABLE] + all_keys = [rk for rks in routing_keys for rk in rks] + assert "node.*.image_generated" in all_keys + assert "node.*.image_failed" in all_keys