Some checks failed
CI / skinny-install (api) (push) Successful in 24s
CI / lint-test (push) Successful in 1m18s
CI / skinny-install (bib) (push) Successful in 30s
CI / skinny-install (bls) (push) Successful in 28s
CI / skinny-install (aco) (push) Successful in 50s
CI / skinny-install (bcda) (push) Successful in 33s
CI / skinny-install (ccw) (push) Successful in 32s
CI / skinny-install (cms) (push) Successful in 24s
CI / skinny-install (cli) (push) Successful in 36s
CI / skinny-install (conf) (push) Successful in 33s
CI / skinny-install (opps) (push) Successful in 30s
CI / skinny-install (perf) (push) Successful in 31s
CI / skinny-install (pfs) (push) Successful in 34s
CI / skinny-install (rex) (push) Has been cancelled
CI / skinny-install (aco) (pull_request) Successful in 51s
CI / lint-test (pull_request) Successful in 1m20s
CI / skinny-install (api) (pull_request) Successful in 29s
CI / skinny-install (bcda) (pull_request) Successful in 30s
CI / skinny-install (bib) (pull_request) Successful in 32s
CI / skinny-install (bls) (pull_request) Successful in 26s
CI / skinny-install (ccw) (pull_request) Successful in 27s
CI / skinny-install (cli) (pull_request) Successful in 29s
CI / skinny-install (cms) (pull_request) Successful in 34s
CI / skinny-install (conf) (pull_request) Successful in 22s
CI / skinny-install (opps) (pull_request) Successful in 33s
CI / skinny-install (perf) (pull_request) Successful in 31s
CI / skinny-install (pfs) (pull_request) Successful in 34s
CI / skinny-install (rex) (pull_request) Successful in 32s
Infra CI / notebooks (push) Successful in 8s
Infra CI / zotero (push) Successful in 6s
Infra CI / docs (push) Failing after 12s
Infra CI / api (push) Successful in 10s
Infra CI / mc (push) Successful in 9s
Infra CI / notebooks (pull_request) Successful in 7s
Infra CI / zotero (pull_request) Successful in 6s
Infra CI / docs (pull_request) Failing after 5s
Infra CI / api (pull_request) Successful in 6s
Infra CI / mc (pull_request) Successful in 6s
Expand Databricks SDK usage from 6 to 18 WorkspaceClient services. Tier 1 — Governance & Compliance: - governance.py: declarative GovernancePolicy with apply + audit drift - UnityClient: grants CRUD (ws.grants), secret management (ws.secrets), table constraints PK/FK (ws.table_constraints) - sync_secrets.py: push env vars to Databricks scopes per stack.toml Tier 2 — Job Orchestration: - jobs.py: JobManager translates Pipeline → Databricks Job with task dependencies, run_now, get_run_status, list/delete - UnityClient: warehouse lookup by name (ws.warehouses), start/stop Tier 3 — Data Quality & Monitoring: - quality.py: setup_monitors creates Lakehouse Monitoring profiles, list_monitors, run_refresh (ws.quality_monitors) - UnityClient: system schema access, table lineage queries, audit log queries (ws.system_schemas, ws.statement_execution) Config: stack.toml gains [databricks.warehouse], [databricks.secrets], [databricks.governance], [databricks.quality] sections. 23 new tests, all passing.
89 lines
2.9 KiB
Python
89 lines
2.9 KiB
Python
"""Tests for aco.lake.jobs — programmatic job management."""
|
|
|
|
from __future__ import annotations
|
|
|
|
from unittest.mock import MagicMock
|
|
|
|
from aco.lake.jobs import JobManager, RunStatus
|
|
|
|
|
|
class TestRunStatus:
|
|
def test_terminal_states(self):
|
|
assert RunStatus(run_id=1, state="TERMINATED").is_terminal
|
|
assert RunStatus(run_id=1, state="SKIPPED").is_terminal
|
|
assert not RunStatus(run_id=1, state="RUNNING").is_terminal
|
|
|
|
def test_success(self):
|
|
assert RunStatus(run_id=1, state="TERMINATED", result_state="SUCCESS").succeeded
|
|
assert not RunStatus(
|
|
run_id=1, state="TERMINATED", result_state="FAILED"
|
|
).succeeded
|
|
|
|
|
|
class TestJobManager:
|
|
def _make_manager(self):
|
|
client = MagicMock()
|
|
return JobManager(client, catalog="aco_dev"), client
|
|
|
|
def test_create_pipeline_job(self):
|
|
mgr, client = self._make_manager()
|
|
client._ws.jobs.create.return_value = MagicMock(job_id=42)
|
|
|
|
job_id = mgr.create_pipeline_job("readmissions")
|
|
assert job_id == 42
|
|
client._ws.jobs.create.assert_called_once()
|
|
|
|
# Verify task structure
|
|
call_kwargs = client._ws.jobs.create.call_args.kwargs
|
|
assert call_kwargs["name"] == "stack_readmissions"
|
|
assert len(call_kwargs["tasks"]) >= 1
|
|
task = call_kwargs["tasks"][0]
|
|
assert task.task_key == "readmissions"
|
|
assert task.python_wheel_task.package_name == "stack"
|
|
|
|
def test_create_with_schedule(self):
|
|
mgr, client = self._make_manager()
|
|
client._ws.jobs.create.return_value = MagicMock(job_id=99)
|
|
|
|
mgr.create_pipeline_job("core", schedule="0 0 6 * * ?")
|
|
call_kwargs = client._ws.jobs.create.call_args.kwargs
|
|
assert call_kwargs["schedule"] is not None
|
|
|
|
def test_run_now(self):
|
|
mgr, client = self._make_manager()
|
|
client._ws.jobs.run_now.return_value = MagicMock(run_id=123)
|
|
|
|
run_id = mgr.run_now(42)
|
|
assert run_id == 123
|
|
client._ws.jobs.run_now.assert_called_once_with(job_id=42)
|
|
|
|
def test_get_run_status(self):
|
|
mgr, client = self._make_manager()
|
|
state = MagicMock()
|
|
state.life_cycle_state = "RUNNING"
|
|
state.result_state = None
|
|
state.state_message = "In progress"
|
|
run = MagicMock()
|
|
run.state = state
|
|
client._ws.jobs.get_run.return_value = run
|
|
|
|
status = mgr.get_run_status(123)
|
|
assert status.state == "RUNNING"
|
|
assert not status.is_terminal
|
|
|
|
def test_list_jobs(self):
|
|
mgr, client = self._make_manager()
|
|
job = MagicMock()
|
|
job.job_id = 1
|
|
job.settings.name = "stack_core"
|
|
client._ws.jobs.list.return_value = [job]
|
|
|
|
jobs = mgr.list_jobs()
|
|
assert len(jobs) == 1
|
|
assert jobs[0]["name"] == "stack_core"
|
|
|
|
def test_delete_job(self):
|
|
mgr, client = self._make_manager()
|
|
mgr.delete_job(42)
|
|
client._ws.jobs.delete.assert_called_once_with(job_id=42)
|