From 934a7beabcd4c538efe61151a4341401b11cda00 Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Sun, 27 Sep 2026 22:56:54 +0800 Subject: [PATCH] fix(purchase): unify SYB session refresh for writeback (#343) --- docs/12-syb-erp-interface.md | 7 + server/app/goauto/migrations/migrate.go | 1 + server/app/goauto/models/schema.go | 11 ++ .../goauto/purchase/order_backfill_test.go | 2 +- server/app/goauto/purchase/order_writeback.go | 5 +- .../goauto/purchase/order_writeback_test.go | 6 +- .../goauto/purchase/order_writeback_worker.go | 58 ++++--- server/app/goauto/sybclient/session.go | 4 + .../app/goauto/sybclient/session_manager.go | 147 ++++++++++++++++++ .../goauto/sybclient/session_manager_test.go | 35 +++++ server/app/goauto/sybimport/sync.go | 52 +------ .../1789801400000_syb_session_auth_lease.go | 26 ++++ 12 files changed, 271 insertions(+), 83 deletions(-) create mode 100644 server/app/goauto/sybclient/session_manager.go create mode 100644 server/app/goauto/sybclient/session_manager_test.go create mode 100644 server/cmd/migrate/migration/version-local/1789801400000_syb_session_auth_lease.go diff --git a/docs/12-syb-erp-interface.md b/docs/12-syb-erp-interface.md index b32725e..b1a7077 100644 --- a/docs/12-syb-erp-interface.md +++ b/docs/12-syb-erp-interface.md @@ -525,6 +525,13 @@ POST /am/stock/detail/updateDetailCode?t=0&id={stockID}&detailId={detailID}&code 它是反复启动的一次性脚本,进程间要传会话;GoAuto 服务端是常驻进程,没有这个需求, 持久化只为重启后免登录,复用现有数据库即可。表见 `syb_session`。 +`[必须]` 同步和采购订单回填共用统一会话获取器:先校验 `syb_session`,只有明确 +会话失效或不存在时才进入登录。登录由 `syb_session_auth_lease` 的单账号租约串行化, +租约内其他任务等待新会话,不重复请求验证码;登录成功后原回填任务重新执行并回读确认。 +网络超时、5xx 或 OCR 不可用不得清除有效会话,必须保留结构化失败状态供人工重试。 +回填终态同时同步到 `purchase_task.writeback_status/writeback_at`,页面不得继续显示旧的 +`not_selected`。 + `[决定已变更]` ~~不引入 OCR 服务。~~ 这条判断在上游工单 #47 里被推翻了, 原文和推翻理由都留在这里,方便后来人知道这个决定变过、为什么变: diff --git a/server/app/goauto/migrations/migrate.go b/server/app/goauto/migrations/migrate.go index 443a094..62c583b 100644 --- a/server/app/goauto/migrations/migrate.go +++ b/server/app/goauto/migrations/migrate.go @@ -41,6 +41,7 @@ func MigratedModels() []any { &models.SYBSpecAIParseRun{}, &models.SYBSpecAIParseWorkItem{}, &models.SYBSession{}, + &models.SYBSessionAuthLease{}, &models.SYBShop{}, &models.SYBProductFilter{}, &models.SYBSyncRun{}, diff --git a/server/app/goauto/models/schema.go b/server/app/goauto/models/schema.go index fa881f4..10861b8 100644 --- a/server/app/goauto/models/schema.go +++ b/server/app/goauto/models/schema.go @@ -504,6 +504,17 @@ type SYBSession struct { UpdatedAt time.Time `json:"updatedAt"` } +// SYBSessionAuthLease serializes re-authentication across sync and writeback +// workers. Cookies and passwords never live in this table. +type SYBSessionAuthLease struct { + ID uint64 `gorm:"primaryKey"` + Username string `gorm:"size:128;not null;uniqueIndex"` + Owner string `gorm:"size:36;not null"` + ExpiresAt time.Time `gorm:"not null"` +} + +func (SYBSessionAuthLease) TableName() string { return "syb_session_auth_lease" } + // SYBProduct is one SYB (顺云宝 ERP) shipment detail line: one order can carry // several Shopee product lines, and the same Shopee product can appear more // than once within one order at different colors/sizes/quantities — each such diff --git a/server/app/goauto/purchase/order_backfill_test.go b/server/app/goauto/purchase/order_backfill_test.go index b22b405..2eb0135 100644 --- a/server/app/goauto/purchase/order_backfill_test.go +++ b/server/app/goauto/purchase/order_backfill_test.go @@ -96,7 +96,7 @@ func TestOrderBackfillMixedBatchAndReplay(t *testing.T) { if saved.StatusVersion != a.StatusVersion+1 || saved.ErrorCode != nil || saved.ErrorMessage != nil || saved.DeviceRunSlot != nil || saved.AccountRunSlot != nil || saved.ActiveSlot == nil || saved.Status != models.PurchaseTaskStatusOrderCreated { t.Fatalf("state metadata: %+v", saved) } - if saved.PaymentReviewStatus != a.PaymentReviewStatus || saved.LogisticsStatus != a.LogisticsStatus || saved.WritebackStatus != a.WritebackStatus || saved.RuleSnapshot != a.RuleSnapshot { + if saved.PaymentReviewStatus != a.PaymentReviewStatus || saved.LogisticsStatus != a.LogisticsStatus || saved.WritebackStatus != models.PurchaseWritebackStatusPending || saved.RuleSnapshot != a.RuleSnapshot { t.Fatal("unrelated business facts changed") } for _, replayID := range []string{rid, uuid.NewString()} { diff --git a/server/app/goauto/purchase/order_writeback.go b/server/app/goauto/purchase/order_writeback.go index 715a845..7319ed5 100644 --- a/server/app/goauto/purchase/order_writeback.go +++ b/server/app/goauto/purchase/order_writeback.go @@ -60,7 +60,10 @@ func ensureOrderWriteback(tx *gorm.DB, t models.PurchaseTask) error { return nil } row := models.PurchaseOrderWriteback{PurchaseTaskID: t.ID, StockID: int64(syb.StockID), DetailID: int64(syb.DetailID), OrderNo: *t.PDDOrderNo, Status: "pending"} - return tx.Session(&gorm.Session{Logger: logger.Default.LogMode(logger.Silent)}).Clauses(clause.OnConflict{DoNothing: true}).Create(&row).Error + if err := tx.Session(&gorm.Session{Logger: logger.Default.LogMode(logger.Silent)}).Clauses(clause.OnConflict{DoNothing: true}).Create(&row).Error; err != nil { + return err + } + return tx.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ? AND writeback_status = ?", t.ID, models.PurchaseWritebackStatusNotSelected).Update("writeback_status", models.PurchaseWritebackStatusPending).Error } func (s *Service) OrderWritebackViews(ctx context.Context, tasks []models.PurchaseTask) (map[uint64]OrderWritebackView, error) { diff --git a/server/app/goauto/purchase/order_writeback_test.go b/server/app/goauto/purchase/order_writeback_test.go index a2dcd40..7d43cd1 100644 --- a/server/app/goauto/purchase/order_writeback_test.go +++ b/server/app/goauto/purchase/order_writeback_test.go @@ -116,7 +116,11 @@ func TestOrderWritebackRemoteOutcomes(t *testing.T) { t.Fatal("automatically repeated write") } after := loadBackfillTask(t, s.DB, task.ID) - if after.PaymentReviewStatus != task.PaymentReviewStatus || after.WritebackStatus != task.WritebackStatus || after.StatusVersion != task.StatusVersion { + wantTaskWriteback := models.PurchaseWritebackStatusFailed + if tc.want == "succeeded" { + wantTaskWriteback = models.PurchaseWritebackStatusSucceeded + } + if after.PaymentReviewStatus != task.PaymentReviewStatus || after.WritebackStatus != wantTaskWriteback || after.StatusVersion != task.StatusVersion { t.Fatal("changed purchase/payment/logistics facts") } }) diff --git a/server/app/goauto/purchase/order_writeback_worker.go b/server/app/goauto/purchase/order_writeback_worker.go index b2421e4..b8ab969 100644 --- a/server/app/goauto/purchase/order_writeback_worker.go +++ b/server/app/goauto/purchase/order_writeback_worker.go @@ -84,38 +84,16 @@ func sessionUnavailableMessage(err error) string { return "SYB会话不可用(" + category + "),将自动重试;如持续失败请恢复登录后重试" } -// 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). +// restoreOrderWritebackClient uses the same session acquisition path as SYB +// sync. A valid cached session is reused; an explicitly invalid session is +// refreshed under the shared database auth lease so concurrent workers do not +// request multiple captcha codes. 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 - } - 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 + return sybclient.AcquireSession(ctx, db, sybclient.LoginConfig{ + BaseURL: cfg.BaseURL, Username: cfg.Username, Password: cfg.Password, + OcrURL: cfg.OcrURL, OcrMaxAttempts: cfg.OcrMaxAttempts, + }) } // One short-lived claim at a time across processes. No business writes occur @@ -184,7 +162,20 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) { if status == "succeeded" { updates["completed_at"] = w.Now() } - return db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error + if err := db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error; err != nil { + return err + } + // Keep the admin task state aligned with the authoritative writeback + // record. Conflicts and unknown outcomes remain failed until a human + // resolves them; they must never appear as successful. + taskStatus := models.PurchaseWritebackStatusFailed + if status == "succeeded" { + taskStatus = models.PurchaseWritebackStatusSucceeded + } + return db.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ?", item.PurchaseTaskID).Updates(map[string]any{ + "writeback_status": taskStatus, + "writeback_at": gorm.Expr("CASE WHEN ? = 'succeeded' THEN ? ELSE writeback_at END", status, w.Now()), + }).Error } // finishSessionUnavailable is the bounded-retry counterpart of finish for // SYB_SESSION_UNAVAILABLE: instead of clearing the lease, it schedules the @@ -201,7 +192,10 @@ func (w *OrderWritebackWorker) RunOnce(ctx context.Context) (bool, error) { } 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 + if err := db.Model(&models.PurchaseOrderWriteback{}).Where("id = ? AND status = 'running' AND lease_owner = ?", item.ID, owner).Updates(updates).Error; err != nil { + return err + } + return db.Session(&gorm.Session{SkipHooks: true}).Model(&models.PurchaseTask{}).Where("id = ?", item.PurchaseTaskID).Update("writeback_status", models.PurchaseWritebackStatusFailed).Error } var task models.PurchaseTask if err = db.First(&task, item.PurchaseTaskID).Error; err != nil { diff --git a/server/app/goauto/sybclient/session.go b/server/app/goauto/sybclient/session.go index 9d8d8a4..1848481 100644 --- a/server/app/goauto/sybclient/session.go +++ b/server/app/goauto/sybclient/session.go @@ -30,6 +30,10 @@ type SessionStore struct{ db *gorm.DB } func NewSessionStore(db *gorm.DB) *SessionStore { return &SessionStore{db: db} } +// DB exposes the store connection to the shared session manager; callers do +// not receive any session data through this accessor. +func (s *SessionStore) DB() *gorm.DB { return s.db } + // Session is one cached SYB login. type Session struct { Username string diff --git a/server/app/goauto/sybclient/session_manager.go b/server/app/goauto/sybclient/session_manager.go new file mode 100644 index 0000000..9ceb4a8 --- /dev/null +++ b/server/app/goauto/sybclient/session_manager.go @@ -0,0 +1,147 @@ +package sybclient + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/google/uuid" + "go-admin/app/goauto/models" + "gorm.io/gorm" + "gorm.io/gorm/clause" +) + +// LoginConfig contains the non-secret connection settings and credentials +// needed to acquire a SYB session. Password is used only during LoginWithOCR. +type LoginConfig struct { + BaseURL string + Username string + Password string + OcrURL string + OcrMaxAttempts int +} + +// AcquireSession reuses a valid cached session and performs at most one +// re-login per account at a time. Callers waiting for another process to +// refresh the session never request another captcha. +func AcquireSession(ctx context.Context, db *gorm.DB, cfg LoginConfig) (*Client, error) { + if cfg.Username == "" { + return nil, fmt.Errorf("顺云宝账号未配置") + } + client, err := New(cfg.BaseURL) + if err != nil { + return nil, err + } + store := NewSessionStore(db) + if cached, loadErr := store.Load(ctx, cfg.Username, time.Now()); loadErr == nil { + if importErr := client.ImportCookiesJSON(cached.CookiesJSON); importErr != nil { + return nil, importErr + } + checkErr := client.CheckSession(ctx, cached.UserID, cfg.Username) + if checkErr == nil { + return client, nil + } + if !errors.Is(checkErr, ErrSessionInvalid) { + return nil, checkErr + } + if cfg.Password == "" { + return nil, checkErr + } + } else if !errors.Is(loadErr, ErrNoSession) { + return nil, loadErr + } else if cfg.Password == "" { + return nil, loadErr + } + + owner := uuid.NewString() + if claimed, err := claimAuthLease(ctx, db, cfg.Username, owner, 2*time.Minute); err != nil { + return nil, err + } else if claimed { + defer releaseAuthLease(context.Background(), db, cfg.Username, owner) + // Another worker may have completed login between our first probe and + // acquiring the lease; always re-check before requesting a captcha. + if c, ok := validCachedSession(ctx, client, store, cfg.Username); ok { + return c, nil + } + if err := refreshSession(ctx, db, client, store, cfg); err != nil { + return nil, err + } + return validSessionOrError(ctx, client, store, cfg.Username) + } + + // A different worker owns the login lease. Wait for its session, bounded by + // the caller's context; do not trigger a second login. + for { + if c, ok := validCachedSession(ctx, client, store, cfg.Username); ok { + return c, nil + } + if err := ctx.Err(); err != nil { + return nil, fmt.Errorf("等待顺云宝会话刷新超时: %w", err) + } + select { + case <-time.After(500 * time.Millisecond): + case <-ctx.Done(): + return nil, ctx.Err() + } + } +} + +func validSessionOrError(ctx context.Context, client *Client, store *SessionStore, username string) (*Client, error) { + s, err := store.Load(ctx, username, time.Now()) + if err != nil { + return nil, err + } + if err := client.ImportCookiesJSON(s.CookiesJSON); err != nil { + return nil, err + } + if err := client.CheckSession(ctx, s.UserID, username); err != nil { + return nil, err + } + return client, nil +} + +func validCachedSession(ctx context.Context, client *Client, store *SessionStore, username string) (*Client, bool) { + s, err := store.Load(ctx, username, time.Now()) + if err != nil { + return nil, false + } + if err := client.ImportCookiesJSON(s.CookiesJSON); err != nil { + return nil, false + } + if err := client.CheckSession(ctx, s.UserID, username); err != nil { + // Network failures are not treated as logout, but they also do not + // authorize a write; the caller will retry through the normal worker. + return nil, false + } + return client, true +} + +func refreshSession(ctx context.Context, db *gorm.DB, client *Client, store *SessionStore, cfg LoginConfig) error { + ocr, err := NewOcrClient(cfg.OcrURL, 0) + if err != nil { + return fmt.Errorf("顺云宝验证码识别服务不可用: %w", err) + } + result, reason := client.LoginWithOCR(ctx, ocr, cfg.Username, cfg.Password, cfg.OcrMaxAttempts) + if result == nil { + return fmt.Errorf("顺云宝自动登录失败,需要手工输入验证码: %s", reason) + } + cookies, err := client.ExportCookiesJSON() + if err != nil { + return err + } + return store.Save(ctx, Session{Username: cfg.Username, UserID: result.User.ID, CookiesJSON: cookies, ExpiresAt: result.ExpiresAt}) +} + +func claimAuthLease(ctx context.Context, db *gorm.DB, username, owner string, ttl time.Duration) (bool, error) { + now := time.Now().UTC() + if err := db.WithContext(ctx).Clauses(clause.OnConflict{DoNothing: true}).Create(&models.SYBSessionAuthLease{ID: 1, Username: username, Owner: "", ExpiresAt: now.Add(-time.Second)}).Error; err != nil { + return false, err + } + r := db.WithContext(ctx).Model(&models.SYBSessionAuthLease{}).Where("id = 1 AND username = ? AND (expires_at <= ? OR owner = '')", username, now).Updates(map[string]any{"owner": owner, "expires_at": now.Add(ttl)}) + return r.RowsAffected == 1, r.Error +} + +func releaseAuthLease(ctx context.Context, db *gorm.DB, username, owner string) { + _ = db.WithContext(ctx).Model(&models.SYBSessionAuthLease{}).Where("id = 1 AND username = ? AND owner = ?", username, owner).Updates(map[string]any{"owner": "", "expires_at": time.Now().UTC().Add(-time.Second)}) +} diff --git a/server/app/goauto/sybclient/session_manager_test.go b/server/app/goauto/sybclient/session_manager_test.go new file mode 100644 index 0000000..4e8e37d --- /dev/null +++ b/server/app/goauto/sybclient/session_manager_test.go @@ -0,0 +1,35 @@ +package sybclient + +import ( + "context" + "testing" + "time" + + "go-admin/app/goauto/models" + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +func TestAuthLeaseAllowsOnlyOneRefreshOwner(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:syb-auth-lease?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatal(err) + } + if err := db.AutoMigrate(&models.SYBSessionAuthLease{}); err != nil { + t.Fatal(err) + } + ctx := context.Background() + first, err := claimAuthLease(ctx, db, "operator", "first", time.Minute) + if err != nil || !first { + t.Fatalf("first owner should claim: %v %v", first, err) + } + second, err := claimAuthLease(ctx, db, "operator", "second", time.Minute) + if err != nil || second { + t.Fatalf("second owner must wait: %v %v", second, err) + } + releaseAuthLease(ctx, db, "operator", "first") + second, err = claimAuthLease(ctx, db, "operator", "second", time.Minute) + if err != nil || !second { + t.Fatalf("lease should be reusable: %v %v", second, err) + } +} diff --git a/server/app/goauto/sybimport/sync.go b/server/app/goauto/sybimport/sync.go index 3cd5a7d..a3a7c6b 100644 --- a/server/app/goauto/sybimport/sync.go +++ b/server/app/goauto/sybimport/sync.go @@ -480,54 +480,10 @@ func splitDateRange(dateFrom, dateTo string) ([]string, error) { // blip as a logout would trigger needless logins and could throw away a // perfectly good session. func Connect(ctx context.Context, store *sybclient.SessionStore, cfg ConnectConfig) (*sybclient.Client, error) { - if cfg.Username == "" || cfg.Password == "" { - return nil, errors.New("顺云宝账号或密码未配置,请设置 GOAUTO_SYB_USERNAME 和 GOAUTO_SYB_PASSWORD") - } - client, err := sybclient.New(cfg.BaseURL) - if err != nil { - return nil, err - } - - cached, err := store.Load(ctx, cfg.Username, time.Now()) - switch { - case err == nil: - if importErr := client.ImportCookiesJSON(cached.CookiesJSON); importErr == nil { - if checkErr := client.CheckSession(ctx, cached.UserID, cfg.Username); checkErr == nil { - return client, nil - } else if errors.Is(checkErr, sybclient.ErrSessionInvalid) { - if delErr := store.Delete(ctx, cfg.Username); delErr != nil { - return nil, delErr - } - } - } - case errors.Is(err, sybclient.ErrNoSession): - // Nothing cached; fall through to a fresh login. - default: - return nil, err - } - - if cfg.OcrURL == "" { - return nil, errors.New("顺云宝会话已失效,且未配置验证码识别服务;请配置 extend.syb.ocrurl 或改用手工登录") - } - ocr, err := sybclient.NewOcrClient(cfg.OcrURL, 0) - if err != nil { - return nil, err - } - result, reason := client.LoginWithOCR(ctx, ocr, cfg.Username, cfg.Password, cfg.OcrMaxAttempts) - if result == nil { - return nil, fmt.Errorf("顺云宝自动登录失败,需要手工输入验证码: %s", reason) - } - jar, err := client.ExportCookiesJSON() - if err != nil { - return nil, err - } - if err := store.Save(ctx, sybclient.Session{ - Username: cfg.Username, UserID: result.User.ID, - CookiesJSON: jar, ExpiresAt: result.ExpiresAt, - }); err != nil { - return nil, err - } - return client, nil + return sybclient.AcquireSession(ctx, store.DB(), sybclient.LoginConfig{ + BaseURL: cfg.BaseURL, Username: cfg.Username, Password: cfg.Password, + OcrURL: cfg.OcrURL, OcrMaxAttempts: cfg.OcrMaxAttempts, + }) } // ConnectConfig carries the login settings. The password is passed through and diff --git a/server/cmd/migrate/migration/version-local/1789801400000_syb_session_auth_lease.go b/server/cmd/migrate/migration/version-local/1789801400000_syb_session_auth_lease.go new file mode 100644 index 0000000..aa5128b --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1789801400000_syb_session_auth_lease.go @@ -0,0 +1,26 @@ +package version_local + +import ( + "runtime" + + "go-admin/app/goauto/models" + "go-admin/cmd/migrate/migration" + common "go-admin/common/models" + "gorm.io/gorm" +) + +func init() { + _, file, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(file), migrateSYBSessionAuthLease) +} + +// The lease is additive and contains no credentials. It serializes automatic +// SYB re-authentication across sync and purchase writeback workers. +func migrateSYBSessionAuthLease(db *gorm.DB, version string) error { + return db.Transaction(func(tx *gorm.DB) error { + if err := tx.AutoMigrate(&models.SYBSessionAuthLease{}); err != nil { + return err + } + return tx.Create(&common.Migration{Version: version}).Error + }) +}