Class: LogStash::Outputs::Tcp

Inherits:
Base
  • Object
show all
Defined in:
lib/logstash/outputs/tcp.rb

Overview

Write events over a TCP socket.

Each event json is separated by a newline.

Can either accept connections from clients or connect to a server, depending on ‘mode`.

Defined Under Namespace

Classes: Client

Instance Method Summary collapse

Instance Method Details

#BaseObject



231
232
233
234
235
236
237
238
239
240
# File 'lib/logstash/outputs/tcp.rb', line 231

def close
  @closed.make_true
  @server_socket.close rescue nil if @server_socket

  return unless @client_threads
  @client_threads.each do |thread|
    client = thread[:client]
    client.close rescue nil if client
  end
end

#log_error(msg, e, backtrace: @logger.info?, **details) ⇒ Object



248
249
250
251
252
# File 'lib/logstash/outputs/tcp.rb', line 248

def log_error(msg, e, backtrace: @logger.info?, **details)
  details = details.merge message: e.message, exception: e.class
  details[:backtrace] = e.backtrace if backtrace
  @logger.error(msg, details)
end

#log_warn(msg, e, backtrace: @logger.debug?, **details) ⇒ Object



242
243
244
245
246
# File 'lib/logstash/outputs/tcp.rb', line 242

def log_warn(msg, e, backtrace: @logger.debug?, **details)
  details = details.merge message: e.message, exception: e.class
  details[:backtrace] = e.backtrace if backtrace
  @logger.warn(msg, details)
end

#BaseObject



226
227
228
# File 'lib/logstash/outputs/tcp.rb', line 226

def receive(event)
  @codec.encode(event)
end

#BaseObject



141
142
143
144
145
146
147
148
149
150
151
152
153
# File 'lib/logstash/outputs/tcp.rb', line 141

def register
  require "socket"
  require "stud/try"
  @closed = Concurrent::AtomicBoolean.new(false)
  @thread_no = Concurrent::AtomicFixnum.new(0)
  setup_ssl if @ssl_enable

  if server?
    run_as_server
  else
    run_as_client
  end
end

#run_as_clientObject



198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
# File 'lib/logstash/outputs/tcp.rb', line 198

def run_as_client
  client_socket = nil
  @codec.on_event do |event, payload|
    begin
      client_socket = connect unless client_socket
      while payload && payload.bytesize > 0
        begin
          written_bytes_size = client_socket.write_nonblock(payload)
          payload = payload.byteslice(written_bytes_size..-1)
        rescue IO::WaitReadable
          IO.select([client_socket])
          retry
        rescue IO::WaitWritable
          IO.select(nil, [client_socket])
          retry
        end
      end
    rescue => e
      log_warn "client socket failed:", e, host: @host, port: @port, socket: (client_socket ? client_socket.to_s : nil)
      client_socket.close rescue nil
      client_socket = nil
      sleep @reconnect_interval
      retry
    end
  end
end

#run_as_serverObject



155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
# File 'lib/logstash/outputs/tcp.rb', line 155

def run_as_server
  @logger.info("Starting tcp output listener", :address => "#{@host}:#{@port}")
  begin
    @server_socket = TCPServer.new(@host, @port)
  rescue Errno::EADDRINUSE
    @logger.error("Could not start tcp server: Address in use", host: @host, port: @port)
    raise
  end
  if @ssl_enable
    @server_socket = OpenSSL::SSL::SSLServer.new(@server_socket, @ssl_context)
  end # @ssl_enable
  @client_threads = Concurrent::Array.new

  @accept_thread = Thread.new(@server_socket) do |server_socket|
    LogStash::Util.set_thread_name("[#{pipeline_id}]|output|tcp|server_accept")
    loop do
      break if @closed.value
      client_socket = server_socket.accept_nonblock exception: false
      if client_socket == :wait_readable
        IO.select [ server_socket ]
        next
      end
      Thread.start(client_socket) do |client_socket|
        # monkeypatch a 'peer' method onto the socket.
        client_socket.extend(::LogStash::Util::SocketPeer)
        @logger.debug("accepted connection", client: client_socket.peer, server: "#{@host}:#{@port}")
        client = Client.new(client_socket, self)
        Thread.current[:client] = client
        LogStash::Util.set_thread_name("[#{pipeline_id}]|output|tcp|client_socket-#{@thread_no.increment}")
        @client_threads << Thread.current
        client.run unless @closed.value
      end
    end
  end

  @codec.on_event do |event, payload|
    @client_threads.select!(&:alive?)
    @client_threads.each do |client_thread|
      client_thread[:client].write(payload)
    end
  end
end