ListPendingCollections
This commit is contained in:
@@ -215,3 +215,51 @@ func (c *Controller) ListTaskExecutionLogs(ctx *gin.Context) {
|
|||||||
c.logger.Infof("%s: 成功, 获取到 %d 条记录, 总计 %d 条", actionType, len(data), total)
|
c.logger.Infof("%s: 成功, 获取到 %d 条记录, 总计 %d 条", actionType, len(data), total)
|
||||||
controller.SendSuccessWithAudit(ctx, controller.CodeSuccess, "获取任务执行日志成功", resp, actionType, "获取任务执行日志成功", req)
|
controller.SendSuccessWithAudit(ctx, controller.CodeSuccess, "获取任务执行日志成功", resp, actionType, "获取任务执行日志成功", req)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ListPendingCollections godoc
|
||||||
|
// @Summary 获取待采集请求列表
|
||||||
|
// @Description 根据提供的过滤条件,分页获取待采集请求
|
||||||
|
// @Tags 数据监控
|
||||||
|
// @Security BearerAuth
|
||||||
|
// @Produce json
|
||||||
|
// @Param query query dto.ListPendingCollectionRequest true "查询参数"
|
||||||
|
// @Success 200 {object} controller.Response{data=dto.ListPendingCollectionResponse}
|
||||||
|
// @Router /api/v1/monitor/pending-collections [get]
|
||||||
|
func (c *Controller) ListPendingCollections(ctx *gin.Context) {
|
||||||
|
const actionType = "获取待采集请求列表"
|
||||||
|
|
||||||
|
var req dto.ListPendingCollectionRequest
|
||||||
|
if err := ctx.ShouldBindQuery(&req); err != nil {
|
||||||
|
c.logger.Errorf("%s: 参数绑定失败: %v", actionType, err)
|
||||||
|
controller.SendErrorWithAudit(ctx, controller.CodeBadRequest, "无效的查询参数: "+err.Error(), actionType, "参数绑定失败", req)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
opts := repository.PendingCollectionListOptions{
|
||||||
|
DeviceID: req.DeviceID,
|
||||||
|
OrderBy: req.OrderBy,
|
||||||
|
StartTime: req.StartTime,
|
||||||
|
EndTime: req.EndTime,
|
||||||
|
}
|
||||||
|
if req.Status != nil {
|
||||||
|
status := models.PendingCollectionStatus(*req.Status)
|
||||||
|
opts.Status = &status
|
||||||
|
}
|
||||||
|
|
||||||
|
data, total, err := c.monitorService.ListPendingCollections(opts, req.Page, req.PageSize)
|
||||||
|
if err != nil {
|
||||||
|
if errors.Is(err, repository.ErrInvalidPagination) {
|
||||||
|
c.logger.Warnf("%s: 无效的分页参数: %v", actionType, err)
|
||||||
|
controller.SendErrorWithAudit(ctx, controller.CodeBadRequest, "无效的分页参数: "+err.Error(), actionType, "无效分页参数", req)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
c.logger.Errorf("%s: 服务层查询失败: %v", actionType, err)
|
||||||
|
controller.SendErrorWithAudit(ctx, controller.CodeInternalError, "获取待采集请求失败: "+err.Error(), actionType, "服务层查询失败", req)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
resp := dto.NewListPendingCollectionResponse(data, total, req.Page, req.PageSize)
|
||||||
|
c.logger.Infof("%s: 成功, 获取到 %d 条记录, 总计 %d 条", actionType, len(data), total)
|
||||||
|
controller.SendSuccessWithAudit(ctx, controller.CodeSuccess, "获取待采集请求成功", resp, actionType, "获取待采集请求成功", req)
|
||||||
|
}
|
||||||
|
|||||||
@@ -247,3 +247,56 @@ func NewListTaskExecutionLogResponse(data []models.TaskExecutionLog, total int64
|
|||||||
},
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --- PendingCollection ---
|
||||||
|
|
||||||
|
// ListPendingCollectionRequest 定义了获取待采集请求列表的请求参数
|
||||||
|
type ListPendingCollectionRequest struct {
|
||||||
|
Page int `form:"page,default=1"`
|
||||||
|
PageSize int `form:"pageSize,default=10"`
|
||||||
|
DeviceID *uint `form:"device_id"`
|
||||||
|
Status *string `form:"status"`
|
||||||
|
StartTime *time.Time `form:"start_time" time_format:"rfc3339"`
|
||||||
|
EndTime *time.Time `form:"end_time" time_format:"rfc3339"`
|
||||||
|
OrderBy string `form:"order_by"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// PendingCollectionDTO 是用于API响应的待采集请求结构
|
||||||
|
type PendingCollectionDTO struct {
|
||||||
|
CorrelationID string `json:"correlation_id"`
|
||||||
|
DeviceID uint `json:"device_id"`
|
||||||
|
CommandMetadata models.UintArray `json:"command_metadata"`
|
||||||
|
Status models.PendingCollectionStatus `json:"status"`
|
||||||
|
FulfilledAt *time.Time `json:"fulfilled_at"`
|
||||||
|
CreatedAt time.Time `json:"created_at"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// ListPendingCollectionResponse 是获取待采集请求列表的响应结构
|
||||||
|
type ListPendingCollectionResponse struct {
|
||||||
|
List []PendingCollectionDTO `json:"list"`
|
||||||
|
Pagination PaginationDTO `json:"pagination"`
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewListPendingCollectionResponse 从模型数据创建列表响应 DTO
|
||||||
|
func NewListPendingCollectionResponse(data []models.PendingCollection, total int64, page, pageSize int) *ListPendingCollectionResponse {
|
||||||
|
dtos := make([]PendingCollectionDTO, len(data))
|
||||||
|
for i, item := range data {
|
||||||
|
dtos[i] = PendingCollectionDTO{
|
||||||
|
CorrelationID: item.CorrelationID,
|
||||||
|
DeviceID: item.DeviceID,
|
||||||
|
CommandMetadata: item.CommandMetadata,
|
||||||
|
Status: item.Status,
|
||||||
|
FulfilledAt: item.FulfilledAt,
|
||||||
|
CreatedAt: item.CreatedAt,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
return &ListPendingCollectionResponse{
|
||||||
|
List: dtos,
|
||||||
|
Pagination: PaginationDTO{
|
||||||
|
Total: total,
|
||||||
|
Page: page,
|
||||||
|
PageSize: pageSize,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ type MonitorService struct {
|
|||||||
sensorDataRepo repository.SensorDataRepository
|
sensorDataRepo repository.SensorDataRepository
|
||||||
deviceCommandLogRepo repository.DeviceCommandLogRepository
|
deviceCommandLogRepo repository.DeviceCommandLogRepository
|
||||||
executionLogRepo repository.ExecutionLogRepository
|
executionLogRepo repository.ExecutionLogRepository
|
||||||
|
pendingCollectionRepo repository.PendingCollectionRepository
|
||||||
// 在这里可以添加其他超表模型的仓库依赖
|
// 在这里可以添加其他超表模型的仓库依赖
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -18,11 +19,13 @@ func NewMonitorService(
|
|||||||
sensorDataRepo repository.SensorDataRepository,
|
sensorDataRepo repository.SensorDataRepository,
|
||||||
deviceCommandLogRepo repository.DeviceCommandLogRepository,
|
deviceCommandLogRepo repository.DeviceCommandLogRepository,
|
||||||
executionLogRepo repository.ExecutionLogRepository,
|
executionLogRepo repository.ExecutionLogRepository,
|
||||||
|
pendingCollectionRepo repository.PendingCollectionRepository,
|
||||||
) *MonitorService {
|
) *MonitorService {
|
||||||
return &MonitorService{
|
return &MonitorService{
|
||||||
sensorDataRepo: sensorDataRepo,
|
sensorDataRepo: sensorDataRepo,
|
||||||
deviceCommandLogRepo: deviceCommandLogRepo,
|
deviceCommandLogRepo: deviceCommandLogRepo,
|
||||||
executionLogRepo: executionLogRepo,
|
executionLogRepo: executionLogRepo,
|
||||||
|
pendingCollectionRepo: pendingCollectionRepo,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -45,3 +48,8 @@ func (s *MonitorService) ListPlanExecutionLogs(opts repository.PlanExecutionLogL
|
|||||||
func (s *MonitorService) ListTaskExecutionLogs(opts repository.TaskExecutionLogListOptions, page, pageSize int) ([]models.TaskExecutionLog, int64, error) {
|
func (s *MonitorService) ListTaskExecutionLogs(opts repository.TaskExecutionLogListOptions, page, pageSize int) ([]models.TaskExecutionLog, int64, error) {
|
||||||
return s.executionLogRepo.ListTaskExecutionLogs(opts, page, pageSize)
|
return s.executionLogRepo.ListTaskExecutionLogs(opts, page, pageSize)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ListPendingCollections 负责处理查询待采集请求列表的业务逻辑
|
||||||
|
func (s *MonitorService) ListPendingCollections(opts repository.PendingCollectionListOptions, page, pageSize int) ([]models.PendingCollection, int64, error) {
|
||||||
|
return s.pendingCollectionRepo.List(opts, page, pageSize)
|
||||||
|
}
|
||||||
|
|||||||
@@ -7,6 +7,15 @@ import (
|
|||||||
"gorm.io/gorm"
|
"gorm.io/gorm"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
// PendingCollectionListOptions 定义了查询待采集请求时的可选参数
|
||||||
|
type PendingCollectionListOptions struct {
|
||||||
|
DeviceID *uint
|
||||||
|
Status *models.PendingCollectionStatus
|
||||||
|
StartTime *time.Time // 基于 created_at 字段
|
||||||
|
EndTime *time.Time // 基于 created_at 字段
|
||||||
|
OrderBy string // 例如 "created_at asc"
|
||||||
|
}
|
||||||
|
|
||||||
// PendingCollectionRepository 定义了与待采集请求相关的数据库操作接口。
|
// PendingCollectionRepository 定义了与待采集请求相关的数据库操作接口。
|
||||||
type PendingCollectionRepository interface {
|
type PendingCollectionRepository interface {
|
||||||
// Create 创建一个新的待采集请求。
|
// Create 创建一个新的待采集请求。
|
||||||
@@ -20,6 +29,9 @@ type PendingCollectionRepository interface {
|
|||||||
|
|
||||||
// MarkAllPendingAsTimedOut 将所有“待处理”请求更新为“已超时”。
|
// MarkAllPendingAsTimedOut 将所有“待处理”请求更新为“已超时”。
|
||||||
MarkAllPendingAsTimedOut() (int64, error)
|
MarkAllPendingAsTimedOut() (int64, error)
|
||||||
|
|
||||||
|
// List 支持分页和过滤的列表查询
|
||||||
|
List(opts PendingCollectionListOptions, page, pageSize int) ([]models.PendingCollection, int64, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
// gormPendingCollectionRepository 是 PendingCollectionRepository 的 GORM 实现。
|
// gormPendingCollectionRepository 是 PendingCollectionRepository 的 GORM 实现。
|
||||||
@@ -65,3 +77,43 @@ func (r *gormPendingCollectionRepository) MarkAllPendingAsTimedOut() (int64, err
|
|||||||
|
|
||||||
return result.RowsAffected, result.Error
|
return result.RowsAffected, result.Error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// List 实现了分页和过滤查询待采集请求的功能
|
||||||
|
func (r *gormPendingCollectionRepository) List(opts PendingCollectionListOptions, page, pageSize int) ([]models.PendingCollection, int64, error) {
|
||||||
|
if page <= 0 || pageSize <= 0 {
|
||||||
|
return nil, 0, ErrInvalidPagination
|
||||||
|
}
|
||||||
|
|
||||||
|
var results []models.PendingCollection
|
||||||
|
var total int64
|
||||||
|
|
||||||
|
query := r.db.Model(&models.PendingCollection{})
|
||||||
|
|
||||||
|
if opts.DeviceID != nil {
|
||||||
|
query = query.Where("device_id = ?", *opts.DeviceID)
|
||||||
|
}
|
||||||
|
if opts.Status != nil {
|
||||||
|
query = query.Where("status = ?", *opts.Status)
|
||||||
|
}
|
||||||
|
if opts.StartTime != nil {
|
||||||
|
query = query.Where("created_at >= ?", *opts.StartTime)
|
||||||
|
}
|
||||||
|
if opts.EndTime != nil {
|
||||||
|
query = query.Where("created_at <= ?", *opts.EndTime)
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := query.Count(&total).Error; err != nil {
|
||||||
|
return nil, 0, err
|
||||||
|
}
|
||||||
|
|
||||||
|
orderBy := "created_at DESC"
|
||||||
|
if opts.OrderBy != "" {
|
||||||
|
orderBy = opts.OrderBy
|
||||||
|
}
|
||||||
|
query = query.Order(orderBy)
|
||||||
|
|
||||||
|
offset := (page - 1) * pageSize
|
||||||
|
err := query.Limit(pageSize).Offset(offset).Find(&results).Error
|
||||||
|
|
||||||
|
return results, total, err
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user