Improve error handling: propagate corrupt-JSON errors, stop swallowing failures
Co-Authored-By: Devin AI <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
parent
fb3a1979dc
commit
3cf939c403
|
|
@ -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():
|
||||
|
|
|
|||
188
server.py
188
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))
|
||||
|
|
|
|||
Loading…
Reference in New Issue