(version uint)
| 217 | } |
| 218 | |
| 219 | func (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 |
no test coverage detected