Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions internal/comparison/comparator.go
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,11 @@ func (c *Comparator) Compare(
if result.Reason() != "" {
attributes = append(attributes, "reason", result.Reason())
}
// dt is the candidate's latency minus the reference's: negative when the
// candidate was faster, positive when it was slower.
if reference.Latency > 0 && candidate.Latency > 0 {
attributes = append(attributes, "dt", candidate.Latency-reference.Latency)
}
c.log.InfoContext(ctx, "Response comparison completed", attributes...)
}()
if excludedRequestPath(requestPath) {
Expand Down
39 changes: 39 additions & 0 deletions internal/comparison/comparator_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -316,6 +316,45 @@ func TestLogsComparisonAndIndividualNormaliserResults(t *testing.T) {
assert.Contains(t, logs, `"level":"INFO","msg":"Response comparison completed","event":"correlation","method":"POST","path":"/test.v1.Service/Get","outcome":"divergent","differences":["$.stable"]`)
}

func TestLogsLatencyDeltaWhenBothResponsesAreTimed(t *testing.T) {
var output bytes.Buffer
log := slog.New(slog.NewJSONHandler(&output, &slog.HandlerOptions{Level: slog.LevelDebug}))
config := newConfig(t, map[string]string{"test.ts": module(`
spectre.field<v1.Response, "stable">((value) => value);
`)})
comparator, err := comparison.New(t.Context(), config, log)
assert.NoError(t, err)
assert.NoError(t, comparator.Configure(t.Context(), descriptorSet()))
reference := connectResponse(`{"stable":"same"}`)
reference.Latency = 243 * time.Millisecond
candidate := connectResponse(`{"stable":"same"}`)
candidate.Latency = 100 * time.Millisecond

result := comparator.Compare(t.Context(), http.MethodPost, "/test.v1.Service/Get", "application/json", reference, candidate)

assert.Equal(t, comparison.Resultf(comparison.Equivalent, ""), result)
// Candidate was 143ms faster, so dt is that negative delta in nanoseconds.
assert.Contains(t, output.String(), `"outcome":"equivalent","dt":-143000000`)
}

func TestOmitsLatencyDeltaWhenEitherResponseIsUntimed(t *testing.T) {
var output bytes.Buffer
log := slog.New(slog.NewJSONHandler(&output, &slog.HandlerOptions{Level: slog.LevelDebug}))
config := newConfig(t, map[string]string{"test.ts": module(`
spectre.field<v1.Response, "stable">((value) => value);
`)})
comparator, err := comparison.New(t.Context(), config, log)
assert.NoError(t, err)
assert.NoError(t, comparator.Configure(t.Context(), descriptorSet()))
reference := connectResponse(`{"stable":"same"}`)
candidate := connectResponse(`{"stable":"same"}`)
candidate.Latency = 100 * time.Millisecond

comparator.Compare(t.Context(), http.MethodPost, "/test.v1.Service/Get", "application/json", reference, candidate)

assert.NotContains(t, output.String(), `"dt":`)
}

func TestAppliesFieldNormaliserToRepeatedMessageElements(t *testing.T) {
comparator := newComparator(t, `
spectre.field<v1.Response, "users[].name">((name) => name.toLowerCase());
Expand Down
3 changes: 3 additions & 0 deletions internal/comparison/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"fmt"
"net/http"
"slices"
"time"
)

// Outcome describes whether a response pair could be compared and matched.
Expand All @@ -32,6 +33,8 @@ type Response struct {
Body []byte
// Overflow reports that Body is truncated and must not be compared.
Overflow bool
// Latency is the backend's time to first response byte, or zero if unmeasured.
Latency time.Duration
}

// Result is the value-safe result of comparing two responses.
Expand Down
38 changes: 38 additions & 0 deletions internal/ingress/handler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -348,6 +348,44 @@ quarantine:
assert.NoError(t, handler.Shutdown(t.Context()))
}

func TestCapturesBackendLatenciesForComparison(t *testing.T) {
type latencies struct {
reference time.Duration
candidate time.Duration
}
captured := make(chan latencies, 1)
comparator := newResponseComparator(
func(
ctx context.Context,
requestMethod string,
requestPath string,
requestContentType string,
reference comparison.Response,
candidate comparison.Response,
) comparison.Result {
_, _, _, _ = ctx, requestMethod, requestPath, requestContentType
captured <- latencies{reference: reference.Latency, candidate: candidate.Latency}
return comparison.Resultf(comparison.Equivalent, "")
},
)
transport := roundTripperFunc(func(request *http.Request) (*http.Response, error) {
return noContentResponse(request), nil
})
config := newTestConfig(t, "http://127.0.0.1:50051", "http://127.0.0.1:50052")
handler, err := ingress.New(config, transport, Some[ingress.DescriptorLoader](matchingDescriptorLoader()), comparator, slog.New(slog.DiscardHandler))
assert.NoError(t, err)
handler.ServeHTTP(httptest.NewRecorder(), httptest.NewRequest(http.MethodGet, "http://proxy.example/forecast", nil))

select {
case got := <-captured:
assert.True(t, got.reference > 0, "reference latency should be measured")
assert.True(t, got.candidate > 0, "candidate latency should be measured")
case <-time.After(time.Second):
t.Fatal("response comparison did not run")
}
assert.NoError(t, handler.Shutdown(t.Context()))
}

func TestServeStopsAfterContextCancellation(t *testing.T) {
reference := httptest.NewServer(http.HandlerFunc(func(writer http.ResponseWriter, _ *http.Request) {
writer.WriteHeader(http.StatusNoContent)
Expand Down
24 changes: 19 additions & 5 deletions internal/ingress/response.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package ingress

import (
"net/http"
"time"

"github.com/alecthomas/errors"

Expand All @@ -12,26 +13,30 @@ import (
// One proxy invocation owns each writer, so capture state needs no synchronization.
type captureResponseWriter struct {
http.ResponseWriter
limit int
status int
body []byte
overflow bool
limit int
status int
body []byte
overflow bool
started time.Time
firstByte time.Time
}

func newCaptureResponseWriter(writer http.ResponseWriter, limit int) *captureResponseWriter {
return &captureResponseWriter{ResponseWriter: writer, limit: limit}
return &captureResponseWriter{ResponseWriter: writer, limit: limit, started: time.Now()}
}

func (w *captureResponseWriter) WriteHeader(status int) {
if w.status == 0 {
w.status = status
w.firstByte = time.Now()
}
w.ResponseWriter.WriteHeader(status)
}

func (w *captureResponseWriter) Write(data []byte) (int, error) {
if w.status == 0 {
w.status = http.StatusOK
w.firstByte = time.Now()
}
remaining := w.limit - len(w.body)
if remaining > 0 {
Expand All @@ -53,12 +58,21 @@ func (w *captureResponseWriter) Unwrap() http.ResponseWriter {
return w.ResponseWriter
}

// Latency returns the backend's time to first response byte, or zero if none was written.
func (w *captureResponseWriter) Latency() time.Duration {
if w.firstByte.IsZero() {
return 0
}
return w.firstByte.Sub(w.started)
}

// Response returns detached headers and body for asynchronous comparison.
func (w *captureResponseWriter) Response() comparison.Response {
return comparison.Response{
StatusCode: w.status,
Header: w.Header().Clone(),
Body: append([]byte(nil), w.body...),
Overflow: w.overflow,
Latency: w.Latency(),
}
}
Loading