WriteStreamObject writes a stream object to RDB file
(key string, stream *model.StreamObject, options ...interface{})
| 335 | |
| 336 | // WriteStreamObject writes a stream object to RDB file |
| 337 | func (enc *Encoder) WriteStreamObject(key string, stream *model.StreamObject, options ...interface{}) error { |
| 338 | err := enc.beforeWriteObject(options...) |
| 339 | if err != nil { |
| 340 | return err |
| 341 | } |
| 342 | |
| 343 | // Write stream type based on version |
| 344 | var streamType byte |
| 345 | switch stream.Version { |
| 346 | case 1: |
| 347 | streamType = typeStreamListPacks |
| 348 | case 2: |
| 349 | streamType = typeStreamListPacks2 |
| 350 | case 3: |
| 351 | streamType = typeStreamListPacks3 |
| 352 | default: |
| 353 | streamType = typeStreamListPacks // default to version 1 |
| 354 | } |
| 355 | |
| 356 | err = enc.write([]byte{streamType}) |
| 357 | if err != nil { |
| 358 | return err |
| 359 | } |
| 360 | |
| 361 | err = enc.writeString(key) |
| 362 | if err != nil { |
| 363 | return err |
| 364 | } |
| 365 | |
| 366 | // Write stream entries |
| 367 | err = enc.writeStreamEntries(stream.Entries) |
| 368 | if err != nil { |
| 369 | return err |
| 370 | } |
| 371 | |
| 372 | // Write stream length |
| 373 | err = enc.writeLength(stream.Length) |
| 374 | if err != nil { |
| 375 | return err |
| 376 | } |
| 377 | |
| 378 | // Write last ID |
| 379 | err = enc.writeStreamId(stream.LastId) |
| 380 | if err != nil { |
| 381 | return err |
| 382 | } |
| 383 | |
| 384 | // Write version 2+ fields if available |
| 385 | if stream.Version >= 2 { |
| 386 | if stream.FirstId != nil { |
| 387 | err = enc.writeStreamId(stream.FirstId) |
| 388 | if err != nil { |
| 389 | return err |
| 390 | } |
| 391 | } else { |
| 392 | // Write zero ID if FirstId is nil |
| 393 | err = enc.writeStreamId(&model.StreamId{Ms: 0, Sequence: 0}) |
| 394 | if err != nil { |