Class: SolidObjects::BroadcastExecutor
- Inherits:
-
Object
- Object
- SolidObjects::BroadcastExecutor
- Defined in:
- lib/solid_objects/broadcast_executor.rb,
sig/generated/lib/solid_objects/broadcast_executor.rbs
Instance Attribute Summary collapse
-
#database_adapter ⇒ Object
readonly
Returns the value of attribute database_adapter.
-
#process_registry ⇒ Object
readonly
Returns the value of attribute process_registry.
Instance Method Summary collapse
- #broadcast_adapter ⇒ Proc
- #claim_next ⇒ Broadcast?
- #complete(broadcast) ⇒ void
- #fail_broadcast(broadcast, error) ⇒ void
-
#initialize(process_registry: ProcessRegistry.new, database_adapter: SolidObjects.database_adapter) ⇒ BroadcastExecutor
constructor
A new instance of BroadcastExecutor.
- #request_shutdown ⇒ void
- #run ⇒ void
- #run_once ⇒ Boolean
- #shutdown_requested? ⇒ Boolean
- #stop ⇒ void
- #stopped? ⇒ Boolean
- #verify_claim!(broadcast) ⇒ void
Constructor Details
#initialize(process_registry: ProcessRegistry.new, database_adapter: SolidObjects.database_adapter) ⇒ BroadcastExecutor
Returns a new instance of BroadcastExecutor.
11 12 13 14 15 16 17 18 19 20 |
# File 'lib/solid_objects/broadcast_executor.rb', line 11 def initialize( process_registry: ProcessRegistry.new, database_adapter: SolidObjects.database_adapter ) @process_registry = process_registry @database_adapter = database_adapter process_registry.register(kind: "broadcast") @stopped = false @shutdown_requested = false end |
Instance Attribute Details
#database_adapter ⇒ Object (readonly)
Returns the value of attribute database_adapter.
77 78 79 |
# File 'lib/solid_objects/broadcast_executor.rb', line 77 def database_adapter @database_adapter end |
#process_registry ⇒ Object (readonly)
Returns the value of attribute process_registry.
77 78 79 |
# File 'lib/solid_objects/broadcast_executor.rb', line 77 def process_registry @process_registry end |
Instance Method Details
#broadcast_adapter ⇒ Proc
102 103 104 105 |
# File 'lib/solid_objects/broadcast_executor.rb', line 102 def broadcast_adapter SolidObjects.configuration.broadcast_adapter || ActionCableBroadcastAdapter.new end |
#claim_next ⇒ Broadcast?
80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 |
# File 'lib/solid_objects/broadcast_executor.rb', line 80 def claim_next database_adapter.transaction do now = database_adapter.database_now stale_at = now - SolidObjects.configuration.process_alive_threshold relation = Broadcast .where(status: "pending", available_at: ..now) .or(Broadcast.where(status: "processing", claimed_at: ..stale_at)) .order(:available_at, :id) broadcast = database_adapter.lock_candidates(relation).first next unless broadcast broadcast.update!( status: "processing", attempt_count: broadcast.attempt_count + 1, claimed_by: process_registry.process_record.id, claimed_at: now ) broadcast end end |
#complete(broadcast) ⇒ void
This method returns an undefined value.
108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 |
# File 'lib/solid_objects/broadcast_executor.rb', line 108 def complete(broadcast) database_adapter.transaction do locked_broadcast = Broadcast.lock.find(broadcast.id) verify_claim!(locked_broadcast) locked_broadcast.update!( status: "delivered", error: nil, delivered_at: database_adapter.database_now, claimed_by: nil, claimed_at: nil ) end SolidObjects.instrument( :"broadcast.delivered", broadcast_id: broadcast.broadcast_id, message_id: broadcast., actor_type: broadcast.instance.actor_type, actor_id: broadcast.instance.actor_id, observable_name: broadcast.observable_name, attempt: broadcast.attempt_count ) end |
#fail_broadcast(broadcast, error) ⇒ void
This method returns an undefined value.
132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 |
# File 'lib/solid_objects/broadcast_executor.rb', line 132 def fail_broadcast(broadcast, error) database_adapter.transaction do locked_broadcast = Broadcast.lock.find(broadcast.id) verify_claim!(locked_broadcast) dead = locked_broadcast.attempt_count >= SolidObjects.configuration.max_attempts locked_broadcast.update!( status: dead ? "dead" : "pending", available_at: database_adapter.database_now + SolidObjects.configuration.retry_delay.call(locked_broadcast.attempt_count), error: { "class" => error.class.name, "message" => error..to_s.byteslice(0, 8_192), "backtrace" => Array(error.backtrace).first(50) }, claimed_by: nil, claimed_at: nil ) end rescue ActiveRecord::RecordNotFound, LostActivation nil end |
#request_shutdown ⇒ void
This method returns an undefined value.
60 61 62 63 |
# File 'lib/solid_objects/broadcast_executor.rb', line 60 def request_shutdown @shutdown_requested = true SolidObjects.wake_up.signal end |
#run ⇒ void
This method returns an undefined value.
48 49 50 51 52 53 54 55 56 57 |
# File 'lib/solid_objects/broadcast_executor.rb', line 48 def run until shutdown_requested? worked = run_once next if worked SolidObjects.wake_up.wait(timeout: SolidObjects.configuration.polling_interval) end ensure stop end |
#run_once ⇒ Boolean
23 24 25 26 27 28 29 30 31 32 33 34 35 36 |
# File 'lib/solid_objects/broadcast_executor.rb', line 23 def run_once return false if stopped? process_registry.heartbeat broadcast = claim_next return false unless broadcast broadcast_adapter.call(broadcast) complete(broadcast) true rescue => error fail_broadcast(broadcast, error) if broadcast false end |
#shutdown_requested? ⇒ Boolean
71 72 73 |
# File 'lib/solid_objects/broadcast_executor.rb', line 71 def shutdown_requested? @shutdown_requested end |
#stop ⇒ void
This method returns an undefined value.
39 40 41 42 43 44 45 |
# File 'lib/solid_objects/broadcast_executor.rb', line 39 def stop return if stopped? @stopped = true process_registry.start_draining process_registry.stop end |
#stopped? ⇒ Boolean
66 67 68 |
# File 'lib/solid_objects/broadcast_executor.rb', line 66 def stopped? @stopped end |
#verify_claim!(broadcast) ⇒ void
This method returns an undefined value.
155 156 157 158 159 160 |
# File 'lib/solid_objects/broadcast_executor.rb', line 155 def verify_claim!(broadcast) return if broadcast.status == "processing" && broadcast.claimed_by == process_registry.process_record.id raise LostActivation, "broadcast claim changed" end |