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

Instance Method Details

#assist(message, deadline:) ⇒ Integer

RBS:

  • (Message, deadline: Float) -> Integer

Parameters:

Returns:

  • (Integer)


96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
# File 'lib/solid_objects/synchronous_invocation.rb', line 96

def assist(message, deadline:)
  process_registry = SolidObjects.caller_process.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)


10
11
12
13
14
15
16
17
18
# File 'lib/solid_objects/synchronous_invocation.rb', line 10

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)


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

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)


71
72
73
74
75
76
77
78
79
80
81
82
# File 'lib/solid_objects/synchronous_invocation.rb', line 71

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:



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

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:



124
125
126
127
128
129
130
# File 'lib/solid_objects/synchronous_invocation.rb', line 124

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:



51
52
53
54
55
56
57
# File 'lib/solid_objects/synchronous_invocation.rb', line 51

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:



60
61
62
63
64
65
66
67
68
# File 'lib/solid_objects/synchronous_invocation.rb', line 60

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)


140
141
142
# File 'lib/solid_objects/synchronous_invocation.rb', line 140

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

#raise_rejection(message) ⇒ bot

RBS:

  • (Message) -> bot

Parameters:

Returns:

  • (bot)


85
86
87
88
89
90
91
92
93
# File 'lib/solid_objects/synchronous_invocation.rb', line 85

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)


117
118
119
120
121
# File 'lib/solid_objects/synchronous_invocation.rb', line 117

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