Implement a multi-agent deep web research subsystem including CLI commands for pruning expired jobs and clearing all artifacts, database tables for sessions/messages/URL cache with indexes, config paths for chroma storage, and internal embed URL for vector operations.
44 lines
1.2 KiB
Python
44 lines
1.2 KiB
Python
# retoor <retoor@molodetz.nl>
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
|
|
MAX_BUFFER = 2000
|
|
|
|
|
|
class ProgressHub:
|
|
def __init__(self) -> None:
|
|
self._subscribers: dict[str, set[asyncio.Queue]] = {}
|
|
self._buffers: dict[str, list[dict]] = {}
|
|
|
|
def subscribe(self, uid: str) -> asyncio.Queue:
|
|
queue: asyncio.Queue = asyncio.Queue()
|
|
self._subscribers.setdefault(uid, set()).add(queue)
|
|
return queue
|
|
|
|
def unsubscribe(self, uid: str, queue: asyncio.Queue) -> None:
|
|
listeners = self._subscribers.get(uid)
|
|
if not listeners:
|
|
return
|
|
listeners.discard(queue)
|
|
if not listeners:
|
|
self._subscribers.pop(uid, None)
|
|
|
|
def publish(self, uid: str, frame: dict) -> None:
|
|
buffer = self._buffers.setdefault(uid, [])
|
|
buffer.append(frame)
|
|
if len(buffer) > MAX_BUFFER:
|
|
del buffer[: len(buffer) - MAX_BUFFER]
|
|
for queue in self._subscribers.get(uid, set()):
|
|
queue.put_nowait(frame)
|
|
|
|
def snapshot(self, uid: str) -> list[dict]:
|
|
return list(self._buffers.get(uid, []))
|
|
|
|
def clear(self, uid: str) -> None:
|
|
self._buffers.pop(uid, None)
|
|
|
|
|
|
hub = ProgressHub()
|