diff --git a/backend/plugins/domain/admin/plugin_test.go b/backend/plugins/domain/admin/plugin_test.go index 23ff5a0f..b0509906 100644 --- a/backend/plugins/domain/admin/plugin_test.go +++ b/backend/plugins/domain/admin/plugin_test.go @@ -38,3 +38,47 @@ func TestAdminPluginUnit(t *testing.T) { require.True(t, ok) assert.Equal(t, "0 4 * * *", setting.Default) } + +func TestAdminMigrationsIncludeTaskExecutionsAndSchedules(t *testing.T) { + ctx := core.NewContext(context.Background()) + p := admin.New() + require.NoError(t, p.Apply(ctx)) + + entry, ok := ctx.Migrations().Get("admin") + require.True(t, ok, "admin plugin must register migrations") + assert.Equal(t, "admin", entry.PluginID) + + // Verify sqlite migration files include w_schedules and w_task_executions + sqliteDir, err := entry.FS.Open("migrations/sqlite/00001_initial.sql") + require.NoError(t, err) + defer sqliteDir.Close() + + stat, err := sqliteDir.Stat() + require.NoError(t, err) + buf := make([]byte, stat.Size()) + _, err = sqliteDir.Read(buf) + require.NoError(t, err) + content := string(buf) + + assert.Contains(t, content, "CREATE TABLE IF NOT EXISTS w_task_executions") + assert.Contains(t, content, "CREATE TABLE IF NOT EXISTS w_schedules") + assert.Contains(t, content, "CREATE TABLE IF NOT EXISTS w_system_configs") + assert.Contains(t, content, "CREATE TABLE IF NOT EXISTS w_templates") + + // Verify postgres migration files include w_schedules and w_task_executions + pgDir, err := entry.FS.Open("migrations/postgres/00001_initial.sql") + require.NoError(t, err) + defer pgDir.Close() + + stat, err = pgDir.Stat() + require.NoError(t, err) + buf = make([]byte, stat.Size()) + _, err = pgDir.Read(buf) + require.NoError(t, err) + pgContent := string(buf) + + assert.Contains(t, pgContent, "CREATE TABLE IF NOT EXISTS w_task_executions") + assert.Contains(t, pgContent, "CREATE TABLE IF NOT EXISTS w_schedules") + assert.Contains(t, pgContent, "CREATE TABLE IF NOT EXISTS w_system_configs") + assert.Contains(t, pgContent, "CREATE TABLE IF NOT EXISTS w_templates") +} diff --git a/backend/plugins/domain/auth/migrations/postgres/00001_initial.sql b/backend/plugins/domain/auth/migrations/postgres/00001_initial.sql index d82265ed..3992a274 100644 --- a/backend/plugins/domain/auth/migrations/postgres/00001_initial.sql +++ b/backend/plugins/domain/auth/migrations/postgres/00001_initial.sql @@ -35,6 +35,7 @@ CREATE TABLE IF NOT EXISTS w_access_tokens ( user_id BIGINT NOT NULL, token_hash VARCHAR(64) NOT NULL UNIQUE, name VARCHAR(128) NOT NULL, + masked_token VARCHAR(64) NOT NULL DEFAULT '', description VARCHAR(255), is_admin BOOLEAN DEFAULT FALSE, expires_at TIMESTAMPTZ, diff --git a/backend/plugins/domain/auth/migrations/sqlite/00001_initial.sql b/backend/plugins/domain/auth/migrations/sqlite/00001_initial.sql index 8995d821..a353bc96 100644 --- a/backend/plugins/domain/auth/migrations/sqlite/00001_initial.sql +++ b/backend/plugins/domain/auth/migrations/sqlite/00001_initial.sql @@ -35,6 +35,7 @@ CREATE TABLE IF NOT EXISTS w_access_tokens ( user_id BIGINT NOT NULL, token_hash VARCHAR(64) NOT NULL UNIQUE, name VARCHAR(128) NOT NULL, + masked_token VARCHAR(64) NOT NULL DEFAULT '', description VARCHAR(255), is_admin BOOLEAN DEFAULT 0, expires_at DATETIME, diff --git a/backend/plugins/domain/domain_test.go b/backend/plugins/domain/domain_test.go index d32c538b..93a57740 100644 --- a/backend/plugins/domain/domain_test.go +++ b/backend/plugins/domain/domain_test.go @@ -151,6 +151,7 @@ func TestUserPlugin(t *testing.T) { require.NoError(t, db.New(db.WithDB(testDB)).Apply(ctx)) require.NoError(t, cache.New().Apply(ctx)) require.NoError(t, logger.New().Apply(ctx)) + require.NoError(t, auth.New().Apply(ctx)) p := user.New() assert.Equal(t, "user", p.Name()) diff --git a/backend/plugins/drivers/driver_asynq_cron/migrations/postgres/00001_initial.sql b/backend/plugins/drivers/driver_asynq_cron/migrations/postgres/00001_initial.sql deleted file mode 100644 index 024b0643..00000000 --- a/backend/plugins/drivers/driver_asynq_cron/migrations/postgres/00001_initial.sql +++ /dev/null @@ -1,24 +0,0 @@ --- +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 TRUE, - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - updated_at TIMESTAMPTZ 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 * * *', '{}', TRUE, 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/migrations/sqlite/00001_initial.sql b/backend/plugins/drivers/driver_asynq_cron/migrations/sqlite/00001_initial.sql deleted file mode 100644 index 3ace525a..00000000 --- a/backend/plugins/drivers/driver_asynq_cron/migrations/sqlite/00001_initial.sql +++ /dev/null @@ -1,24 +0,0 @@ --- +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 5e620bdf..d822d894 100644 --- a/backend/plugins/drivers/driver_asynq_cron/plugin.go +++ b/backend/plugins/drivers/driver_asynq_cron/plugin.go @@ -8,7 +8,6 @@ import ( "Wavelet/core" "Wavelet/core/contracts" "context" - "embed" "encoding/json" "fmt" "sync" @@ -17,9 +16,6 @@ import ( "github.com/hibiken/asynq" ) -//go:embed migrations/*/*.sql -var cronMigrations embed.FS - // Option configures the Asynq cron scheduler driver plugin. type Option func(*Plugin) @@ -148,9 +144,6 @@ func (p *Plugin) Apply(ctx *core.Context) error { return nil }) - // Register migrations for w_schedules table - ctx.Migrations().Register("driver_asynq_cron", cronMigrations) - ctx.OnDispose(func() error { return p.Stop(context.Background()) }) diff --git a/backend/plugins/drivers/driver_asynq_worker/migrations/postgres/00001_initial.sql b/backend/plugins/drivers/driver_asynq_worker/migrations/postgres/00001_initial.sql deleted file mode 100644 index 52a6463a..00000000 --- a/backend/plugins/drivers/driver_asynq_worker/migrations/postgres/00001_initial.sql +++ /dev/null @@ -1,32 +0,0 @@ --- +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 FALSE, - max_retry INTEGER NOT NULL DEFAULT 0, - retry_count INTEGER NOT NULL DEFAULT 0, - log TEXT, - error_message TEXT, - result TEXT, - started_at TIMESTAMPTZ, - finished_at TIMESTAMPTZ, - duration BIGINT, - payload TEXT, - triggered_by VARCHAR(32) NOT NULL DEFAULT 'system', - created_at TIMESTAMPTZ DEFAULT CURRENT_TIMESTAMP, - updated_at TIMESTAMPTZ 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/migrations/sqlite/00001_initial.sql b/backend/plugins/drivers/driver_asynq_worker/migrations/sqlite/00001_initial.sql deleted file mode 100644 index 462af866..00000000 --- a/backend/plugins/drivers/driver_asynq_worker/migrations/sqlite/00001_initial.sql +++ /dev/null @@ -1,32 +0,0 @@ --- +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 cd29567e..8f31d4d9 100644 --- a/backend/plugins/drivers/driver_asynq_worker/plugin.go +++ b/backend/plugins/drivers/driver_asynq_worker/plugin.go @@ -8,7 +8,6 @@ import ( "Wavelet/core" "Wavelet/core/contracts" "context" - "embed" "errors" "fmt" "sync" @@ -23,9 +22,6 @@ const ( defaultShutdownTimeout = 10 * time.Second ) -//go:embed migrations/*/*.sql -var workerMigrations embed.FS - // Option configures the Asynq worker driver plugin. type Option func(*Plugin) @@ -166,9 +162,6 @@ func (p *Plugin) Apply(ctx *core.Context) error { 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) SetRedisClient(nil) diff --git a/docs/WAVELET_WHITE_PAPER.md b/docs/WAVELET_WHITE_PAPER.md index ff1193a9..3f2f99f4 100644 --- a/docs/WAVELET_WHITE_PAPER.md +++ b/docs/WAVELET_WHITE_PAPER.md @@ -152,11 +152,9 @@ Wavelet 贯彻了 Cordis 核心范式,通过形式化保证解决组件系统 | `w_users` | `backend/plugins/domain/user` | `models.go`
`repository.go` | `core/contracts.UserService`
`contracts.UserDTO` | | `w_auth_sources`
`w_external_accounts`
`w_access_tokens` | `backend/plugins/domain/auth` | `models.go`
`service.go` | `core/contracts.AuthService`
`contracts.AuthRegistry` | | `w_uploads`
`w_upload_stats` | `backend/plugins/domain/upload` | `models/models.go`
`repository/repository.go` | `core/contracts.StorageService`
`upload.Ingest` 流水线 | -| `w_system_configs`
`w_templates` | `backend/plugins/domain/admin` | `models.go`
`repository.go` | `ctx.Settings()` / `contracts.ConfigService`
Redis Pub/Sub 广播 | +| `w_system_configs`
`w_templates`
`w_schedules`
`w_task_executions` | `backend/plugins/domain/admin` | `models.go`
`repository.go` | `ctx.Settings()` / `contracts.ConfigService`
`contracts.TaskService` | | `w_message_channels`
`w_message_bindings`
`w_message_pairing_codes`
`w_push_events`
`w_push_channels`
`w_push_histories` | `backend/plugins/domain/message_gateway` | `models.go`
`repository.go` | `EventBus` 强类型事件广播订阅 | | `w_user_access_logs` | `backend/plugins/domain/risk_control` | `logstore/` | `logstore` 门面
ClickHouse PG/SQLite 回落 | -| `w_schedules` | `backend/plugins/drivers/driver_asynq_cron` | `schedule.go` | `ctx.Schedule()` 扩展点 | -| `w_task_executions` | `backend/plugins/drivers/driver_asynq_worker` | `types.go`
`executor.go` | `ctx.Task()` 扩展点 | | `w_schema_versions` | **系统内部** | `backend/cmd/app.go` sharedStore | 自动管理,不归属于任何插件 | ### 5.3 架构防线与单向依赖保障