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
66 changes: 54 additions & 12 deletions src/main/java/com/logtail/logback/LogtailAppender.java
Original file line number Diff line number Diff line change
Expand Up @@ -470,22 +470,34 @@ public void setMdcTypes(String mdcTypes) {
}

/**
* Sets the maximum number of messages in the queue. Messages over the limit will be dropped.
* Sets the maximum number of messages in the queue. Messages over the limit will be dropped. A value below 1 is
* ignored with a warning and the current size is kept.
*
* @param maxQueueSize
* max size of the message queue
* max size of the message queue, 1 or more
*/
public void setMaxQueueSize(int maxQueueSize) {
// No log would ever be queued, so none would be sent
if (maxQueueSize <= 0) {
addWarn("maxQueueSize must be positive, keeping " + this.maxQueueSize + " instead of " + maxQueueSize);
return;
}
this.maxQueueSize = maxQueueSize;
}

/**
* Sets the batch size for the number of messages to be sent via the API
* Sets the batch size for the number of messages to be sent via the API. A value below 1 is ignored with a warning
* and the current size is kept.
*
* @param batchSize
* size of the message batch
* size of the message batch, 1 or more
*/
public void setBatchSize(int batchSize) {
// flush() would loop without end, sending empty batches or failing on every one, and never send the queued logs
if (batchSize <= 0) {
addWarn("batchSize must be positive, keeping " + this.batchSize + " instead of " + batchSize);
return;
}
this.batchSize = batchSize;
}

Expand All @@ -497,12 +509,18 @@ public int getBatchSize() {
}

/**
* Sets the maximum wait time for a batch to be sent via the API, in milliseconds.
* Sets the maximum wait time for a batch to be sent via the API, in milliseconds. A value below 1 is ignored with a
* warning and the current interval is kept.
*
* @param batchInterval
* maximum wait time for message batch [ms]
* maximum wait time for message batch [ms], 1 or more
*/
public void setBatchInterval(int batchInterval) {
// scheduleWithFixedDelay() throws on it, which would fail start() or leave a running sender cancelled
if (batchInterval <= 0) {
addWarn("batchInterval must be positive, keeping " + this.batchInterval + " ms instead of " + batchInterval);
return;
}
this.batchInterval = batchInterval;

// Before start(), which schedules the sender with this interval, there is no sender to reschedule
Expand All @@ -513,32 +531,50 @@ public void setBatchInterval(int batchInterval) {
}

/**
* Sets the connection timeout of the underlying HTTP client, in milliseconds.
* Sets the connection timeout of the underlying HTTP client, in milliseconds. 0 means no timeout. A negative value
* is ignored with a warning and the current timeout is kept.
*
* @param connectTimeout
* client connection timeout [ms]
* client connection timeout [ms], 0 for no timeout
*/
public void setConnectTimeout(int connectTimeout) {
// HttpURLConnection throws on it: every request would fail and every batch be dropped after its retries
if (connectTimeout < 0) {
addWarn("connectTimeout must be 0 (no timeout) or more, keeping " + this.connectTimeout + " ms instead of " + connectTimeout);
return;
}
this.connectTimeout = connectTimeout;
}

/**
* Sets the read timeout of the underlying HTTP client, in milliseconds.
* Sets the read timeout of the underlying HTTP client, in milliseconds. 0 means no timeout. A negative value is
* ignored with a warning and the current timeout is kept.
*
* @param readTimeout
* client read timeout
* client read timeout [ms], 0 for no timeout
*/
public void setReadTimeout(int readTimeout) {
// HttpURLConnection throws on it: every request would fail and every batch be dropped after its retries
if (readTimeout < 0) {
addWarn("readTimeout must be 0 (no timeout) or more, keeping " + this.readTimeout + " ms instead of " + readTimeout);
return;
}
this.readTimeout = readTimeout;
}

/**
* Sets the maximum number of retries for sending logs to Better Stack. After that, current batch of logs will be dropped.
* 0 means no retries. A negative value is ignored with a warning and the current number is kept.
*
* @param maxRetries
* max number of retries for sending logs
* max number of retries for sending logs, 0 or more
*/
public void setMaxRetries(int maxRetries) {
// Every batch would be dropped before its first attempt
if (maxRetries < 0) {
addWarn("maxRetries must be 0 (no retries) or more, keeping " + this.maxRetries + " instead of " + maxRetries);
return;
}
this.maxRetries = maxRetries;
}

Expand All @@ -555,12 +591,18 @@ 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.
* 0 means no limit, as for logback's own AsyncAppender. A negative value is ignored with a warning and the current
* time is kept.
*
* @param maxFlushTime
* maximum time to send queued logs when stopping [ms], 0 for no limit
*/
public void setMaxFlushTime(int maxFlushTime) {
// Thread.join() throws on it: stop() would fail and the JVM would exit without waiting for the queue to be sent
if (maxFlushTime < 0) {
addWarn("maxFlushTime must be 0 (no limit) or more, keeping " + this.maxFlushTime + " ms instead of " + maxFlushTime);
return;
}
this.maxFlushTime = maxFlushTime;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,17 @@
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.Context;
import ch.qos.logback.core.status.Status;
import com.sun.net.httpserver.HttpServer;
import org.junit.Test;

import java.net.InetSocketAddress;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Set;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
Expand All @@ -18,7 +25,8 @@

/**
* start() starts the appender, sender thread included, and does so once; an appender that never starts has no thread
* to leave behind, and the batch interval only reschedules a sender that is running.
* to leave behind, and the batch interval only reschedules a sender that is running. A batch interval below 1 ms, which
* the sender cannot be scheduled with, is ignored with a warning.
*/
public class LogtailAppenderLifecycleTest {

Expand Down Expand Up @@ -97,6 +105,79 @@ public void testARestartedAppenderSendsOnTheScheduleAgain() throws Exception {
}
}

@Test
public void testABatchIntervalBelow1MsInTheConfigKeepsTheDefaultWithAWarning() throws Exception {
CountDownLatch sent = new CountDownLatch(2);
HttpServer endpoint = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0);
endpoint.createContext("/", exchange -> {
sent.countDown();
exchange.sendResponseHeaders(202, -1);
exchange.close();
});
endpoint.start();
LoggerContext context = new LoggerContext();
context.putProperty("ENDPOINT", "http://127.0.0.1:" + endpoint.getAddress().getPort());
try {
JoranConfigurator configurator = new JoranConfigurator();
configurator.setContext(context);
configurator.doConfigure(LogtailAppenderLifecycleTest.class.getResource("/logback-batch-interval.xml"));

assertEquals("Spring Boot refuses to start on a logback error", Collections.emptyList(),
context.getStatusManager().getCopyOfStatusList().stream()
.filter(status -> status.getLevel() == Status.ERROR)
.map(Status::toString)
.collect(Collectors.toList()));
assertEquals(Arrays.asList(
"batchInterval must be positive, keeping 3000 ms instead of 0",
"batchInterval must be positive, keeping 3000 ms instead of -1"), warnings(context));
Logger root = context.getLogger(Logger.ROOT_LOGGER_NAME);
LogtailAppender zero = (LogtailAppender) root.getAppender("Zero");
LogtailAppender negative = (LogtailAppender) root.getAppender("Negative");
assertTrue(zero.isStarted());
assertTrue(negative.isStarted());
assertEquals(3000, zero.batchInterval);
assertEquals(3000, negative.batchInterval);

root.info("Sent by both appenders on the default 3 second schedule");
assertTrue(sent.await(10, TimeUnit.SECONDS));
} finally {
context.stop();
endpoint.stop(0);
}
}

@Test
public void testABatchIntervalBelow1MsKeepsTheIntervalSetBefore() throws Exception {
SendingAppender appender = new SendingAppender();
appender.setBatchInterval(100);
appender.setBatchInterval(0);
appender.start();
try {
appender.doAppend(event("Sent on the 100 ms schedule"));
assertTrue("The default 3 second schedule would not have sent it yet", appender.sent.await(2, TimeUnit.SECONDS));
assertEquals(Collections.singletonList("batchInterval must be positive, keeping 100 ms instead of 0"),
warnings(appender.getContext()));
} finally {
appender.stop();
}
}

@Test
public void testABatchIntervalBelow1MsLeavesARunningSenderOnItsSchedule() throws Exception {
SendingAppender appender = new SendingAppender();
appender.setBatchInterval(100);
appender.start();
try {
appender.setBatchInterval(-5);
appender.doAppend(event("Sent on the 100 ms schedule"));
assertTrue(appender.sent.await(2, TimeUnit.SECONDS));
assertEquals(Collections.singletonList("batchInterval must be positive, keeping 100 ms instead of -5"),
warnings(appender.getContext()));
} finally {
appender.stop();
}
}

/**
* Counts the requests instead of sending them.
*/
Expand Down Expand Up @@ -131,4 +212,11 @@ private static Set<Thread> newSince(Set<Thread> before) {
now.removeAll(before);
return now;
}

private static List<String> warnings(Context context) {
return context.getStatusManager().getCopyOfStatusList().stream()
.filter(status -> status.getLevel() == Status.WARN)
.map(Status::getMessage)
.collect(Collectors.toList());
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@
import ch.qos.logback.classic.joran.JoranConfigurator;
import ch.qos.logback.classic.spi.LoggingEvent;
import ch.qos.logback.core.joran.spi.JoranException;
import ch.qos.logback.core.status.Status;
import org.junit.Test;

import java.io.File;
Expand All @@ -21,6 +22,7 @@
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.stream.Collectors;

import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
Expand All @@ -31,7 +33,8 @@
* 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.
* as for logback's own AsyncAppender. A negative maxFlushTime, which Thread.join() does not take, is ignored with a
* warning.
*/
public class LogtailAppenderMaxFlushTimeTest {

Expand Down Expand Up @@ -205,6 +208,32 @@ protected LogtailResponse callHttpURLConnection(int flushedSize) throws IOExcept
}
}

@Test
public void testANegativeMaxFlushTimeKeepsTheDefaultWithAWarning() {
List<Integer> sentBatchSizes = new CopyOnWriteArrayList<>();
LogtailAppender appender = new LogtailAppender() {
@Override
protected LogtailResponse callHttpURLConnection(int flushedSize) {
sentBatchSizes.add(flushedSize);
return new LogtailResponse(null, 202);
}
};
appender.setContext(new LoggerContext());
appender.setSourceToken("source-token");
appender.setMaxFlushTime(-1);
appender.start();
queue(appender, "Sent by stop()");

appender.stop();

assertEquals(Collections.singletonList(1), sentBatchSizes);
assertEquals(Collections.singletonList("maxFlushTime must be 0 (no limit) or more, keeping 30000 ms instead of -1"),
appender.getContext().getStatusManager().getCopyOfStatusList().stream()
.filter(status -> status.getLevel() == Status.WARN)
.map(Status::getMessage)
.collect(Collectors.toList()));
}

/**
* Run in a JVM of its own: fills a batch, so a flush starts and hangs on the endpoint, and returns from main.
*/
Expand Down
Loading
Loading