Class: SolidObjects::Mailbox

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

Constant Summary collapse

INSTANCE_RETRY_LIMIT =

Returns:

  • (::Integer)
3

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(database_adapter: SolidObjects.database_adapter) ⇒ Mailbox

Returns a new instance of Mailbox.

RBS:

  • (?database_adapter: DatabaseAdapter) -> void

Parameters:

  • database_adapter: (DatabaseAdapter) (defaults to: SolidObjects.database_adapter)


10
11
12
# File 'lib/solid_objects/mailbox.rb', line 10

def initialize(database_adapter: SolidObjects.database_adapter)
  @database_adapter = database_adapter
end

Instance Attribute Details

#database_adapterObject (readonly)

Returns the value of attribute database_adapter.

Returns:

  • (Object)


87
88
89
# File 'lib/solid_objects/mailbox.rb', line 87

def database_adapter
  @database_adapter
end

Instance Method Details

#announce(message) ⇒ void

This method returns an undefined value.

RBS:

  • (Message) -> void

Parameters:



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

def announce(message)
  SolidObjects.instrument(
    :"message.enqueued",
    message_id: message.id,
    actor_type: message.actor_type,
    actor_id: message.actor_id,
    sequence: message.sequence,
    request_id: message.request_id
  )
  SolidObjects.wake_up.signal
end

#create_message(instance, reference, message_name, arguments, kind:, available_at:, idempotency_key:) ⇒ Message

RBS:

  • (Instance, Reference, Symbol | String, untyped, kind: String, available_at: Time?, idempotency_key: String?) -> Message

Parameters:

  • (Instance)
  • (Reference)
  • (Symbol, String)
  • (Object)
  • kind: (String)
  • available_at: (Time, nil)
  • idempotency_key: (String, nil)

Returns:



140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
# File 'lib/solid_objects/mailbox.rb', line 140

def create_message(instance, reference, message_name, arguments, kind:, available_at:, idempotency_key:)
  sequence = instance.next_message_sequence
  now = database_adapter.database_now
  scheduled_at = normalize_availability(available_at, now)
  instance.update!(next_message_sequence: sequence + 1)
  message = Message.create!(
    instance:,
    actor_type: reference.actor_type,
    actor_id: reference.actor_id,
    message_name: message_name.to_s,
    message_kind: kind,
    arguments:,
    sequence:,
    max_attempts: SolidObjects.configuration.max_attempts,
    request_id: SecureRandom.uuid,
    idempotency_key:,
    enqueued_at: now,
    available_at: scheduled_at
  )
  ReadyMessage.create!(
    message:,
    instance:,
    sequence:,
    available_at: scheduled_at
  )
  message
end

#enforce_mailbox_limit!(instance) ⇒ void

This method returns an undefined value.

RBS:

  • (Instance) -> void

Parameters:



131
132
133
134
135
136
137
# File 'lib/solid_objects/mailbox.rb', line 131

def enforce_mailbox_limit!(instance)
  live_count = ReadyMessage.where(instance_id: instance.id).count +
    ClaimedMessage.where(instance_id: instance.id).count
  return if live_count < SolidObjects.configuration.max_mailbox_length

  raise MailboxFull, "mailbox is full for #{instance.actor_type}(#{instance.actor_id.inspect})"
end

#enqueue(reference, message_name, arguments, kind:, available_at: nil, idempotency_key: nil) ⇒ MessageReference

RBS:

  • (Reference, Symbol | String, Hash[Symbol | String, untyped], kind: String, ?available_at: Time?, idempotency_key: String?) -> MessageReference

Parameters:

  • (Reference)
  • (Symbol, String)
  • (Hash[Symbol | String, untyped])
  • kind: (String)
  • idempotency_key: (String, nil) (defaults to: nil)
  • available_at: (Time, nil) (defaults to: nil)

Returns:



15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
# File 'lib/solid_objects/mailbox.rb', line 15

def enqueue(reference, message_name, arguments, kind:, available_at: nil, idempotency_key: nil)
  actor_class = SolidObjects.registry.fetch(reference.actor_type)
  normalized_arguments = Serialization.dump(
    arguments,
    max_bytes: SolidObjects.configuration.max_payload_bytes
  )
  message = with_instance_retry do
    database_adapter.transaction do
      enqueue_in_transaction(
        reference,
        message_name,
        normalized_arguments,
        kind:,
        available_at:,
        idempotency_key:,
        actor_class:
      )
    end
  end

  announce(message)
  MessageReference.from_message(message)
end

#enqueue_in_transaction(reference, message_name, arguments, kind:, available_at: nil, idempotency_key: nil, actor_class: nil) ⇒ Message

RBS:

  • (Reference, Symbol | String, Hash[Symbol | String, untyped], kind: String, ?available_at: Time?, idempotency_key: String?, ?actor_class: Class?) -> Message

Parameters:

  • (Reference)
  • (Symbol, String)
  • (Hash[Symbol | String, untyped])
  • kind: (String)
  • idempotency_key: (String, nil) (defaults to: nil)
  • available_at: (Time, nil) (defaults to: nil)
  • actor_class: (Class, nil) (defaults to: nil)

Returns:



40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
# File 'lib/solid_objects/mailbox.rb', line 40

def enqueue_in_transaction(
  reference,
  message_name,
  arguments,
  kind:,
  available_at: nil,
  idempotency_key: nil,
  actor_class: nil
)
  actor_class ||= SolidObjects.registry.fetch(reference.actor_type)
  normalized_arguments = Serialization.dump(
    arguments,
    max_bytes: SolidObjects.configuration.max_payload_bytes
  )
  instance = find_or_create_instance(reference, actor_class)
  instance.lock!

  existing = find_idempotent_message(instance, idempotency_key)
  return validate_idempotent_message!(existing, message_name, kind, normalized_arguments) if existing

  enforce_mailbox_limit!(instance)
  create_message(
    instance,
    reference,
    message_name,
    normalized_arguments,
    kind:,
    available_at:,
    idempotency_key:
  )
end

#find_idempotent_message(instance, idempotency_key) ⇒ Message?

RBS:

  • (Instance, String?) -> Message?

Parameters:

Returns:



114
115
116
117
118
# File 'lib/solid_objects/mailbox.rb', line 114

def find_idempotent_message(instance, idempotency_key)
  return unless idempotency_key

  Message.find_by(instance_id: instance.id, idempotency_key:)
end

#find_or_create_instance(reference, actor_class) ⇒ Instance

RBS:

  • (Reference, Class) -> Instance

Parameters:

Returns:



103
104
105
106
107
108
109
110
111
# File 'lib/solid_objects/mailbox.rb', line 103

def find_or_create_instance(reference, actor_class)
  Instance.create_or_find_by!(
    actor_type: reference.actor_type,
    actor_id: reference.actor_id
  ) do |instance|
    instance.state = {}
    instance.state_version = actor_class.state_version
  end
end

#normalize_availability(available_at, now) ⇒ Time

RBS:

  • (Time?, Time) -> Time

Parameters:

  • (Time, nil)
  • (Time)

Returns:

  • (Time)


169
170
171
172
173
174
175
176
# File 'lib/solid_objects/mailbox.rb', line 169

def normalize_availability(available_at, now)
  return now unless available_at

  time = available_at.respond_to?(:to_time) ? available_at.to_time : available_at
  raise ArgumentError, "available_at must be a time" unless time.is_a?(Time)

  [ time, now ].max
end

#validate_idempotent_message!(message, message_name, kind, arguments) ⇒ Message

RBS:

  • (Message, Symbol | String, String, untyped) -> Message

Parameters:

  • (Message)
  • (Symbol, String)
  • (String)
  • (Object)

Returns:



121
122
123
124
125
126
127
128
# File 'lib/solid_objects/mailbox.rb', line 121

def validate_idempotent_message!(message, message_name, kind, arguments)
  matches = message.message_name == message_name.to_s &&
    message.message_kind == kind &&
    message.arguments == arguments
  return message if matches

  raise IdempotencyConflict, "idempotency key belongs to a different actor invocation"
end

#with_instance_retry { ... } ⇒ Message

RBS:

  • () { () -> Message } -> Message

Yields:

Yield Returns:

Returns:



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

def with_instance_retry
  attempts = 0
  begin
    yield
  rescue ActiveRecord::RecordNotFound
    attempts += 1
    retry if attempts < INSTANCE_RETRY_LIMIT

    raise ActorDestroyed, "actor changed repeatedly while enqueueing"
  end
end