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)
| 73 | |
| 74 | |
| 75 | async 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"] |
nothing calls this directly
no test coverage detected