4 Commits
17 changed files with 288 additions and 81 deletions
+38 -1
View File
@@ -24,7 +24,20 @@ Questions to answer:
<!-- How should queries be written? Batch operations? --> <!-- How should queries be written? Batch operations? -->
(To be filled by the team) ### 租户级数据删除
- 租户主体与本地关联数据必须在同一 GORM transaction 中按“子记录→父记录”删除,每步错误用 `%w` 返回,不得忽略。
- JSON payload/Setting 键等非外键引用要显式枚举并改写;不能只依赖 `AutoMigrate` 或 ORM association 推断级联范围。
- 与后台任务、Webhook、解析器等并发写入交叉时,先定义全局一致的行锁顺序,并在写入前重新确认父记录存在。
- 进程内 cron/缓存只在事务提交后同步;客户端取消不应中断已提交删除的必需运行时对齐。
- 批量删除/清理遇到**无法解析的 JSON payload** 时记 `log.Printf` 警告并跳过该行(原样保留),不 fail-closed 阻断整个流程——坏数据不应把删除逼到手工修库(tenantdelete.go `logSkippedTask`)。
### 大集合谓词用子查询,禁止展开 IN 列表
- 行数无上界的集合(日志事件、调用日志等)做关联删除/查询时,一律 `WHERE x IN (SELECT …)` 子查询,**不要**先 Pluck ID 再 `IN ?` 展开:绑定变量有硬上限(modernc SQLite 32766,MySQL/PG 65535),数万行即失败且重试无解。
- GORM 写法:把 `tx.Model(&T{}).Select("id").Where(...)` 作为参数传入 `Where("x IN (?)", sub)`(见 tenantdelete.go `deleteAlertHits` / `alertHitRuleIDs`)。
- 若原本的 Pluck 兼有 `FOR UPDATE` 锁定语义,保留锁定 SELECT 本身,只是不再把结果拼进后续 SQL(`lockTenantEventRows`)。
- 行数有小上界的集合(渠道、规则等配置类)可以继续用内存 ID 列表。
--- ---
@@ -48,6 +61,14 @@ Questions to answer:
<!-- Database-related mistakes your team has made --> <!-- Database-related mistakes your team has made -->
### Common Mistake: 用 `Save` 持久化在途任务的陈旧快照
**Symptom**:一条记录已被另一事务删除,在途任务随后执行 `Save` 却把该行重新插入,或覆盖并发更新后的 payload。
**Cause**:GORM `Save``UPDATE` 零命中时会回退到 `CREATE`/upsert,不适合持久化长时运行开始时读取的快照。
**Fix**:用 `WHERE id = ? AND updated_at = ?` 的条件 `Updates`,并严格要求 `RowsAffected == 1`;零命中表示记录已删除或版本已变,不得补做 `Create`
### Common Mistake: gorm 读 NULL 列到已有值的结构体字段不会清零 ### Common Mistake: gorm 读 NULL 列到已有值的结构体字段不会清零
**Symptom**:UPDATE 把可空列(如 `*time.Time`)写成 NULL 后,用**同一个结构体变量**再次 `First()` 读回,该字段仍是旧值;而新变量读取正常。断言/返回值出现「幽灵旧值」。 **Symptom**:UPDATE 把可空列(如 `*time.Time`)写成 NULL 后,用**同一个结构体变量**再次 `First()` 读回,该字段仍是旧值;而新变量读取正常。断言/返回值出现「幽灵旧值」。
@@ -66,3 +87,19 @@ db.Model(&model.AiChannel{}).
Where("id = ? AND (fail_count > 0 OR disabled_until IS NOT NULL)", id). Where("id = ? AND (fail_count > 0 OR disabled_until IS NOT NULL)", id).
Updates(map[string]any{"fail_count": 0, "disabled_until": gorm.Expr("NULL")}) Updates(map[string]any{"fail_count": 0, "disabled_until": gorm.Expr("NULL")})
``` ```
### `serializer:json` 字段走 map Updates 时手动 marshal
- 切片/结构体字段用 `gorm:"serializer:json;type:text"` 声明(如 `AiKey.Models []string`),Create/First/struct 路径自动序列化;
-`Updates(map[string]any{...})` 路径不要依赖 GORM 对 map 值应用 serializer——把值 `json.Marshal` 成 string 放进 map(见 aigateway.go `UpdateKey`),行为版本无关且可测;
- 存储格式与 serializer 一致(JSON 文本),读回仍走自动反序列化;`nil` 切片 marshal 为 `null`,读回 nil,天然表达「空 = 不限」语义。
### Common Mistake: 进程内缓存键漏掉查询维度
**Symptom**:多区间(compartment)租户在前端切换区间后,实例/卷/VCN 列表短暂显示上一个区间的数据(TTL 窗口内)。
**Cause**:`internal/oci/cached.go``ckey` 只拼了租户 OCID+资源名+region,而底层查询按 `cred.EffectiveCompartment()` 过滤——影响结果的维度没有全部进键,不同参数命中同一条缓存。
**Fix**:键值加入 `cred.CompartmentID`(空 = 租户根,天然区分)。
**Prevention**:缓存键必须覆盖影响回源结果的**全部**输入维度(租户、区间、区域、过滤参数);给 Credentials/查询结构体新增会改变结果的字段时,同步检查 `ckey` 调用点;隔离行为写进 `cached_test.go` 的 isolation 用例。
+11
View File
@@ -218,6 +218,17 @@ Responses 合成的最小事件序列为:`response.created` →
</details> </details>
### `multi-agent` 加密推理内容流式断流
> [!WARNING]
> 当 `multi-agent` 请求同时启用 `stream: true` 与
> `include: ["reasoning.encrypted_content"]` 时,上游可能在序列化大体量
> `encrypted_content` 事件期间静默断开连接,且不返回 `error` 或终态事件。
>
> 已复现场景中的断点位于 `output_index: 8` 附近:第三次搜索结束后,上游准备
> 发送较大的加密推理块时连接被中止。`multi-agent` 产生的加密推理内容体积较大,
> 因而更容易暴露该上游流式序列化缺陷。
### ZDR 与文件输入 ### ZDR 与文件输入
Responses 的 `input_file` 内容块(`file_url` / `file_data`)实测会被上游拒绝: Responses 的 `input_file` 内容块(`file_url` / `file_data`)实测会被上游拒绝:
+7
View File
@@ -8654,6 +8654,13 @@ const docTemplate = `{
}, },
"processorDescription": { "processorDescription": {
"type": "string" "type": "string"
},
"quotaNames": {
"description": "QuotaNames 是该 shape 对应的配额名(与 Limits 服务 compute limit name 同名),\n前端据此结合 limits 行判定配额与 AD 可用性。",
"type": "array",
"items": {
"type": "string"
}
} }
} }
}, },
+7
View File
@@ -8647,6 +8647,13 @@
}, },
"processorDescription": { "processorDescription": {
"type": "string" "type": "string"
},
"quotaNames": {
"description": "QuotaNames 是该 shape 对应的配额名(与 Limits 服务 compute limit name 同名),\n前端据此结合 limits 行判定配额与 AD 可用性。",
"type": "array",
"items": {
"type": "string"
}
} }
} }
}, },
+7
View File
@@ -1754,6 +1754,13 @@ definitions:
type: number type: number
processorDescription: processorDescription:
type: string type: string
quotaNames:
description: |-
QuotaNames 是该 shape 对应的配额名(与 Limits 服务 compute limit name 同名),
前端据此结合 limits 行判定配额与 AD 可用性。
items:
type: string
type: array
type: object type: object
oci-portal_internal_oci.ConsoleConnection: oci-portal_internal_oci.ConsoleConnection:
properties: properties:
+3 -1
View File
@@ -32,8 +32,10 @@ func NewCachedClient(inner Client) *CachedClient {
} }
// ckey 组缓存键;租户 OCID 在最前,写失效按前缀一锅端。 // ckey 组缓存键;租户 OCID 在最前,写失效按前缀一锅端。
// compartment 必须参与键值:列表查询按 EffectiveCompartment 过滤,
// 同租户切换区间时若共用键会串到上一个区间的缓存结果。
func ckey(cred Credentials, parts ...string) string { func ckey(cred Credentials, parts ...string) string {
return cred.TenancyOCID + "|" + strings.Join(parts, "|") return cred.TenancyOCID + "|" + cred.CompartmentID + "|" + strings.Join(parts, "|")
} }
// bust 写操作成功后失效该租户全部读缓存。 // bust 写操作成功后失效该租户全部读缓存。
+7
View File
@@ -49,6 +49,13 @@ func TestCachedClientHitAndIsolation(t *testing.T) {
if inner.instCalls != 3 { if inner.instCalls != 3 {
t.Errorf("跨租户/区域回源 %d 次, want 3", inner.instCalls) t.Errorf("跨租户/区域回源 %d 次, want 3", inner.instCalls)
} }
// 同租户不同 compartment 各自回源,不得共用缓存
inCompartment := testCred("t1")
inCompartment.CompartmentID = "ocid1.compartment.a"
_, _ = c.ListInstances(ctx, inCompartment, "r1")
if inner.instCalls != 4 {
t.Errorf("跨 compartment 回源 %d 次, want 4", inner.instCalls)
}
} }
func TestCachedClientWriteBusts(t *testing.T) { func TestCachedClientWriteBusts(t *testing.T) {
+45 -2
View File
@@ -234,8 +234,12 @@ func (c *RealClient) EnsureRelayConnector(ctx context.Context, cred Credentials,
if err != nil { if err != nil {
return RelayResource{}, err return RelayResource{}, err
} }
if res, ok, err := findRelayConnector(ctx, sc, cred.TenancyOCID); err != nil || ok { res, ok, err := findRelayConnector(ctx, sc, cred.TenancyOCID)
return res, err if err != nil {
return RelayResource{}, err
}
if ok {
return res, reconcileRelayCondition(ctx, sc, res.ID, condition)
} }
details := sch.CreateServiceConnectorDetails{ details := sch.CreateServiceConnectorDetails{
DisplayName: common.String(relayConnectorName), DisplayName: common.String(relayConnectorName),
@@ -257,6 +261,45 @@ func (c *RealClient) EnsureRelayConnector(ctx context.Context, cred Credentials,
return waitRelayConnector(ctx, sc, cred.TenancyOCID) return waitRelayConnector(ctx, sc, cred.TenancyOCID)
} }
// reconcileRelayCondition 对齐存量 Connector 的过滤条件:关键事件清单变更
// (如去除 LaunchInstance)后,已建链路经「一键创建」幂等调用原地更新,无需拆除重建。
func reconcileRelayCondition(ctx context.Context, sc sch.ServiceConnectorClient, id, condition string) error {
got, err := sc.GetServiceConnector(ctx, sch.GetServiceConnectorRequest{ServiceConnectorId: &id})
if err != nil {
return fmt.Errorf("get service connector: %w", err)
}
if !relayConditionDiffers(got.Tasks, condition) {
return nil
}
details := sch.UpdateServiceConnectorDetails{Tasks: []sch.TaskDetails{}}
if condition != "" {
details.Tasks = []sch.TaskDetails{sch.LogRuleTaskDetails{Condition: &condition}}
}
_, err = sc.UpdateServiceConnector(ctx, sch.UpdateServiceConnectorRequest{
ServiceConnectorId: &id, UpdateServiceConnectorDetails: details,
})
if err != nil {
return fmt.Errorf("update service connector condition: %w", err)
}
return nil
}
// relayConditionDiffers 判断现有任务集与期望条件是否不一致:
// 期望形态是单条 LogRule 条件;任务数、类型或条件文本不同均视为漂移。
func relayConditionDiffers(tasks []sch.TaskDetailsResponse, want string) bool {
if want == "" {
return len(tasks) > 0
}
if len(tasks) != 1 {
return true
}
rule, ok := tasks[0].(sch.LogRuleTaskDetailsResponse)
if !ok || rule.Condition == nil {
return true
}
return *rule.Condition != want
}
// findRelayConnector 按名查找存活 Connector。 // findRelayConnector 按名查找存活 Connector。
func findRelayConnector(ctx context.Context, sc sch.ServiceConnectorClient, tenancy string) (RelayResource, bool, error) { func findRelayConnector(ctx context.Context, sc sch.ServiceConnectorClient, tenancy string) (RelayResource, bool, error) {
list, err := sc.ListServiceConnectors(ctx, sch.ListServiceConnectorsRequest{ list, err := sc.ListServiceConnectors(ctx, sch.ListServiceConnectorsRequest{
+30
View File
@@ -3,6 +3,8 @@ package oci
import ( import (
"strings" "strings"
"testing" "testing"
"github.com/oracle/oci-go-sdk/v65/sch"
) )
func TestRelayEventCondition(t *testing.T) { func TestRelayEventCondition(t *testing.T) {
@@ -35,3 +37,31 @@ func TestRelayEventConditionQuoting(t *testing.T) {
t.Errorf("单引号数 = %d, want 2 (%s)", strings.Count(cond, "'"), cond) t.Errorf("单引号数 = %d, want 2 (%s)", strings.Count(cond, "'"), cond)
} }
} }
func TestRelayConditionDiffers(t *testing.T) {
cond := "data.eventName='TerminateInstance'"
rule := func(c string) sch.TaskDetailsResponse {
return sch.LogRuleTaskDetailsResponse{Condition: &c}
}
tests := []struct {
name string
tasks []sch.TaskDetailsResponse
want string
diff bool
}{
{name: "条件一致不漂移", tasks: []sch.TaskDetailsResponse{rule(cond)}, want: cond, diff: false},
{name: "条件文本不同漂移", tasks: []sch.TaskDetailsResponse{rule("data.eventName='LaunchInstance'")}, want: cond, diff: true},
{name: "无任务但期望条件漂移", tasks: nil, want: cond, diff: true},
{name: "多任务漂移", tasks: []sch.TaskDetailsResponse{rule(cond), rule(cond)}, want: cond, diff: true},
{name: "期望空且无任务不漂移", tasks: nil, want: "", diff: false},
{name: "期望空但有任务漂移", tasks: []sch.TaskDetailsResponse{rule(cond)}, want: "", diff: true},
{name: "条件缺失漂移", tasks: []sch.TaskDetailsResponse{sch.LogRuleTaskDetailsResponse{}}, want: cond, diff: true},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := relayConditionDiffers(tt.tasks, tt.want); got != tt.diff {
t.Errorf("differs = %v, want %v", got, tt.diff)
}
})
}
}
+4
View File
@@ -26,6 +26,9 @@ type ComputeShape struct {
MemoryMaxGBs float32 `json:"memoryMaxGBs,omitempty"` MemoryMaxGBs float32 `json:"memoryMaxGBs,omitempty"`
Gpus int `json:"gpus,omitempty"` Gpus int `json:"gpus,omitempty"`
ProcessorDescription string `json:"processorDescription,omitempty"` ProcessorDescription string `json:"processorDescription,omitempty"`
// QuotaNames 是该 shape 对应的配额名(与 Limits 服务 compute limit name 同名),
// 前端据此结合 limits 行判定配额与 AD 可用性。
QuotaNames []string `json:"quotaNames,omitempty"`
} }
// shapeCacheEntry 是一份 shape 清单的缓存条目。 // shapeCacheEntry 是一份 shape 清单的缓存条目。
@@ -92,6 +95,7 @@ func toComputeShape(s core.Shape) ComputeShape {
BillingType: string(s.BillingType), BillingType: string(s.BillingType),
IsFlexible: s.OcpuOptions != nil, IsFlexible: s.OcpuOptions != nil,
ProcessorDescription: deref(s.ProcessorDescription), ProcessorDescription: deref(s.ProcessorDescription),
QuotaNames: s.QuotaNames,
} }
out.Ocpus, out.MemoryInGBs = deref32(s.Ocpus), deref32(s.MemoryInGBs) out.Ocpus, out.MemoryInGBs = deref32(s.Ocpus), deref32(s.MemoryInGBs)
if s.Gpus != nil { if s.Gpus != nil {
+26 -64
View File
@@ -27,27 +27,27 @@ type aiCandidate struct {
modelOcid string modelOcid string
} }
// RespPassthrough 编排一次非流式直通调用:选渠道(priority→加权随机)→ 调用 → // routeRetry 统一编排「选渠道(priority→加权随机)→ 调用 → 可重试错误换渠道」,
// 可重试错误换渠道(整请求上限 3 次)并维护熔断;group 非空时只在同分组渠道内路由。 // 整请求上限 3 次并维护熔断;group 非空时只在同分组渠道内路由。
// 上游为 OpenAI-compatible /actions/v1/responses(实测可用,无 Oracle 文档合同)。 func routeRetry[T any](ctx context.Context, s *AiGatewayService, modelName, group, capability string, once func(*aiCandidate) (T, error)) (T, ChatMeta, error) {
func (s *AiGatewayService) RespPassthrough(ctx context.Context, raw []byte, modelName, group string) ([]byte, ChatMeta, error) { var zero T
meta := ChatMeta{} meta := ChatMeta{}
excluded := map[uint]bool{} excluded := map[uint]bool{}
var lastErr error var lastErr error
for attempt := 0; attempt < 3; attempt++ { for attempt := 0; attempt < 3; attempt++ {
cand, err := s.pick(ctx, modelName, group, "CHAT", excluded) cand, err := s.pick(ctx, modelName, group, capability, excluded)
if err != nil { if err != nil {
return nil, meta, firstErr(lastErr, err) return zero, meta, firstErr(lastErr, err)
} }
meta.ChannelID, meta.ChannelName = cand.ch.ID, cand.ch.Name meta.ChannelID, meta.ChannelName = cand.ch.ID, cand.ch.Name
payload, err := s.passthroughOnce(ctx, cand, raw) out, err := once(cand)
if err == nil { if err == nil {
s.markSuccess(ctx, cand.ch.ID) s.markSuccess(ctx, cand.ch.ID)
return payload, meta, nil return out, meta, nil
} }
retry, penalize := switchable(err) retry, penalize := switchable(err)
if !retry { if !retry {
return nil, meta, err return zero, meta, err
} }
if penalize { if penalize {
s.markFailure(ctx, cand.ch.ID) s.markFailure(ctx, cand.ch.ID)
@@ -56,7 +56,15 @@ func (s *AiGatewayService) RespPassthrough(ctx context.Context, raw []byte, mode
meta.Retries++ meta.Retries++
lastErr = err lastErr = err
} }
return nil, meta, lastErr return zero, meta, lastErr
}
// RespPassthrough 编排一次非流式直通调用。
// 上游为 OpenAI-compatible /actions/v1/responses(实测可用,无 Oracle 文档合同)。
func (s *AiGatewayService) RespPassthrough(ctx context.Context, raw []byte, modelName, group string) ([]byte, ChatMeta, error) {
return routeRetry(ctx, s, modelName, group, "CHAT", func(cand *aiCandidate) ([]byte, error) {
return s.passthroughOnce(ctx, cand, raw)
})
} }
func (s *AiGatewayService) passthroughOnce(ctx context.Context, cand *aiCandidate, raw []byte) ([]byte, error) { func (s *AiGatewayService) passthroughOnce(ctx context.Context, cand *aiCandidate, raw []byte) ([]byte, error) {
@@ -70,36 +78,13 @@ func (s *AiGatewayService) passthroughOnce(ctx context.Context, cand *aiCandidat
// RespPassthroughStream 编排流式直通:流建立成功即绑定渠道,建立失败按 switchable // RespPassthroughStream 编排流式直通:流建立成功即绑定渠道,建立失败按 switchable
// 换渠道重试;建立后的中断不重试、不计熔断(与 OpenStream 语义一致)。 // 换渠道重试;建立后的中断不重试、不计熔断(与 OpenStream 语义一致)。
func (s *AiGatewayService) RespPassthroughStream(ctx context.Context, raw []byte, modelName, group string) (io.ReadCloser, ChatMeta, error) { func (s *AiGatewayService) RespPassthroughStream(ctx context.Context, raw []byte, modelName, group string) (io.ReadCloser, ChatMeta, error) {
meta := ChatMeta{} return routeRetry(ctx, s, modelName, group, "CHAT", func(cand *aiCandidate) (io.ReadCloser, error) {
excluded := map[uint]bool{}
var lastErr error
for attempt := 0; attempt < 3; attempt++ {
cand, err := s.pick(ctx, modelName, group, "CHAT", excluded)
if err != nil {
return nil, meta, firstErr(lastErr, err)
}
meta.ChannelID, meta.ChannelName = cand.ch.ID, cand.ch.Name
cred, err := s.configs.credentialsByID(ctx, cand.ch.OciConfigID) cred, err := s.configs.credentialsByID(ctx, cand.ch.OciConfigID)
if err != nil { if err != nil {
return nil, meta, err return nil, err
} }
stream, err := s.client.GenAiCompatResponsesStream(ctx, cred, cand.ch.Region, raw, s.UpstreamWait()) return s.client.GenAiCompatResponsesStream(ctx, cred, cand.ch.Region, raw, s.UpstreamWait())
if err == nil { })
s.markSuccess(ctx, cand.ch.ID)
return stream, meta, nil
}
retry, penalize := switchable(err)
if !retry {
return nil, meta, err
}
if penalize {
s.markFailure(ctx, cand.ch.ID)
}
excluded[cand.ch.ID] = true
meta.Retries++
lastErr = err
}
return nil, meta, lastErr
} }
// firstErr 在换渠道后仍失败时优先返回上游错误(而非「无渠道」)。 // firstErr 在换渠道后仍失败时优先返回上游错误(而非「无渠道」)。
@@ -221,34 +206,11 @@ func weightedPick(chs []model.AiChannel) model.AiChannel {
return chs[len(chs)-1] return chs[len(chs)-1]
} }
// Embeddings 编排向量化调用:按 EMBEDDING 能力选渠道,可重试错误换渠道(整请求上限 3 次) // Embeddings 编排向量化调用:按 EMBEDDING 能力选渠道,可重试错误换渠道。
func (s *AiGatewayService) Embeddings(ctx context.Context, req aiwire.EmbeddingsRequest, group string) (*aiwire.EmbeddingsResponse, ChatMeta, error) { func (s *AiGatewayService) Embeddings(ctx context.Context, req aiwire.EmbeddingsRequest, group string) (*aiwire.EmbeddingsResponse, ChatMeta, error) {
meta := ChatMeta{} return routeRetry(ctx, s, req.Model, group, "EMBEDDING", func(cand *aiCandidate) (*aiwire.EmbeddingsResponse, error) {
excluded := map[uint]bool{} return s.embedOnce(ctx, cand, req)
var lastErr error })
for attempt := 0; attempt < 3; attempt++ {
cand, err := s.pick(ctx, req.Model, group, "EMBEDDING", excluded)
if err != nil {
return nil, meta, firstErr(lastErr, err)
}
meta.ChannelID, meta.ChannelName = cand.ch.ID, cand.ch.Name
resp, err := s.embedOnce(ctx, cand, req)
if err == nil {
s.markSuccess(ctx, cand.ch.ID)
return resp, meta, nil
}
retry, penalize := switchable(err)
if !retry {
return nil, meta, err
}
if penalize {
s.markFailure(ctx, cand.ch.ID)
}
excluded[cand.ch.ID] = true
meta.Retries++
lastErr = err
}
return nil, meta, lastErr
} }
// embedOnce 调用渠道向量化并装配 OpenAI 形态响应。 // embedOnce 调用渠道向量化并装配 OpenAI 形态响应。
+16
View File
@@ -554,6 +554,8 @@ func envActor(env onsEnvelope) string {
// envOutcome 判读事件成败与补充说明:登录事件解析 auditEventMapValue;其余 // envOutcome 判读事件成败与补充说明:登录事件解析 auditEventMapValue;其余
// Audit 事件看 message 后缀,补充说明取 stateChange.current.description(策略描述)。 // Audit 事件看 message 后缀,补充说明取 stateChange.current.description(策略描述)。
// 计算类失败形如 "LaunchInstance failed with response 'NotAuthorizedOrNotFound'",
// 判失败并以引号内错误码作补充说明(无其他说明时)。
func envOutcome(env onsEnvelope) (outcome, detail string) { func envOutcome(env onsEnvelope) (outcome, detail string) {
if env.StateChange != nil { if env.StateChange != nil {
detail = env.StateChange.Current.Description detail = env.StateChange.Current.Description
@@ -566,10 +568,24 @@ func envOutcome(env onsEnvelope) (outcome, detail string) {
outcome = "成功" outcome = "成功"
case strings.HasSuffix(env.Message, " failed"): case strings.HasSuffix(env.Message, " failed"):
outcome = "失败" outcome = "失败"
case strings.Contains(env.Message, " failed with response "):
outcome = "失败"
if detail == "" {
detail = failedResponse(env.Message)
}
} }
return outcome, detail 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 定成败, // ssoOutcome 解析 IDCS 审计负载(JSON 字符串):eventId 含 success/failure 定成败,
// 失败时以其 message 作为原因说明。 // 失败时以其 message 作为原因说明。
func ssoOutcome(raw, detail string) (string, string) { func ssoOutcome(raw, detail string) (string, string) {
+8
View File
@@ -308,6 +308,14 @@ func TestParseLogEvent(t *testing.T) {
wantType: "x", wantType: "x",
wantOutcome: "失败", wantOutcome: "失败",
}, },
{
name: "failed with response 判失败并提取错误码",
payload: `{"type":"com.oraclecloud.computeApi.LaunchInstance.begin","data":{
"message":"LaunchInstance failed with response 'NotAuthorizedOrNotFound'"}}`,
wantType: "com.oraclecloud.computeApi.LaunchInstance.begin",
wantOutcome: "失败",
wantDetail: "NotAuthorizedOrNotFound",
},
} }
for _, tt := range tests { for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
+17 -8
View File
@@ -1,6 +1,7 @@
package service package service
import ( import (
"cmp"
"context" "context"
"errors" "errors"
"fmt" "fmt"
@@ -23,9 +24,11 @@ const (
) )
// RelayCriticalEvents 是回传的关键事件清单(Connector Log Filter 与 P2 告警共用): // RelayCriticalEvents 是回传的关键事件清单(Connector Log Filter 与 P2 告警共用):
// 实例生命周期、用户与凭据、区域订阅、策略变更、控制台登录。 // 实例终止与电源操作、用户与凭据、区域订阅、策略变更、控制台登录。
// 不含 LaunchInstance:创建以面板自身(手动/抢机反复尝试)发起为主,回传即噪声
// (10 个 1min 抢机任务一天即可刷掉 2 万条上限),终止与电源操作才是外部风险信号。
var RelayCriticalEvents = []string{ var RelayCriticalEvents = []string{
"LaunchInstance", "TerminateInstance", "InstanceAction", "TerminateInstance", "InstanceAction",
"CreateUser", "DeleteUser", "UpdateUser", "CreateUser", "DeleteUser", "UpdateUser",
"CreateApiKey", "DeleteApiKey", "UpdateUserCapabilities", "CreateApiKey", "DeleteApiKey", "UpdateUserCapabilities",
"CreateRegionSubscription", "CreateRegionSubscription",
@@ -46,7 +49,7 @@ var relayCriticalSet = func() map[string]bool {
// 通知管理「云端事件」按此细分开关。 // 通知管理「云端事件」按此细分开关。
func relayEventClass(name string) string { func relayEventClass(name string) string {
switch name { switch name {
case "LaunchInstance", "TerminateInstance", "InstanceAction": case "TerminateInstance", "InstanceAction":
return "instance" return "instance"
case "CreateUser", "DeleteUser", "UpdateUser", case "CreateUser", "DeleteUser", "UpdateUser",
"CreateApiKey", "DeleteApiKey", "UpdateUserCapabilities": "CreateApiKey", "DeleteApiKey", "UpdateUserCapabilities":
@@ -313,24 +316,30 @@ func (s *LogEventService) teardownResources(ctx context.Context, cred oci.Creden
// ---- P2 告警联动:解析出关键事件后经 Notifier 推送 ---- // ---- P2 告警联动:解析出关键事件后经 Notifier 推送 ----
// relayEventShortName 取 CloudEvents type 末段(com.oraclecloud.ComputeApi.LaunchInstance → LaunchInstance)。 // relayEventShortName 取 CloudEvents type 中的事件短名:Audit v2 计算类事件带
// .begin/.end 阶段后缀(com.oraclecloud.ComputeApi.LaunchInstance.end),先剥阶段再取末段。
func relayEventShortName(eventType string) string { func relayEventShortName(eventType string) string {
if i := strings.LastIndex(eventType, "."); i >= 0 { t := strings.TrimSuffix(strings.TrimSuffix(eventType, ".begin"), ".end")
return eventType[i+1:] if i := strings.LastIndex(t, "."); i >= 0 {
return t[i+1:]
} }
return eventType return t
} }
// criticalEventVars 判定关键事件并生成模板变量;非关键事件 ok 为 false。 // criticalEventVars 判定关键事件并生成模板变量;非关键事件 ok 为 false。
// 成对事件只推 .begin(携带操作者/IP/成败),.end 无操作者且信息重复,不再告警;
// actor/ip/outcome/resource 缺失时兜底 —,detail 包装为独立行(空则不占行)。 // actor/ip/outcome/resource 缺失时兜底 —,detail 包装为独立行(空则不占行)。
func criticalEventVars(alias string, p parsedEvent) (map[string]string, bool) { func criticalEventVars(alias string, p parsedEvent) (map[string]string, bool) {
if strings.HasSuffix(p.EventType, ".end") {
return nil, false
}
name := relayEventShortName(p.EventType) name := relayEventShortName(p.EventType)
if name == "" || !relayCriticalSet[name] { if name == "" || !relayCriticalSet[name] {
return nil, false return nil, false
} }
return map[string]string{ return map[string]string{
"tenant": alias, "event": name, "tenant": alias, "event": name,
"resource": orDash(p.ResourceName), "actor": orDash(p.Actor), "resource": orDash(cmp.Or(p.ResourceName, p.Source)), "actor": orDash(p.Actor),
"ip": orDash(p.SourceIP), "outcome": orDash(p.Outcome), "ip": orDash(p.SourceIP), "outcome": orDash(p.Outcome),
"detail": detailLine(p.Detail), "detail": detailLine(p.Detail),
}, true }, true
+15 -3
View File
@@ -339,10 +339,22 @@ func TestCriticalEventText(t *testing.T) {
wantOK: true, wantOK: true,
want: map[string]string{"event": "CreatePolicy", "detail": "\n允许发布到 ONS Topic"}, want: map[string]string{"event": "CreatePolicy", "detail": "\n允许发布到 ONS Topic"},
}, },
{
name: "begin 阶段剥后缀命中并回退 source 作资源",
event: parsedEvent{EventType: "com.oraclecloud.computeApi.TerminateInstance.begin",
Source: "instance-20260717-1445", Actor: "IT Team", SourceIP: "137.131.7.136", Outcome: "成功"},
wantOK: true,
want: map[string]string{"event": "TerminateInstance", "resource": "instance-20260717-1445",
"actor": "IT Team", "ip": "137.131.7.136", "outcome": "成功"},
},
{name: "end 阶段不重复告警",
event: parsedEvent{EventType: "com.oraclecloud.ComputeApi.TerminateInstance.end"}, wantOK: false},
{name: "LaunchInstance 已移出清单不告警",
event: parsedEvent{EventType: "com.oraclecloud.computeApi.LaunchInstance.begin"}, wantOK: false},
{name: "List 噪声不推", event: parsedEvent{EventType: "com.oraclecloud.ComputeApi.ListInstances"}, wantOK: false}, {name: "List 噪声不推", event: parsedEvent{EventType: "com.oraclecloud.ComputeApi.ListInstances"}, wantOK: false},
{name: "空类型不推", event: parsedEvent{}, wantOK: false}, {name: "空类型不推", event: parsedEvent{}, wantOK: false},
{name: "短名直接命中", event: parsedEvent{EventType: "LaunchInstance"}, wantOK: true, {name: "短名直接命中", event: parsedEvent{EventType: "InstanceAction"}, wantOK: true,
want: map[string]string{"event": "LaunchInstance"}}, want: map[string]string{"event": "InstanceAction"}},
} }
for _, tt := range tests { for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
@@ -370,7 +382,7 @@ func TestRelayEventClass(t *testing.T) {
name string name string
class string class string
}{ }{
{"LaunchInstance", "instance"}, {"LaunchInstance", ""}, // 已移出回传清单,不归类
{"TerminateInstance", "instance"}, {"TerminateInstance", "instance"},
{"InstanceAction", "instance"}, {"InstanceAction", "instance"},
{"CreateUser", "identity"}, {"CreateUser", "identity"},
+6 -2
View File
@@ -46,6 +46,7 @@ type ConsoleSession struct {
type ConsoleService struct { type ConsoleService struct {
configs *OciConfigService configs *OciConfigService
pollInterval time.Duration // 清理残留连接的轮询间隔,测试注入缩短 pollInterval time.Duration // 清理残留连接的轮询间隔,测试注入缩短
sessionTTL time.Duration // 会话回收检查周期,测试注入缩短
mu sync.Mutex mu sync.Mutex
sessions map[string]*ConsoleSession sessions map[string]*ConsoleSession
@@ -55,6 +56,7 @@ func NewConsoleService(configs *OciConfigService) *ConsoleService {
return &ConsoleService{ return &ConsoleService{
configs: configs, configs: configs,
pollInterval: 2 * time.Second, pollInterval: 2 * time.Second,
sessionTTL: consoleSessionTTL,
sessions: map[string]*ConsoleSession{}, sessions: map[string]*ConsoleSession{},
} }
} }
@@ -161,11 +163,12 @@ func (s *ConsoleService) storeSession(cfgID uint, instanceID, region, typ, connI
s.mu.Lock() s.mu.Lock()
s.sessions[sess.ID] = sess s.sessions[sess.ID] = sess
s.mu.Unlock() s.mu.Unlock()
time.AfterFunc(consoleSessionTTL, func() { s.expire(sess.ID) }) time.AfterFunc(s.sessionTTL, func() { s.expire(sess.ID) })
return sess return sess
} }
// expire TTL 到期回收:正在使用的会话跳过(连接断开后自然停止,无续期)。 // expire TTL 到期回收:正在使用的会话跳过并重挂下一轮检查,
// 连接断开后由后续轮次回收,避免会话与云端连接常驻到进程退出。
func (s *ConsoleService) expire(id string) { func (s *ConsoleService) expire(id string) {
s.mu.Lock() s.mu.Lock()
sess, ok := s.sessions[id] sess, ok := s.sessions[id]
@@ -173,6 +176,7 @@ func (s *ConsoleService) expire(id string) {
sess.mu.Lock() sess.mu.Lock()
if sess.inUse { if sess.inUse {
ok = false ok = false
time.AfterFunc(s.sessionTTL, func() { s.expire(id) })
} else { } else {
delete(s.sessions, id) delete(s.sessions, id)
} }
+41
View File
@@ -116,3 +116,44 @@ func TestCreateConsoleSession(t *testing.T) {
}) })
} }
} }
// TestConsoleSessionExpire 验证过期回收:在用会话跳过本轮,断开后下一轮回收。
func TestConsoleSessionExpire(t *testing.T) {
tests := []struct {
name string
inUse bool
wantAlive bool
}{
{name: "在用会话跳过回收", inUse: true, wantAlive: true},
{name: "空闲会话回收并删云端连接", inUse: false, wantAlive: false},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
client := &consoleFakeClient{}
console, cfgID := newConsoleTestService(t, client)
sess, err := console.CreateSession(context.Background(), cfgID, "ocid1.instance.i", "", ConsoleTypeSerial)
if err != nil {
t.Fatalf("CreateSession: %v", err)
}
sess.MarkUse(tt.inUse)
console.expire(sess.ID)
if alive := console.Get(sess.ID) != nil; alive != tt.wantAlive {
t.Fatalf("session alive = %v, want %v", alive, tt.wantAlive)
}
if wantDel := !tt.wantAlive; (len(client.deleted) == 1) != wantDel {
t.Errorf("cloud connection deleted %v, want %v", client.deleted, wantDel)
}
if tt.inUse {
// 断开后下一轮回收(此前的缺陷:跳过后不再检查,会话常驻)
sess.MarkUse(false)
console.expire(sess.ID)
if console.Get(sess.ID) != nil {
t.Fatal("session still alive after idle expire round")
}
if len(client.deleted) != 1 {
t.Errorf("cloud connection not deleted after idle expire: %v", client.deleted)
}
}
})
}
}