Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 16 additions & 0 deletions projects/agentic_tools/ci_base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -95,8 +97,22 @@ 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, log_file):
env.reset_artifact_dir()

sig_name = signal.Signals(sig).name
logger.info(f"Signal callback: received {sig_name}")
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"
)

def init_vaults(self, phase: str) -> None:
"""Resolve and initialize vaults for *phase*.

Expand Down
92 changes: 50 additions & 42 deletions projects/core/ci_entrypoint/fournos.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
Comment thread
kpouget marked this conversation as resolved.

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:
Expand Down Expand Up @@ -180,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"]
Expand Down
115 changes: 89 additions & 26 deletions projects/core/ci_entrypoint/run_common.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,46 +34,103 @@ 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()

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, 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.txt"
with sig_file.open("a") as f:
f.write(f"{datetime.now()}: {__name__}._forward_signal_and_exit {sig_name} handler\n")
Comment thread
kpouget marked this conversation as resolved.

# 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

sys.exit(130) # Standard exit code for SIGINT
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 signal_handler_sigterm(sig, frame):
"""Handle SIGTERM gracefully."""
logger.info("🛑 Received SIGTERM - Terminating operation...")
def _forward_signal_and_exit(sig, exit_code):
global _signals_forwarded

# 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...")
try:
_child_process.send_signal(signal.SIGTERM)
except (OSError, ProcessLookupError):
pass # Child may have already terminated
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

logger.info(
f"Received {sig_name} in pid={os.getpid()} pgid={os.getpgrp()}, child_pid={child_pid} child_alive={child_alive}"
)

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:
_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")

# Emergency cleanup of dual output
prepare_ci.shutdown_dual_output()

sys.exit(143) # Standard exit code for SIGTERM
sys.exit(exit_code)
Comment thread
kpouget marked this conversation as resolved.


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():
"""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
Expand Down Expand Up @@ -431,8 +488,14 @@ def execute_project_operation(
stderr=None, # Inherit stderr for pdb/debugging
)

# Wait for process to complete
result_code = _child_process.wait()
logger.info(f"Parent pid={os.getpid()} pgid={os.getpgrp()}, child pid={_child_process.pid}")

# 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:
Expand Down
5 changes: 3 additions & 2 deletions projects/core/library/export.py
Original file line number Diff line number Diff line change
Expand Up @@ -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}")
Expand All @@ -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}")
Expand Down Expand Up @@ -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",
]:
Expand Down
15 changes: 12 additions & 3 deletions projects/core/library/export_notifications.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading
Loading