diff --git a/CHANGELOG.md b/CHANGELOG.md index b14d88e2..831c5880 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,9 +3,12 @@ ## [v0.1.1] ### Added - **Multiple pipeline topologies per DAB bundle**: `bundle-add-pipeline` adds independently configured bronze, silver, split, or combined pipelines with their own data-flow groups and target schemas. `bundle-validate` validates each pipeline and its job wiring independently. [Issue #446](https://github.com/databrickslabs/sdp-meta/issues/446) +- **Auto Loader schema-evolution demos**: focused and at-scale demos exercise additive schema evolution, strict-mode failure behavior, and multi-table onboarding from Unity Catalog Volumes. [PR #490](https://github.com/databrickslabs/sdp-meta/pull/490) ### Fixed - **Databricks SDK compatibility floor**: the minimum supported SDK is now `databricks-sdk>=0.138.0,<1` across wheel, runtime, and Databricks App dependencies. CI installs exactly 0.138.0 and exercises the serverless `jobs.JobEnvironment` and `compute.Environment` models used by onboarding jobs. [Issue #457](https://github.com/databrickslabs/sdp-meta/issues/457) +- **[Databricks App] Make demo launches asynchronous and feature the 100-table Auto Loader demo**. [Issue #491](https://github.com/databrickslabs/sdp-meta/issues/491) +- **Data-quality expectations without quarantine targets**: bronze and silver flows can apply DQ expectations when no quarantine table is configured, without referencing a missing quarantine target. [PR #487](https://github.com/databrickslabs/sdp-meta/pull/487) - **Databricks App deployment packaging**: macOS/Linux and Windows deployment staging now includes `examples/`, allowing the App container to build the SDP-META wheel after packaged MCP examples became a required wheel input. - **DAB pipeline group ownership**: `bundle-add-pipeline` now rejects duplicate ownership of the same data-flow group and layer before writing files, while preserving supported bronze/silver split ownership and wiring its job dependency. Its fail-safe preflight checks the complete merged topology and every target override. `bundle-validate` also detects ownership conflicts introduced through manual YAML edits, including conflicts that exist only under a target override. [Issue #458](https://github.com/databrickslabs/sdp-meta/issues/458) - **Customized DAB validation**: `bundle-validate` now honors default and environment-selected targets, accepts transitive bronze-to-silver task dependencies, and reports malformed targets or dependency cycles cleanly. [Issue #459](https://github.com/databrickslabs/sdp-meta/issues/459) diff --git a/databricks_app/README.md b/databricks_app/README.md index 3f86d1ff..52550df9 100644 --- a/databricks_app/README.md +++ b/databricks_app/README.md @@ -243,6 +243,9 @@ export FLASK_DEBUG=true # Optional: pick a non-default profile (matches the Apps "--profile" semantic) # export DATABRICKS_CONFIG_PROFILE= +# Optional: concurrent demo launch limit (default: 8) +# export SDP_META_MAX_ACTIVE_DEMOS=8 + flask --app databricks_app/app.py run --host 127.0.0.1 --port 8000 ``` diff --git a/databricks_app/_jobs.py b/databricks_app/_jobs.py index 3185c736..c3ff60f2 100644 --- a/databricks_app/_jobs.py +++ b/databricks_app/_jobs.py @@ -19,27 +19,141 @@ from __future__ import annotations +import os +import time +import threading import uuid _jobs: dict = {} +_jobs_lock = threading.RLock() +_MAX_LOG_LINES = 5000 +_MAX_COMPLETED_JOBS = 50 +_COMPLETED_JOB_TTL_SECONDS = 60 * 60 -def _new_job_token() -> str: +def _positive_int_env(name: str, default: int) -> int: + """Read a positive integer environment setting with a safe fallback.""" + try: + value = int(os.environ.get(name, str(default))) + except (TypeError, ValueError): + return default + return value if value > 0 else default + + +_MAX_ACTIVE_DEMO_JOBS = _positive_int_env( + 'SDP_META_MAX_ACTIVE_DEMOS', + 8, +) + + +class JobCapacityError(RuntimeError): + """Raised when a bounded job category has no free launch slots.""" + + +def _prune_jobs(now: float | None = None) -> None: + """Evict expired and excess completed jobs from the process store.""" + with _jobs_lock: + now = time.time() if now is None else now + stale = [ + token + for token, job in _jobs.items() + if job.get('done') + and now - job.get('finished_at', now) > _COMPLETED_JOB_TTL_SECONDS + ] + for token in stale: + _jobs.pop(token, None) + + completed = sorted( + ( + (job.get('finished_at', 0), token) + for token, job in _jobs.items() + if job.get('done') + ), + reverse=True, + ) + for _, token in completed[_MAX_COMPLETED_JOBS:]: + _jobs.pop(token, None) + + +def _new_job_token( + *, + kind: str = 'generic', + max_active_for_kind: int | None = None, +) -> str: """Allocate a fresh job-token entry in ``_jobs`` and return the token.""" - token = uuid.uuid4().hex - _jobs[token] = { - 'logs': [], - 'done': False, - 'returncode': None, - 'stdout': '', - 'stderr': '', - 'modal_content': None, - 'error': None, - } - return token + with _jobs_lock: + _prune_jobs() + if max_active_for_kind is not None: + active_count = sum( + 1 + for job in _jobs.values() + if not job.get('done') and job.get('kind') == kind + ) + if active_count >= max_active_for_kind: + raise JobCapacityError( + f"Too many active {kind} jobs " + f"({active_count}/{max_active_for_kind})" + ) + token = uuid.uuid4().hex + _jobs[token] = { + 'kind': kind, + 'logs': [], + 'log_base_offset': 0, + 'done': False, + 'returncode': None, + 'stdout': '', + 'stderr': '', + 'modal_content': None, + 'error': None, + 'created_at': time.time(), + 'finished_at': None, + } + return token + + +def _append_job_log(job: dict, entry: dict) -> None: + """Append one line while bounding retained per-job output.""" + with _jobs_lock: + if len(job['logs']) >= _MAX_LOG_LINES: + job['logs'].pop(0) + job['log_base_offset'] += 1 + job['logs'].append(entry) + + +def _update_job(token: str, **values) -> None: + """Atomically update fields on an existing job.""" + with _jobs_lock: + job = _jobs.get(token) + if job is not None: + job.update(values) + + +def _mark_job_done(token: str) -> None: + """Finalize a job and enforce completed-job retention limits.""" + with _jobs_lock: + job = _jobs.get(token) + if job is None: + return + job['done'] = True + job['finished_at'] = time.time() + _prune_jobs() def _get_job(token: str): """Look up a job entry by token, returning ``None`` if absent.""" - return _jobs.get(token) + with _jobs_lock: + _prune_jobs() + return _jobs.get(token) + + +def _get_job_snapshot(token: str): + """Return a stable copy suitable for request handlers.""" + with _jobs_lock: + _prune_jobs() + job = _jobs.get(token) + if job is None: + return None + snapshot = dict(job) + snapshot['logs'] = [dict(entry) for entry in job['logs']] + return snapshot diff --git a/databricks_app/_subprocess_runner.py b/databricks_app/_subprocess_runner.py index aca75dc6..241f56e0 100644 --- a/databricks_app/_subprocess_runner.py +++ b/databricks_app/_subprocess_runner.py @@ -17,33 +17,108 @@ import logging import os import queue as _queue_module +import signal import subprocess import sys import threading +import time +from collections import deque import _jobs as _jobs_module # noqa: E402 \u2014 absolute import; databricks_app/ is not a package logger = logging.getLogger(__name__) -def _run_cli_json_payload( +def _process_group_exists(process_group_id: int) -> bool: + """Return whether a POSIX process group still has live members.""" + try: + os.killpg(process_group_id, 0) + except ProcessLookupError: + return False + except PermissionError: + return True + return True + + +def _terminate_process_tree( + proc: subprocess.Popen, + grace_seconds: float, +) -> None: + """Terminate a child and every descendant in its isolated session.""" + if os.name == 'posix': + try: + os.killpg(proc.pid, signal.SIGTERM) + except OSError: + pass + else: + deadline = time.monotonic() + grace_seconds + while ( + _process_group_exists(proc.pid) + and time.monotonic() < deadline + ): + proc.poll() + time.sleep(min(0.1, max(0.0, deadline - time.monotonic()))) + + # The leader can exit while descendants retain pipes or + # ignore SIGTERM, so escalation depends on the process + # group—not only on proc.wait(). + if _process_group_exists(proc.pid): + try: + os.killpg(proc.pid, signal.SIGKILL) + except ProcessLookupError: + pass + + try: + proc.wait(timeout=max(1.0, grace_seconds)) + except subprocess.TimeoutExpired: + proc.kill() + try: + proc.wait(timeout=max(1.0, grace_seconds)) + except subprocess.TimeoutExpired: + logger.error( + "Subprocess PID %s did not exit after SIGKILL; " + "continuing job finalization", + getattr(proc, 'pid', '?'), + ) + return + + proc.terminate() + try: + proc.wait(timeout=grace_seconds) + except subprocess.TimeoutExpired: + proc.kill() + try: + proc.wait(timeout=max(1.0, grace_seconds)) + except subprocess.TimeoutExpired: + logger.error( + "Subprocess PID %s did not exit after kill; " + "continuing job finalization", + getattr(proc, 'pid', '?'), + ) + + +def _run_command_in_background( token: str, - json_string: str, + command: list[str], cwd: str, + env: dict | None = None, cleanup_path: str | None = None, + idle_timeout_seconds: int = 600, ) -> None: - """Launch ``python -m databricks.labs.sdp_meta.cli `` in a - background thread and stream its output into ``_jobs[token]``. + """Launch a command in a background thread and stream its output. Arguments: token: a job token previously allocated via ``_new_job_token``. - json_string: the single positional argument the CLI expects. + command: subprocess argv, including the executable. cwd: working directory for the subprocess \u2014 typically ``_repo_root()`` so demo scripts can resolve relative paths. + env: optional child environment. Defaults to the current process. cleanup_path: optional path to ``os.unlink`` after the subprocess completes \u2014 used to clean up the tempfile that ``_resolve_local_onboarding_path`` may have created when the user pointed at a UC Volume / DBFS spec. + idle_timeout_seconds: terminate a child that produces no output for + this long. Long-running demo waiters override this default. Returns immediately after spawning the background thread; the caller is expected to return the token to the client and poll. @@ -53,7 +128,7 @@ def _run_cli_json_payload( # Maximum seconds with no output before we treat the child as # hung. Pulled out so it shows up in stack traces and so a test # can monkey-patch it instead of waiting 10 minutes. - _IDLE_TIMEOUT_S = 600 + _IDLE_TIMEOUT_S = idle_timeout_seconds # Graceful-shutdown grace period after ``terminate()`` before we # escalate to ``kill()``. CLI shells out to ``pip wheel`` which # can take a few seconds to unwind cleanly. @@ -61,18 +136,23 @@ def _run_cli_json_payload( def _run(): proc: subprocess.Popen | None = None - stdout_parts: list[str] = [] - stderr_parts: list[str] = [] + timed_out = False + stdout_parts = deque(maxlen=_jobs_module._MAX_LOG_LINES) + stderr_parts = deque(maxlen=_jobs_module._MAX_LOG_LINES) try: - _env = {**os.environ, 'PYTHONUNBUFFERED': '1'} + _env = {**(env or os.environ), 'PYTHONUNBUFFERED': '1'} proc = subprocess.Popen( - [sys.executable, '-m', 'databricks.labs.sdp_meta.cli', json_string], + command, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, bufsize=1, # line-buffered on our side env=_env, # force child Python to flush every print() cwd=cwd, + # The CLI launches Python, pip, and Databricks helper + # processes. Isolate that tree so an idle timeout can + # terminate every descendant, not only the direct child. + start_new_session=(os.name == 'posix'), ) q: _queue_module.Queue = _queue_module.Queue() @@ -94,14 +174,16 @@ def _reader(pipe, stream_name): # ``stderr`` and never called ``proc.wait()`` \u2014 leaking # the child as a zombie until the worker process exited. done_count = 0 - timed_out = False try: while done_count < 2: stream, line = q.get(timeout=_IDLE_TIMEOUT_S) if line is None: done_count += 1 continue - job['logs'].append({'stream': stream, 'line': line}) + _jobs_module._append_job_log( + job, + {'stream': stream, 'line': line}, + ) if stream == 'stdout': stdout_parts.append(line) else: @@ -113,24 +195,46 @@ def _reader(pipe, stream_name): _IDLE_TIMEOUT_S, getattr(proc, 'pid', '?'), ) - job['error'] = ( - f"Subprocess produced no output for {_IDLE_TIMEOUT_S} " - f"seconds; terminated to avoid leaking a zombie process. " - f"Any output collected before the timeout is preserved " - f"below." + _jobs_module._update_job( + token, + error=( + f"Subprocess produced no output for {_IDLE_TIMEOUT_S} " + f"seconds; terminated to avoid leaking a zombie process. " + f"Any output collected before the timeout is preserved " + f"below." + ), ) # Either the readers drained both pipes (success path) or # we timed out and are about to escalate in ``finally``. - # Calling ``wait()`` here for the success path lets us - # populate ``returncode`` from the natural exit; the - # timeout path will instead see returncode from - # ``terminate`` + ``wait`` below. + # A child can close both pipes without exiting, so keep + # ``wait()`` bounded as well; otherwise it can hold a demo + # capacity slot forever while producing no more output. if not timed_out: - proc.wait() - job['stdout'] = '\n'.join(stdout_parts) - job['stderr'] = '\n'.join(stderr_parts) - job['returncode'] = proc.returncode + try: + proc.wait(timeout=_IDLE_TIMEOUT_S) + except subprocess.TimeoutExpired: + timed_out = True + logger.warning( + "Subprocess closed its output streams but did not " + "exit within %ds; terminating PID %s", + _IDLE_TIMEOUT_S, + getattr(proc, 'pid', '?'), + ) + _jobs_module._update_job( + token, + error=( + "Subprocess closed its output streams but did not " + f"exit within {_IDLE_TIMEOUT_S} seconds; terminated " + "to release its background-job slot." + ), + ) + _jobs_module._update_job( + token, + stdout='\n'.join(stdout_parts), + stderr='\n'.join(stderr_parts), + returncode=proc.returncode, + ) except Exception as exc: logger.exception("Background CLI subprocess thread failed") # Preserve the partial output we collected before the @@ -141,34 +245,30 @@ def _reader(pipe, stream_name): # \u2014 we only skip overwriting when a meaningful value # is already present (i.e. the success path got far # enough to set them). + values = {'returncode': -1} if not job.get('stdout'): - job['stdout'] = '\n'.join(stdout_parts) + values['stdout'] = '\n'.join(stdout_parts) if not job.get('stderr'): - job['stderr'] = '\n'.join(stderr_parts) + values['stderr'] = '\n'.join(stderr_parts) if not job.get('error'): - job['error'] = str(exc) - job['returncode'] = -1 + values['error'] = str(exc) + _jobs_module._update_job(token, **values) finally: # Always reap the child, even when the read loop never # reached ``proc.wait()`` (timeout) or the Popen call # itself raised (proc is None). Without this the child # outlives the gunicorn worker as a zombie. - if proc is not None and proc.poll() is None: + if proc is not None and (timed_out or proc.poll() is None): try: - proc.terminate() - try: - proc.wait(timeout=_TERMINATE_GRACE_S) - except subprocess.TimeoutExpired: - proc.kill() - proc.wait() + _terminate_process_tree(proc, _TERMINATE_GRACE_S) except OSError: # ProcessLookupError etc. \u2014 child already gone. pass # If we had to terminate, refresh returncode from the # signal that landed. if 'returncode' not in job or job['returncode'] is None: - job['returncode'] = proc.returncode - job['done'] = True + _jobs_module._update_job(token, returncode=proc.returncode) + _jobs_module._mark_job_done(token) if cleanup_path and os.path.exists(cleanup_path): try: os.unlink(cleanup_path) @@ -176,3 +276,23 @@ def _reader(pipe, stream_name): pass threading.Thread(target=_run, daemon=True).start() + + +def _run_cli_json_payload( + token: str, + json_string: str, + cwd: str, + cleanup_path: str | None = None, +) -> None: + """Run the SDP-META CLI JSON payload through the shared runner.""" + _run_command_in_background( + token=token, + command=[ + sys.executable, + '-m', + 'databricks.labs.sdp_meta.cli', + json_string, + ], + cwd=cwd, + cleanup_path=cleanup_path, + ) diff --git a/databricks_app/app.py b/databricks_app/app.py index 22c4b209..34f1cb5a 100644 --- a/databricks_app/app.py +++ b/databricks_app/app.py @@ -30,9 +30,10 @@ import logging import os +import uuid import subprocess # noqa: F401 \u2014 re-exported below so tests can mock subprocess.Popen -from flask import Flask, jsonify +from flask import Flask, jsonify, request from werkzeug.exceptions import HTTPException # ── Logging ────────────────────────────────────────────────────────────────── @@ -68,11 +69,32 @@ # logging them as application errors. @app.errorhandler(Exception) def handle_exception(exc): + request_id = uuid.uuid4().hex[:12] if isinstance(exc, HTTPException): - return jsonify({'error': exc.description or str(exc)}), exc.code - logger.exception("Unhandled exception in route: %s", exc) + message = exc.description or str(exc) + return jsonify({ + 'error': message, + 'details': { + 'request_id': request_id, + 'exception_type': type(exc).__name__, + 'message': message, + }, + }), exc.code + logger.exception( + "Unhandled exception in route (request_id=%s): %s", + request_id, + exc, + ) return jsonify({ - 'error': str(exc), + 'error': ( + 'The app could not complete this request. ' + 'Expand Technical details for diagnostics.' + ), + 'details': { + 'request_id': request_id, + 'exception_type': type(exc).__name__, + 'message': str(exc), + }, 'stdout': '', 'stderr': '', 'returncode': -1, @@ -121,9 +143,60 @@ def _inject_app_version(): # ── Security headers ───────────────────────────────────────────────────────── +def _friendly_server_error_message(): + """Return an actionable headline without exposing SDK exception text.""" + path = request.path + if path.startswith('/api/metadata/workspace-'): + return ( + 'The app could not access the requested workspace file. Verify ' + 'the path and the App service principal workspace permissions.' + ) + if path.startswith('/api/metadata/'): + return ( + 'The app could not access Unity Catalog. Verify that the App ' + 'service principal has USE CATALOG, USE SCHEMA, and the required ' + 'table privileges.' + ) + if path.startswith('/api/warehouse'): + return ( + 'The app could not access the SQL warehouse. Verify that the App ' + 'service principal has CAN USE permission on the warehouse.' + ) + if path.startswith('/api/pipelines'): + return ( + 'The app could not complete the pipeline request. Verify the ' + 'pipeline exists and the App service principal can manage it.' + ) + return ( + 'The app could not complete this request. ' + 'Expand Technical details for diagnostics.' + ) + + @app.after_request def add_security_headers(response): - """Attach HTTP security headers to every response (fix M4).""" + """Normalize JSON errors and attach HTTP security headers.""" + if response.status_code >= 400 and response.is_json: + payload = response.get_json(silent=True) + if isinstance(payload, dict) and payload.get('error'): + original_error = str(payload['error']) + details = payload.get('details') + if not isinstance(details, dict): + details = {} + details.setdefault('status', response.status_code) + details.setdefault('message', original_error) + payload['details'] = details + + if response.status_code >= 500: + payload['error'] = _friendly_server_error_message() + elif response.status_code in (401, 403): + payload['error'] = ( + 'The Databricks App service principal does not have ' + 'permission to complete this request. Ask a catalog or ' + 'workspace administrator to grant the required access.' + ) + response.set_data(app.json.dumps(payload)) + response.headers['Content-Security-Policy'] = ( "default-src 'self'; " "script-src 'self' 'unsafe-inline'; " diff --git a/databricks_app/app.yaml b/databricks_app/app.yaml index c1ddc857..97310a25 100644 --- a/databricks_app/app.yaml +++ b/databricks_app/app.yaml @@ -33,3 +33,6 @@ env: # (Warehouse button in the top bar → enter an existing ID or create a new serverless one). - name: DATABRICKS_SQL_WAREHOUSE_ID value: "" + # Maximum demo launcher subprocesses allowed at the same time. + - name: SDP_META_MAX_ACTIVE_DEMOS + value: "8" diff --git a/databricks_app/routes/dataflowspecs.py b/databricks_app/routes/dataflowspecs.py index 00a2f016..37382618 100644 --- a/databricks_app/routes/dataflowspecs.py +++ b/databricks_app/routes/dataflowspecs.py @@ -128,7 +128,19 @@ def _run_query(table_name): return {'columns': columns, 'rows': rows, 'groups': groups, 'error': None} except Exception as exc: logger.exception("DataflowSpec query failed for %s.%s.%s", catalog, schema, table_name) - return {'columns': [], 'rows': [], 'groups': [], 'error': str(exc)} + return { + 'columns': [], + 'rows': [], + 'groups': [], + 'error': ( + 'The app could not read this DataflowSpec table. Verify ' + 'the table name, SQL warehouse, and App service principal permissions.' + ), + 'details': { + 'exception_type': type(exc).__name__, + 'message': str(exc), + }, + } bronze_result = _run_query(bronze_table) silver_result = _run_query(silver_table) diff --git a/databricks_app/routes/demo.py b/databricks_app/routes/demo.py index c749b8d7..49af77ce 100644 --- a/databricks_app/routes/demo.py +++ b/databricks_app/routes/demo.py @@ -14,15 +14,15 @@ import logging import os -import subprocess import sys from dataclasses import asdict from flask import Blueprint, jsonify, request import _jobs as _jobs_module -from _command_output import _parse_command_result, extract_command_output +from _command_output import _parse_command_result from _config import _repo_root +from _subprocess_runner import _run_command_in_background # UC catalog pre-flight (Apps-SP grants). Lives next to app.py so the # probe + GRANT SQL builder ship in the same source tree and can be @@ -92,6 +92,14 @@ def check_uc_grants(): # to the App's service principal in this workspace. Run via the # CLI launcher with ``--profile`` instead. _DEMO_REGISTRY = { + "demo_at_scale_autoloader": { + "file": "demo/launch_at_scale_autoloader_demo.py", + "uc_arg": "--uc_catalog_name", + # App users launch demos for inspection. Preserve the job, pipelines, + # schemas, volume, and output tables instead of deleting them as soon + # as validation succeeds. The standalone CLI keeps its cleanup default. + "extra_args": ["--keep-resources"], + }, "demo_cloudfiles": { "file": "demo/launch_af_cloudfiles_demo.py", "uc_arg": "--uc_catalog_name", @@ -111,13 +119,11 @@ def check_uc_grants(): "demo_interactive": { "file": "demo/launch_interactive_demo.py", "uc_arg": "--uc-catalog-name", - # The interactive launcher submits a serverless job, prints - # the run URL EARLY (before polling), and then blocks on - # ``waiter.result(timeout=timedelta(minutes=N))``. We pass a - # 1-minute timeout so the Flask request unblocks shortly - # after submission with the run URL captured in stdout; the - # actual demo job continues running in the workspace and the - # user clicks through via the surfaced URL. + # Stream upload/submission output and the early run URL through the + # same background progress UI as the 100-table featured demo. Keep + # the launcher's normal completion wait so the UI reports the real + # final state instead of turning a one-minute timeout into a false + # failure while the remote run continues. # # ``--install-source pypi`` makes the spawned job # ``pip install databricks-labs-sdp-meta`` from PyPI on @@ -130,7 +136,6 @@ def check_uc_grants(): # e.g. ``"--pypi-version", "0.1.0"`` to the list below. "extra_args": [ "--install-source", "pypi", - "--timeout-minutes", "1", ], }, } @@ -212,28 +217,65 @@ def run_demo(): pypath_entries.append(existing_pypath) demo_env['PYTHONPATH'] = ':'.join(pypath_entries) - result = subprocess.run( - [ - sys.executable, - os.path.join(current_directory, demo_file), - demo_uc_arg, - uc_name, - *demo_extra_args, - ], - shell=False, - capture_output=True, - text=True, + demo_command = [ + sys.executable, + os.path.join(current_directory, demo_file), + demo_uc_arg, + uc_name, + *demo_extra_args, + ] + # Deployed Apps already have DATABRICKS_APP_PORT, which makes launchers + # use ambient service-principal auth and enables the notebook-path shim + # needed for workspace-synced sources. A local Flask run must NOT fake + # that marker: locally uploaded notebooks retain their ``.py`` suffix. + # Instead, forward the selected CLI profile explicitly so the child stays + # non-interactive while preserving local notebook paths. + if not demo_env.get('DATABRICKS_APP_PORT'): + local_profile = demo_env.get('DATABRICKS_CONFIG_PROFILE') + if local_profile: + demo_command.extend(['--profile', local_profile]) + + # Every demo launcher can wait many minutes for remote jobs or pipelines. + # Always return immediately and stream progress through the shared polling + # UI instead of holding the browser request behind a blocking spinner. + try: + token = _jobs_module._new_job_token( + kind='demo', + max_active_for_kind=_jobs_module._MAX_ACTIVE_DEMO_JOBS, + ) + except _jobs_module.JobCapacityError: + logger.warning( + "Demo launch rejected because the active-demo limit (%d) " + "has been reached", + _jobs_module._MAX_ACTIVE_DEMO_JOBS, + ) + response = jsonify({ + 'error': ( + "Active demo capacity is currently full. Wait for an " + "existing demo to finish before " + "launching another." + ), + }) + response.status_code = 429 + response.headers['Retry-After'] = '30' + return response + _run_command_in_background( + token=token, + command=demo_command, cwd=current_directory, env=demo_env, + # Interactive waits up to 90 minutes and may be silent while the + # remote job runs. Keep the orphan-protection timeout above that. + idle_timeout_seconds=2 * 60 * 60, ) - return extract_command_output(result) + return jsonify({'token': token}), 202 @bp.route('/api/job//logs', methods=['GET']) def get_job_logs(token): """Polling endpoint: returns buffered log lines + done/returncode for the progress UI.""" - job = _jobs_module._get_job(token) + job = _jobs_module._get_job_snapshot(token) if job is None: return jsonify({'error': 'Job not found'}), 404 @@ -254,9 +296,12 @@ def get_job_logs(token): return jsonify({ 'error': f"offset must be non-negative (got {offset})" }), 400 - new_logs = job['logs'][offset:] + base_offset = job.get('log_base_offset', 0) + relative_offset = max(0, offset - base_offset) + new_logs = job['logs'][relative_offset:] payload: dict = { 'logs': new_logs, + 'next_offset': base_offset + len(job['logs']), 'done': job['done'], 'returncode': job.get('returncode'), 'error': job.get('error'), diff --git a/databricks_app/routes/warehouse.py b/databricks_app/routes/warehouse.py index deccbf7d..c5173fb7 100644 --- a/databricks_app/routes/warehouse.py +++ b/databricks_app/routes/warehouse.py @@ -43,7 +43,16 @@ def warehouse_status(): logger.exception("warehouse_status failed for id=%s", wh_id) return jsonify({ 'configured': True, 'warehouse_id': wh_id, - 'name': None, 'state': None, 'error': str(exc), + 'name': None, + 'state': None, + 'error': ( + 'The app could not read the configured SQL warehouse. ' + 'Verify that the App service principal can use the warehouse.' + ), + 'details': { + 'exception_type': type(exc).__name__, + 'message': str(exc), + }, }) @@ -101,9 +110,14 @@ def configure_warehouse(): 'state': state_val, 'message': f"Warehouse \"{wh.name}\" configured successfully.", }) - except Exception as exc: + except Exception: logger.exception("configure_warehouse (existing) failed for id=%s", wh_id) - return jsonify({'error': str(exc)}), 400 + return jsonify({ + 'error': ( + 'The app could not use that SQL warehouse. Verify the ' + 'warehouse ID and the App service principal permissions.' + ), + }), 400 elif mode == 'create': name = (body.get('name') or 'sdp-meta-app-warehouse').strip() diff --git a/databricks_app/services/onboarding/bundled_specs.py b/databricks_app/services/onboarding/bundled_specs.py index 60e0ee4c..39db1c0a 100644 --- a/databricks_app/services/onboarding/bundled_specs.py +++ b/databricks_app/services/onboarding/bundled_specs.py @@ -33,8 +33,7 @@ the Environment field's value at submit time, OR the entry must declare ``env_override`` to force the right value. -The current 4 entries (Cars, Multi-Source CDC, Cloud Files, -DAIS) all satisfy this contract. Keep this registry in lockstep +The current entries all satisfy this contract. Keep this registry in lockstep with ``databricks_app/routes/demo.py::_DEMO_REGISTRY``. """ diff --git a/databricks_app/templates/landingPage.html b/databricks_app/templates/landingPage.html index 7eedc4d0..174ebd48 100644 --- a/databricks_app/templates/landingPage.html +++ b/databricks_app/templates/landingPage.html @@ -860,6 +860,48 @@ opacity: 0; transition: opacity .18s; } .demo-card:hover .demo-cta { opacity: 1; } + .featured-demo-picker { + display: grid; + grid-template-columns: minmax(220px, 320px) 1fr auto; + gap: 18px; + align-items: end; + padding: 22px; + background: var(--surface); + border: 1.5px solid var(--border); + border-radius: var(--radius); + } + .featured-demo-field { + display: flex; + flex-direction: column; + gap: 8px; + } + .featured-demo-field label { + font-size: .74rem; + font-weight: 700; + color: var(--ink-60); + text-transform: uppercase; + letter-spacing: .07em; + } + .featured-demo-details { + display: grid; + grid-template-columns: 42px 1fr; + gap: 4px 12px; + align-items: center; + min-height: 58px; + } + .featured-demo-details .demo-icon { + grid-row: 1 / 3; + } + .featured-demo-details h4 { + font-size: .9rem; + font-weight: 700; + color: var(--ink); + } + .featured-demo-details p { + font-size: .79rem; + color: var(--ink-40); + line-height: 1.5; + } /* ── Loading overlay ─────────────────────────────────────────── */ #loadingOverlay { @@ -939,6 +981,7 @@ .progress-log .log-stdout { color: #9cdcfe; } .progress-log .log-stderr { color: #f4bf4f; } .progress-log .log-info { color: #4ec94e; font-style: italic; } + .progress-log .progress-link { color: inherit; text-decoration: underline; } .progress-foot { padding: 12px 20px; display: flex; @@ -1043,6 +1086,7 @@ .form-grid-2 { grid-template-columns: 1fr; } .demo-grid { grid-template-columns: 1fr 1fr; } .uc-callout { flex-direction: column; align-items: flex-start; } + .featured-demo-picker { grid-template-columns: 1fr; align-items: stretch; } .form-card { padding: 20px; } } @@ -1433,7 +1477,7 @@

Pipeline Onboarding

- +
-
-
-
☁️
-

Cloud Files

-

Ingest streaming data using Auto Loader with cloud file sources

- Run demo → -
-
-
📷
-

Apply Changes Snapshot

-

CDC with SCD Type 1 apply changes from full data snapshots

- Run demo → -
-
-
🔀
-

Silver Fanout

-

Fan-out bronze data into multiple silver tables in one pipeline

- Run demo → -
-
-
🧠
-

DAIS Demo

-

Databricks AI Summit end-to-end pipeline walkthrough

- Run demo → + @@ -1989,14 +2022,33 @@

Response

} /* ── Progress log modal (onboarding) ───────────────────── */ - var _progressPoll = null; - var _dotsTimer = null; + var _progressSession = null; + var _progressSessionCounter = 0; + var _progressPending = false; + var _dotsTimer = null; + + function _claimProgressOperation() { + if (_progressPending || (_progressSession && !_progressSession.finished)) { + /* Keep one recoverable progress stream per browser tab. Reopen the + existing dialog instead of launching an operation whose token + would replace—and make unreachable—the current one. */ + document.getElementById('progressModal').classList.add('visible'); + _startDots(); + return false; + } + _progressPending = true; + return true; + } + + function _releaseProgressClaim() { + _progressPending = false; + } function showProgressModal(title) { document.getElementById('progressTitle').textContent = title || 'Running…'; document.getElementById('progressLog').innerHTML = ''; document.getElementById('progressStatus').textContent = 'Starting…'; - document.getElementById('progressClose').disabled = true; + document.getElementById('progressClose').disabled = false; document.getElementById('progressDots').textContent = '●'; document.getElementById('progressDots').style.color = ''; document.getElementById('progressModal').classList.add('visible'); @@ -2005,7 +2057,8 @@

Response

function closeProgressModal() { document.getElementById('progressModal').classList.remove('visible'); - if (_progressPoll) { clearInterval(_progressPoll); _progressPoll = null; } + /* Hiding progress must not stop tracking. Long-running demos continue + in the background and still surface their completion result later. */ if (_dotsTimer) { clearInterval(_dotsTimer); _dotsTimer = null; } } @@ -2023,42 +2076,148 @@

Response

span.className = stream === 'stderr' ? 'log-stderr' : stream === 'info' ? 'log-info' : 'log-stdout'; - span.textContent = line + '\n'; + var urlPattern = /https:\/\/[^\s<>"']+/g; + var cursor = 0; + var match; + while ((match = urlPattern.exec(line)) !== null) { + span.appendChild(document.createTextNode(line.slice(cursor, match.index))); + var link = document.createElement('a'); + link.href = match[0]; + link.textContent = match[0]; + link.target = '_blank'; + link.rel = 'noopener noreferrer'; + link.className = 'progress-link'; + span.appendChild(link); + cursor = match.index + match[0].length; + } + span.appendChild(document.createTextNode(line.slice(cursor) + '\n')); log.appendChild(span); log.scrollTop = log.scrollHeight; } function _finishProgressModal(success, statusText) { + _releaseProgressClaim(); var dotsEl = document.getElementById('progressDots'); dotsEl.textContent = success ? '✓' : '✗'; dotsEl.style.color = success ? '#4ec94e' : '#f47c7c'; document.getElementById('progressStatus').textContent = statusText || (success ? 'Done' : 'Failed'); document.getElementById('progressClose').disabled = false; if (_dotsTimer) { clearInterval(_dotsTimer); _dotsTimer = null; } - if (_progressPoll) { clearInterval(_progressPoll); _progressPoll = null; } } function _pollJob(token, onDone) { + _releaseProgressClaim(); var offset = 0; - _progressPoll = setInterval(function() { - fetch('/api/job/' + token + '/logs?offset=' + offset) - .then(function(r) { return r.json(); }) + var maxRetries = 4; + var session = { + id: ++_progressSessionCounter, + token: token, + timer: null, + controller: null, + retryCount: 0, + finished: false + }; + _progressSession = session; + + function ownsProgressSurface() { + return _progressSession === session; + } + + function finishOnce(data) { + if (!ownsProgressSurface() || session.finished) return; + session.finished = true; + if (session.timer) clearTimeout(session.timer); + session.controller = null; + _progressSession = null; + onDone(data); + } + + function schedulePoll(delay) { + if (!ownsProgressSurface() || session.finished) return; + session.timer = setTimeout(poll, delay); + } + + function failPolling(err) { + if (!ownsProgressSurface()) return; + var message = err.message || String(err); + var userMessage = err.status === 404 + ? 'Progress tracking expired. Start the operation again.' + : 'The app could not check progress after repeated attempts.'; + _appendProgressLine('stderr', message); + _finishProgressModal(false, 'Unable to check progress'); + finishOnce({ + done: true, + returncode: -1, + error: userMessage, + result: { + error: userMessage, + returncode: -1, + details: { + message: message, + status: err.status || null, + retries: session.retryCount + } + } + }); + } + + function poll() { + if (!ownsProgressSurface() || session.finished) return; + session.controller = new AbortController(); + fetch( + '/api/job/' + token + '/logs?offset=' + offset, + { signal: session.controller.signal } + ) + .then(function(r) { + return r.json().catch(function() { return {}; }).then(function(data) { + if (!r.ok) { + var err = new Error( + data.error || ('Progress request failed (HTTP ' + r.status + ')') + ); + err.status = r.status; + throw err; + } + return data; + }); + }) .then(function(d) { + if (!ownsProgressSurface()) return; + session.controller = null; + session.retryCount = 0; (d.logs || []).forEach(function(entry) { _appendProgressLine(entry.stream, entry.line); }); - offset += (d.logs || []).length; + offset = Number.isInteger(d.next_offset) + ? d.next_offset + : offset + (d.logs || []).length; if (offset > 0) { document.getElementById('progressStatus').textContent = 'Running… (' + offset + ' lines)'; } if (d.done) { - clearInterval(_progressPoll); - _progressPoll = null; - onDone(d); + finishOnce(d); + } else { + schedulePoll(1500); } }) - .catch(function(err) { console.warn('Progress poll error:', err); }); - }, 1500); + .catch(function(err) { + if (!ownsProgressSurface() || err.name === 'AbortError') return; + session.controller = null; + console.warn('Progress poll error:', err); + var transientFailure = !err.status || err.status >= 500; + if (transientFailure && session.retryCount < maxRetries) { + session.retryCount += 1; + var retryDelay = Math.pow(2, session.retryCount - 1) * 1000; + document.getElementById('progressStatus').textContent = + 'Connection interrupted — retrying (' + + session.retryCount + '/' + maxRetries + ')…'; + schedulePoll(retryDelay); + return; + } + failPolling(err); + }); + } + + schedulePoll(0); } /* ── Modal ──────────────────────────────────────────────── */ @@ -2091,9 +2250,36 @@

Response

showModal('
' + escapeHtml(JSON.stringify(data, null, 2)) + '
', defaultTitle); } } + function _technicalDetailsHtml(details) { + if (!details) return ''; + var text; + try { + text = typeof details === 'string' ? details : JSON.stringify(details, null, 2); + } catch (e) { + text = String(details); + } + if (!text.trim()) return ''; + return '
' + + 'Technical details' + + '
' +
+            escapeHtml(text) + '
'; + } + function showErrorResponse(data, title) { + data = data || {}; + var message = data.error || 'The app could not complete this request.'; + showModal( + '

' + + escapeHtml(message) + '

' + _technicalDetailsHtml(data.details), + title + ); + } function handleError(err, title) { hideLoading(); - showModal('

' + escapeHtml(err.message || String(err)) + '

', title); + var payload = err && err.payload ? err.payload : { + error: err && err.message ? err.message : String(err), + details: err && err.details ? err.details : null + }; + showErrorResponse(payload, title); } /* Safe JSON parser — checks Content-Type before calling .json() so an @@ -2127,26 +2313,19 @@

Response

hideLoading(); /* Explicit error field (validation, 4xx, etc.) */ if (data.error) { - showModal( - '

' + - escapeHtml(data.error) + '

', - defaultTitle - ); + showErrorResponse(data, defaultTitle); return; } /* CLI subprocess failed */ if (data.returncode !== undefined && data.returncode !== 0) { - var body = ''; - if (data.stderr && data.stderr.trim()) { - body += '
'
-                     + escapeHtml(data.stderr.trim()) + '
'; - } - if (data.stdout && data.stdout.trim()) { - body += '
'
-                     + escapeHtml(data.stdout.trim()) + '
'; - } - if (!body) body = '

Command exited with code ' + data.returncode + '.

'; - showModal(body, defaultTitle + ' — Failed'); + showErrorResponse({ + error: 'The operation failed before it could complete.', + details: { + returncode: data.returncode, + stderr: data.stderr || '', + stdout: data.stdout || '' + } + }, defaultTitle + ' — Failed'); return; } /* Success with job details */ @@ -2238,6 +2417,21 @@

Response

document.getElementById('onboardingForm').addEventListener('submit', function(e) { e.preventDefault(); var form = this; + var bundledPicker = document.getElementById('bundled_spec_id'); + var sourceMode = (document.querySelector('input[name="source_mode"]:checked') || {}).value; + if (sourceMode === 'bundled' && bundledPicker && + bundledPicker.value === '__featured_at_scale_autoloader__') { + var onboardingCatalog = document.getElementById('unity_catalog_name'); + var catalogName = (onboardingCatalog.value || '').trim(); + if (!catalogName) { + _flashUcInputInvalid('unity_catalog_name'); + return; + } + document.getElementById('demo_uc_name').value = catalogName; + runDemo('demo_at_scale_autoloader'); + return; + } + if (!_claimProgressOperation()) return; showProgressModal('Running Onboarding…'); _appendProgressLine('info', '→ Submitting onboarding request…'); @@ -2257,11 +2451,7 @@

Response

_finishProgressModal(false, 'Validation failed'); setTimeout(function() { closeProgressModal(); - showModal( - '

' + - escapeHtml(d.error) + '

', - 'Onboarding Error' - ); + showErrorResponse(d, 'Onboarding Error'); }, 400); return; } @@ -2290,6 +2480,7 @@

Response

} }) .catch(function(err) { + _finishProgressModal(false, 'Request failed'); closeProgressModal(); handleError(err, 'Onboarding Error'); }); @@ -2388,6 +2579,19 @@

Response

noteHtml; } + function _setOnboardingActionForEntry(entry) { + var submitBtn = document.getElementById('onboardingSubmitBtn'); + var previewBtn = document.getElementById('onboardingPreviewBtn'); + var isLauncher = Boolean(entry && entry.launcher); + submitBtn.innerHTML = isLauncher + ? '▶ Run Featured Demo' + : '▶ Run Onboarding'; + previewBtn.disabled = isLauncher; + previewBtn.title = isLauncher + ? 'The featured demo is an orchestrated launcher and has no single onboarding template to preview.' + : 'Preview the onboarding template with your form values substituted — no side effects, no UC volume created, no job launched.'; + } + /* Canonical schema names that every shipped bundled demo expects via the ``{bronze_schema}`` / ``{silver_schema}`` placeholder substitution. Matches the DAB template defaults @@ -2427,6 +2631,17 @@

Response

function _syncBundledToHidden() { var entry = _bundledSpecById(bundledSelect.value); if (!entry) return; + /* The featured at-scale demo is an orchestrated launcher, not a + static onboarding spec. Keep it discoverable here, but route + it to the Demos workflow instead of writing invalid paths + into this form. */ + if (entry.launcher) { + featuredDemoSelect.value = entry.launcher; + renderFeaturedDemo(); + _renderBundledMeta(entry); + _setOnboardingActionForEntry(entry); + return; + } var fmt = _currentBundledFormat(); /* Prefer the user's chosen format, but fall back to whichever this spec actually ships if the chosen format is missing. */ @@ -2443,6 +2658,7 @@

Response

} _fillBundledSchemaDefaultsIfBlank(); _renderBundledMeta(entry); + _setOnboardingActionForEntry(entry); } /* Push the UC-Volume inputs into the hidden inputs. UC mode @@ -2484,8 +2700,10 @@

Response

_fillBundledSchemaDefaultsIfBlank(); _syncBundledToHidden(); } else if (mode === 'uc_volume') { + _setOnboardingActionForEntry(null); _syncUcToHidden(); } else { + _setOnboardingActionForEntry(null); _syncManualToHidden(); } } @@ -2496,7 +2714,12 @@

Response

fetch('/onboarding/bundled-specs') .then(function(r) { return r.json(); }) .then(function(body) { - bundledData = body.specs || []; + bundledData = [{ + id: '__featured_at_scale_autoloader__', + label: '\u2605 Featured \u2014 At-scale Auto Loader (100 tables) \u2192', + launcher: 'demo_at_scale_autoloader', + description: 'Featured demo: generate, onboard, run, and validate 100 metadata-driven Auto Loader tables. Enter a catalog above, then click Run Featured Demo.' + }].concat(body.specs || []); if (!bundledData.length) { bundledSelect.innerHTML = ''; bundledMeta.innerHTML = 'The App container has no bundled demos. Switch to "Unity Catalog Volume" or "Manual paths" mode.'; @@ -2506,29 +2729,15 @@

Response

return ''; }).join(''); - /* Default-pick priority: the flagship ``onboarding`` - (DAIS end-to-end demo) leads the picker because - it's the demo the App is most often shown with - \u2014 it exercises CDC + DQE + silver transforms - in a single pass and surfaces the most surface - area of SDP-META. Fall back to the simpler - ``onboarding_cars`` if DAIS is missing from the - container, then to the first entry. */ - var defaultIdx = 0; - var preferenceOrder = ['onboarding', 'onboarding_cars']; - for (var p = 0; p < preferenceOrder.length; p++) { - var found = false; - for (var i = 0; i < bundledData.length; i++) { - if (bundledData[i].id === preferenceOrder[p]) { - defaultIdx = i; - found = true; - break; - } - } - if (found) break; - } - bundledSelect.selectedIndex = defaultIdx; - _syncBundledToHidden(); + /* Keep the featured 100-table demo selected by default. + Do not call ``_syncBundledToHidden`` here: launcher + entries intentionally navigate to the Demos panel, + and page initialization must not redirect the user. + Choosing any option explicitly still uses the normal + change handler below. */ + bundledSelect.selectedIndex = 0; + _renderBundledMeta(bundledData[0]); + _setOnboardingActionForEntry(bundledData[0]); }) .catch(function(err) { bundledSelect.innerHTML = ''; @@ -2838,6 +3047,7 @@

Response

return; } + if (!_claimProgressOperation()) return; showProgressModal('Deploying Pipeline…'); _appendProgressLine('info', '→ Submitting deploy request…'); @@ -2856,11 +3066,7 @@

Response

_finishProgressModal(false, 'Validation failed'); setTimeout(function() { closeProgressModal(); - showModal( - '

' + - escapeHtml(d.error) + '

', - 'Deploy Error' - ); + showErrorResponse(d, 'Deploy Error'); }, 400); return; } @@ -2886,6 +3092,7 @@

Response

} }) .catch(function(err) { + _finishProgressModal(false, 'Request failed'); closeProgressModal(); handleError(err, 'Deploy Error'); }); @@ -3098,6 +3305,7 @@

Response

function runDemo(command) { var ucInput = document.getElementById('demo_uc_name'); if (!ucInput.value.trim()) { _flashUcInputInvalid('demo_uc_name'); return; } + if (!_claimProgressOperation()) return; /* Restore the Demo-panel ucContext so the grant-required modal opened by /rundemo (if any) targets the Demo-panel status container, not whichever input the user last tested. */ @@ -3117,22 +3325,89 @@

Response

lacks privileges on the user's catalog. Surface the panel instead of dumping the raw JSON into a generic modal. */ if (resp.status === 400 && resp.body && resp.body.grant_required) { + _releaseProgressClaim(); _setUcStatus('warn', 'Grant required. ' + escapeHtml(resp.body.error || '')); showGrantRequiredModal(resp.body); return; } + if (resp.status === 202 && resp.body && resp.body.token) { + var selectedDemo = featuredDemoCatalog[command]; + var demoTitle = selectedDemo ? selectedDemo.title : 'Demo'; + showProgressModal('Running ' + demoTitle); + _appendProgressLine( + 'info', + '\u2192 Launcher started. Progress and the Databricks run URL will appear here.' + ); + _pollJob(resp.body.token, function(jobData) { + var success = jobData.returncode === 0; + _finishProgressModal( + success, + success ? 'Completed successfully' : 'Failed' + ); + if (jobData.result) { + renderApiResponse(jobData.result, 'Demo Result'); + } + }); + return; + } + _releaseProgressClaim(); renderApiResponse(resp.body, 'Demo Result'); }) - .catch(function(e) { hideLoading(); handleError(e, 'Demo Error'); }); + .catch(function(e) { + _releaseProgressClaim(); + hideLoading(); + handleError(e, 'Demo Error'); + }); } - document.querySelectorAll('.demo-card').forEach(function(card) { - card.addEventListener('click', function() { runDemo(this.dataset.command); }); - card.addEventListener('keydown', function(e) { - if (e.key === 'Enter' || e.key === ' ') { e.preventDefault(); runDemo(this.dataset.command); } - }); + var featuredDemoSelect = document.getElementById('featured_demo_select'); + var featuredDemoCatalog = { + demo_at_scale_autoloader: { + icon: '\uD83D\uDCCA', + title: 'At-scale Auto Loader (100 tables)', + description: 'Generate, onboard, run, and validate 100 metadata-driven Auto Loader tables.' + }, + demo_interactive: { + icon: '\uD83D\uDCD6', + title: 'Interactive Notebook', + description: 'End-to-end notebook walkthrough — submits the SDP-META Interactive Demo as a serverless job.' + }, + demo_cloudfiles: { + icon: '\u2601\uFE0F', + title: 'Cloud Files', + description: 'Ingest streaming data using Auto Loader with cloud file sources.' + }, + demo_acf: { + icon: '\uD83D\uDCF7', + title: 'Apply Changes Snapshot', + description: 'CDC with SCD Type 1 apply changes from full data snapshots.' + }, + demo_silverfanout: { + icon: '\uD83D\uDD00', + title: 'Silver Fanout', + description: 'Fan out bronze data into multiple silver tables in one pipeline.' + }, + demo_dias: { + icon: '\uD83E\uDDE0', + title: 'DAIS Demo', + description: 'Databricks AI Summit end-to-end pipeline walkthrough.' + } + }; + + function renderFeaturedDemo() { + var demo = featuredDemoCatalog[featuredDemoSelect.value]; + if (!demo) return; + document.getElementById('featuredDemoIcon').textContent = demo.icon; + document.getElementById('featuredDemoTitle').textContent = demo.title; + document.getElementById('featuredDemoDescription').textContent = demo.description; + } + + featuredDemoSelect.addEventListener('change', renderFeaturedDemo); + document.getElementById('runFeaturedDemo').addEventListener('click', function() { + runDemo(featuredDemoSelect.value); }); + renderFeaturedDemo(); /* Clear the inline access status when the user edits the catalog name — a stale "OK" indicator for a different catalog name would @@ -3191,6 +3466,7 @@

Response

.then(function(list) { if (list.error) { sel.innerHTML = ''; + showErrorResponse(list, 'Warehouse Error'); return; } if (!list.length) { @@ -3237,6 +3513,7 @@

Response

if (d.error) { msgEl.className = 'err'; msgEl.textContent = 'Error: ' + d.error; + showErrorResponse(d, 'Warehouse Error'); } else { msgEl.className = 'ok'; msgEl.textContent = d.message || 'Configured.'; @@ -3438,6 +3715,7 @@

Response

if (d.error) { statusEl.style.color = 'var(--brand)'; statusEl.textContent = 'Error: ' + d.error; + showErrorResponse(d, 'DataflowSpec Error'); return; } _specsData = d; @@ -3559,6 +3837,7 @@

Response

headEl.innerHTML = ''; bodyEl.innerHTML = '' + escapeHtml(data && data.error ? 'Error: ' + data.error : 'No data') + + (data && data.details ? _technicalDetailsHtml(data.details) : '') + ''; return; } @@ -3669,7 +3948,12 @@

Response

? 'Error: ' + pipelines.error : pipelines.length + ' SDP-META pipeline(s) found'; var tbody = document.getElementById('pipelineTableBody'); - if (!Array.isArray(pipelines) || pipelines.length === 0) { + if (!Array.isArray(pipelines)) { + tbody.innerHTML = 'Unable to load pipelines.'; + showErrorResponse(pipelines, 'Pipeline Monitor Error'); + return; + } + if (pipelines.length === 0) { tbody.innerHTML = 'No SDP-META pipelines found in this workspace.'; document.getElementById('eventsDrawer').style.display = 'none'; return; @@ -3734,8 +4018,14 @@

Response

.then(_parseJSON) .then(function(d) { hideLoading(); - showModal('

' + escapeHtml(d.error || (name + ' ' + (action === 'start' ? 'started' : 'stopped') + '.')) + '

', - action === 'start' ? 'Pipeline Started' : 'Pipeline Stopped'); + if (d.error) { + showErrorResponse(d, 'Pipeline action error'); + return; + } + showModal( + '

' + escapeHtml(name + ' ' + (action === 'start' ? 'started' : 'stopped') + '.') + '

', + action === 'start' ? 'Pipeline Started' : 'Pipeline Stopped' + ); setTimeout(loadPipelines, 1500); }) .catch(function(e) { hideLoading(); handleError(e, 'Pipeline action error'); }); @@ -3754,7 +4044,12 @@

Response

fetch('/api/pipelines/' + encodeURIComponent(pipelineId) + '/events') .then(_parseJSON) .then(function(events) { - if (!Array.isArray(events) || events.length === 0) { + if (!Array.isArray(events)) { + list.innerHTML = '
  • Unable to load events.
  • '; + showErrorResponse(events, 'Pipeline Events Error'); + return; + } + if (events.length === 0) { list.innerHTML = '
  • No events found.
  • '; return; } @@ -3790,6 +4085,7 @@

    Response

    if (!Array.isArray(cats)) { sel.innerHTML = ''; _browseStatus('Error loading catalogs: ' + (cats.error || 'unknown error'), true); + showErrorResponse(cats, 'Catalog Browser Error'); return; } sel.innerHTML = '' + @@ -3822,6 +4118,7 @@

    Response

    if (!Array.isArray(schemas)) { sel.innerHTML = ''; _browseStatus('Error loading schemas: ' + (schemas.error || 'unknown error'), true); + showErrorResponse(schemas, 'Catalog Browser Error'); return; } sel.innerHTML = '' + @@ -3852,6 +4149,7 @@

    Response

    if (!Array.isArray(tables)) { sel.innerHTML = ''; _browseStatus('Error loading tables: ' + (tables.error || 'unknown error'), true); + showErrorResponse(tables, 'Catalog Browser Error'); return; } sel.innerHTML = '' + @@ -3899,6 +4197,7 @@

    Response

    .then(function(d) { if (d.error) { statusEl.textContent = 'Error: ' + d.error; + showErrorResponse(d, 'Table Preview Error'); return; } var cols = d.columns || []; @@ -4021,7 +4320,7 @@

    Response

    .then(_parseJSON) .then(function(d) { hideLoading(); - if (d.error) { showModal('

    ' + escapeHtml(d.error) + '

    ', 'Load error'); return; } + if (d.error) { showErrorResponse(d, 'Load error'); return; } document.getElementById('specEditorTextarea').value = d.content; var fmtSel = document.getElementById('specFormat'); if (d.format) fmtSel.value = d.format; @@ -4048,7 +4347,7 @@

    Response

    .then(_parseJSON) .then(function(d) { hideLoading(); - if (d.error) { showModal('

    ' + escapeHtml(d.error) + '

    ', 'Save error'); return; } + if (d.error) { showErrorResponse(d, 'Save error'); return; } showModal('

    Saved ' + escapeHtml(d.path) + ' (' + d.bytes_written + ' bytes)

    ', 'Saved'); }) .catch(function(e) { hideLoading(); handleError(e, 'Save error'); }); diff --git a/demo/launch_acfs_demo.py b/demo/launch_acfs_demo.py index 3e749790..c87a2fe5 100644 --- a/demo/launch_acfs_demo.py +++ b/demo/launch_acfs_demo.py @@ -31,6 +31,7 @@ def run(self, runner_conf: SDPMetaRunnerConf): except Exception as e: print(e) traceback.print_exc() + raise # finally: # self.clean_up(runner_conf) @@ -60,7 +61,11 @@ def init_runner_conf(self) -> SDPMetaRunnerConf: def launch_workflow(self, runner_conf: SDPMetaRunnerConf): created_job = self.create_workflow_spec(runner_conf) - self.open_job_url(runner_conf, created_job) + self.open_job_url( + runner_conf, + created_job, + wait_for_completion=True, + ) def main(): diff --git a/demo/launch_af_cloudfiles_demo.py b/demo/launch_af_cloudfiles_demo.py index 4530c8de..6a22802a 100644 --- a/demo/launch_af_cloudfiles_demo.py +++ b/demo/launch_af_cloudfiles_demo.py @@ -32,6 +32,7 @@ def run(self, runner_conf: SDPMetaRunnerConf): except Exception as e: print(e) traceback.print_exc() + raise # finally: # self.clean_up(runner_conf) @@ -72,7 +73,11 @@ def init_runner_conf(self) -> SDPMetaRunnerConf: def launch_workflow(self, runner_conf: SDPMetaRunnerConf): created_job = self.create_workflow_spec(runner_conf) - self.open_job_url(runner_conf, created_job) + self.open_job_url( + runner_conf, + created_job, + wait_for_completion=True, + ) def main(): diff --git a/demo/launch_dais_demo.py b/demo/launch_dais_demo.py index 2b705c04..81e34407 100644 --- a/demo/launch_dais_demo.py +++ b/demo/launch_dais_demo.py @@ -71,6 +71,7 @@ def run(self, runner_conf: SDPMetaRunnerConf): except Exception as e: print(e) traceback.print_exc() + raise # finally: # self.clean_up(runner_conf) @@ -82,7 +83,11 @@ def launch_workflow(self, runner_conf: SDPMetaRunnerConf): - runner_conf: SDPMetaRunnerConf object """ created_job = self.create_daisdemo_workflow(runner_conf) - self.open_job_url(runner_conf, created_job) + self.open_job_url( + runner_conf, + created_job, + wait_for_completion=True, + ) def create_daisdemo_workflow(self, runner_conf: SDPMetaRunnerConf): """ diff --git a/demo/launch_silver_fanout_demo.py b/demo/launch_silver_fanout_demo.py index 45028666..ede0b07d 100644 --- a/demo/launch_silver_fanout_demo.py +++ b/demo/launch_silver_fanout_demo.py @@ -51,6 +51,7 @@ def run(self, runner_conf: SDPMetaRunnerConf): except Exception as e: print(e) traceback.print_exc() + raise # finally: # self.clean_up(runner_conf) @@ -92,7 +93,11 @@ def init_runner_conf(self) -> SDPMetaRunnerConf: def launch_workflow(self, runner_conf: SDPMetaRunnerConf): created_job = self.create_sfo_workflow_spec(runner_conf) - self.open_job_url(runner_conf, created_job) + self.open_job_url( + runner_conf, + created_job, + wait_for_completion=True, + ) def create_sfo_workflow_spec(self, runner_conf: SDPMetaRunnerConf): """ diff --git a/docs/docs/changelog.md b/docs/docs/changelog.md index 3fafdf2b..6b53737c 100644 --- a/docs/docs/changelog.md +++ b/docs/docs/changelog.md @@ -13,10 +13,13 @@ sidebar_position: 99 ### New Features - **Multiple pipeline topologies per DAB bundle** — `bundle-add-pipeline` adds independently configured bronze, silver, split, or combined pipelines with their own data-flow groups and target schemas. `bundle-validate` validates each pipeline and its job wiring independently. ([Issue #446](https://github.com/databrickslabs/sdp-meta/issues/446)) +- **Auto Loader schema-evolution demos** — focused and at-scale demos exercise additive schema evolution, strict-mode failure behavior, and multi-table onboarding from Unity Catalog Volumes. ([PR #490](https://github.com/databrickslabs/sdp-meta/pull/490)) ### Fixes - **Databricks SDK compatibility floor** — runtime and App dependencies now require `databricks-sdk>=0.138.0,<1`, and CI tests serverless onboarding environment models against exactly that minimum. ([Issue #457](https://github.com/databrickslabs/sdp-meta/issues/457)) +- **[Databricks App] Make demo launches asynchronous and feature the 100-table Auto Loader demo**. ([Issue #491](https://github.com/databrickslabs/sdp-meta/issues/491)) +- **Data-quality expectations without quarantine targets** — bronze and silver flows can apply DQ expectations when no quarantine table is configured, without referencing a missing quarantine target. ([PR #487](https://github.com/databrickslabs/sdp-meta/pull/487)) - **Databricks App deployment packaging** — deployment staging now includes the `examples/` wheel-build inputs on both macOS/Linux and Windows, preventing remote App startup failures while building the SDP-META wheel. - **DAB pipeline group ownership** — `bundle-add-pipeline` and `bundle-validate` reject duplicate ownership of the same data-flow group and layer, including conflicts introduced through target overrides. ([Issue #458](https://github.com/databrickslabs/sdp-meta/issues/458)) - **Customized DAB validation** — `bundle-validate` now honors default and environment-selected targets, accepts transitive bronze-to-silver task dependencies, and reports malformed targets or dependency cycles cleanly. ([Issue #459](https://github.com/databrickslabs/sdp-meta/issues/459)) diff --git a/integration_tests/run_integration_tests.py b/integration_tests/run_integration_tests.py index c5694474..ed771ed5 100644 --- a/integration_tests/run_integration_tests.py +++ b/integration_tests/run_integration_tests.py @@ -1196,12 +1196,22 @@ def download_test_results(self, runner_conf: SDPMetaRunnerConf): ) as output_file: output_file.write(ws_output_file.read()) - def open_job_url(self, runner_conf, created_job): + def open_job_url( + self, + runner_conf, + created_job, + *, + wait_for_completion=False, + timeout=timedelta(minutes=90), + ): runner_conf.job_id = created_job.job_id url = f"{self.ws.config.host}/jobs/{created_job.job_id}?o={self.ws.get_workspace_id()}" - self.ws.jobs.run_now(job_id=created_job.job_id) + waiter = self.ws.jobs.run_now(job_id=created_job.job_id) webbrowser.open(url) print(f"Job created successfully. job_id={created_job.job_id}, url={url}") + if wait_for_completion: + waiter.result(timeout=timeout) + return waiter def clean_up(self, runner_conf: SDPMetaRunnerConf): print("Cleaning up...") @@ -1476,18 +1486,17 @@ def get_workspace_api_client(profile=None) -> WorkspaceClient: ``~/.databrickscfg`` entry. This is the path ``run_integration_tests.py`` and the ``launch_*_demo.py`` scripts take when invoked from a developer's terminal. - 3. **Interactive fallback** — prompts for host + token via ``input()`` - so ad-hoc local runs without a configured profile still work. - Requires a real stdin; do NOT route App / CI traffic through - this branch (it raises ``EOFError`` when stdin is not a TTY). + 3. **SDK default credential chain** — when no profile is supplied, uses + ``WorkspaceClient()`` so ``DEFAULT`` profile, PAT environment variables, + Azure/GCP credentials, and other SDK-supported non-interactive auth + methods work consistently. Background App/CI launchers must never call + ``input()`` because stdin may not be attached. """ if os.environ.get("DATABRICKS_APP_PORT"): return WorkspaceClient() if profile: return WorkspaceClient(profile=profile) - return WorkspaceClient( - host=input("Databricks Workspace URL: "), token=input("Token: ") - ) + return WorkspaceClient() def main(): diff --git a/tests/test_app_demo_job_logs.py b/tests/test_app_demo_job_logs.py index 94759afb..20819aba 100644 --- a/tests/test_app_demo_job_logs.py +++ b/tests/test_app_demo_job_logs.py @@ -17,6 +17,7 @@ import os import sys import unittest +from unittest import mock _REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) _APP_DIR = os.path.join(_REPO_ROOT, "databricks_app") @@ -25,6 +26,16 @@ import app as app_mod # noqa: E402 (deliberate post-sys.path-insert import) import _jobs as _jobs_module # noqa: E402 +from routes import demo as demo_routes # noqa: E402 +from uc_preflight import PreflightResult # noqa: E402 +from demo.launch_acfs_demo import ApplyChangesFromSnapshotDemo # noqa: E402 +from demo.launch_af_cloudfiles_demo import SDPMETAFCFDemo # noqa: E402 +from demo.launch_dais_demo import SDPMETADAISDemo # noqa: E402 +from demo.launch_silver_fanout_demo import SDPMETASilverFanoutDemo # noqa: E402 +from integration_tests.run_integration_tests import ( # noqa: E402 + SDPMETARunner, + get_workspace_api_client, +) class JobLogsOffsetValidationTests(unittest.TestCase): @@ -94,5 +105,295 @@ def test_missing_token_returns_404_not_400(self): self.assertIn("not found", resp.get_json()["error"].lower()) +class DemoCapacityConfigurationTests(unittest.TestCase): + def test_positive_environment_value_overrides_default(self): + with mock.patch.dict( + os.environ, + {"SDP_META_MAX_ACTIVE_DEMOS": "12"}, + ): + self.assertEqual( + _jobs_module._positive_int_env( + "SDP_META_MAX_ACTIVE_DEMOS", + 8, + ), + 12, + ) + + def test_invalid_environment_values_use_default(self): + for value in ("invalid", "0", "-1"): + with self.subTest(value=value), mock.patch.dict( + os.environ, + {"SDP_META_MAX_ACTIVE_DEMOS": value}, + ): + self.assertEqual( + _jobs_module._positive_int_env( + "SDP_META_MAX_ACTIVE_DEMOS", + 8, + ), + 8, + ) + + +class DemoLauncherAuthenticationTests(unittest.TestCase): + """Demo subprocesses must never fall through to interactive auth.""" + + def setUp(self): + app_mod.app.testing = True + self.client = app_mod.app.test_client() + + def test_interactive_demo_keeps_real_completion_timeout(self): + interactive = demo_routes._DEMO_REGISTRY["demo_interactive"] + self.assertNotIn("--timeout-minutes", interactive["extra_args"]) + + def test_featured_demo_preserves_remote_resources(self): + featured = demo_routes._DEMO_REGISTRY["demo_at_scale_autoloader"] + self.assertIn("--keep-resources", featured["extra_args"]) + + @mock.patch("integration_tests.run_integration_tests.WorkspaceClient") + def test_default_profile_uses_sdk_credential_chain(self, mock_client): + with mock.patch.dict(os.environ, {}, clear=False): + os.environ.pop("DATABRICKS_APP_PORT", None) + get_workspace_api_client(None) + mock_client.assert_called_once_with() + + @mock.patch("integration_tests.run_integration_tests.WorkspaceClient") + def test_pat_environment_uses_sdk_credential_chain(self, mock_client): + with mock.patch.dict( + os.environ, + { + "DATABRICKS_HOST": "https://example.cloud.databricks.com", + "DATABRICKS_TOKEN": "test-token", + }, + clear=False, + ): + os.environ.pop("DATABRICKS_APP_PORT", None) + get_workspace_api_client(None) + mock_client.assert_called_once_with() + + @mock.patch("integration_tests.run_integration_tests.WorkspaceClient") + def test_named_profile_is_forwarded_to_sdk(self, mock_client): + with mock.patch.dict(os.environ, {}, clear=False): + os.environ.pop("DATABRICKS_APP_PORT", None) + get_workspace_api_client("fevm") + mock_client.assert_called_once_with(profile="fevm") + + @mock.patch("integration_tests.run_integration_tests.WorkspaceClient") + def test_app_auth_ignores_local_profile(self, mock_client): + with mock.patch.dict( + os.environ, + {"DATABRICKS_APP_PORT": "8000"}, + clear=False, + ): + get_workspace_api_client("fevm") + mock_client.assert_called_once_with() + + @mock.patch("integration_tests.run_integration_tests.webbrowser.open") + def test_open_job_url_waits_when_requested(self, _mock_browser): + runner = SDPMETARunner.__new__(SDPMETARunner) + runner.ws = mock.Mock() + runner.ws.config.host = "https://example.cloud.databricks.com" + runner.ws.get_workspace_id.return_value = "123" + waiter = runner.ws.jobs.run_now.return_value + runner_conf = mock.Mock() + created_job = mock.Mock(job_id=456) + + runner.open_job_url( + runner_conf, + created_job, + wait_for_completion=True, + ) + + waiter.result.assert_called_once() + + def test_registered_legacy_launchers_wait_for_remote_completion(self): + cases = ( + (SDPMETAFCFDemo, "create_workflow_spec"), + (ApplyChangesFromSnapshotDemo, "create_workflow_spec"), + (SDPMETASilverFanoutDemo, "create_sfo_workflow_spec"), + (SDPMETADAISDemo, "create_daisdemo_workflow"), + ) + for demo_class, create_method in cases: + with self.subTest(demo_class=demo_class.__name__): + demo = demo_class.__new__(demo_class) + created_job = mock.Mock() + setattr(demo, create_method, mock.Mock(return_value=created_job)) + demo.open_job_url = mock.Mock() + runner_conf = mock.Mock() + + demo.launch_workflow(runner_conf) + + demo.open_job_url.assert_called_once_with( + runner_conf, + created_job, + wait_for_completion=True, + ) + + @mock.patch("routes.demo.check_app_sp_grants_on_catalog") + def test_demo_capacity_returns_429(self, mock_preflight): + mock_preflight.return_value = PreflightResult( + ok=True, + uc_name="main", + sp_principal="test-user", + sp_display_name="Test User", + ) + token = _jobs_module._new_job_token(kind="demo") + try: + with mock.patch.object( + _jobs_module, + "_MAX_ACTIVE_DEMO_JOBS", + 1, + ): + resp = self.client.post( + "/rundemo", + json={ + "demo_name": "demo_at_scale_autoloader", + "uc_name": "main", + }, + ) + self.assertEqual(resp.status_code, 429) + self.assertEqual(resp.headers.get("Retry-After"), "30") + error = resp.get_json()["error"] + self.assertIn("active demo", error.lower()) + self.assertNotIn("1/1", error) + finally: + _jobs_module._jobs.pop(token, None) + + @mock.patch("routes.demo._run_command_in_background") + @mock.patch( + "routes.demo._jobs_module._new_job_token", + return_value="demo-token", + ) + @mock.patch("routes.demo.check_app_sp_grants_on_catalog") + def test_every_registered_demo_returns_background_token( + self, + mock_preflight, + _mock_new_token, + mock_run, + ): + mock_preflight.return_value = PreflightResult( + ok=True, + uc_name="main", + sp_principal="test-user", + sp_display_name="Test User", + ) + + for demo_name in demo_routes._DEMO_REGISTRY: + with self.subTest(demo_name=demo_name): + resp = self.client.post( + "/rundemo", + json={"demo_name": demo_name, "uc_name": "main"}, + ) + self.assertEqual(resp.status_code, 202) + self.assertEqual(resp.get_json()["token"], "demo-token") + + self.assertEqual(mock_run.call_count, len(demo_routes._DEMO_REGISTRY)) + self.assertTrue( + all( + call.kwargs["idle_timeout_seconds"] == 2 * 60 * 60 + for call in mock_run.call_args_list + ) + ) + + @mock.patch("routes.demo._run_command_in_background") + @mock.patch( + "routes.demo._jobs_module._new_job_token", + return_value="demo-token", + ) + @mock.patch("routes.demo.check_app_sp_grants_on_catalog") + def test_local_app_forwards_profile_without_enabling_app_path_shim( + self, + mock_preflight, + _mock_new_token, + mock_run, + ): + mock_preflight.return_value = PreflightResult( + ok=True, + uc_name="main", + sp_principal="test-user", + sp_display_name="Test User", + ) + with mock.patch.dict( + os.environ, + {"DATABRICKS_CONFIG_PROFILE": "fevm"}, + clear=False, + ): + os.environ.pop("DATABRICKS_APP_PORT", None) + resp = self.client.post( + "/rundemo", + json={ + "demo_name": "demo_at_scale_autoloader", + "uc_name": "main", + }, + ) + + self.assertEqual(resp.status_code, 202, resp.get_data(as_text=True)) + self.assertEqual(resp.get_json()["token"], "demo-token") + child_command = mock_run.call_args.kwargs["command"] + child_env = mock_run.call_args.kwargs["env"] + self.assertEqual(child_command[-2:], ["--profile", "fevm"]) + self.assertNotIn("DATABRICKS_APP_PORT", child_env) + + +class JobRetentionTests(unittest.TestCase): + """Completed jobs and verbose logs remain bounded.""" + + def test_log_cap_preserves_absolute_polling_offset(self): + token = _jobs_module._new_job_token() + job = _jobs_module._jobs[token] + try: + with mock.patch.object(_jobs_module, "_MAX_LOG_LINES", 2): + for line in ("one", "two", "three"): + _jobs_module._append_job_log( + job, + {"stream": "stdout", "line": line}, + ) + + client = app_mod.app.test_client() + payload = client.get( + f"/api/job/{token}/logs?offset=0" + ).get_json() + self.assertEqual( + [entry["line"] for entry in payload["logs"]], + ["two", "three"], + ) + self.assertEqual(payload["next_offset"], 3) + finally: + _jobs_module._jobs.pop(token, None) + + def test_expired_completed_job_is_evicted(self): + token = _jobs_module._new_job_token() + job = _jobs_module._jobs[token] + job["done"] = True + job["finished_at"] = 10 + + _jobs_module._prune_jobs( + now=10 + _jobs_module._COMPLETED_JOB_TTL_SECONDS + 1 + ) + + self.assertNotIn(token, _jobs_module._jobs) + + +class LegacyDemoFailurePropagationTests(unittest.TestCase): + """Launcher failures must produce a nonzero subprocess exit.""" + + def test_registered_legacy_launchers_reraise_failures(self): + for demo_class in ( + SDPMETAFCFDemo, + ApplyChangesFromSnapshotDemo, + SDPMETASilverFanoutDemo, + SDPMETADAISDemo, + ): + with self.subTest(demo_class=demo_class.__name__): + demo = demo_class.__new__(demo_class) + demo.init_sdp_meta_runner_conf = mock.Mock( + side_effect=RuntimeError("simulated failure") + ) + with self.assertRaisesRegex( + RuntimeError, + "simulated failure", + ): + demo.run(mock.Mock()) + + if __name__ == "__main__": # pragma: no cover unittest.main() diff --git a/tests/test_app_landing_page.py b/tests/test_app_landing_page.py index c20a0993..308c8038 100644 --- a/tests/test_app_landing_page.py +++ b/tests/test_app_landing_page.py @@ -24,6 +24,7 @@ import os import sys import unittest +from unittest import mock _REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) _APP_DIR = os.path.join(_REPO_ROOT, "databricks_app") @@ -33,6 +34,18 @@ import app as app_mod # noqa: E402 +def _raise_test_exception(): + raise RuntimeError("raw SDK failure: principal lacks USE CATALOG") + + +if "_test_unhandled_error" not in app_mod.app.view_functions: + app_mod.app.add_url_rule( + "/__test_unhandled_error", + endpoint="_test_unhandled_error", + view_func=_raise_test_exception, + ) + + class LandingPageVersionInjectionTests(unittest.TestCase): """``GET /`` must render the installed package's version in the sidebar footer.""" @@ -84,5 +97,176 @@ def test_context_processor_exposes_app_version(self): self.assertEqual(ctx["app_version"], __version__) +class LandingPageFeaturedDemoTests(unittest.TestCase): + """The Demos panel exposes one compact launcher dropdown.""" + + def setUp(self): + self.client = app_mod.app.test_client() + + def test_featured_demo_dropdown_lists_every_runnable_demo(self): + html = self.client.get("/").get_data(as_text=True) + + self.assertIn('id="featured_demo_select"', html) + for demo_name in ( + "demo_at_scale_autoloader", + "demo_interactive", + "demo_cloudfiles", + "demo_acf", + "demo_silverfanout", + "demo_dias", + ): + self.assertIn(f'value="{demo_name}"', html) + + def test_at_scale_autoloader_is_the_default_featured_demo(self): + html = self.client.get("/").get_data(as_text=True) + + self.assertIn( + '', + html, + ) + self.assertIn('id="runFeaturedDemo"', html) + + def test_onboarding_picker_links_to_featured_at_scale_demo(self): + html = self.client.get("/").get_data(as_text=True) + + self.assertIn("__featured_at_scale_autoloader__", html) + self.assertIn("Featured \\u2014 At-scale Auto Loader (100 tables)", html) + self.assertIn("launcher: 'demo_at_scale_autoloader'", html) + self.assertIn("bundledSelect.selectedIndex = 0;", html) + self.assertNotIn( + "var preferenceOrder = ['onboarding', 'onboarding_cars'];", + html, + ) + self.assertIn('id="onboardingSubmitBtn"', html) + self.assertIn("bundledPicker.value === '__featured_at_scale_autoloader__'", html) + self.assertIn("runDemo('demo_at_scale_autoloader');", html) + self.assertIn("'▶ Run Featured Demo'", html) + + def test_progress_modal_can_close_while_demo_runs(self): + html = self.client.get("/").get_data(as_text=True) + + self.assertIn( + 'id="progressClose" onclick="closeProgressModal()">Close', + html, + ) + self.assertNotIn('id="progressClose" disabled', html) + self.assertIn( + "document.getElementById('progressClose').disabled = false;", + html, + ) + close_body = html.split( + "function closeProgressModal() {", + 1, + )[1].split( + "function _startDots()", + 1, + )[0] + self.assertNotIn("_cancelProgressSession()", close_body) + self.assertNotIn("clearTimeout", close_body) + + def test_progress_polling_owns_session_and_retries_transient_errors(self): + html = self.client.get("/").get_data(as_text=True) + + self.assertIn("new AbortController()", html) + self.assertIn("_progressSession === session", html) + self.assertNotIn("_cancelProgressSession()", html) + self.assertIn("var maxRetries = 4;", html) + self.assertIn("err.status >= 500", html) + self.assertIn("Math.pow(2, session.retryCount - 1) * 1000", html) + self.assertIn("err.name === 'AbortError'", html) + fail_body = html.split( + "function failPolling(err) {", + 1, + )[1].split( + "function poll()", + 1, + )[0] + self.assertIn("result: {", fail_body) + self.assertIn("details: {", fail_body) + self.assertIn("status: err.status || null", fail_body) + self.assertIn("retries: session.retryCount", fail_body) + + def test_new_operation_reopens_instead_of_abandoning_active_progress(self): + html = self.client.get("/").get_data(as_text=True) + + claim_body = html.split( + "function _claimProgressOperation() {", + 1, + )[1].split( + "function _releaseProgressClaim()", + 1, + )[0] + self.assertIn("_progressSession && !_progressSession.finished", claim_body) + self.assertIn("progressModal", claim_body) + self.assertIn("return false;", claim_body) + self.assertGreaterEqual( + html.count("if (!_claimProgressOperation()) return;"), + 3, + ) + + +class FriendlyErrorTests(unittest.TestCase): + def setUp(self): + self.client = app_mod.app.test_client() + + def test_unhandled_error_has_safe_headline_and_expandable_details(self): + response = self.client.get("/__test_unhandled_error") + + self.assertEqual(response.status_code, 500) + body = response.get_json() + self.assertNotIn("raw SDK failure", body["error"]) + self.assertEqual(body["details"]["exception_type"], "RuntimeError") + self.assertIn("principal lacks USE CATALOG", body["details"]["message"]) + self.assertTrue(body["details"]["request_id"]) + + def test_caught_server_error_is_normalized(self): + with app_mod.app.test_request_context(): + response = app_mod.app.make_response( + (app_mod.jsonify({"error": "raw SDK permission failure"}), 500) + ) + response = app_mod.add_security_headers(response) + + body = response.get_json() + self.assertNotIn("raw SDK permission failure", body["error"]) + self.assertEqual( + body["details"]["message"], + "raw SDK permission failure", + ) + + @mock.patch( + "databricks.sdk.WorkspaceClient", + side_effect=RuntimeError("sensitive SDK credential detail"), + ) + def test_warehouse_configuration_does_not_expose_sdk_exception( + self, + _mock_client, + ): + response = self.client.post( + "/api/warehouse/configure", + json={"mode": "existing", "warehouse_id": "warehouse-id"}, + ) + + self.assertEqual(response.status_code, 400) + response_text = response.get_data(as_text=True) + self.assertNotIn("sensitive SDK credential detail", response_text) + self.assertIn( + "could not use that SQL warehouse", + response.get_json()["error"], + ) + + def test_landing_page_renders_collapsed_technical_details(self): + html = self.client.get("/").get_data(as_text=True) + + self.assertIn("function showErrorResponse(data, title)", html) + self.assertIn("