Files
cid/service/check/material_verify_service.go
2026-08-25 10:14:32 +08:00

847 lines
30 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 check
import (
consts "cid/consts/check"
dao "cid/dao/check"
entity "cid/model/entity/check"
"context"
"encoding/json"
"fmt"
"time"
"gitea.redpowerfuture.com/red-future/common/beans"
"github.com/gogf/gf/v2/frame/g"
)
// 轮询配置常量
const (
// PollBatchSize 每次轮询处理数量
PollBatchSize = 20
)
// MaterialVerifyService 素材校验服务
type MaterialVerifyService struct{}
// MaterialVerify 校验服务单例
var MaterialVerify = new(MaterialVerifyService)
// =============================================================================
// 校验状态转换
// =============================================================================
// SuggestionToVerifyStatus 根据易盾处置建议转换为校验状态
func SuggestionToVerifyStatus(suggestion int) string {
switch suggestion {
case consts.SuggestionPass:
return entity.VerifyStatusVerified // 通过
case consts.SuggestionReview:
return entity.VerifyStatusReview // 嫌疑,需人工复核
case consts.SuggestionBlock:
return entity.VerifyStatusRejected // 不通过
default:
return entity.VerifyStatusPending
}
}
// =============================================================================
// 图片校验
// =============================================================================
// VerifyImageByID 根据图片ID执行校验
// 使用原子 CAS 防止并发重复送检
func (s *MaterialVerifyService) VerifyImageByID(ctx context.Context, imageID string) (*entity.MaterialVerifyLog, error) {
image, err := dao.TencentImage.GetByImageID(ctx, imageID)
if err != nil {
return nil, fmt.Errorf("查询图片数据失败, imageID=%s: %w", imageID, err)
}
if image == nil {
return nil, fmt.Errorf("未找到图片数据, imageID=%s", imageID)
}
// 原子 CASPENDING → SUBMITTING,确保只有第一个调用者能抢到处理权
claimed, err := dao.TencentImage.ClaimPending(ctx, image.Id)
if err != nil {
return nil, fmt.Errorf("认领图片送检失败, imageID=%s: %w", imageID, err)
}
if !claimed {
logs, err := dao.MaterialVerifyLog.GetByMaterialID(ctx, imageID)
if err == nil && len(logs) > 0 {
g.Log().Infof(ctx, "图片已被其他进程送检, imageID=%s, logId=%d", imageID, logs[0].Id)
return &logs[0], nil
}
return nil, fmt.Errorf("图片正在送检中且无校验日志, imageID=%s", imageID)
}
log := s.createVerifyLog(ctx, entity.MaterialTypeImage, imageID, consts.SourceTableTencentImage, image.Id, image.AccountID)
if log == nil {
dao.TencentImage.UpdateStatus(ctx, image.Id, entity.VerifyStatusPending)
return nil, fmt.Errorf("创建校验日志失败")
}
err = s.submitImageCheck(ctx, image, log)
if err != nil {
dao.TencentImage.UpdateStatus(ctx, image.Id, entity.VerifyStatusPending)
return nil, err
}
return log, nil
}
// submitImageCheck 提交图片校验
func (s *MaterialVerifyService) submitImageCheck(ctx context.Context, image *entity.TencentImage, log *entity.MaterialVerifyLog) error {
startTime := time.Now()
callbackMode := g.Cfg().MustGet(ctx, "check.callback_mode").Bool()
requestParams := map[string]interface{}{
"imageURL": image.PreviewURL,
"dataID": image.ImageID,
}
requestParamsJSON, _ := json.Marshal(requestParams)
var (
taskID string
duration int64
)
if callbackMode {
callbackURL := g.Cfg().MustGet(ctx, "check.image.callback_url").String()
requestParams["callbackURL"] = callbackURL
result, err := ImageDetection.DetectImage(ctx, image.PreviewURL, image.ImageID, callbackURL)
duration = time.Since(startTime).Milliseconds()
if err != nil {
dao.MaterialVerifyLog.UpdateError(ctx, log.Id, entity.VerifyStatusPending, err.Error())
dao.MaterialVerifyLog.UpdateDuration(ctx, log.Id, duration)
g.Log().Warningf(ctx, "图片异步检测失败(保持待检验), id=%d, imageId=%s, error=%v", image.Id, image.ImageID, err)
return fmt.Errorf("图片异步检测提交失败, imageId=%s: %w", image.ImageID, err)
}
taskID = result.TaskID
dao.MaterialVerifyLog.UpdateTaskID(ctx, log.Id, taskID)
dao.MaterialVerifyLog.UpdateRequestParams(ctx, log.Id, string(requestParamsJSON))
TencentContentCheck.writeAuditLog(ctx, consts.SourceTableTencentImage, image.Id, image.ImageID, image.PreviewURL, taskID, -1, 0, 0, "", duration)
g.Log().Infof(ctx, "图片异步检测已提交, id=%d, imageId=%s, taskId=%s, duration=%dms",
image.Id, image.ImageID, taskID, duration)
} else {
syncResult, err := ImageDetection.DetectImageSync(ctx, image.PreviewURL, image.ImageID)
duration = time.Since(startTime).Milliseconds()
if err != nil {
dao.MaterialVerifyLog.UpdateError(ctx, log.Id, entity.VerifyStatusPending, err.Error())
dao.MaterialVerifyLog.UpdateDuration(ctx, log.Id, duration)
g.Log().Warningf(ctx, "图片同步检测失败(保持待检验), id=%d, imageId=%s, error=%v", image.Id, image.ImageID, err)
return fmt.Errorf("图片同步检测失败, imageId=%s: %w", image.ImageID, err)
}
taskID = syncResult.TaskID
dao.MaterialVerifyLog.UpdateTaskID(ctx, log.Id, taskID)
dao.MaterialVerifyLog.UpdateRequestParams(ctx, log.Id, string(requestParamsJSON))
verifyStatus := SuggestionToVerifyStatus(syncResult.Suggestion)
responseJSON, _ := json.Marshal(syncResult)
dao.MaterialVerifyLog.UpdateVerifyResult(ctx, log.Id, verifyStatus,
syncResult.Suggestion, syncResult.Label, syncResult.ResultType, string(responseJSON), syncResult.CensorTime)
s.updateImageStatus(ctx, image.Id, verifyStatus)
TencentContentCheck.writeAuditLog(ctx, consts.SourceTableTencentImage, image.Id, image.ImageID, image.PreviewURL, taskID, syncResult.Suggestion, syncResult.Label, syncResult.ResultType, string(responseJSON), duration)
g.Log().Infof(ctx, "图片同步检测完成, id=%d, imageId=%s, taskId=%s, suggestion=%d, verifyStatus=%s, duration=%dms",
image.Id, image.ImageID, taskID, syncResult.Suggestion, verifyStatus, duration)
}
return nil
}
// =============================================================================
// 视频校验
// =============================================================================
// VerifyVideoByID 根据视频ID执行校验
// 使用原子 CAS 防止并发重复送检
func (s *MaterialVerifyService) VerifyVideoByID(ctx context.Context, videoID string) (*entity.MaterialVerifyLog, error) {
video, err := dao.TencentVideo.GetByVideoID(ctx, videoID)
if err != nil {
return nil, fmt.Errorf("查询视频数据失败, videoID=%s: %w", videoID, err)
}
if video == nil {
return nil, fmt.Errorf("未找到视频数据, videoID=%s", videoID)
}
claimed, err := dao.TencentVideo.ClaimPending(ctx, video.Id)
if err != nil {
return nil, fmt.Errorf("认领视频送检失败, videoID=%s: %w", videoID, err)
}
if !claimed {
logs, err := dao.MaterialVerifyLog.GetByMaterialID(ctx, videoID)
if err == nil && len(logs) > 0 {
g.Log().Infof(ctx, "视频已被其他进程送检, videoID=%s, logId=%d", videoID, logs[0].Id)
return &logs[0], nil
}
return nil, fmt.Errorf("视频正在送检中且无校验日志, videoID=%s", videoID)
}
log := s.createVerifyLog(ctx, entity.MaterialTypeVideo, videoID, consts.SourceTableTencentVideo, video.Id, video.AccountID)
if log == nil {
dao.TencentVideo.UpdateStatus(ctx, video.Id, entity.VerifyStatusPending)
return nil, fmt.Errorf("创建校验日志失败")
}
err = s.submitVideoCheck(ctx, video, log)
if err != nil {
dao.TencentVideo.UpdateStatus(ctx, video.Id, entity.VerifyStatusPending)
return nil, err
}
return log, nil
}
// submitVideoCheck 提交视频校验
func (s *MaterialVerifyService) submitVideoCheck(ctx context.Context, video *entity.TencentVideo, log *entity.MaterialVerifyLog) error {
startTime := time.Now()
callbackMode := g.Cfg().MustGet(ctx, "check.callback_mode").Bool()
var callbackURL string
if callbackMode {
callbackURL = g.Cfg().MustGet(ctx, "check.video.callback_url").String()
}
requestParams := map[string]interface{}{
"videoURL": video.PreviewURL,
"dataID": video.VideoID,
"callbackURL": callbackURL,
}
requestParamsJSON, _ := json.Marshal(requestParams)
result, err := VideoDetection.DetectVideo(ctx, video.PreviewURL, video.VideoID, callbackURL)
duration := time.Since(startTime).Milliseconds()
if err != nil {
dao.MaterialVerifyLog.UpdateError(ctx, log.Id, entity.VerifyStatusPending, err.Error())
dao.MaterialVerifyLog.UpdateDuration(ctx, log.Id, duration)
g.Log().Warningf(ctx, "视频校验接口调用失败(保持待检验), id=%d, videoId=%s, error=%v", video.Id, video.VideoID, err)
return fmt.Errorf("视频检测提交失败, videoId=%s: %w", video.VideoID, err)
}
dao.MaterialVerifyLog.UpdateTaskID(ctx, log.Id, result.TaskID)
dao.MaterialVerifyLog.UpdateRequestParams(ctx, log.Id, string(requestParamsJSON))
if !callbackMode {
g.Log().Infof(ctx, "轮询模式:视频检测已提交, taskId=%s, 请通过轮询接口获取结果", result.TaskID)
}
TencentContentCheck.writeAuditLog(ctx, consts.SourceTableTencentVideo, video.Id, video.VideoID, video.PreviewURL, result.TaskID, -1, 0, 0, "", duration)
g.Log().Infof(ctx, "视频校验已提交, id=%d, videoId=%s, taskId=%s, duration=%dms",
video.Id, video.VideoID, result.TaskID, duration)
return nil
}
// =============================================================================
// 回调处理
// =============================================================================
// ProcessImageCallback 处理图片校验回调
func (s *MaterialVerifyService) ProcessImageCallback(ctx context.Context, callbackData string) error {
g.Log().Infof(ctx, "处理图片校验回调, data: %s", callbackData)
var callback ImageCallbackData
if err := json.Unmarshal([]byte(callbackData), &callback); err != nil {
g.Log().Errorf(ctx, "解析图片回调数据失败: %v", err)
return fmt.Errorf("解析图片回调数据失败: %w", err)
}
if callback.Antispam == nil {
return fmt.Errorf("回调数据格式错误:缺少antispam字段")
}
antispam := callback.Antispam
g.Log().Infof(ctx, "处理图片校验结果 - taskId: %s, suggestion: %d, resultType: %d",
antispam.TaskId, antispam.Suggestion, antispam.ResultType)
log, err := dao.MaterialVerifyLog.GetByTaskID(ctx, antispam.TaskId)
if err != nil {
return fmt.Errorf("查询图片校验日志失败, taskId=%s: %w", antispam.TaskId, err)
}
if log == nil {
g.Log().Warningf(ctx, "未找到校验日志, taskId=%s", antispam.TaskId)
return nil
}
verifyStatus := SuggestionToVerifyStatus(antispam.Suggestion)
err = dao.MaterialVerifyLog.UpdateVerifyResult(ctx, log.Id, verifyStatus,
antispam.Suggestion, antispam.Label, antispam.ResultType, callbackData, antispam.CensorTime)
if err != nil {
return fmt.Errorf("更新图片校验日志结果失败, logId=%d: %w", log.Id, err)
}
if log.SourceTable == consts.SourceTableTencentImage {
s.updateImageStatus(ctx, log.SourceID, verifyStatus)
}
// 更新送检审计日志
s.updateCheckLogResult(ctx, antispam.TaskId, antispam.Suggestion, antispam.Label, antispam.ResultType, callbackData)
// 提取风险描述
if antispam.RiskDescription != "" {
_ = dao.MaterialVerifyLog.UpdateRiskDescription(ctx, log.Id, antispam.RiskDescription)
}
g.Log().Infof(ctx, "图片校验回调处理完成, taskId=%s, verifyStatus=%s, suggestion=%d",
antispam.TaskId, verifyStatus, antispam.Suggestion)
return nil
}
// ProcessVideoCallback 处理视频校验回调
func (s *MaterialVerifyService) ProcessVideoCallback(ctx context.Context, callbackData string) error {
g.Log().Infof(ctx, "处理视频校验回调, data: %s", callbackData)
var callback VideoCallbackData
if err := json.Unmarshal([]byte(callbackData), &callback); err != nil {
g.Log().Errorf(ctx, "解析视频回调数据失败: %v", err)
return fmt.Errorf("解析视频回调数据失败: %w", err)
}
if callback.Antispam == nil {
return fmt.Errorf("视频回调数据格式错误:缺少antispam字段")
}
antispam := callback.Antispam
g.Log().Infof(ctx, "处理视频校验结果 - taskId: %s, suggestion: %d, resultType: %d",
antispam.TaskID, antispam.Suggestion, antispam.ResultType)
log, err := dao.MaterialVerifyLog.GetByTaskID(ctx, antispam.TaskID)
if err != nil {
return fmt.Errorf("查询视频校验日志失败, taskId=%s: %w", antispam.TaskID, err)
}
if log == nil {
g.Log().Warningf(ctx, "未找到校验日志, taskId=%s", antispam.TaskID)
return nil
}
verifyStatus := SuggestionToVerifyStatus(antispam.Suggestion)
checkTime := antispam.CensorTime
if checkTime == 0 {
checkTime = antispam.CheckTime
}
err = dao.MaterialVerifyLog.UpdateVerifyResult(ctx, log.Id, verifyStatus,
antispam.Suggestion, antispam.Label, antispam.ResultType, callbackData, checkTime)
if err != nil {
return fmt.Errorf("更新视频校验日志结果失败, logId=%d: %w", log.Id, err)
}
if log.SourceTable == consts.SourceTableTencentVideo {
s.updateVideoStatus(ctx, log.SourceID, verifyStatus)
}
// 更新送检审计日志
s.updateCheckLogResult(ctx, antispam.TaskID, antispam.Suggestion, antispam.Label, antispam.ResultType, callbackData)
// 提取风险描述
if antispam.RiskDescription != "" {
_ = dao.MaterialVerifyLog.UpdateRiskDescription(ctx, log.Id, antispam.RiskDescription)
}
g.Log().Infof(ctx, "视频校验回调处理完成, taskId=%s, verifyStatus=%s, suggestion=%d",
antispam.TaskID, verifyStatus, antispam.Suggestion)
return nil
}
// updateCheckLogResult 更新 tencent_content_check_log 的检测结果
func (s *MaterialVerifyService) updateCheckLogResult(ctx context.Context, taskID string, suggestion, label, resultType int, responseData string) {
checkLog, err := dao.TencentContentCheckLog.GetByTaskID(ctx, taskID)
if err != nil {
g.Log().Warningf(ctx, "查询送检审计日志失败, taskId=%s: %v", taskID, err)
return
}
if checkLog == nil {
g.Log().Debugf(ctx, "送检审计日志不存在, taskId=%s(可能是通过 API 直接提交的)", taskID)
return
}
_ = dao.TencentContentCheckLog.UpdateCheckResult(ctx, checkLog.Id, suggestion, label, resultType, time.Now().UnixMilli())
}
// =============================================================================
// 轮询模式处理
// =============================================================================
// ErrResultPending 表示检测结果尚未就绪,非错误状态
var ErrResultPending = fmt.Errorf("检测结果尚未就绪")
// 易盾检测状态常量
const (
YidunStatusNotStart = 0 // 未开始
YidunStatusProcessing = 1 // 检测中
YidunStatusSuccess = 2 // 检测成功
YidunStatusFailed = 3 // 检测失败
)
// ProcessImageResultByTask 根据任务ID处理图片结果(轮询模式)
// 返回 nil 表示结果已处理完成,返回 ErrResultPending 表示仍未就绪
func (s *MaterialVerifyService) ProcessImageResultByTask(ctx context.Context, taskID string) error {
log, err := dao.MaterialVerifyLog.GetByTaskID(ctx, taskID)
if err != nil {
return fmt.Errorf("查询校验日志失败, taskId=%s: %w", taskID, err)
}
if log == nil {
return fmt.Errorf("未找到校验日志, taskId=%s", taskID)
}
result, err := ImageDetection.GetImageResult(ctx, taskID)
if err != nil {
if err == ErrImageResultNotFound || err == ErrImageStillProcessing {
g.Log().Infof(ctx, "图片检测结果未就绪, taskId=%s, 保持pending状态, err=%v", taskID, err)
return ErrResultPending
}
dao.MaterialVerifyLog.UpdateError(ctx, log.Id, entity.VerifyStatusPending, err.Error())
g.Log().Warningf(ctx, "图片检测查询失败(保持待检验), taskId=%s, error=%v", taskID, err)
return ErrResultPending
}
if result.Status == YidunStatusProcessing || result.Status == YidunStatusNotStart {
g.Log().Infof(ctx, "图片检测仍在进行中, taskId=%s, status=%d, 保持pending状态", taskID, result.Status)
return ErrResultPending
}
if result.Status == YidunStatusFailed {
dao.MaterialVerifyLog.UpdateError(ctx, log.Id, entity.VerifyStatusPending,
fmt.Sprintf("易盾检测失败, status=%d", result.Status))
g.Log().Warningf(ctx, "图片检测失败(保持待检验), taskId=%s, status=%d", taskID, result.Status)
return ErrResultPending
}
verifyStatus := SuggestionToVerifyStatus(result.Suggestion)
responseJSON, _ := json.Marshal(result)
dao.MaterialVerifyLog.UpdateVerifyResult(ctx, log.Id, verifyStatus,
result.Suggestion, result.Label, result.ResultType, string(responseJSON), result.CensorTime)
if log.SourceTable == consts.SourceTableTencentImage {
s.updateImageStatus(ctx, log.SourceID, verifyStatus)
}
// 提取风险描述
if result.Antispam != nil && result.Antispam.RiskDescription != nil {
_ = dao.MaterialVerifyLog.UpdateRiskDescription(ctx, log.Id, *result.Antispam.RiskDescription)
}
// 更新送检审计日志
s.updateCheckLogResult(ctx, taskID, result.Suggestion, result.Label, result.ResultType, string(responseJSON))
g.Log().Infof(ctx, "图片检测结果更新成功, taskId=%s, status=%d, suggestion=%d, verifyStatus=%s",
taskID, result.Status, result.Suggestion, verifyStatus)
return nil
}
// ProcessVideoResultByTask 根据任务ID处理视频结果(轮询模式)
func (s *MaterialVerifyService) ProcessVideoResultByTask(ctx context.Context, taskID string) error {
log, err := dao.MaterialVerifyLog.GetByTaskID(ctx, taskID)
if err != nil {
return fmt.Errorf("查询校验日志失败, taskId=%s: %w", taskID, err)
}
if log == nil {
return fmt.Errorf("未找到校验日志, taskId=%s", taskID)
}
result, err := VideoDetection.GetVideoResult(ctx, taskID)
if err != nil {
if err == ErrVideoResultNotFound || err == ErrVideoStillProcessing {
g.Log().Infof(ctx, "视频检测结果未就绪, taskId=%s, 保持pending状态, err=%v", taskID, err)
return ErrResultPending
}
dao.MaterialVerifyLog.UpdateError(ctx, log.Id, entity.VerifyStatusPending, err.Error())
g.Log().Warningf(ctx, "视频检测查询失败(保持待检验), taskId=%s, error=%v", taskID, err)
return ErrResultPending
}
if result.Status == YidunStatusProcessing || result.Status == YidunStatusNotStart {
g.Log().Infof(ctx, "视频检测仍在进行中, taskId=%s, status=%d, 保持pending状态", taskID, result.Status)
return ErrResultPending
}
if result.Status == YidunStatusFailed {
dao.MaterialVerifyLog.UpdateError(ctx, log.Id, entity.VerifyStatusPending,
fmt.Sprintf("易盾检测失败, status=%d", result.Status))
g.Log().Warningf(ctx, "视频检测失败(保持待检验), taskId=%s, status=%d", taskID, result.Status)
return ErrResultPending
}
verifyStatus := SuggestionToVerifyStatus(result.Suggestion)
responseJSON, _ := json.Marshal(result)
dao.MaterialVerifyLog.UpdateVerifyResult(ctx, log.Id, verifyStatus,
result.Suggestion, result.Label, result.ResultType, string(responseJSON), result.CensorTime)
if log.SourceTable == consts.SourceTableTencentVideo {
s.updateVideoStatus(ctx, log.SourceID, verifyStatus)
}
// 提取风险描述
if result.Antispam != nil && result.Antispam.RiskDescription != nil {
_ = dao.MaterialVerifyLog.UpdateRiskDescription(ctx, log.Id, *result.Antispam.RiskDescription)
}
// 更新送检审计日志
s.updateCheckLogResult(ctx, taskID, result.Suggestion, result.Label, result.ResultType, string(responseJSON))
g.Log().Infof(ctx, "视频检测结果更新成功, taskId=%s, status=%d, suggestion=%d, verifyStatus=%s",
taskID, result.Status, result.Suggestion, verifyStatus)
return nil
}
// =============================================================================
// 辅助方法
// =============================================================================
// createVerifyLog 创建校验日志
func (s *MaterialVerifyService) createVerifyLog(ctx context.Context, materialType, materialID, sourceTable string, sourceID, accountID int64) *entity.MaterialVerifyLog {
var tenantID int64
if user := ctx.Value("user"); user != nil {
if u, ok := user.(*beans.User); ok {
tenantID = int64(u.TenantId)
}
}
log := &entity.MaterialVerifyLog{
TenantID: tenantID,
MaterialType: materialType,
MaterialID: materialID,
SourceTable: sourceTable,
SourceID: sourceID,
AccountID: accountID,
VerifyStatus: entity.VerifyStatusPending,
}
id, err := dao.MaterialVerifyLog.Create(ctx, log)
if err != nil {
g.Log().Errorf(ctx, "创建校验日志失败: %v", err)
return nil
}
log.Id = id
return log
}
// updateImageStatus 更新图片状态
func (s *MaterialVerifyService) updateImageStatus(ctx context.Context, imageID int64, verifyStatus string) {
_, err := dao.TencentImage.UpdateStatus(ctx, imageID, verifyStatus)
if err != nil {
g.Log().Errorf(ctx, "更新图片状态失败: %v", err)
} else {
g.Log().Infof(ctx, "更新图片状态成功, imageID=%d, status=%s", imageID, verifyStatus)
}
}
// updateVideoStatus 更新视频状态
func (s *MaterialVerifyService) updateVideoStatus(ctx context.Context, videoID int64, verifyStatus string) {
_, err := dao.TencentVideo.UpdateStatus(ctx, videoID, verifyStatus)
if err != nil {
g.Log().Errorf(ctx, "更新视频状态失败: %v", err)
} else {
g.Log().Infof(ctx, "更新视频状态成功, videoID=%d, status=%s", videoID, verifyStatus)
}
}
// =============================================================================
// 查询接口
// =============================================================================
// GetLogByID 根据ID获取日志
func (s *MaterialVerifyService) GetLogByID(ctx context.Context, id int64) (*entity.MaterialVerifyLog, error) {
return dao.MaterialVerifyLog.GetByID(ctx, id)
}
// GetLogsByMaterialID 根据素材ID获取日志列表
func (s *MaterialVerifyService) GetLogsByMaterialID(ctx context.Context, materialID string) ([]entity.MaterialVerifyLog, error) {
return dao.MaterialVerifyLog.GetByMaterialID(ctx, materialID)
}
// GetLogsByCondition 条件查询日志
func (s *MaterialVerifyService) GetLogsByCondition(ctx context.Context, condition map[string]interface{}, page, pageSize int) ([]entity.MaterialVerifyLog, int, error) {
return dao.MaterialVerifyLog.GetByCondition(ctx, condition, page, pageSize)
}
// GetStats 获取统计信息
func (s *MaterialVerifyService) GetStats(ctx context.Context) (map[string]int, error) {
return dao.MaterialVerifyLog.GetStats(ctx)
}
// =============================================================================
// 轮询模式 - 批量查询检测结果
// =============================================================================
// PollPendingResults 轮询所有待查询结果的日志
func (s *MaterialVerifyService) PollPendingResults(ctx context.Context) (int, int, error) {
logs, err := dao.MaterialVerifyLog.GetPendingResults(ctx, PollBatchSize)
if err != nil {
return 0, 0, err
}
if len(logs) == 0 {
g.Log().Infof(ctx, "没有待查询结果的日志")
return 0, 0, nil
}
g.Log().Infof(ctx, "开始轮询 %d 条待处理结果", len(logs))
successCount := 0
failCount := 0
var lastErr error
for _, log := range logs {
var err error
if log.SourceTable == consts.SourceTableTencentImage {
err = s.ProcessImageResultByTask(ctx, log.TaskID)
} else if log.SourceTable == consts.SourceTableTencentVideo {
err = s.ProcessVideoResultByTask(ctx, log.TaskID)
} else {
g.Log().Warningf(ctx, "未知的来源表: %s, logId=%d", log.SourceTable, log.Id)
continue
}
if err == ErrResultPending {
g.Log().Infof(ctx, "结果未就绪, logId=%d, taskId=%s", log.Id, log.TaskID)
} else if err != nil {
failCount++
lastErr = err
g.Log().Warningf(ctx, "处理结果失败, logId=%d, taskId=%s, error=%v", log.Id, log.TaskID, err)
} else {
successCount++
g.Log().Infof(ctx, "处理结果成功, logId=%d, taskId=%s", log.Id, log.TaskID)
}
time.Sleep(100 * time.Millisecond)
}
g.Log().Infof(ctx, "轮询完成, 成功=%d, 失败=%d, 未就绪=%d", successCount, failCount, len(logs)-successCount-failCount)
return successCount, failCount, lastErr
}
// PollPendingResultsByType 按类型轮询待查询结果的日志
func (s *MaterialVerifyService) PollPendingResultsByType(ctx context.Context, sourceTable string) (int, int, error) {
logs, err := dao.MaterialVerifyLog.GetPendingResults(ctx, PollBatchSize)
if err != nil {
return 0, 0, err
}
var filteredLogs []entity.MaterialVerifyLog
for _, log := range logs {
if log.SourceTable == sourceTable {
filteredLogs = append(filteredLogs, log)
}
}
if len(filteredLogs) == 0 {
g.Log().Infof(ctx, "没有待查询结果的日志, sourceTable=%s", sourceTable)
return 0, 0, nil
}
successCount := 0
failCount := 0
var lastErr error
for _, log := range filteredLogs {
var err error
if sourceTable == consts.SourceTableTencentImage {
err = s.ProcessImageResultByTask(ctx, log.TaskID)
} else if sourceTable == consts.SourceTableTencentVideo {
err = s.ProcessVideoResultByTask(ctx, log.TaskID)
}
if err != nil {
failCount++
lastErr = err
} else {
successCount++
}
time.Sleep(100 * time.Millisecond)
}
return successCount, failCount, lastErr
}
// PollPendingImageResults 轮询图片待查询结果
func (s *MaterialVerifyService) PollPendingImageResults(ctx context.Context) (int, int, error) {
return s.PollPendingResultsByType(ctx, consts.SourceTableTencentImage)
}
// PollPendingVideoResults 轮询视频待查询结果
func (s *MaterialVerifyService) PollPendingVideoResults(ctx context.Context) (int, int, error) {
return s.PollPendingResultsByType(ctx, consts.SourceTableTencentVideo)
}
// =============================================================================
// 导出服务 - 不通过数据导出
// =============================================================================
// ExportRejectedItem 导出的不通过数据项
type ExportRejectedItem struct {
ID int64 `json:"id"`
MaterialID string `json:"materialId"`
AccountID int64 `json:"accountId"`
CorporationName string `json:"corporationName"`
PreviewURL string `json:"previewUrl"`
Description string `json:"description"`
ErrorMsg string `json:"errorMsg"`
MaterialType string `json:"materialType"`
ImageUsage string `json:"imageUsage,omitempty"`
CreatedAt string `json:"createdAt"`
}
// getFailureReason 获取失败原因
func getFailureReason(log *entity.MaterialVerifyLog) string {
if log == nil {
return "无校验日志"
}
if log.ErrorMsg != "" {
return log.ErrorMsg
}
reasonMap := map[int]string{
0: "内容检测通过",
1: "内容嫌疑(需人工审核)",
2: "内容不通过",
}
suggestionText := reasonMap[log.Suggestion]
if suggestionText == "" {
suggestionText = fmt.Sprintf("未知(suggestion=%d)", log.Suggestion)
}
if log.ResponseResult != "" {
var resultMap map[string]interface{}
if err := json.Unmarshal([]byte(log.ResponseResult), &resultMap); err == nil {
if labels, ok := resultMap["labels"]; ok {
return fmt.Sprintf("%s (labels: %v)", suggestionText, labels)
}
}
return suggestionText
}
return suggestionText
}
const exportBatchSize = 1000
// ExportRejectedData 导出不通过数据
func (s *MaterialVerifyService) ExportRejectedData(ctx context.Context, materialType string) ([]ExportRejectedItem, error) {
var items []ExportRejectedItem
accountMap := make(map[int64]string)
if accounts, err := dao.TencentAccountRelation.GetAll(ctx); err == nil {
for _, acc := range accounts {
if acc.CorporationName != "" {
accountMap[acc.AccountID] = acc.CorporationName
}
}
}
if materialType == "" || materialType == entity.MaterialTypeImage {
condition := map[string]interface{}{
entity.TencentImageCols.VerifyStatus: entity.VerifyStatusRejected,
}
page := 1
for {
images, total, err := dao.TencentImage.GetByCondition(ctx, condition, page, exportBatchSize)
if err != nil {
g.Log().Errorf(ctx, "查询不通过图片失败: %v", err)
return nil, fmt.Errorf("查询不通过图片失败: %w", err)
}
for _, img := range images {
log, _ := dao.MaterialVerifyLog.GetLastRejectedLogByMaterialID(ctx, img.ImageID, entity.VerifyStatusRejected)
var createdAtStr string
if log != nil && log.CreatedAt != nil {
createdAtStr = log.CreatedAt.Format("2006-01-02 15:04:05")
}
items = append(items, ExportRejectedItem{
ID: img.Id, MaterialID: img.ImageID, AccountID: img.AccountID,
CorporationName: accountMap[img.AccountID], PreviewURL: img.PreviewURL,
Description: img.Description, ErrorMsg: getFailureReason(log),
MaterialType: entity.MaterialTypeImage, ImageUsage: img.ImageUsage,
CreatedAt: createdAtStr,
})
}
if page*exportBatchSize >= total {
break
}
page++
}
}
if materialType == "" || materialType == entity.MaterialTypeVideo {
condition := map[string]interface{}{
entity.TencentVideoCols.VerifyStatus: entity.VerifyStatusRejected,
}
page := 1
for {
videos, total, err := dao.TencentVideo.GetByCondition(ctx, condition, page, exportBatchSize)
if err != nil {
g.Log().Errorf(ctx, "查询不通过视频失败: %v", err)
return nil, fmt.Errorf("查询不通过视频失败: %w", err)
}
for _, vid := range videos {
log, _ := dao.MaterialVerifyLog.GetLastRejectedLogByMaterialID(ctx, vid.VideoID, entity.VerifyStatusRejected)
var createdAtStr string
if log != nil && log.CreatedAt != nil {
createdAtStr = log.CreatedAt.Format("2006-01-02 15:04:05")
}
items = append(items, ExportRejectedItem{
ID: vid.Id, MaterialID: vid.VideoID, AccountID: vid.AccountID,
CorporationName: accountMap[vid.AccountID], PreviewURL: vid.PreviewURL,
Description: vid.Description, ErrorMsg: getFailureReason(log),
MaterialType: entity.MaterialTypeVideo, CreatedAt: createdAtStr,
})
}
if page*exportBatchSize >= total {
break
}
page++
}
}
return items, nil
}
// GetPendingResultsCount 获取待查询结果的数量
func (s *MaterialVerifyService) GetPendingResultsCount(ctx context.Context) (int, error) {
return dao.MaterialVerifyLog.CountPendingResults(ctx)
}
// GetPendingResultsDetail 获取待查询结果的明细列表
type PendingResultItem struct {
LogID int64 `json:"logId"`
MaterialID string `json:"materialId"`
MaterialType string `json:"materialType"`
SourceTable string `json:"sourceTable"`
TaskID string `json:"taskId"`
CreatedAt string `json:"createdAt"`
}
func (s *MaterialVerifyService) GetPendingResultsDetail(ctx context.Context, limit int) ([]PendingResultItem, error) {
logs, err := dao.MaterialVerifyLog.GetPendingResults(ctx, limit)
if err != nil {
return nil, err
}
var items []PendingResultItem
for _, l := range logs {
createdAt := ""
if l.CreatedAt != nil {
createdAt = l.CreatedAt.Format("2006-01-02 15:04:05")
}
items = append(items, PendingResultItem{
LogID: l.Id,
MaterialID: l.MaterialID,
MaterialType: l.MaterialType,
SourceTable: l.SourceTable,
TaskID: l.TaskID,
CreatedAt: createdAt,
})
}
return items, nil
}