-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathserver.go
More file actions
147 lines (130 loc) · 2.83 KB
/
Copy pathserver.go
File metadata and controls
147 lines (130 loc) · 2.83 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
package raft
import (
"fmt"
"log"
"net"
"net/rpc"
"sync"
)
type IServer interface {
// Call makes an RPC using the provided service method
Call(id int, service string, args interface{}, res interface{}) error
ConnectToPeer(peerId int, addr net.Addr) error
Serve()
GetListenAddr() net.Addr
}
type Server struct {
mu sync.Mutex
me int
cm *CnsModule
enabled map[interface{}]bool
ready <-chan interface{}
quit chan interface{}
wg sync.WaitGroup
peers map[int]*rpc.Client
//cm CnsModule
peerIds []int
listener net.Listener
rpcServer *rpc.Server
rpcProxy interface{}
commitChan chan<- CommitEntry
storage Persistence
}
func (s *Server) GetListenAddr() net.Addr {
return s.listener.Addr()
}
func (s *Server) Serve() {
s.mu.Lock()
s.cm = NewConsensusModule(s.me, s.peerIds, s, s.ready, s.commitChan, s.storage)
// Create a new RPC server and register a RPCProxy that forwards all methods
// to n.cm
s.rpcServer = rpc.NewServer()
s.rpcProxy = &Proxy{cm: s.cm}
s.rpcServer.RegisterName("CnsModule", s.rpcProxy)
var err error
s.listener, err = net.Listen("tcp", ":0")
if err != nil {
log.Fatal(err)
}
log.Printf("[Server-%v] is listening on port %s", s.me, s.listener.Addr())
s.mu.Unlock()
s.wg.Add(1)
go func() {
defer s.wg.Done()
s.listen()
}()
}
func (s *Server) listen() {
for {
conn, err := s.listener.Accept()
if err != nil {
select {
case <-s.quit:
return
default:
log.Fatal("accept error:", err)
}
}
s.wg.Add(1)
go func() {
s.rpcServer.ServeConn(conn)
s.wg.Done()
}()
}
}
func (s *Server) ConnectToPeer(peerId int, addr net.Addr) error {
s.mu.Lock()
defer s.mu.Unlock()
if s.peers[peerId] == nil {
client, err := rpc.Dial(addr.Network(), addr.String())
if err != nil {
return err
}
s.peers[peerId] = client
}
return nil
}
func (s *Server) Call(id int, service string, args interface{}, res interface{}) error {
s.mu.Lock()
peer := s.peers[id]
s.mu.Unlock()
if peer == nil {
return fmt.Errorf("call client %d after it's closed", id)
} else {
if err := peer.Call(service, args, res); err != nil {
return err
}
}
return nil
}
func (s *Server) DisconnectAllPeers() {
s.mu.Lock()
defer s.mu.Unlock()
for id := range s.peers {
if s.peers[id] != nil {
s.peers[id].Close()
s.peers[id] = nil
}
}
}
func (s *Server) DisconnectPeer(peerId int) error {
s.mu.Lock()
defer s.mu.Unlock()
if s.peers[peerId] != nil {
err := s.peers[peerId].Close()
s.peers[peerId] = nil
return err
}
return nil
}
func NewServer(serverID int, peerIds []int, ready <-chan interface{}, commitChan chan<- CommitEntry, kv Persistence) *Server {
s := new(Server)
s.me = serverID
s.peerIds = peerIds
s.peers = make(map[int]*rpc.Client)
s.ready = ready
s.quit = make(chan interface{})
s.commitChan = commitChan
s.storage = kv
return s
}