tsk-6: Check queue (#12)
Closes #6 Reviewed-on: #12 Co-authored-by: phoenix <kundeng00@pm.me> Co-committed-by: phoenix <kundeng00@pm.me>
This commit is contained in:
@@ -0,0 +1,55 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
|
||||
"git.kundeng.us/phoenix/textsender-models/tx0/token"
|
||||
|
||||
"git.kundeng.us/phoenix/catapult/internal/app"
|
||||
)
|
||||
|
||||
type Auth struct {
|
||||
Application *app.App
|
||||
}
|
||||
|
||||
type service struct {
|
||||
Username string `json:"username"`
|
||||
Passphrase string `json:"passphrase"`
|
||||
}
|
||||
|
||||
type tokenResponse struct {
|
||||
Message string `json:"message"`
|
||||
Data []*token.Login `json:"data"`
|
||||
}
|
||||
|
||||
func (a *Auth) GetToken() (*token.Login, error) {
|
||||
serv := service{
|
||||
Username: a.Application.ServiceUsername,
|
||||
Passphrase: a.Application.ServicePassphrase,
|
||||
}
|
||||
|
||||
jsonData, err := json.Marshal(serv)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
resp, err := http.Post(
|
||||
fmt.Sprintf("%s/api/v1/service/login", a.Application.AuthUrl),
|
||||
"application/json",
|
||||
bytes.NewBuffer(jsonData),
|
||||
)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
var r tokenResponse
|
||||
err = json.NewDecoder(resp.Body).Decode(&r)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return r.Data[0], nil
|
||||
}
|
||||
@@ -0,0 +1,121 @@
|
||||
package core
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"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
|
||||
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
|
||||
}
|
||||
}
|
||||
|
||||
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("Id:", item.Id)
|
||||
} 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")
|
||||
}
|
||||
@@ -0,0 +1,54 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
|
||||
"git.kundeng.us/phoenix/textsender-models/tx0/message/scheduling"
|
||||
"git.kundeng.us/phoenix/textsender-models/tx0/token"
|
||||
|
||||
"git.kundeng.us/phoenix/catapult/internal/app"
|
||||
)
|
||||
|
||||
type Queue struct {
|
||||
Application *app.App
|
||||
Token *token.Login
|
||||
}
|
||||
|
||||
type queueResponse struct {
|
||||
Message string `json:"message"`
|
||||
Data []*scheduling.ScheduledMessage `json:"data"`
|
||||
}
|
||||
|
||||
func (q *Queue) GetQueue() (*scheduling.ScheduledMessage, *bool, error) {
|
||||
req, err := http.NewRequest("GET", fmt.Sprintf("%s/api/v1/schedule/message/fetch", q.Application.ApiUrl), nil)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
token := q.Token.AccessToken
|
||||
req.Header.Set("Authorization", "Bearer "+token)
|
||||
|
||||
client := &http.Client{}
|
||||
resp, err := client.Do(req)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
|
||||
var r queueResponse
|
||||
err = json.NewDecoder(resp.Body).Decode(&r)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
} else {
|
||||
var exists bool
|
||||
if len(r.Data) == 0 {
|
||||
exists = false
|
||||
return nil, &exists, nil
|
||||
} else {
|
||||
exists = true
|
||||
return r.Data[0], &exists, nil
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user