message_set.go 1.6 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586
  1. package protocol
  2. import enc "sarama/encoding"
  3. type MessageBlock struct {
  4. Offset int64
  5. Msg *Message
  6. }
  7. func (msb *MessageBlock) Encode(pe enc.PacketEncoder) error {
  8. pe.PutInt64(msb.Offset)
  9. pe.Push(&enc.LengthField{})
  10. err := msb.Msg.Encode(pe)
  11. if err != nil {
  12. return err
  13. }
  14. return pe.Pop()
  15. }
  16. func (msb *MessageBlock) Decode(pd enc.PacketDecoder) (err error) {
  17. msb.Offset, err = pd.GetInt64()
  18. if err != nil {
  19. return err
  20. }
  21. pd.Push(&enc.LengthField{})
  22. if err != nil {
  23. return err
  24. }
  25. msb.Msg = new(Message)
  26. err = msb.Msg.Decode(pd)
  27. if err != nil {
  28. return err
  29. }
  30. err = pd.Pop()
  31. if err != nil {
  32. return err
  33. }
  34. return nil
  35. }
  36. type MessageSet struct {
  37. PartialTrailingMessage bool // whether the set on the wire contained an incomplete trailing MessageBlock
  38. Messages []*MessageBlock
  39. }
  40. func (ms *MessageSet) Encode(pe enc.PacketEncoder) error {
  41. for i := range ms.Messages {
  42. err := ms.Messages[i].Encode(pe)
  43. if err != nil {
  44. return err
  45. }
  46. }
  47. return nil
  48. }
  49. func (ms *MessageSet) Decode(pd enc.PacketDecoder) (err error) {
  50. ms.Messages = nil
  51. for pd.Remaining() > 0 {
  52. msb := new(MessageBlock)
  53. err = msb.Decode(pd)
  54. switch {
  55. case err == nil:
  56. ms.Messages = append(ms.Messages, msb)
  57. case err == enc.InsufficientData:
  58. // As an optimization the server is allowed to return a partial message at the
  59. // end of the message set. Clients should handle this case. So we just ignore such things.
  60. ms.PartialTrailingMessage = true
  61. return nil
  62. default:
  63. return err
  64. }
  65. }
  66. return nil
  67. }
  68. func (ms *MessageSet) addMessage(msg *Message) {
  69. block := new(MessageBlock)
  70. block.Msg = msg
  71. ms.Messages = append(ms.Messages, block)
  72. }