Compare commits

...
Author SHA1 Message Date
Avishai WeingartenandGitHub 19cdbcaa85 Revert "Allow jupyter to run as root inside Docker image" 2018-01-30 17:25:26 +02:00
Avishai WeingartenandGitHub c3662443e4 Merge pull request #180 from gthouret/jupyter-allow-root-docker
Allow jupyter to run as root inside Docker image
2018-01-30 17:22:47 +02:00
Frederic Fortier 911fb6e934 Merge branch 'gthouret-echo-usage-for-jupyter' 2018-01-24 00:33:52 -05:00
Guy Thouret afa2af3014 Add sys.stdout as second paramter to click.echo calls to prevent 'not writable' when running catalyst from Jupyter
Related to #179

Signed-off-by: Guy Thouret <guy@thouret.uk>
2018-01-22 12:54:48 +00:00
Guy Thouret 55521b87b9 Allow jupyter to run as root inside Docker image
Signed-off-by: Guy Thouret <guy@thouret.uk>
2018-01-22 11:11:37 +00:00
Frederic Fortier 45978ab193 Merge branch 'develop' 2018-01-19 18:18:15 -05:00
Frederic Fortier 575ac4a36d BLD: updated release notes for release 0.4.7 2018-01-19 18:17:07 -05:00
Frederic Fortier db07ac0abb Merge remote-tracking branch 'origin/develop' into develop 2018-01-19 18:13:03 -05:00
Frederic Fortier a754497d65 BLD: completed implementation of issue #60, authentication aliases 2018-01-19 18:12:54 -05:00
Victor Grau Serrat d669419d18 DOC: updating references to catalyst 2018-01-19 14:55:40 -07:00
Frederic Fortier 7d3c53dbef Merge branch 'develop' 2018-01-18 23:12:34 -05:00
Frederic Fortier 04a1513e73 DOC: updated CCXT dependency 2018-01-18 23:09:27 -05:00
Frederic Fortier 03d2c0e306 DOC: updated release notes for 0.4.6 2018-01-18 23:09:04 -05:00
Frederic Fortier 853707dfb2 BLD: updated CCXT and temporarily fixed symbol mapping issue 2018-01-18 21:57:11 -05:00
Frederic Fortier 540dd97dbf BLD: fixed benchmark loader 2018-01-18 19:36:54 -05:00
Frederic Fortier 9f372828b4 BLD: fixed benchmark loader 2018-01-18 19:30:16 -05:00
Frederic Fortier 6a929b6e25 BLD: updating example algos for testing 2018-01-18 19:26:33 -05:00
Frederic Fortier 45f278ab15 Merge remote-tracking branch 'origin/develop' into develop 2018-01-18 19:26:14 -05:00
Frederic Fortier a58fa21234 BLD: fixed benchmark loader 2018-01-18 19:26:07 -05:00
Victor Grau Serrat c8d6e07179 DOC: added update instructions (#156) 2018-01-18 17:02:21 -07:00
Frederic Fortier fdfc3f2ec3 Merge remote-tracking branch 'origin/develop' into develop 2018-01-18 18:58:21 -05:00
Frederic Fortier 3a321eb195 BLD: housekeeping and adjustments 2018-01-18 18:58:14 -05:00
Victor Grau Serrat d21eb3b946 DOC: improved documentation of paper trading mode (#168) 2018-01-18 16:29:09 -07:00
Frederic Fortier 439b5404ae BLD: cleanup and minor adjustments to get_candles() in live mode 2018-01-18 18:06:10 -05:00
Frederic Fortier 821f60897f BUG: troubleshooting issue #169 2018-01-18 18:05:14 -05:00
Frederic Fortier 5ff935723f Merge remote-tracking branch 'origin/develop' into develop 2018-01-18 17:09:43 -05:00
Frederic Fortier 51126fd7ae BLD: improved the bundle test suite and related adjustments 2018-01-18 17:09:37 -05:00
Victor Grau Serrat fec6a159e6 Merge branch 'develop' of github.com:enigmampc/catalyst into develop 2018-01-18 13:10:12 -07:00
Victor Grau Serrat eee9a07f54 DOC: updated timeframe of upcoming features 2018-01-18 13:10:06 -07:00
Frederic Fortier 563fc433d5 BUG: fixed issue with number of superfluous candles at the beginning of bundle history 2018-01-18 15:08:12 -05:00
Frederic Fortier 53a54fde7c Merge remote-tracking branch 'origin/develop' into develop 2018-01-18 14:25:07 -05:00
Frederic Fortier b0f2202b54 BUG: fixed bundle test suite 2018-01-18 14:08:11 -05:00
Victor Grau Serrat cf6c3bb76b BUG: switched benchmark from Poloniex to Bitfinex (#161) 2018-01-18 11:05:21 -07:00
Frederic Fortier 772640e098 BUG: fixed an issue with balancing transactions 2018-01-18 00:05:56 -05:00
Victor Grau Serrat 52d4ced37c BLD: live mode accepts end parameter, when algo finishes 2018-01-16 23:40:18 -07:00
Victor Grau Serrat 3ed44f72ad BLD: preservation of context.state dict between runs 2018-01-16 23:00:53 -07:00
Frederic Fortier 7569f7eb7c BUG: fixed issue #120 with currency substitution 2018-01-15 23:19:18 -05:00
Frederic Fortier 6d28e289c4 Merge remote-tracking branch 'origin/develop' into develop 2018-01-15 23:03:11 -05:00
Frederic Fortier 22154f2337 DOC: code documentation 2018-01-15 23:03:01 -05:00
Victor Grau Serrat 5314d3e1f8 Merge branch 'develop' of github.com:enigmampc/catalyst into develop 2018-01-13 06:17:02 -07:00
Victor Grau Serrat bb16975400 DOC: typo in link 2018-01-13 06:16:40 -07:00
Frederic Fortier 0ce624e6f6 BUG: removing computation of partial order for now to avoid a calculation issue 2018-01-12 22:00:10 -05:00
Frederic Fortier 50dc322230 BLD: minor adjustments and updated release notes 2018-01-12 19:53:26 -05:00
Frederic Fortier fb36435231 BLD: added subfolder to stats exports 2018-01-12 19:52:53 -05:00
Frederic Fortier e688783931 BLD: added auth alias to support more than one api token per exchange 2018-01-12 19:35:32 -05:00
Frederic Fortier 2f3dbeedcd Merge branch 'bfeeser-bfeeser_handle_ticker_errors' into develop 2018-01-12 18:08:45 -05:00
Frederic Fortier cebb1cd6d4 Merge branch 'bfeeser_handle_ticker_errors' of https://github.com/bfeeser/catalyst into bfeeser-bfeeser_handle_ticker_errors 2018-01-12 18:08:30 -05:00
Frederic Fortier 020ec50258 BUG: for issue #159, improved frequency validation in live mode 2018-01-12 17:39:49 -05:00
Ben Feeser df42cc9047 BUG: handle errors more gracefully when fetching tickers 2018-01-12 16:58:39 -05:00
Frederic Fortier 8a9b4e2df7 DOC: updated release notes for next release 2018-01-12 16:56:51 -05:00
Frederic Fortier b9150aab79 BUG: fixed issue with low order amount after adjustment 2018-01-12 16:46:57 -05:00
Ben Feeser d4efed0d3a DEV: add .python-version to .gitignore for pyenv 2018-01-12 16:45:20 -05:00
Frederic Fortier db1ad9aac8 BLD: For issue #151, significantly improved the way in which we are processing order for exchanges supporting "fetch-my-trades" 2018-01-12 16:36:23 -05:00
Frederic Fortier 270c261203 Merge remote-tracking branch 'origin/develop' into develop 2018-01-12 16:33:07 -05:00
Frederic Fortier 3f974b1adf DOC: For issue #151, significantly improved the way in which we are processing order for exchanges supporting "fetch-my-trades" 2018-01-12 16:33:00 -05:00
VictorandGitHub 9eac88344d Merge pull request #148 from treethought/tutorial-fixes
DOC: fix import and support python3 for beginner tutorial
2018-01-12 14:03:03 -07:00
VictorandGitHub 5d1f00bd19 Merge branch 'develop' into tutorial-fixes 2018-01-12 14:02:49 -07:00
VictorandGitHub 8d3f3ba81d Merge branch 'develop' into tutorial-fixes 2018-01-12 14:01:47 -07:00
VictorandGitHub 248299a725 Merge pull request #141 from danim7/master
DOC: wrong year: 2017 --> 2018 ??
2018-01-12 13:56:41 -07:00
VictorandGitHub e916283522 Merge pull request #155 from treethought/virtualenv-install
DOC: present virtualenv info before pip install command
2018-01-12 13:49:18 -07:00
Victor Grau Serrat b7a32656d5 DOC: improved doc of matplotlib error after installation 2018-01-12 12:15:04 -07:00
VictorandGitHub 78eb5d9d64 Update requirements_docs.txt 2018-01-12 09:00:44 -07:00
Victor Grau Serrat 14c7170159 DOC: sphinx/docutils issue when building docs 2018-01-12 08:57:00 -07:00
Victor Grau Serrat 415fceb11c DOC: typo - #158 2018-01-12 08:32:28 -07:00
Frederic Fortier 05c8957c90 BUG: fixed issue with history of multiple assets 2018-01-12 00:48:53 -05:00
Frederic Fortier e4bacb169e Merge branch 'treethought-python3-compatibility' into develop 2018-01-12 00:42:31 -05:00
Frederic Fortier 608dd843d6 Merge branch 'python3-compatibility' of https://github.com/treethought/catalyst into treethought-python3-compatibility 2018-01-12 00:42:19 -05:00
Cam Sweeney e0827fe4ad DOC: present virtualenv info before pip install command
Creating a virtualenv is now explained before the command for install
via pip. Following the instructions in the current state may cause
users to install using the system python, negating the point of the
virtualenv
2018-01-11 12:42:13 -08:00
Cam Sweeney 050fda1bdb MAINT: Convert dictionary .values() to list for python3 2018-01-11 10:36:44 -08:00
Frederic Fortier 713d487808 BUG: troubleshooting and minor fixes 2018-01-10 23:25:27 -05:00
Frederic Fortier ff7d7c5256 Merge branch 'develop' 2018-01-10 13:33:35 -05:00
Frederic Fortier b60b50e99a BUG: fixed python3 issue in run_algo 2018-01-10 13:31:54 -05:00
Cam Sweeney 93ebbf8b1f DOC: fix import and support python3 for beginner tutorial
The beginner tutorial Dual Moving Average example attempted to import
extract_transactions from the wrong location.

The tutorial and corresponding example also would fail using python3
due to indexing the view object returned via context.exchange.values()
2018-01-09 18:25:59 -08:00
Frederic Fortier db37c9c6a7 BLD: updated sample algo 2018-01-09 16:57:57 -05:00
Frederic Fortier a49cb55821 BUG: fixed issue #147 related to python 3 compatibility 2018-01-09 11:32:44 -05:00
Frederic Fortier ac15413af8 DOC: updated release notes for version 0.4.4 2018-01-09 01:09:23 -05:00
Frederic Fortier 141ee65c91 BLD: working on unit tests. 2018-01-09 01:08:57 -05:00
Frederic Fortier 55d1fee82d DOC: added an alpha warning message in response to issue #146 2018-01-08 17:01:58 -05:00
Frederic Fortier 33f94b3ef9 BLD: for issue #144, skipped cash verification when there are open orders. 2018-01-07 02:32:12 -05:00
Frederic Fortier 30448a65c5 BLD: Housekeeping 2018-01-06 19:58:11 -05:00
Frederic Fortier 9c33ee123c BLD: Housekeeping 2018-01-06 19:51:41 -05:00
Frederic Fortier d88710501f DOC: updated release notes in preparation for version 0.4.4. 2018-01-06 19:50:58 -05:00
Frederic Fortier ba0208910f BUG: Fixed issue #142 by removing unecessary capital_base check since we are now validating positions and cash against the exchange 2018-01-06 19:37:53 -05:00
Frederic Fortier 5d9708901d BUG: fixed issue #111 related to positions update after restoring algo state 2018-01-06 19:08:47 -05:00
danim7andGitHub 96b36c6614 DOC: wrong year: 2017 --> 2018 2018-01-06 23:55:07 +01:00
Frederic Fortier 13de3e69ef BLD: refined CCXT error handlers 2018-01-05 22:22:04 -05:00
Frederic Fortier 215de33c35 BUG: Fixed issue with updating positions after restoring the state of an algo 2018-01-05 02:46:11 -05:00
Frederic Fortier 790ac22f8d Merge branch 'develop' 2018-01-05 00:16:26 -05:00
Frederic Fortier 4d8d1e33d0 DOC: updated release notes 2018-01-05 00:15:19 -05:00
Frederic Fortier 595bd82234 BLD: improving unit tests 2018-01-05 00:13:18 -05:00
Frederic Fortier 9515d10cef DOC: updated release notes 2018-01-05 00:11:04 -05:00
Frederic Fortier 56481bbbe0 BUG: rolled back run_algo functions splitting to quickly resolve issue #137. We'll merge it more carefully in the next release. 2018-01-05 00:09:51 -05:00
Frederic Fortier 79e4854973 BUG: fixed potential issue with refreshing the stats in live mode 2018-01-04 22:46:21 -05:00
Frederic Fortier b4e7629e8b BLD: upgraded CCXT 2018-01-04 21:07:15 -05:00
Frederic Fortier a818d10d40 BLD: improving unit tests 2018-01-04 21:06:40 -05:00
Frederic Fortier bd4b0d2756 DOC: Fixed old release header 2018-01-04 02:46:10 -05:00
Frederic Fortier 80e29b2aa4 Merge branch 'develop' 2018-01-04 02:15:35 -05:00
Frederic Fortier 54ebfd6aad DOC: updating release notes with new version number 2018-01-04 02:15:05 -05:00
61 changed files with 2198 additions and 1010 deletions
+1
View File
@@ -40,6 +40,7 @@ develop-eggs
coverage.xml
htmlcov
nosetests.xml
.python-version
# C Extensions
*.o
+40 -17
View File
@@ -3,6 +3,7 @@ import os
from functools import wraps
import click
import sys
import logbook
import pandas as pd
from six import text_type
@@ -257,7 +258,7 @@ def run(ctx,
if capital_base is None:
ctx.fail("must specify a capital base with '--capital-base'")
click.echo('Running in backtesting mode.')
click.echo('Running in backtesting mode.', sys.stdout)
perf = _run(
initialize=None,
@@ -282,13 +283,15 @@ def run(ctx,
exchange=exchange_name,
algo_namespace=algo_namespace,
base_currency=base_currency,
analyze_live=None,
live_graph=False,
simulate_orders=True,
auth_aliases=None,
stats_output=None,
)
if output == '-':
click.echo(str(perf))
click.echo(str(perf), sys.stdout)
elif output != os.devnull: # make the catalyst magic not write any data
perf.to_pickle(output)
@@ -312,11 +315,11 @@ def catalyst_magic(line, cell=None):
'--algotext', cell,
'--output', os.devnull, # don't write the results by default
] + ([
# these options are set when running in line magic mode
# set a non None algo text to use the ipython user_ns
'--algotext', '',
'--local-namespace',
] if cell is None else []) + line.split(),
# these options are set when running in line magic mode
# set a non None algo text to use the ipython user_ns
'--algotext', '',
'--local-namespace',
] if cell is None else []) + line.split(),
'%s%%catalyst' % ((cell or '') and '%'),
# don't use system exit and propogate errors to the caller
standalone_mode=False,
@@ -393,6 +396,12 @@ def catalyst_magic(line, cell=None):
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,
@@ -406,6 +415,15 @@ def catalyst_magic(line, cell=None):
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 live(ctx,
algofile,
@@ -418,7 +436,9 @@ def live(ctx,
exchange_name,
algo_namespace,
base_currency,
end,
live_graph,
auth_aliases,
simulate_orders):
"""Trade live with the given algorithm.
"""
@@ -441,10 +461,10 @@ def live(ctx,
ctx.fail("must specify a capital base with '--capital-base'")
if simulate_orders:
click.echo('Running in paper trading mode.')
click.echo('Running in paper trading mode.', sys.stdout)
else:
click.echo('Running in live trading mode.')
click.echo('Running in live trading mode.', sys.stdout)
perf = _run(
initialize=None,
@@ -460,7 +480,7 @@ def live(ctx,
bundle=None,
bundle_timestamp=None,
start=None,
end=None,
end=end,
output=output,
print_algo=print_algo,
local_namespace=local_namespace,
@@ -470,12 +490,14 @@ def live(ctx,
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))
click.echo(str(perf), sys.stdout)
elif output != os.devnull: # make the catalyst magic not write any data
perf.to_pickle(output)
@@ -557,7 +579,7 @@ def ingest_exchange(ctx, exchange_name, data_frequency, start, end,
exchange_bundle = ExchangeBundle(exchange_name)
click.echo('Ingesting exchange bundle {}...'.format(exchange_name))
click.echo('Ingesting exchange bundle {}...'.format(exchange_name), sys.stdout)
exchange_bundle.ingest(
data_frequency=data_frequency,
include_symbols=include_symbols,
@@ -580,10 +602,11 @@ def ingest_exchange(ctx, exchange_name, data_frequency, start, end,
@click.pass_context
def clean_algo(ctx, algo_namespace):
click.echo(
'Cleaning algo state: {}'.format(algo_namespace)
'Cleaning algo state: {}'.format(algo_namespace),
sys.stdout
)
delete_algo_folder(algo_namespace)
click.echo('Done')
click.echo('Done', sys.stdout)
@main.command(name='clean-exchange')
@@ -610,11 +633,11 @@ def clean_exchange(ctx, exchange_name, data_frequency):
exchange_bundle = ExchangeBundle(exchange_name)
click.echo('Cleaning exchange bundle {}...'.format(exchange_name))
click.echo('Cleaning exchange bundle {}...'.format(exchange_name), sys.stdout)
exchange_bundle.clean(
data_frequency=data_frequency,
)
click.echo('Done')
click.echo('Done', sys.stdout)
@main.command()
@@ -735,7 +758,7 @@ def bundles():
# because there were no entries, print a single message indicating that
# no ingestions have yet been made.
for timestamp in ingestions or ["<no ingestions>"]:
click.echo("%s %s" % (bundle, timestamp))
click.echo("%s %s" % (bundle, timestamp), sys.stdout)
if __name__ == '__main__':
+2 -2
View File
@@ -88,11 +88,11 @@ class AssetDispatchBarReader(with_metaclass(ABCMeta)):
if self._last_available_dt is not None:
return self._last_available_dt
else:
return min(r.last_available_dt for r in self._readers.values())
return min(r.last_available_dt for r in list(self._readers.values()))
@lazyval
def first_trading_day(self):
return max(r.first_trading_day for r in self._readers.values())
return max(r.first_trading_day for r in list(self._readers.values()))
def get_value(self, sid, dt, field):
asset = self._asset_finder.retrieve_asset(sid)
+3 -2
View File
@@ -101,7 +101,7 @@ def load_crypto_market_data(trading_day=None, trading_days=None,
trading_day = get_calendar('OPEN').trading_day
# TODO: consider making configurable
bm_symbol = 'btc_usdt'
bm_symbol = 'btc_usd'
# if trading_days is None:
# trading_days = get_calendar('OPEN').schedule
@@ -144,8 +144,9 @@ def load_crypto_market_data(trading_day=None, trading_days=None,
# breaks things and it's only needed here
from catalyst.exchange.utils.factory import get_exchange
exchange = get_exchange(
exchange_name='poloniex', base_currency='usdt'
exchange_name='bitfinex', base_currency='usd'
)
exchange.init()
benchmark_asset = exchange.get_asset(bm_symbol)
+3 -3
View File
@@ -23,7 +23,7 @@ from catalyst.api import (order_target_value, symbol, record,
def initialize(context):
context.ASSET_NAME = 'btc_usd'
context.ASSET_NAME = 'btc_usdt'
context.TARGET_HODL_RATIO = 0.8
context.RESERVE_RATIO = 1.0 - context.TARGET_HODL_RATIO
@@ -140,9 +140,9 @@ if __name__ == '__main__':
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bitfinex',
exchange_name='poloniex',
algo_namespace='buy_and_hodl',
base_currency='usd',
base_currency='usdt',
start=pd.to_datetime('2015-03-01', utc=True),
end=pd.to_datetime('2017-10-31', utc=True),
)
+3 -3
View File
@@ -27,7 +27,7 @@ import pandas as pd
def initialize(context):
context.asset = symbol('btc_usd')
context.asset = symbol('btc_usdt')
def handle_data(context, data):
@@ -41,9 +41,9 @@ if __name__ == '__main__':
data_frequency='daily',
initialize=initialize,
handle_data=handle_data,
exchange_name='bitfinex',
exchange_name='poloniex',
algo_namespace='buy_and_hodl',
base_currency='usd',
base_currency='usdt',
start=pd.to_datetime('2015-03-01', utc=True),
end=pd.to_datetime('2017-10-31', utc=True),
)
+1 -1
View File
@@ -143,7 +143,7 @@ def analyze(context, stats):
if __name__ == '__main__':
live = False
live = True
if live:
run_algorithm(
capital_base=0.001,
+2 -1
View File
@@ -84,7 +84,8 @@ def handle_data(context, data):
def analyze(context, perf):
# Get the base_currency that was passed as a parameter to the simulation
base_currency = context.exchanges.values()[0].base_currency.upper()
exchange = list(context.exchanges.values())[0]
base_currency = exchange.base_currency.upper()
# First chart: Plot portfolio value using base_currency
ax1 = plt.subplot(411)
+11 -10
View File
@@ -37,14 +37,14 @@ def initialize(context):
context.base_price = None
context.current_day = None
context.RSI_OVERSOLD = 50
context.RSI_OVERBOUGHT = 65
context.CANDLE_SIZE = '5T'
context.RSI_OVERSOLD = 55
context.RSI_OVERBOUGHT = 60
context.CANDLE_SIZE = '15T'
context.start_time = time.time()
# context.set_commission(maker=0.1, taker=0.2)
context.set_slippage(spread=0.0001)
context.set_commission(maker=0.001, taker=0.002)
context.set_slippage(spread=0.001)
def handle_data(context, data):
@@ -114,7 +114,7 @@ def handle_data(context, data):
# TODO: retest with open orders
# Since we are using limit orders, some orders may not execute immediately
# we wait until all orders are executed before considering more trades.
orders = get_open_orders(context.market)
orders = context.blotter.open_orders
if len(orders) > 0:
log.info('exiting because orders are open: {}'.format(orders))
return
@@ -161,7 +161,7 @@ def analyze(context=None, perf=None):
import matplotlib.pyplot as plt
# The base currency of the algo exchange
base_currency = context.exchanges.values()[0].base_currency.upper()
base_currency = list(context.exchanges.values())[0].base_currency.upper()
# Plot the portfolio value over time.
ax1 = plt.subplot(611)
@@ -244,11 +244,11 @@ def analyze(context=None, perf=None):
if __name__ == '__main__':
# The execution mode: backtest or live
live = False
live = True
if live:
run_algorithm(
capital_base=0.03,
capital_base=0.01,
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
@@ -259,6 +259,7 @@ if __name__ == '__main__':
live_graph=False,
simulate_orders=False,
stats_output=None,
# auth_aliases=dict(poloniex='auth2')
)
else:
@@ -280,7 +281,7 @@ if __name__ == '__main__':
analyze=analyze,
exchange_name='bitfinex',
algo_namespace=NAMESPACE,
base_currency='eth',
base_currency='btc',
start=pd.to_datetime('2017-10-01', utc=True),
end=pd.to_datetime('2017-11-10', utc=True),
output=out
@@ -0,0 +1,288 @@
# For this example, we're going to write a simple momentum script. When the
# stock goes up quickly, we're going to buy; when it goes down quickly, we're
# going to sell. Hopefully we'll ride the waves.
import os
import tempfile
import time
import numpy as np
import pandas as pd
import talib
from logbook import Logger
from catalyst import run_algorithm
from catalyst.api import symbol, record, order_target_percent, get_open_orders
from catalyst.exchange.utils.stats_utils import extract_transactions
# We give a name to the algorithm which Catalyst will use to persist its state.
# In this example, Catalyst will create the `.catalyst/data/live_algos`
# directory. If we stop and start the algorithm, Catalyst will resume its
# state using the files included in the folder.
from catalyst.utils.paths import ensure_directory
NAMESPACE = 'mean_reversion_simple'
log = Logger(NAMESPACE)
# To run an algorithm in Catalyst, you need two functions: initialize and
# handle_data.
def initialize(context):
# This initialize function sets any data or variables that you'll use in
# your algorithm. For instance, you'll want to define the trading pair (or
# trading pairs) you want to backtest. You'll also want to define any
# parameters or values you're going to use.
# In our example, we're looking at Neo in Ether.
context.market = symbol('eth_btc')
context.base_price = None
context.current_day = None
context.RSI_OVERSOLD = 50
context.RSI_OVERBOUGHT = 60
context.CANDLE_SIZE = '5T'
context.start_time = time.time()
context.set_commission(maker=0.001, taker=0.002)
# context.set_slippage(spread=0.001)
def handle_data(context, data):
# This handle_data function is where the real work is done. Our data is
# minute-level tick data, and each minute is called a frame. This function
# runs on each frame of the data.
# We flag the first period of each day.
# Since cryptocurrencies trade 24/7 the `before_trading_starts` handle
# would only execute once. This method works with minute and daily
# frequencies.
today = data.current_dt.floor('1D')
if today != context.current_day:
context.traded_today = False
context.current_day = today
# We're computing the volume-weighted-average-price of the security
# defined above, in the context.market variable. For this example, we're
# using three bars on the 15 min bars.
# The frequency attribute determine the bar size. We use this convention
# for the frequency alias:
# http://pandas.pydata.org/pandas-docs/stable/timeseries.html#offset-aliases
prices = data.history(
context.market,
fields='close',
bar_count=50,
frequency=context.CANDLE_SIZE
)
# Ta-lib calculates various technical indicator based on price and
# volume arrays.
# In this example, we are comp
rsi = talib.RSI(prices.values, timeperiod=14)
# We need a variable for the current price of the security to compare to
# the average. Since we are requesting two fields, data.current()
# returns a DataFrame with
current = data.current(context.market, fields=['close', 'volume'])
price = current['close']
# If base_price is not set, we use the current value. This is the
# price at the first bar which we reference to calculate price_change.
if context.base_price is None:
context.base_price = price
price_change = (price - context.base_price) / context.base_price
cash = context.portfolio.cash
# Now that we've collected all current data for this frame, we use
# the record() method to save it. This data will be available as
# a parameter of the analyze() function for further analysis.
record(
volume=current['volume'],
price=price,
price_change=price_change,
rsi=rsi[-1],
cash=cash
)
# We are trying to avoid over-trading by limiting our trades to
# one per day.
if context.traded_today:
return
# TODO: retest with open orders
# Since we are using limit orders, some orders may not execute immediately
# we wait until all orders are executed before considering more trades.
orders = get_open_orders(context.market)
if len(orders) > 0:
log.info('exiting because orders are open: {}'.format(orders))
return
# Exit if we cannot trade
if not data.can_trade(context.market):
return
# Another powerful built-in feature of the Catalyst backtester is the
# portfolio object. The portfolio object tracks your positions, cash,
# cost basis of specific holdings, and more. In this line, we calculate
# how long or short our position is at this minute.
pos_amount = context.portfolio.positions[context.market].amount
if rsi[-1] <= context.RSI_OVERSOLD and pos_amount == 0:
log.info(
'{}: buying - price: {}, rsi: {}'.format(
data.current_dt, price, rsi[-1]
)
)
# Set a style for limit orders,
limit_price = price * 1.005
order_target_percent(
context.market, 1, limit_price=limit_price
)
context.traded_today = True
elif rsi[-1] >= context.RSI_OVERBOUGHT and pos_amount > 0:
log.info(
'{}: selling - price: {}, rsi: {}'.format(
data.current_dt, price, rsi[-1]
)
)
limit_price = price * 0.995
order_target_percent(
context.market, 0, limit_price=limit_price
)
context.traded_today = True
def analyze(context=None, perf=None):
end = time.time()
log.info('elapsed time: {}'.format(end - context.start_time))
import matplotlib.pyplot as plt
# The base currency of the algo exchange
base_currency = list(context.exchanges.values())[0].base_currency.upper()
# Plot the portfolio value over time.
ax1 = plt.subplot(611)
perf.loc[:, 'portfolio_value'].plot(ax=ax1)
ax1.set_ylabel('Portfolio\nValue\n({})'.format(base_currency))
# Plot the price increase or decrease over time.
ax2 = plt.subplot(612, sharex=ax1)
perf.loc[:, 'price'].plot(ax=ax2, label='Price')
ax2.set_ylabel('{asset}\n({base})'.format(
asset=context.market.symbol, base=base_currency
))
transaction_df = extract_transactions(perf)
if not transaction_df.empty:
buy_df = transaction_df[transaction_df['amount'] > 0]
sell_df = transaction_df[transaction_df['amount'] < 0]
ax2.scatter(
buy_df.index.to_pydatetime(),
perf.loc[buy_df.index.floor('1 min'), 'price'],
marker='^',
s=100,
c='green',
label=''
)
ax2.scatter(
sell_df.index.to_pydatetime(),
perf.loc[sell_df.index.floor('1 min'), 'price'],
marker='v',
s=100,
c='red',
label=''
)
ax4 = plt.subplot(613, sharex=ax1)
perf.loc[:, 'cash'].plot(
ax=ax4, label='Base Currency ({})'.format(base_currency)
)
ax4.set_ylabel('Cash\n({})'.format(base_currency))
perf['algorithm'] = perf.loc[:, 'algorithm_period_return']
ax5 = plt.subplot(614, sharex=ax1)
perf.loc[:, ['algorithm', 'price_change']].plot(ax=ax5)
ax5.set_ylabel('Percent\nChange')
ax6 = plt.subplot(615, sharex=ax1)
perf.loc[:, 'rsi'].plot(ax=ax6, label='RSI')
ax6.set_ylabel('RSI')
ax6.axhline(context.RSI_OVERBOUGHT, color='darkgoldenrod')
ax6.axhline(context.RSI_OVERSOLD, color='darkgoldenrod')
if not transaction_df.empty:
ax6.scatter(
buy_df.index.to_pydatetime(),
perf.loc[buy_df.index.floor('1 min'), 'rsi'],
marker='^',
s=100,
c='green',
label=''
)
ax6.scatter(
sell_df.index.to_pydatetime(),
perf.loc[sell_df.index.floor('1 min'), 'rsi'],
marker='v',
s=100,
c='red',
label=''
)
plt.legend(loc=3)
start, end = ax6.get_ylim()
ax6.yaxis.set_ticks(np.arange(0, end, end / 5))
# Show the plot.
plt.gcf().set_size_inches(18, 8)
plt.show()
pass
if __name__ == '__main__':
# The execution mode: backtest or live
live = False
if live:
run_algorithm(
capital_base=0.025,
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='poloniex',
live=True,
algo_namespace=NAMESPACE,
base_currency='btc',
live_graph=False,
simulate_orders=False,
stats_output=None,
)
else:
folder = os.path.join(
tempfile.gettempdir(), 'catalyst', NAMESPACE
)
ensure_directory(folder)
timestr = time.strftime('%Y%m%d-%H%M%S')
out = os.path.join(folder, '{}.p'.format(timestr))
# catalyst run -f catalyst/examples/mean_reversion_simple.py \
# -x bitfinex -s 2017-10-1 -e 2017-11-10 -c usdt -n mean-reversion \
# --data-frequency minute --capital-base 10000
run_algorithm(
capital_base=0.1,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bitfinex',
algo_namespace=NAMESPACE,
base_currency='eth',
start=pd.to_datetime('2017-10-01', utc=True),
end=pd.to_datetime('2017-11-10', utc=True),
output=out
)
log.info('saved perf stats: {}'.format(out))
+1 -1
View File
@@ -175,7 +175,7 @@ def handle_data(context, data):
def analyze(context=None, results=None):
import matplotlib.pyplot as plt
base_currency = context.exchanges.values()[0].base_currency.upper()
base_currency = list(context.exchanges.values())[0].base_currency.upper()
# Plot the portfolio and asset data.
ax1 = plt.subplot(611)
results.loc[:, 'portfolio_value'].plot(ax=ax1)
+1 -1
View File
@@ -57,7 +57,7 @@ def analyze(context, perf):
log.info('the stats: {}'.format(get_pretty_stats(perf)))
# The base currency of the algo exchange
base_currency = context.exchanges.values()[0].base_currency.upper()
base_currency = list(context.exchanges.values())[0].base_currency.upper()
# Plot the portfolio value over time.
ax1 = plt.subplot(611)
+3 -3
View File
@@ -41,8 +41,8 @@ from catalyst.exchange.utils.exchange_utils import get_exchange_symbols
def initialize(context):
context.i = -1 # minute counter
context.exchange = context.exchanges.values()[0].name.lower()
context.base_currency = context.exchanges.values()[0].base_currency.lower()
context.exchange = list(context.exchanges.values())[0].name.lower()
context.base_currency = list(context.exchanges.values())[0].base_currency.lower()
def handle_data(context, data):
@@ -65,7 +65,7 @@ def handle_data(context, data):
minutes = 30
# get lookback_days of history data: that is 'lookback' number of bins
lookback = one_day_in_minutes / minutes * lookback_days
lookback = int(one_day_in_minutes / minutes * lookback_days)
if not context.i % minutes and context.universe:
# we iterate for every pair in the current universe
for coin in context.coins:
+293 -66
View File
@@ -6,23 +6,29 @@ from collections import defaultdict
import ccxt
import pandas as pd
import six
from catalyst.assets._assets import TradingPair
from ccxt import ExchangeNotAvailable, InvalidOrder
from ccxt import InvalidOrder, NetworkError, \
ExchangeError
from logbook import Logger
from six import string_types
from catalyst.algorithm import MarketOrder
from catalyst.assets._assets import TradingPair
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.exchange import Exchange
from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import InvalidHistoryFrequencyError, \
ExchangeSymbolsNotFound, ExchangeRequestError, InvalidOrderStyle, \
ExchangeNotFoundError, CreateOrderError, InvalidHistoryTimeframeError
ExchangeNotFoundError, CreateOrderError, InvalidHistoryTimeframeError, \
UnsupportedHistoryFrequencyError
from catalyst.exchange.exchange_execution import ExchangeLimitOrder
from catalyst.exchange.utils.exchange_utils import mixin_market_params, \
from_ms_timestamp, get_epoch, get_exchange_folder, get_catalyst_symbol, \
get_exchange_folder, get_catalyst_symbol, \
get_exchange_auth
from catalyst.exchange.utils.datetime_utils import from_ms_timestamp, \
get_epoch, \
get_periods_range
from catalyst.finance.order import Order, ORDER_STATUS
from catalyst.finance.transaction import Transaction
log = Logger('CCXT', level=LOG_LEVEL)
@@ -55,6 +61,7 @@ class CCXT(Exchange):
'apiKey': key,
'secret': secret,
})
self.api.enableRateLimit = True
except Exception:
raise ExchangeNotFoundError(exchange_name=exchange_name)
@@ -70,6 +77,7 @@ class CCXT(Exchange):
self.max_requests_per_minute = 60
self.low_balance_threshold = 0.1
self.request_cpt = dict()
self._common_symbols = dict()
self.bundle = ExchangeBundle(self.name)
self.markets = None
@@ -105,7 +113,12 @@ class CCXT(Exchange):
with open(filename, 'w+') as f:
json.dump(self.markets, f, indent=4)
except ExchangeNotAvailable as e:
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to fetch markets {}: {}'.format(
self.name, e
)
)
raise ExchangeRequestError(error=e)
self.load_assets()
@@ -210,6 +223,21 @@ class CCXT(Exchange):
)
return market
def substitute_currency_code(self, currency, source='catalyst'):
if source == 'catalyst':
currency = currency.upper()
key = self.api.common_currency_code(currency)
self._common_symbols[key] = currency.lower()
return key
else:
if currency in self._common_symbols:
return self._common_symbols[currency]
else:
return currency.lower()
def get_symbol(self, asset_or_symbol, source='catalyst'):
"""
The CCXT symbol.
@@ -217,6 +245,7 @@ class CCXT(Exchange):
Parameters
----------
asset_or_symbol
source
Returns
-------
@@ -226,7 +255,13 @@ class CCXT(Exchange):
if source == 'ccxt':
if isinstance(asset_or_symbol, string_types):
parts = asset_or_symbol.split('/')
return '{}_{}'.format(parts[0].lower(), parts[1].lower())
base_currency = self.substitute_currency_code(
parts[0], source
)
quote_currency = self.substitute_currency_code(
parts[1], source
)
return '{}_{}'.format(base_currency, quote_currency)
else:
return asset_or_symbol.symbol
@@ -237,7 +272,13 @@ class CCXT(Exchange):
) else asset_or_symbol.symbol
parts = symbol.split('_')
return '{}/{}'.format(parts[0].upper(), parts[1].upper())
base_currency = self.substitute_currency_code(
parts[0], source
)
quote_currency = self.substitute_currency_code(
parts[1], source
)
return '{}/{}'.format(base_currency, quote_currency)
@staticmethod
def map_frequency(value, source='ccxt', raise_error=True):
@@ -362,7 +403,7 @@ class CCXT(Exchange):
timeframe, source='ccxt', raise_error=raise_error
)
def get_candles(self, freq, assets, bar_count=None, start_dt=None,
def get_candles(self, freq, assets, bar_count=1, start_dt=None,
end_dt=None):
is_single = (isinstance(assets, TradingPair))
if is_single:
@@ -371,17 +412,46 @@ class CCXT(Exchange):
symbols = self.get_symbols(assets)
timeframe = CCXT.get_timeframe(freq)
ms = None
if timeframe not in self.api.timeframes:
freqs = [CCXT.get_frequency(t) for t in self.api.timeframes]
raise UnsupportedHistoryFrequencyError(
exchange=self.name,
freq=freq,
freqs=freqs,
)
if start_dt is not None and end_dt is not None:
raise ValueError(
'Please provide either start_dt or end_dt, not both.'
)
elif end_dt is not None:
# Make sure that end_dt really wants data in the past
# if it's close to now, we skip the 'since' parameters to
# lower the probability of error
bars_to_now = pd.date_range(
end_dt, pd.Timestamp.utcnow(), freq=freq
)
# See: https://github.com/ccxt/ccxt/issues/1360
if len(bars_to_now) > 1 or self.name in ['poloniex']:
dt_range = get_periods_range(
end_dt=end_dt,
periods=bar_count,
freq=freq,
)
start_dt = dt_range[0]
since = None
if start_dt is not None:
delta = start_dt - get_epoch()
ms = int(delta.total_seconds()) * 1000
since = int(delta.total_seconds()) * 1000
candles = dict()
for asset in assets:
for index, asset in enumerate(assets):
ohlcvs = self.api.fetch_ohlcv(
symbol=symbols[0],
symbol=symbols[index],
timeframe=timeframe,
since=ms,
since=since,
limit=bar_count,
params={}
)
@@ -398,6 +468,9 @@ class CCXT(Exchange):
close=ohlcv[4],
volume=ohlcv[5]
))
candles[asset] = sorted(
candles[asset], key=lambda c: c['last_traded']
)
if is_single:
return six.next(six.itervalues(candles))
@@ -408,6 +481,7 @@ class CCXT(Exchange):
def _fetch_symbol_map(self, is_local):
try:
return self.fetch_symbol_map(is_local)
except ExchangeSymbolsNotFound:
return None
@@ -559,8 +633,12 @@ class CCXT(Exchange):
for key in balances:
balances_lower[key.lower()] = balances[key]
except Exception as e:
log.debug('error retrieving balances: {}', e)
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to fetch balance {}: {}'.format(
self.name, e
)
)
raise ExchangeRequestError(error=e)
return balances_lower
@@ -695,6 +773,12 @@ class CCXT(Exchange):
else:
adj_amount = abs(amount)
if adj_amount == 0:
raise CreateOrderError(
exchange=self.name,
e='order amount lower than the smallest lot: {}'.format(amount)
)
try:
result = self.api.create_order(
symbol=symbol,
@@ -703,14 +787,18 @@ class CCXT(Exchange):
amount=adj_amount,
price=price
)
except ExchangeNotAvailable as e:
log.debug('unable to create order: {}'.format(e))
raise ExchangeRequestError(error=e)
except InvalidOrder as e:
log.warn('the exchange rejected the order: {}'.format(e))
raise CreateOrderError(exchange=self.name, error=e)
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to create order {} / {}: {}'.format(
self.name, symbol, e
)
)
raise ExchangeRequestError(error=e)
if 'info' not in result:
raise ValueError('cannot use order without info attribute')
@@ -735,18 +823,121 @@ class CCXT(Exchange):
limit=None,
params=dict()
)
except Exception as e:
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to fetch open orders {} / {}: {}'.format(
self.name, asset.symbol, e
)
)
raise ExchangeRequestError(error=e)
orders = []
for order_status in result:
order, executed_price = self._create_order(order_status)
order, _ = self._create_order(order_status)
if asset is None or asset == order.sid:
orders.append(order)
return orders
def get_order(self, order_id, asset_or_symbol=None):
def _process_order_fallback(self, order):
"""
Fallback method for exchanges which do not play nice with
fetch-my-trades. Apparently, about 60% of exchanges will return
the correct executed values with this method. Others will support
fetch-my-trades.
Parameters
----------
order: Order
Returns
-------
float
"""
exc_order, price = self.get_order(
order.id, order.asset, return_price=True
)
order.status = exc_order.status
order.commission = exc_order.commission
if order.amount != exc_order.amount:
log.warn(
'executed order amount {} differs '
'from original'.format(
exc_order.amount, order.amount
)
)
order.amount = exc_order.amount
if order.status == ORDER_STATUS.FILLED:
transaction = Transaction(
asset=order.asset,
amount=order.amount,
dt=pd.Timestamp.utcnow(),
price=price,
order_id=order.id,
commission=order.commission
)
return [transaction]
def process_order(self, order):
# TODO: move to parent class after tracking features in the parent
if not self.api.hasFetchMyTrades:
return self._process_order_fallback(order)
try:
all_trades = self.get_trades(order.asset)
except ExchangeRequestError as e:
log.warn(
'unable to fetch account trades, trying an alternate '
'method to find executed order {} / {}: {}'.format(
order.id, order.asset.symbol, e
)
)
return self._process_order_fallback(order)
transactions = []
trades = [t for t in all_trades if t['order'] == order.id]
if not trades:
log.debug(
'order {} / {} not found in trades'.format(
order.id, order.asset.symbol
)
)
return transactions
trades.sort(key=lambda t: t['timestamp'], reverse=False)
order.filled = 0
order.commission = 0
for trade in trades:
# status property will update automatically
filled = trade['amount'] * order.direction
order.filled += filled
commission = 0
if 'fee' in trade and 'cost' in trade['fee']:
commission = trade['fee']['cost']
order.commission += commission
order.check_triggers(
price=trade['price'],
dt=pd.to_datetime(trade['timestamp'], unit='ms', utc=True),
)
transaction = Transaction(
asset=order.asset,
amount=filled,
dt=pd.Timestamp.utcnow(),
price=trade['price'],
order_id=order.id,
commission=commission
)
transactions.append(transaction)
order.broker_order_id = ', '.join([t['id'] for t in trades])
return transactions
def get_order(self, order_id, asset_or_symbol=None, return_price=False):
if asset_or_symbol is None:
log.debug(
'order not found in memory, the request might fail '
@@ -758,10 +949,19 @@ class CCXT(Exchange):
order_status = self.api.fetch_order(id=order_id, symbol=symbol)
order, executed_price = self._create_order(order_status)
except Exception as e:
raise ExchangeRequestError(error=e)
if return_price:
return order, executed_price
return order, executed_price
else:
return order
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to fetch order {} / {}: {}'.format(
self.name, order_id, e
)
)
raise ExchangeRequestError(error=e)
def cancel_order(self, order_param, asset_or_symbol=None):
order_id = order_param.id \
@@ -777,7 +977,12 @@ class CCXT(Exchange):
if asset_or_symbol is not None else None
self.api.cancel_order(id=order_id, symbol=symbol)
except Exception as e:
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to cancel order {} / {}: {}'.format(
self.name, order_id, e
)
)
raise ExchangeRequestError(error=e)
def tickers(self, assets):
@@ -793,48 +998,46 @@ class CCXT(Exchange):
list[dict[str, float]
"""
tickers = dict()
try:
for asset in assets:
symbol = self.get_symbol(asset)
# TODO: use fetch_tickers() for efficiency
# I tried using fetch_tickers() but noticed some
# inconsistencies, see issue:
# https://github.com/ccxt/ccxt/issues/870
tickers = {}
for asset in assets:
symbol = self.get_symbol(asset)
# Test the CCXT throttling further to see if we need this
self.ask_request()
# TODO: use fetch_tickers() for efficiency
# I tried using fetch_tickers() but noticed some
# inconsistencies, see issue:
# https://github.com/ccxt/ccxt/issues/870
try:
ticker = self.api.fetch_ticker(symbol=symbol)
if not ticker:
log.warn('ticker not found for {} {}'.format(
self.name, symbol
))
continue
ticker['last_traded'] = from_ms_timestamp(ticker['timestamp'])
if 'last_price' not in ticker:
# TODO: any more exceptions?
ticker['last_price'] = ticker['last']
if 'baseVolume' in ticker and ticker['baseVolume'] is not None:
# Using the volume represented in the base currency
ticker['volume'] = ticker['baseVolume']
elif 'info' in ticker and 'bidQty' in ticker['info'] \
and 'askQty' in ticker['info']:
ticker['volume'] = float(ticker['info']['bidQty']) + \
float(ticker['info']['askQty'])
else:
ticker['volume'] = 0
tickers[asset] = ticker
except ExchangeNotAvailable as e:
log.warn(
'unable to fetch ticker: {} {}'.format(
self.name, asset.symbol
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to fetch ticker {} / {}: {}'.format(
self.name, asset.symbol, e
)
)
)
raise ExchangeRequestError(error=e)
continue
ticker['last_traded'] = from_ms_timestamp(ticker['timestamp'])
if 'last_price' not in ticker:
# TODO: any more exceptions?
ticker['last_price'] = ticker['last']
if 'baseVolume' in ticker and ticker['baseVolume'] is not None:
# Using the volume represented in the base currency
ticker['volume'] = ticker['baseVolume']
elif 'info' in ticker and 'bidQty' in ticker['info'] \
and 'askQty' in ticker['info']:
ticker['volume'] = float(ticker['info']['bidQty']) + \
float(ticker['info']['askQty'])
else:
ticker['volume'] = 0
tickers[asset] = ticker
return tickers
@@ -864,3 +1067,27 @@ class CCXT(Exchange):
))
return result
def get_trades(self, asset, my_trades=True, start_dt=None, limit=None):
if not my_trades:
raise NotImplemented(
'get_trades only supports "my trades"'
)
# TODO: is it possible to sort this? Limit is useless otherwise.
ccxt_symbol = self.get_symbol(asset)
try:
trades = self.api.fetch_my_trades(
symbol=ccxt_symbol,
since=start_dt,
limit=limit,
)
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to fetch trades {} / {}: {}'.format(
self.name, asset.symbol, e
)
)
raise ExchangeRequestError(error=e)
return trades
+66 -32
View File
@@ -5,8 +5,6 @@ from time import sleep
import numpy as np
import pandas as pd
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.data.data_portal import BASE_FIELDS
from catalyst.exchange.exchange_bundle import ExchangeBundle
@@ -15,10 +13,12 @@ from catalyst.exchange.exchange_errors import MismatchingBaseCurrencies, \
PricingDataNotLoadedError, \
NoDataAvailableOnExchange, NoValueForField, LastCandleTooEarlyError, \
TickerNotFoundError, NotEnoughCashError
from catalyst.exchange.utils.bundle_utils import get_start_dt, \
get_delta, get_periods, get_periods_range
from catalyst.exchange.utils.datetime_utils import get_delta, \
get_periods_range, \
get_periods, get_start_dt, get_frequency
from catalyst.exchange.utils.exchange_utils import get_exchange_symbols, \
get_frequency, resample_history_df, has_bundle
resample_history_df, has_bundle
from logbook import Logger
log = Logger('Exchange', level=LOG_LEVEL)
@@ -234,11 +234,15 @@ class Exchange:
"""
asset = None
# TODO: temp mapping, fix to use a single symbol convention
og_symbol = symbol
symbol = self.get_symbol(symbol) if not is_exchange_symbol else symbol
log.debug(
'searching assets for: {} {}'.format(
self.name, symbol
)
)
# TODO: simplify and loose the loop
for a in self.assets:
if asset is not None:
break
@@ -260,7 +264,8 @@ class Exchange:
# The symbol provided may use the Catalyst or the exchange
# convention
key = a.exchange_symbol if is_exchange_symbol else a.symbol
key = a.exchange_symbol if \
is_exchange_symbol else self.get_symbol(a)
if not asset and key.lower() == symbol.lower():
if applies:
asset = a
@@ -276,7 +281,7 @@ class Exchange:
supported_symbols = sorted([a.symbol for a in self.assets])
raise SymbolNotFoundOnExchange(
symbol=symbol,
symbol=og_symbol,
exchange=self.name.title(),
supported_symbols=supported_symbols
)
@@ -434,7 +439,7 @@ class Exchange:
series = pd.Series(values, index=dates)
periods = get_periods_range(
start_dt, end_dt, data_frequency
start_dt=start_dt, end_dt=end_dt, freq=data_frequency
)
# TODO: ensure that this working as expected, if not use fillna
series = series.reindex(
@@ -498,39 +503,37 @@ class Exchange:
freq, candle_size, unit, data_frequency = get_frequency(
frequency, data_frequency
)
adj_bar_count = candle_size * bar_count
start_dt = get_start_dt(end_dt, adj_bar_count, data_frequency)
# The get_history method supports multiple asset
candles = self.get_candles(
freq=freq,
assets=assets,
bar_count=bar_count,
start_dt=start_dt if not is_current else None,
end_dt=end_dt if not is_current else None,
)
series = dict()
for asset in candles:
first_candle = candles[asset][0]
asset_series = self.get_series_from_candles(
candles=candles[asset],
start_dt=start_dt,
start_dt=first_candle['last_traded'],
end_dt=end_dt,
data_frequency=frequency,
field=field,
)
if end_dt is not None:
delta = get_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,
)
# Checking to make sure that the dates match
delta = get_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)
@@ -584,11 +587,11 @@ class Exchange:
A dataframe containing the requested 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
)
adj_bar_count = candle_size * bar_count
try:
series = self.bundle.get_history_window_series_and_load(
assets=assets,
@@ -615,15 +618,14 @@ class Exchange:
# The get_history method supports multiple asset
# Use the original frequency to let each api optimize
# the size of result sets
trailing_bar_count = get_periods(
trailing_bars = get_periods(
trailing_dt, end_dt, freq
)
candles = self.get_candles(
freq=freq,
assets=asset,
bar_count=trailing_bar_count,
start_dt=start_dt,
end_dt=end_dt
end_dt=end_dt,
bar_count=trailing_bars if trailing_bars < 500 else 500,
)
last_value = series[asset].iloc(0) if asset in series \
@@ -900,6 +902,22 @@ class Exchange:
"""
pass
@abstractmethod
def process_order(self, order):
"""
Similar to get_order but looks only for executed orders.
Parameters
----------
order: Order
Returns
-------
float
Avg execution price
"""
@abstractmethod
def cancel_order(self, order_param, symbol_or_asset=None):
"""Cancel an open order.
@@ -914,8 +932,7 @@ class Exchange:
pass
@abstractmethod
def get_candles(self, freq, assets, bar_count=None,
start_dt=None, end_dt=None):
def get_candles(self, freq, assets, bar_count, start_dt=None, end_dt=None):
"""
Retrieve OHLCV candles for the given assets
@@ -980,7 +997,7 @@ class Exchange:
@abc.abstractmethod
def get_orderbook(self, asset, order_type, limit):
"""
Retrieve the the orderbook for the given trading pair.
Retrieve the orderbook for the given trading pair.
Parameters
----------
@@ -994,3 +1011,20 @@ class Exchange:
list[dict[str, float]
"""
pass
@abc.abstractmethod
def get_trades(self, asset, my_trades, start_dt, limit):
"""
Retrieve a list of trades.
Parameters
----------
my_trades: bool
List only my trades.
start_dt
limit
Returns
-------
"""
+100 -44
View File
@@ -18,11 +18,9 @@ from datetime import timedelta
from os import listdir
from os.path import isfile, join
import catalyst.protocol as zp
import logbook
import pandas as pd
from redo import retry
import catalyst.protocol as zp
from catalyst.algorithm import TradingAlgorithm
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.exchange_blotter import ExchangeBlotter
@@ -49,6 +47,7 @@ from catalyst.utils.api_support import api_method
from catalyst.utils.input_validation import error_keywords, ensure_upper_case
from catalyst.utils.math_utils import round_nearest
from catalyst.utils.preprocess import preprocess
from redo import retry
log = logbook.Logger('exchange_algorithm', level=LOG_LEVEL)
@@ -130,7 +129,7 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
@api_method
def set_commission(self, maker=None, taker=None):
key = self.blotter.commission_models.keys()[0]
key = list(self.blotter.commission_models.keys())[0]
if maker is not None:
self.blotter.commission_models[key].maker = maker
@@ -139,7 +138,7 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
@api_method
def set_slippage(self, spread=None):
key = self.blotter.slippage_models.keys()[0]
key = list(self.blotter.slippage_models.keys())[0]
if spread is not None:
self.blotter.slippage_models[key].spread = spread
@@ -271,9 +270,9 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
# Merging latest recorded variables
stats.update(self.recorded_vars)
stats['positions'] = cum.position_tracker.get_positions_list()
period = tracker.todays_performance
stats['positions'] = period.position_tracker.get_positions_list()
# we want the key to be absent, not just empty
# Only include transactions for given dt
stats['transactions'] = []
@@ -304,6 +303,7 @@ class ExchangeTradingAlgorithmBacktest(ExchangeTradingAlgorithmBase):
super(ExchangeTradingAlgorithmBacktest, self).__init__(*args, **kwargs)
self.frame_stats = list()
self.state = {}
log.info('initialized trading algorithm in backtest mode')
def is_last_frame_of_day(self, data):
@@ -351,6 +351,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
self.live_graph = kwargs.pop('live_graph', None)
self.stats_output = kwargs.pop('stats_output', None)
self._analyze_live = kwargs.pop('analyze_live', None)
self.end = kwargs.pop('end', None)
self._clock = None
self.frame_stats = list()
@@ -391,16 +392,20 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
'before exiting the algorithm.')
algo_folder = get_algo_folder(self.algo_namespace)
folder = join(algo_folder, 'daily_perf')
folder = join(algo_folder, 'daily_performance')
files = [f for f in listdir(folder) if isfile(join(folder, f))]
daily_perf_list = []
for item in files:
filename = join(folder, item)
with open(filename, 'rb') as handle:
daily_perf_list.append(pickle.load(handle))
perf_period = pickle.load(handle)
perf_period_dict = perf_period.to_dict()
daily_perf_list.append(perf_period_dict)
stats = pd.DataFrame(daily_perf_list)
stats.set_index('period_close', drop=False, inplace=True)
self.analyze(stats)
@@ -460,43 +465,69 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
return self._clock
def get_generator(self):
if self.trading_client is not None:
return self.trading_client.transform()
def _init_trading_client(self):
"""
This replaces Ziplines `_create_generator` method. The main difference
is that we are restoring performance tracker objects if available.
This allows us to stop/start algos without loosing their state.
"""
self.state = get_algo_object(
algo_name=self.algo_namespace,
key='context.state',
)
if self.state is None:
self.state = {}
perf = None
if self.perf_tracker is None:
# Note from the Zipline dev:
# HACK: When running with the `run` method, we set perf_tracker to
# None so that it will be overwritten here.
tracker = self.perf_tracker = PerformanceTracker(
sim_params=self.sim_params,
trading_calendar=self.trading_calendar,
env=self.trading_environment,
)
# Set the dt initially to the period start by forcing it to change.
self.on_dt_changed(self.sim_params.start_session)
new_position_tracker = tracker.position_tracker
tracker.position_tracker = None
# Unpacking the perf_tracker and positions if available
perf = get_algo_object(
cum_perf = get_algo_object(
algo_name=self.algo_namespace,
key='cumulative_performance',
)
if cum_perf is not None:
tracker.cumulative_performance = cum_perf
# Ensure single common position tracker
tracker.position_tracker = cum_perf.position_tracker
today = pd.Timestamp.utcnow().floor('1D')
todays_perf = get_algo_object(
algo_name=self.algo_namespace,
key=today.strftime('%Y-%m-%d'),
rel_path='daily_performance',
)
if todays_perf is not None:
# Ensure single common position tracker
if tracker.position_tracker is not None:
todays_perf.position_tracker = tracker.position_tracker
else:
tracker.position_tracker = todays_perf.position_tracker
tracker.todays_performance = todays_perf
if tracker.position_tracker is None:
# Use a new position_tracker if not is found in the state
tracker.position_tracker = new_position_tracker
if not self.initialized:
# Calls the initialize function of the algorithm
self.initialize(*self.initialize_args, **self.initialize_kwargs)
self.initialized = True
# Call the simulation trading algorithm for side-effects:
# it creates the perf tracker
# TradingAlgorithm._create_generator(self, self.sim_params)
if perf is not None:
tracker.cumulative_performance = perf
period = self.perf_tracker.todays_performance
period.starting_cash = perf.ending_cash
period.starting_exposure = perf.ending_exposure
period.starting_value = perf.ending_value
period.position_tracker = perf.position_tracker
self.trading_client = ExchangeAlgorithmExecutor(
algo=self,
sim_params=self.sim_params,
@@ -506,6 +537,11 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
restrictions=self.restrictions,
universe_func=self._calculate_universe,
)
def get_generator(self):
if self.trading_client is None:
self._init_trading_client()
return self.trading_client.transform()
def updated_portfolio(self):
@@ -523,10 +559,6 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
positions, returning the available cash, and raising error
if the data goes out of sync.
Parameters
----------
attempt_index: int
Returns
-------
float
@@ -559,10 +591,19 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
if base_currency is None:
base_currency = exchange.base_currency
# Don't check the cash if there are open orders. This could
# results in false positives.
orders = []
for asset in self.blotter.open_orders:
asset_orders = self.blotter.open_orders[asset]
if asset_orders:
orders += asset_orders
required_cash = self.portfolio.cash if not orders else None
cash, positions_value = exchange.sync_positions(
positions=exchange_positions,
check_balances=check_balances,
cash=self.portfolio.cash,
cash=required_cash,
)
total_cash += cash
total_positions_value += positions_value
@@ -670,17 +711,22 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
if not self.is_running:
return
if self.end is not None and self.end < data.current_dt:
log.info('Algorithm has reached specified end time. Finishing...')
self.interrupt_algorithm()
# Resetting the frame stats every day to minimize memory footprint
today = data.current_dt.floor('1D')
if self.current_day is not None and today > self.current_day:
self.frame_stats = list()
self.performance_needs_update = False
new_orders = self.perf_tracker.todays_performance.orders_by_id.keys()
if new_orders != self._last_orders:
orders = list(self.perf_tracker.todays_performance.orders_by_id.keys())
if orders != self._last_orders:
self.performance_needs_update = True
self._last_orders = new_orders
# Saving current orders to detect changes in the next frame
self._last_orders = copy.deepcopy(orders)
if self.performance_needs_update:
self.perf_tracker.update_performance()
@@ -697,7 +743,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
self.portfolio_needs_update = False
log.info(
'got totals from exchanges, cash: {} positions: {}'.format(
'portfolio balances, cash: {}, positions: {}'.format(
cash, positions_value
)
)
@@ -709,18 +755,34 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
# every bar no matter if the algorithm places an order or not.
self.validate_account_controls()
self._save_algo_state(data)
self.current_day = data.current_dt.floor('1D')
def _save_algo_state(self, data):
today = data.current_dt.floor('1D')
try:
self._save_stats_csv(self._process_stats(data))
except Exception as e:
log.warn('unable to calculate performance: {}'.format(e))
log.debug('saving cumulative performance object')
save_algo_object(
algo_name=self.algo_namespace,
key='cumulative_performance',
obj=self.perf_tracker.cumulative_performance,
)
self.current_day = data.current_dt.floor('1D')
log.debug('saving todays performance object')
save_algo_object(
algo_name=self.algo_namespace,
key=today.strftime('%Y-%m-%d'),
obj=self.perf_tracker.todays_performance,
rel_path='daily_performance'
)
log.debug('saving context.state object')
save_algo_object(
algo_name=self.algo_namespace,
key='context.state',
obj=self.state)
def _process_stats(self, data):
today = data.current_dt.floor('1D')
@@ -763,12 +825,6 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
start_dt=today,
end_dt=data.current_dt
)
save_algo_object(
algo_name=self.algo_namespace,
key=today.strftime('%Y-%m-%d'),
obj=daily_stats,
rel_path='daily_perf'
)
return recorded_cols
+1 -2
View File
@@ -1,8 +1,7 @@
import pandas as pd
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.utils.factory import find_exchanges
from logbook import Logger
log = Logger('ExchangeAssetFinder', level=LOG_LEVEL)
+38 -38
View File
@@ -1,8 +1,9 @@
import numpy as np
import pandas as pd
from catalyst.assets._assets import TradingPair
from logbook import Logger
from redo import retry
from catalyst.assets._assets import TradingPair
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.exchange_errors import ExchangeRequestError
from catalyst.finance.blotter import Blotter
@@ -42,6 +43,11 @@ class TradingPairFeeSchedule(CommissionModel):
)
)
def get_maker_taker(self, asset):
maker = self.maker if self.maker is not None else asset.maker
taker = self.taker if self.taker is not None else asset.taker
return maker, taker
def calculate(self, order, transaction):
"""
Calculate the final fee based on the order parameters.
@@ -55,8 +61,7 @@ class TradingPairFeeSchedule(CommissionModel):
cost = abs(transaction.amount) * transaction.price
asset = order.asset
maker = self.maker if self.maker is not None else asset.maker
taker = self.taker if self.taker is not None else asset.taker
maker, taker = self.get_maker_taker(asset)
multiplier = taker
if order.limit is not None:
@@ -90,7 +95,6 @@ class TradingPairFixedSlippage(SlippageModel):
def simulate(self, data, asset, orders_for_asset):
self._volume_for_bar = 0
price = data.current(asset, 'close')
dt = data.current_dt
@@ -100,18 +104,20 @@ class TradingPairFixedSlippage(SlippageModel):
order.check_triggers(price, dt)
if not order.triggered:
log.debug('order has not reached the trigger at current '
'price {}'.format(price))
log.info(
'order has not reached the trigger at current '
'price {}'.format(price)
)
continue
execution_price, execution_volume = self.process_order(data, order)
if execution_price is not None:
transaction = create_transaction(
order, dt, execution_price, execution_volume
)
transaction = create_transaction(
order, dt, execution_price, execution_volume
)
self._volume_for_bar += abs(transaction.amount)
yield order, transaction
self._volume_for_bar += abs(transaction.amount)
yield order, transaction
def process_order(self, data, order):
price = data.current(order.asset, 'close')
@@ -202,34 +208,29 @@ class ExchangeBlotter(Blotter):
for order in self.open_orders[asset]:
log.debug('found open order: {}'.format(order.id))
new_order, executed_price = exchange.get_order(order.id, asset)
log.debug(
'got updated order {} {}'.format(
new_order, executed_price
transactions = exchange.process_order(order)
# This is a temporary measure, we should really update all
# trades, not just when the order gets filled. I just think
# that this is safer until we have a robust way to track
# the trades already processed by the algo. We can't loose
# them if the algo shuts down.
if transactions and order.open_amount == 0:
avg_price = np.average(
a=[t.price for t in transactions],
weights=[t.amount for t in transactions],
)
)
order.status = new_order.status
if order.status == ORDER_STATUS.FILLED:
order.commission = new_order.commission
if order.amount != new_order.amount:
log.warn(
'executed order amount {} differs '
'from original'.format(
new_order.amount, order.amount
)
ostatus = 'filled' if order.open_amount == 0 else 'partial'
log.info(
'{} order {} / {}: {}, avg price: {}'.format(
ostatus,
order.id,
asset.symbol,
order.filled,
avg_price,
)
order.amount = new_order.amount
transaction = Transaction(
asset=order.asset,
amount=order.amount,
dt=pd.Timestamp.utcnow(),
price=executed_price,
order_id=order.id,
commission=order.commission
)
yield order, transaction
for transaction in transactions:
yield order, transaction
elif order.status == ORDER_STATUS.CANCELLED:
yield order, None
@@ -250,7 +251,6 @@ class ExchangeBlotter(Blotter):
for order, txn in self.check_open_orders():
order.dt = txn.dt
transactions.append(txn)
if not order.open:
+25 -26
View File
@@ -8,12 +8,8 @@ from operator import is_not
import numpy as np
import pandas as pd
import pytz
from catalyst.assets._assets import TradingPair
from logbook import Logger
from pytz import UTC
from six import itervalues
from catalyst import get_calendar
from catalyst.assets._assets import TradingPair
from catalyst.constants import DATE_TIME_FORMAT, AUTO_INGEST
from catalyst.constants import LOG_LEVEL
from catalyst.data.minute_bars import BcolzMinuteOverlappingData, \
@@ -25,13 +21,16 @@ from catalyst.exchange.exchange_errors import EmptyValuesInBundleError, \
NoDataAvailableOnExchange, \
PricingDataNotLoadedError, DataCorruptionError, PricingDataValueError
from catalyst.exchange.utils.bundle_utils import range_in_bundle, \
get_bcolz_chunk, get_month_start_end, \
get_year_start_end, get_df_from_arrays, get_start_dt, get_period_label, \
get_delta, get_assets
get_bcolz_chunk, get_df_from_arrays, get_assets
from catalyst.exchange.utils.datetime_utils import get_delta, get_start_dt, \
get_period_label, get_month_start_end, get_year_start_end
from catalyst.exchange.utils.exchange_utils import get_exchange_folder, \
save_exchange_symbols, mixin_market_params, get_catalyst_symbol
from catalyst.utils.cli import maybe_show_progress
from catalyst.utils.paths import ensure_directory
from logbook import Logger
from pytz import UTC
from six import itervalues
log = Logger('exchange_bundle', level=LOG_LEVEL)
@@ -233,12 +232,12 @@ class ExchangeBundle:
problem = '{name} ({start_dt} to {end_dt}) has empty ' \
'periods: {dates}'.format(
name=asset.symbol,
start_dt=asset.start_date.strftime(
DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT),
dates=[date.strftime(
DATE_TIME_FORMAT) for date in dates])
name=asset.symbol,
start_dt=asset.start_date.strftime(
DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT),
dates=[date.strftime(
DATE_TIME_FORMAT) for date in dates])
if empty_rows_behavior == 'warn':
log.warn(problem)
@@ -287,12 +286,12 @@ class ExchangeBundle:
problem = '{name} ({start_dt} to {end_dt}) has {threshold} ' \
'identical close values on: {dates}'.format(
name=asset.symbol,
start_dt=asset.start_date.strftime(DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT),
threshold=threshold,
dates=[pd.to_datetime(date).strftime(DATE_TIME_FORMAT)
for date in dates])
name=asset.symbol,
start_dt=asset.start_date.strftime(DATE_TIME_FORMAT),
end_dt=end_dt.strftime(DATE_TIME_FORMAT),
threshold=threshold,
dates=[pd.to_datetime(date).strftime(DATE_TIME_FORMAT)
for date in dates])
problems.append(problem)
@@ -630,8 +629,8 @@ class ExchangeBundle:
show_progress,
label='Ingesting {frequency} price data on '
'{exchange}'.format(
exchange=self.exchange_name,
frequency=data_frequency,
exchange=self.exchange_name,
frequency=data_frequency,
)) as it:
for chunk in it:
problems += self.ingest_ctable(
@@ -965,15 +964,15 @@ class ExchangeBundle:
data_frequency,
trailing_bar_count=None,
reset_reader=False):
if trailing_bar_count:
delta = get_delta(trailing_bar_count, data_frequency)
end_dt += delta
start_dt = get_start_dt(end_dt, bar_count, data_frequency, False)
start_dt, _ = self.get_adj_dates(
start_dt, end_dt, assets, data_frequency
)
if trailing_bar_count:
delta = get_delta(trailing_bar_count, data_frequency)
end_dt += delta
# This is an attempt to resolve some caching with the reader
# when auto-ingesting data.
# TODO: needs more work
+5 -5
View File
@@ -3,17 +3,16 @@ import abc
import numpy as np
import pandas as pd
from catalyst.assets._assets import TradingPair
from logbook import Logger
from redo import retry
from catalyst.constants import LOG_LEVEL, AUTO_INGEST
from catalyst.data.data_portal import DataPortal
from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import (
ExchangeRequestError,
PricingDataNotLoadedError)
from catalyst.exchange.utils.exchange_utils import get_frequency, \
resample_history_df, group_assets_by_exchange
from catalyst.exchange.utils.exchange_utils import resample_history_df, group_assets_by_exchange
from catalyst.exchange.utils.datetime_utils import get_frequency
from logbook import Logger
from redo import retry
log = Logger('DataPortalExchange', level=LOG_LEVEL)
@@ -292,6 +291,7 @@ class DataPortalExchangeBacktest(DataPortalExchangeBase):
DataFrame
"""
# TODO: verify that the exchange supports the timeframe
bundle = self.exchange_bundles[exchange_name] # type: ExchangeBundle
freq, candle_size, unit, adj_data_frequency = get_frequency(
+7
View File
@@ -100,6 +100,13 @@ class InvalidHistoryFrequencyError(ZiplineError):
).strip()
class UnsupportedHistoryFrequencyError(ZiplineError):
msg = (
'{exchange} does not support candle frequency {freq}, please choose '
'from: {freqs}.'
).strip()
class InvalidHistoryTimeframeError(ZiplineError):
msg = (
'CCXT timeframe {timeframe} not supported by the exchange.'
+1 -2
View File
@@ -1,8 +1,7 @@
import numpy as np
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.protocol import Portfolio, Positions, Position
from logbook import Logger
log = Logger('ExchangePortfolio', level=LOG_LEVEL)
+5 -6
View File
@@ -11,12 +11,6 @@
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
from logbook import Logger
from numpy import (
iinfo,
uint32,
)
from catalyst.constants import LOG_LEVEL
from catalyst.data.us_equity_pricing import BcolzDailyBarReader
from catalyst.errors import NoFurtherDataError
@@ -26,6 +20,11 @@ from catalyst.pipeline.data import DataSet, Column
from catalyst.pipeline.loaders.base import PipelineLoader
from catalyst.utils.calendars import get_calendar
from catalyst.utils.numpy_utils import float64_dtype
from logbook import Logger
from numpy import (
iinfo,
uint32,
)
UINT32_MAX = iinfo(uint32).max
+2 -3
View File
@@ -1,13 +1,12 @@
import pandas as pd
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.utils.stats_utils import prepare_stats
from catalyst.gens.sim_engine import (
BAR,
SESSION_START
)
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.utils.stats_utils import prepare_stats
log = Logger('LiveGraphClock', level=LOG_LEVEL)
+1 -2
View File
@@ -14,14 +14,13 @@
from time import sleep
import pandas as pd
from catalyst.constants import LOG_LEVEL
from catalyst.gens.sim_engine import (
BAR,
SESSION_START
)
from logbook import Logger
from catalyst.constants import LOG_LEVEL
log = Logger('ExchangeClock', level=LOG_LEVEL)
+12 -212
View File
@@ -1,11 +1,18 @@
import calendar
import os
import tarfile
from datetime import timedelta, datetime, date
from datetime import datetime
import numpy as np
import pandas as pd
from catalyst.data.bundles.core import download_without_progress
from catalyst.exchange.utils.exchange_utils import get_exchange_bundles_folder
import os
import tarfile
from datetime import datetime
import numpy as np
import pandas as pd
import pytz
from catalyst.data.bundles.core import download_without_progress
from catalyst.exchange.utils.exchange_utils import get_exchange_bundles_folder
@@ -14,41 +21,6 @@ EXCHANGE_NAMES = ['bitfinex', 'bittrex', 'poloniex']
API_URL = 'http://data.enigma.co/api/v1'
def get_date_from_ms(ms):
"""
The date from the number of miliseconds from the epoch.
Parameters
----------
ms: int
Returns
-------
datetime
"""
return datetime.fromtimestamp(ms / 1000.0)
def get_seconds_from_date(date):
"""
The number of seconds from the epoch.
Parameters
----------
date: datetime
Returns
-------
int
"""
epoch = datetime.utcfromtimestamp(0)
epoch = epoch.replace(tzinfo=pytz.UTC)
return int((date - epoch).total_seconds())
def get_bcolz_chunk(exchange_name, symbol, data_frequency, period):
"""
Download and extract a bcolz bundle.
@@ -78,8 +50,8 @@ def get_bcolz_chunk(exchange_name, symbol, data_frequency, period):
if not os.path.isdir(path):
url = 'https://s3.amazonaws.com/enigmaco/catalyst-bundles/' \
'exchange-{exchange}/{name}.tar.gz'.format(
exchange=exchange_name,
name=name)
exchange=exchange_name,
name=name)
bytes = download_without_progress(url)
with tarfile.open('r', fileobj=bytes) as tar:
@@ -88,178 +60,6 @@ def get_bcolz_chunk(exchange_name, symbol, data_frequency, period):
return path
def get_delta(periods, data_frequency):
"""
Get a time delta based on the specified data frequency.
Parameters
----------
periods: int
data_frequency: str
Returns
-------
timedelta
"""
return timedelta(minutes=periods) \
if data_frequency == 'minute' else timedelta(days=periods)
def get_periods_range(start_dt, end_dt, freq):
"""
Get a date range for the specified parameters.
Parameters
----------
start_dt: datetime
end_dt: datetime
freq: str
Returns
-------
DateTimeIndex
"""
if freq == 'minute':
freq = 'T'
elif freq == 'daily':
freq = 'D'
return pd.date_range(start_dt, end_dt, freq=freq)
def get_periods(start_dt, end_dt, freq):
"""
The number of periods in the specified range.
Parameters
----------
start_dt: datetime
end_dt: datetime
freq: str
Returns
-------
int
"""
return len(get_periods_range(start_dt, end_dt, freq))
def get_start_dt(end_dt, bar_count, data_frequency, include_first=True):
"""
The start date based on specified end date and data frequency.
Parameters
----------
end_dt: datetime
bar_count: int
data_frequency: str
Returns
-------
datetime
"""
periods = bar_count
if periods > 1:
delta = get_delta(periods, data_frequency)
start_dt = end_dt - delta
if not include_first:
start_dt += get_delta(1, data_frequency)
else:
start_dt = end_dt
return start_dt
def get_period_label(dt, data_frequency):
"""
The period label for the specified date and frequency.
Parameters
----------
dt: datetime
data_frequency: str
Returns
-------
str
"""
if data_frequency == 'minute':
return '{}-{:02d}'.format(dt.year, dt.month)
else:
return '{}'.format(dt.year)
def get_month_start_end(dt, first_day=None, last_day=None):
"""
The first and last day of the month for the specified date.
Parameters
----------
dt: datetime
first_day: datetime
last_day: datetime
Returns
-------
datetime, datetime
"""
month_range = calendar.monthrange(dt.year, dt.month)
if first_day:
month_start = first_day
else:
month_start = pd.to_datetime(datetime(
dt.year, dt.month, 1, 0, 0, 0, 0
), utc=True)
if last_day:
month_end = last_day
else:
month_end = pd.to_datetime(datetime(
dt.year, dt.month, month_range[1], 23, 59, 0, 0
), utc=True)
if month_end > pd.Timestamp.utcnow():
month_end = pd.Timestamp.utcnow().floor('1D')
return month_start, month_end
def get_year_start_end(dt, first_day=None, last_day=None):
"""
The first and last day of the year for the specified date.
Parameters
----------
dt: datetime
first_day: datetime
last_day: datetime
Returns
-------
datetime, datetime
"""
year_start = first_day if first_day \
else pd.to_datetime(date(dt.year, 1, 1), utc=True)
year_end = last_day if last_day \
else pd.to_datetime(date(dt.year, 12, 31), utc=True)
if year_end > pd.Timestamp.utcnow():
year_end = pd.Timestamp.utcnow().floor('1D')
return year_start, year_end
def get_df_from_arrays(arrays, periods):
"""
A DataFrame from the specified OHCLV arrays.
+327
View File
@@ -0,0 +1,327 @@
import calendar
import re
from datetime import datetime, timedelta, date
import pandas as pd
import pytz
from catalyst.exchange.exchange_errors import InvalidHistoryFrequencyError, \
InvalidHistoryFrequencyAlias
def get_date_from_ms(ms):
"""
The date from the number of miliseconds from the epoch.
Parameters
----------
ms: int
Returns
-------
datetime
"""
return datetime.fromtimestamp(ms / 1000.0)
def get_seconds_from_date(date):
"""
The number of seconds from the epoch.
Parameters
----------
date: datetime
Returns
-------
int
"""
epoch = datetime.utcfromtimestamp(0)
epoch = epoch.replace(tzinfo=pytz.UTC)
return int((date - epoch).total_seconds())
def get_delta(periods, data_frequency):
"""
Get a time delta based on the specified data frequency.
Parameters
----------
periods: int
data_frequency: str
Returns
-------
timedelta
"""
return timedelta(minutes=periods) \
if data_frequency == 'minute' else timedelta(days=periods)
def get_periods_range(freq, start_dt=None, end_dt=None, periods=None):
"""
Get a date range for the specified parameters.
Parameters
----------
start_dt: datetime
end_dt: datetime
freq: str
Returns
-------
DateTimeIndex
"""
if freq == 'minute':
freq = 'T'
elif freq == 'daily':
freq = 'D'
if start_dt is not None and end_dt is not None and periods is None:
return pd.date_range(start_dt, end_dt, freq=freq)
elif periods is not None and (start_dt is not None or end_dt is not None):
_, unit_periods, unit, _ = get_frequency(freq)
adj_periods = periods * unit_periods
# TODO: standardize time aliases to avoid any mapping
unit = 'd' if unit == 'D' else 'm'
delta = pd.Timedelta(adj_periods, unit)
if start_dt is not None:
return pd.date_range(
start=start_dt,
end=start_dt + delta,
freq=freq,
closed='left',
)
else:
return pd.date_range(
start=end_dt - delta,
end=end_dt,
freq=freq,
)
else:
raise ValueError(
'Choose only two parameters between start_dt, end_dt '
'and periods.'
)
def get_periods(start_dt, end_dt, freq):
"""
The number of periods in the specified range.
Parameters
----------
start_dt: datetime
end_dt: datetime
freq: str
Returns
-------
int
"""
return len(get_periods_range(start_dt=start_dt, end_dt=end_dt, freq=freq))
def get_start_dt(end_dt, bar_count, data_frequency, include_first=True):
"""
The start date based on specified end date and data frequency.
Parameters
----------
end_dt: datetime
bar_count: int
data_frequency: str
include_first
Returns
-------
datetime
"""
periods = bar_count
if periods > 1:
delta = get_delta(periods, data_frequency)
start_dt = end_dt - delta
if not include_first:
start_dt += get_delta(1, data_frequency)
else:
start_dt = end_dt
return start_dt
def get_period_label(dt, data_frequency):
"""
The period label for the specified date and frequency.
Parameters
----------
dt: datetime
data_frequency: str
Returns
-------
str
"""
if data_frequency == 'minute':
return '{}-{:02d}'.format(dt.year, dt.month)
else:
return '{}'.format(dt.year)
def get_month_start_end(dt, first_day=None, last_day=None):
"""
The first and last day of the month for the specified date.
Parameters
----------
dt: datetime
first_day: datetime
last_day: datetime
Returns
-------
datetime, datetime
"""
month_range = calendar.monthrange(dt.year, dt.month)
if first_day:
month_start = first_day
else:
month_start = pd.to_datetime(datetime(
dt.year, dt.month, 1, 0, 0, 0, 0
), utc=True)
if last_day:
month_end = last_day
else:
month_end = pd.to_datetime(datetime(
dt.year, dt.month, month_range[1], 23, 59, 0, 0
), utc=True)
if month_end > pd.Timestamp.utcnow():
month_end = pd.Timestamp.utcnow().floor('1D')
return month_start, month_end
def get_year_start_end(dt, first_day=None, last_day=None):
"""
The first and last day of the year for the specified date.
Parameters
----------
dt: datetime
first_day: datetime
last_day: datetime
Returns
-------
datetime, datetime
"""
year_start = first_day if first_day \
else pd.to_datetime(date(dt.year, 1, 1), utc=True)
year_end = last_day if last_day \
else pd.to_datetime(date(dt.year, 12, 31), utc=True)
if year_end > pd.Timestamp.utcnow():
year_end = pd.Timestamp.utcnow().floor('1D')
return year_start, year_end
def get_frequency(freq, data_frequency=None):
"""
Get the frequency parameters.
Notes
-----
We're trying to use Pandas convention for frequency aliases.
Parameters
----------
freq: str
data_frequency: str
Returns
-------
str, int, str, str
"""
if data_frequency is None:
data_frequency = 'daily' if freq.upper().endswith('D') else 'minute'
if freq == 'minute':
unit = 'T'
candle_size = 1
elif freq == 'daily':
unit = 'D'
candle_size = 1
else:
freq_match = re.match(r'([0-9].*)?(m|M|d|D|h|H|T)', freq, re.M | re.I)
if freq_match:
candle_size = int(freq_match.group(1)) if freq_match.group(1) \
else 1
unit = freq_match.group(2)
else:
raise InvalidHistoryFrequencyError(frequency=freq)
# TODO: some exchanges support H and W frequencies but not bundles
# Find a way to pass-through these parameters to exchanges
# but resample from minute or daily in backtest mode
# see catalyst/exchange/ccxt/ccxt_exchange.py:242 for mapping between
# Pandas offet aliases (used by Catalyst) and the CCXT timeframes
if unit.lower() == 'd':
unit = 'D'
alias = '{}D'.format(candle_size)
if data_frequency == 'minute':
data_frequency = 'daily'
elif unit.lower() == 'm' or unit == 'T':
unit = 'T'
alias = '{}T'.format(candle_size)
if data_frequency == 'daily':
data_frequency = 'minute'
# elif unit.lower() == 'h':
# candle_size = candle_size * 60
#
# alias = '{}T'.format(candle_size)
# if data_frequency == 'daily':
# data_frequency = 'minute'
else:
raise InvalidHistoryFrequencyAlias(freq=freq)
return alias, candle_size, unit, data_frequency
def from_ms_timestamp(ms):
return pd.to_datetime(ms, unit='ms', utc=True)
def get_epoch():
return pd.to_datetime('1970-1-1', utc=True)
+8 -80
View File
@@ -2,7 +2,6 @@ import hashlib
import json
import os
import pickle
import re
import shutil
from datetime import date, datetime
@@ -12,8 +11,7 @@ from six import string_types
from six.moves.urllib import request
from catalyst.constants import DATE_FORMAT, SYMBOLS_URL
from catalyst.exchange.exchange_errors import ExchangeSymbolsNotFound, \
InvalidHistoryFrequencyError, InvalidHistoryFrequencyAlias
from catalyst.exchange.exchange_errors import ExchangeSymbolsNotFound
from catalyst.exchange.utils.serialization_utils import ExchangeJSONEncoder, \
ExchangeJSONDecoder
from catalyst.utils.paths import data_root, ensure_directory, \
@@ -130,7 +128,10 @@ def get_exchange_symbols(exchange_name, is_local=False, environ=None):
if not is_local and (not os.path.isfile(filename) or pd.Timedelta(
pd.Timestamp('now', tz='UTC') - last_modified_time(
filename)).days > 1):
download_exchange_symbols(exchange_name, environ)
try:
download_exchange_symbols(exchange_name, environ)
except Exception as e:
pass
if os.path.isfile(filename):
with open(filename) as data_file:
@@ -190,7 +191,7 @@ def get_symbols_string(assets):
return ', '.join([asset.symbol for asset in array])
def get_exchange_auth(exchange_name, environ=None):
def get_exchange_auth(exchange_name, alias=None, environ=None):
"""
The de-serialized contend of the exchange's auth.json file.
@@ -205,7 +206,8 @@ def get_exchange_auth(exchange_name, environ=None):
"""
exchange_folder = get_exchange_folder(exchange_name, environ)
filename = os.path.join(exchange_folder, 'auth.json')
name = 'auth' if alias is None else alias
filename = os.path.join(exchange_folder, '{}.json'.format(name))
if os.path.isfile(filename):
with open(filename) as data_file:
@@ -510,72 +512,6 @@ def get_common_assets(exchanges):
return assets
def get_frequency(freq, data_frequency):
"""
Get the frequency parameters.
Notes
-----
We're trying to use Pandas convention for frequency aliases.
Parameters
----------
freq: str
data_frequency: str
Returns
-------
str, int, str, str
"""
if freq == 'minute':
unit = 'T'
candle_size = 1
elif freq == 'daily':
unit = 'D'
candle_size = 1
else:
freq_match = re.match(r'([0-9].*)?(m|M|d|D|h|H|T)', freq, re.M | re.I)
if freq_match:
candle_size = int(freq_match.group(1)) if freq_match.group(1) \
else 1
unit = freq_match.group(2)
else:
raise InvalidHistoryFrequencyError(frequency=freq)
# TODO: some exchanges support H and W frequencies but not bundles
# Find a way to pass-through these parameters to exchanges
# but resample from minute or daily in backtest mode
# see catalyst/exchange/ccxt/ccxt_exchange.py:242 for mapping between
# Pandas offet aliases (used by Catalyst) and the CCXT timeframes
if unit.lower() == 'd':
alias = '{}D'.format(candle_size)
if data_frequency == 'minute':
data_frequency = 'daily'
elif unit.lower() == 'm' or unit == 'T':
alias = '{}T'.format(candle_size)
if data_frequency == 'daily':
data_frequency = 'minute'
# elif unit.lower() == 'h':
# candle_size = candle_size * 60
#
# alias = '{}T'.format(candle_size)
# if data_frequency == 'daily':
# data_frequency = 'minute'
else:
raise InvalidHistoryFrequencyAlias(freq=freq)
return alias, candle_size, unit, data_frequency
def resample_history_df(df, freq, field):
"""
Resample the OHCLV DataFrame using the specified frequency.
@@ -649,14 +585,6 @@ def mixin_market_params(exchange_name, params, market):
params['lot'] = params['min_trade_size']
def from_ms_timestamp(ms):
return pd.to_datetime(ms, unit='ms', utc=True)
def get_epoch():
return pd.to_datetime('1970-1-1', utc=True)
def group_assets_by_exchange(assets):
exchange_assets = dict()
for asset in assets:
+3 -4
View File
@@ -1,25 +1,24 @@
import os
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.ccxt.ccxt_exchange import CCXT
from catalyst.exchange.exchange import Exchange
from catalyst.exchange.exchange_errors import ExchangeAuthEmpty
from catalyst.exchange.utils.exchange_utils import get_exchange_auth, \
get_exchange_folder, is_blacklist
from logbook import Logger
log = Logger('factory', level=LOG_LEVEL)
exchange_cache = dict()
def get_exchange(exchange_name, base_currency=None, must_authenticate=False,
skip_init=False):
skip_init=False, auth_alias=None):
key = (exchange_name, base_currency)
if key in exchange_cache:
return exchange_cache[key]
exchange_auth = get_exchange_auth(exchange_name)
exchange_auth = get_exchange_auth(exchange_name, alias=auth_alias)
has_auth = (exchange_auth['key'] != '' and exchange_auth['secret'] != '')
if must_authenticate and not has_auth:
@@ -3,9 +3,8 @@ import re
from json import JSONEncoder
import pandas as pd
from six import string_types
from catalyst.constants import DATE_TIME_FORMAT
from six import string_types
class ExchangeJSONEncoder(json.JSONEncoder):
+35 -13
View File
@@ -8,9 +8,9 @@ import time
import numpy as np
import pandas as pd
from catalyst.assets._assets import TradingPair
from catalyst.exchange.utils.exchange_utils import get_algo_folder
from catalyst.utils.paths import data_root, ensure_directory
from operator import itemgetter
s3_conn = []
mailgun = []
@@ -261,7 +261,14 @@ def prepare_stats(stats, recorded_cols=list()):
return df, columns
def get_pretty_stats(stats, recorded_cols=None, num_rows=10):
def set_print_settings():
pd.set_option('display.expand_frame_repr', False)
pd.set_option('precision', 8)
pd.set_option('display.width', 1000)
pd.set_option('display.max_colwidth', 1000)
def get_pretty_stats(stats, recorded_cols=None, num_rows=10, show_tail=True):
"""
Format and print the last few rows of a statistics DataFrame.
See the pyfolio project for the data structure.
@@ -280,18 +287,18 @@ def get_pretty_stats(stats, recorded_cols=None, num_rows=10):
"""
if isinstance(stats, pd.DataFrame):
stats = stats.T.to_dict().values()
stats = list(stats.T.to_dict().values())
stats.sort(key=itemgetter('period_close'))
if len(stats) > num_rows:
display_stats = stats[-num_rows:] if show_tail else stats[0:num_rows]
else:
display_stats = stats
display_stats = stats[-num_rows:] if len(stats) > num_rows else stats
df, columns = prepare_stats(
display_stats, recorded_cols=recorded_cols
)
pd.set_option('display.expand_frame_repr', False)
pd.set_option('precision', 8)
pd.set_option('display.width', 1000)
pd.set_option('display.max_colwidth', 1000)
set_print_settings()
return df.to_string(columns=columns)
@@ -352,9 +359,13 @@ def stats_to_s3(uri, stats, algo_namespace, recorded_cols=None,
pid = os.getpid()
parts = uri.split('//')
obj = s3.Object(parts[1], '{}/{}-{}-{}.csv'.format(
folder, timestr, algo_namespace, pid
))
path = '{folder}/{algo}/{time}-{algo}-{pid}.csv'.format(
folder=folder,
algo=algo_namespace,
time=timestr,
pid=pid,
)
obj = s3.Object(parts[1], path)
obj.put(Body=bytes_to_write)
@@ -439,6 +450,17 @@ def df_to_string(df):
return df.to_string()
def extract_orders(perf):
order_list = perf.orders.values
all_orders = [t for sublist in order_list for t in sublist]
all_orders.sort(key=lambda o: o['dt'])
orders = pd.DataFrame(all_orders)
if not orders.empty:
orders.set_index('dt', inplace=True, drop=True)
return orders
def extract_transactions(perf):
"""
Compute indexes for buy and sell transactions
+6 -7
View File
@@ -3,7 +3,6 @@ import random
import tempfile
from catalyst.assets._assets import TradingPair
from catalyst.exchange.utils.exchange_utils import get_exchange_folder
from catalyst.exchange.utils.factory import find_exchanges
from catalyst.utils.paths import ensure_directory
@@ -63,14 +62,14 @@ def output_df(df, assets, name=None):
"""
if isinstance(assets, TradingPair):
exchange_folder = assets.exchange
asset_folder = assets.symbol
asset_folder = '{}_{}'.format(assets.exchange, assets.symbol)
else:
exchange_folder = ','.join([asset.exchange for asset in assets])
asset_folder = ','.join([asset.symbol for asset in assets])
asset_folder = ','.join(
['{}_{}'.format(a.exchange, a.symbol) for a in assets]
)
folder = os.path.join(
tempfile.gettempdir(), 'catalyst', exchange_folder, asset_folder
tempfile.gettempdir(), 'catalyst', asset_folder
)
ensure_directory(folder)
@@ -80,4 +79,4 @@ def output_df(df, assets, name=None):
path = os.path.join(folder, '{}.csv'.format(name))
df.to_csv(path)
return path
return path, folder
+2 -2
View File
@@ -142,7 +142,7 @@ class TermGraph(object):
at the end of execution.
"""
refcounts = self.graph.out_degree()
for t in self.outputs.values():
for t in list(self.outputs.values()):
refcounts[t] += 1
for t in initial_terms:
@@ -238,7 +238,7 @@ class ExecutionPlan(TermGraph):
min_extra_rows=0):
super(ExecutionPlan, self).__init__(terms)
for term in terms.values():
for term in list(terms.values()):
self.set_extra_rows(
term,
all_dates,
+2 -2
View File
@@ -144,7 +144,7 @@ class SpecificEquityTrades(object):
for identifier in self.identifiers:
assets_by_identifier[identifier] = env.asset_finder.\
lookup_generic(identifier, datetime.now())[0]
self.sids = [asset.sid for asset in assets_by_identifier.values()]
self.sids = [asset.sid for asset in list(assets_by_identifier.values())]
for event in self.event_list:
event.sid = assets_by_identifier[event.sid].sid
@@ -167,7 +167,7 @@ class SpecificEquityTrades(object):
for identifier in self.identifiers:
assets_by_identifier[identifier] = env.asset_finder.\
lookup_generic(identifier, datetime.now())[0]
self.sids = [asset.sid for asset in assets_by_identifier.values()]
self.sids = [asset.sid for asset in list(assets_by_identifier.values())]
# Hash_value for downstream sorting.
self.arg_string = hash_args(*args, **kwargs)
+57
View File
@@ -0,0 +1,57 @@
import pandas as pd
from catalyst import run_algorithm
def initialize(context):
context.i = -1 # counts the minutes
context.exchange = 'cryptopia'
context.base_currency = 'btc'
context.coins = context.exchanges[context.exchange].assets
context.coins = [c for c in context.coins if
c.quote_currency == context.base_currency]
def handle_data(context, data):
# current date formatted into a string
today = data.current_dt
# update universe everyday
new_day = 60 * 24 # assuming data_frequency='minute'
if not context.i % new_day:
context.coins = context.exchanges[context.exchange].assets
context.coins = [c for c in context.coins if
c.quote_currency == context.base_currency]
# get data every 30 minutes
minutes = 1
if not context.i % minutes:
# we iterate for every pair in the current universe
for coin in context.coins:
pair = str(coin.symbol)
price = data.current(coin, 'price')
print(today, pair, price)
def analyze(context=None, results=None):
pass
if __name__ == '__main__':
start_date = pd.to_datetime('2018-01-17', utc=True)
end_date = pd.to_datetime('2018-01-18', utc=True)
performance = run_algorithm(
capital_base=1.0,
# amount of base_currency, not always in dollars unless usd
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='cryptopia',
data_frequency='minute',
base_currency='btc',
live=True,
live_graph=False,
simulate_orders=True,
algo_namespace='simple_universe'
)
+8
View File
@@ -0,0 +1,8 @@
import ccxt
bitfinex = ccxt.bitfinex()
bitfinex.verbose = True
ohlcvs = bitfinex.fetch_ohlcv('ETH/BTC', '30m', 1504224000000)
dt = bitfinex.iso8601(ohlcvs[0][0])
print(dt) # should print '2017-09-01T00:00:00.000Z'
@@ -0,0 +1,50 @@
import pandas as pd
from catalyst import run_algorithm
from catalyst.api import symbol
def initialize(context):
context.asset1 = symbol('fct_btc')
context.asset2 = symbol('btc_usdt')
context.coins = [context.asset1, context.asset2]
def handle_data(context, data):
df = data.history(context.coins,
'close',
bar_count=10,
frequency='5T',
)
print(df)
print(data.current(context.asset1, 'close'))
print(data.current(context.asset2, 'close'))
exit(0)
if __name__ == '__main__':
LIVE = True
if LIVE:
run_algorithm(
capital_base=1,
initialize=initialize,
handle_data=handle_data,
exchange_name='poloniex',
algo_namespace='test_multi_assets',
base_currency='usdt',
live=True,
simulate_orders=True,
)
else:
run_algorithm(
capital_base=1,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
exchange_name='poloniex',
algo_namespace='test_multi_assets',
base_currency='usdt',
live=False,
start=pd.to_datetime('2017-12-1', utc=True),
end=pd.to_datetime('2017-12-1', utc=True),
)
+44
View File
@@ -0,0 +1,44 @@
from logbook import Logger
from catalyst import run_algorithm
from catalyst.api import order_target_percent
NAMESPACE = 'goose7'
log = Logger(NAMESPACE)
from catalyst.api import record, symbol
def initialize(context):
context.asset = symbol('trx_btc')
def handle_data(context, data):
price = data.current(context.asset, 'price')
record(btc=price)
# Only ordering if it does not have any position to avoid trying some
# tiny orders with the leftover btc
pos_amount = context.portfolio.positions[context.asset].amount
if pos_amount > 0:
return
# Adding a limit price to workaround an issue with performance
# calculations of market orders
order_target_percent(
context.asset, 1, limit_price=price * 1.01
)
if __name__ == '__main__':
run_algorithm(
capital_base=0.003,
initialize=initialize,
handle_data=handle_data,
exchange_name='binance',
live=True,
algo_namespace=NAMESPACE,
base_currency='btc',
live_graph=False,
simulate_orders=False,
)
+44
View File
@@ -0,0 +1,44 @@
import pandas as pd
from catalyst import run_algorithm
from catalyst.api import symbol
def initialize(context):
context.asset = symbol('btc_usdt')
def handle_data(context, data):
df = data.history(context.asset,
'close',
bar_count=10,
frequency='5T',
)
if __name__ == '__main__':
LIVE = True
if LIVE:
run_algorithm(
capital_base=1,
initialize=initialize,
handle_data=handle_data,
exchange_name='poloniex',
algo_namespace='test_algo',
base_currency='usdt',
live=True,
simulate_orders=True,
)
else:
run_algorithm(
capital_base=1,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
exchange_name='poloniex',
algo_namespace='test_algo',
base_currency='usdt',
live=False,
start=pd.to_datetime('2017-12-1', utc=True),
end=pd.to_datetime('2017-12-1', utc=True),
)
+44
View File
@@ -0,0 +1,44 @@
import pandas as pd
from catalyst.utils.run_algo import run_algorithm
from catalyst.api import symbol
from exchange.utils.stats_utils import set_print_settings
def initialize(context):
context.i = 0
context.data = []
def handle_data(context, data):
prices = data.history(
symbol('xlm_eth'),
fields=['open', 'high', 'low', 'close'],
bar_count=50,
frequency='1T'
)
set_print_settings()
print(prices.tail(10))
context.data.append(prices)
context.i = context.i + 1
if context.i == 3:
context.interrupt_algorithm()
def analyze(context, prefs):
for dataset in context.data:
print(dataset[-2:])
if __name__ == '__main__':
run_algorithm(
capital_base=0.1,
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='binance',
algo_namespace='Test candles',
base_currency='eth',
data_frequency='minute',
live=True,
simulate_orders=True)
+196 -265
View File
@@ -8,13 +8,14 @@ from time import sleep
import click
import pandas as pd
from logbook import Logger
from six import string_types
from catalyst.data.bundles import load
from catalyst.data.data_portal import DataPortal
from catalyst.exchange.exchange_pricing_loader import ExchangePricingLoader, \
TradingPairPricing
from catalyst.exchange.utils.factory import get_exchange
from logbook import Logger
try:
from pygments import highlight
@@ -40,9 +41,6 @@ from catalyst.exchange.exchange_algorithm import (
from catalyst.exchange.exchange_data_portal import DataPortalExchangeLive, \
DataPortalExchangeBacktest
from catalyst.exchange.exchange_asset_finder import ExchangeAssetFinder
from catalyst.exchange.exchange_errors import (
ExchangeRequestError, ExchangeRequestErrorTooManyAttempts,
BaseCurrencyNotFoundError, NotEnoughCapitalError)
from catalyst.constants import LOG_LEVEL
@@ -70,7 +68,38 @@ class _RunAlgoError(click.ClickException, ValueError):
return self.pyfunc_msg
def _build_namespace(algotext, local_namespace, defines):
def _run(handle_data,
initialize,
before_trading_start,
analyze,
algofile,
algotext,
defines,
data_frequency,
capital_base,
data,
bundle,
bundle_timestamp,
start,
end,
output,
print_algo,
local_namespace,
environ,
live,
exchange,
algo_namespace,
base_currency,
live_graph,
analyze_live,
simulate_orders,
auth_aliases,
stats_output):
"""Run a backtest for the given algorithm.
This is shared between the cli and :func:`catalyst.run_algo`.
"""
# TODO: refactor for more granularity
if algotext is not None:
if local_namespace:
ip = get_ipython() # noqa
@@ -84,197 +113,145 @@ def _build_namespace(algotext, local_namespace, defines):
except ValueError:
raise ValueError(
'invalid define %r, should be of the form name=value' %
assign)
assign,
)
try:
# evaluate in the same namespace so names may refer to
# eachother
namespace[name] = eval(value, namespace)
except Exception as e:
raise ValueError(
'failed to execute definition for name %r: %s' % (name, e))
'failed to execute definition for name %r: %s' % (name, e),
)
elif defines:
raise _RunAlgoError(
'cannot pass define without `algotext`',
"cannot pass '-D' / '--define' without '-t' / '--algotext'")
"cannot pass '-D' / '--define' without '-t' / '--algotext'",
)
else:
namespace = {}
if algofile is not None:
algotext = algofile.read()
return namespace
if print_algo:
if PYGMENTS:
highlight(
algotext,
PythonLexer(),
TerminalFormatter(),
outfile=sys.stdout,
)
else:
click.echo(algotext)
log.warn(
'Catalyst is currently in ALPHA. It is going through rapid '
'development and it is subject to errors. Please use carefully. '
'We encourage you to report any issue on GitHub: '
'https://github.com/enigmampc/catalyst/issues'
)
sleep(3)
def _mode(simulate_orders, live):
if not live:
return 'backtest'
elif simulate_orders:
return 'paper-trading'
if live:
if simulate_orders:
mode = 'paper-trading'
else:
mode = 'live-trading'
else:
return 'live-trading'
mode = 'backtest'
log.info('running algo in {mode} mode'.format(mode=mode))
def _build_exchanges_dict(exchange, live, simulate_orders, base_currency):
exchange_name = exchange
if exchange_name is None:
raise ValueError('Please specify at least one exchange.')
if isinstance(auth_aliases, string_types):
aliases = auth_aliases.split(',')
if len(aliases) < 2 or len(aliases) % 2 != 0:
raise ValueError(
'the `auth_aliases` parameter must contain an even list '
'of comma-delimited values. For example, '
'"binance,auth2" or "binance,auth2,bittrex,auth2".'
)
auth_aliases = dict(zip(aliases[::2], aliases[1::2]))
exchange_list = [x.strip().lower() for x in exchange.split(',')]
exchanges = {exchange_name: get_exchange(
exchange_name=exchange_name,
base_currency=base_currency,
must_authenticate=(live and not simulate_orders))
for exchange_name in exchange_list}
return exchanges
def _pretty_print_code(algotext):
if PYGMENTS:
highlight(
algotext,
PythonLexer(),
TerminalFormatter(),
outfile=sys.stdout)
else:
click.echo(algotext)
def _choose_loader(data_frequency, column):
bound_cols = TradingPairPricing.columns
if column in bound_cols:
return ExchangePricingLoader(data_frequency)
raise ValueError(
"No PipelineLoader registered for column %s." % column)
def _get_live_time_range():
start = pd.Timestamp.utcnow()
# TODO: fix the end data.
end = start + timedelta(hours=8760)
return start, end
def _data_for_live_trading(sim_params, exchanges, env, open_calendar):
data = DataPortalExchangeLive(
exchanges=exchanges,
asset_finder=env.asset_finder,
trading_calendar=open_calendar,
first_trading_day=pd.to_datetime('today', utc=True))
return data
# TODO use proper retry here
def _fetch_capital_base(base_currency, exchange_name, exchange,
attempt_index=0):
"""
Fetch the base currency amount required to bootstrap
the algorithm against the exchange.
The algorithm cannot continue without this value.
:param exchange: the targeted exchange
:param attempt_index:
:return capital_base: the amount of base currency available for
trading
"""
try:
log.debug('retrieving capital base in {} to bootstrap '
'exchange {}'.format(base_currency, exchange_name))
balances = exchange.get_balances()
except ExchangeRequestError as e:
if attempt_index < 20:
log.warn(
'could not retrieve balances on {}: {}'.format(
exchange.name, e))
sleep(5)
return _fetch_capital_base(base_currency, exchange_name, exchange,
attempt_index + 1)
exchanges = dict()
for name in exchange_list:
if auth_aliases is not None and name in auth_aliases:
auth_alias = auth_aliases[name]
else:
raise ExchangeRequestErrorTooManyAttempts(
attempts=attempt_index,
error=e)
auth_alias = None
if base_currency in balances:
base_currency_available = balances[base_currency]['free']
log.info(
'base currency available in the account: {} {}'.format(
base_currency_available, base_currency))
return base_currency_available
else:
raise BaseCurrencyNotFoundError(
exchanges[name] = get_exchange(
exchange_name=name,
base_currency=base_currency,
exchange=exchange_name)
must_authenticate=(live and not simulate_orders),
skip_init=True,
auth_alias=auth_alias,
)
open_calendar = get_calendar('OPEN')
def _algorithm_class_for_live(algo_namespace, live_graph, stats_output,
analyze_live, base_currency, simulate_orders,
exchanges, capital_base):
if not simulate_orders:
for exchange_name in exchanges:
exchange = exchanges[exchange_name]
balance = _fetch_capital_base(base_currency, exchange_name,
exchange)
env = TradingEnvironment(
load=partial(
load_crypto_market_data,
environ=environ,
start_dt=start,
end_dt=end
),
environ=environ,
exchange_tz='UTC',
asset_db_path=None # We don't need an asset db, we have exchanges
)
env.asset_finder = ExchangeAssetFinder(exchanges=exchanges)
if balance < capital_base:
raise NotEnoughCapitalError(
exchange=exchange_name,
base_currency=base_currency,
balance=balance,
capital_base=capital_base)
algorithm_class = partial(
ExchangeTradingAlgorithmLive,
exchanges=exchanges,
algo_namespace=algo_namespace,
live_graph=live_graph,
simulate_orders=simulate_orders,
stats_output=stats_output,
analyze_live=analyze_live,)
return algorithm_class
def _bundle_trading_environment(bundle_data, environ):
prefix, connstr = re.split(
r'sqlite:///',
str(bundle_data.asset_finder.engine.url),
maxsplit=1)
if prefix:
def choose_loader(column):
bound_cols = TradingPairPricing.columns
if column in bound_cols:
return ExchangePricingLoader(data_frequency)
raise ValueError(
"invalid url %r, must begin with 'sqlite:///'" %
str(bundle_data.asset_finder.engine.url))
"No PipelineLoader registered for column %s." % column
)
return TradingEnvironment(asset_db_path=connstr, environ=environ)
if live:
start = pd.Timestamp.utcnow()
# TODO: fix the end data.
if end is None:
end = start + timedelta(hours=8760)
def _build_live_algo_and_data(sim_params, exchanges, env, open_calendar,
simulate_orders, algo_namespace, capital_base,
live_graph, stats_output, analyze_live,
base_currency, namespace, choose_loader,
algorithm_class_kwargs):
sim_params._arena = 'live' # TODO: use the constructor instead
data = DataPortalExchangeLive(
exchanges=exchanges,
asset_finder=env.asset_finder,
trading_calendar=open_calendar,
first_trading_day=pd.to_datetime('today', utc=True)
)
data = _data_for_live_trading(sim_params, exchanges, env, open_calendar)
sim_params = create_simulation_parameters(
start=start,
end=end,
capital_base=capital_base,
emission_rate='minute',
data_frequency='minute'
)
algorithm_class = _algorithm_class_for_live(
algo_namespace, live_graph, stats_output, analyze_live,
base_currency, simulate_orders, exchanges, capital_base)
# TODO: use the constructor instead
sim_params._arena = 'live'
return data, algorithm_class(
namespace=namespace,
env=env,
get_pipeline_loader=choose_loader,
sim_params=sim_params,
**algorithm_class_kwargs)
def _build_backtest_algo_and_data(
exchanges, bundle, env, environ, bundle_timestamp, open_calendar,
start, end, namespace, choose_loader, sim_params,
algorithm_class_kwargs):
if exchanges:
algorithm_class = partial(
ExchangeTradingAlgorithmLive,
exchanges=exchanges,
algo_namespace=algo_namespace,
live_graph=live_graph,
simulate_orders=simulate_orders,
stats_output=stats_output,
analyze_live=analyze_live,
end=end,
)
elif exchanges:
# Removed the existing Poloniex fork to keep things simple
# We can add back the complexity if required.
@@ -288,19 +265,41 @@ def _build_backtest_algo_and_data(
asset_finder=None,
trading_calendar=open_calendar,
first_trading_day=start,
last_available_session=end)
last_available_session=end
)
sim_params = create_simulation_parameters(
start=start,
end=end,
capital_base=capital_base,
data_frequency=data_frequency,
emission_rate=data_frequency,
)
algorithm_class = partial(
ExchangeTradingAlgorithmBacktest,
exchanges=exchanges)
exchanges=exchanges
)
elif bundle is not None:
# TODO This branch should probably be removed or fixed: it doesn't even
# build `algorithm_class`, so it will break when trying to instantiate
# it.
bundle_data = load(bundle, environ, bundle_timestamp)
bundle_data = load(
bundle,
environ,
bundle_timestamp,
)
env = _bundle_trading_environment(bundle_data, environ)
prefix, connstr = re.split(
r'sqlite:///',
str(bundle_data.asset_finder.engine.url),
maxsplit=1,
)
if prefix:
raise ValueError(
"invalid url %r, must begin with 'sqlite:///'" %
str(bundle_data.asset_finder.engine.url),
)
env = TradingEnvironment(asset_db_path=connstr, environ=environ)
first_trading_day = \
bundle_data.equity_minute_bar_reader.first_trading_day
@@ -309,103 +308,27 @@ def _build_backtest_algo_and_data(
first_trading_day=first_trading_day,
equity_minute_reader=bundle_data.equity_minute_bar_reader,
equity_daily_reader=bundle_data.equity_daily_bar_reader,
adjustment_reader=bundle_data.adjustment_reader)
adjustment_reader=bundle_data.adjustment_reader,
)
return data, algorithm_class(
perf = algorithm_class(
namespace=namespace,
env=env,
get_pipeline_loader=choose_loader,
sim_params=sim_params,
**algorithm_class_kwargs)
def _build_algo_and_data(handle_data, initialize, before_trading_start,
analyze, algofile, algotext, defines, data_frequency,
capital_base, data, bundle, bundle_timestamp, start,
end, output, print_algo, local_namespace, environ,
live, exchange, algo_namespace, base_currency,
live_graph, analyze_live, simulate_orders,
stats_output):
namespace = _build_namespace(algotext, local_namespace, defines)
if algotext is not None:
algotext = algofile.read()
if print_algo:
_pretty_print_code(algotext)
mode = _mode(simulate_orders, live)
log.info('running algo in {mode} mode'.format(mode=mode))
exchanges = _build_exchanges_dict(exchange, live, simulate_orders,
base_currency)
open_calendar = get_calendar('OPEN')
env = TradingEnvironment(
load=partial(load_crypto_market_data, environ=environ, start_dt=start,
end_dt=end),
environ=environ,
exchange_tz='UTC',
asset_db_path=None) # We don't need an asset db, we have exchanges
env.asset_finder = ExchangeAssetFinder(exchanges=exchanges)
choose_loader = partial(_choose_loader, data_frequency)
if live:
start, end = _get_live_time_range()
data_frequency = 'minute' # TODO double check if this is the desired behavior
sim_params = create_simulation_parameters(
start=start,
end=end,
capital_base=capital_base,
emission_rate=data_frequency,
data_frequency=data_frequency)
if algotext is None:
algorithm_class_kwargs = {'initialize': initialize,
'handle_data': handle_data,
'before_trading_start': before_trading_start,
'analyze': analyze}
else:
algorithm_class_kwargs = {'algo_filename': getattr(algofile, 'name',
'<algorithm>'),
'script': algotext}
if live:
return _build_live_algo_and_data(
sim_params, exchanges, env, open_calendar, simulate_orders,
algo_namespace, capital_base, live_graph, stats_output,
analyze_live, base_currency, namespace, choose_loader,
algorithm_class_kwargs)
else:
return _build_backtest_algo_and_data(
exchanges, bundle, env, environ, bundle_timestamp, open_calendar,
start, end, namespace, choose_loader, sim_params,
algorithm_class_kwargs)
def _run(handle_data, initialize, before_trading_start, analyze, algofile,
algotext, defines, data_frequency, capital_base, data, bundle,
bundle_timestamp, start, end, output, print_algo, local_namespace,
environ, live, exchange, algo_namespace, base_currency, live_graph,
analyze_live, simulate_orders, stats_output):
"""Run an algorithm in backtest,
paper-trading or live-trading mode.
This is shared between the cli and :func:`catalyst.run_algo`.
"""
data, algorithm = _build_algo_and_data(
handle_data, initialize, before_trading_start, analyze, algofile,
algotext, defines, data_frequency, capital_base, data, bundle,
bundle_timestamp, start, end, output, print_algo, local_namespace,
environ, live, exchange, algo_namespace, base_currency, live_graph,
analyze_live, simulate_orders, stats_output)
perf = algorithm.run(
**{
'initialize': initialize,
'handle_data': handle_data,
'before_trading_start': before_trading_start,
'analyze': analyze,
} if algotext is None else {
'algo_filename': getattr(algofile, 'name', '<algorithm>'),
'script': algotext,
}
).run(
data,
overwrite_sim_params=False)
overwrite_sim_params=False,
)
if output == '-':
click.echo(str(perf))
@@ -462,7 +385,8 @@ def load_extensions(default, extensions, strict, environ, reload=False):
# without `strict` we should just log the failure
warnings.warn(
'Failed to load extension: %r\n%s' % (ext, e),
stacklevel=2)
stacklevel=2
)
else:
_loaded_extensions.add(ext)
@@ -489,6 +413,7 @@ def run_algorithm(initialize,
live_graph=False,
analyze_live=None,
simulate_orders=True,
auth_aliases=None,
stats_output=None,
output=os.devnull):
"""Run a trading algorithm.
@@ -561,7 +486,8 @@ def run_algorithm(initialize,
catalyst.data.bundles.bundles : The available data bundles.
"""
load_extensions(
default_extension, extensions, strict_extensions, environ)
default_extension, extensions, strict_extensions, environ
)
if capital_base is None:
raise ValueError(
@@ -569,7 +495,8 @@ def run_algorithm(initialize,
'amount of base currency available for trading. For example, '
'if the `capital_base` is 5ETH, the '
'`order_target_percent(asset, 1)` command will order 5ETH worth '
'of the specified asset.')
'of the specified asset.'
)
# I'm not sure that we need this since the modified DataPortal
# does not require extensions to be explicitly loaded.
@@ -587,11 +514,13 @@ def run_algorithm(initialize,
elif len(non_none_data) != 1:
raise ValueError(
'must specify one of `data`, `data_portal`, or `bundle`,'
' got: %r' % non_none_data)
' got: %r' % non_none_data,
)
elif 'bundle' not in non_none_data and bundle_timestamp is not None:
raise ValueError(
'cannot specify `bundle_timestamp` without passing `bundle`')
'cannot specify `bundle_timestamp` without passing `bundle`',
)
return _run(
handle_data=handle_data,
initialize=initialize,
@@ -618,4 +547,6 @@ def run_algorithm(initialize,
live_graph=live_graph,
analyze_live=analyze_live,
simulate_orders=simulate_orders,
stats_output=stats_output)
auth_aliases=auth_aliases,
stats_output=stats_output
)
+5 -5
View File
@@ -23,7 +23,7 @@ I18NSPHINXOPTS = $(PAPEROPT_$(PAPER)) $(SPHINXOPTS) source
help:
@echo "Please use \`make <target>' where <target> is one of"
@echo " build to build the C and Cython extensions for zipline"
@echo " build to build the C and Cython extensions for catalyst"
@echo " html to make standalone HTML files"
@echo " livehtml to run a persistent process that rebuilds the docs"
@echo " dirhtml to make HTML files named index.html in directories"
@@ -96,9 +96,9 @@ qthelp: build
@echo
@echo "Build finished; now you can run "qcollectiongenerator" with the" \
".qhcp project file in $(BUILDDIR)/qthelp, like this:"
@echo "# qcollectiongenerator $(BUILDDIR)/qthelp/zipline.qhcp"
@echo "# qcollectiongenerator $(BUILDDIR)/qthelp/catalyst.qhcp"
@echo "To view the help file:"
@echo "# assistant -collectionFile $(BUILDDIR)/qthelp/zipline.qhc"
@echo "# assistant -collectionFile $(BUILDDIR)/qthelp/catalyst.qhc"
applehelp: build
$(SPHINXBUILD) -b applehelp $(ALLSPHINXOPTS) $(BUILDDIR)/applehelp
@@ -113,8 +113,8 @@ devhelp: build
@echo
@echo "Build finished."
@echo "To view the help file:"
@echo "# mkdir -p $$HOME/.local/share/devhelp/zipline"
@echo "# ln -s $(BUILDDIR)/devhelp $$HOME/.local/share/devhelp/zipline"
@echo "# mkdir -p $$HOME/.local/share/devhelp/catalyst"
@echo "# ln -s $(BUILDDIR)/devhelp $$HOME/.local/share/devhelp/catalyst"
@echo "# devhelp"
epub: build
+5 -5
View File
@@ -8,8 +8,8 @@ from shutil import move, rmtree
from subprocess import check_call
HERE = dirname(abspath(__file__))
ZIPLINE_ROOT = dirname(HERE)
TEMP_LOCATION = '/tmp/zipline-doc'
CATALYST_ROOT = dirname(HERE)
TEMP_LOCATION = '/tmp/catalyst-doc'
TEMP_LOCATION_GLOB = TEMP_LOCATION + '/*'
@@ -46,8 +46,8 @@ def main():
print("Copying built files to temp location.")
move('build/html', TEMP_LOCATION)
print("Moving to '%s'" % ZIPLINE_ROOT)
os.chdir(ZIPLINE_ROOT)
print("Moving to '%s'" % CATALYST_ROOT)
os.chdir(CATALYST_ROOT)
print("Checking out gh-pages branch.")
check_call(
@@ -70,7 +70,7 @@ def main():
os.chdir(old_dir)
print()
print("Updated documentation branch in directory %s" % ZIPLINE_ROOT)
print("Updated documentation branch in directory %s" % CATALYST_ROOT)
print("If you are happy with these changes, commit and push to gh-pages.")
if __name__ == '__main__':
+2 -2
View File
@@ -127,9 +127,9 @@ if "%1" == "qthelp" (
echo.
echo.Build finished; now you can run "qcollectiongenerator" with the ^
.qhcp project file in %BUILDDIR%/qthelp, like this:
echo.^> qcollectiongenerator %BUILDDIR%\qthelp\zipline.qhcp
echo.^> qcollectiongenerator %BUILDDIR%\qthelp\catalyst.qhcp
echo.To view the help file:
echo.^> assistant -collectionFile %BUILDDIR%\qthelp\zipline.ghc
echo.^> assistant -collectionFile %BUILDDIR%\qthelp\catalyst.ghc
goto end
)
+4 -3
View File
@@ -483,7 +483,7 @@ bitcoin price.
Now we will run the simulation again, but this time we extend our original
algorithm with the addition of the ``analyze()`` function. Somewhat analogously
as how ``initialize()`` gets called once before the start of the algorith,
as how ``initialize()`` gets called once before the start of the algorithm,
``analyze()`` gets called once at the end of the algorithm, and receives two
variables: ``context``, which we discussed at the very beginning, and ``perf``,
which is the pandas dataframe containing the performance data for our algorithm
@@ -589,7 +589,7 @@ the ``examples`` directory:
from catalyst import run_algorithm
from catalyst.api import (order, record, symbol, order_target_percent,
get_open_orders)
from catalyst.exchange.stats_utils import extract_transactions
from catalyst.exchange.utils.stats_utils import extract_transactions
NAMESPACE = 'dual_moving_average'
log = Logger(NAMESPACE)
@@ -660,7 +660,8 @@ the ``examples`` directory:
def analyze(context, perf):
# Get the base_currency that was passed as a parameter to the simulation
base_currency = context.exchanges.values()[0].base_currency.upper()
exchange = list(context.exchanges.values())[0]
base_currency = exchange.base_currency.upper()
# First chart: Plot portfolio value using base_currency
ax1 = plt.subplot(411)
+2 -2
View File
@@ -41,11 +41,11 @@ master_doc = 'index'
# General information about the project.
project = u'Catalyst'
copyright = u'2017, Enigma MPC, Inc.'
copyright = u'2018, Enigma MPC, Inc.'
# The full version, including alpha/beta/rc tags, but excluding the commit hash
#release = version.split('+', 1)[0]
release = '0.3'
release = '0.4'
# List of patterns, relative to source directory, that match files and
# directories to ignore when looking for source files.
+19
View File
@@ -84,6 +84,25 @@ To build and view the docs locally, run:
$ {BROWSER} build/html/index.html
There is a `documented issue <https://github.com/sphinx-doc/sphinx/issues/3212>`_
with ``sphinx`` and ``docutils`` that causes the error below when trying to build
the docs.
.. code-block:: text
Exception occurred:
File "(...)/env-c/lib/python2.7/site-packages/docutils/writers/_html_base.py", line 671, in depart_document
assert not self.context, 'len(context) = %s' % len(self.context)
AssertionError: len(context) = 3
If you get this error, you need to downgrade your version of ``docutils`` as
follows, and build the docs again:
.. code-block:: bash
$ pip install docutils==0.12
Commit messages
---------------
+5 -5
View File
@@ -44,11 +44,11 @@ For additional details on the functionality added on recent releases, see the
Upcoming features
~~~~~~~~~~~~~~~~~
* Additional datasets beyond pricing data (Dec. 2017)
* API documentation (Jan. 2017)
* Support for decentralized exchanges (Jan. 2017)
* Support for data ingestion of community-contributed data sets (Jan. 2017)
* Pipeline support (Jan. 2018)
* Additional datasets beyond pricing data (Q1 2018)
* API documentation (Q1 2018)
* Support for decentralized exchanges (Q1 2018)
* Support for data ingestion of community-contributed data sets (Q1 2018)
* Pipeline support (Q1 2018)
* Web UI (Q2 2018)
+70 -30
View File
@@ -180,20 +180,6 @@ use a single tool to install Python and non-Python dependencies, or if you're
already using `Anaconda <http://continuum.io/downloads>`_ as your Python
distribution, refer to the :ref:`Installing with Conda <conda>` section.
Once you've installed the necessary additional dependencies for your system
(see below for your particular platform: :ref:`Linux`, :ref:`MacOS` or
:ref:`Windows`), you should be able to simply run
.. code-block:: bash
$ pip install enigma-catalyst matplotlib
Note that in the command above we install two different packages. The second
one, ``matplotlib`` is a visualization library. While it's not strictly
required to run catalyst simulations or live trading, it comes in very handy
to visualize the performance of your algorithms, and for this reason we
recommend you install it, as well.
If you use Python for anything other than Catalyst, we **strongly** recommend
that you install in a `virtualenv
<https://virtualenv.readthedocs.org/en/latest>`_. The `Hitchhiker's Guide to
@@ -206,8 +192,21 @@ summarized version:
$ pip install virtualenv
$ virtualenv catalyst-venv
$ source ./catalyst-venv/bin/activate
Once you've installed the necessary additional dependencies for your system
(:ref:`Linux`, :ref:`MacOS` or :ref:`Windows`) **and have activated your virtualenv**, you should be able to simply run
.. code-block:: bash
$ pip install enigma-catalyst matplotlib
Note that in the command above we install two different packages. The second
one, ``matplotlib`` is a visualization library. While it's not strictly
required to run catalyst simulations or live trading, it comes in very handy
to visualize the performance of your algorithms, and for this reason we
recommend you install it, as well.
Troubleshooting ``pip`` Install
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
@@ -219,13 +218,13 @@ Troubleshooting ``pip`` Install
.. code-block:: bash
pip install --upgrade pip
$ pip install --upgrade pip
On Windows, the recommended command is:
.. code-block:: bash
python -m pip install --upgrade pip
$ python -m pip install --upgrade pip
----
@@ -251,7 +250,7 @@ Troubleshooting ``pip`` Install
.. code-block:: bash
pip install --pre enigma-catalyst
$ pip install --pre enigma-catalyst
----
@@ -263,7 +262,7 @@ Troubleshooting ``pip`` Install
.. code-block:: bash
pip install --upgrade pip setuptools
$ pip install --upgrade pip setuptools
----
@@ -278,7 +277,7 @@ Troubleshooting ``pip`` Install
.. code-block:: bash
pip install -r requirements.txt
$ pip install -r requirements.txt
----
@@ -294,7 +293,7 @@ Troubleshooting ``pip`` Install
.. code-block:: bash
sudo apt-get install python-dev
$ sudo apt-get install python-dev
.. _pipenv:
@@ -376,14 +375,14 @@ outdated. Thus, you first need to run:
.. code-block:: bash
pip install --upgrade pip setuptools
$ pip install --upgrade pip setuptools
The default installation is also missing the C and C++ compilers, which you
install by:
.. code-block:: bash
sudo yum install gcc gcc-c++
$ sudo yum install gcc gcc-c++
Then you should follow the regular installation instructions outlined at the
beginning of this page.
@@ -408,20 +407,34 @@ following brew packages:
$ brew install freetype pkg-config gcc openssl
MacOS + virtualenv + matplotlib
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
MacOS + virtualenv/conda + matplotlib
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
A note about using matplotlib in virtual enviroments on MacOS: it may be
necessary to run
The first time that you try to run an algorithm that loads the ``matplotlib``
library, you may get the following error:
.. code-block:: text
RuntimeError: Python is not installed as a framework. The Mac OS X backend
will not be able to function correctly if Python is not installed as a
framework. See the Python documentation for more information on installing
Python as a framework on Mac OS X. Please either reinstall Python as a
framework, or try one of the other backends. If you are using (Ana)Conda
please install python.app and replace the use of 'python' with 'pythonw'.
See 'Working with Matplotlib on OSX' in the Matplotlib FAQ for more
information.
This is a ``matplotlib``-specific error, that will go away once you run the
following command:
.. code-block:: bash
echo "backend: TkAgg" > ~/.matplotlib/matplotlibrc
$ echo "backend: TkAgg" > ~/.matplotlib/matplotlibrc
in order to override the default ``MacOS`` backend for your system, which
may not be accessible from inside the virtual environment. This will allow
Catalyst to open matplotlib charts from within a virtual environment, which
is useful for displaying the performance of your backtests. To learn more
may not be accessible from inside the virtual or conda environment. This will
allow Catalyst to open matplotlib charts from within a virtual environment,
which is useful for displaying the performance of your backtests. To learn more
about matplotlib backends, please refer to the
`matplotlib backend documentation <https://matplotlib.org/faq/usage_faq.html#what-is-a-backend>`_.
@@ -475,6 +488,33 @@ mentioned above are as follows:
- ``cd`` into the folder where you downloaded ``VCForPython27.msi``
- Run ``msiexec /i VCForPython27.msi``
Updating Catalyst
-----------------
Catalyst is currently in alpha and in under very active development. We release
new minor versions every few days in response to the thorough battle testing
that our user community puts Catalyst in. As a result, you should expect to
update Catalyst frequently. Once installed, Catalyst can easily be updated as a
``pip`` package regardless of the environemnt used for installation. Make sure
you activate your environment first as you did in your first install, and then
execute:
.. code-block:: bash
$ pip uninstall enigma-catalyst
$ pip install enigma-catalyst
Alternatively, you could update Catalyst issuing the following command:
.. code-block:: bash
$ pip install -U enigma-catalyst
but this command will also upgrade all the Catalyst dependencies to the latest
versions available, and may have unexpected side effects if a newer version of a
dependency inadvertently breaks some functionality that Catalyst relies on.
Thus, the first method is the recommended one.
Getting Help
------------
+64 -9
View File
@@ -4,11 +4,63 @@ This document explains how to get started with live trading.
Supported Exchanges
^^^^^^^^^^^^^^^^^^^
Catalyst can trade against these exchanges:
- Bitfinex, id= ``bitfinex``
- Bittrex, id= ``bittrex``
- Poloniex, id= ``poloniex``
Since version 0.4, Catalyst integrated with `CCXT <https://github.com/ccxt/ccxt>`_,
a cryptocurrency trading library with support for more than 90 exchanges. The
range of CCXT and Catalyst support for each of those exchanges varies greatly.
The most supported exchanges are as follows:
The exchanges available for backtesting are fully supported in live mode:
- Bitfinex, id = ``bitfinex``
- Bittrex, id = ``bittrex``
- Poloniex, id = ``poloniex``
Additionally, we have successfully tested the following exchanges:
- Binance, id = ``binance``
- Bitmex, id = ``bitmex``
- GDAX, id = ``gdax``
As Catalyst is currently in Alpha and in under active development, you are
encouraged to throughly test any exchange in *paper trading* mode before trading
*live* with it.
Paper Trading vs Live Trading modes
^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^
Catalyst currently supports three different modes in which you can execute your
trading algorithm. The first is backtesting, which is covered extensively in the
tutorial, and uses historical data to run your algorithm. There is no
interaction with the exchange in backtesting mode, and this is the first mode
that you should test any new algorithm.
Once you are confident with the simulations that you have obtained with your
algorithm in backtesting, you may switch to live trading, where you have two
different modes:
* *Paper Trading*: The simulated algorithm runs in real time, and fetches
pricing data in real time from the exchange, but the orders never reach the
exchange, and are instead kept within Catalyst and simulated. No real currency
is bought or sold. Think of it as a `backtesting happening in real time`.
* *Live Trading*: This is the proper live trading mode in which an algorithm
runs in real time, fetching pricing data from live exchanges and placing orders
against the exchange. Real currency is transacted on the exchange driven by the
algorithm.
These three modes are controlled by the following variables:
+---------------+-------------------------+
| Mode | Parameters |
+ +-------+-----------------+
| | live | simulate_orders |
+---------------+-------+-----------------+
| backtesting | False | True (default) |
+---------------+-------+-----------------+
| paper trading | True | True |
+---------------+-------+-----------------+
| live trading | True | False |
+---------------+-------+-----------------+
Authentication
^^^^^^^^^^^^^^
@@ -75,7 +127,8 @@ Note that the trading pairs are always referenced in the same manner.
However, not all trading pairs are available on all exchanges. An
error will occur if the specified trading pair is not trading
on the exchange. To check which currency pairs are available on each
of the supported exchanges, see `Catalyst Market Coverage <https://www.enigma.co/catalyst/status`_.
of the supported exchanges, see
`Catalyst Market Coverage <https://www.enigma.co/catalyst/status>`_.
Trading an Algorithm
^^^^^^^^^^^^^^^^^^^^
@@ -105,20 +158,22 @@ What differs are the arguments provided to the catalyst client or
Here is the breakdown of the new arguments:
- ``live``: Boolean flag which enables live trading.
- ``live``: Boolean flag which enables live trading. It defaults to ``False``.
- ``capital_base``: The amount of base_currency assigned to the strategy.
It has to be lower or equal to the amount of base currency available for
trading on the exchange. For illustration, order_target_percent(asset, 1)
will order the capital_base amount specified here of the specified asset.
- ``exchange_name``: The name of the targeted exchange
(supported values: *bitfinex*, *bittrex*).
- ``exchange_name``: The name of the targeted exchange. See the
`CCXT Supported Exchanges <https://github.com/ccxt/ccxt/wiki/Exchange-Markets>`_
for the full list.
- ``algo_namespace``: A arbitrary label assigned to your algorithm for
data storage purposes.
- ``base_currency``: The base currency used to calculate the
statistics of your algorithm. Currently, the base currency of all
trading pairs of your algorithm must match this value.
- ``simulate_orders``: Enables the paper trading mode, in which orders are
simulated in Catalyst instead of processed on the exchange.
simulated in Catalyst instead of processed on the exchange. It defaults to
``True``.
Here is a complete algorithm for reference:
`Buy Low and Sell High <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/buy_low_sell_high_live.py>`_
+64 -3
View File
@@ -2,9 +2,70 @@
Release Notes
=============
Version 0.4.1
Version 0.4.7
^^^^^^^^^^^^^
**Release Date**: 2017-01-03
**Release Date**: 2018-01-19
Bug Fixes
~~~~~~~~~
- Fixing issue :issue:`137` impacting the CLI
Build
~~~~~
- Implemented authentication aliases (:issue:`60`)
Version 0.4.6
^^^^^^^^^^^^^
**Release Date**: 2018-01-18
Bug Fixes
~~~~~~~~~
- Fixed some Python3 issues
- Reading the trade log to get executed order prices on exchanges like Binance (:issue:`151`)
- Fixed issue with market order executing price (:issue:`150` and :issue:`111`)
- Implemented standardized symbol mapping (:issue:`157`)
- Improved error handling for unsupported timeframes (:issue:`159`)
- Using Bitfinex instead of Poloniex to fetch btc_usdt benchmark (:issue:`161`)
Build
~~~~~
- Added a `context.state` dict to keep arbitrary state values between runs
- Added ability to stop live algo at specified end date
Version 0.4.5
^^^^^^^^^^^^^
**Release Date**: 2018-01-12
Bug Fixes
~~~~~~~~~
- Improved order execution for exchanges supporting trade lists (:issue:`151`)
- Fixed an issue where requesting history of multiple assets repeats values
- Raising an error for order amounts smaller than exchange lots
- Handling multiple req errors with tickers more gracefully (:issue:`160`)
Version 0.4.4
^^^^^^^^^^^^^
**Release Date**: 2018-01-09
Bug Fixes
~~~~~~~~~
- Removed redundant capital_base validation (:issue:`142`)
- Fixed portfolio update issue with restored state (:issue:`111`)
- Skipping cash validation where there are open orders (:issue:`144`)
Version 0.4.3
^^^^^^^^^^^^^
**Release Date**: 2018-01-05
Bug Fixes
~~~~~~~~~
- Fixed CLI issue (:issue:`137`)
- Upgraded CCXT
Version 0.4.2
^^^^^^^^^^^^^
**Release Date**: 2018-01-03
Bug Fixes
~~~~~~~~~
@@ -39,7 +100,7 @@ Build
- Added market orders in live mode (:issue:`81`)
Version 0.3.10
^^^^^^^^^^^^^
~~~~~~~~~~~~~~
**Release Date**: 2017-11-28
Bug Fixes
+1 -1
View File
@@ -20,7 +20,7 @@ dependencies:
- bcolz==0.12.1
- bottleneck==1.2.1
- chardet==3.0.4
- ccxt==1.10.283
- ccxt==1.10.774
- click==6.7
- contextlib2==0.5.5
- cycler==0.10.0
+1 -1
View File
@@ -81,6 +81,6 @@ empyrical==0.2.1
tables==3.3.0
#Catalyst dependencies
ccxt==1.10.283
ccxt==1.10.774
boto3==1.4.8
redo==1.6
+1
View File
@@ -1,3 +1,4 @@
Sphinx>=1.3.2
numpydoc>=0.5.0
sphinx-autobuild==0.6.0
docutils==0.12
+2 -1
View File
@@ -10,7 +10,8 @@ from catalyst.exchange.exchange_bcolz import BcolzExchangeBarReader, \
from catalyst.exchange.exchange_bundle import ExchangeBundle, \
BUNDLE_NAME_TEMPLATE
from catalyst.exchange.utils.bundle_utils import get_bcolz_chunk, \
get_start_dt, get_df_from_arrays
get_df_from_arrays
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
+26 -8
View File
@@ -1,7 +1,9 @@
import pandas as pd
from logbook import Logger
from base import BaseExchangeTestCase
from catalyst.testing import ZiplineTestCase
from catalyst.testing.fixtures import WithLogger
from .base import BaseExchangeTestCase
from catalyst.exchange.ccxt.ccxt_exchange import CCXT
from catalyst.exchange.exchange_execution import ExchangeLimitOrder
from catalyst.exchange.utils.exchange_utils import get_exchange_auth
@@ -13,22 +15,22 @@ log = Logger('test_ccxt')
class TestCCXT(BaseExchangeTestCase):
@classmethod
def setup(self):
exchange_name = 'binance'
exchange_name = 'bitfinex'
auth = get_exchange_auth(exchange_name)
self.exchange = CCXT(
exchange_name=exchange_name,
key=auth['key'],
secret=auth['secret'],
base_currency='eth',
base_currency='bnb',
)
self.exchange.init()
def test_order(self):
log.info('creating order')
asset = self.exchange.get_asset('neo_eth')
asset = self.exchange.get_asset('neo_bnb')
order_id = self.exchange.order(
asset=asset,
style=ExchangeLimitOrder(limit_price=0.7),
style=ExchangeLimitOrder(limit_price=10),
amount=1,
)
log.info('order created {}'.format(order_id))
@@ -56,10 +58,10 @@ class TestCCXT(BaseExchangeTestCase):
def test_get_candles(self):
log.info('retrieving candles')
candles = self.exchange.get_candles(
freq='5T',
freq='30T',
assets=[self.exchange.get_asset('eth_btc')],
bar_count=200,
start_dt=pd.to_datetime('2017-01-01', utc=True)
start_dt=pd.to_datetime('2017-09-01', utc=True)
)
for asset in candles:
@@ -70,12 +72,28 @@ class TestCCXT(BaseExchangeTestCase):
def test_tickers(self):
log.info('retrieving tickers')
assets = [
self.exchange.get_asset('eng_eth'),
self.exchange.get_asset('iot_usd'),
]
tickers = self.exchange.tickers(assets)
assert len(tickers) == 1
pass
def test_my_trades(self):
asset = self.exchange.get_asset('dsh_btc')
trades = self.exchange.get_trades(asset)
assert trades
pass
def test_get_executed_order(self):
log.info('retrieving executed order')
asset = self.exchange.get_asset('eng_eth')
order = self.exchange.get_order('165784', asset)
transactions = self.exchange.process_order(order)
assert transactions
pass
def test_get_balances(self):
log.info('testing wallet balances')
# balances = self.exchange.get_balances()
@@ -0,0 +1,79 @@
import importlib
from os.path import join, isfile
import pandas as pd
import os
from catalyst import run_algorithm
from catalyst.exchange.utils.stats_utils import get_pretty_stats, \
extract_transactions, set_print_settings, extract_orders
from catalyst.testing.fixtures import WithLogger, ZiplineTestCase
from logbook import TestHandler, WARNING
from pathtools.path import listdir
filter_algos = [
'buy_and_hodl.py',
'buy_btc_simple.py',
'buy_low_sell_high.py',
'mean_reversion_simple.py',
'rsi_profit_target.py',
'simple_loop.py',
'simple_universe.py',
]
class TestSuiteAlgo(WithLogger, ZiplineTestCase):
@staticmethod
def analyze(context, perf):
set_print_settings()
transaction_df = extract_transactions(perf)
print('the transactions:\n{}'.format(transaction_df))
orders_df = extract_orders(perf)
print('the orders:\n{}'.format(orders_df))
stats = get_pretty_stats(perf, show_tail=False, num_rows=5)
print('the stats:\n{}'.format(stats))
pass
def test_run_examples(self):
folder = join('..', '..', '..', 'catalyst', 'examples')
files = [f for f in listdir(folder) if isfile(join(folder, f))]
algo_list = []
for filename in files:
name = os.path.basename(filename)
if filter_algos and name not in filter_algos:
continue
module_name = 'catalyst.examples.{}'.format(
name.replace('.py', '')
)
algo_list.append(module_name)
for module_name in algo_list:
algo = importlib.import_module(module_name)
namespace = module_name.replace('.', '_')
log_catcher = TestHandler()
with log_catcher:
run_algorithm(
capital_base=0.1,
data_frequency='minute',
initialize=algo.initialize,
handle_data=algo.handle_data,
analyze=TestSuiteAlgo.analyze,
exchange_name='poloniex',
algo_namespace='test_{}'.format(namespace),
base_currency='eth',
start=pd.to_datetime('2017-10-01', utc=True),
end=pd.to_datetime('2017-10-02', utc=True),
# output=out
)
warnings = [record for record in log_catcher.records if
record.level == WARNING]
if len(warnings) > 0:
print('WARNINGS:\n{}'.format(warnings))
pass
@@ -1,7 +1,8 @@
import random
import os
import pandas as pd
from logbook import Logger
from logbook import TestHandler
from pandas.util.testing import assert_frame_equal
from catalyst import get_calendar
@@ -12,8 +13,6 @@ from catalyst.exchange.utils.factory import get_exchange
from catalyst.exchange.utils.test_utils import output_df, \
select_random_assets
log = Logger('TestSuiteExchange')
pd.set_option('display.expand_frame_repr', False)
pd.set_option('precision', 8)
pd.set_option('display.width', 1000)
@@ -22,10 +21,11 @@ pd.set_option('display.max_colwidth', 1000)
class TestSuiteBundle:
@staticmethod
def get_data_portal(exchange_names):
def get_data_portal(exchanges):
open_calendar = get_calendar('OPEN')
asset_finder = ExchangeAssetFinder()
asset_finder = ExchangeAssetFinder(exchanges)
exchange_names = [exchange.name for exchange in exchanges]
data_portal = DataPortalExchangeBacktest(
exchange_names=exchange_names,
asset_finder=asset_finder,
@@ -46,7 +46,9 @@ class TestSuiteBundle:
assets
end_dt
bar_count
sample_minutes
freq
data_frequency
data_portal
Returns
-------
@@ -54,51 +56,60 @@ class TestSuiteBundle:
"""
data = dict()
log.info('creating data sample from bundle')
data['bundle'] = data_portal.get_history_window(
assets=assets,
end_dt=end_dt,
bar_count=bar_count,
frequency=freq,
field='close',
data_frequency=data_frequency,
)
log.info('bundle data:\n{}'.format(
data['bundle'].tail(10))
)
log_catcher = TestHandler()
with log_catcher:
data['bundle'] = data_portal.get_history_window(
assets=assets,
end_dt=end_dt,
bar_count=bar_count,
frequency=freq,
field='close',
data_frequency=data_frequency,
)
candles = exchange.get_candles(
end_dt=end_dt,
freq=freq,
assets=assets,
bar_count=bar_count,
)
data['exchange'] = get_candles_df(
candles=candles,
field='close',
freq=freq,
bar_count=bar_count,
end_dt=end_dt,
)
for source in data:
df = data[source]
path, folder = output_df(
df, assets, '{}_{}'.format(freq, source)
)
log.info('creating data sample from exchange api')
candles = exchange.get_candles(
end_dt=end_dt,
freq=freq,
assets=assets,
bar_count=bar_count,
)
data['exchange'] = get_candles_df(
candles=candles,
field='close',
freq=freq,
bar_count=bar_count,
end_dt=end_dt,
)
log.info('exchange data:\n{}'.format(
data['exchange'].tail(10))
)
for source in data:
df = data[source]
path = output_df(df, assets, '{}_{}'.format(freq, source))
log.info('saved {}:\n{}'.format(source, path))
print('saved {} test results: {}'.format(end_dt, folder))
assert_frame_equal(
right=data['bundle'],
left=data['exchange'],
check_less_precise=True,
)
assert_frame_equal(
right=data['bundle'],
left=data['exchange'],
check_less_precise=1,
)
try:
assert_frame_equal(
right=data['bundle'],
left=data['exchange'],
check_less_precise=min([a.decimals for a in assets]),
)
except Exception as e:
print('Some differences were found within a 1 decimal point '
'interval of confidence: {}'.format(e))
with open(os.path.join(folder, 'compare.txt'), 'w+') as handle:
handle.write(e.args[0])
pass
def test_validate_bundles(self):
# exchange_population = 3
asset_population = 3
data_frequency = random.choice(['minute', 'daily'])
data_frequency = random.choice(['minute'])
# bundle = 'dailyBundle' if data_frequency
# == 'daily' else 'minuteBundle'
@@ -106,11 +117,9 @@ class TestSuiteBundle:
# population=exchange_population,
# features=[bundle],
# ) # Type: list[Exchange]
exchanges = [get_exchange('bitfinex', skip_init=True)]
exchanges = [get_exchange('poloniex', skip_init=True)]
data_portal = TestSuiteBundle.get_data_portal(
[exchange.name for exchange in exchanges]
)
data_portal = TestSuiteBundle.get_data_portal(exchanges)
for exchange in exchanges:
exchange.init()
@@ -1,21 +1,26 @@
import json
import os
import random
from logging import Logger
from logging import Logger, WARNING
from time import sleep
import pandas as pd
from catalyst.assets._assets import TradingPair
from logbook import TestHandler
from catalyst.exchange.exchange_errors import ExchangeRequestError
from catalyst.exchange.exchange_execution import ExchangeLimitOrder
from catalyst.exchange.utils.exchange_utils import get_exchange_folder
from catalyst.exchange.utils.test_utils import select_random_exchanges, \
handle_exchange_error, select_random_assets
from catalyst.testing import ZiplineTestCase
from catalyst.testing.fixtures import WithLogger
from exchange.utils.factory import get_exchanges
log = Logger('TestSuiteExchange')
class TestSuiteExchange:
class TestSuiteExchange(WithLogger, ZiplineTestCase):
def _test_markets_exchange(self, exchange, attempts=0):
assets = None
try:
@@ -79,12 +84,13 @@ class TestSuiteExchange:
def test_tickers(self):
exchange_population = 3
asset_population = 3
asset_population = 15
exchanges = select_random_exchanges(
exchange_population,
features=['fetchTickers'],
) # Type: list[Exchange]
# exchanges = select_random_exchanges(
# exchange_population,
# features=['fetchTickers'],
# ) # Type: list[Exchange]
exchanges = list(get_exchanges(['bitfinex']).values())
for exchange in exchanges:
exchange.init()
@@ -156,34 +162,47 @@ class TestSuiteExchange:
base_currency=quote_currency,
) # Type: list[Exchange]
for exchange in exchanges:
exchange.init()
log_catcher = TestHandler()
with log_catcher:
for exchange in exchanges:
exchange.init()
assets = exchange.get_assets(quote_currency=quote_currency)
asset = select_random_assets(assets, 1)[0]
assert asset
assets = exchange.get_assets(quote_currency=quote_currency)
asset = select_random_assets(assets, 1)[0]
self.assertIsInstance(asset, TradingPair)
tickers = exchange.tickers([asset])
price = tickers[asset]['last_price']
tickers = exchange.tickers([asset])
price = tickers[asset]['last_price']
amount = order_amount / price
amount = order_amount / price
limit_price = price * 0.8
style = ExchangeLimitOrder(limit_price=limit_price)
limit_price = price * 0.8
style = ExchangeLimitOrder(limit_price=limit_price)
order = exchange.order(
asset=asset,
amount=amount,
style=style,
)
sleep(1)
order = exchange.order(
asset=asset,
amount=amount,
style=style,
)
sleep(1)
open_order, _ = exchange.get_order(order.id, asset)
assert open_order.status == 0
open_order = exchange.get_order(order.id, asset)
self.assertEqual(0, open_order.status)
exchange.cancel_order(open_order, asset)
sleep(1)
exchange.cancel_order(open_order, asset)
sleep(1)
canceled_order, _ = exchange.get_order(open_order.id, asset)
assert canceled_order.status == 2
canceled_order = exchange.get_order(open_order.id, asset)
warnings = [record for record in log_catcher.records if
record.level == WARNING]
self.assertEqual(0, len(warnings))
self.assertEqual(2, canceled_order.status)
print(
'tested {exchange} / {symbol}, order: {order}'.format(
exchange=exchange.name,
symbol=asset.symbol,
order=order.id,
)
)
pass