Class: Zillabyte::Command::Flows
- Defined in:
- lib/zillabyte/cli/flows.rb
Overview
manage custom flows
Instance Attribute Summary
Attributes inherited from Base
Class Method Summary collapse
Instance Method Summary collapse
- #handshake(write_stream, read_stream, node, color) ⇒ Object
-
#index ⇒ Object
flows.
-
#info ⇒ Object
flows:info [DIR].
-
#init ⇒ Object
flows:init [LANG] [DIR].
-
#kill ⇒ Object
flows:kill ID.
-
#list ⇒ Object
flows.
-
#live_run ⇒ Object
flows:live_run [DIR].
-
#prep ⇒ Object
flows:prep [DIR].
-
#pull ⇒ Object
flows:pull ID DIR.
-
#push ⇒ Object
flows:push [DIR].
- #read_message(read_stream, color) ⇒ Object
-
#status ⇒ Object
flows:status ID.
-
#test ⇒ Object
flows:test [TEST_DATASET_ID].
- #write_message(write_stream, msg) ⇒ Object
Methods inherited from Base
alias_command, extract_banner, extract_description, extract_help, extract_help_from_caller, extract_options, extract_summary, inherited, #initialize, method_added, namespace
Methods included from Helpers
#ask, #display, #error, #format_with_bang, #longest, #read_multiline, #with_tty
Constructor Details
This class inherits a constructor from Zillabyte::Command::Base
Class Method Details
.get_info(cmd, info_file, options = {}) ⇒ Object
614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 |
# File 'lib/zillabyte/cli/flows.rb', line 614 def self.get_info(cmd, info_file, = {}) response = `#{cmd}` if($?.exitstatus == 1) File.delete("#{info_file}") if File.exists?(info_file) puts "error: #{response}" if [:test] Process.exit 1 end flow_info = {} File.open("#{info_file}", "r+").each do |line| line = JSON.parse(line) if(line["type"]) flow_info["nodes"] << line else flow_info = line flow_info["nodes"] = [] end end File.delete("#{info_file}") flow_info = flow_info.to_json if(File.exists?("info_to_java.in")) java_pipe = open("info_to_java.in","w+") java_pipe.puts(flow_info+"\n") java_pipe.flush java_pipe.close() end flow_info end |
Instance Method Details
#handshake(write_stream, read_stream, node, color) ⇒ Object
247 248 249 250 251 252 253 254 255 |
# File 'lib/zillabyte/cli/flows.rb', line 247 def handshake(write_stream, read_stream, node, color) begin write_stream, "{\"pidDir\": \"/tmp\"}\n" read_stream, color # Read to "end\n" rescue Exception => e puts "Error handshaking node: #{node}" raise e end end |
#index ⇒ Object
flows
17 18 19 |
# File 'lib/zillabyte/cli/flows.rb', line 17 def index self.list end |
#info ⇒ Object
flows:info [DIR]
outputs the info for the flow in the dir.
--pretty # Pretty prints the info output
600 601 602 603 604 605 606 607 608 609 610 611 |
# File 'lib/zillabyte/cli/flows.rb', line 600 def info info_file = SecureRandom.uuid cmd = command("--info --file #{info_file}") flow_info = Zillabyte::Command::Flows.get_info(cmd, info_file) if [:pretty] puts JSON.pretty_generate(JSON.parse(flow_info)) else puts flow_info end exit end |
#init ⇒ Object
flows:init [LANG] [DIR]
initializes a new executable in DIR defaults to a ruby executable for the current directory
Examples:
$ zillabyte flows:init python contact_extractor
186 187 188 189 190 191 192 193 194 195 196 197 198 199 |
# File 'lib/zillabyte/cli/flows.rb', line 186 def init lang = [:lang] || shift_argument || "ruby" dir = [:dir] || shift_argument || Dir.pwd languages = ["ruby","python", "js"] error "Unsupported language #{lang}. We only support #{languages.join(', ')}." if not languages.include? lang display "initializing empty #{lang} flow in #{dir}" FileUtils.cp_r( File.("../templates/#{lang}", __FILE__) + "/." , dir ) end |
#kill ⇒ Object
flows:kill ID
kills the given flow
--config CONFIG_FILE # use the given config file
555 556 557 558 559 560 561 |
# File 'lib/zillabyte/cli/flows.rb', line 555 def kill id = [:id] || shift_argument api.flows.kill(id) display "flow #{id} killed" end |
#list ⇒ Object
flows
list custom flows
26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 |
# File 'lib/zillabyte/cli/flows.rb', line 26 def list headings = ["id", "name", "state"] rows = api.flow.list.map do |row| if headings.size == 0 headings = row.keys headings.delete("rel_dir") end v = row.values_at *headings v end display "flows:\n" + Terminal::Table.new(:headings => headings, :rows => rows).to_s display "Total number of flows: "+rows.length.to_s end |
#live_run ⇒ Object
flows:live_run [DIR]
runs a local flow with live data
--config CONFIG_FILE # use the given config file
574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 |
# File 'lib/zillabyte/cli/flows.rb', line 574 def live_run name = [:name] || shift_argument thread_id = [:thread] || shift_argument || "" dir = [:directory] || shift_argument || Dir.pwd = Zillabyte::CLI::Config.get_config_info(dir) if .nil? throw "could not find meta information for: #{dir}" end if(thread_id == "") exec(command("--execute_live --name #{name.to_s}")) else exec(command("--execute_live --name #{name.to_s} --pipe #{thread_id}")) end end |
#prep ⇒ Object
flows:prep [DIR]
prepares a flow for execution
141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 |
# File 'lib/zillabyte/cli/flows.rb', line 141 def prep dir = [:directory] || shift_argument || Dir.pwd = Zillabyte::CLI::Config.get_config_info(dir) case ["language"] when "ruby" # Execute in the bundler context full_script = File.join(dir, ["script"]) cmd = "cd \"#{['home_dir']}\"; unset BUNDLE_GEMFILE; bundle install" exec(cmd) when "python" vDir = "#{['home_dir']}/vEnv" lock_file = ['home_dir']+"/zillabyte_thread_lock_file" if File.exists?(lock_file) sleep(1) while File.exists?(lock_file) else begin cmd = "touch #{lock_file}; virtualenv --clear --system-site-packages #{vDir}; PYTHONPATH=~/zb1/multilang/python/Zillabyte #{vDir}/bin/pip install -r #{['home_dir']}/requirements.txt" system cmd, :out => :out ensure File.delete(lock_file) end end when "js" end end |
#pull ⇒ Object
flows:pull ID DIR
pulls a flow source to a directory. The target directory must be empty
--force # pulls even if the directory exists
Examples:
$ zillabyte flows:pull .
104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 |
# File 'lib/zillabyte/cli/flows.rb', line 104 def pull id = [:id] || shift_argument dir = [:directory] || shift_argument error "no id given" if id.nil? error "no directory given" if dir.nil? # Create if not exists.. if File.exists?(dir) if Dir.entries(dir).size != 2 and [:force].nil? error "target directory not empty. use --force to override" end else FileUtils.mkdir_p(dir) end res = api.flows.pull_to_directory id, dir, progress if res['error'] display "error: #{res['error_message']}" else display "flow ##{res['id']} pulled to #{dir}" end end |
#push ⇒ Object
flows:push [DIR]
uploads a flow
--config CONFIG_FILE # use the given config file
Examples:
$ zillabyte flows:push .
77 78 79 80 81 82 83 84 85 86 87 88 |
# File 'lib/zillabyte/cli/flows.rb', line 77 def push dir = [:directory] || shift_argument || Dir.pwd res = api.flows.push_directory dir, progress, if res['error'] display "error: #{res['error_message']}" else display "flow ##{res['id']} #{res['action']}" end end |
#read_message(read_stream, color) ⇒ Object
220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 |
# File 'lib/zillabyte/cli/flows.rb', line 220 def (read_stream, color) msg = nil read_stream.each do |line| line.strip! if(line == "end") return msg end begin hash = JSON.parse(line) if(hash["command"] == "done") msg = "done" else msg = hash.to_json end rescue next end end msg end |
#status ⇒ Object
flows:status ID
gets the flows status
--config CONFIG_FILE # use the given config file
536 537 538 539 540 541 542 543 544 545 546 547 548 549 |
# File 'lib/zillabyte/cli/flows.rb', line 536 def status headings = ["name", "implementable", "implemented"] rows = api.flows.list.map do |row| if headings.size == 0 headings = row.keys headings.delete("rel_dir") end v = row.values_at *headings v end display "flow status:\n" + Terminal::Table.new(:headings => headings, :rows => rows).to_s display "Total number of flows in queue: "+rows.length.to_s end |
#test ⇒ Object
flows:test [TEST_DATASET_ID]
tests a local flow with sample data
--config CONFIG_FILE # use the given config file --output OUTPUT_FILE # writes sink output to a file --wait MAX # max time to spend on each operation (default 10 seconds) --batches BATCHES # number of batches to emit (default 1)
214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 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 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 |
# File 'lib/zillabyte/cli/flows.rb', line 214 def test output = [:output] max_seconds = ([:wait] || "30").to_i batches = ([:batches] || "1").to_i def (read_stream, color) msg = nil read_stream.each do |line| line.strip! if(line == "end") return msg end begin hash = JSON.parse(line) if(hash["command"] == "done") msg = "done" else msg = hash.to_json end rescue next end end msg end def (write_stream, msg) write_stream.write msg.strip + "\n" write_stream.write "end\n" write_stream.flush end def handshake(write_stream, read_stream, node, color) begin write_stream, "{\"pidDir\": \"/tmp\"}\n" read_stream, color # Read to "end\n" rescue Exception => e puts "Error handshaking node: #{node}" raise e end end # INIT test_data = [:test_data] || shift_argument dir = [:dir] || Dir.pwd = Zillabyte::API::Flows.(dir, self, {:test => true}) if .nil? error "this is not a valid zillabyte flow directory" exit end # Show the user what we know about their flow... display "inferring your flow details..." colors = {} describe_flow(, colors) # Extract the flow's information.. nodes = ["nodes"] write_to_next_each = [] write_queue = [] = {} default_stream = "_default" split_branches = false # Iterate all nodes sequentially and invoke them in separate processes... nodes.each do |node| # Init type = node["type"] name = node["name"] color = colors[name] || :default op_display = lambda do |msg, override_color = nil| display "#{name} - #{msg}".colorize(override_color || color) end # A Spout? if type == "spout" # A spout from relation? if node['matches'] or node["relation"] matches = node['matches'] || (node["relation"]["query"]) op_display.call "Grabbing remote data" res = api.query.agnostic(matches)["rows"] if(res.nil? or res.length == 0) raise NameError, "Could not find data that matches your 'matches' clause" end res.each do |tuple| values = {} = {} tuple.each do |k, v| if(k == "id") next elsif(k == "confidence" or k == "since" or k == "source") [k] = v else values[k] = v end end read_msg = {"tuple" => values, "meta" => }.to_json op_display.call "emit tuple: #{values} #{}" [default_stream] ||= [] [default_stream] << read_msg end # Done processing... next else # A regular spout.. [default_stream] ||= [] batches.times do |i| [default_stream] << "{\"command\": \"next\"}\n" end end end # A Sink? if type == "sink" if split_branches || .size > 1 sink_stream = node["consumes"] = [sink_stream] || [] else = .values.first || [] end table = Terminal::Table.new :title => name csv_str = CSV.generate do |csv| header_written = false; .each do |msg| obj = JSON.parse(msg) if obj['tuple'] if header_written == false keys = [obj['tuple'].keys, obj['meta'].keys].flatten csv << keys table << keys table << :separator header_written = true end vals = [obj['tuple'].values, obj['meta'].values].flatten csv << vals table << vals end end end display table.to_s.colorize(color) if output filename = "#{output}.csv" f = File.open(filename, "w") f.write(csv_str) f.close() op_display.call "output written to #{filename}" end next end cmd = command("--execute_live --name #{name}") begin # Start the operation... op_display.call "beginning #{type} #{name}" Open3.popen3(cmd) do |stdin, stdout, stderr, wait_thread| begin # Init handshake stdin, stdout, node, color write_queue = [] read_queue = [] # Get the incoming stream if !split_branches && .size == 1 # i.e. the number of streams we're dealing with right now # Assume default stream stream_name = .keys.first write_queue = .values.first.clone .delete(stream_name) else # Multiple streams... split_branches = true; if node['consumes'].nil? error "The node #{name} must declare which stream it 'consumes'" end stream_name = node["consumes"] write_queue = [stream_name].clone() .delete(stream_name) end # Start writing the messages... stuff_to_read = false writing_thread = Thread.start do until(write_queue.empty?) # Make sure we're not reading anything... while(stuff_to_read) # TODO: semaphores sleep 0.5 # spin wait end # Get next mesage write_msg = write_queue.shift # Make it human-understable write_json = JSON.parse(write_msg) if write_json['tuple'] op_display.call "receiving: #{write_json['tuple']}" elsif write_json['command'] == 'next' op_display.call "starting next spout batch" else puts write_json end # Actually send it to the process begin stdin, write_msg stuff_to_read = true sleep 0.1 rescue Exception => e puts "Error running #{cmd}: #{e}" raise e end end end # Start reading messages... reading_thread = Thread.start do while(true) # Get next message read_msg = (stdout, color) if read_msg == "done" || read_msg.nil? stuff_to_read = false if write_queue.empty? break # exit while loop else sleep 0.5 # spin wait next end end stuff_to_read = true # Process message obj = JSON.parse(read_msg) if obj['tuple'] # Conver to a incoming tuple for the next operation next_msg = { :tuple => obj['tuple'], :meta => obj['meta'] } emit_stream = obj['stream'] || default_stream [emit_stream] ||= [] [emit_stream] << next_msg.to_json op_display.call "emitted: #{obj['tuple']} to #{emit_stream}" elsif obj['command'] == 'log' op_display.call "log: #{obj['msg']}" else error "unknown message: #{read_msg}" end end end # stderr thread stderr_thread = Thread.start do stderr.each do |line| op_display.call("stderr: #{line}", :red) end end begin killed = Timeout.timeout(max_seconds) do reading_thread.join() writing_thread.join() stderr_thread.kill() op_display.call "completed #{type} #{name}" end rescue Timeout::Error op_display.call "max time reached. preempting #{type} #{name}. set --wait to increase", :red reading_thread.kill() if reading_thread.alive? writing_thread.kill() if writing_thread.alive? stderr_thread.kill() if stderr_thread.alive? end rescue Errno::EIO puts "Errno:EIO error, but this probably just means " + "that the process has finished giving output" end end rescue PTY::ChildExited puts "The child process exited!" end end end |
#write_message(write_stream, msg) ⇒ Object
241 242 243 244 245 |
# File 'lib/zillabyte/cli/flows.rb', line 241 def (write_stream, msg) write_stream.write msg.strip + "\n" write_stream.write "end\n" write_stream.flush end |