diff --git a/db/schema/schema.go b/db/schema/schema.go index 6fa0b226..d7e13f90 100644 --- a/db/schema/schema.go +++ b/db/schema/schema.go @@ -10,6 +10,11 @@ import ( "fmt" "heckel.io/ntfy/v2/db/pg" + "heckel.io/ntfy/v2/log" +) + +const ( + tag = "schema" ) const ( @@ -68,6 +73,7 @@ func Migrate(db *sql.DB, dialect Dialect, store string, targetVersion int, creat if !ok { return fmt.Errorf("cannot find %s migration step from version %d to %d", store, v, v+1) } + log.Tag(tag).Info("Migrating %s database schema: from %d to %d", store, v, v+1) if err := migrate(tx); err != nil { return fmt.Errorf("%s migration step from version %d to %d failed: %w", store, v, v+1, err) } diff --git a/db/schema/types.go b/db/schema/types.go index 430a02cd..4294feea 100644 --- a/db/schema/types.go +++ b/db/schema/types.go @@ -23,3 +23,9 @@ func AsMigrateFunc(query string) MigrateFunc { return err } } + +// NopMigrateFunc is a migration step that does nothing, for versions where a dialect has no +// work to do (e.g. when only the other dialect's schema changed). +func NopMigrateFunc(_ *sql.Tx) error { + return nil +} diff --git a/docs/releases.md b/docs/releases.md index 2f0e4585..925d1c77 100644 --- a/docs/releases.md +++ b/docs/releases.md @@ -2065,6 +2065,7 @@ and the [ntfy Android app](https://github.com/binwiederhier/ntfy-android/release * Fix Twilio phone calls and phone number verifications failing silently when Twilio rejected the request, and move the Twilio integration into its own `twilio` package * Move the Prometheus metrics into a dedicated `metrics` package +* Message cache databases from ntfy older than v1.10.0 (November 2021) can no longer be migrated; upgrade via an older ntfy version first, or delete the cache database ### ntfy iOS app v1.8.0 (UNRELEASED) diff --git a/message/cache.go b/message/cache.go index 90fbf51d..73aaa076 100644 --- a/message/cache.go +++ b/message/cache.go @@ -17,6 +17,7 @@ import ( const ( tagMessageCache = "message_cache" + schemaStore = "message" // Store name in the schema_version table (see db/schema) ) var errNoRows = errors.New("no rows found") diff --git a/message/cache_postgres.go b/message/cache_postgres.go index 4d7c3f93..e588f8c4 100644 --- a/message/cache_postgres.go +++ b/message/cache_postgres.go @@ -4,6 +4,7 @@ import ( "time" "heckel.io/ntfy/v2/db" + "heckel.io/ntfy/v2/db/schema" ) // PostgreSQL runtime query constants @@ -102,7 +103,7 @@ var postgresQueries = queries{ // NewPostgresStore creates a new PostgreSQL-backed message cache store using an existing database connection pool. func NewPostgresStore(d *db.DB, batchSize int, batchTimeout time.Duration) (*Cache, error) { - if err := setupPostgres(d.Primary()); err != nil { + if err := schema.Migrate(d.Primary(), schema.Postgres, schemaStore, postgresCurrentSchemaVersion, postgresCreateTables, postgresMigrations); err != nil { return nil, err } return newCache(d, postgresQueries, nil, batchSize, batchTimeout, false), nil diff --git a/message/cache_postgres_schema.go b/message/cache_postgres_schema.go index 994df9b0..c52e9f87 100644 --- a/message/cache_postgres_schema.go +++ b/message/cache_postgres_schema.go @@ -1,16 +1,13 @@ package message import ( - "database/sql" - "fmt" - - "heckel.io/ntfy/v2/db" - "heckel.io/ntfy/v2/log" + "heckel.io/ntfy/v2/db/schema" ) // Initial PostgreSQL schema const ( - postgresCreateTablesQuery = ` + postgresCurrentSchemaVersion = 15 + postgresCreateTablesQuery = ` CREATE TABLE IF NOT EXISTS message ( id BIGSERIAL PRIMARY KEY, mid TEXT NOT NULL, @@ -50,21 +47,9 @@ const ( value BIGINT ); INSERT INTO message_stats (key, value) VALUES ('messages', 0); - CREATE TABLE IF NOT EXISTS schema_version ( - store TEXT PRIMARY KEY, - version INT NOT NULL - ); ` ) -// PostgreSQL schema management queries -const ( - postgresCurrentSchemaVersion = 15 - postgresInsertSchemaVersionQuery = `INSERT INTO schema_version (store, version) VALUES ('message', $1)` - postgresUpdateSchemaVersionQuery = `UPDATE schema_version SET version = $1 WHERE store = 'message'` - postgresSelectSchemaVersionQuery = `SELECT version FROM schema_version WHERE store = 'message'` -) - // PostgreSQL schema migrations const ( // 14 -> 15 @@ -73,51 +58,12 @@ const ( ` ) -var postgresMigrations = map[int]func(d *sql.DB) error{ - 14: postgresMigrateFrom14, -} +var ( + postgresCreateTables = schema.AsMigrateFunc(postgresCreateTablesQuery) -func setupPostgres(d *sql.DB) error { - var schemaVersion int - if err := d.QueryRow(postgresSelectSchemaVersionQuery).Scan(&schemaVersion); err != nil { - return setupNewPostgresDB(d) - } else if schemaVersion == postgresCurrentSchemaVersion { - return nil - } else if schemaVersion > postgresCurrentSchemaVersion { - return fmt.Errorf("unexpected schema version: version %d is higher than current version %d", schemaVersion, postgresCurrentSchemaVersion) + // postgresMigrations maps a schema version to the migration upgrading it to the next + // version. Always append migrations at the end, never insert in the middle. + postgresMigrations = map[int]schema.MigrateFunc{ + 14: schema.AsMigrateFunc(postgresMigrate14To15CreateIndexQuery), } - for i := schemaVersion; i < postgresCurrentSchemaVersion; i++ { - fn, ok := postgresMigrations[i] - if !ok { - return fmt.Errorf("cannot find migration step from schema version %d to %d", i, i+1) - } else if err := fn(d); err != nil { - return err - } - } - return nil -} - -func postgresMigrateFrom14(d *sql.DB) error { - log.Tag(tagMessageCache).Info("Migrating message cache database schema: from 14 to 15") - return db.ExecTx(d, func(tx *sql.Tx) error { - if _, err := tx.Exec(postgresMigrate14To15CreateIndexQuery); err != nil { - return err - } - if _, err := tx.Exec(postgresUpdateSchemaVersionQuery, 15); err != nil { - return err - } - return nil - }) -} - -func setupNewPostgresDB(sqlDB *sql.DB) error { - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(postgresCreateTablesQuery); err != nil { - return err - } - if _, err := tx.Exec(postgresInsertSchemaVersionQuery, postgresCurrentSchemaVersion); err != nil { - return err - } - return nil - }) -} +) diff --git a/message/cache_sqlite.go b/message/cache_sqlite.go index b9d7394f..cc01b537 100644 --- a/message/cache_sqlite.go +++ b/message/cache_sqlite.go @@ -9,6 +9,7 @@ import ( _ "github.com/mattn/go-sqlite3" // SQLite driver "heckel.io/ntfy/v2/db" + "heckel.io/ntfy/v2/db/schema" "heckel.io/ntfy/v2/util" ) @@ -113,7 +114,10 @@ func NewSQLiteStore(filename, startupQueries string, cacheDuration time.Duration if err != nil { return nil, err } - if err := setupSQLite(d, startupQueries, cacheDuration); err != nil { + if err := runSQLiteStartupQueries(d, startupQueries); err != nil { + return nil, err + } + if err := schema.Migrate(d, schema.SQLite, schemaStore, sqliteCurrentSchemaVersion, sqliteCreateTables, sqliteMigrations(cacheDuration)); err != nil { return nil, err } return newCache(db.New(&db.Host{DB: d}, nil), sqliteQueries, &sync.Mutex{}, batchSize, batchTimeout, nop), nil diff --git a/message/cache_sqlite_schema.go b/message/cache_sqlite_schema.go index b19bfca1..12c6f937 100644 --- a/message/cache_sqlite_schema.go +++ b/message/cache_sqlite_schema.go @@ -2,16 +2,15 @@ package message import ( "database/sql" - "fmt" "time" - "heckel.io/ntfy/v2/db" - "heckel.io/ntfy/v2/log" + "heckel.io/ntfy/v2/db/schema" ) // Initial SQLite schema const ( - sqliteCreateTablesQuery = ` + sqliteCurrentSchemaVersion = 15 + sqliteCreateTablesQuery = ` CREATE TABLE IF NOT EXISTS messages ( id INTEGER PRIMARY KEY AUTOINCREMENT, mid TEXT NOT NULL, @@ -55,29 +54,9 @@ const ( ` ) -// Schema version management for SQLite +// Schema migrations for SQLite. Databases older than schema version 1 (ntfy < v1.10.0, +// November 2021) can no longer be migrated. const ( - sqliteCurrentSchemaVersion = 15 - sqliteCreateSchemaVersionTableQuery = ` - CREATE TABLE IF NOT EXISTS schemaVersion ( - id INT PRIMARY KEY, - version INT NOT NULL - ); - ` - sqliteInsertSchemaVersionQuery = `INSERT INTO schemaVersion VALUES (1, ?)` - sqliteUpdateSchemaVersionQuery = `UPDATE schemaVersion SET version = ? WHERE id = 1` - sqliteSelectSchemaVersionQuery = `SELECT version FROM schemaVersion WHERE id = 1` -) - -// Schema migrations for SQLite -const ( - // 0 -> 1 - sqliteMigrate0To1AlterMessagesTableQuery = ` - ALTER TABLE messages ADD COLUMN title TEXT NOT NULL DEFAULT(''); - ALTER TABLE messages ADD COLUMN priority INT NOT NULL DEFAULT(0); - ALTER TABLE messages ADD COLUMN tags TEXT NOT NULL DEFAULT(''); - ` - // 1 -> 2 sqliteMigrate1To2AlterMessagesTableQuery = ` ALTER TABLE messages ADD COLUMN published INT NOT NULL DEFAULT(1); @@ -193,67 +172,35 @@ const ( ) var ( - sqliteMigrations = map[int]func(db *sql.DB, cacheDuration time.Duration) error{ - 0: sqliteMigrateFrom0, - 1: sqliteMigrateFrom1, - 2: sqliteMigrateFrom2, - 3: sqliteMigrateFrom3, - 4: sqliteMigrateFrom4, - 5: sqliteMigrateFrom5, - 6: sqliteMigrateFrom6, - 7: sqliteMigrateFrom7, - 8: sqliteMigrateFrom8, - 9: sqliteMigrateFrom9, - 10: sqliteMigrateFrom10, - 11: sqliteMigrateFrom11, - 12: sqliteMigrateFrom12, - 13: sqliteMigrateFrom13, - 14: sqliteMigrateFrom14, - } + sqliteCreateTables = schema.AsMigrateFunc(sqliteCreateTablesQuery) ) -func setupSQLite(db *sql.DB, startupQueries string, cacheDuration time.Duration) error { - if err := runSQLiteStartupQueries(db, startupQueries); err != nil { - return err - } - // If 'messages' table does not exist, this must be a new database - var messagesCount int - if err := db.QueryRow(sqliteSelectMessagesCountQuery).Scan(&messagesCount); err != nil { - return setupNewSQLite(db) - } - // If 'messages' table exists (schema >= 0), check 'schemaVersion' table - var schemaVersion int - db.QueryRow(sqliteSelectSchemaVersionQuery).Scan(&schemaVersion) // Error means schema version is zero! - // Do migrations - if schemaVersion == sqliteCurrentSchemaVersion { - return nil - } else if schemaVersion > sqliteCurrentSchemaVersion { - return fmt.Errorf("unexpected schema version: version %d is higher than current version %d", schemaVersion, sqliteCurrentSchemaVersion) - } - for i := schemaVersion; i < sqliteCurrentSchemaVersion; i++ { - fn, ok := sqliteMigrations[i] - if !ok { - return fmt.Errorf("cannot find migration step from schema version %d to %d", i, i+1) - } else if err := fn(db, cacheDuration); err != nil { +// sqliteMigrations returns the migration steps, keyed by the version they upgrade FROM. The +// cache duration is carried into the 9 -> 10 step via closure (it backfills "expires" from it). +// Always append migrations at the end, never insert in the middle. +func sqliteMigrations(cacheDuration time.Duration) map[int]schema.MigrateFunc { + return map[int]schema.MigrateFunc{ + 1: schema.AsMigrateFunc(sqliteMigrate1To2AlterMessagesTableQuery), + 2: schema.AsMigrateFunc(sqliteMigrate2To3AlterMessagesTableQuery), + 3: schema.AsMigrateFunc(sqliteMigrate3To4AlterMessagesTableQuery), + 4: schema.AsMigrateFunc(sqliteMigrate4To5AlterMessagesTableQuery), + 5: schema.AsMigrateFunc(sqliteMigrate5To6AlterMessagesTableQuery), + 6: schema.AsMigrateFunc(sqliteMigrate6To7AlterMessagesTableQuery), + 7: schema.AsMigrateFunc(sqliteMigrate7To8AlterMessagesTableQuery), + 8: schema.AsMigrateFunc(sqliteMigrate8To9AlterMessagesTableQuery), + 9: func(tx *sql.Tx) error { + if _, err := tx.Exec(sqliteMigrate9To10AlterMessagesTableQuery); err != nil { + return err + } + _, err := tx.Exec(sqliteMigrate9To10UpdateMessageExpiryQuery, int64(cacheDuration.Seconds())) return err - } + }, + 10: schema.AsMigrateFunc(sqliteMigrate10To11AlterMessagesTableQuery), + 11: schema.AsMigrateFunc(sqliteMigrate11To12AlterMessagesTableQuery), + 12: schema.AsMigrateFunc(sqliteMigrate12To13AlterMessagesTableQuery), + 13: schema.AsMigrateFunc(sqliteMigrate13To14AlterMessagesTableQuery), + 14: schema.NopMigrateFunc, // Corresponds to Postgres migration } - return nil -} - -func setupNewSQLite(sqlDB *sql.DB) error { - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteCreateTablesQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteCreateSchemaVersionTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteInsertSchemaVersionQuery, sqliteCurrentSchemaVersion); err != nil { - return err - } - return nil - }) } func runSQLiteStartupQueries(db *sql.DB, startupQueries string) error { @@ -264,203 +211,3 @@ func runSQLiteStartupQueries(db *sql.DB, startupQueries string) error { } return nil } - -func sqliteMigrateFrom0(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 0 to 1") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate0To1AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteCreateSchemaVersionTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteInsertSchemaVersionQuery, 1); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom1(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 1 to 2") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate1To2AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 2); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom2(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 2 to 3") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate2To3AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 3); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom3(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 3 to 4") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate3To4AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 4); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom4(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 4 to 5") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate4To5AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 5); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom5(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 5 to 6") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate5To6AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 6); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom6(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 6 to 7") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate6To7AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 7); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom7(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 7 to 8") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate7To8AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 8); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom8(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 8 to 9") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate8To9AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 9); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom9(sqlDB *sql.DB, cacheDuration time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 9 to 10") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate9To10AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteMigrate9To10UpdateMessageExpiryQuery, int64(cacheDuration.Seconds())); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 10); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom10(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 10 to 11") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate10To11AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 11); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom11(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 11 to 12") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate11To12AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 12); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom12(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 12 to 13") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate12To13AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 13); err != nil { - return err - } - return nil - }) -} - -func sqliteMigrateFrom13(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 13 to 14") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteMigrate13To14AlterMessagesTableQuery); err != nil { - return err - } - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 14); err != nil { - return err - } - return nil - }) -} - -// sqliteMigrateFrom14 is a no-op; the corresponding Postgres migration adds -// idx_message_attachment_expires, which SQLite already has from the initial schema. -func sqliteMigrateFrom14(sqlDB *sql.DB, _ time.Duration) error { - log.Tag(tagMessageCache).Info("Migrating cache database schema: from 14 to 15") - return db.ExecTx(sqlDB, func(tx *sql.Tx) error { - if _, err := tx.Exec(sqliteUpdateSchemaVersionQuery, 15); err != nil { - return err - } - return nil - }) -} diff --git a/message/cache_sqlite_test.go b/message/cache_sqlite_test.go index 95ff7e48..47519163 100644 --- a/message/cache_sqlite_test.go +++ b/message/cache_sqlite_test.go @@ -13,46 +13,6 @@ import ( "heckel.io/ntfy/v2/model" ) -func TestSqliteStore_Migration_From0(t *testing.T) { - filename := newSqliteTestStoreFile(t) - db, err := sql.Open("sqlite3", filename) - require.Nil(t, err) - - // Create "version 0" schema - _, err = db.Exec(` - BEGIN; - CREATE TABLE IF NOT EXISTS messages ( - id VARCHAR(20) PRIMARY KEY, - time INT NOT NULL, - topic VARCHAR(64) NOT NULL, - message VARCHAR(1024) NOT NULL - ); - CREATE INDEX IF NOT EXISTS idx_topic ON messages (topic); - COMMIT; - `) - require.Nil(t, err) - - // Insert a bunch of messages - for i := 0; i < 10; i++ { - _, err = db.Exec(`INSERT INTO messages (id, time, topic, message) VALUES (?, ?, ?, ?)`, - fmt.Sprintf("abcd%d", i), time.Now().Unix(), "mytopic", fmt.Sprintf("some message %d", i)) - require.Nil(t, err) - } - require.Nil(t, db.Close()) - - // Create store to trigger migration - s := newSqliteTestStoreFromFile(t, filename, "") - checkSqliteSchemaVersion(t, filename) - - messages, err := s.Messages("mytopic", model.SinceAllMessages, false) - require.Nil(t, err) - require.Equal(t, 10, len(messages)) - require.Equal(t, "some message 5", messages[5].Message) - require.Equal(t, "", messages[5].Title) - require.Nil(t, messages[5].Tags) - require.Equal(t, 0, messages[5].Priority) -} - func TestSqliteStore_Migration_From1(t *testing.T) { filename := newSqliteTestStoreFile(t) db, err := sql.Open("sqlite3", filename) diff --git a/message/cache_test.go b/message/cache_test.go index 059a1f62..04838abd 100644 --- a/message/cache_test.go +++ b/message/cache_test.go @@ -36,6 +36,60 @@ func newTestPostgresStore(t *testing.T) *message.Cache { return store } +func TestPostgresStore_Migration_From14(t *testing.T) { + // A pre-framework database at version 14: full v14 schema, version tracked in the + // hand-rolled schema_version table, and no idx_message_attachment_expires yet + testDB := dbtest.CreateTestPostgres(t) + _, err := testDB.Exec(` + CREATE TABLE message ( + id BIGSERIAL PRIMARY KEY, + mid TEXT NOT NULL, + sequence_id TEXT NOT NULL, + time BIGINT NOT NULL, + event TEXT NOT NULL, + expires BIGINT NOT NULL, + topic TEXT NOT NULL, + message TEXT NOT NULL, + title TEXT NOT NULL, + priority INT NOT NULL, + tags TEXT NOT NULL, + click TEXT NOT NULL, + icon TEXT NOT NULL, + actions TEXT NOT NULL, + attachment_name TEXT NOT NULL, + attachment_type TEXT NOT NULL, + attachment_size BIGINT NOT NULL, + attachment_expires BIGINT NOT NULL, + attachment_url TEXT NOT NULL, + attachment_deleted BOOLEAN NOT NULL DEFAULT FALSE, + sender TEXT NOT NULL, + user_id TEXT NOT NULL, + content_type TEXT NOT NULL, + encoding TEXT NOT NULL, + published BOOLEAN NOT NULL DEFAULT FALSE + ); + CREATE TABLE message_stats (key TEXT PRIMARY KEY, value BIGINT); + INSERT INTO message_stats (key, value) VALUES ('messages', 0); + CREATE TABLE schema_version (store TEXT PRIMARY KEY, version INT NOT NULL); + INSERT INTO schema_version (store, version) VALUES ('message', 14); + `) + require.Nil(t, err) + store, err := message.NewPostgresStore(testDB, 0, 0) + require.Nil(t, err) + // The 14 -> 15 step ran: version bumped, partial index created + var version int + require.Nil(t, testDB.QueryRow(`SELECT version FROM schema_version WHERE store = 'message'`).Scan(&version)) + require.Equal(t, 15, version) + var indexCount int + require.Nil(t, testDB.QueryRow(`SELECT COUNT(*) FROM pg_indexes WHERE indexname = 'idx_message_attachment_expires' AND schemaname = current_schema()`).Scan(&indexCount)) + require.Equal(t, 1, indexCount) + // And the store works + require.Nil(t, store.AddMessage(model.NewDefaultMessage("mytopic", "hi there"))) + messages, err := store.Messages("mytopic", model.SinceAllMessages, false) + require.Nil(t, err) + require.Len(t, messages, 1) +} + func forEachBackend(t *testing.T, f func(t *testing.T, s *message.Cache)) { t.Run("sqlite", func(t *testing.T) { f(t, newSqliteTestStore(t))