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])


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

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)


122
123
124
125
126
127
128
129
130
131
# File 'lib/solid_objects/sync_diagnostics.rb', line 122

def blocker_details(blocker)
  return unless blocker

  {
    "message_id" => blocker.id,
    "sequence" => blocker.sequence,
    "message_name" => blocker.message_name,
    "status" => message_status(blocker)
  }
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
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
# 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)
  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: blocker_details(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

#earlier_blocker(message) ⇒ Message?

RBS:

  • (Message) -> Message?

Parameters:

Returns:



82
83
84
85
86
87
88
89
90
91
92
93
94
95
# File 'lib/solid_objects/sync_diagnostics.rb', line 82

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)


69
70
71
72
# File 'lib/solid_objects/sync_diagnostics.rb', line 69

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)


75
76
77
78
79
# File 'lib/solid_objects/sync_diagnostics.rb', line 75

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)


44
45
46
47
48
49
50
51
52
# File 'lib/solid_objects/sync_diagnostics.rb', line 44

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)


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

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)


55
56
57
58
59
60
61
62
63
64
65
66
# File 'lib/solid_objects/sync_diagnostics.rb', line 55

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