@@ -10,7 +10,13 @@ class Actor
1010 REMINDER_NAME_LIMIT = 191
1111 REMINDER_KEY_SEPARATOR = ":"
1212
13+ REMINDER_HANDLE_KEY = "reminder_name"
14+
1315 ReminderIntent = Data . define ( :name , :operation , :at , :arguments , :interval_seconds , :missed_policy )
16+ UnscheduleIntent = Data . define ( :name )
17+ UnscheduleAllIntent = Data . define ( :operation )
18+ ReminderStatus = Data . define ( :name , :operation , :key , :next_run_at , :interval_seconds ,
19+ :missed_policy , :occurrence , :status , :handle )
1420 OutboundMessageIntent = Data . define ( :actor_type , :actor_id , :operation , :arguments , :available_at , :idempotency_key )
1521
1622 class << self
@@ -152,10 +158,11 @@ def default_actor_type
152158
153159 attr_reader :actor_id , :state
154160
155- # @rbs (actor_id: String, state: State) -> void
156- def initialize ( actor_id :, state :)
161+ # @rbs (actor_id: String, state: State, ?instance_id: Integer? ) -> void
162+ def initialize ( actor_id :, state :, instance_id : nil )
157163 @actor_id = actor_id
158164 @state = state
165+ @instance_id = instance_id
159166 @effect_intents = [ ]
160167 @effect_recovery_intents = [ ]
161168 @commit_action_intents = [ ]
@@ -255,20 +262,139 @@ def schedule(at:, every: nil, missed: :latest, key: nil)
255262 actor_type : self . class . actor_type ,
256263 handlers : self . class . definition . messages
257264 ) do |operation , arguments |
258- ReminderIntent . new (
259- name : reminder_name ( operation :, key : reminder_key ) ,
265+ name = reminder_name ( operation :, key : reminder_key )
266+ reminder_intents << ReminderIntent . new (
267+ name :,
260268 operation : operation . to_s ,
261269 at :,
262270 arguments : Serialization . dump ( arguments ) ,
263271 interval_seconds :,
264272 missed_policy :
265- ) . tap do |intent |
266- reminder_intents << intent
267- end
268- nil
273+ )
274+ { REMINDER_HANDLE_KEY => name }
269275 end
270276 end
271277
278+ # @rbs (Symbol | String | reminder_handle, ?key: (String | Symbol | Integer)?) -> nil
279+ def unschedule ( operation_or_handle , key : nil )
280+ return unschedule_name ( handle_name ( operation_or_handle , key :) ) if operation_or_handle . is_a? ( Hash )
281+
282+ validated_reminder_operation ( operation_or_handle )
283+ unschedule_name ( reminder_name ( operation : operation_or_handle , key : validated_reminder_key ( key ) ) )
284+ end
285+
286+ # @rbs (Symbol | String) -> nil
287+ def unschedule_all ( operation )
288+ reminder_intents << UnscheduleAllIntent . new ( operation : validated_reminder_operation ( operation ) )
289+ nil
290+ end
291+
292+ # @rbs (Symbol | String | reminder_handle, ?key: (String | Symbol | Integer)?) -> ReminderStatus?
293+ def reminder ( operation_or_handle , key : nil )
294+ return reminder_view [ handle_name ( operation_or_handle , key :) ] if operation_or_handle . is_a? ( Hash )
295+
296+ validated_reminder_operation ( operation_or_handle )
297+ reminder_view [ reminder_name ( operation : operation_or_handle , key : validated_reminder_key ( key ) ) ]
298+ end
299+
300+ # @rbs (Symbol | String) -> Array[ReminderStatus]
301+ def reminders ( operation )
302+ wanted = validated_reminder_operation ( operation )
303+ reminder_view . each_value . select { |status | status . operation == wanted }
304+ end
305+
306+ attr_reader :instance_id
307+
308+ # @rbs (Symbol | String) -> String
309+ def validated_reminder_operation ( operation )
310+ name = operation . to_s
311+ return name if self . class . definition . messages . key? ( name . to_sym )
312+
313+ raise UnknownMessage , "unknown message #{ name . inspect } for #{ self . class . actor_type } "
314+ end
315+
316+ # @rbs (String) -> nil
317+ def unschedule_name ( name )
318+ reminder_intents << UnscheduleIntent . new ( name :)
319+ nil
320+ end
321+
322+ # @rbs (reminder_handle, key: untyped) -> String
323+ def handle_name ( handle , key :)
324+ raise ArgumentError , "a reminder handle already names its key" unless key . nil?
325+
326+ name = handle [ REMINDER_HANDLE_KEY ]
327+ unless name . is_a? ( String ) && !name . empty?
328+ raise InvalidPayload , "expected a reminder handle returned by schedule"
329+ end
330+
331+ name
332+ end
333+
334+ # The view is the committed schedule with this turn's staged intents applied
335+ # in order, so a read agrees with what the commit will leave behind.
336+ # @rbs () -> Hash[String, ReminderStatus]
337+ def reminder_view
338+ reminder_intents . each_with_object ( committed_reminders ) do |intent , view |
339+ apply_reminder_intent ( view , intent )
340+ end
341+ end
342+
343+ # @rbs () -> Hash[String, ReminderStatus]
344+ def committed_reminders
345+ return { } unless instance_id
346+
347+ Reminder . where ( instance_id :) . where . not ( status : "completed" ) . each_with_object ( { } ) do |row , view |
348+ view [ row . name ] = reminder_status (
349+ name : row . name ,
350+ operation : row . operation ,
351+ next_run_at : row . next_run_at ,
352+ interval_seconds : row . interval_seconds ,
353+ missed_policy : row . missed_policy ,
354+ occurrence : row . occurrence ,
355+ status : row . status
356+ )
357+ end
358+ end
359+
360+ # @rbs (Hash[String, ReminderStatus], untyped) -> void
361+ def apply_reminder_intent ( view , intent )
362+ return view . delete_if { |_name , status | status . operation == intent . operation } if intent . is_a? ( UnscheduleAllIntent )
363+ return view . delete ( intent . name ) if intent . is_a? ( UnscheduleIntent )
364+
365+ view [ intent . name ] = reminder_status (
366+ name : intent . name ,
367+ operation : intent . operation ,
368+ next_run_at : intent . at ,
369+ interval_seconds : intent . interval_seconds ,
370+ missed_policy : intent . missed_policy ,
371+ occurrence : view [ intent . name ] &.occurrence || 0 ,
372+ status : "scheduled"
373+ )
374+ end
375+
376+ # @rbs (name: String, operation: String, next_run_at: Time?, interval_seconds: untyped, missed_policy: String, occurrence: Integer, status: String) -> ReminderStatus
377+ def reminder_status ( name :, operation :, next_run_at :, interval_seconds :, missed_policy :, occurrence :, status :)
378+ ReminderStatus . new (
379+ name :,
380+ operation :,
381+ key : reminder_key_of ( name :, operation :) ,
382+ next_run_at :,
383+ interval_seconds : interval_seconds &.to_f ,
384+ missed_policy :,
385+ occurrence :,
386+ status :,
387+ handle : { REMINDER_HANDLE_KEY => name }
388+ )
389+ end
390+
391+ # @rbs (name: String, operation: String) -> String?
392+ def reminder_key_of ( name :, operation :)
393+ return nil if name == operation
394+
395+ name . delete_prefix ( "#{ operation } #{ REMINDER_KEY_SEPARATOR } " )
396+ end
397+
272398 # @rbs ((String | Symbol | Integer)?) -> String?
273399 def validated_reminder_key ( key )
274400 return nil if key . nil?
0 commit comments