目录 · 第 27 / 28 章
EinoPart V · 综合实战与工程化

27综合大项目(生产级 Agent 应用)

chatwitheino:deep.NewTyped 叠中间件 + 审批中断 + TurnLoop 服务化,并做一次全书复盘。

quickstart/chatwitheino/agent.go:39quickstart/chatwitheino/agent.go:122quickstart/chatwitheino/server/server.go:570quickstart/chatwitheino/main.go:69

综合实战:读懂一个生产级 Agent 应用,并复盘全书

我们走到了最后一章之一。前面二十多章,你从 ADK 的世界观出发,理解了 Agent 如何被组合、被赋能、被编排,又下沉到 compose 引擎看清了这一切背后的图运行时。上一章我们把一条 RAG 管线读成了一张 Workflow。现在,我们把最后一块拼图放上桌面——读透 chatwitheino 这个完整可跑的 Agent 应用:它用 deep 预制件搭起一个能读文档、答问题的智能体,用中间件叠上”技能/审批/安全工具”三层能力,用 TurnLoop 把它接成一个支持人工审批、可中断可恢复的 SSE 服务。读完它,你会发现全书讲过的每一个机制,都在这里各就各位。

应用骨架:一个 main,三样东西

先看整个应用是怎么起来的。main 按环境变量选择消息类型(AgenticMessage 或普通 Message),然后走进 runTyped(quickstart/chatwitheino/main.go:38)。runTyped 只干三件事(quickstart/chatwitheino/main.go:69):

func main() {
ctx := context.Background()
// setup cozeloop tracing (optional)
// COZELOOP_WORKSPACE_ID=your workspace id
// COZELOOP_API_TOKEN=your token
cozeloopApiToken := os.Getenv("COZELOOP_API_TOKEN")
cozeloopWorkspaceID := os.Getenv("COZELOOP_WORKSPACE_ID")
if cozeloopApiToken != "" && cozeloopWorkspaceID != "" {
client, err := cozeloop.NewClient(
cozeloop.WithAPIToken(cozeloopApiToken),
cozeloop.WithWorkspaceID(cozeloopWorkspaceID),
)
if err != nil {
log.Fatalf("cozeloop.NewClient failed: %v", err)
}
defer func() {
time.Sleep(5 * time.Second)
client.Close(ctx)
}()
callbacks.AppendGlobalHandlers(clc.NewLoopHandler(client))
}
switch msgops.KindFromEnv() {
case msgops.KindAgentic:
runTyped[*schema.AgenticMessage](ctx)
default:
runTyped[*schema.Message](ctx)
}
}
func runTyped[M adk.MessageType](ctx context.Context) {
agent, err := buildAgentTyped[M](ctx)
if err != nil {
log.Fatalf("failed to build agent: %v", err)
}
checkpointStore := adkstore.NewInMemoryStore()
sessionDir := msgops.DefaultSessionDir(msgops.KindOf[M]())
log.Printf("message kind: %s", msgops.KindOf[M]())
log.Printf("session dir: %s", sessionDir)
workspaceDir := os.Getenv("WORKSPACE_DIR")
if workspaceDir == "" {
workspaceDir = "./data/workspace"
}
store, err := mem.NewStore[M](sessionDir)
if err != nil {
log.Fatalf("failed to create session store: %v", err)
}
port := os.Getenv("PORT")
if port == "" {
port = "8080"
}
projectRoot := os.Getenv("PROJECT_ROOT")
if projectRoot == "" {
// Default: the directory from which the binary is run.
// Override with PROJECT_ROOT=/path/to/repo to give the agent full codebase access.
if cwd, err := os.Getwd(); err == nil {
projectRoot = cwd
}
}
if abs, err := filepath.Abs(projectRoot); err == nil {
projectRoot = abs
}
log.Printf("project root: %s", projectRoot)
// EXAMPLES_DIR points to the root of the eino-examples repository.
// Defaults to PROJECT_ROOT/examples if that directory exists, otherwise PROJECT_ROOT.
examplesDir := os.Getenv("EXAMPLES_DIR")
if examplesDir == "" {
candidate := filepath.Join(projectRoot, "examples")
if fi, err := os.Stat(candidate); err == nil && fi.IsDir() {
examplesDir = candidate
} else {
// … 省略 21 行;完整声明 L69–137,点击上方「浏览完整文件」
func runTyped[M adk.MessageType](ctx context.Context) {
agent, err := buildAgentTyped[M](ctx) // 1. 构建 agent
checkpointStore := adkstore.NewInMemoryStore() // 2. 中断-恢复用的 checkpoint 存储
store, err := mem.NewStore[M](sessionDir) // 3. 会话历史(JSONL 落盘)
srv := server.New[M](server.Config[M]{
Agent: agent,
CheckPointStore: checkpointStore,
Store: store,
// ...
})
srv.Spin()
}

一个 agent、一个 checkpoint 存储、一个会话存储,注入 server 后 Spin 起 HTTP 服务(quickstart/chatwitheino/server/server.go:182)。注意这里的分层:agent 是”大脑”,server 是”外壳”,checkpoint 是”可以随时暂停/续跑的存档”——三者通过接口注入拼在一起,而不是揉成一团。这正是全书反复出现的组合姿态。

// Spin starts the HTTP server (blocking).
func (s *Server[M]) Spin() {
h := hserver.Default(hserver.WithHostPorts(":" + s.cfg.Port))
h.GET("/", func(ctx context.Context, c *app.RequestContext) {
data, err := os.ReadFile("static/index.html")
if err != nil {
c.JSON(consts.StatusNotFound, map[string]string{"error": "index.html not found"})
return
}
c.Data(consts.StatusOK, "text/html; charset=utf-8", data)
})
h.POST("/sessions", func(ctx context.Context, c *app.RequestContext) {
id := uuid.New().String()
if _, err := s.cfg.Store.GetOrCreate(id); err != nil {
c.JSON(consts.StatusInternalServerError, map[string]string{"error": err.Error()})
return
}
c.JSON(consts.StatusOK, map[string]string{"id": id})
})
h.GET("/sessions", func(ctx context.Context, c *app.RequestContext) {
metas, err := s.cfg.Store.List()
if err != nil {
c.JSON(consts.StatusInternalServerError, map[string]string{"error": err.Error()})
return
}
if metas == nil {
metas = []mem.SessionMeta{}
}
c.JSON(consts.StatusOK, metas)
})
h.DELETE("/sessions/:id", func(ctx context.Context, c *app.RequestContext) {
id := c.Param("id")
// Stop any running loop for this session.
ts := s.getTurnState(id)
ts.mu.Lock()
if ts.loop != nil {
ts.loop.Stop(adk.WithImmediate())
ts.loop = nil
}
ts.mu.Unlock()
if err := s.cfg.Store.Delete(id); err != nil {
c.JSON(consts.StatusInternalServerError, map[string]string{"error": err.Error()})
return
// … 省略 27 行;完整声明 L181–255,点击上方「浏览完整文件」

deep.NewTyped:把能力一层层叠到 agent 上

buildAgentTyped 是全书组合哲学的一次集中展示。它先造 ChatModel 和一个本地文件后端,再把上一章的 RAG 工具建出来,然后用中间件把三层能力叠上去(quickstart/chatwitheino/agent.go:39):

var handlers []adk.TypedChatModelAgentMiddleware[M]
if skillsDir, ok := resolveSkillsDir(); ok {
// 可选:文件系统技能中间件
handlers = append(handlers, skillMiddleware)
}
handlers = append(handlers,
newApprovalMiddleware[M](), // 审批闸门(下面详解)
helpers.NewSafeToolMiddleware[M](), // 把工具报错转成字符串而非中断
)
cfg := &deep.TypedConfig[M]{
Name: "ChatWithEinoAgent",
ChatModel: cm,
Backend: backend,
MaxIteration: 50,
Handlers: handlers, // 三层中间件
ToolsConfig: /* Tools: []tool.BaseTool{ragTool} */,
}
helpers.ApplyMessageModelRetry(cfg) // 再叠一层:模型限流自动重试
return deep.NewTyped[M](ctx, cfg)

看这个叠法:业务只有”一个 RAG 工具”,但围绕它,Handlers 数组依次叠上了技能、审批、安全工具三层中间件,ApplyMessageModelRetry 又给模型调用包上一层”429 自动重试”(quickstart/chatwitheino/helpers/retry.go:32)。每一层都是从外部包裹上去的切面,agent 的核心逻辑一行没动。 这就是 Ch10–11 讲的”能力靠中间件注入”在真实应用里的样子:要加能力,就往数组里再 append 一个 handler。

// ApplyMessageModelRetry enables model-call retries for transient rate-limit
// errors.
func ApplyMessageModelRetry[M adk.MessageType](cfg *deep.TypedConfig[M]) {
cfg.ModelRetryConfig = &adk.TypedModelRetryConfig[M]{
MaxRetries: 5,
IsRetryAble: func(_ context.Context, err error) bool {
return strings.Contains(err.Error(), "429") ||
strings.Contains(err.Error(), "Too Many Requests") ||
strings.Contains(err.Error(), "qpm limit")
},
}
}

SafeToolMiddleware 是个小而关键的设计(quickstart/chatwitheino/helpers/middleware.go:34):它把工具执行时的普通错误转成一段文字返回给模型(而不是让整个 agent 崩掉),但对”中断类错误”放行——因为那不是失败,是要暂停等人(quickstart/chatwitheino/helpers/middleware.go:44)。一个工具挂了,模型应该”看到错误、换个思路”,而不是让整轮对话终止。

// NewSafeToolMiddleware converts tool errors into error-message strings so that
// a non-zero exit code or mid-stream failure is returned to the model as a
// readable tool result instead of aborting the agent pipeline.
func NewSafeToolMiddleware[M adk.MessageType]() adk.TypedChatModelAgentMiddleware[M] {
return &safeToolMiddleware[M]{
TypedBaseChatModelAgentMiddleware: &adk.TypedBaseChatModelAgentMiddleware[M]{},
}
}
func (m *safeToolMiddleware[M]) WrapInvokableToolCall(
_ context.Context,
endpoint adk.InvokableToolCallEndpoint,
_ *adk.ToolContext,
) (adk.InvokableToolCallEndpoint, error) {
return func(ctx context.Context, args string, opts ...tool.Option) (string, error) {
result, err := endpoint(ctx, args, opts...)
if err != nil {
if _, ok := compose.IsInterruptRerunError(err); ok {
return "", err
}
return fmt.Sprintf("[tool error] %v", err), nil
}
return result, nil
}, nil
}

🔑 本章的设计钥匙

buildAgentTyped 是全书组合哲学的一次总检阅:它没有为”带审批的文档问答 agent”发明任何新类型,而是把一个 RAG 工具、一串中间件、一个重试包装,像叠积木一样拼进 deep.TypedConfig。能力不是继承来的,是一层层 append 上去的切面;要拿掉审批,就从数组里删一行。这印证了 Eino 最深的一条主张——Agent 不是一种特殊的运行时,而是”模型 + 工具 + 一叠中间件”的组合。当”编排”这一层被打磨到足够通用,“如何造一个生产级 agent”就退化成了”往数组里放哪几个 handler”。这,就是 Eino 的 Agent 设计哲学。

人在环中:一次中断,贯穿全书

这个应用最值得学的动作,是给 answer_from_document 这个 RAG 工具装了一道人工审批闸门:在真的去翻文档之前,先停下来问用户”我可以查这份文档吗?”。看 approvalMiddleware.WrapInvokableToolCall(quickstart/chatwitheino/agent.go:122):

func (m *approvalMiddleware[M]) WrapInvokableToolCall(...) (adk.InvokableToolCallEndpoint, error) {
if tCtx.Name != "answer_from_document" {
return endpoint, nil // 其他工具原样放行
}
return func(ctx context.Context, args string, opts ...tool.Option) (string, error) {
wasInterrupted, _, storedArgs := tool.GetInterruptState[string](ctx)
if !wasInterrupted {
// 第一次进来:还没审批 → 中断,把工具参数存进 checkpoint
return "", tool.StatefulInterrupt(ctx, &commontool.ApprovalInfo{
ToolName: tCtx.Name, ArgumentsInJSON: args,
}, args)
}
// 恢复回来:读出用户的审批结果
isTarget, hasData, data := tool.GetResumeContext[*commontool.ApprovalResult](ctx)
if isTarget && hasData {
if data.Approved {
return endpoint(ctx, storedArgs, opts...) // 批准 → 真的执行工具
}
return fmt.Sprintf("tool '%s' disapproved", tCtx.Name), nil // 拒绝 → 返回一句话
}
// ...
}, nil
}

这是全书中断-恢复机制的一次总演出:第一次调用 answer_from_document 时,中间件不执行工具,而是用 tool.StatefulInterrupt(Ch13)把工具参数存进 checkpoint 并中断;等用户回复”批准/拒绝”,再用 tool.GetResumeContext 读出结果——批准就真的执行 RAG 工具,拒绝就返回一句说明。这套机制,你在 Ch13(HITL 初见)、Ch23(checkpoint 恢复)、Ch24(GraphTool 的嵌套中断)、上一章(RAG 工具内置 checkpoint)已经反复见过。中断-恢复不是某个 agent 的特性,是引擎给所有工具/图的通用能力。

TurnLoop:把一个 agent 接成一个会审批的服务

难点在于:HTTP 是一问一答的,而”中断-等人-恢复”是跨越多个请求的。chatwitheinoadk.TurnLoop 架起这座桥。newLoop 把三个回调拼成一个循环(quickstart/chatwitheino/server/server.go:570):

cfg := adk.TurnLoopConfig[*ChatItem, M]{
GenInput: s.makeGenInput(sess, sessionID), // 从会话历史造 agent 输入
PrepareAgent: s.makePrepareAgent(), // 提供 agent
OnAgentEvents: s.makeOnAgentEvents(sess, sessionID), // 把 agent 事件桥到 SSE
}
if s.cfg.CheckPointStore != nil {
cfg.Store = s.cfg.CheckPointStore
cfg.CheckpointID = sessionID
cfg.GenResume = s.makeGenResume() // 恢复时把审批结果喂回去
}
return adk.NewTurnLoop(cfg)

把三条 HTTP 路由和 TurnLoop 对上,整个人工审批闭环就清楚了:

  • POST /chat(quickstart/chatwitheino/server/server.go:283):用户发消息 → 推入一个 ChatItemGenInput 把它连同历史造成 agent 输入(EnableStreaming: true,quickstart/chatwitheino/server/server.go:587)→ agent 跑到 answer_from_document 时中断 → OnAgentEvents 把”需要审批”的事件通过 channel 桥给 SSE handler 推给前端(quickstart/chatwitheino/server/server.go:643)。整轮的中间消息落盘进会话历史,中断 ID 记在 session 上。
// handleChat handles a new chat message. It creates or reuses a TurnLoop for the session.
// If a loop is already running (busy), it pushes with preempt to cancel the current turn.
func (s *Server[M]) handleChat(ctx context.Context, c *app.RequestContext) {
id := c.Param("id")
body, _ := c.Body()
var req chatRequest
if err := json.Unmarshal(body, &req); err != nil || req.Message == "" {
c.JSON(consts.StatusBadRequest, map[string]string{"error": "message is required"})
return
}
log.Printf("[chat] session=%s msg=%q", id, req.Message)
sess, err := s.cfg.Store.GetOrCreate(id)
if err != nil {
c.JSON(consts.StatusInternalServerError, map[string]string{"error": err.Error()})
return
}
item := &ChatItem{Query: req.Message}
ts := s.getTurnState(id)
// Each handler gets its own local iterReady channel reference and a
// handlerDone channel. This avoids races when multiple preempts replace
// the channels on ts concurrently.
var localIterReady chan iterEnvelope[M]
var localHandlerDone chan struct{}
ts.mu.Lock()
if ts.loop != nil {
// Loop exists — try to push with preempt (AfterToolCalls).
loop := ts.loop
log.Printf("[chat] session=%s preempting current turn", id)
// Signal any previous handler waiting on iterReady to bail.
if ts.handlerDone != nil {
close(ts.handlerDone)
}
ts.iterReady = make(chan iterEnvelope[M], 1)
ts.iterDone = make(chan iterResult[M], 1)
ts.handlerDone = make(chan struct{})
localIterReady = ts.iterReady
localHandlerDone = ts.handlerDone
ts.mu.Unlock()
ok, _ := loop.Push(item, adk.WithPreempt[*ChatItem, M](adk.AfterToolCalls))
if !ok {
// Loop already stopped (e.g. error on previous turn) — create new one.
// … 省略 95 行;完整声明 L281–423,点击上方「浏览完整文件」
// makeGenInput returns the GenInput callback. It builds agent messages from
// session history + workspace context.
func (s *Server[M]) makeGenInput(sess *mem.Session[M], sessionID string) func(ctx context.Context, loop *adk.TurnLoop[*ChatItem, M], items []*ChatItem) (*adk.GenInputResult[*ChatItem, M], error) {
return func(ctx context.Context, loop *adk.TurnLoop[*ChatItem, M], items []*ChatItem) (*adk.GenInputResult[*ChatItem, M], error) {
// Find the first item with a query.
var consumed []*ChatItem
var remaining []*ChatItem
var queryItem *ChatItem
for _, item := range items {
if queryItem == nil && item.Query != "" {
queryItem = item
consumed = append(consumed, item)
} else {
remaining = append(remaining, item)
}
}
if queryItem == nil {
// No query items — stop the loop.
loop.Stop(adk.WithStopCause("no query items"))
return &adk.GenInputResult[*ChatItem, M]{
Input: &adk.TypedAgentInput[M]{Messages: []M{msgops.NewUser[M]("done")}},
Remaining: items,
}, nil
}
// Persist the user message NOW — GenInput fires only after any previous
// turn's OnAgentEvents has finished persisting its intermediates, so the
// session history order is guaranteed correct.
userMsg := msgops.NewUser[M](queryItem.Query)
if appendErr := sess.Append(userMsg); appendErr != nil {
log.Printf("warn: failed to persist user message: %v", appendErr)
}
history := sess.GetMessages()
runMessages := s.buildRunMessages(sessionID, history)
log.Printf("[genInput] session=%s query=%q messages=%d", sessionID, queryItem.Query, len(runMessages))
return &adk.GenInputResult[*ChatItem, M]{
Input: &adk.TypedAgentInput[M]{
Messages: runMessages,
EnableStreaming: true,
},
Consumed: consumed,
Remaining: remaining,
}, nil
}
}
// makeOnAgentEvents returns the OnAgentEvents callback — the bridge between
// the TurnLoop and the HTTP handler.
func (s *Server[M]) makeOnAgentEvents(sess *mem.Session[M], sessionID string) func(ctx context.Context, tc *adk.TurnContext[*ChatItem, M], events *adk.AsyncIterator[*adk.TypedAgentEvent[M]]) error {
return func(ctx context.Context, tc *adk.TurnContext[*ChatItem, M], events *adk.AsyncIterator[*adk.TypedAgentEvent[M]]) error {
ts := s.getTurnState(sessionID)
history := sess.GetMessages()
// Snapshot bridge channels under lock to avoid races with handleChat
// which may recreate them for a preempt.
ts.mu.Lock()
ready := ts.iterReady
done := ts.iterDone
ts.mu.Unlock()
// Send the iterator to the HTTP handler. Include the done channel
// so the handler replies to THIS invocation, not a future one.
select {
case ready <- iterEnvelope[M]{events: events, history: history, done: done}:
case <-ctx.Done():
return ctx.Err()
}
// Wait for the HTTP handler to finish draining. Also select on ctx.Done
// to avoid hanging when a preempt supersedes the handler — in that case
// the old handler bails via handlerDone and nobody sends to our done channel.
var result iterResult[M]
select {
case result = <-done:
case <-ctx.Done():
return ctx.Err()
}
// Persist all intermediate messages (assistant text+tool calls, tool results).
// The intermediates already include the final assistant text message if any,
// so we don't need to persist lastContent separately.
for _, msg := range result.intermediates {
if appendErr := sess.Append(msg); appendErr != nil {
log.Printf("warn: failed to persist intermediate message: %v", appendErr)
}
}
if result.interruptID != "" {
sess.SetPendingInterruptID(result.interruptID)
sess.SetMsgIdx(result.msgIdx)
return errInterrupted
}
return result.err
}
// … 省略 1 行;完整声明 L641–689,点击上方「浏览完整文件」
  • POST /approve(quickstart/chatwitheino/server/server.go:427):用户点”批准/拒绝” → 新建一个带 CheckPointID 的 loop → GenResume 把审批结果通过 ResumeParams.Targets 精确投给那个中断点(quickstart/chatwitheino/server/server.go:692)→ 工具从断点恢复执行 → 结果再流式推回。
// handleApprove resumes an interrupted agent run with the user's approval decision.
// Creates a new TurnLoop with checkpoint/resume to continue from the interrupt.
func (s *Server[M]) handleApprove(ctx context.Context, c *app.RequestContext) {
id := c.Param("id")
sess, err := s.cfg.Store.GetOrCreate(id)
if err != nil {
c.JSON(consts.StatusInternalServerError, map[string]string{"error": err.Error()})
return
}
interruptID := sess.GetPendingInterruptID()
if interruptID == "" {
c.JSON(consts.StatusBadRequest, map[string]string{"error": "no pending interrupt for this session"})
return
}
body, _ := c.Body()
var req approveRequest
if err := json.Unmarshal(body, &req); err != nil {
c.JSON(consts.StatusBadRequest, map[string]string{"error": "invalid request body"})
return
}
var reason *string
if req.Reason != "" {
reason = &req.Reason
}
result := &commontool.ApprovalResult{Approved: req.Approved, DisapproveReason: reason}
// Clear the pending interrupt so a double-approve returns 400.
sess.SetPendingInterruptID("")
log.Printf("[approve] session=%s interruptID=%s approved=%v", id, interruptID, req.Approved)
// Create a new loop with checkpoint resume.
ts := s.getTurnState(id)
ts.mu.Lock()
// Clear any old loop.
if ts.loop != nil {
ts.loop.Stop(adk.WithImmediate())
}
// Signal any previous handler to bail.
if ts.handlerDone != nil {
close(ts.handlerDone)
}
loop := s.newLoop(sess, id, true)
ts.loop = loop
// … 省略 70 行;完整声明 L425–542,点击上方「浏览完整文件」
// makeGenResume returns the GenResume callback for interrupt/resume.
func (s *Server[M]) makeGenResume() func(ctx context.Context, loop *adk.TurnLoop[*ChatItem, M], canceledItems, unhandledItems, newItems []*ChatItem) (*adk.GenResumeResult[*ChatItem, M], error) {
return func(ctx context.Context, loop *adk.TurnLoop[*ChatItem, M], canceledItems, unhandledItems, newItems []*ChatItem) (*adk.GenResumeResult[*ChatItem, M], error) {
// Find the approval item in newItems.
var approvalItem *ChatItem
for _, item := range newItems {
if item.ApprovalResult != nil {
approvalItem = item
break
}
}
if approvalItem == nil {
return nil, errors.New("no approval item found for resume")
}
return &adk.GenResumeResult[*ChatItem, M]{
ResumeParams: &adk.ResumeParams{
Targets: map[string]any{approvalItem.InterruptID: approvalItem.ApprovalResult},
},
Consumed: canceledItems,
Remaining: unhandledItems,
}, nil
}
}
  • POST /abort(quickstart/chatwitheino/server/server.go:545):loop.Stop(adk.WithImmediate()) 立即掐断当前轮。
// handleAbort immediately stops the current TurnLoop for a session.
func (s *Server[M]) handleAbort(_ context.Context, c *app.RequestContext) {
id := c.Param("id")
ts := s.getTurnState(id)
ts.mu.Lock()
loop := ts.loop
ts.loop = nil
ts.mu.Unlock()
if loop == nil {
c.JSON(consts.StatusOK, map[string]string{"status": "no active loop"})
return
}
log.Printf("[abort] session=%s stopping loop immediately", id)
loop.Stop(adk.WithImmediate())
loop.Wait()
log.Printf("[abort] session=%s loop stopped", id)
c.JSON(consts.StatusOK, map[string]string{"status": "aborted"})
}

OnAgentEvents 是这座桥的关键(quickstart/chatwitheino/server/server.go:643):它把 agent 内部的异步事件流,通过一对 channel 交给正在 hold 住 SSE 连接的 HTTP handler,handler 边收边推给浏览器,推完再把结果(是否中断、中断 ID)送回来。TurnLoop 负责”agent 的一生”,SSE handler 负责”这一次 HTTP 请求”,两者用 channel 解耦——这正是 Ch28 要总结的”用 channel 做调度缝”的一个活标本。

sequenceDiagram
  participant U as 前端
  participant S as Server (TurnLoop)
  participant A as Agent + RAG 工具
  participant CP as CheckPointStore
  U->>S: POST /chat "总结这份文档"
  S->>A: GenInput → 运行 agent
  A->>A: 想调用 answer_from_document
  A->>CP: StatefulInterrupt 存档参数
  A-->>S: 事件:需要审批
  S-->>U: SSE 推送"待审批"
  U->>S: POST /approve {approved:true}
  S->>CP: 按 CheckpointID 读档
  S->>A: GenResume 投递审批结果
  A->>A: 从断点恢复,执行 RAG
  A-->>S: 事件:最终答案
  S-->>U: SSE 推送答案

人工审批闭环:两次 HTTP 请求跨越一次中断

记忆:文件式会话,喂给模型的只是一个窗口

对话历史用一个基于 JSONL 的 mem.Store:每条消息 Append 追加一行(quickstart/chatwitheino/mem/store.go:82),整段历史落盘持久化。但真正喂回给 agent 的,是 GetMessages 读出的历史再拼上一段上下文提示(quickstart/chatwitheino/mem/store.go:105)。这印证了 Ch08 的记忆观:持久化要全,喂模型要省——落盘是完整的 JSONL,而每轮 GenInput 只把需要的历史拼进输入。会话状态(待处理的中断 ID、消息游标)也挂在 session 上,正是它让”跨请求恢复”成为可能。

// Append adds a message to memory and persists it to disk.
func (s *Session[M]) Append(msg M) error {
s.mu.Lock()
defer s.mu.Unlock()
msg = msgops.NormalizeForSession(msg)
s.messages = append(s.messages, msg)
data, err := json.Marshal(msg)
if err != nil {
return err
}
f, err := os.OpenFile(s.filePath, os.O_APPEND|os.O_WRONLY, 0o644)
if err != nil {
return err
}
defer f.Close()
_, err = fmt.Fprintf(f, "%s\n", data)
return err
}
// GetMessages returns a snapshot of all messages.
func (s *Session[M]) GetMessages() []M {
s.mu.Lock()
defer s.mu.Unlock()
result := make([]M, len(s.messages))
copy(result, s.messages)
return result
}

从 ch01 到完整应用:一条学习阶梯

chatwitheino 还贴心地留了一条阶梯:cmd/ch01cmd/ch10 是从”最小可跑”到”完整应用”的十级台阶。cmd/ch01/main.go 只有几十行——建一个 ChatModel、Stream 一次、循环收 delta 打印(quickstart/chatwitheino/cmd/ch01/main.go:38),就是本书 Part I 流式那一章的最小复现。顺着 ch01 往上爬到根目录这个带审批、带记忆、带 SSE 的完整应用,你会亲眼看到本书的机制是怎么一级一级叠上来的。这也是我们推荐你动手的起点:从 ch01 开始跑,一路读到 main.go。

func main() {
var instruction string
flag.StringVar(&instruction, "instruction", "You are a helpful assistant.", "")
flag.Parse()
query := strings.TrimSpace(strings.Join(flag.Args(), " "))
if query == "" {
_, _ = fmt.Fprintln(os.Stderr, "usage: go run ./cmd/ch01 -- \"your question\"")
os.Exit(2)
}
ctx := context.Background()
switch msgops.KindFromEnv() {
case msgops.KindAgentic:
runTyped[*schema.AgenticMessage](ctx, instruction, query)
default:
runTyped[*schema.Message](ctx, instruction, query)
}
}

全书复盘:Eino 的 Agent 设计哲学

读完 chatwitheino,正好回望我们走过的路。全书想让你带走的,不是 API 清单,而是几条贯穿始终的设计信念:

  1. Agent 是”模型 + 工具 + 中间件”的组合,不是特殊运行时。(Part I / 本章)deep.NewTyped 把一个 RAG 工具和一叠 handler 拼起来,就是一个生产级 agent。底下还是那个 compose 引擎。

  2. 组合优于继承,一切皆可嵌套。(Ch14 / Ch20 / 上一章)一整张 RAG Workflow 能被 graphtool 包成一个工具;这个工具又被挂进 agent。系统靠”套娃”生长,而不是靠庞大的类层次。

  3. 能力靠切面注入,不靠侵入。(Ch10–11 / Ch26 / 本章)技能、审批、安全工具、重试、可观测,都是中间件/callback 从外部包裹上去的,业务逻辑一行不改。要加能力,append 一个 handler 即可。

  4. 控制权显式流转。(Ch16–17 / 本章)要不要执行敏感工具,不靠隐式魔法,而是把决定权中断出去交给人,再由审批结果精确投回那个中断点。意图是数据,恢复是普通调用。

  5. 中断-恢复是引擎级能力。(Ch13 / Ch23 / 本章)只要状态可序列化,任何工具/图都能在任意点停下、存档、日后带反馈恢复。HITL 因此不是功能,而是免费的副产品。

  6. 类型是脚手架,泛型贯穿始终。(Ch20 / 本章)TypedAgent[M]BaseModel[M]、四种执行范式让编排在编译期就类型安全,AgenticMessage 与普通 Message 靠同一份泛型代码自动适配。

  7. 持久化与运行时解耦。(Ch08 / Ch23 / 本章)会话历史(JSONL)和中断存档(checkpoint)是两个独立的存储,通过接口注入 server——换存储只是换一个实现。

如果只记一句话:Eino 把”编排”打磨成了一等公民,于是 Agent 不再需要被特殊定义——它只是”模型 + 工具 + 一叠中间件”,跑在一个能随时暂停续跑的循环里。 从第一章 ADK 的世界观,到这一章 chatwitheino 的完整应用,你看到的始终是这同一个信念的不同投影。

本章小结

  • chatwitheino 的骨架是”一个 agent + 一个 checkpoint 存储 + 一个会话存储”,注入 server 后 Spin 起服务(quickstart/chatwitheino/main.go:69)。
  • deep.NewTyped 把一个 RAG 工具和技能/审批/安全工具三层中间件、加上模型重试,像叠积木一样拼成 agent(quickstart/chatwitheino/agent.go:39)——能力是 append 上去的切面。
  • 审批中间件用 tool.StatefulInterrupt + GetResumeContextanswer_from_document 装了一道人工闸门(quickstart/chatwitheino/agent.go:122),是引擎级中断-恢复能力的应用。
  • TurnLoop(quickstart/chatwitheino/server/server.go:570)用 GenInput/OnAgentEvents/GenResume 三个回调,把 /chat/approve 两次 HTTP 请求跨越一次中断,拼成一个可审批、可恢复的 SSE 服务。
  • 记忆是文件式 JSONL(quickstart/chatwitheino/mem/store.go:82):落盘要全,喂模型只取需要的窗口(quickstart/chatwitheino/mem/store.go:105)。
  • cmd/ch01ch10 是从最小流式到完整应用的学习阶梯(quickstart/chatwitheino/cmd/ch01/main.go:38),推荐从 ch01 读起。
  • 设计钥匙:Agent 不是特殊运行时,而是”模型 + 工具 + 中间件”的组合——这是贯穿全书的 Eino Agent 设计哲学。

下一章是全书的收束:我们跳出具体应用,看这整套 Agent Runtime 到底是如何长在 Go 语言的原生机制上的。

源码

正在读取完整文件…