· 8 years ago · Jun 21, 2018, 12:14 AM
1require 'eventmachine'
2require 'net/dns'
3require 'net/dns/resolver'
4
5module EM # :nodoc:
6 module Protocols
7
8 include Logger::Severity
9
10 class AsyncResolver < Net::DNS::Resolver
11
12 # Create a new resolver object.
13 def initialize(config = {})
14 # store outstanding requests
15 @outstanding = {}
16 super(config)
17 end
18
19 # Do an asynchronous DNS query. Returns an AsyncQuery object that implements Deferrable
20 # The callback will be passed a DNS::Packet with the returned results
21 def query_async(name,type=Net::DNS::A,cls=Net::DNS::IN)
22 # If the name doesn't contain any dots then append the default domain.
23 if name.class != IPAddr and name !~ /\./ and name !~ /:/ and @config[:defnames]
24 name += "." + @config[:domain]
25 end
26 @logger.debug "Query(#{name},#{Net::DNS::RR::Types.new(type)},#{Net::DNS::RR::Classes.new(cls)})"
27 AsyncQuery.new(send_async(name,type,cls))
28 end
29
30 # Do an asynchronous MX query. Returns a MXQuery object that implements Deferrable
31 # The callback will be passed an array of MX records with the returned results, sorted by preference
32 def mx_async(name,cls=Net::DNS::IN)
33 @logger.debug "Query(#{name},#{Net::DNS::MX},#{Net::DNS::RR::Classes.new(cls)})"
34 MXQuery.new(send_async(name, Net::DNS::MX, cls))
35 end
36
37 # Send a query to the nameservers and return a deferrable that will be called with the response packet
38 def send_async(argument, type = Net::DNS::A, cls = Net::DNS::IN)
39 if @config[:nameservers].size == 0
40 raise ResolverError, "No nameservers specified!"
41 end
42
43 method = :send_udp_async
44 packet = if argument.kind_of? Net::DNS::Packet
45 argument
46 else
47 make_query_packet(argument, type, cls)
48 end
49
50 # Store packet_data for performance improvements,
51 # so methods don't keep on calling Packet#data
52 packet_data = packet.data
53 packet_size = packet_data.size
54
55 # Choose whether use TCP or UDP
56 if packet_size > @config[:packet_size] # Must use TCP, either plain or raw
57 @logger.info "Sending #{packet_size} bytes using TCP"
58 method = :send_tcp_async
59 else # Packet size is inside the boundaries
60 if use_tcp? # User requested TCP
61 @logger.info "Sending #{packet_size} bytes using TCP"
62 method = :send_tcp_async
63 else # Finally use UDP
64 @logger.info "Sending #{packet_size} bytes using UDP"
65 end
66 end
67
68 response = EM::DefaultDeferrable.new
69 result = self.old_send(method,packet,packet_data)
70
71 # handle a successful response
72 result.callback do |packet|
73 response.succeed packet
74 end
75 # return an error message if we fail
76 result.errback do
77 response.fail "No response from nameservers list"
78 end
79
80 return response
81 end
82
83 def receive_datagram(data)
84 response = Net::DNS::Packet.parse(data, nil)
85 if r = @outstanding.delete(response.header.id)
86 r.succeed(response)
87 else
88 @logger.warn "Got datagram with no outstanding request: #{response}"
89 end
90 end
91
92 def resend_udp_packet request
93 ns = @config[:nameservers][ rand(@config[:nameservers].size) ]
94 udp_socket.send_datagram(request.packet.data, ns.to_s, @config[:port])
95 end
96
97 private
98
99 def send_udp_async(packet, packet_data)
100
101 # generate a request
102 request = UDPRequest.new packet, self
103 @logger.warn "ID collision: #{packet.header.id}" if @outstanding[packet.header.id]
104 @outstanding[packet.header.id] = request
105
106 # pick a random nameserver and query it
107 ns = @config[:nameservers][ rand(@config[:nameservers].size) ]
108 @logger.info "Contacting nameserver #{ns} port #{@config[:port]}"
109 udp_socket.send_datagram(packet_data, ns.to_s, @config[:port])
110
111 # return the result
112 request
113 end
114
115 def udp_socket
116 # start listening if we aren't already
117 unless @udp_socket
118 unbind_signaller = proc {@udp_socket = nil}
119 @udp_socket = EM::open_datagram_socket( @config[:source_address].to_s, @config[:source_port], UDPSocket, self ) {|c|
120 c.unbind_signaller = unbind_signaller
121 }
122 end
123 @udp_socket
124 end
125
126 # should implement TCP
127 def send_tcp_async(packet, packet_data)
128 raise NotImplementedError.new "TCP is not yet supported"
129 end
130
131 end
132
133 class UDPSocket < EM::Connection
134 attr_accessor :unbind_signaller
135
136 def initialize resolver
137 @resolver = resolver
138 end
139
140 def receive_data data
141 @resolver.receive_datagram(data)
142 end
143
144 def unbind
145 @unbind_signaller.call if @unbind_signaller
146 end
147 end
148
149 class UDPRequest
150 include EM::Deferrable
151 attr_accessor :attempts, :packet
152
153 def initialize packet, resolver
154 @packet = packet
155 @resolver = resolver
156 @attempts = 0
157 self.timeout @resolver.udp_timeout
158 end
159
160 def fail
161 if @attempts < @resolver.retry_number
162 @attempts += 1
163 @resolver.resend_udp_packet self
164 self.timeout @resolver.udp_timeout
165 else
166 super
167 end
168 end
169
170 end
171
172 # wraps a udp request to return a more useful response
173 class MXQuery
174 include EM::Deferrable
175
176 def initialize(udprequest)
177
178 udprequest.callback do |packet|
179 arr = []
180 packet.answer.each do |entry|
181 arr << entry if entry.type == 'MX'
182 end
183 succeed(arr)
184 end
185
186 udprequest.errback do |error|
187 fail(error)
188 end
189 end
190 end
191
192 # wraps a udp request to return a more useful response
193 class AsyncQuery
194 include EM::Deferrable
195
196 def initialize(udprequest)
197 # handle a successful response
198 udprequest.callback do |packet|
199 succeed packet
200 end
201 # return an error message if we fail
202 udprequest.errback do
203 fail "No response from nameservers"
204 end
205 end
206
207 end
208
209 end
210
211end
212
213if __FILE__ == $0
214 EM.run do
215 res = EM::P::AsyncResolver.new
216 count = 0
217 5.times do
218 ["gmail.com", "yahoo.com", "otherinbox.com", "asdfvaesr.com"].each do |domain|
219 count += 1
220 result = res.mx_async(domain)
221 result.callback { |answer| puts "Got MX records for #{domain}: #{answer.inspect}"; count -= 1; }
222 result.errback { |err| STDERR.puts "Got error for #{domain}: #{err}"; count -= 1; }
223 end
224 end
225 5.times do
226 ["gmail.com", "yahoo.com", "otherinbox.com", "asdfvaesr.com"].each do |domain|
227 count += 1
228 result = res.query_async(domain)
229 result.callback { |packet| puts "Got result for #{domain}: #{packet.answer}"; count -= 1; }
230 result.errback { |err| STDERR.puts "Got error for #{domain}: #{err}"; count -= 1; }
231 end
232 end
233 EM.add_periodic_timer(1) { EM.stop_event_loop if count == 0 }
234 end
235end