触发通知轮询
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m23s

This commit is contained in:
2026-07-28 14:25:18 +08:00
parent 4aff9937e5
commit 9303af3d46
9 changed files with 228 additions and 1 deletions

View File

@@ -8,6 +8,7 @@ import (
cardObservationApp "github.com/break/junhong_cmp_fiber/internal/application/cardobservation"
"github.com/gofiber/fiber/v2"
"github.com/hibiken/asynq"
dto "github.com/break/junhong_cmp_fiber/internal/model/dto"
packageExpiryQuery "github.com/break/junhong_cmp_fiber/internal/query/packageexpiry"
@@ -20,6 +21,7 @@ import (
"github.com/break/junhong_cmp_fiber/pkg/errors"
"github.com/break/junhong_cmp_fiber/pkg/logger"
"github.com/break/junhong_cmp_fiber/pkg/middleware"
"github.com/break/junhong_cmp_fiber/pkg/queue"
"github.com/break/junhong_cmp_fiber/pkg/response"
"go.uber.org/zap"
)
@@ -37,6 +39,7 @@ type AssetHandler struct {
exchangeTraceQuery AssetExchangeTraceResolver
observationSeries cardObservationApp.BestEffortSeriesDispatcher
packageExpiryQuery *packageExpiryQuery.Query
packageExpiryTrigger func(context.Context) error
}
// SetObservationSeriesDispatcher 注入后台实时状态的观测序列端口。
@@ -49,6 +52,23 @@ func (h *AssetHandler) SetPackageExpiryQuery(query *packageExpiryQuery.Query) {
h.packageExpiryQuery = query
}
// SetPackageExpiryQueue 注入套餐临期提醒任务队列。
func (h *AssetHandler) SetPackageExpiryQueue(client *queue.Client) {
if client == nil {
h.packageExpiryTrigger = nil
return
}
h.packageExpiryTrigger = func(ctx context.Context) error {
return client.EnqueueTask(
ctx,
constants.TaskTypePackageExpiryReminder,
struct{}{},
asynq.MaxRetry(3),
asynq.Timeout(10*time.Minute),
)
}
}
// AssetExchangeTraceResolver 定义资产详情换货链路读取用例。
type AssetExchangeTraceResolver interface {
Resolve(ctx context.Context, assetType string, assetID uint) (*dto.AssetExchangeTrace, error)
@@ -135,6 +155,27 @@ func (h *AssetHandler) ListExpiring(c *fiber.Ctx) error {
})
}
// TriggerPackageExpiryReminder 手动提交套餐临期提醒扫描任务。
// POST /api/admin/expiring-assets/reminder-scan
func (h *AssetHandler) TriggerPackageExpiryReminder(c *fiber.Ctx) error {
if middleware.GetUserTypeFromContext(c.UserContext()) != constants.UserTypeSuperAdmin {
return errors.New(errors.CodeForbidden)
}
if h.packageExpiryTrigger == nil {
return errors.New(errors.CodeServiceUnavailable, "套餐临期提醒任务队列未配置")
}
if err := h.packageExpiryTrigger(c.UserContext()); err != nil {
logger.GetAppLogger().Error("手动提交套餐临期提醒扫描任务失败", zap.Error(err))
return errors.Wrap(errors.CodeTaskQueueError, err, "提交套餐临期提醒扫描任务失败")
}
logger.GetAppLogger().Info("已手动提交套餐临期提醒扫描任务",
zap.Uint("operator_id", middleware.GetUserIDFromContext(c.UserContext())))
return response.Success(c, dto.TriggerPackageExpiryReminderResponse{
TaskType: constants.TaskTypePackageExpiryReminder,
Message: "套餐临期提醒扫描任务已提交",
})
}
// RealtimeStatus 获取资产实时状态
// GET /api/admin/assets/:identifier/realtime-status
func (h *AssetHandler) RealtimeStatus(c *fiber.Ctx) error {