Skip to content

Commit 4fb8953

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 cd407eb commit 4fb8953

8 files changed

Lines changed: 69 additions & 108 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`` Delta table keyed by ``run_id`` and passes a
88
tiny stub ``{"__manifest__": true}`` instead. The task runner reads the config
99
from the table using the ``run_id`` and deletes the row once the run finishes.
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 base64
@@ -25,6 +29,12 @@
2529
# Marker key in the inline stub passed to the task runner.
2630
MANIFEST_CONFIG_KEY = "__manifest__"
2731

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

@@ -78,9 +88,10 @@ def build_inline_config_payload(config: dict[str, Any]) -> str:
7888
return _compact_json(config)
7989

8090

81-
def build_manifest_config_payload() -> str:
82-
"""Serialize the stub the task runner uses to load a staged config from the table."""
83-
return _compact_json({MANIFEST_CONFIG_KEY: True})
91+
def build_manifest_config_payload(config: dict[str, Any]) -> str:
92+
"""Serialize the stub the task runner uses to load *config* back from the table."""
93+
summary = {key: config[key] for key in MANIFEST_SUMMARY_KEYS if key in config}
94+
return _compact_json({MANIFEST_CONFIG_KEY: True, **summary})
8495

8596

8697
def stage_config_to_table(sql: SqlExecutor, run_id: str, config: dict[str, Any]) -> None:
@@ -145,7 +156,7 @@ def prepare_config_json(
145156
raise
146157
except Exception as exc:
147158
raise RunConfigStagingError(run_id, sql.fqn(RUN_CONFIGS_TABLE), exc) from exc
148-
manifest = build_manifest_config_payload()
159+
manifest = build_manifest_config_payload(config)
149160
manifest_params = {**job_parameters_without_config, "config_json": manifest}
150161
manifest_size = job_parameters_size(manifest_params)
151162
if manifest_size > JOB_PARAMETERS_CHAR_LIMIT:

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

Lines changed: 1 addition & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -79,7 +79,7 @@ def submit_run(
7979
config=config,
8080
job_parameters_without_config=base_params,
8181
)
82-
staged = config_json == build_manifest_config_payload()
82+
staged = config_json == build_manifest_config_payload(config)
8383

8484
try:
8585
run = self._ws.jobs.run_now(
@@ -125,20 +125,6 @@ async def list_recent_failed_runs(self) -> list[TaskRunnerRun]:
125125
runs = await cached_recent_completed_runs(self._ws, self._job_id)
126126
return [run for run in runs if run.is_failed and not run.is_preview]
127127

128-
def lookup_source_tables(self, table: str, run_ids: list[str]) -> dict[str, str]:
129-
"""Map *run_ids* to their ``source_table_fqn`` from the run table *table*.
130-
131-
Only needed for runs whose config was staged out of the job parameters
132-
(oversized configs), so the Jobs API alone cannot name their table.
133-
"""
134-
if not run_ids:
135-
return {}
136-
in_list = ", ".join(f"'{escape_sql_string(run_id)}'" for run_id in run_ids)
137-
rows = self._sql.query(
138-
f"SELECT DISTINCT run_id, source_table_fqn FROM {table} WHERE run_id IN ({in_list})" # noqa: S608
139-
)
140-
return {row[0]: row[1] for row in rows if row and row[0] and len(row) > 1 and row[1]}
141-
142128
def get_run_creator(self, job_run_id: int) -> str | None:
143129
"""Return the requesting end-user for a job run, or None if unavailable.
144130

‎app/tests/test_job_service.py‎

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

169169
assert await service.list_recent_failed_runs() == []
170170
ws.jobs.list_runs.assert_not_called()
171-
172-
173-
def test_lookup_source_tables_maps_run_ids(sql_executor_mock: MagicMock) -> None:
174-
service, _ws = _job_service(sql_executor_mock)
175-
sql_executor_mock.query.return_value = [("r1", "main.s.t1"), ("r2", None)]
176-
177-
result = service.lookup_source_tables("dqx.dqx_studio.dq_validation_runs", ["r1", "r2", "o'brien"])
178-
179-
assert result == {"r1": "main.s.t1"}
180-
statement = sql_executor_mock.query.call_args.args[0]
181-
assert "FROM dqx.dqx_studio.dq_validation_runs" in statement
182-
assert "IN ('r1', 'r2', 'o''brien')" in statement
183-
184-
185-
def test_lookup_source_tables_skips_the_query_for_no_runs(sql_executor_mock: MagicMock) -> None:
186-
service, _ws = _job_service(sql_executor_mock)
187-
188-
assert service.lookup_source_tables("dqx.dqx_studio.dq_validation_runs", []) == {}
189-
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
@@ -95,6 +95,25 @@ def test_stages_when_over_limit(self) -> None:
9595
assert kwargs["key_cols"] == {"run_id": "run123"}
9696
assert _staged_config(kwargs["value_cols"]["config"]) == config
9797

98+
def test_stub_carries_the_fields_jobs_api_readers_need(self) -> None:
99+
# The failure toast reads a staged run's table and preview flag from
100+
# the job parameters, so they must survive staging (no SQL lookup).
101+
config = {**_big_config(), "source_table_fqn": "main.sales.orders", "skip_history": True, "run_type": "preview"}
102+
103+
result = prepare_config_json(
104+
_sql_mock(),
105+
run_id="run123",
106+
config=config,
107+
job_parameters_without_config=_base_params(),
108+
)
109+
110+
assert json.loads(result) == {
111+
MANIFEST_CONFIG_KEY: True,
112+
"source_table_fqn": "main.sales.orders",
113+
"skip_history": True,
114+
"run_type": "preview",
115+
}
116+
98117
def test_stub_stays_within_the_job_parameter_limit(self) -> None:
99118
# Regression guard: the stub plus base params must always fit, so a
100119
# 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)