diff --git a/internal/comparison/comparator.go b/internal/comparison/comparator.go index 54ff4af..516f5dd 100644 --- a/internal/comparison/comparator.go +++ b/internal/comparison/comparator.go @@ -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) { diff --git a/internal/comparison/comparator_test.go b/internal/comparison/comparator_test.go index a890c99..5fac8de 100644 --- a/internal/comparison/comparator_test.go +++ b/internal/comparison/comparator_test.go @@ -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((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((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((name) => name.toLowerCase()); diff --git a/internal/comparison/types.go b/internal/comparison/types.go index 23875ec..81aeb91 100644 --- a/internal/comparison/types.go +++ b/internal/comparison/types.go @@ -6,6 +6,7 @@ import ( "fmt" "net/http" "slices" + "time" ) // Outcome describes whether a response pair could be compared and matched. @@ -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. diff --git a/internal/ingress/handler_test.go b/internal/ingress/handler_test.go index 5f70d80..6721acf 100644 --- a/internal/ingress/handler_test.go +++ b/internal/ingress/handler_test.go @@ -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) diff --git a/internal/ingress/response.go b/internal/ingress/response.go index 49f4d91..211d023 100644 --- a/internal/ingress/response.go +++ b/internal/ingress/response.go @@ -2,6 +2,7 @@ package ingress import ( "net/http" + "time" "github.com/alecthomas/errors" @@ -12,19 +13,22 @@ 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) } @@ -32,6 +36,7 @@ func (w *captureResponseWriter) WriteHeader(status int) { 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 { @@ -53,6 +58,14 @@ 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{ @@ -60,5 +73,6 @@ func (w *captureResponseWriter) Response() comparison.Response { Header: w.Header().Clone(), Body: append([]byte(nil), w.body...), Overflow: w.overflow, + Latency: w.Latency(), } }