Compare commits

..
Author SHA1 Message Date
Victor Grau Serrat 1bd65397b6 Merge branch 'develop' 2018-03-22 12:02:15 -06:00
Victor Grau Serrat c768b207bc MAINT: [marketplace] output formatting address list 2018-03-22 12:01:21 -06:00
VictorandGitHub e18686d5c5 Merge pull request #278 from izokay/develop
BUG: Error when ingesting marketcap on windows
2018-03-22 12:23:29 -05:00
VictorandGitHub b7779cf363 MAINT: general bug fix for existing path across OS 2018-03-22 11:23:06 -06:00
Victor Grau Serrat 28819b8a32 DOC: updated release notes for 0.5.6 2018-03-21 22:29:21 -06:00
Victor Grau Serrat a56d7f34c7 BLD: [mktplace] support for most wallets, switch to mycrypto 2018-03-21 20:35:18 -05:00
AvishaiW 9f0b3303f1 BUG: #285 #271 changed benchmark to be constant, so it wouldn't ingest data at all, for now 2018-03-21 21:19:40 +02:00
EmbarAlmog 9d7a35658b ENH: when ingesting data of non-existing pair it is now throwing log warning. 2018-03-20 16:12:13 +02:00
Victor Grau Serrat c58cebd1eb ENH: progress on marketplace bundle ingestion 2018-03-19 11:50:39 -06:00
Frederic Fortier 027cdba474 Merge branch 'develop' 2018-03-19 13:12:08 -04:00
Frederic Fortier d223529100 DOC: updated release notes of 0.5.5 2018-03-19 13:11:31 -04:00
Frederic Fortier 1e02506ab4 Merge branch 'develop' 2018-03-19 13:07:06 -04:00
lenak25 2a97ade68e BLD: support hourly freq in live and backtest, as reported on issue #227 and issue #114 2018-03-19 16:44:44 +02:00
lenak25 9648767e9a STY: flake8 fixes 2018-03-19 15:53:06 +02:00
lenak25 98449b2088 BLD: fix issue #274 - a bug in which a wrong bar number was returned when requesting day freq history candles in backtest 2018-03-19 11:06:47 +02:00
Frederic Fortier 91d16aba3b DOC: documented the get_frequency function for additional clarity 2018-03-17 18:32:28 -04:00
izokayandGitHub 9eb649371b BUG: Error when ingesting on windows
Error message: Cannot create a file when that file already exists: '.catalyst\\data\\marketplace\\temp_bundles\\marketcap-hourly-2018' -> '.catalyst\\data\\marketplace\\marketcap'
2018-03-16 16:50:31 -04:00
izokayandGitHub 7f2ded65bc Merge pull request #2 from enigmampc/develop
Develop
2018-03-16 16:45:03 -04:00
Frederic Fortier b76b4458cb Merge branch 'vonpupp-fix_hourly_candles' into develop 2018-03-16 15:49:39 -04:00
Frederic Fortier decbdbf6ea Merge branch 'fix_hourly_candles' of https://github.com/vonpupp/catalyst into vonpupp-fix_hourly_candles 2018-03-16 15:49:28 -04:00
Albert De La Fuente Vigliotti 685ce25b85 Fix H candle support 2018-03-16 16:02:40 -03:00
Avishai WeingartenandGitHub 1cafcc1417 BUG: removed one out of two matplotlib appearences in 2.7 yml 2018-03-16 14:40:05 +02:00
VictorandGitHub 0d77854782 Merge pull request #275 from izokay/patch-1
typo on creating env for python 3.6
2018-03-15 15:18:07 -06:00
izokayandGitHub 4cb8d54d97 typo on creating env for python 3.6 2018-03-15 15:41:30 -04:00
Victor Grau Serrat 7b796a4276 MAINT: [mktplace] sign_msg opens browser window 2018-03-15 12:59:59 -04:00
AvishaiW 41a4c7072f DOC: fixed a mistake on the installation tutorial 2018-03-14 09:44:24 +02:00
Victor Grau Serrat 11302b3af9 Merge branch 'develop' 2018-03-14 00:52:33 -06:00
Victor Grau Serrat 5bb7eed072 MAINT: ref. mktplace to master, updated release notes 0.5.4 2018-03-14 00:51:49 -06:00
Victor Grau Serrat dbf3b6e6b2 MAINT: typo in marketplace help 2018-03-14 00:11:24 -06:00
Victor Grau Serrat 3e69449a6b BLD: marketplace switch to rinkeby post-audit 2018-03-14 00:11:24 -06:00
lenak25 69731b653d BLD: revert hourly freq support reported at issue #227 2018-03-13 18:43:48 +02:00
Victor Grau Serrat 127d779eb1 BUG: fix sanitize_df to min of int32 2018-03-12 16:52:54 -06:00
lenak25 8d86a5548f DOC: add ta_lib troubleshooting to the docs 2018-03-12 18:04:38 +02:00
lenak25 0a37cdec5b BLD: fix 'on the clock' candles fetch and request extra candles using a fixed time interval 2018-03-11 19:36:28 +02:00
38 changed files with 970 additions and 1036 deletions
+2 -2
View File
@@ -580,7 +580,7 @@ def ingest_exchange(ctx, exchange_name, data_frequency, start, end,
exchange_bundle = ExchangeBundle(exchange_name)
click.echo('Ingesting exchange bundle {}...'.format(exchange_name),
click.echo('Trying to ingest exchange bundle {}...'.format(exchange_name),
sys.stdout)
exchange_bundle.ingest(
data_frequency=data_frequency,
@@ -793,7 +793,7 @@ def ls(ctx):
)
@click.pass_context
def subscribe(ctx, dataset):
"""Subscribe to an exisiting dataset.
"""Subscribe to an existing dataset.
"""
marketplace = Marketplace()
marketplace.subscribe(dataset)
+10 -59
View File
@@ -433,7 +433,7 @@ cdef class TradingPair(Asset):
'taker',
'trading_state',
'data_source',
'decimals',
'decimals'
})
def __init__(self,
object symbol,
@@ -455,7 +455,7 @@ cdef class TradingPair(Asset):
float taker=0.0025,
float lot=0,
int decimals = 8,
int trading_state=1,
int trading_state=0,
object data_source='catalyst'):
"""
Replicates the Asset constructor with some built-in conventions
@@ -600,51 +600,14 @@ cdef class TradingPair(Asset):
cpdef to_dict(self):
"""
Convert to a python dict.
Repeat constructor params:
object symbol,
object exchange,
object start_date=None,
object asset_name=None,
int sid=0,
float leverage=1.0,
object end_daily=None,
object end_minute=None,
object end_date=None,
object exchange_symbol=None,
object first_traded=None,
object auto_close_date=None,
object exchange_full=None,
float min_trade_size=0.0001,
float max_trade_size=1000000,
float maker=0.0015,
float taker=0.0025,
float lot=0,
int decimals = 8,
int trading_state=1,
object data_source='catalyst',
"""
trading_pair_dict = dict(
symbol=self.symbol,
exchange=self.exchange,
start_date=self.start_date,
asset_name=self.asset_name,
leverage=self.leverage,
end_daily=self.end_daily,
end_minute=self.end_minute,
end_date=self.end_date,
exchange_symbol=self.exchange_symbol,
exchange_full=self.exchange_full,
min_trade_size=self.min_trade_size,
max_trade_size=self.max_trade_size,
maker=self.maker,
taker=self.taker,
lot=self.lot,
decimals=self.decimals,
trading_state=self.trading_state,
data_source=self.data_source,
)
return trading_pair_dict
#TODO: missing fields
super_dict = super(TradingPair, self).to_dict()
super_dict['end_daily'] = self.end_daily
super_dict['end_minute'] = self.end_minute
super_dict['leverage'] = self.leverage
super_dict['min_trade_size'] = self.min_trade_size
return super_dict
def is_exchange_open(self, dt_minute):
"""
@@ -660,16 +623,6 @@ cdef class TradingPair(Asset):
#TODO: make more dymanic to catch holds
return True
def set_end_date(self, dt, data_frequency):
if data_frequency == 'minute':
self.end_minute = dt
else:
self.end_daily = dt
def set_start_date(self, dt):
self.start_date = dt
cpdef __reduce__(self):
"""
Function used by pickle to determine how to serialize/deserialize this
@@ -693,9 +646,7 @@ cdef class TradingPair(Asset):
self.lot,
self.decimals,
self.taker,
self.maker,
self.trading_state,
self.data_source))
self.maker))
def make_asset_array(int size, Asset asset):
cdef np.ndarray out = np.empty([size], dtype=object)
+7 -8
View File
@@ -11,10 +11,7 @@ LOG_LEVEL = int(os.environ.get('CATALYST_LOG_LEVEL', logbook.INFO))
SYMBOLS_URL = 'https://s3.amazonaws.com/enigmaco/catalyst-exchanges/' \
'{exchange}/symbols.json'
EXCHANGE_CONFIG_URL = 'https://s3.amazonaws.com/enigmaco/ohlcv/' \
'{exchange}/config.json'
BUNDLE_URL = 'https://s3.amazonaws.com/enigmaco/ohlcv/' \
'{exchange}/{data_frequency}/{name}.tar.gz'
DATE_TIME_FORMAT = '%Y-%m-%d %H:%M'
DATE_FORMAT = '%Y-%m-%d'
@@ -28,8 +25,7 @@ AUTO_INGEST = False
AUTH_SERVER = 'https://data.enigma.co'
# TODO: switch to mainnet
ETH_REMOTE_NODE = 'https://ropsten.infura.io/'
ETH_REMOTE_NODE = 'https://rinkeby.infura.io/'
MARKETPLACE_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/master/catalyst/marketplace/' \
@@ -40,10 +36,13 @@ MARKETPLACE_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'contract_marketplace_abi.json'
# TODO: switch to mainnet
ENIGMA_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/catalyst/' \
'master/catalyst/marketplace/' \
ENIGMA_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/master/catalyst/marketplace/' \
'contract_enigma_address.txt'
ENIGMA_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/master/catalyst/marketplace/' \
'contract_enigma_abi.json'
SUPPORTED_WALLETS = ['metamask', 'ledger', 'trezor', 'bitbox', 'keystore',
'key']
+7 -7
View File
@@ -33,12 +33,12 @@ def initialize(context):
# parameters or values you're going to use.
# In our example, we're looking at Neo in Ether.
context.market = symbol('eth_btc')
context.market = symbol('bnb_eth')
context.base_price = None
context.current_day = None
context.RSI_OVERSOLD = 55
context.RSI_OVERBOUGHT = 60
context.RSI_OVERSOLD = 60
context.RSI_OVERBOUGHT = 70
context.CANDLE_SIZE = '15T'
context.start_time = time.time()
@@ -248,14 +248,14 @@ if __name__ == '__main__':
if live:
run_algorithm(
capital_base=0.03,
capital_base=0.1,
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='poloniex',
exchange_name='binance',
live=True,
algo_namespace=NAMESPACE,
base_currency='btc',
base_currency='eth',
live_graph=False,
simulate_orders=False,
stats_output=None,
@@ -274,7 +274,7 @@ if __name__ == '__main__':
# -x bitfinex -s 2017-10-1 -e 2017-11-10 -c usdt -n mean-reversion \
# --data-frequency minute --capital-base 10000
run_algorithm(
capital_base=0.1,
capital_base=0.035,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
+1 -1
View File
@@ -26,7 +26,7 @@ def handle_data(context, data):
context.asset,
fields='price',
bar_count=20,
frequency='2H'
frequency='30T'
)
last_traded = prices.index[-1]
log.info('last candle date: {}'.format(last_traded))
+276 -185
View File
@@ -1,33 +1,34 @@
import json
import os
import re
from collections import defaultdict
import ccxt
import pandas as pd
import six
from catalyst.assets._assets import TradingPair
from redo import retry
from ccxt import InvalidOrder, NetworkError, \
ExchangeError
from logbook import Logger
from six import string_types
from catalyst.algorithm import MarketOrder
from catalyst.assets._assets import TradingPair
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.exchange import Exchange
from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import InvalidHistoryFrequencyError, \
ExchangeSymbolsNotFound, ExchangeRequestError, InvalidOrderStyle, \
UnsupportedHistoryFrequencyError, \
ExchangeNotFoundError, CreateOrderError, InvalidHistoryTimeframeError, \
MarketsNotFoundError, InvalidMarketError
UnsupportedHistoryFrequencyError
from catalyst.exchange.exchange_execution import ExchangeLimitOrder
from catalyst.exchange.utils.ccxt_utils import get_exchange_config
from catalyst.exchange.utils.exchange_utils import mixin_market_params, \
get_exchange_folder, get_catalyst_symbol, \
get_exchange_auth
from catalyst.exchange.utils.datetime_utils import from_ms_timestamp, \
get_epoch, \
get_periods_range
from catalyst.exchange.utils.exchange_utils import get_catalyst_symbol
from catalyst.finance.order import Order, ORDER_STATUS
from catalyst.finance.transaction import Transaction
from ccxt import InvalidOrder, NetworkError, \
ExchangeError
from logbook import Logger
from six import string_types
log = Logger('CCXT', level=LOG_LEVEL)
@@ -43,7 +44,7 @@ SUPPORTED_EXCHANGES = dict(
class CCXT(Exchange):
def __init__(self, exchange_name, key,
secret, password, base_currency, config=None):
secret, password, base_currency):
log.debug(
'finding {} in CCXT exchanges:\n{}'.format(
exchange_name, ccxt.exchanges
@@ -63,8 +64,6 @@ class CCXT(Exchange):
'password': password,
})
self.api.enableRateLimit = True
self.has = self.api.has
self.fees = self.api.fees
except Exception:
raise ExchangeNotFoundError(exchange_name=exchange_name)
@@ -72,7 +71,6 @@ class CCXT(Exchange):
self._symbol_maps = [None, None]
self.name = exchange_name
self.assets = []
self.base_currency = base_currency
self.transactions = defaultdict(list)
@@ -84,123 +82,97 @@ class CCXT(Exchange):
self._common_symbols = dict()
self.bundle = ExchangeBundle(self.name)
self.markets = None
self._is_init = False
self._config = config
def init(self):
if self._is_init:
return
if self._config is None:
self._config = get_exchange_config(self.name)
log.debug(
'got exchange config {}:\n{}'.format(
self.name, self._config
exchange_folder = get_exchange_folder(self.name)
filename = os.path.join(exchange_folder, 'cctx_markets.json')
if os.path.exists(filename):
timestamp = os.path.getmtime(filename)
dt = pd.to_datetime(timestamp, unit='s', utc=True)
if dt >= pd.Timestamp.utcnow().floor('1D'):
with open(filename) as f:
self.markets = json.load(f)
log.debug('loaded markets for {}'.format(self.name))
if self.markets is None:
try:
markets_symbols = self.api.load_markets()
log.debug(
'fetching {} markets:\n{}'.format(
self.name, markets_symbols
)
)
)
self.markets = self.api.fetch_markets()
with open(filename, 'w+') as f:
json.dump(self.markets, f, indent=4)
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to fetch markets {}: {}'.format(
self.name, e
)
)
raise ExchangeRequestError(error=e)
self.load_assets()
self._is_init = True
def load_assets(self):
if self._config is None:
raise ValueError('Exchange config not available.')
@staticmethod
def find_exchanges(features=None, is_authenticated=False):
ccxt_features = []
if features is not None:
for feature in features:
if not feature.endswith('Bundle'):
ccxt_features.append(feature)
self.assets = []
for asset_dict in self._config['assets']:
asset = TradingPair(**asset_dict)
self.assets.append(asset)
exchange_names = []
for exchange_name in ccxt.exchanges:
if is_authenticated:
exchange_auth = get_exchange_auth(exchange_name)
def _fetch_markets(self):
markets_symbols = self.api.load_markets()
log.debug(
'fetching {} markets:\n{}'.format(
self.name, markets_symbols
)
)
try:
markets = self.api.fetch_markets()
has_auth = (exchange_auth['key'] != ''
and exchange_auth['secret'] != '')
except NetworkError as e:
raise ExchangeRequestError(error=e)
if not has_auth:
continue
if not markets:
raise MarketsNotFoundError(
exchange=self.name,
)
log.debug('loading exchange: {}'.format(exchange_name))
exchange = getattr(ccxt, exchange_name)()
for market in markets:
if 'id' not in market:
raise InvalidMarketError(
exchange=self.name,
market=market,
)
return markets
if ccxt_features is None:
has_feature = True
def create_exchange_config(self):
config = dict(
name=self.name,
features=[feature for feature in self.has if self.has[feature]]
)
markets = retry(
action=self._fetch_markets,
attempts=5,
sleeptime=5,
retry_exceptions=(ExchangeRequestError,),
cleanup=lambda: log.warn(
'fetching markets again for {}'.format(self.name)
),
)
else:
try:
has_feature = all(
[exchange.has[feature] for feature in ccxt_features]
)
config['assets'] = []
for market in markets:
asset = self.create_trading_pair(market=market)
config['assets'].append(asset)
except Exception:
has_feature = False
return config
if has_feature:
try:
log.info('initializing {}'.format(exchange_name))
exchange_names.append(exchange_name)
def create_trading_pair(self, market, start_dt=None, end_dt=None,
leverage=1, end_daily=None, end_minute=None):
"""
Creating a TradingPair from market and asset data.
except Exception as e:
log.warn(
'unable to initialize exchange {}: {}'.format(
exchange_name, e
)
)
Parameters
----------
market: dict[str, Object]
start_dt
end_dt
leverage
end_daily
end_minute
Returns
-------
"""
params = dict(
exchange=self.name,
data_source='catalyst',
exchange_symbol=market['id'],
symbol=get_catalyst_symbol(market),
start_date=start_dt,
end_date=end_dt,
leverage=leverage,
asset_name=market['symbol'],
end_daily=end_daily,
end_minute=end_minute,
)
self.apply_conditional_market_params(params, market)
return TradingPair(**params)
def load_assets(self):
if self._config is None or 'error' in self._config:
raise ValueError('Exchange config not available.')
self.assets = []
for asset_dict in self._config['assets']:
asset = TradingPair(**asset_dict)
self.assets.append(asset)
return exchange_names
def account(self):
return None
@@ -218,6 +190,9 @@ class CCXT(Exchange):
if data_frequency == 'minute' and not freq.endswith('T'):
continue
elif data_frequency == 'hourly' and not freq.endswith('D'):
continue
elif data_frequency == 'daily' and not freq.endswith('D'):
continue
@@ -232,11 +207,32 @@ class CCXT(Exchange):
return frequencies
def get_market(self, symbol):
"""
The CCXT market.
Parameters
----------
symbol:
The CCXT symbol.
Returns
-------
dict[str, Object]
"""
s = self.get_symbol(symbol)
market = next(
(market for market in self.markets if market['symbol'] == s),
None,
)
return market
def substitute_currency_code(self, currency, source='catalyst'):
if source == 'catalyst':
currency = currency.upper()
key = self.api.common_currency_code(currency).lower()
key = self.api.common_currency_code(currency)
self._common_symbols[key] = currency.lower()
return key
@@ -264,7 +260,13 @@ class CCXT(Exchange):
if source == 'ccxt':
if isinstance(asset_or_symbol, string_types):
parts = asset_or_symbol.split('/')
return '{}_{}'.format(parts[0].lower(), parts[1].lower())
base_currency = self.substitute_currency_code(
parts[0], source
)
quote_currency = self.substitute_currency_code(
parts[1], source
)
return '{}_{}'.format(base_currency, quote_currency)
else:
return asset_or_symbol.symbol
@@ -275,7 +277,13 @@ class CCXT(Exchange):
) else asset_or_symbol.symbol
parts = symbol.split('_')
return '{}/{}'.format(parts[0].upper(), parts[1].upper())
base_currency = self.substitute_currency_code(
parts[0], source
)
quote_currency = self.substitute_currency_code(
parts[1], source
)
return '{}/{}'.format(base_currency, quote_currency)
@staticmethod
def map_frequency(value, source='ccxt', raise_error=True):
@@ -401,7 +409,7 @@ class CCXT(Exchange):
)
def get_candles(self, freq, assets, bar_count=1, start_dt=None,
end_dt=None, floor_dates=True):
end_dt=None):
is_single = (isinstance(assets, TradingPair))
if is_single:
assets = [assets]
@@ -448,20 +456,16 @@ class CCXT(Exchange):
candles[asset] = []
for ohlcv in ohlcvs:
dt = pd.to_datetime(ohlcv[0], unit='ms', utc=True)
if floor_dates:
dt = dt.floor('1T')
candles[asset].append(
dict(
last_traded=dt,
open=ohlcv[1],
high=ohlcv[2],
low=ohlcv[3],
close=ohlcv[4],
volume=ohlcv[5],
)
)
candles[asset].append(dict(
last_traded=pd.to_datetime(
ohlcv[0], unit='ms', utc=True
),
open=ohlcv[1],
high=ohlcv[2],
low=ohlcv[3],
close=ohlcv[4],
volume=ohlcv[5]
))
candles[asset] = sorted(
candles[asset], key=lambda c: c['last_traded']
)
@@ -479,53 +483,144 @@ class CCXT(Exchange):
except ExchangeSymbolsNotFound:
return None
def apply_conditional_market_params(self, params, market):
def get_asset_defs(self, market):
"""
Applies a CCXT market dict to parameters of TradingPair init.
The local and Catalyst definitions of the specified market.
Parameters
----------
params: dict[Object]
market: dict[Object]
market: dict[str, Object]
The CCXT market dicts.
Returns
-------
dict[str, Object]
The asset definition.
"""
asset_defs = []
for is_local in (False, True):
asset_def = self.get_asset_def(market, is_local)
asset_defs.append((asset_def, is_local))
return asset_defs
def get_asset_def(self, market, is_local=False):
"""
The asset definition (in symbols.json files) corresponding
to the the specified market.
Parameters
----------
market: dict[str, Object]
The CCXT market dict.
is_local
Whether to search in local or Catalyst asset definitions.
Returns
-------
dict[str, Object]
The asset definition.
"""
exchange_symbol = market['id']
symbol_map = self._fetch_symbol_map(is_local)
if symbol_map is not None:
assets_lower = {k.lower(): v for k, v in symbol_map.items()}
key = exchange_symbol.lower()
asset = assets_lower[key] if key in assets_lower else None
if asset is not None:
return asset
else:
return None
else:
return None
def create_trading_pair(self, market, asset_def=None, is_local=False):
"""
Creating a TradingPair from market and asset data.
Parameters
----------
market: dict[str, Object]
asset_def: dict[str, Object]
is_local: bool
Returns
-------
"""
# TODO: make this more externalized / configurable
# Consider representing in some type of JSON structure
if 'active' in market:
params['trading_state'] = 1 if market['active'] else 0
data_source = 'local' if is_local else 'catalyst'
params = dict(
exchange=self.name,
data_source=data_source,
exchange_symbol=market['id'],
)
mixin_market_params(self.name, params, market)
if asset_def is not None:
params['symbol'] = asset_def['symbol']
params['start_date'] = asset_def['start_date'] \
if 'start_date' in asset_def else None
params['end_date'] = asset_def['end_date'] \
if 'end_date' in asset_def else None
params['leverage'] = asset_def['leverage'] \
if 'leverage' in asset_def else 1.0
params['asset_name'] = asset_def['asset_name'] \
if 'asset_name' in asset_def else None
params['end_daily'] = asset_def['end_daily'] \
if 'end_daily' in asset_def \
and asset_def['end_daily'] != 'N/A' else None
params['end_minute'] = asset_def['end_minute'] \
if 'end_minute' in asset_def \
and asset_def['end_minute'] != 'N/A' else None
else:
params['trading_state'] = 1
params['symbol'] = get_catalyst_symbol(market)
# TODO: add as an optional column
params['leverage'] = 1.0
if 'lot' in market:
params['min_trade_size'] = market['lot']
params['lot'] = market['lot']
return TradingPair(**params)
if self.name == 'bitfinex':
params['maker'] = 0.001
params['taker'] = 0.002
def load_assets(self):
log.debug('loading assets for {}'.format(self.name))
self.assets = []
elif 'maker' in market and 'taker' in market \
and market['maker'] is not None \
and market['taker'] is not None:
params['maker'] = market['maker']
params['taker'] = market['taker']
for market in self.markets:
if 'id' not in market:
log.warn('invalid market: {}'.format(market))
continue
else:
# TODO: default commission, make configurable
params['maker'] = 0.0015
params['taker'] = 0.0025
asset_defs = self.get_asset_defs(market)
info = market['info'] if 'info' in market else None
if info:
if 'minimum_order_size' in info:
params['min_trade_size'] = float(info['minimum_order_size'])
asset = None
for asset_def in asset_defs:
if asset_def[0] is not None or not asset_defs[1]:
try:
asset = self.create_trading_pair(
market=market,
asset_def=asset_def[0],
is_local=asset_def[1]
)
self.assets.append(asset)
if 'lot' not in params:
params['lot'] = params['min_trade_size']
except TypeError as e:
log.warn('unable to add asset: {}'.format(e))
if asset is None:
asset = self.create_trading_pair(market=market)
self.assets.append(asset)
def get_balances(self):
try:
@@ -663,14 +758,18 @@ class CCXT(Exchange):
side = 'buy' if amount > 0 else 'sell'
if hasattr(self.api, 'amount_to_lots'):
adj_amount = self.api.amount_to_lots(
symbol=symbol,
amount=abs(amount),
)
if adj_amount != abs(amount):
log.info(
'adjusted order amount {} to {} based on lot size'.format(
abs(amount), adj_amount,
# TODO: is this right?
if self.api.markets is None:
self.api.load_markets()
# https://github.com/ccxt/ccxt/issues/1483
adj_amount = round(abs(amount), asset.decimals)
market = self.api.markets[symbol]
if 'lots' in market and market['lots'] > amount:
raise CreateOrderError(
exchange=self.name,
e='order amount lower than the smallest lot: {}'.format(
amount
)
)
@@ -898,7 +997,7 @@ class CCXT(Exchange):
symbol = self.get_symbol(asset_or_symbol) \
if asset_or_symbol is not None else None
self.api.cancel_order(id=order_id,
symbol=symbol, params=params)
symbol=symbol, params= params)
except (ExchangeError, NetworkError) as e:
log.warn(
@@ -1016,27 +1115,19 @@ class CCXT(Exchange):
return result
def get_trades(self, asset, my_trades=True, start_dt=None, limit=100):
if not my_trades:
raise NotImplemented(
'get_trades only supports "my trades"'
)
# TODO: is it possible to sort this? Limit is useless otherwise.
ccxt_symbol = self.get_symbol(asset)
if start_dt:
delta = start_dt - get_epoch()
since = int(delta.total_seconds()) * 1000
else:
since = None
try:
if my_trades:
trades = self.api.fetch_my_trades(
symbol=ccxt_symbol,
since=since,
limit=limit,
)
else:
trades = self.api.fetch_trades(
symbol=ccxt_symbol,
since=since,
limit=limit,
)
trades = self.api.fetch_my_trades(
symbol=ccxt_symbol,
since=start_dt,
limit=limit,
)
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to fetch trades {} / {}: {}'.format(
+51 -30
View File
@@ -5,8 +5,6 @@ from time import sleep
import numpy as np
import pandas as pd
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.data.data_portal import BASE_FIELDS
from catalyst.exchange.exchange_bundle import ExchangeBundle
@@ -18,9 +16,11 @@ from catalyst.exchange.exchange_errors import MismatchingBaseCurrencies, \
TickerNotFoundError, NotEnoughCashError
from catalyst.exchange.utils.datetime_utils import get_delta, \
get_periods_range, \
get_periods, get_start_dt, get_frequency
from catalyst.exchange.utils.exchange_utils import \
resample_history_df, has_bundle
get_periods, get_start_dt, get_frequency, \
get_candles_number_from_minutes
from catalyst.exchange.utils.exchange_utils import get_exchange_symbols, \
resample_history_df, has_bundle, get_candles_df
from logbook import Logger
log = Logger('Exchange', level=LOG_LEVEL)
@@ -199,12 +199,8 @@ class Exchange:
)
assets.append(asset)
except SymbolNotFoundOnExchange:
log.debug(
'skipping non-existent market {} {}'.format(
self.name, symbol
)
)
except SymbolNotFoundOnExchange as e:
log.warn(e)
return assets
def get_asset(self, symbol, data_frequency=None, is_exchange_symbol=False,
@@ -257,10 +253,10 @@ class Exchange:
elif data_frequency is not None:
applies = (
(
data_frequency == 'minute' and a.end_minute is not None
) or (
data_frequency == 'daily' and a.end_daily is not None
)
data_frequency == 'minute' and
a.end_minute is not None)
or (
data_frequency == 'daily' and a.end_daily is not None)
)
else:
@@ -293,6 +289,16 @@ class Exchange:
log.debug('found asset: {}'.format(asset))
return asset
def fetch_symbol_map(self, is_local=False):
index = 1 if is_local else 0
if self._symbol_maps[index] is not None:
return self._symbol_maps[index]
else:
symbol_map = get_exchange_symbols(self.name, is_local)
self._symbol_maps[index] = symbol_map
return symbol_map
@abstractmethod
def init(self):
"""
@@ -304,13 +310,24 @@ class Exchange:
"""
@abstractmethod
def create_exchange_config(self):
def load_assets(self, is_local=False):
"""
Fetch the exchange market data and generate a config object
Returns
-------
Populate the 'assets' attribute with a dictionary of Assets.
The key of the resulting dictionary is the exchange specific
currency pair symbol. The universal symbol is contained in the
'symbol' attribute of each asset.
Notes
-----
The sid of each asset is calculated based on a numeric hash of the
universal symbol. This simple approach avoids maintaining a mapping
of sids.
This method can be omerridden if an exchange offers equivalent data
via its api.
"""
pass
def get_spot_value(self, assets, field, dt=None, data_frequency='minute'):
"""
@@ -491,7 +508,12 @@ class Exchange:
# so we request more than needed
# TODO: consider defining a const per asset
# and/or some retry mechanism (in each iteration request more data)
requested_bar_count = bar_count + 30
kExtra_minutes_candles = 150
requested_bar_count = bar_count + \
get_candles_number_from_minutes(unit,
candle_size,
kExtra_minutes_candles)
# The get_history method supports multiple asset
candles = self.get_candles(
freq=freq,
@@ -509,11 +531,14 @@ class Exchange:
asset=asset,
exchange=self.name)
# for avoiding unnecessary forward fill end_dt is taken back one second
forward_fill_till_dt = end_dt - timedelta(seconds=1)
series = get_candles_df(candles=candles,
field=field,
freq=frequency,
bar_count=requested_bar_count,
end_dt=end_dt)
end_dt=forward_fill_till_dt)
# TODO: consider how to approach this edge case
# delta_candle_size = candle_size * 60 if unit == 'H' else candle_size
@@ -582,7 +607,7 @@ class Exchange:
# TODO: this function needs some work,
# we're currently using it just for benchmark data
freq, candle_size, unit, data_frequency = get_frequency(
frequency, data_frequency
frequency, data_frequency, supported_freqs=['T', 'D']
)
adj_bar_count = candle_size * bar_count
try:
@@ -606,7 +631,7 @@ class Exchange:
start_dt = get_start_dt(end_dt, adj_bar_count, data_frequency)
trailing_dt = \
series[asset].index[-1] + get_delta(1, data_frequency) \
if asset in series else start_dt
if asset in series else start_dt
# The get_history method supports multiple asset
# Use the original frequency to let each api optimize
@@ -647,20 +672,16 @@ class Exchange:
return df
def _check_low_balance(self, currency, balances, amount, open_orders=None):
def _check_low_balance(self, currency, balances, amount):
free = balances[currency]['free'] if currency in balances else 0.0
if open_orders:
# TODO: make sure that this works
free += sum([order.amount for order in open_orders])
if free < amount:
return free, True
else:
return free, False
def sync_positions(self, positions, open_orders=None, cash=None,
def sync_positions(self, positions, cash=None,
check_balances=False):
"""
Update the portfolio cash and position balances based on the
@@ -690,7 +711,7 @@ class Exchange:
balances=balances,
amount=cash,
)
if is_lower and not open_orders:
if is_lower:
raise NotEnoughCashError(
currency=self.base_currency,
exchange=self.name,
+2 -4
View File
@@ -18,11 +18,9 @@ from datetime import timedelta
from os import listdir
from os.path import isfile, join, exists
import catalyst.protocol as zp
import logbook
import pandas as pd
from redo import retry
import catalyst.protocol as zp
from catalyst.algorithm import TradingAlgorithm
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.exchange_blotter import ExchangeBlotter
@@ -52,6 +50,7 @@ from catalyst.utils.api_support import api_method
from catalyst.utils.input_validation import error_keywords, ensure_upper_case
from catalyst.utils.math_utils import round_nearest
from catalyst.utils.preprocess import preprocess
from redo import retry
log = logbook.Logger('exchange_algorithm', level=LOG_LEVEL)
@@ -671,7 +670,6 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
required_cash = self.portfolio.cash if not orders else None
cash, positions_value = exchange.sync_positions(
positions=exchange_positions,
open_orders=orders,
check_balances=check_balances,
cash=required_cash,
)
+87 -86
View File
@@ -1,4 +1,3 @@
import copy
import os
import shutil
from datetime import timedelta
@@ -9,12 +8,8 @@ from operator import is_not
import numpy as np
import pandas as pd
import pytz
from catalyst.assets._assets import TradingPair
from logbook import Logger
from pytz import UTC
from six import itervalues
from catalyst import get_calendar
from catalyst.assets._assets import TradingPair
from catalyst.constants import DATE_TIME_FORMAT, AUTO_INGEST
from catalyst.constants import LOG_LEVEL
from catalyst.data.minute_bars import BcolzMinuteOverlappingData, \
@@ -28,11 +23,14 @@ from catalyst.exchange.exchange_errors import EmptyValuesInBundleError, \
from catalyst.exchange.utils.bundle_utils import range_in_bundle, \
get_bcolz_chunk, get_df_from_arrays, get_assets
from catalyst.exchange.utils.datetime_utils import get_start_dt, \
get_period_label, get_month_start_end, get_year_start_end, get_period, \
timestr_to_dt
from catalyst.exchange.utils.exchange_utils import get_exchange_folder
get_period_label, get_month_start_end, get_year_start_end
from catalyst.exchange.utils.exchange_utils import get_exchange_folder, \
save_exchange_symbols, mixin_market_params, get_catalyst_symbol
from catalyst.utils.cli import maybe_show_progress
from catalyst.utils.paths import ensure_directory
from logbook import Logger
from pytz import UTC
from six import itervalues
log = Logger('exchange_bundle', level=LOG_LEVEL)
@@ -234,12 +232,12 @@ class ExchangeBundle:
problem = '{name} ({start_dt} to {end_dt}) has empty ' \
'periods: {dates}'.format(
name=asset.symbol,
start_dt=asset.start_date.strftime(
DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT),
dates=[date.strftime(
DATE_TIME_FORMAT) for date in dates])
name=asset.symbol,
start_dt=asset.start_date.strftime(
DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT),
dates=[date.strftime(
DATE_TIME_FORMAT) for date in dates])
if empty_rows_behavior == 'warn':
log.warn(problem)
@@ -288,12 +286,12 @@ class ExchangeBundle:
problem = '{name} ({start_dt} to {end_dt}) has {threshold} ' \
'identical close values on: {dates}'.format(
name=asset.symbol,
start_dt=asset.start_date.strftime(DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT),
threshold=threshold,
dates=[pd.to_datetime(date).strftime(DATE_TIME_FORMAT)
for date in dates])
name=asset.symbol,
start_dt=asset.start_date.strftime(DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT),
threshold=threshold,
dates=[pd.to_datetime(date).strftime(DATE_TIME_FORMAT)
for date in dates])
problems.append(problem)
@@ -460,7 +458,7 @@ class ExchangeBundle:
last_entry = None
if start is None or \
(earliest_trade is not None and earliest_trade > start):
(earliest_trade is not None and earliest_trade > start):
start = earliest_trade
if last_entry is not None and (end is None or end > last_entry):
@@ -514,8 +512,8 @@ class ExchangeBundle:
continue
dates = pd.date_range(
start=get_period(adj_start, data_frequency),
end=get_period(adj_end, data_frequency),
start=get_period_label(adj_start, data_frequency),
end=get_period_label(adj_end, data_frequency),
freq='MS' if data_frequency == 'minute' else 'AS',
tz=UTC
)
@@ -554,9 +552,7 @@ class ExchangeBundle:
# We sort the chunks by end date to ingest most recent data first
chunks[asset].sort(
key=lambda chunk: timestr_to_dt(
chunk['period'], data_frequency
)
key=lambda chunk: pd.to_datetime(chunk['period'])
)
return chunks
@@ -602,17 +598,41 @@ class ExchangeBundle:
# we want to give an end_date far in time
writer = self.get_writer(start_dt, end_dt, data_frequency)
if show_breakdown:
for asset in chunks:
if chunks:
for asset in chunks:
with maybe_show_progress(
chunks[asset],
show_progress,
label='Ingesting {frequency} price data for '
'{symbol} on {exchange}'.format(
exchange=self.exchange_name,
frequency=data_frequency,
symbol=asset.symbol
)) as it:
for chunk in it:
problems += self.ingest_ctable(
asset=chunk['asset'],
data_frequency=data_frequency,
period=chunk['period'],
writer=writer,
empty_rows_behavior='strip',
cleanup=True
)
else:
all_chunks = list(chain.from_iterable(itervalues(chunks)))
# We sort the chunks by end date to ingest most recent data first
if all_chunks:
all_chunks.sort(
key=lambda chunk: pd.to_datetime(chunk['period'])
)
with maybe_show_progress(
chunks[asset],
all_chunks,
show_progress,
label='Ingesting {frequency} price data for '
'{symbol} on {exchange}'.format(
label='Ingesting {frequency} price data on '
'{exchange}'.format(
exchange=self.exchange_name,
frequency=data_frequency,
symbol=asset.symbol
)
) as it:
)) as it:
for chunk in it:
problems += self.ingest_ctable(
asset=chunk['asset'],
@@ -622,33 +642,6 @@ class ExchangeBundle:
empty_rows_behavior='strip',
cleanup=True
)
else:
all_chunks = list(chain.from_iterable(itervalues(chunks)))
# We sort the chunks by end date to ingest most recent data first
all_chunks.sort(
key=lambda chunk: timestr_to_dt(
chunk['period'], data_frequency
)
)
with maybe_show_progress(
all_chunks,
show_progress,
label='Ingesting {frequency} price data on '
'{exchange}'.format(
exchange=self.exchange_name,
frequency=data_frequency,
)
) as it:
for chunk in it:
problems += self.ingest_ctable(
asset=chunk['asset'],
data_frequency=data_frequency,
period=chunk['period'],
writer=writer,
empty_rows_behavior='strip',
cleanup=True
)
if show_report and len(problems) > 0:
log.info('problems during ingestion:{}\n'.format(
@@ -708,36 +701,42 @@ class ExchangeBundle:
for symbol in symbols:
start_dt = df.index.get_level_values(1).min()
end_dt = df.index.get_level_values(1).max()
end_dt_key = 'end_{}'.format(data_frequency)
try:
asset = self.exchange.get_asset(symbol, is_local=True)
except:
asset = copy.deepcopy(self.exchange.get_asset(symbol))
market = self.exchange.get_market(symbol)
if market is None:
raise ValueError('symbol not available in the exchange.')
if asset.data_source == 'local':
asset.start_date = asset.start_date \
if asset.start_date < start_dt else start_dt
params = dict(
exchange=self.exchange.name,
data_source='local',
exchange_symbol=market['id'],
)
mixin_market_params(self.exchange_name, params, market)
if data_frequency == 'daily':
asset.end_date = asset.end_daily = asset.end_daily \
if asset.end_daily > end_dt else end_dt
asset_def = self.exchange.get_asset_def(market, True)
if asset_def is not None:
params['symbol'] = asset_def['symbol']
else:
asset.end_date = asset.end_minute = asset.end_minute \
if asset.end_minute > end_dt else end_dt
params['start_date'] = asset_def['start_date'] \
if asset_def['start_date'] < start_dt else start_dt
params['end_date'] = asset_def[end_dt_key] \
if asset_def[end_dt_key] > end_dt else end_dt
params['end_daily'] = end_dt \
if data_frequency == 'daily' else asset_def['end_daily']
params['end_minute'] = end_dt \
if data_frequency == 'minute' else asset_def['end_minute']
else:
asset.data_source = 'local'
asset.start_date = start_dt
asset.end_dt = end_dt
params['symbol'] = get_catalyst_symbol(market)
if data_frequency == 'daily':
asset.end_daily = end_dt
asset.end_minute = None
else:
asset.end_daily = None
asset.end_minute = end_dt
params['end_daily'] = end_dt \
if data_frequency == 'daily' else 'N/A'
params['end_minute'] = end_dt \
if data_frequency == 'minute' else 'N/A'
if min_start_dt is None or start_dt < min_start_dt:
min_start_dt = start_dt
@@ -745,9 +744,11 @@ class ExchangeBundle:
if max_end_dt is None or end_dt > max_end_dt:
max_end_dt = end_dt
assets[symbol] = asset
asset = TradingPair(**params)
assets[market['id']] = asset
save_exchange_symbols(self.exchange_name, assets, True)
# TODO: update config.json
writer = self.get_writer(
start_dt=min_start_dt.replace(hour=00, minute=00),
end_dt=max_end_dt.replace(hour=23, minute=59),
+2 -2
View File
@@ -296,7 +296,7 @@ class DataPortalExchangeBacktest(DataPortalExchangeBase):
bundle = self.exchange_bundles[exchange_name] # type: ExchangeBundle
freq, candle_size, unit, adj_data_frequency = get_frequency(
frequency, data_frequency
frequency, data_frequency, supported_freqs=['T', 'D']
)
adj_bar_count = candle_size * bar_count
@@ -312,7 +312,7 @@ class DataPortalExchangeBacktest(DataPortalExchangeBase):
algo_end_dt=self._last_available_session,
)
start_dt = get_start_dt(end_dt, adj_bar_count, data_frequency)
start_dt = get_start_dt(end_dt, adj_bar_count, adj_data_frequency)
df = resample_history_df(pd.DataFrame(series), freq, field, start_dt)
return df
-14
View File
@@ -329,17 +329,3 @@ class NoCandlesReceivedFromExchange(ZiplineError):
'Although requesting {bar_count} candles until {end_dt} of asset {asset}, '
'an empty list of candles was received for {exchange}.'
).strip()
class MarketsNotFoundError(ZiplineError):
msg = (
'Exchange {exchange} contains no valid market so it is unusable in '
'Catalyst.'
).strip()
class InvalidMarketError(ZiplineError):
msg = (
'Exchange {exchange} contains at least one incorrectly structured '
'market: {market}, so it is unusable in Catalyst.'
).strip()
+5 -7
View File
@@ -5,7 +5,6 @@ from datetime import datetime
import numpy as np
import pandas as pd
from catalyst.constants import BUNDLE_URL
from catalyst.data.bundles.core import download_without_progress
from catalyst.exchange.utils.exchange_utils import get_exchange_bundles_folder
import os
@@ -49,11 +48,10 @@ def get_bcolz_chunk(exchange_name, symbol, data_frequency, period):
path = os.path.join(root, name)
if not os.path.isdir(path):
url = BUNDLE_URL.format(
url = 'https://s3.amazonaws.com/enigmaco/catalyst-bundles/' \
'exchange-{exchange}/{name}.tar.gz'.format(
exchange=exchange_name,
data_frequency=data_frequency,
name=name,
)
name=name)
bytes = download_without_progress(url)
with tarfile.open('r', fileobj=bytes) as tar:
@@ -77,14 +75,14 @@ def get_df_from_arrays(arrays, periods):
"""
ohlcv = dict()
for index, field in enumerate(['open', 'high', 'low', 'close', 'volume']):
for index, field in enumerate(
['open', 'high', 'low', 'close', 'volume']):
ohlcv[field] = arrays[index].flatten()
df = pd.DataFrame(
data=ohlcv,
index=periods
)
df.index.name = 'last_traded'
return df
-307
View File
@@ -1,307 +0,0 @@
import json
import os
import pandas as pd
from six.moves.urllib import request
from catalyst.assets._assets import TradingPair
from ccxt import NetworkError
from catalyst.constants import LOG_LEVEL, EXCHANGE_CONFIG_URL
from catalyst.exchange.exchange_errors import MarketsNotFoundError, \
InvalidMarketError
from catalyst.exchange.utils.exchange_utils import get_catalyst_symbol, \
get_exchange_folder, get_exchange_auth
from catalyst.exchange.utils.serialization_utils import ExchangeJSONDecoder, \
ExchangeJSONEncoder
from logbook import Logger
from redo import retry
from ccxt.base.exchange import Exchange
from catalyst.utils.paths import last_modified_time, data_root, \
ensure_directory
import ccxt
log = Logger('ccxt_utils', level=LOG_LEVEL)
def scan_exchange_configs(features=None, history=None, is_authenticated=False,
path=None):
"""
Finding exchanges from their config files
Parameters
----------
features
is_authenticated
Returns
-------
"""
for exchange_name in ccxt.exchanges:
config = get_exchange_config(exchange_name, path)
if not config or 'error' in config:
log.info(
'skipping invalid exchange {}'.format(exchange_name)
)
# Check if the exchange has an auth.json file
if is_authenticated:
exchange_auth = get_exchange_auth(exchange_name)
has_auth = (exchange_auth['key'] != ''
and exchange_auth['secret'] != '')
if not has_auth:
continue
if features is None:
has_features = True
else:
try:
supported_features = [
feature for feature in features if
feature in config['features']
]
has_features = len(supported_features) > 0
except Exception:
has_features = False
# TODO: filter by history
if has_features:
yield config
def get_exchange_config(exchange_name, path=None, environ=None,
expiry='1H'):
"""
The de-serialized content of the exchange's config.json.
Parameters
----------
exchange_name: str
The exchange name
filename: str
The target file
environ:
Returns
-------
config: dict[srt, Object]
The config dictionary.
"""
try:
if path is None:
root = data_root(environ)
path = os.path.join(root, 'exchanges')
folder = os.path.join(path, exchange_name)
ensure_directory(folder)
filename = os.path.join(folder, 'config.json')
url = EXCHANGE_CONFIG_URL.format(exchange=exchange_name)
if os.path.isfile(filename):
# If the file exists, only update periodically to avoid
# unnecessary calls
now = pd.Timestamp.utcnow()
limit = pd.Timedelta(expiry)
if pd.Timedelta(now - last_modified_time(filename)) > limit:
try:
request.urlretrieve(url=url, filename=filename)
except Exception as e:
log.warn(
'unable to update config {} => {}: {}'.format(
url, filename, e
)
)
else:
request.urlretrieve(url=url, filename=filename)
with open(filename) as data_file:
data = json.load(data_file, cls=ExchangeJSONDecoder)
return data
except Exception as e:
log.warn(
'unable to download {} config: {}'.format(
exchange_name, e
)
)
return dict(error=e)
def save_exchange_config(config, filename=None, environ=None):
"""
Save assets into an exchange_config file.
Parameters
----------
exchange_name: str
config
environ
Returns
-------
"""
if filename is None:
name = 'config.json'
exchange_folder = get_exchange_folder(config['id'], environ)
filename = os.path.join(exchange_folder, name)
with open(filename, 'w+') as handle:
json.dump(config, handle, indent=4, cls=ExchangeJSONEncoder)
def fetch_markets(ccxt_exchange):
"""
Fetches CCXT market objects.
Parameters
----------
ccxt_exchange: Exchange
Returns
-------
"""
markets_symbols = ccxt_exchange.load_markets()
log.debug(
'fetching {} markets:\n{}'.format(
ccxt_exchange.name, markets_symbols
)
)
markets = ccxt_exchange.fetch_markets()
if not markets:
raise MarketsNotFoundError(
exchange=ccxt_exchange.name,
)
for market in markets:
if 'id' not in market:
raise InvalidMarketError(
exchange=ccxt_exchange.name,
market=market,
)
return markets
def create_exchange_config(ccxt_exchange):
"""
Creates an exchange config structure.
Parameters
----------
ccxt_exchange: Exchange
Returns
-------
"""
exchange_name = ccxt_exchange.__class__.__name__
config = dict(
id=exchange_name,
name=ccxt_exchange.name,
features=[
feature for feature in ccxt_exchange.has if
ccxt_exchange.has[feature]
]
)
markets = retry(
action=fetch_markets,
attempts=5,
sleeptime=5,
retry_exceptions=(NetworkError,),
cleanup=lambda: log.warn(
'fetching markets again for {}'.format(exchange_name)
),
args=(ccxt_exchange,)
)
config['assets'] = []
for market in markets:
asset = create_trading_pair(exchange_name, market)
config['assets'].append(asset)
return config
def create_trading_pair(exchange_name, market, start_dt=None, end_dt=None,
leverage=1, end_daily=None, end_minute=None):
"""
Creating a TradingPair from market and asset data.
Parameters
----------
market: dict[str, Object]
start_dt
end_dt
leverage
end_daily
end_minute
Returns
-------
"""
params = dict(
exchange=exchange_name,
data_source='catalyst',
exchange_symbol=market['id'],
symbol=get_catalyst_symbol(market),
start_date=start_dt,
end_date=end_dt,
leverage=leverage,
asset_name=market['symbol'],
end_daily=end_daily,
end_minute=end_minute,
)
apply_conditional_market_params(exchange_name, params, market)
return TradingPair(**params)
def apply_conditional_market_params(exchange_name, params, market):
"""
Applies a CCXT market dict to parameters of TradingPair init.
Parameters
----------
params: dict[Object]
market: dict[Object]
Returns
-------
"""
# TODO: make this more externalized / configurable
# Consider representing in some type of JSON structure
if 'active' in market:
params['trading_state'] = 1 if market['active'] else 0
else:
params['trading_state'] = 1
if 'lot' in market:
params['min_trade_size'] = market['lot']
params['lot'] = market['lot']
if exchange_name == 'bitfinex':
params['maker'] = 0.001
params['taker'] = 0.002
elif 'maker' in market and 'taker' in market \
and market['maker'] is not None \
and market['taker'] is not None:
params['maker'] = market['maker']
params['taker'] = market['taker']
else:
# TODO: default commission, make configurable
params['maker'] = 0.0015
params['taker'] = 0.0025
info = market['info'] if 'info' in market else None
if info:
if 'minimum_order_size' in info:
params['min_trade_size'] = float(info['minimum_order_size'])
if 'lot' not in params:
params['lot'] = params['min_trade_size']
+38 -30
View File
@@ -1,4 +1,5 @@
import calendar
import math
import re
from datetime import datetime, timedelta, date
@@ -164,12 +165,6 @@ def get_start_dt(end_dt, bar_count, data_frequency, include_first=True):
return start_dt
def timestr_to_dt(timestr, data_frequency):
dt_format = '%Y' if data_frequency == 'daily' else '%Y%m'
dt = pd.to_datetime(timestr, format=dt_format, utc=True)
return dt
def get_period_label(dt, data_frequency):
"""
The period label for the specified date and frequency.
@@ -183,26 +178,6 @@ def get_period_label(dt, data_frequency):
-------
str
"""
if data_frequency == 'minute':
return '{}{:02d}'.format(dt.year, dt.month)
else:
return '{}'.format(dt.year)
def get_period(dt, data_frequency):
"""
The period label for the specified date and frequency.
Parameters
----------
dt: datetime
data_frequency: str
Returns
-------
str
"""
if data_frequency == 'minute':
return '{}-{:02d}'.format(dt.year, dt.month)
@@ -274,9 +249,12 @@ def get_year_start_end(dt, first_day=None, last_day=None):
return year_start, year_end
def get_frequency(freq, data_frequency=None, supported_freqs=['D', 'T']):
def get_frequency(freq, data_frequency=None, supported_freqs=['D', 'H', 'T']):
"""
Get the frequency parameters.
Takes an arbitrary candle size (e.g. 15T) and converts to the lowest
common denominator supported by the data bundles (e.g. 1T). The data
bundles only support 1T and 1D frequencies. If another frequency
is requested, Catalyst must request the underlying data and resample.
Notes
-----
@@ -331,14 +309,14 @@ def get_frequency(freq, data_frequency=None, supported_freqs=['D', 'T']):
data_frequency = 'minute'
elif unit.lower() == 'h':
data_frequency = 'minute'
if 'H' in supported_freqs:
unit = 'H'
alias = '{}H'.format(candle_size)
else:
candle_size = candle_size * 60
alias = '{}T'.format(candle_size)
data_frequency = 'minute'
else:
raise InvalidHistoryFrequencyAlias(freq=freq)
@@ -352,3 +330,33 @@ def from_ms_timestamp(ms):
def get_epoch():
return pd.to_datetime('1970-1-1', utc=True)
def get_candles_number_from_minutes(unit, candle_size, minutes):
"""
Get the number of bars needed for the given time interval
in minutes.
Notes
-----
Supports only "T", "D" and "H" units
Parameters
----------
unit: str
candle_size : int
minutes: int
Returns
-------
int
"""
if unit == "T":
res = (float(minutes) / candle_size)
elif unit == "H":
res = (minutes / 60.0) / candle_size
else: # unit == "D"
res = (minutes / 1440.0) / candle_size
return int(math.ceil(res))
+105 -75
View File
@@ -1,19 +1,19 @@
import hashlib
import os
import shutil
import json
import pandas as pd
import os
import pickle
from catalyst.assets._assets import TradingPair
import shutil
from datetime import date, datetime
import pandas as pd
from catalyst.assets._assets import TradingPair
from six import string_types
from six.moves.urllib import request
from catalyst.constants import EXCHANGE_CONFIG_URL
from catalyst.constants import DATE_FORMAT, SYMBOLS_URL
from catalyst.exchange.exchange_errors import ExchangeSymbolsNotFound
from catalyst.exchange.utils.serialization_utils import ExchangeJSONEncoder, \
ExchangeJSONDecoder, ConfigJSONEncoder
from catalyst.utils.deprecate import deprecated
ExchangeJSONDecoder
from catalyst.utils.paths import data_root, ensure_directory, \
last_modified_time
@@ -69,7 +69,7 @@ def is_blacklist(exchange_name, environ=None):
return os.path.exists(filename)
def get_exchange_config_filename(exchange_name, environ=None):
def get_exchange_symbols_filename(exchange_name, is_local=False, environ=None):
"""
The absolute path of the exchange's symbol.json file.
@@ -83,12 +83,12 @@ def get_exchange_config_filename(exchange_name, environ=None):
str
"""
name = 'config.json'
name = 'symbols.json' if not is_local else 'symbols_local.json'
exchange_folder = get_exchange_folder(exchange_name, environ)
return os.path.join(exchange_folder, name)
def download_exchange_config(exchange_name, filename, environ=None):
def download_exchange_symbols(exchange_name, environ=None):
"""
Downloads the exchange's symbols.json from the repository.
@@ -102,14 +102,15 @@ def download_exchange_config(exchange_name, filename, environ=None):
str
"""
url = EXCHANGE_CONFIG_URL.format(exchange=exchange_name)
request.urlretrieve(url=url, filename=filename)
filename = get_exchange_symbols_filename(exchange_name)
url = SYMBOLS_URL.format(exchange=exchange_name)
response = request.urlretrieve(url=url, filename=filename)
return response
@deprecated
def get_exchange_config(exchange_name, filename=None, environ=None):
def get_exchange_symbols(exchange_name, is_local=False, environ=None):
"""
The de-serialized content of the exchange's config.json.
The de-serialized content of the exchange's symbols.json.
Parameters
----------
@@ -122,48 +123,55 @@ def get_exchange_config(exchange_name, filename=None, environ=None):
Object
"""
if filename is None:
filename = get_exchange_config_filename(exchange_name)
filename = get_exchange_symbols_filename(exchange_name, is_local)
if not is_local and (not os.path.isfile(filename) or pd.Timedelta(
pd.Timestamp('now', tz='UTC') - last_modified_time(
filename)).days > 1):
try:
download_exchange_symbols(exchange_name, environ)
except Exception:
pass
if os.path.isfile(filename):
now = pd.Timestamp.utcnow()
limit = pd.Timedelta('2H')
if pd.Timedelta(now - last_modified_time(filename)) > limit:
download_exchange_config(exchange_name, filename, environ)
with open(filename) as data_file:
try:
data = json.load(data_file, cls=ExchangeJSONDecoder)
return data
except ValueError:
return dict()
else:
download_exchange_config(exchange_name, filename, environ)
with open(filename) as data_file:
try:
data = json.load(data_file, cls=ExchangeJSONDecoder)
return data
except ValueError:
return dict()
raise ExchangeSymbolsNotFound(
exchange=exchange_name,
filename=filename
)
def save_exchange_config(exchange_name, config, filename=None, environ=None):
def save_exchange_symbols(exchange_name, assets, is_local=False, environ=None):
"""
Save assets into an exchange_config file.
Save assets into an exchange_symbols file.
Parameters
----------
exchange_name: str
config
assets: list[dict[str, object]]
is_local: bool
environ
Returns
-------
"""
if filename is None:
name = 'config.json'
exchange_folder = get_exchange_folder(exchange_name, environ)
filename = os.path.join(exchange_folder, name)
asset_dicts = dict()
for symbol in assets:
asset_dicts[symbol] = assets[symbol].to_dict()
with open(filename, 'w+') as handle:
json.dump(config, handle, indent=4, cls=ConfigJSONEncoder)
filename = get_exchange_symbols_filename(
exchange_name, is_local, environ
)
with open(filename, 'wt') as handle:
json.dump(asset_dicts, handle, indent=4, default=symbols_serial)
def get_symbols_string(assets):
@@ -504,6 +512,25 @@ def has_bundle(exchange_name, data_frequency, environ=None):
return os.path.isdir(folder)
def symbols_serial(obj):
"""
JSON serializer for objects not serializable by default json code
Parameters
----------
obj: Object
Returns
-------
str
"""
if isinstance(obj, (datetime, date)):
return obj.floor('1D').strftime(DATE_FORMAT)
raise TypeError("Type %s not serializable" % type(obj))
def perf_serial(obj):
"""
JSON serializer for objects not serializable by default json code
@@ -593,12 +620,46 @@ def resample_history_df(df, freq, field, start_dt=None):
return resampled_df
def from_ms_timestamp(ms):
return pd.to_datetime(ms, unit='ms', utc=True)
def mixin_market_params(exchange_name, params, market):
"""
Applies a CCXT market dict to parameters of TradingPair init.
Parameters
----------
params: dict[Object]
market: dict[Object]
def get_epoch():
return pd.to_datetime('1970-1-1', utc=True)
Returns
-------
"""
# TODO: make this more externalized / configurable
if 'lot' in market:
params['min_trade_size'] = market['lot']
params['lot'] = market['lot']
if exchange_name == 'bitfinex':
params['maker'] = 0.001
params['taker'] = 0.002
elif 'maker' in market and 'taker' in market and \
market['maker'] is not None and market['taker'] is not None:
params['maker'] = market['maker']
params['taker'] = market['taker']
else:
# TODO: default commission, make configurable
params['maker'] = 0.0015
params['taker'] = 0.0025
info = market['info'] if 'info' in market else None
if info:
if 'minimum_order_size' in info:
params['min_trade_size'] = float(info['minimum_order_size'])
if 'lot' not in params:
params['lot'] = params['min_trade_size']
def group_assets_by_exchange(assets):
@@ -687,37 +748,6 @@ def get_candles_df(candles, field, freq, bar_count, end_dt):
all_series[asset] = pd.Series(asset_df[field])
df = pd.DataFrame(all_series)
df.dropna(inplace=True)
return df
def get_trades_df(trades):
df = pd.DataFrame(trades)
df.index = pd.to_datetime(df.pop('datetime'))
df.index = df.index.tz_localize('UTC')
return df
def candles_from_trades(trades_df, freq):
"""
Calculate OHLCV from candles.
Parameters
----------
trades_df
freq
Returns
-------
"""
df = trades_df['price'].resample(freq).ohlc() # type: pd.DataFrame
df['volume'] = trades_df['amount'].resample(freq).sum()
df.dropna(axis=0, how='all', inplace=True)
df.sort_index(inplace=True, ascending=False)
return df
+23 -33
View File
@@ -4,9 +4,8 @@ from catalyst.constants import LOG_LEVEL
from catalyst.exchange.ccxt.ccxt_exchange import CCXT
from catalyst.exchange.exchange import Exchange
from catalyst.exchange.exchange_errors import ExchangeAuthEmpty
from catalyst.exchange.utils.ccxt_utils import scan_exchange_configs
from catalyst.exchange.utils.exchange_utils import get_exchange_auth, \
get_exchange_folder
get_exchange_folder, is_blacklist
from logbook import Logger
log = Logger('factory', level=LOG_LEVEL)
@@ -14,12 +13,9 @@ exchange_cache = dict()
def get_exchange(exchange_name, base_currency=None, must_authenticate=False,
skip_init=False, auth_alias=None, config=None):
skip_init=False, auth_alias=None):
key = (exchange_name, base_currency)
if key in exchange_cache:
if not skip_init:
exchange_cache[key].init()
return exchange_cache[key]
exchange_auth = get_exchange_auth(exchange_name, alias=auth_alias)
@@ -40,7 +36,6 @@ def get_exchange(exchange_name, base_currency=None, must_authenticate=False,
password=exchange_auth['password'] if 'password'
in exchange_auth.keys() else '',
base_currency=base_currency,
config=config,
)
exchange_cache[key] = exchange
@@ -58,8 +53,8 @@ def get_exchanges(exchange_names):
return exchanges
def find_exchanges(features=None, history=None, skip_blacklist=True, path=None,
is_authenticated=False, base_currency=None):
def find_exchanges(features=None, skip_blacklist=True, is_authenticated=False,
base_currency=None):
"""
Find exchanges filtered by a list of feature.
@@ -77,33 +72,28 @@ def find_exchanges(features=None, history=None, skip_blacklist=True, path=None,
list[Exchange]
"""
exchange_names = CCXT.find_exchanges(features, is_authenticated)
return list(
scan_exchanges(
features,
history,
skip_blacklist,
path,
is_authenticated,
base_currency
)
)
def scan_exchanges(features=None, history=None, skip_blacklist=True, path=None,
is_authenticated=False, base_currency=None):
for config in scan_exchange_configs(
features=features,
history=history,
is_authenticated=is_authenticated,
path=path,
):
if skip_blacklist and (config is None or 'error' in config):
exchanges = []
for exchange_name in exchange_names:
if skip_blacklist and is_blacklist(exchange_name):
continue
yield get_exchange(
exchange_name=config['id'],
exchange = get_exchange(
exchange_name=exchange_name,
skip_init=True,
base_currency=base_currency,
config=config,
)
if features is not None:
if 'dailyBundle' in features \
and not exchange.has_bundle('daily'):
continue
elif 'minuteBundle' in features \
and not exchange.has_bundle('minute'):
continue
exchanges.append(exchange)
return exchanges
+1 -34
View File
@@ -3,48 +3,15 @@ import re
from json import JSONEncoder
import pandas as pd
from catalyst.constants import DATE_TIME_FORMAT
from six import string_types
from datetime import date, datetime
from catalyst.constants import DATE_TIME_FORMAT, DATE_FORMAT
from catalyst.assets._assets import TradingPair
class ConfigJSONEncoder(json.JSONEncoder):
def default(self, obj):
"""
JSON serializer for objects not serializable by default json code
Parameters
----------
obj: Object
Returns
-------
str
"""
if isinstance(obj, (datetime, date)):
return obj.floor('1D').strftime(DATE_FORMAT)
elif isinstance(obj, TradingPair):
return obj.to_dict()
class ExchangeJSONEncoder(json.JSONEncoder):
def default(self, obj):
if isinstance(obj, pd.Timestamp):
return obj.strftime(DATE_TIME_FORMAT)
elif isinstance(obj, TradingPair):
asset = obj.to_dict()
asset['maker'] = round(asset['maker'], asset['decimals'])
asset['taker'] = round(asset['taker'], asset['decimals'])
asset['lot'] = round(asset['lot'], 4)
asset['min_trade_size'] = round(asset['min_trade_size'], 4)
asset['max_trade_size'] = round(asset['max_trade_size'], 4)
return asset
# Let the base class default method raise the TypeError
return JSONEncoder.default(self, obj)
+18 -5
View File
@@ -95,11 +95,24 @@ class TradingEnvironment(object):
if not trading_calendar:
trading_calendar = get_calendar("NYSE")
self.benchmark_returns, self.treasury_curves = load(
trading_calendar.day,
trading_calendar.schedule.index,
self.bm_symbol,
)
# todo: uncomment and add a well defined benchmark
# self.benchmark_returns, self.treasury_curves = load(
# trading_calendar.day,
# trading_calendar.schedule.index,
# self.bm_symbol,
# exchange=exchange,
# )
start_data = get_calendar('OPEN').first_trading_session
end_data = pd.Timestamp.utcnow()
treasure_cols = ['1month', '3month', '6month', '1year', '2year',
'3year', '5year', '7year', '10year', '20year', '30year']
self.benchmark_returns = pd.DataFrame(data=0.001,
index=pd.date_range(start_data, end_data),
columns=['close'])
self.treasury_curves = pd.DataFrame(data=0.001,
index=pd.date_range(start_data, end_data),
columns=treasure_cols)
self.exchange_tz = exchange_tz
@@ -1 +1 @@
0x7fAec9aaE31BE428DeAAE1be8195dF609079Fd10
0x39a54f480d922a58c963de8091a6c9afc69db2cf
File diff suppressed because one or more lines are too long
@@ -1 +1 @@
0x3985f5de8fddf2e8f7705cd360b498bf35ebfbc4
0xa2b37c6cd52f60fd4eb46ca59fafcf22d081aebc
+25 -21
View File
@@ -20,6 +20,7 @@ from requests_toolbelt.multipart.decoder import \
from catalyst.constants import (
LOG_LEVEL, AUTH_SERVER, ETH_REMOTE_NODE, MARKETPLACE_CONTRACT,
MARKETPLACE_CONTRACT_ABI, ENIGMA_CONTRACT, ENIGMA_CONTRACT_ABI)
from catalyst.utils.cli import maybe_show_progress
from catalyst.exchange.utils.stats_utils import set_print_settings
from catalyst.marketplace.marketplace_errors import (
MarketplacePubAddressEmpty, MarketplaceDatasetNotFound,
@@ -126,9 +127,10 @@ class Marketplace:
else:
while True:
for i in range(0, len(self.addresses)):
print('{}\t{}\t{}'.format(
print('{}\t{}\t{}\t{}'.format(
i,
self.addresses[i]['pubAddr'],
self.addresses[i]['wallet'].ljust(10),
self.addresses[i]['desc'])
)
address_i = int(input('Choose your address associated with '
@@ -145,7 +147,7 @@ class Marketplace:
def sign_transaction(self, tx):
url = 'https://www.myetherwallet.com/#offline-transaction'
url = 'https://www.mycrypto.com/#offline-transaction'
print('\nVisit {url} and enter the following parameters:\n\n'
'From Address:\t\t{_from}\n'
'\n\tClick the "Generate Information" button\n\n'
@@ -177,10 +179,12 @@ class Marketplace:
def check_transaction(self, tx_hash):
if 'ropsten' in ETH_REMOTE_NODE:
etherscan = 'https://ropsten.etherscan.io/tx/{}'.format(
tx_hash)
etherscan = 'https://ropsten.etherscan.io/tx/'
elif 'rinkeby' in ETH_REMOTE_NODE:
etherscan = 'https://rinkeby.etherscan.io/tx/'
else:
etherscan = 'https://etherscan.io/tx/{}'.format(tx_hash)
etherscan = 'https://etherscan.io/tx/'
etherscan = '{}{}'.format(etherscan, tx_hash)
print('\nYou can check the outcome of your transaction here:\n'
'{}\n\n'.format(etherscan))
@@ -329,9 +333,6 @@ class Marketplace:
'nonce': self.web3.eth.getTransactionCount(address)}
)
if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(tx)
try:
tx_hash = '0x{}'.format(
@@ -371,9 +372,6 @@ class Marketplace:
'from': address,
'nonce': self.web3.eth.getTransactionCount(address)})
if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(tx)
try:
@@ -434,10 +432,9 @@ class Marketplace:
merge_bundles(zsource, ztarget)
else:
shutil.rmtree(bundle_folder, ignore_errors=True)
os.rename(tmp_bundle, bundle_folder)
pass
def ingest(self, ds_name=None, start=None, end=None, force_download=False):
if ds_name is None:
@@ -502,20 +499,29 @@ class Marketplace:
key = self.addresses[address_i]['key']
secret = self.addresses[address_i]['secret']
else:
key, secret = get_key_secret(address)
key, secret = get_key_secret(address,
self.addresses[address_i]['wallet'])
headers = get_signed_headers(ds_name, key, secret)
log.debug('Starting download of dataset for ingestion...')
log.info('Starting download of dataset for ingestion...')
r = requests.post(
'{}/marketplace/ingest'.format(AUTH_SERVER),
headers=headers,
stream=True,
)
if r.status_code == 200:
log.info('Dataset downloaded successfully. Processing dataset...')
target_path = get_temp_bundles_folder()
try:
decoder = MultipartDecoder.from_response(r)
# with maybe_show_progress(
# iter(decoder.parts),
# True,
# label='Processing files') as part:
counter = 0
for part in decoder.parts:
log.info("Processing file {} of {}".format(
counter, len(decoder.parts)))
h = part.headers[b'Content-Disposition'].decode('utf-8')
# Extracting the filename from the header
name = re.search(r'filename="(.*)"', h).group(1)
@@ -529,6 +535,7 @@ class Marketplace:
f.write(part.content)
self.process_temp_bundle(ds_name, filename)
counter += 1
except NonMultipartContentTypeException:
response = r.json()
@@ -596,7 +603,6 @@ class Marketplace:
folder = get_bundle_folder(ds_name, data_frequency)
shutil.rmtree(folder)
pass
def create_metadata(self, key, secret, ds_name, data_frequency, desc,
has_history=True, has_live=True):
@@ -688,7 +694,8 @@ class Marketplace:
key = self.addresses[address_i]['key']
secret = self.addresses[address_i]['secret']
else:
key, secret = get_key_secret(address)
key, secret = get_key_secret(address,
self.addresses[address_i]['wallet'])
grains = to_grains(price)
@@ -701,9 +708,6 @@ class Marketplace:
'nonce': self.web3.eth.getTransactionCount(address)}
)
if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(tx)
try:
@@ -772,7 +776,7 @@ class Marketplace:
key = match['key']
secret = match['secret']
else:
key, secret = get_key_secret(provider_info[0])
key, secret = get_key_secret(provider_info[0], match['wallet'])
headers = get_signed_headers(dataset, key, secret)
filenames = glob.glob(os.path.join(datadir, '*.csv'))
+18 -8
View File
@@ -1,5 +1,6 @@
import hashlib
import hmac
import webbrowser
import requests
import time
@@ -9,10 +10,10 @@ from catalyst.marketplace.marketplace_errors import (
MarketplaceEmptySignature)
from catalyst.marketplace.utils.path_utils import (
get_user_pubaddr, save_user_pubaddr)
from catalyst.constants import AUTH_SERVER
from catalyst.constants import AUTH_SERVER, SUPPORTED_WALLETS
def get_key_secret(pubAddr, wallet='mew'):
def get_key_secret(pubAddr, wallet):
"""
Obtain a new key/secret pair from authentication server
@@ -42,14 +43,22 @@ def get_key_secret(pubAddr, wallet='mew'):
auth_type, auth_info = header.split(None, 1)
d = requests.utils.parse_dict_header(auth_info)
nonce = '0x{}'.format(d['nonce'])
nonce = 'Catalyst nonce: 0x{}'.format(d['nonce'])
if wallet in SUPPORTED_WALLETS:
url = 'https://www.mycrypto.com/signmsg.html'
if wallet == 'mew':
print('\nObtaining a key/secret pair to streamline all future '
'requests with the authentication server.\n'
'Visit https://www.myetherwallet.com/signmsg.html and sign the '
'following message:\n{}'.format(nonce))
signature = input('Copy and Paste the "sig" field from '
'Visit {url} and sign the '
'following message (copy the entire line, without the '
'line break at the end):\n\n{nonce}'.format(
url=url,
nonce=nonce))
webbrowser.open_new(url)
signature = input('\nCopy and Paste the "sig" field from '
'the signature here (without the double quotes, '
'only the HEX value):\n')
else:
@@ -83,7 +92,8 @@ def get_key_secret(pubAddr, wallet='mew'):
addresses = get_user_pubaddr()
match = next((l for l in addresses if
l['pubAddr'] == pubAddr), None)
l['pubAddr'].lower() == pubAddr.lower()), None)
match['key'] = response.json()['key']
match['secret'] = response.json()['secret']
+49 -2
View File
@@ -2,6 +2,7 @@ import os
import json
import tarfile
from catalyst.constants import SUPPORTED_WALLETS
from catalyst.utils.deprecate import deprecated
from catalyst.utils.paths import data_root, ensure_directory
from catalyst.marketplace.marketplace_errors import MarketplaceJSONError
@@ -131,17 +132,63 @@ def get_user_pubaddr(environ=None):
try:
d = data[0]['pubAddr']
except Exception as e:
return [data, ]
data = [data, ]
changed = False
for idx, d in enumerate(data):
try:
if d['wallet'] not in SUPPORTED_WALLETS:
data[idx]['wallet'] = _choose_wallet(
d['pubAddr'], False)
changed = True
except KeyError:
data[idx]['wallet'] = _choose_wallet(
d['pubAddr'], True)
changed = True
if changed:
save_user_pubaddr(data)
return data
else:
data = []
data.append(dict(pubAddr='', desc=''))
data.append(dict(pubAddr='', desc='', wallet=''))
with open(filename, 'w') as f:
json.dump(data, f, sort_keys=False, indent=2,
separators=(',', ':'))
return data
def _choose_wallet(pubAddr, missing):
while True:
if missing:
print('\nYou need to specify a wallet for address '
'{}.'.format(pubAddr))
else:
print('\nThe wallet specified for address {} is not '
'supported.'.format(pubAddr))
print('Please choose among the following options:')
for idx, wallet in enumerate(SUPPORTED_WALLETS):
print('{}\t{}'.format(idx, wallet))
lw = len(SUPPORTED_WALLETS)-1
w = input('Choose a number between 0 and {}: '.format(
lw))
try:
w = int(w)
except ValueError:
print('Enter a number between 0 and {}'.format(lw))
else:
if w not in range(0, lw+1):
print('Enter a number between 0 and '
'{}'.format(lw))
else:
return SUPPORTED_WALLETS[w]
def save_user_pubaddr(data, environ=None):
"""
Saves the user's public addresses and their related metadata in
+49
View File
@@ -0,0 +1,49 @@
import pytz
from datetime import datetime
from catalyst.api import symbol
from catalyst.utils.run_algo import run_algorithm
coin = 'btc'
base_currency = 'usd'
n_candles = 5
def initialize(context):
context.symbol = symbol('%s_%s' % (coin, base_currency))
def handle_data_polo_partial_candles(context, data):
history = data.history(symbol('btc_usdt'), ['volume'],
bar_count=10,
frequency='4H')
print('\nnow: %s\n%s' % (data.current_dt, history))
if not hasattr(context, 'i'):
context.i = 0
context.i += 1
if context.i > 5:
raise Exception('stop')
live = False
if live:
run_algorithm(initialize=lambda ctx: True,
handle_data=handle_data_polo_partial_candles,
exchange_name='poloniex',
base_currency='usdt',
algo_namespace='ns',
live=True,
data_frequency='minute',
capital_base=3000)
else:
run_algorithm(initialize=lambda ctx: True,
handle_data=handle_data_polo_partial_candles,
exchange_name='poloniex',
base_currency='usdt',
algo_namespace='ns',
live=False,
data_frequency='minute',
capital_base=3000,
start=datetime(2018, 2, 2, 0, 0, 0, 0, pytz.utc),
end=datetime(2018, 2, 20, 0, 0, 0, 0, pytz.utc)
)
+14 -16
View File
@@ -1,8 +1,7 @@
from catalyst.api import symbol
from catalyst.utils.run_algo import run_algorithm
coins = ['dash', 'btc', 'dash', 'etc', 'eth', 'ltc', 'nxt', 'rep', 'str',
'xmr', 'xrp', 'zec']
coins = ['dash', 'btc', 'dash', 'etc', 'eth', 'ltc', 'nxt', 'rep', 'str', 'xmr', 'xrp', 'zec']
symbols = None
@@ -14,21 +13,20 @@ def _handle_data(context, data):
global symbols
if symbols is None: symbols = [symbol(c + '_usdt') for c in coins]
print('getting history for: %s' % [s.symbol for s in symbols])
print'getting history for: %s' % [s.symbol for s in symbols]
history = data.history(symbols,
['close', 'volume'],
bar_count=1, # EXCEPTION, Change to 2
frequency='5T')
# print 'history: %s' % history.shape
['close', 'volume'],
bar_count=1, # EXCEPTION, Change to 2
frequency='5T')
#print 'history: %s' % history.shape
run_algorithm(initialize=initialize,
handle_data=_handle_data,
analyze=lambda _, results: True,
exchange_name='poloniex',
base_currency='usdt',
algo_namespace='issue-236',
live=True,
data_frequency='minute',
capital_base=3000,
simulate_orders=True)
analyze=lambda _, results: True,
exchange_name='poloniex',
base_currency='usdt',
algo_namespace='issue-236',
live=True,
data_frequency='minute',
capital_base=3000,
simulate_orders=True)
+35
View File
@@ -0,0 +1,35 @@
import pytz
from datetime import datetime
from catalyst.api import symbol
from catalyst.utils.run_algo import run_algorithm
coin = 'btc'
base_currency = 'usd'
def initialize(context):
context.symbol = symbol('%s_%s' % (coin, base_currency))
def handle_data_polo_partial_candles(context, data):
history = data.history(symbol('btc_usdt'), ['volume'],
bar_count=10,
frequency='1D')
print('\nnow: %s\n%s' % (data.current_dt, history))
if not hasattr(context, 'i'):
context.i = 0
context.i += 1
if context.i > 5:
raise Exception('stop')
run_algorithm(initialize=lambda ctx: True,
handle_data=handle_data_polo_partial_candles,
exchange_name='poloniex',
base_currency='usdt',
algo_namespace='ns',
live=False,
data_frequency='minute',
capital_base=3000,
start=datetime(2018, 2, 2, 0, 0, 0, 0, pytz.utc),
end=datetime(2018, 2, 20, 0, 0, 0, 0, pytz.utc))
+11 -1
View File
@@ -143,7 +143,7 @@ with the following steps:
.. code-block:: bash
conda create --name catalyst python=2.7 scipy zlib
conda create --name catalyst python=3.6 scipy zlib
3. Activate the environment:
@@ -314,6 +314,16 @@ Troubleshooting ``pip`` Install
$ sudo apt-get install python-dev
----
**Issue**:
Missing TA_Lib
**Solution**:
Follow `these instructions
<https://mrjbq7.github.io/ta-lib/install.html>`_ to install the TA_Lib Python wrapper
(and if needed, its underlying C library as well).
.. _pipenv:
Installing with ``pipenv``
+64
View File
@@ -2,6 +2,70 @@
Release Notes
=============
Version 0.5.6
^^^^^^^^^^^^^
**Release Date**: 2018-03-22
Build
~~~~~
- Data Marketplace: ensures compatibility across wallets, now fully supporting
`ledger`, `trezor`, `keystore`, `private key`. Partial support for `metamask`
(includes sign_msg, but not sign_tx). Current support for `Digital Bitbox` is
unknown.
- Data Marketplace: Switched online provider from MyEtherWallet to MyCrypto.
- Data Marketplace: Added progress indicator for data ingestion.
Bug Fixes
~~~~~~~~~
- Changed benchmark to be constant, so it doesn't ingest data at all. Temporary
fix for :issue:`271`, :issue:`285`
Version 0.5.5
^^^^^^^^^^^^^
**Release Date**: 2018-03-19
Bug Fixes
~~~~~~~~~
- Fixed an issue with the data history in daily frequency :issue:`274`
- Fix hourly frequency issues :issue:`227` and :issue:`114`
Version 0.5.4
^^^^^^^^^^^^^
**Release Date**: 2018-03-14
Build
~~~~~
- Switched Data Marketplace from Ropstein testnet to Rinkeby testnet after
incorporating changes resulting from the marketplace contract audit
- Several usability improvements of the Data Marketplace that make the
`--dataset` parameter optional. If it is not included in the command line,
will list available datasets, and let you choose interactively.
Bug Fixes
~~~~~~~~~
- Fix Binance requirement of symbol to be included in the cancelled order
:issue:`204`
- Fix `notenoughcasherror` when an open order is filled minutes later
:issue:`237`
- Properly handle of empty candles received from exchanges :issue:`236`
- Added a function to reduce open orders amount from calculated target/amount
for target orders :issue:`243`
- Fix missing file in live trading mode on date change :issue:`252`,
:issue:`253`
- Upgraded Data Marketplace to Web3==4.0.0b11, which was breaking some
functionality from prior version 4.0.0b7 :issue:`257`
- Always request more data to avoid empty bars and always give the exact bar
number :issue:`260`
Documentation
~~~~~~~~~~~~~
- PyCharm documentation :issue:`195`
- Added TA-Lib troubleshooting instructions
- Added instructions on how to create a Conda environment for Python 3.6, and
updated Visual C++ instructions for Windows and Python 3
- Linking example algorithms in the documentation to their sources
Version 0.5.3
^^^^^^^^^^^^^
**Release Date**: 2018-02-09
+2 -3
View File
@@ -5,7 +5,6 @@ channels:
dependencies:
- certifi=2016.2.28=py27_0
- mkl=2017.0.3
- matplotlib=2.1.2=py36_0
- numpy=1.13.1=py27_0
- openssl=1.0.2l
- pip=9.0.1=py27_1
@@ -22,7 +21,7 @@ dependencies:
- bcolz==0.12.1
- bottleneck==1.2.1
- chardet==3.0.4
- ccxt==1.11.22
- ccxt==1.10.1094
# The Enigma Data Marketplace requires Python3 because it depends on
# web3, which requires Python3, as building its dependencies breaks in Python2
# - web3==4.0.0b7
@@ -40,7 +39,7 @@ dependencies:
- lru-dict==1.1.6
- mako==1.0.7
- markupsafe==1.0
- matplotlib==2.1.0
- matplotlib==2.1.2
- multipledispatch==0.4.9
- networkx==2.0
- numexpr==2.6.4
+1 -1
View File
@@ -31,7 +31,7 @@ dependencies:
- botocore==1.8.41
- bottleneck==1.2.1
- cchardet==2.1.1
- ccxt==1.11.22
- ccxt==1.10.1102
- chardet==3.0.4
- click==6.7
- contextlib2==0.5.5
+1 -1
View File
@@ -81,7 +81,7 @@ empyrical==0.2.1
tables==3.3.0
#Catalyst dependencies
ccxt==1.11.22
ccxt==1.10.1094
boto3==1.4.8
redo==1.6
web3==4.0.0b11; python_version > '3.4'
+12 -30
View File
@@ -14,8 +14,7 @@ from catalyst.exchange.utils.bundle_utils import get_bcolz_chunk, \
from catalyst.exchange.utils.datetime_utils import get_start_dt
from catalyst.exchange.utils.exchange_utils import get_exchange_folder
from catalyst.exchange.utils.factory import get_exchange
from catalyst.exchange.utils.stats_utils import df_to_string, \
set_print_settings
from catalyst.exchange.utils.stats_utils import df_to_string
from catalyst.utils.paths import ensure_directory
log = getLogger('test_exchange_bundle')
@@ -46,9 +45,9 @@ class TestExchangeBundle:
exchange_name = 'binance'
exchange = get_exchange(exchange_name)
exchange_bundle = ExchangeBundle(exchange_name)
exchange_bundle = ExchangeBundle(exchange)
assets = [
exchange.get_asset('bch_eth')
exchange.get_asset('eth_btc')
]
start = pd.to_datetime('2018-03-01', utc=True)
@@ -62,8 +61,7 @@ class TestExchangeBundle:
exclude_symbols=None,
start=start,
end=end,
show_progress=False,
show_breakdown=False
show_progress=True
)
reader = exchange_bundle.get_reader(data_frequency)
@@ -74,15 +72,9 @@ class TestExchangeBundle:
start_dt=start,
end_dt=end
)
periods = exchange_bundle.get_calendar_periods_range(
start, end, data_frequency
print('found {} rows for {} ingestion\n{}'.format(
len(arrays[0]), asset.symbol, arrays[0])
)
dx = get_df_from_arrays(arrays[0], periods)
set_print_settings()
print('found {} rows for last ingestion:\n{}\n{}'.format(
len(dx), dx.head(10), dx.tail(10)
))
pass
def test_ingest_minute_all(self):
@@ -230,14 +222,9 @@ class TestExchangeBundle:
start_dt=start,
end_dt=end
)
periods = exchange_bundle.get_calendar_periods_range(
start, end, data_frequency
print('found {} rows for {} ingestion\n{}'.format(
len(arrays[0]), asset.symbol, arrays[0])
)
dx = get_df_from_arrays(arrays, periods)
print('found {} rows for last ingestion'.format(
len(dx)
))
pass
def test_daily_data_to_minute_table(self):
@@ -303,22 +290,17 @@ class TestExchangeBundle:
for asset in assets:
sid = asset.sid
arrays = reader.load_raw_arrays(
daily_values = reader.load_raw_arrays(
fields=['open', 'high', 'low', 'close', 'volume'],
start_dt=start,
end_dt=end,
sids=[sid],
)
periods = exchange_bundle.get_calendar_periods_range(
start, end, data_frequency
)
dx = get_df_from_arrays(arrays, periods)
print('found {} rows for last ingestion'.format(
len(dx)
))
pass
len(daily_values[0]))
)
pass
def test_minute_bundle(self):
# exchange_name = 'poloniex'
+4 -37
View File
@@ -5,8 +5,7 @@ from catalyst.exchange.utils.stats_utils import set_print_settings
from .base import BaseExchangeTestCase
from catalyst.exchange.ccxt.ccxt_exchange import CCXT
from catalyst.exchange.exchange_execution import ExchangeLimitOrder
from catalyst.exchange.utils.exchange_utils import get_exchange_auth, \
get_trades_df, candles_from_trades
from catalyst.exchange.utils.exchange_utils import get_exchange_auth
from catalyst.finance.order import Order
log = Logger('test_ccxt')
@@ -15,13 +14,12 @@ log = Logger('test_ccxt')
class TestCCXT(BaseExchangeTestCase):
@classmethod
def setup(self):
exchange_name = 'binance'
exchange_name = 'bittrex'
auth = get_exchange_auth(exchange_name)
self.exchange = CCXT(
exchange_name=exchange_name,
key=auth['key'],
secret=auth['secret'],
password=None,
base_currency='usdt',
)
self.exchange.init()
@@ -60,9 +58,9 @@ class TestCCXT(BaseExchangeTestCase):
log.info('retrieving candles')
candles = self.exchange.get_candles(
freq='1T',
assets=[self.exchange.get_asset('eng_eth')],
assets=[self.exchange.get_asset('eth_btc')],
bar_count=200,
start_dt=pd.to_datetime('2017-09-01', utc=True),
# start_dt=pd.to_datetime('2017-09-01', utc=True),
)
for asset in candles:
@@ -92,37 +90,6 @@ class TestCCXT(BaseExchangeTestCase):
assert trades
pass
def test_validate_volume(self):
asset = self.exchange.get_asset('eng_eth')
candles = self.exchange.get_candles(
freq='1T',
assets=[asset],
bar_count=10,
)
df = pd.DataFrame(candles[asset])
df.set_index('last_traded', drop=True, inplace=True)
df.drop_duplicates()
df.sort_index(inplace=True, ascending=False)
assert candles
start_dt = df.index[-1]
trades = self.exchange.get_trades(
asset, start_dt=start_dt, my_trades=False
)
assert trades
trades_df = get_trades_df(trades)
df2 = candles_from_trades(trades_df, '1T')
set_print_settings()
log.info(
'comparing candles / resampled trades:\n{}\n{}'.format(
df, df2
)
)
pass
def test_get_executed_order(self):
log.info('retrieving executed order')
asset = self.exchange.get_asset('eng_eth')
-8
View File
@@ -1,8 +0,0 @@
from catalyst.exchange.utils.factory import get_exchange
class TestConfig:
def test_create_config(self):
exchange = get_exchange('binance', skip_init=True)
config = exchange.create_exchange_config()
pass
+1 -1
View File
@@ -9,7 +9,7 @@ from catalyst.exchange.exchange_data_portal import (
)
from catalyst.exchange.utils.exchange_utils import get_common_assets
from catalyst.exchange.utils.factory import get_exchanges
from .test_utils import rnd_history_date_days, rnd_bar_count
from test_utils import rnd_history_date_days, rnd_bar_count
log = Logger('test_bitfinex')
@@ -197,7 +197,6 @@ class TestSuiteBundle:
# population=exchange_population,
# features=[bundle],
# ) # Type: list[Exchange]
# TODO: currently focusing on Binance, try other exchanges
exchanges = [get_exchange('poloniex', skip_init=True)]
data_portal = TestSuiteBundle.get_data_portal(exchanges)
@@ -205,20 +204,17 @@ class TestSuiteBundle:
exchange.init()
frequencies = exchange.get_candle_frequencies(data_frequency)
# freq = random.sample(frequencies, 1)[0]
freq = '5T'
freq = random.sample(frequencies, 1)[0]
rnd = random.SystemRandom()
# field = rnd.choice(['open', 'high', 'low', 'close', 'volume'])
field = rnd.choice(['close'])
field = rnd.choice(['volume'])
# bar_count = random.randint(3, 6)
bar_count = 5
bar_count = random.randint(3, 6)
# assets = select_random_assets(
# exchange.assets, asset_population
# )
assets = [exchange.get_asset('bch_eth')]
end_dt = pd.to_datetime('2018-03-01', utc=True)
assets = select_random_assets(
exchange.assets, asset_population
)
end_dt = None
for asset in assets:
attribute = 'end_{}'.format(data_frequency)
asset_end_dt = getattr(asset, attribute)
@@ -5,28 +5,63 @@ from logging import Logger, WARNING
from time import sleep
import pandas as pd
from catalyst.assets._assets import TradingPair
from logbook import TestHandler
from catalyst.assets._assets import TradingPair
from catalyst.exchange.exchange_errors import ExchangeRequestError
from catalyst.exchange.exchange_execution import ExchangeLimitOrder
from catalyst.exchange.utils.exchange_utils import get_exchange_folder
from catalyst.exchange.utils.factory import get_exchanges, get_exchange
from catalyst.exchange.utils.test_utils import select_random_exchanges, \
select_random_assets
handle_exchange_error, select_random_assets
from catalyst.testing import ZiplineTestCase
from catalyst.testing.fixtures import WithLogger
from catalyst.exchange.utils.factory import get_exchanges, get_exchange
log = Logger('TestSuiteExchange')
class TestSuiteExchange(WithLogger, ZiplineTestCase):
def _test_markets_exchange(self, exchange, attempts=0):
assets = None
try:
exchange.init()
# Verify that the assets and markets are populated
if not exchange.markets:
raise ValueError(
'no markets found'
)
if not exchange.assets:
raise ValueError(
'no assets derived from markets'
)
assets = exchange.assets
except ExchangeRequestError as e:
sleep(5)
if attempts > 5:
handle_exchange_error(exchange, e)
else:
print(
're-trying an exchange request {} {}'.format(
exchange.name, attempts
)
)
self._test_markets_exchange(exchange, attempts + 1)
except Exception as e:
handle_exchange_error(exchange, e)
return assets
def test_markets(self):
population = 3
results = dict()
exchanges = select_random_exchanges(population) # Type: list[Exchange]
for exchange in exchanges:
exchange.init()
assets = self._test_markets_exchange(exchange)
if assets is not None: