summaryrefslogtreecommitdiff
path: root/backend/index.js
blob: 8480f28f0059b69e57c4359dbf652302ac5c0ff8 (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
#!/usr/bin/node

const R  = require('ramda'),
  Future = require('fluture'),
  moment = require('moment'),
  fsPath = require('fs-path'),
  reduce = require('./reduce'),
  source = require('./source'),
  file   = require('./file'),
  config = require('./config');

const loadFTimestamped = source.loadFTimestamped(config);
const deriveStatistics = reduce.deriveStatistics(config.reduce);
const aggregatePayloads = reduce.aggregatePayloads(config.reduce);
const cleanPayloads = reduce.cleanPayloads(config.reduce);
const saveFPayloads = file.saveFPayloads(config.file);

function main() {
  const start = moment('2018-05-16 16Z');
  const end = moment().utc();
  const duration = moment.duration(end.diff(start));

  const later = R.curry((base, hs) => base.clone().add(hs, 'hours'));
  const hours = R.range(0, Math.floor(duration.asHours()));
  const laterMoments = R.map(later(start), hours);

  const metadata = {
    lastUpdate: moment().toISOString(),
    config,
  };
  const saveMetadata = (path, data) =>
    Future.node((cb) => fsPath.writeFile(path, JSON.stringify(data), cb))
          .map(() => data);

  const futures = R.map(loadFTimestamped, laterMoments);
  Future.parallel(1, futures)
    .map(R.compose(deriveStatistics, aggregatePayloads, cleanPayloads, R.unnest))
    .chain(saveFPayloads(config.file.reportPattern()))
    .chain((data) => saveMetadata(config.file.metadataPattern(), metadata))
    .fork(console.error, (d) => {});
}

main();