diff --git a/backend/cmd/app.go b/backend/cmd/app.go index 9b761016..dae69d79 100644 --- a/backend/cmd/app.go +++ b/backend/cmd/app.go @@ -7,7 +7,9 @@ import ( "context" "database/sql" "fmt" + "io/fs" "log" + "path/filepath" "time" "github.com/pressly/goose/v3" @@ -67,13 +69,13 @@ func newWaveletApp(profile core.Profile) *core.App { ) } - // 3. Register all 8 domain business plugins + // 3. Register all 8 domain business plugins (admin first to ensure schema and base config tables exist) app.Use( - auth.New(), + admin.New(), user.New(), + auth.New(), message_gateway.New(), risk_control.New(), - admin.New(), upload.New(), cap.New(), system.New(), @@ -179,8 +181,7 @@ func (s *sharedStore) ListMigrations(ctx context.Context, db goosedb.DBTxConn) ( var results []*goosedb.ListMigrationsResult for rows.Next() { var r goosedb.ListMigrationsResult - r.IsApplied = true - if err := rows.Scan(&r.Version); err != nil { + if err := rows.Scan(&r.Version, &r.IsApplied); err != nil { return nil, err } results = append(results, &r) @@ -240,7 +241,8 @@ func (e *gooseEngine) Migrate(ctx *core.Context, entries []core.MigrationEntry) dialect: dialectStr, } - provider, err := goose.NewProvider(dialect, sqlDB, entry.FS, goose.WithStore(store)) + migrationFS := findMigrationFS(entry.FS, dialect) + provider, err := goose.NewProvider(goose.DialectCustom, sqlDB, migrationFS, goose.WithStore(store)) if err != nil { return fmt.Errorf("migration %s: create provider: %w", entry.PluginID, err) } @@ -267,3 +269,54 @@ func gooseDialect() goose.Dialect { } return goose.DialectPostgres } + +func findMigrationFS(rootFS fs.FS, dialect goose.Dialect) fs.FS { + dialectDir := "postgres" + if dialect == goose.DialectSQLite3 { + dialectDir = "sqlite" + } + + // 1. Direct search for dialect folder (e.g., "sqlite", "migrations/sqlite", "logstore/migrations/sqlite") + for _, subDir := range []string{ + dialectDir, + "migrations/" + dialectDir, + "logstore/migrations/" + dialectDir, + } { + if sub, err := fs.Sub(rootFS, subDir); err == nil { + if matches, err := fs.Glob(sub, "*.sql"); err == nil && len(matches) > 0 { + return sub + } + } + } + + // 2. Recursive walk to find a directory named dialectDir with *.sql files + var foundDir string + _ = fs.WalkDir(rootFS, ".", func(path string, d fs.DirEntry, err error) error { + if err == nil && d.IsDir() && filepath.Base(path) == dialectDir { + if sub, subErr := fs.Sub(rootFS, path); subErr == nil { + if matches, globErr := fs.Glob(sub, "*.sql"); globErr == nil && len(matches) > 0 { + foundDir = path + return fs.SkipAll + } + } + } + return nil + }) + + if foundDir != "" && foundDir != "." { + if sub, err := fs.Sub(rootFS, foundDir); err == nil { + return sub + } + } + + // 3. Fallback to generic migrations / root if dialect specific is not present + for _, subDir := range []string{"migrations", "logstore/migrations"} { + if sub, err := fs.Sub(rootFS, subDir); err == nil { + if matches, err := fs.Glob(sub, "*.sql"); err == nil && len(matches) > 0 { + return sub + } + } + } + + return rootFS +} diff --git a/backend/cmd/redis_plug_test.go b/backend/cmd/redis_plug_test.go new file mode 100644 index 00000000..184fa433 --- /dev/null +++ b/backend/cmd/redis_plug_test.go @@ -0,0 +1,213 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package cmd + +import ( + "context" + "fmt" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "Wavelet/core" + "Wavelet/core/contracts" + "Wavelet/core/extpoints" + "Wavelet/pkg/config" +) + +func TestRedisPluggability_Simulation(t *testing.T) { + origRedisEnabled := config.Config.Redis.Enabled + defer func() { config.Config.Redis.Enabled = origRedisEnabled }() + + // ══════════════════════════════════════════════════════════════════════════ + // 场景 1: 拔出 Redis (Zero-Redis Monolith 模式) + // ══════════════════════════════════════════════════════════════════════════ + t.Run("Scenario_Unplugged_ZeroRedis_Mode", func(t *testing.T) { + config.Config.Redis.Enabled = false + + app := newWaveletApp(core.ProfileAll) + require.NotNil(t, app) + + // 1. 验证插件挂载形态 + _, ok := app.Plugin("cache_memory") + assert.True(t, ok, "cache_memory 必须挂载") + _, ok = app.Plugin("driver_inproc_worker") + assert.True(t, ok, "driver_inproc_worker 必须挂载") + _, ok = app.Plugin("driver_inproc_cron") + assert.True(t, ok, "driver_inproc_cron 必须挂载") + + _, ok = app.Plugin("cache") + assert.False(t, ok, "分布式 cache 不得挂载") + _, ok = app.Plugin("driver_asynq_worker") + assert.False(t, ok, "asynq_worker 不得挂载") + _, ok = app.Plugin("driver_asynq_cron") + assert.False(t, ok, "asynq_cron 不得挂载") + + // 2. 注册测试任务与 Cron 定时 + var taskExecuted atomic.Int32 + var cronExecuted atomic.Int32 + + app.Context().Tasks().Register("test:inproc_task", func(ctx context.Context, payload []byte) error { + if string(payload) == "payload_unplugged" { + taskExecuted.Add(1) + } + return nil + }, extpoints.WithTaskTimeout(3*time.Second)) + + app.Context().Schedules().RegisterCron("* * * * * *", "test:inproc_cron", []byte("cron_ping")) + app.Context().Tasks().Register("test:inproc_cron", func(ctx context.Context, payload []byte) error { + if string(payload) == "cron_ping" { + cronExecuted.Add(1) + } + return nil + }) + + // 3. 启动应用 + bootCtx, bootCancel := context.WithTimeout(context.Background(), 5*time.Second) + defer bootCancel() + require.NoError(t, app.Start(bootCtx)) + + // 4. 验证 CacheService 操作 + cacheSvc, err := core.Inject[contracts.CacheService](app.Context()) + require.NoError(t, err) + require.NotNil(t, cacheSvc) + + reqCtx := context.Background() + require.NoError(t, cacheSvc.Set(reqCtx, "unplugged_key", "value_123", time.Minute)) + var val string + require.NoError(t, cacheSvc.Get(reqCtx, "unplugged_key", &val)) + assert.Equal(t, "value_123", val) + + // 5. 验证异步 Worker 任务分发与执行 + taskSvc, err := core.Inject[contracts.TaskService](app.Context()) + require.NoError(t, err) + require.NotNil(t, taskSvc) + + taskID, err := taskSvc.Dispatch(reqCtx, "test:inproc_task", []byte("payload_unplugged"), "unit_test") + require.NoError(t, err) + assert.NotEmpty(t, taskID) + + require.Eventually(t, func() bool { + return taskExecuted.Load() >= 1 + }, 3*time.Second, 50*time.Millisecond, "内存 Worker 应在进程内顺利执行任务") + + // 6. 验证 Cron 定时触发 + require.Eventually(t, func() bool { + return cronExecuted.Load() >= 1 + }, 3*time.Second, 100*time.Millisecond, "内存 Cron 驱动应成功触发定时任务") + + // 7. 优雅关闭 + stopCtx, stopCancel := context.WithTimeout(context.Background(), 5*time.Second) + defer stopCancel() + require.NoError(t, app.Stop(stopCtx)) + }) + + // ══════════════════════════════════════════════════════════════════════════ + // 场景 2: 插入 Redis (Distributed Cluster 模式) + // ══════════════════════════════════════════════════════════════════════════ + t.Run("Scenario_Plugged_Redis_Mode", func(t *testing.T) { + config.Config.Redis.Enabled = true + + app := newWaveletApp(core.ProfileAll) + require.NotNil(t, app) + + // 1. 验证插件挂载形态 + _, ok := app.Plugin("cache") + assert.True(t, ok, "分布式 cache 必须挂载") + _, ok = app.Plugin("driver_asynq_worker") + assert.True(t, ok, "driver_asynq_worker 必须挂载") + _, ok = app.Plugin("driver_asynq_cron") + assert.True(t, ok, "driver_asynq_cron 必须挂载") + + _, ok = app.Plugin("cache_memory") + assert.False(t, ok, "纯内存 cache 不得挂载") + _, ok = app.Plugin("driver_inproc_worker") + assert.False(t, ok, "inproc_worker 不得挂载") + _, ok = app.Plugin("driver_inproc_cron") + assert.False(t, ok, "inproc_cron 不得挂载") + + // 2. 注册测试任务 + var asynqTaskExecuted atomic.Int32 + app.Context().Tasks().Register("test:asynq_task", func(ctx context.Context, payload []byte) error { + if string(payload) == "payload_plugged" { + asynqTaskExecuted.Add(1) + } + return nil + }, extpoints.WithTaskTimeout(3*time.Second)) + + // 3. 启动应用 (连接真实运行中的 Redis 6379) + bootCtx, bootCancel := context.WithTimeout(context.Background(), 5*time.Second) + defer bootCancel() + require.NoError(t, app.Start(bootCtx)) + + // 4. 验证 CacheService 操作 (L1 RAM + L2 Redis) + cacheSvc, err := core.Inject[contracts.CacheService](app.Context()) + require.NoError(t, err) + require.NotNil(t, cacheSvc) + + reqCtx := context.Background() + testKey := fmt.Sprintf("plugged_key_%d", time.Now().UnixNano()) + require.NoError(t, cacheSvc.Set(reqCtx, testKey, "value_redis_cluster", time.Minute)) + + var val string + require.NoError(t, cacheSvc.Get(reqCtx, testKey, &val)) + assert.Equal(t, "value_redis_cluster", val) + + // 验证失效广播与删除 + require.NoError(t, cacheSvc.Delete(reqCtx, testKey)) + var valAfterDelete string + err = cacheSvc.Get(reqCtx, testKey, &valAfterDelete) + assert.ErrorIs(t, err, contracts.ErrCacheMiss) + + // 5. 验证 Asynq Worker 任务分发与消费 + taskSvc, err := core.Inject[contracts.TaskService](app.Context()) + require.NoError(t, err) + require.NotNil(t, taskSvc) + + taskID, err := taskSvc.Dispatch(reqCtx, "test:asynq_task", []byte("payload_plugged"), "unit_test") + require.NoError(t, err) + assert.NotEmpty(t, taskID) + + require.Eventually(t, func() bool { + return asynqTaskExecuted.Load() >= 1 + }, 5*time.Second, 100*time.Millisecond, "Asynq Worker 应从 Redis 队列中成功消费并执行任务") + + // 6. 优雅关闭 + stopCtx, stopCancel := context.WithTimeout(context.Background(), 5*time.Second) + defer stopCancel() + require.NoError(t, app.Stop(stopCtx)) + }) + + // ══════════════════════════════════════════════════════════════════════════ + // 场景 3: 往复插拔连续切换 (拔出 → 插入 → 再拔出,验证时空可组合性与零残留) + // ══════════════════════════════════════════════════════════════════════════ + t.Run("Scenario_Dynamic_Plug_Unplug_Sequence", func(t *testing.T) { + for i := 1; i <= 2; i++ { + // 1. 拔出 Redis 运行 + config.Config.Redis.Enabled = false + appUnplugged := newWaveletApp(core.ProfileAll) + require.NoError(t, appUnplugged.Start(context.Background())) + + cacheSvc1, err := core.Inject[contracts.CacheService](appUnplugged.Context()) + require.NoError(t, err) + require.NoError(t, cacheSvc1.Set(context.Background(), fmt.Sprintf("seq_key_%d", i), "seq_val_unplugged", time.Minute)) + + require.NoError(t, appUnplugged.Stop(context.Background())) + + // 2. 插入 Redis 运行 + config.Config.Redis.Enabled = true + appPlugged := newWaveletApp(core.ProfileAll) + require.NoError(t, appPlugged.Start(context.Background())) + + cacheSvc2, err := core.Inject[contracts.CacheService](appPlugged.Context()) + require.NoError(t, err) + require.NoError(t, cacheSvc2.Set(context.Background(), fmt.Sprintf("seq_key_%d", i), "seq_val_plugged", time.Minute)) + + require.NoError(t, appPlugged.Stop(context.Background())) + } + }) +} diff --git a/backend/plugins/domain/admin/migrations/00001_initial.sql b/backend/plugins/domain/admin/migrations/postgres/00001_initial.sql similarity index 99% rename from backend/plugins/domain/admin/migrations/00001_initial.sql rename to backend/plugins/domain/admin/migrations/postgres/00001_initial.sql index 8636c213..e81b9b49 100644 --- a/backend/plugins/domain/admin/migrations/00001_initial.sql +++ b/backend/plugins/domain/admin/migrations/postgres/00001_initial.sql @@ -84,4 +84,4 @@ DELETE FROM w_system_configs WHERE key IN ( ); DROP TABLE IF EXISTS w_templates; DROP TABLE IF EXISTS w_system_configs; --- +goose StatementEnd \ No newline at end of file +-- +goose StatementEnd diff --git a/backend/plugins/domain/admin/migrations/sqlite/00001_initial.sql b/backend/plugins/domain/admin/migrations/sqlite/00001_initial.sql new file mode 100644 index 00000000..6fcd44e1 --- /dev/null +++ b/backend/plugins/domain/admin/migrations/sqlite/00001_initial.sql @@ -0,0 +1,87 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE IF NOT EXISTS w_system_configs ( + key VARCHAR(64) PRIMARY KEY, + value TEXT NOT NULL, + type VARCHAR(32) NOT NULL DEFAULT 'system', + visibility INTEGER NOT NULL DEFAULT 0, + description VARCHAR(255), + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP +); + +CREATE TABLE IF NOT EXISTS w_templates ( + id BIGINT PRIMARY KEY, + key VARCHAR(80) NOT NULL UNIQUE, + name VARCHAR(100) NOT NULL, + type VARCHAR(20) NOT NULL DEFAULT 'email', + subject VARCHAR(255), + content TEXT NOT NULL, + description VARCHAR(255), + is_system BOOLEAN NOT NULL DEFAULT 0, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_templates_is_system ON w_templates (is_system); +CREATE INDEX IF NOT EXISTS idx_w_templates_created_at ON w_templates (created_at); +CREATE INDEX IF NOT EXISTS idx_w_templates_updated_at ON w_templates (updated_at); + +-- Seed system configs (all default platform configs) +INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at) VALUES + ('cap_login_enabled', 'false', 'system', 1, '是否启用登录人机验证(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('cap_auto_solve', 'true', 'system', 1, '打开页面后是否自动开始计算,关闭则需用户手动点击触发', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('cap_challenge_count', '1', 'system', 0, '客户端需求解的 PoW 难题总数,默认 1,推荐 1~5', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('cap_challenge_size', '32', 'system', 0, '人机验证盐值长度', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('cap_challenge_difficulty', '4', 'system', 0, '人机验证 PoW 难度(目标前缀长度)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('cap_challenge_ttl_seconds', '600', 'system', 0, '人机验证难题有效时间(秒)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('cap_token_ttl_seconds', '1200', 'system', 0, '人机验证兑换凭证有效时间(秒)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('server_address', '', 'system', 0, '服务器地址(用于跨域源控制,不设定则允许任意源)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('smtp_host', '', 'system', 0, 'SMTP 服务器地址(例如 smtp.example.com)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('smtp_port', '587', 'system', 0, 'SMTP 端口(例如 587 或 465)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('smtp_username', '', 'system', 0, 'SMTP 账户(如 sender@example.com)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('smtp_password', '', 'system', 0, 'SMTP 访问凭证(授权码/密码)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('upload_allowed_extensions', 'jpg,png,webp', 'system', 1, '允许上传的图片扩展名(逗号分隔)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('site_name', 'Wavelet', 'system', 1, '系统平台的展示名称', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('password_login_enabled', 'true', 'system', 1, '是否允许使用账号密码登录', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('registration_enabled', 'true', 'system', 1, '控制普通用户是否可以自主注册(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('password_register_enabled', 'true', 'system', 1, '是否允许通过密码创建本地账号', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('oidc_login_enabled', 'true', 'system', 1, '是否允许使用第三方 OIDC 认证源登录', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('max_api_keys_per_user', '5', 'business', 1, '限制每个普通用户可以创建的 API Key 最大数量', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('email_login_verification_enabled', 'false', 'system', 1, '是否开启邮箱登录验证(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('email_register_verification_enabled', 'false', 'system', 1, '是否开启邮箱注册验证(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('menu_display_config', '{}', 'system', 1, '目录显示配置(JSON 字符串,格式为 {url: enabled})', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('search_engine_indexing_enabled', 'false', 'system', 1, '是否允许搜索引擎爬取/检索该站点(true/false)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('update_upstream_repository', 'Rain-kl/Wavelet', 'system', 0, 'GitHub Actions Release 上游仓库(owner/repo 或 GitHub 仓库地址)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('storage_config', '{"driver":"local","local":{"root":"."},"s3":{"region":"us-east-1"},"r2":{"region":"auto"},"minio":{"region":"us-east-1","path_style":true},"oss":{},"webdav":{}}', 'system', 0, '文件存储驱动及连接配置(JSON)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('disk_cache_max_size_mb', '1024', 'system', 0, '磁盘缓存最大空间大小(MB)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('disk_cache_ttl_minutes', '1440', 'system', 0, '磁盘缓存默认有效期(分钟)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('disk_cache_lru_enabled', 'true', 'system', 0, '是否启用 LRU 淘汰机制', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('file_access_whitelist', '["avatar"]', 'system', 0, '免登录访问的文件业务类型白名单 (JSON 数组)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('login_session_ttl_hours', '168', 'system', 0, '登录会话过期时间(小时)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('log_database', '', 'system', 0, '当前日志主库(postgres/sqlite/clickhouse),由切换任务写入', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('log_db_migration', '', 'system', 0, '日志库迁移冻结标记(空或 migrating)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) +ON CONFLICT (key) DO NOTHING; + +INSERT INTO w_templates (id, key, name, type, subject, content, description, is_system, created_at, updated_at) VALUES + (1, 'login_email', '登录验证码邮件', 'email', 'Wavelet 登录验证码', '
您的登录验证码为:{{.Code}},5分钟内有效,请勿将验证码泄露给他人。
', '用户密码登录时发送的验证码邮件模板,支持变量:{{.Code}}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + (2, 'register_email', '注册验证码邮件', 'email', 'Wavelet 注册验证码', '您的注册验证码为:{{.Code}},5分钟内有效,请勿泄露给他人。
', '用户注册时发送的验证码邮件模板,支持变量:{{.Code}}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) +ON CONFLICT (key) DO NOTHING; +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DELETE FROM w_templates WHERE key IN ('login_email', 'register_email'); +DELETE FROM w_system_configs WHERE key IN ( + 'cap_login_enabled', 'cap_auto_solve', 'cap_challenge_count', 'cap_challenge_size', + 'cap_challenge_difficulty', 'cap_challenge_ttl_seconds', 'cap_token_ttl_seconds', + 'server_address', 'smtp_host', 'smtp_port', 'smtp_username', 'smtp_password', + 'upload_allowed_extensions', 'site_name', 'password_login_enabled', 'registration_enabled', + 'password_register_enabled', 'oidc_login_enabled', 'max_api_keys_per_user', + 'email_login_verification_enabled', 'email_register_verification_enabled', + 'menu_display_config', 'search_engine_indexing_enabled', 'update_upstream_repository', + 'storage_config', 'disk_cache_max_size_mb', 'disk_cache_ttl_minutes', 'disk_cache_lru_enabled', + 'file_access_whitelist', 'login_session_ttl_hours', 'log_database', 'log_db_migration' +); +DROP TABLE IF EXISTS w_templates; +DROP TABLE IF EXISTS w_system_configs; +-- +goose StatementEnd diff --git a/backend/plugins/domain/admin/plugin.go b/backend/plugins/domain/admin/plugin.go index 09bba20b..4dc4cf54 100644 --- a/backend/plugins/domain/admin/plugin.go +++ b/backend/plugins/domain/admin/plugin.go @@ -17,7 +17,7 @@ import ( "Wavelet/core/extpoints" ) -//go:embed migrations/*.sql +//go:embed migrations/*/*.sql var adminMigrations embed.FS // Option configures the admin plugin. diff --git a/backend/plugins/domain/auth/migrations/00001_initial.sql b/backend/plugins/domain/auth/migrations/postgres/00001_initial.sql similarity index 81% rename from backend/plugins/domain/auth/migrations/00001_initial.sql rename to backend/plugins/domain/auth/migrations/postgres/00001_initial.sql index 88b2a261..d82265ed 100644 --- a/backend/plugins/domain/auth/migrations/00001_initial.sql +++ b/backend/plugins/domain/auth/migrations/postgres/00001_initial.sql @@ -42,17 +42,11 @@ CREATE TABLE IF NOT EXISTS w_access_tokens ( updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP ); CREATE INDEX IF NOT EXISTS idx_w_access_tokens_user_id ON w_access_tokens (user_id); - --- Seed: login session TTL -INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at) -VALUES ('login_session_ttl_hours', '0', 'system', 0, '登录会话过期时间 (小时,0表示浏览器关闭后自动退出,-1表示永不过期)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) -ON CONFLICT (key) DO NOTHING; -- +goose StatementEnd -- +goose Down -- +goose StatementBegin -DELETE FROM w_system_configs WHERE key = 'login_session_ttl_hours'; DROP TABLE IF EXISTS w_access_tokens; DROP TABLE IF EXISTS w_external_accounts; DROP TABLE IF EXISTS w_auth_sources; --- +goose StatementEnd \ No newline at end of file +-- +goose StatementEnd diff --git a/backend/plugins/domain/auth/migrations/sqlite/00001_initial.sql b/backend/plugins/domain/auth/migrations/sqlite/00001_initial.sql new file mode 100644 index 00000000..8995d821 --- /dev/null +++ b/backend/plugins/domain/auth/migrations/sqlite/00001_initial.sql @@ -0,0 +1,52 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE IF NOT EXISTS w_auth_sources ( + id BIGINT PRIMARY KEY, + name VARCHAR(80) NOT NULL UNIQUE, + type VARCHAR(20) NOT NULL, + display_name VARCHAR(100), + is_active BOOLEAN NOT NULL DEFAULT 0, + client_id VARCHAR(255), + client_secret VARCHAR(1024), + openid_discovery_url VARCHAR(1024), + scopes VARCHAR(255), + icon_url VARCHAR(1024), + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_auth_sources_is_active ON w_auth_sources (is_active); + +CREATE TABLE IF NOT EXISTS w_external_accounts ( + id BIGINT PRIMARY KEY, + auth_source_id BIGINT, + user_id BIGINT NOT NULL, + external_id VARCHAR(255) NOT NULL, + external_username VARCHAR(255), + email VARCHAR(255), + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_external_accounts_auth_source_id ON w_external_accounts (auth_source_id); +CREATE INDEX IF NOT EXISTS idx_w_external_accounts_user_id ON w_external_accounts (user_id); +CREATE UNIQUE INDEX IF NOT EXISTS idx_w_external_accounts_source_external ON w_external_accounts (auth_source_id, external_id); + +CREATE TABLE IF NOT EXISTS w_access_tokens ( + id BIGINT PRIMARY KEY, + user_id BIGINT NOT NULL, + token_hash VARCHAR(64) NOT NULL UNIQUE, + name VARCHAR(128) NOT NULL, + description VARCHAR(255), + is_admin BOOLEAN DEFAULT 0, + expires_at DATETIME, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_access_tokens_user_id ON w_access_tokens (user_id); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DROP TABLE IF EXISTS w_access_tokens; +DROP TABLE IF EXISTS w_external_accounts; +DROP TABLE IF EXISTS w_auth_sources; +-- +goose StatementEnd diff --git a/backend/plugins/domain/auth/plugin.go b/backend/plugins/domain/auth/plugin.go index 31e32a70..540af001 100644 --- a/backend/plugins/domain/auth/plugin.go +++ b/backend/plugins/domain/auth/plugin.go @@ -14,7 +14,7 @@ import ( "Wavelet/core/extpoints" ) -//go:embed migrations/*.sql +//go:embed migrations/*/*.sql var authMigrations embed.FS // Option configures the auth plugin. diff --git a/backend/plugins/domain/message_gateway/migrations/00001_initial.sql b/backend/plugins/domain/message_gateway/migrations/postgres/00001_initial.sql similarity index 99% rename from backend/plugins/domain/message_gateway/migrations/00001_initial.sql rename to backend/plugins/domain/message_gateway/migrations/postgres/00001_initial.sql index 091f2645..4a7dc866 100644 --- a/backend/plugins/domain/message_gateway/migrations/00001_initial.sql +++ b/backend/plugins/domain/message_gateway/migrations/postgres/00001_initial.sql @@ -90,4 +90,4 @@ DROP TABLE IF EXISTS w_push_events; DROP TABLE IF EXISTS w_message_pairing_codes; DROP TABLE IF EXISTS w_message_bindings; DROP TABLE IF EXISTS w_message_channels; --- +goose StatementEnd \ No newline at end of file +-- +goose StatementEnd diff --git a/backend/plugins/domain/message_gateway/migrations/sqlite/00001_initial.sql b/backend/plugins/domain/message_gateway/migrations/sqlite/00001_initial.sql new file mode 100644 index 00000000..62c22281 --- /dev/null +++ b/backend/plugins/domain/message_gateway/migrations/sqlite/00001_initial.sql @@ -0,0 +1,93 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE IF NOT EXISTS w_message_channels ( + id BIGINT PRIMARY KEY, + name VARCHAR(128) NOT NULL, + type VARCHAR(32) NOT NULL, + owner_scope VARCHAR(16) NOT NULL DEFAULT 'system', + owner_id BIGINT NULL, + enabled BOOLEAN NOT NULL DEFAULT 1, + credentials TEXT NOT NULL DEFAULT '', + extra TEXT NOT NULL DEFAULT '', + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_message_channels_type ON w_message_channels (type); + +CREATE TABLE IF NOT EXISTS w_message_bindings ( + id BIGINT PRIMARY KEY, + user_id BIGINT NOT NULL, + channel_id BIGINT NOT NULL, + platform_user_id VARCHAR(128) NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); +CREATE UNIQUE INDEX IF NOT EXISTS uniq_w_message_bindings_channel_platform + ON w_message_bindings (channel_id, platform_user_id); +CREATE INDEX IF NOT EXISTS idx_w_message_bindings_user ON w_message_bindings (user_id); + +CREATE TABLE IF NOT EXISTS w_message_pairing_codes ( + code VARCHAR(16) PRIMARY KEY, + channel_id BIGINT NOT NULL, + platform_user_id VARCHAR(128) NOT NULL, + expires_at DATETIME NOT NULL, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_message_pairing_lookup + ON w_message_pairing_codes (channel_id, platform_user_id); + +CREATE TABLE IF NOT EXISTS w_push_events ( + id BIGINT PRIMARY KEY, + event_key VARCHAR(80) NOT NULL, + name VARCHAR(100) NOT NULL, + task_type VARCHAR(100) NOT NULL DEFAULT '', + channels TEXT NOT NULL DEFAULT '', + targets TEXT NOT NULL DEFAULT '', + template TEXT NOT NULL DEFAULT '', + enabled BOOLEAN NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); +CREATE UNIQUE INDEX IF NOT EXISTS uniq_w_push_events_key ON w_push_events(event_key); +CREATE INDEX IF NOT EXISTS idx_w_push_events_enabled ON w_push_events(enabled); +CREATE INDEX IF NOT EXISTS idx_w_push_events_task_type ON w_push_events(task_type); + +CREATE TABLE IF NOT EXISTS w_push_channels ( + id BIGINT PRIMARY KEY, + name VARCHAR(80) NOT NULL, + description VARCHAR(255) NOT NULL DEFAULT '', + type VARCHAR(50) NOT NULL DEFAULT 'custom', + token VARCHAR(100) NOT NULL DEFAULT '', + url TEXT NOT NULL DEFAULT '', + other TEXT NOT NULL DEFAULT '', + enabled BOOLEAN NOT NULL DEFAULT 1, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); +CREATE UNIQUE INDEX IF NOT EXISTS uniq_w_push_channels_name ON w_push_channels(name); +CREATE INDEX IF NOT EXISTS idx_w_push_channels_enabled ON w_push_channels(enabled); + +CREATE TABLE IF NOT EXISTS w_push_histories ( + id BIGINT PRIMARY KEY, + event_key VARCHAR(80) NOT NULL, + channel VARCHAR(50) NOT NULL, + target VARCHAR(255) NOT NULL, + title VARCHAR(255) NOT NULL, + content TEXT NOT NULL, + level VARCHAR(20) NOT NULL, + status VARCHAR(20) NOT NULL, + error_msg TEXT NOT NULL DEFAULT '', + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_push_histories_event ON w_push_histories(event_key); +CREATE INDEX IF NOT EXISTS idx_w_push_histories_created ON w_push_histories(created_at); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DROP TABLE IF EXISTS w_push_histories; +DROP TABLE IF EXISTS w_push_channels; +DROP TABLE IF EXISTS w_push_events; +DROP TABLE IF EXISTS w_message_pairing_codes; +DROP TABLE IF EXISTS w_message_bindings; +DROP TABLE IF EXISTS w_message_channels; +-- +goose StatementEnd diff --git a/backend/plugins/domain/message_gateway/plugin.go b/backend/plugins/domain/message_gateway/plugin.go index 0c3ef847..9ff18809 100644 --- a/backend/plugins/domain/message_gateway/plugin.go +++ b/backend/plugins/domain/message_gateway/plugin.go @@ -17,7 +17,7 @@ import ( "Wavelet/pkg/util" ) -//go:embed migrations/*.sql +//go:embed migrations/*/*.sql var mgMigrations embed.FS // Option configures the message_gateway plugin. diff --git a/backend/plugins/domain/risk_control/logstore/migrations/00001_initial.sql b/backend/plugins/domain/risk_control/logstore/migrations/postgres/00001_initial.sql similarity index 84% rename from backend/plugins/domain/risk_control/logstore/migrations/00001_initial.sql rename to backend/plugins/domain/risk_control/logstore/migrations/postgres/00001_initial.sql index 39b58cde..2cc6d07d 100644 --- a/backend/plugins/domain/risk_control/logstore/migrations/00001_initial.sql +++ b/backend/plugins/domain/risk_control/logstore/migrations/postgres/00001_initial.sql @@ -1,4 +1,5 @@ -- +goose Up +-- +goose StatementBegin CREATE TABLE IF NOT EXISTS w_user_access_logs ( id BIGINT NOT NULL, user_id BIGINT NOT NULL DEFAULT 0, @@ -11,10 +12,13 @@ CREATE TABLE IF NOT EXISTS w_user_access_logs ( latency BIGINT NOT NULL DEFAULT 0, created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (id, created_at) -) PARTITION BY RANGE (created_at); +); CREATE INDEX IF NOT EXISTS idx_w_user_access_logs_user_id ON w_user_access_logs (user_id, created_at DESC); CREATE INDEX IF NOT EXISTS idx_w_user_access_logs_created_at ON w_user_access_logs (created_at DESC); +-- +goose StatementEnd -- +goose Down -DROP TABLE IF EXISTS w_user_access_logs; \ No newline at end of file +-- +goose StatementBegin +DROP TABLE IF EXISTS w_user_access_logs; +-- +goose StatementEnd diff --git a/backend/plugins/domain/risk_control/logstore/migrations/sqlite/00001_initial.sql b/backend/plugins/domain/risk_control/logstore/migrations/sqlite/00001_initial.sql new file mode 100644 index 00000000..d928b8d1 --- /dev/null +++ b/backend/plugins/domain/risk_control/logstore/migrations/sqlite/00001_initial.sql @@ -0,0 +1,24 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE IF NOT EXISTS w_user_access_logs ( + id BIGINT NOT NULL, + user_id BIGINT NOT NULL DEFAULT 0, + path VARCHAR(2048) NOT NULL DEFAULT '', + method VARCHAR(16) NOT NULL DEFAULT '', + ip VARCHAR(128) NOT NULL DEFAULT '', + user_agent TEXT NOT NULL DEFAULT '', + headers TEXT NOT NULL DEFAULT '', + status INTEGER NOT NULL DEFAULT 0, + latency BIGINT NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (id, created_at) +); + +CREATE INDEX IF NOT EXISTS idx_w_user_access_logs_user_id ON w_user_access_logs (user_id, created_at DESC); +CREATE INDEX IF NOT EXISTS idx_w_user_access_logs_created_at ON w_user_access_logs (created_at DESC); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DROP TABLE IF EXISTS w_user_access_logs; +-- +goose StatementEnd diff --git a/backend/plugins/domain/risk_control/plugin.go b/backend/plugins/domain/risk_control/plugin.go index f1ca38dc..5027e4b8 100644 --- a/backend/plugins/domain/risk_control/plugin.go +++ b/backend/plugins/domain/risk_control/plugin.go @@ -17,7 +17,7 @@ import ( "Wavelet/plugins/domain/risk_control/logstore" ) -//go:embed logstore/migrations/*.sql +//go:embed logstore/migrations/*/*.sql var riskControlMigrations embed.FS // Option configures the risk_control plugin. diff --git a/backend/plugins/domain/upload/migrations/00001_initial.sql b/backend/plugins/domain/upload/migrations/00001_initial.sql deleted file mode 100644 index c51815e2..00000000 --- a/backend/plugins/domain/upload/migrations/00001_initial.sql +++ /dev/null @@ -1,49 +0,0 @@ --- +goose Up --- +goose StatementBegin -CREATE TABLE IF NOT EXISTS w_upload_stats ( - dimension VARCHAR(32) NOT NULL, - stat_key VARCHAR(64) NOT NULL DEFAULT '', - file_count BIGINT NOT NULL DEFAULT 0, - file_size BIGINT NOT NULL DEFAULT 0, - updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, - PRIMARY KEY (dimension, stat_key) -); - -CREATE INDEX IF NOT EXISTS idx_w_uploads_status_created_at ON w_uploads (status, created_at); -CREATE INDEX IF NOT EXISTS idx_w_uploads_hash_file_size_status ON w_uploads (hash, file_size, status); - -ALTER TABLE w_uploads ADD COLUMN IF NOT EXISTS access_mode INTEGER NOT NULL DEFAULT 0; -UPDATE w_uploads SET access_mode = 1 WHERE type = 'avatar'; - --- Backfill upload stats -INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) -SELECT 'total', '', COUNT(*), COALESCE(SUM(file_size), 0) -FROM w_uploads -WHERE status != 'deleted' -ON CONFLICT (dimension, stat_key) DO UPDATE SET - file_count = EXCLUDED.file_count, - file_size = EXCLUDED.file_size, - updated_at = CURRENT_TIMESTAMP; - -INSERT INTO w_upload_stats (dimension, stat_key, file_count, file_size) -SELECT - 'type', - COALESCE(NULLIF(type, ''), 'generic'), - COUNT(*), - COALESCE(SUM(file_size), 0) -FROM w_uploads -WHERE status != 'deleted' -GROUP BY COALESCE(NULLIF(type, ''), 'generic') -ON CONFLICT (dimension, stat_key) DO UPDATE SET - file_count = EXCLUDED.file_count, - file_size = EXCLUDED.file_size, - updated_at = CURRENT_TIMESTAMP; --- +goose StatementEnd - --- +goose Down --- +goose StatementBegin -DROP TABLE IF EXISTS w_upload_stats; -DROP INDEX IF EXISTS idx_w_uploads_hash_file_size_status; -DROP INDEX IF EXISTS idx_w_uploads_status_created_at; -ALTER TABLE w_uploads DROP COLUMN IF EXISTS access_mode; --- +goose StatementEnd \ No newline at end of file diff --git a/backend/plugins/domain/upload/migrations/postgres/00001_initial.sql b/backend/plugins/domain/upload/migrations/postgres/00001_initial.sql new file mode 100644 index 00000000..a0b7b69b --- /dev/null +++ b/backend/plugins/domain/upload/migrations/postgres/00001_initial.sql @@ -0,0 +1,40 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE IF NOT EXISTS w_uploads ( + id BIGINT PRIMARY KEY, + user_id BIGINT NOT NULL, + file_name VARCHAR(255) NOT NULL, + file_path VARCHAR(500) NOT NULL, + file_size BIGINT NOT NULL, + mime_type VARCHAR(100) NOT NULL, + extension VARCHAR(50) NOT NULL, + hash VARCHAR(64), + type VARCHAR(50) NOT NULL, + status VARCHAR(20) NOT NULL, + access_mode INTEGER NOT NULL DEFAULT 0, + metadata JSONB, + created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, + updated_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_uploads_user_id ON w_uploads (user_id); +CREATE INDEX IF NOT EXISTS idx_w_uploads_file_path ON w_uploads (file_path); +CREATE INDEX IF NOT EXISTS idx_w_uploads_hash ON w_uploads (hash); +CREATE INDEX IF NOT EXISTS idx_w_uploads_type ON w_uploads (type); +CREATE INDEX IF NOT EXISTS idx_w_uploads_status_created_at ON w_uploads (status, created_at); +CREATE INDEX IF NOT EXISTS idx_w_uploads_hash_file_size_status ON w_uploads (hash, file_size, status); + +CREATE TABLE IF NOT EXISTS w_upload_stats ( + dimension VARCHAR(32) NOT NULL, + stat_key VARCHAR(64) NOT NULL DEFAULT '', + file_count BIGINT NOT NULL DEFAULT 0, + file_size BIGINT NOT NULL DEFAULT 0, + updated_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (dimension, stat_key) +); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DROP TABLE IF EXISTS w_upload_stats; +DROP TABLE IF EXISTS w_uploads; +-- +goose StatementEnd diff --git a/backend/plugins/domain/upload/migrations/sqlite/00001_initial.sql b/backend/plugins/domain/upload/migrations/sqlite/00001_initial.sql new file mode 100644 index 00000000..05ec33e8 --- /dev/null +++ b/backend/plugins/domain/upload/migrations/sqlite/00001_initial.sql @@ -0,0 +1,40 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE IF NOT EXISTS w_uploads ( + id BIGINT PRIMARY KEY, + user_id BIGINT NOT NULL, + file_name VARCHAR(255) NOT NULL, + file_path VARCHAR(500) NOT NULL, + file_size BIGINT NOT NULL, + mime_type VARCHAR(100) NOT NULL, + extension VARCHAR(50) NOT NULL, + hash VARCHAR(64), + type VARCHAR(50) NOT NULL, + status VARCHAR(20) NOT NULL, + access_mode INTEGER NOT NULL DEFAULT 0, + metadata TEXT NOT NULL DEFAULT '{}', + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_uploads_user_id ON w_uploads (user_id); +CREATE INDEX IF NOT EXISTS idx_w_uploads_file_path ON w_uploads (file_path); +CREATE INDEX IF NOT EXISTS idx_w_uploads_hash ON w_uploads (hash); +CREATE INDEX IF NOT EXISTS idx_w_uploads_type ON w_uploads (type); +CREATE INDEX IF NOT EXISTS idx_w_uploads_status_created_at ON w_uploads (status, created_at); +CREATE INDEX IF NOT EXISTS idx_w_uploads_hash_file_size_status ON w_uploads (hash, file_size, status); + +CREATE TABLE IF NOT EXISTS w_upload_stats ( + dimension VARCHAR(32) NOT NULL, + stat_key VARCHAR(64) NOT NULL DEFAULT '', + file_count BIGINT NOT NULL DEFAULT 0, + file_size BIGINT NOT NULL DEFAULT 0, + updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (dimension, stat_key) +); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DROP TABLE IF EXISTS w_upload_stats; +DROP TABLE IF EXISTS w_uploads; +-- +goose StatementEnd diff --git a/backend/plugins/domain/upload/plugin.go b/backend/plugins/domain/upload/plugin.go index 6a4bd195..af076ac1 100644 --- a/backend/plugins/domain/upload/plugin.go +++ b/backend/plugins/domain/upload/plugin.go @@ -21,7 +21,7 @@ import ( "Wavelet/plugins/domain/upload/task" ) -//go:embed migrations/*.sql +//go:embed migrations/*/*.sql var uploadMigrations embed.FS // Plugin implements core.Plugin to provide file upload and media serving domain services. diff --git a/backend/plugins/domain/user/migrations/00001_initial.sql b/backend/plugins/domain/user/migrations/postgres/00001_initial.sql similarity index 98% rename from backend/plugins/domain/user/migrations/00001_initial.sql rename to backend/plugins/domain/user/migrations/postgres/00001_initial.sql index ee891213..d7853608 100644 --- a/backend/plugins/domain/user/migrations/00001_initial.sql +++ b/backend/plugins/domain/user/migrations/postgres/00001_initial.sql @@ -33,4 +33,4 @@ ON CONFLICT (username) DO NOTHING; -- +goose StatementBegin DELETE FROM w_users WHERE username = 'system'; DROP TABLE IF EXISTS w_users; --- +goose StatementEnd \ No newline at end of file +-- +goose StatementEnd diff --git a/backend/plugins/domain/user/migrations/sqlite/00001_initial.sql b/backend/plugins/domain/user/migrations/sqlite/00001_initial.sql new file mode 100644 index 00000000..27c28cd9 --- /dev/null +++ b/backend/plugins/domain/user/migrations/sqlite/00001_initial.sql @@ -0,0 +1,36 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE IF NOT EXISTS w_users ( + id BIGINT PRIMARY KEY, + username VARCHAR(64) NOT NULL UNIQUE, + password VARCHAR(255), + nickname VARCHAR(255), + email VARCHAR(255), + avatar_url VARCHAR(255), + is_active BOOLEAN DEFAULT 1, + is_admin BOOLEAN DEFAULT 0, + bio VARCHAR(500), + phone VARCHAR(32), + gender VARCHAR(16), + website VARCHAR(255), + location VARCHAR(255), + last_login_at DATETIME, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_users_email ON w_users (email); +CREATE INDEX IF NOT EXISTS idx_w_users_is_active ON w_users (is_active); +CREATE INDEX IF NOT EXISTS idx_w_users_last_login_at ON w_users (last_login_at); +CREATE INDEX IF NOT EXISTS idx_w_users_created_at ON w_users (created_at); + +-- Seed system user +INSERT INTO w_users (id, username, password, nickname, avatar_url, is_active, is_admin, last_login_at, created_at, updated_at) +VALUES (999, 'system', '*', '系统', '', 1, 0, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) +ON CONFLICT (username) DO NOTHING; +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DELETE FROM w_users WHERE username = 'system'; +DROP TABLE IF EXISTS w_users; +-- +goose StatementEnd diff --git a/backend/plugins/domain/user/plugin.go b/backend/plugins/domain/user/plugin.go index c0329c31..af17d090 100644 --- a/backend/plugins/domain/user/plugin.go +++ b/backend/plugins/domain/user/plugin.go @@ -17,7 +17,7 @@ import ( "Wavelet/core/extpoints" ) -//go:embed migrations/*.sql +//go:embed migrations/*/*.sql var userMigrations embed.FS // Option configures the user plugin. diff --git a/backend/plugins/drivers/driver_asynq_cron/migrations/00001_initial.sql b/backend/plugins/drivers/driver_asynq_cron/migrations/postgres/00001_initial.sql similarity index 84% rename from backend/plugins/drivers/driver_asynq_cron/migrations/00001_initial.sql rename to backend/plugins/drivers/driver_asynq_cron/migrations/postgres/00001_initial.sql index 8709f302..024b0643 100644 --- a/backend/plugins/drivers/driver_asynq_cron/migrations/00001_initial.sql +++ b/backend/plugins/drivers/driver_asynq_cron/migrations/postgres/00001_initial.sql @@ -1,4 +1,5 @@ -- +goose Up +-- +goose StatementBegin CREATE TABLE IF NOT EXISTS w_schedules ( id BIGINT PRIMARY KEY, name VARCHAR(128) NOT NULL, @@ -15,6 +16,9 @@ CREATE INDEX IF NOT EXISTS idx_w_schedules_is_active ON w_schedules (is_active); INSERT INTO w_schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at) VALUES (1, '系统定期垃圾清理', 'system_cleanup', '0 3 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) ON CONFLICT (id) DO NOTHING; +-- +goose StatementEnd -- +goose Down -DROP TABLE IF EXISTS w_schedules; \ No newline at end of file +-- +goose StatementBegin +DROP TABLE IF EXISTS w_schedules; +-- +goose StatementEnd diff --git a/backend/plugins/drivers/driver_asynq_cron/migrations/sqlite/00001_initial.sql b/backend/plugins/drivers/driver_asynq_cron/migrations/sqlite/00001_initial.sql new file mode 100644 index 00000000..3ace525a --- /dev/null +++ b/backend/plugins/drivers/driver_asynq_cron/migrations/sqlite/00001_initial.sql @@ -0,0 +1,24 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE IF NOT EXISTS w_schedules ( + id BIGINT PRIMARY KEY, + name VARCHAR(128) NOT NULL, + task_type VARCHAR(64) NOT NULL, + cron VARCHAR(64) NOT NULL, + payload TEXT, + is_active BOOLEAN NOT NULL DEFAULT 1, + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_schedules_is_active ON w_schedules (is_active); + +-- Seed initial cleanup task +INSERT INTO w_schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at) +VALUES (1, '系统定期垃圾清理', 'system_cleanup', '0 3 * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) +ON CONFLICT (id) DO NOTHING; +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DROP TABLE IF EXISTS w_schedules; +-- +goose StatementEnd diff --git a/backend/plugins/drivers/driver_asynq_cron/plugin.go b/backend/plugins/drivers/driver_asynq_cron/plugin.go index e61960a6..a8ffd150 100644 --- a/backend/plugins/drivers/driver_asynq_cron/plugin.go +++ b/backend/plugins/drivers/driver_asynq_cron/plugin.go @@ -16,9 +16,10 @@ import ( "Wavelet/core" "Wavelet/core/contracts" + "Wavelet/pkg/config" ) -//go:embed migrations/*.sql +//go:embed migrations/*/*.sql var cronMigrations embed.FS // Option configures the Asynq cron scheduler driver plugin. @@ -148,7 +149,25 @@ func (p *Plugin) Start(_ context.Context) error { opts.Location = p.location } - p.scheduler = asynq.NewScheduler(p.redisOpt, opts) + opt := p.redisOpt + if opt == nil { + if RedisOpt != nil { + opt = RedisOpt + } else { + redisCfg := config.Config.Redis + addr := "127.0.0.1:6379" + if len(redisCfg.Addrs) > 0 && redisCfg.Addrs[0] != "" { + addr = redisCfg.Addrs[0] + } + opt = asynq.RedisClientOpt{ + Addr: addr, + Username: redisCfg.Username, + Password: redisCfg.Password, + DB: redisCfg.DB, + } + } + } + p.scheduler = asynq.NewScheduler(opt, opts) } if p.coreCtx != nil && p.coreCtx.Schedules() != nil { diff --git a/backend/plugins/drivers/driver_asynq_worker/meta.go b/backend/plugins/drivers/driver_asynq_worker/meta.go index 853cfd09..620fb304 100644 --- a/backend/plugins/drivers/driver_asynq_worker/meta.go +++ b/backend/plugins/drivers/driver_asynq_worker/meta.go @@ -5,6 +5,8 @@ package driver_asynq_worker import ( "sync" + + "Wavelet/core/extpoints" ) // TaskParam 任务参数定义 @@ -61,30 +63,65 @@ func GetDispatchableTasks() []TaskMeta { return metas } -// GetTaskMeta 根据任务类型获取元数据 -func GetTaskMeta(taskType string) *TaskMeta { - dispatchableTasksMutex.RLock() - defer dispatchableTasksMutex.RUnlock() - for _, t := range dispatchableTasks { - if t.Type == taskType { - copied := t - return &copied +var ( + activeTaskRegMutex sync.RWMutex + activeTaskReg extpoints.TaskExtension +) + +// SetActiveTaskExtension sets the active task extension registry for task resolution. +func SetActiveTaskExtension(reg extpoints.TaskExtension) { + activeTaskRegMutex.Lock() + defer activeTaskRegMutex.Unlock() + activeTaskReg = reg +} + +func getFromActiveTaskExtension(taskType string) *TaskMeta { + activeTaskRegMutex.RLock() + defer activeTaskRegMutex.RUnlock() + if activeTaskReg == nil { + return nil + } + if td, ok := activeTaskReg.Get(taskType); ok { + return &TaskMeta{ + Type: td.Pattern, + Name: td.Pattern, + AsynqTask: td.Pattern, + Queue: "default", + Retryable: td.Retry > 0, + MaxRetry: td.Retry, } } return nil } -// GetTaskMetaByAsynqTask 根据 Asynq 任务名称获取元数据 -func GetTaskMetaByAsynqTask(asynqTask string) *TaskMeta { +// GetTaskMeta 根据任务类型获取元数据 +func GetTaskMeta(taskType string) *TaskMeta { dispatchableTasksMutex.RLock() - defer dispatchableTasksMutex.RUnlock() for _, t := range dispatchableTasks { - if t.AsynqTask == asynqTask { + if t.Type == taskType { copied := t + dispatchableTasksMutex.RUnlock() return &copied } } - return nil + dispatchableTasksMutex.RUnlock() + + return getFromActiveTaskExtension(taskType) +} + +// GetTaskMetaByAsynqTask 根据 Asynq 任务名称获取元数据 +func GetTaskMetaByAsynqTask(asynqTask string) *TaskMeta { + dispatchableTasksMutex.RLock() + for _, t := range dispatchableTasks { + if t.AsynqTask == asynqTask { + copied := t + dispatchableTasksMutex.RUnlock() + return &copied + } + } + dispatchableTasksMutex.RUnlock() + + return getFromActiveTaskExtension(asynqTask) } // GetRegisteredAsynqTasks 返回所有已注册的 Asynq 任务名称,以便动态注册路由 diff --git a/backend/plugins/drivers/driver_asynq_worker/migrations/00001_initial.sql b/backend/plugins/drivers/driver_asynq_worker/migrations/postgres/00001_initial.sql similarity index 88% rename from backend/plugins/drivers/driver_asynq_worker/migrations/00001_initial.sql rename to backend/plugins/drivers/driver_asynq_worker/migrations/postgres/00001_initial.sql index 3c59c854..52a6463a 100644 --- a/backend/plugins/drivers/driver_asynq_worker/migrations/00001_initial.sql +++ b/backend/plugins/drivers/driver_asynq_worker/migrations/postgres/00001_initial.sql @@ -1,4 +1,5 @@ -- +goose Up +-- +goose StatementBegin CREATE TABLE IF NOT EXISTS w_task_executions ( id BIGINT PRIMARY KEY, task_id VARCHAR(128) NOT NULL UNIQUE, @@ -23,6 +24,9 @@ CREATE INDEX IF NOT EXISTS idx_w_task_executions_task_type ON w_task_executions CREATE INDEX IF NOT EXISTS idx_w_task_executions_status ON w_task_executions (status); CREATE INDEX IF NOT EXISTS idx_w_task_executions_started_at ON w_task_executions (started_at); CREATE INDEX IF NOT EXISTS idx_w_task_executions_created_at ON w_task_executions (created_at); +-- +goose StatementEnd -- +goose Down -DROP TABLE IF EXISTS w_task_executions; \ No newline at end of file +-- +goose StatementBegin +DROP TABLE IF EXISTS w_task_executions; +-- +goose StatementEnd diff --git a/backend/plugins/drivers/driver_asynq_worker/migrations/sqlite/00001_initial.sql b/backend/plugins/drivers/driver_asynq_worker/migrations/sqlite/00001_initial.sql new file mode 100644 index 00000000..462af866 --- /dev/null +++ b/backend/plugins/drivers/driver_asynq_worker/migrations/sqlite/00001_initial.sql @@ -0,0 +1,32 @@ +-- +goose Up +-- +goose StatementBegin +CREATE TABLE IF NOT EXISTS w_task_executions ( + id BIGINT PRIMARY KEY, + task_id VARCHAR(128) NOT NULL UNIQUE, + task_type VARCHAR(64) NOT NULL, + task_name VARCHAR(128), + status VARCHAR(32) NOT NULL, + retryable BOOLEAN NOT NULL DEFAULT 0, + max_retry INTEGER NOT NULL DEFAULT 0, + retry_count INTEGER NOT NULL DEFAULT 0, + log TEXT, + error_message TEXT, + result TEXT, + started_at DATETIME, + finished_at DATETIME, + duration BIGINT, + payload TEXT, + triggered_by VARCHAR(32) NOT NULL DEFAULT 'system', + created_at DATETIME DEFAULT CURRENT_TIMESTAMP, + updated_at DATETIME DEFAULT CURRENT_TIMESTAMP +); +CREATE INDEX IF NOT EXISTS idx_w_task_executions_task_type ON w_task_executions (task_type); +CREATE INDEX IF NOT EXISTS idx_w_task_executions_status ON w_task_executions (status); +CREATE INDEX IF NOT EXISTS idx_w_task_executions_started_at ON w_task_executions (started_at); +CREATE INDEX IF NOT EXISTS idx_w_task_executions_created_at ON w_task_executions (created_at); +-- +goose StatementEnd + +-- +goose Down +-- +goose StatementBegin +DROP TABLE IF EXISTS w_task_executions; +-- +goose StatementEnd diff --git a/backend/plugins/drivers/driver_asynq_worker/plugin.go b/backend/plugins/drivers/driver_asynq_worker/plugin.go index eb5959f8..921743a2 100644 --- a/backend/plugins/drivers/driver_asynq_worker/plugin.go +++ b/backend/plugins/drivers/driver_asynq_worker/plugin.go @@ -16,6 +16,7 @@ import ( "Wavelet/core" "Wavelet/core/contracts" + "Wavelet/pkg/config" ) const ( @@ -23,7 +24,7 @@ const ( defaultShutdownTimeout = 10 * time.Second ) -//go:embed migrations/*.sql +//go:embed migrations/*/*.sql var workerMigrations embed.FS // Option configures the Asynq worker driver plugin. @@ -131,11 +132,13 @@ func (p *Plugin) Apply(ctx *core.Context) error { // 1. Provide contracts.TaskService p.taskSvc = &taskServiceImpl{} core.Provide[contracts.TaskService](ctx, p.taskSvc) + SetActiveTaskExtension(ctx.Tasks()) // 2. Register migrations for w_task_executions table ctx.Migrations().Register("driver_asynq_worker", workerMigrations) ctx.OnDispose(func() error { + SetActiveTaskExtension(nil) shutdownCtx, cancel := context.WithTimeout(context.Background(), p.shutdownTimeout) defer cancel() return p.Stop(shutdownCtx) @@ -175,8 +178,26 @@ func (p *Plugin) Start(_ context.Context) error { } if p.server == nil { + opt := p.redisOpt + if opt == nil { + if RedisOpt != nil { + opt = RedisOpt + } else { + redisCfg := config.Config.Redis + addr := "127.0.0.1:6379" + if len(redisCfg.Addrs) > 0 && redisCfg.Addrs[0] != "" { + addr = redisCfg.Addrs[0] + } + opt = asynq.RedisClientOpt{ + Addr: addr, + Username: redisCfg.Username, + Password: redisCfg.Password, + DB: redisCfg.DB, + } + } + } p.server = asynq.NewServer( - p.redisOpt, + opt, asynq.Config{ Concurrency: p.concurrency, Queues: p.queues, diff --git a/backend/plugins/infra/cache/plugin.go b/backend/plugins/infra/cache/plugin.go index 06531e65..b330d09e 100644 --- a/backend/plugins/infra/cache/plugin.go +++ b/backend/plugins/infra/cache/plugin.go @@ -81,14 +81,16 @@ func (p *Plugin) Name() string { // Apply mounts the multi-layer cache service into the Context. func (p *Plugin) Apply(ctx *core.Context) error { redisClient := p.redisClient - if redisClient == nil && Redis == nil { - var err error - redisClient, err = InitRedis() - if err != nil { - return err + if redisClient == nil { + if Redis == nil { + var err error + redisClient, err = InitRedis() + if err != nil { + return err + } + } else { + redisClient = Redis } - } else if redisClient == nil { - redisClient = Redis } ramCache, err := ram.New[string, ramEntry](ram.Options{ @@ -111,6 +113,7 @@ func (p *Plugin) Apply(ctx *core.Context) error { ctx.OnDispose(func() error { svc.stopPubSubListener() if p.redisClient == nil { + Redis = nil if closeErr := redisClient.Close(); closeErr != nil && !errors.Is(closeErr, redis.ErrClosed) { return closeErr }