Skip to content

Commit 25ee76b

Browse files
author
SqlRush
committed
Add daemon tick endpoint
1 parent cfc7412 commit 25ee76b

6 files changed

Lines changed: 158 additions & 17 deletions

File tree

cmd/claude/main.go

Lines changed: 61 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@ import (
1111
"path/filepath"
1212
"sort"
1313
"strings"
14+
"sync"
1415
"time"
1516

1617
"ccgo/internal/api/anthropic"
@@ -264,6 +265,18 @@ func runDaemon(ctx context.Context, state *bootstrap.State, options daemonOption
264265
}
265266
endpoint := ""
266267
var daemonServer *daemonpkg.Server
268+
var tickMu sync.Mutex
269+
writeHeartbeat := func(now time.Time) error {
270+
return daemonpkg.WriteState(statePath, daemonpkg.BuildState(runner.SessionID, runner.WorkingDirectory, daemonpkg.RuntimeRunning, os.Getpid(), endpoint, now, nil))
271+
}
272+
runTick := func(now time.Time) (contracts.ToolResult, error) {
273+
tickMu.Lock()
274+
defer tickMu.Unlock()
275+
if err := writeHeartbeat(now); err != nil {
276+
return contracts.ToolResult{}, err
277+
}
278+
return runDaemonDueSchedules(ctx, runner, now)
279+
}
267280
if !options.Once {
268281
daemonServer, err = daemonpkg.StartServer(daemonpkg.ServerOptions{
269282
StateFunc: func() daemonpkg.State {
@@ -273,6 +286,13 @@ func runDaemon(ctx context.Context, state *bootstrap.State, options daemonOption
273286
}
274287
return state
275288
},
289+
TickFunc: func(context.Context) daemonpkg.TickResponse {
290+
result, err := runTick(time.Now().UTC())
291+
if err != nil {
292+
return daemonpkg.TickResponse{OK: false, Error: err.Error()}
293+
}
294+
return daemonTickResponse(result)
295+
},
276296
})
277297
if err != nil {
278298
fmt.Fprintf(stderr, "ccgo daemon: %v\n", err)
@@ -285,16 +305,7 @@ func runDaemon(ctx context.Context, state *bootstrap.State, options daemonOption
285305
}()
286306
endpoint = daemonServer.Endpoint()
287307
}
288-
writeHeartbeat := func(now time.Time) error {
289-
return daemonpkg.WriteState(statePath, daemonpkg.BuildState(runner.SessionID, runner.WorkingDirectory, daemonpkg.RuntimeRunning, os.Getpid(), endpoint, now, nil))
290-
}
291-
runTick := func(now time.Time) error {
292-
if err := writeHeartbeat(now); err != nil {
293-
return err
294-
}
295-
return runDaemonDueSchedules(ctx, runner, now)
296-
}
297-
if err := runTick(time.Now().UTC()); err != nil {
308+
if _, err := runTick(time.Now().UTC()); err != nil {
298309
fmt.Fprintf(stderr, "ccgo daemon: %v\n", err)
299310
return 1
300311
}
@@ -307,27 +318,63 @@ func runDaemon(ctx context.Context, state *bootstrap.State, options daemonOption
307318
for {
308319
select {
309320
case <-ctx.Done():
321+
tickMu.Lock()
310322
_ = daemonpkg.WriteState(statePath, daemonpkg.BuildState(runner.SessionID, runner.WorkingDirectory, daemonpkg.RuntimeDisabled, os.Getpid(), endpoint, time.Now().UTC(), ctx.Err()))
323+
tickMu.Unlock()
311324
return 0
312325
case now := <-ticker.C:
313-
if err := runTick(now.UTC()); err != nil {
326+
if _, err := runTick(now.UTC()); err != nil {
314327
fmt.Fprintf(stderr, "ccgo daemon: %v\n", err)
315328
return 1
316329
}
317330
}
318331
}
319332
}
320333

321-
func runDaemonDueSchedules(ctx context.Context, runner conversation.Runner, now time.Time) error {
322-
_, err := tasktools.RunDueSchedules(tool.Context{
334+
func runDaemonDueSchedules(ctx context.Context, runner conversation.Runner, now time.Time) (contracts.ToolResult, error) {
335+
return tasktools.RunDueSchedules(tool.Context{
323336
Context: ctx,
324337
WorkingDirectory: runner.WorkingDirectory,
325338
SessionID: runner.SessionID,
326339
Metadata: map[string]any{
327340
tool.MetadataSessionPathKey: runner.SessionPath,
328341
},
329342
}, "", now, tool.NopProgressSink())
330-
return err
343+
}
344+
345+
func daemonTickResponse(result contracts.ToolResult) daemonpkg.TickResponse {
346+
structured := result.StructuredContent
347+
return daemonpkg.TickResponse{
348+
OK: true,
349+
CheckedAt: stringMapValue(structured, "checked_at"),
350+
TriggeredCount: intMapValue(structured, "triggered_count"),
351+
ErrorCount: intMapValue(structured, "error_count"),
352+
Structured: structured,
353+
}
354+
}
355+
356+
func stringMapValue(values map[string]any, key string) string {
357+
if values == nil {
358+
return ""
359+
}
360+
value, _ := values[key].(string)
361+
return value
362+
}
363+
364+
func intMapValue(values map[string]any, key string) int {
365+
if values == nil {
366+
return 0
367+
}
368+
switch value := values[key].(type) {
369+
case int:
370+
return value
371+
case int64:
372+
return int(value)
373+
case float64:
374+
return int(value)
375+
default:
376+
return 0
377+
}
331378
}
332379

333380
func handleChromeNativeHostMessage(raw json.RawMessage) map[string]any {

cmd/claude/main_test.go

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -192,7 +192,7 @@ func TestRunDaemonDueSchedulesNoopsWithoutSchedules(t *testing.T) {
192192
WorkingDirectory: t.TempDir(),
193193
SessionPath: filepath.Join(t.TempDir(), "session.jsonl"),
194194
}
195-
if err := runDaemonDueSchedules(context.Background(), runner, time.Now().UTC()); err != nil {
195+
if _, err := runDaemonDueSchedules(context.Background(), runner, time.Now().UTC()); err != nil {
196196
t.Fatal(err)
197197
}
198198
}
@@ -232,6 +232,18 @@ func TestRunDaemonServesHealthEndpoint(t *testing.T) {
232232
if !health.OK || health.SessionID != state.SessionID() || health.RuntimeState != daemonpkg.RuntimeRunning {
233233
t.Fatalf("health = %#v", health)
234234
}
235+
tickResp, err := http.Post(daemonState.Endpoint+"/tick", "application/json", nil)
236+
if err != nil {
237+
t.Fatal(err)
238+
}
239+
defer tickResp.Body.Close()
240+
var tick daemonpkg.TickResponse
241+
if err := json.NewDecoder(tickResp.Body).Decode(&tick); err != nil {
242+
t.Fatal(err)
243+
}
244+
if !tick.OK || tick.ErrorCount != 0 {
245+
t.Fatalf("tick = %#v", tick)
246+
}
235247
cancel()
236248
select {
237249
case code := <-done:

docs/cc-100-roadmap.md

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -89,7 +89,7 @@ M10 补充:新增 `ScheduleCron` 工具入口和 session-scoped `schedules.jso
8989

9090
M10 补充:conversation runner 在每轮主请求前会执行一次 `ScheduleCron` due tick,复用 `run_due` 的触发与 last-run 去重逻辑,把到期 schedule 自动注入 running team recipients;无到期任务时不发噪声进度,触发失败会 fail-open 并发 `schedule_due_error` progress。完整常驻后台 daemon 和远端服务接入仍未完成。
9191

92-
M10 补充:新增 session-scoped `daemon-state.json` 状态契约和 `/status show daemon` 审计 section,可记录 runtime_state、PID、endpoint、started_at、heartbeat_at 和错误,并按 heartbeat 超时把 running 状态判定为 stale;CLI 新增 `--daemon` heartbeat loop 和 `--daemon-once` 单次写入模式,`--daemon` 会启动 loopback `/health`/`/status` HTTP endpoint 并在 tick 中复用 `ScheduleCron run_due` 执行到期 schedule。完整 daemon 进程管理/远端托管仍未完成。
92+
M10 补充:新增 session-scoped `daemon-state.json` 状态契约和 `/status show daemon` 审计 section,可记录 runtime_state、PID、endpoint、started_at、heartbeat_at 和错误,并按 heartbeat 超时把 running 状态判定为 stale;CLI 新增 `--daemon` heartbeat loop 和 `--daemon-once` 单次写入模式,`--daemon` 会启动 loopback `/health`/`/status`/`/tick` HTTP endpoint,定时 tick 和 `POST /tick` 都复用 `ScheduleCron run_due` 执行到期 schedule。完整 daemon 进程管理/远端托管仍未完成。
9393

9494
M10 补充:新增 `RemoteTrigger` 工具入口,可把 source/event/message 作为远端触发事件注入到 running team recipients;默认优先发给 coordinator,消息正文保留远端来源和事件类型。`event_id` 可选,提供后会写入 session-scoped `remote_triggers.json` receipt,重复投递会 no-op 并记录 duplicate_count,避免远端重试重复注入。完整 remote websocket/CCR 服务接入仍未完成。
9595

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -734,7 +734,7 @@ test/parity/ # golden tests against TS/official behavior
734734
- 本轮补充:新增 `Brief` 工具入口,可把 summary/title/status/details/next_steps/risks 规范为 structured handoff brief,并支持常见字段别名与单字符串列表项归一化;这为后续远端协作/UI brief surface 提供稳定 payload。完整 remote brief UI/调度接线仍未完成。
735735
- 本轮补充:新增 `ScheduleCron` 工具入口和 session-scoped `schedules.json` manifest,可 create/list/delete/trigger/run_due cron schedule metadata,校验 5-field cron 或常见 `@daily` 类表达式,并可绑定 team_id/target/message;`trigger` 会把保存的 schedule message 发送给绑定 team 的 running recipients,`run_due` 会按当前分钟执行到期且启用的 schedule,并记录 last run 状态避免同一分钟重复触发。当前已有手动与一次性到期执行路径,完整后台 daemon 仍未完成。
736736
- 本轮补充:conversation runner 在每轮主请求前会执行一次 `ScheduleCron` due tick,复用 `run_due` 的触发与 last-run 去重逻辑,把到期 schedule 自动注入 running team recipients;无到期任务时不发噪声进度,触发失败会 fail-open 并发 `schedule_due_error` progress。完整常驻后台 daemon 和远端服务接入仍未完成。
737-
- 本轮补充:新增 session-scoped `daemon-state.json` 状态契约和 `/status show daemon` 审计 section,可记录 runtime_state、PID、endpoint、started_at、heartbeat_at 和错误,并按 heartbeat 超时把 running 状态判定为 stale;CLI 新增 `--daemon` heartbeat loop 和 `--daemon-once` 单次写入模式,`--daemon` 会启动 loopback `/health`/`/status` HTTP endpoint 并在 tick 中复用 `ScheduleCron run_due` 执行到期 schedule。完整 daemon 进程管理/远端托管仍未完成。
737+
- 本轮补充:新增 session-scoped `daemon-state.json` 状态契约和 `/status show daemon` 审计 section,可记录 runtime_state、PID、endpoint、started_at、heartbeat_at 和错误,并按 heartbeat 超时把 running 状态判定为 stale;CLI 新增 `--daemon` heartbeat loop 和 `--daemon-once` 单次写入模式,`--daemon` 会启动 loopback `/health`/`/status`/`/tick` HTTP endpoint,定时 tick 和 `POST /tick` 都复用 `ScheduleCron run_due` 执行到期 schedule。完整 daemon 进程管理/远端托管仍未完成。
738738
- 本轮补充:新增 `RemoteTrigger` 工具入口,可把 source/event/message 作为远端触发事件注入到 running team recipients;默认优先发给 coordinator,消息正文保留远端来源和事件类型。`event_id` 可选,提供后会写入 session-scoped `remote_triggers.json` receipt,重复投递会 no-op 并记录 duplicate_count,避免远端重试重复注入。完整 remote websocket/CCR 服务接入仍未完成。
739739
- 本轮补充:bridge direct server 新增 loopback-only `POST /remote-trigger` HTTP endpoint,并沿用 direct server token guard;conversation runner 会把该 endpoint 接到 `RemoteTrigger` 的校验、注入和 event_id dedupe 逻辑,远端系统可通过受控 HTTP 请求向 running team 注入事件。完整 remote websocket/CCR 长连接服务仍未完成。
740740
- 本轮补充:bridge direct WebSocket JSON 通道新增 `remote_trigger` action,复用 direct remote trigger request/response 结构和同一回调,可在已鉴权的 loopback WebSocket 连接上注入远端事件。完整 CCR 云端长连接协议仍未完成。

internal/daemon/server.go

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ import (
1414
type ServerOptions struct {
1515
Addr string
1616
StateFunc func() State
17+
TickFunc func(context.Context) TickResponse
1718
}
1819

1920
type Server struct {
@@ -28,6 +29,15 @@ type HealthResponse struct {
2829
PID int `json:"pid,omitempty"`
2930
}
3031

32+
type TickResponse struct {
33+
OK bool `json:"ok"`
34+
CheckedAt string `json:"checked_at,omitempty"`
35+
TriggeredCount int `json:"triggered_count,omitempty"`
36+
ErrorCount int `json:"error_count,omitempty"`
37+
Structured map[string]any `json:"structured,omitempty"`
38+
Error string `json:"error,omitempty"`
39+
}
40+
3141
func StartServer(options ServerOptions) (*Server, error) {
3242
addr := strings.TrimSpace(options.Addr)
3343
if addr == "" {
@@ -65,6 +75,22 @@ func StartServer(options ServerOptions) (*Server, error) {
6575
}
6676
writeJSON(w, http.StatusOK, stateFunc())
6777
})
78+
mux.HandleFunc("/tick", func(w http.ResponseWriter, r *http.Request) {
79+
if r.Method != http.MethodPost {
80+
writeJSON(w, http.StatusMethodNotAllowed, map[string]string{"error": "method not allowed"})
81+
return
82+
}
83+
if options.TickFunc == nil {
84+
writeJSON(w, http.StatusNotImplemented, TickResponse{OK: false, Error: "daemon tick is not configured"})
85+
return
86+
}
87+
response := options.TickFunc(r.Context())
88+
status := http.StatusOK
89+
if !response.OK {
90+
status = http.StatusInternalServerError
91+
}
92+
writeJSON(w, status, response)
93+
})
6894
server := &http.Server{Handler: mux}
6995
wrapped := &Server{listener: listener, server: server}
7096
go func() {

internal/daemon/server_test.go

Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import (
44
"context"
55
"encoding/json"
66
"net/http"
7+
"net/http/httptest"
78
"strings"
89
"testing"
910
"time"
@@ -60,3 +61,58 @@ func TestStartServerRejectsNonLoopbackAddress(t *testing.T) {
6061
t.Fatalf("err = %v", err)
6162
}
6263
}
64+
65+
func TestServerTick(t *testing.T) {
66+
var called bool
67+
server, err := StartServer(ServerOptions{
68+
TickFunc: func(context.Context) TickResponse {
69+
called = true
70+
return TickResponse{OK: true, CheckedAt: "2026-06-17T10:00:00Z", TriggeredCount: 2}
71+
},
72+
})
73+
if err != nil {
74+
t.Fatal(err)
75+
}
76+
t.Cleanup(func() {
77+
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
78+
defer cancel()
79+
if err := server.Close(ctx); err != nil {
80+
t.Fatalf("close daemon server: %v", err)
81+
}
82+
})
83+
resp, err := http.Post(server.Endpoint()+"/tick", "application/json", nil)
84+
if err != nil {
85+
t.Fatal(err)
86+
}
87+
defer resp.Body.Close()
88+
if resp.StatusCode != http.StatusOK {
89+
t.Fatalf("tick status = %d", resp.StatusCode)
90+
}
91+
var tick TickResponse
92+
if err := json.NewDecoder(resp.Body).Decode(&tick); err != nil {
93+
t.Fatal(err)
94+
}
95+
if !called || !tick.OK || tick.TriggeredCount != 2 {
96+
t.Fatalf("called=%v tick=%#v", called, tick)
97+
}
98+
}
99+
100+
func TestServerTickRequiresCallback(t *testing.T) {
101+
server, err := StartServer(ServerOptions{})
102+
if err != nil {
103+
t.Fatal(err)
104+
}
105+
t.Cleanup(func() {
106+
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
107+
defer cancel()
108+
if err := server.Close(ctx); err != nil {
109+
t.Fatalf("close daemon server: %v", err)
110+
}
111+
})
112+
recorder := httptest.NewRecorder()
113+
req := httptest.NewRequest(http.MethodPost, "/tick", nil)
114+
server.server.Handler.ServeHTTP(recorder, req)
115+
if recorder.Code != http.StatusNotImplemented {
116+
t.Fatalf("tick status = %d body=%s", recorder.Code, recorder.Body.String())
117+
}
118+
}

0 commit comments

Comments
 (0)