Skip to content

fix: potential invalid balance while upgrading #870

New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Merged
merged 1 commit into from
Apr 14, 2025
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 15 additions & 9 deletions internal/storage/ledger/balances.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,13 +21,13 @@ func (store *Store) GetBalances(ctx context.Context, query ledgercontroller.Bala
}

if isUpToDate {
return store.getBalancesAfterUpgrade(ctx, query)
return store.GetBalancesAfterUpgrade(ctx, query)
} else {
return store.getBalancesWhenUpgrading(ctx, query)
return store.GetBalancesWhenUpgrading(ctx, query)
}
}

func (store *Store) getBalancesWhenUpgrading(ctx context.Context, query ledgercontroller.BalanceQuery) (ledgercontroller.Balances, error) {
func (store *Store) GetBalancesWhenUpgrading(ctx context.Context, query ledgercontroller.BalanceQuery) (ledgercontroller.Balances, error) {
return tracing.TraceWithMetric(
ctx,
"GetBalances",
Expand All @@ -46,12 +46,13 @@ func (store *Store) getBalancesWhenUpgrading(ctx context.Context, query ledgerco
type AccountsVolumesWithLedger struct {
Ledger string `bun:"ledger,type:varchar"`
ledger.AccountsVolumes `bun:",extend"`
Priority int `bun:"priority"` // for ordering (keep at 0)
}

accountsVolumes := make([]AccountsVolumesWithLedger, 0)
defaultAccountsVolumes := make([]AccountsVolumesWithLedger, 0)
for account, assets := range query {
for _, asset := range assets {
accountsVolumes = append(accountsVolumes, AccountsVolumesWithLedger{
defaultAccountsVolumes = append(defaultAccountsVolumes, AccountsVolumesWithLedger{
Ledger: store.ledger.Name,
AccountsVolumes: ledger.AccountsVolumes{
Account: account,
Expand All @@ -64,7 +65,7 @@ func (store *Store) getBalancesWhenUpgrading(ctx context.Context, query ledgerco
}

// prevent deadlocks by sorting the accountsVolumes slice
slices.SortStableFunc(accountsVolumes, func(i, j AccountsVolumesWithLedger) int {
slices.SortStableFunc(defaultAccountsVolumes, func(i, j AccountsVolumesWithLedger) int {
if i.Account < j.Account {
return -1
} else if i.Account > j.Account {
Expand Down Expand Up @@ -94,18 +95,22 @@ func (store *Store) getBalancesWhenUpgrading(ctx context.Context, query ledgerco
Column("ledger", "accounts_address", "asset").
ColumnExpr("(post_commit_volumes).inputs as input").
ColumnExpr("(post_commit_volumes).outputs as output").
ColumnExpr("1 as priority").
UnionAll(
store.db.NewSelect().
TableExpr(
"(?) data",
store.db.NewSelect().NewValues(&accountsVolumes),
store.db.NewSelect().
NewValues(&defaultAccountsVolumes),
).
Column("*"),
)

zeroValueOrMoves := store.db.NewSelect().
TableExpr("(?) data", zeroValuesAndMoves).
Column("ledger", "accounts_address", "asset", "input", "output").
Column("ledger", "accounts_address", "asset").
ColumnExpr("first_value(input) over (partition by ledger, accounts_address, asset order by priority desc) as input").
ColumnExpr("first_value(output) over (partition by ledger, accounts_address, asset order by priority desc) as output").
DistinctOn("ledger, accounts_address, asset")

insertDefaultValue := store.db.NewInsert().
Expand All @@ -122,6 +127,7 @@ func (store *Store) getBalancesWhenUpgrading(ctx context.Context, query ledgerco
// notes(gfyrag): Keep order, it ensures consistent locking order and limit deadlocks
Order("accounts_address", "asset")

accountsVolumes := make([]ledger.AccountsVolumes, 0)
finalQuery := store.db.NewSelect().
With("inserted", insertDefaultValue).
With("existing", selectExistingValues).
Expand Down Expand Up @@ -163,7 +169,7 @@ func (store *Store) getBalancesWhenUpgrading(ctx context.Context, query ledgerco
)
}

func (store *Store) getBalancesAfterUpgrade(ctx context.Context, query ledgercontroller.BalanceQuery) (ledgercontroller.Balances, error) {
func (store *Store) GetBalancesAfterUpgrade(ctx context.Context, query ledgercontroller.BalanceQuery) (ledgercontroller.Balances, error) {
return tracing.TraceWithMetric(
ctx,
"GetBalances",
Expand Down
129 changes: 129 additions & 0 deletions internal/storage/ledger/balances_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ package ledger_test

import (
"database/sql"
"github.com/formancehq/go-libs/v2/bun/bunpaginate"
"math/big"
"testing"

Expand Down Expand Up @@ -129,6 +130,134 @@ func TestBalancesGet(t *testing.T) {
require.NoError(t, err)
require.Equal(t, 2, count)
})

t.Run("with balance from move", func(t *testing.T) {
t.Parallel()

tx := ledger.NewTransaction().WithPostings(
ledger.NewPosting("world", "bank", "USD", big.NewInt(100)),
ledger.NewPosting("world", "bank", "EUR", big.NewInt(200)),
)
err := store.InsertTransaction(ctx, &tx)
require.NoError(t, err)

err = store.UpsertAccounts(ctx,
&ledger.Account{
Address: "world",
},
&ledger.Account{
Address: "bank",
},
&ledger.Account{
Address: "not-existing",
},
)
require.NoError(t, err)

err = store.InsertMoves(ctx, &ledger.Move{
TransactionID: *tx.ID,
IsSource: true,
Account: "world",
Amount: (*bunpaginate.BigInt)(big.NewInt(100)),
Asset: "USD",
InsertionDate: tx.InsertedAt,
EffectiveDate: tx.InsertedAt,
PostCommitVolumes: pointer.For(ledger.NewVolumesInt64(0, 100)),
})
require.NoError(t, err)

err = store.InsertMoves(ctx, &ledger.Move{
TransactionID: *tx.ID,
IsSource: false,
Account: "bank",
Amount: (*bunpaginate.BigInt)(big.NewInt(100)),
Asset: "USD",
InsertionDate: tx.InsertedAt,
EffectiveDate: tx.InsertedAt,
PostCommitVolumes: pointer.For(ledger.NewVolumesInt64(100, 0)),
})
require.NoError(t, err)

err = store.InsertMoves(ctx, &ledger.Move{
TransactionID: *tx.ID,
IsSource: true,
Account: "world",
Amount: (*bunpaginate.BigInt)(big.NewInt(200)),
Asset: "EUR",
InsertionDate: tx.InsertedAt,
EffectiveDate: tx.InsertedAt,
PostCommitVolumes: pointer.For(ledger.NewVolumesInt64(0, 200)),
})
require.NoError(t, err)

err = store.InsertMoves(ctx, &ledger.Move{
TransactionID: *tx.ID,
IsSource: false,
Account: "bank",
Amount: (*bunpaginate.BigInt)(big.NewInt(200)),
Asset: "EUR",
InsertionDate: tx.InsertedAt,
EffectiveDate: tx.InsertedAt,
PostCommitVolumes: pointer.For(ledger.NewVolumesInt64(200, 0)),
})
require.NoError(t, err)

balances, err := store.GetBalancesWhenUpgrading(ctx, ledgercontroller.BalanceQuery{
"bank": {"USD"},
"world": {"USD"},
"not-existing": {"USD"},
})
require.NoError(t, err)

require.NotNil(t, balances["bank"])
RequireEqual(t, big.NewInt(100), balances["bank"]["USD"])
RequireEqual(t, big.NewInt(-100), balances["world"]["USD"])
RequireEqual(t, big.NewInt(0), balances["not-existing"]["USD"])

// Check a new line has been inserted into accounts_volumes table
volumes := &ledger.AccountsVolumes{}
err = store.GetDB().NewSelect().
ModelTableExpr(store.GetPrefixedRelationName("accounts_volumes")).
Where("accounts_address = ? and ledger = ? and asset = 'USD'", "bank", store.GetLedger().Name).
Scan(ctx, volumes)
require.NoError(t, err)

RequireEqual(t, big.NewInt(100), volumes.Input)
RequireEqual(t, big.NewInt(0), volumes.Output)

err = store.GetDB().NewSelect().
ModelTableExpr(store.GetPrefixedRelationName("accounts_volumes")).
Where("accounts_address = ? and ledger = ? and asset = 'USD'", "world", store.GetLedger().Name).
Scan(ctx, volumes)
require.NoError(t, err)

RequireEqual(t, big.NewInt(0), volumes.Input)
RequireEqual(t, big.NewInt(100), volumes.Output)

err = store.GetDB().NewSelect().
ModelTableExpr(store.GetPrefixedRelationName("accounts_volumes")).
Where("accounts_address = ? and ledger = ? and asset = 'USD'", "not-existing", store.GetLedger().Name).
Scan(ctx, volumes)
require.NoError(t, err)

RequireEqual(t, big.NewInt(0), volumes.Input)
RequireEqual(t, big.NewInt(0), volumes.Output)

balances, err = store.GetBalancesWhenUpgrading(ctx, ledgercontroller.BalanceQuery{
"bank": {"USD", "EUR"},
"world": {"USD", "EUR"},
"not-existing": {"USD", "EUR"},
})
require.NoError(t, err)

require.NotNil(t, balances["bank"])
RequireEqual(t, big.NewInt(100), balances["bank"]["USD"])
RequireEqual(t, big.NewInt(200), balances["bank"]["EUR"])
RequireEqual(t, big.NewInt(-100), balances["world"]["USD"])
RequireEqual(t, big.NewInt(-200), balances["world"]["EUR"])
RequireEqual(t, big.NewInt(0), balances["not-existing"]["USD"])
RequireEqual(t, big.NewInt(0), balances["not-existing"]["EUR"])
})
}

func TestBalancesAggregates(t *testing.T) {
Expand Down
Loading