instantiate objects inside async loop
This commit is contained in:
parent
a96ceba6f5
commit
14436b3955
12 changed files with 56 additions and 37 deletions
|
@ -1,4 +1,4 @@
|
|||
from .ledger import Ledger, RegTestLedger, TestNetLedger
|
||||
from .ledger import Ledger, RegTestLedger, TestNetLedger, ledger_class_from_name
|
||||
from .transaction import Transaction, Output, Input
|
||||
from .bcd_data_stream import BCDataStream
|
||||
from .dewies import dewies_to_lbc, lbc_to_dewies, dict_values_to_lbc
|
||||
|
|
|
@ -88,6 +88,7 @@ class Lbrycrd:
|
|||
def temp_regtest(cls):
|
||||
return cls(RegTestLedger(
|
||||
Config.with_same_dir(tempfile.mkdtemp()).set(
|
||||
blockchain="regtest",
|
||||
lbrycrd_rpc_port=9245 + 2, # avoid conflict with default rpc port
|
||||
lbrycrd_peer_port=9246 + 2, # avoid conflict with default peer port
|
||||
lbrycrd_zmq_blocks="tcp://127.0.0.1:29002" # avoid conflict with default port
|
||||
|
|
|
@ -1,6 +1,6 @@
|
|||
import typing
|
||||
from binascii import unhexlify
|
||||
from string import hexdigits
|
||||
from typing import TYPE_CHECKING, Type
|
||||
|
||||
from lbry.crypto.hash import hash160, double_sha256
|
||||
from lbry.crypto.base58 import Base58
|
||||
|
@ -10,7 +10,7 @@ from .checkpoints import HASHES
|
|||
from .dewies import lbc_to_dewies
|
||||
|
||||
|
||||
if typing.TYPE_CHECKING:
|
||||
if TYPE_CHECKING:
|
||||
from lbry.conf import Config
|
||||
|
||||
|
||||
|
@ -164,3 +164,11 @@ class RegTestLedger(Ledger):
|
|||
genesis_bits = 0x207fffff
|
||||
target_timespan = 1
|
||||
checkpoints = {}
|
||||
|
||||
|
||||
def ledger_class_from_name(name) -> Type[Ledger]:
|
||||
return {
|
||||
Ledger.network_name: Ledger,
|
||||
TestNetLedger.network_name: TestNetLedger,
|
||||
RegTestLedger.network_name: RegTestLedger
|
||||
}[name]
|
||||
|
|
18
lbry/cli.py
18
lbry/cli.py
|
@ -10,7 +10,7 @@ from docopt import docopt
|
|||
|
||||
from lbry import __version__
|
||||
from lbry.conf import Config, CLIConfig
|
||||
from lbry.service import Daemon, Client
|
||||
from lbry.service import Daemon, Client, FullNode, LightClient
|
||||
from lbry.service.metadata import interface
|
||||
|
||||
|
||||
|
@ -115,8 +115,7 @@ def get_argument_parser():
|
|||
help='Start LBRY Network interface.'
|
||||
)
|
||||
start.add_argument(
|
||||
'--full-node', dest='full_node', action="store_true",
|
||||
help='Start a full node with local blockchain data, requires lbrycrd.'
|
||||
"service", choices=[LightClient.name, FullNode.name], default=LightClient.name, nargs="?"
|
||||
)
|
||||
start.add_argument(
|
||||
'--quiet', dest='quiet', action="store_true",
|
||||
|
@ -173,10 +172,11 @@ def ensure_directory_exists(path: str):
|
|||
async def execute_command(conf, method, params):
|
||||
client = Client(f"http://{conf.api}/ws")
|
||||
await client.connect()
|
||||
resp = await (await client.send(method, **params)).first
|
||||
print(resp)
|
||||
responses = await client.send(method, **params)
|
||||
result = await responses.first
|
||||
await client.disconnect()
|
||||
return resp
|
||||
print(result)
|
||||
return result
|
||||
|
||||
|
||||
def normalize_value(x, key=None):
|
||||
|
@ -247,10 +247,10 @@ def main(argv=None):
|
|||
elif args.command == 'start':
|
||||
if args.help:
|
||||
args.start_parser.print_help()
|
||||
elif args.full_node:
|
||||
return Daemon.from_config(conf).run()
|
||||
elif args.service == FullNode.name:
|
||||
return Daemon.from_config(FullNode, conf).run()
|
||||
else:
|
||||
print('Only `start --full-node` is currently supported.')
|
||||
print(f'Only `start {FullNode.name}` is currently supported.')
|
||||
elif args.command == 'install':
|
||||
if args.help:
|
||||
args.install_parser.print_help()
|
||||
|
|
|
@ -511,7 +511,7 @@ class Config(CLIConfig):
|
|||
)
|
||||
console = StringChoice(
|
||||
"Basic text console output or advanced colored output with progress bars.",
|
||||
["basic", "advanced"], "advanced"
|
||||
["basic", "advanced", "none"], "advanced"
|
||||
)
|
||||
|
||||
# directories
|
||||
|
@ -615,14 +615,14 @@ class Config(CLIConfig):
|
|||
comment_server = String("Comment server API URL", "https://comments.lbry.com/api")
|
||||
|
||||
# blockchain
|
||||
blockchain = StringChoice("Blockchain network type.", ["mainnet", "regtest", "testnet"], "mainnet")
|
||||
lbrycrd_rpc_user = String("Username for connecting to lbrycrd.", "rpcuser")
|
||||
lbrycrd_rpc_pass = String("Password for connecting to lbrycrd.", "rpcpassword")
|
||||
lbrycrd_rpc_host = String("Hostname for connecting to lbrycrd.", "localhost")
|
||||
lbrycrd_rpc_port = Integer("Port for connecting to lbrycrd.", 9245)
|
||||
lbrycrd_peer_port = Integer("Port for connecting to lbrycrd.", 9246)
|
||||
lbrycrd_peer_port = Integer("Peer port for lbrycrd.", 9246)
|
||||
lbrycrd_zmq_blocks = String("ZMQ block events address.")
|
||||
lbrycrd_dir = Path("Directory containing lbrycrd data.", metavar='DIR')
|
||||
blockchain_name = String("Blockchain name - lbrycrd_main, lbrycrd_regtest, or lbrycrd_testnet", 'lbrycrd_main')
|
||||
spv_address_filters = Toggle(
|
||||
"Generate Golomb-Rice coding filters for blocks and transactions. Enables "
|
||||
"light client to synchronize with a full node.",
|
||||
|
@ -703,7 +703,7 @@ class Config(CLIConfig):
|
|||
def db_url_or_default(self):
|
||||
if self.db_url:
|
||||
return self.db_url
|
||||
return 'sqlite:///'+os.path.join(self.data_dir, f'{self.blockchain_name}.db')
|
||||
return 'sqlite:///'+os.path.join(self.data_dir, f'{self.blockchain}.db')
|
||||
|
||||
|
||||
def get_windows_directories() -> Tuple[str, str, str, str]:
|
||||
|
|
|
@ -3,7 +3,7 @@ import sys
|
|||
import time
|
||||
import itertools
|
||||
import logging
|
||||
from typing import Dict, Any
|
||||
from typing import Dict, Any, Type
|
||||
from tempfile import TemporaryFile
|
||||
|
||||
from tqdm.std import tqdm, Bar
|
||||
|
@ -490,3 +490,7 @@ class Advanced(Basic):
|
|||
# else:
|
||||
# self.stderr.flush(self.bars['read'].write)
|
||||
# self.update_progress(e, d)
|
||||
|
||||
|
||||
def console_class_from_name(name) -> Type[Console]:
|
||||
return {'basic': Basic, 'advanced': Advanced}.get(name, Console)
|
||||
|
|
|
@ -1,4 +1,5 @@
|
|||
from .api import API, Client
|
||||
from .base import Service
|
||||
from .daemon import Daemon, jsonrpc_dumps_pretty
|
||||
from .full_node import FullNode
|
||||
from .light_client import LightClient
|
||||
|
|
|
@ -58,6 +58,7 @@ class Service:
|
|||
This is the programmatic api (as compared to API)
|
||||
"""
|
||||
|
||||
name: str
|
||||
sync: Sync
|
||||
|
||||
def __init__(self, ledger: Ledger):
|
||||
|
|
|
@ -2,19 +2,18 @@ import json
|
|||
import signal
|
||||
import asyncio
|
||||
import logging
|
||||
|
||||
from weakref import WeakSet
|
||||
from asyncio.runners import _cancel_all_tasks
|
||||
from typing import Type
|
||||
|
||||
from aiohttp.web import Application, AppRunner, WebSocketResponse, TCPSite, Response
|
||||
from aiohttp.http_websocket import WSMsgType, WSCloseCode
|
||||
|
||||
from lbry.conf import Config
|
||||
from lbry.console import Console, console_class_from_name
|
||||
from lbry.service import API, Service
|
||||
from lbry.service.json_encoder import JSONResponseEncoder
|
||||
from lbry.blockchain.ledger import Ledger
|
||||
from lbry.service.base import Service
|
||||
from lbry.service.api import API
|
||||
from lbry.service.full_node import FullNode
|
||||
from lbry.console import Console, Advanced as AdvancedConsole, Basic as BasicConsole
|
||||
from lbry.blockchain.ledger import ledger_class_from_name
|
||||
|
||||
|
||||
log = logging.getLogger(__name__)
|
||||
|
@ -74,8 +73,8 @@ class Daemon:
|
|||
Mostly connects API to aiohttp stuff.
|
||||
Handles starting and stopping API
|
||||
"""
|
||||
def __init__(self, loop: asyncio.AbstractEventLoop, service: Service, console: Console):
|
||||
self._loop = loop
|
||||
def __init__(self, service: Service, console: Console):
|
||||
self._loop = asyncio.get_running_loop()
|
||||
self.service = service
|
||||
self.conf = service.conf
|
||||
self.console = console
|
||||
|
@ -95,14 +94,17 @@ class Daemon:
|
|||
self.runner = AppRunner(self.app)
|
||||
|
||||
@classmethod
|
||||
def from_config(cls, conf):
|
||||
loop = asyncio.new_event_loop()
|
||||
service = FullNode(Ledger(conf))
|
||||
if conf.console == "advanced":
|
||||
console = AdvancedConsole(service)
|
||||
else:
|
||||
console = BasicConsole(service)
|
||||
return cls(loop, service, console)
|
||||
def from_config(cls, service_class: Type[Service], conf: Config, ) -> 'Daemon':
|
||||
|
||||
async def setup():
|
||||
ledger_class = ledger_class_from_name(conf.blockchain)
|
||||
ledger = ledger_class(conf)
|
||||
service = service_class(ledger)
|
||||
console_class = console_class_from_name(conf.console)
|
||||
console = console_class(service)
|
||||
return cls(service, console)
|
||||
|
||||
return asyncio.new_event_loop().run_until_complete(setup())
|
||||
|
||||
def run(self):
|
||||
for sig in (signal.SIGHUP, signal.SIGTERM, signal.SIGINT):
|
||||
|
|
|
@ -14,6 +14,8 @@ log = logging.getLogger(__name__)
|
|||
|
||||
class FullNode(Service):
|
||||
|
||||
name = "node"
|
||||
|
||||
sync: BlockchainSync
|
||||
|
||||
def __init__(self, ledger: Ledger, chain: Lbrycrd = None):
|
||||
|
|
|
@ -12,6 +12,8 @@ log = logging.getLogger(__name__)
|
|||
|
||||
class LightClient(Service):
|
||||
|
||||
name = "client"
|
||||
|
||||
sync: SPVSync
|
||||
|
||||
def __init__(self, ledger: Ledger):
|
||||
|
|
|
@ -7,7 +7,6 @@ from unittest import TestCase
|
|||
|
||||
from lbry import Daemon, FullNode
|
||||
from lbry.cli import execute_command
|
||||
from lbry.console import Console
|
||||
from lbry.blockchain.lbrycrd import Lbrycrd
|
||||
from lbry.testcase import CommandTestCase
|
||||
|
||||
|
@ -23,9 +22,8 @@ class TestShutdown(TestCase):
|
|||
chain_loop.run_until_complete(chain.ensure())
|
||||
chain_loop.run_until_complete(chain.start())
|
||||
chain_loop.run_until_complete(chain.generate(1))
|
||||
chain.ledger.conf.set(workers=2)
|
||||
service = FullNode(chain.ledger)
|
||||
daemon = Daemon(service, Console(service))
|
||||
chain.ledger.conf.set(workers=2, console='none')
|
||||
daemon = Daemon.from_config(FullNode, chain.ledger.conf)
|
||||
|
||||
def send_signal():
|
||||
time.sleep(2)
|
||||
|
|
Loading…
Add table
Reference in a new issue