Class: Farce::Transaction

Inherits:
Object
  • Object
show all
Includes:
Unshareable, Internal::MarshalSupport::Reject
Defined in:
lib/farce/transaction.rb,
lib/farce/transaction/map.rb,
lib/farce/transaction/set.rb,
lib/farce/transaction/atom.rb,
lib/farce/transaction/vector.rb,
lib/farce/transaction/mutable.rb,
lib/farce/transaction/wrapper.rb,
lib/farce/transaction/molecule.rb,
lib/farce/transaction/tree_map.rb,
lib/farce/transaction/sorted_set.rb,
lib/farce/integrations/concurrent.rb,
lib/farce/integrations/ractor_tmvar.rb,
lib/farce/transaction/map_operations.rb,
lib/farce/transaction/set_operations.rb,
lib/farce/integrations/ractor_sharing.rb

Overview

A transaction groups multiple changes to Farce objects into a single atomic attempt.

These changes either all succeed together when the transaction is committed or are all discarded if the transaction fails.

Transaction attempts are isolated and can be retried safely.

You must therefore make sure that any side-effects not going through transaction wrappers are either completely avoided, or at least made idempotent.

account1 = Farce::Atom.new(10)
account2 = Farce::Atom.new(20)

# Wire 10 from account1 to account2
Farce.transaction do |tx|
  raise "insufficient funds" unless tx[account1].value >= 10
  tx[account1].update { it - 10 }
  tx[account2].update { it + 10 }
end

Transactions can mix different types and variants, and are therefore ideal for coordinating complex, coordinated changes affecting data within multiple scopes.

Supported objects

Out of the box, transactions support instances of the following classes:

  • Atoms, except weak atoms
  • ConcurrentMap subclasses, except weak maps
  • Molecules
  • Mutable wrappers with frozen snapshots
  • Sets, including sorted sets
  • TreeMap subclasses
  • Vectors
  • Concurrent::TVar from the concurrent-ruby gem
  • Ractor::TVar from the ractor-sharing gem
  • Ractor::TMVar from the ractor-tmvar gem

Defined Under Namespace

Classes: Atom, Map, Molecule, Mutable, RactorTMVar, RactorTVar, Set, SortedSet, TVar, TreeMap, Vector

Constant Summary collapse

ClosedError =

Raised when a wrapper is used outside its transaction attempt.

Class.new(StandardError)
OwnershipError =

Raised when another Fiber uses a transaction.

Class.new(StandardError)

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Methods included from Unshareable

#ractor_shareable?

Constructor Details

#initialize ⇒ Transaction

Create an unused attempt bound to the current Fiber. Call #run once to enroll participants and stage writes. Use run when failed comparisons or conflicts should create and run fresh attempts.



147
148
149
150
151
152
153
154
155
156
# File 'lib/farce/transaction.rb', line 147

def initialize
  @owner     = Fiber.current
  @state     = :new
  @failed    = @aborted = false
  @retryable = true
  @wrappers  = {}.compare_by_identity
  @entries   = {}.compare_by_identity
  @guards    = {}.compare_by_identity
  super
end

Instance Attribute Details

#state ⇒ Symbol (readonly)

Report the attempt's lifecycle state.

Returns:

  • (Symbol) —

    :new before run, :active during the block, :committing during publication, then :committed or :failed



161
162
163
# File 'lib/farce/transaction.rb', line 161

def state
  @state
end

Class Method Details

.define(klass) {|object, transaction| ... } ⇒ BasicObject

Note:

This method must be called from the main Ractor, and the block must be convertible to a shareable proc.

Registers a wrapper for a custom class, allowing them to be used within transactions.

Examples:

Account = Struct.new(:balance) do
  def withdraw(amount) = balance.value -= amount
end

Farce::Transaction.define(Account) do |, tx|
  .class.new(tx[.balance])
end

 = Account.new(Farce::Atom.new(10))
Farce.transaction { |tx| tx[].withdraw(3) }
.balance.value # => 7

Parameters:

  • klass (Class) —

    the class to wrap

Yield Parameters:

  • object (Object) —

    the participant

  • transaction (Transaction) —

    the current attempt

Yield Returns:

  • (Object) —

    the transaction wrapper

Raises:

  • (LocalJumpError)


77
78
79
80
81
# File 'lib/farce/transaction.rb', line 77

def self.define(klass, &)
  raise LocalJumpError, "no block given" unless block_given?
  definition = Internal.prepare_method_definition(&)
  REGISTER.define(klass) { define_method(:call, definition) }
end

.run(*objects, retries: nil, backoff_after: 10, max_backoff: 1.0) ⇒ Boolean

Creates and runs a new transaction attempt.

Automatically retries failed attempts indefinitely unless a retry limit is specified. Starts backing off after the specified number of attempts, up to the maximum delay.

Parameters:

  • objects (Array) —

    list of objects to enroll in the transaction

  • retries (Integer, nil) (defaults to: nil) —

    maximum additional attempts, or nil for unlimited retries

  • backoff_after (Integer) (defaults to: 10) —

    number of attempts before starting to back off

  • max_backoff (Numeric) (defaults to: 1.0) —

    maximum backoff delay in seconds

Returns:

  • (Boolean) —

    whether the transaction committed successfully

See Also:

Raises:

  • (LocalJumpError)


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
# File 'lib/farce/transaction.rb', line 113

def self.run(*, retries: nil, backoff_after: 10, max_backoff: 1.0, &) # rubocop:disable Naming/PredicateMethod
  raise LocalJumpError, "no block given" unless block_given?

  unless retries.nil? || (Integer === retries && retries >= 0)
    raise ArgumentError, "retries must be nil or a non-negative Integer"
  end

  unless Integer === backoff_after && backoff_after >= 0
    raise ArgumentError, "backoff_after must be a non-negative Integer"
  end

  attempt = 0

  while retries.nil? || attempt <= retries
    if attempt > backoff_after
      delay     = (attempt - backoff_after) * 0.01
      delay     = max_backoff if delay > max_backoff
      scheduler = Fiber.scheduler if Fiber.respond_to?(:scheduler) && !Fiber.current.blocking?
      scheduler ? scheduler.kernel_sleep(delay) : sleep(delay)
    end

    transaction = new
    return true if transaction.run(*, &)
    return false unless transaction.retryable?

    attempt += 1
  end

  false
end

Instance Method Details

#[](object) ⇒ Object

Includes an object in this transaction and returns a transaction-aware version of it. Read from and write to the returned object, so the transaction can keep track of inputs and changes.

Parameters:

  • object (Object) —

    the participant or a wrapper from this attempt

Returns:

  • (Object) —

    a wrapper valid only during this attempt

Raises:

  • (TypeError) —

    if the participant is unsupported or belongs to another attempt



229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
# File 'lib/farce/transaction.rb', line 229

def [](object)
  check_open!
  @wrappers.fetch(object) do
    klass = Kernel.instance_method(:class).bind_call(object)
    @wrappers[object] = REGISTER[klass].call(object, self)
  end
rescue Internal::TransactionConflict
  fail! if @state == :active && @owner.equal?(Fiber.current)
  raise
rescue Exception # rubocop:disable Lint/RescueException -- cancellation must poison the attempt too
  if @state == :active && @owner.equal?(Fiber.current)
    @failed = true
    @retryable = false
  end
  raise
end

#abort!

This method returns an undefined value.

Stop the transaction block and discard all staged writes without retrying. The surrounding run call returns false. No writes to participants are published.

Examples:

balance = Farce::Atom.new(100)

Farce.transaction do |tx|
  tx.abort! if tx[balance] < 100
  tx[balance].value -= 100
end

Raises:

  • (Aborted)


257
258
259
260
261
# File 'lib/farce/transaction.rb', line 257

def abort!
  check_open!
  @failed = @aborted = true
  raise Aborted, "transaction aborted"
end

#aborted? ⇒ Boolean

Whether abort! was called during this attempt.

Returns:

  • (Boolean)


165
# File 'lib/farce/transaction.rb', line 165

def aborted? = @aborted

#fail!(retryable: true) ⇒ false

Mark a failed condition, even if the caller ignores its return value.

The block may continue, but commit will discard its staged writes. Custom wrappers can use this for their own conditional operations.

Wrapper operations call this automatically when they raise. If you rescue an error before a wrapper is called, such as an argument computation error, call fail!(retryable: false) to prevent earlier writes from committing.

Parameters:

  • retryable (Boolean) (defaults to: true) —

    whether .run may repeat the block for this failure

Returns:

  • (false)


274
275
276
277
278
279
# File 'lib/farce/transaction.rb', line 274

def fail!(retryable: true) # rubocop:disable Naming/PredicateMethod
  check_open!
  @failed = true
  @retryable &&= retryable
  false
end

#retryable? ⇒ Boolean

Whether .run may retry this attempt after a failed commit or comparison. Explicit aborts, wrapper exceptions, and fail!(retryable: false) disable retries.

Returns:

  • (Boolean)


170
# File 'lib/farce/transaction.rb', line 170

def retryable? = @retryable && !@aborted

#run(*objects) {|transaction, *objects| ... } ⇒ Boolean

Execute this attempt once. Use run for automatic retries.

The block stages work through #[]. After it returns, commit validates the attempt's reads and publishes its staged writes together. Cancellation during commit is deferred until publication finishes. Once committed, an interruption does not undo the changes.

Parameters:

  • objects (Array) —

    the objects to enroll in this transaction attempt

Yields:

  • (transaction, *objects) —

    the current transaction and the enrolled objects

Yield Parameters:

  • transaction (Transaction) —

    this active attempt

  • objects (Array) —

    the enrolled objects

Returns:

  • (Boolean) —

    whether the changes committed

Raises:



186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
# File 'lib/farce/transaction.rb', line 186

def run(*objects)
  started = false
  raise LocalJumpError, "no block given" unless block_given?
  check_owner!
  raise ClosedError, "transaction has already run" unless @state == :new

  @state  = :active
  started = true

  objects.map! { self[it] }
  yield self, *objects

  check_open!

  Thread.handle_interrupt(Internal::INTERRUPT_MASK) do
    @state = :committing
    entries = @entries.values
    unless @failed
      entries = entries.filter_map { Internal::TransactionMapSnapshot === it ? it.prepare_commit : it }
    end
    committed = !@failed && Internal.commit_transaction(entries, @guards.values, self)
    @state    = committed ? :committed : :failed
    Internal.notify_transaction(entries) if committed && entries.none? { it.respond_to?(:commit_group) } &&
      Internal.respond_to?(:notify_transaction)
    committed
  end
rescue Aborted, Internal::TransactionConflict
  false
ensure
  if started
    @state = :failed if %i[active committing].include?(@state)
    Thread.handle_interrupt(Internal::INTERRUPT_MASK) do
      @entries.each_value { it.release if it.respond_to?(:release) }
    end
  end
end