This commit is contained in:
@@ -6,6 +6,7 @@ import (
|
||||
"go.uber.org/zap"
|
||||
|
||||
"github.com/break/junhong_cmp_fiber/internal/polling"
|
||||
assetAuditSvc "github.com/break/junhong_cmp_fiber/internal/service/asset_audit"
|
||||
iotCardSvc "github.com/break/junhong_cmp_fiber/internal/service/iot_card"
|
||||
"github.com/break/junhong_cmp_fiber/internal/store/postgres"
|
||||
"github.com/break/junhong_cmp_fiber/pkg/constants"
|
||||
@@ -21,6 +22,7 @@ type AssetPollingService struct {
|
||||
iotCardService *iotCardSvc.Service
|
||||
queueMgr *polling.PollingQueueManager
|
||||
logger *zap.Logger
|
||||
assetAuditService assetAuditSvc.OperationLogger
|
||||
}
|
||||
|
||||
// NewAssetPollingService 创建资产轮询管控服务
|
||||
@@ -30,6 +32,7 @@ func NewAssetPollingService(
|
||||
iotCardService *iotCardSvc.Service,
|
||||
queueMgr *polling.PollingQueueManager,
|
||||
logger *zap.Logger,
|
||||
assetAuditService assetAuditSvc.OperationLogger,
|
||||
) *AssetPollingService {
|
||||
return &AssetPollingService{
|
||||
deviceStore: deviceStore,
|
||||
@@ -37,9 +40,23 @@ func NewAssetPollingService(
|
||||
iotCardService: iotCardService,
|
||||
queueMgr: queueMgr,
|
||||
logger: logger,
|
||||
assetAuditService: assetAuditService,
|
||||
}
|
||||
}
|
||||
|
||||
func (s *AssetPollingService) logAssetPollingAudit(ctx context.Context, p assetAuditSvc.BuildLogParams) {
|
||||
if s == nil || s.assetAuditService == nil {
|
||||
return
|
||||
}
|
||||
if p.OperationType == "" {
|
||||
p.OperationType = constants.AssetAuditOpAssetPollingStatus
|
||||
}
|
||||
if p.Operator.Type == "" {
|
||||
p.Operator = assetAuditSvc.OperatorFromContext(ctx)
|
||||
}
|
||||
s.assetAuditService.LogOperation(ctx, assetAuditSvc.BuildLog(ctx, p))
|
||||
}
|
||||
|
||||
// UpdatePollingStatus 更新资产轮询状态
|
||||
// assetType: "card" 或 "device"
|
||||
// assetID: 资产ID
|
||||
@@ -47,12 +64,87 @@ func NewAssetPollingService(
|
||||
func (s *AssetPollingService) UpdatePollingStatus(ctx context.Context, assetType string, assetID uint, enablePolling bool) error {
|
||||
switch assetType {
|
||||
case constants.AssetTypeIotCard:
|
||||
beforeData := map[string]any{
|
||||
"asset_type": constants.AssetTypeIotCard,
|
||||
"asset_id": assetID,
|
||||
"enable_polling": "unknown",
|
||||
"source_service": "asset_polling",
|
||||
}
|
||||
afterData := map[string]any{
|
||||
"asset_type": constants.AssetTypeIotCard,
|
||||
"asset_id": assetID,
|
||||
"enable_polling": enablePolling,
|
||||
}
|
||||
// S2 修复:委托给 IotCardService,确保 DB 写入 + PollingCallback 回调一并触发
|
||||
return s.iotCardService.UpdatePollingStatus(ctx, assetID, enablePolling)
|
||||
if err := s.iotCardService.UpdatePollingStatus(ctx, assetID, enablePolling); err != nil {
|
||||
errorCode, errorMsg := assetAuditSvc.BuildErrorInfo(err)
|
||||
s.logAssetPollingAudit(ctx, assetAuditSvc.BuildLogParams{
|
||||
AssetType: constants.AssetTypeIotCard,
|
||||
AssetID: assetID,
|
||||
OperationDesc: "统一入口更新轮询状态失败",
|
||||
ResultStatus: constants.AssetAuditResultFailed,
|
||||
ErrorCode: errorCode,
|
||||
ErrorMsg: errorMsg,
|
||||
BeforeData: beforeData,
|
||||
AfterData: afterData,
|
||||
})
|
||||
return err
|
||||
}
|
||||
s.logAssetPollingAudit(ctx, assetAuditSvc.BuildLogParams{
|
||||
AssetType: constants.AssetTypeIotCard,
|
||||
AssetID: assetID,
|
||||
OperationDesc: "统一入口更新轮询状态",
|
||||
ResultStatus: constants.AssetAuditResultSuccess,
|
||||
BeforeData: beforeData,
|
||||
AfterData: afterData,
|
||||
})
|
||||
return nil
|
||||
|
||||
case constants.AssetTypeDevice:
|
||||
device, getErr := s.deviceStore.GetByID(ctx, assetID)
|
||||
if getErr != nil {
|
||||
errorCode, errorMsg := assetAuditSvc.BuildErrorInfo(getErr)
|
||||
s.logAssetPollingAudit(ctx, assetAuditSvc.BuildLogParams{
|
||||
AssetType: constants.AssetTypeDevice,
|
||||
AssetID: assetID,
|
||||
OperationDesc: "统一入口更新轮询状态失败",
|
||||
ResultStatus: constants.AssetAuditResultFailed,
|
||||
ErrorCode: errorCode,
|
||||
ErrorMsg: errorMsg,
|
||||
AfterData: map[string]any{
|
||||
"asset_type": constants.AssetTypeDevice,
|
||||
"asset_id": assetID,
|
||||
"enable_polling": enablePolling,
|
||||
},
|
||||
})
|
||||
return getErr
|
||||
}
|
||||
beforeData := map[string]any{
|
||||
"asset_type": constants.AssetTypeDevice,
|
||||
"asset_id": device.ID,
|
||||
"asset_identifier": device.VirtualNo,
|
||||
"enable_polling": device.EnablePolling,
|
||||
}
|
||||
afterData := map[string]any{
|
||||
"asset_type": constants.AssetTypeDevice,
|
||||
"asset_id": device.ID,
|
||||
"asset_identifier": device.VirtualNo,
|
||||
"enable_polling": enablePolling,
|
||||
}
|
||||
// 1. 更新设备的 enable_polling 字段
|
||||
if err := s.deviceStore.UpdatePollingStatus(ctx, assetID, enablePolling); err != nil {
|
||||
errorCode, errorMsg := assetAuditSvc.BuildErrorInfo(err)
|
||||
s.logAssetPollingAudit(ctx, assetAuditSvc.BuildLogParams{
|
||||
AssetType: constants.AssetTypeDevice,
|
||||
AssetID: device.ID,
|
||||
AssetIdentifier: device.VirtualNo,
|
||||
OperationDesc: "统一入口更新轮询状态失败",
|
||||
ResultStatus: constants.AssetAuditResultFailed,
|
||||
ErrorCode: errorCode,
|
||||
ErrorMsg: errorMsg,
|
||||
BeforeData: beforeData,
|
||||
AfterData: afterData,
|
||||
})
|
||||
return err
|
||||
}
|
||||
bindings, err := s.deviceBindingStore.ListByDeviceID(ctx, assetID)
|
||||
@@ -81,9 +173,33 @@ func (s *AssetPollingService) UpdatePollingStatus(ctx context.Context, assetType
|
||||
}
|
||||
}
|
||||
}
|
||||
s.logAssetPollingAudit(ctx, assetAuditSvc.BuildLogParams{
|
||||
AssetType: constants.AssetTypeDevice,
|
||||
AssetID: device.ID,
|
||||
AssetIdentifier: device.VirtualNo,
|
||||
OperationDesc: "统一入口更新轮询状态",
|
||||
ResultStatus: constants.AssetAuditResultSuccess,
|
||||
BeforeData: beforeData,
|
||||
AfterData: afterData,
|
||||
})
|
||||
return nil
|
||||
|
||||
default:
|
||||
return errors.New(errors.CodeInvalidParam, "资产类型无效,支持 card 或 device")
|
||||
err := errors.New(errors.CodeInvalidParam, "资产类型无效,支持 card 或 device")
|
||||
errorCode, errorMsg := assetAuditSvc.BuildErrorInfo(err)
|
||||
s.logAssetPollingAudit(ctx, assetAuditSvc.BuildLogParams{
|
||||
AssetType: assetType,
|
||||
AssetID: assetID,
|
||||
OperationDesc: "统一入口更新轮询状态被拒绝",
|
||||
ResultStatus: constants.AssetAuditResultDenied,
|
||||
ErrorCode: errorCode,
|
||||
ErrorMsg: errorMsg,
|
||||
AfterData: map[string]any{
|
||||
"asset_type": assetType,
|
||||
"asset_id": assetID,
|
||||
"enable_polling": enablePolling,
|
||||
},
|
||||
})
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user