From adeb873fa00aa759f05c3a59ed6dfd0fe4e94d77 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Adri=C3=A1n=20Orive?= Date: Thu, 3 May 2018 15:35:05 +0200 Subject: [PATCH] fix Channel writes after close MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ensure Channel.Write returns ErrSinkClosed when called after the channel has been closed, even when the event channel has capacity to accept the event. Keep the close case in the blocking select so that a Write already waiting for a receiver is also unblocked when Close is called. Co-authored-by: Adrián Orive Signed-off-by: Adrián Orive Signed-off-by: Sebastiaan van Stijn --- channel.go | 9 +++++++-- channel_test.go | 39 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 46 insertions(+), 2 deletions(-) diff --git a/channel.go b/channel.go index 0eb7287..c58a2c3 100644 --- a/channel.go +++ b/channel.go @@ -34,10 +34,15 @@ func (ch *Channel) Done() chan struct{} { // the listener. func (ch *Channel) Write(event Event) error { select { - case ch.C <- event: - return nil case <-ch.closed: return ErrSinkClosed + default: + select { + case <-ch.closed: + return ErrSinkClosed + case ch.C <- event: + return nil + } } } diff --git a/channel_test.go b/channel_test.go index e8e9309..f7f9df9 100644 --- a/channel_test.go +++ b/channel_test.go @@ -86,3 +86,42 @@ loop: t.Fatalf("events did not make it through sink: %v != %v", received, nevents) } } + +// TestChannelWriteAfterClose is a regression test for +// https://github.com/docker/go-events/issues/29. Once Close has completed, +// subsequent writes must always return ErrSinkClosed, even if the channel has +// capacity to accept another event. +func TestChannelWriteAfterClose(t *testing.T) { + const nEvents = 100 + + sink := NewChannel(nEvents) + if err := sink.Close(); err != nil { + t.Fatal(err) + } + + for i := range nEvents { + if err := sink.Write(i); err != ErrSinkClosed { + t.Fatalf("Write(%d) error = %v, want %v", i, err, ErrSinkClosed) + } + } +} + +// TestChannelCloseUnblocksWrite verifies that closing a Channel releases a +// Write that is blocked waiting for a receiver, and that the blocked Write +// returns ErrSinkClosed. +func TestChannelCloseUnblocksWrite(t *testing.T) { + sink := NewChannel(0) + + errCh := make(chan error) + go func() { + errCh <- sink.Write("event") + }() + + if err := sink.Close(); err != nil { + t.Fatal(err) + } + + if err := <-errCh; err != ErrSinkClosed { + t.Fatalf("Write() error = %v, want %v", err, ErrSinkClosed) + } +}