From 6773d088c3e82a4e4c51c1f9a7d7502249d9dc95 Mon Sep 17 00:00:00 2001 From: MiMoCode Date: Fri, 10 Jul 2026 19:20:09 +0800 Subject: [PATCH] fix: snapshot queued log entries --- logger/queue.go | 14 ++++++++++-- tests/logger/queue_test.go | 44 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 56 insertions(+), 2 deletions(-) diff --git a/logger/queue.go b/logger/queue.go index 6275de0..eab2870 100644 --- a/logger/queue.go +++ b/logger/queue.go @@ -235,7 +235,8 @@ func (q *Queue) Submit(e *LogEntry) { q.dropped.Add(1) return } - size := EstimatedBytes(e) + entry := cloneLogEntry(e) + size := EstimatedBytes(entry) if !q.mu.TryRLock() { q.dropped.Add(1) return @@ -246,7 +247,7 @@ func (q *Queue) Submit(e *LogEntry) { return } select { - case q.ch <- queuedEntry{entry: e, size: size}: + case q.ch <- queuedEntry{entry: entry, size: size}: q.enq.Add(1) default: q.release(size) @@ -255,6 +256,15 @@ func (q *Queue) Submit(e *LogEntry) { q.mu.RUnlock() } +func cloneLogEntry(entry *LogEntry) *LogEntry { + clone := *entry + clone.RequestHeaders = append([]byte(nil), entry.RequestHeaders...) + clone.RequestBody = append([]byte(nil), entry.RequestBody...) + clone.ResponseHeaders = append([]byte(nil), entry.ResponseHeaders...) + clone.ResponseBody = append([]byte(nil), entry.ResponseBody...) + return &clone +} + func (q *Queue) reserve(size int64) bool { for { used := q.entries.Load() diff --git a/tests/logger/queue_test.go b/tests/logger/queue_test.go index 90e3ac8..f9dbf82 100644 --- a/tests/logger/queue_test.go +++ b/tests/logger/queue_test.go @@ -314,6 +314,50 @@ func TestQueueReleasesSubmittedSizeAfterEntryMutation(t *testing.T) { } } +func TestQueueSubmitOwnsEntrySnapshot(t *testing.T) { + prepareStarted := make(chan struct{}) + releasePrepare := make(chan struct{}) + batch := &fakeBatch{} + backend := newQueueBackend(func(context.Context, string) (logger.Batch, error) { + close(prepareStarted) + <-releasePrepare + return batch, nil + }) + entry := &logger.LogEntry{ + RequestID: "original-id", + RequestHeaders: []byte("original-request-headers"), + RequestBody: []byte("original-request-body"), + ResponseHeaders: []byte("original-response-headers"), + ResponseBody: []byte("original-response-body"), + } + q := logger.NewQueueWithBackend(backend, 1, 1, 1, time.Hour) + q.Start(context.Background()) + q.Submit(entry) + <-prepareStarted + + entry.RequestID = "mutated-id" + entry.RequestHeaders[0] = 'X' + entry.RequestBody[0] = 'X' + entry.ResponseHeaders[0] = 'X' + entry.ResponseBody[0] = 'X' + close(releasePrepare) + + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + if err := q.Stop(ctx); err != nil { + t.Fatalf("Stop: %v", err) + } + if len(batch.rows) != 1 { + t.Fatalf("written rows=%d want 1", len(batch.rows)) + } + row := batch.rows[0] + if row[0] != "original-id" || + row[5] != "original-request-headers" || row[6] != "original-request-body" || + row[9] != "original-response-headers" || row[10] != "original-response-body" { + t.Fatalf("queued entry changed after Submit: %#v", row) + } +} + func TestQueueCanceledStopEventuallyReleasesAllBudget(t *testing.T) { entry := &logger.LogEntry{RequestID: "queued"} q := logger.NewQueue(nil, 32, 32, 1, time.Hour, 32*logger.EstimatedBytes(entry))