Unify RabbitMQ handling with persistent, auto-reconnecting connections#70
Conversation
There was a problem hiding this comment.
Pull request overview
This PR centralizes RabbitMQ connectivity into scansynclib to eliminate per-message connection churn and missed-heartbeat disconnects across ScanSync’s microservices, improving reliability for long-running OCR/upload/naming workloads and SSE updates.
Changes:
- Introduces a unified
RabbitMQClient+ helper API inscansynclib.rabbitmqwith long-lived, auto-reconnecting publish/consume patterns (and a backward-compatibleconnect_rabbitmq). - Refactors services to consume/publish via the shared helpers instead of duplicating connect/reconnect logic.
- Adds focused unit tests for connection reuse and reconnect behavior.
Reviewed changes
Copilot reviewed 11 out of 11 changed files in this pull request and generated 5 comments.
Show a summary per file
| File | Description |
|---|---|
| web_service/src/main.py | Keeps SSE RabbitMQ listener thread alive via reconnect loop. |
| upload_service/main.py | Switches upload consumer to unified consume(...) helper. |
| ocr_service/main.py | Switches OCR consumer to unified consume(...) helper. |
| metadata_service/main.py | Publishes to OCR queue via unified publish(...) and consumes via consume(...). |
| file_naming_service/main.py | Switches consumer to unified consume(...) and changes ack/reconnect behavior. |
| detection_service/main.py | Uses RabbitMQClient + process_events() to keep a single connection alive between scans. |
| scansynclib/scansynclib/rabbitmq.py | New unified, reconnecting RabbitMQ connection/publish/consume implementation. |
| scansynclib/scansynclib/helpers.py | Re-exports RabbitMQ helpers for backward-compatible imports. |
| scansynclib/scansynclib/sqlite_wrapper.py | Routes SSE publish through unified exchange publisher helper. |
| tests/test_rabbitmq.py | Adds tests covering connection reuse and reconnect logic. |
| scansynclib/scansynclib.egg-info/SOURCES.txt | Updates packaged file list (currently introduces duplicates). |
| connection, channel = result | ||
|
|
||
| try: | ||
| # Use fanout as exchange type to broadcast messages to all connected clients | ||
| channel.exchange_declare(exchange=exchange_name, exchange_type='fanout') | ||
| queue_result = channel.queue_declare(queue='', exclusive=True) | ||
| queue_name = queue_result.method.queue | ||
| channel.queue_bind(exchange=exchange_name, queue=queue_name) | ||
|
|
||
| channel.basic_consume(queue=queue_name, on_message_callback=callback, auto_ack=True) | ||
| channel.start_consuming() | ||
| except pika.exceptions.AMQPError as e: | ||
| logger.error(f"RabbitMQ SSE listener connection lost: {e}. Reconnecting in 5 seconds...") | ||
| time.sleep(5) | ||
| except Exception as e: | ||
| logger.exception(f"Unexpected error in SSE listener: {e}. Reconnecting in 5 seconds...") | ||
| time.sleep(5) |
| def declare_queue(self, queue_name: str, durable: bool = True): | ||
| if not queue_name or queue_name in self._declared_queues: | ||
| return | ||
| self._channel.queue_declare(queue=queue_name, durable=durable) | ||
| self._declared_queues.add(queue_name) | ||
|
|
||
| def declare_exchange(self, exchange: str, exchange_type: str = "fanout"): | ||
| if not exchange or exchange in self._declared_exchanges: | ||
| return | ||
| self._channel.exchange_declare(exchange=exchange, exchange_type=exchange_type) | ||
| self._declared_exchanges.add(exchange) |
| client = RabbitMQClient(name="detection") | ||
| client.ensure_connection() | ||
| client.declare_queue(RABBITQUEUE) |
| ch.basic_ack(delivery_tag=method.delivery_tag) | ||
| except pika.exceptions.AMQPConnectionError: | ||
| logger.error("Connection lost while acknowledging message. Reconnecting...") | ||
| connection, channel = connect_rabbitmq([RABBITQUEUE], heartbeat=120) | ||
| channel.basic_ack(delivery_tag=method.delivery_tag) | ||
| except pika.exceptions.AMQPError: | ||
| # The connection was lost before we could acknowledge. The unified | ||
| # consumer will reconnect and the broker will redeliver the message. | ||
| logger.error("Connection lost while acknowledging message. It will be redelivered after reconnect.") |
| scansynclib/ProcessItem.py | ||
| scansynclib/__init__.py | ||
| scansynclib/config.py | ||
| scansynclib/helpers.py | ||
| scansynclib/logging.py | ||
| scansynclib/ollama_helper.py | ||
| scansynclib/onedrive_api.py | ||
| scansynclib/onedrive_smb_manager.py | ||
| scansynclib/openai_helper.py | ||
| scansynclib/settings.py | ||
| scansynclib/settings_schema.py | ||
| scansynclib/sqlite_wrapper.py |
|
remember to update rabbitmq as well. |
| # The connection was lost before we could acknowledge. The unified | ||
| # consumer will reconnect and the broker will redeliver the message. | ||
| logger.error("Connection lost while acknowledging message. It will be redelivered after reconnect.") | ||
| if not acknowledged: | ||
| return |
| _publisher = RabbitMQClient(name="publisher") | ||
| _publisher_lock = threading.Lock() | ||
|
|
| scansynclib/config.py | ||
| scansynclib/helpers.py | ||
| scansynclib/logging.py | ||
| scansynclib/ollama_helper.py | ||
| scansynclib/onedrive_api.py |
| try: | ||
| ch.basic_ack(delivery_tag=method.delivery_tag) | ||
| except pika.exceptions.AMQPConnectionError: | ||
| logger.error("Connection lost while acknowledging message. Reconnecting...") | ||
| connection, channel = connect_rabbitmq([RABBITQUEUE], heartbeat=120) | ||
| channel.basic_ack(delivery_tag=method.delivery_tag) | ||
| except pika.exceptions.AMQPError: | ||
| # The connection was lost before we could acknowledge. The unified | ||
| # consumer will reconnect and the broker will redeliver the message. | ||
| logger.error("Connection lost while acknowledging message. It will be redelivered after reconnect.") |
There was a problem hiding this comment.
Fixed in commit a3aebb8. The finally block now tracks ack_ok and only proceeds with DB updates and upload_queue forwarding when the ack succeeded. The test was also updated to verify that these follow-up actions are skipped on ack failure.
|
@copilot resolve the merge conflicts in this pull request |
Resolve conflicts: - detection_service/main.py: keep RabbitMQClient setup, add cleanup_dangling_documents() from main - upload_service/main.py: keep consume() from helpers, add execute_query import from main
Resolved in merge commit
|
…om/maxi07/ScanSync into copilot/handle-rabbitmq-connection
Services each built their own RabbitMQ connections and rebuilt them constantly, so the broker dropped clients with
missed heartbeats from client, timeout: 60s. The goal: establish connections once, keep them alive, and centralize the logic.Root cause
sqlite_wrapperopened the SSE publisher with the default 60s heartbeat and never serviced it during long OCR/upload/naming work → broker closed it every ~60s.forward_to_rabbitmqopened and closed a new connection per message.Changes
scansynclib/rabbitmq.py— single source of truth for messaging:RabbitMQClient: wraps one long-livedBlockingConnection, lazily connects, re-declares queues/exchanges after a reconnect, and transparently reconnects on publish/consume failures (default heartbeat 600s).publish,forward_to_rabbitmq(shared publisher connection),publish_to_exchange,consume(unified reconnecting consumer loop),process_events(heartbeat servicing for polling publishers), and a backward-compatibleconnect_rabbitmq.helpers.pyre-exports these so existing imports keep working.sqlite_wrapper.py— SSEnotify_sse_clientsnow publishes through the self-healing publisher (fixes the 60s drops).consume(...); metadata publishes toocr_queueviapublish(...); detection usesRabbitMQClient+process_events; the web SSE listener reconnects instead of killing its thread. file_naming's fragile manual re-ack was dropped (broker redelivers after reconnect).tests/test_rabbitmq.py— covers connection reuse, reconnect onStreamLostError, publish-failure handling, exchange declaration, pickling inforward_to_rabbitmq, and consumer reconnection.Before → after
Publishers and consumers use separate connections (recommended by RabbitMQ), and all reconnect automatically, eliminating the connection churn.