From 6fa19ea4f611142bdf86690aa343c080890379b9 Mon Sep 17 00:00:00 2001 From: phoenix Date: Fri, 28 Nov 2025 18:01:18 +0000 Subject: [PATCH] tsk-6: Check queue (#12) Closes #6 Reviewed-on: https://git.kundeng.us/phoenix/catapult/pulls/12 Co-authored-by: phoenix Co-committed-by: phoenix --- cmd/catapult/main.go | 3 +- internal/{auth => service}/auth.go | 6 +-- .../config.go => service/core/service.go} | 17 ++++-- internal/service/queue.go | 54 +++++++++++++++++++ 4 files changed, 73 insertions(+), 7 deletions(-) rename internal/{auth => service}/auth.go (93%) rename internal/{config/config.go => service/core/service.go} (83%) create mode 100644 internal/service/queue.go diff --git a/cmd/catapult/main.go b/cmd/catapult/main.go index f38f94c..f784bf8 100644 --- a/cmd/catapult/main.go +++ b/cmd/catapult/main.go @@ -8,6 +8,7 @@ import ( "git.kundeng.us/phoenix/catapult/internal/app" "git.kundeng.us/phoenix/catapult/internal/config" + "git.kundeng.us/phoenix/catapult/internal/service/core" "git.kundeng.us/phoenix/catapult/internal/version" ) @@ -26,7 +27,7 @@ func main() { return } - service := config.NewService(myApp) + service := core.NewService(myApp) sigChan := make(chan os.Signal, 1) signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM) diff --git a/internal/auth/auth.go b/internal/service/auth.go similarity index 93% rename from internal/auth/auth.go rename to internal/service/auth.go index 0298994..20d3756 100644 --- a/internal/auth/auth.go +++ b/internal/service/auth.go @@ -1,4 +1,4 @@ -package auth +package service import ( "bytes" @@ -20,7 +20,7 @@ type service struct { Passphrase string `json:"passphrase"` } -type response struct { +type tokenResponse struct { Message string `json:"message"` Data []*token.Login `json:"data"` } @@ -45,7 +45,7 @@ func (a *Auth) GetToken() (*token.Login, error) { } defer resp.Body.Close() - var r response + var r tokenResponse err = json.NewDecoder(resp.Body).Decode(&r) if err != nil { return nil, err diff --git a/internal/config/config.go b/internal/service/core/service.go similarity index 83% rename from internal/config/config.go rename to internal/service/core/service.go index e680084..64262f9 100644 --- a/internal/config/config.go +++ b/internal/service/core/service.go @@ -1,4 +1,4 @@ -package config +package core import ( "context" @@ -10,7 +10,7 @@ import ( "git.kundeng.us/phoenix/textsender-models/tx0/token" "git.kundeng.us/phoenix/catapult/internal/app" - "git.kundeng.us/phoenix/catapult/internal/auth" + "git.kundeng.us/phoenix/catapult/internal/service" ) type Service struct { @@ -55,7 +55,7 @@ func (s *Service) worker1() { // Do some work log.Println("Worker 1: Processing...") if s.token == nil { - catapultAuth := auth.Auth{Application: s.App} + catapultAuth := service.Auth{Application: s.App} if token, err := catapultAuth.GetToken(); err != nil { fmt.Println("Error:", err) } else { @@ -63,6 +63,17 @@ func (s *Service) worker1() { 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") + } + } } } } diff --git a/internal/service/queue.go b/internal/service/queue.go new file mode 100644 index 0000000..13f60ea --- /dev/null +++ b/internal/service/queue.go @@ -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 + } + } +}