SWS-564 - CommonsHttpMessageSender no longer properly shuts down MultiThreadedHttpConnectionManager
This commit is contained in:
@@ -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.
|
||||
*
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/*
|
||||
|
||||
@@ -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()) {
|
||||
|
||||
@@ -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));
|
||||
|
||||
@@ -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());
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user