mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-04 23:16:37 +08:00
修改路径
This commit is contained in:
@@ -21,11 +21,11 @@ import (
|
||||
"fmt"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db/idgen"
|
||||
"github.com/Rain-kl/Wavelet/internal/logger"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/otel_trace"
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/linux-do/credit/internal/db/idgen"
|
||||
"github.com/linux-do/credit/internal/logger"
|
||||
"github.com/linux-do/credit/internal/model"
|
||||
"github.com/linux-do/credit/internal/otel_trace"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
"go.opentelemetry.io/otel/codes"
|
||||
"go.opentelemetry.io/otel/trace"
|
||||
|
||||
@@ -22,9 +22,9 @@ import (
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/linux-do/credit/internal/model"
|
||||
"github.com/linux-do/credit/internal/testhelper"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
)
|
||||
@@ -64,10 +64,19 @@ func failHandler() *mockHandler {
|
||||
const testTaskType = "test:mock_task"
|
||||
|
||||
func setupTest(t *testing.T) func() {
|
||||
_, _, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
_, mr, cleanup := testhelper.SetupTestEnvironment(t)
|
||||
AsynqClient = asynq.NewClient(asynq.RedisClientOpt{
|
||||
Addr: mr.Addr(),
|
||||
})
|
||||
// 注册测试用 handler
|
||||
RegisterHandler(testTaskType, successHandler())
|
||||
return cleanup
|
||||
return func() {
|
||||
if AsynqClient != nil {
|
||||
AsynqClient.Close()
|
||||
AsynqClient = nil
|
||||
}
|
||||
cleanup()
|
||||
}
|
||||
}
|
||||
|
||||
func TestRegisterAndGetHandler(t *testing.T) {
|
||||
@@ -156,11 +165,6 @@ func TestProcessTaskSuccess(t *testing.T) {
|
||||
err := model.CreateTaskExecution(ctx, execution)
|
||||
require.NoError(t, err)
|
||||
|
||||
// 模拟 Asynq Task
|
||||
asynqTask := asynq.NewTask(testTaskType, nil)
|
||||
// 使用 ResultWriter 设置 TaskID
|
||||
rw := asynq.NewResultWriter()
|
||||
rw.SetTaskID("process_success_001")
|
||||
// 通过 asynq 的 Task 不能直接设置 taskID,ProcessTask 通过 t.ResultWriter().TaskID() 获取
|
||||
// 但 asynq.Task 在没有经过 asynq server 的情况下 ResultWriter 可能为 nil
|
||||
// 我们需要在 ProcessTask 内部改用 taskID 注入的方式测试
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
/*
|
||||
Copyright 2026 linux.do
|
||||
|
||||
Licensed under the Apache License, Version 2.0 (the "License");
|
||||
you may not use this file except in compliance with the License.
|
||||
You may obtain a copy of the License at
|
||||
|
||||
http://www.apache.org/licenses/LICENSE-2.0
|
||||
|
||||
Unless required by applicable law or agreed to in writing, software
|
||||
distributed under the License is distributed on an "AS IS" BASIS,
|
||||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
See the License for the specific language governing permissions and
|
||||
limitations under the License.
|
||||
*/
|
||||
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/upload"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/user"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
)
|
||||
|
||||
// Register registers all built-in task handlers.
|
||||
func Register() {
|
||||
task.RegisterHandler(task.CleanupUnusedUploadsTask, &upload.CleanupUnusedUploadsHandler{})
|
||||
task.RegisterHandler(task.SendEmailTask, &user.SendEmailHandler{})
|
||||
}
|
||||
@@ -21,8 +21,8 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/linux-do/credit/internal/config"
|
||||
"github.com/linux-do/credit/internal/task"
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
|
||||
"github.com/hibiken/asynq"
|
||||
)
|
||||
|
||||
@@ -17,8 +17,8 @@ limitations under the License.
|
||||
package task
|
||||
|
||||
import (
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/linux-do/credit/internal/config"
|
||||
)
|
||||
|
||||
// RedisOpt asynq Redis 连接配置(兼容 Standalone/Sentinel/Cluster)
|
||||
|
||||
@@ -19,17 +19,15 @@ package worker
|
||||
import (
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/task"
|
||||
taskhandlers "github.com/Rain-kl/Wavelet/internal/task/handlers"
|
||||
"github.com/hibiken/asynq"
|
||||
"github.com/linux-do/credit/internal/apps/upload"
|
||||
"github.com/linux-do/credit/internal/apps/user"
|
||||
"github.com/linux-do/credit/internal/config"
|
||||
"github.com/linux-do/credit/internal/task"
|
||||
)
|
||||
|
||||
func init() {
|
||||
// 注册所有任务处理器
|
||||
task.RegisterHandler(task.CleanupUnusedUploadsTask, &upload.CleanupUnusedUploadsHandler{})
|
||||
task.RegisterHandler(task.SendEmailTask, &user.SendEmailHandler{})
|
||||
taskhandlers.Register()
|
||||
}
|
||||
|
||||
// StartWorker 启动任务处理服务器
|
||||
|
||||
Reference in New Issue
Block a user