From b4fb4e8ed26004189e6a3858d32b61588469bd77 Mon Sep 17 00:00:00 2001 From: jhx <133451314@qq.com> Date: Fri, 18 Aug 2023 11:55:54 +0800 Subject: [PATCH] update --- .gitignore | 2 + utils/common/job/job.go | 80 ++++++++++ utils/common/mqtt/mqtt.go | 91 ++++++++++++ utils/common/receive/receive.go | 39 +++++ utils/common/request/request.go | 70 +++++++++ utils/common/task/task.go | 253 ++++++++++++++++++++++++++++++++ utils/face/mqtt/mqtt.go | 135 +++++++++++++++++ utils/go.mod | 11 ++ utils/go.sum | 12 ++ 9 files changed, 693 insertions(+) create mode 100644 .gitignore create mode 100644 utils/common/job/job.go create mode 100644 utils/common/mqtt/mqtt.go create mode 100644 utils/common/receive/receive.go create mode 100644 utils/common/request/request.go create mode 100644 utils/common/task/task.go create mode 100644 utils/face/mqtt/mqtt.go create mode 100644 utils/go.mod create mode 100644 utils/go.sum diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..a82a80d --- /dev/null +++ b/.gitignore @@ -0,0 +1,2 @@ +/center +/face \ No newline at end of file diff --git a/utils/common/job/job.go b/utils/common/job/job.go new file mode 100644 index 0000000..212a174 --- /dev/null +++ b/utils/common/job/job.go @@ -0,0 +1,80 @@ +package job + +import ( + "encoding/json" + "io/ioutil" + "log" + "os" + + "gitee.com/jhx1219/machine-go/utils/common/task" +) + +type JobConf struct { + JobList []JobListConf `json:"job_list"` +} + +type JobListConf struct { + Name string `json:"name"` +} + +type JobInfo struct { + Task *task.Task +} + +type Job struct { + JobLogPath string + JobMap map[string]*JobInfo +} + +func NewJob() *Job { + + job := &Job{} + job.JobLogPath = os.Getenv("JOB_LOG_PATH") + job.JobMap = make(map[string]*JobInfo) + + confPath := os.Getenv("CONF_PATH") + + data, err := ioutil.ReadFile(confPath + "/job.json") + if err != nil { + log.Println("读取job文件失败: " + err.Error()) + return job + } + + // 解析 JSON 数据 + var jobConf JobConf + err = json.Unmarshal(data, &jobConf) + if err != nil { + log.Println("解析job文件失败: " + err.Error()) + return job + } + + for _, conf := range jobConf.JobList { + t := task.NewTask(job.JobLogPath + "/" + conf.Name) + job.JobMap[conf.Name] = &JobInfo{ + Task: t, + } + } + + return job +} + +func (j *Job) JobTaskServer() { + for _, job := range j.JobMap { + jobCopy := job + jobCopy.Task.TaskServer() + } +} + +func (j *Job) LogJobTask(msg string) { + for _, job := range j.JobMap { + job.Task.LogTask(msg) + } +} + +func (j *Job) SetJobFunc(job string, f func(string, *task.Task)) { + if _, ok := j.JobMap[job]; !ok { + return + } + + j.JobMap[job].Task.TaskFunc = f +} diff --git a/utils/common/mqtt/mqtt.go b/utils/common/mqtt/mqtt.go new file mode 100644 index 0000000..9f20ac3 --- /dev/null +++ b/utils/common/mqtt/mqtt.go @@ -0,0 +1,91 @@ +package mqtt + +import ( + "fmt" + "log" + + mqtt "github.com/eclipse/paho.mqtt.golang" +) + +type MqttClient struct { + Client mqtt.Client + SubscribeList []SubscribeConf +} + +type MqttConf struct { + Host string + Port string + Username string + Password string + ClientId string +} + +type SubscribeConf struct { + Topic string + Handler func(mqtt.Client, mqtt.Message) +} + +func NewMqttClient() *MqttClient { + c := &MqttClient{} + c.Client = nil + c.SubscribeList = make([]SubscribeConf, 0) + + return c +} + +func (c *MqttClient) SubTopic(topic string, f func(mqtt.Client, mqtt.Message)) { + c.SubscribeList = append(c.SubscribeList, SubscribeConf{ + Topic: topic, + Handler: f, + }) +} + +func (c *MqttClient) Connect(conf MqttConf) { + host := conf.Host + port := conf.Port + username := conf.Username + password := conf.Password + clientId := conf.ClientId + + opts := mqtt.NewClientOptions() + opts.AddBroker(fmt.Sprintf("tcp://%s:%s", host, port)) + opts.SetAutoReconnect(true) // 自动重连 + opts.SetUsername(username) + opts.SetPassword(password) + opts.SetClientID(clientId) + + opts.SetOnConnectHandler(func(client mqtt.Client) { + log.Println("mqtt 连接成功") + + // 创建订阅 + log.Println("加载订阅...") + c.CreateSubscribe() + log.Println("加载订阅完成") + }) + + opts.SetConnectionLostHandler(func(client mqtt.Client, err error) { + log.Println("mqtt 连接断开") + }) + + client := mqtt.NewClient(opts) + + if token := client.Connect(); token.Wait() && token.Error() != nil { + panic("创建mqtt客户端错误,错误信息: " + token.Error().Error()) + } + + c.Client = client +} + +func (c *MqttClient) Publish(topic string, payload []byte) { + token := c.Client.Publish(topic, 0, false, payload) + token.Wait() +} + +func (c *MqttClient) CreateSubscribe() { + for _, conf := range c.SubscribeList { + if token := c.Client.Subscribe(conf.Topic, 1, conf.Handler); token.Wait() && token.Error() != nil { + log.Fatal(token.Error()) + } + log.Printf("订阅 %s 完成\n", conf.Topic) + } +} diff --git a/utils/common/receive/receive.go b/utils/common/receive/receive.go new file mode 100644 index 0000000..5c02782 --- /dev/null +++ b/utils/common/receive/receive.go @@ -0,0 +1,39 @@ +package receive + +import ( + "gitee.com/jhx1219/machine-go/utils/common/task" +) + +type Receive struct { + Task *task.Task +} + +func NewReceive(t *task.Task) *Receive { + + t.SetTaskFunc(func(msg string, task *task.Task) { + //log.Println(msg) + }) + + return &Receive{ + Task: t, + } +} + +func (r *Receive) SetTaskFunc(f func(string, *task.Task)) { + r.Task.SetTaskFunc(f) +} + +func (r *Receive) LogReceice(msg string) error { + + err := r.Task.LogTask(msg) + + if err != nil { + return err + } + + return nil +} + +func (r *Receive) ReceiveTaskServer() { + r.Task.TaskServer() +} diff --git a/utils/common/request/request.go b/utils/common/request/request.go new file mode 100644 index 0000000..b568f97 --- /dev/null +++ b/utils/common/request/request.go @@ -0,0 +1,70 @@ +package request + +import ( + "bytes" + "errors" + "fmt" + "io/ioutil" + "log" + "net/http" + "time" +) + +type Request struct { + Timeout time.Duration + MaxRetry int +} + +func NewRequest() *Request { + return &Request{ + Timeout: time.Second * 2, + } +} + +func (r *Request) SetTimeout(t time.Duration) { + r.Timeout = t +} + +func (r *Request) SetMaxRetry(retry int) { + r.MaxRetry = retry +} + +func (r *Request) Post(api string, data []byte) (string, error) { + + client := &http.Client{ + Timeout: r.Timeout, + } + + req, err := http.NewRequest("POST", api, bytes.NewBuffer(data)) + if err != nil { + return "", errors.New("Error creating request:" + err.Error()) + + } + req.Header.Set("Content-Type", "application/json") + + retry := 0 +reqDo: + resp, err := client.Do(req) + if err != nil { + if r.MaxRetry > 0 { + if retry < r.MaxRetry { + retry = retry + 1 + log.Printf("请求api第[%d], 最大重试次数[%d]", retry, r.MaxRetry) + goto reqDo + } + } + return "", errors.New("Error sending request:" + err.Error()) + } + defer resp.Body.Close() + + if resp.StatusCode != 200 { + return "", fmt.Errorf("Error StatusCode: %d", resp.StatusCode) + } + + body, err := ioutil.ReadAll(resp.Body) + if err != nil { + return "", errors.New("Error read body: " + err.Error()) + } + + return string(body), nil +} diff --git a/utils/common/task/task.go b/utils/common/task/task.go new file mode 100644 index 0000000..e9511a1 --- /dev/null +++ b/utils/common/task/task.go @@ -0,0 +1,253 @@ +package task + +import ( + "bufio" + "log" + "os" + "path/filepath" + "strconv" + "strings" + "sync" + "time" +) + +type Task struct { + Path string + FileMutex *sync.Mutex + TaskFunc func(string, *Task) + RuningTaskPath map[string]int +} + +func NewTask(path string) *Task { + + var fileMutex sync.Mutex + return &Task{ + Path: path, + FileMutex: &fileMutex, + TaskFunc: nil, + RuningTaskPath: make(map[string]int), + } +} + +func (task *Task) SetTaskFunc(f func(string, *Task)) { + task.TaskFunc = f +} + +func (task *Task) LogTask(msg string) error { + task.FileMutex.Lock() + defer task.FileMutex.Unlock() + + path := task.Path + task.checkOrCreatePath(path) + currentTime := time.Now().Unix() + + filePath := path + "/" + strconv.FormatInt(currentTime, 10) + ".log" // 目标文件路径 + + file, err := os.OpenFile(filePath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) + if err != nil { + log.Println("打开task文件失败:", err) + return err + } + defer file.Close() + + // 写入内容到文件 + _, err = file.WriteString(msg + "\n") + if err != nil { + log.Println("写入task文件失败:", err) + return err + } + + return nil +} + +func (task *Task) LogFailTask(msg string) error { + task.FileMutex.Lock() + defer task.FileMutex.Unlock() + + path := task.Path + "/fail/" + task.checkOrCreatePath(path) + + currentTime := time.Now() + currentHour := time.Date(currentTime.Year(), currentTime.Month(), currentTime.Day(), currentTime.Hour(), 0, 0, 0, currentTime.Location()) + timestamp := currentHour.Unix() + + filePath := path + strconv.FormatInt(timestamp, 10) + ".log" // 目标文件路径 + + file, err := os.OpenFile(filePath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644) + if err != nil { + log.Println("打开task文件失败:", err) + return err + } + defer file.Close() + + // 写入内容到文件 + _, err = file.WriteString(msg + "\n") + if err != nil { + log.Println("写入task文件失败:", err) + return err + } + + return nil +} + +func (task *Task) checkOrCreatePath(path string) { + dir := path + + // 检查任务目录是否存在 + _, err := os.Stat(dir) + if os.IsNotExist(err) { + // 任务目录不存在,创建任务目录 + err := os.MkdirAll(dir, 0755) + if err != nil { + log.Println(task.Path+" 目录创建失败:", err) + return + } + + log.Println(task.Path + " 目录创建成功") + } +} + +func (task *Task) TaskServer() { + // 扫描接收文件协程 + go func() { + task.ScanTask() + }() +} + +func (task *Task) ScanTask() { + path := task.Path + + log.Println("扫描任务目录: " + path) + for { + currentTime := time.Now().Unix() + + dirRead, err := os.Open(path) + + if err != nil { + time.Sleep(1 * time.Second) + continue + } + for { + files, err := dirRead.Readdir(10) // 每次读取10个文件 + if err != nil && len(files) == 0 { + // 读取完所有文件或发生错误 + time.Sleep(500 * time.Microsecond) + break + } + + for _, file := range files { + fileName := file.Name() + ext := filepath.Ext(fileName) + + fileNameWithoutExt := strings.TrimSuffix(fileName, ext) + + fileTime, err := strconv.ParseInt(fileNameWithoutExt, 10, 64) + if err != nil { + continue // 跳过无法转换为时间戳的文件名 + } + + if ext == ".log" { + if fileTime < currentTime-1 { + filePath := path + "/" + fileName + log.Println("扫描到文件: " + filePath) + newFilePath := path + "/" + fileNameWithoutExt + ".scanned" + err := os.Rename(filePath, newFilePath) + if err != nil { + log.Println(filePath + " 文件重命名失败: " + err.Error()) + } + + task.HandleTask(path + "/" + fileNameWithoutExt) + } + } + + if ext == ".scanned" { + if fileTime < currentTime-600 { + filePath := path + "/" + fileName + // 非执行中的十分钟前未处理的重新处理 + if _, exists := task.RuningTaskPath[filePath]; !exists { + log.Println("扫描到超时scanned需重新执行文件: " + filePath) + task.HandleTask(path + "/" + fileNameWithoutExt) + } + } + } + } + } + + dirRead.Close() + } +} + +func (task *Task) HandleTask(path string) { + filePath := path + ".scanned" + + task.RuningTaskPath[path] = 1 + + log.Println("开始处理文件: " + filePath) + + // 判断文件是否存在 + _, err := os.Stat(filePath) + if os.IsNotExist(err) { + log.Println(filePath + " 文件不存在") + return + } + + // 打开文件 + file, err := os.Open(filePath) + + if err != nil { + log.Println(filePath + " 文件打开失败: " + err.Error()) + return + } + + // 自定义缓冲区大小为 20 * 1MB + const bufferSize = 20 * 1024 * 1024 + // 逐行读取文件内容 + scanner := bufio.NewScanner(file) + // 设置缓冲区大小 + buf := make([]byte, bufferSize) + scanner.Buffer(buf, bufferSize) + + i := 0 // 消息数计数器 + + wg := sync.WaitGroup{} + + lineLimitChan := make(chan bool, 10) // 限制协程数量 + for scanner.Scan() { + lineLimitChan <- true + line := scanner.Text() + if task.TaskFunc != nil { + wg.Add(1) + go func() { + task.TaskFunc(line, task) + <-lineLimitChan + wg.Done() + }() + } + i = i + 1 + } + + wg.Wait() + + if err := scanner.Err(); err != nil { + err := os.Rename(filePath, path+".error") + if err != nil { + log.Println(filePath + " 文件重命名失败: " + err.Error()) + } + log.Println(filePath + " 文件扫描出错: " + err.Error()) + return + } + + // 文件扫描完毕 + file.Close() + + log.Printf(filePath+" 文件处理完毕, 共处理[%d]条信息", i) + + err = os.Remove(filePath) + if err != nil { + log.Println(filePath + " 文件删除失败: " + err.Error()) + return + } + + log.Println(filePath + " 文件处理完成后删除") + delete(task.RuningTaskPath, path) +} diff --git a/utils/face/mqtt/mqtt.go b/utils/face/mqtt/mqtt.go new file mode 100644 index 0000000..754379a --- /dev/null +++ b/utils/face/mqtt/mqtt.go @@ -0,0 +1,135 @@ +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) + } + }() + } +} diff --git a/utils/go.mod b/utils/go.mod new file mode 100644 index 0000000..400e558 --- /dev/null +++ b/utils/go.mod @@ -0,0 +1,11 @@ +module gitee.com/jhx1219/machine-go/utils + +go 1.19 + +require github.com/eclipse/paho.mqtt.golang v1.4.2 + +require ( + github.com/gorilla/websocket v1.4.2 // indirect + golang.org/x/net v0.0.0-20200425230154-ff2c4b7c35a0 // indirect + golang.org/x/sync v0.0.0-20210220032951-036812b2e83c // indirect +) diff --git a/utils/go.sum b/utils/go.sum new file mode 100644 index 0000000..ddd8ce6 --- /dev/null +++ b/utils/go.sum @@ -0,0 +1,12 @@ +github.com/eclipse/paho.mqtt.golang v1.4.2 h1:66wOzfUHSSI1zamx7jR6yMEI5EuHnT1G6rNA5PM12m4= +github.com/eclipse/paho.mqtt.golang v1.4.2/go.mod h1:JGt0RsEwEX+Xa/agj90YJ9d9DH2b7upDZMK9HRbFvCA= +github.com/gorilla/websocket v1.4.2 h1:+/TMaTYc4QFitKJxsQ7Yye35DkWvkdLcvGKqM+x0Ufc= +github.com/gorilla/websocket v1.4.2/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= +golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= +golang.org/x/net v0.0.0-20200425230154-ff2c4b7c35a0 h1:Jcxah/M+oLZ/R4/z5RzfPzGbPXnVDPkEDtf2JnuxN+U= +golang.org/x/net v0.0.0-20200425230154-ff2c4b7c35a0/go.mod h1:qpuaurCH72eLCgpAm/N6yyVIVM9cpaDIP3A8BGJEC5A= +golang.org/x/sync v0.0.0-20210220032951-036812b2e83c h1:5KslGYwFpkhGh+Q16bwMP3cOontH8FOep7tGV86Y7SQ= +golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= +golang.org/x/sys v0.0.0-20200323222414-85ca7c5b95cd/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=