2026-08-03 13:18:11 +08:00
// Package iot 实现设备命令、Outbox 认领和上行报文持久化。
package iot
import (
"crypto/subtle"
"encoding/json"
2026-08-03 16:03:44 +08:00
"errors"
2026-08-03 13:18:11 +08:00
"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"
}
2026-08-03 16:03:44 +08:00
result := 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 } )
if result . Error != nil {
return result . Error
}
if result . RowsAffected > 1 {
return errors . New ( "ambiguous device acknowledgement" )
}
return nil
2026-08-03 13:18:11 +08:00
} )
if err != nil {
ctx . JSON ( 500 , gin . H { "code" : "IOT_MESSAGE_SAVE_FAILED" } )
return
}
ctx . JSON ( http . StatusAccepted , gin . H { "identity" : message . Identity } )
}