RabbitMQ notifications
RabbitMQ is the message channel of Face Matcher. Every service in the deployment talks to the broker, and the same events that the GraphQL API exposes as subscriptions are published as broker messages. Consume them from the broker when you want a durable queue that survives your consumer being down, when you have several consumers, or when your stack already speaks AMQP or MQTT. Smart Corridor consumes Face Matcher this way.
Broker endpoints
| Protocol | From the host | Inside face-matcher-network | Used for |
|---|---|---|---|
| AMQP 0-9-1 | localhost:5672 | rmq:5672 | Notifications and all internal messaging |
| MQTT | localhost:1883 | rmq:1883 | Edge devices publishing frame data; watchlist synchronization to edge streams |
| RabbitMQ streams | localhost:5552 | rmq:5552 | Watchlist update log consumed by the synchronization services |
| Management UI | http://localhost:15672 | rmq:15672 | Inspect exchanges, queues and bindings |
The broker is RabbitMQ 4.3 with the management, MQTT, streams and Prometheus plugins enabled. Credentials are in section 2.2 (RabbitMQ__Username / RabbitMQ__Password, guest/guest by default) and section 2.3 (MQTT__*) of .env; anonymous MQTT is disabled. Change the credentials and restrict the published ports before production, see Network and ports. The management UI is the fastest way to discover the exact exchange and routing keys on your version: open it, go to Exchanges, and watch the message rates while a camera is running.
Notification types
Face Matcher raises a notification for every event in the processing pipeline. There are two kinds, and the difference matters for ordering:
- Direct notifications are published the moment the event happens, before anything is written to the database. They are sent regardless of the save strategy, so even a camera configured to store nothing still produces them. Use them for real-time reactions such as opening a gate.
- Database notifications are published after the entity has been stored, so the data they refer to can immediately be queried through GraphQL. Because storing takes a variable amount of time, database notifications can arrive in a different order than the events occurred, and a match notification typically arrives before the corresponding face-created notification.
| Notification | Kind | When it fires and what it carries |
|---|---|---|
faceProcessed | Direct | A face was detected on a frame: face crop coordinates, frame, match information and liveness result |
pedestrianProcessed | Direct | A pedestrian was detected: pedestrian and frame |
objectProcessed | Direct | An object was detected: object and frame |
identificationEvent | Direct | The identification outcome for a detected object, including stream information and, with Notifications__IncludeTemplates=true, the face template |
frameProcessed | Direct | One event per processed frame, even with no detections: stream, frame, all detected objects with their identification results and the links between them |
matchResult / noMatchResult | Direct | A face did / did not match a watchlist member |
faceCreated | Database | The face has been saved; attributes are not extracted yet |
faceExtracted | Database | Face attributes (age, gender, face mask) have been extracted and saved |
matchResultInsert | Database | The match result has been saved |
trackletCompleted | Database | The tracked face or pedestrian left the scene and its tracklet was closed |
pedestrianInserted | Database | The pedestrian has been saved |
The fields of each notification are documented in the GraphQL schema (http://localhost:8097/graphql?sdl), and Events and face metadata explains the attributes they carry.
Consuming notifications over AMQP
Bind your own queue to the notification exchange with a routing key pattern and consume from it. The example below uses Python with the pika library; the exchange name and routing key are placeholders to fill in from the management UI.
import json
import pika
credentials = pika.PlainCredentials("guest", "guest")
connection = pika.BlockingConnection(
pika.ConnectionParameters(host="rmq", port=5672, virtual_host="/", credentials=credentials)
)
channel = connection.channel()
# Replace with the notification exchange and routing key of your installation.
EXCHANGE = "<notification-exchange>"
ROUTING_KEY = "#" # all notifications; narrow it to one topic in production
queue = channel.queue_declare(queue="my-integration", durable=True).method.queue
channel.queue_bind(exchange=EXCHANGE, queue=queue, routing_key=ROUTING_KEY)
def on_message(ch, method, properties, body):
event = json.loads(body)
print(method.routing_key, event.get("StreamId"))
ch.basic_ack(delivery_tag=method.delivery_tag)
channel.basic_consume(queue=queue, on_message_callback=on_message)
channel.start_consuming()
Declare your queue durable and acknowledge messages only after you have handled them, so a restart of your consumer loses nothing. Do not consume from the platform's own queues, and do not declare exchanges with the platform's names; treat the broker as shared infrastructure. The broker's consumer_timeout is six hours, so long-running handlers must still acknowledge within that window.
Edge devices over MQTT
The MQTT listener is the inbound side of the broker. Smart cameras and AI boxes running the Embedded Stream Processor publish FrameData messages to it, which the edge-stream-processor service consumes from the amq.topic exchange with the routing key edge-stream.*.frame_data (queue MQTT_EDGE_STREAM_CONSUMER); the same broker carries the watchlist synchronization messages back to the devices. To connect a device, point it at rmq:1883 (or the host address and port 1883) with the MQTT__* credentials, as described in Connect to Face Matcher and the MQTT API.