9506a82be9
原逐片脱敏有两个漏:①密钥被切成两片("sk-912cf85b"|"16d0...")逐片都不命中正则而漏检; ②贪婪正则在缓冲末尾凑够最短长度就把半截密钥提前脱敏发走、剩余字符随后明文流出(碎片泄漏)。 有状态 StreamRedactor 跨分片缓冲,切点在「原文」上定且绝不切断任何完整匹配: - opener 暂留末尾仍在增长的疑似密钥(sk/AKIA/JWT/Bearer/手机/邮箱/长数字) - 始终留 16B 尾窗兜底 opener 未覆盖的短模式;勿切断完整匹配(循环至稳定) - rune 边界安全:cut 退到最近 rune 起点,中文不被切成半个发出乱码 - 暂留封顶 256B,防对抗性长串无限暂留 / O(n²) - 新增 PII:手机号 / 邮箱 / 身份证(18 位) 3 个流式点(graph/react_agent/compose_graph)统一接入,逐片 Push + 收尾 Flush。 7 单测(跨片/逐字符 JWT/碎片回归/尾窗内匹配/干净重建/PII/无误伤),-race 干净。 live 实测:26 位密钥(曾泄漏 ijkl90mnop 碎片)与邮箱整条 [已脱敏]。 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
105 lines
3.6 KiB
Go
105 lines
3.6 KiB
Go
package eino
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"io"
|
||
"os"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/cloudwego/eino/compose"
|
||
"github.com/cloudwego/eino/schema"
|
||
|
||
"github.com/sundynix/sundynix-dispatcher/internal/harness"
|
||
)
|
||
|
||
// composeEnabled 报告是否启用 compose.Graph 编排路径(Phase C 灰度开关,默认关 → 走自研 graph.go)。
|
||
// 并存策略:EINO_COMPOSE=1 时对话主流程改走 Eino compose 运行时,行为对齐后再逐步退役 graph.go。
|
||
func composeEnabled() bool { return os.Getenv("EINO_COMPOSE") == "1" }
|
||
|
||
// runConversation 是对话/模型节点的统一入口:按灰度开关选 compose.Graph 或自研 runAgent。
|
||
// 二者对外行为一致(据黑板拼消息 → 流式回流 token → 累计成稿),便于等价回归。
|
||
func (o *Orchestrator) runConversation(ctx context.Context, taskID string, b *board, system string, tr *execTracer, node string) {
|
||
if composeEnabled() {
|
||
o.runComposeConversation(ctx, taskID, b, system, tr, node)
|
||
return
|
||
}
|
||
o.runAgent(ctx, taskID, b, system, tr, node)
|
||
}
|
||
|
||
// runComposeConversation 用 Eino compose.Graph 跑对话主流程:
|
||
// START → ChatModel 节点 → END,编译为 Runnable 后流式执行;可观测经 callbacks 桥到 ExecEvent。
|
||
// 模型未就绪 / 编译失败时降级回自研 runAgent,保证不回归。
|
||
func (o *Orchestrator) runComposeConversation(ctx context.Context, taskID string, b *board, system string, tr *execTracer, node string) {
|
||
cm := o.pool.ChatModel()
|
||
if cm == nil {
|
||
o.runAgent(ctx, taskID, b, system, tr, node) // 无模型 → 走自研路径的降级桩
|
||
return
|
||
}
|
||
|
||
g := compose.NewGraph[[]*schema.Message, *schema.Message]()
|
||
if err := g.AddChatModelNode("model", cm); err != nil {
|
||
tr.info(node, "system", "compose 降级", "建图失败,退回自研路径:"+err.Error())
|
||
o.runAgent(ctx, taskID, b, system, tr, node)
|
||
return
|
||
}
|
||
_ = g.AddEdge(compose.START, "model")
|
||
_ = g.AddEdge("model", compose.END)
|
||
r, err := g.Compile(ctx)
|
||
if err != nil {
|
||
tr.info(node, "system", "compose 降级", "编译失败,退回自研路径:"+err.Error())
|
||
o.runAgent(ctx, taskID, b, system, tr, node)
|
||
return
|
||
}
|
||
|
||
rc := &RunCtx{
|
||
UserID: b.uid, SessionID: b.sid,
|
||
System: firstNonEmpty(system, defaultAgentSystem),
|
||
Query: b.query,
|
||
Profile: b.profile,
|
||
History: b.history,
|
||
ToolOut: append(append([]string{}, b.toolOut...), b.refs...),
|
||
Upstream: append([]string{}, b.agentOut...), // 前序协作 agent 产出 → 接力
|
||
}
|
||
msgs, _ := buildMessages(ctx, rc)
|
||
|
||
t0 := time.Now()
|
||
// ChatModel 的 start/end 由 composeTracer(callbacks)落轨迹,这里不再手写 emit(归一)。
|
||
sr, err := r.Stream(ctx, msgs, compose.WithCallbacks(composeTracer(tr, node)))
|
||
if err != nil {
|
||
tr.emit(node, "model", "error", "compose 图执行", err.Error(), time.Since(t0).Milliseconds())
|
||
return
|
||
}
|
||
defer sr.Close()
|
||
|
||
chunks := 0
|
||
var produced strings.Builder // 本节点产出(供下游 agent 接力)
|
||
red := harness.NewStreamRedactor() // 输出护栏:跨分片脱敏,杜绝密钥被切断而漏检
|
||
emit := func(safe string) {
|
||
if safe == "" {
|
||
return
|
||
}
|
||
_ = o.sink.PublishToken(taskID, []byte(safe))
|
||
produced.WriteString(safe)
|
||
chunks++
|
||
}
|
||
for {
|
||
chunk, rerr := sr.Recv()
|
||
if rerr == io.EOF {
|
||
break
|
||
}
|
||
if rerr != nil {
|
||
tr.emit(node, "model", "error", "compose 图执行", rerr.Error(), time.Since(t0).Milliseconds())
|
||
return
|
||
}
|
||
if chunk.Content == "" {
|
||
continue
|
||
}
|
||
emit(red.Push(chunk.Content))
|
||
}
|
||
emit(red.Flush()) // 吐出暂留尾部
|
||
o.recordAgentOutput(b, produced.String())
|
||
tr.info(node, "system", "compose 图", fmt.Sprintf("%d 段输出 / %d 字(Eino compose 运行时)", chunks, len([]rune(produced.String()))))
|
||
}
|