package service import ( "context" "crypto/rand" "crypto/subtle" "encoding/hex" "encoding/json" "errors" "fmt" "log" "net" "net/http" "strconv" "strings" "sync" "time" "gorm.io/gorm" "gorm.io/gorm/clause" "oci-portal/internal/model" "oci-portal/internal/oci" ) // 日志回传的保留策略与后台节奏(方案A §2.4/§3)。 const ( logEventRetention = 90 * 24 * time.Hour // 过期删除阈值 logEventMaxRows = 20000 // 总量兜底上限(Payload 大,低于系统日志的 5 万) logEventCleanupTick = 24 * time.Hour // 周期清理间隔 logEventParseTick = 30 * time.Second // 解析器轮询间隔 logEventParseBatch = 200 // 单轮解析行数上限 logEventConfirmWait = 10 * time.Second // 订阅确认 GET 超时 ) // 日志回传查询分页默认值与上限。 const ( logEventDefaultPageSize = 20 logEventMaxPageSize = 100 ) // logWebhookSecretPrefix 是每租户回传 secret 的 Setting 键前缀,后接 cfgID。 const logWebhookSecretPrefix = "log_webhook_secret:" // LogEventService 承接 OCI 日志回传:secret 管理、事件入库、异步解析与清理, // 以及 P1 引导创建(SetRelayDeps 注入)与 P2 告警联动(SetNotifier 注入)。 type LogEventService struct { db *gorm.DB wg sync.WaitGroup confirm func(ctx context.Context, url string) error // 订阅确认 GET,测试可注入 configs *OciConfigService // 凭据来源(P1) relayClient oci.Client // 云端链路操作(P1) publicURL string // 面板公网基址,拼接回调 endpoint(P1) notifier *Notifier // 关键事件推送(P2) settings *SettingService // log_event 事件开关(P2) relayPollTick time.Duration // 订阅确认轮询间隔,零值用默认 relayPollTimeout time.Duration // 订阅确认轮询上限,零值用默认 } // NewLogEventService 组装依赖;调用 StartParser / StartCleanup 后台协程后生效。 func NewLogEventService(db *gorm.DB) *LogEventService { s := &LogEventService{db: db} s.confirm = s.confirmSubscription return s } // secretKey 拼接指定配置的 Setting 键。 func secretKey(cfgID uint) string { return fmt.Sprintf("%s%d", logWebhookSecretPrefix, cfgID) } // LogWebhookInfo 是回传回调地址视图;Path 供前端以公网域名拼接完整 URL。 type LogWebhookInfo struct { Path string `json:"path"` Secret string `json:"secret"` CreatedAt time.Time `json:"createdAt"` } // EnsureSecret 为配置生成(或幂等返回)回传 secret。 func (s *LogEventService) EnsureSecret(ctx context.Context, cfgID uint) (LogWebhookInfo, error) { secret, err := generateWebhookSecret() if err != nil { return LogWebhookInfo{}, err } var info LogWebhookInfo err = s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { var err error info, err = ensureSecretTx(tx, cfgID, secret) return err }) if err != nil { return LogWebhookInfo{}, fmt.Errorf("ensure webhook secret: %w", err) } return info, nil } func generateWebhookSecret() (string, error) { buf := make([]byte, 32) if _, err := rand.Read(buf); err != nil { return "", fmt.Errorf("generate webhook secret: %w", err) } return hex.EncodeToString(buf), nil } func ensureSecretTx(tx *gorm.DB, cfgID uint, secret string) (LogWebhookInfo, error) { if err := lockOciConfig(tx, cfgID); err != nil { return LogWebhookInfo{}, err } if info, ok, err := secretInfoTx(tx, cfgID); err != nil || ok { return info, err } st := model.Setting{Key: secretKey(cfgID), Value: secret, UpdatedAt: time.Now()} if err := tx.Create(&st).Error; err != nil { return LogWebhookInfo{}, fmt.Errorf("save webhook secret: %w", err) } return webhookInfo(secret, st.UpdatedAt), nil } // SecretInfo 查询配置是否已生成 secret;未生成时 ok 为 false。 func (s *LogEventService) SecretInfo(ctx context.Context, cfgID uint) (LogWebhookInfo, bool, error) { var info LogWebhookInfo var ok bool err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { if err := lockOciConfig(tx, cfgID); err != nil { return err } var err error info, ok, err = secretInfoTx(tx, cfgID) return err }) return info, ok, err } func secretInfoTx(tx *gorm.DB, cfgID uint) (LogWebhookInfo, bool, error) { var st model.Setting err := tx.First(&st, "key = ?", secretKey(cfgID)).Error if err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { return LogWebhookInfo{}, false, nil } return LogWebhookInfo{}, false, fmt.Errorf("get webhook secret: %w", err) } if st.Value == "" { return LogWebhookInfo{}, false, nil } return webhookInfo(st.Value, st.UpdatedAt), true, nil } // webhookInfo 由 secret 组装回调地址视图。 func webhookInfo(secret string, at time.Time) LogWebhookInfo { return LogWebhookInfo{ Path: "/api/v1/webhooks/oci-logs/" + secret, Secret: secret, CreatedAt: at, } } // RevokeSecret 撤销配置的回传 secret;之后旧回调地址一律 404。 func (s *LogEventService) RevokeSecret(ctx context.Context, cfgID uint) error { err := s.db.WithContext(ctx).Delete(&model.Setting{}, "key = ?", secretKey(cfgID)).Error if err != nil { return fmt.Errorf("revoke webhook secret: %w", err) } return nil } // ResolveSecret 由 secret 反查归属配置;比对使用常数时间比较,防时序侧信道。 func (s *LogEventService) ResolveSecret(ctx context.Context, secret string) (uint, bool) { if secret == "" { return 0, false } var rows []model.Setting err := s.db.WithContext(ctx). Where("key LIKE ?", logWebhookSecretPrefix+"%").Find(&rows).Error if err != nil { log.Printf("resolve webhook secret: %v", err) return 0, false } for _, row := range rows { if subtle.ConstantTimeCompare([]byte(row.Value), []byte(secret)) != 1 { continue } id, ok := webhookConfigID(row.Key) if !ok || !s.configExists(ctx, id) { return 0, false } return id, true } return 0, false } // webhookConfigID 从 secret Setting 键解析租户配置 ID。 func webhookConfigID(key string) (uint, bool) { id, err := strconv.ParseUint(strings.TrimPrefix(key, logWebhookSecretPrefix), 10, 64) return uint(id), err == nil && id > 0 } // configExists 确认 secret 对应租户仍存在;数据库异常按认证失败处理。 func (s *LogEventService) configExists(ctx context.Context, id uint) bool { var cfg model.OciConfig err := s.db.WithContext(ctx).Select("id").First(&cfg, id).Error if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { log.Printf("resolve webhook config: %v", err) } return err == nil } // Ingest 落一条回传事件;MessageID 唯一索引冲突即静默忽略(at-least-once 幂等)。 func (s *LogEventService) Ingest(ctx context.Context, cfgID uint, messageID string, payload []byte, truncated bool) error { err := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { if err := lockOciConfig(tx, cfgID); err != nil { return err } return createLogEvent(tx, cfgID, messageID, payload, truncated) }) if err != nil { return fmt.Errorf("ingest log event: %w", err) } return nil } // lockOciConfig 与租户删除共用行锁,避免并发清理后写入孤儿数据。 func lockOciConfig(tx *gorm.DB, cfgID uint) error { var cfg model.OciConfig err := tx.Clauses(clause.Locking{Strength: "UPDATE"}). Select("id").First(&cfg, cfgID).Error if err != nil { return fmt.Errorf("lock oci config %d: %w", cfgID, err) } return nil } // createLogEvent 在已锁定租户的事务内幂等写入事件。 func createLogEvent(tx *gorm.DB, cfgID uint, messageID string, payload []byte, truncated bool) error { event := model.LogEvent{ OciConfigID: cfgID, MessageID: messageID, Payload: string(payload), Truncated: truncated, ReceivedAt: time.Now(), } return tx.Clauses(clause.OnConflict{ Columns: []clause.Column{{Name: "message_id"}}, DoNothing: true, }).Create(&event).Error } // LogEventQuery 是回传事件查询参数;CfgID 为 0 表示全部租户。 type LogEventQuery struct { CfgID uint Page int PageSize int } // normalize 补分页默认值并钳制上限。 func (q LogEventQuery) normalize() LogEventQuery { if q.Page < 1 { q.Page = 1 } if q.PageSize < 1 { q.PageSize = logEventDefaultPageSize } if q.PageSize > logEventMaxPageSize { q.PageSize = logEventMaxPageSize } return q } // List 按接收顺序倒序分页查询回传事件。 func (s *LogEventService) List(ctx context.Context, q LogEventQuery) ([]model.LogEvent, int64, error) { q = q.normalize() tx := s.db.WithContext(ctx).Model(&model.LogEvent{}) if q.CfgID > 0 { tx = tx.Where("oci_config_id = ?", q.CfgID) } var total int64 if err := tx.Count(&total).Error; err != nil { return nil, 0, fmt.Errorf("count log events: %w", err) } items := make([]model.LogEvent, 0, q.PageSize) err := tx.Order("id DESC").Offset((q.Page - 1) * q.PageSize).Limit(q.PageSize).Find(&items).Error if err != nil { return nil, 0, fmt.Errorf("list log events: %w", err) } return items, total, nil } // Wait 等待在途后台协程(确认/解析/清理)退出,供进程收尾与测试同步。 func (s *LogEventService) Wait() { s.wg.Wait() } // ConfirmAsync 异步 GET 订阅确认链接完成激活;失败仅记系统日志可人工重发, // URL 白名单校验由 api 层完成,此处不再信任外部输入以外的假设。 func (s *LogEventService) ConfirmAsync(confirmURL string) { s.wg.Add(1) go func() { defer s.wg.Done() ctx, cancel := context.WithTimeout(context.Background(), logEventConfirmWait) defer cancel() if err := s.confirm(ctx, confirmURL); err != nil { log.Printf("confirm ons subscription: %v", err) } }() } // confirmSubscription 执行确认 GET;重定向只允许留在 Oracle 域内,响应体丢弃。 func (s *LogEventService) confirmSubscription(ctx context.Context, confirmURL string) error { client := &http.Client{ CheckRedirect: func(req *http.Request, _ []*http.Request) error { if !IsOracleHost(req.URL.Host) { return fmt.Errorf("redirect outside oracle domain") } return nil }, } req, err := http.NewRequestWithContext(ctx, http.MethodGet, confirmURL, nil) if err != nil { return fmt.Errorf("build confirm request: %w", err) } resp, err := client.Do(req) if err != nil { return fmt.Errorf("confirm request: %w", sanitizeURLError(err)) } defer resp.Body.Close() if resp.StatusCode >= http.StatusBadRequest { return fmt.Errorf("confirm request: status %d", resp.StatusCode) } return nil } // IsOracleHost 判定 host 是否属于 Oracle 云域(订阅确认/验签证书源白名单)。 func IsOracleHost(host string) bool { if h, _, err := net.SplitHostPort(host); err == nil { host = h } host = strings.ToLower(host) return host == "oraclecloud.com" || strings.HasSuffix(host, ".oraclecloud.com") } // StartParser 启动解析协程:周期消费未解析事件,回填类型/来源/事件时间。 func (s *LogEventService) StartParser(ctx context.Context) { s.wg.Add(1) go func() { defer s.wg.Done() ticker := time.NewTicker(logEventParseTick) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: s.parseOnce(ctx) } } }() } // parseOnce 消费一批未解析事件;失败只记日志,不中断周期调度。 func (s *LogEventService) parseOnce(ctx context.Context) { var events []model.LogEvent err := s.db.WithContext(ctx). Where("processed = ?", false).Order("id").Limit(logEventParseBatch). Find(&events).Error if err != nil { log.Printf("log event parse load: %v", err) return } if len(events) == 0 { return } for i := range events { s.processLogEvent(ctx, &events[i]) } } // processLogEvent 仅在条件更新命中原行后触发通知,删除并发胜出时静默跳过。 func (s *LogEventService) processLogEvent(ctx context.Context, event *model.LogEvent) { parsed := parseLogEvent([]byte(event.Payload)) if !s.updateParsedEvent(ctx, event, parsed) { return } s.notifyCritical(ctx, event, parsed) } // updateParsedEvent 用条件 UPDATE 禁止 Save 在删除后隐式重建事件。 func (s *LogEventService) updateParsedEvent(ctx context.Context, event *model.LogEvent, parsed parsedEvent) bool { updates := map[string]any{ "event_type": parsed.EventType, "source": parsed.Source, "source_ip": parsed.SourceIP, "event_time": parsed.EventTime, "processed": true, } res := s.db.WithContext(ctx).Model(&model.LogEvent{}). Where("id = ? AND processed = ?", event.ID, false).Updates(updates) if res.Error != nil { log.Printf("log event parse update %d: %v", event.ID, res.Error) return false } if res.RowsAffected != 1 { return false } event.EventType, event.Source, event.SourceIP, event.EventTime = parsed.EventType, parsed.Source, parsed.SourceIP, parsed.EventTime event.Processed = true return true } // onsEnvelope 覆盖 ONS 消息与 CloudEvents 审计事件的常见字段; // 真实格式以联调实测为准,提不出字段时只置 Processed 不回填。 type onsEnvelope struct { EventType string `json:"eventType"` Type string `json:"type"` Source string `json:"source"` EventTime string `json:"eventTime"` ResourceName string `json:"resourceName"` Message string `json:"message"` Data json.RawMessage `json:"data"` Identity *onsIdentity `json:"identity"` AdditionalDetails *onsAddDetails `json:"additionalDetails"` StateChange *onsStateChange `json:"stateChange"` } // onsIdentity 是 Audit 事件 data.identity 中与展示相关的字段。 type onsIdentity struct { IPAddress string `json:"ipAddress"` PrincipalName string `json:"principalName"` } // onsAddDetails 是 IDCS 登录类事件 data.additionalDetails 的补充字段; // AuditEventMapValue 为嵌套的 JSON 字符串(含 eventId 成败与失败原因)。 type onsAddDetails struct { ActorName string `json:"actorName"` ClientIP string `json:"clientIp"` AuditEventMapValue string `json:"auditEventMapValue"` } // onsStateChange 承载 Audit v2 的资源变更快照;description 供策略类事件推送文案。 type onsStateChange struct { Current struct { Description string `json:"description"` } `json:"current"` } // parsedEvent 是从消息原文提取的展示字段集;ResourceName/Actor/Outcome/Detail // 仅供 P2 推送文案,不落库。 type parsedEvent struct { EventType string Source string SourceIP string ResourceName string Actor string // 操作者(identity.principalName,登录事件回退 actorName) Outcome string // 成功 / 失败 / 空(判读不出) Detail string // 补充说明(策略描述、登录失败原因等) EventTime *time.Time } // parseLogEvent 从消息原文宽松提取事件字段;非 JSON 返回零值。 func parseLogEvent(payload []byte) parsedEvent { var env onsEnvelope if err := json.Unmarshal(payload, &env); err != nil { var batch []onsEnvelope if err := json.Unmarshal(payload, &batch); err != nil || len(batch) == 0 { return parsedEvent{} } env = batch[0] } if len(env.Data) > 0 { var inner onsEnvelope if err := json.Unmarshal(env.Data, &inner); err == nil { return envelopeFields(mergeEnvelope(env, inner)) } } return envelopeFields(env) } // mergeEnvelope 外层缺失字段时以 data 内层补齐(Audit 的 identity 等在内层)。 func mergeEnvelope(outer, inner onsEnvelope) onsEnvelope { if outer.EventType == "" { outer.EventType = inner.EventType } if outer.Type == "" { outer.Type = inner.Type } if outer.Source == "" { outer.Source = inner.Source } if outer.EventTime == "" { outer.EventTime = inner.EventTime } if outer.ResourceName == "" { outer.ResourceName = inner.ResourceName } if outer.Message == "" { outer.Message = inner.Message } if outer.Identity == nil { outer.Identity = inner.Identity } if outer.AdditionalDetails == nil { outer.AdditionalDetails = inner.AdditionalDetails } if outer.StateChange == nil { outer.StateChange = inner.StateChange } return outer } // envelopeFields 收敛字段别名并解析事件时间与推送用补充字段。 func envelopeFields(env onsEnvelope) parsedEvent { eventType := env.EventType if eventType == "" { eventType = env.Type } var eventTime *time.Time if env.EventTime != "" { if ts, err := time.Parse(time.RFC3339, env.EventTime); err == nil { eventTime = &ts } } outcome, detail := envOutcome(env) return parsedEvent{ EventType: clip(eventType, 128), Source: clip(env.Source, 64), SourceIP: clip(envIP(env), 64), ResourceName: clip(env.ResourceName, 128), Actor: clipRunes(envActor(env), 64), Outcome: outcome, Detail: clipRunes(detail, 200), EventTime: eventTime, } } // envIP 取发起方 IP:Audit 的 identity.ipAddress,登录事件回退 additionalDetails.clientIp。 func envIP(env onsEnvelope) string { if env.Identity != nil && env.Identity.IPAddress != "" { return env.Identity.IPAddress } if env.AdditionalDetails != nil { return env.AdditionalDetails.ClientIP } return "" } // envActor 取操作者:Audit 的 identity.principalName,登录事件回退 additionalDetails.actorName。 func envActor(env onsEnvelope) string { if env.Identity != nil && env.Identity.PrincipalName != "" { return env.Identity.PrincipalName } if env.AdditionalDetails != nil { return env.AdditionalDetails.ActorName } return "" } // envOutcome 判读事件成败与补充说明:登录事件解析 auditEventMapValue;其余 // Audit 事件看 message 后缀,补充说明取 stateChange.current.description(策略描述)。 // 计算类失败形如 "LaunchInstance failed with response 'NotAuthorizedOrNotFound'", // 判失败并以引号内错误码作补充说明(无其他说明时)。 func envOutcome(env onsEnvelope) (outcome, detail string) { if env.StateChange != nil { detail = env.StateChange.Current.Description } if env.AdditionalDetails != nil && env.AdditionalDetails.AuditEventMapValue != "" { return ssoOutcome(env.AdditionalDetails.AuditEventMapValue, detail) } switch { case strings.HasSuffix(env.Message, " succeeded"): outcome = "成功" case strings.HasSuffix(env.Message, " failed"): outcome = "失败" case strings.Contains(env.Message, " failed with response "): outcome = "失败" if detail == "" { detail = failedResponse(env.Message) } } return outcome, detail } // failedResponse 提取 "X failed with response 'Err'" 中引号内的错误码。 func failedResponse(msg string) string { _, after, ok := strings.Cut(msg, " failed with response ") if !ok { return "" } return strings.Trim(strings.TrimSpace(after), "'") } // ssoOutcome 解析 IDCS 审计负载(JSON 字符串):eventId 含 success/failure 定成败, // 失败时以其 message 作为原因说明。 func ssoOutcome(raw, detail string) (string, string) { var ev struct { EventID string `json:"eventId"` Message string `json:"message"` } if json.Unmarshal([]byte(raw), &ev) != nil { return "", detail } switch { case strings.Contains(ev.EventID, "success"): return "成功", detail case strings.Contains(ev.EventID, "failure"): if ev.Message != "" { detail = ev.Message } return "失败", detail } return "", detail } // clip 按模型列宽截断解析出的字段。 func clip(s string, max int) string { if len(s) > max { return s[:max] } return s } // clipRunes 按字符数截断(通知文案用,中文安全)。 func clipRunes(s string, max int) string { r := []rune(s) if len(r) > max { return string(r[:max]) } return s } // StartCleanup 启动周期清理:启动即清一次,之后每 24h 一次,随 ctx 取消退出。 func (s *LogEventService) StartCleanup(ctx context.Context) { s.wg.Add(1) go func() { defer s.wg.Done() s.cleanupOnce(ctx) ticker := time.NewTicker(logEventCleanupTick) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: s.cleanupOnce(ctx) } } }() } // cleanupOnce 执行一轮清理,失败只记日志、不中断周期调度。 func (s *LogEventService) cleanupOnce(ctx context.Context) { if err := s.cleanup(ctx, logEventRetention, logEventMaxRows); err != nil { log.Printf("log event cleanup: %v", err) } } // cleanup 先删过期记录,再对超量部分删最旧;阈值参数化便于测试。 func (s *LogEventService) cleanup(ctx context.Context, retention time.Duration, maxRows int) error { cutoff := time.Now().Add(-retention) if err := s.db.WithContext(ctx).Where("received_at < ?", cutoff).Delete(&model.LogEvent{}).Error; err != nil { return fmt.Errorf("delete expired log events: %w", err) } var total int64 if err := s.db.WithContext(ctx).Model(&model.LogEvent{}).Count(&total).Error; err != nil { return fmt.Errorf("count log events: %w", err) } overflow := int(total) - maxRows if overflow <= 0 { return nil } oldest := s.db.Model(&model.LogEvent{}).Select("id").Order("id ASC").Limit(overflow) if err := s.db.WithContext(ctx).Where("id IN (?)", oldest).Delete(&model.LogEvent{}).Error; err != nil { return fmt.Errorf("trim log events over cap: %w", err) } return nil }