diff --git a/core/src/main/scala/kafka/server/KafkaConfig.scala b/core/src/main/scala/kafka/server/KafkaConfig.scala index cf22305caf70c..5a1dca395bb8f 100755 --- a/core/src/main/scala/kafka/server/KafkaConfig.scala +++ b/core/src/main/scala/kafka/server/KafkaConfig.scala @@ -782,7 +782,7 @@ object KafkaConfig { .define(SocketSendBufferBytesProp, INT, Defaults.SocketSendBufferBytes, HIGH, SocketSendBufferBytesDoc) .define(SocketReceiveBufferBytesProp, INT, Defaults.SocketReceiveBufferBytes, HIGH, SocketReceiveBufferBytesDoc) .define(SocketRequestMaxBytesProp, INT, Defaults.SocketRequestMaxBytes, atLeast(1), HIGH, SocketRequestMaxBytesDoc) - .define(MaxConnectionsPerIpProp, INT, Defaults.MaxConnectionsPerIp, atLeast(1), MEDIUM, MaxConnectionsPerIpDoc) + .define(MaxConnectionsPerIpProp, INT, Defaults.MaxConnectionsPerIp, atLeast(0), MEDIUM, MaxConnectionsPerIpDoc) .define(MaxConnectionsPerIpOverridesProp, STRING, Defaults.MaxConnectionsPerIpOverrides, MEDIUM, MaxConnectionsPerIpOverridesDoc) .define(ConnectionsMaxIdleMsProp, LONG, Defaults.ConnectionsMaxIdleMs, MEDIUM, ConnectionsMaxIdleMsDoc) diff --git a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala index a057e54bd7edf..0dad3c71230fb 100644 --- a/core/src/test/scala/unit/kafka/network/SocketServerTest.scala +++ b/core/src/test/scala/unit/kafka/network/SocketServerTest.scala @@ -20,8 +20,8 @@ package kafka.network import java.io._ import java.net._ import java.nio.ByteBuffer +import java.util.{HashMap, Properties, Random} import java.nio.channels.SocketChannel -import java.util.{HashMap, Random} import javax.net.ssl._ import com.yammer.metrics.core.{Gauge, Meter} @@ -134,8 +134,8 @@ class SocketServerTest extends JUnitSuite { channel.sendResponse(new RequestChannel.Response(request, Some(send), SendAction, Some(request.header.toString))) } - def connect(s: SocketServer = server, protocol: SecurityProtocol = SecurityProtocol.PLAINTEXT) = { - val socket = new Socket("localhost", s.boundPort(ListenerName.forSecurityProtocol(protocol))) + def connect(s: SocketServer = server, protocol: SecurityProtocol = SecurityProtocol.PLAINTEXT, localAddr: InetAddress = null) = { + val socket = new Socket("localhost", s.boundPort(ListenerName.forSecurityProtocol(protocol)), localAddr, 0) sockets += socket socket } @@ -443,6 +443,43 @@ class SocketServerTest extends JUnitSuite { assertNotNull(request) } + @Test + def testZeroMaxConnectionsPerIp() { + val newProps = TestUtils.createBrokerConfig(0, TestUtils.MockZkConnect, port = 0) + newProps.setProperty(KafkaConfig.MaxConnectionsPerIpProp, "0") + newProps.setProperty(KafkaConfig.MaxConnectionsPerIpOverridesProp, "%s:%s".format("127.0.0.1", "5")) + val server = new SocketServer(KafkaConfig.fromProps(newProps), new Metrics(), Time.SYSTEM, credentialProvider) + try { + server.startup() + // make the maximum allowable number of connections + val conns = (0 until 5).map(_ => connect(server)) + // now try one more (should fail) + val conn = connect(server) + conn.setSoTimeout(3000) + assertEquals(-1, conn.getInputStream.read()) + conn.close() + + // it should succeed after closing one connection + val address = conns.head.getInetAddress + conns.head.close() + TestUtils.waitUntilTrue(() => server.connectionCount(address) < conns.length, + "Failed to decrement connection count after close") + val conn2 = connect(server) + val serializedBytes = producerRequestBytes() + sendRequest(conn2, serializedBytes) + val request = server.requestChannel.receiveRequest(2000) + assertNotNull(request) + + // now try to connect from the external facing interface, which should fail + val conn3 = connect(s = server, localAddr = InetAddress.getLocalHost) + conn3.setSoTimeout(3000) + assertEquals(-1, conn3.getInputStream.read()) + conn3.close() + } finally { + shutdownServerAndMetrics(server) + } + } + @Test def testMaxConnectionsPerIpOverrides() { val overrideNum = server.config.maxConnectionsPerIp + 1 diff --git a/docs/upgrade.html b/docs/upgrade.html index 324f8df1403e6..0c5f5fd3247a2 100644 --- a/docs/upgrade.html +++ b/docs/upgrade.html @@ -67,6 +67,7 @@
offsets.retention.minutes to 1440.max.connections.per.ip minimum to zero and therefore allows IP-based filtering of inbound connections.