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

Function process_external_message

methods/MIRIX/mirix/queue/__init__.py:75–115  ·  view source on GitHub ↗

Process a message consumed by an external system (e.g., Numaflow, custom Kafka consumer). This is the primary high-level API for integrating with external Kafka consumers or event processing systems. It handles all internal details of deserialization and processing. Args:

(raw_message: bytes)

Source from the content-addressed store, hash-verified

73
74
75async def process_external_message(raw_message: bytes) -> None:
76 """
77 Process a message consumed by an external system (e.g., Numaflow, custom Kafka consumer).
78
79 This is the primary high-level API for integrating with external Kafka consumers or event
80 processing systems. It handles all internal details of deserialization and processing.
81
82 Args:
83 raw_message: Raw message bytes from Kafka or event bus (JSON or protobuf format)
84
85 Raises:
86 ValueError: If message parsing fails
87 """
88 if not _manager.is_initialized:
89 logger.info("Queue not initialized, auto-initializing with server for external message processing")
90 from mirix.server.server import AsyncServer
91
92 server = AsyncServer()
93 await _manager.initialize(server=server)
94 logger.info("Queue initialized with server instance")
95
96 workers = _manager._workers
97 if not workers:
98 logger.error("No workers available after initialization - this should not happen!")
99 raise RuntimeError("Failed to create queue workers during initialization")
100
101 worker = workers[0]
102
103 from mirix.queue.config import KAFKA_SERIALIZATION_FORMAT
104 from mirix.queue.queue_util import deserialize_queue_message
105
106 queue_message = deserialize_queue_message(raw_message, format=KAFKA_SERIALIZATION_FORMAT)
107
108 logger.debug(
109 "Processing external message (%s format): agent_id=%s, user_id=%s",
110 KAFKA_SERIALIZATION_FORMAT,
111 queue_message.agent_id,
112 queue_message.user_id if queue_message.HasField("user_id") else "None",
113 )
114
115 await worker.process_external_message(queue_message)
116
117
118__all__ = ["initialize_queue", "save", "process_external_message", "QueueMessage"]

Callers

nothing calls this directly

Calls 7

AsyncServerClass · 0.90
infoMethod · 0.65
errorMethod · 0.65
debugMethod · 0.65
initializeMethod · 0.45

Tested by

no test coverage detected