Files
2023-08-18 11:55:54 +08:00

136 lines
2.7 KiB
Go

package mqtt
import (
"encoding/json"
"log"
"os"
"gitee.com/jhx1219/machine-go/utils/common/job"
mqtt_client "gitee.com/jhx1219/machine-go/utils/common/mqtt"
"gitee.com/jhx1219/machine-go/utils/common/receive"
"gitee.com/jhx1219/machine-go/utils/common/task"
mqtt "github.com/eclipse/paho.mqtt.golang"
)
type Payload map[string]interface{}
type Message struct {
Payload Payload `json:"payload"`
Topic string `json:"topic"`
}
type FaceMqtt struct {
MqttClient *mqtt_client.MqttClient
subTopicList []string
responseFunc func(*FaceMqtt, Message)
ReceiveServer *receive.Receive
JobServer *job.Job
}
func NewFaceMqtt() *FaceMqtt {
fm := &FaceMqtt{}
fm.responseFunc = nil
// 初始化接收数据异步任务
log.Println("初始化异步任务...")
// 初始化job异步任务
job := job.NewJob()
fm.JobServer = job
t := task.NewTask(os.Getenv("RECEIVE_PATH"))
r := receive.NewReceive(t)
r.SetTaskFunc(func(msg string, task *task.Task) {
job.LogJobTask(msg)
})
fm.ReceiveServer = r
log.Println("异步任务初始化完成")
return fm
}
func (fm *FaceMqtt) SubTopic(topic string) {
fm.subTopicList = append(fm.subTopicList, topic)
}
func (fm *FaceMqtt) SetResponseFunc(f func(*FaceMqtt, Message)) {
fm.responseFunc = f
}
func (fm *FaceMqtt) PublishMsg(topic string, payload string) {
fm.MqttClient.Publish(topic, []byte(payload))
}
func (fm *FaceMqtt) MqttServer() {
go func() {
fm.ReceiveServer.ReceiveTaskServer()
}()
// 初始化mqtt
log.Println("初始化mqtt...")
m := mqtt_client.NewMqttClient()
var subMsgChan = make(chan Message, 20)
handler := func(client mqtt.Client, msg mqtt.Message) {
p := make(Payload)
err := json.Unmarshal(msg.Payload(), &p)
if err != nil {
log.Println("mqtt消息解析失败: " + err.Error())
return
}
subMsgChan <- Message{
Payload: p,
Topic: msg.Topic(),
}
}
for _, topic := range fm.subTopicList {
m.SubTopic(topic, handler)
}
m.Connect(mqtt_client.MqttConf{
Host: os.Getenv("MQTT_HOST"),
Port: os.Getenv("MQTT_PORT"),
Username: os.Getenv("MQTT_USERNAME"),
Password: os.Getenv("MQTT_PASSWORD"),
})
fm.MqttClient = m
go func() {
fm.mqttSubMsgServer(subMsgChan)
}()
log.Println("mqtt初始化完成")
}
func (fm *FaceMqtt) mqttSubMsgServer(subMsgChan chan Message) {
for {
msg := <-subMsgChan
go func() {
// 记录到达消息
logBytes, err := json.Marshal(msg)
if err == nil {
err = fm.ReceiveServer.LogReceice(string(logBytes))
if err != nil {
return
}
} else {
log.Printf("消息记录 %s 错误: %s\n", logBytes, err)
return
}
// 处理回复
if fm.responseFunc != nil {
fm.responseFunc(fm, msg)
}
}()
}
}