Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
74 changes: 60 additions & 14 deletions src/main/java/com/logtail/logback/LogtailAppender.java
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ public class LogtailAppender extends UnsynchronizedAppenderBase<ILoggingEvent> {
protected int readTimeout = 10000;
protected int maxRetries = 5;
protected int retrySleepMilliseconds = 300;
protected int maxFlushTime = 30000;

protected PatternLayoutEncoder encoder;

Expand All @@ -60,6 +61,8 @@ public class LogtailAppender extends UnsynchronizedAppenderBase<ILoggingEvent> {
protected ScheduledExecutorService scheduledExecutorService;
protected ScheduledFuture<?> scheduledFuture;
protected Thread shutdownHook;
// Set on the thread whose flush stop() or the shutdown hook wait for: when they give up (System.nanoTime())
protected final ThreadLocal<Long> flushDeadline = new ThreadLocal<>();
protected ObjectMapper dataMapper;
protected Logger logger;
protected int retrySize = 0;
Expand Down Expand Up @@ -126,7 +129,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();
}
}

Expand All @@ -146,6 +152,17 @@ protected void flush() {

try {
do {
Long deadline = flushDeadline.get();
if (deadline != null && System.nanoTime() - deadline >= 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();
Expand Down Expand Up @@ -535,6 +552,18 @@ 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.
* 0 means no limit, as for logback's own AsyncAppender.
*
* @param maxFlushTime
* maximum time to send queued logs when stopping [ms], 0 for no limit
*/
public void setMaxFlushTime(int maxFlushTime) {
this.maxFlushTime = maxFlushTime;
}

/**
* Registers a dynamically loaded Module object to ObjectMapper used for serialization of logged data.
*
Expand Down Expand Up @@ -581,25 +610,42 @@ 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.
* 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() {
flushLock.lock();
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 {
flush();
} finally {
flushLock.unlock();
flushThread.join(maxFlushTime);
} catch (InterruptedException e) {
interrupted = true;
}
if (flushThread.isAlive())
logger.error("Gave up waiting for {} queued logs to be sent (maxFlushTime {} ms).", batch.size(), maxFlushTime);
if (interrupted)
Thread.currentThread().interrupt();
}
}
250 changes: 250 additions & 0 deletions src/test/java/com/logtail/logback/LogtailAppenderMaxFlushTimeTest.java
Original file line number Diff line number Diff line change
@@ -0,0 +1,250 @@
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.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;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;

/**
* 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
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());
}
}

@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<Integer> 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)) {
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());
}
}

@Test
public void testZeroMaxFlushTimeWaitsForAFlushInProgressAsLongAsItTakes() throws Exception {
CountDownLatch requestStarted = new CountDownLatch(1);
CountDownLatch requestMayComplete = new CountDownLatch(1);
List<Integer> 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.
*/
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 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);
context.putProperty("BATCH_SIZE", batchSize);
JoranConfigurator configurator = new JoranConfigurator();
configurator.setContext(context);
configurator.doConfigure(LogtailAppenderMaxFlushTimeTest.class.getResource("/logback-max-flush-time.xml"));
return context;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ public void testLogtailAppenderConfiguration() {

assertEquals(5000, appender.connectTimeout);
assertEquals(10000, appender.readTimeout);
assertEquals(30000, appender.maxFlushTime);

rootLogger.info("I am Groot");
}
Expand Down
16 changes: 16 additions & 0 deletions src/test/resources/logback-max-flush-time.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
<?xml version="1.0" encoding="UTF-8"?>
<configuration>

<appender name="Logtail" class="com.logtail.logback.LogtailAppender">
<sourceToken>source-token</sourceToken>
<ingestUrl>${SILENT_ENDPOINT}</ingestUrl>
<batchSize>${BATCH_SIZE}</batchSize>
<batchInterval>60000</batchInterval>
<maxFlushTime>1000</maxFlushTime>
</appender>

<root level="INFO">
<appender-ref ref="Logtail" />
</root>

</configuration>
Loading