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

Method readStreamGroups

core/stream.go:219–334  ·  view source on GitHub ↗
(version uint)

Source from the content-addressed store, hash-verified

217}
218
219func (dec *Decoder) readStreamGroups(version uint) ([]*model.StreamGroup, error) {
220 groupCount, _, err := dec.readLength()
221 if err != nil {
222 return nil, err
223 }
224 groups := make([]*model.StreamGroup, 0, int(groupCount))
225 for i := uint64(0); i < groupCount; i++ {
226 name, _ := dec.readString()
227 if err != nil {
228 return nil, err
229 }
230 lastId, err := dec.readStreamId()
231 if err != nil {
232 return nil, err
233 }
234
235 var entriesRead uint64
236 if version >= 2 {
237 entriesRead, _, err = dec.readLength()
238 if err != nil {
239 return nil, err
240 }
241 }
242
243 // read pending list
244 pendingCount, _, err := dec.readLength()
245 if err != nil {
246 return nil, err
247 }
248 pending := make([]*model.StreamNAck, 0, int(pendingCount))
249 for j := uint64(0); j < pendingCount; j++ {
250 if err := dec.readFull(dec.buffer); err != nil {
251 return nil, err
252 }
253 ms := binary.BigEndian.Uint64(dec.buffer)
254 if err := dec.readFull(dec.buffer); err != nil {
255 return nil, err
256 }
257 seq := binary.BigEndian.Uint64(dec.buffer)
258 streamId := &model.StreamId{
259 Ms: ms,
260 Sequence: seq,
261 }
262 if err := dec.readFull(dec.buffer); err != nil {
263 return nil, err
264 }
265 deliveryTime := binary.LittleEndian.Uint64(dec.buffer)
266 deliveryCount, _, err := dec.readLength()
267 if err != nil {
268 return nil, err
269 }
270 pending = append(pending, &model.StreamNAck{
271 Id: streamId,
272 DeliveryTime: deliveryTime,
273 DeliveryCount: deliveryCount,
274 })
275 }
276

Callers 1

readStreamListPacksMethod · 0.95

Calls 5

readLengthMethod · 0.95
readStringMethod · 0.95
readStreamIdMethod · 0.95
readFullMethod · 0.95
unsafeBytes2StrFunction · 0.70

Tested by

no test coverage detected