class GxG::Events::EventDispatcher

Provide a Cooperative Processing facility for sessions

Public Class Methods

new(interval=0.333, thread_reservation=100) click to toggle source
Calls superclass method
# File lib/gxg/gxg_events.rb, line 114
def initialize(interval=0.333, thread_reservation=100)
  super()
  @reserved_threads = (thread_reservation.to_i || 100)
  @uuid = ::GxG::uuid_generate().to_sym
  @event_queues = {:root => {:events => [], :settings => {:active => true, :portion => 100.0}}}
  @running = false
  @ticking = false
  # Timer format: {:when => Time_Object.to_f, :what => event_frame, :interval => float, :id => uuid.sym}
  @timers = []
  @scheduler = Rufus::Scheduler.new(:max_work_threads => @reserved_threads)
  def @scheduler.on_error(job, error)
    log_error({:error => error, :parameters => {:job => job}})
  end
  @scheduler.pause
  @tick_timer = nil
  @tick_interval = interval
  #
  self
end

Public Instance Methods

adjust_event_queue(queue=:root,settings={}) click to toggle source
# File lib/gxg/gxg_events.rb, line 259
def adjust_event_queue(queue=:root,settings={})
  # RESEARCH: GxG::Events::EventDispatcher.adjust_event_queue : refactor :portion of each queue algo needed.
  unless queue.to_s.to_sym == :root
    if @event_queues[(queue.to_sym)]
      # LATER: Eventually add :portion as valid setting to accept.
      [:active].each do |the_setting|
        if settings[(the_setting.to_sym)]
          new_setting = {}
          # portion adjustment algo here
          new_setting[(the_setting.to_sym)] = settings[(the_setting.to_sym)]
          #
          @event_queues[(queue.to_sym)][:settings].merge!(new_setting)
        end
      end
    else
      puts "Warning: queue #{queue.inspect} does not exist"
      # For Now: just ignore references to non-existent queues. (reconsider later)
      # raise ArgumentError, "Event Queue :#{queue} does not exist to alter"
    end
  end
end
at(expression="",&block) click to toggle source
# File lib/gxg/gxg_events.rb, line 364
def at(expression="",&block)
  the_time = nil
  case expression.class
  when ::DateTime
    the_time = expression.to_time()
  when ::Time
    the_time = expression
  when ::String
    the_time = ::Chronic::parse(expression)
    unless the_time
      raise ArgumentError, "Invalid time expression: #{expression.inspect}"
    end
  else
    raise ArgumentError, "Invalid time expression: #{expression.inspect}"
  end
  if the_time
    @scheduler.at(the_time.iso8601,&block)
  else
    nil
  end
end
cancel_timer(timer_reference=nil) click to toggle source

Adding Timers See: stackoverflow.com/questions/235504/validating-crontab-entries-with-php Cron REGEX: /^((?:[1-9]?d|*)s*(?:(?:[1-9]?d)|(?:,?d)+)?s*){5}$/ Rufus Duration REGEX: /^(-?)([d.smhdwy]+)$/ at, in, every, cron

# File lib/gxg/gxg_events.rb, line 354
def cancel_timer(timer_reference=nil)
  @scheduler.jobs().each do |the_task|
    if the_task.job_id == timer_reference
      the_task.unschedule
      break
    end
  end
  true
end
create_event_queue(settings={}) click to toggle source
# File lib/gxg/gxg_events.rb, line 281
def create_event_queue(settings={})
  unless settings.is_a?(Hash)
    raise ArgumentError, "you must provide a Hash as a parameter set"
  end
  queue_name = settings.delete(:name).to_s.to_sym
  if queue_name.to_s.size > 0
    if @event_queues[(queue_name)]
      raise ArgumentError, "Event Queue #{queue_name.inspect} already exists"
    else
      @event_queues[(queue_name)] = {:events => [], :settings => {:active => false, :portion => 0.0}}
      # For Now: simply evenly divide the portions equally, only :root can be 100%.
      even_portion = (100.0 / @event_queues.keys.size.to_f)
      @event_queues.keys.each do |the_queue|
        @event_queues[(the_queue)][:settings][:portion] = (even_portion)
      end
      self.adjust_event_queue(queue_name, settings.merge({:active => true}))
      true
    end
  else
    raise ArgumentError, ":name must be provided for Event Queue"
  end
end
cron(expression="",&block) click to toggle source
# File lib/gxg/gxg_events.rb, line 496
def cron(expression="",&block)
  if /^((?:[1-9]?\d|\*)\s*(?:(?:[\/-][1-9]?\d)|(?:,[1-9]?\d)+)?\s*){5}$/.match(expression)
    @scheduler.cron(expression,&block)
  else
    raise ArgumentError, "Invalid CRON expression: #{expression.inspect}"
  end
end
delete_event_queue(queue=nil) click to toggle source
# File lib/gxg/gxg_events.rb, line 304
def delete_event_queue(queue=nil)
  if queue.to_s.to_sym == :root
    false
  else
    @event_queues.delete(queue)
    # For Now: simply evenly divide the portions equally, only :root can be 100%.
    even_portion = (@event_queues.keys.size.to_f / 100.0)
    @event_queues.keys.each do |the_queue|
      @event_queues[(the_queue)][:settings][:portion] = (even_portion)
    end
    true
  end
end
every(expression="", &block) click to toggle source
# File lib/gxg/gxg_events.rb, line 441
def every(expression="", &block)
  if expression.is_a?(::Numeric)
    interval = expression
  else
    if expression.include?(" ")
      # parse to Rufus duration
      duration_type = nil
      interval = nil
      expression.split(" ").each do |entry|
        if duration_type && interval
          break
        else
          if ["s","second", "seconds"].include?(entry)
            duration_type = "s"
          end
          if ["m","minute", "minutes"].include?(entry)
            duration_type = "m"
          end
          if ["h","hour", "hours"].include?(entry)
            duration_type = "h"
          end
          if ["d","day", "days"].include?(entry)
            duration_type = "d"
          end
          if ["w","week", "weeks"].include?(entry)
            duration_type = "w"
          end
          if /[-+0-9.,]*/.match(entry).to_s.size > 0
            if entry.include?(".")
              interval = entry.to_f
            else
              interval = entry.to_i
            end
          end
        end
      end
      the_interval = (interval.to_s + duration_type.to_s)
    else
      # Verify standard Rufus expression
      duration_type = expression.split(/[-+0-9.,]*/)[1].to_s
      if ["s","m","h","d","w"].include?(duration_type.downcase)
        interval = /[-+0-9.,]*/.match(expression).to_s.to_i
        if interval > 0
          the_interval = (interval.to_s + duration_type.to_s)
        else
          raise ArgumentError, "Invalid interval expression: #{expression.inspect}"
        end
      else
        raise ArgumentError, "Invalid interval expression: #{expression.inspect}"
      end
    end
  end
  @scheduler.every(the_interval,&block)
end
in(expression="", &block) click to toggle source
# File lib/gxg/gxg_events.rb, line 386
def in(expression="", &block)
  if expression.is_a?(::Numeric)
    interval = expression
  else
    if expression.include?(" ")
      # parse to Rufus duration
      duration_type = nil
      interval = nil
      expression.split(" ").each do |entry|
        if duration_type && interval
          break
        else
          if ["s","second", "seconds"].include?(entry)
            duration_type = "s"
          end
          if ["m","minute", "minutes"].include?(entry)
            duration_type = "m"
          end
          if ["h","hour", "hours"].include?(entry)
            duration_type = "h"
          end
          if ["d","day", "days"].include?(entry)
            duration_type = "d"
          end
          if ["w","week", "weeks"].include?(entry)
            duration_type = "w"
          end
          if /[-+0-9.,]*/.match(entry).to_s.size > 0
            if entry.include?(".")
              interval = entry.to_f
            else
              interval = entry.to_i
            end
          end
        end
      end
      interval = (interval.to_s + duration_type.to_s)
    else
      # Verify standard Rufus expression
      duration_type = expression.split(/[-+0-9.,]*/)[1].to_s
      if ["s","m","h","d","w"].include?(duration_type.downcase)
        interval = /[-+0-9.,]*/.match(expression).to_s.to_i
        if interval > 0
          interval = (interval.to_s + duration_type.to_s)
        else
          raise ArgumentError, "Invalid interval expression: #{expression.inspect}"
        end
      else
        raise ArgumentError, "Invalid interval expression: #{expression.inspect}"
      end
    end
  end
  @scheduler.in(interval,&block)
end
inspect_queue(queue=:root) click to toggle source
# File lib/gxg/gxg_events.rb, line 255
def inspect_queue(queue=:root)
  @event_queues[(queue)]
end
inspect_timers() click to toggle source
# File lib/gxg/gxg_events.rb, line 252
def inspect_timers()
  @scheduler.jobs()
end
pause_event_queue(queue=nil) click to toggle source
# File lib/gxg/gxg_events.rb, line 318
def pause_event_queue(queue=nil)
  unless queue.to_s.to_sym == :root
    self.adjust_event_queue(queue,{:active => true})
  end
end
post_event(queue_name=:root,&block) click to toggle source
# File lib/gxg/gxg_events.rb, line 341
def post_event(queue_name=:root,&block)
  unless @event_queues[(queue_name.to_sym)]
    self.create_event_queue({:name => queue_name.to_sym})
  end
  self.post_to_event_queue(queue_name.to_sym,block)
  true
end
post_to_event_queue(queue=:root,the_event=nil) click to toggle source
# File lib/gxg/gxg_events.rb, line 330
def post_to_event_queue(queue=:root,the_event=nil)
  if the_event.respond_to?(:call)
    if @event_queues[(queue)]
      @event_queues[(queue)][:events] << the_event
    else
      # issue a warning, and post to :root queue anyways (for now). (link to logger, and output)
      @event_queues[:root][:events] << the_event
    end
  end
end
running?() click to toggle source
# File lib/gxg/gxg_events.rb, line 134
def running?()
  @running
end
shutdown() click to toggle source
# File lib/gxg/gxg_events.rb, line 159
def shutdown()
  @running = false
  if @tick_timer
    self.cancel_timer(@tick_timer)
    @tick_timer = nil
    @scheduler.pause
    @reserved_threads.times do
      ::GxG::Engine::release_event_descriptor()
    end
  end
  true
end
startup() click to toggle source
# File lib/gxg/gxg_events.rb, line 138
def startup()
  @running = true
  if @scheduler.paused?
    @scheduler.resume
  end
  unless @tick_timer
    @tick_timer = @scheduler.every(@tick_interval) do
      tick
    end
  end
  @reserved_threads.times do
    begin
      ::GxG::Engine::reserve_event_descriptor()
    rescue Exception => the_error
      log_error({:error => the_error})
      break
    end
  end
  @running
end
tick() click to toggle source
# File lib/gxg/gxg_events.rb, line 172
def tick()
  if @running
    unless @ticking
      @ticking = true
      events = []
      envelope = ::GxG::Engine::event_allocation_envelope()
      if envelope.last >= envelope.first
        prefetch = {}
        @event_queues.keys.each do |the_queue|
          if (@event_queues[(the_queue)][:settings][:active] && @event_queues[(the_queue)][:events].size > 0)
            ((@event_queues[(the_queue)][:settings][:portion].to_f / 100.0) * envelope.last.to_f).to_i.times do
              op = @event_queues[(the_queue)][:events].shift
              if op.respond_to?(:call)
                (prefetch[(the_queue)] ||= []) << op
              end
            end
          end
        end
        #
        prefetch_count = 0
        prefetch.keys.each do |the_queue|
          prefetch_count += prefetch[(the_queue)].size
        end
        #
        while prefetch_count > 0 do
          prefetch.keys.each do |the_queue|
            op = prefetch[(the_queue)].shift
            if op
              events << op
            end
            prefetch_count -= 1
          end
        end
        #
      end
      # xxx old version:
      # queues = {}
      # load_size = 0
      # @event_queues.keys.each do |the_queue|
      #   if (@event_queues[(the_queue)][:settings][:active] && @event_queues[(the_queue)][:events].size > 0)
      #     queues[(the_queue)] = {:portion => ((@event_queues[(the_queue)][:settings][:portion].to_f / 100.0) * total_slots.to_f).to_i}
      #     if queues[(the_queue)][:portion] == 0
      #       # Fudging a little here in case of real heavy loads and/or too many queues.
      #       queues[(the_queue)][:portion] = 1
      #     end
      #     queues[(the_queue)][:size] = @event_queues[(the_queue)][:events].size
      #     load_size += queues[(the_queue)][:size]
      #   end
      # end
      # queues.keys.each do |the_queue|
      #   # TODO: devise apportioned event dispatch
      #   queues[(the_queue)][:portion].times do
      #     op = @event_queues[(the_queue)][:events].shift
      #     if op.respond_to?(:call)
      #       events << op
      #     else
      #       break
      #     end
      #   end
      # end
      #
      events.each do |the_event|
        if the_event
          Thread.new do
            ::GxG::Engine::reserve_event_descriptor()
            begin
              the_event.call()
            rescue Exception => the_error
              log_error({:error => the_error})
            end
            ::GxG::Engine::release_event_descriptor()
          end
        end
      end
      #
      @ticking = false
    end
  end
end
unpause_event_queue(queue=nil) click to toggle source
# File lib/gxg/gxg_events.rb, line 324
def unpause_event_queue(queue=nil)
  unless queue.to_s.to_sym == :root
    self.adjust_event_queue(queue,{:active => false})
  end
end