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

Method writeStreamEntryContent

core/stream.go:465–556  ·  view source on GitHub ↗

writeStreamEntryContent writes a stream entry content as listpack

(entry *model.StreamEntry)

Source from the content-addressed store, hash-verified

463
464// writeStreamEntryContent writes a stream entry content as listpack
465func (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

Callers 1

writeStreamEntriesMethod · 0.95

Calls 3

writeStringMethod · 0.95
unsafeBytes2StrFunction · 0.70

Tested by

no test coverage detected