2025-02-27 11:49:12 +01:00
|
|
|
package queue
|
|
|
|
|
|
|
|
|
|
import (
|
|
|
|
|
"context"
|
2025-09-01 08:21:10 +03:00
|
|
|
"database/sql"
|
2025-02-27 11:49:12 +01:00
|
|
|
|
|
|
|
|
"github.com/riverqueue/river/riverdriver"
|
2025-09-01 08:21:10 +03:00
|
|
|
"github.com/riverqueue/river/riverdriver/riverdatabasesql"
|
2025-02-27 11:49:12 +01:00
|
|
|
"github.com/riverqueue/river/rivermigrate"
|
|
|
|
|
|
|
|
|
|
"github.com/zitadel/zitadel/internal/database"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
type Migrator struct {
|
2025-09-01 08:21:10 +03:00
|
|
|
driver riverdriver.Driver[*sql.Tx]
|
2025-02-27 11:49:12 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func NewMigrator(client *database.DB) *Migrator {
|
|
|
|
|
return &Migrator{
|
2025-09-01 08:21:10 +03:00
|
|
|
driver: riverdatabasesql.New(client.DB),
|
2025-02-27 11:49:12 +01:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func (m *Migrator) Execute(ctx context.Context) error {
|
2025-09-01 08:21:10 +03:00
|
|
|
err := m.driver.GetExecutor().Exec(ctx, "CREATE SCHEMA IF NOT EXISTS "+schema)
|
2025-02-27 11:49:12 +01:00
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
2025-07-29 09:09:00 +02:00
|
|
|
migrator, err := rivermigrate.New(m.driver, &rivermigrate.Config{Schema: schema})
|
2025-02-27 11:49:12 +01:00
|
|
|
if err != nil {
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
_, err = migrator.Migrate(ctx, rivermigrate.DirectionUp, nil)
|
|
|
|
|
return err
|
|
|
|
|
|
|
|
|
|
}
|