forked from cerc-io/stack-orchestrator
Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
87251ba65b | ||
|
|
873a6d472c | ||
|
|
39df4683ac | ||
|
|
23ca4c4341 | ||
|
|
f64ef5d128 |
@@ -54,6 +54,7 @@ def process_app_deployment_request(
|
|||||||
deployment_record_namespace,
|
deployment_record_namespace,
|
||||||
dns_record_namespace,
|
dns_record_namespace,
|
||||||
default_dns_suffix,
|
default_dns_suffix,
|
||||||
|
dns_value,
|
||||||
deployment_parent_dir,
|
deployment_parent_dir,
|
||||||
kube_config,
|
kube_config,
|
||||||
image_registry,
|
image_registry,
|
||||||
@@ -251,6 +252,7 @@ def process_app_deployment_request(
|
|||||||
dns_record,
|
dns_record,
|
||||||
dns_lrn,
|
dns_lrn,
|
||||||
deployment_dir,
|
deployment_dir,
|
||||||
|
dns_value,
|
||||||
app_deployment_request,
|
app_deployment_request,
|
||||||
webapp_deployer_record,
|
webapp_deployer_record,
|
||||||
logger,
|
logger,
|
||||||
@@ -304,6 +306,7 @@ def dump_known_requests(filename, requests, status="SEEN"):
|
|||||||
help="How to handle requests with an FQDN: prohibit, allow, preexisting",
|
help="How to handle requests with an FQDN: prohibit, allow, preexisting",
|
||||||
default="prohibit",
|
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-dns", help="eg, lrn://laconic/dns", required=True)
|
||||||
@click.option(
|
@click.option(
|
||||||
"--record-namespace-deployments",
|
"--record-namespace-deployments",
|
||||||
@@ -339,6 +342,17 @@ def dump_known_requests(filename, requests, status="SEEN"):
|
|||||||
help="Requests must have a minimum payment to be processed (in alnt)",
|
help="Requests must have a minimum payment to be processed (in alnt)",
|
||||||
default=0,
|
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("--lrn", help="The LRN of this deployer.", required=True)
|
||||||
@click.option(
|
@click.option(
|
||||||
"--all-requests",
|
"--all-requests",
|
||||||
@@ -381,6 +395,7 @@ def command( # noqa: C901
|
|||||||
only_update_state,
|
only_update_state,
|
||||||
dns_suffix,
|
dns_suffix,
|
||||||
fqdn_policy,
|
fqdn_policy,
|
||||||
|
ip,
|
||||||
record_namespace_dns,
|
record_namespace_dns,
|
||||||
record_namespace_deployments,
|
record_namespace_deployments,
|
||||||
dry_run,
|
dry_run,
|
||||||
@@ -390,6 +405,8 @@ def command( # noqa: C901
|
|||||||
recreate_on_deploy,
|
recreate_on_deploy,
|
||||||
log_dir,
|
log_dir,
|
||||||
min_required_payment,
|
min_required_payment,
|
||||||
|
atom_payment_address,
|
||||||
|
min_atom_payment,
|
||||||
lrn,
|
lrn,
|
||||||
config_upload_dir,
|
config_upload_dir,
|
||||||
private_key_file,
|
private_key_file,
|
||||||
@@ -429,6 +446,13 @@ def command( # noqa: C901
|
|||||||
)
|
)
|
||||||
sys.exit(2)
|
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()
|
tempdir = tempfile.mkdtemp()
|
||||||
gpg = gnupg.GPG(gnupghome=tempdir)
|
gpg = gnupg.GPG(gnupghome=tempdir)
|
||||||
|
|
||||||
@@ -619,6 +643,8 @@ def command( # noqa: C901
|
|||||||
payment_address,
|
payment_address,
|
||||||
min_required_payment,
|
min_required_payment,
|
||||||
main_logger,
|
main_logger,
|
||||||
|
atom_payment_address,
|
||||||
|
min_atom_payment,
|
||||||
):
|
):
|
||||||
main_logger.log(f"{r.id}: Payment confirmed.")
|
main_logger.log(f"{r.id}: Payment confirmed.")
|
||||||
requests_to_execute.append(r)
|
requests_to_execute.append(r)
|
||||||
@@ -665,6 +691,7 @@ def command( # noqa: C901
|
|||||||
record_namespace_deployments,
|
record_namespace_deployments,
|
||||||
record_namespace_dns,
|
record_namespace_dns,
|
||||||
dns_suffix,
|
dns_suffix,
|
||||||
|
ip,
|
||||||
os.path.abspath(deployment_parent_dir),
|
os.path.abspath(deployment_parent_dir),
|
||||||
kube_config,
|
kube_config,
|
||||||
image_registry,
|
image_registry,
|
||||||
|
|||||||
@@ -46,6 +46,17 @@ from stack_orchestrator.deploy.webapp.util import LaconicRegistryClient
|
|||||||
help="List the minimum required payment (in alnt) to process a deployment request.",
|
help="List the minimum required payment (in alnt) to process a deployment request.",
|
||||||
default=0,
|
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 ATOM) to process a deployment request.",
|
||||||
|
default=1,
|
||||||
|
type=float,
|
||||||
|
)
|
||||||
@click.option(
|
@click.option(
|
||||||
"--dry-run",
|
"--dry-run",
|
||||||
help="Don't publish anything, just report what would be done.",
|
help="Don't publish anything, just report what would be done.",
|
||||||
@@ -60,6 +71,8 @@ def command( # noqa: C901
|
|||||||
lrn,
|
lrn,
|
||||||
payment_address,
|
payment_address,
|
||||||
min_required_payment,
|
min_required_payment,
|
||||||
|
atom_payment_address,
|
||||||
|
min_atom_payment,
|
||||||
dry_run,
|
dry_run,
|
||||||
):
|
):
|
||||||
laconic = LaconicRegistryClient(laconic_config)
|
laconic = LaconicRegistryClient(laconic_config)
|
||||||
@@ -84,6 +97,10 @@ def command( # noqa: C901
|
|||||||
"minimumPayment"
|
"minimumPayment"
|
||||||
] = f"{min_required_payment}alnt"
|
] = 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:
|
if dry_run:
|
||||||
yaml.dump(webapp_deployer_record, sys.stdout)
|
yaml.dump(webapp_deployer_record, sys.stdout)
|
||||||
return
|
return
|
||||||
|
|||||||
@@ -1,8 +1,59 @@
|
|||||||
import fcntl
|
|
||||||
from functools import wraps
|
from functools import wraps
|
||||||
|
import os
|
||||||
|
import time
|
||||||
|
|
||||||
# Define default file path for the lock
|
# Define default file path for the lock
|
||||||
DEFAULT_LOCK_FILE_PATH = "/tmp/registry_mutex_lock_file"
|
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():
|
def registry_mutex():
|
||||||
@@ -13,18 +64,13 @@ def registry_mutex():
|
|||||||
if self.mutex_lock_file:
|
if self.mutex_lock_file:
|
||||||
lock_file_path = self.mutex_lock_file
|
lock_file_path = self.mutex_lock_file
|
||||||
|
|
||||||
with open(lock_file_path, 'w') as lock_file:
|
# Acquire the lock before running the function
|
||||||
|
acquire_lock(self, lock_file_path, LOCK_TIMEOUT)
|
||||||
try:
|
try:
|
||||||
# Try to acquire the lock
|
return func(self, *args, **kwargs)
|
||||||
fcntl.flock(lock_file, fcntl.LOCK_EX)
|
|
||||||
|
|
||||||
# Call the actual function
|
|
||||||
result = func(self, *args, **kwargs)
|
|
||||||
finally:
|
finally:
|
||||||
# Always release the lock
|
# Release the lock after the function completes
|
||||||
fcntl.flock(lock_file, fcntl.LOCK_UN)
|
release_lock(self, lock_file_path)
|
||||||
|
|
||||||
return result
|
|
||||||
|
|
||||||
return wrapper
|
return wrapper
|
||||||
|
|
||||||
|
|||||||
@@ -117,7 +117,6 @@ class LaconicRegistryClient:
|
|||||||
def __init__(self, config_file, log_file=None, mutex_lock_file=None):
|
def __init__(self, config_file, log_file=None, mutex_lock_file=None):
|
||||||
self.config_file = config_file
|
self.config_file = config_file
|
||||||
self.log_file = log_file
|
self.log_file = log_file
|
||||||
self.mutex_lock_file = mutex_lock_file
|
|
||||||
self.cache = AttrDict(
|
self.cache = AttrDict(
|
||||||
{
|
{
|
||||||
"name_or_id": {},
|
"name_or_id": {},
|
||||||
@@ -126,6 +125,9 @@ class LaconicRegistryClient:
|
|||||||
}
|
}
|
||||||
)
|
)
|
||||||
|
|
||||||
|
self.mutex_lock_file = mutex_lock_file
|
||||||
|
self.mutex_lock_acquired = False
|
||||||
|
|
||||||
def whoami(self, refresh=False):
|
def whoami(self, refresh=False):
|
||||||
if not refresh and "whoami" in self.cache:
|
if not refresh and "whoami" in self.cache:
|
||||||
return self.cache["whoami"]
|
return self.cache["whoami"]
|
||||||
@@ -687,6 +689,7 @@ def publish_deployment(
|
|||||||
dns_record,
|
dns_record,
|
||||||
dns_lrn,
|
dns_lrn,
|
||||||
deployment_dir,
|
deployment_dir,
|
||||||
|
dns_value=None,
|
||||||
app_deployment_request=None,
|
app_deployment_request=None,
|
||||||
webapp_deployer_record=None,
|
webapp_deployer_record=None,
|
||||||
logger=None,
|
logger=None,
|
||||||
@@ -719,6 +722,8 @@ def publish_deployment(
|
|||||||
}
|
}
|
||||||
if app_deployment_request:
|
if app_deployment_request:
|
||||||
new_dns_record["record"]["request"] = app_deployment_request.id
|
new_dns_record["record"]["request"] = app_deployment_request.id
|
||||||
|
if dns_value:
|
||||||
|
new_dns_record["record"]["value"] = dns_value
|
||||||
|
|
||||||
if logger:
|
if logger:
|
||||||
logger.log("Publishing DnsRecord.")
|
logger.log("Publishing DnsRecord.")
|
||||||
@@ -796,7 +801,7 @@ def skip_by_tag(r, include_tags, exclude_tags):
|
|||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
||||||
def confirm_payment(laconic: LaconicRegistryClient, record, payment_address, min_amount, logger):
|
def confirm_payment(laconic: LaconicRegistryClient, record, payment_address, min_amount, logger, atom_payment_address=None, atom_min_amount=None):
|
||||||
req_owner = laconic.get_owner(record)
|
req_owner = laconic.get_owner(record)
|
||||||
if req_owner == payment_address:
|
if req_owner == payment_address:
|
||||||
# No need to confirm payment if the sender and recipient are the same account.
|
# No need to confirm payment if the sender and recipient are the same account.
|
||||||
@@ -806,11 +811,9 @@ def confirm_payment(laconic: LaconicRegistryClient, record, payment_address, min
|
|||||||
logger.log(f"{record.id}: no payment tx info")
|
logger.log(f"{record.id}: no payment tx info")
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
# Try to verify as a laconic payment first
|
||||||
tx = laconic.get_tx(record.attributes.payment)
|
tx = laconic.get_tx(record.attributes.payment)
|
||||||
if not tx:
|
if tx:
|
||||||
logger.log(f"{record.id}: cannot locate payment tx")
|
|
||||||
return False
|
|
||||||
|
|
||||||
if tx.code != 0:
|
if tx.code != 0:
|
||||||
logger.log(
|
logger.log(
|
||||||
f"{record.id}: payment tx {tx.hash} was not successful - code: {tx.code}, log: {tx.log}"
|
f"{record.id}: payment tx {tx.hash} was not successful - code: {tx.code}, log: {tx.log}"
|
||||||
@@ -844,16 +847,21 @@ def confirm_payment(laconic: LaconicRegistryClient, record, payment_address, min
|
|||||||
)
|
)
|
||||||
return False
|
return False
|
||||||
|
|
||||||
# Check if the payment was already used on a
|
# Check if the payment was already used on a deployment
|
||||||
used = laconic.app_deployments(
|
used = laconic.app_deployments(
|
||||||
{"deployer": payment_address, "payment": tx.hash}, all=True
|
{"deployer": record.attributes.deployer, "payment": tx.hash}, all=True
|
||||||
)
|
)
|
||||||
if len(used):
|
if len(used):
|
||||||
logger.log(f"{record.id}: payment {tx.hash} already used on deployment {used}")
|
# Fetch the app name from request record
|
||||||
|
used_request = laconic.get_record(used[0].attributes.request, require=True)
|
||||||
|
|
||||||
|
# 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
|
return False
|
||||||
|
|
||||||
used = laconic.app_deployment_removals(
|
used = laconic.app_deployment_removals(
|
||||||
{"deployer": payment_address, "payment": tx.hash}, all=True
|
{"deployer": record.attributes.deployer, "payment": tx.hash}, all=True
|
||||||
)
|
)
|
||||||
if len(used):
|
if len(used):
|
||||||
logger.log(
|
logger.log(
|
||||||
@@ -863,6 +871,55 @@ def confirm_payment(laconic: LaconicRegistryClient, record, payment_address, min
|
|||||||
|
|
||||||
return True
|
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
|
||||||
|
|
||||||
|
|
||||||
def confirm_auction(laconic: LaconicRegistryClient, record, deployer_lrn, payment_address, logger):
|
def confirm_auction(laconic: LaconicRegistryClient, record, deployer_lrn, payment_address, logger):
|
||||||
auction_id = record.attributes.auction
|
auction_id = record.attributes.auction
|
||||||
|
|||||||
Reference in New Issue
Block a user