Files
cmroubao/backend-api/internal/transport/httpapi/admin_handlers.go
T

577 lines
14 KiB
Go

package httpapi
import (
"encoding/json"
"errors"
"io"
"mime"
"net/http"
"strconv"
"strings"
"time"
"cmroubao/backend-api/internal/domain"
"cmroubao/backend-api/internal/transport/authcommon"
"cmroubao/backend-api/internal/usecase"
"github.com/gin-gonic/gin"
)
const (
localAdminSubject = "local-admin"
maxJSONBodyBytes = 64 << 10
maxMultipartBytes = 21 << 20
)
type AdminServices struct {
Assets *usecase.AssetService
Tasks *usecase.TaskService
}
func (s AdminServices) validate() error {
if s.Assets == nil || s.Tasks == nil {
return errors.New("admin services are required")
}
return nil
}
type adminHandlers struct {
services AdminServices
}
func registerAdminAPI(routes gin.IRoutes, services AdminServices) error {
if err := services.validate(); err != nil {
return err
}
handler := &adminHandlers{services: services}
routes.POST("/api/v1/assets", handler.uploadAsset)
routes.GET("/api/v1/assets/:id/content", handler.assetContent)
routes.POST("/api/v1/tasks", handler.createTask)
routes.GET("/api/v1/tasks", handler.listTasks)
routes.GET("/api/v1/tasks/:id", handler.taskDetail)
routes.POST("/api/v1/tasks/:id/cancel", handler.cancelTask)
return nil
}
func (h *adminHandlers) uploadAsset(ctx *gin.Context) {
if !hasMediaType(ctx, "multipart/form-data") {
writePublicError(
ctx,
http.StatusUnsupportedMediaType,
"UNSUPPORTED_MEDIA_TYPE",
"multipart/form-data is required",
false,
gin.H{},
)
return
}
ctx.Request.Body = http.MaxBytesReader(
ctx.Writer,
ctx.Request.Body,
maxMultipartBytes,
)
if err := ctx.Request.ParseMultipartForm(maxMultipartBytes); err != nil {
writePublicError(
ctx,
http.StatusRequestEntityTooLarge,
"ASSET_TOO_LARGE",
"reference image exceeds the allowed size",
false,
gin.H{},
)
return
}
if ctx.Request.MultipartForm != nil {
defer ctx.Request.MultipartForm.RemoveAll()
}
purpose := strings.TrimSpace(ctx.PostForm("purpose"))
if purpose != domain.AssetPurposeTaskReference {
writePublicError(
ctx,
http.StatusUnprocessableEntity,
"ASSET_PURPOSE_INVALID",
"asset purpose is not supported",
false,
fieldDetails("purpose", "must be TASK_REFERENCE"),
)
return
}
if strings.TrimSpace(ctx.PostForm("task_id")) != "" {
writePublicError(
ctx,
http.StatusUnprocessableEntity,
"ASSET_TASK_ID_INVALID",
"task_id must be empty for a task reference",
false,
fieldDetails("task_id", "must be empty"),
)
return
}
files := ctx.Request.MultipartForm.File["file"]
if len(files) != 1 {
writePublicError(
ctx,
http.StatusBadRequest,
"ASSET_FILE_REQUIRED",
"exactly one reference image is required",
false,
fieldDetails("file", "exactly one file is required"),
)
return
}
content, err := files[0].Open()
if err != nil {
writePublicError(
ctx,
http.StatusUnprocessableEntity,
"ASSET_IMAGE_INVALID",
"reference image is invalid",
false,
fieldDetails("file", "cannot be read"),
)
return
}
defer content.Close()
result, err := h.services.Assets.UploadTaskReference(
ctx.Request.Context(),
usecase.UploadTaskReferenceCommand{
CreatorSubject: localAdminSubject,
IdempotencyKey: ctx.GetHeader("Idempotency-Key"),
DeclaredMediaType: files[0].Header.Get("Content-Type"),
Content: content,
},
)
if err != nil {
writeUsecaseError(ctx, err)
return
}
ctx.JSON(http.StatusCreated, assetResponse(result.Asset))
}
func (h *adminHandlers) assetContent(ctx *gin.Context) {
result, err := h.services.Assets.OpenTaskReference(
ctx.Request.Context(),
localAdminSubject,
ctx.Param("id"),
)
if err != nil {
writeUsecaseError(ctx, err)
return
}
defer result.Content.Close()
writeAssetContent(ctx, result)
}
func (h *adminHandlers) createTask(ctx *gin.Context) {
if !hasMediaType(ctx, "application/json") {
writePublicError(
ctx,
http.StatusUnsupportedMediaType,
"UNSUPPORTED_MEDIA_TYPE",
"application/json is required",
false,
gin.H{},
)
return
}
var request struct {
SourceRef *string `json:"source_ref"`
Title string `json:"title"`
Description string `json:"description"`
SKU string `json:"sku"`
ImageAssetID string `json:"image_asset_id"`
Quantity int `json:"quantity"`
MaxBudget *string `json:"max_budget"`
}
if err := decodeJSON(ctx, &request); err != nil {
writePublicError(
ctx,
http.StatusBadRequest,
"INVALID_JSON",
"request body must be valid JSON",
false,
gin.H{},
)
return
}
result, err := h.services.Tasks.Create(
ctx.Request.Context(),
usecase.CreateTaskCommand{
CreatorSubject: localAdminSubject,
ActorUserID: adminActorUserID(ctx),
IdempotencyKey: ctx.GetHeader("Idempotency-Key"),
SourceRef: request.SourceRef,
Title: request.Title,
Description: request.Description,
SKU: request.SKU,
ImageAssetID: request.ImageAssetID,
Quantity: request.Quantity,
MaxBudget: request.MaxBudget,
},
)
if err != nil {
writeUsecaseError(ctx, err)
return
}
ctx.JSON(http.StatusCreated, taskSummaryResponse(result.Task))
}
func adminActorUserID(ctx *gin.Context) string {
principal, ok := authcommon.Principal(ctx.Request.Context())
if !ok || principal.Role != domain.UserRoleAdmin {
return ""
}
return principal.UserID
}
func (h *adminHandlers) listTasks(ctx *gin.Context) {
query := usecase.ListTasksQuery{
CreatorSubject: localAdminSubject,
Query: ctx.Query("q"),
Cursor: ctx.Query("cursor"),
}
if value := strings.TrimSpace(ctx.Query("status")); value != "" {
query.Status = &value
}
if value := strings.TrimSpace(ctx.Query("created_from")); value != "" {
parsed, err := time.Parse(time.RFC3339, value)
if err != nil {
writeFilterError(ctx, "created_from", "must be RFC3339")
return
}
query.CreatedFrom = &parsed
}
if value := strings.TrimSpace(ctx.Query("created_to")); value != "" {
parsed, err := time.Parse(time.RFC3339, value)
if err != nil {
writeFilterError(ctx, "created_to", "must be RFC3339")
return
}
query.CreatedTo = &parsed
}
if value := strings.TrimSpace(ctx.Query("limit")); value != "" {
parsed, err := strconv.Atoi(value)
if err != nil {
writeFilterError(ctx, "limit", "must be an integer")
return
}
query.Limit = parsed
}
page, err := h.services.Tasks.List(ctx.Request.Context(), query)
if err != nil {
writeUsecaseError(ctx, err)
return
}
items := make([]gin.H, 0, len(page.Items))
for _, task := range page.Items {
items = append(items, taskListItemResponse(task))
}
var nextCursor any
if page.NextCursor != "" {
nextCursor = page.NextCursor
}
ctx.Header("Cache-Control", "no-store")
ctx.JSON(http.StatusOK, gin.H{
"items": items,
"next_cursor": nextCursor,
})
}
func (h *adminHandlers) taskDetail(ctx *gin.Context) {
detail, err := h.services.Tasks.Get(
ctx.Request.Context(),
localAdminSubject,
ctx.Param("id"),
)
if err != nil {
writeUsecaseError(ctx, err)
return
}
events := make([]gin.H, 0, len(detail.Events))
for _, event := range detail.Events {
events = append(events, gin.H{
"id": event.ID,
"actor_user_id": event.ActorUserID,
"actor_device_id": event.ActorDeviceID,
"type": event.Type,
"message": event.Message,
"occurred_at": formatTime(event.OccurredAt),
})
}
var claim any
if detail.Task.ClaimGeneration > 0 {
claim = gin.H{
"user_id": detail.Task.ClaimedByUserID,
"device_id": detail.Task.ClaimedByDeviceID,
"generation": detail.Task.ClaimGeneration,
"issued_at": formatOptionalTime(detail.Task.ClaimIssuedAt),
"expires_at": formatOptionalTime(detail.Task.ClaimExpiresAt),
"cancel_requested_at": formatOptionalTime(
detail.Task.CancelRequestedAt,
),
"cancel_requested_by_user_id": detail.Task.CancelRequestedByUserID,
}
}
var execution any
if detail.Execution != nil {
execution = executionResponse(*detail.Execution)
}
ctx.Header("Cache-Control", "no-store")
ctx.JSON(http.StatusOK, gin.H{
"id": detail.Task.ID,
"status": detail.Task.Status,
"version": detail.Task.Version,
"created_at": formatTime(detail.Task.CreatedAt),
"updated_at": formatTime(detail.Task.UpdatedAt),
"original_requirement": gin.H{
"source_ref": detail.Task.SourceRef,
"title": detail.Task.Title,
"description": detail.Task.Description,
"sku": detail.Task.SKU,
"image_asset_id": detail.Task.ImageAssetID,
"quantity": detail.Task.Quantity,
"max_budget": domain.FormatOptionalCNY(detail.Task.MaxBudgetCents),
"currency": detail.Task.Currency,
},
"derived_requirement": nil,
"claim": claim,
"execution": execution,
"events": events,
"assets": []gin.H{
assetResponse(detail.Asset),
},
})
}
func (h *adminHandlers) cancelTask(ctx *gin.Context) {
if !hasMediaType(ctx, "application/json") {
writePublicError(
ctx,
http.StatusUnsupportedMediaType,
"UNSUPPORTED_MEDIA_TYPE",
"application/json is required",
false,
gin.H{},
)
return
}
var request struct {
Reason string `json:"reason"`
}
if err := decodeJSON(ctx, &request); err != nil {
writePublicError(
ctx,
http.StatusBadRequest,
"INVALID_JSON",
"request body must be valid JSON",
false,
gin.H{},
)
return
}
task, err := h.services.Tasks.Cancel(
ctx.Request.Context(),
usecase.CancelTaskCommand{
CreatorSubject: localAdminSubject,
ActorUserID: adminActorUserID(ctx),
TaskID: ctx.Param("id"),
Reason: request.Reason,
},
)
if err != nil {
writeUsecaseError(ctx, err)
return
}
ctx.JSON(http.StatusOK, taskSummaryResponse(task))
}
func decodeJSON(ctx *gin.Context, target any) error {
ctx.Request.Body = http.MaxBytesReader(
ctx.Writer,
ctx.Request.Body,
maxJSONBodyBytes,
)
decoder := json.NewDecoder(ctx.Request.Body)
decoder.DisallowUnknownFields()
if err := decoder.Decode(target); err != nil {
return err
}
var extra any
if err := decoder.Decode(&extra); !errors.Is(err, io.EOF) {
return errors.New("request must contain one JSON value")
}
return nil
}
func hasMediaType(ctx *gin.Context, expected string) bool {
mediaType, _, err := mime.ParseMediaType(ctx.GetHeader("Content-Type"))
return err == nil && strings.EqualFold(mediaType, expected)
}
func writeFilterError(ctx *gin.Context, field, message string) {
writePublicError(
ctx,
http.StatusBadRequest,
"TASK_LIST_FILTER_INVALID",
"task list filter is invalid",
false,
fieldDetails(field, message),
)
}
func writeUsecaseError(ctx *gin.Context, err error) {
var typed *usecase.Error
if !errors.As(err, &typed) {
writePublicError(
ctx,
http.StatusInternalServerError,
"INTERNAL_ERROR",
"internal server error",
false,
gin.H{},
)
return
}
status := http.StatusInternalServerError
switch typed.Kind {
case usecase.ErrorKindInvalid:
status = http.StatusUnprocessableEntity
switch typed.Code {
case "REQUEST_VALIDATION_FAILED",
"ASSET_FILE_REQUIRED",
"IDEMPOTENCY_KEY_REQUIRED":
status = http.StatusBadRequest
case "ASSET_TOO_LARGE":
status = http.StatusRequestEntityTooLarge
case "ASSET_MEDIA_TYPE_UNSUPPORTED":
status = http.StatusUnsupportedMediaType
}
case usecase.ErrorKindNotFound:
status = http.StatusNotFound
case usecase.ErrorKindForbidden:
status = http.StatusForbidden
case usecase.ErrorKindConflict:
status = http.StatusConflict
case usecase.ErrorKindUnavailable:
status = http.StatusServiceUnavailable
}
details := gin.H{}
if len(typed.Fields) > 0 {
details["fields"] = typed.Fields
}
writePublicError(
ctx,
status,
typed.Code,
typed.Message,
typed.Retryable,
details,
)
}
func writePublicError(
ctx *gin.Context,
status int,
code string,
message string,
retryable bool,
details gin.H,
) {
requestID, _ := ctx.Get(requestIDContextKey)
if details == nil {
details = gin.H{}
}
ctx.Header("Cache-Control", "no-store")
ctx.JSON(status, gin.H{
"error": gin.H{
"code": code,
"message": message,
"retryable": retryable,
"details": details,
},
"request_id": requestID,
})
}
func fieldDetails(field, message string) gin.H {
return gin.H{"fields": gin.H{field: message}}
}
func assetResponse(asset domain.Asset) gin.H {
return gin.H{
"id": asset.ID,
"purpose": asset.Purpose,
"media_type": asset.MediaType,
"size_bytes": asset.SizeBytes,
"sha256": asset.SHA256,
"created_at": formatTime(asset.CreatedAt),
}
}
func taskSummaryResponse(task domain.PurchaseTask) gin.H {
return gin.H{
"id": task.ID,
"status": task.Status,
"title": task.Title,
"sku": task.SKU,
"quantity": task.Quantity,
"max_budget": domain.FormatOptionalCNY(task.MaxBudgetCents),
"created_at": formatTime(task.CreatedAt),
"updated_at": formatTime(task.UpdatedAt),
"version": task.Version,
"cancel_requested": task.CancelRequestedAt != nil,
"cancel_requested_at": formatOptionalTime(
task.CancelRequestedAt,
),
}
}
func taskListItemResponse(task domain.PurchaseTask) gin.H {
response := taskSummaryResponse(task)
response["device_name"] = nil
return response
}
func formatTime(value time.Time) string {
return value.UTC().Format(time.RFC3339Nano)
}
func formatOptionalTime(value *time.Time) any {
if value == nil {
return nil
}
return formatTime(*value)
}
func executionResponse(execution domain.TaskExecution) gin.H {
return gin.H{
"id": execution.ID,
"attempt_no": execution.AttemptNo,
"claim_generation": execution.ClaimGeneration,
"user_id": execution.UserID,
"device_id": execution.DeviceID,
"current_step": execution.CurrentStep,
"order_submitted": execution.OrderSubmitted,
"started_at": formatTime(execution.StartedAt),
"last_heartbeat_at": formatOptionalTime(
execution.LastHeartbeatAt,
),
"finished_at": formatOptionalTime(execution.FinishedAt),
}
}
func writeAssetContent(ctx *gin.Context, result usecase.AssetContent) {
ctx.Header("Cache-Control", "private, no-store")
ctx.Header("Content-Type", result.Asset.MediaType)
ctx.Header("Content-Length", strconv.FormatInt(result.Asset.SizeBytes, 10))
ctx.Header("ETag", `"`+result.Asset.SHA256+`"`)
ctx.Header("X-Content-Type-Options", "nosniff")
ctx.Header(
"Content-Disposition",
`inline; filename="`+result.Asset.ID+`.jpg"`,
)
ctx.Status(http.StatusOK)
_, _ = io.Copy(ctx.Writer, result.Content)
}