-
Notifications
You must be signed in to change notification settings - Fork 2
/
Copy pathpaxos.py
251 lines (187 loc) · 6.46 KB
/
paxos.py
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
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
from sys import stdin
import asyncio
import pickle
from datetime import datetime
loop = asyncio.get_event_loop()
async def async_suppress_exception(coro):
try:
return await coro
except ConnectionError:
pass
def print_debug(*args, **kwargs):
print(datetime.now(), *args, **kwargs)
def rpc_decode(s):
return pickle.loads(s)
def rpc_encode(o):
return pickle.dumps(o)
class RPC:
def __init__(self):
self._cb_table = {}
async def _server_handle(self, reader, writer):
try:
data = await reader.read()
cmd = rpc_decode(data)
event, args, kwargs = cmd
cb = self._cb_table[event]
ret = cb(*args, **kwargs) if not asyncio.iscoroutine(cb) else await cb(*args, **kwargs)
s = rpc_encode(ret)
writer.write(s)
await writer.drain()
writer.close()
print_debug('Receive:', cmd, 'Reply:', ret)
except ConnectionError:
pass
def schedule_server(self, host, port):
coro = asyncio.start_server(self._server_handle, host, port, loop=loop)
asyncio.ensure_future(coro)
print_debug('Start:', host, port)
def add_event_listener(self, event, cb):
self._cb_table[event] = cb
class RPCProxy:
def __init__(self, host, port):
self._host = host
self._port = port
def __getattr__(self, name):
async def call(host, port, event, *args, **kwargs):
reader, writer = await asyncio.open_connection(host, port, loop=loop)
cmd = (event, args, kwargs)
s = rpc_encode(cmd)
writer.write(s)
writer.write_eof()
data = await reader.read()
writer.close()
ret = rpc_decode(data)
print_debug('Request:', cmd, 'Receive:', ret)
return ret
def func(*args, **kwargs):
return call(self._host, self._port, name, *args, **kwargs)
return func
# --- RPC ---
class Storage:
def __init__(self, filename):
self._filename = filename
def load(self):
with open(self._filename, 'rb') as f:
ret = pickle.load(f)
return ret
def dump(self, o):
with open(self._filename, 'wb') as f:
pickle.dump(o, f)
# --- Storage ---
class Paxos:
def __init__(self, host, port, servers, storage, rpc):
self._host = host
self._port = port
self._servers = servers
self._storage = storage
self._over_half_num = len(servers) // 2 + 1
rpc.add_event_listener('prepare', self._on_prepare)
rpc.add_event_listener('accept', self._on_accept)
rpc.add_event_listener('learn', self._on_learn)
# status
try:
self._seq, self._proposal_seq, self._proposal_val, self._proposal_unanimous = self._storage.load()
except FileNotFoundError:
self._seq = (0, self._host, self._port)
self._proposal_seq = self._seq
self._proposal_val = None
self._proposal_unanimous = False
async def propose(self, val):
print_debug('Propose:', val)
if self._proposal_unanimous:
return self._proposal_val
local_seq = self._seq
# When we are `await`ing, the status may be changed
# if so, abort the operation
def is_state_changed_then_handle(futs):
if local_seq < self._seq or self._proposal_unanimous:
for f in futs:
f.cancel()
return True
return False
# prepare
cnt = 0
max_proposal_seq = (0, '', '')
fs = [asyncio.ensure_future(s.prepare(self._seq))
for s in self._servers]
for f in asyncio.as_completed(fs):
try:
r = await f
except ConnectionError:
continue
if is_state_changed_then_handle(fs):
return None
seq, proposal_seq, proposal_val = r
if seq != self._seq:
assert seq > self._seq
self._seq = (seq[0] + 1, self._host, self._port)
self._store()
return None
if proposal_val is not None and proposal_seq > max_proposal_seq:
max_proposal_seq = proposal_seq
val = proposal_val
cnt += 1
if cnt >= self._over_half_num:
break
else:
return None
# accept
cnt = 0
fs = [asyncio.ensure_future(s.accept(self._seq, val))
for s in self._servers]
for f in asyncio.as_completed(fs):
try:
seq = await f
except ConnectionError:
continue
if is_state_changed_then_handle(fs):
return None
if seq != self._seq:
assert seq > self._seq
self._seq = (seq[0] + 1, self._host, self._port)
self._store()
return None
cnt += 1
if cnt >= self._over_half_num:
break
else:
return None
# learn
for s in self._servers:
asyncio.ensure_future(async_suppress_exception(s.learn(val)))
return val
def _store(self):
self._storage.dump((self._seq,
self._proposal_seq, self._proposal_val,
self._proposal_unanimous))
def _on_prepare(self, seq):
if seq > self._seq:
self._seq = seq
self._store()
return self._seq, self._proposal_seq, self._proposal_val
def _on_accept(self, seq, val):
if seq >= self._seq:
self._seq = seq
self._proposal_seq = seq
self._proposal_val = val
self._store()
return self._seq
def _on_learn(self, val):
self._proposal_val = val
self._proposal_unanimous = True
self._store()
# --- Paxos ---
def run_cli(host, port, hosts_ports, filename):
rpc = RPC()
servers = [RPCProxy(h, p) for h, p in hosts_ports]
storage = Storage(filename)
paxos = Paxos(host, port, servers, storage, rpc)
def stdin_handle(*_):
asyncio.ensure_future(paxos.propose(stdin.readline().strip()))
rpc.schedule_server(host, port)
loop.add_reader(stdin.fileno(), stdin_handle)
try:
loop.run_forever()
except KeyboardInterrupt:
pass
# --- CLI ---