newsbox/workers/workers.go
2023-02-15 11:51:54 -08:00

127 lines
2.8 KiB
Go

package workers
import (
"context"
"encoding/json"
"log"
"newsbox/internal"
"newsbox/mailers"
"newsbox/models"
"time"
"github.com/acaloiaro/neoq"
)
var nq neoq.Neoq
func Start() {
log.Println("Firing up workers")
dbURL := internal.DatabaseUrl()
nq, _ = neoq.New(dbURL)
nq.Listen("incoming_email", neoq.NewHandler(func(ctx context.Context) (err error) {
j, err := neoq.JobFromContext(ctx)
if err != nil {
return
}
var msg models.Message
err = models.DB.Eager().Find(&msg, j.Payload["message_id"])
if err != nil {
return
}
var owner models.User
err = models.DB.Eager().Find(&owner, msg.OwnerID)
if err != nil {
return
}
// Re-address the email to the user's email address
msg.To = owner.Email.Interface().(string)
// Determine time to deliver
tzOffset := owner.Preferences.TimeZoneUtcOffset
t := internal.TimeInZone(tzOffset)
et := time.Date(t.Year(), t.Month(), t.Day(), owner.Preferences.EmailHourOfDay, 0, 0, 0, t.Location())
diff := et.Sub(t)
// The message should be delivered later today
if diff > 0 {
et = t.Add(diff)
} else { // The message should be delivered tomorrow
et = time.Date(t.Year(), t.Month(), t.Day()+1, owner.Preferences.EmailHourOfDay, 0, 0, 0, t.Location())
}
msgJSON, _ := json.Marshal(msg)
log.Println("Queueing for:", et)
log.Println("Assumed location:", t.Location())
nq, _ := neoq.New(dbURL)
_, err = nq.Enqueue(neoq.Job{
Queue: "outgoing_email",
Payload: map[string]any{
"message": string(msgJSON),
},
})
return
}))
nq.Listen("welcome_emails", neoq.NewHandler(func(ctx context.Context) (err error) {
j, err := neoq.JobFromContext(ctx)
if err != nil {
return
}
recipient := j.Payload["recipient"].(string)
verificationURL := j.Payload["verification_url"].(string)
log.Println("sending welcome email to:", recipient)
err = mailers.SendWelcomeEmail(recipient, verificationURL)
log.Println("What's my error here?", err)
return
}))
nq.Listen("outgoing_email", neoq.NewHandler(func(ctx context.Context) (err error) {
var msg models.Message
j, err := neoq.JobFromContext(ctx)
if err != nil {
return
}
err = json.Unmarshal([]byte(j.Payload["message"].(string)), &msg)
if err != nil {
return
}
var owner models.User
err = models.DB.Find(&owner, msg.OwnerID)
if err != nil {
return
}
err = mailers.ForwardMessage(owner.Email.Interface().(string), &msg)
return
}))
nq.Listen("email_verifications", neoq.NewHandler(func(ctx context.Context) (err error) {
j, err := neoq.JobFromContext(ctx)
if err != nil {
return
}
recipient := j.Payload["recipient"].(string)
verificationURL := j.Payload["verification_url"].(string)
err = mailers.SendEmailAddressChangedEmail(recipient, verificationURL)
return
}))
}