Class: SolidObjects::ReminderScheduler
- Inherits:
-
Object
- Object
- SolidObjects::ReminderScheduler
- Defined in:
- lib/solid_objects/reminder_scheduler.rb,
sig/generated/lib/solid_objects/reminder_scheduler.rbs
Instance Attribute Summary collapse
-
#database_adapter ⇒ Object
readonly
Returns the value of attribute database_adapter.
-
#polling_backoff ⇒ Object
readonly
Returns the value of attribute polling_backoff.
-
#process_registry ⇒ Object
readonly
Returns the value of attribute process_registry.
Instance Method Summary collapse
- #claim_next(now:) ⇒ Reminder?
- #current_polling_interval ⇒ Float
- #enqueue(reminder, now:) ⇒ MessageReference?
-
#initialize(process_registry: ProcessRegistry.new, database_adapter: SolidObjects.database_adapter) ⇒ ReminderScheduler
constructor
A new instance of ReminderScheduler.
- #next_run_at(reminder, now) ⇒ Time
- #normalize_test_time(now) ⇒ Time?
- #release(reminder) ⇒ void
- #request_shutdown ⇒ void
- #run ⇒ void
- #run_once(now: nil) ⇒ Boolean
- #shutdown_requested? ⇒ Boolean
- #stop ⇒ void
- #stopped? ⇒ Boolean
- #verify_claim!(reminder) ⇒ void
Constructor Details
#initialize(process_registry: ProcessRegistry.new, database_adapter: SolidObjects.database_adapter) ⇒ ReminderScheduler
Returns a new instance of ReminderScheduler.
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_adapter ⇒ Object (readonly)
Returns the value of attribute database_adapter.
107 108 109 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 107 def database_adapter @database_adapter end |
#polling_backoff ⇒ Object (readonly)
Returns the value of attribute polling_backoff.
107 108 109 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 107 def polling_backoff @polling_backoff end |
#process_registry ⇒ Object (readonly)
Returns the value of attribute process_registry.
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?
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_interval ⇒ 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?
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..key?(reminder.operation.to_sym) raise UnknownMessage, "unknown reminder operation #{reminder.operation.inspect}" end mailbox = Mailbox.new(database_adapter:) = 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 = 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 ) end return unless mailbox.announce() SolidObjects.instrument( :"reminder.enqueued", reminder_id: reminder.id, actor_type: reminder.actor_type, actor_id: reminder.actor_id, occurrence: reminder.occurrence ) MessageReference.() end |
#next_run_at(reminder, now) ⇒ 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?
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.
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_shutdown ⇒ void
This method returns an undefined value.
85 86 87 88 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 85 def request_shutdown @shutdown_requested = true SolidObjects.wake_up.signal end |
#run ⇒ void
This method returns an undefined value.
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
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
96 97 98 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 96 def shutdown_requested? @shutdown_requested end |
#stop ⇒ void
This method returns an undefined value.
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
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.
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 |