From 840865f4cce539697ddccaa2e5428137667ebdcd Mon Sep 17 00:00:00 2001 From: Richard O'Shaughnessy Date: Sat, 15 Aug 2026 18:00:40 -0700 Subject: [PATCH 1/3] LISA ILE driver: port the portfolio freeze/allocation policy and NF flow persistence Pass 4 of the catch-up (pass 3, MC-error replicas, deferred). Closes 14 of the 108 remaining gap items (gap 108 -> 94). Pure pass-through to samplers this driver ALREADY wires: it exposes the same ok_lnL_methods as the main driver and builds mcsamplerPortfolio the same way, and its portfolio setup block was byte-identical to main's apart from the missing kwargs. Before this the knobs were reachable only through --sampler-portfolio-args, an eval-able dict; the pipeline passes the named flags. Option definitions are copied verbatim, and a test asserts default/type/action/choices match the main driver's for all 14. A knob that means something different in the two drivers is worse than a missing one: the same pipeline command line would otherwise produce two different integrations. The freeze-policy assembly is inline in both drivers rather than a function, so the tests extract the block and exec it against a fake opts -- testing the real source rather than a paraphrase. That covers the property the assembly exists for: an option left UNSET must stay out of the dict so the sampler keeps its own default, and 0 (which disables probing/reviving) is a REAL value that truthiness would silently drop. NF flow load/save go through _maybe_load_nf_flow / _maybe_save_nf_flow for the same reason the L0 rescue did: this driver has TWO analyze_event variants, and wiring only one would be a silent half-port. A test asserts both get both hooks and that they straddle the integration in the right order. TESTS. test_lisa_sampler_plumbing.py (55), wired into lisa-check. Revert-checked with 7 mutations: a drifted default, a dropped `is not None` guard, inverted VARAHA precedence, a dict that never reaches setup(), a lost hasattr guard, a misplaced save hook, and an unguarded NF exception. The hasattr mutation came back WEAK and the test was rewritten. Asserting "does not raise" was worthless there: the body is wrapped in `except Exception`, so dropping the guard still does not raise -- it announces "loading pre-trained flow", calls a method that does not exist, and swallows the AttributeError, so every non-NF run would log a flow load that never happened. The test now asserts the hook stays SILENT for a sampler with no flow support. Co-Authored-By: Claude Opus 5 --- .travis/test-lisa.sh | 1 + ...egrate_likelihood_extrinsic_batchmode_lisa | 82 ++++- .../integrators/lisa_drift_ledger.json | 56 --- .../integrators/make_lisa_drift_ledger.py | 19 +- .../Code/test/test_lisa_sampler_plumbing.py | 331 ++++++++++++++++++ 5 files changed, 423 insertions(+), 66 deletions(-) create mode 100644 MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py diff --git a/.travis/test-lisa.sh b/.travis/test-lisa.sh index c1db83d6c..45c50178e 100644 --- a/.travis/test-lisa.sh +++ b/.travis/test-lisa.sh @@ -18,4 +18,5 @@ fi MonteCarloMarginalizeCode/Code/test/test_lisa_synthetic_demo.py \ MonteCarloMarginalizeCode/Code/test/test_lisa_fairdraw_weights.py \ MonteCarloMarginalizeCode/Code/test/test_lisa_l0_rescue.py \ + MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py \ MonteCarloMarginalizeCode/Code/test/test_lisa_driver_drift.py diff --git a/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa b/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa index 37d8b465d..0a572018c 100755 --- a/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa +++ b/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa @@ -306,6 +306,24 @@ integration_params.add_option("--internal-use-lnL",action='store_true',help="lik integration_params.add_option("--sampler-method",default="adaptive_cartesian_gpu",help="adaptive_cartesian|GMM|adaptive_cartesian_gpu") integration_params.add_option("--sampler-portfolio",default=None,action='append',type=str,help="comma-separated strings, matching sampler methods other than portfolio") integration_params.add_option("--sampler-portfolio-args",default=None, action='append', type=str, help='eval-able dictionaryo to be passed to that sampler') +# Portfolio freeze/allocation policy and NF flow persistence. Pure pass-through to the +# shared samplers, which this driver already wires (identical ok_lnL_methods). Definitions +# copied verbatim from bin/integrate_likelihood_extrinsic_batchmode; pinned by +# test_lisa_sampler_plumbing.py. +integration_params.add_option("--portfolio-adaptive-alloc",action='store_true',default=False,help="Portfolio: ENABLE (opt-in) adaptive-probe draw allocation -- concentrate draws on the best per-chunk-n_ess member. Good on strongly-correlated targets; NOT recommended for AV-favorable high-SNR events (it starves the slow-contracting AV workhorse). Off by default (legacy n_ess reweighting).") +integration_params.add_option("--portfolio-alloc-exponent",default=None,type=float,help="Portfolio: adaptive allocation ~ member_quality^exponent. Higher concentrates harder on the winner. Sampler default 1.0.") +integration_params.add_option("--portfolio-freeze-wt",default=None,type=float,help="Portfolio: a member whose balance weight is below this stops updating its proposal (subject to grace/revive/VARAHA-exemption). Sampler default 0.05.") +integration_params.add_option("--portfolio-grace-iters",default=None,type=int,help="Portfolio: never freeze ANY member during the first N integration chunks (let slow starters contract). Sampler default 25.") +integration_params.add_option("--portfolio-probe-period",default=None,type=int,help="Portfolio: round-robin probe one member at a raised draw share every N chunks (breaks the under-observation trap). 0 disables probing. Sampler default 4.") +integration_params.add_option("--portfolio-quality-signal",default=None,type=str,help="Portfolio adaptive allocation: which per-member quality signal to rank members by. 'global' (default) = marginal gain in POOLED n_eff per sample (credits weight mass, debits weight variance); 'credit' = q_mix-native MIS credit assignment, sum_i [frac_m q_m/q_mix]_i * w_i per drawn sample (credits a member for COVERING where the integrand is, even if it drew few samples there); 'ness' = legacy per-member Kish n_ess (scale-invariant, misranks a slow-contracting AV -- see DESIGN_portfolio_freeze_policy.md).") +integration_params.add_option("--portfolio-revive-period",default=None,type=int,help="Portfolio: every N chunks, update even a frozen member one step so it can recover. 0 disables. Sampler default 8.") +integration_params.add_option("--portfolio-varaha-can-freeze",action='store_true',default=False,help="Portfolio: DISABLE the VARAHA freeze-exemption, so VARAHA/AV members obey the grace/revive/weight freeze schedule like other members. Use only if a VARAHA member is a known-bad fit and you want to save its selfish-draw eval cycles.") +integration_params.add_option("--portfolio-varaha-max-frac",default=None,type=float,help="Portfolio: CAP the combined DRAW fraction of VARAHA/AV members (0/unset = no cap). Use WITH --portfolio-varaha-min-frac to constrain the VARAHA share to a BAND. Rationale: a floor alone stops the mixture degenerating to peaked-member-only (which strips q_mix of its broad backstop, so a missed mode goes uncovered and lnZ is silently low while n_eff looks GOOD), but the share can then run away the OTHER way to ~1 and the mixture degenerates to VARAHA-only instead. A band (e.g. 0.25/0.75) keeps q_mix genuinely mixed by construction. Unbiased either way (balance heuristic), so it costs at most draws, never correctness.") +integration_params.add_option("--portfolio-varaha-min-frac",default=None,type=float,help="Portfolio: reserve this combined DRAW fraction for VARAHA/AV members (0/unset = off). never-freeze keeps a VARAHA member UPDATING, but both allocation rules score by per-chunk n_ess, which sits at ~1 during VARAHA's slow cumulative contraction -- so a member that looks instantly good can take nearly the whole budget (measured on S250114ax post-#33: GMM took ~0.84 and the portfolio collapsed to n_eff ~2 vs ~100 for standalone AV). Unbiased for any allocation (q_mix); trades efficiency only.") +integration_params.add_option("--portfolio-varaha-never-freeze",action='store_true',default=False,help="Portfolio: VARAHA/AV members always update every chunk past their breakpoint (freeze-exempt). This is the sampler default; the flag is here for explicitness/pipe pass-through.") +integration_params.add_option("--portfolio-weight-clip",default=None,type=float,help="Portfolio: OPT-IN truncated importance sampling applied to the PROPOSAL-FIT INPUT ONLY. Caps the weights fed to member.update_sampling_prior (the GMM covariance fit) at tau = C*sqrt(n)*mean(w) (0/unset = off; C~1 is the standard Ionides choice), so one enormous weight cannot make that fit degenerate. The estimator (ln Z, n_eff), the n_ess report, and the allocation signal all use the TRUE unclipped weights, so they stay exactly unbiased and undistorted. Do NOT clip the estimator (measured on S250114ax: n_eff=100 2x faster than AV but ln Z biased -11.5 nats) or the n_ess report (clipping inflates the clipped member's n_ess and starves the AV workhorse). The withheld tail mass is tracked and reported as a diagnostic. NOTE: if huge weights come from q_mix UNDERFLOW (watch for the warning) they are a numerical artifact, not tail mass.") +integration_params.add_option("--nf-flow-load",default=None,help="NF only: load a pre-trained normalizing flow (.pt from --nf-flow-save); with --n-adapt 0 this reuses it directly (skips training), otherwise it is polished.") +integration_params.add_option("--nf-flow-save",default=None,help="NF only: after integration, serialize the trained normalizing flow (.pt) for reuse across ILE instances.") integration_params.add_option("--sampler-xpy",default=None,help="numpy|cupy if the adaptive_cartesian_gpu sampler is active, use that.") # L0 auto-rescue. Ported from bin/integrate_likelihood_extrinsic_batchmode; defaults and help # text kept IDENTICAL there and here on purpose -- see test_lisa_l0_rescue.py, which pins them. @@ -1447,7 +1465,37 @@ if use_portfolio: if not(isinstance(opts.sampler_portfolio_args[indx], dict)): print(indx,opts.sampler_portfolio_args[indx]) print(" ARGS ", opts.sampler_portfolio_args) - sampler.setup(portfolio_args=opts.sampler_portfolio_args, **pinned_params) # directly pass all parameters set above to low-level portfolios. In particular, GMM setup + # Assemble freeze-policy overrides from the CLI. Only include options the user actually + # set (None = unset) so the sampler keeps its built-in defaults otherwise. The two VARAHA + # flags are mutually exclusive; --portfolio-varaha-can-freeze wins if both are given. + _freeze_policy_kwargs = {} + if opts.portfolio_grace_iters is not None: + _freeze_policy_kwargs['portfolio_grace_iters'] = opts.portfolio_grace_iters + if opts.portfolio_revive_period is not None: + _freeze_policy_kwargs['portfolio_revive_period'] = opts.portfolio_revive_period + if opts.portfolio_freeze_wt is not None: + _freeze_policy_kwargs['portfolio_freeze_wt'] = opts.portfolio_freeze_wt + if opts.portfolio_varaha_can_freeze: + _freeze_policy_kwargs['portfolio_varaha_never_freeze'] = False + elif opts.portfolio_varaha_never_freeze: + _freeze_policy_kwargs['portfolio_varaha_never_freeze'] = True + # adaptive-probe draw allocation (OPT-IN; off by default in the sampler) + if opts.portfolio_adaptive_alloc: + _freeze_policy_kwargs['portfolio_adaptive_alloc'] = True + if opts.portfolio_varaha_min_frac is not None: + _freeze_policy_kwargs['portfolio_varaha_min_frac'] = opts.portfolio_varaha_min_frac + if opts.portfolio_varaha_max_frac is not None: + _freeze_policy_kwargs['portfolio_varaha_max_frac'] = opts.portfolio_varaha_max_frac + if opts.portfolio_weight_clip is not None: + _freeze_policy_kwargs['portfolio_weight_clip'] = opts.portfolio_weight_clip + if opts.portfolio_quality_signal is not None: + _freeze_policy_kwargs['portfolio_quality_signal'] = opts.portfolio_quality_signal + if opts.portfolio_alloc_exponent is not None: + _freeze_policy_kwargs['portfolio_alloc_exponent'] = opts.portfolio_alloc_exponent + if opts.portfolio_probe_period is not None: + _freeze_policy_kwargs['portfolio_probe_period'] = opts.portfolio_probe_period + print(" PORTFOLIO freeze-policy overrides: ", _freeze_policy_kwargs) + sampler.setup(portfolio_args=opts.sampler_portfolio_args, **_freeze_policy_kwargs, **pinned_params) # directly pass all parameters set above to low-level portfolios. In particular, GMM setup # initialize sampler, before we call integrate, so we can seed it if opts.sampler_method == 'adaptive_cartesian_gpu' and opts.skymap_file: @@ -1661,6 +1709,30 @@ def _clear_warm_state(sampler): sampler._warm_applied = False +def _maybe_load_nf_flow(sampler): + """Warm-load a pre-trained normalizing flow so this instance skips/shortens training. + + hasattr-guarded, so it is a no-op for every sampler that is not NF -- including all five + this driver lists in ok_lnL_methods, where NF is reachable only as a portfolio member. + """ + if opts.nf_flow_load and hasattr(sampler, 'load_flow'): + try: + print(" NF: loading pre-trained flow from", opts.nf_flow_load) + sampler.load_flow(opts.nf_flow_load) + except Exception as _e_nf: + print(" NF flow load skipped (", _e_nf, ")") + + +def _maybe_save_nf_flow(sampler): + """Serialize the trained flow for reuse across ILE instances. hasattr-guarded as above.""" + if opts.nf_flow_save and hasattr(sampler, 'save_flow'): + try: + sampler.save_flow(opts.nf_flow_save) + print(" NF: saved trained flow to", opts.nf_flow_save) + except Exception as _e_fs: + print(" NF: could not save flow (", _e_fs, ")") + + def _maybe_l0_rescue(sampler, res, var, neff, dict_return, like_to_integrate, unpinned_params, pinned_params, lnL_offset=0.0): @@ -2040,6 +2112,8 @@ def analyze_event_LISA(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_ lnL_oracles = np.zeros(opts.n_chunk) sampler.update_sampling_prior(lnL_oracles, opts.n_chunk, external_rvs=rvs_train,log_scale_weights=True,floor_integrated_probability=opts.adapt_floor_level) + _maybe_load_nf_flow(sampler) + res, var, neff, dict_return = sampler.integrate(like_to_integrate, *unpinned_params, **pinned_params) # L0 auto-rescue: on a very sharply-peaked (high-amplitude) point a cold AV can stall at @@ -2053,6 +2127,8 @@ def analyze_event_LISA(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_ like_to_integrate, unpinned_params, pinned_params, lnL_offset=manual_avoid_overflow_logarithm) + _maybe_save_nf_flow(sampler) + if not(res): # no resut raise ValueError(" No integral result returned") @@ -2756,6 +2832,8 @@ def analyze_event(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_spec_ lnL_oracles = np.zeros(opts.n_chunk) sampler.update_sampling_prior(lnL_oracles, opts.n_chunk, external_rvs=rvs_train,log_scale_weights=True,floor_integrated_probability=opts.adapt_floor_level) + _maybe_load_nf_flow(sampler) + res, var, neff, dict_return = sampler.integrate(like_to_integrate, *unpinned_params, **pinned_params) # L0 auto-rescue: on a very sharply-peaked (high-amplitude) point a cold AV can stall at @@ -2769,6 +2847,8 @@ def analyze_event(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_spec_ like_to_integrate, unpinned_params, pinned_params, lnL_offset=manual_avoid_overflow_logarithm) + _maybe_save_nf_flow(sampler) + if not(res): # no resut raise ValueError(" No integral result returned") diff --git a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json index ad2ed8672..dab8eac52 100644 --- a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json +++ b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json @@ -309,62 +309,6 @@ "decision": "NA", "reason": "The .dslice export and its placement/tuning knobs. This is a data product for a downstream LIGO CIP distance workflow that the LISA pipeline does not run; there is no consumer. If a LISA distance workflow is ever built, note that the .dslice reweight core was the third Finding-2 site and must not be revived in its pre-#87 form." }, - "OPTION:--nf-flow-load": { - "decision": "PORT", - "reason": "Normalizing-flow persistence. Neither driver lists an NF method in ok_lnL_methods (identical lists), so NF is reached only as a portfolio member -- equally available to LISA. Low priority, but not LISA-specific in any way." - }, - "OPTION:--nf-flow-save": { - "decision": "PORT", - "reason": "Normalizing-flow persistence. Neither driver lists an NF method in ok_lnL_methods (identical lists), so NF is reached only as a portfolio member -- equally available to LISA. Low priority, but not LISA-specific in any way." - }, - "OPTION:--portfolio-adaptive-alloc": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-alloc-exponent": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-freeze-wt": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-grace-iters": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-probe-period": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-quality-signal": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-revive-period": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-varaha-can-freeze": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-varaha-max-frac": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-varaha-min-frac": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-varaha-never-freeze": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, - "OPTION:--portfolio-weight-clip": { - "decision": "PORT", - "reason": "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts --sampler-portfolio-args (an eval-able dict), so these are reachable today via that escape hatch; porting them as first-class flags is pipeline parity, which is what the pipe actually passes. Low risk, no physics." - }, "OPTION:--random-event": { "decision": "PORT", "reason": "Pick a random event from the input file. Detector-agnostic; flagged dangerous in its own help text for oversampling reasons that apply equally to LISA." diff --git a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py index a942beaf2..4ed4c8cd6 100644 --- a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py +++ b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py @@ -160,17 +160,18 @@ "is the ecliptic pair -- the grouping still makes sense, the docstring does not."), # ------------------------------------------------------------------ portfolio plumbing - (r"^OPTION:--portfolio-", "PORT", - "mcsamplerPortfolio tuning. LISA wires the portfolio sampler and already accepts " - "--sampler-portfolio-args (an eval-able dict), so these are reachable today via " - "that escape hatch; porting them as first-class flags is pipeline parity, which is " - "what the pipe actually passes. Low risk, no physics."), + (r"^OPTION:--portfolio-", "PORTED", + "mcsamplerPortfolio freeze/allocation policy. Definitions copied verbatim and the " + "_freeze_policy_kwargs assembly is textually identical to the main driver's, so " + "unset options (None) stay out of the dict and the sampler keeps its own defaults. " + "--portfolio-varaha-can-freeze wins over --portfolio-varaha-never-freeze, as there."), # ------------------------------------------------------------------- NF flow plumbing - (r"^OPTION:--nf-flow-(load|save)$", "PORT", - "Normalizing-flow persistence. Neither driver lists an NF method in " - "ok_lnL_methods (identical lists), so NF is reached only as a portfolio member -- " - "equally available to LISA. Low priority, but not LISA-specific in any way."), + (r"^OPTION:--nf-flow-(load|save)$", "PORTED", + "Normalizing-flow persistence. Neither driver lists an NF method in ok_lnL_methods " + "(identical lists), so NF is reached only as a portfolio member -- equally " + "available to LISA. Both hooks are hasattr-guarded, so they are a no-op for every " + "other sampler."), # --------------------------------------------------------- extrinsic proposal handoff (r"^OPTION:--extrinsic-proposal-output$", "PORT", diff --git a/MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py b/MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py new file mode 100644 index 000000000..5ba3e0ca4 --- /dev/null +++ b/MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py @@ -0,0 +1,331 @@ +#!/usr/bin/env python +""" +Tests for the portfolio freeze/allocation policy and NF flow persistence ported into the +LISA ILE driver (bin/integrate_likelihood_extrinsic_batchmode_lisa). + +This is pure PASS-THROUGH plumbing to samplers the LISA driver already wires -- it exposes +the same ``ok_lnL_methods`` as the main driver (``GMM, adaptive_cartesian, +adaptive_cartesian_gpu, AV, portfolio``, verified identical) and builds mcsamplerPortfolio +the same way. Before this port the knobs were reachable only through +``--sampler-portfolio-args``, an eval-able dict; the pipeline passes the named flags. + +WHAT CAN ACTUALLY GO WRONG HERE, and is therefore what these tests check: + + * a default that differs between the two drivers. Worse than a missing option: the same + command line then means two different things depending on which driver ran it. + * an option that is UNSET leaking into the kwargs as ``None`` and overriding the sampler's + own default with nothing. The assembly's whole shape -- ``if opts.x is not None`` -- + exists for that, and a single dropped guard is invisible until a run behaves oddly. + * the two mutually-exclusive VARAHA flags resolving the wrong way round. + * an NF hook that is not hasattr-guarded, which would break every non-NF sampler. + +The freeze-policy assembly is inline in both drivers (not a function), so it is exercised +here by extracting the block and exec'ing it against a fake ``opts``. That tests the real +source, not a paraphrase of it. +""" + +import ast +import os +import re +import textwrap + +import pytest + +_HERE = os.path.dirname(os.path.abspath(__file__)) +_LISA = os.path.join(_HERE, '..', 'bin', 'integrate_likelihood_extrinsic_batchmode_lisa') +_MAIN = os.path.join(_HERE, '..', 'bin', 'integrate_likelihood_extrinsic_batchmode') + +PORTFOLIO_OPTS = [ + "--portfolio-adaptive-alloc", "--portfolio-alloc-exponent", "--portfolio-freeze-wt", + "--portfolio-grace-iters", "--portfolio-probe-period", "--portfolio-quality-signal", + "--portfolio-revive-period", "--portfolio-varaha-can-freeze", + "--portfolio-varaha-max-frac", "--portfolio-varaha-min-frac", + "--portfolio-varaha-never-freeze", "--portfolio-weight-clip", +] +NF_OPTS = ["--nf-flow-load", "--nf-flow-save"] + + +def _src(path): + with open(path) as fh: + return fh.read() + + +def _option_nodes(path): + """{'--foo': ast.Call} for every add_option in a driver.""" + out = {} + for n in ast.walk(ast.parse(_src(path), filename=path)): + if (isinstance(n, ast.Call) and isinstance(n.func, ast.Attribute) + and n.func.attr in ("add_option", "add_argument")): + names = [a.value for a in n.args + if isinstance(a, ast.Constant) and isinstance(a.value, str)] + if names and names[0].startswith("--"): + out[names[0]] = n + return out + + +def _kwargs_of(node): + out = {} + for kw in node.keywords: + try: + out[kw.arg] = ast.literal_eval(kw.value) + except Exception: + out[kw.arg] = ast.dump(kw.value) + return out + + +@pytest.fixture(scope="module") +def opts_lisa(): + return _option_nodes(_LISA) + + +@pytest.fixture(scope="module") +def opts_main(): + return _option_nodes(_MAIN) + + +# ------------------------------------------------------------------------------ presence +@pytest.mark.parametrize("opt", PORTFOLIO_OPTS + NF_OPTS) +def test_option_is_present_in_the_lisa_driver(opt, opts_lisa): + assert opt in opts_lisa + + +# ------------------------------------------------------------------------------- defaults +@pytest.mark.parametrize("opt", PORTFOLIO_OPTS + NF_OPTS) +def test_option_signature_matches_the_main_driver(opt, opts_lisa, opts_main): + """Same default, same type, same action. + + A knob that means something different in the two drivers is worse than a missing one: + the same pipeline command line would then produce two different integrations. + """ + a, b = _kwargs_of(opts_lisa[opt]), _kwargs_of(opts_main[opt]) + for key in ("default", "type", "action", "choices"): + assert a.get(key) == b.get(key), ( + "%s: %s differs (lisa=%r, main=%r)" % (opt, key, a.get(key), b.get(key))) + + +@pytest.mark.parametrize("opt", [ + "--portfolio-alloc-exponent", "--portfolio-freeze-wt", "--portfolio-grace-iters", + "--portfolio-probe-period", "--portfolio-quality-signal", "--portfolio-revive-period", + "--portfolio-varaha-max-frac", "--portfolio-varaha-min-frac", "--portfolio-weight-clip", +]) +def test_tuning_options_default_to_none_so_the_sampler_keeps_its_own(opt, opts_lisa): + """None is the sentinel the assembly keys on. A default of 0/0.0 would silently + override the sampler's built-in value for every run that never set the flag.""" + assert _kwargs_of(opts_lisa[opt]).get("default") is None + + +@pytest.mark.parametrize("opt", [ + "--portfolio-adaptive-alloc", "--portfolio-varaha-can-freeze", + "--portfolio-varaha-never-freeze", +]) +def test_flags_are_store_true_and_default_false(opt, opts_lisa): + kw = _kwargs_of(opts_lisa[opt]) + assert kw.get("action") == "store_true" and kw.get("default") is False + + +# --------------------------------------------------------- the assembly block, executed +_START = "_freeze_policy_kwargs = {}" +_END = 'print(" PORTFOLIO freeze-policy overrides: "' + + +def _assembly_block(path): + """The inline freeze-policy assembly, dedented so it can be exec'd on its own. + + Slice from the START OF THE LINE holding the sentinel, not from the sentinel itself: + otherwise the first line carries no indentation while the rest do, and dedent finds no + common prefix. + """ + src = _src(path) + i = src.rindex("\n", 0, src.index(_START)) + 1 + j = src.index(_END, i) + j = src.rindex("\n", i, j) + 1 + return textwrap.dedent(src[i:j]) + + +class _Opts(object): + """Every portfolio option at its documented default.""" + portfolio_grace_iters = None + portfolio_revive_period = None + portfolio_freeze_wt = None + portfolio_varaha_can_freeze = False + portfolio_varaha_never_freeze = False + portfolio_adaptive_alloc = False + portfolio_varaha_min_frac = None + portfolio_varaha_max_frac = None + portfolio_weight_clip = None + portfolio_quality_signal = None + portfolio_alloc_exponent = None + portfolio_probe_period = None + + def __init__(self, **kw): + for k, v in kw.items(): + assert hasattr(type(self), k), "unknown option %s" % k + setattr(self, k, v) + + +def _assemble(**kw): + ns = {"opts": _Opts(**kw)} + exec(compile(_assembly_block(_LISA), "freeze_policy", "exec"), ns) + return ns["_freeze_policy_kwargs"] + + +def test_nothing_set_means_nothing_overridden(): + """The important one: an all-defaults run must not touch the sampler's policy at all.""" + assert _assemble() == {} + + +def test_each_tuning_option_passes_through_when_set(): + got = _assemble(portfolio_grace_iters=7, portfolio_revive_period=3, + portfolio_freeze_wt=0.25, portfolio_varaha_min_frac=0.2, + portfolio_varaha_max_frac=0.8, portfolio_weight_clip=1.0, + portfolio_quality_signal='credit', portfolio_alloc_exponent=2.0, + portfolio_probe_period=5) + assert got == {'portfolio_grace_iters': 7, 'portfolio_revive_period': 3, + 'portfolio_freeze_wt': 0.25, 'portfolio_varaha_min_frac': 0.2, + 'portfolio_varaha_max_frac': 0.8, 'portfolio_weight_clip': 1.0, + 'portfolio_quality_signal': 'credit', 'portfolio_alloc_exponent': 2.0, + 'portfolio_probe_period': 5} + + +def test_zero_is_passed_through_not_treated_as_unset(): + """0 disables probing/reviving and is a REAL value; `if x:` would drop it.""" + got = _assemble(portfolio_probe_period=0, portfolio_revive_period=0) + assert got == {'portfolio_probe_period': 0, 'portfolio_revive_period': 0} + + +def test_varaha_never_freeze_sets_true(): + assert _assemble(portfolio_varaha_never_freeze=True) == {'portfolio_varaha_never_freeze': True} + + +def test_varaha_can_freeze_sets_false(): + assert _assemble(portfolio_varaha_can_freeze=True) == {'portfolio_varaha_never_freeze': False} + + +def test_can_freeze_wins_when_both_are_given(): + """Documented precedence; the two flags are mutually exclusive.""" + got = _assemble(portfolio_varaha_can_freeze=True, portfolio_varaha_never_freeze=True) + assert got == {'portfolio_varaha_never_freeze': False} + + +def test_adaptive_alloc_is_opt_in_only(): + assert 'portfolio_adaptive_alloc' not in _assemble() + assert _assemble(portfolio_adaptive_alloc=True) == {'portfolio_adaptive_alloc': True} + + +def test_assembly_block_is_identical_to_the_main_drivers(): + """Deliberate copies in a deliberate fork. Change one, change both.""" + def norm(s): + return re.sub(r"\s+", " ", s).strip() + assert norm(_assembly_block(_LISA)) == norm(_assembly_block(_MAIN)) + + +def test_assembly_result_is_actually_handed_to_setup(): + """Building the dict and not passing it would be a silent no-op.""" + src = _src(_LISA) + assert "sampler.setup(portfolio_args=opts.sampler_portfolio_args, **_freeze_policy_kwargs" in src + + +# ------------------------------------------------------------------------------- NF hooks +def _fn(path, name): + for n in ast.parse(_src(path)).body: + if isinstance(n, ast.FunctionDef) and n.name == name: + return n + raise AssertionError("%s not found in %s" % (name, os.path.basename(path))) + + +def _load_nf(**optkw): + ns = {"opts": type("O", (), dict({"nf_flow_load": None, "nf_flow_save": None}, **optkw))()} + mod = ast.Module(body=[_fn(_LISA, '_maybe_load_nf_flow'), _fn(_LISA, '_maybe_save_nf_flow')], + type_ignores=[]) + exec(compile(ast.fix_missing_locations(mod), "nf", "exec"), ns) + return ns + + +class _NoFlow(object): + """A sampler with no flow support -- i.e. every sampler in ok_lnL_methods.""" + + +class _WithFlow(object): + def __init__(self): + self.loaded = self.saved = None + + def load_flow(self, path): + self.loaded = path + + def save_flow(self, path): + self.saved = path + + +def test_nf_hooks_are_noops_when_the_options_are_unset(): + ns = _load_nf() + s = _WithFlow() + ns['_maybe_load_nf_flow'](s) + ns['_maybe_save_nf_flow'](s) + assert s.loaded is None and s.saved is None + + +def test_nf_hooks_are_noops_for_a_sampler_without_flow_support(capsys): + """hasattr-guarded: must DECLINE for AV/GMM/portfolio/adaptive_cartesian. + + Asserting "does not raise" is not enough and an earlier version of this test made + exactly that mistake: the body is wrapped in `except Exception`, so dropping the + hasattr guard still does not raise -- it announces "loading pre-trained flow", calls a + method that does not exist, and swallows the AttributeError. Every non-NF run would + then log a flow load that never happened. So the observable property is that the hook + says NOTHING and touches nothing when the sampler has no flow support. + """ + ns = _load_nf(nf_flow_load="/x/flow.pt", nf_flow_save="/x/flow.pt") + capsys.readouterr() + ns['_maybe_load_nf_flow'](_NoFlow()) + ns['_maybe_save_nf_flow'](_NoFlow()) + out = capsys.readouterr().out + assert "NF" not in out, ( + "the hook engaged a sampler with no flow support (and the except swallowed it): %r" % out) + + +def test_nf_load_and_save_reach_a_flow_capable_sampler(): + ns = _load_nf(nf_flow_load="/in.pt", nf_flow_save="/out.pt") + s = _WithFlow() + ns['_maybe_load_nf_flow'](s) + ns['_maybe_save_nf_flow'](s) + assert s.loaded == "/in.pt" and s.saved == "/out.pt" + + +def test_nf_failures_do_not_abort_the_event(): + """A missing/corrupt flow file must degrade to a cold run, not kill the point.""" + ns = _load_nf(nf_flow_load="/in.pt", nf_flow_save="/out.pt") + + class _Boom(object): + def load_flow(self, p): + raise IOError("no such file") + + def save_flow(self, p): + raise IOError("read-only") + + ns['_maybe_load_nf_flow'](_Boom()) + ns['_maybe_save_nf_flow'](_Boom()) + + +# ------------------------------------------------------------------------ call-site wiring +def test_both_analyze_event_variants_get_the_nf_hooks(): + """This driver has two; wiring only one is a silent half-port.""" + tree = ast.parse(_src(_LISA)) + fns = {n.name: n for n in tree.body + if isinstance(n, ast.FunctionDef) and n.name in ('analyze_event', 'analyze_event_LISA')} + assert set(fns) == {'analyze_event', 'analyze_event_LISA'} + for name, node in fns.items(): + called = {c.func.id for c in ast.walk(node) + if isinstance(c, ast.Call) and isinstance(c.func, ast.Name)} + assert '_maybe_load_nf_flow' in called, "%s never loads the flow" % name + assert '_maybe_save_nf_flow' in called, "%s never saves the flow" % name + + +def test_flow_is_loaded_before_the_integration_and_saved_after(): + src = _src(_LISA) + pos = 0 + for _ in range(2): + load = src.index("_maybe_load_nf_flow(sampler)", pos) + integ = src.index("sampler.integrate(like_to_integrate", load) + save = src.index("_maybe_save_nf_flow(sampler)", integ) + assert load < integ < save, "flow load/save straddle the integration incorrectly" + pos = save + 1 From 8d58ce8b1b99ecc5a962c1dfef1daa37d7ca0a9d Mon Sep 17 00:00:00 2001 From: Richard O'Shaughnessy Date: Sat, 15 Aug 2026 18:07:37 -0700 Subject: [PATCH 2/3] LISA ILE driver: port AV state persistence, anisotropic bins and the collapse gate Pass 5a of the catch-up. Closes 5 of the 94 remaining gap items (gap 94 -> 89): --sampler-save-state, --sampler-load-state, --sampler-anisotropic-bins, --reject-collapsed-live-volume, and _reject_if_collapsed. All four are sampler-agnostic. The saved state is the AV sampler's own live-volume grid, which carries no detector convention; the bin allocation is per-axis on that same grid. THE ONE THING TO CARRY FORWARD. The main driver calls its collapse gate TWICE -- on the first run AND on the replica pool, because replication can turn a healthy first run into a collapsed POOL, and gating only the first would silently bypass the flag for exactly the case pooling introduces. This driver has no replica pooling yet, so only the first call exists here. That is recorded in the helper's docstring, in the drift ledger, and in a test that asserts the warning is still written where whoever ports --mc-error-replicas will be working. The helpers are hoisted to module level rather than nested (as _reject_if_collapsed is in main), because this driver has TWO analyze_event variants and nesting would mean two copies. A test pins the hoisted body AST-identical to main's nested one. THIS EXPOSED A FALSE POSITIVE IN THE DRIFT AUDIT, now fixed. It compared FUNC items by QUALIFIED name, so main's analyze_event._reject_if_collapsed did not match this driver's correctly-hoisted top-level _reject_if_collapsed, and the item would have sat in the gap forever no matter how well it was ported. A gate that cannot be satisfied is a gate people learn to ignore. FUNC items are now matched on the bare name as well. TESTS. test_lisa_av_state.py (27), wired into lisa-check. Revert-checked with 8 mutations: the lost AV-method restriction on save, bins not reaching portfolio members, bins ceasing to be opt-in, a gate that ignores its flag, a gate that fires on healthy runs, a collapse that is no longer announced, the gate moved before the not(res) guard, and deletion of the second-call-site warning. Each caught by its named test; file restored byte-identical. Co-Authored-By: Claude Opus 5 --- .travis/test-lisa.sh | 1 + ...egrate_likelihood_extrinsic_batchmode_lisa | 80 +++++ .../integrators/audit_lisa_driver_drift.py | 9 + .../integrators/lisa_drift_ledger.json | 20 -- .../integrators/make_lisa_drift_ledger.py | 17 +- .../Code/test/test_lisa_av_state.py | 318 ++++++++++++++++++ 6 files changed, 419 insertions(+), 26 deletions(-) create mode 100644 MonteCarloMarginalizeCode/Code/test/test_lisa_av_state.py diff --git a/.travis/test-lisa.sh b/.travis/test-lisa.sh index 45c50178e..3b4e51633 100644 --- a/.travis/test-lisa.sh +++ b/.travis/test-lisa.sh @@ -19,4 +19,5 @@ fi MonteCarloMarginalizeCode/Code/test/test_lisa_fairdraw_weights.py \ MonteCarloMarginalizeCode/Code/test/test_lisa_l0_rescue.py \ MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py \ + MonteCarloMarginalizeCode/Code/test/test_lisa_av_state.py \ MonteCarloMarginalizeCode/Code/test/test_lisa_driver_drift.py diff --git a/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa b/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa index 0a572018c..f75bb6f8b 100755 --- a/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa +++ b/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa @@ -325,6 +325,12 @@ integration_params.add_option("--portfolio-weight-clip",default=None,type=float, integration_params.add_option("--nf-flow-load",default=None,help="NF only: load a pre-trained normalizing flow (.pt from --nf-flow-save); with --n-adapt 0 this reuses it directly (skips training), otherwise it is polished.") integration_params.add_option("--nf-flow-save",default=None,help="NF only: after integration, serialize the trained normalizing flow (.pt) for reuse across ILE instances.") integration_params.add_option("--sampler-xpy",default=None,help="numpy|cupy if the adaptive_cartesian_gpu sampler is active, use that.") +# AV live-volume state, per-axis bin allocation, and the collapse gate. Copied verbatim +# from bin/integrate_likelihood_extrinsic_batchmode; pinned by test_lisa_av_state.py. +integration_params.add_option("--sampler-save-state",default=None,help="AV only: after integration, write the adapted live-volume state (.npz) for reuse by later instances/iterations. Point --sampler-load-state at the same file across a grid to warm-start each point from the previous one.") +integration_params.add_option("--sampler-load-state",default=None,help="AV only: load a saved live-volume state (.npz from --sampler-save-state) to warm-start this integration. Overrides --sampler-warmstart-samples.") +integration_params.add_option("--sampler-anisotropic-bins",action="store_true",help="AV only: give each extrinsic axis a DIFFERENT number of bins during contraction -- fine where the live points cluster tightly (phase/polarization/sky), coarse where they are broad (distance/inclination) -- instead of the default equal split. Keeps the same total bin budget, so the estimator is unchanged; helps AV wrap a correlated/degenerate posterior more tightly.") +optp.add_option("--reject-collapsed-live-volume",action='store_true',default=False, help="DROP an event whose adaptive-volume live volume degenerated (see the [AV COLLAPSE] report) instead of exporting it: the integration is treated as a failure, so no likelihood row, XML or posterior samples are written for it. Such a run's lnZ and samples describe a single mode of the integrand and are NOT a fair posterior draw, and nothing downstream can distinguish them from a converged export. Default off, because dropping the event silently THINS the posterior in an SNR-dependent way -- that was the pre-fix behaviour, when this case crashed. Left off, the event is exported but announces itself loudly and (with --mc-error-replicas>0) triggers replication. Turn it on when a contaminated point is worse than a missing one.") # L0 auto-rescue. Ported from bin/integrate_likelihood_extrinsic_batchmode; defaults and help # text kept IDENTICAL there and here on purpose -- see test_lisa_l0_rescue.py, which pins them. integration_params.add_option("--sampler-warmstart-retry-neff",type=float,default=None,help="AV or portfolio (L0 auto-rescue): if a pass finishes below this n_eff (i.e. it stalled on a very sharp / high-amplitude peak), automatically re-run a second pass warm-started from THIS point's own highest-likelihood samples. Same-problem reuse in the sense that the seed provably contains the peak the cold pass found -- but NOT that every mode is represented, so the warm pass can be biased low if the seed missed one. The rescue still runs as before; its result is rejected in favour of the cold pass only on positive evidence of lost mass (see --sampler-l0-rescue-reject-dlnZ). A portfolio is unaffected: its GMM member carries a defensive component. Directly targets the high-SNR n_eff LOTTERY (a large fraction of independent runs collapse to n_eff~1 by contracting onto the wrong spot); the rescue re-seeds a collapsed run from the peak it did find. Recommended for high-SNR events; e.g. 5.") @@ -1709,6 +1715,68 @@ def _clear_warm_state(sampler): sampler._warm_applied = False +def _maybe_load_av_state(sampler): + """Warm-start this integration from a saved AV live-volume state (--sampler-load-state).""" + try: + if opts.sampler_load_state and hasattr(sampler, 'load_state'): + print(" warm-start: loading saved sampler state from", opts.sampler_load_state) + sampler.load_state(opts.sampler_load_state) + except Exception as _e_ls: + print(" AV state load skipped (", _e_ls, ")") + + +def _maybe_save_av_state(sampler): + """Persist the adapted live-volume state for reuse by later instances/iterations.""" + if opts.sampler_method == 'AV' and opts.sampler_save_state and hasattr(sampler, 'save_state'): + try: + sampler.save_state(opts.sampler_save_state) + print(" AV: saved live-volume state to", opts.sampler_save_state) + except Exception as _e_ss: + print(" AV: could not save state (", _e_ss, ")") + + +def _maybe_enable_anisotropic_bins(sampler): + """Opt-in per-axis bin allocation, on the AV sampler AND any AV portfolio members.""" + if getattr(opts, 'sampler_anisotropic_bins', False): + _aniso_targets = [sampler] + list(getattr(sampler, 'portfolio_realizations', [])) + for _t in _aniso_targets: + if hasattr(_t, 'anisotropic_bins'): + _t.anisotropic_bins = True + print(" AV: anisotropic per-axis bin allocation ENABLED") + + +def _reject_if_collapsed(dd, stage): + """Apply --reject-collapsed-live-volume to whatever the CURRENT verdict is. + + In the main driver this is called TWICE -- once on the first run and again on the + replica pool, because replication can turn a healthy first run into a collapsed POOL. + This driver has no replica pooling yet, so only the first call exists here; the second + call site must be added WITH --mc-error-replicas, or the flag is silently bypassed for + exactly the case pooling introduces. + """ + if not opts.reject_collapsed_live_volume: + return + if not (isinstance(dd, dict) and dd.get('live_volume_collapsed', False)): + return + # Route through the ordinary failure path, so the caller skips this binary and writes no + # result row -- the pre-fix outcome, but for a stated reason. + _exc = mcsamplerAdaptiveVolume.LiveVolumeCollapse if mcsampler_AV_ok else RuntimeError + raise _exc( + "extrinsic integration collapsed ({}): live volume degenerated ({}); " + "--reject-collapsed-live-volume is set, so this event is being dropped " + "rather than exported".format(stage, dd.get('collapse_reason', ''))) + + +def _report_and_gate_collapse(dict_return, stage="first run"): + """Announce a collapsed live volume, then apply the rejection gate.""" + _collapsed = bool(dict_return.get('live_volume_collapsed', False)) if isinstance(dict_return, dict) else False + _collapse_reason = (dict_return or {}).get('collapse_reason', '') if isinstance(dict_return, dict) else '' + if _collapsed: + print(" [mc error] *** LIVE VOLUME COLLAPSED *** {}".format(_collapse_reason)) + print(" [mc error] this event's lnZ and exported samples are NOT a fair draw from the posterior.") + _reject_if_collapsed(dict_return, stage) + + def _maybe_load_nf_flow(sampler): """Warm-load a pre-trained normalizing flow so this instance skips/shortens training. @@ -2112,6 +2180,8 @@ def analyze_event_LISA(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_ lnL_oracles = np.zeros(opts.n_chunk) sampler.update_sampling_prior(lnL_oracles, opts.n_chunk, external_rvs=rvs_train,log_scale_weights=True,floor_integrated_probability=opts.adapt_floor_level) + _maybe_load_av_state(sampler) + _maybe_enable_anisotropic_bins(sampler) _maybe_load_nf_flow(sampler) res, var, neff, dict_return = sampler.integrate(like_to_integrate, *unpinned_params, **pinned_params) @@ -2127,11 +2197,15 @@ def analyze_event_LISA(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_ like_to_integrate, unpinned_params, pinned_params, lnL_offset=manual_avoid_overflow_logarithm) + _maybe_save_av_state(sampler) _maybe_save_nf_flow(sampler) if not(res): # no resut raise ValueError(" No integral result returned") + # Collapse gate AFTER the result check, matching the main driver's ordering. + _report_and_gate_collapse(dict_return, "first run") + if not(opts.internal_use_lnL): log_res = numpy.log(res) sqrt_var_over_res = numpy.sqrt(var)/res @@ -2832,6 +2906,8 @@ def analyze_event(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_spec_ lnL_oracles = np.zeros(opts.n_chunk) sampler.update_sampling_prior(lnL_oracles, opts.n_chunk, external_rvs=rvs_train,log_scale_weights=True,floor_integrated_probability=opts.adapt_floor_level) + _maybe_load_av_state(sampler) + _maybe_enable_anisotropic_bins(sampler) _maybe_load_nf_flow(sampler) res, var, neff, dict_return = sampler.integrate(like_to_integrate, *unpinned_params, **pinned_params) @@ -2847,11 +2923,15 @@ def analyze_event(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_spec_ like_to_integrate, unpinned_params, pinned_params, lnL_offset=manual_avoid_overflow_logarithm) + _maybe_save_av_state(sampler) _maybe_save_nf_flow(sampler) if not(res): # no resut raise ValueError(" No integral result returned") + # Collapse gate AFTER the result check, matching the main driver's ordering. + _report_and_gate_collapse(dict_return, "first run") + if not(opts.internal_use_lnL): log_res = numpy.log(res) sqrt_var_over_res = numpy.sqrt(var)/res diff --git a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/audit_lisa_driver_drift.py b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/audit_lisa_driver_drift.py index 50c16d832..94ced56c9 100644 --- a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/audit_lisa_driver_drift.py +++ b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/audit_lisa_driver_drift.py @@ -151,8 +151,17 @@ def compute_gap(): main = collect(MAIN) lisa = collect(LISA) gap, extras = [], [] + # A FUNC is satisfied by its BARE name as well as its qualified one. The main driver has + # ONE analyze_event and nests helpers inside it; this driver has TWO (analyze_event_LISA + # and analyze_event), so a helper ported here must be hoisted to module level or else + # duplicated -- and duplicating is the failure mode this audit exists to prevent. Without + # this, every correctly-hoisted port would sit in the gap forever as a false positive, + # which is how a gate gets trained out of people. + _lisa_bare = {n.rsplit(".", 1)[-1] for n in lisa["FUNC"]} for cat in ("FUNC", "OPTION", "CONST", "ATTR"): for name in sorted(set(main[cat]) - set(lisa[cat])): + if cat == "FUNC" and name.rsplit(".", 1)[-1] in _lisa_bare: + continue gap.append({"category": cat, "name": name, "key": "%s:%s" % (cat, name), "main_line": main[cat][name]}) for name in sorted(set(lisa[cat]) - set(main[cat])): diff --git a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json index dab8eac52..031a81ccb 100644 --- a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json +++ b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json @@ -57,10 +57,6 @@ "decision": "PORT", "reason": "Diagnostics for the replica triggers." }, - "FUNC:analyze_event._reject_if_collapsed": { - "decision": "PORT", - "reason": "Implementation of --reject-collapsed-live-volume." - }, "FUNC:dLofz": { "decision": "PHYSICS", "reason": "Cosmology helpers behind --d-prior-redshift. Same question: the interpolation range has to be re-chosen for MBHB redshifts." @@ -313,10 +309,6 @@ "decision": "PORT", "reason": "Pick a random event from the input file. Detector-agnostic; flagged dangerous in its own help text for oversampling reasons that apply equally to LISA." }, - "OPTION:--reject-collapsed-live-volume": { - "decision": "PORT", - "reason": "AV live-volume collapse rejection. AV is wired in the LISA driver identically." - }, "OPTION:--rotation-n-harmonics": { "decision": "NA", "reason": "Sidereal time-dependence of an EARTH-BASED antenna pattern F(t). The LISA constellation's motion is already carried by the LISA response itself (factored_likelihood_LISA + the h5/TDI frames), so this correction is both unnecessary and wrong there -- it would apply Earth rotation to a heliocentric detector." @@ -329,18 +321,6 @@ "decision": "NA", "reason": "Sidereal time-dependence of an EARTH-BASED antenna pattern F(t). The LISA constellation's motion is already carried by the LISA response itself (factored_likelihood_LISA + the h5/TDI frames), so this correction is both unnecessary and wrong there -- it would apply Earth rotation to a heliocentric detector." }, - "OPTION:--sampler-anisotropic-bins": { - "decision": "PORT", - "reason": "AV per-axis bin counts during contraction. AV is wired in LISA, and the argument for it is if anything stronger there: the LISA extrinsic axes are no more isotropic than the ground-based ones, and a sky pair that localizes tightly while distance stays broad is the exact case this exists for." - }, - "OPTION:--sampler-load-state": { - "decision": "PORT", - "reason": "AV live-volume state serialization. AV is wired in LISA; the state is the sampler's own internal grid, so it carries no LIGO-specific convention." - }, - "OPTION:--sampler-save-state": { - "decision": "PORT", - "reason": "AV live-volume state serialization. AV is wired in LISA; the state is the sampler's own internal grid, so it carries no LIGO-specific convention." - }, "OPTION:--sampler-sequential-warmstart": { "decision": "PORT", "reason": "Warm-start each intrinsic point from the previous one's cloud. Applies whenever --n-events-to-analyze>1, which LISA supports. Its snapshot/restore prerequisites (Finding 5) already landed with the L0 rescue, so this is now capture + the event-loop wiring only." diff --git a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py index 4ed4c8cd6..64ee43551 100644 --- a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py +++ b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py @@ -109,10 +109,15 @@ "with that pass rather than with the sequential warm start it is named for."), # --------------------------------------------------------------- L0 rescue / warm start - (r"^OPTION:--reject-collapsed-live-volume$", "PORT", - "AV live-volume collapse rejection. AV is wired in the LISA driver identically."), - (r"^FUNC:analyze_event\._reject_if_collapsed$", "PORT", - "Implementation of --reject-collapsed-live-volume."), + (r"^OPTION:--reject-collapsed-live-volume$", "PORTED", + "AV live-volume collapse rejection. AV is wired in the LISA driver identically. " + "NOTE the main driver calls its gate TWICE -- first run and replica pool -- and only " + "the first call exists here, because there is no pooling yet; the second MUST be " + "added with --mc-error-replicas or the flag is bypassed for the case pooling creates."), + (r"^FUNC:analyze_event\._reject_if_collapsed$", "PORTED", + "Hoisted to module level rather than nested, because this driver has TWO " + "analyze_event variants. The audit matches FUNC items on the bare name for exactly " + "this reason."), (r"^OPTION:--sampler-sequential-warmstart$", "PORT", "Warm-start each intrinsic point from the previous one's cloud. Applies whenever " "--n-events-to-analyze>1, which LISA supports. Its snapshot/restore prerequisites " @@ -120,12 +125,12 @@ "event-loop wiring only."), (r"^OPTION:--sampler-sequential-warmstart-cover-frac$", "PORT", "Coverage floor for the above; meaningless without it, so they travel together."), - (r"^OPTION:--sampler-anisotropic-bins$", "PORT", + (r"^OPTION:--sampler-anisotropic-bins$", "PORTED", "AV per-axis bin counts during contraction. AV is wired in LISA, and the argument " "for it is if anything stronger there: the LISA extrinsic axes are no more " "isotropic than the ground-based ones, and a sky pair that localizes tightly " "while distance stays broad is the exact case this exists for."), - (r"^OPTION:--sampler-(save|load)-state$", "PORT", + (r"^OPTION:--sampler-(save|load)-state$", "PORTED", "AV live-volume state serialization. AV is wired in LISA; the state is the " "sampler's own internal grid, so it carries no LIGO-specific convention."), (r"^OPTION:--sampler-warmstart-(cover-frac|inflate)$", "PORT", diff --git a/MonteCarloMarginalizeCode/Code/test/test_lisa_av_state.py b/MonteCarloMarginalizeCode/Code/test/test_lisa_av_state.py new file mode 100644 index 000000000..adb80c869 --- /dev/null +++ b/MonteCarloMarginalizeCode/Code/test/test_lisa_av_state.py @@ -0,0 +1,318 @@ +#!/usr/bin/env python +""" +Tests for the AV live-volume state, per-axis bin allocation and collapse gate ported into +the LISA ILE driver (bin/integrate_likelihood_extrinsic_batchmode_lisa). + +Four options, all sampler-agnostic: --sampler-save-state / --sampler-load-state (the AV +grid, which carries no detector convention), --sampler-anisotropic-bins, and +--reject-collapsed-live-volume. + +THE ONE THING TO KNOW. The main driver calls its collapse gate TWICE -- once on the first +run, and again on the replica POOL, because replication can turn a healthy first run into a +collapsed pool. This driver has no replica pooling yet, so only the first call exists here. +When --mc-error-replicas is ported the second call MUST come with it, or the flag is +silently bypassed for exactly the case pooling introduces. That is recorded at the helper, +in the drift ledger, and asserted below. +""" + +import ast +import os + +import pytest + +_HERE = os.path.dirname(os.path.abspath(__file__)) +_LISA = os.path.join(_HERE, '..', 'bin', 'integrate_likelihood_extrinsic_batchmode_lisa') +_MAIN = os.path.join(_HERE, '..', 'bin', 'integrate_likelihood_extrinsic_batchmode') + +OPTS = ["--sampler-save-state", "--sampler-load-state", + "--sampler-anisotropic-bins", "--reject-collapsed-live-volume"] + +HELPERS = ['_maybe_load_av_state', '_maybe_save_av_state', + '_maybe_enable_anisotropic_bins', '_reject_if_collapsed', + '_report_and_gate_collapse'] + + +def _src(path): + with open(path) as fh: + return fh.read() + + +def _option_nodes(path): + out = {} + for n in ast.walk(ast.parse(_src(path), filename=path)): + if (isinstance(n, ast.Call) and isinstance(n.func, ast.Attribute) + and n.func.attr in ("add_option", "add_argument")): + names = [a.value for a in n.args + if isinstance(a, ast.Constant) and isinstance(a.value, str)] + if names and names[0].startswith("--"): + out[names[0]] = n + return out + + +def _kwargs_of(node): + out = {} + for kw in node.keywords: + try: + out[kw.arg] = ast.literal_eval(kw.value) + except Exception: + out[kw.arg] = ast.dump(kw.value) + return out + + +class _Collapse(Exception): + pass + + +class _AVModule(object): + LiveVolumeCollapse = _Collapse + + +def _load(**optkw): + base = {"sampler_load_state": None, "sampler_save_state": None, + "sampler_anisotropic_bins": False, "reject_collapsed_live_volume": False, + "sampler_method": "AV"} + base.update(optkw) + defs = {n.name: n for n in ast.parse(_src(_LISA)).body + if isinstance(n, ast.FunctionDef) and n.name in HELPERS} + missing = sorted(set(HELPERS) - set(defs)) + assert not missing, "LISA driver is missing: %s" % missing + mod = ast.Module(body=[defs[n] for n in HELPERS], type_ignores=[]) + ns = {"opts": type("O", (), base)(), + "mcsamplerAdaptiveVolume": _AVModule, "mcsampler_AV_ok": True} + exec(compile(ast.fix_missing_locations(mod), "av_state", "exec"), ns) + return ns + + +# ------------------------------------------------------------------------------- options +@pytest.mark.parametrize("opt", OPTS) +def test_option_present_and_matches_the_main_driver(opt): + a, b = _kwargs_of(_option_nodes(_LISA)[opt]), _kwargs_of(_option_nodes(_MAIN)[opt]) + for key in ("default", "type", "action", "choices"): + assert a.get(key) == b.get(key), "%s: %s differs" % (opt, key) + + +# ---------------------------------------------------------------------------- load / save +class _AV(object): + def __init__(self): + self.loaded = self.saved = None + + def load_state(self, p): + self.loaded = p + + def save_state(self, p): + self.saved = p + + +class _NoState(object): + pass + + +def test_state_hooks_are_noops_when_unset(): + ns = _load() + s = _AV() + ns['_maybe_load_av_state'](s) + ns['_maybe_save_av_state'](s) + assert s.loaded is None and s.saved is None + + +def test_state_round_trip_reaches_the_sampler(): + ns = _load(sampler_load_state="/in.npz", sampler_save_state="/out.npz") + s = _AV() + ns['_maybe_load_av_state'](s) + ns['_maybe_save_av_state'](s) + assert s.loaded == "/in.npz" and s.saved == "/out.npz" + + +def test_save_state_is_restricted_to_the_AV_method(): + """Main gates the save on sampler_method == 'AV'; a portfolio's aggregate has no such grid.""" + ns = _load(sampler_save_state="/out.npz", sampler_method="portfolio") + s = _AV() + ns['_maybe_save_av_state'](s) + assert s.saved is None + + +def test_state_hooks_tolerate_a_sampler_without_state_support(): + ns = _load(sampler_load_state="/in.npz", sampler_save_state="/out.npz") + ns['_maybe_load_av_state'](_NoState()) + ns['_maybe_save_av_state'](_NoState()) + + +def test_a_bad_state_file_degrades_to_a_cold_run(): + """A missing/corrupt state must not kill the point.""" + ns = _load(sampler_load_state="/in.npz", sampler_save_state="/out.npz") + + class _Boom(object): + def load_state(self, p): + raise IOError("nope") + + def save_state(self, p): + raise IOError("read-only") + + ns['_maybe_load_av_state'](_Boom()) + ns['_maybe_save_av_state'](_Boom()) + + +# --------------------------------------------------------------------------- anisotropic +class _Binned(object): + anisotropic_bins = False + + +def test_anisotropic_bins_is_opt_in(): + ns = _load() + s = _Binned() + ns['_maybe_enable_anisotropic_bins'](s) + assert s.anisotropic_bins is False + + +def test_anisotropic_bins_reaches_portfolio_members_too(): + """The grid lives on the MEMBERS; setting it only on the aggregate would do nothing.""" + ns = _load(sampler_anisotropic_bins=True) + m1, m2 = _Binned(), _Binned() + s = _Binned() + s.portfolio_realizations = [m1, m2] + ns['_maybe_enable_anisotropic_bins'](s) + assert s.anisotropic_bins and m1.anisotropic_bins and m2.anisotropic_bins + + +def test_anisotropic_bins_skips_members_that_do_not_support_it(): + ns = _load(sampler_anisotropic_bins=True) + s = _Binned() + s.portfolio_realizations = [_NoState()] + ns['_maybe_enable_anisotropic_bins'](s) # must not raise + assert s.anisotropic_bins is True + + +# -------------------------------------------------------------------------- collapse gate +COLLAPSED = {'live_volume_collapsed': True, 'collapse_reason': 'zero volume'} +HEALTHY = {'live_volume_collapsed': False} + + +def test_gate_is_inert_when_the_flag_is_off(): + _load()['_reject_if_collapsed'](COLLAPSED, "first run") + + +def test_gate_is_inert_on_a_healthy_run(): + _load(reject_collapsed_live_volume=True)['_reject_if_collapsed'](HEALTHY, "first run") + + +def test_gate_raises_when_flag_set_and_run_collapsed(): + with pytest.raises(_Collapse): + _load(reject_collapsed_live_volume=True)['_reject_if_collapsed'](COLLAPSED, "first run") + + +def test_gate_message_names_the_stage_and_reason(): + """The stage is in the message because the main driver calls this at two stages.""" + with pytest.raises(_Collapse) as e: + _load(reject_collapsed_live_volume=True)['_reject_if_collapsed'](COLLAPSED, "pooled") + assert "pooled" in str(e.value) and "zero volume" in str(e.value) + + +@pytest.mark.parametrize("dd", [None, "not a dict", {}]) +def test_gate_tolerates_a_missing_or_malformed_dict_return(dd): + _load(reject_collapsed_live_volume=True)['_reject_if_collapsed'](dd, "first run") + + +def test_report_announces_a_collapse_even_when_the_gate_is_off(capsys): + """Not rejecting is not the same as not telling anyone.""" + ns = _load() + capsys.readouterr() + ns['_report_and_gate_collapse'](COLLAPSED) + out = capsys.readouterr().out + assert "LIVE VOLUME COLLAPSED" in out and "NOT a fair draw" in out + + +def test_report_says_nothing_on_a_healthy_run(capsys): + ns = _load() + capsys.readouterr() + ns['_report_and_gate_collapse'](HEALTHY) + assert "COLLAPSED" not in capsys.readouterr().out + + +def test_report_still_raises_when_gated(): + with pytest.raises(_Collapse): + _load(reject_collapsed_live_volume=True)['_report_and_gate_collapse'](COLLAPSED) + + +def test_gate_falls_back_to_RuntimeError_without_AV(): + """mcsampler_AV_ok False -> the AV exception class is unavailable.""" + defs = {n.name: n for n in ast.parse(_src(_LISA)).body + if isinstance(n, ast.FunctionDef) and n.name == '_reject_if_collapsed'} + mod = ast.Module(body=[defs['_reject_if_collapsed']], type_ignores=[]) + ns = {"opts": type("O", (), {"reject_collapsed_live_volume": True})(), + "mcsamplerAdaptiveVolume": None, "mcsampler_AV_ok": False} + exec(compile(ast.fix_missing_locations(mod), "av_state", "exec"), ns) + with pytest.raises(RuntimeError): + ns['_reject_if_collapsed'](COLLAPSED, "first run") + + +# ------------------------------------------------------------------------- call-site wiring +def test_both_analyze_event_variants_get_every_hook(): + tree = ast.parse(_src(_LISA)) + fns = {n.name: n for n in tree.body + if isinstance(n, ast.FunctionDef) and n.name in ('analyze_event', 'analyze_event_LISA')} + assert set(fns) == {'analyze_event', 'analyze_event_LISA'} + for name, node in fns.items(): + called = {c.func.id for c in ast.walk(node) + if isinstance(c, ast.Call) and isinstance(c.func, ast.Name)} + for hook in ('_maybe_load_av_state', '_maybe_save_av_state', + '_maybe_enable_anisotropic_bins', '_report_and_gate_collapse'): + assert hook in called, "%s does not call %s" % (name, hook) + + +def test_hook_ordering_at_both_call_sites(): + """load/aniso before the integration, save after it, the gate after the result check.""" + src = _src(_LISA) + pos = 0 + for _ in range(2): + load = src.index("_maybe_load_av_state(sampler)", pos) + aniso = src.index("_maybe_enable_anisotropic_bins(sampler)", load) + integ = src.index("sampler.integrate(like_to_integrate", aniso) + save = src.index("_maybe_save_av_state(sampler)", integ) + guard = src.index("if not(res): # no resut", save) + gate = src.index("_report_and_gate_collapse(dict_return", guard) + assert load < aniso < integ < save < guard < gate + pos = gate + 1 + + +def test_the_second_gate_call_site_is_recorded_as_missing(): + """Main gates twice; this driver gates once because it has no replica pooling yet. + + If someone ports --mc-error-replicas without adding the second call, the flag is + silently bypassed for the case pooling creates. This asserts the warning is still + written down where that person will be working. + """ + src = _src(_LISA) + fn = src[src.index("def _reject_if_collapsed"):] + fn = fn[:fn.index("\ndef ")] + assert "mc-error-replicas" in fn and "TWICE" in fn + # Count CALLS, not textual occurrences: the `def` line matches the same substring. + calls = [c for c in ast.walk(ast.parse(src)) + if isinstance(c, ast.Call) and isinstance(c.func, ast.Name) + and c.func.id == '_report_and_gate_collapse'] + assert len(calls) == 2, \ + "expected exactly one gate call per analyze_event variant, found %d" % len(calls) + + +# ---------------------------------------------------------------- anti-drift vs the main driver +def _named(path, name): + for n in ast.walk(ast.parse(_src(path))): + if isinstance(n, ast.FunctionDef) and n.name == name: + return n + raise AssertionError("%s not found in %s" % (name, os.path.basename(path))) + + +def _normalized(fn): + node = ast.parse(ast.unparse(fn)).body[0] if hasattr(ast, "unparse") else fn + body = list(node.body) + if (body and isinstance(body[0], ast.Expr) + and isinstance(getattr(body[0], "value", None), ast.Constant) + and isinstance(body[0].value.value, str)): + body = body[1:] + return ast.dump(ast.fix_missing_locations(ast.Module(body=body, type_ignores=[]))) + + +def test_reject_if_collapsed_body_is_identical_to_the_main_drivers(): + """Hoisted out of analyze_event here, but the body must not have changed with it.""" + assert (_normalized(_named(_LISA, '_reject_if_collapsed')) + == _normalized(_named(_MAIN, '_reject_if_collapsed'))), \ + "_reject_if_collapsed has drifted between the two drivers (docstrings excluded)" From b0e0f5585c8bf3d570250f95d282d7e77b334941 Mon Sep 17 00:00:00 2001 From: Richard O'Shaughnessy Date: Sun, 16 Aug 2026 06:59:48 -0400 Subject: [PATCH 3/3] Fix consolidated LISA sampler state plumbing --- ...egrate_likelihood_extrinsic_batchmode_lisa | 55 +++------ .../integrators/lisa_drift_ledger.json | 8 ++ .../integrators/make_lisa_drift_ledger.py | 10 +- .../Code/test/test_lisa_av_state.py | 20 ++- .../Code/test/test_lisa_l0_rescue.py | 14 +++ .../Code/test/test_lisa_sampler_plumbing.py | 116 +----------------- 6 files changed, 62 insertions(+), 161 deletions(-) diff --git a/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa b/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa index ff07c08d9..f5c669eb0 100755 --- a/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa +++ b/MonteCarloMarginalizeCode/Code/bin/integrate_likelihood_extrinsic_batchmode_lisa @@ -306,9 +306,8 @@ integration_params.add_option("--internal-use-lnL",action='store_true',help="lik integration_params.add_option("--sampler-method",default="adaptive_cartesian_gpu",help="adaptive_cartesian|GMM|adaptive_cartesian_gpu") integration_params.add_option("--sampler-portfolio",default=None,action='append',type=str,help="comma-separated strings, matching sampler methods other than portfolio") integration_params.add_option("--sampler-portfolio-args",default=None, action='append', type=str, help='eval-able dictionaryo to be passed to that sampler') -# Portfolio freeze/allocation policy and NF flow persistence. Pure pass-through to the -# shared samplers, which this driver already wires (identical ok_lnL_methods). Definitions -# copied verbatim from bin/integrate_likelihood_extrinsic_batchmode; pinned by +# Portfolio freeze/allocation policy. Pure pass-through to the shared portfolio sampler. +# Definitions copied verbatim from bin/integrate_likelihood_extrinsic_batchmode; pinned by # test_lisa_sampler_plumbing.py. integration_params.add_option("--portfolio-adaptive-alloc",action='store_true',default=False,help="Portfolio: ENABLE (opt-in) adaptive-probe draw allocation -- concentrate draws on the best per-chunk-n_ess member. Good on strongly-correlated targets; NOT recommended for AV-favorable high-SNR events (it starves the slow-contracting AV workhorse). Off by default (legacy n_ess reweighting).") integration_params.add_option("--portfolio-alloc-exponent",default=None,type=float,help="Portfolio: adaptive allocation ~ member_quality^exponent. Higher concentrates harder on the winner. Sampler default 1.0.") @@ -322,8 +321,6 @@ integration_params.add_option("--portfolio-varaha-max-frac",default=None,type=fl integration_params.add_option("--portfolio-varaha-min-frac",default=None,type=float,help="Portfolio: reserve this combined DRAW fraction for VARAHA/AV members (0/unset = off). never-freeze keeps a VARAHA member UPDATING, but both allocation rules score by per-chunk n_ess, which sits at ~1 during VARAHA's slow cumulative contraction -- so a member that looks instantly good can take nearly the whole budget (measured on S250114ax post-#33: GMM took ~0.84 and the portfolio collapsed to n_eff ~2 vs ~100 for standalone AV). Unbiased for any allocation (q_mix); trades efficiency only.") integration_params.add_option("--portfolio-varaha-never-freeze",action='store_true',default=False,help="Portfolio: VARAHA/AV members always update every chunk past their breakpoint (freeze-exempt). This is the sampler default; the flag is here for explicitness/pipe pass-through.") integration_params.add_option("--portfolio-weight-clip",default=None,type=float,help="Portfolio: OPT-IN truncated importance sampling applied to the PROPOSAL-FIT INPUT ONLY. Caps the weights fed to member.update_sampling_prior (the GMM covariance fit) at tau = C*sqrt(n)*mean(w) (0/unset = off; C~1 is the standard Ionides choice), so one enormous weight cannot make that fit degenerate. The estimator (ln Z, n_eff), the n_ess report, and the allocation signal all use the TRUE unclipped weights, so they stay exactly unbiased and undistorted. Do NOT clip the estimator (measured on S250114ax: n_eff=100 2x faster than AV but ln Z biased -11.5 nats) or the n_ess report (clipping inflates the clipped member's n_ess and starves the AV workhorse). The withheld tail mass is tracked and reported as a diagnostic. NOTE: if huge weights come from q_mix UNDERFLOW (watch for the warning) they are a numerical artifact, not tail mass.") -integration_params.add_option("--nf-flow-load",default=None,help="NF only: load a pre-trained normalizing flow (.pt from --nf-flow-save); with --n-adapt 0 this reuses it directly (skips training), otherwise it is polished.") -integration_params.add_option("--nf-flow-save",default=None,help="NF only: after integration, serialize the trained normalizing flow (.pt) for reuse across ILE instances.") integration_params.add_option("--sampler-xpy",default=None,help="numpy|cupy if the adaptive_cartesian_gpu sampler is active, use that.") # AV live-volume state, per-axis bin allocation, and the collapse gate. Copied verbatim # from bin/integrate_likelihood_extrinsic_batchmode; pinned by test_lisa_av_state.py. @@ -1732,6 +1729,9 @@ def _maybe_load_av_state(sampler): def _maybe_save_av_state(sampler): """Persist the adapted live-volume state for reuse by later instances/iterations.""" if opts.sampler_method == 'AV' and opts.sampler_save_state and hasattr(sampler, 'save_state'): + if not getattr(sampler, '_av_state_reuse_safe', True): + print(" AV: not saving live-volume state from a rejected/failed rescue") + return try: sampler.save_state(opts.sampler_save_state) print(" AV: saved live-volume state to", opts.sampler_save_state) @@ -1781,30 +1781,6 @@ def _report_and_gate_collapse(dict_return, stage="first run"): _reject_if_collapsed(dict_return, stage) -def _maybe_load_nf_flow(sampler): - """Warm-load a pre-trained normalizing flow so this instance skips/shortens training. - - hasattr-guarded, so it is a no-op for every sampler that is not NF -- including all five - this driver lists in ok_lnL_methods, where NF is reachable only as a portfolio member. - """ - if opts.nf_flow_load and hasattr(sampler, 'load_flow'): - try: - print(" NF: loading pre-trained flow from", opts.nf_flow_load) - sampler.load_flow(opts.nf_flow_load) - except Exception as _e_nf: - print(" NF flow load skipped (", _e_nf, ")") - - -def _maybe_save_nf_flow(sampler): - """Serialize the trained flow for reuse across ILE instances. hasattr-guarded as above.""" - if opts.nf_flow_save and hasattr(sampler, 'save_flow'): - try: - sampler.save_flow(opts.nf_flow_save) - print(" NF: saved trained flow to", opts.nf_flow_save) - except Exception as _e_fs: - print(" NF: could not save flow (", _e_fs, ")") - - def _maybe_l0_rescue(sampler, res, var, neff, dict_return, like_to_integrate, unpinned_params, pinned_params, lnL_offset=0.0): @@ -1823,6 +1799,9 @@ def _maybe_l0_rescue(sampler, res, var, neff, dict_return, `lnL_offset` is this event's manual_avoid_overflow_logarithm, used only to print absolute lnZ values. It is a local of the caller in both analyze_event variants, hence a parameter. """ + # Reset per event before ANY early return. The sampler is reused: a rejected rescue on + # one event must not suppress saving a later healthy event that needs no rescue at all. + sampler._av_state_reuse_safe = True # APPLICABILITY FIRST, then n_eff. The main driver evaluates # _neff_val = None if neff is None else float(sampler.identity_convert(neff)) # BEFORE its guard, which is safe there only by luck: identity_convert comes from @@ -1953,6 +1932,7 @@ def _maybe_l0_rescue(sampler, res, var, neff, dict_return, # the warm pass's record while _rvs, the estimate and the diagnostics all # describe the cold one. Nothing reads it in this driver today. res, var, neff, dict_return = _restore_pass_state(sampler, _cold_state_l0) + sampler._av_state_reuse_safe = False _clear_warm_state(sampler) except Exception as _e_l0: # "skipped" is only true if the warm pass never started. If it raised PARTWAY THROUGH @@ -1968,6 +1948,7 @@ def _maybe_l0_rescue(sampler, res, var, neff, dict_return, " samples; restoring the COLD pass so the reported diagnostics and the" " exported samples describe the same integral.") res, var, neff, dict_return = _restore_pass_state(sampler, _cold_state_l0) + sampler._av_state_reuse_safe = False _clear_warm_state(sampler) return res, var, neff, dict_return @@ -2201,8 +2182,6 @@ def analyze_event_LISA(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_ _maybe_load_av_state(sampler) _maybe_enable_anisotropic_bins(sampler) - _maybe_load_nf_flow(sampler) - res, var, neff, dict_return = sampler.integrate(like_to_integrate, *unpinned_params, **pinned_params) # L0 auto-rescue: on a very sharply-peaked (high-amplitude) point a cold AV can stall at @@ -2216,14 +2195,15 @@ def analyze_event_LISA(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_ like_to_integrate, unpinned_params, pinned_params, lnL_offset=manual_avoid_overflow_logarithm) - _maybe_save_av_state(sampler) - _maybe_save_nf_flow(sampler) - if not(res): # no resut raise ValueError(" No integral result returned") # Collapse gate AFTER the result check, matching the main driver's ordering. _report_and_gate_collapse(dict_return, "first run") + # Persist only a result we are actually willing to report. In particular, never write + # a collapsed grid that --reject-collapsed-live-volume just rejected, nor a warm grid + # whose rescue result was rejected/failed and replaced by the cold estimate. + _maybe_save_av_state(sampler) if not(opts.internal_use_lnL): log_res = numpy.log(res) @@ -2927,8 +2907,6 @@ def analyze_event(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_spec_ _maybe_load_av_state(sampler) _maybe_enable_anisotropic_bins(sampler) - _maybe_load_nf_flow(sampler) - res, var, neff, dict_return = sampler.integrate(like_to_integrate, *unpinned_params, **pinned_params) # L0 auto-rescue: on a very sharply-peaked (high-amplitude) point a cold AV can stall at @@ -2942,14 +2920,13 @@ def analyze_event(P_list, indx_event, data_dict, psd_dict, fmax, opts, inv_spec_ like_to_integrate, unpinned_params, pinned_params, lnL_offset=manual_avoid_overflow_logarithm) - _maybe_save_av_state(sampler) - _maybe_save_nf_flow(sampler) - if not(res): # no resut raise ValueError(" No integral result returned") # Collapse gate AFTER the result check, matching the main driver's ordering. _report_and_gate_collapse(dict_return, "first run") + # See the LISA variant above: only persist accepted, reusable AV state. + _maybe_save_av_state(sampler) if not(opts.internal_use_lnL): log_res = numpy.log(res) diff --git a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json index 031a81ccb..39714169e 100644 --- a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json +++ b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/lisa_drift_ledger.json @@ -305,6 +305,14 @@ "decision": "NA", "reason": "The .dslice export and its placement/tuning knobs. This is a data product for a downstream LIGO CIP distance workflow that the LISA pipeline does not run; there is no consumer. If a LISA distance workflow is ever built, note that the .dslice reweight core was the third Finding-2 site and must not be revived in its pre-#87 form." }, + "OPTION:--nf-flow-load": { + "decision": "PORT", + "reason": "Normalizing-flow persistence is detector-agnostic, but the LISA portfolio factory currently constructs only AV, GMM, and adaptive_cartesian_gpu members. Port the NF member construction and route load/save to that member before exposing these flags; hooks on the portfolio aggregate are a silent no-op because it has no flow API." + }, + "OPTION:--nf-flow-save": { + "decision": "PORT", + "reason": "Normalizing-flow persistence is detector-agnostic, but the LISA portfolio factory currently constructs only AV, GMM, and adaptive_cartesian_gpu members. Port the NF member construction and route load/save to that member before exposing these flags; hooks on the portfolio aggregate are a silent no-op because it has no flow API." + }, "OPTION:--random-event": { "decision": "PORT", "reason": "Pick a random event from the input file. Detector-agnostic; flagged dangerous in its own help text for oversampling reasons that apply equally to LISA." diff --git a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py index d0798e3e7..ce51b8cc2 100644 --- a/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py +++ b/MonteCarloMarginalizeCode/Code/test/expensive_before_merging/integrators/make_lisa_drift_ledger.py @@ -175,11 +175,11 @@ "--portfolio-varaha-can-freeze wins over --portfolio-varaha-never-freeze, as there."), # ------------------------------------------------------------------- NF flow plumbing - (r"^OPTION:--nf-flow-(load|save)$", "PORTED", - "Normalizing-flow persistence. Neither driver lists an NF method in ok_lnL_methods " - "(identical lists), so NF is reached only as a portfolio member -- equally " - "available to LISA. Both hooks are hasattr-guarded, so they are a no-op for every " - "other sampler."), + (r"^OPTION:--nf-flow-(load|save)$", "PORT", + "Normalizing-flow persistence is detector-agnostic, but the LISA portfolio factory " + "currently constructs only AV, GMM, and adaptive_cartesian_gpu members. Port the NF " + "member construction and route load/save to that member before exposing these flags; " + "hooks on the portfolio aggregate are a silent no-op because it has no flow API."), # --------------------------------------------------------- extrinsic proposal handoff (r"^OPTION:--extrinsic-proposal-output$", "PORT", diff --git a/MonteCarloMarginalizeCode/Code/test/test_lisa_av_state.py b/MonteCarloMarginalizeCode/Code/test/test_lisa_av_state.py index adb80c869..e578b242f 100644 --- a/MonteCarloMarginalizeCode/Code/test/test_lisa_av_state.py +++ b/MonteCarloMarginalizeCode/Code/test/test_lisa_av_state.py @@ -131,6 +131,16 @@ def test_save_state_is_restricted_to_the_AV_method(): assert s.saved is None +def test_rejected_or_failed_rescue_state_is_not_saved(capsys): + """The warm grid may outlive restoration of the cold result; never persist that mismatch.""" + ns = _load(sampler_save_state="/out.npz") + s = _AV() + s._av_state_reuse_safe = False + ns['_maybe_save_av_state'](s) + assert s.saved is None + assert "not saving" in capsys.readouterr().out + + def test_state_hooks_tolerate_a_sampler_without_state_support(): ns = _load(sampler_load_state="/in.npz", sampler_save_state="/out.npz") ns['_maybe_load_av_state'](_NoState()) @@ -260,18 +270,18 @@ def test_both_analyze_event_variants_get_every_hook(): def test_hook_ordering_at_both_call_sites(): - """load/aniso before the integration, save after it, the gate after the result check.""" + """Only a nonempty, collapse-approved result may persist its live-volume state.""" src = _src(_LISA) pos = 0 for _ in range(2): load = src.index("_maybe_load_av_state(sampler)", pos) aniso = src.index("_maybe_enable_anisotropic_bins(sampler)", load) integ = src.index("sampler.integrate(like_to_integrate", aniso) - save = src.index("_maybe_save_av_state(sampler)", integ) - guard = src.index("if not(res): # no resut", save) + guard = src.index("if not(res): # no resut", integ) gate = src.index("_report_and_gate_collapse(dict_return", guard) - assert load < aniso < integ < save < guard < gate - pos = gate + 1 + save = src.index("_maybe_save_av_state(sampler)", gate) + assert load < aniso < integ < guard < gate < save + pos = save + 1 def test_the_second_gate_call_site_is_recorded_as_missing(): diff --git a/MonteCarloMarginalizeCode/Code/test/test_lisa_l0_rescue.py b/MonteCarloMarginalizeCode/Code/test/test_lisa_l0_rescue.py index 9ad975e79..ca3793ccf 100644 --- a/MonteCarloMarginalizeCode/Code/test/test_lisa_l0_rescue.py +++ b/MonteCarloMarginalizeCode/Code/test/test_lisa_l0_rescue.py @@ -398,6 +398,7 @@ def test_accepted_warm_pass_replaces_the_cold_result(): s.warm_rvs = _rec([0.0, 0.0]) # same lnZ -> no evidence of loss out = _run(H, s) assert out == ('R2', 'V2', 42.0, {'warm': True}) + assert s._av_state_reuse_safe is True def test_warm_pass_far_below_cold_is_rejected_and_cold_is_restored(): @@ -412,6 +413,18 @@ def test_warm_pass_far_below_cold_is_rejected_and_cold_is_restored(): assert out == ('R1', 'V1', 1.0, {'cold': True}), "the warm pass was not rejected" assert s._warm_seed_reserve == {'tag': 'cold'}, "the reserve did not come back (Finding 5)" assert np.allclose(s._rvs['log_integrand'], cold['log_integrand']) + assert s._av_state_reuse_safe is False, "the rejected warm grid could be persisted" + + +def test_a_later_healthy_event_resets_the_state_save_veto(): + """Sampler objects are reused; an earlier rejection must not poison later state saves.""" + H = _load() + s = _Sampler(rvs=_rec([0.0, 0.0]), integrate_result=('R2', 'V2', 42.0, {})) + s.warm_rvs = _rec([-20.0, -20.0]) + _run(H, s) # rejected warm pass + assert s._av_state_reuse_safe is False + _run(H, s, neff=42.0) # healthy next event; returns before attempting a rescue + assert s._av_state_reuse_safe is True def test_reject_message_reports_lnZ_on_the_events_offset_scale(capsys): @@ -469,6 +482,7 @@ def test_a_raising_warm_pass_restores_the_cold_state(): assert np.allclose(s._rvs['log_integrand'], cold['log_integrand']), \ "cold diagnostics were reported beside a warm export" assert s._warm_seed_reserve == {'tag': 'cold'} + assert s._av_state_reuse_safe is False, "the failed warm grid could be persisted" def test_rescue_clears_warm_state_afterwards(): diff --git a/MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py b/MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py index 5ba3e0ca4..976d17169 100644 --- a/MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py +++ b/MonteCarloMarginalizeCode/Code/test/test_lisa_sampler_plumbing.py @@ -1,7 +1,7 @@ #!/usr/bin/env python """ -Tests for the portfolio freeze/allocation policy and NF flow persistence ported into the -LISA ILE driver (bin/integrate_likelihood_extrinsic_batchmode_lisa). +Tests for the portfolio freeze/allocation policy ported into the LISA ILE driver +(bin/integrate_likelihood_extrinsic_batchmode_lisa). This is pure PASS-THROUGH plumbing to samplers the LISA driver already wires -- it exposes the same ``ok_lnL_methods`` as the main driver (``GMM, adaptive_cartesian, @@ -17,7 +17,6 @@ own default with nothing. The assembly's whole shape -- ``if opts.x is not None`` -- exists for that, and a single dropped guard is invisible until a run behaves oddly. * the two mutually-exclusive VARAHA flags resolving the wrong way round. - * an NF hook that is not hasattr-guarded, which would break every non-NF sampler. The freeze-policy assembly is inline in both drivers (not a function), so it is exercised here by extracting the block and exec'ing it against a fake ``opts``. That tests the real @@ -42,7 +41,6 @@ "--portfolio-varaha-max-frac", "--portfolio-varaha-min-frac", "--portfolio-varaha-never-freeze", "--portfolio-weight-clip", ] -NF_OPTS = ["--nf-flow-load", "--nf-flow-save"] def _src(path): @@ -84,13 +82,13 @@ def opts_main(): # ------------------------------------------------------------------------------ presence -@pytest.mark.parametrize("opt", PORTFOLIO_OPTS + NF_OPTS) +@pytest.mark.parametrize("opt", PORTFOLIO_OPTS) def test_option_is_present_in_the_lisa_driver(opt, opts_lisa): assert opt in opts_lisa # ------------------------------------------------------------------------------- defaults -@pytest.mark.parametrize("opt", PORTFOLIO_OPTS + NF_OPTS) +@pytest.mark.parametrize("opt", PORTFOLIO_OPTS) def test_option_signature_matches_the_main_driver(opt, opts_lisa, opts_main): """Same default, same type, same action. @@ -223,109 +221,3 @@ def test_assembly_result_is_actually_handed_to_setup(): """Building the dict and not passing it would be a silent no-op.""" src = _src(_LISA) assert "sampler.setup(portfolio_args=opts.sampler_portfolio_args, **_freeze_policy_kwargs" in src - - -# ------------------------------------------------------------------------------- NF hooks -def _fn(path, name): - for n in ast.parse(_src(path)).body: - if isinstance(n, ast.FunctionDef) and n.name == name: - return n - raise AssertionError("%s not found in %s" % (name, os.path.basename(path))) - - -def _load_nf(**optkw): - ns = {"opts": type("O", (), dict({"nf_flow_load": None, "nf_flow_save": None}, **optkw))()} - mod = ast.Module(body=[_fn(_LISA, '_maybe_load_nf_flow'), _fn(_LISA, '_maybe_save_nf_flow')], - type_ignores=[]) - exec(compile(ast.fix_missing_locations(mod), "nf", "exec"), ns) - return ns - - -class _NoFlow(object): - """A sampler with no flow support -- i.e. every sampler in ok_lnL_methods.""" - - -class _WithFlow(object): - def __init__(self): - self.loaded = self.saved = None - - def load_flow(self, path): - self.loaded = path - - def save_flow(self, path): - self.saved = path - - -def test_nf_hooks_are_noops_when_the_options_are_unset(): - ns = _load_nf() - s = _WithFlow() - ns['_maybe_load_nf_flow'](s) - ns['_maybe_save_nf_flow'](s) - assert s.loaded is None and s.saved is None - - -def test_nf_hooks_are_noops_for_a_sampler_without_flow_support(capsys): - """hasattr-guarded: must DECLINE for AV/GMM/portfolio/adaptive_cartesian. - - Asserting "does not raise" is not enough and an earlier version of this test made - exactly that mistake: the body is wrapped in `except Exception`, so dropping the - hasattr guard still does not raise -- it announces "loading pre-trained flow", calls a - method that does not exist, and swallows the AttributeError. Every non-NF run would - then log a flow load that never happened. So the observable property is that the hook - says NOTHING and touches nothing when the sampler has no flow support. - """ - ns = _load_nf(nf_flow_load="/x/flow.pt", nf_flow_save="/x/flow.pt") - capsys.readouterr() - ns['_maybe_load_nf_flow'](_NoFlow()) - ns['_maybe_save_nf_flow'](_NoFlow()) - out = capsys.readouterr().out - assert "NF" not in out, ( - "the hook engaged a sampler with no flow support (and the except swallowed it): %r" % out) - - -def test_nf_load_and_save_reach_a_flow_capable_sampler(): - ns = _load_nf(nf_flow_load="/in.pt", nf_flow_save="/out.pt") - s = _WithFlow() - ns['_maybe_load_nf_flow'](s) - ns['_maybe_save_nf_flow'](s) - assert s.loaded == "/in.pt" and s.saved == "/out.pt" - - -def test_nf_failures_do_not_abort_the_event(): - """A missing/corrupt flow file must degrade to a cold run, not kill the point.""" - ns = _load_nf(nf_flow_load="/in.pt", nf_flow_save="/out.pt") - - class _Boom(object): - def load_flow(self, p): - raise IOError("no such file") - - def save_flow(self, p): - raise IOError("read-only") - - ns['_maybe_load_nf_flow'](_Boom()) - ns['_maybe_save_nf_flow'](_Boom()) - - -# ------------------------------------------------------------------------ call-site wiring -def test_both_analyze_event_variants_get_the_nf_hooks(): - """This driver has two; wiring only one is a silent half-port.""" - tree = ast.parse(_src(_LISA)) - fns = {n.name: n for n in tree.body - if isinstance(n, ast.FunctionDef) and n.name in ('analyze_event', 'analyze_event_LISA')} - assert set(fns) == {'analyze_event', 'analyze_event_LISA'} - for name, node in fns.items(): - called = {c.func.id for c in ast.walk(node) - if isinstance(c, ast.Call) and isinstance(c.func, ast.Name)} - assert '_maybe_load_nf_flow' in called, "%s never loads the flow" % name - assert '_maybe_save_nf_flow' in called, "%s never saves the flow" % name - - -def test_flow_is_loaded_before_the_integration_and_saved_after(): - src = _src(_LISA) - pos = 0 - for _ in range(2): - load = src.index("_maybe_load_nf_flow(sampler)", pos) - integ = src.index("sampler.integrate(like_to_integrate", load) - save = src.index("_maybe_save_nf_flow(sampler)", integ) - assert load < integ < save, "flow load/save straddle the integration incorrectly" - pos = save + 1