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)


94
95
96
# File 'lib/solid_objects/mailbox.rb', line 94

def database_adapter
  @database_adapter
end

Instance Method Details

#announce(message) ⇒ void

This method returns an undefined value.

RBS:

  • (Message) -> void

Parameters:



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

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:, operation:, arguments:, delivery_mode:, available_at:, idempotency_key:) ⇒ Message

RBS:

  • (instance: Instance, reference: Reference, operation: Symbol | String, arguments: untyped, delivery_mode: String, available_at: Time?, idempotency_key: String?) -> Message

Parameters:

  • instance: (Instance)
  • reference: (Reference)
  • operation: (Symbol, String)
  • arguments: (Object)
  • delivery_mode: (String)
  • available_at: (Time, nil)
  • idempotency_key: (String, nil)

Returns:



147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
# File 'lib/solid_objects/mailbox.rb', line 147

def create_message(instance:, reference:, operation:, arguments:, delivery_mode:, 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,
    operation: operation.to_s,
    delivery_mode:,
    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:



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

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:, operation:, arguments:, delivery_mode:, available_at: nil, idempotency_key: nil) ⇒ MessageReference

RBS:

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

Parameters:

  • reference: (Reference)
  • operation: (Symbol, String)
  • arguments: (Hash[Symbol | String, untyped])
  • delivery_mode: (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:, operation:, arguments:, delivery_mode:, 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:,
        operation:,
        arguments: normalized_arguments,
        delivery_mode:,
        available_at:,
        idempotency_key:,
        actor_class:
      )
    end
  end

  announce(message)
  MessageReference.from_message(message)
end

#enqueue_in_transaction(reference:, operation:, arguments:, delivery_mode:, available_at: nil, idempotency_key: nil, actor_class: nil) ⇒ Message

RBS:

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

Parameters:

  • reference: (Reference)
  • operation: (Symbol, String)
  • arguments: (Hash[Symbol | String, untyped])
  • delivery_mode: (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
71
72
73
74
75
76
77
# File 'lib/solid_objects/mailbox.rb', line 40

def enqueue_in_transaction(
  reference:,
  operation:,
  arguments:,
  delivery_mode:,
  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)
  if existing
    return validate_idempotent_message!(
      existing_message: existing,
      operation:,
      delivery_mode:,
      arguments: normalized_arguments
    )
  end

  enforce_mailbox_limit!(instance)
  create_message(
    instance:,
    reference:,
    operation:,
    arguments: normalized_arguments,
    delivery_mode:,
    available_at:,
    idempotency_key:
  )
end

#find_idempotent_message(instance, idempotency_key) ⇒ Message?

RBS:

  • (Instance, String?) -> Message?

Parameters:

Returns:



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

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:



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

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)


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

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!(existing_message:, operation:, delivery_mode:, arguments:) ⇒ Message

RBS:

  • (existing_message: Message, operation: Symbol | String, delivery_mode: String, arguments: untyped) -> Message

Parameters:

  • existing_message: (Message)
  • operation: (Symbol, String)
  • delivery_mode: (String)
  • arguments: (Object)

Returns:



128
129
130
131
132
133
134
135
# File 'lib/solid_objects/mailbox.rb', line 128

def validate_idempotent_message!(existing_message:, operation:, delivery_mode:, arguments:)
  matches = existing_message.operation == operation.to_s &&
    existing_message.delivery_mode == delivery_mode &&
    existing_message.arguments == arguments
  return existing_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:



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

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