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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
71 changes: 58 additions & 13 deletions src/curate_questions/create_question_set/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,14 +15,14 @@
import json
import logging
import os
import random
import sys
from collections.abc import Callable
from copy import deepcopy
from datetime import datetime, timedelta
from enum import Enum
from fractions import Fraction

import numpy as np
import pandas as pd
from tqdm import tqdm
from utils import gcp
Expand Down Expand Up @@ -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
Expand All @@ -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)

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -789,18 +796,38 @@ 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.

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()

Expand Down Expand Up @@ -857,22 +884,26 @@ 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:
raise ValueError("Stratified sampling produced no results.")
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.

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
Expand All @@ -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=[
Expand All @@ -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.
Expand All @@ -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)


Expand Down Expand Up @@ -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(
Expand All @@ -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")
Expand All @@ -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)
Expand Down
6 changes: 5 additions & 1 deletion src/helpers/env.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
47 changes: 33 additions & 14 deletions src/helpers/keys.py
Original file line number Diff line number Diff line change
@@ -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",
}
Comment thread
nikbpetrov marked this conversation as resolved.

_cache: dict = {}


def get_secret(secret_name, version_id="latest"):
"""
Expand All @@ -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))
52 changes: 33 additions & 19 deletions src/leaderboard/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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:
Expand Down
7 changes: 7 additions & 0 deletions src/tests/test_runtime_requirements.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"],
Expand Down
Loading