Files
catapult/internal/service/core/service.go
T
phoenixandphoenix 66c7fc261f tsk-11: Obtain refresh token (#18)
Closes #11

Reviewed-on: #18
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-12-04 03:13:13 +00:00

208 lines
5.3 KiB
Go

package core
import (
"context"
"encoding/json"
"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...")
catapultAuth := service.Auth{Application: s.App}
if s.token == nil {
log.Println("Token has not been fetched")
log.Println("Fetching token")
if token, err := catapultAuth.GetToken(); err != nil {
log.Println("Error:", err)
} else {
log.Println("Token fetched")
s.token = token
}
} else if catapultAuth.TokenExpired(s.token) {
// Get refresh token
log.Println("Token expired")
log.Println("Fetching refresh token")
refreshToken, err := catapultAuth.GetRefreshToken(s.token)
if err != nil {
log.Println("Error getting refresh token:", err)
} else {
log.Println("Refresh token fetched")
s.token = refreshToken
}
}
queue := service.Queue{Application: s.App, Token: s.token}
if item, exists, err := queue.GetQueue(); err != nil {
log.Println("Error:", err)
} else {
if *exists {
log.Println("Scheduled message Id:", item.Id)
log.Println("Created:", item.Created)
log.Println("Scheduled:", item.Scheduled)
log.Println("Status:", item.Status)
log.Println("User Id:", item.UserId)
if scheduledMessageValid(item, now) {
log.Println("Scheduled Message can be sent")
scheduler := service.Scheduler{Application: s.App, Token: s.token}
if events, err := scheduler.GetEvents(item.Id); err != nil {
log.Println("Error getting event:", err)
} else {
for _, event := range events {
if msg, err := scheduler.GetMessage(event.MessageId); err != nil {
log.Println("Error getting message:", err)
} else {
if c, err := scheduler.GetContact(event.RecipientId); err != nil {
log.Println("Error getting contact:", err)
} else {
log.Println("Message Id:", msg.Id)
log.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 {
log.Println("Invalid scheduled message")
}
} else {
log.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
}
}