diff --git a/VERSION b/VERSION index b056f41..2678ff8 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -0.0.24 +0.0.25 diff --git a/cmd/loop-server/config.go b/cmd/loop-server/config.go index 2e05f59..479f1e1 100644 --- a/cmd/loop-server/config.go +++ b/cmd/loop-server/config.go @@ -6,9 +6,12 @@ import ( "strconv" "strings" "time" + + "github.com/compforge/loopd/server" ) type config struct { + messageTTL time.Duration messageInlineBlocks int messageInlineBytes int messagePartBytes int @@ -55,6 +58,7 @@ func loadConfig() (config, error) { return config{}, fmt.Errorf("unsupported DATABASE_DRIVER %q", databaseDriver) } value := config{ + messageTTL: server.DefaultMessageTTL, address: envOr("SERVER_ADDRESS", ":8080"), databaseDriver: databaseDriver, databaseDSN: databaseDSN, @@ -71,6 +75,7 @@ func loadConfig() (config, error) { name string value *time.Duration }{ + {"MESSAGE_TTL", &value.messageTTL}, {"TASK_CLIENT_TIMEOUT", &value.taskClientTimeout}, {"HTTP_READ_TIMEOUT", &value.readTimeout}, {"HTTP_IDLE_TIMEOUT", &value.idleTimeout}, diff --git a/cmd/loop-server/config_test.go b/cmd/loop-server/config_test.go index 6edaf63..66e00de 100644 --- a/cmd/loop-server/config_test.go +++ b/cmd/loop-server/config_test.go @@ -15,6 +15,9 @@ func TestLoadConfigDefaults(t *testing.T) { if config.address != ":8080" || config.databaseDriver != "sqlite" || config.databaseDSN != "loopd.db" { t.Fatalf("config = %#v", config) } + if config.messageTTL != 24*time.Hour { + t.Fatalf("default message TTL = %s, want 24h", config.messageTTL) + } if config.redisAddress != "127.0.0.1:6379" || config.taskNamespace != "default" || config.taskClientTimeout != 10*time.Second { t.Fatalf("task config = %#v", config) } @@ -28,6 +31,7 @@ func TestLoadConfigUsesUnprefixedEnvironment(t *testing.T) { t.Setenv("REDIS_ADDRESS", "redis:6379") t.Setenv("TASK_NAMESPACE", "loopd-system") t.Setenv("HTTP_IDLE_TIMEOUT", "2m") + t.Setenv("MESSAGE_TTL", "48h") config, err := loadConfig() if err != nil { t.Fatal(err) @@ -38,6 +42,9 @@ func TestLoadConfigUsesUnprefixedEnvironment(t *testing.T) { if config.redisAddress != "redis:6379" || config.taskNamespace != "loopd-system" || config.idleTimeout != 2*time.Minute { t.Fatalf("runtime config = %#v", config) } + if config.messageTTL != 48*time.Hour { + t.Fatal("message TTL not parsed") + } } func TestLoadConfigRejectsInvalidDuration(t *testing.T) { @@ -86,6 +93,7 @@ func TestLoadConfigRejectsLegacyEnvironment(t *testing.T) { func clearConfigEnv(t *testing.T) { t.Helper() for _, name := range []string{ + "MESSAGE_TTL", "SERVER_ADDRESS", "DATABASE_DRIVER", "DATABASE_DSN", "REDIS_ADDRESS", "REDIS_USERNAME", "REDIS_PASSWORD", "TASK_NAMESPACE", "TASK_CLIENT_TIMEOUT", "HTTP_READ_TIMEOUT", "HTTP_IDLE_TIMEOUT", "SHUTDOWN_TIMEOUT", "LOOP_SERVER_MYSQL_DSN", "LOOP_SERVER_SQLITE_PATH", "LOOP_SERVER_ADDR", "LOOP_SERVER_REDIS_ADDR", diff --git a/cmd/loop-server/main.go b/cmd/loop-server/main.go index 643ff7f..de4fe5f 100644 --- a/cmd/loop-server/main.go +++ b/cmd/loop-server/main.go @@ -34,7 +34,8 @@ func run() error { return err } loopServer, err := server.New(server.Config{ - Database: server.DatabaseConfig{MessageInlineBlocks: config.messageInlineBlocks, MessageInlineBytes: config.messageInlineBytes, MessagePartBytes: config.messagePartBytes, Driver: config.databaseDriver, DSN: config.databaseDSN}, + MessageTTL: config.messageTTL, + Database: server.DatabaseConfig{MessageInlineBlocks: config.messageInlineBlocks, MessageInlineBytes: config.messageInlineBytes, MessagePartBytes: config.messagePartBytes, Driver: config.databaseDriver, DSN: config.databaseDSN}, Redis: server.RedisConfig{ Address: config.redisAddress, Username: config.redisUsername, diff --git a/deploy/k8s/loopd/templates/server.yaml b/deploy/k8s/loopd/templates/server.yaml index 9bb077c..bae8a46 100644 --- a/deploy/k8s/loopd/templates/server.yaml +++ b/deploy/k8s/loopd/templates/server.yaml @@ -73,6 +73,8 @@ spec: env: - name: SERVER_ADDRESS value: :8080 + - name: MESSAGE_TTL + value: {{ .Values.server.messageTTL | quote }} {{- if $mysqlEnabled }} - name: DATABASE_DRIVER value: mysql diff --git a/deploy/k8s/loopd/values.yaml b/deploy/k8s/loopd/values.yaml index e49203a..786b60e 100644 --- a/deploy/k8s/loopd/values.yaml +++ b/deploy/k8s/loopd/values.yaml @@ -11,6 +11,8 @@ rbac: create: true server: + # Inactivity timeout for streaming messages and Redis event retention. + messageTTL: 24h # v1 keeps one server replica even when an external MySQL is configured. replicaCount: 1 image: diff --git a/docs/kernel.md b/docs/kernel.md index 54875a9..51fb1f7 100644 --- a/docs/kernel.md +++ b/docs/kernel.md @@ -94,8 +94,8 @@ runtime 不把普通发言自动解释成 steer/followup,也不替 Operator 用户首次提交只创建真实的 user Message;Operator/Harness 回答、Ask、Confirm 在实际发起时 各自创建 Message,不预建空回答。一个执行循环可以接收多次发言,也可以多次发布阶段结果或回应。 -`task_id` 是 UI/Redis 流的身份:不带它提交新发言,带它 replay。它不对应通用 Task CRD 或 task 表, -也不决定消息是新业务工作、补充信息还是确认答复。显式卡片回复给 typed Verb 返回值, +提交返回真实消息及交付标识,页面以 Conv 独立订阅,以 Message 合并流式更新和恢复快照。 +交付标识不对应通用 Task CRD 或 task 表,也不决定消息是新业务工作、补充信息还是确认答复。显式卡片回复给 typed Verb 返回值, 普通发言交给 Operator 判断,不自动解释为批准。 每条 Message 独立寻址、更新和持久化;Speak 可以一次说完,也可以逐步输出后 End。 diff --git a/pkg/contract/chat.go b/pkg/contract/chat.go index 0699df4..a044ba8 100644 --- a/pkg/contract/chat.go +++ b/pkg/contract/chat.go @@ -21,10 +21,11 @@ const ( MessageStatusCompleted MessageStatus = "completed" MessageStatusFailed MessageStatus = "failed" MessageStatusCancelled MessageStatus = "cancelled" + MessageStatusExpired MessageStatus = "expired" ) func (status MessageStatus) Terminal() bool { - return status == MessageStatusCompleted || status == MessageStatusFailed || status == MessageStatusCancelled + return status == MessageStatusCompleted || status == MessageStatusFailed || status == MessageStatusCancelled || status == MessageStatusExpired } type Message struct { diff --git a/server/AGENTS.md b/server/AGENTS.md index 74c2d0f..7b14124 100644 --- a/server/AGENTS.md +++ b/server/AGENTS.md @@ -18,7 +18,10 @@ server/ │ └── harness.go # Harness Registry ├── internal/view/ # API 与 service 共用的 View Model;按领域拆文件,仅定义数据结构 ├── internal/domain/ # Human 消息的纯状态规则,不持有独立存储 -├── internal/delivery/ # Message 寻址与独立流、会话聚合交付及固化 +├── internal/component/ # 有生命周期的运行组件 +│ ├── message_gc.go # 随 server 启停的全局 Message 失活回收 +│ └── conv_listener.go # 随 stream 请求启停的单 Conv 监听 +├── internal/delivery/ # Message 输出固化与独立 Redis 流写入 ├── internal/migrations/ # 已有数据库的 Schema 迁移 ├── internal/model/ # GORM model;一张表一个 Go 文件 │ ├── conversation.go # conversations @@ -36,7 +39,7 @@ server/ │ ├── message.go # MessageService │ ├── message_enrichment.go # 页面消息富化;分页不变、直接引用与卡片投影 │ ├── actor.go # Operator/Harness 注册与 Actor 聚合发现 -│ ├── chat.go # ChatService;输入提交与 UI 流交付 +│ ├── chat.go # ChatService;输入提交与消息输出 │ ├── poll.go # DB 消息接收、提交后通知与重试 │ └── human.go # Human 消息交互、持久到期与类型化答复 └── docs/ # 消息消费、可见事实持久化与用户交互的领域设计 diff --git a/server/const.go b/server/const.go new file mode 100644 index 0000000..826b744 --- /dev/null +++ b/server/const.go @@ -0,0 +1,7 @@ +package server + +import "time" + +// DefaultMessageTTL bounds inactive output and Redis event retention. Each store +// renews independently; Redis eviction never determines a message's SQL status. +const DefaultMessageTTL = 24 * time.Hour diff --git a/server/docs/persistence.md b/server/docs/persistence.md index d53b246..482abae 100644 --- a/server/docs/persistence.md +++ b/server/docs/persistence.md @@ -50,10 +50,10 @@ revision 表示可见快照版本,流式输出对应 AgentUE seq,Human 状 每条消息分别更新,不能因为 block ID 相同就跨消息合并。持久化与 replay 见 [页面交付](ue.md#页面交付)。 -task_id 仅保存在开启 UI/Redis 交付的真实用户 input 上,其他 Actor 发言不需要关联它。 -不再保存页面关闭意图。Message.status 列记录 streaming/completed/failed/cancelled, +task_id 仅保存在真实用户 input 上作为提交交付标识,其他 Actor 发言不需要关联它;页面流不依赖该字段。 +不再保存页面关闭意图。Message.status 列记录 streaming/completed/failed/cancelled/expired, 只表示这条消息的发送状态,不表示业务完成。默认 Speak、用户输入和 Human 卡片直接 completed; -流式输出从 streaming 开始,End 的终态与 Revision 一起保存。 +流式输出从 streaming 开始,End 的终态与 Revision 一起保存;长期失活由 server 按 updated_at + TTL 收口为 expired。 受控 meta.output 只保存最后一次事件指纹,用于辨别响应丢失后的重试,不承担执行检查点。 output、human_request、human_reply 分别表达普通输出、交互问题和卡片答复,不指定唯一主回答。 @@ -132,7 +132,8 @@ reply_to_id 是答复关联的唯一依据,不能用最近消息、相邻位 UUIDv7 的时间顺序不是多节点数据库的全局提交顺序;当前采用人类输入通常有先后的假设, 严格消费顺序的限制见 [Conversation](conversation.md)。 -created_at 与 updated_at 表达首次到最后一次可见活动。Harness 事件携带时间戳, -完成投影不把每条消息的结束时间改成整个页面流的完成时间,重试不缩短活动区间。 +created_at 与 updated_at 表达首次到最后一次可见活动。updated_at 使用 server 实际接受新输出的时间, +不跟随 Harness 的历史或未来时间戳;读取与幂等重试不刷新它。过期只改变 status/revision, +保留最后活动时间与正文。DB 和 Redis 的 TTL 独立推进,容许短暂不一致,详见 [失活与 TTL](ue.md#失活与-ttl)。 时间区间如何用于并行展示见 [消息呈现](ue.md#消息呈现),不由存储层规定页面布局。 diff --git a/server/docs/ue.md b/server/docs/ue.md index 93fc4bc..33e42c3 100644 --- a/server/docs/ue.md +++ b/server/docs/ue.md @@ -71,20 +71,22 @@ service 负责批量读取关联并组装富化结果,api 负责 HTTP 交付 ## 页面交付 -task_id 标识一次页面交付及 Redis 流,不是 Operator 的业务任务,server 不建立 tasks 表。 +task_id 是输入提交的交付标识,不是 Operator 的业务任务;server 不建立 tasks 表。 +页面订阅以 Conv 寻址,Redis 事件流以 Message 寻址。 ### 提交与观察 -用户提交只创建真实 Message 和对应页面流,不预建空回答。Operator/Harness 发言、 +用户提交只创建真实 Message,不预建空回答。Operator/Harness 发言、 Ask/Confirm 在实际发生时各自创建 Message。人可以连续追加,Operator 可以多次回应, 输入与输出数量没有一对一约束。 消息提交只依赖 DB,Redis 不进入输入事务。DB 接收后,即使页面桥暂时不可用,也不要求用户 重新发送。Conv 通知用同事务保存的待通知标记在提交后重试;消费契约见 [Conversation](conversation.md)。 -首次提交在连接页面桥前返回已接受的消息与 task_id;随后断线按该身份重连,不另建输入。 +提交接口返回已接受的消息后结束响应。观察页面使用 +`GET /v1/conversations/:conversation_id/stream`,不需要先发送消息,也不需要 task_id。 +主对话和当前右侧详情分别订阅自身 Conv,不隐式订阅所有子会话;再次发言不替换订阅。 +HTTP/SSE 断开不取消执行。 -带 task_id 的请求只观察其所属会话,不创建输入或通知。一个页面订阅覆盖 User conv 及其直接 -内部会话,包括不同 task_id 或无 task_id 的后续发言。HTTP/SSE 断开不取消执行。 Operator 通过 Poll 接收消息、Read 读取历史;不提供按 task_id 配对输入与回答的业务入口。 ### 消息寻址与快照 @@ -102,33 +104,52 @@ Human 问题与答复由 typed Verb 管理,普通流式写入不能伪造批 桥连续时交付增量,发生版本缺口或乱序时发送最新快照。页面发现晚到的发言和内容, 不依赖写入时命中了哪个 server 实例。 -### 聚合流与 replay +### 聚合流与恢复 -- 有消息身份的事件用 message_id、message、event 外层寻址,客户端按 ID/revision 合并。 -- 没有消息身份的 start/ping 是 UI 连接控制事件,不创建气泡,也没有业务结束信号。 -- Last-Event-ID 是聚合控制流位置;各 Message 在重连时独立重放。 -- 一条 Message 的 end 不关闭 UI 流;已结束消息和 Human 卡片以 SQL 快照交付。 -- 页面仅保留所选会话的一条订阅;再次发言替换连接身份,但仍可收到旧输入对应的后续输出。 -- 切换会话/离开页面主动取消连接,正常 EOF 或断线按保存的 task_id 退避重连。 +页面先分页读取历史,每次 stream 请求创建一个 Conv Listener,聚合当前运行态消息的独立 Redis 流;每个事件仍用 +message_id、message、event 寻址,客户端按消息 ID/revision 合并。AgentUE seq 和 Redis cursor +只在单条消息内有意义,不充当共享 Conv 游标。连接的 ping 不创建消息气泡。 -任一 server 实例都可观察同一交付。Redis 丢失后,已接受的内容可以从 SQL 快照恢复, +server 按 ID 定期增量发现该 Conv 的新消息。新的一次性发言直接交付,流式发言加入监听; +状态检查只读取运行态消息的元数据,revision 变化或增量缺口才加载正文快照,不重复扫描终态历史。 +Ask/Confirm 已发送的卡片仍可能待答,因此其交互状态独立观察。 + +一条 Message 的 end 移除自身监听,不关闭 Conv 连接。切换会话或离开页面主动取消连接; +断线退避重连,从 SQL 恢复运行态快照再接 Redis。每次连接建立后,页面做一次有界的活跃消息 +revision 查询及新增消息查询,补偿首次加载的时间差和断线期间的变化,避免刚结束的消息被遗漏。 +连接期间的增量发现和状态校验由 Listener 承担;浏览器不另开常驻消息轮询。 +Listener 随请求取消,不放入全局注册表,也不负责消息 GC。 + +任一 server 实例都可观察同一 Conv。Redis 丢失后,已接受的内容可以从 SQL 快照恢复, 但不会重新生成每个中间增量;AgentUE Bridge 负责事件协议和续接,server 负责消息寻址与快照。 ### 消息结束与重试 + 默认 Speak 在创建事务中保存完整正文与结束状态。流式 End 与内容事件使用同一顺序和重试契约: 先原子推进 SQL Revision 与 Message.status,再尽力更新消息桥并标记终态。 SQL 失败由句柄重试原事件,不另分配 seq;重复 End 幂等。 -普通 Speak 的内容和 Emit 不能更改消息终态;只有 End 可以结束流式消息。 +普通 Speak 的内容和 Emit 不能更改消息终态;写入者通过 End 结束流式消息,server 还会收口长期失活的输出。 Message.status 表达单条消息的发送生命周期,独立于 AgentUE 内容:streaming 表示仍在输出, -completed 表示发送完成,failed 表示输出失败,cancelled 表示输出被取消。 +completed 表示发送完成,failed 表示输出失败,cancelled 表示输出被取消,expired 表示长期未更新。 End 默认 completed,也可显式传入 failed/cancelled;不同终态不能互相覆盖。 更新请求将 status 与 AgentUE event 并列传入,AgentUE End 本身不携带状态;页面事件的 Message 外层与历史 API 都返回持久化 status。Redis 丢失后也不会把已结束消息重新视为正在输出。 主对话和详情只对 streaming 消息显示“生成中”,不把“连接在线”误标为“Operator 正在执行”。 Ask/Confirm 卡片发送完成即 completed,但交互仍可等待答复;两种生命周期互不替代。 -没有 Delivery.Complete、输入关闭意图或通用页面收尾维护循环。页面拥有订阅生命周期; +没有 Delivery.Complete 或输入关闭意图。页面拥有订阅生命周期; Operator 只表达自己何时说完一条消息。End 不删除 Conv、不自动 Commit、不终止待答问题, 也不禁止任何 Actor 用新 Key 再次发言。 + +### 失活与 TTL + +`MESSAGE_TTL` 统一配置输出失活期限与 Redis 事件保留期限,默认 24h。 +Message GC 随 server 启停,即使没有页面连接也独立执行有界清理。 +它按 DB 的 `updated_at + TTL` 定期将 streaming 消息标记为 expired,并递增 revision; +无需 expires_at 列。保留最后正文与最后活动时间,页面展示“已过期”,迟到写入不能恢复该消息。 + +Redis 按自己的写入时间续期,DB 按 server 实际接受输出的时间续期;读取、心跳和重复事件 +不延长 DB 生命周期。两层允许短暂不一致,不根据 Redis key 是否存在推断消息状态。 +过期只结束页面消息,不取消 Harness 执行、不代替 Operator 判断业务完成;新发言使用新消息。 diff --git a/server/internal/api/chat.go b/server/internal/api/chat.go index 62d7e58..4b5fe65 100644 --- a/server/internal/api/chat.go +++ b/server/internal/api/chat.go @@ -3,13 +3,11 @@ package api import ( "context" "encoding/json" - "errors" hertzapp "github.com/cloudwego/hertz/pkg/app" hertzsse "github.com/cloudwego/hertz/pkg/protocol/sse" ui "github.com/compforge/agentue/sdks/go/ui" "github.com/compforge/loopd/pkg/contract" - "github.com/compforge/loopd/server/internal/delivery" "github.com/compforge/loopd/server/internal/view" ) @@ -21,80 +19,50 @@ func (server *Server) createChatMessages(ctx context.Context, request *hertzapp. return err } conversationID := request.Param("conversation_id") - taskID := input.TaskID - var accepted *contract.Message - if taskID == "" { - if server.Human != nil { - identity, err := server.identity(ctx, request) - if err != nil { - return err - } - input.UserKey = identity - } - message, err := server.chat.Create(ctx, conversationID, input.UserKey, input.Target, input.Content) + if server.Human != nil { + identity, err := server.identity(ctx, request) if err != nil { return err } - taskID = message.TaskID - accepted = &message + input.UserKey = identity + } + message, err := server.chat.Create(ctx, conversationID, input.UserKey, input.Target, input.Content) + if err != nil { + return err } + accepted := &message + taskID := message.TaskID request.Response.Header.Set(taskIDHeader, taskID) - var writer *hertzsse.Writer - // DB acceptance must reach the client before opening the page bridge. - // Otherwise a transient Redis failure could look like a rejected input and - // cause the user to resend instead of reconnecting with this task ID. - if accepted != nil { - start, err := ui.Start(accepted.Content, accepted.Revision) - if err != nil { - return err - } - raw, err := start.Marshal() - if err != nil { - return err - } - data, err := server.messageEventData(ctx, accepted.ID, accepted, raw) - if err != nil { - return err - } - writer = hertzsse.NewWriter(request) - if err := writer.WriteEvent("", "", data); err != nil { - server.logger.WarnContext(ctx, "input accepted but page disconnected", "task_id", taskID, "error", err) - _ = writer.Close() - return nil - } + // Input acknowledgement is independent of the Conv listener and Redis. + start, err := ui.Start(accepted.Content, accepted.Revision) + if err != nil { + return err + } + raw, err := start.Marshal() + if err != nil { + return err + } + data, err := server.messageEventData(ctx, accepted.ID, accepted, raw) + if err != nil { + return err } - streamErr := server.chat.Stream( - ctx, - conversationID, - taskID, - hertzsse.GetLastEventID(&request.Request), - func(event delivery.Event) error { - if writer == nil { - writer = hertzsse.NewWriter(request) - } - data := event.Data - if event.MessageID != "" { - var err error - data, err = server.messageEventData(ctx, event.MessageID, event.Message, event.Data) - if err != nil { - return err - } - } - return writer.WriteEvent(event.ID, "", data) - }, - ) - if errors.Is(streamErr, context.Canceled) { - streamErr = nil + writer := hertzsse.NewWriter(request) + if err := writer.WriteEvent("", "", data); err != nil { + server.logger.WarnContext(ctx, "input accepted but page disconnected", "task_id", taskID, "error", err) + _ = writer.Close() + return nil } - if writer == nil { - return streamErr + // Input submission is acknowledged once. The page independently subscribes + // to its Conv; it does not need an input-owned connection to observe actors. + end, _ := ui.End(accepted.Revision).Marshal() + endData, err := server.messageEventData(ctx, accepted.ID, accepted, end) + if err == nil { + err = writer.WriteEvent("", "", endData) } - 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) + if err != nil { + server.logger.WarnContext(ctx, "input acknowledgement ended early", "message_id", accepted.ID, "error", err) } + _ = writer.Close() return nil } diff --git a/server/internal/api/conversation.go b/server/internal/api/conversation.go index 3d94b1c..35fc0f2 100644 --- a/server/internal/api/conversation.go +++ b/server/internal/api/conversation.go @@ -2,10 +2,13 @@ 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/pkg/contract" + "github.com/compforge/loopd/server/internal/component" "github.com/compforge/loopd/server/internal/view" ) @@ -56,3 +59,35 @@ func (server *Server) listConversations(ctx context.Context, request *hertzapp.R request.JSON(consts.StatusOK, view.Page[contract.Conversation]{Data: conversations}) return nil } + +// streamConversation multiplexes message-addressed events for exactly one Conv. +func (server *Server) streamConversation(ctx context.Context, request *hertzapp.RequestContext) error { + convID := request.Param("conversation_id") + if _, err := server.conversations.GetConversation(ctx, convID); err != nil { + return err + } + var writer *hertzsse.Writer + err := server.Listen(ctx, convID, func(event component.Event) error { + data := event.Data + if event.MessageID != "" { + var err error + data, err = server.messageEventData(ctx, event.MessageID, event.Message, data) + if err != nil { + return err + } + } + if writer == nil { + writer = hertzsse.NewWriter(request) + } + // Redis cursors are message-local; there is no shared Conv event cursor. + return writer.WriteEvent("", "", data) + }) + if writer == nil { + return err + } + closeErr := writer.Close() + if !errors.Is(err, context.Canceled) && errors.Join(err, closeErr) != nil { + server.logger.WarnContext(ctx, "conversation stream stopped", "conversation_id", convID, "error", errors.Join(err, closeErr)) + } + return nil +} diff --git a/server/internal/api/message.go b/server/internal/api/message.go index 1c498a5..72ca421 100644 --- a/server/internal/api/message.go +++ b/server/internal/api/message.go @@ -4,9 +4,11 @@ import ( "context" "fmt" "strconv" + "strings" hertzapp "github.com/cloudwego/hertz/pkg/app" "github.com/cloudwego/hertz/pkg/protocol/consts" + "github.com/compforge/loopd/pkg/contract" "github.com/compforge/loopd/server/internal/service" "github.com/compforge/loopd/server/internal/view" ) @@ -18,9 +20,16 @@ func (server *Server) listMessages(ctx context.Context, request *hertzapp.Reques if err != nil { return err } - messages, err := server.messages.ListMessages( - ctx, request.Param("conversation_id"), string(request.Query("after")), limit, - ) + var messages []contract.Message + if watch := request.Query("watch"); watch != "" { + revisions, parseErr := parseMessageWatch(watch) + if parseErr != nil { + return parseErr + } + messages, err = server.messages.MessageChanges(ctx, request.Param("conversation_id"), revisions) + } else { + messages, err = server.messages.ListMessages(ctx, request.Param("conversation_id"), request.Query("after"), limit) + } if err != nil { return err } @@ -32,6 +41,25 @@ func (server *Server) listMessages(ctx context.Context, request *hertzapp.Reques return nil } +// watch is a bounded set of known active message IDs and their last revisions. +// It does not advance the independent new-message discovery cursor. +func parseMessageWatch(raw string) (map[string]uint64, error) { + items := strings.Split(raw, ",") + if len(items) > 100 { + return nil, service.ErrInvalid + } + result := make(map[string]uint64, len(items)) + for _, item := range items { + id, revision, ok := strings.Cut(item, ":") + value, err := strconv.ParseUint(revision, 10, 64) + if !ok || id == "" || err != nil { + return nil, service.ErrInvalid + } + result[id] = value + } + return result, 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 73a4b07..2deed9c 100644 --- a/server/internal/api/server.go +++ b/server/internal/api/server.go @@ -11,12 +11,14 @@ import ( hertzapp "github.com/cloudwego/hertz/pkg/app" "github.com/cloudwego/hertz/pkg/protocol/consts" "github.com/cloudwego/hertz/pkg/route" + "github.com/compforge/loopd/server/internal/component" "github.com/compforge/loopd/server/internal/repo" "github.com/compforge/loopd/server/internal/service" "github.com/compforge/loopd/server/internal/view" ) type Server struct { + Listen func(context.Context, string, func(component.Event) error) error Poll *service.PollService Human *service.HumanService HumanIdentity HumanIdentity @@ -51,6 +53,7 @@ func (server *Server) Register(engine *route.Engine) { engine.GET("/v1/conversations", server.adapt(server.listConversations)) engine.GET("/v1/conversations/:conversation_id", server.adapt(server.getConversation)) engine.GET("/v1/conversations/:conversation_id/messages", server.adapt(server.listMessages)) + engine.GET("/v1/conversations/:conversation_id/stream", server.adapt(server.streamConversation)) engine.POST("/v1/conversations/:conversation_id/poll", server.adapt(server.pollConversation)) engine.POST("/v1/conversations/:conversation_id/commit", server.adapt(server.commitConversation)) engine.POST("/v1/conversations/:conversation_id/speak", server.adapt(server.publishMessage)) diff --git a/server/internal/api/server_test.go b/server/internal/api/server_test.go index 66c40ab..b1fbf3e 100644 --- a/server/internal/api/server_test.go +++ b/server/internal/api/server_test.go @@ -16,7 +16,8 @@ import ( "github.com/cloudwego/hertz/pkg/route" "github.com/cloudwego/hertz/pkg/route/param" "github.com/compforge/loopd/pkg/contract" - "github.com/compforge/loopd/server/internal/delivery" + "github.com/compforge/loopd/server/internal/component" + "github.com/compforge/loopd/server/internal/model" "github.com/compforge/loopd/server/internal/repo" "github.com/compforge/loopd/server/internal/service" "github.com/compforge/loopd/server/internal/view" @@ -143,10 +144,47 @@ func TestChatHTTPFlow(t *testing.T) { 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 convStreamRunner struct { + completedChatRunner + t *testing.T + convID string +} + +func (runner convStreamRunner) Listen(_ context.Context, convID string, deliver func(component.Event) error) error { + if convID != runner.convID { + runner.t.Fatalf("stream conv = %q", convID) + } + m := contract.Message{ID: "message", ConversationID: convID, Status: contract.MessageStatusStreaming, Kind: contract.ActorKindOperator, Key: "router", Content: json.RawMessage(`{"version":"1.1","biz":"chat","meta":{},"blocks":[]}`)} + return deliver(component.Event{MessageID: m.ID, Message: &m, Data: json.RawMessage(`{"op":"start","seq":1,"model":{"version":"1.1","biz":"chat","meta":{},"blocks":[]}}`)}) +} + +func TestConversationStreamHTTPWithoutUserInput(t *testing.T) { + store, err := repo.Open(repo.Config{Driver: "sqlite", DSN: filepath.Join(t.TempDir(), "stream.db")}) + if err != nil { + t.Fatal(err) + } + defer store.Close() + if _, err := store.CreateConversation(context.Background(), model.Conversation{ID: "conv"}); err != nil { + t.Fatal(err) + } + server := New(service.NewActorService(store, nil), service.NewConversationService(store, nil), service.NewMessageService(store, nil), + service.NewChatService(store, completedChatRunner{}, nil, nil), nil) + server.Listen = (convStreamRunner{t: t, convID: "conv"}).Listen + engine := route.NewEngine(config.NewOptions(nil)) + server.Register(engine) + request := hertzapp.NewContext(1) + request.Request.Header.SetMethod("GET") + request.Request.SetRequestURI("/v1/conversations/conv/stream") + writer := &streamWriter{} + request.Response.HijackWriter(writer) + engine.ServeHTTP(context.Background(), request) + if request.Response.StatusCode() != 200 || !strings.Contains(writer.String(), `"message_id":"message"`) || !strings.Contains(string(request.Response.Header.ContentType()), "text/event-stream") { + t.Fatalf("stream response: %d %s", request.Response.StatusCode(), writer.String()) + } + missing := ut.PerformRequest(engine, "GET", "/v1/conversations/missing/stream", nil).Result() + if missing.StatusCode() != 404 { + t.Fatalf("missing conv status=%d", missing.StatusCode()) + } } type streamWriter struct{ bytes.Buffer } @@ -167,18 +205,6 @@ func performChat(t *testing.T, server *Server, conversationID, body string) (str } return string(request.Response.Header.Peek(taskIDHeader)), writer.String() } -func (completedChatRunner) Stream( - _ context.Context, - _ string, - _ string, - _ string, - deliver func(delivery.Event) error, -) error { - if err := deliver(delivery.Event{ID: "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{ID: "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() @@ -193,11 +219,7 @@ func (completedChatRunner) EmitMessage(context.Context, string, json.RawMessage, type unavailableChatRunner struct{ completedChatRunner } -func (unavailableChatRunner) Stream(context.Context, string, string, string, func(delivery.Event) error) error { - return errors.New("page bridge unavailable") -} - -// +case=`An accepted input returns its message and replay identity even when opening the page stream fails.` +// +case=`An accepted input returns its message and receipt even when opening the page stream fails.` func TestChatAcknowledgesInputBeforePageBridge(t *testing.T) { store, err := repo.Open(repo.Config{Driver: "sqlite", DSN: filepath.Join(t.TempDir(), "chat.db")}) if err != nil { @@ -206,6 +228,9 @@ func TestChatAcknowledgesInputBeforePageBridge(t *testing.T) { defer store.Close() server := New(service.NewActorService(store, nil), service.NewConversationService(store, nil), service.NewMessageService(store, nil), service.NewChatService(store, unavailableChatRunner{}, nil, nil), nil) + server.Listen = func(context.Context, string, func(component.Event) error) error { + return errors.New("page bridge unavailable") + } engine := route.NewEngine(config.NewOptions(nil)) server.Register(engine) created := performJSON(t, engine, "POST", "/v1/conversations", `{"name":"offline bridge"}`) diff --git a/server/internal/component/conv_listener.go b/server/internal/component/conv_listener.go new file mode 100644 index 0000000..0716110 --- /dev/null +++ b/server/internal/component/conv_listener.go @@ -0,0 +1,294 @@ +package component + +import ( + "context" + "encoding/json" + "errors" + "time" + + agentuerunner "github.com/compforge/agentue/sdks/go/runner" + agentueui "github.com/compforge/agentue/sdks/go/ui" + "github.com/compforge/loopd/pkg/contract" + "github.com/compforge/loopd/server/internal/model" + "github.com/compforge/loopd/server/internal/repo" +) + +const discoveryPageSize = 100 + +type Event struct { + MessageID string + Message *contract.Message + Data json.RawMessage +} + +type ConvMessageRepository interface { + LatestMessageID(context.Context, string) (string, error) + ListStreamMessages(context.Context, string, string, string, int) ([]model.Message, error) + GetMessageStates(context.Context, string, []string) ([]repo.MessageState, error) + GetMessage(context.Context, string) (model.Message, error) +} + +type messageReader struct { + message model.Message + cancel context.CancelFunc +} +type messageDelivery struct { + reader *messageReader + delivery agentuerunner.Delivery + done bool +} +type watchedMessage struct { + message model.Message // Only active messages remain in this map. + revision uint64 + reader *messageReader + retryAt time.Time +} + +// ConvListener belongs to one stream request, never to an actor execution. +// +spec=`Only this Conv is observed. GC runs independently; terminal messages leave the active set.` +type ConvListener struct { + repo ConvMessageRepository + events agentuerunner.EventBridge + conversationID string +} + +func NewConvListener(events agentuerunner.EventBridge, repository ConvMessageRepository, conversationID string) *ConvListener { + return &ConvListener{events: events, repo: repository, conversationID: conversationID} +} + +func (listener *ConvListener) Run(ctx context.Context, deliver func(Event) error) error { + conversationID := listener.conversationID + watermark, err := listener.repo.LatestMessageID(ctx, conversationID) + if err != nil { + return err + } + ctx, cancel := context.WithCancel(ctx) + defer cancel() + incoming := make(chan messageDelivery) + watching := map[string]*watchedMessage{} + stop := func(watch *watchedMessage) { + if watch.reader != nil { + watch.reader.cancel() + watch.reader = nil + } + } + start := func(watch *watchedMessage) { + if watch.reader != nil || time.Now().Before(watch.retryAt) { + return + } + watch.retryAt = time.Now().Add(time.Second) + readerCtx, cancelReader := context.WithCancel(ctx) + reader := &messageReader{message: watch.message, cancel: cancelReader} + watch.reader = reader + go func() { + defer cancelReader() + send := func(value messageDelivery) error { + select { + case incoming <- value: + return nil + case <-readerCtx.Done(): + return readerCtx.Err() + } + } + // Readers never provision Redis keys. Missing streams are repaired from SQL. + _ = (agentuerunner.Replayer{Bridge: listener.events}).Stream(readerCtx, "message/"+reader.message.ID, "", func(value agentuerunner.Delivery) error { + return send(messageDelivery{reader: reader, delivery: value}) + }) + _ = send(messageDelivery{reader: reader, done: true}) + }() + } + snapshot := func(row model.Message, watch *watchedMessage) error { + revision := row.Revision + if revision == 0 { + revision = 1 + } + if watch.revision >= revision { + return nil + } + patch, err := agentueui.Start(row.Content, revision) + if err != nil { + return err + } + data, err := patch.Marshal() + if err != nil { + return err + } + message := visibleMessage(row) + if err := deliver(Event{MessageID: row.ID, Message: &message, Data: data}); err != nil { + return err + } + watch.revision = revision + if row.Purpose == "output" && message.Ended() { + data, err := agentueui.End(revision).Marshal() + if err != nil { + return err + } + return deliver(Event{MessageID: row.ID, Message: &message, Data: data}) + } + return nil + } + observe := func(row model.Message) error { + watch := watching[row.ID] + if watch == nil { + watch = &watchedMessage{} + } + if err := snapshot(row, watch); err != nil { + return err + } + watch.message = row + if !visibleMessage(row).Ended() || row.HumanDueAt != nil { + watching[row.ID] = watch + if row.Purpose == "output" && !visibleMessage(row).Ended() { + start(watch) + } + } else { + stop(watch) + delete(watching, row.ID) + } + return nil + } + afterID := "" + discover := func() (int, error) { + // UUIDv7 is an allocation order, not a distributed commit watermark. + // This uses the existing conversation cursor assumption, not Kafka's stronger guarantee. + rows, err := listener.repo.ListStreamMessages(ctx, conversationID, afterID, watermark, discoveryPageSize) + if err != nil { + return 0, err + } + for _, row := range rows { + if err := observe(row); err != nil { + return 0, err + } + afterID = row.ID + } + return len(rows), nil + } + refreshActive := func() error { + ids := make([]string, 0, len(watching)) + for id := range watching { + ids = append(ids, id) + } + for from := 0; from < len(ids); from += discoveryPageSize { + batch := ids[from:min(from+discoveryPageSize, len(ids))] + states, err := listener.repo.GetMessageStates(ctx, conversationID, batch) + if err != nil { + return err + } + found := map[string]bool{} + for _, state := range states { + found[state.ID] = true + watch := watching[state.ID] + if state.Revision > watch.revision || state.Status != visibleMessage(watch.message).Status { + row, err := listener.repo.GetMessage(ctx, state.ID) + if errors.Is(err, repo.ErrNotFound) { + stop(watch) + delete(watching, state.ID) + continue + } + if err != nil { + return err + } + if err := observe(row); err != nil { + return err + } + } else if state.Ended && state.HumanDueAt == nil { + stop(watch) + delete(watching, state.ID) + } else if watch.message.Purpose == "output" { + start(watch) + } + } + for _, id := range batch { + if !found[id] { + stop(watching[id]) + delete(watching, id) + } + } + } + return nil + } + // Conv subscriptions bootstrap active outputs and pending Human cards only. + // Completed history belongs to the paginated messages endpoint. + for { + count, err := discover() + if err != nil { + return err + } + if count < discoveryPageSize { + break + } + } + if afterID < watermark { + afterID = watermark + } + ticker := time.NewTicker(time.Second) + defer ticker.Stop() + heartbeat := time.NewTicker(15 * time.Second) + defer heartbeat.Stop() + ping, _ := agentueui.Ping(0).Marshal() + if err := deliver(Event{Data: ping}); err != nil { + return err + } + for { + select { + case <-ctx.Done(): + return ctx.Err() + case <-heartbeat.C: + if err := deliver(Event{Data: ping}); err != nil { + return err + } + case <-ticker.C: + if _, err := discover(); err != nil { + return err + } + if err := refreshActive(); err != nil { + return err + } + case value := <-incoming: + id := value.reader.message.ID + watch := watching[id] + // Drop queued events from readers cancelled after a SQL terminal snapshot. + if watch == nil || watch.reader != value.reader { + continue + } + if value.done { + stop(watch) + continue + } + message := visibleMessage(watch.message) + patch, err := agentueui.Parse(value.delivery.Data) + if err != nil { + return err + } + + if patch.Op == agentueui.OpEnd || (patch.Op == agentueui.OpStart && patch.Seq > watch.revision) { + // A bridge rebuilt from SQL may already represent a terminal + // snapshot. Its AgentUE Start alone cannot carry Message.status. + row, err := listener.repo.GetMessage(ctx, id) + if err != nil { + return err + } + if err := observe(row); err != nil { + return err + } + continue + } + if patch.Op == agentueui.OpPing || patch.Seq <= watch.revision { + continue + } + // Live deltas may overlap a newer SQL snapshot, but cannot jump a gap. + if patch.Op != agentueui.OpStart && patch.Seq != watch.revision+1 { + continue + } + watch.revision = patch.Seq + event := Event{MessageID: id, Message: &message, Data: value.delivery.Data} + if err := deliver(event); err != nil { + return err + } + } + } +} + +func visibleMessage(m model.Message) contract.Message { + return contract.Message{Status: contract.MessageStatus(m.Status), TargetKind: m.TargetKind, TargetKey: m.TargetKey, ID: m.ID, ConversationID: m.ConversationID, TaskID: m.TaskID, Kind: m.Kind, Key: m.ActorKey, Content: m.Content, ReplyToID: m.ReplyToID, Purpose: m.Purpose, Revision: m.Revision, Timestamped: contract.Timestamped{CreatedAt: m.CreatedAt, UpdatedAt: m.UpdatedAt}} +} diff --git a/server/internal/component/conv_listener_test.go b/server/internal/component/conv_listener_test.go new file mode 100644 index 0000000..97839b6 --- /dev/null +++ b/server/internal/component/conv_listener_test.go @@ -0,0 +1,59 @@ +package component + +import ( + "context" + "errors" + "testing" + "time" + + runner "github.com/compforge/agentue/sdks/go/runner" + "github.com/compforge/loopd/server/internal/model" +) + +type activeRepository struct{ ConvMessageRepository } + +func (activeRepository) LatestMessageID(context.Context, string) (string, error) { return "a", nil } +func (activeRepository) ListStreamMessages(context.Context, string, string, string, int) ([]model.Message, error) { + return []model.Message{{ID: "a", ConversationID: "conv", Purpose: "output", Status: "streaming", Revision: 1, + Content: []byte(`{"version":"1.1","biz":"chat","meta":{},"blocks":[]}`)}}, nil +} + +type blockedBridge struct { + runner.EventBridge + started, stopped chan struct{} +} + +func (b *blockedBridge) State(ctx context.Context, _ string) (runner.State, error) { + close(b.started) + <-ctx.Done() + close(b.stopped) + return runner.State{}, ctx.Err() +} + +// +case=`A blocked Redis reader does not hold up SQL bootstrap/ping and is cancelled when the page leaves.` +func TestConvListenerCancelsItsReaders(t *testing.T) { + bridge := &blockedBridge{started: make(chan struct{}), stopped: make(chan struct{})} + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + disconnected := errors.New("page disconnected") + err := NewConvListener(bridge, activeRepository{}, "conv").Run(ctx, func(event Event) error { + if event.MessageID != "" { + return nil + } + select { + case <-bridge.started: + return disconnected + case <-ctx.Done(): + t.Fatal("Redis setup blocked bootstrap") + return ctx.Err() + } + }) + if !errors.Is(err, disconnected) { + t.Fatalf("listener = %v", err) + } + select { + case <-bridge.stopped: + case <-ctx.Done(): + t.Fatal("reader survived page disconnect") + } +} diff --git a/server/internal/component/message_gc.go b/server/internal/component/message_gc.go new file mode 100644 index 0000000..f9a8610 --- /dev/null +++ b/server/internal/component/message_gc.go @@ -0,0 +1,52 @@ +package component + +import ( + "context" + "log/slog" + "time" +) + +type MessageRepository interface { + ExpireMessages(context.Context, time.Time, int) ([]string, error) +} + +// MessageGC runs independently of page listeners. It closes inactive output, +// retaining conversation history; it never owns Harness execution. +type MessageGC struct { + repo MessageRepository + ttl, interval time.Duration + batchSize int + logger *slog.Logger +} + +func NewMessageGC(repo MessageRepository, ttl, interval time.Duration, batchSize int, logger *slog.Logger) *MessageGC { + if logger == nil { + logger = slog.Default() + } + return &MessageGC{repo: repo, ttl: ttl, interval: interval, batchSize: batchSize, logger: logger} +} + +func (g *MessageGC) Run(ctx context.Context) { + ticker := time.NewTicker(g.interval) + defer ticker.Stop() + for ctx.Err() == nil { + if err := g.Sweep(ctx); err != nil && ctx.Err() == nil { + g.logger.ErrorContext(ctx, "message gc failed", "error", err) + } + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + } +} + +// Sweep performs one bounded pass. Conditional updates make concurrent instances +// safe without a global lease or a second lifecycle table. +func (g *MessageGC) Sweep(ctx context.Context) error { + ids, err := g.repo.ExpireMessages(ctx, time.Now().UTC().Add(-g.ttl), g.batchSize) + if len(ids) > 0 { + g.logger.InfoContext(ctx, "message outputs expired", "count", len(ids), "message_ids", ids) + } + return err +} diff --git a/server/internal/component/message_gc_test.go b/server/internal/component/message_gc_test.go new file mode 100644 index 0000000..ee33b63 --- /dev/null +++ b/server/internal/component/message_gc_test.go @@ -0,0 +1,46 @@ +package component + +import ( + "context" + "errors" + "testing" + "time" +) + +type expiryCall struct { + cutoff time.Time + limit int +} +type expiryRepository struct{ calls chan expiryCall } + +func (r *expiryRepository) ExpireMessages(_ context.Context, cutoff time.Time, limit int) ([]string, error) { + r.calls <- expiryCall{cutoff, limit} + return nil, errors.New("temporary database failure") +} + +// +case=`GC runs with no page listener, retries failed sweeps and stops with the server.` +func TestMessageGCRunsIndependently(t *testing.T) { + repository := &expiryRepository{calls: make(chan expiryCall, 10)} + gc := NewMessageGC(repository, 24*time.Hour, 10*time.Millisecond, 100, nil) + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + done := make(chan struct{}) + go func() { defer close(done); gc.Run(ctx) }() + for range 2 { + select { + case call := <-repository.calls: + age := time.Since(call.cutoff) + if call.limit != 100 || age < 24*time.Hour || age > 24*time.Hour+time.Second { + t.Fatalf("sweep = %+v", call) + } + case <-time.After(time.Second): + t.Fatal("GC did not retry independently") + } + } + cancel() + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("GC did not stop") + } +} diff --git a/server/internal/delivery/conversation_test.go b/server/internal/delivery/conversation_test.go new file mode 100644 index 0000000..f4edf7a --- /dev/null +++ b/server/internal/delivery/conversation_test.go @@ -0,0 +1,130 @@ +package delivery + +import ( + "context" + "errors" + ui "github.com/compforge/agentue/sdks/go/ui" + "github.com/compforge/loopd/pkg/contract" + "github.com/compforge/loopd/server/internal/component" + "github.com/compforge/loopd/server/internal/model" + "github.com/compforge/loopd/server/internal/repo" + "testing" + "time" +) + +type activeQueryCounter struct { + *repo.Store + t *testing.T + activeID string + checks, reads, firstReads int +} + +func (counter *activeQueryCounter) GetMessage(ctx context.Context, id string) (model.Message, error) { + counter.reads++ + return counter.Store.GetMessage(ctx, id) +} + +func (counter *activeQueryCounter) GetMessageStates(ctx context.Context, convID string, ids []string) ([]repo.MessageState, error) { + counter.checks++ + if len(ids) != 1 || ids[0] != counter.activeID { + counter.t.Fatalf("watch includes ended history: %v", ids) + } + if counter.checks == 1 { + counter.firstReads = counter.reads + } + if counter.checks == 2 { + if counter.reads != counter.firstReads { + counter.t.Fatal("unchanged message body reloaded") + } + return nil, errStop + } + return counter.Store.GetMessageStates(ctx, convID, ids) +} + +func TestConversationStreamOnlyChecksActiveMetadata(t *testing.T) { + store, producer, _ := outputFixture(t) + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + active, err := store.Speak(ctx, "root", outputRequest("active")) + if err != nil { + t.Fatal(err) + } + old := outputRequest("ended") + old.Stream = false + if _, err := store.Speak(ctx, "root", old); err != nil { + t.Fatal(err) + } + counter := &activeQueryCounter{Store: store, t: t, activeID: active.ID} + err = component.NewConvListener(producer.events, counter, "root").Run(ctx, func(Event) error { return nil }) + if !errors.Is(err, errStop) || counter.checks != 2 { + t.Fatalf("checks=%d err=%v", counter.checks, err) + } +} + +// +case=`Conv stream skips terminal history and children; expiry ends a Message, not the subscription.` +func TestConversationStreamScopeDiscoveryAndExpiry(t *testing.T) { + store, producer, consumer := outputFixture(t) + ctx, cancel := context.WithTimeout(context.Background(), 8*time.Second) + defer cancel() + old := outputRequest("old") + old.Stream = false + history, err := store.Speak(ctx, "root", old) + if err != nil { + t.Fatal(err) + } + child, err := store.Speak(ctx, "work", outputRequest("child")) + if err != nil { + t.Fatal(err) + } + active, err := store.Speak(ctx, "root", outputRequest("active")) + if err != nil { + t.Fatal(err) + } + var next model.Message + expired := false + err = listen(ctx, consumer, "root", func(event Event) error { + if event.Message == nil { + return nil + } + if event.MessageID == history.ID || event.MessageID == child.ID || event.Message.Purpose == "input" { + t.Fatal("loaded terminal history or child conv") + } + patch, err := ui.Parse(event.Data) + if err != nil { + return err + } + if event.MessageID == active.ID && patch.Op == ui.OpStart && event.Message.Status == contract.MessageStatusStreaming { + // A missing Redis key alone must not expire the SQL Message. + if err := producer.events.Delete(ctx, streamKey(active)); err != nil { + return err + } + row, err := store.GetMessage(ctx, active.ID) + if err != nil || row.Status != "streaming" { + t.Fatal("bridge deletion changed DB status") + } + _, err = store.ExpireMessages(ctx, time.Now().Add(time.Second), 100) + return err + } + if event.MessageID == active.ID && patch.Op == ui.OpEnd { + if event.Message.Status != contract.MessageStatusExpired { + t.Fatal("missing expired status") + } + expired = true + next, err = store.Speak(ctx, "root", outputRequest("after-expiry")) + return err + } + if event.MessageID == next.ID && expired { + return errStop + } + return nil + }) + if !errors.Is(err, errStop) || !expired { + t.Fatalf("stream stopped prematurely: %v", err) + } +} + +type Event = component.Event + +func listen(ctx context.Context, coordinator *Coordinator, convID string, deliver func(Event) error) error { + return component.NewConvListener(coordinator.events, coordinator.repo.(component.ConvMessageRepository), convID).Run(ctx, deliver) +} diff --git a/server/internal/delivery/delivery.go b/server/internal/delivery/delivery.go index bb850a7..549927e 100644 --- a/server/internal/delivery/delivery.go +++ b/server/internal/delivery/delivery.go @@ -1,4 +1,4 @@ -// Package delivery multiplexes message-owned AgentUE streams for a conversation. +// Package delivery projects message output and publishes it to the event bridge. package delivery import ( @@ -20,19 +20,10 @@ var ErrInvalidEvent = errors.New("invalid AgentUE event") type MessageRepository interface { ProjectOutput(context.Context, string, agentueui.Event, ...contract.MessageStatus) error - GetDeliveryInput(context.Context, string) (model.Message, error) - ListDeliveryMessages(context.Context, string) ([]model.Message, error) GetMessage(context.Context, string) (model.Message, error) GetMessageState(context.Context, string) (repo.MessageState, error) } -type Event struct { - MessageID string - Message *contract.Message - ID string - Data json.RawMessage - Persisted bool -} type Coordinator struct { events agentuerunner.EventBridge repo MessageRepository @@ -45,20 +36,6 @@ func New(events agentuerunner.EventBridge, repository MessageRepository, logger } return &Coordinator{events: events, repo: repository, logger: logger} } -func (coordinator *Coordinator) Initialize(ctx context.Context, taskID string, content json.RawMessage) error { - start, err := agentueui.Start(content, 1) - if err != nil { - return err - } - data, err := start.Marshal() - if err != nil { - return err - } - return coordinator.events.Initialize(ctx, taskID, content, data, start.Seq) -} -func (coordinator *Coordinator) Delete(ctx context.Context, taskID string) error { - return coordinator.events.Delete(ctx, taskID) -} // +spec=`Message ID 决定输出归属,block ID 与 seq 只在该 Message 内唯一;Human 状态只能经 typed action 写入` func (coordinator *Coordinator) EmitMessage(ctx context.Context, messageID string, data json.RawMessage, statuses ...contract.MessageStatus) (string, error) { @@ -161,14 +138,7 @@ func (coordinator *Coordinator) publish(ctx context.Context, message repo.Messag return "", nil } -// Only transport control owns the Chat cursor. Every actual Message has an -// independent bridge key. -func streamKey(message model.Message) string { - if message.Purpose == "transport" { - return message.TaskID - } - return "message/" + message.ID -} +func streamKey(message model.Message) string { return "message/" + message.ID } func (coordinator *Coordinator) ensureStream(ctx context.Context, message model.Message) error { key := streamKey(message) if _, err := coordinator.events.State(ctx, key); err == nil { @@ -176,12 +146,11 @@ func (coordinator *Coordinator) ensureStream(ctx context.Context, message model. } else if !errors.Is(err, agentuerunner.ErrNotFound) { return err } - if message.Purpose != "transport" { - var err error - message, err = coordinator.repo.GetMessage(ctx, message.ID) - if err != nil { - return err - } + + var err error + message, err = coordinator.repo.GetMessage(ctx, message.ID) + if err != nil { + return err } revision := message.Revision if revision == 0 { @@ -202,15 +171,6 @@ func (coordinator *Coordinator) ensureStream(ctx context.Context, message model. return err } -func (coordinator *Coordinator) input(ctx context.Context, taskID string) (model.Message, error) { - return coordinator.repo.GetDeliveryInput(ctx, taskID) -} - -func transportMessage(input model.Message) model.Message { - return model.Message{TaskID: input.TaskID, ConversationID: input.ConversationID, Purpose: "transport", Revision: 1, - Content: []byte(`{"version":"1.1","biz":"chat","meta":{},"blocks":[]}`)} -} - func visibleMessage(m model.Message) contract.Message { return contract.Message{Status: contract.MessageStatus(m.Status), TargetKind: m.TargetKind, TargetKey: m.TargetKey, ID: m.ID, ConversationID: m.ConversationID, TaskID: m.TaskID, Kind: m.Kind, Key: m.ActorKey, Content: m.Content, ReplyToID: m.ReplyToID, Purpose: m.Purpose, Revision: m.Revision, Timestamped: contract.Timestamped{CreatedAt: m.CreatedAt, UpdatedAt: m.UpdatedAt}} } diff --git a/server/internal/delivery/delivery_test.go b/server/internal/delivery/delivery_test.go index 203ed73..11f6d24 100644 --- a/server/internal/delivery/delivery_test.go +++ b/server/internal/delivery/delivery_test.go @@ -45,9 +45,6 @@ func TestCoordinatorCompletesAndStreamsAcrossInstances(t *testing.T) { 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) - } if _, err := store.CreateMessage(ctx, model.Message{Status: "streaming", ID: "message-2", ConversationID: "conversation-1", TaskID: "task-1", Kind: "operator", ActorKey: "intent", Purpose: "output", Content: initial, Revision: 1}); err != nil { t.Fatal(err) } @@ -66,9 +63,6 @@ func TestCoordinatorCompletesAndStreamsAcrossInstances(t *testing.T) { if _, err := producer.EmitMessage(ctx, "message-2", appendEvent); err != nil { t.Fatal(err) } - if _, err := producer.EmitMessage(ctx, "message-2", marshalEvent(t, agentueui.End(4))); err != nil { - t.Fatal(err) - } message, err := store.GetMessage(ctx, "message-2") if err != nil { @@ -83,41 +77,31 @@ func TestCoordinatorCompletesAndStreamsAcrossInstances(t *testing.T) { t.Fatalf("persisted snapshot = %s", message.Content) } - var delivered []Event - if err := consumer.Stream(ctx, "task-1", "conversation-1", "", func(event Event) error { - delivered = append(delivered, event) - if event.ID != "" { - return errStop + seen := false + if err := listen(ctx, consumer, "conversation-1", func(event Event) error { + if event.MessageID != "message-2" { + return nil } - return nil - }); !errors.Is(err, errStop) { - t.Fatal(err) - } - // Output event IDs belong to their Message; only the control stream - // supplies a Chat replay cursor. - var cursor string - foundOutput := false - for _, item := range delivered { - if item.MessageID == "message-2" { - foundOutput = true - if item.ID != "" { - t.Fatal("message cursor escaped into Chat transport") - } - } else if item.ID != "" && cursor == "" { - cursor = item.ID + patch, err := agentueui.Parse(event.Data) + if err != nil { + return err } - } - if !foundOutput || cursor == "" { - t.Fatalf("missing snapshot/control: %+v", delivered) - } - if err := consumer.Stream(ctx, "task-1", "conversation-1", cursor, func(event Event) error { - if event.MessageID == "message-2" { + if patch.Op == agentueui.OpStart && !seen { + seen = true + _, err = producer.EmitMessage(ctx, "message-2", marshalEvent(t, agentueui.End(4))) + return err + } + if patch.Op == agentueui.OpEnd { + if event.Message.Status != contract.MessageStatusCompleted { + t.Fatal("missing terminal status") + } return errStop } return nil }); !errors.Is(err, errStop) { - t.Fatalf("resume: %v", err) + t.Fatal(err) } + } var errStop = errors.New("page unsubscribed") @@ -144,25 +128,22 @@ func TestHumanSnapshotsAreMessageAddressedAndRecoverWithoutRedis(t *testing.T) { t.Fatal(err) } initial := []byte(`{"version":"1.0","biz":"chat","meta":{},"blocks":[]}`) - _, err = store.CreateChatInput(ctx, model.Message{ID: "input", ConversationID: "conv", TaskID: "task", Kind: "user", ActorKey: "alice", Content: initial}) + _, err = store.CreateChatInput(ctx, model.Message{ID: "00000000-0000-7000-8000-000000000001", ConversationID: "conv", TaskID: "task", Kind: "user", ActorKey: "alice", Content: initial}) if err != nil { t.Fatal(err) } - r := contract.HumanRequest{ConversationID: "conv", Actor: contract.ActorRef{Kind: contract.ActorKindOperator, Key: "router"}, Target: contract.ActorRef{Kind: contract.ActorKindUser, Key: "alice"}, ReplyToID: "input", EffectKey: "ask", Type: "ask", Title: "Question", Prompt: "Reply", Timeout: time.Minute, AllowOther: true} + r := contract.HumanRequest{ConversationID: "conv", Actor: contract.ActorRef{Kind: contract.ActorKindOperator, Key: "router"}, Target: contract.ActorRef{Kind: contract.ActorKindUser, Key: "alice"}, ReplyToID: "00000000-0000-7000-8000-000000000001", EffectKey: "ask", Type: "ask", Title: "Question", Prompt: "Reply", Timeout: time.Minute, AllowOther: true} q, err := store.CreateHuman(ctx, r) if err != nil { t.Fatal(err) } - if _, err := store.ReplyHuman(ctx, "conv", "alice", contract.HumanReply{ReplyToID: q.Message.ID, Outcome: contract.HumanSuccess, Value: "answer"}); err != nil { - t.Fatal(err) - } redisServer := miniredis.RunT(t) client := redis.NewClient(&redis.Options{Addr: redisServer.Addr()}) defer client.Close() coordinator := New(agentuerunner.NewRedisEventBridge(client, agentuerunner.BridgeOptions{ReadBlock: time.Millisecond}), store, nil) // No Redis stream exists. Observe must recover Human snapshots from Messages. seen := map[string]bool{} - err = coordinator.Stream(ctx, "task", "conv", "", func(value Event) error { + err = listen(ctx, coordinator, "conv", func(value Event) error { if value.MessageID != "" && value.Message == nil { t.Fatal("missing Message envelope") } @@ -171,7 +152,13 @@ func TestHumanSnapshotsAreMessageAddressedAndRecoverWithoutRedis(t *testing.T) { return err } if value.MessageID == q.Message.ID && event.Op == agentueui.OpStart { - seen["question"] = true + if !seen["question"] { + seen["question"] = true + _, err = store.ReplyHuman(ctx, "conv", "alice", contract.HumanReply{ReplyToID: q.Message.ID, Outcome: contract.HumanSuccess, Value: "answer"}) + if err != nil { + return err + } + } } if value.Message != nil && value.Message.Purpose == "human_reply" { seen["reply"] = true diff --git a/server/internal/delivery/details_test.go b/server/internal/delivery/details_test.go index 18707b1..58bfdc7 100644 --- a/server/internal/delivery/details_test.go +++ b/server/internal/delivery/details_test.go @@ -32,7 +32,7 @@ func outputFixture(t *testing.T) (*repo.Store, *Coordinator, *Coordinator) { t.Fatal(err) } initial := json.RawMessage(`{"version":"1.0","biz":"chat","meta":{},"blocks":[]}`) - if _, err := store.CreateChatInput(ctx, model.Message{ID: "input", ConversationID: "root", TaskID: "task", Kind: "user", ActorKey: "human", Content: initial}); err != nil { + if _, err := store.CreateChatInput(ctx, model.Message{ID: "00000000-0000-7000-8000-000000000001", ConversationID: "root", TaskID: "task", Kind: "user", ActorKey: "human", Content: initial}); err != nil { t.Fatal(err) } redisServer := miniredis.RunT(t) @@ -40,9 +40,6 @@ func outputFixture(t *testing.T) (*repo.Store, *Coordinator, *Coordinator) { t.Cleanup(func() { _ = client.Close() }) bridge := agentuerunner.NewRedisEventBridge(client, agentuerunner.BridgeOptions{ReadBlock: time.Millisecond}) producer, consumer := New(bridge, store, nil), New(bridge, store, nil) - if err := producer.Initialize(ctx, "task", initial); err != nil { - t.Fatal(err) - } return store, producer, consumer } func ptr(value string) *string { return &value } @@ -85,7 +82,7 @@ func TestOutputMessagesOwnIdentityAndReplay(t *testing.T) { t.Fatal(err) } seen := map[string]string{} - err = consumer.Stream(ctx, "task", "root", "", func(v Event) error { + err = listen(ctx, consumer, "work", func(v Event) error { e, err := agentueui.Parse(v.Data) if err != nil { return err @@ -116,19 +113,19 @@ func TestOutputMessagesOwnIdentityAndReplay(t *testing.T) { } } count := 0 - if err := consumer.Stream(ctx, "task", "root", "", func(v Event) error { + if err := listen(ctx, consumer, "work", func(v Event) error { if v.MessageID != "" { count++ } - if count == 4 { + if count == 3 { return errStop } return nil }); !errors.Is(err, errStop) { t.Fatal(err) } - if count != 4 { - t.Fatalf("input and three messages=%d", count) + if count != 3 { + t.Fatalf("three active messages=%d", count) } } func TestSpeakConcurrentIdentity(t *testing.T) { @@ -175,12 +172,12 @@ func TestStreamDiscoversOutputDuringDelivery(t *testing.T) { defer cancel() created := "" observed := false - err := consumer.Stream(ctx, "task", "root", "", func(v Event) error { + err := listen(ctx, consumer, "work", func(v Event) error { e, err := agentueui.Parse(v.Data) if err != nil { return err } - if created == "" && v.MessageID == "" && e.Op == agentueui.OpStart { + if created == "" && v.MessageID == "" && e.Op == agentueui.OpPing { m, err := store.Speak(ctx, "work", outputRequest("later")) if err != nil { return err diff --git a/server/internal/delivery/message_test.go b/server/internal/delivery/message_test.go index 31f4ac1..4761e59 100644 --- a/server/internal/delivery/message_test.go +++ b/server/internal/delivery/message_test.go @@ -80,21 +80,23 @@ func TestMessageTerminalStatuses(t *testing.T) { if err := writer.events.Delete(ctx, "message/"+message.ID); err != nil { t.Fatal(err) } - stop := errors.New("observed terminal snapshot") - observeCtx, cancel := context.WithTimeout(ctx, time.Second) - defer cancel() - err = writer.Stream(observeCtx, "task", "root", "", func(event Event) error { - if event.MessageID == message.ID { - if event.Message.Status != status { - t.Fatalf("replay status=%s", event.Message.Status) + rows, err := store.ListMessages(ctx, "root", "", 100) + if err != nil { + t.Fatal(err) + } + found := false + for _, row := range rows { + if row.ID == message.ID { + found = true + if row.Status != string(status) { + t.Fatal("lost terminal status") } - return stop } - return nil - }) - if !errors.Is(err, stop) { - t.Fatalf("restore terminal snapshot: %v", err) } + if !found { + t.Fatal("terminal history missing") + } + }) } } @@ -221,20 +223,11 @@ func TestEndRetriesProjectionAndSurvivesBridgeLoss(t *testing.T) { if _, err := consumer.EmitMessage(ctx, message.ID, outputText(t, 4, "too late")); !errors.Is(err, ErrInvalidEvent) { t.Fatalf("write after end=%v", err) } - watchCtx, cancel := context.WithTimeout(ctx, time.Second) - defer cancel() - err = consumer.Stream(watchCtx, "task", "root", "", func(event Event) error { - if event.MessageID == message.ID { - if !event.Message.Ended() || event.Message.Revision != 3 { - t.Fatalf("lost end on replay: %+v", event.Message) - } - return errStop - } - return nil - }) - if !errors.Is(err, errStop) { - t.Fatal(err) + rows, err := store.ListMessages(ctx, "work", "", 100) + if err != nil || len(rows) != 1 || rows[0].Revision != 3 || !visibleMessage(rows[0]).Ended() { + t.Fatalf("history lost End: %+v %v", rows, err) } + } // +case=`A message End never ends the page subscription; other actors and other/no TaskIDs remain visible.` @@ -249,7 +242,7 @@ func TestSubscriptionContinuesAfterMessageEnd(t *testing.T) { seen := map[string]bool{} sent, ended := false, false var later, unsolicited string - err = consumer.Stream(ctx, "task", "root", "", func(event Event) error { + err = listen(ctx, consumer, "work", func(event Event) error { patch, err := ui.Parse(event.Data) if err != nil { return err @@ -277,7 +270,7 @@ func TestSubscriptionContinuesAfterMessageEnd(t *testing.T) { request = outputRequest("unsolicited") request.Stream = false request.Actor = contract.ActorRef{Kind: contract.ActorKindOperator, Key: "another"} - message, err = store.Speak(ctx, "root", request) + message, err = store.Speak(ctx, "work", request) if err != nil { return err } diff --git a/server/internal/delivery/stream.go b/server/internal/delivery/stream.go deleted file mode 100644 index 63515bd..0000000 --- a/server/internal/delivery/stream.go +++ /dev/null @@ -1,157 +0,0 @@ -package delivery - -import ( - "context" - "errors" - "time" - - agentuerunner "github.com/compforge/agentue/sdks/go/runner" - agentueui "github.com/compforge/agentue/sdks/go/ui" - "github.com/compforge/loopd/server/internal/model" -) - -type messageDelivery struct { - message model.Message - delivery agentuerunner.Delivery - done bool - err error -} - -// Stream multiplexes messages without folding their models together. -// The UI transport advances Last-Event-ID; each message replays its independent -// stream on reconnect. A message's end never terminates the Chat transport. -func (coordinator *Coordinator) Stream(ctx context.Context, taskID, conversationID, after string, deliver func(Event) error) error { - input, err := coordinator.input(ctx, taskID) - if err != nil { - return err - } - if input.ConversationID != conversationID { - return agentuerunner.ErrNotFound - } - ctx, cancel := context.WithCancel(ctx) - defer cancel() - incoming := make(chan messageDelivery) - started := map[string]bool{} - revisions := map[string]uint64{} - start := func(message model.Message, cursor string) error { - if started[message.ID] { - return nil - } - if err := coordinator.ensureStream(ctx, message); err != nil { - return err - } - started[message.ID] = true - go func() { - send := func(value messageDelivery) error { - select { - case incoming <- value: - return nil - case <-ctx.Done(): - return ctx.Err() - } - } - err := (agentuerunner.Replayer{Bridge: coordinator.events}).Stream(ctx, streamKey(message), cursor, func(value agentuerunner.Delivery) error { - return send(messageDelivery{message: message, delivery: value}) - }) - _ = send(messageDelivery{message: message, done: true, err: err}) - }() - return nil - } - snapshot := func(message model.Message) error { - revision := message.Revision - if revision == 0 { - revision = 1 - } - if revisions[message.ID] >= revision { - return nil - } - event, err := agentueui.Start(message.Content, revision) - if err != nil { - return err - } - data, err := event.Marshal() - if err != nil { - return err - } - msg := visibleMessage(message) - if err := deliver(Event{MessageID: message.ID, Message: &msg, Data: data}); err != nil { - return err - } - revisions[message.ID] = revision - return nil - } - discover := func() error { - rows, err := coordinator.repo.ListDeliveryMessages(ctx, conversationID) - if err != nil { - return err - } - for _, row := range rows { - // SQL snapshots repair missed bridge deliveries without asking writers - // or Operators to know whether a browser was connected. - if err := snapshot(row); err != nil { - return err - } - if row.Purpose == "output" && !visibleMessage(row).Ended() { - if err := start(row, ""); err != nil { - return err - } - } - } - return nil - } - if err := discover(); err != nil { - return err - } - control := transportMessage(input) - // A lost bridge can only restart from the durable message snapshot. - if _, err := coordinator.events.State(ctx, streamKey(control)); errors.Is(err, agentuerunner.ErrNotFound) { - after = "" - } else if err != nil { - return err - } - if err := start(control, after); err != nil { - return err - } - ticker := time.NewTicker(250 * time.Millisecond) - defer ticker.Stop() - for { - select { - case <-ctx.Done(): - return ctx.Err() - case <-ticker.C: - if err := discover(); err != nil { - return err - } - case value := <-incoming: - if value.done { - if value.err != nil { - return value.err - } - continue - } - msg := visibleMessage(value.message) - patch, err := agentueui.Parse(value.delivery.Data) - if err != nil { - return err - } - if msg.ID != "" && patch.Op == agentueui.OpEnd { - // AgentUE End carries no business payload. Read the SQL terminal - // status rather than infer successful output from transport closure. - state, err := coordinator.repo.GetMessageState(ctx, msg.ID) - if err != nil { - return err - } - msg.Status = state.Status - } - event := Event{MessageID: msg.ID, Message: &msg, Data: value.delivery.Data} - if msg.ID == control.ID { - event.Message = nil - event.ID = value.delivery.Cursor - event.Persisted = event.ID != "" - } - if err := deliver(event); err != nil { - return err - } - } - } -} diff --git a/server/internal/model/message.go b/server/internal/model/message.go index 3494cd3..0cc7827 100644 --- a/server/internal/model/message.go +++ b/server/internal/model/message.go @@ -7,7 +7,7 @@ import ( ) type Message struct { - Status string `gorm:"size:24;not null;default:completed"` + Status string `gorm:"size:24;not null;default:completed;index:idx_message_expiry,priority:1"` // Empty recipient kind and key explicitly address the conversation. TargetKind contract.ActorKind `gorm:"size:128"` TargetKey string `gorm:"size:128"` @@ -26,7 +26,7 @@ type Message struct { ActorKey string `gorm:"size:128;not null"` Content []byte `gorm:"type:json;not null"` CreatedAt time.Time - UpdatedAt time.Time + UpdatedAt time.Time `gorm:"index:idx_message_expiry,priority:2"` } func (Message) TableName() string { return "messages" } diff --git a/server/internal/repo/message.go b/server/internal/repo/message.go index 3445d13..e01a2f2 100644 --- a/server/internal/repo/message.go +++ b/server/internal/repo/message.go @@ -11,6 +11,8 @@ import ( ) type MessageRepository interface { + ExpireMessages(context.Context, time.Time, int) ([]string, error) + GetMessageStates(context.Context, string, []string) ([]MessageState, error) GetMessages(context.Context, string, []string) ([]model.Message, error) ListHumanReplies(context.Context, string, []string) ([]model.Message, error) Speak(context.Context, string, contract.SpeakRequest) (model.Message, error) @@ -66,6 +68,7 @@ type MessageState struct { Revision uint64 Status contract.MessageStatus Ended bool + HumanDueAt *time.Time } // GetMessageState reads progress without loading message bodies. @@ -120,19 +123,6 @@ func (store *Store) DeleteMessage(ctx context.Context, id string) error { })) } -// ObserveMessageActivity only widens the interval. Conditional updates remain -// safe when different servers deliver accepted events to SQL out of order. -func (store *Store) ObserveMessageActivity(ctx context.Context, id string, at time.Time) error { - ctx, cancel := store.withTimeout(ctx) - defer cancel() - if err := store.db.WithContext(ctx).Model(&model.Message{}). - Where("id = ? AND created_at > ?", id, at).UpdateColumn("created_at", at).Error; err != nil { - return mapError(err) - } - return mapError(store.db.WithContext(ctx).Model(&model.Message{}). - Where("id = ? AND updated_at < ?", id, at).UpdateColumn("updated_at", at).Error) -} - func (store *Store) CreateChatInput(ctx context.Context, input model.Message) (model.Message, error) { ctx, cancel := store.withTimeout(ctx) defer cancel() @@ -153,15 +143,6 @@ func (store *Store) CreateChatInput(ctx context.Context, input model.Message) (m return input, mapError(err) } -// ListDeliveryMessages observes a user conversation and its direct actor workspaces. -func (store *Store) ListDeliveryMessages(ctx context.Context, conversationID string) ([]model.Message, error) { - ctx, cancel := store.withTimeout(ctx) - defer cancel() - return store.readMessages(ctx, func(tx *gorm.DB) *gorm.DB { - return tx.Joins("JOIN conversations ON conversations.id = messages.conversation_id").Where("messages.conversation_id = ? OR conversations.parent_id = ?", conversationID, conversationID).Order("messages.id ASC") - }) -} - // GetMessages resolves IDs only within the already selected conversation. func (store *Store) GetMessages(ctx context.Context, conversationID string, ids []string) ([]model.Message, error) { if len(ids) == 0 { @@ -185,3 +166,77 @@ func (store *Store) ListHumanReplies(ctx context.Context, conversationID string, return tx.Where("conversation_id = ? AND reply_to_id IN ? AND purpose = ?", conversationID, questionIDs, "human_reply").Order("id ASC") }) } + +// GetMessageStates loads only metadata for a bounded set already being watched. +func (store *Store) GetMessageStates(ctx context.Context, conversationID string, ids []string) ([]MessageState, error) { + if len(ids) == 0 { + return nil, nil + } + ctx, cancel := store.withTimeout(ctx) + defer cancel() + var rows []model.Message + q := store.db.WithContext(ctx).Select("id", "conversation_id", "purpose", "revision", "status", "human_due_at").Where("conversation_id = ? AND id IN ?", conversationID, ids) + if err := q.Find(&rows).Error; err != nil { + return nil, mapError(err) + } + states := make([]MessageState, 0, len(rows)) + for _, row := range rows { + status := contract.MessageStatus(row.Status) + states = append(states, MessageState{ID: row.ID, ConversationID: row.ConversationID, Purpose: row.Purpose, Revision: row.Revision, Status: status, Ended: status.Terminal(), HumanDueAt: row.HumanDueAt}) + } + return states, nil +} + +// ExpireMessages closes inactive output without reading or rewriting its body. +// +spec=`Only streaming messages with updated_at <= now-TTL expire. A concurrent accepted write or End wins over stale cleanup candidates.` +func (store *Store) ExpireMessages(ctx context.Context, cutoff time.Time, limit int) ([]string, error) { + ctx, cancel := store.withTimeout(ctx) + defer cancel() + var rows []model.Message + err := store.db.WithContext(ctx).Select("id", "revision", "target_kind"). + Where("status = ? AND updated_at <= ?", contract.MessageStatusStreaming, cutoff). + Order("updated_at ASC, id ASC").Limit(limit).Find(&rows).Error + if err != nil { + return nil, mapError(err) + } + var expired []string + for _, row := range rows { + // Leave updated_at at the last actual output time: expiry is not activity. + result := store.db.WithContext(ctx).Model(&model.Message{}). + Where("id = ? AND status = ? AND revision = ? AND updated_at <= ?", row.ID, contract.MessageStatusStreaming, row.Revision, cutoff). + UpdateColumns(map[string]any{"status": contract.MessageStatusExpired, "revision": row.Revision + 1, "dispatch_pending": row.TargetKind != contract.ActorKindUser}) + if result.Error != nil { + return expired, mapError(result.Error) + } + if result.RowsAffected == 1 { + expired = append(expired, row.ID) + } + } + return expired, nil +} + +// LatestMessageID reads the discovery boundary without loading message bodies. +func (store *Store) LatestMessageID(ctx context.Context, convID string) (string, error) { + ctx, cancel := store.withTimeout(ctx) + defer cancel() + var ids []string + err := store.db.WithContext(ctx).Model(&model.Message{}).Where("conversation_id = ?", convID).Order("id DESC").Limit(1).Pluck("id", &ids).Error + if err != nil || len(ids) == 0 { + return "", err + } + return ids[0], nil +} + +// ListStreamMessages never walks ended history or implicitly includes child convs. +// UUIDv7 follows the existing message allocation-order assumption, not a DB commit log. +func (store *Store) ListStreamMessages(ctx context.Context, convID, after, watermark string, limit int) ([]model.Message, error) { + ctx, cancel := store.withTimeout(ctx) + defer cancel() + return store.readMessages(ctx, func(tx *gorm.DB) *gorm.DB { + q := tx.Where("conversation_id = ? AND id > ?", convID, after) + if after < watermark { + q = q.Where("id > ? OR status = ? OR human_due_at IS NOT NULL", watermark, contract.MessageStatusStreaming) + } + return q.Order("id ASC").Limit(limit) + }) +} diff --git a/server/internal/repo/message_lifecycle_test.go b/server/internal/repo/message_lifecycle_test.go new file mode 100644 index 0000000..11ab2bb --- /dev/null +++ b/server/internal/repo/message_lifecycle_test.go @@ -0,0 +1,99 @@ +package repo + +import ( + "bytes" + "context" + ui "github.com/compforge/agentue/sdks/go/ui" + "github.com/compforge/loopd/pkg/contract" + "github.com/compforge/loopd/server/internal/model" + "gorm.io/gorm" + "testing" + "time" +) + +func TestMessageExpiryPreservesBodyAndRejectsLateOutput(t *testing.T) { + s := partsStore(t) + ctx := context.Background() + m := speech(t, s, 4) + old := time.Now().UTC().Add(-2 * time.Hour) + if err := s.db.Model(&model.Message{}).Where("id = ?", m.ID).UpdateColumn("updated_at", old).Error; err != nil { + t.Fatal(err) + } + before, err := s.GetMessage(ctx, m.ID) + if err != nil { + t.Fatal(err) + } + ids, err := s.ExpireMessages(ctx, time.Now().Add(-time.Hour), 100) + if err != nil || len(ids) != 1 || ids[0] != m.ID { + t.Fatalf("expired=%v err=%v", ids, err) + } + after, err := s.GetMessage(ctx, m.ID) + if err != nil { + t.Fatal(err) + } + if after.Status != "expired" || after.Revision != before.Revision+1 || !after.UpdatedAt.Equal(before.UpdatedAt) || !bytes.Equal(after.Content, before.Content) { + t.Fatalf("expiry changed content/activity or lost status: before=%+v after=%+v", before, after) + } + if err := s.ProjectOutput(ctx, m.ID, ui.End(after.Revision+1)); err == nil { + t.Fatal("late output resurrected expired message") + } + ids, err = s.ExpireMessages(ctx, time.Now(), 100) + if err != nil || len(ids) != 0 { + t.Fatalf("terminal expired again: %v %v", ids, err) + } +} + +func TestExpiryDoesNotOverrideConcurrentWrite(t *testing.T) { + s := partsStore(t) + m := speech(t, s, 1) + old := time.Now().Add(-time.Hour) + if err := s.db.Model(&model.Message{}).Where("id = ?", m.ID).UpdateColumn("updated_at", old).Error; err != nil { + t.Fatal(err) + } + // Inject a committed writer after candidate selection, before expiry's CAS. + callback := "test:refresh_expiry_candidate" + refreshed := false + if err := s.db.Callback().Update().Before("gorm:update").Register(callback, func(tx *gorm.DB) { + if refreshed { + return + } + refreshed = true + if err := tx.Session(&gorm.Session{NewDB: true, SkipDefaultTransaction: true}).Model(&model.Message{}).Where("id = ?", m.ID). + UpdateColumns(map[string]any{"revision": m.Revision + 1, "updated_at": time.Now().UTC()}).Error; err != nil { + t.Error(err) + } + }); err != nil { + t.Fatal(err) + } + defer s.db.Callback().Update().Remove(callback) + ids, err := s.ExpireMessages(context.Background(), time.Now().Add(-time.Minute), 100) + if err != nil || len(ids) != 0 || !refreshed { + t.Fatalf("stale cleanup won: %v %v", ids, err) + } + state, err := s.GetMessageState(context.Background(), m.ID) + if err != nil || state.Status != contract.MessageStatusStreaming { + t.Fatalf("state=%+v err=%v", state, err) + } +} + +func TestAcceptedOutputRefreshesTTLNotEventTimestamp(t *testing.T) { + s := partsStore(t) + m := speech(t, s, 1) + future := time.Now().Add(24 * time.Hour).UnixMilli() + event := ui.Event{Op: ui.OpSet, Seq: 2, Timestamp: &future, Block: map[string]any{"id": "b0", "type": "text", "content": "new"}} + before := time.Now().UTC() + if err := s.ProjectOutput(context.Background(), m.ID, event); err != nil { + t.Fatal(err) + } + row, err := s.GetMessage(context.Background(), m.ID) + if err != nil || row.UpdatedAt.Before(before) || row.UpdatedAt.After(time.Now()) { + t.Fatalf("updated_at follows caller clock: %+v %v", row, err) + } + if err := s.ProjectOutput(context.Background(), m.ID, event); err != nil { + t.Fatal(err) + } + retried, err := s.GetMessage(context.Background(), m.ID) + if err != nil || !retried.UpdatedAt.Equal(row.UpdatedAt) { + t.Fatal("read/retry refreshed TTL") + } +} diff --git a/server/internal/repo/message_part_test.go b/server/internal/repo/message_part_test.go index 3c21898..5be2561 100644 --- a/server/internal/repo/message_part_test.go +++ b/server/internal/repo/message_part_test.go @@ -125,7 +125,7 @@ func TestMessagePartsMixedStorageAndReadPaths(t *testing.T) { reads := []func() ([]model.Message, error){ func() ([]model.Message, error) { return s.ListMessages(ctx, "conv", "", 100) }, func() ([]model.Message, error) { return s.ListInbox(ctx, "conv", "operator", "reader", "", 100) }, - func() ([]model.Message, error) { return s.ListDeliveryMessages(ctx, "conv") }, + func() ([]model.Message, error) { return s.ListStreamMessages(ctx, "conv", "", "", 100) }, func() ([]model.Message, error) { return s.PendingDispatches(ctx, 100) }, } for _, read := range reads { diff --git a/server/internal/repo/projection.go b/server/internal/repo/projection.go index 960c0c0..7cbc750 100644 --- a/server/internal/repo/projection.go +++ b/server/internal/repo/projection.go @@ -97,16 +97,16 @@ func (s *Store) ProjectOutput(ctx context.Context, id string, event agentueui.Ev if event.Op == agentueui.OpEnd && m.TargetKind != contract.ActorKindUser { updates["dispatch_pending"] = true } - at := time.Now().UTC() + // TTL measures server acceptance, not an untrusted or replayed event clock. + now := time.Now().UTC() + at := now if event.Timestamp != nil { at = time.UnixMilli(*event.Timestamp).UTC() } if at.Before(m.CreatedAt) { updates["created_at"] = at } - if at.After(m.UpdatedAt) { - updates["updated_at"] = at - } + updates["updated_at"] = now return tx.Model(&m).UpdateColumns(updates).Error }) } diff --git a/server/internal/service/chat.go b/server/internal/service/chat.go index d6f473b..f807760 100644 --- a/server/internal/service/chat.go +++ b/server/internal/service/chat.go @@ -22,11 +22,10 @@ type ChatRepository interface { type ChatDelivery interface { EmitMessage(context.Context, string, json.RawMessage, ...contract.MessageStatus) (string, error) - Stream(context.Context, string, string, string, func(delivery.Event) error) error } -// ChatService owns one UI chat delivery, not an Operator's business task. -// An input starts the delivery; answers are created only when actors publish. +// ChatService accepts user input independently of Conv page listeners. +// Answers are created only when actors publish. type ChatService struct { notifier MessageNotifier repo ChatRepository @@ -75,31 +74,6 @@ func (service *ChatService) Create( return messageFromModel(message), nil } -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 mapDeliveryError(err error) error { switch { case err == nil: diff --git a/server/internal/service/message.go b/server/internal/service/message.go index d26162c..244df97 100644 --- a/server/internal/service/message.go +++ b/server/internal/service/message.go @@ -142,3 +142,33 @@ func validateContent(content json.RawMessage) error { } return nil } + +// MessageChanges checks small metadata first; unchanged bodies and Parts stay in DB. +func (service *MessageService) MessageChanges(ctx context.Context, convID string, revisions map[string]uint64) ([]contract.Message, error) { + if len(revisions) > maxPageSize { + return nil, ErrInvalid + } + ids := make([]string, 0, len(revisions)) + for id := range revisions { + ids = append(ids, id) + } + states, err := service.repo.GetMessageStates(ctx, convID, ids) + if err != nil { + return nil, err + } + ids = ids[:0] + for _, state := range states { + if state.Revision > revisions[state.ID] { + ids = append(ids, state.ID) + } + } + rows, err := service.repo.GetMessages(ctx, convID, ids) + if err != nil { + return nil, err + } + result := make([]contract.Message, 0, len(rows)) + for _, row := range rows { + result = append(result, messageFromModel(row)) + } + return result, nil +} diff --git a/server/internal/service/message_test.go b/server/internal/service/message_test.go index 72cc075..6616a63 100644 --- a/server/internal/service/message_test.go +++ b/server/internal/service/message_test.go @@ -9,7 +9,6 @@ import ( "testing" "github.com/compforge/loopd/pkg/contract" - "github.com/compforge/loopd/server/internal/delivery" "github.com/compforge/loopd/server/internal/model" "github.com/compforge/loopd/server/internal/repo" ) @@ -133,25 +132,6 @@ type nopChatRunner struct{} type recordingChatRunner struct{} -func (*recordingChatRunner) Stream( - context.Context, - string, - string, - string, - func(delivery.Event) error, -) error { - return nil -} -func (nopChatRunner) Stream( - context.Context, - string, - string, - string, - func(delivery.Event) error, -) error { - return nil -} - type failingCommitRepository struct{} func (failingCommitRepository) CreateChatInput( diff --git a/server/internal/view/chat.go b/server/internal/view/chat.go index b8fb948..d453e98 100644 --- a/server/internal/view/chat.go +++ b/server/internal/view/chat.go @@ -7,7 +7,6 @@ import ( ) type CreateChatMessagesRequest struct { - TaskID string `json:"task_id,omitempty"` UserKey string `json:"user_key,omitempty"` Target contract.ActorRef `json:"target,omitempty"` Content json.RawMessage `json:"content,omitempty"` diff --git a/server/redis.go b/server/redis.go index 60c9695..8ca0cb1 100644 --- a/server/redis.go +++ b/server/redis.go @@ -12,7 +12,6 @@ type RedisConfig struct { 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 3b36364..faf66bb 100644 --- a/server/server.go +++ b/server/server.go @@ -13,6 +13,7 @@ import ( "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/component" "github.com/compforge/loopd/server/internal/delivery" "github.com/compforge/loopd/server/internal/repo" "github.com/compforge/loopd/server/internal/service" @@ -23,6 +24,7 @@ import ( type HumanIdentity func(context.Context, *hertzapp.RequestContext) (string, error) type Config struct { + MessageTTL time.Duration Conversations ConversationCoordinator Database DatabaseConfig Redis RedisConfig @@ -44,15 +46,22 @@ type DatabaseConfig struct { } type Server struct { - poll *service.PollService - store *repo.Store - redis redis.UniversalClient - api *serverapi.Server - human *service.HumanService - chat *service.ChatService + messageGC *component.MessageGC + poll *service.PollService + store *repo.Store + redis redis.UniversalClient + api *serverapi.Server + human *service.HumanService + chat *service.ChatService } func New(config Config) (*Server, error) { + if config.MessageTTL < 0 { + return nil, errors.New("message TTL must be positive") + } + if config.MessageTTL == 0 { + config.MessageTTL = DefaultMessageTTL + } if config.Conversations == nil { return nil, errors.New("conversation coordinator is required") } @@ -68,7 +77,7 @@ func New(config Config) (*Server, error) { if err != nil { return nil, err } - events, redisClient, err := newEventBridge(config.Redis) + events, redisClient, err := newEventBridge(config.Redis, config.MessageTTL) if err != nil { _ = store.Close() return nil, err @@ -81,12 +90,16 @@ func New(config Config) (*Server, error) { chat := service.NewChatService(store, chatDelivery, config.Logger, poll) human := service.NewHumanService(store, config.Logger) api := serverapi.New(actors, conversations, messages, chat, config.Logger) + api.Listen = func(ctx context.Context, convID string, deliver func(component.Event) error) error { + return component.NewConvListener(events, store, convID).Run(ctx, deliver) + } api.Human = human api.Poll = poll api.HumanIdentity = serverapi.HumanIdentity(config.HumanIdentity) return &Server{ - poll: poll, - human: human, chat: chat, + messageGC: component.NewMessageGC(store, config.MessageTTL, time.Second, 100, config.Logger), + poll: poll, + human: human, chat: chat, store: store, redis: redisClient, api: api, @@ -98,13 +111,14 @@ func (server *Server) Run(ctx context.Context) { var workers sync.WaitGroup workers.Go(func() { server.human.Run(ctx) }) workers.Go(func() { server.poll.Run(ctx) }) + workers.Go(func() { server.messageGC.Run(ctx) }) workers.Wait() } func (server *Server) Close() error { return errors.Join(server.redis.Close(), server.store.Close()) } -func newEventBridge(config RedisConfig) (agentuerunner.EventBridge, redis.UniversalClient, error) { +func newEventBridge(config RedisConfig, ttl time.Duration) (agentuerunner.EventBridge, redis.UniversalClient, error) { if config.Address == "" { config.Address = "127.0.0.1:6379" } @@ -123,9 +137,6 @@ func newEventBridge(config RedisConfig) (agentuerunner.EventBridge, redis.Univer 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 } @@ -147,7 +158,7 @@ func newEventBridge(config RedisConfig) (agentuerunner.EventBridge, redis.Univer 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, + KeyPrefix: config.KeyPrefix, TaskTTL: ttl, ReadBlock: config.ReadBlock, ReadCount: config.ReadCount, }), client, nil } diff --git a/web/src/App.tsx b/web/src/App.tsx index 2bb0a52..781a94c 100644 --- a/web/src/App.tsx +++ b/web/src/App.tsx @@ -1,38 +1,24 @@ import { ActorKind, isOperatorKind, operatorRole } from "./actor"; import type { MessageContent } from "./content"; import { useEffect, useRef, useState, type FormEvent } from "react"; -import { - PatchOp, -} from "@compforge/agentue/ui"; import { createConversation, listActors, listConversations, - listMessages, streamMessage, type Actor, - type ActorRef, type Conversation, type Message, } from "./api"; import { MessageBody, ReplyReference } from "./MessageBody"; +import { MessagePoller } from "./message-poll"; import { mergeMessage, applyMessageEvent, messageStatusLabel } from "./message"; import { DetailPanel, detailOrganizer, type DetailSelection } from "./DetailPanel"; -import { readSubscriptions, writeSubscription, type StoredSubscription } from "./streams"; +import { useConversationStream } from "./streams"; const selectedActorKey = "loopd.selected-actor"; const selectedConversationKey = "loopd.selected-conversation"; -type ConnectionStatus = "connecting" | "connected" | "reconnecting" | "failed"; - -interface LiveSubscription { - messages?: Message[]; - conversationID: string; - taskID: string; - lastEventID: string; - status: ConnectionStatus; -} - export function App() { const [conversations, setConversations] = useState([]); const [actors, setActors] = useState([]); @@ -41,17 +27,17 @@ export function App() { const [messages, setMessages] = useState([]); const [selectedMessageID, setSelectedMessageID] = useState(); const [detailSelection, setDetailSelection] = useState(); - const [liveSubscriptions, setLiveSubscriptions] = useState>({}); const [submitting, setSubmitting] = useState(false); const [draft, setDraft] = useState(""); const [loading, setLoading] = useState(true); const [error, setError] = useState(); - const streams = useRef(new Map()); - useEffect(() => () => { for (const controller of streams.current.values()) controller.abort(); }, []); + const messagePoller = useRef<{ conversationID: string; poller: MessagePoller } | undefined>(undefined); + useConversationStream(selectedConversationID, (delivery) => { + if (delivery.messageID) setMessages((current) => applyMessageEvent(current, delivery)); + }, (signal) => refreshMessages(selectedConversationID!, signal, true)); const selectedConversation = conversations.find((item) => item.id === selectedConversationID); const selectedActor = actors.find((actor) => actorIdentity(actor) === selectedActorID); - const subscription = selectedConversationID ? liveSubscriptions[selectedConversationID] : undefined; useEffect(() => { const controller = new AbortController(); @@ -111,26 +97,18 @@ export function App() { if (lastMessage) setDetailSelection((current) => current ?? { parentID: selectedConversationID, organizer: detailOrganizer(lastMessage), }); - const stored = readSubscriptions()[selectedConversationID]; - const input = items.findLast((message) => message.purpose === "input" && message.task_id); - if (!streams.current.has(selectedConversationID) && (stored || input)) { - void observeConversation(selectedConversationID, stored ?? { taskID: input!.task_id, lastEventID: "" }); - } }); return () => controller.abort(); }, [selectedConversationID]); - useEffect(() => { - if (!selectedConversationID) return; - const controller = new AbortController(); - // Actors may publish without an active user Chat; discover their snapshots too. - const timer = window.setInterval(() => { void refreshMessages(selectedConversationID, controller.signal); }, 2000); - return () => { window.clearInterval(timer); controller.abort(); }; - }, [selectedConversationID]); - async function refreshMessages(conversationID: string, signal?: AbortSignal): Promise { + async function refreshMessages(conversationID: string, signal?: AbortSignal, sync = false): Promise { try { - const items = await listMessages(conversationID, signal); + if (messagePoller.current?.conversationID !== conversationID) { + messagePoller.current = { conversationID, poller: new MessagePoller(conversationID) }; + } + const poller = messagePoller.current.poller; + const items = await (sync ? poller.sync(signal) : poller.poll(signal)); if (signal?.aborted) return []; setMessages((current) => { // Equal message revisions may carry refreshed reference previews/cards. @@ -147,10 +125,8 @@ export function App() { function selectConversation(conversationID: string) { if (conversationID === selectedConversationID) return; - for (const controller of streams.current.values()) controller.abort(); - streams.current.clear(); + messagePoller.current = undefined; setMessages([]); - setLiveSubscriptions({}); setSelectedMessageID(undefined); setDetailSelection(undefined); setSelectedConversationID(conversationID); @@ -158,11 +134,9 @@ export function App() { } function startConversation() { - for (const controller of streams.current.values()) controller.abort(); - streams.current.clear(); + messagePoller.current = undefined; setSelectedConversationID(undefined); setMessages([]); - setLiveSubscriptions({}); setSelectedMessageID(undefined); setDetailSelection(undefined); setError(undefined); @@ -214,81 +188,19 @@ export function App() { updated_at: new Date().toISOString(), }, ]); - await observeConversation(conversationID, undefined, text, selectedActor); - } - - async function observeConversation( - conversationID: string, - stored?: StoredSubscription, - text?: string, - requestedTarget?: ActorRef, - ) { - const slot = conversationID; - if (stored && streams.current.has(slot)) return; - streams.current.get(slot)?.abort(); - const controller = new AbortController(); - streams.current.set(slot, controller); - let taskID = stored?.taskID ?? ""; - let lastEventID = stored?.lastEventID ?? ""; - let awaitingID = text !== undefined; - let liveMessages: Message[] = []; - const target = requestedTarget; - const update = (status: ConnectionStatus) => { - if (controller.signal.aborted) return; - setLiveSubscriptions((current) => ({ ...current, [slot]: { - conversationID, taskID, lastEventID, status, - messages: [...liveMessages], - } })); - }; - update("connecting"); try { - for (;;) { - try { - await streamMessage({ - conversationID, taskID: taskID || undefined, lastEventID: lastEventID || undefined, - text, target, signal: controller.signal, - onTaskID: (value) => { - if (controller.signal.aborted) return; - const first = !taskID; - taskID = value; - if (awaitingID) { setSubmitting(false); awaitingID = false; } - if (first) void refreshMessages(conversationID, controller.signal); - writeSubscription(conversationID, { taskID, lastEventID }); - update("connected"); - }, - onEvent: (delivery) => { - if (controller.signal.aborted) return; - const { event: patch, eventId, messageID, message } = delivery; - if (messageID && message) { - liveMessages = applyMessageEvent(liveMessages, delivery); - if (message.conversation_id === conversationID) { - const updated = liveMessages.find((item) => item.id === messageID)!; - setMessages((current) => mergeMessage(current, updated)); - } - // A Message's END closes only that Message. - update("connected"); - return; - } - // Control events carry transport lifecycle only, never a bubble. - if (eventId) lastEventID = eventId; - if (patch.op === PatchOp.ERROR) setError(String(patch.meta.error.message ?? "连接错误")); - writeSubscription(conversationID, { taskID, lastEventID }); - update("connected"); - }, - }); - if (!taskID) throw new Error("chat stream closed before an ID was returned"); - } catch (cause) { - if (isAbort(cause)) return; - if (!taskID) throw cause; - } - update("reconnecting"); - await delay(1_500, controller.signal); - } + await streamMessage({ + conversationID, text, target: selectedActor, + onTaskID: () => setSubmitting(false), + onEvent: (delivery) => { + if (!delivery.messageID) return; + setMessages((current) => applyMessageEvent(current.filter((m) => !m.id.startsWith("local-")), delivery)); + }, + }); } catch (cause) { - if (!isAbort(cause)) { setError(errorMessage(cause)); update("failed"); } + if (!isAbort(cause)) setError(errorMessage(cause)); } finally { - if (streams.current.get(slot) === controller) streams.current.delete(slot); - if (awaitingID) setSubmitting(false); + setSubmitting(false); } } @@ -438,13 +350,13 @@ export function App() { setMessages((current) => { - let next = mergeMessage(current, result.message); - if (result.reply) next = mergeMessage(next, result.reply); + let next = current; + for (const message of [result.message, result.reply]) { + if (message && message.conversation_id === selectedConversationID) next = mergeMessage(next, message); + } return next; })} selection={detailSelection?.parentID === selectedConversationID ? detailSelection : undefined} - liveMessages={subscription?.messages} - running={subscription?.status === "connected"} /> ); @@ -494,12 +406,3 @@ function errorMessage(cause: unknown): string { function isAbort(cause: unknown): boolean { return cause instanceof DOMException && cause.name === "AbortError"; } - -function delay(milliseconds: number, signal: AbortSignal): Promise { - return new Promise((resolve, reject) => { - if (signal.aborted) { reject(new DOMException("aborted", "AbortError")); return; } - const abort = () => { window.clearTimeout(timer); reject(new DOMException("aborted", "AbortError")); }; - const timer = window.setTimeout(() => { signal.removeEventListener("abort", abort); resolve(); }, milliseconds); - signal.addEventListener("abort", abort, { once: true }); - }); -} diff --git a/web/src/DetailPanel.tsx b/web/src/DetailPanel.tsx index 10cb7b3..13526b3 100644 --- a/web/src/DetailPanel.tsx +++ b/web/src/DetailPanel.tsx @@ -1,10 +1,13 @@ import { ActorKind, operatorOwner } from "./actor"; import { MessageBody, ReplyReference } from "./MessageBody"; import type { HumanResult } from "./api"; -import { messageStatusLabel } from "./message"; -import { useEffect, useState, type CSSProperties } from "react"; +import { messageStatusLabel, mergeMessage } from "./message"; +import { useConversationStream } from "./streams"; +import { applyMessageEvent } from "./message"; +import { MessagePoller } from "./message-poll"; +import { useEffect, useRef, useState, type CSSProperties } from "react"; import { parseMessageContent, type MessageContent } from "./content"; -import { findDetailConversation, listMessages, type Conversation, type Message } from "./api"; +import { findDetailConversation, type Conversation, type Message } from "./api"; import { traceColor, traceLabel } from "./trace"; import { groupParallelMessages } from "./parallel"; @@ -21,12 +24,11 @@ export interface DetailSelection { } /** @spec 按父会话/Operator 观察工作会话,不等待主回答;切换参与者不能泄漏上一个查询的结果。 */ -export function DetailPanel({ selection, liveMessages, running, onReply }: { +export function DetailPanel({ selection, onReply }: { selection?: DetailSelection; - liveMessages?: Message[]; - running: boolean; onReply?(result: HumanResult): void; }) { + const history = useRef<{ scope: string; poller: MessagePoller } | undefined>(undefined); const [detail, setDetail] = useState(); const parentID = selection?.parentID; const organizer = selection?.organizer; @@ -37,35 +39,59 @@ export function DetailPanel({ selection, liveMessages, running, onReply }: { if (!parentID || !actorKind || !actorKey) return; const controller = new AbortController(); let timer: ReturnType; + let conversation: Conversation | undefined; + let poller: MessagePoller | undefined; + let messages: Message[] = []; + let loaded = false; // The actor's workspace stays observable beyond any single UI delivery. async function refresh() { try { - const conversation = await findDetailConversation(parentID!, actorKind!, actorKey!, controller.signal); - const messages = conversation ? await listMessages(conversation.id, controller.signal) : []; - if (!controller.signal.aborted) setDetail({ scope, conversation, messages }); + conversation ??= await findDetailConversation(parentID!, actorKind!, actorKey!, controller.signal); + if (conversation) { + poller ??= new MessagePoller(conversation.id); + history.current = { scope, poller }; + for (const message of await poller.poll(controller.signal)) messages = mergeMessage(messages, message); + loaded = true; + } + if (!controller.signal.aborted) setDetail((current) => { + let merged = current?.scope === scope ? current.messages : []; + for (const message of messages) merged = mergeMessage(merged, message); + return { scope, conversation, messages: merged }; + }); } catch (cause) { if (!controller.signal.aborted) { setDetail({ scope, messages: [], error: String(cause) }); } } finally { - if (!controller.signal.aborted) timer = setTimeout(refresh, running ? 1_000 : 2_000); + if (!controller.signal.aborted && !loaded) timer = setTimeout(refresh, 2_000); } } void refresh(); return () => { controller.abort(); clearTimeout(timer); }; - }, [scope, parentID, actorKind, actorKey, running]); + }, [scope, parentID, actorKind, actorKey]); const selected = detail?.scope === scope ? detail : undefined; - const visible = [...(selected?.messages ?? [])]; - for (const item of liveMessages ?? []) { - if (item.conversation_id !== selected?.conversation?.id) continue; - const index = visible.findIndex((value) => value.id === item.id); - if (index < 0) visible.push(item); - else if ((item.revision ?? 0) > (visible[index].revision ?? 0)) { - // The stream carries content updates; polling refreshes the activity interval. - visible[index] = { ...item, created_at: visible[index].created_at, updated_at: visible[index].updated_at }; + useConversationStream(selected?.conversation?.id, (event) => { + if (!event.messageID) return; + setDetail((current) => current?.scope === scope + ? { ...current, messages: applyMessageEvent(current.messages, event) } : current); + }, async (signal) => { + const source = history.current; + if (source?.scope !== scope) return; + try { + const updates = await source.poller.sync(signal); + if (signal.aborted) return; + setDetail((current) => { + if (current?.scope !== scope) return current; + let messages = current.messages; + for (const message of updates) messages = mergeMessage(messages, message); + return { ...current, messages }; + }); + } catch (cause) { + if (!signal.aborted) setDetail((current) => current?.scope === scope ? { ...current, error: String(cause) } : current); } - } + }); + const visible = selected?.messages ?? []; const groups = groupParallelMessages(visible); const indices = new Map(visible.map((item, index) => [item.id, index])); return ( @@ -92,7 +118,15 @@ export function DetailPanel({ selection, liveMessages, running, onReply }: {
{group.columns.map((column) => (
- {column.map((item) => )} + {column.map((item) => { + setDetail((current) => { + if (current?.scope !== scope) return current; + let messages = mergeMessage(current.messages, result.message); + if (result.reply) messages = mergeMessage(messages, result.reply); + return { ...current, messages }; + }); + onReply?.(result); + }} />)}
))}
diff --git a/web/src/api.ts b/web/src/api.ts index 64cd8e5..e06ed35 100644 --- a/web/src/api.ts +++ b/web/src/api.ts @@ -19,7 +19,7 @@ export interface Conversation { export interface Message { card?: MessageCard; reply_to?: { id: string; kind: ActorKind; key: string; preview: string }; - status: "streaming" | "completed" | "failed" | "cancelled"; + status: "streaming" | "completed" | "failed" | "cancelled" | "expired"; id: string; target_kind?: ActorKind; target_key?: string; @@ -69,9 +69,8 @@ export async function createConversation(name: string, signal?: AbortSignal): Pr }); } -export async function listMessages(conversationID: string, signal?: AbortSignal): Promise { +export async function listMessages(conversationID: string, signal?: AbortSignal, after = ""): Promise { const messages: Message[] = []; - let after = ""; for (;;) { const page = await requestJSON>( `/v1/conversations/${encodeURIComponent(conversationID)}/messages?limit=100&after=${encodeURIComponent(after)}`, @@ -83,10 +82,21 @@ export async function listMessages(conversationID: string, signal?: AbortSignal) } } +export async function messageChanges(conversationID: string, revisions: Map, signal?: AbortSignal): Promise { + const entries = [...revisions]; + const messages: Message[] = []; + for (let i = 0; i < entries.length; i += 100) { + const watch = entries.slice(i, i + 100).map(([id, revision]) => `${id}:${revision}`).join(","); + const page = await requestJSON>( + `/v1/conversations/${encodeURIComponent(conversationID)}/messages?watch=${encodeURIComponent(watch)}`, { signal }, + ); + messages.push(...page.data); + } + return messages; +} + export interface StreamRequest { conversationID: string; - taskID?: string; - lastEventID?: string; text?: string; target?: ActorRef; signal?: AbortSignal; @@ -95,19 +105,15 @@ export interface StreamRequest { } export async function streamMessage(request: StreamRequest): Promise { - const replay = Boolean(request.taskID); - const body = replay - ? { task_id: request.taskID } - : { - user_key: "web-user", - target: request.target, - content: textModel(request.text ?? ""), - }; + const body = { + user_key: "web-user", + target: request.target, + content: textModel(request.text ?? ""), + }; const headers: Record = { Accept: "text/event-stream", "Content-Type": "application/json", }; - if (request.lastEventID) headers["Last-Event-ID"] = request.lastEventID; const response = await fetch( `/v1/conversations/${encodeURIComponent(request.conversationID)}/messages`, { method: "POST", headers, body: JSON.stringify(body), signal: request.signal }, @@ -118,6 +124,19 @@ export async function streamMessage(request: StreamRequest): Promise { if (!taskID) throw new Error("loop-server response omitted task ID"); request.onTaskID(taskID); + await readMessageStream(response, request.onEvent); +} + +export async function streamConversation(conversationID: string, signal: AbortSignal, onEvent: (event: MessageEvent) => void): Promise { + const response = await fetch(`/v1/conversations/${encodeURIComponent(conversationID)}/stream`, { + headers: { Accept: "text/event-stream" }, signal, + }); + if (!response.ok) throw await responseError(response); + await readMessageStream(response, onEvent); +} + +async function readMessageStream(response: Response, onEvent: (event: MessageEvent) => void) { + if (!response.body) throw new Error("loop-server returned an empty event stream"); const decoder = new TextDecoder(); const frames = new SseFrameDecoder(); const reader = response.body.getReader(); @@ -125,12 +144,12 @@ export async function streamMessage(request: StreamRequest): Promise { const chunk = await reader.read(); if (chunk.done) break; for (const frame of frames.push(decoder.decode(chunk.value, { stream: true }))) { - request.onEvent(decodeMessageFrame(frame)); + onEvent(decodeMessageFrame(frame)); } } - for (const frame of frames.push(decoder.decode())) request.onEvent(decodeMessageFrame(frame)); + for (const frame of frames.push(decoder.decode())) onEvent(decodeMessageFrame(frame)); const tail = frames.finish(); - if (tail) request.onEvent(decodeMessageFrame(tail)); + if (tail) onEvent(decodeMessageFrame(tail)); } export class SseFrameDecoder { diff --git a/web/src/detail.test.tsx b/web/src/detail.test.tsx index c4ffb0c..e011b01 100644 --- a/web/src/detail.test.tsx +++ b/web/src/detail.test.tsx @@ -31,15 +31,15 @@ it("does not invent an Operator for broadcasts from users or direct Harness mess }); it("shows the selected Operator before a Message or workspace exists", () => { - const html = renderToStaticMarkup(); + const html = renderToStaticMarkup(); expect(html).toContain("处理详情 · router"); expect(html).toContain("正在查找 router 的工作会话"); expect(html).not.toContain("选择一条消息"); }); it("distinguishes no selection from a message with no related Operator", () => { - expect(renderToStaticMarkup()).toContain("发送消息或选择历史消息"); - expect(renderToStaticMarkup()) + expect(renderToStaticMarkup()).toContain("发送消息或选择历史消息"); + expect(renderToStaticMarkup()) .toContain("这条消息未关联 Operator 工作会话"); }); diff --git a/web/src/message-poll.test.ts b/web/src/message-poll.test.ts new file mode 100644 index 0000000..8366173 --- /dev/null +++ b/web/src/message-poll.test.ts @@ -0,0 +1,39 @@ +import { afterEach, expect, it, vi } from "vitest"; +import { MessagePoller } from "./message-poll"; + +afterEach(() => vi.unstubAllGlobals()); + +it("reconciles after an in-flight history read instead of sharing a stale snapshot", async () => { + let finish!: (response: Response) => void; + const fetch = vi.fn() + .mockImplementationOnce(() => new Promise((resolve) => { finish = resolve; })) + .mockResolvedValueOnce(Response.json({ data: [{ id: "b", status: "completed", revision: 1 }] })) + .mockResolvedValueOnce(Response.json({ data: [{ id: "a", status: "expired", revision: 2 }] })); + vi.stubGlobal("fetch", fetch); + const poller = new MessagePoller("conv"); + const history = poller.poll(); + const repair = poller.sync(); + finish(Response.json({ data: [{ id: "a", status: "streaming", revision: 1 }] })); + expect(await history).toHaveLength(1); + expect((await repair).map((message) => [message.id, message.status])).toEqual([["b", "completed"], ["a", "expired"]]); + expect(fetch).toHaveBeenCalledTimes(3); +}); + +it("discovers by ID and stops watching terminal messages", async () => { + const fetch = vi.fn() + .mockResolvedValueOnce(Response.json({ data: [{ id: "a", status: "completed", revision: 1 }, { id: "b", status: "streaming", revision: 1 }] })) + .mockResolvedValueOnce(Response.json({ data: [] })) + .mockResolvedValueOnce(Response.json({ data: [{ id: "b", status: "expired", revision: 2 }] })) + .mockResolvedValueOnce(Response.json({ data: [] })); + vi.stubGlobal("fetch", fetch); + const poller = new MessagePoller("conv"); + expect(await poller.poll()).toHaveLength(2); + expect((await poller.poll())[0].status).toBe("expired"); + expect(await poller.poll()).toEqual([]); + expect(fetch.mock.calls.map(([url]) => url)).toEqual([ + "/v1/conversations/conv/messages?limit=100&after=", + "/v1/conversations/conv/messages?limit=100&after=b", + "/v1/conversations/conv/messages?watch=b%3A1", + "/v1/conversations/conv/messages?limit=100&after=b", + ]); +}); diff --git a/web/src/message-poll.ts b/web/src/message-poll.ts new file mode 100644 index 0000000..fee08c6 --- /dev/null +++ b/web/src/message-poll.ts @@ -0,0 +1,48 @@ +import { listMessages, messageChanges, type Message } from "./api"; +import { parseMessageContent } from "./content"; + +// One poller belongs to one selected conversation. Only discovery advances after; +// receiving an SSE frame must not skip earlier, not-yet-discovered messages. +export class MessagePoller { + private after = ""; + private active = new Map(); + private pending?: Promise; + + constructor(private readonly conversationID: string) {} + + poll(signal?: AbortSignal): Promise { + // Initial load, reconnect and submit callbacks may coincide. Share the request + // instead of racing cursor updates or letting a slow network build a queue. + if (!this.pending) this.pending = this.read(signal).finally(() => { this.pending = undefined; }); + return this.pending; + } + + async sync(signal?: AbortSignal): Promise { + // Do not reuse a history request started before the stream's watermark. + if (this.pending) await this.pending; + signal?.throwIfAborted(); + return this.poll(signal); + } + + private async read(signal?: AbortSignal): Promise { + const added = await listMessages(this.conversationID, signal, this.after); + const changed = await messageChanges(this.conversationID, this.active, signal); + signal?.throwIfAborted(); + if (added.length) this.after = added[added.length - 1].id; + const updates = [...added, ...changed]; + for (const message of updates) { + if (isActive(message)) this.active.set(message.id, message.revision ?? 0); + else this.active.delete(message.id); + } + return updates; + } +} + +function isActive(message: Message): boolean { + if (message.status === "streaming") return true; + if (message.purpose !== "human_request") return false; + if (message.card?.type === "ask" || message.card?.type === "confirm") return message.card.question.status === "pending"; + // The content model remains valid without a server-enriched card. + return parseMessageContent(message.content).blocks.some((block) => + (block.type === "ask" || block.type === "confirm") && block.status === "pending"); +} diff --git a/web/src/message.ts b/web/src/message.ts index bd2547f..5da5525 100644 --- a/web/src/message.ts +++ b/web/src/message.ts @@ -3,7 +3,7 @@ import { applyPatch } from "@compforge/agentue/ui"; import type { Message, MessageEvent } from "./api"; export function messageStatusLabel(status: Message["status"]): string { - return { streaming: "生成中", completed: "已发送", failed: "输出失败", cancelled: "已取消" }[status]; + return { streaming: "生成中", completed: "已发送", failed: "输出失败", cancelled: "已取消", expired: "已过期" }[status]; } // +spec=`Message ID owns the snapshot; equal block IDs in parallel questions never collide` diff --git a/web/src/streams.test.ts b/web/src/streams.test.ts index ccdd2e8..5533432 100644 --- a/web/src/streams.test.ts +++ b/web/src/streams.test.ts @@ -1,21 +1,21 @@ -import { afterEach, beforeEach, expect, it, vi } from "vitest"; -import { readSubscriptions, writeSubscription } from "./streams"; +import { afterEach, expect, it, vi } from "vitest"; +import { streamConversation } from "./api"; -beforeEach(() => { - const values = new Map(); - vi.stubGlobal("localStorage", { - getItem: (key: string) => values.get(key) ?? null, - setItem: (key: string, value: string) => values.set(key, value), - }); -}); afterEach(() => vi.unstubAllGlobals()); -it("replaces the connection identity without accumulating per-input streams", () => { - writeSubscription("conv", { taskID: "first", lastEventID: "1-0" }); - writeSubscription("other", { taskID: "other", lastEventID: "" }); - writeSubscription("conv", { taskID: "followup", lastEventID: "2-0" }); - expect(readSubscriptions()).toEqual({ - conv: { taskID: "followup", lastEventID: "2-0" }, - other: { taskID: "other", lastEventID: "" }, - }); +it("subscribes by Conv without task identity and keeps reading after message End", async () => { + const fetch = vi.fn().mockResolvedValue(new Response( + 'data: {"message_id":"a","event":{"op":"end","seq":2}}\n\n' + + 'data: {"message_id":"b","event":{"op":"set","seq":3,"block":{"id":"text","type":"text","content":"later"}}}\n\n', + { headers: { "Content-Type": "text/event-stream" } }, + )); + vi.stubGlobal("fetch", fetch); + const events: string[] = []; + const controller = new AbortController(); + await streamConversation("conv/one", controller.signal, (event) => events.push(event.messageID!)); + expect(events).toEqual(["a", "b"]); + expect(fetch.mock.calls[0]).toEqual([ + "/v1/conversations/conv%2Fone/stream", + { headers: { Accept: "text/event-stream" }, signal: controller.signal }, + ]); }); diff --git a/web/src/streams.ts b/web/src/streams.ts index ff28f99..7d85f02 100644 --- a/web/src/streams.ts +++ b/web/src/streams.ts @@ -1,17 +1,36 @@ -const subscriptionKey = "loopd.subscriptions"; +import { useEffect, useRef } from "react"; +import { streamConversation, type MessageEvent } from "./api"; -export interface StoredSubscription { - taskID: string; - lastEventID: string; -} - -// One page subscription observes the conversation, including later actor output. -// taskID is a reconnect identity, not a business execution lifetime. -export function readSubscriptions(): Record { - try { return JSON.parse(localStorage.getItem(subscriptionKey) ?? "{}"); } - catch { return {}; } -} - -export function writeSubscription(conversationID: string, value: StoredSubscription) { - localStorage.setItem(subscriptionKey, JSON.stringify({ ...readSubscriptions(), [conversationID]: value })); +// A Conv connection multiplexes Message streams. Reconnect restores active SQL +// snapshots; one message's End never ends this page subscription. +export function useConversationStream(conversationID: string | undefined, onEvent: (event: MessageEvent) => void, onReady?: (signal: AbortSignal) => Promise) { + const callback = useRef(onEvent); + callback.current = onEvent; + const ready = useRef(onReady); + ready.current = onReady; + useEffect(() => { + if (!conversationID) return; + const controller = new AbortController(); + let timer: ReturnType | undefined; + async function connect() { + let connected = false; + try { + await streamConversation(conversationID!, controller.signal, (event) => { + if (controller.signal.aborted) return; + if (!connected) { + connected = true; + // Repair the initial-load gap and messages completed while disconnected. + void ready.current?.(controller.signal); + } + callback.current(event); + }); + } catch { + // Reconnect reconciles SQL once; live discovery belongs to the server listener. + } finally { + if (!controller.signal.aborted) timer = setTimeout(connect, 1500); + } + } + void connect(); + return () => { controller.abort(); clearTimeout(timer); }; + }, [conversationID]); }