|
|
@@ -28,6 +28,12 @@ const (
|
|
|
messageTypeIndexUpdate = 6
|
|
|
)
|
|
|
|
|
|
+const (
|
|
|
+ stateInitial = iota
|
|
|
+ stateCCRcvd
|
|
|
+ stateIdxRcvd
|
|
|
+)
|
|
|
+
|
|
|
const (
|
|
|
FlagDeleted uint32 = 1 << 12
|
|
|
FlagInvalid = 1 << 13
|
|
|
@@ -70,26 +76,27 @@ type Connection interface {
|
|
|
type rawConnection struct {
|
|
|
id NodeID
|
|
|
receiver Model
|
|
|
+ state int
|
|
|
|
|
|
reader io.ReadCloser
|
|
|
cr *countingReader
|
|
|
xr *xdr.Reader
|
|
|
- writer io.WriteCloser
|
|
|
|
|
|
- cw *countingWriter
|
|
|
- wb *bufio.Writer
|
|
|
- xw *xdr.Writer
|
|
|
- wmut sync.Mutex
|
|
|
+ writer io.WriteCloser
|
|
|
+ cw *countingWriter
|
|
|
+ wb *bufio.Writer
|
|
|
+ xw *xdr.Writer
|
|
|
|
|
|
- indexSent map[string]map[string]uint64
|
|
|
- awaiting []chan asyncResult
|
|
|
- imut sync.Mutex
|
|
|
+ awaiting []chan asyncResult
|
|
|
+ awaitingMut sync.Mutex
|
|
|
|
|
|
- idxMut sync.Mutex // ensures serialization of Index calls
|
|
|
+ idxSent map[string]map[string]uint64
|
|
|
+ idxMut sync.Mutex // ensures serialization of Index calls
|
|
|
|
|
|
nextID chan int
|
|
|
outbox chan []encodable
|
|
|
closed chan struct{}
|
|
|
+ once sync.Once
|
|
|
}
|
|
|
|
|
|
type asyncResult struct {
|
|
|
@@ -114,20 +121,21 @@ func NewConnection(nodeID NodeID, reader io.Reader, writer io.Writer, receiver M
|
|
|
wb := bufio.NewWriter(flwr)
|
|
|
|
|
|
c := rawConnection{
|
|
|
- id: nodeID,
|
|
|
- receiver: nativeModel{receiver},
|
|
|
- reader: flrd,
|
|
|
- cr: cr,
|
|
|
- xr: xdr.NewReader(flrd),
|
|
|
- writer: flwr,
|
|
|
- cw: cw,
|
|
|
- wb: wb,
|
|
|
- xw: xdr.NewWriter(wb),
|
|
|
- awaiting: make([]chan asyncResult, 0x1000),
|
|
|
- indexSent: make(map[string]map[string]uint64),
|
|
|
- outbox: make(chan []encodable),
|
|
|
- nextID: make(chan int),
|
|
|
- closed: make(chan struct{}),
|
|
|
+ id: nodeID,
|
|
|
+ receiver: nativeModel{receiver},
|
|
|
+ state: stateInitial,
|
|
|
+ reader: flrd,
|
|
|
+ cr: cr,
|
|
|
+ xr: xdr.NewReader(flrd),
|
|
|
+ writer: flwr,
|
|
|
+ cw: cw,
|
|
|
+ wb: wb,
|
|
|
+ xw: xdr.NewWriter(wb),
|
|
|
+ awaiting: make([]chan asyncResult, 0x1000),
|
|
|
+ idxSent: make(map[string]map[string]uint64),
|
|
|
+ outbox: make(chan []encodable),
|
|
|
+ nextID: make(chan int),
|
|
|
+ closed: make(chan struct{}),
|
|
|
}
|
|
|
|
|
|
go c.indexSerializerLoop()
|
|
|
@@ -148,31 +156,29 @@ func (c *rawConnection) Index(repo string, idx []FileInfo) {
|
|
|
c.idxMut.Lock()
|
|
|
defer c.idxMut.Unlock()
|
|
|
|
|
|
- c.imut.Lock()
|
|
|
var msgType int
|
|
|
- if c.indexSent[repo] == nil {
|
|
|
+ if c.idxSent[repo] == nil {
|
|
|
// This is the first time we send an index.
|
|
|
msgType = messageTypeIndex
|
|
|
|
|
|
- c.indexSent[repo] = make(map[string]uint64)
|
|
|
+ c.idxSent[repo] = make(map[string]uint64)
|
|
|
for _, f := range idx {
|
|
|
- c.indexSent[repo][f.Name] = f.Version
|
|
|
+ c.idxSent[repo][f.Name] = f.Version
|
|
|
}
|
|
|
} else {
|
|
|
// We have sent one full index. Only send updates now.
|
|
|
msgType = messageTypeIndexUpdate
|
|
|
var diff []FileInfo
|
|
|
for _, f := range idx {
|
|
|
- if vs, ok := c.indexSent[repo][f.Name]; !ok || f.Version != vs {
|
|
|
+ if vs, ok := c.idxSent[repo][f.Name]; !ok || f.Version != vs {
|
|
|
diff = append(diff, f)
|
|
|
- c.indexSent[repo][f.Name] = f.Version
|
|
|
+ c.idxSent[repo][f.Name] = f.Version
|
|
|
}
|
|
|
}
|
|
|
idx = diff
|
|
|
}
|
|
|
- c.imut.Unlock()
|
|
|
|
|
|
- if len(idx) > 0 {
|
|
|
+ if msgType == messageTypeIndex || len(idx) > 0 {
|
|
|
c.send(header{0, -1, msgType}, IndexMessage{repo, idx})
|
|
|
}
|
|
|
}
|
|
|
@@ -186,13 +192,13 @@ func (c *rawConnection) Request(repo string, name string, offset int64, size int
|
|
|
return nil, ErrClosed
|
|
|
}
|
|
|
|
|
|
- c.imut.Lock()
|
|
|
+ c.awaitingMut.Lock()
|
|
|
if ch := c.awaiting[id]; ch != nil {
|
|
|
panic("id taken")
|
|
|
}
|
|
|
- rc := make(chan asyncResult)
|
|
|
+ rc := make(chan asyncResult, 1)
|
|
|
c.awaiting[id] = rc
|
|
|
- c.imut.Unlock()
|
|
|
+ c.awaitingMut.Unlock()
|
|
|
|
|
|
ok := c.send(header{0, id, messageTypeRequest},
|
|
|
RequestMessage{repo, name, uint64(offset), uint32(size)})
|
|
|
@@ -221,9 +227,9 @@ func (c *rawConnection) ping() bool {
|
|
|
}
|
|
|
|
|
|
rc := make(chan asyncResult, 1)
|
|
|
- c.imut.Lock()
|
|
|
+ c.awaitingMut.Lock()
|
|
|
c.awaiting[id] = rc
|
|
|
- c.imut.Unlock()
|
|
|
+ c.awaitingMut.Unlock()
|
|
|
|
|
|
ok := c.send(header{0, id, messageTypePing})
|
|
|
if !ok {
|
|
|
@@ -257,21 +263,34 @@ func (c *rawConnection) readerLoop() (err error) {
|
|
|
|
|
|
switch hdr.msgType {
|
|
|
case messageTypeIndex:
|
|
|
+ if c.state < stateCCRcvd {
|
|
|
+ return fmt.Errorf("protocol error: index message in state %d", c.state)
|
|
|
+ }
|
|
|
if err := c.handleIndex(); err != nil {
|
|
|
return err
|
|
|
}
|
|
|
+ c.state = stateIdxRcvd
|
|
|
|
|
|
case messageTypeIndexUpdate:
|
|
|
+ if c.state < stateIdxRcvd {
|
|
|
+ return fmt.Errorf("protocol error: index update message in state %d", c.state)
|
|
|
+ }
|
|
|
if err := c.handleIndexUpdate(); err != nil {
|
|
|
return err
|
|
|
}
|
|
|
|
|
|
case messageTypeRequest:
|
|
|
+ if c.state < stateIdxRcvd {
|
|
|
+ return fmt.Errorf("protocol error: request message in state %d", c.state)
|
|
|
+ }
|
|
|
if err := c.handleRequest(hdr); err != nil {
|
|
|
return err
|
|
|
}
|
|
|
|
|
|
case messageTypeResponse:
|
|
|
+ if c.state < stateIdxRcvd {
|
|
|
+ return fmt.Errorf("protocol error: response message in state %d", c.state)
|
|
|
+ }
|
|
|
if err := c.handleResponse(hdr); err != nil {
|
|
|
return err
|
|
|
}
|
|
|
@@ -283,9 +302,13 @@ func (c *rawConnection) readerLoop() (err error) {
|
|
|
c.handlePong(hdr)
|
|
|
|
|
|
case messageTypeClusterConfig:
|
|
|
+ if c.state != stateInitial {
|
|
|
+ return fmt.Errorf("protocol error: cluster config message in state %d", c.state)
|
|
|
+ }
|
|
|
if err := c.handleClusterConfig(); err != nil {
|
|
|
return err
|
|
|
}
|
|
|
+ c.state = stateCCRcvd
|
|
|
|
|
|
default:
|
|
|
return fmt.Errorf("protocol error: %s: unknown message type %#x", c.id, hdr.msgType)
|
|
|
@@ -309,11 +332,16 @@ func (c *rawConnection) indexSerializerLoop() {
|
|
|
// large index update from the other side. But we must also ensure to
|
|
|
// process the indexes in the order they are received, hence the separate
|
|
|
// routine and buffered channel.
|
|
|
- for ii := range incomingIndexes {
|
|
|
- if ii.update {
|
|
|
- c.receiver.IndexUpdate(ii.id, ii.repo, ii.files)
|
|
|
- } else {
|
|
|
- c.receiver.Index(ii.id, ii.repo, ii.files)
|
|
|
+ for {
|
|
|
+ select {
|
|
|
+ case ii := <-incomingIndexes:
|
|
|
+ if ii.update {
|
|
|
+ c.receiver.IndexUpdate(ii.id, ii.repo, ii.files)
|
|
|
+ } else {
|
|
|
+ c.receiver.Index(ii.id, ii.repo, ii.files)
|
|
|
+ }
|
|
|
+ case <-c.closed:
|
|
|
+ return
|
|
|
}
|
|
|
}
|
|
|
}
|
|
|
@@ -365,32 +393,25 @@ func (c *rawConnection) handleResponse(hdr header) error {
|
|
|
return err
|
|
|
}
|
|
|
|
|
|
- go func(hdr header, err error) {
|
|
|
- c.imut.Lock()
|
|
|
- rc := c.awaiting[hdr.msgID]
|
|
|
+ c.awaitingMut.Lock()
|
|
|
+ if rc := c.awaiting[hdr.msgID]; rc != nil {
|
|
|
c.awaiting[hdr.msgID] = nil
|
|
|
- c.imut.Unlock()
|
|
|
-
|
|
|
- if rc != nil {
|
|
|
- rc <- asyncResult{data, err}
|
|
|
- close(rc)
|
|
|
- }
|
|
|
- }(hdr, c.xr.Error())
|
|
|
+ rc <- asyncResult{data, nil}
|
|
|
+ close(rc)
|
|
|
+ }
|
|
|
+ c.awaitingMut.Unlock()
|
|
|
|
|
|
return nil
|
|
|
}
|
|
|
|
|
|
func (c *rawConnection) handlePong(hdr header) {
|
|
|
- c.imut.Lock()
|
|
|
+ c.awaitingMut.Lock()
|
|
|
if rc := c.awaiting[hdr.msgID]; rc != nil {
|
|
|
- go func() {
|
|
|
- rc <- asyncResult{}
|
|
|
- close(rc)
|
|
|
- }()
|
|
|
-
|
|
|
c.awaiting[hdr.msgID] = nil
|
|
|
+ rc <- asyncResult{}
|
|
|
+ close(rc)
|
|
|
}
|
|
|
- c.imut.Unlock()
|
|
|
+ c.awaitingMut.Unlock()
|
|
|
}
|
|
|
|
|
|
func (c *rawConnection) handleClusterConfig() error {
|
|
|
@@ -434,18 +455,20 @@ func (c *rawConnection) send(h header, es ...encodable) bool {
|
|
|
|
|
|
func (c *rawConnection) writerLoop() {
|
|
|
var err error
|
|
|
- for es := range c.outbox {
|
|
|
- c.wmut.Lock()
|
|
|
- for _, e := range es {
|
|
|
- e.encodeXDR(c.xw)
|
|
|
- }
|
|
|
+ for {
|
|
|
+ select {
|
|
|
+ case es := <-c.outbox:
|
|
|
+ for _, e := range es {
|
|
|
+ e.encodeXDR(c.xw)
|
|
|
+ }
|
|
|
|
|
|
- if err = c.flush(); err != nil {
|
|
|
- c.wmut.Unlock()
|
|
|
- c.close(err)
|
|
|
+ if err = c.flush(); err != nil {
|
|
|
+ c.close(err)
|
|
|
+ return
|
|
|
+ }
|
|
|
+ case <-c.closed:
|
|
|
return
|
|
|
}
|
|
|
- c.wmut.Unlock()
|
|
|
}
|
|
|
}
|
|
|
|
|
|
@@ -470,29 +493,20 @@ func (c *rawConnection) flush() error {
|
|
|
}
|
|
|
|
|
|
func (c *rawConnection) close(err error) {
|
|
|
- c.imut.Lock()
|
|
|
- c.wmut.Lock()
|
|
|
- defer c.imut.Unlock()
|
|
|
- defer c.wmut.Unlock()
|
|
|
-
|
|
|
- select {
|
|
|
- case <-c.closed:
|
|
|
- return
|
|
|
- default:
|
|
|
+ c.once.Do(func() {
|
|
|
close(c.closed)
|
|
|
|
|
|
+ c.awaitingMut.Lock()
|
|
|
for i, ch := range c.awaiting {
|
|
|
if ch != nil {
|
|
|
close(ch)
|
|
|
c.awaiting[i] = nil
|
|
|
}
|
|
|
}
|
|
|
-
|
|
|
- c.writer.Close()
|
|
|
- c.reader.Close()
|
|
|
+ c.awaitingMut.Unlock()
|
|
|
|
|
|
go c.receiver.Close(c.id, err)
|
|
|
- }
|
|
|
+ })
|
|
|
}
|
|
|
|
|
|
func (c *rawConnection) idGenerator() {
|
|
|
@@ -554,8 +568,7 @@ func (c *rawConnection) pingerLoop() {
|
|
|
func (c *rawConnection) processRequest(msgID int, req RequestMessage) {
|
|
|
data, _ := c.receiver.Request(c.id, req.Repository, req.Name, int64(req.Offset), int(req.Size))
|
|
|
|
|
|
- c.send(header{0, msgID, messageTypeResponse},
|
|
|
- encodableBytes(data))
|
|
|
+ c.send(header{0, msgID, messageTypeResponse}, encodableBytes(data))
|
|
|
}
|
|
|
|
|
|
type Statistics struct {
|