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)


11
12
13
14
15
16
17
18
19
20
# File 'lib/solid_objects/reminder_scheduler.rb', line 11

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
end

Instance Attribute Details

#database_adapterObject (readonly)

Returns the value of attribute database_adapter.

Returns:

  • (Object)


76
77
78
# File 'lib/solid_objects/reminder_scheduler.rb', line 76

def database_adapter
  @database_adapter
end

#process_registryObject (readonly)

Returns the value of attribute process_registry.

Returns:

  • (Object)


76
77
78
# File 'lib/solid_objects/reminder_scheduler.rb', line 76

def process_registry
  @process_registry
end

Instance Method Details

#claim_next(now:) ⇒ Reminder?

RBS:

  • (now: Time?) -> Reminder?

Parameters:

  • now: (Time, nil)

Returns:



79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
# File 'lib/solid_objects/reminder_scheduler.rb', line 79

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

#enqueue(reminder, now:) ⇒ MessageReference?

RBS:

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

Parameters:

Returns:



100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
# File 'lib/solid_objects/reminder_scheduler.rb', line 100

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)


161
162
163
164
165
166
167
168
# File 'lib/solid_objects/reminder_scheduler.rb', line 161

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)


171
172
173
174
175
176
177
178
179
180
# File 'lib/solid_objects/reminder_scheduler.rb', line 171

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:



147
148
149
150
151
# File 'lib/solid_objects/reminder_scheduler.rb', line 147

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



59
60
61
62
# File 'lib/solid_objects/reminder_scheduler.rb', line 59

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

#runvoid

This method returns an undefined value.

RBS:

  • () -> void



47
48
49
50
51
52
53
54
55
56
# File 'lib/solid_objects/reminder_scheduler.rb', line 47

def run
  until shutdown_requested?
    worked = run_once
    next if worked

    SolidObjects.wake_up.wait(timeout: SolidObjects.configuration.polling_interval)
  end
ensure
  stop
end

#run_once(now: nil) ⇒ Boolean

RBS:

  • (?now: Time?) -> bool

Parameters:

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

Returns:

  • (Boolean)


23
24
25
26
27
28
29
30
31
32
33
34
35
# File 'lib/solid_objects/reminder_scheduler.rb', line 23

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)


70
71
72
# File 'lib/solid_objects/reminder_scheduler.rb', line 70

def shutdown_requested?
  @shutdown_requested
end

#stopvoid

This method returns an undefined value.

RBS:

  • () -> void



38
39
40
41
42
43
44
# File 'lib/solid_objects/reminder_scheduler.rb', line 38

def stop
  return if stopped?

  @stopped = true
  process_registry.start_draining
  process_registry.stop
end

#stopped?Boolean

RBS:

  • () -> bool

Returns:

  • (Boolean)


65
66
67
# File 'lib/solid_objects/reminder_scheduler.rb', line 65

def stopped?
  @stopped
end

#verify_claim!(reminder) ⇒ void

This method returns an undefined value.

RBS:

  • (Reminder) -> void

Parameters:



154
155
156
157
158
# File 'lib/solid_objects/reminder_scheduler.rb', line 154

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

  raise LostActivation, "reminder claim changed"
end