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
|
#!/usr/bin/node
/* jshint esnext:true */
'use strict';
var amqp = require("amqplib"),
Seq = require("sequelize"),
Bluebird = require("bluebird"),
jsonapi = Bluebird.promisifyAll(require("superagent-jsonapify/common"));
var MADGLORY_TOKEN = process.env.MADGLORY_TOKEN,
RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
DATABASE_URI = process.env.DATABASE_URI || "sqlite:///db.sqlite";
if (MADGLORY_TOKEN == undefined) throw "Need an API token";
(async () => {
let seq = new Seq(DATABASE_URI),
model = require("./model")(seq, Seq),
rabbit = await amqp.connect(RABBITMQ_URI),
ch = await rabbit.createChannel();
/* 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});
await ch.prefetch(1);
ch.consume("process", async (msg) => {
let data = await jsonapi.parse(msg.content);
// TODO commit less often if possible, avoid deadlocks
let transaction = await seq.transaction({ autocommit: false });
data.data.forEach((match_data) => {
function flatten(obj) {
let attrs = obj.attributes || {},
stats = attrs.stats || {},
o = Object.assign({}, obj, attrs, stats);
delete o.type;
delete o.attributes;
delete o.stats;
delete o.relationships;
return o;
}
let match = JSON.parse(JSON.stringify(match_data)); // deep clone
/* bring jsonapi 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 = flatten(match);
/* upsert everything */
model.Match.upsert(match, {
include: [
{
model: model.Roster,
as: "Rosters",
include: [
{
model: model.Participant,
as: "Participants",
include: [
model.Player
]
},
{
model: model.Team,
as: "Team"
}
]
},
{
model: model.Asset,
as: "Assets"
}
]
});
match.rosters.forEach((roster) => {
model.Roster.upsert(roster, {
include: [
{
model: model.Participant,
as: "Participants",
include: [
model.Player
]
},
{
model: model.Team,
as: "Team"
}
]
});
roster.participants.forEach((participant) => {
model.Participant.upsert(participant, {
include: [
model.Player
]
});
model.Player.upsert(participant.player);
});
if (roster.team != null) model.Team.upsert(roster.team);
});
match.assets.forEach((asset) => {
model.Asset.upsert(asset);
});
});
await transaction.commit(); // TODO rollback on err
ch.ack(msg);
}, { noAck: false });
})();
|