Skip to content
This repository was archived by the owner on Jun 5, 2025. It is now read-only.
Merged
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
54 changes: 33 additions & 21 deletions lib/fluent/plugin/in_http_pull.rb
Original file line number Diff line number Diff line change
Expand Up @@ -123,6 +123,7 @@ def configure(conf)
end.to_h)

@http_method = :head if @status_only
@cookies = {}
end

def start
Expand All @@ -136,28 +137,30 @@ def on_timer
record = nil
record_time = Engine.now
emit_stream = Fluent::MultiEventStream.new
site = RestClient::Resource.new(@url, request_options)

base_site = RestClient::Resource.new(@url, request_options)
login_site = @login_path ? base_site[@login_path] : base_site
query_site = @path ? base_site[@path] : base_site
begin
get_session_cookie(login_site)
rescue StandardError => err
record = log_error(login_site.url, err)
emit_stream.add(record_time, record)
router.emit_stream(@tag, emit_stream)
return
end
begin
cookies = get_session_cookie(site)

site = site[@path] if @path
if @payload
res = site.method(@http_method).call(@payload.to_json, :cookies=>cookies)
res = query_site.method(@http_method).call(@payload.to_json, :cookies=>@cookies)
else
res = site.method(@http_method).call(:cookie=>cookies)
res = query_site.method(@http_method).call(:cookie=>@cookies)
end

record, body = get_record(res, site.url)
record, body = get_record(res, query_site.url)
process_events(record, body, emit_stream)
rescue StandardError => err
record = { "url" => site.url, "error" => err.message }
if err.respond_to? :http_code
record["status"] = err.http_code || 0
else
record["status"] = 0
end
log.error(record)
record = log_error(query_site.url, err)
# reset session cookie if it is expired
@cookies = {} if not @cookies.empty? and record["status"] == 401
emit_stream.add(record_time, record)
end

Expand All @@ -169,6 +172,17 @@ def shutdown
end

private
def log_error(url, err)
record = { "url" => url, "error" => err.message }
if err.respond_to? :http_code
record["status"] = err.http_code || 0
else
record["status"] = 0
end
log.error(record)
return record
end

def request_options
options = { timeout: @timeout, headers: @_request_headers }

Expand All @@ -186,17 +200,15 @@ def request_options
end

def get_session_cookie(resource)
cookies = {}
return cookies unless @login_path and @login_payload

login_response = resource[@login_path].post(@login_payload.to_json)
@cookies = {} if @cookies.nil?
return unless @login_path and @login_payload and @cookies.empty?
login_response = resource.post(@login_payload.to_json)
if login_response.code != 200
raise RestClient::ExceptionWithResponse.new(nil, login_response.code)
else
cookies = login_response.cookie_jar
@cookies = login_response.cookie_jar
end

return cookies
end

def get_record(response, url)
Expand Down