diff --git a/apm-agent-core/src/main/java/co/elastic/apm/agent/report/AbstractIntakeApiHandler.java b/apm-agent-core/src/main/java/co/elastic/apm/agent/report/AbstractIntakeApiHandler.java index f10f8dc0d4..27f694d7c4 100644 --- a/apm-agent-core/src/main/java/co/elastic/apm/agent/report/AbstractIntakeApiHandler.java +++ b/apm-agent-core/src/main/java/co/elastic/apm/agent/report/AbstractIntakeApiHandler.java @@ -218,11 +218,15 @@ private void endRequest(boolean isFailed) { onRequestError(-1, writtenBytes, connection.getErrorStream(), e); } } finally { - HttpUtils.consumeAndClose(connection); + HttpURLConnection connectionToClose = connection; connection = null; os = null; countingOs = null; - deflater.reset(); + try { + HttpUtils.consumeAndClose(connectionToClose); + } finally { + deflater.reset(); + } } } } diff --git a/apm-agent-core/src/main/java/co/elastic/apm/agent/report/HttpUtils.java b/apm-agent-core/src/main/java/co/elastic/apm/agent/report/HttpUtils.java index b34aeb24c3..19b3b1ea16 100644 --- a/apm-agent-core/src/main/java/co/elastic/apm/agent/report/HttpUtils.java +++ b/apm-agent-core/src/main/java/co/elastic/apm/agent/report/HttpUtils.java @@ -18,6 +18,8 @@ */ package co.elastic.apm.agent.report; +import co.elastic.apm.agent.sdk.logging.Logger; +import co.elastic.apm.agent.sdk.logging.LoggerFactory; import org.stagemonitor.util.IOUtils; import javax.annotation.Nullable; @@ -29,6 +31,8 @@ public class HttpUtils { + private static final Logger logger = LoggerFactory.getLogger(HttpUtils.class); + private HttpUtils() { } @@ -57,12 +61,23 @@ public static String readToString(final InputStream inputStream) throws IOExcept */ public static void consumeAndClose(@Nullable HttpURLConnection connection) { if (connection != null) { - IOUtils.consumeAndClose(connection.getErrorStream()); + try { + IOUtils.consumeAndClose(connection.getErrorStream()); + } catch (RuntimeException e) { + logSuppressedException(e); + } try { IOUtils.consumeAndClose(connection.getInputStream()); } catch (IOException ignored) { // silently ignored + } catch (RuntimeException e) { + logSuppressedException(e); } } } + + private static void logSuppressedException(RuntimeException e) { + logger.warn("Suppressed exception while consuming APM Server response: {}", e.getMessage()); + logger.debug("Exception while consuming APM Server response", e); + } } diff --git a/apm-agent-core/src/test/java/co/elastic/apm/agent/report/HttpUtilsTest.java b/apm-agent-core/src/test/java/co/elastic/apm/agent/report/HttpUtilsTest.java index 031844d0d2..7e4f06d8af 100644 --- a/apm-agent-core/src/test/java/co/elastic/apm/agent/report/HttpUtilsTest.java +++ b/apm-agent-core/src/test/java/co/elastic/apm/agent/report/HttpUtilsTest.java @@ -24,6 +24,7 @@ import java.io.InputStream; import java.net.HttpURLConnection; +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; @@ -60,6 +61,26 @@ void consumeAndCloseException() throws IOException { verify(errorStream).close(); } + @Test + void consumeAndCloseRuntimeException() throws IOException { + HttpURLConnection connection = mock(HttpURLConnection.class); + doThrow(new NullPointerException("broken connection")).when(connection).getInputStream(); + + assertDoesNotThrow(() -> HttpUtils.consumeAndClose(connection)); + } + + @Test + void consumeAndCloseResponseAfterErrorStreamRuntimeException() throws IOException { + HttpURLConnection connection = mock(HttpURLConnection.class); + doThrow(new NullPointerException("broken error stream")).when(connection).getErrorStream(); + InputStream responseStream = mockEmptyInputStream(); + doReturn(responseStream).when(connection).getInputStream(); + + assertDoesNotThrow(() -> HttpUtils.consumeAndClose(connection)); + + verify(responseStream).close(); + } + @Test void consumeAndCloseResponseContent() throws IOException { HttpURLConnection connection = mock(HttpURLConnection.class); diff --git a/apm-agent-core/src/test/java/co/elastic/apm/agent/report/IntakeV2ReportingEventHandlerTest.java b/apm-agent-core/src/test/java/co/elastic/apm/agent/report/IntakeV2ReportingEventHandlerTest.java index d3469bb3a7..d39a592400 100644 --- a/apm-agent-core/src/test/java/co/elastic/apm/agent/report/IntakeV2ReportingEventHandlerTest.java +++ b/apm-agent-core/src/test/java/co/elastic/apm/agent/report/IntakeV2ReportingEventHandlerTest.java @@ -31,6 +31,7 @@ import co.elastic.apm.agent.impl.transaction.TransactionImpl; import co.elastic.apm.agent.report.processor.ProcessorEventHandler; import co.elastic.apm.agent.report.serialize.DslJsonSerializer; +import co.elastic.apm.agent.report.serialize.SerializationConstants; import com.dslplatform.json.DslJson; import com.dslplatform.json.JsonWriter; import com.fasterxml.jackson.databind.JsonNode; @@ -49,6 +50,7 @@ import java.io.ByteArrayInputStream; import java.io.IOException; import java.io.InputStreamReader; +import java.net.HttpURLConnection; import java.net.MalformedURLException; import java.net.URL; import java.util.Collections; @@ -66,6 +68,8 @@ import static com.github.tomakehurst.wiremock.client.WireMock.serviceUnavailable; import static com.github.tomakehurst.wiremock.client.WireMock.urlEqualTo; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatThrownBy; +import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.mock; class IntakeV2ReportingEventHandlerTest { @@ -104,6 +108,7 @@ void setUp() throws Exception { final ConfigurationRegistry configurationRegistry = SpyConfiguration.createSpyConfig(); final ReporterConfigurationImpl reporterConfiguration = configurationRegistry.getConfig(ReporterConfigurationImpl.class); final CoreConfigurationImpl coreConfiguration = configurationRegistry.getConfig(CoreConfigurationImpl.class); + SerializationConstants.init(coreConfiguration); SystemInfo system = new SystemInfo("x64", "localhost", null, "platform"); final ProcessInfo title = new ProcessInfo("title"); final ServiceImpl service = new ServiceImpl(); @@ -191,6 +196,25 @@ void testNoopWhenNotConnected() throws Exception { assertThat(nonConnectedReportingEventHandler.getBufferSize()).isEqualTo(0); } + @Test + void testCleanupFailureDoesNotRetainConnection() throws Exception { + reportTransaction(reportingEventHandler); + reportingEventHandler.endRequest(); + + HttpURLConnection failedConnection = mock(HttpURLConnection.class); + doThrow(new AssertionError("cleanup failed")).when(failedConnection).getInputStream(); + reportingEventHandler.connection = failedConnection; + + assertThatThrownBy(reportingEventHandler::endRequest) + .isInstanceOf(AssertionError.class) + .hasMessage("cleanup failed"); + assertThat(reportingEventHandler.connection).isNull(); + + reportTransaction(reportingEventHandler); + reportingEventHandler.endRequest(); + mockApmServer1.verify(2, postRequestedFor(urlEqualTo(INTAKE_V2_URL))); + } + @Test void testShutDown() throws Exception { reportTransaction(reportingEventHandler);