diff --git a/modules/src/main/java/org/archive/modules/fetcher/FetchHTTP2.java b/modules/src/main/java/org/archive/modules/fetcher/FetchHTTP2.java
index 46eedea3..3c05ca98 100644
--- a/modules/src/main/java/org/archive/modules/fetcher/FetchHTTP2.java
+++ b/modules/src/main/java/org/archive/modules/fetcher/FetchHTTP2.java
@@ -36,7 +36,6 @@ import org.apache.http.impl.conn.DefaultClientConnectionOperator;
import org.apache.http.impl.conn.SchemeRegistryFactory;
import org.apache.http.impl.io.AbstractSessionInputBuffer;
import org.apache.http.impl.io.AbstractSessionOutputBuffer;
-import org.apache.http.impl.io.IdentityInputStream;
import org.apache.http.impl.io.IdentityOutputStream;
import org.apache.http.io.SessionInputBuffer;
import org.apache.http.io.SessionOutputBuffer;
@@ -166,12 +165,13 @@ public class FetchHTTP2 extends Processor implements Lifecycle {
if (h != null) {
Long.parseLong(h.getValue());
}
-
try {
if (!request.isAborted()) {
// Force read-to-end, so that any socket hangs occur here,
// not in later modules.
- rec.getRecordedInput().readFullyOrUntil(softMax);
+
+ // XXX does it matter that we're circumventing the library here? response.getEntity().getContent()
+ rec.getRecordedInput().readFullyOrUntil(softMax);
}
} catch (RecorderTimeoutException ex) {
doAbort(curi, request, TIMER_TRUNC);
@@ -326,6 +326,7 @@ public class FetchHTTP2 extends Processor implements Lifecycle {
}
protected static class RecordingClientConnectionOperator extends DefaultClientConnectionOperator {
+
protected final Recorder rec;
protected RecordingClientConnectionOperator(SchemeRegistry schemes, Recorder rec) {
@@ -338,9 +339,7 @@ public class FetchHTTP2 extends Processor implements Lifecycle {
return new DefaultClientConnection() {
@Override
protected SessionInputBuffer createSessionInputBuffer(Socket socket, int buffersize, HttpParams params) throws IOException {
- SessionInputBuffer sib = super.createSessionInputBuffer(socket, buffersize, params);
- InputStream ris = rec.inputWrap(new IdentityInputStream(sib));
- return new HcInputWrapper(ris, buffersize, params);
+ return new RecordingSocketInputBuffer(rec, socket, params);
}
@Override
diff --git a/modules/src/main/java/org/archive/modules/fetcher/RecordingSocketInputBuffer.java b/modules/src/main/java/org/archive/modules/fetcher/RecordingSocketInputBuffer.java
new file mode 100644
index 00000000..48baf60c
--- /dev/null
+++ b/modules/src/main/java/org/archive/modules/fetcher/RecordingSocketInputBuffer.java
@@ -0,0 +1,124 @@
+package org.archive.modules.fetcher;
+
+import java.io.IOException;
+import java.io.InputStream;
+import java.net.Socket;
+
+import org.apache.http.impl.io.HttpTransportMetricsImpl;
+import org.apache.http.io.HttpTransportMetrics;
+import org.apache.http.io.SessionInputBuffer;
+import org.apache.http.params.HttpParams;
+import org.apache.http.protocol.HTTP;
+import org.apache.http.util.CharArrayBuffer;
+import org.archive.util.Recorder;
+
+public class RecordingSocketInputBuffer implements SessionInputBuffer {
+
+ protected Socket socket;
+ protected Recorder recorder;
+ protected InputStream in;
+ protected HttpTransportMetricsImpl metrics;
+
+ public RecordingSocketInputBuffer(Recorder rec, Socket socket, HttpParams params) throws IOException {
+ if (socket == null) {
+ throw new IllegalArgumentException("Socket may not be null");
+ }
+ this.recorder = rec;
+ this.socket = socket;
+ this.in = recorder.inputWrap(socket.getInputStream());
+
+ this.metrics = new HttpTransportMetricsImpl();
+ }
+
+ @Override
+ public int read(byte[] b, int off, int len) throws IOException {
+ int n = in.read(b, off, len);
+ if (n > 0) {
+ metrics.incrementBytesTransferred(n);
+ }
+ return n;
+ }
+
+ @Override
+ public int read(byte[] b) throws IOException {
+ int n = in.read(b);
+ if (n > 0) {
+ metrics.incrementBytesTransferred(n);
+ }
+ return n;
+ }
+
+ @Override
+ public int read() throws IOException {
+ int n = in.read();
+ if (n > 0) {
+ metrics.incrementBytesTransferred(n);
+ }
+ return n;
+ }
+
+ /**
+ * Reads a complete line of characters up to a line delimiter from this
+ * session buffer into the given line buffer. The number of chars actually
+ * read is returned as an integer. The line delimiter itself is discarded.
+ * If no char is available because the end of the stream has been reached,
+ * the value -1 is returned. This method blocks until input
+ * data is available, end of file is detected, or an exception is thrown.
+ *
+ * This is only used for http headers, which are all ascii, so input is + * treated as such. + *
+ * This method treats a lone LF as a valid line delimiters in addition to + * CR-LF required by the HTTP specification. + * + * @param charbuffer + * the line buffer. + * @return number of bytes (i.e. ascii characters) read + * @exception IOException + * if an I/O error occurs. + */ + @Override + public int readLine(CharArrayBuffer buffer) throws IOException { + int bytesRead = 0; + int b = in.read(); + while (b >= 0 && b != HTTP.LF) { + bytesRead++; + buffer.append((char) b); + b = in.read(); + } + if (b >= 0) { + bytesRead++; // count LF + } + + // if line ends with CR-LF, get rid of the CR + if (bytesRead > 0 && buffer.charAt(buffer.length() - 1) == HTTP.CR) { + buffer.setLength(buffer.length() - 1); + } + + if (bytesRead > 0) { + metrics.incrementBytesTransferred(bytesRead); + } + return bytesRead; + } + + @Override + public String readLine() throws IOException { + CharArrayBuffer charbuffer = new CharArrayBuffer(64); + int l = readLine(charbuffer); + if (l != -1) { + return charbuffer.toString(); + } else { + return null; + } + } + + @Override + public boolean isDataAvailable(int timeout) throws IOException { + throw new RuntimeException("not implemented"); + } + + @Override + public HttpTransportMetrics getMetrics() { + return metrics; + } +} \ No newline at end of file