64 lines
2.4 KiB
Python
64 lines
2.4 KiB
Python
import asyncio
|
|
import json
|
|
import logging
|
|
|
|
import aiomqtt
|
|
|
|
from odoo_rpc import OdooRPCClient
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
async def handle_incoming_messages(client: aiomqtt.Client, odoo: OdooRPCClient, topics: list):
|
|
"""Subscribes to topics and handles incoming messages."""
|
|
if not topics:
|
|
logger.warning("No 'reply_topics' configured. No messages will be received.")
|
|
return
|
|
|
|
for topic in topics:
|
|
await client.subscribe(topic, qos=2)
|
|
logger.info(f"Subscribed to reply topic: {topic}")
|
|
|
|
async for message in client.messages:
|
|
try:
|
|
payload_str = message.payload.decode('utf-8')
|
|
logger.info(f"Received message on topic '{message.topic}'")
|
|
logger.info(f"Received message on topic '{payload_str}'")
|
|
odoo.execute('iot.message', 'to_reply', json.loads(payload_str))
|
|
logger.info("Successfully forwarded reply to Odoo.")
|
|
except UnicodeDecodeError:
|
|
logger.warning(f"Could not decode message payload on topic {message.topic}: {message.payload}")
|
|
except Exception as e:
|
|
logger.error(f"Error processing incoming MQTT message: {e}", exc_info=True)
|
|
|
|
|
|
async def poll_and_publish(client: aiomqtt.Client, odoo: OdooRPCClient, poll_interval: float):
|
|
"""Periodically polls Odoo for devices and messages to publish."""
|
|
while True:
|
|
try:
|
|
# 1. Get all mesage from Odoo
|
|
results = odoo.execute('iot.message', 'get_draft_records')
|
|
if not results:
|
|
logger.info("No active devices found. Waiting...")
|
|
await asyncio.sleep(poll_interval)
|
|
continue
|
|
|
|
logger.info(f"Found {len(results)} info. Fetching pending messages concurrently.")
|
|
|
|
# 3. Publish fetched messages
|
|
for result in results:
|
|
if isinstance(result, Exception):
|
|
logger.error(f"Error fetching MQTT line from Odoo: {result}")
|
|
continue
|
|
topic = result.get('topic')
|
|
payload = {
|
|
'from': result.get('id'),
|
|
'data': json.loads(result.get('message_data')),
|
|
}
|
|
await client.publish(topic, payload=json.dumps(payload), qos=2)
|
|
|
|
except Exception as e:
|
|
logger.error(f"Error in polling loop: {e}", exc_info=True)
|
|
|
|
await asyncio.sleep(poll_interval)
|