Files
gzth/mqtt/mqtt_service.py
T

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)