376 lines
11 KiB
Go
376 lines
11 KiB
Go
package main
|
||
|
||
import (
|
||
"context"
|
||
"crypto/rand"
|
||
"crypto/sha256"
|
||
"encoding/hex"
|
||
"errors"
|
||
"fmt"
|
||
"io"
|
||
"os"
|
||
"sync"
|
||
"time"
|
||
|
||
"cmsp/internal/erpgo"
|
||
"cmsp/internal/store"
|
||
|
||
"github.com/wailsapp/wails/v2/pkg/runtime"
|
||
)
|
||
|
||
// UploadProgress 仅描述当前这轮批量操作;商品上传状态仍以 SQLite 为准。
|
||
type UploadProgress struct {
|
||
BatchID uint64 `json:"batchId"`
|
||
Total int `json:"total"`
|
||
Done int `json:"done"`
|
||
Uploaded int `json:"uploaded"`
|
||
Skipped int `json:"skipped"`
|
||
Unavailable int `json:"unavailable"`
|
||
Failed int `json:"failed"`
|
||
Unconfirmed int `json:"unconfirmed"`
|
||
Current string `json:"current"`
|
||
ElapsedSec int `json:"elapsedSec"`
|
||
State string `json:"state"`
|
||
}
|
||
|
||
func (a *App) GetUploadProgress() UploadProgress {
|
||
a.uploadProgressMu.Lock()
|
||
defer a.uploadProgressMu.Unlock()
|
||
progress := a.uploadProgress
|
||
if progress.State == "running" {
|
||
progress.ElapsedSec = int(time.Since(a.uploadStarted).Seconds())
|
||
}
|
||
return progress
|
||
}
|
||
|
||
func (a *App) emitUploadProgress(progress UploadProgress) {
|
||
if a.ctx != nil {
|
||
runtime.EventsEmit(a.ctx, "upload:progress", progress)
|
||
}
|
||
}
|
||
|
||
func (a *App) startUploadProgress(total int) {
|
||
a.uploadProgressMu.Lock()
|
||
progress := UploadProgress{BatchID: a.uploadProgress.BatchID + 1, Total: total, State: "running"}
|
||
a.uploadProgress = progress
|
||
a.uploadStarted = time.Now()
|
||
a.uploadProgressMu.Unlock()
|
||
a.emitUploadProgress(progress)
|
||
}
|
||
|
||
func (a *App) finishUploadProduct(id string) {
|
||
status := store.UploadFailed
|
||
current := id
|
||
if product, found, err := a.db.GetProduct(id); err == nil && found {
|
||
status = product.UploadStatus
|
||
if product.ItemID != "" {
|
||
current = product.ItemID
|
||
}
|
||
}
|
||
a.uploadProgressMu.Lock()
|
||
a.uploadProgress.Done++
|
||
a.uploadProgress.Current = current
|
||
a.uploadProgress.ElapsedSec = int(time.Since(a.uploadStarted).Seconds())
|
||
switch status {
|
||
case store.UploadDone:
|
||
a.uploadProgress.Uploaded++
|
||
case store.UploadSkippedExisting:
|
||
a.uploadProgress.Skipped++
|
||
case store.UploadMissingVideo, store.UploadInvalidVideo:
|
||
a.uploadProgress.Unavailable++
|
||
case store.UploadUnconfirmed, store.UploadRunning:
|
||
a.uploadProgress.Unconfirmed++
|
||
default:
|
||
a.uploadProgress.Failed++
|
||
}
|
||
progress := a.uploadProgress
|
||
a.uploadProgressMu.Unlock()
|
||
a.emitUploadProgress(progress)
|
||
}
|
||
|
||
func (a *App) finishUploadProgress() {
|
||
a.uploadProgressMu.Lock()
|
||
a.uploadProgress.State = "completed"
|
||
a.uploadProgress.ElapsedSec = int(time.Since(a.uploadStarted).Seconds())
|
||
progress := a.uploadProgress
|
||
a.uploadProgressMu.Unlock()
|
||
a.emitUploadProgress(progress)
|
||
}
|
||
|
||
// 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
|
||
}
|
||
concurrency := a.cfg.Upload.Concurrency
|
||
if concurrency < 1 || concurrency > 4 {
|
||
return fmt.Errorf("上传并发数不合法,允许范围 1—4")
|
||
}
|
||
if concurrency > len(ids) {
|
||
concurrency = len(ids)
|
||
}
|
||
a.startUploadProgress(len(ids))
|
||
ctx := a.appContext()
|
||
jobs := make(chan string)
|
||
var workers sync.WaitGroup
|
||
for i := 0; i < concurrency; i++ {
|
||
workers.Add(1)
|
||
go func() {
|
||
defer workers.Done()
|
||
for id := range jobs {
|
||
a.uploadOneVideo(ctx, client, id, len(ids))
|
||
a.finishUploadProduct(id)
|
||
}
|
||
}()
|
||
}
|
||
for _, id := range ids {
|
||
jobs <- id
|
||
}
|
||
close(jobs)
|
||
workers.Wait()
|
||
a.finishUploadProgress()
|
||
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
|
||
}
|
||
logID := fmt.Sprintf("%s(货憨憨 ID %s)", p.ItemID, p.ID)
|
||
fail := func(code string) {
|
||
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadFailed, code); err != nil {
|
||
a.log.Error("商品 %s 更新上传状态失败:%v", logID, err)
|
||
}
|
||
a.log.Error("商品 %s 视频上传失败:%s", logID, code)
|
||
}
|
||
pending := func(code string) {
|
||
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadUnconfirmed, code); err != nil {
|
||
a.log.Error("商品 %s 更新上传状态失败:%v", logID, err)
|
||
}
|
||
a.log.Warn("商品 %s 上传结果待确认:%s;再次操作只查询原任务", logID, 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, logID, 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", logID, updateErr)
|
||
}
|
||
a.log.Warn("商品 %s 视频状态无法确认,上传未提交:%s", logID, uploadRequestSummary(err))
|
||
return
|
||
}
|
||
fail(videoCode(err))
|
||
return
|
||
}
|
||
if current.Processing() {
|
||
pending("VIDEO_ALREADY_PROCESSING")
|
||
return
|
||
}
|
||
if current.Confirmed() && count > 1 {
|
||
if err := a.db.UpdateProductStatus(p.ID, "", "", store.UploadSkippedExisting, ""); err != nil {
|
||
a.log.Error("商品 %s 更新上传状态失败:%v", logID, 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", logID, 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", logID, 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", logID, videoCode(putErr))
|
||
}
|
||
a.waitVideoOperation(ctx, client, op, p, logID, pending, fail)
|
||
}
|
||
|
||
func (a *App) waitVideoOperation(ctx context.Context, client *erpgo.Client, op store.UploadOperation, p store.Product, logID string, 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 视频操作已确认成功", logID)
|
||
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
|
||
}
|