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.
-
#process_registry ⇒ Object
readonly
Returns the value of attribute process_registry.
Instance Method Summary collapse
- #claim_next(now:) ⇒ Reminder?
- #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.
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_adapter ⇒ Object (readonly)
Returns the value of attribute database_adapter.
76 77 78 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 76 def database_adapter @database_adapter end |
#process_registry ⇒ Object (readonly)
Returns the value of attribute process_registry.
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?
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?
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..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
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?
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.
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_shutdown ⇒ void
This method returns an undefined value.
59 60 61 62 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 59 def request_shutdown @shutdown_requested = true SolidObjects.wake_up.signal end |
#run ⇒ void
This method returns an undefined value.
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
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
70 71 72 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 70 def shutdown_requested? @shutdown_requested end |
#stop ⇒ void
This method returns an undefined value.
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
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.
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 |