diff --git a/docs/sam-config-docs.md b/docs/sam-config-docs.md index 944a78c2884..1cabad27c4b 100644 --- a/docs/sam-config-docs.md +++ b/docs/sam-config-docs.md @@ -28,6 +28,18 @@ region="us-east-1" profile="srirammv" ``` +### Specifying a Boolean deployment option + +``` +[default.deploy.parameters] +parallel_upload=true +``` + +Setting `parallel_upload` to `true` is equivalent to passing `--parallel-upload` on +`sam deploy`, enabling concurrent S3/ECR uploads during the packaging phase. +Up to 8 uploads run at once; set the `SAM_CLI_PARALLEL_UPLOAD_WORKERS` environment variable +to change the number of concurrent uploads. + Version ------- diff --git a/samcli/commands/deploy/command.py b/samcli/commands/deploy/command.py index 80f9c7e80e3..7be33ab55f2 100644 --- a/samcli/commands/deploy/command.py +++ b/samcli/commands/deploy/command.py @@ -159,6 +159,13 @@ @image_repository_option @image_repositories_option @force_upload_option +@click.option( + "--parallel-upload", + is_flag=True, + default=False, + help="Enable parallel upload of artifacts to S3/ECR during packaging before deployment. " + "Runs up to 8 uploads at once; set the SAM_CLI_PARALLEL_UPLOAD_WORKERS environment variable to change this.", +) @s3_prefix_option @kms_key_id_option @role_arn_option @@ -192,6 +199,7 @@ def cli( image_repository, image_repositories, force_upload, + parallel_upload, no_progressbar, s3_prefix, kms_key_id, @@ -230,6 +238,7 @@ def cli( image_repository, image_repositories, force_upload, + parallel_upload, no_progressbar, s3_prefix, kms_key_id, @@ -267,6 +276,7 @@ def do_cli( image_repository, image_repositories, force_upload, + parallel_upload, no_progressbar, s3_prefix, kms_key_id, @@ -341,6 +351,7 @@ def do_cli( config_file=config_file, disable_rollback=disable_rollback, language_extensions_enabled=language_extensions_enabled, + parallel_upload=parallel_upload, ) guided_context.run() else: @@ -379,6 +390,7 @@ def do_cli( kms_key_id=kms_key_id, use_json=use_json, force_upload=force_upload, + parallel_upload=guided_context.guided_parallel_upload if guided else parallel_upload, no_progressbar=no_progressbar if output_mode is not OutputOption.json else True, metadata=metadata, on_deploy=True, @@ -407,6 +419,7 @@ def do_cli( image_repository=guided_context.guided_image_repository if guided else image_repository, image_repositories=guided_context.guided_image_repositories if guided else image_repositories, force_upload=force_upload, + parallel_upload=guided_context.guided_parallel_upload if guided else parallel_upload, no_progressbar=no_progressbar, s3_prefix=guided_context.guided_s3_prefix if guided else s3_prefix, kms_key_id=kms_key_id, diff --git a/samcli/commands/deploy/core/options.py b/samcli/commands/deploy/core/options.py index d532b167fcb..c1fa6bd7fae 100644 --- a/samcli/commands/deploy/core/options.py +++ b/samcli/commands/deploy/core/options.py @@ -37,6 +37,7 @@ "disable_rollback", "on_failure", "force_upload", + "parallel_upload", "max_wait_duration", "express", ] diff --git a/samcli/commands/deploy/deploy_context.py b/samcli/commands/deploy/deploy_context.py index e44c3246ed0..be4c52fa63d 100644 --- a/samcli/commands/deploy/deploy_context.py +++ b/samcli/commands/deploy/deploy_context.py @@ -81,6 +81,7 @@ def __init__( language_extensions: Optional[bool] = None, express: bool = False, output: str = "text", + parallel_upload: bool = False, ): self.template_file = template_file self.stack_name = stack_name @@ -88,6 +89,7 @@ def __init__( self.image_repository = image_repository self.image_repositories = image_repositories self.force_upload = force_upload + self.parallel_upload = parallel_upload self.no_progressbar = no_progressbar self.s3_prefix = s3_prefix self.kms_key_id = kms_key_id @@ -182,6 +184,7 @@ def run(self): self.signing_profiles, self.use_changeset, self.disable_rollback, + self.parallel_upload, ) return self.deploy( self.stack_name, diff --git a/samcli/commands/deploy/guided_context.py b/samcli/commands/deploy/guided_context.py index 1a3335687ae..3826d181e8b 100644 --- a/samcli/commands/deploy/guided_context.py +++ b/samcli/commands/deploy/guided_context.py @@ -63,6 +63,7 @@ def __init__( config_file=None, disable_rollback=None, language_extensions_enabled: bool = False, + parallel_upload=False, ): self.template_file = template_file self.stack_name = stack_name @@ -97,6 +98,8 @@ def __init__( self.function_provider: Optional[SamFunctionProvider] = None self.disable_rollback = disable_rollback self._language_extensions_enabled = language_extensions_enabled + self.parallel_upload = parallel_upload + self.guided_parallel_upload = None @property def guided_capabilities(self): @@ -207,6 +210,8 @@ def guided_prompts(self, parameter_override_keys): self.guided_s3_prefix = stack_name self.guided_region = region self.guided_profile = self.profile + # Not prompted, so scripted guided deploys keep their answer order; the flag value is saved. + self.guided_parallel_upload = self.parallel_upload self._capabilities = input_capabilities if input_capabilities else default_capabilities self._parameter_overrides = ( input_parameter_overrides if input_parameter_overrides else self.parameter_overrides_from_cmdline @@ -590,6 +595,7 @@ def run(self): capabilities=self._capabilities, signing_profiles=self.signing_profiles, disable_rollback=self.disable_rollback, + parallel_upload=self.guided_parallel_upload, ) @staticmethod diff --git a/samcli/commands/deploy/utils.py b/samcli/commands/deploy/utils.py index 421a16c75df..56270e5ae98 100644 --- a/samcli/commands/deploy/utils.py +++ b/samcli/commands/deploy/utils.py @@ -20,6 +20,7 @@ def print_deploy_args( signing_profiles, use_changeset, disable_rollback, + parallel_upload, ): """ Print a table of the values that are used during a sam deploy. @@ -48,6 +49,7 @@ def print_deploy_args( :param signing_profiles: Signing profile details which will be used to sign functions/layers :param use_changeset: Flag to use or skip the usage of changesets :param disable_rollback: Preserve the state of previously provisioned resources when an operation fails. + :param parallel_upload: Whether artifact uploads run in parallel prior to deployment. """ _parameters = parameter_overrides.copy() @@ -70,6 +72,7 @@ def print_deploy_args( if use_changeset: click.echo(f"\tConfirm changeset : {confirm_changeset}") click.echo(f"\tDisable rollback : {disable_rollback}") + click.echo(f"\tParallel uploads : {parallel_upload}") if image_repository: msg = "Deployment image repository : " # NOTE(sriram-mv): tab length is 8 spaces. diff --git a/samcli/commands/package/package_context.py b/samcli/commands/package/package_context.py index d659aba71f8..5275adf793a 100644 --- a/samcli/commands/package/package_context.py +++ b/samcli/commands/package/package_context.py @@ -86,6 +86,7 @@ def __init__( resolve_image_repos=False, language_extensions=None, output="text", + parallel_upload=False, ): self.template_file = template_file self.s3_bucket = s3_bucket @@ -96,6 +97,7 @@ def __init__( self.output_template_file = output_template_file self.use_json = use_json self.force_upload = force_upload + self.parallel_upload = parallel_upload self.no_progressbar = no_progressbar self.metadata = metadata self.region = region @@ -152,13 +154,15 @@ def run(self): # Pass None instead of validating Docker client upfront - ECRUploader will validate only when needed docker_client = None + # Progress bars redraw the terminal in place and interleave badly across threads. + no_progressbar = self.no_progressbar or self.parallel_upload s3_uploader = S3Uploader( - s3_client, self.s3_bucket, self.s3_prefix, self.kms_key_id, self.force_upload, self.no_progressbar + s3_client, self.s3_bucket, self.s3_prefix, self.kms_key_id, self.force_upload, no_progressbar ) # attach the given metadata to the artifacts to be uploaded s3_uploader.artifact_metadata = self.metadata ecr_uploader = ECRUploader( - docker_client, ecr_client, self.image_repository, self.image_repositories, self.no_progressbar + docker_client, ecr_client, self.image_repository, self.image_repositories, no_progressbar ) self.uploaders = Uploaders(s3_uploader, ecr_uploader) @@ -218,6 +222,7 @@ def _export_without_language_extensions(self, template_path, original_template_d normalize_parameters=True, template_dict=original_template_dict, language_extensions_enabled=False, + parallel_upload=self.parallel_upload, ) return template.export() @@ -265,6 +270,7 @@ def _export_with_language_extensions(self, template_path, original_template_dict template_dict=copy.deepcopy(result.expanded_template), parameter_values=parameter_values, language_extensions_enabled=True, + parallel_upload=self.parallel_upload, ) exported_template = template.export() diff --git a/samcli/lib/package/artifact_exporter.py b/samcli/lib/package/artifact_exporter.py index 05ac43ce314..26335dfa6bb 100644 --- a/samcli/lib/package/artifact_exporter.py +++ b/samcli/lib/package/artifact_exporter.py @@ -17,7 +17,10 @@ import copy import logging import os -from typing import Any, Dict, List, Optional, Sequence, cast +import threading +from collections.abc import MutableMapping +from concurrent.futures import FIRST_EXCEPTION, Future, ThreadPoolExecutor, wait +from typing import Any, Callable, Dict, List, Optional, Sequence, Tuple, cast from botocore.utils import set_value_from_jmespath @@ -65,6 +68,100 @@ # NOTE: sriram-mv, A cyclic dependency on `Template` needs to be broken. +# Uploads are I/O bound, and every in-flight zip upload holds a temporary zip on disk (and member +# files in memory), so the pool is kept small by default and is configurable via the environment. +DEFAULT_PARALLEL_UPLOAD_WORKERS = 8 +PARALLEL_UPLOAD_WORKERS_ENV_VAR = "SAM_CLI_PARALLEL_UPLOAD_WORKERS" + + +def get_parallel_upload_workers() -> int: + """Number of concurrent artifact uploads for --parallel-upload.""" + value = os.environ.get(PARALLEL_UPLOAD_WORKERS_ENV_VAR) + if not value: + return DEFAULT_PARALLEL_UPLOAD_WORKERS + try: + workers = int(value) + except ValueError: + workers = 0 + if workers < 1: + LOG.warning( + "Ignoring invalid %s=%r; using %d parallel upload workers", + PARALLEL_UPLOAD_WORKERS_ENV_VAR, + value, + DEFAULT_PARALLEL_UPLOAD_WORKERS, + ) + return DEFAULT_PARALLEL_UPLOAD_WORKERS + return workers + + +class _UploadAborted(Exception): + """Raised by parallel export work that was skipped because another upload already failed.""" + + +def _is_upload_abort(error: Optional[BaseException]) -> bool: + """ + True if ``error`` is, or wraps, an ``_UploadAborted``. Exporters wrap ``do_export`` failures in + ``ExportFailedError(ex=...)``, so a nested stack that was skipped reaches its parent wrapped, + once per nesting level. + """ + seen = set() + while error is not None and id(error) not in seen: + if isinstance(error, _UploadAborted): + return True + seen.add(id(error)) + error = getattr(error, "ex", None) or error.__cause__ + return False + + +class _ThreadSafeUploadCache(MutableMapping[str, str]): + """ + Thread-safe mapping used to deduplicate uploads across threads. + + With ``key_by_packaging=True`` it is an export-scoped memo keyed by local path *and* how the + path is packaged (zip method and extension), so resources that share a directory but package it + differently, e.g. a Lambda function and an Elastic Beanstalk application version, never receive + each other's artifact. Without it, keys are plain local paths, matching the experimental + PackagePerformance cache. + """ + + def __init__(self, initial: Optional[MutableMapping[str, str]] = None, key_by_packaging: bool = False): + # Copy into a regular dict so we can safely snapshot under a lock + self._cache: Dict[str, str] = dict(initial or {}) + self._key_by_packaging = key_by_packaging + self._lock = threading.Lock() + self._key_locks: Dict[str, threading.Lock] = {} + + def cache_key(self, local_path: str, packaging: str, extension: Optional[str]) -> str: + """Key for an upload of ``local_path`` packaged as ``packaging`` with ``extension``.""" + if not self._key_by_packaging: + return local_path + return f"{packaging}|{extension or ''}|{local_path}" + + def key_lock(self, key: str) -> threading.Lock: + """Lock serializing check-then-upload for a single key; other keys proceed concurrently.""" + with self._lock: + return self._key_locks.setdefault(key, threading.Lock()) + + def __getitem__(self, key: str) -> str: # pragma: no cover - small helper + with self._lock: + return self._cache[key] + + def __setitem__(self, key: str, value: str) -> None: # pragma: no cover - small helper + with self._lock: + self._cache[key] = value + + def __delitem__(self, key: str) -> None: # pragma: no cover - small helper + with self._lock: + del self._cache[key] + + def __iter__(self): # pragma: no cover - small helper + with self._lock: + return iter(dict(self._cache)) + + def __len__(self) -> int: # pragma: no cover - small helper + with self._lock: + return len(self._cache) + def _resolve_nested_stack_parameters(nested_params: Dict, parent_parameter_values: Dict) -> Dict: """ @@ -253,7 +350,16 @@ def do_export(self, resource_id, resource_dict, parent_dir): temporary_file.write(exported_template_str) temporary_file.flush() remote_path = get_uploaded_s3_object_name(file_path=temporary_file.name, extension="template") - url = self.uploader.upload(temporary_file.name, remote_path) + upload_executor = getattr(self, "upload_executor", None) + upload_abort = getattr(self, "upload_abort", None) + if upload_abort is not None and upload_abort.is_set(): + raise _UploadAborted() + if upload_executor is not None: + # Keep the rendered child template upload under the shared upload bound; this runs on + # a nested-stack coordination thread, which may wait on the upload pool. + url = upload_executor.submit(self.uploader.upload, temporary_file.name, remote_path).result() + else: + url = self.uploader.upload(temporary_file.name, remote_path) # TemplateUrl property requires S3 URL to be in path-style format parts = parse_s3_url(url, version_property="Version") @@ -277,6 +383,9 @@ def _do_export_without_language_extensions(self, resource_id: str, template_path normalize_parameters=True, parent_stack_id=resource_id, language_extensions_enabled=False, + parallel_upload=getattr(self, "parallel_upload", False), + upload_executor=getattr(self, "upload_executor", None), + upload_abort=getattr(self, "upload_abort", None), ).export() def _do_export_with_language_extensions( @@ -371,6 +480,9 @@ def _do_export_with_language_extensions( template_dict=copy.deepcopy(result.expanded_template), parameter_values=parameter_values, language_extensions_enabled=self.language_extensions_enabled, + parallel_upload=getattr(self, "parallel_upload", False), + upload_executor=getattr(self, "upload_executor", None), + upload_abort=getattr(self, "upload_abort", None), ) exported_template = template.export() @@ -404,6 +516,9 @@ def _do_export_with_language_extensions( parent_stack_id=resource_id, parameter_values=parameter_values, language_extensions_enabled=self.language_extensions_enabled, + parallel_upload=getattr(self, "parallel_upload", False), + upload_executor=getattr(self, "upload_executor", None), + upload_abort=getattr(self, "upload_abort", None), ).export() return exported_template_dict @@ -486,6 +601,9 @@ def __init__( parameter_values: Optional[Dict] = None, template_dict: Optional[Dict] = None, language_extensions_enabled: bool = False, + parallel_upload: bool = False, + upload_executor: Optional[ThreadPoolExecutor] = None, + upload_abort: Optional[threading.Event] = None, ): """ Reads the template and makes it ready for export @@ -530,6 +648,12 @@ def __init__( # collections that Ref a parameter). None preserves pre-existing behavior. self.parameter_values = parameter_values self.language_extensions_enabled = language_extensions_enabled + self.parallel_upload = parallel_upload + # Shared by the root template and every nested-stack child so that one pool bounds both + # threads and in-flight uploads for the whole export. Created by the root when None. + self.upload_executor = upload_executor + # Set on the first failure anywhere in the export so descendants stop starting new uploads. + self.upload_abort = upload_abort def _export_global_artifacts(self, template_dict: Dict) -> Dict: """See module-level _export_global_artifacts_pass for the canonical @@ -600,10 +724,36 @@ def export(self) -> Dict: self._apply_global_values() self.template_dict = self._export_global_artifacts(self.template_dict) - cache: Optional[Dict] = None + cache: Optional[MutableMapping[str, str]] = None if is_experimental_enabled(ExperimentalFlag.PackagePerformance): cache = {} + if self.parallel_upload: + # Parallel jobs sharing a code path must not all zip and upload it at once. Keep the + # experimental cache's path keys when it is on; otherwise use an export-scoped memo that + # also keys on packaging, so it cannot conflate differently packaged artifacts. + cache = _ThreadSafeUploadCache(cache, key_by_packaging=cache is None) + + if not self.parallel_upload: + for _, job in self._collect_export_jobs(cache, None, None): + job() + elif self.upload_executor is not None: + jobs = self._collect_export_jobs(cache, self.upload_executor, self.upload_abort) + self._run_export_jobs(jobs, self.upload_executor, self.upload_abort) + else: + abort = threading.Event() + with ThreadPoolExecutor(max_workers=get_parallel_upload_workers()) as executor: + self._run_export_jobs(self._collect_export_jobs(cache, executor, abort), executor, abort) + + return self.template_dict + def _collect_export_jobs( + self, + cache: Optional[MutableMapping[str, str]], + executor: Optional[ThreadPoolExecutor], + abort: Optional[threading.Event], + ) -> List[Tuple[bool, Callable[[], None]]]: + """Return (is_nested_stack, job) pairs for every resource that has artifacts to export.""" + jobs: List[Tuple[bool, Callable[[], None]]] = [] for resource_logical_id, resource in iter_regular_resources(self.template_dict): resource_type = resource.get("Type", None) resource_dict = resource.get("Properties", {}) @@ -615,13 +765,95 @@ def export(self) -> Dict: continue if resource_dict.get("PackageType", ZIP) != exporter_class.ARTIFACT_TYPE: continue - # Export code resources - exporter = exporter_class(self.uploaders, self.code_signer, cache) - exporter.parent_parameter_values = self.parameter_values - exporter.language_extensions_enabled = self.language_extensions_enabled - exporter.export(full_path, resource_dict, self.template_dir) - return self.template_dict + is_nested_stack = isinstance(exporter_class, type) and issubclass( + exporter_class, CloudFormationStackResource + ) + job = self._build_export_job(exporter_class, full_path, resource_dict, cache, executor, abort) + jobs.append((is_nested_stack, job)) + return jobs + + def _build_export_job( + self, + exporter_class, + resource_full_path: str, + resource_dict: Dict, + cache: Optional[MutableMapping[str, str]], + executor: Optional[ThreadPoolExecutor] = None, + abort: Optional[threading.Event] = None, + ) -> Callable[[], None]: + def _job() -> None: + if abort is not None and abort.is_set(): + raise _UploadAborted() + # Export code resources + exporter = exporter_class(self.uploaders, self.code_signer, cache) + exporter.parent_parameter_values = self.parameter_values + exporter.language_extensions_enabled = self.language_extensions_enabled + exporter.parallel_upload = self.parallel_upload + exporter.upload_executor = executor + exporter.upload_abort = abort + exporter.export(resource_full_path, resource_dict, self.template_dir) + + return _job + + @staticmethod + def _run_export_jobs( + jobs: List[Tuple[bool, Callable[[], None]]], + executor: ThreadPoolExecutor, + abort: Optional[threading.Event] = None, + ) -> None: + """ + Artifact uploads go to the shared upload executor. Nested-stack exports only coordinate: + they wait on their children's uploads, so they run on a separate pool with one thread per + nested stack at this level, never in the upload pool, where waiting could deadlock it. + Sibling nested stacks therefore proceed concurrently while the shared pool bounds uploads, + and a single fail-fast wait covers both. ``abort`` is shared by the whole export: once any + job fails, work that has not started yet at any level is skipped. + """ + if abort is not None and abort.is_set(): + raise _UploadAborted() + upload_futures = [executor.submit(job) for is_nested_stack, job in jobs if not is_nested_stack] + nested_jobs = [job for is_nested_stack, job in jobs if is_nested_stack] + if not nested_jobs: + Template._wait_fail_fast(upload_futures, abort) + return + with ThreadPoolExecutor(max_workers=len(nested_jobs)) as coordinator: + nested_futures = [coordinator.submit(job) for job in nested_jobs] + Template._wait_fail_fast(upload_futures + nested_futures, abort) + + @staticmethod + def _wait_fail_fast(futures: List[Future], abort: Optional[threading.Event] = None) -> None: + """ + Wait for futures. On the first failure, signal the rest of the export to stop, cancel what + has not started here, wait for running work, log the other failures and re-raise the first + real one (in submission order, so the surfaced error is stable). + """ + done, _ = wait(futures, return_when=FIRST_EXCEPTION) + if not any(future.exception() is not None for future in done): + return + if abort is not None: + abort.set() + for future in futures: + future.cancel() + wait(futures) + failed = [future for future in futures if not future.cancelled() and future.exception() is not None] + real = [future for future in failed if not _is_upload_abort(future.exception())] + surfaced = (real or failed)[0] + Template._log_other_failures(futures, surfaced=surfaced) + raise cast(BaseException, surfaced.exception()) + + @staticmethod + def _log_other_failures(futures: List[Future], surfaced: Optional[Future] = None) -> None: + """Log failures that are not the one being re-raised, so no upload error is silently lost.""" + for future in futures: + if future is surfaced or future.cancelled(): + continue + error = future.exception() + if error is not None and not _is_upload_abort(error): + # One line per extra failure; the surfaced error is reported normally and the full + # traceback of the others is only shown with --debug. + LOG.error("Parallel artifact upload also failed: %s", error) + LOG.debug("Traceback for the parallel upload failure above", exc_info=error) def delete(self, retain_resources: List): """ diff --git a/samcli/lib/package/ecr_uploader.py b/samcli/lib/package/ecr_uploader.py index ae0337cfdb0..af93528999f 100644 --- a/samcli/lib/package/ecr_uploader.py +++ b/samcli/lib/package/ecr_uploader.py @@ -4,6 +4,7 @@ import base64 import logging +import threading from io import StringIO from pathlib import Path from typing import Dict @@ -50,6 +51,7 @@ def __init__( self.stream = StreamWriter(stream=stream, auto_flush=True) self.log_streamer = LogStreamer(stream=self.stream) self.login_session_active = False + self._login_lock = threading.Lock() @property def docker_client(self): @@ -88,8 +90,10 @@ def upload(self, image, resource_name): :return: remote ECR image path that has been uploaded. """ if not self.login_session_active: - self.login() - self.login_session_active = True + with self._login_lock: + if not self.login_session_active: + self.login() + self.login_session_active = True # Sometimes the `resource_name` is used as the `image` parameter to `tag_translation`. # This is because these two cases (directly from an archive or by ID) are effectively diff --git a/samcli/lib/package/utils.py b/samcli/lib/package/utils.py index fb5f22e55bc..73e5881c2a2 100644 --- a/samcli/lib/package/utils.py +++ b/samcli/lib/package/utils.py @@ -186,29 +186,38 @@ def upload_local_artifacts( local_path = make_abs_path(parent_dir, local_path) - if previously_uploaded and local_path in previously_uploaded: - result = previously_uploaded[local_path] - LOG.debug("Skipping upload of %s since is already uploaded to %s", local_path, result) - return cast(str, result) - - # Or, pointing to a folder. Zip the folder and upload (zip_method is changed based on resource type) - if is_local_folder(local_path): - result = zip_and_upload( - local_path, - uploader, - extension, - zip_method=make_zip_with_lambda_permissions if resource_type in LAMBDA_LOCAL_RESOURCES else make_zip, - ) - if previously_uploaded is not None: - previously_uploaded[local_path] = result - return result - - # Path could be pointing to a file. Upload the file - if is_local_file(local_path): - result = uploader.upload_with_dedup(local_path) - if previously_uploaded is not None: - previously_uploaded[local_path] = result - return result + is_lambda_package = resource_type in LAMBDA_LOCAL_RESOURCES + zip_method = make_zip_with_lambda_permissions if is_lambda_package else make_zip + # A thread-safe cache (parallel uploads) keys results by how the path is packaged and exposes a + # per-key lock, so the check-then-upload below is atomic: concurrent jobs sharing a local path + # zip and upload it once and the rest reuse the result. + key_fn = getattr(previously_uploaded, "cache_key", None) + cache_key = key_fn(local_path, "lambda-zip" if is_lambda_package else "zip", extension) if key_fn else local_path + key_lock = getattr(previously_uploaded, "key_lock", None) + with key_lock(cache_key) if key_lock else contextlib.nullcontext(): + if previously_uploaded and cache_key in previously_uploaded: + result = previously_uploaded[cache_key] + LOG.debug("Skipping upload of %s since is already uploaded to %s", local_path, result) + return cast(str, result) + + # Or, pointing to a folder. Zip the folder and upload (zip_method is changed based on resource type) + if is_local_folder(local_path): + result = zip_and_upload( + local_path, + uploader, + extension, + zip_method=zip_method, + ) + if previously_uploaded is not None: + previously_uploaded[cache_key] = result + return result + + # Path could be pointing to a file. Upload the file + if is_local_file(local_path): + result = uploader.upload_with_dedup(local_path) + if previously_uploaded is not None: + previously_uploaded[cache_key] = result + return result raise InvalidLocalPathError(resource_id=resource_id, property_name=property_path, local_path=local_path) diff --git a/schema/samcli.json b/schema/samcli.json index fb48d4346fb..e45140708ad 100644 --- a/schema/samcli.json +++ b/schema/samcli.json @@ -1325,7 +1325,7 @@ "properties": { "parameters": { "title": "Parameters for the deploy command", - "description": "Available parameters for the deploy command:\n* guided:\nSpecify this flag to allow SAM CLI to guide you through the deployment using guided prompts.\n* template_file:\nAWS SAM template which references built artifacts for resources in the template. (if applicable)\n* no_execute_changeset:\nIndicates whether to execute the change set. Specify this flag to view stack changes before executing the change set.\n* fail_on_empty_changeset:\nSpecify whether AWS SAM CLI should return a non-zero exit code if there are no changes to be made to the stack. Defaults to a non-zero exit code.\n* confirm_changeset:\nPrompt to confirm if the computed changeset is to be deployed by SAM CLI.\n* disable_rollback:\nPreserves the state of previously provisioned resources when an operation fails.\n* on_failure:\nProvide an action to determine what will happen when a stack fails to create. Three actions are available:\n\n- ROLLBACK: This will rollback a stack to a previous known good state.\n\n- DELETE: The stack will rollback to a previous state if one exists, otherwise the stack will be deleted.\n\n- DO_NOTHING: The stack will not rollback or delete, this is the same as disabling rollback.\n\nDefault behaviour is ROLLBACK.\n\n\n\nThis option is mutually exclusive with --disable-rollback/--no-disable-rollback. You can provide\n--on-failure or --disable-rollback/--no-disable-rollback but not both at the same time.\n* max_wait_duration:\nMaximum duration in minutes to wait for the deployment to complete.\n* express:\nUse CloudFormation Express mode to speed up deployments by completing once resource configuration is applied, without waiting for full stabilization.\n* stack_name:\nName of the AWS CloudFormation stack.\n* s3_bucket:\nAWS S3 bucket where artifacts referenced in the template are uploaded.\n* image_repository:\nAWS ECR repository URI where artifacts referenced in the template are uploaded.\n* image_repositories:\nMapping of Function Logical ID to AWS ECR Repository URI.\n\nExample: Function_Logical_ID=ECR_Repo_Uri\nThis option can be specified multiple times.\n* force_upload:\nIndicates whether to override existing files in the S3 bucket. Specify this flag to upload artifacts even if they match existing artifacts in the S3 bucket.\n* s3_prefix:\nPrefix name that is added to the artifact's name when it is uploaded to the AWS S3 bucket.\n* kms_key_id:\nThe ID of an AWS KMS key that is used to encrypt artifacts that are at rest in the AWS S3 bucket.\n* role_arn:\nARN of an IAM role that AWS Cloudformation assumes when executing a deployment change set.\n* use_json:\nIndicates whether to use JSON as the format for the output AWS CloudFormation template. YAML is used by default.\n* resolve_s3:\nAutomatically resolve AWS S3 bucket for non-guided deployments. Enabling this option will also create a managed default AWS S3 bucket for you. If one does not provide a --s3-bucket value, the managed bucket will be used. Do not use --guided with this option.\n* resolve_image_repos:\nAutomatically create and delete ECR repositories for image-based functions in non-guided deployments. A companion stack containing ECR repos for each function will be deployed along with the template stack. Automatically created image repositories will be deleted if the corresponding functions are removed.\n* metadata:\nMap of metadata to attach to ALL the artifacts that are referenced in the template.\n* notification_arns:\nARNs of SNS topics that AWS Cloudformation associates with the stack.\n* tags:\nList of tags to associate with the stack.\n* parameter_overrides:\nString that contains AWS CloudFormation parameter overrides encoded as key=value pairs.\n* signing_profiles:\nA string that contains Code Sign configuration parameters as FunctionOrLayerNameToSign=SigningProfileName:SigningProfileOwner Since signing profile owner is optional, it could also be written as FunctionOrLayerNameToSign=SigningProfileName\n* no_progressbar:\nDoes not showcase a progress bar when uploading artifacts to S3 and pushing docker images to ECR\n* capabilities:\nList of capabilities that one must specify before AWS Cloudformation can create certain stacks.\n\nAccepted Values: CAPABILITY_IAM, CAPABILITY_NAMED_IAM, CAPABILITY_RESOURCE_POLICY, CAPABILITY_AUTO_EXPAND.\n\nLearn more at: https://docs.aws.amazon.com/serverlessrepo/latest/devguide/acknowledging-application-capabilities.html\n* language_extensions:\nExpand AWS::LanguageExtensions transforms (Fn::ForEach, Fn::Length, Fn::ToJsonString, Fn::FindInMap with DefaultValue) locally before running SAM transforms. Off by default. Equivalent env var: SAM_CLI_ENABLE_LANGUAGE_EXTENSIONS=1.\n* output:\nOutput the results from the command in a given output format. Supported formats: text (default), json.\n* profile:\nSelect a specific profile from your credential file to get AWS credentials.\n* region:\nSet the AWS Region of the service. (e.g. us-east-1)\n* beta_features:\nEnable/Disable beta features.\n* debug:\nTurn on debug logging to print debug message generated by AWS SAM CLI and display timestamps.\n* save_params:\nSave the parameters provided via the command line to the configuration file.", + "description": "Available parameters for the deploy command:\n* guided:\nSpecify this flag to allow SAM CLI to guide you through the deployment using guided prompts.\n* template_file:\nAWS SAM template which references built artifacts for resources in the template. (if applicable)\n* no_execute_changeset:\nIndicates whether to execute the change set. Specify this flag to view stack changes before executing the change set.\n* fail_on_empty_changeset:\nSpecify whether AWS SAM CLI should return a non-zero exit code if there are no changes to be made to the stack. Defaults to a non-zero exit code.\n* confirm_changeset:\nPrompt to confirm if the computed changeset is to be deployed by SAM CLI.\n* disable_rollback:\nPreserves the state of previously provisioned resources when an operation fails.\n* on_failure:\nProvide an action to determine what will happen when a stack fails to create. Three actions are available:\n\n- ROLLBACK: This will rollback a stack to a previous known good state.\n\n- DELETE: The stack will rollback to a previous state if one exists, otherwise the stack will be deleted.\n\n- DO_NOTHING: The stack will not rollback or delete, this is the same as disabling rollback.\n\nDefault behaviour is ROLLBACK.\n\n\n\nThis option is mutually exclusive with --disable-rollback/--no-disable-rollback. You can provide\n--on-failure or --disable-rollback/--no-disable-rollback but not both at the same time.\n* max_wait_duration:\nMaximum duration in minutes to wait for the deployment to complete.\n* express:\nUse CloudFormation Express mode to speed up deployments by completing once resource configuration is applied, without waiting for full stabilization.\n* stack_name:\nName of the AWS CloudFormation stack.\n* s3_bucket:\nAWS S3 bucket where artifacts referenced in the template are uploaded.\n* image_repository:\nAWS ECR repository URI where artifacts referenced in the template are uploaded.\n* image_repositories:\nMapping of Function Logical ID to AWS ECR Repository URI.\n\nExample: Function_Logical_ID=ECR_Repo_Uri\nThis option can be specified multiple times.\n* force_upload:\nIndicates whether to override existing files in the S3 bucket. Specify this flag to upload artifacts even if they match existing artifacts in the S3 bucket.\n* parallel_upload:\nEnable parallel upload of artifacts to S3/ECR during packaging before deployment. Runs up to 8 uploads at once; set the SAM_CLI_PARALLEL_UPLOAD_WORKERS environment variable to change this.\n* s3_prefix:\nPrefix name that is added to the artifact's name when it is uploaded to the AWS S3 bucket.\n* kms_key_id:\nThe ID of an AWS KMS key that is used to encrypt artifacts that are at rest in the AWS S3 bucket.\n* role_arn:\nARN of an IAM role that AWS Cloudformation assumes when executing a deployment change set.\n* use_json:\nIndicates whether to use JSON as the format for the output AWS CloudFormation template. YAML is used by default.\n* resolve_s3:\nAutomatically resolve AWS S3 bucket for non-guided deployments. Enabling this option will also create a managed default AWS S3 bucket for you. If one does not provide a --s3-bucket value, the managed bucket will be used. Do not use --guided with this option.\n* resolve_image_repos:\nAutomatically create and delete ECR repositories for image-based functions in non-guided deployments. A companion stack containing ECR repos for each function will be deployed along with the template stack. Automatically created image repositories will be deleted if the corresponding functions are removed.\n* metadata:\nMap of metadata to attach to ALL the artifacts that are referenced in the template.\n* notification_arns:\nARNs of SNS topics that AWS Cloudformation associates with the stack.\n* tags:\nList of tags to associate with the stack.\n* parameter_overrides:\nString that contains AWS CloudFormation parameter overrides encoded as key=value pairs.\n* signing_profiles:\nA string that contains Code Sign configuration parameters as FunctionOrLayerNameToSign=SigningProfileName:SigningProfileOwner Since signing profile owner is optional, it could also be written as FunctionOrLayerNameToSign=SigningProfileName\n* no_progressbar:\nDoes not showcase a progress bar when uploading artifacts to S3 and pushing docker images to ECR\n* capabilities:\nList of capabilities that one must specify before AWS Cloudformation can create certain stacks.\n\nAccepted Values: CAPABILITY_IAM, CAPABILITY_NAMED_IAM, CAPABILITY_RESOURCE_POLICY, CAPABILITY_AUTO_EXPAND.\n\nLearn more at: https://docs.aws.amazon.com/serverlessrepo/latest/devguide/acknowledging-application-capabilities.html\n* language_extensions:\nExpand AWS::LanguageExtensions transforms (Fn::ForEach, Fn::Length, Fn::ToJsonString, Fn::FindInMap with DefaultValue) locally before running SAM transforms. Off by default. Equivalent env var: SAM_CLI_ENABLE_LANGUAGE_EXTENSIONS=1.\n* output:\nOutput the results from the command in a given output format. Supported formats: text (default), json.\n* profile:\nSelect a specific profile from your credential file to get AWS credentials.\n* region:\nSet the AWS Region of the service. (e.g. us-east-1)\n* beta_features:\nEnable/Disable beta features.\n* debug:\nTurn on debug logging to print debug message generated by AWS SAM CLI and display timestamps.\n* save_params:\nSave the parameters provided via the command line to the configuration file.", "type": "object", "properties": { "guided": { @@ -1410,6 +1410,11 @@ "type": "boolean", "description": "Indicates whether to override existing files in the S3 bucket. Specify this flag to upload artifacts even if they match existing artifacts in the S3 bucket." }, + "parallel_upload": { + "title": "parallel_upload", + "type": "boolean", + "description": "Enable parallel upload of artifacts to S3/ECR during packaging before deployment. Runs up to 8 uploads at once; set the SAM_CLI_PARALLEL_UPLOAD_WORKERS environment variable to change this." + }, "s3_prefix": { "title": "s3_prefix", "type": "string", diff --git a/tests/unit/commands/deploy/test_command.py b/tests/unit/commands/deploy/test_command.py index aaa8ab9bfda..ae9bc4f2581 100644 --- a/tests/unit/commands/deploy/test_command.py +++ b/tests/unit/commands/deploy/test_command.py @@ -40,6 +40,7 @@ def setUp(self): self.fail_on_empty_changset = True self.role_arn = "role_arn" self.force_upload = False + self.parallel_upload = False self.no_progressbar = False self.metadata = {"abc": "def"} self.region = None @@ -87,6 +88,7 @@ def test_all_args(self, mock_deploy_context, mock_deploy_click, mock_package_con image_repository=self.image_repository, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -123,6 +125,7 @@ def test_all_args(self, mock_deploy_context, mock_deploy_click, mock_package_con image_repository=self.image_repository, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -159,6 +162,7 @@ def _do_cli_with(self, **overrides): image_repository=self.image_repository, image_repositories=None, force_upload=self.force_upload, + parallel_upload=False, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -322,6 +326,7 @@ def test_all_args_guided_no_to_authorization_confirmation_prompt( image_repository=None, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -427,6 +432,7 @@ def test_all_args_guided_use_defaults( image_repository=None, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -462,6 +468,7 @@ def test_all_args_guided_use_defaults( image_repository=None, image_repositories={"HelloWorldFunction": "123456789012.dkr.ecr.us-east-1.amazonaws.com/managed-ecr"}, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix="sam-app", kms_key_id=self.kms_key_id, @@ -505,6 +512,7 @@ def test_all_args_guided_use_defaults( s3_prefix="sam-app", signing_profiles=self.signing_profiles, disable_rollback=True, + parallel_upload=self.parallel_upload, ) mock_managed_stack.assert_called_with(profile=self.profile, region="us-east-1") self.assertEqual(context_mock.run.call_count, 1) @@ -578,6 +586,7 @@ def test_all_args_guided( image_repository=None, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -613,6 +622,7 @@ def test_all_args_guided( image_repository=None, image_repositories={"HelloWorldFunction": "123456789012.dkr.ecr.us-east-1.amazonaws.com/test1"}, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix="sam-app", kms_key_id=self.kms_key_id, @@ -656,6 +666,7 @@ def test_all_args_guided( s3_prefix="sam-app", signing_profiles=self.signing_profiles, disable_rollback=True, + parallel_upload=self.parallel_upload, ) mock_managed_stack.assert_called_with(profile=self.profile, region="us-east-1") self.assertEqual(context_mock.run.call_count, 1) @@ -732,6 +743,7 @@ def test_all_args_guided_no_save_echo_param_to_config( image_repository=None, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -768,6 +780,7 @@ def test_all_args_guided_no_save_echo_param_to_config( image_repository=None, image_repositories={"HelloWorldFunction": "123456789012.dkr.ecr.us-east-1.amazonaws.com/test1"}, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix="sam-app", kms_key_id=self.kms_key_id, @@ -899,6 +912,7 @@ def test_all_args_guided_no_params_save_config( image_repository=None, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -935,6 +949,7 @@ def test_all_args_guided_no_params_save_config( image_repository=None, image_repositories={"HelloWorldFunction": "123456789012.dkr.ecr.us-east-1.amazonaws.com/test1"}, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix="sam-app", kms_key_id=self.kms_key_id, @@ -1046,6 +1061,7 @@ def test_all_args_guided_no_params_no_save_config( image_repository=None, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -1081,6 +1097,7 @@ def test_all_args_guided_no_params_no_save_config( image_repository=None, image_repositories={"HelloWorldFunction": "123456789012.dkr.ecr.us-east-1.amazonaws.com/test1"}, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix="sam-app", kms_key_id=self.kms_key_id, @@ -1129,6 +1146,7 @@ def test_all_args_resolve_s3( image_repository=None, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -1163,6 +1181,7 @@ def test_all_args_resolve_s3( stack_name=self.stack_name, s3_bucket="managed-s3-bucket", force_upload=self.force_upload, + parallel_upload=self.parallel_upload, image_repository=None, image_repositories=None, no_progressbar=self.no_progressbar, @@ -1201,6 +1220,7 @@ def test_resolve_s3_and_s3_bucket_both_set(self): image_repository=None, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -1255,6 +1275,7 @@ def test_all_args_resolve_image_repos( image_repository=None, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -1289,6 +1310,7 @@ def test_all_args_resolve_image_repos( stack_name=self.stack_name, s3_bucket=self.s3_bucket, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, image_repository=None, image_repositories={"HelloWorldFunction1": self.image_repository}, no_progressbar=self.no_progressbar, @@ -1336,6 +1358,7 @@ def test_passing_parameter_overrides_to_context( image_repository=self.image_repository, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -1372,6 +1395,7 @@ def test_passing_parameter_overrides_to_context( image_repository=self.image_repository, image_repositories=None, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, s3_prefix=self.s3_prefix, kms_key_id=self.kms_key_id, @@ -1406,6 +1430,7 @@ def test_passing_parameter_overrides_to_context( kms_key_id=self.kms_key_id, use_json=self.use_json, force_upload=self.force_upload, + parallel_upload=self.parallel_upload, no_progressbar=self.no_progressbar, metadata=self.metadata, on_deploy=True, diff --git a/tests/unit/commands/deploy/test_guided_context.py b/tests/unit/commands/deploy/test_guided_context.py index 052d8a84a76..976faff8320 100644 --- a/tests/unit/commands/deploy/test_guided_context.py +++ b/tests/unit/commands/deploy/test_guided_context.py @@ -104,6 +104,41 @@ def test_guided_prompts_check_defaults_non_public_resources_zips( language_extensions_enabled=False, ) + @parameterized.expand([(True,), (False,)]) + @patch("samcli.commands.deploy.guided_context.get_resource_full_path_by_id") + @patch("samcli.commands.deploy.guided_context.prompt") + @patch("samcli.commands.deploy.guided_context.confirm") + @patch("samcli.commands.deploy.guided_context.manage_stack") + @patch("samcli.commands.deploy.guided_context.auth_per_resource") + @patch("samcli.commands.deploy.guided_context.SamLocalStackProvider.get_stacks") + @patch("samcli.commands.deploy.guided_context.SamFunctionProvider") + @patch("samcli.commands.deploy.guided_context.signer_config_per_function") + def test_guided_prompts_keep_parallel_upload_flag_without_prompting( + self, + parallel_upload, + patched_signer_config_per_function, + patched_sam_function_provider, + patched_get_buildable_stacks, + patchedauth_per_resource, + patched_manage_stack, + patched_confirm, + patched_prompt, + get_resource_full_path_by_id_mock, + ): + patched_signer_config_per_function.return_value = (None, None) + patched_sam_function_provider.return_value.functions = {} + patched_get_buildable_stacks.return_value = (Mock(), []) + patchedauth_per_resource.return_value = [("HelloWorldFunction", False)] + patched_confirm.side_effect = [True, False, False, True, False, True, True] + patched_manage_stack.return_value = "managed_s3_stack" + self.gc.parallel_upload = parallel_upload + + self.gc.guided_prompts(parameter_override_keys=None) + + self.assertEqual(parallel_upload, self.gc.guided_parallel_upload) + prompted = [confirm_call.args[0] for confirm_call in patched_confirm.call_args_list] + self.assertFalse(any("parallel" in text.lower() for text in prompted)) + @patch("samcli.commands.deploy.guided_context.get_resource_full_path_by_id") @patch("samcli.commands.deploy.guided_context.prompt") @patch("samcli.commands.deploy.guided_context.confirm") diff --git a/tests/unit/commands/deploy/test_utils.py b/tests/unit/commands/deploy/test_utils.py new file mode 100644 index 00000000000..14b9aabadb9 --- /dev/null +++ b/tests/unit/commands/deploy/test_utils.py @@ -0,0 +1,70 @@ +from unittest import TestCase +from unittest.mock import patch + +from samcli.commands.deploy.utils import ( + hide_noecho_parameter_overrides, + print_deploy_args, + sanitize_parameter_overrides, +) + + +class TestDeployUtils(TestCase): + @patch("samcli.commands.deploy.utils.click.secho") + @patch("samcli.commands.deploy.utils.click.echo") + def test_print_deploy_args_prints_parallel_upload_and_optional_sections(self, echo_mock, secho_mock): + print_deploy_args( + stack_name="stack", + s3_bucket="bucket", + image_repository={"MyFunc": "123.dkr.ecr.us-east-1.amazonaws.com/repo"}, + region="us-east-1", + capabilities=["CAPABILITY_IAM"], + parameter_overrides={"Param": "Value"}, + confirm_changeset=False, + signing_profiles={"MyFunc": {"profile_name": "pname", "profile_owner": "powner"}}, + use_changeset=True, + disable_rollback=False, + parallel_upload=True, + ) + + echo_texts = [call.args[0] for call in echo_mock.call_args_list] + self.assertTrue(any("Parallel uploads" in text for text in echo_texts)) + self.assertTrue(any("Confirm changeset" in text for text in echo_texts)) + self.assertTrue(any("Deployment image repository" in text for text in echo_texts)) + + # Basic smoke check that we printed the header/footer sections. + self.assertGreaterEqual(secho_mock.call_count, 2) + + @patch("samcli.commands.deploy.utils.click.secho") + @patch("samcli.commands.deploy.utils.click.echo") + def test_print_deploy_args_without_optional_sections(self, echo_mock, secho_mock): + print_deploy_args( + stack_name="stack", + s3_bucket="bucket", + image_repository=None, + region="us-east-1", + capabilities=["CAPABILITY_IAM"], + parameter_overrides={"Param": "Value"}, + confirm_changeset=False, + signing_profiles=None, + use_changeset=False, + disable_rollback=False, + parallel_upload=False, + ) + + echo_texts = [call.args[0] for call in echo_mock.call_args_list] + self.assertFalse(any("Confirm changeset" in text for text in echo_texts)) + self.assertFalse(any("Deployment image repository" in text for text in echo_texts)) + + def test_sanitize_parameter_overrides(self): + self.assertEqual( + {"A": "1", "B": "2"}, + sanitize_parameter_overrides({"A": {"Value": "1"}, "B": "2"}), + ) + + def test_hide_noecho_parameter_overrides(self): + template_parameters = {"Parameters": {"Secret": {"NoEcho": True}, "Visible": {"NoEcho": False}}} + overrides = {"Secret": "shh", "Visible": "ok"} + self.assertEqual( + {"Secret": "*" * 5, "Visible": "ok"}, + hide_noecho_parameter_overrides(template_parameters, overrides), + ) diff --git a/tests/unit/commands/package/test_package_context.py b/tests/unit/commands/package/test_package_context.py index 02ce71243c3..90228250792 100644 --- a/tests/unit/commands/package/test_package_context.py +++ b/tests/unit/commands/package/test_package_context.py @@ -99,6 +99,33 @@ def test_template_path_valid_with_output_template(self, patched_boto, mock_get_v ) package_command_context.run() + @patch("samcli.commands.package.package_context.ECRUploader") + @patch("samcli.commands.package.package_context.S3Uploader") + @patch.object(ResourceMetadataNormalizer, "normalize", MagicMock()) + @patch.object(Template, "export", MagicMock(return_value={})) + @patch("boto3.client") + def test_parallel_upload_disables_progress_bars(self, patched_boto, s3_uploader_mock, ecr_uploader_mock): + with tempfile.NamedTemporaryFile(mode="w", delete=False) as temp_template_file: + PackageContext( + template_file=temp_template_file.name, + s3_bucket="s3-bucket", + s3_prefix="s3-prefix", + image_repository="image-repo", + image_repositories=None, + kms_key_id="kms-key-id", + output_template_file=None, + use_json=True, + force_upload=True, + no_progressbar=False, + metadata={}, + region="us-east-2", + profile=None, + parallel_upload=True, + ).run() + + self.assertTrue(s3_uploader_mock.call_args.args[5]) + self.assertTrue(ecr_uploader_mock.call_args.args[4]) + @patch("samcli.lib.package.ecr_uploader.get_validated_container_client") @patch.object(ResourceMetadataNormalizer, "normalize", MagicMock()) @patch.object(Template, "export", MagicMock(return_value={})) @@ -4513,6 +4540,7 @@ def setUp(self): # Bypass PackageContext.__init__ — wire only the attributes _export reads. self.ctx = PackageContext.__new__(PackageContext) self.ctx.template_file = self.template_path + self.ctx.parallel_upload = False self.ctx.uploaders = self.uploaders self.ctx.code_signer = MagicMock() self.ctx.parameter_overrides = {} @@ -4552,6 +4580,7 @@ def _make_off_path_context(self): """Build a PackageContext with LE disabled and the attributes _export reads.""" ctx = PackageContext.__new__(PackageContext) ctx.template_file = "template.yaml" + ctx.parallel_upload = False ctx.uploaders = MagicMock() ctx.code_signer = MagicMock() ctx.parameter_overrides = {} diff --git a/tests/unit/commands/samconfig/test_samconfig.py b/tests/unit/commands/samconfig/test_samconfig.py index ed8abe6e25c..0136dcb3147 100644 --- a/tests/unit/commands/samconfig/test_samconfig.py +++ b/tests/unit/commands/samconfig/test_samconfig.py @@ -975,6 +975,7 @@ def test_deploy(self, do_cli_mock, template_artifacts_mock1, template_artifacts_ None, True, False, + False, "myprefix", "mykms", {"Key": "Value"}, @@ -1093,6 +1094,7 @@ def test_deploy_different_parameter_override_format( None, True, False, + False, "myprefix", "mykms", {"Key1": "Value1", "Key2": "Multiple spaces in the value"}, diff --git a/tests/unit/hook_packages/terraform/test_main.py b/tests/unit/hook_packages/terraform/test_main.py new file mode 100644 index 00000000000..52020863426 --- /dev/null +++ b/tests/unit/hook_packages/terraform/test_main.py @@ -0,0 +1,12 @@ +from unittest import TestCase +from unittest.mock import patch + +from samcli.hook_packages.terraform import main + + +class TestTerraformHookEntrypoints(TestCase): + @patch("samcli.hook_packages.terraform.main.prepare_hook") + def test_prepare_delegates_to_prepare_hook(self, prepare_hook_mock): + prepare_hook_mock.return_value = {"ok": True} + self.assertEqual({"ok": True}, main.prepare({"hello": "world"})) + prepare_hook_mock.assert_called_once_with({"hello": "world"}) diff --git a/tests/unit/lib/package/test_artifact_exporter.py b/tests/unit/lib/package/test_artifact_exporter.py index 430d775855f..e0265355985 100644 --- a/tests/unit/lib/package/test_artifact_exporter.py +++ b/tests/unit/lib/package/test_artifact_exporter.py @@ -8,14 +8,18 @@ import shutil import string import tempfile +import threading +import time import unittest import zipfile +from concurrent.futures import Future, ThreadPoolExecutor from contextlib import contextmanager, closing from pathlib import Path from typing import Optional, Dict from unittest import mock from unittest.mock import call, patch, Mock, MagicMock +from parameterized import parameterized from samcli.commands._utils.experimental import ExperimentalFlag from samcli.commands.package import exceptions from samcli.commands.package.exceptions import ExportFailedError @@ -34,6 +38,10 @@ _build_child_parameter_values, _export_global_artifacts_pass, _resolve_nested_stack_parameters, + _ThreadSafeUploadCache, + _UploadAborted, + _is_upload_abort, + get_parallel_upload_workers, ) from samcli.lib.package.language_extensions_packaging import merge_language_extensions_s3_uris from samcli.lib.intrinsic_resolver.intrinsics_symbol_table import IntrinsicsSymbolTable @@ -1226,6 +1234,9 @@ def test_export_cloudformation_stack(self, TemplateMock): parent_stack_id="id", parameter_values=mock.ANY, language_extensions_enabled=True, + parallel_upload=False, + upload_executor=None, + upload_abort=None, ) template_instance_mock.export.assert_called_once_with() self.s3_uploader_mock.upload.assert_called_once_with(mock.ANY, mock.ANY) @@ -1409,6 +1420,9 @@ def test_export_serverless_application(self, TemplateMock): parent_stack_id="id", parameter_values=mock.ANY, language_extensions_enabled=True, + parallel_upload=False, + upload_executor=None, + upload_abort=None, ) template_instance_mock.export.assert_called_once_with() self.s3_uploader_mock.upload.assert_called_once_with(mock.ANY, mock.ANY) @@ -1588,6 +1602,452 @@ def test_template_export(self, yaml_parse_mock): resource_type2_class.assert_called_once_with(self.uploaders_mock, self.code_signer_mock, None) resource_type2_instance.export.assert_called_once_with("Resource2", mock.ANY, template_dir) + @patch.object(Template, "_run_export_jobs") + @patch("samcli.lib.package.artifact_exporter.yaml_parse") + def test_template_export_parallel_invokes_executor(self, yaml_parse_mock, run_jobs_mock): + parent_dir = os.path.sep + template_dir = os.path.join(parent_dir, "foo", "bar") + template_path = os.path.join(template_dir, "path") + template_str = self.example_yaml_template() + + resource_type1_class = Mock() + resource_type1_class.RESOURCE_TYPE = "resource_type1" + resource_type1_class.ARTIFACT_TYPE = ZIP + resource_type1_class.EXPORT_DESTINATION = Destination.S3 + resource_type1_class.return_value = Mock() + + resource_type2_class = Mock() + resource_type2_class.RESOURCE_TYPE = "resource_type2" + resource_type2_class.ARTIFACT_TYPE = ZIP + resource_type2_class.EXPORT_DESTINATION = Destination.S3 + resource_type2_class.return_value = Mock() + + resources_to_export = [resource_type1_class, resource_type2_class] + properties = {"foo": "bar"} + template_dict = { + "Resources": { + "Resource1": {"Type": "resource_type1", "Properties": properties}, + "Resource2": {"Type": "resource_type2", "Properties": properties}, + } + } + + yaml_parse_mock.return_value = template_dict + with patch("samcli.lib.package.artifact_exporter.open", mock.mock_open(read_data=template_str)): + template_exporter = Template( + template_path, + parent_dir, + self.uploaders_mock, + self.code_signer_mock, + resources_to_export, + parallel_upload=True, + ) + template_exporter.export() + + run_jobs_mock.assert_called_once() + jobs, executor, abort = run_jobs_mock.call_args[0] + self.assertEqual([False, False], [is_nested_stack for is_nested_stack, _ in jobs]) + self.assertIsInstance(executor, ThreadPoolExecutor) + self.assertIsInstance(abort, threading.Event) + + @patch("samcli.lib.package.artifact_exporter.is_experimental_enabled") + @patch("samcli.lib.package.artifact_exporter.yaml_parse") + def test_template_export_parallel_wraps_cache(self, yaml_parse_mock, is_experimental_enabled_mock): + is_experimental_enabled_mock.side_effect = lambda *args: { + (ExperimentalFlag.PackagePerformance,): True, + }.get(args, False) + + parent_dir = os.path.sep + template_dir = os.path.join(parent_dir, "foo", "bar") + template_path = os.path.join(template_dir, "path") + template_str = self.example_yaml_template() + + resource_type_class = Mock() + resource_type_class.RESOURCE_TYPE = "resource_type1" + resource_type_class.ARTIFACT_TYPE = ZIP + resource_type_class.EXPORT_DESTINATION = Destination.S3 + resource_type_instance = Mock() + resource_type_class.return_value = resource_type_instance + + captured_cache = {} + + def capture_cache(uploaders, code_signer, cache): + nonlocal captured_cache + captured_cache = cache + return resource_type_instance + + resource_type_class.side_effect = capture_cache + + properties = {"foo": "bar"} + template_dict = {"Resources": {"Resource1": {"Type": "resource_type1", "Properties": properties}}} + + yaml_parse_mock.return_value = template_dict + with patch("samcli.lib.package.artifact_exporter.open", mock.mock_open(read_data=template_str)): + template_exporter = Template( + template_path, + parent_dir, + self.uploaders_mock, + self.code_signer_mock, + [resource_type_class], + parallel_upload=True, + ) + template_exporter.export() + + self.assertIsInstance(captured_cache, _ThreadSafeUploadCache) + + @patch("samcli.lib.package.artifact_exporter.yaml_parse") + def test_template_export_skips_mismatched_package_type(self, yaml_parse_mock): + parent_dir = os.path.sep + template_dir = os.path.join(parent_dir, "foo", "bar") + template_path = os.path.join(template_dir, "path") + template_str = self.example_yaml_template() + + resource_type_class = Mock() + resource_type_class.RESOURCE_TYPE = "resource_type1" + resource_type_class.ARTIFACT_TYPE = ZIP + resource_type_class.EXPORT_DESTINATION = Destination.S3 + + # Same resource type, but template says Image package type -> should be skipped + properties = {"PackageType": IMAGE, "foo": "bar"} + template_dict = {"Resources": {"Resource1": {"Type": "resource_type1", "Properties": properties}}} + yaml_parse_mock.return_value = template_dict + + with patch("samcli.lib.package.artifact_exporter.open", mock.mock_open(read_data=template_str)): + template_exporter = Template( + template_path, parent_dir, self.uploaders_mock, self.code_signer_mock, [resource_type_class] + ) + template_exporter.export() + + resource_type_class.assert_not_called() + + @patch("samcli.lib.package.artifact_exporter.yaml_parse") + def test_template_delete_skips_mismatched_package_type(self, yaml_parse_mock): + parent_dir = os.path.sep + template_dir = os.path.join(parent_dir, "foo", "bar") + template_path = os.path.join(template_dir, "path") + template_str = self.example_yaml_template() + + resource_type_class = Mock() + resource_type_class.RESOURCE_TYPE = "resource_type1" + resource_type_class.ARTIFACT_TYPE = ZIP + resource_type_class.EXPORT_DESTINATION = Destination.S3 + resource_type_instance = Mock() + resource_type_class.return_value = resource_type_instance + + properties = {"PackageType": IMAGE, "foo": "bar"} + template_dict = {"Resources": {"Resource1": {"Type": "resource_type1", "Properties": properties}}} + yaml_parse_mock.return_value = template_dict + + with patch("samcli.lib.package.artifact_exporter.open", mock.mock_open(read_data=template_str)): + template_exporter = Template( + template_path, parent_dir, self.uploaders_mock, self.code_signer_mock, [resource_type_class] + ) + template_exporter.delete(retain_resources=[]) + + resource_type_class.assert_not_called() + + def test_run_export_jobs_runs_nested_stacks_outside_the_upload_pool(self): + ran_on = {} + + def record(name): + return lambda: ran_on.setdefault(name, threading.current_thread().name) + + jobs = [(False, record("upload1")), (True, record("nested")), (False, record("upload2"))] + with ThreadPoolExecutor(max_workers=2, thread_name_prefix="upload") as executor: + Template._run_export_jobs(jobs, executor) + + self.assertTrue(ran_on["upload1"].startswith("upload")) + self.assertTrue(ran_on["upload2"].startswith("upload")) + self.assertFalse(ran_on["nested"].startswith("upload")) + + def test_run_export_jobs_runs_sibling_nested_stacks_concurrently(self): + # Each nested stack waits for the other; this only completes if they run at the same time. + barrier = threading.Barrier(2, timeout=5) + jobs = [(True, barrier.wait), (True, barrier.wait)] + with ThreadPoolExecutor(max_workers=1) as executor: + Template._run_export_jobs(jobs, executor) + + @patch("samcli.lib.package.artifact_exporter.LOG") + def test_run_export_jobs_fails_fast_and_reports_all_failures(self, log_mock): + slow_started = threading.Event() + queued_runs = [] + + def fail_first(): + slow_started.wait(5) + raise ValueError("first failure") + + def fail_slow(): + slow_started.set() + time.sleep(0.2) + raise RuntimeError("second failure") + + def queued(): + queued_runs.append(1) + time.sleep(0.05) + + jobs = [(False, fail_first), (False, fail_slow)] + [(False, queued)] * 20 + with ThreadPoolExecutor(max_workers=2) as executor: + with self.assertRaises(ValueError): + Template._run_export_jobs(jobs, executor) + + # Queued jobs are cancelled on the first failure; only jobs a free worker had already + # picked up can still run. + self.assertLess(len(queued_runs), 20) + logged = [log_call.args[1] for log_call in log_mock.error.call_args_list] + self.assertTrue(any(isinstance(error, RuntimeError) for error in logged)) + + @patch("samcli.lib.package.artifact_exporter.LOG") + def test_wait_fail_fast_surfaces_first_failure_in_submission_order(self, log_mock): + first, second = Future(), Future() + second.set_exception(RuntimeError("submitted second")) + first.set_exception(ValueError("submitted first")) + + with self.assertRaises(ValueError): + Template._wait_fail_fast([first, second]) + log_mock.error.assert_called_once() + + def test_run_export_jobs_cancels_uploads_when_nested_stack_fails(self): + release = threading.Event() + blocker_started = threading.Event() + + def blocker(): + blocker_started.set() + release.wait(5) + + queued = Mock() + + def failing_nested_stack(): + blocker_started.wait(5) + raise ValueError("nested failure") + + jobs = [(False, blocker), (False, queued), (True, failing_nested_stack)] + with ThreadPoolExecutor(max_workers=1) as executor: + with self.assertRaises(ValueError): + threading.Timer(0.2, release.set).start() + Template._run_export_jobs(jobs, executor) + + queued.assert_not_called() + + def test_export_job_passes_shared_executor_to_nested_stacks(self): + captured = {} + + class _NestedStack(CloudFormationStackResource): + def __init__(self, *args, **kwargs): + pass + + def export(self, *args, **kwargs): + captured["executor"] = self.upload_executor + + template_exporter = Template.__new__(Template) + template_exporter.uploaders = Mock() + template_exporter.code_signer = Mock() + template_exporter.parameter_values = None + template_exporter.language_extensions_enabled = False + template_exporter.template_dir = "dir" + template_exporter.parallel_upload = True + shared_executor = Mock() + + template_exporter._build_export_job(_NestedStack, "Nested", {}, None, shared_executor)() + + self.assertIs(shared_executor, captured["executor"]) + + def test_packaging_keyed_upload_cache_separates_packaging(self): + cache = _ThreadSafeUploadCache(key_by_packaging=True) + function_key = cache.cache_key("/code", "lambda-zip", None) + layer_key = cache.cache_key("/code", "zip", None) + self.assertNotEqual(function_key, layer_key) + self.assertNotEqual(function_key, cache.cache_key("/code", "lambda-zip", "jar")) + self.assertEqual("/code", _ThreadSafeUploadCache().cache_key("/code", "zip", None)) + + @patch("samcli.lib.package.artifact_exporter.is_experimental_enabled", return_value=False) + @patch("samcli.lib.package.artifact_exporter.yaml_parse") + def test_template_export_parallel_without_experimental_cache_uses_packaging_keyed_cache( + self, yaml_parse_mock, is_experimental_enabled_mock + ): + captured = {} + + def capture_cache(uploaders, code_signer, cache): + captured["cache"] = cache + return Mock() + + resource_type_class = Mock(side_effect=capture_cache) + resource_type_class.RESOURCE_TYPE = "resource_type1" + resource_type_class.ARTIFACT_TYPE = ZIP + resource_type_class.EXPORT_DESTINATION = Destination.S3 + yaml_parse_mock.return_value = {"Resources": {"Resource1": {"Type": "resource_type1", "Properties": {}}}} + + with patch("samcli.lib.package.artifact_exporter.open", mock.mock_open(read_data="")): + Template( + os.path.join(os.path.sep, "foo", "path"), + os.path.sep, + self.uploaders_mock, + self.code_signer_mock, + [resource_type_class], + parallel_upload=True, + ).export() + + self.assertIsInstance(captured["cache"], _ThreadSafeUploadCache) + self.assertNotEqual("/code", captured["cache"].cache_key("/code", "zip", None)) + + @patch("samcli.lib.package.artifact_exporter.LOG") + def test_run_export_jobs_reports_nested_stack_and_upload_failures(self, log_mock): + upload_failed = threading.Event() + + def failing_upload(): + try: + raise RuntimeError("upload failure") + finally: + upload_failed.set() + + def failing_nested_stack(): + upload_failed.wait(5) + raise ValueError("nested failure") + + with ThreadPoolExecutor(max_workers=1) as executor: + with self.assertRaises(RuntimeError): + Template._run_export_jobs([(False, failing_upload), (True, failing_nested_stack)], executor) + + logged = [log_call.args[1] for log_call in log_mock.error.call_args_list] + self.assertTrue(any(isinstance(error, ValueError) for error in logged)) + + @parameterized.expand([(None, 8), ("3", 3), ("0", 8), ("-2", 8), ("abc", 8)]) + def test_parallel_upload_workers_from_environment(self, value, expected): + env = {} if value is None else {"SAM_CLI_PARALLEL_UPLOAD_WORKERS": value} + with patch.dict(os.environ, env, clear=False): + if value is None: + os.environ.pop("SAM_CLI_PARALLEL_UPLOAD_WORKERS", None) + self.assertEqual(expected, get_parallel_upload_workers()) + + @patch("samcli.lib.package.artifact_exporter.Template") + @patch("samcli.lib.package.artifact_exporter.yaml_dump", return_value="template") + def test_nested_stack_template_upload_goes_through_shared_upload_pool(self, yaml_dump_mock, template_mock): + template_mock.return_value.export.return_value = {} + uploader = Mock() + uploader.upload.return_value = "s3://bucket/child.template" + uploader.to_path_style_s3_url.return_value = "https://s3.amazonaws.com/bucket/child.template" + uploaders = Mock() + uploaders.get.return_value = uploader + stack_resource = CloudFormationStackResource(uploaders, Mock()) + stack_resource.language_extensions_enabled = False + stack_resource.parallel_upload = True + + with ( + tempfile.TemporaryDirectory() as parent_dir, + ThreadPoolExecutor(max_workers=1, thread_name_prefix="upload") as executor, + ): + open(os.path.join(parent_dir, "child.yaml"), "w").close() + stack_resource.upload_executor = executor + upload_threads = [] + uploader.upload.side_effect = lambda *args: ( + upload_threads.append(threading.current_thread().name) or "s3://bucket/child.template" + ) + stack_resource.do_export("Child", {"TemplateURL": "child.yaml"}, parent_dir) + + self.assertEqual(1, len(upload_threads)) + self.assertTrue(upload_threads[0].startswith("upload")) + + @patch("samcli.lib.package.artifact_exporter.LOG") + def test_failure_stops_nested_subtrees_from_starting_new_uploads(self, log_mock): + abort = threading.Event() + upload_failed = threading.Event() + child_upload = Mock() + + def failing_upload(): + try: + raise RuntimeError("bucket is gone") + finally: + upload_failed.set() + + def nested_stack(): + # A nested stack that only reaches its own uploads after the root upload has failed. + upload_failed.wait(5) + abort.wait(5) + Template._run_export_jobs([(False, child_upload)], executor, abort) + + with ThreadPoolExecutor(max_workers=2) as executor: + with self.assertRaises(RuntimeError): + Template._run_export_jobs([(False, failing_upload), (True, nested_stack)], executor, abort) + + child_upload.assert_not_called() + self.assertTrue(abort.is_set()) + logged = [log_call.args[1] for log_call in log_mock.error.call_args_list] + self.assertFalse(any(isinstance(error, _UploadAborted) for error in logged)) + + def test_wait_fail_fast_surfaces_real_failure_over_aborted_work(self): + aborted, real = Future(), Future() + aborted.set_exception(_UploadAborted()) + real.set_exception(ValueError("real failure")) + + with self.assertRaises(ValueError): + Template._wait_fail_fast([aborted, real], threading.Event()) + + def test_export_job_is_skipped_once_export_is_aborted(self): + exporter = Mock() + template_exporter = Template.__new__(Template) + template_exporter.uploaders = Mock() + template_exporter.code_signer = Mock() + template_exporter.parameter_values = None + template_exporter.language_extensions_enabled = False + template_exporter.template_dir = "dir" + template_exporter.parallel_upload = True + abort = threading.Event() + abort.set() + + job = template_exporter._build_export_job(Mock(return_value=exporter), "Leaf", {}, None, Mock(), abort) + + with self.assertRaises(_UploadAborted): + job() + exporter.export.assert_not_called() + + @patch("samcli.lib.package.artifact_exporter.LOG") + def test_other_failures_log_one_line_with_traceback_only_at_debug(self, log_mock): + other = Future() + other.set_exception(RuntimeError("other failure")) + + Template._log_other_failures([other]) + + self.assertNotIn("exc_info", log_mock.error.call_args.kwargs) + self.assertIs(other.exception(), log_mock.debug.call_args.kwargs["exc_info"]) + + @patch("samcli.lib.package.artifact_exporter.LOG") + def test_wait_fail_fast_ignores_aborts_wrapped_by_nested_stack_exporters(self, log_mock): + # Sibling nested stack A was skipped after B failed. ResourceZip.export wraps both in + # ExportFailedError, once per nesting level; A comes first in submission order. + aborted_a, failed_b = Future(), Future() + aborted_a.set_exception( + exceptions.ExportFailedError( + resource_id="Outer", + property_name="TemplateURL", + property_value="outer.yaml", + ex=exceptions.ExportFailedError( + resource_id="A", property_name="TemplateURL", property_value="a.yaml", ex=_UploadAborted() + ), + ) + ) + real_error = exceptions.ExportFailedError( + resource_id="B", property_name="TemplateURL", property_value="b.yaml", ex=RuntimeError("bucket is gone") + ) + failed_b.set_exception(real_error) + + with self.assertRaises(exceptions.ExportFailedError) as raised: + Template._wait_fail_fast([aborted_a, failed_b], threading.Event()) + + self.assertIs(real_error, raised.exception) + log_mock.error.assert_not_called() + + def test_is_upload_abort_follows_wrapping(self): + self.assertTrue(_is_upload_abort(_UploadAborted())) + wrapped = exceptions.ExportFailedError( + resource_id="A", property_name="TemplateURL", property_value="a.yaml", ex=_UploadAborted() + ) + self.assertTrue(_is_upload_abort(wrapped)) + self.assertFalse(_is_upload_abort(RuntimeError("real"))) + self.assertFalse(_is_upload_abort(None)) + + def test_thread_safe_upload_cache_key_lock_is_per_key(self): + cache = _ThreadSafeUploadCache() + self.assertIs(cache.key_lock("a"), cache.key_lock("a")) + self.assertIsNot(cache.key_lock("a"), cache.key_lock("b")) + @patch("samcli.lib.package.artifact_exporter.is_experimental_enabled") @patch("samcli.lib.package.artifact_exporter.yaml_parse") def test_template_export_with_experimental_flag(self, yaml_parse_mock, is_experimental_enabled_mock): diff --git a/tests/unit/lib/package/test_stream_cursor_utils.py b/tests/unit/lib/package/test_stream_cursor_utils.py index dc4c66266de..b38c8bee4d6 100644 --- a/tests/unit/lib/package/test_stream_cursor_utils.py +++ b/tests/unit/lib/package/test_stream_cursor_utils.py @@ -1,6 +1,7 @@ from unittest import TestCase from samcli.lib.package.stream_cursor_utils import ( + CursorFormatter, CursorUpFormatter, CursorDownFormatter, CursorLeftFormatter, @@ -14,3 +15,6 @@ def test_cursor_utils(self): self.assertEqual(CursorDownFormatter().cursor_format(count=1), "\x1b[1B") self.assertEqual(CursorLeftFormatter().cursor_format(), "\x1b[0G") self.assertEqual(ClearLineFormatter().cursor_format(), "\x1b[0K") + + def test_base_formatter_is_noop(self): + self.assertIsNone(CursorFormatter().cursor_format(count=0)) diff --git a/tests/unit/lib/package/test_utils.py b/tests/unit/lib/package/test_utils.py index 073802f5314..6eff420db48 100644 --- a/tests/unit/lib/package/test_utils.py +++ b/tests/unit/lib/package/test_utils.py @@ -1,5 +1,8 @@ import tempfile +import threading +import time from unittest import TestCase +from unittest.mock import Mock, patch from parameterized import parameterized @@ -61,3 +64,87 @@ def test_zip_folder_uses_different_path_for_same_file_in_different_run(self): previous_md5_hash = md5_hash else: self.assertEqual(previous_md5_hash, md5_hash) + + def test_upload_local_artifacts_uploads_shared_path_once_across_threads(self): + from samcli.lib.package.artifact_exporter import _ThreadSafeUploadCache + + cache = _ThreadSafeUploadCache({}) + calls = [] + + def slow_zip_and_upload(local_path, *args, **kwargs): + calls.append(local_path) + time.sleep(0.1) + return "s3://bucket/key" + + with ( + tempfile.TemporaryDirectory() as folder, + patch.object(utils, "zip_and_upload", side_effect=slow_zip_and_upload), + ): + results = [] + + def upload(): + results.append( + utils.upload_local_artifacts( + resource_type="AWS::Serverless::Function", + resource_id="Function", + resource_dict={"CodeUri": folder}, + property_path="CodeUri", + parent_dir=folder, + uploader=Mock(), + previously_uploaded=cache, + ) + ) + + threads = [threading.Thread(target=upload) for _ in range(4)] + for thread in threads: + thread.start() + for thread in threads: + thread.join() + + self.assertEqual(1, len(calls)) + self.assertEqual(["s3://bucket/key"] * 4, results) + + def test_upload_local_artifacts_memoizes_by_packaging_across_threads(self): + from samcli.lib.package.artifact_exporter import _ThreadSafeUploadCache + + cache = _ThreadSafeUploadCache(key_by_packaging=True) + calls = [] + + def slow_zip_and_upload(local_path, uploader, extension, zip_method): + calls.append(zip_method) + time.sleep(0.05) + return f"s3://bucket/{len(calls)}" + + with ( + tempfile.TemporaryDirectory() as folder, + patch.object(utils, "zip_and_upload", side_effect=slow_zip_and_upload), + ): + + def upload(resource_type): + return utils.upload_local_artifacts( + resource_type=resource_type, + resource_id="Resource", + resource_dict={"CodeUri": folder}, + property_path="CodeUri", + parent_dir=folder, + uploader=Mock(), + previously_uploaded=cache, + ) + + results = [] + threads = [ + threading.Thread(target=lambda: results.append(upload("AWS::Serverless::Function"))) for _ in range(4) + ] + for thread in threads: + thread.start() + for thread in threads: + thread.join() + bundle_url = upload("AWS::ElasticBeanstalk::ApplicationVersion") + + # Four functions sharing a directory zip and upload it once; a non-Lambda resource using the + # same directory is packaged differently, so it gets its own upload, not the functions' zip. + self.assertEqual(1, len(set(results))) + self.assertEqual(2, len(calls)) + self.assertIs(utils.make_zip_with_lambda_permissions, calls[0]) + self.assertIs(utils.make_zip, calls[1]) + self.assertNotIn(bundle_url, results) diff --git a/tests/unit/lib/utils/test_profile.py b/tests/unit/lib/utils/test_profile.py new file mode 100644 index 00000000000..0c9810fce53 --- /dev/null +++ b/tests/unit/lib/utils/test_profile.py @@ -0,0 +1,11 @@ +from unittest import TestCase +from unittest.mock import patch + +from samcli.lib.utils.profile import list_available_profiles + + +class TestProfileUtils(TestCase): + @patch("samcli.lib.utils.profile.Session") + def test_list_available_profiles(self, session_mock): + session_mock.return_value.available_profiles = ["p1", "p2"] + self.assertEqual(["p1", "p2"], list_available_profiles())