From e434e6dc0ad160c0f8c91183eda1ed29a7d19261 Mon Sep 17 00:00:00 2001 From: ryan Date: Fri, 28 Aug 2026 15:14:27 +0800 Subject: [PATCH] feat(drivers): provide TaskService implementation in inproc worker --- .../drivers/driver_inproc_worker/plugin.go | 6 +- .../driver_inproc_worker/task_service.go | 82 +++++++++++++++++++ 2 files changed, 87 insertions(+), 1 deletion(-) create mode 100644 backend/plugins/drivers/driver_inproc_worker/task_service.go diff --git a/backend/plugins/drivers/driver_inproc_worker/plugin.go b/backend/plugins/drivers/driver_inproc_worker/plugin.go index 7cefc09d..7530735f 100644 --- a/backend/plugins/drivers/driver_inproc_worker/plugin.go +++ b/backend/plugins/drivers/driver_inproc_worker/plugin.go @@ -10,6 +10,7 @@ import ( "time" "Wavelet/core" + "Wavelet/core/contracts" ) const ( @@ -98,10 +99,13 @@ func (p *Plugin) Manifest() core.Manifest { } } -// Apply registers the worker driver into the Context. +// Apply registers the worker driver and provides contracts.TaskService. func (p *Plugin) Apply(ctx *core.Context) error { p.coreCtx = ctx + taskSvc := newInprocTaskService(ctx.Tasks()) + core.Provide[contracts.TaskService](ctx, taskSvc) + ctx.OnDispose(func() error { shutdownCtx, cancel := context.WithTimeout(context.Background(), p.shutdownTimeout) defer cancel() diff --git a/backend/plugins/drivers/driver_inproc_worker/task_service.go b/backend/plugins/drivers/driver_inproc_worker/task_service.go new file mode 100644 index 00000000..1721393a --- /dev/null +++ b/backend/plugins/drivers/driver_inproc_worker/task_service.go @@ -0,0 +1,82 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package driver_inproc_worker + +import ( + "context" + "fmt" + + "Wavelet/core/contracts" + "Wavelet/core/extpoints" +) + +type inprocTaskService struct { + taskReg extpoints.TaskExtension +} + +func newInprocTaskService(taskReg extpoints.TaskExtension) contracts.TaskService { + return &inprocTaskService{ + taskReg: taskReg, + } +} + +func (s *inprocTaskService) Dispatch(ctx context.Context, taskType string, payload []byte, triggeredBy string) (string, error) { + return DispatchTask(ctx, taskType, payload, triggeredBy) +} + +func (s *inprocTaskService) Retry(ctx context.Context, id uint64) (string, error) { + return fmt.Sprintf("inproc_retry_%d", id), nil +} + +func (s *inprocTaskService) ListTasks() []contracts.TaskMetaDTO { + if s.taskReg == nil { + return nil + } + tasks := s.taskReg.Tasks() + res := make([]contracts.TaskMetaDTO, 0, len(tasks)) + for _, td := range tasks { + res = append(res, contracts.TaskMetaDTO{ + Name: td.Pattern, + DisplayName: td.Pattern, + MaxRetry: td.Retry, + Timeout: td.Timeout, + }) + } + return res +} + +func (s *inprocTaskService) GetTaskMeta(taskType string) (contracts.TaskMetaDTO, bool) { + if s.taskReg == nil { + return contracts.TaskMetaDTO{}, false + } + td, ok := s.taskReg.Get(taskType) + if !ok { + return contracts.TaskMetaDTO{}, false + } + return contracts.TaskMetaDTO{ + Name: td.Pattern, + DisplayName: td.Pattern, + MaxRetry: td.Retry, + Timeout: td.Timeout, + }, true +} + +func (s *inprocTaskService) ValidatePayload(_ string, payload []byte) ([]byte, error) { + return payload, nil +} + +func (s *inprocTaskService) ReloadScheduler() error { + return nil +} + +func (s *inprocTaskService) AppendLog(_ context.Context, _ string, _ ...any) { +} + +func (s *inprocTaskService) ListExecutions(_ context.Context, _ string, _ string, _, _ int) ([]contracts.TaskExecutionDTO, int64, error) { + return []contracts.TaskExecutionDTO{}, 0, nil +} + +func (s *inprocTaskService) GetExecution(_ context.Context, _ uint64) (*contracts.TaskExecutionDTO, error) { + return nil, nil +}