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}