writeStreamGroups writes stream groups
(groups []*model.StreamGroup, version uint)
| 636 | |
| 637 | // writeStreamGroups writes stream groups |
| 638 | func (enc *Encoder) writeStreamGroups(groups []*model.StreamGroup, version uint) error { |
| 639 | err := enc.writeLength(uint64(len(groups))) |
| 640 | if err != nil { |
| 641 | return err |
| 642 | } |
| 643 | |
| 644 | for _, group := range groups { |
| 645 | // Write group name |
| 646 | err = enc.writeString(group.Name) |
| 647 | if err != nil { |
| 648 | return err |
| 649 | } |
| 650 | |
| 651 | // Write last ID |
| 652 | err = enc.writeStreamId(group.LastId) |
| 653 | if err != nil { |
| 654 | return err |
| 655 | } |
| 656 | |
| 657 | // Write entries read (version 2+) |
| 658 | if version >= 2 { |
| 659 | err = enc.writeLength(group.EntriesRead) |
| 660 | if err != nil { |
| 661 | return err |
| 662 | } |
| 663 | } |
| 664 | |
| 665 | // Write pending list |
| 666 | err = enc.writeLength(uint64(len(group.Pending))) |
| 667 | if err != nil { |
| 668 | return err |
| 669 | } |
| 670 | |
| 671 | for _, pending := range group.Pending { |
| 672 | // Write message ID |
| 673 | msBytes := make([]byte, 8) |
| 674 | binary.BigEndian.PutUint64(msBytes, pending.Id.Ms) |
| 675 | err = enc.write(msBytes) |
| 676 | if err != nil { |
| 677 | return err |
| 678 | } |
| 679 | |
| 680 | seqBytes := make([]byte, 8) |
| 681 | binary.BigEndian.PutUint64(seqBytes, pending.Id.Sequence) |
| 682 | err = enc.write(seqBytes) |
| 683 | if err != nil { |
| 684 | return err |
| 685 | } |
| 686 | |
| 687 | // Write delivery time |
| 688 | deliveryTimeBytes := make([]byte, 8) |
| 689 | binary.LittleEndian.PutUint64(deliveryTimeBytes, pending.DeliveryTime) |
| 690 | err = enc.write(deliveryTimeBytes) |
| 691 | if err != nil { |
| 692 | return err |
| 693 | } |
| 694 | |
| 695 | // Write delivery count |
no test coverage detected