diff --git a/internal/log/nonblocking_writer.go b/internal/log/nonblocking_writer.go index 269a9766..aa3fc633 100644 --- a/internal/log/nonblocking_writer.go +++ b/internal/log/nonblocking_writer.go @@ -26,7 +26,6 @@ const defaultNonBlockingWriterBufferSize = 1024 type nonBlockingWriter struct { writer io.Writer queue chan []byte - stop chan struct{} workerDone chan struct{} mu sync.RWMutex closed bool @@ -40,27 +39,14 @@ func newNonBlockingWriter(writer io.Writer, queueSize int) *nonBlockingWriter { w := &nonBlockingWriter{ writer: writer, queue: make(chan []byte, queueSize), - stop: make(chan struct{}), workerDone: make(chan struct{}), } go func() { defer close(w.workerDone) - for { - select { - case p := <-w.queue: - _, _ = w.writer.Write(p) - case <-w.stop: - for { - select { - case p := <-w.queue: - _, _ = w.writer.Write(p) - default: - return - } - } - } + for p := range w.queue { + _, _ = w.writer.Write(p) } }() @@ -89,7 +75,7 @@ func (w *nonBlockingWriter) Close() error { w.mu.Lock() if !w.closed { w.closed = true - close(w.stop) + close(w.queue) } w.mu.Unlock() diff --git a/internal/log/nonblocking_writer_test.go b/internal/log/nonblocking_writer_test.go index a78edca5..83d951db 100644 --- a/internal/log/nonblocking_writer_test.go +++ b/internal/log/nonblocking_writer_test.go @@ -111,4 +111,6 @@ func TestNonBlockingWriterDropsInsteadOfBlocking(t *testing.T) { require.Len(t, writes, 2) require.Equal(t, []byte("first"), writes[0]) require.Equal(t, []byte("second"), writes[1]) + require.NotEqual(t, []byte("third"), writes[0]) + require.NotEqual(t, []byte("third"), writes[1]) }