Class: SolidObjects::SynchronousInvocation

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

Instance Method Summary collapse

Constructor Details

#initialize(process_registry: nil) ⇒ SynchronousInvocation

Returns a new instance of SynchronousInvocation.

RBS:

  • (?process_registry: ProcessRegistry?) -> void

Parameters:



12
13
14
# File 'lib/solid_objects/synchronous_invocation.rb', line 12

def initialize(process_registry: nil)
  @dedicated_process_registry = process_registry
end

Instance Method Details

#assist(message, deadline:) ⇒ Integer

RBS:

  • (Message, deadline: Float) -> Integer

Parameters:

Returns:

  • (Integer)


111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
# File 'lib/solid_objects/synchronous_invocation.rb', line 111

def assist(message, deadline:)
  process_registry = self.process_registry
  activation = ActivationManager
    .new(owner_id: process_registry.process_record.id)
    .claim(instance_id: message.instance_id)
  return 0 unless activation

  LeaseRenewer.new(
    activation:,
    process_registry:
  ).around do
    activation.drain_until(message_id: message.id, deadline:)
  end
ensure
  if activation
    activation.yield_ready_messages if activation.pass_exhausted?
    activation.deactivate
  end
end

#call(message_reference, timeout:) ⇒ Object

RBS:

  • (MessageReference, timeout: Numeric) -> untyped

Parameters:

Returns:

  • (Object)


17
18
19
20
21
22
23
24
25
# File 'lib/solid_objects/synchronous_invocation.rb', line 17

def call(message_reference, timeout:)
  return call_before_deadline(message_reference, timeout:) if SyncDeadline.active?

  SyncDeadline.with(timeout:) do
    call_before_deadline(message_reference, timeout:)
  end
rescue ActiveRecord::RecordNotFound
  raise ActorDestroyed, "actor was destroyed while waiting for its result"
end

#call_before_deadline(message_reference, timeout:) ⇒ Object

RBS:

  • (MessageReference, timeout: Numeric) -> untyped

Parameters:

Returns:

  • (Object)


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
# File 'lib/solid_objects/synchronous_invocation.rb', line 30

def call_before_deadline(message_reference, timeout:)
  deadline = monotonic_now + SyncDeadline.remaining
  message = nil

  loop do
    unless (deadline - monotonic_now).positive?
      final_message = load_message_at_deadline(message_reference)
      return completed_result(final_message) if final_message&.completed? || final_message&.dead?

      raise final_message ?
        diagnose_timeout(final_message, timeout:) :
        contention_timeout(message, message_reference, timeout:)
    end

    message = load_message(message_reference)
    return completed_result(message) if message.completed? || message.dead?

    processed = assist(message, deadline:)
    remaining = deadline - monotonic_now
    wait(remaining) if processed.zero? && remaining.positive?
  end
rescue DatabaseDeadlineExceeded
  raise message ?
    diagnose_timeout(message, timeout:) :
    contention_timeout(message, message_reference, timeout:)
end

#completed_result(message) ⇒ Object

RBS:

  • (Message) -> untyped

Parameters:

Returns:

  • (Object)


78
79
80
81
82
83
84
85
86
87
88
89
# File 'lib/solid_objects/synchronous_invocation.rb', line 78

def completed_result(message)
  raise_rejection(message) if message.rejected?
  if message.dead?
    raise MessageFailed.new(
      "actor message failed permanently",
      message_id: message.id,
      details: message.error || {}
    )
  end

  Serialization.readonly_copy(message.result)
end

#contention_timeout(message, message_reference, timeout:) ⇒ SyncTimeout

RBS:

  • (Message?, MessageReference, timeout: Numeric) -> SyncTimeout

Parameters:

Returns:



148
149
150
151
152
# File 'lib/solid_objects/synchronous_invocation.rb', line 148

def contention_timeout(message, message_reference, timeout:)
  return SyncDiagnostics.new.database_contention(message, timeout:) if message

  SyncDiagnostics.new.database_contention_for(message_reference, timeout:)
end

#diagnose_timeout(message, timeout:) ⇒ SyncTimeout

RBS:

  • (Message, timeout: Numeric) -> SyncTimeout

Parameters:

Returns:



139
140
141
142
143
144
145
# File 'lib/solid_objects/synchronous_invocation.rb', line 139

def diagnose_timeout(message, timeout:)
  SolidObjects.database_adapter.with_lock_probe do
    SyncDiagnostics.new.call(message, timeout:)
  end
rescue DatabaseDeadlineExceeded
  SyncDiagnostics.new.database_contention(message, timeout:)
end

#load_message(message_reference) ⇒ Message

RBS:

  • (MessageReference) -> Message

Parameters:

Returns:



58
59
60
61
62
63
64
# File 'lib/solid_objects/synchronous_invocation.rb', line 58

def load_message(message_reference)
  SolidObjects.database_adapter.with_lock_retry do
    Message.uncached do
      Message.includes(:dead_letter).find(message_reference.id)
    end
  end
end

#load_message_at_deadline(message_reference) ⇒ Message?

RBS:

  • (MessageReference) -> Message?

Parameters:

Returns:



67
68
69
70
71
72
73
74
75
# File 'lib/solid_objects/synchronous_invocation.rb', line 67

def load_message_at_deadline(message_reference)
  SolidObjects.database_adapter.with_lock_probe do
    Message.uncached do
      Message.includes(:dead_letter).find(message_reference.id)
    end
  end
rescue DatabaseDeadlineExceeded
  nil
end

#monotonic_nowFloat

RBS:

  • () -> Float

Returns:

  • (Float)


155
156
157
# File 'lib/solid_objects/synchronous_invocation.rb', line 155

def monotonic_now
  ::Process.clock_gettime(::Process::CLOCK_MONOTONIC)
end

#process_registryProcessRegistry

RBS:

  • () -> ProcessRegistry

Returns:



103
104
105
106
107
108
# File 'lib/solid_objects/synchronous_invocation.rb', line 103

def process_registry
  dedicated_registry = @dedicated_process_registry
  return SolidObjects.caller_process.process_registry unless dedicated_registry

  dedicated_registry.tap(&:heartbeat)
end

#raise_rejection(message) ⇒ bot

RBS:

  • (Message) -> bot

Parameters:

Returns:

  • (bot)


92
93
94
95
96
97
98
99
100
# File 'lib/solid_objects/synchronous_invocation.rb', line 92

def raise_rejection(message)
  rejection = message.rejection
  raise Rejected.new(
    code: rejection.fetch("code"),
    message: rejection.fetch("message"),
    details: rejection.fetch("details"),
    message_id: message.id
  )
end

#wait(remaining) ⇒ void

This method returns an undefined value.

RBS:

  • (Numeric) -> void

Parameters:

  • (Numeric)


132
133
134
135
136
# File 'lib/solid_objects/synchronous_invocation.rb', line 132

def wait(remaining)
  SolidObjects.wake_up.wait(
    timeout: [ remaining, SolidObjects.configuration.sync_polling_interval ].min
  )
end