MCPcopy Create free account
hub / github.com/AQ-MedAI/MedMemoryBench / initialize

Method initialize

methods/MIRIX/mirix/queue/manager.py:37–125  ·  view source on GitHub ↗

Initialize the queue and start the background workers. This method is idempotent — calling it multiple times only initializes once. For in-memory queues: - NUM_WORKERS=1 (default): Single queue, single worker - NUM_WORKERS>1: Partitioned queue with N worker

(self, server: Optional[Any] = None)

Source from the content-addressed store, hash-verified

35 self._round_robin = False
36
37 async def initialize(self, server: Optional[Any] = None) -> None:
38 """
39 Initialize the queue and start the background workers.
40
41 This method is idempotent — calling it multiple times only initializes once.
42
43 For in-memory queues:
44 - NUM_WORKERS=1 (default): Single queue, single worker
45 - NUM_WORKERS>1: Partitioned queue with N workers (simulates Kafka)
46
47 Args:
48 server: Optional server instance for workers to invoke APIs on
49 """
50 if self._initialized:
51 logger.warning("Queue manager already initialized - skipping duplicate initialization")
52 worker_count = len(self._workers)
53 running_count = sum(1 for w in self._workers if w._running)
54 logger.info(f" Current state: workers={worker_count}, running={running_count}")
55 if server:
56 logger.info("Updating queue manager with server instance")
57 self._server = server
58 for worker in self._workers:
59 worker.set_server(server)
60 return
61
62 if config.QUEUE_TYPE == "memory":
63 self._num_workers = config.NUM_WORKERS
64 self._round_robin = config.ROUND_ROBIN
65 else:
66 self._num_workers = 1
67 self._round_robin = False
68
69 partition_mode = "round-robin" if self._round_robin else "hash"
70 logger.info(
71 "Initializing queue manager: type=%s, num_workers=%d, partitioning=%s, server=%s",
72 config.QUEUE_TYPE,
73 self._num_workers,
74 partition_mode,
75 "provided" if server else "None",
76 )
77
78 self._server = server
79
80 logger.info("Creating queue instance...")
81 self._queue = self._create_queue()
82 logger.info(f"Queue created: type={type(self._queue).__name__}")
83
84 await self._queue.start()
85
86 self._workers = []
87
88 if self._num_workers > 1 and isinstance(self._queue, PartitionedMemoryQueue):
89 logger.info("Creating %d background workers (partitioned)...", self._num_workers)
90 for partition_id in range(self._num_workers):
91 worker = QueueWorker(self._queue, server=self._server, partition_id=partition_id)
92 self._workers.append(worker)
93 logger.debug("Worker %d created", partition_id)
94 else:

Calls 8

set_serverMethod · 0.95
_create_queueMethod · 0.95
startMethod · 0.95
QueueWorkerClass · 0.90
infoMethod · 0.65
debugMethod · 0.65
errorMethod · 0.65
warningMethod · 0.45