diff --git a/src/curate_questions/create_question_set/main.py b/src/curate_questions/create_question_set/main.py index 87266d5b..426c4477 100644 --- a/src/curate_questions/create_question_set/main.py +++ b/src/curate_questions/create_question_set/main.py @@ -15,7 +15,6 @@ import json import logging import os -import random import sys from collections.abc import Callable from copy import deepcopy @@ -23,6 +22,7 @@ from enum import Enum from fractions import Fraction +import numpy as np import pandas as pd from tqdm import tqdm from utils import gcp @@ -94,15 +94,19 @@ def process_questions( single_generation_func: Callable, show_plots: bool, question_set_target: QuestionSetTarget, + random_state: int | np.random.RandomState | None = None, ) -> dict: """Sample from `questions` to get the number of questions needed. Args: questions (dict): Source questions keyed by source name, each with a "dfq" DataFrame to_questions (dict): Allocation info keyed by source with "num_questions_to_sample" - single_generation_func (Callable): Sampling function taking (values, n) and returning DataFrame + single_generation_func (Callable): Sampling function taking (values, n, random_state) and + returning a DataFrame show_plots (bool): Whether to display distribution plots question_set_target (QuestionSetTarget): Target question set ("llm" or "human") + random_state: Seed/``np.random.RandomState`` threaded to ``single_generation_func`` for + reproducibility. ``None`` (default) samples without a fixed seed. Returns processed_questions (dict): Deep copy of questions with sampled DataFrames @@ -118,7 +122,7 @@ def process_questions( df_available = values["dfq"].copy() # Sample questions for this source - values["dfq"] = single_generation_func(values, num_single) + values["dfq"] = single_generation_func(values, num_single, random_state) df_sampled = values["dfq"] num_found += len(df_sampled) @@ -180,20 +184,23 @@ def process_questions( return processed_questions -def human_sample_questions(values: dict, n_single: int) -> pd.DataFrame: +def human_sample_questions( + values: dict, n_single: int, random_state: int | np.random.RandomState | None = None +) -> pd.DataFrame: """Get questions for the human question set by sampling from LLM questions. Args: values (dict): Source data dict containing "dfq" DataFrame n_single (int): Number of questions to sample + random_state: Seed/``np.random.RandomState`` for reproducible sampling (anything + ``DataFrame.sample`` accepts). ``None`` (default) samples without a fixed seed. Returns dfq (pd.DataFrame): Randomly sampled questions """ dfq = values["dfq"].copy() - indices_to_sample_from = dfq.index.tolist() - indices = random.sample(indices_to_sample_from, min(n_single, len(indices_to_sample_from))) - return dfq.loc[indices] + n = min(n_single, len(dfq)) + return dfq.sample(n=n, random_state=random_state) def get_bin_label(bin_config: dict, bin_type: str) -> str: @@ -789,7 +796,23 @@ def add_line_chart( fig.show() -def stratified_sample_questions(dfq: pd.DataFrame, n_target: int) -> pd.DataFrame: +def _as_random_state( + random_state: int | np.random.RandomState | None, +) -> np.random.RandomState | None: + """Normalize a seed/RandomState/None to a RandomState (or None). + + Returns ``None`` unchanged (preserving unseeded behavior) and a ``RandomState`` unchanged, but + promotes a bare int seed to a ``RandomState`` so successive ``.sample`` calls in a loop + *decorrelate* (advancing one shared generator) instead of reusing the same seed per iteration. + """ + if random_state is None or isinstance(random_state, np.random.RandomState): + return random_state + return np.random.RandomState(random_state) + + +def stratified_sample_questions( + dfq: pd.DataFrame, n_target: int, random_state: int | np.random.RandomState | None = None +) -> pd.DataFrame: """Sample questions using stratified sampling to achieve target distribution. This ensures we get the desired distribution regardless of source data skew. @@ -797,10 +820,14 @@ def stratified_sample_questions(dfq: pd.DataFrame, n_target: int) -> pd.DataFram Args: dfq (pd.DataFrame): DataFrame with bin_weight column and composite bins n_target (int): Number of questions to sample + random_state: Seed/``np.random.RandomState`` for reproducible per-bin sampling (anything + ``DataFrame.sample`` accepts). ``None`` (default) samples without a fixed seed. Thread a + single ``RandomState`` instance through to decorrelate the successive per-bin draws. Returns result (pd.DataFrame): Sampled questions """ + random_state = _as_random_state(random_state) if len(dfq) == 0 or n_target == 0: return pd.DataFrame() @@ -857,7 +884,7 @@ def stratified_sample_questions(dfq: pd.DataFrame, n_target: int) -> pd.DataFram for bin_name, n_samples in bin_samples.items(): if n_samples > 0: bin_df = dfq_weighted[dfq_weighted["composite_bin"] == bin_name] - sampled = bin_df.sample(n=n_samples, replace=False) + sampled = bin_df.sample(n=n_samples, replace=False, random_state=random_state) sampled_dfs.append(sampled) if not sampled_dfs: @@ -865,7 +892,9 @@ def stratified_sample_questions(dfq: pd.DataFrame, n_target: int) -> pd.DataFram return pd.concat(sampled_dfs, ignore_index=True) -def sample_market_questions(dfq: pd.DataFrame, n_target: int) -> pd.DataFrame: +def sample_market_questions( + dfq: pd.DataFrame, n_target: int, random_state: int | np.random.RandomState | None = None +) -> pd.DataFrame: """Sample market questions using multi-dimensional binning strategy. Ensures balanced sampling across market probability values and time horizons. @@ -873,6 +902,8 @@ def sample_market_questions(dfq: pd.DataFrame, n_target: int) -> pd.DataFrame: Args: dfq (pd.DataFrame): Market questions n_target (int): Number of questions to sample + random_state: Seed/``np.random.RandomState`` threaded to ``stratified_sample_questions`` + for reproducibility. ``None`` (default) samples without a fixed seed. Returns df_result (pd.DataFrame): Sampled questions @@ -886,6 +917,7 @@ def sample_market_questions(dfq: pd.DataFrame, n_target: int) -> pd.DataFrame: df_result = stratified_sample_questions( dfq=dfq, n_target=n_target, + random_state=random_state, ) df_result = df_result.drop( columns=[ @@ -898,7 +930,9 @@ def sample_market_questions(dfq: pd.DataFrame, n_target: int) -> pd.DataFrame: return df_result -def llm_sample_questions(values: dict, n_single: int) -> pd.DataFrame: +def llm_sample_questions( + values: dict, n_single: int, random_state: int | np.random.RandomState | None = None +) -> pd.DataFrame: """Generate questions for the LLM question set. For market questions: Sample using binning strategy. @@ -907,23 +941,26 @@ def llm_sample_questions(values: dict, n_single: int) -> pd.DataFrame: Args: values (dict): Source data dict containing "dfq" DataFrame n_single (int): Number of questions to sample + random_state: Seed/``np.random.RandomState`` threaded to the market/category samplers for + reproducibility. ``None`` (default) samples without a fixed seed. Returns df (pd.DataFrame): Sampled questions """ dfq = values["dfq"].copy() source = dfq["source"].iloc[0] + random_state = _as_random_state(random_state) if source in question_curation.MARKET_SOURCES: # Use binning-based sampling for market questions - return sample_market_questions(dfq, n_single) + return sample_market_questions(dfq, n_single, random_state=random_state) else: # Use existing category-based sampling for data sources allocation = allocate_across_categories(num_questions=n_single, dfq=dfq) dfs = [] for key, value in allocation.items(): - dfs.append(dfq[dfq["category"] == key].sample(value)) + dfs.append(dfq[dfq["category"] == key].sample(value, random_state=random_state)) return pd.concat(dfs, ignore_index=True) @@ -1252,6 +1289,12 @@ def driver(_: None) -> None: ) HUMAN_QUESTIONS.update(human_questions_of_question_type) + # For testing/reproducibility only: when `env.RANDOM_SEED` is set, thread a single RandomState + # through both sampling passes so the sampled set is deterministic. It is unset in deployment, + # which preserves the historical unseeded behaviour. + seed = env.RANDOM_SEED + random_state = None if seed is None else np.random.RandomState(seed) + # Sample questions logger.info("LLM SET") LLM_QUESTIONS = process_questions( @@ -1260,6 +1303,7 @@ def driver(_: None) -> None: single_generation_func=llm_sample_questions, show_plots=env.RUNNING_LOCALLY, question_set_target=QuestionSetTarget.LLM, + random_state=random_state, ) logger.info("HUMAN SET") @@ -1269,6 +1313,7 @@ def driver(_: None) -> None: single_generation_func=human_sample_questions, show_plots=False, question_set_target=QuestionSetTarget.HUMAN, + random_state=random_state, ) write_questions(LLM_QUESTIONS, question_set_target=QuestionSetTarget.LLM) diff --git a/src/helpers/env.py b/src/helpers/env.py index 6cd7389b..9c0202a0 100644 --- a/src/helpers/env.py +++ b/src/helpers/env.py @@ -38,10 +38,14 @@ def __getattr__(name): return bool(int(os.environ.get("RUNNING_LOCALLY", False))) if name == "BUCKET_MOUNT_POINT": return os.environ.get("BUCKET_MOUNT_POINT", "") + if name == "RANDOM_SEED": + # Seed for sources of non-determinism for testing/reproducibility; must be unset/None in prod + value = os.environ.get("RANDOM_SEED") + return int(value) if value else None raise AttributeError(f"module {__name__!r} has no attribute {name!r}") def __dir__(): """Expose the lazily-read environment variable names to ``dir()``/autocomplete.""" - extra = {"NUM_CPUS", "RUNNING_LOCALLY", "BUCKET_MOUNT_POINT"} + extra = {"NUM_CPUS", "RUNNING_LOCALLY", "BUCKET_MOUNT_POINT", "RANDOM_SEED"} return sorted(set(globals()) | set(_STR_VARS) | extra) diff --git a/src/helpers/keys.py b/src/helpers/keys.py index 84b68468..cc815368 100644 --- a/src/helpers/keys.py +++ b/src/helpers/keys.py @@ -1,9 +1,31 @@ -"""utils for key-related tasks in llm-benchmark.""" +"""utils for key-related tasks in llm-benchmark. + +Secrets are resolved lazily on first attribute access (PEP 562 module ``__getattr__``) +and memoized. This ensures that merely *importing* this module performs no network/Secret Manager call. +""" from google.cloud import secretmanager from . import env +# Secret Manager secret names; the attribute name is the secret name. +_SECRET_NAMES = { + # QUESTION DATASET SOURCES + "API_EMAIL_ACLED", + "API_PASSWORD_ACLED", + "API_KEY_FRED", + # QUESTION MARKET SOURCES + "API_KEY_METACULUS", + "API_KEY_POLYMARKET", + # WORKFLOW BOT + "API_SLACK_BOT_NOTIFICATION", + "API_SLACK_BOT_CHANNEL", + # GITHUB + "API_GITHUB_DATASET_REPO_URL", +} + +_cache: dict = {} + def get_secret(secret_name, version_id="latest"): """ @@ -27,18 +49,15 @@ def get_secret_that_may_not_exist(secret_name, version_id="latest"): return None -# QUESTION DATASET SOURCES -API_EMAIL_ACLED = get_secret(secret_name="API_EMAIL_ACLED") -API_PASSWORD_ACLED = get_secret(secret_name="API_PASSWORD_ACLED") -API_KEY_FRED = get_secret("API_KEY_FRED") - -# QUESTION MARKET SOURCES -API_KEY_METACULUS = get_secret(secret_name="API_KEY_METACULUS") -API_KEY_POLYMARKET = get_secret("API_KEY_POLYMARKET") +def __getattr__(name): + """Lazily resolve and memoize ``API_*`` secrets on first access (PEP 562).""" + if name not in _SECRET_NAMES: + raise AttributeError(f"module {__name__!r} has no attribute {name!r}") + if name not in _cache: + _cache[name] = get_secret(name) + return _cache[name] -# WORKFLOW BOT -API_SLACK_BOT_NOTIFICATION = get_secret(secret_name="API_SLACK_BOT_NOTIFICATION") -API_SLACK_BOT_CHANNEL = get_secret(secret_name="API_SLACK_BOT_CHANNEL") -# GITHUB -API_GITHUB_DATASET_REPO_URL = get_secret(secret_name="API_GITHUB_DATASET_REPO_URL") +def __dir__(): + """Expose the lazily-resolved secret names to ``dir()``/autocomplete.""" + return sorted(set(globals()) | set(_SECRET_NAMES)) diff --git a/src/leaderboard/main.py b/src/leaderboard/main.py index bb4dbd61..bf27f58f 100644 --- a/src/leaderboard/main.py +++ b/src/leaderboard/main.py @@ -2617,6 +2617,33 @@ def score_models( return df_leaderboard, question_fixed_effects +def _question_level_bootstrap( + df: pd.DataFrame, random_state: int | np.random.RandomState | None = None +) -> pd.DataFrame: + """Resample question_pks with replacement for one bootstrap replicate of a group. + + Args: + df (pd.DataFrame): Rows for one ``(forecast_due_date, source)`` group. + random_state: Seed / ``RandomState`` for reproducible resampling. ``None`` (the default) + draws from fresh entropy, matching the historical non-deterministic behavior. + + Returns: + pd.DataFrame: Rows for the resampled questions, with ``question_pk`` made unique per draw. + """ + questions = df["question_pk"].drop_duplicates() + questions_bs = questions.sample(frac=1, replace=True, random_state=random_state) + sample = questions_bs.to_frame(name="question_pk") + sample["draw"] = sample.groupby("question_pk").cumcount() + retval = pd.merge(sample, df, on="question_pk", how="left") + # `question_pk` must be overwritten with a unique id in case it was sampled more than once. + # This ensures that `two_way_fixed_effects()` treats each drawn question separately (instead + # of treating multiple draws as one question). + retval["question_pk"] = ( + retval["question_pk"].astype(str) + "_sim_id_" + retval["draw"].astype(str) + ) + return retval.drop(columns=["draw"]) + + @decorator.log_runtime def generate_simulated_leaderboards( df: pd.DataFrame, @@ -2643,30 +2670,17 @@ def generate_simulated_leaderboards( df = df.copy() - def question_level_bootstrap(df: pd.DataFrame) -> pd.DataFrame: - questions = df["question_pk"].drop_duplicates() - questions_bs = questions.sample(frac=1, replace=True) - sample = questions_bs.to_frame(name="question_pk") - sample["draw"] = sample.groupby("question_pk").cumcount() - retval = pd.merge( - sample, - df, - on="question_pk", - how="left", - ) - # `question_pk` must be overwritten with a unique id in case it was sampled more than once. - # This ensures that `two_way_fixed_effects()` treats each drawn question separately (instead - # of treating multiple draws as one question). - retval["question_pk"] = ( - retval["question_pk"].astype(str) + "_sim_id_" + retval["draw"].astype(str) - ) - return retval.drop(columns=["draw"]) + # For testing/reproducibility only: when `env.RANDOM_SEED` is set, replicate ``i`` uses + # ``seed + i`` so each loky child process is deterministic independent of process RNG state. + # It is unset in deployment, which preserves the historical non-deterministic behaviour. + seed = env.RANDOM_SEED def bootstrap_and_score(idx): logger.info(f"[replicate {idx+1}/{N}] starting...") + random_state = None if seed is None else np.random.RandomState(seed + idx) df_bs = ( df.groupby(["forecast_due_date", "source"]) - .apply(question_level_bootstrap, include_groups=False) + .apply(_question_level_bootstrap, include_groups=False, random_state=random_state) .reset_index() ) try: diff --git a/src/tests/test_runtime_requirements.py b/src/tests/test_runtime_requirements.py index 47918bfb..1bba57e1 100644 --- a/src/tests/test_runtime_requirements.py +++ b/src/tests/test_runtime_requirements.py @@ -123,6 +123,13 @@ def test_root_makefile_exports_transcript_bucket_to_deployments(): assert "FORECAST_SETS_TRANSCRIPTS_BUCKET=$(FORECAST_SETS_TRANSCRIPTS_BUCKET)" in makefile +def test_root_makefile_does_not_ship_random_seed_to_deployments(): + """`RANDOM_SEED` is for testing/reproducibility only and must never reach a Cloud Run job.""" + makefile = (ROOT / "Makefile").read_text() + + assert "RANDOM_SEED" not in makefile + + def test_root_make_test_bootstraps_python_env_before_pytest(): result = subprocess.run( ["make", "--dry-run", "--always-make", "test", "ARGS=--version"],