From 4d6be2fa777885868bf64533e796febef30b3982 Mon Sep 17 00:00:00 2001 From: ryan Date: Fri, 28 Aug 2026 17:04:44 +0800 Subject: [PATCH] fix(drivers): propagate app-lifetime context through inproc cron/worker drivers MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit cron 触发与 worker 执行的任务现在继承应用生命周期 context(关闭时级联取消,带超时子上下文), 替代裸 context.Background()。contextcheck 清零。 lint_issues 26→24 --- .auto/log.jsonl | 1 + .../domain/upload/shared/test_helpers.go | 2 +- .../drivers/driver_inproc_cron/plugin.go | 7 ++++--- .../drivers/driver_inproc_cron/scheduler.go | 20 +++++++++++-------- .../drivers/driver_inproc_worker/executor.go | 14 ++++++++++--- .../drivers/driver_inproc_worker/plugin.go | 7 ++++--- 6 files changed, 33 insertions(+), 18 deletions(-) diff --git a/.auto/log.jsonl b/.auto/log.jsonl index 9caea7ee..276d010a 100644 --- a/.auto/log.jsonl +++ b/.auto/log.jsonl @@ -2,3 +2,4 @@ {"ts":"2026-08-28","iter":1,"type":"keep","metrics":{"lint_issues":34,"dup_issues":15,"tests_passed":44},"delta":-11,"description":"goconst(9): taskCategoryUpload/taskQueueDefault consts in upload/task; reuse logDBNameSQLite in admin; mnd(1): defaultCleanupInterval in disk cache; staticcheck SA9004: split typed const group in asynq executor","asi":{"lesson":"golangci v2 defaults cap reporting at 50/3 - uncapped via issues:max-issues-per-linter/max-same-issues=0 (strict-only change); formatter war resolved: make format now = golangci-lint fmt (same gate as code-check), 203-file gofumpt normalization committed as infra"}} {"ts":"2026-08-28","iter":2,"type":"keep","metrics":{"lint_issues":33,"dup_issues":15,"tests_passed":44},"delta":-1,"description":"nilerr real bug: FlushTaskExecutionLog swallowed cache faults (non-miss errors) and silently dropped buffered task logs; now propagates wrapped error, ErrCacheMiss stays a no-op. 3 regression tests (fault / miss / persist+clear) with miniredis + stubDBService(in-memory sqlite)","asi":{"lesson":"FlushTaskExecutionLog callers in executor.go only log errors, so returning wrapped err is safe; tests need SetDBService injection since testhelper.SetupTestEnvironment targets infra/database global not admin dbService"}} {"ts":"2026-08-28","iter":3,"type":"keep","metrics":{"lint_issues":26,"dup_issues":15,"tests_passed":44},"delta":-7,"description":"revive cleanup: symmetric accessors ctx->_, unused params->_, doc comments for SetDBServiceForTest and StorageDriver const block","asi":{"lesson":"callers pass targetCfg positionally so _ at def site is safe; accessor ctx removal considered but _ keeps 24 call sites stable"}} +{"ts":"2026-08-28","iter":4,"type":"keep","metrics":{"lint_issues":24,"dup_issues":15,"tests_passed":44},"delta":-2,"description":"contextcheck: thread app-lifetime ctx through inproc drivers. inproc_cron: Plugin.Start(ctx)->scheduler.Start(ctx)->registerJob(ctx,def); cron closures now Dispatch/log/invoke with ctx (child WithTimeout). inproc_worker: InprocQueue.baseCtx captured at Start(ctx); executeTask uses WithTimeout(q.baseCtx). app.go passes signal/base ctx (only cancelled on shutdown) so this adds graceful-shutdown propagation","asi":{"lesson":"Start ctx is lifetime-scoped (signal.NotifyContext or GoContext), NOT startup-scoped - safe for task dispatch; nil-guard baseCtx in queue.Start allows tests constructing queue directly"}} diff --git a/backend/plugins/domain/upload/shared/test_helpers.go b/backend/plugins/domain/upload/shared/test_helpers.go index d0132b5b..3569302b 100644 --- a/backend/plugins/domain/upload/shared/test_helpers.go +++ b/backend/plugins/domain/upload/shared/test_helpers.go @@ -120,7 +120,7 @@ func NewMockStorageService() *MockStorageService { } // Put uploads an object into mock storage. -func (m *MockStorageService) Put(_ context.Context, key string, body io.Reader, _ int64, contentType string) (contracts.StoragePutResult, error) { +func (m *MockStorageService) Put(_ context.Context, key string, body io.Reader, _ int64, _ string) (contracts.StoragePutResult, error) { m.mu.Lock() defer m.mu.Unlock() data, err := io.ReadAll(body) diff --git a/backend/plugins/drivers/driver_inproc_cron/plugin.go b/backend/plugins/drivers/driver_inproc_cron/plugin.go index 196af401..c70dde53 100644 --- a/backend/plugins/drivers/driver_inproc_cron/plugin.go +++ b/backend/plugins/drivers/driver_inproc_cron/plugin.go @@ -56,8 +56,9 @@ func (p *Plugin) Type() core.DriverType { return core.DriverTypeScheduler } -// Start boots the in-process cron scheduler. -func (p *Plugin) Start(_ context.Context) error { +// Start boots the in-process cron scheduler. ctx is the app-lifetime context; +// cron-dispatched tasks carry it so cancellation propagates on shutdown. +func (p *Plugin) Start(ctx context.Context) error { p.mu.Lock() defer p.mu.Unlock() @@ -66,7 +67,7 @@ func (p *Plugin) Start(_ context.Context) error { p.scheduler = newInprocScheduler(p.coreCtx.Schedules(), p.coreCtx.Tasks(), taskSvc) } - return p.scheduler.Start() + return p.scheduler.Start(ctx) } // Stop terminates the in-process cron scheduler. diff --git a/backend/plugins/drivers/driver_inproc_cron/scheduler.go b/backend/plugins/drivers/driver_inproc_cron/scheduler.go index 133c7b5c..05b75a34 100644 --- a/backend/plugins/drivers/driver_inproc_cron/scheduler.go +++ b/backend/plugins/drivers/driver_inproc_cron/scheduler.go @@ -36,7 +36,11 @@ func newInprocScheduler(scheduleReg extpoints.ScheduleExtension, taskReg extpoin } } -func (s *inprocScheduler) Start() error { +func (s *inprocScheduler) Start(ctx context.Context) error { + if err := ctx.Err(); err != nil { + return err + } + s.mu.Lock() defer s.mu.Unlock() @@ -46,7 +50,7 @@ func (s *inprocScheduler) Start() error { if s.scheduleReg != nil { for _, def := range s.scheduleReg.Schedules() { - s.registerJob(def) + s.registerJob(ctx, def) } } @@ -70,7 +74,7 @@ func (s *inprocScheduler) Stop() { const standardCronFields = 5 -func (s *inprocScheduler) registerJob(def extpoints.ScheduleDefinition) { +func (s *inprocScheduler) registerJob(ctx context.Context, def extpoints.ScheduleDefinition) { spec := def.Spec taskType := def.TaskType @@ -94,8 +98,8 @@ func (s *inprocScheduler) registerJob(def extpoints.ScheduleDefinition) { _, err := s.cronRunner.AddFunc(cronSpec, func() { if s.taskSvc != nil { - if _, dispatchErr := s.taskSvc.Dispatch(context.Background(), taskType, payloadBytes, "inproc_cron"); dispatchErr != nil { - logger.ErrorF(context.Background(), "driver_inproc_cron: dispatch task %q failed: %v", taskType, dispatchErr) + if _, dispatchErr := s.taskSvc.Dispatch(ctx, taskType, payloadBytes, "inproc_cron"); dispatchErr != nil { + logger.ErrorF(ctx, "driver_inproc_cron: dispatch task %q failed: %v", taskType, dispatchErr) } return } @@ -107,15 +111,15 @@ func (s *inprocScheduler) registerJob(def extpoints.ScheduleDefinition) { if timeout <= 0 { timeout = 5 * time.Minute } - ctx, cancel := context.WithTimeout(context.Background(), timeout) + runCtx, cancel := context.WithTimeout(ctx, timeout) defer cancel() - _ = invokeHandler(ctx, td.Handler, payloadBytes) + _ = invokeHandler(runCtx, td.Handler, payloadBytes) }) } } }) if err != nil { - logger.ErrorF(context.Background(), "driver_inproc_cron: invalid cron spec %q for task %q: %v", spec, taskType, err) + logger.ErrorF(ctx, "driver_inproc_cron: invalid cron spec %q for task %q: %v", spec, taskType, err) } } diff --git a/backend/plugins/drivers/driver_inproc_worker/executor.go b/backend/plugins/drivers/driver_inproc_worker/executor.go index 5f38a412..6c47acca 100644 --- a/backend/plugins/drivers/driver_inproc_worker/executor.go +++ b/backend/plugins/drivers/driver_inproc_worker/executor.go @@ -35,6 +35,10 @@ type InprocQueue struct { running atomic.Bool stopCh chan struct{} wg sync.WaitGroup + + // baseCtx is the app-lifetime context captured at Start; task handlers + // derive their timeouts from it so shutdown cancellation propagates. + baseCtx context.Context } // NewInprocQueue creates a new InprocQueue with a given concurrency and queue capacity. @@ -83,12 +87,16 @@ func (q *InprocQueue) Enqueue(taskType string, payload []byte, source string) (s } } -// Start begins processing tasks with the worker pool. -func (q *InprocQueue) Start() { +// Start begins processing tasks with the worker pool. ctx is the app-lifetime +// context used as the parent for per-task execution contexts. +func (q *InprocQueue) Start(ctx context.Context) { if !q.running.CompareAndSwap(false, true) { return } + if q.baseCtx == nil { + q.baseCtx = ctx + } for i := 0; i < q.concurrency; i++ { q.wg.Add(1) util.Go(func() { @@ -149,7 +157,7 @@ func (q *InprocQueue) executeTask(msg TaskMessage) { timeout = 5 * time.Minute } - taskCtx, cancel := context.WithTimeout(context.Background(), timeout) + taskCtx, cancel := context.WithTimeout(q.baseCtx, timeout) defer cancel() err := invokeHandler(taskCtx, td.Handler, msg.Payload) diff --git a/backend/plugins/drivers/driver_inproc_worker/plugin.go b/backend/plugins/drivers/driver_inproc_worker/plugin.go index 1ca152c0..0d9ad818 100644 --- a/backend/plugins/drivers/driver_inproc_worker/plugin.go +++ b/backend/plugins/drivers/driver_inproc_worker/plugin.go @@ -119,8 +119,9 @@ func (p *Plugin) Type() core.DriverType { return core.DriverTypeWorker } -// Start initiates task consumption. -func (p *Plugin) Start(_ context.Context) error { +// Start initiates task consumption. ctx is the app-lifetime context handed to +// task executions so shutdown cancellation propagates. +func (p *Plugin) Start(ctx context.Context) error { if p.queue == nil { p.queue = NewInprocQueue(p.concurrency, p.queueCapacity, p.coreCtx.Tasks()) } @@ -129,7 +130,7 @@ func (p *Plugin) Start(_ context.Context) error { globalQueue = p.queue globalMu.Unlock() - p.queue.Start() + p.queue.Start(ctx) return nil }