ruby - 如何在同一个 EventMachine react 器中运行 Net::SSH 和 AMQP?

标签 ruby amqp eventmachine gerrit net-ssh

一些背景:Gerrit exposes an event stream through SSH .这是一个可爱的技巧,但我需要将这些事件转换为 AMQP 消息。我试着用 ruby-amqp 做到这一点和 Net::SSH但是,好吧,AMQP 子组件似乎根本没有在运行。

我是 EventMachine 的新手。有人可以指出我做错了什么吗? “Multiple servers in a single EventMachine reactor”的答案似乎不适用。该程序也可以在 gist 中找到以便于访问,它是:

#!/usr/bin/env ruby                                                                                                                                            

require 'rubygems'
require 'optparse'
require 'net/ssh'
require 'json'
require 'yaml'
require 'amqp'
require 'logger'

trap(:INT) { puts; exit }

options = {
  :logs => 'kili.log',
  :amqp => {
    :host => 'localhost',
    :port => '5672',
  },
  :ssh => {
    :host => 'localhost',
    :port => '22',
    :user => 'nobody',
    :keys => '~/.ssh/id_rsa',
  }
}
optparse = OptionParser.new do|opts|
  opts.banner = "Usage: kili [options]"
  opts.on( '--amqp_host HOST', 'The AMQP host kili will connect to.') do |a|
    options[:amqp][:host] = a
  end
  opts.on( '--amqp_port PORT', 'The port for the AMQP host.') do |ap|
    options[:amqp][:port] = ap
  end
  opts.on( '--ssh_host HOST', 'The SSH host kili will connect to.') do |s|
    options[:ssh][:host] = s
  end
  opts.on( '--ssh_port PORT', 'The SSH port kili will connect on.') do |sp|
    options[:ssh][:port] = sp
  end
  opts.on( '--ssh_keys KEYS', 'Comma delimeted SSH keys for user.') do |sk|
    options[:ssh][:keys] = sk
  end
  opts.on( '--ssh_user USER', 'SSH user for host.') do |su|
    options[:ssh][:user] = su
  end
  opts.on( '-l', '--log LOG', 'The log location of Kili') do |log|
    options[:logs] = log
  end
  opts.on( '-h', '--help', 'Display this screen' ) do
    puts opts
    exit
  end
end


optparse.parse!
log = Logger.new(options[:logs])
log.level = Logger::INFO

amqp = options[:amqp]
sshd = options[:ssh]
queue= EM::Queue.new

EventMachine.run do

  AMQP.connect(:host => amqp[:host], :port => amqp[:port]) do |connection|
    log.info "Connected to AMQP at #{amqp[:host]}:#{amqp[:port]}"
    channel = AMQP::Channel.new(connection)
    exchange = channel.topic("traut", :auto_delete => true)

    queue.pop do |msg|
      log.info("Pulled #{msg} out of queue.")
      exchange.publish(msg[:data], :routing_key => msg[:route]) do
        log.info("On route #{msg[:route]} published:\n#{msg[:data]}")
      end
    end
  end


  Net::SSH.start(sshd[:host], sshd[:user],
    :port => sshd[:port], :keys => sshd[:keys].split(',')) do |ssh|
    log.info "SSH connection to #{sshd[:host]}:#{sshd[:port]} as #{sshd[:user]} made."

    channel = ssh.open_channel do |ch|
      ch.exec "gerrit stream-events" do |ch, success|
        abort "could not stream gerrit events" unless success

        # "on_data" is called when the process writes something to                                                                                             
        # stdout                                                                                                                                               
        ch.on_data do |c, data|
          json = JSON.parse(data)
          if json['type'] == 'change-merged'
            project = json['change']['project']
            route = "com.carepilot.event.code.review.#{project}"
            msg = {:data => data, :route => route}
            queue.push(msg)
            log.info("Pushed #{msg} into queue.")
          else
            log.info("Ignoring event of type #{json['type']}")
          end
        end


    # "on_extended_data" is called when the process writes                                                                                                 
    # something to stderr                                                                                                                                  
    ch.on_extended_data do |c, type, data|
          log.error(data)
    end

    ch.on_close { log.info('Connection closed') }
      end
    end  
  end  

end

最佳答案

Net::SSH 不是异步的,因此您的 EventMachine.run() 永远不会到达 block 的末尾,因此永远不会恢复 react 器线程。这会导致 AMQP 代码永远不会启动。我建议在另一个线程中运行您的 SSH 代码。

关于ruby - 如何在同一个 EventMachine react 器中运行 Net::SSH 和 AMQP?,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/7590758/

相关文章:

ruby - 修改Ruby OptionParser错误消息

ruby-on-rails - 如何使用 strong_parameters 允许除 user_id 之外的所有属性?

nservicebus - NServiceBus相对于普通RabbitMQ的特定优势

jms - Azure 服务总线 : Amqp Idle Timeout condition = amqp:link:detach-forced

ruby - 循环获取 Eventmachine/faye websocket

Ruby 线程和 Websocket

ruby-on-rails - Ruby - 向 Websocket 发送消息

ruby - 使用 Ruby 和 SCP/SSH,如何在上传副本之前确定文件是否存在

ruby-on-rails - 设计:用户未登录

rest - 消息传递 - 所有属性或只是一个 id 指针