package persistence import ( "bufio" "bytes" "compress/gzip" "crypto/rand" "crypto/sha256" "encoding/hex" "errors" "fmt" "io" "log/slog" "net/mail" "os" "path/filepath" "github.com/minio/sio" "golang.org/x/crypto/hkdf" ) // WriteBlob gzips and streams r into the nonce-addressed file and returns the hex-encoded nonce that identifies the blob. // More information about encryption: https://github.com/minio/sio/blob/master/DARE.md func (p *Persistence) WriteBlob(r io.Reader) (string, int64, error) { // Generate a random nonce to derive an encryption key from the master key. var nonce [32]byte if _, err := io.ReadFull(rand.Reader, nonce[:]); err != nil { return "", 0, fmt.Errorf("failed to read random data: %w", err) } // Use the nonce as file name nonceHex := hex.EncodeToString(nonce[:]) dir := filepath.Join(p.blobRoot, nonceHex[:2]) if err := os.MkdirAll(dir, 0755); err != nil { return "", 0, fmt.Errorf("creating blob dir: %w", err) } finalPath := filepath.Join(dir, nonceHex) f, err := os.Create(finalPath) if err != nil { return "", 0, fmt.Errorf("creating file: %w", err) } // Derive an encryption key from the master key and the nonce var key [32]byte kdf := hkdf.New(sha256.New, p.masterkey, nonce[:], nil) if _, err = io.ReadFull(kdf, key[:]); err != nil { return "", 0, fmt.Errorf("failed to derive encryption key: %w", err) } // Create encryption writer encrypted, err := sio.EncryptWriter(f, sio.Config{Key: key[:]}) if err != nil { return "", 0, fmt.Errorf("failed to create encrypted writer: %w", err) } gzipWriter, err := gzip.NewWriterLevel(encrypted, gzip.BestCompression) if err != nil { return "", 0, fmt.Errorf("failed to create gzip writer: %w", err) } size, err := io.Copy(gzipWriter, r) if err != nil { return "", 0, fmt.Errorf("copying data: %w", err) } if err := gzipWriter.Close(); err != nil { return "", 0, fmt.Errorf("closing gzip writer: %w", err) } if err := encrypted.Close(); err != nil { return "", 0, fmt.Errorf("closing encryption writer: %w", err) } return nonceHex, size, nil } // WriteBlobMessage parses the RFC 5322 headers from r, then gzips and streams the entire message to disk. func (p *Persistence) WriteBlobMessage(r io.Reader) (string, int64, *mail.Message, error) { msg, fullReader, err := parseMailHeader(r) if err != nil { return "", 0, nil, fmt.Errorf("parsing message header: %w", err) } nonceHex, size, err := p.WriteBlob(fullReader) if err != nil { return "", 0, nil, err } return nonceHex, size, msg, nil } // parseMailHeader extracts RFC 5322 mail headers line-by-line without reading the whole body into memory. func parseMailHeader(r io.Reader) (*mail.Message, io.Reader, error) { var headerBuf bytes.Buffer br := bufio.NewReader(r) for { line, err := br.ReadBytes('\n') headerBuf.Write(line) // Blank line (\r\n or \n) marks the end of headers if bytes.Equal(line, []byte("\r\n")) || bytes.Equal(line, []byte("\n")) { break } if err != nil { if errors.Is(err, io.EOF) { break // End of message with headers only } return nil, nil, err } } // Parse header structure from the captured header bytes msg, err := mail.ReadMessage(bytes.NewReader(headerBuf.Bytes())) if err != nil { return nil, nil, err } // Reconstruct the exact stream: header bytes first, then remaining body bytes in br fullReader := io.MultiReader(&headerBuf, br) return msg, fullReader, nil } // BlobReader decrypts and unzips the stored blob func (p *Persistence) BlobReader(nonceIdentifier string) (io.ReadCloser, error) { if len(nonceIdentifier) < 8 { return nil, fmt.Errorf("invalid nonce identifier") } dir := filepath.Join(p.blobRoot, nonceIdentifier[:2]) finalPath := filepath.Join(dir, nonceIdentifier) nonce, err := hex.DecodeString(nonceIdentifier) if err != nil { return nil, fmt.Errorf("invalid nonce identifier: %w", err) } // derive an encryption key from the master key and the nonce var key [32]byte kdf := hkdf.New(sha256.New, p.masterkey, nonce, nil) if _, err = io.ReadFull(kdf, key[:]); err != nil { return nil, fmt.Errorf("failed to derive encryption key: %w", err) } f, err := os.Open(finalPath) // leave Close() to sio.DecryptReader if err != nil { return nil, fmt.Errorf("opening file: %w", err) } decrypter, err := sio.DecryptReader(f, sio.Config{Key: key[:]}) if err != nil { return nil, fmt.Errorf("failed to create decrypt reader: %w", err) } gzipReader, err := gzip.NewReader(decrypter) if err != nil { return nil, fmt.Errorf("failed to create gzip reader: %w", err) } slog.Info("reading blob", "path", finalPath) // Return a wrapper that closes both the gzip reader and the underlying file return &blobReadCloser{ Reader: gzipReader, closeFunc: func() error { gzErr := gzipReader.Close() fErr := f.Close() if gzErr != nil { return gzErr } return fErr }, }, nil } type blobReadCloser struct { io.Reader closeFunc func() error } func (b *blobReadCloser) Close() error { return b.closeFunc() } func (p *Persistence) RemoveBlob(nonceIdentifier string) error { if len(nonceIdentifier) < 8 { return fmt.Errorf("invalid nonce identifier") } dir := filepath.Join(p.blobRoot, nonceIdentifier[:2]) finalPath := filepath.Join(dir, nonceIdentifier) err := os.Remove(finalPath) if err != nil { return fmt.Errorf("deleting blob: %w", err) } slog.Info("deleted blob", "path", finalPath) err = os.Remove(dir) if err == nil { // Only logs if the directory was actually empty and successfully deleted slog.Info("deleted empty blob directory", "dir", dir) } return nil }