writeStreamEntryContent writes a stream entry content as listpack
(entry *model.StreamEntry)
| 463 | |
| 464 | // writeStreamEntryContent writes a stream entry content as listpack |
| 465 | func (enc *Encoder) writeStreamEntryContent(entry *model.StreamEntry) error { |
| 466 | // Calculate total messages (including deleted ones) |
| 467 | totalMsgs := len(entry.Msgs) |
| 468 | deletedCount := 0 |
| 469 | for _, msg := range entry.Msgs { |
| 470 | if msg.Deleted { |
| 471 | deletedCount++ |
| 472 | } |
| 473 | } |
| 474 | validCount := totalMsgs - deletedCount |
| 475 | |
| 476 | // Build listpack with proper backlen values |
| 477 | var entries []listpackEntry |
| 478 | |
| 479 | // Add count and deleted count |
| 480 | entries = append(entries, listpackEntry{intVal: int64(validCount)}) |
| 481 | entries = append(entries, listpackEntry{intVal: int64(deletedCount)}) |
| 482 | |
| 483 | // Add master field names |
| 484 | entries = append(entries, listpackEntry{intVal: int64(len(entry.Fields))}) |
| 485 | for _, field := range entry.Fields { |
| 486 | entries = append(entries, listpackEntry{strVal: field}) |
| 487 | } |
| 488 | // Add field count for master entry (this is what the decoder reads as "end flag") |
| 489 | entries = append(entries, listpackEntry{strVal: strconv.Itoa(len(entry.Fields))}) |
| 490 | |
| 491 | // Add messages |
| 492 | for _, msg := range entry.Msgs { |
| 493 | // Calculate flag |
| 494 | flag := StreamItemFlagNone |
| 495 | if msg.Deleted { |
| 496 | flag |= StreamItemFlagDeleted |
| 497 | } |
| 498 | |
| 499 | // Check if message uses same fields as master |
| 500 | if len(msg.Fields) == len(entry.Fields) { |
| 501 | sameFields := true |
| 502 | for _, field := range entry.Fields { |
| 503 | if _, exists := msg.Fields[field]; !exists { |
| 504 | sameFields = false |
| 505 | break |
| 506 | } |
| 507 | } |
| 508 | if sameFields { |
| 509 | flag |= StreamItemFlagSameFields |
| 510 | } |
| 511 | } |
| 512 | |
| 513 | // Add flag |
| 514 | entries = append(entries, listpackEntry{intVal: int64(flag)}) |
| 515 | |
| 516 | // Add message ID (relative to first message ID) |
| 517 | msDiff := int64(msg.Id.Ms) - int64(entry.FirstMsgId.Ms) |
| 518 | seqDiff := int64(msg.Id.Sequence) - int64(entry.FirstMsgId.Sequence) |
| 519 | entries = append(entries, listpackEntry{intVal: msDiff}) |
| 520 | entries = append(entries, listpackEntry{intVal: seqDiff}) |
| 521 | |
| 522 | // Add field count if not same fields |
no test coverage detected