5 Commits
Author SHA1 Message Date
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
21 changed files with 1234 additions and 264 deletions
+1
View File
@@ -20,6 +20,7 @@ Gin + GORM + SQLite(纯 Go 驱动 `glebarez/sqlite`,免 CGO;默认与推荐)+ AE
| [Concurrency](./concurrency.md) | goroutine 生命周期与 context 传递 | 已填 |
| [Testing](./testing.md) | table-driven 测试要求 | 已填 |
| [Database Guidelines](./database-guidelines.md) | ORM 模式、查询、迁移 | 待填 |
| [OCI Audit](./oci-audit.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/)(版本段不记日期),版本号遵循语义化版本。
## [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]
### Added
+1 -1
View File
@@ -1 +1 @@
v0.7.0
v0.7.2
+14 -1
View File
@@ -860,7 +860,7 @@ const docTemplate = `{
"summary": "更新 AI 网关全局设置",
"parameters": [
{
"description": "全量提交;保险丝阈值限 1..1024 KB",
"description": "全量提交;保险丝阈值限 1..1024 KB,上游无响应预算限 30..900 秒",
"name": "body",
"in": "body",
"required": true,
@@ -1572,6 +1572,12 @@ const docTemplate = `{
"description": "单批目标条数,缺省 100,上限 200",
"name": "limit",
"in": "query"
},
{
"type": "string",
"description": "检索关键字(服务端全文匹配,支持 * 通配;仅首查生效)",
"name": "q",
"in": "query"
}
],
"responses": {
@@ -6096,6 +6102,10 @@ const docTemplate = `{
},
"streamGuardKB": {
"type": "integer"
},
"upstreamWaitSeconds": {
"description": "UpstreamWaitSeconds 是 responses 直通的上游无响应预算(秒):非流式为单次\n尝试总超时,流式为等待响应头预算;multi-agent/搜索类模型需远超 60s",
"type": "integer"
}
}
},
@@ -9439,6 +9449,9 @@ const docTemplate = `{
"items": {
"$ref": "#/definitions/oci-portal_internal_oci.AuditEvent"
}
},
"scannedThrough": {
"type": "string"
}
}
},
+14 -1
View File
@@ -853,7 +853,7 @@
"summary": "更新 AI 网关全局设置",
"parameters": [
{
"description": "全量提交;保险丝阈值限 1..1024 KB",
"description": "全量提交;保险丝阈值限 1..1024 KB,上游无响应预算限 30..900 秒",
"name": "body",
"in": "body",
"required": true,
@@ -1565,6 +1565,12 @@
"description": "单批目标条数,缺省 100,上限 200",
"name": "limit",
"in": "query"
},
{
"type": "string",
"description": "检索关键字(服务端全文匹配,支持 * 通配;仅首查生效)",
"name": "q",
"in": "query"
}
],
"responses": {
@@ -6089,6 +6095,10 @@
},
"streamGuardKB": {
"type": "integer"
},
"upstreamWaitSeconds": {
"description": "UpstreamWaitSeconds 是 responses 直通的上游无响应预算(秒):非流式为单次\n尝试总超时,流式为等待响应头预算;multi-agent/搜索类模型需远超 60s",
"type": "integer"
}
}
},
@@ -9432,6 +9442,9 @@
"items": {
"$ref": "#/definitions/oci-portal_internal_oci.AuditEvent"
}
},
"scannedThrough": {
"type": "string"
}
}
},
+12 -1
View File
@@ -65,6 +65,11 @@ definitions:
type: boolean
streamGuardKB:
type: integer
upstreamWaitSeconds:
description: |-
UpstreamWaitSeconds 是 responses 直通的上游无响应预算(秒):非流式为单次
尝试总超时,流式为等待响应头预算;multi-agent/搜索类模型需远超 60s
type: integer
type: object
internal_api.attachBootVolumeRequest:
properties:
@@ -2271,6 +2276,8 @@ definitions:
items:
$ref: '#/definitions/oci-portal_internal_oci.AuditEvent'
type: array
scannedThrough:
type: string
type: object
oci-portal_internal_service.Changes:
additionalProperties:
@@ -3202,7 +3209,7 @@ paths:
- AI 管理
put:
parameters:
- description: 全量提交;保险丝阈值限 1..1024 KB
- description: 全量提交;保险丝阈值限 1..1024 KB,上游无响应预算限 30..900 秒
in: body
name: body
required: true
@@ -3637,6 +3644,10 @@ paths:
in: query
name: limit
type: integer
- description: 检索关键字(服务端全文匹配,支持 * 通配;仅首查生效)
in: query
name: q
type: string
responses:
"200":
description: OK
+19 -6
View File
@@ -3,6 +3,7 @@ package api
import (
"net/http"
"strconv"
"time"
"github.com/gin-gonic/gin"
@@ -450,6 +451,9 @@ type aiSettingsResponse struct {
// 请求 tools 已包含同名工具时不覆盖
GrokWebSearch bool `json:"grokWebSearch"`
GrokXSearch bool `json:"grokXSearch"`
// UpstreamWaitSeconds 是 responses 直通的上游无响应预算(秒):非流式为单次
// 尝试总超时,流式为等待响应头预算;multi-agent/搜索类模型需远超 60s
UpstreamWaitSeconds int `json:"upstreamWaitSeconds"`
}
// currentAiSettings 汇总网关运行时设置为响应体。
@@ -457,11 +461,12 @@ func (h *aiAdminHandler) currentAiSettings() aiSettingsResponse {
guardOn, guardKB := h.gw.StreamGuard()
web, x := h.gw.GrokSearch()
return aiSettingsResponse{
FilterDeprecated: h.gw.FilterDeprecated(),
StreamGuardEnabled: guardOn,
StreamGuardKB: guardKB,
GrokWebSearch: web,
GrokXSearch: x,
FilterDeprecated: h.gw.FilterDeprecated(),
StreamGuardEnabled: guardOn,
StreamGuardKB: guardKB,
GrokWebSearch: web,
GrokXSearch: x,
UpstreamWaitSeconds: int(h.gw.UpstreamWait() / time.Second),
}
}
@@ -480,7 +485,7 @@ func (h *aiAdminHandler) aiSettings(c *gin.Context) {
//
// @Summary 更新 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
// @Failure 400 {object} map[string]string
// @Security BearerAuth
@@ -495,6 +500,10 @@ func (h *aiAdminHandler) updateAiSettings(c *gin.Context) {
c.JSON(http.StatusBadRequest, gin.H{"error": "streamGuardKB 须在 1..1024"})
return
}
if req.UpstreamWaitSeconds < 30 || req.UpstreamWaitSeconds > 900 {
c.JSON(http.StatusBadRequest, gin.H{"error": "upstreamWaitSeconds 须在 30..900"})
return
}
ctx := c.Request.Context()
if err := h.gw.SetFilterDeprecated(ctx, req.FilterDeprecated); err != nil {
respondError(c, err)
@@ -508,5 +517,9 @@ func (h *aiAdminHandler) updateAiSettings(c *gin.Context) {
respondError(c, err)
return
}
if err := h.gw.SetUpstreamWait(ctx, req.UpstreamWaitSeconds); err != nil {
respondError(c, err)
return
}
c.JSON(http.StatusOK, h.currentAiSettings())
}
+4 -2
View File
@@ -68,13 +68,15 @@ func (h *ociConfigHandler) costs(c *gin.Context) {
// ---- 租户审计日志 ----
// getAuditEvents 批式懒加载查询审计事件:cursor 为空自当前时刻首查,
// 非空从上次响应游标继续向更早回溯;limit 单批目标条数(缺省 100,上限 200)
// 非空从上次响应游标继续向更早回溯;limit 单批目标条数(缺省 100,上限 200);
// q 为服务端全文检索关键字,仅首查生效,续查沿用游标内嵌关键字。
//
// @Summary 批式懒加载查询租户 OCI 审计事件
// @Tags 租户 IAM
// @Param id path int true "配置 ID"
// @Param cursor query string false "续查游标(上次响应原样带回)"
// @Param limit query int false "单批目标条数,缺省 100,上限 200"
// @Param q query string false "检索关键字(服务端全文匹配,支持 * 通配;仅首查生效)"
// @Success 200 {object} service.AuditEventsView
// @Security BearerAuth
// @Router /api/v1/oci-configs/{id}/audit-events [get]
@@ -84,7 +86,7 @@ func (h *ociConfigHandler) getAuditEvents(c *gin.Context) {
return
}
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)
if errors.Is(err, service.ErrInvalidAuditCursor) {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
+469 -184
View File
@@ -6,19 +6,29 @@ import (
"fmt"
"net"
"sort"
"strings"
"time"
"github.com/oracle/oci-go-sdk/v65/audit"
"github.com/oracle/oci-go-sdk/v65/common"
"github.com/oracle/oci-go-sdk/v65/loggingsearch"
)
// maxAuditPages 限制单次查询的翻页数:繁忙租户单日事件可上千,
// 到限即返回 Truncated=true,由调用方收窄时间窗
// 默认过滤(噪声事件/内网发起)后有效结果变少,页数放宽到 10 缓解截断。
// maxAuditPages 限制单次查询的翻页数:每页最多 auditSearchPageLimit 条,
// 到限即截断(窗口式回传 Truncated,批式留游标),由调用方续查
const maxAuditPages = 10
// auditBatchTimeBudget 是批式查询的单批耗时预算:全文检索命中稀疏时
// 大窗扫描单页可达十余秒,超时即带游标返回,把长回溯拆成多个有界请求,
// 前端按已回溯位置展示进度并自动续查。
const auditBatchTimeBudget = 20 * time.Second
// auditSearchPageLimit 是 SearchLogs 单页条数(API 上限 1000):批式查询
// 攒满目标条数(~100)即携整页返回,页取 200 兼顾单页凑满一批与响应体量。
const auditSearchPageLimit = 200
// AuditEvent 是审计事件的列表精简视图;EventId 为 CloudEvents 全局唯一 id,
// 详情反查的键。Raw 为 SDK 原始事件的 JSON 序列化,由 service 层剥离进缓存,
// 详情反查的键。Raw 为 _Audit 日志 logContent 原文,由 service 层剥离进缓存,
// 列表响应不再携带(详情接口按 eventId 取回)。
type AuditEvent struct {
EventId string `json:"eventId"`
@@ -35,7 +45,7 @@ type AuditEvent struct {
Raw json.RawMessage `json:"raw,omitempty"`
}
// AuditEventsResult 是一次审计查询的结果;Truncated 表示翻页到限被截断,
// AuditEventsResult 是一次窗口式审计查询的结果;Truncated 表示翻页到限被截断,
// 此时 NextPage 携带 opc-next-page 游标,同一时间窗回传可断点续翻。
type AuditEventsResult struct {
Items []AuditEvent `json:"items"`
@@ -43,6 +53,324 @@ type AuditEventsResult struct {
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) {
ac, err := audit.NewAuditClientWithConfigurationProvider(provider(cred))
if err != nil {
@@ -55,153 +383,10 @@ func (c *RealClient) auditClient(cred Credentials, region string) (audit.AuditCl
return ac, nil
}
// ListAuditEvents 实现 Client:实时查询租户根 compartment 在 [start, end) 内的
// 审计事件,最多翻 maxAuditPages 页,结果按发生时间倒序;纯读不落库。
// page 非空时从该游标断点续翻(必须配同一时间窗);到限截断时回传 NextPage
func (c *RealClient) ListAuditEvents(ctx context.Context, cred Credentials, region string, start, end time.Time, page string) (AuditEventsResult, 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) {
// listAuditPage 拉取窗口内一页 Audit API 原始事件并压平;该 API 窗口内固定
// 正序且只接受分钟粒度(起止秒与毫秒必须为 0)。Raw 为 SDK 事件原文,
// 与 Search 通道的 logContent 形态不同,详情弹窗均按任意 JSON 渲染
func listAuditPage(ctx context.Context, ac audit.AuditClient, tenancyOCID string, cur AuditCursor) ([]AuditEvent, string, error) {
req := audit.ListEventsRequest{
CompartmentId: &tenancyOCID,
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 {
return nil, "", fmt.Errorf("list audit events: %w", err)
}
return resp.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 保留有展示价值的事件: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 := make([]AuditEvent, 0, len(resp.Items))
for _, ev := range resp.Items {
out := toAuditEvent(ev)
if raw, mErr := json.Marshal(ev); mErr == nil {
out.Raw = raw
}
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 {
out := AuditEvent{EventId: deref(ev.EventId), Source: deref(ev.Source)}
if ev.EventTime != nil {
@@ -275,6 +436,130 @@ func toAuditEvent(ev audit.AuditEvent) AuditEvent {
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 时间排最后。
func sortAuditEvents(items []AuditEvent) {
sort.SliceStable(items, func(i, j int) bool {
+294 -27
View File
@@ -1,14 +1,262 @@
package oci
import (
"context"
"encoding/json"
"errors"
"fmt"
"reflect"
"strings"
"testing"
"time"
"github.com/oracle/oci-go-sdk/v65/audit"
"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) {
eventTime := time.Date(2026, 7, 6, 10, 30, 0, 0, time.UTC)
tests := []struct {
@@ -38,44 +286,23 @@ func TestToAuditEvent(t *testing.T) {
},
},
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...",
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: "Data 为 nil 时只保留信封字段",
ev: audit.AuditEvent{
Source: common.String("VcnApi"),
EventTime: &common.SDKTime{Time: eventTime},
},
want: AuditEvent{EventTime: &eventTime, Source: "VcnApi"},
},
{
name: "嵌套局部 nil 各自安全跳过",
ev: audit.AuditEvent{
Data: &audit.Data{
EventName: common.String("GetInstance"),
Identity: nil,
Request: &audit.Request{Path: common.String("/instances")},
Response: nil,
},
},
want: AuditEvent{EventName: "GetInstance", RequestPath: "/instances"},
},
{
name: "空事件全部零值",
ev: audit.AuditEvent{},
want: AuditEvent{},
},
{name: "空事件全部零值", ev: audit.AuditEvent{}, want: AuditEvent{}},
}
for _, tt := range tests {
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 不参与,其余反射比较。
func auditEventEqual(a, b AuditEvent) bool {
if (a.EventTime == nil) != (b.EventTime == nil) {
@@ -149,6 +407,7 @@ func TestAuditCursorAdvance(t *testing.T) {
Start: now.Add(-24 * time.Hour),
End: now,
WindowHours: 24,
Q: "kw",
}
cases := []struct {
name string
@@ -159,8 +418,10 @@ func TestAuditCursorAdvance(t *testing.T) {
}{
{"有事件重置 24h 窗", AuditCursor{Start: base.Start, End: base.End, WindowHours: 96}, false, 24, 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},
{"回退模式有事件重置 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},
}
for _, tc := range cases {
@@ -184,6 +445,12 @@ func TestAuditCursorAdvance(t *testing.T) {
if 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)
}
})
}
}
+6 -4
View File
@@ -102,10 +102,12 @@ type Client interface {
ListGenAiModels(ctx context.Context, cred Credentials, region string) ([]GenAiModel, 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)
// GenAiCompatResponses 直通 OpenAI Responses 请求体到 /actions/v1/responses(xAI 服务端工具通路)
GenAiCompatResponses(ctx context.Context, cred Credentials, region string, body []byte) ([]byte, error)
// GenAiCompatResponsesStream 流式直通 /actions/v1/responses,建立成功返回 SSE body。
GenAiCompatResponsesStream(ctx context.Context, cred Credentials, region string, body []byte) (io.ReadCloser, error)
// GenAiCompatResponses 直通 OpenAI Responses 请求体到 /actions/v1/responses(xAI 服务端工具通路);
// wait 是上游无响应预算(整请求总超时),multi-agent/搜索类模型需远超 SDK 默认 60s。
GenAiCompatResponses(ctx context.Context, cred Credentials, region string, body []byte, wait time.Duration) ([]byte, 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(ctx context.Context, cred Credentials, region string, body []byte) ([]byte, string, error)
// GenAiRerank 文档重排,返回按相关度排序的下标与得分。
+2 -1
View File
@@ -158,7 +158,8 @@ func (c *RealClient) GenAiProbeChat(ctx context.Context, cred Credentials, regio
if err != nil {
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
}
if status, ok := ServiceStatus(err); ok {
+74 -24
View File
@@ -6,6 +6,7 @@ import (
"fmt"
"io"
"net/http"
"time"
"github.com/oracle/oci-go-sdk/v65/common"
)
@@ -13,18 +14,19 @@ import (
// compatResponsesLimit 限制直通响应体大小;web_search 输出含多段引用,给足余量。
const compatResponsesLimit = int64(8 << 20)
// GenAiCompatResponses 实现 Client:把 OpenAI Responses 请求体直通到 OCI
// `/20231130/actions/v1/responses`(IAM 签名)。xAI 服务端工具(web_search /
// x_search / code_interpreter)与 mcp 已被 Oracle 文档正式支持,工具参数与限制
// 遵循 xAI 规格;调用方须自行校验并改写请求体(store/stream)。
func (c *RealClient) GenAiCompatResponses(ctx context.Context, cred Credentials, region string, body []byte) ([]byte, error) {
ic, err := c.genAiInferenceClient(cred, region)
if err != nil {
return nil, err
// dispatcherWithTimeout 把 dispatcher 换成指定总超时的拷贝(保留 Transport,
// 代理链路不受影响);timeout=0 表示无总超时(流式读 body 不能有总时限)。
// 非 *http.Client 的自定义 dispatcher 保持原样,维持既有超时行为。
func dispatcherWithTimeout(d common.HTTPRequestDispatcher, timeout time.Duration) common.HTTPRequestDispatcher {
hc, ok := d.(*http.Client)
if !ok {
return d
}
client := ic.BaseClient
common.UpdateEndpointTemplateForOptions(&client)
common.SetMissingTemplateParams(&client)
return &http.Client{Transport: hc.Transport, Timeout: timeout}
}
// 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))
if err != nil {
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("CompartmentId", 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)
if err != nil {
return nil, err
@@ -44,10 +67,46 @@ func (c *RealClient) GenAiCompatResponses(ctx context.Context, cred Credentials,
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`,
// 建立成功(2xx)返回 SSE body(调用方负责 Close);建立失败返回 SDK ServiceError,
// 与既有渠道切换/熔断错误分类兼容。请求体须由调用方置 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)
if err != nil {
return nil, err
@@ -55,19 +114,10 @@ func (c *RealClient) GenAiCompatResponsesStream(ctx context.Context, cred Creden
client := ic.BaseClient
common.UpdateEndpointTemplateForOptions(&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 {
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 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 应已取消")
}
})
}
+19 -3
View File
@@ -2,6 +2,7 @@ package oci
import (
"fmt"
"net"
"net/http"
"net/url"
"time"
@@ -23,6 +24,13 @@ type ProxySpec struct {
// proxyClientTimeout 与 SDK 默认 HTTPClient 超时保持一致。
const proxyClientTimeout = 60 * time.Second
// 阶段超时对齐 SDK 直连 Transport 模板(transport_template_provider):连不上的
// 代理快速失败,而不是拖满总超时;responses 直通去掉总超时后这是建立阶段的兜底之一。
const (
proxyDialTimeout = 30 * time.Second
proxyTLSHandshakeTimeout = 10 * time.Second
)
// applyProxy 在 SDK client 构造后统一挂出站代理;未关联代理时不动默认配置。
// 所有 New*ClientWithConfigurationProvider 调用点构造成功后都必须经过这里。
func applyProxy(base *common.BaseClient, cred Credentials) {
@@ -52,18 +60,23 @@ func HTTPClientFor(p *ProxySpec) *http.Client {
// transportFor 构造代理 Transport:http / https 走 CONNECT,socks5 走拨号器。
func transportFor(p *ProxySpec) *http.Transport {
addr := fmt.Sprintf("%s:%d", p.Host, p.Port)
dialer := &net.Dialer{Timeout: proxyDialTimeout}
if p.Type == "http" || p.Type == "https" {
u := &url.URL{Scheme: p.Type, Host: addr}
if p.Username != "" {
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
if p.Username != "" {
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 {
return nil
}
@@ -71,5 +84,8 @@ func transportFor(p *ProxySpec) *http.Transport {
if !ok {
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)
}
})
}
}
+29
View File
@@ -68,6 +68,9 @@ type AiGatewayService struct {
// grokWebSearch / grokXSearch 是 xai. 模型服务端搜索工具默认注入开关
grokWebSearch atomic.Bool
grokXSearch atomic.Bool
// upstreamWaitSec 是 responses 直通的上游无响应预算(秒):非流式为单次尝试
// 总超时,流式为等待响应头预算;multi-agent/搜索类模型远超 SDK 默认 60s
upstreamWaitSec atomic.Int64
}
// 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.grokWebSearch.Store(loadBoolSetting(db, settingAiGrokWebSearch, true))
s.grokXSearch.Store(loadBoolSetting(db, settingAiGrokXSearch, true))
s.upstreamWaitSec.Store(int64(loadIntSetting(db, settingAiUpstreamWaitSec, defaultUpstreamWaitSec)))
return s
}
@@ -93,11 +97,17 @@ const (
// 开关,缺省开。
settingAiGrokWebSearch = "ai_grok_web_search"
settingAiGrokXSearch = "ai_grok_x_search"
// settingAiUpstreamWaitSec 是 responses 直通的上游无响应预算(秒)。
settingAiUpstreamWaitSec = "ai_upstream_wait_seconds"
)
// defaultStreamGuardKB 是保险丝阈值缺省值,低于实测断流边界留余量。
const defaultStreamGuardKB = 60
// defaultUpstreamWaitSec 是上游无响应预算缺省值(秒):multi-agent 非流式
// 实测 100~180s 才回响应头,给足余量;上下限见 SetUpstreamWait。
const defaultUpstreamWaitSec = 300
// loadBoolSetting 读 settings 表布尔键,无行或值非法时返回缺省。
func loadBoolSetting(db *gorm.DB, key string, def bool) bool {
var row model.Setting
@@ -165,6 +175,25 @@ func (s *AiGatewayService) SetStreamGuard(ctx context.Context, on bool, kb int)
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)。
func (s *AiGatewayService) GrokSearch() (bool, bool) {
return s.grokWebSearch.Load(), s.grokXSearch.Load()
+2 -2
View File
@@ -64,7 +64,7 @@ func (s *AiGatewayService) passthroughOnce(ctx context.Context, cand *aiCandidat
if err != nil {
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
@@ -83,7 +83,7 @@ func (s *AiGatewayService) RespPassthroughStream(ctx context.Context, raw []byte
if err != nil {
return nil, meta, err
}
stream, err := s.client.GenAiCompatResponsesStream(ctx, cred, cand.ch.Region, raw)
stream, err := s.client.GenAiCompatResponsesStream(ctx, cred, cand.ch.Region, raw, s.UpstreamWait())
if err == nil {
s.markSuccess(ctx, cand.ch.ID)
return stream, meta, nil
+37 -2
View File
@@ -83,7 +83,7 @@ func (f *gatewayStubClient) GenAiApplyGuardrails(ctx context.Context, cred oci.C
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.passRegions = append(f.passRegions, region)
if len(f.passErrs) > 0 {
@@ -96,7 +96,7 @@ func (f *gatewayStubClient) GenAiCompatResponses(ctx context.Context, cred oci.C
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.passRegions = append(f.passRegions, region)
if len(f.passErrs) > 0 {
@@ -1124,3 +1124,38 @@ func TestAggregatedModelsFilterDeprecated(t *testing.T) {
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("原始事件已不可取回,请刷新列表后重试")
// AuditQuery 是批式懒加载查询参数:Cursor 为空表示自当前时刻首查,
// 非空则从上次响应的游标位置继续向更早回溯;Limit 为单批目标条数
// 非空则从上次响应的游标位置继续向更早回溯;Limit 为单批目标条数;
// Q 为检索关键字,仅首查生效(续查沿用游标内嵌的关键字,保证跨批一致)。
type AuditQuery struct {
Region string
Cursor string
Limit int
Q string
}
// AuditEventsView 是批式查询响应:列表不含 raw(详情接口取回);
// Cursor 供下一批续查原样带回,空且 Exhausted 表示已到 365 天保留期尽头
// Cursor 供下一批续查原样带回,空且 Exhausted 表示已到 365 天保留期尽头;
// ScannedThrough 为已完整回溯到的时刻(比它更新的时段已扫完),供前端展示进度。
type AuditEventsView struct {
Items []oci.AuditEvent `json:"items"`
Cursor string `json:"cursor,omitempty"`
Exhausted bool `json:"exhausted"`
Items []oci.AuditEvent `json:"items"`
Cursor string `json:"cursor,omitempty"`
Exhausted bool `json:"exhausted"`
ScannedThrough *time.Time `json:"scannedThrough,omitempty"`
}
// AuditEvents 实时查询租户 OCI 审计事件,纯透传不入库;region 为空时用配置
@@ -52,6 +56,9 @@ func (s *OciConfigService) AuditEvents(ctx context.Context, id uint, q AuditQuer
if err != nil {
return AuditEventsView{}, err
}
if q.Cursor == "" {
cur.Q = oci.SanitizeAuditTerm(q.Q)
}
cred, err := s.credentialsByID(ctx, id)
if err != nil {
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}
if res.Cursor != nil {
view.Cursor = encodeAuditCursor(*res.Cursor)
view.ScannedThrough = &res.Cursor.End
}
return view, nil
}