feat(ledger): add append-only double-entry engine
The Lightning bridge is modelled as the boundary with the outside world and is the one account permitted to go negative; its negative balance is exactly what is owed to players inside the system. All other accounts are floored at zero by both the application and a database trigger. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
247
pkg/ledger/ledger.go
Normal file
247
pkg/ledger/ledger.go
Normal file
@@ -0,0 +1,247 @@
|
||||
// Package ledger implements append-only double-entry accounting.
|
||||
//
|
||||
// Invariants, enforced here and again by database constraints and triggers:
|
||||
// - every transaction's postings sum to exactly zero
|
||||
// - no account balance may go negative
|
||||
// - rows are never updated or deleted; corrections are compensating entries
|
||||
//
|
||||
// Every balance change is explained by a posting that records what happened,
|
||||
// when, which round it belonged to, and the balance either side of it.
|
||||
package ledger
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrUnbalanced = errors.New("ledger: postings do not sum to zero")
|
||||
ErrInsufficientFunds = errors.New("ledger: insufficient funds")
|
||||
ErrEmptyTransaction = errors.New("ledger: transaction has no postings")
|
||||
ErrNonPositiveAmount = errors.New("ledger: amount must be positive")
|
||||
)
|
||||
|
||||
// Posting is a single leg of a transaction. Positive credits, negative debits.
|
||||
type Posting struct {
|
||||
AccountID int64
|
||||
AmountMsat int64
|
||||
}
|
||||
|
||||
// Entry is a posting as seen from one account's history.
|
||||
type Entry struct {
|
||||
TransactionID int64
|
||||
Kind string
|
||||
RoundID *int64
|
||||
AmountMsat int64
|
||||
BalanceBefore int64
|
||||
BalanceAfter int64
|
||||
CreatedAt time.Time
|
||||
}
|
||||
|
||||
type Ledger struct{ pool *pgxpool.Pool }
|
||||
|
||||
func New(pool *pgxpool.Pool) *Ledger { return &Ledger{pool: pool} }
|
||||
|
||||
// Post writes one balanced transaction atomically.
|
||||
//
|
||||
// Accounts are locked in ascending id order so that concurrent transactions
|
||||
// touching the same accounts cannot deadlock, and so a balance read cannot be
|
||||
// stale by the time the posting is written.
|
||||
func (l *Ledger) Post(ctx context.Context, kind string, roundID *int64, postings []Posting) (int64, error) {
|
||||
if len(postings) == 0 {
|
||||
return 0, ErrEmptyTransaction
|
||||
}
|
||||
var sum int64
|
||||
for _, p := range postings {
|
||||
sum += p.AmountMsat
|
||||
}
|
||||
if sum != 0 {
|
||||
return 0, fmt.Errorf("%w: sum is %d", ErrUnbalanced, sum)
|
||||
}
|
||||
|
||||
tx, err := l.pool.Begin(ctx)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
defer tx.Rollback(ctx)
|
||||
|
||||
var txID int64
|
||||
if err := tx.QueryRow(ctx,
|
||||
`INSERT INTO transactions (kind, round_id) VALUES ($1, $2) RETURNING id`,
|
||||
kind, roundID).Scan(&txID); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
|
||||
ordered := append([]Posting(nil), postings...)
|
||||
sort.Slice(ordered, func(i, j int) bool {
|
||||
return ordered[i].AccountID < ordered[j].AccountID
|
||||
})
|
||||
|
||||
for _, p := range ordered {
|
||||
// Lock the account row first, then read its latest balance. Taking the
|
||||
// lock before the read is what serializes concurrent spends.
|
||||
var allowNegative bool
|
||||
if err := tx.QueryRow(ctx,
|
||||
`SELECT allow_negative FROM accounts WHERE id = $1 FOR UPDATE`,
|
||||
p.AccountID).Scan(&allowNegative); err != nil {
|
||||
return 0, fmt.Errorf("locking account %d: %w", p.AccountID, err)
|
||||
}
|
||||
|
||||
var before int64
|
||||
if err := tx.QueryRow(ctx,
|
||||
`SELECT COALESCE(
|
||||
(SELECT balance_after FROM postings
|
||||
WHERE account_id = $1 ORDER BY id DESC LIMIT 1), 0)`,
|
||||
p.AccountID).Scan(&before); err != nil {
|
||||
return 0, fmt.Errorf("reading balance of account %d: %w", p.AccountID, err)
|
||||
}
|
||||
|
||||
after := before + p.AmountMsat
|
||||
if after < 0 && !allowNegative {
|
||||
return 0, fmt.Errorf("%w: account %d holds %d, needs %d",
|
||||
ErrInsufficientFunds, p.AccountID, before, -p.AmountMsat)
|
||||
}
|
||||
|
||||
if _, err := tx.Exec(ctx,
|
||||
`INSERT INTO postings
|
||||
(transaction_id, account_id, amount_msat, balance_before, balance_after)
|
||||
VALUES ($1, $2, $3, $4, $5)`,
|
||||
txID, p.AccountID, p.AmountMsat, before, after); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
}
|
||||
|
||||
if err := tx.Commit(ctx); err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return txID, nil
|
||||
}
|
||||
|
||||
// Transfer moves funds between two accounts. This is the peer-to-peer path.
|
||||
func (l *Ledger) Transfer(ctx context.Context, from, to int64, amountMsat int64) (int64, error) {
|
||||
if amountMsat <= 0 {
|
||||
return 0, ErrNonPositiveAmount
|
||||
}
|
||||
return l.Post(ctx, "transfer", nil, []Posting{
|
||||
{AccountID: from, AmountMsat: -amountMsat},
|
||||
{AccountID: to, AmountMsat: amountMsat},
|
||||
})
|
||||
}
|
||||
|
||||
// Deposit credits a player from the Lightning bridge account.
|
||||
func (l *Ledger) Deposit(ctx context.Context, player int64, amountMsat int64) (int64, error) {
|
||||
if amountMsat <= 0 {
|
||||
return 0, ErrNonPositiveAmount
|
||||
}
|
||||
bridge, err := l.AccountByName(ctx, "lightning_bridge")
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return l.Post(ctx, "deposit", nil, []Posting{
|
||||
{AccountID: bridge, AmountMsat: -amountMsat},
|
||||
{AccountID: player, AmountMsat: amountMsat},
|
||||
})
|
||||
}
|
||||
|
||||
// Withdraw debits a player back to the Lightning bridge account.
|
||||
func (l *Ledger) Withdraw(ctx context.Context, player int64, amountMsat int64) (int64, error) {
|
||||
if amountMsat <= 0 {
|
||||
return 0, ErrNonPositiveAmount
|
||||
}
|
||||
bridge, err := l.AccountByName(ctx, "lightning_bridge")
|
||||
if err != nil {
|
||||
return 0, err
|
||||
}
|
||||
return l.Post(ctx, "withdraw", nil, []Posting{
|
||||
{AccountID: player, AmountMsat: -amountMsat},
|
||||
{AccountID: bridge, AmountMsat: amountMsat},
|
||||
})
|
||||
}
|
||||
|
||||
// Balance returns the account's current balance in millisatoshis.
|
||||
func (l *Ledger) Balance(ctx context.Context, accountID int64) (int64, error) {
|
||||
var bal int64
|
||||
err := l.pool.QueryRow(ctx,
|
||||
`SELECT COALESCE(
|
||||
(SELECT balance_after FROM postings
|
||||
WHERE account_id = $1 ORDER BY id DESC LIMIT 1), 0)`,
|
||||
accountID).Scan(&bal)
|
||||
return bal, err
|
||||
}
|
||||
|
||||
// History returns an account's postings, newest first.
|
||||
func (l *Ledger) History(ctx context.Context, accountID int64, limit int) ([]Entry, error) {
|
||||
rows, err := l.pool.Query(ctx,
|
||||
`SELECT p.transaction_id, t.kind, t.round_id,
|
||||
p.amount_msat, p.balance_before, p.balance_after, p.created_at
|
||||
FROM postings p
|
||||
JOIN transactions t ON t.id = p.transaction_id
|
||||
WHERE p.account_id = $1
|
||||
ORDER BY p.id DESC
|
||||
LIMIT $2`, accountID, limit)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
var out []Entry
|
||||
for rows.Next() {
|
||||
var e Entry
|
||||
if err := rows.Scan(&e.TransactionID, &e.Kind, &e.RoundID,
|
||||
&e.AmountMsat, &e.BalanceBefore, &e.BalanceAfter, &e.CreatedAt); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
out = append(out, e)
|
||||
}
|
||||
return out, rows.Err()
|
||||
}
|
||||
|
||||
// EnsurePlayer returns the account id for a public key, creating it if needed.
|
||||
func (l *Ledger) EnsurePlayer(ctx context.Context, pubkey []byte) (int64, error) {
|
||||
var id int64
|
||||
err := l.pool.QueryRow(ctx,
|
||||
`INSERT INTO accounts (kind, pubkey) VALUES ('player', $1)
|
||||
ON CONFLICT (pubkey) DO UPDATE SET pubkey = EXCLUDED.pubkey
|
||||
RETURNING id`, pubkey).Scan(&id)
|
||||
return id, err
|
||||
}
|
||||
|
||||
// AccountByName resolves a system account such as "house_pot".
|
||||
func (l *Ledger) AccountByName(ctx context.Context, name string) (int64, error) {
|
||||
var id int64
|
||||
err := l.pool.QueryRow(ctx,
|
||||
`SELECT id FROM accounts WHERE name = $1`, name).Scan(&id)
|
||||
if errors.Is(err, pgx.ErrNoRows) {
|
||||
return 0, fmt.Errorf("ledger: no account named %q", name)
|
||||
}
|
||||
return id, err
|
||||
}
|
||||
|
||||
// TotalIssued is the value held inside the system by players and the house —
|
||||
// every account except the external Lightning bridge. It changes only when
|
||||
// funds genuinely enter or leave, never through internal play.
|
||||
func (l *Ledger) TotalIssued(ctx context.Context) (int64, error) {
|
||||
var total int64
|
||||
err := l.pool.QueryRow(ctx,
|
||||
`SELECT COALESCE(SUM(b.balance_msat), 0)
|
||||
FROM account_balances b
|
||||
JOIN accounts a ON a.id = b.account_id
|
||||
WHERE NOT a.allow_negative`).Scan(&total)
|
||||
return total, err
|
||||
}
|
||||
|
||||
// ConservationCheck sums every account including the bridge. Because each
|
||||
// transaction sums to zero, this must always be exactly zero. A non-zero
|
||||
// result means the books are corrupt, and is the top-level audit alarm.
|
||||
func (l *Ledger) ConservationCheck(ctx context.Context) (int64, error) {
|
||||
var total int64
|
||||
err := l.pool.QueryRow(ctx,
|
||||
`SELECT COALESCE(SUM(balance_msat), 0) FROM account_balances`).Scan(&total)
|
||||
return total, err
|
||||
}
|
||||
Reference in New Issue
Block a user