accounts: add migration code from kvdb to SQL

This commit introduces the migration logic for transitioning the
accounts store from kvdb to SQL.

Note that as of this commit, the migration is not yet triggered by any
production code, i.e. only tests execute the migration logic.
This commit is contained in:
Viktor Tigerström 2025-04-16 11:37:39 +02:00
parent 6030f650fd
commit dbb5d73079
No known key found for this signature in database
GPG key ID: B984570980684DCC
3 changed files with 581 additions and 2 deletions

233
accounts/sql_migration.go Normal file
View file

@ -0,0 +1,233 @@
package accounts
import (
"context"
"database/sql"
"errors"
"fmt"
"math"
"reflect"
"time"
"github.com/davecgh/go-spew/spew"
"github.com/lightninglabs/lightning-terminal/db/sqlc"
"github.com/pmezard/go-difflib/difflib"
)
var (
// ErrMigrationMismatch is returned when the migrated account does not
// match the original account.
ErrMigrationMismatch = fmt.Errorf("migrated account does not match " +
"original account")
)
// MigrateAccountStoreToSQL runs the migration of all accounts and indices from
// the KV database to the SQL database. The migration is done in a single
// transaction to ensure that all accounts are migrated or none at all.
func MigrateAccountStoreToSQL(ctx context.Context, kvStore *BoltStore,
tx SQLQueries) error {
log.Infof("Starting migration of the KV accounts store to SQL")
err := migrateAccountsToSQL(ctx, kvStore, tx)
if err != nil {
return fmt.Errorf("unsuccessful migration of accounts to "+
"SQL: %w", err)
}
err = migrateAccountsIndicesToSQL(ctx, kvStore, tx)
if err != nil {
return fmt.Errorf("unsuccessful migration of account indices "+
"to SQL: %w", err)
}
return nil
}
// migrateAccountsToSQL runs the migration of all accounts from the KV database
// to the SQL database. The migration is done in a single transaction to ensure
// that all accounts are migrated or none at all.
func migrateAccountsToSQL(ctx context.Context, kvStore *BoltStore,
tx SQLQueries) error {
log.Infof("Starting migration of accounts from KV to SQL")
kvAccounts, err := kvStore.Accounts(ctx)
if err != nil {
return err
}
for _, kvAccount := range kvAccounts {
migratedAccountID, err := migrateSingleAccountToSQL(
ctx, tx, kvAccount,
)
if err != nil {
return fmt.Errorf("unable to migrate account(%v): %w",
kvAccount.ID, err)
}
migratedAccount, err := getAndMarshalAccount(
ctx, tx, migratedAccountID,
)
if err != nil {
return fmt.Errorf("unable to fetch migrated "+
"account(%v): %w", kvAccount.ID, err)
}
overrideAccountTimeZone(kvAccount)
overrideAccountTimeZone(migratedAccount)
if !reflect.DeepEqual(kvAccount, migratedAccount) {
diff := difflib.UnifiedDiff{
A: difflib.SplitLines(
spew.Sdump(kvAccount),
),
B: difflib.SplitLines(
spew.Sdump(migratedAccount),
),
FromFile: "Expected",
FromDate: "",
ToFile: "Actual",
ToDate: "",
Context: 3,
}
diffText, _ := difflib.GetUnifiedDiffString(diff)
return fmt.Errorf("%w: %v.\n%v", ErrMigrationMismatch,
kvAccount.ID, diffText)
}
}
log.Infof("All accounts migrated from KV to SQL. Total number of "+
"accounts migrated: %d", len(kvAccounts))
return nil
}
// migrateSingleAccountToSQL runs the migration for a single account from the
// KV database to the SQL database.
func migrateSingleAccountToSQL(ctx context.Context,
tx SQLQueries, account *OffChainBalanceAccount) (int64, error) {
accountAlias, err := account.ID.ToInt64()
if err != nil {
return 0, err
}
insertAccountParams := sqlc.InsertAccountParams{
Type: int16(account.Type),
InitialBalanceMsat: int64(account.InitialBalance),
CurrentBalanceMsat: account.CurrentBalance,
LastUpdated: account.LastUpdate.UTC(),
Alias: accountAlias,
Expiration: account.ExpirationDate.UTC(),
Label: sql.NullString{
String: account.Label,
Valid: len(account.Label) > 0,
},
}
sqlId, err := tx.InsertAccount(ctx, insertAccountParams)
if err != nil {
return 0, err
}
for hash := range account.Invoices {
addInvoiceParams := sqlc.AddAccountInvoiceParams{
AccountID: sqlId,
Hash: hash[:],
}
err = tx.AddAccountInvoice(ctx, addInvoiceParams)
if err != nil {
return sqlId, err
}
}
for hash, paymentEntry := range account.Payments {
upsertPaymentParams := sqlc.UpsertAccountPaymentParams{
AccountID: sqlId,
Hash: hash[:],
Status: int16(paymentEntry.Status),
FullAmountMsat: int64(paymentEntry.FullAmount),
}
err = tx.UpsertAccountPayment(ctx, upsertPaymentParams)
if err != nil {
return sqlId, err
}
}
return sqlId, nil
}
// migrateAccountsIndicesToSQL runs the migration for the account indices from
// the KV database to the SQL database.
func migrateAccountsIndicesToSQL(ctx context.Context, kvStore *BoltStore,
tx SQLQueries) error {
log.Infof("Starting migration of accounts indices from KV to SQL")
addIndex, settleIndex, err := kvStore.LastIndexes(ctx)
if errors.Is(err, ErrNoInvoiceIndexKnown) {
log.Infof("No indices found in KV store, skipping migration")
return nil
} else if err != nil {
return err
}
if addIndex > math.MaxInt64 {
return fmt.Errorf("%s:%v is above max int64 value",
addIndexName, addIndex)
}
if settleIndex > math.MaxInt64 {
return fmt.Errorf("%s:%v is above max int64 value",
settleIndexName, settleIndex)
}
setAddIndexParams := sqlc.SetAccountIndexParams{
Name: addIndexName,
Value: int64(addIndex),
}
err = tx.SetAccountIndex(ctx, setAddIndexParams)
if err != nil {
return err
}
setSettleIndexParams := sqlc.SetAccountIndexParams{
Name: settleIndexName,
Value: int64(settleIndex),
}
err = tx.SetAccountIndex(ctx, setSettleIndexParams)
if err != nil {
return err
}
log.Infof("Successfully migratated accounts indices from KV to SQL")
return nil
}
// overrideAccountTimeZone overrides the time zone of the account to the local
// time zone and chops off the nanosecond part for comparison. This is needed
// because KV database stores times as-is which as an unwanted side effect would
// fail migration due to time comparison expecting both the original and
// migrated accounts to be in the same local time zone and in microsecond
// precision. Note that PostgresSQL stores times in microsecond precision while
// SQLite can store times in nanosecond precision if using TEXT storage class.
func overrideAccountTimeZone(account *OffChainBalanceAccount) {
fixTime := func(t time.Time) time.Time {
return t.In(time.Local).Truncate(time.Microsecond)
}
if !account.ExpirationDate.IsZero() {
account.ExpirationDate = fixTime(account.ExpirationDate)
}
if !account.LastUpdate.IsZero() {
account.LastUpdate = fixTime(account.LastUpdate)
}
}

View file

@ -0,0 +1,346 @@
package accounts
import (
"context"
"database/sql"
"testing"
"time"
"github.com/lightninglabs/lightning-terminal/db"
"github.com/lightningnetwork/lnd/clock"
"github.com/lightningnetwork/lnd/fn"
"github.com/lightningnetwork/lnd/lnrpc"
"github.com/lightningnetwork/lnd/lntypes"
"github.com/lightningnetwork/lnd/sqldb"
"github.com/stretchr/testify/require"
)
// TestAccountStoreMigration tests the migration of account store from a bolt
// backed to a SQL database. Note that this test does not attempt to be a
// complete migration test.
func TestAccountStoreMigration(t *testing.T) {
t.Parallel()
ctx := context.Background()
clock := clock.NewTestClock(time.Now())
// When using build tags that creates a kvdb store for NewTestDB, we
// skip this test as it is only applicable for postgres and sqlite tags.
store := NewTestDB(t, clock)
if _, ok := store.(*BoltStore); ok {
t.Skipf("Skipping account store migration test for kvdb build")
}
makeSQLDB := func(t *testing.T) (*SQLStore,
*db.TransactionExecutor[SQLQueries]) {
testDBStore := NewTestDB(t, clock)
t.Cleanup(func() {
require.NoError(t, testDBStore.Close())
})
store, ok := testDBStore.(*SQLStore)
require.True(t, ok)
baseDB := store.BaseDB
genericExecutor := db.NewTransactionExecutor(
baseDB, func(tx *sql.Tx) SQLQueries {
return baseDB.WithTx(tx)
},
)
return store, genericExecutor
}
assertMigrationResults := func(t *testing.T, sqlStore *SQLStore,
kvAccounts []*OffChainBalanceAccount, kvAddIndex uint64,
kvSettleIndex uint64, expectLastIndex bool) {
// The migration function will check if the inserted accounts
// and indices equals the migrated ones, but as a sanity check
// we'll also fetch the accounts and indices from the sql store
// and compare them to the original.
// First we compare the migrated accounts to the original ones.
sqlAccounts, err := sqlStore.Accounts(ctx)
require.NoError(t, err)
require.Equal(t, len(kvAccounts), len(sqlAccounts))
for i := 0; i < len(kvAccounts); i++ {
assertEqualAccounts(t, kvAccounts[i], sqlAccounts[i])
}
// After that we compare the migrated indices. However, if we
// don't expect the last indexes to be set, we don't need to
// compare them.
if !expectLastIndex {
return
}
sqlAddIndex, sqlSettleIndex, err := sqlStore.LastIndexes(ctx)
require.NoError(t, err)
require.Equal(t, kvAddIndex, sqlAddIndex)
require.Equal(t, kvSettleIndex, sqlSettleIndex)
}
tests := []struct {
name string
expectLastIndex bool
populateDB func(t *testing.T, kvStore *BoltStore)
}{
{
name: "empty",
expectLastIndex: false,
populateDB: func(t *testing.T, kvStore *BoltStore) {
// Don't populate the DB.
},
},
{
name: "account no expiry",
expectLastIndex: false,
populateDB: func(t *testing.T, kvStore *BoltStore) {
// Create an account that does not expire.
acct1, err := kvStore.NewAccount(
ctx, 0, time.Time{}, "foo",
)
require.NoError(t, err)
require.False(t, acct1.HasExpired())
},
},
{
name: "account with expiry",
expectLastIndex: false,
populateDB: func(t *testing.T, kvStore *BoltStore) {
// Create an account that does expire.
acct1, err := kvStore.NewAccount(
ctx, 0, time.Now().Add(time.Hour),
"foo",
)
require.NoError(t, err)
require.False(t, acct1.HasExpired())
},
},
{
name: "account with set UpdatedAt",
expectLastIndex: false,
populateDB: func(t *testing.T, kvStore *BoltStore) {
// Create an account that does expire.
acct1, err := kvStore.NewAccount(
ctx, 0, time.Now().Add(time.Hour),
"foo",
)
require.NoError(t, err)
require.False(t, acct1.HasExpired())
err = kvStore.UpdateAccountBalanceAndExpiry(
ctx, acct1.ID, fn.None[int64](),
fn.Some(time.Now().Add(time.Minute)),
)
require.NoError(t, err)
},
},
{
name: "account with balance",
expectLastIndex: false,
populateDB: func(t *testing.T, kvStore *BoltStore) {
// Create an account with balance
acct1, err := kvStore.NewAccount(
ctx, 100000, time.Time{}, "foo",
)
require.NoError(t, err)
require.False(t, acct1.HasExpired())
},
},
{
name: "account with invoices",
expectLastIndex: true,
populateDB: func(t *testing.T, kvStore *BoltStore) {
// Create an account with balance
acct1, err := kvStore.NewAccount(
ctx, 0, time.Time{}, "foo",
)
require.NoError(t, err)
require.False(t, acct1.HasExpired())
hash1 := lntypes.Hash{1, 2, 3, 4}
err = kvStore.AddAccountInvoice(
ctx, acct1.ID, hash1,
)
require.NoError(t, err)
err = kvStore.StoreLastIndexes(ctx, 1, 0)
require.NoError(t, err)
hash2 := lntypes.Hash{1, 2, 3, 4, 5}
err = kvStore.AddAccountInvoice(
ctx, acct1.ID, hash2,
)
require.NoError(t, err)
err = kvStore.StoreLastIndexes(ctx, 2, 1)
require.NoError(t, err)
},
},
{
name: "account with payments",
expectLastIndex: false,
populateDB: func(t *testing.T, kvStore *BoltStore) {
// Create an account with balance
acct1, err := kvStore.NewAccount(
ctx, 0, time.Time{}, "foo",
)
require.NoError(t, err)
require.False(t, acct1.HasExpired())
hash1 := lntypes.Hash{1, 1, 1, 1}
known, err := kvStore.UpsertAccountPayment(
ctx, acct1.ID, hash1, 100,
lnrpc.Payment_UNKNOWN,
)
require.NoError(t, err)
require.False(t, known)
hash2 := lntypes.Hash{2, 2, 2, 2}
known, err = kvStore.UpsertAccountPayment(
ctx, acct1.ID, hash2, 200,
lnrpc.Payment_IN_FLIGHT,
)
require.NoError(t, err)
require.False(t, known)
hash3 := lntypes.Hash{3, 3, 3, 3}
known, err = kvStore.UpsertAccountPayment(
ctx, acct1.ID, hash3, 200,
lnrpc.Payment_SUCCEEDED,
)
require.NoError(t, err)
require.False(t, known)
hash4 := lntypes.Hash{4, 4, 4, 4}
known, err = kvStore.UpsertAccountPayment(
ctx, acct1.ID, hash4, 200,
lnrpc.Payment_FAILED,
)
require.NoError(t, err)
require.False(t, known)
},
},
{
name: "multiple accounts",
expectLastIndex: true,
populateDB: func(t *testing.T, kvStore *BoltStore) {
// Create two accounts with balance and that
// expires.
acct1, err := kvStore.NewAccount(
ctx, 100000, time.Now().Add(time.Hour),
"foo",
)
require.NoError(t, err)
require.False(t, acct1.HasExpired())
acct2, err := kvStore.NewAccount(
ctx, 200000, time.Now().Add(time.Hour),
"bar",
)
require.NoError(t, err)
require.False(t, acct2.HasExpired())
// Create invoices for both accounts.
hash1 := lntypes.Hash{1, 1, 1, 1}
err = kvStore.AddAccountInvoice(
ctx, acct1.ID, hash1,
)
require.NoError(t, err)
err = kvStore.StoreLastIndexes(ctx, 1, 0)
require.NoError(t, err)
hash2 := lntypes.Hash{2, 2, 2, 2}
err = kvStore.AddAccountInvoice(
ctx, acct2.ID, hash2,
)
require.NoError(t, err)
err = kvStore.StoreLastIndexes(ctx, 2, 0)
require.NoError(t, err)
// Create payments for both accounts.
hash3 := lntypes.Hash{3, 3, 3, 3}
known, err := kvStore.UpsertAccountPayment(
ctx, acct1.ID, hash3, 100,
lnrpc.Payment_SUCCEEDED,
)
require.NoError(t, err)
require.False(t, known)
hash4 := lntypes.Hash{4, 4, 4, 4}
known, err = kvStore.UpsertAccountPayment(
ctx, acct2.ID, hash4, 200,
lnrpc.Payment_IN_FLIGHT,
)
require.NoError(t, err)
require.False(t, known)
},
},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
t.Parallel()
// Create a new kvdb store to populate with test data.
kvStore, err := NewBoltStore(
t.TempDir(), DBFilename, clock,
)
require.NoError(t, err)
t.Cleanup(func() {
require.NoError(t, kvStore.db.Close())
})
// Populate the kv store.
test.populateDB(t, kvStore)
// Create the SQL store that we will migrate the data
// to.
sqlStore, txEx := makeSQLDB(t)
// We fetch the accounts and indices from the kvStore
// before migrating them to the SQL store, just to
// ensure that the migration doesn't affect the original
// data.
kvAccounts, err := kvStore.Accounts(ctx)
require.NoError(t, err)
kvAddIndex, kvSettleIndex, err := kvStore.LastIndexes(
ctx,
)
if !test.expectLastIndex {
// If the test expects there to be no invoices
// indices, we also verify that the database
// contains none.
require.ErrorIs(t, err, ErrNoInvoiceIndexKnown)
} else {
require.NoError(t, err)
}
// Perform the migration.
var opts sqldb.MigrationTxOptions
err = txEx.ExecTx(ctx, &opts,
func(tx SQLQueries) error {
return MigrateAccountStoreToSQL(
ctx, kvStore, tx,
)
},
)
require.NoError(t, err)
// Assert migration results.
assertMigrationResults(
t, sqlStore, kvAccounts, kvAddIndex,
kvSettleIndex, test.expectLastIndex,
)
})
}
}

4
go.mod
View file

@ -35,11 +35,13 @@ require (
github.com/lightningnetwork/lnd/fn v1.2.3
github.com/lightningnetwork/lnd/fn/v2 v2.0.8
github.com/lightningnetwork/lnd/kvdb v1.4.16
github.com/lightningnetwork/lnd/sqldb v1.0.9
github.com/lightningnetwork/lnd/tlv v1.3.0
github.com/lightningnetwork/lnd/tor v1.1.6
github.com/mwitkow/go-conntrack v0.0.0-20190716064945-2f068394615f
github.com/mwitkow/grpc-proxy v0.0.0-20230212185441-f345521cb9c9
github.com/ory/dockertest/v3 v3.10.0
github.com/pmezard/go-difflib v1.0.0
github.com/stretchr/testify v1.10.0
github.com/urfave/cli v1.22.14
go.etcd.io/bbolt v1.3.11
@ -143,7 +145,6 @@ require (
github.com/lightningnetwork/lightning-onion v1.2.1-0.20240712235311-98bd56499dfb // indirect
github.com/lightningnetwork/lnd/healthcheck v1.2.6 // indirect
github.com/lightningnetwork/lnd/queue v1.1.1 // indirect
github.com/lightningnetwork/lnd/sqldb v1.0.9 // indirect
github.com/lightningnetwork/lnd/ticker v1.1.1 // indirect
github.com/ltcsuite/ltcd v0.0.0-20190101042124-f37f8bf35796 // indirect
github.com/mattn/go-isatty v0.0.20 // indirect
@ -161,7 +162,6 @@ require (
github.com/opencontainers/image-spec v1.1.0 // indirect
github.com/opencontainers/runc v1.2.0 // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
github.com/prometheus/client_golang v1.14.0 // indirect
github.com/prometheus/client_model v0.4.0 // indirect
github.com/prometheus/common v0.37.0 // indirect