2026-10-04 13:04:24 +00:00
|
|
|
#!/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):
|
2026-10-04 16:47:36 +00:00
|
|
|
"""Return sorted list of (priority, timestamp, ticket_name), purging stale tickets."""
|
2026-10-04 13:04:24 +00:00
|
|
|
tdir = _tickets_dir(node)
|
|
|
|
|
out = []
|
2026-10-04 16:47:36 +00:00
|
|
|
now = time.time()
|
2026-10-04 13:04:24 +00:00
|
|
|
try:
|
|
|
|
|
for name in os.listdir(tdir):
|
|
|
|
|
if not name.endswith(".ticket"):
|
|
|
|
|
continue
|
|
|
|
|
try:
|
|
|
|
|
# ticket name: "<prio>-<timestamp>-<uuid>.ticket"
|
|
|
|
|
parts = name[:-7].split("-")
|
|
|
|
|
prio_s, ts_s = parts[0], parts[1]
|
2026-10-04 16:47:36 +00:00
|
|
|
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))
|
2026-10-04 13:04:24 +00:00
|
|
|
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: "<prio>-<timestamp>-<uuid>.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))
|