Add tests reproducing stale TCP connection issue (#954) - #955
Conversation
TCPSenderTest documents the RST and FIN reconnection behavior at the sender level. FluencyTestWithMockServer demonstrates that ACK mode delivers all records after a simulated Fluentd restart that drops all TCP connections.
There was a problem hiding this comment.
Code Review
This pull request adds new tests to reproduce and verify connection recovery behavior under issue #954, including ACK mode delivery guarantees during connection drops in FluencyTestWithMockServer.java and TCP reconnection behavior after RST and FIN close scenarios in TCPSenderTest.java. The feedback points out a potential resource leak in the new TCPSenderTest helper method, recommending that the mock server and sender be closed within try-finally and try-with-resources blocks to ensure proper cleanup if assertions fail.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
There was a problem hiding this comment.
Pull request overview
This PR adds regression/reproducer tests around stale TCP connections after a Fluentd restart (issue #954), with both low-level TCPSender behavior documentation and an end-to-end Fluency/ACK-mode verification using a mock server.
Changes:
- Add
TCPSenderTestcases covering reconnect behavior when the server closes sockets via RST vs FIN. - Add an end-to-end test in
FluencyTestWithMockServershowing ACK mode delivers records after the server drops existing TCP connections.
Reviewed changes
Copilot reviewed 2 out of 2 changed files in this pull request and generated 3 comments.
| File | Description |
|---|---|
| fluency-fluentd/src/test/java/org/komamitsu/fluency/fluentd/ingester/sender/TCPSenderTest.java | Adds RST/FIN reproducer tests and explanatory Javadoc around stale-connection behavior. |
| fluency-fluentd/src/test/java/org/komamitsu/fluency/fluentd/FluencyTestWithMockServer.java | Adds end-to-end ACK-mode test that simulates a server-side connection drop and verifies delivery. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
- Wrap sender in try-with-resources and server.stop() in finally to prevent resource leaks on assertion failure - Assert connectCount > 1 to verify an actual reconnect occurred, not just consecutive successes on the stale socket - Fix Javadoc cross-reference to point to FluencyTestWithMockServer
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request adds test coverage to reproduce and verify issue #954, introducing tests for connection drops and reconnection behavior under ACK mode, RST, and FIN scenarios. The reviewer suggested a minor improvement in TCPSenderTest.java to extract firstAcceptedSocket.get() into a local variable and assert that it is not null before performing operations on it, avoiding multiple dereferences.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| if (abruptClose) { | ||
| // RST path: SO_LINGER with timeout=0 makes close() send RST instead of FIN | ||
| firstAcceptedSocket.get().setSoLinger(true, 0); | ||
| } | ||
| firstAcceptedSocket.get().close(); |
There was a problem hiding this comment.
To adhere to defensive programming practices and avoid multiple dereferences of the AtomicReference, it is recommended to retrieve the socket into a local variable, assert that it is not null, and then perform the operations on it.
| if (abruptClose) { | |
| // RST path: SO_LINGER with timeout=0 makes close() send RST instead of FIN | |
| firstAcceptedSocket.get().setSoLinger(true, 0); | |
| } | |
| firstAcceptedSocket.get().close(); | |
| Socket socket = firstAcceptedSocket.get(); | |
| assertThat(socket).isNotNull(); | |
| if (abruptClose) { | |
| // RST path: SO_LINGER with timeout=0 makes close() send RST instead of FIN | |
| socket.setSoLinger(true, 0); | |
| } | |
| socket.close(); |
| int consecutiveSuccesses = 0; | ||
| for (int attempt = 0; attempt < 10 && consecutiveSuccesses < 3; attempt++) { | ||
| try { | ||
| sender.send(ByteBuffer.wrap(data)); | ||
| LOG.debug("Send succeeded on attempt {}", attempt); | ||
| if (connectCount.get() > 1) { | ||
| consecutiveSuccesses++; | ||
| } else { | ||
| consecutiveSuccesses = 0; | ||
| } | ||
| } catch (IOException e) { | ||
| LOG.debug("Attempt {} failed (expected on stale socket): {}", attempt, e.getMessage()); | ||
| consecutiveSuccesses = 0; | ||
| } | ||
| } |
| for (int i = 0; i < recordsBeforeDrop; i++) { | ||
| fluency.emit("tag", data); | ||
| } | ||
| fluency.waitUntilAllBufferFlushed(10); |
| for (int i = 0; i < recordsAfterDrop; i++) { | ||
| fluency.emit("tag", data); | ||
| } | ||
| fluency.waitUntilAllBufferFlushed(30); |
- Replace tight attempt loop with a CountDownLatch on the second onConnect event so the reconnect assertion is timing-independent - Add 100ms sleep between send attempts to give the OS time to deliver FIN/RST before the loop exhausts its retries - Extract firstAcceptedSocket.get() to a local variable with null check - Assert waitUntilAllBufferFlushed() return value so flush timeouts surface as clear failures rather than record-count mismatches
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces integration and unit tests to reproduce and verify behavior related to connection drops (issue #954). Specifically, it adds a test in FluencyTestWithMockServer to ensure that ACK mode guarantees message delivery when the server drops connections. Additionally, it adds tests in TCPSenderTest to verify that TCPSender successfully reconnects after both abrupt (RST) and graceful (FIN) connection closures by the server. As there are no review comments, I have no feedback to provide.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| * <p>Scenario B (graceful / FIN): The server closes normally, sending FIN. Due to TCP half-close | ||
| * the first write may silently "succeed" — the OS places data in the send buffer before the RST | ||
| * triggered by the peer's close arrives. That data is silently lost at the TCPSender level. The | ||
| * second write fails ("Broken pipe"), closeSocket() nulls the channel, and the third send | ||
| * reconnects successfully. |
| assertTrue( | ||
| fluency.waitUntilAllBufferFlushed(10), "Buffer should flush before dropping connections"); | ||
|
|
||
| LOG.info("Dropping {} connections to simulate Fluentd restart", acceptedSockets.size()); |
- Clarify Javadoc: silent data loss in the FIN case is a potential concern; this test only verifies reconnection, not data loss - Assert acceptedSockets is non-empty before dropping connections to prevent the drop step from being a silent no-op
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces new integration and unit tests to verify reconnection behavior and message delivery guarantees under connection drop scenarios (both graceful FIN and abrupt RST) for ACK mode and TCPSender. The reviewer feedback highlights critical resource leak risks in both test files where server.start() is called outside the try-finally block, meaning a startup failure would leak background threads and sockets. Additionally, the reviewer suggests increasing the retry limit in the TCPSender test to avoid flakiness in slower CI environments.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| server.start(); | ||
|
|
||
| int recordsBeforeDrop = 200; | ||
| int recordsAfterDrop = 200; | ||
|
|
||
| FluencyBuilderForFluentd builder = new FluencyBuilderForFluentd(); | ||
| builder.setAckResponseMode(true); | ||
| builder.setFlushAttemptIntervalMillis(200); | ||
|
|
||
| try (Fluency fluency = builder.build(server.getLocalPort())) { | ||
| Map<String, Object> data = new HashMap<>(); | ||
| data.put("key", "value"); | ||
|
|
||
| for (int i = 0; i < recordsBeforeDrop; i++) { | ||
| fluency.emit("tag", data); | ||
| } | ||
| assertTrue( | ||
| fluency.waitUntilAllBufferFlushed(10), "Buffer should flush before dropping connections"); | ||
|
|
||
| assertThat(acceptedSockets) | ||
| .as("At least one connection must have been established") | ||
| .isNotEmpty(); | ||
| LOG.info("Dropping {} connections to simulate Fluentd restart", acceptedSockets.size()); | ||
| for (Socket socket : acceptedSockets) { | ||
| try { | ||
| socket.close(); | ||
| } catch (IOException e) { | ||
| LOG.warn("Failed to close socket", e); | ||
| } | ||
| } | ||
|
|
||
| for (int i = 0; i < recordsAfterDrop; i++) { | ||
| fluency.emit("tag", data); | ||
| } | ||
| assertTrue(fluency.waitUntilAllBufferFlushed(30), "Buffer should flush after reconnection"); | ||
| } finally { | ||
| server.stop(); | ||
| } |
There was a problem hiding this comment.
If server.start() throws an exception (e.g., due to port binding issues or timeout) or if FluencyBuilderForFluentd initialization fails, the finally block is never reached. This will leak the background thread pool and server socket of AbstractFluentdServer, potentially causing subsequent tests to fail due to port conflicts or resource exhaustion.
Wrapping the server startup and builder initialization in the outer try-finally block ensures that resources are always cleaned up properly.
try {
server.start();
int recordsBeforeDrop = 200;
int recordsAfterDrop = 200;
FluencyBuilderForFluentd builder = new FluencyBuilderForFluentd();
builder.setAckResponseMode(true);
builder.setFlushAttemptIntervalMillis(200);
try (Fluency fluency = builder.build(server.getLocalPort())) {
Map<String, Object> data = new HashMap<>();
data.put("key", "value");
for (int i = 0; i < recordsBeforeDrop; i++) {
fluency.emit("tag", data);
}
assertTrue(
fluency.waitUntilAllBufferFlushed(10), "Buffer should flush before dropping connections");
assertThat(acceptedSockets)
.as("At least one connection must have been established")
.isNotEmpty();
LOG.info("Dropping {} connections to simulate Fluentd restart", acceptedSockets.size());
for (Socket socket : acceptedSockets) {
try {
socket.close();
} catch (IOException e) {
LOG.warn("Failed to close socket", e);
}
}
for (int i = 0; i < recordsAfterDrop; i++) {
fluency.emit("tag", data);
}
assertTrue(fluency.waitUntilAllBufferFlushed(30), "Buffer should flush after reconnection");
}
} finally {
server.stop();
}| server.start(); | ||
|
|
||
| TCPSender.Config config = new TCPSender.Config(); | ||
| config.setPort(server.getLocalPort()); | ||
|
|
||
| try (TCPSender sender = new TCPSender(config)) { | ||
| byte[] data = "hello".getBytes(StandardCharsets.UTF_8); | ||
|
|
||
| // Establish the connection by sending initial data | ||
| sender.send(ByteBuffer.wrap(data)); | ||
| assertTrue(firstDataReceivedLatch.await(5, TimeUnit.SECONDS)); | ||
|
|
||
| Socket socket = firstAcceptedSocket.get(); | ||
| assertThat(socket).isNotNull(); | ||
| if (abruptClose) { | ||
| // RST path: SO_LINGER with timeout=0 makes close() send RST instead of FIN | ||
| socket.setSoLinger(true, 0); | ||
| } | ||
| socket.close(); | ||
|
|
||
| // Keep sending until the server observes the reconnect (second onConnect). | ||
| // - RST path: the first send fails ("Connection reset"), the next send reconnects. | ||
| // - FIN path: sends may silently "succeed" (TCP half-close) before the RST triggers | ||
| // closeSocket(); after that a send fails ("Broken pipe") and the next reconnects. | ||
| for (int attempt = 0; attempt < 20 && reconnectedLatch.getCount() > 0; attempt++) { | ||
| try { | ||
| sender.send(ByteBuffer.wrap(data)); | ||
| LOG.debug("Send succeeded on attempt {}", attempt); | ||
| } catch (IOException e) { | ||
| LOG.debug("Attempt {} failed (expected on stale socket): {}", attempt, e.getMessage()); | ||
| } | ||
| TimeUnit.MILLISECONDS.sleep(100); | ||
| } | ||
|
|
||
| assertTrue( | ||
| reconnectedLatch.await(10, TimeUnit.SECONDS), | ||
| "TCPSender should reconnect after the server closed the connection"); | ||
|
|
||
| // Verify stable operation on the new connection | ||
| int consecutiveSuccesses = 0; | ||
| for (int attempt = 0; attempt < 10 && consecutiveSuccesses < 3; attempt++) { | ||
| try { | ||
| sender.send(ByteBuffer.wrap(data)); | ||
| consecutiveSuccesses++; | ||
| } catch (IOException e) { | ||
| consecutiveSuccesses = 0; | ||
| } | ||
| } | ||
| assertTrue( | ||
| consecutiveSuccesses >= 3, | ||
| "TCPSender should reach stable operation on the new connection"); | ||
| } finally { | ||
| server.stop(); | ||
| } |
There was a problem hiding this comment.
This addresses two issues:
- Resource Leak: If
server.start()or configuration throws an exception, thefinallyblock is never reached, leaking the background thread pool and socket ofMockTCPServer. - Test Flakiness: In slow or resource-constrained CI environments, the reconnection might take longer than 2 seconds. Since
TCPSenderonly attempts reconnection lazily whensend()is called, stopping thesend()loop after 20 attempts (2 seconds) means no further reconnection attempts will be made, causing the subsequentreconnectedLatch.await(10, TimeUnit.SECONDS)to time out and fail. Increasing the maximum attempts to 100 (10 seconds) matches the latch timeout and prevents flakiness without slowing down successful runs.
try {
server.start();
TCPSender.Config config = new TCPSender.Config();
config.setPort(server.getLocalPort());
try (TCPSender sender = new TCPSender(config)) {
byte[] data = "hello".getBytes(StandardCharsets.UTF_8);
// Establish the connection by sending initial data
sender.send(ByteBuffer.wrap(data));
assertTrue(firstDataReceivedLatch.await(5, TimeUnit.SECONDS));
Socket socket = firstAcceptedSocket.get();
assertThat(socket).isNotNull();
if (abruptClose) {
// RST path: SO_LINGER with timeout=0 makes close() send RST instead of FIN
socket.setSoLinger(true, 0);
}
socket.close();
// Keep sending until the server observes the reconnect (second onConnect).
// - RST path: the first send fails ("Connection reset"), the next send reconnects.
// - FIN path: sends may silently "succeed" (TCP half-close) before the RST triggers
// closeSocket(); after that a send fails ("Broken pipe") and the next reconnects.
for (int attempt = 0; attempt < 100 && reconnectedLatch.getCount() > 0; attempt++) {
try {
sender.send(ByteBuffer.wrap(data));
LOG.debug("Send succeeded on attempt {}", attempt);
} catch (IOException e) {
LOG.debug("Attempt {} failed (expected on stale socket): {}", attempt, e.getMessage());
}
TimeUnit.MILLISECONDS.sleep(100);
}
assertTrue(
reconnectedLatch.await(10, TimeUnit.SECONDS),
"TCPSender should reconnect after the server closed the connection");
// Verify stable operation on the new connection
int consecutiveSuccesses = 0;
for (int attempt = 0; attempt < 10 && consecutiveSuccesses < 3; attempt++) {
try {
sender.send(ByteBuffer.wrap(data));
consecutiveSuccesses++;
} catch (IOException e) {
consecutiveSuccesses = 0;
}
}
assertTrue(
consecutiveSuccesses >= 3,
"TCPSender should reach stable operation on the new connection");
}
} finally {
server.stop();
}| // Keep sending until the server observes the reconnect (second onConnect). | ||
| // - RST path: the first send fails ("Connection reset"), the next send reconnects. | ||
| // - FIN path: sends may silently "succeed" (TCP half-close) before the RST triggers | ||
| // closeSocket(); after that a send fails ("Broken pipe") and the next reconnects. | ||
| for (int attempt = 0; attempt < 20 && reconnectedLatch.getCount() > 0; attempt++) { | ||
| try { | ||
| sender.send(ByteBuffer.wrap(data)); | ||
| LOG.debug("Send succeeded on attempt {}", attempt); | ||
| } catch (IOException e) { | ||
| LOG.debug("Attempt {} failed (expected on stale socket): {}", attempt, e.getMessage()); | ||
| } | ||
| TimeUnit.MILLISECONDS.sleep(100); | ||
| } | ||
|
|
||
| assertTrue( | ||
| reconnectedLatch.await(10, TimeUnit.SECONDS), | ||
| "TCPSender should reconnect after the server closed the connection"); |
Wrap the body after server.start() in try/finally in both TCPSenderTest and FluencyTestWithMockServer so the server is always stopped if builder.build() or TCPSender construction throws.
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request adds tests to verify reconnection behavior and guarantee message delivery in ACK mode when the server drops connections, addressing issue #954. The review feedback suggests increasing the retry attempts in the TCPSenderTest reconnection loop from 20 to 100 to prevent potential test flakiness in slow CI environments.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| // - RST path: the first send fails ("Connection reset"), the next send reconnects. | ||
| // - FIN path: sends may silently "succeed" (TCP half-close) before the RST triggers | ||
| // closeSocket(); after that a send fails ("Broken pipe") and the next reconnects. | ||
| for (int attempt = 0; attempt < 20 && reconnectedLatch.getCount() > 0; attempt++) { |
There was a problem hiding this comment.
In slow or resource-constrained CI environments (such as GitHub Actions), 20 attempts with a 100ms sleep (totaling 2 seconds) might not be sufficient for the TCP handshake and reconnection to complete. Since TCPSender only reconnects when a send operation is actively attempted, if the loop finishes before the reconnection is established, the test will fail.\n\nIncreasing the maximum number of attempts to 100 (allowing up to 10 seconds) to match the subsequent reconnectedLatch.await(10, TimeUnit.SECONDS) timeout will make the test much more resilient and prevent flakiness.
| for (int attempt = 0; attempt < 20 && reconnectedLatch.getCount() > 0; attempt++) { | |
| for (int attempt = 0; attempt < 100 && reconnectedLatch.getCount() > 0; attempt++) { |
There was a problem hiding this comment.
Code Review
This pull request introduces new tests to reproduce and verify reconnection behavior and delivery guarantees after server connection drops (addressing issue #954). Specifically, it adds an ACK mode delivery test in FluencyTestWithMockServer and reconnection tests simulating both abrupt (RST) and graceful (FIN) connection closures in TCPSenderTest. The review feedback suggests moving the server.start() calls inside the try blocks in both test files to guarantee proper resource cleanup and server shutdown in case start() fails.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| server.start(); | ||
| try { |
There was a problem hiding this comment.
To ensure that the mock server is properly stopped and resources (like threads and ports) are cleaned up even if server.start() fails or throws an exception during test execution, it is recommended to move server.start() inside the try block.
| server.start(); | |
| try { | |
| try { | |
| server.start(); |
| server.start(); | ||
| try { |
There was a problem hiding this comment.
| * <p>Scenario B (graceful / FIN): The server closes normally, sending FIN. Due to TCP half-close | ||
| * the first write may silently "succeed" — the OS places data in the send buffer before the RST | ||
| * triggered by the peer's close arrives. That data may be silently lost at the TCPSender level | ||
| * (this test only verifies reconnection, not data loss). The second write fails ("Broken pipe"), | ||
| * closeSocket() nulls the channel, and the third send reconnects successfully. |
- Move server.start() inside the try block in both test files so the finally clause runs server.stop() even if start() throws - Soften the FIN-path Javadoc: "one or more writes may silently succeed" rather than claiming exactly second/third writes fail/reconnect
Ensures the assertion runs (and reports the correct failure) even if fluency.close() or other cleanup throws an exception.
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces new tests to verify connection recovery and delivery guarantees when the Fluentd server drops TCP connections, addressing issue #954. Specifically, it adds tests to ensure ACK mode guarantees delivery after connection drops, and that TCPSender successfully reconnects after both abrupt (RST) and graceful (FIN) server-side closures. The reviewer suggested adding a small backoff delay during the stability check's failure path in TCPSenderTest to prevent rapid exhaustion of the retry budget on transient failures, making the test more robust in CI environments.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| } catch (IOException e) { | ||
| LOG.debug( | ||
| "Send failed in stability check (failure {}): {}", failureCount, e.getMessage()); | ||
| consecutiveSuccesses = 0; | ||
| failureCount++; | ||
| } |
There was a problem hiding this comment.
During the stability check, if a send attempt fails due to a transient socket or network state, retrying immediately without any delay can rapidly exhaust the failureCount budget (up to 10 attempts) in a few microseconds. Adding a small backoff delay (e.g., 100 milliseconds) inside the catch block will make the test significantly more robust and less prone to flakiness in slow or resource-constrained CI environments.
} catch (IOException e) {
LOG.debug(
"Send failed in stability check (failure {}): {}", failureCount, e.getMessage());
consecutiveSuccesses = 0;
failureCount++;
try {
TimeUnit.MILLISECONDS.sleep(100);
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new IOException("Interrupted during stability check backoff", ie);
}
}Reconnection is already confirmed by the latch before this point, so failures here indicate a real problem. No tolerance needed.
ACK mode guarantees at-least-once, but in this test all pre-drop records are fully ACKed before connections are dropped, so duplicates should never occur. Assert that invariant explicitly.
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces reproducer tests for issue #954. Specifically, it adds a test to verify that ACK mode guarantees delivery after connection drops, and adds reconnection tests to TCPSenderTest for both abrupt (RST) and graceful (FIN) connection closures. The review feedback recommends capturing any exceptions thrown in the mock server's background thread using an AtomicReference and asserting that no errors occurred at the end of the test, ensuring that background thread failures are properly propagated to the main test thread.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| void testAckModeGuaranteesDeliveryAfterServerDropsConnections() throws Exception { | ||
| Set<Integer> receivedIds = ConcurrentHashMap.newKeySet(); | ||
| List<Socket> acceptedSockets = new CopyOnWriteArrayList<>(); | ||
| Value idKey = ValueFactory.newString("id"); | ||
|
|
||
| AbstractFluentdServer server = | ||
| new AbstractFluentdServer(false) { | ||
| @Override | ||
| protected EventHandler getFluentdEventHandler() { | ||
| return new EventHandler() { | ||
| @Override | ||
| public void onConnect(Socket socket) { | ||
| acceptedSockets.add(socket); | ||
| } | ||
|
|
||
| @Override | ||
| public void onReceive(String tag, long timestampMillis, MapValue data) { | ||
| Value idValue = data.map().get(idKey); | ||
| if (idValue != null) { | ||
| int id = idValue.asIntegerValue().asInt(); | ||
| if (!receivedIds.add(id)) { | ||
| throw new AssertionError("Duplicate record received: id=" + id); | ||
| } | ||
| } | ||
| } | ||
|
|
||
| @Override | ||
| public void onClose(Socket socket) {} | ||
| }; | ||
| } | ||
| }; |
There was a problem hiding this comment.
The AssertionError thrown inside the mock server's onReceive method (line 354) runs in a background thread pool (serverExecutorService). In JUnit, unhandled exceptions or errors in background threads do not directly fail the main test thread; instead, they are caught and logged by the server task runner, and the test will only fail later due to a timeout or size mismatch. This makes debugging duplicate record issues extremely difficult.
To fix this, we can capture any exception thrown in the background thread using an AtomicReference<Throwable> and assert that it is null at the end of the test (before the size assertion on line 406).
void testAckModeGuaranteesDeliveryAfterServerDropsConnections() throws Exception {
Set<Integer> receivedIds = ConcurrentHashMap.newKeySet();
List<Socket> acceptedSockets = new CopyOnWriteArrayList<>();
Value idKey = ValueFactory.newString("id");
AtomicReference<Throwable> serverError = new AtomicReference<>();
AbstractFluentdServer server =
new AbstractFluentdServer(false) {
@Override
protected EventHandler getFluentdEventHandler() {
return new EventHandler() {
@Override
public void onConnect(Socket socket) {
acceptedSockets.add(socket);
}
@Override
public void onReceive(String tag, long timestampMillis, MapValue data) {
try {
Value idValue = data.map().get(idKey);
if (idValue != null) {
int id = idValue.asIntegerValue().asInt();
if (!receivedIds.add(id)) {
throw new AssertionError("Duplicate record received: id=" + id);
}
}
} catch (Throwable t) {
serverError.compareAndSet(null, t);
throw t;
}
}
@Override
public void onClose(Socket socket) {}
};
}
};| assertThat(receivedIds) | ||
| .as( | ||
| "ACK mode must deliver all %d distinct records despite connection drops", | ||
| recordsBeforeDrop + recordsAfterDrop) | ||
| .hasSize(recordsBeforeDrop + recordsAfterDrop); |
There was a problem hiding this comment.
Assert that no background thread errors occurred in the mock server before asserting the received IDs size. This ensures that any duplicate record assertion failures are immediately propagated to the main test thread with the actual error message.
assertThat(serverError.get()).isNull();
assertThat(receivedIds)
.as(
"ACK mode must deliver all %d distinct records despite connection drops",
recordsBeforeDrop + recordsAfterDrop)
.hasSize(recordsBeforeDrop + recordsAfterDrop);| Value idValue = data.map().get(idKey); | ||
| if (idValue != null) { | ||
| int id = idValue.asIntegerValue().asInt(); | ||
| if (!receivedIds.add(id)) { | ||
| throw new AssertionError("Duplicate record received: id=" + id); | ||
| } | ||
| } |
AbstractFluentdServer's FluentdAcceptTask catches all Throwables, so throwing AssertionError from onReceive was silently swallowed and never failed the test. Capture it in an AtomicReference and rethrow in the main thread after flushing.
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request adds new tests to verify reconnection behavior and delivery guarantees when connections are dropped, addressing issue #954. Specifically, it introduces a test for ACK mode delivery guarantees after server connection drops in FluencyTestWithMockServer, and tests for TCPSender reconnection after abrupt (RST) and graceful (FIN) socket closures in TCPSenderTest. The review feedback suggests wrapping the background event handler's message processing in a try-catch block to capture unexpected exceptions and prevent silent test failures.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| public void onReceive(String tag, long timestampMillis, MapValue data) { | ||
| Value idValue = data.map().get(idKey); | ||
| if (idValue != null) { | ||
| int id = idValue.asIntegerValue().asInt(); | ||
| if (!receivedIds.add(id)) { | ||
| backgroundError.compareAndSet( | ||
| null, new AssertionError("Duplicate record received: id=" + id)); | ||
| } | ||
| } | ||
| } |
There was a problem hiding this comment.
In the background event handler, any unexpected exceptions (such as NullPointerException or ClassCastException) during onReceive will be caught by the server task and logged, but they won't fail the test directly. This can lead to silent failures or generic size mismatch assertions that are hard to debug. Wrapping the processing in a try-catch block and capturing any Throwable in backgroundError ensures that the actual root cause is propagated and fails the test immediately.
public void onReceive(String tag, long timestampMillis, MapValue data) {
try {
Value idValue = data.map().get(idKey);
if (idValue != null) {
int id = idValue.asIntegerValue().asInt();
if (!receivedIds.add(id)) {
backgroundError.compareAndSet(
null, new AssertionError("Duplicate record received: id=" + id));
}
}
} catch (Throwable t) {
backgroundError.compareAndSet(
null, new AssertionError("Unexpected error in background event handler", t));
}
}| for (Socket socket : acceptedSockets) { | ||
| try { | ||
| socket.close(); | ||
| } catch (IOException e) { | ||
| LOG.warn("Failed to close socket", e); | ||
| } | ||
| } |
A single TCPSender to a healthy server always establishes exactly one connection before the flush completes. hasSize(1) makes the precondition explicit and catches unexpected reconnections.
|
/gemini review |
There was a problem hiding this comment.
Code Review
This pull request introduces new test cases to reproduce and verify connection recovery behavior under issue #954. Specifically, it adds tests in FluencyTestWithMockServer to ensure ACK mode guarantees delivery after server connection drops, and in TCPSenderTest to verify reconnection behavior after graceful (FIN) and abrupt (RST) server closures. The review feedback suggests improving test robustness by using server.stop(true) to prevent socket leaks and increasing a latch timeout from 5 to 10 seconds to avoid flaky CI runs.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| recordsBeforeDrop + recordsAfterDrop) | ||
| .hasSize(recordsBeforeDrop + recordsAfterDrop); | ||
| } finally { | ||
| server.stop(); |
There was a problem hiding this comment.
Calling server.stop() (which defaults to server.stop(false)) does not immediately close the accepted sockets or the server socket if the server thread is blocked on accept() or read(). This can lead to thread and socket leaks during test execution. Using server.stop(true) ensures all sockets are closed immediately, releasing resources and preventing potential port-binding conflicts in subsequent tests.
| server.stop(); | |
| server.stop(true); |
|
|
||
| // Establish the connection by sending initial data | ||
| sender.send(ByteBuffer.wrap(data)); | ||
| assertTrue(firstDataReceivedLatch.await(5, TimeUnit.SECONDS)); |
There was a problem hiding this comment.
Using a 5-second timeout for network-based latch synchronization can lead to flaky test failures on slow or heavily loaded CI environments. Increasing this timeout to 10 seconds provides a safer margin without affecting the test duration under normal conditions.
| assertTrue(firstDataReceivedLatch.await(5, TimeUnit.SECONDS)); | |
| assertTrue(firstDataReceivedLatch.await(10, TimeUnit.SECONDS)); |
| } | ||
| } | ||
| } finally { | ||
| server.stop(); |
There was a problem hiding this comment.
Calling server.stop() (which defaults to server.stop(false)) does not immediately close the accepted sockets or the server socket if the server thread is blocked on accept() or read(). This can lead to thread and socket leaks during test execution. Using server.stop(true) ensures all sockets are closed immediately, releasing resources and preventing potential port-binding conflicts in subsequent tests.
| server.stop(); | |
| server.stop(true); |
| int id = idValue.asIntegerValue().asInt(); | ||
| if (!receivedIds.add(id)) { | ||
| backgroundError.compareAndSet( | ||
| null, new AssertionError("Duplicate record received: id=" + id)); | ||
| } |
Summary
TCPSenderTestcases that document RST and FIN reconnection behavior at the sender level, explaining why graceful FIN causes silent data loss (TCP two-write rule) while ACK mode avoids ittestAckModeGuaranteesDeliveryAfterServerDropsConnectionstoFluencyTestWithMockServerto demonstrate end-to-end that ACK mode (setAckResponseMode(true)) delivers all records after a simulated Fluentd restart that drops all TCP connectionsCloses #954