Class: SolidObjects::ProcessRegistry

Inherits:
Object
  • Object
show all
Defined in:
lib/solid_objects/process_registry.rb,
sig/generated/lib/solid_objects/process_registry.rbs

Instance Attribute Summary collapse

Class Method Summary collapse

Instance Method Summary collapse

Constructor Details

#initializeProcessRegistry

Returns a new instance of ProcessRegistry.

RBS:

  • () -> void



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_recordObject (readonly)

RBS:

  • @process_record: Process?

  • @last_heartbeat_at: Float?

Returns:

  • (Object)


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

RBS:

  • (?now: Time) -> Integer

Parameters:

  • now: (Time) (defaults to: SolidObjects.database_adapter.database_now)

Returns:

  • (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

RBS:

  • (Process, ?now: Time) -> bool

Parameters:

  • (Process)
  • now: (Time) (defaults to: SolidObjects.database_adapter.database_now)

Returns:

  • (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.

RBS:

  • () -> void



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_upvoid

This method returns an undefined value.

RBS:

  • () -> void



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_metadataHash[Symbol, String]

RBS:

  • () -> Hash[Symbol, String]

Returns:

  • (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

#heartbeatBoolean

RBS:

  • () -> bool

Returns:

  • (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

RBS:

  • () -> bool

Returns:

  • (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_nowFloat

RBS:

  • () -> Float

Returns:

  • (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

RBS:

  • (?kind: String, ?metadata: Hash[String | Symbol, untyped]) -> Process

Parameters:

  • kind: (String) (defaults to: "worker")
  • metadata: (Hash[String | Symbol, untyped]) (defaults to: {})

Returns:



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_drainingBoolean

RBS:

  • () -> bool

Returns:

  • (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

#stopBoolean

RBS:

  • () -> bool

Returns:

  • (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