Files
cmsp/app_upload.go
T

265 lines
8.0 KiB
Go

package main
import (
"context"
"crypto/rand"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"io"
"os"
"time"
"cmsp/internal/erpgo"
"cmsp/internal/store"
)
// UploadVideos 只在使用者确认后逐件执行;未决操作只能查询,不能生成新键重放写入。
func (a *App) UploadVideos(productIDs []string) (uploadErr error) {
defer func() {
if uploadErr != nil && a.log != nil {
a.log.Error("启动视频上传失败:%s", uploadRequestSummary(uploadErr))
}
}()
if a.db == nil {
return fmt.Errorf("数据库未就绪,请查看运行日志")
}
ids := cleanProductIDs(productIDs)
if len(ids) == 0 {
return fmt.Errorf("请先勾选要上传的商品")
}
if !a.uploadMu.TryLock() {
return fmt.Errorf("已有视频上传正在进行")
}
defer a.uploadMu.Unlock()
client, err := erpgo.NewClient(a.cfg.ERPGo, nil)
if err != nil {
return err
}
for _, id := range ids {
a.uploadOneVideo(a.appContext(), client, id, len(ids))
}
return nil
}
func (a *App) uploadOneVideo(ctx context.Context, client *erpgo.Client, id string, count int) {
p, found, err := a.db.GetProduct(id)
if err != nil || !found {
a.log.Error("商品 %s 无法读取上传信息", id)
return
}
fail := func(code string) {
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadFailed, code); err != nil {
a.log.Error("商品 %s 更新上传状态失败:%v", p.ID, err)
}
a.log.Error("商品 %s 视频上传失败:%s", p.ID, code)
}
pending := func(code string) {
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadPending, code); err != nil {
a.log.Error("商品 %s 更新上传状态失败:%v", p.ID, err)
}
a.log.Warn("商品 %s 上传结果待确认:%s;再次操作只查询原任务", p.ID, code)
}
previous, exists, err := a.db.GetUploadOperation(p.ID)
if err != nil {
fail("LOCAL_OPERATION_ERROR")
return
}
if exists && previous.Status == "succeeded" && p.UploadStatus != store.UploadDone {
if previous.VideoURL == "" {
pending("INVALID_RESPONSE")
return
}
if err := a.db.UpsertUploadedVideo(p.ID, previous.FilePath, previous.FileSize, previous.VideoURL, time.Now().Format("2006-01-02 15:04:05")); err != nil {
pending("LOCAL_VIDEO_ERROR")
return
}
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadDone, ""); err != nil {
pending("LOCAL_STATUS_ERROR")
}
return
}
if exists && previous.Status != "succeeded" && previous.Status != "failed" {
a.waitVideoOperation(ctx, client, previous, p, pending, fail)
return
}
current, err := client.GetCurrentVideo(ctx, p.ItemID)
if err != nil {
var videoErr *erpgo.VideoError
if errors.As(err, &videoErr) && videoErr.Status == 502 && videoErr.Code == "HHH_UPSTREAM_ERROR" && videoErr.Stage == "check" {
lastError := "VIDEO_STATE_UNKNOWN"
if videoErr.RequestID != "" {
lastError += " requestId=" + videoErr.RequestID
}
if updateErr := a.db.UpdateProductStatus(p.ID, "", "", store.UploadUnconfirmed, lastError); updateErr != nil {
a.log.Error("商品 %s 保存视频状态未确认结果失败:%v", p.ID, updateErr)
}
a.log.Warn("商品 %s 视频状态无法确认,上传未提交:%s", p.ID, uploadRequestSummary(err))
return
}
fail(videoCode(err))
return
}
if current.Confirmed() && count > 1 {
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadSkippedExisting, ""); err != nil {
a.log.Error("商品 %s 更新上传状态失败:%v", p.ID, err)
}
return
}
localPath, info, found, err := a.findUploadVideo(p)
if err != nil {
fail("LOCAL_VIDEO_ERROR")
return
}
if !found {
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadMissingVideo, ""); err != nil {
a.log.Error("商品 %s 更新上传状态失败:%v", p.ID, err)
}
return
}
reason, err := a.validateUploadVideo(ctx, localPath, info)
if err != nil {
fail("LOCAL_VIDEO_CHECK_FAILED")
return
}
if reason != "" {
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadInvalidVideo, reason); err != nil {
a.log.Error("商品 %s 更新上传状态失败:%v", p.ID, err)
}
return
}
file, err := os.Open(localPath)
if err != nil {
fail("LOCAL_VIDEO_ERROR")
return
}
defer file.Close()
hash := sha256.New()
readSize, err := io.Copy(hash, file)
if err != nil || readSize != info.Size() {
fail("LOCAL_VIDEO_ERROR")
return
}
if _, err := file.Seek(0, io.SeekStart); err != nil {
fail("LOCAL_VIDEO_ERROR")
return
}
keyBytes := make([]byte, 24)
if _, err := rand.Read(keyBytes); err != nil {
fail("LOCAL_OPERATION_ERROR")
return
}
op := store.UploadOperation{ProductID: p.ID, ShopeeID: p.ItemID, IdempotencyKey: hex.EncodeToString(keyBytes), FilePath: localPath, FileSHA256: hex.EncodeToString(hash.Sum(nil)), FileSize: info.Size(), Status: "created", UpdatedAt: time.Now().Format(time.RFC3339)}
if err := a.db.BeginUploadOperation(op); err != nil {
pending("LOCAL_OPERATION_ERROR")
return
}
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadRunning, ""); err != nil {
pending("LOCAL_STATUS_ERROR")
return
}
result, putErr := client.PutVideo(ctx, p.ItemID, op.IdempotencyKey, file, op.FileSize)
if putErr == nil {
op.Status, op.OperationID, op.VideoURL, op.ErrorCode = result.Status, result.OperationID, result.VideoURL, result.ErrorCode
op.UpdatedAt = time.Now().Format(time.RFC3339)
if err := a.db.UpdateUploadOperation(op); err != nil {
pending("LOCAL_OPERATION_ERROR")
return
}
}
// HTTP/网络异常后的结果有歧义;查询原键,不再次 PUT。
if putErr != nil {
a.log.Warn("商品 %s 提交结果未确认:%s", p.ID, videoCode(putErr))
}
a.waitVideoOperation(ctx, client, op, p, pending, fail)
}
func (a *App) waitVideoOperation(ctx context.Context, client *erpgo.Client, op store.UploadOperation, p store.Product, pending, fail func(string)) {
for attempt := 0; attempt < 6; attempt++ {
if attempt > 0 {
select {
case <-ctx.Done():
pending("REQUEST_CANCELLED")
return
case <-time.After(5 * time.Second):
}
}
remote, err := client.GetVideoOperation(ctx, op.ShopeeID, op.IdempotencyKey)
if err != nil {
pending(videoCode(err))
return
}
op.Status, op.OperationID, op.VideoURL, op.ErrorCode = remote.Status, remote.OperationID, remote.VideoURL, remote.ErrorCode
op.UpdatedAt = time.Now().Format(time.RFC3339)
if err := a.db.UpdateUploadOperation(op); err != nil {
pending("LOCAL_OPERATION_ERROR")
return
}
switch remote.Status {
case "succeeded":
if remote.VideoURL == "" {
pending("INVALID_RESPONSE")
return
}
if err := a.db.UpsertUploadedVideo(p.ID, op.FilePath, op.FileSize, remote.VideoURL, time.Now().Format("2006-01-02 15:04:05")); err != nil {
pending("LOCAL_VIDEO_ERROR")
return
}
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadDone, ""); err != nil {
pending("LOCAL_STATUS_ERROR")
return
}
a.log.Success("商品 %s 的 erpgo 视频操作已确认成功", p.ID)
return
case "failed":
code := remote.ErrorCode
if code == "" {
code = "REMOTE_VIDEO_FAILED"
}
if !safeUploadCode(code) {
code = "REMOTE_VIDEO_FAILED"
}
fail(code)
return
case "unknown":
pending("VIDEO_RESULT_UNKNOWN")
return
}
}
pending("VIDEO_PROCESSING")
}
func videoCode(err error) string {
var videoErr *erpgo.VideoError
if errors.As(err, &videoErr) && safeUploadCode(videoErr.Code) {
return videoErr.Code
}
return "VIDEO_REQUEST_FAILED"
}
// 上传入口错误只记稳定字段;尤其不要把原始上游响应、请求地址或凭据写入日志。
func uploadRequestSummary(err error) string {
var videoErr *erpgo.VideoError
if errors.As(err, &videoErr) {
return fmt.Sprintf("%s (HTTP %d, requestId=%s)", videoCode(err), videoErr.Status, videoErr.RequestID)
}
var bridgeErr *erpgo.Error
if errors.As(err, &bridgeErr) {
return erpgo.LogSummary(err)
}
return "LOCAL_UPLOAD_ERROR"
}
func safeUploadCode(value string) bool {
if len(value) == 0 || len(value) > 64 {
return false
}
for _, r := range value {
if !(r >= 'A' && r <= 'Z' || r >= '0' && r <= '9' || r == '_') {
return false
}
}
return true
}