Implements a comprehensive monitoring system for the admin interface: Backend: - New monitoring package with Redis ring buffer for log storage - Zerolog MultiWriter to capture logs to Redis - System stats collection (CPU, memory, disk, goroutines, GC) - HTTP metrics middleware (request counts, latency, error rates) - Asynq queue stats for worker process - WebSocket endpoint for real-time log streaming - Admin auth middleware now accepts token in query params (for WebSocket) Frontend: - New monitoring page with tabs (Overview, Logs, API Stats, Worker Stats) - Real-time log viewer with level filtering and search - System stats cards showing CPU, memory, goroutines, uptime - HTTP endpoint statistics table - Asynq queue depth visualization - Enable/disable monitoring toggle in settings Memory safeguards: - Max 200 unique endpoints tracked - Hourly stats reset to prevent unbounded growth - Max 1000 log entries in ring buffer - Max 1000 latency samples for P95 calculation 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude <noreply@anthropic.com>
207 lines
7.0 KiB
Go
207 lines
7.0 KiB
Go
package main
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"os"
|
|
"os/signal"
|
|
"syscall"
|
|
|
|
"github.com/hibiken/asynq"
|
|
"github.com/redis/go-redis/v9"
|
|
"github.com/rs/zerolog/log"
|
|
|
|
"github.com/treytartt/casera-api/internal/config"
|
|
"github.com/treytartt/casera-api/internal/database"
|
|
"github.com/treytartt/casera-api/internal/monitoring"
|
|
"github.com/treytartt/casera-api/internal/push"
|
|
"github.com/treytartt/casera-api/internal/repositories"
|
|
"github.com/treytartt/casera-api/internal/services"
|
|
"github.com/treytartt/casera-api/internal/worker/jobs"
|
|
"github.com/treytartt/casera-api/pkg/utils"
|
|
)
|
|
|
|
func main() {
|
|
// Initialize logger
|
|
utils.InitLogger(true)
|
|
|
|
// Load configuration
|
|
cfg, err := config.Load()
|
|
if err != nil {
|
|
log.Fatal().Err(err).Msg("Failed to load configuration")
|
|
}
|
|
|
|
// Initialize database
|
|
db, err := database.Connect(&cfg.Database, cfg.Server.Debug)
|
|
if err != nil {
|
|
log.Fatal().Err(err).Msg("Failed to connect to database")
|
|
}
|
|
log.Info().Msg("Connected to database")
|
|
|
|
// Get underlying *sql.DB for cleanup
|
|
sqlDB, _ := db.DB()
|
|
defer sqlDB.Close()
|
|
|
|
// Initialize push client (APNs + FCM)
|
|
var pushClient *push.Client
|
|
pushClient, err = push.NewClient(&cfg.Push)
|
|
if err != nil {
|
|
log.Warn().Err(err).Msg("Failed to initialize push client - push notifications disabled")
|
|
} else {
|
|
log.Info().
|
|
Bool("ios_enabled", pushClient.IsIOSEnabled()).
|
|
Bool("android_enabled", pushClient.IsAndroidEnabled()).
|
|
Msg("Push notification client initialized")
|
|
}
|
|
|
|
// Initialize email service (optional)
|
|
var emailService *services.EmailService
|
|
if cfg.Email.Host != "" {
|
|
emailService = services.NewEmailService(&cfg.Email)
|
|
log.Info().Str("host", cfg.Email.Host).Msg("Email service initialized")
|
|
}
|
|
|
|
// Initialize notification service for actionable push notifications
|
|
notificationRepo := repositories.NewNotificationRepository(db)
|
|
notificationService := services.NewNotificationService(notificationRepo, pushClient)
|
|
log.Info().Msg("Notification service initialized")
|
|
|
|
// Parse Redis URL for Asynq
|
|
redisOpt, err := asynq.ParseRedisURI(cfg.Redis.URL)
|
|
if err != nil {
|
|
log.Fatal().Err(err).Msg("Failed to parse Redis URL")
|
|
}
|
|
|
|
// Initialize monitoring service (if Redis is available)
|
|
var monitoringService *monitoring.Service
|
|
redisClientOpt, ok := redisOpt.(asynq.RedisClientOpt)
|
|
if ok {
|
|
redisClient := redis.NewClient(&redis.Options{
|
|
Addr: redisClientOpt.Addr,
|
|
Password: redisClientOpt.Password,
|
|
DB: redisClientOpt.DB,
|
|
})
|
|
|
|
// Verify Redis connection
|
|
if err := redisClient.Ping(context.Background()).Err(); err != nil {
|
|
log.Warn().Err(err).Msg("Failed to connect to Redis for monitoring - monitoring disabled")
|
|
} else {
|
|
monitoringService = monitoring.NewService(monitoring.Config{
|
|
Process: "worker",
|
|
RedisClient: redisClient,
|
|
DB: db, // Pass database for enable_monitoring setting sync
|
|
})
|
|
|
|
// Reinitialize logger with monitoring writer
|
|
utils.InitLoggerWithWriter(cfg.Server.Debug, monitoringService.LogWriter())
|
|
|
|
// Create Asynq inspector for queue statistics
|
|
inspector := asynq.NewInspector(redisOpt)
|
|
monitoringService.SetAsynqInspector(inspector)
|
|
|
|
// Start stats collection
|
|
monitoringService.Start()
|
|
defer monitoringService.Stop()
|
|
|
|
log.Info().
|
|
Bool("log_capture_enabled", monitoringService.IsEnabled()).
|
|
Msg("Monitoring service initialized")
|
|
}
|
|
}
|
|
|
|
// Create Asynq server
|
|
srv := asynq.NewServer(
|
|
redisOpt,
|
|
asynq.Config{
|
|
Concurrency: 10,
|
|
Queues: map[string]int{
|
|
"critical": 6,
|
|
"default": 3,
|
|
"low": 1,
|
|
},
|
|
ErrorHandler: asynq.ErrorHandlerFunc(func(ctx context.Context, task *asynq.Task, err error) {
|
|
log.Error().
|
|
Err(err).
|
|
Str("type", task.Type()).
|
|
Bytes("payload", task.Payload()).
|
|
Msg("Task processing failed")
|
|
}),
|
|
},
|
|
)
|
|
|
|
// Create job handler
|
|
jobHandler := jobs.NewHandler(db, pushClient, emailService, notificationService, cfg)
|
|
|
|
// Create Asynq mux and register handlers
|
|
mux := asynq.NewServeMux()
|
|
mux.HandleFunc(jobs.TypeTaskReminder, jobHandler.HandleTaskReminder)
|
|
mux.HandleFunc(jobs.TypeOverdueReminder, jobHandler.HandleOverdueReminder)
|
|
mux.HandleFunc(jobs.TypeDailyDigest, jobHandler.HandleDailyDigest)
|
|
mux.HandleFunc(jobs.TypeSendEmail, jobHandler.HandleSendEmail)
|
|
mux.HandleFunc(jobs.TypeSendPush, jobHandler.HandleSendPush)
|
|
mux.HandleFunc(jobs.TypeOnboardingEmails, jobHandler.HandleOnboardingEmails)
|
|
|
|
// Start scheduler for periodic tasks
|
|
scheduler := asynq.NewScheduler(redisOpt, nil)
|
|
|
|
// Schedule task reminder notifications (runs every hour to support per-user custom times)
|
|
// The job handler filters users based on their preferred notification hour
|
|
if _, err := scheduler.Register("0 * * * *", asynq.NewTask(jobs.TypeTaskReminder, nil)); err != nil {
|
|
log.Fatal().Err(err).Msg("Failed to register task reminder job")
|
|
}
|
|
log.Info().Str("cron", "0 * * * *").Int("default_hour", cfg.Worker.TaskReminderHour).Msg("Registered task reminder job (runs hourly for per-user times)")
|
|
|
|
// Schedule overdue reminder (runs every hour to support per-user custom times)
|
|
if _, err := scheduler.Register("0 * * * *", asynq.NewTask(jobs.TypeOverdueReminder, nil)); err != nil {
|
|
log.Fatal().Err(err).Msg("Failed to register overdue reminder job")
|
|
}
|
|
log.Info().Str("cron", "0 * * * *").Int("default_hour", cfg.Worker.OverdueReminderHour).Msg("Registered overdue reminder job (runs hourly for per-user times)")
|
|
|
|
// Schedule daily digest (runs every hour to support per-user custom times)
|
|
// The job handler filters users based on their preferred notification hour
|
|
if _, err := scheduler.Register("0 * * * *", asynq.NewTask(jobs.TypeDailyDigest, nil)); err != nil {
|
|
log.Fatal().Err(err).Msg("Failed to register daily digest job")
|
|
}
|
|
log.Info().Str("cron", "0 * * * *").Int("default_hour", cfg.Worker.DailyNotifHour).Msg("Registered daily digest job (runs hourly for per-user times)")
|
|
|
|
// Schedule onboarding emails (runs daily at 10:00 AM UTC)
|
|
// Sends emails to users who haven't created residences or tasks after registration
|
|
if _, err := scheduler.Register("0 10 * * *", asynq.NewTask(jobs.TypeOnboardingEmails, nil)); err != nil {
|
|
log.Fatal().Err(err).Msg("Failed to register onboarding emails job")
|
|
}
|
|
log.Info().Str("cron", "0 10 * * *").Msg("Registered onboarding emails job (runs daily at 10:00 AM UTC)")
|
|
|
|
// Handle graceful shutdown
|
|
quit := make(chan os.Signal, 1)
|
|
signal.Notify(quit, syscall.SIGINT, syscall.SIGTERM)
|
|
|
|
// Start scheduler in goroutine
|
|
go func() {
|
|
if err := scheduler.Run(); err != nil {
|
|
log.Fatal().Err(err).Msg("Failed to start scheduler")
|
|
}
|
|
}()
|
|
|
|
// Start worker server in goroutine
|
|
go func() {
|
|
log.Info().Msg("Starting worker server...")
|
|
if err := srv.Run(mux); err != nil {
|
|
log.Fatal().Err(err).Msg("Failed to start worker server")
|
|
}
|
|
}()
|
|
|
|
<-quit
|
|
log.Info().Msg("Shutting down worker...")
|
|
|
|
// Graceful shutdown
|
|
srv.Shutdown()
|
|
scheduler.Shutdown()
|
|
|
|
log.Info().Msg("Worker stopped")
|
|
}
|
|
|
|
// formatCron creates a cron expression for a specific hour (runs at minute 0)
|
|
func formatCron(hour int) string {
|
|
return fmt.Sprintf("0 %02d * * *", hour)
|
|
}
|