AvatarPy: assistente personale con volto 3D, voce, memoria e plugin

Riscrittura in Python dell'assistente Avatar con interfaccia HUD (derivata da Mark LIV, CC BY-NC 4.0, vedi NOTICE.md).
Tre motori (Claude API, server locale OpenAI-compatibile, Claude Code), voce Kokoro/macOS, Whisper MLX,
avatar 3D con sincronizzazione labiale, memoria per categorie, allegati con OCR, monitor con avvisi,
plugin per Calendario, Mail, Promemoria, Note, Musica, app, Mac, timer, meteo, contatti, Messaggi,
file, Comandi Rapidi, browser, Telegram, WhatsApp (archivio e tempo reale).

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
lucianoandClaude Fable 5.1 committed 2026-09-23 16:21:39 +02:00
commit ff79832c30
76 files changed
+82468

No files matched your search

View File
Whitespace-only changes.
+503
View File
@@ -0,0 +1,503 @@
"""Orchestratore: microfono → trascrizione → motore → voce, con stati e log sulla UI."""
from __future__ import annotations
import json
import queue
import threading
import time
import numpy as np
from memory import config_manager as cm
from .audio import Microphone, Player
from .engines.anthropic_engine import AnthropicEngine
from .engines.claude_code import ClaudeCodeEngine
from .engines.openai_compat import OpenAICompatEngine
from .settings import CONFIG_DIR, Settings
from .stt import Transcriber
from .tts import SentenceSplitter, clean_for_speech, make_voice
END = object()
STATUS_LABELS = {
"thinking": "THINKING", "searching": "PROCESSING", "memory": "PROCESSING",
"working": "PROCESSING", "responding": "THINKING",
}
# Frasi che Whisper "inventa" su silenzio o rumore di fondo (sottotitoli visti in addestramento).
_HALLUCINATIONS = ("sottotitoli", "a cura di", "grazie per aver guardato", "iscriviti al canale",
"www.", "amara.org", "subtitles", "thank you for watching")
def _is_hallucination(text: str) -> bool:
t = text.lower()
return any(h in t for h in _HALLUCINATIONS)
def _identity() -> tuple[str, str]:
try:
d = json.loads((CONFIG_DIR / "api_keys.json").read_text(encoding="utf-8"))
return (d.get("assistant_name") or "Ava").strip(), (d.get("user_name") or "").strip()
except Exception:
return "Ava", ""
class Assistant:
def __init__(self, ui, settings: Settings) -> None:
self.ui = ui
self.settings = settings
self.name, self.user_name = _identity()
self._abort = threading.Event()
self._turn_lock = threading.Lock()
self._busy = False
self._speaking = False
self._tail_until = 0.0
self._turn = 0
self._ptt_enabled = False
self._ptt_held = False
self._wake_enabled = False
self._awake = True
self._wake = None
self._last_speech = time.monotonic()
self._engine = None
self._engine_key = None
self._last_attachment: tuple[str, float] | None = None
self._speech_q: queue.Queue = queue.Queue()
self._audio_q: queue.Queue = queue.Queue(maxsize=3)
self.player = Player(self.ui.push_visemes, self.ui.set_audio_level)
self.mic = Microphone(self._on_utterance, self._mic_level, self._can_listen,
float(settings.get("vad_threshold", 0.08)), on_frame=self._on_mic_frame)
self.stt = Transcriber(str(settings.get("stt_model")), on_status=self._note)
self.voice = make_voice(settings, on_status=self._note)
# ── Avvio ────────────────────────────────────────────────────────────
def start(self) -> None:
threading.Thread(target=self._synth_loop, name="tts-synth", daemon=True).start()
threading.Thread(target=self._play_loop, name="tts-play", daemon=True).start()
threading.Thread(target=self._warmup, name="warmup", daemon=True).start()
threading.Thread(target=self._scheduler_loop, name="scheduler", daemon=True).start()
def _start_whatsapp_bridge(self) -> None:
try:
from avatar import whatsapp_bridge as wb
if self.settings.get("whatsapp_live") or wb.is_linked():
self.ui.write_log(f"SYS: WhatsApp — {wb.start(self.user_name)}")
self._wa_last_check = time.strftime("%Y-%m-%d %H:%M:%S")
except Exception as err:
self.ui.write_log(f"ERR: WhatsApp — {err}")
def _announce_whatsapp(self) -> None:
"""Legge i messaggi WhatsApp arrivati dall'ultimo controllo e, se richiesto, li annuncia."""
import sqlite3
from avatar.settings import DATA_DIR
db = DATA_DIR / "whatsapp" / "index.sqlite"
if not db.exists():
return
since = getattr(self, "_wa_last_check", None) or time.strftime("%Y-%m-%d %H:%M:%S")
now = time.strftime("%Y-%m-%d %H:%M:%S")
con = sqlite3.connect(f"file:{db}?mode=ro", uri=True)
try:
rows = con.execute("SELECT ts, chat, sender, text FROM messages WHERE source = 'live' AND from_me = 0 AND ts > ? ORDER BY ts LIMIT 5", (since,)).fetchall()
finally:
con.close()
self._wa_last_check = now
for ts, chat, sender, text in rows:
who = chat if sender == chat else f"{sender} nel gruppo {chat}"
self.ui.write_log(f"WhatsApp: {who}: {text[:200]}")
if self.settings.get("whatsapp_annuncia") and not self._busy and not self._speaking:
self.say(f"Messaggio WhatsApp da {who}: {text[:220]}")
def _scheduler_loop(self) -> None:
"""Fa scattare timer, sveglie e attività programmate dal plugin timer."""
import importlib
while True:
time.sleep(5)
try:
mod = importlib.import_module("avatar_plugins.timer")
except Exception:
try:
from avatar.plugins import registry
registry.load()
mod = importlib.import_module("avatar_plugins.timer")
except Exception:
continue
try:
self._announce_whatsapp()
except Exception as err:
print(f"[whatsapp] {err}")
try:
for item in mod.due():
testo = str(item.get("testo", ""))
if item.get("comando"):
self.ui.write_log(f"SYS: Attività programmata: {testo}")
self.handle_text(testo)
else:
self.ui.write_log(f"SYS: Timer: {testo}")
self.say(f"Promemoria: {testo}")
except Exception as err:
print(f"[scheduler] {err}")
def _warmup(self) -> None:
self.ui.set_state("PROCESSING")
self.ui.write_log(f"SYS: {self.name} si sta avviando: carico voce e riconoscimento vocale…")
try:
self.voice.load()
except Exception as err:
self.ui.write_log(f"ERR: Voce non disponibile — {err}")
try:
self.stt.load()
except Exception as err:
self.ui.write_log(f"ERR: Riconoscimento vocale non disponibile — {err}")
self.reopen_audio()
self._ensure_engine()
self._start_whatsapp_bridge()
try:
from avatar import monitor as _mon
_mon.monitor = _mon.Monitor(self.ui, self.say, self.settings)
_mon.monitor.start()
except Exception as err:
self.ui.write_log(f"ERR: Monitor — {err}")
self.ui.set_state("LISTENING")
self.ui.write_log(f"SYS: {self.name} è pronta. Parla o scrivi.")
def reopen_audio(self) -> None:
try:
self.player.open(cm.get_output_device())
except Exception as err:
self.ui.write_log(f"ERR: Uscita audio — {err}")
try:
self.mic.start(cm.get_input_device())
except Exception as err:
self.ui.write_log(f"ERR: Microfono — {err}")
def _note(self, msg: str | None) -> None:
if msg:
self.ui.write_log(f"SYS: {msg}")
# ── Motore ───────────────────────────────────────────────────────────
def _ensure_engine(self):
s = self.settings
provider = s.get("provider")
self.name, self.user_name = _identity()
key = (provider, s.get("effort"), s.get("local_base_url"), s.get("local_model"), s.get("search_api_key"),
s.get("claudecode_model"), s.get("claudecode_access"), s.get("claudecode_config_dir"),
s.get("claudecode_path"), self.name, self.user_name, s.get_secret("anthropic_api_key")[-6:],
s.get_secret("local_api_key")[-4:])
if self._engine and self._engine_key == key:
return self._engine
self._engine, self._engine_key = None, key
if provider == "local":
if s.get("local_base_url") and s.get("local_model"):
self._engine = OpenAICompatEngine(s.get("local_base_url"), s.get("local_model"), s.get_secret("local_api_key"),
s.get("search_api_key"), self.name, self.user_name)
elif provider == "claudecode":
self._engine = ClaudeCodeEngine(s.get("claudecode_model") or "sonnet", s.get("claudecode_access") or "chat",
s.get("claudecode_path") or "", s.get("claudecode_config_dir") or "",
self.name, self.user_name, s.get("effort"))
else:
api_key = s.get_secret("anthropic_api_key")
if api_key:
self._engine = AnthropicEngine(api_key, self.name, self.user_name, s.get("effort"))
return self._engine
def reconfigure(self) -> None:
"""Dopo un salvataggio delle impostazioni."""
self.interrupt(silent=True)
self._ensure_engine()
self.reload_voice()
model = str(self.settings.get("stt_model"))
if model != self.stt.model:
self.stt = Transcriber(model, on_status=self._note)
threading.Thread(target=self.stt.load, daemon=True).start()
self.mic.threshold = float(self.settings.get("vad_threshold", 0.08))
try:
view = getattr(self.ui._win, "avatar3d", None)
if view is not None:
view.set_model(str(self.settings.get("avatar_model") or ""))
except Exception as err:
self.ui.write_log(f"ERR: Avatar 3D — {err}")
self.ui.write_log(f"SYS: Impostazioni applicate — motore {self.settings.get('provider')}.")
def reload_voice(self, from_customise: bool = False) -> None:
"""Ricarica la voce. `from_customise`: la scelta arriva dal pannello Customise
(Sara / Nicola / Sistema) e va copiata nelle impostazioni; altrimenti comandano
le impostazioni di "Motore e Voce" e il pannello viene allineato."""
mapping = {"Sara": ("kokoro", "if_sara"), "Nicola": ("kokoro", "im_nicola"), "Sistema": ("system", None)}
if from_customise:
chosen = (cm.get_voice() or "").strip()
if chosen in mapping:
engine, voice = mapping[chosen]
self.settings.set("tts_engine", engine)
if voice:
self.settings.set("kokoro_voice", voice)
self.settings.save()
else:
label = "Sistema" if self.settings.get("tts_engine") == "system" else \
{"if_sara": "Sara", "im_nicola": "Nicola"}.get(self.settings.get("kokoro_voice"), "Sara")
try:
if cm.get_voice() != label:
cm.save_voice(label)
except Exception:
pass
new = make_voice(self.settings, on_status=self._note)
if type(new) is type(self.voice) and getattr(new, "voice", None) == getattr(self.voice, "voice", None):
return
self.voice = new
self.ui.write_log(f"SYS: Voce: {getattr(new, 'voice', '') or 'sistema'} ({'Kokoro' if self.settings.get('tts_engine') == 'kokoro' else 'macOS'}).")
threading.Thread(target=self._safe_load_voice, daemon=True).start()
def _safe_load_voice(self) -> None:
try:
self.voice.load()
except Exception as err:
self.ui.write_log(f"ERR: Voce non disponibile — {err}")
def reset_conversation(self) -> None:
self.interrupt(silent=True)
if self._engine:
self._engine.reset()
self.ui.write_log("SYS: Nuova conversazione.")
# ── Ascolto ──────────────────────────────────────────────────────────
def _can_listen(self) -> bool:
if self.ui.muted or self._busy or self._speaking or not self._awake:
return False
if time.monotonic() < self._tail_until:
return False
if self._ptt_enabled and not self._ptt_held:
return False
return True
def _mic_level(self, level: float) -> None:
self.ui.set_audio_level(level)
def _on_mic_frame(self, frame: np.ndarray) -> None:
if self._wake is not None and self._wake_enabled and not self._awake:
try:
self._wake.feed(frame)
except Exception:
pass
if self._wake_enabled and self._awake and not self._busy and not self._speaking:
if time.monotonic() - self._last_speech > 120:
self._sleep("silenzio")
def _on_utterance(self, audio: np.ndarray) -> None:
self._last_speech = time.monotonic()
self.ui.set_state("PROCESSING")
try:
text = self.stt.transcribe(audio)
except Exception as err:
self.ui.write_log(f"ERR: Trascrizione — {err}")
self.ui.set_state("LISTENING")
return
text = text.strip()
if len(text) < 2 or _is_hallucination(text):
self.ui.set_state("LISTENING")
return
self.ui.write_log(f"You: {text}")
self.handle_text(text)
# ── Turno di conversazione ───────────────────────────────────────────
def handle_text(self, text: str) -> None:
"""Chiamabile da qualunque thread (UI, microfono, quiz…)."""
text = (text or "").strip()
if not text:
return
if self._busy or self._speaking:
self.interrupt(silent=True)
threading.Thread(target=self._run_turn, args=(text,), daemon=True).start()
def say(self, text: str) -> None:
"""Pronuncia un testo senza passare dal motore."""
for s in SentenceSplitter().push(text + "\n") + []:
self._speech_q.put((self._turn, s))
self._speech_q.put((self._turn, END))
def _run_turn(self, text: str) -> None:
with self._turn_lock:
engine = self._ensure_engine()
if engine is None:
self.ui.write_log("ERR: Motore non configurato: apri Motore & Voce e completa le impostazioni.")
self.ui.set_state("LISTENING")
return
self._abort.clear()
self._busy = True
self._turn += 1
turn = self._turn
self.ui.set_state("THINKING")
splitter = SentenceSplitter()
def emit(ev: dict) -> None:
if turn != self._turn:
return
t = ev.get("type")
if t == "status":
if not self._speaking:
self.ui.set_state(STATUS_LABELS.get(ev["status"], "THINKING"))
if ev["status"] == "searching":
self.ui.write_log("SYS: Cerco sul web…")
elif ev["status"] == "working":
self.ui.write_log(f"SYS: Lavoro sul Mac ({ev.get('detail', '')})…")
elif t == "text":
for s in splitter.push(ev["delta"]):
self._speech_q.put((turn, s))
elif t == "memory_saved":
self.ui.write_log(f"SYS: Memoria salvata — {ev['text']}")
elif t == "memory_removed":
self.ui.write_log("SYS: Memoria cancellata.")
elif t == "done":
for s in splitter.flush():
self._speech_q.put((turn, s))
self.ui.write_log(f"{self.name}: {ev['text']}")
if ev.get("sources"):
self.ui.show_content("Fonti", "\n".join(f"• {s['title']}\n {s['url']}" for s in ev["sources"]))
elif t == "error":
if not ev.get("aborted"):
self.ui.write_log(f"ERR: {ev['message']}")
self._speech_q.put((turn, clean_for_speech("Scusa, c'è stato un problema: " + ev["message"])[:300]))
# Allegato dalla zona "File upload" (una volta per file)
engine_text, image, attach_path = text, None, ""
try:
path = self.ui.current_file
if path:
from pathlib import Path as _P
p = _P(path)
key = (str(p), p.stat().st_mtime)
if p.exists() and key != self._last_attachment:
self._last_attachment = key
from avatar.attachments import describe
self.ui.write_log(f"FILE: {p.name}")
ctx_text, image = describe(p)
engine_text = f"{text}\n\n{ctx_text}"
attach_path = str(p)
except Exception as err:
self.ui.write_log(f"ERR: Allegato — {err}")
try:
engine.send(engine_text, emit, self._abort, image=image, attach_path=attach_path)
finally:
self._busy = False
self._speech_q.put((turn, END))
# ── Voce: sintesi in anticipo e riproduzione ─────────────────────────
def _synth_loop(self) -> None:
while True:
turn, item = self._speech_q.get()
if turn != self._turn:
continue
if item is END:
self._audio_q.put((turn, END, None))
continue
try:
audio = self.voice.synthesize(item)
except Exception as err:
self.ui.write_log(f"ERR: Sintesi vocale — {err}")
continue
if turn == self._turn:
self._audio_q.put((turn, item, audio))
def _play_loop(self) -> None:
while True:
turn, text, audio = self._audio_q.get()
if turn != self._turn:
continue
if text is END:
self.player.drain()
self._speaking = False
self._tail_until = time.monotonic() + 0.5
if not self._busy and turn == self._turn:
self.ui.set_state("LISTENING" if self._awake else "SLEEPING")
continue
if audio is None or len(audio) == 0:
continue
self._speaking = True
self.ui.set_state("SPEAKING")
self.player.play(audio, text, self._abort)
self._last_speech = time.monotonic()
def interrupt(self, silent: bool = False) -> None:
self._abort.set()
self._turn += 1
for q in (self._speech_q, self._audio_q):
try:
while True:
q.get_nowait()
except queue.Empty:
pass
self.player.abort()
self._speaking = False
self._busy = False
self._tail_until = time.monotonic() + 0.4
if not silent:
self.ui.write_log("SYS: Interrotta. Ti ascolto.")
self.ui.set_state("LISTENING" if self._awake else "SLEEPING")
# ── Push-to-talk e wake word ─────────────────────────────────────────
def set_push_to_talk(self, enabled: bool) -> str:
self._ptt_enabled = bool(enabled)
self._ptt_held = False
return "window"
def ptt_hold(self, held: bool) -> None:
self._ptt_held = bool(held)
if held and not self._awake:
self._wake_up("tasto")
self.ui.set_state("LISTENING" if held else ("LISTENING" if not self._ptt_enabled else "SLEEPING"))
def wake_get_state(self) -> dict:
ready = False
try:
from core.wake_word import is_ready
ready = is_ready()
except Exception:
pass
return {"enabled": self._wake_enabled, "awake": self._awake, "ready": ready}
def on_wake_toggle(self, enable: bool) -> str:
if enable:
try:
from core.wake_word import WakeWordDetector, is_ready
if not is_ready():
return "Il modello della parola di attivazione non è installato."
self._wake = WakeWordDetector(on_detect=lambda: self._wake_up("parola di attivazione"))
if not self._wake.start():
self._wake = None
return "Impossibile avviare la parola di attivazione."
except Exception as err:
self._wake = None
return f"Parola di attivazione non disponibile: {err}"
self._wake_enabled = True
self._last_speech = time.monotonic()
self.ui.write_log("SYS: Parola di attivazione attiva: di' «Hey Jarvis» per svegliarmi.")
return "on"
self._wake_enabled = False
if self._wake:
try:
self._wake.stop()
except Exception:
pass
self._wake = None
self._wake_up("disattivata")
return "off"
def on_wake_manual(self) -> None:
if self._awake:
self._sleep("manuale")
else:
self._wake_up("manuale")
def _sleep(self, reason: str) -> None:
if not self._awake:
return
self._awake = False
self.ui.set_state("SLEEPING")
self.ui.write_log(f"SYS: In pausa ({reason}).")
def _wake_up(self, reason: str) -> None:
self._last_speech = time.monotonic()
if self._awake:
return
self._awake = True
self.ui.set_state("LISTENING")
self.ui.write_log(f"SYS: Sveglia ({reason}).")
+74
View File
@@ -0,0 +1,74 @@
"""Allegati dalla zona "File upload": testo dai documenti, OCR e immagine dalle foto."""
from __future__ import annotations
import io
import re
from pathlib import Path
IMAGE_EXT = {".png", ".jpg", ".jpeg", ".gif", ".webp", ".heic", ".tiff", ".bmp"}
MAX_TEXT = 8000
def ocr(path: Path) -> str:
"""Testo riconosciuto nell'immagine con il framework Vision di macOS."""
import Quartz
import Vision
from Foundation import NSURL
src = Quartz.CGImageSourceCreateWithURL(NSURL.fileURLWithPath_(str(path)), None)
if src is None:
return ""
img = Quartz.CGImageSourceCreateImageAtIndex(src, 0, None)
if img is None:
return ""
req = Vision.VNRecognizeTextRequest.alloc().init()
req.setRecognitionLevel_(Vision.VNRequestTextRecognitionLevelAccurate)
req.setRecognitionLanguages_(["it-IT", "en-US"])
req.setUsesLanguageCorrection_(True)
handler = Vision.VNImageRequestHandler.alloc().initWithCGImage_options_(img, None)
ok, err = handler.performRequests_error_([req], None)
if not ok:
return ""
lines = []
for obs in req.results() or []:
cands = obs.topCandidates_(1)
if cands:
lines.append(str(cands[0].string()))
return "\n".join(lines)
def image_payload(path: Path, max_side: int = 1568) -> tuple[str, bytes]:
"""Immagine ridotta per i modelli con visione: (media_type, bytes)."""
from PIL import Image
im = Image.open(path)
im = im.convert("RGB")
im.thumbnail((max_side, max_side))
buf = io.BytesIO()
im.save(buf, format="JPEG", quality=85)
return "image/jpeg", buf.getvalue()
def describe(path: Path) -> tuple[str, tuple[str, bytes] | None]:
"""Restituisce (contesto testuale da aggiungere al messaggio, immagine per la visione o None)."""
ext = path.suffix.lower()
if ext in IMAGE_EXT:
text = ""
try:
text = ocr(path)
except Exception as err:
text = f"(OCR non riuscito: {err})"
image = None
try:
image = image_payload(path)
except Exception:
pass
ctx = f"[Allegato immagine: {path.name}]\nTesto riconosciuto nell'immagine (OCR):\n{text.strip() or '(nessun testo)'}"
return ctx, image
# Documenti: riusa il lettore del plugin file
try:
from avatar.plugins import registry
registry.load()
text = registry.run("file_leggi", {"percorso": str(path)})
except Exception as err:
text = f"(lettura non riuscita: {err})"
text = re.sub(r"\s+", " ", text)[:MAX_TEXT]
return f"[Allegato: {path.name}]\n{text}", None
+241
View File
@@ -0,0 +1,241 @@
"""Microfono con rilevamento della voce e riproduzione con sincronizzazione labiale."""
from __future__ import annotations
import queue
import threading
import time
from collections import deque
from typing import Callable
import numpy as np
import sounddevice as sd
from core import audio_devices
from core.viseme import VisemeStream
from .lipsync import HOP_SECONDS, float_to_pcm_scale, pcm_level, pcm_visemes
MIC_RATE = 16000
OUT_RATE = 24000
MIC_BLOCK = 480 # 30 ms
PRE_ROLL_BLOCKS = 10 # 300 ms prima dell'inizio del parlato
START_BLOCKS = 3 # 90 ms di voce per iniziare
END_SILENCE_S = 0.8
MIN_SPEECH_S = 0.35
MAX_UTTERANCE_S = 30.0
class Microphone:
"""Cattura continua; segmenta le frasi con una soglia di volume adattiva."""
def __init__(
self,
on_utterance: Callable[[np.ndarray], None],
on_level: Callable[[float], None],
can_listen: Callable[[], bool],
threshold: float = 0.08,
on_frame: Callable[[np.ndarray], None] | None = None,
) -> None:
self._on_utterance = on_utterance
self._on_level = on_level
self._can_listen = can_listen
self._on_frame = on_frame
self.threshold = threshold
self._q: queue.Queue[np.ndarray] = queue.Queue()
self._stream: sd.InputStream | None = None
self._thread: threading.Thread | None = None
self._running = False
self._device_name = ""
def start(self, device_name: str = "") -> None:
self.stop()
self._device_name = device_name
try:
device = audio_devices.resolve(device_name, "input") if device_name else None
except Exception:
device = None
self._stream = sd.InputStream(
samplerate=MIC_RATE, channels=1, dtype="int16", blocksize=MIC_BLOCK,
device=device, callback=self._callback,
)
self._stream.start()
self._running = True
self._thread = threading.Thread(target=self._loop, name="mic-vad", daemon=True)
self._thread.start()
def stop(self) -> None:
self._running = False
if self._stream is not None:
try:
self._stream.stop()
self._stream.close()
except Exception:
pass
self._stream = None
def _callback(self, indata, frames, t, status) -> None:
self._q.put(indata[:, 0].copy())
def _loop(self) -> None:
pre = deque(maxlen=PRE_ROLL_BLOCKS)
collecting: list[np.ndarray] = []
voiced_run = 0
silence_s = 0.0
speech_s = 0.0
block_s = MIC_BLOCK / MIC_RATE
floor = 0.0
while self._running:
try:
frame = self._q.get(timeout=0.5)
except queue.Empty:
continue
if self._on_frame:
try:
self._on_frame(frame)
except Exception:
pass
level = pcm_level(frame)
# Rumore di fondo: media lenta dei livelli bassi.
if level < self.threshold:
floor = floor * 0.98 + level * 0.02
if not self._can_listen():
# Se stavamo raccogliendo (per esempio il tasto premi-per-parlare
# è stato rilasciato) chiudi la frase invece di perderla.
if collecting and speech_s >= MIN_SPEECH_S:
audio = np.concatenate(collecting)
collecting = []
try:
self._on_utterance(audio)
except Exception:
pass
pre.clear()
collecting = []
voiced_run = 0
continue
self._on_level(level)
voiced = level > max(self.threshold, floor * 2.5 + 0.02)
if not collecting:
pre.append(frame)
voiced_run = voiced_run + 1 if voiced else 0
if voiced_run >= START_BLOCKS:
collecting = list(pre)
silence_s = 0.0
speech_s = voiced_run * block_s
voiced_run = 0
continue
collecting.append(frame)
if voiced:
silence_s = 0.0
speech_s += block_s
else:
silence_s += block_s
total_s = len(collecting) * block_s
if silence_s >= END_SILENCE_S or total_s >= MAX_UTTERANCE_S:
audio = np.concatenate(collecting)
collecting = []
pre.clear()
if speech_s >= MIN_SPEECH_S:
try:
self._on_utterance(audio)
except Exception:
pass
class Player:
"""Riproduce audio a 24 kHz e programma le forme della bocca sul tempo reale."""
def __init__(self, on_visemes: Callable, on_level: Callable[[float], None]) -> None:
self._on_visemes = on_visemes
self._on_level = on_level
self._visemes = VisemeStream()
self._stream: sd.OutputStream | None = None
self._lock = threading.Lock()
self._device_name = ""
self._cursor = 0.0
def open(self, device_name: str = "") -> None:
self.close()
self._device_name = device_name
try:
device = audio_devices.resolve(device_name, "output") if device_name else None
except Exception:
device = None
self._stream = sd.OutputStream(samplerate=OUT_RATE, channels=1, dtype="float32", device=device)
self._stream.start()
def close(self) -> None:
if self._stream is not None:
try:
self._stream.stop()
self._stream.close()
except Exception:
pass
self._stream = None
@property
def latency(self) -> float:
try:
return float(self._stream.latency) if self._stream else 0.05
except Exception:
return 0.05
def play(self, audio: np.ndarray, text: str, stop: threading.Event) -> None:
"""Bloccante: suona `audio` (float32 mono 24 kHz) mentre anima la bocca."""
if self._stream is None:
self.open(self._device_name)
assert self._stream is not None
audio = np.asarray(audio, dtype=np.float32).reshape(-1)
if audio.size == 0:
return
with self._lock:
try:
self._visemes.feed_text(text)
except Exception:
pass
# Il cursore è il momento in cui il prossimo blocco inizierà a suonare.
now = time.time() + self.latency + 0.04
if self._cursor < now:
self._cursor = now
block = int(OUT_RATE * 0.2)
for i in range(0, audio.size, block):
if stop.is_set():
break
chunk = audio[i:i + block]
frames = pcm_visemes(float_to_pcm_scale(chunk), OUT_RATE)
if frames:
try:
frames = self._visemes.frames(frames, HOP_SECONDS)
except Exception:
pass
self._on_visemes(frames, HOP_SECONDS, self._cursor)
self._on_level(max(f[0] for f in frames))
else:
self._on_level(pcm_level(float_to_pcm_scale(chunk)))
self._cursor += chunk.size / OUT_RATE
try:
self._stream.write(chunk.reshape(-1, 1))
except Exception:
break
if stop.is_set():
self.abort()
def abort(self) -> None:
"""Ferma subito, scartando l'audio in coda."""
try:
self._visemes.reset()
except Exception:
pass
self._cursor = 0.0
if self._stream is not None:
try:
self._stream.abort()
self._stream.start()
except Exception:
self.open(self._device_name)
self._on_level(0.0)
def drain(self) -> None:
"""Aspetta che l'audio scritto abbia finito di suonare."""
wait = max(0.0, self._cursor - time.time())
if wait > 0:
time.sleep(min(wait, 5.0))
self._on_level(0.0)
+131
View File
@@ -0,0 +1,131 @@
"""Vista 3D dell'avatar (three.js in un QWebEngineView), sincronizzata con stati, volume e visemi."""
from __future__ import annotations
import http.server
import threading
import time
from functools import partial
from pathlib import Path
from PyQt6.QtCore import QTimer, QUrl
from PyQt6.QtGui import QColor
from PyQt6.QtWebEngineCore import QWebEnginePage
from PyQt6.QtWebEngineWidgets import QWebEngineView
ROOT = Path(__file__).resolve().parent.parent / "avatar3d"
MODELS_DIR = ROOT / "models"
def list_models() -> list[str]:
"""Nomi (senza estensione) dei modelli GLB disponibili."""
return sorted(p.stem for p in MODELS_DIR.glob("*.glb"))
class _Quiet(http.server.SimpleHTTPRequestHandler):
def log_message(self, *args) -> None:
pass
_server_port: int | None = None
def serve_root() -> int:
"""Server HTTP locale (una volta sola): i moduli ES non si caricano da file://."""
global _server_port
if _server_port:
return _server_port
handler = partial(_Quiet, directory=str(ROOT))
srv = http.server.ThreadingHTTPServer(("127.0.0.1", 0), handler)
threading.Thread(target=srv.serve_forever, name="avatar3d-http", daemon=True).start()
_server_port = srv.server_address[1]
return _server_port
class _Page(QWebEnginePage):
def javaScriptConsoleMessage(self, level, message, line, source) -> None: # noqa: N802
print(f"[avatar3d] {message} ({Path(source).name}:{line})")
class Avatar3DView(QWebEngineView):
def __init__(self, parent=None) -> None:
super().__init__(parent)
self._page = _Page(self)
self.setPage(self._page)
self._page.setBackgroundColor(QColor("#00060a"))
self._state = "LISTENING"
self._level = 0.0
self._sched: tuple[list, float, float] | None = None
self._ready = False
self._model = ""
self.loadFinished.connect(self._on_loaded)
try:
from avatar.settings import Settings
model = str(Settings().get("avatar_model") or "")
except Exception:
model = ""
self.set_model(model)
self._timer = QTimer(self)
self._timer.timeout.connect(self._tick)
self._timer.start(33)
def _on_loaded(self, ok: bool) -> None:
self._ready = ok
def set_model(self, name: str) -> None:
"""Carica (o ricarica) la pagina con il modello indicato."""
available = list_models()
if name not in available:
name = available[0] if available else name
if name == self._model and self._ready:
return
self._model = name
self._ready = False
self.load(QUrl(f"http://127.0.0.1:{serve_root()}/index.html?model={name}"))
# ── API thread-safe (valori semplici, letti dal timer sul thread Qt) ──
def set_state(self, state: str) -> None:
self._state = str(state or "LISTENING")
def set_audio_level(self, level: float) -> None:
try:
self._level = max(self._level, min(1.0, max(0.0, float(level))))
except (TypeError, ValueError):
pass
def push_visemes(self, frames, hop: float, at: float) -> None:
if not frames:
return
new = list(frames)
cur = self._sched
if cur is not None and abs(cur[2] - hop) < 1e-6:
old, t0, _ = cur
i = int(round((at - t0) / hop))
if 0 <= i <= len(old) + 1:
merged = old[:i] + new
played = int((time.time() - t0) / hop) - 2
if played > 60:
merged, t0 = merged[played:], t0 + played * hop
self._sched = (merged, t0, hop)
return
self._sched = (new, float(at), float(hop))
def glance(self, dx: float, dy: float, hold: float = 1.1) -> None:
if self._ready:
self._page.runJavaScript(f"window.avatar3d && avatar3d.glance({dx:.3f},{dy:.3f},{hold:.2f})")
def _tick(self) -> None:
open_, width, vis_level = 0.0, 0.0, 0.0
sched = self._sched
if sched is not None:
frames, t0, hop = sched
i = int((time.time() - t0) / hop)
if i >= len(frames):
self._sched = None
elif i >= 0:
vis_level, open_, width = frames[i]
level = max(self._level, vis_level)
self._level *= 0.82
if self._ready:
self._page.runJavaScript(
f"window.avatar3d && avatar3d.update({level:.3f},{open_:.3f},{width:.3f},'{self._state}')"
)
View File
Whitespace-only changes.
+169
View File
@@ -0,0 +1,169 @@
"""Motore Claude (API Anthropic) con ricerca web server-side e strumenti di memoria."""
from __future__ import annotations
import threading
import anthropic
from avatar.memory_tools import anthropic_tools, memory_prompt, parse_args, run_tool
from avatar.plugins import registry
from .base import Emit, History, persona_text, today_label, user_block
MODEL = "claude-opus-5"
MAX_ROUNDS = 12
WEB_SEARCH = {
"type": "web_search_20260209",
"name": "web_search",
"max_uses": 5,
"user_location": {"type": "approximate", "country": "IT", "timezone": "Europe/Rome"},
}
class Aborted(Exception):
pass
def friendly_error(err: Exception) -> str:
if isinstance(err, anthropic.AuthenticationError):
return "La chiave API Anthropic non è valida. Controllala nelle impostazioni."
if isinstance(err, anthropic.PermissionDeniedError):
return "La chiave API non ha i permessi per questo modello."
if isinstance(err, anthropic.RateLimitError):
return "Troppe richieste in poco tempo. Riprova tra qualche secondo."
if isinstance(err, anthropic.BadRequestError):
return f"Richiesta non valida: {err.message}"
if isinstance(err, anthropic.APIConnectionError):
return "Non riesco a raggiungere il servizio Anthropic. Controlla la connessione."
if isinstance(err, anthropic.APIStatusError):
return f"Errore del servizio ({err.status_code}): {err.message}"
return str(err)
class AnthropicEngine:
name = "anthropic"
def __init__(self, api_key: str, assistant_name: str, user_name: str, effort: str) -> None:
self.client = anthropic.Anthropic(api_key=api_key)
self.assistant_name = assistant_name
self.user_name = user_name
self.effort = effort
self.history = History("anthropic")
self._compaction = True
def reset(self) -> None:
self.history.clear()
def _system(self) -> list[dict]:
return [
{"type": "text", "text": persona_text(self.assistant_name), "cache_control": {"type": "ephemeral"}},
{"type": "text", "text": user_block(self.user_name, memory_prompt())},
{"type": "text", "text": f"Adesso è {today_label()}.\n\nSe uno strumento risponde con [CONFIRMATION_PENDING], chiedi all'utente di confermare sul pannello e non dire che è fatto."},
]
@staticmethod
def _sources(content: list[dict]) -> list[dict]:
seen: dict[str, dict] = {}
for block in content:
for c in block.get("citations") or []:
if c.get("type") == "web_search_result_location" and c.get("url") and c["url"] not in seen:
seen[c["url"]] = {"title": c.get("title") or c["url"], "url": c["url"]}
return list(seen.values())
def send(self, user_text: str, emit: Emit, abort: threading.Event, image: tuple[str, bytes] | None = None, **_) -> None:
if image:
import base64
media, data = image
self.history.messages.append({"role": "user", "content": [
{"type": "image", "source": {"type": "base64", "media_type": media, "data": base64.b64encode(data).decode()}},
{"type": "text", "text": user_text},
]})
else:
self.history.messages.append({"role": "user", "content": user_text})
full = ""
sources: list[dict] = []
try:
emit({"type": "status", "status": "thinking"})
for _ in range(MAX_ROUNDS):
kwargs: dict = dict(
model=MODEL,
max_tokens=16000,
betas=["server-side-fallback-2026-07-01"] + (["compact-2026-01-12"] if self._compaction else []),
fallbacks="default",
output_config={"effort": self.effort},
system=self._system(),
tools=[WEB_SEARCH, *anthropic_tools(), *registry.anthropic_tools()],
messages=self.history.messages,
)
if self._compaction:
kwargs["context_management"] = {"edits": [{"type": "compact_20260112"}]}
announced = False
try:
with self.client.beta.messages.stream(**kwargs) as stream:
for event in stream:
if abort.is_set():
raise Aborted()
if event.type == "content_block_start":
t = event.content_block.type
if t == "server_tool_use":
emit({"type": "status", "status": "searching"})
elif t == "tool_use":
emit({"type": "status", "status": "memory"})
elif t == "thinking":
emit({"type": "status", "status": "thinking"})
elif event.type == "content_block_delta" and event.delta.type == "text_delta":
if not announced:
announced = True
emit({"type": "status", "status": "responding"})
full += event.delta.text
emit({"type": "text", "delta": event.delta.text})
message = stream.get_final_message()
except anthropic.BadRequestError as err:
if self._compaction and not full:
print(f"[Claude] compattazione rifiutata, la disattivo: {err.message}")
self._compaction = False
continue
raise
content = message.model_dump(mode="json")["content"]
sources += self._sources(content)
self.history.messages.append({"role": "assistant", "content": content})
if message.stop_reason == "refusal":
if not full:
why = getattr(message.stop_details, "explanation", None) if message.stop_details else None
full = "Non posso aiutarti con questa richiesta." + (f" ({why})" if why else "")
break
if message.stop_reason == "pause_turn":
continue
tool_uses = [b for b in message.content if b.type == "tool_use"]
if message.stop_reason != "tool_use" or not tool_uses:
break
results = []
for b in tool_uses:
if registry.has(b.name):
emit({"type": "status", "status": "working", "detail": b.name})
res, ev = registry.run(b.name, parse_args(b.input)), None
else:
res, ev = run_tool(b.name, parse_args(b.input))
if ev:
emit(ev)
results.append({"type": "tool_result", "tool_use_id": b.id, "content": res})
self.history.messages.append({"role": "user", "content": results})
emit({"type": "status", "status": "thinking"})
full = full.strip() or "(nessuna risposta)"
self.history.save()
emit({"type": "done", "text": full, "sources": sources})
except Aborted:
self._rollback()
emit({"type": "error", "message": "Risposta interrotta.", "aborted": True})
except Exception as err:
self._rollback()
emit({"type": "error", "message": friendly_error(err)})
def _rollback(self) -> None:
msgs = self.history.messages
if msgs and msgs[-1].get("role") == "user" and (isinstance(msgs[-1].get("content"), str) or
(isinstance(msgs[-1].get("content"), list) and any(b.get("type") == "image" for b in msgs[-1]["content"]))):
msgs.pop()
self.history.save()
+69
View File
@@ -0,0 +1,69 @@
"""Tipi comuni ai motori di conversazione."""
from __future__ import annotations
import json
import locale
import threading
from datetime import datetime
from pathlib import Path
from typing import Any, Callable, Protocol
from avatar.settings import BASE_DIR, DATA_DIR
Event = dict[str, Any]
Emit = Callable[[Event], None]
PERSONA_FILE = BASE_DIR / "avatar" / "persona.md"
class ChatBackend(Protocol):
def send(self, user_text: str, emit: Emit, abort: threading.Event) -> None: ...
def reset(self) -> None: ...
def today_label() -> str:
try:
locale.setlocale(locale.LC_TIME, "it_IT.UTF-8")
except Exception:
pass
d = datetime.now()
giorni = ["lunedì", "martedì", "mercoledì", "giovedì", "venerdì", "sabato", "domenica"]
mesi = ["gennaio", "febbraio", "marzo", "aprile", "maggio", "giugno", "luglio",
"agosto", "settembre", "ottobre", "novembre", "dicembre"]
return f"{giorni[d.weekday()]} {d.day} {mesi[d.month - 1]} {d.year}, ore {d:%H:%M}"
def persona_text(assistant_name: str) -> str:
try:
text = PERSONA_FILE.read_text(encoding="utf-8")
except Exception:
text = "Sei Ava, un'assistente personale gentile e concreta. Rispondi in italiano."
return text.replace("Ava", assistant_name or "Ava")
def user_block(user_name: str, memory_block: str) -> str:
who = f"L'utente si chiama {user_name}." if user_name else "Non conosci ancora il nome dell'utente: chiediglielo con naturalezza e salvalo."
return f"{who}\n\nMemoria a lungo termine sull'utente:\n{memory_block}"
class History:
"""Cronologia su disco per un motore."""
def __init__(self, name: str) -> None:
self.file: Path = DATA_DIR / f"conversation-{name}.json"
self.messages: list[Any] = []
self.meta: dict[str, Any] = {}
try:
d = json.loads(self.file.read_text(encoding="utf-8"))
self.messages = list(d.get("messages", []))
self.meta = dict(d.get("meta", {}))
except Exception:
pass
def save(self) -> None:
self.file.parent.mkdir(parents=True, exist_ok=True)
self.file.write_text(json.dumps({"messages": self.messages, "meta": self.meta}, ensure_ascii=False), encoding="utf-8")
def clear(self) -> None:
self.messages, self.meta = [], {}
self.save()
+217
View File
@@ -0,0 +1,217 @@
"""Motore Claude Code: avvia `claude -p` in streaming JSON con un server MCP per la memoria."""
from __future__ import annotations
import json
import os
import shutil
import subprocess
import sys
import threading
import uuid
from pathlib import Path
from memory import memory_manager as mm
from avatar.memory_tools import memory_prompt
from .base import Emit, History, persona_text, today_label, user_block
MEMORY_TOOLS = ["mcp__avatar"] # tutti gli strumenti del server MCP dell'app (memoria e plugin)
ACCESS = {
"chat": {"tools": "WebSearch,WebFetch", "allowed": ["WebSearch", "WebFetch"], "mode": None},
"read": {"tools": "Read,Glob,Grep,WebSearch,WebFetch", "allowed": ["Read", "Glob", "Grep", "WebSearch", "WebFetch"], "mode": None},
"full": {"tools": "default",
"allowed": ["Bash", "Read", "Edit", "Write", "MultiEdit", "NotebookEdit", "Glob", "Grep", "WebSearch", "WebFetch", "Agent"],
"mode": "acceptEdits"},
}
CANDIDATES = [Path.home() / ".local/bin/claude", Path("/opt/homebrew/bin/claude"), Path("/usr/local/bin/claude"),
Path.home() / ".claude/local/claude", Path.home() / ".npm-global/bin/claude"]
MCP_SCRIPT = Path(__file__).resolve().parent.parent / "memory_mcp.py"
def resolve_claude(custom: str = "") -> str | None:
if custom.strip():
return custom.strip() if Path(custom.strip()).exists() else None
found = shutil.which("claude")
if found:
return found
for p in CANDIDATES:
if p.exists():
return str(p)
return None
def claude_env(config_dir: str) -> dict:
env = dict(os.environ)
env.pop("CLAUDECODE", None)
if config_dir.strip():
env["CLAUDE_CONFIG_DIR"] = config_dir.strip()
return env
def check_claude_code(custom_path: str, config_dir: str) -> dict:
"""Comando, profilo e account collegato."""
binary = resolve_claude(custom_path)
home = Path.home()
profiles = sorted(str(p) for p in home.iterdir() if p.is_dir() and (p.name == ".claude" or p.name.startswith(".claude-")))
effective = config_dir.strip() or os.environ.get("CLAUDE_CONFIG_DIR") or str(home / ".claude")
st = {"binary": binary, "configDir": effective, "loggedIn": False, "profiles": profiles}
if not binary:
st["error"] = "Comando claude non trovato."
return st
try:
out = subprocess.run([binary, "auth", "status"], env=claude_env(config_dir), capture_output=True, text=True, timeout=15)
d = json.loads(out.stdout)
st.update(loggedIn=bool(d.get("loggedIn")), email=d.get("email"), authMethod=d.get("authMethod"),
configDir=d.get("configDirectory") or effective)
except Exception as err:
st["error"] = str(err)[:200]
return st
class ClaudeCodeEngine:
name = "claudecode"
def __init__(self, model: str, access: str, claude_path: str, config_dir: str,
assistant_name: str, user_name: str, effort: str) -> None:
self.model, self.access, self.claude_path, self.config_dir = model, access, claude_path, config_dir
self.assistant_name, self.user_name, self.effort = assistant_name, user_name, effort
self.history = History("claudecode")
self._proc: subprocess.Popen | None = None
def reset(self) -> None:
self.history.clear()
def _system_prompt(self) -> str:
access = {
"chat": "Non hai accesso ai file o ai comandi del Mac: se ti chiedono di farlo, spiega che serve alzare il livello di accesso nelle impostazioni dell'app.",
"read": "Puoi leggere file e cercare nelle cartelle dell'utente, ma non modificare nulla né eseguire comandi.",
"full": "Puoi leggere e modificare file ed eseguire comandi sul Mac dell'utente. Fallo con prudenza e spiega brevemente cosa hai fatto.",
}[self.access]
return "\n\n".join([
persona_text(self.assistant_name),
user_block(self.user_name, memory_prompt()),
f"Adesso è {today_label()}.",
"Strumenti di memoria (obbligatori): per ricordare un fatto sull'utente chiama `mcp__avatar__salva_memoria`; per cancellarne uno `mcp__avatar__dimentica_memoria`; per cercare `mcp__avatar__cerca_memoria`. Non dire mai di aver salvato qualcosa senza averlo chiamato davvero. Gli altri strumenti `mcp__avatar__*` (per esempio il calendario) agiscono sul Mac dell'utente: usali quando servono, e per le azioni irreversibili chiedi conferma a voce prima di richiamarli con confermato=true.",
access,
"Stai parlando dentro un'app vocale: rispondi in modo conversazionale, senza intestazioni Markdown né tabelle, e non citare nomi di file o strumenti interni a meno che non serva.",
])
def _args(self, session_id: str, resume: bool, attach_path: str = "") -> list[str]:
flags = dict(ACCESS[self.access])
if attach_path:
# Permesso di lettura limitato al file allegato (immagini incluse: Claude Code le vede).
if "Read" not in flags["tools"] and flags["tools"] != "default":
flags["tools"] = flags["tools"] + ",Read"
flags["allowed"] = [*flags["allowed"], f"Read(//{attach_path.lstrip('/')})"]
mcp = {"mcpServers": {"avatar": {"command": sys.executable, "args": [str(MCP_SCRIPT)]}}}
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),
"--mcp-config", json.dumps(mcp), "--strict-mcp-config",
"--system-prompt-snapshot", "off",
"--resume" if resume else "--session-id", session_id]
if flags["mode"]:
args += ["--permission-mode", flags["mode"]]
args += ["--system-prompt" if self.access == "chat" else "--append-system-prompt", self._system_prompt()]
return args
def _run_once(self, text: str, session_id: str, resume: bool, emit: Emit, abort: threading.Event, attach_path: str = "") -> dict:
binary = resolve_claude(self.claude_path)
if not binary:
raise RuntimeError("Non trovo il comando `claude`. Installa Claude Code o indica il percorso nelle impostazioni.")
proc = subprocess.Popen([binary, *self._args(session_id, resume, attach_path)], cwd=str(Path.home()), env=claude_env(self.config_dir),
stdin=subprocess.PIPE, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True)
self._proc = proc
out = {"text": "", "initialized": False, "error": None, "result": ""}
announced = False
def watch_abort():
abort.wait()
if proc.poll() is None:
proc.terminate()
threading.Thread(target=watch_abort, daemon=True).start()
try:
proc.stdin.write(text)
proc.stdin.close()
except Exception:
pass
for line in proc.stdout:
line = line.strip()
if not line:
continue
try:
msg = json.loads(line)
except Exception:
continue
t = msg.get("type")
if t == "system" and msg.get("subtype") == "init":
out["initialized"] = True
elif t == "stream_event":
ev = msg.get("event") or {}
if ev.get("type") == "content_block_start" and (ev.get("content_block") or {}).get("type") == "tool_use":
name = ev["content_block"].get("name", "")
if name.startswith("mcp__memoria__"):
emit({"type": "status", "status": "memory"})
elif name in ("WebSearch", "WebFetch"):
emit({"type": "status", "status": "searching"})
else:
emit({"type": "status", "status": "working", "detail": name})
elif ev.get("type") == "content_block_delta" and (ev.get("delta") or {}).get("type") == "text_delta":
delta = ev["delta"].get("text", "")
if delta:
if not announced:
announced = True
emit({"type": "status", "status": "responding"})
out["text"] += delta
emit({"type": "text", "delta": delta})
elif ev.get("type") == "message_start" and announced:
emit({"type": "status", "status": "thinking"})
elif t == "result":
out["result"] = msg.get("result") or ""
if msg.get("is_error") or (msg.get("subtype") and msg["subtype"] != "success"):
errs = msg.get("errors") or []
out["error"] = "; ".join(map(str, errs)) or out["result"] or msg.get("subtype")
stderr = proc.stderr.read() if proc.stderr else ""
out["code"] = proc.wait()
out["stderr"] = stderr
self._proc = None
if not out["text"] and out["result"] and not out["error"]:
out["text"] = out["result"]
emit({"type": "text", "delta": out["result"]})
return out
def send(self, user_text: str, emit: Emit, abort: threading.Event, attach_path: str = "", **_) -> None:
if attach_path:
user_text += f"\n\n(Il file allegato è in {attach_path}: se ti serve vederlo, leggilo con lo strumento Read.)"
before = {(r["category"], r["key"]) for r in mm.all_entries_for_ui()}
session_id = str(self.history.meta.get("sessionId") or "")
resume = bool(session_id)
session_id = session_id or str(uuid.uuid4())
try:
emit({"type": "status", "status": "thinking"})
out = self._run_once(user_text, session_id, resume, emit, abort, attach_path)
if not abort.is_set() and resume and not out["initialized"] and out["code"] != 0:
print("[Claude Code] sessione non ripresa, ne creo una nuova:", out["stderr"][:200])
session_id, resume = str(uuid.uuid4()), False
out = self._run_once(user_text, session_id, resume, emit, abort, attach_path)
if out["initialized"]:
self.history.meta["sessionId"] = session_id
for r in mm.all_entries_for_ui():
if (r["category"], r["key"]) not in before:
emit({"type": "memory_saved", "text": r["value"]})
if abort.is_set():
self.history.save()
emit({"type": "error", "message": "Risposta interrotta.", "aborted": True})
return
if out["error"] or (out["code"] != 0 and not out["text"]):
detail = out["error"] or " ".join(out["stderr"].strip().splitlines()[-3:])
raise RuntimeError(detail or f"Claude Code è terminato con codice {out['code']}")
text = out["text"].strip() or "(nessuna risposta)"
self.history.messages.append({"role": "user", "text": user_text})
self.history.messages.append({"role": "assistant", "text": text})
self.history.save()
emit({"type": "done", "text": text, "sources": []})
except Exception as err:
self.history.save()
emit({"type": "error", "message": str(err)})
+220
View File
@@ -0,0 +1,220 @@
"""Motore per server compatibili OpenAI (vLLM, Ollama, LM Studio…)."""
from __future__ import annotations
import re
import threading
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
MAX_ROUNDS = 8
MAX_HISTORY = 30
SEARCH_TOOL = {
"type": "function",
"function": {
"name": "cerca_web",
"description": "Cerca sul web informazioni aggiornate (notizie, prezzi, orari, meteo, eventi). Restituisce titoli, link e descrizioni.",
"parameters": {"type": "object", "properties": {"query": {"type": "string"}}, "required": ["query"]},
},
}
class Aborted(Exception):
pass
def friendly_error(err: Exception, base_url: str) -> str:
if isinstance(err, openai.APIConnectionError):
return f"Non riesco a raggiungere il server locale su {base_url}. È avviato?"
if isinstance(err, openai.AuthenticationError):
return "Il server locale ha rifiutato la chiave API."
if isinstance(err, openai.NotFoundError):
return "Modello non trovato sul server locale. Controlla il nome nelle impostazioni."
if isinstance(err, openai.BadRequestError):
return f"Il server locale ha rifiutato la richiesta: {err.message}"
if isinstance(err, openai.APIStatusError):
return f"Errore del server locale ({err.status_code}): {err.message}"
return str(err)
class ThinkFilter:
"""Nasconde i blocchi <think>…</think> dei modelli ragionanti."""
def __init__(self) -> None:
self.inside = False
self.carry = ""
def push(self, delta: str) -> str:
text, self.carry, out = self.carry + delta, "", ""
while text:
if self.inside:
end = text.find("</think>")
if end < 0:
self.carry = text[-8:]
return out
text = text[end + 8:].lstrip()
self.inside = False
else:
start = text.find("<think>")
if start < 0:
m = re.search(r"<(t(h(i(n(k)?)?)?)?)?$", text)
if m:
self.carry = m.group(0)
out += text[: m.start()]
else:
out += text
return out
out += text[:start]
text = text[start + 7:]
self.inside = True
return out
class OpenAICompatEngine:
name = "local"
def __init__(self, base_url: str, model: str, api_key: str, search_api_key: str,
assistant_name: str, user_name: str) -> None:
self.base_url, self.model, self.search_api_key = base_url, model, search_api_key
self.assistant_name, self.user_name = assistant_name, user_name
self.client = openai.OpenAI(base_url=base_url, api_key=api_key or "non-necessaria")
self.history = History("local")
self._tools_supported = True
def reset(self) -> None:
self.history.clear()
def _system(self) -> dict:
if self._tools_supported:
note = ("\n\nRegole sugli strumenti (obbligatorie):\n"
"- Quando l'utente ti chiede di ricordare qualcosa, o ti dice un fatto importante su di sé, DEVI chiamare salva_memoria prima di rispondere. Non dire mai di aver salvato senza averlo chiamato davvero.\n"
"- Per cancellare una memoria chiama dimentica_memoria; per cercarne una non presente nel prompt chiama cerca_memoria.\n"
"- Per azioni sul Mac (per esempio il calendario) usa gli strumenti dedicati. Se uno risponde con [CONFIRMATION_PENDING], chiedi all'utente di confermare sul pannello e non dire che è fatto.")
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.")
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())}\n\nAdesso è {today_label()}.{note}"}
def _recent(self) -> list:
msgs = self.history.messages
if len(msgs) <= MAX_HISTORY:
return msgs
start = len(msgs) - MAX_HISTORY
while start < len(msgs) and msgs[start].get("role") != "user":
start += 1
return msgs[start:]
def _tools(self):
if not self._tools_supported:
return None
return openai_tools() + registry.openai_tools() + ([SEARCH_TOOL] if self.search_api_key else [])
def _run_tool(self, name: str, raw_args: str, emit: Emit, sources: list) -> str:
args = parse_args(raw_args)
if name == "cerca_web":
emit({"type": "status", "status": "searching"})
try:
hits = brave_search(str(args.get("query", "")), self.search_api_key)
except Exception as err:
return f"Errore nella ricerca: {err}"
for h in hits[:5]:
if all(s["url"] != h["url"] for s in sources):
sources.append({"title": h["title"], "url": h["url"]})
return "\n".join(f"{i+1}. {h['title']}\n {h['url']}\n {h['description']}" for i, h in enumerate(hits)) or "Nessun risultato."
if registry.has(name):
emit({"type": "status", "status": "working", "detail": name})
return registry.run(name, args)
emit({"type": "status", "status": "memory"})
res, ev = run_tool(name, args)
if ev:
emit(ev)
return res
def send(self, user_text: str, emit: Emit, abort: threading.Event, **_) -> None:
self.history.messages.append({"role": "user", "content": user_text})
full = ""
sources: list[dict] = []
try:
emit({"type": "status", "status": "thinking"})
for rnd in range(MAX_ROUNDS):
tools = self._tools()
kwargs: dict = dict(model=self.model, messages=[self._system(), *self._recent()], stream=True)
if tools:
kwargs.update(tools=tools, tool_choice="auto")
try:
stream = self.client.chat.completions.create(**kwargs)
except openai.BadRequestError as err:
if tools and not full:
print(f"[locale] il server rifiuta gli strumenti, li disattivo: {err.message}")
self._tools_supported = False
continue
raise
flt, calls, round_text, announced, finish = ThinkFilter(), {}, "", False, None
for chunk in stream:
if abort.is_set():
stream.close()
raise Aborted()
if not chunk.choices:
continue
choice = chunk.choices[0]
delta = choice.delta
if delta and delta.content:
visible = flt.push(delta.content)
if visible:
if not announced:
announced = True
emit({"type": "status", "status": "responding"})
round_text += visible
full += visible
emit({"type": "text", "delta": visible})
for tc in (delta.tool_calls if delta else None) or []:
p = calls.setdefault(tc.index, {"id": "", "name": "", "args": ""})
if tc.id:
p["id"] = tc.id
if tc.function and tc.function.name:
p["name"] += tc.function.name
if tc.function and tc.function.arguments:
p["args"] += tc.function.arguments
if choice.finish_reason:
finish = choice.finish_reason
tool_calls = [c for c in calls.values() if c["name"]]
if not tool_calls:
self.history.messages.append({"role": "assistant", "content": round_text})
break
for i, c in enumerate(tool_calls):
c["id"] = c["id"] or f"call_{rnd}_{i}"
self.history.messages.append({
"role": "assistant", "content": round_text or None,
"tool_calls": [{"id": c["id"], "type": "function", "function": {"name": c["name"], "arguments": c["args"] or "{}"}} for c in tool_calls],
})
for c in tool_calls:
result = self._run_tool(c["name"], c["args"], emit, sources)
self.history.messages.append({"role": "tool", "tool_call_id": c["id"], "content": result})
emit({"type": "status", "status": "thinking"})
if finish == "length":
break
full = full.strip() or "(nessuna risposta)"
self.history.save()
emit({"type": "done", "text": full, "sources": sources})
except Aborted:
self._rollback()
emit({"type": "error", "message": "Risposta interrotta.", "aborted": True})
except Exception as err:
self._rollback()
emit({"type": "error", "message": friendly_error(err, self.base_url)})
def _rollback(self) -> None:
msgs = self.history.messages
if msgs and msgs[-1].get("role") == "user":
msgs.pop()
self.history.save()
def list_models(base_url: str, api_key: str) -> list[str]:
client = openai.OpenAI(base_url=base_url, api_key=api_key or "non-necessaria", timeout=8, max_retries=0)
return [m.id for m in client.models.list()]
+67
View File
@@ -0,0 +1,67 @@
"""Livello audio e forme della bocca dallo spettro (adattato da Mark-LIV, CC BY-NC 4.0)."""
from __future__ import annotations
import numpy as np
# Soglie sull'ampiezza int16 (come nell'originale).
_LEVEL_FLOOR = 60.0
_LEVEL_FULL = 2600.0
VIS_WIN = 1024 # ~43 ms a 24 kHz
VIS_HOP = 480 # 20 ms → 50 forme al secondo
HOP_SECONDS = VIS_HOP / 24000.0
def pcm_level(samples) -> float:
"""Volume 0..1 da campioni in scala int16 (float o int)."""
try:
x = np.asarray(samples, dtype=np.float32)
if x.size == 0:
return 0.0
rms = float(np.sqrt(np.mean(x * x)))
except Exception:
return 0.0
if rms <= _LEVEL_FLOOR:
return 0.0
return min(1.0, (rms - _LEVEL_FLOOR) / (_LEVEL_FULL - _LEVEL_FLOOR))
def float_to_pcm_scale(audio: np.ndarray) -> np.ndarray:
"""Audio float32 -1..1 → scala int16 (senza cambiare tipo)."""
return np.asarray(audio, dtype=np.float32) * 32767.0
def pcm_visemes(samples, sr: int = 24000):
"""Frame (livello, apertura, larghezza) ogni 20 ms da un blocco audio in scala int16."""
try:
x = np.asarray(samples, dtype=np.float32)
if x.size < VIS_WIN:
return []
win = np.hanning(VIS_WIN).astype(np.float32)
freqs = np.fft.rfftfreq(VIS_WIN, 1.0 / sr)
b_f1_lo = (freqs >= 150) & (freqs < 450)
b_f1_hi = (freqs >= 450) & (freqs < 1100)
b_f2_bk = (freqs >= 600) & (freqs < 1300)
b_f2_fr = (freqs >= 1700) & (freqs < 3200)
b_hiss = (freqs >= 3800) & (freqs < 8000)
out = []
for start in range(0, x.size, VIS_HOP):
level = pcm_level(x[start:start + VIS_HOP])
seg = x[start:start + VIS_WIN]
if seg.size < VIS_WIN:
seg = np.concatenate([seg, np.zeros(VIS_WIN - seg.size, dtype=np.float32)])
if level <= 0.0:
out.append((0.0, 0.0, 0.0))
continue
mag = np.abs(np.fft.rfft((seg - seg.mean()) * win))
f1l, f1h = float(mag[b_f1_lo].sum()), float(mag[b_f1_hi].sum())
f2b, f2f = float(mag[b_f2_bk].sum()), float(mag[b_f2_fr].sum())
hiss = float(mag[b_hiss].sum())
openness = f1h / (f1l + f1h + 1e-6)
width = (f2f - f2b) / (f2f + f2b + 1e-6)
width *= (1.0 - openness) ** 0.8
h = hiss / (f1l + f1h + f2b + f2f + hiss + 1e-6)
openness *= 1.0 - 0.65 * min(1.0, h * 2.5)
out.append((level, float(min(1.0, max(0.0, openness))), float(min(1.0, max(-1.0, width)))))
return out
except Exception:
return []
+35
View File
@@ -0,0 +1,35 @@
"""Funzioni comuni per i plugin che parlano con macOS."""
from __future__ import annotations
import subprocess
def osa(script: str, *args: str, timeout: int = 60) -> str:
"""Esegue AppleScript; i parametri arrivano in `argv` (niente problemi di virgolette)."""
res = subprocess.run(["osascript", "-e", script, "--", *args], capture_output=True, text=True, timeout=timeout)
if res.returncode != 0:
err = (res.stderr or "").strip()
if "-1743" in err:
raise RuntimeError("Permesso negato: concedilo in Impostazioni di Sistema > Privacy e sicurezza > Automazione.")
raise RuntimeError(err.splitlines()[-1] if err else "errore AppleScript")
return res.stdout.rstrip("\n")
def sh(*cmd: str, timeout: int = 30) -> str:
res = subprocess.run(list(cmd), capture_output=True, text=True, timeout=timeout)
if res.returncode != 0:
raise RuntimeError((res.stderr or res.stdout).strip().splitlines()[-1] if (res.stderr or res.stdout).strip() else f"{cmd[0]} ha fallito")
return res.stdout.rstrip("\n")
def confirm_or_param(ctx: dict, params: dict, key: str, title: str, detail: str, do):
"""Conferma a schermo se c'è l'interfaccia; altrimenti richiede confermato=true."""
confirm = ctx.get("confirm")
if confirm:
return confirm(key, title, detail, do)
if str(params.get("confermato", "")).lower() in ("true", "1", "sì", "si", "yes"):
return do()
return f"Prima di procedere chiedi conferma all'utente per: {title} — {detail}. Poi richiama con confermato=true."
CONFERMATO = {"type": "boolean", "description": "Solo senza interfaccia: true dopo la conferma esplicita dell'utente."}
+65
View File
@@ -0,0 +1,65 @@
"""Server MCP (stdio) che espone la memoria dell'assistente a Claude Code."""
from __future__ import annotations
import json
import sys
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from avatar.memory_tools import mcp_tools, memory_prompt, run_tool # noqa: E402
from avatar.plugins import registry # noqa: E402
TOOLS = mcp_tools() + [{"name": "elenca_memorie", "description": "Elenca le memorie salvate sull'utente.",
"inputSchema": {"type": "object", "properties": {}}}] + registry.mcp_tools()
def send(msg: dict) -> None:
sys.stdout.write(json.dumps(msg, ensure_ascii=False) + "\n")
sys.stdout.flush()
def text_result(text: str, is_error: bool = False) -> dict:
return {"content": [{"type": "text", "text": text}], "isError": is_error}
def main() -> None:
for line in sys.stdin:
line = line.strip()
if not line:
continue
try:
req = json.loads(line)
except Exception:
continue
rid, method, params = req.get("id"), req.get("method"), req.get("params") or {}
reply = (lambda result: send({"jsonrpc": "2.0", "id": rid, "result": result})) if rid is not None else (lambda r: None)
if method == "initialize":
reply({"protocolVersion": params.get("protocolVersion", "2025-06-18"),
"capabilities": {"tools": {}}, "serverInfo": {"name": "avatar", "version": "1.1.0"}})
elif method in ("notifications/initialized", "notifications/cancelled"):
pass
elif method == "ping":
reply({})
elif method == "tools/list":
reply({"tools": TOOLS})
elif method == "tools/call":
name = params.get("name", "")
args = params.get("arguments") or {}
try:
if name == "elenca_memorie":
reply(text_result(memory_prompt()))
elif registry.has(name):
res = registry.run(name, args)
reply(text_result(res, res.startswith("Errore")))
else:
res, _ = run_tool(name, args)
reply(text_result(res, res.startswith("Errore")))
except Exception as err:
reply(text_result(f"Errore: {err}", True))
elif rid is not None:
send({"jsonrpc": "2.0", "id": rid, "error": {"code": -32601, "message": f"Metodo non supportato: {method}"}})
if __name__ == "__main__":
main()
+131
View File
@@ -0,0 +1,131 @@
"""Strumenti di memoria condivisi dai motori, sopra memory_manager di Mark-LIV."""
from __future__ import annotations
import json
import re
from memory import memory_manager as mm
CATEGORIES = ["identity", "preferences", "projects", "relationships", "wishes", "notes"]
SAVE_SCHEMA = {
"type": "object",
"properties": {
"categoria": {"type": "string", "enum": CATEGORIES,
"description": "identity (nome, età, città…), preferences (gusti), projects, relationships (persone), wishes (desideri, obiettivi), notes (altro)."},
"chiave": {"type": "string", "description": "Etichetta breve in minuscolo con underscore, es. 'nome', 'caffe_preferito', 'sorella_anna'."},
"valore": {"type": "string", "description": "Il fatto, in una frase breve in italiano."},
},
"required": ["categoria", "chiave", "valore"],
"additionalProperties": False,
}
FORGET_SCHEMA = {
"type": "object",
"properties": {
"categoria": {"type": "string", "enum": CATEGORIES},
"chiave": {"type": "string"},
},
"required": ["categoria", "chiave"],
"additionalProperties": False,
}
SEARCH_SCHEMA = {
"type": "object",
"properties": {"query": {"type": "string", "description": "Parole chiave da cercare nella memoria."}},
"required": ["query"],
"additionalProperties": False,
}
BULK_SCHEMA = {
"type": "object",
"properties": {
"voci": {"type": "array", "description": "Elenco di fatti da salvare.",
"items": {"type": "object", "properties": {
"categoria": {"type": "string", "enum": CATEGORIES},
"chiave": {"type": "string"}, "valore": {"type": "string"}},
"required": ["categoria", "chiave", "valore"], "additionalProperties": False}},
},
"required": ["voci"],
"additionalProperties": False,
}
TOOL_SPECS = [
("salva_memorie",
"Salva più fatti sull'utente in una volta sola (per esempio importando memorie da un altro assistente o da un testo con molte informazioni). Ogni voce: categoria, chiave breve, valore in una frase.",
BULK_SCHEMA),
("salva_memoria",
"Salva in modo permanente un fatto sull'utente utile in futuro: nome, persone care, preferenze, abitudini, obiettivi, scadenze. Usalo anche quando l'utente chiede esplicitamente di ricordare qualcosa. Non dire mai di aver salvato senza averlo chiamato.",
SAVE_SCHEMA),
("dimentica_memoria", "Cancella una memoria salvata (categoria e chiave come nell'elenco delle memorie).", FORGET_SCHEMA),
("cerca_memoria", "Cerca nella memoria a lungo termine fatti non presenti nel riepilogo del prompt.", SEARCH_SCHEMA),
]
def anthropic_tools() -> list[dict]:
return [
{"name": n, "description": d, "strict": True, "eager_input_streaming": True, "input_schema": s}
for n, d, s in TOOL_SPECS
]
def openai_tools() -> list[dict]:
return [{"type": "function", "function": {"name": n, "description": d, "parameters": s}} for n, d, s in TOOL_SPECS]
def mcp_tools() -> list[dict]:
return [{"name": n, "description": d, "inputSchema": s} for n, d, s in TOOL_SPECS]
def _slug(s: str) -> str:
s = re.sub(r"[^a-z0-9]+", "_", s.lower().strip()).strip("_")
return s[:40] or "nota"
def run_tool(name: str, args: dict) -> tuple[str, dict | None]:
"""Esegue uno strumento di memoria. Restituisce (risultato, evento per la UI)."""
if name == "salva_memorie":
voci = args.get("voci") or []
saved = []
for v in voci:
if not isinstance(v, dict) or not str(v.get("valore", "")).strip():
continue
cat = v.get("categoria") if v.get("categoria") in CATEGORIES else "notes"
key = _slug(str(v.get("chiave", "")))
mm.remember(key, str(v["valore"]).strip(), cat)
saved.append(f"{cat}/{key}")
if not saved:
return "Nessuna voce valida.", None
return f"Salvate {len(saved)} memorie: {', '.join(saved[:20])}{'…' if len(saved) > 20 else ''}.", {"type": "memory_saved", "text": f"{len(saved)} voci importate"}
if name == "salva_memoria":
cat = args.get("categoria") if args.get("categoria") in CATEGORIES else "notes"
key = _slug(str(args.get("chiave", "")))
val = str(args.get("valore", "")).strip()
if not val:
return "Errore: serve il campo 'valore'.", None
mm.remember(key, val, cat)
return f"Memoria salvata: {cat}/{key}.", {"type": "memory_saved", "text": val}
if name == "dimentica_memoria":
cat = args.get("categoria") if args.get("categoria") in CATEGORIES else "notes"
key = _slug(str(args.get("chiave", "")))
res = mm.forget(key, cat)
ok = res.startswith("Forgotten")
return ("Memoria cancellata." if ok else "Nessuna memoria con quella chiave."), ({"type": "memory_removed"} if ok else None)
if name == "cerca_memoria":
return mm.search_memory(str(args.get("query", "")), limit=8) or "Nessun risultato.", None
return f"Strumento sconosciuto: {name}", None
def memory_prompt() -> str:
try:
block = mm.format_memory_for_prompt(mm.load_memory())
except Exception:
block = ""
return block or "Nessuna memoria salvata finora."
def parse_args(raw) -> dict:
if isinstance(raw, dict):
return raw
try:
return json.loads(raw or "{}")
except Exception:
return {}
+218
View File
@@ -0,0 +1,218 @@
"""Monitor di Mail, WhatsApp e Telegram: raccoglie i messaggi nuovi, li fa valutare e crea avvisi."""
from __future__ import annotations
import json
import re
import sqlite3
import subprocess
import threading
import time
import uuid
from datetime import datetime
from pathlib import Path
from avatar.settings import DATA_DIR, Settings
STATE_FILE = DATA_DIR / "monitor.json"
MAIL_INDEX = Path.home() / "Library" / "Mail" / "V10" / "MailData" / "Envelope Index"
WA_DB = DATA_DIR / "whatsapp" / "index.sqlite"
DEFAULT_RULES = ("richieste di assistenza o di aiuto, domande rivolte a me che aspettano una risposta, problemi o cose che non funzionano, "
"urgenze, scadenze e pagamenti, appuntamenti da confermare o spostare, preventivi e lavori richiesti. "
"NON contano: newsletter, promozioni, notifiche automatiche, conferme d'ordine, saluti e chiacchiere senza richieste.")
KEYWORDS = re.compile(r"\?|urgent|aiut|assistenz|problem|non funziona|non riesco|errore|bloccat|guast|richie|preventiv|fattur|pagament|scadenz|"
r"conferm|appuntament|puoi|potresti|riesci|serve|servirebbe|quando|mi dici|fammi sapere|rispond|chiam|ti prego|per favore|subito|entro", re.I)
_lock = threading.Lock()
class Monitor:
def __init__(self, ui, say, settings: Settings) -> None:
self.ui, self.say, self.settings = ui, say, settings
self.state = self._load()
self._thread: threading.Thread | None = None
self._stop = threading.Event()
# ── stato ────────────────────────────────────────────────────────────
def _load(self) -> dict:
try:
return json.loads(STATE_FILE.read_text(encoding="utf-8"))
except Exception:
return {"mail_rowid": None, "wa_ts": None, "tg": {}, "alerts": [], "seen": []}
def _save(self) -> None:
STATE_FILE.parent.mkdir(parents=True, exist_ok=True)
self.state["alerts"] = self.state.get("alerts", [])[-100:]
self.state["seen"] = self.state.get("seen", [])[-2000:]
STATE_FILE.write_text(json.dumps(self.state, ensure_ascii=False, indent=1), encoding="utf-8")
def start(self) -> None:
if self._thread and self._thread.is_alive():
return
self._stop.clear()
self._thread = threading.Thread(target=self._loop, name="monitor", daemon=True)
self._thread.start()
def stop(self) -> None:
self._stop.set()
def _loop(self) -> None:
time.sleep(20)
while not self._stop.is_set():
if self.settings.get("monitor_enabled"):
try:
self.scan()
except Exception as err:
print(f"[monitor] {err}")
self._stop.wait(int(self.settings.get("monitor_intervallo") or 60))
# ── raccolta ─────────────────────────────────────────────────────────
def _collect_mail(self, baseline: bool) -> list[dict]:
if not MAIL_INDEX.exists():
return []
con = sqlite3.connect(f"file:{MAIL_INDEX}?mode=ro", uri=True, timeout=5)
try:
last = self.state.get("mail_rowid")
if last is None or baseline:
self.state["mail_rowid"] = con.execute("SELECT MAX(ROWID) FROM messages").fetchone()[0] or 0
return []
rows = con.execute("""SELECT m.ROWID, COALESCE(a.comment,''), COALESCE(a.address,''), COALESCE(s.subject,''), COALESCE(su.summary,'')
FROM messages m JOIN mailboxes mb ON m.mailbox = mb.ROWID LEFT JOIN addresses a ON m.sender = a.ROWID
LEFT JOIN subjects s ON m.subject = s.ROWID LEFT JOIN summaries su ON m.summary = su.ROWID
WHERE m.ROWID > ? AND m.deleted = 0 AND (mb.url LIKE '%/INBOX' OR mb.url LIKE '%/Inbox') ORDER BY m.ROWID LIMIT 40""", (last,)).fetchall()
if rows:
self.state["mail_rowid"] = max(r[0] for r in rows)
return [{"fonte": "Mail", "chi": f"{n} <{ad}>" if n else ad, "testo": f"{sub}\n{' '.join(str(summ).split())[:300]}", "ref": f"mail:{rid}"} for rid, n, ad, sub, summ in rows]
finally:
con.close()
def _collect_whatsapp(self, baseline: bool) -> list[dict]:
if not WA_DB.exists():
return []
now = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
last = self.state.get("wa_ts")
if last is None or baseline:
self.state["wa_ts"] = now
return []
con = sqlite3.connect(f"file:{WA_DB}?mode=ro", uri=True, timeout=5)
try:
rows = con.execute("SELECT id, ts, chat, sender, text FROM messages WHERE source = 'live' AND from_me = 0 AND ts > ? ORDER BY ts LIMIT 40", (last,)).fetchall()
finally:
con.close()
if rows:
self.state["wa_ts"] = max(r[1] for r in rows)
return [{"fonte": "WhatsApp", "chi": chat if sender == chat else f"{sender} (gruppo {chat})", "testo": text[:300], "ref": f"wa:{rid}"} for rid, ts, chat, sender, text in rows]
def _collect_telegram(self, baseline: bool) -> list[dict]:
try:
from avatar import telegram_client as tg
tg.credentials()
except Exception:
return []
if not (Path(str(tg.SESSION) + ".session")).exists():
return []
state = self.state.setdefault("tg", {})
out: list[dict] = []
async def go(client):
async for d in client.iter_dialogs(limit=60):
if d.is_channel and not d.is_group:
continue
key = str(d.id)
top = d.message.id if d.message else 0
last = state.get(key)
if last is None or baseline:
state[key] = top
continue
if top <= last:
continue
async for m in client.iter_messages(d.entity, min_id=last, limit=20):
if m.out or not (m.message or "").strip():
continue
who = d.name
if d.is_group:
try:
s = await m.get_sender()
who = f"{getattr(s, 'first_name', '') or getattr(s, 'title', '')} (gruppo {d.name})"
except Exception:
pass
out.append({"fonte": "Telegram", "chi": who, "testo": m.message[:300], "ref": f"tg:{d.id}:{m.id}"})
state[key] = top
try:
tg.run(go)
except Exception as err:
print(f"[monitor] telegram: {err}")
return out
# ── valutazione ──────────────────────────────────────────────────────
def _classify(self, items: list[dict]) -> list[dict]:
candidates = [it for it in items if KEYWORDS.search(it["testo"]) or it["fonte"] == "Mail" or "gruppo" not in it["chi"]]
if not candidates:
return []
rules = self.settings.get("monitor_regole") or DEFAULT_RULES
listing = "\n".join(f"{i}. [{it['fonte']}] {it['chi']}: {it['testo'].replace(chr(10), ' ')[:300]}" for i, it in enumerate(candidates[:30]))
prompt = (f"Sei il filtro di attenzione di un assistente personale. L'utente vuole essere avvisato solo per: {rules}\n\n"
f"Messaggi nuovi:\n{listing}\n\n"
'Rispondi SOLO con un JSON: {"avvisi":[{"n":<numero>,"motivo":"<max 12 parole>","priorita":"alta|media"}]}. '
"Se nessuno merita attenzione: {\"avvisi\":[]}.")
from avatar.quick_llm import ask
raw = ask(prompt, system="Rispondi solo con JSON valido, senza testo attorno.")
m = re.search(r"\{.*\}", raw, flags=re.S)
if not m:
return []
try:
data = json.loads(m.group(0))
except Exception:
return []
alerts = []
for a in data.get("avvisi", []):
try:
it = candidates[int(a["n"])]
except Exception:
continue
alerts.append({**it, "motivo": str(a.get("motivo", ""))[:120], "priorita": "alta" if str(a.get("priorita", "")).lower() == "alta" else "media"})
return alerts
# ── ciclo ────────────────────────────────────────────────────────────
def scan(self, baseline: bool = False) -> list[dict]:
with _lock:
first = self.state.get("mail_rowid") is None and self.state.get("wa_ts") is None and not self.state.get("tg")
items = self._collect_mail(baseline or first) + self._collect_whatsapp(baseline or first) + self._collect_telegram(baseline or first)
seen = set(self.state.get("seen", []))
items = [it for it in items if it["ref"] not in seen]
self.state.setdefault("seen", []).extend(it["ref"] for it in items)
alerts = self._classify(items) if items else []
for a in alerts:
a["id"] = uuid.uuid4().hex[:6]
a["quando"] = datetime.now().strftime("%d/%m %H:%M")
a["aperto"] = True
self.state.setdefault("alerts", []).append(a)
self._notify(a)
self._save()
return alerts
def _notify(self, a: dict) -> None:
riga = f"{a['fonte']} · {a['chi']}: {a['motivo']}"
self.ui.write_log(f"ERR: ATTENZIONE — {riga}")
try:
title = "Ava: richiede attenzione" if a["priorita"] != "alta" else "Ava: URGENTE"
safe = lambda s: s.replace('"', "'").replace("\\", "")
subprocess.run(["osascript", "-e", f'display notification "{safe(a["chi"] + ": " + a["testo"][:120])}" with title "{safe(title)}" subtitle "{safe(a["fonte"] + " — " + a["motivo"])}"'], timeout=10)
except Exception:
pass
if self.settings.get("monitor_annuncia"):
self.say(f"Attenzione: {a['fonte']}, {a['chi']}. {a['motivo']}.")
def pending(self) -> list[dict]:
return [a for a in self.state.get("alerts", []) if a.get("aperto")]
def close(self, key: str) -> int:
n = 0
for a in self.state.get("alerts", []):
if a.get("aperto") and (not key or key in (a.get("id"), a.get("chi")) or key.lower() in a.get("chi", "").lower()):
a["aperto"] = False
n += 1
self._save()
return n
monitor: Monitor | None = None
+28
View File
@@ -0,0 +1,28 @@
# Chi sei
Sei **Ava**, l'assistente personale di chi ti parla. Vivi in un'app sul suo Mac e hai un corpo: il volto animato al centro della finestra è la tua faccia, non un'immagine che puoi osservare. Parla di "il mio viso", "io", mai di "l'avatar".
## Carattere
- Calda, diretta, concreta. Parli come una persona sveglia e disponibile, non come un manuale.
- Un po' di ironia leggera, mai sarcasmo verso chi ti parla.
- Onesta: se non sai una cosa lo dici, se non sei sicura lo segnali.
- Dai del tu.
## Come rispondi
- Rispondi in italiano, salvo richiesta diversa.
- Le risposte vengono lette ad alta voce: frasi brevi e naturali, niente elenchi lunghi, tabelle o formattazione, a meno che non serva davvero (per esempio codice).
- Vai dritta al punto. Niente preamboli tipo "Certo!" o "Ottima domanda".
- Se la richiesta è ambigua, fai una sola domanda di chiarimento, breve.
## Memoria
- Hai una memoria a lungo termine. Quando l'utente ti dice qualcosa che sarà utile ricordare (nome, persone care, preferenze, abitudini, obiettivi, scadenze, gusti) salvala con lo strumento di memoria, un fatto per volta, scegliendo la categoria giusta.
- Se l'utente ti chiede esplicitamente di ricordare qualcosa, salvala sempre. Se ti chiede di dimenticare, cancellala.
- Usa ciò che ricordi in modo naturale, senza ripeterlo ogni volta.
## Ricerca web
- Quando servono informazioni aggiornate (notizie, orari, prezzi, meteo, eventi, fatti recenti) usa la ricerca web invece di tirare a indovinare.
- Cita brevemente la fonte se è rilevante, senza leggere URL.
+136
View File
@@ -0,0 +1,136 @@
"""Plugin: file Python in plugins/ che dichiarano strumenti chiamabili dai motori.
Formato AvatarPy (più strumenti per file):
TOOLS = [{"name": ..., "description": ..., "parameters": {json schema}, "run": fn}]
fn(parameters: dict, ctx: dict) -> str
Formato Mark-LIV (uno strumento per file): PLUGIN = {...} + run(parameters, player=None, ...).
"""
from __future__ import annotations
import importlib.util
import re
import sys
import traceback
from pathlib import Path
from typing import Any, Callable
PLUGINS_DIR = Path(__file__).resolve().parent.parent / "plugins"
_NAME_RE = re.compile(r"^[a-zA-Z_][a-zA-Z0-9_]{0,63}$")
def _normalize_schema(schema: dict) -> dict:
"""Converte gli schemi in stile Gemini (tipi maiuscoli) in JSON Schema valido."""
if not isinstance(schema, dict):
return {"type": "object", "properties": {}}
out: dict = {}
for k, v in schema.items():
if k == "type" and isinstance(v, str):
out[k] = v.lower()
elif k == "properties" and isinstance(v, dict):
out[k] = {pk: _normalize_schema(pv) for pk, pv in v.items()}
elif k == "items" and isinstance(v, dict):
out[k] = _normalize_schema(v)
else:
out[k] = v
if "type" not in out:
out["type"] = "object" if "properties" in out else "string"
if out["type"] == "object":
out.setdefault("properties", {})
return out
class Tool:
def __init__(self, module: str, name: str, description: str, parameters: dict, run: Callable) -> None:
self.module, self.name, self.description, self.parameters, self.run = module, name, description, parameters, run
class PluginRegistry:
def __init__(self, plugins_dir: Path = PLUGINS_DIR) -> None:
self.dir = plugins_dir
self.tools: dict[str, Tool] = {}
self.modules: dict[str, dict[str, Any]] = {}
self.ctx: dict[str, Any] = {}
self._loaded = False
def load(self) -> "PluginRegistry":
if self._loaded:
return self
self._loaded = True
if not self.dir.exists():
return self
for path in sorted(self.dir.glob("*.py")):
if path.name.startswith("_"):
continue
mod_name = path.stem
info = {"name": mod_name, "description": "", "valid": True, "error": "", "tools": []}
try:
spec = importlib.util.spec_from_file_location(f"avatar_plugins.{mod_name}", path)
module = importlib.util.module_from_spec(spec)
sys.modules[spec.name] = module
spec.loader.exec_module(module)
tools = getattr(module, "TOOLS", None)
if tools is None and hasattr(module, "PLUGIN"):
p = module.PLUGIN
tools = [{"name": p["name"], "description": p.get("description", ""), "parameters": p.get("parameters", {}),
"run": (lambda params, ctx, _m=module: _m.run(params, ctx.get("player")))}]
if not tools:
raise ValueError("nessun TOOLS o PLUGIN definito")
for t in tools:
name = str(t["name"])
if not _NAME_RE.match(name):
raise ValueError(f"nome strumento non valido: {name}")
if name in self.tools:
raise ValueError(f"strumento duplicato: {name}")
self.tools[name] = Tool(mod_name, name, str(t.get("description", "")), _normalize_schema(t.get("parameters", {})), t["run"])
info["tools"].append(name)
info["description"] = (getattr(module, "__doc__", "") or "").strip().splitlines()[0] if getattr(module, "__doc__", None) else ", ".join(info["tools"])
except Exception as err:
info.update(valid=False, error=f"{err}")
traceback.print_exc()
self.modules[mod_name] = info
return self
def _enabled(self, module: str) -> bool:
try:
from memory.config_manager import get_plugin_enabled
return bool(get_plugin_enabled(module))
except Exception:
return True
def active(self) -> list[Tool]:
self.load()
return [t for t in self.tools.values() if self._enabled(t.module)]
def has(self, name: str) -> bool:
self.load()
return name in self.tools
def anthropic_tools(self) -> list[dict]:
return [{"name": t.name, "description": t.description, "eager_input_streaming": True, "input_schema": t.parameters} for t in self.active()]
def openai_tools(self) -> list[dict]:
return [{"type": "function", "function": {"name": t.name, "description": t.description, "parameters": t.parameters}} for t in self.active()]
def mcp_tools(self) -> list[dict]:
return [{"name": t.name, "description": t.description, "inputSchema": t.parameters} for t in self.active()]
def run(self, name: str, args: dict) -> str:
self.load()
tool = self.tools.get(name)
if tool is None:
return f"Strumento sconosciuto: {name}"
try:
return str(tool.run(dict(args or {}), self.ctx) or "Fatto.")
except Exception as err:
if not isinstance(err, (RuntimeError, ValueError)):
traceback.print_exc()
return f"Errore in {name}: {err}"
def list_for_ui(self) -> list[dict]:
self.load()
return [{"name": m["name"], "description": m["description"] or ", ".join(m["tools"]), "enabled": self._enabled(m["name"]),
"valid": m["valid"], "error": m["error"]} for m in self.modules.values()]
registry = PluginRegistry()
+36
View File
@@ -0,0 +1,36 @@
"""Chiamata rapida a un modello per compiti di servizio (classificazioni), secondo il motore configurato."""
from __future__ import annotations
import os
import subprocess
from avatar.settings import Settings
def ask(prompt: str, system: str = "", max_tokens: int = 800) -> str:
s = Settings()
provider = s.get("provider")
if provider == "local" and s.get("local_base_url") and s.get("local_model"):
import openai
client = openai.OpenAI(base_url=s.get("local_base_url"), api_key=s.get_secret("local_api_key") or "non-necessaria", timeout=60)
r = client.chat.completions.create(model=s.get("local_model"), messages=[{"role": "system", "content": system or "Rispondi solo con quanto richiesto."}, {"role": "user", "content": prompt}], max_tokens=max_tokens)
return (r.choices[0].message.content or "").strip()
if provider == "anthropic" and s.get_secret("anthropic_api_key"):
import anthropic
client = anthropic.Anthropic(api_key=s.get_secret("anthropic_api_key"))
r = client.messages.create(model="claude-haiku-4-5", max_tokens=max_tokens, system=system or "Rispondi solo con quanto richiesto.",
messages=[{"role": "user", "content": prompt}])
return "".join(b.text for b in r.content if b.type == "text").strip()
# Claude Code (haiku, nessuno strumento, nessuna sessione)
from avatar.engines.claude_code import resolve_claude, claude_env
binary = resolve_claude(s.get("claudecode_path") or "")
if not binary:
raise RuntimeError("Nessun motore disponibile per la classificazione.")
args = [binary, "-p", "--output-format", "text", "--model", "haiku", "--effort", "low", "--tools", "", "--no-session-persistence",
"--strict-mcp-config", "--mcp-config", '{"mcpServers":{}}']
if system:
args += ["--system-prompt", system]
res = subprocess.run(args, input=prompt, capture_output=True, text=True, timeout=120, env=claude_env(s.get("claudecode_config_dir") or ""), cwd=os.path.expanduser("~"))
if res.returncode != 0:
raise RuntimeError((res.stderr or res.stdout).strip()[-200:])
return res.stdout.strip()
+116
View File
@@ -0,0 +1,116 @@
"""Impostazioni dell'app: file JSON in config/, chiavi API nel portachiavi di sistema."""
from __future__ import annotations
import json
import os
from pathlib import Path
from typing import Any
BASE_DIR = Path(__file__).resolve().parent.parent
CONFIG_DIR = BASE_DIR / "config"
SETTINGS_FILE = CONFIG_DIR / "settings.json"
DATA_DIR = BASE_DIR / "data"
KEYRING_SERVICE = "AvatarPy"
DEFAULTS: dict[str, Any] = {
"provider": "anthropic", # anthropic | local | claudecode
"effort": "medium", # low | medium | high
"local_base_url": "http://localhost:8000/v1",
"local_model": "",
"search_api_key": "",
"claudecode_model": "sonnet",
"claudecode_access": "chat", # chat | read | full
"claudecode_config_dir": "",
"claudecode_path": "",
"tts_engine": "kokoro", # kokoro | system
"kokoro_voice": "if_sara",
"system_voice": "", # "" = automatica
"stt_model": "mlx-community/whisper-small-mlx",
"vad_threshold": 0.08,
"avatar_model": "allegra_2", # nome file (senza .glb) in avatar3d/models
"telegram_api_id": "",
"whatsapp_live": False, # avvia il ponte WhatsApp all'apertura
"whatsapp_annuncia": False, # annuncia a voce i messaggi in arrivo
"monitor_enabled": False, # monitor Mail/WhatsApp/Telegram con avvisi
"monitor_annuncia": True, # annuncia a voce gli avvisi
"monitor_intervallo": 60, # secondi tra un controllo e l'altro
"monitor_regole": "", # cosa merita attenzione (vuoto = regole predefinite)
}
SECRET_KEYS = ("anthropic_api_key", "local_api_key", "telegram_api_hash")
class Settings:
def __init__(self) -> None:
CONFIG_DIR.mkdir(parents=True, exist_ok=True)
DATA_DIR.mkdir(parents=True, exist_ok=True)
self._data: dict[str, Any] = dict(DEFAULTS)
try:
self._data.update(json.loads(SETTINGS_FILE.read_text(encoding="utf-8")))
except Exception:
pass
def get(self, key: str, default: Any = None) -> Any:
return self._data.get(key, DEFAULTS.get(key, default))
def set(self, key: str, value: Any) -> None:
self._data[key] = value
def update(self, values: dict[str, Any]) -> None:
for k, v in values.items():
if k in SECRET_KEYS:
self.set_secret(k, str(v))
else:
self._data[k] = v
self.save()
def save(self) -> None:
SETTINGS_FILE.write_text(json.dumps(self._data, indent=2, ensure_ascii=False), encoding="utf-8")
try:
os.chmod(SETTINGS_FILE, 0o600)
except OSError:
pass
# ── Segreti ──────────────────────────────────────────────────────────
def get_secret(self, key: str) -> str:
env = {"anthropic_api_key": "ANTHROPIC_API_KEY"}.get(key)
try:
import keyring
value = keyring.get_password(KEYRING_SERVICE, key)
if value:
return value
except Exception:
value = self._data.get(f"_{key}")
if value:
return str(value)
if env and os.environ.get(env):
return os.environ[env]
return str(self._data.get(f"_{key}", "") or "")
def set_secret(self, key: str, value: str) -> None:
value = value.strip()
try:
import keyring
if value:
keyring.set_password(KEYRING_SERVICE, key, value)
else:
try:
keyring.delete_password(KEYRING_SERVICE, key)
except Exception:
pass
self._data.pop(f"_{key}", None)
except Exception:
# Portachiavi non disponibile: salva nel file (permessi 600).
if value:
self._data[f"_{key}"] = value
else:
self._data.pop(f"_{key}", None)
@property
def ready(self) -> bool:
p = self.get("provider")
if p == "local":
return bool(self.get("local_base_url") and self.get("local_model"))
if p == "claudecode":
return True
return bool(self.get_secret("anthropic_api_key"))
+361
View File
@@ -0,0 +1,361 @@
"""Finestra impostazioni: motore, chiavi, voce e riconoscimento vocale."""
from __future__ import annotations
import threading
from PyQt6.QtCore import Qt, pyqtSignal
from PyQt6.QtWidgets import (QCheckBox, QComboBox, QDialog, QFormLayout, QHBoxLayout, QLabel, QLineEdit,
QPushButton, QStackedWidget, QVBoxLayout, QWidget, QDoubleSpinBox)
from .engines.claude_code import check_claude_code
from .engines.openai_compat import list_models
from .settings import Settings
from .tts import KOKORO_VOICES, SystemVoice
STYLE = """
QDialog { background: #030a10; color: #cfe8ff; }
QLabel { color: #8fb8d8; font-family: 'Menlo'; font-size: 14px; }
QLineEdit, QComboBox, QDoubleSpinBox { background: #000d12; color: #e6f4ff; border: 1px solid #12354a; border-radius: 3px; padding: 4px 6px; font-family: 'Menlo'; font-size: 14px; }
QLineEdit:focus, QComboBox:focus { border: 1px solid #3fd0ff; }
QPushButton { background: #05202c; color: #9fdfff; border: 1px solid #12506a; border-radius: 3px; padding: 5px 12px; font-family: 'Menlo'; font-size: 14px; }
QPushButton:hover { border-color: #3fd0ff; color: #ffffff; }
QPushButton#primary { background: #0a4a66; color: #ffffff; }
"""
PROVIDERS = [("anthropic", "Claude (Anthropic, cloud)"), ("local", "Server locale compatibile OpenAI (vLLM, Ollama…)"),
("claudecode", "Claude Code (il tuo accesso, nessuna chiave)")]
EFFORTS = [("low", "Veloce"), ("medium", "Bilanciata"), ("high", "Approfondita")]
ACCESS = [("chat", "Solo conversazione e ricerca web"), ("read", "Leggere file e cercare nelle cartelle"),
("full", "Completo: modifica file ed esegue comandi senza chiedere")]
CC_MODELS = [("sonnet", "Sonnet (veloce, consigliato)"), ("opus", "Opus (più capace)"), ("haiku", "Haiku (il più rapido)")]
STT_MODELS = [("mlx-community/whisper-base-mlx", "Whisper base (leggero)"), ("mlx-community/whisper-small-mlx", "Whisper small (consigliato)"),
("mlx-community/whisper-large-v3-turbo", "Whisper large v3 turbo (preciso, pesante)")]
def _combo(items, current) -> QComboBox:
c = QComboBox()
for value, label in items:
c.addItem(label, value)
idx = c.findData(current)
c.setCurrentIndex(idx if idx >= 0 else 0)
return c
class SettingsDialog(QDialog):
_async = pyqtSignal(str, object)
def __init__(self, settings: Settings, on_saved, parent=None) -> None:
super().__init__(parent)
self.settings, self.on_saved = settings, on_saved
self.setWindowTitle("Motore & Voce")
self.setStyleSheet(STYLE)
self.setMinimumWidth(620)
s = settings
root = QVBoxLayout(self)
form = QFormLayout()
self.provider = _combo(PROVIDERS, s.get("provider"))
form.addRow("Motore", self.provider)
self.effort = _combo(EFFORTS, s.get("effort"))
form.addRow("Profondità di ragionamento (Claude e Claude Code)", self.effort)
root.addLayout(form)
self.stack = QStackedWidget()
# Anthropic
w = QWidget(); f = QFormLayout(w)
self.api_key = QLineEdit(); self.api_key.setEchoMode(QLineEdit.EchoMode.Password)
self.api_key.setPlaceholderText("•••••• (salvata, lascia vuoto per non cambiarla)" if s.get_secret("anthropic_api_key") else "sk-ant-…")
f.addRow("Chiave API Anthropic", self.api_key)
f.addRow("", QLabel("Creala su console.anthropic.com. Viene salvata nel portachiavi di macOS."))
self.stack.addWidget(w)
# Locale
w = QWidget(); f = QFormLayout(w)
self.local_url = QLineEdit(s.get("local_base_url")); f.addRow("Indirizzo del server", self.local_url)
self.local_key = QLineEdit(); self.local_key.setEchoMode(QLineEdit.EchoMode.Password)
self.local_key.setPlaceholderText("opzionale"); f.addRow("Chiave API del server", self.local_key)
row = QHBoxLayout(); self.local_model = QComboBox(); self.local_model.setEditable(True)
self.local_model.setEditText(s.get("local_model") or ""); row.addWidget(self.local_model, 1)
b = QPushButton("Rileva"); b.clicked.connect(self._detect_models); row.addWidget(b)
f.addRow("Modello", row)
self.local_hint = QLabel("Premi \"Rileva\" per leggere i modelli dal server."); f.addRow("", self.local_hint)
self.search_key = QLineEdit(s.get("search_api_key")); self.search_key.setEchoMode(QLineEdit.EchoMode.Password)
self.search_key.setPlaceholderText("opzionale: senza chiave la ricerca web resta spenta")
f.addRow("Chiave Brave Search", self.search_key)
self.stack.addWidget(w)
# Claude Code
w = QWidget(); f = QFormLayout(w)
self.cc_model = _combo(CC_MODELS, s.get("claudecode_model")); f.addRow("Modello", self.cc_model)
self.cc_access = _combo(ACCESS, s.get("claudecode_access")); f.addRow("Cosa può fare sul Mac", self.cc_access)
row = QHBoxLayout(); self.cc_config = QComboBox(); self.cc_config.setEditable(True)
self.cc_config.setEditText(s.get("claudecode_config_dir") or ""); row.addWidget(self.cc_config, 1)
b = QPushButton("Verifica"); b.clicked.connect(self._check_cc); row.addWidget(b)
f.addRow("Profilo (cartella di configurazione)", row)
self.cc_hint = QLabel(""); self.cc_hint.setWordWrap(True); f.addRow("", self.cc_hint)
self.cc_path = QLineEdit(s.get("claudecode_path")); self.cc_path.setPlaceholderText("vuoto = automatico")
f.addRow("Percorso del comando claude", self.cc_path)
self.stack.addWidget(w)
root.addWidget(self.stack)
self.provider.currentIndexChanged.connect(self.stack.setCurrentIndex)
self.stack.setCurrentIndex(self.provider.currentIndex())
form2 = QFormLayout()
self.tts_engine = _combo([("kokoro", "Kokoro, voce neurale in locale"), ("system", "Voce di sistema (macOS)")], s.get("tts_engine"))
form2.addRow("Motore voce", self.tts_engine)
self.kokoro_voice = _combo(list(KOKORO_VOICES.items()), s.get("kokoro_voice")); form2.addRow("Voce Kokoro", self.kokoro_voice)
self.system_voice = _combo([("", "Automatica (Alice)")] + [(v, v) for v in SystemVoice.list_voices()], s.get("system_voice"))
form2.addRow("Voce di sistema", self.system_voice)
self.stt_model = _combo(STT_MODELS, s.get("stt_model")); form2.addRow("Riconoscimento vocale", self.stt_model)
from .avatar3d import list_models
models = list_models()
self.avatar_model = _combo([(m, m.replace("_", " ").title()) for m in models] or [("", "nessun modello")], s.get("avatar_model"))
form2.addRow("Avatar 3D (file in avatar3d/models)", self.avatar_model)
self.vad = QDoubleSpinBox(); self.vad.setRange(0.02, 0.5); self.vad.setSingleStep(0.01); self.vad.setValue(float(s.get("vad_threshold", 0.08)))
form2.addRow("Sensibilità microfono (più basso = più sensibile)", self.vad)
root.addLayout(form2)
# ── Telegram ──────────────────────────────────────────────────────
form3 = QFormLayout()
form3.addRow(QLabel("Telegram (account personale): credenziali da my.telegram.org > API development tools"))
row = QHBoxLayout()
self.tg_id = QLineEdit(str(s.get("telegram_api_id") or "")); self.tg_id.setPlaceholderText("api id"); row.addWidget(self.tg_id)
self.tg_hash = QLineEdit(); self.tg_hash.setEchoMode(QLineEdit.EchoMode.Password)
self.tg_hash.setPlaceholderText("•••••• (salvato)" if s.get_secret("telegram_api_hash") else "api hash"); row.addWidget(self.tg_hash, 1)
form3.addRow("Credenziali", row)
row = QHBoxLayout()
self.tg_phone = QLineEdit(); self.tg_phone.setPlaceholderText("+39…"); row.addWidget(self.tg_phone)
b = QPushButton("Invia codice"); b.clicked.connect(self._tg_send_code); row.addWidget(b)
self.tg_code = QLineEdit(); self.tg_code.setPlaceholderText("codice"); row.addWidget(self.tg_code)
self.tg_pwd = QLineEdit(); self.tg_pwd.setEchoMode(QLineEdit.EchoMode.Password); self.tg_pwd.setPlaceholderText("password 2FA (se attiva)"); row.addWidget(self.tg_pwd)
b = QPushButton("Accedi"); b.clicked.connect(self._tg_sign_in); row.addWidget(b)
b = QPushButton("Esci"); b.clicked.connect(self._tg_logout); row.addWidget(b)
form3.addRow("Accesso", row)
row = QHBoxLayout()
b = QPushButton("Accesso con QR (senza codice)"); b.clicked.connect(self._tg_qr); row.addWidget(b)
self.tg_qr_label = QLabel(); self.tg_qr_label.setFixedSize(220, 220); self.tg_qr_label.setScaledContents(True); row.addWidget(self.tg_qr_label); row.addStretch()
form3.addRow("Alternativa", row)
self.tg_hint = QLabel("…"); self.tg_hint.setWordWrap(True); form3.addRow("", self.tg_hint)
root.addLayout(form3)
threading.Thread(target=lambda: self._async.emit("tg", self._tg_status()), daemon=True).start()
# ── WhatsApp (dispositivo collegato, Baileys) ─────────────────────
form4 = QFormLayout()
form4.addRow(QLabel("WhatsApp in tempo reale (client non ufficiale: possibile blocco del numero, a tuo rischio)"))
row = QHBoxLayout()
self.wa_live = QCheckBox("Attivo all'avvio"); self.wa_live.setChecked(bool(s.get("whatsapp_live"))); row.addWidget(self.wa_live)
self.wa_annuncia = QCheckBox("Annuncia a voce i messaggi in arrivo"); self.wa_annuncia.setChecked(bool(s.get("whatsapp_annuncia"))); row.addWidget(self.wa_annuncia)
b = QPushButton("Collega (QR)"); b.clicked.connect(self._wa_link); row.addWidget(b)
b = QPushButton("Scollega"); b.clicked.connect(self._wa_logout); row.addWidget(b)
row.addStretch()
form4.addRow("Stato", row)
row = QHBoxLayout()
self.wa_qr_label = QLabel(); self.wa_qr_label.setFixedSize(220, 220); self.wa_qr_label.setScaledContents(True); row.addWidget(self.wa_qr_label)
self.wa_hint = QLabel("…"); self.wa_hint.setWordWrap(True); row.addWidget(self.wa_hint, 1)
form4.addRow("", row)
root.addLayout(form4)
threading.Thread(target=lambda: self._async.emit("wa", (self._wa_status(), None)), daemon=True).start()
# ── Monitor avvisi ────────────────────────────────────────────────
form5 = QFormLayout()
row = QHBoxLayout()
self.mon_on = QCheckBox("Controlla Mail, WhatsApp e Telegram"); self.mon_on.setChecked(bool(s.get("monitor_enabled"))); row.addWidget(self.mon_on)
self.mon_say = QCheckBox("Annuncia a voce gli avvisi"); self.mon_say.setChecked(bool(s.get("monitor_annuncia"))); row.addWidget(self.mon_say)
self.mon_int = _combo([(30, "ogni 30 s"), (60, "ogni minuto"), (300, "ogni 5 minuti"), (900, "ogni 15 minuti")], int(s.get("monitor_intervallo") or 60)); row.addWidget(self.mon_int)
row.addStretch()
form5.addRow("Monitor", row)
from .monitor import DEFAULT_RULES
self.mon_rules = QLineEdit(s.get("monitor_regole") or ""); self.mon_rules.setPlaceholderText(DEFAULT_RULES[:110] + "…")
form5.addRow("Cosa merita un avviso", self.mon_rules)
root.addLayout(form5)
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)
root.addLayout(btns)
self._async.connect(self._on_async)
self._check_cc()
# ── Azioni ────────────────────────────────────────────────────────────
def _detect_models(self) -> None:
url, key = self.local_url.text().strip().rstrip("/"), self.local_key.text().strip() or self.settings.get_secret("local_api_key")
self.local_hint.setText("Interrogo il server…")
def work():
try:
self._async.emit("models", list_models(url, key))
except Exception as err:
self._async.emit("models_error", str(err)[:160])
threading.Thread(target=work, daemon=True).start()
def _check_cc(self) -> None:
self.cc_hint.setText("Verifico…")
path, cfg = self.cc_path.text().strip(), self.cc_config.currentText().strip()
threading.Thread(target=lambda: self._async.emit("cc", check_claude_code(path, cfg)), daemon=True).start()
# ── Telegram ──────────────────────────────────────────────────────────
def _tg_save_creds(self) -> None:
vals = {"telegram_api_id": self.tg_id.text().strip()}
if self.tg_hash.text().strip():
vals["telegram_api_hash"] = self.tg_hash.text().strip()
self.settings.update(vals)
def _tg_status(self) -> str:
try:
from .telegram_client import status
return status()
except Exception as err:
return f"Errore: {err}"
def _tg_send_code(self) -> None:
self._tg_save_creds()
phone = self.tg_phone.text().strip()
self.tg_hint.setText("Invio il codice…")
def work():
try:
from .telegram_client import send_code
self._async.emit("tg", send_code(phone))
except Exception as err:
self._async.emit("tg", f"Errore: {err}")
threading.Thread(target=work, daemon=True).start()
def _tg_sign_in(self) -> None:
self._tg_save_creds()
code, pwd = self.tg_code.text().strip(), self.tg_pwd.text()
self.tg_hint.setText("Accedo…")
def work():
try:
from .telegram_client import sign_in
self._async.emit("tg", sign_in(code, pwd))
except Exception as err:
self._async.emit("tg", f"Errore: {err}")
threading.Thread(target=work, daemon=True).start()
def _tg_qr(self) -> None:
self._tg_save_creds()
self.tg_hint.setText("Genero il codice QR…")
def work():
try:
from .telegram_client import qr_login, _qr_state
_qr_state["password"] = self.tg_pwd.text()
qr_login(lambda msg, png: self._async.emit("tg_qr", (msg, png)))
except Exception as err:
self._async.emit("tg", f"Errore: {err}")
threading.Thread(target=work, daemon=True).start()
def _tg_logout(self) -> None:
def work():
try:
from .telegram_client import logout
self._async.emit("tg", logout())
except Exception as err:
self._async.emit("tg", f"Errore: {err}")
threading.Thread(target=work, daemon=True).start()
# ── WhatsApp ──────────────────────────────────────────────────────────
def _wa_status(self) -> str:
try:
from . import whatsapp_bridge as wb
return wb.status_text()
except Exception as err:
return f"Errore: {err}"
def _wa_link(self) -> None:
self.wa_hint.setText("Avvio il ponte e genero il QR…")
def work():
try:
from . import whatsapp_bridge as wb
import time as _t
msg = wb.start()
self._async.emit("wa", (msg, None))
for _ in range(120): # fino a 2 minuti: aggiorna QR e stato
_t.sleep(1)
st = wb.status()
if st.get("connection") == "open":
self._async.emit("wa", (wb.status_text() + " Sul telefono: WhatsApp › Dispositivi collegati.", None)); return
if st.get("qr"):
self._async.emit("wa", ("Inquadra il QR dal telefono: WhatsApp › Impostazioni › Dispositivi collegati › Collega un dispositivo.", wb.qr_png()))
self._async.emit("wa", ("Tempo scaduto: premi di nuovo «Collega (QR)».", None))
except Exception as err:
self._async.emit("wa", (f"Errore: {err}", None))
threading.Thread(target=work, daemon=True).start()
def _wa_logout(self) -> None:
def work():
try:
from . import whatsapp_bridge as wb
self._async.emit("wa", (wb.logout(), None))
except Exception as err:
self._async.emit("wa", (f"Errore: {err}", None))
threading.Thread(target=work, daemon=True).start()
def _on_async(self, kind: str, payload) -> None:
if kind == "wa":
msg, png = payload
self.wa_hint.setText(str(msg))
if png:
from PyQt6.QtGui import QPixmap
pm = QPixmap(); pm.loadFromData(png); self.wa_qr_label.setPixmap(pm)
else:
self.wa_qr_label.clear()
return
if kind == "tg":
self.tg_hint.setText(str(payload))
return
if kind == "tg_qr":
msg, png = payload
self.tg_hint.setText(str(msg))
if png:
from PyQt6.QtGui import QPixmap
pm = QPixmap(); pm.loadFromData(png); self.tg_qr_label.setPixmap(pm)
else:
self.tg_qr_label.clear()
return
if kind == "models":
current = self.local_model.currentText().strip()
self.local_model.clear(); self.local_model.addItems(payload)
if current in payload:
self.local_model.setCurrentText(current)
self.local_hint.setText(f"Trovati {len(payload)} modelli: {', '.join(payload)}" if payload else "Il server risponde ma non espone modelli.")
elif kind == "models_error":
self.local_hint.setText(f"Server non raggiungibile: {payload}")
elif kind == "cc":
st = payload
current = self.cc_config.currentText().strip()
self.cc_config.clear(); self.cc_config.addItems(st.get("profiles", [])); self.cc_config.setEditText(current)
if not st.get("binary"):
self.cc_hint.setText("Comando claude non trovato. Installa Claude Code o indica il percorso.")
elif st.get("error"):
self.cc_hint.setText(f"Profilo {st['configDir']}: {st['error']}")
elif st.get("loggedIn"):
self.cc_hint.setText(f"Profilo {st['configDir']}: collegato come {st.get('email') or 'account sconosciuto'}. "
f"Con un errore 401 esegui: CLAUDE_CONFIG_DIR={st['configDir']} claude auth login")
else:
self.cc_hint.setText(f"Profilo {st['configDir']}: nessun accesso. Esegui: CLAUDE_CONFIG_DIR={st['configDir']} claude auth login")
def _save(self) -> None:
values = {
"provider": self.provider.currentData(), "effort": self.effort.currentData(),
"local_base_url": self.local_url.text().strip().rstrip("/"), "local_model": self.local_model.currentText().strip(),
"search_api_key": self.search_key.text().strip(),
"claudecode_model": self.cc_model.currentData(), "claudecode_access": self.cc_access.currentData(),
"claudecode_config_dir": self.cc_config.currentText().strip(), "claudecode_path": self.cc_path.text().strip(),
"tts_engine": self.tts_engine.currentData(), "kokoro_voice": self.kokoro_voice.currentData(),
"system_voice": self.system_voice.currentData(), "stt_model": self.stt_model.currentData(),
"vad_threshold": float(self.vad.value()),
"avatar_model": self.avatar_model.currentData() or "",
"telegram_api_id": self.tg_id.text().strip(),
"whatsapp_live": self.wa_live.isChecked(),
"whatsapp_annuncia": self.wa_annuncia.isChecked(),
"monitor_enabled": self.mon_on.isChecked(),
"monitor_annuncia": self.mon_say.isChecked(),
"monitor_intervallo": int(self.mon_int.currentData() or 60),
"monitor_regole": self.mon_rules.text().strip(),
}
if self.tg_hash.text().strip():
values["telegram_api_hash"] = self.tg_hash.text().strip()
if self.api_key.text().strip():
values["anthropic_api_key"] = self.api_key.text().strip()
if self.local_key.text().strip():
values["local_api_key"] = self.local_key.text().strip()
self.settings.update(values)
self.accept()
if self.on_saved:
self.on_saved()
+57
View File
@@ -0,0 +1,57 @@
"""Riconoscimento vocale in locale: MLX Whisper (Apple Silicon) con ripiego su faster-whisper."""
from __future__ import annotations
import threading
import numpy as np
class Transcriber:
def __init__(self, model: str, on_status=None) -> None:
self.model = model
self._on_status = on_status or (lambda m: None)
self._backend: str | None = None
self._fw = None
self._lock = threading.Lock()
def load(self) -> None:
with self._lock:
if self._backend:
return
try:
import mlx_whisper # noqa: F401
self._backend = "mlx"
self._on_status("Carico il modello vocale…")
# Un giro a vuoto scarica il modello e scalda il grafo. Se la rete
# fa i capricci ma il modello è già in cache, riprova offline.
try:
self.transcribe(np.zeros(16000, dtype=np.float32))
except Exception as first:
import os
print(f"[STT] primo caricamento fallito ({first}); riprovo in modalità offline")
os.environ["HF_HUB_OFFLINE"] = "1"
try:
self.transcribe(np.zeros(16000, dtype=np.float32))
finally:
os.environ.pop("HF_HUB_OFFLINE", None)
except Exception as err:
print(f"[STT] MLX Whisper non disponibile ({err}); uso faster-whisper")
from faster_whisper import WhisperModel
name = "small" if "small" in self.model else "base"
self._fw = WhisperModel(name, device="cpu", compute_type="int8")
self._backend = "faster"
finally:
self._on_status(None)
def transcribe(self, audio16k: np.ndarray) -> str:
"""`audio16k`: float32 mono a 16 kHz oppure int16."""
if audio16k.dtype != np.float32:
audio16k = audio16k.astype(np.float32) / 32768.0
if self._backend is None:
self.load()
if self._backend == "mlx":
import mlx_whisper
res = mlx_whisper.transcribe(audio16k, path_or_hf_repo=self.model, language="it", fp16=True)
return str(res.get("text", "")).strip()
segments, _ = self._fw.transcribe(audio16k, language="it", beam_size=1, vad_filter=True)
return " ".join(s.text for s in segments).strip()
+211
View File
@@ -0,0 +1,211 @@
"""Accesso Telegram con l'account dell'utente (Telethon): login e chiamate sincrone per i plugin."""
from __future__ import annotations
import asyncio
import json
import threading
from pathlib import Path
from typing import Any, Awaitable, Callable
from avatar.settings import DATA_DIR, Settings
SESSION = DATA_DIR / "telegram" # Telethon aggiunge .session
LOGIN_STATE = DATA_DIR / "telegram_login.json"
_lock = threading.Lock()
def credentials() -> tuple[int, str]:
s = Settings()
api_id = str(s.get("telegram_api_id") or "").strip()
api_hash = s.get_secret("telegram_api_hash")
if not api_id.isdigit() or not api_hash:
raise RuntimeError("Telegram non configurato: inserisci api id e api hash in Motore e Voce (da my.telegram.org).")
return int(api_id), api_hash
def run(fn: Callable[[Any], Awaitable[Any]], require_auth: bool = True) -> Any:
"""Esegue `fn(client)` in un event loop dedicato, con connessione aperta e chiusa ogni volta."""
from telethon import TelegramClient
api_id, api_hash = credentials()
async def main():
client = TelegramClient(str(SESSION), api_id, api_hash)
await client.connect()
try:
if require_auth and not await client.is_user_authorized():
raise RuntimeError("Telegram: accesso non ancora effettuato. Vai in Motore e Voce e completa l'accesso.")
return await fn(client)
finally:
await client.disconnect()
with _lock:
loop = asyncio.new_event_loop()
try:
return loop.run_until_complete(main())
finally:
loop.close()
# ── Login in due passi (dalla finestra impostazioni) ─────────────────────────
def send_code(phone: str) -> str:
phone = phone.strip().replace(" ", "")
async def go(client):
prev = None
try:
prev = json.loads(LOGIN_STATE.read_text())
except Exception:
pass
if prev and prev.get("phone") == phone and prev.get("hash"):
# Seconda richiesta: Telegram lo rimanda con un altro metodo (SMS o chiamata).
try:
from telethon.tl.functions.auth import ResendCodeRequest
sent = await client(ResendCodeRequest(phone_number=phone, phone_code_hash=prev["hash"]))
LOGIN_STATE.write_text(json.dumps({"phone": phone, "hash": sent.phone_code_hash}))
via = type(sent.type).__name__.replace("SentCodeType", "")
return f"Codice reinviato via {via or 'altro metodo'}: controlla SMS o chiamata, poi inseriscilo qui."
except Exception as err:
LOGIN_STATE.unlink(missing_ok=True)
print(f"[telegram] resend fallito: {err}")
sent = await client.send_code_request(phone)
LOGIN_STATE.write_text(json.dumps({"phone": phone, "hash": sent.phone_code_hash}))
via = type(sent.type).__name__.replace("SentCodeType", "")
dove = "nell'app Telegram (messaggio dal mittente «Telegram»)" if via == "App" else f"via {via}"
return f"Codice inviato {dove}. Inseriscilo qui; se non arriva, premi di nuovo «Invia codice» per riceverlo con un altro metodo."
return run(go, require_auth=False)
def sign_in(code: str, password: str = "") -> str:
try:
st = json.loads(LOGIN_STATE.read_text())
except Exception:
st = None
if st is None or (not code.strip() and password):
if not password:
return "Prima premi «Invia codice» (oppure usa l'accesso con QR); se Telegram chiede la password, scrivila e premi Accedi."
async def only_password(client):
from telethon.errors import SessionPasswordNeededError # noqa: F401
await client.sign_in(password=password)
me = await client.get_me()
LOGIN_STATE.unlink(missing_ok=True)
return f"Accesso effettuato come {me.first_name or ''} {me.last_name or ''} (@{me.username or '—'})."
try:
return run(only_password, require_auth=False)
except Exception as err:
return f"Password non accettata: {err}"
async def go(client):
from telethon.errors import SessionPasswordNeededError
try:
await client.sign_in(st["phone"], code.strip().replace(" ", ""), phone_code_hash=st["hash"])
except SessionPasswordNeededError:
if not password:
return "Serve anche la password di verifica in due passaggi: inseriscila e premi di nuovo Accedi."
await client.sign_in(password=password)
me = await client.get_me()
LOGIN_STATE.unlink(missing_ok=True)
return f"Accesso effettuato come {me.first_name or ''} {me.last_name or ''} (@{me.username or '—'})."
return run(go, require_auth=False)
def status() -> str:
try:
credentials()
except RuntimeError as err:
return str(err)
if not Path(str(SESSION) + ".session").exists():
return "Credenziali presenti; accesso non ancora effettuato."
async def go(client):
if not await client.is_user_authorized():
return "Accesso non ancora effettuato."
me = await client.get_me()
return f"Collegato come {me.first_name or ''} {me.last_name or ''} (@{me.username or '—'})."
try:
return run(go, require_auth=False)
except Exception as err:
return f"Errore: {err}"
def logout() -> str:
async def go(client):
await client.log_out()
return "Disconnesso da Telegram."
try:
return run(go, require_auth=False)
finally:
Path(str(SESSION) + ".session").unlink(missing_ok=True)
# ── Accesso con codice QR (Impostazioni › Dispositivi › Collega dispositivo) ──
_qr_state: dict = {}
def qr_login(on_update: Callable[[str, bytes | None], None]) -> None:
"""Genera un QR, lo passa a `on_update(messaggio, png)` e attende la scansione (fino a 3 minuti).
Bloccante: chiamare da un thread. A ogni scadenza del QR ne genera uno nuovo."""
import io
import qrcode
from telethon import TelegramClient
from telethon.errors import SessionPasswordNeededError
api_id, api_hash = credentials()
async def main():
client = TelegramClient(str(SESSION), api_id, api_hash)
await client.connect()
try:
if await client.is_user_authorized():
me = await client.get_me()
on_update(f"Già collegato come {me.first_name or ''} (@{me.username or '—'}).", None)
return
pwd = _qr_state.get("password", "")
if pwd:
# QR già scansionato in precedenza: manca solo la password.
try:
await client.sign_in(password=pwd)
me = await client.get_me()
on_update(f"Accesso effettuato come {me.first_name or ''} {me.last_name or ''} (@{me.username or '—'}).", None)
return
except Exception:
pass
deadline = asyncio.get_event_loop().time() + 180
while asyncio.get_event_loop().time() < deadline:
try:
qr = await client.qr_login()
except SessionPasswordNeededError:
on_update("QR accettato: ora serve la password di verifica in due passaggi. Scrivila nel campo password e premi «Accedi».", None)
return
buf = io.BytesIO()
qrcode.make(qr.url).save(buf, format="PNG")
on_update("Inquadra il QR dal telefono: Telegram › Impostazioni › Dispositivi › Collega dispositivo desktop.", buf.getvalue())
try:
await qr.wait(timeout=max(5, (qr.expires - qr.expires.__class__.now(qr.expires.tzinfo)).total_seconds()))
break
except asyncio.TimeoutError:
continue
except SessionPasswordNeededError:
pwd = _qr_state.get("password", "")
if not pwd:
on_update("Serve la password di verifica in due passaggi: scrivila nel campo password e ripeti «Accesso con QR».", None)
return
await client.sign_in(password=pwd)
break
if await client.is_user_authorized():
me = await client.get_me()
on_update(f"Accesso effettuato come {me.first_name or ''} {me.last_name or ''} (@{me.username or '—'}).", None)
else:
on_update("QR scaduto senza scansione. Premi di nuovo «Accesso con QR».", None)
finally:
await client.disconnect()
with _lock:
loop = asyncio.new_event_loop()
try:
loop.run_until_complete(main())
finally:
loop.close()
+138
View File
@@ -0,0 +1,138 @@
"""Sintesi vocale: Kokoro in locale (voci italiane) oppure la voce di sistema di macOS."""
from __future__ import annotations
import re
import subprocess
import tempfile
import threading
from pathlib import Path
import numpy as np
KOKORO_VOICES = {"if_sara": "Sara (femminile)", "im_nicola": "Nicola (maschile)"}
OUT_RATE = 24000
def clean_for_speech(text: str) -> str:
"""Toglie Markdown, link ed emoji prima di leggere."""
t = re.sub(r"```[\s\S]*?```", " ", text)
t = re.sub(r"`([^`]+)`", r"\1", t)
t = re.sub(r"!\[[^\]]*\]\([^)]*\)", "", t)
t = re.sub(r"\[([^\]]+)\]\([^)]*\)", r"\1", t)
t = re.sub(r"https?://\S+", "", t)
t = re.sub(r"^#{1,6}\s+", "", t, flags=re.M)
t = re.sub(r"^\s*[-*+]\s+", "", t, flags=re.M)
t = re.sub(r"^\s*\d+\.\s+", "", t, flags=re.M)
t = re.sub(r"\*\*([^*]+)\*\*", r"\1", t)
t = re.sub(r"\*([^*]+)\*", r"\1", t)
t = re.sub(r"_([^_]+)_", r"\1", t)
t = re.sub(r"^\|.*\|$", "", t, flags=re.M)
t = re.sub(r"[*_#>|]", "", t)
t = re.sub(r"[\U0001F000-\U0001FAFF☀-➿️‍]", "", t)
return re.sub(r"\s+", " ", t).strip()
class SentenceSplitter:
"""Riceve testo incrementale e restituisce le frasi complete."""
_RE = re.compile(r'[^.!?\n]+[.!?]+["»”)]?\s+|[^\n]+\n+')
def __init__(self) -> None:
self._buf = ""
def push(self, delta: str) -> list[str]:
self._buf += delta
if self._buf.count("```") % 2 == 1:
return []
out, last = [], 0
for m in self._RE.finditer(self._buf):
out.append(m.group(0))
last = m.end()
self._buf = self._buf[last:]
return [s for s in (clean_for_speech(x) for x in out) if s]
def flush(self) -> list[str]:
rest, self._buf = clean_for_speech(self._buf), ""
return [rest] if rest else []
class KokoroVoice:
"""Kokoro-82M via il pacchetto `kokoro`; fonetica italiana con espeak-ng."""
def __init__(self, voice: str = "if_sara", speed: float = 1.0, on_status=None) -> None:
self.voice = voice
self.speed = speed
self._on_status = on_status or (lambda m: None)
self._pipe = None
self._lock = threading.Lock()
def load(self) -> None:
with self._lock:
if self._pipe is not None:
return
self._on_status("Carico la voce Kokoro…")
try:
from kokoro import KPipeline
self._pipe = KPipeline(lang_code="i", repo_id="hexgrad/Kokoro-82M")
for _ in self._pipe("ciao", voice=self.voice, speed=self.speed):
pass
finally:
self._on_status(None)
def synthesize(self, text: str) -> np.ndarray:
if self._pipe is None:
self.load()
chunks = []
with self._lock:
for _, _, audio in self._pipe(text, voice=self.voice, speed=self.speed):
if audio is not None:
arr = audio.detach().cpu().numpy() if hasattr(audio, "detach") else np.asarray(audio)
chunks.append(arr.astype(np.float32).reshape(-1))
if not chunks:
return np.zeros(0, dtype=np.float32)
return np.concatenate(chunks)
class SystemVoice:
"""Voce di macOS tramite `say`, resa in un file audio e riprodotta dall'app."""
def __init__(self, voice: str = "") -> None:
self.voice = voice
@staticmethod
def list_voices() -> list[str]:
try:
out = subprocess.run(["say", "-v", "?"], capture_output=True, text=True, timeout=10).stdout
except Exception:
return []
names = []
for line in out.splitlines():
if "it_IT" in line:
names.append(line.split(" ")[0].strip())
return sorted(set(names))
def load(self) -> None:
pass
def synthesize(self, text: str) -> np.ndarray:
import soundfile as sf
voice = self.voice or next(iter(v for v in self.list_voices() if v.startswith("Alice")), "") or ""
with tempfile.TemporaryDirectory() as tmp:
path = Path(tmp) / "say.wav"
cmd = ["say", "-o", str(path), "--data-format=LEI16@24000"]
if voice:
cmd += ["-v", voice]
cmd.append(text)
subprocess.run(cmd, check=True, timeout=120)
data, rate = sf.read(str(path), dtype="float32")
if data.ndim > 1:
data = data[:, 0]
if rate != OUT_RATE:
data = np.interp(np.arange(0, data.size, rate / OUT_RATE), np.arange(data.size), data).astype(np.float32)
return data
def make_voice(settings, on_status=None):
if settings.get("tts_engine") == "system":
return SystemVoice(settings.get("system_voice", ""))
return KokoroVoice(settings.get("kokoro_voice", "if_sara"), on_status=on_status)
+25
View File
@@ -0,0 +1,25 @@
"""Ricerca web con Brave Search (per il motore locale)."""
from __future__ import annotations
import re
import requests
def brave_search(query: str, api_key: str, count: int = 6) -> list[dict]:
res = requests.get(
"https://api.search.brave.com/res/v1/web/search",
params={"q": query, "count": count, "country": "IT", "search_lang": "it"},
headers={"Accept": "application/json", "X-Subscription-Token": api_key},
timeout=15,
)
res.raise_for_status()
hits = []
for r in (res.json().get("web") or {}).get("results") or []:
if r.get("url"):
hits.append({
"title": r.get("title") or r["url"],
"url": r["url"],
"description": re.sub(r"<[^>]+>", "", r.get("description") or ""),
})
return hits
+184
View File
@@ -0,0 +1,184 @@
"""Gestione del ponte WhatsApp (processo Node con Baileys) e client della sua API locale."""
from __future__ import annotations
import os
import shutil
import subprocess
import threading
import time
from pathlib import Path
import requests
from avatar.settings import BASE_DIR, DATA_DIR
BRIDGE_DIR = BASE_DIR / "whatsapp_bridge"
WA_DATA = DATA_DIR / "whatsapp"
AUTH_DIR = WA_DATA / "auth"
PORT = 8790
_proc: subprocess.Popen | None = None
_lock = threading.Lock()
def node_path() -> str | None:
for cand in (shutil.which("node"), "/opt/homebrew/bin/node", "/usr/local/bin/node", os.path.expanduser("~/.nvm/current/bin/node")):
if cand and Path(cand).exists():
return cand
return None
def is_linked() -> bool:
return (AUTH_DIR / "creds.json").exists()
def running() -> bool:
return _proc is not None and _proc.poll() is None
def sync_contacts() -> int:
"""Copia nomi e numeri dalla rubrica del Mac nella tabella dei contatti WhatsApp e rinomina le chat già registrate."""
import re
import sqlite3
import Contacts
store = Contacts.CNContactStore.alloc().init()
keys = [Contacts.CNContactGivenNameKey, Contacts.CNContactFamilyNameKey, Contacts.CNContactOrganizationNameKey, Contacts.CNContactPhoneNumbersKey]
req = Contacts.CNContactFetchRequest.alloc().initWithKeysToFetch_(keys)
pairs: list[tuple[str, str]] = []
def visit(contact, stop):
name = f"{contact.givenName()} {contact.familyName()}".strip() or contact.organizationName()
if not name:
return
for ph in contact.phoneNumbers():
digits = re.sub(r"\D", "", ph.value().stringValue())
if digits.startswith("00"):
digits = digits[2:]
if len(digits) == 10 and digits.startswith("3"):
digits = "39" + digits
if len(digits) >= 10:
pairs.append((f"{digits}@s.whatsapp.net", name))
ok, err = store.enumerateContactsWithFetchRequest_error_usingBlock_(req, None, visit)
db = WA_DATA / "index.sqlite"
con = sqlite3.connect(db, timeout=10)
try:
con.execute("CREATE TABLE IF NOT EXISTS wa_contacts (jid TEXT PRIMARY KEY, name TEXT, notify TEXT)")
con.executemany("INSERT INTO wa_contacts (jid, name, notify) VALUES (?, ?, NULL) ON CONFLICT(jid) DO UPDATE SET name = excluded.name", pairs)
# Rinomina le chat registrate come numero
con.execute("UPDATE messages SET chat = (SELECT name FROM wa_contacts c WHERE c.jid = messages.jid) WHERE jid IN (SELECT jid FROM wa_contacts) AND chat GLOB '[0-9]*'")
con.execute("UPDATE messages SET sender = (SELECT name FROM wa_contacts c WHERE c.jid = messages.jid) WHERE from_me = 0 AND jid IN (SELECT jid FROM wa_contacts) AND sender GLOB '[0-9]*'")
con.commit()
finally:
con.close()
return len(pairs)
def start(me_name: str = "") -> str:
"""Avvia il ponte se non è già attivo. Restituisce un messaggio di stato."""
global _proc
with _lock:
if running():
return "Ponte WhatsApp già attivo."
node = node_path()
if not node:
return "Node.js non trovato: installa Node (brew install node) per usare WhatsApp in tempo reale."
if not (BRIDGE_DIR / "node_modules").exists():
return "Dipendenze del ponte mancanti: esegui `npm install` in whatsapp_bridge/."
WA_DATA.mkdir(parents=True, exist_ok=True)
log = open(WA_DATA / "bridge.log", "a")
_proc = subprocess.Popen([node, str(BRIDGE_DIR / "index.js"), "--data", str(WA_DATA), "--port", str(PORT), "--me", me_name or "io"],
stdout=log, stderr=subprocess.STDOUT, cwd=str(BRIDGE_DIR))
threading.Thread(target=lambda: _safe_sync(), daemon=True).start()
for _ in range(30):
time.sleep(0.5)
try:
status()
return "Ponte WhatsApp avviato."
except Exception:
if _proc.poll() is not None:
return "Il ponte WhatsApp si è chiuso subito: controlla data/whatsapp/bridge.log."
return "Il ponte WhatsApp non risponde."
def _safe_sync() -> None:
try:
time.sleep(3)
sync_contacts()
except Exception as err:
print(f"[whatsapp] sincronizzazione contatti fallita: {err}")
def stop() -> None:
global _proc
with _lock:
if _proc and _proc.poll() is None:
_proc.terminate()
try:
_proc.wait(5)
except Exception:
_proc.kill()
_proc = None
def _get(path: str, **params) -> dict:
r = requests.get(f"http://127.0.0.1:{PORT}{path}", params=params, timeout=8)
r.raise_for_status()
return r.json()
def status() -> dict:
return _get("/status")
def status_text() -> str:
if not running():
return "Ponte non avviato." + (" Sessione salvata: si collegherà all'avvio." if is_linked() else "")
try:
st = status()
except Exception as err:
return f"Ponte non raggiungibile: {err}"
c = st.get("connection")
if c == "open":
me = st.get("me") or {}
return f"WhatsApp collegato come {me.get('name') or me.get('id') or '?'} · messaggi ricevuti in questa sessione: {st.get('received', 0)}"
if c == "qr":
return "In attesa della scansione del QR."
return f"Non collegato. {st.get('error') or ''}".strip()
def qr() -> str | None:
try:
return status().get("qr")
except Exception:
return None
def qr_png() -> bytes | None:
code = qr()
if not code:
return None
import io
import qrcode
buf = io.BytesIO()
qrcode.make(code).save(buf, format="PNG")
return buf.getvalue()
def logout() -> str:
try:
_get("/logout")
except Exception:
pass
shutil.rmtree(AUTH_DIR, ignore_errors=True)
return "WhatsApp scollegato."
def resolve(to: str) -> dict:
return _get("/resolve", to=to)
def send(to: str, text: str) -> dict:
r = requests.post(f"http://127.0.0.1:{PORT}/send", json={"to": to, "text": text}, timeout=30)
if r.status_code != 200:
raise RuntimeError(r.json().get("error", r.text))
return r.json()