Merged in feature/baseTemplate (pull request #1)

feat!: base template

* first bit of templating

* codeowners

* linuxbased

* restart

* baseline project

* add grpc api and basic integration test

* startqueue

* queueMsg

* splitscripts

* migrations

* queueintegrationtest

* gateway
This commit is contained in:
Michael McGuinness
2024-12-06 14:38:42 +00:00
parent cccbdaad4f
commit b1f8ac453b
2845 changed files with 832056 additions and 50 deletions
+65
View File
@@ -0,0 +1,65 @@
package queue
import (
"context"
"encoding/json"
"gotemplate/internal/name"
"log"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)
type NameController struct {
name name.Service
}
func NewNameController(nameSvc name.Service) *NameController {
return &NameController{
name: nameSvc,
}
}
type CreateNameBody struct {
Name string `json:"name"`
}
func (s *NameController) Create(ctx context.Context, config *QueueConfig, msg *types.Message) error {
var body CreateNameBody
err := json.Unmarshal([]byte(*msg.Body), &body)
if err != nil {
return err
}
err = s.name.Add(ctx, body.Name)
if err != nil {
return err
}
log.Print("Added: ", body.Name)
_, err = config.Client.SendMessage(ctx, &sqs.SendMessageInput{
MessageAttributes: map[string]types.MessageAttributeValue{
"type": {
DataType: aws.String("String"),
StringValue: aws.String("CREATE_NAME_SUCCESS"),
},
},
QueueUrl: aws.String(config.URL),
MessageBody: aws.String("{}"),
})
if err != nil {
return err
}
_, err = config.Client.DeleteMessage(ctx, &sqs.DeleteMessageInput{
QueueUrl: aws.String(config.URL),
ReceiptHandle: msg.ReceiptHandle,
})
if err != nil {
return err
}
return nil
}
+47
View File
@@ -0,0 +1,47 @@
package queue
import (
"context"
"log"
"github.com/aws/aws-sdk-go-v2/service/sqs"
)
type Controllers struct {
Name NameController
}
type Queue struct {
Config *QueueConfig
Controllers *Controllers
}
func PollMessages(ctx context.Context, queue *Queue) {
for {
result, err := queue.Config.Client.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
QueueUrl: &queue.Config.URL,
MaxNumberOfMessages: 1,
WaitTimeSeconds: 2,
VisibilityTimeout: 2,
MessageAttributeNames: []string{
"type",
},
})
if err != nil {
log.Printf("Message Fetch Fail: %v", err)
continue
}
for _, message := range result.Messages {
go func() {
toProcess, err := processMessage(ctx, queue, message)
if !toProcess {
return
}
if err != nil {
log.Printf("Message Process Fail: %v", err)
}
}()
}
}
}
+24
View File
@@ -0,0 +1,24 @@
package queue
import (
"context"
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
)
type QueueConfig struct {
URL string
Client *sqs.Client
}
func processMessage(ctx context.Context, queue *Queue, message types.Message) (bool, error) {
// Process the message here and return false if not message to process
switch *message.MessageAttributes["type"].StringValue {
case "CREATE_NAME":
err := queue.Controllers.Name.Create(ctx, queue.Config, &message)
return true, err
}
return false, nil
}