Class: SolidObjects::ReminderScheduler

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

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(process_registry: ProcessRegistry.new, database_adapter: SolidObjects.database_adapter) ⇒ ReminderScheduler

Returns a new instance of ReminderScheduler.

RBS:

  • (?process_registry: ProcessRegistry, ?database_adapter: DatabaseAdapter) -> void

Parameters:

  • process_registry: (ProcessRegistry) (defaults to: ProcessRegistry.new)
  • database_adapter: (DatabaseAdapter) (defaults to: SolidObjects.database_adapter)


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

def initialize(
  process_registry: ProcessRegistry.new,
  database_adapter: SolidObjects.database_adapter
)
  @process_registry = process_registry
  @database_adapter = database_adapter
  process_registry.register(kind: "reminder")
  @stopped = false
  @shutdown_requested = false
  @polling_backoff = PollingBackoff.new(
    minimum_interval: SolidObjects.configuration.polling_interval,
    maximum_interval: SolidObjects.configuration.idle_polling_interval,
    on_change: ->(transition) do
      SolidObjects.instrument(
        :"polling.interval_changed",
        role: "reminders",
        **transition
      )
    end
  )
end

Instance Attribute Details

#database_adapterObject (readonly)

Returns the value of attribute database_adapter.

Returns:

  • (Object)


107
108
109
# File 'lib/solid_objects/reminder_scheduler.rb', line 107

def database_adapter
  @database_adapter
end

#polling_backoffObject (readonly)

Returns the value of attribute polling_backoff.

Returns:

  • (Object)


107
108
109
# File 'lib/solid_objects/reminder_scheduler.rb', line 107

def polling_backoff
  @polling_backoff
end

#process_registryObject (readonly)

Returns the value of attribute process_registry.

Returns:

  • (Object)


107
108
109
# File 'lib/solid_objects/reminder_scheduler.rb', line 107

def process_registry
  @process_registry
end

Instance Method Details

#claim_next(now:) ⇒ Reminder?

RBS:

  • (now: Time?) -> Reminder?

Parameters:

  • now: (Time, nil)

Returns:



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

def claim_next(now:)
  database_adapter.transaction do
    database_now = database_adapter.database_now
    due_at = now || database_now
    stale_at = database_now - SolidObjects.configuration.process_alive_threshold
    relation = Reminder
      .where(status: "scheduled", next_run_at: ..due_at)
      .where("claimed_by IS NULL OR claimed_at <= ?", stale_at)
      .order(:next_run_at, :id)
    reminder = database_adapter.lock_candidates(relation).first
    next unless reminder

    reminder.update!(
      claimed_by: process_registry.process_record.id,
      claimed_at: database_now
    )
    reminder
  end
end

#current_polling_intervalFloat

RBS:

  • () -> Float

Returns:

  • (Float)


101
102
103
# File 'lib/solid_objects/reminder_scheduler.rb', line 101

def current_polling_interval
  polling_backoff.current_interval
end

#enqueue(reminder, now:) ⇒ MessageReference?

RBS:

  • (Reminder, now: Time?) -> MessageReference?

Parameters:

Returns:



131
132
133
134
135
136
137
138
139
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
167
168
169
170
171
172
173
174
175
# File 'lib/solid_objects/reminder_scheduler.rb', line 131

def enqueue(reminder, now:)
  actor_class = SolidObjects.registry.fetch(reminder.actor_type)
  unless actor_class.definition.messages.key?(reminder.operation.to_sym)
    raise UnknownMessage, "unknown reminder operation #{reminder.operation.inspect}"
  end

  mailbox = Mailbox.new(database_adapter:)
  message = database_adapter.transaction do
    instance = Instance.lock.find_by(id: reminder.instance_id)
    next unless instance

    locked_reminder = Reminder.lock.find_by(id: reminder.id)
    next unless locked_reminder
    verify_claim!(locked_reminder)
    schedule_now = now || database_adapter.database_now
    message = mailbox.enqueue_in_transaction(
      reference: Reference.new(actor_type: instance.actor_type, actor_id: instance.actor_id),
      operation: locked_reminder.operation,
      arguments: locked_reminder.arguments,
      delivery_mode: "internal",
      idempotency_key: "reminder:#{locked_reminder.id}:#{locked_reminder.occurrence}",
      actor_class:
    )
    recurring = locked_reminder.interval_seconds.present?
    locked_reminder.update!(
      status: recurring ? "scheduled" : "completed",
      occurrence: locked_reminder.occurrence + 1,
      next_run_at: recurring ? next_run_at(locked_reminder, schedule_now) : locked_reminder.next_run_at,
      claimed_by: nil,
      claimed_at: nil
    )
    message
  end
  return unless message

  mailbox.announce(message)
  SolidObjects.instrument(
    :"reminder.enqueued",
    reminder_id: reminder.id,
    actor_type: reminder.actor_type,
    actor_id: reminder.actor_id,
    occurrence: reminder.occurrence
  )
  MessageReference.from_message(message)
end

#next_run_at(reminder, now) ⇒ Time

RBS:

  • (Reminder, Time) -> Time

Parameters:

Returns:

  • (Time)


192
193
194
195
196
197
198
199
# File 'lib/solid_objects/reminder_scheduler.rb', line 192

def next_run_at(reminder, now)
  interval = reminder.interval_seconds.to_f
  next_run = reminder.next_run_at + interval
  return next_run if reminder.missed_policy == "all" || next_run > now

  missed_intervals = ((now - next_run) / interval).floor + 1
  next_run + (missed_intervals * interval)
end

#normalize_test_time(now) ⇒ Time?

RBS:

  • (Time?) -> Time?

Parameters:

  • (Time, nil)

Returns:

  • (Time, nil)


202
203
204
205
206
207
208
209
210
211
# File 'lib/solid_objects/reminder_scheduler.rb', line 202

def normalize_test_time(now)
  return unless now

  time = now.to_time
  return time if time.to_f.finite?

  raise ArgumentError
rescue ArgumentError, NoMethodError, RangeError
  raise ArgumentError, "reminder test time must be a valid time"
end

#release(reminder) ⇒ void

This method returns an undefined value.

RBS:

  • (Reminder) -> void

Parameters:



178
179
180
181
182
# File 'lib/solid_objects/reminder_scheduler.rb', line 178

def release(reminder)
  Reminder
    .where(id: reminder.id, claimed_by: process_registry.process_record.id)
    .update_all(claimed_by: nil, claimed_at: nil)
end

#request_shutdownvoid

This method returns an undefined value.

RBS:

  • () -> void



85
86
87
88
# File 'lib/solid_objects/reminder_scheduler.rb', line 85

def request_shutdown
  @shutdown_requested = true
  SolidObjects.wake_up.signal
end

#runvoid

This method returns an undefined value.

RBS:

  • () -> void



61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
# File 'lib/solid_objects/reminder_scheduler.rb', line 61

def run
  ProcessRegistry.warn_if_polling_is_only_cross_process_wake_up

  until shutdown_requested?
    wake_up = SolidObjects.wake_up
    watch = wake_up.respond_to?(:watch) ? wake_up.watch : wake_up
    worked = run_once
    if worked
      polling_backoff.reset(:work)
      next
    end

    notified = watch.wait(timeout: current_polling_interval)
    if notified == false
      polling_backoff.record_idle
    else
      polling_backoff.reset(:wake_up)
    end
  end
ensure
  stop
end

#run_once(now: nil) ⇒ Boolean

RBS:

  • (?now: Time?) -> bool

Parameters:

  • now: (Time, nil) (defaults to: nil)

Returns:

  • (Boolean)


37
38
39
40
41
42
43
44
45
46
47
48
49
# File 'lib/solid_objects/reminder_scheduler.rb', line 37

def run_once(now: nil)
  return false if stopped?

  now = normalize_test_time(now)
  process_registry.heartbeat
  reminder = claim_next(now:)
  return false unless reminder

  enqueue(reminder, now:).present?
rescue
  release(reminder) if reminder
  raise
end

#shutdown_requested?Boolean

RBS:

  • () -> bool

Returns:

  • (Boolean)


96
97
98
# File 'lib/solid_objects/reminder_scheduler.rb', line 96

def shutdown_requested?
  @shutdown_requested
end

#stopvoid

This method returns an undefined value.

RBS:

  • () -> void



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

def stop
  return if stopped?

  @stopped = true
  process_registry.start_draining
  process_registry.stop
end

#stopped?Boolean

RBS:

  • () -> bool

Returns:

  • (Boolean)


91
92
93
# File 'lib/solid_objects/reminder_scheduler.rb', line 91

def stopped?
  @stopped
end

#verify_claim!(reminder) ⇒ void

This method returns an undefined value.

RBS:

  • (Reminder) -> void

Parameters:



185
186
187
188
189
# File 'lib/solid_objects/reminder_scheduler.rb', line 185

def verify_claim!(reminder)
  return if reminder.claimed_by == process_registry.process_record.id

  raise LostActivation, "reminder claim changed"
end