· 9 years ago · Nov 22, 2016, 10:40 AM
1# encoding: utf-8
2require "logstash/inputs/base"
3require "logstash/namespace"
4require "socket"
5
6
7# Read rows from an sqlite database.
8#
9# This is most useful in cases where you are logging directly to a table.
10# Any tables being watched must have an 'id' column that is monotonically
11# increasing.
12#
13# All tables are read by default except:
14# * ones matching 'sqlite_%' - these are internal/adminstrative tables for sqlite
15# * 'since_table' - this is used by this plugin to track state.
16#
17# ## Example
18#
19# % sqlite /tmp/example.db
20# sqlite> CREATE TABLE weblogs (
21# id INTEGER PRIMARY KEY AUTOINCREMENT,
22# ip STRING,
23# request STRING,
24# response INTEGER);
25# sqlite> INSERT INTO weblogs (ip, request, response)
26# VALUES ("1.2.3.4", "/index.html", 200);
27#
28# Then with this logstash config:
29#
30# input {
31# sqlite {
32# path => "/tmp/example.db"
33# type => weblogs
34# }
35# }
36# output {
37# stdout {
38# debug => true
39# }
40# }
41#
42# Sample output:
43#
44# {
45# "@source" => "sqlite://sadness/tmp/x.db",
46# "@tags" => [],
47# "@fields" => {
48# "ip" => "1.2.3.4",
49# "request" => "/index.html",
50# "response" => 200
51# },
52# "@timestamp" => "2013-05-29T06:16:30.850Z",
53# "@source_host" => "sadness",
54# "@source_path" => "/tmp/x.db",
55# "@message" => "",
56# "@type" => "foo"
57# }
58#
59class LogStash::Inputs::Sqlite < LogStash::Inputs::Base
60 config_name "sqliteonec"
61 milestone 1
62
63 # The path to the sqlite database file.
64 config :path, :validate => :string, :required => true
65
66 # Any tables to exclude by name.
67 # By default all tables are followed.
68 config :exclude_tables, :validate => :array, :default => []
69
70 # How many rows to fetch at a time from each SELECT call.
71 config :batch, :validate => :number, :default => 5
72
73 SINCE_TABLE = :since_table
74
75 public
76 def init_placeholder_table(db)
77 begin
78 db.create_table SINCE_TABLE do
79 String :table
80 Int :place
81 end
82 rescue
83 @logger.debug("since tables already exists")
84 end
85 end
86
87 public
88 def get_placeholder(db, table)
89 since = db[SINCE_TABLE]
90 x = since.where(:table => "#{table}")
91 if x[:place].nil?
92 init_placeholder(db, table)
93 return 0
94 else
95 @logger.debug("placeholder already exists, it is #{x[:place]}")
96 return x[:place][:place]
97 end
98 end
99
100 public
101 def init_placeholder(db, table)
102 @logger.debug("init placeholder for #{table}")
103 since = db[SINCE_TABLE]
104 since.insert(:table => table, :place => 0)
105 end
106
107 public
108 def update_placeholder(db, table, place)
109 @logger.debug("set placeholder to #{place}")
110 since = db[SINCE_TABLE]
111 since.where(:table => table).update(:place => place)
112 end
113
114 public
115 def get_all_tables(db)
116 return db["SELECT * FROM sqlite_master WHERE type = 'table' AND tbl_name != '#{SINCE_TABLE}' AND tbl_name NOT LIKE 'sqlite_%'"].map { |t| t[:name] }.select { |n| !@exclude_tables.include?(n) }
117 end
118
119 public
120 def get_n_rows_from_table(db, table, offset, limit)
121 dataset = db["SELECT * FROM #{table}"]
122 return db["SELECT T1.rowID as id, T1.severity, T1.date, T1.transactionStatus, T1.transactionDate, T1.transactionID, T1.userCode, T2.code, T2.name, T3.code, T3.name, T4.code, T4.name, T5.code, T5.name, T1.comment, T1.data, T1.dataPresentation, T1.workServerCode, T6.code, T6.name, T1.primaryPortCode, T1.secondaryPortCode FROM EventLog T1 LEFT OUTER JOIN UserCodes T2 ON T1.userCode = T2.code LEFT OUTER JOIN ComputerCodes T3 ON T1.computerCode = T3.code LEFT OUTER JOIN AppCodes T4 ON T1.appCode = T4.code LEFT OUTER JOIN EventCodes T5 ON T1.eventCode = T5.code LEFT OUTER JOIN WorkServerCodes T6 ON T1.workServerCode = T6.code WHERE (id > #{offset}) ORDER BY 'id' LIMIT #{limit}"].map { |row| row }
123 end
124
125 public
126 def register
127 require "sequel"
128 require "jdbc/sqlite3"
129 @host = Socket.gethostname
130 @logger.info("Registering sqlite input", :database => @path)
131 @db = Sequel.connect("jdbc:sqlite:#{@path}")
132 @tables = get_all_tables(@db)
133 @table_data = {}
134 @tables.each do |table|
135 init_placeholder_table(@db)
136 last_place = get_placeholder(@db, table)
137 @table_data[table] = { :name => table, :place => last_place }
138 end
139 end # def register
140
141 public
142 def run(queue)
143 sleep_min = 0.01
144 sleep_max = 5
145 sleeptime = sleep_min
146
147 begin
148 @logger.debug("Tailing sqlite db", :path => @path)
149 loop do
150 count = 0
151 @table_data.each do |k, table|
152 table_name = table[:name]
153 offset = table[:place]
154 @logger.debug("offset is #{offset}", :k => k, :table => table_name)
155 next if table_name != "EventLog"
156 rows = get_n_rows_from_table(@db, table_name, offset, @batch)
157 count += rows.count
158 rows.each do |row|
159 event = LogStash::Event.new("host" => @host, "db" => @db)
160 decorate(event)
161 # store each column as a field in the event.
162 row.each do |column, element|
163 next if column == :id
164 event[column.to_s] = element
165 end
166 queue << event
167 @table_data[k][:place] = row[:id]
168 end
169 # Store the last-seen row in the database
170 update_placeholder(@db, table_name, @table_data[k][:place])
171 end
172
173 if count == 0
174 # nothing found in that iteration
175 # sleep a bit
176 @logger.debug("No new rows. Sleeping.", :time => sleeptime)
177 sleeptime = [sleeptime * 2, sleep_max].min
178 sleep(sleeptime)
179 else
180 sleeptime = sleep_min
181 end
182 end # loop
183 end # begin/rescue
184 end #run
185
186end # class Logtstash::Inputs::EventLog