-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathconsensus.js
More file actions
112 lines (104 loc) · 3.56 KB
/
Copy pathconsensus.js
File metadata and controls
112 lines (104 loc) · 3.56 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
var fs = require('fs');
var querystring = require('querystring');
var httpjson = require('./httpjson');
var args = process.argv.slice(2);
var id = args[0];
var config = JSON.parse(fs.readFileSync('config', 'utf8'));
var replicas = config['replicas'];
var numNodes = replicas.length;
var readQuorum = config['read_quorum'];
var writeQuorum = config['write_quorum'];
if (readQuorum + writeQuorum <= numNodes) {
console.log('Error: in config file, read and write quorums added together must be greater than number of replicas');
process.exit(1);
}
if (writeQuorum <= numNodes / 2) {
console.log('Error: write quorum must be greater than numNodes / 2');
process.exit(1);
}
// stagger replicas so not every consensus node is writing to same set of storage nodes
var myIndex = null;
for (var i = 0; i < replicas.length; i++) {
if (replicas[i]['id'] == id) {
myIndex = i;
break;
}
}
if (myIndex != null) {
replicas = replicas.concat(replicas.slice(0, myIndex));
replicas = replicas.slice(myIndex);
}
exports.read = function read(sreq, res, next) {
var responses = [];
var numReadSucceed = 0;
var numReadFail = 0;
var readData = JSON.stringify({
'key' : sreq.query.key
});
function readFromNode(domain, port, readData) {
httpjson.get(domain, port, '/read_vote', readData, function(response) {
responses.push(JSON.parse(response)['value']);
numReadSucceed++;
if (responses.length == readQuorum) {
var mostRecent = null;
responses.forEach(function(r) {
if (r != null && (mostRecent == null || r['timestamp'] < mostRecent['timestamp'])) {
mostRecent = r;
}
});
res.status(200).send({'result' : mostRecent});
}
}, function() {
numReadFail++;
if (numReadFail > numNodes - readQuorum) {
res.status(200).send({'result' : 'error'});
} else {
var replica = replicas[readQuorum - 1 + numReadFail];
readFromNode(replica['domain'], replica['port'], readData);
}
});
}
for (var i = 0; i < readQuorum; i++) {
var replica = replicas[i];
readFromNode(replica['domain'], replica['port'], readData);
}
}
exports.write = function write(sreq, res, next) {
responses = [];
var numWriteSucceed = 0;
var numWriteFail = 0;
var body = sreq.body;
var writeData = JSON.stringify({
'key' : body.key,
'value' : body.value,
'timestamp' : new Date().getTime()
});
function onWrite(response) {
body = JSON.parse(response);
if (body['status'] != 'success') {
onWriteFail();
return;
}
numWriteSucceed++;
responses.push(response);
if (numWriteSucceed == writeQuorum) {
res.send({'status' : 'success'});
}
}
function onWriteFail() {
numWriteFail++;
if (numWriteFail > numNodes - writeQuorum) {
res.send({'status' : 'fail'});
} else {
var replica = replicas[writeQuorum - 1 + numWriteFail];
writeToNode(replica['domain'], replica['port'], writeData);
}
}
function writeToNode(domain, port, writeData) {
httpjson.post(domain, port, '/write_vote', writeData, onWrite, onWriteFail);;
}
for (var i = 0; i < writeQuorum; i++) {
var replica = replicas[i];
writeToNode(replica['domain'], replica['port'], writeData);
}
};