Files
sub2api/backend/internal/handler/batch_image_handler.go
T
shaw 80a229bce5 fix(batch-image): 修复审计发现的计费死锁、状态机与队列原子性缺陷
修复 PR #3768 批量图像 MVP 合并后审计报告中的全部问题:

结算与计费(高危):
- 所有 SETTLEMENT_* 失败(超冻结/计数非法/manifest 冲突/定价缺失/扣费失败)
  统一计入 retry_count 并在耗尽时释放冻结转 failed,消灭 settling 无限
  requeue 导致的冻结余额永久锁死
- 耗尽出口的释放指纹统一为 RequestHash,与 processor/Cancel/recovery 一致;
  release 遇同 request id 指纹冲突视为幂等成功,治愈历史毒消息
- 管理端校验 hold_multiplier >= discount_multiplier,定价快照对存量脏数据钳制
- 释放前校验 per-job hold claim(dedup+归档表),杜绝幻影释放

索引对账(高危):
- provider 输出与提交 custom_id 集对账:未知条目丢弃并记事件,
  漏项补 PROVIDER_RESULT_MISSING 失败行,保证 success+fail == item_count

提交与恢复(高危):
- 提交前转 uploading 并在 provider.Submit 期间心跳刷新 updated_at;
  恢复扫描改为原子复核(FailStaleUnsubmittedBatchImageJob),
  消灭慢提交被误杀退款而上游任务照常计费的孤儿场景
- 上游任务创建成功但本地状态推进失败时,尽力取消上游并清理输入
- recovery 释放失败时入队交由 worker releaseTerminalHold 兜底重试

队列与并发(中危):
- Enqueue(SetNX+LPush)与 Reserve(BRPop+ZAdd)均改为 Lua 原子脚本,
  消灭崩溃窗口导致 job 脱离队列、被 7 天 inflight 键锁死
- 锁冲突按 LockConflictDelay 重新入队(原直接丢弃需等 10 分钟 stale 恢复)
- 处理期间心跳:active zset 续期(ZAddXX 防幽灵成员)+ 锁 TTL 续期
- ReplaceBatchImageItemsForJob 增加 indexing 状态守卫,防掉队 worker 重写账目

存量回归(中危):
- image-only 定价条目(仅图片价无 token 价)恢复 token 计费 fail-closed,
  不再按 $0 计费;图片计费路径不受影响
- 鉴权余额门槛恢复 balance <= 0 语义,MinimumBalanceReserve 不再作硬 403

加固:
- ZIP max_items 钳制到管理员上限;Submit 补齐 Platform==Gemini 校验;
  gemini downloadUri 跟随前做 host 白名单校验
- 批量客户端改用共享 httpclient(拨号/TLS/响应头超时有界)
- 审计点名的忽略错误(MarkDownloaded/SettlementFailed/AppendEvent 等)改为记日志
2026-07-07 18:53:56 +08:00

271 lines
8.0 KiB
Go

package handler
import (
"errors"
"io"
"net/http"
"strconv"
"strings"
infraerrors "github.com/Wei-Shaw/sub2api/internal/pkg/errors"
"github.com/Wei-Shaw/sub2api/internal/pkg/logger"
"github.com/Wei-Shaw/sub2api/internal/server/middleware"
"github.com/Wei-Shaw/sub2api/internal/service"
"github.com/gin-gonic/gin"
"go.uber.org/zap"
)
type BatchImageHandler struct {
service *service.BatchImagePublicService
download *service.BatchImageDownloadService
cleanup *service.BatchImageCleanupService
}
func NewBatchImageHandler(service *service.BatchImagePublicService, download *service.BatchImageDownloadService, cleanup *service.BatchImageCleanupService) *BatchImageHandler {
return &BatchImageHandler{service: service, download: download, cleanup: cleanup}
}
func (h *BatchImageHandler) Submit(c *gin.Context) {
var req service.BatchImageSubmitRequest
if err := c.ShouldBindJSON(&req); err != nil {
batchImageError(c, service.ErrBatchImageInvalidItems)
return
}
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
got, err := h.service.Submit(c.Request.Context(), owner, req, c.GetHeader("Idempotency-Key"))
if err != nil {
batchImageError(c, err)
return
}
c.JSON(http.StatusOK, got)
}
func (h *BatchImageHandler) Get(c *gin.Context) {
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
got, err := h.service.Get(c.Request.Context(), owner, c.Param("id"))
if err != nil {
batchImageError(c, err)
return
}
c.JSON(http.StatusOK, got)
}
func (h *BatchImageHandler) List(c *gin.Context) {
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
limit, _ := strconv.Atoi(c.Query("limit"))
got, err := h.service.List(c.Request.Context(), owner, service.BatchImageJobsQuery{
Status: c.Query("status"),
TaskName: c.Query("task_name"),
Downloaded: c.Query("downloaded"),
From: c.Query("from"),
To: c.Query("to"),
Limit: limit,
Cursor: c.Query("cursor"),
})
if err != nil {
batchImageError(c, err)
return
}
c.JSON(http.StatusOK, got)
}
func (h *BatchImageHandler) Models(c *gin.Context) {
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
got, err := h.service.ListModels(c.Request.Context(), owner)
if err != nil {
batchImageError(c, err)
return
}
c.JSON(http.StatusOK, got)
}
func (h *BatchImageHandler) Items(c *gin.Context) {
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
limit, _ := strconv.Atoi(c.Query("limit"))
got, err := h.service.ListItems(c.Request.Context(), owner, c.Param("id"), service.BatchImageItemsQuery{
Status: c.Query("status"),
Limit: limit,
Cursor: c.Query("cursor"),
})
if err != nil {
batchImageError(c, err)
return
}
c.JSON(http.StatusOK, got)
}
func (h *BatchImageHandler) Cancel(c *gin.Context) {
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
got, err := h.service.Cancel(c.Request.Context(), owner, c.Param("id"))
if err != nil {
batchImageError(c, err)
return
}
c.JSON(http.StatusOK, got)
}
func (h *BatchImageHandler) ItemContent(c *gin.Context) {
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
imageIndex := 0
if raw := c.Query("image_index"); raw != "" {
parsed, err := strconv.Atoi(raw)
if err != nil {
batchImageError(c, service.ErrBatchImageItemImageIndexOutOfRange)
return
}
imageIndex = parsed
}
stream, err := h.download.OpenItemContent(c.Request.Context(), owner, c.Param("id"), c.Param("custom_id"), imageIndex)
if err != nil {
batchImageError(c, err)
return
}
defer func() { _ = stream.Reader.Close() }()
c.Header("Content-Type", stream.ContentType)
c.Header("Content-Disposition", service.BatchImageContentDispositionAttachment(stream.Filename))
c.Header("Cache-Control", "private, max-age=300")
c.Header("X-Content-Type-Options", "nosniff")
if stream.ContentLength != nil && *stream.ContentLength >= 0 {
c.Header("Content-Length", strconv.FormatInt(*stream.ContentLength, 10))
}
c.Status(http.StatusOK)
if _, err := io.Copy(c.Writer, stream.Reader); err != nil {
return
}
h.markDownloadedBestEffort(c, owner)
}
// markDownloadedBestEffort 在响应体已写出后标记下载状态;
// 此时无法再向客户端返回错误,失败只能记日志(不能静默丢弃)。
func (h *BatchImageHandler) markDownloadedBestEffort(c *gin.Context, owner service.BatchImageOwner) {
if err := h.service.MarkDownloaded(c.Request.Context(), owner, c.Param("id")); err != nil {
logger.L().Warn("batch_image.mark_downloaded_failed",
zap.String("batch_id", c.Param("id")),
zap.Error(err),
)
}
}
func (h *BatchImageHandler) Download(c *gin.Context) {
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
maxItems, _ := strconv.Atoi(c.Query("max_items"))
c.Header("Content-Type", "application/zip")
c.Header("Content-Disposition", service.BatchImageContentDispositionAttachment(c.Param("id")+".zip"))
c.Header("Cache-Control", "private, no-store")
c.Header("X-Content-Type-Options", "nosniff")
result, err := h.download.StreamZip(c.Request.Context(), owner, c.Param("id"), service.BatchImageZipOptions{
Status: c.Query("status"),
MaxItems: maxItems,
IncludeManifest: true,
}, c.Writer)
if err != nil {
if result == nil || !c.Writer.Written() {
batchImageError(c, err)
}
return
}
h.markDownloadedBestEffort(c, owner)
}
func (h *BatchImageHandler) DeleteRecord(c *gin.Context) {
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
if err := h.service.DeleteRecord(c.Request.Context(), owner, c.Param("id")); err != nil {
batchImageError(c, err)
return
}
c.Status(http.StatusNoContent)
}
func (h *BatchImageHandler) DeleteOutputs(c *gin.Context) {
owner, ok := batchImageOwnerFromContext(c)
if !ok {
batchImageError(c, infraerrors.New(http.StatusUnauthorized, "API_KEY_REQUIRED", "API key is required"))
return
}
got, err := h.cleanup.DeleteOutputsForOwner(c.Request.Context(), owner, c.Param("id"))
if err != nil {
batchImageError(c, err)
return
}
c.JSON(http.StatusOK, got)
}
func batchImageOwnerFromContext(c *gin.Context) (service.BatchImageOwner, bool) {
apiKey, ok := middleware.GetAPIKeyFromContext(c)
if !ok || apiKey == nil || apiKey.ID <= 0 || apiKey.UserID <= 0 {
return service.BatchImageOwner{}, false
}
return service.BatchImageOwner{
UserID: apiKey.UserID,
APIKeyID: apiKey.ID,
GroupID: apiKey.GroupID,
}, true
}
func batchImageError(c *gin.Context, err error) {
status := infraerrors.Code(err)
code := infraerrors.Reason(err)
message := infraerrors.Message(err)
if err == nil {
status = http.StatusInternalServerError
code = "INTERNAL_ERROR"
message = "internal error"
}
if status == 0 || (status == http.StatusInternalServerError && strings.TrimSpace(code) == "") {
status = http.StatusInternalServerError
code = "INTERNAL_ERROR"
message = "internal error"
}
if errors.Is(err, service.ErrBatchImageJobNotFound) {
status = http.StatusNotFound
code = "BATCH_IMAGE_NOT_FOUND"
message = "batch image job not found"
}
c.JSON(status, gin.H{
"error": gin.H{
"type": "invalid_request_error",
"code": code,
"message": message,
},
})
}