From 09fa93602018128608f58005c43dd15c9f84dc3b Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Wed, 23 Sep 2026 16:28:36 +0200 Subject: [PATCH 1/2] Add failing tests for T-1365: queued logs are lost when the JVM exits or stop() overlaps a flush Co-Authored-By: Claude Fable 5.1 --- .../logback/LogtailAppenderJvmExitTest.java | 127 ++++++++++++++++++ 1 file changed, 127 insertions(+) create mode 100644 src/test/java/com/logtail/logback/LogtailAppenderJvmExitTest.java diff --git a/src/test/java/com/logtail/logback/LogtailAppenderJvmExitTest.java b/src/test/java/com/logtail/logback/LogtailAppenderJvmExitTest.java new file mode 100644 index 0000000..ed68663 --- /dev/null +++ b/src/test/java/com/logtail/logback/LogtailAppenderJvmExitTest.java @@ -0,0 +1,127 @@ +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.spi.LoggingEvent; +import com.fasterxml.jackson.core.type.TypeReference; +import com.fasterxml.jackson.databind.ObjectMapper; +import com.sun.net.httpserver.HttpServer; +import org.junit.Test; + +import java.io.File; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.Scanner; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; + +/** + * 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. + */ +public class LogtailAppenderJvmExitTest { + + @Test + public void testQueuedLogsAreSentWhenTheJvmExitsWithoutStoppingLogback() throws Exception { + List receivedBodies = new CopyOnWriteArrayList<>(); + HttpServer server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.createContext("/", exchange -> { + receivedBodies.add(new Scanner(exchange.getRequestBody(), "UTF-8").useDelimiter("\\A").next()); + exchange.sendResponseHeaders(202, -1); + exchange.close(); + }); + server.start(); + try { + Process app = new ProcessBuilder( + System.getProperty("java.home") + File.separator + "bin" + File.separator + "java", + "-cp", System.getProperty("java.class.path"), + ExitingApp.class.getName(), + "http://127.0.0.1:" + server.getAddress().getPort()) + .inheritIO() + .start(); + if (!app.waitFor(30, TimeUnit.SECONDS)) { + app.destroyForcibly(); + fail("The app did not exit on its own"); + } + assertEquals(0, app.exitValue()); + } finally { + 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())); + } + + /** + * Run in a JVM of its own: logs once and returns from main without stopping logback. + */ + public static class ExitingApp { + public static void main(String[] args) { + LoggerContext context = new LoggerContext(); + LogtailAppender appender = new LogtailAppender(); + appender.setContext(context); + appender.setAppName("ExitingApp"); + appender.setSourceToken("source-token"); + appender.setIngestUrl(args[0]); + appender.start(); + + Logger logger = context.getLogger("ExitingApp"); + logger.addAppender(appender); + logger.info("Logged right before the JVM exits"); + } + } + + @Test + public void testStopWaitsForTheFlushInProgressAndSendsWhatQueuedBehindIt() 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.start(); + Logger logger = new LoggerContext().getLogger(Logger.ROOT_LOGGER_NAME); + + 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(1000); + assertTrue("stop() must wait for the flush in progress", stopping.isAlive()); + + requestMayComplete.countDown(); + stopping.join(5000); + assertFalse(stopping.isAlive()); + assertEquals(Arrays.asList(2, 1), sentBatchSizes); + } +} From 5d057aacde835332c447c7789e3d9d4681d76036 Mon Sep 17 00:00:00 2001 From: Petr Heinz Date: Wed, 23 Sep 2026 16:30:53 +0200 Subject: [PATCH 2/2] T-1365 Send queued logs when the JVM exits Co-Authored-By: Claude Fable 5.1 --- .../com/logtail/logback/LogtailAppender.java | 75 ++++++++++++------- .../logback/LogtailAppenderDecorator.java | 10 --- 2 files changed, 48 insertions(+), 37 deletions(-) diff --git a/src/main/java/com/logtail/logback/LogtailAppender.java b/src/main/java/com/logtail/logback/LogtailAppender.java index d857e83..44cf383 100644 --- a/src/main/java/com/logtail/logback/LogtailAppender.java +++ b/src/main/java/com/logtail/logback/LogtailAppender.java @@ -20,12 +20,12 @@ import java.nio.charset.StandardCharsets; import java.util.*; import java.util.Map.Entry; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.ScheduledFuture; import java.util.concurrent.ThreadFactory; import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.ReentrantLock; import java.util.stream.Collectors; public class LogtailAppender extends UnsynchronizedAppenderBase { @@ -51,13 +51,14 @@ public class LogtailAppender extends UnsynchronizedAppenderBase { // Non-customizable variables protected Vector batch = new Vector<>(); - protected AtomicBoolean isFlushing = new AtomicBoolean(false); + protected ReentrantLock flushLock = new ReentrantLock(); protected boolean mustReflush = false; protected boolean warnAboutMaxQueueSize = true; // Utils protected ScheduledExecutorService scheduledExecutorService; protected ScheduledFuture scheduledFuture; + protected Thread shutdownHook; protected ObjectMapper dataMapper; protected Logger logger; protected int retrySize = 0; @@ -114,7 +115,7 @@ protected void append(ILoggingEvent event) { } if (batch.size() >= batchSize) { - if (isFlushing.get()) + if (flushLock.isLocked()) return; startThread("logtail-appender-flush", new LogtailSender()); @@ -132,29 +133,30 @@ protected void flush() { return; // Guaranteed to not be running concurrently - if (isFlushing.getAndSet(true)) + if (!flushLock.tryLock()) return; - mustReflush = false; + try { + do { + mustReflush = false; - int flushedSize = batch.size(); - if (flushedSize > batchSize) { - flushedSize = batchSize; - mustReflush = true; - } - if (retries > 0 && flushedSize > retrySize) { - flushedSize = retrySize; - mustReflush = true; - } + int flushedSize = batch.size(); + if (flushedSize > batchSize) { + flushedSize = batchSize; + mustReflush = true; + } + if (retries > 0 && flushedSize > retrySize) { + flushedSize = retrySize; + mustReflush = true; + } - if (!flushLogs(flushedSize)) { - mustReflush = true; + if (!flushLogs(flushedSize)) { + mustReflush = true; + } + } while (!batch.isEmpty() && (mustReflush || batch.size() >= batchSize)); + } finally { + flushLock.unlock(); } - - isFlushing.set(false); - - if (mustReflush || batch.size() >= batchSize) - flush(); } protected boolean flushLogs(int flushedSize) { @@ -361,9 +363,6 @@ public void run() { flush(); } catch (Exception e) { logger.error("Error trying to flush : {}", e.getMessage(), e); - if (isFlushing.get()) { - isFlushing.set(false); - } } } } @@ -541,11 +540,33 @@ public boolean isDisabled() { return this.disabled; } + @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"); + Runtime.getRuntime().addShutdownHook(shutdownHook); + super.start(); + } + @Override public void stop() { + if (!isStarted()) + return; + + 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 + } scheduledExecutorService.shutdown(); - mustReflush = true; - flush(); - super.stop(); + + // Waits for a flush in progress on another thread, then sends everything still queued + flushLock.lock(); + try { + super.stop(); + flush(); + } finally { + flushLock.unlock(); + } } } diff --git a/src/test/java/com/logtail/logback/LogtailAppenderDecorator.java b/src/test/java/com/logtail/logback/LogtailAppenderDecorator.java index 9e1be68..77f0040 100644 --- a/src/test/java/com/logtail/logback/LogtailAppenderDecorator.java +++ b/src/test/java/com/logtail/logback/LogtailAppenderDecorator.java @@ -4,15 +4,12 @@ import java.util.List; import java.io.IOException; -import java.util.concurrent.locks.ReentrantLock; public class LogtailAppenderDecorator extends LogtailAppender { private Exception exception; private LogtailResponse response; protected int apiCalls = 0; - private ReentrantLock flushLock = new ReentrantLock(); - @Override protected LogtailResponse callHttpURLConnection(int flushedSize) throws IOException { try { @@ -26,13 +23,6 @@ protected LogtailResponse callHttpURLConnection(int flushedSize) throws IOExcept } } - @Override - public void flush() { - flushLock.lock(); - super.flush(); - flushLock.unlock(); - } - public void awaitFlushCompletion(){ try { // Wait a bit for possible asyncFlush to be initialized