split watcher and http server from sender + add golanci config
This commit is contained in:
+39
-107
@@ -1,51 +1,39 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"context"
|
||||
cfg "mailsrv/config"
|
||||
"mailsrv/mail"
|
||||
"mailsrv/runtime"
|
||||
"os"
|
||||
"os/signal"
|
||||
"path"
|
||||
"strings"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"net/http"
|
||||
"net/smtp"
|
||||
|
||||
"github.com/rs/zerolog/log"
|
||||
)
|
||||
|
||||
const (
|
||||
TickerInterval time.Duration = 10 * time.Second
|
||||
JSONSuffix string = ".json"
|
||||
ErrorSuffix string = ".err"
|
||||
)
|
||||
|
||||
type Sender struct {
|
||||
smtpConfig cfg.SMTPConfig
|
||||
// fetch this directory to collect `.json` e-mail format
|
||||
auth smtp.Auth
|
||||
smtpURL string
|
||||
|
||||
outboxPath string
|
||||
queue *runtime.Queue
|
||||
|
||||
queue *runtime.Queue
|
||||
}
|
||||
|
||||
func NewSender(config cfg.SMTPConfig, outboxPath string) Sender {
|
||||
return Sender{
|
||||
smtpConfig: config,
|
||||
auth: smtp.PlainAuth("", config.User, config.Password, config.URL),
|
||||
smtpURL: config.GetFullURL(),
|
||||
outboxPath: outboxPath,
|
||||
queue: runtime.NewQueue(),
|
||||
}
|
||||
}
|
||||
|
||||
func (s Sender) SendMail(email mail.Email) error {
|
||||
auth := smtp.PlainAuth("", s.smtpConfig.User, s.smtpConfig.Password, s.smtpConfig.Url)
|
||||
log.Debug().Msg("SMTP authentication succeed")
|
||||
|
||||
if err := smtp.SendMail(s.smtpConfig.GetFullUrl(), auth, email.Sender, email.GetReceivers(), email.Generate()); err != nil {
|
||||
func (s Sender) SendMail(email *mail.Email) error {
|
||||
if err := smtp.SendMail(s.smtpURL, s.auth, email.Sender, email.GetReceivers(), email.Generate()); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
@@ -53,80 +41,6 @@ func (s Sender) SendMail(email mail.Email) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s Sender) mailHandler(w http.ResponseWriter, r *http.Request) {
|
||||
content, err := io.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
log.Err(err).Msg("unable to read request body")
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
var mail mail.Email
|
||||
if err := json.Unmarshal(content, &mail); err != nil {
|
||||
log.Err(err).Msg("unable to deserialized request body into mail")
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
s.queue.Add(mail)
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}
|
||||
|
||||
func (s Sender) runHTTPserver() {
|
||||
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/mail", s.mailHandler)
|
||||
|
||||
if err := http.ListenAndServe(":1212", mux); err != nil {
|
||||
log.Err(err).Msg("http server stops listening")
|
||||
}
|
||||
}
|
||||
|
||||
// watchOutbox reads the `outbox` directory every `TickInterval` and put JSON format e-mail in the queue.
|
||||
func (s Sender) watchOutbox() {
|
||||
log.Info().Str("outbox", s.outboxPath).Msg("start watching outbox directory")
|
||||
|
||||
ticker := time.NewTicker(TickerInterval)
|
||||
|
||||
go func() {
|
||||
for range ticker.C {
|
||||
log.Debug().Str("action", "retrieving json e-mail format...").Str("path", s.outboxPath)
|
||||
|
||||
files, err := os.ReadDir(s.outboxPath)
|
||||
if err != nil && !os.IsExist(err) {
|
||||
log.Err(err).Msg("outbox directory does not exist")
|
||||
s.queue.Shutdown()
|
||||
}
|
||||
|
||||
for _, file := range files {
|
||||
filename := file.Name()
|
||||
if !strings.HasSuffix(filename, JSONSuffix) {
|
||||
log.Debug().Str("filename", filename).Msg("incorrect suffix")
|
||||
continue
|
||||
}
|
||||
|
||||
path := path.Join(s.outboxPath, filename)
|
||||
email, err := mail.FromJSON(path)
|
||||
|
||||
if err != nil {
|
||||
log.Err(err).Str("path", path).Msg("unable to parse JSON email")
|
||||
|
||||
// if JSON parsing failed the `path` is renamed with an error suffix to not watch it again
|
||||
newPath := fmt.Sprintf("%s%s", path, ErrorSuffix)
|
||||
if err := os.Rename(path, newPath); err != nil {
|
||||
log.Err(err).Str("path", path).Str("new path", newPath).Msg("unable to rename bad JSON email path")
|
||||
}
|
||||
|
||||
continue
|
||||
}
|
||||
|
||||
email.Path = path
|
||||
s.queue.Add(email)
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
// processNextEmail iterates over the queue and send email.
|
||||
func (s Sender) processNextEmail() bool {
|
||||
item, quit := s.queue.Get()
|
||||
@@ -141,7 +55,7 @@ func (s Sender) processNextEmail() bool {
|
||||
return true
|
||||
}
|
||||
|
||||
if err := s.SendMail(email); err != nil {
|
||||
if err := s.SendMail(&email); err != nil {
|
||||
log.Err(err).Msg("unable to send the email")
|
||||
}
|
||||
|
||||
@@ -169,26 +83,44 @@ func (s Sender) run() <-chan struct{} {
|
||||
return chQueue
|
||||
}
|
||||
|
||||
// Run launches the queue processing and the outbox watcher.
|
||||
// It catches `SIGINT` and `SIGTERM` to properly stopped the queue.
|
||||
// Run launches the queue processing, the outbox watcher and the HTTP server.
|
||||
// It catches `SIGINT` and `SIGTERM` to properly stopped the queue and the services.
|
||||
func (s Sender) Run() {
|
||||
log.Info().Msg("sender service is running")
|
||||
ctx, fnCancel := context.WithCancel(context.Background())
|
||||
|
||||
chSignal := make(chan os.Signal, 1)
|
||||
signal.Notify(chSignal, os.Interrupt, syscall.SIGTERM)
|
||||
|
||||
s.watchOutbox()
|
||||
chQueue := s.run()
|
||||
|
||||
go s.runHTTPserver()
|
||||
server := NewServer(ctx, "1212", s.queue)
|
||||
server.Serve()
|
||||
|
||||
watcher := NewDirectoryWatch(ctx, s.outboxPath, s.queue)
|
||||
watcher.Watch()
|
||||
|
||||
log.Info().Msg("sender service is running...")
|
||||
|
||||
select {
|
||||
case <-chSignal:
|
||||
log.Warn().Msg("stop signal received, stopping e-mail queue...")
|
||||
s.queue.Shutdown()
|
||||
case <-chQueue:
|
||||
log.Info().Msg("e-mail queue stopped successfully")
|
||||
log.Warn().Msg("stop signal received, stopping...")
|
||||
fnCancel()
|
||||
case <-watcher.Done():
|
||||
log.Warn().Msg("watcher is done, stopping...")
|
||||
fnCancel()
|
||||
case <-server.Done():
|
||||
log.Warn().Msg("server is done, stopping...")
|
||||
fnCancel()
|
||||
}
|
||||
|
||||
log.Info().Msg("sender service stopped successfully")
|
||||
<-server.Done()
|
||||
log.Info().Msg("http server stopped successfully")
|
||||
|
||||
<-watcher.Done()
|
||||
log.Info().Msg("watcher stopped successfully")
|
||||
|
||||
s.queue.Shutdown()
|
||||
<-chQueue
|
||||
|
||||
log.Info().Msg("mailsrv stopped gracefully")
|
||||
}
|
||||
|
||||
@@ -0,0 +1,99 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"mailsrv/mail"
|
||||
"mailsrv/runtime"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/rs/zerolog/log"
|
||||
)
|
||||
|
||||
type HTTPServer interface {
|
||||
Serve()
|
||||
Done() <-chan struct{}
|
||||
}
|
||||
|
||||
type Server struct {
|
||||
ctx context.Context
|
||||
fnCancel context.CancelFunc
|
||||
|
||||
port string
|
||||
queue *runtime.Queue
|
||||
|
||||
chDone chan struct{}
|
||||
}
|
||||
|
||||
func NewServer(ctx context.Context, port string, queue *runtime.Queue) Server {
|
||||
ctxChild, fnCancel := context.WithCancel(ctx)
|
||||
|
||||
return Server{
|
||||
ctx: ctxChild,
|
||||
fnCancel: fnCancel,
|
||||
port: port,
|
||||
queue: queue,
|
||||
chDone: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
func (s Server) Done() <-chan struct{} {
|
||||
return s.chDone
|
||||
}
|
||||
|
||||
func (s *Server) Serve() {
|
||||
mux := http.NewServeMux()
|
||||
mux.HandleFunc("/mail", s.handler)
|
||||
|
||||
log.Info().Str("port", s.port).Msg("http server is listening...")
|
||||
|
||||
server := &http.Server{
|
||||
Addr: fmt.Sprintf(":%s", s.port),
|
||||
Handler: mux,
|
||||
ReadTimeout: 10 * time.Second,
|
||||
WriteTimeout: 10 * time.Second,
|
||||
}
|
||||
|
||||
go func() {
|
||||
<-s.ctx.Done()
|
||||
if err := server.Shutdown(s.ctx); err != nil {
|
||||
log.Err(err).Msg("bad server shutdown")
|
||||
}
|
||||
s.chDone <- struct{}{}
|
||||
}()
|
||||
|
||||
go func() {
|
||||
if err := server.ListenAndServe(); err != nil {
|
||||
log.Err(err).Msg("http server stops listening")
|
||||
s.fnCancel()
|
||||
}
|
||||
}()
|
||||
}
|
||||
|
||||
func (s *Server) handler(w http.ResponseWriter, r *http.Request) {
|
||||
content, err := io.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
log.Err(err).Msg("unable to read request body")
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
var email mail.Email
|
||||
if err := json.Unmarshal(content, &email); err != nil {
|
||||
log.Err(err).Msg("unable to deserialized request body into mail")
|
||||
w.WriteHeader(http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
if err := email.Validate(); err != nil {
|
||||
log.Err(err).Msg("email validation failed")
|
||||
w.WriteHeader(http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
s.queue.Add(email)
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}
|
||||
@@ -0,0 +1,100 @@
|
||||
package services
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"mailsrv/mail"
|
||||
"mailsrv/runtime"
|
||||
"os"
|
||||
"path"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/rs/zerolog/log"
|
||||
)
|
||||
|
||||
const (
|
||||
TickerInterval time.Duration = 10 * time.Second
|
||||
JSONSuffix string = ".json"
|
||||
ErrorSuffix string = ".err"
|
||||
)
|
||||
|
||||
type Watcher interface {
|
||||
Watch()
|
||||
Done() <-chan struct{}
|
||||
}
|
||||
|
||||
// DirectoryWatch watches a directory every `tick` interval and collect email files.
|
||||
type DirectoryWatch struct {
|
||||
ctx context.Context
|
||||
fnCancel context.CancelFunc
|
||||
|
||||
outboxPath string
|
||||
queue *runtime.Queue
|
||||
}
|
||||
|
||||
func NewDirectoryWatch(ctx context.Context, outboxPath string, queue *runtime.Queue) DirectoryWatch {
|
||||
ctxChild, fnCancel := context.WithCancel(ctx)
|
||||
|
||||
return DirectoryWatch{
|
||||
ctx: ctxChild,
|
||||
fnCancel: fnCancel,
|
||||
outboxPath: outboxPath,
|
||||
queue: queue,
|
||||
}
|
||||
}
|
||||
|
||||
func (dw DirectoryWatch) Done() <-chan struct{} {
|
||||
return dw.ctx.Done()
|
||||
}
|
||||
|
||||
// Watch reads the `outbox` directory every `TickInterval` and put JSON format e-mail in the queue.
|
||||
func (dw DirectoryWatch) Watch() {
|
||||
log.Info().Str("outbox", dw.outboxPath).Msg("watching outbox directory...")
|
||||
|
||||
ticker := time.NewTicker(TickerInterval)
|
||||
|
||||
go func() {
|
||||
for {
|
||||
select {
|
||||
case <-dw.Done():
|
||||
log.Err(dw.ctx.Err()).Msg("context done")
|
||||
return
|
||||
case <-ticker.C:
|
||||
log.Debug().Str("action", "retrieving json e-mail format...").Str("path", dw.outboxPath)
|
||||
|
||||
files, err := os.ReadDir(dw.outboxPath)
|
||||
if err != nil && !os.IsExist(err) {
|
||||
log.Err(err).Msg("outbox directory does not exist")
|
||||
dw.fnCancel()
|
||||
return
|
||||
}
|
||||
|
||||
for _, file := range files {
|
||||
filename := file.Name()
|
||||
if !strings.HasSuffix(filename, JSONSuffix) {
|
||||
continue
|
||||
}
|
||||
|
||||
emailPath := path.Join(dw.outboxPath, filename)
|
||||
email, err := mail.FromJSON(emailPath)
|
||||
|
||||
if err != nil {
|
||||
log.Err(err).Str("path", emailPath).Msg("unable to parse JSON email")
|
||||
|
||||
// if JSON parsing failed the `path` is renamed with an error suffix to not watch it again
|
||||
newPath := fmt.Sprintf("%s%s", emailPath, ErrorSuffix)
|
||||
if err := os.Rename(emailPath, newPath); err != nil {
|
||||
log.Err(err).Str("path", emailPath).Str("new path", newPath).Msg("unable to rename bad JSON email path")
|
||||
}
|
||||
|
||||
continue
|
||||
}
|
||||
|
||||
email.Path = emailPath
|
||||
dw.queue.Add(email)
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
}
|
||||
Reference in New Issue
Block a user