fix: pg sql

This commit is contained in:
ryan
2026-08-31 15:48:28 +08:00
parent f975c4f7cc
commit 18339ee3d4
7 changed files with 390 additions and 19 deletions
+14 -1
View File
@@ -36,6 +36,7 @@ import (
"github.com/pressly/goose/v3"
goosedb "github.com/pressly/goose/v3/database"
"gorm.io/gorm"
infradb "Wavelet/plugins/infra/database"
)
@@ -279,7 +280,7 @@ func (e *gooseEngine) Migrate(ctx *core.Context, entries []core.MigrationEntry)
return fmt.Errorf("migration: get underlying DB from GORM: %w", err)
}
dialect := gooseDialect(ctx)
dialect := gooseDialectFromGORM(gormDB, ctx)
dialectStr := string(dialect)
goCtx := context.Background()
if ctx != nil {
@@ -346,6 +347,18 @@ func (e *gooseEngine) Migrate(ctx *core.Context, entries []core.MigrationEntry)
return nil
}
// gooseDialectFromGORM prefers the live driver; config is only a fallback when
// GORM has no dialector yet (tests that inject a stub DBService).
func gooseDialectFromGORM(gormDB *gorm.DB, ctx *core.Context) goose.Dialect {
if gormDB != nil && gormDB.Dialector != nil && gormDB.Dialector.Name() == "postgres" {
return goose.DialectPostgres
}
if gormDB != nil && gormDB.Dialector != nil && gormDB.Dialector.Name() == "sqlite" {
return goose.DialectSQLite3
}
return gooseDialect(ctx)
}
// gooseDialect returns the goose dialect based on the configured database engine.
func gooseDialect(ctx *core.Context) goose.Dialect {
if ctx != nil && ctx.Config() != nil && ctx.Config().Bool("database.enabled", false) {
+277 -4
View File
@@ -7,14 +7,20 @@ import (
"Wavelet/core"
"Wavelet/core/contracts"
"context"
"fmt"
"os"
"path/filepath"
"strings"
"testing"
"testing/fstest"
"time"
"github.com/glebarez/sqlite"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"gorm.io/driver/postgres"
"gorm.io/gorm"
gormlogger "gorm.io/gorm/logger"
)
type migrateTestDB struct {
@@ -28,15 +34,21 @@ func (s migrateTestDB) DB(ctx context.Context) *gorm.DB { return s.db.WithContex
func (s migrateTestDB) Named(string) *gorm.DB { return s.db }
type migrateTestPlugin struct {
db *gorm.DB
fs fstest.MapFS
name string
db *gorm.DB
fs fstest.MapFS
}
func (p *migrateTestPlugin) Name() string { return "t" }
func (p *migrateTestPlugin) Name() string {
if p.name != "" {
return p.name
}
return "t"
}
func (p *migrateTestPlugin) Apply(ctx *core.Context) error {
core.Provide[contracts.DBService](ctx, migrateTestDB{db: p.db})
ctx.Migrations().Register("t", p.fs)
ctx.Migrations().Register(p.Name(), p.fs)
return nil
}
@@ -125,3 +137,264 @@ func TestGooseEngineNilBaselineStillMigrates(t *testing.T) {
assert.True(t, sqliteTableExists(t, gdb, "w_schema_versions"))
assert.True(t, sqliteTableExists(t, gdb, "t_up"))
}
func TestGooseEngineUpgradesFrom00001To00002(t *testing.T) {
gdb := openMigrateTestDB(t)
runTestMigrations(t, gdb, testMigrationFS(), "")
require.True(t, sqliteTableExists(t, gdb, "t_up"))
require.False(t, sqliteTableExists(t, gdb, "t_v2"))
require.Equal(t, int64(1), pluginSchemaVersion(t, gdb, "t"))
runTestMigrations(t, gdb, testMigrationFSWithV2("sqlite"), "")
require.True(t, sqliteTableExists(t, gdb, "t_up"), "00001 table must survive 00002")
require.True(t, sqliteTableExists(t, gdb, "t_v2"), "00002 must create t_v2")
require.Equal(t, int64(2), pluginSchemaVersion(t, gdb, "t"))
require.Equal(t, 1, tableRowCount(t, gdb, "t_v2"))
runTestMigrations(t, gdb, testMigrationFSWithV2("sqlite"), "")
require.Equal(t, int64(2), pluginSchemaVersion(t, gdb, "t"), "second 00002 run must be a no-op")
require.Equal(t, 1, tableRowCount(t, gdb, "t_v2"), "00002 INSERT must not run twice")
}
func TestGooseEngineStampedV1AppliesOnly00002(t *testing.T) {
gdb := openMigrateTestDB(t)
require.NoError(t, gdb.Exec(`CREATE TABLE w_schema_versions (
plugin_id VARCHAR(64) NOT NULL,
version_id BIGINT NOT NULL,
applied_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (plugin_id, version_id)
)`).Error)
require.NoError(t, gdb.Exec(`INSERT INTO w_schema_versions (plugin_id, version_id) VALUES ('t', 1)`).Error)
runTestMigrations(t, gdb, testMigrationFSWithV2("sqlite"), "")
require.False(t, sqliteTableExists(t, gdb, "t_up"), "stamped v1 must not re-run 00001")
require.True(t, sqliteTableExists(t, gdb, "t_v2"), "stamped v1 must still apply 00002")
require.Equal(t, int64(2), pluginSchemaVersion(t, gdb, "t"))
}
func TestOpenFlareServerUpgradesFrom00001To00002(t *testing.T) {
dbPath := filepath.Join(t.TempDir(), "of.db")
app := cordisPrepare(t, cordisSQLiteSource(t, dbPath))
require.NoError(t, app.Context().Dispose())
inspect := openInspectDB(t, dbPath, "")
if !pluginHasVersion(t, inspect, false, serverPluginStamp, 1) {
t.Fatal("fresh install did not apply server 00001")
}
_ = inspect.Close()
gdb, err := gorm.Open(sqlite.Open(dbPath), &gorm.Config{Logger: gormlogger.Default.LogMode(gormlogger.Silent)})
require.NoError(t, err)
runTestMigrations(t, gdb, serverFollowupFS("sqlite"), "server")
require.False(t, sqliteTableExists(t, gdb, "should_not_exist_from_00001_rerun"))
require.True(t, sqliteTableExists(t, gdb, "of_upgrade_probe"))
require.Equal(t, int64(2), pluginSchemaVersion(t, gdb, "server"))
require.True(t, sqliteTableExists(t, gdb, "of_zones"), "existing of_* tables must survive 00002")
require.Equal(t, 1, tableRowCount(t, gdb, "of_upgrade_probe"))
}
func TestGooseEngineUpgradesFrom00001To00002Postgres(t *testing.T) {
gdb := openMigratePostgresDB(t)
opts := postgresMigrateOpt()
runTestMigrations(t, gdb, testPostgresMigrationFS(), "", opts)
require.True(t, pgTableExists(t, gdb, "t_up"))
require.False(t, pgTableExists(t, gdb, "t_v2"))
require.Equal(t, int64(1), pluginSchemaVersion(t, gdb, "t"))
runTestMigrations(t, gdb, testMigrationFSWithV2("postgres"), "", opts)
require.True(t, pgTableExists(t, gdb, "t_up"))
require.True(t, pgTableExists(t, gdb, "t_v2"))
require.Equal(t, int64(2), pluginSchemaVersion(t, gdb, "t"))
require.Equal(t, 1, tableRowCount(t, gdb, "t_v2"))
runTestMigrations(t, gdb, testMigrationFSWithV2("postgres"), "", opts)
require.Equal(t, int64(2), pluginSchemaVersion(t, gdb, "t"))
require.Equal(t, 1, tableRowCount(t, gdb, "t_v2"))
}
func TestOpenFlareServerUpgradesFrom00001To00002Postgres(t *testing.T) {
host, port, user, pass, dbName, sslMode, cleanup := createMigratePostgresDB(t)
t.Cleanup(cleanup)
dsn := postgresDSN(host, port, user, pass, dbName, sslMode)
app := cordisPrepare(t, cordisPostgresSource(t, host, port, user, pass, dbName, sslMode))
require.NoError(t, app.Context().Dispose())
inspect := openInspectDB(t, "", dsn)
if !pluginHasVersion(t, inspect, true, serverPluginStamp, 1) {
t.Fatal("fresh install did not apply server 00001")
}
_ = inspect.Close()
gdb, err := gorm.Open(postgres.Open(dsn), &gorm.Config{Logger: gormlogger.Default.LogMode(gormlogger.Silent)})
require.NoError(t, err)
runTestMigrations(t, gdb, serverFollowupFS("postgres"), "server", postgresMigrateOpt())
require.False(t, pgTableExists(t, gdb, "should_not_exist_from_00001_rerun"))
require.True(t, pgTableExists(t, gdb, "of_upgrade_probe"))
require.Equal(t, int64(2), pluginSchemaVersion(t, gdb, "server"))
require.True(t, pgTableExists(t, gdb, "of_zones"))
require.Equal(t, 1, tableRowCount(t, gdb, "of_upgrade_probe"))
}
func runTestMigrations(t *testing.T, gdb *gorm.DB, fs fstest.MapFS, pluginName string, opts ...core.AppOption) {
t.Helper()
plugin := &migrateTestPlugin{name: pluginName, db: gdb, fs: fs}
appOpts := []core.AppOption{
core.WithMigrationEngine(&gooseEngine{}),
core.WithPlugins(plugin),
}
appOpts = append(appOpts, opts...)
app := core.NewApp(appOpts...)
require.NoError(t, app.Prepare())
require.NoError(t, app.ApplyPlugins())
require.NoError(t, app.RunMigrations())
}
func postgresMigrateOpt() core.AppOption {
return core.WithConfigSource(core.NewMapSource(map[string]any{
"database": map[string]any{"enabled": true},
}))
}
func testPostgresMigrationFS() fstest.MapFS {
return fstest.MapFS{
"migrations/postgres/00001_init.sql": &fstest.MapFile{Data: []byte(`-- +goose Up
CREATE TABLE t_up (id BIGINT PRIMARY KEY);
-- +goose Down
DROP TABLE t_up;
`)},
}
}
func testMigrationFSWithV2(dialect string) fstest.MapFS {
v1 := `-- +goose Up
CREATE TABLE t_up (id BIGINT PRIMARY KEY);
-- +goose Down
DROP TABLE t_up;
`
v2 := `-- +goose Up
CREATE TABLE t_v2 (id BIGINT PRIMARY KEY, note TEXT NOT NULL DEFAULT '');
INSERT INTO t_v2 (id, note) VALUES (1, 'from-00002');
-- +goose Down
DROP TABLE t_v2;
`
if dialect == "sqlite" {
v1 = `-- +goose Up
CREATE TABLE t_up (id INTEGER PRIMARY KEY);
-- +goose Down
DROP TABLE t_up;
`
v2 = `-- +goose Up
CREATE TABLE t_v2 (id INTEGER PRIMARY KEY, note TEXT NOT NULL DEFAULT '');
INSERT INTO t_v2 (id, note) VALUES (1, 'from-00002');
-- +goose Down
DROP TABLE t_v2;
`
}
return fstest.MapFS{
"migrations/" + dialect + "/00001_init.sql": &fstest.MapFile{Data: []byte(v1)},
"migrations/" + dialect + "/00002_add_t_v2.sql": &fstest.MapFile{Data: []byte(v2)},
}
}
func serverFollowupFS(dialect string) fstest.MapFS {
v1 := `-- +goose Up
CREATE TABLE should_not_exist_from_00001_rerun (id INTEGER);
-- +goose Down
DROP TABLE should_not_exist_from_00001_rerun;
`
v2 := `-- +goose Up
CREATE TABLE of_upgrade_probe (id INTEGER PRIMARY KEY, note TEXT NOT NULL DEFAULT '');
INSERT INTO of_upgrade_probe (id, note) VALUES (1, 'from-00002');
-- +goose Down
DROP TABLE of_upgrade_probe;
`
if dialect == "postgres" {
v1 = `-- +goose Up
CREATE TABLE should_not_exist_from_00001_rerun (id BIGINT);
-- +goose Down
DROP TABLE should_not_exist_from_00001_rerun;
`
v2 = `-- +goose Up
CREATE TABLE of_upgrade_probe (id BIGINT PRIMARY KEY, note TEXT NOT NULL DEFAULT '');
INSERT INTO of_upgrade_probe (id, note) VALUES (1, 'from-00002');
-- +goose Down
DROP TABLE of_upgrade_probe;
`
}
return fstest.MapFS{
"migrations/" + dialect + "/00001_initial.sql": &fstest.MapFile{Data: []byte(v1)},
"migrations/" + dialect + "/00002_upgrade_probe.sql": &fstest.MapFile{Data: []byte(v2)},
}
}
func pluginSchemaVersion(t *testing.T, db *gorm.DB, pluginID string) int64 {
t.Helper()
var v int64
err := db.Raw(`SELECT COALESCE(MAX(version_id), 0) FROM w_schema_versions WHERE plugin_id = ?`, pluginID).Scan(&v).Error
require.NoError(t, err)
return v
}
func tableRowCount(t *testing.T, db *gorm.DB, name string) int {
t.Helper()
if !safePGIdent(name) {
t.Fatalf("unsafe table name %q", name)
}
var n int
err := db.Raw("SELECT COUNT(*) FROM " + name).Scan(&n).Error
require.NoError(t, err)
return n
}
func pgTableExists(t *testing.T, db *gorm.DB, name string) bool {
t.Helper()
var n int
err := db.Raw(`SELECT COUNT(*) FROM information_schema.tables WHERE table_schema = 'public' AND table_name = ?`, name).Scan(&n).Error
require.NoError(t, err)
return n > 0
}
func openMigratePostgresDB(t *testing.T) *gorm.DB {
t.Helper()
host, port, user, pass, dbName, sslMode, cleanup := createMigratePostgresDB(t)
t.Cleanup(cleanup)
dsn := postgresDSN(host, port, user, pass, dbName, sslMode)
gdb, err := gorm.Open(postgres.Open(dsn), &gorm.Config{Logger: gormlogger.Default.LogMode(gormlogger.Silent)})
require.NoError(t, err)
return gdb
}
func createMigratePostgresDB(t *testing.T) (host string, port int, user, pass, dbName, sslMode string, cleanup func()) {
t.Helper()
dsn := strings.TrimSpace(os.Getenv("TEST_PG_DSN"))
if dsn == "" {
t.Skip("TEST_PG_DSN is not set")
}
host, port, user, pass, adminDB, sslMode := parsePostgresDSN(t, dsn)
adminDSN := postgresDSN(host, port, user, pass, adminDB, sslMode)
admin := openInspectDB(t, "", adminDSN)
dbName = fmt.Sprintf("of_mig_%d", time.Now().UnixNano())
if !safePGIdent(dbName) {
t.Fatalf("generated database name %q is not a safe identifier", dbName)
}
if _, err := admin.Exec("CREATE DATABASE " + dbName); err != nil {
t.Fatalf("CREATE DATABASE %s: %v", dbName, err)
}
cleanup = func() {
_, _ = admin.Exec(`SELECT pg_terminate_backend(pid) FROM pg_stat_activity WHERE datname = $1 AND pid <> pg_backend_pid()`, dbName)
_, _ = admin.Exec("DROP DATABASE IF EXISTS " + dbName)
_ = admin.Close()
}
return host, port, user, pass, dbName, sslMode, cleanup
}
+91
View File
@@ -94,6 +94,97 @@ func TestUpgradePostgresFromGolden(t *testing.T) {
})
}
func TestUpgradePostgresFromExistingDump(t *testing.T) {
dsn := strings.TrimSpace(os.Getenv("TEST_PG_EXISTING_DSN"))
if dsn == "" {
t.Skip("TEST_PG_EXISTING_DSN is not set")
}
host, port, user, pass, dbName, sslMode := parsePostgresDSN(t, dsn)
spec := upgradeDB{
pgDSN: dsn,
source: cordisPostgresSource(t, host, port, user, pass, dbName, sslMode),
}
inspect := openInspectDB(t, "", spec.pgDSN)
beforeCounts := countNamedTables(t, inspect, productionCountTables)
beforeTables := listPublicTables(t, inspect)
_ = inspect.Close()
assertUpgradeFromGolden(t, spec)
inspect = openInspectDB(t, "", spec.pgDSN)
defer func() { _ = inspect.Close() }()
afterCounts := countNamedTables(t, inspect, productionCountTables)
for _, name := range productionCountTables {
if afterCounts[name] < beforeCounts[name] {
t.Errorf("row count dropped for %s: before %d after %d", name, beforeCounts[name], afterCounts[name])
}
}
afterTables := listPublicTables(t, inspect)
for name := range beforeTables {
if !afterTables[name] {
t.Errorf("table %s dropped", name)
}
}
for _, name := range []string{"w_schema_versions", "w_message_channels", "w_message_bindings", "w_message_pairing_codes"} {
if !afterTables[name] {
t.Errorf("expected upgrade to create %s", name)
}
}
var n int
if err := inspect.QueryRow(`SELECT COUNT(*) FROM pg_inherits i JOIN pg_class c ON c.oid = i.inhparent WHERE c.relname IN ('of_node_access_logs', 'w_user_access_logs')`).Scan(&n); err != nil {
t.Fatalf("count partitions: %v", err)
}
if n < 8 {
t.Errorf("partition children = %d, want at least 8", n)
}
}
var productionCountTables = []string{
"of_zones", "of_zone_domains", "of_proxy_routes", "of_nodes", "of_origins",
"of_tls_certificates", "of_waf_rule_groups", "of_pages_projects",
"w_users", "w_schedules", "w_system_configs", "w_templates", "w_uploads",
"of_node_access_logs", "w_user_access_logs",
}
func countNamedTables(t *testing.T, db *sql.DB, tables []string) map[string]int {
t.Helper()
out := make(map[string]int, len(tables))
for _, name := range tables {
if !safePGIdent(name) {
t.Fatalf("unsafe table name %q", name)
}
var n int
if err := db.QueryRow("SELECT COUNT(*) FROM " + name).Scan(&n); err != nil {
t.Fatalf("count %s: %v", name, err)
}
out[name] = n
}
return out
}
func listPublicTables(t *testing.T, db *sql.DB) map[string]bool {
t.Helper()
rows, err := db.Query(`SELECT tablename FROM pg_tables WHERE schemaname = 'public'`)
if err != nil {
t.Fatalf("list public tables: %v", err)
}
defer func() { _ = rows.Close() }()
out := make(map[string]bool)
for rows.Next() {
var name string
if err := rows.Scan(&name); err != nil {
t.Fatalf("scan table name: %v", err)
}
out[name] = true
}
if err := rows.Err(); err != nil {
t.Fatalf("list public tables: %v", err)
}
return out
}
type upgradeDB struct {
sqlitePath string
pgDSN string
@@ -499,13 +499,10 @@ INSERT INTO w_schedules (id, name, task_type, cron, payload, is_active, created_
VALUES
(101, 'OpenFlare SSL 自动续期', 'of_ssl_renew', '0 0 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
(103, 'OpenFlare WAF IP 组同步', 'of_waf_ip_group_sync', '*/5 * * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
(104, 'OpenFlare Uptime Kuma 同步', 'of_uptime_kuma_sync', '* * * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
(104, 'OpenFlare Uptime Kuma 同步', 'of_uptime_kuma_sync', '* * * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
(105, 'OpenFlare Pages 部署源扫描', 'of_pages_source_scan', '0 0 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
ON CONFLICT (id) DO NOTHING;
INSERT INTO w_schedules (name, task_type, cron, payload, is_active, created_at, updated_at)
SELECT 'OpenFlare Pages 部署源扫描', 'of_pages_source_scan', '0 0 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
WHERE NOT EXISTS (SELECT 1 FROM w_schedules WHERE task_type = 'of_pages_source_scan');
INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
VALUES
('agent_heartbeat_interval', '3000', 'business', 0, 'Agent 心跳间隔(毫秒)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
@@ -485,11 +485,8 @@ INSERT OR IGNORE INTO w_schedules (id, name, task_type, cron, payload, is_active
VALUES
(101, 'OpenFlare SSL 自动续期', 'of_ssl_renew', '0 0 * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
(103, 'OpenFlare WAF IP 组同步', 'of_waf_ip_group_sync', '*/5 * * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
(104, 'OpenFlare Uptime Kuma 同步', 'of_uptime_kuma_sync', '* * * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP);
INSERT INTO w_schedules (name, task_type, cron, payload, is_active, created_at, updated_at)
SELECT 'OpenFlare Pages 部署源扫描', 'of_pages_source_scan', '0 0 * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP
WHERE NOT EXISTS (SELECT 1 FROM w_schedules WHERE task_type = 'of_pages_source_scan');
(104, 'OpenFlare Uptime Kuma 同步', 'of_uptime_kuma_sync', '* * * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP),
(105, 'OpenFlare Pages 部署源扫描', 'of_pages_source_scan', '0 0 * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP);
INSERT OR IGNORE INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at)
VALUES
@@ -26,12 +26,12 @@ 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', '*', '系统', '', TRUE, FALSE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
ON CONFLICT (username) DO NOTHING;
ON CONFLICT DO NOTHING;
-- Seed default administrator user (username: admin, password: 12345678)
INSERT INTO w_users (id, username, password, nickname, email, is_active, is_admin, last_login_at, created_at, updated_at)
VALUES (1, 'admin', '12345678', '管理员', 'admin@wavelet.local', TRUE, TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
ON CONFLICT (username) DO NOTHING;
ON CONFLICT DO NOTHING;
-- +goose StatementEnd
-- +goose Down
@@ -26,12 +26,12 @@ 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;
ON CONFLICT DO NOTHING;
-- Seed default administrator user (username: admin, password: 12345678)
INSERT INTO w_users (id, username, password, nickname, email, is_active, is_admin, last_login_at, created_at, updated_at)
VALUES (1, 'admin', '12345678', '管理员', 'admin@wavelet.local', 1, 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP)
ON CONFLICT (username) DO NOTHING;
ON CONFLICT DO NOTHING;
-- +goose StatementEnd
-- +goose Down