summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorschneefux <schneefux+commit@schneefux.xyz>2017-07-14 18:58:52 +0200
committerschneefux <schneefux+commit@schneefux.xyz>2017-07-14 18:58:52 +0200
commit590389198bc0a7d9ce00395287c715a4adbc54b0 (patch)
treeaee0ff1dbf286dee0425d34ca08b047c3bcc1f44
parenta5b43e106ae8a3441690ed7d32c711f0200d4be7 (diff)
downloadprocessor-590389198bc0a7d9ce00395287c715a4adbc54b0.tar.gz
processor-590389198bc0a7d9ce00395287c715a4adbc54b0.zip
2.14.0 rewrite
-rw-r--r--package-lock.json895
-rw-r--r--package.json7
-rw-r--r--sequelize-bulkinsert.patch19
-rw-r--r--worker.js867
4 files changed, 1199 insertions, 589 deletions
diff --git a/package-lock.json b/package-lock.json
new file mode 100644
index 0000000..0840512
--- /dev/null
+++ b/package-lock.json
@@ -0,0 +1,895 @@
+{
+ "name": "processor",
+ "version": "2.0.0",
+ "lockfileVersion": 1,
+ "requires": true,
+ "dependencies": {
+ "@types/geojson": {
+ "version": "1.0.2",
+ "resolved": "https://registry.npmjs.org/@types/geojson/-/geojson-1.0.2.tgz",
+ "integrity": "sha1-sC0QqwKOKSisWSoFGqpJgaGUHQM="
+ },
+ "@types/node": {
+ "version": "6.0.82",
+ "resolved": "https://registry.npmjs.org/@types/node/-/node-6.0.82.tgz",
+ "integrity": "sha512-uFjwSPo9dyoK0bhxMIuCsqIyQW9EHN2RcadPlA2s67TxzGDL9MJKjNAsi+TCn8ProGHcntLhgvIpbZmbtk2uCQ=="
+ },
+ "ajv": {
+ "version": "4.11.8",
+ "resolved": "https://registry.npmjs.org/ajv/-/ajv-4.11.8.tgz",
+ "integrity": "sha1-gv+wKynmYq5TvcIK8VlHcGc5xTY=",
+ "requires": {
+ "co": "4.6.0",
+ "json-stable-stringify": "1.0.1"
+ }
+ },
+ "amqplib": {
+ "version": "0.5.1",
+ "resolved": "https://registry.npmjs.org/amqplib/-/amqplib-0.5.1.tgz",
+ "integrity": "sha1-fMz+ur5WwumE6noiQ/fO/m+/xs8=",
+ "requires": {
+ "bitsyntax": "0.0.4",
+ "bluebird": "3.5.0",
+ "buffer-more-ints": "0.0.2",
+ "readable-stream": "1.1.14"
+ },
+ "dependencies": {
+ "bitsyntax": {
+ "version": "0.0.4",
+ "resolved": "https://registry.npmjs.org/bitsyntax/-/bitsyntax-0.0.4.tgz",
+ "integrity": "sha1-6xDMb4K4xJDj6FaY8H6D1G4MuoI=",
+ "requires": {
+ "buffer-more-ints": "0.0.2"
+ }
+ },
+ "bluebird": {
+ "version": "3.5.0",
+ "resolved": "https://registry.npmjs.org/bluebird/-/bluebird-3.5.0.tgz",
+ "integrity": "sha1-eRQg1/VR7qKJdFOop3ZT+WYG1nw="
+ },
+ "buffer-more-ints": {
+ "version": "0.0.2",
+ "resolved": "https://registry.npmjs.org/buffer-more-ints/-/buffer-more-ints-0.0.2.tgz",
+ "integrity": "sha1-JrOIXRD6E9t/wBquOquHAZngEkw="
+ },
+ "core-util-is": {
+ "version": "1.0.2",
+ "resolved": "https://registry.npmjs.org/core-util-is/-/core-util-is-1.0.2.tgz",
+ "integrity": "sha1-tf1UIgqivFq1eqtxQMlAdUUDwac="
+ },
+ "inherits": {
+ "version": "2.0.3",
+ "resolved": "https://registry.npmjs.org/inherits/-/inherits-2.0.3.tgz",
+ "integrity": "sha1-Yzwsg+PaQqUC9SRmAiSA9CCCYd4="
+ },
+ "isarray": {
+ "version": "0.0.1",
+ "resolved": "https://registry.npmjs.org/isarray/-/isarray-0.0.1.tgz",
+ "integrity": "sha1-ihis/Kmo9Bd+Cav8YDiTmwXR7t8="
+ },
+ "readable-stream": {
+ "version": "1.1.14",
+ "resolved": "https://registry.npmjs.org/readable-stream/-/readable-stream-1.1.14.tgz",
+ "integrity": "sha1-fPTFTvZI44EwhMY23SB54WbAgdk=",
+ "requires": {
+ "core-util-is": "1.0.2",
+ "inherits": "2.0.3",
+ "isarray": "0.0.1",
+ "string_decoder": "0.10.31"
+ }
+ },
+ "string_decoder": {
+ "version": "0.10.31",
+ "resolved": "https://registry.npmjs.org/string_decoder/-/string_decoder-0.10.31.tgz",
+ "integrity": "sha1-YuIDvEF2bGwoyfyEMB2rHFMQ+pQ="
+ }
+ }
+ },
+ "ansicolors": {
+ "version": "0.2.1",
+ "resolved": "https://registry.npmjs.org/ansicolors/-/ansicolors-0.2.1.tgz",
+ "integrity": "sha1-vgiVmQl7dKXJxKhKDNvNtivYeu8="
+ },
+ "asn1": {
+ "version": "0.2.3",
+ "resolved": "https://registry.npmjs.org/asn1/-/asn1-0.2.3.tgz",
+ "integrity": "sha1-2sh4dxPJlmhJ/IGAd36+nB3fO4Y="
+ },
+ "assert-plus": {
+ "version": "0.2.0",
+ "resolved": "https://registry.npmjs.org/assert-plus/-/assert-plus-0.2.0.tgz",
+ "integrity": "sha1-104bh+ev/A24qttwIfP+SBAasjQ="
+ },
+ "asynckit": {
+ "version": "0.4.0",
+ "resolved": "https://registry.npmjs.org/asynckit/-/asynckit-0.4.0.tgz",
+ "integrity": "sha1-x57Zf380y48robyXkLzDZkdLS3k="
+ },
+ "aws-sign2": {
+ "version": "0.6.0",
+ "resolved": "https://registry.npmjs.org/aws-sign2/-/aws-sign2-0.6.0.tgz",
+ "integrity": "sha1-FDQt0428yU0OW4fXY81jYSwOeU8="
+ },
+ "aws4": {
+ "version": "1.6.0",
+ "resolved": "https://registry.npmjs.org/aws4/-/aws4-1.6.0.tgz",
+ "integrity": "sha1-g+9cqGCysy5KDe7e6MdxudtXRx4="
+ },
+ "bcrypt-pbkdf": {
+ "version": "1.0.1",
+ "resolved": "https://registry.npmjs.org/bcrypt-pbkdf/-/bcrypt-pbkdf-1.0.1.tgz",
+ "integrity": "sha1-Y7xdy2EzG5K8Bf1SiVPDNGKgb40=",
+ "optional": true,
+ "requires": {
+ "tweetnacl": "0.14.5"
+ }
+ },
+ "bluebird": {
+ "version": "https://registry.npmjs.org/bluebird/-/bluebird-3.5.0.tgz",
+ "integrity": "sha1-eRQg1/VR7qKJdFOop3ZT+WYG1nw="
+ },
+ "boom": {
+ "version": "2.10.1",
+ "resolved": "https://registry.npmjs.org/boom/-/boom-2.10.1.tgz",
+ "integrity": "sha1-OciRjO/1eZ+D+UkqhI9iWt0Mdm8=",
+ "requires": {
+ "hoek": "2.16.3"
+ }
+ },
+ "cardinal": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/cardinal/-/cardinal-1.0.0.tgz",
+ "integrity": "sha1-UOIcGwqjdyn5N33vGWtanOyTLuk=",
+ "requires": {
+ "ansicolors": "0.2.1",
+ "redeyed": "1.0.1"
+ }
+ },
+ "caseless": {
+ "version": "0.12.0",
+ "resolved": "https://registry.npmjs.org/caseless/-/caseless-0.12.0.tgz",
+ "integrity": "sha1-G2gcIf+EAzyCZUMJBolCDRhxUdw="
+ },
+ "cls-bluebird": {
+ "version": "2.0.1",
+ "resolved": "https://registry.npmjs.org/cls-bluebird/-/cls-bluebird-2.0.1.tgz",
+ "integrity": "sha1-wlmkgK4CwOUGE0MHuxPbMERu4uc=",
+ "requires": {
+ "is-bluebird": "1.0.2",
+ "shimmer": "1.1.0"
+ }
+ },
+ "co": {
+ "version": "4.6.0",
+ "resolved": "https://registry.npmjs.org/co/-/co-4.6.0.tgz",
+ "integrity": "sha1-bqa989hTrlTMuOR7+gvz+QMfsYQ="
+ },
+ "combined-stream": {
+ "version": "1.0.5",
+ "resolved": "https://registry.npmjs.org/combined-stream/-/combined-stream-1.0.5.tgz",
+ "integrity": "sha1-k4NwpXtKUd6ix3wV1cX9+JUWQAk=",
+ "requires": {
+ "delayed-stream": "1.0.0"
+ }
+ },
+ "core-util-is": {
+ "version": "https://registry.npmjs.org/core-util-is/-/core-util-is-1.0.2.tgz",
+ "integrity": "sha1-tf1UIgqivFq1eqtxQMlAdUUDwac="
+ },
+ "cross-env": {
+ "version": "3.2.4",
+ "resolved": "https://registry.npmjs.org/cross-env/-/cross-env-3.2.4.tgz",
+ "integrity": "sha1-ngWF8neGTtQhznVvgamA/w1piro=",
+ "requires": {
+ "cross-spawn": "5.1.0",
+ "is-windows": "1.0.1"
+ }
+ },
+ "cross-spawn": {
+ "version": "5.1.0",
+ "resolved": "https://registry.npmjs.org/cross-spawn/-/cross-spawn-5.1.0.tgz",
+ "integrity": "sha1-6L0O/uWPz/b4+UUQoKVUu/ojVEk=",
+ "requires": {
+ "lru-cache": "4.1.1",
+ "shebang-command": "1.2.0",
+ "which": "1.2.14"
+ }
+ },
+ "cryptiles": {
+ "version": "2.0.5",
+ "resolved": "https://registry.npmjs.org/cryptiles/-/cryptiles-2.0.5.tgz",
+ "integrity": "sha1-O9/s3GCBR8HGcgL6KR59ylnqo7g=",
+ "requires": {
+ "boom": "2.10.1"
+ }
+ },
+ "dashdash": {
+ "version": "1.14.1",
+ "resolved": "https://registry.npmjs.org/dashdash/-/dashdash-1.14.1.tgz",
+ "integrity": "sha1-hTz6D3y+L+1d4gMmuN1YEDX24vA=",
+ "requires": {
+ "assert-plus": "1.0.0"
+ },
+ "dependencies": {
+ "assert-plus": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/assert-plus/-/assert-plus-1.0.0.tgz",
+ "integrity": "sha1-8S4PPF13sLHN2RRpQuTpbB5N1SU="
+ }
+ }
+ },
+ "debug": {
+ "version": "2.6.8",
+ "resolved": "https://registry.npmjs.org/debug/-/debug-2.6.8.tgz",
+ "integrity": "sha1-5zFTHKLt4n0YgiJCfaF4IdaP9Pw=",
+ "requires": {
+ "ms": "2.0.0"
+ }
+ },
+ "delayed-stream": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/delayed-stream/-/delayed-stream-1.0.0.tgz",
+ "integrity": "sha1-3zrhmayt+31ECqrgsp4icrJOxhk="
+ },
+ "denque": {
+ "version": "1.1.1",
+ "resolved": "https://registry.npmjs.org/denque/-/denque-1.1.1.tgz",
+ "integrity": "sha1-ECKcK4juwb0V/4LF/eNW5762254="
+ },
+ "depd": {
+ "version": "1.1.0",
+ "resolved": "https://registry.npmjs.org/depd/-/depd-1.1.0.tgz",
+ "integrity": "sha1-4b2Cxqq2ztlluXuIsX7T5SjKGMM="
+ },
+ "dottie": {
+ "version": "2.0.0",
+ "resolved": "https://registry.npmjs.org/dottie/-/dottie-2.0.0.tgz",
+ "integrity": "sha1-2hkZgci41xPKARXViYzzl8Lw3dA="
+ },
+ "ecc-jsbn": {
+ "version": "0.1.1",
+ "resolved": "https://registry.npmjs.org/ecc-jsbn/-/ecc-jsbn-0.1.1.tgz",
+ "integrity": "sha1-D8c6ntXw1Tw4GTOYUj735UN3dQU=",
+ "optional": true,
+ "requires": {
+ "jsbn": "0.1.1"
+ }
+ },
+ "env-cmd": {
+ "version": "5.1.0",
+ "resolved": "https://registry.npmjs.org/env-cmd/-/env-cmd-5.1.0.tgz",
+ "integrity": "sha1-AjbbOTw/AzAFIE/NCpLuQHI6nJ4=",
+ "requires": {
+ "cross-spawn": "5.1.0"
+ }
+ },
+ "esprima": {
+ "version": "3.0.0",
+ "resolved": "https://registry.npmjs.org/esprima/-/esprima-3.0.0.tgz",
+ "integrity": "sha1-U88kes2ncxPlUcOqLnM0LT+099k="
+ },
+ "extend": {
+ "version": "3.0.1",
+ "resolved": "https://registry.npmjs.org/extend/-/extend-3.0.1.tgz",
+ "integrity": "sha1-p1Xqe8Gt/MWjHOfnYtuq3F5jZEQ="
+ },
+ "extsprintf": {
+ "version": "1.0.2",
+ "resolved": "https://registry.npmjs.org/extsprintf/-/extsprintf-1.0.2.tgz",
+ "integrity": "sha1-4QgOBljjALBilJkMxw4VAiNf1VA="
+ },
+ "forever-agent": {
+ "version": "0.6.1",
+ "resolved": "https://registry.npmjs.org/forever-agent/-/forever-agent-0.6.1.tgz",
+ "integrity": "sha1-+8cfDEGt6zf5bFd60e1C2P2sypE="
+ },
+ "form-data": {
+ "version": "2.1.4",
+ "resolved": "https://registry.npmjs.org/form-data/-/form-data-2.1.4.tgz",
+ "integrity": "sha1-M8GDrPGTJ27KqYFDpp6Uv+4XUNE=",
+ "requires": {
+ "asynckit": "0.4.0",
+ "combined-stream": "1.0.5",
+ "mime-types": "2.1.15"
+ }
+ },
+ "generate-function": {
+ "version": "2.0.0",
+ "resolved": "https://registry.npmjs.org/generate-function/-/generate-function-2.0.0.tgz",
+ "integrity": "sha1-aFj+fAlpt9TpCTM3ZHrHn2DfvnQ="
+ },
+ "generic-pool": {
+ "version": "3.1.7",
+ "resolved": "https://registry.npmjs.org/generic-pool/-/generic-pool-3.1.7.tgz",
+ "integrity": "sha1-2sIrLHp6BOQXMvfY0tJaMDyI9mI="
+ },
+ "getpass": {
+ "version": "0.1.7",
+ "resolved": "https://registry.npmjs.org/getpass/-/getpass-0.1.7.tgz",
+ "integrity": "sha1-Xv+OPmhNVprkyysSgmBOi6YhSfo=",
+ "requires": {
+ "assert-plus": "1.0.0"
+ },
+ "dependencies": {
+ "assert-plus": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/assert-plus/-/assert-plus-1.0.0.tgz",
+ "integrity": "sha1-8S4PPF13sLHN2RRpQuTpbB5N1SU="
+ }
+ }
+ },
+ "har-schema": {
+ "version": "1.0.5",
+ "resolved": "https://registry.npmjs.org/har-schema/-/har-schema-1.0.5.tgz",
+ "integrity": "sha1-0mMTX0MwfALGAq/I/pWXDAFRNp4="
+ },
+ "har-validator": {
+ "version": "4.2.1",
+ "resolved": "https://registry.npmjs.org/har-validator/-/har-validator-4.2.1.tgz",
+ "integrity": "sha1-M0gdDxu/9gDdID11gSpqX7oALio=",
+ "requires": {
+ "ajv": "4.11.8",
+ "har-schema": "1.0.5"
+ }
+ },
+ "hawk": {
+ "version": "3.1.3",
+ "resolved": "https://registry.npmjs.org/hawk/-/hawk-3.1.3.tgz",
+ "integrity": "sha1-B4REvXwWQLD+VA0sm3PVlnjo4cQ=",
+ "requires": {
+ "boom": "2.10.1",
+ "cryptiles": "2.0.5",
+ "hoek": "2.16.3",
+ "sntp": "1.0.9"
+ }
+ },
+ "hoek": {
+ "version": "2.16.3",
+ "resolved": "https://registry.npmjs.org/hoek/-/hoek-2.16.3.tgz",
+ "integrity": "sha1-ILt0A9POo5jpHcRxCo/xuCdKJe0="
+ },
+ "http-signature": {
+ "version": "1.1.1",
+ "resolved": "https://registry.npmjs.org/http-signature/-/http-signature-1.1.1.tgz",
+ "integrity": "sha1-33LiZwZs0Kxn+3at+OE0qPvPkb8=",
+ "requires": {
+ "assert-plus": "0.2.0",
+ "jsprim": "1.4.0",
+ "sshpk": "1.13.1"
+ }
+ },
+ "iconv-lite": {
+ "version": "0.4.18",
+ "resolved": "https://registry.npmjs.org/iconv-lite/-/iconv-lite-0.4.18.tgz",
+ "integrity": "sha512-sr1ZQph3UwHTR0XftSbK85OvBbxe/abLGzEnPENCQwmHf7sck8Oyu4ob3LgBxWWxRoM+QszeUyl7jbqapu2TqA=="
+ },
+ "inflection": {
+ "version": "1.10.0",
+ "resolved": "https://registry.npmjs.org/inflection/-/inflection-1.10.0.tgz",
+ "integrity": "sha1-W//LEZetPoEFD44X4hZoCH7p6y8="
+ },
+ "inherits": {
+ "version": "https://registry.npmjs.org/inherits/-/inherits-2.0.3.tgz",
+ "integrity": "sha1-Yzwsg+PaQqUC9SRmAiSA9CCCYd4="
+ },
+ "is-bluebird": {
+ "version": "1.0.2",
+ "resolved": "https://registry.npmjs.org/is-bluebird/-/is-bluebird-1.0.2.tgz",
+ "integrity": "sha1-CWQ5Bg9KpBGr7hkUOoTWpVNG1uI="
+ },
+ "is-typedarray": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/is-typedarray/-/is-typedarray-1.0.0.tgz",
+ "integrity": "sha1-5HnICFjfDBsR3dppQPlgEfzaSpo="
+ },
+ "is-windows": {
+ "version": "1.0.1",
+ "resolved": "https://registry.npmjs.org/is-windows/-/is-windows-1.0.1.tgz",
+ "integrity": "sha1-MQ23D3QtJZoWo2kgK1GvhCMzENk="
+ },
+ "isexe": {
+ "version": "2.0.0",
+ "resolved": "https://registry.npmjs.org/isexe/-/isexe-2.0.0.tgz",
+ "integrity": "sha1-6PvzdNxVb/iUehDcsFctYz8s+hA="
+ },
+ "isstream": {
+ "version": "https://registry.npmjs.org/isstream/-/isstream-0.1.2.tgz",
+ "integrity": "sha1-R+Y/evVa+m+S4VAOaQ64uFKcCZo="
+ },
+ "jsbn": {
+ "version": "0.1.1",
+ "resolved": "https://registry.npmjs.org/jsbn/-/jsbn-0.1.1.tgz",
+ "integrity": "sha1-peZUwuWi3rXyAdls77yoDA7y9RM=",
+ "optional": true
+ },
+ "json-schema": {
+ "version": "0.2.3",
+ "resolved": "https://registry.npmjs.org/json-schema/-/json-schema-0.2.3.tgz",
+ "integrity": "sha1-tIDIkuWaLwWVTOcnvT8qTogvnhM="
+ },
+ "json-stable-stringify": {
+ "version": "1.0.1",
+ "resolved": "https://registry.npmjs.org/json-stable-stringify/-/json-stable-stringify-1.0.1.tgz",
+ "integrity": "sha1-mnWdOcXy/1A/1TAGRu1EX4jE+a8=",
+ "requires": {
+ "jsonify": "0.0.0"
+ }
+ },
+ "json-stringify-safe": {
+ "version": "5.0.1",
+ "resolved": "https://registry.npmjs.org/json-stringify-safe/-/json-stringify-safe-5.0.1.tgz",
+ "integrity": "sha1-Epai1Y/UXxmg9s4B1lcB4sc1tus="
+ },
+ "jsonify": {
+ "version": "0.0.0",
+ "resolved": "https://registry.npmjs.org/jsonify/-/jsonify-0.0.0.tgz",
+ "integrity": "sha1-LHS27kHZPKUbe1qu6PUDYx0lKnM="
+ },
+ "jsprim": {
+ "version": "1.4.0",
+ "resolved": "https://registry.npmjs.org/jsprim/-/jsprim-1.4.0.tgz",
+ "integrity": "sha1-o7h+QCmNjDgFUtjMdiigu5WiKRg=",
+ "requires": {
+ "assert-plus": "1.0.0",
+ "extsprintf": "1.0.2",
+ "json-schema": "0.2.3",
+ "verror": "1.3.6"
+ },
+ "dependencies": {
+ "assert-plus": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/assert-plus/-/assert-plus-1.0.0.tgz",
+ "integrity": "sha1-8S4PPF13sLHN2RRpQuTpbB5N1SU="
+ }
+ }
+ },
+ "lodash": {
+ "version": "4.17.4",
+ "resolved": "https://registry.npmjs.org/lodash/-/lodash-4.17.4.tgz",
+ "integrity": "sha1-eCA6TRwyiuHYbcpkYONptX9AVa4="
+ },
+ "long": {
+ "version": "3.2.0",
+ "resolved": "https://registry.npmjs.org/long/-/long-3.2.0.tgz",
+ "integrity": "sha1-2CG3E4yhy1gcFymQ7xTbIAtcR0s="
+ },
+ "lru-cache": {
+ "version": "4.1.1",
+ "resolved": "https://registry.npmjs.org/lru-cache/-/lru-cache-4.1.1.tgz",
+ "integrity": "sha512-q4spe4KTfsAS1SUHLO0wz8Qiyf1+vMIAgpRYioFYDMNqKfHQbg+AVDH3i4fvpl71/P1L0dBl+fQi+P37UYf0ew==",
+ "requires": {
+ "pseudomap": "1.0.2",
+ "yallist": "2.1.2"
+ }
+ },
+ "mime-db": {
+ "version": "1.27.0",
+ "resolved": "https://registry.npmjs.org/mime-db/-/mime-db-1.27.0.tgz",
+ "integrity": "sha1-gg9XIpa70g7CXtVeW13oaeVDbrE="
+ },
+ "mime-types": {
+ "version": "2.1.15",
+ "resolved": "https://registry.npmjs.org/mime-types/-/mime-types-2.1.15.tgz",
+ "integrity": "sha1-pOv1BkCUVpI3uM9wBGd20J/JKu0=",
+ "requires": {
+ "mime-db": "1.27.0"
+ }
+ },
+ "moment": {
+ "version": "2.18.1",
+ "resolved": "https://registry.npmjs.org/moment/-/moment-2.18.1.tgz",
+ "integrity": "sha1-w2GT3Tzhwu7SrbfIAtu8d6gbHA8="
+ },
+ "moment-timezone": {
+ "version": "0.5.13",
+ "resolved": "https://registry.npmjs.org/moment-timezone/-/moment-timezone-0.5.13.tgz",
+ "integrity": "sha1-mc5cfYJyYusPH3AgRBd/YHRde5A=",
+ "requires": {
+ "moment": "2.18.1"
+ }
+ },
+ "ms": {
+ "version": "2.0.0",
+ "resolved": "https://registry.npmjs.org/ms/-/ms-2.0.0.tgz",
+ "integrity": "sha1-VgiurfwAvmwpAd9fmGF4jeDVl8g="
+ },
+ "mysql2": {
+ "version": "1.3.6",
+ "resolved": "https://registry.npmjs.org/mysql2/-/mysql2-1.3.6.tgz",
+ "integrity": "sha1-1cFJF94LFYAabBiDWZvPb7J8Gh4=",
+ "requires": {
+ "cardinal": "1.0.0",
+ "denque": "1.1.1",
+ "generate-function": "2.0.0",
+ "iconv-lite": "0.4.18",
+ "long": "3.2.0",
+ "lru-cache": "4.1.1",
+ "named-placeholders": "1.1.1",
+ "object-assign": "4.1.1",
+ "readable-stream": "2.3.2",
+ "safe-buffer": "5.1.1",
+ "seq-queue": "0.0.5",
+ "sqlstring": "https://registry.npmjs.org/sqlstring/-/sqlstring-2.2.0.tgz"
+ },
+ "dependencies": {
+ "isarray": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/isarray/-/isarray-1.0.0.tgz",
+ "integrity": "sha1-u5NdSFgsuhaMBoNJV6VKPgcSTxE="
+ },
+ "readable-stream": {
+ "version": "2.3.2",
+ "resolved": "https://registry.npmjs.org/readable-stream/-/readable-stream-2.3.2.tgz",
+ "integrity": "sha1-WgTfBeT1f+Pw3Gj90R3FyXx+b00=",
+ "requires": {
+ "core-util-is": "https://registry.npmjs.org/core-util-is/-/core-util-is-1.0.2.tgz",
+ "inherits": "https://registry.npmjs.org/inherits/-/inherits-2.0.3.tgz",
+ "isarray": "1.0.0",
+ "process-nextick-args": "1.0.7",
+ "safe-buffer": "5.1.1",
+ "string_decoder": "1.0.3",
+ "util-deprecate": "1.0.2"
+ }
+ },
+ "string_decoder": {
+ "version": "1.0.3",
+ "resolved": "https://registry.npmjs.org/string_decoder/-/string_decoder-1.0.3.tgz",
+ "integrity": "sha512-4AH6Z5fzNNBcH+6XDMfA/BTt87skxqJlO0lAh3Dker5zThcAxG6mKz+iGu308UKoPPQ8Dcqx/4JhujzltRa+hQ==",
+ "requires": {
+ "safe-buffer": "5.1.1"
+ }
+ }
+ }
+ },
+ "named-placeholders": {
+ "version": "1.1.1",
+ "resolved": "https://registry.npmjs.org/named-placeholders/-/named-placeholders-1.1.1.tgz",
+ "integrity": "sha1-O3oNJiA910s6nfTJz7gnsvuQfmQ=",
+ "requires": {
+ "lru-cache": "2.5.0"
+ },
+ "dependencies": {
+ "lru-cache": {
+ "version": "2.5.0",
+ "resolved": "https://registry.npmjs.org/lru-cache/-/lru-cache-2.5.0.tgz",
+ "integrity": "sha1-2COIrpyWC+y+oMc7uet5tsbOmus="
+ }
+ }
+ },
+ "node-loggly-bulk": {
+ "version": "2.0.0",
+ "resolved": "https://registry.npmjs.org/node-loggly-bulk/-/node-loggly-bulk-2.0.0.tgz",
+ "integrity": "sha1-5M+UOY15jroSQwUsLybX3ZuSEE0=",
+ "requires": {
+ "json-stringify-safe": "5.0.1",
+ "request": "2.81.0",
+ "timespan": "2.3.0"
+ }
+ },
+ "oauth-sign": {
+ "version": "0.8.2",
+ "resolved": "https://registry.npmjs.org/oauth-sign/-/oauth-sign-0.8.2.tgz",
+ "integrity": "sha1-Rqarfwrq2N6unsBWV4C31O/rnUM="
+ },
+ "object-assign": {
+ "version": "4.1.1",
+ "resolved": "https://registry.npmjs.org/object-assign/-/object-assign-4.1.1.tgz",
+ "integrity": "sha1-IQmtx5ZYh8/AXLvUQsrIv7s2CGM="
+ },
+ "performance-now": {
+ "version": "0.2.0",
+ "resolved": "https://registry.npmjs.org/performance-now/-/performance-now-0.2.0.tgz",
+ "integrity": "sha1-M+8wxcd9TqIcWlOGnZG1bY8lVeU="
+ },
+ "process-nextick-args": {
+ "version": "1.0.7",
+ "resolved": "https://registry.npmjs.org/process-nextick-args/-/process-nextick-args-1.0.7.tgz",
+ "integrity": "sha1-FQ4gt1ZZCtP5EJPyWk8q2L/zC6M="
+ },
+ "pseudomap": {
+ "version": "1.0.2",
+ "resolved": "https://registry.npmjs.org/pseudomap/-/pseudomap-1.0.2.tgz",
+ "integrity": "sha1-8FKijacOYYkX7wqKw0wa5aaChrM="
+ },
+ "punycode": {
+ "version": "1.4.1",
+ "resolved": "https://registry.npmjs.org/punycode/-/punycode-1.4.1.tgz",
+ "integrity": "sha1-wNWmOycYgArY4esPpSachN1BhF4="
+ },
+ "qs": {
+ "version": "6.4.0",
+ "resolved": "https://registry.npmjs.org/qs/-/qs-6.4.0.tgz",
+ "integrity": "sha1-E+JtKK1rD/qpExLNO/cI7TUecjM="
+ },
+ "redeyed": {
+ "version": "1.0.1",
+ "resolved": "https://registry.npmjs.org/redeyed/-/redeyed-1.0.1.tgz",
+ "integrity": "sha1-6WwZO0DAgWsArshCaY5hGF5VSYo=",
+ "requires": {
+ "esprima": "3.0.0"
+ }
+ },
+ "request": {
+ "version": "2.81.0",
+ "resolved": "https://registry.npmjs.org/request/-/request-2.81.0.tgz",
+ "integrity": "sha1-xpKJRqDgbF+Nb4qTM0af/aRimKA=",
+ "requires": {
+ "aws-sign2": "0.6.0",
+ "aws4": "1.6.0",
+ "caseless": "0.12.0",
+ "combined-stream": "1.0.5",
+ "extend": "3.0.1",
+ "forever-agent": "0.6.1",
+ "form-data": "2.1.4",
+ "har-validator": "4.2.1",
+ "hawk": "3.1.3",
+ "http-signature": "1.1.1",
+ "is-typedarray": "1.0.0",
+ "isstream": "https://registry.npmjs.org/isstream/-/isstream-0.1.2.tgz",
+ "json-stringify-safe": "5.0.1",
+ "mime-types": "2.1.15",
+ "oauth-sign": "0.8.2",
+ "performance-now": "0.2.0",
+ "qs": "6.4.0",
+ "safe-buffer": "5.1.1",
+ "stringstream": "0.0.5",
+ "tough-cookie": "2.3.2",
+ "tunnel-agent": "0.6.0",
+ "uuid": "3.1.0"
+ }
+ },
+ "retry-as-promised": {
+ "version": "2.2.0",
+ "resolved": "https://registry.npmjs.org/retry-as-promised/-/retry-as-promised-2.2.0.tgz",
+ "integrity": "sha1-sEY9f9PPWy/tZFAKtui4pJxbjmw=",
+ "requires": {
+ "bluebird": "https://registry.npmjs.org/bluebird/-/bluebird-3.5.0.tgz",
+ "cross-env": "3.2.4",
+ "debug": "2.6.8"
+ }
+ },
+ "safe-buffer": {
+ "version": "5.1.1",
+ "resolved": "https://registry.npmjs.org/safe-buffer/-/safe-buffer-5.1.1.tgz",
+ "integrity": "sha512-kKvNJn6Mm93gAczWVJg7wH+wGYWNrDHdWvpUmHyEsgCtIwwo3bqPtV4tR5tuPaUhTOo/kvhVwd8XwwOllGYkbg=="
+ },
+ "semver": {
+ "version": "5.3.0",
+ "resolved": "https://registry.npmjs.org/semver/-/semver-5.3.0.tgz",
+ "integrity": "sha1-myzl094C0XxgEq0yaqa00M9U+U8="
+ },
+ "seq-queue": {
+ "version": "0.0.5",
+ "resolved": "https://registry.npmjs.org/seq-queue/-/seq-queue-0.0.5.tgz",
+ "integrity": "sha1-1WgS4cAXpuTnw+Ojeh2m143TyT4="
+ },
+ "sequelize": {
+ "version": "4.3.2",
+ "resolved": "https://registry.npmjs.org/sequelize/-/sequelize-4.3.2.tgz",
+ "integrity": "sha1-FHrRnW6lk++YME6MsLimkrS50aM=",
+ "requires": {
+ "bluebird": "https://registry.npmjs.org/bluebird/-/bluebird-3.5.0.tgz",
+ "cls-bluebird": "2.0.1",
+ "debug": "2.6.8",
+ "depd": "1.1.0",
+ "dottie": "2.0.0",
+ "env-cmd": "5.1.0",
+ "generic-pool": "3.1.7",
+ "inflection": "1.10.0",
+ "lodash": "4.17.4",
+ "moment": "2.18.1",
+ "moment-timezone": "0.5.13",
+ "retry-as-promised": "2.2.0",
+ "semver": "5.3.0",
+ "terraformer-wkt-parser": "1.1.2",
+ "toposort-class": "1.0.1",
+ "uuid": "3.1.0",
+ "validator": "6.3.0",
+ "wkx": "0.4.1"
+ }
+ },
+ "shebang-command": {
+ "version": "1.2.0",
+ "resolved": "https://registry.npmjs.org/shebang-command/-/shebang-command-1.2.0.tgz",
+ "integrity": "sha1-RKrGW2lbAzmJaMOfNj/uXer98eo=",
+ "requires": {
+ "shebang-regex": "1.0.0"
+ }
+ },
+ "shebang-regex": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/shebang-regex/-/shebang-regex-1.0.0.tgz",
+ "integrity": "sha1-2kL0l0DAtC2yypcoVxyxkMmO/qM="
+ },
+ "shimmer": {
+ "version": "1.1.0",
+ "resolved": "https://registry.npmjs.org/shimmer/-/shimmer-1.1.0.tgz",
+ "integrity": "sha1-l9c3cTf/u6tCVSLkKf4KqJpIizU="
+ },
+ "sntp": {
+ "version": "1.0.9",
+ "resolved": "https://registry.npmjs.org/sntp/-/sntp-1.0.9.tgz",
+ "integrity": "sha1-ZUEYTMkK7qbG57NeJlkIJEPGYZg=",
+ "requires": {
+ "hoek": "2.16.3"
+ }
+ },
+ "sqlstring": {
+ "version": "https://registry.npmjs.org/sqlstring/-/sqlstring-2.2.0.tgz",
+ "integrity": "sha1-wxNcTqirzX5+50GklmqJHYak8ZE="
+ },
+ "sshpk": {
+ "version": "1.13.1",
+ "resolved": "https://registry.npmjs.org/sshpk/-/sshpk-1.13.1.tgz",
+ "integrity": "sha1-US322mKHFEMW3EwY/hzx2UBzm+M=",
+ "requires": {
+ "asn1": "0.2.3",
+ "assert-plus": "1.0.0",
+ "bcrypt-pbkdf": "1.0.1",
+ "dashdash": "1.14.1",
+ "ecc-jsbn": "0.1.1",
+ "getpass": "0.1.7",
+ "jsbn": "0.1.1",
+ "tweetnacl": "0.14.5"
+ },
+ "dependencies": {
+ "assert-plus": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/assert-plus/-/assert-plus-1.0.0.tgz",
+ "integrity": "sha1-8S4PPF13sLHN2RRpQuTpbB5N1SU="
+ }
+ }
+ },
+ "stringstream": {
+ "version": "0.0.5",
+ "resolved": "https://registry.npmjs.org/stringstream/-/stringstream-0.0.5.tgz",
+ "integrity": "sha1-TkhM1N5aC7vuGORjB3EKioFiGHg="
+ },
+ "terraformer": {
+ "version": "1.0.8",
+ "resolved": "https://registry.npmjs.org/terraformer/-/terraformer-1.0.8.tgz",
+ "integrity": "sha1-UeCtiXRvzyFh3G9lqnDkI3fItZM=",
+ "requires": {
+ "@types/geojson": "1.0.2"
+ }
+ },
+ "terraformer-wkt-parser": {
+ "version": "1.1.2",
+ "resolved": "https://registry.npmjs.org/terraformer-wkt-parser/-/terraformer-wkt-parser-1.1.2.tgz",
+ "integrity": "sha1-M2oMj8gglKWv+DKI9prt7NNpvww=",
+ "requires": {
+ "terraformer": "1.0.8"
+ }
+ },
+ "timespan": {
+ "version": "2.3.0",
+ "resolved": "https://registry.npmjs.org/timespan/-/timespan-2.3.0.tgz",
+ "integrity": "sha1-SQLOBAvRPYRcj1myfp1ZutbzmSk="
+ },
+ "toposort-class": {
+ "version": "1.0.1",
+ "resolved": "https://registry.npmjs.org/toposort-class/-/toposort-class-1.0.1.tgz",
+ "integrity": "sha1-f/0feMi+KMO6Rc1OGj9e4ZO9mYg="
+ },
+ "tough-cookie": {
+ "version": "2.3.2",
+ "resolved": "https://registry.npmjs.org/tough-cookie/-/tough-cookie-2.3.2.tgz",
+ "integrity": "sha1-8IH3bkyFcg5sN6X6ztc3FQ2EByo=",
+ "requires": {
+ "punycode": "1.4.1"
+ }
+ },
+ "tunnel-agent": {
+ "version": "0.6.0",
+ "resolved": "https://registry.npmjs.org/tunnel-agent/-/tunnel-agent-0.6.0.tgz",
+ "integrity": "sha1-J6XeoGs2sEoKmWZ3SykIaPD8QP0=",
+ "requires": {
+ "safe-buffer": "5.1.1"
+ }
+ },
+ "tweetnacl": {
+ "version": "0.14.5",
+ "resolved": "https://registry.npmjs.org/tweetnacl/-/tweetnacl-0.14.5.tgz",
+ "integrity": "sha1-WuaBd/GS1EViadEIr6k/+HQ/T2Q=",
+ "optional": true
+ },
+ "util-deprecate": {
+ "version": "1.0.2",
+ "resolved": "https://registry.npmjs.org/util-deprecate/-/util-deprecate-1.0.2.tgz",
+ "integrity": "sha1-RQ1Nyfpw3nMnYvvS1KKJgUGaDM8="
+ },
+ "uuid": {
+ "version": "3.1.0",
+ "resolved": "https://registry.npmjs.org/uuid/-/uuid-3.1.0.tgz",
+ "integrity": "sha512-DIWtzUkw04M4k3bf1IcpS2tngXEL26YUD2M0tMDUpnUrz2hgzUBlD55a4FjdLGPvfHxS6uluGWvaVEqgBcVa+g=="
+ },
+ "validator": {
+ "version": "6.3.0",
+ "resolved": "https://registry.npmjs.org/validator/-/validator-6.3.0.tgz",
+ "integrity": "sha1-R84j7Y1Ord+p1LjvAHG2zxB418g="
+ },
+ "verror": {
+ "version": "1.3.6",
+ "resolved": "https://registry.npmjs.org/verror/-/verror-1.3.6.tgz",
+ "integrity": "sha1-z/XfEpRtKX0rqu+qJoniW+AcAFw=",
+ "requires": {
+ "extsprintf": "1.0.2"
+ }
+ },
+ "which": {
+ "version": "1.2.14",
+ "resolved": "https://registry.npmjs.org/which/-/which-1.2.14.tgz",
+ "integrity": "sha1-mofEN48D6CfOyvGs31bHNsAcFOU=",
+ "requires": {
+ "isexe": "2.0.0"
+ }
+ },
+ "winston": {
+ "version": "2.3.1",
+ "resolved": "https://registry.npmjs.org/winston/-/winston-2.3.1.tgz",
+ "integrity": "sha1-C0hCDZeMAYBM8CMLZIhhWYIloRk=",
+ "requires": {
+ "async": "1.0.0",
+ "colors": "1.0.3",
+ "cycle": "1.0.3",
+ "eyes": "0.1.8",
+ "isstream": "0.1.2",
+ "stack-trace": "0.0.10"
+ },
+ "dependencies": {
+ "async": {
+ "version": "1.0.0",
+ "resolved": "https://registry.npmjs.org/async/-/async-1.0.0.tgz",
+ "integrity": "sha1-+PwEyjoTeErenhZBr5hXjPvWR6k="
+ },
+ "colors": {
+ "version": "1.0.3",
+ "resolved": "https://registry.npmjs.org/colors/-/colors-1.0.3.tgz",
+ "integrity": "sha1-BDP0TYCWgP3rYO0mDxsMJi6CpAs="
+ },
+ "cycle": {
+ "version": "1.0.3",
+ "resolved": "https://registry.npmjs.org/cycle/-/cycle-1.0.3.tgz",
+ "integrity": "sha1-IegLK+hYD5i0aPN5QwZisEbDStI="
+ },
+ "eyes": {
+ "version": "0.1.8",
+ "resolved": "https://registry.npmjs.org/eyes/-/eyes-0.1.8.tgz",
+ "integrity": "sha1-Ys8SAjTGg3hdkCNIqADvPgzCC8A="
+ },
+ "isstream": {
+ "version": "0.1.2",
+ "resolved": "https://registry.npmjs.org/isstream/-/isstream-0.1.2.tgz",
+ "integrity": "sha1-R+Y/evVa+m+S4VAOaQ64uFKcCZo="
+ },
+ "stack-trace": {
+ "version": "0.0.10",
+ "resolved": "https://registry.npmjs.org/stack-trace/-/stack-trace-0.0.10.tgz",
+ "integrity": "sha1-VHxws0fo0ytOEI6hoqFZ5f3eGcA="
+ }
+ }
+ },
+ "winston-loggly-bulk": {
+ "version": "2.0.0",
+ "resolved": "https://registry.npmjs.org/winston-loggly-bulk/-/winston-loggly-bulk-2.0.0.tgz",
+ "integrity": "sha1-olOCeBihS7b+CwbMHCLY0+19EoA=",
+ "requires": {
+ "node-loggly-bulk": "2.0.0",
+ "winston": "2.3.1"
+ }
+ },
+ "wkx": {
+ "version": "0.4.1",
+ "resolved": "https://registry.npmjs.org/wkx/-/wkx-0.4.1.tgz",
+ "integrity": "sha1-L8FxtenLVcYlb+9L3h8hvkE77+4=",
+ "requires": {
+ "@types/node": "6.0.82"
+ }
+ },
+ "yallist": {
+ "version": "2.1.2",
+ "resolved": "https://registry.npmjs.org/yallist/-/yallist-2.1.2.tgz",
+ "integrity": "sha1-HBH5IY8HYImkfdUS+TxmmaaoHVI="
+ }
+ }
+}
diff --git a/package.json b/package.json
index ba632fc..469e222 100644
--- a/package.json
+++ b/package.json
@@ -5,11 +5,10 @@
"main": "model.js",
"dependencies": {
"amqplib": "^0.5.1",
- "mysql": "^2.13.0",
- "sequelize": "^3.30.4",
- "sleep-promise": "^2.0.0",
+ "mysql2": "^1.3.6",
+ "sequelize": "^4.3.2",
"winston": "^2.3.1",
- "winston-loggly-bulk": "^1.4.2"
+ "winston-loggly-bulk": "^2.0.0"
},
"devDependencies": {},
"scripts": {
diff --git a/sequelize-bulkinsert.patch b/sequelize-bulkinsert.patch
new file mode 100644
index 0000000..eaf793b
--- /dev/null
+++ b/sequelize-bulkinsert.patch
@@ -0,0 +1,19 @@
+*** /tmp/a 2017-07-14 18:23:00.471411882 +0200
+--- node_modules/sequelize/lib/model.js 2017-07-14 18:25:01.281686417 +0200
+***************
+*** 2349,2360 ****
+--- 2349,2362 ----
+ }
+
+ return this.QueryInterface.bulkInsert(this.getTableName(options), records, options, attributes).then(results => {
++ /*
+ if (Array.isArray(results)) {
+ results.forEach((result, i) => {
+ instances[i].set(this.primaryKeyAttribute, result[this.rawAttributes[this.primaryKeyAttribute].field], {raw: true});
+ });
+ }
+ return results;
++ */
+ });
+ }
+ }).then(() => {
diff --git a/worker.js b/worker.js
index 5dd1411..e190709 100644
--- a/worker.js
+++ b/worker.js
@@ -13,11 +13,7 @@ const amqp = require("amqplib"),
winston = require("winston"),
loggly = require("winston-loggly-bulk"),
Seq = require("sequelize"),
- api_name_mappings = require("../orm/mappings").map,
- isAbility = require("../orm/mappings").isAbility,
- isItem = require("../orm/mappings").isItem,
- isHero = require("../orm/mappings").isHero,
- sleep = require("sleep-promise");
+ api_name_mappings = require("../orm/mappings").map;
const RABBITMQ_URI = process.env.RABBITMQ_URI,
DATABASE_URI = process.env.DATABASE_URI,
@@ -69,8 +65,9 @@ function snakeCaseKeys(obj) {
return obj;
}
-// split an array into arrays of max chunksize
-function* chunks(arr) {
+// split an Set() into arrays of max chunksize
+function* chunks(data) {
+ const arr = [...data]; // TODO maybe slice the Set?
for (let c=0, len=arr.length; c<len; c+=CHUNKSIZE)
yield arr.slice(c, c+CHUNKSIZE);
}
@@ -91,28 +88,26 @@ function flatten(obj) {
return snakeCaseKeys(o);
}
-// main code
-(async () => {
- let seq, model, rabbit, ch;
+amqp.connect(RABBITMQ_URI).then(async (rabbit) => {
+ global.process.on("SIGINT", () => {
+ rabbit.close();
+ global.process.exit();
+ });
// connect to rabbit & db
- while (true) {
- try {
- seq = new Seq(DATABASE_URI, {
- logging: false,
- max: MAXCONNS
- });
- rabbit = await amqp.connect(RABBITMQ_URI, { heartbeat: 120 });
- ch = await rabbit.createChannel();
- await ch.assertQueue(QUEUE, {durable: true});
- break;
- } catch (err) {
- logger.error("Error connecting", err);
- await sleep(5000);
- }
- }
+ const seq = new Seq(DATABASE_URI, {
+ //logging: false,
+ max: MAXCONNS
+ });
+
+ const ch = await rabbit.createChannel();
+ await ch.assertQueue(QUEUE, { durable: true });
+ await ch.assertQueue(QUEUE + "_failed", { 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);
- model = require("../orm/model")(seq, Seq);
+ const model = require("../orm/model")(seq, Seq);
// performance logging
let load_timer = undefined,
@@ -124,8 +119,7 @@ function flatten(obj) {
hero_db_map = new Map(), // "*SAW*" to id
series_db_map = new Map(), // date to series id
game_mode_db_map = new Map(), // "ranked" to id
- role_db_map = new Map(), // "captain" to id
- hero_role_map = new Map(); // SAW.id to "carry"
+ role_db_map = new Map(); // "captain" to id
// populate maps
await Promise.all([
@@ -143,44 +137,61 @@ function flatten(obj) {
model.Role.findAll()
.map((role) => role_db_map.set(role.name, role.id))
]);
-
- // 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);
+ if (item_db_map.size == 0 ||
+ hero_db_map.size == 0 ||
+ series_db_map.size == 0 ||
+ game_mode_db_map.size == 0 ||
+ role_db_map == 0) {
+ logger.error("mapping tables are not seeded!!! quitting");
+ global.process.exit();
+ }
// buffers that will be filled until BATCHSIZE is reached
// to make db transactions more efficient
let player_data = new Set(),
match_data = new Set(),
- telemetry_data = new Set(),
msg_buffer = new Set();
ch.consume(QUEUE, async (msg) => {
- if (msg.properties.type == "player") {
- // bridge sends a single object
- player_data.add(JSON.parse(msg.content));
- msg_buffer.add(msg);
- }
- if (msg.properties.type == "match") {
- // apigrabber sends a single object
- const match = JSON.parse(msg.content);
- // deduplicate and reject immediately
- if (await model.Match.count({ where: { api_id: match.id } }) > 0) {
- if (msg.properties.headers.notify != undefined) {
- await ch.publish("amq.topic",
- msg.properties.headers.notify,
- new Buffer("matches_dupe"))
- // send match_dupe to web player.ign.api_id
- await ch.publish("amq.topic",
- msg.properties.headers.notify + "." + match.id,
- new Buffer("match_dupe"))
+ switch (msg.properties.type) {
+ case "player":
+ // bridge sends a single object
+ let player = JSON.parse(msg.content);
+ // player objects that arrive here came from a search
+ // with search, bridge can't update last_update
+ player.last_update = seq.fn("NOW");
+ player_data.add(player);
+ msg_buffer.add(msg);
+ break;
+ case "match":
+ // apigrabber sends a single object
+ const match = JSON.parse(msg.content);
+ // deduplicate and reject immediately
+ if (await model.Match.count({ where: { api_id: match.id } }) > 0) {
+ if (msg.properties.headers.notify) {
+ await ch.publish("amq.topic",
+ msg.properties.headers.notify,
+ new Buffer("matches_dupe"))
+ // send match_dupe to web player.ign.api_id
+ await ch.publish("amq.topic",
+ msg.properties.headers.notify + "." + match.id,
+ new Buffer("match_dupe"))
+ }
+ await ch.nack(msg, false, false);
+ } else if (match.rosters.length < 2 || match.rosters[0].id == "null") {
+ // it is really `"null"`.
+ // reject invalid matches (handling API bugs)
+ await ch.nack(msg, false, false);
+ await ch.sendToQueue(QUEUE + "_failed", msg.content, {
+ persistent: true,
+ headers: msg.properties.headers
+ });
+ } else {
+ // all good
+ match_data.add(match);
+ msg_buffer.add(msg);
}
- } else match_data.add(match);
- msg_buffer.add(msg);
- }
- if (msg.properties.type == "telemetry") {
- telemetry_data.add(JSON.parse(msg.content));
- msg_buffer.add(msg);
+ break;
}
// fill queue until batchsize or idle
@@ -188,108 +199,136 @@ function flatten(obj) {
if (profiler == undefined) profiler = logger.startTimer();
// timeout after first job
if (load_timer == undefined)
- load_timer = setTimeout(process, LOAD_TIMEOUT);
+ load_timer = setTimeout(tryProcess, LOAD_TIMEOUT);
// timeout after last job
if (idle_timer != undefined)
clearTimeout(idle_timer);
- idle_timer = setTimeout(process, IDLE_TIMEOUT);
+ idle_timer = setTimeout(tryProcess, IDLE_TIMEOUT);
// maximum data pressure
- if (match_data.size + player_data.size + telemetry_data.size == BATCHSIZE)
- await process();
+ if (match_data.size + player_data.size == BATCHSIZE)
+ await tryProcess();
}, { noAck: false });
- // finish a whole batch
- async function process() {
+ // wrap process() in message handler
+ async function tryProcess() {
+ const msgs = new Set(msg_buffer);
+ msg_buffer.clear();
+
profiler.done("buffer filled");
profiler = undefined;
logger.info("processing batch", {
players: player_data.size,
- matches: match_data.size,
- telemetries: telemetry_data.size
+ matches: match_data.size
});
// clean up to allow processor to accept while we wait for db
clearTimeout(idle_timer);
clearTimeout(load_timer);
- const player_objects = new Set(player_data),
- match_objects = new Set(match_data),
- telemetry_objects = new Set(telemetry_data),
- msgs = new Set(msg_buffer);
idle_timer = undefined;
load_timer = undefined;
+
+ if (player_data.size + match_data.size == 0) {
+ logger.info("buffers empty, nothing to do");
+ return;
+ }
+
+ const player_objects = new Set(player_data),
+ match_objects = new Set(match_data);
player_data.clear();
match_data.clear();
- telemetry_data.clear();
- msg_buffer.clear();
- const processed_players = new Set(); // to sort out duplicates
+ try {
+ await process(player_objects, match_objects);
- // aggregate record objects to do a bulk insert
- let match_records = [],
- roster_records = [],
- participant_records = [],
- participant_stats_records = [],
- participant_phase_records = [], // Telemetry
- player_records = [],
- player_records_direct = [], // via `/players`
- asset_records = [];
+ logger.info("acking batch", { size: msgs.size });
+ await Promise.map(msgs, async (m) => await ch.ack(m));
- // populate `_records`
- // data from `/players`
- // `each` executes serially so there are
- // no race conditions within one batch
- await Promise.each(player_objects, async (p) => {
- const player = flatten(p);
- if (processed_players.has(player.api_id)) {
- // duplicate within one batch
- logger.warn("got player in additional region",
- { name: player.name, region: player.shard_id });
- const duplicate = player_records_direct.find((pr) =>
- pr.api_id == player.api_id);
- if (duplicate == undefined) return; // TODO wtf???
- if (duplicate.created_at < player.created_at) {
- // replace by newer one as below
- player_records_direct.splice(
- player_records_direct.indexOf(duplicate), 1);
- } else {
- logger.warn("ignoring player from additional region",
- { name: player.name, region: player.shard_id });
- return;
+ // notify web
+ await Promise.map(msgs, async (m) => {
+ if (m.properties.headers.notify == undefined) return;
+ let notif = "";
+ switch (m.properties.type) {
+ // new match
+ case "match":
+ notif = "matches_update";
+ // notify player.name.api_id about match_update
+ await Promise.map(match_objects, async (mat) =>
+ await ch.publish("amq.topic",
+ m.properties.headers.notify + "." + mat.id,
+ new Buffer("match_update"))
+ );
+ await ch.publish("amq.topic", m.properties.headers.notify,
+ new Buffer("match_update"));
+ break;
+ case "player":
+ // player obj updated
+ notif = "stats_update";
+ break;
}
+ });
+ // …global about new matches
+ if (match_objects.length > 0)
+ await ch.publish("amq.topic", "global", new Buffer("matches_update"));
+ // notify follow up services
+ if (DOANALYZEMATCH)
+ await Promise.each(match_objects, async (m) =>
+ await ch.sendToQueue(ANALYZE_QUEUE, new Buffer(m.id),
+ { persistent: true }));
+ } catch (err) {
+ if (err instanceof Seq.TimeoutError) {
+ // deadlocks / timeout
+ logger.error("SQL error", err);
+ await Promise.map(msgs, async (m) =>
+ await ch.nack(m, false, true)); // retry
} else {
- processed_players.add(player.api_id);
+ // log, move to error queue and NACK
+ logger.error(err);
+ await Promise.map(msgs, async (m) => {
+ await ch.sendToQueue(QUEUE + "_failed", m.content, {
+ persistent: true,
+ headers: m.properties.headers
+ });
+ await ch.nack(m, false, false);
+ });
}
- // player objects that arrive here came from a search
- // with search, updater can't update last_update
- player.last_update = seq.fn("NOW");
+ }
+ }
+
+ // finish a whole batch
+ async function process(player_objects, match_objects) {
+ // aggregate record objects to do a bulk insert
+ let match_records = new Set(),
+ roster_records = new Set(),
+ participant_records = new Set(),
+ participant_stats_records = new Set(),
+ players = new Map(),
+ player_records = new Set(),
+ player_records_dates = new Set(), // from /players
+ asset_records = new Set();
+
+ // populate `_records`
+ // data from `/players`
+ player_objects.forEach((p) => {
+ let player = flatten(p);
+ player.created_at = new Date(Date.parse(player.created_at));
+
logger.info("processing player",
{ name: player.name, region: player.shard_id });
- // duplicate in batch and db
- // check whether there is a player in db
- // that has a more recent `created_at`
- // this is only the case with region changes
- const count = await model.Player.count({ where: {
- api_id: player.api_id,
- created_at: {
- $gt: player.created_at // equal: just update last_update
+ if (!players.has(player.api_id)) {
+ players.set(player.api_id, player);
+ } else { // or a player object that is more recent than the buffer's
+ if (players.get(player.api_id).created_at < player.created_at) {
+ logger.info("buffer has same more recent direct player object, overwriting");
+ players.set(player.api_id, player);
}
- }});
- if (count > 0) {
- logger.warn("ignoring player who seems to have switched from region",
- { name: player.name, region: player.shard_id });
- return;
- } else player_records_direct.push(player);
- });
-
- // reject invalid matches (handling API bugs) TODO should happen in apigrabber
- match_objects.forEach((match, idx) => {
- // it is really `"null"`.
- if (match.rosters[0].id == "null") delete match_objects[idx];
+ }
});
// data from `/matches`
match_objects.forEach((match) => {
+ match.created_at = new Date(Date.parse(match.created_at));
+
// flatten jsonapi nested response into our db structure-like shape
// also, push missing fields
match.rosters = match.rosters.map((roster) => {
@@ -338,7 +377,12 @@ function flatten(obj) {
Object.keys(pas.itemSells)
.map((key) => item_id(key) + ";" + pas.itemSells[key]).join(",");
- participant.player.attributes.shardId = participant.player.attributes.shardId || participant.attributes.shardId;
+ participant.player.attributes.shardId = participant.player.attributes.shardId
+ || participant.attributes.shardId;
+ if (participant.player.attributes.createdAt)
+ participant.player.createdAt =
+ new Date(Date.parse(participant.player.attributes.createdAt));
+ else participant.player.created_at = match.createdAt;
participant.player = flatten(participant.player);
return flatten(participant);
});
@@ -352,476 +396,128 @@ function flatten(obj) {
match = flatten(match);
// after conversion, create the array of records
- match_records.push(match);
+ match_records.add(match);
match.rosters.forEach((r) => {
- roster_records.push(r);
+ roster_records.add(r);
r.participants.forEach((p) => {
const p_pstats = calculate_participant_stats(match, r, p);
// participant gets split into participant and p_stats
- participant_records.push(p_pstats[0]);
- participant_stats_records.push(p_pstats[1]);
- // deduplicate player
- // in a batch, it is very likely that players are duplicated
- // so this improves performance a bit
- if (!processed_players.has(p.player.api_id)) {
- processed_players.add(p.player.api_id);
- player_records.push(p.player);
+ participant_records.add(p_pstats[0]);
+ participant_stats_records.add(p_pstats[1]);
+
+ // if match.included has an unknown player
+ if (!players.has(p.player.api_id))
+ players.set(p.player.api_id, p.player);
+ else {
+ // or a player object that is more recent than the buffer's
+ if (players.get(p.player.api_id).created_at < p.player.created_at) {
+ if (players.get(p.player.api_id).last_update != undefined) {
+ logger.info("buffer has same more recent indirect player object, overwriting direct");
+ // indirect overwrites direct's stats; keep last_update and created_at
+ p.player.created_at = players.get(p.player.api_id).created_at;
+ p.player.last_update = players.get(p.player.api_id).last_update;
+ }
+ // else direct/indirect overwrites indirect
+ players.set(p.player.api_id, p.player);
+ }
}
});
});
- match.assets.forEach((a) => asset_records.push(a));
+ match.assets.forEach((a) => asset_records.add(a));
});
-
- // data from Telemetry, one phase (early/mid/late/…) per obj
- await Promise.map(telemetry_objects, async (telemetry) => {
- if (telemetry.data.length == 0) return; // TODO rm me
- // api -> telemetry format
- const sideToTeam = (s) => s == "left/blue"? "Left" : "Right",
- // yes there is yet another format and yes it's strings
- sideToTeamNo = (s) => s == "left/blue"? "1" : "2";
- // get match participant references
- const participants = await model.Participant.findAll({
- where: { match_api_id: telemetry.match_api_id },
- include: [ { // TODO rm once pushed to participant
- model: model.Roster,
- attributes: [ "side" ]
- } ]
- }).map((p) => { return {
- api_id: p.api_id,
- player_api_id: p.player_api_id,
- actor: p.actor,
- team: sideToTeam(p.roster.side),
- teamNo: sideToTeamNo(p.roster.side)
- } });
-
- // seconds since epoch; first spawn time
- const matchstart = new Date(Date.parse(telemetry.match_start)).getTime() / 1000;
-
- // link participant <-> Telemetry actor/target
- // available as `.actor` or as `.target`
- telemetry.data.forEach((t) => {
- // seconds after this phase's start
- t.offset = new Date(Date.parse(t.time)).getTime() / 1000 - matchstart;
-
- // linking
- if (t.type == "HeroSelect")
- t.actor = participants.filter((p) =>
- p.player_api_id == t.payload.Player)[0];
- if (t.type == "BuyItem"
- || t.type == "SellItem"
- || t.type == "UseItemAbility"
- || t.type == "LearnAbility"
- || t.type == "UseAbility"
- || t.type == "LevelUp")
- t.actor = participants.filter((p) =>
- p.actor == t.payload.Actor
- && p.team == t.payload.Team)[0];
- if (t.type == "UseItemAbility"
- || t.type == "UseAbility")
- t.target = participants.filter((p) =>
- p.actor == t.payload.TargetActor
- && p.team != t.payload.Team)[0];
- // damage actor
- if ((t.type == "DealDamage"
- || t.type == "KillActor")
- && t.payload.IsHero == 1)
- t.actor = participants.filter((p) =>
- p.actor == t.payload.Actor
- && p.team == t.payload.Team)[0];
- // damage target
- if (t.type == "DealDamage"
- && t.payload.TargetIsHero == 1)
- t.target = participants.filter((p) =>
- p.actor == t.payload.Target
- && p.team != t.payload.Team)[0];
- // kill target
- if (t.type == "KillActor"
- && t.payload.TargetIsHero == 1)
- t.target = participants.filter((p) =>
- p.actor == t.payload.Killed
- && p.team == t.payload.KilledTeam)[0];
- });
- /*
- telemetry.data.forEach((ev) => { // TODO debug
- if (ev.payload.Ability == undefined || ev.payload.IsHero == 0) return;
- if (!api_name_mappings.has(ev.payload.Ability))
- console.error("ab to name map missing", ev.payload.Ability);
- });
- */
- const participants_phase = participants.map((p) => { return {
- // TODO ban data workaround
- start: telemetry.start < 0? 0 : telemetry.start, // in seconds
- end: telemetry.end,
- participant_api_id: p.api_id,
-
- kills: telemetry.data.filter((ev) =>
- ev.actor == p
- && ev.type == "KillActor"
- && ev.payload.IsHero == 1
- && ev.payload.TargetIsHero == 1
- ).length,
- deaths: telemetry.data.filter((ev) =>
- ev.target == p
- && ev.type == "KillActor"
- && ev.payload.TargetIsHero == 1
- ).length,
- // assists missing in data
- minion_kills: telemetry.data.filter((ev) =>
- ev.actor == p
- && ev.type == "KillActor"
- && ["*JungleMinion_TreeEnt*",
- "*Neutral_JungleMinion_DefaultBig*",
- "*Neutral_JungleMinion_DefaultSmall*",
- "*LeadMinion*",
- "*RangedMinion*",
- "*TankMinion*"
- ].indexOf(ev.payload.Killed) != -1
- ).length,
- jungle_kills: telemetry.data.filter((ev) =>
- ev.actor == p
- && ev.type == "KillActor"
- && ["*JungleMinion_TreeEnt*",
- "*Neutral_JungleMinion_DefaultBig*",
- "*Neutral_JungleMinion_DefaultSmall*"
- ].indexOf(ev.payload.Killed) != -1
- ).length,
- non_jungle_minion_kills: telemetry.data.filter((ev) =>
- ev.actor == p
- && ev.type == "KillActor"
- && ["*LeadMinion*",
- "*RangedMinion*",
- "*TankMinion*"
- ].indexOf(ev.payload.Killed) != -1
- ).length,
- crystal_mine_captures: telemetry.data.filter((ev) =>
- ev.actor == p
- && ev.type == "KillActor"
- && ev.payload.Killed == "*JungleMinion_CrystalMiner*"
- ).length,
- gold_mine_captures: telemetry.data.filter((ev) =>
- ev.actor == p
- && ev.type == "KillActor"
- && ev.payload.Killed == "*JungleMinion_GoldMiner*"
- ).length,
- kraken_captures: telemetry.data.filter((ev) =>
- ev.actor == p
- && ev.type == "KillActor"
- && ev.payload.Killed == "*Kraken_Jungle*"
- ).length,
- turret_captures: telemetry.data.filter((ev) =>
- ev.actor == p
- && ev.type == "KillActor"
- && (ev.payload.Killed == "*Turret*"
- || ev.payload.Killed == "*VainTurret*")
- ).length,
- // TODO Telemetry does not give accurate LifetimeGold
- gold: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "LevelUp"
- && ev.payload.LifetimeGold > acc
- ? ev.payload.LifetimeGold
- : acc
- , null),
- dmg_true_hero: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ev.payload.TargetIsHero == 1
- ? acc + ev.payload.Damage
- : acc
- , 0),
- dmg_true_kraken: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ["*Kraken_Jungle*",
- "*Kraken_Captured*"
- ].indexOf(ev.payload.Target) != -1
- ? acc + ev.payload.Damage
- : acc
- , 0),
- dmg_true_turret: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ev.payload.Target == "*Turret*"
- ? acc + ev.payload.Damage
- : acc
- , 0),
- dmg_true_vain_turret: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ev.payload.Target == "*VainTurret*"
- ? acc + ev.payload.Damage
- : acc
- , 0),
- dmg_true_others: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ev.payload.TargetIsHero == 0
- ? acc + ev.payload.Damage
- : acc
- , 0),
- dmg_dealt_hero: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ev.payload.TargetIsHero == 1
- ? acc + ev.payload.Delt
- : acc
- , 0),
- dmg_dealt_kraken: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ["*Kraken_Jungle*",
- "*Kraken_Captured*"
- ].indexOf(ev.payload.Target) != -1
- ? acc + ev.payload.Delt
- : acc
- , 0),
- dmg_dealt_turret: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ev.payload.Target == "*Turret*"
- ? acc + ev.payload.Delt
- : acc
- , 0),
- dmg_dealt_vain_turret: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ev.payload.Target == "*VainTurret*"
- ? acc + ev.payload.Delt
- : acc
- , 0),
- dmg_dealt_others: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ev.payload.TargetIsHero == 0
- ? acc + ev.payload.Delt
- : acc
- , 0),
- dmg_rcvd_dealt_hero: telemetry.data.reduce((acc, ev) =>
- ev.target == p
- && ev.type == "DealDamage"
- && ev.payload.IsHero == 1
- ? acc + ev.payload.Delt
- : acc
- , 0),
- dmg_rcvd_true_hero: telemetry.data.reduce((acc, ev) =>
- ev.target == p
- && ev.type == "DealDamage"
- && ev.payload.IsHero == 1
- ? acc + ev.payload.Damage
- : acc
- , 0),
- dmg_rcvd_dealt_others: telemetry.data.reduce((acc, ev) =>
- ev.target == p
- && ev.type == "DealDamage"
- && ev.payload.IsHero == 0
- ? acc + ev.payload.Delt
- : acc
- , 0),
- dmg_rcvd_true_others: telemetry.data.reduce((acc, ev) =>
- ev.target == p
- && ev.type == "DealDamage"
- && ev.payload.IsHero == 0
- ? acc + ev.payload.Damage
- : acc
- , 0),
- hero_level: telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "LevelUp"
- && ev.payload.Level > acc
- ? ev.payload.Level
- : acc
- , -1),
- items: null, // TODO
- item_grants: JSON.stringify(telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "BuyItem"
- ? acc.concat(item_db_map.get(api_name_mappings.get(ev.payload.Item)))
- : acc
- , [])),
- item_sells: JSON.stringify(telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "SellItem"
- ? acc.concat(item_db_map.get(api_name_mappings.get(ev.payload.Item)))
- : acc
- , [])),
- ability_levels: JSON.stringify(telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "LearnAbility"
- ? acc.concat([ [ api_name_mappings.get(ev.payload.Ability).split(" ")[1],
- ev.offset ] ])
- : acc
- , [])),
- ability_uses: JSON.stringify(telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "UseAbility"
- && ["A", "B", "C"].indexOf(
- api_name_mappings.get(ev.payload.Ability).split(" ")[1]
- ) != -1
- ? acc.concat([ [ api_name_mappings.get(ev.payload.Ability).split(" ")[1],
- ev.offset ] ])
- : acc
- , [])),
- ability_damage: JSON.stringify(telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "DealDamage"
- && ev.payload.IsHero == 1
- && isAbility(ev.payload.Source) // TODO
- && api_name_mappings.has(ev.payload.Source)
- && ["A", "B", "C"].indexOf(
- api_name_mappings.get(ev.payload.Source).split(" ")[1]
- ) != -1 // TODO refactor here
- ? acc.concat([ [ api_name_mappings.get(ev.payload.Source).split(" ")[1],
- ev.payload.Damage, ev.offset ] ])
- : acc
- , [])),
- item_uses: JSON.stringify(telemetry.data.reduce((acc, ev) =>
- ev.actor == p
- && ev.type == "UseItemAbility"
- ? acc.concat([ [ item_db_map.get(api_name_mappings.get(ev.payload.Ability)),
- ev.offset ] ])
- : acc
- , [])),
- player_damage: null, // TODO
- items: null, // TODO
- draft_position: telemetry.data.filter((ev) =>
- ev.type == "HeroSelect").indexOf(
- telemetry.data.filter((ev) =>
- ev.type == "HeroSelect"
- && ev.actor == p)[0]),
- ban: hero_db_map.get(api_name_mappings.get(
- telemetry.data
- .filter((ev) => ev.type == "HeroBan" &&
- ev.payload.Team == p.teamNo)
- .map((sel) => sel.payload.Hero)[0] // can be null
- )),
- pick: hero_db_map.get(api_name_mappings.get(
- telemetry.data
- .filter((ev) =>
- ev.type == "HeroSelect"
- && ev.actor == p)
- .map((sel) => sel.payload.Hero)[0] // can be null
- )),// traits calculated later
+ // player.last_update = last time bridge ran an update
+ // player.created_at = last match's created_at
+ // that has been added to the database after fetching full history
+ // Here, all attributes *except* these two will be overwritten
+ // if they are more recent in match.included data
+ await Promise.map(players.values(), async (player) => {
+ const count = await model.Player.count({ where: {
+ api_id: player.api_id,
+ created_at: { $gt: player.created_at }
} });
- participant_phase_records = participant_phase_records.concat(
- participants_phase); // TODO calc stats
+ // update requested from bridge sets both created_at and last_update
+ if (player.last_update == undefined) {
+ // this will not overwrite last_update and created_at
+ if (count == 0) player_records.add(player);
+ // else db is more recent than buffer, skip
+ } else {
+ // this will overwrite dates
+ player_records_dates.add(player);
+ }
});
let transaction_profiler = logger.startTimer();
// now access db
- try {
- // upsert whole batch in parallel
- logger.info("inserting batch into db");
- await seq.transaction({ autocommit: false }, async (transaction) => {
- await Promise.map(chunks(match_records), async (m_r) =>
- model.Match.bulkCreate(m_r, {
- ignoreDuplicates: true, // if this happens, something is wrong
- transaction: transaction
- }), { concurrency: MAXCONNS }
- );
- await Promise.map(chunks(roster_records), async (r_r) =>
- model.Roster.bulkCreate(r_r, {
- ignoreDuplicates: true,
- transaction: transaction
- }), { concurrency: MAXCONNS }
- );
- await Promise.map(chunks(participant_records), async (p_r) =>
- model.Participant.bulkCreate(p_r, {
- ignoreDuplicates: true,
- transaction: transaction
- }), { concurrency: MAXCONNS }
- );
- await Promise.map(chunks(participant_stats_records), async (p_s_r) =>
- model.ParticipantStats.bulkCreate(p_s_r, {
- ignoreDuplicates: true,
- transaction: transaction
- }), { concurrency: MAXCONNS }
- );
- await Promise.map(chunks(participant_phase_records), async (p_p_r) =>
- model.ParticipantPhases.bulkCreate(p_p_r, {
- /* ignoreDuplicates: true, TODO DEBUG */
- ignoreDuplicates: false,
- updateOnDuplicate: [],
- transaction: transaction
- }), { concurrency: MAXCONNS }
- );
- await Promise.map(chunks(player_records), async (pl_r) =>
- model.Player.bulkCreate(pl_r, {
- ignoreDuplicates: true,
- transaction: transaction
- }), { concurrency: MAXCONNS }
- );
- await Promise.map(chunks(player_records_direct), async (p_r_d) =>
- model.Player.bulkCreate(player_records_direct, {
- // if set to [] (all), upsert messes with autoincrement
- updateOnDuplicate: [
- "shard_id", "api_id", "name", "last_update",
- "created_at", "level", "xp", "lifetime_gold",
- "skill_tier"
- ],
- transaction: transaction
- }), { concurrency: MAXCONNS }
- );
- await Promise.map(chunks(asset_records), async (a_r) =>
- model.Asset.bulkCreate(a_r, {
- ignoreDuplicates: true,
- transaction: transaction
- }), { concurrency: MAXCONNS }
- );
- });
-
- logger.info("acking batch", { size: msgs.size });
- await Promise.map(msgs, async (m) => await ch.ack(m));
- } catch (err) {
- // this should only happen for Deadlocks in prod
- // it *must not* fail due to broken schema or missing dependency
- // TODO: eliminate such cases earlier in the chain
- // and immediately NACK those broken matches, requeueing only the rest
- logger.error("SQL error", err);
- await Promise.map(msgs, async (m) => await ch.nack(m, false, true));
- return; // give up
- }
- transaction_profiler.done("database transaction");
-
- // notify web
- await Promise.map(msgs, async (m) => {
- if (m.properties.headers.notify == undefined) return;
- let notif = "error";
- // new match
- if (m.properties.type == "match") {
- notif = "matches_update";
- // TODO this sends match_update for every match in the batch to every player
- // notify player.name.api_id about match_update
- await Promise.map(match_records, async (mat) =>
- await ch.publish("amq.topic",
- m.properties.headers.notify + "." + mat.api_id,
- new Buffer("match_update"))
- );
- await ch.publish("amq.topic", m.properties.headers.notify,
- new Buffer("match_update"));
- }
- // player obj updated
- if (m.properties.type == "player") notif = "stats_update";
- // new phases
- if (m.properties.type == "telemetry") {
- // notify match.api_id about phase_update
- await ch.publish("amq.topic",
- m.properties.headers.notify,
- new Buffer("phase_update"))
- }
-
- if (m.properties.headers.donotify == true) // TODO remove later
- await ch.publish("amq.topic", m.properties.headers.notify,
- new Buffer(notif));
+ // upsert whole batch in parallel
+ logger.info("inserting batch into db");
+ await seq.transaction({ autocommit: false }, async (transaction) => {
+ await Promise.map(chunks(match_records), async (m_r) =>
+ model.Match.bulkCreate(m_r, {
+ ignoreDuplicates: true, // if this happens, something is wrong
+ transaction: transaction
+ }), { concurrency: MAXCONNS }
+ );
+ await Promise.map(chunks(roster_records), async (r_r) =>
+ model.Roster.bulkCreate(r_r, {
+ ignoreDuplicates: true,
+ transaction: transaction
+ }), { concurrency: MAXCONNS }
+ );
+ await Promise.map(chunks(participant_records), async (p_r) =>
+ model.Participant.bulkCreate(p_r, {
+ ignoreDuplicates: true,
+ transaction: transaction
+ }), { concurrency: MAXCONNS }
+ );
+ await Promise.map(chunks(participant_stats_records), async (p_s_r) =>
+ model.ParticipantStats.bulkCreate(p_s_r, {
+ ignoreDuplicates: true,
+ transaction: transaction
+ }), { concurrency: MAXCONNS }
+ );
+ await Promise.map(chunks(player_records), async (p_r) =>
+ model.Player.bulkCreate(p_r, {
+ fields: [
+ // specify fields or Sequelize attempts to update all fields
+ "api_id", "name", "shard_id",
+ "skill_tier",
+ "level", "lifetime_gold", "xp"
+ ],
+ updateOnDuplicate: [
+ // only stats
+ "api_id", "name", "shard_id",
+ "skill_tier",
+ "level", "lifetime_gold", "xp"
+ ],
+ transaction: transaction
+ }), { concurrency: MAXCONNS }
+ );
+ await Promise.map(chunks(player_records_dates), async (p_r_d) =>
+ model.Player.bulkCreate(p_r_d, {
+ fields: [
+ "last_update", "created_at",
+ "api_id", "name", "shard_id",
+ "skill_tier",
+ "level", "lifetime_gold", "xp"
+ ],
+ updateOnDuplicate: [
+ "last_update", "created_at",
+ "api_id", "name", "shard_id",
+ "skill_tier",
+ "level", "lifetime_gold", "xp"
+ ],
+ transaction: transaction
+ }), { concurrency: MAXCONNS }
+ );
+ await Promise.map(chunks(asset_records), async (a_r) =>
+ model.Asset.bulkCreate(a_r, {
+ ignoreDuplicates: true,
+ transaction: transaction
+ }), { concurrency: MAXCONNS }
+ );
});
- // …global about new matches
- if (match_records.length > 0)
- await ch.publish("amq.topic", "global", new Buffer("matches_update"));
- // notify follow up services
- if (DOANALYZEMATCH)
- await Promise.each(match_records, async (m) =>
- await ch.sendToQueue(ANALYZE_QUEUE, new Buffer(m.api_id),
- { persistent: true }));
+ transaction_profiler.done("database transaction");
}
// Split participant API data into participant and participant_stats
@@ -918,8 +614,9 @@ function flatten(obj) {
return "carry";
return "jungler";
}
-})();
+});
-process.on("unhandledRejection", function(reason, promise) {
- logger.error(reason);
+process.on("unhandledRejection", (err) => {
+ logger.error(err);
+ global.process.exit(1); // fail hard and die
});