feat: send warmup data as available

When neither BTC nor LN is running:
send hardware_info

When BTC is running and LN is not:
send hardware_info and btc_info

When both are running:
send the full package

This code doesn't keep track of which connection has received
which data already, so it may send data twice data to the client
when the startup state changes. Especially the hardware info
is rather data intensive. This is OK for now, to keep the code simple.

closes #110
This commit is contained in:
fusion44 2022-06-12 07:59:40 +02:00
parent 15750e8149
commit c08e604d55
No known key found for this signature in database
GPG key ID: 645FA807E935D9D5
2 changed files with 88 additions and 24 deletions

View file

@ -33,8 +33,11 @@ from app.repositories.bitcoin import (
register_bitcoin_zmq_sub,
)
from app.repositories.lightning import initialize_ln_repo, register_lightning_listener
from app.repositories.system import register_hardware_info_gatherer
from app.repositories.utils import get_client_warmup_data
from app.repositories.system import get_hardware_info, register_hardware_info_gatherer
from app.repositories.utils import (
get_bitcoin_client_warmup_data,
get_full_client_warmup_data,
)
from app.routers import apps, bitcoin, lightning, setup, system
from app.utils import SSE, redis_get, send_sse_message
@ -114,10 +117,10 @@ async def _set_startup_status(
api_startup_status.lightning = lightning
if lightning_msg is not None:
api_startup_status.lightning_msg = lightning_msg
await send_sse_message(SSE.SYSTEM_STARTUP_INFO, api_startup_status.dict())
if api_startup_status.is_fully_initialized():
await warmup_new_connections()
loop = asyncio.get_event_loop()
loop.create_task(warmup_new_connections())
await send_sse_message(SSE.SYSTEM_STARTUP_INFO, api_startup_status.dict())
async def _initialize_bitcoin():
@ -214,35 +217,82 @@ async def stream(request: Request):
await q.put(_make_evt_data(SSE.SYSTEM_STARTUP_INFO, api_startup_status.dict()))
if api_startup_status.is_fully_initialized() and len(new_connections) == 1:
loop = asyncio.get_event_loop()
loop.create_task(warmup_new_connections())
loop = asyncio.get_event_loop()
loop.create_task(warmup_new_connections())
return EventSourceResponse(subscribe(request, num_connections - 1, q))
warmup_running = False
async def warmup_new_connections():
# This doesn't keep track of which connection has received
# which data already, so it may send data twice data to the client
# when the startup state changes. Especially the hardware info
# is rather data intensive. This is OK for now, to keep the code simple.
global new_connections
if len(new_connections) == 0:
return
res = await get_client_warmup_data()
global warmup_running
if warmup_running:
logging.debug("Warmup already running, skipping")
return
for c in new_connections:
await asyncio.gather(
*[
c.put(_make_evt_data(SSE.SYSTEM_INFO, res[0].dict())),
c.put(_make_evt_data(SSE.BTC_INFO, res[1].dict())),
c.put(_make_evt_data(SSE.LN_INFO, res[2].dict())),
c.put(_make_evt_data(SSE.LN_INFO_LITE, res[3].dict())),
c.put(_make_evt_data(SSE.LN_FEE_REVENUE, res[4])),
c.put(_make_evt_data(SSE.WALLET_BALANCE, res[5].dict())),
c.put(_make_evt_data(SSE.INSTALLED_APP_STATUS, res[6])),
]
)
warmup_running = True
is_ready = api_startup_status.is_fully_initialized()
new_connections.clear()
if is_ready:
res = await get_full_client_warmup_data()
for c in new_connections:
await asyncio.gather(
*[
c.put(_make_evt_data(SSE.SYSTEM_INFO, res[0].dict())),
c.put(_make_evt_data(SSE.BTC_INFO, res[1].dict())),
c.put(_make_evt_data(SSE.LN_INFO, res[2].dict())),
c.put(_make_evt_data(SSE.LN_INFO_LITE, res[3].dict())),
c.put(_make_evt_data(SSE.LN_FEE_REVENUE, res[4])),
c.put(_make_evt_data(SSE.WALLET_BALANCE, res[5].dict())),
c.put(_make_evt_data(SSE.INSTALLED_APP_STATUS, res[6])),
c.put(_make_evt_data(SSE.HARDWARE_INFO, res[7])),
]
)
new_connections.clear()
if (
api_startup_status.bitcoin == StartupState.DONE
and api_startup_status.lightning != StartupState.DONE
):
res = await get_bitcoin_client_warmup_data()
for c in new_connections:
await asyncio.gather(
*[
c.put(_make_evt_data(SSE.BTC_INFO, res[0].dict())),
c.put(_make_evt_data(SSE.HARDWARE_INFO, res[1])),
]
)
# don't clear new_connections, we'll try again later when api is initialized
if (
api_startup_status.bitcoin != StartupState.DONE
and api_startup_status.lightning != StartupState.DONE
):
# send only the most minimal available data without Bitcoin Core and Lightning running
res = await get_hardware_info()
for c in new_connections:
await c.put(_make_evt_data(SSE.HARDWARE_INFO, res)),
# don't clear new_connections, we'll try again later when api is initialized
warmup_running = False
def _make_evt_data(evt: SSE, data):
d1 = {"event": evt, "data": json.dumps(jsonable_encoder(data))}
# d2 = {"event": evt, "data": jsonable_encoder(data)}
return d1

View file

@ -9,10 +9,23 @@ from app.repositories.lightning import (
get_ln_info_lite,
get_wallet_balance,
)
from app.repositories.system import get_system_info
from app.repositories.system import get_hardware_info, get_system_info
async def get_client_warmup_data() -> List:
async def get_bitcoin_client_warmup_data() -> List:
"""Get the reduced data set needed when the lightning client is not yet ready."""
res = await asyncio.gather(
*[
get_btc_info(),
get_hardware_info(),
]
)
return [*res]
async def get_full_client_warmup_data() -> List:
"""Get the full data set needed when the lightning client is not yet ready."""
res = await asyncio.gather(
*[
get_system_info(),
@ -22,6 +35,7 @@ async def get_client_warmup_data() -> List:
get_fee_revenue(),
get_wallet_balance(),
get_app_status(),
get_hardware_info(),
]
)
return [*res]