syncthing/protocol/protocol.go

570 lines
12 KiB
Go
Raw Normal View History

2014-07-12 22:45:33 +00:00
// Copyright (C) 2014 Jakob Borg and Contributors (see the CONTRIBUTORS file).
// All rights reserved. Use of this source code is governed by an MIT-style
// license that can be found in the LICENSE file.
2014-06-01 20:50:14 +00:00
2013-12-15 10:43:31 +00:00
package protocol
import (
"bufio"
2013-12-15 10:43:31 +00:00
"errors"
"fmt"
2013-12-15 10:43:31 +00:00
"io"
"sync"
"time"
2014-02-15 11:08:55 +00:00
"github.com/calmh/syncthing/xdr"
2013-12-15 10:43:31 +00:00
)
2014-02-20 16:40:15 +00:00
const BlockSize = 128 * 1024
2013-12-15 10:43:31 +00:00
const (
messageTypeClusterConfig = 0
messageTypeIndex = 1
messageTypeRequest = 2
messageTypeResponse = 3
messageTypePing = 4
messageTypePong = 5
messageTypeIndexUpdate = 6
2014-07-26 19:27:55 +00:00
messageTypeClose = 7
2013-12-15 10:43:31 +00:00
)
const (
stateInitial = iota
stateCCRcvd
stateIdxRcvd
)
const (
2014-05-23 10:53:11 +00:00
FlagDeleted uint32 = 1 << 12
FlagInvalid = 1 << 13
FlagDirectory = 1 << 14
FlagNoPermBits = 1 << 15
)
const (
FlagShareTrusted uint32 = 1 << 0
FlagShareReadOnly = 1 << 1
FlagShareBits = 0x000000ff
)
var (
2014-02-24 12:29:30 +00:00
ErrClusterHash = fmt.Errorf("configuration error: mismatched cluster hash")
ErrClosed = errors.New("connection closed")
)
2013-12-15 10:43:31 +00:00
type Model interface {
// An index was received from the peer node
Index(nodeID NodeID, repo string, files []FileInfo)
2013-12-28 13:10:36 +00:00
// An index update was received from the peer node
IndexUpdate(nodeID NodeID, repo string, files []FileInfo)
2013-12-15 10:43:31 +00:00
// A request was made by the peer node
Request(nodeID NodeID, repo string, name string, offset int64, size int) ([]byte, error)
// A cluster configuration message was received
ClusterConfig(nodeID NodeID, config ClusterConfigMessage)
2013-12-15 10:43:31 +00:00
// The peer node closed the connection
Close(nodeID NodeID, err error)
2013-12-15 10:43:31 +00:00
}
type Connection interface {
ID() NodeID
Name() string
Index(repo string, files []FileInfo) error
IndexUpdate(repo string, files []FileInfo) error
Request(repo string, name string, offset int64, size int) ([]byte, error)
ClusterConfig(config ClusterConfigMessage)
Statistics() Statistics
}
type rawConnection struct {
id NodeID
name string
receiver Model
state int
2014-07-14 21:52:11 +00:00
cr *countingReader
xr *xdr.Reader
2014-07-14 21:52:11 +00:00
cw *countingWriter
wb *bufio.Writer
xw *xdr.Writer
2014-07-03 11:37:20 +00:00
awaiting []chan asyncResult
awaitingMut sync.Mutex
idxMut sync.Mutex // ensures serialization of Index calls
nextID chan int
outbox chan []encodable
closed chan struct{}
2014-07-03 11:37:20 +00:00
once sync.Once
2013-12-15 10:43:31 +00:00
}
2013-12-15 14:58:27 +00:00
type asyncResult struct {
val []byte
err error
}
2013-12-31 03:10:54 +00:00
const (
2014-06-14 09:07:34 +00:00
pingTimeout = 30 * time.Second
pingIdleTime = 60 * time.Second
2013-12-31 03:10:54 +00:00
)
func NewConnection(nodeID NodeID, reader io.Reader, writer io.Writer, receiver Model, name string) Connection {
2014-07-25 13:16:23 +00:00
// Byte counters are at the lowest level, counting compressed bytes
cr := &countingReader{Reader: reader}
cw := &countingWriter{Writer: writer}
2014-07-25 13:16:23 +00:00
// Compression is just above counting
zr := newLZ4Reader(cr)
zw := newLZ4Writer(cw)
// We buffer writes on top of compression.
// The LZ4 reader is already internally buffered
wb := bufio.NewWriterSize(zw, 65536)
2013-12-15 10:43:31 +00:00
c := rawConnection{
id: nodeID,
name: name,
receiver: nativeModel{receiver},
state: stateInitial,
cr: cr,
2014-07-25 13:16:23 +00:00
xr: xdr.NewReader(zr),
cw: cw,
wb: wb,
xw: xdr.NewWriter(wb),
awaiting: make([]chan asyncResult, 0x1000),
outbox: make(chan []encodable),
nextID: make(chan int),
closed: make(chan struct{}),
2013-12-15 10:43:31 +00:00
}
go c.readerLoop()
go c.writerLoop()
go c.pingerLoop()
go c.idGenerator()
2013-12-15 10:43:31 +00:00
return wireFormatConnection{&c}
2013-12-15 10:43:31 +00:00
}
func (c *rawConnection) ID() NodeID {
2014-01-09 12:58:35 +00:00
return c.id
}
func (c *rawConnection) Name() string {
return c.name
}
// Index writes the list of file information to the connected peer node
func (c *rawConnection) Index(repo string, idx []FileInfo) error {
select {
case <-c.closed:
return ErrClosed
default:
2013-12-28 13:10:36 +00:00
}
c.idxMut.Lock()
c.send(header{0, -1, messageTypeIndex}, IndexMessage{repo, idx})
c.idxMut.Unlock()
return nil
}
// IndexUpdate writes the list of file information to the connected peer node as an update
func (c *rawConnection) IndexUpdate(repo string, idx []FileInfo) error {
select {
case <-c.closed:
return ErrClosed
default:
}
c.idxMut.Lock()
c.send(header{0, -1, messageTypeIndexUpdate}, IndexMessage{repo, idx})
c.idxMut.Unlock()
return nil
2013-12-15 10:43:31 +00:00
}
// Request returns the bytes for the specified block after fetching them from the connected peer.
func (c *rawConnection) Request(repo string, name string, offset int64, size int) ([]byte, error) {
var id int
select {
case id = <-c.nextID:
case <-c.closed:
2013-12-28 15:30:02 +00:00
return nil, ErrClosed
}
2014-07-03 11:37:20 +00:00
c.awaitingMut.Lock()
if ch := c.awaiting[id]; ch != nil {
panic("id taken")
}
2014-07-03 11:37:20 +00:00
rc := make(chan asyncResult, 1)
c.awaiting[id] = rc
2014-07-03 11:37:20 +00:00
c.awaitingMut.Unlock()
ok := c.send(header{0, id, messageTypeRequest},
RequestMessage{repo, name, uint64(offset), uint32(size)})
if !ok {
return nil, ErrClosed
2013-12-19 23:01:34 +00:00
}
2013-12-15 10:43:31 +00:00
2013-12-15 14:58:27 +00:00
res, ok := <-rc
if !ok {
return nil, ErrClosed
2013-12-15 10:43:31 +00:00
}
2013-12-15 14:58:27 +00:00
return res.val, res.err
2013-12-15 10:43:31 +00:00
}
// ClusterConfig send the cluster configuration message to the peer and returns any error
func (c *rawConnection) ClusterConfig(config ClusterConfigMessage) {
c.send(header{0, -1, messageTypeClusterConfig}, config)
}
func (c *rawConnection) ping() bool {
var id int
select {
case id = <-c.nextID:
case <-c.closed:
2013-12-31 01:52:36 +00:00
return false
2013-12-28 15:30:02 +00:00
}
2013-12-30 20:27:20 +00:00
rc := make(chan asyncResult, 1)
2014-07-03 11:37:20 +00:00
c.awaitingMut.Lock()
c.awaiting[id] = rc
2014-07-03 11:37:20 +00:00
c.awaitingMut.Unlock()
ok := c.send(header{0, id, messageTypePing})
if !ok {
return false
2013-12-19 23:01:34 +00:00
}
2013-12-15 10:43:31 +00:00
2013-12-31 01:52:36 +00:00
res, ok := <-rc
return ok && res.err == nil
2013-12-15 10:43:31 +00:00
}
func (c *rawConnection) readerLoop() (err error) {
defer func() {
c.close(err)
}()
for {
select {
case <-c.closed:
return ErrClosed
default:
}
var hdr header
hdr.decodeXDR(c.xr)
if err := c.xr.Error(); err != nil {
return err
}
if hdr.version != 0 {
return fmt.Errorf("protocol error: %s: unknown message version %#x", c.id, hdr.version)
}
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
}
case messageTypePing:
c.send(header{0, hdr.msgID, messageTypePong})
case messageTypePong:
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
2014-07-26 19:27:55 +00:00
case messageTypeClose:
if err := c.handleClose(); err != nil {
return err
}
default:
return fmt.Errorf("protocol error: %s: unknown message type %#x", c.id, hdr.msgType)
}
}
2013-12-15 10:43:31 +00:00
}
func (c *rawConnection) handleIndex() error {
var im IndexMessage
im.decodeXDR(c.xr)
if err := c.xr.Error(); err != nil {
return err
} else {
2014-07-23 09:54:15 +00:00
if debug {
l.Debugf("Index(%v, %v, %d files)", c.id, im.Repository, len(im.Files))
2014-07-23 09:54:15 +00:00
}
c.receiver.Index(c.id, im.Repository, im.Files)
}
return nil
}
func (c *rawConnection) handleIndexUpdate() error {
var im IndexMessage
im.decodeXDR(c.xr)
if err := c.xr.Error(); err != nil {
return err
} else {
2014-07-23 09:54:15 +00:00
if debug {
l.Debugf("queueing IndexUpdate(%v, %v, %d files)", c.id, im.Repository, len(im.Files))
}
c.receiver.IndexUpdate(c.id, im.Repository, im.Files)
}
return nil
}
func (c *rawConnection) handleRequest(hdr header) error {
var req RequestMessage
req.decodeXDR(c.xr)
if err := c.xr.Error(); err != nil {
return err
}
go c.processRequest(hdr.msgID, req)
return nil
}
func (c *rawConnection) handleResponse(hdr header) error {
data := c.xr.ReadBytesMax(256 * 1024) // Sufficiently larger than max expected block size
if err := c.xr.Error(); err != nil {
return err
2013-12-15 10:43:31 +00:00
}
2014-07-03 11:37:20 +00:00
c.awaitingMut.Lock()
if rc := c.awaiting[hdr.msgID]; rc != nil {
c.awaiting[hdr.msgID] = nil
2014-07-03 11:37:20 +00:00
rc <- asyncResult{data, nil}
close(rc)
}
c.awaitingMut.Unlock()
2013-12-19 23:01:34 +00:00
return nil
2013-12-15 10:43:31 +00:00
}
func (c *rawConnection) handlePong(hdr header) {
2014-07-03 11:37:20 +00:00
c.awaitingMut.Lock()
if rc := c.awaiting[hdr.msgID]; rc != nil {
c.awaiting[hdr.msgID] = nil
2014-07-03 11:37:20 +00:00
rc <- asyncResult{}
close(rc)
}
2014-07-03 11:37:20 +00:00
c.awaitingMut.Unlock()
}
func (c *rawConnection) handleClusterConfig() error {
var cm ClusterConfigMessage
cm.decodeXDR(c.xr)
if err := c.xr.Error(); err != nil {
return err
} else {
go c.receiver.ClusterConfig(c.id, cm)
2013-12-15 10:43:31 +00:00
}
return nil
}
2013-12-15 14:58:27 +00:00
2014-07-26 19:27:55 +00:00
func (c *rawConnection) handleClose() error {
var cm CloseMessage
cm.decodeXDR(c.xr)
if err := c.xr.Error(); err != nil {
return err
}
return errors.New(cm.Reason)
}
type encodable interface {
encodeXDR(*xdr.Writer) (int, error)
2013-12-15 10:43:31 +00:00
}
type encodableBytes []byte
2013-12-15 10:43:31 +00:00
func (e encodableBytes) encodeXDR(xw *xdr.Writer) (int, error) {
return xw.WriteBytes(e)
2013-12-15 10:43:31 +00:00
}
func (c *rawConnection) send(h header, es ...encodable) bool {
if h.msgID < 0 {
select {
case id := <-c.nextID:
h.msgID = id
case <-c.closed:
return false
2013-12-21 07:06:54 +00:00
}
}
msg := append([]encodable{h}, es...)
2013-12-15 10:43:31 +00:00
select {
case c.outbox <- msg:
return true
case <-c.closed:
return false
}
}
2013-12-15 10:43:31 +00:00
func (c *rawConnection) writerLoop() {
for {
select {
case es := <-c.outbox:
for _, e := range es {
e.encodeXDR(c.xw)
}
if err := c.flush(); err != nil {
c.close(err)
return
}
case <-c.closed:
return
}
}
}
2013-12-15 10:43:31 +00:00
func (c *rawConnection) flush() error {
if err := c.xw.Error(); err != nil {
return err
}
if err := c.wb.Flush(); err != nil {
return err
}
return nil
}
2013-12-15 10:43:31 +00:00
func (c *rawConnection) close(err error) {
2014-07-03 11:37:20 +00:00
c.once.Do(func() {
close(c.closed)
2013-12-15 14:58:27 +00:00
2014-07-03 11:37:20 +00:00
c.awaitingMut.Lock()
for i, ch := range c.awaiting {
if ch != nil {
close(ch)
c.awaiting[i] = nil
2013-12-15 10:43:31 +00:00
}
}
2014-07-03 11:37:20 +00:00
c.awaitingMut.Unlock()
go c.receiver.Close(c.id, err)
2014-07-03 11:37:20 +00:00
})
2013-12-15 10:43:31 +00:00
}
func (c *rawConnection) idGenerator() {
nextID := 0
for {
nextID = (nextID + 1) & 0xfff
select {
case c.nextID <- nextID:
case <-c.closed:
return
}
2013-12-15 10:43:31 +00:00
}
}
func (c *rawConnection) pingerLoop() {
2013-12-31 01:52:36 +00:00
var rc = make(chan bool, 1)
ticker := time.Tick(pingIdleTime / 2)
2013-12-31 03:10:54 +00:00
for {
select {
case <-ticker:
if d := time.Since(c.xr.LastRead()); d < pingIdleTime {
if debug {
l.Debugln(c.id, "ping skipped after rd", d)
}
continue
}
if d := time.Since(c.xw.LastWrite()); d < pingIdleTime {
if debug {
l.Debugln(c.id, "ping skipped after wr", d)
}
continue
}
go func() {
if debug {
l.Debugln(c.id, "ping ->")
}
rc <- c.ping()
}()
select {
case ok := <-rc:
2014-05-28 18:45:29 +00:00
if debug {
l.Debugln(c.id, "<- pong")
}
if !ok {
c.close(fmt.Errorf("ping failure"))
}
case <-time.After(pingTimeout):
c.close(fmt.Errorf("ping timeout"))
case <-c.closed:
return
}
case <-c.closed:
return
}
}
}
2013-12-23 17:12:44 +00:00
func (c *rawConnection) processRequest(msgID int, req RequestMessage) {
data, _ := c.receiver.Request(c.id, req.Repository, req.Name, int64(req.Offset), int(req.Size))
2014-07-03 11:37:20 +00:00
c.send(header{0, msgID, messageTypeResponse}, encodableBytes(data))
}
2013-12-23 17:12:44 +00:00
type Statistics struct {
2014-01-05 15:16:37 +00:00
At time.Time
InBytesTotal uint64
OutBytesTotal uint64
2013-12-23 17:12:44 +00:00
}
func (c *rawConnection) Statistics() Statistics {
return Statistics{
2014-01-05 15:16:37 +00:00
At: time.Now(),
InBytesTotal: c.cr.Tot(),
OutBytesTotal: c.cw.Tot(),
2013-12-23 17:12:44 +00:00
}
}
2014-05-23 10:53:26 +00:00
func IsDeleted(bits uint32) bool {
return bits&FlagDeleted != 0
}
func IsInvalid(bits uint32) bool {
return bits&FlagInvalid != 0
}
func IsDirectory(bits uint32) bool {
return bits&FlagDirectory != 0
}
func HasPermissionBits(bits uint32) bool {
return bits&FlagNoPermBits == 0
}