5 Commits
Author SHA1 Message Date
phoenixandphoenix c809815549 tsk-14: Record when message was sent (#16)
Closes #14

Reviewed-on: #16
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-12-01 21:24:05 +00:00
phoenixandphoenix 04da7785c8 tsk-8: Schedule message (#15)
Closes #8

Reviewed-on: #15
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-29 19:28:05 +00:00
phoenixandphoenix 9af11dc961 tsk-7: Verify if scheduled message can be sent (#13)
Closes #7

Reviewed-on: #13
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-28 18:23:47 +00:00
phoenixandphoenix 6fa19ea4f6 tsk-6: Check queue (#12)
Closes #6

Reviewed-on: #12
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-28 18:01:18 +00:00
phoenixandphoenix cef7f8a150 tsk-5: Fetch token (#9)
Closes #5

Reviewed-on: #9
Co-authored-by: phoenix <kundeng00@pm.me>
Co-committed-by: phoenix <kundeng00@pm.me>
2025-11-28 16:54:11 +00:00
16 changed files with 798 additions and 104 deletions
+8
View File
@@ -0,0 +1,8 @@
AUTH_URL=https://auth.txt.com
API_URL=https://txt.com
SERVICE_USERNAME=suave
SERVICE_PASSPHRASE=9238urc9328nr329
TWILIO_AUTH_SID=9M438C93R943U4329MCU43C34U
TWILIO_SERVICE_SID=9M4J3X8439U398NUVT3342MC349C348T
TWILIO_AUTH_TOKEN=OJ8mU98U8UUU098u0U08kd9IDKaeoijtritjerDFISFGOS
TWILIO_PHONE_NUMBER=10123456789
+4
View File
@@ -1 +1,5 @@
catapult
/vendor
.env
.env.local
+10 -2
View File
@@ -2,16 +2,19 @@ package main
import (
"fmt"
"log"
"os"
"os/signal"
"syscall"
"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"
)
func main() {
fmt.Println(config.App_Name)
log.Println(app.App_Name)
versionFlag := config.CheckVersionFlag()
if *versionFlag {
@@ -19,8 +22,13 @@ func main() {
return
}
service := config.NewService()
myApp, err := app.Load()
if err != nil {
log.Println("Error loading app: %w", err)
return
}
service := core.NewService(myApp)
sigChan := make(chan os.Signal, 1)
signal.Notify(sigChan, syscall.SIGINT, syscall.SIGTERM)
+14
View File
@@ -1,3 +1,17 @@
module git.kundeng.us/phoenix/catapult
go 1.25.4
require (
git.kundeng.us/phoenix/swoosh v0.0.6
git.kundeng.us/phoenix/textsender-models v0.0.11
github.com/google/uuid v1.6.0
github.com/joho/godotenv v1.5.1
)
require (
github.com/golang-jwt/jwt/v5 v5.3.0 // indirect
github.com/golang/mock v1.6.0 // indirect
github.com/pkg/errors v0.9.1 // indirect
github.com/twilio/twilio-go v1.28.7 // indirect
)
+60
View File
@@ -0,0 +1,60 @@
git.kundeng.us/phoenix/swoosh v0.0.6 h1:9+BFsmuxufEPJErMJUioyRGb0KoIPdwW+PSbEva5UuU=
git.kundeng.us/phoenix/swoosh v0.0.6/go.mod h1:czsMdVpt7HKkaP1HOl5+j9ePOQCA0r9UW4uAxwAbTPk=
git.kundeng.us/phoenix/textsender-models v0.0.11 h1:kd2FdeZJhJJAXBm8MoyadtgNGyzC+puU1oR8B8N+MfE=
git.kundeng.us/phoenix/textsender-models v0.0.11/go.mod h1:9iPDQJg1Tc6WMNoW5+f8YKmnosMwlWHJ++hmxNLDEe0=
github.com/beevik/etree v1.1.0/go.mod h1:r8Aw8JqVegEf0w2fDnATrX9VpkMcyFeM0FhwO62wh+A=
github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E=
github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c=
github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38=
github.com/golang-jwt/jwt/v5 v5.2.2/go.mod h1:pqrtFR0X4osieyHYxtmOUWsAWrfe1Q5UVIyoH402zdk=
github.com/golang-jwt/jwt/v5 v5.3.0 h1:pv4AsKCKKZuqlgs5sUmn4x8UlGa0kEVt/puTpKx9vvo=
github.com/golang-jwt/jwt/v5 v5.3.0/go.mod h1:fxCRLWMO43lRc8nhHWY6LGqRcf+1gQWArsqaEUEa5bE=
github.com/golang/mock v1.6.0 h1:ErTB+efbowRARo13NNdxyJji2egdxLGQhRaY+DUumQc=
github.com/golang/mock v1.6.0/go.mod h1:p6yTPP+5HYm5mzsMV8JkE6ZKdX+/wYM6Hr+LicevLPs=
github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0=
github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0=
github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4=
github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ=
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
github.com/localtunnel/go-localtunnel v0.0.0-20170326223115-8a804488f275 h1:IZycmTpoUtQK3PD60UYBwjaCUHUP7cML494ao9/O8+Q=
github.com/localtunnel/go-localtunnel v0.0.0-20170326223115-8a804488f275/go.mod h1:zt6UU74K6Z6oMOYJbJzYpYucqdcQwSMPBEdSvGiaUMw=
github.com/niemeyer/pretty v0.0.0-20200227124842-a10e7caefd8e/go.mod h1:zD1mROLANZcx1PVRCS0qkT7pwLkGfwJo4zjcN/Tysno=
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM=
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/testify v1.7.0 h1:nwc3DEeHmmLAfoZucVR881uASk0Mfjw8xYJ99tb5CcY=
github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
github.com/twilio/twilio-go v1.28.7 h1:WzzQDR/rqmNkVs1TwtcHFPYGTdSdPnx/eAZf5UIXzr4=
github.com/twilio/twilio-go v1.28.7/go.mod h1:FpgNWMoD8CFnmukpKq9RNpUSGXC0BwnbeKZj2YHlIkw=
github.com/yuin/goldmark v1.3.5/go.mod h1:mwnBkeHKe2W/ZEtQ+71ViKU8L12m81fl3OWwC1Zlc8k=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/crypto v0.0.0-20191011191535-87dc89f01550/go.mod h1:yigFU9vqHzYiE8UmvKecakEJjdnWj3jj499lnFckfCI=
golang.org/x/mod v0.4.2/go.mod h1:s0Qsj1ACt9ePp/hMypM3fl4fZqREWJwdYDEqhRiZZUA=
golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
golang.org/x/net v0.0.0-20210405180319-a5a99cb37ef4/go.mod h1:p54w0d4576C0XHj96bSt6lcn1PtDYWL6XObtHCRCNQM=
golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM=
golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY=
golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs=
golang.org/x/sys v0.0.0-20210510120138-977fb7262007/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ=
golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/tools v0.1.1/go.mod h1:o0xws9oXOQQZyjljx8fwUC0k7L1pTE6eaCbjGeHmOkk=
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
golang.org/x/xerrors v0.0.0-20191011141410-1b5146add898/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
golang.org/x/xerrors v0.0.0-20200804184101-5ec99f83aff1/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c h1:dUUwHk2QECo/6vqA44rthZ8ie2QXMNeKRTHCNY2nXvo=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
+81
View File
@@ -0,0 +1,81 @@
package app
import (
"fmt"
"os"
"path"
"git.kundeng.us/phoenix/textsender-models/tx0/config"
"github.com/joho/godotenv"
)
const App_Name = "catapult"
type App struct {
ApiUrl string
AuthUrl string
ServiceUsername string
ServicePassphrase string
TwilioConfig *config.TwiloConfig
}
func Load() (*App, error) {
err := godotenv.Load()
if err != nil {
cwd, _ := os.Getwd()
envPath := path.Join(cwd, "../..", ".env")
if err = godotenv.Load(envPath); err != nil {
prevPath := path.Join(envPath, "../..", ".env")
if err = godotenv.Load(prevPath); err != nil {
return nil, fmt.Errorf("Error loading .env file: %w", err)
}
}
}
apiUrl := os.Getenv("API_URL")
authUrl := os.Getenv("AUTH_URL")
serviceUsername := os.Getenv("SERVICE_USERNAME")
servicePassphrase := os.Getenv("SERVICE_PASSPHRASE")
if len(apiUrl) == 0 {
return nil, fmt.Errorf("Api Url not provided")
} else if len(authUrl) == 0 {
return nil, fmt.Errorf("Auth url not provided")
} else if len(serviceUsername) == 0 {
return nil, fmt.Errorf("Service username not provided")
} else if len(servicePassphrase) == 0 {
return nil, fmt.Errorf("Service passphrase not provided")
} else {
if cfg, err := loadTwilioConfig(); err != nil {
return nil, err
} else {
return &App{
ApiUrl: apiUrl, AuthUrl: authUrl, ServiceUsername: serviceUsername, ServicePassphrase: servicePassphrase, TwilioConfig: cfg,
}, nil
}
}
}
func loadTwilioConfig() (*config.TwiloConfig, error) {
authSid := os.Getenv("TWILIO_AUTH_SID")
serviceSid := os.Getenv("TWILIO_SERVICE_SID")
authToken := os.Getenv("TWILIO_AUTH_TOKEN")
phoneNumber := os.Getenv("TWILIO_PHONE_NUMBER")
if len(authSid) == 0 {
return nil, fmt.Errorf("Twilio config auth sid not provided")
} else if len(serviceSid) == 0 {
return nil, fmt.Errorf("Twilio config service sid not provided")
} else if len(authToken) == 0 {
return nil, fmt.Errorf("Twilio config token not provided")
} else if len(phoneNumber) == 0 {
return nil, fmt.Errorf("Twilio config phone number not provided")
} else {
cfg := config.TwiloConfig{}
cfg.AccountSID = authSid
cfg.AuthToken = authToken
cfg.ServiceSID = serviceSid
cfg.Number = phoneNumber
return &cfg, nil
}
}
-102
View File
@@ -1,102 +0,0 @@
package config
import (
"context"
"flag"
"log"
"sync"
"time"
)
const App_Name = "catapult"
func CheckVersionFlag() *bool {
versionFlag := flag.Bool("version", false, "Print version information")
flag.Parse()
return versionFlag
}
type Service struct {
wg sync.WaitGroup
ctx context.Context
cancel context.CancelFunc
}
func NewService() *Service {
ctx, cancel := context.WithCancel(context.Background())
return &Service{
ctx: ctx,
cancel: cancel,
}
}
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...")
}
}
}
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")
}
+12
View File
@@ -0,0 +1,12 @@
package config
import (
"flag"
)
func CheckVersionFlag() *bool {
versionFlag := flag.Bool("version", false, "Print version information")
flag.Parse()
return versionFlag
}
+55
View File
@@ -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
}
+58
View File
@@ -0,0 +1,58 @@
package service
import (
"encoding/json"
"fmt"
"net/http"
"net/url"
"git.kundeng.us/phoenix/textsender-models/tx0/contact"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"github.com/google/uuid"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type Contact struct {
Application *app.App
Token *token.Login
}
type getContactResponse struct {
Message string `json:"message"`
Data []*contact.Contact `json:"data"`
}
func (c *Contact) GetContact(id uuid.UUID) (*contact.Contact, error) {
params := url.Values{}
params.Add("id", id.String())
pm := params.Encode()
fullUrl := fmt.Sprintf("%s/api/v1/contact?%s", c.Application.ApiUrl, pm)
fmt.Println("Url:", fullUrl)
req, err := http.NewRequest("GET", fullUrl, nil)
if err != nil {
return nil, err
}
token := c.Token.AccessToken
req.Header.Set("Authorization", "Bearer "+token)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r getContactResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
} else {
if len(r.Data) == 0 {
return nil, nil
} else {
return r.Data[0], nil
}
}
}
+200
View File
@@ -0,0 +1,200 @@
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
}
}
+58
View File
@@ -0,0 +1,58 @@
package service
import (
"encoding/json"
"fmt"
"net/http"
"net/url"
"git.kundeng.us/phoenix/textsender-models/tx0/message"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"github.com/google/uuid"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type Message struct {
Application *app.App
Token *token.Login
}
type getMessageResponse struct {
Message string `json:"message"`
Data []*message.Message `json:"data"`
}
func (m *Message) GetMessage(id uuid.UUID) (*message.Message, error) {
params := url.Values{}
params.Add("id", id.String())
pm := params.Encode()
fullUrl := fmt.Sprintf("%s/api/v1/message?%s", m.Application.ApiUrl, pm)
fmt.Println("Url:", fullUrl)
req, err := http.NewRequest("GET", fullUrl, nil)
if err != nil {
return nil, err
}
token := m.Token.AccessToken
req.Header.Set("Authorization", "Bearer "+token)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r getMessageResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
} else {
if len(r.Data) == 0 {
return nil, nil
} else {
return r.Data[0], nil
}
}
}
@@ -0,0 +1,58 @@
package service
import (
"encoding/json"
"fmt"
"net/http"
"strings"
"git.kundeng.us/phoenix/textsender-models/tx0/message"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type MessageEventResponse struct {
Application *app.App
Token *token.Login
}
type recordMessageEventResponse struct {
Message string `json:"message"`
Data []*message.MessageEventResponse `json:"data"`
}
func (m *MessageEventResponse) RecordEventResponse(mer *message.MessageEventResponse) ([]*message.MessageEventResponse, error) {
fullUrl := fmt.Sprintf("%s/api/v1/schedule/message/event/response/record", m.Application.ApiUrl)
jsonData, err := json.Marshal(mer)
if err != nil {
return nil, err
}
req, err := http.NewRequest("POST", fullUrl, strings.NewReader(string(jsonData)))
if err != nil {
return nil, err
}
token := m.Token.AccessToken
req.Header.Set("Authorization", "Bearer "+token)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r recordMessageEventResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
} else {
if len(r.Data) == 0 {
return nil, nil
} else {
return r.Data, nil
}
}
}
+54
View File
@@ -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
}
}
}
@@ -0,0 +1,58 @@
package service
import (
"encoding/json"
"fmt"
"net/http"
"net/url"
"git.kundeng.us/phoenix/textsender-models/tx0/message/scheduling"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"github.com/google/uuid"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type ScheduledMessageEvent struct {
Application *app.App
Token *token.Login
}
type getScheduledMessageEventResponse struct {
Message string `json:"message"`
Data []*scheduling.ScheduledMessageEvent `json:"data"`
}
func (s *ScheduledMessageEvent) GetEvents(scheduledMessageId uuid.UUID) ([]*scheduling.ScheduledMessageEvent, error) {
params := url.Values{}
params.Add("scheduled_message_id", scheduledMessageId.String())
pm := params.Encode()
fullUrl := fmt.Sprintf("%s/api/v1/schedule/message/event?%s", s.Application.ApiUrl, pm)
req, err := http.NewRequest("GET", fullUrl, nil)
if err != nil {
return nil, err
}
token := s.Token.AccessToken
req.Header.Set("Authorization", "Bearer "+token)
client := &http.Client{}
resp, err := client.Do(req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
var r getScheduledMessageEventResponse
err = json.NewDecoder(resp.Body).Decode(&r)
if err != nil {
return nil, err
} else {
if len(r.Data) == 0 {
return nil, nil
} else {
return r.Data, nil
}
}
}
+68
View File
@@ -0,0 +1,68 @@
package service
import (
"time"
"git.kundeng.us/phoenix/swoosh/swoop/send"
"git.kundeng.us/phoenix/swoosh/swoop/types"
"git.kundeng.us/phoenix/textsender-models/tx0/contact"
"git.kundeng.us/phoenix/textsender-models/tx0/message"
"git.kundeng.us/phoenix/textsender-models/tx0/message/scheduling"
"git.kundeng.us/phoenix/textsender-models/tx0/token"
"github.com/google/uuid"
"git.kundeng.us/phoenix/catapult/internal/app"
)
type Scheduler struct {
Application *app.App
Token *token.Login
}
func (s *Scheduler) GetEvents(scheduledMessageId uuid.UUID) ([]*scheduling.ScheduledMessageEvent, error) {
sme := ScheduledMessageEvent{Application: s.Application, Token: s.Token}
if events, err := sme.GetEvents(scheduledMessageId); err != nil {
return nil, err
} else {
return events, nil
}
}
func (s *Scheduler) GetMessage(id uuid.UUID) (*message.Message, error) {
msg := Message{Application: s.Application, Token: s.Token}
if letter, err := msg.GetMessage(id); err != nil {
return nil, err
} else {
return letter, nil
}
}
func (s *Scheduler) GetContact(id uuid.UUID) (*contact.Contact, error) {
ctct := Contact{Application: s.Application, Token: s.Token}
if c, err := ctct.GetContact(id); err != nil {
return nil, err
} else {
return c, nil
}
}
func (s *Scheduler) ScheduleMessage(schMsg scheduling.ScheduledMessage, event scheduling.ScheduledMessageEvent, c contact.Contact, msg message.Message) (*types.TwilioResult, *time.Time, error) {
msgSender := send.MessageSender{Config: s.Application.TwilioConfig}
if res, err := msgSender.Send(msg, c, &schMsg.Scheduled); err != nil {
return nil, nil, err
} else {
sent := time.Now()
return res, &sent, nil
}
}
func (s *Scheduler) RecordEventResponse(schMsgEvent *scheduling.ScheduledMessageEvent, userId uuid.UUID, bytes []byte, sent *time.Time) ([]*message.MessageEventResponse, error) {
mre := message.MessageEventResponse{ScheduledMessageEventId: schMsgEvent.Id, UserId: userId, Response: bytes, Sent: *sent}
msgEventResponse := MessageEventResponse{Application: s.Application, Token: s.Token}
if responses, err := msgEventResponse.RecordEventResponse(&mre); err != nil {
return nil, err
} else {
return responses, nil
}
}