Class: AMQP::Client::Connection
- Inherits:
-
Object
- Object
- AMQP::Client::Connection
- Defined in:
- lib/amqp/client/connection.rb,
lib/amqp/client/channel.rb
Overview
Represents a single established AMQP connection
Defined Under Namespace
Classes: Channel
Instance Attribute Summary collapse
-
#frame_max ⇒ Integer
readonly
The max frame size negotiated between the client and the broker.
Callbacks collapse
-
#on_blocked {|String| ... } ⇒ nil
Callback called when client is blocked by the broker.
-
#on_unblocked { ... } ⇒ nil
Callback called when client is unblocked by the broker.
Class Method Summary collapse
- .connect(uri, read_loop_thread: true, **options) ⇒ Object deprecated Deprecated.
Instance Method Summary collapse
-
#blocked? ⇒ Bool
Indicates that the server is blocking publishes.
-
#channel(id = nil) ⇒ Channel
Open an AMQP channel.
-
#close(reason: "", code: 200) ⇒ nil
Gracefully close a connection.
-
#closed? ⇒ Boolean
True if the connection is closed.
-
#initialize(uri = "", read_loop_thread: true, **options) ⇒ Connection
constructor
Establish a connection to an AMQP broker.
-
#read_loop ⇒ nil
Reads from the socket, required for any kind of progress.
-
#update_secret(secret, reason) ⇒ nil
Update authentication secret, for example when an OAuth backend is used.
-
#with_channel {|Channel| ... } ⇒ Object
Declare a new channel, yield, and then close the channel.
Constructor Details
#initialize(uri = "", read_loop_thread: true, **options) ⇒ Connection
Establish a connection to an AMQP broker
29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 |
# File 'lib/amqp/client/connection.rb', line 29 def initialize(uri = "", read_loop_thread: true, **) uri = URI.parse(uri) tls = uri.scheme == "amqps" port = port_from_env || uri.port || (tls ? 5671 : 5672) host = uri.host || "localhost" user = uri.user || "guest" password = uri.password || "guest" vhost = URI.decode_www_form_component(uri.path[1..] || "/") = URI.decode_www_form(uri.query || "").map! { |k, v| [k.to_sym, v] }.to_h.merge() socket = open_socket(host, port, tls, ) channel_max, frame_max, heartbeat = establish(socket, user, password, vhost, ) @socket = socket @channel_max = channel_max.zero? ? 65_536 : channel_max @frame_max = frame_max @heartbeat = heartbeat @channels = {} @channels_lock = Mutex.new @closed = nil @replies = ::Queue.new @write_lock = Mutex.new @blocked = nil @on_blocked = ->(reason) { warn "AMQP-Client blocked by broker: #{reason}" } @on_unblocked = -> { warn "AMQP-Client unblocked by broker" } # Only used with heartbeats @last_activity_time = Process.clock_gettime(Process::CLOCK_MONOTONIC) Thread.new { read_loop } if read_loop_thread end |
Instance Attribute Details
#frame_max ⇒ Integer (readonly)
The max frame size negotiated between the client and the broker
80 81 82 |
# File 'lib/amqp/client/connection.rb', line 80 def frame_max @frame_max end |
Class Method Details
.connect(uri, read_loop_thread: true, **options) ⇒ Object
Alias for #initialize
74 75 76 |
# File 'lib/amqp/client/connection.rb', line 74 def self.connect(uri, read_loop_thread: true, **) new(uri, read_loop_thread: read_loop_thread, **) end |
Instance Method Details
#blocked? ⇒ Bool
Indicates that the server is blocking publishes. If the client keeps publishing the server will stop reading from the socket. Use the #on_blocked callback to get notified when the server is resource constrained.
67 68 69 |
# File 'lib/amqp/client/connection.rb', line 67 def blocked? !@blocked.nil? end |
#channel(id = nil) ⇒ Channel
Open an AMQP channel
92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 |
# File 'lib/amqp/client/connection.rb', line 92 def channel(id = nil) raise ArgumentError, "Channel ID cannot be 0" if id&.zero? raise ArgumentError, "Channel ID higher than connection's channel max #{@channel_max}" if id && id > @channel_max ch = @channels_lock.synchronize do if id @channels[id] ||= Channel.new(self, id) else 1.upto(@channel_max) do |i| break id = i unless @channels.key? i end raise Error, "Max channels reached" if id.nil? @channels[id] = Channel.new(self, id) end end ch.open end |
#close(reason: "", code: 200) ⇒ nil
Gracefully close a connection
127 128 129 130 131 132 133 134 135 136 137 138 139 |
# File 'lib/amqp/client/connection.rb', line 127 def close(reason: "", code: 200) return if @closed @closed = [code, reason] @channels.each_value { |ch| ch.closed!(:connection, code, reason, 0, 0) } if @blocked @socket.close else write_bytes FrameBytes.connection_close(code, reason) expect(:close_ok) end nil end |
#closed? ⇒ Boolean
True if the connection is closed
153 154 155 |
# File 'lib/amqp/client/connection.rb', line 153 def closed? !@closed.nil? end |
#on_blocked {|String| ... } ⇒ nil
Callback called when client is blocked by the broker
162 163 164 165 |
# File 'lib/amqp/client/connection.rb', line 162 def on_blocked(&blk) @on_blocked = blk nil end |
#on_unblocked { ... } ⇒ nil
Callback called when client is unblocked by the broker
170 171 172 173 |
# File 'lib/amqp/client/connection.rb', line 170 def on_unblocked(&blk) @on_unblocked = blk nil end |
#read_loop ⇒ nil
Reads from the socket, required for any kind of progress. Blocks until the connection is closed. Normally run as a background thread automatically.
195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 |
# File 'lib/amqp/client/connection.rb', line 195 def read_loop # read more often than write so that channel errors crop up early Thread.current.priority += 1 socket = @socket frame_max = @frame_max frame_start = String.new(capacity: 7) frame_buffer = String.new(capacity: frame_max) loop do socket.read(7, frame_start) || raise(IOError) type, channel_id, frame_size = frame_start.unpack("C S> L>") frame_max >= frame_size || raise(Error, "Frame size #{frame_size} larger than negotiated max frame size #{frame_max}") # read the frame content socket.read(frame_size, frame_buffer) || raise(IOError) # make sure that the frame end is correct frame_end = socket.readchar.ord raise Error::UnexpectedFrameEnd, frame_end if frame_end != 206 # parse the frame, will return false if a close frame was received parse_frame(type, channel_id, frame_buffer) || return update_last_activity end nil rescue *READ_EXCEPTIONS => e @closed ||= [400, "read error: #{e.}"] nil # ignore read errors ensure @closed ||= [400, "unknown"] @replies.close begin if @write_lock.owned? # if connection is blocked @socket.close else @write_lock.synchronize do @socket.close end end rescue *READ_EXCEPTIONS nil end end |
#update_secret(secret, reason) ⇒ nil
Update authentication secret, for example when an OAuth backend is used
145 146 147 148 149 |
# File 'lib/amqp/client/connection.rb', line 145 def update_secret(secret, reason) write_bytes FrameBytes.update_secret(secret, reason) expect(:update_secret_ok) nil end |
#with_channel {|Channel| ... } ⇒ Object
Declare a new channel, yield, and then close the channel
114 115 116 117 118 119 120 121 |
# File 'lib/amqp/client/connection.rb', line 114 def with_channel ch = channel begin yield ch ensure ch.close end end |