Skip to content

Commit 68ea988

Browse files
committed
added a spec
1 parent e63eb6c commit 68ea988

2 files changed

Lines changed: 20 additions & 2 deletions

File tree

spec/amqproxy/server_spec.cr

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -90,6 +90,19 @@ describe AMQProxy::Server do
9090
end
9191
end
9292

93+
it "creates multiple upstreams when max_upstream_channels is low" do
94+
with_server(max_upstream_channels: 10_u16) do |server, proxy_url|
95+
AMQP::Client.start(proxy_url) do |conn|
96+
25.times { conn.channel }
97+
server.client_connections.should eq 1
98+
server.upstream_connections.should be >= 3
99+
end
100+
sleep 0.1.seconds
101+
server.client_connections.should eq 0
102+
server.upstream_connections.should eq 3
103+
end
104+
end
105+
93106
it "can reconnect if upstream closes" do
94107
with_server do |server, proxy_url|
95108
Fiber.yield

spec/spec_helper.cr

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,14 +10,19 @@ Log.setup_from_env(default_level: :error)
1010
MAYBE_SUDO = (ENV.has_key?("NO_SUDO") || `id -u` == "0\n") ? "" : "sudo "
1111

1212
UPSTREAM_URL = begin
13+
# URI.parse ENV.fetch("UPSTREAM_URL", "amqp://127.0.0.1:5672?idle_connection_timeout=5&max_upstream_channels=65535")
1314
URI.parse ENV.fetch("UPSTREAM_URL", "amqp://127.0.0.1:5672?idle_connection_timeout=5")
1415
rescue e : URI::Error
1516
puts "Invalid UPSTREAM_URL: #{e}"
1617
exit 1
1718
end
1819

19-
def with_server(idle_connection_timeout = 5, &)
20-
server = AMQProxy::Server.new(UPSTREAM_URL)
20+
def with_server(idle_connection_timeout = 5, max_upstream_channels = UInt16::MAX, &)
21+
tls = UPSTREAM_URL.scheme == "amqps"
22+
host = UPSTREAM_URL.host || "127.0.0.1"
23+
port = UPSTREAM_URL.port || 5672
24+
port = 5671 if tls && UPSTREAM_URL.port.nil?
25+
server = AMQProxy::Server.new(host, port, tls, idle_connection_timeout, max_upstream_channels)
2126
tcp_server = TCPServer.new("127.0.0.1", 0)
2227
amqp_url = "amqp://#{tcp_server.local_address}"
2328
spawn { server.listen(tcp_server) }

0 commit comments

Comments
 (0)