-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathloopbroker.coffee
More file actions
109 lines (91 loc) · 2.78 KB
/
Copy pathloopbroker.coffee
File metadata and controls
109 lines (91 loc) · 2.78 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
#!/usr/bin/env coffee
#
ctx = require 'zmq'
EventEmitter = require('events').EventEmitter
class Node extends EventEmitter
constructor: (@client) ->
class Broker extends Node
constructor: (@client) -> # inst var inside constructor need this prefix
@cloudfe = ctx.createSocket('router')
@cloudfe.identity = 'CFE'
@cloudbe = ctx.createSocket('router')
@cloudbe.identity = 'CBE'
@localbe = ctx.createSocket('router')
@localbe.identity = 'LBE'
@bindEvent()
ready: (me, peer) ->
@_ID = me
@_PEER = peer
@cloudfe.identity = me+":CFE"
@cloudfe.bind("ipc://router-" + me + ".ipc", (err) -> console.log err if err)
@cloudbe.identity = me+":CBE"
@cloudbe.connect("ipc://router-"+peer+".ipc")
@localbe.identity = me+":LBE"
@localbe.bind("ipc://router-wk-"+me+".ipc", (err) -> console.log err if err)
console.log @_ID, ' ready...'
self = this
# do invoke the func with passed in arg
tmot = do (self) =>
() => # return a func with self closed over
console.log 'send ready to peer broker'
msg = []
msg.push @_ID+' Ready'
self.sendMsg @cloudbe, @_PEER+":CFE", msg
#setTimeout tmot, 5000
# prototype bind event
bindEvent: ->
@cloudfe.on 'message', (msg) =>
# func created by => can access this property where they are defined.
msg = @parseMsg arguments
from = msg[0]
@processReq(from, msg)
@cloudbe.on 'message', (msg) =>
# func created by => can access this property where they are defined.
msg = @parseMsg arguments
from = msg[0]
@processRep(from, msg)
@localbe.on 'message', (msg) =>
# func created by => can access this property where they are defined.
msg = @parseMsg arguments
from = msg[0]
data = msg[2]
if data is 'ready'
console.log msg.toString()
else
@processReq(from, msg)
process.on 'SIGINT', () =>
@cloudbe.close()
@cloudfe.close()
parseMsg: (args) ->
msg = []
for arg in args
msg.push arg.toString()
return msg
# process msg sent to worker
processReq: (from, msg) ->
if @_ID is 'b1'
# forward thru be, prepend be header
#msg.unshift @_ID+':CBE'
console.log @_ID, ':CFE <<< ', msg
@sendMsg @cloudbe, @_PEER+":CFE", msg
else
console.log @_ID, ':CFE <<< ', msg
from = msg.shift() # get rid of env header if not forward anymore
msg.push 'b2 done'
@sendMsg @cloudfe, from, msg # reply back
processRep: (from, msg) ->
if @_ID is 'b1'
console.log @_ID, ':CBE <<< ', msg
from = msg.shift()
#@sendMsg @cloudfe, @_PEER+":CBE", msg
client = msg.shift()
@sendMsg @cloudfe, client, msg
else
console.log @_ID, ':CBE done done <<< ', msg
sendMsg: (sock, peer, msg) ->
msg.unshift peer
console.log @_ID + ' >>> ', peer, msg
sock.send.apply sock, msg
exports.Broker = Broker
broker = new Broker()
broker.ready(process.argv[2], process.argv[3])