diff --git a/AGENTS.md b/AGENTS.md index 38a88af..a8bbfb4 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -32,9 +32,9 @@ loopd/ 可扩展 blocks 承载文本、可见工具状态和产物,不混入完整执行轨迹。 3. `conversations` 与 `messages` 均使用 go-stdx 生成的 UUIDv7 `id` 作为主键和游标;不再维护平行的 message sequence。 -4. 一次问答由 user Message、responder Message 和同 ID 的 Task CRD 组成。服务端在数据库事务提交前 - 创建 CRD;创建失败则回滚两条 Message 并向 Human 返回错误。 -5. 主会话的 `parent_message_id` 为空;Operator 内部工作会话通过该字段引用主链路 responder Message。 +4. 一次问答由 user Message、目标 Actor 的 response Message 和同 ID 的 Task CRD 组成。服务端在数据库事务提交前 + 初始化 AgentUE Redis Stream 并创建 CRD;任一步失败则回滚两条 Message,并尽力删除已创建的外部资源。 +5. 主会话的 `parent_message_id` 为空;Operator 内部工作会话通过该字段引用主链路 response Message。 v1 不允许详情会话继续嵌套,且同一条 Message 最多关联一个详情会话。 6. v1alpha1 Task CRD 当前保存路由和唤醒信息,详细上下文由 runtime 按 Task ID 向 server 查询。公共协调 字段可按 Kubernetes API 兼容规则增量演进;领域状态复杂时,Operator 创建并拥有自己的 CRD。 diff --git a/README.md b/README.md index 0907e28..bf5231c 100644 --- a/README.md +++ b/README.md @@ -10,7 +10,7 @@ Loop = Resource(spec + status) + Reconcile ``` Every chat request creates a small loopd Task CRD that wakes the selected -responder. A simple Operator can reconcile that Task directly; a complex +Actor. A simple Operator can reconcile that Task directly; a complex Operator may create domain CRDs for its own state and completion semantics. `loopd` keeps visible collaboration in conversations and messages, while full execution history belongs to the Harness and AgentLedger. @@ -26,7 +26,7 @@ execution history belongs to the Harness and AgentLedger. - **Harness adapters** connect loop-server to agentd or another intelligent execution service without leaking provider vocabulary into the public model. -AgentUE is used for the page-visible message model. AgentLedger records complete +AgentUE supplies the page-visible event model and Redis bridge. AgentLedger records complete prompts, model events, tool calls, retries, and costs; it is not the chat database. The public conversation roles are always: @@ -40,15 +40,15 @@ user | harness | operator A question may run for minutes or days. Its lifecycle is independent from any HTTP request or browser connection: -1. loop-server creates a user message and an empty responder message with one - `task_id`, then creates a same-ID Task CRD before committing the messages; +1. loop-server creates a user message and an empty target response message with one + `task_id`, initializes its AgentUE stream, then creates a same-ID Task CRD before committing the messages; 2. the selected Operator or Harness watches the Task and resolves its current input and conversation history through loop-runtime; 3. a complex Operator may create a domain CRD, while a simple Operator handles the shared Task directly; -4. visible progress updates the responder message while full events enter AgentLedger; -5. clients may disconnect and reload the current message snapshot; -6. the selected Harness or Operator finishes the same responder message. +4. visible progress flows through the AgentUE Redis event bridge while full events enter AgentLedger; +5. clients may reconnect to any server instance with the same `task_id` and continue from their last event cursor; +6. completion folds visible events into the selected Actor's response Message. `Harness.Prompt` returns a handle. A Reconciler can inspect it and return, or call `Wait` when waiting is genuinely the only remaining work. Both paths use diff --git a/chat.go b/chat.go index 6576ce3..978b4f2 100644 --- a/chat.go +++ b/chat.go @@ -42,7 +42,7 @@ type Invocation struct { ConversationID string `json:"conversation_id"` InputMessageID string `json:"input_message_id"` OutputMessageID string `json:"output_message_id,omitempty"` - Responder ResponderRef `json:"responder"` + Target ActorRef `json:"target"` ContextThroughMessageID string `json:"context_through_message_id"` Phase InvocationPhase `json:"phase"` Resource *ResourceRef `json:"resource,omitempty"` @@ -72,7 +72,7 @@ type Activity struct { InvocationID string `json:"invocation_id"` Key string `json:"key"` ParentID string `json:"parent_id,omitempty"` - Actor ResponderRef `json:"actor"` + Actor ActorRef `json:"actor"` Kind string `json:"kind"` Title string `json:"title"` Detail string `json:"detail,omitempty"` @@ -97,7 +97,7 @@ type OperatorEvent struct { type ActivityUpdate struct { Key string `json:"key"` ParentID string `json:"parent_id,omitempty"` - Actor ResponderRef `json:"actor"` + Actor ActorRef `json:"actor"` Kind string `json:"kind"` Title string `json:"title"` Detail string `json:"detail,omitempty"` diff --git a/cmd/loop-server/main.go b/cmd/loop-server/main.go index cba4e0b..b7a8bd2 100644 --- a/cmd/loop-server/main.go +++ b/cmd/loop-server/main.go @@ -28,13 +28,20 @@ func main() { func run() error { address := envOr("LOOP_SERVER_ADDR", ":8080") databasePath := envOr("LOOP_SERVER_DB", "loopd.db") + redisAddress := envOr("LOOP_SERVER_REDIS_ADDR", "127.0.0.1:6379") taskNamespace := envOr("LOOP_SERVER_TASK_NAMESPACE", "default") tasks, err := newTaskClient(taskNamespace) if err != nil { return err } loopServer, err := server.New(server.Config{ - Database: server.DatabaseConfig{Path: databasePath}, Tasks: tasks, Logger: slog.Default(), + Database: server.DatabaseConfig{Path: databasePath}, + Redis: server.RedisConfig{ + Address: redisAddress, + Username: os.Getenv("LOOP_SERVER_REDIS_USERNAME"), + Password: os.Getenv("LOOP_SERVER_REDIS_PASSWORD"), + }, + Tasks: tasks, Logger: slog.Default(), }) if err != nil { return err @@ -63,6 +70,7 @@ func run() error { logger.Info("loop-server listening", "address", address, "database", databasePath, + "redis_address", redisAddress, "task_namespace", taskNamespace, ) serveErr <- httpServer.Run() diff --git a/common.go b/common.go index d40075d..9c90925 100644 --- a/common.go +++ b/common.go @@ -21,17 +21,21 @@ func (role Role) Valid() bool { } } -type ResponderRef struct { +type ActorRef struct { Kind Role `json:"kind"` Key string `json:"key"` } -func (ref ResponderRef) Valid() bool { +func (ref ActorRef) Valid() bool { + return ref.Kind.Valid() && ref.Key != "" +} + +func (ref ActorRef) ValidTarget() bool { return (ref.Kind == RoleHarness || ref.Kind == RoleOperator) && ref.Key != "" } -type Responder struct { - ResponderRef +type Actor struct { + ActorRef DisplayName string `json:"display_name,omitempty"` Description string `json:"description,omitempty"` } diff --git a/docs/kernel.md b/docs/kernel.md index e7ac94a..768b84e 100644 --- a/docs/kernel.md +++ b/docs/kernel.md @@ -17,7 +17,7 @@ loopd 的定位分为三个维度: 一句话概括: > Human 在持久 Conversation 中向 Harness 或 Operator 提出目标;loop-server 创建 Task CRD 唤醒 -> responder,Operator 按需建立领域 CRD、调用 Harness,并把过程与结果带回同一 Conversation。 +> 选定的 Actor,Operator 按需建立领域 CRD、调用 Harness,并把过程与结果带回同一 Conversation。 loopd 在四层 Agent 技术体系中的位置如下: @@ -37,7 +37,7 @@ loopd 的公开协作角色只有三类: | 参与者 | 职责 | |---|---| -| **Human** | 提出目标、选择 responder、补充上下文、回答 Interaction,并判断结果是否满足需要 | +| **Human** | 提出目标、选择 Actor、补充上下文、回答 Interaction,并判断结果是否满足需要 | | **Harness** | 接受 prompt 与 tools,执行一次可流式观察的智能任务;既可直接回答 Human,也可被 Operator 调用 | | **Operator** | Reconcile loopd Task,读取外部事实并决定何时调用 Harness、询问 Human、等待或结束;复杂业务可以拥有领域 CRD | @@ -53,6 +53,9 @@ user | harness | operator 外部系统可能把某个执行目标称为 Agent、Assistant 或 Session;接入 loopd 后统一表现为 Harness。 loopd Core、公共 API、数据库和 UI 不再建立一套与 Harness 平行的 Agent 概念。 +公共类型用 `ActorRef {kind, key}` 表达参与者身份,用 `target` 表达本次请求选中的执行目标。`Member` 留给 +未来真实存在的 Operator 团队或可见性成员关系,避免把“能参与一次对话”和“隶属于某个集合”混为一谈。 + ## 3. 核心协作对象 三类参与者通过少量具有稳定身份和生命周期的对象协作: @@ -61,7 +64,7 @@ loopd Core、公共 API、数据库和 UI 不再建立一套与 Harness 平行 |---|---| | **Conversation** | loop-server 拥有的持久协作空间;历史属于 Conversation,不属于某个 Operator 或 Harness | | **Message** | Conversation 中一次页面可见表达;`task_id + kind + key` 标识问答任务与发送者,content 保存 AgentUE semantic model 快照 | -| **Task** | loop-server 为一次问答创建的通用 CRD;只保存 responder 路由与唤醒版本,名称即 Message 使用的 `task_id` | +| **Task** | loop-server 为一次问答创建的通用 CRD;只保存目标 Actor 路由与唤醒版本,名称即 Message 使用的 `task_id` | | **Operator Resource** | 复杂 Operator 按需创建并拥有的领域 CRD;记录该领域独有的目标、状态和完成条件 | | **Harness Execution** | Harness 自己拥有的执行;可耗时、流式返回,完整轨迹由 AgentLedger 记录 | @@ -99,7 +102,7 @@ metadata: spec: target: kind: operator | harness - key: + key: revision: 1 ``` @@ -118,6 +121,7 @@ Auditor,其他 Operator 可以采用完全不同的 Resource;这些领域概 ```text Human → create user Message + empty harness Message (same task_id) + → initialize the same-ID AgentUE event stream → create same-ID Task CRD before the database commit → Harness starts or resumes its execution → project visible stream into the harness Message @@ -133,6 +137,7 @@ Harness 的流式文本可以直接构成主回答,可见工具状态可以投 ```text Human → create user Message + empty operator Message (same task_id) + → initialize the same-ID AgentUE event stream → create same-ID Task CRD before the database commit → Operator Reconciler resolves Conversation context through loop-runtime → optionally create or update an Operator-owned domain CRD @@ -149,12 +154,16 @@ Operator 内部 Harness 的输出默认属于详情会话或 AgentLedger,不 内容才成为 `operator` 角色的最终回答;Operator 也可以选择把某个 Harness 结果显式公开为 Conversation Message。 +首次观察与断线续接使用同一个 Chat API:请求不带 `task_id` 时创建问答和 Task,带 `task_id` 时观察已有 +Task;`Last-Event-ID` 只表示客户端已消费的 Redis transport cursor。HTTP 资源模型不额外暴露 Replay +接口,续接后的 AgentUE delivery 仍然从重建的完整 `start` 开始。 + ## 6. Conversation Context Conversation History 是 loop-server 的事实,不属于任何 Operator 或 Harness。同一 Conversation 可以先后 由不同 Operator 和 Harness 参与,它们都可以在授权范围内读取已有历史。 -v1alpha1 Task CRD 不复制完整对话,当前保存 responder 路由与唤醒版本。loop-server 根据 Task ID 从主 +v1alpha1 Task CRD 不复制完整对话,当前保存目标 Actor 路由与唤醒版本。loop-server 根据 Task ID 从主 Conversation Message 即时组装上下文: ```text @@ -177,10 +186,11 @@ Human interaction 的控制骨架: ```text Chat.Conversation 读取 Conversation Chat.History 显式读取 Conversation 的增量历史 -Chat.Send 创建 user、responder Message 与 Task CRD,并返回 responder Message -Chat.Update 更新 responder Message 的可见语义快照 +Chat.Send 首次请求创建两条 Message 与 Task CRD;带 task_id 时续接同一 AgentUE 事件流 +Chat.Emit Operator 发布 set/append 页面事件 +Chat.Complete 折叠事件并固化 response Message,然后结束事件流 Task.Get 按 Task ID 读取当前 input、response 与 Conversation History -Task.Watch 为指定 responder 注册 Task Reconciler +Task.Watch 为指定 Actor 注册 Task Reconciler Harness.Prompt 以 prompt、tools 和可选 Harness target 发起或恢复一次 Call Ask / Confirm 请求 Human 提供信息或确认决定 ``` @@ -233,9 +243,9 @@ loop-server 是页面协作事实和 Task 分发的 owner: - Conversation 代表一个对话框; - Message 是页面可见对话历史的事实来源; -- 同一次问答的 user、responder 及后续可见交互共享 `task_id`; +- 同一次问答的 user、response 及后续可见交互共享 `task_id`; - Message content 可以投影 Activity、Artifact 等页面需要展示的内容。 -- Task CRD 以同一个 `task_id` 命名;v1alpha1 当前承载 responder 路由与可推进版本。 +- Task CRD 以同一个 `task_id` 命名;v1alpha1 当前承载目标 Actor 路由与可推进版本。 loop-server 不建立 `tasks` 表,也不保存 Operator 领域表。Task 查询是基于 Message 的实时视图;通用 Activity 只承载跨 Operator 都能理解的处理摘要,更丰富的领域详情留在 Operator Resource,必要时通过 @@ -244,26 +254,27 @@ Activity 只承载跨 Operator 都能理解的处理摘要,更丰富的领域 AgentLedger 记录 prompt、模型事件、Harness Call、tool call/result、重试和成本等完整执行事实,用于审计、 回放和分析;它不承担页面 Conversation History,因此不能替代 loop-server 的两表业务模型。 -AgentUE 在 loopd 中定义页面语义模型。AgentUE Runner 当前负责 Python 后台任务、Redis 事件桥和 heartbeat -recovery;直接把它嵌入 Go loop-server 会与 Harness provider、Task 和领域 CRD 形成多个执行 owner。因此 -v1 只复用 AgentUE 的页面模型,不把 AgentUE Runner 作为 loopd 的执行 owner。 +AgentUE 在 loopd 中定义页面语义模型,并提供不拥有业务任务的 Redis Event Bridge。Operator 通过 +loop-runtime 发布 `set/append`;任意 loop-server 实例都可按 `task_id` 和 transport cursor 重建完整 +`start` 快照并继续输出。loop-server 在完成阶段把最终快照写入 Message。AgentUE 不创建 Task、不调用 +Harness,也不成为 Operator 执行的 owner。 ## 10. 关键不变量 1. **公开角色只有三种**:Conversation 中只使用 `user`、`harness`、`operator`;外部 provider 的术语不 扩散到 loopd Core。 -2. **Conversation 独立于 responder**:历史属于 Human 持有的 Conversation,不因切换 Operator 或 Harness +2. **Conversation 独立于 Actor**:历史属于 Human 持有的 Conversation,不因切换 Operator 或 Harness 被复制或切断。 3. **Task 从分发和唤醒起步**:v1alpha1 保持最小;后续只增加通用协调字段,不复制 Query、History、 回答和 Operator 领域状态。 4. **领域 CRD 按复杂度引入**:简单 Operator 直接 Reconcile Task;复杂 Operator 自己拥有领域 Resource。 -5. **一次问答有稳定 task identity**:初始 user/responder Message 和后续反问、确认共享同一 `task_id`。 +5. **一次问答有稳定 task identity**:初始 user/response Message 和后续反问、确认共享同一 `task_id`。 6. **外部动作先获得持久 identity**:同一 owner 与 effect key 的重试必须观察或恢复同一次 Harness Call, 不能重复触发无法证明结果的副作用。 7. **流式响应与完成正交**:有事件表示 Call 有进展,不表示成功;长时间运行也不能被等同于失联。 8. **上下文有明确边界**:TaskContext 给出当前输入、回答和 History 水位;复杂 Operator 可以把所需引用 保存到自己的 Resource。 -9. **最终回答由 responder 收口**:直连 Harness 由 Harness 回答;Operator 内部可以调用多个 Harness,但由 +9. **最终回答由目标 Actor 收口**:直连 Harness 由 Harness 回答;Operator 内部可以调用多个 Harness,但由 Operator 汇总并回答。 10. **可见历史与完整轨迹分层**:Message 只保存页面可见快照,AgentLedger 保存完整执行轨迹。 11. **Harness 差异止于 Adapter**:新增 agentd 或第三方 Harness 不应要求修改 Conversation 或 Operator 模型。 diff --git a/go.mod b/go.mod index f92e1a4..c95c2cb 100644 --- a/go.mod +++ b/go.mod @@ -3,8 +3,11 @@ module github.com/compforge/loopd go 1.26.0 require ( + github.com/alicebob/miniredis/v2 v2.39.0 github.com/cloudwego/hertz v0.10.4 + github.com/compforge/agentue/sdks/go v0.0.0-20260904102512-0ec4d015e66a github.com/qiankunli/go-stdx v0.0.4-0.20260824051808-f7f6d7c53de2 + github.com/redis/go-redis/v9 v9.22.0 gorm.io/driver/sqlite v1.6.0 gorm.io/gorm v1.31.1 k8s.io/apimachinery v0.36.0 @@ -35,7 +38,7 @@ require ( github.com/jinzhu/now v1.1.5 // indirect github.com/josharian/intern v1.0.0 // indirect github.com/json-iterator/go v1.1.12 // indirect - github.com/klauspost/cpuid/v2 v2.2.9 // indirect + github.com/klauspost/cpuid/v2 v2.2.10 // indirect github.com/mailru/easyjson v0.7.7 // indirect github.com/mattn/go-sqlite3 v1.14.22 // indirect github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect @@ -52,6 +55,8 @@ require ( github.com/tidwall/pretty v1.2.0 // indirect github.com/twitchyliquid64/golang-asm v0.15.1 // indirect github.com/x448/float16 v0.8.4 // indirect + github.com/yuin/gopher-lua v1.1.1 // indirect + go.uber.org/atomic v1.11.0 // indirect go.yaml.in/yaml/v2 v2.4.3 // indirect go.yaml.in/yaml/v3 v3.0.4 // indirect golang.org/x/arch v0.0.0-20210923205945-b76863e36670 // indirect diff --git a/go.sum b/go.sum index ac1e00a..ae9824b 100644 --- a/go.sum +++ b/go.sum @@ -1,7 +1,13 @@ github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1Xbatp0= github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= +github.com/alicebob/miniredis/v2 v2.39.0 h1:M7WbmV5BmV56L8KTG0rw6vEQ+woTOghpDgin2xv4A0g= +github.com/alicebob/miniredis/v2 v2.39.0/go.mod h1:TcL7YfarKPGDAthEtl5NBeHZfeUQj6OXMm/+iu5cLMM= github.com/beorn7/perks v1.0.1 h1:VlbKKnNfV8bJzeqoa4cOKqO6bYr3WgKZxO8Z16+hsOM= github.com/beorn7/perks v1.0.1/go.mod h1:G2ZrVWU2WbWT9wwq4/hrbKbnv/1ERSJQ0ibhJ6rlkpw= +github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= +github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= +github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= +github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= github.com/bytedance/gopkg v0.1.1/go.mod h1:576VvJ+eJgyCzdjS+c4+77QF3p7ubbtiKARP3TxducM= github.com/bytedance/gopkg v0.1.3 h1:TPBSwH8RsouGCBcMBktLt1AymVo2TVsBVCY4b6TnZ/M= github.com/bytedance/gopkg v0.1.3/go.mod h1:576VvJ+eJgyCzdjS+c4+77QF3p7ubbtiKARP3TxducM= @@ -19,6 +25,8 @@ github.com/cloudwego/hertz v0.10.4 h1:xJxomApZYR67cROevam6SrtUBDvhcI4ZZhx/WgvpHw github.com/cloudwego/hertz v0.10.4/go.mod h1:tZXEi/4o7R0Ho9yw5V2C+k/wVx3S8+wuuiJGDMopnpg= github.com/cloudwego/netpoll v0.7.2 h1:4qDBGQ6CG2SvEXhZSDxMdtqt/NLDxjAVk0PC/biKiJo= github.com/cloudwego/netpoll v0.7.2/go.mod h1:PI+YrmyS7cIr0+SD4seJz3Eo3ckkXdu2ZVKBLhURLNU= +github.com/compforge/agentue/sdks/go v0.0.0-20260904102512-0ec4d015e66a h1:ZxgLsatqSGYE1WBJ7mUVdEGHN9x6qKw6rIYtXRNIBko= +github.com/compforge/agentue/sdks/go v0.0.0-20260904102512-0ec4d015e66a/go.mod h1:AH0Je8lE9kGZuG8EtdnwExIsj4EqNWwSDFRG17GwFL8= github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -69,8 +77,8 @@ github.com/json-iterator/go v1.1.12 h1:PV8peI4a0ysnczrg+LtxykD8LfKY9ML6u2jnxaEnr github.com/json-iterator/go v1.1.12/go.mod h1:e30LSqwooZae/UwlEbR2852Gd8hjQvJoHmT4TnhNGBo= github.com/klauspost/compress v1.18.0 h1:c/Cqfb0r+Yi+JtIEq73FWXVkRonBlf0CRNYc8Zttxdo= github.com/klauspost/compress v1.18.0/go.mod h1:2Pp+KzxcywXVXMr50+X0Q/Lsb43OQHYWRCY2AiWywWQ= -github.com/klauspost/cpuid/v2 v2.2.9 h1:66ze0taIn2H33fBvCkXuv9BmCwDfafmiIVpKV9kKGuY= -github.com/klauspost/cpuid/v2 v2.2.9/go.mod h1:rqkxqrZ1EhYM9G+hXH7YdowN5R5RGN6NK4QwQ3WMXF8= +github.com/klauspost/cpuid/v2 v2.2.10 h1:tBs3QSyvjDyFTq3uoc/9xFpCuOsJQFNPiAhYdw2skhE= +github.com/klauspost/cpuid/v2 v2.2.10/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= github.com/kr/pretty v0.2.1/go.mod h1:ipq/a2n7PKx3OHsz4KJII5eveXtPO4qwEXGdVfWzfnI= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= @@ -112,6 +120,8 @@ github.com/prometheus/procfs v0.19.2 h1:zUMhqEW66Ex7OXIiDkll3tl9a1ZdilUOd/F6ZXw4 github.com/prometheus/procfs v0.19.2/go.mod h1:M0aotyiemPhBCM0z5w87kL22CxfcH05ZpYlu+b4J7mw= github.com/qiankunli/go-stdx v0.0.4-0.20260824051808-f7f6d7c53de2 h1:WMdZTrjsy7+EBtvHCsrRTl0KR7uVqUbTHgVt5gk6Mmw= github.com/qiankunli/go-stdx v0.0.4-0.20260824051808-f7f6d7c53de2/go.mod h1:cHhuc3z6NOiL/MDOTrFq26y5tzhlSULFa9F584N1FdU= +github.com/redis/go-redis/v9 v9.22.0 h1:laDvpYXTJtZLloinw1fA5Kqd6HAEH2XKxOkG/PDq2F0= +github.com/redis/go-redis/v9 v9.22.0/go.mod h1:y2g0Wj8rQvuK0ELM+oxSudcLtC09JScs98I/X9gRWY4= github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= github.com/spf13/pflag v1.0.9 h1:9exaQaMOCwffKiiiYk6/BndUBv+iRViNW+4lEMi0PvY= @@ -141,6 +151,12 @@ github.com/twitchyliquid64/golang-asm v0.15.1/go.mod h1:a1lVb/DtPvCB8fslRZhAngC2 github.com/x448/float16 v0.8.4 h1:qLwI1I70+NjRFUR3zs1JPUCgaCXSh3SW62uAKT1mSBM= github.com/x448/float16 v0.8.4/go.mod h1:14CWIYCyZA/cWjXOioeEpHeN/83MdbZDRQHoFcYsOfg= github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= +github.com/yuin/gopher-lua v1.1.1 h1:kYKnWBjvbNP4XLT3+bPEwAXJx262OhaHDWDVOPjL46M= +github.com/yuin/gopher-lua v1.1.1/go.mod h1:GBR0iDaNXjAgGg9zfCvksxSRnQx76gclCIb7kdAd1Pw= +github.com/zeebo/xxh3 v1.1.0 h1:s7DLGDK45Dyfg7++yxI0khrfwq9661w9EN78eP/UZVs= +github.com/zeebo/xxh3 v1.1.0/go.mod h1:IisAie1LELR4xhVinxWS5+zf1lA4p0MW4T+w+W07F5s= +go.uber.org/atomic v1.11.0 h1:ZvwS0R+56ePWxUNi+Atn9dWONBPp/AUETXlHW0DxSjE= +go.uber.org/atomic v1.11.0/go.mod h1:LUxbIzbOniOlMKjJjyPfpl4v+PKK2cNJn91OQbhoJI0= go.uber.org/goleak v1.3.0 h1:2K3zAYmnTNqV73imy9J1T3WC+gmCePx2hEGkimedGto= go.uber.org/goleak v1.3.0/go.mod h1:CoHD4mav9JJNrW/WLlf7HGZPjdw8EucARQHekz1X6bE= go.uber.org/multierr v1.11.0 h1:blXXJkSxSSfBVBlC76pxqeO+LN3aDfLQo+309xJstO0= diff --git a/interaction.go b/interaction.go index 00cdbec..da66278 100644 --- a/interaction.go +++ b/interaction.go @@ -32,7 +32,7 @@ type Interaction struct { InvocationID string `json:"invocation_id"` OwnerUID string `json:"owner_uid"` EffectKey string `json:"effect_key"` - Requester ResponderRef `json:"requester"` + Requester ActorRef `json:"requester"` Kind InteractionKind `json:"kind"` Title string `json:"title,omitempty"` Prompt string `json:"prompt"` diff --git a/runtime/chat.go b/runtime/chat.go index 5ed73ef..c61d4dd 100644 --- a/runtime/chat.go +++ b/runtime/chat.go @@ -1,15 +1,24 @@ package runtime import ( + "bufio" + "bytes" "context" "encoding/json" + "errors" + "fmt" + "io" "net/http" "net/url" "strconv" + "strings" + agentueui "github.com/compforge/agentue/sdks/go/ui" loopd "github.com/compforge/loopd" ) +const taskIDHeader = "X-Loopd-Task-ID" + type Chat struct { client *client } @@ -20,9 +29,20 @@ type CreateConversationRequest struct { } type SendMessageRequest struct { - UserKey string `json:"user_key"` - Responder loopd.ResponderRef `json:"responder"` - Content json.RawMessage `json:"content"` + TaskID string `json:"task_id,omitempty"` + UserKey string `json:"user_key,omitempty"` + Target loopd.ActorRef `json:"target,omitempty"` + Content json.RawMessage `json:"content,omitempty"` +} + +type ChatEvent struct { + Cursor string + Data json.RawMessage +} + +type TaskFailure struct { + Code string `json:"code"` + Message string `json:"message"` } func (chat Chat) CreateConversation( @@ -46,22 +66,91 @@ func (chat Chat) Conversation(ctx context.Context, conversationID string) (loopd return result, err } +// Send starts a question when TaskID is empty, or resumes the same task when +// TaskID is present. The HTTP connection only observes execution; closing the +// stream does not cancel the Task or its Operator. func (chat Chat) Send( ctx context.Context, conversationID string, request SendMessageRequest, -) (loopd.Message, error) { - var result loopd.Message - err := chat.client.do( - ctx, - http.MethodPost, - "/v1/conversations/"+url.PathEscape(conversationID)+"/messages", - request, - &result, - ) - return result, err + lastEventID string, +) (*ChatStream, error) { + path := "/v1/conversations/" + url.PathEscape(conversationID) + "/messages" + headers := map[string]string{"Accept": "text/event-stream"} + if lastEventID != "" { + headers["Last-Event-ID"] = lastEventID + } + response, err := chat.client.openWithHeaders(ctx, http.MethodPost, path, request, headers) + if err != nil { + return nil, err + } + if err := decodeResponseError(response); err != nil { + _ = response.Body.Close() + return nil, err + } + if !strings.HasPrefix(response.Header.Get("Content-Type"), "text/event-stream") { + _ = response.Body.Close() + return nil, fmt.Errorf("loop-server returned %q instead of text/event-stream", response.Header.Get("Content-Type")) + } + taskID := response.Header.Get(taskIDHeader) + if taskID == "" { + _ = response.Body.Close() + return nil, errors.New("loop-server response omitted task ID") + } + stream := &ChatStream{ + taskID: taskID, + body: response.Body, + scanner: bufio.NewScanner(response.Body), + cursor: lastEventID, + } + stream.scanner.Buffer(make([]byte, 64<<10), 16<<20) + return stream, nil } +// ChatStream reads one AgentUE event at a time from the Chat SSE response. +type ChatStream struct { + taskID string + body io.ReadCloser + scanner *bufio.Scanner + cursor string +} + +func (stream *ChatStream) TaskID() string { return stream.taskID } + +func (stream *ChatStream) Next() (ChatEvent, error) { + cursor := stream.cursor + var data bytes.Buffer + for stream.scanner.Scan() { + line := stream.scanner.Text() + if line == "" { + if data.Len() == 0 { + continue + } + value := bytes.TrimSuffix(data.Bytes(), []byte{'\n'}) + stream.cursor = cursor + return ChatEvent{Cursor: cursor, Data: append(json.RawMessage(nil), value...)}, nil + } + if strings.HasPrefix(line, "id:") { + cursor = strings.TrimSpace(strings.TrimPrefix(line, "id:")) + continue + } + if strings.HasPrefix(line, "data:") { + value := strings.TrimPrefix(line, "data:") + if strings.HasPrefix(value, " ") { + value = value[1:] + } + data.WriteString(value) + data.WriteByte('\n') + } + } + if err := stream.scanner.Err(); err != nil { + return ChatEvent{}, err + } + return ChatEvent{}, io.EOF +} + +func (stream *ChatStream) Close() error { return stream.body.Close() } + func (chat Chat) History( ctx context.Context, conversationID string, @@ -75,21 +164,33 @@ func (chat Chat) History( return result.Data, err } -func (chat Chat) Update( - ctx context.Context, - conversationID string, - messageID string, - content json.RawMessage, -) (loopd.Message, error) { - var result loopd.Message - err := chat.client.do( - ctx, - http.MethodPut, - "/v1/conversations/"+url.PathEscape(conversationID)+"/messages/"+url.PathEscape(messageID)+"/content", - struct { - Content json.RawMessage `json:"content"` - }{Content: content}, - &result, - ) - return result, err +// Emit publishes one Operator-produced set or append event for a Task. +func (chat Chat) Emit(ctx context.Context, taskID string, event agentueui.Event) (string, error) { + data, err := event.Marshal() + if err != nil { + return "", err + } + var result struct { + Cursor string `json:"cursor"` + } + err = chat.client.do(ctx, http.MethodPost, "/v1/tasks/"+url.PathEscape(taskID)+"/events", struct { + Event json.RawMessage `json:"event"` + }{Event: data}, &result) + return result.Cursor, err +} + +// Complete persists the latest AgentUE snapshot as the response Message and +// closes the task delivery. Repeating the same completion is safe. +func (chat Chat) Complete(ctx context.Context, taskID string, failure *TaskFailure) error { + return chat.client.do(ctx, http.MethodPost, "/v1/tasks/"+url.PathEscape(taskID)+"/complete", struct { + Error *TaskFailure `json:"error,omitempty"` + }{Error: failure}, nil) +} + +func IsEnd(event ChatEvent) (bool, error) { + parsed, err := agentueui.Parse(event.Data) + if err != nil { + return false, err + } + return parsed.Op == agentueui.OpEnd, nil } diff --git a/runtime/chat_test.go b/runtime/chat_test.go new file mode 100644 index 0000000..65ea87f --- /dev/null +++ b/runtime/chat_test.go @@ -0,0 +1,107 @@ +package runtime + +import ( + "context" + "encoding/json" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" + + agentueui "github.com/compforge/agentue/sdks/go/ui" +) + +func TestChatSendStartsAndResumesOneTask(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + if request.URL.Path != "/v1/conversations/conversation-1/messages" { + http.NotFound(response, request) + return + } + var input SendMessageRequest + if err := json.NewDecoder(request.Body).Decode(&input); err != nil { + t.Error(err) + } + if input.TaskID != "" && request.Header.Get("Last-Event-ID") != "2-0" { + t.Errorf("Last-Event-ID = %q, want 2-0", request.Header.Get("Last-Event-ID")) + } + response.Header().Set(taskIDHeader, "task-1") + response.Header().Set("Content-Type", "text/event-stream; charset=utf-8") + _, _ = io.WriteString(response, "data: {\"op\":\"start\",\"seq\":2,\"model\":{\"version\":\"1.0\",\"biz\":\"chat\",\"meta\":{},\"blocks\":[]}}\n\n") + _, _ = io.WriteString(response, "id: 4-0\ndata: {\"op\":\"end\",\"seq\":3}\n\n") + })) + t.Cleanup(server.Close) + runtime, err := New(server.URL, Options{HTTPClient: server.Client()}) + if err != nil { + t.Fatal(err) + } + + stream, err := runtime.Loop.Chat.Send(context.Background(), "conversation-1", SendMessageRequest{ + TaskID: "task-1", + }, "2-0") + if err != nil { + t.Fatal(err) + } + defer stream.Close() + if stream.TaskID() != "task-1" { + t.Fatalf("task ID = %q", stream.TaskID()) + } + first, err := stream.Next() + if err != nil { + t.Fatal(err) + } + parsed, err := agentueui.Parse(first.Data) + if err != nil { + t.Fatal(err) + } + if first.Cursor != "2-0" || parsed.Op != agentueui.OpStart { + t.Fatalf("first event = %#v, parsed=%#v", first, parsed) + } + last, err := stream.Next() + if err != nil { + t.Fatal(err) + } + ended, err := IsEnd(last) + if err != nil || !ended { + t.Fatalf("last event = %#v, ended=%t, error=%v", last, ended, err) + } + if _, err := stream.Next(); err != io.EOF { + t.Fatalf("stream tail error = %v, want EOF", err) + } +} + +func TestChatPublishesAndCompletesTask(t *testing.T) { + var published, completed bool + server := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + switch { + case strings.HasSuffix(request.URL.Path, "/events"): + published = true + response.Header().Set("Content-Type", "application/json") + _, _ = io.WriteString(response, `{"cursor":"2-0"}`) + case strings.HasSuffix(request.URL.Path, "/complete"): + completed = true + response.WriteHeader(http.StatusNoContent) + default: + http.NotFound(response, request) + } + })) + t.Cleanup(server.Close) + runtime, err := New(server.URL, Options{HTTPClient: server.Client()}) + if err != nil { + t.Fatal(err) + } + + cursor, err := runtime.Loop.Chat.Emit(context.Background(), "task-1", agentueui.Event{ + Op: agentueui.OpSet, Seq: 2, + Block: map[string]any{"id": "answer", "type": "text", "content": "done"}, + }) + if err != nil || cursor != "2-0" { + t.Fatalf("Emit cursor=%q error=%v", cursor, err) + } + if err := runtime.Loop.Chat.Complete(context.Background(), "task-1", nil); err != nil { + t.Fatal(err) + } + if !published || !completed { + t.Fatalf("published=%t completed=%t", published, completed) + } +} diff --git a/runtime/interaction.go b/runtime/interaction.go index 8b75c5e..1ac1de2 100644 --- a/runtime/interaction.go +++ b/runtime/interaction.go @@ -13,7 +13,7 @@ type InteractionPrompt struct { InvocationID string OwnerUID string EffectKey string - Requester loopd.ResponderRef + Requester loopd.ActorRef Title string Text string Options []loopd.InteractionOption diff --git a/runtime/runtime.go b/runtime/runtime.go index fc7d469..17451c7 100644 --- a/runtime/runtime.go +++ b/runtime/runtime.go @@ -8,6 +8,7 @@ import ( "errors" "fmt" "io" + "net" "net/http" "net/url" "strings" @@ -15,8 +16,9 @@ import ( ) type Options struct { - HTTPClient *http.Client - PollInterval time.Duration + HTTPClient *http.Client + PollInterval time.Duration + RequestTimeout time.Duration } type Runtime struct { @@ -37,12 +39,25 @@ func New(baseURL string, options Options) (*Runtime, error) { return nil, fmt.Errorf("invalid loop-server URL %q", baseURL) } if options.HTTPClient == nil { - options.HTTPClient = &http.Client{Timeout: 30 * time.Second} + transport := http.DefaultTransport.(*http.Transport).Clone() + transport.DialContext = (&net.Dialer{Timeout: 5 * time.Second, KeepAlive: 30 * time.Second}).DialContext + transport.TLSHandshakeTimeout = 5 * time.Second + transport.ResponseHeaderTimeout = 30 * time.Second + transport.IdleConnTimeout = 90 * time.Second + transport.MaxIdleConns = 100 + transport.MaxIdleConnsPerHost = 20 + options.HTTPClient = &http.Client{Transport: transport} } if options.PollInterval <= 0 { options.PollInterval = time.Second } - c := &client{baseURL: parsed, http: options.HTTPClient, pollInterval: options.PollInterval} + if options.RequestTimeout <= 0 { + options.RequestTimeout = 30 * time.Second + } + c := &client{ + baseURL: parsed, http: options.HTTPClient, + pollInterval: options.PollInterval, requestTimeout: options.RequestTimeout, + } loop := Loop{client: c} loop.Chat = Chat{client: c} loop.Harness = Harness{client: c} @@ -52,32 +67,70 @@ func New(baseURL string, options Options) (*Runtime, error) { } type client struct { - baseURL *url.URL - http *http.Client - pollInterval time.Duration + baseURL *url.URL + http *http.Client + pollInterval time.Duration + requestTimeout time.Duration } func (client *client) do(ctx context.Context, method, path string, input, output any) error { + ctx, cancel := context.WithTimeout(ctx, client.requestTimeout) + defer cancel() + + response, err := client.open(ctx, method, path, input) + if err != nil { + return err + } + defer response.Body.Close() + if err := decodeResponseError(response); err != nil { + return err + } + if output == nil || response.StatusCode == http.StatusNoContent { + return nil + } + if err := json.NewDecoder(io.LimitReader(response.Body, 16<<20)).Decode(output); err != nil { + return fmt.Errorf("decode loop-server response: %w", err) + } + return nil +} + +func (client *client) open(ctx context.Context, method, path string, input any) (*http.Response, error) { + return client.openWithHeaders(ctx, method, path, input, nil) +} + +func (client *client) openWithHeaders( + ctx context.Context, + method string, + path string, + input any, + headers map[string]string, +) (*http.Response, error) { var body io.Reader if input != nil { encoded, err := json.Marshal(input) if err != nil { - return err + return nil, err } body = bytes.NewReader(encoded) } request, err := http.NewRequestWithContext(ctx, method, client.baseURL.String()+path, body) if err != nil { - return err + return nil, err } if input != nil { request.Header.Set("Content-Type", "application/json") } + for key, value := range headers { + request.Header.Set(key, value) + } response, err := client.http.Do(request) if err != nil { - return err + return nil, err } - defer response.Body.Close() + return response, nil +} + +func decodeResponseError(response *http.Response) error { if response.StatusCode < 200 || response.StatusCode >= 300 { var envelope errorResponse if err := json.NewDecoder(io.LimitReader(response.Body, 1<<20)).Decode(&envelope); err == nil && envelope.Error.Message != "" { @@ -85,12 +138,6 @@ func (client *client) do(ctx context.Context, method, path string, input, output } return &Error{StatusCode: response.StatusCode, Message: response.Status} } - if output == nil || response.StatusCode == http.StatusNoContent { - return nil - } - if err := json.NewDecoder(io.LimitReader(response.Body, 16<<20)).Decode(output); err != nil { - return fmt.Errorf("decode loop-server response: %w", err) - } return nil } diff --git a/runtime/task.go b/runtime/task.go index ca666bc..efdad2c 100644 --- a/runtime/task.go +++ b/runtime/task.go @@ -31,8 +31,8 @@ func (task Task) Get(ctx context.Context, taskID string) (loopd.TaskContext, err // Watch registers a controller-runtime Reconciler for Tasks routed to target. // The Operator sees only a Task name and uses Get for current chat context. -func (task Task) Watch(mgr manager.Manager, target loopd.ResponderRef, reconciler reconcile.Reconciler) error { - if !target.Valid() { +func (task Task) Watch(mgr manager.Manager, target loopd.ActorRef, reconciler reconcile.Reconciler) error { + if !target.ValidTarget() { return fmt.Errorf("invalid task target %q/%q", target.Kind, target.Key) } if err := taskv1alpha1.AddToScheme(mgr.GetScheme()); err != nil { diff --git a/runtime/wire.go b/runtime/wire.go index d70af23..5e4403f 100644 --- a/runtime/wire.go +++ b/runtime/wire.go @@ -34,7 +34,7 @@ type promptRequest struct { type interactionRequest struct { OwnerUID string `json:"owner_uid"` EffectKey string `json:"effect_key"` - Requester loopd.ResponderRef `json:"requester"` + Requester loopd.ActorRef `json:"requester"` Kind loopd.InteractionKind `json:"kind"` Title string `json:"title,omitempty"` Prompt string `json:"prompt"` diff --git a/server/AGENTS.md b/server/AGENTS.md index a107839..72e5d6a 100644 --- a/server/AGENTS.md +++ b/server/AGENTS.md @@ -12,6 +12,7 @@ server 是 loop-server 组件,数据库只拥有 Conversation 与 Message 两 server/ ├── server.go # 组件组装与资源生命周期 ├── internal/api/ # Hertz 适配;Conversation、Message、Chat、Task handler 分文件 +├── internal/delivery/ # loopd 完成语义;借助 AgentUE Bridge 续接事件并固化 Message ├── internal/model/ # GORM model;一张表一个 Go 文件 │ ├── conversation.go # conversations │ └── message.go # messages @@ -36,16 +37,18 @@ server/ 不是完整执行日志。 4. 两张表的 `id` 都由 service 使用 go-stdx `uuid.V7()` 生成。消息按 UUIDv7 字典序读取和翻页,不增加 `sequence` 字段。 -5. 一次 Chat 请求由 `ChatService` 在数据库事务中创建 user Message 和空的 responder Message,再以 - `task_id` 创建 Task CRD;CRD 创建失败则事务回滚。Task CRD 创建成功但 DB commit 失败时,server - 尽力删除该 CRD。 -6. `conversations.parent_message_id` 只用于把 Operator 详情会话挂到主链路 responder Message。主会话为空, +5. 一次 Chat 请求由 `ChatService` 在数据库事务中创建 user Message 和目标 Actor 的空 response Message,再以 + `task_id` 初始化 AgentUE Redis Stream 并创建 Task CRD;任一步失败则事务回滚。外部资源创建成功但 DB + commit 失败时,server 尽力删除 Stream 与 CRD。 +6. `conversations.parent_message_id` 只用于把 Operator 详情会话挂到主链路 response Message。主会话为空, 详情会话不可继续嵌套,同一 Message 最多对应一个详情会话。 7. 页面不可见的 prompt、tool call/result、重试和成本等完整轨迹进入 AgentLedger,不扩张 Message schema。 8. Task 查询是基于主 Conversation Message 的实时视图,不增加 `tasks` 表;详情 Conversation 中共享 `task_id` 的内部消息不得覆盖主链路的 input/response。 9. `runtime/api/` 是公共 CRD API,字段按 Kubernetes API 兼容规则增量演进;修改类型后运行 `make generate manifests`,提交生成的 DeepCopy 和 CRD YAML。 +10. AgentUE 拥有事件协议、Reducer、Redis Stream 与续接;server 拥有 Redis client 生命周期、Task 路由和 + 最终 Message 快照。Redis 中的事件是活跃任务投影,不替代 AgentLedger 的完整执行记录。 ## References diff --git a/server/docs/persistence.md b/server/docs/persistence.md index 91b1d75..60f3c65 100644 --- a/server/docs/persistence.md +++ b/server/docs/persistence.md @@ -13,7 +13,7 @@ loop-server 数据库只持久化 Conversation 与 Message 两类页面可见的 - `parent_message_id`:可空;Operator 详情会话引用的主链路 Message ID; - `created_at`、`updated_at`:记录时间。 -主会话的 `parent_message_id` 为 `NULL`。详情会话只允许引用主会话中的 responder Message,不继续嵌套; +主会话的 `parent_message_id` 为 `NULL`。详情会话只允许引用主会话中的 response Message,不继续嵌套; 同一条 Message 在 v1 中最多关联一个详情会话。 ## Message @@ -31,8 +31,8 @@ loop-server 数据库只持久化 Conversation 与 Message 两类页面可见的 一次 Chat 请求在一个数据库事务中创建两条 Message:Human 的问题是 `kind=user`,同时创建一条 `kind=operator` 或 `kind=harness` 的空回答。两条记录共享 `task_id`。提交事务前,ChatService 使用该 ID -和 responder 创建 Task CRD;创建失败则回滚两条 Message 并返回错误。CRD 创建成功但数据库 commit -失败时,server 尽力删除该 CRD 作为补偿。 +初始化 AgentUE Redis Stream 并为目标 Actor 创建 Task CRD;任一步失败则回滚两条 Message。外部资源创建 +成功但数据库 commit 失败时,server 尽力删除 Stream 与 CRD 作为补偿。 Task CRD 可能在数据库 commit 前被 Operator 观察到。此时 Task 查询返回 not found,Reconciler 应稍后 重试。正常提交后,`GET /v1/tasks/:task_id` 从主 Conversation 的 Message 即时组装 input、response 和 @@ -53,8 +53,8 @@ Task CRD 可能在数据库 commit 前被 Operator 观察到。此时 Task 查 ``` block 除 `id` 和 `type` 外的字段由 biz 扩展。页面无需展示的 system prompt、tool call/result、模型原始 -事件、重试与成本不进入 Message,由 AgentLedger 记录。流式 delta 也不作为独立 Message;执行层将它们 -折叠为页面可恢复的 content 快照。 +事件、重试与成本不进入 Message,由 AgentLedger 记录。流式 delta 也不作为独立 Message;AgentUE Redis +Bridge 承载运行中的页面事件,任务完成时由 server 将它们折叠为可恢复的 content 快照。 ## UUIDv7 游标 diff --git a/server/internal/api/chat.go b/server/internal/api/chat.go index 6ac55ee..191444b 100644 --- a/server/internal/api/chat.go +++ b/server/internal/api/chat.go @@ -2,22 +2,54 @@ package api import ( "context" + "errors" hertzapp "github.com/cloudwego/hertz/pkg/app" - "github.com/cloudwego/hertz/pkg/protocol/consts" + hertzsse "github.com/cloudwego/hertz/pkg/protocol/sse" + "github.com/compforge/loopd/server/internal/delivery" ) +const taskIDHeader = "X-Loopd-Task-ID" + func (server *Server) createChatMessages(ctx context.Context, request *hertzapp.RequestContext) error { var input createChatMessagesRequest if err := decodeBody(request, &input); err != nil { return err } - message, err := server.chat.Create( - ctx, request.Param("conversation_id"), input.UserKey, input.Responder, input.Content, + conversationID := request.Param("conversation_id") + taskID := input.TaskID + if taskID == "" { + message, err := server.chat.Create(ctx, conversationID, input.UserKey, input.Target, input.Content) + if err != nil { + return err + } + taskID = message.TaskID + } + request.Response.Header.Set(taskIDHeader, taskID) + var writer *hertzsse.Writer + streamErr := server.chat.Stream( + ctx, + conversationID, + taskID, + hertzsse.GetLastEventID(&request.Request), + func(event delivery.Event) error { + if writer == nil { + writer = hertzsse.NewWriter(request) + } + return writer.WriteEvent(event.Cursor, "", event.Data) + }, ) - if err != nil { - return err + if errors.Is(streamErr, context.Canceled) { + streamErr = nil + } + if writer == nil { + return streamErr + } + closeErr := writer.Close() + if err := errors.Join(streamErr, closeErr); err != nil { + // Headers may already be on the wire. Logging is safe; returning the + // error to the generic adapter would append a JSON error to the SSE body. + server.logger.ErrorContext(ctx, "chat stream stopped", "task_id", taskID, "error", err) } - request.JSON(consts.StatusCreated, message) return nil } diff --git a/server/internal/api/event.go b/server/internal/api/event.go new file mode 100644 index 0000000..e7b3668 --- /dev/null +++ b/server/internal/api/event.go @@ -0,0 +1,40 @@ +package api + +import ( + "context" + + hertzapp "github.com/cloudwego/hertz/pkg/app" + "github.com/cloudwego/hertz/pkg/protocol/consts" + "github.com/compforge/loopd/server/internal/delivery" +) + +func (server *Server) appendTaskEvent(ctx context.Context, request *hertzapp.RequestContext) error { + var input appendTaskEventRequest + if err := decodeBody(request, &input); err != nil { + return err + } + cursor, err := server.chat.Emit(ctx, request.Param("task_id"), input.Event) + if err != nil { + return err + } + request.JSON(consts.StatusAccepted, appendTaskEventResponse{Cursor: cursor}) + return nil +} + +func (server *Server) completeTask(ctx context.Context, request *hertzapp.RequestContext) error { + var input completeTaskRequest + if len(request.Request.Body()) != 0 { + if err := decodeBody(request, &input); err != nil { + return err + } + } + var failure *delivery.Failure + if input.Error != nil { + failure = &delivery.Failure{Code: input.Error.Code, Message: input.Error.Message} + } + if err := server.chat.Complete(ctx, request.Param("task_id"), failure); err != nil { + return err + } + request.SetStatusCode(consts.StatusNoContent) + return nil +} diff --git a/server/internal/api/message.go b/server/internal/api/message.go index ba33815..f513b9f 100644 --- a/server/internal/api/message.go +++ b/server/internal/api/message.go @@ -28,21 +28,6 @@ func (server *Server) listMessages(ctx context.Context, request *hertzapp.Reques return nil } -func (server *Server) updateMessageContent(ctx context.Context, request *hertzapp.RequestContext) error { - var input updateMessageContentRequest - if err := decodeBody(request, &input); err != nil { - return err - } - message, err := server.messages.UpdateMessageContent( - ctx, request.Param("conversation_id"), request.Param("message_id"), input.Content, - ) - if err != nil { - return err - } - request.JSON(consts.StatusOK, message) - return nil -} - func queryLimit(request *hertzapp.RequestContext) (int, error) { value := request.Query("limit") if value == "" { diff --git a/server/internal/api/server.go b/server/internal/api/server.go index 4b1a57f..6e3797c 100644 --- a/server/internal/api/server.go +++ b/server/internal/api/server.go @@ -44,11 +44,9 @@ func (server *Server) Register(engine *route.Engine) { engine.GET("/v1/conversations/:conversation_id", server.adapt(server.getConversation)) engine.GET("/v1/conversations/:conversation_id/messages", server.adapt(server.listMessages)) engine.POST("/v1/conversations/:conversation_id/messages", server.adapt(server.createChatMessages)) - engine.PUT( - "/v1/conversations/:conversation_id/messages/:message_id/content", - server.adapt(server.updateMessageContent), - ) engine.GET("/v1/tasks/:task_id", server.adapt(server.getTask)) + engine.POST("/v1/tasks/:task_id/events", server.adapt(server.appendTaskEvent)) + engine.POST("/v1/tasks/:task_id/complete", server.adapt(server.completeTask)) } type handler func(context.Context, *hertzapp.RequestContext) error diff --git a/server/internal/api/server_test.go b/server/internal/api/server_test.go index 7cfa199..bea3de7 100644 --- a/server/internal/api/server_test.go +++ b/server/internal/api/server_test.go @@ -1,17 +1,21 @@ package api import ( + "bytes" "context" "encoding/json" "path/filepath" "strings" "testing" + hertzapp "github.com/cloudwego/hertz/pkg/app" "github.com/cloudwego/hertz/pkg/common/config" "github.com/cloudwego/hertz/pkg/common/ut" "github.com/cloudwego/hertz/pkg/protocol" "github.com/cloudwego/hertz/pkg/route" + "github.com/cloudwego/hertz/pkg/route/param" loopd "github.com/compforge/loopd" + "github.com/compforge/loopd/server/internal/delivery" "github.com/compforge/loopd/server/internal/repo" "github.com/compforge/loopd/server/internal/service" ) @@ -25,7 +29,7 @@ func TestChatHTTPFlow(t *testing.T) { server := New( service.NewConversationService(store, nil), service.NewMessageService(store, nil), - service.NewChatService(store, nopTaskClient{}, nil), + service.NewChatService(store, nopTaskClient{}, completedChatRunner{}, nil), service.NewTaskService(store, nil), nil, ) @@ -41,22 +45,18 @@ func TestChatHTTPFlow(t *testing.T) { t.Fatal(err) } - sent := performJSON(t, engine, "POST", "/v1/conversations/"+conversation.ID+"/messages", `{ + taskID, stream := performChat(t, server, conversation.ID, `{ "user_key":"user-1", - "responder":{"kind":"operator","key":"intent"}, + "target":{"kind":"operator","key":"intent"}, "content":{"version":"1.0","biz":"chat","meta":{},"blocks":[{"id":"q","type":"text","content":"hello"}]} }`) - if sent.StatusCode() != 201 { - t.Fatalf("send status=%d body=%s", sent.StatusCode(), sent.Body()) + if !strings.Contains(stream, `"op":"start"`) || !strings.Contains(stream, `"op":"end"`) { + t.Fatalf("send body=%s, want AgentUE start and end events", stream) } - var answer loopd.Message - if err := json.Unmarshal(sent.Body(), &answer); err != nil { - t.Fatal(err) - } - if answer.Kind != loopd.RoleOperator || answer.Key != "intent" || answer.TaskID == "" { - t.Fatalf("answer = %#v", answer) + if taskID == "" { + t.Fatal("send response omitted task ID header") } - taskResponse := ut.PerformRequest(engine, "GET", "/v1/tasks/"+answer.TaskID, nil).Result() + taskResponse := ut.PerformRequest(engine, "GET", "/v1/tasks/"+taskID, nil).Result() if taskResponse.StatusCode() != 200 { t.Fatalf("task status=%d body=%s", taskResponse.StatusCode(), taskResponse.Body()) } @@ -64,7 +64,7 @@ func TestChatHTTPFlow(t *testing.T) { if err := json.Unmarshal(taskResponse.Body(), &task); err != nil { t.Fatal(err) } - if task.ID != answer.TaskID || task.Input.Kind != loopd.RoleUser || task.Response.ID != answer.ID { + if task.ID != taskID || task.Input.Kind != loopd.RoleUser || task.Response.Kind != loopd.RoleOperator { t.Fatalf("task = %#v", task) } @@ -76,18 +76,58 @@ func TestChatHTTPFlow(t *testing.T) { if err := json.Unmarshal(history.Body(), &result); err != nil { t.Fatal(err) } - if len(result.Data) != 2 || result.Data[0].TaskID != answer.TaskID || result.Data[1].ID != answer.ID { + if len(result.Data) != 2 || result.Data[0].TaskID != taskID || result.Data[1].ID != task.Response.ID { t.Fatalf("history = %#v", result.Data) } - if response := ut.PerformRequest(engine, "GET", "/v1/responders", nil).Result(); response.StatusCode() != 404 { - t.Fatalf("responders status=%d, want 404", response.StatusCode()) + if response := ut.PerformRequest(engine, "GET", "/v1/actors", nil).Result(); response.StatusCode() != 404 { + t.Fatalf("actors status=%d, want 404", response.StatusCode()) } } type nopTaskClient struct{} -func (nopTaskClient) Create(context.Context, string, loopd.ResponderRef) error { return nil } -func (nopTaskClient) Delete(context.Context, string) error { return nil } +func (nopTaskClient) Create(context.Context, string, loopd.ActorRef) error { return nil } +func (nopTaskClient) Delete(context.Context, string) error { return nil } + +type completedChatRunner struct{} + +func (completedChatRunner) Initialize(context.Context, string, json.RawMessage) error { return nil } +func (completedChatRunner) Delete(context.Context, string) error { return nil } +func (completedChatRunner) Emit(context.Context, string, json.RawMessage) (string, error) { + return "", nil +} + +type streamWriter struct{ bytes.Buffer } + +func (writer *streamWriter) Flush() error { return nil } +func (writer *streamWriter) Finalize() error { return nil } + +func performChat(t *testing.T, server *Server, conversationID, body string) (string, string) { + t.Helper() + request := hertzapp.NewContext(1) + request.Params = param.Params{{Key: "conversation_id", Value: conversationID}} + request.Request.SetBodyString(body) + request.Request.Header.Set("Content-Type", "application/json") + writer := &streamWriter{} + request.Response.HijackWriter(writer) + if err := server.createChatMessages(context.Background(), request); err != nil { + t.Fatal(err) + } + return string(request.Response.Header.Peek(taskIDHeader)), writer.String() +} +func (completedChatRunner) Complete(context.Context, string, *delivery.Failure) error { return nil } +func (completedChatRunner) Stream( + _ context.Context, + _ string, + _ string, + _ string, + deliver func(delivery.Event) error, +) error { + if err := deliver(delivery.Event{Cursor: "1-0", Data: json.RawMessage(`{"op":"start","seq":1,"model":{"version":"1.0","biz":"chat","meta":{},"blocks":[]}}`), Persisted: true}); err != nil { + return err + } + return deliver(delivery.Event{Cursor: "2-0", Data: json.RawMessage(`{"op":"end","seq":2}`), Persisted: true}) +} func performJSON(t *testing.T, engine *route.Engine, method, path, value string) *protocol.Response { t.Helper() diff --git a/server/internal/api/view.go b/server/internal/api/view.go index 899b21c..5f0013c 100644 --- a/server/internal/api/view.go +++ b/server/internal/api/view.go @@ -25,11 +25,25 @@ type createConversationRequest struct { } type createChatMessagesRequest struct { - UserKey string `json:"user_key"` - Responder loopd.ResponderRef `json:"responder"` - Content json.RawMessage `json:"content"` + TaskID string `json:"task_id,omitempty"` + UserKey string `json:"user_key,omitempty"` + Target loopd.ActorRef `json:"target,omitempty"` + Content json.RawMessage `json:"content,omitempty"` } -type updateMessageContentRequest struct { - Content json.RawMessage `json:"content"` +type appendTaskEventRequest struct { + Event json.RawMessage `json:"event"` +} + +type appendTaskEventResponse struct { + Cursor string `json:"cursor"` +} + +type completeTaskRequest struct { + Error *taskFailure `json:"error,omitempty"` +} + +type taskFailure struct { + Code string `json:"code"` + Message string `json:"message"` } diff --git a/server/internal/delivery/delivery.go b/server/internal/delivery/delivery.go new file mode 100644 index 0000000..766b6fc --- /dev/null +++ b/server/internal/delivery/delivery.go @@ -0,0 +1,213 @@ +// Package delivery connects loop-server business completion to AgentUE delivery. +package delivery + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "log/slog" + + agentuerunner "github.com/compforge/agentue/sdks/go/runner" + agentueui "github.com/compforge/agentue/sdks/go/ui" + loopd "github.com/compforge/loopd" + "github.com/compforge/loopd/server/internal/model" + "github.com/compforge/loopd/server/internal/repo" +) + +var ErrInvalidEvent = errors.New("invalid AgentUE event") + +type MessageRepository interface { + ListRootMessagesByTask(context.Context, string) ([]model.Message, error) + UpdateMessageContent(context.Context, string, string, []byte) (model.Message, error) +} + +type Failure struct { + Code string `json:"code"` + Message string `json:"message"` +} + +type Event struct { + Cursor string + Data json.RawMessage + Persisted bool +} + +type Coordinator struct { + events agentuerunner.EventBridge + repo MessageRepository + logger *slog.Logger +} + +func New(events agentuerunner.EventBridge, repository MessageRepository, logger *slog.Logger) *Coordinator { + if logger == nil { + logger = slog.Default() + } + return &Coordinator{events: events, repo: repository, logger: logger} +} + +func (coordinator *Coordinator) Initialize(ctx context.Context, taskID string, model json.RawMessage) error { + start, err := agentueui.Start(model, 1) + if err != nil { + return err + } + data, err := start.Marshal() + if err != nil { + return err + } + return coordinator.events.Initialize(ctx, taskID, model, data, start.Seq) +} + +func (coordinator *Coordinator) Delete(ctx context.Context, taskID string) error { + return coordinator.events.Delete(ctx, taskID) +} + +func (coordinator *Coordinator) Emit(ctx context.Context, taskID string, data json.RawMessage) (string, error) { + event, err := agentueui.Parse(data) + if err != nil { + return "", fmt.Errorf("%w: %v", ErrInvalidEvent, err) + } + if event.Op != agentueui.OpSet && event.Op != agentueui.OpAppend { + return "", fmt.Errorf("%w: only set and append events may be emitted", ErrInvalidEvent) + } + return coordinator.events.Publish(ctx, taskID, data, event.Seq) +} + +func (coordinator *Coordinator) Complete(ctx context.Context, taskID string, failure *Failure) error { + state, err := coordinator.events.State(ctx, taskID) + if err != nil { + return err + } + if state.Status.Terminal() { + return nil + } + conversationID, responseMessageID, err := coordinator.responseMessage(ctx, taskID) + if err != nil { + return err + } + values, err := coordinator.events.EventsThrough(ctx, taskID, "") + if err != nil { + return err + } + snapshot := map[string]any{} + lastSeq := uint64(0) + lastOp := agentueui.Op("") + hasFailure := false + for _, value := range values { + event, parseErr := agentueui.Parse(value.Data) + if parseErr != nil { + return fmt.Errorf("rebuild task %q at cursor %q: %w", taskID, value.Cursor, parseErr) + } + if event.Op != agentueui.OpPing && event.Seq <= lastSeq { + return fmt.Errorf("task %q AgentUE sequence did not increase", taskID) + } + snapshot, err = agentueui.Apply(snapshot, event) + if err != nil { + return fmt.Errorf("rebuild task %q at cursor %q: %w", taskID, value.Cursor, err) + } + if event.Op != agentueui.OpPing { + lastSeq = event.Seq + } + if event.Op == agentueui.OpError { + hasFailure = true + } + lastOp = event.Op + } + status := agentuerunner.StatusCompleted + if failure != nil && lastOp != agentueui.OpEnd && !hasFailure { + status = agentuerunner.StatusFailed + lastSeq++ + failed := agentueui.Failure(lastSeq, failure.Code, failure.Message) + data, marshalErr := failed.Marshal() + if marshalErr != nil { + return marshalErr + } + if _, err = coordinator.events.Publish(ctx, taskID, data, failed.Seq); err != nil { + return err + } + snapshot, err = agentueui.Apply(snapshot, failed) + if err != nil { + return err + } + hasFailure = true + lastOp = failed.Op + } + if hasFailure { + status = agentuerunner.StatusFailed + } + content, err := agentueui.MarshalSnapshot(snapshot) + if err != nil { + return fmt.Errorf("marshal task %q snapshot: %w", taskID, err) + } + if _, err := coordinator.repo.UpdateMessageContent(ctx, conversationID, responseMessageID, content); err != nil { + return fmt.Errorf("persist task %q snapshot: %w", taskID, err) + } + if lastOp != agentueui.OpEnd { + lastSeq++ + end := agentueui.End(lastSeq) + data, marshalErr := end.Marshal() + if marshalErr != nil { + return marshalErr + } + if _, err = coordinator.events.Publish(ctx, taskID, data, end.Seq); err != nil { + return err + } + } + if err := coordinator.events.MarkTerminal(ctx, taskID, status); err != nil { + return err + } + coordinator.logger.InfoContext(ctx, "chat task completed", + "task_id", taskID, + "conversation_id", conversationID, + "message_id", responseMessageID, + "status", status, + ) + return nil +} + +func (coordinator *Coordinator) Stream( + ctx context.Context, + taskID string, + conversationID string, + after string, + deliver func(Event) error, +) error { + taskConversationID, _, err := coordinator.responseMessage(ctx, taskID) + if err != nil { + return err + } + if taskConversationID != conversationID { + return agentuerunner.ErrNotFound + } + replayer := agentuerunner.Replayer{Bridge: coordinator.events} + return replayer.Stream(ctx, taskID, after, func(event agentuerunner.Delivery) error { + return deliver(Event{Cursor: event.Cursor, Data: event.Data, Persisted: event.Cursor != ""}) + }) +} + +func (coordinator *Coordinator) responseMessage(ctx context.Context, taskID string) (string, string, error) { + rows, err := coordinator.repo.ListRootMessagesByTask(ctx, taskID) + if err != nil { + return "", "", err + } + var conversationID, responseMessageID string + hasUser := false + for _, row := range rows { + if conversationID != "" && conversationID != row.ConversationID { + return "", "", repo.ErrNotFound + } + conversationID = row.ConversationID + if row.Kind == string(loopd.RoleUser) { + hasUser = true + continue + } + if responseMessageID != "" { + return "", "", repo.ErrNotFound + } + responseMessageID = row.ID + } + if !hasUser || conversationID == "" || responseMessageID == "" { + return "", "", repo.ErrNotFound + } + return conversationID, responseMessageID, nil +} diff --git a/server/internal/delivery/delivery_test.go b/server/internal/delivery/delivery_test.go new file mode 100644 index 0000000..121a563 --- /dev/null +++ b/server/internal/delivery/delivery_test.go @@ -0,0 +1,115 @@ +package delivery + +import ( + "context" + "encoding/json" + "path/filepath" + "testing" + "time" + + "github.com/alicebob/miniredis/v2" + agentuerunner "github.com/compforge/agentue/sdks/go/runner" + agentueui "github.com/compforge/agentue/sdks/go/ui" + "github.com/compforge/loopd/server/internal/model" + "github.com/compforge/loopd/server/internal/repo" + "github.com/redis/go-redis/v9" +) + +func TestCoordinatorCompletesAndStreamsAcrossInstances(t *testing.T) { + ctx := context.Background() + store, err := repo.Open(repo.Config{Path: filepath.Join(t.TempDir(), "loopd.db")}) + if err != nil { + t.Fatal(err) + } + t.Cleanup(func() { _ = store.Close() }) + if _, err := store.CreateConversation(ctx, model.Conversation{ID: "conversation-1"}); err != nil { + t.Fatal(err) + } + initial := json.RawMessage(`{"version":"1.0","biz":"chat","meta":{},"blocks":[]}`) + _, err = store.CreateChatMessages(ctx, + model.Message{ID: "message-1", ConversationID: "conversation-1", TaskID: "task-1", Kind: "user", Key: "user-1", Content: initial}, + model.Message{ID: "message-2", ConversationID: "conversation-1", TaskID: "task-1", Kind: "operator", Key: "intent", Content: initial}, + nil, + ) + if err != nil { + t.Fatal(err) + } + + redisServer := miniredis.RunT(t) + clientA := redis.NewClient(&redis.Options{Addr: redisServer.Addr()}) + clientB := redis.NewClient(&redis.Options{Addr: redisServer.Addr()}) + t.Cleanup(func() { _ = clientA.Close() }) + t.Cleanup(func() { _ = clientB.Close() }) + options := agentuerunner.BridgeOptions{KeyPrefix: "test", ReadBlock: time.Millisecond} + producer := New(agentuerunner.NewRedisEventBridge(clientA, options), store, nil) + consumer := New(agentuerunner.NewRedisEventBridge(clientB, options), store, nil) + + if err := producer.Initialize(ctx, "task-1", initial); err != nil { + t.Fatal(err) + } + set := marshalEvent(t, agentueui.Event{ + Op: agentueui.OpSet, Seq: 2, + Block: map[string]any{"id": "answer", "type": "text", "content": "hello"}, + }) + cursor, err := producer.Emit(ctx, "task-1", set) + if err != nil { + t.Fatal(err) + } + appendEvent := marshalEvent(t, agentueui.Event{ + Op: agentueui.OpAppend, Seq: 3, Mask: "block.content", + Block: map[string]any{"id": "answer", "type": "text", "content": " world"}, + }) + if _, err := producer.Emit(ctx, "task-1", appendEvent); err != nil { + t.Fatal(err) + } + if err := producer.Complete(ctx, "task-1", nil); err != nil { + t.Fatal(err) + } + + message, err := store.GetMessage(ctx, "message-2") + if err != nil { + t.Fatal(err) + } + var snapshot map[string]any + if err := json.Unmarshal(message.Content, &snapshot); err != nil { + t.Fatal(err) + } + blocks := snapshot["blocks"].([]any) + if blocks[0].(map[string]any)["content"] != "hello world" { + t.Fatalf("persisted snapshot = %s", message.Content) + } + + var delivered []Event + if err := consumer.Stream(ctx, "task-1", "conversation-1", cursor, func(event Event) error { + delivered = append(delivered, event) + return nil + }); err != nil { + t.Fatal(err) + } + if len(delivered) != 3 { + t.Fatalf("got %d deliveries, want reconstructed start, append, and end", len(delivered)) + } + first, err := agentueui.Parse(delivered[0].Data) + if err != nil { + t.Fatal(err) + } + if first.Op != agentueui.OpStart || first.Seq != 2 || delivered[0].Persisted { + t.Fatalf("first delivery = %#v, parsed=%#v", delivered[0], first) + } + last, err := agentueui.Parse(delivered[2].Data) + if err != nil { + t.Fatal(err) + } + if last.Op != agentueui.OpEnd || !delivered[2].Persisted { + t.Fatalf("last delivery = %#v, parsed=%#v", delivered[2], last) + } +} + +func marshalEvent(t *testing.T, event agentueui.Event) json.RawMessage { + t.Helper() + data, err := event.Marshal() + if err != nil { + t.Fatal(err) + } + return data +} diff --git a/server/internal/repo/message.go b/server/internal/repo/message.go index f8d7d43..a311a85 100644 --- a/server/internal/repo/message.go +++ b/server/internal/repo/message.go @@ -33,14 +33,14 @@ func (store *Store) CreateMessage(ctx context.Context, message model.Message) (m func (store *Store) CreateChatMessages( ctx context.Context, userMessage model.Message, - responderMessage model.Message, + responseMessage model.Message, beforeCommit func(context.Context) error, ) (model.Message, error) { ctx, cancel := store.withTimeout(ctx) defer cancel() err := store.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error { - if userMessage.ConversationID != responderMessage.ConversationID || userMessage.TaskID != responderMessage.TaskID { + if userMessage.ConversationID != responseMessage.ConversationID || userMessage.TaskID != responseMessage.TaskID { return ErrConflict } if err := ensureConversation(tx, userMessage.ConversationID); err != nil { @@ -49,7 +49,7 @@ func (store *Store) CreateChatMessages( if err := mapError(tx.Create(&userMessage).Error); err != nil { return err } - if err := mapError(tx.Create(&responderMessage).Error); err != nil { + if err := mapError(tx.Create(&responseMessage).Error); err != nil { return err } if beforeCommit != nil { @@ -60,7 +60,7 @@ func (store *Store) CreateChatMessages( if err != nil { return model.Message{}, err } - return responderMessage, nil + return responseMessage, nil } func (store *Store) ListRootMessagesByTask(ctx context.Context, taskID string) ([]model.Message, error) { diff --git a/server/internal/repo/repo_test.go b/server/internal/repo/repo_test.go index 919f2a8..3aab968 100644 --- a/server/internal/repo/repo_test.go +++ b/server/internal/repo/repo_test.go @@ -81,10 +81,10 @@ func TestCreateChatMessagesRollsBackBothRows(t *testing.T) { TaskID: "01991f3d-1114-7000-8000-000000000000", Kind: "user", Key: "user-1", Content: []byte(`{"version":"1.0","biz":"chat","meta":{},"blocks":[]}`), } - responder := existing - responder.TaskID = user.TaskID - if _, err := store.CreateChatMessages(ctx, user, responder, nil); err == nil { - t.Fatal("CreateChatMessages succeeded with a duplicate responder ID") + response := existing + response.TaskID = user.TaskID + if _, err := store.CreateChatMessages(ctx, user, response, nil); err == nil { + t.Fatal("CreateChatMessages succeeded with a duplicate response ID") } messages, err := store.ListMessages(ctx, conversationID, "", 100) if err != nil { diff --git a/server/internal/service/chat.go b/server/internal/service/chat.go index 5bc7554..1764f24 100644 --- a/server/internal/service/chat.go +++ b/server/internal/service/chat.go @@ -3,12 +3,16 @@ package service import ( "context" "encoding/json" + "errors" "fmt" "log/slog" "strings" + agentuerunner "github.com/compforge/agentue/sdks/go/runner" loopd "github.com/compforge/loopd" + "github.com/compforge/loopd/server/internal/delivery" "github.com/compforge/loopd/server/internal/model" + "github.com/compforge/loopd/server/internal/repo" "github.com/qiankunli/go-stdx/uuid" ) @@ -17,63 +21,91 @@ type ChatRepository interface { } type TaskClient interface { - Create(context.Context, string, loopd.ResponderRef) error + Create(context.Context, string, loopd.ActorRef) error Delete(context.Context, string) error } +type ChatDelivery interface { + Initialize(context.Context, string, json.RawMessage) error + Delete(context.Context, string) error + Emit(context.Context, string, json.RawMessage) (string, error) + Complete(context.Context, string, *delivery.Failure) error + Stream(context.Context, string, string, string, func(delivery.Event) error) error +} + // ChatService owns the transaction boundary for one user question. It writes // the visible message pair, creates the same-ID Task CRD before commit, and -// returns the responder message that will be updated as work progresses. +// returns the selected Actor's message that will be updated as work progresses. type ChatService struct { - repo ChatRepository - tasks TaskClient - logger *slog.Logger + repo ChatRepository + tasks TaskClient + delivery ChatDelivery + logger *slog.Logger } -func NewChatService(repository ChatRepository, tasks TaskClient, logger *slog.Logger) *ChatService { - return &ChatService{repo: repository, tasks: tasks, logger: loggerOrDefault(logger)} +func NewChatService(repository ChatRepository, tasks TaskClient, chatDelivery ChatDelivery, logger *slog.Logger) *ChatService { + return &ChatService{repo: repository, tasks: tasks, delivery: chatDelivery, logger: loggerOrDefault(logger)} } func (service *ChatService) Create( ctx context.Context, conversationID string, userKey string, - responder loopd.ResponderRef, + target loopd.ActorRef, content json.RawMessage, ) (loopd.Message, error) { userKey = strings.TrimSpace(userKey) - responder.Key = strings.TrimSpace(responder.Key) - if userKey == "" || !responder.Valid() || validateContent(content) != nil { + target.Key = strings.TrimSpace(target.Key) + if userKey == "" || !target.ValidTarget() || validateContent(content) != nil { return loopd.Message{}, ErrInvalid } - if service.tasks == nil { + if service.tasks == nil || service.delivery == nil { return loopd.Message{}, ErrUnavailable } taskID := uuid.V7() - responderContent, err := emptyContent(content) + responseContent, err := emptyContent(content) if err != nil { return loopd.Message{}, ErrInvalid } + userMessage := model.Message{ + ID: uuid.V7(), ConversationID: conversationID, TaskID: taskID, + Kind: string(loopd.RoleUser), Key: userKey, Content: content, + } + responseMessage := model.Message{ + ID: uuid.V7(), ConversationID: conversationID, TaskID: taskID, + Kind: string(target.Kind), Key: target.Key, Content: responseContent, + } taskCreated := false + streamCreated := false message, err := service.repo.CreateChatMessages(ctx, - model.Message{ - ID: uuid.V7(), ConversationID: conversationID, TaskID: taskID, - Kind: string(loopd.RoleUser), Key: userKey, Content: content, - }, - model.Message{ - ID: uuid.V7(), ConversationID: conversationID, TaskID: taskID, - Kind: string(responder.Kind), Key: responder.Key, Content: responderContent, - }, + userMessage, + responseMessage, func(txCtx context.Context) error { - if err := service.tasks.Create(txCtx, taskID, responder); err != nil { + if err := service.delivery.Initialize(txCtx, taskID, responseContent); err != nil { + service.logger.ErrorContext(ctx, "initialize chat event stream failed", + "conversation_id", conversationID, + "task_id", taskID, + "error", err, + ) + return fmt.Errorf("%w: %v", ErrUnavailable, err) + } + streamCreated = true + if err := service.tasks.Create(txCtx, taskID, target); err != nil { service.logger.ErrorContext(ctx, "create task CRD failed", "conversation_id", conversationID, "task_id", taskID, - "responder_kind", responder.Kind, - "responder_key", responder.Key, + "target_kind", target.Kind, + "target_key", target.Key, "error", err, ) + if cleanupErr := service.delivery.Delete(context.WithoutCancel(ctx), taskID); cleanupErr != nil { + service.logger.ErrorContext(ctx, "delete rolled back chat event stream failed", + "task_id", taskID, + "error", cleanupErr, + ) + } + streamCreated = false return fmt.Errorf("%w: %v", ErrUnavailable, err) } taskCreated = true @@ -97,18 +129,86 @@ func (service *ChatService) Create( ) } } + if err != nil && streamCreated { + if cleanupErr := service.delivery.Delete(context.WithoutCancel(ctx), taskID); cleanupErr != nil { + service.logger.ErrorContext(ctx, "delete rolled back chat event stream failed", + "task_id", taskID, + "error", cleanupErr, + ) + } + } if err == nil { service.logger.InfoContext(ctx, "chat task created", "conversation_id", conversationID, "task_id", taskID, - "responder_message_id", message.ID, - "responder_kind", responder.Kind, - "responder_key", responder.Key, + "actor_message_id", message.ID, + "target_kind", target.Kind, + "target_key", target.Key, ) } return messageFromModel(message), err } +func (service *ChatService) Stream( + ctx context.Context, + conversationID string, + taskID string, + after string, + deliver func(delivery.Event) error, +) error { + taskID = strings.TrimSpace(taskID) + if taskID == "" { + return ErrInvalid + } + service.logger.InfoContext(ctx, "chat stream opened", + "conversation_id", conversationID, + "task_id", taskID, + "after", after, + ) + err := mapDeliveryError(service.delivery.Stream(ctx, taskID, conversationID, after, deliver)) + service.logger.InfoContext(ctx, "chat stream closed", + "conversation_id", conversationID, + "task_id", taskID, + "error", err, + ) + return err +} + +func (service *ChatService) Emit(ctx context.Context, taskID string, event json.RawMessage) (string, error) { + taskID = strings.TrimSpace(taskID) + if taskID == "" { + return "", ErrInvalid + } + cursor, err := service.delivery.Emit(ctx, taskID, event) + if err == nil { + service.logger.DebugContext(ctx, "chat event published", "task_id", taskID, "cursor", cursor) + } + return cursor, mapDeliveryError(err) +} + +func (service *ChatService) Complete(ctx context.Context, taskID string, failure *delivery.Failure) error { + taskID = strings.TrimSpace(taskID) + if taskID == "" { + return ErrInvalid + } + return mapDeliveryError(service.delivery.Complete(ctx, taskID, failure)) +} + +func mapDeliveryError(err error) error { + switch { + case err == nil: + return nil + case errors.Is(err, agentuerunner.ErrNotFound): + return repo.ErrNotFound + case errors.Is(err, agentuerunner.ErrConflict): + return ErrConflict + case errors.Is(err, delivery.ErrInvalidEvent): + return fmt.Errorf("%w: %v", ErrInvalid, err) + default: + return err + } +} + func emptyContent(content json.RawMessage) (json.RawMessage, error) { var source struct { Version string `json:"version"` diff --git a/server/internal/service/message.go b/server/internal/service/message.go index 021b773..6499234 100644 --- a/server/internal/service/message.go +++ b/server/internal/service/message.go @@ -57,26 +57,6 @@ func (service *MessageService) CreateMessage( return messageFromModel(message), err } -func (service *MessageService) UpdateMessageContent( - ctx context.Context, - conversationID string, - messageID string, - content json.RawMessage, -) (loopd.Message, error) { - if validateContent(content) != nil { - return loopd.Message{}, ErrInvalid - } - message, err := service.repo.UpdateMessageContent(ctx, conversationID, messageID, content) - if err == nil { - service.logger.InfoContext(ctx, "message content updated", - "conversation_id", conversationID, - "message_id", message.ID, - "task_id", message.TaskID, - ) - } - return messageFromModel(message), err -} - func (service *MessageService) ListMessages( ctx context.Context, conversationID string, diff --git a/server/internal/service/message_test.go b/server/internal/service/message_test.go index a3a84e4..13337a7 100644 --- a/server/internal/service/message_test.go +++ b/server/internal/service/message_test.go @@ -9,6 +9,7 @@ import ( "testing" loopd "github.com/compforge/loopd" + "github.com/compforge/loopd/server/internal/delivery" "github.com/compforge/loopd/server/internal/model" "github.com/compforge/loopd/server/internal/repo" ) @@ -20,14 +21,14 @@ func TestChatCreatesVisibleMessagesWithOneTask(t *testing.T) { conversations := NewConversationService(store, nil) messages := NewMessageService(store, nil) tasks := &recordingTaskClient{} - chat := NewChatService(store, tasks, nil) + chat := NewChatService(store, tasks, nopChatRunner{}, nil) ctx := context.Background() conversation, err := conversations.CreateConversation(ctx, "Planning", "") if err != nil { t.Fatal(err) } - answer, err := chat.Create(ctx, conversation.ID, "user-1", loopd.ResponderRef{ + answer, err := chat.Create(ctx, conversation.ID, "user-1", loopd.ActorRef{ Kind: loopd.RoleHarness, Key: "harness-1", }, textContent("question")) @@ -41,7 +42,7 @@ func TestChatCreatesVisibleMessagesWithOneTask(t *testing.T) { if answer.Kind != loopd.RoleHarness || answer.Key != "harness-1" { t.Fatalf("answer identity = %s/%s", answer.Kind, answer.Key) } - if tasks.createdTaskID != answer.TaskID || tasks.createdTarget != (loopd.ResponderRef{ + if tasks.createdTaskID != answer.TaskID || tasks.createdTarget != (loopd.ActorRef{ Kind: loopd.RoleHarness, Key: "harness-1", }) { @@ -74,27 +75,21 @@ func TestChatCreatesVisibleMessagesWithOneTask(t *testing.T) { t.Fatalf("initial answer content = %s", answer.Content) } - updated, err := messages.UpdateMessageContent(ctx, conversation.ID, answer.ID, richContent()) - if err != nil { - t.Fatal(err) - } - if updated.ID != answer.ID || updated.UpdatedAt.Before(answer.UpdatedAt) { - t.Fatalf("updated answer = %#v", updated) - } } func TestChatRollsBackMessagesWhenTaskCreationFails(t *testing.T) { store := openServiceStore(t) conversations := NewConversationService(store, nil) messages := NewMessageService(store, nil) - chat := NewChatService(store, &recordingTaskClient{createErr: errors.New("api unavailable")}, nil) + streams := &recordingChatRunner{} + chat := NewChatService(store, &recordingTaskClient{createErr: errors.New("api unavailable")}, streams, nil) ctx := context.Background() conversation, err := conversations.CreateConversation(ctx, "Planning", "") if err != nil { t.Fatal(err) } - _, err = chat.Create(ctx, conversation.ID, "user-1", loopd.ResponderRef{ + _, err = chat.Create(ctx, conversation.ID, "user-1", loopd.ActorRef{ Kind: loopd.RoleOperator, Key: "operator-1", }, textContent("question")) @@ -108,13 +103,17 @@ func TestChatRollsBackMessagesWhenTaskCreationFails(t *testing.T) { if len(history) != 0 { t.Fatalf("Task creation failure left messages: %#v", history) } + if streams.initializedTaskID == "" || streams.deletedTaskID != streams.initializedTaskID { + t.Fatalf("initialized stream %q, deleted stream %q", streams.initializedTaskID, streams.deletedTaskID) + } } func TestChatDeletesTaskWhenDatabaseCommitFails(t *testing.T) { tasks := &recordingTaskClient{} - chat := NewChatService(failingCommitRepository{}, tasks, nil) + streams := &recordingChatRunner{} + chat := NewChatService(failingCommitRepository{}, tasks, streams, nil) - _, err := chat.Create(context.Background(), "conversation-1", "user-1", loopd.ResponderRef{ + _, err := chat.Create(context.Background(), "conversation-1", "user-1", loopd.ActorRef{ Kind: loopd.RoleOperator, Key: "operator-1", }, textContent("question")) @@ -124,20 +123,23 @@ func TestChatDeletesTaskWhenDatabaseCommitFails(t *testing.T) { if tasks.createdTaskID == "" || tasks.deletedTaskID != tasks.createdTaskID { t.Fatalf("created Task %q, deleted Task %q", tasks.createdTaskID, tasks.deletedTaskID) } + if streams.initializedTaskID == "" || streams.deletedTaskID != streams.initializedTaskID { + t.Fatalf("initialized stream %q, deleted stream %q", streams.initializedTaskID, streams.deletedTaskID) + } } -func TestDetailConversationReferencesResponderMessage(t *testing.T) { +func TestDetailConversationReferencesActorMessage(t *testing.T) { store := openServiceStore(t) conversations := NewConversationService(store, nil) messages := NewMessageService(store, nil) - chat := NewChatService(store, nopTaskClient{}, nil) + chat := NewChatService(store, nopTaskClient{}, nopChatRunner{}, nil) ctx := context.Background() root, err := conversations.CreateConversation(ctx, "Root", "") if err != nil { t.Fatal(err) } - answer, err := chat.Create(ctx, root.ID, "user-1", loopd.ResponderRef{ + answer, err := chat.Create(ctx, root.ID, "user-1", loopd.ActorRef{ Kind: loopd.RoleOperator, Key: "operator-1", }, textContent("question")) @@ -167,17 +169,66 @@ func TestDetailConversationReferencesResponderMessage(t *testing.T) { type nopTaskClient struct{} -func (nopTaskClient) Create(context.Context, string, loopd.ResponderRef) error { return nil } -func (nopTaskClient) Delete(context.Context, string) error { return nil } +func (nopTaskClient) Create(context.Context, string, loopd.ActorRef) error { return nil } +func (nopTaskClient) Delete(context.Context, string) error { return nil } + +type nopChatRunner struct{} + +func (nopChatRunner) Initialize(context.Context, string, json.RawMessage) error { return nil } +func (nopChatRunner) Delete(context.Context, string) error { return nil } +func (nopChatRunner) Emit(context.Context, string, json.RawMessage) (string, error) { + return "", nil +} + +type recordingChatRunner struct { + initializedTaskID string + deletedTaskID string +} + +func (runner *recordingChatRunner) Initialize(_ context.Context, taskID string, _ json.RawMessage) error { + runner.initializedTaskID = taskID + return nil +} + +func (runner *recordingChatRunner) Delete(_ context.Context, taskID string) error { + runner.deletedTaskID = taskID + return nil +} + +func (*recordingChatRunner) Emit(context.Context, string, json.RawMessage) (string, error) { + return "", nil +} + +func (*recordingChatRunner) Complete(context.Context, string, *delivery.Failure) error { return nil } + +func (*recordingChatRunner) Stream( + context.Context, + string, + string, + string, + func(delivery.Event) error, +) error { + return nil +} +func (nopChatRunner) Complete(context.Context, string, *delivery.Failure) error { return nil } +func (nopChatRunner) Stream( + context.Context, + string, + string, + string, + func(delivery.Event) error, +) error { + return nil +} type recordingTaskClient struct { createdTaskID string - createdTarget loopd.ResponderRef + createdTarget loopd.ActorRef deletedTaskID string createErr error } -func (client *recordingTaskClient) Create(_ context.Context, taskID string, target loopd.ResponderRef) error { +func (client *recordingTaskClient) Create(_ context.Context, taskID string, target loopd.ActorRef) error { client.createdTaskID = taskID client.createdTarget = target return client.createErr @@ -221,15 +272,3 @@ func textContent(text string) json.RawMessage { }) return value } - -func richContent() json.RawMessage { - return json.RawMessage(`{ - "version":"1.0", - "biz":"chat", - "meta":{}, - "blocks":[ - {"id":"answer","type":"text","content":"answer"}, - {"id":"tool-1","type":"tool","name":"search","status":"completed"} - ] - }`) -} diff --git a/server/internal/service/task_test.go b/server/internal/service/task_test.go index 76c9b0b..217e454 100644 --- a/server/internal/service/task_test.go +++ b/server/internal/service/task_test.go @@ -10,7 +10,7 @@ import ( func TestTaskContextComesFromMessages(t *testing.T) { store := openServiceStore(t) conversations := NewConversationService(store, nil) - chat := NewChatService(store, nopTaskClient{}, nil) + chat := NewChatService(store, nopTaskClient{}, nopChatRunner{}, nil) tasks := NewTaskService(store, nil) ctx := context.Background() @@ -18,7 +18,7 @@ func TestTaskContextComesFromMessages(t *testing.T) { if err != nil { t.Fatal(err) } - answer, err := chat.Create(ctx, conversation.ID, "user-1", loopd.ResponderRef{ + answer, err := chat.Create(ctx, conversation.ID, "user-1", loopd.ActorRef{ Kind: loopd.RoleOperator, Key: "operator-1", }, textContent("question")) diff --git a/server/internal/task/client.go b/server/internal/task/client.go index d0a7aa5..d68381a 100644 --- a/server/internal/task/client.go +++ b/server/internal/task/client.go @@ -29,7 +29,7 @@ func NewClient(kubeClient controllerclient.Client, namespace string, timeout tim return &Client{kubeClient: kubeClient, namespace: namespace, timeout: timeout} } -func (client *Client) Create(ctx context.Context, taskID string, target loopd.ResponderRef) error { +func (client *Client) Create(ctx context.Context, taskID string, target loopd.ActorRef) error { ctx, cancel := context.WithTimeout(ctx, client.timeout) defer cancel() diff --git a/server/internal/task/client_test.go b/server/internal/task/client_test.go index 84c63db..7875b94 100644 --- a/server/internal/task/client_test.go +++ b/server/internal/task/client_test.go @@ -18,7 +18,7 @@ func TestClientLifecycle(t *testing.T) { } kubeClient := fake.NewClientBuilder().WithScheme(scheme).Build() tasks := NewClient(kubeClient, "loopd-system", 0) - target := loopd.ResponderRef{Kind: loopd.RoleOperator, Key: "operator-1"} + target := loopd.ActorRef{Kind: loopd.RoleOperator, Key: "operator-1"} if err := tasks.Create(context.Background(), "01991f3d-1110-7000-8000-000000000000", target); err != nil { t.Fatal(err) diff --git a/server/redis.go b/server/redis.go new file mode 100644 index 0000000..60c9695 --- /dev/null +++ b/server/redis.go @@ -0,0 +1,19 @@ +package server + +import "time" + +type RedisConfig struct { + Address string + Username string + Password string + DB int + DialTimeout time.Duration + ReadTimeout time.Duration + WriteTimeout time.Duration + PoolSize int + MinIdleConns int + TaskTTL time.Duration // Defaults to 30 days and is refreshed when events arrive. + ReadBlock time.Duration // Defaults to one second. + ReadCount int64 // Defaults to 100 events per read. + KeyPrefix string // Defaults to loopd:agentue. +} diff --git a/server/server.go b/server/server.go index 9daf2a7..330b44c 100644 --- a/server/server.go +++ b/server/server.go @@ -4,17 +4,22 @@ package server import ( "context" "errors" + "fmt" "log/slog" "time" "github.com/cloudwego/hertz/pkg/route" + agentuerunner "github.com/compforge/agentue/sdks/go/runner" serverapi "github.com/compforge/loopd/server/internal/api" + "github.com/compforge/loopd/server/internal/delivery" "github.com/compforge/loopd/server/internal/repo" "github.com/compforge/loopd/server/internal/service" + "github.com/redis/go-redis/v9" ) type Config struct { Database DatabaseConfig + Redis RedisConfig Tasks TaskClient Logger *slog.Logger } @@ -30,6 +35,7 @@ type DatabaseConfig struct { type Server struct { store *repo.Store + redis redis.UniversalClient api *serverapi.Server } @@ -45,16 +51,73 @@ func New(config Config) (*Server, error) { if err != nil { return nil, err } + events, redisClient, err := newEventBridge(config.Redis) + if err != nil { + _ = store.Close() + return nil, err + } conversations := service.NewConversationService(store, config.Logger) messages := service.NewMessageService(store, config.Logger) - chat := service.NewChatService(store, config.Tasks, config.Logger) + chatDelivery := delivery.New(events, store, config.Logger) + chat := service.NewChatService(store, config.Tasks, chatDelivery, config.Logger) tasks := service.NewTaskService(store, config.Logger) return &Server{ store: store, + redis: redisClient, api: serverapi.New(conversations, messages, chat, tasks, config.Logger), }, nil } func (server *Server) Register(engine *route.Engine) { server.api.Register(engine) } func (server *Server) Run(context.Context) {} -func (server *Server) Close() error { return server.store.Close() } +func (server *Server) Close() error { + return errors.Join(server.redis.Close(), server.store.Close()) +} + +func newEventBridge(config RedisConfig) (agentuerunner.EventBridge, redis.UniversalClient, error) { + if config.Address == "" { + config.Address = "127.0.0.1:6379" + } + if config.DialTimeout <= 0 { + config.DialTimeout = 5 * time.Second + } + if config.ReadTimeout <= 0 { + config.ReadTimeout = 5 * time.Second + } + if config.WriteTimeout <= 0 { + config.WriteTimeout = 5 * time.Second + } + if config.PoolSize <= 0 { + config.PoolSize = 20 + } + if config.MinIdleConns < 0 { + config.MinIdleConns = 0 + } + if config.TaskTTL <= 0 { + config.TaskTTL = 30 * 24 * time.Hour + } + if config.ReadBlock <= 0 { + config.ReadBlock = time.Second + } + if config.ReadCount <= 0 { + config.ReadCount = 100 + } + if config.KeyPrefix == "" { + config.KeyPrefix = "loopd:agentue" + } + client := redis.NewClient(&redis.Options{ + Addr: config.Address, Username: config.Username, Password: config.Password, DB: config.DB, + DialTimeout: config.DialTimeout, ReadTimeout: config.ReadTimeout, WriteTimeout: config.WriteTimeout, + PoolSize: config.PoolSize, MinIdleConns: config.MinIdleConns, + }) + ctx, cancel := context.WithTimeout(context.Background(), config.DialTimeout) + defer cancel() + if err := client.Ping(ctx).Err(); err != nil { + _ = client.Close() + return nil, nil, fmt.Errorf("connect to loop-server Redis %q: %w", config.Address, err) + } + return agentuerunner.NewRedisEventBridge(client, agentuerunner.BridgeOptions{ + KeyPrefix: config.KeyPrefix, TaskTTL: config.TaskTTL, + ReadBlock: config.ReadBlock, ReadCount: config.ReadCount, + }), client, nil +} diff --git a/server/server_test.go b/server/server_test.go new file mode 100644 index 0000000..b9a54d5 --- /dev/null +++ b/server/server_test.go @@ -0,0 +1,30 @@ +package server + +import ( + "context" + "path/filepath" + "testing" + + "github.com/alicebob/miniredis/v2" + loopd "github.com/compforge/loopd" +) + +func TestNewConnectsConfiguredRedis(t *testing.T) { + redisServer := miniredis.RunT(t) + server, err := New(Config{ + Database: DatabaseConfig{Path: filepath.Join(t.TempDir(), "loopd.db")}, + Redis: RedisConfig{Address: redisServer.Addr()}, + Tasks: testTaskClient{}, + }) + if err != nil { + t.Fatal(err) + } + if err := server.Close(); err != nil { + t.Fatal(err) + } +} + +type testTaskClient struct{} + +func (testTaskClient) Create(context.Context, string, loopd.ActorRef) error { return nil } +func (testTaskClient) Delete(context.Context, string) error { return nil } diff --git a/server/task.go b/server/task.go index 2c249f3..42a431a 100644 --- a/server/task.go +++ b/server/task.go @@ -12,7 +12,7 @@ import ( // TaskClient is the loop-server boundary for managing Task CRDs inside the // chat submission work unit. type TaskClient interface { - Create(context.Context, string, loopd.ResponderRef) error + Create(context.Context, string, loopd.ActorRef) error Delete(context.Context, string) error }