diff --git a/.gitignore b/.gitignore index 2848c67..33326af 100644 --- a/.gitignore +++ b/.gitignore @@ -8,6 +8,7 @@ whatsapp_bridge/node_modules/ data/ config/api_keys.json config/settings.json +config/mcp.json memory/long_term.json *.session *.session-journal @@ -20,4 +21,5 @@ avatar3d/models_archivio/ *.log .DS_Store .venv-chatterbox/ +.venv-mflux/ tts_chatterbox/test_*.py diff --git a/avatar/engines/base.py b/avatar/engines/base.py index 7219d01..89eff4c 100644 --- a/avatar/engines/base.py +++ b/avatar/engines/base.py @@ -67,3 +67,63 @@ class History: def clear(self) -> None: self.messages, self.meta = [], {} self.save() + + +# ── Compattazione della cronologia (motori locali) ─────────────────────────── +COMPACT_AFTER = 30 # oltre questo numero di messaggi i più vecchi vengono riassunti +COMPACT_KEEP = 12 # messaggi recenti lasciati per esteso + + +def render_messages(msgs: list, per_msg: int = 600) -> str: + out = [] + for m in msgs: + role = m.get("role") + text = str(m.get("content") or "").strip() + if role == "user": + out.append("Utente: " + text[:per_msg]) + elif role == "assistant": + calls = m.get("tool_calls") or [] + if calls: + names = ", ".join(str((c.get("function") or {}).get("name", "?")) for c in calls) + text = (text + f" [usa strumenti: {names}]").strip() + if text: + out.append("Assistente: " + text[:per_msg]) + elif role == "tool": + out.append(f"Risultato di {m.get('name', 'strumento')}: " + text[:200]) + return "\n".join(out) + + +def summary_prompt(previous: str, new_text: str) -> str: + return ("Riassumi in italiano, in modo compatto (al massimo 150 parole), ciò che serve per continuare la conversazione: " + "richieste dell'utente, cose fatte o decise, fatti e preferenze emersi, questioni aperte. Niente saluti, niente frasi generiche.\n\n" + + (f"Riassunto precedente:\n{previous}\n\n" if previous else "") + + f"Nuovi messaggi:\n{new_text}\n\nRiassunto aggiornato:") + + +def compact_history(history: "History", summarize: Callable[[str], str]) -> bool: + """Se la cronologia è lunga, riassume i messaggi più vecchi in history.meta['summary'] e li rimuove.""" + msgs = history.messages + if len(msgs) <= COMPACT_AFTER: + return False + cut = len(msgs) - COMPACT_KEEP + while cut < len(msgs) and msgs[cut].get("role") != "user": + cut += 1 + if cut <= 0 or cut >= len(msgs): + return False + text = render_messages(msgs[:cut]) + try: + summary = (summarize(summary_prompt(str(history.meta.get("summary") or ""), text)) or "").strip() + except Exception as err: + print(f"[cronologia] riassunto fallito: {err}") + return False + if not summary: + return False + history.meta["summary"] = summary[:2000] + history.messages = msgs[cut:] + history.save() + return True + + +def summary_block(history: "History") -> str: + s = str(history.meta.get("summary") or "").strip() + return f"\n\nRiassunto della conversazione precedente (i messaggi più vecchi sono stati compattati):\n{s}" if s else "" diff --git a/avatar/engines/claude_code.py b/avatar/engines/claude_code.py index 096d1f1..54b3385 100644 --- a/avatar/engines/claude_code.py +++ b/avatar/engines/claude_code.py @@ -112,10 +112,33 @@ class ClaudeCodeEngine: flags["tools"] = flags["tools"] + ",Read" flags["allowed"] = [*flags["allowed"], f"Read(//{attach_path.lstrip('/')})"] mcp = {"mcpServers": {"avatar": {"command": sys.executable, "args": [str(MCP_SCRIPT)]}}} + allowed_mcp = list(MEMORY_TOOLS) + disallowed_mcp: list[str] = [] + try: + from avatar.mcp_client import load_config + for name, c in load_config().items(): + if c.get("disabled") or name == "avatar": + continue + mcp["mcpServers"][name] = {k: v for k, v in c.items() if k in ("command", "args", "env", "cwd", "url", "headers", "type")} + if c.get("url") and "type" not in c: + mcp["mcpServers"][name]["type"] = "http" + if c.get("tools") or c.get("exclude"): + # Filtri: consenti solo gli strumenti scelti e nega gli altri (l'elenco completo arriva dal client in-app, se collegato). + from avatar.mcp_client import manager, selected_tools + srv = manager.servers.get(name) + if srv is not None: + chosen = {t["name"] for t in selected_tools(c, srv.tools)} + allowed_mcp += [f"mcp__{name}__{t}" for t in sorted(chosen)] + disallowed_mcp += [f"mcp__{name}__{t['name']}" for t in srv.tools if t["name"] not in chosen] + continue + allowed_mcp.append(f"mcp__{name}") + except Exception as err: + print(f"[claude code] server MCP esterni ignorati: {err}") args = ["-p", "--output-format", "stream-json", "--verbose", "--include-partial-messages", "--model", self.model, "--effort", self.effort, - "--tools", flags["tools"], "--allowedTools", ",".join(flags["allowed"] + MEMORY_TOOLS), + "--tools", flags["tools"], "--allowedTools", ",".join(flags["allowed"] + allowed_mcp), "--mcp-config", json.dumps(mcp), "--strict-mcp-config", + *(["--disallowedTools", ",".join(disallowed_mcp)] if disallowed_mcp else []), "--system-prompt-snapshot", "off", "--resume" if resume else "--session-id", session_id] if flags["mode"]: diff --git a/avatar/engines/mlx_engine.py b/avatar/engines/mlx_engine.py index cd2a4fb..0148ed3 100644 --- a/avatar/engines/mlx_engine.py +++ b/avatar/engines/mlx_engine.py @@ -16,7 +16,7 @@ from pathlib import Path from avatar.memory_tools import memory_prompt, openai_tools, parse_args, run_tool from avatar.plugins import registry from avatar.websearch import brave_search -from .base import Emit, History, persona_text, today_label, user_block +from .base import Emit, History, compact_history, persona_text, summary_block, today_label, user_block from .openai_compat import SEARCH_TOOL, Aborted MAX_ROUNDS = 8 @@ -323,7 +323,7 @@ class MLXEngine: "- Chiama gli strumenti solo nel formato previsto dal tuo template e mai descrivendoli a parole.") note += ("\n- Per informazioni aggiornate chiama cerca_web e rispondi in base ai risultati." if self.search_api_key else "\n- Non hai accesso al web: se ti chiedono informazioni aggiornate, dillo chiaramente.") - return f"{persona_text(self.assistant_name)}\n\n{user_block(self.user_name, memory_prompt())}{note}" + return f"{persona_text(self.assistant_name)}\n\n{user_block(self.user_name, memory_prompt())}{summary_block(self.history)}{note}" def _recent(self) -> list: msgs = self.history.messages @@ -443,6 +443,9 @@ class MLXEngine: full = full.strip() or "(nessuna risposta)" self.history.save() emit({"type": "done", "text": full, "sources": sources}) + with _lock: + if compact_history(self.history, self.summarize): + _loaded["snap"] = None # il prompt cambia: la fotografia non vale più except Aborted: self._rollback() emit({"type": "error", "message": "Risposta interrotta.", "aborted": True}) @@ -450,6 +453,22 @@ class MLXEngine: self._rollback() emit({"type": "error", "message": f"Modello interno: {err}"}) + def summarize(self, prompt: str) -> str: + """Genera un testo breve col modello caricato, senza strumenti né cache condivisa (usato per compattare la cronologia).""" + from mlx_lm import stream_generate + from mlx_lm.sample_utils import make_sampler + from mlx_lm.models.cache import make_prompt_cache + msgs = [{"role": "system", "content": "Sei un assistente che riassume conversazioni in italiano, in modo fedele e compatto."}, + {"role": "user", "content": prompt}] + text = self.tokenizer.apply_chat_template(msgs, add_generation_prompt=True, tokenize=False, enable_thinking=False) + f = FORMATS[self.fmt] + think = TagFilter(*f["think"], keep=False) + think.inside = text.rstrip().endswith(f["think"][0]) + out = "" + for r in stream_generate(self.model, self.tokenizer, text, max_tokens=400, sampler=make_sampler(temp=0.3, top_p=0.9), prompt_cache=make_prompt_cache(self.model)): + out += think.push(r.text) + return (out + think.flush()).strip() + def _rollback(self) -> None: msgs = self.history.messages if msgs and msgs[-1].get("role") == "user": diff --git a/avatar/engines/openai_compat.py b/avatar/engines/openai_compat.py index 03ad19b..666f520 100644 --- a/avatar/engines/openai_compat.py +++ b/avatar/engines/openai_compat.py @@ -9,7 +9,7 @@ import openai from avatar.memory_tools import memory_prompt, openai_tools, parse_args, run_tool from avatar.plugins import registry from avatar.websearch import brave_search -from .base import Emit, History, persona_text, today_label, user_block +from .base import Emit, History, compact_history, persona_text, summary_block, today_label, user_block MAX_ROUNDS = 8 MAX_HISTORY = 30 @@ -98,7 +98,7 @@ class OpenAICompatEngine: else "\n- Non hai accesso al web: se ti chiedono informazioni aggiornate, dillo chiaramente.") else: note = "\n\nNota: in questa modalità non hai strumenti (niente memoria automatica né ricerca web)." - return {"role": "system", "content": f"{persona_text(self.assistant_name)}\n\n{user_block(self.user_name, memory_prompt())}{note}"} + return {"role": "system", "content": f"{persona_text(self.assistant_name)}\n\n{user_block(self.user_name, memory_prompt())}{summary_block(self.history)}{note}"} def _recent(self) -> list: msgs = self.history.messages @@ -210,6 +210,7 @@ class OpenAICompatEngine: full = full.strip() or "(nessuna risposta)" self.history.save() emit({"type": "done", "text": full, "sources": sources}) + compact_history(self.history, self.summarize) except Aborted: self._rollback() emit({"type": "error", "message": "Risposta interrotta.", "aborted": True}) @@ -217,6 +218,13 @@ class OpenAICompatEngine: self._rollback() emit({"type": "error", "message": friendly_error(err, self.base_url)}) + def summarize(self, prompt: str) -> str: + r = self.client.chat.completions.create(model=self.model, max_tokens=400, temperature=0.3, messages=[ + {"role": "system", "content": "Sei un assistente che riassume conversazioni in italiano, in modo fedele e compatto."}, + {"role": "user", "content": prompt}]) + flt = ThinkFilter() + return (flt.push(r.choices[0].message.content or "") + flt.carry).strip() + def _rollback(self) -> None: msgs = self.history.messages if msgs and msgs[-1].get("role") == "user": diff --git a/avatar/mcp_client.py b/avatar/mcp_client.py new file mode 100644 index 0000000..51abfeb --- /dev/null +++ b/avatar/mcp_client.py @@ -0,0 +1,318 @@ +"""Client MCP: collega server esterni (stdio o HTTP) e ne espone gli strumenti ai motori dell'app. + +Configurazione in config/mcp.json, stesso formato di Claude Desktop: +{"mcpServers": {"nome": {"command": "npx", "args": ["-y", "@modelcontextprotocol/server-filesystem", "/cartella"], "env": {}}, + "remoto": {"url": "https://esempio.it/mcp", "headers": {"Authorization": "Bearer …"}}}} +Gli strumenti compaiono come _; con il motore Claude Code i server vengono passati direttamente a `claude`. +Opzioni per server: "tools": ["click", "browser_*"] (solo questi), "exclude": [...], "description_limit": 700 (caratteri +di descrizione per strumento: server come Cua Driver ne espongono decine con descrizioni lunghissime), "disabled": true. +""" +from __future__ import annotations + +import base64 +import fnmatch +import json +import os +import re +import subprocess +import threading +import time +import traceback +import urllib.request +from pathlib import Path +from typing import Any + +from avatar.settings import CONFIG_DIR, DATA_DIR + +CONFIG_FILE = CONFIG_DIR / "mcp.json" +PROTOCOL = "2025-06-18" +CALL_TIMEOUT = 120 +EXTRA_PATH = ("/opt/homebrew/bin", "/opt/homebrew/opt/node@22/bin", "/usr/local/bin", str(Path.home() / ".local/bin")) + + +def load_config() -> dict[str, dict]: + try: + d = json.loads(CONFIG_FILE.read_text(encoding="utf-8")) + return {k: v for k, v in (d.get("mcpServers") or {}).items() if isinstance(v, dict)} + except FileNotFoundError: + return {} + except Exception as err: + raise RuntimeError(f"config/mcp.json non valido: {err}") + + +def save_config(text: str) -> dict[str, dict]: + d = json.loads(text) + if not isinstance(d, dict) or not isinstance(d.get("mcpServers", {}), dict): + raise ValueError("serve un oggetto con la chiave mcpServers") + CONFIG_DIR.mkdir(parents=True, exist_ok=True) + CONFIG_FILE.write_text(json.dumps(d, indent=2, ensure_ascii=False), encoding="utf-8") + return load_config() + + +def _env_with_path(extra: dict | None) -> dict: + env = dict(os.environ) + env["PATH"] = ":".join([*EXTRA_PATH, env.get("PATH", "")]) + env.pop("CLAUDECODE", None) + for k, v in (extra or {}).items(): + env[str(k)] = str(v) + return env + + +class MCPServer: + def __init__(self, name: str, cfg: dict) -> None: + self.name, self.cfg = name, cfg + self.tools: list[dict] = [] + self.error = "" + self._id = 0 + self._lock = threading.Lock() + self._pending: dict[int, dict] = {} + self._proc: subprocess.Popen | None = None + self._session = "" + self.stdio = bool(cfg.get("command")) + + # ── trasporto ────────────────────────────────────────────────────────── + def start(self) -> None: + if self.stdio: + cmd = [str(self.cfg["command"]), *[str(a) for a in self.cfg.get("args", [])]] + self._proc = subprocess.Popen(cmd, stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, + text=True, bufsize=1, env=_env_with_path(self.cfg.get("env")), cwd=self.cfg.get("cwd") or None) + threading.Thread(target=self._reader, daemon=True).start() + threading.Thread(target=self._drain_stderr, daemon=True).start() + res = self._request("initialize", {"protocolVersion": PROTOCOL, "capabilities": {}, "clientInfo": {"name": "LuZa", "version": "1.0"}}, timeout=60) + self._notify("notifications/initialized") + self.server_info = res.get("serverInfo", {}) + self.tools = list(self._request("tools/list", {}, timeout=60).get("tools", [])) + + def stop(self) -> None: + if self._proc: + try: + self._proc.terminate() + self._proc.wait(timeout=5) + except Exception: + try: + self._proc.kill() + except Exception: + pass + self._proc = None + + def _reader(self) -> None: + proc = self._proc + try: + for line in proc.stdout: + line = line.strip() + if not line: + continue + try: + msg = json.loads(line) + except Exception: + continue + mid = msg.get("id") + if mid in self._pending: + self._pending[mid]["msg"] = msg + self._pending[mid]["ev"].set() + except Exception: + pass + for p in list(self._pending.values()): + p["msg"] = {"error": {"message": "server MCP terminato"}} + p["ev"].set() + + def _drain_stderr(self) -> None: + try: + log = (DATA_DIR / "mcp").joinpath(f"{self.name}.log") + log.parent.mkdir(parents=True, exist_ok=True) + with open(log, "a", encoding="utf-8") as f: + for line in self._proc.stderr: + f.write(line) + except Exception: + pass + + def _next_id(self) -> int: + with self._lock: + self._id += 1 + return self._id + + def _notify(self, method: str, params: dict | None = None) -> None: + msg = {"jsonrpc": "2.0", "method": method, "params": params or {}} + if self.stdio: + self._proc.stdin.write(json.dumps(msg) + "\n") + self._proc.stdin.flush() + else: + try: + self._http(msg, timeout=15) + except Exception: + pass + + def _request(self, method: str, params: dict, timeout: float = CALL_TIMEOUT) -> dict: + mid = self._next_id() + msg = {"jsonrpc": "2.0", "id": mid, "method": method, "params": params} + if self.stdio: + if not self._proc or self._proc.poll() is not None: + raise RuntimeError(f"server MCP '{self.name}' non in esecuzione") + slot = {"ev": threading.Event(), "msg": None} + self._pending[mid] = slot + self._proc.stdin.write(json.dumps(msg) + "\n") + self._proc.stdin.flush() + if not slot["ev"].wait(timeout): + self._pending.pop(mid, None) + raise RuntimeError(f"server MCP '{self.name}': nessuna risposta entro {int(timeout)} s") + self._pending.pop(mid, None) + resp = slot["msg"] + else: + resp = self._http(msg, timeout=timeout, want_id=mid) + if resp is None: + raise RuntimeError("risposta vuota") + if "error" in resp: + e = resp["error"] + raise RuntimeError(e.get("message") if isinstance(e, dict) else str(e)) + return resp.get("result") or {} + + def _http(self, msg: dict, timeout: float, want_id: int | None = None) -> dict | None: + headers = {"Content-Type": "application/json", "Accept": "application/json, text/event-stream", **(self.cfg.get("headers") or {})} + if self._session: + headers["Mcp-Session-Id"] = self._session + req = urllib.request.Request(str(self.cfg["url"]), data=json.dumps(msg).encode(), headers=headers, method="POST") + with urllib.request.urlopen(req, timeout=timeout) as r: + sid = r.headers.get("Mcp-Session-Id") + if sid: + self._session = sid + ctype = r.headers.get("Content-Type", "") + body = r.read().decode("utf-8", errors="replace") + if want_id is None: + return None + if "text/event-stream" in ctype: + for line in body.splitlines(): + if line.startswith("data:"): + try: + m = json.loads(line[5:].strip()) + except Exception: + continue + if m.get("id") == want_id: + return m + raise RuntimeError("nessuna risposta nello stream SSE") + return json.loads(body) if body.strip() else None + + # ── strumenti ────────────────────────────────────────────────────────── + def call(self, tool: str, args: dict) -> str: + res = self._request("tools/call", {"name": tool, "arguments": args or {}}) + parts = [] + for c in res.get("content") or []: + t = c.get("type") + if t == "text": + parts.append(str(c.get("text", ""))) + elif t == "image" and c.get("data"): + cap = DATA_DIR / "captures" + cap.mkdir(parents=True, exist_ok=True) + ext = {"image/png": "png", "image/jpeg": "jpg", "image/webp": "webp"}.get(c.get("mimeType", ""), "png") + p = cap / f"mcp-{self.name}-{int(time.time())}.{ext}" + p.write_bytes(base64.b64decode(c["data"])) + parts.append(f"IMMAGINE: {p}") + elif t == "resource": + r = c.get("resource") or {} + parts.append(str(r.get("text") or r.get("uri") or "")) + out = "\n".join(p for p in parts if p).strip() + if res.get("isError"): + return "Errore dello strumento: " + (out or "senza dettagli") + if "structuredContent" in res and not out: + out = json.dumps(res["structuredContent"], ensure_ascii=False) + return out[:20000] or "Fatto." + + +def selected_tools(cfg: dict, tools: list[dict]) -> list[dict]: + """Applica i filtri tools/exclude della configurazione all'elenco degli strumenti di un server.""" + inc, exc = cfg.get("tools") or [], cfg.get("exclude") or [] + out = [] + for t in tools: + n = t["name"] + if inc and not any(fnmatch.fnmatchcase(n, pat) for pat in inc): + continue + if exc and any(fnmatch.fnmatchcase(n, pat) for pat in exc): + continue + out.append(t) + return out + + +def _safe(name: str) -> str: + return re.sub(r"[^a-zA-Z0-9_]", "_", name)[:64].strip("_") or "strumento" + + +class MCPManager: + """Avvia i server configurati e registra i loro strumenti nel registro dei plugin.""" + + def __init__(self) -> None: + self.servers: dict[str, MCPServer] = {} + self.errors: dict[str, str] = {} + self.registry = None + self.log = print + self._lock = threading.Lock() + + def start(self, registry, log=None) -> None: + self.registry = registry + if log: + self.log = log + threading.Thread(target=self.reload, daemon=True).start() + + def reload(self) -> None: + with self._lock: + self._stop_all() + try: + cfg = load_config() + except Exception as err: + self.errors["config"] = str(err) + self.log(f"ERR: MCP — {err}") + return + self.errors.pop("config", None) + for name, c in cfg.items(): + if c.get("disabled"): + continue + srv = MCPServer(name, c) + try: + srv.start() + except Exception as err: + srv.stop() + self.errors[name] = str(err) + self.log(f"ERR: MCP «{name}» — {err}") + traceback.print_exc() + continue + self.errors.pop(name, None) + self.servers[name] = srv + self._register(srv) + self.log(f"SYS: MCP «{name}»: {len(srv.tools)} strumenti.") + + def _register(self, srv: MCPServer) -> None: + if self.registry is None: + return + tools = [] + limit = int(srv.cfg.get("description_limit") or 700) + for t in selected_tools(srv.cfg, srv.tools): + tname = _safe(f"{srv.name}_{t['name']}") + desc = " ".join((t.get("description") or t["name"]).split()) + if len(desc) > limit: + desc = desc[:limit].rsplit(" ", 1)[0] + " …" + tools.append({"name": tname, "description": f"[{srv.name}] {desc}", + "parameters": t.get("inputSchema") or {"type": "object", "properties": {}}, + "run": (lambda params, ctx, _s=srv, _n=t["name"]: _s.call(_n, params))}) + info = getattr(srv, "server_info", {}) or {} + self.registry.add_external(f"mcp:{srv.name}", f"Server MCP {info.get('name') or srv.name}: {len(tools)} strumenti", tools) + + def _stop_all(self) -> None: + for name, srv in list(self.servers.items()): + if self.registry is not None: + self.registry.remove_module(f"mcp:{name}") + srv.stop() + self.servers.clear() + + def stop(self) -> None: + with self._lock: + self._stop_all() + + def status(self) -> str: + lines = [] + for name, srv in self.servers.items(): + sel = selected_tools(srv.cfg, srv.tools) + lines.append(f"• {name}: {len(sel)} strumenti attivi su {len(srv.tools)} — " + ", ".join(t["name"] for t in sel[:8]) + (" …" if len(sel) > 8 else "")) + for name, err in self.errors.items(): + lines.append(f"✗ {name}: {err}") + return "\n".join(lines) or "Nessun server MCP configurato." + + +manager = MCPManager() diff --git a/avatar/persona.md b/avatar/persona.md index fc48458..9586bd5 100644 --- a/avatar/persona.md +++ b/avatar/persona.md @@ -42,3 +42,15 @@ Sei **Ava**, l'assistente personale di chi ti parla. Vivi in un'app sul suo Mac - Con gli strumenti "casa_…" controlli la casa tramite Home Assistant: luci, prese, clima, tapparelle, media, scene e sensori. Se l'utente nomina un dispositivo che non trovi, cerca con casa_dispositivi prima di dire che non esiste; per comandi su più stanze usa casa_chiedi. - Conferma a voce cosa hai fatto in poche parole ("Luce cucina accesa"). Serrature e allarme solo dopo la conferma esplicita dell'utente. + +## Computer + +- Se sono disponibili gli strumenti "cua_driver_…" puoi ispezionare e comandare le app del Mac (aprire finestre, leggere e premere elementi, scrivere nei campi) senza togliere il focus all'utente. Usali quando l'utente chiede un'azione dentro un'app e non esiste un plugin dedicato: per mail, calendario, promemoria, note, musica, file e casa preferisci sempre i plugin. +- Prima di agire osserva: leggi la finestra o l'elemento, poi esegui un passo alla volta e verifica il risultato. Niente sequenze lunghe alla cieca. +- Chiedi conferma prima di azioni difficili da annullare: inviare messaggi o mail, pagare, cancellare, chiudere documenti non salvati, modificare impostazioni. Se un passaggio non riesce due volte, fermati e spiega cosa vedi invece di insistere. +- Non inserire mai password o codici, e non agire in finestre di banche o pagamenti se non su richiesta esplicita e puntuale dell'utente. +- Metodo del driver: prima list_apps o list_windows per trovare pid e finestra, poi get_window_state per leggere gli elementi, poi agisci con click, type_text o set_value indicando l'elemento (element_token o element_index + window_id) invece delle coordinate, infine rileggi lo stato o usa verify_state per controllare. Lancia le app con launch_app (in background), e porta in primo piano con bring_to_front solo se l'utente lo chiede. + +## Immagini + +- Puoi creare immagini con immagine_genera: traduci la richiesta in una descrizione inglese ricca (soggetto, scena, luce, stile, inquadratura) e avvisa che ci vuole qualche decina di secondi. Quando è pronta, dilla in poche parole senza descrivere i dettagli tecnici. diff --git a/avatar/plugins.py b/avatar/plugins.py index a15b997..4084736 100644 --- a/avatar/plugins.py +++ b/avatar/plugins.py @@ -91,6 +91,24 @@ class PluginRegistry: self.modules[mod_name] = info return self + def add_external(self, module: str, description: str, tools: list[dict]) -> None: + """Registra strumenti forniti a runtime (server MCP), sostituendo quelli dello stesso modulo.""" + self.load() + self.remove_module(module) + info = {"name": module, "description": description, "valid": True, "error": "", "tools": []} + for t in tools: + name = str(t["name"]) + if name in self.tools: # nome già usato da un plugin: prefissa + name = f"x_{name}"[:64] + self.tools[name] = Tool(module, name, str(t.get("description", "")), _normalize_schema(t.get("parameters", {})), t["run"]) + info["tools"].append(name) + self.modules[module] = info + + def remove_module(self, module: str) -> None: + for name in [n for n, t in self.tools.items() if t.module == module]: + self.tools.pop(name, None) + self.modules.pop(module, None) + def _enabled(self, module: str) -> bool: try: from memory.config_manager import get_plugin_enabled diff --git a/avatar/settings.py b/avatar/settings.py index 297a803..7a3e811 100644 --- a/avatar/settings.py +++ b/avatar/settings.py @@ -52,6 +52,8 @@ DEFAULTS: dict[str, Any] = { "server_assist_dir": "/opt/aiserverassistance", # cartella del bot sul server ponte (contiene .env) "server_assist_bot": "luzaserver_bot", # username Telegram del bot sysadmin "homeassistant_url": "http://homeassistant.local:8123", # Home Assistant (token nel portachiavi) + "immagini_famiglia": "z-image-turbo", # z-image-turbo | schnell | dev | qwen (comando mflux e passi predefiniti) + "immagini_modello": "", # repo Hugging Face o cartella; vuoto = predefinito della famiglia } SECRET_KEYS = ("anthropic_api_key", "local_api_key", "telegram_api_hash", "elevenlabs_api_key", "homeassistant_token") diff --git a/avatar/settings_dialog.py b/avatar/settings_dialog.py index 86bb76f..4cd000b 100644 --- a/avatar/settings_dialog.py +++ b/avatar/settings_dialog.py @@ -1,6 +1,7 @@ """Finestra impostazioni: motore, chiavi, voce e riconoscimento vocale.""" from __future__ import annotations +import json import sys import threading from pathlib import Path @@ -65,7 +66,7 @@ class SettingsDialog(QDialog): self.tabs.addTab(sc, title) lay.addStretch(1) return lay - pg_motore, pg_voce, pg_msg, pg_mon, pg_casa = page("Motore"), page("Voce e avatar"), page("Messaggistica"), page("Monitor"), page("Casa") + pg_motore, pg_voce, pg_msg, pg_mon, pg_casa, pg_tool = page("Motore"), page("Voce e avatar"), page("Messaggistica"), page("Monitor"), page("Casa"), page("Strumenti") def add(lay, item): # inserisce prima dello stretch finale if isinstance(item, QWidget): lay.insertWidget(lay.count() - 1, item) else: lay.insertLayout(lay.count() - 1, item) @@ -260,6 +261,36 @@ class SettingsDialog(QDialog): self.ha_hint = QLabel(""); self.ha_hint.setWordWrap(True); form_ha.addRow("", self.ha_hint) add(pg_casa, form_ha) + # ── Strumenti (server MCP esterni) ──────────────────────────────── + from PyQt6.QtWidgets import QPlainTextEdit + from .mcp_client import CONFIG_FILE as MCP_FILE, manager as mcp_manager + form_mcp = QFormLayout() + form_mcp.addRow(QLabel("Server MCP esterni: stesso formato di Claude Desktop (mcpServers con command/args/env, oppure url/headers). " + "Gli strumenti diventano disponibili a tutti i motori; con Claude Code i server vengono passati direttamente a claude.")) + self.mcp_text = QPlainTextEdit() + try: + self.mcp_text.setPlainText(MCP_FILE.read_text(encoding="utf-8")) + except Exception: + self.mcp_text.setPlainText('{\n "mcpServers": {\n }\n}\n') + self.mcp_text.setMinimumHeight(220) + form_mcp.addRow(self.mcp_text) + row = QHBoxLayout() + b = QPushButton("Applica e ricollega"); b.clicked.connect(self._mcp_apply); row.addWidget(b) + b = QPushButton("Esempio"); b.clicked.connect(self._mcp_example); row.addWidget(b) + row.addStretch() + form_mcp.addRow(row) + self.mcp_hint = QLabel(mcp_manager.status()); self.mcp_hint.setWordWrap(True); form_mcp.addRow(self.mcp_hint) + add(pg_tool, form_mcp) + + # ── Immagini (mflux) ────────────────────────────────────────────── + form_img = QFormLayout() + form_img.addRow(QLabel("Generazione immagini in locale (mflux, MLX). Il primo uso scarica il modello da Hugging Face.")) + self.img_fam = _combo([("z-image-turbo", "Z-Image Turbo (veloce, 8 passi, consigliato)"), ("schnell", "FLUX.1 schnell (4 passi)"), ("dev", "FLUX.1 dev (lento, licenza non commerciale)"), ("qwen", "Qwen-Image (pesante, ottimo col testo)")], s.get("immagini_famiglia") or "z-image-turbo") + form_img.addRow("Famiglia", self.img_fam) + self.img_model = QLineEdit(str(s.get("immagini_modello") or "")); self.img_model.setPlaceholderText("vuoto = predefinito (es. mflux-community/z-image-turbo-mflux-q4)") + form_img.addRow("Modello (repo Hugging Face)", self.img_model) + add(pg_tool, form_img) + btns = QHBoxLayout(); btns.addStretch() cancel = QPushButton("Annulla"); cancel.clicked.connect(self.reject); btns.addWidget(cancel) save = QPushButton("Salva"); save.setObjectName("primary"); save.clicked.connect(self._save); btns.addWidget(save) @@ -407,6 +438,25 @@ class SettingsDialog(QDialog): self._async.emit("el_error", str(err)[:120]) threading.Thread(target=work, daemon=True).start() + def _mcp_example(self) -> None: + self.mcp_text.setPlainText(json.dumps({"mcpServers": { + "filesystem": {"command": "npx", "args": ["-y", "@modelcontextprotocol/server-filesystem", str(Path.home() / "Documents")]}, + "remoto": {"url": "https://esempio.it/mcp", "headers": {"Authorization": "Bearer TOKEN"}, "disabled": True}, + }}, indent=2, ensure_ascii=False)) + + def _mcp_apply(self) -> None: + from .mcp_client import save_config, manager as mcp_manager + try: + save_config(self.mcp_text.toPlainText()) + except Exception as err: + self.mcp_hint.setText(f"JSON non valido: {err}"); return + self.mcp_hint.setText("Collego i server…") + + def go(): + mcp_manager.reload() + self._async.emit("mcp", mcp_manager.status()) + threading.Thread(target=go, daemon=True).start() + def _ha_check(self) -> None: vals = {"homeassistant_url": self.ha_url.text().strip().rstrip("/")} if self.ha_token.text().strip(): @@ -427,6 +477,8 @@ class SettingsDialog(QDialog): def _on_async(self, kind: str, payload) -> None: if kind == "ha": self.ha_hint.setText(str(payload)); return + if kind == "mcp": + self.mcp_hint.setText(str(payload)); return if kind == "vb": self.vb_profile.clear() for pid, label in payload: @@ -517,6 +569,7 @@ class SettingsDialog(QDialog): "server_assist_dir": self.srv_dir.text().strip() or "/opt/aiserverassistance", "server_assist_bot": self.srv_bot.text().strip().lstrip("@") or "luzaserver_bot", "homeassistant_url": self.ha_url.text().strip().rstrip("/"), + "immagini_famiglia": self.img_fam.currentData() or "z-image-turbo", "immagini_modello": self.img_model.text().strip(), } if self.ha_token.text().strip(): values["homeassistant_token"] = self.ha_token.text().strip() diff --git a/config/mcp.example.json b/config/mcp.example.json new file mode 100644 index 0000000..775336c --- /dev/null +++ b/config/mcp.example.json @@ -0,0 +1,52 @@ +{ + "mcpServers": { + "filesystem": { + "command": "npx", + "args": [ + "-y", + "@modelcontextprotocol/server-filesystem", + "/Users/tuonome/Documents" + ] + }, + "remoto": { + "url": "https://esempio.it/mcp", + "headers": { + "Authorization": "Bearer TOKEN" + }, + "disabled": true + }, + "cua-driver": { + "command": "/Users/tuonome/.local/bin/cua-driver", + "args": [ + "mcp" + ], + "tools": [ + "list_apps", + "list_windows", + "get_window_state", + "verify_state", + "launch_app", + "bring_to_front", + "invoke_menu", + "click", + "double_click", + "right_click", + "type_text", + "press_key", + "hotkey", + "set_value", + "scroll", + "get_desktop_state", + "zoom", + "clipboard_read", + "clipboard_write", + "get_browser_state", + "browser_navigate", + "browser_click", + "browser_type" + ], + "description_limit": 800, + "disabled": true + } + } +} diff --git a/main.py b/main.py index dc06e3c..0fc268e 100644 --- a/main.py +++ b/main.py @@ -62,6 +62,9 @@ def main() -> None: confirm.bind(ui.show_confirm, ui.hide_confirm, ui.write_log) registry.ctx = {"say": assistant.say, "log": ui.write_log, "confirm": confirm.request, "player": ui} ui.get_plugins = registry.list_for_ui + from avatar.mcp_client import manager as mcp_manager + mcp_manager.start(registry, ui.write_log) + ui._app.aboutToQuit.connect(mcp_manager.stop) ui.get_plugin_settings = lambda: [] ui.request_say = assistant.say # Questi due sono attributi della finestra, non della facciata JarvisUI. diff --git a/memory/memory_manager.py b/memory/memory_manager.py index b2d18d8..5c78632 100644 --- a/memory/memory_manager.py +++ b/memory/memory_manager.py @@ -1,479 +1,513 @@ -import json -import re -from datetime import datetime -from threading import Lock -from pathlib import Path -import sys - - -def get_base_dir() -> Path: - if getattr(sys, "frozen", False): - return Path(sys.executable).parent - return Path(__file__).resolve().parent.parent - - -BASE_DIR = get_base_dir() -MEMORY_PATH = BASE_DIR / "memory" / "long_term.json" -_lock = Lock() -MAX_VALUE_LENGTH = 380 - -# ── Why there are two very different numbers here ──────────────────────────── -# -# There used to be one: MEMORY_MAX_CHARS = 2200, applied to the whole store. It -# was a *storage* limit, and it existed only because the entire memory was -# pasted into the system prompt on every connect — so growing the memory grew -# every single request. When it filled, _trim_to_limit() deleted the oldest -# entries and printed one line to a console nobody reads. A memory described as -# "deeply remembers projects, preferences and personal context" was in practice -# two pages long, and quietly forgot your sister's name after a few weeks. -# -# Storage and prompt budget are now separate concerns: -# -# MEMORY_MAX_CHARS — a runaway guard, not a feature limit. Nothing normal -# reaches it; a bug writing in a loop does. -# PROMPT_CORE_CHARS — what actually rides in the system prompt every session. -# Smaller than the old whole-memory dump, so sessions -# start *faster* than before, not slower. -# -# Everything above the core stays on disk and is fetched on demand by the -# recall_memory tool — see search_memory() and format_memory_for_prompt(). -MEMORY_MAX_CHARS = 200_000 -PROMPT_CORE_CHARS = 900 -PROMPT_INDEX_CHARS = 420 -# Most entries any one category may contribute to the core block, so a person -# with forty stored preferences still gets their sister into the prompt. -PROMPT_MAX_PER_CATEGORY = 6 - -def _empty_memory() -> dict: - return { - "identity": {}, - "preferences": {}, - "projects": {}, - "relationships": {}, - "wishes": {}, - "notes": {}, - } - -def load_memory() -> dict: - if not MEMORY_PATH.exists(): - return _empty_memory() - with _lock: - try: - data = json.loads(MEMORY_PATH.read_text(encoding="utf-8")) - if isinstance(data, dict): - base = _empty_memory() - for key in base: - if key not in data: - data[key] = {} - return data - return _empty_memory() - except Exception as e: - print(f"[Memory] ⚠️ Load error: {e}") - return _empty_memory() - -def _all_entries(memory: dict) -> list[tuple]: - entries = [] - for cat, items in memory.items(): - if not isinstance(items, dict): - continue - for key, entry in items.items(): - if isinstance(entry, dict) and "value" in entry: - entries.append((cat, key, entry)) - return entries - - -# Set by main.py so a trim can reach the activity log. Deleting something a -# person told you and mentioning it only on stdout is how a memory loses trust. -_trim_notifier = None - - -def set_trim_notifier(fn) -> None: - """Register a callable(str) that surfaces trims to the user.""" - global _trim_notifier - _trim_notifier = fn - - -def _trim_to_limit(memory: dict) -> dict: - if len(json.dumps(memory, ensure_ascii=False)) <= MEMORY_MAX_CHARS: - return memory - entries = _all_entries(memory) - entries.sort(key=lambda t: t[2].get("updated", "0000-00-00")) - dropped = [] - for cat, key, _ in entries: - if len(json.dumps(memory, ensure_ascii=False)) <= MEMORY_MAX_CHARS: - break - del memory[cat][key] - dropped.append(f"{cat}/{key}") - print(f"[Memory] 🗑️ Trimmed {cat}/{key}") - if dropped and _trim_notifier: - try: - _trim_notifier( - f"SYS: Memory full — forgot {len(dropped)} oldest entries " - f"({', '.join(dropped[:3])}{'…' if len(dropped) > 3 else ''})" - ) - except Exception: - pass - return memory - -def save_memory(memory: dict) -> None: - if not isinstance(memory, dict): - return - memory = _trim_to_limit(memory) - MEMORY_PATH.parent.mkdir(parents=True, exist_ok=True) - with _lock: - MEMORY_PATH.write_text( - json.dumps(memory, indent=2, ensure_ascii=False), - encoding="utf-8", - ) - - -def _truncate_value(val: str) -> str: - if isinstance(val, str) and len(val) > MAX_VALUE_LENGTH: - return val[:MAX_VALUE_LENGTH].rstrip() + "…" - return val - - -def _recursive_update(target: dict, updates: dict) -> bool: - changed = False - for key, value in updates.items(): - if value is None: - continue - if isinstance(value, str) and not value.strip(): - continue - if isinstance(value, dict) and "value" not in value: - if key not in target or not isinstance(target[key], dict): - target[key] = {} - changed = True - if _recursive_update(target[key], value): - changed = True - else: - new_val = _truncate_value(str(value["value"] if isinstance(value, dict) else value)) - entry = {"value": new_val, "updated": datetime.now().strftime("%Y-%m-%d")} - existing = target.get(key, {}) - if not isinstance(existing, dict) or existing.get("value") != new_val: - target[key] = entry - changed = True - return changed - - -def update_memory(memory_update: dict) -> dict: - if not isinstance(memory_update, dict) or not memory_update: - return load_memory() - memory = load_memory() - if _recursive_update(memory, memory_update): - save_memory(memory) - print(f"[Memory] 💾 Saved: {list(memory_update.keys())}") - return memory - -def _entry_value(entry) -> str: - """Accept both the {'value': ..., 'updated': ...} shape and a bare string, - because early versions of the store wrote plain strings.""" - if isinstance(entry, dict): - return str(entry.get("value", "") or "").strip() - return str(entry or "").strip() - - -def _pretty(key: str) -> str: - return key.replace("_", " ").strip() - - -# Identity is always in the prompt; these categories compete for the remaining -# budget by recency. -_CATEGORY_LABELS = { - "preferences": "Preferences", - "projects": "Active projects / goals", - "relationships": "People in their life", - "wishes": "Wishes / plans", - "notes": "Notes", -} - -_IDENTITY_FIELDS = ["name", "age", "birthday", "city", "job", - "language", "school", "nationality"] - - -def format_memory_for_prompt(memory: dict | None) -> str: - """Build the memory block that goes into the system prompt. - - This used to dump everything. It now sends three things: - - 1. IDENTITY - always, in full. It is small, and it is wrong for the - assistant to have to look up your name. - 2. RECENT - the most recently updated entries from every other - category, up to PROMPT_CORE_CHARS. Recency is the cheapest useful - relevance signal available without embeddings. - 3. AN INDEX - the *keys* of everything else, values omitted. - - Point 3 is what makes recall work at all. A model cannot decide to look - something up if it does not know the thing exists: with only points 1 and 2, - "who is Ayse?" would get "I don't know" while ayse_sister sat on disk - unread. The index costs a few hundred characters and turns recall from a - gamble into a lookup. - - Net effect on latency: this block is SMALLER than the old full dump, so - every session connects with fewer tokens. Occasionally the model spends one - extra round trip on recall_memory - covered by the acknowledgment it - already speaks before any slow step.""" - if not memory: - return "" - - core_lines: list[str] = [] - - # 1. Identity - always, in full - identity = memory.get("identity", {}) or {} - for field in _IDENTITY_FIELDS: - val = _entry_value(identity.get(field)) - if not val: - continue - if field == "language": - # Labelled as an observation, not a setting. A bare "Language: - # English" line written months ago reads like a standing order and - # was one of the reasons a Turkish question came back in English. - core_lines.append( - f"Has spoken to you in: {val} (an observation about the past — " - f"always answer in the language of their CURRENT message)") - else: - core_lines.append(f"{field.title()}: {val}") - for key, entry in identity.items(): - if key in _IDENTITY_FIELDS: - continue - val = _entry_value(entry) - if val: - core_lines.append(f"{_pretty(key).title()}: {val}") - - # 2. Everything else, most recently updated first - rest: list[tuple[str, str, str, str]] = [] # (updated, cat, key, value) - for cat in _CATEGORY_LABELS: - for key, entry in (memory.get(cat, {}) or {}).items(): - val = _entry_value(entry) - if not val: - continue - updated = (entry.get("updated", "") if isinstance(entry, dict) else "") or "0000-00-00" - rest.append((updated, cat, key, val)) - rest.sort(key=lambda t: t[0], reverse=True) - - used = sum(len(l) + 1 for l in core_lines) - shown: dict[str, list[str]] = {} - overflow: dict[str, list[str]] = {} - - # Recency decides order, but no single category may take the whole budget. - # Without the cap, someone with forty stored preferences gets a prompt that - # is forty preferences and not one person's name — the categories that - # matter most in conversation are also the ones that change least often, so - # pure recency systematically buries them. - per_cat_used: dict[str, int] = {} - for _updated, cat, key, val in rest: - line = f" - {_pretty(key).title()}: {val}" - if (per_cat_used.get(cat, 0) < PROMPT_MAX_PER_CATEGORY - and used + len(line) + 1 <= PROMPT_CORE_CHARS): - shown.setdefault(cat, []).append(line) - per_cat_used[cat] = per_cat_used.get(cat, 0) + 1 - used += len(line) + 1 - else: - overflow.setdefault(cat, []).append(_pretty(key)) - - # The index is a table of contents, so it is interleaved across categories - # rather than continuing in recency order. Sorted by recency it would list - # twenty-four preferences before the first relationship, and the one entry - # the index exists for — the old fact the model has no other way to know - # about — would fall off the end. - indexed: list[str] = [] - if overflow: - cats = [c for c in _CATEGORY_LABELS if overflow.get(c)] - cursor = {c: 0 for c in cats} - while cats: - for cat in list(cats): - i = cursor[cat] - if i >= len(overflow[cat]): - cats.remove(cat) - continue - indexed.append(overflow[cat][i]) - cursor[cat] = i + 1 - - for cat, label in _CATEGORY_LABELS.items(): - if shown.get(cat): - core_lines.append("") - core_lines.append(f"{label}:") - core_lines.extend(shown[cat]) - - if not core_lines and not indexed: - return "" - - out = [ - "[WHAT YOU KNOW ABOUT THIS PERSON — use naturally, never recite like a list]", - *core_lines, - ] - - # 3. The index of what is on disk but not in this prompt - if indexed: - budget, names = PROMPT_INDEX_CHARS, [] - for n in indexed: - if budget - len(n) - 2 < 0: - break - names.append(n) - budget -= len(n) + 2 - if names: - out.append("") - out.append( - "[ALSO REMEMBERED — values not shown here. Call recall_memory " - "with a keyword to read any of these before saying you do not know]" - ) - out.append(", ".join(names) - + (f" (+{len(indexed) - len(names)} more)" - if len(indexed) > len(names) else "")) - - return "\n".join(out) + "\n" - - -# ── Recall ──────────────────────────────────────────────────────────────────── - -def _score(query_words: list[str], cat: str, key: str, value: str) -> int: - """Cheap lexical relevance. No embeddings, no network, no model call - this - runs in well under a millisecond, which is the entire point: recall must - cost one model round trip, never two.""" - hay_key = _pretty(key).lower() - hay_val = value.lower() - score = 0 - for w in query_words: - if not w: - continue - if w == hay_key: - score += 10 - elif w in hay_key: - score += 6 - if w in hay_val: - score += 3 - if w in cat: - score += 1 - return score - - -def search_memory(query: str, limit: int = 8) -> str: - """Find stored facts matching `query`. Backs the recall_memory tool. - - An empty query is treated as "show me everything you know", capped - the - model asks that when the user says "what do you remember about me?".""" - memory = load_memory() - words = [w for w in re.split(r"[^\w]+", (query or "").lower()) if len(w) > 1] - - rows: list[tuple[int, str, str, str]] = [] - for cat, items in memory.items(): - if not isinstance(items, dict): - continue # skip 'sessions', which is a list - for key, entry in items.items(): - val = _entry_value(entry) - if not val: - continue - s = _score(words, cat, key, val) if words else 1 - if s > 0: - rows.append((s, cat, key, val)) - - if not rows: - return (f"Nothing stored about '{query}'." if query - else "I have not stored anything about this person yet.") - - rows.sort(key=lambda r: (-r[0], r[2])) - lines = [f"{cat}/{_pretty(key)}: {val}" for _s, cat, key, val in rows[:max(1, limit)]] - head = (f"Stored facts matching '{query}':" if query - else "Everything currently stored:") - more = (f"\n(+{len(rows) - len(lines)} more — search with a narrower keyword)" - if len(rows) > len(lines) else "") - return head + "\n" + "\n".join(lines) + more - - -def all_entries_for_ui() -> list[dict]: - """Flat list for the memory panel: what JARVIS knows, and when it learned it. - Sorted newest first so the panel opens on what changed most recently.""" - memory = load_memory() - rows = [] - for cat, items in memory.items(): - if not isinstance(items, dict): - continue - for key, entry in items.items(): - val = _entry_value(entry) - if not val: - continue - rows.append({ - "category": cat, - "key": key, - "value": val, - "updated": (entry.get("updated", "") if isinstance(entry, dict) else ""), - }) - rows.sort(key=lambda r: (r["updated"] or "0000-00-00"), reverse=True) - return rows - -def remember(key: str, value: str, category: str = "notes") -> str: - valid = {"identity", "preferences", "projects", "relationships", "wishes", "notes"} - if category not in valid: - category = "notes" - update_memory({category: {key: {"value": value}}}) - return f"Remembered: {category}/{key} = {value}" - - -def forget(key: str, category: str = "notes") -> str: - memory = load_memory() - cat = memory.get(category, {}) - if key in cat: - del cat[key] - memory[category] = cat - save_memory(memory) - return f"Forgotten: {category}/{key}" - return f"Not found: {category}/{key}" - - -forget_memory = forget - - -# ── Session memory ───────────────────────────────────────────────────────────── - -_SESSION_MAX = 3 # safety cap — in practice 0-1 entries after pop - - -def save_session_summary(summary: str, language: str = "") -> None: - """Append a 1-2 sentence session summary to long_term.json['sessions'].""" - summary = (summary or "").strip() - if not summary: - return - memory = load_memory() - sessions = memory.get("sessions", []) - if not isinstance(sessions, list): - sessions = [] - entry: dict = { - "date": datetime.now().strftime("%Y-%m-%d"), - "summary": summary[:280], - } - if language: - entry["language"] = language - sessions.append(entry) - memory["sessions"] = sessions[-_SESSION_MAX:] - with _lock: - MEMORY_PATH.parent.mkdir(parents=True, exist_ok=True) - MEMORY_PATH.write_text( - json.dumps(memory, indent=2, ensure_ascii=False), - encoding="utf-8", - ) - print(f"[Memory] 📝 Session saved ({entry['date']}): {summary[:60]}…") - - -def pop_last_session() -> dict | None: - """ - Return AND remove the most recent session entry. - Calling this consumes the entry so it is never repeated in future briefings. - """ - with _lock: - if not MEMORY_PATH.exists(): - return None - try: - memory = json.loads(MEMORY_PATH.read_text(encoding="utf-8")) - sessions = memory.get("sessions", []) - if not isinstance(sessions, list) or not sessions: - return None - entry = sessions.pop() # remove the last entry - memory["sessions"] = sessions - MEMORY_PATH.write_text( - json.dumps(memory, indent=2, ensure_ascii=False), - encoding="utf-8", - ) - return entry - except Exception as e: - print(f"[Memory] ⚠️ pop_last_session error: {e}") +import json +import re +from datetime import datetime +from threading import Lock +from pathlib import Path +import sys + + +def get_base_dir() -> Path: + if getattr(sys, "frozen", False): + return Path(sys.executable).parent + return Path(__file__).resolve().parent.parent + + +BASE_DIR = get_base_dir() +MEMORY_PATH = BASE_DIR / "memory" / "long_term.json" +_lock = Lock() +MAX_VALUE_LENGTH = 380 + +# ── Why there are two very different numbers here ──────────────────────────── +# +# There used to be one: MEMORY_MAX_CHARS = 2200, applied to the whole store. It +# was a *storage* limit, and it existed only because the entire memory was +# pasted into the system prompt on every connect — so growing the memory grew +# every single request. When it filled, _trim_to_limit() deleted the oldest +# entries and printed one line to a console nobody reads. A memory described as +# "deeply remembers projects, preferences and personal context" was in practice +# two pages long, and quietly forgot your sister's name after a few weeks. +# +# Storage and prompt budget are now separate concerns: +# +# MEMORY_MAX_CHARS — a runaway guard, not a feature limit. Nothing normal +# reaches it; a bug writing in a loop does. +# PROMPT_CORE_CHARS — what actually rides in the system prompt every session. +# Smaller than the old whole-memory dump, so sessions +# start *faster* than before, not slower. +# +# Everything above the core stays on disk and is fetched on demand by the +# recall_memory tool — see search_memory() and format_memory_for_prompt(). +MEMORY_MAX_CHARS = 200_000 +PROMPT_CORE_CHARS = 900 +PROMPT_INDEX_CHARS = 420 +# Most entries any one category may contribute to the core block, so a person +# with forty stored preferences still gets their sister into the prompt. +PROMPT_MAX_PER_CATEGORY = 6 + +def _empty_memory() -> dict: + return { + "identity": {}, + "preferences": {}, + "projects": {}, + "relationships": {}, + "wishes": {}, + "notes": {}, + } + +def load_memory() -> dict: + if not MEMORY_PATH.exists(): + return _empty_memory() + with _lock: + try: + data = json.loads(MEMORY_PATH.read_text(encoding="utf-8")) + if isinstance(data, dict): + base = _empty_memory() + for key in base: + if key not in data: + data[key] = {} + return data + return _empty_memory() + except Exception as e: + print(f"[Memory] ⚠️ Load error: {e}") + return _empty_memory() + +def _all_entries(memory: dict) -> list[tuple]: + entries = [] + for cat, items in memory.items(): + if not isinstance(items, dict): + continue + for key, entry in items.items(): + if isinstance(entry, dict) and "value" in entry: + entries.append((cat, key, entry)) + return entries + + +# Set by main.py so a trim can reach the activity log. Deleting something a +# person told you and mentioning it only on stdout is how a memory loses trust. +_trim_notifier = None + + +def set_trim_notifier(fn) -> None: + """Register a callable(str) that surfaces trims to the user.""" + global _trim_notifier + _trim_notifier = fn + + +def _trim_to_limit(memory: dict) -> dict: + if len(json.dumps(memory, ensure_ascii=False)) <= MEMORY_MAX_CHARS: + return memory + entries = _all_entries(memory) + entries.sort(key=lambda t: t[2].get("updated", "0000-00-00")) + dropped = [] + for cat, key, _ in entries: + if len(json.dumps(memory, ensure_ascii=False)) <= MEMORY_MAX_CHARS: + break + del memory[cat][key] + dropped.append(f"{cat}/{key}") + print(f"[Memory] 🗑️ Trimmed {cat}/{key}") + if dropped and _trim_notifier: + try: + _trim_notifier( + f"SYS: Memory full — forgot {len(dropped)} oldest entries " + f"({', '.join(dropped[:3])}{'…' if len(dropped) > 3 else ''})" + ) + except Exception: + pass + return memory + +def save_memory(memory: dict) -> None: + if not isinstance(memory, dict): + return + memory = _trim_to_limit(memory) + MEMORY_PATH.parent.mkdir(parents=True, exist_ok=True) + with _lock: + MEMORY_PATH.write_text( + json.dumps(memory, indent=2, ensure_ascii=False), + encoding="utf-8", + ) + try: + export_markdown(memory) + except Exception as err: + print(f"[Memory] esportazione Markdown fallita: {err}") + + +EXPORT_DIR = get_base_dir() / "data" / "memoria" +_EXPORT_LABELS = {"identity": "Identità", "preferences": "Preferenze", "projects": "Progetti", "relationships": "Persone", + "wishes": "Desideri e obiettivi", "notes": "Note"} + + +def export_markdown(memory: dict, folder: Path | None = None) -> Path: + """Specchia la memoria in una cartella di file Markdown (apribile come vault Obsidian): un file per categoria più un indice.""" + folder = folder or EXPORT_DIR + folder.mkdir(parents=True, exist_ok=True) + index = ["# Memoria di LuZa", "", f"Aggiornata: {datetime.now():%Y-%m-%d %H:%M}", ""] + for cat, label in _EXPORT_LABELS.items(): + entries = memory.get(cat) or {} + if not isinstance(entries, dict): + continue + lines = ["---", f"categoria: {cat}", f"voci: {len(entries)}", "---", "", f"# {label}", ""] + for key, entry in sorted(entries.items(), key=lambda kv: (kv[1].get("updated", "") if isinstance(kv[1], dict) else ""), reverse=True): + val = _entry_value(entry) + upd = entry.get("updated", "") if isinstance(entry, dict) else "" + lines.append(f"- **{_pretty(key)}**: {val}" + (f" _(aggiornato {upd})_" if upd else "")) + (folder / f"{label}.md").write_text("\n".join(lines) + "\n", encoding="utf-8") + index.append(f"- [[{label}]] — {len(entries)} voci") + sessions = memory.get("sessions") or [] + if isinstance(sessions, list) and sessions: + lines = ["# Sessioni recenti", ""] + [f"- {e.get('date', '')}: {e.get('summary', '')}" for e in sessions] + (folder / "Sessioni recenti.md").write_text("\n".join(lines) + "\n", encoding="utf-8") + index.append("- [[Sessioni recenti]]") + (folder / "Memoria.md").write_text("\n".join(index) + "\n", encoding="utf-8") + return folder + + +def _truncate_value(val: str) -> str: + if isinstance(val, str) and len(val) > MAX_VALUE_LENGTH: + return val[:MAX_VALUE_LENGTH].rstrip() + "…" + return val + + +def _recursive_update(target: dict, updates: dict) -> bool: + changed = False + for key, value in updates.items(): + if value is None: + continue + if isinstance(value, str) and not value.strip(): + continue + if isinstance(value, dict) and "value" not in value: + if key not in target or not isinstance(target[key], dict): + target[key] = {} + changed = True + if _recursive_update(target[key], value): + changed = True + else: + new_val = _truncate_value(str(value["value"] if isinstance(value, dict) else value)) + entry = {"value": new_val, "updated": datetime.now().strftime("%Y-%m-%d")} + existing = target.get(key, {}) + if not isinstance(existing, dict) or existing.get("value") != new_val: + target[key] = entry + changed = True + return changed + + +def update_memory(memory_update: dict) -> dict: + if not isinstance(memory_update, dict) or not memory_update: + return load_memory() + memory = load_memory() + if _recursive_update(memory, memory_update): + save_memory(memory) + print(f"[Memory] 💾 Saved: {list(memory_update.keys())}") + return memory + +def _entry_value(entry) -> str: + """Accept both the {'value': ..., 'updated': ...} shape and a bare string, + because early versions of the store wrote plain strings.""" + if isinstance(entry, dict): + return str(entry.get("value", "") or "").strip() + return str(entry or "").strip() + + +def _pretty(key: str) -> str: + return key.replace("_", " ").strip() + + +# Identity is always in the prompt; these categories compete for the remaining +# budget by recency. +_CATEGORY_LABELS = { + "preferences": "Preferences", + "projects": "Active projects / goals", + "relationships": "People in their life", + "wishes": "Wishes / plans", + "notes": "Notes", +} + +_IDENTITY_FIELDS = ["name", "age", "birthday", "city", "job", + "language", "school", "nationality"] + + +def format_memory_for_prompt(memory: dict | None) -> str: + """Build the memory block that goes into the system prompt. + + This used to dump everything. It now sends three things: + + 1. IDENTITY - always, in full. It is small, and it is wrong for the + assistant to have to look up your name. + 2. RECENT - the most recently updated entries from every other + category, up to PROMPT_CORE_CHARS. Recency is the cheapest useful + relevance signal available without embeddings. + 3. AN INDEX - the *keys* of everything else, values omitted. + + Point 3 is what makes recall work at all. A model cannot decide to look + something up if it does not know the thing exists: with only points 1 and 2, + "who is Ayse?" would get "I don't know" while ayse_sister sat on disk + unread. The index costs a few hundred characters and turns recall from a + gamble into a lookup. + + Net effect on latency: this block is SMALLER than the old full dump, so + every session connects with fewer tokens. Occasionally the model spends one + extra round trip on recall_memory - covered by the acknowledgment it + already speaks before any slow step.""" + if not memory: + return "" + + core_lines: list[str] = [] + + # 1. Identity - always, in full + identity = memory.get("identity", {}) or {} + for field in _IDENTITY_FIELDS: + val = _entry_value(identity.get(field)) + if not val: + continue + if field == "language": + # Labelled as an observation, not a setting. A bare "Language: + # English" line written months ago reads like a standing order and + # was one of the reasons a Turkish question came back in English. + core_lines.append( + f"Has spoken to you in: {val} (an observation about the past — " + f"always answer in the language of their CURRENT message)") + else: + core_lines.append(f"{field.title()}: {val}") + for key, entry in identity.items(): + if key in _IDENTITY_FIELDS: + continue + val = _entry_value(entry) + if val: + core_lines.append(f"{_pretty(key).title()}: {val}") + + # 2. Everything else, most recently updated first + rest: list[tuple[str, str, str, str]] = [] # (updated, cat, key, value) + for cat in _CATEGORY_LABELS: + for key, entry in (memory.get(cat, {}) or {}).items(): + val = _entry_value(entry) + if not val: + continue + updated = (entry.get("updated", "") if isinstance(entry, dict) else "") or "0000-00-00" + rest.append((updated, cat, key, val)) + rest.sort(key=lambda t: t[0], reverse=True) + + used = sum(len(l) + 1 for l in core_lines) + shown: dict[str, list[str]] = {} + overflow: dict[str, list[str]] = {} + + # Recency decides order, but no single category may take the whole budget. + # Without the cap, someone with forty stored preferences gets a prompt that + # is forty preferences and not one person's name — the categories that + # matter most in conversation are also the ones that change least often, so + # pure recency systematically buries them. + per_cat_used: dict[str, int] = {} + for _updated, cat, key, val in rest: + line = f" - {_pretty(key).title()}: {val}" + if (per_cat_used.get(cat, 0) < PROMPT_MAX_PER_CATEGORY + and used + len(line) + 1 <= PROMPT_CORE_CHARS): + shown.setdefault(cat, []).append(line) + per_cat_used[cat] = per_cat_used.get(cat, 0) + 1 + used += len(line) + 1 + else: + overflow.setdefault(cat, []).append(_pretty(key)) + + # The index is a table of contents, so it is interleaved across categories + # rather than continuing in recency order. Sorted by recency it would list + # twenty-four preferences before the first relationship, and the one entry + # the index exists for — the old fact the model has no other way to know + # about — would fall off the end. + indexed: list[str] = [] + if overflow: + cats = [c for c in _CATEGORY_LABELS if overflow.get(c)] + cursor = {c: 0 for c in cats} + while cats: + for cat in list(cats): + i = cursor[cat] + if i >= len(overflow[cat]): + cats.remove(cat) + continue + indexed.append(overflow[cat][i]) + cursor[cat] = i + 1 + + for cat, label in _CATEGORY_LABELS.items(): + if shown.get(cat): + core_lines.append("") + core_lines.append(f"{label}:") + core_lines.extend(shown[cat]) + + if not core_lines and not indexed: + return "" + + out = [ + "[WHAT YOU KNOW ABOUT THIS PERSON — use naturally, never recite like a list]", + *core_lines, + ] + + # 3. The index of what is on disk but not in this prompt + if indexed: + budget, names = PROMPT_INDEX_CHARS, [] + for n in indexed: + if budget - len(n) - 2 < 0: + break + names.append(n) + budget -= len(n) + 2 + if names: + out.append("") + out.append( + "[ALSO REMEMBERED — values not shown here. Call recall_memory " + "with a keyword to read any of these before saying you do not know]" + ) + out.append(", ".join(names) + + (f" (+{len(indexed) - len(names)} more)" + if len(indexed) > len(names) else "")) + + return "\n".join(out) + "\n" + + +# ── Recall ──────────────────────────────────────────────────────────────────── + +def _score(query_words: list[str], cat: str, key: str, value: str) -> int: + """Cheap lexical relevance. No embeddings, no network, no model call - this + runs in well under a millisecond, which is the entire point: recall must + cost one model round trip, never two.""" + hay_key = _pretty(key).lower() + hay_val = value.lower() + score = 0 + for w in query_words: + if not w: + continue + if w == hay_key: + score += 10 + elif w in hay_key: + score += 6 + if w in hay_val: + score += 3 + if w in cat: + score += 1 + return score + + +def search_memory(query: str, limit: int = 8) -> str: + """Find stored facts matching `query`. Backs the recall_memory tool. + + An empty query is treated as "show me everything you know", capped - the + model asks that when the user says "what do you remember about me?".""" + memory = load_memory() + words = [w for w in re.split(r"[^\w]+", (query or "").lower()) if len(w) > 1] + + rows: list[tuple[int, str, str, str]] = [] + for cat, items in memory.items(): + if not isinstance(items, dict): + continue # skip 'sessions', which is a list + for key, entry in items.items(): + val = _entry_value(entry) + if not val: + continue + s = _score(words, cat, key, val) if words else 1 + if s > 0: + rows.append((s, cat, key, val)) + + if not rows: + return (f"Nothing stored about '{query}'." if query + else "I have not stored anything about this person yet.") + + rows.sort(key=lambda r: (-r[0], r[2])) + lines = [f"{cat}/{_pretty(key)}: {val}" for _s, cat, key, val in rows[:max(1, limit)]] + head = (f"Stored facts matching '{query}':" if query + else "Everything currently stored:") + more = (f"\n(+{len(rows) - len(lines)} more — search with a narrower keyword)" + if len(rows) > len(lines) else "") + return head + "\n" + "\n".join(lines) + more + + +def all_entries_for_ui() -> list[dict]: + """Flat list for the memory panel: what JARVIS knows, and when it learned it. + Sorted newest first so the panel opens on what changed most recently.""" + memory = load_memory() + rows = [] + for cat, items in memory.items(): + if not isinstance(items, dict): + continue + for key, entry in items.items(): + val = _entry_value(entry) + if not val: + continue + rows.append({ + "category": cat, + "key": key, + "value": val, + "updated": (entry.get("updated", "") if isinstance(entry, dict) else ""), + }) + rows.sort(key=lambda r: (r["updated"] or "0000-00-00"), reverse=True) + return rows + +def remember(key: str, value: str, category: str = "notes") -> str: + valid = {"identity", "preferences", "projects", "relationships", "wishes", "notes"} + if category not in valid: + category = "notes" + update_memory({category: {key: {"value": value}}}) + return f"Remembered: {category}/{key} = {value}" + + +def forget(key: str, category: str = "notes") -> str: + memory = load_memory() + cat = memory.get(category, {}) + if key in cat: + del cat[key] + memory[category] = cat + save_memory(memory) + return f"Forgotten: {category}/{key}" + return f"Not found: {category}/{key}" + + +forget_memory = forget + + +# ── Session memory ───────────────────────────────────────────────────────────── + +_SESSION_MAX = 3 # safety cap — in practice 0-1 entries after pop + + +def save_session_summary(summary: str, language: str = "") -> None: + """Append a 1-2 sentence session summary to long_term.json['sessions'].""" + summary = (summary or "").strip() + if not summary: + return + memory = load_memory() + sessions = memory.get("sessions", []) + if not isinstance(sessions, list): + sessions = [] + entry: dict = { + "date": datetime.now().strftime("%Y-%m-%d"), + "summary": summary[:280], + } + if language: + entry["language"] = language + sessions.append(entry) + memory["sessions"] = sessions[-_SESSION_MAX:] + with _lock: + MEMORY_PATH.parent.mkdir(parents=True, exist_ok=True) + MEMORY_PATH.write_text( + json.dumps(memory, indent=2, ensure_ascii=False), + encoding="utf-8", + ) + print(f"[Memory] 📝 Session saved ({entry['date']}): {summary[:60]}…") + + +def pop_last_session() -> dict | None: + """ + Return AND remove the most recent session entry. + Calling this consumes the entry so it is never repeated in future briefings. + """ + with _lock: + if not MEMORY_PATH.exists(): + return None + try: + memory = json.loads(MEMORY_PATH.read_text(encoding="utf-8")) + sessions = memory.get("sessions", []) + if not isinstance(sessions, list) or not sessions: + return None + entry = sessions.pop() # remove the last entry + memory["sessions"] = sessions + MEMORY_PATH.write_text( + json.dumps(memory, indent=2, ensure_ascii=False), + encoding="utf-8", + ) + return entry + except Exception as e: + print(f"[Memory] ⚠️ pop_last_session error: {e}") return None \ No newline at end of file diff --git a/plugins/immagini.py b/plugins/immagini.py new file mode 100644 index 0000000..6e16c6e --- /dev/null +++ b/plugins/immagini.py @@ -0,0 +1,98 @@ +"""Generazione di immagini in locale con mflux (MLX su Apple Silicon): Z-Image Turbo, FLUX e altri. + +Gira in un ambiente separato (.venv-mflux) come processo esterno, così la memoria viene liberata a fine generazione. +Il modello si sceglie nelle impostazioni (scheda Strumenti); il primo uso scarica i pesi da Hugging Face. +""" +from __future__ import annotations + +import re +import subprocess +import sys +import time +from pathlib import Path + +sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) +from avatar.settings import BASE_DIR, DATA_DIR, Settings # noqa: E402 + +VENV = BASE_DIR / ".venv-mflux" / "bin" +OUT_DIR = DATA_DIR / "immagini" +# comando mflux, passi e guida consigliati per ciascuna famiglia +MODELLI = { + "z-image-turbo": {"cmd": "mflux-generate-z-image-turbo", "steps": 8, "guidance": None, "default_repo": "mflux-community/z-image-turbo-mflux-q4"}, + "schnell": {"cmd": "mflux-generate", "steps": 4, "guidance": None, "default_repo": "dhairyashil/FLUX.1-schnell-mflux-4bit", "base": "schnell"}, + "dev": {"cmd": "mflux-generate", "steps": 20, "guidance": 3.5, "default_repo": "", "base": "dev"}, + "qwen": {"cmd": "mflux-generate-qwen", "steps": 20, "guidance": 4.0, "default_repo": ""}, +} +FORMATI = {"quadrata": (1024, 1024), "orizzontale": (1280, 768), "verticale": (768, 1280), "schermo": (1536, 864), "piccola": (768, 768)} +_last: dict = {"path": ""} + + +def _cfg() -> tuple[str, str]: + s = Settings() + fam = str(s.get("immagini_famiglia") or "z-image-turbo") + repo = str(s.get("immagini_modello") or "").strip() or MODELLI.get(fam, MODELLI["z-image-turbo"])["default_repo"] + return fam, repo + + +def genera(params: dict, ctx: dict) -> str: + prompt = str(params.get("descrizione", "")).strip() + if not prompt: + return "Errore: serve la descrizione dell'immagine (in inglese, dettagliata)." + if not (VENV / "mflux-generate").exists(): + return "Generazione immagini non disponibile: manca l'ambiente .venv-mflux (uv venv .venv-mflux && uv pip install mflux)." + fam, repo = _cfg() + m = MODELLI.get(fam, MODELLI["z-image-turbo"]) + w, h = FORMATI.get(str(params.get("formato", "quadrata")).lower(), FORMATI["quadrata"]) + steps = int(params.get("passi") or m["steps"]) + OUT_DIR.mkdir(parents=True, exist_ok=True) + stem = re.sub(r"[^a-z0-9]+", "-", prompt.lower())[:40].strip("-") or "immagine" + out = OUT_DIR / f"{time.strftime('%Y%m%d-%H%M%S')}-{stem}.png" + cmd = [str(VENV / m["cmd"]), "--prompt", prompt, "--steps", str(steps), "--width", str(w), "--height", str(h), "--output", str(out)] + if repo: + cmd += ["--model", repo] + if m.get("base") and "/" in repo: + cmd += ["--base-model", m["base"]] + if m["guidance"] is not None: + cmd += ["--guidance", str(params.get("guida") or m["guidance"])] + if params.get("seme") is not None: + cmd += ["--seed", str(int(params["seme"]))] + if params.get("immagine_base"): + cmd += ["--image-path", str(params["immagine_base"]), "--image-strength", str(params.get("forza") or 0.5)] + if ctx.get("log"): + ctx["log"](f"[immagini] {fam} {w}x{h} {steps} passi: {prompt[:80]}") + t = time.time() + try: + proc = subprocess.run(cmd, capture_output=True, text=True, timeout=1800) + except subprocess.TimeoutExpired: + return "Generazione interrotta: troppo lenta (oltre 30 minuti)." + if proc.returncode != 0 or not out.exists(): + err = (proc.stderr or proc.stdout).strip().splitlines() + tail = " | ".join(err[-3:]) if err else f"codice {proc.returncode}" + if "not found" in tail.lower() or "401" in tail or "403" in tail: + tail += " (controlla il nome del modello nelle impostazioni, o accetta la licenza su Hugging Face)" + return f"Generazione fallita: {tail[:400]}" + _last["path"] = str(out) + secs = time.time() - t + if params.get("mostra", True) not in (False, "false", "no"): + subprocess.Popen(["open", str(out)]) + return f"Immagine generata in {secs:.0f} s e salvata in {out} ({w}x{h}, {fam}).\nIMMAGINE: {out}" + + +def ultima(params: dict, ctx: dict) -> str: + files = sorted(OUT_DIR.glob("*.png"), key=lambda p: p.stat().st_mtime, reverse=True) if OUT_DIR.exists() else [] + if not files: + return "Nessuna immagine generata finora." + n = max(1, min(int(params.get("numero") or 5), 20)) + if str(params.get("apri", "")).lower() in ("true", "1", "sì", "si"): + subprocess.Popen(["open", str(files[0])]) + return "Ultime immagini:\n" + "\n".join(f"- {p.name}" for p in files[:n]) + f"\nIMMAGINE: {files[0]}" + + +TOOLS = [ + {"name": "immagine_genera", "description": "Genera un'immagine in locale (MLX) da una descrizione. Scrivi la descrizione in INGLESE, dettagliata (soggetto, ambiente, luce, stile, inquadratura), anche se l'utente parla italiano. formato: quadrata (default), orizzontale, verticale, schermo, piccola. L'immagine viene salvata e aperta in Anteprima; richiede 10-60 secondi, il primo uso scarica il modello. Con immagine_base (percorso) e forza (0-1) parte da un'immagine esistente.", + "parameters": {"type": "object", "properties": {"descrizione": {"type": "string"}, "formato": {"type": "string", "enum": ["quadrata", "orizzontale", "verticale", "schermo", "piccola"]}, + "passi": {"type": "integer"}, "seme": {"type": "integer"}, "immagine_base": {"type": "string"}, "forza": {"type": "number"}, "mostra": {"type": "boolean"}}, + "required": ["descrizione"]}, "run": genera}, + {"name": "immagine_ultime", "description": "Elenca le ultime immagini generate (cartella data/immagini); con apri=true apre la più recente.", + "parameters": {"type": "object", "properties": {"numero": {"type": "integer"}, "apri": {"type": "boolean"}}}, "run": ultima}, +]