252 lines
9.7 KiB
Go
252 lines
9.7 KiB
Go
package grpc
|
||
|
||
import (
|
||
"context"
|
||
|
||
activityv1 "hyapp.local/api/proto/activity/v1"
|
||
"hyapp/pkg/appcode"
|
||
"hyapp/pkg/xerr"
|
||
taskdomain "hyapp/services/activity-service/internal/domain/task"
|
||
taskservice "hyapp/services/activity-service/internal/service/task"
|
||
)
|
||
|
||
// TaskServer 把任务查询、事件消费和领奖用例适配为 activity-service gRPC。
|
||
type TaskServer struct {
|
||
activityv1.UnimplementedTaskServiceServer
|
||
|
||
svc *taskservice.Service
|
||
}
|
||
|
||
// NewTaskServer 创建 App 任务服务 gRPC adapter。
|
||
func NewTaskServer(svc *taskservice.Service) *TaskServer {
|
||
return &TaskServer{svc: svc}
|
||
}
|
||
|
||
// ListUserTasks 返回当前用户 daily/exclusive 任务页。
|
||
func (s *TaskServer) ListUserTasks(ctx context.Context, req *activityv1.ListUserTasksRequest) (*activityv1.ListUserTasksResponse, error) {
|
||
ctx = appcode.WithContext(ctx, req.GetMeta().GetAppCode())
|
||
result, err := s.svc.ListUserTasks(ctx, req.GetUserId())
|
||
if err != nil {
|
||
return nil, xerr.ToGRPCError(err)
|
||
}
|
||
return taskListToProto(result), nil
|
||
}
|
||
|
||
// ClaimTaskReward 领取已完成任务奖励,金币入账仍由 wallet-service 完成。
|
||
func (s *TaskServer) ClaimTaskReward(ctx context.Context, req *activityv1.ClaimTaskRewardRequest) (*activityv1.ClaimTaskRewardResponse, error) {
|
||
ctx = appcode.WithContext(ctx, req.GetMeta().GetAppCode())
|
||
claim, err := s.svc.ClaimTaskReward(ctx, taskdomain.ClaimCommand{
|
||
UserID: req.GetUserId(),
|
||
TaskID: req.GetTaskId(),
|
||
TaskType: req.GetTaskType(),
|
||
CycleKey: req.GetTaskDay(),
|
||
CommandID: req.GetCommandId(),
|
||
})
|
||
if err != nil {
|
||
return nil, xerr.ToGRPCError(err)
|
||
}
|
||
return claimToProto(claim), nil
|
||
}
|
||
|
||
// ConsumeTaskEvent 消费由 outbox worker 投递的服务端事实事件。
|
||
func (s *TaskServer) ConsumeTaskEvent(ctx context.Context, req *activityv1.ConsumeTaskEventRequest) (*activityv1.ConsumeTaskEventResponse, error) {
|
||
ctx = appcode.WithContext(ctx, req.GetMeta().GetAppCode())
|
||
result, err := s.svc.ConsumeTaskEvent(ctx, taskdomain.Event{
|
||
EventID: req.GetEventId(),
|
||
EventType: req.GetEventType(),
|
||
SourceService: req.GetSourceService(),
|
||
UserID: req.GetUserId(),
|
||
MetricType: req.GetMetricType(),
|
||
Value: req.GetValue(),
|
||
OccurredAtMS: req.GetOccurredAtMs(),
|
||
DimensionsJSON: req.GetDimensionsJson(),
|
||
})
|
||
if err != nil {
|
||
return nil, xerr.ToGRPCError(err)
|
||
}
|
||
return &activityv1.ConsumeTaskEventResponse{
|
||
EventId: result.EventID,
|
||
Status: result.Status,
|
||
MatchedTaskCount: result.MatchedTaskCount,
|
||
}, nil
|
||
}
|
||
|
||
// AdminTaskServer 暴露后台任务配置管理入口。
|
||
type AdminTaskServer struct {
|
||
activityv1.UnimplementedAdminTaskServiceServer
|
||
|
||
svc *taskservice.Service
|
||
}
|
||
|
||
// NewAdminTaskServer 创建后台任务配置 gRPC adapter。
|
||
func NewAdminTaskServer(svc *taskservice.Service) *AdminTaskServer {
|
||
return &AdminTaskServer{svc: svc}
|
||
}
|
||
|
||
// ListTaskDefinitions 返回后台任务定义列表。
|
||
func (s *AdminTaskServer) ListTaskDefinitions(ctx context.Context, req *activityv1.ListTaskDefinitionsRequest) (*activityv1.ListTaskDefinitionsResponse, error) {
|
||
ctx = appcode.WithContext(ctx, req.GetMeta().GetAppCode())
|
||
items, total, err := s.svc.ListTaskDefinitions(ctx, taskdomain.DefinitionQuery{
|
||
TaskType: req.GetTaskType(),
|
||
Category: req.GetCategory(),
|
||
Status: req.GetStatus(),
|
||
Keyword: req.GetKeyword(),
|
||
Page: req.GetPage(),
|
||
PageSize: req.GetPageSize(),
|
||
})
|
||
if err != nil {
|
||
return nil, xerr.ToGRPCError(err)
|
||
}
|
||
resp := &activityv1.ListTaskDefinitionsResponse{Tasks: make([]*activityv1.TaskDefinition, 0, len(items)), Total: total}
|
||
for _, item := range items {
|
||
resp.Tasks = append(resp.Tasks, definitionToProto(item))
|
||
}
|
||
return resp, nil
|
||
}
|
||
|
||
// UpsertTaskDefinition 创建或编辑任务配置。
|
||
func (s *AdminTaskServer) UpsertTaskDefinition(ctx context.Context, req *activityv1.UpsertTaskDefinitionRequest) (*activityv1.UpsertTaskDefinitionResponse, error) {
|
||
ctx = appcode.WithContext(ctx, req.GetMeta().GetAppCode())
|
||
item, created, err := s.svc.UpsertTaskDefinition(ctx, taskdomain.DefinitionCommand{
|
||
TaskID: req.GetTaskId(),
|
||
TaskType: req.GetTaskType(),
|
||
Category: req.GetCategory(),
|
||
MetricType: req.GetMetricType(),
|
||
Title: req.GetTitle(),
|
||
Description: req.GetDescription(),
|
||
AudienceType: req.GetAudienceType(),
|
||
IconKey: req.GetIconKey(),
|
||
IconURL: req.GetIconUrl(),
|
||
ActionType: req.GetActionType(),
|
||
ActionParam: req.GetActionParam(),
|
||
ActionPayloadJSON: req.GetActionPayloadJson(),
|
||
DimensionFilterJSON: req.GetDimensionFilterJson(),
|
||
TargetValue: req.GetTargetValue(),
|
||
TargetUnit: req.GetTargetUnit(),
|
||
RewardCoinAmount: req.GetRewardCoinAmount(),
|
||
RewardAssetType: req.GetRewardAssetType(),
|
||
Status: req.GetStatus(),
|
||
SortOrder: req.GetSortOrder(),
|
||
EffectiveFromMS: req.GetEffectiveFromMs(),
|
||
EffectiveToMS: req.GetEffectiveToMs(),
|
||
OperatorAdminID: req.GetOperatorAdminId(),
|
||
})
|
||
if err != nil {
|
||
return nil, xerr.ToGRPCError(err)
|
||
}
|
||
return &activityv1.UpsertTaskDefinitionResponse{Task: definitionToProto(item), Created: created}, nil
|
||
}
|
||
|
||
// SetTaskDefinitionStatus 启停或归档任务。
|
||
func (s *AdminTaskServer) SetTaskDefinitionStatus(ctx context.Context, req *activityv1.SetTaskDefinitionStatusRequest) (*activityv1.SetTaskDefinitionStatusResponse, error) {
|
||
ctx = appcode.WithContext(ctx, req.GetMeta().GetAppCode())
|
||
item, err := s.svc.SetTaskDefinitionStatus(ctx, req.GetTaskId(), req.GetStatus(), req.GetOperatorAdminId())
|
||
if err != nil {
|
||
return nil, xerr.ToGRPCError(err)
|
||
}
|
||
return &activityv1.SetTaskDefinitionStatusResponse{Task: definitionToProto(item)}, nil
|
||
}
|
||
|
||
// PublishTaskRewardPolicy 保留旧 RPC 的 wire contract,但拒绝继续写入已退役的默认奖励政策。
|
||
// 调用方必须把资产类型显式保存到每条任务定义,避免一次 App 级发布静默改写既有任务语义。
|
||
func (s *AdminTaskServer) PublishTaskRewardPolicy(context.Context, *activityv1.PublishTaskRewardPolicyRequest) (*activityv1.PublishTaskRewardPolicyResponse, error) {
|
||
return nil, xerr.ToGRPCError(xerr.New(xerr.Conflict, "task reward policy publishing is retired; configure reward_asset_type on each task definition"))
|
||
}
|
||
|
||
func taskListToProto(result taskdomain.ListResult) *activityv1.ListUserTasksResponse {
|
||
resp := &activityv1.ListUserTasksResponse{
|
||
Sections: make([]*activityv1.TaskSection, 0, len(result.Sections)),
|
||
ServerTimeMs: result.ServerTimeMS,
|
||
NextRefreshAtMs: result.NextRefreshAtMS,
|
||
}
|
||
for _, section := range result.Sections {
|
||
target := &activityv1.TaskSection{
|
||
Section: section.Section,
|
||
Items: make([]*activityv1.TaskItem, 0, len(section.Items)),
|
||
ServerTimeMs: section.ServerTimeMS,
|
||
NextRefreshAtMs: section.NextRefreshAtMS,
|
||
}
|
||
for _, item := range section.Items {
|
||
target.Items = append(target.Items, itemToProto(item))
|
||
}
|
||
resp.Sections = append(resp.Sections, target)
|
||
}
|
||
return resp
|
||
}
|
||
|
||
func itemToProto(item taskdomain.Item) *activityv1.TaskItem {
|
||
return &activityv1.TaskItem{
|
||
TaskId: item.TaskID,
|
||
TaskType: item.TaskType,
|
||
Category: item.Category,
|
||
MetricType: item.MetricType,
|
||
Title: item.Title,
|
||
Description: item.Description,
|
||
AudienceType: item.AudienceType,
|
||
IconKey: item.IconKey,
|
||
IconUrl: item.IconURL,
|
||
ActionType: item.ActionType,
|
||
ActionParam: item.ActionParam,
|
||
ActionPayloadJson: item.ActionPayloadJSON,
|
||
DimensionFilterJson: item.DimensionFilterJSON,
|
||
TargetValue: item.TargetValue,
|
||
TargetUnit: item.TargetUnit,
|
||
ProgressValue: item.ProgressValue,
|
||
RewardCoinAmount: item.RewardCoinAmount,
|
||
RewardAssetType: item.RewardAssetType,
|
||
Status: item.UserStatus,
|
||
Claimable: item.Claimable,
|
||
TaskDay: item.TaskDay,
|
||
ServerTimeMs: item.ServerTimeMS,
|
||
NextRefreshAtMs: item.NextRefreshAtMS,
|
||
SortOrder: item.SortOrder,
|
||
Version: item.Version,
|
||
}
|
||
}
|
||
|
||
func claimToProto(claim taskdomain.Claim) *activityv1.ClaimTaskRewardResponse {
|
||
return &activityv1.ClaimTaskRewardResponse{
|
||
ClaimId: claim.ClaimID,
|
||
TaskId: claim.TaskID,
|
||
TaskType: claim.TaskType,
|
||
TaskDay: claim.CycleKey,
|
||
RewardCoinAmount: claim.RewardCoinAmount,
|
||
RewardAssetType: claim.RewardAssetType,
|
||
Status: claim.Status,
|
||
WalletTransactionId: claim.WalletTransactionID,
|
||
GrantedAtMs: claim.UpdatedAtMS,
|
||
Claimed: claim.Status == taskdomain.ClaimStatusGranted,
|
||
}
|
||
}
|
||
|
||
func definitionToProto(item taskdomain.Definition) *activityv1.TaskDefinition {
|
||
return &activityv1.TaskDefinition{
|
||
TaskId: item.TaskID,
|
||
TaskType: item.TaskType,
|
||
Category: item.Category,
|
||
MetricType: item.MetricType,
|
||
Title: item.Title,
|
||
Description: item.Description,
|
||
AudienceType: item.AudienceType,
|
||
IconKey: item.IconKey,
|
||
IconUrl: item.IconURL,
|
||
ActionType: item.ActionType,
|
||
ActionParam: item.ActionParam,
|
||
ActionPayloadJson: item.ActionPayloadJSON,
|
||
DimensionFilterJson: item.DimensionFilterJSON,
|
||
TargetValue: item.TargetValue,
|
||
TargetUnit: item.TargetUnit,
|
||
RewardCoinAmount: item.RewardCoinAmount,
|
||
RewardAssetType: item.RewardAssetType,
|
||
Status: item.Status,
|
||
SortOrder: item.SortOrder,
|
||
Version: item.Version,
|
||
EffectiveFromMs: item.EffectiveFromMS,
|
||
EffectiveToMs: item.EffectiveToMS,
|
||
CreatedByAdminId: item.CreatedByAdminID,
|
||
UpdatedByAdminId: item.UpdatedByAdminID,
|
||
CreatedAtMs: item.CreatedAtMS,
|
||
UpdatedAtMs: item.UpdatedAtMS,
|
||
}
|
||
}
|