Skip to content

Commit 9ca0c92

Browse files
author
SqlRush
committed
Avoid duplicate websocket reads in daemon ticks
1 parent 3ae298e commit 9ca0c92

4 files changed

Lines changed: 82 additions & 4 deletions

File tree

cmd/claude/main.go

Lines changed: 43 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -66,6 +66,10 @@ type daemonOptions struct {
6666
HeartbeatInterval time.Duration
6767
}
6868

69+
type daemonTickOptions struct {
70+
SkipRemoteWhenWebSocket bool
71+
}
72+
6973
type daemonControlOptions struct {
7074
StatePath string
7175
Status bool
@@ -690,6 +694,7 @@ func runDaemon(ctx context.Context, state *bootstrap.State, options daemonOption
690694
var tickMu sync.Mutex
691695
stopRequested := make(chan struct{})
692696
var stopOnce sync.Once
697+
remoteStreamMode := false
693698
writeHeartbeat := func(now time.Time) error {
694699
return daemonpkg.WriteState(statePath, daemonpkg.BuildState(runner.SessionID, runner.WorkingDirectory, daemonpkg.RuntimeRunning, os.Getpid(), endpoint, now, nil))
695700
}
@@ -699,7 +704,7 @@ func runDaemon(ctx context.Context, state *bootstrap.State, options daemonOption
699704
if err := writeHeartbeat(now); err != nil {
700705
return contracts.ToolResult{}, err
701706
}
702-
return runDaemonTick(ctx, runner, now)
707+
return runDaemonTickWithOptions(ctx, runner, now, daemonTickOptions{SkipRemoteWhenWebSocket: remoteStreamMode})
703708
}
704709
if !options.Once {
705710
daemonServer, err = daemonpkg.StartServer(daemonpkg.ServerOptions{
@@ -737,6 +742,7 @@ func runDaemon(ctx context.Context, state *bootstrap.State, options daemonOption
737742
fmt.Fprintf(stderr, "ccgo daemon: %v\n", err)
738743
return 1
739744
}
745+
remoteStreamMode = true
740746
fmt.Fprintf(stdout, "ccgo daemon running\nsession_id=%s\nstate_path=%s\n", runner.SessionID, statePath)
741747
if options.Once {
742748
return 0
@@ -921,11 +927,20 @@ func runDaemonDueSchedules(ctx context.Context, runner conversation.Runner, now
921927
}
922928

923929
func runDaemonTick(ctx context.Context, runner conversation.Runner, now time.Time) (contracts.ToolResult, error) {
930+
return runDaemonTickWithOptions(ctx, runner, now, daemonTickOptions{})
931+
}
932+
933+
func runDaemonTickWithOptions(ctx context.Context, runner conversation.Runner, now time.Time, options daemonTickOptions) (contracts.ToolResult, error) {
924934
scheduleResult, err := runDaemonDueSchedules(ctx, runner, now)
925935
if err != nil {
926936
return contracts.ToolResult{}, err
927937
}
928-
remoteResult := runDaemonRemotePoll(ctx, runner, now)
938+
var remoteResult contracts.ToolResult
939+
if options.SkipRemoteWhenWebSocket {
940+
remoteResult = runDaemonRemotePollUnlessStream(ctx, runner, now)
941+
} else {
942+
remoteResult = runDaemonRemotePoll(ctx, runner, now)
943+
}
929944
schedule := scheduleResult.StructuredContent
930945
remotePoll := remoteResult.StructuredContent
931946
scheduleTriggered := intMapValue(schedule, "triggered_count")
@@ -953,6 +968,32 @@ func runDaemonTick(ctx context.Context, runner conversation.Runner, now time.Tim
953968
}, nil
954969
}
955970

971+
func runDaemonRemotePollUnlessStream(ctx context.Context, runner conversation.Runner, now time.Time) contracts.ToolResult {
972+
registrationPath := remotepkg.SessionRegistrationPath(runner.SessionPath, runner.SessionID)
973+
if registrationPath == "" {
974+
return runDaemonRemotePoll(ctx, runner, now)
975+
}
976+
registration, err := remotepkg.LoadRegistrationState(registrationPath)
977+
if err != nil || registration.RuntimeState != remotepkg.RegistrationRegistered || strings.TrimSpace(registration.WebSocketURL) == "" {
978+
return runDaemonRemotePoll(ctx, runner, now)
979+
}
980+
structured := map[string]any{
981+
"type": "remote_poll",
982+
"checked_at": now.UTC().Format(time.RFC3339Nano),
983+
"runtime_state": remotepkg.PumpRunning,
984+
"transport": "websocket_stream",
985+
"skipped": true,
986+
"event_count": 0,
987+
"delivered_count": 0,
988+
"duplicate_count": 0,
989+
"error_count": 0,
990+
}
991+
return contracts.ToolResult{
992+
Content: "Remote stream is handling websocket events.",
993+
StructuredContent: structured,
994+
}
995+
}
996+
956997
func runDaemonRemotePoll(ctx context.Context, runner conversation.Runner, now time.Time) contracts.ToolResult {
957998
structured := map[string]any{
958999
"type": "remote_poll",

cmd/claude/main_test.go

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -513,6 +513,43 @@ func TestRunDaemonRemoteStreamInjectsRemoteTriggers(t *testing.T) {
513513
}
514514
}
515515

516+
func TestRunDaemonTickSkipsRemotePollWhenWebSocketStreamRegistered(t *testing.T) {
517+
dir := t.TempDir()
518+
transcriptPath := filepath.Join(dir, "session.jsonl")
519+
sessionID := contracts.ID("sess_daemon_stream_skip")
520+
if err := remotepkg.WriteRegistrationState(remotepkg.SessionRegistrationPath(transcriptPath, sessionID), remotepkg.RegistrationState{
521+
SessionID: sessionID,
522+
RuntimeState: remotepkg.RegistrationRegistered,
523+
WebSocketURL: "ws://127.0.0.1:1/stream?token=secret",
524+
PollURL: "https://poll.example.invalid/events?token=secret",
525+
}); err != nil {
526+
t.Fatal(err)
527+
}
528+
runner := conversation.Runner{
529+
SessionID: sessionID,
530+
SessionPath: transcriptPath,
531+
WorkingDirectory: dir,
532+
}
533+
result, err := runDaemonTickWithOptions(context.Background(), runner, time.Unix(200, 0).UTC(), daemonTickOptions{SkipRemoteWhenWebSocket: true})
534+
if err != nil {
535+
t.Fatal(err)
536+
}
537+
remotePoll, ok := result.StructuredContent["remote_poll"].(map[string]any)
538+
if !ok {
539+
t.Fatalf("remote poll = %#v", result.StructuredContent["remote_poll"])
540+
}
541+
if remotePoll["transport"] != "websocket_stream" || remotePoll["skipped"] != true || remotePoll["error_count"] != 0 {
542+
t.Fatalf("remote poll = %#v", remotePoll)
543+
}
544+
pump, err := remotepkg.LoadPumpState(remotepkg.SessionPumpPath(transcriptPath, sessionID))
545+
if err != nil {
546+
t.Fatal(err)
547+
}
548+
if pump.RuntimeState != "" {
549+
t.Fatalf("pump should not be overwritten by skipped tick: %#v", pump)
550+
}
551+
}
552+
516553
func TestRunDaemonServesHealthEndpoint(t *testing.T) {
517554
t.Setenv("CLAUDE_CONFIG_DIR", t.TempDir())
518555
cwd := t.TempDir()

docs/cc-100-roadmap.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -119,7 +119,7 @@ M10 补充:remote WebSocket pump 现在支持单次 tick 内读取多帧事件
119119

120120
M10 补充:`internal/remote` 新增 callback 型 `StreamWebSocketEvents` primitive,可保持 WebSocket 连接逐帧解码并把事件批次交给调用方,支持 context 取消、可选帧上限、handler 错误传播、异常 close/读错后的 backoff 重连以及 `ReconnectAttempts < 0` 无限重连语义;该能力为 daemon 常驻 stream 托管接线打底。完整云端协议 hardening 仍未完成。
121121

122-
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 stop/context cancel 会取消 stream。完整云端协议 hardening 和更细的 stream lifecycle 审计仍未完成。
122+
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。完整云端协议 hardening 和更细的 stream lifecycle 审计仍未完成。
123123

124124
M7 补充:interaction script paste payload 现在接受 ClipboardItem 风格的 `items[].getAsString`/`get_as_string` 以及 `stringData`/`textData` 文本字段,DOM clipboard 录制脚本可直接恢复 pasted text。
125125

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -748,7 +748,7 @@ test/parity/ # golden tests against TS/official behavior
748748
- 本轮补充:新增 session-scoped `remote-pump.json` 和 daemon remote 消息泵;daemon tick 会在 ScheduleCron due tick 后优先读取 registered `websocket_url`,通过 Bearer auth 建立 WebSocket、读取单帧事件并复用 poll 解码/`RemoteTrigger` 注入路径;无 WebSocket 或 WebSocket 失败且存在 poll URL 时回退到带 cursor 的 poll 拉取。pump 状态会记录 transport、websocket/poll URL 脱敏值、cursor、HTTP status、event/delivered/duplicate/error 计数,`/status show remote` 可审计当前传输。完整 CCR WebSocket 常驻持久 stream 和云端协议 hardening 仍未完成。
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 仍未完成。
751-
- 本轮补充:`--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 stop/context cancel 会取消 stream。完整云端协议 hardening 和更细的 stream lifecycle 审计仍未完成。
751+
- 本轮补充:`--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。完整云端协议 hardening 和更细的 stream lifecycle 审计仍未完成。
752752

753753
### M11: Bridge 和高级集成
754754

0 commit comments

Comments
 (0)