diff --git a/bin/cdp_queue.py b/bin/cdp_queue.py new file mode 100644 index 0000000..eda98be --- /dev/null +++ b/bin/cdp_queue.py @@ -0,0 +1,230 @@ +#!/usr/bin/env python3 +""" +Per-browser CDP operation queue for the NetVM fleet. + +Problem: nothing coordinates browser operations. DM sends (dm.py), tab +operations, agent reads (muse-chat-api.py), and watchdog restarts all hit +the same Chromium with zero scheduling. Functional tests pass under load, +but there is no backpressure — under real fleet concurrency this will +overwhelm the machine. + +Solution: per-node FIFO queue with priority levels and a cap on concurrent +CDP operations per browser. + +Cross-process design: dm.py shells out to muse-chat-api.py via subprocess, +so in-process locks (threading.Semaphore) alone cannot coordinate. This +module uses flock'd slot files + ticket files in /tmp, which work across +processes AND threads. + +Usage: + from cdp_queue import cdp_slot, PRIORITY_HIGH + + with cdp_slot("opm", priority=PRIORITY_HIGH): + ... do CDP work ... + +Explicit acquire/release: + from cdp_queue import acquire, release, QueueTimeout + token = acquire("opm", priority=PRIORITY_HIGH, timeout=60) + try: + ... + finally: + release(token) + +Priority levels (lower number = higher priority): + PRIORITY_HIGH = 0 # DM sends (user-facing) + PRIORITY_NORMAL = 1 # tab opens, reads + PRIORITY_LOW = 2 # background scans + +Fairness: tickets are ordered by (priority, arrival time). A waiter only +proceeds when its ticket is first in line AND a slot is free. No starvation: +a low-priority ticket eventually becomes the oldest and gets served. +""" + +import fcntl +import logging +import os +import time +import uuid + +log = logging.getLogger("cdp_queue") + +# ---- Tunables ---- +PRIORITY_HIGH = 0 +PRIORITY_NORMAL = 1 +PRIORITY_LOW = 2 + +MAX_CONCURRENT = 2 # max simultaneous CDP ops per browser +ACQUIRE_TIMEOUT = 60.0 # fail loud instead of hanging forever +WARN_AFTER = 10.0 # log a warning when a waiter waits this long +POLL_INTERVAL = 0.05 # ticket/slot poll cadence + +QUEUE_DIR = "/tmp/cdp-queue" +VALID_NODES = ("muse", "pip", "646", "opm") + + +class QueueTimeout(Exception): + """Raised when a slot cannot be acquired within the timeout.""" + pass + + +def _node_dir(node): + return os.path.join(QUEUE_DIR, node) + + +def _tickets_dir(node): + return os.path.join(_node_dir(node), "tickets") + + +def _slot_path(node, i): + return os.path.join(_node_dir(node), "slot-%d.lock" % i) + + +def _ensure_dirs(node): + os.makedirs(_tickets_dir(node), exist_ok=True) + # Pre-create slot files so flock targets always exist + for i in range(MAX_CONCURRENT): + p = _slot_path(node, i) + if not os.path.exists(p): + open(p, "a").close() + + +def _read_tickets(node): + """Return sorted list of (priority, timestamp, ticket_name).""" + tdir = _tickets_dir(node) + out = [] + try: + for name in os.listdir(tdir): + if not name.endswith(".ticket"): + continue + try: + # ticket name: "--.ticket" + parts = name[:-7].split("-") + prio_s, ts_s = parts[0], parts[1] + out.append((int(prio_s), float(ts_s), name)) + except (ValueError, IndexError): + continue + except FileNotFoundError: + pass + out.sort() + return out + + +class _Slot: + """A held queue slot. Release via .release() or context manager.""" + + def __init__(self, node, ticket_name, fh, waited): + self.node = node + self.ticket_name = ticket_name + self.fh = fh + self.waited = waited + self._released = False + + def release(self): + if self._released: + return + self._released = True + try: + fcntl.flock(self.fh, fcntl.LOCK_UN) + self.fh.close() + except Exception: + pass + # Remove our ticket (best effort — a stale ticket is harmless; + # the next waiter re-reads the directory each poll) + try: + os.unlink(os.path.join(_tickets_dir(self.node), self.ticket_name)) + except Exception: + pass + log.debug("cdp_queue: released slot for node=%s (waited %.1fs)", + self.node, self.waited) + + def __enter__(self): + return self + + def __exit__(self, *exc): + self.release() + + +def acquire(node, priority=PRIORITY_NORMAL, timeout=ACQUIRE_TIMEOUT): + """ + Block until a CDP slot is free for `node`, then return a _Slot. + Raises QueueTimeout after `timeout` seconds. Raises ValueError for + unknown nodes. + """ + if node not in VALID_NODES: + raise ValueError("unknown node: %r (valid: %s)" % (node, VALID_NODES)) + if priority not in (PRIORITY_HIGH, PRIORITY_NORMAL, PRIORITY_LOW): + raise ValueError("invalid priority: %r" % (priority,)) + + _ensure_dirs(node) + tdir = _tickets_dir(node) + + # Our ticket: "--.ticket", sorted by (prio, ts) + ticket = "%d-%f-%s.ticket" % (priority, time.time(), uuid.uuid4().hex[:8]) + open(os.path.join(tdir, ticket), "w").close() + + start = time.time() + warned = False + try: + while True: + elapsed = time.time() - start + if elapsed >= timeout: + raise QueueTimeout( + "node=%s: no CDP slot free after %.0fs (priority=%d)" % + (node, timeout, priority)) + if elapsed >= WARN_AFTER and not warned: + warned = True + depth = len(_read_tickets(node)) + log.warning("cdp_queue: node=%s waiting %.0fs for slot " + "(queue depth %d, priority %d)", + node, elapsed, depth, priority) + + tickets = _read_tickets(node) + # Am I among the first MAX_CONCURRENT in line? + # (priority, then arrival time). The first N tickets are all + # eligible to grab slots; they distribute via non-blocking flock. + my_pos = next((i for i, (_, _, name) in enumerate(tickets) + if name == ticket), None) + if my_pos is not None and my_pos < MAX_CONCURRENT: + # Try each slot file non-blocking + for i in range(MAX_CONCURRENT): + fh = open(_slot_path(node, i), "w") + try: + fcntl.flock(fh, fcntl.LOCK_EX | fcntl.LOCK_NB) + except (BlockingIOError, OSError): + fh.close() + continue + # Got it + waited = time.time() - start + if waited > 1.0: + log.debug("cdp_queue: node=%s acquired slot after " + "%.1fs (priority %d)", node, waited, priority) + return _Slot(node, ticket, fh, waited) + time.sleep(POLL_INTERVAL) + except BaseException: + # On timeout or interrupt, remove our ticket so we don't block others + try: + os.unlink(os.path.join(tdir, ticket)) + except Exception: + pass + raise + + +def release(slot): + """Release a slot returned by acquire().""" + slot.release() + + +def cdp_slot(node, priority=PRIORITY_NORMAL, timeout=ACQUIRE_TIMEOUT): + """ + Context manager. Usage: + with cdp_slot("opm", priority=PRIORITY_HIGH): + ... CDP work ... + """ + return acquire(node, priority=priority, timeout=timeout) + + +def queue_depth(node): + """Current number of waiters for a node (for monitoring).""" + if node not in VALID_NODES: + raise ValueError("unknown node: %r" % node) + return len(_read_tickets(node))