diff --git a/bin/dm.py b/bin/dm.py index d6193c8..66f5a64 100755 --- a/bin/dm.py +++ b/bin/dm.py @@ -27,6 +27,7 @@ Usage: dm.py thread --from opm --to 646 --target main "message" """ import argparse +import os import re import subprocess import time @@ -47,8 +48,12 @@ def log_event(event): with open(LOG_FILE, "a") as f: f.write(json.dumps(event) + "\n") -def run(cmd, timeout=60): - result = subprocess.run(cmd, shell=True, capture_output=True, text=True, timeout=timeout) +def run(cmd, timeout=60, priority=None): + # 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() 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 delivered = False 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 # Verify by reading recipient's chat (independent check, not local echo) diff --git a/bin/muse-chat-api.py b/bin/muse-chat-api.py index 9cb0947..750b7c9 100755 --- a/bin/muse-chat-api.py +++ b/bin/muse-chat-api.py @@ -25,6 +25,11 @@ try: HAS_SIDECHAT_MANAGER = True except ImportError: 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(): @@ -493,52 +498,65 @@ def main(): node, cdp_url = ACCOUNTS[args.account] page = get_page(node, cdp_url) - ws = websocket.create_connection(page['webSocketDebuggerUrl'], timeout=15) - try: - if args.command == 'send': - if not args.arg: - print("ERROR: send requires a message", file=sys.stderr) - sys.exit(1) - cmd_send(ws, " ".join(args.arg)) - elif args.command == 'messages': - n = int(args.arg[0]) if args.arg else 5 - w = int(args.arg[1]) if len(args.arg) > 1 else 200 - cmd_messages(ws, n, w) - elif args.command == 'wait': - t = int(args.arg[0]) if args.arg else 30 - cmd_wait(ws, t) - elif args.command == 'approvals': - cmd_approvals(ws) - elif args.command == 'sidechat': - if not args.arg: - print("ERROR: sidechat requires subcommand (use|main|list|create)", 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) + # CDP operation queue: serialize browser access across processes. + # Priority from CDP_PRIORITY env (high/normal/low); default normal. + # Falls back to unqueued if cdp_queue is unavailable. + import contextlib + if HAS_CDP_QUEUE: + _prio_name = os.environ.get("CDP_PRIORITY", "normal").lower() + _prio = {"high": PRIORITY_HIGH, "low": PRIORITY_LOW}.get( + _prio_name, PRIORITY_NORMAL) + _slot = cdp_slot(node, priority=_prio) + else: + _slot = contextlib.nullcontext() + with _slot: + ws = websocket.create_connection(page['webSocketDebuggerUrl'], timeout=15) + + try: + if args.command == 'send': + if not args.arg: + print("ERROR: send requires a message", 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() + cmd_send(ws, " ".join(args.arg)) + elif args.command == 'messages': + n = int(args.arg[0]) if args.arg else 5 + w = int(args.arg[1]) if len(args.arg) > 1 else 200 + cmd_messages(ws, n, w) + elif args.command == 'wait': + t = int(args.arg[0]) if args.arg else 30 + cmd_wait(ws, t) + elif args.command == 'approvals': + cmd_approvals(ws) + elif args.command == 'sidechat': + if not args.arg: + print("ERROR: sidechat requires subcommand (use|main|list|create)", 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) + 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__': main()