Author: dkulp
Date: Fri Oct 19 18:37:16 2012
New Revision: 1400214
URL: http://svn.apache.org/viewvc?rev=1400214&view=rev
Log:
Minor update so on a loaded system, if it failed to put the item on the queue,
try again on the first data packet
Modified:
cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/AsyncHTTPConduit.java
cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/CXFHttpAsyncResponseConsumer.java
Modified:
cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/AsyncHTTPConduit.java
URL:
http://svn.apache.org/viewvc/cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/AsyncHTTPConduit.java?rev=1400214&r1=1400213&r2=1400214&view=diff
==============================================================================
---
cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/AsyncHTTPConduit.java
(original)
+++
cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/AsyncHTTPConduit.java
Fri Oct 19 18:37:16 2012
@@ -474,11 +474,16 @@ public class AsyncHTTPConduit extends UR
c.getCredentialsProvider().setCredentials(AuthScope.ANY,
creds);
}
c.execute(new CXFHttpAsyncRequestProducer(entity, outbuf),
- new CXFHttpAsyncResponseConsumer(inbuf,
responseCallback),
+ new CXFHttpAsyncResponseConsumer(this, inbuf,
responseCallback),
ctx,
callback);
}
+ protected void retrySetHttpResponse(HttpResponse r) {
+ if (httpResponse == null && isAsync) {
+ setHttpResponse(r);
+ }
+ }
protected synchronized void setHttpResponse(HttpResponse r) {
httpResponse = r;
if (isAsync) {
@@ -487,7 +492,7 @@ public class AsyncHTTPConduit extends UR
handleResponseOnWorkqueue(false, true);
isAsync = false; // don't trigger another start on next
block. :-)
} catch (Exception ex) {
- ex.printStackTrace();
+ //ignore, we'll try again on the next consume;
}
}
notifyAll();
Modified:
cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/CXFHttpAsyncResponseConsumer.java
URL:
http://svn.apache.org/viewvc/cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/CXFHttpAsyncResponseConsumer.java?rev=1400214&r1=1400213&r2=1400214&view=diff
==============================================================================
---
cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/CXFHttpAsyncResponseConsumer.java
(original)
+++
cxf/trunk/rt/transports/http-hc/src/main/java/org/apache/cxf/transport/http/asyncclient/CXFHttpAsyncResponseConsumer.java
Fri Oct 19 18:37:16 2012
@@ -21,6 +21,7 @@ package org.apache.cxf.transport.http.as
import java.io.IOException;
+import
org.apache.cxf.transport.http.asyncclient.AsyncHTTPConduit.AsyncWrappedOutputStream;
import org.apache.http.HttpException;
import org.apache.http.HttpResponse;
import org.apache.http.nio.ContentDecoder;
@@ -31,16 +32,21 @@ import org.apache.http.protocol.HttpCont
public class CXFHttpAsyncResponseConsumer implements
HttpAsyncResponseConsumer<Boolean> {
private final SharedInputBuffer buf;
+ private final AsyncWrappedOutputStream outstream;
private final CXFResponseCallback responseCallback;
private volatile boolean completed;
private volatile Exception exception;
+ private volatile HttpResponse response;
public CXFHttpAsyncResponseConsumer(
- final SharedInputBuffer buf, final CXFResponseCallback
responseCallback) {
+ final AsyncWrappedOutputStream asyncWrappedOutputStream,
+ final SharedInputBuffer buf,
+ final CXFResponseCallback responseCallback) {
super();
- this.buf = buf;
+ this.outstream = asyncWrappedOutputStream;
this.responseCallback = responseCallback;
+ this.buf = buf;
}
@Override
@@ -56,12 +62,14 @@ public class CXFHttpAsyncResponseConsume
}
@Override
- public void responseReceived(final HttpResponse response) throws
IOException, HttpException {
+ public void responseReceived(final HttpResponse resp) throws IOException,
HttpException {
+ response = resp;
responseCallback.responseReceived(response);
}
@Override
public void consumeContent(final ContentDecoder dec, final IOControl ioc)
throws IOException {
+ outstream.retrySetHttpResponse(response);
buf.consumeContent(dec, ioc);
}