class Telemetry::TCP::Server

Constants

TIME_TO_LIVE

Public Class Methods

new(port) click to toggle source
Calls superclass method
# File lib/telemetry/tcp/server.rb, line 9
def initialize(port)
    super(self.class.name,::Orocos::Async.event_loop)
    @server = TCPServer.open(port)
    @server.setsockopt(:SOCKET, :REUSEADDR, true)
    @server.setsockopt(Socket::IPPROTO_TCP, Socket::TCP_NODELAY,false)
    @clients = Hash.new do |h,key|
        h[key] = []
    end

    #wait for client connections
    p = Proc.new do |client,error|
        event_loop.async_with_options(@server.method(:accept),{:sync_key => @server},&p) unless @server.closed?

        # do not do handshake to prevent traffic to the robot
        # which is not allowed by the spacebot rules !!!
        if !error && client
            client.setsockopt(Socket::IPPROTO_TCP, Socket::TCP_NODELAY, false)
            @clients[client.remote_address.ip_address] << client
            emit_connected client
            emit_port_connected client.remote_address.ip_port
        else
            Telemetry.error error
        end
    end
    event_loop.async_with_options(@server.method(:accept),{:sync_key => @server},&p)

    on_disconnected do |client|
        begin
            @clients.each do |key,ports|
                ports.delete(client)
                emit_ip_disconnected key if ports.empty?
            end
            emit_port_disconnected client.remote_address.ip_port
        rescue Errno::ENOTCONN,IOError
        end
    end

    #watch dog
    event_loop.every(1) do
        closed?
    end
end

Public Instance Methods

close() click to toggle source
# File lib/telemetry/tcp/server.rb, line 52
def close
    @clients.each_value do |clients|
        clients.each do |c|
            emit_disconnected c
        end
    end
    @server.close
end
closed?() click to toggle source
# File lib/telemetry/tcp/server.rb, line 70
def closed?
    @clients.each_pair do |key,sockets|
        s = sockets.find do |s|
            if !event_loop.thread_pool.sync_keys.include?(s) && s.closed?
                emit_disconnected s
            else
                s
            end
        end
        return false if s
    end
    true
end
eof?() click to toggle source
# File lib/telemetry/tcp/server.rb, line 61
def eof?
    @clients.each_pair do |key,sockets|
        s = sockets.find do |s|
            !event_loop.thread_pool.sync_keys.include?(s) && !s.eof?
        end
        return s if s
    end
end
write(id,string) click to toggle source
# File lib/telemetry/tcp/server.rb, line 84
def write(id,string)
    frame = WebSocket::Frame::Outgoing::Server.new(:data => string.to_s,:type => :binary)
    raise "WebSocket frame type is not supported by the selected draft!" unless frame.support_type?
    msg = frame.to_s

    #write string to each ip address
    total_bytes = 0
    @clients.each_pair do |key,client|
        next if client.empty?

        # select socket
        id_temp = id%(client.size)
        s = client[id_temp]
        bytes = msg.size

        # get unique access to the socket
        event_loop.sync_timeout(s,TIME_TO_LIVE) do
            # send hole message
            while bytes != 0
                begin
                    b = s.syswrite(msg[msg.size-bytes..-1])
                    bytes -= b
                    total_bytes += b
                rescue Errno::EBADF,Errno::EPIPE,Errno::ECONNRESET,IOError => e
                    Telemetry.warn "error on io: #{e}"
                    s.close unless s.closed?
                rescue Errno::EAGAIN
                    Telemetry.warn "try again"
                    sleep 0.1
                    next
                end
                if s.closed?
                    event_loop.call do
                        emit_disconnected s
                    end
                    break
                end
            end
        end
    end
    total_bytes
end
write_nonblock(id,string,&block) click to toggle source
# File lib/telemetry/tcp/server.rb, line 127
def write_nonblock(id,string,&block)
    #write string to each ip address
    event_loop.async_with_options(self.method(:write),{:known_errors=>[Timeout::Error]},id,string) do |bok,error|
        block.call(bok,error) if block_given?
    end
end