-
Notifications
You must be signed in to change notification settings - Fork 954
ARTEMIS-1924 Test consumer cleanup after socket connection reset #2149
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -16,15 +16,23 @@ | |
| */ | ||
| package org.apache.activemq.artemis.tests.integration.amqp; | ||
|
|
||
| import java.net.URI; | ||
| import java.util.Arrays; | ||
| import java.util.Collection; | ||
| import java.util.Map; | ||
| import java.util.concurrent.TimeUnit; | ||
|
|
||
| import org.apache.activemq.artemis.core.server.ActiveMQServer; | ||
| import org.apache.activemq.artemis.tests.util.SpawnedVMSupport; | ||
| import org.apache.activemq.transport.amqp.client.AmqpClient; | ||
| import org.apache.activemq.transport.amqp.client.AmqpConnection; | ||
| import org.apache.activemq.transport.amqp.client.AmqpMessage; | ||
| import org.apache.activemq.transport.amqp.client.AmqpReceiver; | ||
| import org.apache.activemq.transport.amqp.client.AmqpSender; | ||
| import org.apache.activemq.transport.amqp.client.AmqpSession; | ||
| import org.apache.activemq.transport.amqp.client.AmqpValidator; | ||
| import org.apache.qpid.proton.engine.Connection; | ||
| import org.junit.Assert; | ||
| import org.junit.Test; | ||
| import org.junit.runner.RunWith; | ||
| import org.junit.runners.Parameterized; | ||
|
|
@@ -80,4 +88,68 @@ public void inspectOpenedResource(Connection connection) { | |
| connection.close(); | ||
| } | ||
|
|
||
| private static final String QUEUE_NAME = "queue://testHeartless"; | ||
|
|
||
| // This test is validating a scenario where the client will leave with connection reset | ||
| // This is done by setting soLinger=0 on the socket, which will make the system to issue a connection.reset instead of sending a | ||
| // disconnect. | ||
| @Test(timeout = 60000) | ||
| public void testCloseConsumerOnConnectionReset() throws Exception { | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I dont really understand why this test is in the AmqpNoHearbeatsTest class, it feels like it should operate about the same regardless whether AMQP idle-timeout handling is enabled or not. Its definitely not where I would go looking for such a test. It also seems far more appropriate to testing the changes made in #2139 for https://issues.apache.org/jira/browse/ARTEMIS-1927 (without test).
Contributor
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The two are related... you should be able to cleanup consumer, disable heartbeats and call disconnect on windows. ARTEMIS-1927 was found after I disabled heart beats and killed a consumer on windows which issued a connection reset.
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. To me it seems that the broker behaviour should essentially be the same here regardless whether idle-timeout is enabled or not, i.e a connection explicitly goes away, the broker cleans up a consumer on it, and its not really testing idle-timeout behaviour or the lack of it, making it feels a bit odd it being in here. |
||
|
|
||
| AmqpClient client = createAmqpClient(); | ||
| assertNotNull(client); | ||
|
|
||
| client.setValidator(new AmqpValidator() { | ||
|
|
||
| @Override | ||
| public void inspectOpenedResource(Connection connection) { | ||
| assertEquals("idle timeout was not disabled", 0, connection.getTransport().getRemoteIdleTimeout()); | ||
| } | ||
| }); | ||
|
|
||
| AmqpConnection connection = addConnection(client.connect()); | ||
| assertNotNull(connection); | ||
|
|
||
| connection.getStateInspector().assertValid(); | ||
| AmqpSession session = connection.createSession(); | ||
| AmqpReceiver receiver = session.createReceiver(QUEUE_NAME); | ||
|
|
||
| // This test needs a remote process exiting without closing the socket | ||
| // with soLinger=0 on the socket so it will issue a connection.reset | ||
| Process p = SpawnedVMSupport.spawnVM(AmqpNoHearbeatsTest.class.getName(), "testConnectionReset"); | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Would using the test name be a better idea here to make it align better? |
||
| Assert.assertEquals(33, p.waitFor()); | ||
|
|
||
| AmqpSender sender = session.createSender(QUEUE_NAME); | ||
|
|
||
| for (int i = 0; i < 10; i++) { | ||
| AmqpMessage msg = new AmqpMessage(); | ||
| msg.setBytes(new byte[] {0}); | ||
| sender.send(msg); | ||
| } | ||
|
|
||
| receiver.flow(20); | ||
|
|
||
| for (int i = 0; i < 10; i++) { | ||
| AmqpMessage msg = receiver.receive(1, TimeUnit.SECONDS); | ||
| Assert.assertNotNull(msg); | ||
| msg.accept(); | ||
| } | ||
| } | ||
|
|
||
| public static void main(String[] arg) { | ||
| if (arg.length > 0 && arg[0].equals("testConnectionReset")) { | ||
| try { | ||
| AmqpClient client = new AmqpClient(new URI("tcp://127.0.0.1:5672?transport.soLinger=0"), null, null); | ||
| AmqpConnection connection = client.connect(); | ||
| AmqpSession session = connection.createSession(); | ||
| AmqpReceiver receiver = session.createReceiver(QUEUE_NAME); | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Linked with above comments, using the actual test name as the argument (instead of the differerent "testConnectionReset" literal) would mean this queue name constant wouldnt be needed and the argument could be used directly. |
||
| receiver.flow(10); | ||
| System.exit(33); | ||
| } catch (Throwable e) { | ||
| e.printStackTrace(); | ||
| System.exit(-1); | ||
| } | ||
| } | ||
| } | ||
|
|
||
| } | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Why does a test named "testCloseConsumerOnConnectionReset" use a queue named "queue://testHeartless", when "testHeartless" is the name of a different test in the class? Does it need a "queue://" prefix?
I'd just use getTestName() inside the test get a helpful value.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
I made a mistake. I created a static property because this needs a Spawned VM. Getting it from the static property seemed easier.
Do you really care about this? if so we can change the test.