147 lines
3.6 KiB
Python
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()
|