message.go (1890B)
1 package eventstream 2 3 import ( 4 "bytes" 5 "encoding/binary" 6 "hash/crc32" 7 ) 8 9 const preludeLen = 8 10 const preludeCRCLen = 4 11 const msgCRCLen = 4 12 const minMsgLen = preludeLen + preludeCRCLen + msgCRCLen 13 14 var crc32IEEETable = crc32.MakeTable(crc32.IEEE) 15 16 // A Message provides the eventstream message representation. 17 type Message struct { 18 Headers Headers 19 Payload []byte 20 } 21 22 func (m *Message) rawMessage() (rawMessage, error) { 23 var raw rawMessage 24 25 if len(m.Headers) > 0 { 26 var headers bytes.Buffer 27 if err := EncodeHeaders(&headers, m.Headers); err != nil { 28 return rawMessage{}, err 29 } 30 raw.Headers = headers.Bytes() 31 raw.HeadersLen = uint32(len(raw.Headers)) 32 } 33 34 raw.Length = raw.HeadersLen + uint32(len(m.Payload)) + minMsgLen 35 36 hash := crc32.New(crc32IEEETable) 37 binaryWriteFields(hash, binary.BigEndian, raw.Length, raw.HeadersLen) 38 raw.PreludeCRC = hash.Sum32() 39 40 binaryWriteFields(hash, binary.BigEndian, raw.PreludeCRC) 41 42 if raw.HeadersLen > 0 { 43 hash.Write(raw.Headers) 44 } 45 46 // Read payload bytes and update hash for it as well. 47 if len(m.Payload) > 0 { 48 raw.Payload = m.Payload 49 hash.Write(raw.Payload) 50 } 51 52 raw.CRC = hash.Sum32() 53 54 return raw, nil 55 } 56 57 // Clone returns a deep copy of the message. 58 func (m Message) Clone() Message { 59 var payload []byte 60 if m.Payload != nil { 61 payload = make([]byte, len(m.Payload)) 62 copy(payload, m.Payload) 63 } 64 65 return Message{ 66 Headers: m.Headers.Clone(), 67 Payload: payload, 68 } 69 } 70 71 type messagePrelude struct { 72 Length uint32 73 HeadersLen uint32 74 PreludeCRC uint32 75 } 76 77 func (p messagePrelude) PayloadLen() uint32 { 78 return p.Length - p.HeadersLen - minMsgLen 79 } 80 81 func (p messagePrelude) ValidateLens() error { 82 if p.Length == 0 { 83 return LengthError{ 84 Part: "message prelude", 85 Want: minMsgLen, 86 Have: int(p.Length), 87 } 88 } 89 return nil 90 } 91 92 type rawMessage struct { 93 messagePrelude 94 95 Headers []byte 96 Payload []byte 97 98 CRC uint32 99 }