From 740958005a36b3fc14891b5bc9d54492237ecb70 Mon Sep 17 00:00:00 2001 From: Ivan Cheung Date: Fri, 12 Nov 2021 18:31:36 -0500 Subject: [PATCH] Parallel custom job notebook execution (#120) * Initial commit * Added parallel execution --- .cloud-build/ExecuteChangedNotebooks.py | 229 ++++++++++++------ .cloud-build/ExecuteNotebook.py | 173 ++++--------- .cloud-build/execute_notebook_cli.py | 38 +++ .cloud-build/execute_notebook_remote.py | 72 ++++++ ...book-execution-test-cloudbuild-single.yaml | 26 ++ .../notebook-execution-test-cloudbuild.yaml | 18 +- .cloud-build/requirements.txt | 4 + .../{ => utils}/NotebookProcessors.py | 4 +- .../{ => utils}/UpdateNotebookVariables.py | 0 .cloud-build/utils/__init__.py | 0 .cloud-build/utils/util.py | 60 +++++ 11 files changed, 411 insertions(+), 213 deletions(-) create mode 100644 .cloud-build/execute_notebook_cli.py create mode 100644 .cloud-build/execute_notebook_remote.py create mode 100644 .cloud-build/notebook-execution-test-cloudbuild-single.yaml rename .cloud-build/{ => utils}/NotebookProcessors.py (94%) rename .cloud-build/{ => utils}/UpdateNotebookVariables.py (100%) create mode 100644 .cloud-build/utils/__init__.py create mode 100644 .cloud-build/utils/util.py diff --git a/.cloud-build/ExecuteChangedNotebooks.py b/.cloud-build/ExecuteChangedNotebooks.py index 8ad72a040..9a09a7da4 100755 --- a/.cloud-build/ExecuteChangedNotebooks.py +++ b/.cloud-build/ExecuteChangedNotebooks.py @@ -14,18 +14,22 @@ # limitations under the License. import argparse +import concurrent import dataclasses import datetime import functools -import pathlib import os +import pathlib +import nbformat +import re import subprocess -from pathlib import Path from typing import List, Optional -import concurrent from tabulate import tabulate +import operator -import ExecuteNotebook +import execute_notebook_remote +from utils import util, NotebookProcessors +from google.cloud.devtools.cloudbuild_v1.types import BuildOperationMetadata def str2bool(v): @@ -62,49 +66,135 @@ def format_timedelta(delta: datetime.timedelta) -> str: @dataclasses.dataclass class NotebookExecutionResult: - notebook: str + name: str duration: datetime.timedelta is_pass: bool + log_url: str + output_uri: str + build_id: str error_message: Optional[str] +def _process_notebook( + notebook_path: str, + variable_project_id: str, + variable_region: str, +): + # Read notebook + with open(notebook_path) as f: + nb = nbformat.read(f, as_version=4) + + # Create preprocessors + remove_no_execute_cells_preprocessor = NotebookProcessors.RemoveNoExecuteCells() + update_variables_preprocessor = NotebookProcessors.UpdateVariablesPreprocessor( + replacement_map={ + "PROJECT_ID": variable_project_id, + "REGION": variable_region, + }, + ) + + # Use no-execute preprocessor + ( + nb, + resources, + ) = remove_no_execute_cells_preprocessor.preprocess(nb) + + (nb, resources) = update_variables_preprocessor.preprocess(nb, resources) + + with open(notebook_path, mode="w", encoding="utf-8") as new_file: + nbformat.write(nb, new_file) + + +def _create_tag(filepath: str) -> str: + tag = os.path.basename(os.path.normpath(filepath)) + tag = re.sub("[^0-9a-zA-Z_.-]+", "-", tag) + + if tag.startswith(".") or tag.startswith("-"): + tag = tag[1:] + + return tag + + def execute_notebook( - artifacts_path: str, + container_uri: str, + staging_bucket: str, + artifacts_bucket: str, variable_project_id: str, variable_region: str, - should_log_output: bool, - should_use_new_kernel: bool, notebook: str, + should_get_tail_logs: bool = False, ) -> NotebookExecutionResult: print(f"Running notebook: {notebook}") + # Create paths + notebook_output_uri = "/".join([artifacts_bucket, pathlib.Path(notebook).name]) + + # Create tag from notebook + tag = _create_tag(filepath=notebook) + result = NotebookExecutionResult( - notebook=notebook, + name=tag, duration=datetime.timedelta(seconds=0), is_pass=False, + output_uri=notebook_output_uri, + log_url="", + build_id="", error_message=None, ) # TODO: Handle cases where multiple notebooks have the same name time_start = datetime.datetime.now() + operation = None try: - ExecuteNotebook.execute_notebook( - notebook_file_path=notebook, - output_file_folder=artifacts_path, - replacement_map={ - "PROJECT_ID": variable_project_id, - "REGION": variable_region, - }, - should_log_output=should_log_output, - should_use_new_kernel=should_use_new_kernel, + # Pre-process notebook by substituting variable names + _process_notebook( + notebook_path=notebook, + variable_project_id=variable_project_id, + variable_region=variable_region, ) + + # Upload the pre-processed code to a GCS bucket + code_archive_uri = util.archive_code_and_upload(staging_bucket=staging_bucket) + + operation = execute_notebook_remote.execute_notebook_remote( + code_archive_uri=code_archive_uri, + notebook_uri=notebook, + notebook_output_uri=notebook_output_uri, + container_uri=container_uri, + tag=tag, + ) + + operation_metadata = BuildOperationMetadata(mapping=operation.metadata) + result.build_id = operation_metadata.build.id + result.log_url = operation_metadata.build.log_url + + # Block and wait for the result + operation_result = operation.result() + result.duration = datetime.datetime.now() - time_start result.is_pass = True print(f"{notebook} PASSED in {format_timedelta(result.duration)}.") except Exception as error: + result.error_message = str(error) + + if operation and should_get_tail_logs: + # Extract the logs + logs_bucket = operation_metadata.build.logs_bucket + + # Download tail end of logs file + log_file_uri = f"{logs_bucket}/log-{result.build_id}.txt" + + # Use gcloud to get tail + try: + result.error_message = subprocess.check_output( + ["gsutil", "cat", "-r", "-1000", log_file_uri], encoding="UTF-8" + ) + except Exception as error: + result.error_message = str(error) + result.duration = datetime.datetime.now() - time_start result.is_pass = False - result.error_message = str(error) + print( f"{notebook} FAILED in {format_timedelta(result.duration)}: {result.error_message}" ) @@ -114,18 +204,19 @@ def execute_notebook( def run_changed_notebooks( test_paths_file: str, - base_branch: Optional[str], - output_folder: str, + container_uri: str, + staging_bucket: str, + artifacts_bucket: str, variable_project_id: str, variable_region: str, should_parallelize: bool, - should_use_separate_kernels: bool, + base_branch: Optional[str] = None, ): """ Run the notebooks that exist under the folders defined in the test_paths_file. It only runs notebooks that have differences from the Git base_branch. - The executed notebooks are saved in the output_folder. + The executed notebooks are saved in the artifacts_bucket. Variables are also injected into the notebooks such as the variable_project_id and variable_region. @@ -136,19 +227,16 @@ def run_changed_notebooks( base_branch (str): Optional. If provided, only the files that have changed from the base_branch will be checked. If not provided, all files will be checked. - output_folder (str): - Required. The folder to write executed notebooks to. + staging_bucket (str): + Required. The GCS staging bucket to write source code to. + artifacts_bucket (str): + Required. The GCS staging bucket to write executed notebooks to. variable_project_id (str): Required. The value for PROJECT_ID to inject into notebooks. variable_region (str): Required. The value for REGION to inject into notebooks. should_parallelize (bool): Required. Should run notebooks in parallel using a thread pool as opposed to in sequence. - should_use_separate_kernels (bool): - Note: Dependencies don't install correctly when this is set to True - See https://github.com/nteract/papermill/issues/625 - - Required. Should run each notebook in a separate and independent virtual environment. """ test_paths = [] @@ -176,13 +264,7 @@ def run_changed_notebooks( notebooks = notebooks.decode("utf-8").split("\n") notebooks = [notebook for notebook in notebooks if notebook.endswith(".ipynb")] notebooks = [notebook for notebook in notebooks if len(notebook) > 0] - notebooks = [notebook for notebook in notebooks if Path(notebook).exists()] - - # Create paths - artifacts_path = Path(output_folder) - artifacts_path.mkdir(parents=True, exist_ok=True) - artifacts_path.joinpath("success").mkdir(parents=True, exist_ok=True) - artifacts_path.joinpath("failure").mkdir(parents=True, exist_ok=True) + notebooks = [notebook for notebook in notebooks if pathlib.Path(notebook).exists()] notebook_execution_results: List[NotebookExecutionResult] = [] @@ -198,11 +280,11 @@ def run_changed_notebooks( executor.map( functools.partial( execute_notebook, - artifacts_path, + container_uri, + staging_bucket, + artifacts_bucket, variable_project_id, variable_region, - False, - should_use_separate_kernels, ), notebooks, ) @@ -210,12 +292,12 @@ def run_changed_notebooks( else: notebook_execution_results = [ execute_notebook( - artifacts_path=artifacts_path, + container_uri=container_uri, + staging_bucket=staging_bucket, + artifacts_bucket=artifacts_bucket, variable_project_id=variable_project_id, variable_region=variable_region, notebook=notebook, - should_log_output=True, - should_use_new_kernel=should_use_separate_kernels, ) for notebook in notebooks ] @@ -224,29 +306,40 @@ def run_changed_notebooks( print("\n=== RESULTS ===\n") - notebooks_sorted = sorted( + results_sorted = sorted( notebook_execution_results, key=lambda result: result.is_pass, reverse=True, ) + # Print results print( tabulate( [ [ - os.path.basename(os.path.normpath(result.notebook)), + result.name, "PASSED" if result.is_pass else "FAILED", format_timedelta(result.duration), - result.error_message or "--", + result.log_url, ] - for result in notebooks_sorted + for result in results_sorted ], - headers=["file", "status", "duration", "error"], + headers=["build_tag", "status", "duration", "log_url"], ) ) print("\n=== END RESULTS===\n") + total_notebook_duration = functools.reduce( + operator.add, [result.duration for result in results_sorted] + ) + + print(f"Cumulative notebook duration: {format_timedelta(total_notebook_duration)}") + + # Raise error if any notebooks failed + if not all([result.is_pass for result in results_sorted]): + raise RuntimeError("Notebook failures detected. See logs for details") + parser = argparse.ArgumentParser(description="Run changed notebooks.") parser.add_argument( @@ -261,9 +354,9 @@ parser.add_argument( required=False, ) parser.add_argument( - "--output_folder", - type=pathlib.Path, - help="The path to the folder to store executed notebooks.", + "--container_uri", + type=str, + help="The container uri to run each notebook in.", required=True, ) parser.add_argument( @@ -278,35 +371,35 @@ parser.add_argument( help="The GCP region. This is used to inject a variable value into the notebook before running.", required=True, ) - -# Note: Dependencies don't install correctly when this is set to True +parser.add_argument( + "--staging_bucket", + type=str, + help="The GCP directory for staging temporary files.", + required=True, +) +parser.add_argument( + "--artifacts_bucket", + type=str, + help="The GCP directory for storing executed notebooks.", + required=True, +) parser.add_argument( "--should_parallelize", type=str2bool, nargs="?", const=True, - default=False, + default=True, help="Should run notebooks in parallel.", ) -# Note: This isn't guaranteed to work correctly due to existing Papermill issue -# See https://github.com/nteract/papermill/issues/625 -parser.add_argument( - "--should_use_separate_kernels", - type=str2bool, - nargs="?", - const=True, - default=False, - help="(Experimental) Should run each notebook in a separate and independent virtual environment.", -) - args = parser.parse_args() run_changed_notebooks( test_paths_file=args.test_paths_file, - base_branch=args.base_branch, - output_folder=args.output_folder, + container_uri=args.container_uri, + staging_bucket=args.staging_bucket, + artifacts_bucket=args.artifacts_bucket, variable_project_id=args.variable_project_id, variable_region=args.variable_region, should_parallelize=args.should_parallelize, - should_use_separate_kernels=args.should_use_separate_kernels, + base_branch=args.base_branch, ) diff --git a/.cloud-build/ExecuteNotebook.py b/.cloud-build/ExecuteNotebook.py index 029be3750..189549e7b 100644 --- a/.cloud-build/ExecuteNotebook.py +++ b/.cloud-build/ExecuteNotebook.py @@ -13,163 +13,76 @@ # See the License for the specific language governing permissions and # limitations under the License. -import json import sys -import nbformat import os import errno -from NotebookProcessors import RemoveNoExecuteCells, UpdateVariablesPreprocessor -from typing import Dict, Tuple import papermill as pm import shutil -import virtualenv -import uuid -from jupyter_client.kernelspecapp import KernelSpecManager + +from utils import util +from google.cloud.aiplatform import utils # This script is used to execute a notebook and write out the output notebook. -# The replaces calling the nbconvert via command-line, which doesn't write the output notebook correctly when there are errors during execution. - -STAGING_FOLDER = "staging" -ENVIRONMENTS_PATH = "environments" -KERNELS_SPECS_PATH = "kernel_specs" - - -def create_and_install_kernel() -> Tuple[str, str]: - # Create environment - kernel_name = str(uuid.uuid4()) - env_name = f"{ENVIRONMENTS_PATH}/{kernel_name}" - # venv.create(env_name, system_site_packages=True, with_pip=True) - virtualenv.cli_run([env_name, "--system-site-packages"]) - - # Create kernel spec - kernel_spec = { - "argv": [ - f"{env_name}/bin/python", - "-m", - "ipykernel_launcher", - "-f", - "{connection_file}", - ], - "display_name": "Python 3", - "language": "python", - } - kernel_spec_folder = os.path.join(KERNELS_SPECS_PATH, kernel_name) - kernel_spec_file = os.path.join(kernel_spec_folder, "kernel.json") - - # Create kernel spec folder - if not os.path.exists(os.path.dirname(kernel_spec_file)): - try: - os.makedirs(os.path.dirname(kernel_spec_file)) - except OSError as exc: # Guard against race condition - if exc.errno != errno.EEXIST: - raise - - with open(kernel_spec_file, mode="w", encoding="utf-8") as f: - json.dump(kernel_spec, f) - - # Install kernel - kernel_spec_manager = KernelSpecManager() - kernel_spec_manager.install_kernel_spec( - source_dir=kernel_spec_folder, kernel_name=kernel_name - ) - - return kernel_name, env_name def execute_notebook( - notebook_file_path: str, - output_file_folder: str, - replacement_map: Dict[str, str], + notebook_source: str, + output_file_or_uri: str, should_log_output: bool, - should_use_new_kernel: bool, ): - # Create staging directory if it doesn't exist - staging_file_path = f"{STAGING_FOLDER}/{notebook_file_path}" - if not os.path.exists(os.path.dirname(staging_file_path)): - try: - os.makedirs(os.path.dirname(staging_file_path)) - except OSError as exc: # Guard against race condition - if exc.errno != errno.EEXIST: - raise + file_name = os.path.basename(os.path.normpath(notebook_source)) - file_name = os.path.basename(os.path.normpath(notebook_file_path)) + # Download notebook if it's a GCS URI + if notebook_source.startswith("gs://"): + # Extract uri components + bucket_name, prefix = utils.extract_bucket_and_prefix_from_gcs_path( + notebook_source + ) - # Create environments folder - if not os.path.exists(ENVIRONMENTS_PATH): - try: - os.makedirs(ENVIRONMENTS_PATH) - except OSError as exc: # Guard against race condition - if exc.errno != errno.EEXIST: - raise + # Download remote notebook to local file system + notebook_source = file_name + util.download_file( + bucket_name=bucket_name, blob_name=prefix, destination_file=notebook_source + ) - # Create and install kernel - kernel_name = next( - iter(KernelSpecManager().find_kernel_specs().keys()), None - ) # Find first existing kernel and use as default - env_name = None - if should_use_new_kernel: - kernel_name, env_name = create_and_install_kernel() - - # Read notebook - with open(notebook_file_path) as f: - nb = nbformat.read(f, as_version=4) - - has_error = False + execution_exception = None # Execute notebook try: - # Create preprocessors - remove_no_execute_cells_preprocessor = RemoveNoExecuteCells() - update_variables_preprocessor = UpdateVariablesPreprocessor( - replacement_map=replacement_map - ) - - # Use no-execute preprocessor - ( - nb, - resources, - ) = remove_no_execute_cells_preprocessor.preprocess(nb) - - (nb, resources) = update_variables_preprocessor.preprocess(nb, resources) - - # print(f"Staging modified notebook to: {staging_file_path}") - with open(staging_file_path, mode="w", encoding="utf-8") as f: - nbformat.write(nb, f) - # Execute notebook pm.execute_notebook( - input_path=staging_file_path, - output_path=staging_file_path, - kernel_name=kernel_name, + input_path=notebook_source, + output_path=notebook_source, progress_bar=should_log_output, request_save_on_cell_execute=should_log_output, log_output=should_log_output, stdout_file=sys.stdout if should_log_output else None, stderr_file=sys.stderr if should_log_output else None, ) - except Exception: - # print(f"Error executing the notebook: {notebook_file_path}.\n\n") - has_error = True - - raise - + except Exception as exception: + execution_exception = exception finally: - # Clear env - if env_name is not None: - shutil.rmtree(path=env_name) + # Copy executed notebook + if output_file_or_uri.startswith("gs://"): + # Upload to GCS path + util.upload_file(notebook_source, remote_file_path=output_file_or_uri) - # Copy execute notebook - output_file_path = os.path.join( - output_file_folder, "failure" if has_error else "success", file_name - ) + print("\n=== EXECUTION FINISHED ===\n") + print( + f"Please debug the executed notebook by downloading: {output_file_or_uri}" + ) + print("\n======\n") + else: + # Create directories if they don't exist + if not os.path.exists(os.path.dirname(output_file_or_uri)): + try: + os.makedirs(os.path.dirname(output_file_or_uri)) + except OSError as exc: # Guard against race condition + if exc.errno != errno.EEXIST: + raise - # Create directories if they don't exist - if not os.path.exists(os.path.dirname(output_file_path)): - try: - os.makedirs(os.path.dirname(output_file_path)) - except OSError as exc: # Guard against race condition - if exc.errno != errno.EEXIST: - raise + print(f"Writing output to: {output_file_or_uri}") + shutil.move(notebook_source, output_file_or_uri) - # print(f"Writing output to: {output_file_path}") - shutil.move(staging_file_path, output_file_path) + if execution_exception: + raise execution_exception diff --git a/.cloud-build/execute_notebook_cli.py b/.cloud-build/execute_notebook_cli.py new file mode 100644 index 000000000..6cd999225 --- /dev/null +++ b/.cloud-build/execute_notebook_cli.py @@ -0,0 +1,38 @@ +#!/usr/bin/env python +# Copyright 2021 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +import argparse +import ExecuteNotebook + +parser = argparse.ArgumentParser(description="Run changed notebooks.") +parser.add_argument( + "--notebook_source", + type=str, + help="Local filepath or GCS URI to notebook.", + required=True, +) +parser.add_argument( + "--output_file_or_uri", + type=str, + help="Local file or GCS URI to save executed notebook to.", + required=True, +) + +args = parser.parse_args() +ExecuteNotebook.execute_notebook( + notebook_source=args.notebook_source, + output_file_or_uri=args.output_file_or_uri, + should_log_output=True, +) diff --git a/.cloud-build/execute_notebook_remote.py b/.cloud-build/execute_notebook_remote.py new file mode 100644 index 000000000..98123e107 --- /dev/null +++ b/.cloud-build/execute_notebook_remote.py @@ -0,0 +1,72 @@ +from google.protobuf import duration_pb2 +from yaml.loader import FullLoader + +import google.auth +from google.cloud.devtools import cloudbuild_v1 +from google.cloud.devtools.cloudbuild_v1.types import Source, StorageSource + +from typing import Optional +import yaml + +from google.cloud.aiplatform import utils +from google.api_core import operation + +CLOUD_BUILD_FILEPATH = ".cloud-build/notebook-execution-test-cloudbuild-single.yaml" +TIMEOUT_IN_SECONDS = 86400 + + +def execute_notebook_remote( + code_archive_uri: str, + notebook_uri: str, + notebook_output_uri: str, + container_uri: str, + tag: Optional[str], +) -> operation.Operation: + """Create and execute a simple Google Cloud Build configuration, + print the in-progress status and print the completed status.""" + + # Authorize the client with Google defaults + credentials, project_id = google.auth.default() + client = cloudbuild_v1.services.cloud_build.CloudBuildClient() + + build = cloudbuild_v1.Build() + + # The following build steps will output "hello world" + # For more information on build configuration, see + # https://cloud.google.com/build/docs/configuring-builds/create-basic-configuration + cloudbuild_config = yaml.load(open(CLOUD_BUILD_FILEPATH), Loader=FullLoader) + + substitutions = { + "_PYTHON_IMAGE": container_uri, + "_NOTEBOOK_GCS_URI": notebook_uri, + "_NOTEBOOK_OUTPUT_GCS_URI": notebook_output_uri, + } + + ( + source_archived_file_gcs_bucket, + source_archived_file_gcs_object, + ) = utils.extract_bucket_and_prefix_from_gcs_path(code_archive_uri) + + build.source = Source( + storage_source=StorageSource( + bucket=source_archived_file_gcs_bucket, + object_=source_archived_file_gcs_object, + ) + ) + + build.steps = cloudbuild_config["steps"] + build.substitutions = substitutions + build.timeout = duration_pb2.Duration(seconds=TIMEOUT_IN_SECONDS) + build.queue_ttl = duration_pb2.Duration(seconds=TIMEOUT_IN_SECONDS) + + if tag: + build.tags = [tag] + + operation = client.create_build(project_id=project_id, build=build) + # Print the in-progress operation + # print("IN PROGRESS:") + # print(operation.metadata) + + # Print the completed status + # print("RESULT:", result.status) + return operation diff --git a/.cloud-build/notebook-execution-test-cloudbuild-single.yaml b/.cloud-build/notebook-execution-test-cloudbuild-single.yaml new file mode 100644 index 000000000..9dc1b2032 --- /dev/null +++ b/.cloud-build/notebook-execution-test-cloudbuild-single.yaml @@ -0,0 +1,26 @@ +steps: + # Show the gcloud info and check if gcloud exists + - name: ${_PYTHON_IMAGE} + entrypoint: /bin/sh + args: + - -c + - 'gcloud config list' + # Check the Python version + - name: ${_PYTHON_IMAGE} + entrypoint: /bin/sh + args: + - -c + - 'python3 .cloud-build/CheckPythonVersion.py' + # Install Python dependencies + - name: ${_PYTHON_IMAGE} + entrypoint: pip + args: ['install', '--upgrade', '--user', '--requirement', '.cloud-build/requirements.txt'] + # Install Python dependencies and run testing script + - name: ${_PYTHON_IMAGE} + entrypoint: /bin/sh + args: + - -c + - 'python3 -m pip freeze && python3 .cloud-build/execute_notebook_cli.py --notebook_source "${_NOTEBOOK_GCS_URI}" --output_file_or_uri "${_NOTEBOOK_OUTPUT_GCS_URI}"' + env: + - 'IS_TESTING=1' +timeout: 86400s diff --git a/.cloud-build/notebook-execution-test-cloudbuild.yaml b/.cloud-build/notebook-execution-test-cloudbuild.yaml index 2a8ffea79..b3fa12a27 100644 --- a/.cloud-build/notebook-execution-test-cloudbuild.yaml +++ b/.cloud-build/notebook-execution-test-cloudbuild.yaml @@ -5,6 +5,10 @@ steps: args: - -c - 'gcloud config list' + # # Clone the Git repo + # - name: ${_PYTHON_IMAGE} + # entrypoint: git + # args: ['clone', "${_GIT_REPO}", "--branch", "${_GIT_BRANCH_NAME}", "."] # Check the Python version - name: ${_PYTHON_IMAGE} entrypoint: /bin/sh @@ -26,19 +30,7 @@ steps: entrypoint: /bin/sh args: - -c - - 'python3 -m pip freeze && python3 .cloud-build/ExecuteChangedNotebooks.py --test_paths_file "${_TEST_PATHS_FILE}" --base_branch "${_FORCED_BASE_BRANCH}" --output_folder ${BUILD_ID} --variable_project_id ${PROJECT_ID} --variable_region ${_GCP_REGION}' + - 'python3 -m pip freeze && python3 .cloud-build/ExecuteChangedNotebooks.py --test_paths_file "${_TEST_PATHS_FILE}" --base_branch "${_FORCED_BASE_BRANCH}" --container_uri ${_PYTHON_IMAGE} --staging_bucket ${_GCS_STAGING_BUCKET} --artifacts_bucket ${_GCS_STAGING_BUCKET}/executed_notebooks/PR_${_PR_NUMBER}/BUILD_${BUILD_ID} --variable_project_id ${PROJECT_ID} --variable_region ${_GCP_REGION}' env: - 'IS_TESTING=1' - # Manually copy artifacts to GCS - - name: gcr.io/cloud-builders/gsutil - entrypoint: /bin/sh - args: - - -c - - 'if [ $(ls -pR "/workspace/${BUILD_ID}" | grep -v / | grep -v ^$ | wc -l) -ne 0 ]; then gsutil -m -q rsync -r "/workspace/${BUILD_ID}" "gs://${_GCS_ARTIFACTS_BUCKET}/test-artifacts/PR_${_PR_NUMBER}/BUILD_${BUILD_ID}/"; else echo "No artifacts to copy."; fi' - # Fail if there is anything in the failure folder - - name: ${_PYTHON_IMAGE} - entrypoint: /bin/sh - args: - - -c - - 'echo "Download executed notebooks with this command: \"mkdir -p artifacts && gsutil rsync -r gs://${_GCS_ARTIFACTS_BUCKET}/test-artifacts/PR_${_PR_NUMBER}/BUILD_${BUILD_ID} artifacts/\"" && if [ "$(ls -A /workspace/${BUILD_ID}/failure | wc -l)" -ne 0 ]; then exit 1; else exit 0; fi' timeout: 86400s diff --git a/.cloud-build/requirements.txt b/.cloud-build/requirements.txt index f803b8033..c2cad3439 100644 --- a/.cloud-build/requirements.txt +++ b/.cloud-build/requirements.txt @@ -6,3 +6,7 @@ numpy>=1.19 pandas>=1.2 matplotlib>=3.4 tabulate +google-cloud-aiplatform +google-cloud-storage +google-cloud-build +gcloud diff --git a/.cloud-build/NotebookProcessors.py b/.cloud-build/utils/NotebookProcessors.py similarity index 94% rename from .cloud-build/NotebookProcessors.py rename to .cloud-build/utils/NotebookProcessors.py index 22e184aae..04f10906c 100644 --- a/.cloud-build/NotebookProcessors.py +++ b/.cloud-build/utils/NotebookProcessors.py @@ -15,7 +15,7 @@ from nbconvert.preprocessors import Preprocessor from typing import Dict -import UpdateNotebookVariables +from . import UpdateNotebookVariables as update_notebook_variables class RemoveNoExecuteCells(Preprocessor): @@ -41,7 +41,7 @@ class UpdateVariablesPreprocessor(Preprocessor): # VARIABLE_NAME = '[description]' for variable_name, variable_value in replacement_map.items(): - content = UpdateNotebookVariables.get_updated_value( + content = update_notebook_variables.get_updated_value( content=content, variable_name=variable_name, variable_value=variable_value, diff --git a/.cloud-build/UpdateNotebookVariables.py b/.cloud-build/utils/UpdateNotebookVariables.py similarity index 100% rename from .cloud-build/UpdateNotebookVariables.py rename to .cloud-build/utils/UpdateNotebookVariables.py diff --git a/.cloud-build/utils/__init__.py b/.cloud-build/utils/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/.cloud-build/utils/util.py b/.cloud-build/utils/util.py new file mode 100644 index 000000000..5f235189e --- /dev/null +++ b/.cloud-build/utils/util.py @@ -0,0 +1,60 @@ +from datetime import datetime +from typing import Optional +from google.cloud import storage +from google.cloud.aiplatform import utils +from google.auth import credentials as auth_credentials +import os + +import subprocess +import tarfile +import uuid + + +def download_file(bucket_name: str, blob_name: str, destination_file: str) -> str: + """Copies a remote GCS file to a local path.""" + remote_file_path = "".join(["gs://", "/".join([bucket_name, blob_name])]) + + subprocess.check_output( + ["gsutil", "cp", remote_file_path, destination_file], encoding="UTF-8" + ) + + return destination_file + + +def upload_file( + local_file_path: str, + remote_file_path: str, +) -> str: + """Copies a local file to a GCS path.""" + subprocess.check_output( + ["gsutil", "cp", local_file_path, remote_file_path], encoding="UTF-8" + ) + + return remote_file_path + + +def archive_code_and_upload(staging_bucket: str): + # Archive all source in current directory + unique_id = uuid.uuid4() + source_archived_file = f"source_archived_{unique_id}.tar.gz" + + git_files = subprocess.check_output( + ["git", "ls-tree", "-r", "HEAD", "--name-only"], encoding="UTF-8" + ).split("\n") + + with tarfile.open(source_archived_file, "w:gz") as tar: + for file in git_files: + if len(file) > 0 and os.path.exists(file): + tar.add(file) + + # Upload archive to GCS bucket + source_archived_file_gcs = upload_file( + local_file_path=f"{source_archived_file}", + remote_file_path="/".join( + [staging_bucket, "code_archives", source_archived_file] + ), + ) + + print(f"Uploaded source code archive to {source_archived_file_gcs}") + + return source_archived_file_gcs \ No newline at end of file