decode.go 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199
  1. package eventstream
  2. import (
  3. "bytes"
  4. "encoding/binary"
  5. "encoding/hex"
  6. "encoding/json"
  7. "fmt"
  8. "hash"
  9. "hash/crc32"
  10. "io"
  11. "github.com/aws/aws-sdk-go/aws"
  12. )
  13. // Decoder provides decoding of an Event Stream messages.
  14. type Decoder struct {
  15. r io.Reader
  16. logger aws.Logger
  17. }
  18. // NewDecoder initializes and returns a Decoder for decoding event
  19. // stream messages from the reader provided.
  20. func NewDecoder(r io.Reader) *Decoder {
  21. return &Decoder{
  22. r: r,
  23. }
  24. }
  25. // Decode attempts to decode a single message from the event stream reader.
  26. // Will return the event stream message, or error if Decode fails to read
  27. // the message from the stream.
  28. func (d *Decoder) Decode(payloadBuf []byte) (m Message, err error) {
  29. reader := d.r
  30. if d.logger != nil {
  31. debugMsgBuf := bytes.NewBuffer(nil)
  32. reader = io.TeeReader(reader, debugMsgBuf)
  33. defer func() {
  34. logMessageDecode(d.logger, debugMsgBuf, m, err)
  35. }()
  36. }
  37. crc := crc32.New(crc32IEEETable)
  38. hashReader := io.TeeReader(reader, crc)
  39. prelude, err := decodePrelude(hashReader, crc)
  40. if err != nil {
  41. return Message{}, err
  42. }
  43. if prelude.HeadersLen > 0 {
  44. lr := io.LimitReader(hashReader, int64(prelude.HeadersLen))
  45. m.Headers, err = decodeHeaders(lr)
  46. if err != nil {
  47. return Message{}, err
  48. }
  49. }
  50. if payloadLen := prelude.PayloadLen(); payloadLen > 0 {
  51. buf, err := decodePayload(payloadBuf, io.LimitReader(hashReader, int64(payloadLen)))
  52. if err != nil {
  53. return Message{}, err
  54. }
  55. m.Payload = buf
  56. }
  57. msgCRC := crc.Sum32()
  58. if err := validateCRC(reader, msgCRC); err != nil {
  59. return Message{}, err
  60. }
  61. return m, nil
  62. }
  63. // UseLogger specifies the Logger that that the decoder should use to log the
  64. // message decode to.
  65. func (d *Decoder) UseLogger(logger aws.Logger) {
  66. d.logger = logger
  67. }
  68. func logMessageDecode(logger aws.Logger, msgBuf *bytes.Buffer, msg Message, decodeErr error) {
  69. w := bytes.NewBuffer(nil)
  70. defer func() { logger.Log(w.String()) }()
  71. fmt.Fprintf(w, "Raw message:\n%s\n",
  72. hex.Dump(msgBuf.Bytes()))
  73. if decodeErr != nil {
  74. fmt.Fprintf(w, "Decode error: %v\n", decodeErr)
  75. return
  76. }
  77. rawMsg, err := msg.rawMessage()
  78. if err != nil {
  79. fmt.Fprintf(w, "failed to create raw message, %v\n", err)
  80. return
  81. }
  82. decodedMsg := decodedMessage{
  83. rawMessage: rawMsg,
  84. Headers: decodedHeaders(msg.Headers),
  85. }
  86. fmt.Fprintf(w, "Decoded message:\n")
  87. encoder := json.NewEncoder(w)
  88. if err := encoder.Encode(decodedMsg); err != nil {
  89. fmt.Fprintf(w, "failed to generate decoded message, %v\n", err)
  90. }
  91. }
  92. func decodePrelude(r io.Reader, crc hash.Hash32) (messagePrelude, error) {
  93. var p messagePrelude
  94. var err error
  95. p.Length, err = decodeUint32(r)
  96. if err != nil {
  97. return messagePrelude{}, err
  98. }
  99. p.HeadersLen, err = decodeUint32(r)
  100. if err != nil {
  101. return messagePrelude{}, err
  102. }
  103. if err := p.ValidateLens(); err != nil {
  104. return messagePrelude{}, err
  105. }
  106. preludeCRC := crc.Sum32()
  107. if err := validateCRC(r, preludeCRC); err != nil {
  108. return messagePrelude{}, err
  109. }
  110. p.PreludeCRC = preludeCRC
  111. return p, nil
  112. }
  113. func decodePayload(buf []byte, r io.Reader) ([]byte, error) {
  114. w := bytes.NewBuffer(buf[0:0])
  115. _, err := io.Copy(w, r)
  116. return w.Bytes(), err
  117. }
  118. func decodeUint8(r io.Reader) (uint8, error) {
  119. type byteReader interface {
  120. ReadByte() (byte, error)
  121. }
  122. if br, ok := r.(byteReader); ok {
  123. v, err := br.ReadByte()
  124. return uint8(v), err
  125. }
  126. var b [1]byte
  127. _, err := io.ReadFull(r, b[:])
  128. return uint8(b[0]), err
  129. }
  130. func decodeUint16(r io.Reader) (uint16, error) {
  131. var b [2]byte
  132. bs := b[:]
  133. _, err := io.ReadFull(r, bs)
  134. if err != nil {
  135. return 0, err
  136. }
  137. return binary.BigEndian.Uint16(bs), nil
  138. }
  139. func decodeUint32(r io.Reader) (uint32, error) {
  140. var b [4]byte
  141. bs := b[:]
  142. _, err := io.ReadFull(r, bs)
  143. if err != nil {
  144. return 0, err
  145. }
  146. return binary.BigEndian.Uint32(bs), nil
  147. }
  148. func decodeUint64(r io.Reader) (uint64, error) {
  149. var b [8]byte
  150. bs := b[:]
  151. _, err := io.ReadFull(r, bs)
  152. if err != nil {
  153. return 0, err
  154. }
  155. return binary.BigEndian.Uint64(bs), nil
  156. }
  157. func validateCRC(r io.Reader, expect uint32) error {
  158. msgCRC, err := decodeUint32(r)
  159. if err != nil {
  160. return err
  161. }
  162. if msgCRC != expect {
  163. return ChecksumError{}
  164. }
  165. return nil
  166. }