MCPcopy Create free account
hub / github.com/HDT3213/rdb / writeStreamGroups

Method writeStreamGroups

core/stream.go:638–759  ·  view source on GitHub ↗

writeStreamGroups writes stream groups

(groups []*model.StreamGroup, version uint)

Source from the content-addressed store, hash-verified

636
637// writeStreamGroups writes stream groups
638func (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

Callers 1

WriteStreamObjectMethod · 0.95

Calls 4

writeLengthMethod · 0.95
writeStringMethod · 0.95
writeStreamIdMethod · 0.95
writeMethod · 0.95

Tested by

no test coverage detected