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
40 changed files with 439 additions and 1001 deletions
+2 -2
View File
@@ -1,11 +1,11 @@
# #
# Dockerfile for an image with the currently checked out version of catalyst installed. To build: # Dockerfile for an image with the currently checked out version of catalyst installed. To build:
# #
# docker build -t enigmampc/catalyst . # docker build -t quantopian/catalyst .
# #
# To run the container: # To run the container:
# #
# docker run -v /path/to/your/notebooks:/projects -v ~/.catalyst:/root/.catalyst -p 8888:8888/tcp --name catalyst -it enigmampc/catalyst # docker run -v /path/to/your/notebooks:/projects -v ~/.catalyst:/root/.catalyst -p 8888:8888/tcp --name catalyst -it quantopian/catalyst
# #
# To access Jupyter when running docker locally (you may need to add NAT rules): # To access Jupyter when running docker locally (you may need to add NAT rules):
# #
+5 -5
View File
@@ -1,15 +1,15 @@
# #
# Dockerfile for an image with the currently checked out version of catalyst installed. To build: # Dockerfile for an image with the currently checked out version of catalyst installed. To build:
# #
# docker build -t enigmampc/catalystdev -f Dockerfile-dev . # docker build -t quantopian/catalystdev -f Dockerfile-dev .
# #
# Note: the dev build requires a enigmampc/catalyst image, which you can build as follows: # Note: the dev build requires a quantopian/catalyst image, which you can build as follows:
# #
# docker build -t enigmampc/catalyst -f Dockerfile . # docker build -t quantopian/catalyst -f Dockerfile .
# #
# To run the container: # To run the container:
# #
# docker run -v /path/to/your/notebooks:/projects -v ~/.catalyst:/root/.catalyst -p 8888:8888/tcp --name catalystdev -it enigmampc/catalystdev # docker run -v /path/to/your/notebooks:/projects -v ~/.catalyst:/root/.catalyst -p 8888:8888/tcp --name catalystdev -it quantopian/catalystdev
# #
# To access Jupyter when running docker locally (you may need to add NAT rules): # To access Jupyter when running docker locally (you may need to add NAT rules):
# #
@@ -25,7 +25,7 @@
# #
# docker exec -it catalystdev catalyst run -f /projects/my_algo.py --start 2015-1-1 --end 2016-1-1 /projects/result.pickle # docker exec -it catalystdev catalyst run -f /projects/my_algo.py --start 2015-1-1 --end 2016-1-1 /projects/result.pickle
# #
FROM enigmampc/catalyst FROM quantopian/catalyst
WORKDIR /catalyst WORKDIR /catalyst
+3 -9
View File
@@ -5,7 +5,6 @@
|version tag| |version tag|
|version status| |version status|
|forum|
|discord| |discord|
|twitter| |twitter|
@@ -23,11 +22,9 @@ visit `enigma.co <https://www.enigma.co>`_ to learn more about Catalyst.
Catalyst builds on top of the well-established Catalyst builds on top of the well-established
`Zipline <https://github.com/quantopian/zipline>`_ project. We did our best to `Zipline <https://github.com/quantopian/zipline>`_ project. We did our best to
minimize structural changes to the general API to maximize compatibility with minimize structural changes to the general API to maximize compatibility with
existing trading algorithms, developer knowledge, and tutorials. Join us on the existing trading algorithms, developer knowledge, and tutorials. Join us on
`Catalyst Forum <https://catalyst.enigma.co/>`_ for questions around Catalyst, `Discord <https://discord.gg/SJK32GY>`_ where we have a *#catalyst_dev* channel
algorithmic trading and technical support. We also have a for questions around Catalyst, algorithmic trading and technical support.
`Discord <https://discord.gg/SJK32GY>`_ group with the *#catalyst_dev* and
*#catalyst_setup* dedicated channels.
Overview Overview
======== ========
@@ -64,9 +61,6 @@ Go to our `Documentation Website <https://enigmampc.github.io/catalyst/>`_.
.. |version status| image:: https://img.shields.io/pypi/pyversions/enigma-catalyst.svg .. |version status| image:: https://img.shields.io/pypi/pyversions/enigma-catalyst.svg
:target: https://pypi.python.org/pypi/enigma-catalyst :target: https://pypi.python.org/pypi/enigma-catalyst
.. |forum| image:: https://img.shields.io/badge/forum-join-green.svg
:target: https://catalyst.enigma.co/
.. |discord| image:: https://img.shields.io/badge/discord-join%20chat-green.svg .. |discord| image:: https://img.shields.io/badge/discord-join%20chat-green.svg
:target: https://discordapp.com/invite/SJK32GY :target: https://discordapp.com/invite/SJK32GY
+190 -17
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.exchange.utils.exchange_utils import delete_algo_folder
from catalyst.utils.cli import Date, Timestamp from catalyst.utils.cli import Date, Timestamp
from catalyst.utils.run_algo import _run, load_extensions from catalyst.utils.run_algo import _run, load_extensions
from catalyst.utils.run_server import run_server
try: try:
__IPYTHON__ __IPYTHON__
@@ -505,6 +506,179 @@ def live(ctx,
return perf 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') @main.command(name='ingest-exchange')
@click.option( @click.option(
'-x', '-x',
@@ -580,7 +754,7 @@ def ingest_exchange(ctx, exchange_name, data_frequency, start, end,
exchange_bundle = ExchangeBundle(exchange_name) exchange_bundle = ExchangeBundle(exchange_name)
click.echo('Trying to ingest exchange bundle {}...'.format(exchange_name), click.echo('Ingesting exchange bundle {}...'.format(exchange_name),
sys.stdout) sys.stdout)
exchange_bundle.ingest( exchange_bundle.ingest(
data_frequency=data_frequency, data_frequency=data_frequency,
@@ -767,18 +941,12 @@ def bundles():
@main.group() @main.group()
@click.pass_context @click.pass_context
def marketplace(ctx): 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 pass
@marketplace.command() @marketplace.command()
@click.pass_context @click.pass_context
def ls(ctx): def ls(ctx):
"""List all available datasets.
"""
click.echo('Listing of available data sources on the marketplace:', click.echo('Listing of available data sources on the marketplace:',
sys.stdout) sys.stdout)
marketplace = Marketplace() marketplace = Marketplace()
@@ -793,8 +961,10 @@ def ls(ctx):
) )
@click.pass_context @click.pass_context
def subscribe(ctx, dataset): 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 = Marketplace()
marketplace.subscribe(dataset) marketplace.subscribe(dataset)
@@ -829,8 +999,11 @@ def subscribe(ctx, dataset):
) )
@click.pass_context @click.pass_context
def ingest(ctx, dataset, data_frequency, start, end): 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 = Marketplace()
marketplace.ingest(dataset, data_frequency, start, end) marketplace.ingest(dataset, data_frequency, start, end)
@@ -843,17 +1016,19 @@ def ingest(ctx, dataset, data_frequency, start, end):
) )
@click.pass_context @click.pass_context
def clean(ctx, dataset): 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 = Marketplace()
marketplace.clean(dataset) marketplace.clean(dataset)
click.echo('Done', sys.stdout)
@marketplace.command() @marketplace.command()
@click.pass_context @click.pass_context
def register(ctx): def register(ctx):
"""Register a new dataset.
"""
marketplace = Marketplace() marketplace = Marketplace()
marketplace.register() marketplace.register()
@@ -877,8 +1052,6 @@ def register(ctx):
) )
@click.pass_context @click.pass_context
def publish(ctx, dataset, datadir, watch): def publish(ctx, dataset, datadir, watch):
"""Publish data for a registered dataset.
"""
marketplace = Marketplace() marketplace = Marketplace()
if dataset is None: if dataset is None:
ctx.fail("must specify a dataset to publish data for " ctx.fail("must specify a dataset to publish data for "
+3 -8
View File
@@ -630,28 +630,23 @@ cdef class TradingPair(Asset):
and whose second element is a tuple of all the attributes that should and whose second element is a tuple of all the attributes that should
be serialized/deserialized during pickling. be serialized/deserialized during pickling.
""" """
# added arguments for catalyst #TODO: make sure that all fields set there
return (self.__class__, (self.symbol, return (self.__class__, (self.symbol,
self.exchange, self.exchange,
self.start_date, self.start_date,
self.asset_name, self.asset_name,
self.sid, self.sid,
self.leverage, self.leverage,
self.end_daily,
self.end_minute,
self.end_date, self.end_date,
self.exchange_symbol,
self.first_traded, self.first_traded,
self.auto_close_date, self.auto_close_date,
self.exchange_full, self.exchange_full,
self.min_trade_size, self.min_trade_size,
self.max_trade_size, self.max_trade_size,
self.maker,
self.taker,
self.lot, self.lot,
self.decimals, self.decimals,
self.trading_state, self.taker,
self.data_source)) self.maker))
def make_asset_array(int size, Asset asset): def make_asset_array(int size, Asset asset):
cdef np.ndarray out = np.empty([size], dtype=object) cdef np.ndarray out = np.empty([size], dtype=object)
+7 -9
View File
@@ -25,24 +25,22 @@ AUTO_INGEST = False
AUTH_SERVER = 'https://data.enigma.co' AUTH_SERVER = 'https://data.enigma.co'
# TODO: switch to mainnet # 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/' \ MARKETPLACE_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/master/catalyst/marketplace/' \ 'catalyst/develop/catalyst/marketplace/' \
'contract_marketplace_address.txt' 'contract_marketplace_address.txt'
MARKETPLACE_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \ MARKETPLACE_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/master/catalyst/marketplace/' \ 'catalyst/develop/catalyst/marketplace/' \
'contract_marketplace_abi.json' 'contract_marketplace_abi.json'
# TODO: switch to mainnet # TODO: switch to mainnet
ENIGMA_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/' \ ENIGMA_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/catalyst/' \
'catalyst/master/catalyst/marketplace/' \ 'develop/catalyst/marketplace/' \
'contract_enigma_address.txt' 'contract_enigma_address.txt'
ENIGMA_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \ ENIGMA_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/master/catalyst/marketplace/' \ 'catalyst/develop/catalyst/marketplace/' \
'contract_enigma_abi.json' 'contract_enigma_abi.json'
SUPPORTED_WALLETS = ['metamask', 'ledger', 'trezor', 'bitbox', 'keystore',
'key']
+1
View File
@@ -7,6 +7,7 @@ from catalyst.api import (
order_target_percent, order_target_percent,
symbol, symbol,
record, record,
get_open_orders,
) )
from catalyst.exchange.utils.stats_utils import get_pretty_stats from catalyst.exchange.utils.stats_utils import get_pretty_stats
from catalyst.utils.run_algo import run_algorithm 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 # Define portfolio optimization parameters
n_portfolios = 50000 n_portfolios = 50000
results_array = np.zeros((3 + context.nassets, n_portfolios)) 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.random.random(context.nassets)
weights /= np.sum(weights) weights /= np.sum(weights)
w = np.asmatrix(weights) w = np.asmatrix(weights)
+1 -1
View File
@@ -26,7 +26,7 @@ def handle_data(context, data):
context.asset, context.asset,
fields='price', fields='price',
bar_count=20, bar_count=20,
frequency='30T' frequency='2H'
) )
last_traded = prices.index[-1] last_traded = prices.index[-1]
log.info('last candle date: {}'.format(last_traded)) 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'): if data_frequency == 'minute' and not freq.endswith('T'):
continue continue
elif data_frequency == 'hourly' and not freq.endswith('D'):
continue
elif data_frequency == 'daily' and not freq.endswith('D'): elif data_frequency == 'daily' and not freq.endswith('D'):
continue continue
+42 -52
View File
@@ -1,4 +1,5 @@
import abc import abc
import pytz
from abc import ABCMeta, abstractmethod, abstractproperty from abc import ABCMeta, abstractmethod, abstractproperty
from datetime import timedelta from datetime import timedelta
from time import sleep from time import sleep
@@ -11,15 +12,13 @@ from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import MismatchingBaseCurrencies, \ from catalyst.exchange.exchange_errors import MismatchingBaseCurrencies, \
SymbolNotFoundOnExchange, \ SymbolNotFoundOnExchange, \
PricingDataNotLoadedError, \ PricingDataNotLoadedError, \
NoDataAvailableOnExchange, NoValueForField, \ NoDataAvailableOnExchange, NoValueForField, LastCandleTooEarlyError, \
NoCandlesReceivedFromExchange, \
TickerNotFoundError, NotEnoughCashError TickerNotFoundError, NotEnoughCashError
from catalyst.exchange.utils.datetime_utils import get_delta, \ from catalyst.exchange.utils.datetime_utils import get_delta, \
get_periods_range, \ get_periods_range, \
get_periods, get_start_dt, get_frequency, \ get_periods, get_start_dt, get_frequency
get_candles_number_from_minutes
from catalyst.exchange.utils.exchange_utils import get_exchange_symbols, \ 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 from logbook import Logger
log = Logger('Exchange', level=LOG_LEVEL) log = Logger('Exchange', level=LOG_LEVEL)
@@ -199,8 +198,12 @@ class Exchange:
) )
assets.append(asset) assets.append(asset)
except SymbolNotFoundOnExchange as e: except SymbolNotFoundOnExchange:
log.warn(e) log.debug(
'skipping non-existent market {} {}'.format(
self.name, symbol
)
)
return assets return assets
def get_asset(self, symbol, data_frequency=None, is_exchange_symbol=False, def get_asset(self, symbol, data_frequency=None, is_exchange_symbol=False,
@@ -253,8 +256,7 @@ class Exchange:
elif data_frequency is not None: elif data_frequency is not None:
applies = ( applies = (
( (
data_frequency == 'minute' and data_frequency == 'minute' and a.end_minute is not None)
a.end_minute is not None)
or ( or (
data_frequency == 'daily' and a.end_daily is not None) data_frequency == 'daily' and a.end_daily is not None)
) )
@@ -503,60 +505,49 @@ class Exchange:
freq, candle_size, unit, data_frequency = get_frequency( freq, candle_size, unit, data_frequency = get_frequency(
frequency, data_frequency, supported_freqs=['T', 'D', 'H'] 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 # The get_history method supports multiple asset
candles = self.get_candles( candles = self.get_candles(
freq=freq, freq=freq,
assets=assets, assets=assets,
bar_count=requested_bar_count, bar_count=bar_count,
end_dt=end_dt if not is_current else None, 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: for asset in candles:
if not candles[asset]: if candles[asset]:
raise NoCandlesReceivedFromExchange( first_candle = candles[asset][0]
bar_count=requested_bar_count, asset_series = self.get_series_from_candles(
candles=candles[asset],
start_dt=first_candle['last_traded'],
end_dt=end_dt, end_dt=end_dt,
asset=asset, data_frequency=frequency,
exchange=self.name) field=field,
)
# for avoiding unnecessary forward fill end_dt is taken back one second delta_candle_size = candle_size * 60 if unit == 'H' else candle_size
forward_fill_till_dt = end_dt - timedelta(seconds=1) # 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, if last_traded < adj_end_dt:
field=field, raise LastCandleTooEarlyError(
freq=frequency, last_traded=last_traded,
bar_count=requested_bar_count, end_dt=adj_end_dt,
end_dt=forward_fill_till_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 series[asset] = asset_series
# 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,
# )
df = pd.DataFrame(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, def get_history_window_with_bundle(self,
assets, assets,
@@ -604,10 +595,9 @@ class Exchange:
A dataframe containing the requested data. A dataframe containing the requested data.
""" """
# TODO: this function needs some work, # TODO: this function needs some work, we're currently using it just for benchmark data
# we're currently using it just for benchmark data
freq, candle_size, unit, data_frequency = get_frequency( 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 adj_bar_count = candle_size * bar_count
try: try:
@@ -631,7 +621,7 @@ class Exchange:
start_dt = get_start_dt(end_dt, adj_bar_count, data_frequency) start_dt = get_start_dt(end_dt, adj_bar_count, data_frequency)
trailing_dt = \ trailing_dt = \
series[asset].index[-1] + get_delta(1, data_frequency) \ 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 # The get_history method supports multiple asset
# Use the original frequency to let each api optimize # Use the original frequency to let each api optimize
-19
View File
@@ -163,25 +163,6 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
style) style)
return amount, 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): def round_order(self, amount, asset):
""" """
We need fractions with cryptocurrencies We need fractions with cryptocurrencies
+1 -1
View File
@@ -68,7 +68,7 @@ class TradingPairFeeSchedule(CommissionModel):
multiplier = maker \ multiplier = maker \
if ((order.amount > 0 and order.limit < transaction.price) if ((order.amount > 0 and order.limit < transaction.price)
or (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 fee = cost * multiplier
return fee return fee
+44 -45
View File
@@ -22,7 +22,7 @@ from catalyst.exchange.exchange_errors import EmptyValuesInBundleError, \
PricingDataNotLoadedError, DataCorruptionError, PricingDataValueError PricingDataNotLoadedError, DataCorruptionError, PricingDataValueError
from catalyst.exchange.utils.bundle_utils import range_in_bundle, \ from catalyst.exchange.utils.bundle_utils import range_in_bundle, \
get_bcolz_chunk, get_df_from_arrays, get_assets get_bcolz_chunk, get_df_from_arrays, get_assets
from catalyst.exchange.utils.datetime_utils import get_start_dt, \ from catalyst.exchange.utils.datetime_utils import get_delta, get_start_dt, \
get_period_label, get_month_start_end, get_year_start_end get_period_label, get_month_start_end, get_year_start_end
from catalyst.exchange.utils.exchange_utils import get_exchange_folder, \ from catalyst.exchange.utils.exchange_utils import get_exchange_folder, \
save_exchange_symbols, mixin_market_params, get_catalyst_symbol save_exchange_symbols, mixin_market_params, get_catalyst_symbol
@@ -232,12 +232,12 @@ class ExchangeBundle:
problem = '{name} ({start_dt} to {end_dt}) has empty ' \ problem = '{name} ({start_dt} to {end_dt}) has empty ' \
'periods: {dates}'.format( 'periods: {dates}'.format(
name=asset.symbol, name=asset.symbol,
start_dt=asset.start_date.strftime( start_dt=asset.start_date.strftime(
DATE_TIME_FORMAT), DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT), end_dt=end_dt.strftime(DATE_TIME_FORMAT),
dates=[date.strftime( dates=[date.strftime(
DATE_TIME_FORMAT) for date in dates]) DATE_TIME_FORMAT) for date in dates])
if empty_rows_behavior == 'warn': if empty_rows_behavior == 'warn':
log.warn(problem) log.warn(problem)
@@ -286,12 +286,12 @@ class ExchangeBundle:
problem = '{name} ({start_dt} to {end_dt}) has {threshold} ' \ problem = '{name} ({start_dt} to {end_dt}) has {threshold} ' \
'identical close values on: {dates}'.format( 'identical close values on: {dates}'.format(
name=asset.symbol, name=asset.symbol,
start_dt=asset.start_date.strftime(DATE_TIME_FORMAT), start_dt=asset.start_date.strftime(DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT), end_dt=end_dt.strftime(DATE_TIME_FORMAT),
threshold=threshold, threshold=threshold,
dates=[pd.to_datetime(date).strftime(DATE_TIME_FORMAT) dates=[pd.to_datetime(date).strftime(DATE_TIME_FORMAT)
for date in dates]) for date in dates])
problems.append(problem) problems.append(problem)
@@ -458,7 +458,7 @@ class ExchangeBundle:
last_entry = None last_entry = None
if start is None or \ 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 start = earliest_trade
if last_entry is not None and (end is None or end > last_entry): if last_entry is not None and (end is None or end > last_entry):
@@ -598,41 +598,16 @@ class ExchangeBundle:
# we want to give an end_date far in time # we want to give an end_date far in time
writer = self.get_writer(start_dt, end_dt, data_frequency) writer = self.get_writer(start_dt, end_dt, data_frequency)
if show_breakdown: if show_breakdown:
if chunks: for asset in 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( with maybe_show_progress(
all_chunks, chunks[asset],
show_progress, show_progress,
label='Ingesting {frequency} price data on ' label='Ingesting {frequency} price data for '
'{exchange}'.format( '{symbol} on {exchange}'.format(
exchange=self.exchange_name, exchange=self.exchange_name,
frequency=data_frequency, frequency=data_frequency,
)) as it: symbol=asset.symbol
)) as it:
for chunk in it: for chunk in it:
problems += self.ingest_ctable( problems += self.ingest_ctable(
asset=chunk['asset'], asset=chunk['asset'],
@@ -642,6 +617,30 @@ class ExchangeBundle:
empty_rows_behavior='strip', empty_rows_behavior='strip',
cleanup=True 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: pd.to_datetime(chunk['period'])
)
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: if show_report and len(problems) > 0:
log.info('problems during ingestion:{}\n'.format( log.info('problems during ingestion:{}\n'.format(
+2 -2
View File
@@ -296,7 +296,7 @@ class DataPortalExchangeBacktest(DataPortalExchangeBase):
bundle = self.exchange_bundles[exchange_name] # type: ExchangeBundle bundle = self.exchange_bundles[exchange_name] # type: ExchangeBundle
freq, candle_size, unit, adj_data_frequency = get_frequency( 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 adj_bar_count = candle_size * bar_count
@@ -312,7 +312,7 @@ class DataPortalExchangeBacktest(DataPortalExchangeBase):
algo_end_dt=self._last_available_session, 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) df = resample_history_df(pd.DataFrame(series), freq, field, start_dt)
return df return df
-7
View File
@@ -322,10 +322,3 @@ class BalanceTooLowError(ZiplineError):
'add positions to hold a free amount greater than {amount}, or clean ' 'add positions to hold a free amount greater than {amount}, or clean '
'the state of this algo and restart.' 'the state of this algo and restart.'
).strip() ).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 calendar
import math
import re import re
from datetime import datetime, timedelta, date 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 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 Get the frequency parameters.
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 Notes
----- -----
@@ -309,14 +305,14 @@ def get_frequency(freq, data_frequency=None, supported_freqs=['D', 'H', 'T']):
data_frequency = 'minute' data_frequency = 'minute'
elif unit.lower() == 'h': elif unit.lower() == 'h':
data_frequency = 'minute'
if 'H' in supported_freqs: if 'H' in supported_freqs:
unit = 'H' unit = 'H'
alias = '{}H'.format(candle_size) alias = '{}H'.format(candle_size)
else: else:
candle_size = candle_size * 60 candle_size = candle_size * 60
alias = '{}T'.format(candle_size) alias = '{}T'.format(candle_size)
data_frequency = 'minute'
else: else:
raise InvalidHistoryFrequencyAlias(freq=freq) raise InvalidHistoryFrequencyAlias(freq=freq)
@@ -330,33 +326,3 @@ def from_ms_timestamp(ms):
def get_epoch(): def get_epoch():
return pd.to_datetime('1970-1-1', utc=True) 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): def get_candles_df(candles, field, freq, bar_count, end_dt,
df = df.reindex(periods) previous_value=None):
# 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):
all_series = dict() all_series = dict()
for asset in candles: for asset in candles:
asset_df = transform_candles_to_df(candles[asset]) periods = pd.date_range(end=end_dt, periods=bar_count, freq=freq)
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)
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 = pd.DataFrame(all_series)
df.dropna(inplace=True) df.dropna(inplace=True)
+5 -18
View File
@@ -95,24 +95,11 @@ class TradingEnvironment(object):
if not trading_calendar: if not trading_calendar:
trading_calendar = get_calendar("NYSE") trading_calendar = get_calendar("NYSE")
# todo: uncomment and add a well defined benchmark self.benchmark_returns, self.treasury_curves = load(
# self.benchmark_returns, self.treasury_curves = load( trading_calendar.day,
# trading_calendar.day, trading_calendar.schedule.index,
# trading_calendar.schedule.index, self.bm_symbol,
# 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 self.exchange_tz = exchange_tz
@@ -1 +1 @@
0xf0ee6b27b759c9893ce4f094b49ad28fd15a23e4 0x7fAec9aaE31BE428DeAAE1be8195dF609079Fd10
File diff suppressed because one or more lines are too long
@@ -1 +1 @@
0xa64927358a82254be92eb1f1cb01de68d1787004 0x3985f5de8fddf2e8f7705cd360b498bf35ebfbc4
+72 -181
View File
@@ -7,7 +7,6 @@ import re
import shutil import shutil
import sys import sys
import time import time
import webbrowser
import bcolz import bcolz
import logbook import logbook
@@ -24,7 +23,7 @@ from catalyst.exchange.utils.stats_utils import set_print_settings
from catalyst.marketplace.marketplace_errors import ( from catalyst.marketplace.marketplace_errors import (
MarketplacePubAddressEmpty, MarketplaceDatasetNotFound, MarketplacePubAddressEmpty, MarketplaceDatasetNotFound,
MarketplaceNoAddressMatch, MarketplaceHTTPRequest, MarketplaceNoAddressMatch, MarketplaceHTTPRequest,
MarketplaceNoCSVFiles, MarketplaceRequiresPython3) MarketplaceNoCSVFiles)
from catalyst.marketplace.utils.auth_utils import get_key_secret, \ from catalyst.marketplace.utils.auth_utils import get_key_secret, \
get_signed_headers get_signed_headers
from catalyst.marketplace.utils.bundle_utils import merge_bundles 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, \ from catalyst.marketplace.utils.path_utils import get_bundle_folder, \
get_data_source_folder, get_marketplace_folder, \ get_data_source_folder, get_marketplace_folder, \
get_user_pubaddr, get_temp_bundles_folder, extract_bundle get_user_pubaddr, get_temp_bundles_folder, extract_bundle
from catalyst.utils.paths import ensure_directory
if sys.version_info.major < 3: if sys.version_info.major < 3:
import urllib import urllib
@@ -46,10 +44,7 @@ log = logbook.Logger('Marketplace', level=LOG_LEVEL)
class Marketplace: class Marketplace:
def __init__(self): def __init__(self):
global Web3 global Web3
try: from web3 import Web3, HTTPProvider
from web3 import Web3, HTTPProvider
except ImportError:
raise MarketplaceRequiresPython3()
self.addresses = get_user_pubaddr() self.addresses = get_user_pubaddr()
@@ -65,14 +60,10 @@ class Marketplace:
contract_url = urllib.urlopen(MARKETPLACE_CONTRACT) contract_url = urllib.urlopen(MARKETPLACE_CONTRACT)
self.mkt_contract_address = Web3.toChecksumAddress( self.mkt_contract_address = Web3.toChecksumAddress(
contract_url.readline().decode( contract_url.readline().strip())
contract_url.info().get_content_charset()).strip())
abi_url = urllib.urlopen(MARKETPLACE_CONTRACT_ABI) abi_url = urllib.urlopen(MARKETPLACE_CONTRACT_ABI)
abi_url = abi_url.read().decode( abi = json.load(abi_url)
abi_url.info().get_content_charset())
abi = json.loads(abi_url)
self.mkt_contract = self.web3.eth.contract( self.mkt_contract = self.web3.eth.contract(
self.mkt_contract_address, self.mkt_contract_address,
@@ -82,14 +73,10 @@ class Marketplace:
contract_url = urllib.urlopen(ENIGMA_CONTRACT) contract_url = urllib.urlopen(ENIGMA_CONTRACT)
self.eng_contract_address = Web3.toChecksumAddress( self.eng_contract_address = Web3.toChecksumAddress(
contract_url.readline().decode( contract_url.readline().strip())
contract_url.info().get_content_charset()).strip())
abi_url = urllib.urlopen(ENIGMA_CONTRACT_ABI) abi_url = urllib.urlopen(ENIGMA_CONTRACT_ABI)
abi_url = abi_url.read().decode( abi = json.load(abi_url)
abi_url.info().get_content_charset())
abi = json.loads(abi_url)
self.eng_contract = self.web3.eth.contract( self.eng_contract = self.web3.eth.contract(
self.eng_contract_address, self.eng_contract_address,
@@ -132,10 +119,9 @@ class Marketplace:
else: else:
while True: while True:
for i in range(0, len(self.addresses)): for i in range(0, len(self.addresses)):
print('{}\t{}\t{}\t{}'.format( print('{}\t{}\t{}'.format(
i, i,
self.addresses[i]['pubAddr'], self.addresses[i]['pubAddr'],
self.addresses[i]['wallet'].ljust(10),
self.addresses[i]['desc']) self.addresses[i]['desc'])
) )
address_i = int(input('Choose your address associated with ' address_i = int(input('Choose your address associated with '
@@ -150,10 +136,10 @@ class Marketplace:
return address, address_i return address, address_i
def sign_transaction(self, tx): def sign_transaction(self, from_address, tx):
url = 'https://www.mycrypto.com/#offline-transaction' print('\nVisit https://www.myetherwallet.com/#offline-transaction and '
print('\nVisit {url} and enter the following parameters:\n\n' 'enter the following parameters:\n\n'
'From Address:\t\t{_from}\n' 'From Address:\t\t{_from}\n'
'\n\tClick the "Generate Information" button\n\n' '\n\tClick the "Generate Information" button\n\n'
'To Address:\t\t{to}\n' 'To Address:\t\t{to}\n'
@@ -162,16 +148,13 @@ class Marketplace:
'Gas Price:\t\t[Accept the default value]\n' 'Gas Price:\t\t[Accept the default value]\n'
'Nonce:\t\t\t{nonce}\n' 'Nonce:\t\t\t{nonce}\n'
'Data:\t\t\t{data}\n'.format( 'Data:\t\t\t{data}\n'.format(
url=url, _from=from_address,
_from=tx['from'], to=tx['to'],
to=tx['to'], value=tx['value'],
value=tx['value'], gas=tx['gas'],
gas=tx['gas'], nonce=tx['nonce'],
nonce=tx['nonce'], data=tx['data'], )
data=tx['data'], ) )
)
webbrowser.open_new(url)
signed_tx = input('Copy and Paste the "Signed Transaction" ' signed_tx = input('Copy and Paste the "Signed Transaction" '
'field here:\n') 'field here:\n')
@@ -184,17 +167,16 @@ class Marketplace:
def check_transaction(self, tx_hash): def check_transaction(self, tx_hash):
if 'ropsten' in ETH_REMOTE_NODE: if 'ropsten' in ETH_REMOTE_NODE:
etherscan = 'https://ropsten.etherscan.io/tx/' etherscan = 'https://ropsten.etherscan.io/tx/{}'.format(
elif 'rinkeby' in ETH_REMOTE_NODE: tx_hash)
etherscan = 'https://rinkeby.etherscan.io/tx/'
else: else:
etherscan = 'https://etherscan.io/tx/' etherscan = 'https://etherscan.io/tx/{}'.format(tx_hash)
etherscan = '{}{}'.format(etherscan, tx_hash)
print('\nYou can check the outcome of your transaction here:\n' print('\nYou can check the outcome of your transaction here:\n'
'{}\n\n'.format(etherscan)) '{}\n\n'.format(etherscan))
def _list(self): def list(self):
data_sources = self.mkt_contract.functions.getAllProviders().call() data_sources = self.mkt_contract.functions.getAllProviders().call()
data = [] data = []
@@ -206,44 +188,15 @@ class Marketplace:
dataset=self.to_text(data_source) dataset=self.to_text(data_source)
) )
) )
return pd.DataFrame(data)
def list(self):
df = self._list()
df = pd.DataFrame(data)
set_print_settings() set_print_settings()
if df.empty: if df.empty:
print('There are no datasets available yet.') print('There are no datasets available yet.')
else: else:
print(df) print(df)
def subscribe(self, dataset=None): def subscribe(self, dataset):
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
dataset = dataset.lower() dataset = dataset.lower()
@@ -306,14 +259,14 @@ class Marketplace:
'buy: {} ENG. Get enough ENG to cover the costs of the ' 'buy: {} ENG. Get enough ENG to cover the costs of the '
'monthly\nsubscription for what you are trying to buy, ' 'monthly\nsubscription for what you are trying to buy, '
'and try again.'.format( 'and try again.'.format(
address, from_grains(balance), price)) address, from_grains(balance), price))
return return
while True: while True:
agree_pay = input('Please confirm that you agree to pay {} ENG ' agree_pay = input('Please confirm that you agree to pay {} ENG '
'for a monthly subscription to the dataset "{}" ' 'for a monthly subscription to the dataset "{}" '
'starting today. [default: Y] '.format( 'starting today. [default: Y] '.format(
price, dataset)) or 'y' price, dataset)) or 'y'
if agree_pay.lower() not in ('y', 'n'): if agree_pay.lower() not in ('y', 'n'):
print("Please answer Y or N.") print("Please answer Y or N.")
else: else:
@@ -334,11 +287,13 @@ class Marketplace:
self.mkt_contract_address, self.mkt_contract_address,
grains, grains,
).buildTransaction( ).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: try:
tx_hash = '0x{}'.format( tx_hash = '0x{}'.format(
bin_hex(self.web3.eth.sendRawTransaction(signed_tx)) bin_hex(self.web3.eth.sendRawTransaction(signed_tx))
@@ -373,11 +328,13 @@ class Marketplace:
tx = self.mkt_contract.functions.subscribe( tx = self.mkt_contract.functions.subscribe(
Web3.toHex(dataset), Web3.toHex(dataset),
).buildTransaction({ ).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: try:
tx_hash = '0x{}'.format(bin_hex( tx_hash = '0x{}'.format(bin_hex(
@@ -412,7 +369,7 @@ class Marketplace:
'You can now ingest this dataset anytime during the ' 'You can now ingest this dataset anytime during the '
'next month by running the following command:\n' 'next month by running the following command:\n'
'catalyst marketplace ingest --dataset={}'.format( 'catalyst marketplace ingest --dataset={}'.format(
dataset, address, dataset)) dataset, address, dataset))
def process_temp_bundle(self, ds_name, path): def process_temp_bundle(self, ds_name, path):
""" """
@@ -430,43 +387,17 @@ class Marketplace:
""" """
tmp_bundle = extract_bundle(path) tmp_bundle = extract_bundle(path)
bundle_folder = get_data_source_folder(ds_name) bundle_folder = get_data_source_folder(ds_name)
ensure_directory(bundle_folder)
if os.listdir(bundle_folder): if os.listdir(bundle_folder):
zsource = bcolz.ctable(rootdir=tmp_bundle, mode='r') zsource = bcolz.ctable(rootdir=tmp_bundle, mode='r')
ztarget = bcolz.ctable(rootdir=bundle_folder, mode='r') ztarget = bcolz.ctable(rootdir=bundle_folder, mode='r')
merge_bundles(zsource, ztarget) merge_bundles(zsource, ztarget)
else: else:
shutil.rmtree(bundle_folder, ignore_errors=True)
os.rename(tmp_bundle, bundle_folder) os.rename(tmp_bundle, bundle_folder)
def ingest(self, ds_name=None, start=None, end=None, force_download=False): pass
if ds_name is None: def ingest(self, ds_name, start=None, end=None, force_download=False):
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
# ds_name = ds_name.lower() # ds_name = ds_name.lower()
@@ -495,38 +426,29 @@ class Marketplace:
print('Your subscription to dataset "{}" expired on {} UTC.' print('Your subscription to dataset "{}" expired on {} UTC.'
'Please renew your subscription by running:\n' 'Please renew your subscription by running:\n'
'catalyst marketplace subscribe --dataset={}'.format( 'catalyst marketplace subscribe --dataset={}'.format(
ds_name, ds_name,
pd.to_datetime(check_sub[4], unit='s', utc=True), pd.to_datetime(check_sub[4], unit='s', utc=True),
ds_name) ds_name)
) )
if 'key' in self.addresses[address_i]: if 'key' in self.addresses[address_i]:
key = self.addresses[address_i]['key'] key = self.addresses[address_i]['key']
secret = self.addresses[address_i]['secret'] secret = self.addresses[address_i]['secret']
else: 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) headers = get_signed_headers(ds_name, key, secret)
log.info('Starting download of dataset for ingestion...') log.debug('Starting download of dataset for ingestion...')
r = requests.post( r = requests.post(
'{}/marketplace/ingest'.format(AUTH_SERVER), '{}/marketplace/ingest'.format(AUTH_SERVER),
headers=headers, headers=headers,
stream=True, stream=True,
) )
if r.status_code == 200: if r.status_code == 200:
log.info('Dataset downloaded successfully. Processing dataset...')
target_path = get_temp_bundles_folder() target_path = get_temp_bundles_folder()
try: try:
decoder = MultipartDecoder.from_response(r) decoder = MultipartDecoder.from_response(r)
# with maybe_show_progress(
# iter(decoder.parts),
# True,
# label='Processing files') as part:
counter = 1
for part in decoder.parts: for part in decoder.parts:
log.info("Processing file {} of {}".format(
counter, len(decoder.parts)))
h = part.headers[b'Content-Disposition'].decode('utf-8') h = part.headers[b'Content-Disposition'].decode('utf-8')
# Extracting the filename from the header # Extracting the filename from the header
name = re.search(r'filename="(.*)"', h).group(1) name = re.search(r'filename="(.*)"', h).group(1)
@@ -540,7 +462,6 @@ class Marketplace:
f.write(part.content) f.write(part.content)
self.process_temp_bundle(ds_name, filename) self.process_temp_bundle(ds_name, filename)
counter += 1
except NonMultipartContentTypeException: except NonMultipartContentTypeException:
response = r.json() response = r.json()
@@ -572,42 +493,17 @@ class Marketplace:
return df return df
def clean(self, ds_name=None, data_frequency=None): def clean(self, data_source_name, data_frequency=None):
data_source_name = data_source_name.lower()
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()
if data_frequency is None: if data_frequency is None:
folder = get_data_source_folder(ds_name) folder = get_data_source_folder(data_source_name)
else: else:
folder = get_bundle_folder(ds_name, data_frequency) folder = get_bundle_folder(data_source_name, data_frequency)
shutil.rmtree(folder) shutil.rmtree(folder)
pass
def create_metadata(self, key, secret, ds_name, data_frequency, desc, def create_metadata(self, key, secret, ds_name, data_frequency, desc,
has_history=True, has_live=True): has_history=True, has_live=True):
@@ -643,7 +539,7 @@ class Marketplace:
def register(self): def register(self):
while True: while True:
desc = input('Enter the name of the dataset to register: ') desc = input('Enter the name of the dataset to register: ')
dataset = desc.lower().strip() dataset = desc.lower()
provider_info = self.mkt_contract.functions.getDataProviderInfo( provider_info = self.mkt_contract.functions.getDataProviderInfo(
Web3.toHex(dataset) Web3.toHex(dataset)
).call() ).call()
@@ -699,8 +595,7 @@ class Marketplace:
key = self.addresses[address_i]['key'] key = self.addresses[address_i]['key']
secret = self.addresses[address_i]['secret'] secret = self.addresses[address_i]['secret']
else: else:
key, secret = get_key_secret(address, key, secret = get_key_secret(address)
self.addresses[address_i]['wallet'])
grains = to_grains(price) grains = to_grains(price)
@@ -709,11 +604,13 @@ class Marketplace:
grains, grains,
address, address,
).buildTransaction( ).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: try:
tx_hash = '0x{}'.format( tx_hash = '0x{}'.format(
@@ -724,7 +621,7 @@ class Marketplace:
) )
except Exception as e: except Exception as e:
print('Unable to register the requested dataset: {}'.format(e)) print('Unable to subscribe to data source: {}'.format(e))
return return
self.check_transaction(tx_hash) self.check_transaction(tx_hash)
@@ -781,34 +678,28 @@ class Marketplace:
key = match['key'] key = match['key']
secret = match['secret'] secret = match['secret']
else: else:
key, secret = get_key_secret(provider_info[0], match['wallet']) key, secret = get_key_secret(provider_info[0])
headers = get_signed_headers(dataset, key, secret)
filenames = glob.glob(os.path.join(datadir, '*.csv')) filenames = glob.glob(os.path.join(datadir, '*.csv'))
if not filenames: if not filenames:
raise MarketplaceNoCSVFiles(datadir=datadir) raise MarketplaceNoCSVFiles(datadir=datadir)
files = [] files = []
for idx, file in enumerate(filenames): for file in filenames:
log.info('Uploading file {} of {}: {}'.format(
idx+1, len(filenames), file))
files = []
files.append(('file', open(file, 'rb'))) files.append(('file', open(file, 'rb')))
headers = get_signed_headers(dataset, key, secret) r = requests.post('{}/marketplace/publish'.format(AUTH_SERVER),
r = requests.post('{}/marketplace/publish'.format(AUTH_SERVER), files=files,
files=files, headers=headers)
headers=headers)
if r.status_code != 200: if r.status_code != 200:
raise MarketplaceHTTPRequest(request='upload file', raise MarketplaceHTTPRequest(request='upload file',
error=r.status_code) error=r.status_code)
if 'error' in r.json(): if 'error' in r.json():
raise MarketplaceHTTPRequest(request='upload file', raise MarketplaceHTTPRequest(request='upload file',
error=r.json()['error']) error=r.json()['error'])
log.info('File processed successfully.') print('Dataset {} uploaded successfully.'.format(dataset))
print('\nDataset {} uploaded and processed successfully.'.format(
dataset))
+1 -10
View File
@@ -9,8 +9,7 @@ def silent_except_hook(exctype, excvalue, exctraceback):
MarketplaceNoAddressMatch, MarketplaceHTTPRequest, MarketplaceNoAddressMatch, MarketplaceHTTPRequest,
MarketplaceNoCSVFiles, MarketplaceContractDataNoMatch, MarketplaceNoCSVFiles, MarketplaceContractDataNoMatch,
MarketplaceSubscriptionExpired, MarketplaceJSONError, MarketplaceSubscriptionExpired, MarketplaceJSONError,
MarketplaceWalletNotSupported, MarketplaceEmptySignature, MarketplaceWalletNotSupported, MarketplaceEmptySignature]:
MarketplaceRequiresPython3]:
fn = traceback.extract_tb(exctraceback)[-1][0] fn = traceback.extract_tb(exctraceback)[-1][0]
ln = traceback.extract_tb(exctraceback)[-1][1] ln = traceback.extract_tb(exctraceback)[-1][1]
print("Error traceback: {1} (line {2})\n" print("Error traceback: {1} (line {2})\n"
@@ -87,11 +86,3 @@ class MarketplaceJSONError(ZiplineError):
'The configuration file {file} is malformed. Please correct ' 'The configuration file {file} is malformed. Please correct '
'the following error:\n{error}' '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.')
+9 -19
View File
@@ -1,6 +1,5 @@
import hashlib import hashlib
import hmac import hmac
import webbrowser
import requests import requests
import time import time
@@ -10,10 +9,10 @@ from catalyst.marketplace.marketplace_errors import (
MarketplaceEmptySignature) MarketplaceEmptySignature)
from catalyst.marketplace.utils.path_utils import ( from catalyst.marketplace.utils.path_utils import (
get_user_pubaddr, save_user_pubaddr) get_user_pubaddr, save_user_pubaddr)
from catalyst.constants import AUTH_SERVER, SUPPORTED_WALLETS from catalyst.constants import AUTH_SERVER
def get_key_secret(pubAddr, wallet): def get_key_secret(pubAddr, wallet='mew'):
""" """
Obtain a new key/secret pair from authentication server Obtain a new key/secret pair from authentication server
@@ -43,22 +42,14 @@ def get_key_secret(pubAddr, wallet):
auth_type, auth_info = header.split(None, 1) auth_type, auth_info = header.split(None, 1)
d = requests.utils.parse_dict_header(auth_info) d = requests.utils.parse_dict_header(auth_info)
nonce = 'Catalyst nonce: 0x{}'.format(d['nonce']) 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 ' print('\nObtaining a key/secret pair to streamline all future '
'requests with the authentication server.\n' 'requests with the authentication server.\n'
'Visit {url} and sign the ' 'Visit https://www.myetherwallet.com/signmsg.html and sign the '
'following message (copy the entire line, without the ' 'following message:\n{}'.format(nonce))
'line break at the end):\n\n{nonce}'.format( signature = input('Copy and Paste the "sig" field from '
url=url,
nonce=nonce))
webbrowser.open_new(url)
signature = input('\nCopy and Paste the "sig" field from '
'the signature here (without the double quotes, ' 'the signature here (without the double quotes, '
'only the HEX value):\n') 'only the HEX value):\n')
else: else:
@@ -92,8 +83,7 @@ def get_key_secret(pubAddr, wallet):
addresses = get_user_pubaddr() addresses = get_user_pubaddr()
match = next((l for l in addresses if match = next((l for l in addresses if
l['pubAddr'].lower() == pubAddr.lower()), None) l['pubAddr'] == pubAddr), None)
match['key'] = response.json()['key'] match['key'] = response.json()['key']
match['secret'] = response.json()['secret'] match['secret'] = response.json()['secret']
@@ -123,7 +113,7 @@ def get_signed_headers(ds_name, key, secret):
------- -------
""" """
nonce = str(int(time.time() * 1000)) nonce = str(int(time.time()))
signature = hmac.new( signature = hmac.new(
secret.encode('utf-8'), secret.encode('utf-8'),
+1 -59
View File
@@ -1,12 +1,8 @@
import os import os
import random
import re
import shutil import shutil
import bcolz import bcolz
import numpy as np
import pandas as pd import pandas as pd
from six import string_types
def merge_bundles(zsource, ztarget): def merge_bundles(zsource, ztarget):
@@ -31,64 +27,10 @@ def merge_bundles(zsource, ztarget):
df.drop_duplicates(inplace=True) df.drop_duplicates(inplace=True)
df.set_index(['date', 'symbol'], drop=False, inplace=True) df.set_index(['date', 'symbol'], drop=False, inplace=True)
sanitize_df(df)
dirname = os.path.basename(ztarget.rootdir) dirname = os.path.basename(ztarget.rootdir)
bak_dir = ztarget.rootdir.replace(dirname, '.{}'.format(dirname)) 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) z = bcolz.ctable.fromdataframe(df=df, rootdir=ztarget.rootdir)
shutil.rmtree(bak_dir) shutil.rmtree(bak_dir)
return z 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)
+2 -49
View File
@@ -2,7 +2,6 @@ import os
import json import json
import tarfile import tarfile
from catalyst.constants import SUPPORTED_WALLETS
from catalyst.utils.deprecate import deprecated from catalyst.utils.deprecate import deprecated
from catalyst.utils.paths import data_root, ensure_directory from catalyst.utils.paths import data_root, ensure_directory
from catalyst.marketplace.marketplace_errors import MarketplaceJSONError from catalyst.marketplace.marketplace_errors import MarketplaceJSONError
@@ -132,63 +131,17 @@ def get_user_pubaddr(environ=None):
try: try:
d = data[0]['pubAddr'] d = data[0]['pubAddr']
except Exception as e: except Exception as e:
data = [data, ] return [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 return data
else: else:
data = [] data = []
data.append(dict(pubAddr='', desc='', wallet='')) data.append(dict(pubAddr='', desc=''))
with open(filename, 'w') as f: with open(filename, 'w') as f:
json.dump(data, f, sort_keys=False, indent=2, json.dump(data, f, sort_keys=False, indent=2,
separators=(',', ':')) separators=(',', ':'))
return data 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): def save_user_pubaddr(data, environ=None):
""" """
Saves the user's public addresses and their related metadata in Saves the user's public addresses and their related metadata in
-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 import pandas as pd
from six import string_types from six import string_types
import catalyst
from catalyst.data.bundles import load from catalyst.data.bundles import load
from catalyst.data.data_portal import DataPortal from catalyst.data.data_portal import DataPortal
from catalyst.exchange.exchange_pricing_loader import ExchangePricingLoader, \ from catalyst.exchange.exchange_pricing_loader import ExchangePricingLoader, \
@@ -24,7 +23,7 @@ try:
from pygments.formatters import TerminalFormatter from pygments.formatters import TerminalFormatter
PYGMENTS = True PYGMENTS = True
except ImportError: except:
PYGMENTS = False PYGMENTS = False
from toolz import valfilter, concatv from toolz import valfilter, concatv
from functools import partial from functools import partial
@@ -152,7 +151,6 @@ def _run(handle_data,
'We encourage you to report any issue on GitHub: ' 'We encourage you to report any issue on GitHub: '
'https://github.com/enigmampc/catalyst/issues' 'https://github.com/enigmampc/catalyst/issues'
) )
log.info('Catalyst version {}'.format(catalyst.__version__))
sleep(3) sleep(3)
if live: if live:
@@ -263,15 +261,6 @@ def _run(handle_data,
# We still need to support bundles for other misc data, but we # We still need to support bundles for other misc data, but we
# can handle this later. # 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( data = DataPortalExchangeBacktest(
exchange_names=[exchange_name for exchange_name in exchanges], exchange_names=[exchange_name for exchange_name in exchanges],
asset_finder=None, asset_finder=None,
+1 -14
View File
@@ -145,6 +145,7 @@ with the following steps:
conda create --name catalyst python=3.6 scipy zlib conda create --name catalyst python=3.6 scipy zlib
3. Activate the environment: 3. Activate the environment:
**Linux or MacOS:** **Linux or MacOS:**
@@ -314,16 +315,6 @@ Troubleshooting ``pip`` Install
$ sudo apt-get install python-dev $ 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: .. _pipenv:
Installing with ``pipenv`` Installing with ``pipenv``
@@ -562,10 +553,6 @@ If after following the instructions above, and going through the
*Troubleshooting* sections, you still experience problems installing Catalyst, *Troubleshooting* sections, you still experience problems installing Catalyst,
you can seek additional help through the following channels: you can seek additional help through the following channels:
- Join our `Catalyst Forum <https://catalyst.enigma.co/>`_, and browse a variety
of topics and conversations around common issues that others face when using
Catalyst, and how to resolve them. And join the conversation!
- Join our `Discord community <https://discord.gg/SJK32GY>`_, and head over - Join our `Discord community <https://discord.gg/SJK32GY>`_, and head over
the #catalyst_dev channel where many other users (as well as the project the #catalyst_dev channel where many other users (as well as the project
developers) hang out, and can assist you with your particular issue. The developers) hang out, and can assist you with your particular issue. The
-79
View File
@@ -2,85 +2,6 @@
Release Notes Release Notes
============= =============
Version 0.5.7
^^^^^^^^^^^^^
**Release Date**: 2018-03-29
Build
~~~~~
- Data Marketplace deployed on mainnet.
- Added progress indicators for publishing data, and made the data publishing
synchronous to provide feedback to the publisher.
Bug Fixes
~~~~~~~~~
- Added arguments to the ``reduce`` function in tha Asset class :issue:`214`,
:issue:`287`
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, but believed to be supported.
- 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 Version 0.5.3
^^^^^^^^^^^^^ ^^^^^^^^^^^^^
**Release Date**: 2018-02-09 **Release Date**: 2018-02-09
+3 -4
View File
@@ -5,6 +5,7 @@ channels:
dependencies: dependencies:
- certifi=2016.2.28=py27_0 - certifi=2016.2.28=py27_0
- mkl=2017.0.3 - mkl=2017.0.3
- matplotlib=2.1.2=py36_0
- numpy=1.13.1=py27_0 - numpy=1.13.1=py27_0
- openssl=1.0.2l - openssl=1.0.2l
- pip=9.0.1=py27_1 - pip=9.0.1=py27_1
@@ -22,9 +23,7 @@ dependencies:
- bottleneck==1.2.1 - bottleneck==1.2.1
- chardet==3.0.4 - chardet==3.0.4
- ccxt==1.10.1094 - ccxt==1.10.1094
# The Enigma Data Marketplace requires Python3 because it depends on - web3==4.0.0b7
# web3, which requires Python3, as building its dependencies breaks in Python2
# - web3==4.0.0b7
- requests-toolbelt==0.8.0 - requests-toolbelt==0.8.0
- click==6.7 - click==6.7
- contextlib2==0.5.5 - contextlib2==0.5.5
@@ -39,7 +38,7 @@ dependencies:
- lru-dict==1.1.6 - lru-dict==1.1.6
- mako==1.0.7 - mako==1.0.7
- markupsafe==1.0 - markupsafe==1.0
- matplotlib==2.1.2 - matplotlib==2.1.0
- multipledispatch==0.4.9 - multipledispatch==0.4.9
- networkx==2.0 - networkx==2.0
- numexpr==2.6.4 - numexpr==2.6.4
+1 -1
View File
@@ -84,5 +84,5 @@ tables==3.3.0
ccxt==1.10.1094 ccxt==1.10.1094
boto3==1.4.8 boto3==1.4.8
redo==1.6 redo==1.6
web3==4.0.0b11; python_version > '3.4' web3==4.0.0b7
requests-toolbelt==0.8.0 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 BUNDLE_NAME_TEMPLATE
from catalyst.exchange.utils.bundle_utils import get_bcolz_chunk, \ from catalyst.exchange.utils.bundle_utils import get_bcolz_chunk, \
get_df_from_arrays 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.exchange_utils import get_exchange_folder
from catalyst.exchange.utils.factory import get_exchange from catalyst.exchange.utils.factory import get_exchange
from catalyst.exchange.utils.stats_utils import df_to_string from catalyst.exchange.utils.stats_utils import df_to_string
@@ -42,7 +42,7 @@ class TestExchangeBundle:
def test_ingest_minute(self): def test_ingest_minute(self):
data_frequency = 'minute' data_frequency = 'minute'
exchange_name = 'binance' exchange_name = 'poloniex'
exchange = get_exchange(exchange_name) exchange = get_exchange(exchange_name)
exchange_bundle = ExchangeBundle(exchange) exchange_bundle = ExchangeBundle(exchange)
@@ -50,8 +50,8 @@ class TestExchangeBundle:
exchange.get_asset('eth_btc') exchange.get_asset('eth_btc')
] ]
start = pd.to_datetime('2018-03-01', utc=True) start = pd.to_datetime('2016-03-01', utc=True)
end = pd.to_datetime('2018-03-8', utc=True) end = pd.to_datetime('2017-11-1', utc=True)
log.info('ingesting exchange bundle {}'.format(exchange_name)) log.info('ingesting exchange bundle {}'.format(exchange_name))
exchange_bundle.ingest( exchange_bundle.ingest(
@@ -101,7 +101,7 @@ class TestExchangeBundle:
# data_frequency = 'daily' # data_frequency = 'daily'
# include_symbols = 'neo_btc,bch_btc,eth_btc' # include_symbols = 'neo_btc,bch_btc,eth_btc'
exchange_name = 'binance' exchange_name = 'bitfinex'
data_frequency = 'minute' data_frequency = 'minute'
exchange = get_exchange(exchange_name) 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)) print('saved {} test results: {}'.format(end_dt, folder))
assert_frame_equal( assert_frame_equal(
right=data['bundle'][:-1], right=data['bundle'],
left=data['exchange'][:-1], left=data['exchange'],
check_less_precise=1, check_less_precise=1,
) )
try: try:
assert_frame_equal( assert_frame_equal(
right=data['bundle'][:-1], right=data['bundle'],
left=data['exchange'][:-1], left=data['exchange'],
check_less_precise=min([a.decimals for a in assets]), check_less_precise=min([a.decimals for a in assets]),
) )
except Exception as e: except Exception as e:
+3 -2
View File
@@ -1,5 +1,6 @@
from catalyst.marketplace.marketplace import Marketplace from catalyst.marketplace.marketplace import Marketplace
from catalyst.testing.fixtures import WithLogger, ZiplineTestCase from catalyst.testing.fixtures import WithLogger, ZiplineTestCase
import pandas as pd
class TestMarketplace(WithLogger, ZiplineTestCase): class TestMarketplace(WithLogger, ZiplineTestCase):
@@ -15,12 +16,12 @@ class TestMarketplace(WithLogger, ZiplineTestCase):
def test_subscribe(self): def test_subscribe(self):
marketplace = Marketplace() marketplace = Marketplace()
marketplace.subscribe('marketcap') marketplace.subscribe('marketcap2222')
pass pass
def test_ingest(self): def test_ingest(self):
marketplace = Marketplace() marketplace = Marketplace()
ds_def = marketplace.ingest('marketcap') ds_def = marketplace.ingest('github')
pass pass
def test_publish(self): def test_publish(self):