feat(protocol): implement unified protocol_muse package, wheel caching, and rebuild recovery
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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}")
|
||||
```
|
||||
@@ -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",
|
||||
]
|
||||
@@ -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)
|
||||
@@ -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)
|
||||
@@ -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"
|
||||
@@ -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)
|
||||
@@ -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 <operator@muse-dev.online>
|
||||
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}")
|
||||
```
|
||||
@@ -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
|
||||
@@ -0,0 +1,6 @@
|
||||
curl-cffi>=0.7.0
|
||||
noiseprotocol>=0.3.1
|
||||
protobuf>=4.21.0
|
||||
|
||||
[test]
|
||||
pytest>=7.0.0
|
||||
@@ -0,0 +1 @@
|
||||
protocol_muse
|
||||
@@ -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",
|
||||
]
|
||||
@@ -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)
|
||||
@@ -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)
|
||||
@@ -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"
|
||||
@@ -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)
|
||||
@@ -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"
|
||||
]
|
||||
@@ -0,0 +1,7 @@
|
||||
from setuptools import setup, find_packages
|
||||
|
||||
setup(
|
||||
name="protocol_muse",
|
||||
version="0.1.0",
|
||||
packages=find_packages(),
|
||||
)
|
||||
@@ -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()
|
||||
Reference in New Issue
Block a user