Remove wrongly-specced opp-dm.py and dev-dm.py; dm.py is the headless DM tool

This commit is contained in:
operator
2026-10-03 20:34:33 +00:00
parent 7272ccfeea
commit 9d6805060e
21 changed files with 739 additions and 177 deletions
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
Binary file not shown.
+81
View File
@@ -0,0 +1,81 @@
#!/usr/bin/env python3
"""Per-account session vitality check (runs INSIDE the node's netns).
Usage: sudo ip netns exec warp-<node> python3 accounts-health.py <cdp_port>
Probes the account's browser via CDP:
- browser reachable
- muse.ai tab present
- login markers (heuristic: "Log in" button vs user content)
Prints JSON to stdout. Exit 0 on success, 1 if the browser is unreachable.
"""
import json, sys, urllib.request
CDP_PORT = int(sys.argv[1]) if len(sys.argv) > 1 else 9410
def http(path, timeout=5):
with urllib.request.urlopen(f"http://127.0.0.1:{CDP_PORT}{path}",
timeout=timeout) as r:
return json.loads(r.read())
result = {"browser_up": False, "muse_tab": False,
"session_alive": None, "title": "", "url": "", "detail": ""}
try:
http("/json/version")
result["browser_up"] = True
except Exception as e:
result["detail"] = f"CDP unreachable: {e}"
print(json.dumps(result)); sys.exit(1)
try:
tabs = http("/json/list")
except Exception as e:
result["detail"] = f"tab list failed: {e}"
print(json.dumps(result)); sys.exit(0)
muse_tab = None
for t in tabs:
url = t.get("url", "")
if "muse.ai" in url and t.get("type") == "page":
muse_tab = t
break
if not muse_tab:
result["detail"] = "no muse.ai tab open"
print(json.dumps(result)); sys.exit(0)
result["muse_tab"] = True
result["url"] = muse_tab.get("url", "")
result["title"] = muse_tab.get("title", "")
# DOM heuristic via the tab's debugger socket
try:
import websocket
ws = websocket.create_connection(muse_tab["webSocketDebuggerUrl"], timeout=10)
js = """JSON.stringify({
loginButtons: [...document.querySelectorAll('button')].filter(
b => /^\\s*log\\s*in\\s*$/i.test(b.innerText)).map(b => b.innerText.trim()),
hasAvatar: !!document.querySelector(
'img[alt*="avatar" i], [data-testid*="avatar" i], [aria-label*="profile" i]'),
title: document.title,
url: location.href
})"""
ws.send(json.dumps({"id": 1, "method": "Runtime.evaluate",
"params": {"expression": js, "returnByValue": True}}))
resp = json.loads(ws.recv())
ws.close()
dom = json.loads(resp["result"]["result"]["value"])
# Heuristic: login buttons present + no avatar => logged out.
# No login buttons (or avatar present) => likely logged in.
if dom["loginButtons"] and not dom["hasAvatar"]:
result["session_alive"] = False
result["detail"] = f"login wall visible: {dom['loginButtons'][:2]}"
else:
result["session_alive"] = True
result["detail"] = "no login wall detected"
except Exception as e:
result["detail"] = f"DOM check failed: {e}"
print(json.dumps(result))
+97
View File
@@ -0,0 +1,97 @@
#!/usr/bin/env bash
# accounts-health.sh — account session vitality reporter for the front-door network.
#
# Reads ACCOUNTS.md, probes each account's browser via CDP (inside its NetVM
# netns) for session liveness, and emits a JSON report. Optionally signs and
# POSTs it to the board health ingest, following the health-report.sh
# convention: payload is <machine>\n<ts>\n<facts-json>, namespace "health".
#
# Usage: accounts-health.sh [--no-post]
#
# Cron (on bl, every 15 min):
# */15 * * * * ~/Projects/NetVM/bin/accounts-health.sh >/dev/null 2>&1
#
# Health key setup (once, on bl):
# ssh-keygen -t ed25519 -N "" -f ~/.ssh/muse-health
# # operator registers the pubkey on the VM:
# echo "bl $(cat ~/.ssh/muse-health.pub)" \
# | ssh super@34.139.37.135 "sudo tee -a /srv/board/health_signers"
set -u
NETVM_DIR="${NETVM_DIR:-$HOME/Projects/NetVM}"
ACCOUNTS="$NETVM_DIR/ACCOUNTS.md"
CHECKER="$NETVM_DIR/bin/accounts-health.py"
MACHINE="${MUSE_MACHINE:-bl}"
KEY="${HEALTH_KEY:-$HOME/.ssh/muse-health}"
ENDPOINT="${HEALTH_ENDPOINT:-https://board.muse-dev.online/api/health/report}"
POST=1
[ "${1:-}" = "--no-post" ] && POST=0
[ -f "$ACCOUNTS" ] || { echo "accounts-health: $ACCOUNTS missing" >&2; exit 1; }
[ -f "$CHECKER" ] || { echo "accounts-health: $CHECKER missing" >&2; exit 1; }
TS="$(date +%s)"
TMP="$(mktemp -d)"
trap 'rm -rf "$TMP"' EXIT
# Parse ACCOUNTS.md pipe table, probe each account inside its netns
NETVM_DIR="$NETVM_DIR" python3 - > "$TMP/facts.json" <<'PYEOF'
import json, os, subprocess, time
netvm = os.environ["NETVM_DIR"]
rows = []
for line in open(os.path.join(netvm, "ACCOUNTS.md")):
line = line.strip()
if not line.startswith("|"):
continue
cells = [c.strip() for c in line.strip("|").split("|")]
if len(cells) < 13 or cells[0] in ("agent", "-------", ""):
continue
if cells[10] not in ("", "-"):
rows.append({"agent": cells[0], "node": cells[1],
"status": cells[8], "cdp_port": cells[10]})
checker = os.path.join(netvm, "bin", "accounts-health.py")
out = {}
for a in rows:
try:
r = subprocess.run(
["sudo", "-n", "ip", "netns", "exec", f"warp-{a['node']}",
"python3", checker, a["cdp_port"]],
capture_output=True, text=True, timeout=60)
res = json.loads(r.stdout.strip().splitlines()[-1])
res["registry_status"] = a["status"]
out[a["agent"]] = res
except Exception as e:
out[a["agent"]] = {"browser_up": False, "session_alive": None,
"detail": f"probe failed: {e}",
"registry_status": a["status"]}
print(json.dumps({"accounts": out, "checked_at": int(time.time())}, indent=2))
PYEOF
if [ "$POST" -eq 0 ]; then
cat "$TMP/facts.json"
exit 0
fi
if [ ! -f "$KEY" ]; then
echo "accounts-health: $KEY missing — printing JSON, not posting (see header for key setup)" >&2
cat "$TMP/facts.json"
exit 0
fi
printf '%s\n%s\n' "$MACHINE" "$TS" > "$TMP/payload"
FACTS_JSON="$(cat "$TMP/facts.json")"
printf '%s' "$FACTS_JSON" >> "$TMP/payload"
ssh-keygen -Y sign -f "$KEY" -n health "$TMP/payload" >/dev/null 2>&1
SIG="$(cat "$TMP/payload.sig")"
python3 - "$MACHINE" "$TS" "$FACTS_JSON" "$SIG" <<'PYEOF' > "$TMP/body.json"
import json, sys
machine, ts, facts_json, sig = sys.argv[1], int(sys.argv[2]), sys.argv[3], sys.argv[4]
body = {"machine": machine, "ts": ts,
"facts": json.loads(facts_json), "facts_json": facts_json,
"signature": sig}
print(json.dumps(body))
PYEOF
curl -s -X POST "$ENDPOINT" -H 'Content-Type: application/json' \
--data @"$TMP/body.json" | head -c 300
echo
-71
View File
@@ -1,71 +0,0 @@
#!/usr/bin/env python3
"""
dev-dm.py: Developer Direct Messages
For operator-to-developer communication (muse, pip task agents).
Developers are browser-only, no SSH. Uses headless via dm.py backend.
Distinct from opp-dm.py (for operators: 646, operator-main).
Usage:
dev-dm.py send --agent muse --target main "message"
dev-dm.py read --agent muse --target main
"""
import argparse
import subprocess
import sys
DM_PY = "/home/super/Projects/NetVM/bin/dm.py"
def run(cmd, timeout=90):
result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout)
return result.stdout.strip()
def dev_send(agent, target, message):
"""Send a dev DM via headless (dm.py backend)."""
if agent not in ["muse", "pip"]:
print(f"ERROR: dev-dm only for muse, pip (not {agent})", file=sys.stderr)
sys.exit(1)
# Delegate to dm.py
safe = message.replace('"', '\\"')[:1000]
cmd = f"{DM_PY} send --agent {agent} --target {target} \"{safe}\""
# Run on bl
bl_cmd = f"ssh -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -o BatchMode=yes super@100.123.153.75 '{cmd}'"
vm_cmd = f"ssh -i ~/.ssh/vm_to_gcp -o ProxyCommand=\"$HOME/workspace/bin/ssh-via-proxy %h %p\" -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -o BatchMode=yes super@34.139.37.135 \"{bl_cmd}\""
result = run(vm_cmd)
print(f"DEV-DM sent to {agent}/{target}")
return result
def dev_read(agent, target, n=5):
"""Read dev DMs via headless."""
if agent not in ["muse", "pip"]:
print(f"ERROR: dev-dm only for muse, pip", file=sys.stderr)
sys.exit(1)
cmd = f"{DM_PY} read --agent {agent} --target {target} --n {n}"
bl_cmd = f"ssh -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -o BatchMode=yes super@100.123.153.75 '{cmd}'"
vm_cmd = f"ssh -i ~/.ssh/vm_to_gcp -o ProxyCommand=\"$HOME/workspace/bin/ssh-via-proxy %h %p\" -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -o BatchMode=yes super@34.139.37.135 \"{bl_cmd}\""
result = run(vm_cmd)
print(result)
return result
def main():
p = argparse.ArgumentParser(description="DEV-DM: Developer Direct Messages")
sub = p.add_subparsers(dest='cmd', required=True)
ps = sub.add_parser('send', help='Send dev DM')
ps.add_argument('--agent', required=True, choices=['muse', 'pip'])
ps.add_argument('--target', required=True, help='main or chat_id')
ps.add_argument('message', help='Message')
ps.set_defaults(func=lambda a: dev_send(a.agent, a.target, a.message))
pr = sub.add_parser('read', help='Read dev DMs')
pr.add_argument('--agent', required=True, choices=['muse', 'pip'])
pr.add_argument('--target', required=True)
pr.add_argument('--n', type=int, default=5)
pr.set_defaults(func=lambda a: dev_read(a.agent, a.target, a.n))
args = p.parse_args()
args.func(args)
if __name__ == '__main__':
main()
+205
View File
@@ -0,0 +1,205 @@
#!/usr/bin/env python3
"""Meta Accounts Center change-detection harness.
Captures a structural snapshot of the accountscenter.meta.com auth flow
via CDP inside a NetVM netns, diffs against the stored baseline.
Outcomes:
PASS - matches baseline (or first run establishes it)
CHANGED - structural diff detected; needs human review, baseline untouched
FAIL - automation itself broke (browser/CDP/network error)
Usage:
meta-ac-snapshot.py [--node NAME] [--promote] [--snapshot-dir DIR]
--node NetVM node to run in (default: phone)
--promote after human review, promote the latest snapshot to baseline
--snapshot-dir where snapshots live (default: ~/Projects/NetVM/snapshots/meta-ac)
"""
import argparse, base64, datetime, json, os, subprocess, sys, time
import urllib.parse, urllib.request
CDP_PORT = 19744
def log(*a):
print(*a, flush=True)
def ns_exec(node, cmd):
return subprocess.run(
["sudo", "-n", "ip", "netns", "exec", f"warp-{node}"] + cmd,
capture_output=True, text=True)
def norm_url(u):
"""Strip query/fragment — nonces change every visit."""
p = urllib.parse.urlparse(u)
return f"{p.scheme}://{p.host}{p.path}" if hasattr(p, 'host') else f"{p.scheme}://{p.hostname}{p.path}"
def main():
ap = argparse.ArgumentParser()
ap.add_argument("--node", default="phone")
ap.add_argument("--promote", action="store_true")
ap.add_argument("--snapshot-dir", default=os.path.expanduser(
"~/Projects/NetVM/snapshots/meta-ac"))
args = ap.parse_args()
os.makedirs(args.snapshot_dir, exist_ok=True)
baseline_path = os.path.join(args.snapshot_dir, "baseline.json")
if args.promote:
snaps = sorted(f for f in os.listdir(args.snapshot_dir)
if f.startswith("snap-") and f.endswith(".json"))
if not snaps:
log("no snapshots to promote"); return 2
latest = os.path.join(args.snapshot_dir, snaps[-1])
data = json.load(open(latest))
data["promoted_at"] = datetime.datetime.now(datetime.timezone.utc).isoformat()
json.dump(data, open(baseline_path, "w"), indent=2)
log(f"promoted {snaps[-1]} -> baseline.json")
return 0
profile_dir = "/tmp/meta-ac-snap-profile"
subprocess.run(["rm", "-rf", profile_dir])
os.makedirs(profile_dir, exist_ok=True)
def http(path):
with urllib.request.urlopen(
f"http://127.0.0.1:{CDP_PORT}{path}", timeout=5) as r:
return json.loads(r.read())
# launch chromium inside the netns via a wrapper script
wrapper = "/tmp/meta-ac-snap-run.py"
open(wrapper, "w").write(WRAPPER_SRC)
log(f"launching chromium in warp-{args.node} (CDP {CDP_PORT})...")
proc = subprocess.Popen(
["sudo", "-n", "ip", "netns", "exec", f"warp-{args.node}",
"python3", wrapper, str(CDP_PORT), profile_dir],
stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True)
try:
out, _ = proc.communicate(timeout=120)
except subprocess.TimeoutExpired:
proc.kill(); log("FAIL: harness timed out"); return 1
print(out)
# wrapper prints SNAPSHOT_JSON=<json> on success
snap = None
for line in out.splitlines():
if line.startswith("SNAPSHOT_JSON="):
snap = json.loads(line[len("SNAPSHOT_JSON="):])
if not snap:
log("FAIL: no snapshot captured"); return 1
snap["node"] = args.node
snap["captured_at"] = datetime.datetime.now(datetime.timezone.utc).isoformat()
# egress ip for context
try:
r = ns_exec(args.node, ["curl", "-s", "--max-time", "8",
"https://api.ipify.org"])
snap["egress_ip"] = r.stdout.strip()
except Exception:
snap["egress_ip"] = "unknown"
ts = datetime.datetime.now(datetime.timezone.utc).strftime("%Y%m%d-%H%M%S")
snap_path = os.path.join(args.snapshot_dir, f"snap-{ts}.json")
json.dump(snap, open(snap_path, "w"), indent=2)
log(f"snapshot saved: {snap_path}")
if not os.path.exists(baseline_path):
json.dump(snap, open(baseline_path, "w"), indent=2)
log("PASS: baseline established (first run)")
return 0
baseline = json.load(open(baseline_path))
diffs = diff_snapshots(baseline, snap)
if not diffs:
log("PASS: matches baseline")
return 0
log("CHANGED: structural diff detected (baseline untouched):")
for d in diffs:
log(f" - {d}")
log("review with: diff baseline.json snap-<ts>.json")
log("promote after review with: --promote")
return 3
def diff_snapshots(base, snap):
diffs = []
b_chain = [norm_url(u) for u in base.get("redirect_chain", [])]
s_chain = [norm_url(u) for u in snap.get("redirect_chain", [])]
if b_chain != s_chain:
diffs.append(f"redirect_chain changed: {b_chain} -> {s_chain}")
for key in ("forms", "inputs", "buttons"):
b = sorted(base.get("dom_markers", {}).get(key, []))
s = sorted(snap.get("dom_markers", {}).get(key, []))
if b != s:
added = [x for x in s if x not in b]
removed = [x for x in b if x not in s]
diffs.append(f"dom_markers.{key}: added={added} removed={removed}")
if base.get("final_title") != snap.get("final_title"):
diffs.append(f"final_title: {base.get('final_title')!r} -> {snap.get('final_title')!r}")
return diffs
WRAPPER_SRC = '''
import json, subprocess, sys, time, os, urllib.request, base64
CDP_PORT = int(sys.argv[1])
PROFILE_DIR = sys.argv[2]
def http(path):
with urllib.request.urlopen(f"http://127.0.0.1:{CDP_PORT}{path}", timeout=5) as r:
return json.loads(r.read())
logf = open("/tmp/meta-ac-snap-chrome.log", "w")
proc = subprocess.Popen(["chromium", "--headless=new", "--disable-gpu", "--no-sandbox",
"--disable-dev-shm-usage", f"--user-data-dir={PROFILE_DIR}",
f"--remote-debugging-port={CDP_PORT}", "--remote-allow-origins=*", "about:blank"],
stdout=logf, stderr=subprocess.STDOUT)
try:
for i in range(30):
try:
ver = http("/json/version")
if "webSocketDebuggerUrl" in ver: break
except Exception: pass
time.sleep(1)
else:
print("FAIL: CDP never came up"); sys.exit(1)
import websocket
bws = websocket.create_connection(ver["webSocketDebuggerUrl"], timeout=20)
bws.send(json.dumps({"id": 1, "method": "Target.createTarget",
"params": {"url": "https://accountscenter.meta.com"}}))
target_id = json.loads(bws.recv())["result"]["targetId"]
bws.close()
# redirect chain: seed with the navigation target (we always start
# there), then poll for where Meta sends us. Seeding fixes the race
# where a fast redirect is missed by the poll interval.
START_URL = "https://accountscenter.meta.com/"
chain, seen = [START_URL], {START_URL}
for _ in range(24):
time.sleep(2)
for t in http("/json/list"):
if t.get("id") == target_id or "meta.com" in t.get("url", ""):
u = t["url"]
if u not in seen:
seen.add(u); chain.append(u)
title = t.get("title", "")
break
# dom markers from the final tab
tab_ws = None
for t in http("/json/list"):
if t.get("id") == target_id or "meta.com" in t.get("url", ""):
tab_ws = t["webSocketDebuggerUrl"]; final_url = t["url"]; break
ws = websocket.create_connection(tab_ws, timeout=20)
js = """JSON.stringify({
forms: [...document.forms].map(f => f.id || f.name || '(anon)'),
inputs: [...document.querySelectorAll('input')].map(i => i.name || i.type || '(anon)'),
buttons: [...document.querySelectorAll('button, [role=button]')].map(b => (b.innerText||'').trim()).filter(Boolean)
})"""
ws.send(json.dumps({"id": 1, "method": "Runtime.evaluate",
"params": {"expression": js, "returnByValue": True}}))
markers = json.loads(json.loads(ws.recv())["result"]["result"]["value"])
# dedupe buttons, keep order
markers["buttons"] = list(dict.fromkeys(markers["buttons"]))
ws.close()
snap = {"redirect_chain": chain, "final_url": final_url,
"final_title": title, "dom_markers": markers}
print("SNAPSHOT_JSON=" + json.dumps(snap))
finally:
proc.terminate()
'''
if __name__ == "__main__":
sys.exit(main())
+52
View File
@@ -0,0 +1,52 @@
#!/usr/bin/env python3
"""
NetVM Docs HTTP Server (bl)
Serves markdown documentation for agents via HTTPS (tailnet).
- Central repo: ~/Projects/NetVM/
- Serves: *.md, docs/*.png
- Read-only, no auth (tailnet is the auth boundary)
GOLDEN PATH: container -> VM (34.139.37.135) -> bl (100.123.153.75) -> this server
Usage:
python3 netvm-docs-server.py [--port 8080]
Agents fetch via:
curl http://0.0.0.0:8080/ACCOUNTS.md
"""
import http.server
import socketserver
import os
import argparse
from pathlib import Path
class DocsHandler(http.server.SimpleHTTPRequestHandler):
def __init__(self, *args, **kwargs):
# Serve from NetVM repo root
self.base_dir = Path.home() / "Projects" / "NetVM"
super().__init__(*args, directory=str(self.base_dir), **kwargs)
def log_message(self, format, *args):
# Quiet logging
pass
def end_headers(self):
# Allow CORS for browser agents
self.send_header('Access-Control-Allow-Origin', '*')
super().end_headers()
def main():
p = argparse.ArgumentParser()
p.add_argument('--port', type=int, default=8080)
args = p.parse_args()
# Only bind to tailnet interface (not public)
# 100.123.153.75 is bl's tailnet IP
with socketserver.TCPServer(("0.0.0.0", args.port), DocsHandler) as httpd:
print(f"Serving NetVM docs on http://0.0.0.0:{args.port}/")
print(f"Directory: {Path.home()}/Projects/NetVM/")
httpd.serve_forever()
if __name__ == '__main__':
main()
+86
View File
@@ -0,0 +1,86 @@
#!/usr/bin/env python3
"""onboard-driver.py — bl-side OTP onboarding driver.
Runs INSIDE the node's netns (via netvm-exec.sh). Reads the identifier
(line 1) and OTP code (line 2, submit step only) from stdin — never argv.
onboard-driver.py --node muse --service muse --id-type email --step initiate [--dry-run]
onboard-driver.py --node muse --service muse --id-type email --step submit
Exit codes: 0 = step done, 2 = APPROVAL_NEEDED (code sent, awaiting OTP),
1 = failed. The identifier/code are passed to the local signin script as
argv (transient, same trust domain — bl is operator infrastructure);
they never cross a network boundary except inside the already-encrypted
VM->bl SSH stdin pipe.
Part of the cred onboarding module (front-door repo, docs/CRED-MODULE.md).
"""
import argparse
import json
import subprocess
import sys
import urllib.request
SIGNIN = "/home/super/Projects/NetVM/bin/muse-signin.py"
CDP_PORTS = {"muse": "9410", "pip": "9420"}
def cdp_ok(port):
try:
ts = json.load(urllib.request.urlopen(
"http://127.0.0.1:%s/json/list" % port, timeout=5))
return any(t.get("type") == "page" for t in ts)
except Exception:
return False
def main():
p = argparse.ArgumentParser()
p.add_argument("--node", required=True)
p.add_argument("--service", required=True)
p.add_argument("--id-type", required=True)
p.add_argument("--step", required=True, choices=["initiate", "submit"])
p.add_argument("--dry-run", action="store_true")
args = p.parse_args()
lines = sys.stdin.read().splitlines()
identifier = lines[0].strip() if lines else ""
code = lines[1].strip() if len(lines) > 1 else ""
if args.service != "muse" or args.id_type != "email":
print("ERROR: unsupported service/id_type "
"(muse+email only for now)", file=sys.stderr)
return 1
port = CDP_PORTS.get(args.node)
if not port:
print("ERROR: unknown node", file=sys.stderr)
return 1
if args.dry_run:
# Walk the chain without sending anything: netns + CDP + page.
if cdp_ok(port):
print("dry-run ok: node=%s cdp=%s reachable, page present"
% (args.node, port))
return 0
print("ERROR: CDP unreachable on %s" % port, file=sys.stderr)
return 1
if not identifier:
print("ERROR: no identifier on stdin", file=sys.stderr)
return 1
cmd = [sys.executable, SIGNIN, "--email", identifier]
if args.step == "submit":
if not code:
print("ERROR: no code on stdin", file=sys.stderr)
return 1
cmd += ["--otp", code]
r = subprocess.run(cmd, capture_output=True, text=True, timeout=220)
# Propagate the signin script's contract: 2 = OTP prompt reached.
sys.stdout.write(r.stdout)
sys.stderr.write(r.stderr)
return r.returncode
if __name__ == "__main__":
sys.exit(main())
-102
View File
@@ -1,102 +0,0 @@
#!/usr/bin/env python3
"""
opp-dm.py: Operator Direct Messages (fixed)
Works from container OR from bl directly.
Detects environment and acts accordingly.
Verifies writes.
For operator-to-operator: 646, operator-main.
"""
import argparse
import subprocess
import sys
import os
import datetime
from pathlib import Path
# Detect if we're on bl (has /home/super/Projects/NetVM)
ON_BL = Path("/home/super/Projects/NetVM").exists()
BL_MSG_DIR = Path("/home/super/Projects/NetVM/opp-dms")
def bl_run_container(cmd):
"""Run on bl via SSH chain (from container)."""
full = f"ssh -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -o BatchMode=yes super@100.123.153.75 '{cmd}'"
vm_cmd = f"ssh -i ~/.ssh/vm_to_gcp -o ProxyCommand=\"$HOME/workspace/bin/ssh-via-proxy %h %p\" -o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null -o BatchMode=yes super@34.139.37.135 \"{full}\""
result = subprocess.run(vm_cmd, shell=True, capture_output=True, text=True, timeout=30)
return result.stdout.strip(), result.returncode
def opp_send(to, message, from_who="operator-main"):
"""Send operator DM. Works from container or bl."""
if to not in ["646", "operator-main"]:
print(f"ERROR: Unknown operator {to}", file=sys.stderr)
sys.exit(1)
ts = datetime.datetime.now().isoformat()
safe = message.replace("'", "'\"'\"'")[:1000]
line = f"{ts} [{from_who}]: {safe}"
if ON_BL:
# Direct write (we're on bl)
BL_MSG_DIR.mkdir(parents=True, exist_ok=True)
logfile = BL_MSG_DIR / f"{to}.log"
with open(logfile, "a") as f:
f.write(line + "\n")
# Verify
if logfile.exists():
print(f"OPP-DM sent to {to} (verified)")
return True
else:
print(f"ERROR: Write failed", file=sys.stderr)
sys.exit(1)
else:
# Via SSH (from container)
cmd = f"mkdir -p {BL_MSG_DIR} && echo '{line}' >> {BL_MSG_DIR}/{to}.log && test -f {BL_MSG_DIR}/{to}.log && echo OK"
out, rc = bl_run_container(cmd)
if rc == 0 and "OK" in out:
print(f"OPP-DM sent to {to} (verified)")
return True
else:
print(f"ERROR: Send failed: {out}", file=sys.stderr)
sys.exit(1)
def opp_read(who="operator-main", n=5):
"""Read operator DMs. Works from container or bl."""
if who not in ["646", "operator-main"]:
print(f"ERROR: Unknown operator {who}", file=sys.stderr)
sys.exit(1)
if ON_BL:
logfile = BL_MSG_DIR / f"{who}.log"
if not logfile.exists():
print("no messages")
return
with open(logfile) as f:
lines = f.readlines()
for l in lines[-n:]:
print(l.strip())
else:
cmd = f"tail -n {n} {BL_MSG_DIR}/{who}.log 2>/dev/null || echo 'no messages'"
out, rc = bl_run_container(cmd)
print(out)
def main():
p = argparse.ArgumentParser(description="OPP-DM: Operator DMs (fixed)")
sub = p.add_subparsers(dest='cmd', required=True)
ps = sub.add_parser('send', help='Send operator DM')
ps.add_argument('--to', required=True, choices=['646', 'operator-main'])
ps.add_argument('--from', dest='from_who', default='operator-main')
ps.add_argument('message', help='Message')
ps.set_defaults(func=lambda a: opp_send(a.to, a.message, a.from_who))
pr = sub.add_parser('read', help='Read operator DMs')
pr.add_argument('--who', default='operator-main', choices=['646', 'operator-main'])
pr.add_argument('--n', type=int, default=5)
pr.set_defaults(func=lambda a: opp_read(a.who, a.n))
args = p.parse_args()
args.func(args)
if __name__ == '__main__':
main()