Files
bj_power/bj_power_workstation/internal/syncer/syncer.go
T
2026-08-28 15:06:01 +08:00

142 lines
3.6 KiB
Go

package syncer
import (
"context"
"encoding/json"
"time"
"bj_power_workstation/internal/mes"
"bj_power_workstation/internal/store"
"github.com/zeromicro/go-zero/core/logx"
)
const (
reportBatchLimit = 200 // 单轮最多处理条数
defaultInterval = 15 * time.Second
)
// Syncer 后台同步器:每 15 秒把未同步的 report_queue 与 torque_results 上报 MES。
type Syncer struct {
st *store.Store
mes *mes.Client
interval time.Duration
}
func New(st *store.Store, mc *mes.Client) *Syncer {
return &Syncer{st: st, mes: mc, interval: defaultInterval}
}
// Run 周期冲刷循环,随 ctx 取消退出。
func (s *Syncer) Run(ctx context.Context) {
logx.Infof("后台同步协程已启动:每 %s 冲刷一次未上报数据", s.interval)
ticker := time.NewTicker(s.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
logx.Info("后台同步协程已停止")
return
case <-ticker.C:
if n := s.RunOnce(ctx); n > 0 {
logx.Infof("本轮同步成功 %d 条", n)
}
}
}
}
// RunOnce 立即执行一轮同步(先队列后拧紧结果),返回本次成功条数。
func (s *Syncer) RunOnce(ctx context.Context) int {
n := s.flushQueue(ctx)
n += s.flushTorque(ctx)
return n
}
// flushTorque 上报未同步拧紧结果:
// POST /api/internal/tightening/report;成功标记 synced=1,失败保持 0 等待下轮重试。
func (s *Syncer) flushTorque(ctx context.Context) int {
rows, err := s.st.ListTorqueResults("", true, reportBatchLimit)
if err != nil {
logx.Errorf("查询未同步拧紧结果失败: %v", err)
return 0
}
var done []int64
count := 0
for _, t := range rows {
select {
case <-ctx.Done():
return count
default:
}
body := map[string]any{
"workOrderNo": t.WorkOrderNo,
"sn": t.Sn,
"dockCode": t.DockCode,
"screwNo": t.ScrewNo,
"torque": t.Torque,
"angle": t.Angle,
"result": t.Result,
"operator": t.Operator,
}
if _, _, err := s.mes.Post("/api/internal/tightening/report", body); err != nil {
logx.Errorf("上报拧紧结果失败(id=%d sn=%s): %v", t.ID, t.Sn, err)
continue
}
done = append(done, t.ID)
count++
}
if len(done) > 0 {
if err := s.st.MarkTorqueSynced(done); err != nil {
logx.Errorf("标记拧紧结果已同步失败: %v", err)
}
}
return count
}
// flushQueue 上报未同步的离线缓存队列,payload 原样透传:
// - kind=process_done → POST /api/internal/station/done
// - kind=temp_store → POST /api/internal/station/checkin
//
// 成功标记 synced=1;失败 synced 保持 0 并累计 retry_count。
func (s *Syncer) flushQueue(ctx context.Context) int {
items, err := s.st.PendingQueueItems(reportBatchLimit)
if err != nil {
logx.Errorf("查询待上报队列失败: %v", err)
return 0
}
count := 0
for _, it := range items {
select {
case <-ctx.Done():
return count
default:
}
var path string
switch it.Kind {
case "process_done":
path = "/api/internal/station/done"
case "temp_store":
path = "/api/internal/station/checkin"
default:
logx.Errorf("未知队列类型(kind=%s id=%d),跳过", it.Kind, it.ID)
_ = s.st.BumpQueueRetry(it.ID)
continue
}
if _, _, err := s.mes.Post(path, json.RawMessage(it.Payload)); err != nil {
logx.Errorf("上报队列失败(kind=%s id=%d 第 %d 次): %v", it.Kind, it.ID, it.RetryCount+1, err)
_ = s.st.BumpQueueRetry(it.ID)
continue
}
if err := s.st.MarkQueueSynced(it.ID); err != nil {
logx.Errorf("标记队列已同步失败(id=%d): %v", it.ID, err)
continue
}
count++
}
return count
}