diff --git a/backend/internal/service/delivery.go b/backend/internal/service/delivery.go index 7701614..6f2e2ce 100644 --- a/backend/internal/service/delivery.go +++ b/backend/internal/service/delivery.go @@ -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 diff --git a/backend/internal/service/delivery_test.go b/backend/internal/service/delivery_test.go index b423c53..e92e6f0 100644 --- a/backend/internal/service/delivery_test.go +++ b/backend/internal/service/delivery_test.go @@ -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)