summaryrefslogtreecommitdiff
path: root/api.js
blob: 37ed6aec68d70c778f5acdf030c49b42dd122534 (plain)
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
#!/usr/bin/node
/* jshint esnext:true */
"use strict";

const request = require("request-promise-native"),
    WebSocket = require("ws"),
    webstomp = require("webstomp-client"),
    Channel = require("async-csp").Channel;

const API_FE_URL = process.env.API_FE_URL || "http://vainsocial.dev/bots/api",
      API_WS_URL = process.env.API_WS_URL || "ws://vainsocial.dev/ws",
      API_BE_URL = process.env.API_BE_URL || "http://vainsocial.dev/bridge";

const notif = webstomp.over(new WebSocket(API_WS_URL,
    { perMessageDeflate: false }));

(function connect() {
    notif.connect("web", "web",
        () => console.log("connected to queue"),
        (err) => connect()
    );
})();

// TODO use keepalive / connection pool
function getFE(url) {
    return request({
        uri: API_FE_URL + url,
        json: true
    });
}

function postBE(url) {
    return request.post({
        uri: API_BE_URL + url,
        json: true
    });
}

function subscribe(topic, channel) {
    return notif.subscribe("/topic/" + topic, (msg) => {
        channel.put(msg.body);
        msg.ack();
    }, {"ack": "client"});
}

// be an async generator
// next() returns player data whenever an update is available
module.exports.searchPlayer = async function (name, timeout=60) {
    let channel = new Channel(),
        subscription = subscribe("player." + name, channel);
    await postBE("/player/" + name + "/update");
    // stop updates after timeout
    setTimeout(() => {
        channel.close();
        subscription.unsubscribe();
    }, timeout*1000);
    channel.put("initial");  // initial fetch

    let generator = async () => {
        let msg;
        while (true) {
            msg = await channel.take();
            if (["initial", "search_fail", "stats_update",
                "matches_update"].indexOf(msg) != -1)
                break;
        }
        if (msg == "search_fail") {
            throw "not found";
        }
        if (msg == Channel.DONE) {
            throw "exhausted";
        }
        return getFE("/player/" + name);
    }

    return { next: generator };
}

// return matches
module.exports.searchMatches = async function (name) {
    return getFE("/player/" + name + "/matches/1.1.1.1");
}

// return single match
module.exports.searchMatch = async function (id) {
    return getFE("/match/" + id);
}