Skip to content

Commit fe9e29d

Browse files
committed
added structured logging
Signed-off-by: redpinecube <tara.chakkithara@icloud.com>
1 parent 4295037 commit fe9e29d

4 files changed

Lines changed: 75 additions & 39 deletions

File tree

src/api.go

Lines changed: 44 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -2,8 +2,9 @@ package main
22

33
import (
44
"context"
5+
"errors"
56
"fmt"
6-
"log"
7+
"log/slog"
78
"net/http"
89
"os"
910
"os/signal"
@@ -20,44 +21,50 @@ type App struct {
2021
supervisor *Supervisor
2122
httpServer *http.Server
2223
wg sync.WaitGroup
24+
log *slog.Logger
2325
}
2426

25-
func NewApp(redisAddr, gpuType string) *App {
27+
func NewApp(redisAddr, gpuType string, log *slog.Logger) *App {
2628
client := redis.NewClient(&redis.Options{Addr: redisAddr})
27-
scheduler := NewScheduler(redisAddr)
29+
scheduler := NewScheduler(redisAddr, log)
2830

2931
consumerID := fmt.Sprintf("worker_%d", os.Getpid())
30-
supervisor := NewSupervisor(redisAddr, consumerID, gpuType)
32+
supervisor := NewSupervisor(redisAddr, consumerID, gpuType, log)
3133

3234
mux := http.NewServeMux()
3335
a := &App{
3436
redisClient: client,
3537
scheduler: scheduler,
3638
supervisor: supervisor,
3739
httpServer: &http.Server{Addr: ":3000", Handler: mux},
40+
log: log,
3841
}
3942

4043
mux.HandleFunc("/auth/login", a.login)
4144
mux.HandleFunc("/auth/refresh", a.refresh)
4245
mux.HandleFunc("/jobs", a.enqueueJob)
4346
mux.HandleFunc("/jobs/status", a.getJobStatus)
4447

48+
a.log.Info("new app initialized", "redis_address", redisAddr,
49+
"gpu_type", gpuType, "http_address", a.httpServer.Addr)
50+
4551
return a
4652
}
4753

4854
func (a *App) Start() error {
4955
// Connect to redis
5056
if err := a.redisClient.Ping(context.Background()).Err(); err != nil {
51-
return fmt.Errorf("redis ping failed: %w", err)
57+
a.log.Error("redis ping failed", "err", err)
58+
return err
5259
}
53-
60+
5461
// Launch HTTP server
5562
a.wg.Add(1)
5663
go func() {
5764
defer a.wg.Done()
58-
log.Println("HTTP server listening on", a.httpServer.Addr)
65+
slog.Info("http server started", "address", a.httpServer.Addr)
5966
if err := a.httpServer.ListenAndServe(); err != nil && err != http.ErrServerClosed {
60-
log.Fatalf("HTTP server error: %v", err)
67+
a.log.Error("HTTP server error", "err", err)
6168
}
6269
}()
6370

@@ -66,7 +73,7 @@ func (a *App) Start() error {
6673

6774
func (a *App) Shutdown(ctx context.Context) error {
6875
if err := a.httpServer.Shutdown(ctx); err != nil {
69-
log.Printf("error shutting down HTTP server: %v", err)
76+
a.log.Error("error shutting down HTTP server", "err", err)
7077
}
7178

7279
// Wait for ListenAndServe goroutine to finish
@@ -75,44 +82,61 @@ func (a *App) Shutdown(ctx context.Context) error {
7582
a.supervisor.Stop()
7683

7784
if err := a.scheduler.Close(); err != nil {
78-
log.Printf("error closing scheduler: %v", err)
85+
a.log.Error("error closing scheduler", "err", err)
86+
87+
} else {
88+
a.log.Info("scheduler closed successfully")
7989
}
8090

8191
if err := a.redisClient.Close(); err != nil {
82-
log.Printf("error closing redis client: %v", err)
92+
a.log.Error("error closing redis client", "err", err)
93+
} else {
94+
a.log.Info("redis client closed successfully")
8395
}
8496

97+
a.log.Info("shutdown completed")
98+
8599
return nil
86100
}
87101

88102
func main() {
89-
app := NewApp("localhost:6379", "AMD")
103+
log := slog.New(slog.NewJSONHandler(os.Stdout, nil))
104+
app := NewApp("localhost:6379", "AMD", log)
90105

91106
if err := app.Start(); err != nil {
92-
log.Fatalf("failed to start app: %v", err)
107+
log.Error("failed to start app", "err", err)
108+
os.Exit(1)
93109
}
94110

95111
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
96112
defer stop()
97113
<-ctx.Done()
98-
log.Println("shutdown signal received")
114+
log.Info("shutdown signal received")
99115

100116
shutdownCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second)
101117
defer cancel()
102118
if err := app.Shutdown(shutdownCtx); err != nil {
103-
log.Fatalf("shutdown error: %v", err)
119+
log.Error("shutdown error", "err", err)
104120
}
105121

106-
log.Println("all services stopped cleanly")
122+
log.Info("all services stopped cleanly")
107123
}
108124

109125
func (a *App) login(w http.ResponseWriter, r *http.Request) {
110126
ctx := r.Context()
127+
a.log.Info("login handler accessed", "remote_address", r.RemoteAddr)
111128
val, err := a.redisClient.Get(ctx, "some:key").Result()
112-
if err != nil || err != redis.Nil {
129+
if errors.Is(err, redis.Nil) {
130+
a.log.Info("redis key not found")
131+
http.Error(w, "redis key not found", http.StatusNotFound)
132+
return
133+
}
134+
if err != nil {
135+
a.log.Error("redis error on login", "err", err)
113136
http.Error(w, "redis error", http.StatusInternalServerError)
114137
return
115138
}
139+
a.log.Info("login success", "remote_address", r.RemoteAddr)
116140
fmt.Fprintf(w, "login page; redis says: %q\n", val)
117141
}
118142

@@ -121,14 +145,17 @@ func (a *App) refresh(w http.ResponseWriter, r *http.Request) {
121145
}
122146

123147
func (a *App) enqueueJob(w http.ResponseWriter, r *http.Request) {
148+
a.log.Info("enqueueJob handler accessed", "remote_address", r.RemoteAddr)
124149
payload := map[string]interface{}{
125150
"task_id": 123,
126151
"data": "test_data_123",
127152
}
128153
if err := a.scheduler.Enqueue("jobType", payload); err != nil {
154+
a.log.Error("enqueue failed", "err", err, "payload", payload)
129155
http.Error(w, "enqueue failed", http.StatusInternalServerError)
130156
return
131157
}
158+
a.log.Info("job enqueued", "payload", payload)
132159
w.WriteHeader(http.StatusAccepted)
133160
fmt.Fprint(w, "enqueued")
134161
}

src/int_test.go

Lines changed: 7 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -3,9 +3,10 @@ package main
33
import (
44
"context"
55
"fmt"
6+
"log/slog"
67
"os"
7-
"sync"
88
"os/signal"
9+
"sync"
910
"syscall"
1011
"testing"
1112

@@ -14,18 +15,19 @@ import (
1415

1516
func TestIntegration(t *testing.T) {
1617
redisAddr := "localhost:6379"
18+
log := slog.New(slog.NewJSONHandler(os.Stdout, nil))
1719

1820
client := redis.NewClient(&redis.Options{Addr: redisAddr})
1921
defer client.Close()
2022
if err := client.Ping(context.Background()).Err(); err != nil {
2123
t.Errorf("Failed to connect to Redis: %v", err)
2224
}
2325

24-
scheduler := NewScheduler(redisAddr)
26+
scheduler := NewScheduler(redisAddr, log)
2527
defer scheduler.Close()
2628

2729
consumerID := fmt.Sprintf("worker_%d", os.Getpid())
28-
supervisor := NewSupervisor(redisAddr, consumerID, "AMD")
30+
supervisor := NewSupervisor(redisAddr, consumerID, "AMD", log)
2931

3032
if err := supervisor.Start(); err != nil {
3133
t.Errorf("Failed to start supervisor: %v", err)
@@ -36,7 +38,7 @@ func TestIntegration(t *testing.T) {
3638

3739
// test jobs
3840
var wg sync.WaitGroup
39-
wg.Add(1)
41+
wg.Add(1)
4042
go func() {
4143
defer wg.Done()
4244
jobTypes := []string{"a", "b", "c"}
@@ -52,7 +54,7 @@ func TestIntegration(t *testing.T) {
5254
}
5355
}
5456
}()
55-
57+
5658
wg.Wait()
5759
supervisor.Stop()
5860
}

src/scheduler.go

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,7 +4,7 @@ import (
44
"context"
55
"encoding/json"
66
"fmt"
7-
"log"
7+
"log/slog"
88
"time"
99

1010
"github.com/redis/go-redis/v9"
@@ -13,16 +13,18 @@ import (
1313
type Scheduler struct {
1414
client *redis.Client
1515
ctx context.Context
16+
log *slog.Logger
1617
}
1718

18-
func NewScheduler(redisAddr string) *Scheduler {
19+
func NewScheduler(redisAddr string, log *slog.Logger) *Scheduler {
1920
client := redis.NewClient(&redis.Options{
2021
Addr: redisAddr,
2122
})
2223

2324
return &Scheduler{
2425
client: client,
2526
ctx: context.Background(),
27+
log: log,
2628
}
2729
}
2830

@@ -52,7 +54,7 @@ func (s *Scheduler) Enqueue(jobType string, payload map[string]interface{}) erro
5254
return fmt.Errorf("failed to enqueue job: %w", result.Err())
5355
}
5456

55-
log.Printf("Enqueued job %s of type %s", job.ID, job.Type)
57+
s.log.Info("enqueued job", "job_id", job.ID, "job_type", job.Type)
5658
return nil
5759
}
5860

src/supervisor.go

Lines changed: 19 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -3,8 +3,9 @@ package main
33
import (
44
"context"
55
"encoding/json"
6+
"errors"
67
"fmt"
7-
"log"
8+
"log/slog"
89
"sync"
910
"time"
1011

@@ -18,9 +19,10 @@ type Supervisor struct {
1819
consumerID string
1920
gpuType string
2021
wg sync.WaitGroup
22+
log *slog.Logger
2123
}
2224

23-
func NewSupervisor(redisAddr, consumerID, gpuType string) *Supervisor {
25+
func NewSupervisor(redisAddr, consumerID, gpuType string, log *slog.Logger) *Supervisor {
2426
client := redis.NewClient(&redis.Options{
2527
Addr: redisAddr,
2628
})
@@ -33,6 +35,7 @@ func NewSupervisor(redisAddr, consumerID, gpuType string) *Supervisor {
3335
cancel: cancel,
3436
consumerID: consumerID,
3537
gpuType: gpuType,
38+
log: log,
3639
}
3740
}
3841

@@ -46,7 +49,7 @@ func (s *Supervisor) Start() error {
4649
s.wg.Add(1)
4750
go s.processJobs()
4851

49-
log.Printf("Supervisor %s started with GPU: %s", s.consumerID, s.gpuType)
52+
s.log.Info("supervisor started", "consumer_id", s.consumerID, "gpu_type", s.gpuType)
5053
return nil
5154
}
5255

@@ -63,6 +66,8 @@ func (s *Supervisor) createConsumerGroup() error {
6366

6467
func (s *Supervisor) processJobs() {
6568
defer s.wg.Done()
69+
s.log.Info("job processor started", "consumer_id", s.consumerID)
70+
defer s.log.Info("job processor stopped", "consumer_id", s.consumerID)
6671

6772
for {
6873
select {
@@ -79,8 +84,8 @@ func (s *Supervisor) processJobs() {
7984
})
8085

8186
if result.Err() != nil {
82-
if result.Err() != redis.Nil {
83-
log.Printf("Error reading from stream: %v", result.Err())
87+
if !errors.Is(result.Err(), redis.Nil) {
88+
s.log.Error("error reading from stream", "error", result.Err())
8489
}
8590
continue
8691
}
@@ -98,36 +103,36 @@ func (s *Supervisor) processJobs() {
98103
func (s *Supervisor) handleMessage(message redis.XMessage) {
99104
jobData, ok := message.Values["data"].(string)
100105
if !ok {
101-
log.Printf("Invalid job data in message %s", message.ID)
106+
s.log.Error("invalid job data in message", "message_id", message.ID)
102107
s.ackMessage(message.ID)
103108
return
104109
}
105110

106111
var job Job
107112
if err := json.Unmarshal([]byte(jobData), &job); err != nil {
108-
log.Printf("Failed to unmarshal job data: %v", err)
113+
s.log.Error("failed to unmarshal job data", "error", err, "message_id", message.ID)
109114
s.ackMessage(message.ID)
110115
return
111116
}
112117

113118
// certain jobs require a specific GPU
114119
if !s.canHandleJob(job) {
115-
log.Printf("Job %s requires GPU type %s, but supervisor has %s - skipping",
116-
job.ID, job.RequiredGPU, s.gpuType)
120+
s.log.Info("skipping job due to GPU mismatch",
121+
"job_id", job.ID, "required_gpu", job.RequiredGPU, "supervisor_gpu", s.gpuType)
117122
// let another supervisor can pick it up
118123
return
119124
}
120125

121-
log.Printf("Processing job %s:%s", job.ID, job.Type)
126+
s.log.Info("processing job", "job_id", job.ID, "job_type", job.Type)
122127

123128
// Simulate job processing
124129
success := s.processJob(job)
125130

126131
if success {
127132
s.ackMessage(message.ID)
128-
log.Printf("Job %s completed successfully", job.ID)
133+
s.log.Info("job completed successfully", "job_id", job.ID)
129134
} else {
130-
log.Printf("Job %s failed", job.ID)
135+
s.log.Error("job failed", "job_id", job.ID)
131136
s.ackMessage(message.ID) // TODO: change this once we have docker support
132137
}
133138
}
@@ -151,12 +156,12 @@ func (s *Supervisor) processJob(job Job) bool {
151156
func (s *Supervisor) ackMessage(messageID string) {
152157
result := s.redisClient.XAck(s.ctx, StreamName, ConsumerGroup, messageID)
153158
if result.Err() != nil {
154-
log.Printf("Failed to ack message %s: %v", messageID, result.Err())
159+
s.log.Error("failed to ack message", "message_id", messageID, "error", result.Err())
155160
}
156161
}
157162

158163
func (s *Supervisor) Stop() {
159-
log.Printf("Stopping supervisor %s", s.consumerID)
164+
s.log.Info("stopping supervisor", "consumer_id", s.consumerID)
160165
s.cancel()
161166
s.wg.Wait()
162167
s.redisClient.Close()

0 commit comments

Comments
 (0)