140 lines
3.7 KiB
Ruby
140 lines
3.7 KiB
Ruby
require "socket"
|
|
require "openssl"
|
|
|
|
module MeowClient
|
|
TransportEvent = Struct.new(:type, :payload, keyword_init: true)
|
|
|
|
class Transport
|
|
CONNECT_TIMEOUT = 15
|
|
|
|
attr_reader :events
|
|
|
|
def self.connect(profile)
|
|
tcp = nil
|
|
socket = nil
|
|
begin
|
|
tcp = Socket.tcp(profile.host, profile.port, :connect_timeout => CONNECT_TIMEOUT)
|
|
if profile.tls
|
|
context = EltenAPI::TLS.client_context
|
|
socket = OpenSSL::SSL::SSLSocket.new(tcp, context)
|
|
socket.hostname = profile.host if socket.respond_to?(:hostname=)
|
|
socket.sync_close = true
|
|
socket.connect
|
|
else
|
|
socket = tcp
|
|
end
|
|
new(socket)
|
|
rescue Exception
|
|
begin
|
|
socket.close if socket != nil
|
|
rescue Exception
|
|
end
|
|
begin
|
|
tcp.close if tcp != nil && tcp != socket
|
|
rescue Exception
|
|
end
|
|
raise
|
|
end
|
|
end
|
|
|
|
def initialize(socket)
|
|
@socket = socket
|
|
@events = Queue.new
|
|
@outbound = Queue.new
|
|
@write_buffer = +"".b
|
|
@close_mutex = Mutex.new
|
|
@closed = false
|
|
@thread = Thread.new { run }
|
|
@thread.report_on_exception = false
|
|
end
|
|
|
|
def send_bytes(bytes)
|
|
return false if closed?
|
|
@outbound << bytes.to_s.b
|
|
true
|
|
end
|
|
|
|
def close
|
|
thread = nil
|
|
@close_mutex.synchronize do
|
|
return if @closed
|
|
@closed = true
|
|
thread = @thread
|
|
begin
|
|
@socket.close
|
|
rescue Exception
|
|
end
|
|
end
|
|
thread.join(1.0) if thread != nil && thread != Thread.current
|
|
nil
|
|
end
|
|
|
|
def closed?
|
|
@close_mutex.synchronize { @closed }
|
|
end
|
|
|
|
private
|
|
|
|
def run
|
|
@events << TransportEvent.new(:type => :connected)
|
|
loop do
|
|
break if closed?
|
|
fill_write_buffer
|
|
readable = [@socket]
|
|
writable = @write_buffer.empty? ? [] : [@socket]
|
|
selected = IO.select(readable, writable, nil, 0.1)
|
|
next if selected == nil
|
|
read_available if selected[0].include?(@socket)
|
|
write_available if selected[1].include?(@socket)
|
|
end
|
|
rescue EOFError
|
|
@events << TransportEvent.new(:type => :disconnected, :payload => _("The server closed the connection."))
|
|
rescue IOError, SystemCallError, OpenSSL::SSL::SSLError => error
|
|
unless closed?
|
|
@events << TransportEvent.new(
|
|
:type => :error,
|
|
:payload => _("Connection error: %{message}") % {:message => error.message}
|
|
)
|
|
end
|
|
rescue Exception => error
|
|
unless closed?
|
|
Log.error("Meow transport failure: #{error.class}: #{error.message}") if defined?(Log)
|
|
@events << TransportEvent.new(
|
|
:type => :error,
|
|
:payload => _("Unexpected connection error: %{message}") % {:message => error.message}
|
|
)
|
|
end
|
|
ensure
|
|
@close_mutex.synchronize { @closed = true }
|
|
begin
|
|
@socket.close
|
|
rescue Exception
|
|
end
|
|
end
|
|
|
|
def fill_write_buffer
|
|
loop do
|
|
@write_buffer << @outbound.pop(true)
|
|
end
|
|
rescue ThreadError
|
|
end
|
|
|
|
def read_available
|
|
data = @socket.read_nonblock(16_384)
|
|
raise EOFError if data == nil || data == ""
|
|
@events << TransportEvent.new(:type => :bytes, :payload => data.b)
|
|
rescue IO::WaitReadable, OpenSSL::SSL::SSLErrorWaitReadable
|
|
rescue IO::WaitWritable, OpenSSL::SSL::SSLErrorWaitWritable
|
|
end
|
|
|
|
def write_available
|
|
return if @write_buffer.empty?
|
|
count = @socket.write_nonblock(@write_buffer)
|
|
@write_buffer = @write_buffer.byteslice(count..-1).to_s.b
|
|
rescue IO::WaitWritable, OpenSSL::SSL::SSLErrorWaitWritable
|
|
rescue IO::WaitReadable, OpenSSL::SSL::SSLErrorWaitReadable
|
|
end
|
|
end
|
|
end
|
|
|