diff --git a/compute_worker/compute_worker.py b/compute_worker/compute_worker.py index f99b07458..67b99525c 100644 --- a/compute_worker/compute_worker.py +++ b/compute_worker/compute_worker.py @@ -113,13 +113,17 @@ def to_bool(val): COMPETITION_CONTAINER_HTTP_PROXY = get("COMPETITION_CONTAINER_HTTP_PROXY", "") COMPETITION_CONTAINER_HTTPS_PROXY = get("COMPETITION_CONTAINER_HTTPS_PROXY", "") - CODALAB_IGNORE_CLEANUP_STEP = to_bool(get("CODALAB_IGNORE_CLEANUP_STEP")) + COMPUTE_WORKER_NO_CLEANUP = to_bool(get("COMPUTE_WORKER_NO_CLEANUP", "False")) WORKER_BUNDLE_URL_REWRITE = get("WORKER_BUNDLE_URL_REWRITE", "").strip() HUMAN_IN_THE_LOOP = ( get("HUMAN_IN_THE_LOOP", "false").lower() == "true" ) + COMPUTE_WORKER_DISABLE_LOG_UPLOAD = to_bool(get("COMPUTE_WORKER_DISABLE_LOG_UPLOAD", "False")) + COMPUTE_WORKER_DISABLE_PREDICTION_UPLOAD = to_bool(get("COMPUTE_WORKER_DISABLE_PREDICTION_UPLOAD", "False")) + + # ----------------------------------------------- # Program Kind @@ -171,6 +175,10 @@ class SubmissionStatus: f"{'with GPU capabilities: ' + Settings.GPU_DEVICE if Settings.USE_GPU else 'without GPU capabilities'}. " f"Network disabled for the competition container is set to {Settings.COMPETITION_CONTAINER_NETWORK_DISABLED}" ) +if Settings.COMPUTE_WORKER_DISABLE_PREDICTION_UPLOAD: + logger.warning("COMPUTE_WORKER_DISABLE_PREDICTION_UPLOAD is set to True, setting COMPUTE_WORKER_NO_CLEANUP to True") + Settings.COMPUTE_WORKER_NO_CLEANUP = True + # Intializing client # NOTE: CONTAINER_SOCKET is set in Settings based on CONTAINER_ENGINE_EXECUTABLE which must has either podman or docker @@ -372,10 +380,15 @@ def run_wrapper(run_args): except SubmissionException as e: msg = str(e).strip() if msg: - msg = f"Submission failed: {msg}. See logs for more details." + if Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: + msg = f"Submission failed: {msg}. Contact the Organizer(s) for more details." + else: + msg = f"Submission failed: {msg}. See logs for more details." else: - msg = "Submission failed. See logs for more details." - + if Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: + msg = "Submission failed. Contact the Organizer(s) for more details." + else: + msg = "Submission failed. See logs for more details." run._update_status(SubmissionStatus.FAILED, extra_information=msg) raise @@ -493,10 +506,19 @@ def __init__(self, run_args): self.run_related_name = ( f"uPK-{run_args['user_pk']}_sID-{run_args['id']}" ) + if run_args["is_scoring"]: + task_type = "scoring_program" + else: + task_type = "ingestion" + # Directories for the run self.watch = True self.completed_program_counter = 0 - self.root_dir = tempfile.mkdtemp(prefix=f'{self.run_related_name}__', dir=Settings.BASE_DIR) + # Create the folder then save the path to root_dir + self.submission_run_directory = Settings.BASE_DIR + f'{self.run_related_name}__' + os.mkdir(self.submission_run_directory + task_type, mode=0o700) + self.root_dir = self.submission_run_directory + task_type + self.bundle_dir = os.path.join(self.root_dir, "bundles") self.input_dir = os.path.join(self.root_dir, "input") self.output_dir = os.path.join(self.root_dir, "output") @@ -510,7 +532,10 @@ def __init__(self, run_args): self.submissions_api_url = run_args["submissions_api_url"] self.container_image = run_args["docker_image"] self.secret = run_args["secret"] - self.prediction_result = run_args["prediction_result"] + if Settings.COMPUTE_WORKER_DISABLE_PREDICTION_UPLOAD: + self.prediction_result = "Prediction Upload Disabled." + else: + self.prediction_result = run_args["prediction_result"] self.scoring_result = run_args.get("scoring_result") self.execution_time_limit = run_args["execution_time_limit"] # ----- HITL ------ @@ -601,18 +626,36 @@ async def watch_detailed_results(self): def push_logs(self): """Upload any collected logs, even in case of crash. """ - try: - for kind, logs in (self.logs or {}).items(): - for stream_key in ("stdout", "stderr"): - entry = logs.get(stream_key) if isinstance(logs, dict) else None - if not entry: - continue - location = entry.get("location") - data = entry.get("data") or b"" - if location: - self._put_file(location, raw_data=data) - except Exception as e: - logger.exception(f"Failed best-effort log upload: {e}") + if not Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: + try: + for kind, logs in (self.logs or {}).items(): + for stream_key in ("stdout", "stderr"): + entry = logs.get(stream_key) if isinstance(logs, dict) else None + if not entry: + continue + location = entry.get("location") + data = entry.get("data") or b"" + if location: + self._put_file(location, raw_data=data) + except Exception as e: + logger.exception(f"Failed best-effort log upload: {e}") + + else: + try: + logs_path = os.path.join(self.root_dir, "logs") + with open(logs_path, "w") as f: + for kind, logs in (self.logs or {}).items(): + for stream_key in ("stdout", "stderr"): + entry = logs.get(stream_key) if isinstance(logs, dict) else None + if not entry: + continue + location = entry.get("location") + data = entry.get("data") or b"" + if location: + f.write(str(data)) + except Exception as e: + logger.exception(f"Failed best-effort log file creation: {e}") + def get_detailed_results_file_path(self): default_detailed_results_path = os.path.join( @@ -767,9 +810,14 @@ def _get_container_image(self, image_name): self._update_submission(docker_pull_fail_data) # Send error through web socket to the frontend asyncio.run(self._send_data_through_socket(str(pull_error))) - raise DockerImagePullException( - f"Pull for {image_name} failed! Check the logs for more information" - ) + if Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: + raise DockerImagePullException( + f"Pull for {image_name} failed! Contact the Organizer(s) for more details." + ) + else: + raise DockerImagePullException( + f"Pull for {image_name} failed! Check the logs for more information" + ) else: logger.warning("Failed. Retrying in 5 seconds...") time.sleep(5) # Wait 5 seconds before retrying @@ -908,7 +956,8 @@ def _create_container( "SYS_CHROOT", ] - # Configure whether or not we use the GPU. Also setting auto_remove to False because + # Configure whether or not we use the GPU. Also setting auto_remove to False because removing too fast + # can bug out the worker (can't get the logs fast enough) if Settings.CONTAINER_ENGINE_EXECUTABLE == Settings.DOCKER: security_options = ["no-new-privileges"] else: @@ -982,21 +1031,24 @@ async def _run_container_engine_cmd(self, container, kind): # Create a websocket to send the logs in real time to the codabench instance # We need to set a timeout for the websocket connection otherwise the program will get stuck if he websocket does not connect. websocket = None - try: - websocket_url = f"{self.websocket_url}?kind={kind}" - logger.debug(f"Connecting to {websocket_url} for container {str(container.get('Id'))}") - websocket = await asyncio.wait_for( - websockets.connect(websocket_url), timeout=10.0 - ) - logger.debug(f"connected to {websocket_url} for container {str(container.get('Id'))}") - except Exception as e: - logger.error( - f"There was an error trying to connect to the websocket on the codabench instance: {e}" - ) + # Do not create a websocket if the real time logs are not wanted (Silent Compute Worker) + if not Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: + try: + websocket_url = f"{self.websocket_url}?kind={kind}" + logger.debug(f"Connecting to {websocket_url} for container {str(container.get('Id'))}") + websocket = await asyncio.wait_for( + websockets.connect(websocket_url), timeout=10.0 + ) + logger.debug(f"connected to {websocket_url} for container {str(container.get('Id'))}") - if Settings.LOG_LEVEL == Settings.LOG_LEVEL_DEBUG: - logger.exception(e) + except Exception as e: + logger.error( + f"There was an error trying to connect to the websocket on the codabench instance: {e}" + ) + + if Settings.LOG_LEVEL == Settings.LOG_LEVEL_DEBUG: + logger.exception(e) start = time.time() @@ -1012,6 +1064,7 @@ async def _run_container_engine_cmd(self, container, kind): ) # If we enter the for loop after the container exited, the program will get stuck + # Do not send the real time logs if they are not wanted (Silent Compute Worker) if client.inspect_container(container)["State"]["Status"].lower() == "running": logger.debug( "Show the logs and stream them to codabench " + container.get("Id") @@ -1021,25 +1074,27 @@ async def _run_container_engine_cmd(self, container, kind): if log[0] is not None: stdout_chunks.append(log[0]) logger.info(log[0].decode()) - try: - if websocket is not None: - await websocket.send( - json.dumps({"kind": kind, "message": log[0].decode()}) - ) - except Exception as e: - logger.error(e) + if not Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: + try: + if websocket is not None: + await websocket.send( + json.dumps({"kind": kind, "message": log[0].decode()}) + ) + except Exception as e: + logger.error(e) # Errors elif log[1] is not None: stderr_chunks.append(log[1]) logger.error(log[1].decode()) - try: - if websocket is not None: - await websocket.send( - json.dumps({"kind": kind, "message": log[1].decode()}) - ) - except Exception as e: - logger.error(e) + if not Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: + try: + if websocket is not None: + await websocket.send( + json.dumps({"kind": kind, "message": log[1].decode()}) + ) + except Exception as e: + logger.error(e) except (docker.errors.NotFound, docker.errors.APIError) as e: logger.error(e) @@ -1055,15 +1110,17 @@ async def _run_container_engine_cmd(self, container, kind): # Gets the logs of the container, sperating stdout and stderr (first and second position) thanks for demux=True return_Code = client.wait(container) logs_Unified = (b"".join(stdout_chunks), b"".join(stderr_chunks)) - logger.debug( - f"WORKER_MARKER: Disconnecting from {websocket_url}, program counter = {self.completed_program_counter}" - ) - if websocket is not None: - try: - await websocket.close() - await websocket.wait_closed() - except Exception as e: - logger.error(e) + + if not Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: + logger.debug( + f"WORKER_MARKER: Disconnecting from {websocket_url}, program counter = {self.completed_program_counter}" + ) + if websocket is not None: + try: + await websocket.close() + await websocket.wait_closed() + except Exception as e: + logger.error(e) client.remove_container(container, v=True, force=True) logger.debug(f"Container {container.get('Id')} exited with status code : {str(return_Code['StatusCode'])}") @@ -1345,15 +1402,22 @@ def prepare(self): (self.input_data, "input_data"), (self.reference_data, "input/ref"), ] - if self.is_scoring: + if self.is_scoring and not Settings.COMPUTE_WORKER_DISABLE_PREDICTION_UPLOAD: # Send along submission result so scoring_program can get access bundles += [(self.prediction_result, "input/res")] + elif self.is_scoring: + bundles += [("local_prediction_results", "submission")] for url, path in bundles: if url is not None: # At the moment let's just cache input & reference data cache_this_bundle = path in ("input_data", "input/ref") - zip_file = self._get_bundle(url, path, cache=cache_this_bundle) + if url == "local_prediction_results" and self.is_scoring: + submission_run_directory_ingestion = self.submission_run_directory + "ingestion/output/" + submission_run_directory_scoring = self.submission_run_directory + "scoring_program/submission/" + shutil.copytree(submission_run_directory_ingestion, submission_run_directory_scoring, dirs_exist_ok=True) + else: + zip_file = self._get_bundle(url, path, cache=cache_this_bundle) # Computing checksum of the submission file during ingestion run if url == self.submission_data and not self.is_scoring: @@ -1423,6 +1487,7 @@ def start(self): self._run_program_directory(kind=ProgramKind.INGESTION_PROGRAM, program_dir=ingestion_program_dir), ]) + logger.info(tasks) gathered_tasks = asyncio.gather(*tasks, return_exceptions=True) task_results = [] # will store results/exceptions from gather @@ -1520,12 +1585,14 @@ def start(self): self.ingestion_program_exit_code = return_code self.ingestion_program_elapsed_time = elapsed_time logger.info(f"[exited with {logs['returncode']}]") - for key, value in logs.items(): - if key not in ["stdout", "stderr"]: - continue - if value["data"]: - logger.info(f"[{key}]\n{value['data']}") - self._put_file(value["location"], raw_data=value["data"]) + if Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD == False: + for key, value in logs.items(): + if key not in ["stdout", "stderr"]: + continue + if value["data"]: + logger.info(f"[{key}]\n{value['data']}") + self._put_file(value["location"], raw_data=value["data"]) + # set logs of this kind to None, since we handled them already logger.info("Program finished") @@ -1710,15 +1777,16 @@ def push_output(self): raise SubmissionException("Failed to write metadata file.") if not self.is_scoring: - self._put_dir(self.prediction_result, self.output_dir) + if not Settings.COMPUTE_WORKER_DISABLE_PREDICTION_UPLOAD: + self._put_dir(self.prediction_result, self.output_dir) else: self._put_dir(self.scoring_result, self.output_dir) def clean_up(self): self.stop_hitl_http_server() - if Settings.CODALAB_IGNORE_CLEANUP_STEP: + if Settings.COMPUTE_WORKER_NO_CLEANUP: logger.warning( - f"CODALAB_IGNORE_CLEANUP_STEP mode enabled, ignoring clean up of: {self.root_dir}" + f"COMPUTE_WORKER_NO_CLEANUP mode enabled, ignoring clean up of: {self.root_dir}" ) return diff --git a/docker-compose.yml b/docker-compose.yml index 5b66662c7..0bd754e22 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -244,7 +244,7 @@ services: environment: - BROKER_URL=pyamqp://${RABBITMQ_DEFAULT_USER}:${RABBITMQ_DEFAULT_PASS}@${RABBITMQ_HOST}:${RABBITMQ_PORT}// # Make the worker leave behind the submission so we can examine it - - CODALAB_IGNORE_CLEANUP_STEP=1 + - COMPUTE_WORKER_NO_CLEANUP=True tty: true logging: options: diff --git a/documentation/docs/Organizers/Running_a_benchmark/Compute-Worker-Management---Setup.md b/documentation/docs/Organizers/Running_a_benchmark/Compute-Worker-Management---Setup.md index fc0576c1c..8909834dc 100644 --- a/documentation/docs/Organizers/Running_a_benchmark/Compute-Worker-Management---Setup.md +++ b/documentation/docs/Organizers/Running_a_benchmark/Compute-Worker-Management---Setup.md @@ -60,6 +60,18 @@ CONTAINER_ENGINE_EXECUTABLE=docker #USE_GPU=True #GPU_DEVICE=nvidia.com/gpu=all #HUMAN_IN_THE_LOOP=False +# This option removes the ability of the compute worker to send logs to +# codabench, instead writing them locally on disk. Combine with +# COMPUTE_WORKER_NO_CLEANUP=true to stop the worker's cleanup to keep +# all the logs locally only +#COMPUTE_WORKER_DISABLE_LOG_UPLOAD=False +# Stop the predictions from being sent to Codabench. +# This option requires only having one compute worker for ingestion +# and scoring. +#COMPUTE_WORKER_DISABLE_PREDICTION_UPLOAD=False + +#COMPUTE_WORKER_NO_CLEANUP=False + ####################################################################### # Network # ####################################################################### diff --git a/documentation/docs/Organizers/Running_a_benchmark/Compute-worker-installation-with-Podman.md b/documentation/docs/Organizers/Running_a_benchmark/Compute-worker-installation-with-Podman.md index 6e2729157..153aa063c 100644 --- a/documentation/docs/Organizers/Running_a_benchmark/Compute-worker-installation-with-Podman.md +++ b/documentation/docs/Organizers/Running_a_benchmark/Compute-worker-installation-with-Podman.md @@ -37,6 +37,18 @@ HOST_DIRECTORY=/codabench CONTAINER_ENGINE_EXECUTABLE=podman #USE_GPU=True #GPU_DEVICE=nvidia.com/gpu=all +#HUMAN_IN_THE_LOOP=False +# This option removes the ability of the compute worker to send logs to +# codabench, instead writing them locally on disk. Combine with +# COMPUTE_WORKER_NO_CLEANUP=true to stop the worker's cleanup to keep +# all the logs locally only +#COMPUTE_WORKER_DISABLE_LOG_UPLOAD=False +# Stop the predictions from being sent to Codabench. +# This option requires only having one compute worker for ingestion +# and scoring. +#COMPUTE_WORKER_DISABLE_PREDICTION_UPLOAD=False + +#COMPUTE_WORKER_NO_CLEANUP=False ####################################################################### # Network # diff --git a/src/settings/logs_loguru.py b/src/settings/logs_loguru.py index 28b2cf075..137ff83c9 100644 --- a/src/settings/logs_loguru.py +++ b/src/settings/logs_loguru.py @@ -130,7 +130,17 @@ def colorize_run_args(json_str): json_str, ) json_str = re.sub( - r'("ingestion_program": ")(.*?)(",)', + r'("ingestion_program_data": ")(.*?)(",)', + rf"\1{yellow}\2{reset}\3{lineskip}", + json_str, + ) + json_str = re.sub( + r'("submission_data": ")(.*?)(",)', + rf"\1{yellow}\2{reset}\3{lineskip}", + json_str, + ) + json_str = re.sub( + r'("scoring_program_data": ")(.*?)(",)', rf"\1{yellow}\2{reset}\3{lineskip}", json_str, )