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 是一问一答的,而”中断-等人-恢复”是跨越多个请求的。chatwitheino 用 adk.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):用户发消息 → 推入一个ChatItem→GenInput把它连同历史造成 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/ch01–cmd/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 清单,而是几条贯穿始终的设计信念:
-
Agent 是”模型 + 工具 + 中间件”的组合,不是特殊运行时。(Part I / 本章)
deep.NewTyped把一个 RAG 工具和一叠 handler 拼起来,就是一个生产级 agent。底下还是那个 compose 引擎。 -
组合优于继承,一切皆可嵌套。(Ch14 / Ch20 / 上一章)一整张 RAG Workflow 能被
graphtool包成一个工具;这个工具又被挂进 agent。系统靠”套娃”生长,而不是靠庞大的类层次。 -
能力靠切面注入,不靠侵入。(Ch10–11 / Ch26 / 本章)技能、审批、安全工具、重试、可观测,都是中间件/callback 从外部包裹上去的,业务逻辑一行不改。要加能力,append 一个 handler 即可。
-
控制权显式流转。(Ch16–17 / 本章)要不要执行敏感工具,不靠隐式魔法,而是把决定权中断出去交给人,再由审批结果精确投回那个中断点。意图是数据,恢复是普通调用。
-
中断-恢复是引擎级能力。(Ch13 / Ch23 / 本章)只要状态可序列化,任何工具/图都能在任意点停下、存档、日后带反馈恢复。HITL 因此不是功能,而是免费的副产品。
-
类型是脚手架,泛型贯穿始终。(Ch20 / 本章)
TypedAgent[M]、BaseModel[M]、四种执行范式让编排在编译期就类型安全,AgenticMessage与普通Message靠同一份泛型代码自动适配。 -
持久化与运行时解耦。(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+GetResumeContext给answer_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/ch01–ch10是从最小流式到完整应用的学习阶梯(quickstart/chatwitheino/cmd/ch01/main.go:38),推荐从 ch01 读起。- 设计钥匙:Agent 不是特殊运行时,而是”模型 + 工具 + 中间件”的组合——这是贯穿全书的 Eino Agent 设计哲学。
下一章是全书的收束:我们跳出具体应用,看这整套 Agent Runtime 到底是如何长在 Go 语言的原生机制上的。