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
75 changes: 48 additions & 27 deletions src/main/java/com/logtail/logback/LogtailAppender.java
Original file line number Diff line number Diff line change
Expand Up @@ -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<ILoggingEvent> {
Expand All @@ -51,13 +51,14 @@ public class LogtailAppender extends UnsynchronizedAppenderBase<ILoggingEvent> {

// Non-customizable variables
protected Vector<ILoggingEvent> 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;
Expand Down Expand Up @@ -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());
Expand All @@ -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) {
Expand Down Expand Up @@ -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);
}
}
}
}
Expand Down Expand Up @@ -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();
}
}
}
10 changes: 0 additions & 10 deletions src/test/java/com/logtail/logback/LogtailAppenderDecorator.java
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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
Expand Down
127 changes: 127 additions & 0 deletions src/test/java/com/logtail/logback/LogtailAppenderJvmExitTest.java
Original file line number Diff line number Diff line change
@@ -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<String> 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<Map<String, Object>> lines = new ObjectMapper().readValue(receivedBodies.get(0), new TypeReference<List<Map<String, Object>>>() {});
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<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.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);
}
}
Loading