From eaca3d0ac8d2264a89f6143b13214931c41b696e Mon Sep 17 00:00:00 2001 From: zxr <271055687@qq.com> Date: Tue, 21 Jul 2026 21:21:47 +0800 Subject: [PATCH] fix: preserve log alert recovery events --- internal/ingest/alert_forward.go | 2 +- internal/ingest/alert_outbox.go | 50 ++++++++++++++++++++++++++++---- internal/ingest/engine.go | 31 ++++++++++++++++++-- 3 files changed, 74 insertions(+), 9 deletions(-) diff --git a/internal/ingest/alert_forward.go b/internal/ingest/alert_forward.go index f6ec8a7..0f9661c 100644 --- a/internal/ingest/alert_forward.go +++ b/internal/ingest/alert_forward.go @@ -108,7 +108,7 @@ func forwardAlert(body AlertReceiveBody) error { } func postAlertPayload(cfg *config.AlertForwardConf, path string, payload []byte, traceID string) (*alertForwardResponse, error) { - req, err := http.NewRequest(http.MethodPost, cfg.BaseURL+path, bytes.NewReader(payload)) + req, err := http.NewRequest(http.MethodPost, strings.TrimRight(strings.TrimSpace(cfg.BaseURL), "/")+path, bytes.NewReader(payload)) if err != nil { return nil, fmt.Errorf("创建 Alert 转发请求失败:%w", err) } diff --git a/internal/ingest/alert_outbox.go b/internal/ingest/alert_outbox.go index cf133dd..b33c435 100644 --- a/internal/ingest/alert_outbox.go +++ b/internal/ingest/alert_outbox.go @@ -156,14 +156,19 @@ func processOneOutbox(db *gorm.DB, row models.AlertOutbox) error { } return markOutboxDead(db, row, row.RetryCount, "invalid_payload: "+err.Error(), now) } - forwardErr := forwardOutboxPayload(row.PayloadJSON, body, fmt.Sprintf("logs-outbox:%d", row.ID)) + fallbackKey := fmt.Sprintf("logs-outbox:%d", row.ID) + occurredAt := row.CreatedAt.UTC() + if occurredAt.IsZero() { + occurredAt = time.Now().UTC() + } + forwardErr := forwardOutboxPayload(row.PayloadJSON, body, fallbackKey, occurredAt) now, err := databaseNow(db) if err != nil { return err } if forwardErr != nil { if errors.Is(forwardErr, errAlertForwardDisabled) { - return markOutboxDead(db, row, row.RetryCount, forwardErr.Error(), now) + return markOutboxWaiting(db, row, forwardErr.Error(), now) } return markOutboxRetry(db, row, forwardErr.Error(), now) } @@ -178,16 +183,35 @@ func databaseNow(db *gorm.DB) (time.Time, error) { return now.UTC(), nil } -func forwardOutboxPayload(payloadJSON string, legacyBody AlertReceiveBody, fallbackSourceEventKey string) error { +func forwardOutboxPayload(payloadJSON string, legacyBody AlertReceiveBody, fallbackSourceEventKey string, occurredAt time.Time) error { var rawEvent RawEventIngestBody if err := json.Unmarshal([]byte(payloadJSON), &rawEvent); err == nil && rawEvent.SourceType != "" && len(rawEvent.RawPayload) > 0 { + if strings.TrimSpace(rawEvent.SourceEventKey) == "" { + rawEvent.SourceEventKey = fallbackSourceEventKey + } + if rawEvent.EventTime.IsZero() { + rawEvent.EventTime = occurredAt + } rawEvent.TraceID = ensureAlertTraceID(rawEvent.TraceID, firstNonEmpty(rawEvent.SourceEventKey, fallbackSourceEventKey)) return forwardRawEvent(rawEvent) } + if strings.TrimSpace(legacyBody.SourceEventKey) == "" { + legacyBody.SourceEventKey = fallbackSourceEventKey + } + if legacyBody.OccurredAt.IsZero() { + legacyBody.OccurredAt = occurredAt + } legacyBody.TraceID = ensureAlertTraceID(legacyBody.TraceID, firstNonEmpty(legacyBody.SourceEventKey, fallbackSourceEventKey)) return forwardAlert(legacyBody) } +func markOutboxWaiting(db *gorm.DB, row models.AlertOutbox, msg string, now time.Time) error { + return updateClaimedOutbox(db, row, map[string]interface{}{ + "status": outboxStatusRetrying, "next_retry_at": now.Add(30 * time.Second), + "last_error": truncateError(msg, 1024), "lease_until": nil, "lease_owner": "", + }, "retrying") +} + func forwardRawEvent(body RawEventIngestBody) error { cfg := config.Spec.AlertForward if cfg == nil || !cfg.Enabled || cfg.BaseURL == "" { @@ -246,9 +270,16 @@ func markOutboxSent(db *gorm.DB, row models.AlertOutbox, now time.Time) error { if result.RowsAffected == 0 { return nil } - return tx.Model(&models.LogEvent{}).Where("id = ? AND dispatch_outbox_id = ?", row.LogEventID, row.ID).Updates(map[string]interface{}{ + eventResult := tx.Model(&models.LogEvent{}).Where("id = ? AND dispatch_outbox_id = ?", row.LogEventID, row.ID).Updates(map[string]interface{}{ "alert_sent": true, "dispatch_status": "sent", - }).Error + }) + if eventResult.Error != nil { + return eventResult.Error + } + if eventResult.RowsAffected != 1 { + return fmt.Errorf("log event %d does not belong to outbox %d", row.LogEventID, row.ID) + } + return nil }) } @@ -286,7 +317,14 @@ func updateClaimedOutbox(db *gorm.DB, row models.AlertOutbox, updates map[string if result.RowsAffected == 0 { return nil } - return tx.Model(&models.LogEvent{}).Where("id = ? AND dispatch_outbox_id = ?", row.LogEventID, row.ID).Update("dispatch_status", eventStatus).Error + eventResult := tx.Model(&models.LogEvent{}).Where("id = ? AND dispatch_outbox_id = ?", row.LogEventID, row.ID).Update("dispatch_status", eventStatus) + if eventResult.Error != nil { + return eventResult.Error + } + if eventResult.RowsAffected != 1 { + return fmt.Errorf("log event %d does not belong to outbox %d", row.LogEventID, row.ID) + } + return nil }) } diff --git a/internal/ingest/engine.go b/internal/ingest/engine.go index 52a550b..18c2bff 100644 --- a/internal/ingest/engine.go +++ b/internal/ingest/engine.go @@ -267,11 +267,15 @@ func (e *Engine) HandleSyslog(addr *net.UDPAddr, payload []byte) { if err != nil { return 0, err } + state := "firing" + if isRecoverySignal(parsed.Message, parsed.RawLine) { + state = "resolved" + } body := AlertReceiveBody{ AlertName: matched.AlertName, Summary: summary, Description: summary, SeverityCode: firstNonEmpty(matchDetails.SeverityCode, firstNonEmpty(matched.SeverityCode, sev)), Value: parsed.Message, Labels: labels, Agent: "logs-syslog", PolicyID: matched.PolicyID, - State: "firing", SourceEventKey: logSourceEventKey(stored, "alert"), OccurredAt: occurredAt, RawData: rawBytes, + State: state, SourceEventKey: logSourceEventKey(stored, "alert"), OccurredAt: occurredAt, RawData: rawBytes, } return enqueueAlertWithDB(tx, stored.ID, body) }); err != nil { @@ -584,10 +588,18 @@ func (e *Engine) HandleTrap(addr *net.UDPAddr, pkt *gosnmp.SnmpPacket) { if err != nil { return 0, err } + state := "firing" + dictText := "" + if dict != nil { + dictText = firstNonEmpty(dict.Name, dict.Title) + } + if isRecoverySignal(readable, dictText, matched.Name, matched.AlertName) { + state = "resolved" + } body := AlertReceiveBody{ AlertName: firstNonEmpty(matched.AlertName, "SNMP Trap"), Summary: readable, Description: desc, SeverityCode: firstNonEmpty(matched.SeverityCode, sev), Value: string(vbJSON), Labels: labels, - Agent: "logs-trap", PolicyID: matched.PolicyID, State: "firing", + Agent: "logs-trap", PolicyID: matched.PolicyID, State: state, SourceEventKey: logSourceEventKey(stored, "alert"), OccurredAt: occurredAt, RawData: rawBytes, } return enqueueAlertWithDB(tx, stored.ID, body) @@ -681,6 +693,21 @@ func firstNonEmpty(a, b string) string { return b } +func isRecoverySignal(values ...string) bool { + for _, value := range values { + text := strings.ToLower(strings.TrimSpace(value)) + if text == "" { + continue + } + for _, marker := range []string{"恢复", "recovered", "recovery", "ifup", "link up", "port up", "interface up", "normal"} { + if strings.Contains(text, marker) { + return true + } + } + } + return false +} + func (e *Engine) resolveResource(sourceIP, hostname string) (resourceRef, string) { e.mu.RLock() ipMap := e.resourceByIP