Files
box/bin/cdp_queue.py
T

240 lines
7.8 KiB
Python

#!/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: "<prio>-<timestamp>-<uuid>.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: "<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))