From 62cff9b52623285a2ba13397593e0dbfbab808ad Mon Sep 17 00:00:00 2001 From: Carlo Marinangeli Date: Fri, 2 Oct 2026 10:47:35 +0400 Subject: [PATCH 01/16] fix(okhttp): Don't block on response bodies of unknown length Capturing a response body for session replay peeked it up front with `Response.peekBody(MAX_NETWORK_BODY_SIZE + 1)`. OkHttp implements a peek as `request(byteCount)`, which keeps reading until that many bytes are buffered or the stream ends. A response with no Content-Length never satisfies either condition, so for server-sent events, long-poll and any chunked endpoint that stays open the interceptor never returned and the caller never received the response at all. Bodies with a known length still end on their own, so they keep the existing up-front capture. Bodies of unknown length are now wrapped in a body that copies what the application consumes into a capped buffer and reports it once no more bytes can arrive: when the stream ends, when the application closes the body, or when the cap is reached. Co-Authored-By: Claude Opus 4.8 (1M context) --- CHANGELOG.md | 4 + .../NetworkBodyCapturingResponseBody.kt | 91 +++++ .../sentry/okhttp/SentryOkHttpInterceptor.kt | 82 +++- .../NetworkBodyCapturingResponseBodyTest.kt | 229 ++++++++++++ .../SentryOkHttpInterceptorStreamingTest.kt | 353 ++++++++++++++++++ 5 files changed, 745 insertions(+), 14 deletions(-) create mode 100644 sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt create mode 100644 sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt create mode 100644 sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt diff --git a/CHANGELOG.md b/CHANGELOG.md index 4126d1706d..1c0bbde33d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,10 @@ ## Unreleased +### Fixes + +- Fix `SentryOkHttpInterceptor` hanging forever on responses whose body has no known length, such as Server-Sent Events ([#XXXX](https://github.com/getsentry/sentry-java/pull/XXXX)) + ### Features - Report the cellular network technology generation in `device.connection_effective_type`, for example `4g` or `5g` ([#6146](https://github.com/getsentry/sentry-java/pull/6146)) diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt new file mode 100644 index 0000000000..3675a7bf24 --- /dev/null +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt @@ -0,0 +1,91 @@ +package io.sentry.okhttp + +import java.util.concurrent.atomic.AtomicBoolean +import okhttp3.MediaType +import okhttp3.ResponseBody +import okio.Buffer +import okio.BufferedSource +import okio.ForwardingSource +import okio.Source +import okio.buffer + +/** + * A [ResponseBody] that copies the bytes the application consumes into a capped buffer. + * + * Capturing a body by peeking at it up front does not work for every response. OkHttp implements + * [okhttp3.Response.peekBody] as `request(byteCount)`, which keeps reading until the requested + * number of bytes is buffered or the stream ends. A response of unknown length may never end — a + * server-sent-events stream, a long-poll, a chunked endpoint that stays open — so the peek never + * returns and the thread that called `execute()` never gets the response at all. + * + * This body instead captures what is actually consumed. It forwards every byte to the application + * untouched and hands the captured bytes to [onCaptured] as soon as no more bytes can arrive, which + * is when the stream ends, when the application closes the body, or when the cap is reached. + * + * @param delegate the body to capture from. + * @param maxBytes the maximum number of bytes to retain; capture stops once it is reached. + * @param onCaptured invoked exactly once with the captured bytes. + */ +internal class NetworkBodyCapturingResponseBody( + private val delegate: ResponseBody, + private val maxBytes: Long, + private val onCaptured: (ByteArray) -> Unit, +) : ResponseBody() { + + private val captured = Buffer() + private val reported = AtomicBoolean(false) + private val capturingSource: BufferedSource by lazy { + CapturingSource(delegate.source()).buffer() + } + + override fun contentType(): MediaType? = delegate.contentType() + + override fun contentLength(): Long = delegate.contentLength() + + override fun source(): BufferedSource = capturingSource + + override fun close() { + // Report before closing the delegate: closing it notifies listeners, which may serialize the + // breadcrumb this capture belongs to. + reportCaptured() + delegate.close() + } + + private inner class CapturingSource(source: Source) : ForwardingSource(source) { + override fun read(sink: Buffer, byteCount: Long): Long { + val sinkBefore = sink.size + val read = super.read(sink, byteCount) + + if (read > 0L) { + val reachedCap = + synchronized(captured) { + val room = maxBytes - captured.size + val toTake = minOf(room, sink.size - sinkBefore) + if (toTake > 0L) { + sink.copyTo(captured, sinkBefore, toTake) + } + captured.size >= maxBytes + } + if (reachedCap) { + reportCaptured() + } + } else if (read == -1L) { + // End of stream: nothing more will ever be captured. + reportCaptured() + } + return read + } + + override fun close() { + super.close() + reportCaptured() + } + } + + private fun reportCaptured() { + if (reported.compareAndSet(false, true)) { + val bytes = synchronized(captured) { captured.clone().readByteArray() } + onCaptured(bytes) + } + } +} diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt index 71a43590e5..04e227aae9 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt @@ -191,7 +191,8 @@ public open class SentryOkHttpInterceptor( } } - return response + // The response is returned after the finally block: a body of unknown length may be wrapped + // there so its content can be captured while it is streamed. } catch (e: IOException) { span?.apply { this.throwable = e @@ -203,19 +204,7 @@ public open class SentryOkHttpInterceptor( // this only works correctly if SentryOkHttpInterceptor is the last one in the chain okHttpEvent?.setRequest(request) - response?.let { - networkDetailData?.setResponseDetails( - it.code, - NetworkDetailCaptureUtils.createResponse( - it, - it.body?.contentLength(), - scopes.options.sessionReplay.isNetworkCaptureBodies, - { resp: Response -> resp.extractResponseBody(scopes.options.logger) }, - scopes.options.sessionReplay.networkResponseHeaders, - { resp: Response -> resp.headers.toMap() }, - ), - ) - } + response = response?.withNetworkBodyCapture(networkDetailData) // Set network details on the OkHttpEvent so it can include them in the breadcrumb hint okHttpEvent?.setNetworkDetails(networkDetailData) @@ -227,6 +216,8 @@ public open class SentryOkHttpInterceptor( sendBreadcrumb(request, code, response, startTimestamp, networkDetailData) } } + + return response!! } private fun isIgnored(): Boolean = @@ -322,6 +313,69 @@ public open class SentryOkHttpInterceptor( } } + /** + * Returns this response with a body that can be captured for replay, without blocking on streams + * that never end. + * + * A body with a known length ends on its own, so it is captured up front, as before. A body of + * unknown length may never end (server-sent events, long-poll, a chunked endpoint that stays + * open): peeking it would block the calling thread until the connection dies, so it is wrapped + * and captured while it is being consumed instead. + */ + private fun Response.withNetworkBodyCapture(networkDetailData: NetworkRequestData?): Response { + val data = networkDetailData ?: return this + val responseBody = body ?: return this + val logger = scopes.options.logger + val captureBodies = scopes.options.sessionReplay.isNetworkCaptureBodies + val responseHeaders = scopes.options.sessionReplay.networkResponseHeaders + + if (!captureBodies || responseBody.contentLength() >= 0) { + data.setResponseDetails( + code, + NetworkDetailCaptureUtils.createResponse( + this, + responseBody.contentLength(), + captureBodies, + { resp: Response -> resp.extractResponseBody(logger) }, + responseHeaders, + { resp: Response -> resp.headers.toMap() }, + ), + ) + return this + } + + val maxBodySize = SentryReplayOptions.MAX_NETWORK_BODY_SIZE + val contentType = responseBody.contentType() + val contentTypeString = contentType?.toString() + val charset = contentType?.charset(Charsets.UTF_8)?.name() ?: "UTF-8" + + return newBuilder() + .body( + NetworkBodyCapturingResponseBody(responseBody, maxBodySize.toLong()) { capturedBytes -> + data.setResponseDetails( + code, + NetworkDetailCaptureUtils.createResponse( + this, + capturedBytes.size.toLong(), + true, + { + NetworkBodyParser.fromBytes( + capturedBytes, + contentTypeString, + charset, + maxBodySize, + logger, + ) + }, + responseHeaders, + { resp: Response -> resp.headers.toMap() }, + ), + ) + } + ) + .build() + } + /** Extracts the body content from an OkHttp Response safely */ private fun Response.extractResponseBody(logger: ILogger): NetworkBody? { return body?.let { responseBody -> diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt new file mode 100644 index 0000000000..ac750c4b69 --- /dev/null +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt @@ -0,0 +1,229 @@ +package io.sentry.okhttp + +import java.io.IOException +import kotlin.test.Test +import kotlin.test.assertContentEquals +import kotlin.test.assertEquals +import kotlin.test.assertFailsWith +import kotlin.test.assertFalse +import kotlin.test.assertTrue +import okhttp3.MediaType +import okhttp3.MediaType.Companion.toMediaType +import okhttp3.ResponseBody +import okio.Buffer +import okio.BufferedSource +import okio.Source +import okio.Timeout +import okio.buffer + +class NetworkBodyCapturingResponseBodyTest { + + /** Emits one chunk per read, then EOF. Records how many reads happened. */ + private class ChunkedSource(private val chunks: List) : Source { + var reads = 0 + private set + + private var index = 0 + private var closed = false + + override fun read(sink: Buffer, byteCount: Long): Long { + reads++ + if (index >= chunks.size) return -1L + val chunk = chunks[index++] + sink.write(chunk) + return chunk.size.toLong() + } + + override fun timeout(): Timeout = Timeout.NONE + + override fun close() { + closed = true + } + + val isClosed: Boolean + get() = closed + } + + private class FailingSource(private val bytesBeforeFailure: Int) : Source { + private var written = 0 + + override fun read(sink: Buffer, byteCount: Long): Long { + if (written >= bytesBeforeFailure) throw IOException("connection reset") + sink.write(ByteArray(bytesBeforeFailure) { 'x'.code.toByte() }) + written += bytesBeforeFailure + return bytesBeforeFailure.toLong() + } + + override fun timeout(): Timeout = Timeout.NONE + + override fun close() {} + } + + private fun bodyOf( + source: Source, + type: String? = "text/plain", + length: Long = -1L, + ): ResponseBody = + object : ResponseBody() { + override fun contentType(): MediaType? = type?.toMediaType() + + override fun contentLength(): Long = length + + override fun source(): BufferedSource = source.buffer() + } + + private fun capture( + source: Source, + maxBytes: Long, + type: String? = "text/plain", + length: Long = -1L, + ): Pair> { + val captured = mutableListOf() + val wrapper = + NetworkBodyCapturingResponseBody(bodyOf(source, type, length), maxBytes) { + captured.add(it) + } + return wrapper to captured + } + + @Test + fun `does not read anything before the application does`() { + val source = ChunkedSource(listOf("hello".toByteArray())) + val (wrapper, captured) = capture(source, 1024) + + assertEquals(0, source.reads, "the wrapper must not read the body up front") + assertTrue(captured.isEmpty(), "nothing can be captured before the application reads") + assertEquals("text/plain", wrapper.contentType()?.toString(), "content type is delegated") + } + + @Test + fun `forwards every byte to the application unchanged`() { + val payload = "data: event-0\n\ndata: event-1\n\n".toByteArray() + val (wrapper, _) = capture(ChunkedSource(listOf(payload)), 1024) + + assertContentEquals(payload, wrapper.source().readByteArray()) + } + + @Test + fun `captures the whole body when it is consumed and ends`() { + val (wrapper, captured) = + capture( + ChunkedSource(listOf("hello ".toByteArray(), "world".toByteArray())), + 1024, + ) + + wrapper.source().readByteArray() + + assertEquals(1, captured.size, "the capture must be reported exactly once") + assertEquals("hello world", captured.single()?.decodeToString()) + } + + @Test + fun `passes small events through as they arrive and captures them on close`() { + val events = (0 until 3).map { "data: event-$it\n\n".toByteArray() } + val source = ChunkedSource(events) + val (wrapper, captured) = capture(source, 1024) + + val application = wrapper.source() + for (expected in events) { + assertContentEquals(expected, application.readByteArray(expected.size.toLong())) + } + + assertTrue(captured.isEmpty(), "an open stream has nothing final to report yet") + + wrapper.close() + + assertEquals(1, captured.size) + assertEquals( + "data: event-0\n\ndata: event-1\n\ndata: event-2\n\n", + captured.single()?.decodeToString(), + ) + } + + @Test + fun `captures what was read when the body is closed early`() { + val (wrapper, captured) = + capture(ChunkedSource(listOf("first ".toByteArray(), "second".toByteArray())), 1024) + + val application = wrapper.source() + assertEquals("first ", application.readUtf8(6)) + + wrapper.close() + + assertEquals(1, captured.size) + assertEquals("first ", captured.single()?.decodeToString()) + } + + @Test + fun `reports an empty capture for an empty body`() { + val (wrapper, captured) = capture(ChunkedSource(emptyList()), 1024) + + wrapper.source().readByteArray() + wrapper.close() + + assertEquals(1, captured.size) + assertEquals(0, captured.single()?.size) + } + + @Test + fun `never captures more than the cap but still delivers the whole body`() { + val payload = "0123456789abcdefghij".toByteArray() + val (wrapper, captured) = capture(ChunkedSource(listOf(payload)), 8) + + assertContentEquals( + payload, + wrapper.source().readByteArray(), + "the application is not truncated", + ) + + assertEquals(1, captured.size, "reaching the cap is final, so it is reported once") + assertEquals("01234567", captured.single()?.decodeToString()) + } + + @Test + fun `reports only once when the body ends and is then closed`() { + val (wrapper, captured) = capture(ChunkedSource(listOf("done".toByteArray())), 1024) + + wrapper.source().readByteArray() + wrapper.close() + wrapper.source().close() + + assertEquals(1, captured.size, "the capture is reported exactly once") + } + + @Test + fun `delegates content type and content length`() { + val type = "text/event-stream".toMediaType() + val delegate = + object : ResponseBody() { + override fun contentType(): MediaType? = type + + override fun contentLength(): Long = -1L + + override fun source(): BufferedSource = ChunkedSource(emptyList()).buffer() + } + + val wrapper = NetworkBodyCapturingResponseBody(delegate, 16) {} + + assertEquals(type, wrapper.contentType()) + assertEquals(-1L, wrapper.contentLength()) + } + + @Test + fun `propagates read failures to the application`() { + val (wrapper, captured) = capture(FailingSource(bytesBeforeFailure = 4), 1024) + + assertFailsWith { wrapper.source().readByteArray() } + assertTrue(captured.isEmpty(), "a failed stream has no final capture") + } + + @Test + fun `closes the delegate body`() { + val source = ChunkedSource(listOf("x".toByteArray())) + val (wrapper, _) = capture(source, 1024) + + assertFalse(source.isClosed) + wrapper.close() + assertTrue(source.isClosed, "the application must still be able to release the connection") + } +} diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt new file mode 100644 index 0000000000..e79a1e0d3e --- /dev/null +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt @@ -0,0 +1,353 @@ +package io.sentry.okhttp + +import io.sentry.Hint +import io.sentry.IScopes +import io.sentry.SentryOptions +import io.sentry.TypeCheckHint +import io.sentry.util.network.NetworkRequestData +import java.io.ByteArrayOutputStream +import java.io.Closeable +import java.io.IOException +import java.net.InetAddress +import java.net.ServerSocket +import java.net.Socket +import java.util.concurrent.CopyOnWriteArrayList +import java.util.concurrent.Executors +import java.util.concurrent.TimeUnit +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNotNull +import kotlin.test.assertNull +import kotlin.test.assertTrue +import okhttp3.Call +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.Response +import okhttp3.mockwebserver.MockResponse +import okhttp3.mockwebserver.MockWebServer +import org.mockito.kotlin.any +import org.mockito.kotlin.argumentCaptor +import org.mockito.kotlin.mock +import org.mockito.kotlin.verify +import org.mockito.kotlin.whenever + +/** + * Response bodies of unknown length: server-sent events, long-poll, chunked endpoints that stay + * open. + * + * These used to hang. Capturing the body means peeking it, and OkHttp implements a peek as "keep + * reading until the peek size is reached or the stream ends". For these responses neither happens, + * so the thread that called `execute()` never received the response at all. + */ +class SentryOkHttpInterceptorStreamingTest { + + private val scopes = mock() + private lateinit var options: SentryOptions + private lateinit var sut: OkHttpClient + + private fun setUpSut(captureBodies: Boolean = true) { + options = + SentryOptions().apply { + dsn = "https://key@sentry.io/proj" + sessionReplay.setNetworkDetailAllowUrls(listOf(".*")) + sessionReplay.setNetworkCaptureBodies(captureBodies) + } + whenever(scopes.options).thenReturn(options) + sut = OkHttpClient.Builder().addInterceptor(SentryOkHttpInterceptor(scopes)).build() + } + + private fun requestTo(url: String) = Request.Builder().url(url).build() + + /** The network details the breadcrumb points at; read after the body has been consumed. */ + private fun networkDetails(): NetworkRequestData { + val hint = argumentCaptor() + verify(scopes).addBreadcrumb(any(), hint.capture()) + return assertNotNull( + hint.firstValue.getAs( + TypeCheckHint.SENTRY_REPLAY_NETWORK_DETAILS, + NetworkRequestData::class.java, + ) + ) + } + + private fun capturedBody(): String? = networkDetails().response?.body?.body as String? + + // --------------------------------------------------------------------------------------- + // the regression + // --------------------------------------------------------------------------------------- + + @Test + fun `returns the response even though the body never ends`() { + setUpSut() + EventStreamServer(listOf("data: event-0\n\n"), closeStream = false).use { stream -> + val executor = Executors.newSingleThreadExecutor { r -> Thread(r).apply { isDaemon = true } } + val call: Call = sut.newCall(requestTo(stream.url)) + try { + val response = executor.submit { call.execute() }.get(5, TimeUnit.SECONDS) + assertEquals(200, response.code) + response.close() + } finally { + executor.shutdownNow() + } + } + } + + @Test + fun `delivers each small event to the application as it arrives`() { + setUpSut() + val events = (0 until 3).map { "data: event-$it\n\n" } + EventStreamServer(events, closeStream = true).use { stream -> + sut.newCall(requestTo(stream.url)).execute().use { response -> + val source = assertNotNull(response.body).source() + for (event in events) { + assertEquals( + event, + source.readUtf8(event.length.toLong()), + "the application must see every event", + ) + } + assertTrue(source.exhausted(), "the stream ended cleanly") + } + } + } + + // --------------------------------------------------------------------------------------- + // what ends up captured + // --------------------------------------------------------------------------------------- + + @Test + fun `captures a streamed body once the stream ends`() { + setUpSut() + val events = listOf("data: event-0\n\n", "data: event-1\n\n") + EventStreamServer(events, closeStream = true).use { stream -> + sut.newCall(requestTo(stream.url)).execute().use { it.body?.string() } + val details = networkDetails() + assertEquals(200, details.statusCode) + assertEquals(events.joinToString(""), capturedBody()) + assertEquals(events.joinToString("").length.toLong(), details.responseBodySize) + } + } + + @Test + fun `captures a streamed body when the application closes it before the stream ends`() { + setUpSut() + EventStreamServer(listOf("data: event-0\n\n", "data: event-1\n\n"), closeStream = false).use { + stream -> + sut.newCall(requestTo(stream.url)).execute().use { response -> + assertEquals("data: event-0\n\n", assertNotNull(response.body).source().readUtf8(15)) + } + + assertEquals("data: event-0\n\n", capturedBody()) + } + } + + @Test + fun `records the response of a streamed body the application never reads`() { + setUpSut() + EventStreamServer(listOf("data: event-0\n\n"), closeStream = true).use { stream -> + sut.newCall(requestTo(stream.url)).execute().close() + + val details = networkDetails() + assertEquals(200, details.statusCode) + assertNull( + capturedBody(), + "an unread stream is never peeked, so there is nothing to report; the response is still recorded", + ) + } + } + + @Test + fun `captures a streamed error response body`() { + setUpSut() + EventStreamServer(listOf("data: boom\n\n"), closeStream = true, statusCode = 500).use { stream + -> + sut.newCall(requestTo(stream.url)).execute().use { it.body?.string() } + + val details = networkDetails() + assertEquals(500, details.statusCode) + assertEquals("data: boom\n\n", capturedBody()) + } + } + + @Test + fun `caps the captured body without truncating the body the application reads`() { + setUpSut() + val events = (0 until 200).map { "data: " + "x".repeat(1000) + "\n\n" } + val whole = events.joinToString("") + EventStreamServer(events, closeStream = true).use { stream -> + val received = sut.newCall(requestTo(stream.url)).execute().use { it.body?.string() } + + assertEquals(whole.length, received?.length, "the application must receive the whole stream") + val captured = assertNotNull(capturedBody()) + assertTrue( + captured.length <= io.sentry.SentryReplayOptions.MAX_NETWORK_BODY_SIZE, + "the capture must stop at the cap, was ${captured.length}", + ) + assertEquals(whole.substring(0, captured.length), captured) + } + } + + @Test + fun `does not capture bodies when network body capture is turned off`() { + setUpSut(captureBodies = false) + EventStreamServer(listOf("data: event-0\n\n"), closeStream = true).use { stream -> + sut.newCall(requestTo(stream.url)).execute().use { it.body?.string() } + + assertEquals(200, networkDetails().statusCode) + assertNull(capturedBody()) + } + } + + // --------------------------------------------------------------------------------------- + // responses with a known length keep the existing behaviour + // --------------------------------------------------------------------------------------- + + @Test + fun `still captures a response with a known length up front`() { + setUpSut() + MockWebServer().use { server -> + server.enqueue(MockResponse().setBody("response body").setResponseCode(200)) + + sut.newCall(requestTo(server.url("/hello").toString())).execute().close() + + assertEquals("response body", capturedBody()) + } + } + + @Test + fun `still captures an error response with a known length the application never reads`() { + setUpSut() + MockWebServer().use { server -> + server.enqueue(MockResponse().setBody("failure").setResponseCode(500)) + + sut.newCall(requestTo(server.url("/hello").toString())).execute().close() + + assertEquals(500, networkDetails().statusCode) + assertEquals("failure", capturedBody()) + } + } + + @Test + fun `handles a zero length body`() { + setUpSut() + MockWebServer().use { server -> + server.enqueue(MockResponse().setResponseCode(204)) + + sut.newCall(requestTo(server.url("/hello").toString())).execute().use { response -> + assertEquals(204, response.code) + assertEquals(0, assertNotNull(response.body).contentLength()) + } + + assertEquals(204, networkDetails().statusCode) + assertNull(capturedBody()) + } + } + + @Test + fun `handles a zero length streamed body`() { + setUpSut() + EventStreamServer(emptyList(), closeStream = true).use { stream -> + sut.newCall(requestTo(stream.url)).execute().use { response -> + assertEquals("", response.body?.string()) + } + + assertEquals(200, networkDetails().statusCode) + assertNull(capturedBody()) + } + } + + @Test + fun `keeps the connection usable after a streamed response is closed`() { + setUpSut() + EventStreamServer(listOf("data: event-0\n\n"), closeStream = true).use { stream -> + // a second call on the same client proves the connection was released properly + repeat(3) { + sut.newCall(requestTo(stream.url)).execute().use { assertEquals(200, it.code) } + } + } + } + + /** + * Minimal HTTP/1.1 origin that streams one chunk per event and can leave the stream open, so the + * response body never ends. + */ + private class EventStreamServer( + private val events: List, + private val closeStream: Boolean, + private val statusCode: Int = 200, + ) : Closeable { + private val server = ServerSocket(0, 8, InetAddress.getLoopbackAddress()) + private val connections = CopyOnWriteArrayList() + + init { + Thread({ acceptLoop() }, "event-stream-server").apply { isDaemon = true }.start() + } + + val url: String + get() = "http://127.0.0.1:${server.localPort}/events" + + private fun acceptLoop() { + while (!server.isClosed) { + val socket = + try { + server.accept() + } catch (e: IOException) { + return + } + connections.add(socket) + Thread({ serve(socket) }, "event-stream-connection").apply { isDaemon = true }.start() + } + } + + private fun serve(socket: Socket) { + try { + socket.use { + readRequestHead(it.getInputStream()) + val output = it.getOutputStream() + output.write( + ("HTTP/1.1 $statusCode OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nTransfer-Encoding: chunked\r\n\r\n") + .toByteArray() + ) + output.flush() + for (event in events) { + val bytes = event.toByteArray() + output.write("${bytes.size.toString(16)}\r\n".toByteArray()) + output.write(bytes) + output.write("\r\n".toByteArray()) + output.flush() + } + if (closeStream) { + output.write("0\r\n\r\n".toByteArray()) + output.flush() + } else { + // hold the connection open: the body stays open ended + Thread.sleep(30_000) + } + } + } catch (e: InterruptedException) { + Thread.currentThread().interrupt() + } catch (e: IOException) { + // the client went away + } + } + + private fun readRequestHead(input: java.io.InputStream) { + val head = ByteArrayOutputStream() + while (true) { + val byte = input.read() + if (byte == -1) { + return + } + head.write(byte) + if (head.size() >= 4 && head.toString(Charsets.ISO_8859_1).endsWith("\r\n\r\n")) { + return + } + } + } + + override fun close() { + runCatching { server.close() } + connections.forEach { runCatching { it.close() } } + } + } +} From 4d7c686082df6e8db4bca7c2d12b100213c418da Mon Sep 17 00:00:00 2001 From: Carlo Marinangeli Date: Fri, 2 Oct 2026 10:48:16 +0400 Subject: [PATCH 02/16] changelog: reference the pull request --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1c0bbde33d..bba665ffd1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,7 +4,7 @@ ### Fixes -- Fix `SentryOkHttpInterceptor` hanging forever on responses whose body has no known length, such as Server-Sent Events ([#XXXX](https://github.com/getsentry/sentry-java/pull/XXXX)) +- Fix `SentryOkHttpInterceptor` hanging forever on responses whose body has no known length, such as Server-Sent Events ([#1](https://github.com/getsentry/sentry-java/pull/1)) ### Features From 554d12149b6d00580422ef3ff444cce9688f16be Mon Sep 17 00:00:00 2001 From: Carlo Marinangeli Date: Fri, 2 Oct 2026 10:51:34 +0400 Subject: [PATCH 03/16] fix(okhttp): Satisfy detekt for the streaming body capture Split the capture wrapping into its own function to keep a single return, and suppress TooManyFunctions on the interceptor as done elsewhere in this module. --- .../sentry/okhttp/SentryOkHttpInterceptor.kt | 48 ++++++++++++------- .../NetworkBodyCapturingResponseBodyTest.kt | 14 +++++- .../SentryOkHttpInterceptorStreamingTest.kt | 12 +++-- 3 files changed, 51 insertions(+), 23 deletions(-) diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt index 04e227aae9..f926fec022 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt @@ -34,6 +34,7 @@ import okhttp3.Interceptor import okhttp3.Request import okhttp3.RequestBody.Companion.toRequestBody import okhttp3.Response +import okhttp3.ResponseBody import org.jetbrains.annotations.VisibleForTesting /** @@ -50,6 +51,7 @@ import org.jetbrains.annotations.VisibleForTesting * @param failedRequestTargets The SDK will only capture HTTP Client errors if the HTTP Request URL * is a match for any of the defined targets. */ +@Suppress("TooManyFunctions") // one function per concern of the interception pipeline public open class SentryOkHttpInterceptor( private val scopes: IScopes = ScopesAdapter.getInstance(), private val beforeSpan: BeforeSpanCallback? = null, @@ -323,27 +325,39 @@ public open class SentryOkHttpInterceptor( * and captured while it is being consumed instead. */ private fun Response.withNetworkBodyCapture(networkDetailData: NetworkRequestData?): Response { - val data = networkDetailData ?: return this - val responseBody = body ?: return this - val logger = scopes.options.logger + val responseBody = body val captureBodies = scopes.options.sessionReplay.isNetworkCaptureBodies val responseHeaders = scopes.options.sessionReplay.networkResponseHeaders + val logger = scopes.options.logger - if (!captureBodies || responseBody.contentLength() >= 0) { - data.setResponseDetails( - code, - NetworkDetailCaptureUtils.createResponse( - this, - responseBody.contentLength(), - captureBodies, - { resp: Response -> resp.extractResponseBody(logger) }, - responseHeaders, - { resp: Response -> resp.headers.toMap() }, - ), - ) - return this + return when { + networkDetailData == null || responseBody == null -> this + // a body with a known length ends on its own, so capturing it up front is bounded + !captureBodies || responseBody.contentLength() >= 0 -> { + networkDetailData.setResponseDetails( + code, + NetworkDetailCaptureUtils.createResponse( + this, + responseBody.contentLength(), + captureBodies, + { resp: Response -> resp.extractResponseBody(logger) }, + responseHeaders, + { resp: Response -> resp.headers.toMap() }, + ), + ) + this + } + else -> wrapForCapturedStreamingBody(networkDetailData, responseBody, responseHeaders, logger) } + } + /** Wraps a body that may never end, capturing the bytes as the application consumes them. */ + private fun Response.wrapForCapturedStreamingBody( + networkDetailData: NetworkRequestData, + responseBody: ResponseBody, + responseHeaders: List, + logger: ILogger, + ): Response { val maxBodySize = SentryReplayOptions.MAX_NETWORK_BODY_SIZE val contentType = responseBody.contentType() val contentTypeString = contentType?.toString() @@ -352,7 +366,7 @@ public open class SentryOkHttpInterceptor( return newBuilder() .body( NetworkBodyCapturingResponseBody(responseBody, maxBodySize.toLong()) { capturedBytes -> - data.setResponseDetails( + networkDetailData.setResponseDetails( code, NetworkDetailCaptureUtils.createResponse( this, diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt index ac750c4b69..98695d9278 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt @@ -45,6 +45,9 @@ class NetworkBodyCapturingResponseBodyTest { } private class FailingSource(private val bytesBeforeFailure: Int) : Source { + var isClosed = false + private set + private var written = 0 override fun read(sink: Buffer, byteCount: Long): Long { @@ -56,7 +59,9 @@ class NetworkBodyCapturingResponseBodyTest { override fun timeout(): Timeout = Timeout.NONE - override fun close() {} + override fun close() { + isClosed = true + } } private fun bodyOf( @@ -211,10 +216,15 @@ class NetworkBodyCapturingResponseBodyTest { @Test fun `propagates read failures to the application`() { - val (wrapper, captured) = capture(FailingSource(bytesBeforeFailure = 4), 1024) + val source = FailingSource(bytesBeforeFailure = 4) + val captured = mutableListOf() + val wrapper = NetworkBodyCapturingResponseBody(bodyOf(source), 1024) { captured.add(it) } assertFailsWith { wrapper.source().readByteArray() } assertTrue(captured.isEmpty(), "a failed stream has no final capture") + + wrapper.close() + assertTrue(source.isClosed, "the connection must still be releasable after a failure") } @Test diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt index e79a1e0d3e..4bcfeaeaf8 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt @@ -286,6 +286,7 @@ class SentryOkHttpInterceptorStreamingTest { val url: String get() = "http://127.0.0.1:${server.localPort}/events" + @Suppress("SwallowedException") // the server socket is closed when the test finishes private fun acceptLoop() { while (!server.isClosed) { val socket = @@ -299,15 +300,18 @@ class SentryOkHttpInterceptorStreamingTest { } } + @Suppress("SwallowedException") // a client that hangs up mid-stream is normal here private fun serve(socket: Socket) { try { socket.use { readRequestHead(it.getInputStream()) val output = it.getOutputStream() - output.write( - ("HTTP/1.1 $statusCode OK\r\nContent-Type: text/event-stream\r\nCache-Control: no-cache\r\nTransfer-Encoding: chunked\r\n\r\n") - .toByteArray() - ) + val headers = + "HTTP/1.1 $statusCode OK\r\n" + + "Content-Type: text/event-stream\r\n" + + "Cache-Control: no-cache\r\n" + + "Transfer-Encoding: chunked\r\n\r\n" + output.write(headers.toByteArray()) output.flush() for (event in events) { val bytes = event.toByteArray() From 0d5e6dffd9470020a842990e40248908305a8c0b Mon Sep 17 00:00:00 2001 From: Carlo Marinangeli Date: Fri, 2 Oct 2026 11:25:11 +0400 Subject: [PATCH 04/16] refactor(okhttp): Drop unnecessary locking from the capturing source A response body has a single consumer, and neither Http1ExchangeCodec.cancel() nor Http2ExchangeCodec.cancel() closes the body, so nothing reaches this class from another thread. The AtomicBoolean already guarantees the capture is reported exactly once. --- .../NetworkBodyCapturingResponseBody.kt | 27 ++++++++++--------- 1 file changed, 14 insertions(+), 13 deletions(-) diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt index 3675a7bf24..ddefd29972 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt @@ -51,26 +51,28 @@ internal class NetworkBodyCapturingResponseBody( delegate.close() } + /** + * Copies the bytes passing through into [captured]. + * + * A response body has a single consumer, as OkHttp requires of one, so the buffer is only ever + * touched by the thread reading it and needs no locking. + */ private inner class CapturingSource(source: Source) : ForwardingSource(source) { override fun read(sink: Buffer, byteCount: Long): Long { val sinkBefore = sink.size val read = super.read(sink, byteCount) if (read > 0L) { - val reachedCap = - synchronized(captured) { - val room = maxBytes - captured.size - val toTake = minOf(room, sink.size - sinkBefore) - if (toTake > 0L) { - sink.copyTo(captured, sinkBefore, toTake) - } - captured.size >= maxBytes - } - if (reachedCap) { + val toTake = minOf(maxBytes - captured.size, sink.size - sinkBefore) + if (toTake > 0L) { + sink.copyTo(captured, sinkBefore, toTake) + } + if (captured.size >= maxBytes) { + // the capture can never grow again, so this is the moment it becomes final reportCaptured() } } else if (read == -1L) { - // End of stream: nothing more will ever be captured. + // end of stream: nothing more can ever arrive reportCaptured() } return read @@ -84,8 +86,7 @@ internal class NetworkBodyCapturingResponseBody( private fun reportCaptured() { if (reported.compareAndSet(false, true)) { - val bytes = synchronized(captured) { captured.clone().readByteArray() } - onCaptured(bytes) + onCaptured(captured.clone().readByteArray()) } } } From a982ceb1eab6a1ba6ec865ebfc75d3a1110436c5 Mon Sep 17 00:00:00 2001 From: Carlo Marinangeli Date: Fri, 2 Oct 2026 11:34:50 +0400 Subject: [PATCH 05/16] fix: Publish late response details safely A streamed body is only known once it has been consumed, which can be after the NetworkRequestData carrying it was handed to the scope, so the replay thread can read it while it is still being written. Keeping the three response values behind one volatile reference means a reader sees either nothing or the complete set. --- .../util/network/NetworkRequestData.java | 44 ++++++++--- .../util/network/NetworkRequestDataTest.kt | 74 +++++++++++++++++++ 2 files changed, 106 insertions(+), 12 deletions(-) create mode 100644 sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt diff --git a/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java b/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java index 997e001a35..bebf6189c8 100644 --- a/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java +++ b/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java @@ -13,11 +13,14 @@ @ApiStatus.Internal public final class NetworkRequestData { private @Nullable final String method; - private @Nullable Integer statusCode; private @Nullable Long requestBodySize; - private @Nullable Long responseBodySize; private @Nullable ReplayNetworkRequestOrResponse request; - private @Nullable ReplayNetworkRequestOrResponse response; + + // The response can be filled in after this instance was handed to the scope: an integration that + // captures a streamed body only knows it once the stream has been consumed. Keeping the three + // response values behind one volatile reference means a reader on the replay thread sees either + // nothing or the complete set, never a mix of the two. + private volatile @Nullable ResponseDetails responseDetails; public NetworkRequestData(@Nullable final String method) { this.method = method; @@ -28,7 +31,8 @@ public NetworkRequestData(@Nullable final String method) { } public @Nullable Integer getStatusCode() { - return statusCode; + final ResponseDetails details = responseDetails; + return details == null ? null : details.statusCode; } public @Nullable Long getRequestBodySize() { @@ -36,7 +40,8 @@ public NetworkRequestData(@Nullable final String method) { } public @Nullable Long getResponseBodySize() { - return responseBodySize; + final ResponseDetails details = responseDetails; + return details == null ? null : details.bodySize; } public @Nullable ReplayNetworkRequestOrResponse getRequest() { @@ -44,7 +49,8 @@ public NetworkRequestData(@Nullable final String method) { } public @Nullable ReplayNetworkRequestOrResponse getResponse() { - return response; + final ResponseDetails details = responseDetails; + return details == null ? null : details.response; } /** @@ -62,9 +68,7 @@ public void setRequestDetails(@NotNull final ReplayNetworkRequestOrResponse requ */ public void setResponseDetails( final int statusCode, @NotNull final ReplayNetworkRequestOrResponse responseData) { - this.statusCode = statusCode; - this.response = responseData; - this.responseBodySize = responseData.getSize(); + this.responseDetails = new ResponseDetails(statusCode, responseData.getSize(), responseData); } @Override @@ -74,15 +78,31 @@ public String toString() { + method + '\'' + ", statusCode=" - + statusCode + + getStatusCode() + ", requestBodySize=" + requestBodySize + ", responseBodySize=" - + responseBodySize + + getResponseBodySize() + ", request=" + request + ", response=" - + response + + getResponse() + '}'; } + + /** Immutable, so publishing one reference publishes all three values. */ + private static final class ResponseDetails { + private final int statusCode; + private final @Nullable Long bodySize; + private final @NotNull ReplayNetworkRequestOrResponse response; + + ResponseDetails( + final int statusCode, + final @Nullable Long bodySize, + final @NotNull ReplayNetworkRequestOrResponse response) { + this.statusCode = statusCode; + this.bodySize = bodySize; + this.response = response; + } + } } diff --git a/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt b/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt new file mode 100644 index 0000000000..89079449ff --- /dev/null +++ b/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt @@ -0,0 +1,74 @@ +package io.sentry.util.network + +import java.util.concurrent.CountDownLatch +import java.util.concurrent.atomic.AtomicBoolean +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.test.assertTrue + +class NetworkRequestDataTest { + + private fun response(statusCode: Int, size: Long): ReplayNetworkRequestOrResponse = + ReplayNetworkRequestOrResponse(size, null, emptyMap()) + + @Test + fun `response details are empty until they are set`() { + val data = NetworkRequestData("GET") + + assertNull(data.getStatusCode()) + assertNull(data.getResponse()) + assertNull(data.getResponseBodySize()) + } + + @Test + fun `response details are all readable once set`() { + val data = NetworkRequestData("GET") + + data.setResponseDetails(200, response(200, 42L)) + + assertEquals(200, data.getStatusCode()) + assertEquals(42L, data.getResponseBodySize()) + assertEquals(42L, data.getResponse()?.getSize()) + } + + @Test + fun `response details can be filled in after the instance was published`() { + // A streamed body is only known once it has been consumed, which can be after the breadcrumb + // holding this instance reached the scope. + val data = NetworkRequestData("GET") + assertNull(data.getStatusCode()) + + data.setResponseDetails(500, response(500, 7L)) + + assertEquals(500, data.getStatusCode()) + assertEquals(7L, data.getResponseBodySize()) + } + + @Test + fun `a reader never observes a partially updated response`() { + val data = NetworkRequestData("GET") + val started = CountDownLatch(1) + val stop = AtomicBoolean(false) + var torn = false + + val reader = Thread { + started.await() + while (!stop.get()) { + // status, size and response must move together: a reader that saw the status code before + // the response would report an inconsistent request in the replay + if (data.getStatusCode() != null && data.getResponse() == null) { + torn = true + } + } + } + + reader.start() + started.countDown() + repeat(200_000) { data.setResponseDetails(200, response(200, it.toLong())) } + stop.set(true) + reader.join(5_000) + + assertTrue(!torn, "observed statusCode without a matching response") + } +} From 1c6174f4a2dae1b63c7545d3b4b070ab294c0b02 Mon Sep 17 00:00:00 2001 From: Carlo Marinangeli Date: Fri, 2 Oct 2026 11:37:05 +0400 Subject: [PATCH 06/16] test: Cover late-filled network request data --- .../util/network/NetworkRequestDataTest.kt | 34 +++++-------------- 1 file changed, 9 insertions(+), 25 deletions(-) diff --git a/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt b/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt index 89079449ff..3665457637 100644 --- a/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt +++ b/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt @@ -1,11 +1,8 @@ package io.sentry.util.network -import java.util.concurrent.CountDownLatch -import java.util.concurrent.atomic.AtomicBoolean import kotlin.test.Test import kotlin.test.assertEquals import kotlin.test.assertNull -import kotlin.test.assertTrue class NetworkRequestDataTest { @@ -35,40 +32,27 @@ class NetworkRequestDataTest { @Test fun `response details can be filled in after the instance was published`() { // A streamed body is only known once it has been consumed, which can be after the breadcrumb - // holding this instance reached the scope. + // holding this instance already reached the scope, so the values must stay consistent for a + // reader that arrives late. val data = NetworkRequestData("GET") assertNull(data.getStatusCode()) + assertNull(data.getResponse()) data.setResponseDetails(500, response(500, 7L)) assertEquals(500, data.getStatusCode()) assertEquals(7L, data.getResponseBodySize()) + assertEquals(7L, data.getResponse()?.getSize()) } @Test - fun `a reader never observes a partially updated response`() { + fun `request details are unaffected`() { val data = NetworkRequestData("GET") - val started = CountDownLatch(1) - val stop = AtomicBoolean(false) - var torn = false - - val reader = Thread { - started.await() - while (!stop.get()) { - // status, size and response must move together: a reader that saw the status code before - // the response would report an inconsistent request in the replay - if (data.getStatusCode() != null && data.getResponse() == null) { - torn = true - } - } - } + val request = ReplayNetworkRequestOrResponse(11L, null, mapOf("Accept" to "application/json")) - reader.start() - started.countDown() - repeat(200_000) { data.setResponseDetails(200, response(200, it.toLong())) } - stop.set(true) - reader.join(5_000) + data.setRequestDetails(request) - assertTrue(!torn, "observed statusCode without a matching response") + assertEquals(11L, data.getRequestBodySize()) + assertEquals(request, data.getRequest()) } } From 9377594db7530b5f7e1f8d74eac4e6e872220eba Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Tue, 6 Oct 2026 14:24:35 +0200 Subject: [PATCH 07/16] fix(okhttp): Harden the capture of response bodies with no known length --- .../DefaultReplayBreadcrumbConverterTest.kt | 42 ++++--- .../NetworkBodyCapturingResponseBody.kt | 61 +++++---- .../sentry/okhttp/SentryOkHttpInterceptor.kt | 118 +++++++++++------- .../NetworkBodyCapturingResponseBodyTest.kt | 31 +++++ ...ttpInterceptorUnknownContentLengthTest.kt} | 71 +++++++++-- sentry/api/sentry.api | 7 +- .../util/network/NetworkRequestData.java | 51 ++++---- .../util/network/NetworkRequestDataTest.kt | 14 ++- 8 files changed, 276 insertions(+), 119 deletions(-) rename sentry-okhttp/src/test/java/io/sentry/okhttp/{SentryOkHttpInterceptorStreamingTest.kt => SentryOkHttpInterceptorUnknownContentLengthTest.kt} (81%) diff --git a/sentry-android-replay/src/test/java/io/sentry/android/replay/DefaultReplayBreadcrumbConverterTest.kt b/sentry-android-replay/src/test/java/io/sentry/android/replay/DefaultReplayBreadcrumbConverterTest.kt index 749d349669..ee5944283c 100644 --- a/sentry-android-replay/src/test/java/io/sentry/android/replay/DefaultReplayBreadcrumbConverterTest.kt +++ b/sentry-android-replay/src/test/java/io/sentry/android/replay/DefaultReplayBreadcrumbConverterTest.kt @@ -429,12 +429,14 @@ class DefaultReplayBreadcrumbConverterTest { ) ) fakeOkHttpNetworkDetails.setResponseDetails( - 200, - ReplayNetworkRequestOrResponse( - 500L, - NetworkBody(mapOf("status" to "success", "message" to "OK")), - mapOf("Content-Type" to "text/plain"), - ), + NetworkRequestData.ResponseDetails( + 200, + ReplayNetworkRequestOrResponse( + 500L, + NetworkBody(mapOf("status" to "success", "message" to "OK")), + mapOf("Content-Type" to "text/plain"), + ), + ) ) val hintWithFakeOKHttpNetworkDetails = Hint() hintWithFakeOKHttpNetworkDetails.set(SENTRY_REPLAY_NETWORK_DETAILS, fakeOkHttpNetworkDetails) @@ -491,12 +493,14 @@ class DefaultReplayBreadcrumbConverterTest { ) ) fakeOkHttpNetworkDetails.setResponseDetails( - 404, - ReplayNetworkRequestOrResponse( - 550L, - NetworkBody(mapOf("status" to "success", "message" to "OK")), - mapOf("Content-Type" to "text/plain"), - ), + NetworkRequestData.ResponseDetails( + 404, + ReplayNetworkRequestOrResponse( + 550L, + NetworkBody(mapOf("status" to "success", "message" to "OK")), + mapOf("Content-Type" to "text/plain"), + ), + ) ) val hintWithFakeOKHttpNetworkDetails = Hint() hintWithFakeOKHttpNetworkDetails.set(SENTRY_REPLAY_NETWORK_DETAILS, fakeOkHttpNetworkDetails) @@ -546,12 +550,14 @@ class DefaultReplayBreadcrumbConverterTest { ) ) networkRequestData.setResponseDetails( - 200, - ReplayNetworkRequestOrResponse( - 100L, - NetworkBody("response body content"), - mapOf("Content-Type" to "application/json"), - ), + NetworkRequestData.ResponseDetails( + 200, + ReplayNetworkRequestOrResponse( + 100L, + NetworkBody("response body content"), + mapOf("Content-Type" to "application/json"), + ), + ) ) hint.set(SENTRY_REPLAY_NETWORK_DETAILS, networkRequestData) diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt index ddefd29972..3a0149a06b 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt @@ -19,12 +19,13 @@ import okio.buffer * returns and the thread that called `execute()` never gets the response at all. * * This body instead captures what is actually consumed. It forwards every byte to the application - * untouched and hands the captured bytes to [onCaptured] as soon as no more bytes can arrive, which - * is when the stream ends, when the application closes the body, or when the cap is reached. + * untouched and hands the captured bytes to [onCaptured] as soon as the capture can no longer grow, + * which is when the stream ends, when the application closes the body, or when the cap is reached. * * @param delegate the body to capture from. * @param maxBytes the maximum number of bytes to retain; capture stops once it is reached. - * @param onCaptured invoked exactly once with the captured bytes. + * @param onCaptured invoked at most once, synchronously, on the thread that finishes the body. It + * is never invoked for a body that is abandoned without being read to the end or closed. */ internal class NetworkBodyCapturingResponseBody( private val delegate: ResponseBody, @@ -34,6 +35,13 @@ internal class NetworkBodyCapturingResponseBody( private val captured = Buffer() private val reported = AtomicBoolean(false) + + // A ResponseBody must hand out a BufferedSource, so the capturing source is buffered. That cannot + // re-introduce the blocking this class exists to avoid: BufferedSource.read(sink, byteCount) + // issues at most one segment-sized read on the source below it and returns with whatever arrived, + // so a short event is still forwarded on its own. The reads that loop until a byte count is + // reached — request, require, readByteArray() — are the application's own choice, and the capture + // never calls them. private val capturingSource: BufferedSource by lazy { CapturingSource(delegate.source()).buffer() } @@ -45,29 +53,31 @@ internal class NetworkBodyCapturingResponseBody( override fun source(): BufferedSource = capturingSource override fun close() { - // Report before closing the delegate: closing it notifies listeners, which may serialize the - // breadcrumb this capture belongs to. - reportCaptured() - delegate.close() + try { + // Report before closing the delegate: closing it notifies listeners, which may serialize the + // breadcrumb this capture belongs to. + reportCaptured() + } finally { + delegate.close() + } } - /** - * Copies the bytes passing through into [captured]. - * - * A response body has a single consumer, as OkHttp requires of one, so the buffer is only ever - * touched by the thread reading it and needs no locking. - */ + /** Copies the bytes passing through into [captured]. */ private inner class CapturingSource(source: Source) : ForwardingSource(source) { override fun read(sink: Buffer, byteCount: Long): Long { val sinkBefore = sink.size val read = super.read(sink, byteCount) if (read > 0L) { - val toTake = minOf(maxBytes - captured.size, sink.size - sinkBefore) - if (toTake > 0L) { - sink.copyTo(captured, sinkBefore, toTake) - } - if (captured.size >= maxBytes) { + val capFull = + synchronized(captured) { + val toTake = minOf(maxBytes - captured.size, sink.size - sinkBefore) + if (toTake > 0L) { + sink.copyTo(captured, sinkBefore, toTake) + } + captured.size >= maxBytes + } + if (capFull) { // the capture can never grow again, so this is the moment it becomes final reportCaptured() } @@ -79,14 +89,23 @@ internal class NetworkBodyCapturingResponseBody( } override fun close() { - super.close() - reportCaptured() + try { + reportCaptured() + } finally { + super.close() + } } } + /** + * The consuming thread writes [captured] from [CapturingSource.read] while [close] may be called + * by another thread — cancelling a stream from elsewhere is ordinary use — and [Buffer] is not + * thread-safe, so both accesses are guarded. [reported] keeps the callback to a single + * invocation. + */ private fun reportCaptured() { if (reported.compareAndSet(false, true)) { - onCaptured(captured.clone().readByteArray()) + onCaptured(synchronized(captured) { captured.clone().readByteArray() }) } } } diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt index f926fec022..7a512241cf 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt @@ -173,28 +173,32 @@ public open class SentryOkHttpInterceptor( ) request = requestBuilder.build() - response = chain.proceed(request) - code = response.code + val rawResponse = chain.proceed(request) + response = rawResponse + code = rawResponse.code span?.setData(SpanDataConvention.HTTP_STATUS_CODE_KEY, code) span?.status = SpanStatus.fromHttpStatusCode(code) // OkHttp errors (4xx, 5xx) don't throw, so it's safe to call within this block. // breadcrumbs are added on the finally block because we'd like to know if the device // had an unstable connection or something similar - if (shouldCaptureClientError(request, response)) { + if (shouldCaptureClientError(request, rawResponse)) { // If we capture the client error directly, it could be associated with the // currently running span by the backend. In case the listener is in use, that is // an inner span. So, if the listener is in use, we let it capture the client // error, to shown it in the http root call span in the dashboard. if (isFromEventListener && okHttpEvent != null) { - okHttpEvent.setClientErrorResponse(response) + okHttpEvent.setClientErrorResponse(rawResponse) } else { - SentryOkHttpUtils.captureClientError(scopes, request, response) + SentryOkHttpUtils.captureClientError(scopes, request, rawResponse) } } - // The response is returned after the finally block: a body of unknown length may be wrapped - // there so its content can be captured while it is streamed. + // Wrapping comes last, after the client error was captured from the untouched response. The + // returned value is computed here, so the finally block below only reads it. + val capturingResponse = rawResponse.withNetworkBodyCapture(networkDetailData) + response = capturingResponse + return capturingResponse } catch (e: IOException) { span?.apply { this.throwable = e @@ -206,8 +210,6 @@ public open class SentryOkHttpInterceptor( // this only works correctly if SentryOkHttpInterceptor is the last one in the chain okHttpEvent?.setRequest(request) - response = response?.withNetworkBodyCapture(networkDetailData) - // Set network details on the OkHttpEvent so it can include them in the breadcrumb hint okHttpEvent?.setNetworkDetails(networkDetailData) @@ -218,8 +220,6 @@ public open class SentryOkHttpInterceptor( sendBreadcrumb(request, code, response, startTimestamp, networkDetailData) } } - - return response!! } private fun isIgnored(): Boolean = @@ -335,56 +335,88 @@ public open class SentryOkHttpInterceptor( // a body with a known length ends on its own, so capturing it up front is bounded !captureBodies || responseBody.contentLength() >= 0 -> { networkDetailData.setResponseDetails( - code, - NetworkDetailCaptureUtils.createResponse( - this, - responseBody.contentLength(), - captureBodies, - { resp: Response -> resp.extractResponseBody(logger) }, - responseHeaders, - { resp: Response -> resp.headers.toMap() }, - ), + NetworkRequestData.ResponseDetails( + code, + NetworkDetailCaptureUtils.createResponse( + this, + responseBody.contentLength(), + captureBodies, + { resp: Response -> resp.extractResponseBody(logger) }, + responseHeaders, + { resp: Response -> resp.headers.toMap() }, + ), + ) ) this } - else -> wrapForCapturedStreamingBody(networkDetailData, responseBody, responseHeaders, logger) + else -> wrapUnknownContentLengthBody(networkDetailData, responseBody, responseHeaders, logger) } } /** Wraps a body that may never end, capturing the bytes as the application consumes them. */ - private fun Response.wrapForCapturedStreamingBody( + private fun Response.wrapUnknownContentLengthBody( networkDetailData: NetworkRequestData, responseBody: ResponseBody, responseHeaders: List, logger: ILogger, ): Response { - val maxBodySize = SentryReplayOptions.MAX_NETWORK_BODY_SIZE + // One byte above the limit, as the peek path does: NetworkBodyParser.fromBytes tells a + // truncated body from one that happens to match the limit exactly by that byte. + val captureCap = SentryReplayOptions.MAX_NETWORK_BODY_SIZE.toLong() + 1 val contentType = responseBody.contentType() val contentTypeString = contentType?.toString() val charset = contentType?.charset(Charsets.UTF_8)?.name() ?: "UTF-8" + // The status code and the headers are known now, so record them before anything is consumed. + // The capture below replaces them with the same values plus the body: for a stream that stays + // open for minutes, or a body the application never reads, this is all the replay ever gets. + networkDetailData.setResponseDetails( + NetworkRequestData.ResponseDetails( + code, + NetworkDetailCaptureUtils.createResponse( + this, + null, + false, + { null }, + responseHeaders, + { resp: Response -> resp.headers.toMap() }, + ), + ) + ) + return newBuilder() .body( - NetworkBodyCapturingResponseBody(responseBody, maxBodySize.toLong()) { capturedBytes -> - networkDetailData.setResponseDetails( - code, - NetworkDetailCaptureUtils.createResponse( - this, - capturedBytes.size.toLong(), - true, - { - NetworkBodyParser.fromBytes( - capturedBytes, - contentTypeString, - charset, - maxBodySize, - logger, - ) - }, - responseHeaders, - { resp: Response -> resp.headers.toMap() }, - ), - ) + NetworkBodyCapturingResponseBody(responseBody, captureCap) { capturedBytes -> + // This runs on the thread consuming the body, inside its read or close, so a failure in + // here must stay in here. The peek path guards the same parsing work the same way. + try { + networkDetailData.setResponseDetails( + NetworkRequestData.ResponseDetails( + code, + NetworkDetailCaptureUtils.createResponse( + this, + capturedBytes.size.toLong(), + true, + { + NetworkBodyParser.fromBytes( + capturedBytes, + contentTypeString, + charset, + SentryReplayOptions.MAX_NETWORK_BODY_SIZE, + logger, + ) + }, + responseHeaders, + { resp: Response -> resp.headers.toMap() }, + ), + ) + ) + } catch (e: Exception) { + logger.log( + io.sentry.SentryLevel.ERROR, + "Failed to capture the http response body for Network Details: ${e.message}", + ) + } } ) .build() diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt index 98695d9278..f79b66530f 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt @@ -185,6 +185,37 @@ class NetworkBodyCapturingResponseBodyTest { assertEquals("01234567", captured.single()?.decodeToString()) } + @Test + fun `captures every byte when the application reads less than a chunk at a time`() { + // The capture copies out of the application's sink at an offset, because a BufferedSource keeps + // what the application did not take yet in that same sink. + val chunks = (0 until 5).map { i -> "chunk-$i--".toByteArray() } + val whole = chunks.joinToString("") { it.decodeToString() } + val (wrapper, captured) = capture(ChunkedSource(chunks), 1024) + + val source = wrapper.source() + val read = StringBuilder() + while (!source.exhausted()) { + read.append(source.readUtf8(minOf(3L, source.buffer.size.coerceAtLeast(1L)))) + } + + assertEquals(whole, read.toString(), "the application receives the whole body") + assertEquals(whole, captured.single()?.decodeToString(), "and the capture holds the same bytes") + } + + @Test + fun `a failing capture callback does not reach the application`() { + val source = ChunkedSource(listOf("payload".toByteArray())) + val body = bodyOf(source) + val wrapper = + NetworkBodyCapturingResponseBody(body, 1024) { throw IllegalStateException("parse failed") } + + // The interceptor guards its own callback; the wrapper must not swallow the failure silently + // while leaving the delegate open. + assertFailsWith { wrapper.close() } + assertTrue(source.isClosed, "the delegate is closed even though the callback threw") + } + @Test fun `reports only once when the body ends and is then closed`() { val (wrapper, captured) = capture(ChunkedSource(listOf("done".toByteArray())), 1024) diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt similarity index 81% rename from sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt rename to sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt index 4bcfeaeaf8..c6ac9c1190 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorStreamingTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt @@ -4,6 +4,7 @@ import io.sentry.Hint import io.sentry.IScopes import io.sentry.SentryOptions import io.sentry.TypeCheckHint +import io.sentry.util.network.NetworkBody import io.sentry.util.network.NetworkRequestData import java.io.ByteArrayOutputStream import java.io.Closeable @@ -25,6 +26,9 @@ import okhttp3.Request import okhttp3.Response import okhttp3.mockwebserver.MockResponse import okhttp3.mockwebserver.MockWebServer +import okio.Buffer +import okio.GzipSink +import okio.buffer import org.mockito.kotlin.any import org.mockito.kotlin.argumentCaptor import org.mockito.kotlin.mock @@ -39,7 +43,7 @@ import org.mockito.kotlin.whenever * reading until the peek size is reached or the stream ends". For these responses neither happens, * so the thread that called `execute()` never received the response at all. */ -class SentryOkHttpInterceptorStreamingTest { +class SentryOkHttpInterceptorUnknownContentLengthTest { private val scopes = mock() private lateinit var options: SentryOptions @@ -72,6 +76,8 @@ class SentryOkHttpInterceptorStreamingTest { private fun capturedBody(): String? = networkDetails().response?.body?.body as String? + private fun capturedBodyValue(): Any? = networkDetails().response?.body?.body + // --------------------------------------------------------------------------------------- // the regression // --------------------------------------------------------------------------------------- @@ -85,6 +91,14 @@ class SentryOkHttpInterceptorStreamingTest { try { val response = executor.submit { call.execute() }.get(5, TimeUnit.SECONDS) assertEquals(200, response.code) + + // The stream is still open and nothing has been consumed, so the capture cannot have + // happened yet. Everything that is already known must still reach the breadcrumb. + val details = networkDetails() + assertEquals(200, details.statusCode) + assertNotNull(details.response, "the response is recorded before the body is consumed") + assertNull(details.response?.body, "but its body is not known yet") + response.close() } finally { executor.shutdownNow() @@ -116,7 +130,7 @@ class SentryOkHttpInterceptorStreamingTest { // --------------------------------------------------------------------------------------- @Test - fun `captures a streamed body once the stream ends`() { + fun `captures a body of unknown length once the stream ends`() { setUpSut() val events = listOf("data: event-0\n\n", "data: event-1\n\n") EventStreamServer(events, closeStream = true).use { stream -> @@ -129,7 +143,7 @@ class SentryOkHttpInterceptorStreamingTest { } @Test - fun `captures a streamed body when the application closes it before the stream ends`() { + fun `captures a body of unknown length when the application closes it before the stream ends`() { setUpSut() EventStreamServer(listOf("data: event-0\n\n", "data: event-1\n\n"), closeStream = false).use { stream -> @@ -142,7 +156,7 @@ class SentryOkHttpInterceptorStreamingTest { } @Test - fun `records the response of a streamed body the application never reads`() { + fun `records the response of a body of unknown length the application never reads`() { setUpSut() EventStreamServer(listOf("data: event-0\n\n"), closeStream = true).use { stream -> sut.newCall(requestTo(stream.url)).execute().close() @@ -157,7 +171,7 @@ class SentryOkHttpInterceptorStreamingTest { } @Test - fun `captures a streamed error response body`() { + fun `captures an error response body of unknown length`() { setUpSut() EventStreamServer(listOf("data: boom\n\n"), closeStream = true, statusCode = 500).use { stream -> @@ -179,11 +193,17 @@ class SentryOkHttpInterceptorStreamingTest { assertEquals(whole.length, received?.length, "the application must receive the whole stream") val captured = assertNotNull(capturedBody()) - assertTrue( - captured.length <= io.sentry.SentryReplayOptions.MAX_NETWORK_BODY_SIZE, - "the capture must stop at the cap, was ${captured.length}", + assertEquals( + io.sentry.SentryReplayOptions.MAX_NETWORK_BODY_SIZE, + captured.length, + "the capture must stop at the cap", ) assertEquals(whole.substring(0, captured.length), captured) + assertEquals( + listOf(NetworkBody.NetworkBodyWarning.TEXT_TRUNCATED), + networkDetails().response?.body?.warnings, + "a capped capture must be reported as truncated, not as a complete body", + ) } } @@ -202,6 +222,37 @@ class SentryOkHttpInterceptorStreamingTest { // responses with a known length keep the existing behaviour // --------------------------------------------------------------------------------------- + @Test + fun `captures a gzipped body, which okhttp hands over without a known length`() { + // OkHttp asks for gzip on its own and decompresses transparently, dropping Content-Length on + // the + // way, so an application interceptor sees -1 for an ordinary gzipped JSON response. That makes + // this the common path through the wrapper, not an exotic one. + setUpSut() + val json = """{"hello":"world"}""" + val gzipped = Buffer() + GzipSink(gzipped).buffer().use { it.writeUtf8(json) } + MockWebServer().use { server -> + server.enqueue( + MockResponse() + .setBody(gzipped) + .setHeader("Content-Encoding", "gzip") + .setHeader("Content-Type", "application/json") + ) + + val response = sut.newCall(requestTo(server.url("/json").toString())).execute() + assertEquals(json, response.use { it.body?.string() }) + + assertEquals( + -1L, + response.body?.contentLength(), + "okhttp reports no length for a gzipped body", + ) + assertEquals(mapOf("hello" to "world"), capturedBodyValue()) + assertEquals(200, networkDetails().statusCode) + } + } + @Test fun `still captures a response with a known length up front`() { setUpSut() @@ -244,7 +295,7 @@ class SentryOkHttpInterceptorStreamingTest { } @Test - fun `handles a zero length streamed body`() { + fun `handles a zero length body of unknown length`() { setUpSut() EventStreamServer(emptyList(), closeStream = true).use { stream -> sut.newCall(requestTo(stream.url)).execute().use { response -> @@ -257,7 +308,7 @@ class SentryOkHttpInterceptorStreamingTest { } @Test - fun `keeps the connection usable after a streamed response is closed`() { + fun `keeps the connection usable after a response of unknown length is closed`() { setUpSut() EventStreamServer(listOf("data: event-0\n\n"), closeStream = true).use { stream -> // a second call on the same client proves the connection was released properly diff --git a/sentry/api/sentry.api b/sentry/api/sentry.api index 01b068ee68..844347b7b4 100644 --- a/sentry/api/sentry.api +++ b/sentry/api/sentry.api @@ -8313,7 +8313,12 @@ public final class io/sentry/util/network/NetworkRequestData { public fun getResponseBodySize ()Ljava/lang/Long; public fun getStatusCode ()Ljava/lang/Integer; public fun setRequestDetails (Lio/sentry/util/network/ReplayNetworkRequestOrResponse;)V - public fun setResponseDetails (ILio/sentry/util/network/ReplayNetworkRequestOrResponse;)V + public fun setResponseDetails (Lio/sentry/util/network/NetworkRequestData$ResponseDetails;)V + public fun toString ()Ljava/lang/String; +} + +public final class io/sentry/util/network/NetworkRequestData$ResponseDetails { + public fun (ILio/sentry/util/network/ReplayNetworkRequestOrResponse;)V public fun toString ()Ljava/lang/String; } diff --git a/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java b/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java index bebf6189c8..b9e7e3f1c2 100644 --- a/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java +++ b/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java @@ -17,7 +17,7 @@ public final class NetworkRequestData { private @Nullable ReplayNetworkRequestOrResponse request; // The response can be filled in after this instance was handed to the scope: an integration that - // captures a streamed body only knows it once the stream has been consumed. Keeping the three + // captures a body of unknown length only knows it once the body has been consumed. Keeping the // response values behind one volatile reference means a reader on the replay thread sees either // nothing or the complete set, never a mix of the two. private volatile @Nullable ResponseDetails responseDetails; @@ -62,13 +62,9 @@ public void setRequestDetails(@NotNull final ReplayNetworkRequestOrResponse requ this.requestBodySize = requestData.getSize(); } - /** - * Populates this instance with request details obtained via {@link - * NetworkDetailCaptureUtils#createResponse} - */ - public void setResponseDetails( - final int statusCode, @NotNull final ReplayNetworkRequestOrResponse responseData) { - this.responseDetails = new ResponseDetails(statusCode, responseData.getSize(), responseData); + /** Populates this instance with the response details assembled by the caller. */ + public void setResponseDetails(@NotNull final ResponseDetails details) { + this.responseDetails = details; } @Override @@ -77,32 +73,45 @@ public String toString() { + "method='" + method + '\'' - + ", statusCode=" - + getStatusCode() + ", requestBodySize=" + requestBodySize - + ", responseBodySize=" - + getResponseBodySize() + ", request=" + request - + ", response=" - + getResponse() + + ", responseDetails=" + + responseDetails + '}'; } - /** Immutable, so publishing one reference publishes all three values. */ - private static final class ResponseDetails { + /** + * The response side of a {@link NetworkRequestData}, immutable so one reference publishes all. + */ + public static final class ResponseDetails { private final int statusCode; private final @Nullable Long bodySize; private final @NotNull ReplayNetworkRequestOrResponse response; - ResponseDetails( - final int statusCode, - final @Nullable Long bodySize, - final @NotNull ReplayNetworkRequestOrResponse response) { + /** + * @param statusCode the HTTP status code of the response. + * @param response the response details obtained via {@link + * NetworkDetailCaptureUtils#createResponse} + */ + public ResponseDetails( + final int statusCode, final @NotNull ReplayNetworkRequestOrResponse response) { this.statusCode = statusCode; - this.bodySize = bodySize; + this.bodySize = response.getSize(); this.response = response; } + + @Override + public String toString() { + return "ResponseDetails{" + + "statusCode=" + + statusCode + + ", bodySize=" + + bodySize + + ", response=" + + response + + '}'; + } } } diff --git a/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt b/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt index 3665457637..55f97a8382 100644 --- a/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt +++ b/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt @@ -6,8 +6,11 @@ import kotlin.test.assertNull class NetworkRequestDataTest { - private fun response(statusCode: Int, size: Long): ReplayNetworkRequestOrResponse = - ReplayNetworkRequestOrResponse(size, null, emptyMap()) + private fun responseDetails(statusCode: Int, size: Long): NetworkRequestData.ResponseDetails = + NetworkRequestData.ResponseDetails( + statusCode, + ReplayNetworkRequestOrResponse(size, null, emptyMap()), + ) @Test fun `response details are empty until they are set`() { @@ -22,7 +25,7 @@ class NetworkRequestDataTest { fun `response details are all readable once set`() { val data = NetworkRequestData("GET") - data.setResponseDetails(200, response(200, 42L)) + data.setResponseDetails(responseDetails(200, 42L)) assertEquals(200, data.getStatusCode()) assertEquals(42L, data.getResponseBodySize()) @@ -31,14 +34,15 @@ class NetworkRequestDataTest { @Test fun `response details can be filled in after the instance was published`() { - // A streamed body is only known once it has been consumed, which can be after the breadcrumb + // A body of unknown length is only known once it has been consumed, which can be after the + // breadcrumb // holding this instance already reached the scope, so the values must stay consistent for a // reader that arrives late. val data = NetworkRequestData("GET") assertNull(data.getStatusCode()) assertNull(data.getResponse()) - data.setResponseDetails(500, response(500, 7L)) + data.setResponseDetails(responseDetails(500, 7L)) assertEquals(500, data.getStatusCode()) assertEquals(7L, data.getResponseBodySize()) From a93c694ebb567b21fd0eba3f785f30d296f52521 Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Tue, 6 Oct 2026 15:34:16 +0200 Subject: [PATCH 08/16] test(okhttp): Cover async, cancelled and HTTP/2 bodies of unknown length --- .../NetworkBodyCapturingResponseBodyTest.kt | 56 +++++++++ ...HttpInterceptorUnknownContentLengthTest.kt | 108 ++++++++++++++++++ 2 files changed, 164 insertions(+) diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt index f79b66530f..47fcd9329a 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt @@ -1,6 +1,8 @@ package io.sentry.okhttp import java.io.IOException +import java.util.concurrent.CountDownLatch +import java.util.concurrent.TimeUnit import kotlin.test.Test import kotlin.test.assertContentEquals import kotlin.test.assertEquals @@ -64,6 +66,30 @@ class NetworkBodyCapturingResponseBodyTest { } } + /** Emits one chunk, then parks until it is closed, like a stream that has gone quiet. */ + private class ParkingSource(private val chunk: ByteArray) : Source { + val emitted = CountDownLatch(1) + private val closed = CountDownLatch(1) + private var sent = false + + override fun read(sink: Buffer, byteCount: Long): Long { + if (!sent) { + sent = true + sink.write(chunk) + emitted.countDown() + return chunk.size.toLong() + } + closed.await(10, TimeUnit.SECONDS) + return -1L + } + + override fun timeout(): Timeout = Timeout.NONE + + override fun close() { + closed.countDown() + } + } + private fun bodyOf( source: Source, type: String? = "text/plain", @@ -216,6 +242,36 @@ class NetworkBodyCapturingResponseBodyTest { assertTrue(source.isClosed, "the delegate is closed even though the callback threw") } + @Test + fun `reports once when the body is closed while another thread is reading it`() { + // Cancelling a stream from elsewhere closes the body from a thread other than the consumer, so + // the capture is read and written at the same time. + val chunk = "data: event-0\n\n".toByteArray() + val source = ParkingSource(chunk) + val (wrapper, captured) = capture(source, 1024) + val readDone = CountDownLatch(1) + val reader = Thread { + wrapper.source().readByteArray() + readDone.countDown() + } + reader.isDaemon = true + reader.start() + + assertTrue(source.emitted.await(10, TimeUnit.SECONDS), "the reader must have taken the chunk") + wrapper.close() + + assertTrue( + readDone.await(10, TimeUnit.SECONDS), + "the reader must finish once the body is closed", + ) + assertEquals( + 1, + captured.size, + "the capture is reported once, by whichever thread got there first", + ) + assertContentEquals(chunk, captured.single()) + } + @Test fun `reports only once when the body ends and is then closed`() { val (wrapper, captured) = capture(ChunkedSource(listOf("done".toByteArray())), 1024) diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt index c6ac9c1190..6bfdf12cf4 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt @@ -13,15 +13,20 @@ import java.net.InetAddress import java.net.ServerSocket import java.net.Socket import java.util.concurrent.CopyOnWriteArrayList +import java.util.concurrent.CountDownLatch import java.util.concurrent.Executors import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicReference import kotlin.test.Test import kotlin.test.assertEquals +import kotlin.test.assertNotEquals import kotlin.test.assertNotNull import kotlin.test.assertNull import kotlin.test.assertTrue import okhttp3.Call +import okhttp3.Callback import okhttp3.OkHttpClient +import okhttp3.Protocol import okhttp3.Request import okhttp3.Response import okhttp3.mockwebserver.MockResponse @@ -218,6 +223,109 @@ class SentryOkHttpInterceptorUnknownContentLengthTest { } } + // --------------------------------------------------------------------------------------- + // how the body is actually consumed + // --------------------------------------------------------------------------------------- + + @Test + fun `captures a body consumed from the callback of an asynchronous call`() { + // enqueue() hands the response to a dispatcher thread, so the capture is completed by a thread + // other than the one that started the call. That is the publication NetworkRequestData guards. + setUpSut() + val events = (0 until 3).map { "data: event-$it\n\n" } + EventStreamServer(events, closeStream = true).use { stream -> + val body = AtomicReference() + val thread = AtomicReference() + val done = CountDownLatch(1) + + sut + .newCall(requestTo(stream.url)) + .enqueue( + object : Callback { + override fun onFailure(call: Call, e: IOException) = done.countDown() + + override fun onResponse(call: Call, response: Response) { + response.use { body.set(it.body?.string()) } + thread.set(Thread.currentThread().name) + done.countDown() + } + } + ) + + assertTrue(done.await(10, TimeUnit.SECONDS), "the callback must be reached") + assertEquals(events.joinToString(""), body.get()) + assertNotEquals( + Thread.currentThread().name, + thread.get(), + "the body is consumed off the calling thread", + ) + assertEquals(events.joinToString(""), capturedBody()) + assertEquals(200, networkDetails().statusCode) + } + } + + @Test + fun `records the response of a stream that is cancelled while a reader is parked on it`() { + // How an SSE client consumes a stream: line by line on its own thread, until the call is + // cancelled. Nothing closes the body afterwards, so the capture never completes - but what was + // known when the response arrived must still reach the breadcrumb. + setUpSut() + val events = (0 until 3).map { "data: event-$it\n\n" } + EventStreamServer(events, closeStream = false).use { stream -> + val call = sut.newCall(requestTo(stream.url)) + val response = call.execute() + val readThree = CountDownLatch(3) + val failed = CountDownLatch(1) + val reader = Thread { + @Suppress("SwallowedException") // cancelling the call is how this read is meant to end + try { + val source = response.body!!.source() + while (true) { + val line = source.readUtf8Line() ?: break + if (line.startsWith("data:")) readThree.countDown() + } + } catch (e: IOException) { + failed.countDown() + } + } + reader.isDaemon = true + reader.start() + + assertTrue(readThree.await(10, TimeUnit.SECONDS), "the reader must receive the events") + call.cancel() + + assertTrue(failed.await(10, TimeUnit.SECONDS), "the parked read must end with an IOException") + assertEquals(200, networkDetails().statusCode) + assertNull(capturedBody(), "a cancelled stream never completes its capture") + } + } + + @Test + fun `captures a body of unknown length over HTTP2`() { + // HTTP/2 has no chunked encoding and ends a body with an empty DATA frame, so it reaches the + // wrapper by a different route than the HTTP/1.1 tests above. + setUpSut() + val payload = "data: over-h2\n\n" + MockWebServer().use { server -> + server.protocols = listOf(Protocol.H2_PRIOR_KNOWLEDGE) + server.enqueue( + MockResponse() + .setBody(payload) + .removeHeader("Content-Length") + .setHeader("Content-Type", "text/event-stream") + ) + sut = sut.newBuilder().protocols(listOf(Protocol.H2_PRIOR_KNOWLEDGE)).build() + + val response = sut.newCall(requestTo(server.url("/events").toString())).execute() + assertEquals(payload, response.use { it.body?.string() }) + + assertEquals(Protocol.H2_PRIOR_KNOWLEDGE, response.protocol) + assertEquals(-1L, response.body?.contentLength(), "no known length on the h2 body") + assertEquals(payload, capturedBody()) + assertEquals(200, networkDetails().statusCode) + } + } + // --------------------------------------------------------------------------------------- // responses with a known length keep the existing behaviour // --------------------------------------------------------------------------------------- From b77f96a2a7a4d9c81c347e0dbdbe44eb8b3ecb20 Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Tue, 6 Oct 2026 15:34:16 +0200 Subject: [PATCH 09/16] ref: Publish both sides of NetworkRequestData the same way --- .../util/network/NetworkRequestData.java | 40 +++++++++---------- 1 file changed, 18 insertions(+), 22 deletions(-) diff --git a/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java b/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java index b9e7e3f1c2..026f4880aa 100644 --- a/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java +++ b/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java @@ -13,13 +13,13 @@ @ApiStatus.Internal public final class NetworkRequestData { private @Nullable final String method; - private @Nullable Long requestBodySize; - private @Nullable ReplayNetworkRequestOrResponse request; - // The response can be filled in after this instance was handed to the scope: an integration that - // captures a body of unknown length only knows it once the body has been consumed. Keeping the - // response values behind one volatile reference means a reader on the replay thread sees either - // nothing or the complete set, never a mix of the two. + // Both sides are filled in by the thread running the http call and read by the replay thread, so + // both are published through a volatile write. The response side needs it most: an integration + // that captures a body of unknown length only knows the response once the body has been consumed, + // which can be after this instance was handed to the scope. Keeping its values behind one + // reference to an immutable object means a reader sees either nothing or the complete set. + private volatile @Nullable ReplayNetworkRequestOrResponse request; private volatile @Nullable ResponseDetails responseDetails; public NetworkRequestData(@Nullable final String method) { @@ -36,12 +36,13 @@ public NetworkRequestData(@Nullable final String method) { } public @Nullable Long getRequestBodySize() { - return requestBodySize; + final ReplayNetworkRequestOrResponse requestData = request; + return requestData == null ? null : requestData.getSize(); } public @Nullable Long getResponseBodySize() { final ResponseDetails details = responseDetails; - return details == null ? null : details.bodySize; + return details == null ? null : details.response.getSize(); } public @Nullable ReplayNetworkRequestOrResponse getRequest() { @@ -59,10 +60,16 @@ public NetworkRequestData(@Nullable final String method) { */ public void setRequestDetails(@NotNull final ReplayNetworkRequestOrResponse requestData) { this.request = requestData; - this.requestBodySize = requestData.getSize(); } - /** Populates this instance with the response details assembled by the caller. */ + /** + * Populates this instance with the response details assembled by the caller. + * + *

May be called from another thread than the one that created this instance, and after the + * instance was handed to the scope. A later call replaces the details of an earlier one, so an + * integration can record the status code and the headers as soon as the response arrives and add + * the body once it has been consumed. + */ public void setResponseDetails(@NotNull final ResponseDetails details) { this.responseDetails = details; } @@ -73,8 +80,6 @@ public String toString() { + "method='" + method + '\'' - + ", requestBodySize=" - + requestBodySize + ", request=" + request + ", responseDetails=" @@ -87,7 +92,6 @@ public String toString() { */ public static final class ResponseDetails { private final int statusCode; - private final @Nullable Long bodySize; private final @NotNull ReplayNetworkRequestOrResponse response; /** @@ -98,20 +102,12 @@ public static final class ResponseDetails { public ResponseDetails( final int statusCode, final @NotNull ReplayNetworkRequestOrResponse response) { this.statusCode = statusCode; - this.bodySize = response.getSize(); this.response = response; } @Override public String toString() { - return "ResponseDetails{" - + "statusCode=" - + statusCode - + ", bodySize=" - + bodySize - + ", response=" - + response - + '}'; + return "ResponseDetails{" + "statusCode=" + statusCode + ", response=" + response + '}'; } } } From 17d513455bc4d8e08e698f9cebf02bf36da2ef58 Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Tue, 6 Oct 2026 16:41:39 +0200 Subject: [PATCH 10/16] fix(okhttp): Keep the bytes received before a stream breaks --- .../NetworkBodyCapturingResponseBody.kt | 14 ++++++++-- .../NetworkBodyCapturingResponseBodyTest.kt | 7 ++++- ...HttpInterceptorUnknownContentLengthTest.kt | 28 ++++++++++++++++--- 3 files changed, 42 insertions(+), 7 deletions(-) diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt index 3a0149a06b..64df558ae3 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt @@ -1,5 +1,6 @@ package io.sentry.okhttp +import java.io.IOException import java.util.concurrent.atomic.AtomicBoolean import okhttp3.MediaType import okhttp3.ResponseBody @@ -20,7 +21,8 @@ import okio.buffer * * This body instead captures what is actually consumed. It forwards every byte to the application * untouched and hands the captured bytes to [onCaptured] as soon as the capture can no longer grow, - * which is when the stream ends, when the application closes the body, or when the cap is reached. + * which is when the stream ends or breaks, when the application closes the body, or when the cap is + * reached. * * @param delegate the body to capture from. * @param maxBytes the maximum number of bytes to retain; capture stops once it is reached. @@ -66,7 +68,15 @@ internal class NetworkBodyCapturingResponseBody( private inner class CapturingSource(source: Source) : ForwardingSource(source) { override fun read(sink: Buffer, byteCount: Long): Long { val sinkBefore = sink.size - val read = super.read(sink, byteCount) + val read = + try { + super.read(sink, byteCount) + } catch (e: IOException) { + // The stream broke, so nothing more can arrive. What did arrive is the evidence for this + // very failure, and a caller handling the error is not obliged to close the body. + reportCaptured() + throw e + } if (read > 0L) { val capFull = diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt index 47fcd9329a..9a87d0bcf3 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt @@ -308,9 +308,14 @@ class NetworkBodyCapturingResponseBodyTest { val wrapper = NetworkBodyCapturingResponseBody(bodyOf(source), 1024) { captured.add(it) } assertFailsWith { wrapper.source().readByteArray() } - assertTrue(captured.isEmpty(), "a failed stream has no final capture") + assertEquals( + "xxxx", + captured.single()?.decodeToString(), + "what arrived before the failure is reported, since nothing more can arrive", + ) wrapper.close() + assertEquals(1, captured.size, "closing afterwards does not report a second time") assertTrue(source.isClosed, "the connection must still be releasable after a failure") } diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt index 6bfdf12cf4..dc3ec57be3 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt @@ -19,6 +19,7 @@ import java.util.concurrent.TimeUnit import java.util.concurrent.atomic.AtomicReference import kotlin.test.Test import kotlin.test.assertEquals +import kotlin.test.assertFailsWith import kotlin.test.assertNotEquals import kotlin.test.assertNotNull import kotlin.test.assertNull @@ -265,10 +266,10 @@ class SentryOkHttpInterceptorUnknownContentLengthTest { } @Test - fun `records the response of a stream that is cancelled while a reader is parked on it`() { + fun `captures what arrived when a stream is cancelled while a reader is parked on it`() { // How an SSE client consumes a stream: line by line on its own thread, until the call is - // cancelled. Nothing closes the body afterwards, so the capture never completes - but what was - // known when the response arrived must still reach the breadcrumb. + // cancelled. Nothing closes the body afterwards, so the broken read is what completes the + // capture, with the events that did arrive. setUpSut() val events = (0 until 3).map { "data: event-$it\n\n" } EventStreamServer(events, closeStream = false).use { stream -> @@ -296,7 +297,21 @@ class SentryOkHttpInterceptorUnknownContentLengthTest { assertTrue(failed.await(10, TimeUnit.SECONDS), "the parked read must end with an IOException") assertEquals(200, networkDetails().statusCode) - assertNull(capturedBody(), "a cancelled stream never completes its capture") + assertEquals(events.joinToString(""), capturedBody()) + } + } + + @Test + fun `captures what arrived when the server drops the connection mid stream`() { + setUpSut() + val events = (0 until 3).map { "data: event-$it\n\n" } + EventStreamServer(events, closeStream = false, dropAfterEvents = true).use { stream -> + val response = sut.newCall(requestTo(stream.url)).execute() + + assertFailsWith { response.body?.string() } + + assertEquals(200, networkDetails().statusCode) + assertEquals(events.joinToString(""), capturedBody()) } } @@ -434,6 +449,7 @@ class SentryOkHttpInterceptorUnknownContentLengthTest { private val events: List, private val closeStream: Boolean, private val statusCode: Int = 200, + private val dropAfterEvents: Boolean = false, ) : Closeable { private val server = ServerSocket(0, 8, InetAddress.getLoopbackAddress()) private val connections = CopyOnWriteArrayList() @@ -479,6 +495,10 @@ class SentryOkHttpInterceptorUnknownContentLengthTest { output.write("\r\n".toByteArray()) output.flush() } + if (dropAfterEvents) { + // hang up without the terminating chunk, the way a mobile connection dies + return + } if (closeStream) { output.write("0\r\n\r\n".toByteArray()) output.flush() From f9e79c0d69c21e189a67a6e98e40473a7cff7302 Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Tue, 6 Oct 2026 18:49:57 +0200 Subject: [PATCH 11/16] ref: Read the response details of a network breadcrumb as one snapshot --- .../DefaultReplayBreadcrumbConverter.kt | 10 ++++--- sentry/api/sentry.api | 3 +++ .../util/network/NetworkRequestData.java | 27 ++++++++++++++++--- .../util/network/NetworkRequestDataTest.kt | 17 ++++++++++++ 4 files changed, 50 insertions(+), 7 deletions(-) diff --git a/sentry-android-replay/src/main/java/io/sentry/android/replay/DefaultReplayBreadcrumbConverter.kt b/sentry-android-replay/src/main/java/io/sentry/android/replay/DefaultReplayBreadcrumbConverter.kt index 5405b33eea..7e6771f00e 100644 --- a/sentry-android-replay/src/main/java/io/sentry/android/replay/DefaultReplayBreadcrumbConverter.kt +++ b/sentry-android-replay/src/main/java/io/sentry/android/replay/DefaultReplayBreadcrumbConverter.kt @@ -233,10 +233,14 @@ public open class DefaultReplayBreadcrumbConverter() : ReplayBreadcrumbConverter // Add Network Details data when available networkDetailData?.let { networkData -> + // One snapshot of the response: an integration that captures a streamed body may fill it in + // while this runs, and the status code must then belong to the body next to it. + val responseDetails = networkData.responseDetails + networkData.method?.let { breadcrumbData["method"] = it } - networkData.statusCode?.let { breadcrumbData["statusCode"] = it } + responseDetails?.let { breadcrumbData["statusCode"] = it.statusCode } networkData.requestBodySize?.let { breadcrumbData["requestBodySize"] = it } - networkData.responseBodySize?.let { breadcrumbData["responseBodySize"] = it } + responseDetails?.response?.size?.let { breadcrumbData["responseBodySize"] = it } networkData.request?.let { request -> val requestData = mutableMapOf() @@ -257,7 +261,7 @@ public open class DefaultReplayBreadcrumbConverter() : ReplayBreadcrumbConverter } } - networkData.response?.let { response -> + responseDetails?.response?.let { response -> val responseData = mutableMapOf() response.size?.let { responseData["size"] = it } response.body?.let { diff --git a/sentry/api/sentry.api b/sentry/api/sentry.api index 844347b7b4..8a53d2ee76 100644 --- a/sentry/api/sentry.api +++ b/sentry/api/sentry.api @@ -8311,6 +8311,7 @@ public final class io/sentry/util/network/NetworkRequestData { public fun getRequestBodySize ()Ljava/lang/Long; public fun getResponse ()Lio/sentry/util/network/ReplayNetworkRequestOrResponse; public fun getResponseBodySize ()Ljava/lang/Long; + public fun getResponseDetails ()Lio/sentry/util/network/NetworkRequestData$ResponseDetails; public fun getStatusCode ()Ljava/lang/Integer; public fun setRequestDetails (Lio/sentry/util/network/ReplayNetworkRequestOrResponse;)V public fun setResponseDetails (Lio/sentry/util/network/NetworkRequestData$ResponseDetails;)V @@ -8319,6 +8320,8 @@ public final class io/sentry/util/network/NetworkRequestData { public final class io/sentry/util/network/NetworkRequestData$ResponseDetails { public fun (ILio/sentry/util/network/ReplayNetworkRequestOrResponse;)V + public fun getResponse ()Lio/sentry/util/network/ReplayNetworkRequestOrResponse; + public fun getStatusCode ()I public fun toString ()Ljava/lang/String; } diff --git a/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java b/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java index 026f4880aa..9a55c0f0bb 100644 --- a/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java +++ b/sentry/src/main/java/io/sentry/util/network/NetworkRequestData.java @@ -32,7 +32,7 @@ public NetworkRequestData(@Nullable final String method) { public @Nullable Integer getStatusCode() { final ResponseDetails details = responseDetails; - return details == null ? null : details.statusCode; + return details == null ? null : details.getStatusCode(); } public @Nullable Long getRequestBodySize() { @@ -42,7 +42,7 @@ public NetworkRequestData(@Nullable final String method) { public @Nullable Long getResponseBodySize() { final ResponseDetails details = responseDetails; - return details == null ? null : details.response.getSize(); + return details == null ? null : details.getResponse().getSize(); } public @Nullable ReplayNetworkRequestOrResponse getRequest() { @@ -51,7 +51,7 @@ public NetworkRequestData(@Nullable final String method) { public @Nullable ReplayNetworkRequestOrResponse getResponse() { final ResponseDetails details = responseDetails; - return details == null ? null : details.response; + return details == null ? null : details.getResponse(); } /** @@ -62,6 +62,17 @@ public void setRequestDetails(@NotNull final ReplayNetworkRequestOrResponse requ this.request = requestData; } + /** + * The response details as one snapshot, or {@code null} while the response is not known yet. + * + *

Prefer this over {@link #getStatusCode()}, {@link #getResponseBodySize()} and {@link + * #getResponse()} when more than one of them is needed: the details may be replaced between two + * of those calls, which would mix one response with the next. + */ + public @Nullable ResponseDetails getResponseDetails() { + return responseDetails; + } + /** * Populates this instance with the response details assembled by the caller. * @@ -105,9 +116,17 @@ public ResponseDetails( this.response = response; } + public int getStatusCode() { + return statusCode; + } + + public @NotNull ReplayNetworkRequestOrResponse getResponse() { + return response; + } + @Override public String toString() { - return "ResponseDetails{" + "statusCode=" + statusCode + ", response=" + response + '}'; + return "ResponseDetails{statusCode=" + statusCode + ", response=" + response + '}'; } } } diff --git a/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt b/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt index 55f97a8382..3f42e1cbae 100644 --- a/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt +++ b/sentry/src/test/java/io/sentry/util/network/NetworkRequestDataTest.kt @@ -2,6 +2,7 @@ package io.sentry.util.network import kotlin.test.Test import kotlin.test.assertEquals +import kotlin.test.assertNotNull import kotlin.test.assertNull class NetworkRequestDataTest { @@ -49,6 +50,22 @@ class NetworkRequestDataTest { assertEquals(7L, data.getResponse()?.getSize()) } + @Test + fun `the response details snapshot keeps the status code and the response together`() { + val data = NetworkRequestData("GET") + assertNull(data.responseDetails) + + data.setResponseDetails(responseDetails(200, 42L)) + val first = assertNotNull(data.responseDetails) + + data.setResponseDetails(responseDetails(500, 7L)) + + assertEquals(200, first.statusCode, "a snapshot is unaffected by a later call") + assertEquals(42L, first.response.size) + assertEquals(500, data.responseDetails?.statusCode) + assertEquals(7L, data.responseDetails?.response?.size) + } + @Test fun `request details are unaffected`() { val data = NetworkRequestData("GET") From 5b2bf58fd562b1152ed3651da7bef70f4cd5cc75 Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Wed, 7 Oct 2026 06:51:00 +0200 Subject: [PATCH 12/16] ref(okhttp): Capture every response body through one path --- .../NetworkBodyCapturingResponseBody.kt | 24 +++ .../sentry/okhttp/SentryOkHttpInterceptor.kt | 164 ++++++------------ .../NetworkBodyCapturingResponseBodyTest.kt | 23 +++ 3 files changed, 103 insertions(+), 108 deletions(-) diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt index 64df558ae3..77def564b9 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt @@ -24,6 +24,10 @@ import okio.buffer * which is when the stream ends or breaks, when the application closes the body, or when the cap is * reached. * + * A body the application never reads is still captured if its length is known: such a body ends on + * its own, so taking it while closing is bounded. That is how the body of a response whose status + * code was the only thing of interest reaches the capture. + * * @param delegate the body to capture from. * @param maxBytes the maximum number of bytes to retain; capture stops once it is reached. * @param onCaptured invoked at most once, synchronously, on the thread that finishes the body. It @@ -37,6 +41,7 @@ internal class NetworkBodyCapturingResponseBody( private val captured = Buffer() private val reported = AtomicBoolean(false) + @Volatile private var readStarted = false // A ResponseBody must hand out a BufferedSource, so the capturing source is buffered. That cannot // re-introduce the blocking this class exists to avoid: BufferedSource.read(sink, byteCount) @@ -56,6 +61,7 @@ internal class NetworkBodyCapturingResponseBody( override fun close() { try { + captureUnreadBody() // Report before closing the delegate: closing it notifies listeners, which may serialize the // breadcrumb this capture belongs to. reportCaptured() @@ -64,9 +70,27 @@ internal class NetworkBodyCapturingResponseBody( } } + /** + * Reads a body nobody touched, so it is captured as well. Only for a body of known length, which + * ends on its own, and only while no one is reading: a concurrent read of the same body is what + * OkHttp forbids. + */ + @Suppress("SwallowedException") // the body is being closed; a failure here leaves it as it was + private fun captureUnreadBody() { + if (readStarted || reported.get() || delegate.contentLength() < 0L) { + return + } + try { + capturingSource.request(maxBytes) + } catch (e: IOException) { + // nothing more to capture + } + } + /** Copies the bytes passing through into [captured]. */ private inner class CapturingSource(source: Source) : ForwardingSource(source) { override fun read(sink: Buffer, byteCount: Long): Long { + readStarted = true val sinkBefore = sink.size val read = try { diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt index 7a512241cf..0447c0efca 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt @@ -34,7 +34,6 @@ import okhttp3.Interceptor import okhttp3.Request import okhttp3.RequestBody.Companion.toRequestBody import okhttp3.Response -import okhttp3.ResponseBody import org.jetbrains.annotations.VisibleForTesting /** @@ -316,101 +315,58 @@ public open class SentryOkHttpInterceptor( } /** - * Returns this response with a body that can be captured for replay, without blocking on streams - * that never end. + * Returns this response with a body whose content is captured for replay while the application + * consumes it. * - * A body with a known length ends on its own, so it is captured up front, as before. A body of - * unknown length may never end (server-sent events, long-poll, a chunked endpoint that stays - * open): peeking it would block the calling thread until the connection dies, so it is wrapped - * and captured while it is being consumed instead. + * Capturing a body by peeking at it up front does not work for every response. OkHttp implements + * [Response.peekBody] as `request(byteCount)`, which keeps reading until that many bytes are + * buffered or the stream ends. A response of unknown length may never end — server-sent events, a + * long-poll, a chunked endpoint that stays open, and every response OkHttp decompressed itself — + * so the peek would never return and the caller would never get the response at all. + * + * Every body therefore goes through [NetworkBodyCapturingResponseBody], which copies what passes + * through it. A body with a known length that the application never reads is taken when it is + * closed, where reading it is bounded. */ private fun Response.withNetworkBodyCapture(networkDetailData: NetworkRequestData?): Response { val responseBody = body - val captureBodies = scopes.options.sessionReplay.isNetworkCaptureBodies - val responseHeaders = scopes.options.sessionReplay.networkResponseHeaders - val logger = scopes.options.logger + if (networkDetailData == null || responseBody == null) { + return this + } - return when { - networkDetailData == null || responseBody == null -> this - // a body with a known length ends on its own, so capturing it up front is bounded - !captureBodies || responseBody.contentLength() >= 0 -> { - networkDetailData.setResponseDetails( - NetworkRequestData.ResponseDetails( - code, - NetworkDetailCaptureUtils.createResponse( - this, - responseBody.contentLength(), - captureBodies, - { resp: Response -> resp.extractResponseBody(logger) }, - responseHeaders, - { resp: Response -> resp.headers.toMap() }, - ), - ) - ) - this - } - else -> wrapUnknownContentLengthBody(networkDetailData, responseBody, responseHeaders, logger) + val knownBodySize = responseBody.contentLength().takeIf { it >= 0L } + + // The status code and the headers are known now, so record them before anything is consumed. A + // stream that stays open for minutes never gets further than this. + recordResponse(networkDetailData, knownBodySize) { null } + + if (!scopes.options.sessionReplay.isNetworkCaptureBodies) { + return this } - } - /** Wraps a body that may never end, capturing the bytes as the application consumes them. */ - private fun Response.wrapUnknownContentLengthBody( - networkDetailData: NetworkRequestData, - responseBody: ResponseBody, - responseHeaders: List, - logger: ILogger, - ): Response { - // One byte above the limit, as the peek path does: NetworkBodyParser.fromBytes tells a - // truncated body from one that happens to match the limit exactly by that byte. + val logger = scopes.options.logger + // One byte above the limit, so NetworkBodyParser.fromBytes can tell a truncated body from one + // that happens to match the limit exactly. val captureCap = SentryReplayOptions.MAX_NETWORK_BODY_SIZE.toLong() + 1 val contentType = responseBody.contentType() val contentTypeString = contentType?.toString() val charset = contentType?.charset(Charsets.UTF_8)?.name() ?: "UTF-8" - // The status code and the headers are known now, so record them before anything is consumed. - // The capture below replaces them with the same values plus the body: for a stream that stays - // open for minutes, or a body the application never reads, this is all the replay ever gets. - networkDetailData.setResponseDetails( - NetworkRequestData.ResponseDetails( - code, - NetworkDetailCaptureUtils.createResponse( - this, - null, - false, - { null }, - responseHeaders, - { resp: Response -> resp.headers.toMap() }, - ), - ) - ) - return newBuilder() .body( NetworkBodyCapturingResponseBody(responseBody, captureCap) { capturedBytes -> // This runs on the thread consuming the body, inside its read or close, so a failure in - // here must stay in here. The peek path guards the same parsing work the same way. + // here must stay in here. try { - networkDetailData.setResponseDetails( - NetworkRequestData.ResponseDetails( - code, - NetworkDetailCaptureUtils.createResponse( - this, - capturedBytes.size.toLong(), - true, - { - NetworkBodyParser.fromBytes( - capturedBytes, - contentTypeString, - charset, - SentryReplayOptions.MAX_NETWORK_BODY_SIZE, - logger, - ) - }, - responseHeaders, - { resp: Response -> resp.headers.toMap() }, - ), + recordResponse(networkDetailData, knownBodySize ?: capturedBytes.size.toLong()) { + NetworkBodyParser.fromBytes( + capturedBytes, + contentTypeString, + charset, + SentryReplayOptions.MAX_NETWORK_BODY_SIZE, + logger, ) - ) + } } catch (e: Exception) { logger.log( io.sentry.SentryLevel.ERROR, @@ -422,36 +378,28 @@ public open class SentryOkHttpInterceptor( .build() } - /** Extracts the body content from an OkHttp Response safely */ - private fun Response.extractResponseBody(logger: ILogger): NetworkBody? { - return body?.let { responseBody -> - try { - val contentType = responseBody.contentType() - val contentTypeString = contentType?.toString() - val maxBodySize = SentryReplayOptions.MAX_NETWORK_BODY_SIZE - - // Peek at the body (doesn't consume it) - // We +1 here in order to properly truncate within NetworkBodyParser.fromBytes - // and be able to distinguish from an oversized request and a request matching maxBodySize - val peekBody = peekBody(maxBodySize.toLong() + 1) - val bodyBytes = peekBody.bytes() - - val charset = contentType?.charset(Charsets.UTF_8)?.name() ?: "UTF-8" - return NetworkBodyParser.fromBytes( - bodyBytes, - contentTypeString, - charset, - maxBodySize, - logger, - ) - } catch (e: Exception) { - logger.log( - io.sentry.SentryLevel.ERROR, - "Failed to read http response body for Network Details: ${e.message}", - ) - null - } - } + /** + * Records what is known about this response, replacing anything recorded for it before. Called + * once before the body is consumed and once with the captured body. + */ + private fun Response.recordResponse( + networkDetailData: NetworkRequestData, + bodySize: Long?, + body: () -> NetworkBody?, + ) { + networkDetailData.setResponseDetails( + NetworkRequestData.ResponseDetails( + code, + NetworkDetailCaptureUtils.createResponse( + this, + bodySize, + true, + { body() }, + scopes.options.sessionReplay.networkResponseHeaders, + { resp: Response -> resp.headers.toMap() }, + ), + ) + ) } private fun finishSpan( diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt index 9a87d0bcf3..8b1c2822a8 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt @@ -185,6 +185,29 @@ class NetworkBodyCapturingResponseBodyTest { assertEquals("first ", captured.single()?.decodeToString()) } + @Test + fun `captures a body of known length that the application never read`() { + val payload = "failure".toByteArray() + val source = ChunkedSource(listOf(payload)) + val (wrapper, captured) = capture(source, 1024, length = payload.size.toLong()) + + wrapper.close() + + assertContentEquals(payload, captured.single(), "a bounded body is taken while closing") + assertTrue(source.isClosed) + } + + @Test + fun `does not read a body of unknown length that the application never read`() { + val source = ChunkedSource(listOf("data: event-0\n\n".toByteArray())) + val (wrapper, captured) = capture(source, 1024) + + wrapper.close() + + assertEquals(0, source.reads, "a body that may never end must not be read while closing") + assertTrue(captured.single()!!.isEmpty()) + } + @Test fun `reports an empty capture for an empty body`() { val (wrapper, captured) = capture(ChunkedSource(emptyList()), 1024) From 8879de292c6d1c2f9279a3c3d4756f962fe2ed72 Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Wed, 7 Oct 2026 07:55:57 +0200 Subject: [PATCH 13/16] changelog --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index bba665ffd1..0f0620c9fa 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,7 +4,7 @@ ### Fixes -- Fix `SentryOkHttpInterceptor` hanging forever on responses whose body has no known length, such as Server-Sent Events ([#1](https://github.com/getsentry/sentry-java/pull/1)) +- Fix `SentryOkHttpInterceptor` hanging forever on responses whose body has no known length, such as Server-Sent Events or a gzipped response ([#6231](https://github.com/getsentry/sentry-java/pull/6231)) ### Features From 09687ef98c23a5ba697cbfa415148447579ebd65 Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Wed, 7 Oct 2026 08:39:21 +0200 Subject: [PATCH 14/16] test: Silence R8 for the okhttp types the new response body extends --- .../test-app-size/proguard-rules.pro | 2 ++ 1 file changed, 2 insertions(+) diff --git a/sentry-android-integration-tests/test-app-size/proguard-rules.pro b/sentry-android-integration-tests/test-app-size/proguard-rules.pro index 4db0031ab0..e26a4ca513 100644 --- a/sentry-android-integration-tests/test-app-size/proguard-rules.pro +++ b/sentry-android-integration-tests/test-app-size/proguard-rules.pro @@ -47,5 +47,7 @@ -dontwarn kotlin.math.MathKt -dontwarn okhttp3.EventListener -dontwarn okhttp3.Interceptor +-dontwarn okhttp3.ResponseBody +-dontwarn okio.ForwardingSource # Assume all classes are used to not strip them out, e.g. integrations like Compose or Sqlite -keep class io.sentry.** From ead659c516c21ec75e4e7f7ff9913979ae2c6aec Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Thu, 8 Oct 2026 13:27:57 +0200 Subject: [PATCH 15/16] ref(okhttp): Rename NetworkBodyCapturingResponseBody to CapturedResponseBody --- ...uringResponseBody.kt => CapturedResponseBody.kt} | 2 +- .../io/sentry/okhttp/SentryOkHttpInterceptor.kt | 8 ++++---- ...ponseBodyTest.kt => CapturedResponseBodyTest.kt} | 13 ++++++------- 3 files changed, 11 insertions(+), 12 deletions(-) rename sentry-okhttp/src/main/java/io/sentry/okhttp/{NetworkBodyCapturingResponseBody.kt => CapturedResponseBody.kt} (99%) rename sentry-okhttp/src/test/java/io/sentry/okhttp/{NetworkBodyCapturingResponseBodyTest.kt => CapturedResponseBodyTest.kt} (96%) diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/CapturedResponseBody.kt similarity index 99% rename from sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt rename to sentry-okhttp/src/main/java/io/sentry/okhttp/CapturedResponseBody.kt index 77def564b9..3bb7ea5612 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/NetworkBodyCapturingResponseBody.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/CapturedResponseBody.kt @@ -33,7 +33,7 @@ import okio.buffer * @param onCaptured invoked at most once, synchronously, on the thread that finishes the body. It * is never invoked for a body that is abandoned without being read to the end or closed. */ -internal class NetworkBodyCapturingResponseBody( +internal class CapturedResponseBody( private val delegate: ResponseBody, private val maxBytes: Long, private val onCaptured: (ByteArray) -> Unit, diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt index 0447c0efca..65aed1a3de 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt @@ -324,9 +324,9 @@ public open class SentryOkHttpInterceptor( * long-poll, a chunked endpoint that stays open, and every response OkHttp decompressed itself — * so the peek would never return and the caller would never get the response at all. * - * Every body therefore goes through [NetworkBodyCapturingResponseBody], which copies what passes - * through it. A body with a known length that the application never reads is taken when it is - * closed, where reading it is bounded. + * Every body therefore goes through [CapturedResponseBody], which copies what passes through it. + * A body with a known length that the application never reads is taken when it is closed, where + * reading it is bounded. */ private fun Response.withNetworkBodyCapture(networkDetailData: NetworkRequestData?): Response { val responseBody = body @@ -354,7 +354,7 @@ public open class SentryOkHttpInterceptor( return newBuilder() .body( - NetworkBodyCapturingResponseBody(responseBody, captureCap) { capturedBytes -> + CapturedResponseBody(responseBody, captureCap) { capturedBytes -> // This runs on the thread consuming the body, inside its read or close, so a failure in // here must stay in here. try { diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/CapturedResponseBodyTest.kt similarity index 96% rename from sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt rename to sentry-okhttp/src/test/java/io/sentry/okhttp/CapturedResponseBodyTest.kt index 8b1c2822a8..cb355c59ed 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/NetworkBodyCapturingResponseBodyTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/CapturedResponseBodyTest.kt @@ -18,7 +18,7 @@ import okio.Source import okio.Timeout import okio.buffer -class NetworkBodyCapturingResponseBodyTest { +class CapturedResponseBodyTest { /** Emits one chunk per read, then EOF. Records how many reads happened. */ private class ChunkedSource(private val chunks: List) : Source { @@ -108,10 +108,10 @@ class NetworkBodyCapturingResponseBodyTest { maxBytes: Long, type: String? = "text/plain", length: Long = -1L, - ): Pair> { + ): Pair> { val captured = mutableListOf() val wrapper = - NetworkBodyCapturingResponseBody(bodyOf(source, type, length), maxBytes) { + CapturedResponseBody(bodyOf(source, type, length), maxBytes) { captured.add(it) } return wrapper to captured @@ -256,8 +256,7 @@ class NetworkBodyCapturingResponseBodyTest { fun `a failing capture callback does not reach the application`() { val source = ChunkedSource(listOf("payload".toByteArray())) val body = bodyOf(source) - val wrapper = - NetworkBodyCapturingResponseBody(body, 1024) { throw IllegalStateException("parse failed") } + val wrapper = CapturedResponseBody(body, 1024) { throw IllegalStateException("parse failed") } // The interceptor guards its own callback; the wrapper must not swallow the failure silently // while leaving the delegate open. @@ -318,7 +317,7 @@ class NetworkBodyCapturingResponseBodyTest { override fun source(): BufferedSource = ChunkedSource(emptyList()).buffer() } - val wrapper = NetworkBodyCapturingResponseBody(delegate, 16) {} + val wrapper = CapturedResponseBody(delegate, 16) {} assertEquals(type, wrapper.contentType()) assertEquals(-1L, wrapper.contentLength()) @@ -328,7 +327,7 @@ class NetworkBodyCapturingResponseBodyTest { fun `propagates read failures to the application`() { val source = FailingSource(bytesBeforeFailure = 4) val captured = mutableListOf() - val wrapper = NetworkBodyCapturingResponseBody(bodyOf(source), 1024) { captured.add(it) } + val wrapper = CapturedResponseBody(bodyOf(source), 1024) { captured.add(it) } assertFailsWith { wrapper.source().readByteArray() } assertEquals( From 8193c3632971d76bb5df2451f7ee6a14500803cb Mon Sep 17 00:00:00 2001 From: Markus Hintersteiner Date: Thu, 8 Oct 2026 13:28:40 +0200 Subject: [PATCH 16/16] ref(okhttp): Capture only the response bytes the application reads --- .../io/sentry/okhttp/CapturedResponseBody.kt | 75 ++++++------------- .../sentry/okhttp/SentryOkHttpInterceptor.kt | 16 ++-- .../sentry/okhttp/CapturedResponseBodyTest.kt | 20 +---- ...HttpInterceptorUnknownContentLengthTest.kt | 18 ++--- 4 files changed, 42 insertions(+), 87 deletions(-) diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/CapturedResponseBody.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/CapturedResponseBody.kt index 3bb7ea5612..70c282bd24 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/CapturedResponseBody.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/CapturedResponseBody.kt @@ -11,27 +11,21 @@ import okio.Source import okio.buffer /** - * A [ResponseBody] that copies the bytes the application consumes into a capped buffer. + * A [ResponseBody] that copies the bytes the application reads into a capped buffer. * - * Capturing a body by peeking at it up front does not work for every response. OkHttp implements - * [okhttp3.Response.peekBody] as `request(byteCount)`, which keeps reading until the requested - * number of bytes is buffered or the stream ends. A response of unknown length may never end — a - * server-sent-events stream, a long-poll, a chunked endpoint that stays open — so the peek never - * returns and the thread that called `execute()` never gets the response at all. + * Reading the body up front instead, with [okhttp3.Response.peekBody], deadlocks on a response of + * unknown length: a peek is `request(byteCount)`, which waits for that many bytes or for the end of + * the stream, and a server-sent-events stream, a long-poll or an open chunked endpoint delivers + * neither. The caller would never receive the response at all. * - * This body instead captures what is actually consumed. It forwards every byte to the application - * untouched and hands the captured bytes to [onCaptured] as soon as the capture can no longer grow, - * which is when the stream ends or breaks, when the application closes the body, or when the cap is - * reached. - * - * A body the application never reads is still captured if its length is known: such a body ends on - * its own, so taking it while closing is bounded. That is how the body of a response whose status - * code was the only thing of interest reaches the capture. + * Nothing is therefore read on the capture's own account. Only what the application reads is + * captured, which also keeps the cost of the instrumentation to a copy of those bytes. * * @param delegate the body to capture from. * @param maxBytes the maximum number of bytes to retain; capture stops once it is reached. - * @param onCaptured invoked at most once, synchronously, on the thread that finishes the body. It - * is never invoked for a body that is abandoned without being read to the end or closed. + * @param onCaptured invoked at most once, synchronously, on the thread that finishes the body — so + * it must not block. A body that is abandoned without being read to the end or closed never + * reaches it. */ internal class CapturedResponseBody( private val delegate: ResponseBody, @@ -41,14 +35,11 @@ internal class CapturedResponseBody( private val captured = Buffer() private val reported = AtomicBoolean(false) - @Volatile private var readStarted = false - // A ResponseBody must hand out a BufferedSource, so the capturing source is buffered. That cannot - // re-introduce the blocking this class exists to avoid: BufferedSource.read(sink, byteCount) - // issues at most one segment-sized read on the source below it and returns with whatever arrived, - // so a short event is still forwarded on its own. The reads that loop until a byte count is - // reached — request, require, readByteArray() — are the application's own choice, and the capture - // never calls them. + // Buffering the capturing source cannot re-introduce the deadlock above: a BufferedSource read + // takes at most one segment from the source below it and returns with whatever arrived, so a + // short event is still forwarded on its own. Only request/require/readByteArray() wait for a byte + // count, and those are the application's own calls. private val capturingSource: BufferedSource by lazy { CapturingSource(delegate.source()).buffer() } @@ -61,43 +52,23 @@ internal class CapturedResponseBody( override fun close() { try { - captureUnreadBody() - // Report before closing the delegate: closing it notifies listeners, which may serialize the - // breadcrumb this capture belongs to. + // Closing the delegate notifies listeners, which may serialize the breadcrumb this capture + // belongs to, so report first. reportCaptured() } finally { delegate.close() } } - /** - * Reads a body nobody touched, so it is captured as well. Only for a body of known length, which - * ends on its own, and only while no one is reading: a concurrent read of the same body is what - * OkHttp forbids. - */ - @Suppress("SwallowedException") // the body is being closed; a failure here leaves it as it was - private fun captureUnreadBody() { - if (readStarted || reported.get() || delegate.contentLength() < 0L) { - return - } - try { - capturingSource.request(maxBytes) - } catch (e: IOException) { - // nothing more to capture - } - } - - /** Copies the bytes passing through into [captured]. */ private inner class CapturingSource(source: Source) : ForwardingSource(source) { override fun read(sink: Buffer, byteCount: Long): Long { - readStarted = true val sinkBefore = sink.size val read = try { super.read(sink, byteCount) } catch (e: IOException) { - // The stream broke, so nothing more can arrive. What did arrive is the evidence for this - // very failure, and a caller handling the error is not obliged to close the body. + // What arrived is the evidence for this very failure, and a caller handling it is not + // obliged to close the body. reportCaptured() throw e } @@ -105,6 +76,7 @@ internal class CapturedResponseBody( if (read > 0L) { val capFull = synchronized(captured) { + // the application may not have taken everything in the sink yet, hence the offset val toTake = minOf(maxBytes - captured.size, sink.size - sinkBefore) if (toTake > 0L) { sink.copyTo(captured, sinkBefore, toTake) @@ -112,11 +84,9 @@ internal class CapturedResponseBody( captured.size >= maxBytes } if (capFull) { - // the capture can never grow again, so this is the moment it becomes final reportCaptured() } } else if (read == -1L) { - // end of stream: nothing more can ever arrive reportCaptured() } return read @@ -132,10 +102,9 @@ internal class CapturedResponseBody( } /** - * The consuming thread writes [captured] from [CapturingSource.read] while [close] may be called - * by another thread — cancelling a stream from elsewhere is ordinary use — and [Buffer] is not - * thread-safe, so both accesses are guarded. [reported] keeps the callback to a single - * invocation. + * [captured] is written by the thread reading the body and read here, which [close] may reach + * from another thread — cancelling a stream from elsewhere is ordinary use — and [Buffer] is not + * thread-safe. */ private fun reportCaptured() { if (reported.compareAndSet(false, true)) { diff --git a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt index 65aed1a3de..820dc8b044 100644 --- a/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt +++ b/sentry-okhttp/src/main/java/io/sentry/okhttp/SentryOkHttpInterceptor.kt @@ -316,17 +316,15 @@ public open class SentryOkHttpInterceptor( /** * Returns this response with a body whose content is captured for replay while the application - * consumes it. + * reads it. * - * Capturing a body by peeking at it up front does not work for every response. OkHttp implements - * [Response.peekBody] as `request(byteCount)`, which keeps reading until that many bytes are - * buffered or the stream ends. A response of unknown length may never end — server-sent events, a - * long-poll, a chunked endpoint that stays open, and every response OkHttp decompressed itself — - * so the peek would never return and the caller would never get the response at all. + * Reading the body up front, as [Response.peekBody] does, deadlocks on a response of unknown + * length: server-sent events, a long-poll, a chunked endpoint that stays open, and every response + * OkHttp decompressed itself. The caller would never receive the response at all. * - * Every body therefore goes through [CapturedResponseBody], which copies what passes through it. - * A body with a known length that the application never reads is taken when it is closed, where - * reading it is bounded. + * Every body therefore goes through [CapturedResponseBody], which captures only what the + * application reads. A body nobody reads leaves the status code and the headers recorded below, + * and no body. */ private fun Response.withNetworkBodyCapture(networkDetailData: NetworkRequestData?): Response { val responseBody = body diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/CapturedResponseBodyTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/CapturedResponseBodyTest.kt index cb355c59ed..9b740f1789 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/CapturedResponseBodyTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/CapturedResponseBodyTest.kt @@ -186,25 +186,13 @@ class CapturedResponseBodyTest { } @Test - fun `captures a body of known length that the application never read`() { - val payload = "failure".toByteArray() - val source = ChunkedSource(listOf(payload)) - val (wrapper, captured) = capture(source, 1024, length = payload.size.toLong()) + fun `does not read a body the application never read`() { + val source = ChunkedSource(listOf("failure".toByteArray())) + val (wrapper, captured) = capture(source, 1024, length = 7L) wrapper.close() - assertContentEquals(payload, captured.single(), "a bounded body is taken while closing") - assertTrue(source.isClosed) - } - - @Test - fun `does not read a body of unknown length that the application never read`() { - val source = ChunkedSource(listOf("data: event-0\n\n".toByteArray())) - val (wrapper, captured) = capture(source, 1024) - - wrapper.close() - - assertEquals(0, source.reads, "a body that may never end must not be read while closing") + assertEquals(0, source.reads, "the capture must not read on its own account") assertTrue(captured.single()!!.isEmpty()) } diff --git a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt index dc3ec57be3..fb21bebb62 100644 --- a/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt +++ b/sentry-okhttp/src/test/java/io/sentry/okhttp/SentryOkHttpInterceptorUnknownContentLengthTest.kt @@ -342,15 +342,13 @@ class SentryOkHttpInterceptorUnknownContentLengthTest { } // --------------------------------------------------------------------------------------- - // responses with a known length keep the existing behaviour + // responses with a known length // --------------------------------------------------------------------------------------- @Test fun `captures a gzipped body, which okhttp hands over without a known length`() { - // OkHttp asks for gzip on its own and decompresses transparently, dropping Content-Length on - // the - // way, so an application interceptor sees -1 for an ordinary gzipped JSON response. That makes - // this the common path through the wrapper, not an exotic one. + // OkHttp asks for gzip itself and decompresses transparently, dropping Content-Length on the + // way, so an application interceptor sees -1 for an ordinary gzipped JSON response. setUpSut() val json = """{"hello":"world"}""" val gzipped = Buffer() @@ -377,19 +375,21 @@ class SentryOkHttpInterceptorUnknownContentLengthTest { } @Test - fun `still captures a response with a known length up front`() { + fun `captures a body with a known length as the application reads it`() { setUpSut() MockWebServer().use { server -> server.enqueue(MockResponse().setBody("response body").setResponseCode(200)) - sut.newCall(requestTo(server.url("/hello").toString())).execute().close() + sut.newCall(requestTo(server.url("/hello").toString())).execute().use { it.body?.string() } assertEquals("response body", capturedBody()) } } @Test - fun `still captures an error response with a known length the application never reads`() { + fun `records an error response the application never reads, without its body`() { + // Nothing is read on the capture's own account, not even a body that would end on its own, so + // the interceptor costs a consumer no more than the bytes it asked for itself. setUpSut() MockWebServer().use { server -> server.enqueue(MockResponse().setBody("failure").setResponseCode(500)) @@ -397,7 +397,7 @@ class SentryOkHttpInterceptorUnknownContentLengthTest { sut.newCall(requestTo(server.url("/hello").toString())).execute().close() assertEquals(500, networkDetails().statusCode) - assertEquals("failure", capturedBody()) + assertNull(capturedBody()) } }