2015-10-06 13:24:49 -04:00
|
|
|
package broadcaster
|
2015-08-11 13:12:47 -04:00
|
|
|
|
|
|
|
import (
|
|
|
|
"errors"
|
|
|
|
"io"
|
|
|
|
"sync"
|
|
|
|
)
|
|
|
|
|
2015-10-06 13:24:49 -04:00
|
|
|
// Buffered keeps track of one or more observers watching the progress
|
2015-08-11 13:12:47 -04:00
|
|
|
// of an operation. For example, if multiple clients are trying to pull an
|
2015-10-06 13:24:49 -04:00
|
|
|
// image, they share a Buffered struct for the download operation.
|
|
|
|
type Buffered struct {
|
2015-08-11 13:12:47 -04:00
|
|
|
sync.Mutex
|
|
|
|
// c is a channel that observers block on, waiting for the operation
|
|
|
|
// to finish.
|
|
|
|
c chan struct{}
|
|
|
|
// cond is a condition variable used to wake up observers when there's
|
|
|
|
// new data available.
|
|
|
|
cond *sync.Cond
|
|
|
|
// history is a buffer of the progress output so far, so a new observer
|
2015-08-28 13:09:00 -04:00
|
|
|
// can catch up. The history is stored as a slice of separate byte
|
|
|
|
// slices, so that if the writer is a WriteFlusher, the flushes will
|
|
|
|
// happen in the right places.
|
|
|
|
history [][]byte
|
2015-08-11 13:12:47 -04:00
|
|
|
// wg is a WaitGroup used to wait for all writes to finish on Close
|
|
|
|
wg sync.WaitGroup
|
2015-08-25 17:23:52 -04:00
|
|
|
// result is the argument passed to the first call of Close, and
|
|
|
|
// returned to callers of Wait
|
|
|
|
result error
|
2015-08-11 13:12:47 -04:00
|
|
|
}
|
|
|
|
|
2015-10-06 13:24:49 -04:00
|
|
|
// NewBuffered returns an initialized Buffered structure.
|
|
|
|
func NewBuffered() *Buffered {
|
|
|
|
b := &Buffered{
|
2015-08-11 13:12:47 -04:00
|
|
|
c: make(chan struct{}),
|
|
|
|
}
|
|
|
|
b.cond = sync.NewCond(b)
|
|
|
|
return b
|
|
|
|
}
|
|
|
|
|
|
|
|
// closed returns true if and only if the broadcaster has been closed
|
2015-10-06 13:24:49 -04:00
|
|
|
func (broadcaster *Buffered) closed() bool {
|
2015-08-11 13:12:47 -04:00
|
|
|
select {
|
|
|
|
case <-broadcaster.c:
|
|
|
|
return true
|
|
|
|
default:
|
|
|
|
return false
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// receiveWrites runs as a goroutine so that writes don't block the Write
|
|
|
|
// function. It writes the new data in broadcaster.history each time there's
|
|
|
|
// activity on the broadcaster.cond condition variable.
|
2015-10-06 13:24:49 -04:00
|
|
|
func (broadcaster *Buffered) receiveWrites(observer io.Writer) {
|
2015-08-11 13:12:47 -04:00
|
|
|
n := 0
|
|
|
|
|
|
|
|
broadcaster.Lock()
|
|
|
|
|
|
|
|
// The condition variable wait is at the end of this loop, so that the
|
|
|
|
// first iteration will write the history so far.
|
|
|
|
for {
|
2015-08-28 13:09:00 -04:00
|
|
|
newData := broadcaster.history[n:]
|
2015-08-11 13:12:47 -04:00
|
|
|
// Make a copy of newData so we can release the lock
|
2015-08-28 13:09:00 -04:00
|
|
|
sendData := make([][]byte, len(newData), len(newData))
|
2015-08-11 13:12:47 -04:00
|
|
|
copy(sendData, newData)
|
|
|
|
broadcaster.Unlock()
|
|
|
|
|
2015-08-28 13:09:00 -04:00
|
|
|
for len(sendData) > 0 {
|
|
|
|
_, err := observer.Write(sendData[0])
|
2015-08-11 13:12:47 -04:00
|
|
|
if err != nil {
|
|
|
|
broadcaster.wg.Done()
|
|
|
|
return
|
|
|
|
}
|
2015-08-28 13:09:00 -04:00
|
|
|
n++
|
|
|
|
sendData = sendData[1:]
|
2015-08-11 13:12:47 -04:00
|
|
|
}
|
|
|
|
|
|
|
|
broadcaster.Lock()
|
|
|
|
|
2015-09-10 17:58:06 -04:00
|
|
|
// If we are behind, we need to catch up instead of waiting
|
|
|
|
// or handling a closure.
|
|
|
|
if len(broadcaster.history) != n {
|
|
|
|
continue
|
|
|
|
}
|
|
|
|
|
2015-08-11 13:12:47 -04:00
|
|
|
// detect closure of the broadcast writer
|
|
|
|
if broadcaster.closed() {
|
|
|
|
broadcaster.Unlock()
|
|
|
|
broadcaster.wg.Done()
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2015-09-10 17:58:06 -04:00
|
|
|
broadcaster.cond.Wait()
|
2015-08-11 13:12:47 -04:00
|
|
|
|
|
|
|
// Mutex is still locked as the loop continues
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// Write adds data to the history buffer, and also writes it to all current
|
|
|
|
// observers.
|
2015-10-06 13:24:49 -04:00
|
|
|
func (broadcaster *Buffered) Write(p []byte) (n int, err error) {
|
2015-08-11 13:12:47 -04:00
|
|
|
broadcaster.Lock()
|
|
|
|
defer broadcaster.Unlock()
|
|
|
|
|
|
|
|
// Is the broadcaster closed? If so, the write should fail.
|
|
|
|
if broadcaster.closed() {
|
2015-10-06 13:24:49 -04:00
|
|
|
return 0, errors.New("attempted write to a closed broadcaster.Buffered")
|
2015-08-11 13:12:47 -04:00
|
|
|
}
|
|
|
|
|
2015-08-28 13:09:00 -04:00
|
|
|
// Add message in p to the history slice
|
|
|
|
newEntry := make([]byte, len(p), len(p))
|
|
|
|
copy(newEntry, p)
|
|
|
|
broadcaster.history = append(broadcaster.history, newEntry)
|
|
|
|
|
2015-08-11 13:12:47 -04:00
|
|
|
broadcaster.cond.Broadcast()
|
|
|
|
|
|
|
|
return len(p), nil
|
|
|
|
}
|
|
|
|
|
2015-10-06 13:24:49 -04:00
|
|
|
// Add adds an observer to the broadcaster. The new observer receives the
|
2015-08-11 13:12:47 -04:00
|
|
|
// data from the history buffer, and also all subsequent data.
|
2015-10-06 13:24:49 -04:00
|
|
|
func (broadcaster *Buffered) Add(w io.Writer) error {
|
2015-08-11 13:12:47 -04:00
|
|
|
// The lock is acquired here so that Add can't race with Close
|
|
|
|
broadcaster.Lock()
|
|
|
|
defer broadcaster.Unlock()
|
|
|
|
|
|
|
|
if broadcaster.closed() {
|
2015-10-06 13:24:49 -04:00
|
|
|
return errors.New("attempted to add observer to a closed broadcaster.Buffered")
|
2015-08-11 13:12:47 -04:00
|
|
|
}
|
|
|
|
|
|
|
|
broadcaster.wg.Add(1)
|
|
|
|
go broadcaster.receiveWrites(w)
|
|
|
|
|
|
|
|
return nil
|
|
|
|
}
|
|
|
|
|
2015-08-25 17:23:52 -04:00
|
|
|
// CloseWithError signals to all observers that the operation has finished. Its
|
|
|
|
// argument is a result that should be returned to waiters blocking on Wait.
|
2015-10-06 13:24:49 -04:00
|
|
|
func (broadcaster *Buffered) CloseWithError(result error) {
|
2015-08-11 13:12:47 -04:00
|
|
|
broadcaster.Lock()
|
2015-09-10 13:08:02 -04:00
|
|
|
if broadcaster.closed() {
|
2015-08-11 13:12:47 -04:00
|
|
|
broadcaster.Unlock()
|
|
|
|
return
|
|
|
|
}
|
2015-08-25 17:23:52 -04:00
|
|
|
broadcaster.result = result
|
2015-08-11 13:12:47 -04:00
|
|
|
close(broadcaster.c)
|
|
|
|
broadcaster.cond.Broadcast()
|
|
|
|
broadcaster.Unlock()
|
|
|
|
|
2015-08-25 17:23:52 -04:00
|
|
|
// Don't return until all writers have caught up.
|
2015-08-11 13:12:47 -04:00
|
|
|
broadcaster.wg.Wait()
|
|
|
|
}
|
|
|
|
|
2015-08-25 17:23:52 -04:00
|
|
|
// Close signals to all observers that the operation has finished. It causes
|
|
|
|
// all calls to Wait to return nil.
|
2015-10-06 13:24:49 -04:00
|
|
|
func (broadcaster *Buffered) Close() {
|
2015-08-25 17:23:52 -04:00
|
|
|
broadcaster.CloseWithError(nil)
|
|
|
|
}
|
|
|
|
|
2015-09-10 13:08:02 -04:00
|
|
|
// Wait blocks until the operation is marked as completed by the Close method,
|
|
|
|
// and all writer goroutines have completed. It returns the argument that was
|
|
|
|
// passed to Close.
|
2015-10-06 13:24:49 -04:00
|
|
|
func (broadcaster *Buffered) Wait() error {
|
2015-08-11 13:12:47 -04:00
|
|
|
<-broadcaster.c
|
2015-09-10 13:08:02 -04:00
|
|
|
broadcaster.wg.Wait()
|
2015-08-25 17:23:52 -04:00
|
|
|
return broadcaster.result
|
2015-08-11 13:12:47 -04:00
|
|
|
}
|