Module: PgLedger

Defined in:
lib/pg_ledger.rb,
lib/pg_ledger.rb,
lib/pg_ledger/batch.rb,
lib/pg_ledger/models.rb,
lib/pg_ledger/posting.rb,
lib/pg_ledger/railtie.rb,
lib/pg_ledger/version.rb,
lib/pg_ledger/ledger_scope.rb,
lib/pg_ledger/entry_builder.rb,
lib/generators/pg_ledger/install/install_generator.rb

Defined Under Namespace

Modules: Generators Classes: Account, AccountDay, Balance, Batch, CrossLedger, Entry, EntryBuilder, Error, FrozenAccount, IdempotencyConflict, ImmutableRecord, InsufficientBalance, InvalidReversal, Ledger, LedgerScope, Line, OverReversal, PeriodClosed, Posting, Railtie, Record, Unbalanced, UnknownAccount

Constant Summary collapse

DEFAULT_LEDGER_ID =
1
VERSION =
"0.1.0"

Class Attribute Summary collapse

Class Method Summary collapse

Class Attribute Details

.trading_account_code ⇒ Object



266
267
268
# File 'lib/pg_ledger.rb', line 266

def 
   ||= "trading"
end

Class Method Details

.autoscale_shards!(up_at: 12_000, down_at: up_at / 4, window: 60, max_shards: 64) ⇒ Object

One pass of shard autoscaling: doubles shards on accounts whose posting rate over the window exceeds up_at (per minute), halves when below down_at. Hysteresis (default down_at = up_at / 4) prevents flapping. Only accounts without min_balance participate. Run it from a periodic job, like rollup!. Returns the applied changes.



229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
# File 'lib/pg_ledger.rb', line 229

def autoscale_shards!(up_at: 12_000, down_at: up_at / 4, window: 60, max_shards: 64)
  rates = Record.connection.select_rows(Record.sanitize_sql_array(["    SELECT account_id, count(*) * (60.0 / :window) AS per_minute FROM (\n      SELECT debit_account_id AS account_id FROM pg_ledger_lines\n      WHERE posted_at >= now() - make_interval(secs => :window)\n      UNION ALL\n      SELECT credit_account_id FROM pg_ledger_lines\n      WHERE posted_at >= now() - make_interval(secs => :window)\n    ) t GROUP BY 1\n  SQL\n  rates = rates.to_h { |id, rate| [id, rate.to_f] }\n\n  changes = []\n  # Busy accounts are candidates to grow; every sharded account is a\n  # candidate to shrink \u2014 a fully idle one has no rows in the window.\n  candidates = Account.where(min_balance: nil)\n    .where(\"id IN (?) OR balance_shards > 1\", rates.keys.presence || [0])\n  candidates.find_each do |account|\n    rate = rates.fetch(account.id, 0.0)\n    from = account.balance_shards\n    target =\n      if rate >= up_at && from < max_shards\n        [from * 2, max_shards].min\n      elsif rate < down_at && from > 1\n        from / 2\n      end\n    next unless target\n\n    resize_shards!(account, target)\n    changes << { account: account.code, currency: account.currency,\n                 from: from, to: target, per_minute: rate.round(1) }\n  end\n  changes\nend\n", { window: window }]))

.balances(owner: nil, ledger_id: nil, at: nil) ⇒ Object

Balances per currency, optionally for one owner's accounts / one ledger.

Raises:

  • (ArgumentError)


112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
# File 'lib/pg_ledger.rb', line 112

def balances(owner: nil, ledger_id: nil, at: nil)
  unless at
    scope = Balance.joins(:account)
    scope = scope.where(pg_ledger_accounts: { owner: owner }) if owner
    scope = scope.where(pg_ledger_accounts: { ledger_id: ledger_id }) if ledger_id
    return scope.group("pg_ledger_accounts.currency").sum(:amount)
  end

  raise ArgumentError, "owner is not supported with at: — resolve accounts and use balance(at:)" if owner

  # As-of balances: daily rollups below the global watermark plus the line
  # tail up to :at. rollup! advances every account to one shared boundary,
  # so MAX(day) + 1 is a valid watermark for the whole table. Bound as a
  # literal — a subquery here blinds the planner (see Account#balance).
  wm_day = AccountDay.maximum(:day)
  wm_next = wm_day ? wm_day + 1 : Date.new(1970, 1, 1)
  ledger_cond = ledger_id ? "AND a.ledger_id = :ledger_id" : ""
  sql = "    SELECT currency, SUM(v)::bigint FROM (\n      SELECT a.currency,\n             SUM(CASE WHEN a.normal_balance = 'debit'\n                      THEN d.debits - d.credits ELSE d.credits - d.debits END) AS v\n      FROM pg_ledger_account_days d\n      JOIN pg_ledger_accounts a ON a.id = d.account_id\n      WHERE d.day < LEAST(:wm_next::date, :at::timestamptz::date) \#{ledger_cond}\n      GROUP BY a.currency\n      UNION ALL\n      SELECT a.currency,\n             SUM(CASE WHEN a.normal_balance = 'debit' THEN l.amount ELSE -l.amount END)\n      FROM pg_ledger_lines l\n      JOIN pg_ledger_accounts a ON a.id = l.debit_account_id\n      WHERE l.posted_at >= LEAST(:wm_next::date, :at::timestamptz::date)::timestamptz\n        AND l.posted_at <= :at \#{ledger_cond}\n      GROUP BY a.currency\n      UNION ALL\n      SELECT a.currency,\n             SUM(CASE WHEN a.normal_balance = 'credit' THEN l.amount ELSE -l.amount END)\n      FROM pg_ledger_lines l\n      JOIN pg_ledger_accounts a ON a.id = l.credit_account_id\n      WHERE l.posted_at >= LEAST(:wm_next::date, :at::timestamptz::date)::timestamptz\n        AND l.posted_at <= :at \#{ledger_cond}\n      GROUP BY a.currency\n    ) t GROUP BY currency ORDER BY currency\n  SQL\n  Record.connection.select_rows(Record.sanitize_sql_array(\n    [sql, { at: at, wm_next: wm_next, ledger_id: ledger_id }]\n  )).to_h { |currency, amount| [currency, Integer(amount || 0)] }\nend\n"

.batch {|collector| ... } ⇒ Object

N postings in one statement and one fsync. Atomic as a whole.

Yields:

  • (collector)


301
302
303
304
305
# File 'lib/pg_ledger.rb', line 301

def batch
  collector = Batch.new
  yield collector
  collector.commit!
end

.build_fx_pair(entry, account, trading, amount, decrease:) ⇒ Object



95
96
97
98
99
100
101
102
103
104
# File 'lib/pg_ledger.rb', line 95

def build_fx_pair(entry, , trading, amount, decrease:)
  direction =
    if decrease
      .normal_balance == "debit" ? :credit : :debit
    else
      .normal_balance == "debit" ? :debit : :credit
    end
  entry.public_send(direction, , amount)
  entry.public_send(direction == :debit ? :credit : :debit, trading, amount)
end

.build_transfer(entry, from, to, amount) ⇒ Object



106
107
108
109
# File 'lib/pg_ledger.rb', line 106

def build_transfer(entry, from, to, amount)
  entry.public_send(from.normal_balance == "debit" ? :credit : :debit, from, amount)
  entry.public_send(to.normal_balance == "debit" ? :debit : :credit, to, amount)
end

.close_period!(before:, ledger_id: DEFAULT_LEDGER_ID) ⇒ Object

Entries with posted_at earlier than before are rejected from now on.



337
338
339
340
341
342
343
# File 'lib/pg_ledger.rb', line 337

def close_period!(before:, ledger_id: DEFAULT_LEDGER_ID)
  Record.connection.execute(Record.sanitize_sql_array(["    INSERT INTO pg_ledger_config (ledger_id, closed_before) VALUES (:ledger, :before)\n    ON CONFLICT (ledger_id) DO UPDATE SET closed_before = EXCLUDED.closed_before\n  SQL\n  nil\nend\n", { before: before, ledger: ledger_id }]))

.create_account!(code:, currency:, normal_balance:, owner: nil, min_balance: 0, balance_shards: 1, metadata: {}, ledger_id: DEFAULT_LEDGER_ID) ⇒ Object



375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
# File 'lib/pg_ledger.rb', line 375

def create_account!(code:, currency:, normal_balance:, owner: nil, min_balance: 0,
                    balance_shards: 1, metadata: {}, ledger_id: DEFAULT_LEDGER_ID)
  if min_balance && balance_shards > 1
    raise ArgumentError, "sharded balances cannot enforce min_balance — pass min_balance: nil"
  end

  .transaction do
     = .create!(
      ledger_id: ledger_id,
      code: code,
      currency: currency.to_s.upcase,
      normal_balance: normal_balance.to_s,
      owner: owner,
      min_balance: min_balance,
      balance_shards: balance_shards,
      metadata: 
    )
    Balance.insert_all(
      balance_shards.times.map { |shard| { account_id: .id, shard: shard, amount: 0 } }
    )
    
  end
end

.freeze_account!(account, ledger_id: DEFAULT_LEDGER_ID) ⇒ Object



328
329
330
# File 'lib/pg_ledger.rb', line 328

def freeze_account!(, ledger_id: DEFAULT_LEDGER_ID)
  .resolve!(, ledger_id: ledger_id).update!(frozen_at: Time.current)
end

.ledger(name = "main", owner: nil) ⇒ Object

Scoped facade: every operation of the module, pinned to one ledger. billing = PgLedger.ledger("billing") billing.transfer!(from: "a", to: "b", amount: 100)



164
165
166
167
# File 'lib/pg_ledger.rb', line 164

def ledger(name = "main", owner: nil)
  record = Ledger.find_or_create_by!(name: name.to_s) { |l| l.owner = owner }
  LedgerScope.new(record)
end

.merge_excess_shards!(account) ⇒ Object

Idempotent: folds balances of shards >= balance_shards into shard 0. Locks rows one at a time in ascending shard order — same global lock order as posting, so it cannot deadlock against apply_deltas.



201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
# File 'lib/pg_ledger.rb', line 201

def merge_excess_shards!()
   = .resolve!()
  Record.transaction do
    moved = 0
    Balance.where(account_id: .id).where("shard >= ? AND shard > 0", .balance_shards)
      .order(:shard).pluck(:shard).each_with_index do |shard, i|
      if i.zero?
        Balance.where(account_id: .id, shard: 0).lock.pick(:id)
      end
      row = Balance.where(account_id: .id, shard: shard).lock.first
      next if row.amount.zero?

      moved += row.amount
      row.update!(amount: 0)
    end
    if moved != 0
      Balance.where(account_id: .id, shard: 0)
        .update_all(["amount = amount + ?, updated_at = now()", moved])
    end
  end
  nil
end

.post!(idempotency_key: nil, metadata: {}, posted_at: nil, reverses: nil, ledger_id: DEFAULT_LEDGER_ID) {|builder| ... } ⇒ Object

reverses: link this entry as a (possibly partial) reversal of another. The database enforces that per account, the sum of all reversals never exceeds the original legs, and that legs oppose original directions.

Yields:

  • (builder)


32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
# File 'lib/pg_ledger.rb', line 32

def post!(idempotency_key: nil, metadata: {}, posted_at: nil, reverses: nil,
          ledger_id: DEFAULT_LEDGER_ID)
  builder = EntryBuilder.new
  yield builder
  posting = Posting.new(
    legs: builder.legs,
    idempotency_key: idempotency_key,
    metadata: ,
    posted_at: posted_at,
    reverses: reverses.is_a?(Entry) ? reverses.id : reverses,
    ledger_id: ledger_id
  )
  ActiveSupport::Notifications.instrument("post.pg_ledger") do |payload|
    started_at = Time.now
    entry = posting.call
    payload.merge!(
      entry_id: entry.id,
      legs: builder.legs.size,
      idempotency_key: idempotency_key,
      reverses_entry_id: entry.reverses_entry_id,
      # An entry that predates this call is an idempotent replay, not a post.
      replayed: idempotency_key ? entry.created_at < started_at - 1 : false
    )
    entry
  end
end

.resize_shards!(account, target) ⇒ Object

Changes the shard count of an account's balance. Expansion is instant and safe under concurrent postings (a posting holding a stale shard count still hits a valid row; the balance is always the SUM of all rows). Shrinking lowers the count so new postings use the narrow range, then merges leftovers on excess shards into shard 0 — concurrent postings may land on an excess shard with a stale count; a later merge (or the next autoscale pass) folds them in. The sum is correct at every moment.

Raises:

  • (ArgumentError)


177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
# File 'lib/pg_ledger.rb', line 177

def resize_shards!(, target)
   = .resolve!()
  raise ArgumentError, "shard count must be >= 1, got #{target}" if target < 1
  if .min_balance && target > 1
    raise ArgumentError, "sharded balances cannot enforce min_balance"
  end
  return  if target == .balance_shards

  .transaction do
    if target > .balance_shards
      Balance.insert_all(
        (.balance_shards...target).map { |s| { account_id: .id, shard: s, amount: 0 } },
        unique_by: [:account_id, :shard]
      )
    end
    .update!(balance_shards: target)
  end
  merge_excess_shards!()
  
end

.reverse!(entry, amount: nil, idempotency_key: nil, metadata: {}) ⇒ Object

Full reversal mirrors every leg. Partial reversal (amount:) is only defined for two-leg entries — reversing part of a multi-leg entry means deciding which legs shrink, and that is the caller's call: build it with post!(reverses: entry) and explicit legs.



311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
# File 'lib/pg_ledger.rb', line 311

def reverse!(entry, amount: nil, idempotency_key: nil, metadata: {})
  lines = entry.lines.includes(:debit_account, :credit_account).to_a
  if amount && lines.size != 1
    raise ArgumentError,
      "amount: is only supported for single-pair entries; use post!(reverses:) with explicit legs"
  end

  post!(idempotency_key: idempotency_key, metadata: , reverses: entry.id,
        ledger_id: entry.ledger_id) do |reversal|
    lines.each do |line|
      # Storno mirrors the pair: money flows back the way it came.
      reversal.debit(line., amount || line.amount)
      reversal.credit(line., amount || line.amount)
    end
  end
end

.rollup!(upto: Date.today) ⇒ Object

Rolls completed days into pg_ledger_account_days. Idempotent; run from a scheduled job. Days from the current watermark up to (not including) upto are aggregated from lines in one pass.



348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
# File 'lib/pg_ledger.rb', line 348

def rollup!(upto: Date.today)
  Record.transaction do
    # The backfill aggregates the whole lines table; without enough
    # work_mem the hash aggregate spills gigabytes to disk.
    Record.connection.execute("SET LOCAL work_mem = '256MB'")
    Record.connection.execute(Record.sanitize_sql_array(["    INSERT INTO pg_ledger_account_days (account_id, day, debits, credits)\n    SELECT t.account_id, t.day, SUM(t.d), SUM(t.c) FROM (\n      SELECT l.debit_account_id AS account_id, date_trunc('day', l.posted_at)::date AS day,\n             l.amount AS d, 0 AS c\n      FROM pg_ledger_lines l\n      WHERE l.posted_at >= (SELECT COALESCE(max(day), '-infinity'::date) FROM pg_ledger_account_days) + interval '1 day'\n        AND l.posted_at < :upto::date\n      UNION ALL\n      SELECT l.credit_account_id, date_trunc('day', l.posted_at)::date, 0, l.amount\n      FROM pg_ledger_lines l\n      WHERE l.posted_at >= (SELECT COALESCE(max(day), '-infinity'::date) FROM pg_ledger_account_days) + interval '1 day'\n        AND l.posted_at < :upto::date\n    ) t\n    GROUP BY 1, 2\n    ON CONFLICT (account_id, day) DO UPDATE\n      SET debits = EXCLUDED.debits, credits = EXCLUDED.credits\n    SQL\n  end\n  nil\nend\n", { upto: upto }]))

.table_name_prefix ⇒ Object



4
5
6
# File 'lib/pg_ledger/models.rb', line 4

def self.table_name_prefix
  "pg_ledger_"
end

.transfer!(from:, to:, amount:, idempotency_key: nil, metadata: {}, posted_at: nil, rate: nil, to_amount: nil, ledger_id: DEFAULT_LEDGER_ID) ⇒ Object

Decreases the natural balance of from and increases the natural balance of to. Same-currency transfers balance only for accounts of the same polarity — mixed movements (deposits, issuance) need post!.

Cross-currency transfers route through per-currency trading accounts (one balanced pair of legs per currency). Pass either the exact to_amount, or rate as a String/Rational — never a Float — which is rounded with banker's rounding and recorded in the entry metadata.



67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
# File 'lib/pg_ledger.rb', line 67

def transfer!(from:, to:, amount:, idempotency_key: nil, metadata: {}, posted_at: nil,
              rate: nil, to_amount: nil, ledger_id: DEFAULT_LEDGER_ID)
  from = .resolve!(from, ledger_id: ledger_id)
  to = .resolve!(to, ledger_id: ledger_id)
  # Account objects carry their own ledger; it wins over the keyword.
  ledger_id = from.ledger_id

  if from.currency == to.currency
    raise ArgumentError, "rate/to_amount are for cross-currency transfers" if rate || to_amount
    return post!(idempotency_key: idempotency_key, metadata: , posted_at: posted_at,
                 ledger_id: ledger_id) do |entry|
      build_transfer(entry, from, to, amount)
    end
  end

  to_amount ||= convert(amount, rate)
   = .merge(
    "fx" => { "rate" => rate&.to_s, "from_amount" => amount, "from_currency" => from.currency,
              "to_amount" => to_amount, "to_currency" => to.currency }
  )
  post!(idempotency_key: idempotency_key, metadata: , posted_at: posted_at,
        ledger_id: ledger_id) do |entry|
    # One balanced pair per currency: the trading leg mirrors the user leg.
    build_fx_pair(entry, from, (from.currency, ledger_id: ledger_id), amount, decrease: true)
    build_fx_pair(entry, to, (to.currency, ledger_id: ledger_id), to_amount, decrease: false)
  end
end

.unfreeze_account!(account, ledger_id: DEFAULT_LEDGER_ID) ⇒ Object



332
333
334
# File 'lib/pg_ledger.rb', line 332

def unfreeze_account!(, ledger_id: DEFAULT_LEDGER_ID)
  .resolve!(, ledger_id: ledger_id).update!(frozen_at: nil)
end