Files

147 lines
3.6 KiB
Python

import configparser
import paho.mqtt.client as mqtt
import sys
import logging
from odoo_rpc import OdooRPCClient
logging.basicConfig(
stream=sys.stdout,
level=logging.INFO,
format='%(asctime)s [%(levelname)s] %(message)s'
)
logger = logging.getLogger(__name__)
import json
import random
import urllib.request
HOST = '127.0.0.1'
PORT = 8069
DB = 'odoo17_jinan'
USER = 'thtzjt'
PASS = 'gzsthtz@126.com'
def json_rpc(url, method, params):
data = {
"jsonrpc": "2.0",
"method": method,
"params": params,
"id": random.randint(0, 1000000000),
}
req = urllib.request.Request(url=url, data=json.dumps(data).encode(), headers={
"Content-Type": "application/json",
})
reply = json.loads(urllib.request.urlopen(req).read().decode('UTF-8'))
if reply.get("error"):
raise Exception(reply["error"])
if reply.get('result') != None:
return reply["result"]
return reply
def call(url, service, method, *args):
return json_rpc(url, "call", {"service": service, "method": method, "args": args})
# log in the given database
url = "http://%s:%s/jsonrpc" % (HOST, PORT)
uid = call(url, "common", "login", DB, USER, PASS)
# create a new note
# args = {
# 'color': 8,
# 'memo': 'This is another note',
# 'create_uid': uid,
# }
# note_id = call(url, "object", "execute", DB, uid, PASS, 'note.note', 'create', args)
def get_mqtt_message():
# 获取所有的下发消息
return call(url, 'object', 'execute', DB, uid, PASS, 'iot.message', 'get_draft_records', [], {})
def to_reply(result_reply):
return call(url, 'object', 'execute', DB, uid, PASS, 'iot.message', 'to_reply', result_reply)
replay_topic = 'iot/service/deviceid'
# 连接成功回调
def on_connect(client, userdata, flags, rc):
# logger.info('Connected with result code '+str(rc))
client.subscribe(replay_topic, qos=2)
# 消息接收回调
def on_message(client, userdata, msg):
# redisClient.rpush('personCreate_reply',msg.payload.decode('utf-8'))
# logger.info(msg.topic+" "+str(msg.payload.decode('utf-8')))
# 数据直接传回
#to_reply(msg.payload.decode('utf-8'))
print(userdata)
print(msg)
client = mqtt.Client(client_id='Dest')
# 指定回调函数
client.on_connect = on_connect
client.on_message = on_message
# 建立连接
client.connect('1.12.37.24', 1883, 60)
# 发布消息
# client.publish(send_topic,payload=json.dumps(payload),qos=2)
client.loop_start()
# 循环读取 redis 数据
import time
logger = logging.getLogger(__name__)
config = configparser.ConfigParser()
config.read_file(open('config.ini'))
odoo_cfg = config['odoo']
mqtt_cfg = config['mqtt']
app_cfg = config['app']
odoo_url = f"{odoo_cfg['protocol']}://{odoo_cfg['host']}:{odoo_cfg['port']}"
poll_interval = app_cfg.getfloat('poll_interval', 2.0)
odoo = OdooRPCClient(
url=odoo_url,
db=odoo_cfg['db'],
username=odoo_cfg['user'],
password=odoo_cfg['password']
)
# 新线程执行的代码:
while True:
try:
# 获取所有的下发消息
messages = results = odoo.execute('iot.message', 'get_draft_records')
#logger.info(f'获取所有设备:{messages}')
# 根据设备信息拉取未执行完成的
for msg in messages:
logger.info(f'推送消息:{msg}')
data = {
'from': msg.get('device_sn'),
'data': json.loads(msg.get('message_data')),
}
client.publish(msg.get('topic'), payload=json.dumps(data), qos=2)
# logger.info('执行查询数据')
time.sleep(300)
except Exception as e:
logger.info(f'抛出异常,异常原因:{e}')
client.loop_stop()