package service import ( "cmp" "context" "errors" "fmt" "log" "strings" "time" "oci-portal/internal/model" "oci-portal/internal/oci" ) // ErrRelayNotConfigured 表示 P1 引导创建不可用:面板地址未配置或依赖未注入。 var ErrRelayNotConfigured = errors.New("日志回传引导创建需要配置面板地址(设置 → 安全 → 面板地址,或 PUBLIC_URL 环境变量)") // 订阅确认轮询节奏:创建订阅后 ONS 向面板投递确认消息,面板回访后订阅转 ACTIVE。 const ( relaySubPollTick = 5 * time.Second relaySubPollTimeout = 90 * time.Second relayRollbackWait = 60 * time.Second ) // RelayCriticalEvents 是回传的关键事件清单(Connector Log Filter 与 P2 告警共用): // 实例终止与电源操作、用户与凭据、区域订阅、策略变更、控制台登录。 // 不含 LaunchInstance:创建以面板自身(手动/抢机反复尝试)发起为主,回传即噪声 // (10 个 1min 抢机任务一天即可刷掉 2 万条上限),终止与电源操作才是外部风险信号。 var RelayCriticalEvents = []string{ "TerminateInstance", "InstanceAction", "CreateUser", "DeleteUser", "UpdateUser", "CreateApiKey", "DeleteApiKey", "UpdateUserCapabilities", "CreateRegionSubscription", "CreatePolicy", "UpdatePolicy", "DeletePolicy", "InteractiveLogin", } // relayCriticalSet 供 P2 按事件短名 O(1) 判定。 var relayCriticalSet = func() map[string]bool { m := make(map[string]bool, len(RelayCriticalEvents)) for _, e := range RelayCriticalEvents { m[e] = true } return m }() // relayEventClass 把关键事件归入通知子类,对应设置键 notify_event_log_event_; // 通知管理「云端事件」按此细分开关。 func relayEventClass(name string) string { switch name { case "TerminateInstance", "InstanceAction": return "instance" case "CreateUser", "DeleteUser", "UpdateUser", "CreateApiKey", "DeleteApiKey", "UpdateUserCapabilities": return "identity" case "CreatePolicy", "UpdatePolicy", "DeletePolicy": return "policy" case "CreateRegionSubscription": return "region" case "InteractiveLogin", "FederatedInteractiveLogin": return "login" } return "" } // RelayView 是 OCI 侧链路状态视图;Events 为关键事件清单,Ready 表示全链路可用。 type RelayView struct { Webhook *LogWebhookInfo `json:"webhook,omitempty"` Endpoint string `json:"endpoint,omitempty"` Topic oci.RelayResource `json:"topic"` Subscription oci.RelayResource `json:"subscription"` Connector oci.RelayResource `json:"connector"` Policy oci.RelayResource `json:"policy"` Ready bool `json:"ready"` Events []string `json:"events"` } // SetRelayDeps 注入 P1 引导创建依赖;publicURL 为面板公网基址(如 https://demo.example.com)。 func (s *LogEventService) SetRelayDeps(configs *OciConfigService, client oci.Client, publicURL string) { s.configs = configs s.relayClient = client s.publicURL = strings.TrimRight(publicURL, "/") } // SetNotifier 注入 P2 告警联动依赖;settings 控制 log_event 事件开关。 func (s *LogEventService) SetNotifier(n *Notifier, settings *SettingService) { s.notifier = n s.settings = settings } // relayBase 返回生效的公网基址:设置里的 app_url 优先(经 settings),回退启动时的 PUBLIC_URL。 func (s *LogEventService) relayBase() string { if s.settings != nil { if u := s.settings.EffectiveAppURL(); u != "" { return u } } return s.publicURL } // relayEndpoint 以公网基址拼接回调完整 URL。 func (s *LogEventService) relayEndpoint(path string) string { return s.relayBase() + path } // relayCreds 加载配置凭据与 home region(IAM 写操作路由用)。 func (s *LogEventService) relayCreds(ctx context.Context, cfgID uint) (oci.Credentials, string, error) { var cfg model.OciConfig if err := s.db.WithContext(ctx).First(&cfg, cfgID).Error; err != nil { return oci.Credentials{}, "", fmt.Errorf("find oci config %d: %w", cfgID, err) } cred, err := s.configs.credentialsOf(&cfg) if err != nil { return oci.Credentials{}, "", err } return cred, cfg.HomeRegionKey, nil } // SetupRelay 一键建立回传链路:secret → Topic → 订阅(等待面板自动确认)→ Policy → Connector; // 任一步失败逆序回滚本次新建的资源。 func (s *LogEventService) SetupRelay(ctx context.Context, cfgID uint) (RelayView, error) { if s.relayClient == nil || s.relayBase() == "" { return RelayView{}, ErrRelayNotConfigured } info, err := s.EnsureSecret(ctx, cfgID) if err != nil { return RelayView{}, err } cred, home, err := s.relayCreds(ctx, cfgID) if err != nil { return RelayView{}, err } if err := s.buildRelay(ctx, cred, home, s.relayEndpoint(info.Path)); err != nil { return RelayView{}, err } return s.RelayStatus(ctx, cfgID) } // buildRelay 依序创建链路资源;失败时逆序删除本次新建的部分后返回原错误。 func (s *LogEventService) buildRelay(ctx context.Context, cred oci.Credentials, home, endpoint string) (err error) { var undo []func(context.Context) error defer func() { if err != nil { s.rollbackRelay(undo) } }() topic, err := s.relayClient.EnsureRelayTopic(ctx, cred) undo = appendRelayUndo(undo, topic, func(c context.Context) error { return s.relayClient.DeleteRelayTopic(c, cred, topic.ID) }) if err != nil { return fmt.Errorf("创建 Topic: %w", err) } sub, err := s.relayClient.EnsureRelaySubscription(ctx, cred, topic.ID, endpoint) undo = appendRelayUndo(undo, sub, func(c context.Context) error { return s.relayClient.DeleteRelaySubscription(c, cred, sub.ID) }) if err != nil { return fmt.Errorf("创建订阅: %w", err) } if err = s.waitSubscriptionActive(ctx, cred, sub.ID); err != nil { return err } policy, err := s.relayClient.EnsureRelayPolicy(ctx, cred, home) undo = appendRelayUndo(undo, policy, func(c context.Context) error { return s.relayClient.DeleteRelayPolicy(c, cred, home, policy.ID) }) if err != nil { return fmt.Errorf("创建 Policy: %w", err) } conn, err := s.relayClient.EnsureRelayConnector(ctx, cred, topic.ID, oci.RelayEventCondition(RelayCriticalEvents)) undo = appendRelayUndo(undo, conn, func(c context.Context) error { return s.relayClient.DeleteRelayConnector(c, cred, conn.ID) }) if err != nil { return fmt.Errorf("创建 Connector: %w", err) } return nil } // appendRelayUndo 只为本次新建且拿到 ID 的资源登记回滚动作。 func appendRelayUndo(undo []func(context.Context) error, res oci.RelayResource, del func(context.Context) error) []func(context.Context) error { if res.Created && res.ID != "" { return append(undo, del) } return undo } // rollbackRelay 逆序执行回滚;使用独立超时上下文,原请求取消不影响清理。 func (s *LogEventService) rollbackRelay(undo []func(context.Context) error) { ctx, cancel := context.WithTimeout(context.Background(), relayRollbackWait) defer cancel() for i := len(undo) - 1; i >= 0; i-- { if err := undo[i](ctx); err != nil { log.Printf("relay rollback: %v", err) } } } // waitSubscriptionActive 轮询订阅至 ACTIVE;确认动作由 webhook 端点收到 ONS 消息后自动完成, // 超时通常意味着面板公网不可达或回调地址配置有误。 func (s *LogEventService) waitSubscriptionActive(ctx context.Context, cred oci.Credentials, subID string) error { deadline := time.Now().Add(s.subPollTimeout()) for { res, err := s.relayClient.GetRelaySubscription(ctx, cred, subID) if err != nil { return fmt.Errorf("查询订阅状态: %w", err) } if res.State == "ACTIVE" { return nil } if time.Now().After(deadline) { return fmt.Errorf("订阅确认超时(当前 %s):请确认面板公网可达后重试", res.State) } select { case <-ctx.Done(): return ctx.Err() case <-time.After(s.subPollTick()): } } } // subPollTick / subPollTimeout 返回订阅轮询节奏;测试注入短间隔。 func (s *LogEventService) subPollTick() time.Duration { if s.relayPollTick > 0 { return s.relayPollTick } return relaySubPollTick } func (s *LogEventService) subPollTimeout() time.Duration { if s.relayPollTimeout > 0 { return s.relayPollTimeout } return relaySubPollTimeout } // RelayStatus 查询 OCI 侧链路现状与关键事件清单。 func (s *LogEventService) RelayStatus(ctx context.Context, cfgID uint) (RelayView, error) { if s.relayClient == nil { return RelayView{}, ErrRelayNotConfigured } info, exists, err := s.SecretInfo(ctx, cfgID) if err != nil { return RelayView{}, err } view := RelayView{Events: RelayCriticalEvents} endpoint := "" if exists && s.relayBase() != "" { endpoint = s.relayEndpoint(info.Path) view.Webhook, view.Endpoint = &info, endpoint } cred, _, err := s.relayCreds(ctx, cfgID) if err != nil { return RelayView{}, err } st, err := s.relayClient.RelayState(ctx, cred, endpoint) if err != nil { return RelayView{}, fmt.Errorf("查询链路状态: %w", err) } view.Topic, view.Subscription, view.Connector, view.Policy = st.Topic, st.Subscription, st.Connector, st.Policy view.Ready = relayReady(exists, st) return view, nil } // relayReady 判定全链路可用:secret 已生成且四资源均处于活跃状态。 func relayReady(secretExists bool, st oci.RelayState) bool { return secretExists && st.Topic.State == "ACTIVE" && st.Subscription.State == "ACTIVE" && st.Connector.State == "ACTIVE" && st.Policy.ID != "" } // TeardownRelay 逆序销毁链路(Connector → Policy → 订阅 → Topic)并撤销 secret。 func (s *LogEventService) TeardownRelay(ctx context.Context, cfgID uint) error { if s.relayClient == nil { return ErrRelayNotConfigured } info, exists, err := s.SecretInfo(ctx, cfgID) if err != nil { return err } cred, home, err := s.relayCreds(ctx, cfgID) if err != nil { return err } endpoint := "" if exists { endpoint = s.relayEndpoint(info.Path) } st, err := s.relayClient.RelayState(ctx, cred, endpoint) if err != nil { return fmt.Errorf("查询链路状态: %w", err) } if err := s.teardownResources(ctx, cred, home, st); err != nil { return err } return s.RevokeSecret(ctx, cfgID) } // teardownResources 按依赖逆序删除存在的链路资源。 func (s *LogEventService) teardownResources(ctx context.Context, cred oci.Credentials, home string, st oci.RelayState) error { if st.Connector.ID != "" { if err := s.relayClient.DeleteRelayConnector(ctx, cred, st.Connector.ID); err != nil { return fmt.Errorf("删除 Connector: %w", err) } } if st.Policy.ID != "" { if err := s.relayClient.DeleteRelayPolicy(ctx, cred, home, st.Policy.ID); err != nil { return fmt.Errorf("删除 Policy: %w", err) } } if st.Subscription.ID != "" { if err := s.relayClient.DeleteRelaySubscription(ctx, cred, st.Subscription.ID); err != nil { return fmt.Errorf("删除订阅: %w", err) } } if st.Topic.ID != "" { if err := s.relayClient.DeleteRelayTopic(ctx, cred, st.Topic.ID); err != nil { return fmt.Errorf("删除 Topic: %w", err) } } return nil } // ---- P2 告警联动:解析出关键事件后经 Notifier 推送 ---- // relayEventShortName 取 CloudEvents type 中的事件短名:Audit v2 计算类事件带 // .begin/.end 阶段后缀(com.oraclecloud.ComputeApi.LaunchInstance.end),先剥阶段再取末段。 func relayEventShortName(eventType string) string { t := strings.TrimSuffix(strings.TrimSuffix(eventType, ".begin"), ".end") if i := strings.LastIndex(t, "."); i >= 0 { return t[i+1:] } return t } // criticalEventVars 判定关键事件并生成模板变量;非关键事件 ok 为 false。 // 成对事件只推 .begin(携带操作者/IP/成败),.end 无操作者且信息重复,不再告警; // actor/ip/outcome/resource 缺失时兜底 —,detail 包装为独立行(空则不占行)。 func criticalEventVars(alias string, p parsedEvent) (map[string]string, bool) { if strings.HasSuffix(p.EventType, ".end") { return nil, false } name := relayEventShortName(p.EventType) if name == "" || !relayCriticalSet[name] { return nil, false } return map[string]string{ "tenant": alias, "event": name, "resource": orDash(cmp.Or(p.ResourceName, p.Source)), "actor": orDash(p.Actor), "ip": orDash(p.SourceIP), "outcome": orDash(p.Outcome), "detail": detailLine(p.Detail), }, true } // orDash 空值兜底为 —,避免模板出现悬空标签。 func orDash(v string) string { if v == "" { return "—" } return v } // detailLine 把补充说明包装为独立行;为空时不产生多余空行。 func detailLine(v string) string { if v == "" { return "" } return "\n" + v } // notifyCritical 对关键事件推送告警;依赖未注入或对应子类开关关闭时跳过。 // MessageID 幂等入库保证同一事件只解析一次,天然防重复推送。 func (s *LogEventService) notifyCritical(ctx context.Context, e *model.LogEvent, p parsedEvent) { if s.notifier == nil { return } vars, ok := criticalEventVars(s.configAlias(ctx, e.OciConfigID), p) if !ok { return } class := relayEventClass(relayEventShortName(p.EventType)) if s.settings != nil && !s.settings.NotifyEventEnabled(ctx, "log_event_"+class) { return } s.notifier.SendTemplateAsync("log_event_"+class, vars) } // configAlias 查配置别名,失败时退化为编号占位。 func (s *LogEventService) configAlias(ctx context.Context, cfgID uint) string { var cfg model.OciConfig if err := s.db.WithContext(ctx).Select("alias").First(&cfg, cfgID).Error; err != nil { return fmt.Sprintf("租户#%d", cfgID) } return cfg.Alias }