summaryrefslogtreecommitdiff
path: root/worker.js
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-03-29 20:31:05 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-03-29 20:31:05 +0200
commit5c869aa26666215ec4ea5315b40fea3025a1ae83 (patch)
tree7d1468587aa899fd364a5e96d193665e0bba61bb /worker.js
parent202d7641749b817c8893cc07918e23b211c92ba7 (diff)
downloadprocessor-5c869aa26666215ec4ea5315b40fea3025a1ae83.tar.gz
processor-5c869aa26666215ec4ea5315b40fea3025a1ae83.zip
import db schema from vainsocial
Diffstat (limited to 'worker.js')
-rw-r--r--worker.js123
1 files changed, 44 insertions, 79 deletions
diff --git a/worker.js b/worker.js
index aec51b9..7cefbdb 100644
--- a/worker.js
+++ b/worker.js
@@ -3,9 +3,7 @@
'use strict';
var amqp = require("amqplib"),
- Seq = require("sequelize"),
- Bluebird = require("bluebird"),
- jsonapi = Bluebird.promisifyAll(require("superagent-jsonapify/common"));
+ Seq = require("sequelize");
var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
DATABASE_URI = process.env.DATABASE_URI || "sqlite:///db.sqlite";
@@ -20,100 +18,67 @@ var RABBITMQ_URI = process.env.RABBITMQ_URI || "amqp://localhost",
await seq.query("SET FOREIGN_KEY_CHECKS=0");
await seq.sync({force: true});
*/
- await seq.sync();
+ //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);
+ let match = JSON.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
+ 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;
+ }
- /* 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);
+ /* bring 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);
});
- match = flatten(match);
+ return flatten(roster);
+ });
+ match = flatten(match);
+ console.log(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"
- }
- ]
- });
+ /* upsert everything */
+ await match.rosters.forEach(async (roster) => {
+ await roster.participants.forEach(async (participant) => {
+ await model.Player.upsert(participant.player);
- match.rosters.forEach((roster) => {
- model.Roster.upsert(roster, {
- include: [
- {
- model: model.Participant,
- as: "Participants",
- include: [
- model.Player
- ]
- },
- {
- model: model.Team,
- as: "Team"
- }
- ]
+ participant.roster_api_id = roster.api_id;
+ participant.player_api_id = participant.player.api_id;
+ await model.Participant.upsert(participant, {
+ include: [ model.Player ]
});
+ });
- roster.participants.forEach((participant) => {
- model.Participant.upsert(participant, {
- include: [
- model.Player
- ]
- });
+ //if (roster.team != null) model.Team.upsert(roster.team);
- model.Player.upsert(participant.player);
- });
-
- if (roster.team != null) model.Team.upsert(roster.team);
+ roster.match_api_id = match.api_id;
+ await model.Roster.upsert(roster, {
+ include: [ model.Participant/*, model.Team*/ ]
});
+ });
- match.assets.forEach((asset) => {
- model.Asset.upsert(asset);
- });
+ /*match.assets.forEach((asset) => {
+ model.Asset.upsert(asset);
+ });*/
+
+ await model.Match.upsert(match, {
+ include: [ model.Roster/*, model.Asset*/ ]
});
await transaction.commit(); // TODO rollback on err