Class: SolidObjects::SyncDiagnostics

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

Instance Method Summary collapse

Instance Method Details

#activation_details(instance) ⇒ Hash[String, untyped]

RBS:

  • (Instance) -> Hash[String, untyped]

Parameters:

Returns:

  • (Hash[String, untyped])


152
153
154
155
156
157
158
159
160
# File 'lib/solid_objects/sync_diagnostics.rb', line 152

def activation_details(instance)
  process_record = Process.find_by(id: instance.activation_owner_id)
  {
    "owner_id" => instance.activation_owner_id,
    "generation" => instance.activation_generation,
    "expires_at" => instance.activation_expires_at&.iso8601(6),
    "process" => process_details(process_record)
  }
end

#blocker_details(blocker) ⇒ Hash[String, untyped]?

RBS:

  • (Message?) -> Hash[String, untyped]?

Parameters:

Returns:

  • (Hash[String, untyped], nil)


176
177
178
179
180
181
182
183
184
185
# File 'lib/solid_objects/sync_diagnostics.rb', line 176

def blocker_details(blocker)
  return unless blocker

  {
    "message_id" => blocker.id,
    "sequence" => blocker.sequence,
    "message_name" => blocker.message_name,
    "status" => message_status(blocker)
  }
end

#build_error(message, timeout:, status:, waiting_on:, activation:, blocker:) ⇒ SyncTimeout

RBS:

  • (Message, timeout: Numeric, status: String, waiting_on: String, activation: Hash[String, untyped], blocker: Hash[String, untyped]?) -> SyncTimeout

Parameters:

  • (Message)
  • timeout: (Numeric)
  • status: (String)
  • waiting_on: (String)
  • activation: (Hash[String, untyped])
  • blocker: (Hash[String, untyped], nil)

Returns:



68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
# File 'lib/solid_objects/sync_diagnostics.rb', line 68

def build_error(message, timeout:, status:, waiting_on:, activation:, blocker:)
  error = SyncTimeout.new(
    timeout:,
    actor_type: message.actor_type,
    actor_id: message.actor_id,
    message_name: message.message_name,
    message_id: message.id,
    request_id: message.request_id,
    sequence: message.sequence,
    status:,
    waiting_on:,
    activation:,
    blocker:
  )
  SolidObjects.instrument(
    :"sync.timeout",
    message_id: message.id,
    request_id: message.request_id,
    actor_type: message.actor_type,
    actor_id: message.actor_id,
    sequence: message.sequence,
    status:,
    waiting_on:,
    activation_owner_id: activation["owner_id"],
    activation_generation: activation["generation"]
  )
  error
end

#call(message, timeout:) ⇒ SyncTimeout

RBS:

  • (Message, timeout: Numeric) -> SyncTimeout

Parameters:

Returns:



6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
# File 'lib/solid_objects/sync_diagnostics.rb', line 6

def call(message, timeout:)
  message = Message.uncached { Message.find(message.id) }
  instance = Instance.uncached { Instance.find(message.instance_id) }
  blocker = earlier_blocker(message)
  status = message_status(message)
  waiting_on = waiting_reason(message, instance, blocker)
  activation = activation_details(instance)
  build_error(
    message,
    timeout:,
    status:,
    waiting_on:,
    activation:,
    blocker: blocker_details(blocker)
  )
end

#database_contention(message, timeout:) ⇒ SyncTimeout

RBS:

  • (Message, timeout: Numeric) -> SyncTimeout

Parameters:

Returns:



24
25
26
27
28
29
30
31
32
33
# File 'lib/solid_objects/sync_diagnostics.rb', line 24

def database_contention(message, timeout:)
  build_error(
    message,
    timeout:,
    status: "unknown",
    waiting_on: "database_contention",
    activation: {},
    blocker: nil
  )
end

#database_contention_for(message_reference, timeout:) ⇒ SyncTimeout

RBS:

  • (MessageReference, timeout: Numeric) -> SyncTimeout

Parameters:

Returns:



36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
# File 'lib/solid_objects/sync_diagnostics.rb', line 36

def database_contention_for(message_reference, timeout:)
  error = SyncTimeout.new(
    timeout:,
    actor_type: message_reference.actor_type,
    actor_id: message_reference.actor_id,
    message_name: "unknown",
    message_id: message_reference.id,
    request_id: message_reference.request_id,
    sequence: message_reference.sequence,
    status: "unknown",
    waiting_on: "database_contention",
    activation: {},
    blocker: nil
  )
  SolidObjects.instrument(
    :"sync.timeout",
    message_id: message_reference.id,
    request_id: message_reference.request_id,
    actor_type: message_reference.actor_type,
    actor_id: message_reference.actor_id,
    sequence: message_reference.sequence,
    status: "unknown",
    waiting_on: "database_contention",
    activation_owner_id: nil,
    activation_generation: nil
  )
  error
end

#earlier_blocker(message) ⇒ Message?

RBS:

  • (Message) -> Message?

Parameters:

Returns:



136
137
138
139
140
141
142
143
144
145
146
147
148
149
# File 'lib/solid_objects/sync_diagnostics.rb', line 136

def earlier_blocker(message)
  ready = Message
    .joins(:ready_message)
    .where(instance_id: message.instance_id, sequence: ...message.sequence)
    .order(:sequence)
    .first
  claimed = Message
    .joins(:claimed_message)
    .where(instance_id: message.instance_id, sequence: ...message.sequence)
    .order(:sequence)
    .first

  [ ready, claimed ].compact.min_by(&:sequence)
end

#future?(ready_message) ⇒ Boolean

RBS:

  • (ReadyMessage?) -> bool

Parameters:

Returns:

  • (Boolean)


123
124
125
126
# File 'lib/solid_objects/sync_diagnostics.rb', line 123

def future?(ready_message)
  ready_message&.available_at.present? &&
    ready_message.available_at > SolidObjects.database_adapter.database_now
end

#live_activation?(instance) ⇒ Boolean

RBS:

  • (Instance) -> bool

Parameters:

Returns:

  • (Boolean)


129
130
131
132
133
# File 'lib/solid_objects/sync_diagnostics.rb', line 129

def live_activation?(instance)
  instance.activation_owner_id.present? &&
    instance.activation_expires_at.present? &&
    instance.activation_expires_at > SolidObjects.database_adapter.database_now
end

#message_status(message) ⇒ String

RBS:

  • (Message) -> String

Parameters:

Returns:

  • (String)


98
99
100
101
102
103
104
105
106
# File 'lib/solid_objects/sync_diagnostics.rb', line 98

def message_status(message)
  return "rejected" if message.rejected?
  return "completed" if message.completed?
  return "dead" if message.dead?
  return "claimed" if ClaimedMessage.where(message_id: message.id).exists?
  return "ready" if ReadyMessage.where(message_id: message.id).exists?

  "unknown"
end

#process_details(process_record) ⇒ Hash[String, untyped]?

RBS:

  • (Process?) -> Hash[String, untyped]?

Parameters:

Returns:

  • (Hash[String, untyped], nil)


163
164
165
166
167
168
169
170
171
172
173
# File 'lib/solid_objects/sync_diagnostics.rb', line 163

def process_details(process_record)
  return unless process_record

  {
    "kind" => process_record.kind,
    "hostname" => process_record.hostname,
    "pid" => process_record.pid,
    "last_heartbeat_at" => process_record.last_heartbeat_at&.iso8601(6),
    "shutdown_state" => process_record.shutdown_state
  }
end

#waiting_reason(message, instance, blocker) ⇒ String

RBS:

  • (Message, Instance, Message?) -> String

Parameters:

Returns:

  • (String)


109
110
111
112
113
114
115
116
117
118
119
120
# File 'lib/solid_objects/sync_diagnostics.rb', line 109

def waiting_reason(message, instance, blocker)
  return "actor_paused" if instance.paused_at
  return "activation_held" if live_activation?(instance)
  return "earlier_message" if blocker
  return "message_claimed" if ClaimedMessage.where(message_id: message.id).exists?

  ready_message = ReadyMessage.find_by(message_id: message.id)
  return "not_yet_available" if future?(ready_message)
  return "ready_unclaimed" if ready_message

  "unknown"
end