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)


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

def database_adapter
  @database_adapter
end

#process_registryObject (readonly)

Returns the value of attribute process_registry.

Returns:

  • (Object)


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

def process_registry
  @process_registry
end

Instance Method Details

#claim_nextReminder?

RBS:

  • () -> Reminder?

Returns:



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

def claim_next
  database_adapter.transaction do
    now = database_adapter.database_now
    stale_at = now - SolidObjects.configuration.process_alive_threshold
    relation = Reminder
      .where(status: "scheduled", next_run_at: ..now)
      .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: now
    )
    reminder
  end
end

#enqueue(reminder) ⇒ MessageReference?

RBS:

  • (Reminder) -> MessageReference?

Parameters:

Returns:



98
99
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
# File 'lib/solid_objects/reminder_scheduler.rb', line 98

def enqueue(reminder)
  actor_class = SolidObjects.registry.fetch(reminder.actor_type)
  unless actor_class.definition.messages.key?(reminder.message_name.to_sym)
    raise UnknownMessage, "unknown reminder message #{reminder.message_name.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)
    now = database_adapter.database_now
    message = mailbox.enqueue_in_transaction(
      Reference.new(actor_type: instance.actor_type, actor_id: instance.actor_id),
      locked_reminder.message_name,
      locked_reminder.arguments,
      kind: "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, 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)


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

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

#release(reminder) ⇒ void

This method returns an undefined value.

RBS:

  • (Reminder) -> void

Parameters:



145
146
147
148
149
# File 'lib/solid_objects/reminder_scheduler.rb', line 145

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



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

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

#runvoid

This method returns an undefined value.

RBS:

  • () -> void



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

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_onceBoolean

RBS:

  • () -> bool

Returns:

  • (Boolean)


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

def run_once
  return false if stopped?

  process_registry.heartbeat
  reminder = claim_next
  return false unless reminder

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

#shutdown_requested?Boolean

RBS:

  • () -> bool

Returns:

  • (Boolean)


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

def shutdown_requested?
  @shutdown_requested
end

#stopvoid

This method returns an undefined value.

RBS:

  • () -> void



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

def stop
  return if stopped?

  @stopped = true
  process_registry.start_draining
  process_registry.stop
end

#stopped?Boolean

RBS:

  • () -> bool

Returns:

  • (Boolean)


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

def stopped?
  @stopped
end

#verify_claim!(reminder) ⇒ void

This method returns an undefined value.

RBS:

  • (Reminder) -> void

Parameters:



152
153
154
155
156
# File 'lib/solid_objects/reminder_scheduler.rb', line 152

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

  raise LostActivation, "reminder claim changed"
end