- Reset master to upstream/main (16,697 commits) - Overlay 2,271 local-only files (skills, tools, workspace, configs, apps) - Restore IDENTITY.md and USER.md templates - Build verified, gateway running, Discord working Co-Authored-By: Claude Opus 4.6 <[email protected]>
228 lines
7.3 KiB
Python
228 lines
7.3 KiB
Python
#!/usr/bin/env python3
|
|
import json
|
|
import os
|
|
import sys
|
|
import time
|
|
import uuid
|
|
from urllib.error import HTTPError, URLError
|
|
from urllib.request import Request, urlopen
|
|
|
|
from james.lib.log import get_logger
|
|
from james.lib.schema import validate_event
|
|
|
|
log = get_logger("analyst")
|
|
|
|
|
|
def redis_url() -> str:
|
|
url = os.environ.get("REDIS_URL")
|
|
if url:
|
|
return url
|
|
host = os.environ.get("REDIS_HOST", "127.0.0.1")
|
|
port = os.environ.get("REDIS_PORT", "6379")
|
|
password = os.environ.get("REDIS_PASSWORD")
|
|
auth = f":{password}@" if password else ""
|
|
return f"redis://{auth}{host}:{port}/0"
|
|
|
|
|
|
def base_event(event_type: str, severity: str, payload: dict) -> dict:
|
|
return {
|
|
"id": f"evt-{time.strftime('%Y%m%d')}-{uuid.uuid4().hex[:8]}",
|
|
"source": os.environ.get("JAMES_NODE", "analyst"),
|
|
"type": event_type,
|
|
"timestamp": int(time.time()),
|
|
"severity": severity,
|
|
"payload": payload,
|
|
}
|
|
|
|
|
|
def severity_to_priority(severity: str) -> str:
|
|
mapping = {
|
|
"low": "low",
|
|
"medium": "medium",
|
|
"high": "high",
|
|
"critical": "critical",
|
|
}
|
|
return mapping.get(severity, "low")
|
|
|
|
|
|
def ollama_generate(prompt: str) -> dict | None:
|
|
url = os.environ.get("OLLAMA_URL", "http://127.0.0.1:11434")
|
|
model = os.environ.get("OLLAMA_MODEL", "qwen2.5-coder")
|
|
timeout = float(os.environ.get("OLLAMA_TIMEOUT", "20"))
|
|
|
|
payload = {
|
|
"model": model,
|
|
"prompt": prompt,
|
|
"stream": False,
|
|
"format": "json",
|
|
"options": {
|
|
"temperature": 0.2,
|
|
},
|
|
}
|
|
data = json.dumps(payload).encode("utf-8")
|
|
request = Request(f"{url}/api/generate", data=data, headers={"Content-Type": "application/json"})
|
|
|
|
try:
|
|
with urlopen(request, timeout=timeout) as response:
|
|
body = response.read().decode("utf-8")
|
|
except (HTTPError, URLError, TimeoutError, ValueError) as exc:
|
|
log.warning("ollama request failed", model=model, error=str(exc))
|
|
return None
|
|
|
|
try:
|
|
raw = json.loads(body)
|
|
except json.JSONDecodeError as exc:
|
|
log.warning("ollama response not valid JSON", error=str(exc))
|
|
return None
|
|
|
|
text = raw.get("response") if isinstance(raw, dict) else None
|
|
if not text:
|
|
log.warning("ollama returned empty response", model=model)
|
|
return None
|
|
|
|
try:
|
|
return json.loads(text)
|
|
except json.JSONDecodeError as exc:
|
|
log.warning("ollama inner response not valid JSON", error=str(exc), raw=text[:200])
|
|
return None
|
|
|
|
|
|
def heuristic_decision(event: dict) -> dict:
|
|
event_type = event.get("type", "unknown")
|
|
severity = event.get("severity", "low")
|
|
payload = event.get("payload", {})
|
|
actions = []
|
|
|
|
if event_type in {"container.stopped", "service.unhealthy"}:
|
|
container = payload.get("container") or payload.get("service")
|
|
if container:
|
|
actions.append(
|
|
{
|
|
"action": "restart_container",
|
|
"target": container,
|
|
"reason": f"{event_type} detected",
|
|
}
|
|
)
|
|
elif event_type in {"disk.high_usage", "memory.high_usage", "cpu.high_usage"}:
|
|
actions.append(
|
|
{
|
|
"action": "refresh_status",
|
|
"target": event.get("source", "unknown"),
|
|
"reason": f"{event_type} threshold exceeded",
|
|
}
|
|
)
|
|
elif event_type.startswith("mail."):
|
|
actions.append(
|
|
{
|
|
"action": "send_discord",
|
|
"target": "mail_summary",
|
|
"reason": f"New mail event: {event_type}",
|
|
}
|
|
)
|
|
|
|
summary = payload.get("summary") or f"{event_type} from {event.get('source', 'unknown')}"
|
|
return {
|
|
"summary": summary,
|
|
"priority": severity_to_priority(severity),
|
|
"confidence": 0.55,
|
|
"suggested_actions": actions,
|
|
"event": {
|
|
"id": event.get("id"),
|
|
"type": event_type,
|
|
"source": event.get("source"),
|
|
"severity": severity,
|
|
},
|
|
}
|
|
|
|
|
|
def build_prompt(event: dict) -> str:
|
|
return (
|
|
"You are James Analyst. Output JSON only with keys: summary, priority, confidence, "
|
|
"suggested_actions (array of objects with action, target, reason), and event (id,type,source,severity). "
|
|
"Use only allowed actions: restart_container, refresh_status, reassign_workload, collect_logs, "
|
|
"update_metrics, organize_files, send_discord, send_email, publish_message, "
|
|
"home_automation_changes, create_github_issue, post_to_social. "
|
|
"If no action is needed, return an empty suggested_actions array.\n\n"
|
|
f"Event JSON:\n{json.dumps(event)}"
|
|
)
|
|
|
|
|
|
def analyze_event(event: dict) -> dict:
|
|
response = ollama_generate(build_prompt(event))
|
|
if response is None:
|
|
log.debug("falling back to heuristic", event_id=event.get("id"))
|
|
return heuristic_decision(event)
|
|
|
|
# Basic sanity checks and defaulting
|
|
response.setdefault("summary", f"{event.get('type', 'event')} detected")
|
|
response.setdefault("priority", severity_to_priority(event.get("severity", "low")))
|
|
response.setdefault("confidence", 0.7)
|
|
response.setdefault("suggested_actions", [])
|
|
response.setdefault(
|
|
"event",
|
|
{
|
|
"id": event.get("id"),
|
|
"type": event.get("type"),
|
|
"source": event.get("source"),
|
|
"severity": event.get("severity"),
|
|
},
|
|
)
|
|
return response
|
|
|
|
|
|
def main() -> int:
|
|
channels = os.environ.get("JAMES_ANALYST_CHANNELS")
|
|
if channels:
|
|
channels = [chan.strip() for chan in channels.split(",") if chan.strip()]
|
|
else:
|
|
channels = ["james.events.system", "james.events.mail"]
|
|
|
|
try:
|
|
import redis
|
|
except Exception as exc: # pragma: no cover - runtime guard
|
|
log.error("redis is required", exc_info=True)
|
|
return 2
|
|
|
|
url = redis_url()
|
|
try:
|
|
r = redis.Redis.from_url(url)
|
|
r.ping()
|
|
except Exception:
|
|
log.error("redis connection failed", exc_info=True)
|
|
return 1
|
|
|
|
pubsub = r.pubsub()
|
|
pubsub.subscribe(*channels)
|
|
log.info("started", channels=channels)
|
|
|
|
for message in pubsub.listen():
|
|
if message.get("type") != "message":
|
|
continue
|
|
data = message.get("data")
|
|
if isinstance(data, bytes):
|
|
data = data.decode("utf-8", errors="replace")
|
|
try:
|
|
event = json.loads(data)
|
|
except json.JSONDecodeError as exc:
|
|
log.warning("skipping malformed event", error=str(exc), raw=data[:200])
|
|
continue
|
|
|
|
decision_payload = analyze_event(event)
|
|
decision_type = "decision.action_suggested" if decision_payload.get("suggested_actions") else "decision.no_action"
|
|
decision_event = base_event(decision_type, event.get("severity", "low"), decision_payload)
|
|
|
|
try:
|
|
validate_event(decision_event)
|
|
except Exception as exc:
|
|
log.error("decision event failed validation", event_id=event.get("id"), exc_info=True)
|
|
continue
|
|
|
|
r.publish("james.decisions", json.dumps(decision_event))
|
|
log.debug("published decision", event_id=event.get("id"), decision_type=decision_type)
|
|
|
|
return 0
|
|
|
|
|
|
if __name__ == "__main__":
|
|
raise SystemExit(main())
|