Compare commits

..
Author SHA1 Message Date
TudorandClaude Fable 5 d677b54533 fix(pipeline): cache invalidation for IDACI DAG too; curl timeouts
PR Checks / Frontend Typecheck + Tests (pull_request) Successful in 9m37s
PR Checks / Backend Smoke (pull_request) Successful in 6s
PR Checks / Build Backend (no push) (pull_request) Successful in 18s
PR Checks / Build Frontend (no push) (pull_request) Successful in 48s
PR Checks / Build Pipeline (no push) (pull_request) Successful in 35s
PR Checks / AI Code Review (Claude) (pull_request) Successful in 1m40s
Addresses AI-review findings: the annual IDACI DAG also rebuilds a mart
(fact_deprivation) and needs the reload; curl gets connect/max timeouts
so an unreachable backend fails fast instead of hanging the task.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-09 21:26:45 +01:00
TudorandClaude Fable 5 c353e36072 fix(pipeline): actually invalidate the backend cache after data rebuilds
PR Checks / Frontend Typecheck + Tests (pull_request) Successful in 9m40s
PR Checks / Backend Smoke (pull_request) Successful in 6s
PR Checks / Build Backend (no push) (pull_request) Successful in 16s
PR Checks / Build Frontend (no push) (pull_request) Successful in 51s
PR Checks / Build Pipeline (no push) (pull_request) Successful in 36s
PR Checks / AI Code Review (Claude) (pull_request) Failing after 1m20s
The daily/monthly/annual DAG docstring promised an Invalidate Cache step
that never existed — after a marts rebuild the backend kept serving its
startup-cached (possibly empty) DataFrame until a container restart.
Add a POST /api/admin/reload task at the end of each pipeline DAG,
mirroring the sitemap DAG's admin-call pattern.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-09 21:12:01 +01:00
4 changed files with 64 additions and 168 deletions
+11 -68
View File
@@ -4,7 +4,6 @@ Provides efficient queries with caching.
""" """
import logging import logging
import re
import pandas as pd import pandas as pd
import numpy as np import numpy as np
@@ -263,81 +262,25 @@ assert "NULL AS has_sixth_form" in str(_MAIN_QUERY_NO_SIXTH_FORM), (
"expected replacement of 's.has_sixth_form,' to have taken effect" "expected replacement of 's.has_sixth_form,' to have taken effect"
) )
# Fallback used when marts.dim_school predates the GIAS code-dictionary
# migration (i.e. the nightly dbt pipeline hasn't rebuilt the mart yet on
# this DB, so it still has the old name columns instead of *_code columns).
_MAIN_QUERY_LEGACY_NAMES = str(_MAIN_QUERY)
_LEGACY_NAME_REPLACEMENTS = [
("s.phase_code,", "s.phase,"),
("s.school_type_code,", "s.school_type,"),
(
"s.religious_character_code,",
"s.religious_character AS religious_denomination,",
),
("s.status_code,", "s.status,"),
("s.admissions_policy_code,", "s.admissions_policy,"),
]
for _old, _new in _LEGACY_NAME_REPLACEMENTS:
assert _old in _MAIN_QUERY_LEGACY_NAMES, (
f"expected {_old!r} to be present in _MAIN_QUERY before replacement"
)
_MAIN_QUERY_LEGACY_NAMES = _MAIN_QUERY_LEGACY_NAMES.replace(_old, _new)
_MAIN_QUERY_LEGACY_NAMES = text(_MAIN_QUERY_LEGACY_NAMES)
_GIAS_CODE_COLUMN_NAMES = (
"phase_code",
"school_type_code",
"religious_character_code",
"status_code",
"admissions_policy_code",
)
_MISSING_COLUMN_RE = re.compile(r'column "?(?:s\.)?(\w+)"? does not exist')
def _missing_column_name(exc: Exception) -> Optional[str]:
"""Name of the missing column from a psycopg2 UndefinedColumn error.
Inspects exc.orig (the DBAPI error), whose message names only the
offending column — str(exc) also embeds the full SQL statement, which
contains every column name and therefore must not be matched against.
"""
orig = getattr(exc, "orig", None)
match = _MISSING_COLUMN_RE.search(str(orig) if orig is not None else str(exc))
return match.group(1) if match else None
def load_school_data_as_dataframe() -> pd.DataFrame: def load_school_data_as_dataframe() -> pd.DataFrame:
"""Load all school + KS2 data as a pandas DataFrame.""" """Load all school + KS2 data as a pandas DataFrame."""
try: try:
df = pd.read_sql(_MAIN_QUERY, engine) df = pd.read_sql(_MAIN_QUERY, engine)
except sqlalchemy.exc.ProgrammingError as exc: except sqlalchemy.exc.ProgrammingError as exc:
missing = _missing_column_name(exc) if "has_sixth_form" not in str(exc):
if missing in _GIAS_CODE_COLUMN_NAMES:
logging.getLogger(__name__).warning(
"marts predate the GIAS code migration — falling back to "
"legacy name-column query: %s",
exc,
)
try:
df = pd.read_sql(_MAIN_QUERY_LEGACY_NAMES, engine)
except Exception as exc2:
print(f"Warning: Could not load school data from marts: {exc2}")
return pd.DataFrame()
elif missing == "has_sixth_form":
logging.getLogger(__name__).warning(
"marts.dim_school is missing has_sixth_form (pipeline hasn't "
"rebuilt the mart yet on this DB) — retrying without it: %s",
exc,
)
try:
df = pd.read_sql(_MAIN_QUERY_NO_SIXTH_FORM, engine)
except Exception as exc2:
print(f"Warning: Could not load school data from marts: {exc2}")
return pd.DataFrame()
else:
print(f"Warning: Could not load school data from marts: {exc}") print(f"Warning: Could not load school data from marts: {exc}")
return pd.DataFrame() return pd.DataFrame()
logging.getLogger(__name__).warning(
"marts.dim_school is missing has_sixth_form (pipeline hasn't "
"rebuilt the mart yet on this DB) — retrying without it: %s",
exc,
)
try:
df = pd.read_sql(_MAIN_QUERY_NO_SIXTH_FORM, engine)
except Exception as exc2:
print(f"Warning: Could not load school data from marts: {exc2}")
return pd.DataFrame()
except Exception as exc: except Exception as exc:
print(f"Warning: Could not load school data from marts: {exc}") print(f"Warning: Could not load school data from marts: {exc}")
return pd.DataFrame() return pd.DataFrame()
+1 -89
View File
@@ -4,7 +4,7 @@ rest of the backend sees must carry today's name strings."""
import numpy as np import numpy as np
import pandas as pd import pandas as pd
from backend.data_loader import _missing_column_name, translate_gias_code_columns from backend.data_loader import translate_gias_code_columns
from backend.gias_codes import ESTABLISHMENT_STATUS, PHASE_OF_EDUCATION from backend.gias_codes import ESTABLISHMENT_STATUS, PHASE_OF_EDUCATION
@@ -42,91 +42,3 @@ def test_missing_code_columns_are_a_noop():
out = translate_gias_code_columns(df) out = translate_gias_code_columns(df)
assert out.iloc[0]["phase"] == "Primary" assert out.iloc[0]["phase"] == "Primary"
assert out.iloc[0]["status"] == "Open" assert out.iloc[0]["status"] == "Open"
def _fake_exc(orig_message):
"""A stand-in for sqlalchemy.exc.ProgrammingError: str(exc) embeds the
full SQL statement (deliberately containing every column name below, to
prove the matcher doesn't fall back to it), while .orig carries the real
DBAPI error message naming only the offending column."""
exc = Exception(
"SELECT s.phase_code, s.school_type_code, s.religious_character_code, "
"s.status_code, s.admissions_policy_code, s.has_sixth_form FROM ... "
f"[SQL: ...] (Background on this error at: https://...)"
)
exc.orig = Exception(orig_message) if orig_message is not None else None
return exc
def test_missing_column_name_quoted():
assert _missing_column_name(_fake_exc('column "phase_code" does not exist')) == "phase_code"
def test_missing_column_name_unquoted():
assert _missing_column_name(_fake_exc("column phase_code does not exist")) == "phase_code"
def test_missing_column_name_table_prefixed():
assert (
_missing_column_name(_fake_exc("column s.has_sixth_form does not exist"))
== "has_sixth_form"
)
def test_missing_column_name_no_match_returns_none():
assert _missing_column_name(_fake_exc("relation \"marts.dim_school\" does not exist")) is None
def test_load_school_data_survives_premigration_marts(monkeypatch):
"""Real prod state until the nightly pipeline first rebuilds the mart with
the GIAS code columns: marts.dim_school still has the old name columns
(phase, school_type, religious_character, status, admissions_policy)
instead of the new *_code columns. The first query raises UndefinedColumn
on s.phase_code; load_school_data_as_dataframe must retry with the
legacy name-column query rather than swallow the error and return (and
then have load_school_data cache) an empty DataFrame."""
import sqlalchemy.exc
from backend import data_loader
data_loader._df_cache = None
data_loader._df_latest_cache = None
good_df = pd.DataFrame(
[
{
"urn": 1,
"school_name": "Legacy School",
"phase": "Primary",
"school_type": "Academy",
"status": "Open",
}
]
)
calls = []
def fake_read_sql(query, con):
calls.append(query)
if len(calls) == 1:
raise sqlalchemy.exc.ProgrammingError(
statement=str(data_loader._MAIN_QUERY),
params=None,
orig=Exception(
"(psycopg2.errors.UndefinedColumn) column s.phase_code "
"does not exist\nLINE 5: s.phase_code,"
),
)
return good_df.copy()
monkeypatch.setattr(data_loader.pd, "read_sql", fake_read_sql)
try:
df = data_loader.load_school_data_as_dataframe()
finally:
data_loader._df_cache = None
data_loader._df_latest_cache = None
assert len(calls) == 2, "must retry with the legacy name-column query variant"
assert calls[1] is data_loader._MAIN_QUERY_LEGACY_NAMES
assert not df.empty
assert df["phase"].iloc[0] == "Primary"
assert df["status"].iloc[0] == "Open"
+3 -7
View File
@@ -148,14 +148,10 @@ def test_load_school_data_survives_missing_has_sixth_form_column(monkeypatch):
def fake_read_sql(query, con): def fake_read_sql(query, con):
calls.append(query) calls.append(query)
if len(calls) == 1: if len(calls) == 1:
# The statement text still contains phase_code, school_type_code,
# etc. (it's the full _MAIN_QUERY SELECT list) — that's exactly
# the collision this test guards against: matching must be done
# against exc.orig (the DBAPI error), not str(exc)/the statement.
raise sqlalchemy.exc.ProgrammingError( raise sqlalchemy.exc.ProgrammingError(
statement=str(data_loader._MAIN_QUERY), "SELECT ...",
params=None, None,
orig=Exception( Exception(
"(psycopg2.errors.UndefinedColumn) column s.has_sixth_form " "(psycopg2.errors.UndefinedColumn) column s.has_sixth_form "
"does not exist" "does not exist"
), ),
+49 -4
View File
@@ -38,6 +38,31 @@ default_args = {
"retry_delay": timedelta(minutes=5), "retry_delay": timedelta(minutes=5),
} }
# The backend caches the marts DataFrame at startup; after any rebuild the
# cache must be invalidated or the API serves stale (or empty) data until the
# container restarts.
INVALIDATE_CACHE_CMD = """
set -e
BACKEND_URL="${BACKEND_URL:-http://backend:80}"
ADMIN_KEY="${ADMIN_API_KEY:-changeme}"
echo "Calling $BACKEND_URL/api/admin/reload ..."
response=$(curl -s -o /tmp/reload_response.json -w "%{http_code}" \\
--connect-timeout 10 --max-time 120 \\
-X POST "$BACKEND_URL/api/admin/reload" \\
-H "X-API-Key: $ADMIN_KEY" \\
-H "Content-Type: application/json")
echo "HTTP status: $response"
cat /tmp/reload_response.json
if [ "$response" != "200" ]; then
echo "ERROR: backend cache reload failed (HTTP $response)"
exit 1
fi
"""
# ── Daily DAG (GIAS + downstream) ────────────────────────────────────── # ── Daily DAG (GIAS + downstream) ──────────────────────────────────────
@@ -91,7 +116,12 @@ print(f'Validation passed: {{count}} GIAS rows')
bash_command=f"cd {PIPELINE_DIR} && python scripts/sync_typesense.py", bash_command=f"cd {PIPELINE_DIR} && python scripts/sync_typesense.py",
) )
extract_group >> validate_raw >> dbt_build >> sync_typesense invalidate_cache = BashOperator(
task_id="invalidate_cache",
bash_command=INVALIDATE_CACHE_CMD,
)
extract_group >> validate_raw >> dbt_build >> sync_typesense >> invalidate_cache
# ── Monthly DAG (Ofsted) ─────────────────────────────────────────────── # ── Monthly DAG (Ofsted) ───────────────────────────────────────────────
@@ -121,7 +151,12 @@ with DAG(
bash_command=f"cd {PIPELINE_DIR} && python scripts/sync_typesense.py", bash_command=f"cd {PIPELINE_DIR} && python scripts/sync_typesense.py",
) )
extract_ofsted >> dbt_build_ofsted >> sync_typesense_ofsted invalidate_cache_ofsted = BashOperator(
task_id="invalidate_cache",
bash_command=INVALIDATE_CACHE_CMD,
)
extract_ofsted >> dbt_build_ofsted >> sync_typesense_ofsted >> invalidate_cache_ofsted
# ── Annual DAG (EES: KS2, KS4, Census, Admissions) ─────────────────── # ── Annual DAG (EES: KS2, KS4, Census, Admissions) ───────────────────
@@ -153,7 +188,12 @@ with DAG(
bash_command=f"cd {PIPELINE_DIR} && python scripts/sync_typesense.py", bash_command=f"cd {PIPELINE_DIR} && python scripts/sync_typesense.py",
) )
extract_ees_group >> dbt_build_ees >> sync_typesense_ees invalidate_cache_ees = BashOperator(
task_id="invalidate_cache",
bash_command=INVALIDATE_CACHE_CMD,
)
extract_ees_group >> dbt_build_ees >> sync_typesense_ees >> invalidate_cache_ees
# ── Annual DAG (IDACI Deprivation) ──────────────────────────────────── # ── Annual DAG (IDACI Deprivation) ────────────────────────────────────
@@ -178,4 +218,9 @@ with DAG(
bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_idaci+ fact_deprivation+", bash_command=f"cd {PIPELINE_DIR}/transform && {DBT_BIN} build --profiles-dir . --target production --select stg_idaci+ fact_deprivation+",
) )
extract_idaci >> dbt_build_idaci invalidate_cache_idaci = BashOperator(
task_id="invalidate_cache",
bash_command=INVALIDATE_CACHE_CMD,
)
extract_idaci >> dbt_build_idaci >> invalidate_cache_idaci