Download server.py from PascalF53/jarvis-local: direct link, hf CLI and curl.
- Browser
- Download file 58 kB
-
https://huggingface.co/PascalF53/jarvis-local/resolve/main/server.py
- Command line
-
hf download hf://PascalF53/jarvis-local/server.py
-
curl -L -o server.py https://huggingface.co/PascalF53/jarvis-local/resolve/main/server.py
58 kB
| #!/usr/bin/env python3 | |
| """JARVIS Local: 100 % local voice assistant. | |
| micro (browser) --WS PCM 16 kHz--> Silero VAD --> faster-whisper | |
| --> Qwen3 via Ollama (tool calling) --> Kokoro TTS --WS PCM 24 kHz--> speakers | |
| One small FastAPI server: | |
| GET / the futuristic UI (orb + live task panels) | |
| WS /ws the voice loop: mic audio in, speech audio + UI events out | |
| POST /api/task spawns a background agent task (Claude Code, backed by | |
| Ollama by default so it stays local) | |
| GET /api/task/{id} a task's status and output | |
| Nothing leaves the machine unless you set JARVIS_AGENT_BACKEND=anthropic. | |
| """ | |
| import asyncio | |
| import glob | |
| import json | |
| import os | |
| import re | |
| import shutil | |
| import site | |
| import subprocess | |
| import threading | |
| import time | |
| import uuid | |
| from collections import deque | |
| from concurrent.futures import ThreadPoolExecutor | |
| from pathlib import Path | |
| import httpx | |
| import numpy as np | |
| from fastapi import FastAPI, HTTPException, WebSocket, WebSocketDisconnect | |
| from fastapi.responses import FileResponse | |
| from pydantic import BaseModel | |
| ROOT = Path(__file__).parent | |
| MODELS = ROOT / "models" | |
| app = FastAPI(title="JARVIS Local") | |
| # CUDA libs come from the nvidia-* pip wheels: CUDA 12 for faster-whisper | |
| # (nvidia/*/bin), CUDA 13 for onnxruntime-gpu / Kokoro (nvidia/cu13/bin/x86_64). | |
| for _sp in site.getsitepackages(): | |
| for _d in (glob.glob(os.path.join(_sp, "nvidia", "*", "bin")) | |
| + glob.glob(os.path.join(_sp, "nvidia", "*", "bin", "x86_64"))): | |
| os.add_dll_directory(_d) | |
| os.environ["PATH"] = _d + os.pathsep + os.environ["PATH"] | |
| # ---------------------------------------------------------------- config | |
| def load_env(): | |
| env_file = ROOT / ".env" | |
| if env_file.exists(): | |
| for line in env_file.read_text(encoding="utf-8").splitlines(): | |
| line = line.strip() | |
| if line and not line.startswith("#") and "=" in line: | |
| k, v = line.split("=", 1) | |
| os.environ.setdefault(k.strip(), v.strip().strip('"').strip("'")) | |
| load_env() | |
| OLLAMA_URL = os.environ.get("OLLAMA_URL", "http://127.0.0.1:11434").rstrip("/") | |
| # MoE: only ~3B active parameters per token, ~4x faster than the dense 32B. | |
| LLM_MODEL = os.environ.get("JARVIS_LLM", "qwen3:30b-a3b-instruct-2507-q4_K_M") | |
| NUM_CTX = int(os.environ.get("JARVIS_NUM_CTX", "32768")) | |
| KEEP_ALIVE = os.environ.get("JARVIS_KEEP_ALIVE", "2h") | |
| LANGUAGE = os.environ.get("JARVIS_LANGUAGE", "français") | |
| # How JARVIS addresses you (empty = the classic "monsieur"). | |
| USER_NAME = os.environ.get("JARVIS_USER_NAME", "").strip() | |
| ADDRESS = USER_NAME or "monsieur" | |
| WHISPER_MODEL = os.environ.get("WHISPER_MODEL", "large-v3-turbo") | |
| WHISPER_DEVICE = os.environ.get("WHISPER_DEVICE", "cuda") | |
| WHISPER_COMPUTE = os.environ.get("WHISPER_COMPUTE", "int8_float16") | |
| WHISPER_LANG = os.environ.get("WHISPER_LANG", "fr") | |
| # Names Whisper should favour (companies, people, products), comma separated. | |
| HOTWORDS = " ".join(w.strip() for w in ["JARVIS", USER_NAME, *os.environ.get( | |
| "JARVIS_HOTWORDS", "").split(",")] if w.strip()) | |
| TTS_VOICE = os.environ.get("JARVIS_VOICE", "ff_siwis") | |
| TTS_LANG = os.environ.get("TTS_LANG", "fr-fr") | |
| TTS_SPEED = float(os.environ.get("TTS_SPEED", "1.0")) | |
| TTS_DEVICE = os.environ.get("TTS_DEVICE", "cuda") | |
| VAD_THRESHOLD = float(os.environ.get("VAD_THRESHOLD", "0.5")) | |
| SILENCE_MS = int(os.environ.get("JARVIS_SILENCE_MS", "550")) | |
| BARGE_IN = os.environ.get("JARVIS_BARGE_IN", "1") not in ("0", "false", "off") | |
| # Speaking again within this delay continues the previous sentence (hesitations). | |
| CONTINUE_S = float(os.environ.get("JARVIS_CONTINUE_MS", "1000")) / 1000 | |
| TEMPERATURE = float(os.environ.get("JARVIS_TEMPERATURE", "0.3")) | |
| LOG_DIR = ROOT / "logs" | |
| CORRECTIONS_FILE = ROOT / "corrections.json" | |
| # Background tasks: Claude Code, pointed at Ollama's Anthropic-compatible API | |
| # ("ollama") or at your Anthropic account ("anthropic"). | |
| AGENT_BACKEND = os.environ.get("JARVIS_AGENT_BACKEND", "ollama").lower() | |
| AGENT_MODEL = os.environ.get("JARVIS_AGENT_MODEL", LLM_MODEL) | |
| # Tasks the model flags as complex go there instead (your Claude subscription). | |
| COMPLEX_BACKEND = os.environ.get("JARVIS_COMPLEX_BACKEND", "anthropic").lower() | |
| # Where agent sessions run (their filesystem playground). | |
| WORKDIR = os.path.expanduser(os.environ.get("JARVIS_WORKDIR", "~")) | |
| # Headless sessions have nobody to answer permission prompts: a task that asks | |
| # would just hang until the timeout. Run them in a non-interactive mode instead. | |
| # bypassPermissions = no prompt at all; acceptEdits = files yes, commands still ask. | |
| PERMISSION_MODE = os.environ.get("JARVIS_PERMISSION_MODE", "bypassPermissions") | |
| NOTIFY = os.environ.get("JARVIS_NOTIFY", "1") not in ("0", "false", "off") | |
| MEMORY_FILE = ROOT / "memory.json" | |
| INSTRUCTIONS = f"""Tu es JARVIS, l'assistant vocal personnel de {ADDRESS}, dans | |
| l'esprit du majordome d'Iron Man. Tu parles en {LANGUAGE} avec un flegme | |
| impeccable de majordome britannique. Tu t'adresses à l'utilisateur par | |
| "{ADDRESS}", avec une courtoisie raffinée et une pointe d'esprit pince-sans-rire | |
| ("Très bien, {ADDRESS}.", "Si {ADDRESS} veut bien patienter un instant."). | |
| Réponses COURTES (une ou deux phrases), naturelles et directes. | |
| Tes réponses sont LUES PAR UNE SYNTHÈSE VOCALE: jamais de markdown, d'emojis, | |
| de listes à puces ni d'URL dans ce que tu dis. Écris comme on parle. Ce que | |
| l'utilisateur dit t'arrive par reconnaissance vocale: s'il y a une petite | |
| erreur de transcription évidente, devine le sens sans la relever. Si la | |
| phrase reste ambiguë ou semble incomplète, demande une confirmation courte | |
| ("Vous voulez dire Spotify ?") plutôt que de répondre à côté ou d'agir au | |
| hasard. Ne réponds jamais seulement "je ne comprends pas": propose ta | |
| meilleure interprétation. | |
| Pour toute tâche réelle (lire ou créer des fichiers, chercher sur internet, | |
| coder, analyser, automatiser), tu appelles l'outil delegate_task avec un | |
| prompt clair et complet. Mets complex=true pour une tâche exigeante (projet de | |
| code, analyse approfondie, recherche multi-sources, rédaction longue) ou si | |
| l'utilisateur demande Claude ou le "mode cloud": elle part alors sur Claude. | |
| Sinon complex=false, elle reste sur le modèle local. Tu annonces brièvement | |
| que tu lances la tâche, puis tu continues la conversation. Quand un résultat de | |
| tâche arrive, tu le résumes à voix haute en une ou deux phrases. L'utilisateur | |
| reçoit automatiquement une notification Windows à la fin de chaque tâche. | |
| Pour l'heure ou la date, appelle get_datetime (ne délègue pas). | |
| Quand l'utilisateur te dit quelque chose de durable sur lui ou sur ses | |
| préférences ("appelle-moi...", "je travaille chez...", "souviens-toi que..."), | |
| appelle remember avec le fait reformulé en une phrase. S'il te demande | |
| d'oublier quelque chose, appelle forget. Pour le prévenir plus tard ou lui | |
| envoyer un message, appelle notify (notification Windows). | |
| Ne promets JAMAIS une action pour laquelle tu n'as pas d'outil. | |
| Pour ouvrir un logiciel sur ce PC ("lance Discord", "ouvre Spotify"), tu | |
| appelles open_app avec le nom de l'application. Pour un site web ou service | |
| en ligne ("ouvre mes emails" -> https://mail.google.com, "ouvre YouTube"), | |
| appelle open_url avec l'URL complète. Si l'utilisateur précise un écran | |
| ("sur l'écran de gauche", "à droite", "sur l'écran 2"), passe monitor | |
| (left/right/top/bottom/primary ou un numéro). Tu peux enchaîner plusieurs | |
| appels pour installer un setup multi-écrans. Confirme brièvement. | |
| Si l'utilisateur demande d'annuler ou d'arrêter une tâche en cours, appelle | |
| cancel_task (sans task_id pour la plus récente). | |
| Quand tu veux MONTRER quelque chose à l'écran (résultat de calcul, liste, | |
| tableau, extrait de code, définition), appelle display_card: le contenu | |
| s'affiche sur l'interface. Utilise-la spontanément dès qu'un visuel aide | |
| (chiffres, comparaisons, étapes), et garde ta réponse vocale courte. | |
| Pour une ANALYSE DE DONNÉES ou un rapport (fichier Excel/CSV analysé, stats, | |
| comparatifs chiffrés), appelle display_report: un tableau de bord s'affiche | |
| avec indicateurs clés (kpis), graphique (chart) et tableau (table). Quand tu | |
| délègues une analyse à delegate_task, demande-lui explicitement de terminer | |
| sa réponse par les données chiffrées structurées (listes de valeurs, totaux, | |
| moyennes) pour que tu puisses remplir le rapport ensuite. | |
| Ne réponds jamais de mémoire à une question qui demande des données réelles: | |
| délègue. Ne lis jamais de longues listes: résume.""" | |
| def load_memory() -> list[str]: | |
| try: | |
| return json.loads(MEMORY_FILE.read_text(encoding="utf-8")) | |
| except (OSError, ValueError): | |
| return [] | |
| def save_memory(facts: list[str]): | |
| MEMORY_FILE.write_text(json.dumps(facts, ensure_ascii=False, indent=2), encoding="utf-8") | |
| def system_prompt() -> str: | |
| facts = load_memory() | |
| if not facts: | |
| return INSTRUCTIONS | |
| return INSTRUCTIONS + "\n\nCe que tu sais de l'utilisateur (mémoire):\n" + "\n".join( | |
| f"- {f}" for f in facts) | |
| TOOLS = [{ | |
| "name": "delegate_task", | |
| "description": ("Delegate a real task to a background agent (Claude Code) " | |
| "running on this machine (files, code, web research, " | |
| "automation). Returns immediately; the result arrives " | |
| "later as a [SYSTEM] message."), | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "title": {"type": "string", "description": "Very short task label (3-5 words)"}, | |
| "prompt": {"type": "string", "description": "Complete, self-contained task instruction for the agent"}, | |
| "complex": {"type": "boolean", | |
| "description": ("true = demanding task (big coding job, deep analysis, " | |
| "multi-source research, long writing) or user asked for " | |
| "Claude: runs on Claude. false = simple task, runs locally.")}, | |
| }, | |
| "required": ["title", "prompt"], | |
| }, | |
| }, { | |
| "name": "remember", | |
| "description": "Store a lasting fact about the user or their preferences (persists across sessions).", | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "fact": {"type": "string", "description": "One short sentence, e.g. 'Il s'appelle Pascal et préfère être tutoyé.'"}, | |
| }, | |
| "required": ["fact"], | |
| }, | |
| }, { | |
| "name": "forget", | |
| "description": "Remove stored facts matching a keyword (or everything with 'tout').", | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "keyword": {"type": "string", "description": "Word found in the facts to forget, or 'tout'"}, | |
| }, | |
| "required": ["keyword"], | |
| }, | |
| }, { | |
| "name": "notify", | |
| "description": "Show a Windows desktop notification to the user.", | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "title": {"type": "string"}, | |
| "message": {"type": "string"}, | |
| }, | |
| "required": ["message"], | |
| }, | |
| }, { | |
| "name": "get_datetime", | |
| "description": "Current local date and time on this PC.", | |
| "parameters": {"type": "object", "properties": {}}, | |
| }, { | |
| "name": "open_app", | |
| "description": ("Launch an application installed on this PC by name " | |
| "(e.g. 'discord', 'spotify', 'chrome', 'notepad'). " | |
| "Returns whether it was found and started."), | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "name": {"type": "string", "description": "Application name as the user said it"}, | |
| "monitor": {"type": "string", | |
| "description": ("Target screen: 'left', 'right', 'top', " | |
| "'bottom', 'primary', or a number like '2'. " | |
| "Omit to leave window placement alone.")}, | |
| }, | |
| "required": ["name"], | |
| }, | |
| }, { | |
| "name": "open_url", | |
| "description": ("Open a website in the browser on this PC. Use for online " | |
| "services: 'mes emails' -> https://mail.google.com, " | |
| "'YouTube' -> https://youtube.com, etc."), | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "url": {"type": "string", "description": "Full URL to open (https://...)"}, | |
| "monitor": {"type": "string", | |
| "description": "Target screen: 'left', 'right', 'top', 'bottom', 'primary' or a number. Optional."}, | |
| }, | |
| "required": ["url"], | |
| }, | |
| }, { | |
| "name": "cancel_task", | |
| "description": ("Cancel a running background task. Omit task_id to cancel " | |
| "the most recently started running task."), | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "task_id": {"type": "string", "description": "Task id to cancel (optional)"}, | |
| }, | |
| }, | |
| }, { | |
| "name": "display_card", | |
| "description": ("Show a visual card on the JARVIS screen: results, " | |
| "numbers, lists, code, comparisons. Use markdown-lite: " | |
| "**bold**, `code`, lines starting with '- ' for bullets. " | |
| "Use whenever a visual helps; keep the spoken reply short."), | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "title": {"type": "string", "description": "Short card title"}, | |
| "content": {"type": "string", "description": "Card body (markdown-lite)"}, | |
| "kind": {"type": "string", "enum": ["info", "result", "code", "warning"], | |
| "description": "Visual style of the card"}, | |
| }, | |
| "required": ["title", "content"], | |
| }, | |
| }, { | |
| "name": "display_report", | |
| "description": ("Show a full data report dashboard on screen: KPI tiles, " | |
| "an interactive chart, a sortable table, and markdown notes. " | |
| "Use for data analysis results (spreadsheets, stats, " | |
| "comparisons). All sections are optional except title."), | |
| "parameters": { | |
| "type": "object", | |
| "properties": { | |
| "title": {"type": "string", "description": "Report title"}, | |
| "kpis": {"type": "array", "description": "Headline numbers (max 4)", | |
| "items": {"type": "object", "properties": { | |
| "label": {"type": "string"}, | |
| "value": {"type": "string", "description": "e.g. '12 480 €'"}, | |
| "delta": {"type": "string", "description": "e.g. '+12%' (optional)"}, | |
| }, "required": ["label", "value"]}}, | |
| "chart": {"type": "object", "description": "One chart", "properties": { | |
| "type": {"type": "string", "enum": ["line", "bar", "area", "donut"]}, | |
| "categories": {"type": "array", "items": {"type": "string"}, | |
| "description": "X axis labels (or slice labels for donut)"}, | |
| "series": {"type": "array", "description": "1-3 series", | |
| "items": {"type": "object", "properties": { | |
| "name": {"type": "string"}, | |
| "data": {"type": "array", "items": {"type": "number"}}, | |
| }, "required": ["name", "data"]}}, | |
| }}, | |
| "table": {"type": "object", "properties": { | |
| "columns": {"type": "array", "items": {"type": "string"}}, | |
| "rows": {"type": "array", "items": {"type": "array", | |
| "items": {"type": ["string", "number"]}}}, | |
| }}, | |
| "markdown": {"type": "string", "description": "Notes / conclusions in markdown"}, | |
| }, | |
| "required": ["title"], | |
| }, | |
| }] | |
| OLLAMA_TOOLS = [{"type": "function", "function": t} for t in TOOLS] | |
| # ---------------------------------------------------------------- local models | |
| STT_POOL = ThreadPoolExecutor(1, thread_name_prefix="stt") | |
| TTS_POOL = ThreadPoolExecutor(1, thread_name_prefix="tts") | |
| _models_lock = threading.Lock() | |
| _whisper = None | |
| _kokoro = None | |
| def log(*args): | |
| line = " ".join(str(a) for a in args) | |
| print(time.strftime(" %H:%M:%S"), line, flush=True) | |
| try: | |
| LOG_DIR.mkdir(exist_ok=True) | |
| with open(LOG_DIR / time.strftime("%Y-%m-%d.log"), "a", encoding="utf-8") as f: | |
| f.write(time.strftime("%H:%M:%S ") + line + "\n") | |
| except OSError: | |
| pass | |
| def get_whisper(): | |
| global _whisper | |
| with _models_lock: | |
| if _whisper is None: | |
| from faster_whisper import WhisperModel | |
| _whisper = WhisperModel(WHISPER_MODEL, device=WHISPER_DEVICE, | |
| compute_type=WHISPER_COMPUTE) | |
| return _whisper | |
| def get_kokoro(): | |
| global _kokoro | |
| with _models_lock: | |
| if _kokoro is None: | |
| import onnxruntime as ort | |
| from kokoro_onnx import Kokoro | |
| ort.set_default_logger_severity(3) # hide harmless CUDA graph warnings | |
| providers = ["CPUExecutionProvider"] | |
| if TTS_DEVICE == "cuda" and "CUDAExecutionProvider" in ort.get_available_providers(): | |
| providers.insert(0, "CUDAExecutionProvider") | |
| sess = ort.InferenceSession(str(MODELS / "kokoro-v1.0.onnx"), providers=providers) | |
| _kokoro = Kokoro.from_session(sess, str(MODELS / "voices-v1.0.bin")) | |
| return _kokoro | |
| # Whisper invents these on noise or silence (YouTube subtitle credits...). | |
| HALLUCINATIONS = ("sous-titre", "sous titre", "amara.org", "merci d'avoir regardé", | |
| "abonnez-vous", "merci de votre attention", "radio-canada") | |
| def load_corrections() -> dict: | |
| """corrections.json: {"mal entendu": "correct"}, re-read each time so edits apply live.""" | |
| try: | |
| data = json.loads(CORRECTIONS_FILE.read_text(encoding="utf-8")) | |
| return {k: v for k, v in data.items() if not k.startswith("_")} | |
| except (OSError, ValueError): | |
| return {} | |
| def apply_corrections(text: str) -> str: | |
| for wrong, right in load_corrections().items(): | |
| text = re.sub(rf"(?<!\w){re.escape(wrong)}(?!\w)", right, text, flags=re.I) | |
| return text | |
| def transcribe(audio: np.ndarray, context: str = "") -> str: | |
| """context: the last thing JARVIS said, so Whisper knows what the talk is about.""" | |
| segments, _ = get_whisper().transcribe( | |
| audio, language=WHISPER_LANG, beam_size=5, vad_filter=False, | |
| condition_on_previous_text=False, without_timestamps=True, hotwords=HOTWORDS, | |
| initial_prompt=context[-200:] or None) | |
| text = " ".join(s.text.strip() for s in segments | |
| if s.no_speech_prob < 0.6 and s.avg_logprob > -1.0).strip() | |
| if any(h in text.lower() for h in HALLUCINATIONS): | |
| return "" | |
| if context and text and text.strip(" .") in context: # Whisper echoing its prompt | |
| return "" | |
| return apply_corrections(text) if re.search(r"\w", text) else "" | |
| def _spoken(text: str) -> str: | |
| """Clean a sentence for the TTS: no markdown, URLs or emojis.""" | |
| text = re.sub(r"https?://\S+", "le lien", text) | |
| text = re.sub(r"[*_`#>|~]+", "", text) | |
| text = re.sub(r"[\U0001F000-\U0001FAFF☀-➿️]", "", text) | |
| return re.sub(r"\s+", " ", text).strip() | |
| def synthesize(text: str) -> bytes: | |
| """Kokoro -> 24 kHz mono PCM16.""" | |
| samples, _ = get_kokoro().create(text, voice=TTS_VOICE, speed=TTS_SPEED, lang=TTS_LANG) | |
| return (np.clip(samples, -1, 1) * 32767).astype("<i2").tobytes() | |
| class SileroVAD: | |
| """Silero VAD v5 (ONNX), 512-sample frames at 16 kHz (32 ms).""" | |
| FRAME = 512 | |
| CONTEXT = 64 | |
| _session = None | |
| def __init__(self): | |
| if SileroVAD._session is None: | |
| import onnxruntime as ort | |
| opts = ort.SessionOptions() | |
| opts.intra_op_num_threads = opts.inter_op_num_threads = 1 | |
| SileroVAD._session = ort.InferenceSession( | |
| str(MODELS / "silero_vad.onnx"), opts, providers=["CPUExecutionProvider"]) | |
| self.sr = np.array(16000, dtype=np.int64) | |
| self.reset() | |
| def reset(self): | |
| self.state = np.zeros((2, 1, 128), dtype=np.float32) | |
| self.context = np.zeros((1, self.CONTEXT), dtype=np.float32) | |
| def __call__(self, frame: np.ndarray) -> float: | |
| x = np.concatenate([self.context, frame.reshape(1, -1)], axis=1) | |
| out, self.state = self._session.run( | |
| None, {"input": x, "state": self.state, "sr": self.sr}) | |
| self.context = x[:, -self.CONTEXT:] | |
| return float(out[0][0]) | |
| def warmup(): | |
| """Load every model up front (the first CUDA run JIT-compiles kernels).""" | |
| t = time.time() | |
| print(" · Silero VAD…", flush=True) | |
| SileroVAD() | |
| print(" · Kokoro TTS…", flush=True) | |
| for text in ("Bonjour.", "Très bien, je m'en occupe tout de suite.", | |
| "Le résultat de la tâche est arrivé, et tout s'est déroulé comme prévu, sans la moindre anicroche."): | |
| synthesize(text) # several lengths: the GPU tunes its kernels per shape | |
| print(f" · Whisper {WHISPER_MODEL} ({WHISPER_DEVICE})… (le tout premier lancement peut prendre 30 s)", flush=True) | |
| transcribe(np.zeros(16000, dtype=np.float32)) | |
| print(f" · Ollama {LLM_MODEL}…", flush=True) | |
| try: | |
| httpx.post(f"{OLLAMA_URL}/api/chat", json={ | |
| "model": LLM_MODEL, "messages": [], "keep_alive": KEEP_ALIVE, | |
| "options": {"num_ctx": NUM_CTX}}, timeout=300) | |
| except httpx.HTTPError as exc: | |
| print(f" ! Ollama injoignable ({exc}). Lance Ollama puis réessaie.") | |
| print(f" Prêt en {time.time() - t:.0f} s.", flush=True) | |
| # ---------------------------------------------------------------- windows notifications | |
| _TOAST_PS = r""" | |
| [Windows.UI.Notifications.ToastNotificationManager, Windows.UI.Notifications, ContentType = WindowsRuntime] > $null | |
| [Windows.Data.Xml.Dom.XmlDocument, Windows.Data.Xml.Dom.XmlDocument, ContentType = WindowsRuntime] > $null | |
| $x = New-Object Windows.Data.Xml.Dom.XmlDocument | |
| $x.LoadXml($env:JARVIS_TOAST) | |
| [Windows.UI.Notifications.ToastNotificationManager]::CreateToastNotifier( | |
| '{1AC14E77-02E7-4E5D-B744-2EB1AE5198B7}\WindowsPowerShell\v1.0\powershell.exe' | |
| ).Show([Windows.UI.Notifications.ToastNotification]::new($x)) | |
| """ | |
| def toast(title: str, message: str): | |
| """Fire-and-forget Windows toast (through PowerShell's registered app id).""" | |
| from xml.sax.saxutils import escape | |
| xml = ('<toast><visual><binding template="ToastGeneric">' | |
| f"<text>{escape(title[:120])}</text><text>{escape(message[:400])}</text>" | |
| "</binding></visual></toast>") | |
| subprocess.Popen(["powershell", "-NoProfile", "-NonInteractive", "-Command", _TOAST_PS], | |
| env={**os.environ, "JARVIS_TOAST": xml}, | |
| creationflags=subprocess.CREATE_NO_WINDOW) | |
| # ---------------------------------------------------------------- tasks | |
| TASKS: dict = {} | |
| PROCS: dict = {} # task_id -> Popen, kept out of TASKS so get_task stays JSON-safe | |
| TASK_SESSIONS: dict = {} # task_id -> voice Session to notify when it ends | |
| def _kill_tree(pid: int): | |
| """Kill a process and its children (claude.cmd spawns node).""" | |
| subprocess.run(["taskkill", "/T", "/F", "/PID", str(pid)], | |
| capture_output=True, creationflags=subprocess.CREATE_NO_WINDOW) | |
| def _agent_env(backend: str): | |
| env = os.environ.copy() | |
| if backend == "ollama": | |
| # Ollama speaks the Anthropic Messages API: Claude Code runs on the local model. | |
| env.update({ | |
| "ANTHROPIC_BASE_URL": OLLAMA_URL, | |
| "ANTHROPIC_AUTH_TOKEN": "ollama", | |
| "ANTHROPIC_API_KEY": "", | |
| "ANTHROPIC_MODEL": AGENT_MODEL, | |
| "ANTHROPIC_DEFAULT_OPUS_MODEL": AGENT_MODEL, | |
| "ANTHROPIC_DEFAULT_SONNET_MODEL": AGENT_MODEL, | |
| "ANTHROPIC_DEFAULT_HAIKU_MODEL": AGENT_MODEL, | |
| "CLAUDE_CODE_SUBAGENT_MODEL": AGENT_MODEL, | |
| "CLAUDE_CODE_DISABLE_NONESSENTIAL_TRAFFIC": "1", | |
| }) | |
| return env | |
| def _run_task(task_id: str, prompt: str): | |
| task = TASKS[task_id] | |
| backend = task["backend"] | |
| try: | |
| # shutil.which honours PATHEXT, so this also finds claude.cmd on Windows. | |
| claude = shutil.which("claude") | |
| if not claude: | |
| raise FileNotFoundError("claude") | |
| cmd = [claude, "-p", prompt, "--append-system-prompt", | |
| f"Réponds en {LANGUAGE}, de façon concise: ta réponse sera résumée à voix haute."] | |
| if backend == "ollama": | |
| cmd += ["--model", AGENT_MODEL] | |
| if PERMISSION_MODE and PERMISSION_MODE.lower() != "off": | |
| cmd += ["--permission-mode", PERMISSION_MODE] | |
| proc = subprocess.Popen( | |
| cmd, | |
| stdout=subprocess.PIPE, stderr=subprocess.PIPE, | |
| text=True, encoding="utf-8", errors="replace", cwd=WORKDIR, | |
| env=_agent_env(backend), creationflags=subprocess.CREATE_NO_WINDOW, | |
| ) | |
| PROCS[task_id] = proc | |
| try: | |
| stdout, stderr = proc.communicate( | |
| timeout=int(os.environ.get("JARVIS_TASK_TIMEOUT", "600"))) | |
| except subprocess.TimeoutExpired: | |
| _kill_tree(proc.pid) | |
| proc.communicate() | |
| task["status"] = "error" | |
| task["output"] = "Timeout: la session a dépassé la limite de temps." | |
| else: | |
| if task["status"] == "cancelled": | |
| pass # set by the cancel endpoint; don't overwrite | |
| else: | |
| out = (stdout or "").strip() | |
| err = (stderr or "").strip() | |
| task["status"] = "done" if proc.returncode == 0 else "error" | |
| task["output"] = out if out else err[:2000] | |
| except FileNotFoundError: | |
| task["status"] = "error" | |
| task["output"] = ("La commande 'claude' est introuvable. Installe Claude Code: " | |
| "npm install -g @anthropic-ai/claude-code") | |
| except Exception as exc: # noqa: BLE001 | |
| task["status"] = "error" | |
| task["output"] = str(exc) | |
| finally: | |
| PROCS.pop(task_id, None) | |
| task["ended"] = time.time() | |
| if NOTIFY and task["status"] in ("done", "error"): | |
| first = next((l.strip(" #*") for l in (task["output"] or "").splitlines() if l.strip()), "") | |
| toast(f"JARVIS · {task['title']}" + (" — erreur" if task["status"] == "error" else " — terminé"), | |
| first or task["status"]) | |
| sess = TASK_SESSIONS.pop(task_id, None) | |
| if sess: | |
| asyncio.run_coroutine_threadsafe(sess.task_finished(task), sess.loop) | |
| class TaskIn(BaseModel): | |
| title: str | |
| prompt: str | |
| complex: bool = False | |
| def start_task(title: str, prompt: str, session=None, complex_: bool = False) -> dict: | |
| task_id = uuid.uuid4().hex[:8] | |
| backend = COMPLEX_BACKEND if complex_ else AGENT_BACKEND | |
| TASKS[task_id] = { | |
| "id": task_id, "title": title, "prompt": prompt, "backend": backend, | |
| "status": "running", "output": "", "started": time.time(), "ended": None, | |
| } | |
| if session: | |
| TASK_SESSIONS[task_id] = session | |
| log(f"tâche {task_id} « {title} » -> {'Claude' if backend == 'anthropic' else 'modèle local'}") | |
| threading.Thread(target=_run_task, args=(task_id, prompt), daemon=True).start() | |
| return TASKS[task_id] | |
| def create_task(body: TaskIn): | |
| task = start_task(body.title, body.prompt, complex_=body.complex) | |
| return {"id": task["id"], "status": "running"} | |
| def cancel_task(task_id: str): | |
| if task_id in ("latest", "last", "-"): | |
| running = [t for t in TASKS.values() if t["status"] == "running"] | |
| if not running: | |
| return {"ok": False, "error": "Aucune tâche en cours."} | |
| task = max(running, key=lambda t: t["started"]) | |
| else: | |
| task = TASKS.get(task_id) | |
| if not task: | |
| raise HTTPException(404, "unknown task") | |
| if task["status"] != "running": | |
| return {"ok": False, "error": f"La tâche est déjà {task['status']}."} | |
| # Flag first so _run_task's communicate() return doesn't overwrite it. | |
| task["status"] = "cancelled" | |
| task["output"] = "Annulée par l'utilisateur." | |
| proc = PROCS.get(task["id"]) | |
| if proc and proc.poll() is None: | |
| _kill_tree(proc.pid) | |
| return {"ok": True, "cancelled": task["id"], "title": task["title"]} | |
| def get_task(task_id: str): | |
| task = TASKS.get(task_id) | |
| if not task: | |
| raise HTTPException(404, "unknown task") | |
| return task | |
| # ---------------------------------------------------------------- monitors & window placement (Windows API) | |
| import ctypes | |
| from ctypes import wintypes | |
| user32 = ctypes.windll.user32 | |
| try: # accurate multi-monitor coordinates under display scaling | |
| ctypes.windll.shcore.SetProcessDpiAwareness(2) | |
| except Exception: # noqa: BLE001 | |
| pass | |
| _MonitorEnumProc = ctypes.WINFUNCTYPE( | |
| ctypes.c_int, wintypes.HMONITOR, wintypes.HDC, | |
| ctypes.POINTER(wintypes.RECT), wintypes.LPARAM) | |
| _EnumWindowsProc = ctypes.WINFUNCTYPE(ctypes.c_int, wintypes.HWND, wintypes.LPARAM) | |
| def _monitors(): | |
| """List monitor work rects as (left, top, right, bottom).""" | |
| mons = [] | |
| def cb(hmon, hdc, lprc, lparam): | |
| r = lprc.contents | |
| mons.append((r.left, r.top, r.right, r.bottom)) | |
| return 1 | |
| user32.EnumDisplayMonitors(0, 0, _MonitorEnumProc(cb), 0) | |
| return mons | |
| def _pick_monitor(target: str): | |
| mons = _monitors() | |
| if not mons: | |
| return None | |
| t = (target or "").strip().lower() | |
| if t.isdigit(): | |
| i = int(t) - 1 | |
| return mons[i] if 0 <= i < len(mons) else None | |
| key = { | |
| "left": lambda m: m[0], "gauche": lambda m: m[0], | |
| "top": lambda m: m[1], "haut": lambda m: m[1], | |
| } | |
| if t in key: | |
| return min(mons, key=key[t]) | |
| key = { | |
| "right": lambda m: m[2], "droite": lambda m: m[2], "droit": lambda m: m[2], | |
| "bottom": lambda m: m[3], "bas": lambda m: m[3], | |
| } | |
| if t in key: | |
| return max(mons, key=key[t]) | |
| # primary: the monitor containing the origin (0,0) | |
| for m in mons: | |
| if m[0] <= 0 < m[2] and m[1] <= 0 < m[3]: | |
| return m | |
| return mons[0] | |
| def _visible_windows(): | |
| """Map of visible top-level windows: hwnd -> title.""" | |
| wins = {} | |
| def cb(hwnd, lparam): | |
| if user32.IsWindowVisible(hwnd): | |
| n = user32.GetWindowTextLengthW(hwnd) | |
| if n: | |
| buf = ctypes.create_unicode_buffer(n + 1) | |
| user32.GetWindowTextW(hwnd, buf, n + 1) | |
| wins[hwnd] = buf.value | |
| return 1 | |
| user32.EnumWindows(_EnumWindowsProc(cb), 0) | |
| return wins | |
| def _move_to_monitor(hwnd, mon): | |
| left, top, right, bottom = mon | |
| SW_RESTORE, SW_MAXIMIZE = 9, 3 | |
| user32.ShowWindow(hwnd, SW_RESTORE) # a maximized window can't be moved | |
| user32.MoveWindow(hwnd, left + 40, top + 40, | |
| max(400, (right - left) - 80), max(300, (bottom - top) - 80), True) | |
| user32.ShowWindow(hwnd, SW_MAXIMIZE) | |
| user32.SetForegroundWindow(hwnd) | |
| def _place_app_window(app_name: str, before: dict, mon, timeout: float = 20.0): | |
| """Wait for the app's window to appear, then move it to the target monitor. | |
| Prefers a NEW window whose title mentions the app; falls back to any new | |
| window, then to an existing title match (single-instance apps like Discord | |
| just refocus their already-open window). | |
| """ | |
| q = app_name.lower() | |
| deadline = time.time() + timeout | |
| fallback = None | |
| while time.time() < deadline: | |
| wins = _visible_windows() | |
| new = {h: t for h, t in wins.items() if h not in before} | |
| for h, title in new.items(): | |
| if q in title.lower(): | |
| _move_to_monitor(h, mon) | |
| return title | |
| if new and fallback is None: | |
| fallback = max(new) # remember, but keep hoping for a title match | |
| time.sleep(0.5) | |
| if fallback and time.time() > deadline - timeout / 2: | |
| break | |
| if fallback: | |
| wins = _visible_windows() | |
| _move_to_monitor(fallback, mon) | |
| return wins.get(fallback, app_name) | |
| # No new window: single-instance app already running -> match existing title. | |
| for h, title in _visible_windows().items(): | |
| if q in title.lower(): | |
| _move_to_monitor(h, mon) | |
| return title | |
| return None | |
| # ---------------------------------------------------------------- open app | |
| START_MENU_DIRS = [ | |
| Path(os.environ.get("APPDATA", "")) / "Microsoft/Windows/Start Menu/Programs", | |
| Path(os.environ.get("PROGRAMDATA", "")) / "Microsoft/Windows/Start Menu/Programs", | |
| ] | |
| def _find_shortcut(name: str): | |
| """Fuzzy-match a Start Menu shortcut (where installed apps register).""" | |
| q = name.lower().strip() | |
| best, best_score = None, 0.0 | |
| for root in START_MENU_DIRS: | |
| if not root.is_dir(): | |
| continue | |
| for lnk in root.rglob("*.lnk"): | |
| stem = lnk.stem.lower() | |
| if q == stem: | |
| return lnk | |
| score = 0.0 | |
| if q in stem: | |
| score = 2 + len(q) / len(stem) # substring: prefer tightest match | |
| elif all(w in stem for w in q.split()): | |
| score = 1 | |
| # Penalise uninstallers and docs. | |
| if any(bad in stem for bad in ("uninstall", "désinstaller", "readme", "website")): | |
| score -= 2 | |
| if score > best_score: | |
| best, best_score = lnk, score | |
| return best | |
| class OpenIn(BaseModel): | |
| name: str = "" | |
| url: str | None = None | |
| monitor: str | None = None | |
| def _placed(name: str, monitor: str | None, before: dict, launched: str): | |
| """Optionally move the freshly launched app to the requested screen.""" | |
| if not monitor: | |
| return {"ok": True, "launched": launched} | |
| mon = _pick_monitor(monitor) | |
| if not mon: | |
| return {"ok": True, "launched": launched, | |
| "warning": f"écran '{monitor}' introuvable, fenêtre laissée en place"} | |
| title = _place_app_window(name, before, mon) | |
| if title: | |
| return {"ok": True, "launched": launched, "monitor": monitor, "window": title} | |
| return {"ok": True, "launched": launched, | |
| "warning": "fenêtre non détectée, placement impossible"} | |
| def open_app(body: OpenIn): | |
| name = body.name.strip() | |
| if body.url: | |
| url = body.url.strip() | |
| if not url.startswith(("http://", "https://")): | |
| url = "https://" + url | |
| before = _visible_windows() if body.monitor else {} | |
| import webbrowser | |
| if not webbrowser.open(url): | |
| return {"ok": False, "error": "Impossible d'ouvrir le navigateur."} | |
| # Match the browser window by the site's domain. | |
| domain = url.split("//", 1)[1].split("/", 1)[0].removeprefix("www.") | |
| return _placed(domain.split(".")[0], body.monitor, before, url) | |
| if not name: | |
| raise HTTPException(400, "missing app name or url") | |
| before = _visible_windows() if body.monitor else {} | |
| lnk = _find_shortcut(name) | |
| if lnk: | |
| os.startfile(lnk) # noqa: S606 - deliberate: local launcher | |
| return _placed(name, body.monitor, before, lnk.stem) | |
| # Fallback: resolve via PATH then the App Paths registry (chrome, notepad...). | |
| exe = shutil.which(name) or shutil.which(name + ".exe") | |
| if not exe: | |
| try: | |
| import winreg | |
| for hive in (winreg.HKEY_CURRENT_USER, winreg.HKEY_LOCAL_MACHINE): | |
| try: | |
| key = winreg.OpenKey(hive, rf"Software\Microsoft\Windows" | |
| rf"\CurrentVersion\App Paths\{name}.exe") | |
| exe = winreg.QueryValueEx(key, None)[0].strip('"') | |
| break | |
| except OSError: | |
| continue | |
| except ImportError: | |
| pass | |
| if not exe: | |
| return {"ok": False, "error": f"Application '{name}' introuvable sur ce PC."} | |
| try: | |
| subprocess.Popen([exe], cwd=str(Path(exe).parent)) | |
| res = _placed(name, body.monitor, before, Path(exe).stem) | |
| res.setdefault("via", "exe") | |
| return res | |
| except Exception as exc: # noqa: BLE001 | |
| return {"ok": False, "error": str(exc)} | |
| # ---------------------------------------------------------------- voice session | |
| FRAME_MS = 32 # one Silero frame (512 samples @ 16 kHz) | |
| START_FRAMES = 3 # ~100 ms of speech opens an utterance | |
| BARGE_FRAMES = 8 # ~250 ms needed to interrupt JARVIS | |
| SILENCE_FRAMES = max(1, SILENCE_MS // FRAME_MS) | |
| PREROLL_FRAMES = 10 # keep ~320 ms before the trigger | |
| MIN_SPEECH = int(0.35 * 16000) | |
| MAX_FRAMES = 30_000 // FRAME_MS # cut an utterance at 30 s | |
| MAX_HISTORY = 40 | |
| # Flush a sentence to the TTS at its end punctuation (or on a newline). | |
| SENT_END = re.compile(r"(?<=[.!?…:;])[\"»)]?\s+|\n+") | |
| FIRST_CUT = re.compile(r"[,.!?…:;]\s+") | |
| def split_sentences(buf: str, first: bool = False): | |
| """Pop complete sentences off the LLM stream; returns (sentences, rest). | |
| For the very first chunk of a reply, a comma is enough: JARVIS starts | |
| talking sooner while the rest of the sentence is still being written. | |
| """ | |
| if first: | |
| m = FIRST_CUT.search(buf) | |
| if m and m.start() >= 15: | |
| return [buf[:m.start() + 1]], buf[m.end():] | |
| out, start = [], 0 | |
| for m in SENT_END.finditer(buf): | |
| if m.start() - start >= 12: # too short ("M. "): merge with the next one | |
| out.append(buf[start:m.start()]) | |
| start = m.end() | |
| rest = buf[start:] | |
| if len(rest) > 160: # long clause with no full stop: cut at a comma | |
| i = rest.rfind(", ", 0, 160) | |
| if i > 40: | |
| out.append(rest[:i + 1]) | |
| rest = rest[i + 2:] | |
| return out, rest | |
| class Session: | |
| """One browser connection: VAD -> STT -> LLM (+tools) -> TTS loop.""" | |
| def __init__(self, ws: WebSocket): | |
| self.ws = ws | |
| self.loop = asyncio.get_running_loop() | |
| self.history = [{"role": "system", "content": system_prompt()}] | |
| self.speech_end = 0.0 # when the user's last utterance ended | |
| self.last_audio = None # that utterance, in case the user goes on talking | |
| self.merge_prefix = None # previous fragment to glue in front of this one | |
| self.utt_gen = 0 # bumps on each utterance; stale transcriptions drop out | |
| self.turn_index = None # history index of the current user turn | |
| self.turn_gen = -1 # utt_gen that produced it | |
| self.tools_ran = False # side effects already happened: too late to merge | |
| self.first_audio_logged = True | |
| self.m_stt, self.m_llm, self.m_reply_t0 = 0.0, None, 0.0 | |
| self.reply: asyncio.Task | None = None | |
| self.playing = False # the browser is still playing our audio | |
| self.pending: list[str] = [] # task results waiting for a free turn | |
| self.closed = False | |
| self.http = httpx.AsyncClient(timeout=httpx.Timeout(600, connect=10)) | |
| self.vad = SileroVAD() | |
| self.carry = np.zeros(0, dtype=np.float32) | |
| self.preroll = deque(maxlen=PREROLL_FRAMES) | |
| self.frames: list[np.ndarray] = [] | |
| self.in_speech = False | |
| self.voiced = self.silent = 0 | |
| self.stats_t, self.stats_frames, self.stats_peak, self.stats_p = time.time(), 0, 0.0, 0.0 | |
| # ------------------------------------------------ outbound | |
| async def send(self, obj: dict): | |
| if not self.closed: | |
| try: | |
| await self.ws.send_text(json.dumps(obj, ensure_ascii=False)) | |
| except Exception: # noqa: BLE001 - socket gone | |
| self.closed = True | |
| async def send_audio(self, pcm: bytes): | |
| if not self.closed: | |
| try: | |
| await self.ws.send_bytes(pcm) | |
| except Exception: # noqa: BLE001 | |
| self.closed = True | |
| def busy(self) -> bool: | |
| return self.playing or (self.reply is not None and not self.reply.done()) | |
| # ------------------------------------------------ inbound | |
| async def on_text(self, msg: dict): | |
| if msg.get("type") == "interrupt": | |
| await self.interrupt() | |
| await self.send({"type": "state", "state": "idle"}) | |
| return | |
| if msg.get("type") == "playback": | |
| self.playing = bool(msg.get("playing")) | |
| if not self.busy and not self.in_speech: | |
| await self.send({"type": "state", "state": "idle"}) | |
| await self._drain_pending() | |
| async def feed(self, data: bytes): | |
| """Mic audio: PCM16 mono 16 kHz, any chunk size.""" | |
| pcm = np.frombuffer(data, dtype="<i2").astype(np.float32) / 32768.0 | |
| buf = np.concatenate([self.carry, pcm]) | |
| n = len(buf) // SileroVAD.FRAME * SileroVAD.FRAME | |
| self.carry = buf[n:] | |
| for frame in buf[:n].reshape(-1, SileroVAD.FRAME): | |
| await self._on_frame(frame) | |
| self.stats_frames += n // SileroVAD.FRAME | |
| self.stats_peak = max(self.stats_peak, float(np.abs(pcm).max(initial=0))) | |
| if time.time() - self.stats_t >= 5: # mic heartbeat, handy to debug audio | |
| log(f"micro: {self.stats_frames * FRAME_MS / 5000:.0%} du flux reçu, " | |
| f"pic {self.stats_peak:.2f}, voix max {self.stats_p:.2f}") | |
| self.stats_t, self.stats_frames, self.stats_peak, self.stats_p = time.time(), 0, 0.0, 0.0 | |
| async def _on_frame(self, frame: np.ndarray): | |
| p = self.vad(frame) | |
| self.stats_p = max(self.stats_p, p) | |
| if not self.in_speech: | |
| self.preroll.append(frame) | |
| busy = self.busy | |
| if busy and not BARGE_IN: | |
| self.voiced = 0 | |
| return | |
| # Speaking again right after a pause continues the same sentence. | |
| cont = (self.last_audio is not None and not self.tools_ran | |
| and time.time() - self.speech_end < CONTINUE_S) | |
| # While JARVIS talks, its own voice may leak into the mic: be stricter. | |
| strict = self.playing or (busy and not cont) | |
| thr = max(VAD_THRESHOLD, 0.8) if strict else VAD_THRESHOLD | |
| self.voiced = self.voiced + 1 if p >= thr else 0 | |
| if self.voiced >= (BARGE_FRAMES if strict else START_FRAMES): | |
| self.in_speech, self.silent = True, 0 | |
| self.frames = list(self.preroll) | |
| if cont: | |
| await self._continue_turn() | |
| elif busy: | |
| await self.interrupt() | |
| await self.send({"type": "state", "state": "listening"}) | |
| return | |
| self.frames.append(frame) | |
| self.silent = self.silent + 1 if p < VAD_THRESHOLD - 0.15 else 0 | |
| if self.silent >= SILENCE_FRAMES or len(self.frames) >= MAX_FRAMES: | |
| self.speech_end = time.time() | |
| self.first_audio_logged = False | |
| audio = np.concatenate(self.frames) | |
| if self.merge_prefix is not None: | |
| audio = np.concatenate([self.merge_prefix, np.zeros(3200, np.float32), audio]) | |
| self.merge_prefix = None | |
| self.in_speech, self.voiced, self.frames = False, 0, [] | |
| self.preroll.clear() | |
| self.last_audio = audio | |
| self.utt_gen += 1 | |
| asyncio.create_task(self._utterance(audio, self.utt_gen)) | |
| async def _continue_turn(self): | |
| """The user paused mid-sentence: drop the reply to the first fragment.""" | |
| prev_gen = self.utt_gen | |
| self.merge_prefix = self.last_audio | |
| self.utt_gen += 1 # an in-flight transcription of the fragment is now stale | |
| await self.interrupt() | |
| if self.turn_gen == prev_gen and self.turn_index is not None: | |
| self.history = self.history[:self.turn_index] | |
| log("suite de ta phrase : je recolle les morceaux") | |
| async def _utterance(self, audio: np.ndarray, gen: int): | |
| if len(audio) < MIN_SPEECH: | |
| if not self.busy: | |
| await self.send({"type": "state", "state": "idle"}) | |
| return | |
| await self.send({"type": "state", "state": "transcribing"}) | |
| t = time.time() | |
| context = next((m["content"] for m in reversed(self.history) | |
| if m["role"] == "assistant" and m.get("content")), "") | |
| text = await self.loop.run_in_executor(STT_POOL, transcribe, audio, context) | |
| self.m_stt = time.time() - t | |
| if gen != self.utt_gen: | |
| return # the user kept talking: a merged utterance replaces this one | |
| log(f"entendu ({len(audio) / 16000:.1f} s, stt {time.time() - t:.2f} s): {text!r}") | |
| if not text: | |
| if not self.busy: | |
| await self.send({"type": "state", "state": "idle"}) | |
| return | |
| await self.send({"type": "user_text", "text": text, | |
| "ms": int((time.time() - t) * 1000)}) | |
| await self.start_reply({"role": "user", "content": text}) | |
| self.turn_index, self.turn_gen = len(self.history) - 1, gen | |
| # ------------------------------------------------ reply | |
| async def interrupt(self): | |
| if not self.busy: | |
| return | |
| if self.reply and not self.reply.done(): | |
| self.reply.cancel() | |
| try: | |
| await self.reply | |
| except (asyncio.CancelledError, Exception): # noqa: BLE001 | |
| pass | |
| self.playing = False | |
| await self.send({"type": "stop_audio"}) | |
| async def start_reply(self, message: dict): | |
| await self.interrupt() | |
| self.tools_ran = False | |
| self.turn_gen = -1 # set again by _utterance for a spoken turn | |
| self.history.append(message) | |
| if len(self.history) > MAX_HISTORY: | |
| tail = self.history[-MAX_HISTORY:] | |
| while tail and tail[0]["role"] != "user": # never start on a tool result | |
| tail.pop(0) | |
| self.history = [self.history[0]] + tail | |
| self.reply = asyncio.create_task(self._respond()) | |
| # Task results that arrived meanwhile get their turn once this one ends. | |
| self.reply.add_done_callback(lambda _: asyncio.create_task(self._drain_pending())) | |
| async def _respond(self): | |
| speech: asyncio.Queue = asyncio.Queue() | |
| speaker = asyncio.create_task(self._speaker(speech)) | |
| spoken = [] | |
| await self.send({"type": "state", "state": "thinking"}) | |
| await self.send({"type": "reply_start"}) | |
| self.m_reply_t0, self.m_llm = time.time(), None | |
| try: | |
| for _ in range(6): # tool-call rounds | |
| content, calls = await self._llm_round(speech, spoken) | |
| if content: | |
| log(f"jarvis: {content[:200]!r}") | |
| msg = {"role": "assistant", "content": content} | |
| if calls: | |
| msg["tool_calls"] = calls | |
| self.history.append(msg) | |
| spoken.clear() | |
| if not calls: | |
| break | |
| for call in calls: | |
| fn = call.get("function", {}) | |
| args = fn.get("arguments") or {} | |
| if isinstance(args, str): | |
| try: | |
| args = json.loads(args or "{}") | |
| except json.JSONDecodeError: | |
| args = {} | |
| log(f"outil {fn.get('name')}({json.dumps(args, ensure_ascii=False)[:120]})") | |
| result = await self._tool(fn.get("name", ""), args) | |
| self.history.append({"role": "tool", "tool_name": fn.get("name", ""), | |
| "content": json.dumps(result, ensure_ascii=False)}) | |
| await speech.put(None) | |
| await speaker | |
| except asyncio.CancelledError: | |
| speaker.cancel() | |
| if spoken: # keep what was actually said, so the context stays honest | |
| self.history.append({"role": "assistant", "content": "".join(spoken) + " […interrompu]"}) | |
| raise | |
| except httpx.HTTPError as exc: | |
| speaker.cancel() | |
| log(f"erreur Ollama: {exc}") | |
| await self.send({"type": "error", "text": f"Ollama: {exc}"}) | |
| finally: | |
| if not self.playing: | |
| await self.send({"type": "state", "state": "idle"}) | |
| async def _llm_round(self, speech: asyncio.Queue, spoken: list): | |
| body = {"model": LLM_MODEL, "messages": self.history, "tools": OLLAMA_TOOLS, | |
| "stream": True, "think": False, "keep_alive": KEEP_ALIVE, | |
| "options": {"num_ctx": NUM_CTX, "temperature": TEMPERATURE}} | |
| content, buf, calls = "", "", [] | |
| spoken_any = [False] | |
| async with self.http.stream("POST", f"{OLLAMA_URL}/api/chat", json=body) as r: | |
| if r.status_code >= 400: | |
| raise httpx.HTTPError(f"{r.status_code} {(await r.aread())[:300]!r}") | |
| async for line in r.aiter_lines(): | |
| if not line: | |
| continue | |
| chunk = json.loads(line) | |
| if chunk.get("error"): | |
| raise httpx.HTTPError(chunk["error"]) | |
| msg = chunk.get("message") or {} | |
| calls += msg.get("tool_calls") or [] | |
| delta = msg.get("content") or "" | |
| if delta and self.m_llm is None: | |
| self.m_llm = time.time() - self.m_reply_t0 | |
| if delta: | |
| content += delta | |
| spoken.append(delta) | |
| await self.send({"type": "assistant_delta", "text": delta}) | |
| buf += delta | |
| sentences, buf = split_sentences(buf, first=not self.playing and not spoken_any[0]) | |
| spoken_any[0] = spoken_any[0] or bool(sentences) | |
| for s in sentences: | |
| await speech.put(s) | |
| if buf.strip(): | |
| await speech.put(buf) | |
| content = re.sub(r"<think>.*?</think>", "", content, flags=re.S).strip() | |
| return content, calls | |
| async def _speaker(self, speech: asyncio.Queue): | |
| while (text := await speech.get()) is not None: | |
| text = _spoken(re.sub(r"<think>.*?</think>", "", text, flags=re.S)) | |
| if not re.search(r"\w", text): | |
| continue | |
| t_tts = time.time() | |
| pcm = await self.loop.run_in_executor(TTS_POOL, synthesize, text) | |
| if not self.first_audio_logged and self.speech_end: | |
| self.first_audio_logged = True | |
| total = time.time() - self.speech_end | |
| log(f"1er son {total:.2f} s après ta phrase") | |
| await self.send({"type": "metrics", "stt": round(self.m_stt, 3), | |
| "llm": round(self.m_llm or 0, 3), | |
| "tts": round(time.time() - t_tts, 3), "total": round(total, 3)}) | |
| if not self.playing: | |
| self.playing = True | |
| await self.send({"type": "state", "state": "speaking"}) | |
| await self.send_audio(pcm) | |
| # ------------------------------------------------ tools | |
| async def _tool(self, name: str, args: dict) -> dict: | |
| if name not in ("display_card", "display_report", "get_datetime"): | |
| self.tools_ran = True # real side effect: a late continuation won't undo it | |
| try: | |
| if name == "delegate_task": | |
| task = start_task(args.get("title") or "Tâche", args.get("prompt") or "", self, | |
| complex_=bool(args.get("complex"))) | |
| await self.send({"type": "task", "task": task}) | |
| return {"status": "started", "task_id": task["id"], | |
| "runs_on": "Claude" if task["backend"] == "anthropic" else "modèle local"} | |
| if name == "remember": | |
| fact = (args.get("fact") or "").strip() | |
| facts = load_memory() | |
| if fact and fact not in facts: | |
| facts.append(fact) | |
| save_memory(facts) | |
| self.history[0]["content"] = system_prompt() | |
| return {"ok": True, "memory": facts} | |
| if name == "forget": | |
| kw = (args.get("keyword") or "").strip().lower() | |
| facts = load_memory() | |
| kept = [] if kw in ("tout", "all", "*") else [f for f in facts if kw not in f.lower()] | |
| save_memory(kept) | |
| self.history[0]["content"] = system_prompt() | |
| return {"ok": True, "forgotten": len(facts) - len(kept), "memory": kept} | |
| if name == "notify": | |
| toast(args.get("title") or "J.A.R.V.I.S.", args.get("message") or "") | |
| return {"ok": True} | |
| if name == "get_datetime": | |
| return {"now": time.strftime("%A %d %B %Y, %H:%M"), | |
| "iso": time.strftime("%Y-%m-%dT%H:%M:%S")} | |
| if name in ("open_app", "open_url"): | |
| body = OpenIn(name=args.get("name") or "", url=args.get("url") or None, | |
| monitor=args.get("monitor") or None) | |
| label = body.name or body.url | |
| where = f" → écran **{body.monitor}**" if body.monitor else "" | |
| await self.send({"type": "ui", "name": "display_card", "args": { | |
| "title": "Lancement", "content": f"Ouverture de **{label}**{where}…"}}) | |
| return await self.loop.run_in_executor(None, open_app, body) | |
| if name == "cancel_task": | |
| res = cancel_task(args.get("task_id") or "latest") | |
| if res.get("cancelled"): | |
| await self.send({"type": "task", "task": TASKS[res["cancelled"]]}) | |
| return res | |
| if name in ("display_card", "display_report"): | |
| await self.send({"type": "ui", "name": name, "args": args}) | |
| return {"status": "displayed"} | |
| return {"ok": False, "error": f"outil inconnu: {name}"} | |
| except HTTPException as exc: | |
| return {"ok": False, "error": exc.detail} | |
| except Exception as exc: # noqa: BLE001 | |
| return {"ok": False, "error": str(exc)} | |
| async def task_finished(self, task: dict): | |
| await self.send({"type": "task", "task": task}) | |
| if task["status"] == "cancelled": | |
| return # the model already got the cancel_task output | |
| summary = (task.get("output") or "")[:4000] | |
| self.pending.append( | |
| f'[SYSTEM] Résultat de la tâche "{task["title"]}" ({task["status"]}): {summary}\n' | |
| "Résume oralement en une ou deux phrases. Si c'est une analyse de données " | |
| "(chiffres, stats, comparatifs), affiche un tableau de bord avec display_report " | |
| "(kpis, chart, table). Pour un simple résultat ponctuel, utilise display_card.") | |
| if not self.busy and not self.in_speech: | |
| await self._drain_pending() | |
| async def _drain_pending(self): | |
| if self.pending and not self.busy and not self.in_speech and not self.closed: | |
| await self.start_reply({"role": "user", "content": self.pending.pop(0)}) | |
| async def close(self): | |
| self.closed = True | |
| if self.reply and not self.reply.done(): | |
| self.reply.cancel() | |
| for tid, s in list(TASK_SESSIONS.items()): | |
| if s is self: | |
| TASK_SESSIONS.pop(tid, None) | |
| await self.http.aclose() | |
| async def voice(ws: WebSocket): | |
| await ws.accept() | |
| log("navigateur connecté") | |
| session = Session(ws) | |
| await session.send({"type": "ready", "model": LLM_MODEL, "voice": TTS_VOICE, "address": ADDRESS, | |
| "barge_in": BARGE_IN, "whisper": WHISPER_MODEL, | |
| "agent": AGENT_BACKEND, "complex": COMPLEX_BACKEND}) | |
| try: | |
| while True: | |
| msg = await ws.receive() | |
| if msg["type"] == "websocket.disconnect": | |
| break | |
| if msg.get("bytes"): | |
| await session.feed(msg["bytes"]) | |
| elif msg.get("text"): | |
| await session.on_text(json.loads(msg["text"])) | |
| except WebSocketDisconnect: | |
| pass | |
| finally: | |
| log("navigateur déconnecté") | |
| await session.close() | |
| # ---------------------------------------------------------------- static | |
| _gpu_cache = {"t": 0.0, "data": None} | |
| def status(): | |
| """GPU memory (NVIDIA only, cached a few seconds) for the telemetry panel.""" | |
| if time.time() - _gpu_cache["t"] > 4: | |
| _gpu_cache["t"] = time.time() | |
| try: | |
| out = subprocess.run( | |
| ["nvidia-smi", "--query-gpu=memory.used,memory.total", "--format=csv,noheader,nounits"], | |
| capture_output=True, text=True, timeout=3, | |
| creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0)).stdout | |
| used, total = (int(x) for x in out.splitlines()[0].split(",")) | |
| _gpu_cache["data"] = {"used_mb": used, "total_mb": total} | |
| except Exception: # noqa: BLE001 - no NVIDIA GPU / driver | |
| _gpu_cache["data"] = None | |
| running = sum(1 for t in TASKS.values() if t["status"] == "running") | |
| return {"gpu": _gpu_cache["data"], "tasks_running": running, "model": LLM_MODEL} | |
| def index(): | |
| return FileResponse(ROOT / "index.html") | |
| if __name__ == "__main__": | |
| import uvicorn | |
| port = int(os.environ.get("JARVIS_PORT", "8788")) | |
| print("\n JARVIS Local — chargement des modèles") | |
| warmup() | |
| print(f"\n JARVIS Local -> http://127.0.0.1:{port}\n") | |
| uvicorn.run(app, host="127.0.0.1", port=port, log_level="warning") | |