Integrate cdp_queue into CDP path

All CDP sessions now go through per-node queue (max 2 concurrent, priority levels). DM sends use high priority. Graceful fallback if module unavailable.

Session: sidechat/chromebox-ops
This commit is contained in:
operator-main
2026-10-04 13:15:59 +00:00
parent 0d25696860
commit 5f0a77d04b
2 changed files with 71 additions and 47 deletions
+9 -3
View File
@@ -27,6 +27,7 @@ Usage:
dm.py thread --from opm --to 646 --target main "message" dm.py thread --from opm --to 646 --target main "message"
""" """
import argparse import argparse
import os
import re import re
import subprocess import subprocess
import time import time
@@ -47,8 +48,12 @@ def log_event(event):
with open(LOG_FILE, "a") as f: with open(LOG_FILE, "a") as f:
f.write(json.dumps(event) + "\n") f.write(json.dumps(event) + "\n")
def run(cmd, timeout=60): def run(cmd, timeout=60, priority=None):
result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout) # priority: CDP queue priority for the subprocess (high/normal/low).
# DM sends are user-facing -> high. Reads stay at default (normal).
env = dict(os.environ, CDP_PRIORITY=priority) if priority else None
result = subprocess.run(cmd, shell=True, capture_output=True, text=True,
timeout=timeout, env=env)
return result.stdout.strip() return result.stdout.strip()
def run_full(cmd, timeout=60): def run_full(cmd, timeout=60):
@@ -105,7 +110,8 @@ def dm_send(agent, target, message, verify=True, raw=False, to_agent=None):
max_retries = 3 max_retries = 3
delivered = False delivered = False
for attempt in range(max_retries): for attempt in range(max_retries):
run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} send "{safe}"') run(f'{NETVM_EXEC} {recipient} -- python3 {API} --account {recipient} send "{safe}"',
priority="high")
time.sleep(3) # Wait for message to propagate time.sleep(3) # Wait for message to propagate
# Verify by reading recipient's chat (independent check, not local echo) # Verify by reading recipient's chat (independent check, not local echo)
+62 -44
View File
@@ -25,6 +25,11 @@ try:
HAS_SIDECHAT_MANAGER = True HAS_SIDECHAT_MANAGER = True
except ImportError: except ImportError:
HAS_SIDECHAT_MANAGER = False HAS_SIDECHAT_MANAGER = False
try:
from cdp_queue import cdp_slot, PRIORITY_HIGH, PRIORITY_NORMAL, PRIORITY_LOW
HAS_CDP_QUEUE = True
except ImportError:
HAS_CDP_QUEUE = False
def _load_accounts(): def _load_accounts():
@@ -493,52 +498,65 @@ def main():
node, cdp_url = ACCOUNTS[args.account] node, cdp_url = ACCOUNTS[args.account]
page = get_page(node, cdp_url) page = get_page(node, cdp_url)
ws = websocket.create_connection(page['webSocketDebuggerUrl'], timeout=15)
try: # CDP operation queue: serialize browser access across processes.
if args.command == 'send': # Priority from CDP_PRIORITY env (high/normal/low); default normal.
if not args.arg: # Falls back to unqueued if cdp_queue is unavailable.
print("ERROR: send requires a message", file=sys.stderr) import contextlib
sys.exit(1) if HAS_CDP_QUEUE:
cmd_send(ws, " ".join(args.arg)) _prio_name = os.environ.get("CDP_PRIORITY", "normal").lower()
elif args.command == 'messages': _prio = {"high": PRIORITY_HIGH, "low": PRIORITY_LOW}.get(
n = int(args.arg[0]) if args.arg else 5 _prio_name, PRIORITY_NORMAL)
w = int(args.arg[1]) if len(args.arg) > 1 else 200 _slot = cdp_slot(node, priority=_prio)
cmd_messages(ws, n, w) else:
elif args.command == 'wait': _slot = contextlib.nullcontext()
t = int(args.arg[0]) if args.arg else 30 with _slot:
cmd_wait(ws, t) ws = websocket.create_connection(page['webSocketDebuggerUrl'], timeout=15)
elif args.command == 'approvals':
cmd_approvals(ws) try:
elif args.command == 'sidechat': if args.command == 'send':
if not args.arg: if not args.arg:
print("ERROR: sidechat requires subcommand (use|main|list|create)", file=sys.stderr) print("ERROR: send requires a message", file=sys.stderr)
sys.exit(1)
sub = args.arg[0]
if sub == "create":
cmd_sidechat_create(ws)
elif sub == "use":
if len(args.arg) < 2:
print("ERROR: sidechat use requires chat_id", file=sys.stderr)
sys.exit(1) sys.exit(1)
cmd_sidechat_use(ws, args.arg[1]) cmd_send(ws, " ".join(args.arg))
elif sub == "main": elif args.command == 'messages':
cmd_sidechat_main(ws) n = int(args.arg[0]) if args.arg else 5
elif sub == "list": w = int(args.arg[1]) if len(args.arg) > 1 else 200
cmd_sidechat_list(ws) cmd_messages(ws, n, w)
else: elif args.command == 'wait':
print(f"ERROR: unknown sidechat subcommand: {sub}", file=sys.stderr) t = int(args.arg[0]) if args.arg else 30
sys.exit(1) cmd_wait(ws, t)
elif args.command == 'url': elif args.command == 'approvals':
cmd_url(ws) cmd_approvals(ws)
elif args.command == 'upload': elif args.command == 'sidechat':
if not args.arg: if not args.arg:
print("ERROR: upload requires a filepath", file=sys.stderr) print("ERROR: sidechat requires subcommand (use|main|list|create)", file=sys.stderr)
sys.exit(1) sys.exit(1)
cmd_upload(ws, args.arg[0], message=args.message, sub = args.arg[0]
dry_run=args.dry_run) if sub == "create":
finally: cmd_sidechat_create(ws)
ws.close() elif sub == "use":
if len(args.arg) < 2:
print("ERROR: sidechat use requires chat_id", file=sys.stderr)
sys.exit(1)
cmd_sidechat_use(ws, args.arg[1])
elif sub == "main":
cmd_sidechat_main(ws)
elif sub == "list":
cmd_sidechat_list(ws)
else:
print(f"ERROR: unknown sidechat subcommand: {sub}", file=sys.stderr)
sys.exit(1)
elif args.command == 'url':
cmd_url(ws)
elif args.command == 'upload':
if not args.arg:
print("ERROR: upload requires a filepath", file=sys.stderr)
sys.exit(1)
cmd_upload(ws, args.arg[0], message=args.message,
dry_run=args.dry_run)
finally:
ws.close()
if __name__ == '__main__': if __name__ == '__main__':
main() main()