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 |
maxi07
commented
Jul 2, 2026
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.
maxi07
commented
Jul 4, 2026
@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
Uh oh!
There was an error while loading. Please reload this page.
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.