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
|
#!/usr/bin/node
/* jshint esnext:true */
'use strict';
var amqp = require("amqplib"),
Seq = require("sequelize");
var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
DATABASE_URI = process.env.DATABASE_URI || "sqlite:///db.sqlite",
BATCHSIZE = process.env.PROCESSOR_BATCH || 50, // matches
IDLE_TIMEOUT = process.env.PROCESSOR_IDLETIMEOUT || 500; // ms
(async () => {
let seq = new Seq(DATABASE_URI),
model = require("./model")(seq, Seq),
rabbit = await amqp.connect(RABBITMQ_URI),
ch = await rabbit.createChannel();
let queue = [],
timer = undefined;
/* recreate for debugging
await seq.query("SET FOREIGN_KEY_CHECKS=0");
await seq.sync({force: true});
*/
await seq.sync();
await ch.assertQueue("process", {durable: true});
// as long as the queue is filled, msg are not ACKed
// server sends as long as there are less than `prefetch` unACKed
await ch.prefetch(BATCHSIZE);
ch.consume("process", async (msg) => {
queue.push(msg);
// fill queue until batchsize or idle
if (timer === undefined)
timer = setTimeout(process, IDLE_TIMEOUT)
if (queue.length == BATCHSIZE)
await process();
}, { noAck: false });
async function process() {
console.log("processing batch", queue.length);
// clean up to allow processor to accept while we wait for db
let matchmsgs = queue.slice();
queue = [];
clearTimeout(timer);
timer = undefined;
// BEGIN
let transaction = await seq.transaction({ autocommit: false });
// helper to convert API response into flat JSON
// db structure is (almost) 1:1 the API structure
// so we can insert the flat API response as-is
function flatten(obj) {
let attrs = obj.attributes || {},
stats = attrs.stats || {},
o = Object.assign({}, obj, attrs, stats);
o.api_id = o.id; // rename
delete o.id;
delete o.type;
delete o.attributes;
delete o.stats;
delete o.relationships;
return o;
}
// UPSERT
try {
let matches = matchmsgs.map((msg) => JSON.parse(msg.content));
await Promise.all(matches.map(async (match) => {
console.log("processing match", match.id);
// flatten jsonapi nested response into our db structure-like shape
match.rosters = match.rosters.map((roster) => {
roster.participants = roster.participants.map((participant) => {
participant.player = flatten(participant.player);
return flatten(participant);
});
return flatten(roster);
});
match.assets = match.assets.map((asset) => flatten(asset));
match = flatten(match);
// upsert match
await model.Match.upsert(match, {
include: [ model.Roster, model.Asset ]
});
// upsert children
// before, add foreign keys and other missing information (shardId)
await Promise.all(match.rosters.map(async (roster) => {
roster.match_api_id = match.api_id;
roster.shard_id = match.shard_id;
await model.Roster.upsert(roster, {
include: [ model.Participant/*, model.Team */]
});
await Promise.all(roster.participants.map(async (participant) => {
participant.player.shard_id = participant.shard_id;
await model.Player.upsert(participant.player);
participant.shard_id = roster.shard_id;
participant.roster_api_id = roster.api_id;
participant.player_api_id = participant.player.api_id;
await model.Participant.upsert(participant, {
include: [ model.Player ]
});
}));
if (roster.team != null) {
roster.team.shard_id = roster.shard_id;
roster.team.roster_api_id = roster.api_id;
/*await model.Team.upsert(roster.team);*/
}
}));
await Promise.all(match.assets.map(async (asset) => {
asset.match_api_id = match.api_id;
asset.shard_id = match.shard_id;
asset.url = asset.uRL; // camelcasization failed ;)
await model.Asset.upsert(asset);
}));
}));
// COMMIT
await transaction.commit();
console.log("acking batch");
await Promise.all(matchmsgs.map(async (msg) => {
await ch.ack(msg);
}));
// request child jobs, notify player
await Promise.all(matches.map(async (m) => {
await Promise.all(m.rosters.map(async (r) => {
await Promise.all(r.participants.map(async (p) => {
await ch.publish("amq.topic", p.player.name, new Buffer("process_commit"));
}));
}));
}));
} catch (err) { // TODO catch only SQL error, also catch errors in the promises
console.error(err);
await ch.nack(matchmsgs.pop(), true, true); // nack all messages until the last and requeue
// TODO don't requeue broken records
}
}
})();
|