Skip to content

Keep JDK Unix socket timeouts from stopping trace delivery - #12322

Merged
gh-worker-dd-mergequeue-cf854d[bot] merged 8 commits into
masterfrom
bdu/harden-jdk-uds-socket
Sep 30, 2026
Merged

gh-worker-dd-mergequeue-cf854d[bot] merged 8 commits into
masterfrom
bdu/harden-jdk-uds-socket

Conversation

@bric3

@bric3 bric3 commented Aug 27, 2026 •

Copy link
Copy Markdown
Contributor

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:

  1. It reports asynchronous selector closure as an I/O failure instead of letting an unchecked selector exception escape the socket boundary.
  2. It serializes the compound socket lifecycle transitions that can race with Okio closing a timed-out connection.

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 uses AsyncTimeout.Watchdog to enforce the read timeout. When a response stalls, the watchdog closes the socket from a different thread to interrupt the blocked read.

Closing TunnelingJdkSocket closes its selector. Depending on where the reader is, JDK NIO can then throw ClosedSelectorException or CancelledKeyException. These are unchecked exceptions, so Okio does not handle them as transport failures. They can escape the send path and stop the single dd-trace-processor thread, 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 --> udsRead
Loading

The 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 by close() 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 lets close() interrupt that write. It then revalidates the state and publishes the lazily created selector under the monitor so close() can snapshot it.
  • shutdownInput() and shutdownOutput() keep validation, the channel operation, and the state update atomic with close().
  • 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.
  • The channel is created with the socket and remains stable for its lifetime. A failed connection closes that channel and leaves the socket terminally closed.
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 --> channel
Loading

Normal 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_RESPONSE policy 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,

  • two forks
  • 3 warmup iterations,
  • five iterations,
  • GC profiler

Times are means more or less JMH's reported "99.9% confidence error", in microseconds per operation.

Operation Before After Mean change
Setup and close 8.27 +/- 0.97 8.56 +/- 0.33 +3.6%
Round trip, 256 B 5.12 +/- 0.28 5.13 +/- 0.43 +0.3%
Round trip, 8 KiB 9.19 +/- 0.15 9.09 +/- 0.36 -1.0%
Round trip, 64 KiB 68.52 +/- 8.70 68.32 +/- 2.78 -0.3%

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.

./gradlew :utils:socket-utils:jmh -PtestJvm=17 -Pjmh.includes=TunnelingJdkSocketBenchmark -Pjmh.profilers=gc

Contributor Checklist

bric3 added 3 commits August 27, 2026 18:36
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.
@bric3 bric3 added type: bug fix Bug fix comp: core Tracer core tag: ai generated Largely based on code generated by an AI or LLM tag: concurrency Virtual Threads, Coroutines, Async, RX, Executors and removed tag: concurrency Virtual Threads, Coroutines, Async, RX, Executors labels Aug 27, 2026
@datadog-prod-us1-4

datadog-prod-us1-4 Bot commented Aug 27, 2026 •

Copy link
Copy Markdown

🎯 Code Coverage (details)
• Patch Coverage: 100.00%
• Overall Coverage: 59.08% (-0.17%)

This comment will be updated automatically if new data arrives.
🔗 Commit SHA: dc676ab | Docs | Give us feedback!

@dd-octo-sts

dd-octo-sts Bot commented Aug 27, 2026 •

Copy link
Copy Markdown
Contributor

🟢 Java Benchmark SLOs — All performance SLOs passed

Suite Status
Startup 🟢 pass

SLO thresholds are defined here based on automatically generated metrics. A warning is raised when results are within 5% of the threshold.

PR vs. master results
Scenario Candidate master Δ (95% CI of mean)
startup:insecure-bank:iast:Agent 14.04 s 14.11 s [-1.2%; +0.3%] (no difference)
startup:insecure-bank:tracing:Agent 12.96 s 13.08 s [-1.7%; -0.1%] (maybe better)
startup:petclinic:appsec:Agent 16.91 s 16.79 s [-0.5%; +1.9%] (no difference)
startup:petclinic:iast:Agent 16.88 s 16.98 s [-1.5%; +0.3%] (no difference)
startup:petclinic:profiling:Agent 16.56 s 16.42 s [-3.9%; +5.6%] (no difference)
startup:petclinic:sca:Agent 16.55 s 16.74 s [-5.6%; +3.4%] (no difference)
startup:petclinic:tracing:Agent 15.77 s 15.90 s [-6.6%; +5.0%] (unstable)

Commit: dc676ab0 · CI Pipeline · Benchmarking Platform UI


Load and DaCapo benchmarks can be triggered manually in the GitLab pipeline. Results will appear in the Benchmarking Platform UI after completion.

@bric3
bric3 marked this pull request as ready for review September 28, 2026 20:39
@bric3
bric3 requested review from a team as code owners September 28, 2026 20:39
@bric3
bric3 requested review from dougqh and removed request for a team September 28, 2026 20:39
@chatgpt-codex-connector

chatgpt-codex-connector Bot commented Sep 28, 2026 •

Copy link
Copy Markdown

Codex Review Summary

This comment shows the latest Codex review activity on this pull request.

Review Status Commit Review trigger
📝 Code Review ✅ Completed 2026-09-28T20:43:17.191224Z d1b6434 Draft marked ready
🔒 Security Review ✅ Completed 2026-09-28T20:43:36.758451Z d1b6434 Draft marked ready
ℹ️ About Codex in GitHub

Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you

  • Open a pull request for review
  • Mark a draft as ready
  • Comment "@codex review" or "@codex security review".

Codex reacts with 👀 while any review is running, comments if it has suggestions, and reacts with 👍 once all reviews finish with no findings.

@bric3
bric3 requested a review from sarahchen6 September 28, 2026 20:40
@dougqh

dougqh commented Sep 28, 2026

Copy link
Copy Markdown
Contributor

Review comment from Claude (AI assistant, posted at Doug's request)

UnixDomainSocketFactory.createSocket(): narrowing catch (Throwable ignore) to catch (IOException | UnsupportedOperationException ignore) skips the JNR fallback for more failures than intended.

The TunnelingJdkSocket constructor now does real work (path.toPath() and SocketChannel.open(StandardProtocolFamily.UNIX)). Failures such as SecurityException, InvalidPathException, or a LinkageError / ServiceConfigurationError from the NIO provider used to fall through to TunnelingUnixSocket (JNR). Now they escape the inner catch and land in the outer catch (Throwable e). That handler either fails over to the TCP port (when agentConfiguredUsingDefault) or rethrows, and it never tries JNR. That contradicts the // fall back to jnr-unixsocket library comment.

Suggestion: keep catch (Throwable ignore) here, or at minimum widen to IOException | RuntimeException. The narrowing doesn't seem to be needed for the goal of this PR (selector-closure handling and lifecycle serialization).

I reviewed the diff by reading it only; I did not build or run the tests.

@dougqh

dougqh commented Sep 28, 2026

Copy link
Copy Markdown
Contributor

Review comment from Claude (AI assistant, posted at Doug's request)

TunnelingJdkSocket.getInputStream(): the pre-check and configureBlocking(false) run outside the monitor, so a concurrent close() can surface a different exception than the locked block would.

if (!isClosed() && isConnected() && !isInputShutdown() && selector == null) {
  unixSocketChannel.configureBlocking(false);
}
synchronized (this) { ... }

If close() lands between the check and configureBlocking, the caller gets a raw ClosedChannelException instead of the SocketException("Socket is closed") thrown from the locked section. It is still an IOException, so OkHttp handles it, but it is inconsistent with the rest of the class and with the PR's goal of translating close races into SocketException.

Suggestion: catch ClosedChannelException around configureBlocking and rethrow as SocketException("Socket is closed") (with the cause attached).

Related, minor: the selector == null check is unlocked, so two concurrent first callers can both call configureBlocking(false). That is idempotent and harmless, but worth a short comment so a future reader doesn't "fix" it.

I reviewed the diff by reading it only; I did not build or run the tests.

@dougqh

dougqh commented Sep 28, 2026

Copy link
Copy Markdown
Contributor

Review comment from Claude (AI assistant, posted at Doug's request)

TunnelingJdkSocket: the number of volatile fields makes the memory-ordering story hard to follow. Reading the diff, they fall into three groups:

  • Written under the monitor, read unlocked: closed, shutIn, shutOut. close() writes them in one block (shutIn, shutOut, closed), and shutdownInput()/shutdownOutput() write them under the same lock. Unlocked readers only need the latest value and there is no multi-field invariant. A reader seeing shutIn == true while closed is still false just gets -1 slightly early, which is benign.
  • Written once (null to non-null): inetSocketAddress and selector. selector is never reset now that close() no longer nulls it.
  • Independent latest-value: timeout. Volatile is fine; the cost is negligible next to select.

Suggestion 1 (this PR): drop the volatile selector requirement for the read path. getInputStream() already publishes the selector inside synchronized. The anonymous InputStream can capture the local currentSelector as a final field instead of re-reading this.selector in read(). Final-field semantics then give safe publication, and the currentSelector == null branch in read() becomes dead code. close() already reads selector under the lock, so the only remaining reader that forces volatile is the @VisibleForTesting access; a test-only getter that takes the lock would cover that.

Suggestion 2 (follow-up is fine): collapse the lifecycle flags into one state word. A single volatile int state (NEW, CONNECTED, CLOSED, plus half-close bits), written only under the monitor, would replace closed, shutIn, shutOut and the "connected" role currently played by inetSocketAddress != null. That gives one ordering story instead of four interacting volatiles. It would also make connect() lifecycle-synchronized, which addresses the fact that connect() is currently unsynchronized and two concurrent calls can reach SocketChannel.connect, which throws the unchecked AlreadyConnectedException outside the translation added here. The cost is rewriting the isX() helpers, so it may be better as a separate PR than folded into this one.

Minor: sendBufferSize and receiveBufferSize stayed non-volatile while everything around them became volatile. Either make them volatile or note why they don't need to be.

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.

@datadog-prod-us1-4 datadog-prod-us1-4 Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Bits Code Review: PASS

More details

The revised lifecycle publishes terminal state before closing resources and translates selector-close races into SocketException, keeping timeouts on the writer’s recoverable I/O path.

Was this helpful? React 👍 or 👎

Open Bits AI session

🤖 Bits Code Review · Commit d1b6434 · @DataDog review to ask questions

bric3 commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor Author

@dougqh It seems your comments didn't landed on code, and as such can't start a discussion.

  1. On the JNR fallback: I deliberately narrowed the exception as Throwable is too wide and do not rethrow other unrelated exception/errors.

  2. On the getInputStream() close race, I agree that a concurrent close() can make configureBlocking(false) throw ClosedChannelException. But it is already an IOException, so OkHttp treats it as a transport failure. The call stays outside the socket monitor because it can wait for an active write, while close() must be able to interrupt that write. Wrapping the exception as SocketException would not change recovery.

  3. On the volatile fields, I looked at both suggestions.

    1. Capturing selector in the stream would work, but there's already a local copy which is sufficient in this case. OkHttp creates one InputStream per physical connection, where each read ends up calling select().
      That being said, removing one volatile is useful so I di go with it

    2. On the idea of state, I believe this will over-engineer for a case that is not needed. Each OkHttp connection owns its socket (connect() is called only once), the socket fields describe its lifecycle properties already, and method reading those do not need those properties combined into one atomic state.

@bric3
bric3 force-pushed the bdu/harden-jdk-uds-socket branch from 1c3bf3e to 255e37f Compare September 29, 2026 12:33
@bric3
bric3 force-pushed the bdu/harden-jdk-uds-socket branch from 255e37f to 1863ea5 Compare September 29, 2026 13:11
@dougqh

dougqh commented Sep 29, 2026

Copy link
Copy Markdown
Contributor

@bric3 Re: the JNR fallback: okay, that's fine with me.

(Posted by Claude at Doug's request.)

@dougqh

dougqh commented Sep 29, 2026

Copy link
Copy Markdown
Contributor

@bric3 Re: the getInputStream() close race: okay, that's fine with me.

(Posted by Claude at Doug's request.)

@dougqh

dougqh commented Sep 29, 2026

Copy link
Copy Markdown
Contributor

@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 dougqh left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Overall, looks good to me. As indicated in comments, I'm fine with not changing the areas I/Claude raised yesterday.

@sarahchen6 sarahchen6 left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

small cleanup comment, but otherwise LGTM!

@gh-worker-dd-mergequeue-cf854d
gh-worker-dd-mergequeue-cf854d Bot merged commit bd24f91 into master Sep 30, 2026
606 checks passed
@gh-worker-dd-mergequeue-cf854d
gh-worker-dd-mergequeue-cf854d Bot deleted the bdu/harden-jdk-uds-socket branch September 30, 2026 14:29
@github-actions github-actions Bot added this to the 1.67.0 milestone Sep 30, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

comp: core Tracer core tag: ai generated Largely based on code generated by an AI or LLM type: bug fix Bug fix

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants