Compare commits

..
23 Commits
Author SHA1 Message Date
AvishaiW 8a3cc7e5fa BLD: added a cmd for running on the cloud (WIP) 2018-03-01 01:30:18 +02:00
Frederic Fortier fec829b82e BLD: adjustments from testing 2018-02-27 23:10:09 -05:00
Frederic Fortier 25dc3ee737 DOC: fixed issue with exchange config script 2018-02-27 22:29:21 -05:00
Frederic Fortier 5de67a5a61 BLD: fixed some issues with the scripts 2018-02-27 20:00:04 -05:00
Frederic Fortier b2f042e2c2 Merge remote-tracking branch 'origin/new_exchange_config' into new_exchange_config
# Conflicts:
#	catalyst/examples/mean_reversion_simple.py
#	catalyst/exchange/ccxt/ccxt_exchange.py
#	catalyst/exchange/exchange.py
#	catalyst/exchange/utils/serialization_utils.py
#	etc/python2.7-environment.yml
#	etc/requirements.txt
#	tests/exchange/test_suites/test_suite_exchange.py
2018-02-27 00:29:40 -05:00
Frederic Fortier 704c93dac9 BLD: completed rebasing 2018-02-27 00:21:22 -05:00
Frederic Fortier 478579ed8c BUG: for issue #237, checking considering open orders when verifying the exchange balance for each positions 2018-02-27 00:18:43 -05:00
Frederic Fortier c65a976b81 working on trades collector 2018-02-27 00:18:37 -05:00
Frederic Fortier f9fa28c103 BLD: created a report with the start and end time of all collected candles 2018-02-27 00:17:10 -05:00
Frederic Fortier 14c5ef3006 BLD: using an iterable to yield exchanges instead of populating a list 2018-02-27 00:17:04 -05:00
Frederic Fortier de2d3f6f54 BLD: replacing symbols.json and fetching markets with a single config 2018-02-27 00:17:01 -05:00
Frederic Fortier b271b372d8 BLD: made some adjustment to generate and use the exchange config more efficiently and without mapping. Currently testing. 2018-02-27 00:11:43 -05:00
Frederic Fortier 3f9a0727c0 BLD: replacing symbols.json and fetching markets with a single config 2018-02-27 00:10:38 -05:00
Frederic Fortier 271a51a393 BLD: adjusting candle computation 2018-02-21 12:52:31 -05:00
Frederic Fortier 69153295f0 BUG: for issue #237, checking considering open orders when verifying the exchange balance for each positions 2018-02-20 15:47:08 -05:00
Frederic Fortier 1aed7c71f6 working on trades collector 2018-02-12 12:46:09 -05:00
Frederic Fortier 62d21f1aca BLD: updated CCXT 2018-01-29 23:17:09 -05:00
Frederic Fortier 48a89ad521 BLD: created a report with the start and end time of all collected candles 2018-01-29 23:14:37 -05:00
Frederic Fortier 2225c40b76 BLD: using an iterable to yield exchanges instead of populating a list 2018-01-29 22:14:44 -05:00
Frederic Fortier ad0bc5c41a Merge remote-tracking branch 'origin/new_exchange_config' into new_exchange_config
# Conflicts:
#	catalyst/exchange/ccxt/ccxt_exchange.py
#	catalyst/exchange/exchange.py
#	catalyst/exchange/utils/exchange_utils.py
#	catalyst/exchange/utils/serialization_utils.py
#	etc/requirements.txt
2018-01-27 20:59:15 -05:00
Frederic Fortier 2a239fd5bb BLD: made some adjustment to generate and use the exchange config more efficiently and without mapping. Currently testing. 2018-01-27 20:53:38 -05:00
Frederic Fortier 6b3f59ff76 BLD: replacing symbols.json and fetching markets with a single config 2018-01-26 17:16:42 -05:00
Frederic Fortier e872b1fc82 BLD: replacing symbols.json and fetching markets with a single config 2018-01-06 17:28:32 -05:00
33 changed files with 339 additions and 769 deletions
+189 -16
View File
@@ -14,6 +14,7 @@ from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.utils.exchange_utils import delete_algo_folder
from catalyst.utils.cli import Date, Timestamp
from catalyst.utils.run_algo import _run, load_extensions
from catalyst.utils.run_server import run_server
try:
__IPYTHON__
@@ -505,6 +506,179 @@ def live(ctx,
return perf
@main.command(name='serve-live')
@click.option(
'-f',
'--algofile',
default=None,
type=click.File('r'),
help='The file that contains the algorithm to run.',
)
@click.option(
'--capital-base',
type=float,
show_default=True,
help='The amount of capital (in base_currency) allocated to trading.',
)
@click.option(
'-t',
'--algotext',
help='The algorithm script to run.',
)
@click.option(
'-D',
'--define',
multiple=True,
help="Define a name to be bound in the namespace before executing"
" the algotext. For example '-Dname=value'. The value may be"
" any python expression. These are evaluated in order so they"
" may refer to previously defined names.",
)
@click.option(
'-o',
'--output',
default='-',
metavar='FILENAME',
show_default=True,
help="The location to write the perf data. If this is '-' the perf will"
" be written to stdout.",
)
@click.option(
'--print-algo/--no-print-algo',
is_flag=True,
default=False,
help='Print the algorithm to stdout.',
)
@ipython_only(click.option(
'--local-namespace/--no-local-namespace',
is_flag=True,
default=None,
help='Should the algorithm methods be resolved in the local namespace.'
))
@click.option(
'-x',
'--exchange-name',
help='The name of the targeted exchange.',
)
@click.option(
'-n',
'--algo-namespace',
help='A label assigned to the algorithm for data storage purposes.'
)
@click.option(
'-c',
'--base-currency',
help='The base currency used to calculate statistics '
'(e.g. usd, btc, eth).',
)
@click.option(
'-e',
'--end',
type=Date(tz='utc', as_timestamp=True),
help='An optional end date at which to stop the execution.',
)
@click.option(
'--live-graph/--no-live-graph',
is_flag=True,
default=False,
help='Display live graph.',
)
@click.option(
'--simulate-orders/--no-simulate-orders',
is_flag=True,
default=True,
help='Simulating orders enable the paper trading mode. No orders will be '
'sent to the exchange unless set to false.',
)
@click.option(
'--auth-aliases',
default=None,
help='Authentication file aliases for the specified exchanges. By default,'
'each exchange uses the "auth.json" file in the exchange folder. '
'Specifying an "auth2" alias would use "auth2.json". It should be '
'specified like this: "[exchange_name],[alias],..." For example, '
'"binance,auth2" or "binance,auth2,bittrex,auth2".',
)
@click.pass_context
def serve_live(ctx,
algofile,
capital_base,
algotext,
define,
output,
print_algo,
local_namespace,
exchange_name,
algo_namespace,
base_currency,
end,
live_graph,
auth_aliases,
simulate_orders):
"""Trade live with the given algorithm on the server.
"""
if (algotext is not None) == (algofile is not None):
ctx.fail(
"must specify exactly one of '-f' / '--algofile' or"
" '-t' / '--algotext'",
)
if exchange_name is None:
ctx.fail("must specify an exchange name '-x'")
if algo_namespace is None:
ctx.fail("must specify an algorithm name '-n' in live execution mode")
if base_currency is None:
ctx.fail("must specify a base currency '-c' in live execution mode")
if capital_base is None:
ctx.fail("must specify a capital base with '--capital-base'")
if simulate_orders:
click.echo('Running in paper trading mode.', sys.stdout)
else:
click.echo('Running in live trading mode.', sys.stdout)
perf = run_server(
initialize=None,
handle_data=None,
before_trading_start=None,
analyze=None,
algofile=algofile,
algotext=algotext,
defines=define,
data_frequency=None,
capital_base=capital_base,
data=None,
bundle=None,
bundle_timestamp=None,
start=None,
end=end,
output=output,
print_algo=print_algo,
local_namespace=local_namespace,
environ=os.environ,
live=True,
exchange=exchange_name,
algo_namespace=algo_namespace,
base_currency=base_currency,
live_graph=live_graph,
analyze_live=None,
simulate_orders=simulate_orders,
auth_aliases=auth_aliases,
stats_output=None,
)
if output == '-':
click.echo(str(perf), sys.stdout)
elif output != os.devnull: # make the catalyst magic not write any data
perf.to_pickle(output)
return perf
@main.command(name='ingest-exchange')
@click.option(
'-x',
@@ -767,18 +941,12 @@ def bundles():
@main.group()
@click.pass_context
def marketplace(ctx):
"""Access the Enigma Data Marketplace to:\n
- Register and Publish new datasets (seller-side)\n
- Subscribe and Ingest premium datasets (buyer-side)\n
"""
pass
@marketplace.command()
@click.pass_context
def ls(ctx):
"""List all available datasets.
"""
click.echo('Listing of available data sources on the marketplace:',
sys.stdout)
marketplace = Marketplace()
@@ -793,8 +961,10 @@ def ls(ctx):
)
@click.pass_context
def subscribe(ctx, dataset):
"""Subscribe to an existing dataset.
"""
if dataset is None:
ctx.fail("must specify a dataset to subscribe to with '--dataset'\n"
"List available dataset on the marketplace with "
"'catalyst marketplace ls'")
marketplace = Marketplace()
marketplace.subscribe(dataset)
@@ -829,8 +999,11 @@ def subscribe(ctx, dataset):
)
@click.pass_context
def ingest(ctx, dataset, data_frequency, start, end):
"""Ingest a dataset (requires subscription).
"""
if dataset is None:
ctx.fail("must specify a dataset to clean with '--dataset'\n"
"List available dataset on the marketplace with "
"'catalyst marketplace ls'")
click.echo('Ingesting data: {}'.format(dataset), sys.stdout)
marketplace = Marketplace()
marketplace.ingest(dataset, data_frequency, start, end)
@@ -843,17 +1016,19 @@ def ingest(ctx, dataset, data_frequency, start, end):
)
@click.pass_context
def clean(ctx, dataset):
"""Clean/Remove local data for a given dataset.
"""
if dataset is None:
ctx.fail("must specify a dataset to ingest with '--dataset'\n"
"List available dataset on the marketplace with "
"'catalyst marketplace ls'")
click.echo('Cleaning data source: {}'.format(dataset), sys.stdout)
marketplace = Marketplace()
marketplace.clean(dataset)
click.echo('Done', sys.stdout)
@marketplace.command()
@click.pass_context
def register(ctx):
"""Register a new dataset.
"""
marketplace = Marketplace()
marketplace.register()
@@ -877,8 +1052,6 @@ def register(ctx):
)
@click.pass_context
def publish(ctx, dataset, datadir, watch):
"""Publish data for a registered dataset.
"""
marketplace = Marketplace()
if dataset is None:
ctx.fail("must specify a dataset to publish data for "
+7 -6
View File
@@ -25,21 +25,22 @@ AUTO_INGEST = False
AUTH_SERVER = 'https://data.enigma.co'
# TODO: switch to mainnet
ETH_REMOTE_NODE = 'https://rinkeby.infura.io/'
ETH_REMOTE_NODE = 'https://ropsten.infura.io/'
# TODO: move to MASTER branch on github
MARKETPLACE_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/master/catalyst/marketplace/' \
'catalyst/develop/catalyst/marketplace/' \
'contract_marketplace_address.txt'
MARKETPLACE_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/master/catalyst/marketplace/' \
'catalyst/develop/catalyst/marketplace/' \
'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/' \
'develop/catalyst/marketplace/' \
'contract_enigma_address.txt'
ENIGMA_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/master/catalyst/marketplace/' \
'catalyst/develop/catalyst/marketplace/' \
'contract_enigma_abi.json'
+1
View File
@@ -7,6 +7,7 @@ from catalyst.api import (
order_target_percent,
symbol,
record,
get_open_orders,
)
from catalyst.exchange.utils.stats_utils import get_pretty_stats
from catalyst.utils.run_algo import run_algorithm
+1 -1
View File
@@ -66,7 +66,7 @@ def handle_data(context, data):
# Define portfolio optimization parameters
n_portfolios = 50000
results_array = np.zeros((3 + context.nassets, n_portfolios))
for p in range(n_portfolios):
for p in xrange(n_portfolios):
weights = np.random.random(context.nassets)
weights /= np.sum(weights)
w = np.asmatrix(weights)
+1 -1
View File
@@ -26,7 +26,7 @@ def handle_data(context, data):
context.asset,
fields='price',
bar_count=20,
frequency='30T'
frequency='2H'
)
last_traded = prices.index[-1]
log.info('last candle date: {}'.format(last_traded))
-3
View File
@@ -190,9 +190,6 @@ 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
+36 -50
View File
@@ -1,4 +1,5 @@
import abc
import pytz
from abc import ABCMeta, abstractmethod, abstractproperty
from datetime import timedelta
from time import sleep
@@ -11,15 +12,13 @@ from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import MismatchingBaseCurrencies, \
SymbolNotFoundOnExchange, \
PricingDataNotLoadedError, \
NoDataAvailableOnExchange, NoValueForField, \
NoCandlesReceivedFromExchange, \
NoDataAvailableOnExchange, NoValueForField, LastCandleTooEarlyError, \
TickerNotFoundError, NotEnoughCashError
from catalyst.exchange.utils.datetime_utils import get_delta, \
get_periods_range, \
get_periods, get_start_dt, get_frequency, \
get_candles_number_from_minutes
get_periods, get_start_dt, get_frequency
from catalyst.exchange.utils.exchange_utils import get_exchange_symbols, \
resample_history_df, has_bundle, get_candles_df
resample_history_df, has_bundle
from logbook import Logger
log = Logger('Exchange', level=LOG_LEVEL)
@@ -257,8 +256,7 @@ class Exchange:
elif data_frequency is not None:
applies = (
(
data_frequency == 'minute' and
a.end_minute is not None)
data_frequency == 'minute' and a.end_minute is not None)
or (
data_frequency == 'daily' and a.end_daily is not None)
)
@@ -507,60 +505,49 @@ class Exchange:
freq, candle_size, unit, data_frequency = get_frequency(
frequency, data_frequency, supported_freqs=['T', 'D', 'H']
)
# we want to avoid receiving empty candles
# so we request more than needed
# TODO: consider defining a const per asset
# and/or some retry mechanism (in each iteration request more data)
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,
assets=assets,
bar_count=requested_bar_count,
bar_count=bar_count,
end_dt=end_dt if not is_current else None,
)
# candles sanity check - verify no empty candles were received:
series = dict()
for asset in candles:
if not candles[asset]:
raise NoCandlesReceivedFromExchange(
bar_count=requested_bar_count,
if candles[asset]:
first_candle = candles[asset][0]
asset_series = self.get_series_from_candles(
candles=candles[asset],
start_dt=first_candle['last_traded'],
end_dt=end_dt,
asset=asset,
exchange=self.name)
data_frequency=frequency,
field=field,
)
# for avoiding unnecessary forward fill end_dt is taken back one second
forward_fill_till_dt = end_dt - timedelta(seconds=1)
delta_candle_size = candle_size * 60 if unit == 'H' else candle_size
# Checking to make sure that the dates match
delta = get_delta(delta_candle_size, data_frequency)
adj_end_dt = end_dt - delta
last_traded = asset_series.index[-1]
series = get_candles_df(candles=candles,
field=field,
freq=frequency,
bar_count=requested_bar_count,
end_dt=forward_fill_till_dt)
if last_traded < adj_end_dt:
raise LastCandleTooEarlyError(
last_traded=last_traded,
end_dt=adj_end_dt,
exchange=self.name,
)
else: # empty candle received
# because other assets are tz-aware, we need its tz to be set as well
asset_series = pd.Series([], index=pd.DatetimeIndex([], tz=pytz.utc))
# TODO: consider how to approach this edge case
# delta_candle_size = candle_size * 60 if unit == 'H' else candle_size
# Checking to make sure that the dates match
# delta = get_delta(delta_candle_size, data_frequency)
# adj_end_dt = end_dt - delta
# last_traded = asset_series.index[-1]
# if last_traded < adj_end_dt:
# raise LastCandleTooEarlyError(
# last_traded=last_traded,
# end_dt=adj_end_dt,
# exchange=self.name,
# )
series[asset] = asset_series
df = pd.DataFrame(series)
df.dropna(inplace=True)
#df.dropna(inplace=True) # commented out due to issue 236
return df.tail(bar_count)
return df
def get_history_window_with_bundle(self,
assets,
@@ -608,10 +595,9 @@ class Exchange:
A dataframe containing the requested data.
"""
# TODO: this function needs some work,
# we're currently using it just for benchmark data
# 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, supported_freqs=['T', 'D']
frequency, data_frequency
)
adj_bar_count = candle_size * bar_count
try:
@@ -635,7 +621,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
-19
View File
@@ -163,25 +163,6 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
style)
return amount, style
def _calculate_order_target_amount(self, asset, target):
"""
removes order amounts so we won't run into issues
when two orders are placed one after the other.
it then proceeds to removing positions amount at TradingAlgorithm
:param asset:
:param target:
:return: target
"""
if asset in self.blotter.open_orders:
for open_order in self.blotter.open_orders[asset]:
current_amount = open_order.amount
target -= current_amount
target = super(ExchangeTradingAlgorithmBase, self). \
_calculate_order_target_amount(asset, target)
return target
def round_order(self, amount, asset):
"""
We need fractions with cryptocurrencies
+1 -1
View File
@@ -68,7 +68,7 @@ class TradingPairFeeSchedule(CommissionModel):
multiplier = maker \
if ((order.amount > 0 and order.limit < transaction.price)
or (order.amount < 0 and order.limit > transaction.price)) \
and order.limit_reached else taker
and order.limit_reached else taker
fee = cost * multiplier
return fee
+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, supported_freqs=['T', 'D']
frequency, data_frequency
)
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, adj_data_frequency)
start_dt = get_start_dt(end_dt, adj_bar_count, data_frequency)
df = resample_history_df(pd.DataFrame(series), freq, field, start_dt)
return df
-7
View File
@@ -322,10 +322,3 @@ class BalanceTooLowError(ZiplineError):
'add positions to hold a free amount greater than {amount}, or clean '
'the state of this algo and restart.'
).strip()
class NoCandlesReceivedFromExchange(ZiplineError):
msg = (
'Although requesting {bar_count} candles until {end_dt} of asset {asset}, '
'an empty list of candles was received for {exchange}.'
).strip()
+4 -38
View File
@@ -1,5 +1,4 @@
import calendar
import math
import re
from datetime import datetime, timedelta, date
@@ -249,12 +248,9 @@ 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', 'H', 'T']):
def get_frequency(freq, data_frequency=None, supported_freqs=['D', 'T']):
"""
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.
Get the frequency parameters.
Notes
-----
@@ -309,14 +305,14 @@ def get_frequency(freq, data_frequency=None, supported_freqs=['D', 'H', '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)
@@ -330,33 +326,3 @@ 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))
+16 -27
View File
@@ -716,36 +716,25 @@ def save_asset_data(folder, df, decimals=8):
)
def forward_fill_df_if_needed(df, periods):
df = df.reindex(periods)
# volume should always be 0 (if there were no trades in this interval)
df['volume'] = df['volume'].fillna(0.0)
# ie pull the last close into this close
df['close'] = df.fillna(method='pad')
# now copy the close that was pulled down from the last timestep
# into this row, across into o/h/l
df['open'] = df['open'].fillna(df['close'])
df['low'] = df['low'].fillna(df['close'])
df['high'] = df['high'].fillna(df['close'])
return df
def transform_candles_to_df(candles):
return pd.DataFrame(candles).set_index('last_traded')
def get_candles_df(candles, field, freq, bar_count, end_dt):
def get_candles_df(candles, field, freq, bar_count, end_dt,
previous_value=None):
all_series = dict()
for asset in candles:
asset_df = transform_candles_to_df(candles[asset])
rounded_end_dt = end_dt.floor(freq)
periods = pd.date_range(end=rounded_end_dt,
periods=bar_count,
freq=freq)
asset_df = forward_fill_df_if_needed(asset_df, periods)
periods = pd.date_range(end=end_dt, periods=bar_count, freq=freq)
all_series[asset] = pd.Series(asset_df[field])
dates = [candle['last_traded'] for candle in candles[asset]]
values = [candle[field] for candle in candles[asset]]
series = pd.Series(values, index=dates)
"""
series = series.reindex(
periods,
method='ffill',
fill_value=previous_value,
)
series.sort_index(inplace=True)
"""
all_series[asset] = series
df = pd.DataFrame(all_series)
df.dropna(inplace=True)
@@ -1 +1 @@
0x39a54f480d922a58c963de8091a6c9afc69db2cf
0x7fAec9aaE31BE428DeAAE1be8195dF609079Fd10
File diff suppressed because one or more lines are too long
@@ -1 +1 @@
0xa2b37c6cd52f60fd4eb46ca59fafcf22d081aebc
0x3985f5de8fddf2e8f7705cd360b498bf35ebfbc4
+50 -137
View File
@@ -7,7 +7,6 @@ import re
import shutil
import sys
import time
import webbrowser
import bcolz
import logbook
@@ -24,7 +23,7 @@ from catalyst.exchange.utils.stats_utils import set_print_settings
from catalyst.marketplace.marketplace_errors import (
MarketplacePubAddressEmpty, MarketplaceDatasetNotFound,
MarketplaceNoAddressMatch, MarketplaceHTTPRequest,
MarketplaceNoCSVFiles, MarketplaceRequiresPython3)
MarketplaceNoCSVFiles)
from catalyst.marketplace.utils.auth_utils import get_key_secret, \
get_signed_headers
from catalyst.marketplace.utils.bundle_utils import merge_bundles
@@ -33,7 +32,6 @@ from catalyst.marketplace.utils.eth_utils import bin_hex, from_grains, \
from catalyst.marketplace.utils.path_utils import get_bundle_folder, \
get_data_source_folder, get_marketplace_folder, \
get_user_pubaddr, get_temp_bundles_folder, extract_bundle
from catalyst.utils.paths import ensure_directory
if sys.version_info.major < 3:
import urllib
@@ -46,10 +44,7 @@ log = logbook.Logger('Marketplace', level=LOG_LEVEL)
class Marketplace:
def __init__(self):
global Web3
try:
from web3 import Web3, HTTPProvider
except ImportError:
raise MarketplaceRequiresPython3()
from web3 import Web3, HTTPProvider
self.addresses = get_user_pubaddr()
@@ -65,8 +60,7 @@ class Marketplace:
contract_url = urllib.urlopen(MARKETPLACE_CONTRACT)
self.mkt_contract_address = Web3.toChecksumAddress(
contract_url.readline().decode(
contract_url.info().get_content_charset()).strip())
contract_url.readline().strip())
abi_url = urllib.urlopen(MARKETPLACE_CONTRACT_ABI)
abi = json.load(abi_url)
@@ -79,8 +73,7 @@ class Marketplace:
contract_url = urllib.urlopen(ENIGMA_CONTRACT)
self.eng_contract_address = Web3.toChecksumAddress(
contract_url.readline().decode(
contract_url.info().get_content_charset()).strip())
contract_url.readline().strip())
abi_url = urllib.urlopen(ENIGMA_CONTRACT_ABI)
abi = json.load(abi_url)
@@ -143,10 +136,10 @@ class Marketplace:
return address, address_i
def sign_transaction(self, tx):
def sign_transaction(self, from_address, tx):
url = 'https://www.myetherwallet.com/#offline-transaction'
print('\nVisit {url} and enter the following parameters:\n\n'
print('\nVisit https://www.myetherwallet.com/#offline-transaction and '
'enter the following parameters:\n\n'
'From Address:\t\t{_from}\n'
'\n\tClick the "Generate Information" button\n\n'
'To Address:\t\t{to}\n'
@@ -155,16 +148,13 @@ class Marketplace:
'Gas Price:\t\t[Accept the default value]\n'
'Nonce:\t\t\t{nonce}\n'
'Data:\t\t\t{data}\n'.format(
url=url,
_from=tx['from'],
to=tx['to'],
value=tx['value'],
gas=tx['gas'],
nonce=tx['nonce'],
data=tx['data'], )
)
webbrowser.open_new(url)
_from=from_address,
to=tx['to'],
value=tx['value'],
gas=tx['gas'],
nonce=tx['nonce'],
data=tx['data'], )
)
signed_tx = input('Copy and Paste the "Signed Transaction" '
'field here:\n')
@@ -177,17 +167,16 @@ class Marketplace:
def check_transaction(self, tx_hash):
if 'ropsten' in ETH_REMOTE_NODE:
etherscan = 'https://ropsten.etherscan.io/tx/'
elif 'rinkeby' in ETH_REMOTE_NODE:
etherscan = 'https://rinkeby.etherscan.io/tx/'
etherscan = 'https://ropsten.etherscan.io/tx/{}'.format(
tx_hash)
else:
etherscan = 'https://etherscan.io/tx/'
etherscan = '{}{}'.format(etherscan, tx_hash)
etherscan = 'https://etherscan.io/tx/{}'.format(tx_hash)
print('\nYou can check the outcome of your transaction here:\n'
'{}\n\n'.format(etherscan))
def _list(self):
def list(self):
data_sources = self.mkt_contract.functions.getAllProviders().call()
data = []
@@ -199,44 +188,15 @@ class Marketplace:
dataset=self.to_text(data_source)
)
)
return pd.DataFrame(data)
def list(self):
df = self._list()
df = pd.DataFrame(data)
set_print_settings()
if df.empty:
print('There are no datasets available yet.')
else:
print(df)
def subscribe(self, dataset=None):
if dataset is None:
df_sets = self._list()
if df_sets.empty:
print('There are no datasets available yet.')
return
set_print_settings()
while True:
print(df_sets)
dataset_num = input('Choose the dataset you want to '
'subscribe to [0..{}]: '.format(
df_sets.size - 1))
try:
dataset_num = int(dataset_num)
except ValueError:
print('Enter a number between 0 and {}'.format(
df_sets.size - 1))
else:
if dataset_num not in range(0, df_sets.size):
print('Enter a number between 0 and {}'.format(
df_sets.size - 1))
else:
dataset = df_sets.iloc[dataset_num]['dataset']
break
def subscribe(self, dataset):
dataset = dataset.lower()
@@ -299,14 +259,14 @@ class Marketplace:
'buy: {} ENG. Get enough ENG to cover the costs of the '
'monthly\nsubscription for what you are trying to buy, '
'and try again.'.format(
address, from_grains(balance), price))
address, from_grains(balance), price))
return
while True:
agree_pay = input('Please confirm that you agree to pay {} ENG '
'for a monthly subscription to the dataset "{}" '
'starting today. [default: Y] '.format(
price, dataset)) or 'y'
price, dataset)) or 'y'
if agree_pay.lower() not in ('y', 'n'):
print("Please answer Y or N.")
else:
@@ -327,11 +287,13 @@ class Marketplace:
self.mkt_contract_address,
grains,
).buildTransaction(
{'from': address,
'nonce': self.web3.eth.getTransactionCount(address)}
{'nonce': self.web3.eth.getTransactionCount(address)}
)
signed_tx = self.sign_transaction(tx)
if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(address, tx)
try:
tx_hash = '0x{}'.format(
bin_hex(self.web3.eth.sendRawTransaction(signed_tx))
@@ -366,11 +328,13 @@ class Marketplace:
tx = self.mkt_contract.functions.subscribe(
Web3.toHex(dataset),
).buildTransaction({
'from': address,
'nonce': self.web3.eth.getTransactionCount(address)})
).buildTransaction(
{'nonce': self.web3.eth.getTransactionCount(address)})
signed_tx = self.sign_transaction(tx)
if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(address, tx)
try:
tx_hash = '0x{}'.format(bin_hex(
@@ -405,7 +369,7 @@ class Marketplace:
'You can now ingest this dataset anytime during the '
'next month by running the following command:\n'
'catalyst marketplace ingest --dataset={}'.format(
dataset, address, dataset))
dataset, address, dataset))
def process_temp_bundle(self, ds_name, path):
"""
@@ -423,7 +387,6 @@ class Marketplace:
"""
tmp_bundle = extract_bundle(path)
bundle_folder = get_data_source_folder(ds_name)
ensure_directory(bundle_folder)
if os.listdir(bundle_folder):
zsource = bcolz.ctable(rootdir=tmp_bundle, mode='r')
ztarget = bcolz.ctable(rootdir=bundle_folder, mode='r')
@@ -434,33 +397,7 @@ class Marketplace:
pass
def ingest(self, ds_name=None, start=None, end=None, force_download=False):
if ds_name is None:
df_sets = self._list()
if df_sets.empty:
print('There are no datasets available yet.')
return
set_print_settings()
while True:
print(df_sets)
dataset_num = input('Choose the dataset you want to '
'ingest [0..{}]: '.format(
df_sets.size - 1))
try:
dataset_num = int(dataset_num)
except ValueError:
print('Enter a number between 0 and {}'.format(
df_sets.size - 1))
else:
if dataset_num not in range(0, df_sets.size):
print('Enter a number between 0 and {}'.format(
df_sets.size - 1))
else:
ds_name = df_sets.iloc[dataset_num]['dataset']
break
def ingest(self, ds_name, start=None, end=None, force_download=False):
# ds_name = ds_name.lower()
@@ -489,10 +426,10 @@ class Marketplace:
print('Your subscription to dataset "{}" expired on {} UTC.'
'Please renew your subscription by running:\n'
'catalyst marketplace subscribe --dataset={}'.format(
ds_name,
pd.to_datetime(check_sub[4], unit='s', utc=True),
ds_name)
)
ds_name,
pd.to_datetime(check_sub[4], unit='s', utc=True),
ds_name)
)
if 'key' in self.addresses[address_i]:
key = self.addresses[address_i]['key']
@@ -556,40 +493,14 @@ class Marketplace:
return df
def clean(self, ds_name=None, data_frequency=None):
if ds_name is None:
mktplace_root = get_marketplace_folder()
folders = [os.path.basename(f.rstrip('/'))
for f in glob.glob('{}/*/'.format(mktplace_root))
if 'temp_bundles' not in f]
while True:
for idx, f in enumerate(folders):
print('{}\t{}'.format(idx, f))
dataset_num = input('Choose the dataset you want to '
'clean [0..{}]: '.format(
len(folders) - 1))
try:
dataset_num = int(dataset_num)
except ValueError:
print('Enter a number between 0 and {}'.format(
len(folders) - 1))
else:
if dataset_num not in range(0, len(folders)):
print('Enter a number between 0 and {}'.format(
len(folders) - 1))
else:
ds_name = folders[dataset_num]
break
ds_name = ds_name.lower()
def clean(self, data_source_name, data_frequency=None):
data_source_name = data_source_name.lower()
if data_frequency is None:
folder = get_data_source_folder(ds_name)
folder = get_data_source_folder(data_source_name)
else:
folder = get_bundle_folder(ds_name, data_frequency)
folder = get_bundle_folder(data_source_name, data_frequency)
shutil.rmtree(folder)
pass
@@ -693,11 +604,13 @@ class Marketplace:
grains,
address,
).buildTransaction(
{'from': address,
'nonce': self.web3.eth.getTransactionCount(address)}
{'nonce': self.web3.eth.getTransactionCount(address)}
)
signed_tx = self.sign_transaction(tx)
if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(address, tx)
try:
tx_hash = '0x{}'.format(
@@ -708,7 +621,7 @@ class Marketplace:
)
except Exception as e:
print('Unable to register the requested dataset: {}'.format(e))
print('Unable to subscribe to data source: {}'.format(e))
return
self.check_transaction(tx_hash)
+1 -10
View File
@@ -9,8 +9,7 @@ def silent_except_hook(exctype, excvalue, exctraceback):
MarketplaceNoAddressMatch, MarketplaceHTTPRequest,
MarketplaceNoCSVFiles, MarketplaceContractDataNoMatch,
MarketplaceSubscriptionExpired, MarketplaceJSONError,
MarketplaceWalletNotSupported, MarketplaceEmptySignature,
MarketplaceRequiresPython3]:
MarketplaceWalletNotSupported, MarketplaceEmptySignature]:
fn = traceback.extract_tb(exctraceback)[-1][0]
ln = traceback.extract_tb(exctraceback)[-1][1]
print("Error traceback: {1} (line {2})\n"
@@ -87,11 +86,3 @@ class MarketplaceJSONError(ZiplineError):
'The configuration file {file} is malformed. Please correct '
'the following error:\n{error}'
)
class MarketplaceRequiresPython3(ZiplineError):
msg = (
'\nCatalyst requires Python3 to access the Enigma Data Marketplace.\n'
'If you want to use the Data Marketplace, you need to reinstall '
'Catalyst\nwith Python3. See the documentation website for additional '
'information.')
+2 -10
View File
@@ -1,6 +1,5 @@
import hashlib
import hmac
import webbrowser
import requests
import time
@@ -46,17 +45,10 @@ def get_key_secret(pubAddr, wallet='mew'):
nonce = '0x{}'.format(d['nonce'])
if wallet == 'mew':
url = 'https://www.myetherwallet.com/signmsg.html'
print('\nObtaining a key/secret pair to streamline all future '
'requests with the authentication server.\n'
'Visit {url} and sign the '
'following message:\n{nonce}'.format(
url=url,
nonce=nonce))
webbrowser.open_new(url)
'Visit https://www.myetherwallet.com/signmsg.html and sign the '
'following message:\n{}'.format(nonce))
signature = input('Copy and Paste the "sig" field from '
'the signature here (without the double quotes, '
'only the HEX value):\n')
+1 -59
View File
@@ -1,12 +1,8 @@
import os
import random
import re
import shutil
import bcolz
import numpy as np
import pandas as pd
from six import string_types
def merge_bundles(zsource, ztarget):
@@ -31,64 +27,10 @@ def merge_bundles(zsource, ztarget):
df.drop_duplicates(inplace=True)
df.set_index(['date', 'symbol'], drop=False, inplace=True)
sanitize_df(df)
dirname = os.path.basename(ztarget.rootdir)
bak_dir = ztarget.rootdir.replace(dirname, '.{}'.format(dirname))
shutil.move(ztarget.rootdir, bak_dir)
os.rename(ztarget.rootdir, bak_dir)
z = bcolz.ctable.fromdataframe(df=df, rootdir=ztarget.rootdir)
shutil.rmtree(bak_dir)
return z
def sanitize_df(df):
# Using a sampling method to identify dates for efficiency with
# large datasets
if len(df) > 100:
indexes = random.sample(range(0, len(df) - 1), 100)
elif len(df) > 1:
indexes = range(0, len(df) - 1)
else:
indexes = [0, ]
for column in df.columns:
is_date = False
for index in indexes:
value = df[column].iloc[index]
if not isinstance(value, string_types):
continue
# TODO: assuming that the date is at least daily
exp = re.compile(r'^\d{4}-\d{2}-\d{2}.*$')
matches = exp.findall(value)
if matches:
is_date = True
break
if is_date:
df[column] = pd.to_datetime(df[column])
else:
try:
ser = safely_reduce_dtype(df[column])
df[column] = ser
except Exception:
pass
return df
def safely_reduce_dtype(ser): # pandas.Series or numpy.array
orig_dtype = "".join(
[x for x in ser.dtype.name if x.isalpha()]) # float/int
mx = 1
for val in ser.values:
new_itemsize = np.min_scalar_type(val).itemsize
if mx < new_itemsize:
mx = new_itemsize
if orig_dtype == 'int':
mx = max(mx, 4)
new_dtype = orig_dtype + str(mx * 8)
return ser.astype(new_dtype)
-49
View File
@@ -1,49 +0,0 @@
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)
)
-32
View File
@@ -1,32 +0,0 @@
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']
symbols = None
def initialize(context):
pass
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]
history = data.history(symbols,
['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)
-35
View File
@@ -1,35 +0,0 @@
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))
+1 -12
View File
@@ -10,7 +10,6 @@ import click
import pandas as pd
from six import string_types
import catalyst
from catalyst.data.bundles import load
from catalyst.data.data_portal import DataPortal
from catalyst.exchange.exchange_pricing_loader import ExchangePricingLoader, \
@@ -24,7 +23,7 @@ try:
from pygments.formatters import TerminalFormatter
PYGMENTS = True
except ImportError:
except:
PYGMENTS = False
from toolz import valfilter, concatv
from functools import partial
@@ -152,7 +151,6 @@ def _run(handle_data,
'We encourage you to report any issue on GitHub: '
'https://github.com/enigmampc/catalyst/issues'
)
log.info('Catalyst version {}'.format(catalyst.__version__))
sleep(3)
if live:
@@ -263,15 +261,6 @@ def _run(handle_data,
# We still need to support bundles for other misc data, but we
# can handle this later.
if start != pd.tslib.normalize_date(start) or \
end != pd.tslib.normalize_date(end):
# todo: add to Sim_Params the option to start & end at specific times
log.warn(
"Catalyst currently starts and ends on the start and "
"end of the dates specified, respectively. We hope to "
"Modify this and support specific times in a future release."
)
data = DataPortalExchangeBacktest(
exchange_names=[exchange_name for exchange_name in exchanges],
asset_finder=None,
+5 -14
View File
@@ -132,19 +132,20 @@ with the following steps:
conda env remove --name catalyst
2. Create the environment:
for python 2.7:
for python 2.7:
.. code-block:: bash
conda create --name catalyst python=2.7 scipy zlib
or for python 3.6:
.. code-block:: bash
conda create --name catalyst python=3.6 scipy zlib
3. Activate the environment:
**Linux or MacOS:**
@@ -314,16 +315,6 @@ 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``
-46
View File
@@ -2,52 +2,6 @@
Release Notes
=============
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
+3 -4
View File
@@ -5,6 +5,7 @@ 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,9 +23,7 @@ dependencies:
- bottleneck==1.2.1
- chardet==3.0.4
- 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
- web3==4.0.0b7
- requests-toolbelt==0.8.0
- click==6.7
- contextlib2==0.5.5
@@ -39,7 +38,7 @@ dependencies:
- lru-dict==1.1.6
- mako==1.0.7
- markupsafe==1.0
- matplotlib==2.1.2
- matplotlib==2.1.0
- multipledispatch==0.4.9
- networkx==2.0
- numexpr==2.6.4
+1 -1
View File
@@ -84,5 +84,5 @@ tables==3.3.0
ccxt==1.10.1094
boto3==1.4.8
redo==1.6
web3==4.0.0b11; python_version > '3.4'
web3==4.0.0b7
requests-toolbelt==0.8.0
+2
View File
@@ -0,0 +1,2 @@
web3==4.0.0b7
requests-toolbelt==0.8.0
+5 -5
View File
@@ -11,7 +11,7 @@ from catalyst.exchange.exchange_bundle import ExchangeBundle, \
BUNDLE_NAME_TEMPLATE
from catalyst.exchange.utils.bundle_utils import get_bcolz_chunk, \
get_df_from_arrays
from catalyst.exchange.utils.datetime_utils import get_start_dt
from 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
@@ -42,7 +42,7 @@ class TestExchangeBundle:
def test_ingest_minute(self):
data_frequency = 'minute'
exchange_name = 'binance'
exchange_name = 'poloniex'
exchange = get_exchange(exchange_name)
exchange_bundle = ExchangeBundle(exchange)
@@ -50,8 +50,8 @@ class TestExchangeBundle:
exchange.get_asset('eth_btc')
]
start = pd.to_datetime('2018-03-01', utc=True)
end = pd.to_datetime('2018-03-8', utc=True)
start = pd.to_datetime('2016-03-01', utc=True)
end = pd.to_datetime('2017-11-1', utc=True)
log.info('ingesting exchange bundle {}'.format(exchange_name))
exchange_bundle.ingest(
@@ -101,7 +101,7 @@ class TestExchangeBundle:
# data_frequency = 'daily'
# include_symbols = 'neo_btc,bch_btc,eth_btc'
exchange_name = 'binance'
exchange_name = 'bitfinex'
data_frequency = 'minute'
exchange = get_exchange(exchange_name)
-175
View File
@@ -1,175 +0,0 @@
from catalyst.exchange.utils.exchange_utils import transform_candles_to_df, \
forward_fill_df_if_needed, get_candles_df
from catalyst.testing.fixtures import WithLogger, ZiplineTestCase
from datetime import timedelta
from pandas import Timestamp, DataFrame, concat
import numpy as np
class TestExchangeUtils(WithLogger, ZiplineTestCase):
@classmethod
def get_specific_field_from_df(cls, df, field, asset):
new_df = DataFrame(df[field])
new_df.columns = [asset]
new_df.index.name = None
return new_df
@classmethod
def verify_forward_fill_df_if_needed(cls, candles, periods, expected_df):
observed_df = forward_fill_df_if_needed(
transform_candles_to_df(candles),
periods)
assert (expected_df.equals(observed_df))
@classmethod
def verify_get_candles_df(cls, assets, candles, end_fixed_dt,
expected_df, check_next_candle=False):
# run on all the fields
for field in ['volume', 'open', 'close', 'high', 'low']:
field_dt = cls.get_specific_field_from_df(expected_df,
field,
assets[0])
# run on several timestamps
for delta in range(5):
end_dt = end_fixed_dt + timedelta(minutes=delta)
assert (field_dt.equals(get_candles_df({assets[0]: candles},
field, '5T', 3,
end_dt=end_dt)))
field_dt_a1 = cls.get_specific_field_from_df(expected_df,
field,
assets[0])
field_dt_a2 = cls.get_specific_field_from_df(expected_df,
field,
assets[1])
observed_df = get_candles_df({assets[0]: candles,
assets[1]: candles},
field, '5T', 3,
end_dt=end_dt)
assert (observed_df.equals(concat([field_dt_a1, field_dt_a2],
axis=1)))
if check_next_candle:
# one candle forward
end_dt = end_fixed_dt + timedelta(minutes=6)
observed_df = get_candles_df({assets[0]: candles,
assets[1]: candles},
field, '5T', 3,
end_dt=end_dt)
assert (not observed_df.equals(concat([field_dt_a1,
field_dt_a2],
axis=1)))
assert (concat([field_dt_a1, field_dt_a2],
axis=1)[1:].equals(observed_df[:-1]))
def test_get_candles_df(self):
assets = ['btc_usdt', 'eth_usdt']
# test forward fill in the end
candles = [{'high': 595, 'volume': 10, 'low': 594,
'close': 595, 'open': 594,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 594, 'volume': 108, 'low': 592,
'close': 593, 'open': 592,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
}]
expected = [{'high': 595.0, 'volume': 10.0, 'low': 594.0,
'close': 595.0, 'open': 594.0,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 594.0, 'volume': 108.0, 'low': 592.0,
'close': 593.0, 'open': 592.0,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
},
{'high': 593.0, 'volume': 0.0, 'low': 593.0,
'close': 593.0, 'open': 593.0,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
periods = [Timestamp('2018-03-01 09:45:00+0000', tz='UTC'),
Timestamp('2018-03-01 09:50:00+0000', tz='UTC'),
Timestamp('2018-03-01 09:55:00+0000', tz='UTC')]
expected_df = transform_candles_to_df(expected)
self.verify_forward_fill_df_if_needed(candles, periods,
expected_df)
self.verify_get_candles_df(assets, candles, periods[2],
expected_df, True)
# test forward fill in the middle
candles = [{'high': 595, 'volume': 10, 'low': 594,
'close': 595, 'open': 594,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 594, 'volume': 108, 'low': 592,
'close': 593, 'open': 592,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
expected = [{'high': 595.0, 'volume': 10.0, 'low': 594.0,
'close': 595.0, 'open': 594.0,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 595.0, 'volume': 0.0, 'low': 595.0,
'close': 595.0, 'open': 595.0,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
},
{'high': 594.0, 'volume': 108.0, 'low': 592.0,
'close': 593.0, 'open': 592.0,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
expected_df = transform_candles_to_df(expected)
self.verify_forward_fill_df_if_needed(candles, periods, expected_df)
self.verify_get_candles_df(assets, candles, periods[2], expected_df)
# test "forward fill" at the beginning
candles = [{'high': 595, 'volume': 10, 'low': 594,
'close': 595, 'open': 594,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
},
{'high': 594, 'volume': 108, 'low': 592,
'close': 593, 'open': 592,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
expected = [{'high': np.NaN, 'volume': 0.0, 'low': np.NaN,
'close': np.NaN, 'open': np.NaN,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 595, 'volume': 10, 'low': 594,
'close': 595, 'open': 594,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
},
{'high': 594, 'volume': 108, 'low': 592,
'close': 593, 'open': 592,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
expected_df = transform_candles_to_df(expected)
self.verify_forward_fill_df_if_needed(candles, periods, expected_df)
# Not the same due to dropna - commenting out for now
# self.verify_get_candles_df(assets, candles, periods[2], expected_df)
@@ -107,14 +107,14 @@ class TestSuiteBundle:
print('saved {} test results: {}'.format(end_dt, folder))
assert_frame_equal(
right=data['bundle'][:-1],
left=data['exchange'][:-1],
right=data['bundle'],
left=data['exchange'],
check_less_precise=1,
)
try:
assert_frame_equal(
right=data['bundle'][:-1],
left=data['exchange'][:-1],
right=data['bundle'],
left=data['exchange'],
check_less_precise=min([a.decimals for a in assets]),
)
except Exception as e:
+3 -2
View File
@@ -1,5 +1,6 @@
from catalyst.marketplace.marketplace import Marketplace
from catalyst.testing.fixtures import WithLogger, ZiplineTestCase
import pandas as pd
class TestMarketplace(WithLogger, ZiplineTestCase):
@@ -15,12 +16,12 @@ class TestMarketplace(WithLogger, ZiplineTestCase):
def test_subscribe(self):
marketplace = Marketplace()
marketplace.subscribe('marketcap')
marketplace.subscribe('marketcap2222')
pass
def test_ingest(self):
marketplace = Marketplace()
ds_def = marketplace.ingest('marketcap')
ds_def = marketplace.ingest('github')
pass
def test_publish(self):