Skip to content

Commit 1802564

Browse files
OGordon100Isaac
andcommitted
fix(app): carry table and preview flag in staged-config stub
The staged-config stub now includes source_table_fqn, skip_history and run_type, so the Jobs API failure feed can name a staged run's table and exclude staged previews without a SQL lookup. This removes the last warehouse query from the recent-failures endpoints. Co-authored-by: Isaac <no-reply@databricks.com>
1 parent 6f5825d commit 1802564

8 files changed

Lines changed: 69 additions & 109 deletions

File tree

‎app/src/databricks_labs_dqx_app/backend/routes/v1/dryrun.py‎

Lines changed: 2 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,3 @@
1-
import asyncio
21
import json
32
from collections.abc import Callable
43
from typing import Annotated, Any
@@ -241,7 +240,6 @@ async def list_validation_runs(
241240
async def list_recent_validation_failures(
242241
job_svc: Annotated[JobService, Depends(get_job_service)],
243242
user_catalogs: Annotated[frozenset[str], Depends(get_user_catalog_names)],
244-
sql: Annotated[SqlExecutor, Depends(get_sp_sql_executor)],
245243
) -> list[RunFailureOut]:
246244
"""Return recently-failed validation runs, bounded to the most recent *N*.
247245
@@ -253,16 +251,11 @@ async def list_recent_validation_failures(
253251
"""
254252
try:
255253
failed = [run for run in await job_svc.list_recent_failed_runs() if run.task_type in _VALIDATION_TASK_TYPES]
256-
# Oversized configs are staged out of the job parameters, so only the
257-
# run table can name their source table — a rare, failure-only lookup.
258-
staged = [run.app_run_id for run in failed if not run.source_table_fqn]
259-
staged_tables = (
260-
await asyncio.to_thread(job_svc.lookup_source_tables, sql.fqn(_DRYRUN_TABLE), staged) if staged else {}
261-
)
262254

263255
results: list[RunFailureOut] = []
264256
for run in failed:
265-
fqn = run.source_table_fqn or staged_tables.get(run.app_run_id) or ""
257+
# A run with no table cannot pass the catalog filter, so it is skipped.
258+
fqn = run.source_table_fqn or ""
266259
if not fqn or (not fqn.startswith(_SQL_CHECK_PREFIX) and _catalog_of(fqn) not in user_catalogs):
267260
continue
268261
results.append(

‎app/src/databricks_labs_dqx_app/backend/routes/v1/profiler.py‎

Lines changed: 1 addition & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,3 @@
1-
import asyncio
21
import json
32
from typing import Annotated
43
from uuid import uuid4
@@ -262,19 +261,13 @@ async def list_recent_profile_failures(
262261
"""
263262
try:
264263
failed = [run for run in await job_svc.list_recent_failed_runs() if run.task_type == "profile"]
265-
# Oversized configs are staged out of the job parameters, so only the
266-
# run table can name their source table — a rare, failure-only lookup.
267-
staged = [run.app_run_id for run in failed if not run.source_table_fqn]
268-
staged_tables = (
269-
await asyncio.to_thread(job_svc.lookup_source_tables, _run_table_fqn(), staged) if staged else {}
270-
)
271264

272265
results: list[RunFailureOut] = []
273266
for run in failed:
274267
results.append(
275268
RunFailureOut(
276269
run_id=run.app_run_id,
277-
source_table_fqn=run.source_table_fqn or staged_tables.get(run.app_run_id) or "",
270+
source_table_fqn=run.source_table_fqn or "",
278271
status="FAILED",
279272
created_at=run.created_at,
280273
)

‎app/src/databricks_labs_dqx_app/backend/run_config_store.py‎

Lines changed: 15 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,6 +7,10 @@
77
into the ``dq_run_configs`` Lakebase table keyed by ``run_id`` and passes a
88
tiny stub ``{"__manifest__": true}`` instead. The task runner reads the config
99
from Lakebase using the ``run_id``.
10+
11+
The stub also carries the few small config fields the app reads back from
12+
the Jobs API (see *MANIFEST_SUMMARY_KEYS*), so a staged run can still be
13+
attributed to its table and recognised as a preview without a SQL lookup.
1014
"""
1115

1216
import json
@@ -24,6 +28,12 @@
2428
# Marker key in the inline stub passed to the task runner.
2529
MANIFEST_CONFIG_KEY = "__manifest__"
2630

31+
# Small config fields copied into the stub so Jobs API readers
32+
# (services.task_runner_runs) can name a staged run's table and tell previews
33+
# apart without reading the staged config. The runner ignores them: it loads
34+
# the full config from the table whenever the manifest key is set.
35+
MANIFEST_SUMMARY_KEYS = ("source_table_fqn", "skip_history", "run_type")
36+
2737
# Lakebase OLTP table holding run configs that are too large to inline.
2838
RUN_CONFIGS_TABLE = "dq_run_configs"
2939

@@ -91,9 +101,10 @@ def build_inline_config_payload(config: dict[str, Any]) -> str:
91101
return _compact_json(config)
92102

93103

94-
def build_manifest_config_payload() -> str:
95-
"""Serialize the stub the task runner uses to load a staged config from the table."""
96-
return _compact_json({MANIFEST_CONFIG_KEY: True})
104+
def build_manifest_config_payload(config: dict[str, Any]) -> str:
105+
"""Serialize the stub the task runner uses to load *config* back from the table."""
106+
summary = {key: config[key] for key in MANIFEST_SUMMARY_KEYS if key in config}
107+
return _compact_json({MANIFEST_CONFIG_KEY: True, **summary})
97108

98109

99110
def stage_config_to_table(sql: OltpExecutorProtocol, run_id: str, config: dict[str, Any]) -> None:
@@ -163,7 +174,7 @@ def prepare_config_json(
163174
raise
164175
except Exception as exc:
165176
raise RunConfigStagingError(run_id, sql.fqn(RUN_CONFIGS_TABLE), exc) from exc
166-
manifest = build_manifest_config_payload()
177+
manifest = build_manifest_config_payload(config)
167178
manifest_params = {**job_parameters_without_config, "config_json": manifest}
168179
manifest_size = job_parameters_size(manifest_params)
169180
if manifest_size > JOB_PARAMETERS_CHAR_LIMIT:

‎app/src/databricks_labs_dqx_app/backend/services/job_service.py‎

Lines changed: 1 addition & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@
1717
)
1818
from databricks_labs_dqx_app.backend.services.task_runner_runs import TaskRunnerRun, cached_recent_completed_runs
1919
from databricks_labs_dqx_app.backend.sql_executor import OltpExecutorProtocol, SqlExecutor
20-
from databricks_labs_dqx_app.backend.sql_utils import escape_sql_string
2120

2221
logger = logging.getLogger(__name__)
2322

@@ -103,7 +102,7 @@ def submit_run(
103102
config=config,
104103
job_parameters_without_config=base_params,
105104
)
106-
staged = config_json == build_manifest_config_payload()
105+
staged = config_json == build_manifest_config_payload(config)
107106

108107
try:
109108
run = self._ws.jobs.run_now(
@@ -149,20 +148,6 @@ async def list_recent_failed_runs(self) -> list[TaskRunnerRun]:
149148
runs = await cached_recent_completed_runs(self._ws, self._job_id)
150149
return [run for run in runs if run.is_failed and not run.is_preview]
151150

152-
def lookup_source_tables(self, table: str, run_ids: list[str]) -> dict[str, str]:
153-
"""Map *run_ids* to their ``source_table_fqn`` from the run table *table*.
154-
155-
Only needed for runs whose config was staged out of the job parameters
156-
(oversized configs), so the Jobs API alone cannot name their table.
157-
"""
158-
if not run_ids:
159-
return {}
160-
in_list = ", ".join(f"'{escape_sql_string(run_id)}'" for run_id in run_ids)
161-
rows = self._sql.query(
162-
f"SELECT DISTINCT run_id, source_table_fqn FROM {table} WHERE run_id IN ({in_list})" # noqa: S608
163-
)
164-
return {row[0]: row[1] for row in rows if row and row[0] and len(row) > 1 and row[1]}
165-
166151
def get_run_creator(self, job_run_id: int) -> str | None:
167152
"""Return the requesting end-user for a job run, or None if unavailable.
168153

‎app/tests/test_job_service.py‎

Lines changed: 0 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -159,22 +159,3 @@ async def test_recent_failed_runs_empty_without_a_job(sql_executor_mock: MagicMo
159159

160160
assert await service.list_recent_failed_runs() == []
161161
ws.jobs.list_runs.assert_not_called()
162-
163-
164-
def test_lookup_source_tables_maps_run_ids(sql_executor_mock: MagicMock) -> None:
165-
service, _ws = _job_service(sql_executor_mock)
166-
sql_executor_mock.query.return_value = [("r1", "main.s.t1"), ("r2", None)]
167-
168-
result = service.lookup_source_tables("dqx.dqx_studio.dq_validation_runs", ["r1", "r2", "o'brien"])
169-
170-
assert result == {"r1": "main.s.t1"}
171-
statement = sql_executor_mock.query.call_args.args[0]
172-
assert "FROM dqx.dqx_studio.dq_validation_runs" in statement
173-
assert "IN ('r1', 'r2', 'o''brien')" in statement
174-
175-
176-
def test_lookup_source_tables_skips_the_query_for_no_runs(sql_executor_mock: MagicMock) -> None:
177-
service, _ws = _job_service(sql_executor_mock)
178-
179-
assert service.lookup_source_tables("dqx.dqx_studio.dq_validation_runs", []) == {}
180-
sql_executor_mock.query.assert_not_called()

‎app/tests/test_recent_failures.py‎

Lines changed: 24 additions & 53 deletions
Original file line numberDiff line numberDiff line change
@@ -8,8 +8,7 @@
88
1. Only the endpoint's own task types are returned (validation vs. profile).
99
2. The validation feed honours the caller's catalog access.
1010
3. The result is bounded at _RECENT_FAILURES_LIMIT with the minimal RunFailureOut shape.
11-
4. The run table is queried only for runs whose config was staged out of the
12-
job parameters, and never otherwise.
11+
4. Neither feed ever queries the SQL warehouse.
1312
"""
1413

1514
from unittest.mock import MagicMock, create_autospec
@@ -27,7 +26,6 @@
2726
)
2827
from databricks_labs_dqx_app.backend.services.job_service import JobService
2928
from databricks_labs_dqx_app.backend.services.task_runner_runs import TaskRunnerRun
30-
from databricks_labs_dqx_app.backend.sql_executor import SqlExecutor
3129

3230

3331
def _failed_run(
@@ -50,67 +48,58 @@ def _failed_run(
5048

5149
@pytest.fixture
5250
def job_service_mock() -> MagicMock:
53-
svc = create_autospec(JobService, instance=True)
54-
svc.lookup_source_tables.return_value = {}
55-
return svc
51+
return create_autospec(JobService, instance=True)
5652

5753

58-
@pytest.fixture
59-
def sql_executor() -> MagicMock:
60-
sql = create_autospec(SqlExecutor, instance=True)
61-
sql.fqn.side_effect = lambda name: f"dqx.dqx_studio.{name}"
62-
return sql
63-
64-
65-
async def _validation(job_svc: MagicMock, sql: MagicMock, catalogs: frozenset[str] = frozenset({"main"})):
66-
return await list_recent_validation_failures(job_svc=job_svc, user_catalogs=catalogs, sql=sql)
54+
async def _validation(job_svc: MagicMock, catalogs: frozenset[str] = frozenset({"main"})):
55+
return await list_recent_validation_failures(job_svc=job_svc, user_catalogs=catalogs)
6756

6857

6958
class TestListRecentValidationFailures:
70-
async def test_returns_only_validation_task_types(self, job_service_mock, sql_executor):
59+
async def test_returns_only_validation_task_types(self, job_service_mock):
7160
job_service_mock.list_recent_failed_runs.return_value = [
7261
_failed_run("run-dryrun"),
7362
_failed_run("run-scheduled", task_type="scheduled"),
7463
_failed_run("run-profile", task_type="profile"),
7564
]
7665

77-
result = await _validation(job_service_mock, sql_executor)
66+
result = await _validation(job_service_mock)
7867

7968
assert [r.run_id for r in result] == ["run-dryrun", "run-scheduled"]
8069
assert all(r.status == "FAILED" for r in result)
8170

82-
async def test_excludes_runs_from_inaccessible_catalogs(self, job_service_mock, sql_executor):
71+
async def test_excludes_runs_from_inaccessible_catalogs(self, job_service_mock):
8372
job_service_mock.list_recent_failed_runs.return_value = [
8473
_failed_run("run-visible", fqn="main.public.orders"),
8574
_failed_run("run-hidden", fqn="restricted.public.orders"),
8675
]
8776

88-
result = await _validation(job_service_mock, sql_executor)
77+
result = await _validation(job_service_mock)
8978

9079
assert [r.run_id for r in result] == ["run-visible"]
9180

92-
async def test_includes_sql_check_prefix_runs(self, job_service_mock, sql_executor):
81+
async def test_includes_sql_check_prefix_runs(self, job_service_mock):
9382
job_service_mock.list_recent_failed_runs.return_value = [
9483
_failed_run("run-sql", fqn="__sql_check__/my_check"),
9584
]
9685

97-
result = await _validation(job_service_mock, sql_executor, catalogs=frozenset())
86+
result = await _validation(job_service_mock, catalogs=frozenset())
9887

9988
assert [r.source_table_fqn for r in result] == ["__sql_check__/my_check"]
10089

101-
async def test_result_bounded_at_limit(self, job_service_mock, sql_executor):
90+
async def test_result_bounded_at_limit(self, job_service_mock):
10291
job_service_mock.list_recent_failed_runs.return_value = [
10392
_failed_run(f"f-{i}") for i in range(DRYRUN_LIMIT + 10)
10493
]
10594

106-
result = await _validation(job_service_mock, sql_executor)
95+
result = await _validation(job_service_mock)
10796

10897
assert len(result) == DRYRUN_LIMIT
10998

110-
async def test_returns_minimal_fields_only(self, job_service_mock, sql_executor):
99+
async def test_returns_minimal_fields_only(self, job_service_mock):
111100
job_service_mock.list_recent_failed_runs.return_value = [_failed_run("run-failed")]
112101

113-
result = await _validation(job_service_mock, sql_executor)
102+
result = await _validation(job_service_mock)
114103

115104
assert result == [
116105
RunFailureOut(
@@ -121,39 +110,22 @@ async def test_returns_minimal_fields_only(self, job_service_mock, sql_executor)
121110
)
122111
]
123112

124-
async def test_never_queries_the_warehouse_when_configs_are_inline(self, job_service_mock, sql_executor):
125-
job_service_mock.list_recent_failed_runs.return_value = [_failed_run("run-failed")]
126-
127-
await _validation(job_service_mock, sql_executor)
128-
129-
job_service_mock.lookup_source_tables.assert_not_called()
130-
sql_executor.query.assert_not_called()
131-
sql_executor.query_dicts.assert_not_called()
132-
133-
async def test_staged_config_runs_resolve_their_table_from_the_run_table(self, job_service_mock, sql_executor):
113+
async def test_runs_without_a_source_table_are_skipped(self, job_service_mock):
114+
# The catalog filter cannot be applied to a run with no table, so it
115+
# is dropped rather than shown to everyone.
134116
job_service_mock.list_recent_failed_runs.return_value = [
135-
_failed_run("run-inline"),
136-
_failed_run("run-staged", fqn=None),
117+
_failed_run("run-known"),
137118
_failed_run("run-unknown", fqn=None),
138119
]
139-
job_service_mock.lookup_source_tables.return_value = {"run-staged": "main.public.big"}
140-
141-
result = await _validation(job_service_mock, sql_executor)
142-
143-
job_service_mock.lookup_source_tables.assert_called_once_with(
144-
"dqx.dqx_studio.dq_validation_runs", ["run-staged", "run-unknown"]
145-
)
146-
# A staged run whose table cannot be resolved is dropped (the catalog
147-
# filter cannot be applied to it).
148-
assert [(r.run_id, r.source_table_fqn) for r in result] == [
149-
("run-inline", "main.public.orders"),
150-
("run-staged", "main.public.big"),
151-
]
152120

153-
async def test_empty_list_when_no_failures(self, job_service_mock, sql_executor):
121+
result = await _validation(job_service_mock)
122+
123+
assert [r.run_id for r in result] == ["run-known"]
124+
125+
async def test_empty_list_when_no_failures(self, job_service_mock):
154126
job_service_mock.list_recent_failed_runs.return_value = []
155127

156-
assert await _validation(job_service_mock, sql_executor) == []
128+
assert await _validation(job_service_mock) == []
157129

158130

159131
class TestListRecentProfileFailures:
@@ -189,7 +161,6 @@ async def test_returns_minimal_fields_only(self, job_service_mock):
189161
created_at="2026-09-21T14:13:20+00:00",
190162
)
191163
]
192-
job_service_mock.lookup_source_tables.assert_not_called()
193164

194165
async def test_empty_list_when_no_failures(self, job_service_mock):
195166
job_service_mock.list_recent_failed_runs.return_value = []

‎app/tests/test_run_config_store.py‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -92,6 +92,25 @@ def test_stages_when_over_limit(self) -> None:
9292
assert kwargs["key_cols"] == {"run_id": "run123"}
9393
assert json.loads(kwargs["value_cols"]["config"]) == config
9494

95+
def test_stub_carries_the_fields_jobs_api_readers_need(self) -> None:
96+
# The failure toast reads a staged run's table and preview flag from
97+
# the job parameters, so they must survive staging (no SQL lookup).
98+
config = {**_big_config(), "source_table_fqn": "main.sales.orders", "skip_history": True, "run_type": "preview"}
99+
100+
result = prepare_config_json(
101+
_oltp_mock(),
102+
run_id="run123",
103+
config=config,
104+
job_parameters_without_config=_base_params(),
105+
)
106+
107+
assert json.loads(result) == {
108+
MANIFEST_CONFIG_KEY: True,
109+
"source_table_fqn": "main.sales.orders",
110+
"skip_history": True,
111+
"run_type": "preview",
112+
}
113+
95114
def test_stub_stays_within_the_job_parameter_limit(self) -> None:
96115
# Regression guard: the stub plus base params must always fit, so a
97116
# staged config never re-trips the limit it was meant to dodge.

‎app/tests/test_task_runner_runs.py‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,13 @@ def test_staged_config_has_no_source_table(self):
6363

6464
assert run is not None and run.source_table_fqn is None
6565

66+
def test_staged_config_stub_still_names_its_table_and_preview_flag(self):
67+
run = parse_task_runner_run(
68+
make_run(config={"__manifest__": True, "source_table_fqn": "main.sales.big", "skip_history": True})
69+
)
70+
71+
assert run is not None and run.source_table_fqn == "main.sales.big" and run.is_preview
72+
6673
def test_unparseable_config_has_no_source_table(self):
6774
run = parse_task_runner_run(make_run(config="not json"))
6875

0 commit comments

Comments
 (0)