rpc.py 2.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475
  1. import asyncio, json, random
  2. from flask import abort
  3. # TODO have a single channel for the whole app
  4. # Class representing the channel with the JSON-RPC server.
  5. class Channel:
  6. def __init__(self, reader, writer):
  7. self.reader = reader
  8. self.writer = writer
  9. async def readline(self):
  10. if not (line := await self.reader.readline()):
  11. self.writer.close()
  12. return None
  13. # Strip the newline
  14. return line[:-1].decode()
  15. async def receive(self):
  16. if (plaintext := await self.readline()) is None:
  17. return None
  18. message = plaintext
  19. response = json.loads(message)
  20. return response
  21. async def send(self, obj):
  22. message = json.dumps(obj)
  23. data = message.encode()
  24. self.writer.write(data + b"\n")
  25. await self.writer.drain()
  26. # Create a Channel for given server.
  27. async def create_channel(server_name, port):
  28. try:
  29. reader, writer = await asyncio.open_connection(server_name, port)
  30. except ConnectionRefusedError:
  31. print(f"Error: Connection Refused to '{server_name}:{port}', Either because the daemon is down, is currently syncing or wrong url.")
  32. abort(500)
  33. channel = Channel(reader, writer)
  34. return channel
  35. # Execute a request towards the JSON-RPC server
  36. async def query(method, params, server_name, port):
  37. channel = await create_channel(server_name, port)
  38. request = {
  39. "id": random.randint(0, 2**32),
  40. "method": method,
  41. "params": params,
  42. "jsonrpc": "2.0",
  43. }
  44. await channel.send(request)
  45. response = await channel.receive()
  46. # Closed connect returns None
  47. if response is None:
  48. print("error: connection with server was closed", file=sys.stderr)
  49. abort(500)
  50. if "error" in response:
  51. error = response["error"]
  52. errcode, errmsg = error["code"], error["message"]
  53. print(f"error: {errcode} - {errmsg}", file=sys.stderr)
  54. abort(500)
  55. return response["result"]
  56. # Retrieve last n blocks from blockchain-explorer daemon
  57. async def get_last_n_blocks(n, server_name, port):
  58. return await query("blocks.get_last_n_blocks", [str(n)], server_name, int(port))
  59. # Retrieve basic statistics from blockchain-explorer daemon
  60. async def get_basic_statistics(server_name, port):
  61. return await query("statistics.get_basic_statistics", [], server_name, int(port))