From 032d8daac9bf11cc8b71befb925d8edb7d26b5fc Mon Sep 17 00:00:00 2001 From: Arjen Poutsma Date: Tue, 10 Nov 2009 14:44:12 +0000 Subject: [PATCH] SWS-564 - CommonsHttpMessageSender no longer properly shuts down MultiThreadedHttpConnectionManager --- .../AbstractWebServiceConnection.java | 11 ++++ .../transport/http/CommonsHttpConnection.java | 20 +++++- .../http/CommonsHttpMessageSender.java | 12 +++- ...rviceMessageSenderIntegrationTestCase.java | 17 ++++-- ...mmonsHttpMessageSenderIntegrationTest.java | 61 +++++++++++++++++++ 5 files changed, 113 insertions(+), 8 deletions(-) diff --git a/core/src/main/java/org/springframework/ws/transport/AbstractWebServiceConnection.java b/core/src/main/java/org/springframework/ws/transport/AbstractWebServiceConnection.java index 8f8934df..6bab2bc8 100644 --- a/core/src/main/java/org/springframework/ws/transport/AbstractWebServiceConnection.java +++ b/core/src/main/java/org/springframework/ws/transport/AbstractWebServiceConnection.java @@ -33,7 +33,10 @@ public abstract class AbstractWebServiceConnection implements WebServiceConnecti private TransportOutputStream tos; + private boolean closed = false; + public final void send(WebServiceMessage message) throws IOException { + checkClosed(); onSendBeforeWrite(message); tos = createTransportOutputStream(); if (tos == null) { @@ -78,6 +81,7 @@ public abstract class AbstractWebServiceConnection implements WebServiceConnecti } public final WebServiceMessage receive(WebServiceMessageFactory messageFactory) throws IOException { + checkClosed(); onReceiveBeforeRead(); tis = createTransportInputStream(); if (tis == null) { @@ -138,11 +142,18 @@ public abstract class AbstractWebServiceConnection implements WebServiceConnecti } } onClose(); + closed = true; if (ioex != null) { throw ioex; } } + private void checkClosed() { + if (closed) { + throw new IllegalStateException("Connection has been closed and cannot be reused."); + } + } + /** * Template method invoked from {@link #close()}. Default implementation is empty. * diff --git a/core/src/main/java/org/springframework/ws/transport/http/CommonsHttpConnection.java b/core/src/main/java/org/springframework/ws/transport/http/CommonsHttpConnection.java index 3648efae..eff285ed 100644 --- a/core/src/main/java/org/springframework/ws/transport/http/CommonsHttpConnection.java +++ b/core/src/main/java/org/springframework/ws/transport/http/CommonsHttpConnection.java @@ -28,6 +28,7 @@ import java.util.Iterator; import org.apache.commons.httpclient.Header; import org.apache.commons.httpclient.HttpClient; import org.apache.commons.httpclient.URIException; +import org.apache.commons.httpclient.MultiThreadedHttpConnectionManager; import org.apache.commons.httpclient.methods.ByteArrayRequestEntity; import org.apache.commons.httpclient.methods.PostMethod; @@ -50,6 +51,8 @@ public class CommonsHttpConnection extends AbstractHttpSenderConnection { private ByteArrayOutputStream requestBuffer; + private MultiThreadedHttpConnectionManager connectionManager; + protected CommonsHttpConnection(HttpClient httpClient, PostMethod postMethod) { Assert.notNull(httpClient, "httpClient must not be null"); Assert.notNull(postMethod, "postMethod must not be null"); @@ -63,6 +66,9 @@ public class CommonsHttpConnection extends AbstractHttpSenderConnection { public void onClose() throws IOException { postMethod.releaseConnection(); + if (connectionManager != null) { + connectionManager.shutdown(); + } } /* @@ -97,7 +103,19 @@ public class CommonsHttpConnection extends AbstractHttpSenderConnection { protected void onSendAfterWrite(WebServiceMessage message) throws IOException { postMethod.setRequestEntity(new ByteArrayRequestEntity(requestBuffer.toByteArray())); requestBuffer = null; - httpClient.executeMethod(postMethod); + try { + httpClient.executeMethod(postMethod); + } catch (IllegalStateException ex) { + if ("Connection factory has been shutdown.".equals(ex.getMessage())) { + // The application context has been closed, resulting in a connection factory shutdown and an ISE. + // Let's create a new connection factory for this connection only. + connectionManager = new MultiThreadedHttpConnectionManager(); + httpClient.setHttpConnectionManager(connectionManager); + httpClient.executeMethod(postMethod); + } else { + throw ex; + } + } } /* diff --git a/core/src/main/java/org/springframework/ws/transport/http/CommonsHttpMessageSender.java b/core/src/main/java/org/springframework/ws/transport/http/CommonsHttpMessageSender.java index 85fbd14b..5efbfcc1 100644 --- a/core/src/main/java/org/springframework/ws/transport/http/CommonsHttpMessageSender.java +++ b/core/src/main/java/org/springframework/ws/transport/http/CommonsHttpMessageSender.java @@ -24,6 +24,7 @@ import java.util.Properties; import org.apache.commons.httpclient.Credentials; import org.apache.commons.httpclient.HostConfiguration; import org.apache.commons.httpclient.HttpClient; +import org.apache.commons.httpclient.HttpConnectionManager; import org.apache.commons.httpclient.HttpURL; import org.apache.commons.httpclient.HttpsURL; import org.apache.commons.httpclient.MultiThreadedHttpConnectionManager; @@ -33,6 +34,7 @@ import org.apache.commons.httpclient.UsernamePasswordCredentials; import org.apache.commons.httpclient.auth.AuthScope; import org.apache.commons.httpclient.methods.PostMethod; +import org.springframework.beans.factory.DisposableBean; import org.springframework.beans.factory.InitializingBean; import org.springframework.util.Assert; import org.springframework.ws.transport.WebServiceConnection; @@ -51,7 +53,8 @@ import org.springframework.ws.transport.WebServiceConnection; * @see #setCredentials(Credentials) * @since 1.0.0 */ -public class CommonsHttpMessageSender extends AbstractHttpWebServiceMessageSender implements InitializingBean { +public class CommonsHttpMessageSender extends AbstractHttpWebServiceMessageSender + implements InitializingBean, DisposableBean { private static final int DEFAULT_CONNECTION_TIMEOUT_MILLISECONDS = (60 * 1000); @@ -214,6 +217,13 @@ public class CommonsHttpMessageSender extends AbstractHttpWebServiceMessageSende } } + public void destroy() throws Exception { + HttpConnectionManager connectionManager = getHttpClient().getHttpConnectionManager(); + if (connectionManager instanceof MultiThreadedHttpConnectionManager) { + ((MultiThreadedHttpConnectionManager) connectionManager).shutdown(); + } + } + public WebServiceConnection createConnection(URI uri) throws IOException { PostMethod postMethod = new PostMethod(uri.toString()); if (isAcceptGzipEncoding()) { diff --git a/core/src/test/java/org/springframework/ws/transport/http/AbstractHttpWebServiceMessageSenderIntegrationTestCase.java b/core/src/test/java/org/springframework/ws/transport/http/AbstractHttpWebServiceMessageSenderIntegrationTestCase.java index 715ca575..6aaab067 100644 --- a/core/src/test/java/org/springframework/ws/transport/http/AbstractHttpWebServiceMessageSenderIntegrationTestCase.java +++ b/core/src/test/java/org/springframework/ws/transport/http/AbstractHttpWebServiceMessageSenderIntegrationTestCase.java @@ -81,15 +81,16 @@ public abstract class AbstractHttpWebServiceMessageSenderIntegrationTestCase ext private Context jettyContext; - private static final String URI_STRING = "http://localhost:8888/"; - private MessageFactory saajMessageFactory; private TransformerFactory transformerFactory; private WebServiceMessageFactory messageFactory; + private URI connectionUri; + protected final void setUp() throws Exception { + connectionUri = new URI("http://localhost:8888/"); jettyServer = new Server(8888); jettyContext = new Context(jettyServer, "/"); messageSender = createMessageSender(); @@ -111,7 +112,7 @@ public abstract class AbstractHttpWebServiceMessageSenderIntegrationTestCase ext } public void testSupports() throws URISyntaxException { - assertTrue("Message sender does not support HTTP url", messageSender.supports(new URI(URI_STRING))); + assertTrue("Message sender does not support HTTP url", messageSender.supports(connectionUri)); } public void testSendAndReceiveResponse() throws Exception { @@ -159,7 +160,7 @@ public abstract class AbstractHttpWebServiceMessageSenderIntegrationTestCase ext jettyContext.addServlet(new ServletHolder(servlet), "/"); jettyServer.start(); FaultAwareWebServiceConnection connection = - (FaultAwareWebServiceConnection) messageSender.createConnection(new URI(URI_STRING)); + (FaultAwareWebServiceConnection) messageSender.createConnection(connectionUri); SOAPMessage request = createRequest(); try { connection.send(new SaajSoapMessage(request)); @@ -171,11 +172,15 @@ public abstract class AbstractHttpWebServiceMessageSenderIntegrationTestCase ext } } + public void testReuseClosedConnection() throws Exception { + + } + private void validateResponse(Servlet servlet) throws Exception { jettyContext.addServlet(new ServletHolder(servlet), "/"); jettyServer.start(); FaultAwareWebServiceConnection connection = - (FaultAwareWebServiceConnection) messageSender.createConnection(new URI(URI_STRING)); + (FaultAwareWebServiceConnection) messageSender.createConnection(connectionUri); SOAPMessage request = createRequest(); try { connection.send(new SaajSoapMessage(request)); @@ -201,7 +206,7 @@ public abstract class AbstractHttpWebServiceMessageSenderIntegrationTestCase ext jettyContext.addServlet(new ServletHolder(servlet), "/"); jettyServer.start(); - WebServiceConnection connection = messageSender.createConnection(new URI(URI_STRING)); + WebServiceConnection connection = messageSender.createConnection(connectionUri); SOAPMessage request = createRequest(); try { connection.send(new SaajSoapMessage(request)); diff --git a/core/src/test/java/org/springframework/ws/transport/http/CommonsHttpMessageSenderIntegrationTest.java b/core/src/test/java/org/springframework/ws/transport/http/CommonsHttpMessageSenderIntegrationTest.java index b8a53845..e4a980e8 100644 --- a/core/src/test/java/org/springframework/ws/transport/http/CommonsHttpMessageSenderIntegrationTest.java +++ b/core/src/test/java/org/springframework/ws/transport/http/CommonsHttpMessageSenderIntegrationTest.java @@ -16,15 +16,28 @@ package org.springframework.ws.transport.http; +import java.io.IOException; import java.net.URI; import java.net.URISyntaxException; import java.util.Properties; +import javax.servlet.ServletException; +import javax.servlet.http.HttpServlet; +import javax.servlet.http.HttpServletRequest; +import javax.servlet.http.HttpServletResponse; +import javax.xml.soap.MessageFactory; import org.apache.commons.httpclient.ConnectTimeoutException; import org.apache.commons.httpclient.URIException; +import org.mortbay.jetty.Server; +import org.mortbay.jetty.servlet.Context; +import org.mortbay.jetty.servlet.ServletHolder; +import org.springframework.context.support.StaticApplicationContext; +import org.springframework.util.FileCopyUtils; import org.springframework.ws.MockWebServiceMessage; import org.springframework.ws.WebServiceMessage; +import org.springframework.ws.soap.saaj.SaajSoapMessage; +import org.springframework.ws.soap.saaj.SaajSoapMessageFactory; import org.springframework.ws.transport.WebServiceConnection; public class CommonsHttpMessageSenderIntegrationTest extends AbstractHttpWebServiceMessageSenderIntegrationTestCase { @@ -58,4 +71,52 @@ public class CommonsHttpMessageSenderIntegrationTest extends AbstractHttpWebServ messageSender.setMaxConnectionsPerHost(maxConnectionsPerHost); } + public void testContextClose() throws Exception { + MessageFactory messageFactory = MessageFactory.newInstance(); + Server jettyServer = new Server(8888); + Context jettyContext = new Context(jettyServer, "/"); + jettyContext.addServlet(new ServletHolder(new EchoServlet()), "/"); + jettyServer.start(); + WebServiceConnection connection = null; + try { + + StaticApplicationContext appContext = new StaticApplicationContext(); + appContext.registerSingleton("messageSender", CommonsHttpMessageSender.class); + appContext.refresh(); + + CommonsHttpMessageSender messageSender = (CommonsHttpMessageSender) appContext + .getBean("messageSender", CommonsHttpMessageSender.class); + connection = messageSender.createConnection(new URI("http://localhost:8888/")); + + appContext.close(); + + connection.send(new SaajSoapMessage(messageFactory.createMessage())); + connection.receive(new SaajSoapMessageFactory(messageFactory)); + } + finally { + if (connection != null) { + try { + connection.close(); + } catch (IOException ex) { + // ignore + } + } + if (jettyServer.isRunning()) { + jettyServer.stop(); + } + } + + } + + private class EchoServlet extends HttpServlet { + + protected void doPost(HttpServletRequest request, HttpServletResponse response) + throws ServletException, IOException { + response.setContentType("text/xml"); + FileCopyUtils.copy(request.getInputStream(), response.getOutputStream()); + + } + } + + }