Skip to content

Commit bd6fc3f

Browse files
feat: new queue and workers
1 parent cf1f11a commit bd6fc3f

12 files changed

Lines changed: 241 additions & 3 deletions

File tree

apps/api/.env.dev.example

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
DATABASE_URL="postgres://postgres:postgres@postgres:5432/coredb"
22
DATABASE_URL_MIGRATION="postgres://postgres:postgres@localhost:5432/coredb"
3+
REDIS_URL="redis:6379"
34

45
# For OAuth
56
AUTH_DISCORD_CLIENT_ID=

apps/api/cmd/api/main.go

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import (
44
"net/http"
55
"time"
66

7+
"github.com/hibiken/asynq"
78
"github.com/rs/zerolog/log"
89
"github.com/swamphacks/core/apps/api/internal/api"
910
"github.com/swamphacks/core/apps/api/internal/api/handlers"
@@ -38,6 +39,11 @@ func main() {
3839
},
3940
}
4041

42+
// Create asynq client
43+
taskQueueClient := asynq.NewClient(asynq.RedisClientOpt{
44+
Addr: cfg.RedisURL,
45+
})
46+
4147
// Create new middleware injectable
4248
mw := middleware.NewMiddleware(database, logger, cfg)
4349

@@ -50,9 +56,10 @@ func main() {
5056
// Injections into services
5157
authService := services.NewAuthService(userRepo, accountRepo, sessionRepo, txm, client, logger, &cfg.Auth)
5258
eventInterestService := services.NewEventInterestService(eventInterestRepo, logger)
59+
emailService := services.NewEmailService(taskQueueClient, logger)
5360

5461
// Injections into handlers
55-
apiHandlers := handlers.NewHandlers(authService, eventInterestService, cfg, logger)
62+
apiHandlers := handlers.NewHandlers(authService, eventInterestService, emailService, cfg, logger)
5663

5764
api := api.NewAPI(&logger, apiHandlers, mw)
5865

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,27 @@
1+
FROM golang:1.24-alpine AS base
2+
3+
WORKDIR /app
4+
5+
COPY ../../go.mod ../../go.sum ./
6+
7+
RUN go mod download
8+
9+
COPY ../../ ./
10+
11+
# Dev
12+
FROM base AS dev
13+
14+
RUN go install github.com/air-verse/air@latest
15+
16+
CMD ["air"]
17+
18+
# Production
19+
FROM base as prod
20+
21+
RUN CGO_ENABLED=0 GOOS=linux go build -ldflags="-s -w" -o email_worker ./
22+
23+
RUN apk --no-cache add ca-certificates
24+
25+
CMD ["./email_worker"]
26+
27+

apps/api/cmd/email_worker/main.go

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,21 +1,39 @@
11
package main
22

33
import (
4+
"fmt"
45
"log"
56

67
"github.com/hibiken/asynq"
8+
"github.com/swamphacks/core/apps/api/internal/config"
9+
"github.com/swamphacks/core/apps/api/internal/logger"
10+
"github.com/swamphacks/core/apps/api/internal/services"
11+
"github.com/swamphacks/core/apps/api/internal/tasks"
12+
"github.com/swamphacks/core/apps/api/internal/workers"
713
)
814

915
func main() {
16+
logger := logger.New()
17+
cfg := config.Load()
18+
1019
srv := asynq.NewServer(
11-
asynq.RedisClientOpt{Addr: "redis:6379"},
20+
asynq.RedisClientOpt{Addr: cfg.RedisURL},
1221
asynq.Config{
1322
Concurrency: 10,
23+
Queues: map[string]int{
24+
"email": 10,
25+
},
1426
},
1527
)
1628

29+
emailService := services.NewEmailService(nil, logger)
30+
emailWorker := workers.NewEmailWorker(emailService, logger)
31+
1732
mux := asynq.NewServeMux()
1833

34+
mux.HandleFunc(tasks.TypeSendEmail, emailWorker.HandleSendEmailTask)
35+
fmt.Println("Starting email worker")
36+
1937
if err := srv.Run(mux); err != nil {
2038
log.Fatalf("Failed to run email worker")
2139
}

apps/api/internal/api/api.go

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,11 @@ func (api *API) setupRoutes(mw *mw.Middleware) {
7171
r.Post("/{eventId}/interest", api.Handlers.EventInterest.AddEmailToEvent)
7272
})
7373

74+
// Email routes
75+
api.Router.Route("/email", func(r chi.Router) {
76+
r.Post("/queue", api.Handlers.Email.QueueEmail)
77+
})
78+
7479
// Protected test routes
7580
api.Router.Route("/protected", func(r chi.Router) {
7681
r.Use(mw.Auth.RequireAuth)
Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
package handlers
2+
3+
import (
4+
"encoding/json"
5+
"net/http"
6+
7+
"github.com/rs/zerolog"
8+
res "github.com/swamphacks/core/apps/api/internal/api/response"
9+
"github.com/swamphacks/core/apps/api/internal/email"
10+
"github.com/swamphacks/core/apps/api/internal/services"
11+
)
12+
13+
type EmailHandler struct {
14+
emailService *services.EmailService
15+
logger zerolog.Logger
16+
}
17+
18+
func NewEmailHandler(emailService *services.EmailService, logger zerolog.Logger) *EmailHandler {
19+
return &EmailHandler{
20+
emailService: emailService,
21+
logger: logger.With().Str("handler", "EmailHandler").Str("component", "email").Logger(),
22+
}
23+
}
24+
25+
type QueueEmailRequest struct {
26+
To string `json:"to"`
27+
From string `json:"from"`
28+
Body string `json:"body"`
29+
}
30+
31+
func (h *EmailHandler) QueueEmail(w http.ResponseWriter, r *http.Request) {
32+
var req QueueEmailRequest
33+
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
34+
res.SendError(w, http.StatusBadRequest, res.NewError("invalid_request", "Could not parse request body"))
35+
return
36+
}
37+
38+
if !email.IsValidEmail(req.To) || !email.IsValidEmail(req.From) {
39+
res.SendError(w, http.StatusBadRequest, res.NewError("malformed_email", "To and/or From email is malformed or missing"))
40+
return
41+
}
42+
43+
if req.Body == "" {
44+
res.SendError(w, http.StatusBadRequest, res.NewError("missing_body", "Body is missing or is an empty string."))
45+
return
46+
}
47+
48+
taskInfo, err := h.emailService.QueueSendEmail(req.To, req.From, req.Body)
49+
if err != nil {
50+
h.logger.Err(err).Msg("Failed to queue email sending from EmailHandler")
51+
res.SendError(w, http.StatusInternalServerError, res.NewError("internal_err", "The server went kaput while queueing email sending"))
52+
return
53+
}
54+
55+
h.logger.Info().Str("TaskID", taskInfo.ID).Str("Task Queue", taskInfo.Queue).Str("Task Type", taskInfo.Type).Msg("Queued Send Email task!")
56+
57+
w.WriteHeader(http.StatusCreated)
58+
}

apps/api/internal/api/handlers/handlers.go

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -9,11 +9,13 @@ import (
99
type Handlers struct {
1010
Auth *AuthHandler
1111
EventInterest *EventInterestHandler
12+
Email *EmailHandler
1213
}
1314

14-
func NewHandlers(authService *services.AuthService, eventInterestService *services.EventInterestService, cfg *config.Config, logger zerolog.Logger) *Handlers {
15+
func NewHandlers(authService *services.AuthService, eventInterestService *services.EventInterestService, emailService *services.EmailService, cfg *config.Config, logger zerolog.Logger) *Handlers {
1516
return &Handlers{
1617
Auth: NewAuthHandler(authService, cfg, logger),
1718
EventInterest: NewEventInterestHandler(eventInterestService, cfg, logger),
19+
Email: NewEmailHandler(emailService, logger),
1820
}
1921
}

apps/api/internal/config/config.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,7 @@ type AuthConfig struct {
2626

2727
type Config struct {
2828
DatabaseURL string `env:"DATABASE_URL"`
29+
RedisURL string `env:"REDIS_URL"`
2930
Port string `env:"PORT" envDefault:"8080"`
3031

3132
Auth AuthConfig `envPrefix:"AUTH_"`
Lines changed: 45 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,45 @@
1+
package services
2+
3+
import (
4+
"github.com/hibiken/asynq"
5+
"github.com/rs/zerolog"
6+
"github.com/swamphacks/core/apps/api/internal/tasks"
7+
)
8+
9+
type EmailService struct {
10+
logger zerolog.Logger
11+
taskQueue *asynq.Client
12+
// Add SMTP client later
13+
}
14+
15+
func NewEmailService(taskQueue *asynq.Client, logger zerolog.Logger) *EmailService {
16+
return &EmailService{
17+
logger: logger.With().Str("service", "EmailService").Str("component", "email").Logger(),
18+
taskQueue: taskQueue,
19+
}
20+
}
21+
22+
func (s *EmailService) SendEmail(to, from, body string) error {
23+
s.logger.Info().Msgf("Sending from %s to %s with body %s", from, to, body)
24+
return nil
25+
}
26+
27+
func (s *EmailService) QueueSendEmail(to, from, body string) (*asynq.TaskInfo, error) {
28+
task, err := tasks.NewTaskSendEmail(tasks.SendEmailPayload{
29+
To: to,
30+
From: from,
31+
Body: body,
32+
})
33+
if err != nil {
34+
s.logger.Err(err).Msg("Failed to create Send Email task")
35+
return nil, err
36+
}
37+
38+
info, err := s.taskQueue.Enqueue(task, asynq.Queue("email"))
39+
if err != nil {
40+
s.logger.Err(err).Msg("Failed to queue Send Email task")
41+
return nil, err
42+
}
43+
44+
return info, nil
45+
}

apps/api/internal/tasks/email.go

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,26 @@
1+
package tasks
2+
3+
import (
4+
"encoding/json"
5+
6+
"github.com/hibiken/asynq"
7+
)
8+
9+
const (
10+
TypeSendEmail = "email:send"
11+
)
12+
13+
type SendEmailPayload struct {
14+
To string
15+
From string
16+
Body string
17+
}
18+
19+
func NewTaskSendEmail(payload SendEmailPayload) (*asynq.Task, error) {
20+
data, err := json.Marshal(payload)
21+
if err != nil {
22+
return nil, err
23+
}
24+
25+
return asynq.NewTask(TypeSendEmail, data), nil
26+
}

0 commit comments

Comments
 (0)