From 694745c674c94c3be27db53c1505344d3c34b898 Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Tue, 29 Sep 2026 10:58:17 +0200 Subject: [PATCH 1/9] Add failing test for logs written while a framework shuts down Spring Boot and Quarkus keep logging while they shut down and stop logback at the very end. Since 0.3.8 the appender's own JVM shutdown hook stops the appender as soon as the JVM starts shutting down, so those lines are rejected. FrameworkApp reproduces that in a child JVM. Co-Authored-By: Claude Opus 5.5 --- .../logback/LogtailAppenderJvmExitTest.java | 60 +++++++++++++++++-- 1 file changed, 54 insertions(+), 6 deletions(-) diff --git a/src/test/java/com/logtail/logback/LogtailAppenderJvmExitTest.java b/src/test/java/com/logtail/logback/LogtailAppenderJvmExitTest.java index ed68663..6606983 100644 --- a/src/test/java/com/logtail/logback/LogtailAppenderJvmExitTest.java +++ b/src/test/java/com/logtail/logback/LogtailAppenderJvmExitTest.java @@ -12,6 +12,7 @@ import java.io.File; import java.io.IOException; import java.net.InetSocketAddress; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; import java.util.List; @@ -29,12 +30,27 @@ /** * Logs queued in the appender must reach Better Stack even when the application never stops logback and - * simply lets the JVM exit, and stop() must not return while a flush is still in progress on another thread. + * simply lets the JVM exit, logs written while a framework shuts down must still be sent when it stops logback, + * and stop() must not return while a flush is still in progress on another thread. */ public class LogtailAppenderJvmExitTest { @Test public void testQueuedLogsAreSentWhenTheJvmExitsWithoutStoppingLogback() throws Exception { + assertEquals(Collections.singletonList(Collections.singletonList("Logged right before the JVM exits")), + messagesSentByApp(ExitingApp.class)); + } + + @Test + public void testLogsWrittenWhileAFrameworkShutsDownAreSentWhenItStopsLogback() throws Exception { + assertEquals(Arrays.asList("Logged right before the JVM exits", "Logged while the framework shuts down"), + messagesSentByApp(FrameworkApp.class).stream().flatMap(List::stream).collect(Collectors.toList())); + } + + /** + * Runs the app in a JVM of its own against a local endpoint and returns the messages of each request it sent. + */ + private List> messagesSentByApp(Class appClass) throws Exception { List receivedBodies = new CopyOnWriteArrayList<>(); HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); server.createContext("/", exchange -> { @@ -47,7 +63,7 @@ public void testQueuedLogsAreSentWhenTheJvmExitsWithoutStoppingLogback() throws Process app = new ProcessBuilder( System.getProperty("java.home") + File.separator + "bin" + File.separator + "java", "-cp", System.getProperty("java.class.path"), - ExitingApp.class.getName(), + appClass.getName(), "http://127.0.0.1:" + server.getAddress().getPort()) .inheritIO() .start(); @@ -60,10 +76,12 @@ public void testQueuedLogsAreSentWhenTheJvmExitsWithoutStoppingLogback() throws server.stop(0); } - assertEquals(1, receivedBodies.size()); - List> lines = new ObjectMapper().readValue(receivedBodies.get(0), new TypeReference>>() {}); - assertEquals(Collections.singletonList("Logged right before the JVM exits"), - lines.stream().map(line -> line.get("message")).collect(Collectors.toList())); + List> messages = new ArrayList<>(); + for (String body : receivedBodies) { + List> lines = new ObjectMapper().readValue(body, new TypeReference>>() {}); + messages.add(lines.stream().map(line -> line.get("message")).collect(Collectors.toList())); + } + return messages; } /** @@ -85,6 +103,36 @@ public static void main(String[] args) { } } + /** + * Run in a JVM of its own: like Spring Boot or Quarkus, its shutdown hook keeps logging while it shuts down and + * stops logback at the very end. + */ + public static class FrameworkApp { + public static void main(String[] args) { + LoggerContext context = new LoggerContext(); + LogtailAppender appender = new LogtailAppender(); + appender.setContext(context); + appender.setAppName("FrameworkApp"); + appender.setSourceToken("source-token"); + appender.setIngestUrl(args[0]); + appender.start(); + + Logger logger = context.getLogger("FrameworkApp"); + logger.addAppender(appender); + Runtime.getRuntime().addShutdownHook(new Thread(() -> { + try { + // All shutdown hooks start together - by now the appender's own hook has done its part + Thread.sleep(300); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + logger.info("Logged while the framework shuts down"); + context.stop(); + })); + logger.info("Logged right before the JVM exits"); + } + } + @Test public void testStopWaitsForTheFlushInProgressAndSendsWhatQueuedBehindIt() throws Exception { CountDownLatch requestStarted = new CountDownLatch(1); From 968954a60cebb71d29ee8575658aa27434e708ca Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Tue, 29 Sep 2026 10:59:41 +0200 Subject: [PATCH 2/9] Keep the appender running while the JVM shuts down The shutdown hook now only sends what is queued, after waiting for a flush in progress, instead of stopping the appender. Lines logged afterwards by the application's own shutdown are queued as usual and sent when the framework (or logback's ) stops logback. Co-Authored-By: Claude Opus 5.5 --- .../com/logtail/logback/LogtailAppender.java | 20 ++++++++++++++++--- 1 file changed, 17 insertions(+), 3 deletions(-) diff --git a/src/main/java/com/logtail/logback/LogtailAppender.java b/src/main/java/com/logtail/logback/LogtailAppender.java index 0e59543..09a844f 100644 --- a/src/main/java/com/logtail/logback/LogtailAppender.java +++ b/src/main/java/com/logtail/logback/LogtailAppender.java @@ -551,8 +551,10 @@ public boolean isDisabled() { @Override public void start() { - // The sender runs on a daemon thread, so a JVM exiting on its own would take the queued logs with it - shutdownHook = new Thread(this::stop, "logtail-appender-shutdown"); + // The sender runs on a daemon thread, so a JVM exiting on its own would take the queued logs with it. The hook + // only sends the queue and leaves the appender running: all shutdown hooks run at once, and frameworks such as + // Spring Boot and Quarkus keep logging while they shut down and stop logback themselves at the very end + shutdownHook = new Thread(this::flushQueue, "logtail-appender-shutdown"); Runtime.getRuntime().addShutdownHook(shutdownHook); super.start(); } @@ -565,7 +567,7 @@ public void stop() { try { Runtime.getRuntime().removeShutdownHook(shutdownHook); } catch (IllegalStateException e) { - // The JVM is already shutting down - stop() is running from the hook itself or from logback's + // The JVM is already shutting down - stop() is running from logback's or a framework's shutdown hook } scheduledExecutorService.shutdown(); @@ -578,4 +580,16 @@ public void stop() { flushLock.unlock(); } } + + /** + * Waits for a flush in progress on another thread, then sends everything still queued. + */ + protected void flushQueue() { + flushLock.lock(); + try { + flush(); + } finally { + flushLock.unlock(); + } + } } From dbc28344cdd0e8411b3303b834828da6c1466270 Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Tue, 29 Sep 2026 11:01:42 +0200 Subject: [PATCH 3/9] Add failing tests for stopping against an endpoint that never answers stop() must give up after maxFlushTime, and a JVM whose flush hangs on the endpoint must still exit once main returns. Co-Authored-By: Claude Opus 5.5 --- .../LogtailAppenderMaxFlushTimeTest.java | 83 +++++++++++++++++++ src/test/resources/logback-max-flush-time.xml | 16 ++++ 2 files changed, 99 insertions(+) create mode 100644 src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java create mode 100644 src/test/resources/logback-max-flush-time.xml diff --git a/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java b/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java new file mode 100644 index 0000000..9333eae --- /dev/null +++ b/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java @@ -0,0 +1,83 @@ +package com.logtail.logback; + +import ch.qos.logback.classic.Logger; +import ch.qos.logback.classic.LoggerContext; +import ch.qos.logback.classic.joran.JoranConfigurator; +import ch.qos.logback.core.joran.spi.JoranException; +import org.junit.Test; + +import java.io.File; +import java.net.ServerSocket; +import java.util.concurrent.TimeUnit; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +/** + * stop() and the JVM shutdown hook send what is queued, but an endpoint that never answers must not hold the + * application's shutdown for longer than maxFlushTime (1 second in logback-max-flush-time.xml). + */ +public class LogtailAppenderMaxFlushTimeTest { + + @Test + public void testStopGivesUpOnAnEndpointThatNeverAnswers() throws Exception { + // Accepts connections but never answers: every attempt would wait for the full 10 second read timeout + try (ServerSocket silentEndpoint = new ServerSocket(0)) { + LoggerContext context = configure("http://127.0.0.1:" + silentEndpoint.getLocalPort(), "1000"); + Logger logger = context.getLogger("MaxFlushTimeTest"); + logger.info("Never sent 1"); + logger.info("Never sent 2"); + logger.info("Never sent 3"); + LogtailAppender appender = (LogtailAppender) context.getLogger(Logger.ROOT_LOGGER_NAME).getAppender("Logtail"); + + Thread stopping = new Thread(appender::stop); + stopping.setDaemon(true); + stopping.start(); + stopping.join(3000); + + assertFalse("stop() must give up after maxFlushTime", stopping.isAlive()); + assertTrue(appender.batch.isEmpty()); + } + } + + @Test + public void testJvmExitsWithinMaxFlushTimeWhenAFlushHangsOnTheEndpoint() throws Exception { + try (ServerSocket silentEndpoint = new ServerSocket(0)) { + Process app = new ProcessBuilder( + System.getProperty("java.home") + File.separator + "bin" + File.separator + "java", + "-cp", System.getProperty("java.class.path"), + SilentEndpointApp.class.getName(), + "http://127.0.0.1:" + silentEndpoint.getLocalPort()) + .inheritIO() + .start(); + if (!app.waitFor(10, TimeUnit.SECONDS)) { + app.destroyForcibly(); + fail("The app did not exit within maxFlushTime"); + } + assertEquals(0, app.exitValue()); + } + } + + /** + * Run in a JVM of its own: fills a batch, so a flush starts and hangs on the endpoint, and returns from main. + */ + public static class SilentEndpointApp { + public static void main(String[] args) throws JoranException { + Logger logger = configure(args[0], "2").getLogger("SilentEndpointApp"); + logger.info("First of a full batch"); + logger.info("Second of a full batch"); + } + } + + private static LoggerContext configure(String silentEndpoint, String batchSize) throws JoranException { + LoggerContext context = new LoggerContext(); + context.putProperty("SILENT_ENDPOINT", silentEndpoint); + context.putProperty("BATCH_SIZE", batchSize); + JoranConfigurator configurator = new JoranConfigurator(); + configurator.setContext(context); + configurator.doConfigure(LogtailAppenderMaxFlushTimeTest.class.getResource("/logback-max-flush-time.xml")); + return context; + } +} diff --git a/src/test/resources/logback-max-flush-time.xml b/src/test/resources/logback-max-flush-time.xml new file mode 100644 index 0000000..ac68825 --- /dev/null +++ b/src/test/resources/logback-max-flush-time.xml @@ -0,0 +1,16 @@ + + + + + source-token + ${SILENT_ENDPOINT} + ${BATCH_SIZE} + 60000 + 1000 + + + + + + + From 54489ad0a975e67d41da7d9b5b4210978ccc1d1e Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Tue, 29 Sep 2026 11:02:46 +0200 Subject: [PATCH 4/9] Give up sending after maxFlushTime when the application stops stop() and the shutdown hook wait at most maxFlushTime (default 10 s) for a flush in progress, clamp request timeouts and retry pauses to what is left, and drop what they could not send once it is up. Flush threads started by a full batch are daemon threads now, so they no longer keep a JVM whose main returned alive while they retry. Co-Authored-By: Claude Opus 5.5 --- .../com/logtail/logback/LogtailAppender.java | 76 +++++++++++++++---- 1 file changed, 60 insertions(+), 16 deletions(-) diff --git a/src/main/java/com/logtail/logback/LogtailAppender.java b/src/main/java/com/logtail/logback/LogtailAppender.java index 09a844f..572110b 100644 --- a/src/main/java/com/logtail/logback/LogtailAppender.java +++ b/src/main/java/com/logtail/logback/LogtailAppender.java @@ -46,6 +46,7 @@ public class LogtailAppender extends UnsynchronizedAppenderBase { protected int readTimeout = 10000; protected int maxRetries = 5; protected int retrySleepMilliseconds = 300; + protected int maxFlushTime = 10000; protected PatternLayoutEncoder encoder; @@ -59,6 +60,8 @@ public class LogtailAppender extends UnsynchronizedAppenderBase { protected ScheduledExecutorService scheduledExecutorService; protected ScheduledFuture scheduledFuture; protected Thread shutdownHook; + // Deadline (System.nanoTime()) of the flush that stop() or the shutdown hook runs on this thread + protected final ThreadLocal flushDeadline = new ThreadLocal<>(); protected ObjectMapper dataMapper; protected Logger logger; protected int retrySize = 0; @@ -122,7 +125,10 @@ protected void append(ILoggingEvent event) { if (flushLock.isLocked()) return; - startThread("logtail-appender-flush", new LogtailSender()); + // A daemon thread like the scheduled sender: at exit, the shutdown hook sends what it leaves behind + Thread flushThread = threadFactory.newThread(new LogtailSender()); + flushThread.setName("logtail-appender-flush"); + flushThread.start(); } } @@ -142,6 +148,16 @@ protected void flush() { try { do { + if (millisLeftToFlush() <= 0) { + int dropped; + synchronized (batch) { + dropped = batch.size(); + batch.clear(); + } + logger.error("Dropped {} logs that could not be sent within maxFlushTime ({} ms).", dropped, maxFlushTime); + retries = 0; + return; + } mustReflush = false; int flushedSize = batch.size(); @@ -181,7 +197,7 @@ protected boolean flushLogs(int flushedSize) { if (retries > 0) { logger.info("Retrying to send {} logs to Better Stack ({} / {})", flushedSize, retries, maxRetries); try { - TimeUnit.MILLISECONDS.sleep(retrySleepMilliseconds); + TimeUnit.MILLISECONDS.sleep(Math.min(retrySleepMilliseconds, millisLeftToFlush())); } catch (InterruptedException e) { // Continue } @@ -251,8 +267,10 @@ protected HttpURLConnection getHttpURLConnection() throws IOException { httpURLConnection.setRequestProperty("Charset", "UTF-8"); httpURLConnection.setRequestProperty("Authorization", String.format("Bearer %s", this.sourceToken)); httpURLConnection.setRequestMethod("POST"); - httpURLConnection.setConnectTimeout(this.connectTimeout); - httpURLConnection.setReadTimeout(this.readTimeout); + // Within stop() or the shutdown hook, no request may outlast maxFlushTime + long millisLeft = Math.max(1, millisLeftToFlush()); + httpURLConnection.setConnectTimeout((int) Math.min(this.connectTimeout, millisLeft)); + httpURLConnection.setReadTimeout((int) Math.min(this.readTimeout, millisLeft)); return httpURLConnection; } @@ -525,6 +543,17 @@ public void setRetrySleepMilliseconds(int retrySleepMilliseconds) { this.retrySleepMilliseconds = retrySleepMilliseconds; } + /** + * Sets the maximum time stop() and the JVM shutdown hook wait for queued logs to be sent, in milliseconds. Logs + * that could not be sent by then are dropped, so an endpoint that cannot be reached does not hold the shutdown. + * + * @param maxFlushTime + * maximum time to send queued logs when stopping [ms] + */ + public void setMaxFlushTime(int maxFlushTime) { + this.maxFlushTime = maxFlushTime; + } + /** * Registers a dynamically loaded Module object to ObjectMapper used for serialization of logged data. * @@ -571,25 +600,40 @@ public void stop() { } scheduledExecutorService.shutdown(); - // Waits for a flush in progress on another thread, then sends everything still queued - flushLock.lock(); - try { - super.stop(); - flush(); - } finally { - flushLock.unlock(); - } + // Like logback's own AsyncAppender: stop taking events, then send what is queued within maxFlushTime + super.stop(); + flushQueue(); } /** - * Waits for a flush in progress on another thread, then sends everything still queued. + * Waits for a flush in progress on another thread, then sends everything still queued - giving up once + * maxFlushTime has passed, so an endpoint that cannot be reached does not hold the application's shutdown. */ protected void flushQueue() { - flushLock.lock(); + flushDeadline.set(System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(maxFlushTime)); try { - flush(); + if (!flushLock.tryLock(maxFlushTime, TimeUnit.MILLISECONDS)) { + logger.error("Gave up waiting for a flush in progress after maxFlushTime ({} ms), {} logs not sent.", maxFlushTime, batch.size()); + return; + } + try { + flush(); + } finally { + flushLock.unlock(); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); } finally { - flushLock.unlock(); + flushDeadline.remove(); } } + + /** + * Milliseconds left until the flush that stop() or the shutdown hook runs on this thread gives up, Long.MAX_VALUE + * for any other flush. + */ + protected long millisLeftToFlush() { + Long deadline = flushDeadline.get(); + return deadline == null ? Long.MAX_VALUE : TimeUnit.NANOSECONDS.toMillis(deadline - System.nanoTime()); + } } From f1719b9e7db76c547142bbc53e74510121d230f1 Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Tue, 29 Sep 2026 12:35:04 +0200 Subject: [PATCH 5/9] Add failing tests for maxFlushTime 0 and a 30 second default As for logback's AsyncAppender, which hands maxFlushTime to Thread.join, 0 must mean no limit rather than dropping everything at once. The default becomes 30 seconds, like logtail-python's flush_timeout. Co-Authored-By: Claude Opus 5.5 --- .../LogtailAppenderMaxFlushTimeTest.java | 55 ++++++++++++++++++- .../logback/LogtailAppenderXmlConfigTest.java | 1 + 2 files changed, 55 insertions(+), 1 deletion(-) diff --git a/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java b/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java index 9333eae..843da38 100644 --- a/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java +++ b/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java @@ -1,13 +1,20 @@ package com.logtail.logback; +import ch.qos.logback.classic.Level; import ch.qos.logback.classic.Logger; import ch.qos.logback.classic.LoggerContext; import ch.qos.logback.classic.joran.JoranConfigurator; +import ch.qos.logback.classic.spi.LoggingEvent; import ch.qos.logback.core.joran.spi.JoranException; import org.junit.Test; import java.io.File; +import java.io.IOException; import java.net.ServerSocket; +import java.util.Arrays; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import static org.junit.Assert.assertEquals; @@ -17,7 +24,8 @@ /** * stop() and the JVM shutdown hook send what is queued, but an endpoint that never answers must not hold the - * application's shutdown for longer than maxFlushTime (1 second in logback-max-flush-time.xml). + * application's shutdown for longer than maxFlushTime (1 second in logback-max-flush-time.xml) - unless it is 0, + * which means no limit, as for logback's own AsyncAppender. */ public class LogtailAppenderMaxFlushTimeTest { @@ -60,6 +68,51 @@ public void testJvmExitsWithinMaxFlushTimeWhenAFlushHangsOnTheEndpoint() throws } } + @Test + public void testZeroMaxFlushTimeWaitsForAFlushInProgressAsLongAsItTakes() throws Exception { + CountDownLatch requestStarted = new CountDownLatch(1); + CountDownLatch requestMayComplete = new CountDownLatch(1); + List sentBatchSizes = new CopyOnWriteArrayList<>(); + LogtailAppender appender = new LogtailAppender() { + @Override + protected LogtailResponse callHttpURLConnection(int flushedSize) throws IOException { + requestStarted.countDown(); + try { + requestMayComplete.await(); + } catch (InterruptedException e) { + throw new IOException(e); + } + sentBatchSizes.add(flushedSize); + return new LogtailResponse(null, 202); + } + }; + appender.setContext(new LoggerContext()); + appender.setSourceToken("source-token"); + appender.setBatchSize(2); + appender.setMaxFlushTime(0); + appender.start(); + Logger logger = new LoggerContext().getLogger(Logger.ROOT_LOGGER_NAME); + + try { + appender.doAppend(new LoggingEvent(Logger.FQCN, logger, Level.INFO, "First", null, new Object[]{})); + appender.doAppend(new LoggingEvent(Logger.FQCN, logger, Level.INFO, "Second", null, new Object[]{})); + assertTrue("A full batch starts a flush", requestStarted.await(5, TimeUnit.SECONDS)); + appender.doAppend(new LoggingEvent(Logger.FQCN, logger, Level.INFO, "Third", null, new Object[]{})); + + Thread stopping = new Thread(appender::stop); + stopping.start(); + stopping.join(1500); + assertTrue("stop() must keep waiting for the flush in progress", stopping.isAlive()); + + requestMayComplete.countDown(); + stopping.join(5000); + assertFalse(stopping.isAlive()); + assertEquals(Arrays.asList(2, 1), sentBatchSizes); + } finally { + requestMayComplete.countDown(); + } + } + /** * Run in a JVM of its own: fills a batch, so a flush starts and hangs on the endpoint, and returns from main. */ diff --git a/src/test/java/com/logtail/logback/LogtailAppenderXmlConfigTest.java b/src/test/java/com/logtail/logback/LogtailAppenderXmlConfigTest.java index d1a9180..6d7c909 100644 --- a/src/test/java/com/logtail/logback/LogtailAppenderXmlConfigTest.java +++ b/src/test/java/com/logtail/logback/LogtailAppenderXmlConfigTest.java @@ -52,6 +52,7 @@ public void testLogtailAppenderConfiguration() { assertEquals(5000, appender.connectTimeout); assertEquals(10000, appender.readTimeout); + assertEquals(30000, appender.maxFlushTime); rootLogger.info("I am Groot"); } From 13fabacf6d8883cc7da4ef25840a983b233b8c70 Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Tue, 29 Sep 2026 12:36:32 +0200 Subject: [PATCH 6/9] Default maxFlushTime to 30 seconds and treat 0 as no limit 30 seconds matches logtail-python's flush_timeout. 0 now waits for as long as the flush takes, like logback's AsyncAppender, instead of dropping the queue at once. Co-Authored-By: Claude Opus 5.5 --- .../java/com/logtail/logback/LogtailAppender.java | 11 +++++++---- 1 file changed, 7 insertions(+), 4 deletions(-) diff --git a/src/main/java/com/logtail/logback/LogtailAppender.java b/src/main/java/com/logtail/logback/LogtailAppender.java index 572110b..18652be 100644 --- a/src/main/java/com/logtail/logback/LogtailAppender.java +++ b/src/main/java/com/logtail/logback/LogtailAppender.java @@ -46,7 +46,7 @@ public class LogtailAppender extends UnsynchronizedAppenderBase { protected int readTimeout = 10000; protected int maxRetries = 5; protected int retrySleepMilliseconds = 300; - protected int maxFlushTime = 10000; + protected int maxFlushTime = 30000; protected PatternLayoutEncoder encoder; @@ -546,9 +546,10 @@ public void setRetrySleepMilliseconds(int retrySleepMilliseconds) { /** * Sets the maximum time stop() and the JVM shutdown hook wait for queued logs to be sent, in milliseconds. Logs * that could not be sent by then are dropped, so an endpoint that cannot be reached does not hold the shutdown. + * 0 means no limit, as for logback's own AsyncAppender. * * @param maxFlushTime - * maximum time to send queued logs when stopping [ms] + * maximum time to send queued logs when stopping [ms], 0 for no limit */ public void setMaxFlushTime(int maxFlushTime) { this.maxFlushTime = maxFlushTime; @@ -610,9 +611,11 @@ public void stop() { * maxFlushTime has passed, so an endpoint that cannot be reached does not hold the application's shutdown. */ protected void flushQueue() { - flushDeadline.set(System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(maxFlushTime)); + // 0 means no limit, as for logback's own AsyncAppender, which hands maxFlushTime to Thread.join() + if (maxFlushTime > 0) + flushDeadline.set(System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(maxFlushTime)); try { - if (!flushLock.tryLock(maxFlushTime, TimeUnit.MILLISECONDS)) { + if (!flushLock.tryLock(millisLeftToFlush(), TimeUnit.MILLISECONDS)) { logger.error("Gave up waiting for a flush in progress after maxFlushTime ({} ms), {} logs not sent.", maxFlushTime, batch.size()); return; } From ff912a0c74d19a7f0033c0831d38f604aaaa1be6 Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Tue, 29 Sep 2026 13:29:43 +0200 Subject: [PATCH 7/9] Add failing tests for stops that maxFlushTime does not bound yet stop() must give up after maxFlushTime on a request without a read timeout and on an endpoint that stopped taking data, and it must still send the queue when it is called on an interrupted thread. Also pins that logs not sent in time are dropped once the request in progress is over, and no longer expects the queue to be empty the moment stop() gives up. Co-Authored-By: Claude Fable 5.1 --- .../LogtailAppenderMaxFlushTimeTest.java | 122 +++++++++++++++++- 1 file changed, 118 insertions(+), 4 deletions(-) diff --git a/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java b/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java index 843da38..cf15118 100644 --- a/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java +++ b/src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java @@ -10,12 +10,17 @@ import java.io.File; import java.io.IOException; +import java.net.InetSocketAddress; import java.net.ServerSocket; +import java.net.SocketTimeoutException; import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; @@ -23,12 +28,15 @@ import static org.junit.Assert.fail; /** - * stop() and the JVM shutdown hook send what is queued, but an endpoint that never answers must not hold the - * application's shutdown for longer than maxFlushTime (1 second in logback-max-flush-time.xml) - unless it is 0, - * which means no limit, as for logback's own AsyncAppender. + * stop() and the JVM shutdown hook send what is queued, but nothing may hold the application's shutdown for longer + * than maxFlushTime (1 second in logback-max-flush-time.xml): not an endpoint that never answers, not one that + * stopped taking data, not a request without a timeout of its own - unless maxFlushTime is 0, which means no limit, + * as for logback's own AsyncAppender. */ public class LogtailAppenderMaxFlushTimeTest { + private static final Logger LOGGER = new LoggerContext().getLogger(Logger.ROOT_LOGGER_NAME); + @Test public void testStopGivesUpOnAnEndpointThatNeverAnswers() throws Exception { // Accepts connections but never answers: every attempt would wait for the full 10 second read timeout @@ -46,10 +54,94 @@ public void testStopGivesUpOnAnEndpointThatNeverAnswers() throws Exception { stopping.join(3000); assertFalse("stop() must give up after maxFlushTime", stopping.isAlive()); - assertTrue(appender.batch.isEmpty()); } } + @Test + public void testStopGivesUpOnAnEndpointThatNeverAnswersARequestWithoutReadTimeout() throws Exception { + // readTimeout 0 turns the request's own timeout off, so only maxFlushTime can end the wait for an answer + try (ServerSocket silentEndpoint = new ServerSocket(0)) { + LogtailAppender appender = new LogtailAppender(); + appender.setReadTimeout(0); + start(appender, "http://127.0.0.1:" + silentEndpoint.getLocalPort()); + queue(appender, "Never sent"); + + assertStopGivesUp(appender); + } + } + + @Test + public void testStopGivesUpOnAnEndpointThatStoppedTakingData() throws Exception { + // The batch is bigger than what the connection's buffers take and nobody reads it on the other end. Writing + // to a socket has no timeout at all, so only maxFlushTime can end the wait + try (ServerSocket silentEndpoint = new ServerSocket()) { + silentEndpoint.setReceiveBufferSize(1024); + silentEndpoint.bind(new InetSocketAddress("127.0.0.1", 0)); + LogtailAppender appender = new LogtailAppender(); + start(appender, "http://127.0.0.1:" + silentEndpoint.getLocalPort()); + char[] line = new char[4 * 1024]; + Arrays.fill(line, 'x'); + String message = new String(line); + // One short of the default batchSize, so that no flush starts before stop() + for (int i = 0; i < 999; i++) + queue(appender, message); + + assertStopGivesUp(appender); + } + } + + @Test + public void testStopSendsTheQueueOnAnInterruptedThread() throws Exception { + List sentBatchSizes = new CopyOnWriteArrayList<>(); + LogtailAppender appender = new LogtailAppender() { + @Override + protected LogtailResponse callHttpURLConnection(int flushedSize) { + sentBatchSizes.add(flushedSize); + return new LogtailResponse(null, 202); + } + }; + start(appender, "http://127.0.0.1"); + queue(appender, "Sent by an interrupted thread"); + + AtomicBoolean interruptedAfterStop = new AtomicBoolean(); + Thread stopping = new Thread(() -> { + Thread.currentThread().interrupt(); + appender.stop(); + interruptedAfterStop.set(Thread.currentThread().isInterrupted()); + }); + stopping.start(); + stopping.join(5000); + + assertEquals(Collections.singletonList(1), sentBatchSizes); + assertTrue("The interrupt is left for the caller to handle", interruptedAfterStop.get()); + } + + @Test + public void testLogsNotSentInTimeAreDroppedAfterTheRequestInProgress() throws Exception { + AtomicInteger requests = new AtomicInteger(); + LogtailAppender appender = new LogtailAppender() { + @Override + protected LogtailResponse callHttpURLConnection(int flushedSize) throws IOException { + requests.incrementAndGet(); + try { + Thread.sleep(1000); + } catch (InterruptedException e) { + throw new IOException(e); + } + throw new SocketTimeoutException("Read timed out"); + } + }; + start(appender, "http://127.0.0.1"); + queue(appender, "Never sent"); + + appender.stop(); + for (int waited = 0; !appender.batch.isEmpty() && waited < 5000; waited += 50) + Thread.sleep(50); + + assertTrue(appender.batch.isEmpty()); + assertEquals("No retry once the time is up", 1, requests.get()); + } + @Test public void testJvmExitsWithinMaxFlushTimeWhenAFlushHangsOnTheEndpoint() throws Exception { try (ServerSocket silentEndpoint = new ServerSocket(0)) { @@ -124,6 +216,28 @@ public static void main(String[] args) throws JoranException { } } + private static void start(LogtailAppender appender, String ingestUrl) { + appender.setContext(new LoggerContext()); + appender.setSourceToken("source-token"); + appender.setIngestUrl(ingestUrl); + appender.setBatchInterval(60000); + appender.setMaxFlushTime(500); + appender.start(); + } + + private static void queue(LogtailAppender appender, String message) { + appender.doAppend(new LoggingEvent(Logger.FQCN, LOGGER, Level.INFO, message, null, new Object[]{})); + } + + private static void assertStopGivesUp(LogtailAppender appender) throws InterruptedException { + Thread stopping = new Thread(appender::stop); + stopping.setDaemon(true); + stopping.start(); + stopping.join(3000); + + assertFalse("stop() must give up after maxFlushTime", stopping.isAlive()); + } + private static LoggerContext configure(String silentEndpoint, String batchSize) throws JoranException { LoggerContext context = new LoggerContext(); context.putProperty("SILENT_ENDPOINT", silentEndpoint); From ba35122e832fadcb9a37729236617a1ef2b3dfed Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Tue, 29 Sep 2026 13:33:50 +0200 Subject: [PATCH 8/9] Leave the last flush to a daemon thread and wait for it for maxFlushTime stop() and the shutdown hook no longer send the queue themselves with timeouts cut down to the time left. They start a daemon thread for it and join it for maxFlushTime, as logback's AsyncAppender does with its worker, so the wait ends on time whatever holds the request: a read or connect timeout set to 0, a socket that stopped taking data, a name server that does not answer. An interrupt pending on the calling thread no longer skips the flush. The thread left behind drops what is queued once its request in progress is over, instead of retrying on. Co-Authored-By: Claude Fable 5.1 --- .../com/logtail/logback/LogtailAppender.java | 57 +++++++++---------- 1 file changed, 28 insertions(+), 29 deletions(-) diff --git a/src/main/java/com/logtail/logback/LogtailAppender.java b/src/main/java/com/logtail/logback/LogtailAppender.java index 18652be..2044c1c 100644 --- a/src/main/java/com/logtail/logback/LogtailAppender.java +++ b/src/main/java/com/logtail/logback/LogtailAppender.java @@ -60,7 +60,7 @@ public class LogtailAppender extends UnsynchronizedAppenderBase { protected ScheduledExecutorService scheduledExecutorService; protected ScheduledFuture scheduledFuture; protected Thread shutdownHook; - // Deadline (System.nanoTime()) of the flush that stop() or the shutdown hook runs on this thread + // Set on the thread whose flush stop() or the shutdown hook wait for: when they give up (System.nanoTime()) protected final ThreadLocal flushDeadline = new ThreadLocal<>(); protected ObjectMapper dataMapper; protected Logger logger; @@ -148,7 +148,8 @@ protected void flush() { try { do { - if (millisLeftToFlush() <= 0) { + Long deadline = flushDeadline.get(); + if (deadline != null && System.nanoTime() - deadline >= 0) { int dropped; synchronized (batch) { dropped = batch.size(); @@ -197,7 +198,7 @@ protected boolean flushLogs(int flushedSize) { if (retries > 0) { logger.info("Retrying to send {} logs to Better Stack ({} / {})", flushedSize, retries, maxRetries); try { - TimeUnit.MILLISECONDS.sleep(Math.min(retrySleepMilliseconds, millisLeftToFlush())); + TimeUnit.MILLISECONDS.sleep(retrySleepMilliseconds); } catch (InterruptedException e) { // Continue } @@ -267,10 +268,8 @@ protected HttpURLConnection getHttpURLConnection() throws IOException { httpURLConnection.setRequestProperty("Charset", "UTF-8"); httpURLConnection.setRequestProperty("Authorization", String.format("Bearer %s", this.sourceToken)); httpURLConnection.setRequestMethod("POST"); - // Within stop() or the shutdown hook, no request may outlast maxFlushTime - long millisLeft = Math.max(1, millisLeftToFlush()); - httpURLConnection.setConnectTimeout((int) Math.min(this.connectTimeout, millisLeft)); - httpURLConnection.setReadTimeout((int) Math.min(this.readTimeout, millisLeft)); + httpURLConnection.setConnectTimeout(this.connectTimeout); + httpURLConnection.setReadTimeout(this.readTimeout); return httpURLConnection; } @@ -607,36 +606,36 @@ public void stop() { } /** - * Waits for a flush in progress on another thread, then sends everything still queued - giving up once - * maxFlushTime has passed, so an endpoint that cannot be reached does not hold the application's shutdown. + * Sends everything still queued, after a flush in progress on another thread, and waits for that for at most + * maxFlushTime. The sending is left to a daemon thread, as in logback's own AsyncAppender, so whatever holds it + * cannot hold the application's shutdown for longer: an endpoint that cannot be reached, a name server that does + * not answer, a connection that stopped taking data. */ protected void flushQueue() { - // 0 means no limit, as for logback's own AsyncAppender, which hands maxFlushTime to Thread.join() - if (maxFlushTime > 0) - flushDeadline.set(System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(maxFlushTime)); - try { - if (!flushLock.tryLock(millisLeftToFlush(), TimeUnit.MILLISECONDS)) { - logger.error("Gave up waiting for a flush in progress after maxFlushTime ({} ms), {} logs not sent.", maxFlushTime, batch.size()); - return; - } + Thread flushThread = threadFactory.newThread(() -> { + // 0 means no limit, here as well as for Thread.join() below + if (maxFlushTime > 0) + flushDeadline.set(System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(maxFlushTime)); + flushLock.lock(); try { flush(); } finally { flushLock.unlock(); } + }); + flushThread.setName("logtail-appender-flush"); + flushThread.start(); + + // An interrupt from before must not skip the wait, it is put back for the caller in the end + boolean interrupted = Thread.interrupted(); + try { + flushThread.join(maxFlushTime); } catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } finally { - flushDeadline.remove(); + interrupted = true; } - } - - /** - * Milliseconds left until the flush that stop() or the shutdown hook runs on this thread gives up, Long.MAX_VALUE - * for any other flush. - */ - protected long millisLeftToFlush() { - Long deadline = flushDeadline.get(); - return deadline == null ? Long.MAX_VALUE : TimeUnit.NANOSECONDS.toMillis(deadline - System.nanoTime()); + if (flushThread.isAlive()) + logger.error("Gave up waiting for {} queued logs to be sent (maxFlushTime {} ms).", batch.size(), maxFlushTime); + if (interrupted) + Thread.currentThread().interrupt(); } } From c13e18612f01def61a45bb5ad1601a558148c63f Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Tue, 29 Sep 2026 14:19:49 +0200 Subject: [PATCH 9/9] Drop the flushQueue() the merge of main left behind Merging main brought in flushQueue() as #34 added it, next to the version this branch replaces it with, and the class did not compile. Co-Authored-By: Claude Fable 5.1 --- .../java/com/logtail/logback/LogtailAppender.java | 12 ------------ 1 file changed, 12 deletions(-) diff --git a/src/main/java/com/logtail/logback/LogtailAppender.java b/src/main/java/com/logtail/logback/LogtailAppender.java index fa8600d..d259301 100644 --- a/src/main/java/com/logtail/logback/LogtailAppender.java +++ b/src/main/java/com/logtail/logback/LogtailAppender.java @@ -648,16 +648,4 @@ protected void flushQueue() { if (interrupted) Thread.currentThread().interrupt(); } - - /** - * Waits for a flush in progress on another thread, then sends everything still queued. - */ - protected void flushQueue() { - flushLock.lock(); - try { - flush(); - } finally { - flushLock.unlock(); - } - } }