Files
break c64f3d8b80
All checks were successful
构建并部署到测试环境(无 SSH) / build-and-deploy (push) Successful in 8m31s
全局审计完成
2026-08-07 11:02:52 +08:00

318 lines
15 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package audit
import (
"context"
"fmt"
"sort"
"strconv"
"time"
"github.com/break/junhong_cmp_fiber/internal/model"
retentionquery "github.com/break/junhong_cmp_fiber/internal/query/retention"
"github.com/break/junhong_cmp_fiber/pkg/constants"
"github.com/break/junhong_cmp_fiber/pkg/errors"
)
// LinkTimeline 是 request 或 correlation 的跨事实只读时间线。
type LinkTimeline struct {
RequestID *string `json:"request_id" description:"本次按请求查询的稳定ID"`
CorrelationID *string `json:"correlation_id" description:"本次按业务关联查询的稳定ID"`
AccessLogLookupRequestID *string `json:"access_log_lookup_request_id" description:"可复制到Access Log检索的request_id本接口自身不扫描Access Log"`
Nodes []LinkTimelineNode `json:"nodes" description:"跨事实来源按发生时间稳定排序的节点"`
Retention retentionquery.Info `json:"retention" description:"Audit与Integration共同在线留存边界"`
}
// LinkTimelineNode 是保留各事实源权威边界的时间线节点。
type LinkTimelineNode struct {
RecordSource string `json:"record_source" enum:"audit_event,integration_log,outbox_event,asynq_task,domain_ledger_ref" description:"事实来源 (audit_event:审计事件, integration_log:外部交互, outbox_event:可靠事件引用, asynq_task:异步任务引用, domain_ledger_ref:业务账本引用)"`
NodeID string `json:"node_id" description:"该事实来源内的稳定节点ID"`
OccurredAt time.Time `json:"occurred_at" description:"节点发生时间"`
Code string `json:"code" description:"来源内稳定动作、操作或事件编码"`
Title string `json:"title" description:"code对应的中文展示名称"`
Result string `json:"result" description:"来源内原始结果稳定编码"`
ResultName string `json:"result_name" description:"result对应的中文展示名称"`
Summary string `json:"summary" description:"已脱敏节点摘要"`
ReferenceOnly bool `json:"reference_only" description:"true表示仅保存其他事实的引用不代表该来源独立完成业务状态变更"`
RequestID *string `json:"request_id" description:"HTTP请求关联ID"`
CorrelationID *string `json:"correlation_id" description:"跨请求业务链路ID"`
ParentEventID *string `json:"parent_event_id" description:"父审计事件ID"`
Resources []InvestigationResourceRef `json:"resources" description:"节点可稳定定位的资源引用"`
InvestigationRefs InvestigationRefs `json:"investigation_refs" description:"可继续跳转的稳定调查引用"`
Fidelity LinkageFidelity `json:"fidelity" description:"历史字段完整度和可关联能力"`
}
// LinkageFidelity 明确节点已有的稳定关联能力,不补猜历史缺失字段。
type LinkageFidelity struct {
RequestAvailable bool `json:"request_available" description:"是否有稳定request_id"`
CorrelationAvailable bool `json:"correlation_available" description:"是否有稳定correlation_id"`
ParentEventAvailable bool `json:"parent_event_available" description:"是否有稳定parent_event_id"`
DirectAuditLinkAvailable bool `json:"direct_audit_link_available" description:"是否可直接跳转审计事件详情"`
StableResourceAvailable bool `json:"stable_resource_available" description:"是否至少有一个含resource_id的稳定资源引用"`
}
// RequestTimeline 按精确 request ID 组合已持久化事实,不扫描 Access Log。
func (q *Query) RequestTimeline(ctx context.Context, requestID string) (*LinkTimeline, error) {
return q.linkTimeline(ctx, "request_id", requestID)
}
// CorrelationTimeline 按精确 correlation ID 组合跨请求业务链路。
func (q *Query) CorrelationTimeline(ctx context.Context, correlationID string) (*LinkTimeline, error) {
return q.linkTimeline(ctx, "correlation_id", correlationID)
}
func (q *Query) linkTimeline(ctx context.Context, column, value string) (*LinkTimeline, error) {
if err := q.authorize(ctx); err != nil {
return nil, err
}
if value == "" || (column != "request_id" && column != "correlation_id") {
return nil, errors.New(errors.CodeInvalidParam)
}
retention, err := retentionquery.Load(ctx, q.db, retentionquery.SourceAudit, retentionquery.SourceIntegration)
if err != nil {
return nil, err
}
auditRows, integrationRows, outboxRows, err := q.loadLinkRows(ctx, column, value, retention.OnlineFrom)
if err != nil {
return nil, err
}
events, err := q.project(ctx, auditRows)
if err != nil {
return nil, err
}
nodes := make([]LinkTimelineNode, 0, len(events)+len(integrationRows)+len(outboxRows))
integrationByAudit := integrationRefsByAuditID(integrationRows)
for index, event := range events {
refs := event.InvestigationRefs
refs.IntegrationRefs = append(refs.IntegrationRefs, integrationByAudit[auditRows[index].ID]...)
refs.IntegrationRefs = append(refs.IntegrationRefs, integrationResourceRefs(event.Resources)...)
refs.IntegrationRefs = uniqueIntegrationRefs(refs.IntegrationRefs)
nodes = append(nodes, auditTimelineNode(event, refs))
nodes = append(nodes, resourceReferenceNodes(event, refs)...)
}
for _, row := range integrationRows {
nodes = append(nodes, integrationTimelineNode(row))
}
for _, row := range outboxRows {
nodes = append(nodes, outboxTimelineNode(row))
}
sort.Slice(nodes, func(i, j int) bool {
if nodes[i].OccurredAt.Equal(nodes[j].OccurredAt) {
if nodes[i].RecordSource == nodes[j].RecordSource {
return nodes[i].NodeID < nodes[j].NodeID
}
return nodes[i].RecordSource < nodes[j].RecordSource
}
return nodes[i].OccurredAt.Before(nodes[j].OccurredAt)
})
timeline := &LinkTimeline{Nodes: nodes, Retention: retention}
if timeline.Nodes == nil {
timeline.Nodes = []LinkTimelineNode{}
}
if column == "request_id" {
timeline.RequestID = stringPointer(value)
timeline.AccessLogLookupRequestID = stringPointer(value)
} else {
timeline.CorrelationID = stringPointer(value)
}
return timeline, nil
}
func (q *Query) loadLinkRows(ctx context.Context, column, value string, onlineFrom time.Time) ([]model.AuditEvent, []model.IntegrationLog, []model.OutboxEvent, error) {
auditRows := []model.AuditEvent{}
if err := q.db.WithContext(ctx).Where(column+" = ? AND occurred_at >= ?", value, onlineFrom.UTC()).Find(&auditRows).Error; err != nil {
return nil, nil, nil, errors.Wrap(errors.CodeDatabaseError, err, "查询链路审计事件失败")
}
integrationRows := []model.IntegrationLog{}
if err := q.db.WithContext(ctx).Where(column+" = ? AND created_at >= ?", value, onlineFrom.UTC()).Find(&integrationRows).Error; err != nil {
return nil, nil, nil, errors.Wrap(errors.CodeDatabaseError, err, "查询链路外部交互失败")
}
outboxRows := []model.OutboxEvent{}
if err := q.db.WithContext(ctx).Where(column+" = ? AND created_at >= ?", value, onlineFrom.UTC()).Find(&outboxRows).Error; err != nil {
return nil, nil, nil, errors.Wrap(errors.CodeDatabaseError, err, "查询链路可靠事件失败")
}
return auditRows, integrationRows, outboxRows, nil
}
func auditTimelineNode(event EventView, refs InvestigationRefs) LinkTimelineNode {
return LinkTimelineNode{
RecordSource: constants.AuditRecordSourceAuditEvent, NodeID: event.EventID,
OccurredAt: event.OccurredAt, Code: event.ActionCode, Title: event.ActionName,
Result: event.Result, Summary: event.Summary,
RequestID: stringPointer(event.RequestID), CorrelationID: stringPointer(event.CorrelationID),
ParentEventID: stringPointer(event.ParentEventID), Resources: refs.ResourceRefs, InvestigationRefs: refs,
Fidelity: linkageFidelity(event.RequestID, event.CorrelationID, event.ParentEventID, true, len(refs.ResourceRefs) > 0),
}
}
func integrationTimelineNode(row model.IntegrationLog) LinkTimelineNode {
resource := integrationResourceRef(row)
resources := make([]InvestigationResourceRef, 0, 1)
if resource != nil {
resources = append(resources, *resource)
}
refs := InvestigationRefs{
ResourceRefs: resources, RequestID: row.RequestID, CorrelationID: row.CorrelationID,
IntegrationRefs: []IntegrationRef{{IntegrationID: row.IntegrationID}},
}
return LinkTimelineNode{
RecordSource: constants.AuditRecordSourceIntegrationLog, NodeID: row.IntegrationID,
OccurredAt: row.CreatedAt, Code: row.Operation,
Title: constants.IntegrationProviderName(row.Provider) + " · " + constants.IntegrationOperationName(row.Operation),
Result: row.Result, ResultName: constants.IntegrationResultName(row.Result), Summary: "外部交互事实",
RequestID: row.RequestID, CorrelationID: row.CorrelationID, Resources: resources, InvestigationRefs: refs,
Fidelity: linkageFidelity(pointerValue(row.RequestID), pointerValue(row.CorrelationID), "", row.AuditEventID != nil, resource != nil),
}
}
func outboxTimelineNode(row model.OutboxEvent) LinkTimelineNode {
resourceType := row.ResourceType
if resourceType == "" {
resourceType = row.AggregateType
}
resourceID := row.ResourceID
if resourceID == "" {
resourceID = row.AggregateID
}
resourceKey := row.BusinessKey
if resourceKey == "" {
resourceKey = row.AggregateID
}
resource := InvestigationResourceRef{ResourceType: resourceType, ResourceID: stringPointer(resourceID), ResourceKey: resourceKey, DisplayName: resourceKey}
refs := InvestigationRefs{
ResourceRefs: []InvestigationResourceRef{resource}, RequestID: stringPointer(row.RequestID),
CorrelationID: stringPointer(row.CorrelationID), IntegrationRefs: []IntegrationRef{},
}
return LinkTimelineNode{
RecordSource: constants.AuditRecordSourceOutboxEvent, NodeID: row.EventID,
OccurredAt: row.CreatedAt, Code: row.EventType, Title: "可靠事件:" + row.EventType,
Result: strconv.Itoa(row.Status), ResultName: constants.GetOutboxStatusName(row.Status),
Summary: fmt.Sprintf("%s/%s重试 %d 次", row.AggregateType, row.AggregateID, row.RetryCount),
RequestID: stringPointer(row.RequestID), CorrelationID: stringPointer(row.CorrelationID),
ParentEventID: stringPointer(row.ParentEventID), Resources: refs.ResourceRefs, InvestigationRefs: refs,
Fidelity: linkageFidelity(row.RequestID, row.CorrelationID, row.ParentEventID, row.ParentEventID != "", resourceType != "" && resourceID != ""),
}
}
func resourceReferenceNodes(event EventView, refs InvestigationRefs) []LinkTimelineNode {
nodes := make([]LinkTimelineNode, 0, len(event.Resources))
for _, resource := range event.Resources {
recordSource := ""
titlePrefix := ""
summary := ""
switch {
case isAsynqTaskResource(resource.ResourceType):
recordSource = constants.AuditRecordSourceAsynqTask
titlePrefix = "异步任务:"
summary = "持久化任务资源摘要;不读取或推断 Redis 队列历史"
case isDomainLedgerResource(resource.ResourceType):
recordSource = constants.AuditRecordSourceDomainLedgerRef
titlePrefix = "业务账本引用:"
summary = "状态、金额及业务结论以对应业务表为准"
default:
continue
}
resourceRef := InvestigationResourceRef{ResourceType: resource.ResourceType, ResourceID: resource.ResourceID, ResourceKey: resource.ResourceKey, DisplayName: resource.DisplayName}
nodeRefs := refs
nodeRefs.ResourceRefs = []InvestigationResourceRef{resourceRef}
nodes = append(nodes, LinkTimelineNode{
RecordSource: recordSource,
NodeID: fmt.Sprintf("%s:%s:%s:%s:%s", event.EventID, resource.ResourceType, pointerValue(resource.ResourceID), resource.ResourceKey, resource.Role),
OccurredAt: event.OccurredAt, Code: resource.ResourceType, Title: titlePrefix + resource.DisplayName,
Result: event.Result, Summary: summary, ReferenceOnly: true,
RequestID: stringPointer(event.RequestID), CorrelationID: stringPointer(event.CorrelationID),
ParentEventID: stringPointer(event.ParentEventID), Resources: nodeRefs.ResourceRefs, InvestigationRefs: nodeRefs,
Fidelity: linkageFidelity(event.RequestID, event.CorrelationID, event.ParentEventID, true, resource.ResourceID != nil || resource.ResourceKey != ""),
})
}
return nodes
}
func integrationRefsByAuditID(rows []model.IntegrationLog) map[uint][]IntegrationRef {
refs := make(map[uint][]IntegrationRef)
for _, row := range rows {
if row.AuditEventID != nil {
refs[*row.AuditEventID] = append(refs[*row.AuditEventID], IntegrationRef{IntegrationID: row.IntegrationID})
}
}
return refs
}
func integrationResourceRefs(resources []ResourceView) []IntegrationRef {
refs := make([]IntegrationRef, 0)
for _, resource := range resources {
if resource.ResourceType == constants.AuditResourceIntegrationLog && resource.ResourceKey != "" {
refs = append(refs, IntegrationRef{IntegrationID: resource.ResourceKey})
}
}
return refs
}
func uniqueIntegrationRefs(refs []IntegrationRef) []IntegrationRef {
unique := make([]IntegrationRef, 0, len(refs))
seen := make(map[string]bool, len(refs))
for _, ref := range refs {
if ref.IntegrationID == "" || seen[ref.IntegrationID] {
continue
}
seen[ref.IntegrationID] = true
unique = append(unique, ref)
}
return unique
}
func integrationResourceRef(row model.IntegrationLog) *InvestigationResourceRef {
if row.ResourceType == nil || *row.ResourceType == "" {
return nil
}
ref := InvestigationResourceRef{ResourceType: *row.ResourceType, ResourceID: row.ResourceID}
if row.ResourceKey != nil {
ref.ResourceKey = *row.ResourceKey
ref.DisplayName = *row.ResourceKey
}
return &ref
}
func linkageFidelity(requestID, correlationID, parentEventID string, directAuditLink, stableResource bool) LinkageFidelity {
return LinkageFidelity{
RequestAvailable: requestID != "", CorrelationAvailable: correlationID != "",
ParentEventAvailable: parentEventID != "", DirectAuditLinkAvailable: directAuditLink,
StableResourceAvailable: stableResource,
}
}
func pointerValue(value *string) string {
if value == nil {
return ""
}
return *value
}
func isAsynqTaskResource(resourceType string) bool {
switch resourceType {
case constants.AuditResourceDeviceBatchTask, constants.AuditResourceIotCardImportTask,
constants.AuditResourceDeviceImportTask, constants.AuditResourceAssetPackageBatchOrderTask,
constants.AuditResourceOrderPackageInvalidateTask, constants.AuditResourceExportTask:
return true
default:
return false
}
}
func isDomainLedgerResource(resourceType string) bool {
switch resourceType {
case constants.AuditResourceOrder, constants.AuditResourcePayment, constants.AuditResourceRefund,
constants.AuditResourceAgentRecharge, constants.AuditResourceRechargeOrder,
constants.AuditResourceAssetWallet, constants.AuditResourceAssetWalletTransaction,
constants.AuditResourceAgentWallet, constants.AuditResourceAgentWalletTransaction,
constants.AuditResourceAgentWalletReservation, constants.AuditResourcePackageUsage,
constants.AuditResourceApprovalInstance, constants.AuditResourceCommissionRecord,
constants.AuditResourceCommissionWithdrawal:
return true
default:
return false
}
}