Files
junhong_cmp_fiber/internal/task/audit_monthly_retention.go
break c64f3d8b80
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m31s
全局审计完成
2026-08-07 11:02:52 +08:00

103 lines
3.9 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package task
import (
"context"
"fmt"
"time"
"github.com/bytedance/sonic"
"github.com/hibiken/asynq"
"go.uber.org/zap"
"github.com/break/junhong_cmp_fiber/internal/application/auditarchive"
)
// AuditMonthlyRetentionPayload 是人工补跑月度清理时可选的任务载荷。
type AuditMonthlyRetentionPayload struct {
ArchiveMonth string `json:"archive_month"`
}
// AuditMonthlyRetentionHandler 处理归档完整性门禁与上月在线日志物理清理。
type AuditMonthlyRetentionHandler struct {
service *auditarchive.Service
logger *zap.Logger
cleanupEnabled bool
}
// NewAuditMonthlyRetentionHandler 创建月度日志留存清理处理器。
func NewAuditMonthlyRetentionHandler(service *auditarchive.Service, logger *zap.Logger, cleanupEnabled bool) *AuditMonthlyRetentionHandler {
return &AuditMonthlyRetentionHandler{service: service, logger: logger, cleanupEnabled: cleanupEnabled}
}
// Handle 校验整月归档后按固定顺序分批物理删除 PostgreSQL 在线日志。
func (h *AuditMonthlyRetentionHandler) Handle(ctx context.Context, task *asynq.Task) error {
if !h.cleanupEnabled {
return h.handleDryRun(ctx, task)
}
if h.service == nil {
return fmt.Errorf("月度日志留存清理服务未配置")
}
startedAt := time.Now()
var result auditarchive.RetentionResult
var err error
if len(task.Payload()) == 0 {
result, err = h.service.CleanupPreviousMonth(ctx)
} else {
var payload AuditMonthlyRetentionPayload
if unmarshalErr := sonic.Unmarshal(task.Payload(), &payload); unmarshalErr != nil {
return fmt.Errorf("解析月度日志留存清理任务载荷失败: %w", unmarshalErr)
}
month, parseErr := parseArchiveMonth(payload.ArchiveMonth)
if parseErr != nil {
return parseErr
}
result, err = h.service.CleanupMonth(ctx, month)
}
fields := []zap.Field{
zap.String("archive_month", result.Month), zap.Int64("audit_event_count", result.EventCount),
zap.Int64("event_resource_count", result.ResourceCount), zap.Int64("integration_log_count", result.IntegrationCount),
zap.Duration("duration", time.Since(startedAt)), zap.Int("manifest_count", len(result.ManifestKeys)),
}
if err != nil {
fields = append(fields, zap.String("severity", "critical"), zap.Error(err))
h.logger.Error("月度日志留存清理失败PostgreSQL 整月清理已阻断或等待断点续跑", fields...)
return err
}
h.logger.Info("月度日志留存清理完成", fields...)
return nil
}
func (h *AuditMonthlyRetentionHandler) handleDryRun(ctx context.Context, task *asynq.Task) error {
if h.service == nil {
return fmt.Errorf("月度日志留存演练服务未配置")
}
var result auditarchive.RetentionResult
var err error
if len(task.Payload()) == 0 {
result, err = h.service.ValidatePreviousMonth(ctx)
} else {
var payload AuditMonthlyRetentionPayload
if unmarshalErr := sonic.Unmarshal(task.Payload(), &payload); unmarshalErr != nil {
return fmt.Errorf("解析月度日志留存演练任务载荷失败: %w", unmarshalErr)
}
month, parseErr := parseArchiveMonth(payload.ArchiveMonth)
if parseErr != nil {
return parseErr
}
result, err = h.service.ValidateMonth(ctx, month)
}
fields := []zap.Field{
zap.Bool("cleanup_enabled", false), zap.String("archive_month", result.Month),
zap.Int64("audit_event_count", result.EventCount), zap.Int64("event_resource_count", result.ResourceCount),
zap.Int64("integration_log_count", result.IntegrationCount), zap.Int64("estimated_cleanup_batches", result.EstimatedBatches),
zap.Duration("validation_duration", result.Duration), zap.Int("manifest_count", len(result.ManifestKeys)),
}
if err != nil {
fields = append(fields, zap.String("severity", "critical"), zap.Error(err))
h.logger.Error("审计日志月度只读演练失败,物理清理保持关闭", fields...)
return err
}
h.logger.Info("审计日志月度只读演练通过,物理清理保持关闭", fields...)
return nil
}