diff --git a/tai.py b/tai.py index 258c18d..5b5a392 100755 --- a/tai.py +++ b/tai.py @@ -37,17 +37,110 @@ try: except ImportError: readline = None -VERSION = "1.17.0" +try: + import fcntl +except ImportError: + fcntl = None + +# Python 3.6 compatibility notes: +# - subprocess 'text' and 'capture_output' kwargs were added in 3.7, so this +# module only uses 'universal_newlines' plus explicit PIPEs (works on 3.6+). +# - datetime.fromisoformat was added in 3.7, so all parsing goes through +# _fromisoformat below (native when available, strptime fallback on 3.6). + + +def _fromisoformat(text): + """datetime.fromisoformat backport for Python 3.6 (handles Z and +HH:MM).""" + native = getattr(datetime, "fromisoformat", None) + if native is not None: + raw = str(text or "").strip() + if raw[-1:] in ("Z", "z"): + raw = raw[:-1] + "+00:00" + return native(raw) + raw = str(text or "").strip() + if not raw: + raise ValueError("empty datetime") + if raw[-1:] in ("Z", "z"): + raw = raw[:-1] + "+00:00" + tzinfo = None + core = raw + match = re.match(r"^(\d{4}-\d{2}-\d{2}(?:[T ]\d{2}:\d{2}(?::\d{2}(?:\.\d{1,6})?)?)?)([+-]\d{2}:?\d{2}(?::\d{2})?|[+-]\d{4}|[+-]\d{2})$", raw) + if match is not None and match.group(2): + core, tzpart = match.group(1), match.group(2) + sign = 1 if tzpart[0] == "+" else -1 + digits = tzpart[1:].replace(":", "") + hours = int(digits[:2]) + minutes = int(digits[2:4]) if len(digits) >= 4 else 0 + tzinfo = timezone(sign * timedelta(hours=hours, minutes=minutes)) + for fmt in ("%Y-%m-%dT%H:%M:%S.%f", "%Y-%m-%dT%H:%M:%S", "%Y-%m-%dT%H:%M", + "%Y-%m-%d %H:%M:%S.%f", "%Y-%m-%d %H:%M:%S", "%Y-%m-%d %H:%M", + "%Y-%m-%d"): + try: + parsed = datetime.strptime(core, fmt) + except ValueError: + continue + return parsed.replace(tzinfo=tzinfo) + raise ValueError("invalid datetime, use ISO like 2026-10-08T09:00") + +VERSION = "1.22.0" WORKER_STEPS = 12 CREATE_SKILL_STEPS = 40 FORK_MAX_DEPTH = 2 AGENTS = {} AGENTS_LOCK = threading.Lock() AGENTS_NEXT = [1] -PRIMARY_BASE = "https://model.cloud.pravda.education" -PRIMARY_LABEL = "primary" -FALLBACK_BASE = "https://devplace.net/openai/v1" -FALLBACK_LABEL = "fallback" +MUSE_LABEL = "muse" +MUSE_CHAT_URL = "muse://exec" +MUSE_MODELS_URL = "muse://models" +GROK_LABEL = "grok" +GROK_CHAT_URL = "grok://exec" +GROK_MODELS_URL = "grok://models" +CLAUDE_LABEL = "claude" +CLAUDE_CHAT_URL = "claude://exec" +CLAUDE_MODELS_URL = "claude://models" +CODEX_LABEL = "codex" +CODEX_CHAT_URL = "codex://exec" +CODEX_MODELS_URL = "codex://models" +GEMINI_LABEL = "gemini" +GEMINI_CHAT_URL = "gemini://exec" +GEMINI_MODELS_URL = "gemini://models" +OLLAMA_LABEL = "ollama" +OLLAMA_BASE = "http://localhost:11434/v1" +OLLAMA_LIST_TTL = 600 +CLI_CHAT_URLS = (MUSE_CHAT_URL, GROK_CHAT_URL, CLAUDE_CHAT_URL, CODEX_CHAT_URL, GEMINI_CHAT_URL) +MUSE_MAX_STEPS = 1 +MUSE_SESSIONS_FILE = "muse_sessions.json" +MUSE_LOCKS_DIR = "muse_locks" +MUSE_WORKER_TTL = 86400 +MUSE_REAP_GRACE = 30 +OPENROUTER_LABEL = "openrouter" +OPENROUTER_BASE = "https://openrouter.ai/api/v1" +OPENROUTER_MODEL = "deepseek/deepseek-v4.1-flash" +OPENCODE_LABEL = "opencode" +OPENCODE_CHAT_URL = "opencode://run" +OPENCODE_MODELS_URL = "opencode://models" +OPENCODE_AGENT = "plan" +DEVPLACE_LABEL = "devplace" +DEVPLACE_BASE = "https://devplace.net/openai/v1" +DEVPLACE_FREE_ROUTES = ("free", "devii") +ZEN_LABEL = "zen" +ZEN_BASE = "https://opencode.ai/zen/v1" +ZEN_FREE_MODELS = ("mimo-v2.6-flash-free", "nemotron-3.5-lightning-free", "nemotron-3-ultra-free", "muse-spark-1.3-contributor-free", "muse-spark-1.2-contributor-free", "jev-1.13-free", "exo-free", "space-bunny-free", "longcat-2.5-preview-free", "ling-3.1-flash-free", "ling-3.0-flash-fin-free", "fledge-alpha-free") +POLLINATIONS_LABEL = "pollinations" +POLLINATIONS_BASE = "https://text.pollinations.ai/openai/v1" +POLLINATIONS_MODELS = ("openai",) +OPENCODE_USER_AGENT = "opencode/1.18.29 ai-sdk/provider-utils/4.0.23 runtime/bun/1.3.15" +OPENCODE_CLIENT_NAME = "cli" +MODEL_HEALTH_FILE = "model_health.json" +MODEL_NEUTRAL_REWARD = 0.5 +MODEL_WEIGHT_PRIOR = 2.0 +MODEL_SPEED_REF_TPS = 25.0 +MODEL_LATENCY_REF_MS = 8000.0 +MODEL_CIRCUIT_THRESHOLD = 3 +MODEL_CIRCUIT_COOLDOWN = 300.0 +RETRY_BACKOFF = (2, 5, 10, 20, 30, 60) +MODEL_LIST_TIMEOUT = 8 +OPENROUTER_FREE_CACHE_TTL = 3600 RSEARCH_BASE = "https://rsearch.app.molodetz.nl" FETCH_METHODS = ("GET", "POST", "PUT", "PATCH", "DELETE", "HEAD", "OPTIONS") CURL_HINT = "hint: use the web_fetch tool for HTTP(S) (methods, headers, bodies, status) instead of curl/wget via shell" @@ -57,7 +150,7 @@ def curl_hint(command): if re.search(r"\bcurl\b|\bwget\b", command or "") and "http" in (command or "").lower(): return "\n" + CURL_HINT return "" -DEFAULT_MODEL = "openrouter/free" +DEFAULT_MODEL = "muse" DEFAULT_VOICE = "en-US-EmmaMultilingualNeural" DEFAULT_PASSPHRASE = "tai-default-insecure-change-me" BOX_IMAGE = "tai-box:latest" @@ -159,7 +252,7 @@ SAFE_COMMANDS = ("ls", "pwd", "echo", "cat", "head", "tail", "grep", "find", "wc RECORDERS = (("arecord", ("arecord", "-q", "-d", "{seconds}", "-f", "cd", "-t", "wav", "{path}")), ("rec", ("rec", "-q", "{path}", "trim", "0", "{seconds}")), ("ffmpeg", ("ffmpeg", "-y", "-v", "quiet", "-f", "alsa", "-i", "default", "-t", "{seconds}", "{path}"))) TRANSCRIBERS = ("whisper-cpp", "whisper-cli", "whisper", "faster-whisper") SYSINFO_TIMEOUT = 10 -SYSINFO_BINARIES = ("tmux", "git", "ffmpeg", "arecord", "sox", "whisper-cpp", "whisper", "systemctl", "podman", "docker") +SYSINFO_BINARIES = ("tmux", "git", "ffmpeg", "arecord", "sox", "whisper-cpp", "whisper", "systemctl", "podman", "docker", "muse", "grok", "opencode", "claude", "codex", "gemini", "ollama") DEFAULT_SYSTEM = ( "You are tai, a professional autonomous assistant. You act with tools, verify results, and report concisely. " @@ -201,6 +294,7 @@ HELP_TEXT = ( " /skills list loaded skill files\n" " /secret manage sealed secrets (set|list|delete)\n" " /sysinfo show host environment checks\n" + " /models show model roster with speed health\n" " /schedule at|every schedule a prompt for later\n" " /schedules list scheduled prompts\n" " /unschedule delete a scheduled prompt\n" @@ -661,7 +755,7 @@ class LiveStream: def run_live(argv, timeout=SHELL_TIMEOUT, shell=False, cwd=None, env=None, stream=None): started = time.time() - proc = subprocess.Popen(argv, shell=shell, cwd=cwd, env=env, stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, errors="replace", bufsize=1) + proc = subprocess.Popen(argv, shell=shell, cwd=cwd, env=env, stdin=subprocess.DEVNULL, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, universal_newlines=True, errors="replace", bufsize=1) chunks = [] pending = queue.Queue() @@ -1391,7 +1485,7 @@ def container_engine(): def box_state(engine): try: - done = subprocess.run([engine, "inspect", "-f", "{{.State.Running}}", BOX_NAME], capture_output=True, text=True, timeout=15) + done = subprocess.run([engine, "inspect", "-f", "{{.State.Running}}", BOX_NAME], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=15) except (OSError, subprocess.SubprocessError): return "missing" if done.returncode != 0: @@ -1407,7 +1501,7 @@ def ensure_box(): if state == "running": return engine if state == "stopped": - done = subprocess.run([engine, "start", BOX_NAME], capture_output=True, text=True, timeout=60) + done = subprocess.run([engine, "start", BOX_NAME], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=60) if done.returncode != 0: raise BackendError("cannot start box: " + (done.stderr or "").strip()[:200]) return engine @@ -1419,14 +1513,14 @@ def ensure_box(): done = subprocess.run([engine, "build", "-t", BOX_IMAGE, build], timeout=1800) if done.returncode != 0: raise BackendError("box image build failed") - done = subprocess.run([engine, "run", "-d", "--name", BOX_NAME, "--restart", "unless-stopped", BOX_IMAGE], capture_output=True, text=True, timeout=120) + done = subprocess.run([engine, "run", "-d", "--name", BOX_NAME, "--restart", "unless-stopped", BOX_IMAGE], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=120) if done.returncode != 0: raise BackendError("cannot start box: " + (done.stderr or "").strip()[:200]) return engine def box_exec(engine, argv, extra=(), input_data=None, timeout=180): - return subprocess.run([engine, "exec", *extra, "-i", BOX_NAME, *argv], input=input_data, capture_output=True, timeout=timeout) + return subprocess.run([engine, "exec", *extra, "-i", BOX_NAME, *argv], input=input_data, stdout=subprocess.PIPE, stderr=subprocess.PIPE, timeout=timeout) def box_python(engine): @@ -1454,6 +1548,29 @@ def box_speak(engine, text, voice=""): return done.stdout +def split_models(spec): + return [item.strip() for item in (spec or "").replace("\n", ",").split(",") if item.strip()] + + +def parse_backends(spec): + found = [] + for entry in (spec or "").replace("\n", ";").split(";"): + fields = [field.strip() for field in entry.split("|")] + if len(fields) < 3: + continue + label, base, key = fields[0], fields[1].rstrip("/"), fields[2] + if key.startswith("$"): + key = os.environ.get(key[1:], "") + model = fields[3].strip() if len(fields) > 3 else "" + profile = fields[4].strip().lower() if len(fields) > 4 else "" + if profile not in ("", "opencode"): + profile = "" + if not label or not base: + continue + found.append({"label": label, "base": base, "key": key, "model": model or "", "profile": profile}) + return found + + class Config: def __init__(self, args): self.home = os.environ.get("TAI_HOME") or os.path.join(os.path.expanduser("~"), ".tai") @@ -1462,7 +1579,21 @@ class Config: self.db_path = os.path.join(self.home, "memory.db") self.history_path = os.path.join(self.home, "history") self.model = os.environ.get("TAI_MODEL") or DEFAULT_MODEL + self.explicit_model = os.environ.get("TAI_MODEL") or "" + self.backends = parse_backends(os.environ.get("TAI_BACKENDS")) + self.openrouter_key = os.environ.get("OPENROUTER_API_KEY") or "" + self.openrouter_models = split_models(os.environ.get("TAI_OPENROUTER_MODELS")) + self.retry_seconds = float(os.environ.get("TAI_RETRY_SECONDS") or 0) if re.match(r"^\d+(\.\d+)?$", (os.environ.get("TAI_RETRY_SECONDS") or "0").strip()) else 0.0 + self.opencode_enabled = (os.environ.get("TAI_OPENCODE") or "1").strip().lower() not in ("0", "false", "no", "off") + self.opencode_models = split_models(os.environ.get("TAI_OPENCODE_MODELS")) self.devplace_key = os.environ.get("DEVPLACE_API_KEY") or "" + self.devplace_base = (os.environ.get("TAI_DEVPLACE_BASE") or DEVPLACE_BASE).rstrip("/") + self.devplace_models = split_models(os.environ.get("TAI_DEVPLACE_MODELS")) + self.openrouter_enabled = (os.environ.get("TAI_OPENROUTER") or "1").strip().lower() not in ("0", "false", "no", "off") + self.pollinations_enabled = (os.environ.get("TAI_POLLINATIONS") or "1").strip().lower() not in ("0", "false", "no", "off") + self.pollinations_models = split_models(os.environ.get("TAI_POLLINATIONS_MODELS")) or list(POLLINATIONS_MODELS) + self.zen_enabled = (os.environ.get("TAI_ZEN") or "").strip().lower() in ("1", "true", "yes", "on") + self.zen_models = split_models(os.environ.get("TAI_ZEN_MODELS")) or list(ZEN_FREE_MODELS) self.voice = os.environ.get("TAI_VOICE") or DEFAULT_VOICE self.auto_approve = bool(args.yes) self.profile = args.profile or "default" @@ -2428,39 +2559,1157 @@ class BackendError(Exception): self.status = status +def muse_binary(): + return os.environ.get("MUSE_BIN") or shutil.which("muse") + + +def opencode_binary(): + return os.environ.get("OPENCODE_BIN") or shutil.which("opencode") + + +def grok_binary(): + return os.environ.get("GROK_BIN") or shutil.which("grok") or shutil.which("grokcli") + + +def claude_binary(): + return os.environ.get("CLAUDE_BIN") or shutil.which("claude") + + +def codex_binary(): + return os.environ.get("CODEX_BIN") or shutil.which("codex") + + +def gemini_binary(): + return os.environ.get("GEMINI_BIN") or shutil.which("gemini") + + +def opencode_prompt(messages, tools): + schema = json.dumps(muse_output_schema()) + rules = ( + "You are the model behind another agent. The transcript above is the whole conversation; lines starting with TOOL are results of tool calls that already ran. " + "Do not use any of your own tools and do not repeat a tool call whose result is already shown. " + "If the results answer the user, put the final answer in content with an empty tool_calls list; otherwise request the next tool calls from TOOLS. " + "Reply with one JSON object only, no code fence, matching this schema: " + ) + return muse_prompt(messages, tools) + "\n\n" + rules + schema + + +def opencode_parse(text): + text = (text or "").strip() + fenced = re.match(r"^```(?:json)?\s*(.*?)\s*```$", text, re.S) + if fenced: + text = fenced.group(1) + start, end = text.find("{"), text.rfind("}") + if start != -1 and end > start: + try: + answer = json.loads(text[start:end + 1]) + if isinstance(answer, dict) and ("content" in answer or "tool_calls" in answer): + return answer + except ValueError: + pass + return {"content": text, "tool_calls": []} + + +def opencode_run(model, prompt, timeout): + binary = opencode_binary() + if not binary: + raise BackendError("opencode CLI not found (OPENCODE_BIN or PATH)") + argv = [binary, "run", "--pure", "--agent", OPENCODE_AGENT, "--format", "json", "-m", "opencode/" + model] + with tempfile.TemporaryDirectory(prefix="tai-opencode-") as folder: + try: + done = subprocess.run(argv, input=prompt, stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, cwd=folder, timeout=timeout) + except subprocess.TimeoutExpired: + raise BackendError("opencode run timed out after %ds" % timeout) + except OSError as exc: + raise BackendError("opencode run failed: %s" % short_error(exc)) + texts = [] + for line in done.stdout.splitlines(): + try: + event = json.loads(line) + except ValueError: + continue + if not isinstance(event, dict): + continue + if event.get("type") == "error": + data = (event.get("error") or {}).get("data") or {} + raise BackendError("opencode: %s" % (data.get("message") or json.dumps(event.get("error"))[:300]), data.get("statusCode") or 0) + part = event.get("part") or {} + if event.get("type") == "text" and part.get("text"): + texts.append(part["text"]) + if not texts: + raise BackendError("opencode returned no text (exit %d): %s" % (done.returncode, (done.stderr or done.stdout)[-300:])) + return "".join(texts) + + +def opencode_free_models(): + binary = opencode_binary() + if not binary: + return None + try: + done = subprocess.run([binary, "models", "opencode"], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=MODEL_LIST_TIMEOUT * 4) + except (OSError, subprocess.TimeoutExpired): + return None + models = [line.strip().split("/", 1)[1] for line in done.stdout.splitlines() if line.strip().startswith("opencode/")] + return [item for item in models if is_free_model(item)] or None + + +def _cli_unknown_flag(stderr): + return bool(re.search(r"unknown (option|flag|command)|unrecognised|unrecognized|invalid flag|not found|no such option", stderr or "", re.IGNORECASE)) + + +def _cli_auth_failure(text): + return bool(re.search(r"not (logged in|authenticated)|auth|api key|unauthorized|401|login required|please (login|log in|authenticate)", text or "", re.IGNORECASE)) + + +def _cli_run(argv, prompt_input, timeout, label): + try: + done = subprocess.run(argv, input=prompt_input, stdin=subprocess.DEVNULL if prompt_input is None else None, stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=timeout) + except subprocess.TimeoutExpired: + raise BackendError("%s run timed out after %ds" % (label, timeout)) + except OSError as exc: + raise BackendError("%s run failed: %s" % (label, short_error(exc))) + return done + + +def grok_parse_output(stdout): + try: + data = json.loads(stdout or "") + except ValueError: + return opencode_parse(stdout) + if isinstance(data, dict): + if data.get("is_error") or data.get("error"): + detail = data.get("error") or data.get("result") or data.get("text") or "unknown error" + raise BackendError("grok: %s" % str(detail)[:300]) + if isinstance(data.get("content"), str) or isinstance(data.get("tool_calls"), list): + if "content" in data or "tool_calls" in data: + return data + for key in ("text", "result", "output", "answer", "content"): + value = data.get(key) + if isinstance(value, str) and value.strip(): + return opencode_parse(value) + return opencode_parse(stdout) + + +def grok_run(prompt, timeout): + binary = grok_binary() + if not binary: + raise BackendError("grok CLI not found (GROK_BIN or PATH)") + with tempfile.TemporaryDirectory(prefix="tai-grok-") as folder: + prompt_path = os.path.join(folder, "prompt.txt") + with open(prompt_path, "w", encoding="utf-8") as handle: + handle.write(prompt) + variants = [ + [binary, "--prompt-file", prompt_path, "--output-format", "json", "--no-auto-update"], + [binary, "--prompt-file", prompt_path, "--output-format", "json"], + [binary, "-p", prompt, "--output-format", "json"], + [binary, "-p", prompt], + ] + last_error = "" + for argv in variants: + done = _cli_run(argv, None, timeout, "grok") + if done.returncode != 0 and _cli_unknown_flag(done.stderr) and argv is not variants[-1]: + last_error = (done.stderr or "").strip()[:200] + continue + if done.returncode != 0 and not (done.stdout or "").strip(): + blob = ((done.stderr or "") + "\n" + (done.stdout or "")).strip() + if _cli_auth_failure(blob): + raise BackendError("grok not authenticated (set XAI_API_KEY or run grok login): %s" % blob[:200]) + raise BackendError("grok exit %d: %s" % (done.returncode, blob[:300] or last_error or "no output")) + if not (done.stdout or "").strip(): + blob = (done.stderr or "").strip() + raise BackendError("grok returned no output: %s" % (blob[:300] or last_error or "empty stdout")) + return done.stdout or "" + raise BackendError("grok failed: %s" % (last_error or "no usable output")) + + +def claude_parse_output(stdout): + try: + data = json.loads(stdout or "") + except ValueError: + return opencode_parse(stdout) + if not isinstance(data, dict): + return opencode_parse(stdout) + if data.get("is_error"): + raise BackendError("claude: %s" % str(data.get("result") or "unknown error")[:300]) + structured = data.get("structured_output") + if isinstance(structured, dict) and ("content" in structured or "tool_calls" in structured): + return structured + result = data.get("result") + if isinstance(result, str) and result.strip(): + return opencode_parse(result) + return opencode_parse(stdout) + + +def claude_run(prompt, timeout): + binary = claude_binary() + if not binary: + raise BackendError("claude CLI not found (CLAUDE_BIN or PATH)") + schema = json.dumps(muse_output_schema()) + base = [binary, "-p"] + if len(prompt.encode("utf-8", "replace")) > 120000: + variants = [ + (base + ["--output-format", "json", "--json-schema", schema], prompt), + (base + ["--output-format", "json"], prompt), + ] + else: + variants = [ + (base + [prompt, "--output-format", "json", "--json-schema", schema], None), + (base + [prompt, "--output-format", "json"], None), + (base + [prompt], None), + ] + last_error = "" + for argv, prompt_input in variants: + done = _cli_run(argv, prompt_input, timeout, "claude") + if done.returncode != 0 and _cli_unknown_flag(done.stderr) and (argv, prompt_input) != variants[-1]: + last_error = (done.stderr or "").strip()[:200] + continue + if done.returncode != 0 and not (done.stdout or "").strip(): + blob = ((done.stderr or "") + "\n" + (done.stdout or "")).strip() + if _cli_auth_failure(blob): + raise BackendError("claude not authenticated (run claude login): %s" % blob[:200]) + raise BackendError("claude exit %d: %s" % (done.returncode, blob[:300] or last_error or "no output")) + if not (done.stdout or "").strip(): + blob = (done.stderr or "").strip() + raise BackendError("claude returned no output: %s" % (blob[:300] or last_error or "empty stdout")) + return done.stdout or "" + raise BackendError("claude failed: %s" % (last_error or "no usable output")) + + +def codex_jsonl_texts(stdout): + texts = [] + last_error = "" + for line in (stdout or "").splitlines(): + line = line.strip() + if not line.startswith("{"): + continue + try: + event = json.loads(line) + except ValueError: + continue + if not isinstance(event, dict): + continue + etype = event.get("type") + if etype == "item.completed": + item = event.get("item") or {} + if item.get("type") == "agent_message" and item.get("text"): + texts.append(item["text"]) + elif etype in ("turn.failed", "error"): + msg = event.get("message") or event.get("error") or "" + if isinstance(msg, dict): + msg = msg.get("message") or json.dumps(msg) + if msg: + last_error = str(msg)[:300] + return texts, last_error + + +def codex_parse_output(stdout, outfile_text="", used_schema=False): + if isinstance(outfile_text, str) and outfile_text.strip(): + if used_schema: + try: + answer = json.loads(outfile_text) + except ValueError: + answer = None + if isinstance(answer, dict) and ("content" in answer or "tool_calls" in answer): + return answer + return opencode_parse(outfile_text) + texts, last_error = codex_jsonl_texts(stdout) + if texts: + return opencode_parse("".join(texts)) + if last_error: + raise BackendError("codex: %s" % last_error) + return opencode_parse(stdout) + + +def codex_run(prompt, timeout): + binary = codex_binary() + if not binary: + raise BackendError("codex CLI not found (CODEX_BIN or PATH)") + schema = json.dumps(muse_output_schema()) + with tempfile.TemporaryDirectory(prefix="tai-codex-") as folder: + schema_path = os.path.join(folder, "schema.json") + out_path = os.path.join(folder, "last.txt") + with open(schema_path, "w", encoding="utf-8") as handle: + handle.write(schema) + base = [binary, "exec", "--json", "--ephemeral", "--sandbox", "read-only", "--skip-git-repo-check"] + if len(prompt.encode("utf-8", "replace")) > 120000: + variants = [ + (base + ["--output-schema", schema_path, "-o", out_path, "-"], prompt, True), + (base + ["-o", out_path, "-"], prompt, False), + ] + else: + variants = [ + (base + ["--output-schema", schema_path, "-o", out_path, prompt], None, True), + (base + ["-o", out_path, prompt], None, False), + (base + [prompt], None, False), + ] + last_error = "" + for argv, prompt_input, used_schema in variants: + try: + done = subprocess.run(argv, input=prompt_input, stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, cwd=folder, timeout=timeout) + except subprocess.TimeoutExpired: + raise BackendError("codex run timed out after %ds" % timeout) + except OSError as exc: + raise BackendError("codex run failed: %s" % short_error(exc)) + if done.returncode != 0 and _cli_unknown_flag(done.stderr) and (argv, prompt_input, used_schema) != variants[-1]: + last_error = (done.stderr or "").strip()[:200] + continue + outfile_text = "" + if "-o" in argv: + try: + with open(out_path, "r", encoding="utf-8", errors="replace") as handle: + outfile_text = handle.read() + except OSError: + outfile_text = "" + if done.returncode != 0 and not (outfile_text.strip() or (done.stdout or "").strip()): + blob = ((done.stderr or "") + "\n" + (done.stdout or "")).strip() + if _cli_auth_failure(blob): + raise BackendError("codex not authenticated (run codex login): %s" % blob[:200]) + raise BackendError("codex exit %d: %s" % (done.returncode, blob[:300] or last_error or "no output")) + if not (outfile_text.strip() or (done.stdout or "").strip()): + blob = (done.stderr or "").strip() + raise BackendError("codex returned no output: %s" % (blob[:300] or last_error or "empty stdout")) + return done.stdout or "", outfile_text, used_schema + raise BackendError("codex failed: %s" % (last_error or "no usable output")) + + +def gemini_parse_output(stdout): + try: + data = json.loads(stdout or "") + except ValueError: + return opencode_parse(stdout) + if not isinstance(data, dict): + return opencode_parse(stdout) + err = data.get("error") + if err: + raise BackendError("gemini: %s" % (json.dumps(err) if isinstance(err, dict) else str(err))[:300]) + response = data.get("response") + if isinstance(response, str) and response.strip(): + return opencode_parse(response) + return opencode_parse(stdout) + + +def gemini_run(prompt, timeout): + binary = gemini_binary() + if not binary: + raise BackendError("gemini CLI not found (GEMINI_BIN or PATH)") + base = [binary, "-p", prompt] + variants = [ + base + ["--output-format", "json", "--approval-mode", "plan"], + base + ["--output-format", "json"], + base, + ] + last_error = "" + for argv in variants: + done = _cli_run(argv, None, timeout, "gemini") + if done.returncode != 0 and _cli_unknown_flag(done.stderr) and argv is not variants[-1]: + last_error = (done.stderr or "").strip()[:200] + continue + if done.returncode != 0 and not (done.stdout or "").strip(): + blob = ((done.stderr or "") + "\n" + (done.stdout or "")).strip() + if _cli_auth_failure(blob): + raise BackendError("gemini not authenticated (run gemini login or set GEMINI_API_KEY): %s" % blob[:200]) + hint = {42: "input error", 53: "turn limit exceeded"}.get(done.returncode, "") + raise BackendError("gemini exit %d%s: %s" % (done.returncode, " (%s)" % hint if hint else "", blob[:300] or last_error or "no output")) + if not (done.stdout or "").strip(): + blob = (done.stderr or "").strip() + raise BackendError("gemini returned no output: %s" % (blob[:300] or last_error or "empty stdout")) + return done.stdout or "" + raise BackendError("gemini failed: %s" % (last_error or "no usable output")) + + +def muse_output_schema(): + call = {"type": "object", "additionalProperties": False, "properties": {"name": {"type": "string"}, "arguments": {"type": "string"}}, "required": ["name", "arguments"]} + return {"type": "object", "additionalProperties": False, "properties": {"content": {"type": "string"}, "tool_calls": {"type": "array", "items": call}}, "required": ["content", "tool_calls"]} + + +def muse_prompt(messages, tools): + parts = [] + for message in messages or []: + role = (message.get("role") or "user").upper() + body = message.get("content") or "" + if isinstance(body, list): + body = " ".join(str(part.get("text") or "") for part in body if isinstance(part, dict)) + if message.get("tool_call_id"): + parts.append("%s %s:\n%s" % (role, message["tool_call_id"], body)) + else: + parts.append("%s:\n%s" % (role, body)) + for call in message.get("tool_calls") or []: + func = call.get("function") or {} + parts.append("TOOL CALL %s:\n%s %s" % (call.get("id") or "", func.get("name") or "", func.get("arguments") or "")) + if tools: + parts.append("TOOLS:\n%s" % json.dumps(tools)) + parts.append("Answer with content text plus tool_calls for every tool to invoke now. Empty tool_calls ends the turn.") + return "\n\n".join(part for part in parts if part.strip()) + + +def muse_argv(config, schema_path, prompt_path, session_id=None): + argv = [muse_binary(), "exec", "--json", "--max-model-steps", str(MUSE_MAX_STEPS), "--disable-web-tools", "--disable-reminders", "--allow-workspace-switch", "--output-schema", schema_path, "--prompt-file", prompt_path] + if session_id is None: + argv.append("--no-session-log") + else: + argv.extend(["--session-id", session_id]) + if config.explicit_model: + argv.extend(["--model", config.explicit_model]) + return argv + + +def muse_write_file(path, text): + with open(path, "w", encoding="utf-8") as handle: + handle.write(text) + handle.flush() + os.fsync(handle.fileno()) + + +def muse_message_hash(message): + blob = json.dumps(message, sort_keys=True, separators=(",", ":"), default=str) + return hashlib.sha256(blob.encode("utf-8")).hexdigest() + + +def muse_map_path(home): + return os.path.join(home, MUSE_SESSIONS_FILE) + + +def muse_load_map(home): + path = muse_map_path(home) + try: + with open(path, "r", encoding="utf-8") as handle: + data = json.load(handle) + except FileNotFoundError: + return {} + except (OSError, ValueError): + quarantine(path) + return {} + if not isinstance(data, dict): + quarantine(path) + return {} + now = time.time() + pruned = False + for key in list(data): + entry = data[key] + if not isinstance(entry, dict) or not isinstance(entry.get("id"), str) or not isinstance(entry.get("hashes"), list): + del data[key] + pruned = True + elif key.startswith("worker/"): + try: + stamp = _fromisoformat(entry.get("updated") or "").timestamp() + except ValueError: + stamp = 0 + if now - stamp > MUSE_WORKER_TTL: + del data[key] + pruned = True + if pruned: + muse_save_map(home, data) + return data + + +def muse_save_map(home, data): + muse_write_file(muse_map_path(home), json.dumps(data, sort_keys=True) + "\n") + + +_MUSE_THREAD_LOCKS = {} +_MUSE_THREAD_GUARD = threading.Lock() + + +def muse_acquire(home, name): + with _MUSE_THREAD_GUARD: + thread_lock = _MUSE_THREAD_LOCKS.setdefault(name, threading.Lock()) + thread_lock.acquire() + handle = None + if fcntl is not None: + try: + folder = os.path.join(home, MUSE_LOCKS_DIR) + os.makedirs(folder, exist_ok=True) + handle = open(os.path.join(folder, hashlib.sha256(name.encode("utf-8")).hexdigest() + ".lock"), "w") + fcntl.flock(handle.fileno(), fcntl.LOCK_EX) + except OSError: + if handle is not None: + try: + handle.close() + except OSError: + pass + handle = None + + def release(): + try: + if handle is not None: + if fcntl is not None: + try: + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + except OSError: + pass + try: + handle.close() + except OSError: + pass + finally: + thread_lock.release() + + return release + + +def muse_transact(home, func): + release = muse_acquire(home, "map") + try: + data = muse_load_map(home) + result = func(data) + muse_save_map(home, data) + return result + finally: + release() + + +def muse_session_root(): + base = os.environ.get("XDG_DATA_HOME") or os.path.join(os.path.expanduser("~"), ".local", "share") + return os.path.join(base, "muse", "sessions") + + +def muse_session_alive(session_id): + try: + pattern = os.path.join(muse_session_root(), "*", "*", "*", session_id, "session.jsonl") + return bool(glob.glob(pattern)) + except (OSError, ValueError): + return True + + +def muse_session_split(messages, hashes): + if len(hashes) > len(messages): + return None + for pos, digest in enumerate(hashes): + if muse_message_hash(messages[pos]) != digest: + return None + return list(messages[len(hashes):]) + + +def muse_valid_answer(text): + try: + data = json.loads(text) + except ValueError: + return None + if not isinstance(data, dict): + return None + if not isinstance(data.get("content"), str) or not isinstance(data.get("tool_calls"), list): + return None + for call in data["tool_calls"]: + if not isinstance(call, dict) or not isinstance(call.get("name"), str) or not isinstance(call.get("arguments"), str): + return None + return data + + +def muse_gate(func): + def gated(ok): + func(ok) + + gated.armed = threading.Event() + return gated + + +def muse_reap(proc, stderr_done, on_reap): + def target(): + ok = False + try: + try: + for _line in proc.stdout: + pass + except (OSError, ValueError): + pass + try: + proc.stdout.close() + except (OSError, ValueError): + pass + try: + ok = proc.wait(timeout=MUSE_REAP_GRACE) == 0 + except subprocess.TimeoutExpired: + try: + proc.kill() + except OSError: + pass + try: + proc.wait(timeout=10) + except (OSError, subprocess.SubprocessError): + pass + try: + stderr_done.wait(timeout=10) + except (OSError, RuntimeError): + pass + finally: + if on_reap is not None: + try: + on_reap(ok) + except Exception: + pass + + thread = threading.Thread(target=target, daemon=True) + thread.start() + return thread + + +def muse_run(argv, cwd, stream_sink, timeout, on_reap=None): + try: + proc = subprocess.Popen(argv, stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, cwd=cwd, bufsize=1) + except OSError as exc: + raise BackendError("muse exec failed: %s" % short_error(exc)) + gate = getattr(on_reap, "armed", None) + if gate is not None: + gate.set() + stderr_parts = [] + stderr_done = threading.Event() + + def collect_stderr(): + try: + for chunk in proc.stderr: + stderr_parts.append(chunk) + except (OSError, ValueError): + pass + finally: + stderr_done.set() + + threading.Thread(target=collect_stderr, daemon=True).start() + finished = threading.Event() + expired = threading.Event() + + def watchdog(): + if not finished.wait(timeout): + expired.set() + try: + proc.kill() + except OSError: + pass + + threading.Thread(target=watchdog, daemon=True).start() + answer = None + terminal = None + buffer = [] + try: + for line in proc.stdout: + text = line.strip() + if not text.startswith("{"): + continue + try: + event = json.loads(text) + except ValueError: + continue + kind = event.get("payload_type") or "" + payload = event.get("payload") or {} + if kind == "run.output.delta" and payload.get("text"): + piece = payload["text"] + if stream_sink is not None: + stream_sink(piece) + buffer.append(piece) + if answer is None: + answer = muse_valid_answer("".join(buffer)) + if answer is not None: + break + elif kind.startswith("run.terminal."): + terminal = payload + break + finally: + finished.set() + muse_reap(proc, stderr_done, on_reap) + if terminal is not None: + if terminal.get("terminal") != "completed": + raise BackendError("muse %s: %s" % (terminal.get("terminal") or "failed", (terminal.get("reason") or "unknown")[:300])) + parsed = muse_valid_answer(terminal.get("text") or "") + if parsed is None: + raise BackendError("muse returned invalid structured answer") + return parsed + if answer is not None: + return answer + if expired.is_set(): + raise BackendError("muse timed out after %ds" % timeout) + try: + code = proc.returncode + except (OSError, ValueError): + code = "?" + raise BackendError("muse exit %s: %s" % (code, "".join(stderr_parts).strip()[-300:] or "no terminal event")) + + +def muse_health(): + binary = muse_binary() + if not binary: + return "missing: muse CLI not found" + try: + done = subprocess.run([binary, "--version"], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=15) + except (OSError, subprocess.SubprocessError) as exc: + return "fail: %s" % short_error(exc)[:60] + if done.returncode != 0: + return "fail: exit %d" % done.returncode + return "ok, subscription (%s)" % (done.stdout.strip().splitlines()[0][:60] if done.stdout.strip() else "muse") + + +def grok_health(): + binary = grok_binary() + if not binary: + return "missing: grok CLI not found" + try: + done = subprocess.run([binary, "--version"], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=15) + except (OSError, subprocess.SubprocessError) as exc: + return "fail: %s" % short_error(exc)[:60] + if done.returncode != 0: + return "fail: exit %d" % done.returncode + return "ok, subscription (%s)" % ((done.stdout.strip() or done.stderr.strip()).splitlines()[0][:60] if (done.stdout.strip() or done.stderr.strip()) else "grok") + + +def claude_health(): + binary = claude_binary() + if not binary: + return "missing: claude CLI not found" + try: + done = subprocess.run([binary, "--version"], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=15) + except (OSError, subprocess.SubprocessError) as exc: + return "fail: %s" % short_error(exc)[:60] + if done.returncode != 0: + return "fail: exit %d" % done.returncode + return "ok, subscription (%s)" % ((done.stdout.strip() or done.stderr.strip()).splitlines()[0][:60] if (done.stdout.strip() or done.stderr.strip()) else "claude") + + +def codex_health(): + binary = codex_binary() + if not binary: + return "missing: codex CLI not found" + try: + done = subprocess.run([binary, "--version"], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=15) + except (OSError, subprocess.SubprocessError) as exc: + return "fail: %s" % short_error(exc)[:60] + if done.returncode != 0: + return "fail: exit %d" % done.returncode + return "ok, subscription (%s)" % ((done.stdout.strip() or done.stderr.strip()).splitlines()[0][:60] if (done.stdout.strip() or done.stderr.strip()) else "codex") + + +def gemini_health(): + binary = gemini_binary() + if not binary: + return "missing: gemini CLI not found" + try: + done = subprocess.run([binary, "--version"], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=15) + except (OSError, subprocess.SubprocessError) as exc: + return "fail: %s" % short_error(exc)[:60] + if done.returncode != 0: + return "fail: exit %d" % done.returncode + return "ok, subscription (%s)" % ((done.stdout.strip() or done.stderr.strip()).splitlines()[0][:60] if (done.stdout.strip() or done.stderr.strip()) else "gemini") + + +_OPENCODE_ID_ALPHABET = "0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz" +_OPENCODE_SESSION = "ses_" + "".join(random.choice(_OPENCODE_ID_ALPHABET) for _ in range(26)) + + +def opencode_headers(): + return { + "User-Agent": os.environ.get("OPENCODE_CLIENT_USER_AGENT") or OPENCODE_USER_AGENT, + "x-opencode-client": os.environ.get("OPENCODE_CLIENT_NAME") or OPENCODE_CLIENT_NAME, + "x-opencode-project": "global", + "x-opencode-session": _OPENCODE_SESSION, + "x-opencode-request": "msg_" + "".join(random.choice(_OPENCODE_ID_ALPHABET) for _ in range(26)), + } + + +def model_speed_reward(tokens_per_second): + if not tokens_per_second or tokens_per_second <= 0: + return MODEL_NEUTRAL_REWARD + return tokens_per_second / (tokens_per_second + MODEL_SPEED_REF_TPS) + + +def model_latency_reward(latency_ms): + if not latency_ms or latency_ms <= 0: + return 1.0 + return 1.0 / (1.0 + latency_ms / MODEL_LATENCY_REF_MS) + + +def model_outcome_reward(latency_ms=None, tokens_per_second=None): + return model_speed_reward(tokens_per_second) * model_latency_reward(latency_ms) + + +def model_weight(entry): + entry = entry or {} + attempts = (entry.get("s") or 0) + (entry.get("f") or 0) + return ((entry.get("r") or 0.0) + MODEL_NEUTRAL_REWARD * MODEL_WEIGHT_PRIOR) / (attempts + MODEL_WEIGHT_PRIOR) + + +class ModelHealth: + def __init__(self, home): + self.path = os.path.join(home, MODEL_HEALTH_FILE) if home else "" + self.data = {"models": {}, "lists": {}} + if self.path: + try: + with open(self.path, "r", encoding="utf-8") as handle: + raw = json.load(handle) + if isinstance(raw, dict): + if isinstance(raw.get("models"), dict): + self.data["models"] = raw["models"] + if isinstance(raw.get("lists"), dict): + self.data["lists"] = raw["lists"] + except (OSError, ValueError): + pass + + def _save(self): + if not self.path: + return + try: + tmp = self.path + ".tmp" + with open(tmp, "w", encoding="utf-8") as handle: + json.dump(self.data, handle) + os.replace(tmp, self.path) + except OSError: + pass + + def entry(self, provider, model): + return self.data["models"].get("%s|%s" % (provider or "default", model or "")) or {} + + def weight(self, provider, model): + return model_weight(self.entry(provider, model)) + + def circuit_open(self, provider, model): + return (self.entry(provider, model).get("until") or 0.0) > time.time() + + def record(self, provider, model, success, latency_ms=None, tokens_per_second=None): + if not model: + return + key = "%s|%s" % (provider or "default", model or "") + entry = self.data["models"].setdefault(key, {"s": 0, "f": 0, "r": 0.0, "lat": 0.0, "tps": 0.0, "cf": 0, "until": 0.0, "used": 0.0, "n": 0}) + entry["n"] += 1 + entry["used"] = time.time() + if success: + reward = MODEL_NEUTRAL_REWARD if latency_ms is None and tokens_per_second is None else model_outcome_reward(latency_ms, tokens_per_second) + entry["s"] += 1 + entry["cf"] = 0 + entry["until"] = 0.0 + entry["r"] += reward + if latency_ms is not None: + entry["lat"] += latency_ms + if tokens_per_second is not None: + entry["tps"] += tokens_per_second + else: + entry["f"] += 1 + entry["cf"] += 1 + if entry["cf"] >= MODEL_CIRCUIT_THRESHOLD: + entry["until"] = time.time() + MODEL_CIRCUIT_COOLDOWN + self._save() + + def cached_list(self, name, ttl=OPENROUTER_FREE_CACHE_TTL): + row = self.data["lists"].get(name) + if not isinstance(row, dict): + return None + models = row.get("models") + if not models or (time.time() - (row.get("at") or 0)) > ttl: + return None + return [str(item) for item in models if str(item).strip()] + + def store_list(self, name, models): + cleaned = [str(item).strip() for item in models if str(item).strip()] + if not cleaned: + return + self.data["lists"][name] = {"at": time.time(), "models": cleaned} + self._save() + + def snapshot(self): + rows = [] + for key, entry in self.data["models"].items(): + provider, _, model = key.partition("|") + ok = entry.get("s") or 0 + rows.append({ + "provider": provider, + "model": model, + "weight": round(model_weight(entry), 4), + "ok": ok, + "fail": entry.get("f") or 0, + "avg_ms": round((entry.get("lat") or 0.0) / ok) if ok and entry.get("lat") else None, + "avg_tps": round((entry.get("tps") or 0.0) / ok, 1) if ok and entry.get("tps") else None, + "circuit_open": (entry.get("until") or 0.0) > time.time(), + "used": entry.get("used") or 0.0, + }) + rows.sort(key=lambda row: (-row["weight"], row["provider"], row["model"])) + return rows + + +def fetch_model_ids(models_url, key="", extra_headers=None, timeout=MODEL_LIST_TIMEOUT): + headers = dict(extra_headers or {}) + if key: + headers["Authorization"] = "Bearer " + key + try: + request = urllib.request.Request(models_url, headers=headers, method="GET") + with urllib.request.urlopen(request, timeout=timeout) as response: + payload = json.loads(response.read().decode("utf-8", "replace")) + except (OSError, ValueError): + return None + rows = payload.get("data") if isinstance(payload, dict) else None + if not isinstance(rows, list): + return None + ids = [str(row.get("id")).strip() for row in rows if isinstance(row, dict) and str(row.get("id") or "").strip()] + return ids or None + + +def embedded_failure_message(data): + if not isinstance(data, dict): + return "" + error = data.get("error") + if isinstance(error, dict): + message = error.get("message") or error.get("type") or "" + return str(message).strip() + if isinstance(error, str) and error.strip(): + return error.strip() + if data.get("type") == "error": + return str(data.get("message") or "upstream error").strip() + return "" + + +def is_free_model(model): + model = (model or "").lower() + return model.endswith(":free") or model.endswith("-free") + + +def prefer_muse(models): + return sorted(models, key=lambda item: 0 if MUSE_LABEL in item.lower() else 1) + + +def openrouter_free_tier(cand): + model = (cand.get("model") or "").lower() + muse = 0 if MUSE_LABEL in model else 1 + if cand.get("provider") == OPENCODE_LABEL: + return muse + if cand.get("provider") == DEVPLACE_LABEL: + return 2 + muse + if cand.get("provider") == OPENROUTER_LABEL and is_free_model(model): + return 4 + muse + return 6 + + +def models_report(chat): + lines = ["label/model | weight ok/fail avg-ms avg-tps status"] + seen = set() + for cand in chat.ordered_candidates(): + key = (cand["provider"], cand["model"]) + if key in seen: + continue + seen.add(key) + entry = chat.health.entry(cand["provider"], cand["model"]) + ok = entry.get("s") or 0 + fail = entry.get("f") or 0 + avg_ms = str(round((entry.get("lat") or 0.0) / ok)) if ok and entry.get("lat") else "-" + avg_tps = str(round((entry.get("tps") or 0.0) / ok, 1)) if ok and entry.get("tps") else "-" + if cand["chat_url"] == MUSE_CHAT_URL: + status = muse_health() + elif chat.health.circuit_open(cand["provider"], cand["model"]): + status = "circuit open" + elif not ok and not fail: + status = "untested" + else: + status = "ready" + lines.append("%s/%s | %.4f %d/%d %s %s %s" % (cand["provider"], cand["model"], model_weight(entry), ok, fail, avg_ms, avg_tps, status)) + return "\n".join(lines) + + class ChatClient: def __init__(self, config): self.config = config + self.scope = ("default", "main") + self.fixed_key = None + self.health = ModelHealth(getattr(config, "home", "") or "") + + def muse_key(self): + if self.fixed_key is not None: + return self.fixed_key + return "%s/%s" % (self.scope[0], self.scope[1]) + + def reset_muse(self): + home = getattr(self.config, "home", None) + if home is None: + return + key = self.muse_key() + + def drop(data): + data.pop(key, None) + + muse_transact(home, drop) + + def openrouter_models(self): + models = [item for item in getattr(self.config, "openrouter_models", []) if item] + explicit = getattr(self.config, "explicit_model", "") + if explicit and explicit not in models: + models.insert(0, explicit) + free = self.health.cached_list("openrouter:free") + if free is None: + fetched = fetch_model_ids(OPENROUTER_BASE + "/models", self.config.openrouter_key) + free = [item for item in fetched if item.endswith(":free")] if fetched else [] + if free: + self.health.store_list("openrouter:free", free) + free = prefer_muse(free or []) + models = [item for item in free if item not in models] + models + if not models: + models = [OPENROUTER_MODEL] + return models + + def opencode_models(self): + models = list(getattr(self.config, "opencode_models", []) or []) + if not models: + models = self.health.cached_list("opencode:free") + if models is None: + models = opencode_free_models() or [item for item in ZEN_FREE_MODELS if is_free_model(item)] + self.health.store_list("opencode:free", models) + return prefer_muse(models) + + def ollama_models(self): + if getattr(self.config, "ollama_enabled", True) is False: + return [] + models = list(getattr(self.config, "ollama_models", []) or []) + if not models: + models = self.health.cached_list("ollama:local", OLLAMA_LIST_TTL) + if models is None: + models = fetch_model_ids(OLLAMA_BASE + "/models", "", timeout=3) or [] + if models: + self.health.store_list("ollama:local", models) + return models + + def devplace_models(self): + models = list(getattr(self.config, "devplace_models", []) or []) + if models: + return prefer_muse(models) + free = self.health.cached_list("devplace:free") + if free is None: + fetched = fetch_model_ids(self.config.devplace_base + "/models", self.config.devplace_key) + if fetched is None: + return list(DEVPLACE_FREE_ROUTES) + free = [item for item in fetched if item in DEVPLACE_FREE_ROUTES or (MUSE_LABEL in item.lower() and is_free_model(item))] + if free: + self.health.store_list("devplace:free", free) + return prefer_muse(free or []) + + def candidates(self): + found = [] + if getattr(self.config, "opencode_enabled", True) and opencode_binary(): + for model in self.opencode_models(): + found.append({"provider": OPENCODE_LABEL, "chat_url": OPENCODE_CHAT_URL, "models_url": OPENCODE_MODELS_URL, "key": "", "model": model, "profile": ""}) + for model in self.ollama_models(): + found.append({"provider": OLLAMA_LABEL, "chat_url": OLLAMA_BASE + "/chat/completions", "models_url": OLLAMA_BASE + "/models", "key": "", "model": model, "profile": ""}) + if getattr(self.config, "devplace_key", ""): + for model in self.devplace_models(): + found.append({"provider": DEVPLACE_LABEL, "chat_url": self.config.devplace_base + "/chat/completions", "models_url": self.config.devplace_base + "/models", "key": self.config.devplace_key, "model": model, "profile": ""}) + for entry in self.config.backends: + for model in split_models(entry.get("model")) or [self.config.model]: + found.append({"provider": entry["label"], "chat_url": entry["base"] + "/chat/completions", "models_url": entry["base"] + "/models", "key": entry["key"], "model": model, "profile": entry.get("profile") or ""}) + if getattr(self.config, "openrouter_enabled", True): + for model in self.openrouter_models(): + found.append({"provider": OPENROUTER_LABEL, "chat_url": OPENROUTER_BASE + "/chat/completions", "models_url": OPENROUTER_BASE + "/models", "key": self.config.openrouter_key, "model": model, "profile": ""}) + if getattr(self.config, "pollinations_enabled", True): + for model in getattr(self.config, "pollinations_models", []) or list(POLLINATIONS_MODELS): + found.append({"provider": POLLINATIONS_LABEL, "chat_url": POLLINATIONS_BASE + "/chat/completions", "models_url": POLLINATIONS_BASE + "/models", "key": "", "model": model, "profile": ""}) + if getattr(self.config, "zen_enabled", False): + for model in getattr(self.config, "zen_models", []) or list(ZEN_FREE_MODELS): + found.append({"provider": ZEN_LABEL, "chat_url": ZEN_BASE + "/chat/completions", "models_url": ZEN_BASE + "/models", "key": "", "model": model, "profile": "opencode"}) + found.append({"provider": MUSE_LABEL, "chat_url": MUSE_CHAT_URL, "models_url": MUSE_MODELS_URL, "key": "", "model": self.config.model, "profile": ""}) + if grok_binary(): + found.append({"provider": GROK_LABEL, "chat_url": GROK_CHAT_URL, "models_url": GROK_MODELS_URL, "key": "", "model": self.config.model, "profile": ""}) + if claude_binary(): + found.append({"provider": CLAUDE_LABEL, "chat_url": CLAUDE_CHAT_URL, "models_url": CLAUDE_MODELS_URL, "key": "", "model": self.config.model, "profile": ""}) + if codex_binary(): + found.append({"provider": CODEX_LABEL, "chat_url": CODEX_CHAT_URL, "models_url": CODEX_MODELS_URL, "key": "", "model": self.config.model, "profile": ""}) + if gemini_binary(): + found.append({"provider": GEMINI_LABEL, "chat_url": GEMINI_CHAT_URL, "models_url": GEMINI_MODELS_URL, "key": "", "model": self.config.model, "profile": ""}) + return found + + def ordered_candidates(self): + http = [cand for cand in self.candidates() if cand["chat_url"] not in CLI_CHAT_URLS] + cli = [cand for cand in self.candidates() if cand["chat_url"] in CLI_CHAT_URLS] + ranked = sorted(enumerate(http), key=lambda pair: (openrouter_free_tier(pair[1]), -self.health.weight(pair[1]["provider"], pair[1]["model"]), pair[0])) + usable = [cand for _, cand in ranked if not self.health.circuit_open(cand["provider"], cand["model"])] + return (usable or [cand for _, cand in ranked]) + cli def backends(self): - items = [(PRIMARY_LABEL, PRIMARY_BASE + "/v1/chat/completions", PRIMARY_BASE + "/v1/models", "x")] - if self.config.devplace_key: - items.append((FALLBACK_LABEL, FALLBACK_BASE + "/chat/completions", FALLBACK_BASE + "/models", self.config.devplace_key)) - return items + rows = [] + seen = set() + for cand in self.candidates(): + key = (cand["provider"], cand["chat_url"]) + if key in seen: + continue + seen.add(key) + rows.append((cand["provider"], cand["chat_url"], cand["models_url"], cand["key"], cand["model"])) + return rows - def headers(self, key): - return {"Content-Type": "application/json", "Authorization": "Bearer " + key} + def headers(self, key, extra=None): + headers = {"Content-Type": "application/json"} + if key: + headers["Authorization"] = "Bearer " + key + if extra: + headers.update(extra) + return headers - def complete(self, messages, tools=None, stream_sink=None): + def complete(self, messages, tools=None, stream_sink=None, ephemeral=False, deadline=None): payload = {"model": self.config.model, "messages": messages} if tools: payload["tools"] = tools payload["tool_choice"] = "auto" - errors = [] - for label, chat_url, _models_url, key in self.backends(): - try: - if stream_sink is None: - return self.single_shot(chat_url, key, payload, label) - return self.streaming(chat_url, key, payload, label, stream_sink) - except BackendError as exc: - errors.append("%s: %s" % (label, exc)) - except (OSError, ValueError) as exc: - errors.append("%s: %s" % (label, short_error(exc))) - raise BackendError("all backends failed (%s)" % "; ".join(errors)) + limit = getattr(self.config, "retry_seconds", 0) + give_up = time.time() + limit if limit > 0 else None + if deadline is not None: + give_up = deadline if give_up is None else min(give_up, deadline) + round_no = 0 + while True: + errors = [] + cands = self.ordered_candidates() if round_no == 0 else self.all_candidates() + for cand in cands: + reply = self.try_candidate(cand, messages, tools, payload, stream_sink, ephemeral, errors) + if reply is not None: + return reply + pause = RETRY_BACKOFF[min(round_no, len(RETRY_BACKOFF) - 1)] + round_no += 1 + if give_up is not None and time.time() + pause > give_up: + raise BackendError("all backends failed after %d rounds (%s)" % (round_no, "; ".join(errors[-6:]))) + print(paint("all %d backends failed (round %d), retrying in %ds: %s" % (len(cands), round_no, pause, (errors[-1] if errors else "")[:160]), Ansi.DIM), file=sys.stderr) + time.sleep(pause) - def single_shot(self, url, key, payload, label): + def all_candidates(self): + http = [cand for cand in self.candidates() if cand["chat_url"] not in CLI_CHAT_URLS] + cli = [cand for cand in self.candidates() if cand["chat_url"] in CLI_CHAT_URLS] + ranked = sorted(enumerate(http), key=lambda pair: (openrouter_free_tier(pair[1]), -self.health.weight(pair[1]["provider"], pair[1]["model"]), pair[0])) + return [cand for _, cand in ranked] + cli + + def call_candidate(self, cand, messages, tools, attempt, stream_sink, ephemeral): + label, chat_url, key, model = cand["provider"], cand["chat_url"], cand["key"], cand["model"] + extra = opencode_headers() if cand["profile"] == "opencode" else None + timeout = STREAM_TIMEOUT if stream_sink is not None else HTTP_TIMEOUT + if chat_url == MUSE_CHAT_URL: + return self.muse_turn(messages, tools, stream_sink, label, timeout=timeout, ephemeral=ephemeral) + if chat_url == GROK_CHAT_URL: + return self.grok_turn(messages, tools, stream_sink, label, timeout=timeout) + if chat_url == CLAUDE_CHAT_URL: + return self.claude_turn(messages, tools, stream_sink, label, timeout=timeout) + if chat_url == CODEX_CHAT_URL: + return self.codex_turn(messages, tools, stream_sink, label, timeout=timeout) + if chat_url == GEMINI_CHAT_URL: + return self.gemini_turn(messages, tools, stream_sink, label, timeout=timeout) + if chat_url == OPENCODE_CHAT_URL: + return self.opencode_turn(messages, tools, stream_sink, label, model, timeout=timeout) + if stream_sink is None: + return self.single_shot(chat_url, key, attempt, label, extra) + return self.streaming(chat_url, key, attempt, label, stream_sink, extra) + + def try_candidate(self, cand, messages, tools, payload, stream_sink, ephemeral, errors): + label, model = cand["provider"], cand["model"] + attempt = dict(payload, model=model) + tries = [(tools, attempt)] + if attempt.get("tools"): + tries.append((None, {name: value for name, value in attempt.items() if name not in ("tools", "tool_choice")})) + for pos, (use_tools, body) in enumerate(tries): + started = time.time() + try: + reply = self.call_candidate(cand, messages, use_tools, body, stream_sink, ephemeral) + if not isinstance(reply, dict): + raise BackendError("backend returned no reply") + self.health.record(label, model, True, (time.time() - started) * 1000, max(1, len(reply.get("content") or "") // 4) / max(time.time() - started, 0.001)) + return reply + except BackendError as exc: + if pos == 0 and exc.status == 400 and len(tries) > 1: + continue + self.health.record(label, model, False) + errors.append("%s/%s: %s" % (label, model, exc)) + return None + except Exception as exc: + self.health.record(label, model, False) + errors.append("%s/%s: %s" % (label, model, short_error(exc))) + return None + return None + + def single_shot(self, url, key, payload, label, extra=None): body = json.dumps(payload).encode("utf-8") - request = urllib.request.Request(url, data=body, headers=self.headers(key), method="POST") + request = urllib.request.Request(url, data=body, headers=self.headers(key, extra), method="POST") try: with urllib.request.urlopen(request, timeout=HTTP_TIMEOUT) as response: data = json.loads(response.read().decode("utf-8", "replace")) @@ -2468,18 +3717,25 @@ class ChatClient: raise BackendError("HTTP %d: %s" % (exc.code, exc.read().decode("utf-8", "replace")[:300]), exc.code) except urllib.error.URLError as exc: raise BackendError(short_error(exc)) - return self.normalize(data["choices"][0]["message"], label) + failure = embedded_failure_message(data) + if failure: + raise BackendError("upstream error: %s" % failure[:300]) + try: + message = data["choices"][0]["message"] + except (KeyError, IndexError, TypeError): + raise BackendError("upstream returned no choices: %s" % json.dumps(data)[:300]) + return self.normalize(message, label) def normalize(self, message, label): calls = [] for call in message.get("tool_calls") or []: func = call.get("function") or {} calls.append({"id": call.get("id") or uuid.uuid4().hex, "name": func.get("name") or "", "arguments": func.get("arguments") or "{}"}) - return {"role": "assistant", "content": message.get("content") or "", "reasoning": message.get("reasoning") or "", "tool_calls": calls, "backend": label} + return {"role": "assistant", "content": message.get("content") or "", "reasoning": message.get("reasoning") or message.get("reasoning_content") or "", "tool_calls": calls, "backend": label} - def streaming(self, url, key, payload, label, sink): + def streaming(self, url, key, payload, label, sink, extra=None): body = json.dumps(dict(payload, stream=True)).encode("utf-8") - request = urllib.request.Request(url, data=body, headers=self.headers(key), method="POST") + request = urllib.request.Request(url, data=body, headers=self.headers(key, extra), method="POST") try: response = urllib.request.urlopen(request, timeout=STREAM_TIMEOUT) except urllib.error.HTTPError as exc: @@ -2509,7 +3765,7 @@ class ChatClient: if piece: content_parts.append(piece) sink(piece) - think = delta.get("reasoning") + think = delta.get("reasoning") or delta.get("reasoning_content") if think: reasoning_parts.append(think) for call in delta.get("tool_calls") or []: @@ -2532,13 +3788,132 @@ class ChatClient: slot["id"] = uuid.uuid4().hex return {"role": "assistant", "content": "".join(content_parts), "reasoning": "".join(reasoning_parts), "tool_calls": ordered, "backend": label} + def muse_turn(self, messages, tools, stream_sink, label, timeout=HTTP_TIMEOUT, ephemeral=False): + binary = muse_binary() + if not binary: + raise BackendError("muse CLI not found (MUSE_BIN or PATH)") + home = getattr(self.config, "home", None) + if ephemeral or home is None: + with tempfile.TemporaryDirectory(prefix="tai-muse-") as folder: + schema_path = os.path.join(folder, "schema.json") + prompt_path = os.path.join(folder, "prompt.txt") + muse_write_file(schema_path, json.dumps(muse_output_schema())) + muse_write_file(prompt_path, muse_prompt(messages, tools)) + answer = muse_run(muse_argv(self.config, schema_path, prompt_path), folder, stream_sink, timeout) + return self.muse_reply(answer, label) + key = self.muse_key() + release = muse_acquire(home, "session:" + key) + + def on_reap(ok): + try: + if not ok: + self.reset_muse() + finally: + release() + + gated = muse_gate(on_reap) + try: + return self.muse_persistent(messages, tools, stream_sink, label, timeout, home, key, gated) + except BaseException: + if not gated.armed.is_set(): + release() + raise + + def muse_persistent(self, messages, tools, stream_sink, label, timeout, home, key, gated): + def plan(data): + entry = data.get(key) + pending = None + session_id = None + if isinstance(entry, dict) and isinstance(entry.get("id"), str) and isinstance(entry.get("hashes"), list): + pending = muse_session_split(messages, entry["hashes"]) + if pending is not None and not muse_session_alive(entry["id"]): + pending = None + if pending is not None: + session_id = entry["id"] + if pending is None or not pending: + session_id = str(uuid.uuid4()) + pending = list(messages) + return session_id, pending + + session_id, pending = muse_transact(home, plan) + try: + with tempfile.TemporaryDirectory(prefix="tai-muse-") as folder: + schema_path = os.path.join(folder, "schema.json") + prompt_path = os.path.join(folder, "prompt.txt") + muse_write_file(schema_path, json.dumps(muse_output_schema())) + muse_write_file(prompt_path, muse_prompt(pending, tools)) + answer = muse_run(muse_argv(self.config, schema_path, prompt_path, session_id), folder, stream_sink, timeout, gated) + except BackendError: + self.reset_muse() + raise + + def commit(data): + data[key] = {"id": session_id, "hashes": [muse_message_hash(item) for item in messages], "updated": now_iso()} + + muse_transact(home, commit) + return self.muse_reply(answer, label) + + def opencode_turn(self, messages, tools, stream_sink, label, model, timeout=HTTP_TIMEOUT): + answer = opencode_parse(opencode_run(model, opencode_prompt(messages, tools), timeout)) + if stream_sink is not None and answer.get("content"): + stream_sink(answer["content"]) + return self.muse_reply(answer, label) + + def grok_turn(self, messages, tools, stream_sink, label, timeout=HTTP_TIMEOUT): + answer = grok_parse_output(grok_run(opencode_prompt(messages, tools), timeout)) + if stream_sink is not None and answer.get("content"): + stream_sink(answer["content"]) + return self.muse_reply(answer, label) + + def claude_turn(self, messages, tools, stream_sink, label, timeout=HTTP_TIMEOUT): + answer = claude_parse_output(claude_run(opencode_prompt(messages, tools), timeout)) + if stream_sink is not None and answer.get("content"): + stream_sink(answer["content"]) + return self.muse_reply(answer, label) + + def codex_turn(self, messages, tools, stream_sink, label, timeout=HTTP_TIMEOUT): + stdout, outfile_text, used_schema = codex_run(opencode_prompt(messages, tools), timeout) + answer = codex_parse_output(stdout, outfile_text, used_schema) + if stream_sink is not None and answer.get("content"): + stream_sink(answer["content"]) + return self.muse_reply(answer, label) + + def gemini_turn(self, messages, tools, stream_sink, label, timeout=HTTP_TIMEOUT): + answer = gemini_parse_output(gemini_run(opencode_prompt(messages, tools), timeout)) + if stream_sink is not None and answer.get("content"): + stream_sink(answer["content"]) + return self.muse_reply(answer, label) + + def muse_reply(self, answer, label): + calls = [] + for call in answer.get("tool_calls") or []: + calls.append({"id": uuid.uuid4().hex, "name": call.get("name") or "", "arguments": call.get("arguments") or "{}"}) + return {"role": "assistant", "content": answer.get("content") or "", "reasoning": "", "tool_calls": calls, "backend": label} + def probe_backend(models_url, key): + if models_url == MUSE_MODELS_URL: + return muse_health() + if models_url == GROK_MODELS_URL: + return grok_health() + if models_url == CLAUDE_MODELS_URL: + return claude_health() + if models_url == CODEX_MODELS_URL: + return codex_health() + if models_url == GEMINI_MODELS_URL: + return gemini_health() + if models_url == OPENCODE_MODELS_URL: + return "ok, %s" % opencode_binary() if opencode_binary() else "opencode CLI not found" request = urllib.request.Request(models_url, headers={"Authorization": "Bearer " + key}) started = time.time() try: with urllib.request.urlopen(request, timeout=6) as response: - count = len(json.loads(response.read().decode("utf-8", "replace")).get("data", [])) + try: + payload = json.loads(response.read().decode("utf-8", "replace")) + except ValueError: + return "bad response" + data = payload.get("data", []) if isinstance(payload, dict) else [] + count = len(data) if isinstance(data, list) else 0 return "ok, %d models, %dms" % (count, (time.time() - started) * 1000) except urllib.error.HTTPError as exc: if exc.code == 401: @@ -2546,6 +3921,8 @@ def probe_backend(models_url, key): return "http %d" % exc.code except OSError as exc: return "fail: %s" % short_error(exc)[:60] + except Exception as exc: + return "fail: %s" % short_error(exc)[:60] def edge_token(skew=0): @@ -2749,7 +4126,7 @@ def listen_audio(seconds, workdir): command = [part.replace("{seconds}", str(seconds)).replace("{path}", path) for part in recorder[1]] try: subprocess.run(command, timeout=seconds + 15, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, check=False) - done = subprocess.run([transcriber, path], capture_output=True, text=True, timeout=180) + done = subprocess.run([transcriber, path], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=180) except (OSError, subprocess.SubprocessError) as exc: return "listening failed: " + short_error(exc) finally: @@ -3543,7 +4920,7 @@ class Tools: spinner = Spinner("merging memory") spinner.start() try: - merged = self.app.chat.complete(prompt) + merged = self.app.chat.complete(prompt, ephemeral=True) except BackendError as exc: return "error: memory merge failed: " + short_error(exc) finally: @@ -3618,13 +4995,13 @@ class Tools: return "tmux not available" header = [] try: - info = subprocess.run(["tmux", "display-message", "-p", "#{session_name}:#{window_index}.#{pane_index} #{pane_current_command}"], capture_output=True, text=True, timeout=10) + info = subprocess.run(["tmux", "display-message", "-p", "#{session_name}:#{window_index}.#{pane_index} #{pane_current_command}"], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=10) if info.returncode == 0 and info.stdout.strip(): header.append(info.stdout.strip()) except (OSError, subprocess.SubprocessError): pass try: - done = subprocess.run(["tmux", "capture-pane", "-p", "-S", "-%d" % lines], capture_output=True, text=True, timeout=10) + done = subprocess.run(["tmux", "capture-pane", "-p", "-S", "-%d" % lines], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=10) except (OSError, subprocess.SubprocessError) as exc: return "terminal capture failed: " + short_error(exc) if done.returncode != 0: @@ -4269,7 +5646,7 @@ def compact_messages(messages, chat, keep=KEEP_TURNS): {"role": "system", "content": COMPACT_SYSTEM}, {"role": "user", "content": "Compress this history:\n" + "\n".join(digest)}, ] - summary = chat.complete(request)["content"].strip() + summary = chat.complete(request, ephemeral=True)["content"].strip() rebuilt = [ system, {"role": "user", "content": "Previous session summary:\n" + summary}, @@ -4291,6 +5668,8 @@ class Agent: self.yolo = yolo self.store = store self.chat = ChatClient(config) + if not persist or depth > 0: + self.chat.fixed_key = "worker/%s" % uuid.uuid4().hex self.tools = Tools(self) self.profile = "" self.bot = "main" @@ -4317,6 +5696,9 @@ class Agent: content += "\n\n" + AUTONOMOUS_NOTE self.messages[0] = {"role": "system", "content": content} + def sync_chat_scope(self): + self.chat.scope = (self.profile, self.bot) + def reset_history(self): self.messages = [{"role": "system", "content": self.system_message}] self.apply_system() @@ -4340,6 +5722,7 @@ class Agent: restored = self.store.load_session(name) if self.persist else [] self.messages = [{"role": "system", "content": self.system_message}] + restored self.apply_system() + self.sync_chat_scope() if not silent: print(paint("profile: %s (%d chars knowledge, %d restored messages)" % (name, len(self.system_message), len(restored)), Ansi.GREEN)) @@ -4370,6 +5753,7 @@ class Agent: restored = self.store.load_session(self.profile, target) if self.persist else [] self.messages = [{"role": "system", "content": self.system_message}] + restored self.apply_system() + self.sync_chat_scope() if not silent: print(paint("bot: %s (%d restored messages)" % (target, len(restored)), Ansi.GREEN)) return "switched to bot '%s'" % target @@ -4383,6 +5767,7 @@ class Agent: args = (recorded.get("function") or {}).get("arguments") if isinstance(args, str) and args: recorded["function"]["arguments"] = self.store.redact(args, self.profile) + self.chat.reset_muse() def ask_approval(self, command, guidance=True): if self.yolo or self.auto: @@ -4506,7 +5891,7 @@ class Agent: spinner = Spinner("thinking") spinner.start() try: - reply = self.chat.complete(self.messages, payload, stream_sink=sink) + reply = self.chat.complete(self.messages, payload, stream_sink=sink, deadline=self.deadline) except BackendError as exc: if loud: print(paint("backend error: %s" % exc, Ansi.RED)) @@ -4521,7 +5906,7 @@ class Agent: print(render_markdown(reply["content"])) if reply["backend"] not in backends: backends.append(reply["backend"]) - if loud and (len(backends) > 1 or reply["backend"] != PRIMARY_LABEL): + if loud and (len(backends) > 1 or reply["backend"] != MUSE_LABEL): print(paint("(via %s)" % reply["backend"], Ansi.DIM)) if loud and reply["reasoning"] and not reply["content"] and not reply["tool_calls"]: print(paint(reply["reasoning"][:2000], Ansi.GRAY)) @@ -4625,6 +6010,7 @@ class Agent: self.system_message = system self.bot = target self.apply_system() + self.sync_chat_scope() try: answer = self.run_turn(rest, capture=capture, max_steps=max_steps, _routed=True) finally: @@ -4632,6 +6018,7 @@ class Agent: self.store.save_session(self.profile, self.messages, target) self.messages, self.system_message, self.bot = saved_messages, saved_system, saved_bot self.apply_system() + self.sync_chat_scope() self.store.log_event(self.profile, "user", "message", self.store.redact(original, self.profile)) self.messages.append({"role": "user", "content": original}) self.messages.append({"role": "assistant", "content": answer, "reasoning": "", "tool_calls": [], "backend": "mention"}) @@ -4727,7 +6114,7 @@ def parse_schedule_at(raw): if text[-1:] in ("Z", "z"): text = text[:-1] + "+00:00" try: - moment = datetime.fromisoformat(text) + moment = _fromisoformat(text) except ValueError: raise ValueError("invalid datetime, use ISO like 2026-10-08T09:00") return moment.astimezone(timezone.utc).isoformat() @@ -4755,7 +6142,7 @@ def format_delay(seconds): def local_display(utc_iso): try: - return datetime.fromisoformat(utc_iso).astimezone().strftime("%Y-%m-%d %H:%M") + return _fromisoformat(utc_iso).astimezone().strftime("%Y-%m-%d %H:%M") except ValueError: return utc_iso @@ -4766,7 +6153,7 @@ def format_schedule_lines(items, now): lines = [] for item in items: try: - due_in = (datetime.fromisoformat(item["next_run"]) - now).total_seconds() + due_in = (_fromisoformat(item["next_run"]) - now).total_seconds() except ValueError: due_in = -1 when = "every %s" % format_delay(item["every"]) if item["every"] else "at %s" % local_display(item["next_run"]) @@ -4809,7 +6196,7 @@ def scheduler_tick(config, seal, runner=None): stamp = datetime.now(timezone.utc).isoformat() if every_sec: try: - claimed_next = datetime.fromisoformat(next_run) + claimed_next = _fromisoformat(next_run) except ValueError: continue while claimed_next <= datetime.now(timezone.utc): @@ -4851,7 +6238,7 @@ def start_scheduler(config, seal): return stop_flag -COMMANDS = ("help", "profile", "profiles", "bots", "bot", "env", "skills", "sysinfo", "secret", "schedule", "schedules", "unschedule", "records", "record", "graph", "tags", "tools", "install", "search", "audit", "restore", "release", "fork", "agents", "agent", "compact", "clear", "quit", "exit") +COMMANDS = ("help", "profile", "profiles", "bots", "bot", "env", "skills", "sysinfo", "models", "secret", "schedule", "schedules", "unschedule", "records", "record", "graph", "tags", "tools", "install", "search", "audit", "restore", "release", "fork", "agents", "agent", "compact", "clear", "quit", "exit") def complete_command(text, state): @@ -4919,6 +6306,8 @@ def handle_command(agent, text): print("- %s: %s" % (skill_name, SKILL_BLUEPRINTS[skill_name]["hint"])) elif name == "sysinfo": print(collect_sysinfo()) + elif name == "models": + print(models_report(agent.chat)) elif name == "secret": sub = arg.split(None, 1) action = sub[0].lower() if sub else "" @@ -5185,7 +6574,7 @@ def run_telegram_bot(agent): def run_host(argv, timeout=120): try: - done = subprocess.run(argv, capture_output=True, text=True, timeout=timeout) + done = subprocess.run(argv, stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=timeout) return done.returncode == 0, (done.stdout or "") + (done.stderr or "") except (OSError, subprocess.SubprocessError) as exc: return False, short_error(exc) @@ -5217,7 +6606,7 @@ def ensure_venv(home): if not venv_available(): return False, "python venv module unavailable" try: - done = subprocess.run([sys.executable, "-m", "venv", venv_dir(home)], capture_output=True, text=True, timeout=300) + done = subprocess.run([sys.executable, "-m", "venv", venv_dir(home)], stdout=subprocess.PIPE, stderr=subprocess.PIPE, universal_newlines=True, timeout=300) except (OSError, subprocess.SubprocessError) as exc: return False, short_error(exc) if done.returncode != 0 or not venv_python(home): @@ -5695,10 +7084,8 @@ def banner(agent): if agent.store.seal.enabled: seal_state = "on (default key)" if agent.store.seal.default_key else "on (TAI_PASSPHRASE)" print(" seal: %s" % seal_state) - for label, _chat_url, models_url, key in agent.chat.backends(): + for label, _chat_url, models_url, key, _model in agent.chat.backends(): print(" %s: %s" % (label, probe_backend(models_url, key))) - if not agent.config.devplace_key: - print(" fallback: no DEVPLACE_API_KEY") if agent.store.load_turn_state(agent.profile, agent.bot).get("open"): print(paint("interrupted turn available, restart with --continue to pick it up", Ansi.YELLOW)) print(paint("type /help for commands", Ansi.DIM)) @@ -5777,6 +7164,9 @@ def main(argv=None): return 1 try: agent = boot(args) + except KeyboardInterrupt: + print(paint("\ninterrupted", Ansi.YELLOW)) + return 130 except SealError as exc: print("seal error: %s" % exc, file=sys.stderr) return 2 @@ -5805,9 +7195,18 @@ def main(argv=None): if args.prompt: try: agent.run_turn(" ".join(args.prompt)) + except KeyboardInterrupt: + print(paint("\ninterrupted", Ansi.YELLOW)) + return 130 finally: - agent.checkpoint() - agent.store.close() + try: + agent.checkpoint() + except Exception: + pass + try: + agent.store.close() + except Exception: + pass return 0 if args.resume and not args.telegram and not args.scheduler: resume_turn(agent)