From 789512cdff6e8c8ae76e4fe8c6dd1979f75061f5 Mon Sep 17 00:00:00 2001 From: Sylvester Kaczmarek <16242628+sylvesterkaczmarek@users.noreply.github.com> Date: Tue, 18 Aug 2026 16:34:44 +0100 Subject: [PATCH] fix: prevent duplicate phantom-reachable cleanup --- .../com/openai/core/PhantomReachable.kt | 19 ++++++++++++++---- .../openai/core/PhantomReachableSleeper.kt | 6 ++---- ...ntomReachableClosingAsyncStreamResponse.kt | 5 ++--- .../http/PhantomReachableClosingHttpClient.kt | 6 ++---- ...eachableClosingHttpRequestAuthenticator.kt | 6 ++---- .../PhantomReachableClosingStreamResponse.kt | 6 ++---- .../com/openai/core/PhantomReachableTest.kt | 15 ++++++++++++++ .../PhantomReachableClosingHttpClientTest.kt | 20 +++++++++++++++++++ 8 files changed, 60 insertions(+), 23 deletions(-) create mode 100644 openai-java-core/src/test/kotlin/com/openai/core/http/PhantomReachableClosingHttpClientTest.kt diff --git a/openai-java-core/src/main/kotlin/com/openai/core/PhantomReachable.kt b/openai-java-core/src/main/kotlin/com/openai/core/PhantomReachable.kt index a1662c173..d1a459382 100644 --- a/openai-java-core/src/main/kotlin/com/openai/core/PhantomReachable.kt +++ b/openai-java-core/src/main/kotlin/com/openai/core/PhantomReachable.kt @@ -4,28 +4,39 @@ package com.openai.core import com.openai.errors.OpenAIException import java.lang.reflect.InvocationTargetException +import java.util.concurrent.atomic.AtomicBoolean /** * Closes [closeable] when [observed] becomes only phantom reachable. * * This is a wrapper around a Java 9+ [java.lang.ref.Cleaner], or a no-op in older Java versions. + * The returned handle performs the same cleanup explicitly, at most once across both paths. */ @JvmSynthetic -internal fun closeWhenPhantomReachable(observed: Any, closeable: AutoCloseable) { +internal fun closeWhenPhantomReachable(observed: Any, closeable: AutoCloseable): AutoCloseable { check(observed !== closeable) { "`observed` cannot be the same object as `closeable` because it would never become phantom reachable" } - closeWhenPhantomReachable(observed, closeable::close) + return closeWhenPhantomReachable(observed, closeable::close) } /** * Calls [close] when [observed] becomes only phantom reachable. * * This is a wrapper around a Java 9+ [java.lang.ref.Cleaner], or a no-op in older Java versions. + * Calling the returned handle performs the same cleanup explicitly, at most once across both paths. */ @JvmSynthetic -internal fun closeWhenPhantomReachable(observed: Any, close: () -> Unit) { - closeWhenPhantomReachable?.let { it(observed, close) } +internal fun closeWhenPhantomReachable(observed: Any, close: () -> Unit): AutoCloseable { + val closed = AtomicBoolean(false) + val closeOnce = { + if (closed.compareAndSet(false, true)) { + close() + } + } + + closeWhenPhantomReachable?.let { it(observed, closeOnce) } + return AutoCloseable { closeOnce() } } private val closeWhenPhantomReachable: ((Any, () -> Unit) -> Unit)? by lazy { diff --git a/openai-java-core/src/main/kotlin/com/openai/core/PhantomReachableSleeper.kt b/openai-java-core/src/main/kotlin/com/openai/core/PhantomReachableSleeper.kt index 65ec40080..4dc882d04 100644 --- a/openai-java-core/src/main/kotlin/com/openai/core/PhantomReachableSleeper.kt +++ b/openai-java-core/src/main/kotlin/com/openai/core/PhantomReachableSleeper.kt @@ -10,14 +10,12 @@ import java.util.concurrent.CompletableFuture */ internal class PhantomReachableSleeper(private val sleeper: Sleeper) : Sleeper { - init { - closeWhenPhantomReachable(this, sleeper) - } + private val closeHandle = closeWhenPhantomReachable(this, sleeper) override fun sleep(duration: Duration) = sleeper.sleep(duration) override fun sleepAsync(duration: Duration): CompletableFuture = sleeper.sleepAsync(duration) - override fun close() = sleeper.close() + override fun close() = closeHandle.close() } diff --git a/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingAsyncStreamResponse.kt b/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingAsyncStreamResponse.kt index 8c32247d8..3aade8e43 100644 --- a/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingAsyncStreamResponse.kt +++ b/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingAsyncStreamResponse.kt @@ -21,9 +21,8 @@ internal class PhantomReachableClosingAsyncStreamResponse( */ private val reachabilityTracker = Object() - init { + private val closeHandle = closeWhenPhantomReachable(reachabilityTracker, asyncStreamResponse::close) - } override fun subscribe(handler: Handler): AsyncStreamResponse = apply { asyncStreamResponse.subscribe(TrackedHandler(handler, reachabilityTracker)) @@ -37,7 +36,7 @@ internal class PhantomReachableClosingAsyncStreamResponse( override fun onCompleteFuture(): CompletableFuture = asyncStreamResponse.onCompleteFuture() - override fun close() = asyncStreamResponse.close() + override fun close() = closeHandle.close() } /** diff --git a/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingHttpClient.kt b/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingHttpClient.kt index 4d891cc69..c58ce999c 100644 --- a/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingHttpClient.kt +++ b/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingHttpClient.kt @@ -10,9 +10,7 @@ import java.util.concurrent.CompletableFuture * This class ensures the `HttpClient` is closed even if the user forgets to close it. */ internal class PhantomReachableClosingHttpClient(private val httpClient: HttpClient) : HttpClient { - init { - closeWhenPhantomReachable(this, httpClient) - } + private val closeHandle = closeWhenPhantomReachable(this, httpClient) override fun execute(request: HttpRequest, requestOptions: RequestOptions): HttpResponse = httpClient.execute(request, requestOptions) @@ -22,5 +20,5 @@ internal class PhantomReachableClosingHttpClient(private val httpClient: HttpCli requestOptions: RequestOptions, ): CompletableFuture = httpClient.executeAsync(request, requestOptions) - override fun close() = httpClient.close() + override fun close() = closeHandle.close() } diff --git a/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingHttpRequestAuthenticator.kt b/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingHttpRequestAuthenticator.kt index c4679e598..be59abfb5 100644 --- a/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingHttpRequestAuthenticator.kt +++ b/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingHttpRequestAuthenticator.kt @@ -12,9 +12,7 @@ import java.util.concurrent.CompletableFuture internal class PhantomReachableClosingHttpRequestAuthenticator( private val authenticator: HttpRequestAuthenticator ) : HttpRequestAuthenticator { - init { - closeWhenPhantomReachable(this, authenticator) - } + private val closeHandle = closeWhenPhantomReachable(this, authenticator) override fun authenticate(request: HttpRequest): HttpRequest = authenticator.authenticate(request) @@ -22,5 +20,5 @@ internal class PhantomReachableClosingHttpRequestAuthenticator( override fun authenticateAsync(request: HttpRequest): CompletableFuture = authenticator.authenticateAsync(request) - override fun close() = authenticator.close() + override fun close() = closeHandle.close() } diff --git a/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingStreamResponse.kt b/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingStreamResponse.kt index 4321edb49..d59695ece 100644 --- a/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingStreamResponse.kt +++ b/openai-java-core/src/main/kotlin/com/openai/core/http/PhantomReachableClosingStreamResponse.kt @@ -11,11 +11,9 @@ import java.util.stream.Stream internal class PhantomReachableClosingStreamResponse( private val streamResponse: StreamResponse ) : StreamResponse { - init { - closeWhenPhantomReachable(this, streamResponse) - } + private val closeHandle = closeWhenPhantomReachable(this, streamResponse) override fun stream(): Stream = streamResponse.stream() - override fun close() = streamResponse.close() + override fun close() = closeHandle.close() } diff --git a/openai-java-core/src/test/kotlin/com/openai/core/PhantomReachableTest.kt b/openai-java-core/src/test/kotlin/com/openai/core/PhantomReachableTest.kt index 7ca627384..a217f48d8 100644 --- a/openai-java-core/src/test/kotlin/com/openai/core/PhantomReachableTest.kt +++ b/openai-java-core/src/test/kotlin/com/openai/core/PhantomReachableTest.kt @@ -1,5 +1,6 @@ package com.openai.core +import java.util.concurrent.atomic.AtomicInteger import org.assertj.core.api.Assertions.assertThat import org.junit.jupiter.api.Test @@ -24,4 +25,18 @@ internal class PhantomReachableTest { assertThat(closed).isTrue() } + + @Test + fun closeWhenPhantomReachable_explicitHandleClosesAtMostOnce() { + val closeCount = AtomicInteger() + val observed = Any() + val handle = closeWhenPhantomReachable(observed) { closeCount.incrementAndGet() } + + handle.close() + handle.close() + + assertThat(closeCount.get()).isEqualTo(1) + // Keep the observed object strongly reachable until after both explicit closes. + assertThat(observed).isNotNull() + } } diff --git a/openai-java-core/src/test/kotlin/com/openai/core/http/PhantomReachableClosingHttpClientTest.kt b/openai-java-core/src/test/kotlin/com/openai/core/http/PhantomReachableClosingHttpClientTest.kt new file mode 100644 index 000000000..58b8ab72a --- /dev/null +++ b/openai-java-core/src/test/kotlin/com/openai/core/http/PhantomReachableClosingHttpClientTest.kt @@ -0,0 +1,20 @@ +package com.openai.core.http + +import org.junit.jupiter.api.Test +import org.mockito.kotlin.mock +import org.mockito.kotlin.times +import org.mockito.kotlin.verify + +internal class PhantomReachableClosingHttpClientTest { + + @Test + fun close_closesDelegateAtMostOnce() { + val delegate = mock() + val client = PhantomReachableClosingHttpClient(delegate) + + client.close() + client.close() + + verify(delegate, times(1)).close() + } +}