fix(purchase): bounded auto-retry for SYB writeback session failures (#330)
SYB order-number writeback silently gave up on session-class failures (SYB_SESSION_UNAVAILABLE), requiring manual resubmit even though the hourly sync job refreshes the session on its own. This adds a bounded, backoff-scheduled auto-retry for that error code only: - restoreOrderWritebackClient now actively probes the cached cookie jar with sybclient.CheckSession after import, so a remotely-expired session is classified as retryable up front instead of surfacing later as SYB_READ_FAILED. It never logs in, never triggers OCR and never deletes the cached session. - The dropped Factory error is now categorized into a safe message (no cookies/tokens) and recorded in error_message. - The worker's claim query additionally picks up failed rows with error_code=SYB_SESSION_UNAVAILABLE once their backoff (lease_expires_at) has elapsed and attempt_count is below maxSessionRetryAttempts=6 (1m/2m/4m/8m/15m growing backoff, chosen to span the hourly sync window); other failure codes are unchanged. - CanSubmit no longer hides manual resubmit during that backoff window; manual resubmit resets attempt_count to 0 and clears the lease so the worker cannot double-claim the same row. Diff is limited to the purchase package; sybimport/sybclient/ sybinnercode are untouched. Tests: go test ./app/goauto/purchase/... (new order_writeback_session_retry_test.go covers backoff scheduling, reclaim timing, max-attempt cutoff, CheckSession invalid/network classification with no session deletion, CanSubmit during backoff, manual resubmit reset, and non-session codes being excluded). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01NTDbDcwbDw1TSAcE6wfh2F
This commit is contained in:
@@ -108,7 +108,12 @@ func (s *Service) OrderWritebackViews(ctx context.Context, tasks []models.Purcha
|
|||||||
v.Reason = r.ErrorMessage
|
v.Reason = r.ErrorMessage
|
||||||
}
|
}
|
||||||
v.CompletedAt = r.CompletedAt
|
v.CompletedAt = r.CompletedAt
|
||||||
v.CanSubmit = v.CanSubmit && (r.Status == "failed" || r.Status == "unknown") && (r.LeaseExpiresAt == nil || !r.LeaseExpiresAt.After(s.Now()))
|
// A session-class failure's LeaseExpiresAt is the automatic-retry backoff
|
||||||
|
// deadline (#330 修订2, order_writeback_worker.go finishSessionUnavailable),
|
||||||
|
// not an in-flight write lease — manual "resubmit" must stay available
|
||||||
|
// during that window instead of being hidden until it expires.
|
||||||
|
sessionBackoff := r.Status == "failed" && r.ErrorCode == "SYB_SESSION_UNAVAILABLE"
|
||||||
|
v.CanSubmit = v.CanSubmit && (r.Status == "failed" || r.Status == "unknown") && (sessionBackoff || r.LeaseExpiresAt == nil || !r.LeaseExpiresAt.After(s.Now()))
|
||||||
out[r.PurchaseTaskID] = v
|
out[r.PurchaseTaskID] = v
|
||||||
}
|
}
|
||||||
return out, nil
|
return out, nil
|
||||||
@@ -197,7 +202,12 @@ func (s *Service) RequestOrderWriteback(ctx context.Context, req OrderWritebackR
|
|||||||
case row.Status == "pending":
|
case row.Status == "pending":
|
||||||
a.Result, a.Reason = "pending", "已加入回填"
|
a.Result, a.Reason = "pending", "已加入回填"
|
||||||
default:
|
default:
|
||||||
if err := tx.Model(&row).Updates(map[string]any{"status": "pending", "write_started": false, "lease_owner": "", "lease_expires_at": nil, "error_code": "", "error_message": ""}).Error; err != nil {
|
// Manual resubmit resets attempt_count to 0 (#330 修订3) so a stale
|
||||||
|
// history of automatic session-class retries never eats into a
|
||||||
|
// fresh manual attempt budget, and clears lease_expires_at so the
|
||||||
|
// worker's bounded auto-retry claim (which requires it non-nil)
|
||||||
|
// does not race a double-claim against this manual pending row.
|
||||||
|
if err := tx.Model(&row).Updates(map[string]any{"status": "pending", "write_started": false, "lease_owner": "", "lease_expires_at": nil, "error_code": "", "error_message": "", "attempt_count": 0}).Error; err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
a.Result, a.Reason = "pending", "已加入回填,将先回读SYB"
|
a.Result, a.Reason = "pending", "已加入回填,将先回读SYB"
|
||||||
|
|||||||
@@ -0,0 +1,229 @@
|
|||||||
|
package purchase
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/google/uuid"
|
||||||
|
"go-admin/app/goauto/models"
|
||||||
|
"go-admin/app/goauto/sybclient"
|
||||||
|
"gorm.io/gorm"
|
||||||
|
)
|
||||||
|
|
||||||
|
// wbFactoryWorker builds a worker whose Factory itself fails, exercising the
|
||||||
|
// restoreOrderWritebackClient failure path (session unavailable) rather than
|
||||||
|
// a remote read/write failure on an otherwise-working client.
|
||||||
|
func wbFactoryWorker(s *Service, err error) *OrderWritebackWorker {
|
||||||
|
return &OrderWritebackWorker{DB: s.DB, Now: s.Now, Factory: func(context.Context, *gorm.DB) (OrderNumberClient, error) { return nil, err }}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestOrderWritebackSessionFailureSchedulesBoundedRetry(t *testing.T) {
|
||||||
|
s, task := orderWritebackFixture(t)
|
||||||
|
if ok, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil || !ok {
|
||||||
|
t.Fatalf("run %v %v", ok, err)
|
||||||
|
}
|
||||||
|
row := loadOrderWriteback(t, s.DB, task.ID)
|
||||||
|
if row.Status != "failed" || row.ErrorCode != "SYB_SESSION_UNAVAILABLE" {
|
||||||
|
t.Fatalf("status=%s code=%s", row.Status, row.ErrorCode)
|
||||||
|
}
|
||||||
|
if row.ErrorMessage == "" || len(row.ErrorMessage) > 300 {
|
||||||
|
t.Fatalf("error message not recorded safely: %q", row.ErrorMessage)
|
||||||
|
}
|
||||||
|
if row.LeaseExpiresAt == nil || !row.LeaseExpiresAt.After(s.Now()) {
|
||||||
|
t.Fatal("no backoff scheduled for first session-class failure")
|
||||||
|
}
|
||||||
|
if row.AttemptCount != 1 {
|
||||||
|
t.Fatalf("attempt_count=%d", row.AttemptCount)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestOrderWritebackSessionFailureNotReclaimedBeforeBackoffExpires(t *testing.T) {
|
||||||
|
s, _ := orderWritebackFixture(t)
|
||||||
|
if _, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
f := &fakeOrderNumberClient{apply: true}
|
||||||
|
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || ok {
|
||||||
|
t.Fatalf("claimed before backoff expired: ok=%v err=%v", ok, err)
|
||||||
|
}
|
||||||
|
if f.writes != 0 {
|
||||||
|
t.Fatal("wrote while still inside backoff window")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// loadOrderWriteback in order_writeback_test.go takes (t, db, id); provide a
|
||||||
|
// small adapter so this file reads naturally when task id is already in hand.
|
||||||
|
func loadOrderWritebackByTask(t *testing.T, s *Service, id uint64) models.PurchaseOrderWriteback {
|
||||||
|
return loadOrderWriteback(t, s.DB, id)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestOrderWritebackSessionFailureReclaimedAfterBackoffExpires(t *testing.T) {
|
||||||
|
s, task := orderWritebackFixture(t)
|
||||||
|
if _, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||||
|
s.Now = func() time.Time { return row.LeaseExpiresAt.Add(time.Second) }
|
||||||
|
f := &fakeOrderNumberClient{apply: true}
|
||||||
|
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || !ok {
|
||||||
|
t.Fatalf("not reclaimed after backoff expired: ok=%v err=%v", ok, err)
|
||||||
|
}
|
||||||
|
if f.writes != 1 {
|
||||||
|
t.Fatal("did not write after successful reclaim")
|
||||||
|
}
|
||||||
|
after := loadOrderWritebackByTask(t, s, task.ID)
|
||||||
|
if after.Status != "succeeded" {
|
||||||
|
t.Fatalf("status=%s", after.Status)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestOrderWritebackSessionFailureStopsRetryingAtMaxAttempts(t *testing.T) {
|
||||||
|
s, task := orderWritebackFixture(t)
|
||||||
|
now := s.Now()
|
||||||
|
for i := 0; i < maxSessionRetryAttempts; i++ {
|
||||||
|
s.Now = func() time.Time { return now }
|
||||||
|
if ok, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil || !ok {
|
||||||
|
t.Fatalf("attempt %d: ok=%v err=%v", i+1, ok, err)
|
||||||
|
}
|
||||||
|
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||||
|
if row.AttemptCount != i+1 {
|
||||||
|
t.Fatalf("attempt %d: attempt_count=%d", i+1, row.AttemptCount)
|
||||||
|
}
|
||||||
|
if i+1 < maxSessionRetryAttempts {
|
||||||
|
if row.LeaseExpiresAt == nil {
|
||||||
|
t.Fatalf("attempt %d: no backoff scheduled", i+1)
|
||||||
|
}
|
||||||
|
now = row.LeaseExpiresAt.Add(time.Second)
|
||||||
|
} else {
|
||||||
|
if row.LeaseExpiresAt != nil {
|
||||||
|
t.Fatal("lease still scheduled at max attempts")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
// One more tick past any plausible backoff: the claim query must exclude
|
||||||
|
// attempt_count >= maxSessionRetryAttempts, so nothing is claimed.
|
||||||
|
s.Now = func() time.Time { return now.Add(24 * time.Hour) }
|
||||||
|
f := &fakeOrderNumberClient{apply: true}
|
||||||
|
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || ok {
|
||||||
|
t.Fatalf("claimed a row past max attempts: ok=%v err=%v", ok, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestOrderWritebackCheckSessionInvalidIsSessionClassAndNeverDeletesSession(t *testing.T) {
|
||||||
|
s, task := orderWritebackFixture(t)
|
||||||
|
seedWritebackSession(t, s)
|
||||||
|
if ok, err := wbFactoryWorker(s, sybclient.ErrSessionInvalid).RunOnce(context.Background()); err != nil || !ok {
|
||||||
|
t.Fatalf("run %v %v", ok, err)
|
||||||
|
}
|
||||||
|
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||||
|
if row.Status != "failed" || row.ErrorCode != "SYB_SESSION_UNAVAILABLE" {
|
||||||
|
t.Fatalf("status=%s code=%s", row.Status, row.ErrorCode)
|
||||||
|
}
|
||||||
|
if row.LeaseExpiresAt == nil {
|
||||||
|
t.Fatal("ErrSessionInvalid was not scheduled for retry")
|
||||||
|
}
|
||||||
|
assertWritebackSessionUntouched(t, s)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestOrderWritebackCheckSessionNetworkErrorIsSessionClassAndNeverDeletesSession(t *testing.T) {
|
||||||
|
s, task := orderWritebackFixture(t)
|
||||||
|
seedWritebackSession(t, s)
|
||||||
|
if ok, err := wbFactoryWorker(s, errors.New("dial tcp: i/o timeout")).RunOnce(context.Background()); err != nil || !ok {
|
||||||
|
t.Fatalf("run %v %v", ok, err)
|
||||||
|
}
|
||||||
|
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||||
|
if row.Status != "failed" || row.ErrorCode != "SYB_SESSION_UNAVAILABLE" {
|
||||||
|
t.Fatalf("status=%s code=%s", row.Status, row.ErrorCode)
|
||||||
|
}
|
||||||
|
if row.LeaseExpiresAt == nil {
|
||||||
|
t.Fatal("network error was not scheduled for retry")
|
||||||
|
}
|
||||||
|
assertWritebackSessionUntouched(t, s)
|
||||||
|
}
|
||||||
|
|
||||||
|
func seedWritebackSession(t *testing.T, s *Service) {
|
||||||
|
t.Helper()
|
||||||
|
store := sybclient.NewSessionStore(s.DB)
|
||||||
|
if err := store.Save(context.Background(), sybclient.Session{
|
||||||
|
Username: "syb-writeback-test", UserID: 555, CookiesJSON: `[{"name":"SESSION","value":"x"}]`, ExpiresAt: s.Now().Add(time.Hour),
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func assertWritebackSessionUntouched(t *testing.T, s *Service) {
|
||||||
|
t.Helper()
|
||||||
|
var count int64
|
||||||
|
if err := s.DB.Model(&models.SYBSession{}).Where("username = ?", "syb-writeback-test").Count(&count).Error; err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if count != 1 {
|
||||||
|
t.Fatal("writeback worker deleted or otherwise removed the cached SYB session")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestOrderWritebackCanSubmitDuringSessionBackoff(t *testing.T) {
|
||||||
|
s, task := orderWritebackFixture(t)
|
||||||
|
if _, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
updated := loadBackfillTask(t, s.DB, task.ID)
|
||||||
|
views, err := s.OrderWritebackViews(context.Background(), []models.PurchaseTask{updated})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
v := views[task.ID]
|
||||||
|
if v.Status != "failed" || !v.CanSubmit {
|
||||||
|
t.Fatalf("expected resubmit available during backoff: %+v", v)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestOrderWritebackManualResubmitResetsAttemptCountAndLease(t *testing.T) {
|
||||||
|
s, task := orderWritebackFixture(t)
|
||||||
|
now := s.Now()
|
||||||
|
for i := 0; i < 3; i++ {
|
||||||
|
s.Now = func() time.Time { return now }
|
||||||
|
if _, err := wbFactoryWorker(s, sybclient.ErrNoSession).RunOnce(context.Background()); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||||
|
now = row.LeaseExpiresAt.Add(time.Second)
|
||||||
|
}
|
||||||
|
before := loadOrderWritebackByTask(t, s, task.ID)
|
||||||
|
if before.AttemptCount != 3 {
|
||||||
|
t.Fatalf("attempt_count=%d", before.AttemptCount)
|
||||||
|
}
|
||||||
|
if _, err := s.RequestOrderWriteback(context.Background(), OrderWritebackRequest{uuid.NewString(), []uint64{task.ID}}); err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
after := loadOrderWritebackByTask(t, s, task.ID)
|
||||||
|
if after.Status != "pending" || after.AttemptCount != 0 || after.LeaseExpiresAt != nil {
|
||||||
|
t.Fatalf("manual resubmit did not reset state: %+v", after)
|
||||||
|
}
|
||||||
|
f := &fakeOrderNumberClient{apply: true}
|
||||||
|
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || !ok || f.writes != 1 {
|
||||||
|
t.Fatalf("worker could not process post-resubmit row: ok=%v err=%v writes=%d", ok, err, f.writes)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestOrderWritebackOtherFailureCodesAreNotAutoRetried(t *testing.T) {
|
||||||
|
s, task := orderWritebackFixture(t)
|
||||||
|
f := &fakeOrderNumberClient{readErr: errors.New("offline")}
|
||||||
|
if ok, err := wbWorker(s, f).RunOnce(context.Background()); err != nil || !ok {
|
||||||
|
t.Fatalf("run %v %v", ok, err)
|
||||||
|
}
|
||||||
|
row := loadOrderWritebackByTask(t, s, task.ID)
|
||||||
|
if row.Status != "failed" || row.ErrorCode != "SYB_READ_FAILED" {
|
||||||
|
t.Fatalf("status=%s code=%s", row.Status, row.ErrorCode)
|
||||||
|
}
|
||||||
|
if row.LeaseExpiresAt != nil {
|
||||||
|
t.Fatal("non-session failure code must not be scheduled for automatic retry")
|
||||||
|
}
|
||||||
|
s.Now = func() time.Time { return row.CreatedAt.Add(24 * time.Hour) }
|
||||||
|
again := &fakeOrderNumberClient{apply: true}
|
||||||
|
if ok, err := wbWorker(s, again).RunOnce(context.Background()); err != nil || ok {
|
||||||
|
t.Fatalf("a non-session failure code was auto-reclaimed: ok=%v err=%v", ok, err)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -24,12 +24,75 @@ type OrderWritebackWorker struct {
|
|||||||
Factory func(context.Context, *gorm.DB) (OrderNumberClient, error)
|
Factory func(context.Context, *gorm.DB) (OrderNumberClient, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// errSessionUserIDMissing marks a cached session whose UserID column is not a
|
||||||
|
// positive SYB account id. SessionStore.Save (session.go) rejects UserID<=0
|
||||||
|
// before it is ever persisted, so this should be unreachable in practice; it
|
||||||
|
// exists so a corrupted/legacy row fails loudly and safely instead of calling
|
||||||
|
// CheckSession with id=0 (#330 修订1).
|
||||||
|
var errSessionUserIDMissing = errors.New("SYB 会话记录缺少有效 user id")
|
||||||
|
|
||||||
|
// Bounded auto-retry for session-class writeback failures (#330). A session
|
||||||
|
// outage self-heals once GoAutoSYBHourlySync refreshes syb_session, so a
|
||||||
|
// short-lived backoff schedule — growing up to ~15 minutes — comfortably
|
||||||
|
// spans that hourly cadence without hammering SYB while the session is down.
|
||||||
|
// maxSessionRetryAttempts caps the automatic attempts so a session that never
|
||||||
|
// recovers still lands back in "failed" for a human instead of retrying
|
||||||
|
// forever.
|
||||||
|
const maxSessionRetryAttempts = 6
|
||||||
|
|
||||||
|
var sessionRetryBackoff = []time.Duration{
|
||||||
|
1 * time.Minute,
|
||||||
|
2 * time.Minute,
|
||||||
|
4 * time.Minute,
|
||||||
|
8 * time.Minute,
|
||||||
|
15 * time.Minute,
|
||||||
|
}
|
||||||
|
|
||||||
|
// sessionRetryDelay returns the backoff before the next automatic attempt,
|
||||||
|
// given the attempt number (1-based, i.e. the count already recorded for the
|
||||||
|
// attempt that just failed).
|
||||||
|
func sessionRetryDelay(attempt int) time.Duration {
|
||||||
|
idx := attempt - 1
|
||||||
|
if idx < 0 {
|
||||||
|
idx = 0
|
||||||
|
}
|
||||||
|
if idx >= len(sessionRetryBackoff) {
|
||||||
|
idx = len(sessionRetryBackoff) - 1
|
||||||
|
}
|
||||||
|
return sessionRetryBackoff[idx]
|
||||||
|
}
|
||||||
|
|
||||||
|
// sessionUnavailableMessage classifies why the cached SYB session could not
|
||||||
|
// be used, without ever including cookies, tokens or other credential
|
||||||
|
// material (#330 修订1点3). The category — not the raw error text — is what
|
||||||
|
// gets persisted to error_message.
|
||||||
|
func sessionUnavailableMessage(err error) string {
|
||||||
|
switch {
|
||||||
|
case errors.Is(err, sybclient.ErrNoSession):
|
||||||
|
return "会话缺失/已过期"
|
||||||
|
case errors.Is(err, errSessionUserIDMissing):
|
||||||
|
return "会话记录异常,缺少 user id"
|
||||||
|
case errors.Is(err, sybclient.ErrSessionInvalid):
|
||||||
|
return "会话校验失效"
|
||||||
|
default:
|
||||||
|
return "会话校验网络错误"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// restoreOrderWritebackClient rebuilds a SYB client from the cached session
|
||||||
|
// only. It never logs in, never triggers OCR and never deletes the cached
|
||||||
|
// session (that stays the exclusive responsibility of sybimport.Connect's
|
||||||
|
// login/refresh path) — it only reports whether the cached cookies still
|
||||||
|
// work, via CheckSession, so the caller can classify the failure (#330).
|
||||||
func restoreOrderWritebackClient(ctx context.Context, db *gorm.DB) (OrderNumberClient, error) {
|
func restoreOrderWritebackClient(ctx context.Context, db *gorm.DB) (OrderNumberClient, error) {
|
||||||
cfg := config.ExtConfig.SYB.Resolved()
|
cfg := config.ExtConfig.SYB.Resolved()
|
||||||
session, err := sybclient.NewSessionStore(db).Load(ctx, cfg.Username, time.Now())
|
session, err := sybclient.NewSessionStore(db).Load(ctx, cfg.Username, time.Now())
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
if session.UserID <= 0 {
|
||||||
|
return nil, errSessionUserIDMissing
|
||||||
|
}
|
||||||
c, err := sybclient.New(cfg.BaseURL)
|
c, err := sybclient.New(cfg.BaseURL)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
@@ -37,6 +100,16 @@ func restoreOrderWritebackClient(ctx context.Context, db *gorm.DB) (OrderNumberC
|
|||||||
if err = c.ImportCookiesJSON(session.CookiesJSON); err != nil {
|
if err = c.ImportCookiesJSON(session.CookiesJSON); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
// Active probe (#330 修订1): without this, a remotely-expired cookie jar
|
||||||
|
// imports cleanly and only fails later inside read(), which would record
|
||||||
|
// it as SYB_READ_FAILED instead of the retryable session-class outcome.
|
||||||
|
// Any error here — ErrSessionInvalid or network/format — is treated as
|
||||||
|
// session-class; only ErrSessionInvalid is a confirmed logout, but a
|
||||||
|
// network/format error is not confirmed-valid either, so it is still
|
||||||
|
// retried rather than attempted as a write.
|
||||||
|
if err = c.CheckSession(ctx, session.UserID, cfg.Username); err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
return c, nil
|
return c, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -78,7 +151,10 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
|
|||||||
var item models.PurchaseOrderWriteback
|
var item models.PurchaseOrderWriteback
|
||||||
recovering := false
|
recovering := false
|
||||||
err := db.Transaction(func(tx *gorm.DB) error {
|
err := db.Transaction(func(tx *gorm.DB) error {
|
||||||
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where("status = ? OR (status = ? AND lease_expires_at <= ?)", "pending", "running", now).Order("id").First(&item).Error; err != nil {
|
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).Where(
|
||||||
|
"status = ? OR (status = ? AND lease_expires_at <= ?) OR (status = ? AND error_code = ? AND lease_expires_at IS NOT NULL AND lease_expires_at <= ? AND attempt_count < ?)",
|
||||||
|
"pending", "running", now, "failed", "SYB_SESSION_UNAVAILABLE", now, maxSessionRetryAttempts,
|
||||||
|
).Order("id").First(&item).Error; err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
recovering = item.Status == "running"
|
recovering = item.Status == "running"
|
||||||
@@ -90,6 +166,11 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return false, err
|
return false, err
|
||||||
}
|
}
|
||||||
|
// tx.Model(&item).Updates used gorm.Expr("attempt_count + 1") above, which
|
||||||
|
// GORM does not read back into the struct; sync it here so downstream
|
||||||
|
// bounded-retry math (finishSessionUnavailable) sees the true post-claim
|
||||||
|
// count instead of being off by one.
|
||||||
|
item.AttemptCount++
|
||||||
finish := func(status, code, message string) error {
|
finish := func(status, code, message string) error {
|
||||||
updates := map[string]any{"status": status, "error_code": code, "error_message": message, "lease_owner": ""}
|
updates := map[string]any{"status": status, "error_code": code, "error_message": message, "lease_owner": ""}
|
||||||
if status != "unknown" {
|
if status != "unknown" {
|
||||||
@@ -100,6 +181,23 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
|
|||||||
}
|
}
|
||||||
return db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error
|
return db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error
|
||||||
}
|
}
|
||||||
|
// finishSessionUnavailable is the bounded-retry counterpart of finish for
|
||||||
|
// SYB_SESSION_UNAVAILABLE: instead of clearing the lease, it schedules the
|
||||||
|
// next automatic attempt (item.AttemptCount was already incremented by the
|
||||||
|
// claim above) until maxSessionRetryAttempts is reached, at which point it
|
||||||
|
// behaves like finish("failed", ...) and stops retrying (#330).
|
||||||
|
finishSessionUnavailable := func(err error) error {
|
||||||
|
updates := map[string]any{
|
||||||
|
"status": "failed", "error_code": "SYB_SESSION_UNAVAILABLE",
|
||||||
|
"error_message": sessionUnavailableMessage(err), "lease_owner": "",
|
||||||
|
}
|
||||||
|
if item.AttemptCount < maxSessionRetryAttempts {
|
||||||
|
updates["lease_expires_at"] = w.Now().Add(sessionRetryDelay(item.AttemptCount))
|
||||||
|
} else {
|
||||||
|
updates["lease_expires_at"] = nil
|
||||||
|
}
|
||||||
|
return db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error
|
||||||
|
}
|
||||||
var task models.PurchaseTask
|
var task models.PurchaseTask
|
||||||
if err = db.First(&task, item.PurchaseTaskID).Error; err != nil {
|
if err = db.First(&task, item.PurchaseTaskID).Error; err != nil {
|
||||||
return true, finish("failed", "TASK_UNAVAILABLE", "采购任务不可用,请人工核对")
|
return true, finish("failed", "TASK_UNAVAILABLE", "采购任务不可用,请人工核对")
|
||||||
@@ -115,7 +213,7 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
|
|||||||
client, err := w.Factory(callCtx, db)
|
client, err := w.Factory(callCtx, db)
|
||||||
cancel()
|
cancel()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return true, finish("failed", "SYB_SESSION_UNAVAILABLE", "SYB会话不可用,请恢复登录后重试")
|
return true, finishSessionUnavailable(err)
|
||||||
}
|
}
|
||||||
read := func() (string, string, error) {
|
read := func() (string, string, error) {
|
||||||
readCtx, stop := context.WithTimeout(ctx, 20*time.Second)
|
readCtx, stop := context.WithTimeout(ctx, 20*time.Second)
|
||||||
|
|||||||
Reference in New Issue
Block a user