Module: Thrift::Processor

Included in:
Types::Plugin::Plugin::Processor, Types::Plugin::Plugin::Processor
Defined in:
lib/thrift/processor.rb

Defined Under Namespace

Classes: BaseProcessor, BaseProcessorFunction, BaseStreamProcessor, BidiStreamProcessor, BinaryProcessor, BinaryProcessorFunction, InboundStreamProcessor, OutboundStreamProcessor, UnaryProcessor, UnaryProcessorFunction, UnkwonFunctionProcessor

Instance Method Summary collapse

Instance Method Details

#build_processor(name, info) ⇒ Object



147
148
149
150
151
152
153
154
155
156
157
158
159
# File 'lib/thrift/processor.rb', line 147

def build_processor(name, info)
  if info[:result_klass].nil?
    UnaryProcessor
  elsif info[:stream_klass].nil? && info[:sink_klass].nil?
    BinaryProcessor
  elsif info[:sink_klass].nil?
    OutboundStreamProcessor
  elsif info[:stream_klass].nil?
    InboundStreamProcessor
  else
    BidiStreamProcessor
  end.new(name, info, @middleware, @handler)
end

#initialize(handler, middlewares = [], logger = nil) ⇒ Object



100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
# File 'lib/thrift/processor.rb', line 100

def initialize(handler, middlewares = [], logger=nil)
  @handler = handler
  if logger.nil?
    @logger = Logger.new(STDERR)
    @logger.level = Logger::WARN
  else
    @logger = logger
  end
  @middleware = Middleware.wrap(middlewares)

  @processors = if self.class.const_defined? :METHODS
                  self.class::METHODS.reduce({}) do |acc, (name, info)|
                    acc.merge(name => build_processor(name, info))
                  end
                else
                  {}
                end
end

#process(iprot, oprot) ⇒ Object



133
134
135
136
137
138
139
140
141
142
143
144
145
# File 'lib/thrift/processor.rb', line 133

def process(iprot, oprot)
  name, _type, seqid = iprot.read_message_begin

  mth = "process_#{name}"
  if respond_to?(mth)
    send(mth, seqid, iprot, oprot)
    return true
  end

  (
    @processors[name] || UnkwonFunctionProcessor.new(name)
  ).process(seqid, iprot, oprot)
end

#read_args(iprot, args_class) ⇒ Object



119
120
121
122
123
124
# File 'lib/thrift/processor.rb', line 119

def read_args(iprot, args_class)
  args = args_class.new
  args.read(iprot)
  iprot.read_message_end
  args
end

#write_result(result, oprot, name, seqid) ⇒ Object



126
127
128
129
130
131
# File 'lib/thrift/processor.rb', line 126

def write_result(result, oprot, name, seqid)
  oprot.write_message_begin(name, MessageTypes::REPLY, seqid)
  result.write(oprot)
  oprot.write_message_end
  oprot.trans.flush
end