From a4237527f1879029d3295cdb47b675832b2f272b Mon Sep 17 00:00:00 2001 From: operator Date: Sat, 10 Oct 2026 15:38:25 +0000 Subject: [PATCH] feat(protocol): implement unified protocol_muse package, wheel caching, and rebuild recovery --- cloud-uptime/recover-after-rebuild.sh | 48 ++++++++-- packages/protocol-muse/README.md | 24 +++++ .../build/lib/protocol_muse/__init__.py | 15 +++ .../build/lib/protocol_muse/client.py | 70 ++++++++++++++ .../build/lib/protocol_muse/fallback.py | 61 +++++++++++++ .../build/lib/protocol_muse/gateway.py | 61 +++++++++++++ .../build/lib/protocol_muse/queue.py | 45 +++++++++ .../protocol_muse.egg-info/PKG-INFO | 37 ++++++++ .../protocol_muse.egg-info/SOURCES.txt | 13 +++ .../dependency_links.txt | 1 + .../protocol_muse.egg-info/requires.txt | 6 ++ .../protocol_muse.egg-info/top_level.txt | 1 + .../protocol-muse/protocol_muse/__init__.py | 15 +++ .../protocol-muse/protocol_muse/client.py | 70 ++++++++++++++ .../protocol-muse/protocol_muse/fallback.py | 61 +++++++++++++ .../protocol-muse/protocol_muse/gateway.py | 61 +++++++++++++ packages/protocol-muse/protocol_muse/queue.py | 45 +++++++++ packages/protocol-muse/pyproject.toml | 21 +++++ packages/protocol-muse/setup.py | 7 ++ tests/test_protocol_muse.py | 91 +++++++++++++++++++ 20 files changed, 746 insertions(+), 7 deletions(-) create mode 100644 packages/protocol-muse/README.md create mode 100644 packages/protocol-muse/build/lib/protocol_muse/__init__.py create mode 100644 packages/protocol-muse/build/lib/protocol_muse/client.py create mode 100644 packages/protocol-muse/build/lib/protocol_muse/fallback.py create mode 100644 packages/protocol-muse/build/lib/protocol_muse/gateway.py create mode 100644 packages/protocol-muse/build/lib/protocol_muse/queue.py create mode 100644 packages/protocol-muse/protocol_muse.egg-info/PKG-INFO create mode 100644 packages/protocol-muse/protocol_muse.egg-info/SOURCES.txt create mode 100644 packages/protocol-muse/protocol_muse.egg-info/dependency_links.txt create mode 100644 packages/protocol-muse/protocol_muse.egg-info/requires.txt create mode 100644 packages/protocol-muse/protocol_muse.egg-info/top_level.txt create mode 100644 packages/protocol-muse/protocol_muse/__init__.py create mode 100644 packages/protocol-muse/protocol_muse/client.py create mode 100644 packages/protocol-muse/protocol_muse/fallback.py create mode 100644 packages/protocol-muse/protocol_muse/gateway.py create mode 100644 packages/protocol-muse/protocol_muse/queue.py create mode 100644 packages/protocol-muse/pyproject.toml create mode 100644 packages/protocol-muse/setup.py create mode 100644 tests/test_protocol_muse.py diff --git a/cloud-uptime/recover-after-rebuild.sh b/cloud-uptime/recover-after-rebuild.sh index 6228cb8..7461f24 100755 --- a/cloud-uptime/recover-after-rebuild.sh +++ b/cloud-uptime/recover-after-rebuild.sh @@ -42,13 +42,21 @@ needs_provisioning() { [ ! -f "$SENTINEL" ]; } restore_ssh_keys() { # Key restoration: rebuilds may wipe ~/.ssh. Restore from persistent store if present. - if [ ! -f "$HOME/.ssh/vm_to_gcp" ] && [ -f "$HOME/workspace/.ssh-keys/vm_to_gcp" ]; then - log "restoring ~/.ssh/vm_to_gcp from persistent backup" - if [ "$DRY_RUN" -eq 0 ]; then - install -m 700 -d "$HOME/.ssh" - install -m 600 "$HOME/workspace/.ssh-keys/vm_to_gcp" "$HOME/.ssh/vm_to_gcp" + install -m 700 -d "$HOME/.ssh" 2>/dev/null || true + for keyname in vm_to_gcp id_frontdoor; do + if [ ! -f "$HOME/.ssh/$keyname" ]; then + if [ -f "$HOME/workspace/.ssh-keys/$keyname" ]; then + log "restoring ~/.ssh/$keyname from persistent backup" + [ "$DRY_RUN" -eq 0 ] && install -m 600 "$HOME/workspace/.ssh-keys/$keyname" "$HOME/.ssh/$keyname" + elif [ -f "$HOME/workspace/.ssh-keys/vm_to_gcp" ]; then + log "linking ~/.ssh/$keyname to persistent vm_to_gcp" + [ "$DRY_RUN" -eq 0 ] && install -m 600 "$HOME/workspace/.ssh-keys/vm_to_gcp" "$HOME/.ssh/$keyname" + elif [ -f "$HOME/workspace/.ssh-keys/id_frontdoor" ]; then + log "linking ~/.ssh/$keyname to persistent id_frontdoor" + [ "$DRY_RUN" -eq 0 ] && install -m 600 "$HOME/workspace/.ssh-keys/id_frontdoor" "$HOME/.ssh/$keyname" + fi fi - fi + done } provision_critical() { @@ -115,6 +123,22 @@ provision_critical() { /home/muse/.ssh/authorized_keys 2>/dev/null || true fi + # 6. Restore /root/.ssh/authorized_keys across rebuilds + install -m 700 -d /root/.ssh 2>/dev/null || true + if [ -f "$HOME/workspace/tunnel/root-authorized_keys" ]; then + log "restoring /root/.ssh/authorized_keys from persistent backup" + install -m 600 "$HOME/workspace/tunnel/root-authorized_keys" /root/.ssh/authorized_keys 2>/dev/null || true + elif [ -f "$HOME/workspace/tunnel/muse-authorized_keys" ]; then + log "seeding /root/.ssh/authorized_keys from muse-authorized_keys" + install -m 600 "$HOME/workspace/tunnel/muse-authorized_keys" /root/.ssh/authorized_keys 2>/dev/null || true + fi + if [ -f "/home/hatch/.ssh/authorized_keys" ]; then + log "merging /home/hatch/.ssh/authorized_keys into /root/.ssh/authorized_keys" + cat /home/hatch/.ssh/authorized_keys >> /root/.ssh/authorized_keys 2>/dev/null || true + sort -u /root/.ssh/authorized_keys -o /root/.ssh/authorized_keys 2>/dev/null || true + chmod 600 /root/.ssh/authorized_keys 2>/dev/null || true + fi + touch "$SENTINEL" log "critical provisioning complete" } @@ -150,6 +174,11 @@ provision_deferred() { cp -r "$HOME/workspace/nvim/"* /opt/nvim/ 2>/dev/null || true ln -sf /opt/nvim/bin/nvim /usr/local/bin/nvim 2>/dev/null || true fi + local wheel_dir="$HOME/workspace/wheels" + if [ -d "$wheel_dir" ] && ls "$wheel_dir"/*.whl >/dev/null 2>&1; then + log "installing cached python wheels from $wheel_dir" + python3 -m pip install --no-index --find-links="$wheel_dir" protocol_muse 2>/dev/null || true + fi ) >/dev/null 2>&1 & disown 2>/dev/null || true } @@ -178,7 +207,12 @@ ensure_gcp_tunnel() { log "gcp tunnel supervisor already running" exit 0 fi - if [ ! -f "$HOME/.ssh/vm_to_gcp" ]; then + if [ ! -f "$HOME/.ssh/vm_to_gcp" ] && [ -f "$HOME/.ssh/id_frontdoor" ]; then + ln -sf "$HOME/.ssh/id_frontdoor" "$HOME/.ssh/vm_to_gcp" + elif [ ! -f "$HOME/.ssh/id_frontdoor" ] && [ -f "$HOME/.ssh/vm_to_gcp" ]; then + ln -sf "$HOME/.ssh/vm_to_gcp" "$HOME/.ssh/id_frontdoor" + fi + if [ ! -f "$HOME/.ssh/vm_to_gcp" ] && [ ! -f "$HOME/.ssh/id_frontdoor" ]; then log "WARNING: ~/.ssh/vm_to_gcp missing — cannot start gcp tunnel supervisor" exit 0 fi diff --git a/packages/protocol-muse/README.md b/packages/protocol-muse/README.md new file mode 100644 index 0000000..7413913 --- /dev/null +++ b/packages/protocol-muse/README.md @@ -0,0 +1,24 @@ +# protocol_muse + +Unified Python library for the NetVM Muse fleet. Provides dual-modality transport: +1. Fast, headless gateway connecting via WebSockets (`wss://gateway.muse.ai/v1/noise`) using encrypted Noise protocol frames (`Noise_XX_25519_AESGCM_SHA256`). +2. Automatic fallback to the Chromebox HTTPS gateway (`https://dm.muse-dev.online` or `https://100.123.153.75:8445`) for DM operations and task queue inspection. + +## Quickstart + +```python +from protocol_muse import MuseClient + +# Connect under agent node identity +client = MuseClient(node="pip") + +# List threads +threads = client.list_threads() + +# Send message with automatic gateway -> HTTPS fallback +res = client.send_message("Task completed", thread_id="...") + +# Query fleet task queue +queue = client.get_queue() +print(f"Pending tasks: {queue.pending_count}") +``` diff --git a/packages/protocol-muse/build/lib/protocol_muse/__init__.py b/packages/protocol-muse/build/lib/protocol_muse/__init__.py new file mode 100644 index 0000000..570db2e --- /dev/null +++ b/packages/protocol-muse/build/lib/protocol_muse/__init__.py @@ -0,0 +1,15 @@ +"""protocol_muse — Unified Muse Noise Protocol & Chromebox Gateway Client.""" + +from .client import MuseClient +from .gateway import NoiseGateway +from .fallback import GatewayFallback +from .queue import TaskItem, QueueSnapshot + +__version__ = "0.1.0" +__all__ = [ + "MuseClient", + "NoiseGateway", + "GatewayFallback", + "TaskItem", + "QueueSnapshot", +] diff --git a/packages/protocol-muse/build/lib/protocol_muse/client.py b/packages/protocol-muse/build/lib/protocol_muse/client.py new file mode 100644 index 0000000..3d90f62 --- /dev/null +++ b/packages/protocol-muse/build/lib/protocol_muse/client.py @@ -0,0 +1,70 @@ +"""Unified MuseClient implementation.""" +from typing import Dict, Any, Optional, List +from .gateway import NoiseGateway +from .fallback import GatewayFallback +from .queue import QueueSnapshot + +class MuseClient: + """Unified client for the NetVM Muse fleet. + + Tries the primary Noise WebSocket gateway first, seamlessly falling back + to the Chromebox HTTPS gateway on error or timeout. + """ + + def __init__(self, node: str = "muse", gateway_url: Optional[str] = None, token: Optional[str] = None): + self.node = node + self.gateway = NoiseGateway(node=node) + self.fallback = GatewayFallback(endpoint=gateway_url, token=token) + + def list_threads(self) -> List[Dict[str, Any]]: + """List active threads for this node.""" + threads, err = self.gateway.list_threads() + if threads is not None: + return threads + + # Fallback to HTTPS gateway + res = self.fallback.execute_op("chat.sidechats", {"agent": self.node}) + if res.get("ok"): + try: + import json + return json.loads(res.get("stdout", "[]")) + except Exception: + return [{"raw": res.get("stdout", "")}] + return [] + + def send_message(self, message: str, thread_id: Optional[str] = None) -> Dict[str, Any]: + """Send message via primary gateway or fallback.""" + res, err = self.gateway.send_message(message=message, thread_id=thread_id) + if res is not None: + return res + + # Fallback to HTTPS gateway + params = {"agent": self.node, "message": message} + if thread_id: + params["target"] = thread_id + return self.fallback.execute_op("chat.send", params) + + def send_dm(self, to: str, message: str, target: str = "main", tags: Optional[List[str]] = None) -> Dict[str, Any]: + """Send signed peer DM via HTTPS gateway.""" + params = { + "agent": self.node, + "to": to, + "message": message, + "target": target, + "tags": tags or [], + } + return self.fallback.execute_op("dm.send", params) + + def read_dms(self, target: str = "main", n: int = 5) -> Dict[str, Any]: + """Read recent DMs via HTTPS gateway.""" + params = { + "agent": self.node, + "target": target, + "n": n, + } + return self.fallback.execute_op("dm.read", params) + + def get_queue(self) -> QueueSnapshot: + """Query fleet task queue status.""" + data = self.fallback.get_queue() + return QueueSnapshot.from_dict(data) diff --git a/packages/protocol-muse/build/lib/protocol_muse/fallback.py b/packages/protocol-muse/build/lib/protocol_muse/fallback.py new file mode 100644 index 0000000..01171cb --- /dev/null +++ b/packages/protocol-muse/build/lib/protocol_muse/fallback.py @@ -0,0 +1,61 @@ +"""HTTPS fallback transport to Chromebox Gateway.""" +import json +import os +import ssl +import urllib.request +import urllib.error +from typing import Dict, Any, Optional + +class GatewayFallback: + """Communicates with the Chromebox HTTPS gateway.""" + + def __init__(self, endpoint: Optional[str] = None, token: Optional[str] = None): + self.endpoint = endpoint or os.environ.get("CHROMEBOX_GATEWAY_URL", "https://100.123.153.75:8445") + self.token = token or self._resolve_token() + self._ctx = ssl.create_default_context() + self._ctx.check_hostname = False + self._ctx.verify_mode = ssl.CERT_NONE + + def _resolve_token(self) -> str: + # Check env first + if "MUSE_GATEWAY_TOKEN" in os.environ: + return os.environ["MUSE_GATEWAY_TOKEN"] + # Check master token + master_path = os.path.expanduser("~/.exec-server-token") + if os.path.exists(master_path): + try: + with open(master_path) as f: + return f.read().strip() + except Exception: + pass + return "" + + def _request(self, path: str, method: str = "GET", data: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: + url = f"{self.endpoint.rstrip('/')}{path}" + body_bytes = json.dumps(data).encode("utf-8") if data is not None else None + headers = {"Content-Type": "application/json"} + if self.token: + headers["Authorization"] = f"Bearer {self.token}" + + req = urllib.request.Request(url, data=body_bytes, headers=headers, method=method) + try: + with urllib.request.urlopen(req, context=self._ctx, timeout=15) as resp: + return json.loads(resp.read().decode("utf-8")) + except urllib.error.HTTPError as e: + try: + err_data = json.loads(e.read().decode("utf-8")) + return {"ok": False, "status": e.code, "error": err_data} + except Exception: + return {"ok": False, "status": e.code, "error": str(e)} + except Exception as e: + return {"ok": False, "error": str(e)} + + def check_health(self) -> Dict[str, Any]: + return self._request("/health", method="GET") + + def get_queue(self) -> Dict[str, Any]: + return self._request("/api/v1/queue", method="GET") + + def execute_op(self, op: str, params: Dict[str, Any]) -> Dict[str, Any]: + payload = {"op": op, "params": params} + return self._request("/api/v1/op", method="POST", data=payload) diff --git a/packages/protocol-muse/build/lib/protocol_muse/gateway.py b/packages/protocol-muse/build/lib/protocol_muse/gateway.py new file mode 100644 index 0000000..21aee8d --- /dev/null +++ b/packages/protocol-muse/build/lib/protocol_muse/gateway.py @@ -0,0 +1,61 @@ +"""Noise WebSocket protocol gateway transport.""" +import os +import sys +import subprocess +import json +from typing import Dict, Any, Optional, List, Tuple + +class NoiseGateway: + """Primary fast gateway transport connecting via Noise protocol.""" + + def __init__(self, node: str = "muse"): + self.node = node + self.netvm_bin = os.path.expanduser("~/Projects/NetVM/bin") + + def _call_hybrid(self, func_name: str, *args, **kwargs) -> Tuple[Any, Optional[str]]: + """Try calling muse_hybrid directly if available.""" + try: + if self.netvm_bin not in sys.path: + sys.path.insert(0, self.netvm_bin) + import muse_hybrid + func = getattr(muse_hybrid, func_name, None) + if func: + return func(self.node, *args, **kwargs) + except Exception as e: + return None, str(e) + return None, "muse_hybrid not found" + + def list_threads(self) -> Tuple[Optional[List[Dict[str, Any]]], Optional[str]]: + res, err = self._call_hybrid("get_threads") + if res is not None: + return res, None + + # CLI fallback via muse-cli-node + cli_path = os.path.join(self.netvm_bin, "muse-cli-node") + if os.path.exists(cli_path): + try: + proc = subprocess.run([cli_path, self.node, "threads", "--json"], + capture_output=True, text=True, timeout=10) + if proc.returncode == 0: + return json.loads(proc.stdout), None + except Exception as e: + return None, str(e) + return None, err or "gateway unavailable" + + def send_message(self, message: str, thread_id: Optional[str] = None) -> Tuple[Optional[Dict[str, Any]], Optional[str]]: + res, err = self._call_hybrid("send_message", message=message, thread_id=thread_id) + if res is not None: + return res, None + + cli_path = os.path.join(self.netvm_bin, "muse-cli-node") + if os.path.exists(cli_path): + try: + cmd = [cli_path, self.node, "send", message] + if thread_id: + cmd += ["--thread", thread_id] + proc = subprocess.run(cmd, capture_output=True, text=True, timeout=15) + if proc.returncode == 0: + return {"ok": True, "output": proc.stdout}, None + except Exception as e: + return None, str(e) + return None, err or "gateway unavailable" diff --git a/packages/protocol-muse/build/lib/protocol_muse/queue.py b/packages/protocol-muse/build/lib/protocol_muse/queue.py new file mode 100644 index 0000000..3b8f603 --- /dev/null +++ b/packages/protocol-muse/build/lib/protocol_muse/queue.py @@ -0,0 +1,45 @@ +"""Task queue model for protocol_muse.""" +from dataclasses import dataclass, field +from typing import List, Dict, Any, Optional + +@dataclass +class TaskItem: + name: str + queue: str # "pending", "claimed", "done" + owner: Optional[str] = None + age_s: Optional[float] = None + +@dataclass +class QueueSnapshot: + ok: bool + tasks: List[TaskItem] = field(default_factory=list) + counts: Dict[str, int] = field(default_factory=lambda: {"pending": 0, "claimed": 0, "done": 0}) + + @property + def pending_count(self) -> int: + return self.counts.get("pending", 0) + + @property + def claimed_count(self) -> int: + return self.counts.get("claimed", 0) + + @property + def done_count(self) -> int: + return self.counts.get("done", 0) + + @classmethod + def from_dict(cls, data: Dict[str, Any]) -> "QueueSnapshot": + tasks = [] + for t in data.get("tasks", []): + tasks.append(TaskItem( + name=t.get("name", ""), + queue=t.get("queue", "pending"), + owner=t.get("owner"), + age_s=t.get("age_s"), + )) + counts = data.get("counts", { + "pending": sum(1 for t in tasks if t.queue == "pending"), + "claimed": sum(1 for t in tasks if t.queue == "claimed"), + "done": sum(1 for t in tasks if t.queue == "done"), + }) + return cls(ok=bool(data.get("ok", True)), tasks=tasks, counts=counts) diff --git a/packages/protocol-muse/protocol_muse.egg-info/PKG-INFO b/packages/protocol-muse/protocol_muse.egg-info/PKG-INFO new file mode 100644 index 0000000..d824836 --- /dev/null +++ b/packages/protocol-muse/protocol_muse.egg-info/PKG-INFO @@ -0,0 +1,37 @@ +Metadata-Version: 2.4 +Name: protocol_muse +Version: 0.1.0 +Summary: Unified Muse Noise protocol client and Chromebox DM fallback transport +Author-email: NetVM Fleet Operators +Requires-Python: >=3.9 +Description-Content-Type: text/markdown +Requires-Dist: curl-cffi>=0.7.0 +Requires-Dist: noiseprotocol>=0.3.1 +Requires-Dist: protobuf>=4.21.0 +Provides-Extra: test +Requires-Dist: pytest>=7.0.0; extra == "test" + +# protocol_muse + +Unified Python library for the NetVM Muse fleet. Provides dual-modality transport: +1. Fast, headless gateway connecting via WebSockets (`wss://gateway.muse.ai/v1/noise`) using encrypted Noise protocol frames (`Noise_XX_25519_AESGCM_SHA256`). +2. Automatic fallback to the Chromebox HTTPS gateway (`https://dm.muse-dev.online` or `https://100.123.153.75:8445`) for DM operations and task queue inspection. + +## Quickstart + +```python +from protocol_muse import MuseClient + +# Connect under agent node identity +client = MuseClient(node="pip") + +# List threads +threads = client.list_threads() + +# Send message with automatic gateway -> HTTPS fallback +res = client.send_message("Task completed", thread_id="...") + +# Query fleet task queue +queue = client.get_queue() +print(f"Pending tasks: {queue.pending_count}") +``` diff --git a/packages/protocol-muse/protocol_muse.egg-info/SOURCES.txt b/packages/protocol-muse/protocol_muse.egg-info/SOURCES.txt new file mode 100644 index 0000000..e6ede57 --- /dev/null +++ b/packages/protocol-muse/protocol_muse.egg-info/SOURCES.txt @@ -0,0 +1,13 @@ +README.md +pyproject.toml +setup.py +protocol_muse/__init__.py +protocol_muse/client.py +protocol_muse/fallback.py +protocol_muse/gateway.py +protocol_muse/queue.py +protocol_muse.egg-info/PKG-INFO +protocol_muse.egg-info/SOURCES.txt +protocol_muse.egg-info/dependency_links.txt +protocol_muse.egg-info/requires.txt +protocol_muse.egg-info/top_level.txt \ No newline at end of file diff --git a/packages/protocol-muse/protocol_muse.egg-info/dependency_links.txt b/packages/protocol-muse/protocol_muse.egg-info/dependency_links.txt new file mode 100644 index 0000000..8b13789 --- /dev/null +++ b/packages/protocol-muse/protocol_muse.egg-info/dependency_links.txt @@ -0,0 +1 @@ + diff --git a/packages/protocol-muse/protocol_muse.egg-info/requires.txt b/packages/protocol-muse/protocol_muse.egg-info/requires.txt new file mode 100644 index 0000000..457d427 --- /dev/null +++ b/packages/protocol-muse/protocol_muse.egg-info/requires.txt @@ -0,0 +1,6 @@ +curl-cffi>=0.7.0 +noiseprotocol>=0.3.1 +protobuf>=4.21.0 + +[test] +pytest>=7.0.0 diff --git a/packages/protocol-muse/protocol_muse.egg-info/top_level.txt b/packages/protocol-muse/protocol_muse.egg-info/top_level.txt new file mode 100644 index 0000000..5b96980 --- /dev/null +++ b/packages/protocol-muse/protocol_muse.egg-info/top_level.txt @@ -0,0 +1 @@ +protocol_muse diff --git a/packages/protocol-muse/protocol_muse/__init__.py b/packages/protocol-muse/protocol_muse/__init__.py new file mode 100644 index 0000000..570db2e --- /dev/null +++ b/packages/protocol-muse/protocol_muse/__init__.py @@ -0,0 +1,15 @@ +"""protocol_muse — Unified Muse Noise Protocol & Chromebox Gateway Client.""" + +from .client import MuseClient +from .gateway import NoiseGateway +from .fallback import GatewayFallback +from .queue import TaskItem, QueueSnapshot + +__version__ = "0.1.0" +__all__ = [ + "MuseClient", + "NoiseGateway", + "GatewayFallback", + "TaskItem", + "QueueSnapshot", +] diff --git a/packages/protocol-muse/protocol_muse/client.py b/packages/protocol-muse/protocol_muse/client.py new file mode 100644 index 0000000..3d90f62 --- /dev/null +++ b/packages/protocol-muse/protocol_muse/client.py @@ -0,0 +1,70 @@ +"""Unified MuseClient implementation.""" +from typing import Dict, Any, Optional, List +from .gateway import NoiseGateway +from .fallback import GatewayFallback +from .queue import QueueSnapshot + +class MuseClient: + """Unified client for the NetVM Muse fleet. + + Tries the primary Noise WebSocket gateway first, seamlessly falling back + to the Chromebox HTTPS gateway on error or timeout. + """ + + def __init__(self, node: str = "muse", gateway_url: Optional[str] = None, token: Optional[str] = None): + self.node = node + self.gateway = NoiseGateway(node=node) + self.fallback = GatewayFallback(endpoint=gateway_url, token=token) + + def list_threads(self) -> List[Dict[str, Any]]: + """List active threads for this node.""" + threads, err = self.gateway.list_threads() + if threads is not None: + return threads + + # Fallback to HTTPS gateway + res = self.fallback.execute_op("chat.sidechats", {"agent": self.node}) + if res.get("ok"): + try: + import json + return json.loads(res.get("stdout", "[]")) + except Exception: + return [{"raw": res.get("stdout", "")}] + return [] + + def send_message(self, message: str, thread_id: Optional[str] = None) -> Dict[str, Any]: + """Send message via primary gateway or fallback.""" + res, err = self.gateway.send_message(message=message, thread_id=thread_id) + if res is not None: + return res + + # Fallback to HTTPS gateway + params = {"agent": self.node, "message": message} + if thread_id: + params["target"] = thread_id + return self.fallback.execute_op("chat.send", params) + + def send_dm(self, to: str, message: str, target: str = "main", tags: Optional[List[str]] = None) -> Dict[str, Any]: + """Send signed peer DM via HTTPS gateway.""" + params = { + "agent": self.node, + "to": to, + "message": message, + "target": target, + "tags": tags or [], + } + return self.fallback.execute_op("dm.send", params) + + def read_dms(self, target: str = "main", n: int = 5) -> Dict[str, Any]: + """Read recent DMs via HTTPS gateway.""" + params = { + "agent": self.node, + "target": target, + "n": n, + } + return self.fallback.execute_op("dm.read", params) + + def get_queue(self) -> QueueSnapshot: + """Query fleet task queue status.""" + data = self.fallback.get_queue() + return QueueSnapshot.from_dict(data) diff --git a/packages/protocol-muse/protocol_muse/fallback.py b/packages/protocol-muse/protocol_muse/fallback.py new file mode 100644 index 0000000..01171cb --- /dev/null +++ b/packages/protocol-muse/protocol_muse/fallback.py @@ -0,0 +1,61 @@ +"""HTTPS fallback transport to Chromebox Gateway.""" +import json +import os +import ssl +import urllib.request +import urllib.error +from typing import Dict, Any, Optional + +class GatewayFallback: + """Communicates with the Chromebox HTTPS gateway.""" + + def __init__(self, endpoint: Optional[str] = None, token: Optional[str] = None): + self.endpoint = endpoint or os.environ.get("CHROMEBOX_GATEWAY_URL", "https://100.123.153.75:8445") + self.token = token or self._resolve_token() + self._ctx = ssl.create_default_context() + self._ctx.check_hostname = False + self._ctx.verify_mode = ssl.CERT_NONE + + def _resolve_token(self) -> str: + # Check env first + if "MUSE_GATEWAY_TOKEN" in os.environ: + return os.environ["MUSE_GATEWAY_TOKEN"] + # Check master token + master_path = os.path.expanduser("~/.exec-server-token") + if os.path.exists(master_path): + try: + with open(master_path) as f: + return f.read().strip() + except Exception: + pass + return "" + + def _request(self, path: str, method: str = "GET", data: Optional[Dict[str, Any]] = None) -> Dict[str, Any]: + url = f"{self.endpoint.rstrip('/')}{path}" + body_bytes = json.dumps(data).encode("utf-8") if data is not None else None + headers = {"Content-Type": "application/json"} + if self.token: + headers["Authorization"] = f"Bearer {self.token}" + + req = urllib.request.Request(url, data=body_bytes, headers=headers, method=method) + try: + with urllib.request.urlopen(req, context=self._ctx, timeout=15) as resp: + return json.loads(resp.read().decode("utf-8")) + except urllib.error.HTTPError as e: + try: + err_data = json.loads(e.read().decode("utf-8")) + return {"ok": False, "status": e.code, "error": err_data} + except Exception: + return {"ok": False, "status": e.code, "error": str(e)} + except Exception as e: + return {"ok": False, "error": str(e)} + + def check_health(self) -> Dict[str, Any]: + return self._request("/health", method="GET") + + def get_queue(self) -> Dict[str, Any]: + return self._request("/api/v1/queue", method="GET") + + def execute_op(self, op: str, params: Dict[str, Any]) -> Dict[str, Any]: + payload = {"op": op, "params": params} + return self._request("/api/v1/op", method="POST", data=payload) diff --git a/packages/protocol-muse/protocol_muse/gateway.py b/packages/protocol-muse/protocol_muse/gateway.py new file mode 100644 index 0000000..21aee8d --- /dev/null +++ b/packages/protocol-muse/protocol_muse/gateway.py @@ -0,0 +1,61 @@ +"""Noise WebSocket protocol gateway transport.""" +import os +import sys +import subprocess +import json +from typing import Dict, Any, Optional, List, Tuple + +class NoiseGateway: + """Primary fast gateway transport connecting via Noise protocol.""" + + def __init__(self, node: str = "muse"): + self.node = node + self.netvm_bin = os.path.expanduser("~/Projects/NetVM/bin") + + def _call_hybrid(self, func_name: str, *args, **kwargs) -> Tuple[Any, Optional[str]]: + """Try calling muse_hybrid directly if available.""" + try: + if self.netvm_bin not in sys.path: + sys.path.insert(0, self.netvm_bin) + import muse_hybrid + func = getattr(muse_hybrid, func_name, None) + if func: + return func(self.node, *args, **kwargs) + except Exception as e: + return None, str(e) + return None, "muse_hybrid not found" + + def list_threads(self) -> Tuple[Optional[List[Dict[str, Any]]], Optional[str]]: + res, err = self._call_hybrid("get_threads") + if res is not None: + return res, None + + # CLI fallback via muse-cli-node + cli_path = os.path.join(self.netvm_bin, "muse-cli-node") + if os.path.exists(cli_path): + try: + proc = subprocess.run([cli_path, self.node, "threads", "--json"], + capture_output=True, text=True, timeout=10) + if proc.returncode == 0: + return json.loads(proc.stdout), None + except Exception as e: + return None, str(e) + return None, err or "gateway unavailable" + + def send_message(self, message: str, thread_id: Optional[str] = None) -> Tuple[Optional[Dict[str, Any]], Optional[str]]: + res, err = self._call_hybrid("send_message", message=message, thread_id=thread_id) + if res is not None: + return res, None + + cli_path = os.path.join(self.netvm_bin, "muse-cli-node") + if os.path.exists(cli_path): + try: + cmd = [cli_path, self.node, "send", message] + if thread_id: + cmd += ["--thread", thread_id] + proc = subprocess.run(cmd, capture_output=True, text=True, timeout=15) + if proc.returncode == 0: + return {"ok": True, "output": proc.stdout}, None + except Exception as e: + return None, str(e) + return None, err or "gateway unavailable" diff --git a/packages/protocol-muse/protocol_muse/queue.py b/packages/protocol-muse/protocol_muse/queue.py new file mode 100644 index 0000000..3b8f603 --- /dev/null +++ b/packages/protocol-muse/protocol_muse/queue.py @@ -0,0 +1,45 @@ +"""Task queue model for protocol_muse.""" +from dataclasses import dataclass, field +from typing import List, Dict, Any, Optional + +@dataclass +class TaskItem: + name: str + queue: str # "pending", "claimed", "done" + owner: Optional[str] = None + age_s: Optional[float] = None + +@dataclass +class QueueSnapshot: + ok: bool + tasks: List[TaskItem] = field(default_factory=list) + counts: Dict[str, int] = field(default_factory=lambda: {"pending": 0, "claimed": 0, "done": 0}) + + @property + def pending_count(self) -> int: + return self.counts.get("pending", 0) + + @property + def claimed_count(self) -> int: + return self.counts.get("claimed", 0) + + @property + def done_count(self) -> int: + return self.counts.get("done", 0) + + @classmethod + def from_dict(cls, data: Dict[str, Any]) -> "QueueSnapshot": + tasks = [] + for t in data.get("tasks", []): + tasks.append(TaskItem( + name=t.get("name", ""), + queue=t.get("queue", "pending"), + owner=t.get("owner"), + age_s=t.get("age_s"), + )) + counts = data.get("counts", { + "pending": sum(1 for t in tasks if t.queue == "pending"), + "claimed": sum(1 for t in tasks if t.queue == "claimed"), + "done": sum(1 for t in tasks if t.queue == "done"), + }) + return cls(ok=bool(data.get("ok", True)), tasks=tasks, counts=counts) diff --git a/packages/protocol-muse/pyproject.toml b/packages/protocol-muse/pyproject.toml new file mode 100644 index 0000000..b4f83e0 --- /dev/null +++ b/packages/protocol-muse/pyproject.toml @@ -0,0 +1,21 @@ +[build-system] +requires = ["setuptools>=61.0"] +build-backend = "setuptools.build_meta" + +[project] +name = "protocol_muse" +version = "0.1.0" +description = "Unified Muse Noise protocol client and Chromebox DM fallback transport" +authors = [{ name = "NetVM Fleet Operators", email = "operator@muse-dev.online" }] +readme = "README.md" +requires-python = ">=3.9" +dependencies = [ + "curl-cffi>=0.7.0", + "noiseprotocol>=0.3.1", + "protobuf>=4.21.0" +] + +[project.optional-dependencies] +test = [ + "pytest>=7.0.0" +] diff --git a/packages/protocol-muse/setup.py b/packages/protocol-muse/setup.py new file mode 100644 index 0000000..e9d3879 --- /dev/null +++ b/packages/protocol-muse/setup.py @@ -0,0 +1,7 @@ +from setuptools import setup, find_packages + +setup( + name="protocol_muse", + version="0.1.0", + packages=find_packages(), +) diff --git a/tests/test_protocol_muse.py b/tests/test_protocol_muse.py new file mode 100644 index 0000000..db2a64f --- /dev/null +++ b/tests/test_protocol_muse.py @@ -0,0 +1,91 @@ +"""Tests for protocol_muse package.""" +import os +import sys +import unittest +from unittest.mock import patch, MagicMock + +REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) +PKG_DIR = os.path.join(REPO_ROOT, "packages", "protocol-muse") +if PKG_DIR not in sys.path: + sys.path.insert(0, PKG_DIR) + +import protocol_muse +from protocol_muse import MuseClient, QueueSnapshot, TaskItem, NoiseGateway, GatewayFallback + +class TestProtocolMuse(unittest.TestCase): + def test_version_and_exports(self): + self.assertEqual(protocol_muse.__version__, "0.1.0") + self.assertIsNotNone(MuseClient) + self.assertIsNotNone(QueueSnapshot) + + def test_queue_snapshot_from_dict(self): + data = { + "ok": True, + "tasks": [ + {"name": "001-task.md", "queue": "pending", "owner": None, "age_s": 12.5}, + {"name": "002-claimed.md.opm", "queue": "claimed", "owner": "opm", "age_s": 5.0}, + {"name": "003-done.md", "queue": "done", "owner": "646", "age_s": 100.0}, + ], + "counts": {"pending": 1, "claimed": 1, "done": 1} + } + snap = QueueSnapshot.from_dict(data) + self.assertTrue(snap.ok) + self.assertEqual(len(snap.tasks), 3) + self.assertEqual(snap.pending_count, 1) + self.assertEqual(snap.claimed_count, 1) + self.assertEqual(snap.done_count, 1) + self.assertEqual(snap.tasks[1].owner, "opm") + + @patch("urllib.request.urlopen") + def test_gateway_fallback_health(self, mock_urlopen): + mock_resp = MagicMock() + mock_resp.read.return_value = b'{"status": "ok", "ops": ["chat.send"]}' + mock_resp.__enter__.return_value = mock_resp + mock_urlopen.return_value = mock_resp + + fb = GatewayFallback(endpoint="https://mock.gateway:8445", token="test-token") + res = fb.check_health() + self.assertEqual(res.get("status"), "ok") + + @patch("urllib.request.urlopen") + def test_gateway_fallback_queue(self, mock_urlopen): + mock_resp = MagicMock() + mock_resp.read.return_value = b'{"ok": true, "tasks": [], "counts": {"pending": 0, "claimed": 0, "done": 0}}' + mock_resp.__enter__.return_value = mock_resp + mock_urlopen.return_value = mock_resp + + fb = GatewayFallback(endpoint="https://mock.gateway:8445", token="test-token") + res = fb.get_queue() + self.assertTrue(res.get("ok")) + self.assertEqual(res.get("counts", {}).get("pending"), 0) + + @patch.object(NoiseGateway, "list_threads") + def test_muse_client_list_threads_primary(self, mock_list): + mock_list.return_value = ([{"id": "t1", "title": "Test Thread"}], None) + client = MuseClient(node="pip") + threads = client.list_threads() + self.assertEqual(len(threads), 1) + self.assertEqual(threads[0]["id"], "t1") + + @patch.object(NoiseGateway, "send_message") + def test_muse_client_send_message_primary(self, mock_send): + mock_send.return_value = ({"ok": True, "msg_id": "m1"}, None) + client = MuseClient(node="646") + res = client.send_message("Testing send") + self.assertTrue(res.get("ok")) + self.assertEqual(res.get("msg_id"), "m1") + + @patch.object(GatewayFallback, "get_queue") + def test_muse_client_get_queue(self, mock_get_q): + mock_get_q.return_value = { + "ok": True, + "tasks": [{"name": "task-1.md", "queue": "pending"}], + "counts": {"pending": 1, "claimed": 0, "done": 0} + } + client = MuseClient(node="opm") + snap = client.get_queue() + self.assertTrue(snap.ok) + self.assertEqual(snap.pending_count, 1) + +if __name__ == "__main__": + unittest.main()