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 ⇒ Reminder?
- #enqueue(reminder) ⇒ MessageReference?
-
#initialize(process_registry: ProcessRegistry.new, database_adapter: SolidObjects.database_adapter) ⇒ ReminderScheduler
constructor
A new instance of ReminderScheduler.
- #next_run_at(reminder, now) ⇒ Time
- #release(reminder) ⇒ void
- #request_shutdown ⇒ void
- #run ⇒ void
- #run_once ⇒ 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.
75 76 77 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 75 def database_adapter @database_adapter end |
#process_registry ⇒ Object (readonly)
Returns the value of attribute process_registry.
75 76 77 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 75 def process_registry @process_registry end |
Instance Method Details
#claim_next ⇒ Reminder?
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?
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..key?(reminder..to_sym) raise UnknownMessage, "unknown reminder message #{reminder..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) now = database_adapter.database_now = mailbox.enqueue_in_transaction( Reference.new(actor_type: instance.actor_type, actor_id: instance.actor_id), locked_reminder., 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 ) 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
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.
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_shutdown ⇒ void
This method returns an undefined value.
58 59 60 61 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 58 def request_shutdown @shutdown_requested = true SolidObjects.wake_up.signal end |
#run ⇒ void
This method returns an undefined value.
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_once ⇒ 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
69 70 71 |
# File 'lib/solid_objects/reminder_scheduler.rb', line 69 def shutdown_requested? @shutdown_requested end |
#stop ⇒ void
This method returns an undefined value.
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
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.
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 |