Encode encodes a single EventStream message to the io.Writer the Encoder was created with. An error is returned if writing the message fails.
(w io.Writer, msg Message)
| 45 | // Encode encodes a single EventStream message to the io.Writer the Encoder |
| 46 | // was created with. An error is returned if writing the message fails. |
| 47 | func (e *Encoder) Encode(w io.Writer, msg Message) (err error) { |
| 48 | e.headersBuf.Reset() |
| 49 | e.messageBuf.Reset() |
| 50 | |
| 51 | var writer io.Writer = e.messageBuf |
| 52 | if e.options.Logger != nil && e.options.LogMessages { |
| 53 | encodeMsgBuf := bytes.NewBuffer(nil) |
| 54 | writer = io.MultiWriter(writer, encodeMsgBuf) |
| 55 | defer func() { |
| 56 | logMessageEncode(e.options.Logger, encodeMsgBuf, msg, err) |
| 57 | }() |
| 58 | } |
| 59 | |
| 60 | if err = EncodeHeaders(e.headersBuf, msg.Headers); err != nil { |
| 61 | return err |
| 62 | } |
| 63 | |
| 64 | crc := crc32.New(crc32IEEETable) |
| 65 | hashWriter := io.MultiWriter(writer, crc) |
| 66 | |
| 67 | headersLen := uint32(e.headersBuf.Len()) |
| 68 | payloadLen := uint32(len(msg.Payload)) |
| 69 | |
| 70 | if err = encodePrelude(hashWriter, crc, headersLen, payloadLen); err != nil { |
| 71 | return err |
| 72 | } |
| 73 | |
| 74 | if headersLen > 0 { |
| 75 | if _, err = io.Copy(hashWriter, e.headersBuf); err != nil { |
| 76 | return err |
| 77 | } |
| 78 | } |
| 79 | |
| 80 | if payloadLen > 0 { |
| 81 | if _, err = hashWriter.Write(msg.Payload); err != nil { |
| 82 | return err |
| 83 | } |
| 84 | } |
| 85 | |
| 86 | msgCRC := crc.Sum32() |
| 87 | if err := binary.Write(writer, binary.BigEndian, msgCRC); err != nil { |
| 88 | return err |
| 89 | } |
| 90 | |
| 91 | _, err = io.Copy(w, e.messageBuf) |
| 92 | |
| 93 | return err |
| 94 | } |
| 95 | |
| 96 | func logMessageEncode(logger logging.Logger, msgBuf *bytes.Buffer, msg Message, encodeErr error) { |
| 97 | w := bytes.NewBuffer(nil) |