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}