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:
QiuSW
2026-09-21 16:05:26 +08:00
co-authored by Claude Opus 5
parent 4261a542ca
commit 01510a85dc
3 changed files with 341 additions and 4 deletions
+12 -2
View File
@@ -108,7 +108,12 @@ func (s *Service) OrderWritebackViews(ctx context.Context, tasks []models.Purcha
v.Reason = r.ErrorMessage
}
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
}
return out, nil
@@ -197,7 +202,12 @@ func (s *Service) RequestOrderWriteback(ctx context.Context, req OrderWritebackR
case row.Status == "pending":
a.Result, a.Reason = "pending", "已加入回填"
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
}
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)
}
// 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) {
cfg := config.ExtConfig.SYB.Resolved()
session, err := sybclient.NewSessionStore(db).Load(ctx, cfg.Username, time.Now())
if err != nil {
return nil, err
}
if session.UserID <= 0 {
return nil, errSessionUserIDMissing
}
c, err := sybclient.New(cfg.BaseURL)
if err != nil {
return nil, err
@@ -37,6 +100,16 @@ func restoreOrderWritebackClient(ctx context.Context, db *gorm.DB) (OrderNumberC
if err = c.ImportCookiesJSON(session.CookiesJSON); err != nil {
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
}
@@ -78,7 +151,10 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
var item models.PurchaseOrderWriteback
recovering := false
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
}
recovering = item.Status == "running"
@@ -90,6 +166,11 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) {
if err != nil {
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 {
updates := map[string]any{"status": status, "error_code": code, "error_message": message, "lease_owner": ""}
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
}
// 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
if err = db.First(&task, item.PurchaseTaskID).Error; err != nil {
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)
cancel()
if err != nil {
return true, finish("failed", "SYB_SESSION_UNAVAILABLE", "SYB会话不可用,请恢复登录后重试")
return true, finishSessionUnavailable(err)
}
read := func() (string, string, error) {
readCtx, stop := context.WithTimeout(ctx, 20*time.Second)