diff --git a/.gitignore b/.gitignore index 0d7a819..b43bb3a 100644 --- a/.gitignore +++ b/.gitignore @@ -22,6 +22,7 @@ setup !front/tiny-engine-platform/env/.env.production !front/tiny-engine-platform/env/.env.production.example tmp/ +.cache/ *.env.local *.secret.* *.local.yaml diff --git a/AGENTS.md b/AGENTS.md index bf5dbbd..5a9d319 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -13,7 +13,7 @@ - 验收范围:`docs/12-验收与迭代规划.md` 2. 检查 `git status`,保留用户已有改动;不要顺手格式化、重写或删除无关文件。 3. 区分三类内容: - - `frontend/platform_admin`、`backend/{api,worker,iot}` 是当前业务实现。 + - `frontend/platform_admin`、`backend/{api,worker,iot-client,iot-server}` 是当前业务实现。 - `docs/10-技术实现规划.md` 中尚不存在的目录和系统属于规划,不得为完成局部任务擅自创建空壳工程。 4. 需求不明确时以文档中的系统主责、数据主责和安全规则为边界;不得把标为“待确认”的事项自行固化为政策。 5. 文档之间出现冲突时按以下顺序处理: @@ -29,7 +29,8 @@ - `backend/api`:Gin HTTP API/BFF、鉴权、同步业务事务、数据库事实和审计。 - `backend/worker`:Outbox 投递、Redis Streams 消费、告警、派单、通知、对账和超时任务。当前主要是 Mock 边界,未确定持久化消息契约前不要伪造完整实现。 -- `backend/iot`:MQTT 会话、协议适配、遥测、设备命令和回执。当前使用 Mock 适配,不得假定已连接真实设备或 Broker。 +- `backend/iot-client`:管理系统/App 的设备上下行 HTTP 边界;不得绕过 API 的业务鉴权和安全状态校验。 +- `backend/iot-server`:外部 MQTT Broker 会话、厂商协议、遥测、设备命令和回执;本地配置不代表已完成真实设备联调。 - `frontend/platform_admin`:平台级治理、运营、财务、安全、主数据和审计界面;后端资源契约是其资源和操作能力的依据。 - `docs`:产品和技术基线;涉及业务口径、实体、状态、权限或跨端流程的变更必须同步更新相应文档。 - `doc`:原始设计图和协议附件,仅作需求来源;不要批量改名、压缩或重写。 @@ -39,7 +40,7 @@ - 平台总后台不承接气站、配送点、生产单位的日常重复操作。 - 气站、配送点、生产、API 中心的业务事实各有唯一主责系统;其他系统通过 API、事件或只读聚合使用。 - API 中心只管理接口产品和调用方,不能绕过业务服务直接写业务数据库。 -- `api` 不应承担长期运行的异步循环;`worker` 不应复制同步领域事务;`iot` 不应内嵌交易、资金或后台页面逻辑。 +- `api` 不应承担长期运行的异步循环;`worker` 不应复制同步领域事务;两个 IoT 进程不应内嵌交易、资金或后台页面逻辑。 ### 目录归属规则 @@ -59,7 +60,8 @@ platforms/ backend/ api/ # 当前:同步 HTTP API/BFF worker/ # 当前:异步任务和事件消费 - iot/ # 当前:MQTT 与设备协议边界 + iot-client/ # 当前:管理系统/App 的设备上下行接口 + iot-server/ # 当前:MQTT 与设备协议边界 contracts/ openapi/ # 规划:HTTP 契约 asyncapi/ # 规划:MQTT/事件契约 @@ -145,7 +147,8 @@ go build ./cmd/main/main.go # Worker / IoT(改动对应模块时) cd backend/worker && go test ./... && go build ./cmd/main/main.go -cd backend/iot && go test ./... && go build ./cmd/main/main.go +cd backend/iot-client && go test ./... && go build ./cmd/main/main.go +cd backend/iot-server && go test ./... && go build ./cmd/main/main.go # 平台总后台 cd frontend/platform_admin diff --git a/backend/README.md b/backend/README.md index 3608e6f..af2d29a 100644 --- a/backend/README.md +++ b/backend/README.md @@ -5,7 +5,8 @@ | 目录 | 进程与职责 | | --- | --- | | `api` | Gin HTTP API / BFF、同步事务、JWT 鉴权与统一响应 | -| `worker` | Redis Streams 消费、Outbox 投递、超时扫描与 Mock 外部适配 | -| `iot` | MQTT 协议适配边界、遥测/命令契约校验与 Mock 设备接入 | +| `worker` | Outbox 投递、超时扫描与可重试异步任务 | +| `iot-client` | 管理系统/App 的设备上下行 HTTP 边界、幂等受理与状态查询 | +| `iot-server` | 外部 MQTT Broker 会话、厂商二进制协议、遥测、命令与回执 | -三个进程沿用统一的工程机制;领域模型使用仓库约定的自定义 `models.Entity`,并遵守 `id` 自增主键与 `identity` 对外身份分离规则。 +四个进程沿用统一的工程机制;API 是命令、Outbox 和上行报文的数据库事实主责,Worker 负责投递,IoT Client 负责系统接口,IoT Server 负责设备协议。 diff --git a/backend/api/etc/heqi_dev.yaml b/backend/api/etc/heqi_dev.yaml index f80ff09..ccb20e2 100644 --- a/backend/api/etc/heqi_dev.yaml +++ b/backend/api/etc/heqi_dev.yaml @@ -47,3 +47,7 @@ Payment: OfficialAccountAppID: "" MiniProgramAppID: "" NotifyURL: "https://example.invalid/heqi/payment-return/v1/wechat/notify" + +IoT: + InternalServiceToken: change-me-iot-internal-token + ClientBaseURL: http://127.0.0.1:12429 diff --git a/backend/api/internal/config/config.go b/backend/api/internal/config/config.go index 6c850bc..1312244 100644 --- a/backend/api/internal/config/config.go +++ b/backend/api/internal/config/config.go @@ -66,6 +66,12 @@ type PaymentConfig struct { Wechat WechatPayConfig `yaml:"Wechat"` } +// IoTConfig 保存 API、Worker 与 IoT Client 间的内部认证和地址。 +type IoTConfig struct { + InternalServiceToken string `yaml:"InternalServiceToken"` + ClientBaseURL string `yaml:"ClientBaseURL"` +} + // SrvConfig 与仓库现有进程的配置结构保持一致。 type SrvConfig struct { conf.Base `yaml:",inline"` @@ -75,6 +81,7 @@ type SrvConfig struct { Global GlobalConfig `yaml:"Global"` Wallet WalletConfig `yaml:"-"` Payment PaymentConfig `yaml:"Payment"` + IoT IoTConfig `yaml:"IoT"` } // New 初始化 BSM 配置并校验服务监听地址。 @@ -103,6 +110,9 @@ func New(srvKey string) { if Spec.Payment.ExpireMinutes <= 0 || Spec.Payment.RefundWindowDays <= 0 { panic("Payment expiration and refund window must be greater than zero") } + if len(strings.TrimSpace(Spec.IoT.InternalServiceToken)) < 16 || !strings.HasPrefix(Spec.IoT.ClientBaseURL, "http") { + panic("IoT internal token and client base URL are required") + } registerURL, err := url.ParseRequestURI(Spec.Global.UserRegisterURL) if err != nil || (registerURL.Scheme != "http" && registerURL.Scheme != "https") || registerURL.Host == "" { panic("Global.UserRegisterURL must be a valid HTTP or HTTPS URL") diff --git a/backend/api/internal/logic/iot/iot.go b/backend/api/internal/logic/iot/iot.go new file mode 100644 index 0000000..64132dd --- /dev/null +++ b/backend/api/internal/logic/iot/iot.go @@ -0,0 +1,197 @@ +// Package iot 实现设备命令、Outbox 认领和上行报文持久化。 +package iot + +import ( + "crypto/subtle" + "encoding/json" + "net/http" + "strings" + "time" + + "git.apinb.com/heqiapp/platforms/backend/api/internal/config" + "git.apinb.com/heqiapp/platforms/backend/api/internal/impl" + "git.apinb.com/heqiapp/platforms/backend/api/internal/models" + "github.com/gin-gonic/gin" + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +type commandRequest struct { + DeviceIdentity string `json:"device_identity"` + DeviceID string `json:"device_id"` + IdempotencyKey string `json:"idempotency_key"` + Action string `json:"action"` + ExpiresAt time.Time `json:"expires_at"` + KeyID byte `json:"key_id"` + DeviceKind byte `json:"device_kind"` + DeviceType byte `json:"device_type"` + DeviceModel [3]byte `json:"device_model"` + Controller byte `json:"controller"` + Loop byte `json:"loop"` + Component byte `json:"component"` +} + +func RequireInternalToken() gin.HandlerFunc { + return func(ctx *gin.Context) { + actual := ctx.GetHeader("X-Heqi-Iot-Token") + expected := config.Spec.IoT.InternalServiceToken + if len(actual) != len(expected) || subtle.ConstantTimeCompare([]byte(actual), []byte(expected)) != 1 { + ctx.AbortWithStatusJSON(http.StatusUnauthorized, gin.H{"code": "IOT_UNAUTHORIZED"}) + return + } + ctx.Next() + } +} + +func CreateCommand(ctx *gin.Context) { + var request commandRequest + if err := ctx.ShouldBindJSON(&request); err != nil || request.DeviceIdentity == "" || len(request.DeviceID) != 16 || request.IdempotencyKey == "" || (request.Action != "open_valve" && request.Action != "close_valve") { + ctx.JSON(422, gin.H{"code": "IOT_COMMAND_INVALID"}) + return + } + if request.ExpiresAt.IsZero() { + request.ExpiresAt = time.Now().Add(30 * time.Second) + } + var existing models.IotCommand + if err := impl.DBService.Where("idempotency_key = ?", request.IdempotencyKey).First(&existing).Error; err == nil { + ctx.JSON(200, existing) + return + } else if err != gorm.ErrRecordNotFound { + ctx.JSON(500, gin.H{"code": "IOT_COMMAND_QUERY_FAILED"}) + return + } + payload, _ := json.Marshal(request) + command := models.IotCommand{Entity: models.Entity{Identity: models.NewIdentity()}, DeviceIdentity: request.DeviceIdentity, DeviceID: request.DeviceID, IdempotencyKey: request.IdempotencyKey, Action: request.Action, RequestPayload: string(payload), CommandStatus: "accepted", ExpiresAt: request.ExpiresAt} + outbox := models.IotOutbox{Identity: models.NewIdentity(), CommandIdentity: command.Identity, EventType: "iot.command.accepted", Payload: string(payload), OutboxStatus: "pending", AvailableAt: time.Now()} + err := impl.DBService.Transaction(func(tx *gorm.DB) error { + if err := tx.Create(&command).Error; err != nil { + return err + } + return tx.Create(&outbox).Error + }) + if err != nil { + ctx.JSON(500, gin.H{"code": "IOT_COMMAND_CREATE_FAILED"}) + return + } + ctx.JSON(http.StatusAccepted, command) +} + +func GetCommand(ctx *gin.Context) { + var command models.IotCommand + if err := impl.DBService.Where("identity = ?", ctx.Param("identity")).First(&command).Error; err != nil { + ctx.JSON(404, gin.H{"code": "IOT_COMMAND_NOT_FOUND"}) + return + } + ctx.JSON(200, command) +} + +func ClaimOutbox(ctx *gin.Context) { + var outbox models.IotOutbox + err := impl.DBService.Transaction(func(tx *gorm.DB) error { + now := time.Now() + err := tx.Clauses(clause.Locking{Strength: "UPDATE", Options: "SKIP LOCKED"}).Where("(outbox_status IN ? AND available_at <= ?) OR (outbox_status = ? AND updated_at < ?)", []string{"pending", "retry"}, now, "processing", now.Add(-time.Minute)).Order("id").First(&outbox).Error + if err != nil { + return err + } + return tx.Model(&outbox).Updates(map[string]any{"outbox_status": "processing", "attempts": gorm.Expr("attempts + 1"), "updated_at": time.Now()}).Error + }) + if err == gorm.ErrRecordNotFound { + ctx.Status(http.StatusNoContent) + return + } + if err != nil { + ctx.JSON(500, gin.H{"code": "IOT_OUTBOX_CLAIM_FAILED"}) + return + } + ctx.JSON(200, outbox) +} + +func CompleteOutbox(ctx *gin.Context) { + var request struct { + Success bool `json:"success"` + PacketNumber uint16 `json:"packet_number"` + ErrorCode string `json:"error_code"` + } + if ctx.ShouldBindJSON(&request) != nil { + ctx.JSON(400, gin.H{"code": "IOT_OUTBOX_RESULT_INVALID"}) + return + } + var outbox models.IotOutbox + if err := impl.DBService.Where("identity = ?", ctx.Param("identity")).First(&outbox).Error; err != nil { + ctx.JSON(404, gin.H{"code": "IOT_OUTBOX_NOT_FOUND"}) + return + } + now := time.Now() + err := impl.DBService.Transaction(func(tx *gorm.DB) error { + if request.Success { + if err := tx.Model(&outbox).Updates(map[string]any{"outbox_status": "published", "updated_at": now}).Error; err != nil { + return err + } + return tx.Model(&models.IotCommand{}).Where("identity = ?", outbox.CommandIdentity).Updates(map[string]any{"command_status": "pending_confirmation", "packet_number": request.PacketNumber, "dispatched_at": now, "updated_at": now}).Error + } + delay := time.Duration(outbox.Attempts+1) * time.Second + if delay > time.Minute { + delay = time.Minute + } + if err := tx.Model(&outbox).Updates(map[string]any{"outbox_status": "retry", "available_at": now.Add(delay), "updated_at": now}).Error; err != nil { + return err + } + return tx.Model(&models.IotCommand{}).Where("identity = ?", outbox.CommandIdentity).Updates(map[string]any{"command_status": "dispatch_failed", "error_code": request.ErrorCode, "updated_at": now}).Error + }) + if err != nil { + ctx.JSON(500, gin.H{"code": "IOT_OUTBOX_RESULT_FAILED"}) + return + } + ctx.Status(204) +} + +func SaveDeviceMessage(ctx *gin.Context) { + data, err := ctx.GetRawData() + if err != nil || len(data) == 0 { + ctx.JSON(400, gin.H{"code": "IOT_MESSAGE_INVALID"}) + return + } + var envelope struct { + Type, Topic, DeviceID, ReceivedAt, PayloadHex string + Frame json.RawMessage `json:"frame"` + } + if json.Unmarshal(data, &envelope) != nil || envelope.ReceivedAt == "" { + ctx.JSON(422, gin.H{"code": "IOT_MESSAGE_INVALID"}) + return + } + received, err := time.Parse(time.RFC3339Nano, envelope.ReceivedAt) + if err != nil { + ctx.JSON(422, gin.H{"code": "IOT_RECEIVED_AT_INVALID"}) + return + } + message := models.IotDeviceMessage{Identity: models.NewIdentity(), DeviceID: envelope.DeviceID, Topic: envelope.Topic, MessageType: envelope.Type, PayloadHex: strings.ToLower(envelope.PayloadHex), DecodedFrame: string(envelope.Frame), ReceivedAt: received} + var frame struct { + Control byte `json:"Control"` + PacketNumber uint16 `json:"PacketNumber"` + DeviceTime string `json:"DeviceTime"` + } + _ = json.Unmarshal(envelope.Frame, &frame) + if frame.DeviceTime != "" { + if occurred, parseErr := time.Parse(time.RFC3339, frame.DeviceTime); parseErr == nil { + message.DeviceOccurredAt = &occurred + } + } + err = impl.DBService.Transaction(func(tx *gorm.DB) error { + if createErr := tx.Create(&message).Error; createErr != nil { + return createErr + } + if envelope.DeviceID == "" || frame.PacketNumber == 0 || frame.Control&0x08 == 0 { + return nil + } + status, errorCode := "succeeded", "" + if frame.Control&0x04 != 0 { + status, errorCode = "failed", "DEVICE_REPORTED_FAILURE" + } + return tx.Model(&models.IotCommand{}).Where("device_id = ? AND packet_number = ? AND command_status IN ?", envelope.DeviceID, frame.PacketNumber, []string{"dispatched", "pending_confirmation"}).Updates(map[string]any{"command_status": status, "error_code": errorCode, "acknowledged_at": received, "updated_at": received}).Error + }) + if err != nil { + ctx.JSON(500, gin.H{"code": "IOT_MESSAGE_SAVE_FAILED"}) + return + } + ctx.JSON(http.StatusAccepted, gin.H{"identity": message.Identity}) +} diff --git a/backend/api/internal/models/iot_command.go b/backend/api/internal/models/iot_command.go new file mode 100644 index 0000000..6c00f1f --- /dev/null +++ b/backend/api/internal/models/iot_command.go @@ -0,0 +1,64 @@ +package models + +import ( + "time" + + "git.apinb.com/bsm-sdk/core/database" +) + +// IotCommand 是设备下行命令事实;受理、投递与设备执行状态必须分离。 +type IotCommand struct { + Entity `gorm:"embedded;comment:设备命令公共实体字段"` // 命令公共实体字段 + DeviceIdentity string `gorm:"column:device_identity;type:varchar(36);not null;index" json:"device_identity"` // 平台设备 identity + DeviceID string `gorm:"column:device_id;type:varchar(16);not null;index" json:"device_id"` // 厂商 8-byte BCD 标识 + IdempotencyKey string `gorm:"column:idempotency_key;type:varchar(128);not null;uniqueIndex;comment:命令幂等键" json:"idempotency_key"` // 命令幂等键 + Action string `gorm:"column:action;type:varchar(32);not null;comment:设备控制动作" json:"action"` // 设备控制动作 + RequestPayload string `gorm:"column:request_payload;type:text;not null;comment:原始命令请求快照" json:"request_payload"` // 原始请求快照 + CommandStatus string `gorm:"column:command_status;type:varchar(32);not null;index;comment:命令受理投递执行状态" json:"command_status"` // 命令处理状态 + PacketNumber uint16 `gorm:"column:packet_number;not null;default:0;comment:厂商协议包号" json:"packet_number"` // 厂商协议包号 + ErrorCode string `gorm:"column:error_code;type:varchar(64);not null;default:'';comment:稳定错误码" json:"error_code"` // 稳定错误码 + ExpiresAt time.Time `gorm:"column:expires_at;type:timestamptz;not null;index;comment:命令失效时间" json:"expires_at"` // 命令失效时间 + DispatchedAt *time.Time `gorm:"column:dispatched_at;type:timestamptz;comment:下发到消息代理时间" json:"dispatched_at"` // 下发时间 + AcknowledgedAt *time.Time `gorm:"column:acknowledged_at;type:timestamptz;comment:设备回执时间" json:"acknowledged_at"` // 设备回执时间 +} + +func (*IotCommand) TableName() string { return "iot_command" } + +// IotOutbox 保证命令事实与待投递事件在同一数据库事务提交。 +type IotOutbox struct { + ID uint64 `gorm:"column:id;primaryKey;autoIncrement;comment:数据库自增主键" json:"-"` // 数据库主键 + Identity string `gorm:"column:identity;type:varchar(36);not null;uniqueIndex;comment:Outbox业务标识" json:"identity"` // Outbox业务标识 + CommandIdentity string `gorm:"column:command_identity;type:varchar(36);not null;uniqueIndex;comment:设备命令业务标识" json:"command_identity"` // 命令业务标识 + EventType string `gorm:"column:event_type;type:varchar(64);not null;comment:事件类型" json:"event_type"` // 事件类型 + Payload string `gorm:"column:payload;type:text;not null;comment:待投递事件载荷" json:"payload"` // 事件载荷 + OutboxStatus string `gorm:"column:outbox_status;type:varchar(24);not null;index;comment:Outbox投递状态" json:"outbox_status"` // 投递状态 + Attempts int `gorm:"column:attempts;not null;default:0;comment:投递尝试次数" json:"attempts"` // 尝试次数 + AvailableAt time.Time `gorm:"column:available_at;type:timestamptz;not null;index;comment:下次可投递时间" json:"available_at"` // 可投递时间 + CreatedAt time.Time `gorm:"column:created_at;type:timestamptz;not null;index;comment:创建时间" json:"created_at"` // 创建时间 + UpdatedAt time.Time `gorm:"column:updated_at;type:timestamptz;not null;comment:更新时间" json:"updated_at"` // 更新时间 +} + +func (*IotOutbox) TableName() string { return "iot_outbox" } + +// IotDeviceMessage 保存原始上行/回执和服务端接收时间,原始事实不可由前端修改。 +type IotDeviceMessage struct { + ID uint64 `gorm:"column:id;primaryKey;autoIncrement;comment:数据库自增主键" json:"-"` // 数据库主键 + Identity string `gorm:"column:identity;type:varchar(36);not null;uniqueIndex;comment:上行消息业务标识" json:"identity"` // 上行消息标识 + DeviceID string `gorm:"column:device_id;type:varchar(16);not null;index;comment:厂商设备标识" json:"device_id"` // 厂商设备标识 + Topic string `gorm:"column:topic;type:varchar(255);not null;comment:接收消息主题" json:"topic"` // MQTT主题 + MessageType string `gorm:"column:message_type;type:varchar(32);not null;index;comment:上行或回执消息类型" json:"message_type"` // 消息类型 + PayloadHex string `gorm:"column:payload_hex;type:text;not null;comment:原始报文十六进制" json:"payload_hex"` // 原始报文 + DecodedFrame string `gorm:"column:decoded_frame;type:text;not null;comment:协议解析结果快照" json:"decoded_frame"` // 解析快照 + DeviceOccurredAt *time.Time `gorm:"column:device_occurred_at;type:timestamptz;comment:设备原始采集时间" json:"device_occurred_at"` // 设备采集时间 + ReceivedAt time.Time `gorm:"column:received_at;type:timestamptz;not null;index;comment:服务端接收时间" json:"received_at"` // 服务端接收时间 + CreatedAt time.Time `gorm:"column:created_at;type:timestamptz;not null;index;comment:创建时间" json:"created_at"` // 创建时间 + UpdatedAt time.Time `gorm:"column:updated_at;type:timestamptz;not null;comment:更新时间" json:"updated_at"` // 更新时间 +} + +func (*IotDeviceMessage) TableName() string { return "iot_device_message" } + +func init() { + database.AppendMigrate(&IotCommand{}) + database.AppendMigrate(&IotOutbox{}) + database.AppendMigrate(&IotDeviceMessage{}) +} diff --git a/backend/api/internal/routers/iot.go b/backend/api/internal/routers/iot.go new file mode 100644 index 0000000..aa929ad --- /dev/null +++ b/backend/api/internal/routers/iot.go @@ -0,0 +1,17 @@ +package routers + +import ( + "fmt" + iotlogic "git.apinb.com/heqiapp/platforms/backend/api/internal/logic/iot" + "github.com/gin-gonic/gin" +) + +func RegisterIoT(serviceKey string, engine *gin.Engine) { + group := engine.Group(fmt.Sprintf("/%s/internal/v1/iot", serviceKey)) + group.Use(iotlogic.RequireInternalToken()) + group.POST("/commands", iotlogic.CreateCommand) + group.GET("/commands/:identity", iotlogic.GetCommand) + group.GET("/outbox/next", iotlogic.ClaimOutbox) + group.POST("/outbox/:identity/result", iotlogic.CompleteOutbox) + group.POST("/device-messages", iotlogic.SaveDeviceMessage) +} diff --git a/backend/api/internal/routers/register.go b/backend/api/internal/routers/register.go index 443606b..034871e 100644 --- a/backend/api/internal/routers/register.go +++ b/backend/api/internal/routers/register.go @@ -10,5 +10,6 @@ func Register(serviceKey string, engine *gin.Engine) { RegisterDelivery(serviceKey, engine) RegisterClient(serviceKey, engine) RegisterPaymentReturn(serviceKey, engine) + RegisterIoT(serviceKey, engine) registerUploadRoute(serviceKey, engine) } diff --git a/backend/go.work b/backend/go.work index d0f4f71..6e604f8 100644 --- a/backend/go.work +++ b/backend/go.work @@ -3,5 +3,6 @@ go 1.26.1 use ( ./api ./worker - ./iot + ./iot-client + ./iot-server ) diff --git a/backend/go.work.sum b/backend/go.work.sum index 233d4c5..7a636ae 100644 --- a/backend/go.work.sum +++ b/backend/go.work.sum @@ -13,6 +13,8 @@ github.com/cncf/xds/go v0.0.0-20260202195803-dba9d589def2 h1:aBangftG7EVZoUb69Os github.com/cncf/xds/go v0.0.0-20260202195803-dba9d589def2/go.mod h1:qwXFYgsP6T7XnJtbKlf1HP8AjxZZyzxMmc+Lq5GjlU4= github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= +github.com/eclipse/paho.mqtt.golang v1.5.1 h1:/VSOv3oDLlpqR2Epjn1Q7b2bSTplJIeV2ISgCl2W7nE= +github.com/eclipse/paho.mqtt.golang v1.5.1/go.mod h1:1/yJCneuyOoCOzKSsOTUc0AJfpsItBGWvYpBLimhArU= github.com/elastic/elastic-transport-go/v8 v8.9.0 h1:KeT/2P54F0xS0S8Y3Pf+tFDg4HmBgReQMB+BMz8dDAs= github.com/elastic/elastic-transport-go/v8 v8.9.0/go.mod h1:ssMTvNS2hwf7CaiGsRRsx4gQHFZ/jS/DkLcISxekWzc= github.com/elastic/go-elasticsearch/v9 v9.4.1 h1:pEF8xlnL8D2WdVp4HRHhGNUwC5QHdgk0DZKrJY0WCNs= @@ -33,6 +35,8 @@ github.com/godbus/dbus/v5 v5.0.4 h1:9349emZab16e7zQvpmsbtjc18ykshndd8y2PG3sgJbA= github.com/golang/glog v1.2.5 h1:DrW6hGnjIhtvhOIiAKT6Psh/Kd/ldepEa81DKeiRJ5I= github.com/golang/glog v1.2.5/go.mod h1:6AhwSGph0fcJtXVM/PEHPqZlFeoLxhs7/t5UDAwmO+w= github.com/google/gofuzz v1.0.0 h1:A8PeW59pxE9IoFRqBp37U+mSNaQoZ46F1f0f863XSXw= +github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= +github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= github.com/grpc-ecosystem/go-grpc-middleware/providers/prometheus v1.0.1 h1:qnpSQwGEnkcRpTqNOIR6bJbR0gAorgP9CSALpRcKoAA= github.com/grpc-ecosystem/go-grpc-middleware/providers/prometheus v1.0.1/go.mod h1:lXGCsh6c22WGtjr+qGHj1otzZpV/1kwTMAqkwZsnWRU= github.com/grpc-ecosystem/go-grpc-middleware/v2 v2.1.0 h1:pRhl55Yx1eC7BZ1N+BBWwnKaMyD8uC+34TLdndZMAKk= diff --git a/backend/iot-client/README.md b/backend/iot-client/README.md new file mode 100644 index 0000000..1e9bfdb --- /dev/null +++ b/backend/iot-client/README.md @@ -0,0 +1,5 @@ +# IoT Client + +管理系统与 App 的上下行设备接口边界。外部业务请求应先经过 Platform API 的 JWT、角色、对象归属和安全状态校验;API/Worker 使用内部令牌调用本服务。命令要求 `identity`、`idempotency_key`、过期时间,并返回可查询状态,不能把 HTTP 受理视为设备执行成功。 + +设备上行报文由 IoT Server 回调本服务,再转交 Platform API 持久化和审计。 diff --git a/backend/iot-client/cmd/main/main.go b/backend/iot-client/cmd/main/main.go new file mode 100644 index 0000000..ed77a0a --- /dev/null +++ b/backend/iot-client/cmd/main/main.go @@ -0,0 +1,25 @@ +package main + +import ( + "context" + "flag" + "git.apinb.com/heqiapp/platforms/backend/iot-client/internal/config" + "git.apinb.com/heqiapp/platforms/backend/iot-client/internal/service" + "log" + "os/signal" + "syscall" +) + +func main() { + path := flag.String("config", "etc/platform_iot_client_dev.yaml", "YAML 配置文件") + flag.Parse() + cfg, err := config.Load(*path) + if err != nil { + log.Fatal(err) + } + ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer cancel() + if err = service.New(cfg).Run(ctx); err != nil { + log.Fatal(err) + } +} diff --git a/backend/iot-client/etc/platform_iot_client_dev.yaml b/backend/iot-client/etc/platform_iot_client_dev.yaml new file mode 100644 index 0000000..fc4f838 --- /dev/null +++ b/backend/iot-client/etc/platform_iot_client_dev.yaml @@ -0,0 +1,8 @@ +Service: platform-iot-client +HTTP: + Address: 127.0.0.1:12429 + InternalToken: change-me-iot-internal-token +Upstream: + IoTServerURL: http://127.0.0.1:12428 + PlatformAPIURL: http://127.0.0.1:12426 + Token: change-me-iot-internal-token diff --git a/backend/iot-client/go.mod b/backend/iot-client/go.mod new file mode 100644 index 0000000..bd440ab --- /dev/null +++ b/backend/iot-client/go.mod @@ -0,0 +1,5 @@ +module git.apinb.com/heqiapp/platforms/backend/iot-client + +go 1.26.1 + +require gopkg.in/yaml.v3 v3.0.1 diff --git a/backend/iot-client/go.sum b/backend/iot-client/go.sum new file mode 100644 index 0000000..6497cf3 --- /dev/null +++ b/backend/iot-client/go.sum @@ -0,0 +1,4 @@ +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= diff --git a/backend/iot-client/internal/config/config.go b/backend/iot-client/internal/config/config.go new file mode 100644 index 0000000..69ef69c --- /dev/null +++ b/backend/iot-client/internal/config/config.go @@ -0,0 +1,44 @@ +package config + +import ( + "fmt" + "gopkg.in/yaml.v3" + "os" +) + +type HTTP struct { + Address string `yaml:"Address"` + InternalToken string `yaml:"InternalToken"` +} +type Upstream struct { + IoTServerURL string `yaml:"IoTServerURL"` + PlatformAPIURL string `yaml:"PlatformAPIURL"` + Token string `yaml:"Token"` +} +type Config struct { + Service string `yaml:"Service"` + HTTP HTTP `yaml:"HTTP"` + Upstream Upstream `yaml:"Upstream"` +} + +func Load(path string) (Config, error) { + var cfg Config + data, err := os.ReadFile(path) + if err != nil { + return cfg, err + } + if err = yaml.Unmarshal(data, &cfg); err != nil { + return cfg, err + } + override(&cfg.HTTP.InternalToken, "HEQI_IOT_INTERNAL_TOKEN") + override(&cfg.Upstream.Token, "HEQI_IOT_INTERNAL_TOKEN") + if cfg.HTTP.Address == "" || cfg.HTTP.InternalToken == "" || cfg.Upstream.IoTServerURL == "" || cfg.Upstream.PlatformAPIURL == "" { + return cfg, fmt.Errorf("HTTP 和 Upstream 配置不完整") + } + return cfg, nil +} +func override(target *string, name string) { + if value := os.Getenv(name); value != "" { + *target = value + } +} diff --git a/backend/iot-client/internal/service/service.go b/backend/iot-client/internal/service/service.go new file mode 100644 index 0000000..b1e30ec --- /dev/null +++ b/backend/iot-client/internal/service/service.go @@ -0,0 +1,218 @@ +// Package service 提供管理系统/App 与 IoT 设备链路之间的稳定上下行 HTTP 边界。 +package service + +import ( + "bytes" + "context" + "crypto/subtle" + "encoding/json" + "fmt" + "io" + "net/http" + "strings" + "sync" + "time" + + "git.apinb.com/heqiapp/platforms/backend/iot-client/internal/config" +) + +type Service struct { + cfg config.Config + client *http.Client + mu sync.RWMutex + byIdentity map[string]Command + byIdempotency map[string]string + http *http.Server +} +type Command struct { + Identity string `json:"identity"` + IdempotencyKey string `json:"idempotency_key"` + DeviceID string `json:"device_id"` + Action string `json:"action"` + ExpiresAt time.Time `json:"expires_at"` + KeyID byte `json:"key_id"` + DeviceKind byte `json:"device_kind"` + DeviceType byte `json:"device_type"` + DeviceModel [3]byte `json:"device_model"` + Controller byte `json:"controller"` + Loop byte `json:"loop"` + Component byte `json:"component"` + Status string `json:"status"` + PacketNumber uint16 `json:"packet_number,omitempty"` + ErrorCode string `json:"error_code,omitempty"` + CreatedAt time.Time `json:"created_at"` + UpdatedAt time.Time `json:"updated_at"` +} +type deviceEnvelope struct { + Type, Topic, DeviceID, ReceivedAt, PayloadHex string + Frame json.RawMessage `json:"frame,omitempty"` +} + +func New(cfg config.Config) *Service { + return &Service{cfg: cfg, client: &http.Client{Timeout: 12 * time.Second}, byIdentity: map[string]Command{}, byIdempotency: map[string]string{}} +} +func (s *Service) Run(ctx context.Context) error { + mux := http.NewServeMux() + mux.HandleFunc("/health", func(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusNoContent) }) + mux.HandleFunc("/v1/device-commands", s.commands) + mux.HandleFunc("/v1/device-commands/", s.commandStatus) + mux.HandleFunc("/internal/v1/device-messages", s.deviceMessage) + s.http = &http.Server{Addr: s.cfg.HTTP.Address, Handler: mux, ReadHeaderTimeout: 5 * time.Second} + go func() { + <-ctx.Done() + shutdown, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + _ = s.http.Shutdown(shutdown) + }() + err := s.http.ListenAndServe() + if err == http.ErrServerClosed { + return nil + } + return err +} +func (s *Service) commands(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + http.Error(w, "method not allowed", 405) + return + } + if !s.authorized(r) { + http.Error(w, "unauthorized", 401) + return + } + var cmd Command + if err := json.NewDecoder(io.LimitReader(r.Body, 64<<10)).Decode(&cmd); err != nil { + http.Error(w, "invalid json", 400) + return + } + if cmd.Identity == "" || cmd.IdempotencyKey == "" || cmd.DeviceID == "" || (cmd.Action != "open_valve" && cmd.Action != "close_valve") { + http.Error(w, "invalid command", 422) + return + } + if cmd.ExpiresAt.IsZero() { + cmd.ExpiresAt = time.Now().Add(30 * time.Second) + } + if time.Now().After(cmd.ExpiresAt) { + http.Error(w, "command expired", 422) + return + } + s.mu.Lock() + if existingID, ok := s.byIdempotency[cmd.IdempotencyKey]; ok { + existing := s.byIdentity[existingID] + s.mu.Unlock() + writeJSON(w, http.StatusOK, existing) + return + } + cmd.Status = "accepted" + cmd.CreatedAt = time.Now().UTC() + cmd.UpdatedAt = cmd.CreatedAt + s.byIdentity[cmd.Identity] = cmd + s.byIdempotency[cmd.IdempotencyKey] = cmd.Identity + s.mu.Unlock() + status, packet, err := s.dispatch(r.Context(), cmd) + s.mu.Lock() + current := s.byIdentity[cmd.Identity] + current.UpdatedAt = time.Now().UTC() + if err != nil { + current.Status = "dispatch_failed" + current.ErrorCode = "IOT_DISPATCH_FAILED" + } else { + current.Status = status + current.PacketNumber = packet + } + s.byIdentity[cmd.Identity] = current + s.mu.Unlock() + if err != nil { + writeJSON(w, http.StatusBadGateway, current) + return + } + writeJSON(w, http.StatusAccepted, current) +} +func (s *Service) commandStatus(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + http.Error(w, "method not allowed", 405) + return + } + if !s.authorized(r) { + http.Error(w, "unauthorized", 401) + return + } + identity := strings.TrimPrefix(r.URL.Path, "/v1/device-commands/") + s.mu.RLock() + cmd, ok := s.byIdentity[identity] + s.mu.RUnlock() + if !ok { + http.Error(w, "not found", 404) + return + } + writeJSON(w, 200, cmd) +} +func (s *Service) deviceMessage(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost || !s.authorized(r) { + http.Error(w, "unauthorized", 401) + return + } + data, err := io.ReadAll(io.LimitReader(r.Body, 1<<20)) + if err != nil { + http.Error(w, "invalid body", 400) + return + } + var envelope deviceEnvelope + if json.Unmarshal(data, &envelope) != nil || envelope.ReceivedAt == "" { + http.Error(w, "invalid envelope", 422) + return + } + request, err := http.NewRequestWithContext(r.Context(), http.MethodPost, s.cfg.Upstream.PlatformAPIURL+"/heqi/internal/v1/iot/device-messages", bytes.NewReader(data)) + if err != nil { + http.Error(w, "upstream request failed", 502) + return + } + request.Header.Set("Content-Type", "application/json") + request.Header.Set("X-Heqi-Iot-Token", s.cfg.Upstream.Token) + response, err := s.client.Do(request) + if err != nil { + http.Error(w, "platform api unavailable", 502) + return + } + defer response.Body.Close() + if response.StatusCode >= 300 { + http.Error(w, "platform api rejected message", 502) + return + } + w.WriteHeader(http.StatusAccepted) +} +func (s *Service) dispatch(ctx context.Context, cmd Command) (string, uint16, error) { + data, _ := json.Marshal(cmd) + request, err := http.NewRequestWithContext(ctx, http.MethodPost, s.cfg.Upstream.IoTServerURL+"/internal/v1/commands", bytes.NewReader(data)) + if err != nil { + return "", 0, err + } + request.Header.Set("Content-Type", "application/json") + request.Header.Set("X-Heqi-Iot-Token", s.cfg.Upstream.Token) + response, err := s.client.Do(request) + if err != nil { + return "", 0, err + } + defer response.Body.Close() + if response.StatusCode != http.StatusAccepted { + body, _ := io.ReadAll(io.LimitReader(response.Body, 4096)) + return "", 0, fmt.Errorf("iot server %d: %s", response.StatusCode, body) + } + var result struct { + Status string `json:"status"` + PacketNumber uint16 `json:"packet_number"` + } + if err = json.NewDecoder(response.Body).Decode(&result); err != nil { + return "", 0, err + } + return result.Status, result.PacketNumber, nil +} +func (s *Service) authorized(r *http.Request) bool { + value := r.Header.Get("X-Heqi-Iot-Token") + expected := s.cfg.HTTP.InternalToken + return len(value) == len(expected) && subtle.ConstantTimeCompare([]byte(value), []byte(expected)) == 1 +} +func writeJSON(w http.ResponseWriter, status int, value any) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(value) +} diff --git a/backend/iot-client/internal/service/service_test.go b/backend/iot-client/internal/service/service_test.go new file mode 100644 index 0000000..4f1de02 --- /dev/null +++ b/backend/iot-client/internal/service/service_test.go @@ -0,0 +1,54 @@ +package service + +import ( + "bytes" + "context" + "net/http" + "net/http/httptest" + "testing" + "time" + + "git.apinb.com/heqiapp/platforms/backend/iot-client/internal/config" +) + +func TestDispatchPreservesIdempotency(t *testing.T) { + calls := 0 + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls++ + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusAccepted) + _, _ = w.Write([]byte(`{"status":"dispatched","packet_number":9}`)) + })) + defer upstream.Close() + srv := New(config.Config{HTTP: config.HTTP{InternalToken: "token"}, Upstream: config.Upstream{IoTServerURL: upstream.URL, Token: "token"}}) + cmd := Command{Identity: "one", IdempotencyKey: "same", DeviceID: "1234567890123456", Action: "close_valve", ExpiresAt: time.Now().Add(time.Minute)} + status, packet, err := srv.dispatch(context.Background(), cmd) + if err != nil || status != "dispatched" || packet != 9 || calls != 1 { + t.Fatalf("status=%s packet=%d calls=%d err=%v", status, packet, calls, err) + } +} + +func TestCommandsDoesNotDispatchDuplicateIdempotencyKey(t *testing.T) { + calls := 0 + upstream := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + calls++ + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusAccepted) + _, _ = w.Write([]byte(`{"status":"dispatched","packet_number":9}`)) + })) + defer upstream.Close() + srv := New(config.Config{HTTP: config.HTTP{InternalToken: "token"}, Upstream: config.Upstream{IoTServerURL: upstream.URL, Token: "token"}}) + body := []byte(`{"identity":"one","idempotency_key":"same","device_id":"1234567890123456","action":"close_valve","expires_at":"2099-01-01T00:00:00Z"}`) + for range 2 { + request := httptest.NewRequest(http.MethodPost, "/v1/device-commands", bytes.NewReader(body)) + request.Header.Set("X-Heqi-Iot-Token", "token") + response := httptest.NewRecorder() + srv.commands(response, request) + if response.Code != http.StatusAccepted && response.Code != http.StatusOK { + t.Fatalf("unexpected status %d: %s", response.Code, response.Body.String()) + } + } + if calls != 1 { + t.Fatalf("duplicate command dispatched %d times", calls) + } +} diff --git a/backend/iot-server/README.md b/backend/iot-server/README.md new file mode 100644 index 0000000..b6cfde0 --- /dev/null +++ b/backend/iot-server/README.md @@ -0,0 +1,5 @@ +# IoT Server + +设备侧 MQTT 适配进程。连接外部 EMQX/Mosquitto Broker,订阅 `devices/+/up` 与 `devices/+/ack`,按《气体探测器通讯协议》V1.8 校验、解密和解析原始二进制帧;下行命令由受保护的内部 HTTP 接口接收后发布到 `devices/{deviceId}/down`。 + +生产环境必须启用 MQTT TLS、每设备凭证并通过环境变量或密钥注入覆盖 YAML 开发占位密钥。 diff --git a/backend/iot-server/cmd/main/main.go b/backend/iot-server/cmd/main/main.go new file mode 100644 index 0000000..e723877 --- /dev/null +++ b/backend/iot-server/cmd/main/main.go @@ -0,0 +1,30 @@ +package main + +import ( + "context" + "flag" + "log" + "os/signal" + "syscall" + + "git.apinb.com/heqiapp/platforms/backend/iot-server/internal/config" + "git.apinb.com/heqiapp/platforms/backend/iot-server/internal/service" +) + +func main() { + path := flag.String("config", "etc/platform_iot_server_dev.yaml", "YAML 配置文件") + flag.Parse() + cfg, err := config.Load(*path) + if err != nil { + log.Fatal(err) + } + srv, err := service.New(cfg) + if err != nil { + log.Fatal(err) + } + ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer cancel() + if err = srv.Run(ctx); err != nil { + log.Fatal(err) + } +} diff --git a/backend/iot-server/etc/platform_iot_server_dev.yaml b/backend/iot-server/etc/platform_iot_server_dev.yaml new file mode 100644 index 0000000..71bfdc8 --- /dev/null +++ b/backend/iot-server/etc/platform_iot_server_dev.yaml @@ -0,0 +1,23 @@ +Service: platform-iot-server +MQTT: + Broker: tcp://127.0.0.1:1883 + ClientID: heqi-iot-server-dev + Username: heqi-dev + Password: change-me + UpTopic: devices/+/up + DownTopic: devices/{deviceId}/down + AckTopic: devices/+/ack + QoS: 1 + TLS: false + CAFile: "" + CertificateFile: "" + PrivateKeyFile: "" +HTTP: + Address: 127.0.0.1:12428 + InternalToken: change-me-iot-internal-token + CallbackURL: http://127.0.0.1:12429/internal/v1/device-messages +Protocol: + # 仅为本地开发占位。生产通过 HEQI_IOT_KEY_1/2/3 或密钥注入覆盖。 + Key1: "00000000000000000000000000000000" + Key2: "11111111111111111111111111111111" + Key3: "22222222222222222222222222222222" diff --git a/backend/iot-server/go.mod b/backend/iot-server/go.mod new file mode 100644 index 0000000..2641db8 --- /dev/null +++ b/backend/iot-server/go.mod @@ -0,0 +1,14 @@ +module git.apinb.com/heqiapp/platforms/backend/iot-server + +go 1.26.1 + +require ( + github.com/eclipse/paho.mqtt.golang v1.5.1 + gopkg.in/yaml.v3 v3.0.1 +) + +require ( + github.com/gorilla/websocket v1.5.3 // indirect + golang.org/x/net v0.52.0 // indirect + golang.org/x/sync v0.20.0 // indirect +) diff --git a/backend/iot-server/go.sum b/backend/iot-server/go.sum new file mode 100644 index 0000000..74aed36 --- /dev/null +++ b/backend/iot-server/go.sum @@ -0,0 +1,12 @@ +github.com/eclipse/paho.mqtt.golang v1.5.1 h1:/VSOv3oDLlpqR2Epjn1Q7b2bSTplJIeV2ISgCl2W7nE= +github.com/eclipse/paho.mqtt.golang v1.5.1/go.mod h1:1/yJCneuyOoCOzKSsOTUc0AJfpsItBGWvYpBLimhArU= +github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= +github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= +golang.org/x/net v0.52.0 h1:He/TN1l0e4mmR3QqHMT2Xab3Aj3L9qjbhRm78/6jrW0= +golang.org/x/net v0.52.0/go.mod h1:R1MAz7uMZxVMualyPXb+VaqGSa3LIaUqk0eEt3w36Sw= +golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= +golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/backend/iot-server/internal/config/config.go b/backend/iot-server/internal/config/config.go new file mode 100644 index 0000000..60aeb9f --- /dev/null +++ b/backend/iot-server/internal/config/config.go @@ -0,0 +1,83 @@ +// Package config 加载 IoT Server 的 MQTT、协议密钥和内部接口配置。 +package config + +import ( + "encoding/hex" + "fmt" + "os" + "strings" + + "gopkg.in/yaml.v3" +) + +type MQTT struct { + Broker string `yaml:"Broker"` + ClientID string `yaml:"ClientID"` + Username string `yaml:"Username"` + Password string `yaml:"Password"` + UpTopic string `yaml:"UpTopic"` + DownTopic string `yaml:"DownTopic"` + AckTopic string `yaml:"AckTopic"` + QoS byte `yaml:"QoS"` + TLS bool `yaml:"TLS"` + CAFile string `yaml:"CAFile"` + CertificateFile string `yaml:"CertificateFile"` + PrivateKeyFile string `yaml:"PrivateKeyFile"` +} +type HTTP struct { + Address string `yaml:"Address"` + InternalToken string `yaml:"InternalToken"` + CallbackURL string `yaml:"CallbackURL"` +} +type Protocol struct { + Key1 string `yaml:"Key1"` + Key2 string `yaml:"Key2"` + Key3 string `yaml:"Key3"` +} +type Config struct { + Service string `yaml:"Service"` + MQTT MQTT `yaml:"MQTT"` + HTTP HTTP `yaml:"HTTP"` + Protocol Protocol `yaml:"Protocol"` +} + +func Load(path string) (Config, error) { + var cfg Config + data, err := os.ReadFile(path) + if err != nil { + return cfg, err + } + if err = yaml.Unmarshal(data, &cfg); err != nil { + return cfg, err + } + override(&cfg.MQTT.Password, "HEQI_IOT_MQTT_PASSWORD") + override(&cfg.HTTP.InternalToken, "HEQI_IOT_INTERNAL_TOKEN") + override(&cfg.Protocol.Key1, "HEQI_IOT_KEY_1") + override(&cfg.Protocol.Key2, "HEQI_IOT_KEY_2") + override(&cfg.Protocol.Key3, "HEQI_IOT_KEY_3") + if cfg.MQTT.Broker == "" || cfg.HTTP.Address == "" || cfg.HTTP.InternalToken == "" { + return cfg, fmt.Errorf("MQTT.Broker、HTTP.Address 和 HTTP.InternalToken 必填") + } + return cfg, nil +} + +func (cfg Config) Keys() (map[byte][]byte, error) { + result := map[byte][]byte{} + for id, value := range map[byte]string{1: cfg.Protocol.Key1, 2: cfg.Protocol.Key2, 3: cfg.Protocol.Key3} { + if strings.TrimSpace(value) == "" { + continue + } + decoded, err := hex.DecodeString(value) + if err != nil || len(decoded) != 16 { + return nil, fmt.Errorf("Protocol.Key%d 必须是 32 位十六进制 AES-128 密钥", id) + } + result[id] = decoded + } + return result, nil +} + +func override(target *string, name string) { + if value := os.Getenv(name); value != "" { + *target = value + } +} diff --git a/backend/iot-server/internal/protocol/frame.go b/backend/iot-server/internal/protocol/frame.go new file mode 100644 index 0000000..43b74d4 --- /dev/null +++ b/backend/iot-server/internal/protocol/frame.go @@ -0,0 +1,173 @@ +// Package protocol 实现《气体探测器通讯协议》V1.8 的二进制帧编解码。 +package protocol + +import ( + "crypto/aes" + "encoding/binary" + "errors" + "fmt" + "time" +) + +const ( + StartByte = byte(0x5E) + EndByte = byte(0x5B) + FixedBodyBytes = 29 // key 到有效数据长度,不含载荷和校验。 +) + +var ( + ErrFrameTooShort = errors.New("数据帧长度不足") + ErrBoundary = errors.New("起始符或结束符无效") + ErrLength = errors.New("帧长度不匹配") + ErrChecksum = errors.New("LRC8 校验失败") + ErrPayloadLength = errors.New("数据包有效长度无效") + ErrEncryptedLength = errors.New("密文长度不是 16 字节的倍数") +) + +// Frame 是设备原始帧的强类型表示;多字节整数均按大端序传输。 +type Frame struct { + KeyID byte + Version byte + Control byte + DeviceKind byte + DeviceType byte + DeviceModel [3]byte + DeviceID [8]byte + PacketNumber uint16 + Sequence byte + Final bool + DeviceTime time.Time + Payload []byte +} + +// Keyring 按协议支持 0 号明文和 1/2/3 号 AES-128 密钥。 +type Keyring map[byte][]byte + +// Encode 构造可直接作为 MQTT payload 发布的厂商二进制帧。 +func Encode(frame Frame, keys Keyring) ([]byte, error) { + payload, err := cryptPayload(frame.KeyID, frame.Payload, keys, false) + if err != nil { + return nil, err + } + body := make([]byte, FixedBodyBytes+len(payload)) + body[0], body[1], body[2] = frame.KeyID, frame.Version, frame.Control + body[3], body[4] = frame.DeviceKind, frame.DeviceType + copy(body[5:8], frame.DeviceModel[:]) + copy(body[8:16], frame.DeviceID[:]) + binary.BigEndian.PutUint16(body[16:18], frame.PacketNumber) + body[18] = frame.Sequence + if frame.Final { + body[19] = 1 + } + encodeBCDTime(body[20:27], frame.DeviceTime) + binary.BigEndian.PutUint16(body[27:29], uint16(len(frame.Payload))) + copy(body[29:], payload) + + frameLength := len(body) + 1 // 加上校验字节,不含起始符、长度字段和结束符。 + if frameLength > 0xffff { + return nil, fmt.Errorf("帧过长: %d", frameLength) + } + result := make([]byte, 0, frameLength+4) + result = append(result, StartByte, byte(frameLength>>8), byte(frameLength)) + result = append(result, body...) + result = append(result, LRC8(result[1:]), EndByte) + return result, nil +} + +// Decode 校验边界、长度、LRC8 和 AES 后返回有效载荷。 +func Decode(raw []byte, keys Keyring) (Frame, error) { + var frame Frame + if len(raw) < FixedBodyBytes+5 { + return frame, ErrFrameTooShort + } + if raw[0] != StartByte || raw[len(raw)-1] != EndByte { + return frame, ErrBoundary + } + declared := int(binary.BigEndian.Uint16(raw[1:3])) + if declared+4 != len(raw) { + return frame, ErrLength + } + if LRC8(raw[1:len(raw)-2]) != raw[len(raw)-2] { + return frame, ErrChecksum + } + body := raw[3 : len(raw)-2] + frame.KeyID, frame.Version, frame.Control = body[0], body[1], body[2] + frame.DeviceKind, frame.DeviceType = body[3], body[4] + copy(frame.DeviceModel[:], body[5:8]) + copy(frame.DeviceID[:], body[8:16]) + frame.PacketNumber = binary.BigEndian.Uint16(body[16:18]) + frame.Sequence, frame.Final = body[18], body[19] == 1 + frame.DeviceTime = decodeBCDTime(body[20:27]) + validLength := int(binary.BigEndian.Uint16(body[27:29])) + plain, err := cryptPayload(frame.KeyID, body[29:], keys, true) + if err != nil { + return frame, err + } + if validLength > len(plain) { + return frame, ErrPayloadLength + } + frame.Payload = append([]byte(nil), plain[:validLength]...) + return frame, nil +} + +// LRC8 返回连续字节和的二进制补码低字节。 +func LRC8(data []byte) byte { + var sum byte + for _, value := range data { + sum += value + } + return ^sum + 1 +} + +func cryptPayload(keyID byte, input []byte, keys Keyring, decrypt bool) ([]byte, error) { + if keyID == 0 { + size := len(input) + if !decrypt && size%aes.BlockSize != 0 { + size += aes.BlockSize - size%aes.BlockSize + } + result := make([]byte, size) + copy(result, input) + return result, nil + } + key := keys[keyID] + if len(key) != aes.BlockSize { + return nil, fmt.Errorf("%d 号 AES 密钥必须为 16 字节", keyID) + } + if decrypt && len(input)%aes.BlockSize != 0 { + return nil, ErrEncryptedLength + } + size := len(input) + if !decrypt && size%aes.BlockSize != 0 { + size += aes.BlockSize - size%aes.BlockSize + } + output := make([]byte, size) + copy(output, input) + block, err := aes.NewCipher(key) + if err != nil { + return nil, err + } + for offset := 0; offset < len(output); offset += aes.BlockSize { + if decrypt { + block.Decrypt(output[offset:offset+aes.BlockSize], input[offset:offset+aes.BlockSize]) + } else { + block.Encrypt(output[offset:offset+aes.BlockSize], output[offset:offset+aes.BlockSize]) + } + } + return output, nil +} + +func encodeBCDTime(dst []byte, value time.Time) { + if value.IsZero() { + value = time.Now() + } + parts := []int{value.Year() / 100, value.Year() % 100, int(value.Month()), value.Day(), value.Hour(), value.Minute(), value.Second()} + for index, part := range parts { + dst[index] = byte((part/10)<<4 | part%10) + } +} + +func decodeBCDTime(src []byte) time.Time { + n := func(value byte) int { return int(value>>4)*10 + int(value&0x0f) } + year := n(src[0])*100 + n(src[1]) + return time.Date(year, time.Month(n(src[2])), n(src[3]), n(src[4]), n(src[5]), n(src[6]), 0, time.Local) +} diff --git a/backend/iot-server/internal/protocol/frame_test.go b/backend/iot-server/internal/protocol/frame_test.go new file mode 100644 index 0000000..99ee485 --- /dev/null +++ b/backend/iot-server/internal/protocol/frame_test.go @@ -0,0 +1,55 @@ +package protocol + +import ( + "bytes" + "testing" + "time" +) + +func TestFrameRoundTripEncrypted(t *testing.T) { + keys := Keyring{2: []byte("0123456789abcdef")} + payload, _ := ValveCommand(1, 2, 3, false) + want := Frame{KeyID: 2, Version: 1, Control: 0x10, DeviceKind: 1, DeviceType: 1, PacketNumber: 7, Sequence: 1, Final: true, DeviceTime: time.Date(2025, 5, 7, 8, 9, 10, 0, time.Local), Payload: payload} + copy(want.DeviceID[:], []byte{0x12, 0x34, 0x56, 0x78, 0x90, 0x12, 0x34, 0x56}) + raw, err := Encode(want, keys) + if err != nil { + t.Fatal(err) + } + got, err := Decode(raw, keys) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(got.Payload, want.Payload) || got.DeviceID != want.DeviceID || got.PacketNumber != want.PacketNumber { + t.Fatalf("round trip mismatch: %#v", got) + } +} + +func TestDecodeRejectsTampering(t *testing.T) { + raw, err := Encode(Frame{Version: 1, Final: true, Payload: []byte{1, 0, 0}}, nil) + if err != nil { + t.Fatal(err) + } + raw[5] ^= 1 + if _, err = Decode(raw, nil); err != ErrChecksum { + t.Fatalf("got %v", err) + } +} + +func TestValveCommand(t *testing.T) { + got, _ := ValveCommand(1, 2, 3, true) + want := []byte{0x40, 1, 0x0A, 1, 2, 3, 0xA0, 0} + if !bytes.Equal(got, want) { + t.Fatalf("got %x want %x", got, want) + } +} + +func TestDecodeRealtimeSensor(t *testing.T) { + payload := []byte{0x50, 1, 0x01, 0, 1, 0, 0, 0, 1, 0x08, 0, 0xE4, 0, 123, 0x0C, 2, 1, 0} + got, err := DecodeRealtime(payload) + if err != nil { + t.Fatal(err) + } + if len(got.Sensors) != 1 || got.Sensors[0].Number != 1 || got.Sensors[0].Value != 123 || got.Sensors[0].Decimal != 2 { + t.Fatalf("unexpected result: %#v", got) + } +} diff --git a/backend/iot-server/internal/protocol/payload.go b/backend/iot-server/internal/protocol/payload.go new file mode 100644 index 0000000..99998c2 --- /dev/null +++ b/backend/iot-server/internal/protocol/payload.go @@ -0,0 +1,155 @@ +package protocol + +import ( + "encoding/binary" + "fmt" +) + +const ( + MainBasicInfo = byte(0x01) + MainRuntime = byte(0x20) + MainQuery = byte(0x30) + MainSetting = byte(0x40) + MainRealtime = byte(0x50) + MainHistorical = byte(0x60) + SubValve = byte(0x0A) +) + +// DataBlock 是主标识下的一个子标识数据块。 +type DataBlock struct { + SubID byte + Data []byte +} + +type RealtimeData struct { + Sensors []SensorReading `json:"sensors,omitempty"` + Events []EventReading `json:"events,omitempty"` + IO []IOReading `json:"io,omitempty"` + Parameters []ParameterReading `json:"parameters,omitempty"` +} +type SensorReading struct { + Number uint32 `json:"number"` + SensorType byte `json:"sensor_type"` + Object uint16 `json:"object"` + Value int16 `json:"value"` + Unit byte `json:"unit"` + Decimal byte `json:"decimal"` + Status byte `json:"status"` +} +type EventReading struct { + Type uint16 `json:"type"` + Number uint32 `json:"number"` + Value int16 `json:"value"` +} +type IOReading struct { + Type uint16 `json:"type"` + Number uint32 `json:"number"` + Status byte `json:"status"` +} +type ParameterReading struct { + Code byte `json:"code"` + Value int16 `json:"value"` +} + +// EncodePayload 按“主标识、块数量、子标识、定长数据”编码。 +// 由于厂商协议没有携带块长度,本函数用于已知命令;上行解析由业务标识专用解析器完成。 +func EncodePayload(mainID byte, blocks ...DataBlock) ([]byte, error) { + if len(blocks) > 255 { + return nil, fmt.Errorf("数据块数量超过 255") + } + result := []byte{mainID, byte(len(blocks))} + for _, block := range blocks { + result = append(result, block.SubID) + result = append(result, block.Data...) + } + result = append(result, 0x00) + return result, nil +} + +// ValveCommand 编码 0x40/0x0A 电磁阀控制;0x01 关闭,0xA0 开启。 +func ValveCommand(controller, loop, component byte, open bool) ([]byte, error) { + action := byte(0x01) + if open { + action = 0xA0 + } + return EncodePayload(MainSetting, DataBlock{SubID: SubValve, Data: []byte{controller, loop, component, action}}) +} + +// DecodeRealtime 解析协议 0x50 的传感器、事件、IO 和 AI 阀参数定长数据块。 +func DecodeRealtime(payload []byte) (RealtimeData, error) { + var result RealtimeData + if len(payload) < 3 || payload[0] != MainRealtime { + return result, fmt.Errorf("不是实时数据包") + } + offset := 2 + for blockIndex := 0; blockIndex < int(payload[1]); blockIndex++ { + if offset >= len(payload) { + return result, fmt.Errorf("实时数据块被截断") + } + subID := payload[offset] + offset++ + switch subID { + case 0x01: + count, next, err := readCount(payload, offset, 12) + if err != nil { + return result, err + } + offset = next + for range count { + item := payload[offset : offset+12] + result.Sensors = append(result.Sensors, SensorReading{Number: binary.BigEndian.Uint32(item[0:4]), SensorType: item[4], Object: binary.BigEndian.Uint16(item[5:7]), Value: int16(binary.BigEndian.Uint16(item[7:9])), Unit: item[9], Decimal: item[10], Status: item[11]}) + offset += 12 + } + case 0x02: + count, next, err := readCount(payload, offset, 8) + if err != nil { + return result, err + } + offset = next + for range count { + item := payload[offset : offset+8] + result.Events = append(result.Events, EventReading{Type: binary.BigEndian.Uint16(item[0:2]), Number: binary.BigEndian.Uint32(item[2:6]), Value: int16(binary.BigEndian.Uint16(item[6:8]))}) + offset += 8 + } + case 0x03: + count, next, err := readCount(payload, offset, 7) + if err != nil { + return result, err + } + offset = next + for range count { + item := payload[offset : offset+7] + result.IO = append(result.IO, IOReading{Type: binary.BigEndian.Uint16(item[0:2]), Number: binary.BigEndian.Uint32(item[2:6]), Status: item[6]}) + offset += 7 + } + case 0x04: + if offset >= len(payload) { + return result, fmt.Errorf("参数块被截断") + } + count := int(payload[offset]) + offset++ + if offset+count*3 > len(payload) { + return result, fmt.Errorf("参数块长度无效") + } + for range count { + result.Parameters = append(result.Parameters, ParameterReading{Code: payload[offset], Value: int16(binary.BigEndian.Uint16(payload[offset+1 : offset+3]))}) + offset += 3 + } + default: + return result, fmt.Errorf("未知实时数据子标识 0x%02X", subID) + } + } + return result, nil +} + +func readCount(payload []byte, offset, recordSize int) (int, int, error) { + if offset+2 > len(payload) { + return 0, offset, fmt.Errorf("数据块数量被截断") + } + count := int(binary.BigEndian.Uint16(payload[offset : offset+2])) + offset += 2 + if offset+count*recordSize > len(payload) { + return 0, offset, fmt.Errorf("数据块记录长度无效") + } + return count, offset, nil +} diff --git a/backend/iot-server/internal/service/service.go b/backend/iot-server/internal/service/service.go new file mode 100644 index 0000000..eec365c --- /dev/null +++ b/backend/iot-server/internal/service/service.go @@ -0,0 +1,236 @@ +// Package service 连接外部 MQTT Broker,并在内部 HTTP 边界接收待下发命令。 +package service + +import ( + "bytes" + "context" + "crypto/subtle" + "crypto/tls" + "crypto/x509" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "net/http" + "os" + "strings" + "sync/atomic" + "time" + + "git.apinb.com/heqiapp/platforms/backend/iot-server/internal/config" + "git.apinb.com/heqiapp/platforms/backend/iot-server/internal/protocol" + mqtt "github.com/eclipse/paho.mqtt.golang" +) + +type Service struct { + cfg config.Config + keys protocol.Keyring + mqtt mqtt.Client + packet atomic.Uint32 + http *http.Server +} +type Command struct { + Identity string `json:"identity"` + IdempotencyKey string `json:"idempotency_key"` + DeviceID string `json:"device_id"` + Action string `json:"action"` + ExpiresAt time.Time `json:"expires_at"` + KeyID byte `json:"key_id"` + DeviceKind byte `json:"device_kind"` + DeviceType byte `json:"device_type"` + DeviceModel [3]byte `json:"device_model"` + Controller byte `json:"controller"` + Loop byte `json:"loop"` + Component byte `json:"component"` +} +type Envelope struct { + Type, Topic, DeviceID, ReceivedAt string + PayloadHex string + Frame *DecodedFrame `json:"frame,omitempty"` +} +type DecodedFrame struct { + KeyID, Version, Control, MainID byte + PacketNumber uint16 + Sequence byte + Final bool + DeviceTime string + PayloadHex string + Realtime *protocol.RealtimeData `json:"realtime,omitempty"` +} + +func New(cfg config.Config) (*Service, error) { + keys, err := cfg.Keys() + if err != nil { + return nil, err + } + options := mqtt.NewClientOptions().AddBroker(cfg.MQTT.Broker).SetClientID(cfg.MQTT.ClientID).SetUsername(cfg.MQTT.Username).SetPassword(cfg.MQTT.Password).SetAutoReconnect(true).SetConnectRetry(true) + if cfg.MQTT.TLS { + tlsConfig, tlsErr := makeTLSConfig(cfg) + if tlsErr != nil { + return nil, tlsErr + } + options.SetTLSConfig(tlsConfig) + } + client := mqtt.NewClient(options) + return &Service{cfg: cfg, keys: keys, mqtt: client}, nil +} + +func makeTLSConfig(cfg config.Config) (*tls.Config, error) { + roots, err := x509.SystemCertPool() + if err != nil { + roots = x509.NewCertPool() + } + if cfg.MQTT.CAFile != "" { + data, readErr := os.ReadFile(cfg.MQTT.CAFile) + if readErr != nil { + return nil, readErr + } + if !roots.AppendCertsFromPEM(data) { + return nil, fmt.Errorf("MQTT CA 证书无效") + } + } + result := &tls.Config{MinVersion: tls.VersionTLS12, RootCAs: roots} + if cfg.MQTT.CertificateFile != "" || cfg.MQTT.PrivateKeyFile != "" { + certificate, loadErr := tls.LoadX509KeyPair(cfg.MQTT.CertificateFile, cfg.MQTT.PrivateKeyFile) + if loadErr != nil { + return nil, loadErr + } + result.Certificates = []tls.Certificate{certificate} + } + return result, nil +} + +func (s *Service) Run(ctx context.Context) error { + if token := s.mqtt.Connect(); !token.WaitTimeout(15 * time.Second) { + return fmt.Errorf("MQTT 连接超时") + } else if token.Error() != nil { + return token.Error() + } + for _, topic := range []string{s.cfg.MQTT.UpTopic, s.cfg.MQTT.AckTopic} { + if token := s.mqtt.Subscribe(topic, s.cfg.MQTT.QoS, s.onMessage); token.Wait() && token.Error() != nil { + return token.Error() + } + } + mux := http.NewServeMux() + mux.HandleFunc("/health", s.health) + mux.HandleFunc("/internal/v1/commands", s.command) + s.http = &http.Server{Addr: s.cfg.HTTP.Address, Handler: mux, ReadHeaderTimeout: 5 * time.Second} + go func() { + <-ctx.Done() + shutdown, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + _ = s.http.Shutdown(shutdown) + s.mqtt.Disconnect(250) + }() + err := s.http.ListenAndServe() + if err == http.ErrServerClosed { + return nil + } + return err +} + +func (s *Service) health(w http.ResponseWriter, _ *http.Request) { + if !s.mqtt.IsConnectionOpen() { + http.Error(w, "mqtt disconnected", http.StatusServiceUnavailable) + return + } + w.WriteHeader(http.StatusNoContent) +} + +func (s *Service) command(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost { + http.Error(w, "method not allowed", http.StatusMethodNotAllowed) + return + } + if !secureEqual(r.Header.Get("X-Heqi-Iot-Token"), s.cfg.HTTP.InternalToken) { + http.Error(w, "unauthorized", http.StatusUnauthorized) + return + } + var cmd Command + if err := json.NewDecoder(io.LimitReader(r.Body, 64<<10)).Decode(&cmd); err != nil { + http.Error(w, "invalid json", http.StatusBadRequest) + return + } + if cmd.Identity == "" || cmd.IdempotencyKey == "" || cmd.DeviceID == "" || time.Now().After(cmd.ExpiresAt) { + http.Error(w, "invalid or expired command", http.StatusUnprocessableEntity) + return + } + deviceID, err := decodeDeviceID(cmd.DeviceID) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + var payload []byte + switch cmd.Action { + case "open_valve": + payload, err = protocol.ValveCommand(cmd.Controller, cmd.Loop, cmd.Component, true) + case "close_valve": + payload, err = protocol.ValveCommand(cmd.Controller, cmd.Loop, cmd.Component, false) + default: + err = fmt.Errorf("unsupported action") + } + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + n := s.packet.Add(1) + frame := protocol.Frame{KeyID: cmd.KeyID, Version: 1, Control: 0x10, DeviceKind: cmd.DeviceKind, DeviceType: cmd.DeviceType, DeviceModel: cmd.DeviceModel, DeviceID: deviceID, PacketNumber: uint16(n%65535 + 1), Sequence: 1, Final: true, DeviceTime: time.Now(), Payload: payload} + raw, err := protocol.Encode(frame, s.keys) + if err != nil { + http.Error(w, err.Error(), http.StatusUnprocessableEntity) + return + } + topic := strings.ReplaceAll(s.cfg.MQTT.DownTopic, "{deviceId}", cmd.DeviceID) + token := s.mqtt.Publish(topic, s.cfg.MQTT.QoS, false, raw) + if !token.WaitTimeout(10*time.Second) || token.Error() != nil { + http.Error(w, "mqtt publish failed", http.StatusBadGateway) + return + } + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusAccepted) + _ = json.NewEncoder(w).Encode(map[string]any{"command_identity": cmd.Identity, "packet_number": frame.PacketNumber, "status": "dispatched"}) +} + +func (s *Service) onMessage(_ mqtt.Client, message mqtt.Message) { + raw := append([]byte(nil), message.Payload()...) + envelope := Envelope{Type: "device_message", Topic: message.Topic(), ReceivedAt: time.Now().UTC().Format(time.RFC3339Nano), PayloadHex: hex.EncodeToString(raw)} + if frame, err := protocol.Decode(raw, s.keys); err == nil { + deviceID := hex.EncodeToString(frame.DeviceID[:]) + envelope.DeviceID = deviceID + main := byte(0) + if len(frame.Payload) > 0 { + main = frame.Payload[0] + } + decoded := &DecodedFrame{KeyID: frame.KeyID, Version: frame.Version, Control: frame.Control, MainID: main, PacketNumber: frame.PacketNumber, Sequence: frame.Sequence, Final: frame.Final, DeviceTime: frame.DeviceTime.Format(time.RFC3339), PayloadHex: hex.EncodeToString(frame.Payload)} + if main == protocol.MainRealtime { + if realtime, parseErr := protocol.DecodeRealtime(frame.Payload); parseErr == nil { + decoded.Realtime = &realtime + } + } + envelope.Frame = decoded + } + data, _ := json.Marshal(envelope) + request, err := http.NewRequest(http.MethodPost, s.cfg.HTTP.CallbackURL, bytes.NewReader(data)) + if err != nil { + return + } + request.Header.Set("Content-Type", "application/json") + request.Header.Set("X-Heqi-Iot-Token", s.cfg.HTTP.InternalToken) + response, err := http.DefaultClient.Do(request) + if err == nil { + _ = response.Body.Close() + } +} + +func decodeDeviceID(value string) ([8]byte, error) { + var result [8]byte + decoded, err := hex.DecodeString(value) + if err != nil || len(decoded) != 8 { + return result, fmt.Errorf("device_id 必须是 16 位 BCD/十六进制字符串") + } + copy(result[:], decoded) + return result, nil +} +func secureEqual(left, right string) bool { + return len(left) == len(right) && subtle.ConstantTimeCompare([]byte(left), []byte(right)) == 1 +} diff --git a/backend/iot/README.md b/backend/iot/README.md deleted file mode 100644 index aac1e03..0000000 --- a/backend/iot/README.md +++ /dev/null @@ -1,3 +0,0 @@ -# Platform IoT - -独立 IoT 进程预留 MQTT TLS 会话、协议适配、遥测校验与命令回执边界。首期使用 Mock 适配,不连接真实设备或 MQTT Broker。 diff --git a/backend/iot/cmd/main/main.go b/backend/iot/cmd/main/main.go deleted file mode 100644 index b3e9df4..0000000 --- a/backend/iot/cmd/main/main.go +++ /dev/null @@ -1,21 +0,0 @@ -// IoT 进程入口;真实 MQTT TLS 适配将在设备协议和证书契约确认后接入。 -package main - -import ( - "os" - "os/signal" - "syscall" - - "git.apinb.com/bsm-sdk/core/printer" - "git.apinb.com/heqiapp/platforms/backend/iot/internal/config" - "git.apinb.com/heqiapp/platforms/backend/iot/internal/impl" -) - -func main() { - config.New("PlatformIot") - impl.NewImpl() - printer.Info("[BSM - PlatformIot] Mock IoT adapter started; MQTT broker is not connected") - quit := make(chan os.Signal, 1) - signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) - <-quit -} diff --git a/backend/iot/etc/platform_iot_dev.yaml b/backend/iot/etc/platform_iot_dev.yaml deleted file mode 100644 index fc586f7..0000000 --- a/backend/iot/etc/platform_iot_dev.yaml +++ /dev/null @@ -1,9 +0,0 @@ -Service: platform-iot -Port: 12428 -Databases: - Driver: postgres - Source: - - host=127.0.0.1 user=postgres password=change-me dbname=agent_dev port=5432 sslmode=disable TimeZone=Asia/Shanghai -Cache: redis://default:change-me@127.0.0.1:6379/0 -OnMicroService: false -SecretKey: change-me-to-a-random-string diff --git a/backend/iot/go.mod b/backend/iot/go.mod deleted file mode 100644 index 3987a48..0000000 --- a/backend/iot/go.mod +++ /dev/null @@ -1,48 +0,0 @@ -module git.apinb.com/heqiapp/platforms/backend/iot - -go 1.26.1 - -require ( - git.apinb.com/bsm-sdk/core v0.2.0 - github.com/patrickmn/go-cache v2.1.0+incompatible - gorm.io/gorm v1.31.1 -) - -require ( - filippo.io/edwards25519 v1.1.0 // indirect - github.com/cespare/xxhash/v2 v2.3.0 // indirect - github.com/coreos/go-semver v0.3.1 // indirect - github.com/coreos/go-systemd/v22 v22.5.0 // indirect - github.com/go-sql-driver/mysql v1.8.1 // indirect - github.com/gogo/protobuf v1.3.2 // indirect - github.com/golang/protobuf v1.5.4 // indirect - github.com/google/uuid v1.6.0 // indirect - github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 // indirect - github.com/jackc/pgpassfile v1.0.0 // indirect - github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect - github.com/jackc/pgx/v5 v5.6.0 // indirect - github.com/jackc/puddle/v2 v2.2.2 // indirect - github.com/jinzhu/inflection v1.0.0 // indirect - github.com/jinzhu/now v1.1.5 // indirect - github.com/oklog/ulid/v2 v2.1.1 // indirect - github.com/redis/go-redis/v9 v9.19.0 // indirect - github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e // indirect - go.etcd.io/etcd/api/v3 v3.6.11 // indirect - go.etcd.io/etcd/client/pkg/v3 v3.6.11 // indirect - go.etcd.io/etcd/client/v3 v3.6.11 // indirect - go.uber.org/atomic v1.11.0 // indirect - go.uber.org/multierr v1.11.0 // indirect - go.uber.org/zap v1.27.0 // indirect - golang.org/x/crypto v0.49.0 // indirect - golang.org/x/net v0.52.0 // indirect - golang.org/x/sync v0.20.0 // indirect - golang.org/x/sys v0.44.0 // indirect - golang.org/x/text v0.36.0 // indirect - google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478 // indirect - google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478 // indirect - google.golang.org/grpc v1.81.1 // indirect - google.golang.org/protobuf v1.36.11 // indirect - gopkg.in/yaml.v3 v3.0.1 // indirect - gorm.io/driver/mysql v1.6.0 // indirect - gorm.io/driver/postgres v1.6.0 // indirect -) diff --git a/backend/iot/go.sum b/backend/iot/go.sum deleted file mode 100644 index 3417f6d..0000000 --- a/backend/iot/go.sum +++ /dev/null @@ -1,159 +0,0 @@ -filippo.io/edwards25519 v1.1.0 h1:FNf4tywRC1HmFuKW5xopWpigGjJKiJSV0Cqo0cJWDaA= -filippo.io/edwards25519 v1.1.0/go.mod h1:BxyFTGdWcka3PhytdK4V28tE5sGfRvvvRV7EaN4VDT4= -git.apinb.com/bsm-sdk/core v0.2.0 h1:/e9yqpsbKBrRgMiGpS3KX2O4qDLxo5V5GpPSLNGxEKw= -git.apinb.com/bsm-sdk/core v0.2.0/go.mod h1:E9T6Eboo/0Zb36BjkKbIgvFzq4fQ2Q8P/7y5zmYTI6Y= -github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= -github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= -github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= -github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= -github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= -github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/coreos/go-semver v0.3.1 h1:yi21YpKnrx1gt5R+la8n5WgS0kCrsPp33dmEyHReZr4= -github.com/coreos/go-semver v0.3.1/go.mod h1:irMmmIw/7yzSRPWryHsK7EYSg09caPQL03VsM8rvUec= -github.com/coreos/go-systemd/v22 v22.5.0 h1:RrqgGjYQKalulkV8NGVIfkXQf6YYmOyiJKk8iXXhfZs= -github.com/coreos/go-systemd/v22 v22.5.0/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSVTIJ3seZv2GcEnc= -github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= -github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= -github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= -github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= -github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= -github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= -github.com/go-sql-driver/mysql v1.8.1 h1:LedoTUt/eveggdHS9qUFC1EFSa8bU2+1pZjSRpvNJ1Y= -github.com/go-sql-driver/mysql v1.8.1/go.mod h1:wEBSXgmK//2ZFJyE+qWnIsVGmvmEKlqwuVSjsCm7DZg= -github.com/godbus/dbus/v5 v5.0.4/go.mod h1:xhWf0FNVPg57R7Z0UbKHbJfkEywrmjJnf7w5xrFpKfA= -github.com/gogo/protobuf v1.3.2 h1:Ov1cvc58UF3b5XjBnZv7+opcTcQFZebYjWzi34vdm4Q= -github.com/gogo/protobuf v1.3.2/go.mod h1:P1XiOD3dCwIKUDQYPy72D8LYyHL2YPYrpS2s69NZV8Q= -github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= -github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= -github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= -github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= -github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= -github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= -github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0 h1:5VipnvEpbqr2gA2VbM+nYVbkIF28c5ZQfqCBQ5g2xfk= -github.com/grpc-ecosystem/grpc-gateway/v2 v2.29.0/go.mod h1:Hyl3n6Twe1hvtd9XUXDec4pTvgMSEixRuQKPTMH2bNs= -github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= -github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= -github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= -github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= -github.com/jackc/pgx/v5 v5.6.0 h1:SWJzexBzPL5jb0GEsrPMLIsi/3jOo7RHlzTjcAeDrPY= -github.com/jackc/pgx/v5 v5.6.0/go.mod h1:DNZ/vlrUnhWCoFGxHAG8U2ljioxukquj7utPDgtQdTw= -github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= -github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= -github.com/jinzhu/inflection v1.0.0 h1:K317FqzuhWc8YvSVlFMCCUb36O/S9MCKRDI7QkRKD/E= -github.com/jinzhu/inflection v1.0.0/go.mod h1:h+uFLlag+Qp1Va5pdKtLDYj+kHp5pxUVkryuEj+Srlc= -github.com/jinzhu/now v1.1.5 h1:/o9tlHleP7gOFmsnYNz3RGnqzefHA47wQpKrrdTIwXQ= -github.com/jinzhu/now v1.1.5/go.mod h1:d3SSVoowX0Lcu0IBviAWJpolVfI5UJVZZ7cO71lE/z8= -github.com/kisielk/errcheck v1.5.0/go.mod h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8= -github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck= -github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y= -github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= -github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= -github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= -github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= -github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= -github.com/oklog/ulid/v2 v2.1.1 h1:suPZ4ARWLOJLegGFiZZ1dFAkqzhMjL3J1TzI+5wHz8s= -github.com/oklog/ulid/v2 v2.1.1/go.mod h1:rcEKHmBBKfef9DhnvX7y1HZBYxjXb0cP5ExxNsTT1QQ= -github.com/patrickmn/go-cache v2.1.0+incompatible h1:HRMgzkcYKYpi3C8ajMPV8OFXaaRUnok+kx1WdO15EQc= -github.com/patrickmn/go-cache v2.1.0+incompatible/go.mod h1:3Qf8kWWT7OJRJbdiICTKqZju1ZixQ/KpMGzzAfe6+WQ= -github.com/pborman/getopt v0.0.0-20170112200414-7148bc3a4c30/go.mod h1:85jBQOZwpVEaDAr341tbn15RS4fCAsIst0qp7i8ex1o= -github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= -github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/redis/go-redis/v9 v9.19.0 h1:XPVaaPSnG6RhYf7p+rmSa9zZfeVAnWsH5h3lxthOm/k= -github.com/redis/go-redis/v9 v9.19.0/go.mod h1:v/M13XI1PVCDcm01VtPFOADfZtHf8YW3baQf57KlIkA= -github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= -github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= -github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e h1:MRM5ITcdelLK2j1vwZ3Je0FKVCfqOLp5zO6trqMLYs0= -github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e/go.mod h1:XV66xRDqSt+GTGFMVlhk3ULuV0y9ZmzeVGR4mloJI3M= -github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= -github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= -github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= -github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= -github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= -github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= -github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= -github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= -github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= -go.etcd.io/etcd/api/v3 v3.6.11 h1:XFGTgrJ8nak3kB4NgMG8t7NT+lEeuuvKQAqUHKVgkWQ= -go.etcd.io/etcd/api/v3 v3.6.11/go.mod h1:HYfTh0jyh+uFgp6gMbxJteIDYY97yMuYz85Rnw6Gy9o= -go.etcd.io/etcd/client/pkg/v3 v3.6.11 h1:e41mp315Yn3QMGPmEzCyLsMINgJXTY/dX8kM++1csxU= -go.etcd.io/etcd/client/pkg/v3 v3.6.11/go.mod h1:DysuMe/inqRyC/1tjRR6hReH/VV9Lufs27YKSKBWWJg= -go.etcd.io/etcd/client/v3 v3.6.11 h1:LAByD96VmmeuairkvdAcE0RZnrmGz/q3ceeWePo9bwc= -go.etcd.io/etcd/client/v3 v3.6.11/go.mod h1:vOTDMCo+fGPEClJqcFEFSqZ+8e7WKV7AyqJjX//HR2w= -go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= -go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= -go.opentelemetry.io/otel v1.43.0 h1:mYIM03dnh5zfN7HautFE4ieIig9amkNANT+xcVxAj9I= -go.opentelemetry.io/otel v1.43.0/go.mod h1:JuG+u74mvjvcm8vj8pI5XiHy1zDeoCS2LB1spIq7Ay0= -go.opentelemetry.io/otel/metric v1.43.0 h1:d7638QeInOnuwOONPp4JAOGfbCEpYb+K6DVWvdxGzgM= -go.opentelemetry.io/otel/metric v1.43.0/go.mod h1:RDnPtIxvqlgO8GRW18W6Z/4P462ldprJtfxHxyKd2PY= -go.opentelemetry.io/otel/sdk v1.43.0 h1:pi5mE86i5rTeLXqoF/hhiBtUNcrAGHLKQdhg4h4V9Dg= -go.opentelemetry.io/otel/sdk v1.43.0/go.mod h1:P+IkVU3iWukmiit/Yf9AWvpyRDlUeBaRg6Y+C58QHzg= -go.opentelemetry.io/otel/sdk/metric v1.43.0 h1:S88dyqXjJkuBNLeMcVPRFXpRw2fuwdvfCGLEo89fDkw= -go.opentelemetry.io/otel/sdk/metric v1.43.0/go.mod h1:C/RJtwSEJ5hzTiUz5pXF1kILHStzb9zFlIEe85bhj6A= -go.opentelemetry.io/otel/trace v1.43.0 h1:BkNrHpup+4k4w+ZZ86CZoHHEkohws8AY+WTX09nk+3A= -go.opentelemetry.io/otel/trace v1.43.0/go.mod h1:/QJhyVBUUswCphDVxq+8mld+AvhXZLhe+8WVFxiFff0= -go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= -go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= -go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= -go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= -go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= -go.uber.org/multierr v1.11.0/go.mod h1:20+QtiLqy0Nd6FdQB9TLXag12DsQkrbs3htMFfDN80Y= -go.uber.org/zap v1.27.0 h1:aJMhYGrd5QSmlpLMr2MftRKl7t8J8PTZPA732ud/XR8= -go.uber.org/zap v1.27.0/go.mod h1:GB2qFLM7cTU87MWRP2mPIjqfIDnGu+VIO4V/SdhGo2E= -golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= -golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI= -golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= -golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= -golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA= -golang.org/x/mod v0.2.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= -golang.org/x/mod v0.3.0/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA= -golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= -golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= -golang.org/x/net v0.0.0-20200226121028-0de0cce0169b/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= -golang.org/x/net v0.0.0-20201021035429-f5854403a974/go.mod h1:sp8m0HH+o8qH0wwXwYZr8TS3Oi6o0r6Gce1SSxlDquU= -golang.org/x/net v0.52.0 h1:He/TN1l0e4mmR3QqHMT2Xab3Aj3L9qjbhRm78/6jrW0= -golang.org/x/net v0.52.0/go.mod h1:R1MAz7uMZxVMualyPXb+VaqGSa3LIaUqk0eEt3w36Sw= -golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20190911185100-cd5d95a43a6e/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.0.0-20201020160332-67f06af15bc9/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= -golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= -golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= -golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20200930185726-fdedc70b468f/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.44.0 h1:ildZl3J4uzeKP07r2F++Op7E9B29JRUy+a27EibtBTQ= -golang.org/x/sys v0.44.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= -golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= -golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= -golang.org/x/text v0.36.0 h1:JfKh3XmcRPqZPKevfXVpI1wXPTqbkE5f7JA92a55Yxg= -golang.org/x/text v0.36.0/go.mod h1:NIdBknypM8iqVmPiuco0Dh6P5Jcdk8lJL0CUebqK164= -golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= -golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= -golang.org/x/tools v0.0.0-20200619180055-7c47624df98f/go.mod h1:EkVYQZoAsY45+roYkvgYkIh4xh/qjgUK9TdY2XT94GE= -golang.org/x/tools v0.0.0-20210106214847-113979e3529a/go.mod h1:emZCQorbCU4vsT4fOWvOPXz4eW1wZW4PmDk9uLelYpA= -golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= -gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= -gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E= -google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478 h1:yQugLulqltosq0B/f8l4w9VryjV+N/5gcW0jQ3N8Qec= -google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478/go.mod h1:C6ADNqOxbgdUUeRTU+LCHDPB9ttAMCTff6auwCVa4uc= -google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478 h1:RmoJA1ujG+/lRGNfUnOMfhCy5EipVMyvUE+KNbPbTlw= -google.golang.org/genproto/googleapis/rpc v0.0.0-20260414002931-afd174a4e478/go.mod h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8= -google.golang.org/grpc v1.81.1 h1:VnnIIZ88UzOOKLukQi+ImGz8O1Wdp8nAGGnvOfEIWQQ= -google.golang.org/grpc v1.81.1/go.mod h1:xGH9GfzOyMTGIOXBJmXt+BX/V0kcdQbdcuwQ/zNw42I= -google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE= -google.golang.org/protobuf v1.36.11/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= -gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= -gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= -gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= -gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= -gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= -gorm.io/driver/mysql v1.6.0 h1:eNbLmNTpPpTOVZi8MMxCi2aaIm0ZpInbORNXDwyLGvg= -gorm.io/driver/mysql v1.6.0/go.mod h1:D/oCC2GWK3M/dqoLxnOlaNKmXz8WNTfcS9y5ovaSqKo= -gorm.io/driver/postgres v1.6.0 h1:2dxzU8xJ+ivvqTRph34QX+WrRaJlmfyPqXmoGVjMBa4= -gorm.io/driver/postgres v1.6.0/go.mod h1:vUw0mrGgrTK+uPHEhAdV4sfFELrByKVGnaVRkXDhtWo= -gorm.io/gorm v1.31.1 h1:7CA8FTFz/gRfgqgpeKIBcervUn3xSyPUmr6B2WXJ7kg= -gorm.io/gorm v1.31.1/go.mod h1:XyQVbO2k6YkOis7C2437jSit3SsDK72s7n7rsSHd+Gs= diff --git a/backend/iot/internal/config/config.go b/backend/iot/internal/config/config.go deleted file mode 100644 index 5a49183..0000000 --- a/backend/iot/internal/config/config.go +++ /dev/null @@ -1,29 +0,0 @@ -// Package config 使用仓库统一的 BSM 运行配置与地址校验。 -package config - -import ( - "net" - - "git.apinb.com/bsm-sdk/core/conf" -) - -// Spec 是 IoT 进程运行配置。 -var Spec SrvConfig - -// SrvConfig 保持与 API 进程相同的 BSM 配置结构。 -type SrvConfig struct { - conf.Base `yaml:",inline"` - Databases *conf.DBConf `yaml:"Databases"` - Rpc map[string]conf.RpcConf `yaml:"Rpc"` - Apm *conf.ApmConf `yaml:"APM"` -} - -// New 初始化 IoT 进程配置。 -func New(srvKey string) { - conf.New(srvKey, &Spec) - Spec.Port = conf.CheckPort(Spec.Port) - Spec.BindIP = conf.CheckIP(Spec.BindIP) - Spec.Addr = net.JoinHostPort(Spec.BindIP, Spec.Port) - conf.NotNil(Spec.Service, Spec.Cache) - conf.PrintInfo(Spec.Addr) -} diff --git a/backend/iot/internal/impl/new.go b/backend/iot/internal/impl/new.go deleted file mode 100644 index cbc6957..0000000 --- a/backend/iot/internal/impl/new.go +++ /dev/null @@ -1,25 +0,0 @@ -// Package impl 按仓库统一约定创建 IoT 进程所需的基础设施实例。 -package impl - -import ( - "git.apinb.com/bsm-sdk/core/cache/redis" - "git.apinb.com/bsm-sdk/core/logger" - "git.apinb.com/bsm-sdk/core/with" - "git.apinb.com/heqiapp/platforms/backend/iot/internal/config" - "github.com/patrickmn/go-cache" - "gorm.io/gorm" -) - -var ( - RedisService *redis.RedisClient - DBService *gorm.DB - MemoryService *cache.Cache -) - -// NewImpl 初始化 MQTT Mock 进程共用的数据库、缓存、Redis 与日志。 -func NewImpl() { - MemoryService = with.Memory(nil) - RedisService = with.RedisCache(config.Spec.Cache) - DBService = with.Databases(config.Spec.Databases, nil) - logger.New(nil) -} diff --git a/backend/worker/README.md b/backend/worker/README.md index 52e10da..fb5809e 100644 --- a/backend/worker/README.md +++ b/backend/worker/README.md @@ -1,3 +1,3 @@ -# Platform Worker +# Worker -独立 Worker 进程沿用仓库统一的 BSM 配置与基础设施创建方式。首期只提供 Mock 事件循环,真实 Redis Streams 消费将在 Outbox、死信、幂等消费与积压监控契约确定后接入。 +独立 Worker 负责支付超时任务,并通过 API 的锁定认领接口投递 IoT Outbox。数据库事实仍由 API 维护,Worker 不复制同步领域事务;后续可在保持 Outbox 状态契约的前提下用 Redis Streams 唤醒替代短轮询。 diff --git a/backend/worker/cmd/main/main.go b/backend/worker/cmd/main/main.go index 1e792f2..819bd3a 100644 --- a/backend/worker/cmd/main/main.go +++ b/backend/worker/cmd/main/main.go @@ -1,18 +1,18 @@ -// Worker 进程入口;保留独立扩缩和 Redis Streams 消费边界。 package main import ( "bytes" "context" + "encoding/json" + "git.apinb.com/bsm-sdk/core/printer" + "git.apinb.com/heqiapp/platforms/backend/worker/internal/config" + "git.apinb.com/heqiapp/platforms/backend/worker/internal/impl" + "io" "net/http" "os" "os/signal" "syscall" "time" - - "git.apinb.com/bsm-sdk/core/printer" - "git.apinb.com/heqiapp/platforms/backend/worker/internal/config" - "git.apinb.com/heqiapp/platforms/backend/worker/internal/impl" ) func main() { @@ -21,18 +21,115 @@ func main() { ctx, cancel := context.WithCancel(context.Background()) defer cancel() go closeExpiredPayments(ctx) - printer.Info("[BSM - PlatformWorker] payment timeout scheduler started") + go dispatchIoTOutbox(ctx) + printer.Info("[BSM - PlatformWorker] payment scheduler and IoT Outbox dispatcher started") quit := make(chan os.Signal, 1) signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM) <-quit cancel() } - func closeExpiredPayments(ctx context.Context) { - ticker := time.NewTicker(time.Duration(config.Spec.PaymentAPI.IntervalSeconds) * time.Second); defer ticker.Stop() - for { select { case <-ctx.Done(): return; case <-ticker.C: - request, err := http.NewRequestWithContext(ctx, http.MethodPost, config.Spec.PaymentAPI.BaseURL+"/heqi/internal/v1/payment/close-expired", bytes.NewReader(nil)); if err != nil { continue } - request.Header.Set("X-Heqi-Worker-Token", config.Spec.PaymentAPI.Token) - response, err := http.DefaultClient.Do(request); if err == nil { _ = response.Body.Close() } - } } + ticker := time.NewTicker(time.Duration(config.Spec.PaymentAPI.IntervalSeconds) * time.Second) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + request, err := http.NewRequestWithContext(ctx, http.MethodPost, config.Spec.PaymentAPI.BaseURL+"/heqi/internal/v1/payment/close-expired", bytes.NewReader(nil)) + if err != nil { + continue + } + request.Header.Set("X-Heqi-Worker-Token", config.Spec.PaymentAPI.Token) + response, err := http.DefaultClient.Do(request) + if err == nil { + _ = response.Body.Close() + } + } + } +} + +type iotOutbox struct { + Identity string `json:"identity"` + CommandIdentity string `json:"command_identity"` + Payload string `json:"payload"` +} + +func dispatchIoTOutbox(ctx context.Context) { + ticker := time.NewTicker(time.Duration(config.Spec.IoTAPI.IntervalMilliseconds) * time.Millisecond) + defer ticker.Stop() + client := &http.Client{Timeout: 15 * time.Second} + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + request, err := http.NewRequestWithContext(ctx, http.MethodGet, config.Spec.IoTAPI.PlatformBaseURL+"/heqi/internal/v1/iot/outbox/next", nil) + if err != nil { + continue + } + request.Header.Set("X-Heqi-Iot-Token", config.Spec.IoTAPI.Token) + response, err := client.Do(request) + if err != nil { + continue + } + if response.StatusCode == http.StatusNoContent { + _ = response.Body.Close() + continue + } + if response.StatusCode != http.StatusOK { + _ = response.Body.Close() + continue + } + var outbox iotOutbox + err = json.NewDecoder(io.LimitReader(response.Body, 1<<20)).Decode(&outbox) + _ = response.Body.Close() + if err != nil { + continue + } + var payload map[string]any + if json.Unmarshal([]byte(outbox.Payload), &payload) != nil { + completeIoTOutbox(ctx, client, outbox.Identity, false, 0, "IOT_OUTBOX_PAYLOAD_INVALID") + continue + } + payload["identity"] = outbox.CommandIdentity + body, _ := json.Marshal(payload) + dispatch, err := http.NewRequestWithContext(ctx, http.MethodPost, config.Spec.IoTAPI.ClientBaseURL+"/v1/device-commands", bytes.NewReader(body)) + if err != nil { + continue + } + dispatch.Header.Set("Content-Type", "application/json") + dispatch.Header.Set("X-Heqi-Iot-Token", config.Spec.IoTAPI.Token) + dispatchResponse, err := client.Do(dispatch) + if err != nil { + completeIoTOutbox(ctx, client, outbox.Identity, false, 0, "IOT_CLIENT_UNAVAILABLE") + continue + } + var result struct { + PacketNumber uint16 `json:"packet_number"` + } + decodeErr := json.NewDecoder(io.LimitReader(dispatchResponse.Body, 1<<20)).Decode(&result) + status := dispatchResponse.StatusCode + _ = dispatchResponse.Body.Close() + success := (status == http.StatusAccepted || status == http.StatusOK) && decodeErr == nil + errorCode := "" + if !success { + errorCode = "IOT_CLIENT_REJECTED" + } + completeIoTOutbox(ctx, client, outbox.Identity, success, result.PacketNumber, errorCode) + } + } +} +func completeIoTOutbox(ctx context.Context, client *http.Client, identity string, success bool, packet uint16, errorCode string) { + body, _ := json.Marshal(map[string]any{"success": success, "packet_number": packet, "error_code": errorCode}) + request, err := http.NewRequestWithContext(ctx, http.MethodPost, config.Spec.IoTAPI.PlatformBaseURL+"/heqi/internal/v1/iot/outbox/"+identity+"/result", bytes.NewReader(body)) + if err != nil { + return + } + request.Header.Set("Content-Type", "application/json") + request.Header.Set("X-Heqi-Iot-Token", config.Spec.IoTAPI.Token) + response, err := client.Do(request) + if err == nil { + _ = response.Body.Close() + } } diff --git a/backend/worker/etc/platform_worker_dev.yaml b/backend/worker/etc/platform_worker_dev.yaml index cbd227a..3c4243a 100644 --- a/backend/worker/etc/platform_worker_dev.yaml +++ b/backend/worker/etc/platform_worker_dev.yaml @@ -11,4 +11,9 @@ PaymentAPI: BaseURL: http://localhost:12426 Token: change-me-payment-worker-token IntervalSeconds: 60 +IoTAPI: + PlatformBaseURL: http://127.0.0.1:12426 + ClientBaseURL: http://127.0.0.1:12429 + Token: change-me-iot-internal-token + IntervalMilliseconds: 500 SecretKey: change-me-to-a-random-string diff --git a/backend/worker/internal/config/config.go b/backend/worker/internal/config/config.go index 323e6db..cafdc9c 100644 --- a/backend/worker/internal/config/config.go +++ b/backend/worker/internal/config/config.go @@ -2,31 +2,32 @@ package config import ( - "net" - "git.apinb.com/bsm-sdk/core/conf" + "net" ) -// Spec 是 Worker 运行配置。 var Spec SrvConfig -// SrvConfig 保持与 API 进程相同的 BSM 配置结构。 type SrvConfig struct { conf.Base `yaml:",inline"` Databases *conf.DBConf `yaml:"Databases"` Rpc map[string]conf.RpcConf `yaml:"Rpc"` Apm *conf.ApmConf `yaml:"APM"` PaymentAPI PaymentAPIConfig `yaml:"PaymentAPI"` + IoTAPI IoTAPIConfig `yaml:"IoTAPI"` } - -// PaymentAPIConfig 保存 Worker 调用支付内部动作所需的最小配置。 type PaymentAPIConfig struct { BaseURL string `yaml:"BaseURL"` Token string `yaml:"Token"` IntervalSeconds int `yaml:"IntervalSeconds"` } +type IoTAPIConfig struct { + PlatformBaseURL string `yaml:"PlatformBaseURL"` + ClientBaseURL string `yaml:"ClientBaseURL"` + Token string `yaml:"Token"` + IntervalMilliseconds int `yaml:"IntervalMilliseconds"` +} -// New 初始化 Worker 配置。 func New(srvKey string) { conf.New(srvKey, &Spec) Spec.Port = conf.CheckPort(Spec.Port) @@ -36,5 +37,8 @@ func New(srvKey string) { if Spec.PaymentAPI.BaseURL == "" || Spec.PaymentAPI.Token == "" || Spec.PaymentAPI.IntervalSeconds <= 0 { panic("PaymentAPI configuration is required") } + if Spec.IoTAPI.PlatformBaseURL == "" || Spec.IoTAPI.ClientBaseURL == "" || Spec.IoTAPI.Token == "" || Spec.IoTAPI.IntervalMilliseconds <= 0 { + panic("IoTAPI configuration is required") + } conf.PrintInfo(Spec.Addr) } diff --git a/checking/02-总体测试计划.md b/checking/02-总体测试计划.md index ed57406..e7175ad 100644 --- a/checking/02-总体测试计划.md +++ b/checking/02-总体测试计划.md @@ -23,7 +23,7 @@ | Web | `frontend/site` | 官网内容、导航、响应式、静态资源、可访问性、SEO 基础、构建与托管 Worker | | 后端 | `backend/api` | Gin API/BFF、JWT、权限、同步事务、状态动作、金额、幂等、审计、上传和数据库事实 | | 后端 | `backend/worker` | 配置、启动、Mock 循环、优雅退出;未来 Outbox/Streams 的投递、重试、死信和积压 | -| 后端 | `backend/iot` | 配置、启动、Mock 适配、命令/回执边界;未来 MQTT TLS、身份、遥测和断线重连 | +| 后端 | `backend/iot-client`、`backend/iot-server` | 配置、启动、幂等命令、协议编解码、MQTT TLS、身份、遥测、回执和断线重连 | 不修改 `sample/`。真实支付、短信、地图、对象存储、电子合同、MQTT 厂商接入等,以受控沙箱或 Mock 验证;生产联调另设上线前检查点。 @@ -138,4 +138,3 @@ - 所有受影响模块的静态检查、单元、契约和构建通过;关键链路 E2E 通过。 - 设备、安全、订单、支付、提现、检查、整改、敏感访问的审计可按 `request_id`/`identity` 串联。 - 安全、性能、弱网/离线和恢复专项完成;未具备真实依赖的项目明确标记“受限通过”,不得写成已完成生产验收。 - diff --git a/checking/05-执行批次与质量门禁.md b/checking/05-执行批次与质量门禁.md index caee43e..0c277b4 100644 --- a/checking/05-执行批次与质量门禁.md +++ b/checking/05-执行批次与质量门禁.md @@ -56,7 +56,8 @@ npm run build ```bash cd backend/api && go test ./... && go vet ./... && go build ./cmd/main/main.go cd backend/worker && go test ./... && go vet ./... && go build ./cmd/main/main.go -cd backend/iot && go test ./... && go vet ./... && go build ./cmd/main/main.go +cd backend/iot-client && go test ./... && go vet ./... && go build ./cmd/main/main.go +cd backend/iot-server && go test ./... && go vet ./... && go build ./cmd/main/main.go ``` 建议 CI 增加 `go test -race ./...`(运行环境支持时)、覆盖率输出、依赖/密钥扫描和 API E2E 作业。命令只是门禁,不等价于完整业务验收。 @@ -122,4 +123,3 @@ cd backend/iot && go test ./... && go vet ./... && go build ./cmd/main/main.go - 自动化通过率、偶发失败率、覆盖率变化和执行时长。 - API 性能分位数/错误率、资源峰值、安全发现、恢复实测值。 - 环境限制、未测范围、残余风险、责任人和下一步日期。 - diff --git a/checking/README.md b/checking/README.md index 2f816eb..0759151 100644 --- a/checking/README.md +++ b/checking/README.md @@ -24,7 +24,6 @@ ## 计划边界 -- `backend/worker` 和 `backend/iot` 当前仍是 Mock 边界。计划会验证其进程、配置和失败行为,但真实 Redis Streams、MQTT Broker、设备证书及厂商协议的生产验收必须等契约和环境就绪后执行。 +- `backend/worker` 已投递 IoT Outbox;`backend/iot-client` 与 `backend/iot-server` 已落地系统接口和 MQTT 协议边界。真实 Broker、设备证书和厂商硬件仍需在联调环境验收。 - `frontend/site` 是官网,按展示、响应式、可访问性、链接、构建和托管 Worker 测试;不把它当成业务事实写入端。 - 本计划不把尚未实现或文档中“待确认”的政策当作通过标准。此类项进入阻塞/待决清单,由产品、安全、财务或法务确认后再固化用例。 - diff --git a/docs/10-技术实现规划.md b/docs/10-技术实现规划.md index 14e109c..2182289 100644 --- a/docs/10-技术实现规划.md +++ b/docs/10-技术实现规划.md @@ -53,7 +53,8 @@ flowchart LR | --- | --- | --- | | `api` | 用户端和五个管理系统的 HTTP API、鉴权、同步业务事务 | 无状态部署;所有写操作支持幂等键、事务和审计 | | `worker` | 事件消费、告警、派单、推送、对账、超时扫描、轨迹异常识别 | Redis Streams 消费组;重试、死信、幂等消费和可观测的积压告警 | -| `iot` | MQTT 设备会话、协议适配、遥测校验、命令下发与回执 | 保持设备会话一致性;命令/回执持久化;协议版本和设备身份校验 | +| `iot-client` | 管理系统/App 的设备上下行 HTTP 边界、幂等受理与状态查询 | 不绕过 API 的 JWT、角色、对象归属与安全状态校验 | +| `iot-server` | 外部 MQTT Broker 会话、厂商协议适配、遥测校验、命令下发与回执 | 保持设备会话一致性;协议版本、设备身份、LRC8 与 AES 校验 | 关键业务采用“数据库事务 + Outbox 事件表 + Worker 投递”的模式:先在 PostgreSQL 提交业务事实与待投递事件,再异步写入 Redis Streams。这样 Redis 故障或 Worker 重启不会丢失订单、告警、支付或设备命令的业务事实。 @@ -62,7 +63,7 @@ flowchart LR | 基线 | 路径 | 使用要求 | | --- | --- | --- | | 前端工程基线 | 现有 Vue 管理端 | Vue 管理系统统一前端框架、路由、状态管理、请求封装、权限指令、表格表单、主题、错误处理、国际化与测试规范 | -| 后端工程基线 | `backend/{api,worker,iot}` | Go API、Worker、IoT 进程统一沿用配置、日志、错误码、认证、数据库访问、任务、测试和发布规范 | +| 后端工程基线 | `backend/{api,worker,iot-client,iot-server}` | Go API、Worker 与两个 IoT 进程统一沿用配置、日志、错误码、认证、任务、测试和发布规范 | 业务项目应通过共享包、模板或上游同步机制复用标准库,禁止将标准库目录复制到每个子项目后自行漂移。标准库升级需要记录版本、影响范围、兼容策略和回滚方式。 @@ -97,7 +98,8 @@ platforms/ backend/ api/ # Go HTTP API、BFF、同步领域事务 worker/ # Go 异步任务:派单、告警、通知、对账、超时扫描 - iot/ # Go MQTT 协议适配、设备命令、遥测与回执 + iot-client/ # 管理系统/App 的设备上下行接口 + iot-server/ # MQTT 会话、厂商协议、设备命令、遥测与回执 contracts/ openapi/ # HTTP API 契约及生成配置 asyncapi/ # MQTT/Redis Streams 事件契约与 Schema @@ -112,7 +114,7 @@ platforms/ performance/ # 遥测、订单、轨迹与消息积压压测 ``` -其中 `apps/user_app`、`apps/service_app`、`backend/{api,worker,iot}` 与当前管理端目录已经落地;其余标记为规划的目录仍不得因局部任务提前创建空壳。 +其中 `apps/user_app`、`apps/service_app`、`backend/{api,worker,iot-client,iot-server}` 与当前管理端目录已经落地;其余标记为规划的目录仍不得因局部任务提前创建空壳。 ## 6. 后端领域划分 diff --git a/docs/11-数据接口与安全.md b/docs/11-数据接口与安全.md index 9a8f223..a0c8239 100644 --- a/docs/11-数据接口与安全.md +++ b/docs/11-数据接口与安全.md @@ -43,6 +43,10 @@ ## 3. IoT 协议与可靠性 +- 设备厂商 V1.8 二进制帧保持 `0x5E` 起始、`0x5B` 结束、大端序、数据包 AES-128 与 LRC8 规则不变,并作为 MQTT payload 传输;Topic 使用 `devices/{deviceId}/{up|down|ack}`,QoS 1,下行命令禁止 retained。 +- `iot-server` 只处理 MQTT 会话和设备协议,`iot-client` 只提供系统侧上下行接口;命令、Outbox、原始上行和回执事实由 API 持久化,Worker 负责可重试投递。 +- 协议封面版本与变更记录冲突时以最新 V1.8 变更记录和绿色标注为兼容实现依据;重复子标识等歧义必须保留原始报文并按设备型号配置解析,不得静默猜测。 + - 设备采用 MQTT over TLS,设备身份使用每设备证书或短期轮换令牌;禁止共享默认密钥。 - 上行消息至少包含设备 ID、协议版本、消息 ID、设备时间、服务端接收时间、指标值、质量标记和固件版本。 - 下行命令包含命令 ID、幂等键、期望状态、过期时间、签名/鉴权信息;设备回传已收到、执行中、成功/失败与错误码。 diff --git a/docs/12-验收与迭代规划.md b/docs/12-验收与迭代规划.md index 8b4e8c6..f1da0a4 100644 --- a/docs/12-验收与迭代规划.md +++ b/docs/12-验收与迭代规划.md @@ -39,6 +39,7 @@ | AC-22 | 设备共享与紧急联系人 | 设备所有者可独立授予或撤销成员查看/控制权限;高风险告警自动关阀后仅通知已授权紧急联系人,且审计完整 | | AC-23 | 二维码与轨迹隐私 | 二维码不包含用户隐私、账号凭证或接口密钥;用户只看本人订单简化轨迹,精确轨迹回放/导出须经审批、水印和审计,非履约位置不可访问 | | AC-24 | 钱包实体命名一致性 | 数据库表、Go/Flutter/Vue 模型和契约均使用 `wal_wallet_ledger`;不得出现 `wal_ledger`、复数表名或同义钱包流水实体 | +| AC-25 | MQTT 设备命令与协议回执 | API 事务写入命令和 Outbox;Worker 可重试投递;MQTT 下行包含命令标识、幂等键和过期时间;无设备回执保持待确认,重复请求不重复下发;原始上行、设备时间和服务端接收时间可追溯 | ## 3. 非功能验收 diff --git a/scripts/build-backend.sh b/scripts/build-backend.sh index 139a513..139dce6 100644 --- a/scripts/build-backend.sh +++ b/scripts/build-backend.sh @@ -1,5 +1,7 @@ #!/bin/bash -cd ./backend/api - -GOARCH=amd64 GOOS=linux go build -o ../../output/heqiapp ./cmd/main/main.go +set -Eeuo pipefail +mkdir -p ./output +for module in api worker iot-client iot-server; do + (cd "./backend/${module}" && GOARCH=amd64 GOOS=linux go build -o "../../output/${module}" ./cmd/main/main.go) +done diff --git a/scripts/run_iot.sh b/scripts/run_iot.sh new file mode 100644 index 0000000..a4e9383 --- /dev/null +++ b/scripts/run_iot.sh @@ -0,0 +1,11 @@ +#!/usr/bin/env bash +set -Eeuo pipefail +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +PROJECT_ROOT="$(cd "${SCRIPT_DIR}/.." && pwd)" +PIDS=() +cleanup(){ trap - EXIT INT TERM; for pid in "${PIDS[@]}"; do kill "${pid}" 2>/dev/null || true; done; for pid in "${PIDS[@]}"; do wait "${pid}" 2>/dev/null || true; done; } +trap cleanup EXIT INT TERM +(cd "${PROJECT_ROOT}/backend/iot-server" && go run ./cmd/main/main.go -config etc/platform_iot_server_dev.yaml) & PIDS+=("$!") +(cd "${PROJECT_ROOT}/backend/iot-client" && go run ./cmd/main/main.go -config etc/platform_iot_client_dev.yaml) & PIDS+=("$!") +echo "IoT Server :12428 与 IoT Client :12429 已启动。" +wait -n "${PIDS[@]}"