From 5e96490254aa9caee604e9403ef89f63ee00b843 Mon Sep 17 00:00:00 2001 From: tiennm99 Date: Thu, 8 Oct 2026 10:35:21 +0700 Subject: [PATCH] refactor: remove completed one-time startup migrations Production data is migrated and every marker is recorded in the system collection, so the stock dividend-history, sticker legacy-pack and stats stock_dividend migrations no longer do anything. Stats keeps its startup index creation and still hides retired stock_dividend rows. --- cmd/server/main.go | 58 +------ cmd/server/main_test.go | 51 ------ docs/deploy-coolify-selfhosted.md | 10 +- docs/sticker-packs.md | 7 +- internal/modules/stats/startup.go | 137 ++------------- internal/modules/stats/startup_mongo_test.go | 56 ++----- internal/modules/stats/startup_test.go | 98 ----------- internal/modules/stats/usage_store.go | 7 +- internal/modules/sticker/startup.go | 90 ---------- internal/modules/sticker/startup_test.go | 146 ---------------- internal/modules/stock/startup.go | 148 ----------------- internal/modules/stock/startup_mongo_test.go | 88 ---------- internal/modules/stock/startup_test.go | 165 ------------------- 13 files changed, 31 insertions(+), 1030 deletions(-) delete mode 100644 internal/modules/stats/startup_test.go delete mode 100644 internal/modules/sticker/startup.go delete mode 100644 internal/modules/sticker/startup_test.go delete mode 100644 internal/modules/stock/startup.go delete mode 100644 internal/modules/stock/startup_mongo_test.go delete mode 100644 internal/modules/stock/startup_test.go diff --git a/cmd/server/main.go b/cmd/server/main.go index 7184bbf..d813c73 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -42,7 +42,6 @@ import ( "github.com/tiennm99/tiennm99bot/internal/modules/wordle" "github.com/tiennm99/tiennm99bot/internal/server" "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" "github.com/tiennm99/tiennm99bot/internal/telegram" ) @@ -117,10 +116,6 @@ func factories() map[string]modules.Factory { // container; 10s leaves headroom without hiding a wedged cluster. const mongodbInitTimeout = 10 * time.Second -// stockMigrationTimeout bounds the one-time scan of persisted stock -// portfolios without tying it to the shorter MongoDB connection timeout. -const stockMigrationTimeout = 2 * time.Minute - func main() { rootCtx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) defer stop() @@ -145,21 +140,12 @@ func main() { } defer closeProvider() - if err := initStatsStore(rootCtx, provider); err != nil { + if err := stats.InitStore(rootCtx, provider.Collection("stats")); err != nil { log.Fatal("stats storage init failed", "err", err) } if err := lol.InitStore(rootCtx, provider.Collection(lol.CollectionName)); err != nil { log.Fatal("lol storage init failed", "err", err) } - if err := initStickerStore(rootCtx, provider); err != nil { - log.Fatal("sticker storage init failed", "err", err) - } - migrationCtx, cancelMigration := context.WithTimeout(rootCtx, stockMigrationTimeout) - if err := initStockStore(migrationCtx, provider); err != nil { - cancelMigration() - log.Fatal("stock storage init failed", "err", err) - } - cancelMigration() b, err := telegram.NewBot(cfg.TelegramBotToken) if err != nil { @@ -261,48 +247,6 @@ func main() { metrics.Flush() } -func initStockStore(ctx context.Context, provider storage.Provider) error { - return initStockStoreWith(ctx, provider, stock.InitStore) -} - -func initStatsStore(ctx context.Context, provider storage.Provider) error { - return initStatsStoreWith(ctx, provider, stats.InitStore) -} - -type statsStoreInitializer func(context.Context, storage.Collection, storage.Collection) error - -func initStatsStoreWith(ctx context.Context, provider storage.Provider, init statsStoreInitializer) error { - return init( - ctx, - provider.Collection("stats"), - provider.Collection(systemstate.CollectionName), - ) -} - -func initStickerStore(ctx context.Context, provider storage.Provider) error { - return initStickerStoreWith(ctx, provider, sticker.InitStore) -} - -type stickerStoreInitializer func(context.Context, storage.Collection, storage.Collection) error - -func initStickerStoreWith(ctx context.Context, provider storage.Provider, init stickerStoreInitializer) error { - return init( - ctx, - provider.Collection(sticker.CollectionName), - provider.Collection(systemstate.CollectionName), - ) -} - -type stockStoreInitializer func(context.Context, storage.Collection, storage.Collection) error - -func initStockStoreWith(ctx context.Context, provider storage.Provider, init stockStoreInitializer) error { - return init( - ctx, - provider.Collection(stock.CollectionName), - provider.Collection(systemstate.CollectionName), - ) -} - // buildProvider picks the storage backend. Selection order: // 1. Explicit KV_PROVIDER env (memory|mongodb) wins. // 2. Auto-detect: MONGO_URL set → mongodb; otherwise memory. diff --git a/cmd/server/main_test.go b/cmd/server/main_test.go index 96fe3e6..dc2e322 100644 --- a/cmd/server/main_test.go +++ b/cmd/server/main_test.go @@ -1,17 +1,13 @@ package main import ( - "context" - "errors" "os" "path/filepath" "strings" "testing" "github.com/tiennm99/tiennm99bot/internal/modules" - "github.com/tiennm99/tiennm99bot/internal/modules/stock" "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" ) func TestResolveCommitSHA(t *testing.T) { @@ -36,53 +32,6 @@ func TestResolveCommitSHA(t *testing.T) { } } -func TestInitStockStoreRunsStartupMigration(t *testing.T) { - ctx := context.Background() - provider := storage.NewMemoryProvider() - if err := initStockStore(ctx, provider); err != nil { - t.Fatalf("initStockStore: %v", err) - } - marker, exists, err := systemstate.New(provider.Collection(systemstate.CollectionName)).Get(ctx, "migration:stock-dividend-history-v1") - if err != nil || !exists || marker.Status != "completed" { - t.Fatalf("marker=%+v exists=%v err=%v", marker, exists, err) - } - if _, _, err := storage.Typed[stock.Portfolio](provider.Collection(stock.CollectionName)).Get(ctx, "user:1"); !errors.Is(err, storage.ErrNotFound) { - t.Fatalf("unexpected portfolio lookup error: %v", err) - } -} - -func TestInitStatsStoreRunsStartupMigration(t *testing.T) { - ctx := context.Background() - provider := storage.NewMemoryProvider() - if err := initStatsStore(ctx, provider); err != nil { - t.Fatalf("initStatsStore: %v", err) - } - marker, exists, err := systemstate.New(provider.Collection(systemstate.CollectionName)).Get(ctx, "migration:stats-delete-stock-dividend-v1") - if err != nil || !exists || marker.Status != "completed" { - t.Fatalf("marker=%+v exists=%v err=%v", marker, exists, err) - } -} - -func TestInitStockStorePropagatesMigrationError(t *testing.T) { - want := errors.New("migration failed") - err := initStockStoreWith(context.Background(), storage.NewMemoryProvider(), func(context.Context, storage.Collection, storage.Collection) error { - return want - }) - if !errors.Is(err, want) { - t.Fatalf("initStockStoreWith error=%v, want %v", err, want) - } -} - -func TestInitStatsStorePropagatesMigrationError(t *testing.T) { - want := errors.New("migration failed") - err := initStatsStoreWith(context.Background(), storage.NewMemoryProvider(), func(context.Context, storage.Collection, storage.Collection) error { - return want - }) - if !errors.Is(err, want) { - t.Fatalf("initStatsStoreWith error=%v, want %v", err, want) - } -} - func TestShortCommitSHA(t *testing.T) { if got := shortCommitSHA(" 0123456789abcdef "); got != "0123456" { t.Errorf("full revision: got %q, want %q", got, "0123456") diff --git a/docs/deploy-coolify-selfhosted.md b/docs/deploy-coolify-selfhosted.md index e1df256..e822acd 100644 --- a/docs/deploy-coolify-selfhosted.md +++ b/docs/deploy-coolify-selfhosted.md @@ -149,13 +149,9 @@ MP4, with the same text fallback. - **`coin`** stores cash as `usd` and embeds positions as `assets..{quantity,base}`. - **`system`** holds one marker per completed one-time startup migration. Keep - those records as audit history. The current markers are - `migration:stats-delete-stock-dividend-v1` (retires historical - `/stock_dividend` stats rows without erasing them), - `migration:stock-dividend-history-v1` (removes the retired dividend cursor - and hashed applied-event ledger), and - `migration:sticker-drop-legacy-packs-v1` (removes records left by the retired - per-user sticker pack commands). + those records as audit history. No one-time migration runs at startup now; + the completed ones were removed from the code once production data was + verified migrated. ## 2. Coolify diff --git a/docs/sticker-packs.md b/docs/sticker-packs.md index 87989e3..ec1560c 100644 --- a/docs/sticker-packs.md +++ b/docs/sticker-packs.md @@ -6,12 +6,7 @@ storage, and no per-user packs. **The module stores nothing.** Its factory ignores the collection it is handed: the pack comes from the environment and the set owner from `OWNER_ID`, so there -is nothing per-user to key. A one-time startup cleanup -(`migration:sticker-drop-legacy-packs-v1`) removes the records the retired -per-user pack commands left in the `sticker` collection — pack documents keyed -by owner ID, `slug:` name reservations, and `pending-delete:` confirmations. -It is marker-guarded, so it scans once per database and never touches anything -written afterwards. +is nothing per-user to key. | Command | Parameters | Reply to | What it does | |---|---|---|---| diff --git a/internal/modules/stats/startup.go b/internal/modules/stats/startup.go index 47dda85..4136e90 100644 --- a/internal/modules/stats/startup.go +++ b/internal/modules/stats/startup.go @@ -2,146 +2,29 @@ package stats import ( "context" - "errors" "fmt" - "time" "go.mongodb.org/mongo-driver/v2/bson" "go.mongodb.org/mongo-driver/v2/mongo" "go.mongodb.org/mongo-driver/v2/mongo/options" "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" ) -// The retired stock_dividend command's usage rows are soft-deleted (marked -// deleted: true) rather than removed, and the run is recorded under -// deletedStockDividendMarkerKey in the system collection. const ( - statsCommandUsersIndexName = "stats_cmd_n_user" - statsUserCommandsIndexName = "stats_uid_n_cmd" - statsUsernameLookupIndexName = "stats_user_uid" - deletedStockDividendMarkerKey = "migration:stats-delete-stock-dividend-v1" - deletedStockDividendCommand = "stock_dividend" - deletedCommandMigrationRetries = 5 + statsCommandUsersIndexName = "stats_cmd_n_user" + statsUserCommandsIndexName = "stats_uid_n_cmd" + statsUsernameLookupIndexName = "stats_user_uid" ) -// InitStore performs stats collection startup maintenance. MongoDB index -// creation and the retired-command migration are both idempotent and run on -// every boot. -func InitStore(ctx context.Context, statsColl, systemColl storage.Collection) error { - if mongoColl, ok := storage.MongoCollection(statsColl); ok { - if err := ensureUsageIndexes(ctx, mongoColl); err != nil { - return err - } +// InitStore creates the stats collection's MongoDB indexes. Index creation is +// idempotent, so it runs on every boot; the memory backend needs none. +func InitStore(ctx context.Context, statsColl storage.Collection) error { + mongoColl, ok := storage.MongoCollection(statsColl) + if !ok { + return nil } - return markDeletedCommand(ctx, statsColl, systemColl, deletedStockDividendCommand, deletedStockDividendMarkerKey) -} - -// markDeletedCommand soft-deletes every usage row of a retired command and -// records the outcome in the system-state marker under markerKey. -// -// It runs on every boot rather than stopping once the marker exists: an older -// build still running alongside this one (during a rolling deploy, say) can -// write fresh rows for the command after the first run, and each later boot -// sweeps those up. CompletedAt keeps the first run's time; Count and UpdatedAt -// reflect the latest one. -func markDeletedCommand(ctx context.Context, statsColl, systemColl storage.Collection, command, markerKey string) error { - system := systemstate.New(systemColl) - marker, exists, err := system.Get(ctx, markerKey) - if err != nil { - return fmt.Errorf("stats deleted-command migration: read marker: %w", err) - } - var matched int64 - if mongoColl, ok := storage.MongoCollection(statsColl); ok { - matched, err = markMongoUsageEntriesDeleted(ctx, mongoColl, command) - } else { - matched, err = markDocUsageEntriesDeleted(ctx, storage.Typed[usageEntry](statsColl), command) - } - if err != nil { - return err - } - - now := time.Now().UnixMilli() - if !exists { - marker = systemstate.Record{Kind: "migration", Name: "stats delete " + command + " v1"} - } - if marker.CompletedAt == 0 { - marker.CompletedAt = now - } - marker.Status = "completed" - marker.Count = matched - marker.UpdatedAt = now - if err := system.Put(ctx, markerKey, marker); err != nil { - return fmt.Errorf("stats deleted-command migration: write marker: %w", err) - } - return nil -} - -func markMongoUsageEntriesDeleted(ctx context.Context, coll *mongo.Collection, command string) (int64, error) { - filter := bson.M{"cmd": command} - if _, err := coll.UpdateMany(ctx, - bson.M{"cmd": command, "deleted": bson.M{"$ne": true}}, - bson.M{ - "$set": bson.M{"deleted": true}, - "$inc": bson.M{"version": int64(1)}, - "$currentDate": bson.M{"updatedAt": true}, - }, - ); err != nil { - return 0, fmt.Errorf("stats deleted-command migration: update %s: %w", command, err) - } - count, err := coll.CountDocuments(ctx, filter) - if err != nil { - return 0, fmt.Errorf("stats deleted-command migration: count %s: %w", command, err) - } - return count, nil -} - -func markDocUsageEntriesDeleted(ctx context.Context, docs storage.DocStore[usageEntry], command string) (int64, error) { - keys, err := docs.List(ctx, command) - if err != nil { - return 0, fmt.Errorf("stats deleted-command migration: list %s: %w", command, err) - } - var matched int64 - for _, key := range keys { - isMatch, err := markUsageEntryDeleted(ctx, docs, key, command) - if err != nil { - return 0, err - } - if isMatch { - matched++ - } - } - return matched, nil -} - -// markUsageEntryDeleted flags one row as deleted with a versioned write, -// retrying on conflict so a concurrent writer's update is never overwritten. It -// reports whether the row belongs to command at all, since the key prefix -// alone also matches longer command names. -func markUsageEntryDeleted(ctx context.Context, docs storage.DocStore[usageEntry], key, command string) (bool, error) { - for attempt := 0; attempt < deletedCommandMigrationRetries; attempt++ { - entry, version, err := docs.Get(ctx, key) - if errors.Is(err, storage.ErrNotFound) { - return false, nil - } - if err != nil { - return false, fmt.Errorf("stats deleted-command migration: read %s: %w", key, err) - } - if entry.Cmd != command { - return false, nil - } - if entry.Deleted { - return true, nil - } - entry.Deleted = true - if err := docs.PutVersioned(ctx, key, version, entry); err == nil { - return true, nil - } else if !errors.Is(err, storage.ErrConflict) { - return false, fmt.Errorf("stats deleted-command migration: write %s: %w", key, err) - } - } - return false, fmt.Errorf("stats deleted-command migration: write %s: %w", key, storage.ErrConflict) + return ensureUsageIndexes(ctx, mongoColl) } func ensureUsageIndexes(ctx context.Context, coll *mongo.Collection) error { diff --git a/internal/modules/stats/startup_mongo_test.go b/internal/modules/stats/startup_mongo_test.go index 8ac89ee..f7044d4 100644 --- a/internal/modules/stats/startup_mongo_test.go +++ b/internal/modules/stats/startup_mongo_test.go @@ -8,7 +8,6 @@ import ( "time" "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" "github.com/tiennm99/tiennm99bot/internal/testutil/mongotest" ) @@ -19,7 +18,7 @@ func TestMain(m *testing.M) { } func TestInitStore_MongoCreatesIndexes(t *testing.T) { - ctx, statsColl, systemColl := setupMongoStatsTest(t) + ctx, statsColl := setupMongoStatsTest(t) rawStatsColl, ok := storage.MongoCollection(statsColl) if !ok { @@ -27,65 +26,32 @@ func TestInitStore_MongoCreatesIndexes(t *testing.T) { } docs := storage.Typed[usageEntry](statsColl) - if err := docs.Put(ctx, "stock_dividend", usageEntry{Cmd: "stock_dividend", N: 9}); err != nil { + if err := docs.Put(ctx, "stock_dividend", usageEntry{Cmd: "stock_dividend", N: 9, Deleted: true}); err != nil { t.Fatalf("seed anonymous stats: %v", err) } - if err := docs.Put(ctx, "stock_dividend:7", usageEntry{Cmd: "stock_dividend", UserID: 7, Username: "alice", N: 4}); err != nil { + if err := docs.Put(ctx, "stock_dividend:7", usageEntry{Cmd: "stock_dividend", UserID: 7, Username: "alice", N: 4, Deleted: true}); err != nil { t.Fatalf("seed user stats: %v", err) } if err := docs.Put(ctx, "stock_dividend_extra", usageEntry{Cmd: "stock_dividend_extra", N: 3}); err != nil { t.Fatalf("seed prefix stats: %v", err) } - if err := InitStore(ctx, statsColl, systemColl); err != nil { + if err := InitStore(ctx, statsColl); err != nil { t.Fatalf("InitStore: %v", err) } - if err := InitStore(ctx, statsColl, systemColl); err != nil { + if err := InitStore(ctx, statsColl); err != nil { t.Fatalf("InitStore second run: %v", err) } - for _, key := range []string{"stock_dividend", "stock_dividend:7"} { - entry, _, err := docs.Get(ctx, key) - if err != nil || !entry.Deleted { - t.Fatalf("retained %s = %+v, err=%v", key, entry, err) - } - } - prefixEntry, _, err := docs.Get(ctx, "stock_dividend_extra") - if err != nil || prefixEntry.Deleted { - t.Fatalf("prefix entry = %+v, err=%v", prefixEntry, err) - } - marker, exists, err := systemstate.New(systemColl).Get(ctx, deletedStockDividendMarkerKey) - if err != nil || !exists || marker.Status != "completed" || marker.Count != 2 { - t.Fatalf("marker=%+v exists=%v err=%v", marker, exists, err) - } - legacy, _, err := docs.Get(ctx, "stock_dividend:7") - if err != nil { - t.Fatal(err) - } - legacy.Deleted = false - legacy.N++ - if err := docs.Put(ctx, "stock_dividend:7", legacy); err != nil { - t.Fatal(err) - } store := newUsageStore(statsColl) - if err := store.Increment(ctx, "stock_dividend", usageUser{ID: 7, Username: "alice"}, true); err != nil { - t.Fatalf("retired Increment: %v", err) - } if rows, err := store.TopCommands(ctx, 10); err != nil || len(rows) != 1 || rows[0].display != "/stock_dividend_extra" || rows[0].n != 3 { - t.Fatalf("top commands after legacy write = %+v, err=%v", rows, err) + t.Fatalf("top commands with deleted rows = %+v, err=%v", rows, err) } if rows, err := store.TopUsers(ctx, 10); err != nil || len(rows) != 0 { - t.Fatalf("retired top users = %+v, err=%v", rows, err) + t.Fatalf("deleted top users = %+v, err=%v", rows, err) } if rows, err := store.UsersByCommand(ctx, "stock_dividend", 10); err != nil || len(rows) != 0 { - t.Fatalf("retired users = %+v, err=%v", rows, err) - } - if err := InitStore(ctx, statsColl, systemColl); err != nil { - t.Fatalf("InitStore reconciliation: %v", err) - } - legacy, _, err = docs.Get(ctx, "stock_dividend:7") - if err != nil || !legacy.Deleted || legacy.N != 5 { - t.Fatalf("reconciled legacy entry = %+v, err=%v", legacy, err) + t.Fatalf("deleted users = %+v, err=%v", rows, err) } cur, err := rawStatsColl.Indexes().List(ctx) @@ -114,7 +80,7 @@ func TestInitStore_MongoCreatesIndexes(t *testing.T) { } } -func setupMongoStatsTest(t *testing.T) (context.Context, storage.Collection, storage.Collection) { +func setupMongoStatsTest(t *testing.T) (context.Context, storage.Collection) { t.Helper() uri := mongoTests.URI(t) @@ -136,10 +102,10 @@ func setupMongoStatsTest(t *testing.T) (context.Context, storage.Collection, sto }) provider := storage.NewMongoProvider(db) - return ctx, provider.Collection("stats"), provider.Collection(systemstate.CollectionName) + return ctx, provider.Collection("stats") } func TestInc_MongoUsernameMoveClearsOldHolder(t *testing.T) { - _, statsColl, _ := setupMongoStatsTest(t) + _, statsColl := setupMongoStatsTest(t) assertUsernameMoveClearsOldHolder(t, statsColl) } diff --git a/internal/modules/stats/startup_test.go b/internal/modules/stats/startup_test.go deleted file mode 100644 index 8f91e3d..0000000 --- a/internal/modules/stats/startup_test.go +++ /dev/null @@ -1,98 +0,0 @@ -package stats - -import ( - "context" - "testing" - - "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" -) - -func TestInitStoreMarksStockDividendStatsDeletedAndRetainsHistory(t *testing.T) { - ctx := context.Background() - provider := storage.NewMemoryProvider() - statsColl := provider.Collection("stats") - systemColl := provider.Collection(systemstate.CollectionName) - docs := storage.Typed[usageEntry](statsColl) - - seed := map[string]usageEntry{ - "stock_dividend": {Cmd: "stock_dividend", N: 12}, - "stock_dividend:7": {Cmd: "stock_dividend", UserID: 7, Username: "alice", N: 5, Deleted: true}, - "stock_dividend_extra": {Cmd: "stock_dividend_extra", N: 3}, - "stock_cash_dividend:7": {Cmd: "stock_cash_dividend", UserID: 7, Username: "alice", N: 2}, - } - for key, entry := range seed { - if err := docs.Put(ctx, key, entry); err != nil { - t.Fatalf("seed %s: %v", key, err) - } - } - - if err := InitStore(ctx, statsColl, systemColl); err != nil { - t.Fatalf("InitStore: %v", err) - } - if err := InitStore(ctx, statsColl, systemColl); err != nil { - t.Fatalf("InitStore second run: %v", err) - } - - for _, key := range []string{"stock_dividend", "stock_dividend:7"} { - entry, _, err := docs.Get(ctx, key) - if err != nil { - t.Fatalf("get retained %s: %v", key, err) - } - if !entry.Deleted || entry.N != seed[key].N || entry.Username != seed[key].Username { - t.Fatalf("retained %s = %+v, want deleted history %+v", key, entry, seed[key]) - } - } - for _, key := range []string{"stock_dividend_extra", "stock_cash_dividend:7"} { - entry, _, err := docs.Get(ctx, key) - if err != nil || entry.Deleted { - t.Fatalf("unrelated %s = %+v, err=%v", key, entry, err) - } - } - - marker, exists, err := systemstate.New(systemColl).Get(ctx, deletedStockDividendMarkerKey) - if err != nil || !exists || marker.Status != "completed" || marker.Count != 2 { - t.Fatalf("marker=%+v exists=%v err=%v", marker, exists, err) - } - - rows, err := newUsageStore(statsColl).TopCommands(ctx, 10) - if err != nil { - t.Fatalf("TopCommands: %v", err) - } - for _, row := range rows { - if row.display == "/stock_dividend" { - t.Fatalf("retired command remains visible: %+v", rows) - } - } - - // Simulate a legacy instance writing after the first startup completed. - legacy, _, err := docs.Get(ctx, "stock_dividend:7") - if err != nil { - t.Fatal(err) - } - legacy.Deleted = false - legacy.N++ - if err := docs.Put(ctx, "stock_dividend:7", legacy); err != nil { - t.Fatal(err) - } - store := newUsageStore(statsColl) - if err := store.Increment(ctx, "stock_dividend", usageUser{ID: 7, Username: "alice"}, true); err != nil { - t.Fatalf("retired Increment: %v", err) - } - if rows, err := store.UsersByCommand(ctx, "stock_dividend", 10); err != nil || len(rows) != 0 { - t.Fatalf("retired users = %+v, err=%v", rows, err) - } - if rows, err := store.TopUsers(ctx, 10); err != nil || len(rows) != 1 || rows[0].display != "@alice" || rows[0].n != 2 { - t.Fatalf("top users after legacy write = %+v, err=%v", rows, err) - } - if rows, found, err := store.CommandsByUser(ctx, "alice", 10); err != nil || !found || len(rows) != 1 || rows[0].display != "/stock_cash_dividend" { - t.Fatalf("commands after legacy write = %+v, found=%v err=%v", rows, found, err) - } - if err := InitStore(ctx, statsColl, systemColl); err != nil { - t.Fatalf("InitStore reconciliation: %v", err) - } - legacy, _, err = docs.Get(ctx, "stock_dividend:7") - if err != nil || !legacy.Deleted || legacy.N != 6 { - t.Fatalf("reconciled legacy entry = %+v, err=%v", legacy, err) - } -} diff --git a/internal/modules/stats/usage_store.go b/internal/modules/stats/usage_store.go index 5db42c0..99bf701 100644 --- a/internal/modules/stats/usage_store.go +++ b/internal/modules/stats/usage_store.go @@ -504,9 +504,12 @@ func (s *mongoUsageStore) userIDByUsername(ctx context.Context, username string) return doc.UserID, true, nil } +// deletedStockDividendCommand names the retired /stock_dividend command. Its +// rows are kept with deleted: true as history. +const deletedStockDividendCommand = "stock_dividend" + // isRetiredCommand reports whether cmd has been removed from the bot. Its rows -// are never incremented or shown, even ones written before the migration -// marked them deleted. +// are never incremented or shown, even ones missing the deleted flag. func isRetiredCommand(cmd string) bool { return cmd == deletedStockDividendCommand } diff --git a/internal/modules/sticker/startup.go b/internal/modules/sticker/startup.go deleted file mode 100644 index f02ae84..0000000 --- a/internal/modules/sticker/startup.go +++ /dev/null @@ -1,90 +0,0 @@ -package sticker - -import ( - "context" - "fmt" - "time" - - "github.com/tiennm99/tiennm99bot/internal/log" - "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" -) - -const legacyPackCleanupMarkerKey = "migration:sticker-drop-legacy-packs-v1" - -// legacyRecord is a placeholder type, not a schema. -// -// The retired design wrote three different shapes into this collection — pack -// records keyed by owner ID, "slug:" name reservations, and "pending-delete:" -// confirmations. Cleanup only lists keys and deletes them, and neither -// operation decodes a document, so one empty type serves for all three rather -// than resurrecting structs whose only remaining purpose would be deletion. -type legacyRecord struct{} - -// InitStore removes the per-user sticker pack records the retired pack -// commands left behind. -// -// The module no longer stores anything: /addsticker writes to one shared, -// env-configured set and takes the set owner's user ID from OWNER_ID, so there -// is nothing per-user to key. Its factory ignores the collection entirely. -// These documents are therefore unreachable by any code path — not stale data -// that some handler might still read, but orphans. -// -// Guarded by a completion marker so the scan runs once per database rather than -// on every boot, matching the stock and stats migrations. Safe to run against a -// collection that is already empty, and safe on the memory backend where the -// collection never had anything in it. -func InitStore(ctx context.Context, stickerColl, systemColl storage.Collection) error { - system := systemstate.New(systemColl) - marker, exists, err := system.Get(ctx, legacyPackCleanupMarkerKey) - if err != nil { - return fmt.Errorf("sticker legacy pack cleanup: read marker: %w", err) - } - if exists && marker.Status == "completed" { - return nil - } - - docs := storage.Typed[legacyRecord](stickerColl) - // Empty prefix: the retired design used three disjoint key spaces (bare - // owner IDs, "slug:", "pending-delete:") and all of them are dead, so - // listing everything is both correct and cheaper than three scans. - keys, err := docs.List(ctx, "") - if err != nil { - return fmt.Errorf("sticker legacy pack cleanup: list records: %w", err) - } - - var deleted int64 - for _, key := range keys { - if err := docs.Delete(ctx, key); err != nil { - // Abort without writing the marker, so the next boot retries the - // rest. Deletes are idempotent, so a partial run is safe to repeat. - return fmt.Errorf("sticker legacy pack cleanup: delete %s: %w", key, err) - } - deleted++ - } - if deleted > 0 { - log.Info("sticker legacy pack records removed", "count", deleted) - } - - marker = completedLegacyPackCleanup(marker, exists, deleted, time.Now().UnixMilli()) - if err := system.Put(ctx, legacyPackCleanupMarkerKey, marker); err != nil { - return fmt.Errorf("sticker legacy pack cleanup: write marker: %w", err) - } - return nil -} - -func completedLegacyPackCleanup(marker systemstate.Record, exists bool, deleted, now int64) systemstate.Record { - if !exists { - marker = systemstate.Record{ - Kind: "migration", - Name: "sticker drop legacy packs v1", - } - } - if marker.CompletedAt == 0 { - marker.CompletedAt = now - } - marker.Status = "completed" - marker.Count += deleted - marker.UpdatedAt = now - return marker -} diff --git a/internal/modules/sticker/startup_test.go b/internal/modules/sticker/startup_test.go deleted file mode 100644 index 61e8718..0000000 --- a/internal/modules/sticker/startup_test.go +++ /dev/null @@ -1,146 +0,0 @@ -package sticker - -import ( - "context" - "testing" - - "github.com/tiennm99/tiennm99bot/internal/modules" - "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" -) - -// legacyShape stands in for the retired records: a pack keyed by owner ID, a -// "slug:" reservation, and a "pending-delete:" confirmation. Only the keys -// matter to the cleanup, so one loose shape covers all three. -type legacyShape struct { - Slug string `bson:"slug"` - OwnerID int64 `bson:"ownerId"` -} - -func seedLegacy(t *testing.T, coll storage.Collection, keys ...string) { - t.Helper() - docs := storage.Typed[legacyShape](coll) - for _, k := range keys { - if err := docs.Put(context.Background(), k, legacyShape{Slug: "old", OwnerID: 42}); err != nil { - t.Fatalf("seed %s: %v", k, err) - } - } -} - -func remainingKeys(t *testing.T, coll storage.Collection) []string { - t.Helper() - keys, err := storage.Typed[legacyShape](coll).List(context.Background(), "") - if err != nil { - t.Fatalf("list: %v", err) - } - return keys -} - -// All three retired key spaces go, in one pass. -func TestInitStore_RemovesEveryLegacyKeySpace(t *testing.T) { - provider := storage.NewMemoryProvider() - stickerColl := provider.Collection(CollectionName) - systemColl := provider.Collection(systemstate.CollectionName) - - seedLegacy(t, stickerColl, - "123456789", // pack record, keyed by owner ID - "987654321", // another owner's pack - "slug:mypack", // name reservation - "pending-delete:1234567", // /delpack confirmation - ) - - if err := InitStore(context.Background(), stickerColl, systemColl); err != nil { - t.Fatalf("InitStore: %v", err) - } - - if got := remainingKeys(t, stickerColl); len(got) != 0 { - t.Errorf("collection still holds %v, want it emptied", got) - } -} - -// The marker records how many were removed, so the count is auditable after -// the fact. -func TestInitStore_MarksCompletionWithCount(t *testing.T) { - provider := storage.NewMemoryProvider() - stickerColl := provider.Collection(CollectionName) - systemColl := provider.Collection(systemstate.CollectionName) - seedLegacy(t, stickerColl, "1", "2", "slug:x") - - if err := InitStore(context.Background(), stickerColl, systemColl); err != nil { - t.Fatalf("InitStore: %v", err) - } - - rec, found, err := systemstate.New(systemColl).Get(context.Background(), legacyPackCleanupMarkerKey) - if err != nil || !found { - t.Fatalf("marker: found=%v err=%v", found, err) - } - if rec.Status != "completed" { - t.Errorf("status = %q, want completed", rec.Status) - } - if rec.Count != 3 { - t.Errorf("count = %d, want 3", rec.Count) - } - if rec.CompletedAt == 0 || rec.UpdatedAt == 0 { - t.Errorf("timestamps unset: %+v", rec) - } -} - -// Once marked, the scan must not run again — and specifically must not delete -// anything written to this collection later. -func TestInitStore_MarkerStopsASecondPass(t *testing.T) { - provider := storage.NewMemoryProvider() - stickerColl := provider.Collection(CollectionName) - systemColl := provider.Collection(systemstate.CollectionName) - seedLegacy(t, stickerColl, "1") - - if err := InitStore(context.Background(), stickerColl, systemColl); err != nil { - t.Fatalf("first InitStore: %v", err) - } - - // Whatever a future version of this module might store. - seedLegacy(t, stickerColl, "something-new") - - if err := InitStore(context.Background(), stickerColl, systemColl); err != nil { - t.Fatalf("second InitStore: %v", err) - } - got := remainingKeys(t, stickerColl) - if len(got) != 1 || got[0] != "something-new" { - t.Errorf("remaining = %v, want only the newly written key", got) - } -} - -// An already-clean database is the normal case on a fresh deploy, and on the -// memory backend where the collection never held anything. -func TestInitStore_EmptyCollectionIsFine(t *testing.T) { - provider := storage.NewMemoryProvider() - stickerColl := provider.Collection(CollectionName) - systemColl := provider.Collection(systemstate.CollectionName) - - if err := InitStore(context.Background(), stickerColl, systemColl); err != nil { - t.Fatalf("InitStore on an empty collection: %v", err) - } - - rec, found, err := systemstate.New(systemColl).Get(context.Background(), legacyPackCleanupMarkerKey) - if err != nil || !found { - t.Fatalf("marker: found=%v err=%v", found, err) - } - if rec.Count != 0 { - t.Errorf("count = %d, want 0", rec.Count) - } -} - -// The module itself must keep using none of this: if a future edit gives the -// factory a store, the cleanup above would start deleting live data. -func TestNew_UsesNoStorage(t *testing.T) { - provider := storage.NewMemoryProvider() - coll := provider.Collection(CollectionName) - seedLegacy(t, coll, "sentinel") - - mod := New(modules.Deps{Store: coll}) - if len(mod.Commands) != 1 || mod.Commands[0].Name != "addsticker" { - t.Fatalf("module commands = %+v, want only /addsticker", mod.Commands) - } - if got := remainingKeys(t, coll); len(got) != 1 { - t.Errorf("factory touched storage; remaining = %v", got) - } -} diff --git a/internal/modules/stock/startup.go b/internal/modules/stock/startup.go deleted file mode 100644 index 95bdfd7..0000000 --- a/internal/modules/stock/startup.go +++ /dev/null @@ -1,148 +0,0 @@ -package stock - -import ( - "context" - "errors" - "fmt" - "time" - - "github.com/tiennm99/tiennm99bot/internal/log" - "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" -) - -const ( - dividendHistoryMigrationMarkerKey = "migration:stock-dividend-history-v1" - dividendHistoryMigrationRetries = 5 -) - -// legacyDividendAssetPosition contains the retired discovery cursor so a -// whole-document replacement can physically remove it from MongoDB. -type legacyDividendAssetPosition struct { - Quantity int64 `json:"quantity" bson:"quantity"` - Base float64 `json:"base" bson:"base"` - DividendCheckedAt *int64 `json:"dividendCheckedAt,omitempty" bson:"dividendCheckedAt,omitempty"` - OpenedAt int64 `json:"openedAt,omitempty" bson:"openedAt,omitempty"` -} - -// legacyDividendPortfolio mirrors both the current portfolio fields and the -// retired applied-event ledger. Keeping Dividends here is important: migration -// must never discard history already written by the new runtime. -type legacyDividendPortfolio struct { - VND float64 `json:"vnd" bson:"vnd"` - Assets map[string]legacyDividendAssetPosition `json:"assets" bson:"assets"` - Dividends map[string]map[string]DividendRecord `json:"dividends,omitempty" bson:"dividends,omitempty"` - AppliedDividendEvents map[string]int64 `json:"appliedDividendEvents,omitempty" bson:"appliedDividendEvents,omitempty"` - Meta PortfolioMeta `json:"meta" bson:"meta"` -} - -// InitStore removes the cursor-era dividend fields from stock user portfolios. -// The completion marker makes the one-time migration idempotent; it is written -// only after every listed user document has been handled successfully. -func InitStore(ctx context.Context, portfolioColl, systemColl storage.Collection) error { - system := systemstate.New(systemColl) - marker, exists, err := system.Get(ctx, dividendHistoryMigrationMarkerKey) - if err != nil { - return fmt.Errorf("stock dividend history migration: read marker: %w", err) - } - if exists && marker.Status == "completed" { - return nil - } - - docs := storage.Typed[legacyDividendPortfolio](portfolioColl) - keys, err := docs.List(ctx, "user:") - if err != nil { - return fmt.Errorf("stock dividend history migration: list portfolios: %w", err) - } - - var migrated int64 - for index, key := range keys { - changed, err := migrateDividendHistorySchema(ctx, docs, key) - if err != nil { - return err - } - if changed { - migrated++ - log.Info("stock dividend history migrated", "portfolio", index+1, "total", len(keys)) - } - } - - marker = completedDividendHistoryMigration(marker, exists, migrated, time.Now().UnixMilli()) - if err := system.Put(ctx, dividendHistoryMigrationMarkerKey, marker); err != nil { - return fmt.Errorf("stock dividend history migration: write marker: %w", err) - } - return nil -} - -func completedDividendHistoryMigration(marker systemstate.Record, exists bool, migrated, now int64) systemstate.Record { - if !exists { - marker = systemstate.Record{ - Kind: "migration", - Name: "stock dividend history v1", - } - } - if marker.CompletedAt == 0 { - marker.CompletedAt = now - } - marker.Status = "completed" - marker.Count += migrated - marker.UpdatedAt = now - return marker -} - -func migrateDividendHistorySchema(ctx context.Context, docs storage.DocStore[legacyDividendPortfolio], key string) (bool, error) { - for attempt := 0; attempt < dividendHistoryMigrationRetries; attempt++ { - doc, version, err := docs.Get(ctx, key) - if err != nil { - return false, fmt.Errorf("stock dividend history migration: read %s: %w", key, err) - } - if !doc.hasLegacyDividendFields() { - return false, nil - } - current := doc.currentPortfolio() - if err := current.Validate(); err != nil { - return false, fmt.Errorf("stock dividend history migration: validate %s: %w", key, err) - } - err = docs.PutVersioned(ctx, key, version, doc.withoutLegacyDividendFields()) - if err == nil { - return true, nil - } - if !errors.Is(err, storage.ErrConflict) { - return false, fmt.Errorf("stock dividend history migration: write %s: %w", key, err) - } - } - return false, fmt.Errorf("stock dividend history migration: write %s: %w", key, storage.ErrConflict) -} - -func (p legacyDividendPortfolio) hasLegacyDividendFields() bool { - if p.AppliedDividendEvents != nil { - return true - } - for _, position := range p.Assets { - if position.DividendCheckedAt != nil { - return true - } - } - return false -} - -func (p legacyDividendPortfolio) withoutLegacyDividendFields() legacyDividendPortfolio { - p.AppliedDividendEvents = nil - for ticker, position := range p.Assets { - position.DividendCheckedAt = nil - p.Assets[ticker] = position - } - return p -} - -func (p legacyDividendPortfolio) currentPortfolio() Portfolio { - assets := make(map[string]AssetPosition, len(p.Assets)) - for ticker, position := range p.Assets { - assets[ticker] = AssetPosition{ - Quantity: position.Quantity, - Base: position.Base, - OpenedAt: position.OpenedAt, - } - } - return Portfolio{VND: p.VND, Assets: assets, Dividends: p.Dividends, Meta: p.Meta} -} diff --git a/internal/modules/stock/startup_mongo_test.go b/internal/modules/stock/startup_mongo_test.go deleted file mode 100644 index c038ea9..0000000 --- a/internal/modules/stock/startup_mongo_test.go +++ /dev/null @@ -1,88 +0,0 @@ -package stock - -import ( - "context" - "fmt" - "os" - "testing" - "time" - - "go.mongodb.org/mongo-driver/v2/bson" - - "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" - "github.com/tiennm99/tiennm99bot/internal/testutil/mongotest" -) - -var stockMongoTests mongotest.Manager - -func TestMain(m *testing.M) { - os.Exit(stockMongoTests.Run(m)) -} - -func TestInitStorePhysicallyDropsLegacyDividendFieldsInMongoDB(t *testing.T) { - ctx, portfolioColl, systemColl := setupMongoStockMigrationTest(t) - docs := storage.Typed[legacyDividendPortfolio](portfolioColl) - if err := docs.Put(ctx, "user:7", legacyDividendPortfolio{ - VND: 50_000, - Assets: map[string]legacyDividendAssetPosition{ - "TCB": {Quantity: 100, Base: 3_000_000, DividendCheckedAt: legacyCursor(123), OpenedAt: 99}, - }, - Dividends: map[string]map[string]DividendRecord{ - "TCB": {"2612974": legacyPreservedDividend()}, - }, - AppliedDividendEvents: map[string]int64{"old-hash": 456}, - Meta: PortfolioMeta{Invested: 3_000_000, CreatedAt: 1}, - }); err != nil { - t.Fatal(err) - } - - if err := InitStore(ctx, portfolioColl, systemColl); err != nil { - t.Fatalf("InitStore: %v", err) - } - if err := InitStore(ctx, portfolioColl, systemColl); err != nil { - t.Fatalf("second InitStore: %v", err) - } - - rawColl, ok := storage.MongoCollection(portfolioColl) - if !ok { - t.Fatal("portfolio collection is not MongoDB-backed") - } - var raw bson.Raw - if err := rawColl.FindOne(ctx, bson.M{"_id": "user:7"}).Decode(&raw); err != nil { - t.Fatal(err) - } - if _, err := raw.LookupErr("appliedDividendEvents"); err == nil { - t.Fatalf("root legacy ledger remains in %v", raw) - } - position := raw.Lookup("assets").Document().Lookup("TCB").Document() - if _, err := position.LookupErr("dividendCheckedAt"); err == nil { - t.Fatalf("asset legacy cursor remains in %v", position) - } - if position.Lookup("openedAt").Int64() != 99 || raw.Lookup("vnd").Double() != 50_000 { - t.Fatalf("current portfolio fields changed: %v", raw) - } - if _, err := raw.LookupErr("dividends"); err != nil { - t.Fatalf("new dividend history was discarded: %v", raw) - } -} - -func setupMongoStockMigrationTest(t *testing.T) (context.Context, storage.Collection, storage.Collection) { - t.Helper() - uri := stockMongoTests.URI(t) - ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) - t.Cleanup(cancel) - client, err := storage.NewMongoClient(ctx, uri) - if err != nil { - t.Fatal(err) - } - db := client.Database(fmt.Sprintf("tiennm99bot_stock_migration_test_%d", time.Now().UnixNano())) - t.Cleanup(func() { - cleanupCtx, cleanupCancel := context.WithTimeout(context.Background(), 10*time.Second) - defer cleanupCancel() - _ = db.Drop(cleanupCtx) - _ = client.Disconnect(cleanupCtx) - }) - provider := storage.NewMongoProvider(db) - return ctx, provider.Collection(CollectionName), provider.Collection(systemstate.CollectionName) -} diff --git a/internal/modules/stock/startup_test.go b/internal/modules/stock/startup_test.go deleted file mode 100644 index c4340e4..0000000 --- a/internal/modules/stock/startup_test.go +++ /dev/null @@ -1,165 +0,0 @@ -package stock - -import ( - "context" - "errors" - "testing" - - "github.com/tiennm99/tiennm99bot/internal/storage" - "github.com/tiennm99/tiennm99bot/internal/systemstate" -) - -func legacyCursor(value int64) *int64 { return &value } - -func legacyPreservedDividend() DividendRecord { - return DividendRecord{ - Kind: DividendKindCash, - PublishedAt: 1, - RecordDate: 2, - VNDPerShare: 1_500, - Title: "Cash dividend", - SourceURL: "https://example.test/event/2612974", - } -} - -func TestInitStoreRemovesLegacyDividendFields(t *testing.T) { - ctx := context.Background() - provider := storage.NewMemoryProvider() - portfolioColl := provider.Collection(CollectionName) - systemColl := provider.Collection(systemstate.CollectionName) - docs := storage.Typed[legacyDividendPortfolio](portfolioColl) - - legacy := legacyDividendPortfolio{ - VND: 50_000, - Assets: map[string]legacyDividendAssetPosition{ - "TCB": {Quantity: 100, Base: 3_000_000, DividendCheckedAt: legacyCursor(123), OpenedAt: 99}, - }, - Dividends: map[string]map[string]DividendRecord{ - "TCB": {"2612974": legacyPreservedDividend()}, - }, - AppliedDividendEvents: map[string]int64{"old-hash": 456}, - Meta: PortfolioMeta{Invested: 3_000_000, CreatedAt: 1}, - } - if err := docs.Put(ctx, "user:7", legacy); err != nil { - t.Fatal(err) - } - if err := docs.Put(ctx, "pending-dividend:leave-me", legacy); err != nil { - t.Fatal(err) - } - - if err := InitStore(ctx, portfolioColl, systemColl); err != nil { - t.Fatalf("InitStore: %v", err) - } - - got, _, err := docs.Get(ctx, "user:7") - if err != nil { - t.Fatal(err) - } - if got.AppliedDividendEvents != nil || got.Assets["TCB"].DividendCheckedAt != nil { - t.Fatalf("legacy fields remain: %+v", got) - } - if got.VND != legacy.VND || got.Assets["TCB"].Quantity != 100 || got.Assets["TCB"].Base != 3_000_000 || got.Assets["TCB"].OpenedAt != 99 || got.Meta != legacy.Meta { - t.Fatalf("portfolio data changed: got=%+v want=%+v", got, legacy) - } - if event, ok := got.Dividends["TCB"]["2612974"]; !ok || event != legacyPreservedDividend() { - t.Fatalf("dividend history was discarded: %+v", got.Dividends) - } - - pending, _, err := docs.Get(ctx, "pending-dividend:leave-me") - if err != nil { - t.Fatal(err) - } - if pending.AppliedDividendEvents == nil || pending.Assets["TCB"].DividendCheckedAt == nil { - t.Fatalf("non-user document was rewritten: %+v", pending) - } - - marker, exists, err := systemstate.New(systemColl).Get(ctx, dividendHistoryMigrationMarkerKey) - if err != nil || !exists || marker.Status != "completed" || marker.Count != 1 { - t.Fatalf("marker=%+v exists=%v err=%v", marker, exists, err) - } - - if err := InitStore(ctx, portfolioColl, systemColl); err != nil { - t.Fatalf("second InitStore: %v", err) - } - marker, _, _ = systemstate.New(systemColl).Get(ctx, dividendHistoryMigrationMarkerKey) - if marker.Count != 1 { - t.Fatalf("idempotent marker count=%d, want 1", marker.Count) - } -} - -type conflictOnceDividendMigrationStore struct { - storage.DocStore[legacyDividendPortfolio] - conflicted bool -} - -func (s *conflictOnceDividendMigrationStore) PutVersioned(ctx context.Context, key string, version int64, value legacyDividendPortfolio) error { - if !s.conflicted { - s.conflicted = true - return storage.ErrConflict - } - return s.DocStore.PutVersioned(ctx, key, version, value) -} - -func TestDividendHistoryMigrationRetriesVersionConflict(t *testing.T) { - ctx := context.Background() - provider := storage.NewMemoryProvider() - docs := storage.Typed[legacyDividendPortfolio](provider.Collection(CollectionName)) - if err := docs.Put(ctx, "user:7", legacyDividendPortfolio{ - Assets: map[string]legacyDividendAssetPosition{ - "TCB": {Quantity: 10, Base: 300_000, DividendCheckedAt: legacyCursor(123), OpenedAt: 99}, - }, - }); err != nil { - t.Fatal(err) - } - store := &conflictOnceDividendMigrationStore{DocStore: docs} - changed, err := migrateDividendHistorySchema(ctx, store, "user:7") - if err != nil || !changed || !store.conflicted { - t.Fatalf("changed=%v conflicted=%v err=%v", changed, store.conflicted, err) - } -} - -type alwaysConflictDividendMigrationStore struct { - storage.DocStore[legacyDividendPortfolio] -} - -func (s alwaysConflictDividendMigrationStore) PutVersioned(context.Context, string, int64, legacyDividendPortfolio) error { - return storage.ErrConflict -} - -func TestDividendHistoryMigrationReturnsExhaustedConflict(t *testing.T) { - ctx := context.Background() - provider := storage.NewMemoryProvider() - docs := storage.Typed[legacyDividendPortfolio](provider.Collection(CollectionName)) - if err := docs.Put(ctx, "user:7", legacyDividendPortfolio{ - Assets: map[string]legacyDividendAssetPosition{ - "TCB": {Quantity: 10, Base: 300_000, DividendCheckedAt: legacyCursor(123), OpenedAt: 99}, - }, - }); err != nil { - t.Fatal(err) - } - changed, err := migrateDividendHistorySchema(ctx, alwaysConflictDividendMigrationStore{docs}, "user:7") - if changed || !errors.Is(err, storage.ErrConflict) { - t.Fatalf("changed=%v err=%v, want wrapped conflict", changed, err) - } -} - -func TestInitStoreDoesNotMarkFailedMigrationComplete(t *testing.T) { - ctx := context.Background() - provider := storage.NewMemoryProvider() - // A legacy portfolio with an invalid position must fail validation before - // its retired fields are removed. - if err := storage.Typed[legacyDividendPortfolio](provider.Collection(CollectionName)).Put(ctx, "user:7", legacyDividendPortfolio{ - Assets: map[string]legacyDividendAssetPosition{ - "TCB": {Quantity: 10, Base: -1, DividendCheckedAt: legacyCursor(123), OpenedAt: 99}, - }, - }); err != nil { - t.Fatal(err) - } - if err := InitStore(ctx, provider.Collection(CollectionName), provider.Collection(systemstate.CollectionName)); err == nil { - t.Fatal("InitStore accepted invalid legacy portfolio") - } - _, exists, err := systemstate.New(provider.Collection(systemstate.CollectionName)).Get(ctx, dividendHistoryMigrationMarkerKey) - if err != nil || exists { - t.Fatalf("completion marker exists=%v after failure, err=%v", exists, err) - } -}