consumer.go 5.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200
  1. package services
  2. import (
  3. "context"
  4. "crypto/md5"
  5. "encoding/hex"
  6. "encoding/json"
  7. "fmt"
  8. "io"
  9. "log"
  10. "strconv"
  11. "time"
  12. "github.com/aws/aws-sdk-go-v2/service/sqs"
  13. "github.com/aws/aws-sdk-go-v2/service/sqs/types"
  14. "github.com/davecgh/go-spew/spew"
  15. "github.com/fgm/izidic"
  16. "github.com/google/uuid"
  17. )
  18. type Event struct {
  19. MessageID uuid.UUID `json:"MessageId"`
  20. BodySum string `json:"MD5OfBody"`
  21. AttrSum string `json:"MD5MofMessageAttributes"`
  22. EventAttributes
  23. MessageAttributes map[string]types.MessageAttributeValue
  24. Body []byte
  25. }
  26. func (e Event) IsRetryable() bool {
  27. maybeRetry, ok := e.MessageAttributes["retry"]
  28. if !ok {
  29. return false
  30. }
  31. if maybeRetry.DataType == nil || *maybeRetry.DataType != "String" || maybeRetry.StringValue == nil {
  32. return false
  33. }
  34. return *maybeRetry.StringValue == "1"
  35. }
  36. type EventAttributes struct {
  37. SenderID string `json:"SenderId"`
  38. SentTime time.Time `json:"SentTimestamp"`
  39. ApproximateReceiveCount int `json:"ApproximateReceiveCount"`
  40. ApproximateFirstReceiveTimestamp time.Time `json:"ApproximateFirstReceiveTimestamp"`
  41. }
  42. type Handler func(ctx context.Context, enc *json.Encoder, msgID uuid.UUID, sent time.Time, input []byte, meta map[string]types.MessageAttributeValue) error
  43. type message types.Message
  44. func (m message) String() string {
  45. if m.MessageId == nil {
  46. return "0"
  47. }
  48. return *m.MessageId
  49. }
  50. func ConsumerService(dic *izidic.Container) (any, error) {
  51. cli := dic.MustService("sqs").(*sqs.Client)
  52. w := dic.MustParam(PWriter).(io.Writer)
  53. hdl := dic.MustParam(PHandler).(Handler)
  54. enc := json.NewEncoder(w)
  55. return func(ctx context.Context, qURL string) error {
  56. return consumeMessage(ctx, w, enc, cli, qURL, hdl)
  57. }, nil
  58. }
  59. func ReceiverService(dic *izidic.Container) (any, error) {
  60. cli := dic.MustService("sqs").(*sqs.Client)
  61. w := dic.MustParam(PWriter).(io.Writer)
  62. return func(ctx context.Context, qURL string) {
  63. receiveMessage(ctx, w, cli, qURL)
  64. }, nil
  65. }
  66. func consumeMessage(ctx context.Context, _ io.Writer, enc *json.Encoder, client *sqs.Client, qURL string, hdl Handler) error {
  67. rmi := sqs.ReceiveMessageInput{
  68. QueueUrl: &qURL,
  69. AttributeNames: []types.QueueAttributeName{"All"},
  70. MessageAttributeNames: []string{"All"},
  71. MaxNumberOfMessages: 1, // Default, also used when set to 0
  72. VisibilityTimeout: 1,
  73. WaitTimeSeconds: 5,
  74. }
  75. for {
  76. recv, err := client.ReceiveMessage(ctx, &rmi)
  77. if err != nil {
  78. return fmt.Errorf("failed receiving from queue: %w, aborting", err)
  79. }
  80. if len(recv.Messages) == 0 {
  81. log.Printf("No message with %d seconds timeout\n", rmi.WaitTimeSeconds)
  82. continue
  83. }
  84. if len(recv.Messages) != 1 {
  85. return fmt.Errorf("unexpected number of messages: %d, expected 0 or 1, aborting", len(recv.Messages))
  86. }
  87. msg := message(recv.Messages[0])
  88. evt, err := validateMessage(msg)
  89. if err != nil {
  90. log.Printf("invalid message %s: %v, dropping it anyway", msg, err)
  91. } else {
  92. if err := hdl(ctx, enc, evt.MessageID, evt.SentTime, evt.Body, evt.MessageAttributes); err != nil {
  93. log.Printf("message %s failed processing : %v, dropping it anyway\n", msg, err)
  94. } else {
  95. log.Printf("message %s processed successfully\n", msg)
  96. }
  97. }
  98. if evt.IsRetryable() {
  99. log.Printf("message %s not deleted, for retry", msg)
  100. } else {
  101. dmi := sqs.DeleteMessageInput{
  102. QueueUrl: &qURL,
  103. ReceiptHandle: msg.ReceiptHandle,
  104. }
  105. _, err = client.DeleteMessage(ctx, &dmi)
  106. if err != nil {
  107. log.Printf("Error deleting message %s after successful processing: %v\n", msg, err)
  108. continue
  109. }
  110. log.Printf("message %s deleted after processing\n", msg)
  111. }
  112. }
  113. }
  114. func receiveMessage(ctx context.Context, w io.Writer, client *sqs.Client, qURL string) {
  115. rmi := sqs.ReceiveMessageInput{
  116. QueueUrl: &qURL,
  117. AttributeNames: []types.QueueAttributeName{"All"},
  118. MessageAttributeNames: []string{"All"},
  119. VisibilityTimeout: 1,
  120. WaitTimeSeconds: 5,
  121. }
  122. msg, err := client.ReceiveMessage(ctx, &rmi)
  123. if err != nil {
  124. log.Fatalf("failed receiving from queue %s: %v", qURL, err)
  125. }
  126. spew.Fdump(w, msg.Messages)
  127. }
  128. func validateMessage(msg message) (*Event, error) {
  129. var (
  130. err error
  131. evt = Event{EventAttributes: EventAttributes{SenderID: msg.Attributes["SenderId"]}}
  132. )
  133. // Top-level fields
  134. if msg.MessageId == nil {
  135. return nil, fmt.Errorf("error: MessageId is nil")
  136. }
  137. evt.MessageID, err = uuid.Parse(*msg.MessageId)
  138. if err != nil {
  139. return nil, fmt.Errorf("error parsing MessageId as a UUID: %w", err)
  140. }
  141. if msg.MD5OfBody == nil {
  142. return nil, fmt.Errorf("error: MD5OfBody is nil")
  143. }
  144. evt.BodySum = *msg.MD5OfBody
  145. if msg.MD5OfMessageAttributes != nil {
  146. evt.AttrSum = *msg.MD5OfMessageAttributes
  147. }
  148. // EventAttributes fields
  149. msec, err := strconv.Atoi(msg.Attributes["SentTimestamp"])
  150. if err != nil {
  151. return nil, fmt.Errorf("error parsing SentTimestamp as milliseconds: %w", err)
  152. }
  153. evt.EventAttributes.SentTime = time.Unix(int64(msec)/1000, int64(msec)%1000)
  154. evt.EventAttributes.ApproximateReceiveCount, err = strconv.Atoi(msg.Attributes["ApproximateReceiveCount"])
  155. if err != nil {
  156. return nil, fmt.Errorf("error parsing ApproximateReceiveCount: %w", err)
  157. }
  158. msec, err = strconv.Atoi(msg.Attributes["ApproximateFirstReceiveTimestamp"])
  159. if err != nil {
  160. return nil, fmt.Errorf("error parsing ApproximateFirstReceiveTimestam as milliseconds: %w", err)
  161. }
  162. evt.EventAttributes.ApproximateFirstReceiveTimestamp = time.Unix(int64(msec)/1000, int64(msec)*1000)
  163. evt.MessageAttributes = msg.MessageAttributes
  164. // EventBody field
  165. if msg.Body == nil {
  166. return nil, fmt.Errorf("message body is nil")
  167. }
  168. body := *msg.Body
  169. bs, err := hex.DecodeString(evt.BodySum)
  170. if err != nil {
  171. return nil, fmt.Errorf("error parsy body sum as a hex string: %w", err)
  172. }
  173. expected := *(*[16]byte)(bs)
  174. actual := md5.Sum([]byte(body))
  175. if actual != expected {
  176. return nil, fmt.Errorf("error parsing body sum as a MD5 sum")
  177. }
  178. evt.Body = []byte(body)
  179. return &evt, nil
  180. }