Class: SolidObjects::ProcessRegistry
- Inherits:
-
Object
- Object
- SolidObjects::ProcessRegistry
- Defined in:
- lib/solid_objects/process_registry.rb,
sig/generated/lib/solid_objects/process_registry.rbs
Instance Attribute Summary collapse
- #process_record ⇒ Object readonly
Class Method Summary collapse
- .cleanup_dead(now: SolidObjects.database_adapter.database_now) ⇒ Integer
- .deregister(process_record, now: SolidObjects.database_adapter.database_now) ⇒ Boolean
- .reset_polling_warning! ⇒ void
- .warn_if_polling_is_only_cross_process_wake_up ⇒ void
Instance Method Summary collapse
- #default_metadata ⇒ Hash[Symbol, String]
- #heartbeat ⇒ Boolean
- #heartbeat_recent? ⇒ Boolean
-
#initialize ⇒ ProcessRegistry
constructor
A new instance of ProcessRegistry.
- #monotonic_now ⇒ Float
- #register(kind: "worker", metadata: {}) ⇒ Process
- #start_draining ⇒ Boolean
- #stop ⇒ Boolean
Constructor Details
#initialize ⇒ ProcessRegistry
Returns a new instance of ProcessRegistry.
127 128 129 130 |
# File 'lib/solid_objects/process_registry.rb', line 127 def initialize @process_record = nil @last_heartbeat_at = nil end |
Instance Attribute Details
#process_record ⇒ Object (readonly)
124 125 126 |
# File 'lib/solid_objects/process_registry.rb', line 124 def process_record @process_record end |
Class Method Details
.cleanup_dead(now: SolidObjects.database_adapter.database_now) ⇒ Integer
12 13 14 15 16 17 18 19 20 21 |
# File 'lib/solid_objects/process_registry.rb', line 12 def cleanup_dead(now: SolidObjects.database_adapter.database_now) stale_at = now - SolidObjects.configuration.process_alive_threshold dead_processes = Process .where.not(shutdown_state: "stopped") .where(last_heartbeat_at: ..stale_at) .to_a dead_processes.each { |process_record| cleanup_process(process_record, now) } dead_processes.length end |
.deregister(process_record, now: SolidObjects.database_adapter.database_now) ⇒ Boolean
24 25 26 27 28 29 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 56 57 |
# File 'lib/solid_objects/process_registry.rb', line 24 def deregister(process_record, now: SolidObjects.database_adapter.database_now) SolidObjects.database_adapter.transaction do Instance.where(activation_owner_id: process_record.id).update_all( activation_owner_id: nil, activation_token: nil, activation_expires_at: nil ) ClaimedMessage.where(process_id: process_record.id).update_all( process_id: nil, activation_token: nil ) Effect.where(claimed_by: process_record.id).update_all( status: "pending", claimed_by: nil, claimed_at: nil, available_at: now ) Reminder.where(claimed_by: process_record.id).update_all( claimed_by: nil, claimed_at: nil ) Broadcast.where(claimed_by: process_record.id).update_all( status: "pending", claimed_by: nil, claimed_at: nil, available_at: now ) process_record.update!( shutdown_state: "stopped", stopped_at: now ) end true end |
.reset_polling_warning! ⇒ void
This method returns an undefined value.
82 83 84 |
# File 'lib/solid_objects/process_registry.rb', line 82 def reset_polling_warning! polling_warning_mutex.synchronize { @polling_warning_emitted = false } end |
.warn_if_polling_is_only_cross_process_wake_up ⇒ void
This method returns an undefined value.
60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 |
# File 'lib/solid_objects/process_registry.rb', line 60 def warn_if_polling_is_only_cross_process_wake_up return if SolidObjects.configuration.wake_up_adapter polling_warning_mutex.synchronize do return if polling_warning_emitted? return unless another_live_process? payload = { event: "solid_objects.polling_only_cross_process_wake_up", polling_interval: SolidObjects.configuration.polling_interval, idle_polling_interval: SolidObjects.configuration.idle_polling_interval } SolidObjects.configuration.logger.warn(payload) SolidObjects.instrument( :"polling.only_cross_process_wake_up", **payload.except(:event) ) @polling_warning_emitted = true end end |
Instance Method Details
#default_metadata ⇒ Hash[Symbol, String]
185 186 187 188 189 190 |
# File 'lib/solid_objects/process_registry.rb', line 185 def { solid_objects_version: SolidObjects::VERSION, ruby_version: RUBY_VERSION } end |
#heartbeat ⇒ Boolean
152 153 154 155 156 157 158 159 160 161 162 163 |
# File 'lib/solid_objects/process_registry.rb', line 152 def heartbeat return false unless process_record return false if heartbeat_recent? updated = SolidObjects.database_adapter.with_lock_retry do process_record.update( last_heartbeat_at: SolidObjects.database_adapter.database_now ) end @last_heartbeat_at = monotonic_now if updated updated end |
#heartbeat_recent? ⇒ Boolean
193 194 195 196 197 198 |
# File 'lib/solid_objects/process_registry.rb', line 193 def heartbeat_recent? return false unless @last_heartbeat_at monotonic_now - @last_heartbeat_at < SolidObjects.configuration.process_heartbeat_interval end |
#monotonic_now ⇒ Float
201 202 203 |
# File 'lib/solid_objects/process_registry.rb', line 201 def monotonic_now ::Process.clock_gettime(::Process::CLOCK_MONOTONIC) end |
#register(kind: "worker", metadata: {}) ⇒ Process
133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 |
# File 'lib/solid_objects/process_registry.rb', line 133 def register(kind: "worker", metadata: {}) process_record = SolidObjects.database_adapter.with_lock_retry do now = SolidObjects.database_adapter.database_now Process.create!( id: SecureRandom.uuid, kind:, hostname: Socket.gethostname, pid: ::Process.pid, started_at: now, last_heartbeat_at: now, metadata: Serialization.dump(.merge()) ) end @process_record = process_record @last_heartbeat_at = monotonic_now process_record end |
#start_draining ⇒ Boolean
166 167 168 169 170 171 172 173 |
# File 'lib/solid_objects/process_registry.rb', line 166 def start_draining return false unless process_record process_record.update( shutdown_state: "draining", shutdown_requested_at: SolidObjects.database_adapter.database_now ) end |
#stop ⇒ Boolean
176 177 178 179 180 |
# File 'lib/solid_objects/process_registry.rb', line 176 def stop return false unless process_record self.class.deregister(process_record) end |