Repository navigation
Keep JDK Unix socket timeouts from stopping trace delivery - #12322
Conversation
Create the `SocketChannel` with `TunnelingJdkSocket` instead of waiting until `connect` so they share the same lifetime. This gives each socket one stable channel and ensures a failed connection closes all associated resources. Okio may close a socket from its _timeout watchdog_ while another thread is _initializing the read selector_ or _performing a half-close_. Coordinate these compound lifecycle operations on the socket monitor: - `getInputStream` publishes the selector before `close` can snapshot it, also it runs once per physical OkHttp connection - `shutdownInput` and `shutdownOutput` complete their channel operation and state update without interleaving with `close`. Those are lifecycle operations. - `close` publishes the terminal state and snapshots the selector atomically, then releases the monitor before closing resources and waking blocked reads. This is a lifecycle operation. - The normal bulk reads and writes only perform volatile state reads. - Previous synchronization on the selection key is unchanged. So, this introduces the socket monitor only for these lifecycle transitions. This keeps normal reads and writes outside the socket monitor. Only selector publication and socket lifecycle transitions require that coordination.
🟢 Java Benchmark SLOs — All performance SLOs passed
PR vs. master results
Commit: Load and DaCapo benchmarks can be triggered manually in the GitLab pipeline. Results will appear in the Benchmarking Platform UI after completion. |
Codex Review SummaryThis comment shows the latest Codex review activity on this pull request.
ℹ️ About Codex in GitHubYour team has set up Codex to review pull requests in this repo. Reviews are triggered when you
Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings. |
|
Review comment from Claude (AI assistant, posted at Doug's request)
The Suggestion: keep I reviewed the diff by reading it only; I did not build or run the tests. |
|
Review comment from Claude (AI assistant, posted at Doug's request)
if (!isClosed() && isConnected() && !isInputShutdown() && selector == null) {
unixSocketChannel.configureBlocking(false);
}
synchronized (this) { ... }If Suggestion: catch Related, minor: the I reviewed the diff by reading it only; I did not build or run the tests. |
|
Review comment from Claude (AI assistant, posted at Doug's request)
Suggestion 1 (this PR): drop the volatile Suggestion 2 (follow-up is fine): collapse the lifecycle flags into one state word. A single Minor: I reviewed the diff by reading it only; I did not build or run the tests, and I did not verify every reader in the unchanged parts of the class. |
There was a problem hiding this comment.
|
@dougqh It seems your comments didn't landed on code, and as such can't start a discussion.
|
1c3bf3e to
255e37f
Compare
255e37f to
1863ea5
Compare
|
@bric3 Re: the JNR fallback: okay, that's fine with me. (Posted by Claude at Doug's request.) |
|
@bric3 Re: the (Posted by Claude at Doug's request.) |
|
@bric3 Re: the volatile fields: okay, I just get a little worried when I see a lot of volatile fields. It can be hard to keep track of the appropriate access order, but I'll leave the final decision to you. (Posted by Claude at Doug's request.) |
dougqh
left a comment
There was a problem hiding this comment.
Overall, looks good to me. As indicated in comments, I'm fine with not changing the areas I/Claude raised yesterday.
sarahchen6
left a comment
There was a problem hiding this comment.
small cleanup comment, but otherwise LGTM!
What Does This Do
Keeps trace delivery alive when an OkHttp request over a JDK Unix-domain socket times out.
This pull request makes two related but distinct changes:
It also adds focused socket tests, a writer-level timeout-then-reconnect test over a real Unix-domain socket, and JMH benchmarks for connection setup and bulk I/O.
Motivation
1. Keep selector closure on the I/O failure path
OkHttp creates the Agent connection with
UnixDomainSocketFactory. Okio then wraps that socket and usesAsyncTimeout.Watchdogto enforce the read timeout. When a response stalls, the watchdog closes the socket from a different thread to interrupt the blocked read.Closing
TunnelingJdkSocketcloses its selector. Depending on where the reader is, JDK NIO can then throwClosedSelectorExceptionorCancelledKeyException. These are unchecked exceptions, so Okio does not handle them as transport failures. They can escape the send path and stop the singledd-trace-processorthread, after which later traces are no longer sent.flowchart LR subgraph DD["dd-trace-java"] worker["TraceProcessingWorker"] api["DDAgentApi"] udsRead["TunnelingJdkSocket.read()"] udsClose["TunnelingJdkSocket.close()"] nextPayload["next payload<br/>fresh connection"] end subgraph OkHttp["OkHttp 3.x"] codec["Http1Codec"] connection["RealConnection"] end subgraph Okio["Okio 1.x"] source["Okio.source(socket)"] watchdog["AsyncTimeout.Watchdog"] end subgraph JDK["JDK NIO"] select["Selector.select()"] lifecycleException["ClosedSelectorException<br/>or CancelledKeyException"] end worker --> api --> codec --> source --> udsRead --> select watchdog -->|"read timeout closes socket"| udsClose udsClose -->|"closes selector"| select select --> lifecycleException lifecycleException -->|"translated to SocketException"| udsRead udsRead -->|"IOException"| source source -->|"failed exchange"| codec codec -->|"failed send; worker survives"| worker worker --> nextPayload --> connection --> udsReadThe socket now translates only these selector lifecycle exceptions to
SocketException. The timed-out POST is not retried because the Agent may already have received it. Instead, that payload fails normally, the processor keeps running, and the next payload uses a fresh connection.2. Serialize socket lifecycle transitions
Tip
Here, serialize means allowing only one thread at a time to validate and update the socket's lifecycle state under its monitor (
synchronized (this)). For example, publishing the selector and taking the snapshot used byclose()cannot interleave. Bulk reads and writes stay outside that monitor.The exception translation makes asynchronous closure survivable, but the socket also needs a coherent lifetime while OkHttp and Okio use it from different threads.
The synchronization is deliberately narrow. It protects compound lifecycle operations rather than normal I/O:
getInputStream()configures the channel outside the monitor because switching blocking mode can wait for an active write. This letsclose()interrupt that write. It then revalidates the state and publishes the lazily created selector under the monitor soclose()can snapshot it.shutdownInput()andshutdownOutput()keep validation, the channel operation, and the state update atomic withclose().close()publishes the terminal state and snapshots the selector under the socket monitor, then closes the selector and channel after releasing the monitor so a blocked read can wake without extending the critical section.flowchart LR subgraph OkHttp["OkHttp 3.x"] realConnection["RealConnection<br/>one socket per physical connection"] halfClose["connection shutdown"] end subgraph Okio["Okio 1.x"] sourceSetup["Okio.source(socket)"] timeoutClose["AsyncTimeout.Watchdog<br/>socket.close()"] end subgraph DD["dd-trace-java"] factory["UnixDomainSocketFactory"] socket["TunnelingJdkSocket"] inputLifecycle["getInputStream()<br/>validate + publish"] halfLifecycle["shutdownInput/Output()<br/>validate + update"] closeLifecycle["close()<br/>terminal state + snapshot"] monitor["synchronized (this)<br/>short lifecycle transition"] resourceClose["close resources<br/>outside monitor"] end subgraph JDK["JDK NIO"] channel["final SocketChannel"] selector["lazy Selector"] end factory --> socket -->|"created together"| channel realConnection -->|"connect"| socket sourceSetup --> inputLifecycle --> monitor -->|"publish"| selector halfClose --> halfLifecycle --> monitor -->|"native half-close"| channel timeoutClose --> closeLifecycle --> monitor -->|"snapshot"| resourceClose resourceClose -->|"wake blocked read"| selector resourceClose --> channelNormal bulk reads and writes do not acquire the socket lifecycle monitor. They continue to observe volatile state and operate on the stable channel.
Additional Notes
The writer-level test uses periodic flushing and MockWebServer's
NO_RESPONSEpolicy over a JNR Unix-domain server adapter. It exercises the worker-ending path: the first request times out, the second request arrives on a new connection, and both callbacks run on the same surviving writer thread. Setup and cleanup I/O failures fail the test.A socket regression test starts a blocked write before input-stream initialization and verifies that
close()interrupts both operations.JMH comparison
Added benchmarks comparing the pre-PR revision (cd0e187) with a0282ae on macOS
aarch64, JDK 17.0.20,Times are means more or less JMH's reported "99.9% confidence error", in microseconds per operation.
Allocations were essentially unchanged. In this results the confidence intervals overlap, which I think doesn;t not establish a regression nor an equivalence. No clear bulk I/O regression was detected in these workloads. Round trip benchmarks reuse the connection and include peer-thread scheduling.
Contributor Checklist
type:and (comp:orinst:) labels in addition to any other useful labelsclose,fix, or any linking keywords when referencing an issueUse
solvesinstead, and assign the PR milestone to the issue