# -*- coding: utf-8 -*- """ swarm_llm.py — Local CPU LLM Swarm (Qwen2.5-1.5B-Instruct, Q4_K_M) ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ Why a local model? The NIM API is rate-limited (tokens/minute). Every small sub-task (JSON formatting, log summarization, brainstorm question generation, Bell Curve trimming) that goes to NIM wastes quota needed for coding. This module runs Qwen2.5-1.5B-Instruct at Q4_K_M quantization: - RAM: ~1.1 GB (leaves 14 GB free for FastAPI + data) - Speed: ~45 tok/s on 2 vCPUs (good enough for short tasks) - Model: downloaded from HF Hub on first run → cached in /tmp/models/ NIM is used ONLY for heavy coding tasks (forge_execute). SwarmLLM handles everything else. If llama-cpp-python is not installed (e.g. first boot before pip): the module silently degrades and returns a FALLBACK_STUB response, so the rest of the system still works. """ import os import logging import asyncio import time from typing import Optional logger = logging.getLogger("swarm_llm") MODEL_CACHE_DIR = os.environ.get("MODEL_CACHE_DIR", "/tmp/models") MODEL_REPO = "Qwen/Qwen2.5-1.5B-Instruct-GGUF" MODEL_FILENAME = "qwen2.5-1.5b-instruct-q4_k_m.gguf" N_CTX = 2048 # context window — matches our Bell Curve budget N_THREADS = int(os.environ.get("SWARM_THREADS", "2")) MAX_TOKENS = int(os.environ.get("SWARM_MAX_TOKENS", "256")) ENABLE_SWARM = os.environ.get("ENABLE_SWARM_LLM", "true").lower() == "true" _llm = None # loaded lazily on first call _llm_lock = None # asyncio.Lock initialised at first call def _get_lock(): global _llm_lock if _llm_lock is None: _llm_lock = asyncio.Lock() return _llm_lock def _load_model() -> Optional[object]: """Download (if needed) and load the GGUF model. Blocking — call in executor.""" global _llm if _llm is not None: return _llm if not ENABLE_SWARM: logger.info("[SwarmLLM] Disabled via ENABLE_SWARM_LLM=false") return None try: from llama_cpp import Llama except ImportError: logger.warning("[SwarmLLM] llama-cpp-python not installed. Running in stub mode.") return None model_path = os.path.join(MODEL_CACHE_DIR, MODEL_FILENAME) if not os.path.exists(model_path): logger.info(f"[SwarmLLM] Downloading {MODEL_FILENAME} from HF Hub…") try: from huggingface_hub import hf_hub_download os.makedirs(MODEL_CACHE_DIR, exist_ok=True) model_path = hf_hub_download( repo_id=MODEL_REPO, filename=MODEL_FILENAME, local_dir=MODEL_CACHE_DIR, local_dir_use_symlinks=False, ) logger.info(f"[SwarmLLM] Downloaded → {model_path}") except Exception as e: logger.error(f"[SwarmLLM] Download failed: {e}") return None logger.info(f"[SwarmLLM] Loading model (n_ctx={N_CTX}, threads={N_THREADS}, type_k=8 (Q8_0), type_v=8 (Q8_0))…") t0 = time.time() try: _llm = Llama( model_path=model_path, n_ctx=N_CTX, n_threads=N_THREADS, n_gpu_layers=0, # CPU only — HF free spaces have no GPU verbose=False, chat_format="chatml", type_k=8, # 8-bit quantization for Key Cache (Turbo Quant) type_v=8, # 8-bit quantization for Value Cache (Turbo Quant) ) logger.info(f"[SwarmLLM] Model loaded in {time.time()-t0:.1f}s") return _llm except Exception as e: logger.error(f"[SwarmLLM] Failed to load model: {e}") return None class SwarmLLM: """ Async wrapper around the local Qwen model. All inference runs in a thread executor so the FastAPI event loop is never blocked. """ def __init__(self): self._ready = False async def warm_up(self): """Pre-load the model at startup so first inference is instant.""" loop = asyncio.get_event_loop() model = await loop.run_in_executor(None, _load_model) self._ready = model is not None if self._ready: logger.info("[SwarmLLM] Warm-up complete. Ready for inference.") else: logger.warning("[SwarmLLM] Running in stub mode (model not available).") async def infer(self, prompt: str, system: str = "", max_tokens: int = MAX_TOKENS) -> str: """ Run inference on the local model. Returns the generated text. Falls back to a stub if model is not loaded. """ if not self._ready: return self._stub(prompt) lock = _get_lock() async with lock: loop = asyncio.get_event_loop() result = await loop.run_in_executor( None, lambda: self._sync_infer(prompt, system, max_tokens), ) return result def _sync_infer(self, prompt: str, system: str, max_tokens: int) -> str: global _llm if _llm is None: return self._stub(prompt) try: messages = [] if system: messages.append({"role": "system", "content": system}) messages.append({"role": "user", "content": prompt}) response = _llm.create_chat_completion( messages=messages, max_tokens=max_tokens, temperature=0.3, stop=["<|im_end|>", ""], ) text = response["choices"][0]["message"]["content"].strip() logger.debug("[SwarmLLM] Infer complete: %d chars", len(text)) return text except Exception as e: logger.error(f"[SwarmLLM] Inference error: {e}") return self._stub(prompt) @staticmethod def _stub(prompt: str) -> str: """Fallback when model is unavailable — returns a safe placeholder.""" return "[SwarmLLM unavailable — NIM will handle this task]" # ── High-Level Task Shortcuts ────────────────────────────────────────────── async def summarize(self, text: str, max_words: int = 80) -> str: """Summarise a long text into ≤ max_words words. Used to enforce Bell Curve budget.""" prompt = ( f"Summarise the following in ≤ {max_words} words. " f"Be dense with information. No filler sentences.\n\n{text[:3000]}" ) return await self.infer(prompt, system="You are a precise technical summariser.") async def format_json(self, raw: str) -> str: """Extract and clean a JSON object from a messy LLM response.""" prompt = ( "Extract the JSON object from the following text. " "Return ONLY the JSON, no markdown fences, no explanation.\n\n" + raw[:2000] ) return await self.infer(prompt, system="You are a JSON extractor. Output only valid JSON.") async def generate_brainstorm_questions(self, goal: str, n: int = 5) -> list: """Generate n 'What if…' brainstorm questions for the hourly swarm cycle.""" prompt = ( f"Generate exactly {n} creative 'What if…' questions to improve: '{goal}'. " f"Each question should be a novel technical idea. " f"Format: one question per line, no numbering." ) raw = await self.infer(prompt, system="You are a creative technical brainstormer.", max_tokens=400) questions = [q.strip() for q in raw.strip().splitlines() if q.strip()] return questions[:n] async def smart_trim(self, text: str, char_limit: int) -> str: """ If text exceeds char_limit, ask the local model to summarise it to fit. Better than hard-cutting. Used in Bell Curve budget enforcement. """ if len(text) <= char_limit: return text target_words = char_limit // 6 # rough chars-to-words ratio summary = await self.summarize(text, max_words=target_words) return summary # Singleton — import and use directly swarm = SwarmLLM()