fix(purchase): unify SYB session refresh for writeback (#343)

This commit is contained in:
QiuSW
2026-09-27 22:56:54 +08:00
parent 24100f04d4
commit 934a7beabc
12 changed files with 271 additions and 83 deletions
+7
View File
@@ -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 里被推翻了,
原文和推翻理由都留在这里,方便后来人知道这个决定变过、为什么变:
+1
View File
@@ -41,6 +41,7 @@ func MigratedModels() []any {
&models.SYBSpecAIParseRun{},
&models.SYBSpecAIParseWorkItem{},
&models.SYBSession{},
&models.SYBSessionAuthLease{},
&models.SYBShop{},
&models.SYBProductFilter{},
&models.SYBSyncRun{},
+11
View File
@@ -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
@@ -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()} {
@@ -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) {
@@ -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")
}
})
@@ -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 {
+4
View File
@@ -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
@@ -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)})
}
@@ -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)
}
}
+4 -48
View File
@@ -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
@@ -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
})
}