src

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

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 }