389 lines
13 KiB
Go
389 lines
13 KiB
Go
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_<class>;
|
|
// 通知管理「云端事件」按此细分开关。
|
|
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
|
|
}
|