Files
logs/internal/logic/otlp/query.go
2026-09-21 10:46:34 +08:00

305 lines
11 KiB
Go

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
}