Skip to content

Agent 模块 #4

Description

@SATA260

Agent 模块

功能职责

Agent 把用户的一句话跑成一次可恢复的对话执行:装上下文、调模型、执行工具、处理审批,并记下事件和用量。

  • 用户发一条消息,创建或排队一次 Run;interrupt 先取消再开,queue 等当前结束再领
  • 按 checkpoint、历史消息、可见工具、冻结目录装出本轮上下文;超预算则压缩
  • 按本次冻结的配置流式调用模型,产出文字、Tool Call 和用量
  • 校验并执行工具;一批待批工具合成一条审批,审完再继续
  • 事件先落库再推送;断线按序号回放
  • 跨 Session 保存目录和专题;写消息时建索引,供按工作区检索
  • 用户可看/删记忆,不可从这边直接写

边界

  • 同一 Session 同时只有一个 active Run
  • 工具之后的下一轮模型调用仍属同一个 Run,不新建 Session
  • 拒绝审批不当整次失败,把失败结果喂回模型
  • 压缩只改默认装载视图,不删原始消息和事件
  • 当前 Session 已冻结的目录前缀,写入新目录后不立刻改
  • 用户侧记忆只能看/删;写入只走记忆工具。记什么、要不要建专题,不由本模块决定
  • 断开事件订阅不取消 Run

内部拆分

会话(Session)

长期对话容器。创建、查询、更新、归档;记住当前 active Run 和事件游标。不管消息正文,不管一次执行怎么跑。

type Session struct {
    ID            string
    TenantID      string
    UserID        string
    AgentID       string
    WorkspaceID   string
    Status        SessionStatus // active | archived
    ActiveRunID   *string       // 当前正在执行的 Run;空表示空闲。
    LastEventSeq  int64         // 已写入的最大事件序号。
    CompactionSeq int64         // 最近一次压缩 checkpoint 的基准序号。
}

func Create(tenantID, userID, agentID, workspaceID string) (Session, error) // 创建会话。
func Get(sessionID string) (Session, error)                                // 按 ID 取会话。
func Update(sessionID, agentID string, status SessionStatus) (Session, error) // 改 Agent 或状态。
func Archive(sessionID string) error                                       // 归档;有 active Run 时不可归档。
func ClaimActiveRun(sessionID, runID string) error                         // 标成当前执行;同时只能有一个。
func ClearActiveRun(sessionID, runID string) error                         // 结束或取消后清空。

消息(Message)

用户、助手、工具消息的写入与查询。给上下文装载提供历史,不决定装哪些、不压缩。

type Message struct {
    ID          string
    SessionID   string
    RunID       *string
    TurnID      *string
    Role        MessageRole     // user | assistant | tool | system
    Content     json.RawMessage // 文本、Tool Call 或 Tool Result。
    Attachments []Attachment
    ToolCalls   []ToolCall
    EventSeq    int64           // 对应事件流里的序号,用来和 SSE 对齐。
}

func Insert(msg Message) (Message, error)                                                 // 写入一条消息。
func List(sessionID string) (messages []Message, asOfEventSeq int64, err error)            // 列出会话消息;asOfEventSeq 给前端对上事件流。
func Get(messageID string) (Message, error)                                               // 按 ID 取消息。
func Delete(sessionID, messageID string) error                                            // 删除一条消息。

Run 控制(Run)

一次用户触发的执行:开始、排队、继续、重试、取消、终态。interrupt 先取消再开新 Run;queue 只入队,当前结束后再领。不装上下文,不调模型。

type Run struct {
    ID               string
    SessionID        string
    TriggerMessageID string            // 触发本次执行的用户消息。
    Mode             AgentMode
    Config           RunConfigSnapshot // 创建时冻结,后续 Turn 只读这份。
    Status           RunStatus
    CurrentTurnID    *string
    StopReason       *StopReason
    CancelRequested  bool
}

type InputMode string // interrupt 先取消再开;queue 只入队。

func Start(sessionID, content string, inputMode InputMode, mode AgentMode) (runID string, err error) // 写入用户消息并创建或排队 Run。
func Continue(runID string) error                                     // 从审批或恢复点继续。
func Retry(runID string) error                                        // 失败后按原配置再跑。
func Cancel(runID string) error                                       // 请求停止,并向模型和工具传播取消。
func Get(runID string) (Run, error)                                   // 按 ID 取 Run。
func Terminate(runID string, status RunStatus, reason StopReason) error // 进入终态并记下原因。

上下文(Context)

按压缩 checkpoint、摘要、之后的消息、可见工具和冻结目录,装出本轮要给模型看的内容,并组装成模型请求。不调用模型,不压缩。

type ContextSnapshot struct {
    SessionID       string
    BaseEventSeq    int64              // 摘要覆盖历史后,从这条之后重新装消息。
    Summary         *CompactionSummary
    Messages        []Message
    Tools           []ToolDefinition
    SystemPrompt    string
    MemoryIndexes   []string           // 冻结的用户目录、工作区目录。
    EstimatedTokens int64
}

type Chat struct {
    SessionID       string
    RunID           string
    TurnID          string
    Model           ModelConfig
    SystemPrompt    string
    Messages        []Message
    Tools           []ToolDefinition
    MaxInputTokens  int64
    MaxOutputTokens int64
}

func Load(run Run, turn Turn, checkpoint *CompactionCheckpoint, messages []Message, tools []ToolDefinition, prompt string) (ContextSnapshot, error) // 装出本轮上下文。
func Build(run Run, turn Turn, snapshot ContextSnapshot) (Chat, error) // 组装成模型请求。

上下文压缩(Compaction)

上下文超 token 预算时生成摘要 checkpoint,再交给上下文重新装载。超限目录改短盖写。只改默认装载视图,不删原始消息和事件。不改当前 Session 已冻结的目录前缀。

type CompactionCheckpoint struct {
    ID           string
    SessionID    string
    BaseEventSeq int64  // 之后的消息才重新装进上下文。
    Summary      string
    CreatedByRun string
}

func Needs(run Run, snapshot ContextSnapshot) bool                                          // 是否超过 token 预算。
func CompactIfNeeded(run Run, turn Turn, snapshot ContextSnapshot) (ContextSnapshot, error) // 超预算才压缩,否则原样返回。
func CompactIndex(model ModelConfig, content string) (string, error)                        // 把超限目录改短。

模型调用(Model)

按本次 Run 冻结的配置发起流式调用,产出文字增量、完整助手消息、Tool Call 和用量。模型按配置创建,不由其他模块注入。不执行工具。

type ModelStreamResult struct {
    Message   Message
    ToolCalls []ToolCall
    Usage     ProviderUsage
}

type UsageRecord struct {
    SessionID   string
    RunID       string
    TurnID      string
    Provider    string
    Model       string
    UsageType   string // generation | compaction
    TotalTokens int64
    Estimated   bool   // Provider 没给精确用量时为估算。
}

func Stream(chat Chat) (ModelStream, error) // 按 chat.Model 创建并发起流式调用。

type ModelStream interface {
    Events() <-chan ModelStreamEvent    // 文字或 Tool Call 增量。
    Result() (ModelStreamResult, error) // 流结束后的完整助手消息、Tool Call、用量。
    Close() error                       // 取消请求并释放连接。
}

func CountTokens(text string) int64 // 估算文本 token 数。

Tool 调用(Tool)

模型给出 Tool Call 之后:校验、按模式过滤可见工具、权限、是否要审批、串行或并行执行、回填结果。本阶段工具:pingmemory_readmemory_writememory_search。不实现文件 / Shell / Git。

运行模式提供 read / write / memory,须覆盖工具声明的全部能力才可调用。记忆工具声明 memory。审批由工具声明,模式决定是否暂停。

type ToolCall struct {
    ID        string
    Name      string
    Arguments json.RawMessage
}

type Invocation struct {
    Calls            []ToolCall
    Mode             ToolExecutionMode // serial | parallel
    FailurePolicy    ToolFailurePolicy
    PermissionPolicy PermissionPolicy
    ApprovalPolicy   ApprovalPolicy
    AgentMode        AgentMode
    ApprovedCallIDs  []string          // 本轮已批准、可直接执行的调用。
    DeniedCallIDs    []string          // 本轮已拒绝,当失败结果回填。
}

type DispatchResult struct {
    Results         []ToolResult
    WaitingApproval bool           // 有待批则整批不执行。
    PendingCalls    []ToolCall
    ApprovalCalls   []ToolCall
}

func VisibleDefinitions(all []ToolDefinition, names []string, mode AgentMode) []ToolDefinition // 按绑定名和模式能力,算出本轮对模型可见的工具。
func Dispatch(inv Invocation) (DispatchResult, error)                                          // 校验、授权、审批判定、执行。
func CancelInFlight()                                                                          // 取消进行中的工具调用。

type Tool interface {
    Definition() ToolDefinition                      // 给模型看的名字、说明、schema、权限。
    Execute(input ToolInput) (ToolResult, error)     // 执行一次调用。
}

审批(Approval)

一次模型回复里的待批工具合成一条审批。Run 暂停;用户一次提交对每条批或拒,全部裁定后再继续。拒绝当工具失败结果喂回模型,不打死 Run。

type Approval struct {
    ID        string
    SessionID string
    RunID     string
    ToolCalls []ApprovalToolCall // 一次模型回复里的待批工具。
    Scope     ApprovalScope      // once | run | session
    Status    ApprovalStatus     // pending | approved | denied | expired
}

type ApprovalDecision struct {
    ToolCallID string
    Status     ApprovalStatus // approved | denied
    Reason     string
}

type ToolCheckpoint struct {
    RunID     string
    Completed []string   // 已执行。
    Approved  []string   // 已批准。
    Denied    []string   // 已拒绝。
    Pending   []ToolCall
    Results   []ToolResult
}

func Create(run Run, calls []ToolCall) (Approval, error)                       // 一批待批工具合成一条审批,并暂停 Run。
func Decide(approvalID string, decisions []ApprovalDecision) (Approval, error) // 一次提交审完再流转。
func HasCheckpoint(runID string) bool                                          // 是否已有可恢复的工具裁决。

记忆(Memory)

跨 Session 的热层目录 + 专题,以及同一工作区内已落库消息的检索。给新 Session 或压缩后提供冻结目录;写消息时建索引。不定义工具,不负责提示词和对话压缩,不自动建专题。

type TextMemory struct {
    Scope      TextMemoryScope // user | workspace
    ScopeID    string
    Kind       TextMemoryKind  // index 为目录,topic 为专题。
    Name       string          // 目录固定为 index。
    Content    string
    ByteLen    int
    OverBudget bool            // 目录超过 200 行或 25KB。
}

type ContextMessage struct {
    ID          string
    WorkspaceID string
    SessionID   string
    RunID       string
    Role        string
    Content     string
}

func Get(key TextMemoryKey) (TextMemory, error)                                            // 读一篇目录或专题。
func Upsert(item TextMemory) (TextMemory, error)                                           // 覆盖写入;用户侧不调,只给记忆工具用。
func Delete(key TextMemoryKey) error                                                       // 删一篇目录或专题。
func List(scope TextMemoryScope, scopeID string) ([]TextMemory, error)                     // 列出某 scope 下的目录和专题。
func SearchMessages(search Search) ([]MessageHit, error)                                   // 按工作区检索已索引消息。
func IndexMessage(msg ContextMessage) error                                                // 写消息时建检索索引。
func FrozenIndexes(session Session) (userIndex, workspaceIndex string)                     // 新 Session 或压缩后装一次,之后复用。

事件(Event)

Session 内严格递增的运行事实。先落库再推送。订阅先回放游标之后的事件,再接实时增量。客户端断开不取消 Run。

type AgentEvent struct {
    EventID   string
    SessionID string
    RunID     string
    TurnID    *string
    Seq       int64      // Session 内严格递增,重连用这个游标。
    Type      EventType
    Payload   json.RawMessage
}

func Append(event AgentEvent) (AgentEvent, error)                           // 先落库再推送。
func Subscribe(sessionID string, afterSeq int64) (<-chan AgentEvent, error) // 先回放 afterSeq 之后,再接实时增量。

流程

用户想发起一个会话,输入文本;模型若要改记忆,再审批。也可以打断、排队、取消,或从断点继续。

// 用户发起会话
Session.Create(...)

// 用户输入文本
Message.Insert(...)
Run.Start(...)
Session.ClaimActiveRun(...)
Event.Append(run.created)

Context.Load(...)
Compaction.CompactIfNeeded(...)
Context.Build(...)
Model.Stream(...)
Message.Insert(...)          // 助手消息
Memory.IndexMessage(...)
Event.Append(assistant.delta)

Tool.Dispatch(...)
Approval.Create(...)         // 需要审批时
Event.Append(tool.approval_required)

// 用户审批
Approval.Decide(...)
Run.Continue(...)
Tool.Dispatch(...)           // 按裁决继续执行
Message.Insert(...)          // 工具结果
Memory.IndexMessage(...)
Event.Append(tool.approval_decided)

// 用户取消
Run.Cancel(...)
Tool.CancelInFlight(...)
Event.Append(run.cancelled)
Session.ClearActiveRun(...)

// 用户从断点继续 / 重试
Run.Continue(...)
Run.Retry(...)

// 用户打开会话或断线重连
Event.Subscribe(...)

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    moduleSingle-module objects, interfaces, and design

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions