Skip to content

Commit 1c9e530

Browse files
author
SqlRush
committed
Accept remote event wrappers
1 parent 82de698 commit 1c9e530

4 files changed

Lines changed: 62 additions & 9 deletions

File tree

docs/cc-100-roadmap.md

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -121,6 +121,8 @@ M10 补充:`internal/remote` 新增 callback 型 `StreamWebSocketEvents` primi
121121

122122
M10 补充:`--daemon` 常驻模式现在会在初始 tick 后启动 remote WebSocket stream goroutine,按 heartbeat 间隔重试注册状态,复用 `StreamWebSocketEvents``RemoteTrigger` delivery/dedupe,把推送事件实时注入 running team,并在 `remote-pump.json` 中持续更新 `websocket_stream` transport、frame/connect/reconnect、delivered/duplicate/error 计数;daemon heartbeat/tick 在已注册 `websocket_url` 时会跳过短 WebSocket 读取,避免 stream 和 tick 双连接/重复写 pump state,poll-only 注册仍走原 tick 路径;daemon stop/context cancel 会取消 stream,并在 pump state 与 `/status show remote` 中记录 stream start/end/stop reason。完整云端协议 hardening 仍未完成。
123123

124+
M10 补充:remote poll/WebSocket 共用解码器现在兼容 `data``event``remote_event``delivery``payload` 包裹的单条事件,以及这些 wrapper 下的 `events/items/messages/deliveries` 列表;云端可以用 envelope 协议携带 cursor 和事件内容,而无需强制把事件字段铺在顶层。更深的鉴权刷新、ack/lease 和服务端协议协商仍未完成。
125+
124126
M7 补充:interaction script paste payload 现在接受 ClipboardItem 风格的 `items[].getAsString`/`get_as_string` 以及 `stringData`/`textData` 文本字段,DOM clipboard 录制脚本可直接恢复 pasted text。
125127

126128
M7 补充:scripted task runtime payload 和 task expectation 现在接受 `taskID``jobId``runId``label``displayName``phase``taskState``message``currentStep``percent`/`percentage`/`pct` 等相邻字段,并支持数字 task ID 与数字字符串 progress。

docs/claude-code-go-rewrite-plan.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -749,6 +749,7 @@ test/parity/ # golden tests against TS/official behavior
749749
- 本轮补充:remote WebSocket pump 现在支持单次 tick 内读取多帧事件,并在握手/读帧失败或非正常 close 时按可配置 backoff 重连;daemon 默认读取最多 8 帧、最多重连 2 次,并把 frame/connect/reconnect 计数写入 `remote-pump.json``/status show remote`。完整 CCR 云端 WebSocket 常驻持久 stream 与更深协议 hardening 仍未完成。
750750
- 本轮补充:`internal/remote` 新增 callback 型 `StreamWebSocketEvents` primitive,可保持 WebSocket 连接逐帧解码并把事件批次交给调用方,支持 context 取消、可选帧上限、handler 错误传播、异常 close/读错后的 backoff 重连以及 `ReconnectAttempts < 0` 无限重连语义;该能力为 daemon 常驻 stream 托管接线打底。完整云端协议 hardening 仍未完成。
751751
- 本轮补充:`--daemon` 常驻模式现在会在初始 tick 后启动 remote WebSocket stream goroutine,按 heartbeat 间隔重试注册状态,复用 `StreamWebSocketEvents``RemoteTrigger` delivery/dedupe,把推送事件实时注入 running team,并在 `remote-pump.json` 中持续更新 `websocket_stream` transport、frame/connect/reconnect、delivered/duplicate/error 计数;daemon heartbeat/tick 在已注册 `websocket_url` 时会跳过短 WebSocket 读取,避免 stream 和 tick 双连接/重复写 pump state,poll-only 注册仍走原 tick 路径;daemon stop/context cancel 会取消 stream,并在 pump state 与 `/status show remote` 中记录 stream start/end/stop reason。完整云端协议 hardening 仍未完成。
752+
- 本轮补充:remote poll/WebSocket 共用解码器现在兼容 `data``event``remote_event``delivery``payload` 包裹的单条事件,以及这些 wrapper 下的 `events/items/messages/deliveries` 列表;云端可以用 envelope 协议携带 cursor 和事件内容,而无需强制把事件字段铺在顶层。更深的鉴权刷新、ack/lease 和服务端协议协商仍未完成。
752753

753754
### M11: Bridge 和高级集成
754755

internal/remote/pump.go

Lines changed: 38 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -145,16 +145,21 @@ func DecodePollEvents(data []byte) ([]PollEvent, string, error) {
145145
case map[string]any:
146146
cursor := firstString(value, "next_cursor", "nextCursor", "cursor", "after", "after_id", "afterId")
147147
if nested, ok := firstAny(value, "events", "items", "messages", "deliveries", "data"); ok {
148-
if nestedMap, ok := nested.(map[string]any); ok {
149-
if nestedEvents, ok := firstAny(nestedMap, "events", "items", "messages", "deliveries"); ok {
150-
nested = nestedEvents
151-
}
152-
if cursor == "" {
153-
cursor = firstString(nestedMap, "next_cursor", "nextCursor", "cursor", "after", "after_id", "afterId")
154-
}
148+
events, nestedCursor, decoded := decodePollEventValue(nested)
149+
if cursor == "" {
150+
cursor = nestedCursor
155151
}
156-
if list, ok := nested.([]any); ok {
157-
return decodePollEventList(list), cursor, nil
152+
if decoded && (len(events) > 0 || nestedCursor != "") {
153+
return events, cursor, nil
154+
}
155+
}
156+
if nested, ok := firstAny(value, "event", "remote_event", "remoteEvent", "delivery", "payload"); ok {
157+
events, nestedCursor, decoded := decodePollEventValue(nested)
158+
if cursor == "" {
159+
cursor = nestedCursor
160+
}
161+
if decoded && (len(events) > 0 || nestedCursor != "") {
162+
return events, cursor, nil
158163
}
159164
}
160165
if event, ok := decodePollEventMap(value); ok {
@@ -166,6 +171,30 @@ func DecodePollEvents(data []byte) ([]PollEvent, string, error) {
166171
}
167172
}
168173

174+
func decodePollEventValue(value any) ([]PollEvent, string, bool) {
175+
switch typed := value.(type) {
176+
case []any:
177+
return decodePollEventList(typed), "", true
178+
case map[string]any:
179+
cursor := firstString(typed, "next_cursor", "nextCursor", "cursor", "after", "after_id", "afterId")
180+
if nested, ok := firstAny(typed, "events", "items", "messages", "deliveries"); ok {
181+
events, nestedCursor, decoded := decodePollEventValue(nested)
182+
if cursor == "" {
183+
cursor = nestedCursor
184+
}
185+
if decoded {
186+
return events, cursor, true
187+
}
188+
}
189+
if event, ok := decodePollEventMap(typed); ok {
190+
return []PollEvent{event}, cursor, true
191+
}
192+
return nil, cursor, cursor != ""
193+
default:
194+
return nil, "", false
195+
}
196+
}
197+
169198
func WritePumpState(path string, state PumpState) error {
170199
if path == "" {
171200
return os.ErrInvalid

internal/remote/pump_test.go

Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -64,6 +64,27 @@ func TestDecodePollEventsAcceptsNestedDataAndArray(t *testing.T) {
6464
if cursor != "" || len(events) != 1 || events[0].EventID != "evt-single" || events[0].Message != "single payload" {
6565
t.Fatalf("single events=%#v cursor=%q", events, cursor)
6666
}
67+
events, cursor, err = DecodePollEvents([]byte(`{"cursor":"c4","data":{"id":"evt-data","team":"team","payload":{"message":"data payload"}}}`))
68+
if err != nil {
69+
t.Fatal(err)
70+
}
71+
if cursor != "c4" || len(events) != 1 || events[0].EventID != "evt-data" || events[0].Message != "data payload" {
72+
t.Fatalf("data wrapper events=%#v cursor=%q", events, cursor)
73+
}
74+
events, cursor, err = DecodePollEvents([]byte(`{"type":"remote_trigger","event":{"delivery_id":"evt-event","team_id":"team","target":"coordinator","message":"event wrapper"}}`))
75+
if err != nil {
76+
t.Fatal(err)
77+
}
78+
if cursor != "" || len(events) != 1 || events[0].EventID != "evt-event" || events[0].Target != "coordinator" || events[0].Message != "event wrapper" {
79+
t.Fatalf("event wrapper events=%#v cursor=%q", events, cursor)
80+
}
81+
events, cursor, err = DecodePollEvents([]byte(`{"kind":"delivery","payload":{"id":"evt-payload","team_id":"team","event_type":"deploy","message":"payload wrapper"}}`))
82+
if err != nil {
83+
t.Fatal(err)
84+
}
85+
if cursor != "" || len(events) != 1 || events[0].EventID != "evt-payload" || events[0].Event != "deploy" || events[0].Message != "payload wrapper" {
86+
t.Fatalf("payload wrapper events=%#v cursor=%q", events, cursor)
87+
}
6788
events, cursor, err = DecodePollEvents([]byte(`{"cursor":"c3","message":"ok"}`))
6889
if err != nil {
6990
t.Fatal(err)

0 commit comments

Comments
 (0)