346 lines
12 KiB
Python
346 lines
12 KiB
Python
#!/usr/bin/env python3
|
|
"""Generic thread keepalive for bl side chats.
|
|
|
|
Pattern proven by the heartbeat job overnight: a state file maps a stable
|
|
key -> thread UUID, and each tick navigates directly to
|
|
https://muse.ai/thread/<uuid> (UUID reuse fix in muse-chat-api.py
|
|
cmd_sidechat_use) instead of name-based sidebar lookup.
|
|
|
|
This module generalizes that pattern:
|
|
|
|
- ANY side chat can register for keepalive (not just job sidechats).
|
|
- Each tick: ensure the thread is reachable; ping it if idle past threshold;
|
|
recreate it if the stored UUID is dead.
|
|
- State lives in keepalive-threads.json (atomic write via tmp+rename).
|
|
|
|
Usage:
|
|
from keepalive import Keepalive
|
|
|
|
ka = Keepalive(state_path="/home/super/Projects/NetVM/keepalive-threads.json")
|
|
ka.ensure("ops-watch", agent="opm") # reuse or create
|
|
ka.tick("ops-watch", agent="opm",
|
|
ping_message="[keepalive] ops-watch {ts}",
|
|
idle_threshold_s=3600) # ping if idle
|
|
|
|
The timer entry point (keepalive-timer.py) drives this from a JSON config.
|
|
"""
|
|
|
|
import json
|
|
import os
|
|
import re
|
|
import subprocess
|
|
import sys
|
|
import time
|
|
from datetime import datetime, timezone
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Paths (overridable for tests)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
NETVM_ROOT = os.environ.get("NETVM_ROOT", "/home/super/Projects/NetVM")
|
|
NETVM_EXEC = os.path.join(NETVM_ROOT, "bin", "netvm-exec.sh")
|
|
CHAT_API = os.path.join(NETVM_ROOT, "bin", "muse-chat-api.py")
|
|
DEFAULT_STATE = os.path.join(NETVM_ROOT, "keepalive-threads.json")
|
|
DEFAULT_LOG = os.path.join(NETVM_ROOT, "keepalive-log.jsonl")
|
|
|
|
UUID_RE = re.compile(
|
|
r"/thread/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-"
|
|
r"[0-9a-f]{4}-[0-9a-f]{12})"
|
|
)
|
|
|
|
|
|
def _utcnow():
|
|
return datetime.now(timezone.utc).isoformat()
|
|
|
|
|
|
def _ts():
|
|
return time.time()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Registry
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class Keepalive:
|
|
"""Thread keepalive registry + ensure/ping/check operations."""
|
|
|
|
def __init__(self, state_path=DEFAULT_STATE, log_path=DEFAULT_LOG,
|
|
dry_run=False):
|
|
self.state_path = state_path
|
|
self.log_path = log_path
|
|
self.dry_run = dry_run
|
|
|
|
# -- state ---------------------------------------------------------
|
|
|
|
def load_state(self):
|
|
if os.path.exists(self.state_path):
|
|
try:
|
|
with open(self.state_path) as f:
|
|
data = json.load(f)
|
|
return data if isinstance(data, dict) else {}
|
|
except Exception:
|
|
return {}
|
|
return {}
|
|
|
|
def save_state(self, state):
|
|
if self.dry_run:
|
|
return
|
|
tmp = self.state_path + ".tmp"
|
|
with open(tmp, "w") as f:
|
|
json.dump(state, f, indent=2)
|
|
os.replace(tmp, self.state_path)
|
|
|
|
def log_event(self, event_type, data):
|
|
if self.dry_run:
|
|
return
|
|
entry = {"ts": _utcnow(), "type": event_type, **data}
|
|
try:
|
|
with open(self.log_path, "a") as f:
|
|
f.write(json.dumps(entry) + "\n")
|
|
except Exception:
|
|
pass
|
|
|
|
# -- browser ops (mirrors job-dispatch.py; UUID navigation, not names) --
|
|
|
|
def _run(self, agent, *api_args, timeout=60):
|
|
"""Run muse-chat-api.py in the agent's netns via netvm-exec.sh."""
|
|
cmd = [NETVM_EXEC, agent, "--", "python3", CHAT_API,
|
|
"--account", agent] + list(api_args)
|
|
try:
|
|
r = subprocess.run(cmd, capture_output=True, text=True,
|
|
timeout=timeout)
|
|
return r.returncode, r.stdout.strip(), r.stderr.strip()
|
|
except Exception as e:
|
|
return -1, "", str(e)
|
|
|
|
def navigate_to_uuid(self, agent, thread_uuid):
|
|
"""Go directly to https://muse.ai/thread/<uuid>.
|
|
|
|
This is the heartbeat UUID-reuse fix: cmd_sidechat_use navigates by
|
|
URL for UUID args instead of searching sidebar titles by name.
|
|
Returns True if the command succeeded.
|
|
"""
|
|
if self.dry_run:
|
|
print(f"[DRY] navigate {agent} -> {thread_uuid}")
|
|
return True
|
|
rc, out, err = self._run(agent, "sidechat", "use", thread_uuid,
|
|
timeout=45)
|
|
return rc == 0
|
|
|
|
def current_thread_uuid(self, agent):
|
|
"""Read the browser's current URL and extract the thread UUID."""
|
|
if self.dry_run:
|
|
return None
|
|
rc, out, err = self._run(agent, "url", timeout=30)
|
|
if rc != 0:
|
|
return None
|
|
m = UUID_RE.search(out or "")
|
|
return m.group(1) if m else None
|
|
|
|
def create_sidechat(self, agent, name=None):
|
|
"""Create a new sidechat in the agent's account. Returns True."""
|
|
if self.dry_run:
|
|
print(f"[DRY] create sidechat for {agent}")
|
|
return True
|
|
rc, out, err = self._run(agent, "sidechat", "create", timeout=90)
|
|
return rc == 0 and "Created:" in (out or "")
|
|
|
|
def send_message(self, agent, message):
|
|
"""Send a message to the current chat. Returns True on success."""
|
|
if self.dry_run:
|
|
print(f"[DRY] send to {agent}: {message[:80]}")
|
|
return True
|
|
rc, out, err = self._run(agent, "send", message, timeout=60)
|
|
return rc == 0
|
|
|
|
# -- keepalive ops ---------------------------------------------------
|
|
|
|
def check(self, key, agent):
|
|
"""Verify the registered thread is reachable.
|
|
|
|
Navigates to the stored UUID and confirms the browser lands on it.
|
|
Returns (ok, thread_uuid). Updates last_verified_ts on success.
|
|
"""
|
|
state = self.load_state()
|
|
rec = state.get(key)
|
|
if not rec or not rec.get("thread_uuid"):
|
|
return False, None
|
|
uuid = rec["thread_uuid"]
|
|
if not self.navigate_to_uuid(agent, uuid):
|
|
self.log_event("keepalive_check_failed",
|
|
{"key": key, "agent": agent, "thread_uuid": uuid,
|
|
"reason": "navigate_failed"})
|
|
return False, uuid
|
|
cur = self.current_thread_uuid(agent)
|
|
if cur == uuid:
|
|
rec["last_verified_ts"] = _ts()
|
|
rec["consecutive_failures"] = 0
|
|
state[key] = rec
|
|
self.save_state(state)
|
|
self.log_event("keepalive_check_ok",
|
|
{"key": key, "agent": agent, "thread_uuid": uuid})
|
|
return True, uuid
|
|
self.log_event("keepalive_check_failed",
|
|
{"key": key, "agent": agent, "thread_uuid": uuid,
|
|
"reason": "url_mismatch", "current": cur})
|
|
return False, uuid
|
|
|
|
def ensure(self, key, agent, name=None):
|
|
"""Ensure the thread exists and is reachable; create if needed.
|
|
|
|
Returns (ok, thread_uuid). Mirrors the job-dispatch.py reuse-or-
|
|
create flow: try stored UUID first, fall back to creation, then
|
|
capture the new UUID from the browser URL.
|
|
"""
|
|
state = self.load_state()
|
|
rec = state.get(key)
|
|
if rec and rec.get("thread_uuid"):
|
|
ok, uuid = self.check(key, agent)
|
|
if ok:
|
|
return True, uuid
|
|
# Stored thread is dead — fall through to recreate.
|
|
print(f"keepalive: stored thread for {key} unreachable, "
|
|
f"recreating", file=sys.stderr)
|
|
# Create a fresh sidechat.
|
|
if not self.create_sidechat(agent, name=name):
|
|
self.log_event("keepalive_create_failed",
|
|
{"key": key, "agent": agent})
|
|
return False, None
|
|
# Capture the new thread UUID from the browser URL.
|
|
uuid = None
|
|
if not self.dry_run:
|
|
for _ in range(15):
|
|
time.sleep(1)
|
|
uuid = self.current_thread_uuid(agent)
|
|
if uuid:
|
|
break
|
|
else:
|
|
uuid = "dry-run-uuid"
|
|
if not uuid:
|
|
self.log_event("keepalive_create_failed",
|
|
{"key": key, "agent": agent,
|
|
"reason": "uuid_capture_failed"})
|
|
return False, None
|
|
now = _ts()
|
|
state[key] = {
|
|
"thread_uuid": uuid,
|
|
"agent": agent,
|
|
"created_ts": now,
|
|
"last_ping_ts": 0,
|
|
"last_verified_ts": now,
|
|
"last_activity_ts": now,
|
|
"ping_count": 0,
|
|
"consecutive_failures": 0,
|
|
}
|
|
self.save_state(state)
|
|
self.log_event("keepalive_created",
|
|
{"key": key, "agent": agent, "thread_uuid": uuid})
|
|
return True, uuid
|
|
|
|
def ping(self, key, agent, message=None):
|
|
"""Send a keepalive ping into the thread.
|
|
|
|
Navigates to the thread first (cheap no-op if already there),
|
|
then sends the ping message. Updates last_ping_ts / ping_count.
|
|
"""
|
|
state = self.load_state()
|
|
rec = state.get(key)
|
|
if not rec or not rec.get("thread_uuid"):
|
|
return False
|
|
uuid = rec["thread_uuid"]
|
|
if not self.navigate_to_uuid(agent, uuid):
|
|
return False
|
|
msg = (message or "[keepalive:{key}] tick {ts} {uuid}").format(
|
|
key=key, ts=_utcnow(), uuid=uuid)
|
|
if not self.send_message(agent, msg):
|
|
self.log_event("keepalive_ping_failed",
|
|
{"key": key, "agent": agent, "thread_uuid": uuid})
|
|
return False
|
|
now = _ts()
|
|
rec["last_ping_ts"] = now
|
|
rec["last_activity_ts"] = now
|
|
rec["ping_count"] = rec.get("ping_count", 0) + 1
|
|
state[key] = rec
|
|
self.save_state(state)
|
|
self.log_event("keepalive_ping",
|
|
{"key": key, "agent": agent, "thread_uuid": uuid,
|
|
"ping_count": rec["ping_count"]})
|
|
return True
|
|
|
|
def note_activity(self, key):
|
|
"""Record external activity (e.g. a job just sent to the thread)
|
|
so the idle timer doesn't ping unnecessarily."""
|
|
state = self.load_state()
|
|
rec = state.get(key)
|
|
if rec:
|
|
rec["last_activity_ts"] = _ts()
|
|
state[key] = rec
|
|
self.save_state(state)
|
|
|
|
def tick(self, key, agent, idle_threshold_s=3600, max_retries=3,
|
|
ping_message=None):
|
|
"""One keepalive tick for a registered thread.
|
|
|
|
- Ensures the thread exists (recreate if dead).
|
|
- Pings only if idle longer than idle_threshold_s.
|
|
- After max_retries consecutive failures, forces recreation.
|
|
Returns a status string: ok | pinged | recreated | failed.
|
|
"""
|
|
state = self.load_state()
|
|
rec = state.get(key)
|
|
|
|
if not rec or not rec.get("thread_uuid"):
|
|
ok, _ = self.ensure(key, agent)
|
|
return "recreated" if ok else "failed"
|
|
|
|
failures = rec.get("consecutive_failures", 0)
|
|
if failures >= max_retries:
|
|
# Force recreation: drop the dead UUID and re-ensure.
|
|
self.log_event("keepalive_force_recreate",
|
|
{"key": key, "agent": agent,
|
|
"failures": failures,
|
|
"old_uuid": rec.get("thread_uuid")})
|
|
rec["thread_uuid"] = None
|
|
state[key] = rec
|
|
self.save_state(state)
|
|
ok, _ = self.ensure(key, agent)
|
|
return "recreated" if ok else "failed"
|
|
|
|
ok, _ = self.check(key, agent)
|
|
if not ok:
|
|
state = self.load_state()
|
|
rec = state.get(key, {})
|
|
rec["consecutive_failures"] = failures + 1
|
|
state[key] = rec
|
|
self.save_state(state)
|
|
return "failed"
|
|
|
|
now = _ts()
|
|
idle_for = now - rec.get("last_activity_ts", 0)
|
|
if idle_for >= idle_threshold_s:
|
|
if self.ping(key, agent, message=ping_message):
|
|
return "pinged"
|
|
return "failed"
|
|
return "ok"
|
|
|
|
def unregister(self, key):
|
|
"""Remove a thread from keepalive (does not delete the sidechat)."""
|
|
state = self.load_state()
|
|
if key in state:
|
|
del state[key]
|
|
self.save_state(state)
|
|
self.log_event("keepalive_unregistered", {"key": key})
|
|
return True
|
|
return False
|
|
|
|
def status(self):
|
|
"""Return the full registry with computed idle times."""
|
|
state = self.load_state()
|
|
now = _ts()
|
|
out = {}
|
|
for key, rec in state.items():
|
|
r = dict(rec)
|
|
r["idle_s"] = int(now - rec.get("last_activity_ts", now))
|
|
out[key] = r
|
|
return out
|