package otlp import ( "crypto/sha256" "encoding/hex" "encoding/json" "errors" "fmt" "sort" "strconv" "strings" "time" "git.apinb.com/bsm-sdk/core/infra" "git.apinb.com/ops/logs/internal/impl" "git.apinb.com/ops/logs/internal/models" "github.com/gin-gonic/gin" "gorm.io/gorm" ) type traceSummary struct { TraceID string `json:"trace_id"` StartTime time.Time `json:"start_time"` EndTime time.Time `json:"end_time"` DurationMs float64 `json:"duration_ms"` SpanCount int64 `json:"span_count"` ErrorCount int64 `json:"error_count"` Services string `json:"services"` BusinessSystemID *uint `json:"business_system_id,omitempty"` ResourceUID string `json:"resource_uid,omitempty"` } type dependencyEdge struct { SourceService string `json:"source_service"` TargetService string `json:"target_service"` CallCount int64 `json:"call_count"` ErrorCount int64 `json:"error_count"` AverageMs float64 `json:"average_ms"` } type dependencyAggregate struct { Calls int64 Errors int64 DurationMs float64 } type middlewareCandidate struct { CandidateKey string `json:"candidate_key"` Product string `json:"product"` DisplayName string `json:"display_name"` Address string `json:"address,omitempty"` SourceService string `json:"source_service"` BusinessSystemID uint `json:"business_system_id"` FirstSeenAt time.Time `json:"first_seen_at"` LastSeenAt time.Time `json:"last_seen_at"` SpanCount int64 `json:"span_count"` ErrorCount int64 `json:"error_count"` } func ListTraces(ctx *gin.Context) { page, size := pageParams(ctx) query := applyTraceFilters(impl.DBService.Model(&models.TraceSpan{}), ctx) var total int64 if err := query.Distinct("trace_id").Count(&total).Error; err != nil { infra.Response.Error(ctx, err) return } var rows []traceSummary err := query.Select(`trace_id, MIN(start_time) AS start_time, MAX(end_time) AS end_time, EXTRACT(EPOCH FROM (MAX(end_time) - MIN(start_time))) * 1000 AS duration_ms, COUNT(*) AS span_count, COUNT(*) FILTER (WHERE status_code = 'error') AS error_count, STRING_AGG(DISTINCT service_name, ', ' ORDER BY service_name) AS services, MAX(business_system_id) AS business_system_id, MAX(resource_uid) AS resource_uid`). Group("trace_id").Order("start_time DESC").Offset((page - 1) * size).Limit(size).Scan(&rows).Error if err != nil { infra.Response.Error(ctx, err) return } infra.Response.Success(ctx, gin.H{"total": total, "page": page, "page_size": size, "items": rows}) } func GetTrace(ctx *gin.Context) { traceID := strings.ToLower(strings.TrimSpace(ctx.Param("trace_id"))) if _, err := normalizeID(traceID, 16); err != nil { infra.Response.Error(ctx, errors.New("trace_id 无效")) return } var rows []models.TraceSpan if err := impl.DBService.Where("trace_id = ?", traceID).Order("start_time ASC, id ASC").Find(&rows).Error; err != nil { infra.Response.Error(ctx, err) return } if len(rows) == 0 { infra.Response.Error(ctx, gorm.ErrRecordNotFound) return } infra.Response.Success(ctx, gin.H{"trace_id": traceID, "spans": rows, "count": len(rows)}) } func ListDependencies(ctx *gin.Context) { query := applyTraceFilters(impl.DBService.Model(&models.TraceSpan{}), ctx) var rows []models.TraceSpan if err := query.Order("start_time DESC").Limit(100000).Find(&rows).Error; err != nil { infra.Response.Error(ctx, err) return } spanIndex := make(map[string]models.TraceSpan, len(rows)) for _, row := range rows { spanIndex[row.TraceID+":"+row.SpanID] = row } aggregates := make(map[string]*dependencyAggregate) for _, row := range rows { target := "" if row.ParentSpanID != "" { if parent, exists := spanIndex[row.TraceID+":"+row.ParentSpanID]; exists && parent.ServiceName != row.ServiceName { target = row.ServiceName addDependency(aggregates, parent.ServiceName, target, row) continue } } if row.SpanKind == "client" { var attrs map[string]interface{} _ = json.Unmarshal([]byte(row.Attributes), &attrs) target = firstTextAttribute(attrs, nil, "peer.service", "server.address", "network.peer.address") if target != "" && target != row.ServiceName { addDependency(aggregates, row.ServiceName, target, row) } } } edges := make([]dependencyEdge, 0, len(aggregates)) for key, value := range aggregates { parts := strings.SplitN(key, "\x00", 2) edges = append(edges, dependencyEdge{ SourceService: parts[0], TargetService: parts[1], CallCount: value.Calls, ErrorCount: value.Errors, AverageMs: value.DurationMs / float64(value.Calls), }) } sort.Slice(edges, func(left, right int) bool { if edges[left].SourceService != edges[right].SourceService { return edges[left].SourceService < edges[right].SourceService } return edges[left].TargetService < edges[right].TargetService }) infra.Response.Success(ctx, gin.H{"items": edges, "count": len(edges)}) } // ListMiddlewareCandidates 从调用链标准属性中识别尚待管理员确认的中间件候选,不自动创建设备。 func ListMiddlewareCandidates(ctx *gin.Context) { businessSystemID, err := strconv.ParseUint(strings.TrimSpace(ctx.Query("business_system_id")), 10, 32) if err != nil || businessSystemID == 0 { infra.Response.Error(ctx, errors.New("business_system_id 无效")) return } start := time.Now().UTC().Add(-30 * 24 * time.Hour) if value, parseErr := time.Parse(time.RFC3339, strings.TrimSpace(ctx.Query("start_time"))); parseErr == nil { start = value.UTC() } var rows []models.TraceSpan if err := impl.DBService.Where("business_system_id = ? AND start_time >= ?", uint(businessSystemID), start). Order("start_time DESC").Limit(100000).Find(&rows).Error; err != nil { infra.Response.Error(ctx, err) return } aggregates := make(map[string]*middlewareCandidate) for _, row := range rows { product, displayName, address, ok := middlewareIdentity(row) if !ok { continue } key := middlewareCandidateKey(uint(businessSystemID), product, displayName, address) item := aggregates[key] if item == nil { item = &middlewareCandidate{ CandidateKey: key, Product: product, DisplayName: displayName, Address: address, SourceService: row.ServiceName, BusinessSystemID: uint(businessSystemID), FirstSeenAt: row.StartTime, LastSeenAt: row.StartTime, } aggregates[key] = item } item.SpanCount++ if row.StatusCode == "error" { item.ErrorCount++ } if row.StartTime.Before(item.FirstSeenAt) { item.FirstSeenAt = row.StartTime } if row.StartTime.After(item.LastSeenAt) { item.LastSeenAt = row.StartTime item.SourceService = row.ServiceName } } items := make([]middlewareCandidate, 0, len(aggregates)) for _, item := range aggregates { items = append(items, *item) } sort.Slice(items, func(left, right int) bool { if !items[left].LastSeenAt.Equal(items[right].LastSeenAt) { return items[left].LastSeenAt.After(items[right].LastSeenAt) } return items[left].CandidateKey < items[right].CandidateKey }) infra.Response.Success(ctx, gin.H{"items": items, "count": len(items), "start_time": start}) } func middlewareIdentity(row models.TraceSpan) (string, string, string, bool) { attrs := make(map[string]interface{}) resourceAttrs := make(map[string]interface{}) _ = json.Unmarshal([]byte(row.Attributes), &attrs) _ = json.Unmarshal([]byte(row.ResourceAttrs), &resourceAttrs) rawProduct := firstTextAttribute(attrs, resourceAttrs, "messaging.system", "db.system.name", "db.system", "peer.service") product := knownMiddlewareProduct(rawProduct) if product == "" { return "", "", "", false } displayName := firstTextAttribute(attrs, resourceAttrs, "peer.service", "server.address", "network.peer.address", "net.peer.name") if displayName == "" { displayName = product } address := firstTextAttribute(attrs, resourceAttrs, "server.address", "network.peer.address", "net.peer.name") port := firstTextAttribute(attrs, resourceAttrs, "server.port", "network.peer.port", "net.peer.port") if address != "" && port != "" && !strings.Contains(address, ":") { address += ":" + port } return product, limitCharacters(displayName, 255), limitCharacters(address, 255), true } func knownMiddlewareProduct(value string) string { value = strings.ToLower(strings.TrimSpace(value)) for _, product := range []string{"nginx", "apache", "tomcat", "redis", "kafka", "rabbitmq", "elasticsearch", "activemq", "rocketmq", "websphere"} { if strings.Contains(value, product) { return product } } return "" } func middlewareCandidateKey(businessSystemID uint, product, name, address string) string { raw := fmt.Sprintf("%d\x00%s\x00%s\x00%s", businessSystemID, product, strings.ToLower(name), strings.ToLower(address)) digest := sha256.Sum256([]byte(raw)) return hex.EncodeToString(digest[:]) } func addDependency(aggregates map[string]*dependencyAggregate, source, target string, span models.TraceSpan) { key := source + "\x00" + target item := aggregates[key] if item == nil { item = &dependencyAggregate{} aggregates[key] = item } item.Calls++ item.DurationMs += span.DurationMs if span.StatusCode == "error" { item.Errors++ } } func applyTraceFilters(query *gorm.DB, ctx *gin.Context) *gorm.DB { if value := strings.TrimSpace(ctx.Query("trace_id")); value != "" { query = query.Where("trace_id = ?", strings.ToLower(value)) } if value := strings.TrimSpace(ctx.Query("service_name")); value != "" { scope := impl.DBService.Model(&models.TraceSpan{}).Select("trace_id").Where("service_name ILIKE ?", "%"+value+"%") query = query.Where("trace_id IN (?)", scope) } if value := strings.TrimSpace(ctx.Query("resource_uid")); value != "" { scope := impl.DBService.Model(&models.TraceSpan{}).Select("trace_id").Where("resource_uid = ?", value) query = query.Where("trace_id IN (?)", scope) } if value, err := strconv.ParseUint(ctx.Query("business_system_id"), 10, 32); err == nil && value > 0 { scope := impl.DBService.Model(&models.TraceSpan{}).Select("trace_id").Where("business_system_id = ?", uint(value)) query = query.Where("trace_id IN (?)", scope) } if value := strings.ToLower(strings.TrimSpace(ctx.Query("status"))); value == "error" { errorScope := impl.DBService.Model(&models.TraceSpan{}).Select("trace_id").Where("status_code = ?", "error") query = query.Where("trace_id IN (?)", errorScope) } else if value == "ok" { errorScope := impl.DBService.Model(&models.TraceSpan{}).Select("trace_id").Where("status_code = ?", "error") query = query.Where("trace_id NOT IN (?)", errorScope) } else if value == "unset" { setScope := impl.DBService.Model(&models.TraceSpan{}).Select("trace_id").Where("status_code IN ?", []string{"ok", "error"}) query = query.Where("trace_id NOT IN (?)", setScope) } start := time.Now().UTC().Add(-24 * time.Hour) if value, err := time.Parse(time.RFC3339, ctx.Query("start_time")); err == nil { start = value.UTC() } timeScope := impl.DBService.Model(&models.TraceSpan{}).Select("trace_id").Where("start_time >= ?", start) if value, err := time.Parse(time.RFC3339, ctx.Query("end_time")); err == nil { timeScope = timeScope.Where("start_time <= ?", value.UTC()) } return query.Where("trace_id IN (?)", timeScope) } func pageParams(ctx *gin.Context) (int, int) { page, _ := strconv.Atoi(ctx.DefaultQuery("page", "1")) size, _ := strconv.Atoi(ctx.DefaultQuery("page_size", "50")) if page < 1 { page = 1 } if size < 1 || size > 500 { size = 50 } return page, size }