message_set.go 2.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102
  1. package sarama
  2. type MessageBlock struct {
  3. Offset int64
  4. Msg *Message
  5. }
  6. // Messages convenience helper which returns either all the
  7. // messages that are wrapped in this block
  8. func (msb *MessageBlock) Messages() []*MessageBlock {
  9. if msb.Msg.Set != nil {
  10. return msb.Msg.Set.Messages
  11. }
  12. return []*MessageBlock{msb}
  13. }
  14. func (msb *MessageBlock) encode(pe packetEncoder) error {
  15. pe.putInt64(msb.Offset)
  16. pe.push(&lengthField{})
  17. err := msb.Msg.encode(pe)
  18. if err != nil {
  19. return err
  20. }
  21. return pe.pop()
  22. }
  23. func (msb *MessageBlock) decode(pd packetDecoder) (err error) {
  24. if msb.Offset, err = pd.getInt64(); err != nil {
  25. return err
  26. }
  27. if err = pd.push(&lengthField{}); err != nil {
  28. return err
  29. }
  30. msb.Msg = new(Message)
  31. if err = msb.Msg.decode(pd); err != nil {
  32. return err
  33. }
  34. if err = pd.pop(); err != nil {
  35. return err
  36. }
  37. return nil
  38. }
  39. type MessageSet struct {
  40. PartialTrailingMessage bool // whether the set on the wire contained an incomplete trailing MessageBlock
  41. Messages []*MessageBlock
  42. }
  43. func (ms *MessageSet) encode(pe packetEncoder) error {
  44. for i := range ms.Messages {
  45. err := ms.Messages[i].encode(pe)
  46. if err != nil {
  47. return err
  48. }
  49. }
  50. return nil
  51. }
  52. func (ms *MessageSet) decode(pd packetDecoder) (err error) {
  53. ms.Messages = nil
  54. for pd.remaining() > 0 {
  55. magic, err := magicValue(pd)
  56. if err != nil {
  57. if err == ErrInsufficientData {
  58. ms.PartialTrailingMessage = true
  59. return nil
  60. }
  61. return err
  62. }
  63. if magic > 1 {
  64. return nil
  65. }
  66. msb := new(MessageBlock)
  67. err = msb.decode(pd)
  68. switch err {
  69. case nil:
  70. ms.Messages = append(ms.Messages, msb)
  71. case ErrInsufficientData:
  72. // As an optimization the server is allowed to return a partial message at the
  73. // end of the message set. Clients should handle this case. So we just ignore such things.
  74. ms.PartialTrailingMessage = true
  75. return nil
  76. default:
  77. return err
  78. }
  79. }
  80. return nil
  81. }
  82. func (ms *MessageSet) addMessage(msg *Message) {
  83. block := new(MessageBlock)
  84. block.Msg = msg
  85. ms.Messages = append(ms.Messages, block)
  86. }