Class: PgLedger::Batch

Inherits:
Object
  • Object
show all
Defined in:
lib/pg_ledger/batch.rb

Overview

Collects postings and commits them as ONE statement / one fsync. Atomic: if any posting fails, the whole batch rolls back.

Constant Summary collapse

CALL_SQL =

Typed arrays instead of one big jsonb: the server unnests pre-parsed values straight into the batch temp tables — no jsonb walking.

<<~SQL
  SELECT * FROM pg_ledger_post_batch(
    $1::bigint[], $2::text[], $3::text[], $4::jsonb[], $5::timestamptz[], $6::bigint[],
    $7::int[], $8::bigint[], $9::bigint[], $10::text[], $11::bigint[],
    $12::int[], $13::bigint[], $14::int[], $15::bigint[], $16::bigint[])
SQL

Instance Method Summary collapse

Constructor Details

#initialize(default_ledger_id: DEFAULT_LEDGER_ID) ⇒ Batch

Returns a new instance of Batch.



16
17
18
19
20
21
22
23
# File 'lib/pg_ledger/batch.rb', line 16

def initialize(default_ledger_id: DEFAULT_LEDGER_ID)
  @default_ledger_id = default_ledger_id
  @postings = []
  # Shard affinity: every posting in this batch hits the same shard per
  # account, so the batch pre-locks one row per account and concurrent
  # batches (with different hints) don't contend.
  @shard_hint = rand
end

Instance Method Details

#commit! ⇒ Object



51
52
53
54
55
56
57
58
59
60
61
62
63
# File 'lib/pg_ledger/batch.rb', line 51

def commit!
  return [] if @postings.empty?

  encoder = PG::TextEncoder::Array.new
  result = Entry.connection.raw_connection.exec_params(
    CALL_SQL, column_arrays.map { |a| encoder.encode(a) }
  )
  entries = result.map { |row| Entry.instantiate(row) }
  result.clear
  entries
rescue PG::Error, ActiveRecord::StatementInvalid => e
  raise Posting.map_pg_error(e)
end

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

Yields:

  • (builder)


25
26
27
28
29
30
31
32
33
34
35
36
37
38
# File 'lib/pg_ledger/batch.rb', line 25

def post!(idempotency_key: nil, metadata: {}, posted_at: nil, reverses: nil,
          ledger_id: nil)
  ledger_id ||= @default_ledger_id
  builder = EntryBuilder.new
  yield builder
  @postings << Posting.new(
    legs: builder.legs, idempotency_key: idempotency_key,
    metadata: , posted_at: posted_at,
    reverses: reverses.is_a?(Entry) ? reverses.id : reverses,
    shard_hint: @shard_hint,
    ledger_id: ledger_id
  )
  nil
end

#transfer!(from:, to:, amount:, idempotency_key: nil, metadata: {}, posted_at: nil, ledger_id: nil) ⇒ Object



40
41
42
43
44
45
46
47
48
49
# File 'lib/pg_ledger/batch.rb', line 40

def transfer!(from:, to:, amount:, idempotency_key: nil, metadata: {}, posted_at: nil,
              ledger_id: nil)
  ledger_id ||= @default_ledger_id
  from = Account.resolve!(from, ledger_id: ledger_id)
  to = Account.resolve!(to, ledger_id: ledger_id)
  post!(idempotency_key: idempotency_key, metadata: , posted_at: posted_at,
        ledger_id: from.ledger_id) do |entry|
    PgLedger.build_transfer(entry, from, to, amount)
  end
end