class GxG::Networking::ZmqNode
xxx
class AbstractZMQSocket # Abstract class, never instantiated def self.supported_protocols() if ::GxG::SYSTEM.platform()[:platform] == :windows ["inproc", "tcp", "tcp4", "tcp6", "pgm", "epgm"] else ["inproc", "ipc", "tcp", "tcp4", "tcp6", "pgm", "epgm"] end end # def initialize(settings={}) # {:context => nil, :type => nil, :bind => {:thisone => "", :thatone => ""}, :connect => nil} unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end @thread_safety = ::Mutex.new if settings[:context].is_a?(::ZMQ::Context) @context = settings[:context] else @context = ::GxG::Networking::zmq_context() end if settings[:type].is_a?(::Numeric) @socket = @context.socket(settings[:type]) else raise ArgumentError, "you must supply a type for the socket" end @name = ::Celluloid::UUID.random_generate().to_sym socket_id = (settings[:identity] || @name).to_s.dup if socket_id.bytesize() > 255 socket_id.force_encoding(::Encoding::ASCII_8BIT) socket_id = socket_id[(0..255)] end @socket.identity = socket_id # @recv_buffer_size = 1024 @send_buffer_size = 1024 # @bindings = {} if settings[:bind].is_a?(::Hash) # settings[:bind].keys.to_enum(:each).each do |binding_key| if settings[:bind][(binding_key)].is_a?(::URI::Generic) address = settings[:bind][(binding_key)] else address = ::URI::parse(settings[:bind][(binding_key)]) end if ::GxG::Networking::AbstractZMQSocket::supported_protocols().include?(address.scheme()) address.resolve_host() if ["tcp", "tcp4", "tcp6", "udp", "udp4", "udp6"].include?(address.scheme()) if ::GxG::SYSTEM.network_port_used?(::Addrinfo::parse(address)) raise Errno::EADDRINUSE, "Address and port already in use: #{address.to_s}" end end begin @socket.bind(address.to_s) @bindings[(binding_key)] = address rescue Exception @socket.close raise end else raise ArgumentError, "#{address.scheme().to_s} is an unsupported protocol" end end # end # @connections = {} if settings[:connect].is_a?(::Hash) settings[:connect].keys.to_enum(:each).each do |connection_key| if settings[:connect][(connection_key)].is_a?(::URI::Generic) address = settings[:connect][(connection_key)] else address = ::URI::parse(settings[:connect][(connection_key)]) end # address = ::URI::parse(settings[:connect][(connection_key)]) if ::GxG::Networking::AbstractZMQSocket::supported_protocols().include?(address.scheme()) address.resolve_host() unless ::GxG::SYSTEM.network_port_used?(::Addrinfo::parse(address)) raise Errno::EADDRNOTAVAIL, "Address and port not bound: #{address.to_s}" end begin @socket.connect(address.to_s) @connections[(connection_key)] = address rescue Exception @socket.close raise end else raise ArgumentError, "#{address.scheme().to_s} is an unsupported protocol" end end end # @recv_buffer_size = this.recv_buffer_size() @send_buffer_size = this.send_buffer_size() # this end # def context() @thread_safety.synchronize { @context } end # def socket() @thread_safety.synchronize { @socket } end # def identity() @socket.identity() end # def recv_buffer_size(*args) # if args.size > 0 if args[0].is_a?(::Numeric) the_size = args.first.to_i if the_size > 0 valid_min = [16384] valid_max = [65536] [@bindings, @connections].to_enum(:each).each do |the_set| the_set.keys.to_enum(:each).each do |the_terminus| setting = ::GxG::Networking::buffer_limits(:in, (the_set[(the_terminus)].scheme).to_sym)[:valid] valid_min << (setting.min) valid_max << (setting.max) end end if ((valid_min.min)..(valid_max.max)).include?(the_size) @thread_safety.synchronize { @recv_buffer_size = the_size @socket.setsockopt(::ZMQ::RCVBUF,the_size) } else raise Exception, "New buffer size #{the_size} is outside acceptable range: #{((valid_min.min)..(valid_max.max)).inspect}" end end end else data = [] @thread_safety.synchronize { @socket.getsockopt(::ZMQ::RCVBUF,data) } the_size = (data.first || 0).to_i if the_size == 0 # find scheme of connection/binding : the_scheme # buffer_size = (::GxG::Networking::buffer_limits(:out, the_scheme.to_sym)[:initial] || 65536) buffer_sizes = [65536] [@bindings, @connections].to_enum(:each).each do |the_set| the_set.keys.to_enum(:each).each do |the_terminus| buffer_sizes << (::GxG::Networking::buffer_limits(:in, (the_set[(the_terminus)].scheme).to_sym)[:valid].min) end end # FORNOW: use the smallest buffer size found across all bindings and/or connections @thread_safety.synchronize { @recv_buffer_size = buffer_sizes.min } else @thread_safety.synchronize { @recv_buffer_size = the_size } end end @thread_safety.synchronize { @recv_buffer_size.dup } end # def send_buffer_size(*args) # if args.size > 0 if args[0].is_a?(::Numeric) the_size = args.first.to_i if the_size > 0 valid_min = [16384] valid_max = [65536] [@bindings, @connections].to_enum(:each).each do |the_set| the_set.keys.to_enum(:each).each do |the_terminus| setting = ::GxG::Networking::buffer_limits(:out, (the_set[(the_terminus)].scheme).to_sym)[:valid] valid_min << (setting.min) valid_max << (setting.max) end end if ((valid_min.min)..(valid_max.max)).include?(the_size) @thread_safety.synchronize { @send_buffer_size = the_size @socket.setsockopt(::ZMQ::SNDBUF,the_size) } else raise Exception, "New buffer size #{the_size} is outside acceptable range: #{((valid_min.min)..(valid_max.max)).inspect}" end end end else data = [] @thread_safety.synchronize { @socket.getsockopt(::ZMQ::SNDBUF,data) } the_size = (data.first || 0).to_i if the_size == 0 # find scheme of connection/binding : the_scheme # buffer_size = (::GxG::Networking::buffer_limits(:out, the_scheme.to_sym)[:initial] || 65536) buffer_sizes = [65536] [@bindings, @connections].to_enum(:each).each do |the_set| the_set.keys.to_enum(:each).each do |the_terminus| buffer_sizes << (::GxG::Networking::buffer_limits(:out, (the_set[(the_terminus)].scheme).to_sym)[:valid].min) end end # FORNOW: use the smallest buffer size found across all bindings and/or connections @thread_safety.synchronize { @send_buffer_size = buffer_sizes.min } else @thread_safety.synchronize { @send_buffer_size = the_size } end end @thread_safety.synchronize { @send_buffer_size.dup } end # def binding(binding_key=nil) if binding_key @thread_safety.synchronize { @bindings[(binding_key)] } else @thread_safety.synchronize { @bindings } end end # def bind(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of keyed bind point(s) for the binding" end result = false settings.keys.to_enum(:each).each do |binding_key| if settings[(binding_key)].is_a?(::URI::Generic) address = settings[(binding_key)] else address = ::URI::parse(settings[(binding_key)]) end # address = ::URI::parse(settings[(binding_key)]) if ::GxG::Networking::AbstractZMQSocket::supported_protocols().include?(address.scheme()) address.resolve_host() if ["tcp", "tcp4", "tcp6", "udp", "udp4", "udp6"].include?(address.scheme()) if ::GxG::SYSTEM.network_port_used?(::Addrinfo::parse(address)) raise Errno::EADDRINUSE, "Address and port already in use: #{address.to_s}" end end begin @thread_safety.synchronize { @socket.bind(address.to_s) @bindings[(binding_key)] = address } result = true rescue Exception @thread_safety.synchronize { @socket.close } raise end else raise ArgumentError, "#{address.scheme().to_s} is an unsupported protocol" end end @recv_buffer_size = this.recv_buffer_size() @send_buffer_size = this.send_buffer_size() # result end # def connection(connection_key=nil) if connection_key @thread_safety.synchronize { @connections[(connection_key)] } else @thread_safety.synchronize { @connections } end end # def connect(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of keyed address(es) for the connection" end result = false settings.keys.to_enum(:each).each do |connection_key| if settings[(connection_key)].is_a?(::URI::Generic) address = settings[(connection_key)] else address = ::URI::parse(settings[(connection_key)]) end # address = ::URI::parse(settings[(connection_key)]) if ::GxG::Networking::AbstractZMQSocket::supported_protocols().include?(address.scheme()) address.resolve_host() if ["tcp", "tcp4", "tcp6"].include?(address.scheme()) unless ::GxG::SYSTEM.network_port_used?(::Addrinfo::parse(address)) raise Errno::EADDRNOTAVAIL, "Address and port not bound: #{address.to_s}" end end begin @thread_safety.synchronize { @socket.connect(address.to_s) @connections[(connection_key)] = address } result = true rescue Exception @thread_safety.synchronize { @socket.close() } raise end else raise ArgumentError, "#{address.scheme().to_s} is an unsupported protocol" end end @recv_buffer_size = this.recv_buffer_size() @send_buffer_size = this.send_buffer_size() # result end # def close() @thread_safety.synchronize { @socket.close() } end # def triage_strings(data=[]) # takes raw output from @socket.recv_strings([]) and separates routing from content result = {:route => [], :content => []} prior_element = nil data.to_enum(:each).each do |the_element| # Empty strings terminate routes -- Do NOT use in payload as parse element if the_element.to_s.size > 0 if prior_element.to_s.size > 0 result[:content] << prior_element end prior_element = the_element else if prior_element.to_s.size > 0 result[:route] << prior_element end prior_element = nil end end if prior_element.to_s.size > 0 result[:content] << prior_element end result end # def recv(options={:autojoin => true}) unless options.is_a?(::Hash) options = {:autojoin => true} end payload = [] # @thread_safety.lock status = @socket.recv_strings(payload,0) if ::ZMQ::Util.resultcode_ok?(status) data = this.triage_strings(payload) if (data[:route].size > 0 || data[:content].size > 0) if options[:autojoin] ::GxG::Events::Message.new({:sender => this, :body => data[:content].join(), :route => data[:route]}) else ::GxG::Events::Message.new({:sender => this, :body => data[:content], :route => data[:route]}) end else nil end else begin raise Exception, ("ZMQ - " + ::ZMQ::Util.error_string()) rescue => the_error log_error({:error => the_error}) end nil end # @thread_safety.unlock end # def send(data=nil, route=[]) if data # FORNOW: use the smallest buffer size found across all bindings and/or connections buffer_size = @thread_safety.synchronize { @send_buffer_size } payload = [] messages = [] unless route.is_a?(::Array) route = [(route)] end # if data.is_a?(::GxG::Events::Message) if data.body().is_a?(::String) data = [(data.body())] else data = [(data.body().serialize())] end end if data.is_a?(::String) data = [(data)] end # if route.is_a?(::Array) route.flatten.to_enum(:each).each do |the_route| if the_route.is_a?(::String) payload << the_route.slice_bytes(0..255) payload << "" end end end payload << @socket.identity() payload << "" # if data.is_a?(::Array) data.flatten.to_enum(:each).each do |the_element| if the_element.is_a?(::String) if the_element.size > 0 # Disallow parsing payload with empty strings : used to parse routes # Break into appropriate buffer size as needed : respect buffer sizes if the_element.bytesize > buffer_size length = the_element.bytesize() passes = ::GxG::Networking::passes_needed(length, buffer_size) total = (length - 1) starting = 0 ending = ([length,buffer_size].min - 1) passes.times do messages << the_element.slice_bytes(starting..ending) # starting = ending + 1 ending = (starting + (buffer_size - 1)) if ending > total ending = ending - (ending - total) end end # else messages << the_element end end end end if messages.size > 0 payload << messages payload.flatten! status = @thread_safety.synchronize { @socket.send_strings(payload) } if status == -1 begin raise Exception, "ZMQ - Message could not be enqueued." rescue => the_error log_error({:error => the_error, :parameters => {:data => data, :route => route}}) end end end end # end end # end class ZMQRouter < ::GxG::Networking::AbstractZMQSocket # def initialize(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end unless settings[:context].is_a?(::ZMQ::Context) settings[:context] = ::GxG::Networking::zmq_context() end settings[:type] = ::ZMQ::ROUTER bindings = (settings.delete(:bind) || {}) super(settings) unless bindings.keys.size > 0 bindings[:inproc] = ::URI::parse("inproc://#{@name}") end this.bind(bindings) # this end # end class ZMQDealer < ::GxG::Networking::AbstractZMQSocket # def initialize(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end settings[:type] = ::ZMQ::DEALER super(settings) # this end # end class ZMQRequest < ::GxG::Networking::AbstractZMQSocket # def initialize(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end settings[:type] = ::ZMQ::REQ super(settings) # this end # end class ZMQReply < ::GxG::Networking::AbstractZMQSocket # def initialize(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end settings[:type] = ::ZMQ::REP super(settings) # this end # end class ZMQPuller < ::GxG::Networking::AbstractZMQSocket # def initialize(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end settings[:type] = ::ZMQ::PULL super(settings) # this end # def send(*args) # no-op end # end class ZMQPusher < ::GxG::Networking::AbstractZMQSocket # def initialize(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end settings[:type] = ::ZMQ::PUSH super(settings) # this end # def recv() nil end # end class ZMQPair < ::GxG::Networking::AbstractZMQSocket # def initialize(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end settings[:type] = ::ZMQ::PAIR super(settings) # this end # end class ZMQSubscriber < ::GxG::Networking::AbstractZMQSocket # def initialize(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end settings[:type] = ::ZMQ::SUB super(settings) @topics = [] # this end # def topics() @thread_safety.synchronize { @topics.dup } end # def subscribe_topic(the_topic=nil) the_topic = the_topic.to_s if the_topic.bytesize > 255 the_topic = the_topic.slice_bytes(0..255).to_s end if the_topic.bytesize > 0 @thread_safety.synchronize { unless @topics.include?(the_topic) @topics << the_topic @socket.setsockopt(::ZMQ::SUBSCRIBE,the_topic) end } end end # def unsubscribe_topic(the_topic=nil) the_topic = the_topic.to_s if the_topic.bytesize > 255 the_topic = the_topic.slice_bytes(0..255).to_s end if the_topic.bytesize > 0 @thread_safety.synchronize { if @topics.include?(the_topic) @topics.delete(the_topic) @socket.setsockopt(::ZMQ::UNSUBSCRIBE,the_topic) end } end end # def send(*args) # no-op end # end class ZMQPublisher < ::GxG::Networking::AbstractZMQSocket # def initialize(settings={}) # unless settings.is_a?(::Hash) raise ArgumentError, "you must supply a hash of settings for the socket" end settings[:type] = ::ZMQ::PUB bindings = (settings.delete(:bind) || {}) super(settings) unless bindings.keys.size > 0 bindings[:inproc] = ::URI::parse("inproc://#{@name}") end this.bind(bindings) # this end # def recv() nil end # def publish(the_data=nil,the_topic=nil) if the_data the_topic = the_topic.to_s if the_topic.bytesize > 255 the_topic = the_topic.slice_bytes(0..255).to_s end if the_topic.bytesize > 0 this.send(the_data,[(the_topic)]) else this.send(the_data) end end end end class ZMQ_RequestReplyBroker < ::GxG::Events::Actor # class ZMQ_RequestReplyBroker is a derivative work based upon: http://zguide.zeromq.org/rb:rrbroker Copyright (c) 2013 iMatix Corporation # accordingly, class ZMQ_RequestReplyBroker is licensed under this license: http://creativecommons.org/licenses/by-sa/3.0/legalcode def initialize(settings={}) unless settings.is_a?(::Hash) raise ArgumentError, "You must provide a hash of settings for the Broker." end unless settings[:frontend].is_a?(::URI::Generic) raise ArgumentError, "You must provide a URI as :frontend." end unless settings[:backend].is_a?(::URI::Generic) raise ArgumentError, "You must provide a URI as :backend." end super() # @context = ::GxG::Networking::zmq_context() @frontend = ::GxG::Networking::ZMQRouter.new({:context => @context, :bind => {:frontend => settings[:frontend]}}) @backend = ::GxG::Networking::ZMQDealer.new({:context => @context, :bind => {:backend => settings[:backend]}}) # @queue = ::ZMQ::Device.new(::ZMQ::QUEUE, @frontend.socket(), @backend.socket()) @queue = ::ZMQ::Poller.new @queue.register(@frontend.socket(), ::ZMQ::POLLIN) @queue.register(@backend.socket(), ::ZMQ::POLLIN) # @circulating = false @is_circulating = false @circulate_safety = ::Mutex.new # this.resume_circulating() ::Celluloid::current_actor() end # def terminate() this.halt_circulating() @queue.deregister(@frontend.socket()) @queue.deregister(@backend.socket()) @frontend.close() @backend.close() # super() end # def circulating?() @circulate_safety.synchronize { @circulating } end # def halt_circulating() @circulate_safety.synchronize { @circulating = false } end # def resume_circulating() unless this.circulating?() @circulate_safety.synchronize { @circulating = true } this.async.circulate end end # def circulate() # Circulate Data as required if this.alive?() front_socket = @frontend.socket() front_id = @frontend.identity() back_socket = @backend.socket() back_id = @frontend.identity() busy = @circulate_safety.synchronize { @is_circulating } while (this.circulating?() && !busy) do # @circulate_safety.synchronize { @is_circulating = true } @circulate_safety.lock # @queue.poll(:blocking) @queue.readables.each do |the_socket| if the_socket === front_socket data = @frontend.triage_strings(the_socket.recv_strings([])) payload = [] if (data[:route].size > 0 || data[:content].size > 0) data[:route].to_enum(:each).each do |the_route| payload << the_route payload << "" end payload << back_id payload << "" payload << data[:content] payload.flatten! if payload.size > 0 back_socket.send_strings(payload) end end end if the_socket === back_socket data = @backend.triage_strings(the_socket.recv_strings([])) payload = [] if (data[:route].size > 0 || data[:content].size > 0) data[:route].to_enum(:each).each do |the_route| payload << the_route payload << "" end payload << front_id payload << "" payload << data[:content] payload.flatten! if payload.size > 0 front_socket.send_strings(payload) end end end end # @circulate_safety.unlock @circulate_safety.synchronize { @is_circulating = false } busy = @circulate_safety.synchronize { @is_circulating } # if this.alive?() if this.circulating?() sleep 0.033 end end # end end end # end
Deprecated: completely rethink this.
Public Class Methods
new(*args)
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1026 def initialize(*args) # @zmq_thread_safety = ::Mutex.new @zmq_processing = false @zmq_is_processing = false @zmq_processing_handler = nil @zmq_processing_error_handler = nil @zmq_inputs = {} @zmq_outputs = {} # self end
Public Instance Methods
add_zmq_input(the_key=nil, the_socket=nil)
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1062 def add_zmq_input(the_key=nil, the_socket=nil) valid = [::GxG::Networking::ZMQPair, ::GxG::Networking::ZMQPuller, ::GxG::Networking::ZMQSubscriber, ::GxG::Networking::ZMQReply, ::GxG::Networking::ZMQRequest] # unless the_key.is_a?(::Symbol) raise ArgumentError, "you need so supply a unique symbol key for the socket" end if (@zmq_thread_safety.synchronize { @zmq_inputs[(the_key)] }) raise ArgumentError, "inputs already exists, you need so supply a unique symbol key for the socket" end unless the_socket.is_any?(valid) raise ArgumentError, "you supplied a #{the_socket.class.inspect}, you need so supply one of these as socket type: #{valid.inspect}" end @zmq_thread_safety.synchronize { @zmq_inputs[(the_key)] = the_socket } true end
add_zmq_output(the_key=nil, the_socket=nil)
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1100 def add_zmq_output(the_key=nil, the_socket=nil) valid = [::GxG::Networking::ZMQPair, ::GxG::Networking::ZMQPusher, ::GxG::Networking::ZMQPublisher, ::GxG::Networking::ZMQReply, ::GxG::Networking::ZMQRequest] # unless the_key.is_a?(::Symbol) raise ArgumentError, "you need so supply a unique symbol key for the socket" end if (@zmq_thread_safety.synchronize { @zmq_outputs[(the_key)] }) raise ArgumentError, "inputs already exists, you need so supply a unique symbol key for the socket" end unless the_socket.is_any?(valid) raise ArgumentError, "you supplied a #{the_socket.class.inspect}, you need so supply one of these as socket type: #{valid.inspect}" end @zmq_thread_safety.synchronize { @zmq_outputs[(the_key)] = the_socket } true end
halt_zmq_processing()
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1167 def halt_zmq_processing() @zmq_thread_safety.synchronize { @zmq_processing = false;@zmq_is_processing = false } end
remove_zmq_input(the_key=nil)
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1078 def remove_zmq_input(the_key=nil) # result = false the_inputs = this.zmq_inputs() if the_inputs[(the_key)] @zmq_thread_safety.synchronize { @zmq_inputs.delete(the_key) } result = true end result end
remove_zmq_output(the_key=nil)
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1116 def remove_zmq_output(the_key=nil) # result = false the_outputs = this.zmq_outputs() if the_outputs[(the_key)] @zmq_thread_safety.synchronize { @zmq_outputs.delete(the_key) } result = true end result end
resume_zmq_processing()
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1171 def resume_zmq_processing() unless this.processing?() @zmq_thread_safety.synchronize { @zmq_processing = true } this.async.zmq_process end end
terminate()
click to toggle source
Calls superclass method
# File lib/gxg/gxg_zmq.rb, line 1039 def terminate() # Halt any zmq processing first this.halt_zmq_processing() # Remove any references to input and output sockets # Question: should I close all input/output sockets here??? @zmq_thread_safety.synchronize { @zmq_inputs = {} @zmq_outputs = {} } super() end
zmq_error_handler(params={})
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1145 def zmq_error_handler(params={}) result = {:result => false} begin unless params.is_a?(:Hash) raise ArgumentError, "you must provide a Hash as parameter container" end if params[:block].respond_to(:call) @zmq_thread_safety.synchronize { @zmq_processing_error_handler = params[:block] } result[:result] = true else raise ArgumentError, "you must provide something that responds to :call with the :block parameter" end rescue Exception => the_error log_error({:error => the_error, :parameters => params}) end result end
zmq_handler(params={})
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1127 def zmq_handler(params={}) result = {:result => false} begin unless params.is_a?(:Hash) raise ArgumentError, "you must provide a Hash as parameter container" end if params[:block].respond_to(:call) @zmq_thread_safety.synchronize { @zmq_processing_handler = params[:block] } result[:result] = true else raise ArgumentError, "you must provide something that responds to :call with the :block parameter" end rescue Exception => the_error log_error({:error => the_error, :parameters => params}) end result end
zmq_inputs()
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1051 def zmq_inputs() result = {} @zmq_thread_safety.synchronize { @zmq_inputs.keys.to_enum(:each).each do |the_key| result[(the_key)] = @zmq_inputs[(the_key)] end # } result end
zmq_outputs()
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1089 def zmq_outputs() result = {} @zmq_thread_safety.synchronize { @zmq_outputs.keys.to_enum(:each).each do |the_key| result[(the_key)] = @zmq_outputs[(the_key)] end # } result end
zmq_process()
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1178 def zmq_process() # Process items on the zmq inputs and potentially send things via outputs while (this.zmq_processing?() && !(@zmq_thread_safety.synchronize { @zmq_is_processing })) do # @zmq_thread_safety.synchronize { @zmq_is_processing = true } # Process event item begin my_handler = @zmq_thread_safety.synchronize { @zmq_processing_error_handler } begin processor = @zmq_thread_safety.synchronize { @zmq_processing_handler } if processor.respond_to?(:call) processor.call(this, this.zmq_inputs(), this.zmq_outputs()) end rescue Exception => the_error # TODO: Revise logging/error-handling system generally, after display manager construction this.halt_zmq_processing() if my_handler.respond_to?(:call) # error handler provided my_handler.call({:error => the_error}) else # log error and forget about it log_error({:error => the_error}) end end rescue Exception => error_handling log_error({:error => error_handling}) end # @zmq_thread_safety.synchronize { @zmq_is_processing = false } if this.alive?() if this.zmq_processing?() sleep 0.033 end end # end # end
zmq_processing?()
click to toggle source
# File lib/gxg/gxg_zmq.rb, line 1163 def zmq_processing?() @zmq_thread_safety.synchronize { @zmq_processing } end