001/* 002 * Copyright (c) 2023-2026, Agents-Flex (fuhai999@gmail.com). 003 * <p> 004 * Licensed under the Apache License, Version 2.0 (the "License"); 005 * you may not use this file except in compliance with the License. 006 * You may obtain a copy of the License at 007 * <p> 008 * http://www.apache.org/licenses/LICENSE-2.0 009 * <p> 010 * Unless required by applicable law or agreed to in writing, software 011 * distributed under the License is distributed on an "AS IS" BASIS, 012 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 013 * See the License for the specific language governing permissions and 014 * limitations under the License. 015 */ 016package com.agentsflex.core.observability; 017 018import com.agentsflex.core.Consts; 019import io.opentelemetry.api.GlobalOpenTelemetry; 020import io.opentelemetry.api.OpenTelemetry; 021import io.opentelemetry.api.common.AttributeKey; 022import io.opentelemetry.api.common.Attributes; 023import io.opentelemetry.api.metrics.Meter; 024import io.opentelemetry.api.trace.Tracer; 025import io.opentelemetry.exporter.logging.LoggingMetricExporter; 026import io.opentelemetry.exporter.logging.LoggingSpanExporter; 027import io.opentelemetry.exporter.otlp.metrics.OtlpGrpcMetricExporter; 028import io.opentelemetry.exporter.otlp.trace.OtlpGrpcSpanExporter; 029import io.opentelemetry.sdk.OpenTelemetrySdk; 030import io.opentelemetry.sdk.metrics.SdkMeterProvider; 031import io.opentelemetry.sdk.metrics.export.MetricExporter; 032import io.opentelemetry.sdk.metrics.export.PeriodicMetricReader; 033import io.opentelemetry.sdk.resources.Resource; 034import io.opentelemetry.sdk.trace.SdkTracerProvider; 035import io.opentelemetry.sdk.trace.SpanProcessor; 036import io.opentelemetry.sdk.trace.export.BatchSpanProcessor; 037import io.opentelemetry.sdk.trace.export.SimpleSpanProcessor; 038import io.opentelemetry.sdk.trace.export.SpanExporter; 039import org.slf4j.Logger; 040import org.slf4j.LoggerFactory; 041 042import java.time.Duration; 043import java.util.concurrent.TimeUnit; 044 045/** 046 * OpenTelemetry 可观测性统一入口。 047 * 支持 Tracing(链路追踪)和 Metrics(指标)。 048 */ 049public final class Observability { 050 private static final Logger logger = LoggerFactory.getLogger(Observability.class); 051 052 053 // === Tracing 相关 === 054 private static volatile Tracer globalTracer; 055 private static volatile SdkTracerProvider tracerProvider; 056 057 // === Metrics 相关 === 058 private static volatile Meter globalMeter; 059 private static volatile SdkMeterProvider meterProvider; 060 061 // === 共享状态 === 062 private static volatile boolean initialized = false; 063 private static volatile boolean shutdownHookRegistered = false; 064 private static volatile Throwable initError = null; 065 066 private static volatile Boolean observabilityEnabled; 067 private static volatile java.util.Set<String> excludedTools; 068 069 070 // === 自定义 Exporter 支持 === 071 private static volatile SpanExporter customSpanExporter; 072 private static volatile MetricExporter customMetricExporter; 073 074 075 private Observability() { 076 } 077 078 /** 079 * 注入自定义 Exporter 080 */ 081 public static void setCustomExporters(SpanExporter spanExporter, MetricExporter metricExporter) { 082 customSpanExporter = spanExporter; 083 customMetricExporter = metricExporter; 084 } 085 086 private static void init() { 087 if (initialized) return; 088 synchronized (Observability.class) { 089 if (initialized) return; 090 if (initError != null) { 091 throw new IllegalStateException("OpenTelemetry already failed to initialize", initError); 092 } 093 094 try { 095 OpenTelemetry openTelemetry = GlobalOpenTelemetry.get(); 096 String propagatorsClassName = openTelemetry.getPropagators().getTextMapPropagator().getClass().getSimpleName(); 097 098 // 检查是否已经被其他组件(如 SpringBoot)注册了 OpenTelemetry SDK 099 if (!"NoopTextMapPropagator".equals(propagatorsClassName)) { 100 logger.info("OpenTelemetry SDK already registered globally. Reusing existing instance."); 101 globalTracer = GlobalOpenTelemetry.getTracer("agents-flex"); 102 globalMeter = GlobalOpenTelemetry.getMeter("agents-flex"); 103 initialized = true; 104 // 注意:此处不需要注册 shutdown hook,由 GlobalOpenTelemetry 初始化的组件去关闭,比如 Spring 会管理生命周期 105 return; 106 } else { 107 GlobalOpenTelemetry.resetForTest(); 108 } 109 110 111 Resource resource = Resource.getDefault() 112 .merge(Resource.create(Attributes.of( 113 AttributeKey.stringKey("service.name"), "agents-flex", 114 AttributeKey.stringKey("service.version"), Consts.VERSION 115 ))); 116 117 118 // 1. 创建 Span 相关组件 119 SpanExporter spanExporter = customSpanExporter != null ? customSpanExporter : createSpanExporter(); 120 SpanProcessor spanProcessor = createSpanProcessor(spanExporter); 121 tracerProvider = SdkTracerProvider.builder() 122 .addSpanProcessor(spanProcessor) 123 .setResource(resource) 124 .build(); 125 126 // 2. 创建 Metric 相关组件 127 MetricExporter metricExporter = customMetricExporter != null ? customMetricExporter : createMetricExporter(); 128 meterProvider = SdkMeterProvider.builder() 129 .registerMetricReader(PeriodicMetricReader.builder(metricExporter) 130 .setInterval(Duration.ofSeconds(getMetricExportIntervalSeconds())) // 默认每60秒导出一次 131 .build()) 132 .setResource(resource) 133 .build(); 134 135 // 3. 构建并注册全局 OpenTelemetry 实例(同时包含 Trace 和 Metrics) 136 OpenTelemetrySdk.builder() 137 .setTracerProvider(tracerProvider) 138 .setMeterProvider(meterProvider) 139 .buildAndRegisterGlobal(); 140 141 // 4. 获取全局 Tracer 和 Meter 142 globalTracer = GlobalOpenTelemetry.getTracer("agents-flex"); 143 globalMeter = GlobalOpenTelemetry.getMeter("agents-flex"); 144 145 initialized = true; 146 147 // 5. 注册 Shutdown Hook(统一关闭) 148 if (!shutdownHookRegistered) { 149 Runtime.getRuntime().addShutdownHook(new Thread(Observability::shutdown)); 150 shutdownHookRegistered = true; 151 } 152 } catch (Throwable e) { 153 initError = e; 154 throw new IllegalStateException("Failed to initialize OpenTelemetry", e); 155 } 156 } 157 } 158 159 160 /** 161 * 全局可观测性开关。默认开启。 162 * 可通过系统属性 {@code agentsflex.otel.enabled} 控制(true/false)。 163 */ 164 public static boolean isEnabled() { 165 if (observabilityEnabled != null) { 166 return observabilityEnabled; 167 } 168 synchronized (Observability.class) { 169 if (observabilityEnabled != null) { 170 return observabilityEnabled; 171 } 172 // 默认 true,与现有行为一致 173 String prop = System.getProperty("agentsflex.otel.enabled", "true"); 174 observabilityEnabled = Boolean.parseBoolean(prop); 175 return observabilityEnabled; 176 } 177 } 178 179 /** 180 * 判断指定工具是否被排除在可观测性之外。 181 * 可通过系统属性 {@code agentsflex.otel.tool.excluded} 配置(逗号分隔,如 "heartbeat,debug")。 182 */ 183 public static boolean isToolExcluded(String toolName) { 184 if (toolName == null || toolName.isEmpty()) { 185 return false; 186 } 187 188 java.util.Set<String> excluded = excludedTools; 189 if (excluded != null) { 190 return excluded.contains(toolName); 191 } 192 193 synchronized (Observability.class) { 194 if (excludedTools != null) { 195 return excludedTools.contains(toolName); 196 } 197 198 String prop = System.getProperty("agentsflex.otel.tool.excluded", ""); 199 java.util.Set<String> set = new java.util.HashSet<>(); 200 if (!prop.trim().isEmpty()) { 201 for (String name : prop.split(",")) { 202 name = name.trim(); 203 if (!name.isEmpty()) { 204 set.add(name); 205 } 206 } 207 } 208 excludedTools = java.util.Collections.unmodifiableSet(set); 209 return excludedTools.contains(toolName); 210 } 211 } 212 213 214 private static long getMetricExportIntervalSeconds() { 215 String prop = System.getProperty("agentsflex.otel.metric.export.interval", "60"); 216 try { 217 long interval = Long.parseLong(prop); 218 if (interval > 0) { 219 return interval; 220 } 221 } catch (NumberFormatException ignored) { 222 // fall through 223 } 224 return 60; // 默认 60 秒 225 } 226 227 private static String getExporterType() { 228 return System.getProperty("agentsflex.otel.exporter.type", "logging").toLowerCase(); 229 } 230 231 private static SpanExporter createSpanExporter() { 232 String exporterType = getExporterType(); 233 switch (exporterType) { 234 case "otlp": 235 return OtlpGrpcSpanExporter.getDefault(); 236 case "logging": 237 return LoggingSpanExporter.create(); 238 default: 239 return createSpanExporterByClassName(exporterType); 240 } 241 } 242 243 private static SpanExporter createSpanExporterByClassName(String className) { 244 try { 245 Class<?> clazz = Class.forName(className); 246 return (SpanExporter) clazz.getDeclaredConstructor().newInstance(); 247 } catch (Exception e) { 248 logger.warn("Failed to create MetricExporter by className: " + className + ", use LoggingSpanExporter to replaced", e); 249 return LoggingSpanExporter.create(); 250 } 251 } 252 253 private static SpanProcessor createSpanProcessor(SpanExporter exporter) { 254 String exporterType = getExporterType(); 255 if ("otlp".equals(exporterType)) { 256 return BatchSpanProcessor.builder(exporter) 257 .setScheduleDelay(Duration.ofSeconds(2)) 258 .setMaxQueueSize(4096) 259 .setMaxExportBatchSize(512) 260 .setExporterTimeout(Duration.ofSeconds(10)) 261 .build(); 262 } else { 263 return SimpleSpanProcessor.create(exporter); 264 } 265 } 266 267 private static MetricExporter createMetricExporter() { 268 String exporterType = getExporterType(); 269 switch (exporterType) { 270 case "otlp": 271 return OtlpGrpcMetricExporter.getDefault(); 272 case "logging": 273 return LoggingMetricExporter.create(); 274 default: 275 return createMetricExporterByClassName(exporterType); 276 } 277 } 278 279 public static MetricExporter createMetricExporterByClassName(String className) { 280 try { 281 Class<?> clazz = Class.forName(className); 282 return (MetricExporter) clazz.getDeclaredConstructor().newInstance(); 283 } catch (Exception e) { 284 logger.warn("Failed to create MetricExporter by className: " + className + ", use LoggingMetricExporter to replaced", e); 285 return LoggingMetricExporter.create(); 286 } 287 } 288 289 290 private static void shutdown() { 291 // 先关闭 TracerProvider 292 if (tracerProvider != null) { 293 tracerProvider.shutdown().join(10, TimeUnit.SECONDS); 294 } 295 // 再关闭 MeterProvider 296 if (meterProvider != null) { 297 meterProvider.shutdown().join(10, TimeUnit.SECONDS); 298 } 299 } 300 301 302 /** 303 * 获取全局 Tracer 实例(用于链路追踪)。 304 * 首次调用时自动初始化 OpenTelemetry。 305 */ 306 public static Tracer getTracer() { 307 Tracer tracer = globalTracer; 308 if (tracer != null) { 309 return tracer; 310 } 311 synchronized (Observability.class) { 312 if (globalTracer != null) { 313 return globalTracer; 314 } 315 if (initError != null) { 316 throw new IllegalStateException("OpenTelemetry initialization failed", initError); 317 } 318 init(); 319 return globalTracer; 320 } 321 } 322 323 /** 324 * 获取全局 Meter 实例(用于指标收集)。 325 * 首次调用时自动初始化 OpenTelemetry。 326 */ 327 public static Meter getMeter() { 328 Meter meter = globalMeter; 329 if (meter != null) { 330 return meter; 331 } 332 synchronized (Observability.class) { 333 if (globalMeter != null) { 334 return globalMeter; 335 } 336 if (initError != null) { 337 throw new IllegalStateException("OpenTelemetry initialization failed", initError); 338 } 339 init(); 340 return globalMeter; 341 } 342 } 343}