-
Notifications
You must be signed in to change notification settings - Fork 52
Expand file tree
/
Copy path1-broker.js
More file actions
109 lines (98 loc) · 2.5 KB
/
1-broker.js
File metadata and controls
109 lines (98 loc) · 2.5 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
'use strict';
global.api = {};
api.net = require('net');
api.os = require('os');
api.metasync = require('metasync');
api.queue = require('./queue.js');
const workersQueue = new api.queue.Queue();
let workerId = 0;
let taskId = 0;
const newWorker = (socket) => {
socket.on('data', (data) => {
const connParams = JSON.parse(data);
//add new worker
workersQueue.put({
id: workerId++,
port: connParams.port,
host: socket.remoteAddress
});
console.dir(workersQueue.queue);
socket.end();
});
};
const createConnection = (task) => (data, datacb) => {
workersQueue.use((workerParams, qcb) => {
const workerSocket = new api.net.Socket();
workerSocket.on('data', (readData) => {
console.log('Data received (by broker): ' + readData);
const res = JSON.parse(readData);
data[parseInt(res.index)] = res.answer;
qcb();
datacb(data);
});
workerSocket.connect({
port: workerParams.port,
host: workerParams.host
}, () => {
console.log('Data send (by broker): ' + JSON.stringify(task));
workerSocket.write(JSON.stringify(task));
});
});
};
const mergeResult = (data) => {
console.log('Merging - ');
console.dir(data);
//merge array of results in one
const res = data.reduce((a, b) => a.concat(b));
// console.log(task);
console.log(res);
return res;
};
const createTasks = (arr, workers) => {
const tasks = [];
const elemsByTask = Math.ceil(arr.length / workers);
let i = 0;
while (arr.length > 0) {
tasks.push({
index: i++,
funcId: 1,
task: arr.splice(0, elemsByTask),
taskId: taskId++
});
}
return tasks;
};
const sendTasks = (tasks, clientSocket) => {
api.metasync.map(
tasks,
(curr, cb) => {
cb(null, createConnection(curr));
},
(err, res) => {
if (err) console.error(err);
api.metasync.parallel(
res,
(data) => {
const res = mergeResult(data);
const response = { answer: res };
clientSocket.end(JSON.stringify(response));
},
[]
);
});
};
const newTasks = (data, clientSocket) => {
const obj = JSON.parse(data);
sendTasks(createTasks(obj.task, workersQueue.size()), clientSocket);
};
//for clients
api.net.createServer((clientSocket) => {
clientSocket.on('error', (err) => {
console.error(err);
});
clientSocket.on('data', (data) => {
newTasks(data, clientSocket);
});
}).listen(50000);
//for workers
api.net.createServer(newWorker).listen(51000);