Add cdp_queue.py: per-browser CDP operation queue
Per-node FIFO, max 2 concurrent, priority levels (high/normal/low), flock-based cross-process coordination. DM sends use high priority. Session: sidechat/chromebox-ops
This commit is contained in:
@@ -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: "<prio>-<timestamp>-<uuid>.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: "<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))
|
||||
Reference in New Issue
Block a user