etcd/rafthttp/transport.go

158 lines
3.2 KiB
Go
Raw Normal View History

package rafthttp
import (
2014-12-29 04:44:26 +03:00
"log"
"net/http"
2014-12-29 04:44:26 +03:00
"net/url"
"path"
"sync"
"github.com/coreos/etcd/etcdserver/stats"
"github.com/coreos/etcd/pkg/types"
"github.com/coreos/etcd/raft/raftpb"
"github.com/coreos/etcd/Godeps/_workspace/src/golang.org/x/net/context"
)
2014-12-29 04:44:26 +03:00
type Raft interface {
Process(ctx context.Context, m raftpb.Message) error
}
type Transporter interface {
Handler() http.Handler
Send(m []raftpb.Message)
AddPeer(id types.ID, urls []string)
RemovePeer(id types.ID)
UpdatePeer(id types.ID, urls []string)
Stop()
}
2014-12-29 23:20:52 +03:00
type transport struct {
roundTripper http.RoundTripper
id types.ID
clusterID types.ID
raft Raft
serverStats *stats.ServerStats
leaderStats *stats.LeaderStats
2015-01-03 07:00:29 +03:00
mu sync.RWMutex // protect the peer map
peers map[types.ID]*peer // remote peers
errorc chan error
}
2015-01-03 07:00:29 +03:00
func NewTransporter(rt http.RoundTripper, id, cid types.ID, r Raft, errorc chan error, ss *stats.ServerStats, ls *stats.LeaderStats) Transporter {
2014-12-29 23:20:52 +03:00
return &transport{
roundTripper: rt,
id: id,
clusterID: cid,
raft: r,
serverStats: ss,
leaderStats: ls,
peers: make(map[types.ID]*peer),
2015-01-03 07:00:29 +03:00
errorc: errorc,
}
2014-12-29 04:44:26 +03:00
}
2014-12-29 23:20:52 +03:00
func (t *transport) Handler() http.Handler {
h := NewHandler(t.raft, t.clusterID)
sh := NewStreamHandler(t, t.id, t.clusterID)
mux := http.NewServeMux()
mux.Handle(RaftPrefix, h)
mux.Handle(RaftStreamPrefix+"/", sh)
2014-12-29 04:44:26 +03:00
return mux
}
2014-12-29 23:20:52 +03:00
func (t *transport) Peer(id types.ID) *peer {
2014-12-29 04:44:26 +03:00
t.mu.RLock()
defer t.mu.RUnlock()
return t.peers[id]
}
2014-12-29 23:20:52 +03:00
func (t *transport) Send(msgs []raftpb.Message) {
2014-12-29 04:44:26 +03:00
for _, m := range msgs {
// intentionally dropped message
if m.To == 0 {
continue
}
to := types.ID(m.To)
p, ok := t.peers[to]
if !ok {
log.Printf("etcdserver: send message to unknown receiver %s", to)
continue
}
if m.Type == raftpb.MsgApp {
t.serverStats.SendAppendReq(m.Size())
2014-12-29 04:44:26 +03:00
}
p.Send(m)
}
}
2014-12-29 23:20:52 +03:00
func (t *transport) Stop() {
2014-12-29 04:44:26 +03:00
for _, p := range t.peers {
p.Stop()
}
if tr, ok := t.roundTripper.(*http.Transport); ok {
2014-12-29 04:44:26 +03:00
tr.CloseIdleConnections()
}
}
2014-12-29 23:20:52 +03:00
func (t *transport) AddPeer(id types.ID, urls []string) {
2014-12-29 04:44:26 +03:00
t.mu.Lock()
defer t.mu.Unlock()
if _, ok := t.peers[id]; ok {
return
}
// TODO: considering how to switch between all available peer urls
peerURL := urls[0]
u, err := url.Parse(peerURL)
if err != nil {
log.Panicf("unexpect peer url %s", peerURL)
}
2014-12-31 00:48:07 +03:00
u.Path = path.Join(u.Path, RaftPrefix)
fs := t.leaderStats.Follower(id.String())
2015-01-03 07:00:29 +03:00
t.peers[id] = NewPeer(t.roundTripper, u.String(), id, t.clusterID, t.raft, fs, t.errorc)
2014-12-29 04:44:26 +03:00
}
2014-12-29 23:20:52 +03:00
func (t *transport) RemovePeer(id types.ID) {
2014-12-29 04:44:26 +03:00
t.mu.Lock()
defer t.mu.Unlock()
t.peers[id].Stop()
delete(t.peers, id)
}
2014-12-29 23:20:52 +03:00
func (t *transport) UpdatePeer(id types.ID, urls []string) {
2014-12-29 04:44:26 +03:00
t.mu.Lock()
defer t.mu.Unlock()
// TODO: return error or just panic?
if _, ok := t.peers[id]; !ok {
return
}
peerURL := urls[0]
u, err := url.Parse(peerURL)
if err != nil {
log.Panicf("unexpect peer url %s", peerURL)
}
2014-12-31 00:48:07 +03:00
u.Path = path.Join(u.Path, RaftPrefix)
2014-12-29 04:44:26 +03:00
t.peers[id].Update(u.String())
}
2014-12-29 23:20:52 +03:00
type Pausable interface {
Pause()
Resume()
}
2014-12-29 04:44:26 +03:00
// for testing
2014-12-29 23:20:52 +03:00
func (t *transport) Pause() {
2014-12-29 04:44:26 +03:00
for _, p := range t.peers {
p.Pause()
}
}
2014-12-29 23:20:52 +03:00
func (t *transport) Resume() {
2014-12-29 04:44:26 +03:00
for _, p := range t.peers {
p.Resume()
}
}