diff --git a/server/app/goauto/purchase/order_writeback.go b/server/app/goauto/purchase/order_writeback.go index 0b75be3..715a845 100644 --- a/server/app/goauto/purchase/order_writeback.go +++ b/server/app/goauto/purchase/order_writeback.go @@ -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" diff --git a/server/app/goauto/purchase/order_writeback_session_retry_test.go b/server/app/goauto/purchase/order_writeback_session_retry_test.go new file mode 100644 index 0000000..fdb60d7 --- /dev/null +++ b/server/app/goauto/purchase/order_writeback_session_retry_test.go @@ -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) + } +} diff --git a/server/app/goauto/purchase/order_writeback_worker.go b/server/app/goauto/purchase/order_writeback_worker.go index fdc03e5..9ab4b12 100644 --- a/server/app/goauto/purchase/order_writeback_worker.go +++ b/server/app/goauto/purchase/order_writeback_worker.go @@ -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)