0% found this document useful (0 votes)
2 views158 pages

AI Code Ruled

Its a manuel

Uploaded by

George
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd
0% found this document useful (0 votes)
2 views158 pages

AI Code Ruled

Its a manuel

Uploaded by

George
Copyright
© All Rights Reserved
We take content rights seriously. If you suspect this is your content, claim it here.
Available Formats
Download as PDF, TXT or read online on Scribd

CODE FINAL MS

# =============================================================================

# MS DMT ANALYSIS (Loaded-Dataframe-Only Pipeline)

# NIH All of Us Researcher Workbench

# =============================================================================

# This pipeline uses ONLY preloaded dataframes:

# - dataset_64931192_person_df

# - dataset_64931192_survey_df

# - dataset_64931192_observation_df

# - dataset_64931192_measurement_df

# - dataset_64931192_condition_df

# - dataset_64931192_visit_df

# - dataset_64931192_drug_df

# - dataset_64931192_zip_code_socioeconomic_df

# No new SQL pulls are used in this script.

# =============================================================================
import os

import re

import numpy as np

import pandas as pd

try:

from [Link] import display, Markdown

except Exception:

def display(x):

print(x)

def Markdown(x):

return x

try:

from [Link] import chi2_contingency

SCIPY_OK = True

except Exception:

SCIPY_OK = False
# -----------------------------------------------------------------------------

# SECTION 0 - GLOBALS / HELPERS

# -----------------------------------------------------------------------------

SMALL_CELL_THRESHOLD = 20

LOOKBACK_DAYS = 365

DISCONTINUATION_GAP_DAYS = 180

SEED = 42

[Link](SEED)

OUTPUT_DIR = "outputs/ms_dmt_loaded_data_only"

[Link](OUTPUT_DIR, exist_ok=True)

def audit_df(name, x, head_n=5, top_null_n=20):

print(f"\n[{name}] shape = {[Link]}")

if isinstance(x, [Link]):

print(f"[{name}] null counts (top {top_null_n}):")

print([Link]().sum().sort_values(ascending=False).head(top_null_n))
print(f"[{name}] head({head_n}):")

display([Link](head_n))

else:

print(f"[{name}] object type: {type(x)}")

display(x)

def ensure_datetime(df, cols):

out = [Link]()

for c in cols:

if c in [Link]:

out[c] = pd.to_datetime(out[c], errors="coerce")

return out

def as_bool_mask(mask, index, name):

"""

Coerce any mask-like object into a 1D bool numpy array aligned to index.

Fills NA as False and prints compact debug diagnostics.


"""

if isinstance(mask, [Link]):

aligned = [Link](index)

na_count = int([Link]().sum())

print(f"[mask debug] {name}: Series dtype={[Link]}, shape={[Link]}, na={na_count}")

return [Link](False).to_numpy(dtype=bool)

arr = [Link](mask)

print(f"[mask debug] {name}: ndarray dtype={[Link]}, shape={[Link]}")

if [Link] == 0:

raise ValueError(f"{name} is scalar; expected 1D mask of length {len(index)}")

if [Link] != 1:

raise ValueError(f"{name} has ndim={[Link]}; expected 1D mask")

if [Link][0] != len(index):

raise ValueError(f"{name} length {[Link][0]} != len(index) {len(index)}")

if [Link] != bool:

arr = [Link](arr, index=index).fillna(False).to_numpy(dtype=bool)

return arr
def normalize_column_names(df, df_name="df"):

"""

Standardize columns to lowercase snake-like names so mixed-case

dataset-builder outputs (e.g., PERSON_ID) are handled consistently.

"""

out = [Link]()

original_cols = [str(c) for c in [Link]]

normalized = [[Link](r" ", "_", str(c).strip().lower()) for c in original_cols]

seen = {}

deduped = []

for c in normalized:

if c not in seen:

seen[c] = 1

[Link](c)

else:

seen[c] += 1

[Link](f"{c}__dup{seen[c]}")
[Link] = deduped

if original_cols != deduped:

print(f"[normalize_column_names] {df_name}: normalized column names")

print(" sample mapping:")

for o, n in list(zip(original_cols, deduped))[:12]:

if o != n:

print(f" {o} -> {n}")

if len(original_cols) > 12:

print(" ...")

return out

def exact_median(x):

s = pd.to_numeric(x, errors="coerce").dropna()

return [Link] if [Link] else float([Link]())

def mask_small_cells(df_in, n_col_candidates=None, threshold=SMALL_CELL_THRESHOLD):


out = df_in.copy().astype(object)

if n_col_candidates is None:

n_col_candidates = ["n", "n_group", "n_total", "n_initiators"]

n_col = None

for c in n_col_candidates:

if c in [Link]:

n_col = c

break

if n_col is None:

return out

n_vals = pd.to_numeric(out[n_col], errors="coerce")

small_mask = n_vals < threshold

keep_cols = {n_col}

for c in [Link]:

if [Link]().startswith("denominator"):

keep_cols.add(c)
if c in {"sex_clean", "race_eth_harmonized", "race_3group", "ethnicity_2group", "first_dmt_group", "delay_category", "index_era", "group_type", "group_valu

e", "intersection_group"}:

keep_cols.add(c)

for c in [Link]:

if c in keep_cols:

continue

numeric_c = pd.to_numeric(out[c], errors="coerce")

if numeric_c.notna().any():

[Link][small_mask, c] = f"<{threshold}"

out["suppression_note"] = [Link](small_mask, f"Masked: {n_col} < {threshold}", "")

return out

def suppress_crosstab_cells(df_in, id_cols, threshold=SMALL_CELL_THRESHOLD):

out = df_in.copy().astype(object)

for c in [Link]:

if c in id_cols:
continue

numeric_c = pd.to_numeric(out[c], errors="coerce")

mask = numeric_c.notna() & (numeric_c < threshold)

[Link][mask, c] = f"<{threshold}"

return out

def print_restriction(label, n_before, n_after):

print(f"[{label}] N before={n_before:,} | N after={n_after:,} | N excluded={n_before - n_after:,}")

def categorize_delay(days_to_dmt, initiated):

if initiated != 1 or [Link](days_to_dmt):

return "Never initiated"

if days_to_dmt <= 90:

return "<=90 days"

if days_to_dmt <= 365:

return "91-365 days"

return ">365 days"


def era_from_year(y):

if [Link](y):

return "Unknown"

y = int(y)

if y <= 2009:

return "2000-2009"

if y <= 2016:

return "2010-2016"

return "2017+"

def print_issue_cause_fix(items, title="ROOT-CAUSE PASS"):

"""

Print compact diagnostics in required format: Issue -> Cause -> Fix.

"""

print(f"\n=== {title} ===")

for i, item in enumerate(items, 1):


issue = [Link]("issue", "unspecified issue")

cause = [Link]("cause", "unspecified cause")

fix = [Link]("fix", "unspecified fix")

print(f"{i}. Issue: {issue}")

print(f" Cause: {cause}")

print(f" Fix: {fix}")

def make_gate(check_name, passed, detail):

return {

"gate": check_name,

"status": "PASS" if bool(passed) else "FAIL",

"detail": str(detail),

def require_columns(df, required_cols, df_name):

missing = [c for c in required_cols if c not in [Link]]

if missing:
raise KeyError(f"{df_name} is missing required columns: {missing}")

return True

# -----------------------------------------------------------------------------

# SECTION 1 - LOAD AND AUDIT ALREADY-LOADED DATAFRAMES

# -----------------------------------------------------------------------------

required_frames = {

"person_df": "dataset_64931192_person_df",

"survey_df": "dataset_64931192_survey_df",

"observation_df": "dataset_64931192_observation_df",

"measurement_df": "dataset_64931192_measurement_df",

"condition_df": "dataset_64931192_condition_df",

"visit_df": "dataset_64931192_visit_df",

"drug_df": "dataset_64931192_drug_df",

"zip_ses_df": "dataset_64931192_zip_code_socioeconomic_df",

}
missing_frames = [gname for gname in required_frames.values() if gname not in globals()]

if missing_frames:

raise ValueError(

"Missing required preloaded dataframes. Please run the dataset-builder cells first: "

+ ", ".join(missing_frames)

person_df = globals()[required_frames["person_df"]].copy()

survey_df = globals()[required_frames["survey_df"]].copy()

observation_df = globals()[required_frames["observation_df"]].copy()

measurement_df = globals()[required_frames["measurement_df"]].copy()

condition_df = globals()[required_frames["condition_df"]].copy()

visit_df = globals()[required_frames["visit_df"]].copy()

drug_df = globals()[required_frames["drug_df"]].copy()

zip_ses_df = globals()[required_frames["zip_ses_df"]].copy()

# Normalize column names first (handles mixed-case exports like PERSON_ID).

person_df = normalize_column_names(person_df, "person_df")

survey_df = normalize_column_names(survey_df, "survey_df")


observation_df = normalize_column_names(observation_df, "observation_df")

measurement_df = normalize_column_names(measurement_df, "measurement_df")

condition_df = normalize_column_names(condition_df, "condition_df")

visit_df = normalize_column_names(visit_df, "visit_df")

drug_df = normalize_column_names(drug_df, "drug_df")

zip_ses_df = normalize_column_names(zip_ses_df, "zip_ses_df")

person_df = ensure_datetime(person_df, ["date_of_birth"])

survey_df = ensure_datetime(survey_df, ["survey_datetime"])

observation_df = ensure_datetime(observation_df, ["observation_datetime"])

measurement_df = ensure_datetime(measurement_df, ["measurement_datetime"])

condition_df = ensure_datetime(condition_df, ["condition_start_datetime", "condition_end_datetime"])

visit_df = ensure_datetime(visit_df, ["visit_start_datetime", "visit_end_datetime"])

drug_df = ensure_datetime(drug_df, ["drug_exposure_start_datetime", "drug_exposure_end_datetime", "verbatim_end_date"])

zip_ses_df = ensure_datetime(zip_ses_df, ["observation_datetime"])

print("[columns] person_df:", person_df.[Link]()[:20])

print("[columns] visit_df:", visit_df.[Link]()[:20])

print("[columns] drug_df:", drug_df.[Link]()[:20])


# Root-cause diagnostics (explicit Issue -> Cause -> Fix).

root_cause_items = [

"issue": "Potential DMT row inflation during concept mapping merge.",

"cause": "Non-canonical concept table can create many-to-many merges.",

"fix": "Enforce one-row-per-drug_concept_id mapping and merge with validate='many_to_one'.",

},

"issue": "Potential exposure duplication distorting sequence/history analyses.",

"cause": "Exact duplicate rows or duplicate (person_id, drug_concept_id, drug_date) records.",

"fix": "Detect and remove exact duplicates and key-level duplicates with before/after counts.",

},

"issue": "Baseline/index inconsistency in MS-only extracts.",

"cause": "Applying lookback exclusion while index is earliest observed activity and using mixed baselines across sections.",

"fix": "Retain lookback only as QC; no eligibility exclusion. Use index_date for whole-cohort analyses and first_dmt_date for initiator-only treatment-path

way analyses.",

},
{

"issue": "Silent empty utilization/comorbidity outputs.",

"cause": "Window/index mismatches can filter all rows without hard failure.",

"fix": "Add explicit candidate/window counts and fail loudly when emptiness is inconsistent with available data.",

},

"issue": "Ambiguous high-efficacy inference columns in a descriptive-only output.",

"cause": "Residual log-rank columns can be interpreted as inferential results when not intended.",

"fix": "Keep high_efficacy_summary descriptive-only and enforce a non-empty summary gate.",

},

print_issue_cause_fix(root_cause_items)

# Mandatory audits

for _name, _df in [

("person_df", person_df),

("survey_df", survey_df),

("observation_df", observation_df),

("measurement_df", measurement_df),
("condition_df", condition_df),

("visit_df", visit_df),

("drug_df", drug_df),

("zip_ses_df", zip_ses_df),

]:

audit_df(_name, _df)

# -----------------------------------------------------------------------------

# SECTION 2 - CLEAN DMT MAPPING TABLE (from loaded drug_df only)

# -----------------------------------------------------------------------------

# Candidate ingredient rules for true MS DMTs.

DMT_RULES = [

"ingredient_name": "interferon beta-1a",

"keywords": ["interferon beta-1a"],

"dmt_class": "Injectable",

"efficacy_tier": "Moderate",
"route_group": "Injectable",

},

"ingredient_name": "interferon beta-1b",

"keywords": ["interferon beta-1b"],

"dmt_class": "Injectable",

"efficacy_tier": "Moderate",

"route_group": "Injectable",

},

"ingredient_name": "peginterferon beta-1a",

"keywords": ["peginterferon beta-1a"],

"dmt_class": "Injectable",

"efficacy_tier": "Moderate",

"route_group": "Injectable",

},

"ingredient_name": "glatiramer",

"keywords": ["glatiramer"],
"dmt_class": "Injectable",

"efficacy_tier": "Moderate",

"route_group": "Injectable",

},

"ingredient_name": "fingolimod",

"keywords": ["fingolimod"],

"dmt_class": "Oral moderate-efficacy",

"efficacy_tier": "Moderate",

"route_group": "Oral",

},

"ingredient_name": "siponimod",

"keywords": ["siponimod"],

"dmt_class": "Oral high-efficacy",

"efficacy_tier": "High",

"route_group": "Oral",

},

{
"ingredient_name": "ozanimod",

"keywords": ["ozanimod"],

"dmt_class": "Oral high-efficacy",

"efficacy_tier": "High",

"route_group": "Oral",

},

"ingredient_name": "ponesimod",

"keywords": ["ponesimod"],

"dmt_class": "Oral high-efficacy",

"efficacy_tier": "High",

"route_group": "Oral",

},

"ingredient_name": "dimethyl fumarate",

"keywords": ["dimethyl fumarate"],

"dmt_class": "Oral moderate-efficacy",

"efficacy_tier": "Moderate",

"route_group": "Oral",
},

"ingredient_name": "diroximel fumarate",

"keywords": ["diroximel fumarate"],

"dmt_class": "Oral moderate-efficacy",

"efficacy_tier": "Moderate",

"route_group": "Oral",

},

"ingredient_name": "monomethyl fumarate",

"keywords": ["monomethyl fumarate"],

"dmt_class": "Oral moderate-efficacy",

"efficacy_tier": "Moderate",

"route_group": "Oral",

},

"ingredient_name": "teriflunomide",

"keywords": ["teriflunomide"],

"dmt_class": "Oral moderate-efficacy",


"efficacy_tier": "Moderate",

"route_group": "Oral",

},

"ingredient_name": "cladribine",

"keywords": ["cladribine"],

"dmt_class": "Oral high-efficacy",

"efficacy_tier": "High",

"route_group": "Oral",

},

"ingredient_name": "natalizumab",

"keywords": ["natalizumab"],

"dmt_class": "Infused high-efficacy",

"efficacy_tier": "High",

"route_group": "Infusion",

},

"ingredient_name": "ocrelizumab",
"keywords": ["ocrelizumab"],

"dmt_class": "Infused high-efficacy",

"efficacy_tier": "High",

"route_group": "Infusion",

},

"ingredient_name": "ofatumumab",

"keywords": ["ofatumumab"],

"dmt_class": "Infused high-efficacy",

"efficacy_tier": "High",

"route_group": "Injectable",

},

"ingredient_name": "alemtuzumab",

"keywords": ["alemtuzumab"],

"dmt_class": "Infused high-efficacy",

"efficacy_tier": "High",

"route_group": "Infusion",

},
{

"ingredient_name": "rituximab",

"keywords": ["rituximab"],

"dmt_class": "Off-label historical",

"efficacy_tier": "Off-label",

"route_group": "Infusion",

},

SUSPICIOUS_EXCLUDE = {

"ibuprofen": "Ibuprofen is not an MS DMT",

"azacitidine": "Azacitidine reviewed and excluded (not an MS DMT)",

def map_dmt_concept(concept_name, route_name):

txt = str(concept_name).lower() if [Link](concept_name) else ""


# Explicit suspicious non-DMT exclusions

for k, reason in SUSPICIOUS_EXCLUDE.items():

if k in txt:

return {

"ingredient_name": [Link],

"dmt_class": [Link],

"efficacy_tier": [Link],

"route_group": [Link],

"include_flag": 0,

"exclusion_reason": reason,

for rule in DMT_RULES:

if any(kw in txt for kw in rule["keywords"]):

return {

"ingredient_name": rule["ingredient_name"],

"dmt_class": rule["dmt_class"],

"efficacy_tier": rule["efficacy_tier"],

"route_group": rule["route_group"],
"include_flag": 1,

"exclusion_reason": "",

# If no explicit match, do not include to avoid concept contamination.

fallback_reason = "Not mapped to recognized MS DMT ingredient"

if [Link](route_name):

route_lc = str(route_name).lower()

if "oral" in route_lc or "tablet" in route_lc or "capsule" in route_lc:

fallback_reason += " (oral route alone is insufficient for inclusion)"

return {

"ingredient_name": [Link],

"dmt_class": [Link],

"efficacy_tier": [Link],

"route_group": [Link],

"include_flag": 0,

"exclusion_reason": fallback_reason,

}
def mode_or_first_nonnull(series):

s = [Link]()

if [Link]:

return [Link]

m = [Link]()

if not [Link]:

return [Link][0]

return [Link][0]

# Canonical concept table: exactly one row per drug_concept_id.

concept_base = (

drug_df.groupby("drug_concept_id", as_index=False)

.agg(

standard_concept_name=("standard_concept_name", mode_or_first_nonnull),

route_concept_name_mode=("route_concept_name", mode_or_first_nonnull),

.copy()
)

concept_base["route_for_mapping"] = concept_base["route_concept_name_mode"]

mapped_rows = concept_base.apply(

lambda r: map_dmt_concept(r["standard_concept_name"], r["route_for_mapping"]),

axis=1,

result_type="expand",

cleaned_dmt_mapping = [Link]([concept_base, mapped_rows], axis=1)

cleaned_dmt_mapping = cleaned_dmt_mapping[

"drug_concept_id",

"standard_concept_name",

"ingredient_name",

"dmt_class",

"efficacy_tier",

"route_group",

"include_flag",
"exclusion_reason",

"route_concept_name_mode",

].copy()

cleaned_dmt_mapping = cleaned_dmt_mapping.sort_values(["include_flag", "dmt_class", "ingredient_name", "drug_concept_id"], ascending=[False, True, True, True]).res

et_index(drop=True)

# Enforce one canonical mapping row per concept_id.

cleaned_dmt_mapping = cleaned_dmt_mapping.drop_duplicates(subset=["drug_concept_id"], keep="first").reset_index(drop=True)

assert cleaned_dmt_mapping["drug_concept_id"].is_unique, "cleaned_dmt_mapping must be unique by drug_concept_id"

dmt_mapping_unique_ok = bool(cleaned_dmt_mapping["drug_concept_id"].is_unique)

original_unique_concepts_n = int(drug_df["drug_concept_id"].nunique())

cleaned_included_n = int(cleaned_dmt_mapping.loc[cleaned_dmt_mapping["include_flag"] == 1, "drug_concept_id"].nunique())

excluded_dmt_concepts = cleaned_dmt_mapping[cleaned_dmt_mapping["include_flag"] == 0].copy()

print("\n=== STEP 2: DMT CONCEPT MAPPING AUDIT ===")

print(f"Original unique drug concept count: {original_unique_concepts_n:,}")


print(f"Cleaned included concept count: {cleaned_included_n:,}")

print(f"Rows in cleaned mapping (must equal unique concepts): {len(cleaned_dmt_mapping):,}")

audit_df("cleaned_dmt_mapping", cleaned_dmt_mapping)

audit_df("excluded_dmt_concepts", excluded_dmt_concepts)

# Map exposures using cleaned table only.

drug_mapped = drug_df.merge(

cleaned_dmt_mapping[

"drug_concept_id",

"ingredient_name",

"dmt_class",

"efficacy_tier",

"route_group",

"include_flag",

"exclusion_reason",

],

on="drug_concept_id",
how="left",

validate="many_to_one",

drug_mapped = ensure_datetime(drug_mapped, ["drug_exposure_start_datetime", "drug_exposure_end_datetime", "verbatim_end_date"])

assert len(drug_mapped) == len(drug_df), "drug_mapped row count should match drug_df row count"

drug_merge_no_inflation_ok = bool(len(drug_mapped) == len(drug_df))

print(f"Rows in drug_df: {len(drug_df):,}")

print(f"Rows in drug_mapped: {len(drug_mapped):,}")

audit_df("drug_mapped", drug_mapped)

included_exposures = drug_mapped[drug_mapped["include_flag"] == 1].copy()

included_exposures = ensure_datetime(

included_exposures,

["drug_exposure_start_datetime", "drug_exposure_end_datetime", "verbatim_end_date"],

included_exposures["drug_date"] = pd.to_datetime(

included_exposures["drug_exposure_start_datetime"], errors="coerce"

).[Link]()

included_exposures["drug_end_date"] = pd.to_datetime(
included_exposures["drug_exposure_end_datetime"], errors="coerce"

).[Link]()

# Exposure deduplication integrity checks.

included_rows_before_dedup = int(len(included_exposures))

exact_dup_mask = included_exposures.duplicated(keep="first")

exact_dupe_rows_n = int(exact_dup_mask.sum())

if exact_dupe_rows_n > 0:

included_exposures = included_exposures.loc[~exact_dup_mask].copy()

required_key_cols = ["person_id", "drug_concept_id", "drug_date"]

require_columns(included_exposures, required_key_cols, "included_exposures")

included_exposures = included_exposures.sort_values(

["person_id", "drug_concept_id", "drug_date", "drug_exposure_start_datetime"]

).copy()

key_dup_mask = included_exposures.duplicated(subset=required_key_cols, keep="first")

key_dupe_rows_n = int(key_dup_mask.sum())

if key_dupe_rows_n > 0:

included_exposures = included_exposures.loc[~key_dup_mask].copy()
included_rows_after_dedup = int(len(included_exposures))

print(

"[dmt dedup] included_exposures before="

f"{included_rows_before_dedup:,} | exact_dupes_removed={exact_dupe_rows_n:,} | "

f"key_dupes_removed={key_dupe_rows_n:,} | after={included_rows_after_dedup:,}"

assert included_rows_after_dedup <= included_rows_before_dedup, "Dedup should never increase included exposure rows"

dmt_exposure_dedup_summary = [Link](

{"metric": "included_rows_before_dedup", "value": included_rows_before_dedup},

{"metric": "exact_duplicate_rows_removed", "value": exact_dupe_rows_n},

{"metric": "key_duplicate_rows_removed", "value": key_dupe_rows_n},

{"metric": "included_rows_after_dedup", "value": included_rows_after_dedup},

audit_df("dmt_exposure_dedup_summary", dmt_exposure_dedup_summary)
# Person-level exposure indicator: any included DMT exposure marks person as exposed.

ever_dmt_exposed = (

included_exposures.groupby("person_id", as_index=False)

.size()

.rename(columns={"size": "n_included_dmt_exposures"})

ever_dmt_exposed["ever_dmt_exposed_flag"] = 1

audit_df("ever_dmt_exposed", ever_dmt_exposed)

dmt_class_freq = (

included_exposures["dmt_class"]

.fillna("Missing")

.value_counts(dropna=False)

.rename_axis("dmt_class")

.reset_index(name="n_exposures")

audit_df("dmt_class_freq", dmt_class_freq)

cleaned_dmt_mapping.to_csv(f"{OUTPUT_DIR}/cleaned_dmt_mapping.csv", index=False)
# -----------------------------------------------------------------------------

# SECTION 3 - REBUILD HARMONIZED MS COHORT

# -----------------------------------------------------------------------------

print("\n=== STEP 3: COHORT REBUILD ===")

# ============================================================

# STEP 3: COHORT REBUILD

# ASSUMPTION: the supplied extract is already restricted to

# confirmed MS patients only.

# Therefore, all persons in person_df are considered to have MS.

# We do NOT require re-identification of MS diagnosis rows from

# condition_df for cohort eligibility in this extract.

# ============================================================

EXTRACT_IS_MS_ONLY = True
if EXTRACT_IS_MS_ONLY:

print("[INFO] Extract declared as MS-only by study design.")

print("[INFO] All persons in person_df are treated as confirmed MS cases.")

# Build an index date proxy from earliest available clinical activity.

activity_parts = []

if "condition_start_datetime" in condition_df.columns and not condition_df.empty:

tmp = condition_df[["person_id", "condition_start_datetime"]].copy()

tmp = [Link](columns={"condition_start_datetime": "activity_datetime"})

activity_parts.append(tmp)

if "visit_start_datetime" in visit_df.columns and not visit_df.empty:

tmp = visit_df[["person_id", "visit_start_datetime"]].copy()

tmp = [Link](columns={"visit_start_datetime": "activity_datetime"})

activity_parts.append(tmp)

if "drug_exposure_start_datetime" in drug_df.columns and not drug_df.empty:

tmp = drug_df[["person_id", "drug_exposure_start_datetime"]].copy()


tmp = [Link](columns={"drug_exposure_start_datetime": "activity_datetime"})

activity_parts.append(tmp)

if "measurement_datetime" in measurement_df.columns and not measurement_df.empty:

tmp = measurement_df[["person_id", "measurement_datetime"]].copy()

tmp = [Link](columns={"measurement_datetime": "activity_datetime"})

activity_parts.append(tmp)

if "observation_datetime" in observation_df.columns and not observation_df.empty:

tmp = observation_df[["person_id", "observation_datetime"]].copy()

tmp = [Link](columns={"observation_datetime": "activity_datetime"})

activity_parts.append(tmp)

if len(activity_parts) == 0:

raise ValueError("No activity tables available to derive index_date for MS-only extract.")

activity_df_ms_only = [Link](activity_parts, ignore_index=True)

activity_df_ms_only["activity_datetime"] = pd.to_datetime(activity_df_ms_only["activity_datetime"], errors="coerce")

activity_df_ms_only = activity_df_ms_only.dropna(subset=["person_id", "activity_datetime"]).copy()


audit_df("activity_df_ms_only", activity_df_ms_only)

index_df = (

activity_df_ms_only.groupby("person_id", as_index=False)["activity_datetime"]

.min()

.rename(columns={"activity_datetime": "index_date"})

index_df["last_ms_dx_date"] = index_df["index_date"]

index_df["n_ms_dx_records"] = [Link]

index_df["n_ms_dx_dates"] = [Link]

index_df["ms_two_dx_flag"] = 1

index_df["ms_dx_verifiable_flag"] = 0

index_df["ms_index_method"] = "ms_only_extract_earliest_activity"

# Ensure all patients in person_df are retained.

index_df = person_df[["person_id"]].merge(index_df, on="person_id", how="left")

if index_df["index_date"].isna().any():

missing_n = int(index_df["index_date"].isna().sum())

raise ValueError(f"{missing_n} MS-only cohort patients have no activity-derived index_date.")


ms_condition_concept_audit = [Link](

columns=["condition_concept_id", "standard_concept_name", "source_concept_name", "condition_source_value", "source_concept_code"]

ms_dx_df = [Link](columns=[

"person_id", "condition_concept_id", "standard_concept_name",

"condition_start_datetime", "source_concept_name",

"condition_source_value", "source_concept_code", "dx_date"

])

ms_dx_verification_summary = [Link]({

"ms_dx_verifiable_flag": [0],

"n": [len(index_df)],

})

else:

# Diagnosis-based logic for non-MS-only extracts.

condition_df["condition_name_lc"] = condition_df["standard_concept_name"].astype(str).[Link]()

condition_df["source_name_lc"] = condition_df["source_concept_name"].astype(str).[Link]() if "source_concept_name" in condition_df.columns else ""

condition_df["source_value_lc"] = condition_df["condition_source_value"].astype(str).[Link]() if "condition_source_value" in condition_df.columns else ""

condition_df["source_code_lc"] = condition_df["source_concept_code"].astype(str).[Link]() if "source_concept_code" in condition_df.columns else ""


ms_regex = (

r"multiple sclerosis|disseminated sclerosis|"

r"\bg35\b|\b340\b|"

r"\bms\b"

condition_df["ms_text_blob"] = (

condition_df["condition_name_lc"].fillna("").astype(str)

+ " "

+ condition_df["source_name_lc"].fillna("").astype(str)

+ " "

+ condition_df["source_value_lc"].fillna("").astype(str)

+ " "

+ condition_df["source_code_lc"].fillna("").astype(str)

).[Link]()

condition_df["is_ms_dx"] = condition_df["ms_text_blob"].[Link](ms_regex, regex=True, na=False)

ms_condition_concept_audit = (

condition_df.loc[
condition_df["is_ms_dx"],

["condition_concept_id", "standard_concept_name", "source_concept_name", "condition_source_value", "source_concept_code"],

.drop_duplicates()

.sort_values(["standard_concept_name", "condition_concept_id"])

.reset_index(drop=True)

ms_dx_df = condition_df.loc[

condition_df["is_ms_dx"] & condition_df["condition_start_datetime"].notna(),

["person_id", "condition_concept_id", "standard_concept_name", "condition_start_datetime", "source_concept_name", "condition_source_value", "source_concept

_code"],

].copy()

ms_dx_df["dx_date"] = pd.to_datetime(ms_dx_df["condition_start_datetime"], errors="coerce").[Link]()

index_df = (

ms_dx_df.groupby("person_id", as_index=False)

.agg(

index_date=("dx_date", "min"),
last_ms_dx_date=("dx_date", "max"),

n_ms_dx_records=("dx_date", "size"),

n_ms_dx_dates=("dx_date", "nunique"),

index_df["ms_two_dx_flag"] = (index_df["n_ms_dx_records"] >= 2).astype(int)

index_df["ms_dx_verifiable_flag"] = 1

index_df["ms_index_method"] = "ms_condition_rows"

ms_dx_verification_summary = (

index_df["ms_dx_verifiable_flag"]

.value_counts(dropna=False)

.rename_axis("ms_dx_verifiable_flag")

.reset_index(name="n")

audit_df("ms_condition_concept_audit", ms_condition_concept_audit)

audit_df("ms_dx_df", ms_dx_df)

audit_df("index_df", index_df)

audit_df("ms_dx_verification_summary", ms_dx_verification_summary)
# Start flow from all persons with >=1 MS diagnosis row.

cohort_flow = []

cohort0 = index_df.merge(

person_df[[

"person_id",

"date_of_birth",

"gender",

"sex_at_birth",

"race",

"ethnicity",

"self_reported_category",

]],

on="person_id",

how="left",

cohort0 = ensure_datetime(cohort0, ["index_date", "date_of_birth"])

cohort0["age_at_index"] = (
cohort0["index_date"].[Link]

- cohort0["date_of_birth"].[Link]

- (

(cohort0["index_date"].[Link] < cohort0["date_of_birth"].[Link])

| (

(cohort0["index_date"].[Link] == cohort0["date_of_birth"].[Link])

& (cohort0["index_date"].[Link] < cohort0["date_of_birth"].[Link])

).astype(float)

audit_df("cohort0", cohort0)

_start_rule = "MS-only extract persons from person_df" if EXTRACT_IS_MS_ONLY else ">=1 MS diagnosis"

cohort_flow.append({"step": "start_ms_population", "n_before": len(cohort0), "n_after": len(cohort0), "n_excluded": 0, "rule": _start_rule})

# Age >=18

n_before = len(cohort0)

cohort1 = cohort0[cohort0["age_at_index"] >= 18].copy()

print_restriction("age>=18", n_before, len(cohort1))


cohort_flow.append({"step": "age>=18", "n_before": n_before, "n_after": len(cohort1), "n_excluded": n_before - len(cohort1), "rule": "age_at_index >= 18"})

audit_df("cohort1_age_filtered", cohort1)

# >=2 MS diagnosis records (treated as already satisfied in MS-only extract mode)

n_before = len(cohort1)

cohort2 = cohort1[cohort1["ms_two_dx_flag"] == 1].copy()

print_restriction(">=2 MS dx", n_before, len(cohort2))

_two_dx_rule = ">=2 MS diagnosis records" if not EXTRACT_IS_MS_ONLY else "MS-only extract assumption (not re-verified from condition rows)"

cohort_flow.append({"step": "ms_two_dx", "n_before": n_before, "n_after": len(cohort2), "n_excluded": n_before - len(cohort2), "rule": _two_dx_rule})

audit_df("cohort2_two_dx", cohort2)

# Build activity table for lookback and censoring.

activity_parts = []

if not visit_df.empty:

activity_parts.append(

visit_df[["person_id", "visit_start_datetime"]]

.rename(columns={"visit_start_datetime": "activity_datetime"})

if not condition_df.empty:
activity_parts.append(

condition_df[["person_id", "condition_start_datetime"]]

.rename(columns={"condition_start_datetime": "activity_datetime"})

if not observation_df.empty:

activity_parts.append(

observation_df[["person_id", "observation_datetime"]]

.rename(columns={"observation_datetime": "activity_datetime"})

if not measurement_df.empty:

activity_parts.append(

measurement_df[["person_id", "measurement_datetime"]]

.rename(columns={"measurement_datetime": "activity_datetime"})

if not drug_mapped.empty:

activity_parts.append(

drug_mapped[["person_id", "drug_exposure_start_datetime"]]

.rename(columns={"drug_exposure_start_datetime": "activity_datetime"})

)
activity_df = [Link](activity_parts, ignore_index=True)

activity_df["activity_date"] = pd.to_datetime(activity_df["activity_datetime"], errors="coerce").[Link]()

activity_df = activity_df.dropna(subset=["activity_date"]).copy()

audit_df("activity_df", activity_df)

activity_bounds = (

activity_df.groupby("person_id", as_index=False)["activity_date"]

.agg(first_activity_date="min", last_activity_date="max")

audit_df("activity_bounds", activity_bounds)

lookback_check = cohort2[["person_id", "index_date"]].merge(activity_df[["person_id", "activity_date"]], on="person_id", how="left")

lookback_check["lookback_hit"] = (

(lookback_check["activity_date"] < lookback_check["index_date"])

& (lookback_check["activity_date"] >= lookback_check["index_date"] - pd.to_timedelta(LOOKBACK_DAYS, unit="D"))

lookback_flag = (

lookback_check.groupby("person_id", as_index=False)["lookback_hit"]
.max()

.rename(columns={"lookback_hit": "lookback_flag"})

lookback_flag["lookback_flag"] = lookback_flag["lookback_flag"].fillna(False).astype(int)

audit_df("lookback_flag", lookback_flag)

# Keep lookback as QC only; do not restrict cohort eligibility.

cohort3 = [Link](lookback_flag, on="person_id", how="left")

cohort3["lookback_flag"] = cohort3["lookback_flag"].fillna(0).astype(int)

n_before = len(cohort3)

cohort4 = [Link]()

print_restriction("lookback_not_applied", n_before, len(cohort4))

cohort_flow.append({

"step": "lookback_removed_for_ms_only",

"n_before": n_before,

"n_after": len(cohort4),

"n_excluded": 0,

"rule": "365-day lookback computed for QC only; no exclusion applied",

})
audit_df("cohort4_no_lookback_exclusion", cohort4)

# Restrict DMT exposures to cleaned mapping included concepts only.

included_exposures = included_exposures[included_exposures["drug_exposure_start_datetime"].notna()].copy()

included_exposures["drug_date"] = pd.to_datetime(included_exposures["drug_exposure_start_datetime"], errors="coerce").[Link]()

included_exposures["drug_end_date"] = pd.to_datetime(included_exposures["drug_exposure_end_datetime"], errors="coerce").[Link]()

audit_df("included_exposures", included_exposures)

# Pre-index DMT handling:

# In this MS-only extract, pre-index DMT is not treated as a separate exclusion state.

# Same-day timestamp normalization can create false pre-index flags.

# Therefore pre_index_dmt_flag is forced to 0 for all persons and no exclusion is applied.

drug_idx = included_exposures.merge(cohort4[["person_id", "index_date"]], on="person_id", how="inner")

pre_index_dmt = cohort4[["person_id"]].drop_duplicates().copy()

pre_index_dmt["pre_index_dmt_flag"] = 0

print("[INFO] pre_index_dmt_flag forced to 0 for all persons in MS-only extract.")

print("[INFO] Reason: pre-index DMT can arise from timestamp/date normalization artifacts and is not used as an exclusion criterion.")

audit_df("pre_index_dmt", pre_index_dmt)
post_index_exposure = drug_idx[drug_idx["drug_date"] >= drug_idx["index_date"]].copy()

post_index_exposure = post_index_exposure.sort_values(["person_id", "drug_date", "drug_concept_id"]).copy()

first_post_dmt = (

post_index_exposure.drop_duplicates("person_id", keep="first")

[["person_id", "drug_date", "dmt_class", "ingredient_name", "efficacy_tier"]]

.rename(

columns={

"drug_date": "first_dmt_date",

"dmt_class": "first_dmt_group",

"ingredient_name": "first_dmt_ingredient",

"efficacy_tier": "efficacy_tier_first",

audit_df("first_post_dmt", first_post_dmt)

# Pre-index 365-day visit count

visit_pre = visit_df[["person_id", "visit_start_datetime", "standard_concept_name", "visit_concept_id"]].copy()

visit_pre["visit_date"] = pd.to_datetime(visit_pre["visit_start_datetime"], errors="coerce").[Link]()


visit_pre = visit_pre.dropna(subset=["visit_date"]).copy()

visit_pre = visit_pre.merge(cohort4[["person_id", "index_date"]], on="person_id", how="inner")

visit_pre = visit_pre[

(visit_pre["visit_date"] < visit_pre["index_date"])

& (visit_pre["visit_date"] >= visit_pre["index_date"] - pd.to_timedelta(LOOKBACK_DAYS, unit="D"))

].copy()

visit_count_pre = (

visit_pre.groupby("person_id", as_index=False)

.agg(visit_count_365d_preindex=("visit_date", "count"))

audit_df("visit_count_pre", visit_count_pre)

# Assemble primary cohort

cohort_primary = (

cohort4

.merge(activity_bounds, on="person_id", how="left")

.merge(ever_dmt_exposed[["person_id", "ever_dmt_exposed_flag"]], on="person_id", how="left")

.merge(pre_index_dmt, on="person_id", how="left")


.merge(first_post_dmt, on="person_id", how="left")

.merge(visit_count_pre, on="person_id", how="left")

cohort_primary["ever_dmt_exposed_flag"] = cohort_primary["ever_dmt_exposed_flag"].fillna(0).astype(int)

cohort_primary["pre_index_dmt_flag"] = cohort_primary["pre_index_dmt_flag"].fillna(0).astype(int)

cohort_primary["visit_count_365d_preindex"] = cohort_primary["visit_count_365d_preindex"].fillna(0).astype(int)

cohort_primary["initiated"] = cohort_primary["first_dmt_date"].notna().astype(int)

# Baseline standardization rule:

# - Whole-cohort analyses (initiators + non-initiators): use index_date.

# - Initiator-only treatment-pathway analyses: use first_dmt_date.

cohort_primary["cohort_baseline_date"] = pd.to_datetime(

cohort_primary["index_date"], errors="coerce"

).[Link]()

cohort_primary["treatment_start_date"] = pd.to_datetime(

cohort_primary["first_dmt_date"], errors="coerce"

).[Link]()

# Backward-compatible alias retained, but no longer mixed.


cohort_primary["analysis_index_date"] = cohort_primary["cohort_baseline_date"]

cohort_primary["analysis_index_source"] = "cohort_baseline_index_date"

cohort_primary["days_to_dmt"] = [Link](

cohort_primary["initiated"] == 1,

(pd.to_datetime(cohort_primary["first_dmt_date"]) - pd.to_datetime(cohort_primary["index_date"])).[Link],

[Link],

cohort_primary["followup_end_date"] = pd.to_datetime(cohort_primary["last_activity_date"], errors="coerce")

cohort_primary["followup_end_date"] = [Link](

pd.to_datetime(cohort_primary["followup_end_date"], errors="coerce") < pd.to_datetime(cohort_primary["cohort_baseline_date"], errors="coerce"),

cohort_primary["cohort_baseline_date"],

cohort_primary["followup_end_date"],

cohort_primary["followup_end_date"] = pd.to_datetime(cohort_primary["followup_end_date"], errors="coerce")

cohort_primary["time_to_event_days"] = [Link](

cohort_primary["initiated"] == 1,

cohort_primary["days_to_dmt"],

(cohort_primary["followup_end_date"] - cohort_primary["cohort_baseline_date"]).[Link],
)

cohort_primary["time_to_event_days"] = pd.to_numeric(cohort_primary["time_to_event_days"], errors="coerce").clip(lower=0)

audit_df("cohort_primary_pre_demo", cohort_primary)

# Sensitivity cohort keeps all persons; no pre-index exclusion in MS-only extract.

n_before = len(cohort_primary)

cohort_sensitivity = cohort_primary.copy()

print_restriction("sensitivity_exclude_pre_index_dmt_removed", n_before, len(cohort_sensitivity))

cohort_flow.append({

"step": "pre_index_dmt_exclusion_removed",

"n_before": n_before,

"n_after": len(cohort_sensitivity),

"n_excluded": 0,

"rule": "pre_index_dmt_flag not used for exclusion in MS-only extract",

})

audit_df("cohort_sensitivity_pre_demo", cohort_sensitivity)

cohort_flow_summary = [Link](cohort_flow)
audit_df("cohort_flow_summary", cohort_flow_summary)

cohort_flow_summary.to_csv(f"{OUTPUT_DIR}/cohort_flow_summary.csv", index=False)

# -----------------------------------------------------------------------------

# SECTION 4 - DEMOGRAPHIC HARMONIZATION

# -----------------------------------------------------------------------------

print("\n=== STEP 4: DEMOGRAPHICS HARMONIZATION ===")

# Sex

sex_map = {

"male": "Male",

"female": "Female",

def clean_sex(row):

sx = str([Link]("sex_at_birth", "")).strip().lower()

gd = str([Link]("gender", "")).strip().lower()
if sx in sex_map:

return sex_map[sx]

if gd in sex_map:

return sex_map[gd]

return "Unknown/Other"

# Race (3-group) + Ethnicity (2-group).

def map_text_to_race3(x):

t = str(x).strip().lower()

if "black" in t or "african" in t:

return "Black"

if "white" in t:

return "White"

return "Other"

def survey_race3_map(ans):

t = str(ans).strip().lower()
if "black" in t or "african" in t:

return "Black"

if "white" in t:

return "White"

return "Other"

def survey_hisp_flag(ans):

t = str(ans).strip().lower()

if "not hispanic" in t or "non-hispanic" in t:

return 0

if "hispanic" in t or "latino" in t:

return 1

return [Link]

cohort_demo = cohort_primary.copy()

cohort_demo["sex_clean"] = cohort_demo.apply(clean_sex, axis=1)

cohort_demo["race_3group"] = cohort_demo["self_reported_category"].apply(map_text_to_race3)
cohort_demo["race_3group_source"] = "self_reported_category"

# Ethnicity initialized from [Link] then refined with survey.

eth_lc = cohort_demo.get("ethnicity", [Link](index=cohort_demo.index, dtype=object)).astype(str).[Link]()

cohort_demo["ethnicity_2group"] = [Link](

eth_lc.[Link]("hispanic|latino", regex=True, na=False)

& ~eth_lc.[Link]("not hispanic|non-hispanic", regex=True, na=False),

"Hispanic",

"Non-Hispanic",

cohort_demo["ethnicity_2group_source"] = "person_ethnicity"

# Survey support for Hispanic/Latino restoration and unknown fallback.

survey_re = survey_df[["person_id", "survey_datetime", "question", "answer"]].copy()

survey_re["q_txt"] = survey_re["question"].astype(str).[Link]()

survey_re["a_txt"] = survey_re["answer"].astype(str).[Link]()

survey_re = survey_re[

survey_re["q_txt"].[Link]("race|ethnic|hispanic|latino", regex=True, na=False)

| survey_re["a_txt"].[Link]("hispanic|latino|black|african|white|asian|american indian|alaska|hawaiian|pacific|middle eastern|north african|multiracial",


regex=True, na=False)

].copy()

audit_df("survey_re_candidates", survey_re)

if not survey_re.empty:

survey_re = survey_re.merge(cohort_demo[["person_id", "index_date"]], on="person_id", how="inner")

survey_re["delta_days"] = (survey_re["survey_datetime"] - survey_re["index_date"]).[Link]

survey_re["abs_delta_days"] = survey_re["delta_days"].abs()

survey_re = survey_re.sort_values(["person_id", "abs_delta_days", "survey_datetime"]).copy()

survey_nearest = survey_re.drop_duplicates("person_id", keep="first").copy()

survey_nearest["survey_race3_candidate"] = survey_nearest["answer"].apply(survey_race3_map)

survey_nearest["survey_hispanic_flag"] = survey_nearest["answer"].apply(survey_hisp_flag)

audit_df("survey_nearest", survey_nearest)

# Optional race fallback to upgrade Other -> White/Black when survey gives direct evidence.

survey_race3_series = survey_nearest.set_index("person_id")["survey_race3_candidate"]

mask_other = cohort_demo["race_3group"].eq("Other") & cohort_demo["person_id"].isin(survey_race3_series.index)

candidate_race = cohort_demo.loc[mask_other, "person_id"].map(survey_race3_series)

mask_upgrade = mask_other & candidate_race.isin(["White", "Black"])


cohort_demo.loc[mask_upgrade, "race_3group"] = cohort_demo.loc[mask_upgrade, "person_id"].map(survey_race3_series)

cohort_demo.loc[mask_upgrade, "race_3group_source"] = "survey_race_fallback"

# Ethnicity from all available survey responses per person.

survey_hisp_by_person = (

survey_re.assign(survey_hispanic_flag=survey_re["answer"].apply(survey_hisp_flag))

.groupby("person_id")["survey_hispanic_flag"]

.apply(lambda x: 1 if (x == 1).any() else (0 if (x == 0).any() else [Link]))

hisp1_ids = survey_hisp_by_person[survey_hisp_by_person == 1].index

hisp0_ids = survey_hisp_by_person[survey_hisp_by_person == 0].index

cohort_demo.loc[cohort_demo["person_id"].isin(hisp1_ids), "ethnicity_2group"] = "Hispanic"

cohort_demo.loc[cohort_demo["person_id"].isin(hisp1_ids), "ethnicity_2group_source"] = "survey_hispanic_any"

cohort_demo.loc[

cohort_demo["person_id"].isin(hisp0_ids) & ~cohort_demo["person_id"].isin(hisp1_ids),

"ethnicity_2group"

] = "Non-Hispanic"

cohort_demo.loc[

cohort_demo["person_id"].isin(hisp0_ids) & ~cohort_demo["person_id"].isin(hisp1_ids),


"ethnicity_2group_source"

] = "survey_non_hispanic_any"

# Backward-compatible aliases retained for existing pipeline sections.

cohort_demo["race_eth_harmonized"] = cohort_demo["race_3group"]

cohort_demo["race_eth_source"] = cohort_demo["race_3group_source"]

cohort_demo["race_eth_strict"] = cohort_demo["race_3group"]

cohort_demo["unknown_demo_flag"] = (

(cohort_demo["sex_clean"] == "Unknown/Other")

).astype(int)

audit_df("cohort_demo", cohort_demo)

race_codebook_md = """

## Race/Ethnicity Codebook

Race analysis variable: `race_3group`

- `White`

- `Black`
- `Other` (anyone not classified as White or Black)

Ethnicity analysis variable: `ethnicity_2group`

- `Hispanic`

- `Non-Hispanic` (everyone not classified as Hispanic)

Assignment notes:

1. Race starts from `self_reported_category` and may use survey fallback for White/Black.

2. Ethnicity uses person `ethnicity` with survey Hispanic/non-Hispanic overrides.

"""

display(Markdown(race_codebook_md))

# Counts

race_counts = cohort_demo["race_3group"].value_counts(dropna=False).rename_axis("race_3group").reset_index(name="n")

ethnicity_counts = cohort_demo["ethnicity_2group"].value_counts(dropna=False).rename_axis("ethnicity_2group").reset_index(name="n")

sex_counts = cohort_demo["sex_clean"].value_counts(dropna=False).rename_axis("sex_clean").reset_index(name="n")

strict_counts = cohort_demo["race_eth_strict"].value_counts(dropna=False).rename_axis("race_eth_strict").reset_index(name="n")

audit_df("race_counts", race_counts)
audit_df("ethnicity_counts", ethnicity_counts)

audit_df("sex_counts", sex_counts)

audit_df("strict_counts", strict_counts)

# Attach demographics to primary and sensitivity cohorts.

cohort_primary = cohort_demo.copy()

cohort_sensitivity = cohort_sensitivity.drop(

columns=[c for c in ["sex_clean", "race_eth_harmonized", "race_eth_strict", "unknown_demo_flag", "race_eth_source", "race_3group", "race_3group_source", "ethni

city_2group", "ethnicity_2group_source"] if c in cohort_sensitivity.columns],

errors="ignore",

).merge(

cohort_demo[["person_id", "sex_clean", "race_eth_harmonized", "race_eth_strict", "unknown_demo_flag", "race_eth_source", "race_3group", "race_3group_source", "

ethnicity_2group", "ethnicity_2group_source"]],

on="person_id",

how="left",

audit_df("cohort_sensitivity_demo", cohort_sensitivity)
# -----------------------------------------------------------------------------

# SECTION 5 - PATIENT-LEVEL ANALYTIC COHORT + SES PROXY MERGE

# -----------------------------------------------------------------------------

print("\n=== STEP 5: ANALYTIC COHORT BUILD ===")

analytic_cohort = cohort_primary.copy()

analytic_cohort["index_year"] = pd.to_datetime(analytic_cohort["index_date"], errors="coerce").[Link]

analytic_cohort["index_era"] = analytic_cohort["index_year"].apply(era_from_year)

analytic_cohort["delay_category"] = analytic_cohort.apply(lambda r: categorize_delay(r["days_to_dmt"], r["initiated"]), axis=1)

# SES merge from zip_code_socioeconomic, nearest to index_date preferring on/before index.

ses_cols = [

"person_id",

"observation_datetime",

"deprivation_index",

"median_income",

"poverty",

"no_health_insurance",
"assisted_income",

"high_school_education",

ses_available = [c for c in ses_cols if c in zip_ses_df.columns]

ses_df = zip_ses_df[ses_available].copy()

for c in ["deprivation_index", "median_income", "poverty", "no_health_insurance", "assisted_income", "high_school_education"]:

if c in ses_df.columns:

ses_df[c] = pd.to_numeric(ses_df[c], errors="coerce")

ses_df = ses_df.merge(analytic_cohort[["person_id", "index_date"]], on="person_id", how="inner")

ses_df["ses_date"] = pd.to_datetime(ses_df["observation_datetime"], errors="coerce")

ses_df = ses_df.dropna(subset=["ses_date", "index_date"]).copy()

ses_df["is_after_index"] = (ses_df["ses_date"] > ses_df["index_date"]).astype(int)

ses_df["abs_days_from_index"] = (ses_df["ses_date"] - ses_df["index_date"]).abs().[Link]

ses_df = ses_df.sort_values(

["person_id", "is_after_index", "abs_days_from_index", "ses_date"],

ascending=[True, True, True, False],


)

ses_nearest = ses_df.drop_duplicates("person_id", keep="first").copy()

keep_ses_person_cols = ["person_id", "deprivation_index", "median_income", "poverty", "no_health_insurance", "assisted_income", "high_school_education", "ses_date"

keep_ses_person_cols = [c for c in keep_ses_person_cols if c in ses_nearest.columns]

ses_person = ses_nearest[keep_ses_person_cols].copy()

audit_df("ses_person", ses_person)

analytic_cohort = analytic_cohort.merge(ses_person, on="person_id", how="left")

# SES quartiles where feasible.

def safe_quartile(series):

s = pd.to_numeric(series, errors="coerce")

n_nonnull = int([Link]().sum())

if n_nonnull < 4:

return [Link]([Link], index=[Link])

try:

return [Link](s, 4, labels=["Q1", "Q2", "Q3", "Q4"], duplicates="drop")


except Exception:

return [Link]([Link], index=[Link])

for c in ["deprivation_index", "median_income", "poverty", "no_health_insurance", "assisted_income", "high_school_education"]:

if c in analytic_cohort.columns:

analytic_cohort[f"{c}_quartile"] = safe_quartile(analytic_cohort[c]).astype(object)

# Keep required analytic columns.

required_analytic_cols = [

"person_id",

"index_date",

"cohort_baseline_date",

"treatment_start_date",

"analysis_index_date",

"analysis_index_source",

"age_at_index",

"index_year",

"index_era",

"sex_clean",
"race_3group",

"ethnicity_2group",

"race_eth_harmonized",

"race_eth_strict",

"unknown_demo_flag",

"lookback_flag",

"ever_dmt_exposed_flag",

"pre_index_dmt_flag",

"initiated",

"first_dmt_date",

"first_dmt_group",

"first_dmt_ingredient",

"efficacy_tier_first",

"days_to_dmt",

"delay_category",

"followup_end_date",

"time_to_event_days",

"visit_count_365d_preindex",

"deprivation_index",
"median_income",

"poverty",

"no_health_insurance",

"assisted_income",

"high_school_education",

for c in required_analytic_cols:

if c not in analytic_cohort.columns:

analytic_cohort[c] = [Link]

analytic_cohort = analytic_cohort[required_analytic_cols + [c for c in analytic_cohort.columns if [Link]("_quartile")]].copy()

# Follow-up field QC after column pruning.

analytic_cohort["index_date"] = pd.to_datetime(analytic_cohort["index_date"], errors="coerce").[Link]()

analytic_cohort["cohort_baseline_date"] = pd.to_datetime(analytic_cohort["cohort_baseline_date"], errors="coerce").[Link]()

analytic_cohort["cohort_baseline_date"] = analytic_cohort["cohort_baseline_date"].fillna(analytic_cohort["index_date"])

analytic_cohort["treatment_start_date"] = pd.to_datetime(analytic_cohort["treatment_start_date"], errors="coerce").[Link]()

analytic_cohort["analysis_index_date"] = pd.to_datetime(analytic_cohort["analysis_index_date"], errors="coerce").[Link]()


analytic_cohort["analysis_index_date"] = analytic_cohort["analysis_index_date"].fillna(analytic_cohort["cohort_baseline_date"])

analytic_cohort["followup_end_date"] = pd.to_datetime(analytic_cohort["followup_end_date"], errors="coerce").[Link]()

analytic_cohort["followup_end_date"] = analytic_cohort["followup_end_date"].fillna(analytic_cohort["cohort_baseline_date"])

_bad_followup = (

analytic_cohort["followup_end_date"].notna()

& analytic_cohort["cohort_baseline_date"].notna()

& (analytic_cohort["followup_end_date"] < analytic_cohort["cohort_baseline_date"])

analytic_cohort.loc[_bad_followup, "followup_end_date"] = analytic_cohort.loc[_bad_followup, "cohort_baseline_date"]

followup_qc = [Link](

"metric": [

"n_rows",

"index_date_missing",

"cohort_baseline_date_missing",

"followup_end_missing",

"followup_equals_cohort_baseline",

"followup_before_index_fixed_rows",
],

"value": [

len(analytic_cohort),

int(analytic_cohort["index_date"].isna().sum()),

int(analytic_cohort["cohort_baseline_date"].isna().sum()),

int(analytic_cohort["followup_end_date"].isna().sum()),

int((analytic_cohort["followup_end_date"] == analytic_cohort["cohort_baseline_date"]).sum()),

int(_bad_followup.sum()),

],

audit_df("followup_qc", followup_qc)

assert int(followup_qc.loc[followup_qc["metric"] == "cohort_baseline_date_missing", "value"].iloc[0]) == 0, "cohort_baseline_date has missing values"

assert int(followup_qc.loc[followup_qc["metric"] == "followup_end_missing", "value"].iloc[0]) == 0, "followup_end_date has missing values"

assert int((_bad_followup).sum()) >= 0, "followup QC internal check failed"

audit_df("analytic_cohort", analytic_cohort)

# Missingness summary

missingness_summary = [Link]({
"variable": analytic_cohort.columns,

"n_total": len(analytic_cohort),

"n_missing": [int(analytic_cohort[c].isna().sum()) for c in analytic_cohort.columns],

})

missingness_summary["n_non_missing"] = missingness_summary["n_total"] - missingness_summary["n_missing"]

missingness_summary["pct_missing"] = 100 * missingness_summary["n_missing"] / missingness_summary["n_total"].replace(0, [Link])

audit_df("missingness_summary", missingness_summary)

# -----------------------------------------------------------------------------

# SECTION 6 - EXACT DESCRIPTIVE TABLES

# -----------------------------------------------------------------------------

print("\n=== STEP 6: EXACT DESCRIPTIVE TABLES ===")

def make_group_table(df_in, group_col, table_name):

d = df_in.copy()

d[group_col] = d[group_col].fillna("Missing").astype(str)
out = (

[Link](group_col, dropna=False)

.agg(

n=("person_id", "nunique"),

n_initiated=("initiated", "sum"),

n_pre_index_dmt=("pre_index_dmt_flag", "sum"),

mean_pre_index_visits=("visit_count_365d_preindex", "mean"),

.reset_index()

out["denominator_total_n"] = int(d["person_id"].nunique())

out["pct_of_total"] = 100 * out["n"] / out["denominator_total_n"].replace(0, [Link])

out["pct_initiated_within_group"] = 100 * out["n_initiated"] / out["n"].replace(0, [Link])

d_init = d[d["initiated"] == 1].copy()

dmt_stats = (

d_init.groupby(group_col, dropna=False)["days_to_dmt"]
.agg(mean_days_to_dmt="mean", median_days_to_dmt=exact_median)

.reset_index()

out = [Link](dmt_stats, on=group_col, how="left")

[Link](0, "table", table_name)

return out.sort_values(["n", group_col], ascending=[False, True]).reset_index(drop=True)

def make_first_dmt_table(df_in):

d = df_in[df_in["initiated"] == 1].copy()

d = d[d["first_dmt_group"].notna()].copy()

d["first_dmt_group"] = d["first_dmt_group"].astype(str)

out = (

[Link]("first_dmt_group", dropna=False)

.agg(

n=("person_id", "nunique"),

mean_days_to_dmt=("days_to_dmt", "mean"),

median_days_to_dmt=("days_to_dmt", exact_median),
)

.reset_index()

out["denominator_initiators_n"] = int(d["person_id"].nunique())

out["pct_of_initiators"] = 100 * out["n"] / out["denominator_initiators_n"].replace(0, [Link])

return out.sort_values(["n", "first_dmt_group"], ascending=[False, True]).reset_index(drop=True)

def make_race_dmt_crosstab(df_in):

d = df_in[(df_in["initiated"] == 1) & (df_in["first_dmt_group"].notna())].copy()

d["race_3group"] = d["race_3group"].fillna("Other").astype(str)

d["first_dmt_group"] = d["first_dmt_group"].fillna("Missing").astype(str)

ct = [Link](d["race_3group"], d["first_dmt_group"], dropna=False)

ct = ct.reset_index().rename(columns={"race_3group": "race_3group"})

ct["denominator_initiators_n"] = int(d["person_id"].nunique())

return ct
descriptive_table_sex = make_group_table(analytic_cohort, "sex_clean", "by_sex")

descriptive_table_race_eth = make_group_table(analytic_cohort, "race_3group", "by_race3")

descriptive_table_ethnicity = make_group_table(analytic_cohort, "ethnicity_2group", "by_ethnicity")

descriptive_table_first_dmt = make_first_dmt_table(analytic_cohort)

race_by_dmt_crosstab = make_race_dmt_crosstab(analytic_cohort)

# Exact global summary metrics

summary_exact = [Link]([

"total_n": int(analytic_cohort["person_id"].nunique()),

"initiated_n": int(analytic_cohort["initiated"].sum()),

"initiated_pct": 100 * analytic_cohort["initiated"].mean(),

"pre_index_dmt_n": int(analytic_cohort["pre_index_dmt_flag"].sum()),

"mean_days_to_dmt": float(analytic_cohort.loc[analytic_cohort["initiated"] == 1, "days_to_dmt"].mean()) if int(analytic_cohort["initiated"].sum()) > 0 else

[Link],

"median_days_to_dmt": exact_median(analytic_cohort.loc[analytic_cohort["initiated"] == 1, "days_to_dmt"]),

"mean_pre_index_visits": float(analytic_cohort["visit_count_365d_preindex"].mean()) if len(analytic_cohort) else [Link],

])
audit_df("summary_exact", summary_exact)

# Policy-safe display tables

display_descriptive_table_sex = mask_small_cells(descriptive_table_sex, n_col_candidates=["n"])

display_descriptive_table_race_eth = mask_small_cells(descriptive_table_race_eth, n_col_candidates=["n"])

display_descriptive_table_ethnicity = mask_small_cells(descriptive_table_ethnicity, n_col_candidates=["n"])

display_descriptive_table_first_dmt = mask_small_cells(descriptive_table_first_dmt, n_col_candidates=["n"])

display_race_by_dmt_crosstab = suppress_crosstab_cells(race_by_dmt_crosstab, id_cols=["race_3group", "denominator_initiators_n"])

print("\nDisplayed policy-safe table: sex")

display(display_descriptive_table_sex)

print("\nDisplayed policy-safe table: race (White/Black/Other)")

display(display_descriptive_table_race_eth)

print("\nDisplayed policy-safe table: ethnicity")

display(display_descriptive_table_ethnicity)

print("\nDisplayed policy-safe table: first DMT")

display(display_descriptive_table_first_dmt)

print("\nDisplayed policy-safe table: race x first DMT")

display(display_race_by_dmt_crosstab)
# Validation compare vs existing in-memory prior if available

validation_compare = []

validation_compare.append({"metric": "total_n", "current": int(summary_exact.loc[0, "total_n"])})

validation_compare.append({"metric": "initiated_n", "current": int(summary_exact.loc[0, "initiated_n"])})

validation_compare.append({"metric": "initiated_pct", "current": float(summary_exact.loc[0, "initiated_pct"])})

validation_compare.append({"metric": "pre_index_dmt_n", "current": int(summary_exact.loc[0, "pre_index_dmt_n"])})

if "cohort_raw" in globals() and isinstance(globals()["cohort_raw"], [Link]):

prev = globals()["cohort_raw"].copy()

prev_total = int(prev["person_id"].nunique()) if "person_id" in [Link] else [Link]

prev_init = int(prev["dmt_event"].sum()) if "dmt_event" in [Link] else [Link]

prev_pre = int(prev["pre_index_dmt_flag"].sum()) if "pre_index_dmt_flag" in [Link] else [Link]

for r in validation_compare:

if r["metric"] == "total_n":

r["prior_existing_analysis"] = prev_total

elif r["metric"] == "initiated_n":

r["prior_existing_analysis"] = prev_init

elif r["metric"] == "initiated_pct":


r["prior_existing_analysis"] = 100 * prev_init / prev_total if prev_total and not [Link](prev_total) and not [Link](prev_init) else [Link]

elif r["metric"] == "pre_index_dmt_n":

r["prior_existing_analysis"] = prev_pre

else:

for r in validation_compare:

r["prior_existing_analysis"] = [Link]

validation_compare_df = [Link](validation_compare)

validation_compare_df["difference"] = validation_compare_df["current"] - validation_compare_df["prior_existing_analysis"]

validation_compare_df["discrepancy_explanation"] = [Link](

validation_compare_df["prior_existing_analysis"].isna(),

"No prior in-memory cohort found for direct comparison.",

"Differences may reflect strict cleaned DMT mapping, loaded-data-only constraints, and revised demographic/SES logic.",

audit_df("validation_compare_df", validation_compare_df)

# -----------------------------------------------------------------------------

# SECTION 7 - SUPPORTED NOVEL ANALYSES


# -----------------------------------------------------------------------------

print("\n=== STEP 7: SUPPORTED NOVEL ANALYSES ===")

# 7A. Early vs delayed treatment by sex/race (3-group) and ethnicity (2-group).

delay_category_summary = (

analytic_cohort.groupby(["race_3group", "sex_clean", "delay_category"], dropna=False)

.agg(n=("person_id", "nunique"))

.reset_index()

delay_tot = (

delay_category_summary.groupby(["race_3group", "sex_clean"], dropna=False)["n"]

.sum()

.reset_index()

.rename(columns={"n": "denominator_group_n"})

delay_category_summary = delay_category_summary.merge(delay_tot, on=["race_3group", "sex_clean"], how="left")

delay_category_summary["pct_within_group"] = 100 * delay_category_summary["n"] / delay_category_summary["denominator_group_n"].replace(0, [Link])

audit_df("delay_category_summary", delay_category_summary)
delay_category_ethnicity_summary = (

analytic_cohort.groupby(["ethnicity_2group", "delay_category"], dropna=False)

.agg(n=("person_id", "nunique"))

.reset_index()

delay_eth_tot = (

delay_category_ethnicity_summary.groupby(["ethnicity_2group"], dropna=False)["n"]

.sum()

.reset_index()

.rename(columns={"n": "denominator_group_n"})

delay_category_ethnicity_summary = delay_category_ethnicity_summary.merge(delay_eth_tot, on=["ethnicity_2group"], how="left")

delay_category_ethnicity_summary["pct_within_group"] = 100 * delay_category_ethnicity_summary["n"] / delay_category_ethnicity_summary["denominator_group_n"].replac

e(0, [Link])

audit_df("delay_category_ethnicity_summary", delay_category_ethnicity_summary)

# 7B. Calendar-era analysis.

def era_summary_by(df_in, subgroup_col=None):


d = df_in.copy()

if subgroup_col is None:

d["group_type"] = "overall"

d["group_value"] = "All"

else:

d["group_type"] = subgroup_col

d["group_value"] = d[subgroup_col].fillna("Missing").astype(str)

out = (

[Link](["index_era", "group_type", "group_value"], dropna=False)

.agg(

n=("person_id", "nunique"),

initiated_n=("initiated", "sum"),

median_days_to_dmt=("days_to_dmt", exact_median),

.reset_index()

out["initiation_pct"] = 100 * out["initiated_n"] / out["n"].replace(0, [Link])


# Most frequent first DMT class among initiators in that era/group.

d_init = d[(d["initiated"] == 1) & (d["first_dmt_group"].notna())].copy()

if not d_init.empty:

mode_class = (

d_init.groupby(["index_era", "group_type", "group_value", "first_dmt_group"], dropna=False)

.agg(n_class=("person_id", "nunique"))

.reset_index()

.sort_values(["index_era", "group_type", "group_value", "n_class", "first_dmt_group"], ascending=[True, True, True, False, True])

.drop_duplicates(["index_era", "group_type", "group_value"], keep="first")

.rename(columns={"first_dmt_group": "most_common_first_dmt_group"})

[["index_era", "group_type", "group_value", "most_common_first_dmt_group"]]

out = [Link](mode_class, on=["index_era", "group_type", "group_value"], how="left")

else:

out["most_common_first_dmt_group"] = [Link]

return out

calendar_era_summary = [Link](
[

era_summary_by(analytic_cohort, None),

era_summary_by(analytic_cohort, "sex_clean"),

era_summary_by(analytic_cohort, "race_3group"),

era_summary_by(analytic_cohort, "ethnicity_2group"),

],

ignore_index=True,

audit_df("calendar_era_summary", calendar_era_summary)

# 7C. Time to high-efficacy DMT (descriptive only).

high_eff_classes = {"Oral high-efficacy", "Infused high-efficacy"}

# Whole-cohort baseline for high-efficacy uptake analyses: cohort_baseline_date (index_date).

if "followup_end_date" not in analytic_cohort.columns:

print("[WARNING] `followup_end_date` missing in analytic_cohort; reconstructing from cohort_baseline_date for this step.")

analytic_cohort["followup_end_date"] = pd.to_datetime(analytic_cohort["cohort_baseline_date"], errors="coerce")

else:

analytic_cohort["followup_end_date"] = pd.to_datetime(analytic_cohort["followup_end_date"], errors="coerce")


analytic_cohort["cohort_baseline_date"] = pd.to_datetime(analytic_cohort["cohort_baseline_date"], errors="coerce")

analytic_cohort["followup_end_date"] = analytic_cohort["followup_end_date"].fillna(analytic_cohort["cohort_baseline_date"])

analytic_cohort.loc[

analytic_cohort["followup_end_date"] < analytic_cohort["cohort_baseline_date"],

"followup_end_date",

] = analytic_cohort.loc[

analytic_cohort["followup_end_date"] < analytic_cohort["cohort_baseline_date"],

"cohort_baseline_date",

high_eff_exp = included_exposures.merge(

analytic_cohort[["person_id", "cohort_baseline_date", "followup_end_date", "race_3group", "ethnicity_2group", "sex_clean"]],

on="person_id",

how="inner",

high_eff_exp = high_eff_exp[high_eff_exp["drug_date"] >= high_eff_exp["cohort_baseline_date"]].copy()

high_eff_exp = high_eff_exp[high_eff_exp["dmt_class"].isin(high_eff_classes)].copy()

high_eff_first = (
high_eff_exp.sort_values(["person_id", "drug_date", "drug_concept_id"])

.drop_duplicates("person_id", keep="first")

[["person_id", "drug_date"]]

.rename(columns={"drug_date": "first_high_efficacy_dmt_date"})

audit_df("high_eff_first", high_eff_first)

high_efficacy_df = analytic_cohort[["person_id", "cohort_baseline_date", "followup_end_date", "race_3group", "ethnicity_2group", "sex_clean"]].merge(

high_eff_first,

on="person_id",

how="left",

high_efficacy_df["high_efficacy_initiated"] = high_efficacy_df["first_high_efficacy_dmt_date"].notna().astype(int)

high_efficacy_df["days_to_high_efficacy_dmt"] = [Link](

high_efficacy_df["high_efficacy_initiated"] == 1,

(high_efficacy_df["first_high_efficacy_dmt_date"] - high_efficacy_df["cohort_baseline_date"]).[Link],

[Link],

)
high_efficacy_df["time_to_high_eff_event_or_censor"] = [Link](

high_efficacy_df["high_efficacy_initiated"] == 1,

high_efficacy_df["days_to_high_efficacy_dmt"],

(high_efficacy_df["followup_end_date"] - high_efficacy_df["cohort_baseline_date"]).[Link],

high_efficacy_df["time_to_high_eff_event_or_censor"] = pd.to_numeric(

high_efficacy_df["time_to_high_eff_event_or_censor"], errors="coerce"

).clip(lower=0)

audit_df("high_efficacy_df", high_efficacy_df)

high_efficacy_summary = (

high_efficacy_df.groupby(["race_3group", "sex_clean"], dropna=False)

.agg(

n=("person_id", "count"),

high_efficacy_initiated_n=("high_efficacy_initiated", "sum"),

median_days_to_high_efficacy=("days_to_high_efficacy_dmt", "median"),

.reset_index()

)
high_efficacy_summary["high_efficacy_initiation_pct"] = (

100 * high_efficacy_summary["high_efficacy_initiated_n"] / high_efficacy_summary["n"].replace(0, [Link])

cols_to_drop = [

"logrank_p_race_3group",

"logrank_status_race_3group",

"logrank_p_ethnicity_2group",

"logrank_status_ethnicity_2group",

"logrank_p_sex_clean",

"logrank_status_sex_clean",

high_efficacy_summary = high_efficacy_summary.drop(

columns=[c for c in cols_to_drop if c in high_efficacy_summary.columns],

errors="ignore",

audit_df("high_efficacy_summary", high_efficacy_summary)
# 7D. Intersectional analysis (race x sex), policy-compliant collapse.

intersectional_df = analytic_cohort[["person_id", "race_3group", "sex_clean", "initiated", "delay_category", "first_dmt_group"]].copy()

intersectional_df["race_3group"] = intersectional_df["race_3group"].fillna("Other")

intersectional_df["sex_clean"] = intersectional_df["sex_clean"].fillna("Missing")

intersectional_df["intersection_group_raw"] = intersectional_df["race_3group"] + " | " + intersectional_df["sex_clean"]

group_n = intersectional_df["intersection_group_raw"].value_counts()

small_groups = group_n[group_n < SMALL_CELL_THRESHOLD].[Link]()

intersectional_df["intersection_group"] = [Link](

intersectional_df["intersection_group_raw"].isin(small_groups),

f"Suppressed (<{SMALL_CELL_THRESHOLD})",

intersectional_df["intersection_group_raw"],

intersectional_summary = (

intersectional_df.groupby(["intersection_group", "delay_category"], dropna=False)

.agg(n=("person_id", "nunique"), initiated_n=("initiated", "sum"))

.reset_index()

)
intersection_tot = (

intersectional_summary.groupby("intersection_group", dropna=False)["n"]

.sum()

.reset_index()

.rename(columns={"n": "denominator_group_n"})

intersectional_summary = intersectional_summary.merge(intersection_tot, on="intersection_group", how="left")

intersectional_summary["pct_within_group"] = 100 * intersectional_summary["n"] / intersectional_summary["denominator_group_n"].replace(0, [Link])

audit_df("intersectional_summary", intersectional_summary)

# 7E. Healthcare utilization (post-baseline 365-day)

UTIL_WINDOW_DAYS_POST = 365

visit_util = visit_df[["person_id", "visit_start_datetime", "standard_concept_name"]].copy()

visit_util["visit_date"] = pd.to_datetime(visit_util["visit_start_datetime"], errors="coerce").[Link]()

visit_util["visit_name_lc"] = visit_util["standard_concept_name"].astype(str).[Link]()

visit_util = visit_util.dropna(subset=["visit_date"]).copy()

visit_util = visit_util.merge(analytic_cohort[["person_id", "cohort_baseline_date"]], on="person_id", how="inner")


util_candidates_before_window = int(len(visit_util))

util_people_candidates_before_window = int(visit_util["person_id"].nunique())

visit_util = visit_util[

(visit_util["visit_date"] >= visit_util["cohort_baseline_date"])

& (visit_util["visit_date"] < visit_util["cohort_baseline_date"] + pd.to_timedelta(UTIL_WINDOW_DAYS_POST, unit="D"))

].copy()

util_rows_after_window = int(len(visit_util))

util_people_after_window = int(visit_util["person_id"].nunique())

print(f"[utilization validation] rows in visit_util: {len(visit_util):,}")

print(f"[utilization validation] persons in visit_util: {visit_util['person_id'].nunique():,}")

print(

"[utilization validation] candidates before window rows="

f"{util_candidates_before_window:,}, persons={util_people_candidates_before_window:,}"

if util_rows_after_window == 0:

if util_candidates_before_window == 0:

print("[utilization proof] no cohort-linked visit rows available before post-index windowing.")

else:

raise AssertionError(
"visit_util is empty after post-index window despite non-empty candidate rows. "

"Check cohort_baseline_date and window logic."

visit_util["visit_type_group"] = [Link](

visit_util["visit_name_lc"].[Link]("emergency", regex=False, na=False),

visit_util["visit_name_lc"].[Link]("inpatient|hospital", regex=True, na=False),

visit_util["visit_name_lc"].[Link]("outpatient|office|ambulatory", regex=True, na=False),

],

["ED", "Inpatient", "Outpatient"],

default="Other",

audit_df("visit_util", visit_util)

util_counts = (

visit_util.groupby(["person_id", "visit_type_group"], dropna=False)

.size()
.reset_index(name="n")

util_wide = util_counts.pivot_table(index="person_id", columns="visit_type_group", values="n", fill_value=0).reset_index()

for c in ["Outpatient", "Inpatient", "ED", "Other"]:

if c not in util_wide.columns:

util_wide[c] = 0

util_wide = util_wide.rename(

columns={

"Outpatient": "outpatient_count_365d_postindex",

"Inpatient": "inpatient_count_365d_postindex",

"ED": "ED_count_365d_postindex",

"Other": "other_visit_count_365d_postindex",

util_total = (

visit_util.groupby("person_id", dropna=False)

.size()

.reset_index(name="visit_count_365d_postindex")
)

util_person = util_total.merge(util_wide, on="person_id", how="left")

utilization_df = analytic_cohort.merge(util_person, on="person_id", how="left")

for c in [

"visit_count_365d_postindex",

"outpatient_count_365d_postindex",

"inpatient_count_365d_postindex",

"ED_count_365d_postindex",

"other_visit_count_365d_postindex",

]:

if c not in utilization_df.columns:

utilization_df[c] = 0

utilization_df[c] = utilization_df[c].fillna(0)

audit_df("utilization_df", utilization_df)

utilization_summary = [Link](

[
utilization_df.groupby("race_3group", dropna=False)

.agg(

n=("person_id", "nunique"),

initiated_pct=("initiated", "mean"),

mean_visits=("visit_count_365d_postindex", "mean"),

mean_outpatient=("outpatient_count_365d_postindex", "mean"),

mean_inpatient=("inpatient_count_365d_postindex", "mean"),

mean_ed=("ED_count_365d_postindex", "mean"),

.reset_index()

.rename(columns={"race_3group": "group_value"})

.assign(group_type="race_3group"),

utilization_df.groupby("ethnicity_2group", dropna=False)

.agg(

n=("person_id", "nunique"),

initiated_pct=("initiated", "mean"),

mean_visits=("visit_count_365d_postindex", "mean"),

mean_outpatient=("outpatient_count_365d_postindex", "mean"),

mean_inpatient=("inpatient_count_365d_postindex", "mean"),
mean_ed=("ED_count_365d_postindex", "mean"),

.reset_index()

.rename(columns={"ethnicity_2group": "group_value"})

.assign(group_type="ethnicity_2group"),

utilization_df.groupby("sex_clean", dropna=False)

.agg(

n=("person_id", "nunique"),

initiated_pct=("initiated", "mean"),

mean_visits=("visit_count_365d_postindex", "mean"),

mean_outpatient=("outpatient_count_365d_postindex", "mean"),

mean_inpatient=("inpatient_count_365d_postindex", "mean"),

mean_ed=("ED_count_365d_postindex", "mean"),

.reset_index()

.rename(columns={"sex_clean": "group_value"})

.assign(group_type="sex_clean"),

],

ignore_index=True,
)

utilization_summary["initiated_pct"] = 100 * utilization_summary["initiated_pct"]

audit_df("utilization_summary", utilization_summary)

# 7F. Race (3-group) vs Ethnicity (2-group) association analysis.

race_eth_ct = [Link](

analytic_cohort["race_3group"].fillna("Other"),

analytic_cohort["ethnicity_2group"].fillna("Non-Hispanic"),

dropna=False,

race_eth_ct = race_eth_ct.reset_index()

audit_df("race_ethnicity_contingency", race_eth_ct)

if SCIPY_OK:

_ct = [Link](

analytic_cohort["race_3group"].fillna("Other"),

analytic_cohort["ethnicity_2group"].fillna("Non-Hispanic"),

dropna=False,

)
if _ct.shape[0] >= 2 and _ct.shape[1] >= 2 and _ct.[Link]() > 0:

chi2, pval, dof, _ = chi2_contingency(_ct.values)

n = _ct.[Link]()

min_dim = min(_ct.shape[0] - 1, _ct.shape[1] - 1)

cramers_v = [Link]((chi2 / n) / min_dim) if min_dim > 0 else [Link]

race_ethnicity_association = [Link](

[{

"status": "computed",

"chi2": float(chi2),

"p_value": float(pval),

"dof": int(dof),

"n": int(n),

"cramers_v": float(cramers_v) if [Link](cramers_v) else [Link],

"interpretation": "Association test between race_3group and ethnicity_2group.",

}]

else:

race_ethnicity_association = [Link](

[{
"status": "not computed: insufficient contingency dimensions",

"chi2": [Link],

"p_value": [Link],

"dof": [Link],

"n": int(_ct.[Link]()),

"cramers_v": [Link],

"interpretation": "Need at least 2x2 non-empty table.",

}]

else:

race_ethnicity_association = [Link](

[{

"status": "not computed: scipy unavailable",

"chi2": [Link],

"p_value": [Link],

"dof": [Link],

"n": int(len(analytic_cohort)),

"cramers_v": [Link],

"interpretation": "Install scipy to compute chi-square association.",


}]

audit_df("race_ethnicity_association", race_ethnicity_association)

# -----------------------------------------------------------------------------

# SECTION 8 - LONGITUDINAL DMT HISTORY (switching / escalation / discontinuation)

# -----------------------------------------------------------------------------

print("\n=== STEP 8: LONGITUDINAL DMT HISTORY ===")

# Restrict to mapped included DMT exposures in analytic cohort.

dmt_history = included_exposures.merge(

analytic_cohort[["person_id", "initiated", "treatment_start_date", "followup_end_date", "first_dmt_date", "race_3group", "ethnicity_2group", "sex_clean"]],

on="person_id",

how="inner",

# Initiator-only treatment-pathway clock: treatment_start_date (first DMT exposure date).


dmt_history = dmt_history[dmt_history["initiated"].eq(1)].copy()

dmt_history = dmt_history[dmt_history["treatment_start_date"].notna()].copy()

dmt_history = dmt_history[dmt_history["drug_date"] >= dmt_history["treatment_start_date"]].copy()

# Build expected end date using available fields and conservative assumptions.

dmt_history["end_date_from_drug_end"] = pd.to_datetime(dmt_history["drug_exposure_end_datetime"], errors="coerce").[Link]()

dmt_history["end_date_from_verbatim"] = pd.to_datetime(dmt_history["verbatim_end_date"], errors="coerce").[Link]()

dmt_history["days_supply_num"] = pd.to_numeric(dmt_history.get("days_supply", [Link]), errors="coerce")

dmt_history["expected_end_date"] = dmt_history["end_date_from_drug_end"]

mask_need_verbatim = (

dmt_history["expected_end_date"].isna()

& dmt_history["end_date_from_verbatim"].notna()

mask_need_verbatim_np = as_bool_mask(mask_need_verbatim, dmt_history.index, "mask_need_verbatim")

dmt_history.loc[mask_need_verbatim_np, "expected_end_date"] = dmt_history.loc[

mask_need_verbatim_np, "end_date_from_verbatim"

]
mask_need_days_supply = (

dmt_history["expected_end_date"].isna()

& dmt_history["days_supply_num"].notna()

& (dmt_history["days_supply_num"] > 0)

mask_need_days_supply_np = as_bool_mask(mask_need_days_supply, dmt_history.index, "mask_need_days_supply")

dmt_history.loc[mask_need_days_supply_np, "expected_end_date"] = (

dmt_history.loc[mask_need_days_supply_np, "drug_date"]

+ pd.to_timedelta(dmt_history.loc[mask_need_days_supply_np, "days_supply_num"], unit="D")

cond_end_from_drug_end_np = as_bool_mask(

dmt_history["end_date_from_drug_end"].notna(),

dmt_history.index,

"cond_end_from_drug_end",

assert isinstance(cond_end_from_drug_end_np, [Link]), "cond_end_from_drug_end must be [Link]"

assert isinstance(mask_need_verbatim_np, [Link]), "mask_need_verbatim must be [Link]"


assert isinstance(mask_need_days_supply_np, [Link]), "mask_need_days_supply must be [Link]"

assert cond_end_from_drug_end_np.dtype == bool, "cond_end_from_drug_end dtype must be bool"

assert mask_need_verbatim_np.dtype == bool, "mask_need_verbatim dtype must be bool"

assert mask_need_days_supply_np.dtype == bool, "mask_need_days_supply dtype must be bool"

assert cond_end_from_drug_end_np.shape[0] == len(dmt_history), "cond_end_from_drug_end length mismatch"

assert mask_need_verbatim_np.shape[0] == len(dmt_history), "mask_need_verbatim length mismatch"

assert mask_need_days_supply_np.shape[0] == len(dmt_history), "mask_need_days_supply length mismatch"

# If still missing expected_end_date, discontinuation cannot be reliably assessed for that episode.

dmt_history["end_date_source"] = [Link](

cond_end_from_drug_end_np,

mask_need_verbatim_np,

mask_need_days_supply_np,

],

["drug_exposure_end_datetime", "verbatim_end_date", "days_supply_imputation"],

default="missing_unassessable",

)
dmt_history = dmt_history.sort_values(["person_id", "drug_date", "drug_concept_id"]).copy()

dmt_history["exposure_order"] = dmt_history.groupby("person_id").cumcount() + 1

dmt_history["prior_dmt_class"] = dmt_history.groupby("person_id")["dmt_class"].shift(1)

dmt_history["switched_class"] = (

dmt_history["prior_dmt_class"].notna() & (dmt_history["dmt_class"] != dmt_history["prior_dmt_class"])

).astype(int)

# Escalation/de-escalation by efficacy rank.

efficacy_rank = {"Moderate": 1, "High": 2, "Off-label": 1}

dmt_history["efficacy_rank"] = dmt_history["efficacy_tier"].map(efficacy_rank)

dmt_history["prior_efficacy_rank"] = dmt_history.groupby("person_id")["efficacy_rank"].shift(1)

dmt_history["escalated"] = (

dmt_history["switched_class"].eq(1)

& dmt_history["efficacy_rank"].notna()

& dmt_history["prior_efficacy_rank"].notna()

& (dmt_history["efficacy_rank"] > dmt_history["prior_efficacy_rank"])

).astype(int)

dmt_history["deescalated"] = (

dmt_history["switched_class"].eq(1)
& dmt_history["efficacy_rank"].notna()

& dmt_history["prior_efficacy_rank"].notna()

& (dmt_history["efficacy_rank"] < dmt_history["prior_efficacy_rank"])

).astype(int)

dmt_history["is_high_efficacy_class"] = dmt_history["dmt_class"].isin(high_eff_classes).astype(int)

# Discontinuation: final gap > 180 days after expected end with no subsequent DMT.

last_episode = (

dmt_history.sort_values(["person_id", "drug_date", "expected_end_date"])

.groupby("person_id", as_index=False)

.tail(1)

.copy()

# Normalize to day-level before subtraction to avoid timezone/time-of-day artifacts.

for col in ["drug_date", "expected_end_date", "followup_end_date"]:

if col in last_episode.columns:

last_episode[col] = pd.to_datetime(last_episode[col], errors="coerce").[Link]()


last_episode["discontinuation_assessable"] = (

last_episode["expected_end_date"].notna() & last_episode["followup_end_date"].notna()

).astype(int)

last_episode["gap_after_last_end_days"] = (

last_episode["followup_end_date"] - last_episode["expected_end_date"]

).[Link]

last_episode["discontinued"] = [Link](

last_episode["discontinuation_assessable"].eq(1)

& (last_episode["gap_after_last_end_days"] > DISCONTINUATION_GAP_DAYS),

1,

0,

last_episode["discontinuation_date"] = [Link](

last_episode["discontinued"].eq(1),

last_episode["expected_end_date"] + pd.to_timedelta(DISCONTINUATION_GAP_DAYS, unit="D"),

[Link],

last_episode["discontinuation_date"] = pd.to_datetime(last_episode["discontinuation_date"], errors="coerce")

audit_df("last_episode", last_episode)
switching_person = (

dmt_history.groupby("person_id", as_index=False)

.agg(

first_dmt_class=("dmt_class", "first"),

ever_switched=("switched_class", "max"),

ever_escalated=("escalated", "max"),

ever_deescalated=("deescalated", "max"),

ever_high_efficacy=("is_high_efficacy_class", "max"),

.merge(

last_episode[["person_id", "discontinuation_assessable", "discontinued", "discontinuation_date", "followup_end_date", "first_dmt_date", "race_3group", "eth

nicity_2group", "sex_clean"]],

on="person_id",

how="left",

switching_person["event_or_censor_date"] = switching_person["followup_end_date"]
switching_person.loc[switching_person["discontinued"] == 1, "event_or_censor_date"] = switching_person.loc[

switching_person["discontinued"] == 1, "discontinuation_date"

switching_person["time_to_discontinuation"] = (

pd.to_datetime(switching_person["event_or_censor_date"], errors="coerce").[Link]()

- pd.to_datetime(switching_person["first_dmt_date"], errors="coerce").[Link]()

).[Link]

switching_person["time_to_discontinuation"] = pd.to_numeric(switching_person["time_to_discontinuation"], errors="coerce").clip(lower=0)

audit_df("dmt_history", dmt_history)

audit_df("switching_person", switching_person)

switching_summary = [Link](

switching_person.groupby("race_3group", dropna=False)

.agg(

n=("person_id", "nunique"),

ever_switched_n=("ever_switched", "sum"),

ever_escalated_n=("ever_escalated", "sum"),
ever_deescalated_n=("ever_deescalated", "sum"),

ever_high_efficacy_n=("ever_high_efficacy", "sum"),

discontinuation_assessable_n=("discontinuation_assessable", "sum"),

discontinued_n=("discontinued", "sum"),

median_time_to_discontinuation=("time_to_discontinuation", exact_median),

.reset_index()

.rename(columns={"race_3group": "group_value"})

.assign(group_type="race_3group"),

switching_person.groupby("ethnicity_2group", dropna=False)

.agg(

n=("person_id", "nunique"),

ever_switched_n=("ever_switched", "sum"),

ever_escalated_n=("ever_escalated", "sum"),

ever_deescalated_n=("ever_deescalated", "sum"),

ever_high_efficacy_n=("ever_high_efficacy", "sum"),

discontinuation_assessable_n=("discontinuation_assessable", "sum"),

discontinued_n=("discontinued", "sum"),

median_time_to_discontinuation=("time_to_discontinuation", exact_median),
)

.reset_index()

.rename(columns={"ethnicity_2group": "group_value"})

.assign(group_type="ethnicity_2group"),

switching_person.groupby("sex_clean", dropna=False)

.agg(

n=("person_id", "nunique"),

ever_switched_n=("ever_switched", "sum"),

ever_escalated_n=("ever_escalated", "sum"),

ever_deescalated_n=("ever_deescalated", "sum"),

ever_high_efficacy_n=("ever_high_efficacy", "sum"),

discontinuation_assessable_n=("discontinuation_assessable", "sum"),

discontinued_n=("discontinued", "sum"),

median_time_to_discontinuation=("time_to_discontinuation", exact_median),

.reset_index()

.rename(columns={"sex_clean": "group_value"})

.assign(group_type="sex_clean"),

],
ignore_index=True,

for n_col, p_col in [

("ever_switched_n", "ever_switched_pct"),

("ever_escalated_n", "ever_escalated_pct"),

("ever_deescalated_n", "ever_deescalated_pct"),

("ever_high_efficacy_n", "ever_high_efficacy_pct"),

]:

switching_summary[p_col] = 100 * switching_summary[n_col] / switching_summary["n"].replace(0, [Link])

switching_summary["discontinued_pct_among_assessable"] = 100 * switching_summary["discontinued_n"] / switching_summary["discontinuation_assessable_n"].replace(0, n

[Link])

audit_df("switching_summary", switching_summary)

print("\nDiscontinuation assumption note:")

print(

"Used drug_exposure_end_datetime first, then verbatim_end_date, then days_supply imputation. "

"If expected end date remained missing, discontinuation was marked unassessable."
)

# -----------------------------------------------------------------------------

# SECTION 9 - COMORBIDITY CONTEXT

# -----------------------------------------------------------------------------

print("\n=== STEP 9: COMORBIDITY CONTEXT ===")

COMORB_WINDOW = "all" # "all", "post_baseline_365", "pre_baseline_365"

MIN_OTHER_CONDITION_PEOPLE = 200

comorb = condition_df[["person_id", "condition_concept_id", "condition_start_datetime", "standard_concept_name"]].copy()

comorb = [Link](

analytic_cohort[["person_id", "cohort_baseline_date", "initiated", "race_3group", "ethnicity_2group", "sex_clean"]],

on="person_id",

how="inner",

comorb_candidates_after_merge = int(len(comorb))
comorb_people_candidates_after_merge = int(comorb["person_id"].nunique())

comorb["cond_date"] = pd.to_datetime(comorb["condition_start_datetime"], errors="coerce").[Link]()

comorb["name_lc"] = comorb["standard_concept_name"].astype(str).[Link]()

comorb = comorb[comorb["cond_date"].notna()].copy()

if COMORB_WINDOW == "post_baseline_365":

comorb = comorb[

(comorb["cond_date"] >= comorb["cohort_baseline_date"])

& (comorb["cond_date"] < comorb["cohort_baseline_date"] + pd.to_timedelta(365, unit="D"))

].copy()

elif COMORB_WINDOW == "pre_baseline_365":

comorb = comorb[

(comorb["cond_date"] < comorb["cohort_baseline_date"])

& (comorb["cond_date"] >= comorb["cohort_baseline_date"] - pd.to_timedelta(365, unit="D"))

].copy()

comorb_rows_after_window = int(len(comorb))

comorb_people_after_window = int(comorb["person_id"].nunique())

audit_df("comorb_rows_windowed", comorb)
print(f"[comorb validation] Rows after cohort merge/window: {len(comorb):,}")

print(f"[comorb validation] Unique people: {comorb['person_id'].nunique():,}")

print(

"[comorb validation] candidates after merge rows="

f"{comorb_candidates_after_merge:,}, people={comorb_people_candidates_after_merge:,}"

if comorb_rows_after_window == 0:

if comorb_candidates_after_merge == 0:

print("[comorb proof] no condition rows linked to analytic cohort after merge.")

elif COMORB_WINDOW == "all":

raise AssertionError(

"comorb is empty with COMORB_WINDOW='all' despite non-empty merged candidates. "

"Check condition datetime parsing and filters."

else:

print(

f"[comorb proof] no rows fell in requested window '{COMORB_WINDOW}' despite "

f"{comorb_candidates_after_merge:,} merged candidate rows."

)
# Known disease groups requested.

comorb["known_group"] = "other"

[Link][comorb["name_lc"].[Link](r"diabet|hyperglyc", regex=True, na=False), "known_group"] = "diabetes"

[Link][comorb["name_lc"].[Link](r"hypertension|high blood pressure", regex=True, na=False), "known_group"] = "hypertension"

[Link][comorb["name_lc"].[Link](r"coronary|ischemi|myocard|heart failure|atrial fib|stroke|cerebrovascular|angina|atheroscler|cardiomyopathy|arrhythm", r

egex=True, na=False), "known_group"] = "cvd"

cohort_n = int(analytic_cohort["person_id"].nunique())

known_group_summary = (

[Link]("known_group", dropna=False)

.agg(

n_people=("person_id", "nunique"),

n_rows=("person_id", "size"),

.reset_index()

.sort_values(["n_people", "n_rows"], ascending=False)

)
known_group_summary["pct_people"] = 100 * known_group_summary["n_people"] / max(cohort_n, 1)

audit_df("known_group_summary", known_group_summary)

top_other_conditions = (

comorb[comorb["known_group"] == "other"]

.groupby(["condition_concept_id", "standard_concept_name"], dropna=False)

.agg(

n_people=("person_id", "nunique"),

n_rows=("person_id", "size"),

.reset_index()

.sort_values(["n_people", "n_rows", "condition_concept_id"], ascending=[False, False, True])

top_other_conditions["pct_people"] = 100 * top_other_conditions["n_people"] / max(cohort_n, 1)

audit_df("top_other_conditions", top_other_conditions)

audit_df("top_other_conditions_100", top_other_conditions.head(100))

selected_other_conditions = top_other_conditions[top_other_conditions["n_people"] >= MIN_OTHER_CONDITION_PEOPLE].copy()

audit_df("selected_other_conditions_n>=200", selected_other_conditions)
# Person-level flags: diabetes, hypertension, cvd + high-prevalence other conditions.

cond_flags = analytic_cohort[["person_id"]].drop_duplicates().copy()

for col in ["has_diabetes", "has_hypertension", "has_cvd"]:

cond_flags[col] = 0

for grp, col in [("diabetes", "has_diabetes"), ("hypertension", "has_hypertension"), ("cvd", "has_cvd")]:

ids = [Link][comorb["known_group"] == grp, "person_id"].unique()

cond_flags.loc[cond_flags["person_id"].isin(ids), col] = 1

selected_flag_rows = []

selected_flag_cols = []

for _, r in selected_other_conditions.iterrows():

cid = r["condition_concept_id"]

cname = str(r["standard_concept_name"])

try:

cid_int = int(cid)

flag_col = f"has_other_{cid_int}"

except Exception:
flag_col = f"has_other_{len(selected_flag_cols)+1}"

if flag_col in cond_flags.columns:

continue

cond_flags[flag_col] = 0

ids = [Link][comorb["condition_concept_id"] == cid, "person_id"].unique()

cond_flags.loc[cond_flags["person_id"].isin(ids), flag_col] = 1

selected_flag_cols.append(flag_col)

selected_flag_rows.append(

"condition_concept_id": cid,

"standard_concept_name": cname,

"flag_column": flag_col,

"n_people": int(r["n_people"]),

"n_rows": int(r["n_rows"]),

"pct_people": float(r["pct_people"]),

selected_condition_flags_table = [Link](selected_flag_rows)
audit_df("selected_condition_flags_table", selected_condition_flags_table)

audit_df("comorbidity_flags", cond_flags)

comorbidity_person = analytic_cohort[["person_id", "race_3group", "ethnicity_2group", "sex_clean", "initiated"]].merge(

cond_flags,

on="person_id",

how="left",

for c in ["has_diabetes", "has_hypertension", "has_cvd"] + selected_flag_cols:

comorbidity_person[c] = comorbidity_person[c].fillna(0).astype(int)

comorbidity_person["other_selected_condition_count"] = (

comorbidity_person[selected_flag_cols].sum(axis=1) if selected_flag_cols else 0

comorbidity_person["comorbidity_burden_simple"] = (

comorbidity_person[["has_diabetes", "has_hypertension", "has_cvd"]].sum(axis=1)

+ comorbidity_person["other_selected_condition_count"]

audit_df("comorbidity_person", comorbidity_person)
comorbidity_summary = [Link](

comorbidity_person.groupby("race_3group", dropna=False)

.agg(

n=("person_id", "nunique"),

mean_burden=("comorbidity_burden_simple", "mean"),

median_burden=("comorbidity_burden_simple", "median"),

diabetes_prev=("has_diabetes", "mean"),

hypertension_prev=("has_hypertension", "mean"),

cvd_prev=("has_cvd", "mean"),

mean_other_selected_conditions=("other_selected_condition_count", "mean"),

.reset_index()

.rename(columns={"race_3group": "group_value"})

.assign(group_type="race_3group"),

comorbidity_person.groupby("ethnicity_2group", dropna=False)

.agg(

n=("person_id", "nunique"),
mean_burden=("comorbidity_burden_simple", "mean"),

median_burden=("comorbidity_burden_simple", "median"),

diabetes_prev=("has_diabetes", "mean"),

hypertension_prev=("has_hypertension", "mean"),

cvd_prev=("has_cvd", "mean"),

mean_other_selected_conditions=("other_selected_condition_count", "mean"),

.reset_index()

.rename(columns={"ethnicity_2group": "group_value"})

.assign(group_type="ethnicity_2group"),

comorbidity_person.groupby("sex_clean", dropna=False)

.agg(

n=("person_id", "nunique"),

mean_burden=("comorbidity_burden_simple", "mean"),

median_burden=("comorbidity_burden_simple", "median"),

diabetes_prev=("has_diabetes", "mean"),

hypertension_prev=("has_hypertension", "mean"),

cvd_prev=("has_cvd", "mean"),

mean_other_selected_conditions=("other_selected_condition_count", "mean"),
)

.reset_index()

.rename(columns={"sex_clean": "group_value"})

.assign(group_type="sex_clean"),

comorbidity_person.assign(initiation_group=[Link](comorbidity_person["initiated"] == 1, "Initiator", "Non-initiator"))

.groupby("initiation_group", dropna=False)

.agg(

n=("person_id", "nunique"),

mean_burden=("comorbidity_burden_simple", "mean"),

median_burden=("comorbidity_burden_simple", "median"),

diabetes_prev=("has_diabetes", "mean"),

hypertension_prev=("has_hypertension", "mean"),

cvd_prev=("has_cvd", "mean"),

mean_other_selected_conditions=("other_selected_condition_count", "mean"),

.reset_index()

.rename(columns={"initiation_group": "group_value"})

.assign(group_type="initiation_group"),

],
ignore_index=True,

for c in ["diabetes_prev", "hypertension_prev", "cvd_prev"]:

comorbidity_summary[c] = 100 * comorbidity_summary[c]

audit_df("comorbidity_summary", comorbidity_summary)

comorbidity_limitations = [Link]([

{"note": "Comorbidity groups are keyword-based from loaded condition names and may under/over-capture true prevalence."},

{"note": f"Additional 'other condition' flags include only conditions with n_people >= {MIN_OTHER_CONDITION_PEOPLE} in the analytic cohort."},

{"note": "No formal validated Charlson Comorbidity Index is implemented in this loaded-data-only pipeline."},

])

audit_df("comorbidity_limitations", comorbidity_limitations)

# -----------------------------------------------------------------------------

# DATA ASSURANCE AND ROBUSTNESS CHECKS

# -----------------------------------------------------------------------------
print("\n=== DATA ASSURANCE AND ROBUSTNESS CHECKS ===")

# A) Cohort integrity and timeline plausibility

cohort_n = int(len(analytic_cohort))

dob_lookup = person_df[["person_id", "date_of_birth"]].copy()

dob_lookup["date_of_birth"] = pd.to_datetime(dob_lookup["date_of_birth"], errors="coerce").[Link]()

drug_with_dob = dmt_history[["person_id", "drug_date"]].merge(dob_lookup, on="person_id", how="left")

def _qc_row(check, violations, denominator):

denom = int(max(denominator, 1))

v = int(violations)

pct = 100.0 * v / denom

return {

"check": check,

"n_violations": v,

"denominator_n": int(denominator),

"pct_violations": pct,

"status": "PASS" if v == 0 else "FAIL",

}
cohort_timeline_qc_rows = []

cohort_timeline_qc_rows.append(

_qc_row(

"analytic_cohort_one_row_per_person",

int(len(analytic_cohort) - analytic_cohort["person_id"].nunique()),

cohort_n,

cohort_timeline_qc_rows.append(

_qc_row(

"cohort_baseline_date_gte_index_date",

int(

analytic_cohort["cohort_baseline_date"].notna()

& analytic_cohort["index_date"].notna()

& (analytic_cohort["cohort_baseline_date"] < analytic_cohort["index_date"])

).sum()

),
cohort_n,

cohort_timeline_qc_rows.append(

_qc_row(

"followup_end_date_gte_cohort_baseline_date",

int(

analytic_cohort["followup_end_date"].notna()

& analytic_cohort["cohort_baseline_date"].notna()

& (analytic_cohort["followup_end_date"] < analytic_cohort["cohort_baseline_date"])

).sum()

),

cohort_n,

cohort_timeline_qc_rows.append(

_qc_row(

"age_at_index_gte_18",
int((pd.to_numeric(analytic_cohort["age_at_index"], errors="coerce") < 18).fillna(False).sum()),

cohort_n,

cohort_timeline_qc_rows.append(

_qc_row(

"drug_date_not_before_date_of_birth",

int(

drug_with_dob["drug_date"].notna()

& drug_with_dob["date_of_birth"].notna()

& (drug_with_dob["drug_date"] < drug_with_dob["date_of_birth"])

).sum()

),

int(drug_with_dob["drug_date"].notna().sum()),

initiator_df = analytic_cohort[analytic_cohort["initiated"] == 1].copy()

cohort_timeline_qc_rows.append(
_qc_row(

"first_dmt_date_not_before_treatment_start_date_for_initiators",

int(

initiator_df["first_dmt_date"].notna()

& initiator_df["treatment_start_date"].notna()

& (initiator_df["first_dmt_date"] < initiator_df["treatment_start_date"])

).sum()

),

int(len(initiator_df)),

cohort_timeline_qc_rows.append(

_qc_row(

"days_to_dmt_nonnegative",

int((pd.to_numeric(analytic_cohort["days_to_dmt"], errors="coerce") < 0).fillna(False).sum()),

int(pd.to_numeric(analytic_cohort["days_to_dmt"], errors="coerce").notna().sum()),

)
cohort_timeline_qc_rows.append(

_qc_row(

"time_to_event_days_nonnegative",

int((pd.to_numeric(analytic_cohort["time_to_event_days"], errors="coerce") < 0).fillna(False).sum()),

int(pd.to_numeric(analytic_cohort["time_to_event_days"], errors="coerce").notna().sum()),

cohort_timeline_qc_rows.append(

_qc_row(

"time_to_discontinuation_nonnegative",

int((pd.to_numeric(switching_person["time_to_discontinuation"], errors="coerce") < 0).fillna(False).sum()),

int(pd.to_numeric(switching_person["time_to_discontinuation"], errors="coerce").notna().sum()),

cohort_timeline_qc_rows.append(

_qc_row(

"days_to_high_efficacy_dmt_nonnegative",

int((pd.to_numeric(high_efficacy_df["days_to_high_efficacy_dmt"], errors="coerce") < 0).fillna(False).sum()),

int(pd.to_numeric(high_efficacy_df["days_to_high_efficacy_dmt"], errors="coerce").notna().sum()),
)

data_assurance_cohort_timeline_qc = [Link](cohort_timeline_qc_rows)

audit_df("data_assurance_cohort_timeline_qc", data_assurance_cohort_timeline_qc)

# B) DMT exposure construction audit

dmt_end_date_source_summary = (

dmt_history["end_date_source"]

.fillna("missing_unassessable")

.value_counts(dropna=False)

.rename_axis("end_date_source")

.reset_index(name="n")

dmt_end_date_source_summary["pct"] = (

100 * dmt_end_date_source_summary["n"] / max(len(dmt_history), 1)

audit_df("dmt_end_date_source_summary", dmt_end_date_source_summary)

dmt_end_date_source_by_class = (
dmt_history.assign(

dmt_class=dmt_history["dmt_class"].fillna("Missing"),

end_date_source=dmt_history["end_date_source"].fillna("missing_unassessable"),

.groupby(["dmt_class", "end_date_source"], dropna=False)

.size()

.reset_index(name="n")

_class_tot = (

dmt_end_date_source_by_class.groupby("dmt_class", dropna=False)["n"]

.sum()

.reset_index(name="class_total_n")

dmt_end_date_source_by_class = dmt_end_date_source_by_class.merge(

_class_tot, on="dmt_class", how="left"

dmt_end_date_source_by_class["pct_within_class"] = (

100 * dmt_end_date_source_by_class["n"] / dmt_end_date_source_by_class["class_total_n"].replace(0, [Link])

)
audit_df("dmt_end_date_source_by_class", dmt_end_date_source_by_class)

_end_before_start_n = int(

dmt_history["expected_end_date"].notna()

& dmt_history["drug_date"].notna()

& (dmt_history["expected_end_date"] < dmt_history["drug_date"])

).sum()

_ord_sorted = dmt_history.sort_values(["person_id", "drug_date", "drug_concept_id", "exposure_order"]).copy()

_ord_sorted["prev_exposure_order"] = _ord_sorted.groupby("person_id")["exposure_order"].shift(1)

_order_not_strict_n = int(

_ord_sorted["prev_exposure_order"].notna()

& (_ord_sorted["exposure_order"] <= _ord_sorted["prev_exposure_order"])

).sum()

)
_order_span = (

dmt_history.groupby("person_id", dropna=False)["exposure_order"]

.agg(min_order="min", max_order="max", nunique_order="nunique", n_rows="count")

.reset_index()

_order_sequence_gap_person_n = int(

(_order_span["min_order"] != 1)

| (_order_span["max_order"] != _order_span["n_rows"])

| (_order_span["nunique_order"] != _order_span["n_rows"])

).sum()

_first_exposure = (

_ord_sorted.groupby("person_id", as_index=False)

.head(1)

.copy()

_first_prior_present_n = int(_first_exposure["prior_dmt_class"].notna().sum())
dmt_construction_qc_rows = [

_qc_row(

"expected_end_date_gte_drug_date_when_both_present",

_end_before_start_n,

int((dmt_history["expected_end_date"].notna() & dmt_history["drug_date"].notna()).sum()),

),

_qc_row(

"exposure_order_strictly_increasing_within_person",

_order_not_strict_n,

int(len(dmt_history)),

),

_qc_row(

"exposure_order_sequence_has_no_person_level_gaps",

_order_sequence_gap_person_n,

int(dmt_history["person_id"].nunique()),

),

_qc_row(

"first_exposure_has_no_prior_dmt_class",
_first_prior_present_n,

int(len(_first_exposure)),

),

data_assurance_dmt_construction_qc = [Link](dmt_construction_qc_rows)

audit_df("data_assurance_dmt_construction_qc", data_assurance_dmt_construction_qc)

# C) Discontinuation robustness sensitivity (descriptive only)

disc_sensitivity_rows = []

for _gap in [90, 180, 365]:

_tmp = last_episode.copy()

_tmp["discontinuation_assessable_tmp"] = (

_tmp["expected_end_date"].notna() & _tmp["followup_end_date"].notna()

_tmp["gap_after_last_end_days_tmp"] = (

_tmp["followup_end_date"] - _tmp["expected_end_date"]

).[Link]

_tmp["discontinued_tmp"] = [Link](

_tmp["discontinuation_assessable_tmp"]
& (_tmp["gap_after_last_end_days_tmp"] > _gap),

1,

0,

_tmp["discontinuation_date_tmp"] = [Link]

_m = _tmp["discontinued_tmp"].eq(1) & _tmp["expected_end_date"].notna()

_tmp.loc[_m, "discontinuation_date_tmp"] = (

_tmp.loc[_m, "expected_end_date"] + pd.to_timedelta(_gap, unit="D")

_tmp["event_or_censor_date_tmp"] = _tmp["followup_end_date"]

_tmp.loc[_tmp["discontinued_tmp"].eq(1), "event_or_censor_date_tmp"] = _tmp.loc[

_tmp["discontinued_tmp"].eq(1), "discontinuation_date_tmp"

_tmp["time_to_discontinuation_tmp"] = (

pd.to_datetime(_tmp["event_or_censor_date_tmp"], errors="coerce").[Link]()

- pd.to_datetime(_tmp["first_dmt_date"], errors="coerce").[Link]()

).[Link]

_tmp["time_to_discontinuation_tmp"] = pd.to_numeric(

_tmp["time_to_discontinuation_tmp"], errors="coerce"
).clip(lower=0)

_assess_n = int(_tmp["discontinuation_assessable_tmp"].sum())

_disc_n = int(_tmp["discontinued_tmp"].sum())

disc_sensitivity_rows.append(

"threshold_days": int(_gap),

"n_people": int(len(_tmp)),

"discontinuation_assessable_n": _assess_n,

"discontinued_n": _disc_n,

"discontinued_pct_among_assessable": (

100 * _disc_n / _assess_n if _assess_n > 0 else [Link]

),

"median_time_to_discontinuation": exact_median(

_tmp.loc[_tmp["discontinued_tmp"].eq(1), "time_to_discontinuation_tmp"]

),

)
discontinuation_sensitivity_summary = [Link](disc_sensitivity_rows)

audit_df("discontinuation_sensitivity_summary", discontinuation_sensitivity_summary)

# D) Baseline definition sensitivity (cohort_baseline_date vs index_date)

def _baseline_sensitivity(df_cohort, exposures_df, baseline_col, baseline_name):

_base = df_cohort[["person_id", baseline_col]].copy()

_base["baseline_date"] = pd.to_datetime(_base[baseline_col], errors="coerce").[Link]()

_exp = exposures_df[["person_id", "drug_date"]].copy()

_cand = _base.merge(_exp, on="person_id", how="left")

_cand = _cand[

_cand["baseline_date"].notna()

& _cand["drug_date"].notna()

& (_cand["drug_date"] >= _cand["baseline_date"])

].copy()

_first = (

_cand.sort_values(["person_id", "drug_date"])

.drop_duplicates("person_id", keep="first")

[["person_id", "drug_date"]]
.rename(columns={"drug_date": "first_dmt_date_sensitivity"})

_res = _base.merge(_first, on="person_id", how="left")

_res["initiated"] = _res["first_dmt_date_sensitivity"].notna().astype(int)

_res["days_to_dmt"] = (

_res["first_dmt_date_sensitivity"] - _res["baseline_date"]

).[Link]

return {

"baseline_definition": baseline_name,

"cohort_n": int(len(_res)),

"initiated_n": int(_res["initiated"].sum()),

"initiated_pct": float(100 * _res["initiated"].mean()) if len(_res) > 0 else [Link],

"median_days_to_dmt": exact_median(_res.loc[_res["initiated"].eq(1), "days_to_dmt"]),

"mean_days_to_dmt": float(pd.to_numeric(_res.loc[_res["initiated"].eq(1), "days_to_dmt"], errors="coerce").mean()),

baseline_definition_sensitivity_summary = [Link](

_baseline_sensitivity(analytic_cohort, included_exposures, "cohort_baseline_date", "cohort_baseline_date"),


_baseline_sensitivity(analytic_cohort, included_exposures, "index_date", "index_date"),

audit_df("baseline_definition_sensitivity_summary", baseline_definition_sensitivity_summary)

# E) Manual audit sample output (reproducible n=25 from exposed persons)

MANUAL_AUDIT_SAMPLE_N = 25

MANUAL_AUDIT_SEED = 42

manual_cols = [

"person_id",

"drug_date",

"ingredient_name",

"dmt_class",

"expected_end_date",

"end_date_source",

"exposure_order",

"prior_dmt_class",

"switched_class",

"escalated",
"deescalated",

_exposed_ids = [Link](sorted(dmt_history["person_id"].dropna().unique()))

if len(_exposed_ids) > 0:

_rng = [Link].default_rng(MANUAL_AUDIT_SEED)

_sample_n = int(min(MANUAL_AUDIT_SAMPLE_N, len(_exposed_ids)))

_sample_ids = _rng.choice(_exposed_ids, size=_sample_n, replace=False)

dmt_manual_audit_sample_25 = (

dmt_history[dmt_history["person_id"].isin(_sample_ids)][manual_cols]

.sort_values(["person_id", "drug_date", "exposure_order"])

.reset_index(drop=True)

else:

dmt_manual_audit_sample_25 = [Link](columns=manual_cols)

print("[INFO] Manual audit sample is empty because no DMT exposures were available.")

audit_df("dmt_manual_audit_sample_25", dmt_manual_audit_sample_25, head_n=25)

# Baseline standardization audit.


baseline_definition_audit = [Link](

{"section_name": "analytic_cohort", "population_type": "whole_cohort", "baseline_used": "cohort_baseline_date(index_date)"},

{"section_name": "delay_category_summary", "population_type": "whole_cohort", "baseline_used": "cohort_baseline_date(index_date)"},

{"section_name": "calendar_era_summary", "population_type": "whole_cohort", "baseline_used": "cohort_baseline_date(index_date)"},

{"section_name": "utilization_df", "population_type": "whole_cohort", "baseline_used": "cohort_baseline_date(index_date)"},

{"section_name": "utilization_summary", "population_type": "whole_cohort", "baseline_used": "cohort_baseline_date(index_date)"},

{"section_name": "comorbidity_person", "population_type": "whole_cohort", "baseline_used": "cohort_baseline_date(index_date)"},

{"section_name": "comorbidity_summary", "population_type": "whole_cohort", "baseline_used": "cohort_baseline_date(index_date)"},

{"section_name": "high_efficacy_df", "population_type": "whole_cohort", "baseline_used": "cohort_baseline_date(index_date)"},

{"section_name": "high_efficacy_summary", "population_type": "whole_cohort", "baseline_used": "cohort_baseline_date(index_date)"},

{"section_name": "dmt_history", "population_type": "initiator_only", "baseline_used": "treatment_start_date(first_dmt_date)"},

{"section_name": "switching_person", "population_type": "initiator_only", "baseline_used": "treatment_start_date(first_dmt_date)"},

{"section_name": "switching_summary", "population_type": "initiator_only", "baseline_used": "treatment_start_date(first_dmt_date)"},

baseline_definition_audit["status"] = [Link](

(baseline_definition_audit["population_type"] == "whole_cohort")
& baseline_definition_audit["baseline_used"].[Link]("index_date", regex=False)

| (

(baseline_definition_audit["population_type"] == "initiator_only")

& baseline_definition_audit["baseline_used"].[Link]("first_dmt_date", regex=False)

),

"PASS",

"FAIL",

audit_df("baseline_definition_audit", baseline_definition_audit)

print(

"For analyses including both initiators and non-initiators, index_date was used as the common baseline. "

"For initiator-only treatment-pathway analyses, first DMT exposure date was used as the treatment-specific baseline."

# -----------------------------------------------------------------------------

# SECTION 10 - FINAL QC REVIEW / SAVE DELIVERABLES


# -----------------------------------------------------------------------------

print("\n=== STEP 10: FINAL QC REVIEW ===")

# Basic QC checks

qc_checks = []

qc_checks.append({"check": "analytic_cohort_unique_person", "value": int(analytic_cohort["person_id"].nunique()), "status": "ok" if int(analytic_cohort["person_id"

].nunique()) == len(analytic_cohort) else "review"})

qc_checks.append({"check": "dmt_mapping_unique_by_concept", "value": int(dmt_mapping_unique_ok), "status": "ok" if dmt_mapping_unique_ok else "critical"})

qc_checks.append({"check": "drug_merge_no_row_inflation", "value": int(drug_merge_no_inflation_ok), "status": "ok" if drug_merge_no_inflation_ok else "critical"})

qc_checks.append({"check": "sensitivity_n", "value": int(cohort_sensitivity["person_id"].nunique()), "status": "ok"})

qc_checks.append({"check": "initiated_n", "value": int(analytic_cohort["initiated"].sum()), "status": "ok"})

qc_checks.append({"check": "pre_index_dmt_n", "value": int(analytic_cohort["pre_index_dmt_flag"].sum()), "status": "ok"})

qc_checks.append({"check": "lookback_restriction_removed", "value": int((cohort_flow_summary["step"] == "lookback_removed_for_ms_only").any()), "status": "ok"})

qc_checks.append({"check": "utilization_rows_nonempty", "value": int(len(visit_util)), "status": "ok" if len(visit_util) > 0 else "review"})

qc_checks.append({"check": "comorb_rows_nonempty", "value": int(len(comorb)), "status": "ok" if len(comorb) > 0 else "review"})

qc_checks.append({"check": "selected_other_conditions_n>=200", "value": int(len(selected_other_conditions)), "status": "ok" if len(selected_other_conditions) > 0 e

lse "review"})

qc_checks.append({"check": "race_group_other_n", "value": int((analytic_cohort["race_3group"] == "Other").sum()), "status": "ok"})


qc_checks.append({"check": "ethnicity_hispanic_n", "value": int((analytic_cohort["ethnicity_2group"] == "Hispanic").sum()), "status": "ok"})

qc_checks.append({"check": "unknown_sex_n", "value": int((analytic_cohort["sex_clean"] == "Unknown/Other").sum()), "status": "ok"})

qc_checks.append({

"check": "baseline_definition_audit_all_pass",

"value": int((baseline_definition_audit["status"] == "PASS").all()),

"status": "ok" if bool((baseline_definition_audit["status"] == "PASS").all()) else "critical",

})

qc_df = [Link](qc_checks)

audit_df("qc_df", qc_df)

# Required PASS/FAIL gates.

cohort_flow_ok = bool(

(cohort_flow_summary["step"] == "lookback_removed_for_ms_only").any()

and int(

cohort_flow_summary.loc[

cohort_flow_summary["step"] == "lookback_removed_for_ms_only", "n_excluded"

].sum()

)
== 0

utilization_ok = bool((len(visit_util) > 0) or (util_candidates_before_window == 0))

comorb_ok = bool(

(len(comorb) > 0 or comorb_candidates_after_merge == 0)

and {"has_diabetes", "has_hypertension", "has_cvd"}.issubset(set(comorbidity_person.columns))

and (len(selected_other_conditions) > 0 or len(comorb) == 0)

high_eff_ok = bool(len(high_efficacy_summary) > 0)

baseline_standardization_ok = bool((baseline_definition_audit["status"] == "PASS").all())

qc_gates = [Link](

make_gate(

"mapping_unique_and_no_inflation",

dmt_mapping_unique_ok and drug_merge_no_inflation_ok,

f"dmt_mapping_unique={dmt_mapping_unique_ok}, "

f"drug_rows={len(drug_df):,}, drug_mapped_rows={len(drug_mapped):,}"
),

),

make_gate(

"cohort_flow_consistency",

cohort_flow_ok,

"lookback_removed_for_ms_only present with n_excluded=0",

),

make_gate(

"utilization_nonempty_or_proven_none",

utilization_ok,

f"candidates_before_window={util_candidates_before_window:,}, "

f"rows_after_window={len(visit_util):,}"

),

),

make_gate(

"comorbidity_nonempty_and_required_flags",

comorb_ok,

(
f"comorb_candidates={comorb_candidates_after_merge:,}, "

f"rows_after_window={len(comorb):,}, selected_other_conditions={len(selected_other_conditions):,}"

),

),

make_gate(

"high_efficacy_summary_nonempty",

high_eff_ok,

f"rows={len(high_efficacy_summary):,}",

),

make_gate(

"baseline_standardization_consistent",

baseline_standardization_ok,

"whole-cohort sections use cohort_baseline_date(index_date); "

"initiator-only sections use treatment_start_date(first_dmt_date)"

),

),

)
audit_df("qc_gates", qc_gates)

rowcount_fix_summary = [Link](

{"section": "dmt_mapping", "before": original_unique_concepts_n, "after": len(cleaned_dmt_mapping), "detail": "unique concepts vs canonical mapping rows"},

{"section": "drug_merge_rows", "before": len(drug_df), "after": len(drug_mapped), "detail": "must be unchanged"},

{"section": "included_exposure_dedup", "before": included_rows_before_dedup, "after": included_rows_after_dedup, "detail": f"exact_removed={exact_dupe_rows

_n}, key_removed={key_dupe_rows_n}"},

{"section": "cohort_age_filter", "before": len(cohort0), "after": len(cohort1), "detail": "age>=18"},

{"section": "lookback_exclusion_removed", "before": len(cohort3), "after": len(cohort4), "detail": "no exclusion"},

{"section": "baseline_standardization", "before": 1, "after": 1, "detail": "whole_cohort=index_date, initiator_only=first_dmt_date"},

{"section": "utilization_window", "before": util_candidates_before_window, "after": len(visit_util), "detail": f"post-index {UTIL_WINDOW_DAYS_POST}d"},

{"section": "comorbidity_window", "before": comorb_candidates_after_merge, "after": len(comorb), "detail": f"COMORB_WINDOW={COMORB_WINDOW}"},

audit_df("rowcount_fix_summary", rowcount_fix_summary)

critical_issues_fixed = [

"Removed dependence on new SQL pulls; pipeline now uses loaded dataframes only.",
"DMT mapping is now canonical one-row-per-concept with many_to_one merge validation and no row inflation.",

"MS-only cohort no longer excludes on 365-day lookback; lookback is retained as QC only.",

"Pre-index DMT is no longer used as an exclusion concept; pre_index_dmt_flag is forced to 0 in this MS-only extract.",

"Baseline standardization is explicit: whole-cohort analyses use index_date, initiator-only treatment-pathway analyses use first_dmt_date.",

"Utilization rebuilt on post-baseline 365-day window with date normalization and explicit non-empty validation.",

"Comorbidity rebuilt with diabetes/hypertension/cvd plus additional high-prevalence other conditions (n_people >= 200).",

"High-efficacy summary is descriptive-only and log-rank columns are removed to avoid ambiguous inference outputs.",

"Race analyses now use fixed 3-group race (White, Black, Other), and ethnicity analyses use Hispanic vs Non-Hispanic.",

"Added race-vs-ethnicity association analysis (chi-square + Cramer's V when scipy is available).",

"Unsupported inferred person-level covariates (insurance/state/rural/employment) are not required for final outputs.",

remaining_limitations = [

"MS diagnosis identification is based on loaded condition concept names (no live concept hierarchy expansion in this script).",

"Discontinuation depends on end-date availability; rows without end-date evidence remain unassessable.",

"Comorbidity definitions are rule-based from condition names and not a fully validated phenotype library.",

removed_or_deferred = [
"Insurance/payer type as a patient-level payer variable (not robustly derivable from loaded tables only).",

"State/Census region/rural-urban classification (not directly available in loaded dataframes).",

"Formal mediation models requiring unavailable validated mediators.",

print("\nCritical issues fixed:")

for x in critical_issues_fixed:

print(f"- {x}")

print("\nRemaining limitations:")

for x in remaining_limitations:

print(f"- {x}")

print("\nAnalyses removed/deferred:")

for x in removed_or_deferred:

print(f"- {x}")

# Save required outputs

cohort_flow_summary.to_csv(f"{OUTPUT_DIR}/cohort_flow_summary.csv", index=False)
analytic_cohort.to_csv(f"{OUTPUT_DIR}/analytic_cohort.csv", index=False)

# Save policy-safe displayed versions under required deliverable names.

display_descriptive_table_sex.to_csv(f"{OUTPUT_DIR}/descriptive_table_sex.csv", index=False)

display_descriptive_table_race_eth.to_csv(f"{OUTPUT_DIR}/descriptive_table_race_eth.csv", index=False)

display_descriptive_table_ethnicity.to_csv(f"{OUTPUT_DIR}/descriptive_table_ethnicity.csv", index=False)

display_descriptive_table_first_dmt.to_csv(f"{OUTPUT_DIR}/descriptive_table_first_dmt.csv", index=False)

display_race_by_dmt_crosstab.to_csv(f"{OUTPUT_DIR}/race_by_dmt_crosstab.csv", index=False)

delay_category_summary.to_csv(f"{OUTPUT_DIR}/delay_category_summary.csv", index=False)

delay_category_ethnicity_summary.to_csv(f"{OUTPUT_DIR}/delay_category_ethnicity_summary.csv", index=False)

calendar_era_summary.to_csv(f"{OUTPUT_DIR}/calendar_era_summary.csv", index=False)

high_efficacy_summary.to_csv(f"{OUTPUT_DIR}/high_efficacy_summary.csv", index=False)

race_ethnicity_association.to_csv(f"{OUTPUT_DIR}/race_ethnicity_association.csv", index=False)

race_eth_ct.to_csv(f"{OUTPUT_DIR}/race_ethnicity_contingency.csv", index=False)

intersectional_summary.to_csv(f"{OUTPUT_DIR}/intersectional_summary.csv", index=False)

utilization_summary.to_csv(f"{OUTPUT_DIR}/utilization_summary.csv", index=False)

dmt_history.to_csv(f"{OUTPUT_DIR}/dmt_history_analytic.csv", index=False)
switching_summary.to_csv(f"{OUTPUT_DIR}/switching_summary.csv", index=False)

comorbidity_summary.to_csv(f"{OUTPUT_DIR}/comorbidity_summary.csv", index=False)

known_group_summary.to_csv(f"{OUTPUT_DIR}/comorbidity_known_group_summary.csv", index=False)

top_other_conditions.to_csv(f"{OUTPUT_DIR}/comorbidity_top_other_conditions.csv", index=False)

selected_condition_flags_table.to_csv(f"{OUTPUT_DIR}/comorbidity_selected_other_condition_flags.csv", index=False)

missingness_summary.to_csv(f"{OUTPUT_DIR}/missingness_summary.csv", index=False)

dmt_exposure_dedup_summary.to_csv(f"{OUTPUT_DIR}/dmt_exposure_dedup_summary.csv", index=False)

rowcount_fix_summary.to_csv(f"{OUTPUT_DIR}/rowcount_fix_summary.csv", index=False)

qc_gates.to_csv(f"{OUTPUT_DIR}/qc_gates_pass_fail.csv", index=False)

[Link](root_cause_items).to_csv(f"{OUTPUT_DIR}/root_cause_issue_cause_fix.csv", index=False)

data_assurance_cohort_timeline_qc.to_csv(f"{OUTPUT_DIR}/data_assurance_cohort_timeline_qc.csv", index=False)

data_assurance_dmt_construction_qc.to_csv(f"{OUTPUT_DIR}/data_assurance_dmt_construction_qc.csv", index=False)

dmt_end_date_source_summary.to_csv(f"{OUTPUT_DIR}/dmt_end_date_source_summary.csv", index=False)

dmt_end_date_source_by_class.to_csv(f"{OUTPUT_DIR}/dmt_end_date_source_by_class.csv", index=False)

discontinuation_sensitivity_summary.to_csv(f"{OUTPUT_DIR}/discontinuation_sensitivity_summary.csv", index=False)

baseline_definition_sensitivity_summary.to_csv(f"{OUTPUT_DIR}/baseline_definition_sensitivity_summary.csv", index=False)

baseline_definition_audit.to_csv(f"{OUTPUT_DIR}/baseline_definition_audit.csv", index=False)

dmt_manual_audit_sample_25.to_csv(f"{OUTPUT_DIR}/dmt_manual_audit_sample_25.csv", index=False)
# Save exact (unsuppressed) descriptive tables as additional reproducibility artifacts.

descriptive_table_sex.to_csv(f"{OUTPUT_DIR}/descriptive_table_sex_exact_unsuppressed.csv", index=False)

descriptive_table_race_eth.to_csv(f"{OUTPUT_DIR}/descriptive_table_race_eth_exact_unsuppressed.csv", index=False)

descriptive_table_ethnicity.to_csv(f"{OUTPUT_DIR}/descriptive_table_ethnicity_exact_unsuppressed.csv", index=False)

descriptive_table_first_dmt.to_csv(f"{OUTPUT_DIR}/descriptive_table_first_dmt_exact_unsuppressed.csv", index=False)

race_by_dmt_crosstab.to_csv(f"{OUTPUT_DIR}/race_by_dmt_crosstab_exact_unsuppressed.csv", index=False)

# Save fixed-issues summary notes for reproducibility.

[Link]({"critical_issues_fixed": critical_issues_fixed}).to_csv(f"{OUTPUT_DIR}/_qc_critical_issues_fixed.csv", index=False)

[Link]({"remaining_limitations": remaining_limitations}).to_csv(f"{OUTPUT_DIR}/_qc_remaining_limitations.csv", index=False)

[Link]({"removed_or_deferred": removed_or_deferred}).to_csv(f"{OUTPUT_DIR}/_qc_removed_or_deferred.csv", index=False)

qc_df.to_csv(f"{OUTPUT_DIR}/_qc_checks.csv", index=False)

validation_compare_df.to_csv(f"{OUTPUT_DIR}/_validation_compare_existing_analysis.csv", index=False)

summary_exact.to_csv(f"{OUTPUT_DIR}/_summary_exact_metrics.csv", index=False)

print("\nSaved required deliverables to:", OUTPUT_DIR)

print("- cleaned_dmt_mapping.csv")

print("- cohort_flow_summary.csv")

print("- analytic_cohort.csv")
print("- descriptive_table_sex.csv")

print("- descriptive_table_race_eth.csv")

print("- descriptive_table_ethnicity.csv")

print("- descriptive_table_first_dmt.csv")

print("- race_by_dmt_crosstab.csv")

print("- delay_category_summary.csv")

print("- delay_category_ethnicity_summary.csv")

print("- calendar_era_summary.csv")

print("- high_efficacy_summary.csv")

print("- race_ethnicity_association.csv")

print("- race_ethnicity_contingency.csv")

print("- intersectional_summary.csv")

print("- utilization_summary.csv")

print("- dmt_history_analytic.csv")

print("- switching_summary.csv")

print("- comorbidity_summary.csv")

print("- comorbidity_known_group_summary.csv")

print("- comorbidity_top_other_conditions.csv")

print("- comorbidity_selected_other_condition_flags.csv")
print("- missingness_summary.csv")

print("- dmt_exposure_dedup_summary.csv")

print("- rowcount_fix_summary.csv")

print("- qc_gates_pass_fail.csv")

print("- root_cause_issue_cause_fix.csv")

print("- data_assurance_cohort_timeline_qc.csv")

print("- data_assurance_dmt_construction_qc.csv")

print("- dmt_end_date_source_summary.csv")

print("- dmt_end_date_source_by_class.csv")

print("- discontinuation_sensitivity_summary.csv")

print("- baseline_definition_sensitivity_summary.csv")

print("- baseline_definition_audit.csv")

print("- dmt_manual_audit_sample_25.csv")

print("\nPipeline complete.")

You might also like