Skip to content
Merged
Show file tree
Hide file tree
Changes from 14 commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
87ae0c0
Use a Delta table to store oversized run configs
ghanse Sep 23, 2026
477b9b2
Fix formatting and typing
ghanse Sep 23, 2026
ed255cb
Merge remote-tracking branch 'origin/main' into studio/run-config-table
ghanse Sep 23, 2026
2784b81
Use Lakebase to stage large job run config
ghanse Sep 28, 2026
6071bc5
Format tests
ghanse Sep 28, 2026
e6b0eac
Merge branch 'main' into dqx-studio-run-config-table
ghanse Sep 28, 2026
99c5735
Add OLTP resource to startup and seed_demo
ghanse Sep 28, 2026
548fbe1
Merge branch 'dqx-studio-run-config-table' of https://github.com/data…
ghanse Sep 28, 2026
14ed6c5
Merge branch 'main' into dqx-studio-run-config-table
ghanse Sep 28, 2026
df21cb0
Move Lakebase connection into shared helper
ghanse Sep 28, 2026
dffaf5d
Merge branch 'dqx-studio-run-config-table' of https://github.com/data…
ghanse Sep 28, 2026
81bf3a0
Merge branch 'main' into dqx-studio-run-config-table
ghanse Sep 28, 2026
c82b8ca
Handle grant failures
ghanse Sep 28, 2026
db292d0
Merge branch 'dqx-studio-run-config-table' of https://github.com/data…
ghanse Sep 28, 2026
76b4d02
Handle flaky LLM-generated rules
ghanse Sep 28, 2026
e3fa811
fix(app): fail fast on oversized run configs when Lakebase is disable…
mwojtyczka Sep 28, 2026
b44c444
test(app): pin the failing-submit JobService OLTP mock to the postgre…
mwojtyczka Sep 29, 2026
e4b3373
fix(app): harden run-config staging error handling and coordinate wiring
mwojtyczka Sep 29, 2026
487db79
style(app): reformat test_monitored_tables_routes.py with black
mwojtyczka Sep 29, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 6 additions & 6 deletions .build-constraints.txt
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
hatch-fancy-pypi-readme==25.1.0 \
--hash=sha256:9c58ed3dff90d51f43414ce37009ad1d5b0f08ffc9fc216998a06380f01c0045 \
--hash=sha256:ce0134c40d63d874ac48f48ccc678b8f3b62b8e50e9318520d2bffc752eedaf3
hatchling==1.32.0 \
--hash=sha256:0bdbde4a52b06c37e3eca395f85a762bf0ef06fe374fd8ae429dc6be10230f5f \
--hash=sha256:0e17c9c3b9aa7c625acc8d0f5b622f107d5049af9ecf5ada4de1aada5be7cdbc
hatchling==1.32.4 \
--hash=sha256:08ecf7548fb48205e7f213d70c71e67b8271b7242093dc3f1da578b42c734a2c \
--hash=sha256:c4468f73144c054d2aab4ef0f0378c43b9878bf07f8ffd6b79690e970d375f07
# via hatch-fancy-pypi-readme
packaging==26.3 \
--hash=sha256:94edc256424af38762eb31306eed28beb9f0efc50a8837492c9d6fd6004aed79 \
Expand All @@ -21,7 +21,7 @@ tomlkit==0.15.1 \
--hash=sha256:177a05aece5a8ca5266fd3c448abb47b8d352f09d477d3ca8332db4d89b24304 \
--hash=sha256:e25bbf38843005246210a12982776f27f99cb9be67160e14434d0c0d21ee1e97
# via hatchling
trove-classifiers==2026.6.1.19 \
--hash=sha256:ab4c4ec93cc4a4e7815fa759906e05e6bb3f2fbd92ea0f897288c6a43efd15b3 \
--hash=sha256:c5132b4b61a829d11cfbd2d72e97f20a45ed6edb95e45c5efdeb5e00836b2745
trove-classifiers==2026.9.21.13 \
--hash=sha256:0a9ebc8d4e2f3e8a22848c5258033035bec17a3012ac3fea16dbaa764489eb71 \
--hash=sha256:8b1ff4f9c191b1040b71c37f1e445ab99732911e3cd91de52838453a854d7a17
# via hatchling
6 changes: 3 additions & 3 deletions app/.build-constraints.txt
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@ pluggy==1.6.0 \
--hash=sha256:7dcc130b76258d33b90f61b658791dede3486c3e6bfb003ee5c9bfb396dd22f3 \
--hash=sha256:e920276dd6813095e9377c0bc5566d94c932c33b27a3e3945d8389c374dd4746
# via hatchling
trove-classifiers==2026.6.1.19 \
--hash=sha256:ab4c4ec93cc4a4e7815fa759906e05e6bb3f2fbd92ea0f897288c6a43efd15b3 \
--hash=sha256:c5132b4b61a829d11cfbd2d72e97f20a45ed6edb95e45c5efdeb5e00836b2745
trove-classifiers==2026.9.21.13 \
--hash=sha256:0a9ebc8d4e2f3e8a22848c5258033035bec17a3012ac3fea16dbaa764489eb71 \
--hash=sha256:8b1ff4f9c191b1040b71c37f1e445ab99732911e3cd91de52838453a854d7a17
# via hatchling
10 changes: 6 additions & 4 deletions app/AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -82,16 +82,18 @@ but is protected transitively by the instance-level guard.
│ ├── dq_schedule_configs (OLTP*) per-schedule config (cron/interval, target rules)
│ ├── dq_schedule_configs_history (OLTP*) schedule config change audit log
│ ├── dq_schedule_runs (OLTP*) scheduler last/next run state (survives restarts)
│ ├── dq_run_configs (OLTP*) staged run configs too large to inline in job params
│ └── dq_migrations (Delta) Delta migration version tracker
├── dqx_studio_tmp ← temp views created via OBO for profiler/dryrun jobs
└── dqx_studio.wheels (volume) ← DQX + task-runner wheels uploaded at app startup

Lakebase project (when enabled, default `lakebase_project_id` = `dqx-studio-db`):
└── databricks_postgres (database — always-present admin DB; no per-app DB provisioned)
└── dqx_studio (schema — created by PgMigrationRunner on first start; configurable via DQX_LAKEBASE_SCHEMA)
├── dq_app_settings, dq_role_mappings, dq_resolved_rules,
│ dq_resolved_rules_history, dq_comments, dq_schedule_configs,
│ dq_schedule_configs_history, dq_schedule_runs
├── dq_app_settings, dq_role_mappings, dq_quality_rules,
| dq_resolved_rules, dq_quality_rules_history, dq_comments,
| dq_schedule_configs, dq_schedule_configs_history, dq_schedule_runs,
| dq_run_configs
└── dq_migrations (Postgres migration version tracker)
```

Expand Down Expand Up @@ -512,7 +514,7 @@ choice is driven entirely by `databricks.yml`:
| Backend | Tables | Why |
|---------|--------|-----|
| **Delta Lake** (always) | `dq_validation_runs`, `dq_profiling_results`, `dq_quarantine_records`, `dq_metrics` | Spark task runner writes these; high-volume append-mostly; columnar reads. |
| **Lakebase Postgres** *(default — opt-out via `lakebase_endpoint="-"`)* | `dq_app_settings`, `dq_role_mappings`, `dq_resolved_rules`, `dq_resolved_rules_history`, `dq_comments`, `dq_schedule_configs`, `dq_schedule_configs_history`, `dq_schedule_runs` | Low-latency point reads/writes from FastAPI request handlers; row-level upserts; primary-key/foreign-key semantics. |
| **Lakebase Postgres** *(default — opt-out via `lakebase_endpoint="-"`)* | `dq_app_settings`, `dq_role_mappings`, `dq_quality_rules`, `dq_resolved_rules`, `dq_resolved_rules_history`, `dq_comments`, `dq_schedule_configs`, `dq_schedule_configs_history`, `dq_schedule_runs`, `dq_run_configs` | Low-latency point reads/writes from FastAPI request handlers; row-level upserts; primary-key/foreign-key semantics. |

When Lakebase is **disabled** (no `lakebase_endpoint` set), the OLTP
tables fall back to Delta — `MigrationRunner` runs both
Expand Down
42 changes: 42 additions & 0 deletions app/databricks.yml
Comment thread
mwojtyczka marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -149,6 +149,10 @@ variables:
value: "${var.lakebase_database_name}"
- name: "DQX_LAKEBASE_SCHEMA"
value: "${var.lakebase_schema_name}"
# Task-runner SP's Postgres role — granted read/delete on dq_run_configs
# at startup so the runner can read staged (oversized) run configs.
- name: "DQX_TASK_RUNNER_POSTGRES_ROLE"
value: "${var.dqx_service_principal_application_id}"
# Starter Insights dashboard; admins can override at runtime via Configuration.
- name: "DQX_DEFAULT_DASHBOARD_ID"
value: "${resources.dashboards.dqx_quality_overview.id}"
Expand Down Expand Up @@ -412,6 +416,16 @@ resources:
bypassrls: true
membership_roles:
- DATABRICKS_SUPERUSER
# Login role for the task-runner SP (a different SP than the app). It reads
# and deletes staged run configs in ``dq_run_configs``; the app grants it
# USAGE + SELECT/DELETE at startup (see _grant_task_runner_run_config_access).
# Least privilege: no createdb/createrole/superuser — only the explicit grant.
task_runner_sp:
parent: "projects/${var.lakebase_project_id}/branches/${var.lakebase_branch}"
role_id: "sp-${var.dqx_service_principal_application_id}"
postgres_role: ${var.dqx_service_principal_application_id}
identity_type: SERVICE_PRINCIPAL
auth_method: LAKEBASE_OAUTH_V1

dashboards:
dqx_quality_overview:
Expand Down Expand Up @@ -455,6 +469,18 @@ resources:
- "{{job.parameters.run_id}}"
- "--requesting_user"
- "{{job.parameters.requesting_user}}"
- "--lakebase_endpoint"
- "{{job.parameters.lakebase_endpoint}}"
- "--lakebase_database"
- "{{job.parameters.lakebase_database}}"
- "--lakebase_schema"
- "{{job.parameters.lakebase_schema}}"
- "--lakebase_host"
- "{{job.parameters.lakebase_host}}"
- "--lakebase_port"
- "{{job.parameters.lakebase_port}}"
- "--lakebase_username"
- "{{job.parameters.lakebase_username}}"
environment_key: "default"
environments:
- environment_key: "default"
Expand Down Expand Up @@ -482,6 +508,22 @@ resources:
default: ""
- name: "requesting_user"
default: "unknown"
# Lakebase settings for reading staged run configs. ``submit_run`` always
# passes the app-resolved values; these defaults only apply to a manual
# job run. ``host`` is resolved at runtime so it has no static default;
# ``username`` is the task-runner SP's Postgres role (its client id).
- name: "lakebase_endpoint"
default: "${var.lakebase_endpoint}"
- name: "lakebase_database"
default: "${var.lakebase_database_name}"
- name: "lakebase_schema"
default: "${var.lakebase_schema_name}"
- name: "lakebase_host"
default: ""
- name: "lakebase_port"
default: "5432"
- name: "lakebase_username"
default: "${var.dqx_service_principal_application_id}"

# Deploy targets are defined in the untracked ``target.dev.yml`` (see the
# ``include`` above and ``target.dev.yml.example``).
1 change: 1 addition & 0 deletions app/scripts/seed_demo.py
Comment thread
mwojtyczka marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -276,6 +276,7 @@ def main() -> int:
ws=ws,
job_id=conf.job_id,
sql=sp_sql,
oltp_sql=oltp,
warehouse_id=warehouse_id,
)
binding_run = BindingRunService(
Expand Down
9 changes: 9 additions & 0 deletions app/src/databricks_labs_dqx_app/backend/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,15 @@ def validate_admin_group(cls, value: str) -> str:
validation_alias="DQX_LAKEBASE_SCHEMA",
description="Postgres schema for app tables. Created at startup if missing.",
)
# The task-runner job runs as a separate service principal and reads staged
# run configs from ``dq_run_configs`` over Postgres. Its Postgres role (the
# SP client id) is granted USAGE + SELECT/DELETE at startup by the app (the
# table owner). Empty when there is no separate runner SP.
task_runner_postgres_role: str = Field(
default="",
validation_alias="DQX_TASK_RUNNER_POSTGRES_ROLE",
description="Postgres role (service principal client id for the task runner). Granted read/delete on dq_run_configs.",
)
# Default 0 so the pool can drain to zero idle connections and let a
# scale-to-zero Lakebase endpoint suspend. A held-open connection (min_size
# >= 1) is periodically re-established after suspension kills it, nudging
Expand Down
22 changes: 21 additions & 1 deletion app/src/databricks_labs_dqx_app/backend/dependencies.py
Original file line number Diff line number Diff line change
Expand Up @@ -781,6 +781,7 @@ def get_check_validator() -> Callable[[list[Any]], ChecksValidationStatus]:
async def get_job_service(
sp_ws: Annotated[WorkspaceClient, Depends(get_sp_ws)],
sql: Annotated[SqlExecutor, Depends(get_sp_sql_executor)],
oltp: Annotated[OltpExecutorProtocol, Depends(get_sp_oltp_executor)],
app_settings: Annotated[AppSettingsService, Depends(get_app_settings_service)],
) -> JobService:
"""Create a JobService using app (SP) credentials.
Expand All @@ -789,13 +790,32 @@ async def get_job_service(
admin-configured SQL warehouse (``dq_app_settings``) is resolved here and
threaded into the submitted run so the task runner's temp-view cleanup path
honours it (env fallback when unset).

Oversized run configs are staged in the ``dq_run_configs`` Lakebase table via
the OLTP executor; the Lakebase connection settings are passed to the runner
as job parameters so tasks can read the staged configs.
"""
lakebase = rt.require_resources().lakebase
# Prefer the endpoint, host, and port the live OLTP executor already
# resolved, falling back to the configured connection values.
resolved_endpoint = getattr(oltp, "endpoint", None) or lakebase.endpoint or ""
resolved_host = getattr(oltp, "host", None) or lakebase.host or ""
resolved_port = getattr(oltp, "port", None) or lakebase.port or 5432
resolved_username = (
conf.task_runner_postgres_role.strip() or getattr(oltp, "username", None) or lakebase.username or ""
)
return JobService(
ws=sp_ws,
job_id=str(_require_resolved_job_id()),
sql=sql,
oltp_sql=oltp,
warehouse_id=resolve_warehouse_id(app_settings),
wheels_volume=rt.require_resources().volume.path,
lakebase_endpoint=resolved_endpoint,
lakebase_database=lakebase.database or "",
lakebase_schema=lakebase.schema or "",
lakebase_host=resolved_host,
lakebase_port=resolved_port,
lakebase_username=resolved_username,
)


Expand Down
18 changes: 18 additions & 0 deletions app/src/databricks_labs_dqx_app/backend/migrations/postgres.py
Original file line number Diff line number Diff line change
Expand Up @@ -735,6 +735,24 @@ class PgMigration:
" updated_at TIMESTAMPTZ"
");"
# ----------------------------------------------------------
# dq_run_configs — staging table for task-runner run configs
# too large to inline in the job's ``config_json`` parameter.
# Databricks caps job parameters at 10,000 characters. The app
# upserts the fully-resolved config keyed by ``run_id`` and
# passes a tiny ``{"__manifest__": true}`` stub; the serverless
# task runner reads the row back over Postgres by ``run_id``,
# then deletes it. ``run_id`` is a real enforced PRIMARY KEY so
# a resubmit replaces the row and the reader always sees exactly
# one. ``config`` is the compact JSON payload stored as text
# (same convention as ``dq_schedule_configs.config_json``); the
# daily retention sweep is the backstop for rows a crashed run
# left behind.
# ----------------------------------------------------------
f"CREATE TABLE IF NOT EXISTS {_S}.dq_run_configs ("
" run_id TEXT PRIMARY KEY,"
" config TEXT NOT NULL,"
" created_at TIMESTAMPTZ NOT NULL"
");"
# dq_rules_core — dqx-core-compatible READ view over the
# APPROVED rules in ``dq_resolved_rules``. It reshapes Studio's
# storage into exactly the columns dqx-core's
Expand Down
15 changes: 15 additions & 0 deletions app/src/databricks_labs_dqx_app/backend/pg_executor.py
Original file line number Diff line number Diff line change
Expand Up @@ -390,6 +390,21 @@ def username(self) -> str:
def database(self) -> str:
return self._database

@property
def endpoint(self) -> str | None:
"""Resolved Lakebase endpoint path, or ``None`` for a static-password connection."""
return self._endpoint

@property
def host(self) -> str:
"""Resolved read/write Postgres host."""
return self._host

@property
def port(self) -> int:
"""Resolved Postgres port (honours ``PGPORT`` in platform-bound mode)."""
return self._port

# ------------------------------------------------------------------
# Token-refresh observability
# ------------------------------------------------------------------
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,7 @@
MissingSnapshotError,
NeverApprovedError,
)
from databricks_labs_dqx_app.backend.run_config_store import RunConfigTooLargeError
from databricks_labs_dqx_app.backend.run_config_store import RunConfigError
from databricks_labs_dqx_app.backend.services.discovery import DiscoveryService
from databricks_labs_dqx_app.backend.services.materializer import MaterializationError, Materializer
from databricks_labs_dqx_app.backend.services.monitored_table_service import (
Expand Down Expand Up @@ -707,7 +707,7 @@ def run_monitored_table(
raise HTTPException(status_code=422, detail=str(e))
except BindingRunError as e:
raise HTTPException(status_code=400, detail=str(e))
except RunConfigTooLargeError as e:
except RunConfigError as e:
raise HTTPException(status_code=400, detail=str(e))
except Exception as e:
logger.error(f"Failed to run monitored table {binding_id}: {e}", exc_info=True)
Expand Down
Loading
Loading