signer.go (2170B)
1 package eventstream 2 3 import ( 4 "bytes" 5 "io" 6 "time" 7 ) 8 9 // MessageSigner signs event stream message header and payload byte pairs. 10 // Each invocation chains off the previous signature. 11 type MessageSigner interface { 12 SignMessage(headers, payload []byte, signingTime time.Time) ([]byte, error) 13 } 14 15 // SigningWriter wraps an io.WriteCloser and signs each event stream message 16 // frame written to it. Each Write call MUST contain exactly one complete 17 // encoded event stream message frame. 18 // 19 // The signing writer wraps each incoming frame in an outer event stream 20 // message with :date and :chunk-signature headers, then encodes the outer 21 // message to the underlying writer. 22 // 23 // Close sends a signed empty message to signal end-of-stream, then closes 24 // the underlying writer. 25 type SigningWriter struct { 26 writer io.WriteCloser 27 signer MessageSigner 28 encoder *Encoder 29 30 headersBuf bytes.Buffer 31 } 32 33 // NewSigningWriter returns a SigningWriter that signs frames and writes them 34 // to w. 35 func NewSigningWriter(w io.WriteCloser, signer MessageSigner) *SigningWriter { 36 return &SigningWriter{ 37 writer: w, 38 signer: signer, 39 encoder: NewEncoder(), 40 } 41 } 42 43 // Write signs a complete event stream message frame and writes the signed 44 // outer envelope to the underlying writer. 45 func (s *SigningWriter) Write(frame []byte) (int, error) { 46 if err := s.signAndWrite(frame); err != nil { 47 return 0, err 48 } 49 return len(frame), nil 50 } 51 52 // Close sends a signed empty message to signal end-of-stream, then closes 53 // the underlying writer. 54 func (s *SigningWriter) Close() error { 55 if err := s.signAndWrite([]byte{}); err != nil { 56 _ = s.writer.Close() 57 return err 58 } 59 return s.writer.Close() 60 } 61 62 func (s *SigningWriter) signAndWrite(payload []byte) error { 63 now := time.Now().UTC() 64 65 var msg Message 66 msg.Headers.Set(DateHeader, TimestampValue(now)) 67 msg.Payload = payload 68 69 s.headersBuf.Reset() 70 if err := EncodeHeaders(&s.headersBuf, msg.Headers); err != nil { 71 return err 72 } 73 74 sig, err := s.signer.SignMessage(s.headersBuf.Bytes(), payload, now) 75 if err != nil { 76 return err 77 } 78 79 msg.Headers.Set(ChunkSignatureHeader, BytesValue(sig)) 80 81 return s.encoder.Encode(s.writer, msg) 82 }