From 4a7f493fb7c414d669947e867f380583949f5e24 Mon Sep 17 00:00:00 2001 From: Obada Haddad Date: Fri, 3 Jul 2026 16:09:49 +0200 Subject: [PATCH 1/8] add option in compute worker env to not send logs to the instance, instead writing them in a local file --- compute_worker/compute_worker.py | 161 +++++++++++++++++++------------ 1 file changed, 101 insertions(+), 60 deletions(-) diff --git a/compute_worker/compute_worker.py b/compute_worker/compute_worker.py index f99b07458..56f56c436 100644 --- a/compute_worker/compute_worker.py +++ b/compute_worker/compute_worker.py @@ -120,6 +120,9 @@ def to_bool(val): get("HUMAN_IN_THE_LOOP", "false").lower() == "true" ) + SILENT_COMPUTE_WORKER = get("SILENT_COMPUTE_WORKER", "false").lower() + + # ----------------------------------------------- # Program Kind @@ -372,10 +375,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.SILENT_COMPUTE_WORKER == "true": + msg = f"Submission failed: {msg}. Contact the Organizer 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.SILENT_COMPUTE_WORKER == "true": + msg = "Submission failed. Contact the Organizer for more details." + else: + msg = "Submission failed. See logs for more details." run._update_status(SubmissionStatus.FAILED, extra_information=msg) raise @@ -601,18 +609,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 Settings.SILENT_COMPUTE_WORKER == "false": + 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 +793,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.SILENT_COMPUTE_WORKER == "true": + raise DockerImagePullException( + f"Pull for {image_name} failed! Contact the Organizer 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 @@ -982,21 +1013,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 Settings.SILENT_COMPUTE_WORKER == "false": + 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 +1046,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 +1056,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 Settings.SILENT_COMPUTE_WORKER == "false": + 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 Settings.SILENT_COMPUTE_WORKER == "false": + 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 +1092,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 Settings.SILENT_COMPUTE_WORKER == "false": + 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'])}") @@ -1520,12 +1559,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.SILENT_COMPUTE_WORKER == "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") From c9ff45acb236b9ac8d320fc655d5c8a60f68202c Mon Sep 17 00:00:00 2001 From: Obada Haddad Date: Fri, 3 Jul 2026 16:18:16 +0200 Subject: [PATCH 2/8] rename the No Cleanup env variable, add documentation --- compute_worker/compute_worker.py | 6 +++--- docker-compose.yml | 2 +- .../Compute-Worker-Management---Setup.md | 7 +++++++ 3 files changed, 11 insertions(+), 4 deletions(-) diff --git a/compute_worker/compute_worker.py b/compute_worker/compute_worker.py index 56f56c436..46509a5c0 100644 --- a/compute_worker/compute_worker.py +++ b/compute_worker/compute_worker.py @@ -113,7 +113,7 @@ 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")) WORKER_BUNDLE_URL_REWRITE = get("WORKER_BUNDLE_URL_REWRITE", "").strip() HUMAN_IN_THE_LOOP = ( @@ -1757,9 +1757,9 @@ def push_output(self): 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..d4458f70d 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..858dea24c 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,13 @@ 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 +#SILENT_COMPUTE_WORKER=False +#COMPUTE_WORKER_NO_CLEANUP=false + ####################################################################### # Network # ####################################################################### From ab3323ece88522eac153e2b2c10b4522972e4be8 Mon Sep 17 00:00:00 2001 From: Obada Haddad Date: Mon, 6 Jul 2026 15:53:50 +0200 Subject: [PATCH 3/8] use real boolean values --- compute_worker/compute_worker.py | 23 ++++++++++++----------- 1 file changed, 12 insertions(+), 11 deletions(-) diff --git a/compute_worker/compute_worker.py b/compute_worker/compute_worker.py index 46509a5c0..3260023f8 100644 --- a/compute_worker/compute_worker.py +++ b/compute_worker/compute_worker.py @@ -113,14 +113,14 @@ def to_bool(val): COMPETITION_CONTAINER_HTTP_PROXY = get("COMPETITION_CONTAINER_HTTP_PROXY", "") COMPETITION_CONTAINER_HTTPS_PROXY = get("COMPETITION_CONTAINER_HTTPS_PROXY", "") - COMPUTE_WORKER_NO_CLEANUP = to_bool(get("COMPUTE_WORKER_NO_CLEANUP")) + 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" ) - SILENT_COMPUTE_WORKER = get("SILENT_COMPUTE_WORKER", "false").lower() + SILENT_COMPUTE_WORKER = to_bool(get("SILENT_COMPUTE_WORKER", "False")) @@ -375,12 +375,12 @@ def run_wrapper(run_args): except SubmissionException as e: msg = str(e).strip() if msg: - if Settings.SILENT_COMPUTE_WORKER == "true": + if Settings.SILENT_COMPUTE_WORKER: msg = f"Submission failed: {msg}. Contact the Organizer for more details." else: msg = f"Submission failed: {msg}. See logs for more details." else: - if Settings.SILENT_COMPUTE_WORKER == "true": + if Settings.SILENT_COMPUTE_WORKER: msg = "Submission failed. Contact the Organizer for more details." else: msg = "Submission failed. See logs for more details." @@ -609,7 +609,7 @@ async def watch_detailed_results(self): def push_logs(self): """Upload any collected logs, even in case of crash. """ - if Settings.SILENT_COMPUTE_WORKER == "false": + if Settings.SILENT_COMPUTE_WORKER == False: try: for kind, logs in (self.logs or {}).items(): for stream_key in ("stdout", "stderr"): @@ -793,7 +793,7 @@ 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))) - if Settings.SILENT_COMPUTE_WORKER == "true": + if Settings.SILENT_COMPUTE_WORKER: raise DockerImagePullException( f"Pull for {image_name} failed! Contact the Organizer for more details." ) @@ -1015,7 +1015,7 @@ async def _run_container_engine_cmd(self, container, kind): websocket = None # Do not create a websocket if the real time logs are not wanted (Silent Compute Worker) - if Settings.SILENT_COMPUTE_WORKER == "false": + if Settings.SILENT_COMPUTE_WORKER == False: try: websocket_url = f"{self.websocket_url}?kind={kind}" logger.debug(f"Connecting to {websocket_url} for container {str(container.get('Id'))}") @@ -1056,7 +1056,7 @@ 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()) - if Settings.SILENT_COMPUTE_WORKER == "false": + if Settings.SILENT_COMPUTE_WORKER == False: try: if websocket is not None: await websocket.send( @@ -1069,7 +1069,7 @@ async def _run_container_engine_cmd(self, container, kind): elif log[1] is not None: stderr_chunks.append(log[1]) logger.error(log[1].decode()) - if Settings.SILENT_COMPUTE_WORKER == "false": + if Settings.SILENT_COMPUTE_WORKER == False: try: if websocket is not None: await websocket.send( @@ -1093,7 +1093,7 @@ async def _run_container_engine_cmd(self, container, kind): return_Code = client.wait(container) logs_Unified = (b"".join(stdout_chunks), b"".join(stderr_chunks)) - if Settings.SILENT_COMPUTE_WORKER == "false": + if Settings.SILENT_COMPUTE_WORKER == False: logger.debug( f"WORKER_MARKER: Disconnecting from {websocket_url}, program counter = {self.completed_program_counter}" ) @@ -1462,6 +1462,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 @@ -1559,7 +1560,7 @@ def start(self): self.ingestion_program_exit_code = return_code self.ingestion_program_elapsed_time = elapsed_time logger.info(f"[exited with {logs['returncode']}]") - if Settings.SILENT_COMPUTE_WORKER == "false": + if Settings.SILENT_COMPUTE_WORKER == False: for key, value in logs.items(): if key not in ["stdout", "stderr"]: continue From 7e0dd7899087b54fa0c8015cf47d78af41e56228 Mon Sep 17 00:00:00 2001 From: Obada Haddad Date: Mon, 6 Jul 2026 16:27:44 +0200 Subject: [PATCH 4/8] docs: update value to use Boolean False --- .../Running_a_benchmark/Compute-Worker-Management---Setup.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 858dea24c..2472667ea 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 @@ -65,7 +65,7 @@ CONTAINER_ENGINE_EXECUTABLE=docker # COMPUTE_WORKER_NO_CLEANUP=true to stop the worker's cleanup to keep # all the logs locally only #SILENT_COMPUTE_WORKER=False -#COMPUTE_WORKER_NO_CLEANUP=false +#COMPUTE_WORKER_NO_CLEANUP=False ####################################################################### # Network # From 1b9848f2637bf1509d3983f096147ce79ef6cf1f Mon Sep 17 00:00:00 2001 From: Obada Haddad Date: Mon, 6 Jul 2026 16:50:20 +0200 Subject: [PATCH 5/8] update logs_loguru to inclue new tasks variable names to color them --- src/settings/logs_loguru.py | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) 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, ) From 075c943903925eda269d5a07c8bf950641ee7ce3 Mon Sep 17 00:00:00 2001 From: Obada Haddad Date: Mon, 6 Jul 2026 16:51:04 +0200 Subject: [PATCH 6/8] change variable name to use boolean --- docker-compose.yml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docker-compose.yml b/docker-compose.yml index d4458f70d..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 - - COMPUTE_WORKER_NO_CLEANUP="true" + - COMPUTE_WORKER_NO_CLEANUP=True tty: true logging: options: From 8fc59fc8ce7c5c50fe9f0b0056dc883c9633c128 Mon Sep 17 00:00:00 2001 From: Obada Haddad Date: Tue, 7 Jul 2026 14:44:15 +0200 Subject: [PATCH 7/8] fix some syntax --- compute_worker/compute_worker.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/compute_worker/compute_worker.py b/compute_worker/compute_worker.py index 3260023f8..4b4b55c3c 100644 --- a/compute_worker/compute_worker.py +++ b/compute_worker/compute_worker.py @@ -376,12 +376,12 @@ def run_wrapper(run_args): msg = str(e).strip() if msg: if Settings.SILENT_COMPUTE_WORKER: - msg = f"Submission failed: {msg}. Contact the Organizer for more details." + msg = f"Submission failed: {msg}. Contact the Organizer(s) for more details." else: msg = f"Submission failed: {msg}. See logs for more details." else: if Settings.SILENT_COMPUTE_WORKER: - msg = "Submission failed. Contact the Organizer for more details." + 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) @@ -795,7 +795,7 @@ def _get_container_image(self, image_name): asyncio.run(self._send_data_through_socket(str(pull_error))) if Settings.SILENT_COMPUTE_WORKER: raise DockerImagePullException( - f"Pull for {image_name} failed! Contact the Organizer for more details." + f"Pull for {image_name} failed! Contact the Organizer(s) for more details." ) else: raise DockerImagePullException( @@ -939,7 +939,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: From b9427babb08c68ea406045850fa0e1e3d765600f Mon Sep 17 00:00:00 2001 From: Obada Haddad Date: Thu, 27 Aug 2026 15:51:02 +0200 Subject: [PATCH 8/8] add option to forbid the compute worker from sending predictions to the codabench instance storage --- compute_worker/compute_worker.py | 55 ++++++++++++++----- .../Compute-Worker-Management---Setup.md | 7 ++- ...Compute-worker-installation-with-Podman.md | 12 ++++ 3 files changed, 58 insertions(+), 16 deletions(-) diff --git a/compute_worker/compute_worker.py b/compute_worker/compute_worker.py index 4b4b55c3c..67b99525c 100644 --- a/compute_worker/compute_worker.py +++ b/compute_worker/compute_worker.py @@ -120,7 +120,8 @@ def to_bool(val): get("HUMAN_IN_THE_LOOP", "false").lower() == "true" ) - SILENT_COMPUTE_WORKER = to_bool(get("SILENT_COMPUTE_WORKER", "False")) + 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")) @@ -174,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 @@ -375,12 +380,12 @@ def run_wrapper(run_args): except SubmissionException as e: msg = str(e).strip() if msg: - if Settings.SILENT_COMPUTE_WORKER: + 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: - if Settings.SILENT_COMPUTE_WORKER: + 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." @@ -501,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") @@ -518,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 ------ @@ -609,7 +626,7 @@ async def watch_detailed_results(self): def push_logs(self): """Upload any collected logs, even in case of crash. """ - if Settings.SILENT_COMPUTE_WORKER == False: + if not Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: try: for kind, logs in (self.logs or {}).items(): for stream_key in ("stdout", "stderr"): @@ -793,7 +810,7 @@ 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))) - if Settings.SILENT_COMPUTE_WORKER: + if Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: raise DockerImagePullException( f"Pull for {image_name} failed! Contact the Organizer(s) for more details." ) @@ -1016,7 +1033,7 @@ async def _run_container_engine_cmd(self, container, kind): websocket = None # Do not create a websocket if the real time logs are not wanted (Silent Compute Worker) - if Settings.SILENT_COMPUTE_WORKER == False: + 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'))}") @@ -1057,7 +1074,7 @@ 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()) - if Settings.SILENT_COMPUTE_WORKER == False: + if not Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: try: if websocket is not None: await websocket.send( @@ -1070,7 +1087,7 @@ async def _run_container_engine_cmd(self, container, kind): elif log[1] is not None: stderr_chunks.append(log[1]) logger.error(log[1].decode()) - if Settings.SILENT_COMPUTE_WORKER == False: + if not Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: try: if websocket is not None: await websocket.send( @@ -1094,7 +1111,7 @@ async def _run_container_engine_cmd(self, container, kind): return_Code = client.wait(container) logs_Unified = (b"".join(stdout_chunks), b"".join(stderr_chunks)) - if Settings.SILENT_COMPUTE_WORKER == False: + if not Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD: logger.debug( f"WORKER_MARKER: Disconnecting from {websocket_url}, program counter = {self.completed_program_counter}" ) @@ -1385,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: @@ -1561,7 +1585,7 @@ def start(self): self.ingestion_program_exit_code = return_code self.ingestion_program_elapsed_time = elapsed_time logger.info(f"[exited with {logs['returncode']}]") - if Settings.SILENT_COMPUTE_WORKER == False: + if Settings.COMPUTE_WORKER_DISABLE_LOG_UPLOAD == False: for key, value in logs.items(): if key not in ["stdout", "stderr"]: continue @@ -1753,7 +1777,8 @@ 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) 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 2472667ea..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 @@ -64,7 +64,12 @@ CONTAINER_ENGINE_EXECUTABLE=docker # 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 -#SILENT_COMPUTE_WORKER=False +#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 ####################################################################### 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 #