bump to go 1.20 + fix some code issues
This commit is contained in:
+43
-37
@@ -51,14 +51,14 @@ func (s Sender) SendMail(email mail.Email) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// watchOutbox reads the `outbox` directory every `TickInterval` and put JSON format e-mail in the queue
|
||||
// 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 {
|
||||
for range ticker.C {
|
||||
log.Debug().Str("action", "retrieving json e-mail format...").Str("path", s.outboxPath)
|
||||
|
||||
files, err := os.ReadDir(s.outboxPath)
|
||||
@@ -69,18 +69,34 @@ func (s Sender) watchOutbox() {
|
||||
|
||||
for _, file := range files {
|
||||
filename := file.Name()
|
||||
if strings.HasSuffix(filename, JSONSuffix) {
|
||||
s.queue.Add(path.Join(s.outboxPath, filename))
|
||||
if !strings.HasSuffix(filename, JSONSuffix) {
|
||||
log.Debug().Str("filename", filename).Msg("incorrect suffix")
|
||||
continue
|
||||
}
|
||||
|
||||
log.Debug().Str("filename", filename).Msg("incorrect suffix")
|
||||
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 loops over the queue and send email
|
||||
// processNextEmail iterates over the queue and send email.
|
||||
func (s Sender) processNextEmail() bool {
|
||||
item, quit := s.queue.Get()
|
||||
if quit {
|
||||
@@ -88,66 +104,56 @@ func (s Sender) processNextEmail() bool {
|
||||
}
|
||||
defer s.queue.Done(item)
|
||||
|
||||
path, ok := item.(string)
|
||||
email, ok := item.(mail.Email)
|
||||
if !ok {
|
||||
log.Error().Any("item", item).Msg("unable to cast queue item into mail.Email")
|
||||
return true
|
||||
}
|
||||
|
||||
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 avoid enqueued 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")
|
||||
s.queue.Shutdown()
|
||||
}
|
||||
return true
|
||||
if err := s.SendMail(email); err != nil {
|
||||
log.Err(err).Msg("unable to send the email")
|
||||
}
|
||||
|
||||
// whatever the return, the email will be not enqueued again
|
||||
s.SendMail(email)
|
||||
|
||||
if err := os.Remove(path); err != nil {
|
||||
// this is a fatal error, can't send same e-mail indefinitely
|
||||
if !os.IsExist(err) {
|
||||
log.Err(err).Str("path", path).Msg("unable to remove the JSON email")
|
||||
s.queue.Shutdown()
|
||||
if path := email.Path; path != "" {
|
||||
if err := os.Remove(path); err != nil {
|
||||
// this is a fatal error, can't send same e-mail indefinitely
|
||||
if !os.IsExist(err) {
|
||||
log.Err(err).Str("path", path).Msg("unable to remove the JSON email")
|
||||
s.queue.Shutdown()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
// run starts processing the queue
|
||||
// run starts processing the queue.
|
||||
func (s Sender) run() <-chan struct{} {
|
||||
queueCh := make(chan struct{})
|
||||
chQueue := make(chan struct{})
|
||||
go func() {
|
||||
for s.processNextEmail() {
|
||||
}
|
||||
queueCh <- struct{}{}
|
||||
chQueue <- struct{}{}
|
||||
}()
|
||||
return queueCh
|
||||
return chQueue
|
||||
}
|
||||
|
||||
// Run launches the queue processing and the outbox watcher
|
||||
// catches `SIGINT` and `SIGTERM` to properly stopped the queue
|
||||
// Run launches the queue processing and the outbox watcher.
|
||||
// It catches `SIGINT` and `SIGTERM` to properly stopped the queue.
|
||||
func (s Sender) Run() {
|
||||
log.Info().Msg("sender service is running")
|
||||
|
||||
sigCh := make(chan os.Signal, 1)
|
||||
signal.Notify(sigCh, os.Interrupt, syscall.SIGTERM)
|
||||
chSignal := make(chan os.Signal, 1)
|
||||
signal.Notify(chSignal, os.Interrupt, syscall.SIGTERM)
|
||||
|
||||
s.watchOutbox()
|
||||
queueCh := s.run()
|
||||
chQueue := s.run()
|
||||
|
||||
select {
|
||||
case <-sigCh:
|
||||
case <-chSignal:
|
||||
log.Warn().Msg("stop signal received, stopping e-mail queue...")
|
||||
s.queue.Shutdown()
|
||||
case <-queueCh:
|
||||
case <-chQueue:
|
||||
log.Info().Msg("e-mail queue stopped successfully")
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user