src

Go monorepo.
git clone git://code.dwrz.net/src
Log | Files | Refs

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 }