Class: SolidObjects::SynchronousInvocation
- Inherits:
-
Object
- Object
- SolidObjects::SynchronousInvocation
- Defined in:
- lib/solid_objects/synchronous_invocation.rb,
sig/generated/lib/solid_objects/synchronous_invocation.rbs
Instance Method Summary collapse
- #assist(message, deadline:) ⇒ Integer
- #call(message_reference, timeout:) ⇒ Object
- #call_before_deadline(message_reference, timeout:) ⇒ Object
- #completed_result(message) ⇒ Object
- #contention_timeout(message, message_reference, timeout:) ⇒ SyncTimeout
- #diagnose_timeout(message, timeout:) ⇒ SyncTimeout
-
#initialize(process_registry: nil) ⇒ SynchronousInvocation
constructor
A new instance of SynchronousInvocation.
- #load_message(message_reference) ⇒ Message
- #load_message_at_deadline(message_reference) ⇒ Message?
- #monotonic_now ⇒ Float
- #process_registry ⇒ ProcessRegistry
- #raise_rejection(message) ⇒ bot
- #wait(remaining) ⇒ void
Constructor Details
#initialize(process_registry: nil) ⇒ SynchronousInvocation
Returns a new instance of SynchronousInvocation.
12 13 14 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 12 def initialize(process_registry: nil) @dedicated_process_registry = process_registry end |
Instance Method Details
#assist(message, deadline:) ⇒ Integer
111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 111 def assist(, deadline:) process_registry = self.process_registry activation = ActivationManager .new(owner_id: process_registry.process_record.id) .claim(instance_id: .instance_id) return 0 unless activation LeaseRenewer.new( activation:, process_registry: ).around do activation.drain_until(message_id: .id, deadline:) end ensure if activation activation. if activation.pass_exhausted? activation.deactivate end end |
#call(message_reference, timeout:) ⇒ Object
17 18 19 20 21 22 23 24 25 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 17 def call(, timeout:) return call_before_deadline(, timeout:) if SyncDeadline.active? SyncDeadline.with(timeout:) do call_before_deadline(, timeout:) end rescue ActiveRecord::RecordNotFound raise ActorDestroyed, "actor was destroyed while waiting for its result" end |
#call_before_deadline(message_reference, timeout:) ⇒ Object
30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 30 def call_before_deadline(, timeout:) deadline = monotonic_now + SyncDeadline.remaining = nil loop do unless (deadline - monotonic_now).positive? = () return completed_result() if &.completed? || &.dead? raise ? diagnose_timeout(, timeout:) : contention_timeout(, , timeout:) end = () return completed_result() if .completed? || .dead? processed = assist(, deadline:) remaining = deadline - monotonic_now wait(remaining) if processed.zero? && remaining.positive? end rescue DatabaseDeadlineExceeded raise ? diagnose_timeout(, timeout:) : contention_timeout(, , timeout:) end |
#completed_result(message) ⇒ Object
78 79 80 81 82 83 84 85 86 87 88 89 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 78 def completed_result() raise_rejection() if .rejected? if .dead? raise MessageFailed.new( "actor message failed permanently", message_id: .id, details: .error || {} ) end Serialization.readonly_copy(.result) end |
#contention_timeout(message, message_reference, timeout:) ⇒ SyncTimeout
148 149 150 151 152 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 148 def contention_timeout(, , timeout:) return SyncDiagnostics.new.database_contention(, timeout:) if SyncDiagnostics.new.database_contention_for(, timeout:) end |
#diagnose_timeout(message, timeout:) ⇒ SyncTimeout
139 140 141 142 143 144 145 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 139 def diagnose_timeout(, timeout:) SolidObjects.database_adapter.with_lock_probe do SyncDiagnostics.new.call(, timeout:) end rescue DatabaseDeadlineExceeded SyncDiagnostics.new.database_contention(, timeout:) end |
#load_message(message_reference) ⇒ Message
58 59 60 61 62 63 64 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 58 def () SolidObjects.database_adapter.with_lock_retry do Message.uncached do Message.includes(:dead_letter).find(.id) end end end |
#load_message_at_deadline(message_reference) ⇒ Message?
67 68 69 70 71 72 73 74 75 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 67 def () SolidObjects.database_adapter.with_lock_probe do Message.uncached do Message.includes(:dead_letter).find(.id) end end rescue DatabaseDeadlineExceeded nil end |
#monotonic_now ⇒ Float
155 156 157 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 155 def monotonic_now ::Process.clock_gettime(::Process::CLOCK_MONOTONIC) end |
#process_registry ⇒ ProcessRegistry
103 104 105 106 107 108 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 103 def process_registry dedicated_registry = @dedicated_process_registry return SolidObjects.caller_process.process_registry unless dedicated_registry dedicated_registry.tap(&:heartbeat) end |
#raise_rejection(message) ⇒ bot
92 93 94 95 96 97 98 99 100 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 92 def raise_rejection() rejection = .rejection raise Rejected.new( code: rejection.fetch("code"), message: rejection.fetch("message"), details: rejection.fetch("details"), message_id: .id ) end |
#wait(remaining) ⇒ void
This method returns an undefined value.
132 133 134 135 136 |
# File 'lib/solid_objects/synchronous_invocation.rb', line 132 def wait(remaining) SolidObjects.wake_up.wait( timeout: [ remaining, SolidObjects.configuration.sync_polling_interval ].min ) end |