FazBrowse GitHub Viewer
|
Trending
|
URL:
|
Home
Tools:
[Download Repo ZIP]
[View Raw Code]
[Original HTTPS Page]
python-bitshares/bitsharesapi/websocket.py at master · allrob23/python-bitshares · GitHub
allrob23
/
python-bitshares
Public
forked from
bitshares/python-bitshares
Notifications
You must be signed in to change notification settings
Fork
0
Star
0
Code
Pull requests
0
Actions
Projects
Security and quality
0
Insights
Additional navigation options
Code
Pull requests
Actions
Projects
Security and quality
Insights
Expand file tree
Breadcrumbs
python-bitshares
/
bitsharesapi
/
websocket.py
Copy path
More file actions
More file actions
Latest commit
History
History
History
395 lines (325 loc) · 13 KB
Breadcrumbs
python-bitshares
/
bitsharesapi
/
websocket.py
Copy path
File metadata and controls
395 lines (325 loc) · 13 KB
Raw
Copy raw file
Download raw file
Open symbols panel
Edit and raw actions
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
# -*- coding: utf-8 -*-
import
json
import
time
import
signal
import
logging
import
threading
import
websocket
import
traceback
from
itertools
import
cycle
from
events
import
Events
from
.
exceptions
import
NumRetriesReached
# This restores the default Ctrl+C signal handler, which just kills the process
signal
.
signal
(
signal
.
SIGINT
,
signal
.
SIG_DFL
)
log
=
logging
.
getLogger
(
__name__
)
# logging.basicConfig(level=logging.DEBUG)
class
BitSharesWebsocket
(
Events
):
"""
Create a websocket connection and request push notifications.
:param str urls: Either a single Websocket URL, or a list of URLs
:param str user: Username for Authentication
:param str password: Password for Authentication
:param list accounts: list of account names or ids to get push notifications for
:param list markets: list of asset_ids, e.g. ``[['1.3.0', '1.3.121']]``
:param list objects: list of objects id's you'd like to be notified when changing
:param int keep_alive: seconds between a ping to the backend (defaults to 25seconds)
After instanciating this class, you can add event slots for:
* ``on_tx``
* ``on_object``
* ``on_block``
* ``on_account``
* ``on_market``
which will be called accordingly with the notification
message received from the BitShares node:
.. code-block:: python
ws = BitSharesWebsocket(
"wss://node.testnet.bitshares.eu",
objects=["2.0.x", "2.1.x", "1.3.x"]
)
ws.on_object += print
ws.run_forever()
Notices:
* ``on_account``:
.. code-block:: js
{'id': '2.6.29',
'lifetime_fees_paid': '44257768405',
'most_recent_op': '2.9.1195638',
'owner': '1.2.29',
'pending_fees': 0,
'pending_vested_fees': 100,
'total_core_in_orders': '6788960277634',
'total_ops': 505865}
* ``on_block``:
.. code-block:: js
'0062f19df70ecf3a478a84b4607d9ad8b3e3b607'
* ``on_tx``:
.. code-block:: js
{'expiration': '2017-02-23T09:33:22',
'extensions': [],
'operations': [[0,
{'amount': {'amount': 100000, 'asset_id': '1.3.0'},
'extensions': [],
'fee': {'amount': 100, 'asset_id': '1.3.0'},
'from': '1.2.29',
'to': '1.2.17'}]],
'ref_block_num': 62001,
'ref_block_prefix': 390951726,
'signatures': ['20784246dc1064ed5f87dbbb9aaff3fcce052135269a8653fb500da46e7068bec56e85ea997b8d250a9cc926777c700eed41e34ba1cabe65940965ebe133ff9098']}
* ``on_market``:
.. code-block:: js
['1.7.68612']
"""
__events__
=
[
"on_tx"
,
"on_object"
,
"on_block"
,
"on_account"
,
"on_market"
]
def
__init__
(
self
,
urls
,
user
=
""
,
password
=
""
,
*
args
,
accounts
=
None
,
markets
=
None
,
objects
=
None
,
on_tx
=
None
,
on_object
=
None
,
on_block
=
None
,
on_account
=
None
,
on_market
=
None
,
keep_alive
=
25
,
num_retries
=
-
1
,
**
kwargs
):
self
.
num_retries
=
num_retries
self
.
keepalive
=
None
self
.
_request_id
=
0
self
.
ws
=
None
self
.
user
=
user
self
.
password
=
password
self
.
keep_alive
=
keep_alive
self
.
run_event
=
threading
.
Event
()
if
isinstance
(
urls
,
cycle
):
self
.
urls
=
urls
elif
isinstance
(
urls
,
list
):
self
.
urls
=
cycle
(
urls
)
else
:
self
.
urls
=
cycle
([
urls
])
# Instanciate Events
Events
.
__init__
(
self
)
self
.
events
=
Events
()
# Store the objects we are interested in
self
.
subscription_accounts
=
accounts
or
[]
self
.
subscription_markets
=
markets
or
[]
self
.
subscription_objects
=
objects
or
[]
if
on_tx
:
self
.
on_tx
+=
on_tx
if
on_object
:
self
.
on_object
+=
on_object
if
on_block
:
self
.
on_block
+=
on_block
if
on_account
:
self
.
on_account
+=
on_account
if
on_market
:
self
.
on_market
+=
on_market
def
cancel_subscriptions
(
self
):
self
.
cancel_all_subscriptions
()
def
on_open
(
self
,
*
args
,
**
kwargs
):
"""
This method will be called once the websocket connection is established. It
will.
* login,
* register to the database api, and
* subscribe to the objects defined if there is a
callback/slot available for callbacks
"""
self
.
login
(
self
.
user
,
self
.
password
,
api_id
=
1
)
self
.
database
(
api_id
=
1
)
self
.
__set_subscriptions
()
self
.
keepalive
=
threading
.
Thread
(
target
=
self
.
_ping
)
self
.
keepalive
.
start
()
def
reset_subscriptions
(
self
,
accounts
=
None
,
markets
=
None
,
objects
=
None
):
self
.
subscription_accounts
=
accounts
or
[]
self
.
subscription_markets
=
markets
or
[]
self
.
subscription_objects
=
objects
or
[]
self
.
__set_subscriptions
()
def
__set_subscriptions
(
self
):
self
.
cancel_all_subscriptions
()
# Subscribe to events on the Backend and give them a
# callback number that allows us to identify the event
if
len
(
self
.
on_object
)
or
len
(
self
.
subscription_accounts
):
self
.
set_subscribe_callback
(
self
.
__events__
.
index
(
"on_object"
),
False
)
if
self
.
subscription_accounts
and
self
.
on_account
:
# Unfortunately, account subscriptions don't have their own
# callback number
log
.
debug
(
"Subscribing to accounts %s"
%
str
(
self
.
subscription_accounts
))
self
.
get_full_accounts
(
self
.
subscription_accounts
,
True
)
if
self
.
subscription_markets
and
self
.
on_market
:
log
.
debug
(
"Subscribing to markets %s"
%
str
(
self
.
subscription_markets
))
for
market
in
self
.
subscription_markets
:
# Technially, every market could have it's own
# callback number
self
.
subscribe_to_market
(
self
.
__events__
.
index
(
"on_market"
),
market
[
0
],
market
[
1
]
)
if
len
(
self
.
on_tx
):
self
.
set_pending_transaction_callback
(
self
.
__events__
.
index
(
"on_tx"
))
if
len
(
self
.
on_block
):
self
.
set_block_applied_callback
(
self
.
__events__
.
index
(
"on_block"
))
def
_ping
(
self
):
# We keep the connection alive by requesting a short object
while
not
self
.
run_event
.
wait
(
self
.
keep_alive
):
log
.
debug
(
"Sending ping"
)
self
.
get_objects
([
"2.8.0"
])
def
process_notice
(
self
,
notice
):
"""
This method is called on notices that need processing.
Here, we call ``on_object`` and ``on_account`` slots.
"""
id
=
notice
[
"id"
]
_a
,
_b
,
_
=
id
.
split
(
"."
)
if
id
in
self
.
subscription_objects
:
self
.
on_object
(
notice
)
elif
"."
.
join
([
_a
,
_b
,
"x"
])
in
self
.
subscription_objects
:
self
.
on_object
(
notice
)
elif
id
[:
4
]
==
"2.6."
:
# Treat account updates separately
self
.
on_account
(
notice
)
def
on_message
(
self
,
reply
,
*
args
,
**
kwargs
):
"""
This method is called by the websocket connection on every message that is
received.
If we receive a ``notice``, we hand over post-processing and signalling of
events to ``process_notice``.
"""
if
isinstance
(
reply
,
websocket
.
WebSocketApp
):
reply
=
args
[
0
]
log
.
debug
(
"Received message: %s"
%
str
(
reply
))
data
=
{}
try
:
data
=
json
.
loads
(
reply
,
strict
=
False
)
except
ValueError
:
raise
ValueError
(
"API node returned invalid format. Expected JSON!"
)
if
data
.
get
(
"method"
)
==
"notice"
:
id
=
data
[
"params"
][
0
]
if
id
>=
len
(
self
.
__events__
):
log
.
critical
(
"Received an id that is out of range
\n
\n
"
+
str
(
data
))
return
# This is a "general" object change notification
if
id
==
self
.
__events__
.
index
(
"on_object"
):
# Let's see if a specific object has changed
for
notice
in
data
[
"params"
][
1
]:
try
:
if
"id"
in
notice
:
self
.
process_notice
(
notice
)
else
:
for
obj
in
notice
:
if
"id"
in
obj
:
self
.
process_notice
(
obj
)
except
Exception
as
e
:
log
.
critical
(
"Error in process_notice: {}
\n
\n
{}"
.
format
(
str
(
e
),
traceback
.
format_exc
)
)
else
:
try
:
callbackname
=
self
.
__events__
[
id
]
log
.
debug
(
"Patching through to call %s"
%
callbackname
)
[
getattr
(
self
.
events
,
callbackname
)(
x
)
for
x
in
data
[
"params"
][
1
]]
except
Exception
as
e
:
log
.
critical
(
"Error in {}: {}
\n
\n
{}"
.
format
(
callbackname
,
str
(
e
),
traceback
.
format_exc
()
)
)
def
on_error
(
self
,
error
,
*
args
,
**
kwargs
):
"""Called on websocket errors."""
log
.
exception
(
error
)
def
on_close
(
self
,
*
args
,
**
kwargs
):
"""Called when websocket connection is closed."""
log
.
debug
(
"Closing WebSocket connection with {}"
.
format
(
self
.
url
))
def
run_forever
(
self
,
*
args
,
**
kwargs
):
"""
This method is used to run the websocket app continuously.
It will execute callbacks as defined and try to stay connected with the provided
APIs
"""
cnt
=
0
while
not
self
.
run_event
.
is_set
():
cnt
+=
1
self
.
url
=
next
(
self
.
urls
)
log
.
debug
(
"Trying to connect to node %s"
%
self
.
url
)
try
:
# websocket.enableTrace(True)
self
.
ws
=
websocket
.
WebSocketApp
(
self
.
url
,
on_message
=
self
.
on_message
,
on_error
=
self
.
on_error
,
on_close
=
self
.
on_close
,
on_open
=
self
.
on_open
,
)
self
.
ws
.
run_forever
()
except
websocket
.
WebSocketException
:
if
self
.
num_retries
>=
0
and
cnt
>
self
.
num_retries
:
raise
NumRetriesReached
()
sleeptime
=
(
cnt
-
1
)
*
2
if
cnt
<
10
else
10
if
sleeptime
:
log
.
warning
(
"Lost connection to node during wsconnect(): %s (%d/%d) "
%
(
self
.
url
,
cnt
,
self
.
num_retries
)
+
"Retrying in %d seconds"
%
sleeptime
)
time
.
sleep
(
sleeptime
)
except
KeyboardInterrupt
:
self
.
ws
.
keep_running
=
False
return
except
Exception
as
e
:
log
.
critical
(
"{}
\n
\n
{}"
.
format
(
str
(
e
),
traceback
.
format_exc
()))
def
close
(
self
,
*
args
,
**
kwargs
):
"""Closes the websocket connection and waits for the ping thread to close."""
self
.
run_event
.
set
()
self
.
ws
.
close
()
if
self
.
keepalive
and
self
.
keepalive
.
is_alive
():
self
.
keepalive
.
join
()
def
get_request_id
(
self
):
self
.
_request_id
+=
1
return
self
.
_request_id
""" RPC Calls
"""
def
rpcexec
(
self
,
payload
):
"""
Execute a call by sending the payload.
:param dict payload: Payload data
:raises ValueError: if the server does not respond in proper JSON format
:raises RPCError: if the server returns an error
"""
log
.
debug
(
json
.
dumps
(
payload
))
self
.
ws
.
send
(
json
.
dumps
(
payload
,
ensure_ascii
=
False
).
encode
(
"utf8"
))
def
__getattr__
(
self
,
name
):
"""Map all methods to RPC calls and pass through the arguments."""
if
name
in
self
.
__events__
:
return
getattr
(
self
.
events
,
name
)
def
method
(
*
args
,
**
kwargs
):
# Sepcify the api to talk to
if
"api_id"
not
in
kwargs
:
if
"api"
in
kwargs
:
if
kwargs
[
"api"
]
in
self
.
api_id
and
self
.
api_id
[
kwargs
[
"api"
]]:
api_id
=
self
.
api_id
[
kwargs
[
"api"
]]
else
:
raise
ValueError
(
"Unknown API! "
"Verify that you have registered to %s"
%
kwargs
[
"api"
]
)
else
:
api_id
=
0
else
:
api_id
=
kwargs
[
"api_id"
]
# let's be able to define the num_retries per query
self
.
num_retries
=
kwargs
.
get
(
"num_retries"
,
self
.
num_retries
)
query
=
{
"method"
:
"call"
,
"params"
: [
api_id
,
name
,
list
(
args
)],
"jsonrpc"
:
"2.0"
,
"id"
:
self
.
get_request_id
(),
}
r
=
self
.
rpcexec
(
query
)
return
r
return
method
Back
|
FazBrowse Home
|
New Git URL