From c47f70c046221fe954d9d03c21a301d04e437ddf Mon Sep 17 00:00:00 2001 From: Alex Edmondson Date: Fri, 11 Sep 2026 09:09:38 -0400 Subject: [PATCH] Fail the run when an R statistics step fails or times out On the FORGE treatment re-run, roi_directed (in all three source arms) and roi_cross_freq's AAC/PPC tier were killed by a hard-coded one-hour R timeout. Each module logged it and the command exited 0, so region tables from the previous code version went on looking current. Only checking the outputs found it. - R statistics steps have no default time limit. They were hard-coded to 3600 s, or 600 s for the electrode and evoked modules; r_timeout_sec in a module's config block sets one. - A failed or timed-out step raises RStepFailed and the CLI exits 1. A run over a paradigm, or over every paradigm, finishes its other modules first and then lists what failed. A missing Rscript still only logs, as before. - roi_cross_freq attempts PAC and AAC/PPC before failing, so one tier failing does not cost the other its tables. - roi_cross_freq's PAC mosaics named legacy columns (hedges_g, region, contrast, freq_pair) that the native table does not have, so none was ever drawn. They now read effect_size / spatial / hypothesis / band, and the figures step draws them (summary() did only when figures were requested in the same run). The three vertex modules lose their limit too but still only log a failure; that is left for the vertex split. Tests: 230 pass (8 new), covering the missing default, a configured limit, a timeout and a failure each raising, both cross-frequency tiers attempted, the PAC mosaic columns and figures pass, and the CLI finishing a batch before exiting 1. Co-Authored-By: Claude Opus 5 Claude-Session: https://claude.ai/code/session_01RibM8Zep2YEjjUbc2LLkgj --- CHANGELOG.md | 24 +++ README.md | 1 + src/source_analytics/analyses/base.py | 26 +++ .../analyses/electrode_analysis.py | 6 +- .../analyses/electrode_aperiodic_analysis.py | 8 +- .../analyses/electrode_comparison_analysis.py | 6 +- .../analyses/electrode_evoked_analysis.py | 8 +- .../analyses/roi_aperiodic_analysis.py | 8 +- .../analyses/roi_connectivity_analysis.py | 6 +- .../analyses/roi_cross_freq_analysis.py | 49 ++++-- .../analyses/roi_directed_analysis.py | 6 +- .../analyses/roi_evoked_analysis.py | 6 +- .../analyses/roi_psd_analysis.py | 8 +- .../analyses/vertex_cluster_analysis.py | 6 +- .../analyses/vertex_connectivity_analysis.py | 2 +- .../analyses/vertex_signature_analysis.py | 2 +- .../analyses/vertex_specparam_analysis.py | 2 +- src/source_analytics/cli.py | 42 ++++- tests/test_r_step_failures.py | 158 ++++++++++++++++++ 19 files changed, 315 insertions(+), 59 deletions(-) create mode 100644 tests/test_r_step_failures.py diff --git a/CHANGELOG.md b/CHANGELOG.md index beb9aa9..9e83b40 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,29 @@ # Changelog +## Unreleased + +### Behaviour changes (read these before re-running a study) + +- **A failed or timed-out R statistics step now fails the run** (exit 1). A run over a + whole paradigm still finishes its other modules first, then lists the failures. It + used to log the failure and exit 0: on the FORGE treatment re-run, roi_directed (all + three source arms) and roi_cross_freq's AAC/PPC tier were killed by a one-hour limit, + and region tables from the previous code version went on looking current. +- **R statistics steps have no default time limit.** They were hard-coded to 3600 s + (roi_psd, roi_aperiodic, roi_connectivity, roi_directed, roi_cross_freq, + vertex_cluster) or 600 s (electrode and evoked modules). Set `r_timeout_sec` in a + module's config block to impose one. The three vertex modules with their own limit + lose it too; unlike the rest they still only log a failed step, which is left for + the vertex split. + +### Fixed + +- **roi_cross_freq draws its PAC mosaics.** The mosaic call named legacy columns + (`hedges_g`, `region`, `contrast`, `freq_pair`) that the native hypothesis table + does not have, so none was ever drawn. It now reads `effect_size` / `spatial` / + `hypothesis` / `band`, and the figures step draws them; `summary()` only did so + when figures were requested in the same run. + ## v0.7.0 — 2026-09-10 (audit remediation, per-atlas resolution, roi_signature) A repo audit (2026-09-04) compared the README/CLAUDE.md against the code and diff --git a/README.md b/README.md index 53d3150..4a05b04 100644 --- a/README.md +++ b/README.md @@ -395,6 +395,7 @@ paradigms: | `bands` | all spectral/connectivity | frequency bands analysed | | `pipeline.atlas`, `atlas_files`, `atlas_dir` | atlas I/O, R region tier, mosaics | which parcellation the ROI data use: resolved by name through source-localization's `registry.yaml` to that atlas's own labels / mapping / categories / anatomy; `atlas_files` names the files of an unregistered atlas. The 10× voxel convention is read from each NIfTI header, never inferred from its filename | | `roi_categories` | region tier (Python + R), mosaics | category → ROI map. The study's map (or a profile's narrowing) always wins over the atlas default, on the Python and R sides alike | +| `.r_timeout_sec` | R-backed analyses | wall-clock limit for that module's R statistics step; unset = no limit. A step that fails or times out fails the run (exit 1); a paradigm-wide run finishes its other modules first | | `epoch_sampling` | spectral/connectivity, all levels | random-epoch resampling (`n_bootstrap: 0` = full timeseries). Precedence: global → `vertex.epoch_sampling` → per-analysis block | | `jobs` | `run --jobs` default | worker count when `--jobs` is not given | | `.{include_analyses, include_hypotheses, bands, rois}` | `run --profile` | a narrowed study written to its own tree (see below) | diff --git a/src/source_analytics/analyses/base.py b/src/source_analytics/analyses/base.py index 241da1d..29069b4 100644 --- a/src/source_analytics/analyses/base.py +++ b/src/source_analytics/analyses/base.py @@ -12,6 +12,14 @@ logger = logging.getLogger(__name__) + +class RStepFailed(RuntimeError): + """An R statistics step failed or timed out, so its tables were not (re)written. + + Raised rather than logged: a failed or timed-out R step used to leave the run + exiting 0, so a batch reported success over tables that were never rewritten. + """ + VALID_STEPS = {"setup", "process", "aggregate", "statistics", "figures", "summary"} DEFAULT_RUN_STEPS = VALID_STEPS - {"figures"} @@ -509,6 +517,24 @@ def _call_r_figures_only(self, r_script_name: str, data_csv: str) -> bool: logger.error("R figures-only timed out after 300s") return False + @property + def _r_timeout(self) -> float | None: + """Wall-clock limit for this module's R statistics step, in seconds, or None. + + Set per module with ``r_timeout_sec`` in its config block. There is no + default: these steps are long -- roi_directed's directed-edge tier runs past + an hour on 26 parcels -- and a limit that silently kills them is worse than + waiting for them. + """ + value = self.config.raw.get(self.name, {}).get("r_timeout_sec") + return float(value) if value else None + + def _r_step_failed(self, msg: str, *args) -> None: + """Log an R statistics-step failure and raise :class:`RStepFailed`.""" + text = msg % args if args else msg + logger.error(text) + raise RStepFailed(f"{self.name}: {text}") + def _r_roi_categories_flags(self) -> list[str]: """``['--roi-categories', path]`` for the resolved atlas's own category file. diff --git a/src/source_analytics/analyses/electrode_analysis.py b/src/source_analytics/analyses/electrode_analysis.py index b1e984e..500bcb9 100644 --- a/src/source_analytics/analyses/electrode_analysis.py +++ b/src/source_analytics/analyses/electrode_analysis.py @@ -337,7 +337,7 @@ def summary(self) -> None: logger.info("Calling R: %s", " ".join(cmd)) try: result = subprocess.run( - cmd, capture_output=True, text=True, timeout=600, + cmd, capture_output=True, text=True, timeout=self._r_timeout, ) if result.stdout: for line in result.stdout.strip().split("\n"): @@ -347,8 +347,8 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error("R script failed with exit code %d", result.returncode) + self._r_step_failed("R script failed with exit code %d", result.returncode) except FileNotFoundError: logger.error("Rscript not found. Install R to enable statistics and visualization.") except subprocess.TimeoutExpired: - logger.error("R script timed out after 600 seconds") + self._r_step_failed("R script timed out after %s s", self._r_timeout) diff --git a/src/source_analytics/analyses/electrode_aperiodic_analysis.py b/src/source_analytics/analyses/electrode_aperiodic_analysis.py index c0e2678..55a6eed 100644 --- a/src/source_analytics/analyses/electrode_aperiodic_analysis.py +++ b/src/source_analytics/analyses/electrode_aperiodic_analysis.py @@ -280,7 +280,7 @@ def summary(self) -> None: logger.info("Calling R: %s", " ".join(cmd)) try: result = subprocess.run( - cmd, capture_output=True, text=True, timeout=600, + cmd, capture_output=True, text=True, timeout=self._r_timeout, ) if result.stdout: for line in result.stdout.strip().split("\n"): @@ -290,12 +290,10 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error( - "R script failed with exit code %d", result.returncode, - ) + self._r_step_failed("R script failed with exit code %d", result.returncode) except FileNotFoundError: logger.error( "Rscript not found. Install R to enable statistics.", ) except subprocess.TimeoutExpired: - logger.error("R script timed out after 600 seconds") + self._r_step_failed("R script timed out after %s s", self._r_timeout) diff --git a/src/source_analytics/analyses/electrode_comparison_analysis.py b/src/source_analytics/analyses/electrode_comparison_analysis.py index 08ef7de..974db67 100644 --- a/src/source_analytics/analyses/electrode_comparison_analysis.py +++ b/src/source_analytics/analyses/electrode_comparison_analysis.py @@ -675,7 +675,7 @@ def summary(self) -> None: logger.info("Calling R: %s", " ".join(cmd)) try: result = subprocess.run( - cmd, capture_output=True, text=True, timeout=600, + cmd, capture_output=True, text=True, timeout=self._r_timeout, ) if result.stdout: for line in result.stdout.strip().split("\n"): @@ -685,8 +685,8 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error("R script failed with exit code %d", result.returncode) + self._r_step_failed("R script failed with exit code %d", result.returncode) except FileNotFoundError: logger.error("Rscript not found. Install R to enable comparison report.") except subprocess.TimeoutExpired: - logger.error("R script timed out after 600 seconds") + self._r_step_failed("R script timed out after %s s", self._r_timeout) diff --git a/src/source_analytics/analyses/electrode_evoked_analysis.py b/src/source_analytics/analyses/electrode_evoked_analysis.py index bdcbc62..b514d71 100644 --- a/src/source_analytics/analyses/electrode_evoked_analysis.py +++ b/src/source_analytics/analyses/electrode_evoked_analysis.py @@ -428,7 +428,7 @@ def summary(self) -> None: cmd, capture_output=True, text=True, - timeout=600, + timeout=self._r_timeout, ) if result.stdout: for line in result.stdout.strip().split("\n"): @@ -438,12 +438,10 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error( - "R script failed with exit code %d", result.returncode - ) + self._r_step_failed("R script failed with exit code %d", result.returncode) except FileNotFoundError: logger.error( "Rscript not found. Install R to enable statistics and visualization." ) except subprocess.TimeoutExpired: - logger.error("R script timed out after 600 seconds") + self._r_step_failed("R script timed out after %s s", self._r_timeout) diff --git a/src/source_analytics/analyses/roi_aperiodic_analysis.py b/src/source_analytics/analyses/roi_aperiodic_analysis.py index 8db164e..77cdbc2 100644 --- a/src/source_analytics/analyses/roi_aperiodic_analysis.py +++ b/src/source_analytics/analyses/roi_aperiodic_analysis.py @@ -195,7 +195,7 @@ def summary(self) -> None: cmd, capture_output=True, text=True, - timeout=3600, + timeout=self._r_timeout, ) if result.stdout: for line in result.stdout.strip().split("\n"): @@ -205,11 +205,11 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error("R script failed with exit code %d", result.returncode) + self._r_step_failed("R script failed with exit code %d", result.returncode) except FileNotFoundError: logger.error("Rscript not found. Install R to enable statistics and visualization.") except subprocess.TimeoutExpired: - logger.error("R script timed out after 3600 seconds") + self._r_step_failed("R script timed out after %s s", self._r_timeout) # Render brain mosaics from posthoc effect sizes if self._generate_figures: @@ -244,5 +244,5 @@ def _render_brain_mosaics(self) -> None: correction_label="FDR", facet_cols=["hypothesis", "dv"], colorbar_label="Hedges' g", - auto_slices=True,atlas=self._atlas_dir + auto_slices=True, atlas=self._atlas_dir ) diff --git a/src/source_analytics/analyses/roi_connectivity_analysis.py b/src/source_analytics/analyses/roi_connectivity_analysis.py index 25ed606..95ea7b6 100644 --- a/src/source_analytics/analyses/roi_connectivity_analysis.py +++ b/src/source_analytics/analyses/roi_connectivity_analysis.py @@ -343,7 +343,7 @@ def summary(self) -> None: cmd, capture_output=True, text=True, - timeout=3600, + timeout=self._r_timeout, ) if result.stdout: for line in result.stdout.strip().split("\n"): @@ -353,10 +353,10 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error("R script failed with exit code %d", result.returncode) + self._r_step_failed("R script failed with exit code %d", result.returncode) except FileNotFoundError: logger.error( "Rscript not found. Install R to enable statistics and visualization." ) except subprocess.TimeoutExpired: - logger.error("R script timed out after 3600 seconds") + self._r_step_failed("R script timed out after %s s", self._r_timeout) diff --git a/src/source_analytics/analyses/roi_cross_freq_analysis.py b/src/source_analytics/analyses/roi_cross_freq_analysis.py index 4f80b12..c00e3f6 100644 --- a/src/source_analytics/analyses/roi_cross_freq_analysis.py +++ b/src/source_analytics/analyses/roi_cross_freq_analysis.py @@ -251,34 +251,46 @@ def figures(self) -> None: for m in self._edge_metrics_on_disk(): self._call_r_figures_only("roi_cross_freq_edges_analysis.R", f"{m}_edges.csv") break # one call handles every edge CSV present + if "pac" in self._metrics: + # From the tables on disk, as roi_psd does: summary() drew these only + # when figures were requested in the same run, so a separate figures + # pass (the usual way to redraw) never did. + self._render_brain_mosaics() # -------------------------------------------------------------- summary def summary(self) -> None: """Statistics + figures + report via R (PAC first, then AAC/PPC).""" - if "pac" in self._metrics: - self._run_pac_r() + failed = [] + if "pac" in self._metrics and not self._run_pac_r(): + failed.append("PAC") edge_metrics = self._edge_metrics_on_disk() - if edge_metrics: - self._run_edges_r(edge_metrics) - - def _run_edges_r(self, metrics: list[str]) -> None: + if edge_metrics and not self._run_edges_r(edge_metrics): + failed.append("AAC/PPC") + if failed: + # Both tiers were attempted first, so one failing does not cost the + # other its tables. + self._r_step_failed("R step(s) failed: %s", ", ".join(failed)) + + def _run_edges_r(self, metrics: list[str]) -> bool: """AAC/PPC hypotheses (global / cells / region pairs) via R.""" - self._run_r_script( + return self._run_r_script( "roi_cross_freq_edges_analysis.R", extra_args=["--metric", ",".join(metrics)], label="AAC/PPC", ) - def _run_pac_r(self) -> None: + def _run_pac_r(self) -> bool: data_dir = self.output_dir / "data" if not (data_dir / "pac_values.csv").exists(): logger.error("pac_values.csv not found -- skipping PAC R analysis") - return - if self._run_r_script("roi_pac_analysis.R", label="PAC") and self._generate_figures: + return True # nothing to test is a skip, not a failure + ok = self._run_r_script("roi_pac_analysis.R", label="PAC") + if ok and self._generate_figures: self._render_brain_mosaics() + return ok def _run_r_script(self, script_name: str, *, extra_args: list[str] | None = None, - label: str = "R", timeout: int = 3600) -> bool: + label: str = "R", timeout: float | None = None) -> bool: """Run one of this module's R scripts with the standard argument set. Returns True when the script exited 0. @@ -320,9 +332,10 @@ def _run_r_script(self, script_name: str, *, extra_args: list[str] | None = None if extra_args: cmd.extend(extra_args) + limit = timeout if timeout is not None else self._r_timeout logger.info("Calling R (%s): %s", label, " ".join(cmd)) try: - result = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout) + result = subprocess.run(cmd, capture_output=True, text=True, timeout=limit) for stream in (result.stdout, result.stderr): if stream: for line in stream.strip().split("\n"): @@ -335,7 +348,7 @@ def _run_r_script(self, script_name: str, *, extra_args: list[str] | None = None except FileNotFoundError: logger.error("Rscript not found. Install R to enable %s statistics.", label) except subprocess.TimeoutExpired: - logger.error("%s R script timed out after %d seconds", label, timeout) + logger.error("%s R script timed out after %s s", label, limit) return False def _render_brain_mosaics(self) -> None: @@ -353,7 +366,11 @@ def _render_brain_mosaics(self) -> None: render_posthoc_mosaics( posthoc_csv, roi_cats, self.fig_dir, analysis_name="roi_cross_freq", - effect_col="hedges_g", roi_col="region", - facet_cols=["contrast", "freq_pair"], - colorbar_label="Hedges' g",atlas=self._atlas_dir + # The native hypothesis schema: `spatial` holds the region and `band` + # the frequency pair. The legacy names (hedges_g / region / contrast / + # freq_pair) matched no column, so no PAC mosaic was ever drawn. + effect_col="effect_size", roi_col="spatial", + p_col="p_value", q_col="q_value", correction_label="FDR", + facet_cols=["hypothesis", "band"], + colorbar_label="Hedges' g", atlas=self._atlas_dir, ) diff --git a/src/source_analytics/analyses/roi_directed_analysis.py b/src/source_analytics/analyses/roi_directed_analysis.py index cd8ab84..f61fb4c 100644 --- a/src/source_analytics/analyses/roi_directed_analysis.py +++ b/src/source_analytics/analyses/roi_directed_analysis.py @@ -234,7 +234,7 @@ def summary(self) -> None: if wanted_hyp: cmd.extend(["--hypothesis", ",".join(sorted(wanted_hyp))]) - r_timeout = 3600 + r_timeout = self._r_timeout logger.info("Calling R: %s", " ".join(cmd)) try: result = subprocess.run( @@ -251,10 +251,10 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error("R script failed with exit code %d", result.returncode) + self._r_step_failed("R script failed with exit code %d", result.returncode) except FileNotFoundError: logger.error( "Rscript not found. Install R to enable statistics and visualization." ) except subprocess.TimeoutExpired: - logger.error("R script timed out after %d seconds", r_timeout) + self._r_step_failed("R script timed out after %s s", self._r_timeout) diff --git a/src/source_analytics/analyses/roi_evoked_analysis.py b/src/source_analytics/analyses/roi_evoked_analysis.py index d03414a..1ba1f87 100644 --- a/src/source_analytics/analyses/roi_evoked_analysis.py +++ b/src/source_analytics/analyses/roi_evoked_analysis.py @@ -321,7 +321,7 @@ def summary(self) -> None: cmd, capture_output=True, text=True, - timeout=600, + timeout=self._r_timeout, ) if result.stdout: for line in result.stdout.strip().split("\n"): @@ -331,8 +331,8 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error("R script failed with exit code %d", result.returncode) + self._r_step_failed("R script failed with exit code %d", result.returncode) except FileNotFoundError: logger.error("Rscript not found. Install R to enable statistics and visualization.") except subprocess.TimeoutExpired: - logger.error("R script timed out after 600 seconds") + self._r_step_failed("R script timed out after %s s", self._r_timeout) diff --git a/src/source_analytics/analyses/roi_psd_analysis.py b/src/source_analytics/analyses/roi_psd_analysis.py index 9309bca..18c6eae 100644 --- a/src/source_analytics/analyses/roi_psd_analysis.py +++ b/src/source_analytics/analyses/roi_psd_analysis.py @@ -231,7 +231,7 @@ def summary(self) -> None: cmd, capture_output=True, text=True, - timeout=3600, + timeout=self._r_timeout, ) if result.stdout: for line in result.stdout.strip().split("\n"): @@ -241,11 +241,11 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error("R script failed with exit code %d", result.returncode) + self._r_step_failed("R script failed with exit code %d", result.returncode) except FileNotFoundError: logger.error("Rscript not found. Install R to enable statistics and visualization.") except subprocess.TimeoutExpired: - logger.error("R script timed out after 3600 seconds") + self._r_step_failed("R script timed out after %s s", self._r_timeout) # Render brain mosaics from posthoc effect sizes if self._generate_figures: @@ -280,5 +280,5 @@ def _render_brain_mosaics(self) -> None: correction_label="FDR", facet_cols=["hypothesis", "band", "dv"], colorbar_label="Hedges' g", - auto_slices=True,atlas=self._atlas_dir + auto_slices=True, atlas=self._atlas_dir ) diff --git a/src/source_analytics/analyses/vertex_cluster_analysis.py b/src/source_analytics/analyses/vertex_cluster_analysis.py index 88a8bfb..9e1ae13 100644 --- a/src/source_analytics/analyses/vertex_cluster_analysis.py +++ b/src/source_analytics/analyses/vertex_cluster_analysis.py @@ -697,7 +697,7 @@ def summary(self) -> None: logger.info("Calling R: %s", " ".join(cmd)) try: result = subprocess.run( - cmd, capture_output=True, text=True, timeout=3600, + cmd, capture_output=True, text=True, timeout=self._r_timeout, ) if result.stdout: for line in result.stdout.strip().split("\n"): @@ -707,13 +707,13 @@ def summary(self) -> None: if line.strip(): logger.info("[R] %s", line) if result.returncode != 0: - logger.error("R script failed with exit code %d", result.returncode) + self._r_step_failed("R script failed with exit code %d", result.returncode) self._write_python_summary() except FileNotFoundError: logger.warning("Rscript not found — writing Python summary") self._write_python_summary() except subprocess.TimeoutExpired: - logger.error("R script timed out after 3600 seconds") + self._r_step_failed("R script timed out after %s s", self._r_timeout) self._write_python_summary() def _write_python_summary(self) -> None: diff --git a/src/source_analytics/analyses/vertex_connectivity_analysis.py b/src/source_analytics/analyses/vertex_connectivity_analysis.py index 7c63c18..adcdb95 100644 --- a/src/source_analytics/analyses/vertex_connectivity_analysis.py +++ b/src/source_analytics/analyses/vertex_connectivity_analysis.py @@ -434,7 +434,7 @@ def summary(self) -> None: ] cmd.extend(self._r_no_figures_flags()) result = subprocess.run( - cmd, capture_output=True, text=True, timeout=3600, + cmd, capture_output=True, text=True, timeout=self._r_timeout, ) if result.returncode == 0: return diff --git a/src/source_analytics/analyses/vertex_signature_analysis.py b/src/source_analytics/analyses/vertex_signature_analysis.py index 838bfac..95ea920 100644 --- a/src/source_analytics/analyses/vertex_signature_analysis.py +++ b/src/source_analytics/analyses/vertex_signature_analysis.py @@ -389,7 +389,7 @@ def summary(self) -> None: "--tbl-dir", str(self.tbl_dir), ] cmd.extend(self._r_no_figures_flags()) - result = subprocess.run(cmd, capture_output=True, text=True, timeout=3600) + result = subprocess.run(cmd, capture_output=True, text=True, timeout=self._r_timeout) if result.returncode == 0: return except (FileNotFoundError, subprocess.TimeoutExpired): diff --git a/src/source_analytics/analyses/vertex_specparam_analysis.py b/src/source_analytics/analyses/vertex_specparam_analysis.py index 4cfc789..a63bcb5 100644 --- a/src/source_analytics/analyses/vertex_specparam_analysis.py +++ b/src/source_analytics/analyses/vertex_specparam_analysis.py @@ -719,7 +719,7 @@ def summary(self) -> None: "--tbl-dir", str(self.tbl_dir), ] cmd.extend(self._r_no_figures_flags()) - result = subprocess.run(cmd, capture_output=True, text=True, timeout=3600) + result = subprocess.run(cmd, capture_output=True, text=True, timeout=self._r_timeout) if result.returncode == 0: return except (FileNotFoundError, subprocess.TimeoutExpired): diff --git a/src/source_analytics/cli.py b/src/source_analytics/cli.py index cf22262..c6f1152 100644 --- a/src/source_analytics/cli.py +++ b/src/source_analytics/cli.py @@ -11,6 +11,7 @@ from .config import StudyConfig from .core import StudyAnalyzer, ANALYSIS_REGISTRY, ANALYSIS_METADATA, canonical_analysis_name +from .analyses.base import RStepFailed from .analyses.base import VALID_STEPS, BaseAnalysis @@ -36,6 +37,20 @@ def _print_study_summary(config: StudyConfig, analyzer: StudyAnalyzer): print() +def _exit_if_failed(failures: list[str]) -> None: + """Exit 1 after listing the R statistics steps that failed. + + A failed or timed-out R step used to be logged while the run exited 0, so a + batch reported success over tables that were never rewritten. + """ + if not failures: + return + print(f"\n{len(failures)} R statistics step(s) failed; their tables were not (re)written:") + for failure in failures: + print(f" - {failure}") + sys.exit(1) + + def _run_single( config: StudyConfig, analysis_name: str, @@ -218,27 +233,38 @@ def cmd_run(args): # Scope to one paradigm + one analysis aconfig = config.for_paradigm_analysis(args.paradigm, args.analysis) _prepare_output(aconfig, args.analysis, strict=strict, force=force, steps=steps) - _run_single(aconfig, args.analysis, steps=steps, select=select, jobs=jobs) + try: + _run_single(aconfig, args.analysis, steps=steps, select=select, jobs=jobs) + except RStepFailed as e: + _exit_if_failed([str(e)]) else: # Run all analyses listed for this paradigm analyses = config.get_paradigm_analyses(args.paradigm) if not analyses: print(f"No analyses listed for paradigm '{args.paradigm}'") sys.exit(1) + # One module failing does not stop the rest; the run still exits 1. + failures: list[str] = [] for analysis_name in analyses: print(f"{'='*60}") print(f"Paradigm: {args.paradigm} | Analysis: {analysis_name}") print(f"{'='*60}") aconfig = config.for_paradigm_analysis(args.paradigm, analysis_name) _prepare_output(aconfig, analysis_name, strict=strict, force=force, steps=steps) - _run_single(aconfig, analysis_name, steps=steps, select=select, jobs=jobs) + try: + _run_single(aconfig, analysis_name, steps=steps, select=select, jobs=jobs) + except RStepFailed as e: + failures.append(str(e)) + print(f"FAILED: {e}") print() + _exit_if_failed(failures) else: if args.analysis: print("ERROR: --analysis without --paradigm is ambiguous in multi-paradigm mode.") print("Specify --paradigm or omit --analysis to run everything.") sys.exit(1) # Run all paradigms, all their analyses + failures: list[str] = [] for pname in config.paradigms: analyses = config.get_paradigm_analyses(pname) or [] if not analyses: @@ -250,8 +276,13 @@ def cmd_run(args): print(f"{'='*60}") aconfig = config.for_paradigm_analysis(pname, analysis_name) _prepare_output(aconfig, analysis_name, strict=strict, force=force, steps=steps) - _run_single(aconfig, analysis_name, steps=steps, select=select, jobs=jobs) + try: + _run_single(aconfig, analysis_name, steps=steps, select=select, jobs=jobs) + except RStepFailed as e: + failures.append(f"{pname}/{e}") + print(f"FAILED: {e}") print() + _exit_if_failed(failures) else: # Legacy single-paradigm config if not args.analysis: @@ -260,7 +291,10 @@ def cmd_run(args): _prepare_output(config, args.analysis, strict=strict, force=force, steps=steps) analyzer = StudyAnalyzer(config) _print_study_summary(config, analyzer) - analyzer.run_analysis(args.analysis, steps=steps, select=select, jobs=jobs) + try: + analyzer.run_analysis(args.analysis, steps=steps, select=select, jobs=jobs) + except RStepFailed as e: + _exit_if_failed([str(e)]) print(f"\nDone. Output: {config.output_dir / canonical_analysis_name(args.analysis)}") diff --git a/tests/test_r_step_failures.py b/tests/test_r_step_failures.py new file mode 100644 index 0000000..774801d --- /dev/null +++ b/tests/test_r_step_failures.py @@ -0,0 +1,158 @@ +"""A failed or timed-out R statistics step fails the run instead of exiting 0. + +Found on the FORGE treatment re-run. roi_directed (in all three source arms) and +roi_cross_freq's AAC/PPC tier hit a hard-coded one-hour R timeout. The module +logged it and the command exited 0, so region tables from the previous code +version went on looking current. There is now no default limit (``r_timeout_sec`` +sets one), and a failure raises ``RStepFailed``. + +Also locks the roi_cross_freq PAC mosaics, whose call named columns the native +hypothesis table does not have, so none was ever drawn. +""" + +from __future__ import annotations + +import subprocess +from types import SimpleNamespace + +import pytest + +import source_analytics.viz.brain_roi as brain_roi +from source_analytics.analyses.base import RStepFailed +from source_analytics.analyses.roi_cross_freq_analysis import ROICrossFreqAnalysis +from source_analytics.analyses.roi_directed_analysis import ROIDirectedAnalysis +from source_analytics.analyses.roi_psd_analysis import ROIPsdAnalysis +from source_analytics.config import StudyConfig + + +def _module(cls, sample_config_yaml, tmp_path, block=None): + config = StudyConfig.from_yaml(sample_config_yaml) + if block is not None: + config.raw[cls.name] = block + module = cls(config, tmp_path / "analytics" / cls.name) + if not hasattr(module, "_sfreq"): + module._sfreq = None + return module + + +def test_there_is_no_default_r_timeout_and_one_can_be_set(sample_config_yaml, tmp_path): + assert _module(ROIPsdAnalysis, sample_config_yaml, tmp_path)._r_timeout is None + limited = _module(ROIPsdAnalysis, sample_config_yaml, tmp_path, {"r_timeout_sec": 7200}) + assert limited._r_timeout == 7200.0 + + +def test_a_timed_out_r_step_raises_and_uses_no_limit_by_default(sample_config_yaml, tmp_path, monkeypatch): + module = _module(ROIDirectedAnalysis, sample_config_yaml, tmp_path) + (module.output_dir / "data" / "roi_transfer_entropy_edges.csv").write_text("x\n") + seen = {} + + def fake_run(cmd, **kw): + seen.update(kw) + raise subprocess.TimeoutExpired(cmd, 1) + + monkeypatch.setattr(subprocess, "run", fake_run) + with pytest.raises(RStepFailed, match="roi_directed: R script timed out"): + module.summary() + assert seen["timeout"] is None + + +def test_a_failing_r_step_raises(sample_config_yaml, tmp_path, monkeypatch): + module = _module(ROIPsdAnalysis, sample_config_yaml, tmp_path) + for name in ("band_power.csv", "psd_curves.csv"): + (module.output_dir / "data" / name).write_text("x\n") + monkeypatch.setattr(subprocess, "run", + lambda cmd, **kw: SimpleNamespace(returncode=2, stdout="", stderr="boom")) + with pytest.raises(RStepFailed, match="roi_psd: R script failed with exit code 2"): + module.summary() + + +def test_cross_freq_attempts_both_tiers_before_failing(sample_config_yaml, tmp_path): + module = _module(ROICrossFreqAnalysis, sample_config_yaml, tmp_path) + module._metrics = ["pac", "aac"] + edges = [] + module._run_pac_r = lambda: False + module._run_edges_r = lambda metrics: edges.append(metrics) or True + module._edge_metrics_on_disk = lambda: ["aac"] + with pytest.raises(RStepFailed, match="PAC"): + module.summary() + assert edges == [["aac"]], "the AAC/PPC tier must still run when PAC fails" + + +def _pac_table(module): + cols = ["hypothesis", "band", "spatial", "effect_size", "p_value", "q_value", "significant", "dv"] + path = module.tbl_dir / "roi_pac_posthoc_region.csv" + path.write_text(",".join(cols) + "\ndisease_effect,Theta-Low Gamma,Motor,0.9,0.001,0.01,TRUE,z_score\n") + return set(cols) + + +def test_pac_mosaics_name_columns_the_table_has(sample_config_yaml, tmp_path, monkeypatch): + module = _module(ROICrossFreqAnalysis, sample_config_yaml, tmp_path) + columns = _pac_table(module) + seen = {} + monkeypatch.setattr(brain_roi, "render_posthoc_mosaics", lambda *a, **kw: seen.update(kw) or []) + module._render_brain_mosaics() + named = {seen["effect_col"], seen["roi_col"], seen["p_col"], seen["q_col"], *seen["facet_cols"]} + assert named <= columns, f"mosaic asks for columns the table lacks: {sorted(named - columns)}" + + +def test_the_figures_pass_draws_pac_mosaics(sample_config_yaml, tmp_path, monkeypatch): + module = _module(ROICrossFreqAnalysis, sample_config_yaml, tmp_path) + _pac_table(module) + module._metrics = ["pac"] + module._call_r_figures_only = lambda *a, **kw: None + module._edge_metrics_on_disk = lambda: [] + drawn = [] + monkeypatch.setattr(brain_roi, "render_posthoc_mosaics", lambda *a, **kw: drawn.append(kw) or []) + module.figures() + assert drawn, "a figures-only pass must draw the PAC mosaics" + + +# ---- the CLI: a failure fails the run, but a batch finishes its other modules -------- + +def _two_module_study(tmp_path): + import yaml + + (tmp_path / "data").mkdir() + path = tmp_path / "study.yaml" + path.write_text(yaml.safe_dump({ + "name": "t", "groups": {"WT_VEH": "WT", "KO_VEH": "KO"}, "bands": {"Theta": [4, 8]}, + "paths": {"analytics": str(tmp_path / "a"), "results": str(tmp_path / "r")}, + "paradigms": {"p1": {"data_dir": str(tmp_path / "data"), + "analyses": {"roi_psd": {}, "roi_aperiodic": {}}}}, + })) + return path + + +def _run_cli(monkeypatch, argv, fail_on): + import sys + + from source_analytics import cli + + ran = [] + + def fake_run_single(aconfig, name, **kw): + ran.append(name) + if name == fail_on: + raise RStepFailed(f"{name}: R script failed with exit code 2") + + monkeypatch.setattr(cli, "_run_single", fake_run_single) + monkeypatch.setattr(cli, "_prepare_output", lambda *a, **kw: None) + monkeypatch.setattr(sys, "argv", ["source-analytics", *argv]) + with pytest.raises(SystemExit) as exit_info: + cli.main() + return ran, exit_info.value.code + + +def test_a_paradigm_run_finishes_the_other_modules_then_exits_1(tmp_path, monkeypatch, capsys): + study = _two_module_study(tmp_path) + ran, code = _run_cli(monkeypatch, ["run", "--study", str(study), "--paradigm", "p1"], fail_on="roi_psd") + # Order follows the config; what matters is that the second module still ran. + assert sorted(ran) == ["roi_aperiodic", "roi_psd"] and code == 1 + assert "1 R statistics step(s) failed" in capsys.readouterr().out + + +def test_a_single_module_run_that_fails_exits_1(tmp_path, monkeypatch): + study = _two_module_study(tmp_path) + ran, code = _run_cli(monkeypatch, ["run", "--study", str(study), "--paradigm", "p1", + "--analysis", "roi_psd"], fail_on="roi_psd") + assert ran == ["roi_psd"] and code == 1