mirror of
https://github.com/Wei-Shaw/sub2api.git
synced 2026-09-24 16:05:44 +08:00
修复 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 等)改为记日志
271 lines
8.0 KiB
Go
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,
|
|
},
|
|
})
|
|
}
|