Compare commits

..
Author SHA1 Message Date
prathamesh e1e7075bf3 Fix arg 2025-02-10 14:10:16 +05:30
prathamesh e491545354 Use existing certificates if available and update status command 2025-02-10 11:44:45 +05:30
prathamesh 9bc8ce4866 Update ingress creation for multiple host names 2025-02-10 10:56:00 +05:30
prathamesh f98e75c81f Fix deployment record publication 2025-02-07 11:30:07 +05:30
prathamesh 8c77af0fec Fix types usage 2025-02-06 20:30:45 +05:30
prathamesh 4ca96bbf8c Fix types usage 2025-02-06 19:13:58 +05:30
prathamesh 831e4a14ff Handle lint errors 2025-02-06 16:56:21 +05:30
prathamesh af5ec54b8d Add deployer version to deployer record 2025-02-06 16:13:43 +05:30
prathamesh 50508644be Delete deployment and DNS names for removed domains 2025-02-06 16:12:23 +05:30
prathamesh 42103ef551 Publish multiple DNS records and set all deployment names 2025-02-06 15:14:12 +05:30
prathamesh ef6f5db743 Update image tag for subsequent deployments of same app 2025-02-05 19:18:49 +05:30
prathamesh cdd431d8a5 Update existing deployment spec with new urls 2025-02-05 19:09:23 +05:30
prathamesh 3d1a455344 Configure spec with all given URLs 2025-02-05 18:48:10 +05:30
prathamesh f5947a8a65 Check for existing deployment 2025-02-05 18:39:18 +05:30
prathamesh a8e7235267 Handle multiple domains in webapp deployment request 2025-02-05 17:57:18 +05:30
prathamesh 873a6d472c Update webapp deployment flow for supporting custom domains (#963)
Part of https://www.notion.so/Support-custom-domains-in-deploy-laconic-com-18aa6b22d4728067a44ae27090c02ce5 and cerc-io/snowballtools-base#47

- Set `value` (IP address for `A` resource) in the DNS records
- Update `deploy-webapp-from-registry` command with an option to pass k8s cluster IP address (only required with fqdn policy `allow`)

Reviewed-on: cerc-io/stack-orchestrator#963
Reviewed-by: ashwin <ashwin@noreply.git.vdb.to>
Co-authored-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
Co-committed-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
2025-02-04 13:26:31 +00:00
prathamesh 39df4683ac Allow payment reuse for same app LRN (#961)
Part of [Service provider auctions for web deployments](https://www.notion.so/Service-provider-auctions-for-web-deployments-104a6b22d47280dbad51d28aa3a91d75)

Reviewed-on: cerc-io/stack-orchestrator#961
Reviewed-by: ashwin <ashwin@noreply.git.vdb.to>
Co-authored-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
Co-committed-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
2024-10-29 11:30:03 +00:00
prathamesh 23ca4c4341 Allow payment reuse for application redeployment (#960)
Part of [Service provider auctions for web deployments](https://www.notion.so/Service-provider-auctions-for-web-deployments-104a6b22d47280dbad51d28aa3a91d75)

Reviewed-on: cerc-io/stack-orchestrator#960
Reviewed-by: ashwin <ashwin@noreply.git.vdb.to>
Co-authored-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
Co-committed-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
2024-10-29 06:51:48 +00:00
prathamesh f64ef5d128 Use file existence for registry mutex (#959)
Part of [Service provider auctions for web deployments](https://www.notion.so/Service-provider-auctions-for-web-deployments-104a6b22d47280dbad51d28aa3a91d75)

Reviewed-on: cerc-io/stack-orchestrator#959
Reviewed-by: ashwin <ashwin@noreply.git.vdb.to>
Co-authored-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
Co-committed-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
2024-10-29 04:05:35 +00:00
prathamesh 5f8e809b2d Add mutex lock file path to registry CLI wrapper class (#958)
Part of [Service provider auctions for web deployments](https://www.notion.so/Service-provider-auctions-for-web-deployments-104a6b22d47280dbad51d28aa3a91d75)
Follows cerc-io/stack-orchestrator#957

Reviewed-on: cerc-io/stack-orchestrator#958
Reviewed-by: ashwin <ashwin@noreply.git.vdb.to>
Co-authored-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
Co-committed-by: Prathamesh Musale <prathamesh.musale0@gmail.com>
2024-10-28 06:03:13 +00:00
7 changed files with 402 additions and 200 deletions
+31 -26
View File
@@ -114,22 +114,27 @@ class ClusterInfo:
nodeports.append(service) nodeports.append(service)
return nodeports return nodeports
def get_ingress(self, use_tls=False, certificate=None, cluster_issuer="letsencrypt-prod"): def get_ingress(self, use_tls=False, certificate_by_host={}, cluster_issuer="letsencrypt-prod"):
# No ingress for a deployment that has no http-proxy defined, for now # No ingress for a deployment that has no http-proxy defined, for now
http_proxy_info_list = self.spec.get_http_proxy() http_proxy_info_list = self.spec.get_http_proxy()
ingress = None if not http_proxy_info_list:
if http_proxy_info_list: return None
# TODO: handle multiple definitions
http_proxy_info = http_proxy_info_list[0] tls = [] if use_tls else None
rules = []
for http_proxy_info in http_proxy_info_list:
if opts.o.debug: if opts.o.debug:
print(f"http-proxy: {http_proxy_info}") print(f"http-proxy: {http_proxy_info}")
# TODO: good enough parsing for webapp deployment for now # TODO: good enough parsing for webapp deployment for now
host_name = http_proxy_info["host-name"] host_name = http_proxy_info["host-name"]
rules = [] certificate = certificate_by_host[host_name] if host_name in certificate_by_host else None
tls = [client.V1IngressTLS(
hosts=certificate["spec"]["dnsNames"] if certificate else [host_name], if use_tls:
secret_name=certificate["spec"]["secretName"] if certificate else f"{self.app_name}-tls" tls.append(client.V1IngressTLS(
)] if use_tls else None hosts=certificate["spec"]["dnsNames"] if certificate else [host_name],
secret_name=certificate["spec"]["secretName"] if certificate else f"{self.app_name}-{host_name}-tls"
))
paths = [] paths = []
for route in http_proxy_info["routes"]: for route in http_proxy_info["routes"]:
path = route["path"] path = route["path"]
@@ -156,24 +161,24 @@ class ClusterInfo:
paths=paths paths=paths
) )
)) ))
spec = client.V1IngressSpec( spec = client.V1IngressSpec(
tls=tls, tls=tls,
rules=rules rules=rules
) )
ingress_annotations = { ingress_annotations = {
"kubernetes.io/ingress.class": "nginx", "kubernetes.io/ingress.class": "nginx",
} }
if not certificate: if not certificate:
ingress_annotations["cert-manager.io/cluster-issuer"] = cluster_issuer ingress_annotations["cert-manager.io/cluster-issuer"] = cluster_issuer
ingress = client.V1Ingress( ingress = client.V1Ingress(
metadata=client.V1ObjectMeta( metadata=client.V1ObjectMeta(
name=f"{self.app_name}-ingress", name=f"{self.app_name}-ingress",
annotations=ingress_annotations annotations=ingress_annotations
), ),
spec=spec spec=spec
) )
return ingress return ingress
# TODO: suppoprt multiple services # TODO: suppoprt multiple services
+40 -27
View File
@@ -227,15 +227,18 @@ class K8sDeployer(Deployer):
self._create_volume_data() self._create_volume_data()
self._create_deployment() self._create_deployment()
http_proxy_info = self.cluster_info.spec.get_http_proxy() http_proxy_info_list = self.cluster_info.spec.get_http_proxy()
# Note: at present we don't support tls for kind (and enabling tls causes errors) # Note: at present we don't support tls for kind (and enabling tls causes errors)
use_tls = http_proxy_info and not self.is_kind() use_tls = http_proxy_info_list and not self.is_kind()
certificate = self._find_certificate_for_host_name(http_proxy_info[0]["host-name"]) if use_tls else None certificate_by_host = {}
if opts.o.debug: if use_tls:
if certificate: for http_proxy_info in http_proxy_info_list:
print(f"Using existing certificate: {certificate}") certificate = self._find_certificate_for_host_name(http_proxy_info["host-name"])
if opts.o.debug and certificate:
print(f"Using existing certificate: {certificate}")
certificate_by_host[http_proxy_info["host-name"]] = certificate
ingress: client.V1Ingress = self.cluster_info.get_ingress(use_tls=use_tls, certificate=certificate) ingress: client.V1Ingress = self.cluster_info.get_ingress(use_tls=use_tls, certificate_by_host=certificate_by_host)
if ingress: if ingress:
if opts.o.debug: if opts.o.debug:
print(f"Sending this ingress: {ingress}") print(f"Sending this ingress: {ingress}")
@@ -381,36 +384,46 @@ class K8sDeployer(Deployer):
if not pods: if not pods:
return return
hostname = "?" tls_by_host = {}
ip = "?"
tls = "?"
try: try:
ingress = self.networking_api.read_namespaced_ingress(namespace=self.k8s_namespace, ingress = self.networking_api.read_namespaced_ingress(namespace=self.k8s_namespace,
name=self.cluster_info.get_ingress().metadata.name) name=self.cluster_info.get_ingress().metadata.name)
cert = self.custom_obj_api.get_namespaced_custom_object(
group="cert-manager.io",
version="v1",
namespace=self.k8s_namespace,
plural="certificates",
name=ingress.spec.tls[0].secret_name
)
hostname = ingress.spec.rules[0].host
ip = ingress.status.load_balancer.ingress[0].ip ip = ingress.status.load_balancer.ingress[0].ip
tls = "notBefore: %s; notAfter: %s; names: %s" % ( for rule in ingress.spec.rules:
cert["status"]["notBefore"], cert["status"]["notAfter"], ingress.spec.tls[0].hosts hostname = rule.host
) tls_spec = next((tls for tls in ingress.spec.tls if hostname in tls.hosts), None)
if tls_spec:
cert = self.custom_obj_api.get_namespaced_custom_object(
group="cert-manager.io",
version="v1",
namespace=self.k8s_namespace,
plural="certificates",
name=tls_spec.secret_name
)
tls = "notBefore: %s; notAfter: %s; names: %s" % (
cert["status"]["notBefore"], cert["status"]["notAfter"], tls_spec.hosts
)
tls_by_host[hostname] = tls
else:
tls_by_host[hostname] = None
except: # noqa: E722 except: # noqa: E722
pass pass
print("Ingress:") print("Ingress:")
print("\tHostname:", hostname) if len(tls_by_host) == 0:
print("\tIP:", ip) print("\tHostname:", "?")
print("\tTLS:", tls) print("\tIP:", "?")
print("") print("\tTLS:", "?")
print("Pods:") print("")
for hostname, tls in tls_by_host.items():
print("\tHostname:", hostname)
print("\tIP:", ip)
print("\tTLS:", tls)
print("")
print("Pods:")
for p in pods: for p in pods:
if p.metadata.deletion_timestamp: if p.metadata.deletion_timestamp:
print(f"\t{p.metadata.namespace}/{p.metadata.name}: Terminating ({p.metadata.deletion_timestamp})") print(f"\t{p.metadata.namespace}/{p.metadata.name}: Terminating ({p.metadata.deletion_timestamp})")
@@ -13,12 +13,16 @@
# You should have received a copy of the GNU Affero General Public License # You should have received a copy of the GNU Affero General Public License
# along with this program. If not, see <http:#www.gnu.org/licenses/>. # along with this program. If not, see <http:#www.gnu.org/licenses/>.
import shutil
from typing import List
import click import click
import os import os
import yaml
from pathlib import Path from pathlib import Path
from urllib.parse import urlparse from urllib.parse import urlparse
from tempfile import NamedTemporaryFile from tempfile import NamedTemporaryFile
from stack_orchestrator import constants
from stack_orchestrator.util import error_exit, global_options2 from stack_orchestrator.util import error_exit, global_options2
from stack_orchestrator.deploy.deployment_create import init_operation, create_operation from stack_orchestrator.deploy.deployment_create import init_operation, create_operation
from stack_orchestrator.deploy.deploy import create_deploy_context from stack_orchestrator.deploy.deploy import create_deploy_context
@@ -36,25 +40,38 @@ def _fixup_container_tag(deployment_dir: str, image: str):
wfile.write(contents) wfile.write(contents)
def _fixup_url_spec(spec_file_name: str, url: str): def _fixup_url_spec(spec_file_name: str, urls: List[str]):
# url is like: https://example.com/path
parsed_url = urlparse(url)
http_proxy_spec = f'''
http-proxy:
- host-name: {parsed_url.hostname}
routes:
- path: '{parsed_url.path if parsed_url.path else "/"}'
proxy-to: webapp:80
'''
spec_file_path = Path(spec_file_name) spec_file_path = Path(spec_file_name)
with open(spec_file_path) as rfile:
contents = rfile.read() # Load existing spec
contents = contents + http_proxy_spec with open(spec_file_path, "r") as file:
with open(spec_file_path, "w") as wfile: spec_data = yaml.safe_load(file) or {}
wfile.write(contents)
# Build new http-proxy entries
http_proxy_entries = []
for url in urls:
parsed_url = urlparse(url)
http_proxy_entries.append({
"host-name": parsed_url.hostname,
"routes": [
{
"path": parsed_url.path if parsed_url.path else "/",
"proxy-to": "webapp:80"
}
]
})
# Update the spec
if "network" not in spec_data:
spec_data["network"] = {}
spec_data["network"]["http-proxy"] = http_proxy_entries
# Write back the updated YAML
with open(spec_file_path, "w") as file:
yaml.dump(spec_data, file, default_flow_style=False, sort_keys=False)
def create_deployment(ctx, deployment_dir, image, url, kube_config, image_registry, env_file): def create_deployment(ctx, deployment_dir, image, urls, kube_config, image_registry, env_file):
# Do the equivalent of: # Do the equivalent of:
# 1. laconic-so --stack webapp-template deploy --deploy-to k8s init --output webapp-spec.yml # 1. laconic-so --stack webapp-template deploy --deploy-to k8s init --output webapp-spec.yml
# --config (eqivalent of the contents of my-config.env) # --config (eqivalent of the contents of my-config.env)
@@ -86,7 +103,7 @@ def create_deployment(ctx, deployment_dir, image, url, kube_config, image_regist
None None
) )
# Add the TLS and DNS spec # Add the TLS and DNS spec
_fixup_url_spec(spec_file_name, url) _fixup_url_spec(spec_file_name, urls)
create_operation( create_operation(
deploy_command_context, deploy_command_context,
spec_file_name, spec_file_name,
@@ -99,6 +116,20 @@ def create_deployment(ctx, deployment_dir, image, url, kube_config, image_regist
os.remove(spec_file_name) os.remove(spec_file_name)
def update_deployment(deployment_dir, image, urls, env_file):
# Update config if required
if env_file:
deployment_config_file = os.path.join(deployment_dir, "config.env")
shutil.copyfile(env_file, deployment_config_file)
# Update existing deployment spec with new urls
_fixup_url_spec(os.path.join(deployment_dir, constants.spec_file_name), urls)
# Update the image name if required
if image:
_fixup_container_tag(deployment_dir, image)
@click.group() @click.group()
@click.pass_context @click.pass_context
def command(ctx): def command(ctx):
@@ -120,4 +151,4 @@ def command(ctx):
def create(ctx, deployment_dir, image, url, kube_config, image_registry, env_file): def create(ctx, deployment_dir, image, url, kube_config, image_registry, env_file):
'''create a deployment for the specified webapp container''' '''create a deployment for the specified webapp container'''
return create_deployment(ctx, deployment_dir, image, url, kube_config, image_registry, env_file) return create_deployment(ctx, deployment_dir, image, [url], kube_config, image_registry, env_file)
@@ -20,6 +20,7 @@ import shutil
import sys import sys
import tempfile import tempfile
import time import time
from typing import List
import uuid import uuid
import yaml import yaml
@@ -38,7 +39,8 @@ from stack_orchestrator.deploy.webapp.util import (
file_hash, file_hash,
deploy_to_k8s, deploy_to_k8s,
publish_deployment, publish_deployment,
hostname_for_deployment_request, get_requested_names,
hostnames_for_deployment_request,
generate_hostname_for_app, generate_hostname_for_app,
match_owner, match_owner,
skip_by_tag, skip_by_tag,
@@ -47,13 +49,14 @@ from stack_orchestrator.deploy.webapp.util import (
) )
def process_app_deployment_request( def process_app_deployment_request( # noqa
ctx, ctx,
laconic: LaconicRegistryClient, laconic: LaconicRegistryClient,
app_deployment_request, 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,
@@ -75,44 +78,53 @@ def process_app_deployment_request(
logger.log(f"Retrieved app record {app_deployment_request.attributes.application}") logger.log(f"Retrieved app record {app_deployment_request.attributes.application}")
# 2. determine dns # 2. determine dns
requested_name = hostname_for_deployment_request(app_deployment_request, laconic) requested_names = hostnames_for_deployment_request(app_deployment_request, laconic)
logger.log(f"Determined requested name: {requested_name}") logger.log(f"Determined requested name(s): {','.join(str(x) for x in requested_names)}")
if "." in requested_name: fqdns = []
if "allow" == fqdn_policy or "preexisting" == fqdn_policy: for requested_name in requested_names:
fqdn = requested_name if "." in requested_name:
if "allow" == fqdn_policy or "preexisting" == fqdn_policy:
fqdns.append(requested_name)
else:
raise Exception(
f"{requested_name} is invalid: only unqualified hostnames are allowed."
)
else: else:
raise Exception( fqdns.append(f"{requested_name}.{default_dns_suffix}")
f"{requested_name} is invalid: only unqualified hostnames are allowed."
)
else:
fqdn = f"{requested_name}.{default_dns_suffix}"
# Normalize case (just in case) # Normalize case (just in case)
fqdn = fqdn.lower() fqdns = [fqdn.lower() for fqdn in fqdns]
# 3. check ownership of existing dnsrecord vs this request # 3. check ownership of existing dnsrecord(s) vs this request
dns_lrn = f"{dns_record_namespace}/{fqdn}" dns_lrns = []
dns_record = laconic.get_record(dns_lrn) existing_dns_records_by_lrns = {}
if dns_record: for fqdn in fqdns:
matched_owner = match_owner(app_deployment_request, dns_record) dns_lrn = f"{dns_record_namespace}/{fqdn}"
if not matched_owner and dns_record.attributes.request: dns_lrns.append(dns_lrn)
matched_owner = match_owner(
app_deployment_request,
laconic.get_record(dns_record.attributes.request, require=True),
)
if matched_owner: dns_record = laconic.get_record(dns_lrn)
logger.log(f"Matched DnsRecord ownership: {matched_owner}") existing_dns_records_by_lrns[dns_lrn] = dns_record
else:
if dns_record:
matched_owner = match_owner(app_deployment_request, dns_record)
if not matched_owner and dns_record.attributes.request:
matched_owner = match_owner(
app_deployment_request,
laconic.get_record(dns_record.attributes.request, require=True),
)
if matched_owner:
logger.log(f"Matched DnsRecord ownership for {fqdn}: {matched_owner}")
else:
raise Exception(
"Unable to confirm ownership of DnsRecord %s for request %s"
% (dns_lrn, app_deployment_request.id)
)
elif "preexisting" == fqdn_policy:
raise Exception( raise Exception(
"Unable to confirm ownership of DnsRecord %s for request %s" f"No pre-existing DnsRecord {dns_lrn} could be found for request {app_deployment_request.id}."
% (dns_lrn, app_deployment_request.id)
) )
elif "preexisting" == fqdn_policy:
raise Exception(
f"No pre-existing DnsRecord {dns_lrn} could be found for request {app_deployment_request.id}."
)
# 4. get build and runtime config from request # 4. get build and runtime config from request
env = {} env = {}
@@ -144,30 +156,49 @@ def process_app_deployment_request(
# 5. determine new or existing deployment # 5. determine new or existing deployment
# a. check for deployment lrn # a. check for deployment lrn
app_deployment_lrn = f"{deployment_record_namespace}/{fqdn}" app_deployment_lrns = [f"{deployment_record_namespace}/{fqdn}" for fqdn in fqdns]
if app_deployment_request.attributes.deployment: if app_deployment_request.attributes.deployment:
app_deployment_lrn = app_deployment_request.attributes.deployment app_deployment_lrns = [app_deployment_request.attributes.deployment]
if not app_deployment_lrn.startswith(deployment_record_namespace): if not app_deployment_lrns[0].startswith(deployment_record_namespace):
raise Exception( raise Exception(
"Deployment LRN %s is not in a supported namespace" "Deployment LRN %s is not in a supported namespace"
% app_deployment_request.attributes.deployment % app_deployment_request.attributes.deployment
) )
deployment_record = laconic.get_record(app_deployment_lrn) # Target deployment dir
deployment_dir = os.path.join(deployment_parent_dir, fqdn) deployment_dir = os.path.join(deployment_parent_dir, fqdns[-1])
# Existing deployment record: take the first lrn that resolves
deployment_record = None
fqdns_to_release = []
existing_deployment_dir = deployment_dir # Default to target dir in case the app had been undeployed
for app_deployment_lrn in app_deployment_lrns:
deployment_record = laconic.get_record(app_deployment_lrn)
if deployment_record:
# Determine the deployment dir for existing deployment
dir_name = deployment_record.attributes.url.replace("https://", "")
existing_deployment_dir = os.path.join(deployment_parent_dir, dir_name)
if not os.path.exists(existing_deployment_dir):
raise Exception(
"Deployment record %s exists, but not deployment dir %s. Please remove name."
% (app_deployment_lrn, existing_deployment_dir)
)
previous_app_deployment_lrns: List[str] = deployment_record.names
previous_fqdns = [lrn.removeprefix(f"{deployment_record_namespace}/") for lrn in previous_app_deployment_lrns]
fqdns_to_release = list(set(previous_fqdns) - set(fqdns))
break
# Use the last fqdn for unique deployment container tag
# At present we use this to generate a unique but stable ID for the app's host container # At present we use this to generate a unique but stable ID for the app's host container
# TODO: implement support to derive this transparently from the already-unique deployment id # TODO: implement support to derive this transparently from the already-unique deployment id
unique_deployment_id = hashlib.md5(fqdn.encode()).hexdigest()[:16] unique_deployment_id = hashlib.md5(fqdns[-1].encode()).hexdigest()[:16]
deployment_config_file = os.path.join(deployment_dir, "config.env")
deployment_container_tag = "laconic-webapp/%s:local" % unique_deployment_id deployment_container_tag = "laconic-webapp/%s:local" % unique_deployment_id
app_image_shared_tag = f"laconic-webapp/{app.id}:local" app_image_shared_tag = f"laconic-webapp/{app.id}:local"
# b. check for deployment directory (create if necessary) # b. check for deployment directory (create if necessary)
if not os.path.exists(deployment_dir): if not os.path.exists(existing_deployment_dir):
if deployment_record:
raise Exception(
"Deployment record %s exists, but not deployment dir %s. Please remove name."
% (app_deployment_lrn, deployment_dir)
)
logger.log( logger.log(
f"Creating webapp deployment in: {deployment_dir} with container id: {deployment_container_tag}" f"Creating webapp deployment in: {deployment_dir} with container id: {deployment_container_tag}"
) )
@@ -175,13 +206,23 @@ def process_app_deployment_request(
ctx, ctx,
deployment_dir, deployment_dir,
deployment_container_tag, deployment_container_tag,
f"https://{fqdn}", [f"https://{fqdn}" for fqdn in fqdns],
kube_config, kube_config,
image_registry, image_registry,
env_filename, env_filename,
) )
elif env_filename: else:
shutil.copyfile(env_filename, deployment_config_file) # Rename deployment dir according to new request (last fqdn from given dns)
os.rename(existing_deployment_dir, deployment_dir)
# Update the image name deployment_container_tag
# Skip for redeployment as deployment_container_tag won't get built
updated_image = None
if deployment_record.attributes.application != app.id:
updated_image = deployment_container_tag
# Update the existing deployment
deploy_webapp.update_deployment(deployment_dir, updated_image, [f"https://{fqdn}" for fqdn in fqdns], env_filename)
needs_k8s_deploy = False needs_k8s_deploy = False
if force_rebuild: if force_rebuild:
@@ -231,10 +272,11 @@ def process_app_deployment_request(
else: else:
logger.log("Requested app is already deployed, skipping build and image push") logger.log("Requested app is already deployed, skipping build and image push")
# 7. update config (if needed) # 7. restart deployment on config or url spec change
if ( if (
not deployment_record not deployment_record
or file_hash(deployment_config_file) != deployment_record.attributes.meta.config or file_hash(os.path.join(deployment_dir, "config.env")) != deployment_record.attributes.meta.config
or len(fqdns_to_release) != 0
): ):
needs_k8s_deploy = True needs_k8s_deploy = True
@@ -247,15 +289,28 @@ def process_app_deployment_request(
laconic, laconic,
app, app,
deployment_record, deployment_record,
app_deployment_lrn, app_deployment_lrns,
dns_record, existing_dns_records_by_lrns,
dns_lrn, dns_lrns,
dns_record_namespace,
deployment_dir, deployment_dir,
dns_value,
app_deployment_request, app_deployment_request,
webapp_deployer_record, webapp_deployer_record,
logger, logger,
) )
logger.log("Publication complete.") logger.log("Publication complete.")
# 9. delete unused names from previous deployment (app and dns)
for fqdn in fqdns_to_release:
# Delete app deployment name and DNS name
deployment_name = "{deployment_record_namespace}/{fqdn}"
dns_name = "{deployment_record_namespace}/{fqdn}"
logger.log(f"Removing names {deployment_name} and {dns_name}")
laconic.delete_name(deployment_name)
laconic.delete_name(dns_name)
logger.log("END - process_app_deployment_request") logger.log("END - process_app_deployment_request")
@@ -304,6 +359,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",
@@ -381,6 +437,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,
@@ -429,6 +486,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)
@@ -506,21 +570,26 @@ def command( # noqa: C901
result = "ERROR" result = "ERROR"
continue continue
requested_name = r.attributes.dns requested_names = get_requested_names(r)
if not requested_name: if len(requested_names) == 0:
requested_name = generate_hostname_for_app(app) requested_names = [generate_hostname_for_app(app)]
main_logger.log( main_logger.log(
"Generating name %s for request %s." % (requested_name, r.id) "Generating name %s for request %s." % (requested_names[0], r.id)
) )
if ( # Skip request if any of the names is superseded
requested_name in skipped_by_name for requested_name in requested_names:
or requested_name in requests_by_name if (
): requested_name in skipped_by_name
main_logger.log( or requested_name in requests_by_name
"Ignoring request %s, it has been superseded." % r.id ):
) main_logger.log(
result = "SKIP" "Ignoring request %s, it has been superseded." % r.id
)
result = "SKIP"
break
if result == "SKIP":
continue continue
if skip_by_tag(r, include_tags, exclude_tags): if skip_by_tag(r, include_tags, exclude_tags):
@@ -528,15 +597,20 @@ def command( # noqa: C901
"Skipping request %s, filtered by tag (include %s, exclude %s, present %s)" "Skipping request %s, filtered by tag (include %s, exclude %s, present %s)"
% (r.id, include_tags, exclude_tags, r.attributes.tags) % (r.id, include_tags, exclude_tags, r.attributes.tags)
) )
skipped_by_name[requested_name] = r
for requested_name in requested_names:
skipped_by_name[requested_name] = r
result = "SKIP" result = "SKIP"
continue continue
main_logger.log( main_logger.log(
"Found pending request %s to run application %s on %s." "Found pending request %s to run application %s on %s."
% (r.id, r.attributes.application, requested_name) % (r.id, r.attributes.application, ','.join(str(x) for x in requested_names))
) )
requests_by_name[requested_name] = r
# Set request for on of the names
requests_by_name[requested_names[-1]] = r
except Exception as e: except Exception as e:
result = "ERROR" result = "ERROR"
main_logger.log(f"ERROR examining request {r.id}: " + str(e)) main_logger.log(f"ERROR examining request {r.id}: " + str(e))
@@ -665,6 +739,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,
@@ -76,6 +76,7 @@ def command( # noqa: C901
"name": hostname, "name": hostname,
"publicKey": pub_key, "publicKey": pub_key,
"paymentAddress": payment_address, "paymentAddress": payment_address,
"deployerVersion": "1.0.0",
} }
} }
@@ -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
try: acquire_lock(self, lock_file_path, LOCK_TIMEOUT)
# Try to acquire the lock try:
fcntl.flock(lock_file, fcntl.LOCK_EX) return func(self, *args, **kwargs)
finally:
# Call the actual function # Release the lock after the function completes
result = func(self, *args, **kwargs) release_lock(self, lock_file_path)
finally:
# Always release the lock
fcntl.flock(lock_file, fcntl.LOCK_UN)
return result
return wrapper return wrapper
+75 -44
View File
@@ -21,6 +21,7 @@ import random
import subprocess import subprocess
import sys import sys
import tempfile import tempfile
from typing import List
import uuid import uuid
import yaml import yaml
@@ -117,7 +118,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 +126,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"]
@@ -683,10 +686,12 @@ def publish_deployment(
laconic: LaconicRegistryClient, laconic: LaconicRegistryClient,
app_record, app_record,
deploy_record, deploy_record,
deployment_lrn, deployment_lrns,
dns_record, existing_dns_records_by_lrns,
dns_lrn, dns_lrns: List[str],
dns_record_namespace,
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,
@@ -698,40 +703,48 @@ def publish_deployment(
int(deploy_record.attributes.version.split(".")[-1]) + 1 int(deploy_record.attributes.version.split(".")[-1]) + 1
) )
if not dns_record: dns_ids = []
dns_ver = "0.0.1" for dns_lrn in dns_lrns:
else: dns_record = existing_dns_records_by_lrns[dns_lrn]
dns_ver = "0.0.%d" % (int(dns_record.attributes.version.split(".")[-1]) + 1) if not dns_record:
dns_ver = "0.0.1"
else:
dns_ver = "0.0.%d" % (int(dns_record.attributes.version.split(".")[-1]) + 1)
fqdn = dns_lrn.removeprefix(f"{dns_record_namespace}/")
uniq = uuid.uuid4()
new_dns_record = {
"record": {
"type": "DnsRecord",
"version": dns_ver,
"name": fqdn,
"resource_type": "A",
"meta": {"so": uniq.hex},
}
}
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.")
dns_id = laconic.publish(new_dns_record, [dns_lrn])
dns_ids.append(dns_id)
spec = yaml.full_load(open(os.path.join(deployment_dir, "spec.yml"))) spec = yaml.full_load(open(os.path.join(deployment_dir, "spec.yml")))
fqdn = spec["network"]["http-proxy"][0]["host-name"] last_fqdn = spec["network"]["http-proxy"][-1]["host-name"]
last_dns_id = dns_ids[-1]
uniq = uuid.uuid4()
new_dns_record = {
"record": {
"type": "DnsRecord",
"version": dns_ver,
"name": fqdn,
"resource_type": "A",
"meta": {"so": uniq.hex},
}
}
if app_deployment_request:
new_dns_record["record"]["request"] = app_deployment_request.id
if logger:
logger.log("Publishing DnsRecord.")
dns_id = laconic.publish(new_dns_record, [dns_lrn])
new_deployment_record = { new_deployment_record = {
"record": { "record": {
"type": "ApplicationDeploymentRecord", "type": "ApplicationDeploymentRecord",
"version": deploy_ver, "version": deploy_ver,
"url": f"https://{fqdn}", "url": f"https://{last_fqdn}",
"name": app_record.attributes.name, "name": app_record.attributes.name,
"application": app_record.id, "application": app_record.id,
"dns": dns_id, "dns": last_dns_id,
"meta": { "meta": {
"config": file_hash(os.path.join(deployment_dir, "config.env")), "config": file_hash(os.path.join(deployment_dir, "config.env")),
"so": uniq.hex, "so": uniq.hex,
@@ -753,21 +766,34 @@ def publish_deployment(
if logger: if logger:
logger.log("Publishing ApplicationDeploymentRecord.") logger.log("Publishing ApplicationDeploymentRecord.")
deployment_id = laconic.publish(new_deployment_record, [deployment_lrn]) deployment_id = laconic.publish(new_deployment_record, deployment_lrns)
return {"dns": dns_id, "deployment": deployment_id} return {"dns": dns_ids, "deployment": deployment_id}
def hostname_for_deployment_request(app_deployment_request, laconic): def get_requested_names(app_deployment_request):
dns_name = app_deployment_request.attributes.dns request_dns = app_deployment_request.attributes.dns
if not dns_name: return request_dns.split(",") if request_dns else []
def hostnames_for_deployment_request(app_deployment_request, laconic):
requested_names = get_requested_names(app_deployment_request)
if len(requested_names) == 0:
app = laconic.get_record( app = laconic.get_record(
app_deployment_request.attributes.application, require=True app_deployment_request.attributes.application, require=True
) )
dns_name = generate_hostname_for_app(app) return [generate_hostname_for_app(app)]
elif dns_name.startswith("lrn://"):
record = laconic.get_record(dns_name, require=True) dns_names = []
dns_name = record.attributes.name for requested_name in requested_names:
return dns_name dns_name = requested_name
if dns_name.startswith("lrn://"):
record = laconic.get_record(dns_name, require=True)
dns_name = record.attributes.name
dns_names.append(dns_name)
return dns_names
def generate_hostname_for_app(app): def generate_hostname_for_app(app):
@@ -844,16 +870,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
return False 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
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(