Skip to content

Commit f3f61bf

Browse files
authored
Merge pull request #154 from shelltime/feat/daemon-coding-heartbeat
feat(daemon): add coding heartbeat tracking with offline persistence
2 parents 7aa6099 + b86fa51 commit f3f61bf

10 files changed

Lines changed: 448 additions & 11 deletions

File tree

cmd/daemon/main.go

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,17 @@ func main() {
105105
}
106106
}
107107

108+
// Start heartbeat resync service if codeTracking is enabled
109+
if cfg.CodeTracking != nil && cfg.CodeTracking.Enabled != nil && *cfg.CodeTracking.Enabled {
110+
heartbeatResyncService := daemon.NewHeartbeatResyncService(cfg)
111+
if err := heartbeatResyncService.Start(ctx); err != nil {
112+
slog.Error("Failed to start heartbeat resync service", slog.Any("err", err))
113+
} else {
114+
slog.Info("Heartbeat resync service started")
115+
defer heartbeatResyncService.Stop()
116+
}
117+
}
118+
108119
// Create processor instance
109120
processor := daemon.NewSocketHandler(&cfg, pubsub)
110121

daemon/handlers.go

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -17,16 +17,26 @@ func SocketTopicProccessor(messages <-chan *message.Message) {
1717
if err := json.Unmarshal(msg.Payload, &socketMsg); err != nil {
1818
slog.ErrorContext(ctx, "failed to parse socket message", slog.Any("err", err))
1919
msg.Nack()
20+
continue
2021
}
2122

22-
if socketMsg.Type == SocketMessageTypeSync {
23-
err := handlePubSubSync(ctx, socketMsg.Payload)
24-
if err != nil {
25-
slog.ErrorContext(ctx, "failed to parse socket message", slog.Any("err", err))
26-
msg.Nack()
27-
} else {
28-
msg.Ack()
29-
}
23+
var err error
24+
switch socketMsg.Type {
25+
case SocketMessageTypeSync:
26+
err = handlePubSubSync(ctx, socketMsg.Payload)
27+
case SocketMessageTypeHeartbeat:
28+
err = handlePubSubHeartbeat(ctx, socketMsg.Payload)
29+
default:
30+
slog.ErrorContext(ctx, "unknown socket message type", slog.String("type", string(socketMsg.Type)))
31+
msg.Nack()
32+
continue
33+
}
34+
35+
if err != nil {
36+
slog.ErrorContext(ctx, "failed to handle socket message", slog.Any("err", err), slog.String("type", string(socketMsg.Type)))
37+
msg.Nack()
38+
} else {
39+
msg.Ack()
3040
}
3141
}
3242
}

daemon/handlers.heartbeat.go

Lines changed: 80 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,80 @@
1+
package daemon
2+
3+
import (
4+
"context"
5+
"encoding/json"
6+
"fmt"
7+
"log/slog"
8+
"os"
9+
10+
"github.com/malamtime/cli/model"
11+
)
12+
13+
func handlePubSubHeartbeat(ctx context.Context, socketMsgPayload interface{}) error {
14+
pb, err := json.Marshal(socketMsgPayload)
15+
if err != nil {
16+
slog.Error("Failed to marshal the heartbeat payload again for unmarshal", slog.Any("payload", socketMsgPayload))
17+
return err
18+
}
19+
20+
var heartbeatPayload model.HeartbeatPayload
21+
err = json.Unmarshal(pb, &heartbeatPayload)
22+
if err != nil {
23+
slog.Error("Failed to parse heartbeat payload", slog.Any("payload", socketMsgPayload))
24+
return err
25+
}
26+
27+
if len(heartbeatPayload.Heartbeats) == 0 {
28+
slog.Debug("Empty heartbeat payload, skipping")
29+
return nil
30+
}
31+
32+
cfg, err := stConfig.ReadConfigFile(ctx)
33+
if err != nil {
34+
slog.Error("Failed to read config file", slog.Any("err", err))
35+
return err
36+
}
37+
38+
// Try to send to server
39+
err = model.SendHeartbeatsToServer(ctx, cfg, heartbeatPayload)
40+
if err != nil {
41+
slog.Warn("Failed to send heartbeats to server, saving to local file", slog.Any("err", err))
42+
// On failure, save to local file
43+
if saveErr := saveHeartbeatToFile(heartbeatPayload); saveErr != nil {
44+
slog.Error("Failed to save heartbeat to local file", slog.Any("err", saveErr))
45+
return saveErr
46+
}
47+
// Return nil because we saved the data locally - don't nack the message
48+
return nil
49+
}
50+
51+
slog.Info("Successfully sent heartbeats to server", slog.Int("count", len(heartbeatPayload.Heartbeats)))
52+
return nil
53+
}
54+
55+
// saveHeartbeatToFile appends a heartbeat payload as a single JSON line to the log file
56+
func saveHeartbeatToFile(payload model.HeartbeatPayload) error {
57+
logFilePath := os.ExpandEnv(fmt.Sprintf("%s/%s", "$HOME", model.HEARTBEAT_LOG_FILE))
58+
59+
// Open file for appending, create if not exists
60+
file, err := os.OpenFile(logFilePath, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0644)
61+
if err != nil {
62+
return fmt.Errorf("failed to open heartbeat log file: %w", err)
63+
}
64+
defer file.Close()
65+
66+
// Marshal payload to JSON
67+
data, err := json.Marshal(payload)
68+
if err != nil {
69+
return fmt.Errorf("failed to marshal heartbeat payload: %w", err)
70+
}
71+
72+
// Write as single line with newline
73+
_, err = file.Write(append(data, '\n'))
74+
if err != nil {
75+
return fmt.Errorf("failed to write heartbeat to file: %w", err)
76+
}
77+
78+
slog.Debug("Saved heartbeat to local file", slog.String("path", logFilePath), slog.Int("count", len(payload.Heartbeats)))
79+
return nil
80+
}

daemon/heartbeat_resync.go

Lines changed: 181 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,181 @@
1+
package daemon
2+
3+
import (
4+
"bufio"
5+
"context"
6+
"encoding/json"
7+
"fmt"
8+
"log/slog"
9+
"os"
10+
"sync"
11+
"time"
12+
13+
"github.com/malamtime/cli/model"
14+
)
15+
16+
const (
17+
// HeartbeatResyncInterval is the interval for retrying failed heartbeats
18+
HeartbeatResyncInterval = 30 * time.Minute
19+
)
20+
21+
// HeartbeatResyncService handles periodic resync of failed heartbeats
22+
type HeartbeatResyncService struct {
23+
config model.ShellTimeConfig
24+
ticker *time.Ticker
25+
stopChan chan struct{}
26+
wg sync.WaitGroup
27+
}
28+
29+
// NewHeartbeatResyncService creates a new heartbeat resync service
30+
func NewHeartbeatResyncService(config model.ShellTimeConfig) *HeartbeatResyncService {
31+
return &HeartbeatResyncService{
32+
config: config,
33+
stopChan: make(chan struct{}),
34+
}
35+
}
36+
37+
// Start begins the periodic resync job
38+
func (s *HeartbeatResyncService) Start(ctx context.Context) error {
39+
s.ticker = time.NewTicker(HeartbeatResyncInterval)
40+
s.wg.Add(1)
41+
42+
go func() {
43+
defer s.wg.Done()
44+
45+
// Run once at startup
46+
s.resync(ctx)
47+
48+
for {
49+
select {
50+
case <-s.ticker.C:
51+
s.resync(ctx)
52+
case <-s.stopChan:
53+
return
54+
case <-ctx.Done():
55+
return
56+
}
57+
}
58+
}()
59+
60+
slog.Info("Heartbeat resync service started", slog.Duration("interval", HeartbeatResyncInterval))
61+
return nil
62+
}
63+
64+
// Stop stops the resync service
65+
func (s *HeartbeatResyncService) Stop() {
66+
if s.ticker != nil {
67+
s.ticker.Stop()
68+
}
69+
close(s.stopChan)
70+
s.wg.Wait()
71+
slog.Info("Heartbeat resync service stopped")
72+
}
73+
74+
// resync reads failed heartbeats from the log file and attempts to send them
75+
func (s *HeartbeatResyncService) resync(ctx context.Context) {
76+
logFilePath := os.ExpandEnv(fmt.Sprintf("%s/%s", "$HOME", model.HEARTBEAT_LOG_FILE))
77+
78+
// Check if file exists
79+
if _, err := os.Stat(logFilePath); os.IsNotExist(err) {
80+
slog.Debug("No heartbeat log file found, nothing to resync")
81+
return
82+
}
83+
84+
// Read the file
85+
file, err := os.Open(logFilePath)
86+
if err != nil {
87+
slog.Error("Failed to open heartbeat log file for resync", slog.Any("err", err))
88+
return
89+
}
90+
91+
var lines []string
92+
scanner := bufio.NewScanner(file)
93+
for scanner.Scan() {
94+
line := scanner.Text()
95+
if line != "" {
96+
lines = append(lines, line)
97+
}
98+
}
99+
file.Close()
100+
101+
if err := scanner.Err(); err != nil {
102+
slog.Error("Error reading heartbeat log file", slog.Any("err", err))
103+
return
104+
}
105+
106+
if len(lines) == 0 {
107+
slog.Debug("No failed heartbeats to resync")
108+
return
109+
}
110+
111+
slog.Info("Starting heartbeat resync", slog.Int("pendingCount", len(lines)))
112+
113+
// Process each line
114+
var failedLines []string
115+
successCount := 0
116+
117+
for _, line := range lines {
118+
var payload model.HeartbeatPayload
119+
if err := json.Unmarshal([]byte(line), &payload); err != nil {
120+
slog.Error("Failed to parse heartbeat line, discarding", slog.Any("err", err), slog.String("line", line))
121+
continue
122+
}
123+
124+
// Try to send to server
125+
if err := model.SendHeartbeatsToServer(ctx, s.config, payload); err != nil {
126+
slog.Warn("Failed to resync heartbeat, keeping for next retry", slog.Any("err", err))
127+
failedLines = append(failedLines, line)
128+
} else {
129+
successCount++
130+
}
131+
}
132+
133+
// Rewrite the file with only failed lines
134+
if err := s.rewriteLogFile(logFilePath, failedLines); err != nil {
135+
slog.Error("Failed to update heartbeat log file", slog.Any("err", err))
136+
return
137+
}
138+
139+
slog.Info("Heartbeat resync completed",
140+
slog.Int("success", successCount),
141+
slog.Int("remaining", len(failedLines)))
142+
}
143+
144+
// rewriteLogFile atomically rewrites the log file with the given lines
145+
func (s *HeartbeatResyncService) rewriteLogFile(logFilePath string, lines []string) error {
146+
// If no lines remaining, remove the file
147+
if len(lines) == 0 {
148+
if err := os.Remove(logFilePath); err != nil && !os.IsNotExist(err) {
149+
return fmt.Errorf("failed to remove empty log file: %w", err)
150+
}
151+
return nil
152+
}
153+
154+
// Write to temp file first
155+
tempFile := logFilePath + ".tmp"
156+
file, err := os.OpenFile(tempFile, os.O_CREATE|os.O_WRONLY|os.O_TRUNC, 0644)
157+
if err != nil {
158+
return fmt.Errorf("failed to create temp file: %w", err)
159+
}
160+
161+
for _, line := range lines {
162+
if _, err := file.WriteString(line + "\n"); err != nil {
163+
file.Close()
164+
os.Remove(tempFile)
165+
return fmt.Errorf("failed to write to temp file: %w", err)
166+
}
167+
}
168+
169+
if err := file.Close(); err != nil {
170+
os.Remove(tempFile)
171+
return fmt.Errorf("failed to close temp file: %w", err)
172+
}
173+
174+
// Atomic rename
175+
if err := os.Rename(tempFile, logFilePath); err != nil {
176+
os.Remove(tempFile)
177+
return fmt.Errorf("failed to rename temp file: %w", err)
178+
}
179+
180+
return nil
181+
}

daemon/socket.go

Lines changed: 18 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,8 @@ import (
1414
type SocketMessageType string
1515

1616
const (
17-
SocketMessageTypeSync SocketMessageType = "sync"
17+
SocketMessageTypeSync SocketMessageType = "sync"
18+
SocketMessageTypeHeartbeat SocketMessageType = "heartbeat"
1819
)
1920

2021
type SocketMessage struct {
@@ -112,6 +113,22 @@ func (p *SocketHandler) handleConnection(conn net.Conn) {
112113
if err := p.channel.Publish(PubSubTopic, chMsg); err != nil {
113114
slog.Error("Error to publish topic", slog.Any("err", err))
114115
}
116+
case SocketMessageTypeHeartbeat:
117+
// Only process heartbeat if codeTracking is enabled
118+
if p.config.CodeTracking == nil || p.config.CodeTracking.Enabled == nil || !*p.config.CodeTracking.Enabled {
119+
slog.Debug("Heartbeat message received but codeTracking is disabled, ignoring")
120+
return
121+
}
122+
buf, err := json.Marshal(msg)
123+
if err != nil {
124+
slog.Error("Error encoding heartbeat message", slog.Any("err", err))
125+
return
126+
}
127+
128+
chMsg := message.NewMessage(watermill.NewUUID(), buf)
129+
if err := p.channel.Publish(PubSubTopic, chMsg); err != nil {
130+
slog.Error("Error publishing heartbeat topic", slog.Any("err", err))
131+
}
115132
default:
116133
slog.Error("Unknown message type:", slog.String("messageType", string(msg.Type)))
117134
}

model/api_heartbeat.go

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
package model
2+
3+
import (
4+
"context"
5+
"net/http"
6+
"time"
7+
)
8+
9+
// SendHeartbeatsToServer sends heartbeat data to the server
10+
func SendHeartbeatsToServer(ctx context.Context, cfg ShellTimeConfig, payload HeartbeatPayload) error {
11+
ctx, span := modelTracer.Start(ctx, "api.sendHeartbeats")
12+
defer span.End()
13+
14+
endpoint := Endpoint{
15+
Token: cfg.Token,
16+
APIEndpoint: cfg.APIEndpoint,
17+
}
18+
19+
var response HeartbeatResponse
20+
err := SendHTTPRequestJSON(HTTPRequestOptions[HeartbeatPayload, HeartbeatResponse]{
21+
Context: ctx,
22+
Endpoint: endpoint,
23+
Method: http.MethodPost,
24+
Path: "/api/v1/heartbeats",
25+
Payload: payload,
26+
Response: &response,
27+
Timeout: 10 * time.Second,
28+
})
29+
30+
if err != nil {
31+
return err
32+
}
33+
34+
return nil
35+
}

0 commit comments

Comments
 (0)