putting in base query and collector framework
This commit is contained in:
@@ -0,0 +1,53 @@
|
||||
package queue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"gotemplate/internal/document"
|
||||
"gotemplate/internal/queue"
|
||||
|
||||
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
|
||||
)
|
||||
|
||||
type DocumentController struct {
|
||||
document document.Service
|
||||
}
|
||||
|
||||
func NewDocumentController(svc document.Service) *DocumentController {
|
||||
return &DocumentController{
|
||||
document: svc,
|
||||
}
|
||||
}
|
||||
|
||||
type DocumentQueryEvent struct {
|
||||
ID string `json:"id"`
|
||||
}
|
||||
|
||||
func (s *DocumentController) Sync(ctx context.Context, config *queue.QueueConfig, msg *types.Message) error {
|
||||
var body document.Document
|
||||
err := json.Unmarshal([]byte(*msg.Body), &body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
err = s.document.Sync(ctx, body)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
queryEvent := DocumentQueryEvent{
|
||||
ID: body.ID,
|
||||
}
|
||||
|
||||
err = queue.Send(ctx, config, "DOCQUERY", queryEvent)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
err = queue.Delete(ctx, config, msg)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -1,65 +0,0 @@
|
||||
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
|
||||
}
|
||||
+48
-47
@@ -1,47 +1,48 @@
|
||||
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)
|
||||
}
|
||||
}()
|
||||
}
|
||||
}
|
||||
}
|
||||
package queue
|
||||
|
||||
import (
|
||||
"context"
|
||||
"gotemplate/internal/queue"
|
||||
"log"
|
||||
|
||||
"github.com/aws/aws-sdk-go-v2/service/sqs"
|
||||
)
|
||||
|
||||
type Controllers struct {
|
||||
Document DocumentController
|
||||
}
|
||||
|
||||
type Queue struct {
|
||||
Config *queue.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)
|
||||
}
|
||||
}()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+18
-24
@@ -1,24 +1,18 @@
|
||||
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
|
||||
}
|
||||
package queue
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
|
||||
)
|
||||
|
||||
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 "DOCTEXT":
|
||||
err := queue.Controllers.Document.Sync(ctx, queue.Config, &message)
|
||||
return true, err
|
||||
}
|
||||
|
||||
return false, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user