From f7cf4fb893be4706d7f54abb4e3f5f8a6972ee68 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Fri, 25 Sep 2026 14:56:20 +0200 Subject: [PATCH 01/14] [core] ci_entrypoint: fournos: properly unset the KUBECONFIG if already set --- projects/core/ci_entrypoint/fournos.py | 63 +++++++++++++++----------- 1 file changed, 36 insertions(+), 27 deletions(-) diff --git a/projects/core/ci_entrypoint/fournos.py b/projects/core/ci_entrypoint/fournos.py index 1a63db723..1473d35d5 100644 --- a/projects/core/ci_entrypoint/fournos.py +++ b/projects/core/ci_entrypoint/fournos.py @@ -36,39 +36,48 @@ def check_fjob_resolver_error(): logger.info(f"Checking FournosJob resolver status for fjob/{job_name} in {namespace}") + # Unset KUBECONFIG to use the pod SA access (fjob lives on the Fournos cluster) + original_kubeconfig = os.environ.get("KUBECONFIG") + if "KUBECONFIG" in os.environ: + del os.environ["KUBECONFIG"] + try: - result = run.run( - f"oc get fjob/{job_name} -n {namespace} -o json", - capture_stdout=True, - check=True, - ) - fjob_data = json.loads(result.stdout) - except Exception as e: - logger.warning(f"Could not fetch FournosJob for resolver error check: {e}") - return + try: + result = run.run( + f"oc get fjob/{job_name} -n {namespace} -o json", + capture_stdout=True, + check=True, + ) + fjob_data = json.loads(result.stdout) + except Exception as e: + logger.warning(f"Could not fetch FournosJob for resolver error check: {e}") + return - resolver_error = ( - fjob_data.get("status", {}) - .get("engineStatus", {}) - .get("forge", {}) - .get("resolver", {}) - .get("error") - ) + resolver_error = ( + fjob_data.get("status", {}) + .get("engineStatus", {}) + .get("forge", {}) + .get("resolver", {}) + .get("error") + ) - resolver_status = ( - fjob_data.get("status", {}).get("engineStatus", {}).get("forge", {}).get("resolver", {}) - ) - resolver_pod = resolver_status.get("pod") - logs_captured = resolver_status.get("logsCaptured", False) + resolver_status = ( + fjob_data.get("status", {}).get("engineStatus", {}).get("forge", {}).get("resolver", {}) + ) + resolver_pod = resolver_status.get("pod") + logs_captured = resolver_status.get("logsCaptured", False) - if resolver_pod and not logs_captured: - logger.info(f"Capturing logs from resolver pod: {resolver_pod}") - _capture_resolver_pod_logs(job_name, namespace, resolver_pod) + if resolver_pod and not logs_captured: + logger.info(f"Capturing logs from resolver pod: {resolver_pod}") + _capture_resolver_pod_logs(job_name, namespace, resolver_pod) - if resolver_error: - raise RuntimeError(f"FournosJob resolver failed: {resolver_error}") + if resolver_error: + raise RuntimeError(f"FournosJob resolver failed: {resolver_error}") - logger.info("FournosJob resolver status: OK (no errors)") + logger.info("FournosJob resolver status: OK (no errors)") + finally: + if original_kubeconfig is not None: + os.environ["KUBECONFIG"] = original_kubeconfig def _capture_resolver_pod_logs(job_name: str, namespace: str, pod_name: str) -> None: From 9f74cebc97118ba103011052947ea82e902ea8d5 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 11:31:58 +0200 Subject: [PATCH 02/14] [core] library: export: export the artifacts-export status.yaml --- projects/core/library/export.py | 1 + 1 file changed, 1 insertion(+) diff --git a/projects/core/library/export.py b/projects/core/library/export.py index d5341b91e..e9e9c2d76 100644 --- a/projects/core/library/export.py +++ b/projects/core/library/export.py @@ -557,6 +557,7 @@ def _update_artifacts( for fpath in [ artifact_dir_path / "run.log", + artifact_dir_path / "status.yaml", artifact_dir_path / "COMPLETION-NOTIFICATION.md", artifact_dir_path / "000__ci_metadata" / "fournos_fjob.yaml", ]: From 5f764f2d57a5e6914a33c08596e59f39cba6e7f1 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 11:32:16 +0200 Subject: [PATCH 03/14] [core] library: export_notifications: lookup the fjob in the right step directory --- projects/core/library/export.py | 4 ++-- projects/core/library/export_notifications.py | 15 ++++++++++++--- 2 files changed, 14 insertions(+), 5 deletions(-) diff --git a/projects/core/library/export.py b/projects/core/library/export.py index e9e9c2d76..1b70dc955 100644 --- a/projects/core/library/export.py +++ b/projects/core/library/export.py @@ -374,7 +374,7 @@ def caliper_export_entrypoint( ) # Check for job shutdown/abort status and add to mock status - shutdown_status = _check_job_shutdown_status(artifact_dir) + shutdown_status = _check_job_shutdown_status() if shutdown_status: status.job_shutdown = JobShutdown.from_dict(shutdown_status) logger.info(f"Added job shutdown status to dry run mock status: {shutdown_status}") @@ -386,7 +386,7 @@ def caliper_export_entrypoint( ) # Check for job shutdown/abort status and add to main export status - shutdown_status = _check_job_shutdown_status(artifact_dir) + shutdown_status = _check_job_shutdown_status() if shutdown_status: status.job_shutdown = JobShutdown.from_dict(shutdown_status) logger.info(f"Added job shutdown status to main export status: {shutdown_status}") diff --git a/projects/core/library/export_notifications.py b/projects/core/library/export_notifications.py index 06603b695..9869318fb 100644 --- a/projects/core/library/export_notifications.py +++ b/projects/core/library/export_notifications.py @@ -1312,18 +1312,27 @@ def _process_step_details(step_dir: Path, mlflow_run_url: str | None = None) -> return step_details -def _check_job_shutdown_status(artifact_dir: Path) -> dict[str, Any] | None: +def _check_job_shutdown_status(artifact_dir: Path | None = None) -> dict[str, Any] | None: """Check if the job has been aborted via spec.shutdown field.""" try: - metadata_dir = ci_lib.get_ci_metadata_dir(artifact_dir, any_level=True) + metadata_dir = ( + ci_lib.get_ci_metadata_dir(artifact_dir, any_level=True) + if artifact_dir + else ci_lib.get_ci_metadata_dir_location() + ) + fournos_fjob_path = metadata_dir / "fournos_fjob.yaml" + logger.info(f"Checking job shutdown status from {fournos_fjob_path}") if not fournos_fjob_path.exists(): + logger.info(f"fournos_fjob.yaml not found at {fournos_fjob_path}") return None with open(fournos_fjob_path, encoding="utf-8") as f: fjob_data = yaml.safe_load(f) - shutdown_value = fjob_data.get("spec", {}).get("shutdown") + spec = fjob_data.get("spec", {}) + shutdown_value = spec.get("shutdown") + logger.info(f"spec keys: {list(spec.keys())}, shutdown_value: {shutdown_value!r}") if shutdown_value: return { "shutdown_detected": True, From 9b6e487a08315be2b7b507d2ad3339f36df4ef5c Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 11:36:45 +0200 Subject: [PATCH 04/14] [core] ci_entrypoint: run_common: fix the signal handler --- projects/core/ci_entrypoint/run_common.py | 71 +++++++++++++++-------- 1 file changed, 48 insertions(+), 23 deletions(-) diff --git a/projects/core/ci_entrypoint/run_common.py b/projects/core/ci_entrypoint/run_common.py index 96e1d6bd8..0c4749657 100644 --- a/projects/core/ci_entrypoint/run_common.py +++ b/projects/core/ci_entrypoint/run_common.py @@ -35,40 +35,63 @@ def setup_logging(): _child_process = None -def signal_handler_sigint(sig, frame): - """Handle SIGINT (Ctrl+C) gracefully.""" - logger.info("🚫 Received SIGINT (Ctrl+C) - Interrupting operation...") +CHILD_SIGNAL_TIMEOUT = 30 - # Forward signal to child process first - if _child_process and _child_process.poll() is None: # Child is still running - logger.info("📡 Forwarding SIGINT to child process...") - try: - _child_process.send_signal(signal.SIGINT) - except (OSError, ProcessLookupError): - pass # Child may have already terminated - # Emergency cleanup of dual output - prepare_ci.shutdown_dual_output() +def _write_signal_file(sig_name): + from datetime import datetime + + artifact_dir = os.environ.get("ARTIFACT_DIR") + if not artifact_dir: + return + sig_file = Path(artifact_dir) / f"{sig_name}_interrupted" + with sig_file.open("a") as f: + f.write(f"{datetime.now()}: {__name__}._forward_signal_and_exit {sig_name} handler\n") - sys.exit(130) # Standard exit code for SIGINT +def _forward_signal_and_exit(sig, exit_code): + sig_name = signal.Signals(sig).name + child_pid = _child_process.pid if _child_process else None + child_alive = _child_process.poll() is None if _child_process else False -def signal_handler_sigterm(sig, frame): - """Handle SIGTERM gracefully.""" - logger.info("🛑 Received SIGTERM - Terminating operation...") + logger.info( + f"Received {sig_name} in pid={os.getpid()} pgid={os.getpgrp()}, child_pid={child_pid} child_alive={child_alive}" + ) - # Forward signal to child process first - if _child_process and _child_process.poll() is None: # Child is still running - logger.info("📡 Forwarding SIGTERM to child process...") + _write_signal_file(sig_name) + + if _child_process and child_alive: try: - _child_process.send_signal(signal.SIGTERM) - except (OSError, ProcessLookupError): - pass # Child may have already terminated + _child_process.send_signal(sig) + logger.info(f"Forwarded {sig_name} to child pid={child_pid}") + except (OSError, ProcessLookupError) as e: + logger.info(f"Failed to forward {sig_name} to child pid={child_pid}: {e}") + else: + logger.info( + f"Waiting up to {CHILD_SIGNAL_TIMEOUT}s for child pid={child_pid} to exit ..." + ) + try: + _child_process.wait(timeout=CHILD_SIGNAL_TIMEOUT) + logger.info(f"Child pid={child_pid} exited after {sig_name}") + except subprocess.TimeoutExpired: + logger.warning( + f"Child pid={child_pid} did not exit within {CHILD_SIGNAL_TIMEOUT}s after {sig_name}" + ) + else: + logger.info(f"No child to forward {sig_name} to") # Emergency cleanup of dual output prepare_ci.shutdown_dual_output() - sys.exit(143) # Standard exit code for SIGTERM + sys.exit(exit_code) + + +def signal_handler_sigint(sig, frame): + _forward_signal_and_exit(sig, 130) + + +def signal_handler_sigterm(sig, frame): + _forward_signal_and_exit(sig, 143) def setup_signal_handlers(): @@ -431,6 +454,8 @@ def execute_project_operation( stderr=None, # Inherit stderr for pdb/debugging ) + logger.info(f"Parent pid={os.getpid()} pgid={os.getpgrp()}, child pid={_child_process.pid}") + # Wait for process to complete result_code = _child_process.wait() From b082b33519d88799a23648320ae0875cc984a1e0 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 11:32:38 +0200 Subject: [PATCH 05/14] [skeleton] orchestration: add signal handler placeholder --- projects/skeleton/orchestration/ci.py | 15 ++++++++++++ .../skeleton/orchestration/test_skeleton.py | 23 ------------------- 2 files changed, 15 insertions(+), 23 deletions(-) diff --git a/projects/skeleton/orchestration/ci.py b/projects/skeleton/orchestration/ci.py index 58fcd8d7c..08fa6a3a0 100755 --- a/projects/skeleton/orchestration/ci.py +++ b/projects/skeleton/orchestration/ci.py @@ -9,7 +9,9 @@ import logging import pathlib +import signal import types +from datetime import datetime import click import prepare_skeleton @@ -34,9 +36,22 @@ def init(): env.init() run.init() + run.register_signal_callback(_signal_callback) config.init(pathlib.Path(__file__).parent) +def _signal_callback(sig, frame): + env.reset_artifact_dir() + + sig_name = signal.Signals(sig).name + logger.info(f"Signal callback: received {sig_name}") + sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" + with sig_file.open("a") as f: + f.write( + f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" + ) + + @click.group(cls=ci_lib.HelpfulGroup) @click.pass_context @ci_lib.safe_ci_function diff --git a/projects/skeleton/orchestration/test_skeleton.py b/projects/skeleton/orchestration/test_skeleton.py index eb267f5d9..355e745f8 100644 --- a/projects/skeleton/orchestration/test_skeleton.py +++ b/projects/skeleton/orchestration/test_skeleton.py @@ -1,7 +1,6 @@ import json import logging import pathlib -import signal import time from datetime import UTC, datetime @@ -95,28 +94,6 @@ def seed_skeleton_caliper_artifacts_with_data( return demo_dir -def _signal_handler_sigint(sig, frame): - """Sample SIGINT signal handler for skeleton project.""" - env.reset_artifact_dir() - # Sample handler - does nothing else - - -def _signal_handler_sigterm(sig, frame): - """Sample SIGTERM signal handler for skeleton project.""" - env.reset_artifact_dir() - # Sample handler - does nothing else - - -def _setup_sample_signal_handlers(): - """Set up sample signal handlers for demonstration.""" - try: - signal.signal(signal.SIGINT, _signal_handler_sigint) - signal.signal(signal.SIGTERM, _signal_handler_sigterm) - logger.debug("Sample signal handlers installed") - except Exception as e: - logger.warning(f"Failed to set up sample signal handlers: {e}") - - def test(): """Main test function that wraps do_test() with outcome postprocessing.""" return run_and_postprocess(do_test) From c690b230f1d42656ee52705548d90295edb244eb Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Fri, 25 Sep 2026 15:48:37 +0200 Subject: [PATCH 06/14] [llm_d] orchestration: test_phase: add signal handlers with logs --- projects/llm_d/orchestration/test_phase.py | 25 ++++++++++++++++++++++ 1 file changed, 25 insertions(+) diff --git a/projects/llm_d/orchestration/test_phase.py b/projects/llm_d/orchestration/test_phase.py index c7f7bff4e..7b9014a32 100644 --- a/projects/llm_d/orchestration/test_phase.py +++ b/projects/llm_d/orchestration/test_phase.py @@ -2,6 +2,7 @@ import logging import shutil +import signal import sys from datetime import UTC, datetime from pathlib import Path @@ -46,6 +47,30 @@ logger = logging.getLogger(__name__) +def _signal_handler_sigint(sig, frame): + """Sample SIGINT signal handler for skeleton project.""" + env.reset_artifact_dir() + logger.error("Sigint received.") + (env.ARTIFACT_DIR / "SIGINT").touch() + + +def _signal_handler_sigterm(sig, frame): + """Sample SIGTERM signal handler for skeleton project.""" + env.reset_artifact_dir() + logger.error("Sigterm received.") + (env.ARTIFACT_DIR / "SIGTERM").touch() + + +def _setup_sample_signal_handlers(): + """Set up sample signal handlers for demonstration.""" + try: + signal.signal(signal.SIGINT, _signal_handler_sigint) + signal.signal(signal.SIGTERM, _signal_handler_sigterm) + logger.debug("Sample signal handlers installed") + except Exception as e: + logger.warning(f"Failed to set up sample signal handlers: {e}") + + def _delete_resources_by_type(resource_type: str, namespace: str, description: str) -> None: """Delete all resources of a given type in the namespace. From 506121ecfc7c1186212259b3c75e7df58ed6d4e9 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 10:47:36 +0200 Subject: [PATCH 07/14] [llm_d] orchestration: enhance the signal handlers --- projects/llm_d/orchestration/ci.py | 17 +++++++++++++++ projects/llm_d/orchestration/test_phase.py | 25 ---------------------- 2 files changed, 17 insertions(+), 25 deletions(-) diff --git a/projects/llm_d/orchestration/ci.py b/projects/llm_d/orchestration/ci.py index 3d49d8467..22e70c842 100755 --- a/projects/llm_d/orchestration/ci.py +++ b/projects/llm_d/orchestration/ci.py @@ -5,7 +5,9 @@ """ import logging +import signal import types +from datetime import datetime from pathlib import Path import click @@ -43,6 +45,8 @@ def init(presets=None): env.init() run.init() + run.register_signal_callback(_signal_callback) + # Set presets configuration if provided if presets: config.write_variables_override(presets=presets) @@ -60,6 +64,19 @@ def init_vaults_for_phase(phase: str): vault.phase_vault_init(phase, extra_mandatory=(extra_mandatory or None)) +def _signal_callback(sig, frame): + + env.reset_artifact_dir() + + sig_name = signal.Signals(sig).name + logger.info(f"Signal callback: received {sig_name}") + sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" + with sig_file.open("a") as f: + f.write( + f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" + ) + + @click.group(cls=ci_lib.HelpfulGroup) @click.option("--preset", multiple=True, help="Set preset configuration before starting") @click.pass_context diff --git a/projects/llm_d/orchestration/test_phase.py b/projects/llm_d/orchestration/test_phase.py index 7b9014a32..c7f7bff4e 100644 --- a/projects/llm_d/orchestration/test_phase.py +++ b/projects/llm_d/orchestration/test_phase.py @@ -2,7 +2,6 @@ import logging import shutil -import signal import sys from datetime import UTC, datetime from pathlib import Path @@ -47,30 +46,6 @@ logger = logging.getLogger(__name__) -def _signal_handler_sigint(sig, frame): - """Sample SIGINT signal handler for skeleton project.""" - env.reset_artifact_dir() - logger.error("Sigint received.") - (env.ARTIFACT_DIR / "SIGINT").touch() - - -def _signal_handler_sigterm(sig, frame): - """Sample SIGTERM signal handler for skeleton project.""" - env.reset_artifact_dir() - logger.error("Sigterm received.") - (env.ARTIFACT_DIR / "SIGTERM").touch() - - -def _setup_sample_signal_handlers(): - """Set up sample signal handlers for demonstration.""" - try: - signal.signal(signal.SIGINT, _signal_handler_sigint) - signal.signal(signal.SIGTERM, _signal_handler_sigterm) - logger.debug("Sample signal handlers installed") - except Exception as e: - logger.warning(f"Failed to set up sample signal handlers: {e}") - - def _delete_resources_by_type(resource_type: str, namespace: str, description: str) -> None: """Delete all resources of a given type in the namespace. From 05088259eda1420b5f2f5e4bd941051c8c4c89ea Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 13:28:49 +0200 Subject: [PATCH 08/14] [core] library: run: expose a signal handler --- projects/core/library/run.py | 27 +++++++++++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/projects/core/library/run.py b/projects/core/library/run.py index 132e28391..0404aaeaa 100644 --- a/projects/core/library/run.py +++ b/projects/core/library/run.py @@ -5,6 +5,16 @@ logger = logging.getLogger(__name__) +_signal_callbacks = {signal.SIGINT: [], signal.SIGTERM: []} + + +def register_signal_callback(fn, *, sig=None): + if sig is None: + for callbacks in _signal_callbacks.values(): + callbacks.append(fn) + else: + _signal_callbacks[sig].append(fn) + def init(): signal.signal(signal.SIGINT, raise_signal) @@ -30,6 +40,23 @@ def __str__(self): def raise_signal(sig, frame): + logger.info(f"Raising signal {sig}") + + for cb in _signal_callbacks.get(sig, []): + try: + cb(sig, frame) + except Exception: + logger.exception(f"Signal callback {cb} failed") + + from datetime import datetime + + from projects.core.library import env + + sig_name = signal.Signals(sig).name + sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" + with sig_file.open("a") as f: + f.write(f"{datetime.now()}: {__name__}.{raise_signal.__qualname__} {sig_name} handler\n") + raise SignalInterrupt(sig, frame) From 322a892b3bc14feb43db5639a2505ee39dc3b5c1 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 14:09:22 +0200 Subject: [PATCH 09/14] [agentic_tools] ci_base: add signal handler placeholder --- projects/agentic_tools/ci_base.py | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/projects/agentic_tools/ci_base.py b/projects/agentic_tools/ci_base.py index 12a8c0f70..4b0e83ff6 100644 --- a/projects/agentic_tools/ci_base.py +++ b/projects/agentic_tools/ci_base.py @@ -11,9 +11,11 @@ import functools import importlib import logging +import signal import types from collections.abc import Callable from dataclasses import dataclass +from datetime import datetime from pathlib import Path import click @@ -95,8 +97,20 @@ def init(self, phase: str | None = None) -> None: """Bootstrap env, run, and config. Override to add custom init.""" env.init() run.init() + run.register_signal_callback(self._signal_callback) config.init(self.config_dir) + def _signal_callback(self, sig, frame): + env.reset_artifact_dir() + + sig_name = signal.Signals(sig).name + logger.info(f"Signal callback: received {sig_name}") + sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" + with sig_file.open("a") as f: + f.write( + f"{datetime.now()}: {__name__}.{type(self).__name__}.{self._signal_callback.__name__} {sig_name} handler\n" + ) + def init_vaults(self, phase: str) -> None: """Resolve and initialize vaults for *phase*. From e7a436f88998badfc913b30124c6172d14960df3 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 14:10:26 +0200 Subject: [PATCH 10/14] [inference_playbooks] orchestration: add signal handler placeholder --- projects/inference_playbooks/orchestration/ci.py | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/projects/inference_playbooks/orchestration/ci.py b/projects/inference_playbooks/orchestration/ci.py index 48c490720..2f64255b1 100755 --- a/projects/inference_playbooks/orchestration/ci.py +++ b/projects/inference_playbooks/orchestration/ci.py @@ -5,7 +5,9 @@ import logging import pathlib +import signal import types +from datetime import datetime import click import prepare_phase @@ -30,9 +32,22 @@ def init(): env.init() run.init() + run.register_signal_callback(_signal_callback) config.init(pathlib.Path(__file__).parent) +def _signal_callback(sig, frame): + env.reset_artifact_dir() + + sig_name = signal.Signals(sig).name + logger.info(f"Signal callback: received {sig_name}") + sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" + with sig_file.open("a") as f: + f.write( + f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" + ) + + @click.group(cls=ci_lib.HelpfulGroup) @click.pass_context @ci_lib.safe_ci_function From 4b587ae99598a79e634ec6d45265dc28312aabb2 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 14:11:06 +0200 Subject: [PATCH 11/14] [minimal] orchestration: add signal handler placeholder --- projects/minimal/orchestration/ci.py | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/projects/minimal/orchestration/ci.py b/projects/minimal/orchestration/ci.py index 66802ad81..76dbe3a37 100755 --- a/projects/minimal/orchestration/ci.py +++ b/projects/minimal/orchestration/ci.py @@ -5,7 +5,9 @@ import logging import pathlib +import signal import types +from datetime import datetime import click import prepare_phase @@ -30,9 +32,22 @@ def init(): env.init() run.init() + run.register_signal_callback(_signal_callback) config.init(pathlib.Path(__file__).parent) +def _signal_callback(sig, frame): + env.reset_artifact_dir() + + sig_name = signal.Signals(sig).name + logger.info(f"Signal callback: received {sig_name}") + sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" + with sig_file.open("a") as f: + f.write( + f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" + ) + + @click.group(cls=ci_lib.HelpfulGroup) @click.pass_context @ci_lib.safe_ci_function From 05639272614d3471100c103be06dd8751ad5fddb Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 14:11:37 +0200 Subject: [PATCH 12/14] [rhaiis] orchestration: add signal handler placeholder --- projects/rhaiis/orchestration/runtime_config.py | 15 +++++++++++++++ 1 file changed, 15 insertions(+) diff --git a/projects/rhaiis/orchestration/runtime_config.py b/projects/rhaiis/orchestration/runtime_config.py index ee0d3b228..1e0e5ea3b 100644 --- a/projects/rhaiis/orchestration/runtime_config.py +++ b/projects/rhaiis/orchestration/runtime_config.py @@ -2,6 +2,8 @@ import logging import pathlib +import signal +from datetime import datetime from projects.core.library import config, env, run @@ -13,9 +15,22 @@ def init() -> None: env.init() run.init() + run.register_signal_callback(_signal_callback) config.init(CONFIG_DIR) +def _signal_callback(sig, frame): + env.reset_artifact_dir() + + sig_name = signal.Signals(sig).name + logger.info(f"Signal callback: received {sig_name}") + sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" + with sig_file.open("a") as f: + f.write( + f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" + ) + + def get_namespace() -> str: return config.project.get_config("rhaiis.namespace") From 5a0f80f011eb8cf2e9e783c7f7de4f0263807675 Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 14:31:32 +0200 Subject: [PATCH 13/14] [core] ci_entrypoint: fournos: reduce the indentation level --- projects/core/ci_entrypoint/fournos.py | 29 +++++++++++++------------- 1 file changed, 14 insertions(+), 15 deletions(-) diff --git a/projects/core/ci_entrypoint/fournos.py b/projects/core/ci_entrypoint/fournos.py index 1473d35d5..d36b44a70 100644 --- a/projects/core/ci_entrypoint/fournos.py +++ b/projects/core/ci_entrypoint/fournos.py @@ -189,25 +189,24 @@ def transform_fournos_config_to_variable_overrides(fjob: dict) -> dict: fjob_engine = fjob_spec.get("executionEngine") forge_config = fjob_engine.get(FJOB_FORGE_ENGINE_NAME) - if forge_config: - # Process forge configuration - # Transform project -> project.name - if "project" in forge_config: - variable_overrides["project.name"] = forge_config["project"] - - # Transform args -> project.args - if "args" in forge_config: - variable_overrides["project.args"] = forge_config["args"] - - # Add all configOverrides entries directly (flatten them) - config_overrides = forge_config.get("configOverrides", {}) - variable_overrides.update(config_overrides) - - else: + if not forge_config: raise ValueError( f"Forge received an invalid fjob: spec.executionEngine.{FJOB_FORGE_ENGINE_NAME} not defined. Got {', '.join(fjob_engine.keys())}." ) + # Process forge configuration + # Transform project -> project.name + if "project" in forge_config: + variable_overrides["project.name"] = forge_config["project"] + + # Transform args -> project.args + if "args" in forge_config: + variable_overrides["project.args"] = forge_config["args"] + + # Add all configOverrides entries directly (flatten them) + config_overrides = forge_config.get("configOverrides", {}) + variable_overrides.update(config_overrides) + # Add ci_job mappings from spec if "exclusive" in fjob_spec: variable_overrides["ci_job.exclusive"] = fjob_spec["exclusive"] From cdc232ccc111cea4bd96ae7b83e4a6a6b7003a3f Mon Sep 17 00:00:00 2001 From: Kevin Pouget Date: Mon, 28 Sep 2026 14:50:22 +0200 Subject: [PATCH 14/14] [projects] orchestration: ci: improve the signal handler --- projects/agentic_tools/ci_base.py | 8 +- projects/core/ci_entrypoint/run_common.py | 80 ++++++++++++++----- projects/core/library/run.py | 39 ++++++--- .../inference_playbooks/orchestration/ci.py | 17 +++- projects/llm_d/orchestration/ci.py | 14 ++-- projects/minimal/orchestration/ci.py | 17 +++- .../rhaiis/orchestration/runtime_config.py | 8 +- projects/skeleton/orchestration/ci.py | 17 +++- 8 files changed, 145 insertions(+), 55 deletions(-) diff --git a/projects/agentic_tools/ci_base.py b/projects/agentic_tools/ci_base.py index 4b0e83ff6..f5f282340 100644 --- a/projects/agentic_tools/ci_base.py +++ b/projects/agentic_tools/ci_base.py @@ -100,13 +100,15 @@ def init(self, phase: str | None = None) -> None: run.register_signal_callback(self._signal_callback) config.init(self.config_dir) - def _signal_callback(self, sig, frame): + def _signal_callback(self, sig, frame, log_file): env.reset_artifact_dir() sig_name = signal.Signals(sig).name logger.info(f"Signal callback: received {sig_name}") - sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" - with sig_file.open("a") as f: + if not log_file: + return + + with log_file.open("a") as f: f.write( f"{datetime.now()}: {__name__}.{type(self).__name__}.{self._signal_callback.__name__} {sig_name} handler\n" ) diff --git a/projects/core/ci_entrypoint/run_common.py b/projects/core/ci_entrypoint/run_common.py index 0c4749657..8efcdfbfa 100644 --- a/projects/core/ci_entrypoint/run_common.py +++ b/projects/core/ci_entrypoint/run_common.py @@ -34,22 +34,58 @@ def setup_logging(): # Global reference to child process for signal forwarding _child_process = None +# Track which signals have already been forwarded to prevent re-entrant loops +_signals_forwarded = set() CHILD_SIGNAL_TIMEOUT = 30 -def _write_signal_file(sig_name): +def _write_signal_file(sig_name, exit_code): from datetime import datetime + import yaml + artifact_dir = os.environ.get("ARTIFACT_DIR") if not artifact_dir: return - sig_file = Path(artifact_dir) / f"{sig_name}_interrupted" + sig_file = Path(artifact_dir) / f"{sig_name}_interrupted.txt" with sig_file.open("a") as f: f.write(f"{datetime.now()}: {__name__}._forward_signal_and_exit {sig_name} handler\n") + # Write exit_status.yaml to the current step's ci_metadata + metadata_dir = Path(artifact_dir) / "000__ci_metadata" + metadata_dir.mkdir(parents=True, exist_ok=True) + exit_status_file = metadata_dir / "exit_status.yaml" + exit_status_data = { + "return_code": exit_code, + "reason": f"Aborted by signal {sig_name}", + } + with open(exit_status_file, "w", encoding="utf-8") as f: + yaml.dump(exit_status_data, f, default_flow_style=False) + logger.info(f"Wrote abort exit status to {exit_status_file}") + + +def _forward_signal_to_child(child_pid, sig, sig_name): + try: + os.killpg(child_pid, sig) + logger.info(f"Forwarded {sig_name} to child process group pgid={child_pid}") + except (OSError, ProcessLookupError) as e: + logger.info(f"Failed to forward {sig_name} to child pid={child_pid}: {e}") + return + + logger.info(f"Waiting up to {CHILD_SIGNAL_TIMEOUT}s for child pid={child_pid} to exit ...") + try: + _child_process.wait(timeout=CHILD_SIGNAL_TIMEOUT) + logger.info(f"Child pid={child_pid} exited after {sig_name}") + except subprocess.TimeoutExpired: + logger.warning( + f"Child pid={child_pid} did not exit within {CHILD_SIGNAL_TIMEOUT}s after {sig_name}" + ) + def _forward_signal_and_exit(sig, exit_code): + global _signals_forwarded + sig_name = signal.Signals(sig).name child_pid = _child_process.pid if _child_process else None child_alive = _child_process.poll() is None if _child_process else False @@ -58,25 +94,22 @@ def _forward_signal_and_exit(sig, exit_code): f"Received {sig_name} in pid={os.getpid()} pgid={os.getpgrp()}, child_pid={child_pid} child_alive={child_alive}" ) - _write_signal_file(sig_name) + if sig in _signals_forwarded: + logger.info(f"{sig_name} already forwarded, escalating to SIGTERM") + if _child_process and child_alive: + _forward_signal_to_child(child_pid, signal.SIGTERM, "SIGTERM") + prepare_ci.shutdown_dual_output() + sys.exit(exit_code) + + _signals_forwarded.add(sig) + _write_signal_file(sig_name, exit_code) if _child_process and child_alive: - try: - _child_process.send_signal(sig) - logger.info(f"Forwarded {sig_name} to child pid={child_pid}") - except (OSError, ProcessLookupError) as e: - logger.info(f"Failed to forward {sig_name} to child pid={child_pid}: {e}") - else: - logger.info( - f"Waiting up to {CHILD_SIGNAL_TIMEOUT}s for child pid={child_pid} to exit ..." - ) - try: - _child_process.wait(timeout=CHILD_SIGNAL_TIMEOUT) - logger.info(f"Child pid={child_pid} exited after {sig_name}") - except subprocess.TimeoutExpired: - logger.warning( - f"Child pid={child_pid} did not exit within {CHILD_SIGNAL_TIMEOUT}s after {sig_name}" - ) + _forward_signal_to_child(child_pid, sig, sig_name) + + if _child_process.poll() is None: + logger.info(f"Child still alive after {sig_name}, escalating to SIGTERM") + _forward_signal_to_child(child_pid, signal.SIGTERM, "SIGTERM") else: logger.info(f"No child to forward {sig_name} to") @@ -97,6 +130,7 @@ def signal_handler_sigterm(sig, frame): def setup_signal_handlers(): """Set up signal handlers for graceful interruption.""" try: + logger.info(f"Installing parent signal handlers for pid={os.getpid()} pgid={os.getpgrp()}") signal.signal(signal.SIGINT, signal_handler_sigint) signal.signal(signal.SIGTERM, signal_handler_sigterm) # SIGPIPE handling for broken pipes @@ -456,8 +490,12 @@ def execute_project_operation( logger.info(f"Parent pid={os.getpid()} pgid={os.getpgrp()}, child pid={_child_process.pid}") - # Wait for process to complete - result_code = _child_process.wait() + # Poll instead of blocking wait: Python's blocking waitpid uses + # SA_RESTART, which prevents signal handlers from firing. Polling + # with sleep allows SIGINT/SIGTERM handlers to run between checks. + while _child_process.poll() is None: + time.sleep(1) + result_code = _child_process.returncode # Create result object similar to subprocess.run() class Result: diff --git a/projects/core/library/run.py b/projects/core/library/run.py index 0404aaeaa..d47447cf0 100644 --- a/projects/core/library/run.py +++ b/projects/core/library/run.py @@ -10,13 +10,18 @@ def register_signal_callback(fn, *, sig=None): if sig is None: - for callbacks in _signal_callbacks.values(): + for sig_type, callbacks in _signal_callbacks.items(): callbacks.append(fn) + logger.info( + f"Registered signal callback {fn.__qualname__} for {signal.Signals(sig_type).name}" + ) else: _signal_callbacks[sig].append(fn) + logger.info(f"Registered signal callback {fn.__qualname__} for {signal.Signals(sig).name}") def init(): + logger.info(f"Installing signal handlers for pid={os.getpid()} pgid={os.getpgrp()}") signal.signal(signal.SIGINT, raise_signal) signal.signal(signal.SIGTERM, raise_signal) @@ -40,22 +45,34 @@ def __str__(self): def raise_signal(sig, frame): - logger.info(f"Raising signal {sig}") + sig_name = signal.Signals(sig).name + # Write to stderr immediately, before any import or env access + import sys - for cb in _signal_callbacks.get(sig, []): - try: - cb(sig, frame) - except Exception: - logger.exception(f"Signal callback {cb} failed") + print(f"raise_signal: {sig_name} in pid={os.getpid()}", file=sys.stderr, flush=True) + logger.info(f"Received {sig_name} in pid={os.getpid()}") from datetime import datetime from projects.core.library import env - sig_name = signal.Signals(sig).name - sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" - with sig_file.open("a") as f: - f.write(f"{datetime.now()}: {__name__}.{raise_signal.__qualname__} {sig_name} handler\n") + log_file = None + if env.BASE_ARTIFACT_DIR and env.BASE_ARTIFACT_DIR.is_dir(): + log_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted.txt" + else: + logger.warning("BASE_ARTIFACT_DIR not available, cannot write signal file") + + for cb in _signal_callbacks.get(sig, []): + try: + cb(sig, frame, log_file) + except Exception: + logger.exception(f"Signal callback {cb} failed") + + if log_file: + with log_file.open("a") as f: + f.write( + f"{datetime.now()}: {__name__}.{raise_signal.__qualname__} {sig_name} handler\n" + ) raise SignalInterrupt(sig, frame) diff --git a/projects/inference_playbooks/orchestration/ci.py b/projects/inference_playbooks/orchestration/ci.py index 2f64255b1..8d2de9cc8 100755 --- a/projects/inference_playbooks/orchestration/ci.py +++ b/projects/inference_playbooks/orchestration/ci.py @@ -36,15 +36,24 @@ def init(): config.init(pathlib.Path(__file__).parent) -def _signal_callback(sig, frame): +def _signal_callback(sig, frame, log_file): env.reset_artifact_dir() sig_name = signal.Signals(sig).name logger.info(f"Signal callback: received {sig_name}") - sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" - with sig_file.open("a") as f: + if not log_file: + return + + module_name = ( + pathlib.Path(__file__) + .relative_to(env.FORGE_HOME) + .with_suffix("") + .as_posix() + .replace("/", ".") + ) + with log_file.open("a") as f: f.write( - f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" + f"{datetime.now()}: {module_name}.{_signal_callback.__qualname__} {sig_name} handler\n" ) diff --git a/projects/llm_d/orchestration/ci.py b/projects/llm_d/orchestration/ci.py index 22e70c842..668746d72 100755 --- a/projects/llm_d/orchestration/ci.py +++ b/projects/llm_d/orchestration/ci.py @@ -64,16 +64,20 @@ def init_vaults_for_phase(phase: str): vault.phase_vault_init(phase, extra_mandatory=(extra_mandatory or None)) -def _signal_callback(sig, frame): - +def _signal_callback(sig, frame, log_file): env.reset_artifact_dir() sig_name = signal.Signals(sig).name logger.info(f"Signal callback: received {sig_name}") - sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" - with sig_file.open("a") as f: + if not log_file: + return + + module_name = ( + Path(__file__).relative_to(env.FORGE_HOME).with_suffix("").as_posix().replace("/", ".") + ) + with log_file.open("a") as f: f.write( - f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" + f"{datetime.now()}: {module_name}.{_signal_callback.__qualname__} {sig_name} handler\n" ) diff --git a/projects/minimal/orchestration/ci.py b/projects/minimal/orchestration/ci.py index 76dbe3a37..d29d7cb9a 100755 --- a/projects/minimal/orchestration/ci.py +++ b/projects/minimal/orchestration/ci.py @@ -36,15 +36,24 @@ def init(): config.init(pathlib.Path(__file__).parent) -def _signal_callback(sig, frame): +def _signal_callback(sig, frame, log_file): env.reset_artifact_dir() sig_name = signal.Signals(sig).name logger.info(f"Signal callback: received {sig_name}") - sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" - with sig_file.open("a") as f: + if not log_file: + return + + module_name = ( + pathlib.Path(__file__) + .relative_to(env.FORGE_HOME) + .with_suffix("") + .as_posix() + .replace("/", ".") + ) + with log_file.open("a") as f: f.write( - f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" + f"{datetime.now()}: {module_name}.{_signal_callback.__qualname__} {sig_name} handler\n" ) diff --git a/projects/rhaiis/orchestration/runtime_config.py b/projects/rhaiis/orchestration/runtime_config.py index 1e0e5ea3b..ffc75cac8 100644 --- a/projects/rhaiis/orchestration/runtime_config.py +++ b/projects/rhaiis/orchestration/runtime_config.py @@ -19,13 +19,15 @@ def init() -> None: config.init(CONFIG_DIR) -def _signal_callback(sig, frame): +def _signal_callback(sig, frame, log_file): env.reset_artifact_dir() sig_name = signal.Signals(sig).name logger.info(f"Signal callback: received {sig_name}") - sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" - with sig_file.open("a") as f: + if not log_file: + return + + with log_file.open("a") as f: f.write( f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" ) diff --git a/projects/skeleton/orchestration/ci.py b/projects/skeleton/orchestration/ci.py index 08fa6a3a0..d059f5667 100755 --- a/projects/skeleton/orchestration/ci.py +++ b/projects/skeleton/orchestration/ci.py @@ -40,15 +40,24 @@ def init(): config.init(pathlib.Path(__file__).parent) -def _signal_callback(sig, frame): +def _signal_callback(sig, frame, log_file): env.reset_artifact_dir() sig_name = signal.Signals(sig).name logger.info(f"Signal callback: received {sig_name}") - sig_file = env.BASE_ARTIFACT_DIR / f"{sig_name}_interrupted" - with sig_file.open("a") as f: + if not log_file: + return + + module_name = ( + pathlib.Path(__file__) + .relative_to(env.FORGE_HOME) + .with_suffix("") + .as_posix() + .replace("/", ".") + ) + with log_file.open("a") as f: f.write( - f"{datetime.now()}: {__name__}.{_signal_callback.__qualname__} {sig_name} handler\n" + f"{datetime.now()}: {module_name}.{_signal_callback.__qualname__} {sig_name} handler\n" )