Closes #14 Reviewed-on: #16 Co-authored-by: phoenix <kundeng00@pm.me> Co-committed-by: phoenix <kundeng00@pm.me>
201 lines
5.2 KiB
Go
201 lines
5.2 KiB
Go
package core
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log"
|
|
"sync"
|
|
"time"
|
|
|
|
"git.kundeng.us/phoenix/textsender-models/tx0/message/scheduling"
|
|
"git.kundeng.us/phoenix/textsender-models/tx0/token"
|
|
|
|
"git.kundeng.us/phoenix/catapult/internal/app"
|
|
"git.kundeng.us/phoenix/catapult/internal/service"
|
|
)
|
|
|
|
type Service struct {
|
|
wg sync.WaitGroup
|
|
ctx context.Context
|
|
cancel context.CancelFunc
|
|
App *app.App
|
|
token *token.Login
|
|
}
|
|
|
|
func NewService(application *app.App) *Service {
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
return &Service{
|
|
ctx: ctx,
|
|
cancel: cancel,
|
|
App: application,
|
|
}
|
|
}
|
|
|
|
func (s *Service) Start() {
|
|
log.Println("Starting service...")
|
|
|
|
// Start multiple background workers
|
|
s.wg.Add(3)
|
|
go s.worker1()
|
|
go s.worker2()
|
|
go s.healthChecker()
|
|
}
|
|
|
|
func (s *Service) worker1() {
|
|
defer s.wg.Done()
|
|
|
|
ticker := time.NewTicker(5 * time.Second)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-s.ctx.Done():
|
|
log.Println("Worker 1 shutting down...")
|
|
return
|
|
case <-ticker.C:
|
|
// Do some work
|
|
now := time.Now()
|
|
log.Println("Worker 1: Processing...")
|
|
if s.token == nil {
|
|
catapultAuth := service.Auth{Application: s.App}
|
|
if token, err := catapultAuth.GetToken(); err != nil {
|
|
fmt.Println("Error:", err)
|
|
} else {
|
|
fmt.Println("Access token:", token.AccessToken)
|
|
s.token = token
|
|
}
|
|
}
|
|
|
|
log.Println("Twilio config account sid:", s.App.TwilioConfig.AccountSID)
|
|
log.Println("Twilio config auth token:", s.App.TwilioConfig.AuthToken)
|
|
log.Println("Twilio config service sid:", s.App.TwilioConfig.ServiceSID)
|
|
log.Println("Twilio config number:", s.App.TwilioConfig.Number)
|
|
|
|
queue := service.Queue{Application: s.App, Token: s.token}
|
|
if item, exists, err := queue.GetQueue(); err != nil {
|
|
fmt.Println("Error:", err)
|
|
} else {
|
|
if *exists {
|
|
fmt.Println("Scheduled message Id:", item.Id)
|
|
fmt.Println("Created:", item.Created)
|
|
fmt.Println("Scheduled:", item.Scheduled)
|
|
fmt.Println("Status:", item.Status)
|
|
log.Println("User Id:", item.UserId)
|
|
if scheduledMessageValid(item, now) {
|
|
fmt.Println("Scheduled Message can be sent")
|
|
scheduler := service.Scheduler{Application: s.App, Token: s.token}
|
|
if events, err := scheduler.GetEvents(item.Id); err != nil {
|
|
fmt.Println("Error getting event:", err)
|
|
} else {
|
|
for _, event := range events {
|
|
if msg, err := scheduler.GetMessage(event.MessageId); err != nil {
|
|
fmt.Println("Error getting message:", err)
|
|
} else {
|
|
if c, err := scheduler.GetContact(event.RecipientId); err != nil {
|
|
fmt.Println("Error getting contact:", err)
|
|
} else {
|
|
fmt.Println("Message Id:", msg.Id)
|
|
fmt.Println("Contact Id:", c.Id)
|
|
if res, sent, err := scheduler.ScheduleMessage(*item, *event, *c, *msg); err != nil {
|
|
log.Println("Failure with scheduling the message:", err)
|
|
} else {
|
|
if res != nil {
|
|
bytes, err := json.Marshal(res)
|
|
if err != nil {
|
|
log.Println("Error parsing result:", err)
|
|
} else {
|
|
jsonValue := string(bytes)
|
|
log.Println("Result:", jsonValue)
|
|
|
|
if eventResponse, err := scheduler.RecordEventResponse(event, item.UserId, bytes, sent); err != nil {
|
|
log.Println("Failure recording event response:", err)
|
|
} else {
|
|
if len(eventResponse) == 0 {
|
|
log.Println("No event responses")
|
|
} else {
|
|
eventResponse := eventResponse[0]
|
|
log.Println("Event response saved. Id:", eventResponse.Id)
|
|
sameJsonValue := string(eventResponse.Response)
|
|
|
|
if jsonValue == sameJsonValue {
|
|
log.Println("Matches")
|
|
} else if len(jsonValue) == len(sameJsonValue) {
|
|
log.Println("Same size")
|
|
log.Println(jsonValue)
|
|
log.Println(sameJsonValue)
|
|
} else {
|
|
log.Println("Not the same")
|
|
}
|
|
}
|
|
}
|
|
}
|
|
} else {
|
|
log.Println("Result should not be empty")
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
} else {
|
|
fmt.Println("Invalid scheduled message")
|
|
}
|
|
} else {
|
|
fmt.Println("Empty queue")
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Service) worker2() {
|
|
defer s.wg.Done()
|
|
|
|
ticker := time.NewTicker(10 * time.Second)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-s.ctx.Done():
|
|
log.Println("Worker 2 shutting down...")
|
|
return
|
|
case <-ticker.C:
|
|
// Do some work
|
|
log.Println("Worker 2: Processing...")
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Service) healthChecker() {
|
|
defer s.wg.Done()
|
|
|
|
ticker := time.NewTicker(30 * time.Second)
|
|
defer ticker.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-s.ctx.Done():
|
|
log.Println("Health checker shutting down...")
|
|
return
|
|
case <-ticker.C:
|
|
log.Println("Health check: Service is healthy")
|
|
}
|
|
}
|
|
}
|
|
|
|
func (s *Service) Stop() {
|
|
log.Println("Shutting down service...")
|
|
s.cancel()
|
|
s.wg.Wait()
|
|
log.Println("Service stopped gracefully")
|
|
}
|
|
|
|
func scheduledMessageValid(schMsg *scheduling.ScheduledMessage, now time.Time) bool {
|
|
if schMsg.Scheduled.Before(now) {
|
|
return false
|
|
} else {
|
|
return true
|
|
}
|
|
}
|