001package com.logfire.logback;
002
003import ch.qos.logback.classic.encoder.PatternLayoutEncoder;
004import ch.qos.logback.classic.spi.ILoggingEvent;
005import ch.qos.logback.classic.spi.IThrowableProxy;
006import ch.qos.logback.core.UnsynchronizedAppenderBase;
007import com.fasterxml.jackson.annotation.JsonInclude;
008import com.fasterxml.jackson.core.JsonProcessingException;
009import com.fasterxml.jackson.databind.Module;
010import com.fasterxml.jackson.databind.ObjectMapper;
011import com.fasterxml.jackson.databind.PropertyNamingStrategies;
012import org.slf4j.Logger;
013import org.slf4j.LoggerFactory;
014
015import java.io.IOException;
016import java.io.OutputStream;
017import java.net.HttpURLConnection;
018import java.net.URL;
019import java.nio.charset.StandardCharsets;
020import java.util.*;
021import java.util.Map.Entry;
022import java.util.concurrent.atomic.AtomicBoolean;
023import java.util.concurrent.Executors;
024import java.util.concurrent.ScheduledExecutorService;
025import java.util.concurrent.ScheduledFuture;
026import java.util.concurrent.ThreadFactory;
027import java.util.concurrent.TimeUnit;
028import java.util.stream.Collectors;
029
030public class LogfireAppender extends UnsynchronizedAppenderBase<ILoggingEvent> {
031
032    // Customizable variables
033    protected String appName;
034//    protected String ingestUrl = "https://in.logfire.ai";
035
036    protected String ingestUrl = "http://localhost:8888/logfire.sh";
037    protected String sourceToken;
038    protected String userAgent = "Logfire Logback Appender";
039
040    protected List<String> mdcFields = new ArrayList<>();
041    protected List<String> mdcTypes = new ArrayList<>();
042
043    protected int maxQueueSize = 100000;
044    protected int batchSize = 1000;
045    protected int batchInterval = 3000;
046    protected int connectTimeout = 5000;
047    protected int readTimeout = 10000;
048    protected int maxRetries = 5;
049    protected int retrySleepMilliseconds = 300;
050
051    protected PatternLayoutEncoder encoder;
052
053    // Non-customizable variables
054    protected Vector<ILoggingEvent> batch = new Vector<>();
055    protected AtomicBoolean isFlushing = new AtomicBoolean(false);
056    protected boolean mustReflush = false;
057    protected boolean warnAboutMaxQueueSize = true;
058
059    // Utils
060    protected ScheduledExecutorService scheduledExecutorService;
061    protected ScheduledFuture<?> scheduledFuture;
062    protected ObjectMapper dataMapper;
063    protected Logger logger;
064    protected int retrySize = 0;
065    protected int retries = 0;
066    protected boolean disabled = false;
067
068    protected ThreadFactory threadFactory = r -> {
069        Thread thread = Executors.defaultThreadFactory().newThread(r);
070        thread.setName("logfire-appender");
071        thread.setDaemon(true);
072        return thread;
073    };
074
075    public LogfireAppender() {
076        logger = LoggerFactory.getLogger(LogfireAppender.class);
077
078        dataMapper = new ObjectMapper()
079                .setSerializationInclusion(JsonInclude.Include.NON_NULL)
080                .setPropertyNamingStrategy(PropertyNamingStrategies.UPPER_CAMEL_CASE);
081
082        scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(threadFactory);
083        scheduledFuture = scheduledExecutorService.scheduleWithFixedDelay(new LogfireSender(), batchInterval, batchInterval, TimeUnit.MILLISECONDS);
084    }
085
086    @Override
087    protected void append(ILoggingEvent event) {
088        if (disabled)
089            return;
090
091        if (event.getLoggerName().equals(LogfireAppender.class.getName()))
092            return;
093
094        if (this.ingestUrl.isEmpty() || this.sourceToken == null || this.sourceToken.isEmpty()) {
095            // Prevent potential deadlock, when a blocking logger is configured - avoid using logger directly in append
096            startThread("logfire-warning-logger", () -> {
097                logger.warn("Missing Source token for Logfire - disabling LogfireAppender. Find out how to fix this at: https://logfire.ai/docs/logs/java ");
098            });
099            this.disabled = true;
100            return;
101        }
102
103        if (batch.size() < maxQueueSize) {
104            batch.add(event);
105        }
106
107        if (warnAboutMaxQueueSize && batch.size() == maxQueueSize) {
108            this.warnAboutMaxQueueSize = false;
109            // Prevent potential deadlock, when a blocking logger is configured - avoid using logger directly in append
110            startThread("logfire-error-logger", () -> {
111                logger.error("Maximum number of messages in queue reached ({}). New messages will be dropped.", maxQueueSize);
112            });
113        }
114
115        if (batch.size() >= batchSize) {
116            if (isFlushing.get())
117                return;
118
119            startThread("logfire-appender-flush", new LogfireSender());
120        }
121    }
122
123    protected void startThread(String threadName, Runnable runnable) {
124        Thread thread = Executors.defaultThreadFactory().newThread(runnable);
125        thread.setName(threadName);
126        thread.start();
127    }
128
129    protected void flush() {
130        if (batch.isEmpty())
131            return;
132
133        // Guaranteed to not be running concurrently
134        if (isFlushing.getAndSet(true))
135            return;
136
137        mustReflush = false;
138
139        int flushedSize = batch.size();
140        if (flushedSize > batchSize) {
141            flushedSize = batchSize;
142            mustReflush = true;
143        }
144        if (retries > 0 && flushedSize > retrySize) {
145            flushedSize = retrySize;
146            mustReflush = true;
147        }
148
149        if (!flushLogs(flushedSize)) {
150            mustReflush = true;
151        }
152
153        isFlushing.set(false);
154
155        if (mustReflush || batch.size() >= batchSize) {
156            flush();
157        }
158    }
159
160    protected boolean flushLogs(int flushedSize) {
161        retrySize = flushedSize;
162
163        try {
164            if (retries > maxRetries) {
165                batch.subList(0, flushedSize).clear();
166                logger.error("Dropped batch of {} logs.", flushedSize);
167                warnAboutMaxQueueSize = true;
168                retries = 0;
169
170                return true;
171            }
172
173            if (retries > 0) {
174                logger.info("Retrying to send {} logs to Logfire ({} / {})", flushedSize, retries, maxRetries);
175                try {
176                    TimeUnit.MILLISECONDS.sleep(retrySleepMilliseconds);
177                } catch (InterruptedException e) {
178                    // Continue
179                }
180            }
181
182            LogfireResponse response = callHttpURLConnection(flushedSize);
183
184            if (response.getStatus() >= 300 || response.getStatus() < 200) {
185                logger.error("Error calling Logfire : {} ({})", response.getError(), response.getStatus());
186                retries++;
187
188                return false;
189            }
190
191            batch.subList(0, flushedSize).clear();
192            warnAboutMaxQueueSize = true;
193            retries = 0;
194
195            return true;
196
197        } catch (ConcurrentModificationException e) {
198            logger.error("Error clearing {} logs from batch, will retry immediately.", flushedSize, e);
199            retries = maxRetries; // No point in retrying to send the data
200
201        } catch (JsonProcessingException e) {
202            logger.error("Error processing JSON data : {}", e.getMessage(), e);
203            retries = maxRetries; // No point in retrying when batch cannot be processed into JSON
204
205        } catch (Exception e) {
206            logger.error("Error trying to call Logfire : {}", e.getMessage(), e);
207        }
208
209        retries++;
210
211        return false;
212    }
213
214    protected LogfireResponse callHttpURLConnection(int flushedSize) throws IOException {
215        HttpURLConnection connection = getHttpURLConnection();
216
217        try {
218            connection.connect();
219        } catch (Exception e) {
220            logger.error("Error trying to call Logfire : {}", e.getMessage(), e);
221        }
222
223        try (OutputStream os = connection.getOutputStream()) {
224            byte[] input = batchToJson(flushedSize).getBytes(StandardCharsets.UTF_8);
225            os.write(input, 0, input.length);
226            os.flush();
227        }
228
229        connection.disconnect();
230
231        return new LogfireResponse(connection.getResponseMessage(), connection.getResponseCode());
232    }
233
234    protected HttpURLConnection getHttpURLConnection() throws IOException {
235        HttpURLConnection httpURLConnection = (HttpURLConnection) new URL(this.ingestUrl).openConnection();
236        httpURLConnection.setDoOutput(true);
237        httpURLConnection.setDoInput(true);
238        httpURLConnection.setRequestProperty("User-Agent", this.userAgent);
239        httpURLConnection.setRequestProperty("Accept", "application/json");
240        httpURLConnection.setRequestProperty("Content-Type", "application/json");
241        httpURLConnection.setRequestProperty("Charset", "UTF-8");
242        httpURLConnection.setRequestProperty("Authorization", String.format("Bearer %s", this.sourceToken));
243        httpURLConnection.setRequestMethod("POST");
244        httpURLConnection.setConnectTimeout(this.connectTimeout);
245        httpURLConnection.setReadTimeout(this.readTimeout);
246        return httpURLConnection;
247    }
248
249    protected String batchToJson(int flushedSize) throws JsonProcessingException {
250        return this.dataMapper.writeValueAsString(
251            new ArrayList<>(batch.subList(0, flushedSize))
252                .stream()
253                .map(this::buildPostData)
254                .collect(Collectors.toList())
255        );
256    }
257
258    protected Map<String, Object> buildPostData(ILoggingEvent event) {
259        Map<String, Object> logLine = new HashMap<>();
260        logLine.put("dt", Long.toString(event.getTimeStamp()));
261        logLine.put("level", event.getLevel().toString());
262        logLine.put("app", this.appName);
263        logLine.put("message", generateLogMessage(event));
264        logLine.put("meta", generateLogMeta(event));
265        logLine.put("runtime", generateLogRuntime(event));
266        logLine.put("args", event.getArgumentArray());
267        if (event.getThrowableProxy() != null) {
268            logLine.put("throwable", generateLogThrowable(event.getThrowableProxy()));
269        }
270
271        return logLine;
272    }
273
274    protected String generateLogMessage(ILoggingEvent event) {
275        return this.encoder != null ? new String(this.encoder.encode(event)) : event.getFormattedMessage();
276    }
277
278    protected Map<String, Object> generateLogMeta(ILoggingEvent event) {
279        Map<String, Object> logMeta = new HashMap<>();
280        logMeta.put("logger", event.getLoggerName());
281
282        if (!mdcFields.isEmpty() && !event.getMDCPropertyMap().isEmpty()) {
283            for (Entry<String, String> entry : event.getMDCPropertyMap().entrySet()) {
284                if (mdcFields.contains(entry.getKey())) {
285                    String type = mdcTypes.get(mdcFields.indexOf(entry.getKey()));
286                    logMeta.put(entry.getKey(), getMetaValue(type, entry.getValue()));
287                }
288            }
289        }
290
291        return logMeta;
292    }
293
294    protected Map<String, Object> generateLogRuntime(ILoggingEvent event) {
295        Map<String, Object> logRuntime = new HashMap<>();
296        logRuntime.put("thread", event.getThreadName());
297
298        if (event.hasCallerData()) {
299            StackTraceElement[] callerData = event.getCallerData();
300
301            if (callerData.length > 0) {
302                StackTraceElement callerContext = callerData[0];
303
304                logRuntime.put("class", callerContext.getClassName());
305                logRuntime.put("method", callerContext.getMethodName());
306                logRuntime.put("file", callerContext.getFileName());
307                logRuntime.put("line", callerContext.getLineNumber());
308            }
309        }
310
311        return logRuntime;
312    }
313
314    protected Map<String, Object> generateLogThrowable(IThrowableProxy throwable) {
315        Map<String, Object> logThrowable = new HashMap<>();
316        logThrowable.put("message", throwable.getMessage());
317        logThrowable.put("class", throwable.getClassName());
318        logThrowable.put("stackTrace", throwable.getStackTraceElementProxyArray());
319        if (throwable.getCause() != null) {
320            logThrowable.put("cause", generateLogThrowable(throwable.getCause()));
321        }
322
323        return logThrowable;
324    }
325
326    protected Object getMetaValue(String type, String value) {
327        try {
328            switch (type) {
329                case "int":
330                    return Integer.valueOf(value);
331                case "long":
332                    return Long.valueOf(value);
333                case "boolean":
334                    return Boolean.valueOf(value);
335            }
336        } catch (NumberFormatException e) {
337            logger.error("Error getting meta value - {}", e.getMessage(), e);
338        }
339
340        return value;
341    }
342
343    public class LogfireSender implements Runnable {
344        @Override
345        public void run() {
346            try {
347                flush();
348            } catch (Exception e) {
349                logger.error("Error trying to flush : {}", e.getMessage(), e);
350                if (isFlushing.get()) {
351                    isFlushing.set(false);
352                }
353            }
354        }
355    }
356
357    /**
358     * Sets the application name for Logfire indexation.
359     *
360     * @param appName
361     *            application name
362     */
363    public void setAppName(String appName) {
364        this.appName = appName;
365    }
366
367    /**
368     * Sets the Logfire ingest API url.
369     *
370     * @param ingestUrl
371     *            Logfire ingest url
372     */
373    public void setIngestUrl(String ingestUrl) {
374        this.ingestUrl = ingestUrl;
375    }
376
377    /**
378     * Sets your Logfire source token.
379     *
380     * @param sourceToken
381     *            your Logfire source token
382     */
383    public void setSourceToken(String sourceToken) {
384        this.sourceToken = sourceToken;
385    }
386
387    /**
388     * Deprecated! Kept for backward compatibility.
389     * Sets your Logfire source token if unset.
390     *
391     * @param ingestKey
392     *            your Logfire source token
393     */
394    public void setIngestKey(String ingestKey) {
395        if (this.sourceToken == null) {
396            return;
397        }
398        this.sourceToken = ingestKey;
399    }
400
401    public void setUserAgent(String userAgent) {
402        this.userAgent = userAgent;
403    }
404
405    /**
406     * Sets the MDC fields that will be sent as metadata, separated by a comma.
407     *
408     * @param mdcFields
409     *            MDC fields to include in structured logs
410     */
411    public void setMdcFields(String mdcFields) {
412        this.mdcFields = Arrays.asList(mdcFields.split(","));
413    }
414
415    /**
416     * Sets the MDC fields types that will be sent as metadata, in the same order as <i>mdcFields</i> are set
417     * up, separated by a comma. Possible values are <i>string</i>, <i>boolean</i>, <i>int</i> and <i>long</i>.
418     *
419     * @param mdcTypes
420     *            MDC fields types
421     */
422    public void setMdcTypes(String mdcTypes) {
423        this.mdcTypes = Arrays.asList(mdcTypes.split(","));
424    }
425
426    /**
427     * Sets the maximum number of messages in the queue. Messages over the limit will be dropped.
428     *
429     * @param maxQueueSize
430     *            max size of the message queue
431     */
432    public void setMaxQueueSize(int maxQueueSize) {
433        this.maxQueueSize = maxQueueSize;
434    }
435
436    /**
437     * Sets the batch size for the number of messages to be sent via the API
438     *
439     * @param batchSize
440     *            size of the message batch
441     */
442    public void setBatchSize(int batchSize) {
443        this.batchSize = batchSize;
444    }
445
446    /**
447     * Get the batch size for the number of messages to be sent via the API
448     */
449    public int getBatchSize() {
450        return batchSize;
451    }
452
453    /**
454     * Sets the maximum wait time for a batch to be sent via the API, in milliseconds.
455     *
456     * @param batchInterval
457     *            maximum wait time for message batch [ms]
458     */
459    public void setBatchInterval(int batchInterval) {
460        scheduledFuture.cancel(false);
461        scheduledFuture = scheduledExecutorService.scheduleWithFixedDelay(new LogfireSender(), batchInterval, batchInterval, TimeUnit.MILLISECONDS);
462
463        this.batchInterval = batchInterval;
464    }
465
466    /**
467     * Sets the connection timeout of the underlying HTTP client, in milliseconds.
468     *
469     * @param connectTimeout
470     *            client connection timeout [ms]
471     */
472    public void setConnectTimeout(int connectTimeout) {
473        this.connectTimeout = connectTimeout;
474    }
475
476    /**
477     * Sets the read timeout of the underlying HTTP client, in milliseconds.
478     *
479     * @param readTimeout
480     *            client read timeout
481     */
482    public void setReadTimeout(int readTimeout) {
483        this.readTimeout = readTimeout;
484    }
485
486    /**
487     * Sets the maximum number of retries for sending logs to Logfire. After that, current batch of logs will be dropped.
488     *
489     * @param maxRetries
490     *            max number of retries for sending logs
491     */
492    public void setMaxRetries(int maxRetries) {
493        this.maxRetries = maxRetries;
494    }
495
496    /**
497     * Sets the number of milliseconds to sleep before retrying to send logs to Logfire.
498     *
499     * @param retrySleepMilliseconds
500     *            number of milliseconds to sleep before retry
501     */
502    public void setRetrySleepMilliseconds(int retrySleepMilliseconds) {
503        this.retrySleepMilliseconds = retrySleepMilliseconds;
504    }
505
506    /**
507     * Registers a dynamically loaded Module object to ObjectMapper used for serialization of logged data.
508     *
509     * @param className
510     *            fully qualified class name of the module, eg. "com.fasterxml.jackson.datatype.jsr310.JavaTimeModule"
511     */
512    public void setObjectMapperModule(String className) {
513        try {
514            Module module = (Module) Class.forName(className).newInstance();
515            dataMapper.registerModule(module);
516            logger.info("Module '{}' successfully registered in ObjectMapper.", className);
517        } catch (ClassNotFoundException|InstantiationException|IllegalAccessException e) {
518            logger.error("Module '{}' couldn't be registered in ObjectMapper : ", className, e);
519        }
520    }
521
522    public void setEncoder(PatternLayoutEncoder encoder) {
523        this.encoder = encoder;
524    }
525
526    public boolean isDisabled() {
527        return this.disabled;
528    }
529
530    @Override
531    public void stop() {
532        scheduledExecutorService.shutdown();
533        mustReflush = true;
534        flush();
535        super.stop();
536    }
537}