9 Commits
Author SHA1 Message Date
wangdefa deea8629e0 实例规格透传 quotaNames,重新生成 swagger
CI / test (push) Successful in 31s
2026-07-17 16:52:49 +08:00
wangdefa 6cf9465fea 云端事件:修复实例通知并移出 LaunchInstance 回传 2026-07-17 16:52:49 +08:00
wangdefa 882eeade1e 修复跨区间缓存串数据与会话回收泄漏,收敛网关重试
CI / test (push) Successful in 29s
2026-07-17 12:19:21 +08:00
wangdefa 7019d4c5a6 文档:记录 multi-agent 加密推理流式断流问题 2026-07-17 12:19:21 +08:00
wangdefa 18e63d2dbd 补提 DASH_VERSION 至 v0.7.2
CI / test (push) Successful in 21s
Release / release (push) Successful in 52s
2026-07-16 21:33:38 +08:00
wangdefa 91999205e2 发布 v0.7.2
CI / test (push) Successful in 22s
Release / release (push) Successful in 48s
2026-07-16 21:26:39 +08:00
wangdefa cb66567256 responses 直通超时可配,流式去总超时,代理补阶段超时
CI / test (push) Successful in 31s
2026-07-16 21:20:02 +08:00
wangdefa d56678e1de 发布 v0.7.1
CI / test (push) Successful in 21s
Release / release (push) Successful in 40s
2026-07-16 16:42:29 +08:00
wangdefa 8897c847a1 审计改用日志搜索倒序,支持检索与无配额租户回退
CI / test (push) Successful in 31s
2026-07-16 16:36:31 +08:00
34 changed files with 1521 additions and 344 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 用例。
+1
View File
@@ -20,6 +20,7 @@ Gin + GORM + SQLite(纯 Go 驱动 `glebarez/sqlite`,免 CGO;默认与推荐)+ AE
| [Concurrency](./concurrency.md) | goroutine 生命周期与 context 传递 | 已填 | | [Concurrency](./concurrency.md) | goroutine 生命周期与 context 传递 | 已填 |
| [Testing](./testing.md) | table-driven 测试要求 | 已填 | | [Testing](./testing.md) | table-driven 测试要求 | 已填 |
| [Database Guidelines](./database-guidelines.md) | ORM 模式、查询、迁移 | 待填 | | [Database Guidelines](./database-guidelines.md) | ORM 模式、查询、迁移 | 待填 |
| [OCI Audit](./oci-audit.md) | 审计事件双通道数据源、检索语义与预算纪律 | 已填 |
| [Logging Guidelines](./logging-guidelines.md) | 结构化日志、日志级别 | 待填 | | [Logging Guidelines](./logging-guidelines.md) | 结构化日志、日志级别 | 待填 |
--- ---
+19
View File
@@ -0,0 +1,19 @@
# OCI 审计事件集成约定
> 2026-07 审计日志重构(数据源切换 + 检索 + 配额回退)沉淀;实现见 `internal/oci/audit.go`。
## 数据源:双通道,Search 主路 + Audit API 回退
- **Audit API(`audit.ListEvents`)无排序参数,窗口内固定按处理时间正序分页**。任何"从最新往更早"的列表需求禁止直接用它凑批——首批会拿到窗口内最旧的一段(2026-07-16 曾以此形态上线出 bug)。
- 倒序列表一律走 **Logging Search**(`loggingsearch.SearchLogs`,`search "<tenancy>/_Audit" | ... | sort by datetime desc`)。硬约束:单次查询时间窗 ≤ 14 天、limit ≤ 1000、时间过滤基于**处理时间**而非发生时间。
- **部分免费租户 Logging Search 服务配额为零**(报错含 `Rate limit exceeded` + `maxQueriesPerMinute: 0`,SDK 解析该错误体还会失败),属永久不可用,须自动回退 Audit API(小窗正序 + 前端全局重排);普通限流(配额非零)不回退。游标携带通道模式,续查不再试错。
## 检索语义
- `logContent = '*词*'` 是对整条日志 JSON **所有字段值**的包含匹配,会命中隐藏认证元数据(如 `opc-principal` 头里的 `ttype: login`),只可作服务端粗筛;**用户可见语义必须再做客户端精筛**(只匹配列表可见字段,不区分大小写,`*` 通配分段)。
- 用户输入进检索语句前必须消毒(去引号/反斜杠/控制字符、截断),见 `SanitizeAuditTerm`
## 批式回溯的预算纪律
- 单批双预算:页数(`maxAuditPages`)+ 时间(`auditBatchTimeBudget`≈20s)。全文检索命中稀疏时大窗扫描单页可达十余秒,没有时间预算会出现 3 分钟级单请求。
- 空窗按倍增扩窗(上限受 14 天查询窗约束);响应回传 `scannedThrough` 供前端展示回溯进度,前端自动补批必须封顶,由用户显式继续。
+30
View File
@@ -2,6 +2,36 @@
格式参考 [Keep a Changelog](https://keepachangelog.com/zh-CN/1.1.0/)(版本段不记日期),版本号遵循语义化版本。 格式参考 [Keep a Changelog](https://keepachangelog.com/zh-CN/1.1.0/)(版本段不记日期),版本号遵循语义化版本。
## [0.7.2]
### Added
- AI 网关设置新增「上游无响应预算」(`upstreamWaitSeconds`30..900 秒,缺省 300,持久化、即时生效):非流式为单次尝试总超时,流式为等待响应头上限,供 multi-agent / 搜索类慢模型调宽
### Changed
- 出站代理 Transport 补齐阶段超时(连接 30 秒 / TLS 握手 10 秒,对齐 SDK 直连模板):连不上的代理快速失败,不再拖满总超时
### Fixed
- 修复 Responses 直通调用 multi-agent / 搜索类模型必然超时:SDK 默认 `http.Client` 60 秒总超时覆盖到 body 读完,流式恰好 60 秒断流、非流式(响应头 >60 秒才返回)重试耗尽后约 121 秒报错。非流式改用预算总超时;流式去掉总超时,以定时取消模拟等待响应头预算,响应头到达后流时长不限、生命周期由客户端连接决定。附带消除此类超时对渠道熔断计数的误伤
## [0.7.1]
### Added
- 审计事件接口新增检索参数 `q`:服务端 `logContent` 全文粗筛 + 可见字段(事件名 / 资源 / 操作者 / IP / 请求路径等)精筛,不区分大小写、支持 `*` 通配;关键字内嵌续查游标,跨批过滤口径一致。全文粗筛不作最终判定——`logContent` 会命中隐藏认证元数据(如 `opc-principal` 头里的 `ttype: login`
- 审计批式响应新增 `scannedThrough`(已完整回溯到的时刻),供前端展示回溯进度
- Logging Search 服务配额为零的租户(报错含 `maxQueriesPerMinute: 0`,部分免费租户如此)自动回退 Audit API 小窗回溯:游标携带通道模式、续查不再试错,搜索降级为客户端可见字段匹配,审计页不再报错
### Changed
- 审计事件数据源由 Audit API 切换为 Logging Search`_Audit` 日志按 `datetime` 倒序,单页 200 条);空窗倍增上限由 30 天收紧到 14 天(单次查询时间窗硬限),单批新增约 20 秒时间预算,命中稀疏的深回溯拆成多个有界请求由前端接力
### Fixed
- 修复审计日志固定显示旧事件、刷新也看不到最新记录:Audit API 无排序参数且窗口内固定按处理时间正序,原实现凑满一批即返回,首批永远是 24h 窗口内最旧的一段,「向更早加载」实际在向更新方向翻页
## [0.7.0] ## [0.7.0]
### Added ### Added
+1 -1
View File
@@ -1 +1 @@
v0.7.0 v0.7.2
+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`)实测会被上游拒绝:
+21 -1
View File
@@ -860,7 +860,7 @@ const docTemplate = `{
"summary": "更新 AI 网关全局设置", "summary": "更新 AI 网关全局设置",
"parameters": [ "parameters": [
{ {
"description": "全量提交;保险丝阈值限 1..1024 KB", "description": "全量提交;保险丝阈值限 1..1024 KB,上游无响应预算限 30..900 秒",
"name": "body", "name": "body",
"in": "body", "in": "body",
"required": true, "required": true,
@@ -1572,6 +1572,12 @@ const docTemplate = `{
"description": "单批目标条数,缺省 100,上限 200", "description": "单批目标条数,缺省 100,上限 200",
"name": "limit", "name": "limit",
"in": "query" "in": "query"
},
{
"type": "string",
"description": "检索关键字(服务端全文匹配,支持 * 通配;仅首查生效)",
"name": "q",
"in": "query"
} }
], ],
"responses": { "responses": {
@@ -6096,6 +6102,10 @@ const docTemplate = `{
}, },
"streamGuardKB": { "streamGuardKB": {
"type": "integer" "type": "integer"
},
"upstreamWaitSeconds": {
"description": "UpstreamWaitSeconds 是 responses 直通的上游无响应预算(秒):非流式为单次\n尝试总超时,流式为等待响应头预算;multi-agent/搜索类模型需远超 60s",
"type": "integer"
} }
} }
}, },
@@ -8644,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"
}
} }
} }
}, },
@@ -9439,6 +9456,9 @@ const docTemplate = `{
"items": { "items": {
"$ref": "#/definitions/oci-portal_internal_oci.AuditEvent" "$ref": "#/definitions/oci-portal_internal_oci.AuditEvent"
} }
},
"scannedThrough": {
"type": "string"
} }
} }
}, },
+21 -1
View File
@@ -853,7 +853,7 @@
"summary": "更新 AI 网关全局设置", "summary": "更新 AI 网关全局设置",
"parameters": [ "parameters": [
{ {
"description": "全量提交;保险丝阈值限 1..1024 KB", "description": "全量提交;保险丝阈值限 1..1024 KB,上游无响应预算限 30..900 秒",
"name": "body", "name": "body",
"in": "body", "in": "body",
"required": true, "required": true,
@@ -1565,6 +1565,12 @@
"description": "单批目标条数,缺省 100,上限 200", "description": "单批目标条数,缺省 100,上限 200",
"name": "limit", "name": "limit",
"in": "query" "in": "query"
},
{
"type": "string",
"description": "检索关键字(服务端全文匹配,支持 * 通配;仅首查生效)",
"name": "q",
"in": "query"
} }
], ],
"responses": { "responses": {
@@ -6089,6 +6095,10 @@
}, },
"streamGuardKB": { "streamGuardKB": {
"type": "integer" "type": "integer"
},
"upstreamWaitSeconds": {
"description": "UpstreamWaitSeconds 是 responses 直通的上游无响应预算(秒):非流式为单次\n尝试总超时,流式为等待响应头预算;multi-agent/搜索类模型需远超 60s",
"type": "integer"
} }
} }
}, },
@@ -8637,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"
}
} }
} }
}, },
@@ -9432,6 +9449,9 @@
"items": { "items": {
"$ref": "#/definitions/oci-portal_internal_oci.AuditEvent" "$ref": "#/definitions/oci-portal_internal_oci.AuditEvent"
} }
},
"scannedThrough": {
"type": "string"
} }
} }
}, },
+19 -1
View File
@@ -65,6 +65,11 @@ definitions:
type: boolean type: boolean
streamGuardKB: streamGuardKB:
type: integer type: integer
upstreamWaitSeconds:
description: |-
UpstreamWaitSeconds 是 responses 直通的上游无响应预算(秒):非流式为单次
尝试总超时,流式为等待响应头预算;multi-agent/搜索类模型需远超 60s
type: integer
type: object type: object
internal_api.attachBootVolumeRequest: internal_api.attachBootVolumeRequest:
properties: properties:
@@ -1749,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:
@@ -2271,6 +2283,8 @@ definitions:
items: items:
$ref: '#/definitions/oci-portal_internal_oci.AuditEvent' $ref: '#/definitions/oci-portal_internal_oci.AuditEvent'
type: array type: array
scannedThrough:
type: string
type: object type: object
oci-portal_internal_service.Changes: oci-portal_internal_service.Changes:
additionalProperties: additionalProperties:
@@ -3202,7 +3216,7 @@ paths:
- AI 管理 - AI 管理
put: put:
parameters: parameters:
- description: 全量提交;保险丝阈值限 1..1024 KB - description: 全量提交;保险丝阈值限 1..1024 KB,上游无响应预算限 30..900 秒
in: body in: body
name: body name: body
required: true required: true
@@ -3637,6 +3651,10 @@ paths:
in: query in: query
name: limit name: limit
type: integer type: integer
- description: 检索关键字(服务端全文匹配,支持 * 通配;仅首查生效)
in: query
name: q
type: string
responses: responses:
"200": "200":
description: OK description: OK
+19 -6
View File
@@ -3,6 +3,7 @@ package api
import ( import (
"net/http" "net/http"
"strconv" "strconv"
"time"
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
@@ -450,6 +451,9 @@ type aiSettingsResponse struct {
// 请求 tools 已包含同名工具时不覆盖 // 请求 tools 已包含同名工具时不覆盖
GrokWebSearch bool `json:"grokWebSearch"` GrokWebSearch bool `json:"grokWebSearch"`
GrokXSearch bool `json:"grokXSearch"` GrokXSearch bool `json:"grokXSearch"`
// UpstreamWaitSeconds 是 responses 直通的上游无响应预算(秒):非流式为单次
// 尝试总超时,流式为等待响应头预算;multi-agent/搜索类模型需远超 60s
UpstreamWaitSeconds int `json:"upstreamWaitSeconds"`
} }
// currentAiSettings 汇总网关运行时设置为响应体。 // currentAiSettings 汇总网关运行时设置为响应体。
@@ -457,11 +461,12 @@ func (h *aiAdminHandler) currentAiSettings() aiSettingsResponse {
guardOn, guardKB := h.gw.StreamGuard() guardOn, guardKB := h.gw.StreamGuard()
web, x := h.gw.GrokSearch() web, x := h.gw.GrokSearch()
return aiSettingsResponse{ return aiSettingsResponse{
FilterDeprecated: h.gw.FilterDeprecated(), FilterDeprecated: h.gw.FilterDeprecated(),
StreamGuardEnabled: guardOn, StreamGuardEnabled: guardOn,
StreamGuardKB: guardKB, StreamGuardKB: guardKB,
GrokWebSearch: web, GrokWebSearch: web,
GrokXSearch: x, GrokXSearch: x,
UpstreamWaitSeconds: int(h.gw.UpstreamWait() / time.Second),
} }
} }
@@ -480,7 +485,7 @@ func (h *aiAdminHandler) aiSettings(c *gin.Context) {
// //
// @Summary 更新 AI 网关全局设置 // @Summary 更新 AI 网关全局设置
// @Tags AI 管理 // @Tags AI 管理
// @Param body body aiSettingsResponse true "全量提交;保险丝阈值限 1..1024 KB" // @Param body body aiSettingsResponse true "全量提交;保险丝阈值限 1..1024 KB,上游无响应预算限 30..900 秒"
// @Success 200 {object} aiSettingsResponse // @Success 200 {object} aiSettingsResponse
// @Failure 400 {object} map[string]string // @Failure 400 {object} map[string]string
// @Security BearerAuth // @Security BearerAuth
@@ -495,6 +500,10 @@ func (h *aiAdminHandler) updateAiSettings(c *gin.Context) {
c.JSON(http.StatusBadRequest, gin.H{"error": "streamGuardKB 须在 1..1024"}) c.JSON(http.StatusBadRequest, gin.H{"error": "streamGuardKB 须在 1..1024"})
return return
} }
if req.UpstreamWaitSeconds < 30 || req.UpstreamWaitSeconds > 900 {
c.JSON(http.StatusBadRequest, gin.H{"error": "upstreamWaitSeconds 须在 30..900"})
return
}
ctx := c.Request.Context() ctx := c.Request.Context()
if err := h.gw.SetFilterDeprecated(ctx, req.FilterDeprecated); err != nil { if err := h.gw.SetFilterDeprecated(ctx, req.FilterDeprecated); err != nil {
respondError(c, err) respondError(c, err)
@@ -508,5 +517,9 @@ func (h *aiAdminHandler) updateAiSettings(c *gin.Context) {
respondError(c, err) respondError(c, err)
return return
} }
if err := h.gw.SetUpstreamWait(ctx, req.UpstreamWaitSeconds); err != nil {
respondError(c, err)
return
}
c.JSON(http.StatusOK, h.currentAiSettings()) c.JSON(http.StatusOK, h.currentAiSettings())
} }
+4 -2
View File
@@ -68,13 +68,15 @@ func (h *ociConfigHandler) costs(c *gin.Context) {
// ---- 租户审计日志 ---- // ---- 租户审计日志 ----
// getAuditEvents 批式懒加载查询审计事件:cursor 为空自当前时刻首查, // getAuditEvents 批式懒加载查询审计事件:cursor 为空自当前时刻首查,
// 非空从上次响应游标继续向更早回溯;limit 单批目标条数(缺省 100,上限 200) // 非空从上次响应游标继续向更早回溯;limit 单批目标条数(缺省 100,上限 200);
// q 为服务端全文检索关键字,仅首查生效,续查沿用游标内嵌关键字。
// //
// @Summary 批式懒加载查询租户 OCI 审计事件 // @Summary 批式懒加载查询租户 OCI 审计事件
// @Tags 租户 IAM // @Tags 租户 IAM
// @Param id path int true "配置 ID" // @Param id path int true "配置 ID"
// @Param cursor query string false "续查游标(上次响应原样带回)" // @Param cursor query string false "续查游标(上次响应原样带回)"
// @Param limit query int false "单批目标条数,缺省 100,上限 200" // @Param limit query int false "单批目标条数,缺省 100,上限 200"
// @Param q query string false "检索关键字(服务端全文匹配,支持 * 通配;仅首查生效)"
// @Success 200 {object} service.AuditEventsView // @Success 200 {object} service.AuditEventsView
// @Security BearerAuth // @Security BearerAuth
// @Router /api/v1/oci-configs/{id}/audit-events [get] // @Router /api/v1/oci-configs/{id}/audit-events [get]
@@ -84,7 +86,7 @@ func (h *ociConfigHandler) getAuditEvents(c *gin.Context) {
return return
} }
limit, _ := strconv.Atoi(c.Query("limit")) limit, _ := strconv.Atoi(c.Query("limit"))
q := service.AuditQuery{Region: c.Query("region"), Cursor: c.Query("cursor"), Limit: limit} q := service.AuditQuery{Region: c.Query("region"), Cursor: c.Query("cursor"), Limit: limit, Q: c.Query("q")}
result, err := h.svc.AuditEvents(c.Request.Context(), id, q) result, err := h.svc.AuditEvents(c.Request.Context(), id, q)
if errors.Is(err, service.ErrInvalidAuditCursor) { if errors.Is(err, service.ErrInvalidAuditCursor) {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()}) c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
+469 -184
View File
@@ -6,19 +6,29 @@ import (
"fmt" "fmt"
"net" "net"
"sort" "sort"
"strings"
"time" "time"
"github.com/oracle/oci-go-sdk/v65/audit" "github.com/oracle/oci-go-sdk/v65/audit"
"github.com/oracle/oci-go-sdk/v65/common" "github.com/oracle/oci-go-sdk/v65/common"
"github.com/oracle/oci-go-sdk/v65/loggingsearch"
) )
// maxAuditPages 限制单次查询的翻页数:繁忙租户单日事件可上千, // maxAuditPages 限制单次查询的翻页数:每页最多 auditSearchPageLimit 条,
// 到限即返回 Truncated=true,由调用方收窄时间窗 // 到限即截断(窗口式回传 Truncated,批式留游标),由调用方续查
// 默认过滤(噪声事件/内网发起)后有效结果变少,页数放宽到 10 缓解截断。
const maxAuditPages = 10 const maxAuditPages = 10
// auditBatchTimeBudget 是批式查询的单批耗时预算:全文检索命中稀疏时
// 大窗扫描单页可达十余秒,超时即带游标返回,把长回溯拆成多个有界请求,
// 前端按已回溯位置展示进度并自动续查。
const auditBatchTimeBudget = 20 * time.Second
// auditSearchPageLimit 是 SearchLogs 单页条数(API 上限 1000):批式查询
// 攒满目标条数(~100)即携整页返回,页取 200 兼顾单页凑满一批与响应体量。
const auditSearchPageLimit = 200
// AuditEvent 是审计事件的列表精简视图;EventId 为 CloudEvents 全局唯一 id, // AuditEvent 是审计事件的列表精简视图;EventId 为 CloudEvents 全局唯一 id,
// 详情反查的键。Raw 为 SDK 原始事件的 JSON 序列化,由 service 层剥离进缓存, // 详情反查的键。Raw 为 _Audit 日志 logContent 原文,由 service 层剥离进缓存,
// 列表响应不再携带(详情接口按 eventId 取回)。 // 列表响应不再携带(详情接口按 eventId 取回)。
type AuditEvent struct { type AuditEvent struct {
EventId string `json:"eventId"` EventId string `json:"eventId"`
@@ -35,7 +45,7 @@ type AuditEvent struct {
Raw json.RawMessage `json:"raw,omitempty"` Raw json.RawMessage `json:"raw,omitempty"`
} }
// AuditEventsResult 是一次审计查询的结果;Truncated 表示翻页到限被截断, // AuditEventsResult 是一次窗口式审计查询的结果;Truncated 表示翻页到限被截断,
// 此时 NextPage 携带 opc-next-page 游标,同一时间窗回传可断点续翻。 // 此时 NextPage 携带 opc-next-page 游标,同一时间窗回传可断点续翻。
type AuditEventsResult struct { type AuditEventsResult struct {
Items []AuditEvent `json:"items"` Items []AuditEvent `json:"items"`
@@ -43,6 +53,324 @@ type AuditEventsResult struct {
NextPage string `json:"nextPage,omitempty"` NextPage string `json:"nextPage,omitempty"`
} }
// auditSearchClient 构造区域化的日志搜索客户端。审计数据源为 Logging Search
// 的 _Audit 日志:Audit API 无排序参数、窗口内固定按处理时间正序,首批只能
// 拿到窗口内最旧的一段;Logging Search 支持 datetime 倒序,才能从最新回溯。
func (c *RealClient) auditSearchClient(cred Credentials, region string) (loggingsearch.LogSearchClient, error) {
sc, err := loggingsearch.NewLogSearchClientWithConfigurationProvider(provider(cred))
if err != nil {
return sc, fmt.Errorf("new logging search client: %w", err)
}
applyProxy(&sc.BaseClient, cred)
if region != "" {
sc.SetRegion(normalizeRegion(region))
}
return sc, nil
}
// auditSearchQuery 组装租户根 compartment 审计日志的倒序检索语句;
// SummarizeMetricsData 遥测噪声占比高,服务端先滤一道减少无效翻页。
// q 非空时追加 logContent 全文包含匹配——它扫的是整条 JSON 的所有值,只当
// 粗筛;可见字段的精筛由 filterAuditTerm 兜底,避免隐藏元数据误命中。
func auditSearchQuery(tenancyOCID, q string) string {
query := fmt.Sprintf("search %q | where data.eventName != 'SummarizeMetricsData'", tenancyOCID+"/_Audit")
if term := SanitizeAuditTerm(q); term != "" {
query += fmt.Sprintf(" and logContent = '*%s*'", term)
}
return query + " | sort by datetime desc"
}
// auditTermMaxLen 限制检索关键字长度,防止游标与查询语句被撑爆。
const auditTermMaxLen = 100
// SanitizeAuditTerm 归一检索关键字:去除引号/反斜杠/控制字符防语句注入
// (查询目标已锁定本租户 _Audit 流,注入最坏只是语法错),截断超长输入;
// 保留 * 供用户通配。返回空串表示不追加过滤子句。
// service 构造首查游标与本包组装语句共用,对篡改游标二次消毒兜底。
func SanitizeAuditTerm(q string) string {
out := make([]rune, 0, len(q))
for _, r := range q {
if r == '\'' || r == '"' || r == '\\' || r < 0x20 {
continue
}
out = append(out, r)
if len(out) >= auditTermMaxLen {
break
}
}
return strings.TrimSpace(string(out))
}
// ListAuditEvents 实现 Client:实时查询租户根 compartment 在 [start, end) 内的
// 审计事件,最多翻 maxAuditPages 页,结果按发生时间倒序;纯读不落库。
// page 非空时从该游标断点续翻(必须配同一时间窗);到限截断时回传 NextPage。
// Search 配额为零的租户自动回退 Audit API 重查同一窗口(page 跨通道失效,重头吃窗)。
func (c *RealClient) ListAuditEvents(ctx context.Context, cred Credentials, region string, start, end time.Time, page string) (AuditEventsResult, error) {
f := c.newAuditFetchers(cred, region)
res, err := listAuditWindow(ctx, f.search, AuditCursor{Start: start, End: end, Page: page})
if err != nil && isSearchQuotaZero(err) {
res, err = listAuditWindow(ctx, f.audit, AuditCursor{Start: start, End: end})
}
return res, err
}
// listAuditWindow 用给定取页函数吃一个固定时间窗,最多 maxAuditPages 页;
// 页预算耗尽即截断,NextPage 携带未消费的窗内游标。
func listAuditWindow(ctx context.Context, fetch auditPageFetch, cur AuditCursor) (AuditEventsResult, error) {
result := AuditEventsResult{Items: []AuditEvent{}}
for i := 0; i < maxAuditPages; i++ {
items, next, err := fetch(ctx, cur)
if err != nil {
return AuditEventsResult{}, err
}
result.Items = appendKeptAuditEvents(result.Items, items)
if next == "" {
sortAuditEvents(result.Items)
return result, nil
}
cur.Page = next
}
result.Truncated = true
result.NextPage = cur.Page
sortAuditEvents(result.Items)
return result, nil
}
// auditPageFetch 拉取游标位置的一页已映射事件,返回窗内下一页游标。
type auditPageFetch func(ctx context.Context, cur AuditCursor) ([]AuditEvent, string, error)
// auditFetchers 汇集两条数据通道:search 为 Logging Search 倒序主路,
// audit 为 Search 配额为零租户的 Audit API 回退路。
type auditFetchers struct {
search auditPageFetch
audit auditPageFetch
}
// newAuditFetchers 构造两条通道的取页闭包;客户端惰性初始化,
// 各模式的续查不会白建用不到的客户端。
func (c *RealClient) newAuditFetchers(cred Credentials, region string) auditFetchers {
return auditFetchers{search: c.searchFetcher(cred, region), audit: c.auditAPIFetcher(cred, region)}
}
// searchFetcher 构造 Logging Search 通道的取页闭包。
func (c *RealClient) searchFetcher(cred Credentials, region string) auditPageFetch {
var sc *loggingsearch.LogSearchClient
return func(ctx context.Context, cur AuditCursor) ([]AuditEvent, string, error) {
if sc == nil {
cli, err := c.auditSearchClient(cred, region)
if err != nil {
return nil, "", err
}
sc = &cli
}
return searchAuditPage(ctx, *sc, cred.TenancyOCID, cur)
}
}
// auditAPIFetcher 构造 Audit API 回退通道的取页闭包。
func (c *RealClient) auditAPIFetcher(cred Credentials, region string) auditPageFetch {
var ac *audit.AuditClient
return func(ctx context.Context, cur AuditCursor) ([]AuditEvent, string, error) {
if ac == nil {
cli, err := c.auditClient(cred, region)
if err != nil {
return nil, "", err
}
ac = &cli
}
return listAuditPage(ctx, *ac, cred.TenancyOCID, cur)
}
}
// isSearchQuotaZero 识别「租户 Logging Search 配额为零」的失败:此类租户该
// 服务永久不可用(maxQueriesPerMinute/maxConcurrentQueries 均为 0),应回退
// Audit API;普通限流(配额非零)不回退,避免数据通道来回切换。
func isSearchQuotaZero(err error) bool {
if err == nil {
return false
}
msg := strings.ToLower(strings.ReplaceAll(err.Error(), " ", ""))
return strings.Contains(msg, "ratelimitexceeded") && strings.Contains(msg, "maxqueriesperminute:0,")
}
// appendKeptAuditEvents 过滤噪声后追加一页已映射事件;窗口式与批式查询共用。
func appendKeptAuditEvents(dst []AuditEvent, items []AuditEvent) []AuditEvent {
for _, ev := range items {
if keepAuditEvent(ev) {
dst = append(dst, ev)
}
}
return dst
}
// ---- 批式懒加载查询:分窗回溯 + 游标续查 ----
// 批式查询参数:单批翻页预算沿用 maxAuditPages;首窗 24h,连续空窗倍增
// 加速跨越闲置期,上限 14 天(Logging Search 单次查询时间窗硬限);
// 回溯下限为审计事件保留期 365 天。
const (
auditWindowHours = 24
auditWindowMaxHours = 336
auditRetentionDays = 365
)
// auditModeFallback 标记游标处于 Audit API 回退模式:部分租户的
// Logging Search 服务配额为零(maxQueriesPerMinute: 0),永久不可用。
const auditModeFallback = "a"
// auditFallbackWindowHours 是回退模式的基准窗宽:Audit API 窗口内固定按
// 处理时间正序且无排序参数,只能小窗回溯 + 前端全局重排保住从新到旧的体验。
const auditFallbackWindowHours = 1
// AuditCursor 是批式查询的续查位置:当前时间窗、窗内翻页游标、当前窗宽
// (小时,空窗倍增的记忆)、检索关键字(随游标续查,保证跨批过滤一致)
// 与数据通道模式(空为 Search 主路,"a" 为 Audit API 回退,续查沿用不再试错)。
// 序列化为不透明 cursor 由 service 层负责。
type AuditCursor struct {
Start time.Time `json:"s"`
End time.Time `json:"e"`
Page string `json:"p,omitempty"`
WindowHours int `json:"w"`
Q string `json:"q,omitempty"`
M string `json:"m,omitempty"`
}
// toFallback 把游标切到 Audit API 回退模式:Search 页游标跨通道失效须清空;
// 首窗收窄到基准窗宽,避免大窗正序分页又回到「首批全是窗口内最旧事件」的老问题。
func (cur AuditCursor) toFallback() AuditCursor {
cur.M = auditModeFallback
cur.Page = ""
cur.WindowHours = auditFallbackWindowHours
if cur.End.Sub(cur.Start) > auditFallbackWindowHours*time.Hour {
cur.Start = cur.End.Add(-auditFallbackWindowHours * time.Hour)
}
return cur
}
// NewAuditCursor 构造首查游标:自 now 起回溯第一个 24h 窗。
func NewAuditCursor(now time.Time) AuditCursor {
end := now.UTC().Truncate(time.Minute)
return AuditCursor{Start: end.Add(-auditWindowHours * time.Hour), End: end, WindowHours: auditWindowHours}
}
// advance 推进到紧邻更早的窗;empty 表示刚结束的窗无保留事件,窗宽倍增,
// 否则重置为该模式基准窗宽。done 为 true 表示已越过保留期尽头。
func (cur AuditCursor) advance(now time.Time, empty bool) (AuditCursor, bool) {
base := auditWindowHours
if cur.M == auditModeFallback {
base = auditFallbackWindowHours
}
w := cur.WindowHours
if w <= 0 {
w = base
}
if empty {
if w *= 2; w > auditWindowMaxHours {
w = auditWindowMaxHours
}
} else {
w = base
}
end := cur.Start
if end.Before(now.UTC().AddDate(0, 0, -auditRetentionDays)) {
return cur, true
}
return AuditCursor{Start: end.Add(-time.Duration(w) * time.Hour), End: end, WindowHours: w, Q: cur.Q, M: cur.M}, false
}
// AuditBatchResult 是一批懒加载结果;Cursor 为 nil 且 Exhausted 为 true
// 表示已回溯到保留期尽头,无更早数据。
type AuditBatchResult struct {
Items []AuditEvent
Cursor *AuditCursor
Exhausted bool
}
// ListAuditEventsBatch 实现 Client:从 cur 位置向更早方向收集约 limit 条
// 保留事件;单批受页预算与时间预算双重约束,不足额也返回,由前端按需续查。
// 倒序返回下,窗口不重叠 + 窗内游标续翻保证跨批不重不漏。
func (c *RealClient) ListAuditEventsBatch(ctx context.Context, cred Credentials, region string, cur AuditCursor, limit int) (AuditBatchResult, error) {
return listAuditBatch(ctx, c.newAuditFetchers(cred, region), cur, limit)
}
// listAuditBatch 是批式回溯的通道无关内核,取页函数注入便于测试。
func listAuditBatch(ctx context.Context, f auditFetchers, cur AuditCursor, limit int) (AuditBatchResult, error) {
res := AuditBatchResult{Items: []AuditEvent{}}
windowHasKept := false
deadline := time.Now().Add(auditBatchTimeBudget)
for budget := maxAuditPages; budget > 0 && len(res.Items) < limit && time.Now().Before(deadline); budget-- {
items, next, nextCur, err := fetchAuditPage(ctx, f, cur)
if err != nil {
return AuditBatchResult{}, err
}
cur = nextCur
before := len(res.Items)
res.Items = appendKeptAuditEvents(res.Items, filterAuditTerm(items, cur))
windowHasKept = windowHasKept || len(res.Items) > before
if next != "" {
cur.Page = next
continue
}
adv, done := cur.advance(time.Now(), !windowHasKept)
if done {
res.Exhausted = true
sortAuditEvents(res.Items)
return res, nil
}
cur, windowHasKept = adv, false
}
sortAuditEvents(res.Items)
res.Cursor = &cur
return res, nil
}
// fetchAuditPage 按游标模式取一页;Search 主路报「配额为零」时切到回退游标
// 并立即用 Audit API 重试,后续批次凭游标模式直达回退通道不再试错。
func fetchAuditPage(ctx context.Context, f auditFetchers, cur AuditCursor) ([]AuditEvent, string, AuditCursor, error) {
if cur.M == auditModeFallback {
items, next, err := f.audit(ctx, cur)
return items, next, cur, err
}
items, next, err := f.search(ctx, cur)
if err != nil && isSearchQuotaZero(err) {
cur = cur.toFallback()
items, next, err = f.audit(ctx, cur)
}
return items, next, cur, err
}
// filterAuditTerm 关键字精筛:只认列表可见字段(matchesAuditTerm),两条通道
// 语义一致。Search 主路的 logContent 全文条件会命中隐藏认证元数据(如
// opc-principal 头里的 ttype:login),只作粗筛减少翻页,不作为最终判定。
func filterAuditTerm(items []AuditEvent, cur AuditCursor) []AuditEvent {
if cur.Q == "" {
return items
}
out := items[:0]
for _, ev := range items {
if matchesAuditTerm(ev, cur.Q) {
out = append(out, ev)
}
}
return out
}
// matchesAuditTerm 判断事件是否命中关键字:不区分大小写的包含匹配,
// * 作为通配分段、各段都出现即命中,近似 Search 通道的 logContent 语义。
func matchesAuditTerm(ev AuditEvent, q string) bool {
hay := strings.ToLower(strings.Join([]string{
ev.EventName, ev.Source, ev.ResourceName, ev.CompartmentName,
ev.PrincipalName, ev.IPAddress, ev.Status, ev.RequestAction, ev.RequestPath,
}, "\n"))
for _, part := range strings.Split(strings.ToLower(q), "*") {
if part != "" && !strings.Contains(hay, part) {
return false
}
}
return true
}
// auditClient 构造区域化的 Audit API 客户端(回退通道)。
func (c *RealClient) auditClient(cred Credentials, region string) (audit.AuditClient, error) { func (c *RealClient) auditClient(cred Credentials, region string) (audit.AuditClient, error) {
ac, err := audit.NewAuditClientWithConfigurationProvider(provider(cred)) ac, err := audit.NewAuditClientWithConfigurationProvider(provider(cred))
if err != nil { if err != nil {
@@ -55,153 +383,10 @@ func (c *RealClient) auditClient(cred Credentials, region string) (audit.AuditCl
return ac, nil return ac, nil
} }
// ListAuditEvents 实现 Client:实时查询租户根 compartment 在 [start, end) 内的 // listAuditPage 拉取窗口内一页 Audit API 原始事件并压平;该 API 窗口内固定
// 审计事件,最多翻 maxAuditPages 页,结果按发生时间倒序;纯读不落库。 // 正序且只接受分钟粒度(起止秒与毫秒必须为 0)。Raw 为 SDK 事件原文,
// page 非空时从该游标断点续翻(必须配同一时间窗);到限截断时回传 NextPage // 与 Search 通道的 logContent 形态不同,详情弹窗均按任意 JSON 渲染
func (c *RealClient) ListAuditEvents(ctx context.Context, cred Credentials, region string, start, end time.Time, page string) (AuditEventsResult, error) { func listAuditPage(ctx context.Context, ac audit.AuditClient, tenancyOCID string, cur AuditCursor) ([]AuditEvent, string, error) {
ac, err := c.auditClient(cred, region)
if err != nil {
return AuditEventsResult{}, err
}
// Audit API 只接受分钟粒度:起止时间的秒与毫秒必须为 0
req := audit.ListEventsRequest{
CompartmentId: &cred.TenancyOCID,
StartTime: &common.SDKTime{Time: start.UTC().Truncate(time.Minute)},
EndTime: &common.SDKTime{Time: end.UTC().Truncate(time.Minute)},
}
if page != "" {
req.Page = &page
}
result := AuditEventsResult{Items: []AuditEvent{}}
for i := 0; i < maxAuditPages; i++ {
resp, err := ac.ListEvents(ctx, req)
if err != nil {
return AuditEventsResult{}, fmt.Errorf("list audit events: %w", err)
}
appendAuditEvents(&result, resp.Items)
if resp.OpcNextPage == nil {
sortAuditEvents(result.Items)
return result, nil
}
req.Page = resp.OpcNextPage
}
result.Truncated = true
result.NextPage = deref(req.Page)
sortAuditEvents(result.Items)
return result, nil
}
// appendAuditEvents 过滤噪声后追加一页事件;原始事件只对保留条目序列化。
func appendAuditEvents(result *AuditEventsResult, items []audit.AuditEvent) {
result.Items = appendKeptAuditEvents(result.Items, items)
}
// appendKeptAuditEvents 是过滤追加的通用形态,窗口式与批式查询共用。
func appendKeptAuditEvents(dst []AuditEvent, items []audit.AuditEvent) []AuditEvent {
for _, ev := range items {
out := toAuditEvent(ev)
if !keepAuditEvent(out) {
continue
}
if raw, mErr := json.Marshal(ev); mErr == nil {
out.Raw = raw
}
dst = append(dst, out)
}
return dst
}
// ---- 批式懒加载查询:分窗回溯 + 游标续查 ----
// 批式查询参数:单批 OCI 翻页预算沿用 maxAuditPages;首窗 24h,
// 连续空窗倍增(上限 30 天)加速跨越闲置期;回溯下限为事件保留期 365 天。
const (
auditWindowHours = 24
auditWindowMaxHours = 720
auditRetentionDays = 365
)
// AuditCursor 是批式查询的续查位置:当前时间窗、窗内 OCI 翻页游标
// 与当前窗宽(小时,空窗倍增的记忆)。序列化为不透明 cursor 由 service 层负责。
type AuditCursor struct {
Start time.Time `json:"s"`
End time.Time `json:"e"`
Page string `json:"p,omitempty"`
WindowHours int `json:"w"`
}
// NewAuditCursor 构造首查游标:自 now 起回溯第一个 24h 窗。
func NewAuditCursor(now time.Time) AuditCursor {
end := now.UTC().Truncate(time.Minute)
return AuditCursor{Start: end.Add(-auditWindowHours * time.Hour), End: end, WindowHours: auditWindowHours}
}
// advance 推进到紧邻更早的窗;empty 表示刚结束的窗无保留事件,窗宽倍增,
// 否则重置 24h。done 为 true 表示已越过保留期尽头。
func (cur AuditCursor) advance(now time.Time, empty bool) (AuditCursor, bool) {
w := cur.WindowHours
if w <= 0 {
w = auditWindowHours
}
if empty {
if w *= 2; w > auditWindowMaxHours {
w = auditWindowMaxHours
}
} else {
w = auditWindowHours
}
end := cur.Start
if end.Before(now.UTC().AddDate(0, 0, -auditRetentionDays)) {
return cur, true
}
return AuditCursor{Start: end.Add(-time.Duration(w) * time.Hour), End: end, WindowHours: w}, false
}
// AuditBatchResult 是一批懒加载结果;Cursor 为 nil 且 Exhausted 为 true
// 表示已回溯到保留期尽头,无更早数据。
type AuditBatchResult struct {
Items []AuditEvent
Cursor *AuditCursor
Exhausted bool
}
// ListAuditEventsBatch 实现 Client:从 cur 位置向更早方向收集约 limit 条
// 保留事件;单批最多消费 maxAuditPages 页 OCI 调用,不足额也按预算返回,
// 由前端按需续查。窗口不重叠 + 窗内游标续翻保证跨批不重不漏。
func (c *RealClient) ListAuditEventsBatch(ctx context.Context, cred Credentials, region string, cur AuditCursor, limit int) (AuditBatchResult, error) {
ac, err := c.auditClient(cred, region)
if err != nil {
return AuditBatchResult{}, err
}
res := AuditBatchResult{Items: []AuditEvent{}}
windowHasKept := false
for budget := maxAuditPages; budget > 0 && len(res.Items) < limit; budget-- {
items, next, err := listAuditPage(ctx, ac, cred.TenancyOCID, cur)
if err != nil {
return AuditBatchResult{}, err
}
before := len(res.Items)
res.Items = appendKeptAuditEvents(res.Items, items)
windowHasKept = windowHasKept || len(res.Items) > before
if next != "" {
cur.Page = next
continue
}
nextCur, done := cur.advance(time.Now(), !windowHasKept)
if done {
res.Exhausted = true
sortAuditEvents(res.Items)
return res, nil
}
cur, windowHasKept = nextCur, false
}
sortAuditEvents(res.Items)
res.Cursor = &cur
return res, nil
}
// listAuditPage 拉取当前游标位置的一页原始事件。
func listAuditPage(ctx context.Context, ac audit.AuditClient, tenancyOCID string, cur AuditCursor) ([]audit.AuditEvent, string, error) {
req := audit.ListEventsRequest{ req := audit.ListEventsRequest{
CompartmentId: &tenancyOCID, CompartmentId: &tenancyOCID,
StartTime: &common.SDKTime{Time: cur.Start.UTC().Truncate(time.Minute)}, StartTime: &common.SDKTime{Time: cur.Start.UTC().Truncate(time.Minute)},
@@ -214,42 +399,18 @@ func listAuditPage(ctx context.Context, ac audit.AuditClient, tenancyOCID string
if err != nil { if err != nil {
return nil, "", fmt.Errorf("list audit events: %w", err) return nil, "", fmt.Errorf("list audit events: %w", err)
} }
return resp.Items, deref(resp.OpcNextPage), nil items := make([]AuditEvent, 0, len(resp.Items))
} for _, ev := range resp.Items {
out := toAuditEvent(ev)
// auditInternalCIDRs 是 OCI 服务内部互调的发起方网段(RFC1918 + CGNAT)。 if raw, mErr := json.Marshal(ev); mErr == nil {
var auditInternalCIDRs = func() []*net.IPNet { out.Raw = raw
out := make([]*net.IPNet, 0, 4)
for _, cidr := range []string{"10.0.0.0/8", "172.16.0.0/12", "192.168.0.0/16", "100.64.0.0/10"} {
_, block, _ := net.ParseCIDR(cidr)
out = append(out, block)
}
return out
}()
// keepAuditEvent 保留有展示价值的事件:Audit API 无服务端过滤参数(仅时间窗),
// 在翻页循环内排除高频遥测噪声(SummarizeMetricsData)与内网地址发起的服务互调;
// 无 IP 的事件(控制面内部)保留。
func keepAuditEvent(ev AuditEvent) bool {
if ev.EventName == "SummarizeMetricsData" {
return false
}
if ev.IPAddress == "" {
return true
}
ip := net.ParseIP(ev.IPAddress)
if ip == nil {
return true
}
for _, block := range auditInternalCIDRs {
if block.Contains(ip) {
return false
} }
items = append(items, out)
} }
return true return items, deref(resp.OpcNextPage), nil
} }
// toAuditEvent 把 SDK 审计事件压平为列表 DTOSDK 字段全为指针逐层判 nil。 // toAuditEvent 把 Audit SDK 事件压平为列表 DTO;SDK 字段全为指针,逐层判 nil。
func toAuditEvent(ev audit.AuditEvent) AuditEvent { func toAuditEvent(ev audit.AuditEvent) AuditEvent {
out := AuditEvent{EventId: deref(ev.EventId), Source: deref(ev.Source)} out := AuditEvent{EventId: deref(ev.EventId), Source: deref(ev.Source)}
if ev.EventTime != nil { if ev.EventTime != nil {
@@ -275,6 +436,130 @@ func toAuditEvent(ev audit.AuditEvent) AuditEvent {
return out return out
} }
// searchAuditPage 拉取游标窗口内按 datetime 倒序的一页审计事件(已映射未过滤)。
func searchAuditPage(ctx context.Context, sc loggingsearch.LogSearchClient, tenancyOCID string, cur AuditCursor) ([]AuditEvent, string, error) {
req := loggingsearch.SearchLogsRequest{
SearchLogsDetails: loggingsearch.SearchLogsDetails{
TimeStart: &common.SDKTime{Time: cur.Start.UTC().Truncate(time.Minute)},
TimeEnd: &common.SDKTime{Time: cur.End.UTC().Truncate(time.Minute)},
SearchQuery: common.String(auditSearchQuery(tenancyOCID, cur.Q)),
},
Limit: common.Int(auditSearchPageLimit),
}
if cur.Page != "" {
req.Page = &cur.Page
}
resp, err := sc.SearchLogs(ctx, req)
if err != nil {
return nil, "", fmt.Errorf("search audit logs: %w", err)
}
items := make([]AuditEvent, 0, len(resp.Results))
for _, r := range resp.Results {
if ev, ok := toSearchAuditEvent(r); ok {
items = append(items, ev)
}
}
return items, deref(resp.OpcNextPage), nil
}
// auditInternalCIDRs 是 OCI 服务内部互调的发起方网段(RFC1918 + CGNAT)。
var auditInternalCIDRs = func() []*net.IPNet {
out := make([]*net.IPNet, 0, 4)
for _, cidr := range []string{"10.0.0.0/8", "172.16.0.0/12", "192.168.0.0/16", "100.64.0.0/10"} {
_, block, _ := net.ParseCIDR(cidr)
out = append(out, block)
}
return out
}()
// keepAuditEvent 保留有展示价值的事件:SummarizeMetricsData 已在检索语句里
// 先滤(此处兜底),内网地址发起的服务互调用 CIDR 判断(查询语言不便表达);
// 无 IP 的事件(控制面内部)保留。
func keepAuditEvent(ev AuditEvent) bool {
if ev.EventName == "SummarizeMetricsData" {
return false
}
if ev.IPAddress == "" {
return true
}
ip := net.ParseIP(ev.IPAddress)
if ip == nil {
return true
}
for _, block := range auditInternalCIDRs {
if block.Contains(ip) {
return false
}
}
return true
}
// searchAuditContent 是 _Audit 日志 logContent 的字段投影,只取列表展示所需;
// identity / request / response 可能为 null,零值即缺省。
type searchAuditContent struct {
ID string `json:"id"`
Time *time.Time `json:"time"`
Source string `json:"source"`
Data struct {
EventName string `json:"eventName"`
ResourceName string `json:"resourceName"`
CompartmentName string `json:"compartmentName"`
Identity struct {
PrincipalName string `json:"principalName"`
IPAddress string `json:"ipAddress"`
} `json:"identity"`
Request struct {
Action string `json:"action"`
Path string `json:"path"`
} `json:"request"`
Response struct {
Status string `json:"status"`
} `json:"response"`
} `json:"data"`
}
// toSearchAuditEvent 把日志搜索结果压平为列表 DTO;Raw 即 logContent 原文。
// 结构不符的条目丢弃(返回 false),不因单条脏数据整页失败。
func toSearchAuditEvent(r loggingsearch.SearchResult) (AuditEvent, bool) {
if r.Data == nil {
return AuditEvent{}, false
}
b, err := json.Marshal(r.Data)
if err != nil {
return AuditEvent{}, false
}
var hit struct {
LogContent json.RawMessage `json:"logContent"`
}
if err := json.Unmarshal(b, &hit); err != nil || len(hit.LogContent) == 0 {
return AuditEvent{}, false
}
var content searchAuditContent
if err := json.Unmarshal(hit.LogContent, &content); err != nil {
return AuditEvent{}, false
}
ev := searchContentToEvent(content)
ev.Raw = hit.LogContent
return ev, true
}
// searchContentToEvent 把投影字段填入列表 DTO;Raw 由调用方设置。
func searchContentToEvent(c searchAuditContent) AuditEvent {
return AuditEvent{
EventId: c.ID,
EventTime: c.Time,
EventName: c.Data.EventName,
Source: c.Source,
ResourceName: c.Data.ResourceName,
CompartmentName: c.Data.CompartmentName,
PrincipalName: c.Data.Identity.PrincipalName,
IPAddress: c.Data.Identity.IPAddress,
Status: c.Data.Response.Status,
RequestAction: c.Data.Request.Action,
RequestPath: c.Data.Request.Path,
}
}
// sortAuditEvents 按发生时间倒序排列;服务端返回顺序不保证,nil 时间排最后。 // sortAuditEvents 按发生时间倒序排列;服务端返回顺序不保证,nil 时间排最后。
func sortAuditEvents(items []AuditEvent) { func sortAuditEvents(items []AuditEvent) {
sort.SliceStable(items, func(i, j int) bool { sort.SliceStable(items, func(i, j int) bool {
+294 -27
View File
@@ -1,14 +1,262 @@
package oci package oci
import ( import (
"context"
"encoding/json"
"errors"
"fmt"
"reflect" "reflect"
"strings"
"testing" "testing"
"time" "time"
"github.com/oracle/oci-go-sdk/v65/audit" "github.com/oracle/oci-go-sdk/v65/audit"
"github.com/oracle/oci-go-sdk/v65/common" "github.com/oracle/oci-go-sdk/v65/common"
"github.com/oracle/oci-go-sdk/v65/loggingsearch"
) )
// quotaZeroErr 复刻 Search 配额为零租户的真实报错(SDK 解析错误体失败后带原文)。
var quotaZeroErr = errors.New(`search audit logs: Failed to parse json from response body due to: json: cannot unmarshal number into Go struct field servicefailure.code of type string. With response body { "code" : 500, "message" : "Rate limit exceeded for ocid: ocid1.tenancy..x, maxQueriesPerMinute: 0, maxConcurrentQueries: 0" }.`)
func TestIsSearchQuotaZero(t *testing.T) {
cases := []struct {
name string
err error
want bool
}{
{"配额为零真实报错", quotaZeroErr, true},
{"普通限流不回退", errors.New(`Rate limit exceeded for ocid: x, maxQueriesPerMinute: 60, maxConcurrentQueries: 2`), false},
{"其他错误", errors.New("service unavailable"), false},
{"nil", nil, false},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if got := isSearchQuotaZero(tc.err); got != tc.want {
t.Fatalf("isSearchQuotaZero() = %v, want %v", got, tc.want)
}
})
}
}
func TestListAuditBatchFallback(t *testing.T) {
et := time.Now().UTC().Add(-10 * time.Minute)
searchCalls, auditCalls := 0, 0
f := auditFetchers{
search: func(context.Context, AuditCursor) ([]AuditEvent, string, error) {
searchCalls++
return nil, "", quotaZeroErr
},
audit: func(_ context.Context, cur AuditCursor) ([]AuditEvent, string, error) {
auditCalls++
if cur.M != auditModeFallback || cur.Page != "" {
t.Fatalf("回退通道应携带模式标记且清空页游标, got %+v", cur)
}
ev := AuditEvent{EventId: fmt.Sprint(auditCalls), EventName: "GetInstance", EventTime: &et}
return []AuditEvent{ev}, "", nil
},
}
res, err := listAuditBatch(context.Background(), f, NewAuditCursor(time.Now()), 3)
if err != nil {
t.Fatalf("配额为零应回退成功, got %v", err)
}
if searchCalls != 1 {
t.Fatalf("Search 只应试错一次, got %d", searchCalls)
}
if len(res.Items) < 3 || auditCalls < 3 {
t.Fatalf("回退后应继续凑批, items=%d auditCalls=%d", len(res.Items), auditCalls)
}
if res.Cursor == nil || res.Cursor.M != auditModeFallback || res.Cursor.WindowHours != auditFallbackWindowHours {
t.Fatalf("续查游标应保持回退模式与基准窗宽, got %+v", res.Cursor)
}
}
func TestListAuditBatchFallbackCursorSkipsSearch(t *testing.T) {
et := time.Now().UTC().Add(-10 * time.Minute)
f := auditFetchers{
search: func(context.Context, AuditCursor) ([]AuditEvent, string, error) {
t.Fatal("回退模式游标不应再调用 Search 通道")
return nil, "", nil
},
audit: func(context.Context, AuditCursor) ([]AuditEvent, string, error) {
return []AuditEvent{{EventId: "e1", EventName: "GetVcn", EventTime: &et}}, "", nil
},
}
cur := NewAuditCursor(time.Now()).toFallback()
if _, err := listAuditBatch(context.Background(), f, cur, 1); err != nil {
t.Fatalf("回退模式续查失败: %v", err)
}
}
func TestListAuditBatchSearchErrorNoFallback(t *testing.T) {
f := auditFetchers{
search: func(context.Context, AuditCursor) ([]AuditEvent, string, error) {
return nil, "", errors.New("search audit logs: timeout")
},
audit: func(context.Context, AuditCursor) ([]AuditEvent, string, error) {
t.Fatal("普通错误不应触发回退")
return nil, "", nil
},
}
if _, err := listAuditBatch(context.Background(), f, NewAuditCursor(time.Now()), 1); err == nil {
t.Fatal("普通错误应原样上抛")
}
}
func TestFilterAuditTerm(t *testing.T) {
login := AuditEvent{EventId: "e1", EventName: "InteractiveLogin"}
noise := AuditEvent{EventId: "e2", EventName: "ListRecommendations"}
items := []AuditEvent{login, noise}
cases := []struct {
name string
cur AuditCursor
want int
}{
{"无关键字原样放行", AuditCursor{}, 2},
{"Search 主路也精筛可见字段", AuditCursor{Q: "login"}, 1},
{"回退模式精筛", AuditCursor{Q: "login", M: auditModeFallback}, 1},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if got := filterAuditTerm(append([]AuditEvent{}, items...), tc.cur); len(got) != tc.want {
t.Fatalf("filterAuditTerm() 保留 %d 条, want %d", len(got), tc.want)
}
})
}
}
func TestListAuditBatchSearchTermPrecision(t *testing.T) {
et := time.Now().UTC().Add(-10 * time.Minute)
// 模拟 Search 主路粗筛后仍混入的隐藏元数据误命中(如 ttype:login)
f := auditFetchers{
search: func(_ context.Context, cur AuditCursor) ([]AuditEvent, string, error) {
return []AuditEvent{
{EventId: "hit", EventName: "InteractiveLogin", EventTime: &et},
{EventId: "noise1", EventName: "ListRecommendations", EventTime: &et},
{EventId: "noise2", EventName: "SearchLogs", EventTime: &et},
}, "", nil
},
audit: func(context.Context, AuditCursor) ([]AuditEvent, string, error) {
t.Fatal("Search 正常时不应走回退")
return nil, "", nil
},
}
cur := NewAuditCursor(time.Now())
cur.Q = "login"
res, err := listAuditBatch(context.Background(), f, cur, 1)
if err != nil {
t.Fatalf("listAuditBatch() err = %v", err)
}
if len(res.Items) != 1 || res.Items[0].EventId != "hit" {
t.Fatalf("应只保留可见字段命中的事件, got %+v", res.Items)
}
}
func TestMatchesAuditTerm(t *testing.T) {
ev := AuditEvent{EventName: "ListVnicAttachments", ResourceName: "web-1", PrincipalName: "Vivien", IPAddress: "1.2.3.4"}
cases := []struct {
name string
q string
want bool
}{
{"不区分大小写", "listvnic", true},
{"通配分段都出现", "List*Attachments", true},
{"资源名命中", "WEB-1", true},
{"未命中", "TerminateInstance", false},
{"通配缺段不命中", "List*Volume", false},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if got := matchesAuditTerm(ev, tc.q); got != tc.want {
t.Fatalf("matchesAuditTerm(%q) = %v, want %v", tc.q, got, tc.want)
}
})
}
}
// searchResultFromJSON 把 JSON 文本构造成 SearchLogs 单条结果(Data 为 interface{})。
func searchResultFromJSON(t *testing.T, s string) loggingsearch.SearchResult {
t.Helper()
var v interface{}
if err := json.Unmarshal([]byte(s), &v); err != nil {
t.Fatalf("fixture 不是合法 JSON: %v", err)
}
return loggingsearch.SearchResult{Data: &v}
}
func TestToSearchAuditEvent(t *testing.T) {
eventTime := time.Date(2026, 7, 6, 10, 30, 0, 0, time.UTC)
tests := []struct {
name string
data string
wantOK bool
want AuditEvent
}{
{
name: "全字段齐全",
data: `{"datetime":1783074600000,"logContent":{
"id":"evt-abc","time":"2026-07-06T10:30:00Z","source":"ComputeApi",
"data":{"eventName":"TerminateInstance","resourceName":"web-1","compartmentName":"prod",
"identity":{"principalName":"api-admin","ipAddress":"1.2.3.4"},
"request":{"action":"DELETE","path":"/20160918/instances/ocid1..."},
"response":{"status":"204"}}}}`,
wantOK: true,
want: AuditEvent{
EventId: "evt-abc",
EventTime: &eventTime,
EventName: "TerminateInstance",
Source: "ComputeApi",
ResourceName: "web-1",
CompartmentName: "prod",
PrincipalName: "api-admin",
IPAddress: "1.2.3.4",
Status: "204",
RequestAction: "DELETE",
RequestPath: "/20160918/instances/ocid1...",
},
},
{
name: "identity/request/response 为 null 时只保留信封字段",
data: `{"logContent":{"id":"evt-x","time":"2026-07-06T10:30:00Z","source":"VcnApi",
"data":{"eventName":"GetVcn","identity":null,"request":null,"response":null}}}`,
wantOK: true,
want: AuditEvent{EventId: "evt-x", EventTime: &eventTime, Source: "VcnApi", EventName: "GetVcn"},
},
{
name: "缺 logContent 丢弃",
data: `{"datetime":1783074600000}`,
wantOK: false,
},
{
name: "logContent 结构不符丢弃",
data: `{"logContent":"plain-text"}`,
wantOK: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, ok := toSearchAuditEvent(searchResultFromJSON(t, tt.data))
if ok != tt.wantOK {
t.Fatalf("ok = %v, want %v", ok, tt.wantOK)
}
if !ok {
return
}
if len(got.Raw) == 0 {
t.Fatalf("Raw 应携带 logContent 原文")
}
if !auditEventEqual(got, tt.want) {
t.Errorf("toSearchAuditEvent() = %+v, want %+v", got, tt.want)
}
})
}
}
func TestToSearchAuditEventNilData(t *testing.T) {
if _, ok := toSearchAuditEvent(loggingsearch.SearchResult{}); ok {
t.Fatal("Data 为 nil 应丢弃")
}
}
func TestToAuditEvent(t *testing.T) { func TestToAuditEvent(t *testing.T) {
eventTime := time.Date(2026, 7, 6, 10, 30, 0, 0, time.UTC) eventTime := time.Date(2026, 7, 6, 10, 30, 0, 0, time.UTC)
tests := []struct { tests := []struct {
@@ -38,44 +286,23 @@ func TestToAuditEvent(t *testing.T) {
}, },
}, },
want: AuditEvent{ want: AuditEvent{
EventId: "evt-abc", EventId: "evt-abc", EventTime: &eventTime, EventName: "TerminateInstance",
EventTime: &eventTime, Source: "ComputeApi", ResourceName: "web-1", CompartmentName: "prod",
EventName: "TerminateInstance", PrincipalName: "api-admin", IPAddress: "1.2.3.4", Status: "204",
Source: "ComputeApi", RequestAction: "DELETE", RequestPath: "/20160918/instances/ocid1...",
ResourceName: "web-1",
CompartmentName: "prod",
PrincipalName: "api-admin",
IPAddress: "1.2.3.4",
Status: "204",
RequestAction: "DELETE",
RequestPath: "/20160918/instances/ocid1...",
}, },
}, },
{
name: "Data 为 nil 时只保留信封字段",
ev: audit.AuditEvent{
Source: common.String("VcnApi"),
EventTime: &common.SDKTime{Time: eventTime},
},
want: AuditEvent{EventTime: &eventTime, Source: "VcnApi"},
},
{ {
name: "嵌套局部 nil 各自安全跳过", name: "嵌套局部 nil 各自安全跳过",
ev: audit.AuditEvent{ ev: audit.AuditEvent{
Data: &audit.Data{ Data: &audit.Data{
EventName: common.String("GetInstance"), EventName: common.String("GetInstance"),
Identity: nil,
Request: &audit.Request{Path: common.String("/instances")}, Request: &audit.Request{Path: common.String("/instances")},
Response: nil,
}, },
}, },
want: AuditEvent{EventName: "GetInstance", RequestPath: "/instances"}, want: AuditEvent{EventName: "GetInstance", RequestPath: "/instances"},
}, },
{ {name: "空事件全部零值", ev: audit.AuditEvent{}, want: AuditEvent{}},
name: "空事件全部零值",
ev: audit.AuditEvent{},
want: AuditEvent{},
},
} }
for _, tt := range tests { for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) { t.Run(tt.name, func(t *testing.T) {
@@ -86,6 +313,37 @@ func TestToAuditEvent(t *testing.T) {
} }
} }
func TestAuditSearchQuery(t *testing.T) {
const prefix = `search "ocid1.tenancy.oc1..aaa/_Audit" | where data.eventName != 'SummarizeMetricsData'`
cases := []struct {
name string
term string
want string
}{
{"无关键字", "", prefix + ` | sort by datetime desc`},
{"带关键字追加全文匹配", "TerminateInstance", prefix + ` and logContent = '*TerminateInstance*' | sort by datetime desc`},
{"引号与反斜杠被消毒", `O'Brien\"x`, prefix + ` and logContent = '*OBrienx*' | sort by datetime desc`},
{"纯引号消毒后为空不追加", `'"`, prefix + ` | sort by datetime desc`},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
if got := auditSearchQuery("ocid1.tenancy.oc1..aaa", tc.term); got != tc.want {
t.Fatalf("auditSearchQuery() = %q, want %q", got, tc.want)
}
})
}
}
func TestSanitizeAuditTerm(t *testing.T) {
if got := SanitizeAuditTerm(" Get*Instance\t "); got != "Get*Instance" {
t.Fatalf("应保留 * 并去除首尾空白与控制字符, got %q", got)
}
long := strings.Repeat("a", 300)
if got := SanitizeAuditTerm(long); len(got) != auditTermMaxLen {
t.Fatalf("超长应截断到 %d, got %d", auditTermMaxLen, len(got))
}
}
// auditEventEqual 比较两个 DTOEventTime 按值比较,Raw 不参与,其余反射比较。 // auditEventEqual 比较两个 DTOEventTime 按值比较,Raw 不参与,其余反射比较。
func auditEventEqual(a, b AuditEvent) bool { func auditEventEqual(a, b AuditEvent) bool {
if (a.EventTime == nil) != (b.EventTime == nil) { if (a.EventTime == nil) != (b.EventTime == nil) {
@@ -149,6 +407,7 @@ func TestAuditCursorAdvance(t *testing.T) {
Start: now.Add(-24 * time.Hour), Start: now.Add(-24 * time.Hour),
End: now, End: now,
WindowHours: 24, WindowHours: 24,
Q: "kw",
} }
cases := []struct { cases := []struct {
name string name string
@@ -159,8 +418,10 @@ func TestAuditCursorAdvance(t *testing.T) {
}{ }{
{"有事件重置 24h 窗", AuditCursor{Start: base.Start, End: base.End, WindowHours: 96}, false, 24, false}, {"有事件重置 24h 窗", AuditCursor{Start: base.Start, End: base.End, WindowHours: 96}, false, 24, false},
{"空窗倍增", base, true, 48, false}, {"空窗倍增", base, true, 48, false},
{"倍增封顶 720h", AuditCursor{Start: base.Start, End: base.End, WindowHours: 512}, true, 720, false}, {"倍增封顶 336h(14 天查询窗硬限)", AuditCursor{Start: base.Start, End: base.End, WindowHours: 256}, true, 336, false},
{"窗宽缺省按 24h 起算", AuditCursor{Start: base.Start, End: base.End}, true, 48, false}, {"窗宽缺省按 24h 起算", AuditCursor{Start: base.Start, End: base.End}, true, 48, false},
{"回退模式有事件重置 1h 基准窗", AuditCursor{Start: base.Start, End: base.End, WindowHours: 8, M: auditModeFallback}, false, 1, false},
{"回退模式空窗照常倍增", AuditCursor{Start: base.Start, End: base.End, WindowHours: 1, M: auditModeFallback}, true, 2, false},
{"越过保留期即尽头", AuditCursor{Start: now.AddDate(0, 0, -366), End: now.AddDate(0, 0, -365), WindowHours: 24}, false, 0, true}, {"越过保留期即尽头", AuditCursor{Start: now.AddDate(0, 0, -366), End: now.AddDate(0, 0, -365), WindowHours: 24}, false, 0, true},
} }
for _, tc := range cases { for _, tc := range cases {
@@ -184,6 +445,12 @@ func TestAuditCursorAdvance(t *testing.T) {
if next.Page != "" { if next.Page != "" {
t.Fatalf("新窗应清空窗内游标, got %q", next.Page) t.Fatalf("新窗应清空窗内游标, got %q", next.Page)
} }
if next.Q != tc.cur.Q {
t.Fatalf("新窗应继承检索关键字, got %q want %q", next.Q, tc.cur.Q)
}
if next.M != tc.cur.M {
t.Fatalf("新窗应继承通道模式, got %q want %q", next.M, tc.cur.M)
}
}) })
} }
} }
+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) {
+6 -4
View File
@@ -102,10 +102,12 @@ type Client interface {
ListGenAiModels(ctx context.Context, cred Credentials, region string) ([]GenAiModel, error) ListGenAiModels(ctx context.Context, cred Credentials, region string) ([]GenAiModel, error)
GenAiProbeChat(ctx context.Context, cred Credentials, region, modelOcid, modelName string) (int, error) GenAiProbeChat(ctx context.Context, cred Credentials, region, modelOcid, modelName string) (int, error)
GenAiEmbed(ctx context.Context, cred Credentials, region, modelOcid string, inputs []string, dimensions *int) ([][]float32, *aiwire.Usage, error) GenAiEmbed(ctx context.Context, cred Credentials, region, modelOcid string, inputs []string, dimensions *int) ([][]float32, *aiwire.Usage, error)
// GenAiCompatResponses 直通 OpenAI Responses 请求体到 /actions/v1/responses(xAI 服务端工具通路) // GenAiCompatResponses 直通 OpenAI Responses 请求体到 /actions/v1/responses(xAI 服务端工具通路);
GenAiCompatResponses(ctx context.Context, cred Credentials, region string, body []byte) ([]byte, error) // wait 是上游无响应预算(整请求总超时),multi-agent/搜索类模型需远超 SDK 默认 60s。
// GenAiCompatResponsesStream 流式直通 /actions/v1/responses,建立成功返回 SSE body。 GenAiCompatResponses(ctx context.Context, cred Credentials, region string, body []byte, wait time.Duration) ([]byte, error)
GenAiCompatResponsesStream(ctx context.Context, cred Credentials, region string, body []byte) (io.ReadCloser, error) // GenAiCompatResponsesStream 流式直通 /actions/v1/responses,建立成功返回 SSE body;
// wait 仅约束等待响应头阶段,建立后的流生命周期由 ctx 决定。
GenAiCompatResponsesStream(ctx context.Context, cred Credentials, region string, body []byte, wait time.Duration) (io.ReadCloser, error)
// GenAiCompatSpeech 直通 OpenAI Audio Speech 请求体到 /openai/v1/audio/speech,返回音频与 Content-Type。 // GenAiCompatSpeech 直通 OpenAI Audio Speech 请求体到 /openai/v1/audio/speech,返回音频与 Content-Type。
GenAiCompatSpeech(ctx context.Context, cred Credentials, region string, body []byte) ([]byte, string, error) GenAiCompatSpeech(ctx context.Context, cred Credentials, region string, body []byte) ([]byte, string, error)
// GenAiRerank 文档重排,返回按相关度排序的下标与得分。 // GenAiRerank 文档重排,返回按相关度排序的下标与得分。
+2 -1
View File
@@ -158,7 +158,8 @@ func (c *RealClient) GenAiProbeChat(ctx context.Context, cred Credentials, regio
if err != nil { if err != nil {
return 0, err return 0, err
} }
if _, err = c.GenAiCompatResponses(ctx, cred, region, body); err == nil { // 探测追求快速失败,沿用 SDK 默认量级的 60s 预算即可
if _, err = c.GenAiCompatResponses(ctx, cred, region, body, 60*time.Second); err == nil {
return http.StatusOK, nil return http.StatusOK, nil
} }
if status, ok := ServiceStatus(err); ok { if status, ok := ServiceStatus(err); ok {
+74 -24
View File
@@ -6,6 +6,7 @@ import (
"fmt" "fmt"
"io" "io"
"net/http" "net/http"
"time"
"github.com/oracle/oci-go-sdk/v65/common" "github.com/oracle/oci-go-sdk/v65/common"
) )
@@ -13,18 +14,19 @@ import (
// compatResponsesLimit 限制直通响应体大小;web_search 输出含多段引用,给足余量。 // compatResponsesLimit 限制直通响应体大小;web_search 输出含多段引用,给足余量。
const compatResponsesLimit = int64(8 << 20) const compatResponsesLimit = int64(8 << 20)
// GenAiCompatResponses 实现 Client:把 OpenAI Responses 请求体直通到 OCI // dispatcherWithTimeout 把 dispatcher 换成指定总超时的拷贝(保留 Transport,
// `/20231130/actions/v1/responses`(IAM 签名)。xAI 服务端工具(web_search / // 代理链路不受影响);timeout=0 表示无总超时(流式读 body 不能有总时限)。
// x_search / code_interpreter)与 mcp 已被 Oracle 文档正式支持,工具参数与限制 // 非 *http.Client 的自定义 dispatcher 保持原样,维持既有超时行为。
// 遵循 xAI 规格;调用方须自行校验并改写请求体(store/stream)。 func dispatcherWithTimeout(d common.HTTPRequestDispatcher, timeout time.Duration) common.HTTPRequestDispatcher {
func (c *RealClient) GenAiCompatResponses(ctx context.Context, cred Credentials, region string, body []byte) ([]byte, error) { hc, ok := d.(*http.Client)
ic, err := c.genAiInferenceClient(cred, region) if !ok {
if err != nil { return d
return nil, err
} }
client := ic.BaseClient return &http.Client{Transport: hc.Transport, Timeout: timeout}
common.UpdateEndpointTemplateForOptions(&client) }
common.SetMissingTemplateParams(&client)
// newCompatResponsesRequest 构造 /actions/v1/responses 直通请求。
func newCompatResponsesRequest(ctx context.Context, cred Credentials, body []byte) (*http.Request, error) {
request, err := http.NewRequestWithContext(ctx, http.MethodPost, "/actions/v1/responses", bytes.NewReader(body)) request, err := http.NewRequestWithContext(ctx, http.MethodPost, "/actions/v1/responses", bytes.NewReader(body))
if err != nil { if err != nil {
return nil, fmt.Errorf("build compat responses request: %w", err) return nil, fmt.Errorf("build compat responses request: %w", err)
@@ -32,6 +34,27 @@ func (c *RealClient) GenAiCompatResponses(ctx context.Context, cred Credentials,
request.Header.Set("Content-Type", "application/json") request.Header.Set("Content-Type", "application/json")
request.Header.Set("CompartmentId", cred.TenancyOCID) request.Header.Set("CompartmentId", cred.TenancyOCID)
request.Header.Set("opc-compartment-id", cred.TenancyOCID) request.Header.Set("opc-compartment-id", cred.TenancyOCID)
return request, nil
}
// GenAiCompatResponses 实现 Client:把 OpenAI Responses 请求体直通到 OCI
// `/20231130/actions/v1/responses`(IAM 签名)。xAI 服务端工具(web_search /
// x_search / code_interpreter)与 mcp 已被 Oracle 文档正式支持,工具参数与限制
// 遵循 xAI 规格;调用方须自行校验并改写请求体(store/stream)。
// wait 为整请求总超时:非流式上游要等全部生成完才回响应头,SDK 默认 60s 会掐断慢模型。
func (c *RealClient) GenAiCompatResponses(ctx context.Context, cred Credentials, region string, body []byte, wait time.Duration) ([]byte, error) {
ic, err := c.genAiInferenceClient(cred, region)
if err != nil {
return nil, err
}
client := ic.BaseClient
common.UpdateEndpointTemplateForOptions(&client)
common.SetMissingTemplateParams(&client)
client.HTTPClient = dispatcherWithTimeout(client.HTTPClient, wait)
request, err := newCompatResponsesRequest(ctx, cred, body)
if err != nil {
return nil, err
}
response, err := client.Call(ctx, request) response, err := client.Call(ctx, request)
if err != nil { if err != nil {
return nil, err return nil, err
@@ -44,10 +67,46 @@ func (c *RealClient) GenAiCompatResponses(ctx context.Context, cred Credentials,
return payload, nil return payload, nil
} }
// cancelReadCloser 在流关闭时同步取消建立阶段派生的 ctx,避免其随流生命周期泄漏。
type cancelReadCloser struct {
io.ReadCloser
cancel context.CancelFunc
}
func (c *cancelReadCloser) Close() error {
c.cancel()
return c.ReadCloser.Close()
}
// httpCaller 抽象 BaseClient.Call,便于对预算逻辑做无签名单测。
type httpCaller interface {
Call(ctx context.Context, request *http.Request) (*http.Response, error)
}
// callWithHeaderBudget 以 wait 为等待响应头预算执行调用:预算内未返回则取消
// 请求(SDK Call 会把 ctx 重绑到请求);响应头到达即解除预算,之后流的生命
// 周期由 ctx 决定,返回的流 Close 时同步取消派生 ctx。
func callWithHeaderBudget(ctx context.Context, c httpCaller, req *http.Request, wait time.Duration) (io.ReadCloser, error) {
callCtx, cancel := context.WithCancel(ctx)
timer := time.AfterFunc(wait, cancel)
response, err := c.Call(callCtx, req)
timer.Stop()
if err != nil {
if response != nil && response.Body != nil {
response.Body.Close()
}
cancel()
return nil, err
}
return &cancelReadCloser{ReadCloser: response.Body, cancel: cancel}, nil
}
// GenAiCompatResponsesStream 实现 Client:以流式直通 OCI `/actions/v1/responses`, // GenAiCompatResponsesStream 实现 Client:以流式直通 OCI `/actions/v1/responses`,
// 建立成功(2xx)返回 SSE body(调用方负责 Close);建立失败返回 SDK ServiceError, // 建立成功(2xx)返回 SSE body(调用方负责 Close);建立失败返回 SDK ServiceError,
// 与既有渠道切换/熔断错误分类兼容。请求体须由调用方置 stream:true。 // 与既有渠道切换/熔断错误分类兼容。请求体须由调用方置 stream:true。
func (c *RealClient) GenAiCompatResponsesStream(ctx context.Context, cred Credentials, region string, body []byte) (io.ReadCloser, error) { // 总超时置 0(SSE 读 body 不能有总时限);wait 以定时取消模拟等待响应头预算,
// 响应头到达即解除,此后流的生命周期完全由 ctx(下游客户端断开)决定。
func (c *RealClient) GenAiCompatResponsesStream(ctx context.Context, cred Credentials, region string, body []byte, wait time.Duration) (io.ReadCloser, error) {
ic, err := c.genAiInferenceClient(cred, region) ic, err := c.genAiInferenceClient(cred, region)
if err != nil { if err != nil {
return nil, err return nil, err
@@ -55,19 +114,10 @@ func (c *RealClient) GenAiCompatResponsesStream(ctx context.Context, cred Creden
client := ic.BaseClient client := ic.BaseClient
common.UpdateEndpointTemplateForOptions(&client) common.UpdateEndpointTemplateForOptions(&client)
common.SetMissingTemplateParams(&client) common.SetMissingTemplateParams(&client)
request, err := http.NewRequestWithContext(ctx, http.MethodPost, "/actions/v1/responses", bytes.NewReader(body)) client.HTTPClient = dispatcherWithTimeout(client.HTTPClient, 0)
request, err := newCompatResponsesRequest(ctx, cred, body)
if err != nil { if err != nil {
return nil, fmt.Errorf("build compat responses stream request: %w", err)
}
request.Header.Set("Content-Type", "application/json")
request.Header.Set("CompartmentId", cred.TenancyOCID)
request.Header.Set("opc-compartment-id", cred.TenancyOCID)
response, err := client.Call(ctx, request)
if err != nil {
if response != nil && response.Body != nil {
response.Body.Close()
}
return nil, err return nil, err
} }
return response.Body, nil return callWithHeaderBudget(ctx, client, request, wait)
} }
+141
View File
@@ -0,0 +1,141 @@
package oci
import (
"context"
"errors"
"io"
"net/http"
"strings"
"testing"
"time"
"github.com/oracle/oci-go-sdk/v65/common"
)
// staticDispatcher 是非 *http.Client 的自定义 dispatcher,用于降级分支。
type staticDispatcher struct{}
func (staticDispatcher) Do(*http.Request) (*http.Response, error) { return nil, nil }
func TestDispatcherWithTimeout(t *testing.T) {
tr := &http.Transport{}
tests := []struct {
name string
in common.HTTPRequestDispatcher
timeout time.Duration
check func(t *testing.T, out common.HTTPRequestDispatcher)
}{
{
name: "http.Client 换总超时并保留 Transport", in: &http.Client{Transport: tr, Timeout: 60 * time.Second},
timeout: 300 * time.Second,
check: func(t *testing.T, out common.HTTPRequestDispatcher) {
hc, ok := out.(*http.Client)
if !ok || hc.Timeout != 300*time.Second || hc.Transport != tr {
t.Fatalf("期望拷贝 client 且 Timeout=300s、Transport 保留, 得到 %#v", out)
}
},
},
{
name: "timeout=0 表示无总超时", in: &http.Client{Timeout: 60 * time.Second}, timeout: 0,
check: func(t *testing.T, out common.HTTPRequestDispatcher) {
if hc := out.(*http.Client); hc.Timeout != 0 {
t.Fatalf("期望 Timeout=0, 得到 %v", hc.Timeout)
}
},
},
{
name: "非 http.Client 原样返回", in: staticDispatcher{}, timeout: 300 * time.Second,
check: func(t *testing.T, out common.HTTPRequestDispatcher) {
if _, ok := out.(staticDispatcher); !ok {
t.Fatalf("期望原样返回自定义 dispatcher, 得到 %#v", out)
}
},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) { tt.check(t, dispatcherWithTimeout(tt.in, tt.timeout)) })
}
}
// ctxReader 模拟真实 HTTP body:请求 ctx 取消后读即失败。
type ctxReader struct {
ctx context.Context
r io.Reader
}
func (c *ctxReader) Read(p []byte) (int, error) {
if err := c.ctx.Err(); err != nil {
return 0, err
}
return c.r.Read(p)
}
// fakeCaller 在 delay 后返回响应;期间 ctx 取消则按真实行为返回 ctx 错误。
type fakeCaller struct {
delay time.Duration
body string
gotCtx context.Context
failErr error
}
func (f *fakeCaller) Call(ctx context.Context, _ *http.Request) (*http.Response, error) {
f.gotCtx = ctx
if f.failErr != nil {
return nil, f.failErr
}
select {
case <-time.After(f.delay):
body := &ctxReader{ctx: ctx, r: strings.NewReader(f.body)}
return &http.Response{StatusCode: http.StatusOK, Body: io.NopCloser(body)}, nil
case <-ctx.Done():
return nil, ctx.Err()
}
}
func TestCallWithHeaderBudget(t *testing.T) {
req, _ := http.NewRequest(http.MethodPost, "/actions/v1/responses", nil)
t.Run("预算内返回响应头后长读不受预算影响", func(t *testing.T) {
fc := &fakeCaller{delay: 0, body: "data: hello"}
stream, err := callWithHeaderBudget(context.Background(), fc, req, 30*time.Millisecond)
if err != nil {
t.Fatalf("callWithHeaderBudget = %v", err)
}
defer stream.Close()
time.Sleep(90 * time.Millisecond) // 远超预算,验证响应头到达后预算已解除
payload, err := io.ReadAll(stream)
if err != nil || string(payload) != "data: hello" {
t.Fatalf("预算解除后读流 = %q, %v; 期望完整 body", payload, err)
}
})
t.Run("超预算未返回响应头即取消", func(t *testing.T) {
fc := &fakeCaller{delay: time.Minute}
start := time.Now()
_, err := callWithHeaderBudget(context.Background(), fc, req, 30*time.Millisecond)
if !errors.Is(err, context.Canceled) {
t.Fatalf("期望 context.Canceled, 得到 %v", err)
}
if elapsed := time.Since(start); elapsed > 5*time.Second {
t.Fatalf("取消耗时 %v, 未受预算约束", elapsed)
}
})
t.Run("关闭流时取消派生 ctx", func(t *testing.T) {
fc := &fakeCaller{delay: 0, body: "x"}
stream, err := callWithHeaderBudget(context.Background(), fc, req, time.Minute)
if err != nil {
t.Fatalf("callWithHeaderBudget = %v", err)
}
stream.Close()
if fc.gotCtx.Err() == nil {
t.Fatal("Close 后派生 ctx 应已取消")
}
})
t.Run("建立失败时同样取消派生 ctx", func(t *testing.T) {
fc := &fakeCaller{failErr: errors.New("boom")}
if _, err := callWithHeaderBudget(context.Background(), fc, req, time.Minute); err == nil {
t.Fatal("期望建立失败")
}
if fc.gotCtx.Err() == nil {
t.Fatal("失败路径派生 ctx 应已取消")
}
})
}
+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)
}
})
}
}
+19 -3
View File
@@ -2,6 +2,7 @@ package oci
import ( import (
"fmt" "fmt"
"net"
"net/http" "net/http"
"net/url" "net/url"
"time" "time"
@@ -23,6 +24,13 @@ type ProxySpec struct {
// proxyClientTimeout 与 SDK 默认 HTTPClient 超时保持一致。 // proxyClientTimeout 与 SDK 默认 HTTPClient 超时保持一致。
const proxyClientTimeout = 60 * time.Second const proxyClientTimeout = 60 * time.Second
// 阶段超时对齐 SDK 直连 Transport 模板(transport_template_provider):连不上的
// 代理快速失败,而不是拖满总超时;responses 直通去掉总超时后这是建立阶段的兜底之一。
const (
proxyDialTimeout = 30 * time.Second
proxyTLSHandshakeTimeout = 10 * time.Second
)
// applyProxy 在 SDK client 构造后统一挂出站代理;未关联代理时不动默认配置。 // applyProxy 在 SDK client 构造后统一挂出站代理;未关联代理时不动默认配置。
// 所有 New*ClientWithConfigurationProvider 调用点构造成功后都必须经过这里。 // 所有 New*ClientWithConfigurationProvider 调用点构造成功后都必须经过这里。
func applyProxy(base *common.BaseClient, cred Credentials) { func applyProxy(base *common.BaseClient, cred Credentials) {
@@ -52,18 +60,23 @@ func HTTPClientFor(p *ProxySpec) *http.Client {
// transportFor 构造代理 Transport:http / https 走 CONNECT,socks5 走拨号器。 // transportFor 构造代理 Transport:http / https 走 CONNECT,socks5 走拨号器。
func transportFor(p *ProxySpec) *http.Transport { func transportFor(p *ProxySpec) *http.Transport {
addr := fmt.Sprintf("%s:%d", p.Host, p.Port) addr := fmt.Sprintf("%s:%d", p.Host, p.Port)
dialer := &net.Dialer{Timeout: proxyDialTimeout}
if p.Type == "http" || p.Type == "https" { if p.Type == "http" || p.Type == "https" {
u := &url.URL{Scheme: p.Type, Host: addr} u := &url.URL{Scheme: p.Type, Host: addr}
if p.Username != "" { if p.Username != "" {
u.User = url.UserPassword(p.Username, p.Password) u.User = url.UserPassword(p.Username, p.Password)
} }
return &http.Transport{Proxy: http.ProxyURL(u)} return &http.Transport{
Proxy: http.ProxyURL(u),
DialContext: dialer.DialContext,
TLSHandshakeTimeout: proxyTLSHandshakeTimeout,
}
} }
var auth *proxy.Auth var auth *proxy.Auth
if p.Username != "" { if p.Username != "" {
auth = &proxy.Auth{User: p.Username, Password: p.Password} auth = &proxy.Auth{User: p.Username, Password: p.Password}
} }
d, err := proxy.SOCKS5("tcp", addr, auth, proxy.Direct) d, err := proxy.SOCKS5("tcp", addr, auth, dialer)
if err != nil { if err != nil {
return nil return nil
} }
@@ -71,5 +84,8 @@ func transportFor(p *ProxySpec) *http.Transport {
if !ok { if !ok {
return nil return nil
} }
return &http.Transport{DialContext: cd.DialContext} return &http.Transport{
DialContext: cd.DialContext,
TLSHandshakeTimeout: proxyTLSHandshakeTimeout,
}
} }
+34
View File
@@ -0,0 +1,34 @@
package oci
import (
"testing"
)
func TestTransportForStageTimeouts(t *testing.T) {
tests := []struct {
name string
spec *ProxySpec
wantProxy bool // CONNECT 分支应设 Proxy 函数
}{
{name: "http CONNECT 代理", spec: &ProxySpec{Type: "http", Host: "127.0.0.1", Port: 8080}, wantProxy: true},
{name: "https CONNECT 代理", spec: &ProxySpec{Type: "https", Host: "127.0.0.1", Port: 8443}, wantProxy: true},
{name: "socks5 代理", spec: &ProxySpec{Type: "socks5", Host: "127.0.0.1", Port: 1080, Username: "u", Password: "p"}},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
tr := transportFor(tt.spec)
if tr == nil {
t.Fatal("transportFor 返回 nil")
}
if (tr.Proxy != nil) != tt.wantProxy {
t.Fatalf("Proxy 函数存在性 = %v, 期望 %v", tr.Proxy != nil, tt.wantProxy)
}
if tr.DialContext == nil {
t.Fatal("应设置带超时的 DialContext")
}
if tr.TLSHandshakeTimeout != proxyTLSHandshakeTimeout {
t.Fatalf("TLSHandshakeTimeout = %v, 期望 %v", tr.TLSHandshakeTimeout, proxyTLSHandshakeTimeout)
}
})
}
}
+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 {
+29
View File
@@ -68,6 +68,9 @@ type AiGatewayService struct {
// grokWebSearch / grokXSearch 是 xai. 模型服务端搜索工具默认注入开关 // grokWebSearch / grokXSearch 是 xai. 模型服务端搜索工具默认注入开关
grokWebSearch atomic.Bool grokWebSearch atomic.Bool
grokXSearch atomic.Bool grokXSearch atomic.Bool
// upstreamWaitSec 是 responses 直通的上游无响应预算(秒):非流式为单次尝试
// 总超时,流式为等待响应头预算;multi-agent/搜索类模型远超 SDK 默认 60s
upstreamWaitSec atomic.Int64
} }
// NewAiGatewayService 组装依赖;调用 StartCleanup 后开始调用日志周期清理。 // NewAiGatewayService 组装依赖;调用 StartCleanup 后开始调用日志周期清理。
@@ -78,6 +81,7 @@ func NewAiGatewayService(db *gorm.DB, configs *OciConfigService, client oci.Clie
s.streamGuardKB.Store(int64(loadIntSetting(db, settingAiStreamGuardKB, defaultStreamGuardKB))) s.streamGuardKB.Store(int64(loadIntSetting(db, settingAiStreamGuardKB, defaultStreamGuardKB)))
s.grokWebSearch.Store(loadBoolSetting(db, settingAiGrokWebSearch, true)) s.grokWebSearch.Store(loadBoolSetting(db, settingAiGrokWebSearch, true))
s.grokXSearch.Store(loadBoolSetting(db, settingAiGrokXSearch, true)) s.grokXSearch.Store(loadBoolSetting(db, settingAiGrokXSearch, true))
s.upstreamWaitSec.Store(int64(loadIntSetting(db, settingAiUpstreamWaitSec, defaultUpstreamWaitSec)))
return s return s
} }
@@ -93,11 +97,17 @@ const (
// 开关,缺省开。 // 开关,缺省开。
settingAiGrokWebSearch = "ai_grok_web_search" settingAiGrokWebSearch = "ai_grok_web_search"
settingAiGrokXSearch = "ai_grok_x_search" settingAiGrokXSearch = "ai_grok_x_search"
// settingAiUpstreamWaitSec 是 responses 直通的上游无响应预算(秒)。
settingAiUpstreamWaitSec = "ai_upstream_wait_seconds"
) )
// defaultStreamGuardKB 是保险丝阈值缺省值,低于实测断流边界留余量。 // defaultStreamGuardKB 是保险丝阈值缺省值,低于实测断流边界留余量。
const defaultStreamGuardKB = 60 const defaultStreamGuardKB = 60
// defaultUpstreamWaitSec 是上游无响应预算缺省值(秒):multi-agent 非流式
// 实测 100~180s 才回响应头,给足余量;上下限见 SetUpstreamWait。
const defaultUpstreamWaitSec = 300
// loadBoolSetting 读 settings 表布尔键,无行或值非法时返回缺省。 // loadBoolSetting 读 settings 表布尔键,无行或值非法时返回缺省。
func loadBoolSetting(db *gorm.DB, key string, def bool) bool { func loadBoolSetting(db *gorm.DB, key string, def bool) bool {
var row model.Setting var row model.Setting
@@ -165,6 +175,25 @@ func (s *AiGatewayService) SetStreamGuard(ctx context.Context, on bool, kb int)
return nil return nil
} }
// UpstreamWait 返回 responses 直通的上游无响应预算。
func (s *AiGatewayService) UpstreamWait() time.Duration {
return time.Duration(s.upstreamWaitSec.Load()) * time.Second
}
// SetUpstreamWait 持久化并即时生效上游无响应预算;sec 限定 30..900。
func (s *AiGatewayService) SetUpstreamWait(ctx context.Context, sec int) error {
if sec < 30 || sec > 900 {
return fmt.Errorf("上游无响应预算须在 30..900 秒, 收到 %d", sec)
}
err := s.db.WithContext(ctx).
Save(&model.Setting{Key: settingAiUpstreamWaitSec, Value: strconv.Itoa(sec)}).Error
if err != nil {
return fmt.Errorf("保存上游无响应预算: %w", err)
}
s.upstreamWaitSec.Store(int64(sec))
return nil
}
// GrokSearch 返回 grok 服务端搜索工具默认注入开关(web_search, x_search)。 // GrokSearch 返回 grok 服务端搜索工具默认注入开关(web_search, x_search)。
func (s *AiGatewayService) GrokSearch() (bool, bool) { func (s *AiGatewayService) GrokSearch() (bool, bool) {
return s.grokWebSearch.Load(), s.grokXSearch.Load() return s.grokWebSearch.Load(), s.grokXSearch.Load()
+27 -65
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) {
@@ -64,42 +72,19 @@ func (s *AiGatewayService) passthroughOnce(ctx context.Context, cand *aiCandidat
if err != nil { if err != nil {
return nil, err return nil, err
} }
return s.client.GenAiCompatResponses(ctx, cred, cand.ch.Region, raw) return s.client.GenAiCompatResponses(ctx, cred, cand.ch.Region, raw, s.UpstreamWait())
} }
// 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) 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 形态响应。
+37 -2
View File
@@ -83,7 +83,7 @@ func (f *gatewayStubClient) GenAiApplyGuardrails(ctx context.Context, cred oci.C
return f.guardOutcome, f.guardErr return f.guardOutcome, f.guardErr
} }
func (f *gatewayStubClient) GenAiCompatResponses(ctx context.Context, cred oci.Credentials, region string, body []byte) ([]byte, error) { func (f *gatewayStubClient) GenAiCompatResponses(ctx context.Context, cred oci.Credentials, region string, body []byte, wait time.Duration) ([]byte, error) {
f.passCalls++ f.passCalls++
f.passRegions = append(f.passRegions, region) f.passRegions = append(f.passRegions, region)
if len(f.passErrs) > 0 { if len(f.passErrs) > 0 {
@@ -96,7 +96,7 @@ func (f *gatewayStubClient) GenAiCompatResponses(ctx context.Context, cred oci.C
return f.passPayload, nil return f.passPayload, nil
} }
func (f *gatewayStubClient) GenAiCompatResponsesStream(ctx context.Context, cred oci.Credentials, region string, body []byte) (io.ReadCloser, error) { func (f *gatewayStubClient) GenAiCompatResponsesStream(ctx context.Context, cred oci.Credentials, region string, body []byte, wait time.Duration) (io.ReadCloser, error) {
f.passCalls++ f.passCalls++
f.passRegions = append(f.passRegions, region) f.passRegions = append(f.passRegions, region)
if len(f.passErrs) > 0 { if len(f.passErrs) > 0 {
@@ -1124,3 +1124,38 @@ func TestAggregatedModelsFilterDeprecated(t *testing.T) {
t.Fatalf("开关开:弃用模型应被过滤, got %+v", items) t.Fatalf("开关开:弃用模型应被过滤, got %+v", items)
} }
} }
func TestUpstreamWaitSetting(t *testing.T) {
gw, _ := newTestGateway(t, &gatewayStubClient{fakeClient: &fakeClient{}})
ctx := context.Background()
if got := gw.UpstreamWait(); got != 300*time.Second {
t.Fatalf("缺省上游无响应预算 = %v, 期望 300s", got)
}
tests := []struct {
name string
sec int
wantErr bool
}{
{name: "下界 30 有效", sec: 30},
{name: "上界 900 有效", sec: 900},
{name: "低于下界拒绝", sec: 29, wantErr: true},
{name: "高于上界拒绝", sec: 901, wantErr: true},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
err := gw.SetUpstreamWait(ctx, tt.sec)
if (err != nil) != tt.wantErr {
t.Fatalf("SetUpstreamWait(%d) = %v, wantErr=%v", tt.sec, err, tt.wantErr)
}
if !tt.wantErr && gw.UpstreamWait() != time.Duration(tt.sec)*time.Second {
t.Fatalf("UpstreamWait = %v, 期望 %ds", gw.UpstreamWait(), tt.sec)
}
})
}
// 持久化后新实例应加载已存值(最后一次成功设置为 900)
gw2 := NewAiGatewayService(gw.db, nil, nil)
if got := gw2.UpstreamWait(); got != 900*time.Second {
t.Fatalf("重建服务加载预算 = %v, 期望 900s", got)
}
}
+13 -5
View File
@@ -30,19 +30,23 @@ var ErrInvalidAuditCursor = errors.New("audit events: invalid cursor, refresh to
var ErrAuditEventGone = errors.New("原始事件已不可取回,请刷新列表后重试") var ErrAuditEventGone = errors.New("原始事件已不可取回,请刷新列表后重试")
// AuditQuery 是批式懒加载查询参数:Cursor 为空表示自当前时刻首查, // AuditQuery 是批式懒加载查询参数:Cursor 为空表示自当前时刻首查,
// 非空则从上次响应的游标位置继续向更早回溯;Limit 为单批目标条数 // 非空则从上次响应的游标位置继续向更早回溯;Limit 为单批目标条数;
// Q 为检索关键字,仅首查生效(续查沿用游标内嵌的关键字,保证跨批一致)。
type AuditQuery struct { type AuditQuery struct {
Region string Region string
Cursor string Cursor string
Limit int Limit int
Q string
} }
// AuditEventsView 是批式查询响应:列表不含 raw(详情接口取回); // AuditEventsView 是批式查询响应:列表不含 raw(详情接口取回);
// Cursor 供下一批续查原样带回,空且 Exhausted 表示已到 365 天保留期尽头 // Cursor 供下一批续查原样带回,空且 Exhausted 表示已到 365 天保留期尽头;
// ScannedThrough 为已完整回溯到的时刻(比它更新的时段已扫完),供前端展示进度。
type AuditEventsView struct { type AuditEventsView struct {
Items []oci.AuditEvent `json:"items"` Items []oci.AuditEvent `json:"items"`
Cursor string `json:"cursor,omitempty"` Cursor string `json:"cursor,omitempty"`
Exhausted bool `json:"exhausted"` Exhausted bool `json:"exhausted"`
ScannedThrough *time.Time `json:"scannedThrough,omitempty"`
} }
// AuditEvents 实时查询租户 OCI 审计事件,纯透传不入库;region 为空时用配置 // AuditEvents 实时查询租户 OCI 审计事件,纯透传不入库;region 为空时用配置
@@ -52,6 +56,9 @@ func (s *OciConfigService) AuditEvents(ctx context.Context, id uint, q AuditQuer
if err != nil { if err != nil {
return AuditEventsView{}, err return AuditEventsView{}, err
} }
if q.Cursor == "" {
cur.Q = oci.SanitizeAuditTerm(q.Q)
}
cred, err := s.credentialsByID(ctx, id) cred, err := s.credentialsByID(ctx, id)
if err != nil { if err != nil {
return AuditEventsView{}, err return AuditEventsView{}, err
@@ -63,6 +70,7 @@ func (s *OciConfigService) AuditEvents(ctx context.Context, id uint, q AuditQuer
view := AuditEventsView{Items: s.stripAuditRaw(id, res.Items), Exhausted: res.Exhausted} view := AuditEventsView{Items: s.stripAuditRaw(id, res.Items), Exhausted: res.Exhausted}
if res.Cursor != nil { if res.Cursor != nil {
view.Cursor = encodeAuditCursor(*res.Cursor) view.Cursor = encodeAuditCursor(*res.Cursor)
view.ScannedThrough = &res.Cursor.End
} }
return view, nil return view, nil
} }
+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)
}
}
})
}
}