Merge remote-tracking branch 'riccardobl/master'
This commit is contained in:
commit
60676ff15b
4 changed files with 32 additions and 7 deletions
3
.gitignore
vendored
3
.gitignore
vendored
|
|
@ -9,4 +9,5 @@ node_modules
|
||||||
data
|
data
|
||||||
.vscode
|
.vscode
|
||||||
package.json
|
package.json
|
||||||
package-lock.json
|
package-lock.json
|
||||||
|
dump
|
||||||
|
|
@ -2,7 +2,7 @@
|
||||||
"repos": [
|
"repos": [
|
||||||
{
|
{
|
||||||
"id": "nwcprovider",
|
"id": "nwcprovider",
|
||||||
"organisation": "lnbits",
|
"organisation": "riccardobl",
|
||||||
"repository": "nwcprovider"
|
"repository": "nwcprovider"
|
||||||
}
|
}
|
||||||
]
|
]
|
||||||
|
|
|
||||||
32
nwcp.py
32
nwcp.py
|
|
@ -17,7 +17,12 @@ from loguru import logger
|
||||||
from pydantic import BaseModel
|
from pydantic import BaseModel
|
||||||
|
|
||||||
|
|
||||||
class MainSubscription(BaseModel):
|
class RateLimit:
|
||||||
|
backoff: int = 0
|
||||||
|
last_attempt_time: int = 0
|
||||||
|
|
||||||
|
|
||||||
|
class MainSubscription:
|
||||||
def __init__(self):
|
def __init__(self):
|
||||||
self.requests_sub_id: Optional[str] = None
|
self.requests_sub_id: Optional[str] = None
|
||||||
self.responses_sub_id: Optional[str] = None
|
self.responses_sub_id: Optional[str] = None
|
||||||
|
|
@ -90,6 +95,7 @@ class NWCServiceProvider(BaseModel):
|
||||||
|
|
||||||
# Subscription
|
# Subscription
|
||||||
self.sub = None
|
self.sub = None
|
||||||
|
self.rate_limit: Dict[str, RateLimit] = {}
|
||||||
|
|
||||||
# websocket connection
|
# websocket connection
|
||||||
self.ws = None
|
self.ws = None
|
||||||
|
|
@ -203,6 +209,23 @@ class NWCServiceProvider(BaseModel):
|
||||||
logger.debug("Waiting for connection...")
|
logger.debug("Waiting for connection...")
|
||||||
await asyncio.sleep(1)
|
await asyncio.sleep(1)
|
||||||
|
|
||||||
|
async def _ratelimit(self, unit: str, max_sleep_time: int = 120) -> None:
|
||||||
|
limit: Optional[RateLimit] = self.rate_limit.get(unit)
|
||||||
|
if not limit:
|
||||||
|
self.rate_limit[unit] = limit = RateLimit()
|
||||||
|
|
||||||
|
if time.time() - limit.last_attempt_time > max_sleep_time:
|
||||||
|
# reset backoff if action lasted more than max_sleep_time
|
||||||
|
limit.backoff = 0
|
||||||
|
else:
|
||||||
|
# increase backoff
|
||||||
|
limit.backoff = (
|
||||||
|
min(limit.backoff * 2, max_sleep_time) if limit.backoff > 0 else 1
|
||||||
|
)
|
||||||
|
logger.debug("Sleeping for " + str(limit.backoff) + " seconds before " + unit)
|
||||||
|
await asyncio.sleep(limit.backoff)
|
||||||
|
limit.last_attempt_time = int(time.time())
|
||||||
|
|
||||||
async def _subscribe(self):
|
async def _subscribe(self):
|
||||||
"""
|
"""
|
||||||
[Re]Subscribe to receive nip 47 requests and responses from the relay
|
[Re]Subscribe to receive nip 47 requests and responses from the relay
|
||||||
|
|
@ -387,6 +410,7 @@ class NWCServiceProvider(BaseModel):
|
||||||
+ info
|
+ info
|
||||||
+ " ... resubscribing..."
|
+ " ... resubscribing..."
|
||||||
)
|
)
|
||||||
|
await self._ratelimit("subscribing")
|
||||||
await self._subscribe()
|
await self._subscribe()
|
||||||
|
|
||||||
async def _on_message(self, ws, message: str):
|
async def _on_message(self, ws, message: str):
|
||||||
|
|
@ -407,7 +431,7 @@ class NWCServiceProvider(BaseModel):
|
||||||
elif msg[0] == "OK":
|
elif msg[0] == "OK":
|
||||||
pass
|
pass
|
||||||
else:
|
else:
|
||||||
raise Exception("Unknown message type")
|
raise Exception("Unknown message type " + str(msg[0]))
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error("Error parsing event: " + str(e))
|
logger.error("Error parsing event: " + str(e))
|
||||||
|
|
||||||
|
|
@ -447,8 +471,8 @@ class NWCServiceProvider(BaseModel):
|
||||||
self.connected = False
|
self.connected = False
|
||||||
if not self._is_shutting_down():
|
if not self._is_shutting_down():
|
||||||
# Wait some time before reconnecting
|
# Wait some time before reconnecting
|
||||||
logger.debug("Reconnecting to NWC relay in 5 seconds...")
|
logger.debug("Reconnecting to NWC relay...")
|
||||||
await asyncio.sleep(5)
|
await self._ratelimit("connecting")
|
||||||
|
|
||||||
def _encrypt_content(
|
def _encrypt_content(
|
||||||
self, content: str, pubkey_hex: str, iv_seed: Optional[int] = None
|
self, content: str, pubkey_hex: str, iv_seed: Optional[int] = None
|
||||||
|
|
|
||||||
|
|
@ -615,7 +615,7 @@
|
||||||
try {
|
try {
|
||||||
const response = await LNbits.api.request(
|
const response = await LNbits.api.request(
|
||||||
"GET",
|
"GET",
|
||||||
"/nwcprovider/api/v1/nwc?includeExpired=true&calculateSpendBudget=true",
|
"/nwcprovider/api/v1/nwc?include_expired=true&calculate_spent_budget=true",
|
||||||
wallet.adminkey,
|
wallet.adminkey,
|
||||||
);
|
);
|
||||||
this.nwcs = response.data;
|
this.nwcs = response.data;
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue