mirror of
https://github.com/tiennm99/tiennm99bot.git
synced 2026-10-11 03:13:46 +00:00
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.
This commit is contained in:
1 parent
5344111787
commit
5e96490254
13 files changed
+31
-1030
No files matched your search
@@ -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 {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user