From a2cc9a144f0c4293f8e2e1a4f88c0a9fcef4814a Mon Sep 17 00:00:00 2001 From: QiuSW <105186638@qq.com> Date: Mon, 24 Aug 2026 16:02:16 +0800 Subject: [PATCH] =?UTF-8?q?feat(#75):=20=E5=A2=9E=E5=8A=A0=20SYB=20?= =?UTF-8?q?=E6=AF=8F=E5=B0=8F=E6=97=B6=E8=87=AA=E5=8A=A8=E5=90=8C=E6=AD=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/02-architecture-and-code-map.md | 5 +- docs/03-business-rules-and-glossary.md | 13 +- docs/04-local-development-and-verification.md | 20 ++- server/app/goauto/sybimport/import_handler.go | 137 +++++++++++------- server/app/goauto/sybimport/scheduled_job.go | 69 +++++++++ .../goauto/sybimport/scheduled_job_test.go | 60 ++++++++ server/app/jobs/examples.go | 5 +- server/app/jobs/jobbase.go | 8 +- server/app/jobs/service/sys_job.go | 4 + server/app/jobs/type.go | 16 +- server/app/jobs/type_test.go | 29 ++++ .../1786701600000_syb_hourly_sync_job.go | 44 ++++++ .../1786701600000_syb_hourly_sync_job_test.go | 38 +++++ web/src/api/goauto/syb-products.js | 4 +- web/src/views/goauto/syb-products/index.vue | 79 +++------- web/tests/e2e/syb-product-layout.spec.ts | 31 +++- 16 files changed, 429 insertions(+), 133 deletions(-) create mode 100644 server/app/goauto/sybimport/scheduled_job.go create mode 100644 server/app/goauto/sybimport/scheduled_job_test.go create mode 100644 server/app/jobs/type_test.go create mode 100644 server/cmd/migrate/migration/version-local/1786701600000_syb_hourly_sync_job.go create mode 100644 server/cmd/migrate/migration/version-local/1786701600000_syb_hourly_sync_job_test.go diff --git a/docs/02-architecture-and-code-map.md b/docs/02-architecture-and-code-map.md index db7a30a..0410cd2 100644 --- a/docs/02-architecture-and-code-map.md +++ b/docs/02-architecture-and-code-map.md @@ -2,8 +2,8 @@ generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) wiki_page: Architecture-and-Code-Map wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Architecture-and-Code-Map.- -wiki_revision: 5f88de723d9ca9124383cd4abcfa5c27ddcd6d48 -synchronized_at: 2026-08-24T07:03:25Z +wiki_revision: ce37cbe86bad24c38c801955e301ff8e90ec86b5 +synchronized_at: 2026-08-24T07:56:32Z # 架构与代码地图 @@ -120,6 +120,7 @@ Android Portal/Agent | SYB 商品明细增量迁移 | `server/cmd/migrate/migration/version-local/1786700700000_syb_product_import.go` | | SYB 店铺准入、发现与过滤 | `server/app/goauto/sybshop/`、`server/app/goauto/sybimport/`;迁移 `server/cmd/migrate/migration/version-local/1786700900000_syb_shop.go` | | SYB 后台导入任务、进度、单任务互斥与启动恢复 | `server/app/goauto/sybimport/sync_run.go`、`sync_run_handler.go`;表 `syb_sync_run`,迁移 `server/cmd/migrate/migration/version-local/1786701000000_syb_sync_run.go` | +| SYB 每小时自动同步 | `server/app/goauto/sybimport/import_handler.go` 提供手动/定时共用的 `StartImport` 服务;`scheduled_job.go` 通过 go-admin 定时任务调度,固定按 Asia/Shanghai 取今天和昨天;任务注册迁移为 `server/cmd/migrate/migration/version-local/1786701600000_syb_hourly_sync_job.go` | | 采购任务数据与类型化规则契约 | `server/app/goauto/models/purchase.go`、`server/app/goauto/purchasecontract/`;迁移 `server/cmd/migrate/migration/version-local/1786701100000_purchase_contract.go` | | 采购任务单条/批量预检与创建、租约、attempt 幂等、Admin 只读查询和人工处置状态机 | `server/app/goauto/purchase/`;`batch-preview` 通过批量预加载 SYB、蝦皮、PDD 与最新任务执行快速只读预检,不访问 AI;`batch-create` 仍按最新数据逐项完整复核并可在必要时调用 AI;既有追加迁移为 `server/cmd/migrate/migration/version-local/1786701200000_purchase_state_machine.go` | | 采购规格标准化匹配与 AI Provider 设置 | `server/app/goauto/aimatching/`;服务端先做繁简、空白/全半角/大小写和公斤/斤的唯一确定性匹配,再按需调用单一 OpenAI-compatible Provider;`1786701300000_ai_matching_setting.go` 创建设置表,`1786701400000_ai_matching_setting_plain_api_key.go` 将原加密列迁移为内部明文 `api_key`,仅管理员读取 | diff --git a/docs/03-business-rules-and-glossary.md b/docs/03-business-rules-and-glossary.md index 0402d8d..259b90a 100644 --- a/docs/03-business-rules-and-glossary.md +++ b/docs/03-business-rules-and-glossary.md @@ -2,8 +2,8 @@ generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) wiki_page: Business-Rules-and-Glossary wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Business-Rules-and-Glossary.- -wiki_revision: 5adb7b0f05c4562e23fe92fe7ff89625d697c2d4 -synchronized_at: 2026-08-24T07:03:44Z +wiki_revision: bb1dea3110c827ee6516f74eb440417a0551f9b8 +synchronized_at: 2026-08-24T07:56:42Z # 业务规则与术语 @@ -184,6 +184,15 @@ synchronized_at: 2026-08-24T07:03:44Z | spec_source | 采购任务的规格来源标记:`manual_mapping` / `exact_match` / `ai_match`,用于事后批量追溯 | +## SYB 自动同步 + +- 服务端内置 go-admin 定时任务“SYB 每小时自动同步”,默认每小时第 5 分钟执行;同步日期按 Asia/Shanghai 计算,覆盖当天和前一天。 +- 定时同步与保留的手动导入 API 共用同一个导入服务和 `syb_sync_run` 记录;定时任务的操作人显示为“系统定时同步”。 +- 同一时刻只允许一个 SYB 同步任务运行。上一次仍在运行时,本次定时触发直接跳过,不排队、不并发,也不立即重试;等待下一小时再次触发。 +- 定时任务失败时明确记录失败原因,不自动重试。服务重启后由既有启动恢复逻辑处理遗留的运行中记录。 +- SYB 商品页不再提供“导入”和“同步记录”快捷按钮;页面自动显示最近一次同步状态,运行中轮询进度,完成后刷新商品列表。同步记录页面和手动导入 API 继续保留。 + + ## cmautobuy 商品导入 - 这是人工触发的单向离线导入,不是持续双写或双库同步;不随服务启动执行。 diff --git a/docs/04-local-development-and-verification.md b/docs/04-local-development-and-verification.md index 05c487d..fdcf966 100644 --- a/docs/04-local-development-and-verification.md +++ b/docs/04-local-development-and-verification.md @@ -2,8 +2,8 @@ generated: true (请先修改 Gitea Wiki,禁止直接编辑本文件) wiki_page: Local-Development-and-Verification wiki_url: https://git.ilapage.cn/OPC/goauto/wiki/Local-Development-and-Verification.- -wiki_revision: 2075bb8a28fa01dcdbbcf105f12a29afc9b2e6f3 -synchronized_at: 2026-08-24T02:59:16Z +wiki_revision: 465390216d69d28a642854fd946d2d818a834f9c +synchronized_at: 2026-08-24T07:56:50Z # 本地开发与验证 @@ -77,6 +77,22 @@ go run -tags sqlite3 . server -c config/settings.sqlite.yml SQLite 只用于测试;正式运行和最终迁移目标仍为 MySQL 8.4。真实 MySQL 连接串继续通过 `GOAUTO_DB_DSN` 注入。 + +## SYB 定时同步 + +迁移 `1786701600000_syb_hourly_sync_job.go` 会幂等写入 go-admin 的 `sys_job`,调用目标为 `GoAutoSYBHourlySync`,默认 Cron 为 `0 5 * * * *`、状态为启用。已有同调用目标的任务不会被迁移覆盖;管理员可在 go-admin“定时任务”中调整 Cron、启停状态和参数。 + +默认参数为 `{"lookbackDays":2,"timezone":"Asia/Shanghai"}`。当前实现允许回看 1~7 天;默认 2 天即当天和前一天。修改后需要让调度器重新加载任务(通常重启 Admin API)。 + +本地验证不要为了检查迁移而启动 Admin API:迁移本身不会访问 SYB,但启用状态的定时任务会在服务启动并到达下一次调度时间后访问已配置的 SYB。可只执行迁移并再次执行确认幂等: + +```powershell +Set-Location server +go run . migrate -c config/settings.yml +go run . migrate -c config/settings.yml +``` + + ## Web 验证 ```powershell diff --git a/server/app/goauto/sybimport/import_handler.go b/server/app/goauto/sybimport/import_handler.go index c3126c0..3b1bb35 100644 --- a/server/app/goauto/sybimport/import_handler.go +++ b/server/app/goauto/sybimport/import_handler.go @@ -5,6 +5,7 @@ import ( "errors" "net/http" "strconv" + "strings" "sync" "time" @@ -36,11 +37,80 @@ type ImportRequest struct { DateTo string `json:"dateTo"` } +type ImportActor struct { + ID uint64 + Name string +} + +type StartImportResult struct { + RunID uint64 + Status string + Skipped bool +} + +// StartImport is the single entry point shared by the authenticated Admin +// handler and the built-in hourly job. It deliberately owns every preflight +// and both single-flight guards so a scheduled run can never bypass the same +// safety boundary as a manual run. +func StartImport(ctx context.Context, db *gorm.DB, request ImportRequest, actor ImportActor, skipIfRunning bool) (StartImportResult, error) { + if request.DateFrom == "" || request.DateTo == "" { + return StartImportResult{}, invalidRequest("dateFrom 和 dateTo 不能为空,格式为 YYYY-MM-DD") + } + if _, err := splitDateRange(request.DateFrom, request.DateTo); err != nil { + return StartImportResult{}, invalidRequest(err.Error()) + } + + // This check must stay before credentials and Connect: an invalid sync must + // not consume a login attempt or send a captcha to OCR. + enabled, err := sybshop.EnabledNames(ctx, db) + if err != nil { + return StartImportResult{}, &ServiceError{Code: CodeSyncShopPreflightFailed, Message: "服务端处理失败", Cause: err} + } + if len(enabled) == 0 { + return StartImportResult{}, invalidRequest(ErrNoEnabledShop.Error()) + } + + settings := config.ExtConfig.SYB.Resolved() + if !settings.HasCredentials() { + return StartImportResult{}, invalidRequest(credentialHint()) + } + + importGate.Lock() + if importGate.running { + importGate.Unlock() + if skipIfRunning { + return StartImportResult{Skipped: true}, nil + } + return StartImportResult{}, invalidRequest("已有一个导入任务正在执行,请等它结束后再试") + } + run, err := NewSyncRunService(db).Create(ctx, CreateSyncRunInput{ + DateFrom: request.DateFrom, DateTo: request.DateTo, ShopFilterHash: enabledShopHash(enabled), + OperatorID: actor.ID, OperatorName: strings.TrimSpace(actor.Name), + }) + if err != nil { + importGate.Unlock() + if skipIfRunning && isImportAlreadyRunning(err) { + return StartImportResult{Skipped: true}, nil + } + return StartImportResult{}, err + } + importGate.running = true + importGate.Unlock() + + go runImport(db, run.ID, request, settings) + return StartImportResult{RunID: run.ID, Status: run.Status}, nil +} + +func isImportAlreadyRunning(err error) bool { + var serviceErr *ServiceError + return errors.As(err, &serviceErr) && serviceErr.Code == CodeInvalidRequest && strings.Contains(serviceErr.Message, "已有") +} + // Import pulls shipment orders from SYB for a date range and folds them into // the archive. // -// `[必须]` This is the only endpoint that reaches out to SYB. It performs reads -// only — no SYB write endpoint is called from anywhere in GoAuto. +// `[必须]` This remains the authenticated manual entry point. Both it and the +// scheduler delegate to StartImport, and every downstream SYB call is read-only. func (handler Handler) Import(c *gin.Context) { if claimString(jwt.ExtractClaims(c)["rolekey"]) != "admin" { c.JSON(http.StatusForbidden, gin.H{"code": "FORBIDDEN", "message": "只有管理员可以开始导入"}) @@ -51,69 +121,28 @@ func (handler Handler) Import(c *gin.Context) { writeError(c, invalidRequest("请求体必须是合法 JSON,且包含 dateFrom 和 dateTo")) return } - if request.DateFrom == "" || request.DateTo == "" { - writeError(c, invalidRequest("dateFrom 和 dateTo 不能为空,格式为 YYYY-MM-DD")) - return - } - if _, err := splitDateRange(request.DateFrom, request.DateTo); err != nil { - writeError(c, invalidRequest(err.Error())) - return - } - service, ok := handler.service(c) if !ok { return } - // `[必须]` Refuse before Connect: Connect may log in and send a captcha to - // the configured OCR service. With no enabled shop there is no valid sync - // to run, so consuming either remote service would be wasteful and would - // violate #49's pre-flight boundary. Sync checks again to cover a shop being - // disabled between this pre-flight and the actual run. - enabled, err := sybshop.EnabledNames(c.Request.Context(), service.DB) - if err != nil { - handler.logInternalFailure(c, "shop_preflight", err) - writeError(c, &ServiceError{Code: CodeSyncShopPreflightFailed, Message: "服务端处理失败", Cause: err}) - return - } - if len(enabled) == 0 { - writeError(c, invalidRequest(ErrNoEnabledShop.Error())) - return - } - - settings := config.ExtConfig.SYB.Resolved() - if !settings.HasCredentials() { - // Name the file that was actually consulted. The previous wording only - // mentioned environment variables, which sent operators looking in the - // wrong place once config.yaml became the normal way to configure this. - writeError(c, invalidRequest(credentialHint())) - return - } - - importGate.Lock() - if importGate.running { - importGate.Unlock() - writeError(c, &ServiceError{Code: CodeInvalidRequest, Message: "已有一个导入任务正在执行,请等它结束后再试"}) - return - } claims := jwt.ExtractClaims(c) - run, err := NewSyncRunService(service.DB).Create(c.Request.Context(), CreateSyncRunInput{ - DateFrom: request.DateFrom, DateTo: request.DateTo, ShopFilterHash: enabledShopHash(enabled), - OperatorID: claimUint64(claims["identity"]), OperatorName: claimString(claims["nice"]), - }) + result, err := StartImport(c.Request.Context(), service.DB, request, ImportActor{ + ID: claimUint64(claims["identity"]), Name: claimString(claims["nice"]), + }, false) if err != nil { - importGate.Unlock() var serviceErr *ServiceError - if errors.As(err, &serviceErr) && serviceErr.Code == CodeSyncRunCreateFailed { - handler.logInternalFailure(c, "sync_run_create", serviceErr.Cause) + if errors.As(err, &serviceErr) { + switch serviceErr.Code { + case CodeSyncShopPreflightFailed: + handler.logInternalFailure(c, "shop_preflight", serviceErr.Cause) + case CodeSyncRunCreateFailed: + handler.logInternalFailure(c, "sync_run_create", serviceErr.Cause) + } } writeError(c, err) return } - importGate.running = true - importGate.Unlock() - - go runImport(service.DB, run.ID, request, settings) - c.JSON(http.StatusAccepted, gin.H{"code": 200, "data": gin.H{"runId": run.ID, "status": run.Status}}) + c.JSON(http.StatusAccepted, gin.H{"code": 200, "data": gin.H{"runId": result.RunID, "status": result.Status}}) } func runImport(db *gorm.DB, runID uint64, request ImportRequest, settings config.SYB) { diff --git a/server/app/goauto/sybimport/scheduled_job.go b/server/app/goauto/sybimport/scheduled_job.go new file mode 100644 index 0000000..f18665c --- /dev/null +++ b/server/app/goauto/sybimport/scheduled_job.go @@ -0,0 +1,69 @@ +package sybimport + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "strings" + "time" + + "gorm.io/gorm" +) + +const ( + HourlySyncInvokeTarget = "GoAutoSYBHourlySync" + defaultHourlyTimezone = "Asia/Shanghai" + defaultLookbackDays = 2 +) + +type hourlySyncArgs struct { + LookbackDays int `json:"lookbackDays"` + Timezone string `json:"timezone"` +} + +// HourlySyncJob is registered in go-admin's ExecJob map. ExecWithDB is the +// production path; Exec exists only to satisfy legacy jobs.JobExec and fails +// closed if an older caller forgets to provide the current database. +type HourlySyncJob struct{} + +func (HourlySyncJob) Exec(_ interface{}) error { + return errors.New("SYB 定时同步缺少数据库连接") +} + +func (HourlySyncJob) ExecWithDB(db *gorm.DB, arg interface{}) error { + args, err := parseHourlySyncArgs(arg) + if err != nil { + return err + } + location, err := time.LoadLocation(args.Timezone) + if err != nil { + return fmt.Errorf("SYB 定时同步时区无效: %w", err) + } + dateFrom, dateTo := hourlySyncDateRange(time.Now().In(location), args.LookbackDays) + _, err = StartImport(context.Background(), db, ImportRequest{DateFrom: dateFrom, DateTo: dateTo}, ImportActor{Name: "系统定时同步"}, true) + return err +} + +func parseHourlySyncArgs(arg interface{}) (hourlySyncArgs, error) { + result := hourlySyncArgs{LookbackDays: defaultLookbackDays, Timezone: defaultHourlyTimezone} + raw, _ := arg.(string) + if strings.TrimSpace(raw) != "" { + if err := json.Unmarshal([]byte(raw), &result); err != nil { + return hourlySyncArgs{}, fmt.Errorf("SYB 定时同步参数不是合法 JSON: %w", err) + } + } + if result.LookbackDays < 1 || result.LookbackDays > 7 { + return hourlySyncArgs{}, errors.New("SYB 定时同步 lookbackDays 必须在 1 到 7 之间") + } + if strings.TrimSpace(result.Timezone) == "" { + result.Timezone = defaultHourlyTimezone + } + return result, nil +} + +func hourlySyncDateRange(now time.Time, lookbackDays int) (string, string) { + to := now.Format("2006-01-02") + from := now.AddDate(0, 0, -(lookbackDays - 1)).Format("2006-01-02") + return from, to +} diff --git a/server/app/goauto/sybimport/scheduled_job_test.go b/server/app/goauto/sybimport/scheduled_job_test.go new file mode 100644 index 0000000..3998bc2 --- /dev/null +++ b/server/app/goauto/sybimport/scheduled_job_test.go @@ -0,0 +1,60 @@ +package sybimport + +import ( + "context" + "testing" + "time" + + "go-admin/app/goauto/models" + "go-admin/config" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +func TestHourlySyncDateRangeUsesConfiguredLookback(t *testing.T) { + from, to := hourlySyncDateRange(time.Date(2026, 8, 24, 0, 5, 0, 0, time.FixedZone("CST", 8*60*60)), 2) + if from != "2026-08-23" || to != "2026-08-24" { + t.Fatalf("range = %s..%s", from, to) + } +} + +func TestHourlySyncArgsDefaultsAndBounds(t *testing.T) { + args, err := parseHourlySyncArgs("") + if err != nil || args.LookbackDays != 2 || args.Timezone != "Asia/Shanghai" { + t.Fatalf("defaults = %+v, err = %v", args, err) + } + if _, err := parseHourlySyncArgs(`{"lookbackDays":8,"timezone":"Asia/Shanghai"}`); err == nil { + t.Fatal("lookbackDays above safety bound should fail") + } +} + +func TestScheduledImportSkipsExistingRunWithoutStartingNetworkWork(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:scheduled-import-busy?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatalf("open db: %v", err) + } + if err := db.AutoMigrate(&models.SYBShop{}, &models.SYBSyncRun{}); err != nil { + t.Fatalf("migrate: %v", err) + } + if err := db.Create(&models.SYBShop{DisplayName: "测试店铺", NormalizedName: "测试店铺", Enabled: true}).Error; err != nil { + t.Fatalf("create shop: %v", err) + } + if _, err := NewSyncRunService(db).Create(context.Background(), CreateSyncRunInput{DateFrom: "2026-08-23", DateTo: "2026-08-24", ShopFilterHash: "existing"}); err != nil { + t.Fatalf("create active run: %v", err) + } + + originalSYB := config.ExtConfig.SYB + config.ExtConfig.SYB.Username = "test-user" + config.ExtConfig.SYB.Password = "test-password" + t.Cleanup(func() { config.ExtConfig.SYB = originalSYB }) + + result, err := StartImport(context.Background(), db, ImportRequest{DateFrom: "2026-08-23", DateTo: "2026-08-24"}, ImportActor{Name: "系统定时同步"}, true) + if err != nil || !result.Skipped { + t.Fatalf("result = %+v, err = %v", result, err) + } + var count int64 + if err := db.Model(&models.SYBSyncRun{}).Count(&count).Error; err != nil || count != 1 { + t.Fatalf("sync run count = %d, err = %v", count, err) + } +} diff --git a/server/app/jobs/examples.go b/server/app/jobs/examples.go index ed233cc..5937819 100644 --- a/server/app/jobs/examples.go +++ b/server/app/jobs/examples.go @@ -3,6 +3,8 @@ package jobs import ( "fmt" "time" + + "go-admin/app/goauto/sybimport" ) // InitJob @@ -10,7 +12,8 @@ import ( // 字典 key 可以配置到 自动任务 调用目标 中; func InitJob() { jobList = map[string]JobExec{ - "ExamplesOne": ExamplesOne{}, + "ExamplesOne": ExamplesOne{}, + sybimport.HourlySyncInvokeTarget: sybimport.HourlySyncJob{}, // ... } } diff --git a/server/app/jobs/jobbase.go b/server/app/jobs/jobbase.go index 2ec5d94..3b3a613 100644 --- a/server/app/jobs/jobbase.go +++ b/server/app/jobs/jobbase.go @@ -37,6 +37,7 @@ type HttpJob struct { type ExecJob struct { JobCore + DB *gorm.DB } func (e *ExecJob) Run() { @@ -46,10 +47,10 @@ func (e *ExecJob) Run() { log.Warn("[Job] ExecJob Run job nil") return } - err := CallExec(obj.(JobExec), e.Args) + err := CallExecWithDB(obj.(JobExec), e.DB, e.Args) if err != nil { - // 如果失败暂停一段时间重试 - fmt.Println(time.Now().Format(timeFormat), " [ERROR] mission failed! ", err) + log.Errorf("[Job] JobCore %s failed: %v", e.Name, err) + return } // 结束时间 endTime := time.Now() @@ -134,6 +135,7 @@ func setup(key string, db *gorm.DB) { sysJob.EntryId, err = AddJob(crontab, j) } else if jobList[i].JobType == 2 { j := &ExecJob{} + j.DB = db j.InvokeTarget = jobList[i].InvokeTarget j.CronExpression = jobList[i].CronExpression j.JobId = jobList[i].JobId diff --git a/server/app/jobs/service/sys_job.go b/server/app/jobs/service/sys_job.go index a1ab810..356ea00 100644 --- a/server/app/jobs/service/sys_job.go +++ b/server/app/jobs/service/sys_job.go @@ -1,6 +1,7 @@ package service import ( + "context" "errors" "time" @@ -71,6 +72,9 @@ func (e *SysJob) StartJob(c *dto.GeneralGetDto) error { } } else { var j = &jobs.ExecJob{} + // A manually restarted job outlives the HTTP request that created this + // service. Detach it from any request context before storing it in cron. + j.DB = e.Orm.WithContext(context.Background()) j.InvokeTarget = data.InvokeTarget j.CronExpression = data.CronExpression j.JobId = data.JobId diff --git a/server/app/jobs/type.go b/server/app/jobs/type.go index 1b5c30b..9bd9ee3 100644 --- a/server/app/jobs/type.go +++ b/server/app/jobs/type.go @@ -1,6 +1,9 @@ package jobs -import "github.com/robfig/cron/v3" +import ( + "github.com/robfig/cron/v3" + "gorm.io/gorm" +) type Job interface { Run() @@ -11,6 +14,17 @@ type JobExec interface { Exec(arg interface{}) error } +type JobExecWithDB interface { + ExecWithDB(db *gorm.DB, arg interface{}) error +} + func CallExec(e JobExec, arg interface{}) error { return e.Exec(arg) } + +func CallExecWithDB(e JobExec, db *gorm.DB, arg interface{}) error { + if withDB, ok := e.(JobExecWithDB); ok { + return withDB.ExecWithDB(db, arg) + } + return e.Exec(arg) +} diff --git a/server/app/jobs/type_test.go b/server/app/jobs/type_test.go new file mode 100644 index 0000000..483c5c8 --- /dev/null +++ b/server/app/jobs/type_test.go @@ -0,0 +1,29 @@ +package jobs + +import ( + "testing" + + "gorm.io/gorm" +) + +type dbAwareExec struct { + gotDB *gorm.DB + gotArg interface{} +} + +func (*dbAwareExec) Exec(interface{}) error { return nil } +func (job *dbAwareExec) ExecWithDB(db *gorm.DB, arg interface{}) error { + job.gotDB, job.gotArg = db, arg + return nil +} + +func TestCallExecWithDBUsesDatabaseAwarePath(t *testing.T) { + db := &gorm.DB{} + job := &dbAwareExec{} + if err := CallExecWithDB(job, db, "args"); err != nil { + t.Fatalf("call: %v", err) + } + if job.gotDB != db || job.gotArg != "args" { + t.Fatalf("db-aware call = db:%p arg:%v", job.gotDB, job.gotArg) + } +} diff --git a/server/cmd/migrate/migration/version-local/1786701600000_syb_hourly_sync_job.go b/server/cmd/migrate/migration/version-local/1786701600000_syb_hourly_sync_job.go new file mode 100644 index 0000000..856eca6 --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1786701600000_syb_hourly_sync_job.go @@ -0,0 +1,44 @@ +package version_local + +import ( + "errors" + "runtime" + + "go-admin/app/goauto/sybimport" + jobsmodels "go-admin/app/jobs/models" + "go-admin/cmd/migrate/migration" + common "go-admin/common/models" + + "gorm.io/gorm" +) + +func init() { + _, fileName, _, _ := runtime.Caller(0) + migration.Migrate.SetVersion(migration.GetFilename(fileName), migrateSYBHourlySyncJob) +} + +func migrateSYBHourlySyncJob(db *gorm.DB, version string) error { + return db.Transaction(func(tx *gorm.DB) error { + if err := ensureSYBHourlySyncJob(tx); err != nil { + return err + } + return tx.Create(&common.Migration{Version: version}).Error + }) +} + +func ensureSYBHourlySyncJob(db *gorm.DB) error { + var existing jobsmodels.SysJob + err := db.Where("invoke_target = ?", sybimport.HourlySyncInvokeTarget).First(&existing).Error + if err == nil { + return nil + } + if !errors.Is(err, gorm.ErrRecordNotFound) { + return err + } + return db.Create(&jobsmodels.SysJob{ + JobName: "SYB 每小时自动同步", JobGroup: "GoAuto", JobType: 2, + CronExpression: "0 5 * * * *", InvokeTarget: sybimport.HourlySyncInvokeTarget, + Args: `{"lookbackDays":2,"timezone":"Asia/Shanghai"}`, + MisfirePolicy: 1, Concurrent: 1, Status: 2, + }).Error +} diff --git a/server/cmd/migrate/migration/version-local/1786701600000_syb_hourly_sync_job_test.go b/server/cmd/migrate/migration/version-local/1786701600000_syb_hourly_sync_job_test.go new file mode 100644 index 0000000..ad9ab50 --- /dev/null +++ b/server/cmd/migrate/migration/version-local/1786701600000_syb_hourly_sync_job_test.go @@ -0,0 +1,38 @@ +package version_local + +import ( + "testing" + + "go-admin/app/goauto/sybimport" + jobsmodels "go-admin/app/jobs/models" + + "gorm.io/driver/sqlite" + "gorm.io/gorm" +) + +func TestEnsureSYBHourlySyncJobIsIdempotentAndPreservesAdminChanges(t *testing.T) { + db, err := gorm.Open(sqlite.Open("file:syb-hourly-job-migration?mode=memory&cache=shared"), &gorm.Config{}) + if err != nil { + t.Fatalf("open db: %v", err) + } + if err := db.AutoMigrate(&jobsmodels.SysJob{}); err != nil { + t.Fatalf("migrate: %v", err) + } + if err := ensureSYBHourlySyncJob(db); err != nil { + t.Fatalf("first ensure: %v", err) + } + if err := db.Model(&jobsmodels.SysJob{}).Where("invoke_target = ?", sybimport.HourlySyncInvokeTarget). + Updates(map[string]any{"cron_expression": "0 15 * * * *", "status": 1}).Error; err != nil { + t.Fatalf("customize: %v", err) + } + if err := ensureSYBHourlySyncJob(db); err != nil { + t.Fatalf("second ensure: %v", err) + } + var rows []jobsmodels.SysJob + if err := db.Where("invoke_target = ?", sybimport.HourlySyncInvokeTarget).Find(&rows).Error; err != nil { + t.Fatalf("list: %v", err) + } + if len(rows) != 1 || rows[0].CronExpression != "0 15 * * * *" || rows[0].Status != 1 { + t.Fatalf("rows = %+v", rows) + } +} diff --git a/web/src/api/goauto/syb-products.js b/web/src/api/goauto/syb-products.js index 486c873..5c95c81 100644 --- a/web/src/api/goauto/syb-products.js +++ b/web/src/api/goauto/syb-products.js @@ -24,8 +24,8 @@ export function importSybProducts(data) { return request({ url: '/api/admin/v1/syb-products/import', method: 'post', data }) } -export function listSybSyncRuns(params) { - return request({ url: '/api/admin/v1/syb-products/sync-runs', method: 'get', params }) +export function listSybSyncRuns(params, options = {}) { + return request({ url: '/api/admin/v1/syb-products/sync-runs', method: 'get', params, ...options }) } export function getSybSyncRun(runId) { diff --git a/web/src/views/goauto/syb-products/index.vue b/web/src/views/goauto/syb-products/index.vue index b11b496..34d5e7a 100644 --- a/web/src/views/goauto/syb-products/index.vue +++ b/web/src/views/goauto/syb-products/index.vue @@ -3,8 +3,6 @@