修复发货提交卡死恢复

This commit is contained in:
yml2213
2026-08-03 13:24:47 +08:00
parent 6eb05b322c
commit fe55585f4d
2 changed files with 135 additions and 19 deletions
+48 -17
View File
@@ -34,6 +34,8 @@ const (
deliveryStageCreateQueue = "create_queue" // 创建上游 queue 失败
deliveryStagePatchQueue = "patch_queue" // 补充 queue 账号信息失败
deliveryStageCreateOrder = "create_upstream" // 创建上游正式订单失败
deliverySubmissionStaleTimeout = 10 * time.Minute
)
// deliverySubmittedUpstream 判断订单是否已成功提交到上游正式订单。
@@ -315,16 +317,10 @@ func (s *DeliveryService) submit(orderNo, gameAccount, bindUUID string, apiClien
s.markDeliverySubmissionFailed(claimed, apiClientID, deliveryStageCreateQueue, err.Error())
return nil, err
}
// 建 queue 后立即持久化,缩小崩溃窗口:即使后续崩溃,也能识别到已建 queue 的阶段
if _, err := s.fulfillment.UpdateFulfillment(FulfillmentUpdateInput{
MerchantID: order.MerchantID,
APIClientID: apiClientID,
OrderNo: order.OrderNo,
Status: model.OrderStatusDelivering,
ResultData: mergeResultData(claimed.ResultData, map[string]interface{}{
"provider_order_stage": deliveryStageQueueCreated,
"queue_order_id": queueOrderID,
}),
// 建 queue 后持久化阶段信息,不提前改成 delivering
// 上游创建正式订单前会回查开放接口,此时订单必须仍然 can_ship=true。
if err := s.recordDeliverySubmissionStage(claimed, deliveryStageQueueCreated, map[string]interface{}{
"queue_order_id": queueOrderID,
}); err != nil {
return nil, err
}
@@ -401,14 +397,19 @@ func (s *DeliveryService) claimDeliverySubmission(merchantID, apiClientID uint,
if err := validateOrderStatusTransition(&order, model.OrderStatusDelivering, fulfillmentTransitionUpdate); err != nil {
return newDeliveryHTTPError(http.StatusConflict, "订单暂不可发货")
}
now := time.Now()
stage := resultDataString(order.ResultData, "provider_order_stage")
if deliverySubmissionInProgress(stage) && !deliverySubmissionStale(&order, now) {
updated = order
return nil
}
patch := map[string]interface{}{
"source": "delivery_proxy",
"submit_started_at": timeutil.FormatAPITime(time.Now()),
"submit_started_at": timeutil.FormatAPITime(now),
"provider_order_stage": deliveryStageClaimed,
"ship_attempts": resultDataNumber(order.ResultData, "ship_attempts") + 1,
}
if err := tx.Model(&order).Updates(map[string]interface{}{
"order_status": model.OrderStatusDelivering,
"result_data": mergeResultData(order.ResultData, patch),
"failure_reason": "",
}).Error; err != nil {
@@ -422,11 +423,6 @@ func (s *DeliveryService) claimDeliverySubmission(merchantID, apiClientID uint,
if err := writeAudit(tx, &merchantID, nil, apiClientIDPtr, "delivery.submit.claim", "fulfillment_order", order.OrderNo, nil); err != nil {
return err
}
if s.fulfillment.callbacks != nil {
if err := s.fulfillment.callbacks.Enqueue(tx, merchantID, "order.shipping.updated", orderCallbackData(&updated)); err != nil {
return err
}
}
return nil
})
if err != nil {
@@ -435,6 +431,24 @@ func (s *DeliveryService) claimDeliverySubmission(merchantID, apiClientID uint,
return &updated, claimedNow, nil
}
func (s *DeliveryService) recordDeliverySubmissionStage(order *model.FulfillmentOrder, stage string, patch map[string]interface{}) error {
if order == nil {
return errors.New("订单不存在")
}
if patch == nil {
patch = map[string]interface{}{}
}
patch["provider_order_stage"] = stage
resultData := mergeResultData(order.ResultData, patch)
if err := s.fulfillment.db.Model(&model.FulfillmentOrder{}).
Where("id = ?", order.ID).
Update("result_data", resultData).Error; err != nil {
return err
}
order.ResultData = resultData
return nil
}
func (s *DeliveryService) markDeliverySubmissionFailed(order *model.FulfillmentOrder, apiClientID uint, stage, reason string) {
if order == nil {
return
@@ -585,6 +599,23 @@ func upstreamDeliverySucceeded(order map[string]interface{}) bool {
strings.EqualFold(stringFromMap(order, "send_status"), "SUCCESS")
}
func deliverySubmissionInProgress(stage string) bool {
return stage == deliveryStageClaimed || stage == deliveryStageQueueCreated
}
func deliverySubmissionStale(order *model.FulfillmentOrder, now time.Time) bool {
if order == nil {
return false
}
startedAtRaw := resultDataString(order.ResultData, "submit_started_at")
if startedAtRaw != "" {
if startedAt, err := time.Parse(timeutil.APITimeLayout, startedAtRaw); err == nil {
return now.Sub(startedAt) >= deliverySubmissionStaleTimeout
}
}
return !order.UpdatedAt.IsZero() && now.Sub(order.UpdatedAt) >= deliverySubmissionStaleTimeout
}
func optionalUint(value uint) *uint {
if value == 0 {
return nil
+87 -2
View File
@@ -7,8 +7,10 @@ import (
"strings"
"sync/atomic"
"testing"
"time"
"affiliate_dash/internal/model"
"affiliate_dash/internal/pkg/timeutil"
)
func TestDeliveryLinkGenerateAuthorizeAndRevoke(t *testing.T) {
@@ -112,7 +114,9 @@ func TestDeliveryMerchantApiBindAndSubmit(t *testing.T) {
t.Fatalf("create merchant product: %v", err)
}
var bindCalls, boundCalls, queueCreateCalls, queuePatchCalls, upstreamCalls int32
fulfillmentSvc := NewFulfillmentService(db, nil)
var bindCalls, boundCalls, queueCreateCalls, queuePatchCalls, upstreamCalls, upstreamSawCanShip int32
orderNoForUpstreamCheck := ""
var requestLog []string
mockBFF := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != "/sign-proxy" {
@@ -167,6 +171,14 @@ func TestDeliveryMerchantApiBindAndSubmit(t *testing.T) {
case "/public/users/orders-queue":
if payload.Method == http.MethodPost {
atomic.AddInt32(&queueCreateCalls, 1)
if orderNoForUpstreamCheck != "" {
openOrder, err := fulfillmentSvc.QueryOpenOrder(orderNoForUpstreamCheck)
if err == nil && openOrder.Status == model.OrderStatusPaid && openOrder.CanShip {
atomic.StoreInt32(&upstreamSawCanShip, 1)
} else {
atomic.StoreInt32(&upstreamSawCanShip, -1)
}
}
_ = json.NewEncoder(w).Encode(map[string]interface{}{
"code": 0,
"data": map[string]interface{}{
@@ -202,7 +214,6 @@ func TestDeliveryMerchantApiBindAndSubmit(t *testing.T) {
}))
defer mockBFF.Close()
fulfillmentSvc := NewFulfillmentService(db, nil)
deliverySvc := NewDeliveryService(fulfillmentSvc, mockBFF.URL, "dlc", "https://shop.example", "link-secret", 60)
created, err := fulfillmentSvc.CreateOrder(CreateFulfillmentOrderInput{
@@ -214,6 +225,7 @@ func TestDeliveryMerchantApiBindAndSubmit(t *testing.T) {
if err != nil {
t.Fatalf("create order: %v", err)
}
orderNoForUpstreamCheck = created.Order.OrderNo
info, err := deliverySvc.GetMerchantOrder(merchantID, created.Order.OrderNo)
if err != nil {
@@ -248,6 +260,9 @@ func TestDeliveryMerchantApiBindAndSubmit(t *testing.T) {
if atomic.LoadInt32(&queueCreateCalls) != 1 || atomic.LoadInt32(&upstreamCalls) != 1 {
t.Fatalf("first submit should only create one queue, queue=%d upstream=%d requests=%v", queueCreateCalls, upstreamCalls, requestLog)
}
if atomic.LoadInt32(&upstreamSawCanShip) != 1 {
t.Fatalf("upstream re-query during submit should still see paid/can_ship=true")
}
second, err := deliverySvc.SubmitForMerchant(merchantID, 88, created.Order.OrderNo, "4808146277", bind.BindUUID)
if err != nil {
t.Fatalf("repeat submit for merchant: %v", err)
@@ -263,6 +278,76 @@ func TestDeliveryMerchantApiBindAndSubmit(t *testing.T) {
}
}
func TestDeliveryClaimBlocksRecentInProgressStage(t *testing.T) {
recentStartedAt := time.Now().Add(-time.Minute)
merchantID, deliverySvc, order := createDeliveryOrderWithSubmitStage(t, "delivery-claim-recent", deliveryStageClaimed, recentStartedAt, 1)
claimed, claimedNow, err := deliverySvc.claimDeliverySubmission(merchantID, 88, order.OrderNo)
if err != nil {
t.Fatalf("claim recent in-progress order: %v", err)
}
if claimedNow {
t.Fatalf("recent in-progress stage should still block duplicate claim")
}
if got := resultDataNumber(claimed.ResultData, "ship_attempts"); got != 1 {
t.Fatalf("blocked claim should keep attempts, got %d", got)
}
}
func TestDeliveryClaimRecoversStaleInProgressStage(t *testing.T) {
staleStartedAt := time.Now().Add(-deliverySubmissionStaleTimeout - time.Minute)
merchantID, deliverySvc, order := createDeliveryOrderWithSubmitStage(t, "delivery-claim-stale", deliveryStageQueueCreated, staleStartedAt, 1)
claimed, claimedNow, err := deliverySvc.claimDeliverySubmission(merchantID, 88, order.OrderNo)
if err != nil {
t.Fatalf("claim stale in-progress order: %v", err)
}
if !claimedNow {
t.Fatalf("stale in-progress stage should allow a fresh claim")
}
if got := resultDataString(claimed.ResultData, "provider_order_stage"); got != deliveryStageClaimed {
t.Fatalf("stale claim should reset stage to claimed, got %q", got)
}
if got := resultDataNumber(claimed.ResultData, "ship_attempts"); got != 2 {
t.Fatalf("stale claim should count a new attempt, got %d", got)
}
if got := resultDataString(claimed.ResultData, "submit_started_at"); got == "" || got == timeutil.FormatAPITime(staleStartedAt) {
t.Fatalf("stale claim should refresh submit_started_at, got %q", got)
}
}
func createDeliveryOrderWithSubmitStage(t *testing.T, merchantCode, stage string, startedAt time.Time, attempts int64) (uint, *DeliveryService, *model.FulfillmentOrder) {
t.Helper()
db := newServiceTestDB(t)
merchantID, product := seedFulfillmentMerchant(t, db, merchantCode, 5000, 5, 100)
fulfillmentSvc := NewFulfillmentService(db, nil)
created, err := fulfillmentSvc.CreateOrder(CreateFulfillmentOrderInput{
MerchantID: merchantID,
APIClientID: 1,
ClientOrderNo: merchantCode + "-001",
SKU: product.SKU,
})
if err != nil {
t.Fatalf("create order: %v", err)
}
resultData := mergeResultData(created.Order.ResultData, map[string]interface{}{
"provider_order_stage": stage,
"submit_started_at": timeutil.FormatAPITime(startedAt),
"ship_attempts": attempts,
})
if err := db.Model(&model.FulfillmentOrder{}).
Where("id = ?", created.Order.ID).
Update("result_data", resultData).Error; err != nil {
t.Fatalf("seed submit stage: %v", err)
}
var order model.FulfillmentOrder
if err := db.First(&order, created.Order.ID).Error; err != nil {
t.Fatalf("reload order: %v", err)
}
deliverySvc := NewDeliveryService(fulfillmentSvc, "https://bff.example", "dlc", "https://shop.example", "link-secret", 60)
return merchantID, deliverySvc, &order
}
func TestDeliveryResubmitBlockedAfterUpstreamSubmission(t *testing.T) {
db := newServiceTestDB(t)
merchantID, _ := seedFulfillmentMerchant(t, db, "delivery-guard", 5000, 5, 100)