#!/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", "dev", "def") 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), purging stale tickets.""" tdir = _tickets_dir(node) out = [] now = time.time() 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] ts = float(ts_s) # Purge stale ticket if process crashed or timed out ungracefully if now - ts > (ACQUIRE_TIMEOUT * 2): try: os.unlink(os.path.join(tdir, name)) except Exception: pass continue out.append((int(prio_s), ts, 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))