Class: Norikra::Engine

Inherits:
Object
  • Object
show all
Defined in:
lib/norikra/engine.rb

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(output_pool, typedef_manager, opts = {}) ⇒ Engine

Returns a new instance of Engine.



24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
# File 'lib/norikra/engine.rb', line 24

def initialize(output_pool, typedef_manager, opts={})
  @statistics = {
    started: Time.now,
    events: { input: 0, processed: 0, output: 0, },
  }

  @output_pool = output_pool
  @typedef_manager = typedef_manager

  conf = configuration(opts)
  @service = com.espertech.esper.client.EPServiceProviderManager.getDefaultProvider(conf)
  @config = @service.getEPAdministrator.getConfiguration

  @mutex = Mutex.new

  # fieldsets already registered into @runtime
  @registered_fieldsets = {} # {target => {fieldset_summary => Fieldset}

  @targets = []
  @queries = []
  @suspended_queries = []

  @waiting_queries = []

  @listeners = {} # Listener.label => Listener
  @running_listeners = {} # query_name => listener
end

Instance Attribute Details

#output_poolObject (readonly)

Returns the value of attribute output_pool.



22
23
24
# File 'lib/norikra/engine.rb', line 22

def output_pool
  @output_pool
end

#queriesObject (readonly)

Returns the value of attribute queries.



22
23
24
# File 'lib/norikra/engine.rb', line 22

def queries
  @queries
end

#suspended_queriesObject (readonly)

Returns the value of attribute suspended_queries.



22
23
24
# File 'lib/norikra/engine.rb', line 22

def suspended_queries
  @suspended_queries
end

#targetsObject (readonly)

Returns the value of attribute targets.



22
23
24
# File 'lib/norikra/engine.rb', line 22

def targets
  @targets
end

#typedef_managerObject (readonly)

Returns the value of attribute typedef_manager.



22
23
24
# File 'lib/norikra/engine.rb', line 22

def typedef_manager
  @typedef_manager
end

Instance Method Details

#camelize(sym) ⇒ Object



115
116
117
# File 'lib/norikra/engine.rb', line 115

def camelize(sym)
  sym.to_s.split(/_/).map(&:capitalize).join
end

#close(target_name) ⇒ Object



167
168
169
170
171
172
173
174
175
176
# File 'lib/norikra/engine.rb', line 167

def close(target_name)
  info "closing target", target: target_name
  targets = @targets.select{|t| t.name == target_name}
  return false if targets.size != 1
  target = targets.first
  @queries.select{|q| q.targets.include?(target.name)}.each do |query|
    deregister_query(query)
  end
  close_target(target)
end

#configuration(opts) ⇒ Object



119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
# File 'lib/norikra/engine.rb', line 119

def configuration(opts)
  conf = com.espertech.esper.client.Configuration.new
  defaults = conf.getEngineDefaults

  if opts[:thread]
    t = opts[:thread]
    threading = defaults.getThreading
    [:inbound, :outbound, :route_exec, :timer_exec].each do |sym|
      next unless t[sym] && t[sym][:threads] && t[sym][:threads] > 0

      threads = t[sym][:threads].to_i
      capacity = t[sym][:capacity].to_i
      info "Engine #{sym} thread pool enabling", threads: threads, capacity: (capacity == 0 ? 'default' : capacity)

      cam = camelize(sym)
      threading.send("setThreadPool#{cam}".to_sym, true)
      threading.send("setThreadPool#{cam}NumThreads".to_sym, threads)
      if t[sym][:capacity] && t[sym][:capacity] > 0
        threading.send("setThreadPool#{cam}Capacity".to_sym, capacity)
      end
    end
  end

  conf
end

#deregister(query_name) ⇒ Object



206
207
208
209
210
211
212
213
214
215
216
217
218
219
# File 'lib/norikra/engine.rb', line 206

def deregister(query_name)
  info "de-registering query", name: query_name
  queries = @queries.select{|q| q.name == query_name }
  s_queries = @suspended_queries.select{|q| q.name == query_name }

  if queries.size == 1
    deregister_query(queries.first)
  elsif s_queries.size == 1
    @suspended_queries.delete(s_queries.first)
    true
  else
    nil # just ignore for 'not found'
  end
end

#event_filter(event) ⇒ Object



248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
# File 'lib/norikra/engine.rb', line 248

def event_filter(event)
  unless event.is_a?(Hash)
    error "Invalid input event: Non-Hash (JSON) object: #{event.class}"
    return nil
  end
  event.keys.each do |k|
    if ! k.is_a?(String)
      warn "Invalid key in event: Non-String field key: #{k.class}"
      event.delete(k)
    elsif k =~ /^\d/
      warn "Invalid key in event: Starting with numeric char: #{k.to_s}"
      event.delete(k)
    end
  end
  event
end

#gc_statisticsObject



100
101
102
103
104
105
106
107
108
109
110
111
112
113
# File 'lib/norikra/engine.rb', line 100

def gc_statistics
  gcBeans = Java::JavaLangManagement::ManagementFactory.getGarbageCollectorMXBeans()

  gc = {}
  gcBeans.each do |bean|
    name = bean.getName()
    gc[name] = {
      total_count: bean.getCollectionCount(),
      total_time: bean.getCollectionTime(),
    }
  end

  gc
end

#load(type, plugin_klass) ⇒ Object



326
327
328
329
330
331
332
333
# File 'lib/norikra/engine.rb', line 326

def load(type, plugin_klass)
  case type
  when :udf then load_udf(plugin_klass)
  when :listener then load_listener(plugin_klass)
  else
    raise "BUG: unknown plugin type: #{type}"
  end
end

#memory_statisticsObject



76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
# File 'lib/norikra/engine.rb', line 76

def memory_statistics
  mb = 1024 * 1024

  memoryBean = Java::JavaLangManagement::ManagementFactory.getMemoryMXBean()

  usage = memoryBean.getHeapMemoryUsage()
  total = usage.getMax() / mb
  committed = usage.getCommitted() / mb
  committed_percent = (committed.to_f / total * 1000).floor / 10.0
  used = usage.getUsed() / mb
  used_percent = (used.to_f / total * 1000).floor / 10.0
  heap = { max: total, committed: committed, committed_percent: committed_percent, used: used, used_percent: used_percent }

  usage = memoryBean.getNonHeapMemoryUsage()
  total = usage.getMax() / mb
  committed = usage.getCommitted() / mb
  committed_percent = (committed.to_f / total * 1000).floor / 10.0
  used = usage.getUsed() / mb
  used_percent = (used.to_f / total * 1000).floor / 10.0
  non_heap = { max: total, committed: committed, committed_percent: committed_percent, used: used, used_percent: used_percent }

  { heap: heap, nonheap: non_heap }
end

#modify(target_name, auto_field) ⇒ Object



178
179
180
181
182
183
184
185
186
# File 'lib/norikra/engine.rb', line 178

def modify(target_name, auto_field)
  info "modify target", target: target_name, auto_field: auto_field
  targets = @targets.select{|t| t.name == target_name}
  if targets.size != 1
    raise Norikra::ArgumentError, "target name '#{target_name}' not found"
  end
  target = targets.first
  target.auto_field = auto_field
end

#open(target_name, fields = nil, auto_field = true) ⇒ Object



157
158
159
160
161
162
163
164
165
# File 'lib/norikra/engine.rb', line 157

def open(target_name, fields=nil, auto_field=true)
  # fields nil || [] => lazy
  # fields {'fieldname' => 'type'} : type 'string', 'boolean', 'int', 'long', 'float', 'double'
  info "opening target", target: target_name, fields: fields, auto_field: auto_field
  raise Norikra::ArgumentError, "invalid target name" unless Norikra::Target.valid?(target_name)
  target = Norikra::Target.new(target_name, fields, auto_field)
  return false if @targets.include?(target)
  open_target(target)
end

#register(query) ⇒ Object



192
193
194
195
196
197
198
199
200
201
202
203
204
# File 'lib/norikra/engine.rb', line 192

def register(query)
  info "registering query", name: query.name, targets: query.targets, expression: query.expression
  raise Norikra::ClientError, "query name '#{query.name}' already exists" if @queries.select{|q| q.name == query.name }.size > 0
  raise Norikra::ClientError, "query name '#{query.name}' already exists in suspended" if @suspended_queries.select{|q| q.name == query.name }.size > 0
  if reason = query.invalid?
    raise Norikra::ClientError, "invalid query '#{query.name}': #{reason}"
  end

  query.targets.each do |target_name|
    open(target_name) unless @targets.any?{|t| t.name == target_name}
  end
  register_query(query)
end

#reserve(target_name, field, type) ⇒ Object



188
189
190
# File 'lib/norikra/engine.rb', line 188

def reserve(target_name, field, type)
  @typedef_manager.reserve(target_name, field, type)
end

#resume(query_name) ⇒ Object



233
234
235
236
237
238
239
240
241
242
243
244
245
246
# File 'lib/norikra/engine.rb', line 233

def resume(query_name)
  info "resuming query", name: query_name
  queries = @suspended_queries.select{|q| q.name == query_name }
  return nil unless queries.size == 1 # just ignore

  suspended_query = queries.first
  query = suspended_query.create # suspended query -> query object

  query.targets.each do |target_name|
    open(target_name) unless @targets.any?{|t| t.name == target_name}
  end
  register_query(query)
  remove_suspended_query(suspended_query)
end

#send(target_name, events) ⇒ Object



265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
# File 'lib/norikra/engine.rb', line 265

def send(target_name, events)
  trace "send messages", target: target_name, events: events

  @statistics[:events][:input] += events.size

  unless @targets.any?{|t| t.name == target_name} # discard events for target not registered
    trace "messages skipped for non-opened target", target: target_name
    return
  end
  return if events.size < 1

  target = @targets.select{|t| t.name == target_name}.first

  if @typedef_manager.lazy?(target.name)
    info "opening lazy target", target: target

    first_event = event_filter(events.first)
    if first_event.nil? # non-hash object
      raise Norikra::ClientError, "Input data must be JSON object"
    end
    debug "generating base fieldset from event", target: target.name, event: first_event
    base_fieldset = @typedef_manager.generate_base_fieldset(target.name, first_event)

    debug "registering base fieldset", target: target.name, base: base_fieldset
    register_base_fieldset(target.name, base_fieldset)

    info "target successfully opened with fieldset", target: target, base: base_fieldset
  end

  registered_data_fieldset = @registered_fieldsets[target_name][:data]

  strict_refer = (not target.auto_field?)

  events.each do |input_event|
    event = event_filter(input_event)
    next if event.nil? # non-hash object

    fieldset = @typedef_manager.refer(target_name, event, strict_refer)

    unless registered_data_fieldset[fieldset.summary]
      # register waiting queries including this fieldset, and this fieldset itself
      debug "registering unknown fieldset", target: target_name, fieldset: fieldset
      register_fieldset(target_name, fieldset)
      debug "successfully registered"

      # fieldset should be refined, when waiting_queries rewrite inheritance structure and data fieldset be renewed.
      fieldset = @typedef_manager.refer(target_name, event, strict_refer)
      debug "re-referred data fieldset", target: target_name, fieldset: fieldset
    end

    trace "calling sendEvent with bound fieldset (w/ valid event_type_name)", target: target_name, event: event
    trace("This is assert for valid event_type_name"){ { event_type_name: fieldset.event_type_name } }
    formed = fieldset.format(event)
    trace "sendEvent", data: formed
    @runtime.sendEvent(formed.to_java, fieldset.event_type_name)
  end
  target.update!
  @statistics[:events][:processed] += events.size
  nil
end

#startObject



145
146
147
148
149
# File 'lib/norikra/engine.rb', line 145

def start
  debug "norikra engine starting: creating esper runtime"
  @runtime = @service.getEPRuntime
  debug "norikra engine started"
end

#statisticsObject



52
53
54
55
56
57
58
59
60
61
62
63
64
65
# File 'lib/norikra/engine.rb', line 52

def statistics
  s = @statistics
  {
    started: s[:started].rfc2822,
    uptime: self.uptime,
    memory: self.memory_statistics,
    garbage_collector: self.gc_statistics,
    input_events: s[:events][:input],
    processed_events: s[:events][:processed],
    output_events: s[:events][:output],
    queries: @queries.size,
    targets: @targets.size,
  }
end

#stopObject



151
152
153
154
155
# File 'lib/norikra/engine.rb', line 151

def stop
  debug "stopping norikra engine: stop all statements on esper"
  @service.getEPAdministrator.stopAllStatements
  debug "norikra engine stopped"
end

#suspend(query_name) ⇒ Object



221
222
223
224
225
226
227
228
229
230
231
# File 'lib/norikra/engine.rb', line 221

def suspend(query_name)
  info "suspending query", name: query_name
  queries = @queries.select{|q| q.name == query_name }
  return nil unless queries.size == 1 # just ignore for 'not found'

  suspending_query = queries.first
  suspended_query = Norikra::SuspendedQuery.new(suspending_query)

  deregister_query(suspending_query)
  add_suspended_query(suspended_query)
end

#uptimeObject



67
68
69
70
71
72
73
74
# File 'lib/norikra/engine.rb', line 67

def uptime
  # up 239 days, 20:40
  seconds = (Time.now - @statistics[:started]).to_i
  days = seconds / (24*60*60)
  hours = (seconds - days * (24*60*60)) / (60*60)
  minutes = (seconds - days * (24*60*60) - hours * (60*60)) / 60
  "#{days} days, #{sprintf("%02d", hours)}:#{sprintf("%02d", minutes)}"
end