WIPP
This commit is contained in:
parent
4184e99273
commit
493e669e6d
5 changed files with 198 additions and 196 deletions
197
lib/plugins.js
197
lib/plugins.js
|
|
@ -1,5 +1,3 @@
|
|||
'use strict'
|
||||
|
||||
const uuid = require('uuid')
|
||||
const R = require('ramda')
|
||||
const BigNumber = require('bignumber.js')
|
||||
|
|
@ -12,16 +10,14 @@ const db = require('./postgresql_interface')
|
|||
const logger = require('./logger')
|
||||
const notifier = require('./notifier')
|
||||
const T = require('./time')
|
||||
const settingsLoader = require('./settings')
|
||||
const configManager = require('./config-manager')
|
||||
const settingsLoader = require('./settings-loader')
|
||||
const ticker = require('./ticker')
|
||||
const wallet = require('./wallet')
|
||||
const exchange = require('./exchange')
|
||||
const sms = require('./sms')
|
||||
const email = require('./email')
|
||||
|
||||
const tradeIntervals = {}
|
||||
|
||||
const CHECK_NOTIFICATION_INTERVAL = T.minute
|
||||
const ALERT_SEND_INTERVAL = T.hour
|
||||
const INCOMING_TX_INTERVAL = 30 * T.seconds
|
||||
|
|
@ -36,6 +32,8 @@ const SWEEP_LIVE_HD_INTERVAL = T.minute
|
|||
const SWEEP_OLD_HD_INTERVAL = 2 * T.minutes
|
||||
const TRADE_INTERVAL = T.minute
|
||||
const TRADE_TTL = 5 * T.minutes
|
||||
const STALE_TICKER = 3 * 60 * 1000
|
||||
const STALE_BALANCE = 3 * 60 * 1000
|
||||
|
||||
const tradesQueues = {}
|
||||
|
||||
|
|
@ -47,13 +45,46 @@ const coins = {
|
|||
let alertFingerprint = null
|
||||
let lastAlertTime = null
|
||||
|
||||
function getConfig (machineId) {
|
||||
const config = settingsLoader.settings().config
|
||||
return configManager.machineScoped(machineId, config)
|
||||
exports.logEvent = db.recordDeviceEvent
|
||||
|
||||
function buildRates (settings, deviceId, tickers) {
|
||||
const config = configManager.machineScoped(deviceId, settings.config)
|
||||
const cryptoCodes = config.currencies.cryptoCurrencies
|
||||
|
||||
const cashInCommission = new BigNumber(config.commissions.cashInCommission).div(100).plus(1)
|
||||
const cashOutCommission = new BigNumber(config.commissions.cashOutCommission).div(100).plus(1)
|
||||
|
||||
const rates = {}
|
||||
|
||||
cryptoCodes.forEach((cryptoCode, i) => {
|
||||
const rateRec = tickers[i]
|
||||
if (Date.now() - rateRec.timestamp > STALE_TICKER) return logger.warn('Stale rate for ' + cryptoCode)
|
||||
const rate = rateRec.rates
|
||||
rates[cryptoCode] = {
|
||||
cashIn: rate.ask.times(cashInCommission),
|
||||
cashOut: rate.bid.div(cashOutCommission)
|
||||
}
|
||||
})
|
||||
|
||||
return rates
|
||||
}
|
||||
|
||||
exports.getConfig = getConfig
|
||||
exports.logEvent = db.recordDeviceEvent
|
||||
function buildBalances (settings, deviceId, balanceRecs) {
|
||||
const config = configManager.machineScoped(deviceId, settings.config)
|
||||
const cryptoCodes = config.currencies.cryptoCurrencies
|
||||
|
||||
const balances = {}
|
||||
|
||||
cryptoCodes.forEach((cryptoCode, i) => {
|
||||
const balanceRec = balanceRecs[i]
|
||||
if (!balanceRec) return logger.warn('No balance for ' + cryptoCode + ' yet')
|
||||
if (Date.now() - balanceRec.timestamp > STALE_BALANCE) return logger.warn('Stale balance for ' + cryptoCode)
|
||||
|
||||
balances[cryptoCode] = balanceRec.balance
|
||||
})
|
||||
|
||||
return balances
|
||||
}
|
||||
|
||||
function buildCartridges (cartridges, virtualCartridges, rec) {
|
||||
return {
|
||||
|
|
@ -71,16 +102,32 @@ function buildCartridges (cartridges, virtualCartridges, rec) {
|
|||
}
|
||||
}
|
||||
|
||||
exports.pollQueries = function pollQueries (deviceId) {
|
||||
const config = getConfig(deviceId)
|
||||
exports.pollQueries = function pollQueries (settings, deviceTime, deviceId, deviceRec) {
|
||||
const config = configManager.machineScoped(deviceId, settings.config)
|
||||
const fiatCode = config.currencies.fiatCurrency
|
||||
const cryptoCodes = config.currencies.cryptoCurrencies
|
||||
const cartridges = [ config.currencies.topCashOutDenomination,
|
||||
config.currencies.bottomCashOutDenomination ]
|
||||
const virtualCartridges = [config.currencies.virtualCashOutDenomination]
|
||||
|
||||
return db.cartridgeCounts(deviceId)
|
||||
.then(result => ({
|
||||
cartridges: buildCartridges(cartridges, virtualCartridges, result)
|
||||
}))
|
||||
const tickerPromises = cryptoCodes.map(c => ticker.getRates(settings, fiatCode, c))
|
||||
const balancePromises = cryptoCodes.map(c => wallet.balance(settings, c))
|
||||
const pingPromise = recordPing(deviceId, deviceTime, deviceRec)
|
||||
|
||||
const promises = [db.cartridgeCounts(deviceId), pingPromise].concat(tickerPromises, balancePromises)
|
||||
|
||||
return Promise.all(promises)
|
||||
.then(arr => {
|
||||
const cartridgeCounts = arr[0]
|
||||
const tickers = arr.slice(2, cryptoCodes.length + 2)
|
||||
const balances = arr.slice(cryptoCodes.length + 2)
|
||||
|
||||
return {
|
||||
cartridges: buildCartridges(cartridges, virtualCartridges, cartridgeCounts),
|
||||
rates: buildRates(settings, deviceId, tickers),
|
||||
balances: buildBalances(settings, deviceId, balances)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// NOTE: This will fail if we have already sent coins because there will be
|
||||
|
|
@ -137,7 +184,7 @@ exports.stateChange = function stateChange (deviceId, deviceTime, rec, cb) {
|
|||
return db.machineEvent(event)
|
||||
}
|
||||
|
||||
exports.recordPing = function recordPing (deviceId, deviceTime, rec, cb) {
|
||||
function recordPing (deviceId, deviceTime, rec) {
|
||||
const event = {
|
||||
id: uuid.v4(),
|
||||
deviceId: deviceId,
|
||||
|
|
@ -177,17 +224,16 @@ exports.cashOut = function cashOut (deviceId, tx) {
|
|||
})
|
||||
}
|
||||
|
||||
exports.dispenseAck = function (deviceId, tx) {
|
||||
const config = getConfig(deviceId)
|
||||
exports.dispenseAck = function (settings, deviceId, tx) {
|
||||
const config = configManager.machineScoped(deviceId, settings.config)
|
||||
const cartridges = [ config.currencies.topCashOutDenomination,
|
||||
config.currencies.bottomCashOutDenomination ]
|
||||
|
||||
return db.addDispense(deviceId, tx, cartridges)
|
||||
}
|
||||
|
||||
exports.fiatBalance = function fiatBalance (fiatCode, cryptoCode, deviceId) {
|
||||
const _config = settingsLoader.settings().config
|
||||
const config = configManager.scoped(cryptoCode, deviceId, _config)
|
||||
function fiatBalance (settings, fiatCode, cryptoCode, deviceId) {
|
||||
const config = configManager.scoped(cryptoCode, deviceId, settings.config)
|
||||
|
||||
return Promise.all([ticker.ticker(cryptoCode), wallet.balance(cryptoCode)])
|
||||
.then(([rates, balanceRec]) => {
|
||||
|
|
@ -257,11 +303,7 @@ function monitorUnnotified () {
|
|||
exports.startPolling = function startPolling () {
|
||||
executeTrades()
|
||||
|
||||
const cryptoCodes = getAllCryptoCodes()
|
||||
cryptoCodes.forEach(cryptoCode => {
|
||||
startTrader(cryptoCode)
|
||||
})
|
||||
|
||||
setInterval(executeTrades, TRADE_INTERVAL)
|
||||
setInterval(monitorLiveIncoming, LIVE_INCOMING_TX_INTERVAL)
|
||||
setInterval(monitorIncoming, INCOMING_TX_INTERVAL)
|
||||
setInterval(monitorUnnotified, UNNOTIFIED_INTERVAL)
|
||||
|
|
@ -275,14 +317,6 @@ exports.startPolling = function startPolling () {
|
|||
sweepOldHD()
|
||||
}
|
||||
|
||||
function startTrader (cryptoCode) {
|
||||
if (tradeIntervals[cryptoCode]) return
|
||||
|
||||
logger.debug('[%s] startTrader', cryptoCode)
|
||||
|
||||
tradeIntervals[cryptoCode] = setInterval(() => executeTrades(cryptoCode), TRADE_INTERVAL)
|
||||
}
|
||||
|
||||
/*
|
||||
* Trader functions
|
||||
*/
|
||||
|
|
@ -325,7 +359,30 @@ function consolidateTrades (cryptoCode, fiatCode) {
|
|||
return consolidatedTrade
|
||||
}
|
||||
|
||||
function executeTrades (cryptoCode, fiatCode) {
|
||||
function executeTrades () {
|
||||
return settingsLoader()
|
||||
.then(settings => {
|
||||
const config = settings.config
|
||||
return db.devices()
|
||||
.then(devices => {
|
||||
const deviceIds = devices.map(device => device.device_id)
|
||||
const lists = deviceIds.map(deviceId => {
|
||||
const currencies = configManager.machineScoped(deviceId, config).currencies
|
||||
const fiatCode = currencies.fiatCurrency
|
||||
const cryptoCodes = currencies.cryptoCurrencies
|
||||
return cryptoCodes.map(cryptoCode => ({fiatCode, cryptoCode}))
|
||||
})
|
||||
|
||||
const tradesPromises = R.uniq(R.flatten(lists))
|
||||
.map(r => executeTradesForMarket(settings, r.fiatCode, r.cryptoCode))
|
||||
|
||||
return Promise.all(tradesPromises)
|
||||
})
|
||||
})
|
||||
.catch(logger.error)
|
||||
}
|
||||
|
||||
function executeTradesForMarket (settings, fiatCode, cryptoCode) {
|
||||
const market = [fiatCode, cryptoCode].join('')
|
||||
logger.debug('[%s] checking for trades', market)
|
||||
|
||||
|
|
@ -405,59 +462,35 @@ function checkNotification () {
|
|||
})
|
||||
}
|
||||
|
||||
function getCryptoCodes (deviceId) {
|
||||
return settingsLoader.settings()
|
||||
.then(settings => {
|
||||
return configManager.machineScoped(deviceId, settings.config).currencies.cryptoCurrencies
|
||||
})
|
||||
}
|
||||
exports.getCryptoCodes = getCryptoCodes
|
||||
function checkDeviceBalances (settings, deviceId) {
|
||||
const config = configManager.machineScoped(deviceId, settings.config)
|
||||
const cryptoCodes = config.currencies.cryptoCurrencies
|
||||
const fiatCode = config.currencies.fiatCurrency
|
||||
const fiatBalancePromises = cryptoCodes.map(c => fiatBalance(settings, fiatCode, c, deviceId))
|
||||
|
||||
// Get union of all cryptoCodes from all machines
|
||||
function getAllCryptoCodes () {
|
||||
return Promise.all([db.devices(), settingsLoader.settings()])
|
||||
.then(([rows, settings]) => {
|
||||
return rows.reduce((acc, r) => {
|
||||
const cryptoCodes = configManager.machineScoped(r.device_id, settings.config).currencies.cryptoCurrencies
|
||||
cryptoCodes.forEach(c => acc.add(c))
|
||||
return acc
|
||||
}, new Set())
|
||||
return Promise.all(fiatBalancePromises)
|
||||
.then(arr => {
|
||||
return arr.map((balance, i) => ({
|
||||
fiatBalance: balance,
|
||||
cryptoCode: cryptoCodes[i],
|
||||
fiatCode,
|
||||
deviceId
|
||||
}))
|
||||
})
|
||||
}
|
||||
|
||||
function getAllMarkets () {
|
||||
return Promise.all([db.devices(), settingsLoader.settings()])
|
||||
.then(([rows, settings]) => {
|
||||
return rows.reduce((acc, r) => {
|
||||
const currencies = configManager.machineScoped(r.device_id, settings.config).currencies
|
||||
const cryptoCodes = currencies.cryptoCurrencies
|
||||
const fiatCodes = currencies.fiatCodes
|
||||
fiatCodes.forEach(fiatCode => cryptoCodes.forEach(cryptoCode => acc.add(fiatCode + cryptoCode)))
|
||||
return acc
|
||||
}, new Set())
|
||||
})
|
||||
}
|
||||
|
||||
function checkBalances () {
|
||||
return Promise.all([getAllMarkets(), db.devices()])
|
||||
.then(([markets, devices]) => {
|
||||
function checkBalances (settings) {
|
||||
return db.devices()
|
||||
.then(devices => {
|
||||
const deviceIds = devices.map(r => r.device_id)
|
||||
const balances = []
|
||||
const deviceBalancePromises = deviceIds.map(deviceId => checkDeviceBalances(settings, deviceId))
|
||||
|
||||
markets.forEach(market => {
|
||||
const fiatCode = market.fiatCode
|
||||
const cryptoCode = market.cryptoCode
|
||||
const minBalance = deviceIds.map(deviceId => {
|
||||
const fiatBalanceRec = exports.fiatBalance(fiatCode, cryptoCode, deviceId)
|
||||
return fiatBalanceRec ? fiatBalanceRec.balance : Infinity
|
||||
})
|
||||
.reduce((min, cur) => Math.min(min, cur), Infinity)
|
||||
|
||||
const rec = {fiatBalance: minBalance, cryptoCode, fiatCode}
|
||||
balances.push(rec)
|
||||
return Promise.all(deviceBalancePromises)
|
||||
.then(arr => {
|
||||
const toMarket = r => r.fiatBalance + r.cryptoCode
|
||||
const min = R.minBy(r => r.fiatBalance)
|
||||
return R.reduceBy(min, Infinity, toMarket, R.flatten(arr))
|
||||
})
|
||||
|
||||
return balances
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue