-
Notifications
You must be signed in to change notification settings - Fork 3
/
Copy pathworker.py
executable file
·130 lines (117 loc) · 4.4 KB
/
worker.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
#!/usr/bin/env python
"""MapReduce Worker (Container).
Usage:
worker.py [-br] [(-d <dr> -t <ttl>)] <IP> <PORT> <MASTER_IP> <MASTER_PORT>
Options:
-h --help Show this screen.
-b --background Run service in background.
-r --random Randomize completion time.
-d --die=<dr> Probability that the worker will die [default: 0.0].
-t --timeToLive=<ttl> Average time worker lives before death in tasks.
"""
from docopt import docopt
from rpc import RPCManager, RPC
from session import WorkerSessionManager
import time
import daemon
from multiprocessing import Queue
from random import gauss, random
import sys
RAND = 0
DIE = False
TTL = 0
def run(IP, PORT, mIP, mPORT):
processQ = Queue()
sessionManager = WorkerSessionManager(IP, PORT, mIP, mPORT, processQ)
rpcManager = RPCManager(sessionManager, processQ)
working = None
state = "IDLE";
doneTime = 0
while True:
sessionManager.poll()
rpcManager.poll()
# Convert incoming RPCs (events) into state changes
for rpc in rpcManager.inRPC.values():
if rpc.status == "pending":
rpcType, payload = rpc.msg
if rpcType == "LAUNCH":
if state == "IDLE":
state = "RUNNING"
working = rpc
rpc.status = "working"
continue
elif rpcType == "COMMIT":
if state == "COMPLETE":
state = "COMMITTING"
working = rpc
rpc.status = "working"
continue
elif rpcType == "CONTAINER_REMOTE_CLEANUP":
state = "CLEANUP"
if working != None:
working.reply = "failed"
working.status = "send"
working = rpc
rpc.status = "working"
continue
elif rpcType == "DIE":
sys.exit(0)
rpc.reply = "failed"
rpc.status = "send"
# Rules engine
if state == "IDLE":
pass
elif state == "RUNNING":
# simulate random completion time
if doneTime == 0:
doneTime = time.time() + 5.0 + RAND * gauss(0.0, 1.0)
if time.time() > doneTime:
working.reply = working.msg
working.status = "send"
print "Work Finished: ", working.msg, working.locator
working = None
doneTime = 0
state = "COMPLETE"
elif state == "COMPLETE":
pass
elif state == "COMMITTING":
# simulate random commit time
if doneTime == 0:
doneTime = time.time() + 2.0 + RAND * gauss(0.0, 1.0)
if time.time() > doneTime:
print "Work Committed: ", working.msg, working.locator
doneTime = 0
state = "CLEANUP"
elif state == "CLEANUP":
# simulate random cleaup time
if doneTime == 0:
doneTime = time.time() + 1.0 + RAND * gauss(0.0, 1.0)
if time.time() > doneTime:
working.reply = working.msg
working.status = "send"
print "Container Clean: ", working.msg, working.locator
working = None
doneTime = 0
state = "IDLE"
# Kill failed RPCs
for key in rpcManager.outRPC.keys():
if rpcManager.outRPC[key].status == "failed":
del rpcManager.outRPC[key]
if DIE and time.time() > TTL:
sys.exit(0)
if __name__ == '__main__':
args = docopt(__doc__)
print(args)
RAND = int(args['--random'])
if random() < float(args['--die']):
DIE = True
TTL = time.time() + 5.0 * gauss(float(args['--timeToLive']),
float(args['--timeToLive']) / 2)
if (args['--background']):
with daemon.DaemonContext():
run(args['<IP>'], int(args['<PORT>']),
args['<MASTER_IP>'], int(args['<MASTER_PORT>']))
print "Done"
else:
run(args['<IP>'], int(args['<PORT>']),
args['<MASTER_IP>'], int(args['<MASTER_PORT>']))