-
Notifications
You must be signed in to change notification settings - Fork 0
/
ipfs_pubsub.py
executable file
·138 lines (101 loc) · 4.67 KB
/
ipfs_pubsub.py
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
#!/usr/bin/python3
###############################################################################
# IPFS PUBSUB DUMP1090 v0.1
# Copyright (C) 2022 Juan Benitez
# Distributed under GPLv3
###############################################################################
from collections import OrderedDict
import requests, json, logging
from base64 import b64decode, b64encode
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry
###############################################################################
# LOGGING
###############################################################################
loglevel = 'info'
numeric_level = getattr(logging, loglevel.upper())
fmt = '[%(levelname)s] %(message)s'
logging.basicConfig(level=numeric_level, format=fmt)
###############################################################################
# API CLASS
###############################################################################
class IPFS_API:
def __init__(self, host="localhost", port=5001, proto='http'):
self.host = host
self.port = port
self.proto = proto
self.auth = False
self.session = requests.Session()
retry = Retry(connect=10, backoff_factor=2)
adapter = HTTPAdapter(max_retries=retry)
if(proto == 'http'):
self.session.mount('http://', adapter)
else:
self.session.mount('https://', adapter)
self.base_url = self.proto + "://" + self.host + ":" + str(self.port)
def setHttpAuth(self, user, passwd):
self.user = user
self.passwd = passwd
self.auth = True
def setHost(self, newHost):
self.host = newHost
self.base_url = self.proto + "://" + self.host + ":" + str(self.port)
def setPort(self, newPort):
self.port = newPort
self.base_url = self.proto + "://" + self.host + ":" + str(self.port)
@staticmethod
def ipfsb64decode( data ):
d = str(data[1:]) + '==='
return b64decode(d).decode('utf-8')
@staticmethod
def ipfsb64encode( data ):
return 'u' + b64encode( bytes(data.encode('utf-8')) ).decode('utf-8').replace('=', '')
def getPeers( self, topic ):
endpoint = self.base_url + "/api/v0/pubsub/peers?arg="
if(not self.auth):
peers = self.session.post( endpoint + self.ipfsb64encode(topic) )
else:
peers = self.session.post( endpoint + self.ipfsb64encode(topic), auth=(self.user, self.passwd) )
return json.loads( peers.text )['Strings']
def printPeers( self, topic ):
endpoint = self.base_url + "/api/v0/pubsub/peers?arg="
if(not self.auth):
req = self.session.post( endpoint + self.ipfsb64encode(topic) )
else:
req = self.session.post( endpoint + self.ipfsb64encode(topic), auth=(self.user, self.passwd) )
# Check response status codes
if req.status_code != 200:
logging.error("HTTP Request Error Code: %s", req.status_code)
return
print( json.loads( req.text )['Strings'] )
def publishNDJSON( self, topic, dataIN, delimiter='\n' ):
endpoint = self.base_url + "/api/v0/pubsub/pub?arg="
data = { 'file': ( json.dumps(dataIN) + delimiter ) }
req = self.session.post( endpoint + self.ipfsb64encode(topic), files=data )
return req.status_code
def publishOrderedNDJSON( self, topic, data, keyOrder, delimiter='\n' ):
endpoint = self.base_url + "/api/v0/pubsub/pub?arg="
jsonOrderedData = OrderedDict( (k, data[k]) for k in keyOrder )
data = { 'file': ( json.dumps(jsonOrderedData) + delimiter ) }
req = self.session.post( endpoint + self.ipfsb64encode(topic), files=data )
return req.status_code
def subscribe( self, topic, callback ):
endpoint = self.base_url + "/api/v0/pubsub/sub?arg="
# Send request w/wo auth headers
if(not self.auth):
req = self.session.post( endpoint + self.ipfsb64encode(topic), stream=True )
else:
req = self.session.post( endpoint + self.ipfsb64encode(topic), auth=(self.user, self.passwd), stream=True )
# Check response status codes
if req.status_code != 200:
logging.error("HTTP Request Error Code: %s", req.status_code)
return
# Get data and send to callback
for line in req.iter_lines():
json_data = json.loads( line.decode('utf-8') )
callback( json_data['from'], self.ipfsb64decode( json_data['data'] ) )
###############################################################################
# Main
###############################################################################
if __name__ == '__main__':
pass