// Command content-decoders contains copy-oriented Brotli and Zstandard decoder
// examples. It is deliberately a main package so applications cannot depend on
// it as a library.
package main
import (
"io"
"github.com/andybalholm/brotli"
"github.com/klauspost/compress/zstd"
"github.com/mgurevin/recorder"
)
func main() {}
// Brotli opens a streaming Brotli decoder suitable for Content-Encoding: br.
func Brotli(r io.Reader) (io.ReadCloser, error) {
return io.NopCloser(brotli.NewReader(r)), nil
}
// Zstandard opens a streaming Zstandard decoder suitable for
// Content-Encoding: zstd.
func Zstandard(r io.Reader) (io.ReadCloser, error) {
decoder, err := zstd.NewReader(r)
if err != nil {
return nil, err
}
return decoder.IOReadCloser(), nil
}
// Decoders returns the content decoders to add to recorder.Config.
func Decoders() map[string]recorder.ContentDecoder {
return map[string]recorder.ContentDecoder{"br": Brotli, "zstd": Zstandard}
}
// Command csv-redactor contains a copy-oriented example of a streaming
// recorder.BodyRedactor for CSV request and response bodies. It is deliberately
// a main package so applications cannot depend on it as a library.
package main
import (
"encoding/csv"
"errors"
"fmt"
"io"
"strings"
"sync"
"github.com/mgurevin/recorder"
)
// Redactor replaces values in columns selected by their header names.
// A Redactor is safe for concurrent use; every call to Redact creates an
// independent parser and writer.
type Redactor struct {
Columns []string
Comma rune
}
func main() {}
// Redact implements recorder.BodyRedactor. Recorder supplies protector, so
// this redactor never handles keys or constructs protected token formats.
func (r Redactor) Redact(dst io.Writer, _ string, protector recorder.BodyValueProtector) (io.WriteCloser, error) {
if protector == nil {
return nil, errors.New("csv redactor: nil value protector")
}
return r.open(dst, protector)
}
func (r Redactor) open(dst io.Writer, protector recorder.BodyValueProtector) (io.WriteCloser, error) {
if dst == nil {
return nil, errors.New("csv redactor: nil destination")
}
columns := make(map[string]struct{}, len(r.Columns))
for _, column := range r.Columns {
name := strings.ToLower(strings.TrimSpace(column))
if name == "" {
return nil, errors.New("csv redactor: column names must not be empty")
}
columns[name] = struct{}{}
}
if len(columns) == 0 {
return nil, errors.New("csv redactor: at least one column is required")
}
comma := r.Comma
if comma == 0 {
comma = ','
}
pr, pw := io.Pipe()
w := &writer{pw: pw, done: make(chan result, 1)}
go w.process(pr, dst, columns, protector, comma)
return w, nil
}
type result struct {
err error
}
type writer struct {
pw *io.PipeWriter
done chan result
closeOnce sync.Once
closeErr error
}
func (w *writer) Write(p []byte) (int, error) {
return w.pw.Write(p)
}
func (w *writer) Close() error {
w.closeOnce.Do(func() {
closeErr := w.pw.Close()
processed := <-w.done
if processed.err != nil {
w.closeErr = processed.err
} else {
w.closeErr = closeErr
}
})
return w.closeErr
}
func (w *writer) process(pr *io.PipeReader, dst io.Writer, columns map[string]struct{}, protector recorder.BodyValueProtector, comma rune) {
var processed result
defer func() {
_ = pr.CloseWithError(processed.err)
w.done <- processed
}()
reader := csv.NewReader(pr)
reader.Comma = comma
reader.FieldsPerRecord = 0
output := csv.NewWriter(dst)
output.Comma = comma
header, err := reader.Read()
if err != nil {
processed.err = fmt.Errorf("csv redactor: read header: %w", err)
return
}
selected := make([]int, 0, len(columns))
found := make(map[string]struct{}, len(columns))
for index, field := range header {
name := strings.ToLower(strings.TrimSpace(field))
if _, ok := columns[name]; ok {
selected = append(selected, index)
found[name] = struct{}{}
}
}
if len(found) != len(columns) {
processed.err = errors.New("csv redactor: one or more configured columns are missing from the header")
return
}
if err := output.Write(header); err != nil {
processed.err = fmt.Errorf("csv redactor: write header: %w", err)
return
}
for {
record, readErr := reader.Read()
if readErr == io.EOF {
break
}
if readErr != nil {
processed.err = fmt.Errorf("csv redactor: read record: %w", readErr)
return
}
for _, index := range selected {
if index >= len(record) {
processed.err = errors.New("csv redactor: record has fewer fields than its header")
return
}
value := protector.NewValue()
if _, err := io.WriteString(value, record[index]); err != nil {
processed.err = fmt.Errorf("csv redactor: buffer protected value: %w", err)
return
}
var protected strings.Builder
if err := value.FinishTo(&protected); err != nil {
processed.err = fmt.Errorf("csv redactor: finish protected value: %w", err)
return
}
record[index] = protected.String()
}
if err := output.Write(record); err != nil {
processed.err = fmt.Errorf("csv redactor: write record: %w", err)
return
}
}
output.Flush()
if err := output.Error(); err != nil {
processed.err = fmt.Errorf("csv redactor: flush output: %w", err)
}
}
// Command debug-stream demonstrates ephemeral, single-browser live inspection
// during local development. It is deliberately not a production evidence sink.
package main
import (
"context"
"errors"
"fmt"
"io"
"log"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/mgurevin/recorder"
)
func main() {
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
output, err := os.Create("debug-entries.ndjson")
if err != nil {
log.Fatal(err)
}
stream, async, err := newExampleRecorders(output)
if err != nil {
log.Fatal(err)
}
server := &http.Server{
Addr: "127.0.0.1:7070",
Handler: stream,
ReadHeaderTimeout: 5 * time.Second,
}
serverErrors := make(chan error, 1)
go func() {
serverErrors <- server.ListenAndServe()
}()
config := recorder.DefaultConfig()
if err := config.Validate(); err != nil {
log.Fatal(err)
}
client := &http.Client{
Transport: recorder.NewTransport(http.DefaultTransport, async, config),
}
log.Printf("open the Inspector, choose live, and connect to http://127.0.0.1:7070")
// This request stands in for application traffic. The exchange appears
// after the response body reaches EOF or is closed. The bounded debug queue
// keeps it available when the Inspector connects after this request.
response, err := client.Get("https://example.com/")
if err != nil {
log.Printf("example request: %v", err)
} else {
_, _ = io.Copy(io.Discard, response.Body)
_ = response.Body.Close()
}
select {
case <-ctx.Done():
case err := <-serverErrors:
if !errors.Is(err, http.ErrServerClosed) {
log.Printf("debug server: %v", err)
}
}
shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
if err := server.Shutdown(shutdownCtx); err != nil {
log.Printf("shut down debug server: %v", err)
}
if err := async.Close(shutdownCtx); err != nil {
log.Printf("close async recorder: %v", err)
}
if err := output.Sync(); err != nil {
log.Printf("sync NDJSON output: %v", err)
}
if err := output.Close(); err != nil {
log.Printf("close NDJSON output: %v", err)
}
if err := stream.Close(); err != nil {
log.Printf("close debug stream: %v", err)
}
}
func newExampleRecorders(output io.Writer) (*recorder.DebugStreamRecorder, *recorder.AsyncRecorder, error) {
stream, err := recorder.NewDebugStreamRecorder(recorder.DefaultDebugStreamRecorderConfig())
if err != nil {
return nil, nil, fmt.Errorf("create debug stream recorder: %w", err)
}
fileSink := recorder.NewJSONStreamRecorder(output)
fanout := recorder.NewMultiRecorder(fileSink, stream)
asyncConfig := recorder.DefaultAsyncRecorderConfig()
asyncConfig.QueueCapacity = 256
asyncConfig.BatchSize = 32
asyncConfig.FlushInterval = 50 * time.Millisecond
async, err := recorder.NewAsyncRecorder(fanout, asyncConfig)
if err != nil {
_ = stream.Close()
return nil, nil, fmt.Errorf("create async recorder: %w", err)
}
return stream, async, nil
}
// Command har-fixture demonstrates an optional, network-free HTTP fixture
// backed by recorder HAR evidence.
package main
import (
"fmt"
"io"
"net/http"
"os"
"github.com/mgurevin/recorder/hario"
"github.com/mgurevin/recorder/hartest"
)
func main() {
if len(os.Args) != 2 {
fmt.Fprintln(os.Stderr, "usage: har-fixture capture.har")
os.Exit(2)
}
if err := run(os.Stdout, os.Args[1]); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
}
func run(output io.Writer, path string) error {
file, err := os.Open(path)
if err != nil {
return fmt.Errorf("open fixture: %w", err)
}
document, err := hario.ReadHAR(file, hario.DefaultReadConfig())
closeErr := file.Close()
if err != nil {
return fmt.Errorf("read fixture: %w", err)
}
if closeErr != nil {
return fmt.Errorf("close fixture: %w", closeErr)
}
fixture, err := hartest.NewTransport(document.Log.Entries, hartest.DefaultConfig())
if err != nil {
return fmt.Errorf("create fixture transport: %w", err)
}
client := &http.Client{Transport: fixture}
response, err := client.Get("https://api.example.test/orders/42")
if err != nil {
return fmt.Errorf("execute fixture request: %w", err)
}
body, readErr := io.ReadAll(response.Body)
bodyCloseErr := response.Body.Close()
if readErr != nil {
return fmt.Errorf("read fixture response: %w", readErr)
}
if bodyCloseErr != nil {
return fmt.Errorf("close fixture response: %w", bodyCloseErr)
}
if _, err := fmt.Fprintf(output, "%s\n%s\n", response.Status, body); err != nil {
return fmt.Errorf("write fixture result: %w", err)
}
return fixture.Verify()
}