Compare commits

..
Author SHA1 Message Date
prathamesh ed4ed48375 Add mutex lock file to registry CLI wrapper class
Lint Checks / Run linter (pull_request) Successful in 44s
Deploy Test / Run deploy test suite (pull_request) Successful in 4m57s
K8s Deploy Test / Run deploy test suite on kind/k8s (pull_request) Successful in 8m22s
K8s Deployment Control Test / Run deployment control suite on kind/k8s (pull_request) Successful in 5m51s
Webapp Test / Run webapp test suite (pull_request) Successful in 5m6s
Smoke Test / Run basic test suite (pull_request) Successful in 4m36s
2024-10-28 11:00:44 +05:30
4 changed files with 65 additions and 212 deletions
@@ -54,7 +54,6 @@ def process_app_deployment_request(
deployment_record_namespace,
dns_record_namespace,
default_dns_suffix,
dns_value,
deployment_parent_dir,
kube_config,
image_registry,
@@ -252,7 +251,6 @@ def process_app_deployment_request(
dns_record,
dns_lrn,
deployment_dir,
dns_value,
app_deployment_request,
webapp_deployer_record,
logger,
@@ -306,7 +304,6 @@ def dump_known_requests(filename, requests, status="SEEN"):
help="How to handle requests with an FQDN: prohibit, allow, preexisting",
default="prohibit",
)
@click.option("--ip", help="IP address of the k8s deployment (to be set in DNS record)", default=None)
@click.option("--record-namespace-dns", help="eg, lrn://laconic/dns", required=True)
@click.option(
"--record-namespace-deployments",
@@ -342,17 +339,6 @@ def dump_known_requests(filename, requests, status="SEEN"):
help="Requests must have a minimum payment to be processed (in alnt)",
default=0,
)
@click.option(
"--atom-payment-address",
help="Cosmos ATOM address to receive payments",
default=None,
)
@click.option(
"--min-atom-payment",
help="Minimum required ATOM payment amount",
default=1,
type=float,
)
@click.option("--lrn", help="The LRN of this deployer.", required=True)
@click.option(
"--all-requests",
@@ -395,7 +381,6 @@ def command( # noqa: C901
only_update_state,
dns_suffix,
fqdn_policy,
ip,
record_namespace_dns,
record_namespace_deployments,
dry_run,
@@ -405,8 +390,6 @@ def command( # noqa: C901
recreate_on_deploy,
log_dir,
min_required_payment,
atom_payment_address,
min_atom_payment,
lrn,
config_upload_dir,
private_key_file,
@@ -446,13 +429,6 @@ def command( # noqa: C901
)
sys.exit(2)
if fqdn_policy == "allow" and not ip:
print(
"--ip is required with 'allow' fqdn-policy",
file=sys.stderr,
)
sys.exit(2)
tempdir = tempfile.mkdtemp()
gpg = gnupg.GPG(gnupghome=tempdir)
@@ -643,8 +619,6 @@ def command( # noqa: C901
payment_address,
min_required_payment,
main_logger,
atom_payment_address,
min_atom_payment,
):
main_logger.log(f"{r.id}: Payment confirmed.")
requests_to_execute.append(r)
@@ -691,7 +665,6 @@ def command( # noqa: C901
record_namespace_deployments,
record_namespace_dns,
dns_suffix,
ip,
os.path.abspath(deployment_parent_dir),
kube_config,
image_registry,
@@ -46,17 +46,6 @@ from stack_orchestrator.deploy.webapp.util import LaconicRegistryClient
help="List the minimum required payment (in alnt) to process a deployment request.",
default=0,
)
@click.option(
"--atom-payment-address",
help="The Cosmos ATOM address to which payments should be made.",
default=None,
)
@click.option(
"--min-atom-payment",
help="List the minimum required payment (in uatom) to process a deployment request.",
default="1000000uatom",
type=str,
)
@click.option(
"--dry-run",
help="Don't publish anything, just report what would be done.",
@@ -71,8 +60,6 @@ def command( # noqa: C901
lrn,
payment_address,
min_required_payment,
atom_payment_address,
min_atom_payment,
dry_run,
):
laconic = LaconicRegistryClient(laconic_config)
@@ -96,10 +83,6 @@ def command( # noqa: C901
webapp_deployer_record["record"][
"minimumPayment"
] = f"{min_required_payment}alnt"
if atom_payment_address:
webapp_deployer_record["record"]["atomPaymentAddress"] = atom_payment_address
webapp_deployer_record["record"]["minimumAtomPayment"] = min_atom_payment
if dry_run:
yaml.dump(webapp_deployer_record, sys.stdout)
@@ -1,59 +1,8 @@
import fcntl
from functools import wraps
import os
import time
# Define default file path for the lock
DEFAULT_LOCK_FILE_PATH = "/tmp/registry_mutex_lock_file"
LOCK_TIMEOUT = 30
LOCK_RETRY_INTERVAL = 3
def acquire_lock(client, lock_file_path, timeout):
# Lock alreay acquired by the current client
if client.mutex_lock_acquired:
return
while True:
try:
# Check if lock file exists and is potentially stale
if os.path.exists(lock_file_path):
with open(lock_file_path, 'r') as lock_file:
timestamp = float(lock_file.read().strip())
# If lock is stale, remove the lock file
if time.time() - timestamp > timeout:
print(f"Stale lock detected, removing lock file {lock_file_path}")
os.remove(lock_file_path)
else:
print(f"Lock file {lock_file_path} exists and is recent, waiting...")
time.sleep(LOCK_RETRY_INTERVAL)
continue
# Try to create a new lock file with the current timestamp
fd = os.open(lock_file_path, os.O_CREAT | os.O_EXCL | os.O_RDWR)
with os.fdopen(fd, 'w') as lock_file:
lock_file.write(str(time.time()))
client.mutex_lock_acquired = True
print(f"Registry lock acquired, {lock_file_path}")
# Lock successfully acquired
return
except FileExistsError:
print(f"Lock file {lock_file_path} exists, waiting...")
time.sleep(LOCK_RETRY_INTERVAL)
def release_lock(client, lock_file_path):
try:
os.remove(lock_file_path)
client.mutex_lock_acquired = False
print(f"Registry lock released, {lock_file_path}")
except FileNotFoundError:
# Lock file already removed
pass
def registry_mutex():
@@ -64,13 +13,18 @@ def registry_mutex():
if self.mutex_lock_file:
lock_file_path = self.mutex_lock_file
# Acquire the lock before running the function
acquire_lock(self, lock_file_path, LOCK_TIMEOUT)
try:
return func(self, *args, **kwargs)
finally:
# Release the lock after the function completes
release_lock(self, lock_file_path)
with open(lock_file_path, 'w') as lock_file:
try:
# Try to acquire the lock
fcntl.flock(lock_file, fcntl.LOCK_EX)
# Call the actual function
result = func(self, *args, **kwargs)
finally:
# Always release the lock
fcntl.flock(lock_file, fcntl.LOCK_UN)
return result
return wrapper
+52 -109
View File
@@ -117,6 +117,7 @@ class LaconicRegistryClient:
def __init__(self, config_file, log_file=None, mutex_lock_file=None):
self.config_file = config_file
self.log_file = log_file
self.mutex_lock_file = mutex_lock_file
self.cache = AttrDict(
{
"name_or_id": {},
@@ -125,9 +126,6 @@ class LaconicRegistryClient:
}
)
self.mutex_lock_file = mutex_lock_file
self.mutex_lock_acquired = False
def whoami(self, refresh=False):
if not refresh and "whoami" in self.cache:
return self.cache["whoami"]
@@ -689,7 +687,6 @@ def publish_deployment(
dns_record,
dns_lrn,
deployment_dir,
dns_value=None,
app_deployment_request=None,
webapp_deployer_record=None,
logger=None,
@@ -722,8 +719,6 @@ def publish_deployment(
}
if app_deployment_request:
new_dns_record["record"]["request"] = app_deployment_request.id
if dns_value:
new_dns_record["record"]["value"] = dns_value
if logger:
logger.log("Publishing DnsRecord.")
@@ -801,7 +796,7 @@ def skip_by_tag(r, include_tags, exclude_tags):
return False
def confirm_payment(laconic: LaconicRegistryClient, record, payment_address, min_amount, logger, atom_payment_address=None, atom_min_amount=None):
def confirm_payment(laconic: LaconicRegistryClient, record, payment_address, min_amount, logger):
req_owner = laconic.get_owner(record)
if req_owner == payment_address:
# No need to confirm payment if the sender and recipient are the same account.
@@ -811,114 +806,62 @@ def confirm_payment(laconic: LaconicRegistryClient, record, payment_address, min
logger.log(f"{record.id}: no payment tx info")
return False
# Try to verify as a laconic payment first
tx = laconic.get_tx(record.attributes.payment)
if tx:
if tx.code != 0:
logger.log(
f"{record.id}: payment tx {tx.hash} was not successful - code: {tx.code}, log: {tx.log}"
)
return False
if not tx:
logger.log(f"{record.id}: cannot locate payment tx")
return False
if tx.sender != req_owner:
logger.log(
f"{record.id}: payment sender {tx.sender} in tx {tx.hash} does not match deployment "
f"request owner {req_owner}"
)
return False
if tx.recipient != payment_address:
logger.log(
f"{record.id}: payment recipient {tx.recipient} in tx {tx.hash} does not match {payment_address}"
)
return False
pay_denom = "".join([i for i in tx.amount if not i.isdigit()])
if pay_denom != "alnt":
logger.log(
f"{record.id}: {pay_denom} in tx {tx.hash} is not an expected payment denomination"
)
return False
pay_amount = int("".join([i for i in tx.amount if i.isdigit()]))
if pay_amount < min_amount:
logger.log(
f"{record.id}: payment amount {tx.amount} is less than minimum {min_amount}"
)
return False
# Check if the payment was already used on a deployment
used = laconic.app_deployments(
{"deployer": record.attributes.deployer, "payment": tx.hash}, all=True
if tx.code != 0:
logger.log(
f"{record.id}: payment tx {tx.hash} was not successful - code: {tx.code}, log: {tx.log}"
)
if len(used):
# Fetch the app name from request record
used_request = laconic.get_record(used[0].attributes.request, require=True)
return False
# Check that payment was used for deployment of same application
if record.attributes.application != used_request.attributes.application:
logger.log(f"{record.id}: payment {tx.hash} already used on a different application deployment {used}")
return False
used = laconic.app_deployment_removals(
{"deployer": record.attributes.deployer, "payment": tx.hash}, all=True
if tx.sender != req_owner:
logger.log(
f"{record.id}: payment sender {tx.sender} in tx {tx.hash} does not match deployment "
f"request owner {req_owner}"
)
if len(used):
logger.log(
f"{record.id}: payment {tx.hash} already used on deployment removal {used}"
)
return False
return False
return True
# If we get here, the transaction hash wasn't found in the laconic testnet
# Let's check if it's a valid Cosmos ATOM payment if configuration is available
if atom_payment_address:
logger.log(f"{record.id}: checking if payment is a valid Cosmos ATOM transaction")
try:
import requests
# Use the webapp-deployment-status-api to verify the ATOM payment
deployer_record = laconic.get_record(record.attributes.deployer)
if not deployer_record or not deployer_record.attributes.apiUrl:
logger.log(f"{record.id}: cannot find deployer API URL to verify ATOM payment")
return False
api_url = deployer_record.attributes.apiUrl
verify_url = f"{api_url}/verify/atom-payment"
# Make a request to the API to verify the ATOM payment
# Pass markAsUsed=true to prevent this transaction from being used again
response = requests.post(
verify_url,
json={
"txHash": record.attributes.payment,
"minAmount": atom_min_amount,
"markAsUsed": True
},
timeout=10
)
if response.status_code != 200:
logger.log(f"{record.id}: ATOM payment verification API request failed with status {response.status_code}")
return False
result = response.json()
if not result.get("valid", False):
logger.log(f"{record.id}: ATOM payment verification failed: {result.get('reason', 'unknown reason')}")
return False
# Payment is valid
logger.log(f"{record.id}: ATOM payment verified successfully, amount: {result.get('amount')} ATOM")
return True
except Exception as e:
logger.log(f"{record.id}: error verifying ATOM payment: {str(e)}")
return False
logger.log(f"{record.id}: payment tx {record.attributes.payment} not found in laconic testnet and ATOM payment verification not configured")
return False
if tx.recipient != payment_address:
logger.log(
f"{record.id}: payment recipient {tx.recipient} in tx {tx.hash} does not match {payment_address}"
)
return False
pay_denom = "".join([i for i in tx.amount if not i.isdigit()])
if pay_denom != "alnt":
logger.log(
f"{record.id}: {pay_denom} in tx {tx.hash} is not an expected payment denomination"
)
return False
pay_amount = int("".join([i for i in tx.amount if i.isdigit()]))
if pay_amount < min_amount:
logger.log(
f"{record.id}: payment amount {tx.amount} is less than minimum {min_amount}"
)
return False
# Check if the payment was already used on a
used = laconic.app_deployments(
{"deployer": payment_address, "payment": tx.hash}, all=True
)
if len(used):
logger.log(f"{record.id}: payment {tx.hash} already used on deployment {used}")
return False
used = laconic.app_deployment_removals(
{"deployer": payment_address, "payment": tx.hash}, all=True
)
if len(used):
logger.log(
f"{record.id}: payment {tx.hash} already used on deployment removal {used}"
)
return False
return True
def confirm_auction(laconic: LaconicRegistryClient, record, deployer_lrn, payment_address, logger):