summaryrefslogtreecommitdiff
path: root/worker.js
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-04-01 12:03:39 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-04-01 12:03:39 +0200
commitc78435c7b68fbaa92a36d366f11fd5db9699ddf4 (patch)
tree7a15aa72da487ae4e3f89b7a3a40adff2eb8b47d /worker.js
parent2f5ce153253164b50aff24a96d8b81abf2b17a90 (diff)
downloadprocessor-c78435c7b68fbaa92a36d366f11fd5db9699ddf4.tar.gz
processor-c78435c7b68fbaa92a36d366f11fd5db9699ddf4.zip
request compile jobs; use transactions
Diffstat (limited to 'worker.js')
-rw-r--r--worker.js34
1 files changed, 28 insertions, 6 deletions
diff --git a/worker.js b/worker.js
index 29e1d03..9aa7cd0 100644
--- a/worker.js
+++ b/worker.js
@@ -7,7 +7,7 @@ var amqp = require("amqplib"),
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
+ BATCHSIZE = process.env.PROCESSOR_BATCH || 50 * (6 + 5), // objects
IDLE_TIMEOUT = process.env.PROCESSOR_IDLETIMEOUT || 500; // ms
(async () => {
@@ -26,6 +26,7 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
await seq.sync();
await ch.assertQueue("process", {durable: true});
+ await ch.assertQueue("compile", {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);
@@ -100,7 +101,8 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
// upsert match
await model.Match.upsert(match, {
- include: [ model.Roster, model.Asset ]
+ include: [ model.Roster, model.Asset ],
+ transaction: transaction
});
// upsert children
@@ -110,7 +112,8 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
roster.shard_id = match.shard_id;
await model.Roster.upsert(roster, {
- include: [ model.Participant/*, model.Team */]
+ include: [ model.Participant/*, model.Team */],
+ transaction: transaction
});
await Promise.all(roster.participants.map(async (participant) => {
@@ -118,7 +121,8 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
participant.roster_api_id = roster.api_id;
participant.player_api_id = participant.player.api_id;
await model.Participant.upsert(participant, {
- include: [ model.Player ]
+ include: [ model.Player ],
+ transaction: transaction
});
}));
@@ -132,7 +136,7 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
await Promise.all(match.assets.map(async (asset) => {
asset.match_api_id = match.api_id;
asset.shard_id = match.shard_id;
- await model.Asset.upsert(asset);
+ await model.Asset.upsert(asset, { transaction: transaction });
}));
}));
@@ -140,7 +144,7 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
// because they are duplicated among a page of matches
// as provided by apigrabber
console.log("processing", players.length, "players", teams.length, "teams");
- await Promise.all(players.map(async (p) => await model.Player.upsert(p) ));
+ await Promise.all(players.map(async (p) => await model.Player.upsert(p, { transaction: transaction }) ));
//await Promise.all(teams.map(async (t) => await model.Team.upsert(flatten(t)) ));
// COMMIT
@@ -150,8 +154,26 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
// notify web
await Promise.all(players.map(async (p) => await ch.publish("amq.topic", p.name, new Buffer("process_commit")) ));
+ // notify compiler
+ 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.sendToQueue("compile", new Buffer(JSON.stringify(p)), {
+ persistent: true,
+ type: "participant"
+ });
+ }));
+ }));
+ }));
+ await Promise.all(players.map(async (p) =>
+ await ch.sendToQueue("compile", new Buffer(JSON.stringify(p)), {
+ persistent: true,
+ type: "player"
+ })
+ ));
} catch (err) { // TODO catch only SQL error, also catch errors in the promises
console.error(err);
+ await transaction.rollback();
await ch.nack(msgs.pop(), true, true); // nack all messages until the last and requeue
// TODO don't requeue broken records
}