diff --git a/scheduler/scheduler.py b/scheduler/scheduler.py index 61638d3..b8cba8a 100644 --- a/scheduler/scheduler.py +++ b/scheduler/scheduler.py @@ -24,24 +24,40 @@ def run_skill(skill_name: str): "skill": skill_name, "timestamp": datetime.now(timezone.utc).isoformat(), } - with open(audit_file, "a") as f: - f.write(json.dumps(entry) + "\n") + try: + audit_file.parent.mkdir(parents=True, exist_ok=True) + with open(audit_file, "a", encoding="utf-8") as f: + f.write(json.dumps(entry) + "\n") + except OSError as e: + print(f" [audit] failed to record run of {skill_name!r}: {e}") print(f"[{datetime.now().isoformat()}] Ran skill: {skill_name}") def load_jobs(scheduler: BackgroundScheduler): - """Load job definitions from jobs/ directory.""" + """Load job definitions from jobs/ directory. + + A single malformed job file is logged and skipped rather than being allowed + to abort loading of every other job. + """ for job_file in JOBS_DIR.glob("*.json"): - data = json.loads(job_file.read_text()) + try: + data = json.loads(job_file.read_text()) + except (json.JSONDecodeError, OSError) as e: + print(f" Skipping {job_file.name}: could not read job ({e})") + continue if not data.get("enabled", True): continue - scheduler.add_job( - run_skill, - CronTrigger.from_crontab(data["cron"]), - args=[data["skill"]], - id=data.get("id", data["name"]), - name=data["name"], - replace_existing=True, - ) + try: + scheduler.add_job( + run_skill, + CronTrigger.from_crontab(data["cron"]), + args=[data["skill"]], + id=data.get("id", data["name"]), + name=data["name"], + replace_existing=True, + ) + except (KeyError, ValueError) as e: + print(f" Skipping {job_file.name}: invalid job definition ({e})") + continue print(f" Scheduled: {data['name']} ({data['cron']})") def main(): diff --git a/server.py b/server.py index 41bfe56..4a139c2 100644 --- a/server.py +++ b/server.py @@ -115,6 +115,24 @@ def write_file(path: Path, content: str): path.write_text(content, encoding="utf-8") return True +_MISSING = object() + +def load_json_file(path: Path, default=_MISSING): + """Read and parse a JSON file. + + Raises a descriptive HTTPException instead of leaking an opaque 500 when the + file is missing or corrupt, so callers propagate a clear error to the client. + If ``default`` is provided it is returned when the file does not exist. + """ + if not path.exists(): + if default is not _MISSING: + return default + raise HTTPException(404, f"{path.name} not found") + try: + return json.loads(path.read_text(encoding="utf-8")) + except (json.JSONDecodeError, OSError, UnicodeDecodeError) as e: + raise HTTPException(500, f"Failed to read {path.name}: {e}") + def list_dir(path: Path): if not path.exists(): return [] @@ -127,29 +145,37 @@ def append_audit(entry: dict): audit_file = BASE_DIR / "audit" / "audit.log" entry["timestamp"] = get_timestamp() entry["id"] = str(uuid.uuid4())[:8] - with open(audit_file, "a") as f: - f.write(json.dumps(entry) + "\n") + try: + audit_file.parent.mkdir(parents=True, exist_ok=True) + with open(audit_file, "a", encoding="utf-8") as f: + f.write(json.dumps(entry) + "\n") + except OSError as e: + # Auditing is best-effort: never let a logging failure abort the + # underlying operation, but surface it on the server console. + print(f"[audit] failed to write entry {entry.get('action')!r}: {e}") # ─── Agent Discovery (instant filesystem checks) ──────────────────── def check_agent(name: str) -> dict: """Instant filesystem-based check. No subprocess needed.""" - try: - if name == "opencode": - exists = shutil.which("opencode") is not None - status = "online" if exists else "offline" - elif name == "hermes": - exists = shutil.which("hermes") is not None - status = "online" if exists else "offline" - elif name == "gemini": - # Gemini has valid OAuth tokens logged in - oauth = Path.home() / ".gemini" / "oauth_creds.json" - exists = shutil.which("gemini") is not None - logged_in = oauth.exists() and "ya29" in oauth.read_text() - status = "online" if exists and logged_in else "offline" if not exists else "warning" - else: - status = "offline" - except Exception: + if name == "opencode": + status = "online" if shutil.which("opencode") is not None else "offline" + elif name == "hermes": + status = "online" if shutil.which("hermes") is not None else "offline" + elif name == "gemini": + exists = shutil.which("gemini") is not None + # Gemini needs a valid OAuth token on disk to be usable. + logged_in = False + oauth = Path.home() / ".gemini" / "oauth_creds.json" + if oauth.exists(): + try: + logged_in = "ya29" in oauth.read_text() + except OSError as e: + # Distinguish an unreadable credential file from "not logged in" + # instead of silently reporting the agent as offline. + print(f"[agent-health] could not read gemini credentials: {e}") + status = "online" if exists and logged_in else "offline" if not exists else "warning" + else: status = "offline" return {"name": name, "status": status} @@ -201,14 +227,8 @@ def list_skills(): if d.is_dir() and not d.name.startswith("_"): skill_md = read_file(d / "SKILL.md") learnings = read_file(d / "learnings.md") - eval_data = {} - eval_path = d / "eval.json" - if eval_path.exists(): - eval_data = json.loads(eval_path.read_text()) - score_history = [] - score_path = d / "score-history.json" - if score_path.exists(): - score_history = json.loads(score_path.read_text()) + eval_data = load_json_file(d / "eval.json", default={}) + score_history = load_json_file(d / "score-history.json", default=[]) skills.append({ "name": d.name, "description": skill_md[:200] if skill_md else "", @@ -227,8 +247,8 @@ def get_skill(name: str): "name": name, "skill": read_file(path / "SKILL.md"), "learnings": read_file(path / "learnings.md"), - "eval": json.loads((path / "eval.json").read_text()) if (path / "eval.json").exists() else {}, - "score_history": json.loads((path / "score-history.json").read_text()) if (path / "score-history.json").exists() else [], + "eval": load_json_file(path / "eval.json", default={}), + "score_history": load_json_file(path / "score-history.json", default=[]), "context": [f.name for f in (path / "context").iterdir()] if (path / "context").exists() else [], } @@ -318,9 +338,7 @@ def run_skill(name: str, req: Optional[SkillRunRequest] = None): @app.get("/api/skills/{name}/eval") def get_skill_eval(name: str): path = BASE_DIR / "skills" / name / "score-history.json" - if not path.exists(): - return {"scores": []} - return {"scores": json.loads(path.read_text())} + return {"scores": load_json_file(path, default=[])} # ─── Routes: Scheduler ──────────────────────────────────────────── @@ -329,7 +347,7 @@ def list_jobs(): jobs_dir = BASE_DIR / "scheduler" / "jobs" jobs = [] for f in sorted(jobs_dir.glob("*.json")): - jobs.append(json.loads(f.read_text())) + jobs.append(load_json_file(f)) return jobs @app.post("/api/scheduler/jobs") @@ -356,7 +374,7 @@ def create_job(job: ScheduleJobRequest): def delete_job(job_id: str): jobs_dir = BASE_DIR / "scheduler" / "jobs" for f in jobs_dir.glob("*.json"): - data = json.loads(f.read_text()) + data = load_json_file(f) if data.get("id") == job_id: f.unlink() append_audit({"action": "job_deleted", "job_id": job_id}) @@ -371,7 +389,15 @@ def get_audit(limit: int = Query(100, le=500)): if not audit_file.exists(): return {"entries": []} lines = audit_file.read_text().strip().split("\n") - entries = [json.loads(l) for l in lines if l.strip()] + entries = [] + for l in lines: + if not l.strip(): + continue + try: + entries.append(json.loads(l)) + except json.JSONDecodeError: + # Skip a corrupt line rather than failing the whole audit view. + continue return {"entries": entries[-limit:]} # ─── Routes: Cost Analytics ─────────────────────────────────────── @@ -379,15 +405,18 @@ def get_audit(limit: int = Query(100, le=500)): @app.get("/api/cost") def get_cost(): cost_file = BASE_DIR / "data" / "cost-history.json" - if not cost_file.exists(): - return {"entries": [], "daily_totals": {}, "monthly_projection": 0, "free_tier_alerts": []} - return json.loads(cost_file.read_text()) + return load_json_file( + cost_file, + default={"entries": [], "daily_totals": {}, "monthly_projection": 0, "free_tier_alerts": []}, + ) @app.post("/api/cost/record") def record_cost(data: dict): cost_file = BASE_DIR / "data" / "cost-history.json" - cost_data = json.loads(cost_file.read_text()) if cost_file.exists() else \ - {"entries": [], "daily_totals": {}, "monthly_projection": 0, "free_tier_alerts": []} + cost_data = load_json_file( + cost_file, + default={"entries": [], "daily_totals": {}, "monthly_projection": 0, "free_tier_alerts": []}, + ) cost_data["entries"].append({ "timestamp": get_timestamp(), "agent": data.get("agent", "unknown"), @@ -403,9 +432,7 @@ def record_cost(data: dict): @app.get("/api/plugins") def list_plugins(): reg_file = BASE_DIR / "registry" / "plugins.json" - if not reg_file.exists(): - return {"plugins": []} - return json.loads(reg_file.read_text()) + return load_json_file(reg_file, default={"plugins": []}) @app.post("/api/plugins/install") def install_plugin(data: dict): @@ -413,7 +440,7 @@ def install_plugin(data: dict): if not name: raise HTTPException(400, "Plugin name required") reg_file = BASE_DIR / "registry" / "plugins.json" - reg = json.loads(reg_file.read_text()) if reg_file.exists() else {"plugins": []} + reg = load_json_file(reg_file, default={"plugins": []}) if any(p["name"] == name for p in reg["plugins"]): return {"status": "already_installed"} reg["plugins"].append({ @@ -478,15 +505,13 @@ def list_prompts(): @app.get("/api/settings") def get_settings(): sf = BASE_DIR / "data" / "settings.json" - if not sf.exists(): - return {} - return json.loads(sf.read_text()) + return load_json_file(sf, default={}) @app.put("/api/settings") def update_settings(data: SettingsUpdate): sf = BASE_DIR / "data" / "settings.json" # Merge with existing - existing = json.loads(sf.read_text()) if sf.exists() else {} + existing = load_json_file(sf, default={}) existing.update(data.settings) sf.write_text(json.dumps(existing, indent=2)) append_audit({"action": "settings_updated"}) @@ -520,9 +545,7 @@ def discover_standards(): CHAT_HISTORY_FILE = BASE_DIR / "data" / "chat-history.json" def load_chat_history(): - if CHAT_HISTORY_FILE.exists(): - return json.loads(CHAT_HISTORY_FILE.read_text()) - return {"messages": []} + return load_json_file(CHAT_HISTORY_FILE, default={"messages": []}) def save_chat_message(msg: dict): history = load_chat_history() @@ -771,7 +794,7 @@ def load_kanban_tasks(): ensure_dir(KANBAN_DIR) tasks = [] for f in sorted(KANBAN_DIR.glob("*.json")): - tasks.append(json.loads(f.read_text())) + tasks.append(load_json_file(f)) return tasks KANBAN_ID_RE = re.compile(r"^[0-9a-f]{6,16}$") @@ -809,38 +832,51 @@ def dispatch_kanban_task(task_id: str): threading.Thread(target=_run_kanban_agent, args=(task_id,), daemon=True).start() def _run_kanban_agent(task_id: str): + # Runs in a daemon thread: any unhandled exception would be lost and leave + # the task stuck in "in_progress" forever, so catch failures and surface + # them by marking the task blocked with the error. path = kanban_task_path(task_id) if not path.exists(): return - task = json.loads(path.read_text()) - agent = task.get("assignee") - prompt = task["title"] if not task.get("body") else f"{task['title']}\n\n{task['body']}" + try: + task = json.loads(path.read_text()) + agent = task.get("assignee") + prompt = task["title"] if not task.get("body") else f"{task['title']}\n\n{task['body']}" - response = execute_agent(agent, prompt) - failed = response.startswith(("⏱", "⚠", "Unknown agent")) + response = execute_agent(agent, prompt) + failed = response.startswith(("⏱", "⚠", "Unknown agent")) - task = json.loads(path.read_text()) # reload in case it changed while the agent ran - task.setdefault("comments", []).append({ - "id": str(uuid.uuid4())[:8], - "message": f"🤖 **{agent}**\n\n{response}", - "timestamp": get_timestamp(), - }) - if failed: - task["status"] = "blocked" - task["block_reason"] = response[:300] - append_audit({"action": "kanban_task_dispatch_failed", "task_id": task_id, "agent": agent}) - else: - task["status"] = "done" - task["summary"] = response[:300] - task["completed_at"] = get_timestamp() - append_audit({"action": "kanban_task_dispatch_completed", "task_id": task_id, "agent": agent}) - task["updated"] = get_timestamp() - save_kanban_task(task) + task = json.loads(path.read_text()) # reload in case it changed while the agent ran + task.setdefault("comments", []).append({ + "id": str(uuid.uuid4())[:8], + "message": f"🤖 **{agent}**\n\n{response}", + "timestamp": get_timestamp(), + }) + if failed: + task["status"] = "blocked" + task["block_reason"] = response[:300] + append_audit({"action": "kanban_task_dispatch_failed", "task_id": task_id, "agent": agent}) + else: + task["status"] = "done" + task["summary"] = response[:300] + task["completed_at"] = get_timestamp() + append_audit({"action": "kanban_task_dispatch_completed", "task_id": task_id, "agent": agent}) + task["updated"] = get_timestamp() + save_kanban_task(task) + except Exception as e: + print(f"[kanban] dispatch for task {task_id} crashed: {e}") + try: + task = json.loads(path.read_text()) + task["status"] = "blocked" + task["block_reason"] = f"Dispatch crashed: {e}"[:300] + task["updated"] = get_timestamp() + save_kanban_task(task) + append_audit({"action": "kanban_task_dispatch_error", "task_id": task_id, "error": str(e)[:200]}) + except Exception as inner: + print(f"[kanban] could not mark task {task_id} as blocked: {inner}") def load_goals(): - if GOALS_FILE.exists(): - return json.loads(GOALS_FILE.read_text()) - return [] + return load_json_file(GOALS_FILE, default=[]) def save_goals(goals: list): GOALS_FILE.write_text(json.dumps(goals, indent=2))