Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 0 additions & 28 deletions lib/fluent/plugin/in_stream.rb
Original file line number Diff line number Diff line change
Expand Up @@ -160,34 +160,6 @@ def on_close
end
end


# obsolete
# ForwardInput is backward compatible with TcpInput
#class TcpInput < StreamInput
# Plugin.register_input('tcp', self)
#
# config_param :port, :integer, :default => DEFAULT_LISTEN_PORT
# config_param :bind, :string, :default => '0.0.0.0'
#
# def configure(conf)
# super
# end
#
# def listen
# log.debug "listening fluent socket on #{@bind}:#{@port}"
# Coolio::TCPServer.new(@bind, @port, Handler, method(:on_message))
# end
#end
class TcpInput < ForwardInput
Plugin.register_input('tcp', self)

def initialize
super
$log.warn "'tcp' input is obsoleted and will be removed soon. Use 'forward' instead."
end
end


class UnixInput < StreamInput
Plugin.register_input('unix', self)

Expand Down
64 changes: 5 additions & 59 deletions lib/fluent/plugin/in_syslog.rb
Original file line number Diff line number Diff line change
Expand Up @@ -122,7 +122,7 @@ def run
end

protected
def receive_data_parser(data)
def receive_data_parser(data, addr)
m = SYSLOG_REGEXP.match(data)
unless m
log.warn "invalid syslog message: #{data.dump}"
Expand All @@ -144,7 +144,7 @@ def receive_data_parser(data)
log.error_backtrace
end

def receive_data(data)
def receive_data(data, addr)
m = SYSLOG_ALL_REGEXP.match(data)
unless m
log.warn "invalid syslog message", :data=>data
Expand Down Expand Up @@ -183,9 +183,10 @@ def listen(callback)
if @protocol_type == :udp
@usock = SocketUtil.create_udp_socket(@bind)
@usock.bind(@bind, @port)
UdpHandler.new(@usock, callback)
SocketUtil::UdpHandler.new(@usock, log, 2048, callback)
else
Coolio::TCPServer.new(@bind, @port, TcpHandler, log, callback)
# syslog family add "\n" to each message and this seems only way to split messages in tcp stream
Coolio::TCPServer.new(@bind, @port, SocketUtil::TcpHandler, log, "\n", callback)
end
end

Expand All @@ -199,60 +200,5 @@ def emit(pri, time, record)
rescue => e
log.error "syslog failed to emit", :error => e.to_s, :error_class => e.class.to_s, :tag => tag, :record => Yajl.dump(record)
end

class UdpHandler < Coolio::IO
def initialize(io, callback)
super(io)
@io = io
@callback = callback
end

def on_readable
msg, addr = @io.recvfrom_nonblock(2048)
#host = addr[3]
#port = addr[1]
#@callback.call(host, port, msg)
@callback.call(msg)
rescue
# TODO log?
end
end

class TcpHandler < Coolio::Socket
def initialize(io, log, on_message)
super(io)
if io.is_a?(TCPSocket)
opt = [1, @timeout.to_i].pack('I!I!') # { int l_onoff; int l_linger; }
io.setsockopt(Socket::SOL_SOCKET, Socket::SO_LINGER, opt)
end
@on_message = on_message
@log = log
@log.trace { "accepted fluent socket object_id=#{self.object_id}" }
@buffer = "".force_encoding('ASCII-8BIT')
end

def on_connect
end

def on_read(data)
@buffer << data
pos = 0

# syslog family add "\n" to each message and this seems only way to split messages in tcp stream
while i = @buffer.index("\n", pos)
msg = @buffer[pos..i]
@on_message.call(msg)
pos = i + 1
end
@buffer.slice!(0, pos) if pos > 0
rescue => e
@log.error "syslog error", :error => e, :error_class => e.class
close
end

def on_close
@log.trace { "closed fluent socket object_id=#{self.object_id}" }
end
end
end
end
15 changes: 15 additions & 0 deletions lib/fluent/plugin/in_tcp.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
require 'fluent/plugin/socket_util'

module Fluent
class TcpInput < SocketUtil::BaseInput
Plugin.register_input('tcp', self)

config_set_default :port, 5170
config_param :delimiter, :string, :default => "\n" # syslog family add "\n" to each message and this seems only way to split messages in tcp stream

def listen(callback)
log.debug "listening tcp socket on #{@bind}:#{@port}"
Coolio::TCPServer.new(@bind, @port, SocketUtil::TcpHandler, log, @delimiter, callback)
end
end
end
17 changes: 17 additions & 0 deletions lib/fluent/plugin/in_udp.rb
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
require 'fluent/plugin/socket_util'

module Fluent
class UdpInput < SocketUtil::BaseInput
Plugin.register_input('udp', self)

config_set_default :port, 5160
config_param :body_size_limit, :size, :default => 4096

def listen(callback)
log.debug "listening udp socket on #{@bind}:#{@port}"
@usock = SocketUtil.create_udp_socket(@bind)
@usock.bind(@bind, @port)
SocketUtil::UdpHandler.new(@usock, log, @body_size_limit, callback)
end
end
end
119 changes: 119 additions & 0 deletions lib/fluent/plugin/socket_util.rb
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
require 'cool.io'

module Fluent
module SocketUtil
def create_udp_socket(host)
Expand All @@ -10,5 +12,122 @@ def create_udp_socket(host)
end
end
module_function :create_udp_socket

class UdpHandler < Coolio::IO
def initialize(io, log, body_size_limit, callback)
super(io)
@io = io
@log = log
@body_size_limit = body_size_limit
@callback = callback
end

def on_readable
msg, addr = @io.recvfrom_nonblock(@body_size_limit)
msg.chomp!
@callback.call(msg, addr)
rescue => e
@log.error "unexpected error", :error => e, :error_class => e.class
end
end

class TcpHandler < Coolio::Socket
PEERADDR_FAILED = ["?", "?", "name resolusion failed", "?"]

def initialize(io, log, delimiter, callback)
super(io)
if io.is_a?(TCPSocket)
@addr = (io.peeraddr rescue PEERADDR_FAILED)

opt = [1, @timeout.to_i].pack('I!I!') # { int l_onoff; int l_linger; }
io.setsockopt(Socket::SOL_SOCKET, Socket::SO_LINGER, opt)
end
@delimiter = delimiter
@callback = callback
@log = log
@log.trace { "accepted fluent socket object_id=#{self.object_id}" }
@buffer = "".force_encoding('ASCII-8BIT')
end

def on_connect
end

def on_read(data)
@buffer << data
pos = 0

while i = @buffer.index(@delimiter, pos)
msg = @buffer[pos...i]
@callback.call(msg, @addr)
pos = i + @delimiter.length
end
@buffer.slice!(0, pos) if pos > 0
rescue => e
@log.error "unexpected error", :error => e, :error_class => e.class
close
end

def on_close
@log.trace { "closed fluent socket object_id=#{self.object_id}" }
end
end

class BaseInput < Fluent::Input
def initialize
super
require 'fluent/parser'
end

config_param :tag, :string
config_param :format, :string
config_param :port, :integer, :default => 5150
config_param :bind, :string, :default => '0.0.0.0'
config_param :source_host_key, :string, :default => nil

def configure(conf)
super

@parser = TextParser.new
@parser.configure(conf)
end

def start
@loop = Coolio::Loop.new
@handler = listen(method(:on_message))
@loop.attach(@handler)
@thread = Thread.new(&method(:run))
end

def shutdown
@loop.watchers.each { |w| w.detach }
@loop.stop
@handler.close
@thread.join
end

def run
@loop.run
rescue => e
log.error "unexpected error", :error => e, :error_class => e.class
log.error_backtrace
end

private

def on_message(msg, addr)
@parser.parse(msg) { |time, record|
unless time && record
log.warn "pattern not match: #{msg.inspect}"
return
end

record[@source_host_key] = addr[3] if @source_host_key
Engine.emit(@tag, time, record)
}
rescue => e
log.error msg.dump, :error => e, :error_class => e.class, :host => addr[3]
log.error_backtrace
end
end
end
end
24 changes: 0 additions & 24 deletions test/plugin/test_in_stream.rb
Original file line number Diff line number Diff line change
Expand Up @@ -105,30 +105,6 @@ def send_data(data)
end
end

class TcpInputTest < Test::Unit::TestCase
include StreamInputTest

PORT = unused_port
CONFIG = %[
port #{PORT}
bind 127.0.0.1
]

def create_driver(conf=CONFIG)
super(Fluent::TcpInput, conf)
end

def test_configure
d = create_driver
assert_equal PORT, d.instance.port
assert_equal '127.0.0.1', d.instance.bind
end

def connect
TCPSocket.new('127.0.0.1', PORT)
end
end

class UnixInputTest < Test::Unit::TestCase
include StreamInputTest

Expand Down
Loading