Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion VERSION
Original file line number Diff line number Diff line change
@@ -1 +1 @@
0.0.24
0.0.25
5 changes: 5 additions & 0 deletions cmd/loop-server/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
Expand All @@ -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},
Expand Down
8 changes: 8 additions & 0 deletions cmd/loop-server/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand All @@ -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)
Expand All @@ -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) {
Expand Down Expand Up @@ -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",
Expand Down
3 changes: 2 additions & 1 deletion cmd/loop-server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
2 changes: 2 additions & 0 deletions deploy/k8s/loopd/templates/server.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions deploy/k8s/loopd/values.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
4 changes: 2 additions & 2 deletions docs/kernel.md
Original file line number Diff line number Diff line change
Expand Up @@ -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。
Expand Down
3 changes: 2 additions & 1 deletion pkg/contract/chat.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
7 changes: 5 additions & 2 deletions server/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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/ # 消息消费、可见事实持久化与用户交互的领域设计
Expand Down
7 changes: 7 additions & 0 deletions server/const.go
Original file line number Diff line number Diff line change
@@ -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
11 changes: 6 additions & 5 deletions server/docs/persistence.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 分别表达普通输出、交互问题和卡片答复,不指定唯一主回答。

Expand Down Expand Up @@ -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#消息呈现),不由存储层规定页面布局。
53 changes: 37 additions & 16 deletions server/docs/ue.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 配对输入与回答的业务入口。

### 消息寻址与快照
Expand All @@ -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 判断业务完成;新发言使用新消息。
Loading