diff --git a/org.ektorp/src/main/java/org/ektorp/http/StdHttpResponse.java b/org.ektorp/src/main/java/org/ektorp/http/StdHttpResponse.java index a84248ee..96672b99 100644 --- a/org.ektorp/src/main/java/org/ektorp/http/StdHttpResponse.java +++ b/org.ektorp/src/main/java/org/ektorp/http/StdHttpResponse.java @@ -72,7 +72,11 @@ public String getContentType() { public InputStream getContent() { try { - return new ConnectionReleasingInputStream(entity.getContent()); + InputStream content = entity.getContent(); + if(content == null) { + return null; + } + return new ConnectionReleasingInputStream(content); } catch (Exception e) { throw Exceptions.propagate(e); } diff --git a/org.ektorp/src/main/java/org/ektorp/http/StdResponseHandler.java b/org.ektorp/src/main/java/org/ektorp/http/StdResponseHandler.java index 2038fc71..8bf7ff62 100644 --- a/org.ektorp/src/main/java/org/ektorp/http/StdResponseHandler.java +++ b/org.ektorp/src/main/java/org/ektorp/http/StdResponseHandler.java @@ -68,10 +68,17 @@ protected static JsonNode responseBodyAsNode(InputStream inputStream, ObjectMapp return mapper.readTree(inputStream); } + @edu.umd.cs.findbugs.annotations.SuppressWarnings( + value="NP_LOAD_OF_KNOWN_NULL_VALUE", + justification="The response should be null only if it is not a valid OK message. " + + "In that case, it is safe to die. In any case, this method should be owerwritten if the class is extended.") protected static T checkResponseBodyOkAndReturnDefaultValue(HttpResponse hr, T defaultValue, ObjectMapper mapper) throws IOException { - InputStream content = hr.getContent(); + InputStream content = null; try { content = hr.getContent(); + if(content == null){ + return defaultValue; + } JsonNode body = responseBodyAsNode(content, MAPPER); JsonNode okNode = body.get("ok"); if (okNode != null) { @@ -90,7 +97,7 @@ protected static T checkResponseBodyOkAndReturnDefaultValue(HttpResponse hr, public T error(HttpResponse hr) { throw StdResponseHandler.createDbAccessException(hr); } - + public T success(HttpResponse hr) throws Exception { return checkResponseBodyOkAndReturnDefaultValue(hr, null, MAPPER); } diff --git a/org.ektorp/src/main/java/org/ektorp/impl/changes/ContinuousChangesFeed.java b/org.ektorp/src/main/java/org/ektorp/impl/changes/ContinuousChangesFeed.java index 1750ceab..1ae12d9a 100644 --- a/org.ektorp/src/main/java/org/ektorp/impl/changes/ContinuousChangesFeed.java +++ b/org.ektorp/src/main/java/org/ektorp/impl/changes/ContinuousChangesFeed.java @@ -1,18 +1,23 @@ package org.ektorp.impl.changes; -import java.io.*; -import java.util.*; -import java.util.concurrent.*; -import java.util.concurrent.atomic.*; - import com.fasterxml.jackson.core.JsonParseException; import com.fasterxml.jackson.databind.JsonMappingException; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.node.NullNode; -import org.ektorp.changes.*; +import org.ektorp.changes.ChangesFeed; +import org.ektorp.changes.DocumentChange; import org.ektorp.http.HttpResponse; -import org.ektorp.util.*; -import org.slf4j.*; +import org.ektorp.util.Exceptions; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import java.io.*; +import java.util.HashSet; +import java.util.Set; +import java.util.concurrent.BlockingQueue; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; /** * * @author henrik lundgren @@ -117,7 +122,11 @@ public void run() { String reason = !shouldRun ? "Cancelled" : "EOF"; LOG.info("Changes feed stopped. Reason: " + reason); } catch (Exception e) { - handleException(e); + if(!shouldRun) { + LOG.info("Changes feed was interrupted"); + } else { + handleException(e); + } } finally { sendInterruptMarker(); httpResponse.abort();