Compare commits

...
124 Commits
Author SHA1 Message Date
Frederic Fortier 49b6792399 Merge branch 'develop' 2018-02-09 12:01:57 -05:00
Frederic Fortier 6b1982274e DOC: updated release notes for 0.5.3 2018-02-09 02:44:45 -05:00
Frederic Fortier bc58640a3a Merge remote-tracking branch 'origin/develop' into develop 2018-02-09 02:41:08 -05:00
Frederic Fortier 511af7946e BUG: for issue #219, created another unit test which compares current price against last candle 2018-02-09 02:40:53 -05:00
Victor Grau Serrat dcdb3d9d7a MAINT: removing many warnings from building docs 2018-02-09 00:07:38 -07:00
Frederic Fortier 5c77cfc68d BUG: for issue #219, fixed a resampling issue 2018-02-09 01:38:22 -05:00
Frederic Fortier 11f8e36ff0 Merge remote-tracking branch 'origin/develop' into develop 2018-02-09 00:56:45 -05:00
Frederic Fortier 44e4e66f5c BLD: adding more log messages 2018-02-09 00:56:35 -05:00
Victor Grau Serrat 3c4c6c3dfd DOC: small edits, eliminating sphinx warnings 2018-02-08 22:10:57 -07:00
Victor Grau Serrat bebebc46dc DOC: small edits, eliminating sphinx warnings 2018-02-08 22:07:58 -07:00
Frederic Fortier 403d7f9c29 Merge branch 'develop' 2018-02-08 17:32:47 -05:00
Frederic Fortier 09ad397ee6 DOC: adjusted the release notes 2018-02-08 17:31:09 -05:00
Frederic Fortier e56e1f8e21 BUG: fixed an issue with open orders 2018-02-08 17:26:48 -05:00
Frederic Fortier 00f232e2d7 BUG: fixed sample algo 2018-02-08 17:10:41 -05:00
Frederic Fortier a820f66bdc BLD: adjusted unit test 2018-02-08 16:57:23 -05:00
Frederic Fortier 82d318a1c0 BUG: fixed issue #216 with bad candle data 2018-02-08 15:51:55 -05:00
Frederic Fortier 97afbf781e Merge branch 'master' into develop 2018-02-08 01:13:29 -05:00
Frederic Fortier 43a5f8858b DOC: updated release number 2018-02-08 01:11:09 -05:00
Frederic Fortier dc04cd8781 Merge branch 'master' into develop 2018-02-08 00:59:46 -05:00
Frederic Fortier 793b6e92d8 BLD: included marketplace dependencies 2018-02-08 00:47:25 -05:00
Frederic Fortier c0b9939580 Merge branch 'develop' 2018-02-08 00:33:55 -05:00
Frederic Fortier e248831719 BLD: adjusted sample algo 2018-02-08 00:27:56 -05:00
Frederic Fortier 25f1f6e641 BLD: upgraded CCXT 2018-02-08 00:18:48 -05:00
Frederic Fortier d2f9762fbf BLD: adjusting the marketplace sample algo 2018-02-08 00:02:40 -05:00
Victor Grau Serrat 7a89cfc02b BUG: marketplace: balance returns different types 2018-02-07 21:01:21 -07:00
Frederic Fortier 64e22ba27d BLD: Updated release notes for 5.0 2018-02-07 22:28:08 -05:00
Frederic Fortier 8c58916bb1 Merge remote-tracking branch 'origin/develop' into develop 2018-02-07 22:17:05 -05:00
Frederic Fortier 18cdd51680 BUG: fixed issue with order processing 2018-02-07 22:16:32 -05:00
Victor Grau Serrat db7b7639a0 MAINT: not listing test datasets 2018-02-07 17:34:11 -07:00
Frederic Fortier 1ad26b0c39 BLD: trying to format a line 2018-02-07 19:29:38 -05:00
Victor Grau Serrat e1abecb556 MAINT: updated catalyst logo 2018-02-07 17:09:49 -07:00
Victor Grau Serrat 6e642f45d8 MAINT: marketplace prompts 2018-02-07 14:44:11 -07:00
Frederic Fortier 9f9bfc9df0 BLD: dropping dataset index cols to avoid duplicates 2018-02-07 16:40:43 -05:00
Frederic Fortier 866b92910f BLD: fixed the marketplace api in the algo runtime 2018-02-07 13:10:07 -05:00
Victor Grau Serrat e071f6ec8d BLD: implemented grains, catch JSON malformed in addresses.json 2018-02-07 10:57:07 -07:00
lenak25 3eff58fcf9 Merge branch 'develop' of github.com:enigmampc/catalyst into develop 2018-02-07 12:49:12 +02:00
lenak25 540358d0f9 BLD: new marketplace contract address 2018-02-07 12:48:24 +02:00
Victor Grau Serrat 0e3be98d24 BLD: mmarketplace: switched encoding to Web3.toHex() 2018-02-07 01:44:16 -07:00
Victor Grau Serrat c39e075766 MAINT: constants points to contract address/abi in develop branch 2018-02-06 23:50:24 -07:00
Frederic Fortier 54935dd6f6 Merge remote-tracking branch 'origin/develop' into develop 2018-02-07 01:46:45 -05:00
Frederic Fortier 093660ab1a BLD: simplified unit tests until the next release 2018-02-07 01:46:33 -05:00
Victor Grau Serrat fb65839032 BLD: marketplace: getkeysecret is authenticated 2018-02-06 23:33:42 -07:00
Frederic Fortier 461a5942fb BUG: fixed Python 2 issue EXPERIMENTAL 2018-02-07 01:17:00 -05:00
Frederic Fortier eee2e1be88 BUG: fixed Python 2 issue 2018-02-07 00:48:17 -05:00
Frederic Fortier 2d43955abf BUG: fixed Python 2 issue 2018-02-07 00:45:33 -05:00
Frederic Fortier c5bed6e8c4 BLD: made some adjustments during testing 2018-02-06 21:15:30 -05:00
Frederic Fortier 2320a1432a BLD: Merge remote-tracking branch 'remotes/origin/data-marketplace' into develop 2018-02-06 14:34:07 -05:00
Frederic Fortier 5d74cd6f89 BLD: Merge remote-tracking branch 'remotes/origin/data-marketplace' into develop 2018-02-06 14:32:59 -05:00
Frederic Fortier 273408e9ed Merge remote-tracking branch 'origin/data-marketplace' into data-marketplace 2018-02-06 14:18:00 -05:00
Frederic Fortier 5ccddf6d21 BLD: removed dummy smart contract 2018-02-06 14:17:48 -05:00
Victor Grau Serrat d97890fc4e BLD: marketplace: moving AUTH_SERVER paths to /marketplace/* 2018-02-06 11:43:41 -07:00
Victor Grau Serrat 4c1a9b1dd7 BLD: marketplace added listing of datasets 2018-02-05 22:58:40 -07:00
Victor Grau Serrat 2168fb0d5b BUG: address json initialized with array of 1 dict, instead of dict 2018-02-05 22:17:21 -07:00
Frederic Fortier 3bc54a6c2e BLD: finalized the register implementation 2018-02-05 23:46:09 -05:00
Frederic Fortier 444fcbb2b3 BLD: adjustments for the new contract and developing register 2018-02-05 23:04:23 -05:00
Isan-Rivkin 881e3e5953 REV:constants.py revert 2018-02-04 08:51:27 -08:00
Isan-Rivkin 5889f4c74b BLD: added smart contract addresses 2018-02-04 08:45:06 -08:00
Isan-Rivkin ee93b16558 BLD:Solidity contract addressed added to constants 2018-02-04 08:42:42 -08:00
lenak25 e94c9db2ec Updated contract 2018-02-04 18:29:11 +02:00
Frederic Fortier 9f88e7a003 BUG: for issue #183, added more logic to catch order amount adjustments 2018-02-03 18:24:28 -05:00
Frederic Fortier b9bedfda21 BLD: enhancing registration features 2018-02-02 17:23:13 -05:00
Frederic Fortier d1529afeb3 BLD: enhancing registration features 2018-02-02 17:10:10 -05:00
Frederic Fortier 99dece953f BLD: ingesting marketplace bundles 2018-02-02 16:20:33 -05:00
Victor Grau Serrat 9ef000e529 BUG: balance returns different types 2018-02-02 12:06:57 -07:00
Victor Grau Serrat db4ee4c01e BLD: check for existence of dataset in ingestion & subscription 2018-02-01 23:42:40 -07:00
Victor Grau Serrat 0decdf5aab BLD: check for name duplicates when registering dataset 2018-02-01 22:30:24 -07:00
Victor Grau Serrat e460290056 BLD: better exception handling for failed multipart download 2018-02-01 21:59:30 -07:00
Frederic Fortier b9130fd968 BLD: added missing requirement 2018-02-01 22:03:16 -05:00
Frederic Fortier 37980f3580 BLD: working on marketplace integration 2018-02-01 20:43:48 -05:00
Victor Grau Serrat 56bd3db7a2 BLD: marketplace - ingestion downloads multiple files 2018-02-01 16:54:24 -07:00
Victor Grau Serrat 18123d7f7e BLD: marketplace constants in constants.py 2018-02-01 11:47:25 -07:00
Victor Grau Serrat a04a99373d BLD: marketplace: subscription to dataset 2018-02-01 08:07:17 -07:00
Victor Grau Serrat e906315969 BLD: marketplace: enigma contract files 2018-01-31 22:06:52 -07:00
Frederic Fortier a360a5fe3a Housekeeping 2018-01-31 23:57:22 -05:00
Frederic Fortier fcbdc131ec BLD: trying to catch as many trades as possible 2018-01-31 22:51:41 -05:00
Frederic Fortier 857a5d8a91 BUG: for issue #178, modified logic to track adjusted order amount 2018-01-31 21:57:33 -05:00
Frederic Fortier e914481325 BLD: minor adjustment 2018-01-31 17:19:09 -05:00
Frederic Fortier c0ba8b2ebb BUG: for issue #178, adjusted the fallback processing of orders for exchanges lacking a "my trades" api 2018-01-31 17:12:05 -05:00
Frederic Fortier 69507d1b00 BLD: for issue #178, fixed the "set" issue when fetching a ticker from positions 2018-01-31 16:32:22 -05:00
Frederic Fortier 311e357451 BLD: updated CCXT 2018-01-31 16:10:39 -05:00
Victor Grau Serrat d1cd95c492 BLD: marketplace subscribe + refactoring 2018-01-31 13:39:59 -07:00
Victor Grau Serrat 0f55f0e9a6 BLD: passing dataset to marketplace publish request 2018-01-31 10:50:10 -07:00
Victor Grau Serrat c26d78cf05 BLD: set AUTH_SERVER as a Catalyst constant 2018-01-30 11:20:06 -07: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 18c96303a6 BLD: improved unit test 2018-01-29 19:20:55 -05:00
Frederic Fortier 5c236a65f7 BUG: for issue #178, checking the order status instead of relying on the open amount 2018-01-29 19:20:34 -05:00
Frederic Fortier 7f021acb2e BUG: fixed catalyst import 2018-01-29 15:58:34 -05:00
Victor Grau Serrat 68faad4098 BLD: publish dataset in marketplace 2018-01-26 15:32:53 -07:00
Victor Grau Serrat ff989ba524 BLD: register dataset in marketplace 2018-01-25 22:39:54 -07:00
Frederic Fortier fa60457a0f BLD: misc adjustments to fetch all candles on the Binance exchange 2018-01-25 23:38:37 -05:00
Frederic Fortier 43737a9730 Merge branch 'alexiri-patch-1' into develop 2018-01-25 17:52:41 -05:00
Frederic Fortier e6fb708c1c Merge branch 'patch-1' of https://github.com/alexiri/catalyst into alexiri-patch-1 2018-01-25 17:52:22 -05:00
Frederic Fortier 8ff8e6458e BUG: fixed stats output issue #171 by adding orders and transactions in the header 2018-01-25 17:43:00 -05:00
Frederic Fortier 7d24433a42 Merge remote-tracking branch 'origin/develop' into develop 2018-01-25 17:41:02 -05:00
Frederic Fortier 7ae23e2340 BLD: conditionally fetching single or multi tickers for performance reasons 2018-01-25 17:40:45 -05:00
Avishai WeingartenandVictor Grau Serrat 940f625ec1 fix for click.echo
added sys.stdout to click.echo to prevent errors on jupyter
2018-01-25 12:46:27 -07:00
Victor Grau Serrat 2ec9aa2ca9 BLD: CLI implementation for the marketplace 2018-01-25 11:55:25 -07:00
Avishai WeingartenandGitHub c6dfd502d2 fix for click.echo
added sys.stdout to click.echo to prevent errors on jupyter
2018-01-25 16:02:44 +02:00
Frederic Fortier 0aa8c91577 BLD: for issue #174, re-implemented the fetch_tickers approach 2018-01-24 22:43:08 -05:00
Frederic Fortier c09f53449a BUG: fixed issue #176 with ignoring ticker errors instead of raising 2018-01-24 22:26:20 -05:00
Frederic Fortier f34a66e9c6 BUG: trying to fix issue #178 with Binance lot sizes 2018-01-24 21:54:05 -05:00
Victor Grau Serrat 04fed4140c BLD: sourcing contract address+abi from github 2018-01-24 14:12:59 -07:00
Victor Grau Serrat 76e183b5f7 BLD: contract address+abi on testnet 2018-01-24 12:48:42 -07:00
Frederic Fortier 911fb6e934 Merge branch 'gthouret-echo-usage-for-jupyter' 2018-01-24 00:33:52 -05:00
Frederic Fortier e8f98825e0 Merge branch 'treethought-empyrical-errors' into develop 2018-01-24 00:31:59 -05:00
Frederic Fortier 5619f6f451 Merge branch 'empyrical-errors' of https://github.com/treethought/catalyst into treethought-empyrical-errors 2018-01-24 00:31:49 -05:00
Cam Sweeney 8223d06d98 BUG: Use empyrical patches for persisting issue #126
The referenced issue was addressed via importing a set of patches
for empyrical. However the same error occurs occasionally when calling
the empyrical functions inside "/catalyst/finance/risk/period.py".

This PR simply applies 2 of the same patches in period.py.
I have only experienced problems with cum_returns and max_drawdown
thus far, but it is likely the other patches may be needed.
2018-01-23 02:42:55 -08:00
Frederic Fortier b47884a504 BLD: code formatting 2018-01-22 17:27:59 -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
Alex IribarrenandGitHub df4dd0f5a4 Update example-algos.rst
exchange_utils was moved.
2018-01-14 13:55:16 +01:00
Frederic Fortier 70b9fd8e03 BLD: completed data market place first integration and created an algo to test it 2018-01-12 00:39:14 -05:00
Frederic Fortier 21ad753fcf Merge remote-tracking branch 'origin/data-marketplace' into data-marketplace 2018-01-11 20:05:01 -05:00
Frederic Fortier 88aa7154db BLD: ingesting and cleaning marketplace data sources 2018-01-11 20:04:53 -05:00
VictorandGitHub b1d63c9a5b Update README.md 2018-01-11 12:15:49 -07:00
Frederic Fortier 3ea00e783c DOC: improved code comments 2018-01-11 14:08:52 -05:00
Frederic Fortier 9361570895 DOC: improved code comments 2018-01-11 14:08:02 -05:00
Frederic Fortier 50ded59c3c DOC: documented the marketplace structure 2018-01-11 14:01:00 -05:00
Frederic Fortier 354e422be8 BLD: working on ingesting alt data sources 2018-01-10 21:54:09 -05:00
Frederic Fortier 1e4af12e3d BLD: improved smart contract and catalyst commands 2018-01-10 13:31:08 -05:00
Frederic Fortier e0c8178932 BLD: improved smart contract and catalyst commands 2018-01-09 22:42:09 -05:00
Frederic Fortier 22a850ab14 BLD: testing basic smart contract 2018-01-09 20:41:03 -05:00
Frederic Fortier aab501aa17 BLD: first test of the marketplace smart contract 2018-01-09 19:08:53 -05:00
Frederic Fortier 94d5b4a4d5 BLD: defined first version of commands and marketplace class 2018-01-09 17:16:54 -05:00
53 changed files with 3022 additions and 411 deletions
+2 -4
View File
@@ -1,4 +1,4 @@
.. image:: https://s3.amazonaws.com/enigmaco-docs/enigma-catalyst.jpg
.. image:: https://s3.amazonaws.com/enigmaco-docs/enigma-catalyst.png
:target: https://enigmampc.github.io/catalyst
:align: center
:alt: Enigma | Catalyst
@@ -17,9 +17,7 @@ insights regarding a particular strategy's performance. Catalyst also supports
live-trading of crypto-assets starting with three exchanges (Bitfinex, Bittrex,
and Poloniex) with more being added over time. Catalyst empowers users to share
and curate data and build profitable, data-driven investment strategies. Please
visit `enigma.co <https://www.enigma.co>`_ to learn more about Catalyst, or
refer to the `whitepaper <https://www.enigma.co/enigma_catalyst.pdf>`_ for
further technical details.
visit `enigma.co <https://www.enigma.co>`_ to learn more about Catalyst.
Catalyst builds on top of the well-established
`Zipline <https://github.com/quantopian/zipline>`_ project. We did our best to
+140 -11
View File
@@ -3,8 +3,10 @@ import os
from functools import wraps
import click
import sys
import logbook
import pandas as pd
from catalyst.marketplace.marketplace import Marketplace
from six import text_type
from catalyst.data import bundles as bundles_module
@@ -257,7 +259,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,
@@ -290,7 +292,7 @@ def run(ctx,
)
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)
@@ -460,10 +462,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,
@@ -496,7 +498,7 @@ def live(ctx,
)
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)
@@ -578,7 +580,8 @@ 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,
@@ -601,10 +604,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')
@@ -631,11 +635,12 @@ 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()
@@ -756,7 +761,131 @@ 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)
@main.group()
@click.pass_context
def marketplace(ctx):
pass
@marketplace.command()
@click.pass_context
def ls(ctx):
click.echo('Listing of available data sources on the marketplace:',
sys.stdout)
marketplace = Marketplace()
marketplace.list()
@marketplace.command()
@click.option(
'--dataset',
default=None,
help='The name of the dataset to ingest from the Data Marketplace.',
)
@click.pass_context
def subscribe(ctx, dataset):
if dataset is None:
ctx.fail("must specify a dataset to subscribe to with '--dataset'\n"
"List available dataset on the marketplace with "
"'catalyst marketplace ls'")
marketplace = Marketplace()
marketplace.subscribe(dataset)
@marketplace.command()
@click.option(
'--dataset',
default=None,
help='The name of the dataset to ingest from the Data Marketplace.',
)
@click.option(
'-f',
'--data-frequency',
type=click.Choice({'daily', 'minute', 'daily,minute', 'minute,daily'}),
default='daily',
show_default=True,
help='The data frequency of the desired OHLCV bars.',
)
@click.option(
'-s',
'--start',
default=None,
type=Date(tz='utc', as_timestamp=True),
help='The start date of the data range. (default: one year from end date)',
)
@click.option(
'-e',
'--end',
default=None,
type=Date(tz='utc', as_timestamp=True),
help='The end date of the data range. (default: today)',
)
@click.pass_context
def ingest(ctx, dataset, data_frequency, start, end):
if dataset is None:
ctx.fail("must specify a dataset to clean with '--dataset'\n"
"List available dataset on the marketplace with "
"'catalyst marketplace ls'")
click.echo('Ingesting data: {}'.format(dataset), sys.stdout)
marketplace = Marketplace()
marketplace.ingest(dataset, data_frequency, start, end)
@marketplace.command()
@click.option(
'--dataset',
default=None,
help='The name of the dataset to ingest from the Data Marketplace.',
)
@click.pass_context
def clean(ctx, dataset):
if dataset is None:
ctx.fail("must specify a dataset to ingest with '--dataset'\n"
"List available dataset on the marketplace with "
"'catalyst marketplace ls'")
click.echo('Cleaning data source: {}'.format(dataset), sys.stdout)
marketplace = Marketplace()
marketplace.clean(dataset)
click.echo('Done', sys.stdout)
@marketplace.command()
@click.pass_context
def register(ctx):
marketplace = Marketplace()
marketplace.register()
@marketplace.command()
@click.option(
'--dataset',
default=None,
help='The name of the Marketplace dataset to publish data for.',
)
@click.option(
'--datadir',
default=None,
help='The folder that contains the CSV data files to publish.',
)
@click.option(
'--watch/--no-watch',
is_flag=True,
default=False,
help='Whether to watch the datadir for live data.',
)
@click.pass_context
def publish(ctx, dataset, datadir, watch):
marketplace = Marketplace()
if dataset is None:
ctx.fail("must specify a dataset to publish data for "
" with '--dataset'\n")
if datadir is None:
ctx.fail("must specify a datadir where to find the files to publish "
" with '--datadir'\n")
marketplace.publish(dataset, datadir, watch)
if __name__ == '__main__':
+5 -5
View File
@@ -939,7 +939,7 @@ class TradingAlgorithm(object):
The field to query. The options have the following meanings:
arena : str
The arena from the simulation parameters. This will normally
be ``'backtest'`` but some systems may use this distinguish
be ``backtest`` but some systems may use this distinguish
live trading from backtesting.
data_frequency : {'daily', 'minute'}
data_frequency tells the algorithm if it is running with
@@ -954,7 +954,7 @@ class TradingAlgorithm(object):
The platform that the code is running on. By default this
will be the string 'catalyst'. This can allow algorithms to
know if they are running on the Quantopian platform instead.
* : dict[str -> any]
\* : dict[str -> any]
Returns all of the fields in a dictionary.
Returns
@@ -1032,7 +1032,7 @@ class TradingAlgorithm(object):
argument is the name of the column in the preprocessed dataframe
containing the symbols. This will be used along with the date
information to map the sids in the asset finder.
**kwargs
\*\*kwargs
Forwarded to :func:`pandas.read_csv`.
Returns
@@ -1156,7 +1156,7 @@ class TradingAlgorithm(object):
Parameters
----------
**kwargs
\*\*kwargs
The names and values to record.
Notes
@@ -1273,7 +1273,7 @@ class TradingAlgorithm(object):
Parameters
----------
*args : iterable[str]
\*args : iterable[str]
The ticker symbols to lookup.
Returns
+65 -8
View File
@@ -34,6 +34,7 @@ def attach_pipeline(pipeline, name, chunks=None):
:func:`catalyst.api.pipeline_output`
"""
def batch_market_order(share_counts):
"""Place a batch market order for multiple assets.
@@ -48,6 +49,7 @@ def batch_market_order(share_counts):
Index of ids for newly-created orders.
"""
def cancel_order(order_param):
"""Cancel an open order.
@@ -57,7 +59,9 @@ def cancel_order(order_param):
The order_id or order object to cancel.
"""
def continuous_future(root_symbol_str, offset=0, roll='volume', adjustment='mul'):
def continuous_future(root_symbol_str, offset=0, roll='volume',
adjustment='mul'):
"""Create a specifier for a continuous contract.
Parameters
@@ -81,7 +85,10 @@ def continuous_future(root_symbol_str, offset=0, roll='volume', adjustment='mul'
The continuous future specifier.
"""
def fetch_csv(url, pre_func=None, post_func=None, date_column='date', date_format=None, timezone='UTC', symbol=None, mask=True, symbol_column=None, special_params_checker=None, **kwargs):
def fetch_csv(url, pre_func=None, post_func=None, date_column='date',
date_format=None, timezone='UTC', symbol=None, mask=True,
symbol_column=None, special_params_checker=None, **kwargs):
"""Fetch a csv from a remote url and register the data so that it is
queryable from the ``data`` object.
@@ -125,6 +132,7 @@ def fetch_csv(url, pre_func=None, post_func=None, date_column='date', date_forma
A requests source that will pull data from the url specified.
"""
def future_symbol(symbol):
"""Lookup a futures contract with a given symbol.
@@ -144,6 +152,7 @@ def future_symbol(symbol):
Raised when no contract named 'symbol' is found.
"""
def get_datetime(tz=None):
"""
Returns the current simulation datetime.
@@ -159,6 +168,7 @@ dt : datetime
The current simulation datetime converted to ``tz``.
"""
def get_environment(field='platform'):
"""Query the execution environment.
@@ -198,6 +208,7 @@ def get_environment(field='platform'):
Raised when ``field`` is not a valid option.
"""
def get_order(order_id):
"""Lookup an order based on the order id returned from one of the
order functions.
@@ -213,10 +224,12 @@ def get_order(order_id):
The order object.
"""
def history(bar_count, frequency, field, ffill=True):
"""DEPRECATED: use ``data.history`` instead.
"""
def order(asset, amount, limit_price=None, stop_price=None, style=None):
"""Place an order.
@@ -258,7 +271,9 @@ def order(asset, amount, limit_price=None, stop_price=None, style=None):
:func:`catalyst.api.order_percent`
"""
def order_percent(asset, percent, limit_price=None, stop_price=None, style=None):
def order_percent(asset, percent, limit_price=None, stop_price=None,
style=None):
"""Place an order in the specified asset corresponding to the given
percent of the current portfolio value.
@@ -293,6 +308,7 @@ def order_percent(asset, percent, limit_price=None, stop_price=None, style=None)
:func:`catalyst.api.order_value`
"""
def order_target(asset, target, limit_price=None, stop_price=None, style=None):
"""Place an order to adjust a position to a target number of shares. If
the position doesn't already exist, this is equivalent to placing a new
@@ -344,7 +360,9 @@ def order_target(asset, target, limit_price=None, stop_price=None, style=None):
:func:`catalyst.api.order_target_value`
"""
def order_target_percent(asset, target, limit_price=None, stop_price=None, style=None):
def order_target_percent(asset, target, limit_price=None, stop_price=None,
style=None):
"""Place an order to adjust a position to a target percent of the
current portfolio value. If the position doesn't already exist, this is
equivalent to placing a new order. If the position does exist, this is
@@ -396,7 +414,9 @@ def order_target_percent(asset, target, limit_price=None, stop_price=None, style
:func:`catalyst.api.order_target_value`
"""
def order_target_value(asset, target, limit_price=None, stop_price=None, style=None):
def order_target_value(asset, target, limit_price=None, stop_price=None,
style=None):
"""Place an order to adjust a position to a target value. If
the position doesn't already exist, this is equivalent to placing a new
order. If the position does exist, this is equivalent to placing an
@@ -448,6 +468,7 @@ def order_target_value(asset, target, limit_price=None, stop_price=None, style=N
:func:`catalyst.api.order_target_percent`
"""
def order_value(asset, value, limit_price=None, stop_price=None, style=None):
"""Place an order by desired value rather than desired number of
shares.
@@ -488,6 +509,7 @@ def order_value(asset, value, limit_price=None, stop_price=None, style=None):
:func:`catalyst.api.order_percent`
"""
def pipeline_output(name):
"""Get the results of the pipeline that was attached with the name:
``name``.
@@ -514,6 +536,7 @@ def pipeline_output(name):
:meth:`catalyst.pipeline.engine.PipelineEngine.run_pipeline`
"""
def record(*args, **kwargs):
"""Track and record values each day.
@@ -529,7 +552,9 @@ def record(*args, **kwargs):
:func:`~catalyst.run_algorithm`.
"""
def schedule_function(func, date_rule=None, time_rule=None, half_days=True, calendar=None):
def schedule_function(func, date_rule=None, time_rule=None, half_days=True,
calendar=None):
"""Schedules a function to be called according to some timed rules.
Parameters
@@ -549,6 +574,7 @@ def schedule_function(func, date_rule=None, time_rule=None, half_days=True, cale
:class:`catalyst.api.time_rules`
"""
def set_asset_restrictions(restrictions, on_error='fail'):
"""Set a restriction on which assets can be ordered.
@@ -562,6 +588,7 @@ def set_asset_restrictions(restrictions, on_error='fail'):
catalyst.finance.asset_restrictions.Restrictions
"""
def set_benchmark(benchmark):
"""Set the benchmark asset.
@@ -576,6 +603,7 @@ def set_benchmark(benchmark):
automatically reinvested.
"""
def set_cancel_policy(cancel_policy):
"""Sets the order cancellation policy for the simulation.
@@ -590,6 +618,7 @@ def set_cancel_policy(cancel_policy):
:class:`catalyst.api.NeverCancel`
"""
def set_commission(commission):
"""Sets the commission model for the simulation.
@@ -605,6 +634,7 @@ def set_commission(commission):
:class:`catalyst.finance.commission.PerDollar`
"""
def set_do_not_order_list(restricted_list, on_error='fail'):
"""Set a restriction on which assets can be ordered.
@@ -614,11 +644,13 @@ def set_do_not_order_list(restricted_list, on_error='fail'):
The assets that cannot be ordered.
"""
def set_long_only(on_error='fail'):
"""Set a rule specifying that this algorithm cannot take short
positions.
"""
def set_max_leverage(max_leverage):
"""Set a limit on the maximum leverage of the algorithm.
@@ -629,6 +661,7 @@ def set_max_leverage(max_leverage):
be no maximum.
"""
def set_max_order_count(max_count, on_error='fail'):
"""Set a limit on the number of orders that can be placed in a single
day.
@@ -639,7 +672,9 @@ def set_max_order_count(max_count, on_error='fail'):
The maximum number of orders that can be placed on any single day.
"""
def set_max_order_size(asset=None, max_shares=None, max_notional=None, on_error='fail'):
def set_max_order_size(asset=None, max_shares=None, max_notional=None,
on_error='fail'):
"""Set a limit on the number of shares and/or dollar value of any single
order placed for sid. Limits are treated as absolute values and are
enforced at the time that the algo attempts to place an order for sid.
@@ -658,7 +693,9 @@ def set_max_order_size(asset=None, max_shares=None, max_notional=None, on_error=
The maximum value that can be ordered at one time.
"""
def set_max_position_size(asset=None, max_shares=None, max_notional=None, on_error='fail'):
def set_max_position_size(asset=None, max_shares=None, max_notional=None,
on_error='fail'):
"""Set a limit on the number of shares and/or dollar value held for the
given sid. Limits are treated as absolute values and are enforced at
the time that the algo attempts to place an order for sid. This means
@@ -681,6 +718,7 @@ def set_max_position_size(asset=None, max_shares=None, max_notional=None, on_err
The maximum value to hold for an asset.
"""
def set_slippage(slippage):
"""Set the slippage model for the simulation.
@@ -694,6 +732,7 @@ def set_slippage(slippage):
:class:`catalyst.finance.slippage.SlippageModel`
"""
def set_symbol_lookup_date(dt):
"""Set the date for which symbols will be resolved to their assets
(symbols may map to different firms or underlying assets at
@@ -705,6 +744,7 @@ def set_symbol_lookup_date(dt):
The new symbol lookup date.
"""
def sid(sid):
"""Lookup an Asset by its unique asset identifier.
@@ -724,6 +764,7 @@ def sid(sid):
When a requested ``sid`` does not map to any asset.
"""
def symbol(symbol_str):
"""Lookup an Equity by its ticker symbol.
@@ -748,6 +789,7 @@ def symbol(symbol_str):
:func:`catalyst.api.set_symbol_lookup_date`
"""
def symbols(*args):
"""Lookup multuple Equities as a list.
@@ -773,3 +815,18 @@ def symbols(*args):
:func:`catalyst.api.set_symbol_lookup_date`
"""
def get_dataset(ds_name, start=None, end=None):
"""
Lookup a data source from the marketplace
Parameters
----------
ds_name: str
start: pd.Timestamp
end: pd.Timestamp
Returns
-------
"""
+28
View File
@@ -15,4 +15,32 @@ SYMBOLS_URL = 'https://s3.amazonaws.com/enigmaco/catalyst-exchanges/' \
DATE_TIME_FORMAT = '%Y-%m-%d %H:%M'
DATE_FORMAT = '%Y-%m-%d'
try:
ROOT_DIR = os.path.dirname(os.path.abspath(__file__))
except Exception as e:
print('unable to get catalyst path: {}'.format(e))
AUTO_INGEST = False
AUTH_SERVER = 'https://data.enigma.co'
# TODO: switch to mainnet
ETH_REMOTE_NODE = 'https://ropsten.infura.io/'
# TODO: move to MASTER branch on github
MARKETPLACE_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/develop/catalyst/marketplace/' \
'contract_marketplace_address.txt'
MARKETPLACE_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/develop/catalyst/marketplace/' \
'contract_marketplace_abi.json'
# TODO: switch to mainnet
ENIGMA_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/catalyst/' \
'develop/catalyst/marketplace/' \
'contract_enigma_address.txt'
ENIGMA_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/develop/catalyst/marketplace/' \
'contract_enigma_abi.json'
+3 -3
View File
@@ -60,7 +60,7 @@ def _handle_data(context, data):
rsi=rsi,
)
orders = get_open_orders(context.asset)
orders = context.blotter.open_orders
if orders:
log.info('skipping bar until all open orders execute')
return
@@ -146,11 +146,11 @@ if __name__ == '__main__':
live = True
if live:
run_algorithm(
capital_base=0.001,
capital_base=1000,
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='binance',
exchange_name='bittrex',
live=True,
algo_namespace=algo_namespace,
base_currency='btc',
+38 -25
View File
@@ -20,8 +20,8 @@ def initialize(context):
def handle_data(context, data):
# define the windows for the moving averages
short_window = 50
long_window = 200
short_window = 2
long_window = 2
# Skip as many bars as long_window to properly compute the average
context.i += 1
@@ -32,16 +32,18 @@ def handle_data(context, data):
# moving average with the appropriate parameters. We choose to use
# minute bars for this simulation -> freq="1m"
# Returns a pandas dataframe.
short_mavg = data.history(context.asset,
short_data = data.history(context.asset,
'price',
bar_count=short_window,
frequency="1m",
).mean()
long_mavg = data.history(context.asset,
frequency="1T",
)
short_mavg = short_data.mean()
long_data = data.history(context.asset,
'price',
bar_count=long_window,
frequency="1m",
).mean()
frequency="1T",
)
long_mavg = long_data.mean()
# Let's keep the price of our asset in a more handy variable
price = data.current(context.asset, 'price')
@@ -82,7 +84,6 @@ def handle_data(context, data):
def analyze(context, perf):
# Get the base_currency that was passed as a parameter to the simulation
exchange = list(context.exchanges.values())[0]
base_currency = exchange.base_currency.upper()
@@ -93,7 +94,7 @@ def analyze(context, perf):
ax1.legend_.remove()
ax1.set_ylabel('Portfolio Value\n({})'.format(base_currency))
start, end = ax1.get_ylim()
ax1.yaxis.set_ticks(np.arange(start, end, (end-start)/5))
ax1.yaxis.set_ticks(np.arange(start, end, (end - start) / 5))
# Second chart: Plot asset price, moving averages and buys/sells
ax2 = plt.subplot(412, sharex=ax1)
@@ -104,9 +105,9 @@ def analyze(context, perf):
ax2.set_ylabel('{asset}\n({base})'.format(
asset=context.asset.symbol,
base=base_currency
))
))
start, end = ax2.get_ylim()
ax2.yaxis.set_ticks(np.arange(start, end, (end-start)/5))
ax2.yaxis.set_ticks(np.arange(start, end, (end - start) / 5))
transaction_df = extract_transactions(perf)
if not transaction_df.empty:
@@ -136,28 +137,40 @@ def analyze(context, perf):
ax3.legend_.remove()
ax3.set_ylabel('Percent Change')
start, end = ax3.get_ylim()
ax3.yaxis.set_ticks(np.arange(start, end, (end-start)/5))
ax3.yaxis.set_ticks(np.arange(start, end, (end - start) / 5))
# Fourth chart: Plot our cash
ax4 = plt.subplot(414, sharex=ax1)
perf.cash.plot(ax=ax4)
ax4.set_ylabel('Cash\n({})'.format(base_currency))
start, end = ax4.get_ylim()
ax4.yaxis.set_ticks(np.arange(0, end, end/5))
ax4.yaxis.set_ticks(np.arange(0, end, end / 5))
plt.show()
if __name__ == '__main__':
run_algorithm(
capital_base=1000,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bitfinex',
algo_namespace=NAMESPACE,
base_currency='usd',
start=pd.to_datetime('2017-9-22', utc=True),
end=pd.to_datetime('2017-9-23', utc=True),
)
capital_base=1000,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bitfinex',
algo_namespace=NAMESPACE,
base_currency='usd',
simulate_orders=True,
live=True,
)
# run_algorithm(
# capital_base=1000,
# data_frequency='minute',
# initialize=initialize,
# handle_data=handle_data,
# analyze=analyze,
# exchange_name='bitfinex',
# algo_namespace=NAMESPACE,
# base_currency='usd',
# start=pd.to_datetime('2017-9-22', utc=True),
# end=pd.to_datetime('2017-9-23', utc=True),
# )
@@ -0,0 +1,237 @@
# 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 pandas as pd
import talib
from logbook import Logger
from catalyst import run_algorithm
from catalyst.api import symbol, record, order_target_percent, get_dataset
from catalyst.exchange.utils.stats_utils import set_print_settings, \
get_pretty_stats
# 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.
df = get_dataset('testmarketcap2') # type: pd.DataFrame
# Picking a specific date in our DataFrame
first_dt = df.index.get_level_values(0)[0]
# Since we use a MultiIndex with date / symbol, picking a date will
# result in a new DataFrame for the selected date with a single
# symbol index
df = df.xs(first_dt, level=0)
# Keep only the top coins by market cap
df = df.loc[df['market_cap_usd'].isin(df['market_cap_usd'].nlargest(100))]
set_print_settings()
df.sort_values(by=['market_cap_usd'], ascending=True, inplace=True)
print('the marketplace data:\n{}'.format(df))
# Pick the 5 assets with the lowest market cap for trading
quote_currency = 'eth'
exchange = context.exchanges[next(iter(context.exchanges))]
symbols = [a.symbol for a in exchange.assets
if a.start_date < context.datetime]
context.assets = []
for currency, price in df['market_cap_usd'].iteritems():
if len(context.assets) >= 5:
break
s = '{}_{}'.format(currency.decode('utf-8'), quote_currency)
if s in symbols:
context.assets.append(symbol(s))
context.base_price = None
context.current_day = None
context.RSI_OVERSOLD = 55
context.RSI_OVERBOUGHT = 60
context.CANDLE_SIZE = '5T'
context.start_time = time.time()
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 = dict()
context.current_day = today
# Preparing dictionaries for asset-level data points
volumes = dict()
rsis = dict()
price_values = dict()
cash = context.portfolio.cash
for asset in context.assets:
# We're computing the volume-weighted-average-price of the security
# defined above, in the context.assets 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(
asset,
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(asset, 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 asset not in context.base_price:
# context.base_price[asset] = price
#
# base_price = context.base_price[asset]
# price_change = (price - base_price) / base_price
# Tracking the relevant data
volumes[asset] = current['volume']
rsis[asset] = rsi[-1]
price_values[asset] = price
# price_changes[asset] = price_change
# We are trying to avoid over-trading by limiting our trades to
# one per day.
if asset in context.traded_today:
continue
# Exit if we cannot trade
if not data.can_trade(asset):
continue
# 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[asset].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
target = 1.0 / len(context.assets)
order_target_percent(
asset, target, limit_price=limit_price
)
context.traded_today[asset] = 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(
asset, 0, limit_price=limit_price
)
context.traded_today[asset] = True
# 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(
current_price=price_values,
volume=volumes,
rsi=rsis,
cash=cash,
)
def analyze(context=None, perf=None):
stats = get_pretty_stats(perf)
print('the algo stats:\n{}'.format(stats))
pass
if __name__ == '__main__':
# The execution mode: backtest or live
live = False
if live:
run_algorithm(
capital_base=0.1,
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=100,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='poloniex',
algo_namespace=NAMESPACE,
base_currency='eth',
start=pd.to_datetime('2017-10-01', utc=True),
end=pd.to_datetime('2017-10-15', utc=True),
)
log.info('saved perf stats: {}'.format(out))
+6 -6
View File
@@ -33,11 +33,11 @@ def initialize(context):
# parameters or values you're going to use.
# In our example, we're looking at Neo in Ether.
context.market = symbol('eth_btc')
context.market = symbol('bnb_eth')
context.base_price = None
context.current_day = None
context.RSI_OVERSOLD = 55
context.RSI_OVERSOLD = 40
context.RSI_OVERBOUGHT = 60
context.CANDLE_SIZE = '15T'
@@ -248,14 +248,14 @@ if __name__ == '__main__':
if live:
run_algorithm(
capital_base=0.01,
capital_base=0.1,
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='poloniex',
exchange_name='binance',
live=True,
algo_namespace=NAMESPACE,
base_currency='btc',
base_currency='eth',
live_graph=False,
simulate_orders=False,
stats_output=None,
@@ -274,7 +274,7 @@ if __name__ == '__main__':
# -x bitfinex -s 2017-10-1 -e 2017-11-10 -c usdt -n mean-reversion \
# --data-frequency minute --capital-base 10000
run_algorithm(
capital_base=0.1,
capital_base=0.035,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
+102 -56
View File
@@ -425,15 +425,12 @@ class CCXT(Exchange):
'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']:
if start_dt is None:
# TODO: determine why binance is failing
if end_dt is None and self.name not in ['binance']:
end_dt = pd.Timestamp.utcnow()
if end_dt is not None:
dt_range = get_periods_range(
end_dt=end_dt,
periods=bar_count,
@@ -441,10 +438,13 @@ class CCXT(Exchange):
)
start_dt = dt_range[0]
since = None
if start_dt is not None:
# Convert out start date to a UNIX timestamp, then translate to
# milliseconds
delta = start_dt - get_epoch()
since = int(delta.total_seconds()) * 1000
else:
since = None
candles = dict()
for index, asset in enumerate(assets):
@@ -760,24 +760,23 @@ class CCXT(Exchange):
side = 'buy' if amount > 0 else 'sell'
if hasattr(self.api, 'amount_to_lots'):
adj_amount = self.api.amount_to_lots(
symbol=symbol,
amount=abs(amount),
)
if adj_amount != abs(amount):
log.info(
'adjusted order amount {} to {} based on lot size'.format(
abs(amount), adj_amount,
# TODO: is this right?
if self.api.markets is None:
self.api.load_markets()
# https://github.com/ccxt/ccxt/issues/1483
adj_amount = round(abs(amount), asset.decimals)
market = self.api.markets[symbol]
if 'lots' in market and market['lots'] > amount:
raise CreateOrderError(
exchange=self.name,
e='order amount lower than the smallest lot: {}'.format(
amount
)
)
else:
adj_amount = abs(amount)
if adj_amount == 0:
raise CreateOrderError(
exchange=self.name,
e='order amount lower than the smallest lot: {}'.format(amount)
)
else:
adj_amount = round(abs(amount), asset.decimals)
try:
result = self.api.create_order(
@@ -799,6 +798,22 @@ class CCXT(Exchange):
)
raise ExchangeRequestError(error=e)
exchange_amount = None
if 'amount' in result and result['amount'] != adj_amount:
exchange_amount = result['amount']
elif 'info' in result:
if 'origQty' in result['info']:
exchange_amount = float(result['info']['origQty'])
if exchange_amount:
log.info(
'order amount adjusted by {} from {} to {}'.format(
self.name, adj_amount, exchange_amount
)
)
adj_amount = exchange_amount
if 'info' not in result:
raise ValueError('cannot use order without info attribute')
@@ -859,31 +874,38 @@ class CCXT(Exchange):
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
order.filled = exc_order.amount
if order.status == ORDER_STATUS.FILLED:
transactions = []
if exc_order.status == ORDER_STATUS.FILLED:
if order.amount > exc_order.amount:
log.warn(
'executed order amount {} differs '
'from original'.format(
exc_order.amount, order.amount
)
)
order.check_triggers(
price=price,
dt=exc_order.dt,
)
transaction = Transaction(
asset=order.asset,
amount=order.amount,
dt=pd.Timestamp.utcnow(),
price=price,
order_id=order.id,
commission=order.commission
commission=order.commission,
)
return [transaction]
transactions.append(transaction)
return transactions
def process_order(self, order):
# TODO: move to parent class after tracking features in the parent
if not self.api.hasFetchMyTrades:
if not self.api.has['fetchMyTrades']:
return self._process_order_fallback(order)
try:
@@ -985,7 +1007,7 @@ class CCXT(Exchange):
)
raise ExchangeRequestError(error=e)
def tickers(self, assets):
def tickers(self, assets, on_ticker_error='raise'):
"""
Retrieve current tick data for the given assets
@@ -998,27 +1020,51 @@ class CCXT(Exchange):
list[dict[str, float]
"""
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
if len(assets) == 1:
try:
ticker = self.api.fetch_ticker(symbol=symbol)
except (ExchangeError, NetworkError) as e:
symbol = self.get_symbol(assets[0])
log.debug('fetching single ticker: {}'.format(symbol))
results = dict()
results[symbol] = self.api.fetch_ticker(symbol=symbol)
except (ExchangeError, NetworkError,) as e:
log.warn(
'unable to fetch ticker {} / {}: {}'.format(
self.name, asset.symbol, e
self.name, symbol, e
)
)
continue
raise ExchangeRequestError(error=e)
elif len(assets) > 1:
symbols = self.get_symbols(assets)
try:
log.debug('fetching multiple tickers: {}'.format(symbols))
results = self.api.fetch_tickers(symbols=symbols)
except (ExchangeError, NetworkError) as e:
log.warn(
'unable to fetch tickers {} / {}: {}'.format(
self.name, symbols, e
)
)
raise ExchangeRequestError(error=e)
else:
raise ValueError('Cannot request tickers with not assets.')
tickers = dict()
for asset in assets:
symbol = self.get_symbol(asset)
if symbol not in results:
msg = 'ticker not found {} / {}'.format(
self.name, symbol
)
log.warn(msg)
if on_ticker_error == 'warn':
continue
else:
raise ExchangeRequestError(error=msg)
ticker = results[symbol]
ticker['last_traded'] = from_ms_timestamp(ticker['timestamp'])
if 'last_price' not in ticker:
@@ -1030,7 +1076,7 @@ class CCXT(Exchange):
ticker['volume'] = ticker['baseVolume']
elif 'info' in ticker and 'bidQty' in ticker['info'] \
and 'askQty' in ticker['info']:
and 'askQty' in ticker['info']:
ticker['volume'] = float(ticker['info']['bidQty']) + \
float(ticker['info']['askQty'])
@@ -1068,7 +1114,7 @@ class CCXT(Exchange):
return result
def get_trades(self, asset, my_trades=True, start_dt=None, limit=None):
def get_trades(self, asset, my_trades=True, start_dt=None, limit=100):
if not my_trades:
raise NotImplemented(
'get_trades only supports "my trades"'
+6 -3
View File
@@ -178,6 +178,7 @@ class Exchange:
if symbols is None:
# Make a distinct list of all symbols
symbols = list(set([asset.symbol for asset in self.assets]))
symbols.sort()
if quote_currency is not None:
for symbol in symbols[:]:
@@ -701,8 +702,8 @@ class Exchange:
)
positions_value = 0.0
if positions is not None:
assets = set([position.asset for position in positions])
if positions:
assets = list(set([position.asset for position in positions]))
tickers = self.tickers(assets)
for position in positions:
@@ -972,13 +973,15 @@ class Exchange:
pass
@abc.abstractmethod
def tickers(self, assets):
def tickers(self, assets, on_ticker_error='raise'):
"""
Retrieve current tick data for the given assets
Parameters
----------
assets: list[TradingPair]
on_ticker_error: str [raise|warn]
How to handle an error when retrieving a single ticker.
Returns
-------
+24 -5
View File
@@ -43,6 +43,7 @@ from catalyst.finance.execution import MarketOrder
from catalyst.finance.performance import PerformanceTracker
from catalyst.finance.performance.period import calc_period_stats
from catalyst.gens.tradesimulation import AlgorithmSimulator
from catalyst.marketplace.marketplace import Marketplace
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
@@ -67,7 +68,7 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
self.current_day = None
if self.simulate_orders is None \
and self.sim_params.arena == 'backtest':
and self.sim_params.arena == 'backtest':
self.simulate_orders = True
# Operations with retry features
@@ -92,6 +93,8 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
attempts=self.attempts,
)
self._marketplace = None
@staticmethod
def __convert_order_params_for_blotter(limit_price, stop_price, style):
"""
@@ -115,7 +118,7 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
# be in-line with CXXT and many exchanges. We'll consider
# adding more order types in the future.
if not isinstance(style, ExchangeLimitOrder) or \
not isinstance(style, MarketOrder):
not isinstance(style, MarketOrder):
raise OrderTypeNotSupported(
order_type=style.__class__.__name__
)
@@ -167,6 +170,15 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
"""
return round_nearest(amount, asset.min_trade_size)
@api_method
def get_dataset(self, data_source_name, start=None, end=None):
if self._marketplace is None:
self._marketplace = Marketplace()
return self._marketplace.get_dataset(
data_source_name, start, end,
)
@api_method
@preprocess(symbol_str=ensure_upper_case)
def symbol(self, symbol_str, exchange_name=None):
@@ -379,8 +391,6 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
log.warn("Can't initialize signal handler inside another thread."
"Exit should be handled by the user.")
log.info('initialized trading algorithm in live mode')
def interrupt_algorithm(self):
self.is_running = False
@@ -862,6 +872,13 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
raise NotImplementedError()
def _get_open_orders(self, asset=None):
if self.simulate_orders:
raise ValueError(
'The get_open_orders() method only works in live mode. '
'The purpose is to list open orders on the exchange '
'regardless who placed them. To list the open orders of '
'this algo, use `context.blotter.open_orders`.'
)
if asset:
exchange = self.exchanges[asset.exchange]
return exchange.get_open_orders(asset)
@@ -895,13 +912,15 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
If an asset is passed then this will return a list of the open
orders for this asset.
"""
# TODO: should this be a shortcut to the open orders in the blotter?
return retry(
action=self._get_open_orders,
attempts=self.attempts['get_open_orders_attempts'],
sleeptime=self.attempts['retry_sleeptime'],
retry_exceptions=(ExchangeRequestError,),
cleanup=lambda: log.warn('Fetching open orders again.'),
args=(asset,))
args=(asset,)
)
@api_method
def get_order(self, order_id, exchange_name):
+1 -1
View File
@@ -214,7 +214,7 @@ class ExchangeBlotter(Blotter):
# 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:
if transactions and order.status == ORDER_STATUS.FILLED:
avg_price = np.average(
a=[t.price for t in transactions],
weights=[t.amount for t in transactions],
+16 -25
View File
@@ -458,7 +458,7 @@ class ExchangeBundle:
last_entry = None
if start is None or \
(earliest_trade is not None and earliest_trade > start):
(earliest_trade is not None and earliest_trade > start):
start = earliest_trade
if last_entry is not None and (end is None or end > last_entry):
@@ -600,14 +600,14 @@ class ExchangeBundle:
if show_breakdown:
for asset in chunks:
with maybe_show_progress(
chunks[asset],
show_progress,
label='Ingesting {frequency} price data for '
'{symbol} on {exchange}'.format(
exchange=self.exchange_name,
frequency=data_frequency,
symbol=asset.symbol
)) as it:
chunks[asset],
show_progress,
label='Ingesting {frequency} price data for '
'{symbol} on {exchange}'.format(
exchange=self.exchange_name,
frequency=data_frequency,
symbol=asset.symbol
)) as it:
for chunk in it:
problems += self.ingest_ctable(
asset=chunk['asset'],
@@ -625,13 +625,13 @@ class ExchangeBundle:
key=lambda chunk: pd.to_datetime(chunk['period'])
)
with maybe_show_progress(
all_chunks,
show_progress,
label='Ingesting {frequency} price data on '
'{exchange}'.format(
exchange=self.exchange_name,
frequency=data_frequency,
)) as it:
all_chunks,
show_progress,
label='Ingesting {frequency} price data on '
'{exchange}'.format(
exchange=self.exchange_name,
frequency=data_frequency,
)) as it:
for chunk in it:
problems += self.ingest_ctable(
asset=chunk['asset'],
@@ -830,7 +830,6 @@ class ExchangeBundle:
field,
data_frequency,
algo_end_dt=None,
trailing_bar_count=None,
force_auto_ingest=False
):
"""
@@ -858,7 +857,6 @@ class ExchangeBundle:
bar_count=bar_count,
field=field,
data_frequency=data_frequency,
trailing_bar_count=trailing_bar_count,
)
return pd.DataFrame(series)
@@ -887,7 +885,6 @@ class ExchangeBundle:
field=field,
data_frequency=data_frequency,
reset_reader=True,
trailing_bar_count=trailing_bar_count,
)
return series
@@ -898,7 +895,6 @@ class ExchangeBundle:
bar_count=bar_count,
field=field,
data_frequency=data_frequency,
trailing_bar_count=trailing_bar_count,
)
return pd.DataFrame(series)
@@ -962,12 +958,7 @@ class ExchangeBundle:
bar_count,
field,
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
@@ -298,7 +298,6 @@ class DataPortalExchangeBacktest(DataPortalExchangeBase):
frequency, data_frequency
)
adj_bar_count = candle_size * bar_count
trailing_bar_count = candle_size - 1
if data_frequency == 'minute' and adj_data_frequency == 'daily':
end_dt = end_dt.floor('1D')
@@ -310,7 +309,6 @@ class DataPortalExchangeBacktest(DataPortalExchangeBase):
field=field,
data_frequency=adj_data_frequency,
algo_end_dt=self._last_available_session,
trailing_bar_count=trailing_bar_count,
)
df = resample_history_df(pd.DataFrame(series), freq, field)
+1 -1
View File
@@ -540,7 +540,7 @@ def resample_history_df(df, freq, field):
else:
raise ValueError('Invalid field.')
resampled_df = df.resample(freq).agg(agg)
resampled_df = df.resample(freq, closed='left', label='left').agg(agg)
return resampled_df
+8 -10
View File
@@ -44,7 +44,7 @@ def crossover(source, target):
"""
if isinstance(target, numbers.Number):
if source[-1] is np.nan or source[-2] is np.nan \
or target is np.nan:
or target is np.nan:
return False
if source[-1] >= target > source[-2]:
@@ -54,7 +54,7 @@ def crossover(source, target):
else:
if source[-1] is np.nan or source[-2] is np.nan \
or target[-1] is np.nan or target[-2] is np.nan:
or target[-1] is np.nan or target[-2] is np.nan:
return False
if source[-1] > target[-1] and source[-2] < target[-2]:
@@ -81,7 +81,7 @@ def crossunder(source, target):
"""
if isinstance(target, numbers.Number):
if source[-1] is np.nan or source[-2] is np.nan \
or target is np.nan:
or target is np.nan:
return False
if source[-1] < target <= source[-2]:
@@ -90,7 +90,7 @@ def crossunder(source, target):
return False
else:
if source[-1] is np.nan or source[-2] is np.nan \
or target[-1] is np.nan or target[-2] is np.nan:
or target[-1] is np.nan or target[-2] is np.nan:
return False
if source[-1] < target[-1] and source[-2] >= target[-2]:
@@ -229,7 +229,10 @@ def prepare_stats(stats, recorded_cols=list()):
asset_values)
df = pd.DataFrame(stats)
df['orders'] = df['orders'].apply(lambda orders: len(orders))
df['transactions'] = df['transactions'].apply(
lambda transactions: len(transactions)
)
index_cols = [
'period_close', 'starting_cash', 'ending_cash', 'portfolio_value',
'pnl', 'long_exposure', 'short_exposure', 'orders', 'transactions',
@@ -241,11 +244,6 @@ def prepare_stats(stats, recorded_cols=list()):
for column in recorded_cols:
index_cols.append(column)
df['orders'] = df['orders'].apply(lambda orders: len(orders))
df['transactions'] = df['transactions'].apply(
lambda transactions: len(transactions)
)
if asset_cols:
columns = asset_cols
df.set_index(index_cols, drop=True, inplace=True)
+4 -2
View File
@@ -29,13 +29,15 @@ from .risk import check_entry
from empyrical import (
alpha_beta_aligned,
annual_volatility,
cum_returns,
downside_risk,
information_ratio,
max_drawdown,
sharpe_ratio,
sortino_ratio
)
from catalyst.patches.stats import (
max_drawdown,
cum_returns,
)
from catalyst.constants import LOG_LEVEL
View File
@@ -0,0 +1,302 @@
[
{
"constant": true,
"inputs": [],
"name": "name",
"outputs": [
{
"name": "",
"type": "string"
}
],
"payable": false,
"stateMutability": "view",
"type": "function"
},
{
"constant": false,
"inputs": [
{
"name": "_spender",
"type": "address"
},
{
"name": "_value",
"type": "uint256"
}
],
"name": "approve",
"outputs": [
{
"name": "",
"type": "bool"
}
],
"payable": false,
"stateMutability": "nonpayable",
"type": "function"
},
{
"constant": true,
"inputs": [],
"name": "totalSupply",
"outputs": [
{
"name": "",
"type": "uint256"
}
],
"payable": false,
"stateMutability": "view",
"type": "function"
},
{
"constant": false,
"inputs": [
{
"name": "_from",
"type": "address"
},
{
"name": "_to",
"type": "address"
},
{
"name": "_value",
"type": "uint256"
}
],
"name": "transferFrom",
"outputs": [
{
"name": "",
"type": "bool"
}
],
"payable": false,
"stateMutability": "nonpayable",
"type": "function"
},
{
"constant": true,
"inputs": [],
"name": "INITIAL_SUPPLY",
"outputs": [
{
"name": "",
"type": "uint256"
}
],
"payable": false,
"stateMutability": "view",
"type": "function"
},
{
"constant": true,
"inputs": [],
"name": "decimals",
"outputs": [
{
"name": "",
"type": "uint8"
}
],
"payable": false,
"stateMutability": "view",
"type": "function"
},
{
"constant": false,
"inputs": [
{
"name": "_spender",
"type": "address"
},
{
"name": "_subtractedValue",
"type": "uint256"
}
],
"name": "decreaseApproval",
"outputs": [
{
"name": "success",
"type": "bool"
}
],
"payable": false,
"stateMutability": "nonpayable",
"type": "function"
},
{
"constant": false,
"inputs": [],
"name": "getAfterApproveTest",
"outputs": [
{
"name": "",
"type": "uint256"
}
],
"payable": false,
"stateMutability": "nonpayable",
"type": "function"
},
{
"constant": true,
"inputs": [
{
"name": "_owner",
"type": "address"
}
],
"name": "balanceOf",
"outputs": [
{
"name": "balance",
"type": "uint256"
}
],
"payable": false,
"stateMutability": "view",
"type": "function"
},
{
"constant": true,
"inputs": [],
"name": "symbol",
"outputs": [
{
"name": "",
"type": "string"
}
],
"payable": false,
"stateMutability": "view",
"type": "function"
},
{
"constant": false,
"inputs": [
{
"name": "_to",
"type": "address"
},
{
"name": "_value",
"type": "uint256"
}
],
"name": "transfer",
"outputs": [
{
"name": "",
"type": "bool"
}
],
"payable": false,
"stateMutability": "nonpayable",
"type": "function"
},
{
"constant": false,
"inputs": [
{
"name": "_spender",
"type": "address"
},
{
"name": "_addedValue",
"type": "uint256"
}
],
"name": "increaseApproval",
"outputs": [
{
"name": "success",
"type": "bool"
}
],
"payable": false,
"stateMutability": "nonpayable",
"type": "function"
},
{
"constant": true,
"inputs": [
{
"name": "_owner",
"type": "address"
},
{
"name": "_spender",
"type": "address"
}
],
"name": "allowance",
"outputs": [
{
"name": "",
"type": "uint256"
}
],
"payable": false,
"stateMutability": "view",
"type": "function"
},
{
"inputs": [
{
"name": "testValue",
"type": "address"
}
],
"payable": false,
"stateMutability": "nonpayable",
"type": "constructor"
},
{
"anonymous": false,
"inputs": [
{
"indexed": true,
"name": "owner",
"type": "address"
},
{
"indexed": true,
"name": "spender",
"type": "address"
},
{
"indexed": false,
"name": "value",
"type": "uint256"
}
],
"name": "Approval",
"type": "event"
},
{
"anonymous": false,
"inputs": [
{
"indexed": true,
"name": "from",
"type": "address"
},
{
"indexed": true,
"name": "to",
"type": "address"
},
{
"indexed": false,
"name": "value",
"type": "uint256"
}
],
"name": "Transfer",
"type": "event"
}
]
@@ -0,0 +1 @@
0x7fAec9aaE31BE428DeAAE1be8195dF609079Fd10
File diff suppressed because one or more lines are too long
@@ -0,0 +1 @@
0x3985f5de8fddf2e8f7705cd360b498bf35ebfbc4
+705
View File
@@ -0,0 +1,705 @@
from __future__ import print_function
import glob
import json
import os
import re
import shutil
import sys
import time
import bcolz
import logbook
import pandas as pd
import requests
from requests_toolbelt import MultipartDecoder
from requests_toolbelt.multipart.decoder import \
NonMultipartContentTypeException
from catalyst.constants import (
LOG_LEVEL, AUTH_SERVER, ETH_REMOTE_NODE, MARKETPLACE_CONTRACT,
MARKETPLACE_CONTRACT_ABI, ENIGMA_CONTRACT, ENIGMA_CONTRACT_ABI)
from catalyst.exchange.utils.stats_utils import set_print_settings
from catalyst.marketplace.marketplace_errors import (
MarketplacePubAddressEmpty, MarketplaceDatasetNotFound,
MarketplaceNoAddressMatch, MarketplaceHTTPRequest,
MarketplaceNoCSVFiles)
from catalyst.marketplace.utils.auth_utils import get_key_secret, \
get_signed_headers
from catalyst.marketplace.utils.bundle_utils import merge_bundles
from catalyst.marketplace.utils.eth_utils import bin_hex, from_grains, \
to_grains
from catalyst.marketplace.utils.path_utils import get_bundle_folder, \
get_data_source_folder, get_marketplace_folder, \
get_user_pubaddr, get_temp_bundles_folder, extract_bundle
if sys.version_info.major < 3:
import urllib
else:
import urllib.request as urllib
log = logbook.Logger('Marketplace', level=LOG_LEVEL)
class Marketplace:
def __init__(self):
global Web3
from web3 import Web3, HTTPProvider
self.addresses = get_user_pubaddr()
if self.addresses[0]['pubAddr'] == '':
raise MarketplacePubAddressEmpty(
filename=os.path.join(
get_marketplace_folder(), 'addresses.json')
)
self.default_account = self.addresses[0]['pubAddr']
self.web3 = Web3(HTTPProvider(ETH_REMOTE_NODE))
contract_url = urllib.urlopen(MARKETPLACE_CONTRACT)
self.mkt_contract_address = Web3.toChecksumAddress(
contract_url.readline().strip())
abi_url = urllib.urlopen(MARKETPLACE_CONTRACT_ABI)
abi = json.load(abi_url)
self.mkt_contract = self.web3.eth.contract(
self.mkt_contract_address,
abi=abi,
)
contract_url = urllib.urlopen(ENIGMA_CONTRACT)
self.eng_contract_address = Web3.toChecksumAddress(
contract_url.readline().strip())
abi_url = urllib.urlopen(ENIGMA_CONTRACT_ABI)
abi = json.load(abi_url)
self.eng_contract = self.web3.eth.contract(
self.eng_contract_address,
abi=abi,
)
# def get_data_sources_map(self):
# return [
# dict(
# name='Marketcap',
# desc='The marketcap value in USD.',
# start_date=pd.to_datetime('2017-01-01'),
# end_date=pd.to_datetime('2018-01-15'),
# data_frequencies=['daily'],
# ),
# dict(
# name='GitHub',
# desc='The rate of development activity on GitHub.',
# start_date=pd.to_datetime('2017-01-01'),
# end_date=pd.to_datetime('2018-01-15'),
# data_frequencies=['daily', 'hour'],
# ),
# dict(
# name='Influencers',
# desc='Tweets & related sentiments by selected influencers.',
# start_date=pd.to_datetime('2017-01-01'),
# end_date=pd.to_datetime('2018-01-15'),
# data_frequencies=['daily', 'hour', 'minute'],
# ),
# ]
def to_text(self, hex):
return Web3.toText(hex).rstrip('\0')
def choose_pubaddr(self):
if len(self.addresses) == 1:
address = self.addresses[0]['pubAddr']
address_i = 0
print('Using {} for this transaction.'.format(address))
else:
while True:
for i in range(0, len(self.addresses)):
print('{}\t{}\t{}'.format(
i,
self.addresses[i]['pubAddr'],
self.addresses[i]['desc'])
)
address_i = int(input('Choose your address associated with '
'this transaction: [default: 0] ') or 0)
if not (0 <= address_i < len(self.addresses)):
print('Please choose a number between 0 and {}\n'.format(
len(self.addresses) - 1))
else:
address = Web3.toChecksumAddress(
self.addresses[address_i]['pubAddr'])
break
return address, address_i
def sign_transaction(self, from_address, tx):
print('\nVisit https://www.myetherwallet.com/#offline-transaction and '
'enter the following parameters:\n\n'
'From Address:\t\t{_from}\n'
'\n\tClick the "Generate Information" button\n\n'
'To Address:\t\t{to}\n'
'Value / Amount to Send:\t{value}\n'
'Gas Limit:\t\t{gas}\n'
'Gas Price:\t\t[Accept the default value]\n'
'Nonce:\t\t\t{nonce}\n'
'Data:\t\t\t{data}\n'.format(
_from=from_address,
to=tx['to'],
value=tx['value'],
gas=tx['gas'],
nonce=tx['nonce'],
data=tx['data'], )
)
signed_tx = input('Copy and Paste the "Signed Transaction" '
'field here:\n')
if signed_tx.startswith('0x'):
signed_tx = signed_tx[2:]
return signed_tx
def check_transaction(self, tx_hash):
if 'ropsten' in ETH_REMOTE_NODE:
etherscan = 'https://ropsten.etherscan.io/tx/{}'.format(
tx_hash)
else:
etherscan = 'https://etherscan.io/tx/{}'.format(tx_hash)
print('\nYou can check the outcome of your transaction here:\n'
'{}\n\n'.format(etherscan))
def list(self):
data_sources = self.mkt_contract.functions.getAllProviders().call()
data = []
for index, data_source in enumerate(data_sources):
if index > 0:
if 'test' not in Web3.toText(data_source).lower():
data.append(
dict(
dataset=self.to_text(data_source)
)
)
df = pd.DataFrame(data)
set_print_settings()
if df.empty:
print('There are no datasets available yet.')
else:
print(df)
def subscribe(self, dataset):
dataset = dataset.lower()
address = self.choose_pubaddr()[0]
provider_info = self.mkt_contract.functions.getDataProviderInfo(
Web3.toHex(dataset)
).call()
if not provider_info[4]:
print('The requested "{}" dataset is not registered in '
'the Data Marketplace.'.format(dataset))
return
grains = provider_info[1]
price = from_grains(grains)
subscribed = self.mkt_contract.functions.checkAddressSubscription(
address, Web3.toHex(dataset)
).call()
if subscribed[5]:
print(
'\nYou are already subscribed to the "{}" dataset.\n'
'Your subscription started on {} UTC, and is valid until '
'{} UTC.'.format(
dataset,
pd.to_datetime(subscribed[3], unit='s', utc=True),
pd.to_datetime(subscribed[4], unit='s', utc=True)
)
)
return
print('\nThe price for a monthly subscription to this dataset is'
' {} ENG'.format(price))
print(
'Checking that the ENG balance in {} is greater than {} '
'ENG... '.format(address, price), end=''
)
wallet_address = address[2:]
balance = self.web3.eth.call({
'from': address,
'to': self.eng_contract_address,
'data': '0x70a08231000000000000000000000000{}'.format(
wallet_address
)
})
try:
balance = Web3.toInt(balance) # web3 >= 4.0.0b7
except TypeError:
balance = Web3.toInt(hexstr=balance) # web3 <= 4.0.0b6
if balance > grains:
print('OK.')
else:
print('FAIL.\n\nAddress {} balance is {} ENG,\nwhich is lower '
'than the price of the dataset that you are trying to\n'
'buy: {} ENG. Get enough ENG to cover the costs of the '
'monthly\nsubscription for what you are trying to buy, '
'and try again.'.format(
address, from_grains(balance), price))
return
while True:
agree_pay = input('Please confirm that you agree to pay {} ENG '
'for a monthly subscription to the dataset "{}" '
'starting today. [default: Y] '.format(
price, dataset)) or 'y'
if agree_pay.lower() not in ('y', 'n'):
print("Please answer Y or N.")
else:
if agree_pay.lower() == 'y':
break
else:
return
print('Ready to subscribe to dataset {}.\n'.format(dataset))
print('In order to execute the subscription, you will need to sign '
'two different transactions:\n'
'1. First transaction is to authorize the Marketplace contract '
'to spend {} ENG on your behalf.\n'
'2. Second transaction is the actual subscription for the '
'desired dataset'.format(price))
tx = self.eng_contract.functions.approve(
self.mkt_contract_address,
grains,
).buildTransaction(
{'nonce': self.web3.eth.getTransactionCount(address)}
)
if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(address, tx)
try:
tx_hash = '0x{}'.format(
bin_hex(self.web3.eth.sendRawTransaction(signed_tx))
)
print(
'\nThis is the TxHash for this transaction: {}'.format(tx_hash)
)
except Exception as e:
print('Unable to subscribe to data source: {}'.format(e))
return
self.check_transaction(tx_hash)
print('Waiting for the first transaction to succeed...')
while True:
try:
if self.web3.eth.getTransactionReceipt(tx_hash).status:
break
else:
print('\nTransaction failed. Aborting...')
return
except AttributeError:
pass
for i in range(0, 10):
print('.', end='', flush=True)
time.sleep(1)
print('\nFirst transaction successful!\n'
'Now processing second transaction.')
tx = self.mkt_contract.functions.subscribe(
Web3.toHex(dataset),
).buildTransaction(
{'nonce': self.web3.eth.getTransactionCount(address)})
if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(address, tx)
try:
tx_hash = '0x{}'.format(bin_hex(
self.web3.eth.sendRawTransaction(signed_tx)))
print('\nThis is the TxHash for this transaction: '
'{}'.format(tx_hash))
except Exception as e:
print('Unable to subscribe to data source: {}'.format(e))
return
self.check_transaction(tx_hash)
print('Waiting for the second transaction to succeed...')
while True:
try:
if self.web3.eth.getTransactionReceipt(tx_hash).status:
break
else:
print('\nTransaction failed. Aborting...')
return
except AttributeError:
pass
for i in range(0, 10):
print('.', end='', flush=True)
time.sleep(1)
print('\nSecond transaction successful!\n'
'You have successfully subscribed to dataset {} with'
'address {}.\n'
'You can now ingest this dataset anytime during the '
'next month by running the following command:\n'
'catalyst marketplace ingest --dataset={}'.format(
dataset, address, dataset))
def process_temp_bundle(self, ds_name, path):
"""
Merge the temp bundle into the main bundle for the specified
data source.
Parameters
----------
ds_name
path
Returns
-------
"""
tmp_bundle = extract_bundle(path)
bundle_folder = get_data_source_folder(ds_name)
if os.listdir(bundle_folder):
zsource = bcolz.ctable(rootdir=tmp_bundle, mode='r')
ztarget = bcolz.ctable(rootdir=bundle_folder, mode='r')
merge_bundles(zsource, ztarget)
else:
os.rename(tmp_bundle, bundle_folder)
pass
def ingest(self, ds_name, start=None, end=None, force_download=False):
# ds_name = ds_name.lower()
# TODO: catch error conditions
provider_info = self.mkt_contract.functions.getDataProviderInfo(
Web3.toHex(ds_name)
).call()
if not provider_info[4]:
print('The requested "{}" dataset is not registered in '
'the Data Marketplace.'.format(ds_name))
return
address, address_i = self.choose_pubaddr()
fns = self.mkt_contract.functions
check_sub = fns.checkAddressSubscription(
address, Web3.toHex(ds_name)
).call()
if check_sub[0] != address or self.to_text(check_sub[1]) != ds_name:
print('You are not subscribed to dataset "{}" with address {}. '
'Plese subscribe first.'.format(ds_name, address))
return
if not check_sub[5]:
print('Your subscription to dataset "{}" expired on {} UTC.'
'Please renew your subscription by running:\n'
'catalyst marketplace subscribe --dataset={}'.format(
ds_name,
pd.to_datetime(check_sub[4], unit='s', utc=True),
ds_name)
)
if 'key' in self.addresses[address_i]:
key = self.addresses[address_i]['key']
secret = self.addresses[address_i]['secret']
else:
key, secret = get_key_secret(address)
headers = get_signed_headers(ds_name, key, secret)
log.debug('Starting download of dataset for ingestion...')
r = requests.post(
'{}/marketplace/ingest'.format(AUTH_SERVER),
headers=headers,
stream=True,
)
if r.status_code == 200:
target_path = get_temp_bundles_folder()
try:
decoder = MultipartDecoder.from_response(r)
for part in decoder.parts:
h = part.headers[b'Content-Disposition'].decode('utf-8')
# Extracting the filename from the header
name = re.search(r'filename="(.*)"', h).group(1)
filename = os.path.join(target_path, name)
with open(filename, 'wb') as f:
# for chunk in part.content.iter_content(
# chunk_size=1024):
# if chunk: # filter out keep-alive new chunks
# f.write(chunk)
f.write(part.content)
self.process_temp_bundle(ds_name, filename)
except NonMultipartContentTypeException:
response = r.json()
raise MarketplaceHTTPRequest(
request='ingest dataset',
error=response,
)
else:
raise MarketplaceHTTPRequest(
request='ingest dataset',
error=r.status_code,
)
log.info('{} ingested successfully'.format(ds_name))
def get_dataset(self, ds_name, start=None, end=None):
ds_name = ds_name.lower()
# TODO: filter ctable by start and end date
bundle_folder = get_data_source_folder(ds_name)
z = bcolz.ctable(rootdir=bundle_folder, mode='r')
df = z.todataframe() # type: pd.DataFrame
df.set_index(['date', 'symbol'], drop=True, inplace=True)
# TODO: implement the filter more carefully
# if start and end is None:
# df = df.xs(start, level=0)
return df
def clean(self, data_source_name, data_frequency=None):
data_source_name = data_source_name.lower()
if data_frequency is None:
folder = get_data_source_folder(data_source_name)
else:
folder = get_bundle_folder(data_source_name, data_frequency)
shutil.rmtree(folder)
pass
def create_metadata(self, key, secret, ds_name, data_frequency, desc,
has_history=True, has_live=True):
"""
Returns
-------
"""
headers = get_signed_headers(ds_name, key, secret)
r = requests.post(
'{}/marketplace/register'.format(AUTH_SERVER),
json=dict(
ds_name=ds_name,
desc=desc,
data_frequency=data_frequency,
has_history=has_history,
has_live=has_live,
),
headers=headers,
)
if r.status_code != 200:
raise MarketplaceHTTPRequest(
request='register', error=r.status_code
)
if 'error' in r.json():
raise MarketplaceHTTPRequest(
request='upload file', error=r.json()['error']
)
def register(self):
while True:
desc = input('Enter the name of the dataset to register: ')
dataset = desc.lower()
provider_info = self.mkt_contract.functions.getDataProviderInfo(
Web3.toHex(dataset)
).call()
if provider_info[4]:
print('There is already a dataset registered under '
'the name "{}". Please choose a different '
'name.'.format(dataset))
else:
break
price = int(
input(
'Enter the price for a monthly subscription to '
'this dataset in ENG: '
)
)
while True:
freq = input('Enter the data frequency [daily, hourly, minute]: ')
if freq.lower() not in ('daily', 'hourly', 'minute'):
print('Not a valid frequency.')
else:
break
while True:
reg_pub = input(
'Does it include historical data? [default: Y]: '
) or 'y'
if reg_pub.lower() not in ('y', 'n'):
print('Please answer Y or N.')
else:
if reg_pub.lower() == 'y':
has_history = True
else:
has_history = False
break
while True:
reg_pub = input(
'Doest it include live data? [default: Y]: '
) or 'y'
if reg_pub.lower() not in ('y', 'n'):
print('Please answer Y or N.')
else:
if reg_pub.lower() == 'y':
has_live = True
else:
has_live = False
break
address, address_i = self.choose_pubaddr()
if 'key' in self.addresses[address_i]:
key = self.addresses[address_i]['key']
secret = self.addresses[address_i]['secret']
else:
key, secret = get_key_secret(address)
grains = to_grains(price)
tx = self.mkt_contract.functions.register(
Web3.toHex(dataset),
grains,
address,
).buildTransaction(
{'nonce': self.web3.eth.getTransactionCount(address)}
)
if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(address, tx)
try:
tx_hash = '0x{}'.format(
bin_hex(self.web3.eth.sendRawTransaction(signed_tx))
)
print(
'\nThis is the TxHash for this transaction: {}'.format(tx_hash)
)
except Exception as e:
print('Unable to subscribe to data source: {}'.format(e))
return
self.check_transaction(tx_hash)
print('Waiting for the transaction to succeed...')
while True:
try:
if self.web3.eth.getTransactionReceipt(tx_hash).status:
break
else:
print('\nTransaction failed. Aborting...')
return
except AttributeError:
pass
for i in range(0, 10):
print('.', end='', flush=True)
time.sleep(1)
print('\nWarming up the {} dataset'.format(dataset))
self.create_metadata(
key=key,
secret=secret,
ds_name=dataset,
data_frequency=freq,
desc=desc,
has_history=has_history,
has_live=has_live,
)
print('\n{} registered successfully'.format(dataset))
def publish(self, dataset, datadir, watch):
dataset = dataset.lower()
provider_info = self.mkt_contract.functions.getDataProviderInfo(
Web3.toHex(dataset)
).call()
if not provider_info[4]:
raise MarketplaceDatasetNotFound(dataset=dataset)
match = next(
(l for l in self.addresses if l['pubAddr'] == provider_info[0]),
None
)
if not match:
raise MarketplaceNoAddressMatch(
dataset=dataset,
address=provider_info[0])
print('Using address: {} to publish this dataset.'.format(
provider_info[0]))
if 'key' in match:
key = match['key']
secret = match['secret']
else:
key, secret = get_key_secret(provider_info[0])
headers = get_signed_headers(dataset, key, secret)
filenames = glob.glob(os.path.join(datadir, '*.csv'))
if not filenames:
raise MarketplaceNoCSVFiles(datadir=datadir)
files = []
for file in filenames:
files.append(('file', open(file, 'rb')))
r = requests.post('{}/marketplace/publish'.format(AUTH_SERVER),
files=files,
headers=headers)
if r.status_code != 200:
raise MarketplaceHTTPRequest(request='upload file',
error=r.status_code)
if 'error' in r.json():
raise MarketplaceHTTPRequest(request='upload file',
error=r.json()['error'])
print('Dataset {} uploaded successfully.'.format(dataset))
@@ -0,0 +1,88 @@
import sys
import traceback
from catalyst.errors import ZiplineError
def silent_except_hook(exctype, excvalue, exctraceback):
if exctype in [MarketplacePubAddressEmpty, MarketplaceDatasetNotFound,
MarketplaceNoAddressMatch, MarketplaceHTTPRequest,
MarketplaceNoCSVFiles, MarketplaceContractDataNoMatch,
MarketplaceSubscriptionExpired, MarketplaceJSONError,
MarketplaceWalletNotSupported, MarketplaceEmptySignature]:
fn = traceback.extract_tb(exctraceback)[-1][0]
ln = traceback.extract_tb(exctraceback)[-1][1]
print("Error traceback: {1} (line {2})\n"
"{0.__name__}: {3}".format(exctype, fn, ln, excvalue))
else:
sys.__excepthook__(exctype, excvalue, exctraceback)
sys.excepthook = silent_except_hook
class MarketplacePubAddressEmpty(ZiplineError):
msg = (
'Please enter your public address to use in the Data Marketplace '
'in the following file: {filename}'
).strip()
class MarketplaceDatasetNotFound(ZiplineError):
msg = (
'The dataset "{dataset}" is not registered in the Data Marketplace.'
).strip()
class MarketplaceNoAddressMatch(ZiplineError):
msg = (
'The address registered with the dataset {dataset}: {address} '
'does not match any of your addresses.'
).strip()
class MarketplaceHTTPRequest(ZiplineError):
msg = (
'Request to remote server to {request} failed: {error}'
).strip()
class MarketplaceNoCSVFiles(ZiplineError):
msg = (
'No CSV files found on {datadir} to upload.'
)
class MarketplaceContractDataNoMatch(ZiplineError):
msg = (
'The information found on the contract does not match the '
'requested data:\n{params}.'
)
class MarketplaceSubscriptionExpired(ZiplineError):
msg = (
'Your subscription to dataset "{dataset}" expired on {date} '
'and is no longer active. You have to subscribe again running the '
'following command:\n'
'catalyst marketplace subscribe --dataset={dataset}'
)
class MarketplaceWalletNotSupported(ZiplineError):
msg = (
'Wallet {wallet} is not supported.'
)
class MarketplaceEmptySignature(ZiplineError):
msg = (
'Signature cannot be empty.'
)
class MarketplaceJSONError(ZiplineError):
msg = (
'The configuration file {file} is malformed. Please correct '
'the following error:\n{error}'
)
+131
View File
@@ -0,0 +1,131 @@
import hashlib
import hmac
import requests
import time
from catalyst.marketplace.marketplace_errors import (
MarketplaceHTTPRequest, MarketplaceWalletNotSupported,
MarketplaceEmptySignature)
from catalyst.marketplace.utils.path_utils import (
get_user_pubaddr, save_user_pubaddr)
from catalyst.constants import AUTH_SERVER
def get_key_secret(pubAddr, wallet='mew'):
"""
Obtain a new key/secret pair from authentication server
Parameters
----------
pubAddr: str
dataset: str
Returns
-------
key: str
secret: str
"""
session = requests.Session()
response = session.get('{}/marketplace/getkeysecret'.format(AUTH_SERVER),
headers={
'Authorization': 'Digest username="{0}"'.format(
pubAddr)})
if response.status_code != 401:
raise MarketplaceHTTPRequest(request=str('obtain key/secret'),
error='Unexpected response code: '
'{}'.format(response.status_code))
header = response.headers.get('WWW-Authenticate')
auth_type, auth_info = header.split(None, 1)
d = requests.utils.parse_dict_header(auth_info)
nonce = '0x{}'.format(d['nonce'])
if wallet == 'mew':
print('\nObtaining a key/secret pair to streamline all future '
'requests with the authentication server.\n'
'Visit https://www.myetherwallet.com/signmsg.html and sign the'
'following message:\n{}'.format(nonce))
signature = input('Copy and Paste the "sig" field from '
'the signature here (without the double quotes, '
'only the HEX value:\n')
else:
raise MarketplaceWalletNotSupported(wallet=wallet)
if signature is None:
raise MarketplaceEmptySignature()
signature = signature[2:]
r = int(signature[0:64], base=16)
s = int(signature[64:128], base=16)
v = int(signature[128:130], base=16)
vrs = [v, r, s]
response = session.get('{}/marketplace/getkeysecret'.format(AUTH_SERVER),
headers={
'Authorization': 'Digest username="{0}",realm="{1}",'
'nonce="{2}",uri="/marketplace/getkeysecret",response="{3}",'
'opaque="{4}"'.format(pubAddr,
d['realm'],
d['nonce'],
','.join(str(e) for e in vrs+[wallet]),
d['opaque'])})
if response.status_code == 200:
if 'error' in response.json():
raise MarketplaceHTTPRequest(request=str('obtain key/secret'),
error=str(response.json()['error']))
else:
addresses = get_user_pubaddr()
match = next((l for l in addresses if
l['pubAddr'] == pubAddr), None)
match['key'] = response.json()['key']
match['secret'] = response.json()['secret']
addresses[addresses.index(match)] = match
save_user_pubaddr(addresses)
print('Key/secret pair retrieved successfully from server.')
return match['key'], match['secret']
else:
raise MarketplaceHTTPRequest(request=str('obtain key/secret'),
error=response.status_code)
def get_signed_headers(ds_name, key, secret):
"""
Return a new request header including the key / secret signature
Parameters
----------
ds_name
key
secret
Returns
-------
"""
nonce = str(int(time.time()))
signature = hmac.new(
secret.encode('utf-8'),
'{}{}'.format(ds_name, nonce).encode('utf-8'),
hashlib.sha512
).hexdigest()
headers = {
'Sign': signature,
'Key': key,
'Nonce': nonce,
'Dataset': ds_name,
}
return headers
@@ -0,0 +1,36 @@
import os
import shutil
import bcolz
def merge_bundles(zsource, ztarget):
"""
Merge
Parameters
----------
zsource
ztarget
Returns
-------
"""
# TODO: find a way to do this iteratively instead of in-memory
df_source = zsource.todataframe()
df_source.set_index('date', drop=False, inplace=True)
df_target = ztarget.todataframe()
df_target.set_index('date', drop=False, inplace=True)
df = df_target.merge(
right=df_source,
how='right',
) # type: pd.DataFrame
dirname = os.path.basename(ztarget.rootdir)
bak_dir = ztarget.rootdir.replace(dirname, '.{}'.format(dirname))
os.rename(ztarget.rootdir, bak_dir)
z = bcolz.ctable.fromdataframe(df=df, rootdir=ztarget.rootdir)
shutil.rmtree(bak_dir)
return z
+82
View File
@@ -0,0 +1,82 @@
import binascii
# def bytes32(string):
# """
# Convert string to bytes32 data type for smart contract
# Parameters
# ----------
# string: str
# Returns
# -------
# list
# """
# return binascii.hexlify(string.encode('utf-8'))
# def b32_str(bytes32):
# """
# Convert bytes32 to string
# Parameters
# ----------
# input: bytes object
# Returns
# -------
# str
# """
# return binascii.unhexlify(
# bytes32.decode('utf-8').rstrip('\0')).decode('ascii')
def bin_hex(binary):
"""
Convert bytes32 to string
Parameters
----------
input: bytes object
Returns
-------
str
"""
return binascii.hexlify(binary).decode('utf-8')
def from_grains(amount):
"""
Convert from grains to cryptocurrency
Parameters
----------
input: amount
Returns
-------
int
"""
return amount // 10 ** 8
def to_grains(amount):
"""
Convert from cryptocurrency to grains
Parameters
----------
input: amount
Returns
-------
int
"""
return amount * 10 ** 8
+166
View File
@@ -0,0 +1,166 @@
import os
import json
import tarfile
from catalyst.utils.deprecate import deprecated
from catalyst.utils.paths import data_root, ensure_directory
from catalyst.marketplace.marketplace_errors import MarketplaceJSONError
def get_marketplace_folder(environ=None):
"""
The root path of the marketplace folder.
Parameters
----------
environ:
Returns
-------
str
"""
if not environ:
environ = os.environ
root = data_root(environ)
marketplace_folder = os.path.join(root, 'marketplace')
ensure_directory(marketplace_folder)
return marketplace_folder
def get_data_source_folder(data_source_name, environ=None):
"""
The root path of an data_source folder.
Parameters
----------
data_source_name: str
environ:
Returns
-------
str
"""
if not environ:
environ = os.environ
root = data_root(environ)
data_source_folder = os.path.join(root, 'marketplace', data_source_name)
ensure_directory(data_source_folder)
return data_source_folder
@deprecated
def get_bundle_folder(data_source_name, data_frequency, environ=None):
data_source_folder = get_data_source_folder(data_source_name, environ)
bundle_folder = os.path.join(data_source_folder, data_frequency)
ensure_directory(bundle_folder)
return bundle_folder
def get_temp_bundles_folder(environ=None):
"""
The temp folder for bundle downloads by algo name.
Parameters
----------
ds_name: str
environ:
Returns
-------
str
"""
root = data_root(environ)
folder = os.path.join(root, 'marketplace', 'temp_bundles')
ensure_directory(folder)
return folder
def extract_bundle(tar_filename):
"""
Extract a bcolz bundle.
Parameters
----------
ds_name
Returns
-------
str
"""
target_path = tar_filename.replace('.tar.gz', '')
with tarfile.open(tar_filename, 'r') as tar:
tar.extractall(target_path)
return target_path
def get_user_pubaddr(environ=None):
"""
The de-serialized contend of the user's addresses.json file.
Parameters
----------
environ:
Returns
-------
Object
"""
marketplace_folder = get_marketplace_folder(environ)
filename = os.path.join(marketplace_folder, 'addresses.json')
if os.path.isfile(filename):
with open(filename) as data_file:
try:
data = json.load(data_file)
except json.decoder.JSONDecodeError as e:
raise MarketplaceJSONError(file=filename, error=e)
try:
d = data[0]['pubAddr']
except Exception as e:
return [data, ]
return data
else:
data = []
data.append(dict(pubAddr='', desc=''))
with open(filename, 'w') as f:
json.dump(data, f, sort_keys=False, indent=2,
separators=(',', ':'))
return data
def save_user_pubaddr(data, environ=None):
"""
Saves the user's public addresses and their related metadata in
the corresponding addresses.json file.
Parameters
----------
data: dict
Returns
-------
True
"""
marketplace_folder = get_marketplace_folder(environ)
filename = os.path.join(marketplace_folder, 'addresses.json')
with open(filename, 'w') as f:
json.dump(data, f, sort_keys=False, indent=2,
separators=(',', ':'))
return True
+376
View File
@@ -0,0 +1,376 @@
# -*- coding: utf-8 -*-
# !/usr/bin/env python2
import sys
import os
import pandas as pd
import signal
# import talib
from logbook import Logger
from catalyst import run_algorithm
from catalyst.api import (
symbol,
record,
order,
order_target,
order_target_percent,
get_open_orders
)
from catalyst.finance import commission
# from base.telegrambot import TelegramBot
class GracefulKiller:
# Source: https://stackoverflow.com/a/31464349
def __init__(self, context):
self.kill_now = False
self.signal = 0
self.context = context
signal.signal(signal.SIGINT, self.exit_gracefully)
def exit_gracefully(self, signum, frame):
self.kill_now = True
self.signal = signum
if hasattr(self.context,
'telegram_bot') and self.context.telegram_bot is not None:
self.context.telegram_bot.updater.stop()
sys.exit(0)
def exit(self):
return self.kill_now
class SimulationParameters:
MODE = 'paper'
CAPITAL_BASE = 1000
"""
Capital base used on this simulation
"""
DATA_FREQUECY = 'minute'
EXCHANGE_NAME = 'bitfinex'
# EXCHANGE_NAME = 'binance'
"""
Exchange used on this simulation
"""
DATA_DIR = '/home/av/Dropbox/simulations/data'
ALGO_NAMESPACE = os.path.basename(__file__).split('.')[0]
ALGO_NAMESPACE_IMAGE = '{}/{}/{}.png'.format(DATA_DIR, 'images',
ALGO_NAMESPACE)
ALGO_NAMESPACE_RESULTS_TABLE = '{}/{}/{}.csv'.format(DATA_DIR, 'tables',
ALGO_NAMESPACE + '_results')
ALGO_NAMESPACE_TRANSACTIONS_TABLE = '{}/{}/{}.csv'.format(DATA_DIR,
'tables',
ALGO_NAMESPACE + '_transactions')
BASE_CURRENCY = 'usd'
# BASE_CURRENCY = 'usdt'
# SHORT PERIOD
START_DATE = '2017-09-07'
"""
Start date used on this simulation
"""
END_DATE = '2017-12-12'
"""
End date used on this simulation
"""
SKIP_FIRST_CANDLES = 0
# CANDLES_SAMPLE_RATE = 60
# CANDLES_SAMPLE_RATE = 30
CANDLES_SAMPLE_RATE = 1
"""
Candle interval used on this simulation (in minutes)
"""
# http://pandas.pydata.org/pandas-docs/stable/timeseries.html#offset-aliases
# 30 minute interval ohlcv data (the standard data required for candlestick or
# indicators/signals)
# 30T means 30 minutes re-sampling of one minute data.
# CANDLES_FREQUENCY = '60T'
# CANDLES_FREQUENCY = '30T'
CANDLES_FREQUENCY = '1T'
CANDLES_BUFFER_SIZE = 48
COIN_PAIR = 'btc_usd'
# COIN_PAIR = 'btc_usdt'
"""
Coin pair used on this simulation
"""
# TRANSACTIONS
COMMISSION_FEE = 0.0030
BUY_MIN_AMOUNT = 5 # i.e: USD
SELL_MIN_AMOUNT = 0.001 # i.e: USD
BUY_SELL_PERCENTAGE = 1 # 0.50
BUY_PERCENTAGE = BUY_SELL_PERCENTAGE
SELL_PERCENTAGE = BUY_SELL_PERCENTAGE
BASE_PRICE = 'close'
"""
Base price used (close / Heiken Ashi)
"""
log = None
parameters = None
def print_facts(context):
context.log.info("""
Index: {}
Date: {}
Candle:
O: {}
H: {}
L: {}
C: {}
V: {}
Metrics:
...
Portfolio:
Base price: {}
Base coin (coin2/usd): {}
Amount (coin1/btc): {}
""".format(
# Facts
context.i,
context.curr_minute,
context.candles_open[-1],
context.candles_high[-1],
context.candles_low[-1],
context.candles_close[-1],
context.candles_volume[-1],
# Metrics
# ...
# Portfolio
context.curr_base_price,
context.portfolio.cash,
context.portfolio.positions[context.coin_pair].amount,
))
def print_facts_telegram(context):
price = context.curr_base_price
amount = context.portfolio.positions[context.coin_pair].amount
pnl = context.portfolio.pnl
capital_used = context.portfolio.capital_used
portfolio_value = context.portfolio.portfolio_value
portfolio_returns = context.portfolio.returns
starting_cash = context.portfolio.starting_cash
cash = context.portfolio.cash
msg = """
Status...
Price: {}
Starting cash: {}
Cash: {}
Capital used: {}
Amount: {}
Portfolio value: {}
Returns: {}
PnL: {}
""".format(
price,
starting_cash,
cash,
capital_used,
amount,
portfolio_value,
portfolio_returns,
pnl,
)
if hasattr(context, 'telegram_bot') and context.telegram_bot is not None:
context.telegram_bot.msg(msg)
def default_initialize(context):
# FIXME: set_benchmark
# set_benchmark(symbol(context.parameters.COIN_PAIR))
context.coin_pair = symbol(context.parameters.COIN_PAIR)
context.base_price = None
context.current_day = None
context.counter = -1
context.i = 0
context.candles_sample_rate = context.parameters.CANDLES_SAMPLE_RATE
context.candles_frequency = context.parameters.CANDLES_FREQUENCY
context.candles_buffer_size = context.parameters.CANDLES_BUFFER_SIZE
context.set_commission(
commission.PerShare(cost=context.parameters.COMMISSION_FEE))
def default_handle_data(context, data):
context.curr_minute = data.current_dt
context.counter += 1
if context.candles_sample_rate == 1:
context.i += 1
elif context.counter % context.candles_sample_rate != 0:
context.i += 1
return
if context.i < context.parameters.SKIP_FIRST_CANDLES:
return
context.candles_open = data.history(
context.coin_pair,
'open',
bar_count=context.candles_buffer_size,
frequency=context.candles_frequency)
context.candles_high = data.history(
context.coin_pair,
'high',
bar_count=context.candles_buffer_size,
frequency=context.candles_frequency)
context.candles_low = data.history(
context.coin_pair,
'low',
bar_count=context.candles_buffer_size,
frequency=context.candles_frequency)
context.candles_close = data.history(
context.coin_pair,
'price',
bar_count=context.candles_buffer_size,
frequency=context.candles_frequency)
context.candles_volume = data.history(
context.coin_pair,
'volume',
bar_count=context.candles_buffer_size,
frequency=context.candles_frequency)
# FIXME: Here is the error!
# The candles_close frame shows more or less always a value of 94, while
# bitcoin price is very different from that
print(context.candles_close)
context.base_prices = context.candles_close
cash = context.portfolio.cash
amount = context.portfolio.positions[context.coin_pair].amount
price = data.current(context.coin_pair, 'price')
order_id = None
context.last_base_price = context.base_prices[-2]
context.curr_base_price = context.base_prices[-1]
# TA calculations
# ...
# Sanity checks
# assert cash >= 0
if cash < 0:
import ipdb;
ipdb.set_trace() # BREAKPOINT
print_facts(context)
print_facts_telegram(context)
# Order management
net_shares = 0
if context.counter == 2:
brute_shares = (cash / price) * context.parameters.BUY_PERCENTAGE
share_commission_fee = brute_shares * context.parameters.COMMISSION_FEE
net_shares = brute_shares - share_commission_fee
buy_order_id = order(context.coin_pair, net_shares)
if context.counter == 3:
brute_shares = amount * context.parameters.SELL_PERCENTAGE
share_commission_fee = brute_shares * context.parameters.COMMISSION_FEE
net_shares = -(brute_shares - share_commission_fee)
sell_order_id = order(context.coin_pair, net_shares)
# Record
record(
price=price,
foo='bar',
# volume=current['volume'],
# price_change=price_change,
# Metrics
cash=cash,
# buy=context.buy,
# sell=context.sell
)
def default_analyze(context=None, perf=None):
pass
def initialize(context):
global log
context.parameters = parameters
context.log = Logger(context.parameters.ALGO_NAMESPACE)
log = context.log
default_initialize(context)
context.killer = GracefulKiller(context)
context.telegram_bot = None
# TELEGRAM_TOKEN='token'
# context.telegram_bot = TelegramBot()
# context.telegram_bot.initialize(TELEGRAM_TOKEN, context)
if __name__ == '__main__':
# Parameters:
parameters = SimulationParameters()
start_date = pd.to_datetime(parameters.START_DATE, utc=True)
end_date = pd.to_datetime(parameters.END_DATE, utc=True)
if parameters.MODE == 'backtest':
results = run_algorithm(
capital_base=parameters.CAPITAL_BASE,
data_frequency=parameters.DATA_FREQUECY,
initialize=initialize,
handle_data=default_handle_data,
analyze=default_analyze,
exchange_name=parameters.EXCHANGE_NAME,
algo_namespace=parameters.ALGO_NAMESPACE,
base_currency=parameters.BASE_CURRENCY,
start=start_date,
end=end_date,
live=False,
live_graph=False
)
returns_daily = results
results.to_csv('{}'.format(parameters.ALGO_NAMESPACE_RESULTS_TABLE))
# returns_daily = returns_minutely.add(1).groupby(pd.TimeGrouper('24H')).prod().add(-1)
# FIXME: pyfolio integration
# pf_data = pyfolio.utils.extract_rets_pos_txn_from_zipline(results)
# pf_data = pyfolio.utils.extract_rets_pos_txn_from_zipline(results[:'2017-01-01'])
# pyfolio.create_full_tear_sheet(*pf_data)
elif parameters.MODE == 'paper':
results = run_algorithm(
capital_base=parameters.CAPITAL_BASE,
data_frequency=parameters.DATA_FREQUECY,
initialize=initialize,
handle_data=default_handle_data,
analyze=default_analyze,
exchange_name=parameters.EXCHANGE_NAME,
algo_namespace=parameters.ALGO_NAMESPACE,
base_currency=parameters.BASE_CURRENCY,
live=True,
simulate_orders=True,
live_graph=False
)
elif parameters.MODE == 'live':
results = run_algorithm(
initialize=initialize,
handle_data=default_handle_data,
analyze=default_analyze,
exchange_name=parameters.EXCHANGE_NAME,
algo_namespace=parameters.ALGO_NAMESPACE,
base_currency=parameters.BASE_CURRENCY,
live=True,
live_graph=True
)
+6 -8
View File
@@ -55,6 +55,7 @@ class _RunAlgoError(click.ClickException, ValueError):
----------
pyfunc_msg : str
The message that will be shown when called as a python function.
cmdline_msg : str
The message that will be shown on the command line.
"""
@@ -416,7 +417,8 @@ def run_algorithm(initialize,
auth_aliases=None,
stats_output=None,
output=os.devnull):
"""Run a trading algorithm.
"""
Run a trading algorithm.
Parameters
----------
@@ -458,7 +460,7 @@ def run_algorithm(initialize,
This argument is mutually exclusive with ``data``.
default_extension : bool, optional
Should the default catalyst extension be loaded. This is found at
``$ZIPLINE_ROOT/extension.py``
``$CATALYST_ROOT/extension.py``
extensions : iterable[str], optional
The names of any other extensions to load. Each element may either be
a dotted module path like ``a.b.c`` or a path to a python file ending
@@ -469,12 +471,8 @@ def run_algorithm(initialize,
environ : mapping[str -> str], optional
The os environment to use. Many extensions use this to get parameters.
This defaults to ``os.environ``.
live: execute live trading
exchange_conn: The exchange connection parameters
Supported Exchanges
-------------------
bitfinex
live : bool, optional
Execute algorithm in live trading mode.
Returns
-------
+173 -179
View File
@@ -4,7 +4,7 @@ API Reference
Running a Backtest
~~~~~~~~~~~~~~~~~~
.. autofunction:: zipline.run_algorithm(...)
.. autofunction:: catalyst.run_algorithm(...)
Algorithm API
~~~~~~~~~~~~~
@@ -18,341 +18,335 @@ currently-executing :class:`~zipline.algorithm.TradingAlgorithm` instance.
Data Object
```````````
.. autoclass:: zipline.protocol.BarData
.. autoclass:: catalyst.protocol.BarData
:members:
Scheduling Functions
````````````````````
.. autofunction:: zipline.api.schedule_function
.. autofunction:: catalyst.api.schedule_function
.. autoclass:: zipline.api.date_rules
.. autoclass:: catalyst.api.date_rules
:members:
:undoc-members:
.. autoclass:: zipline.api.time_rules
.. autoclass:: catalyst.api.time_rules
:members:
Orders
``````
.. autofunction:: zipline.api.order
.. autofunction:: catalyst.api.order
.. autofunction:: zipline.api.order_value
.. autofunction:: catalyst.api.order_value
.. autofunction:: zipline.api.order_percent
.. autofunction:: catalyst.api.order_percent
.. autofunction:: zipline.api.order_target
.. autofunction:: catalyst.api.order_target
.. autofunction:: zipline.api.order_target_value
.. autofunction:: catalyst.api.order_target_value
.. autofunction:: zipline.api.order_target_percent
.. autofunction:: catalyst.api.order_target_percent
.. autoclass:: zipline.finance.execution.ExecutionStyle
.. autoclass:: catalyst.finance.execution.ExecutionStyle
:members:
.. autoclass:: zipline.finance.execution.MarketOrder
.. autoclass:: catalyst.finance.execution.MarketOrder
.. autoclass:: zipline.finance.execution.LimitOrder
.. autoclass:: catalyst.finance.execution.LimitOrder
.. autoclass:: zipline.finance.execution.StopOrder
.. autoclass:: catalyst.finance.execution.StopOrder
.. autoclass:: zipline.finance.execution.StopLimitOrder
.. autoclass:: catalyst.finance.execution.StopLimitOrder
.. autofunction:: zipline.api.get_order
.. autofunction:: catalyst.api.get_order
.. autofunction:: zipline.api.get_open_orders
.. autofunction:: catalyst.api.get_open_orders
.. autofunction:: zipline.api.cancel_order
.. autofunction:: catalyst.api.cancel_order
Order Cancellation Policies
'''''''''''''''''''''''''''
.. autofunction:: zipline.api.set_cancel_policy
.. autofunction:: catalyst.api.set_cancel_policy
.. autoclass:: zipline.finance.cancel_policy.CancelPolicy
.. autoclass:: catalyst.finance.cancel_policy.CancelPolicy
:members:
.. autofunction:: zipline.api.EODCancel
.. autofunction:: catalyst.api.EODCancel
.. autofunction:: zipline.api.NeverCancel
.. autofunction:: catalyst.api.NeverCancel
Assets
``````
.. autofunction:: zipline.api.symbol
.. autofunction:: catalyst.api.symbol
.. autofunction:: zipline.api.symbols
.. autofunction:: catalyst.api.symbols
.. autofunction:: zipline.api.future_symbol
.. autofunction:: catalyst.api.set_symbol_lookup_date
.. autofunction:: zipline.api.set_symbol_lookup_date
.. autofunction:: zipline.api.sid
.. autofunction:: catalyst.api.sid
Trading Controls
````````````````
Zipline provides trading controls to help ensure that the algorithm is
zipline provides trading controls to help ensure that the algorithm is
performing as expected. The functions help protect the algorithm from certian
bugs that could cause undesirable behavior when trading with real money.
.. autofunction:: zipline.api.set_do_not_order_list
.. autofunction:: catalyst.api.set_do_not_order_list
.. autofunction:: zipline.api.set_long_only
.. autofunction:: catalyst.api.set_long_only
.. autofunction:: zipline.api.set_max_leverage
.. autofunction:: catalyst.api.set_max_leverage
.. autofunction:: zipline.api.set_max_order_count
.. autofunction:: catalyst.api.set_max_order_count
.. autofunction:: zipline.api.set_max_order_size
.. autofunction:: catalyst.api.set_max_order_size
.. autofunction:: zipline.api.set_max_position_size
.. autofunction:: catalyst.api.set_max_position_size
Simulation Parameters
`````````````````````
.. autofunction:: zipline.api.set_benchmark
.. autofunction:: catalyst.api.set_benchmark
Commission Models
'''''''''''''''''
.. autofunction:: zipline.api.set_commission
.. autofunction:: catalyst.api.set_commission
.. autoclass:: zipline.finance.commission.CommissionModel
.. autoclass:: catalyst.finance.commission.CommissionModel
:members:
.. autoclass:: zipline.finance.commission.PerShare
.. autoclass:: catalyst.finance.commission.PerShare
.. autoclass:: zipline.finance.commission.PerTrade
.. autoclass:: catalyst.finance.commission.PerTrade
.. autoclass:: zipline.finance.commission.PerDollar
.. autoclass:: catalyst.finance.commission.PerDollar
Slippage Models
'''''''''''''''
.. autofunction:: zipline.api.set_slippage
.. autofunction:: catalyst.api.set_slippage
.. autoclass:: zipline.finance.slippage.SlippageModel
.. autoclass:: catalyst.finance.slippage.SlippageModel
:members:
.. autoclass:: zipline.finance.slippage.FixedSlippage
.. autoclass:: catalyst.finance.slippage.FixedSlippage
.. autoclass:: zipline.finance.slippage.VolumeShareSlippage
.. autoclass:: catalyst.finance.slippage.VolumeShareSlippage
Pipeline
````````
For more information, see :ref:`pipeline-api`
Not supported yet.
.. autofunction:: zipline.api.attach_pipeline
.. For more information, see :ref:`pipeline-api`
.. autofunction:: zipline.api.pipeline_output
.. .. autofunction:: catalyst.api.attach_pipeline
.. .. autofunction:: catalyst.api.pipeline_output
Miscellaneous
`````````````
.. autofunction:: zipline.api.record
.. autofunction:: catalyst.api.record
.. autofunction:: zipline.api.get_environment
.. autofunction:: catalyst.api.get_environment
.. autofunction:: zipline.api.fetch_csv
.. autofunction:: catalyst.api.fetch_csv
.. _pipeline-api:
Pipeline API
~~~~~~~~~~~~
.. Pipeline API
.. ~~~~~~~~~~~~
.. autoclass:: zipline.pipeline.Pipeline
:members:
:member-order: groupwise
.. .. autoclass:: zipline.pipeline.Pipeline
.. :members:
.. :member-order: groupwise
.. autoclass:: zipline.pipeline.CustomFactor
:members:
:member-order: groupwise
.. .. autoclass:: zipline.pipeline.CustomFactor
.. :members:
.. :member-order: groupwise
.. autoclass:: zipline.pipeline.filters.Filter
:members: __and__, __or__
:exclude-members: dtype
.. .. autoclass:: zipline.pipeline.filters.Filter
.. :members: __and__, __or__
.. :exclude-members: dtype
.. autoclass:: zipline.pipeline.factors.Factor
:members: bottom, deciles, demean, linear_regression, pearsonr,
percentile_between, quantiles, quartiles, quintiles, rank,
spearmanr, top, winsorize, zscore, isnan, notnan, isfinite, eq,
__add__, __sub__, __mul__, __div__, __mod__, __pow__, __lt__,
__le__, __ne__, __ge__, __gt__
:exclude-members: dtype
:member-order: bysource
.. .. autoclass:: zipline.pipeline.factors.Factor
.. :members: bottom, deciles, demean, linear_regression, pearsonr,
.. percentile_between, quantiles, quartiles, quintiles, rank,
.. spearmanr, top, winsorize, zscore, isnan, notnan, isfinite, eq,
.. \__add__, \__sub__, \__mul__, \__div__, \__mod__, \__pow__,
.. \__lt__, \__le__, \__ne__, \__ge__, \__gt__
.. :exclude-members: dtype
.. :member-order: bysource
.. autoclass:: zipline.pipeline.term.Term
:members:
:exclude-members: compute_extra_rows, dependencies, inputs, mask, windowed
.. .. autoclass:: zipline.pipeline.term.Term
.. :members:
.. :exclude-members: compute_extra_rows, dependencies, inputs, mask, windowed
.. autoclass:: zipline.pipeline.data.USEquityPricing
:members: open, high, low, close, volume
:undoc-members:
.. .. autoclass:: zipline.pipeline.data.USEquityPricing
.. :members: open, high, low, close, volume
.. :undoc-members:
Built-in Factors
````````````````
.. Built-in Factors
.. ````````````````
.. autoclass:: zipline.pipeline.factors.AverageDollarVolume
:members:
.. .. autoclass:: zipline.pipeline.factors.AverageDollarVolume
.. :members:
.. autoclass:: zipline.pipeline.factors.BollingerBands
:members:
.. .. autoclass:: zipline.pipeline.factors.BollingerBands
.. :members:
.. autoclass:: zipline.pipeline.factors.BusinessDaysSincePreviousEvent
:members:
.. .. autoclass:: zipline.pipeline.factors.BusinessDaysSincePreviousEvent
.. :members:
.. autoclass:: zipline.pipeline.factors.BusinessDaysUntilNextEvent
:members:
.. .. autoclass:: zipline.pipeline.factors.BusinessDaysUntilNextEvent
.. :members:
.. autoclass:: zipline.pipeline.factors.ExponentialWeightedMovingAverage
:members:
.. .. autoclass:: zipline.pipeline.factors.ExponentialWeightedMovingAverage
.. :members:
.. autoclass:: zipline.pipeline.factors.ExponentialWeightedMovingStdDev
:members:
.. .. autoclass:: zipline.pipeline.factors.ExponentialWeightedMovingStdDev
.. :members:
.. autoclass:: zipline.pipeline.factors.Latest
:members:
.. .. autoclass:: zipline.pipeline.factors.Latest
.. :members:
.. autoclass:: zipline.pipeline.factors.MaxDrawdown
:members:
.. .. autoclass:: zipline.pipeline.factors.MaxDrawdown
.. :members:
.. autoclass:: zipline.pipeline.factors.Returns
:members:
.. .. autoclass:: zipline.pipeline.factors.Returns
.. :members:
.. autoclass:: zipline.pipeline.factors.RollingLinearRegressionOfReturns
:members:
.. .. autoclass:: zipline.pipeline.factors.RollingLinearRegressionOfReturns
.. :members:
.. autoclass:: zipline.pipeline.factors.RollingPearsonOfReturns
:members:
.. .. autoclass:: zipline.pipeline.factors.RollingPearsonOfReturns
.. :members:
.. autoclass:: zipline.pipeline.factors.RollingSpearmanOfReturns
:members:
.. .. autoclass:: zipline.pipeline.factors.RollingSpearmanOfReturns
.. :members:
.. autoclass:: zipline.pipeline.factors.RSI
:members:
.. .. autoclass:: zipline.pipeline.factors.RSI
.. :members:
.. autoclass:: zipline.pipeline.factors.SimpleMovingAverage
:members:
.. .. autoclass:: zipline.pipeline.factors.SimpleMovingAverage
.. :members:
.. autoclass:: zipline.pipeline.factors.VWAP
:members:
.. .. autoclass:: zipline.pipeline.factors.VWAP
.. :members:
.. autoclass:: zipline.pipeline.factors.WeightedAverageValue
:members:
.. .. autoclass:: zipline.pipeline.factors.WeightedAverageValue
.. :members:
Pipeline Engine
```````````````
.. Pipeline Engine
.. ```````````````
.. autoclass:: zipline.pipeline.engine.PipelineEngine
:members: run_pipeline, run_chunked_pipeline
:member-order: bysource
.. .. autoclass:: zipline.pipeline.engine.PipelineEngine
.. :members: run_pipeline, run_chunked_pipeline
.. :member-order: bysource
.. autoclass:: zipline.pipeline.engine.SimplePipelineEngine
:members: __init__, run_pipeline, run_chunked_pipeline
:member-order: bysource
.. .. autoclass:: zipline.pipeline.engine.SimplePipelineEngine
.. :members: __init__, run_pipeline, run_chunked_pipeline
.. :member-order: bysource
.. autofunction:: zipline.pipeline.engine.default_populate_initial_workspace
.. .. autofunction:: zipline.pipeline.engine.default_populate_initial_workspace
Data Loaders
````````````
.. Data Loaders
.. ````````````
.. autoclass:: zipline.pipeline.loaders.equity_pricing_loader.USEquityPricingLoader
:members: __init__, from_files, load_adjusted_array
:member-order: bysource
.. .. autoclass:: zipline.pipeline.loaders.equity_pricing_loader.USEquityPricingLoader
.. :members: __init__, from_files, load_adjusted_array
.. :member-order: bysource
Asset Metadata
~~~~~~~~~~~~~~
.. autoclass:: zipline.assets.Asset
.. autoclass:: catalyst.assets.Asset
:members:
.. autoclass:: zipline.assets.Equity
:members:
.. autoclass:: zipline.assets.Future
:members:
.. autoclass:: zipline.assets.AssetConvertible
.. autoclass:: catalyst.assets.AssetConvertible
:members:
Trading Calendar API
~~~~~~~~~~~~~~~~~~~~
.. autofunction:: zipline.utils.calendars.get_calendar
.. autofunction:: catalyst.utils.calendars.get_calendar
.. autoclass:: zipline.utils.calendars.TradingCalendar
.. autoclass:: catalyst.utils.calendars.TradingCalendar
:members:
.. autofunction:: zipline.utils.calendars.register_calendar
.. autofunction:: catalyst.utils.calendars.register_calendar
.. autofunction:: zipline.utils.calendars.register_calendar_type
.. autofunction:: catalyst.utils.calendars.register_calendar_type
.. autofunction:: zipline.utils.calendars.deregister_calendar
.. autofunction:: catalyst.utils.calendars.deregister_calendar
.. autofunction:: zipline.utils.calendars.clear_calendars
.. autofunction:: catalyst.utils.calendars.clear_calendars
Data API
~~~~~~~~
Writers
```````
.. autoclass:: zipline.data.minute_bars.BcolzMinuteBarWriter
:members:
.. Writers
.. ```````
.. .. autoclass:: zipline.data.minute_bars.BcolzMinuteBarWriter
.. :members:
.. autoclass:: zipline.data.us_equity_pricing.BcolzDailyBarWriter
:members:
.. .. autoclass:: zipline.data.us_equity_pricing.BcolzDailyBarWriter
.. :members:
.. autoclass:: zipline.data.us_equity_pricing.SQLiteAdjustmentWriter
:members:
.. .. autoclass:: zipline.data.us_equity_pricing.SQLiteAdjustmentWriter
.. :members:
.. autoclass:: zipline.assets.AssetDBWriter
:members:
.. .. autoclass:: zipline.assets.AssetDBWriter
.. :members:
Readers
```````
.. autoclass:: zipline.data.minute_bars.BcolzMinuteBarReader
:members:
.. Readers
.. ```````
.. .. autoclass:: zipline.data.minute_bars.BcolzMinuteBarReader
.. :members:
.. autoclass:: zipline.data.us_equity_pricing.BcolzDailyBarReader
:members:
.. .. autoclass:: zipline.data.us_equity_pricing.BcolzDailyBarReader
.. :members:
.. autoclass:: zipline.data.us_equity_pricing.SQLiteAdjustmentReader
:members:
.. .. autoclass:: zipline.data.us_equity_pricing.SQLiteAdjustmentReader
.. :members:
.. autoclass:: zipline.assets.AssetFinder
:members:
.. .. autoclass:: zipline.assets.AssetFinder
.. :members:
.. autoclass:: zipline.data.data_portal.DataPortal
:members:
.. .. autoclass:: zipline.data.data_portal.DataPortal
.. :members:
Bundles
```````
.. autofunction:: zipline.data.bundles.register
.. Bundles
.. ```````
.. .. autofunction:: zipline.data.bundles.register
.. autofunction:: zipline.data.bundles.ingest(name, environ=os.environ, date=None, show_progress=True)
.. .. autofunction:: zipline.data.bundles.ingest(name, environ=os.environ, date=None, show_progress=True)
.. autofunction:: zipline.data.bundles.load(name, environ=os.environ, date=None)
.. .. autofunction:: zipline.data.bundles.load(name, environ=os.environ, date=None)
.. autofunction:: zipline.data.bundles.unregister
.. .. autofunction:: zipline.data.bundles.unregister
.. data:: zipline.data.bundles.bundles
.. .. data:: zipline.data.bundles.bundles
The bundles that have been registered as a mapping from bundle name to bundle
data. This mapping is immutable and should only be updated through
:func:`~zipline.data.bundles.register` or
:func:`~zipline.data.bundles.unregister`.
.. The bundles that have been registered as a mapping from bundle name to bundle
.. data. This mapping is immutable and should only be updated through
.. :func:`~zipline.data.bundles.register` or
.. :func:`~zipline.data.bundles.unregister`.
.. autofunction:: zipline.data.bundles.yahoo_equities
.. .. autofunction:: zipline.data.bundles.yahoo_equities
@@ -362,16 +356,16 @@ Utilities
Caching
```````
.. autoclass:: zipline.utils.cache.CachedObject
.. autoclass:: catalyst.utils.cache.CachedObject
.. autoclass:: zipline.utils.cache.ExpiringCache
.. autoclass:: catalyst.utils.cache.ExpiringCache
.. autoclass:: zipline.utils.cache.dataframe_cache
.. autoclass:: catalyst.utils.cache.dataframe_cache
.. autoclass:: zipline.utils.cache.working_file
.. autoclass:: catalyst.utils.cache.working_file
.. autoclass:: zipline.utils.cache.working_dir
.. autoclass:: catalyst.utils.cache.working_dir
Command Line
````````````
.. autofunction:: zipline.utils.cli.maybe_show_progress
.. autofunction:: catalyst.utils.cli.maybe_show_progress
+7 -5
View File
@@ -168,7 +168,7 @@ We'll start with the CLI, and introduce the ``run_algorithm()`` in the last
example of this tutorial. Some of the :doc:`example algorithms <example-algos>`
provide instructions on how to run them both from the CLI, and using the
:func:`~catalyst.run_algorithm` function. For the third method, refer to the
corresponding section on :doc:`Catalyst & Jupyter Notebook <jupyter>` after you
corresponding section on :ref:`Catalyst & Jupyter Notebook <jupyter>` after you
have assimilated the contents of this tutorial.
Command line interface
@@ -473,6 +473,7 @@ Which we execute by running:
</div>
|
There is a row for each trading day, starting on the first day of our
simulation Jan 1st, 2016. In the columns you can find various
information about the state of your algorithm. The column
@@ -518,7 +519,7 @@ alongside enigma-catalyst (with the exception of the ``Conda`` install, where it
was included by default inside the conda environment we created). If for any
reason you don't have it installed, you can add it by running:
.. code-block:: python
.. code-block:: bash
(catalyst)$ pip install matplotlib
@@ -806,6 +807,7 @@ the ``scikit-learn`` functions require ``numpy.ndarray``\ s rather than
``pandas.DataFrame``\ s, so you can simply pass the underlying
``ndarray`` of a ``DataFrame`` via ``.values``).
.. _jupyter:
Jupyter Notebook
~~~~~~~~~~~~~~~~
@@ -826,13 +828,13 @@ In order to use Jupyter Notebook, you first have to install it inside your
environment. It's available as ``pip`` package, so regardless of how you
installed Catalyst, go inside your catalyst environemnt and run:
.. code:: bash
.. code-block:: bash
(catalyst)$ pip install jupyter
Once you have Jupyter Notebook installed, every time you want to use it run:
.. code:: bash
.. code-block:: bash
(catalyst)$ jupyter notebook
@@ -846,7 +848,7 @@ Before running your algorithms inside the Jupyter Notebook, remember to ingest
the data from the command line interface (CLI). In the example below, you would
need to run first:
.. code:: bash
.. code-block:: bash
catalyst ingest-exchange -x bitfinex -i btc_usd
+5 -2
View File
@@ -27,8 +27,8 @@ extlinks = {
# -- Docstrings ---------------------------------------------------------------
#extensions += ['numpydoc']
#numpydoc_show_class_members = False
extensions += ['numpydoc']
numpydoc_show_class_members = False
# Add any paths that contain templates here, relative to this directory.
templates_path = ['.templates']
@@ -97,3 +97,6 @@ intersphinx_mapping = {
doctest_global_setup = "import catalyst"
todo_include_todos = True
suppress_warnings = ['image.nonlocal_uri']
+7 -17
View File
@@ -36,25 +36,15 @@ Finally, you can build the C extensions by running:
$ python setup.py build_ext --inplace
.. To finish, make sure `tests`__ pass.
Development with Docker
-----------------------
.. __ #style-guide-running-tests
If you want to work with zipline using a `Docker`__ container, you'll need to
build the ``Dockerfile`` in the Zipline root directory, and then build
``Dockerfile-dev``. Instructions for building both containers can be found in
``Dockerfile`` and ``Dockerfile-dev``, respectively.
.. If you get an error running nosetests after setting up a fresh virtualenv, please try running
.. code-block
.. # where zipline is the name of your virtualenv
.. $ deactivate zipline
.. $ workon zipline
.. Development with Docker
.. -----------------------
..If you want to work with zipline using a `Docker`__ container, you'll need to build the ``Dockerfile`` in the Zipline root directory, and then build ``Dockerfile-dev``. Instructions for building both containers can be found in ``Dockerfile`` and ``Dockerfile-dev``, respectively.
.. __ https://docs.docker.com/get-started/
__ https://docs.docker.com/get-started/
Git Branching Structure
-----------------------
+2 -1
View File
@@ -1,4 +1,5 @@
|
Example Algorithms
==================
@@ -805,7 +806,7 @@ Credits: This code was originally submitted by `Abner Ayala-Acevedo
import pandas as pd
from catalyst import run_algorithm
from catalyst.exchange.exchange_utils import get_exchange_symbols
from catalyst.exchange.utils.exchange_utils import get_exchange_symbols
from catalyst.api import (symbols, )
+2
View File
@@ -1,6 +1,8 @@
.. include:: ../../README.rst
|
|
Table of Contents
-----------------
+2 -1
View File
@@ -298,7 +298,7 @@ Troubleshooting ``pip`` Install
.. _pipenv:
Installing with ``pipenv``
-------------------------
--------------------------
Installing Catalyst via ``pipenv`` is perhaps easier that installing it via
``pip`` itself but you need to install ``pipenv`` first via ``pip``.
@@ -476,6 +476,7 @@ mentioned above are as follows:
default you get 0 as the Value Data)
|
- **The installer has encountered an unexpected error installing this package.
This may indicate a problem with this package. The error code is 2503.**
+1 -1
View File
@@ -113,7 +113,7 @@ Currency symbols (e.g. btc, eth, ltc) follow the Bittrex convention.
Here are some examples:
.. code-block:: json
.. code:: python
# With Bitfinex
bitcoin_usd_asset = symbol('btc_usd')
+32
View File
@@ -2,6 +2,38 @@
Release Notes
=============
Version 0.5.3
^^^^^^^^^^^^^
**Release Date**: 2018-02-09
Bug Fixes
~~~~~~~~~
- Fixed an issue with last candle in backtesting :issue:`219`
Version 0.5.2
^^^^^^^^^^^^^
**Release Date**: 2018-02-08
Bug Fixes
~~~~~~~~~
- Fixed an issue with live candle values :issue:`216` and :issue:`199`
Version 0.5.1
^^^^^^^^^^^^^
**Release Date**: 2018-02-07
Bug Fixes
~~~~~~~~~
- Fixed an issue with orders that stay open :issue:`211`
- Fixed Jupyter issues :issue:`179`
- Fetching multiple tickers in one call to minimize rate limit risks :issue:`174`
- Improved live state presentation :issue:`171`
Build
~~~~~
- Introducing the Enigma Marketplace
Version 0.4.7
^^^^^^^^^^^^^
**Release Date**: 2018-01-19
+6 -1
View File
@@ -11,6 +11,7 @@ Installation: MacOS
|
|
Installation: Windows
---------------------
@@ -21,6 +22,7 @@ Where things go smoothly:
<iframe width="560" height="315" src="https://www.youtube.com/embed/H8HqcEbZmkk" frameborder="0" allowfullscreen></iframe>
|
Where things don't:
.. raw:: html
@@ -29,6 +31,7 @@ Where things don't:
|
|
Backtesting a Strategy
----------------------
@@ -44,6 +47,7 @@ sell. Hopefully, well ride the waves.
|
|
Live Trading a Strategy
-----------------------
@@ -54,5 +58,6 @@ in the previous video, we now take it to trade live against the Bittrex exchange
.. raw:: html
<iframe width="560" height="315" src="https://www.youtube.com/embed/NupiE-Xuglw" frameborder="0" allowfullscreen></iframe>
|
|
|
+1 -1
View File
@@ -16,4 +16,4 @@ fi
jupyter notebook -y --no-browser --notebook-dir=${PROJECT_DIR} \
--certfile=${SSL_CERT_PEM} --keyfile=${SSL_CERT_KEY} --ip='*' \
--config=${CONFIG_PATH}
--config=${CONFIG_PATH} --allow-root
+3 -1
View File
@@ -20,7 +20,9 @@ dependencies:
- bcolz==0.12.1
- bottleneck==1.2.1
- chardet==3.0.4
- ccxt==1.10.774
- ccxt==1.10.1049
- web3==4.0.0b7
- requests-toolbelt==0.8.0
- click==6.7
- contextlib2==0.5.5
- cycler==0.10.0
+3 -1
View File
@@ -81,6 +81,8 @@ empyrical==0.2.1
tables==3.3.0
#Catalyst dependencies
ccxt==1.10.774
ccxt==1.10.1049
boto3==1.4.8
redo==1.6
web3==4.0.0b7
requests-toolbelt==0.8.0
+1 -1
View File
@@ -16,7 +16,7 @@ babel==1.3
docutils==0.12
snowballstemmer==1.2.0
sphinx-rtd-theme==0.1.8
sphinx==1.3.4
sphinx==1.6.7
pbr==1.10.0
mock==2.0.0
+1 -1
View File
@@ -1,4 +1,4 @@
Sphinx>=1.3.2
Sphinx==1.6.7
numpydoc>=0.5.0
sphinx-autobuild==0.6.0
docutils==0.12
+2
View File
@@ -0,0 +1,2 @@
web3==4.0.0b7
requests-toolbelt==0.8.0
+16 -11
View File
@@ -1,8 +1,7 @@
import pandas as pd
from logbook import Logger
from catalyst.testing import ZiplineTestCase
from catalyst.testing.fixtures import WithLogger
from catalyst.exchange.utils.stats_utils import set_print_settings
from .base import BaseExchangeTestCase
from catalyst.exchange.ccxt.ccxt_exchange import CCXT
from catalyst.exchange.exchange_execution import ExchangeLimitOrder
@@ -15,23 +14,23 @@ log = Logger('test_ccxt')
class TestCCXT(BaseExchangeTestCase):
@classmethod
def setup(self):
exchange_name = 'bitfinex'
exchange_name = 'bittrex'
auth = get_exchange_auth(exchange_name)
self.exchange = CCXT(
exchange_name=exchange_name,
key=auth['key'],
secret=auth['secret'],
base_currency='bnb',
base_currency='usdt',
)
self.exchange.init()
def test_order(self):
log.info('creating order')
asset = self.exchange.get_asset('neo_bnb')
asset = self.exchange.get_asset('eth_usdt')
order_id = self.exchange.order(
asset=asset,
style=ExchangeLimitOrder(limit_price=10),
amount=1,
style=ExchangeLimitOrder(limit_price=1000),
amount=1.01,
)
log.info('order created {}'.format(order_id))
assert order_id is not None
@@ -58,24 +57,30 @@ class TestCCXT(BaseExchangeTestCase):
def test_get_candles(self):
log.info('retrieving candles')
candles = self.exchange.get_candles(
freq='30T',
freq='1T',
assets=[self.exchange.get_asset('eth_btc')],
bar_count=200,
start_dt=pd.to_datetime('2017-09-01', utc=True)
# start_dt=pd.to_datetime('2017-09-01', utc=True),
)
for asset in candles:
df = pd.DataFrame(candles[asset])
df.set_index('last_traded', drop=True, inplace=True)
set_print_settings()
print('got {} candles'.format(len(df)))
print(df.head(10))
print(df.tail(10))
pass
def test_tickers(self):
log.info('retrieving tickers')
assets = [
self.exchange.get_asset('iot_usd'),
self.exchange.get_asset('ada_eth'),
self.exchange.get_asset('zrx_eth'),
]
tickers = self.exchange.tickers(assets)
assert len(tickers) == 1
assert len(tickers) == 2
pass
def test_my_trades(self):
@@ -2,6 +2,7 @@ import random
import os
import pandas as pd
from datetime import timedelta
from logbook import TestHandler
from pandas.util.testing import assert_frame_equal
@@ -12,6 +13,7 @@ from catalyst.exchange.utils.exchange_utils import get_candles_df
from catalyst.exchange.utils.factory import get_exchange
from catalyst.exchange.utils.test_utils import output_df, \
select_random_assets
from catalyst.exchange.utils.stats_utils import set_print_settings
pd.set_option('display.expand_frame_repr', False)
pd.set_option('precision', 8)
@@ -58,6 +60,12 @@ class TestSuiteBundle:
log_catcher = TestHandler()
with log_catcher:
symbols = [asset.symbol for asset in assets]
print(
'comparing data for {}/{} with {} timeframe until {}'.format(
exchange.name, symbols, freq, end_dt
)
)
data['bundle'] = data_portal.get_history_window(
assets=assets,
end_dt=end_dt,
@@ -66,6 +74,12 @@ class TestSuiteBundle:
field='close',
data_frequency=data_frequency,
)
set_print_settings()
print(
'the bundle first / last row:\n{}'.format(
data['bundle'].iloc[[-1, 0]]
)
)
candles = exchange.get_candles(
end_dt=end_dt,
freq=freq,
@@ -79,6 +93,11 @@ class TestSuiteBundle:
bar_count=bar_count,
end_dt=end_dt,
)
print(
'the exchange first / last row:\n{}'.format(
data['exchange'].iloc[[-1, 0]]
)
)
for source in data:
df = data[source]
path, folder = output_df(
@@ -106,6 +125,65 @@ class TestSuiteBundle:
pass
def compare_current_with_last_candle(self, exchange, assets, end_dt,
freq, data_frequency, data_portal):
"""
Creates DataFrames from the bundle and exchange for the specified
data set.
Parameters
----------
exchange: Exchange
assets
end_dt
bar_count
freq
data_frequency
data_portal
Returns
-------
"""
data = dict()
assets = sorted(assets, key=lambda a: a.symbol)
log_catcher = TestHandler()
with log_catcher:
symbols = [asset.symbol for asset in assets]
print(
'comparing data for {}/{} with {} timeframe on {}'.format(
exchange.name, symbols, freq, end_dt
)
)
data['candle'] = data_portal.get_history_window(
assets=assets,
end_dt=end_dt,
bar_count=1,
frequency=freq,
field='close',
data_frequency=data_frequency,
)
set_print_settings()
print(
'the bundle first / last row:\n{}'.format(
data['candle'].iloc[[-1]]
)
)
current = data_portal.get_spot_value(
assets=assets,
field='close',
dt=end_dt,
data_frequency=data_frequency,
)
data['current'] = pd.Series(data=current, index=assets)
print(
'the current price:\n{}'.format(
data['current']
)
)
pass
def test_validate_bundles(self):
# exchange_population = 3
asset_population = 3
@@ -139,6 +217,7 @@ class TestSuiteBundle:
if end_dt is None or asset_end_dt < end_dt:
end_dt = asset_end_dt
end_dt = end_dt + timedelta(minutes=3)
dt_range = pd.date_range(
end=end_dt, periods=bar_count, freq=freq
)
@@ -152,3 +231,45 @@ class TestSuiteBundle:
data_portal=data_portal,
)
pass
def test_validate_last_candle(self):
# exchange_population = 3
asset_population = 3
data_frequency = random.choice(['minute'])
# bundle = 'dailyBundle' if data_frequency
# == 'daily' else 'minuteBundle'
# exchanges = select_random_exchanges(
# population=exchange_population,
# features=[bundle],
# ) # Type: list[Exchange]
exchanges = [get_exchange('poloniex', skip_init=True)]
data_portal = TestSuiteBundle.get_data_portal(exchanges)
for exchange in exchanges:
exchange.init()
frequencies = exchange.get_candle_frequencies(data_frequency)
freq = random.sample(frequencies, 1)[0]
assets = select_random_assets(
exchange.assets, asset_population
)
end_dt = None
for asset in assets:
attribute = 'end_{}'.format(data_frequency)
asset_end_dt = getattr(asset, attribute)
if end_dt is None or asset_end_dt < end_dt:
end_dt = asset_end_dt
end_dt = end_dt + timedelta(minutes=3)
self.compare_current_with_last_candle(
exchange=exchange,
assets=assets,
end_dt=end_dt,
freq=freq,
data_frequency=data_frequency,
data_portal=data_portal,
)
pass
@@ -15,7 +15,7 @@ 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
from catalyst.exchange.utils.factory import get_exchanges, get_exchange
log = Logger('TestSuiteExchange')
@@ -90,7 +90,7 @@ class TestSuiteExchange(WithLogger, ZiplineTestCase):
# exchange_population,
# features=['fetchTickers'],
# ) # Type: list[Exchange]
exchanges = list(get_exchanges(['bitfinex']).values())
exchanges = list(get_exchanges(['binance']).values())
for exchange in exchanges:
exchange.init()
@@ -113,10 +113,11 @@ class TestSuiteExchange(WithLogger, ZiplineTestCase):
exchange_population = 3
asset_population = 3
exchanges = select_random_exchanges(
population=exchange_population,
features=['fetchOHLCV'],
) # Type: list[Exchange]
# exchanges = select_random_exchanges(
# population=exchange_population,
# features=['fetchOHLCV'],
# ) # Type: list[Exchange]
exchanges = list(get_exchanges(['binance']).values())
for exchange in exchanges:
exchange.init()
@@ -138,7 +139,6 @@ class TestSuiteExchange(WithLogger, ZiplineTestCase):
assets=assets,
bar_count=bar_count,
start_dt=dt_range[0],
end_dt=dt_range[-1],
)
assert len(candles) == asset_population
@@ -155,13 +155,20 @@ class TestSuiteExchange(WithLogger, ZiplineTestCase):
quote_currency = 'eth'
order_amount = 0.1
exchanges = select_random_exchanges(
population=population,
features=['fetchOrder'],
is_authenticated=True,
base_currency=quote_currency,
) # Type: list[Exchange]
# exchanges = select_random_exchanges(
# population=population,
# features=['fetchOrder'],
# is_authenticated=True,
# base_currency=quote_currency,
# ) # Type: list[Exchange]
exchanges = [
get_exchange(
'binance',
base_currency=quote_currency,
must_authenticate=True,
)
]
log_catcher = TestHandler()
with log_catcher:
for exchange in exchanges:
View File
+36
View File
@@ -0,0 +1,36 @@
from catalyst.marketplace.marketplace import Marketplace
from catalyst.testing.fixtures import WithLogger, ZiplineTestCase
import pandas as pd
class TestMarketplace(WithLogger, ZiplineTestCase):
def test_list(self):
marketplace = Marketplace()
marketplace.list()
pass
def test_register(self):
marketplace = Marketplace()
marketplace.register()
pass
def test_subscribe(self):
marketplace = Marketplace()
marketplace.subscribe('marketcap2222')
pass
def test_ingest(self):
marketplace = Marketplace()
ds_def = marketplace.ingest('marketcap1234')
pass
def test_publish(self):
marketplace = Marketplace()
datadir = '/Users/fredfortier/Downloads/marketcap_test_single'
marketplace.publish('marketcap1234', datadir, False)
pass
def test_clean(self):
marketplace = Marketplace()
marketplace.clean('marketcap')
pass