726 lines
24 KiB
Go
726 lines
24 KiB
Go
package service
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"net/url"
|
||
"os"
|
||
"path/filepath"
|
||
"time"
|
||
|
||
"github.com/gogf/gf/v2/database/gdb"
|
||
"github.com/gogf/gf/v2/errors/gerror"
|
||
"github.com/gogf/gf/v2/frame/g"
|
||
"github.com/gogf/gf/v2/os/gtime"
|
||
|
||
"observer-server/biz/consts"
|
||
"observer-server/biz/dao"
|
||
"observer-server/biz/model/dto"
|
||
"observer-server/biz/model/entity"
|
||
"observer-server/common"
|
||
)
|
||
|
||
// annotateService 标注众包与时长激励(技术设计.md「App 标注众包与时长激励」):
|
||
// 管理端下发数据集级动态池任务 → App 用户领取(pending 锁)/提交(进待审核)→
|
||
// 每累计 N 张即时发时长(日上限封顶)→ 管理端审核,通过比例过低自动冻结领取资格。
|
||
type annotateService struct{}
|
||
|
||
var Annotate = &annotateService{}
|
||
|
||
// annotateRewardCfg annotateReward 配置(缺失/非法回退 consts 默认)
|
||
type annotateRewardCfg struct {
|
||
ClaimSize int
|
||
ClaimTimeoutHour int
|
||
RewardPerImages int
|
||
RewardMinutes int
|
||
DailyCapMinutes int
|
||
FreezeRatio float64
|
||
FreezeMinReviewed int
|
||
FreezeHours int
|
||
}
|
||
|
||
func annotateRewardConfig(ctx context.Context) annotateRewardCfg {
|
||
get := func(key string, def int) int {
|
||
v := g.Cfg().MustGet(ctx, "annotateReward."+key, def).Int()
|
||
if v <= 0 {
|
||
return def
|
||
}
|
||
return v
|
||
}
|
||
ratio := g.Cfg().MustGet(ctx, "annotateReward.freezeRatio", consts.AnnotateFreezeRatio).Float64()
|
||
if ratio <= 0 || ratio > 1 {
|
||
ratio = consts.AnnotateFreezeRatio
|
||
}
|
||
return annotateRewardCfg{
|
||
ClaimSize: get("claimSize", consts.AnnotateClaimSize),
|
||
ClaimTimeoutHour: get("claimTimeoutHour", consts.AnnotateClaimTimeoutHour),
|
||
RewardPerImages: get("rewardPerImages", consts.AnnotateRewardPerImages),
|
||
RewardMinutes: get("rewardMinutes", consts.AnnotateRewardMinutes),
|
||
DailyCapMinutes: get("dailyCapMinutes", consts.AnnotateDailyCapMinutes),
|
||
FreezeRatio: ratio,
|
||
FreezeMinReviewed: get("freezeMinReviewed", consts.AnnotateFreezeMinReviewed),
|
||
FreezeHours: get("freezeHours", consts.AnnotateFreezeHours),
|
||
}
|
||
}
|
||
|
||
// speciesOf 任务物种展示名:gen_species 空回退数据集名(单物种规则)
|
||
func annotateSpecies(d *entity.Dataset) string {
|
||
if d.GenSpecies != "" {
|
||
return d.GenSpecies
|
||
}
|
||
return d.Name
|
||
}
|
||
|
||
// annotateImageUrl App 端图片访问地址(登录态 Bearer 鉴权,controller 直写响应体,
|
||
// 与管理端 /admin/datasets/image 的 X-Admin-Token 鉴权隔离——App 用户无管理 token)
|
||
func annotateImageUrl(ctx context.Context, datasetId int64, filename string) string {
|
||
return fmt.Sprintf("/api/v1/annotate/image?%s",
|
||
url.Values{"datasetId": {fmt.Sprintf("%d", datasetId)}, "filename": {filename}}.Encode())
|
||
}
|
||
|
||
func annotateFrozen(lic *entity.License) *gtime.Time {
|
||
if lic == nil || lic.AnnotateFrozenUntil == nil {
|
||
return nil
|
||
}
|
||
if lic.AnnotateFrozenUntil.Timestamp() <= gtime.Now().Timestamp() {
|
||
return nil // 已过期 = 自动解冻
|
||
}
|
||
return lic.AnnotateFrozenUntil
|
||
}
|
||
|
||
// ---------- 管理端 ----------
|
||
|
||
// AdminCreateTask 下发任务(2026-09-07 图片粒度):勾选的具体图片集合即任务图集——
|
||
// dataset_image.annotate_task_id 写任务 id 作占用标记(已下发图即从管理端「未标注」tab 消失),
|
||
// 停用任务时释放未领取图回未标注。同数据集同时只允许一个 published 任务;
|
||
// Serial 内查重 + 逐图校验(属本数据集/未标注/未占用,任一非法整单拒绝)+ 事务内插任务与批量占用。
|
||
func (s *annotateService) AdminCreateTask(ctx context.Context, req *dto.AdminAnnotateTaskCreateReq) (*dto.AdminAnnotateTaskCreateRes, error) {
|
||
dataset, err := dao.Dataset.GetById(ctx, req.DatasetId)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if dataset == nil {
|
||
return nil, gerror.NewCode(common.CodeDatasetNotFound)
|
||
}
|
||
if dataset.Source == consts.DatasetSourceNegative {
|
||
return nil, gerror.New("负样本库不参与标注众包")
|
||
}
|
||
images, err := dao.DatasetImage.GetByIds(ctx, req.ImageIds)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if len(images) != len(req.ImageIds) {
|
||
return nil, gerror.New("部分图片不存在,请刷新后重试")
|
||
}
|
||
for _, img := range images {
|
||
if img.DatasetId != dataset.Id {
|
||
return nil, gerror.Newf("图片 %d 不属于该数据集", img.Id)
|
||
}
|
||
if img.ReviewStatus != consts.ReviewImageNone {
|
||
return nil, gerror.Newf("图片 %d 已标注,只能下发未标注图", img.Id)
|
||
}
|
||
if img.AnnotateTaskId != 0 {
|
||
return nil, gerror.Newf("图片 %d 已下发给其他任务,请刷新后重试", img.Id)
|
||
}
|
||
}
|
||
var taskId int64
|
||
err = common.Serial().Submit(ctx, func() error {
|
||
running, err := dao.AnnotateTask.GetPublishedByDataset(ctx, dataset.Id)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if running != nil {
|
||
return gerror.Newf("该数据集已有进行中的标注任务(id=%d),请先停用", running.Id)
|
||
}
|
||
return g.DB().Transaction(ctx, func(ctx context.Context, tx gdb.TX) error {
|
||
taskId, err = dao.AnnotateTask.InsertInTx(ctx, tx, &entity.AnnotateTask{
|
||
DatasetId: dataset.Id,
|
||
Name: req.Name,
|
||
Status: consts.AnnotateTaskPublished,
|
||
CreatedAt: gtime.Now(),
|
||
})
|
||
if err != nil {
|
||
return err
|
||
}
|
||
return dao.DatasetImage.OccupyByTask(ctx, req.ImageIds, taskId)
|
||
})
|
||
})
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &dto.AdminAnnotateTaskCreateRes{Id: taskId}, nil
|
||
}
|
||
|
||
// AdminListTasks 任务分页:组装数据集名/池余量/各状态记录数
|
||
func (s *annotateService) AdminListTasks(ctx context.Context, req *dto.AdminAnnotateTaskListReq) (*dto.AdminAnnotateTaskListRes, error) {
|
||
page, size := common.NormalizePage(req.Page, req.Size)
|
||
list, total, err := dao.AnnotateTask.Page(ctx, req.DatasetId, page, size)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
res := &dto.AdminAnnotateTaskListRes{Total: total, List: make([]*dto.AdminAnnotateTaskItem, 0, len(list))}
|
||
if len(list) == 0 {
|
||
return res, nil
|
||
}
|
||
taskIds := make([]int64, 0, len(list))
|
||
for _, t := range list {
|
||
taskIds = append(taskIds, t.Id)
|
||
}
|
||
names := Training.datasetNameMap(ctx)
|
||
poolCnt, err := dao.DatasetImage.CountPoolByTaskIds(ctx, taskIds)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
recCnt, err := dao.AnnotateRecord.CountByTaskIds(ctx, taskIds)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
for _, t := range list {
|
||
item := &dto.AdminAnnotateTaskItem{
|
||
Id: t.Id,
|
||
DatasetId: t.DatasetId,
|
||
DatasetName: names[t.DatasetId],
|
||
Name: t.Name,
|
||
Status: t.Status,
|
||
PoolRemain: poolCnt[t.Id],
|
||
CreatedAt: t.CreatedAt,
|
||
}
|
||
if m := recCnt[t.Id]; m != nil {
|
||
item.Pending = m[consts.AnnotateRecordPending]
|
||
item.Submitted = m[consts.AnnotateRecordSubmitted]
|
||
item.Approved = m[consts.AnnotateRecordApproved]
|
||
item.Rejected = m[consts.AnnotateRecordRejected]
|
||
}
|
||
res.List = append(res.List, item)
|
||
}
|
||
return res, nil
|
||
}
|
||
|
||
// AdminStopTask 停用任务(整体不可再领取;已领取未提交的可继续提交)。
|
||
// Serial 内停用 + 释放:未被领取过的图清占用标记回未标注(重新出现在未标注 tab 可再次下发);
|
||
// 有 pending 锁的图保留占用(用户可继续提交)。过期锁先按领取同语义惰性清理,防已放弃的锁
|
||
// 永久占图(停用任务不再有领取触发清理)。
|
||
func (s *annotateService) AdminStopTask(ctx context.Context, req *dto.AdminAnnotateTaskStopReq) (*dto.AdminAnnotateTaskStopRes, error) {
|
||
t, err := dao.AnnotateTask.GetById(ctx, req.Id)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if t == nil {
|
||
return nil, gerror.New("标注任务不存在")
|
||
}
|
||
cfg := annotateRewardConfig(ctx)
|
||
err = common.Serial().Submit(ctx, func() error {
|
||
if err := dao.AnnotateTask.Stop(ctx, req.Id); err != nil {
|
||
return err
|
||
}
|
||
if err := dao.AnnotateRecord.DeleteExpiredPending(ctx, gtime.Now().Add(-time.Duration(cfg.ClaimTimeoutHour)*time.Hour)); err != nil {
|
||
return err
|
||
}
|
||
pool, err := dao.DatasetImage.ListPoolByTask(ctx, req.Id)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
lockedIds, err := dao.AnnotateRecord.ListPendingImageIdsByTask(ctx, req.Id)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
locked := make(map[int64]bool, len(lockedIds))
|
||
for _, id := range lockedIds {
|
||
locked[id] = true
|
||
}
|
||
release := make([]int64, 0, len(pool))
|
||
for _, img := range pool {
|
||
if !locked[img.Id] {
|
||
release = append(release, img.Id)
|
||
}
|
||
}
|
||
return dao.DatasetImage.ReleaseByTask(ctx, release, req.Id)
|
||
})
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &dto.AdminAnnotateTaskStopRes{}, nil
|
||
}
|
||
|
||
// AdminReview 批量审核:通过 → review_status=2;拒绝 → 清标注回未标注池。
|
||
// 多表变更(dataset_image + annotate_record + license 冻结)走 Serial + 事务;
|
||
// 审核后对涉及用户重算通过比例,低于阈值自动冻结(解冻后按 stats_since 重新累计)。
|
||
func (s *annotateService) AdminReview(ctx context.Context, req *dto.AdminAnnotateReviewReq) (*dto.AdminAnnotateReviewRes, error) {
|
||
res := &dto.AdminAnnotateReviewRes{FrozenPhones: []string{}}
|
||
// 只审待审核图(重复审核/未标注图直接忽略)
|
||
images, err := dao.DatasetImage.GetByIds(ctx, req.ImageIds)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
ids := make([]int64, 0, len(images))
|
||
for _, img := range images {
|
||
if img.ReviewStatus == consts.ReviewImagePending {
|
||
ids = append(ids, img.Id)
|
||
}
|
||
}
|
||
if len(ids) == 0 {
|
||
return res, nil
|
||
}
|
||
now := gtime.Now()
|
||
recordStatus := consts.AnnotateRecordRejected
|
||
if req.Approve {
|
||
recordStatus = consts.AnnotateRecordApproved
|
||
}
|
||
var phones []string
|
||
err = common.Serial().Submit(ctx, func() error {
|
||
return g.DB().Transaction(ctx, func(ctx context.Context, tx gdb.TX) error {
|
||
if req.Approve {
|
||
// 通过 = 已标注(全部框原样保留——通过即人工背书,疑似框是合法的 class 1 训练标注)
|
||
if err := dao.DatasetImage.SetReviewStatus(ctx, ids, consts.ReviewImageApproved); err != nil {
|
||
return err
|
||
}
|
||
} else {
|
||
// 拒绝:清标注(不变式 0 ⟹ 无框)回未标注池,可被再次领取
|
||
for _, id := range ids {
|
||
if err := dao.DatasetImage.UpdateLabelsAndReview(ctx, id, "", consts.ReviewImageNone); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
}
|
||
if err := dao.AnnotateRecord.MarkReviewedByImages(ctx, ids, recordStatus, now); err != nil {
|
||
return err
|
||
}
|
||
phones, err = dao.AnnotateRecord.ListSubmittedPhonesByImages(ctx, ids)
|
||
return err
|
||
})
|
||
})
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
res.Affected = int64(len(ids))
|
||
// 审核后逐用户重算通过比例(IO 次数 = 涉及用户数,审核批量 ≤ 每页图片数,可控)
|
||
for _, phone := range phones {
|
||
frozen, err := s.maybeFreeze(ctx, phone)
|
||
if err != nil {
|
||
g.Log().Errorf(ctx, "标注低质冻结判定失败(%s): %+v", phone, err)
|
||
continue
|
||
}
|
||
if frozen {
|
||
res.FrozenPhones = append(res.FrozenPhones, phone)
|
||
}
|
||
}
|
||
return res, nil
|
||
}
|
||
|
||
// maybeFreeze 通过比例过低自动冻结:已审核样本(stats_since 基线后)≥ FreezeMinReviewed 且
|
||
// approved/(approved+rejected) < FreezeRatio → frozen_until = now + FreezeHours、基线重置(解冻后重新累计)。
|
||
// 已到账时长不追溯扣回。返回是否触发冻结。
|
||
func (s *annotateService) maybeFreeze(ctx context.Context, phone string) (bool, error) {
|
||
lic, err := dao.License.GetByPhone(ctx, phone)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
if lic == nil {
|
||
return false, nil
|
||
}
|
||
stats, err := dao.AnnotateRecord.CountUserStats(ctx, phone, lic.AnnotateStatsSince)
|
||
if err != nil {
|
||
return false, err
|
||
}
|
||
reviewed := stats.Approved + stats.Rejected
|
||
if reviewed < int64(annotateRewardConfig(ctx).FreezeMinReviewed) {
|
||
return false, nil // 样本不足不判(防小样本误冻)
|
||
}
|
||
cfg := annotateRewardConfig(ctx)
|
||
if float64(stats.Approved)/float64(reviewed) >= cfg.FreezeRatio {
|
||
return false, nil
|
||
}
|
||
frozenUntil := gtime.Now().Add(time.Duration(cfg.FreezeHours) * time.Hour)
|
||
if err := dao.License.UpdateAnnotateState(ctx, phone, frozenUntil, gtime.Now()); err != nil {
|
||
return false, err
|
||
}
|
||
g.Log().Infof(ctx, "标注低质冻结: %s 通过 %d/%d,冻结至 %s", phone, stats.Approved, reviewed, frozenUntil)
|
||
return true, nil
|
||
}
|
||
|
||
// AdminListRecords 用户标注记录分页(附图片文件/URL 与数据集名,内存组装)
|
||
func (s *annotateService) AdminListRecords(ctx context.Context, req *dto.AdminAnnotateRecordListReq) (*dto.AdminAnnotateRecordListRes, error) {
|
||
page, size := common.NormalizePage(req.Page, req.Size)
|
||
list, total, err := dao.AnnotateRecord.PageByFilter(ctx, req.Phone, req.TaskId, req.DatasetId, req.Status, page, size)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
res := &dto.AdminAnnotateRecordListRes{Total: total, List: make([]*dto.AdminAnnotateRecordItem, 0, len(list))}
|
||
if len(list) == 0 {
|
||
return res, nil
|
||
}
|
||
imageIds := make([]int64, 0, len(list))
|
||
for _, r := range list {
|
||
imageIds = append(imageIds, r.ImageId)
|
||
}
|
||
images, err := dao.DatasetImage.GetByIds(ctx, imageIds)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
imageMap := make(map[int64]*entity.DatasetImage, len(images))
|
||
datasetIds := make([]int64, 0)
|
||
for _, img := range images {
|
||
imageMap[img.Id] = img
|
||
datasetIds = append(datasetIds, img.DatasetId)
|
||
}
|
||
names := Training.datasetNameMap(ctx)
|
||
for _, r := range list {
|
||
item := &dto.AdminAnnotateRecordItem{
|
||
Id: r.Id,
|
||
PhoneNum: r.PhoneNum,
|
||
TaskId: r.TaskId,
|
||
ImageId: r.ImageId,
|
||
Status: r.Status,
|
||
CreatedAt: r.CreatedAt,
|
||
SubmittedAt: r.SubmittedAt,
|
||
ReviewedAt: r.ReviewedAt,
|
||
}
|
||
if img := imageMap[r.ImageId]; img != nil {
|
||
item.Filename = img.Filename
|
||
item.DatasetName = names[img.DatasetId]
|
||
item.ImageUrl = datasetImageUrl(ctx, img.DatasetId, img.Filename) // 管理端页面展示:admin token 鉴权
|
||
}
|
||
if r.LabelsJson != "" {
|
||
var boxes []*dto.AdminLabelBox
|
||
if json.Unmarshal([]byte(r.LabelsJson), &boxes) == nil {
|
||
item.Boxes = boxes
|
||
}
|
||
}
|
||
res.List = append(res.List, item)
|
||
}
|
||
return res, nil
|
||
}
|
||
|
||
// AdminUnfreeze 手动解冻:清冻结标记与统计基线(历史重新累计)
|
||
func (s *annotateService) AdminUnfreeze(ctx context.Context, req *dto.AdminAnnotateUnfreezeReq) (*dto.AdminAnnotateUnfreezeRes, error) {
|
||
lic, err := dao.License.GetByPhone(ctx, req.Phone)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if lic == nil {
|
||
return nil, gerror.New("账号不存在")
|
||
}
|
||
if err := dao.License.UpdateAnnotateState(ctx, req.Phone, nil, nil); err != nil {
|
||
return nil, err
|
||
}
|
||
return &dto.AdminAnnotateUnfreezeRes{}, nil
|
||
}
|
||
|
||
// ---------- App ----------
|
||
|
||
// AppTaskList 可领任务列表 + 我的统计(冻结中仍展示,领取时拒绝)
|
||
func (s *annotateService) AppTaskList(ctx context.Context, req *dto.AnnotateTaskListReq) (*dto.AnnotateTaskListRes, error) {
|
||
phone := common.PhoneFromCtx(ctx)
|
||
tasks, err := dao.AnnotateTask.ListPublished(ctx)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
res := &dto.AnnotateTaskListRes{List: make([]*dto.AnnotateTaskItem, 0, len(tasks))}
|
||
stats, err := s.buildStats(ctx, phone)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
res.Stats = stats
|
||
if len(tasks) == 0 {
|
||
return res, nil
|
||
}
|
||
datasetIds := make([]int64, 0, len(tasks))
|
||
taskIds := make([]int64, 0, len(tasks))
|
||
for _, t := range tasks {
|
||
datasetIds = append(datasetIds, t.DatasetId)
|
||
taskIds = append(taskIds, t.Id)
|
||
}
|
||
datasets, err := dao.Dataset.GetByIds(ctx, datasetIds)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
datasetMap := make(map[int64]*entity.Dataset, len(datasets))
|
||
for _, d := range datasets {
|
||
datasetMap[d.Id] = d
|
||
}
|
||
names := Training.datasetNameMap(ctx)
|
||
poolCnt, err := dao.DatasetImage.CountPoolByTaskIds(ctx, taskIds)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
for _, t := range tasks {
|
||
item := &dto.AnnotateTaskItem{
|
||
Id: t.Id,
|
||
Name: t.Name,
|
||
DatasetId: t.DatasetId,
|
||
DatasetName: names[t.DatasetId],
|
||
PoolRemain: poolCnt[t.Id],
|
||
}
|
||
if d := datasetMap[t.DatasetId]; d != nil {
|
||
item.Species = annotateSpecies(d)
|
||
}
|
||
res.List = append(res.List, item)
|
||
}
|
||
return res, nil
|
||
}
|
||
|
||
// AppClaim 领取:惰性释放过期锁 → 任务占用池内挑未处理过且未被锁定的图 → Serial+事务插入 pending 记录
|
||
func (s *annotateService) AppClaim(ctx context.Context, req *dto.AnnotateClaimReq) (*dto.AnnotateClaimRes, error) {
|
||
phone := common.PhoneFromCtx(ctx)
|
||
cfg := annotateRewardConfig(ctx)
|
||
task, err := dao.AnnotateTask.GetById(ctx, req.TaskId)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if task == nil {
|
||
return nil, gerror.New("标注任务不存在")
|
||
}
|
||
if task.Status != consts.AnnotateTaskPublished {
|
||
return nil, gerror.New("任务已停用")
|
||
}
|
||
lic, err := dao.License.GetByPhone(ctx, phone)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if lic == nil {
|
||
return nil, gerror.New("账号不存在")
|
||
}
|
||
if annotateFrozen(lic) != nil {
|
||
return nil, gerror.Newf("标注资格已冻结(%s 解冻),暂不可领取任务", lic.AnnotateFrozenUntil)
|
||
}
|
||
dataset, err := dao.Dataset.GetById(ctx, task.DatasetId)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if dataset == nil {
|
||
return nil, gerror.NewCode(common.CodeDatasetNotFound)
|
||
}
|
||
res := &dto.AnnotateClaimRes{Images: []*dto.AnnotateClaimImage{}, Species: annotateSpecies(dataset)}
|
||
// 惰性释放:超时未提交的领取锁直接删除(未产生任何标注,非处理记录)
|
||
if err := dao.AnnotateRecord.DeleteExpiredPending(ctx, gtime.Now().Add(-time.Duration(cfg.ClaimTimeoutHour)*time.Hour)); err != nil {
|
||
return nil, err
|
||
}
|
||
// 候选挑选放在 Serial 临界区内(挑图 + 插锁全程串行):两个用户并发领取时,
|
||
// 后者读到的 pending 锁一定包含前者刚插入的,同一张图不会分给多人——
|
||
// partial unique index(image_id WHERE pending) 是跨实例兜底的最后一道防线
|
||
now := gtime.Now()
|
||
var picked []*entity.DatasetImage
|
||
err = common.Serial().Submit(ctx, func() error {
|
||
pool, err := dao.DatasetImage.ListPoolByTask(ctx, task.Id)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
mine, err := dao.AnnotateRecord.ListImageIdsByPhone(ctx, phone)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
locked, err := dao.AnnotateRecord.ListPendingImageIds(ctx)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
exclude := make(map[int64]bool, len(mine)+len(locked))
|
||
for _, id := range mine {
|
||
exclude[id] = true
|
||
}
|
||
for _, id := range locked {
|
||
exclude[id] = true
|
||
}
|
||
picked = make([]*entity.DatasetImage, 0, cfg.ClaimSize)
|
||
for _, img := range pool {
|
||
if len(picked) >= cfg.ClaimSize {
|
||
break
|
||
}
|
||
if exclude[img.Id] {
|
||
continue
|
||
}
|
||
picked = append(picked, img)
|
||
}
|
||
if len(picked) == 0 {
|
||
return nil
|
||
}
|
||
return g.DB().Transaction(ctx, func(ctx context.Context, tx gdb.TX) error {
|
||
for _, img := range picked {
|
||
if err := dao.AnnotateRecord.InsertInTx(ctx, tx, &entity.AnnotateRecord{
|
||
PhoneNum: phone,
|
||
TaskId: task.Id,
|
||
DatasetId: task.DatasetId,
|
||
ImageId: img.Id,
|
||
Status: consts.AnnotateRecordPending,
|
||
CreatedAt: now,
|
||
}); err != nil {
|
||
return err
|
||
}
|
||
}
|
||
return nil
|
||
})
|
||
})
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if len(picked) == 0 {
|
||
return res, nil
|
||
}
|
||
dir := common.DatasetImagesDir(ctx, dataset.Name)
|
||
for _, img := range picked {
|
||
item := &dto.AnnotateClaimImage{
|
||
ImageId: img.Id,
|
||
Url: annotateImageUrl(ctx, dataset.Id, img.Filename),
|
||
}
|
||
if data, rErr := os.ReadFile(filepath.Join(dir, img.Filename)); rErr == nil {
|
||
item.Width, item.Height = imageSize(data)
|
||
}
|
||
res.Images = append(res.Images, item)
|
||
}
|
||
return res, nil
|
||
}
|
||
|
||
// AppSubmit 提交标注:写 labels_json + 待审核;每累计 RewardPerImages 张有效提交即时发时长
|
||
// (expires_at = max(now, expires_at) + RewardMinutes 分钟级顺延;当日累计超上限不再发)
|
||
func (s *annotateService) AppSubmit(ctx context.Context, req *dto.AnnotateSubmitReq) (*dto.AnnotateSubmitRes, error) {
|
||
phone := common.PhoneFromCtx(ctx)
|
||
cfg := annotateRewardConfig(ctx)
|
||
for _, box := range req.Boxes {
|
||
if box.Cx < 0 || box.Cy < 0 || box.W <= 0 || box.H <= 0 || box.Cx > 1 || box.Cy > 1 {
|
||
return nil, gerror.New("标注框坐标非法(需 0~1 归一化)")
|
||
}
|
||
}
|
||
raw, err := json.Marshal(req.Boxes)
|
||
if err != nil {
|
||
return nil, gerror.Wrap(err, "标注序列化失败")
|
||
}
|
||
record, err := dao.AnnotateRecord.GetMyPending(ctx, phone, req.ImageId)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if record == nil {
|
||
return nil, gerror.New("未领取该图片或已提交")
|
||
}
|
||
img, err := dao.DatasetImage.GetById(ctx, req.ImageId)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if img == nil {
|
||
return nil, gerror.NewCode(common.CodeImageNotFound)
|
||
}
|
||
now := gtime.Now()
|
||
// Serial 内先提交(提交计数与发放判定串行化,防并发双发)
|
||
err = common.Serial().Submit(ctx, func() error {
|
||
return g.DB().Transaction(ctx, func(ctx context.Context, tx gdb.TX) error {
|
||
if err := dao.AnnotateRecord.MarkSubmitted(ctx, record.Id, string(raw), now); err != nil {
|
||
return err
|
||
}
|
||
// App 提交一律进待审核(空框=「画面无目标」的主张,同样待审核确认)
|
||
if err := dao.DatasetImage.UpdateLabelsAndReview(ctx, img.Id, string(raw), consts.ReviewImagePending); err != nil {
|
||
return err
|
||
}
|
||
// 奖励判定:累计提交数整跨 RewardPerImages → 尝试发放
|
||
stats, err := dao.AnnotateRecord.CountUserStats(ctx, phone, nil)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if stats.SubmittedTotal%int64(cfg.RewardPerImages) != 0 {
|
||
return nil
|
||
}
|
||
dayStart := gtime.New(now.Format("Y-m-d") + " 00:00:00")
|
||
todaySum, err := dao.RewardLog.SumMinutesByPhoneSince(ctx, phone, dayStart)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if todaySum+cfg.RewardMinutes > cfg.DailyCapMinutes {
|
||
g.Log().Infof(ctx, "标注奖励触达日上限不发: %s 今日已得 %d 分钟", phone, todaySum)
|
||
return nil
|
||
}
|
||
lic, err := dao.License.GetByPhoneInTx(ctx, tx, phone)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if lic == nil {
|
||
return gerror.New("账号不存在")
|
||
}
|
||
base := now
|
||
if lic.ExpiresAt != nil && lic.ExpiresAt.Timestamp() > now.Timestamp() {
|
||
base = lic.ExpiresAt
|
||
}
|
||
newExpires := base.Add(time.Duration(cfg.RewardMinutes) * time.Minute)
|
||
if err := dao.License.ExtendExpiresInTx(ctx, tx, phone, newExpires); err != nil {
|
||
return err
|
||
}
|
||
if err := dao.RewardLog.InsertInTx(ctx, tx, phone, cfg.RewardMinutes, record.TaskId, now); err != nil {
|
||
return err
|
||
}
|
||
common.ClearCache(ctx, "license:"+phone)
|
||
g.Log().Infof(ctx, "标注奖励发放: %s +%d 分钟,到期 %s", phone, cfg.RewardMinutes, newExpires)
|
||
return nil
|
||
})
|
||
})
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
// 回程统计(展示口径,允许轻微滞后)
|
||
dayStart := gtime.New(now.Format("Y-m-d") + " 00:00:00")
|
||
todaySum, _ := dao.RewardLog.SumMinutesByPhoneSince(ctx, phone, dayStart)
|
||
lic, _ := dao.License.GetByPhone(ctx, phone)
|
||
stats, _ := dao.AnnotateRecord.CountUserStats(ctx, phone, nil)
|
||
var expiresAt *gtime.Time
|
||
if lic != nil {
|
||
expiresAt = lic.ExpiresAt
|
||
}
|
||
return &dto.AnnotateSubmitRes{
|
||
Granted: stats.SubmittedTotal > 0 && stats.SubmittedTotal%int64(cfg.RewardPerImages) == 0,
|
||
Minutes: cfg.RewardMinutes,
|
||
TodayEarnedMinutes: todaySum,
|
||
ExpiresAt: expiresAt,
|
||
ProgressDone: stats.SubmittedTotal % int64(cfg.RewardPerImages),
|
||
}, nil
|
||
}
|
||
|
||
// AppMe 我的标注统计
|
||
func (s *annotateService) AppMe(ctx context.Context, req *dto.AnnotateMeReq) (*dto.AnnotateMeRes, error) {
|
||
stats, err := s.buildStats(ctx, common.PhoneFromCtx(ctx))
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
return &dto.AnnotateMeRes{Stats: stats}, nil
|
||
}
|
||
|
||
// buildStats 我的统计组装(任务列表/我的页/提交回程共用)
|
||
func (s *annotateService) buildStats(ctx context.Context, phone string) (*dto.AnnotateStats, error) {
|
||
cfg := annotateRewardConfig(ctx)
|
||
stats := &dto.AnnotateStats{
|
||
RewardPerImages: cfg.RewardPerImages,
|
||
RewardMinutes: cfg.RewardMinutes,
|
||
DailyCapMinutes: cfg.DailyCapMinutes,
|
||
ApproveRatio: 1, // 无样本不惩罚
|
||
ProgressDone: 0,
|
||
}
|
||
lic, err := dao.License.GetByPhone(ctx, phone)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
var since *gtime.Time
|
||
if lic != nil {
|
||
since = lic.AnnotateStatsSince
|
||
stats.FrozenUntil = annotateFrozen(lic)
|
||
stats.ExpiresAt = lic.ExpiresAt
|
||
}
|
||
rec, err := dao.AnnotateRecord.CountUserStats(ctx, phone, since)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
stats.SubmittedTotal = rec.SubmittedTotal
|
||
stats.Approved = rec.Approved
|
||
stats.Rejected = rec.Rejected
|
||
if reviewed := rec.Approved + rec.Rejected; reviewed > 0 {
|
||
stats.ApproveRatio = float64(rec.Approved) / float64(reviewed)
|
||
}
|
||
stats.ProgressDone = rec.SubmittedTotal % int64(cfg.RewardPerImages)
|
||
total, err := dao.RewardLog.SumMinutesByPhone(ctx, phone)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
stats.TotalEarnedMinutes = total
|
||
dayStart := gtime.New(gtime.Now().Format("Y-m-d") + " 00:00:00")
|
||
today, err := dao.RewardLog.SumMinutesByPhoneSince(ctx, phone, dayStart)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
stats.TodayEarnedMinutes = today
|
||
return stats, nil
|
||
}
|