Skip to content

插件闸门:go-plugin 事件驱动(子进程 RPC) #16

Description

@SATA260

Parent

#4

What

在 Agent Loop 的闸门上改成 pi 式的 typed func(例如 BeforeRunStart),原操作发生前先调钩子;nil 则与现在一样直接执行。钩子可以放行、改写载荷、或一票否决。

后期用 Compose / Hooks.Then 把实现拼上去:测试注入假函数,插件是常驻子进程,用 go-plugin 的双向 gRPC 把某个 typed func 做成 OnEvent。Loop 出现事件字符串。

钩子能做的事:拦下整轮用户输入(不建 Run)、改即将发给模型的上下文、改 Provider HTTP 的 header/body、在审批和 Execute 之前阻断或改工具参数、在入库前改工具结果、以及插件注册自己的 Tool(Execute 走 RPC,仍走同一条 Dispatch)。

钩子不能做的事:当前端时间线用。PluginEvent 不是 #4AgentEvent。SSE 仍是闸门走完后先落库再推送。

Why

#4 的 Loop 已经闭环,但闸门处都是直接调用:校验完立刻建 Run,Build 完立刻 Stream,组好 HTTP 立刻 Do,权限过了立刻审批并 Execute。只订阅 SSE 看到的是已经发生的事实,拦不住。

所以每处原操作前要有一个可空的函数字段。对齐 pi 的 AgentLoopConfigbeforeToolCalltransformContextonPayload):Loop 只写 if hook != nil,不写 OnEvent("tool_call")。这样测试、指标、插件都能 Compose 上去,不必改 runner。

pkg/agent 必须继续无状态、不依赖 internal / go-plugin。因此类型和 Compose 放在 Agent;internal/plugin 只是这些 func 的一种实现,启动时 hooks.Then(plugin.Hooks())

BeforeInput 必须在 interrupt 取消当前 Run 之前:handled 不应误杀正在跑的 Run。BeforeToolCall 必须在审批之前:危险调用应在打扰用户之前被拦下。BeforeRunStart 每个 Run 一次(lease 之后、进 loop 之前),对齐 pi 的 before_agent_start;每次 LLM 前改 messages 走 TransformContext,不要把二者混在每个 Turn 上。

多实现串行、后者看见前者改过的载荷,才能做策略叠加。未订阅的插件事件不发 RPC,避免把 assistant.delta 打成逐 token 往返。

插件注册的 Tool 必须进现有 Registry 和 VisibleDefinitions:否则模型看不见,也不会经过 BeforeToolCallSetActiveTools 不能摘掉内置 ping / memory_*

How

Agent:typed func + Compose

闸门(等返回值再做原操作)。nil = 跳过,与现在相同:

type BeforeInput func(ctx context.Context, in InputEvent) (InputResult, error)
// Start 校验之后、Cancel / Insert / 建 Run 之前。Action: continue | transform | handled。

type BeforeRunStart func(ctx context.Context, in RunStartEvent) (RunStartResult, error)
// Execute 拿到 lease 之后、进 loop 之前,每个 Run 一次。可改 systemPrompt、追加隐藏消息。

type TransformContext func(ctx context.Context, messages []Message) ([]Message, error)
// 每次 Build 之后、Stream 之前,链式改写 messages。

type BeforeProviderHeaders func(ctx context.Context, header http.Header) error
type BeforeProviderRequest func(ctx context.Context, body []byte) ([]byte, error)
type AfterProviderResponse func(ctx context.Context, status int, header http.Header) error
// 仅 openai。Do 前改请求;Do 后、读 body 前通知。

type BeforeToolCall func(ctx context.Context, in BeforeToolCallEvent) (BeforeToolCallResult, error)
// Validate + CheckPermission 之后、审批之前。Block 则失败 Result,不审批、不 Execute。

type AfterToolCall func(ctx context.Context, in AfterToolCallEvent) (AfterToolCallResult, error)
// Execute 之后、入库之前,按字段覆盖 Output / Error / Success。

通知(都调,忽略返回值):OnRunStartOnTurnStartOnToolExecutionStartOnTurnEndOnRunEndOnRunEnd 打在 Terminate 入口,completed / failed / cancelled / timeout 都覆盖。

拼接头,便于后期接插件或别的实现:

func ComposeBeforeRunStart(fns ...BeforeRunStart) BeforeRunStart
// 按参数顺序串行。后者看见前者改过的 Event。任一 Block 立刻停,不再调后面的 fn。

type Hooks struct {
    BeforeInput            BeforeInput
    BeforeRunStart         BeforeRunStart
    TransformContext       TransformContext
    BeforeProviderHeaders  BeforeProviderHeaders
    BeforeProviderRequest  BeforeProviderRequest
    AfterProviderResponse  AfterProviderResponse
    BeforeToolCall         BeforeToolCall
    AfterToolCall          AfterToolCall
    OnRunStart             OnRunStart
    OnTurnStart            OnTurnStart
    OnToolExecutionStart   OnToolExecutionStart
    OnTurnEnd              OnTurnEnd
    OnRunEnd               OnRunEnd
}

func (h Hooks) Then(next Hooks) Hooks // 逐字段 Compose

Chat 只留 OnHeaders / OnPayload / OnResponseInvocation 只留 BeforeCall / AfterCall。Runtime 在 Stream / Dispatch 前从 Hooks 填进去。

串联(直接调用改成先钩子再原操作)

// 用户发一条消息。校验 session / content 之后,先 BeforeInput,再决定要不要建 Run。
out, err := hooks.BeforeInput(ctx, InputEvent{Content: content})
// handled:不 Cancel 当前 Run,不 Insert 用户消息,不 Start Run,HTTP 200 且 run_id 为空。
// transform:content = out.Content 再往下走。
// continue 或 hook == nil:content 不变。

Run.Cancel(...)                         // 仅 interrupt;handled 时禁止走到这里
Message.Insert(user, content)
Run.Start(...)                          // 写 queued Run
Session.ClaimActiveRun(...)
Event.Append("run.created")             // 仍是 #4 的前端事件,不是 PluginEvent

// Worker 领取后。lease 拿到之后、进 loop 之前:先 BeforeRunStart,再通知 OnRunStart。
Runtime.Execute(...)
out, err = hooks.BeforeRunStart(ctx, RunStartEvent{...}) // 每个 Run 一次
hooks.OnRunStart(ctx, ...)                               // 通知,忽略返回值

// 每个 Turn 都通知;新建 Turn 会写 turn.started,resume 跳过模型时也要调 OnTurnStart。
Turn.Ensure(...)
Event.Append("turn.started")            // 仅新建 Turn 时
hooks.OnTurnStart(ctx, ...)             // 新建与 resume 都调

// 装上下文、压缩,仍按 #4。压缩闸门第一期不做。
snapshot := Context.Load(...)
snapshot = Compaction.CompactIfNeeded(...)

// 每次调模型:Build 之后用 TransformContext 改 messages,再 Stream。
chat := Context.Build(snapshot)
chat.Messages, err = hooks.TransformContext(ctx, chat.Messages)
stream := Model.Stream(chat)
// Stream 内部(仅 openai;fake 不调这三条):
//   hooks.BeforeProviderHeaders(ctx, req.Header)
//   body = hooks.BeforeProviderRequest(ctx, body)
//   resp = HTTP.Do(req)                    // 必须用改过的 header / body
//   hooks.AfterProviderResponse(ctx, resp.Status, resp.Header)

Message.Insert(assistant)
Event.Append("assistant.delta")         // 不接到 Hooks,默认不向插件广播

// 工具:Lookup / Validate / Permission 仍直接做;之后不再直接审批或 Execute。
Tool.Lookup(...)
Tool.Validate(arguments)
Tool.CheckPermission(...)
out, err = hooks.BeforeToolCall(ctx, BeforeToolCallEvent{Arguments: arguments})
// Block:failResult(reason),skip Execute,不 Approval.Create,把失败结果喂回模型。
// 放行:arguments 写回 call,再走原来的审批。
Approval.Requires(...)                  // 仅钩子放行后
// 待批则暂停 Run,与 #4 相同。

hooks.OnToolExecutionStart(ctx, ...)    // 通知;现有 emit(execution_started) 旁再转一次
output := Tool.Execute(...)
output, err = hooks.AfterToolCall(ctx, AfterToolCallEvent{Result: output})
// 可覆盖 Output / Error / Success,再入库。每次重试都调;Success 被改成 true 则不再重试。
Event.Append("execution_result")

hooks.OnTurnEnd(ctx, ...)               // 写 turn.completed 之前
Event.Append("turn.completed")

hooks.OnRunEnd(ctx, ...)                // 必须在 Terminate 入口,覆盖一切终态
Run.Terminate(...)

审批顺序固定为:

Validate → CheckPermission → hooks.BeforeToolCall
  → Block:失败 Result,不审批、不执行
  → 放行:RequiresApproval → Execute → hooks.AfterToolCall

启动时:无插件 Hooks{};有插件 base.Then(plugin.Hooks())。还可以再 Then 指标或其他实现,不必改 Loop。

插件:typed func 的一种实现

插件不接入 Loop,只实现 Hooks 里的函数:内部再走 gRPC OnEvent。未订阅的类型不发 RPC。多插件按启动顺序串行,由 Compose* 完成;插件进程不直连。

plugin.BeforeToolCall = ComposeBeforeToolCall(p1.BeforeToolCall, p2.BeforeToolCall)
// 每个 pN.BeforeToolCall 内部:Plugin.OnEvent("tool_call", ...)

gRPC 双向通道,不要 gob。RPC 超时建议 2s。否决类超时/崩溃视为 Block(reason 标明 host);通知类忽略。

type PluginEvent struct {
    Type      string          // 只在插件 RPC 里用,Loop 看不见
    SessionID string
    RunID     string
    TurnID    string
    Payload   json.RawMessage
}
type EventResult struct {
    Block     bool
    Reason    string
    Terminate bool
    Action    string          // input: continue | transform | handled
    Payload   json.RawMessage
}

Host → Plugin:Bootstrap(注入 Host、声明 Subscribe、RegisterTool)、OnEventExecuteTool(仅插件自己注册的 tool)。

Plugin → Host 第一期要实现:RegisterToolGetActiveToolsSetActiveToolsGetSystemPrompt。接口先留空:RegisterCommandRegisterFlagSendUserMessageAppendEntryEventsEmit

需要实现的地方

Agent:

  • 定义上述 typed func、Compose*HooksHooks.Then。Runtime 持有 HooksExecute / EnsureTurn / Build 之后 / Terminate 只调字段。
  • Handler StartBeforeInput(interrupt 之前)。
  • Stream(openai)调 Chat 上的 OnHeaders / OnPayload / OnResponse
  • Dispatch 调 Invocation 上的 BeforeCall / AfterCall。Runtime 调用前从 Hooks 填入。

插件(后期拼到 Hooks 上):

  • Plugin.Load:发现二进制(CODEDOCK_PLUGIN_DIR~/.codedock/plugins/、项目 .codedock/plugins/),go-plugin 握手,常驻进程,启动顺序 = 发现顺序,同名后者覆盖。
  • 每个 typed func 的 RPC 封装、订阅过滤、超时。
  • Host.RegisterTool:写入现有 Tool.RegistryExecutePlugin.ExecuteTool。live 工具名与 Profile.Tools.Names 取并集后再 Tool.VisibleDefinitions
  • 热重载:Kill + 再 Bootstrap。第一期可只做启动加载。

插件进程 SDK(不依赖 db):

func main() { Plugin.Serve(...) }

func (p *MyPlugin) Bootstrap(host Host) {
    host.RegisterTool(def) // 并声明 Subscribe
}

func (p *MyPlugin) OnEvent(ev PluginEvent) EventResult {
    if ev.Type == "tool_call" && dangerous(ev) {
        return EventResult{Block: true, Reason: "policy"}
    }
    return EventResult{} // 沉默 = 放行
}

文档:写清 Loop 只认 Hooks;Plugin 与 AgentEvent 分属两条总线;插件 Tool 仍走 Tool.Dispatch

第一期范围

必须:BeforeInput handled、TransformContext / BeforeRunStart 改写、BeforeProvider*(仅 openai)、BeforeToolCall 阻断与改参、AfterToolCall 改结果、Compose 串行多实现、进程常驻、RegisterTool 进 Registry。

不做:压缩闸门、自定义 UI、替换整个 Provider stream、逐 token RPC、EventsEmit 实装、沙箱。assistant.delta 不进 Hooks。不改 #4 的 Session 并发、压缩不删原文、断线不取消 Run。

测试

  • ComposeBeforeToolCall:一票否决短路、改参写回、nil 跳过。
  • Run.StartBeforeInput handled 不 Insert / 不 Start、不 Cancel active Run、响应 run_id 为空;transform 后入库的是改写后的 content。
  • Tool.DispatchBeforeCall Block 时 Tool.Execute 次数为 0,有失败 Result,且不进审批。
  • Model.Stream(openai):mock HTTP,断言 OnPayload 替换后的 body 才 DoOnHeaders 改过的 header 出现在发出的请求上。

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

    enhancementNew feature or request

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions