From 2b7701d845c1b0db3b26c26e37b83bf1e43d32b7 Mon Sep 17 00:00:00 2001 From: Rada Kamysheva Date: Tue, 28 Jul 2026 08:09:43 +0000 Subject: [PATCH 1/2] testserver: roll task outcomes up into the run state The fake workspace reported every run as TERMINATED SUCCESS, overwriting the FAILED state it had just recorded for a task it executed locally. Runs now report the outcome their tasks add up to, so job_runs/failed_run can exercise a real failing run instead of stubbing runs/get through test.toml. Tasks whose code the fake workspace does not have (an immutable deployment uploads a snapshot zip it never unpacks) are left successful: that gap is in the fake workspace, not in the job under test. --- .../job_runs/failed_run/databricks.yml | 12 +++- .../resources/job_runs/failed_run/fail.py | 4 ++ .../resources/job_runs/failed_run/output.txt | 7 ++- .../resources/job_runs/failed_run/test.toml | 29 --------- libs/testserver/jobs.go | 62 ++++++++++++++----- libs/testserver/jobs_test.go | 60 +++++++++++++++++- 6 files changed, 123 insertions(+), 51 deletions(-) create mode 100644 acceptance/bundle/resources/job_runs/failed_run/fail.py diff --git a/acceptance/bundle/resources/job_runs/failed_run/databricks.yml b/acceptance/bundle/resources/job_runs/failed_run/databricks.yml index faad7eaf049..659976f5765 100644 --- a/acceptance/bundle/resources/job_runs/failed_run/databricks.yml +++ b/acceptance/bundle/resources/job_runs/failed_run/databricks.yml @@ -6,9 +6,17 @@ resources: my_job: name: my-job tasks: + # The test server runs this task locally; the script exits non-zero, so + # the task and with it the whole run fail. - task_key: main - notebook_task: - notebook_path: /Workspace/missing-notebook + spark_python_task: + python_file: ./fail.py + environment_key: default + + environments: + - environment_key: default + spec: + client: "2" # Depends on my_run's result_state, so the failing run aborts the deploy # before this job is created. diff --git a/acceptance/bundle/resources/job_runs/failed_run/fail.py b/acceptance/bundle/resources/job_runs/failed_run/fail.py new file mode 100644 index 00000000000..3262aa05529 --- /dev/null +++ b/acceptance/bundle/resources/job_runs/failed_run/fail.py @@ -0,0 +1,4 @@ +import sys + +print("intentional failure", file=sys.stderr) +sys.exit(1) diff --git a/acceptance/bundle/resources/job_runs/failed_run/output.txt b/acceptance/bundle/resources/job_runs/failed_run/output.txt index 8012b0aafe8..ea941facaec 100644 --- a/acceptance/bundle/resources/job_runs/failed_run/output.txt +++ b/acceptance/bundle/resources/job_runs/failed_run/output.txt @@ -3,9 +3,14 @@ >>> [CLI] bundle deploy Uploading bundle files to /Workspace/Users/[USERNAME]/.bundle/job-runs-failed-run/default/files... Deploying resources... +resources.job_runs.my_run: Run URL: [DATABRICKS_URL]/jobs/[NUMID]/runs/[NUMID]?o=[NUMID] +resources.job_runs.my_run: [TIMESTAMP] "my-job" RUNNING resources.job_runs.my_run: [TIMESTAMP] "my-job" TERMINATED FAILED task main failed Error: cannot create resources.job_runs.my_run: job run [NUMID] did not succeed: FAILED: task main failed -task "main": notebook not found in workspace: /Workspace/missing-notebook +task "main": spark python task execution failed: exit status 1 +intentional failure + +run page: [DATABRICKS_URL]/jobs/[NUMID]/runs/[NUMID]?o=[NUMID] redeploying reports this same run; set rerun_token to a new value to run the job again Error: cannot create resources.jobs.downstream_job: dependency failed: resources.job_runs.my_run diff --git a/acceptance/bundle/resources/job_runs/failed_run/test.toml b/acceptance/bundle/resources/job_runs/failed_run/test.toml index 998e7d1ce39..1c9be9d5bb9 100644 --- a/acceptance/bundle/resources/job_runs/failed_run/test.toml +++ b/acceptance/bundle/resources/job_runs/failed_run/test.toml @@ -1,31 +1,2 @@ # The deploy fails mid-way, leaving local deployment state behind. Ignore = [".databricks"] - -# The fake workspace reports every run as succeeded, so a failing run has to be -# stubbed. run_page_url is omitted because a static stub cannot know the test -# server's address; wait_output covers the run page link. -[[Server]] -Pattern = "GET /api/2.2/jobs/runs/get" -Response.Body = ''' -{ - "run_id": 12345678, - "job_id": 87654321, - "run_name": "my-job", - "state": { - "life_cycle_state": "TERMINATED", - "result_state": "FAILED", - "state_message": "task main failed" - }, - "tasks": [ - { - "task_key": "main", - "run_id": 12345679, - "state": {"life_cycle_state": "TERMINATED", "result_state": "FAILED"} - } - ] -} -''' - -[[Server]] -Pattern = "GET /api/2.2/jobs/runs/get-output" -Response.Body = '{"error": "notebook not found in workspace: /Workspace/missing-notebook"}' diff --git a/libs/testserver/jobs.go b/libs/testserver/jobs.go index a83e60ee354..6730468364a 100644 --- a/libs/testserver/jobs.go +++ b/libs/testserver/jobs.go @@ -20,6 +20,13 @@ import ( const missingJobGitProviderMessage = "git_source.git_provider must be one of: github,gitlab,bitbucketcloud,gitlabenterpriseedition,bitbucketserver,azuredevopsservices,githubenterprise,awscodecommit" +// errNoCodeInWorkspace marks a task whose code the fake workspace does not have, +// so there is nothing to execute locally: an immutable deployment, for example, +// uploads the bundle as a snapshot zip that the fake workspace never unpacks. +// Such a task is left successful, because the gap is in the fake workspace +// rather than in the job under test. +var errNoCodeInWorkspace = errors.New("task code is not in the workspace") + // venvPython returns the path to the Python executable in a venv. // On Unix: venv/bin/python // On Windows: venv\Scripts\python.exe @@ -401,12 +408,15 @@ func (s *FakeWorkspace) JobsRunNow(req Request) Response { logs, err = s.executeSparkPythonTask(t) } - if err != nil { + switch { + case errors.Is(err, errNoCodeInWorkspace): + // Nothing was executed; see errNoCodeInWorkspace. + case err != nil: taskRun.State.ResultState = jobs.RunResultStateFailed s.JobRunOutputs[taskRunId] = jobs.RunOutput{ Error: err.Error(), } - } else if logs != "" { + case logs != "": s.JobRunOutputs[taskRunId] = jobs.RunOutput{ Logs: logs, } @@ -617,7 +627,7 @@ func (s *FakeWorkspace) executePythonWheelTask(jobSettings *jobs.JobSettings, ta } data := s.files[whlPath].Data if len(data) == 0 { - return "", fmt.Errorf("wheel file not found in workspace: %s", whlPath) + return "", fmt.Errorf("%w: wheel file not found in workspace: %s", errNoCodeInWorkspace, whlPath) } localPath := filepath.Join(env.dir, filepath.Base(whlPath)) if err := os.WriteFile(localPath, data, 0o644); err != nil { @@ -636,7 +646,7 @@ func (s *FakeWorkspace) executePythonWheelTask(jobSettings *jobs.JobSettings, ta } if len(env.installedLibs) == 0 { - return "", errors.New("no wheel libraries found in task") + return "", fmt.Errorf("%w: no wheel libraries found in task", errNoCodeInWorkspace) } // Run the entry point using runpy with sys.argv[0] set to the package name, @@ -682,7 +692,7 @@ func (s *FakeWorkspace) executeNotebookTask(task jobs.Task, notebookParams map[s notebookData = s.files[notebookPath+".py"].Data } if len(notebookData) == 0 { - return "", fmt.Errorf("notebook not found in workspace: %s (also tried .py)", notebookPath) + return "", fmt.Errorf("%w: notebook not found in workspace: %s (also tried .py)", errNoCodeInWorkspace, notebookPath) } // Create a temporary Python environment for notebook execution @@ -768,7 +778,7 @@ func (s *FakeWorkspace) executeSparkPythonTask(task jobs.Task) (string, error) { pythonData := s.files[pythonPath].Data if len(pythonData) == 0 { - return "", fmt.Errorf("python file not found in workspace: %s", pythonPath) + return "", fmt.Errorf("%w: python file not found in workspace: %s", errNoCodeInWorkspace, pythonPath) } env, cleanup, err := s.getOrCreateClusterEnv(task) @@ -866,6 +876,32 @@ func sparkVersionToPython(task jobs.Task) string { return "3.10" } +// terminateRun completes the run and rolls its task outcomes up into the +// run-level state, the way the Jobs API does: a task that JobsRunNow executed +// and that failed fails the whole run. +func terminateRun(run *jobs.Run) { + for i := range run.Tasks { + // Tasks the fake workspace does not execute (jobs/runs/submit) are still + // running at this point; complete them before rolling the run up. + if run.Tasks[i].State.LifeCycleState != jobs.RunLifeCycleStateTerminated { + run.Tasks[i].State.LifeCycleState = jobs.RunLifeCycleStateTerminated + run.Tasks[i].State.ResultState = jobs.RunResultStateSuccess + } + } + + run.State = &jobs.RunState{ + LifeCycleState: jobs.RunLifeCycleStateTerminated, + ResultState: jobs.RunResultStateSuccess, + } + for _, task := range run.Tasks { + if task.State.ResultState != jobs.RunResultStateSuccess { + run.State.ResultState = task.State.ResultState + run.State.StateMessage = fmt.Sprintf("task %s failed", task.TaskKey) + return + } + } +} + func (s *FakeWorkspace) JobsGetRun(req Request) Response { runId := req.URL.Query().Get("run_id") runIdInt, err := strconv.ParseInt(runId, 10, 64) @@ -883,19 +919,11 @@ func (s *FakeWorkspace) JobsGetRun(req Request) Response { return Response{StatusCode: 404} } - // Simulate cloud behavior: first poll returns RUNNING, next returns TERMINATED SUCCESS. + // Simulate cloud behavior: first poll returns RUNNING, next returns the + // terminal state the tasks add up to. if run.State.LifeCycleState == jobs.RunLifeCycleStateRunning { // Transition stored state to TERMINATED for the next poll. - run.State = &jobs.RunState{ - LifeCycleState: jobs.RunLifeCycleStateTerminated, - ResultState: jobs.RunResultStateSuccess, - } - for i := range run.Tasks { - run.Tasks[i].State = &jobs.RunState{ - LifeCycleState: jobs.RunLifeCycleStateTerminated, - ResultState: jobs.RunResultStateSuccess, - } - } + terminateRun(&run) s.JobRuns[runIdInt] = run // Return RUNNING for this poll (before the transition). diff --git a/libs/testserver/jobs_test.go b/libs/testserver/jobs_test.go index 63aa0335290..76eb0910ff5 100644 --- a/libs/testserver/jobs_test.go +++ b/libs/testserver/jobs_test.go @@ -87,9 +87,9 @@ func TestJobsSubmit_RunReachesTerminalStateOnPoll(t *testing.T) { assert.Equal(t, jobs.RunResultStateSuccess, second.State.ResultState) } -func createJob(t *testing.T, workspace *FakeWorkspace) int64 { +func createJob(t *testing.T, workspace *FakeWorkspace, tasks ...jobs.Task) int64 { t.Helper() - body, err := json.Marshal(jobs.CreateJob{Name: "my-job"}) + body, err := json.Marshal(jobs.CreateJob{Name: "my-job", Tasks: tasks}) require.NoError(t, err) response := workspace.JobsCreate(Request{Body: body}) @@ -134,6 +134,62 @@ func TestJobsRunNow_IdempotencyTokenTombstonedAfterDelete(t *testing.T) { assert.Contains(t, third.Body.(string), "has been deleted") } +func terminatedTask(taskKey string, result jobs.RunResultState) jobs.RunTask { + return jobs.RunTask{ + TaskKey: taskKey, + State: &jobs.RunState{ + LifeCycleState: jobs.RunLifeCycleStateTerminated, + ResultState: result, + }, + } +} + +// A run whose task failed must not report SUCCESS: the failure is what a test +// waiting on the run is there to observe. +func TestTerminateRun_FailedTaskFailsTheRun(t *testing.T) { + run := jobs.Run{Tasks: []jobs.RunTask{ + terminatedTask("first", jobs.RunResultStateSuccess), + terminatedTask("second", jobs.RunResultStateFailed), + }} + + terminateRun(&run) + + assert.Equal(t, jobs.RunLifeCycleStateTerminated, run.State.LifeCycleState) + assert.Equal(t, jobs.RunResultStateFailed, run.State.ResultState) + assert.Equal(t, "task second failed", run.State.StateMessage) +} + +func TestTerminateRun_CompletesTasksThatAreStillRunning(t *testing.T) { + // jobs/runs/submit records its tasks as running: they are never executed. + run := jobs.Run{Tasks: []jobs.RunTask{ + {TaskKey: "main", State: &jobs.RunState{LifeCycleState: jobs.RunLifeCycleStateRunning}}, + }} + + terminateRun(&run) + + assert.Equal(t, jobs.RunResultStateSuccess, run.State.ResultState) + assert.Empty(t, run.State.StateMessage) + assert.Equal(t, jobs.RunResultStateSuccess, run.Tasks[0].State.ResultState) +} + +// The fake workspace has nothing to execute for a task whose code it does not +// have, so the task is left successful (see errNoCodeInWorkspace). +func TestJobsGetRun_TaskWithoutCodeDoesNotFailTheRun(t *testing.T) { + workspace := NewFakeWorkspace("http://test", "dbapi123") + jobID := createJob(t, workspace, jobs.Task{ + TaskKey: "main", + NotebookTask: &jobs.NotebookTask{NotebookPath: "/missing-notebook"}, + }) + + response := runNow(t, workspace, jobs.RunNow{JobId: jobID}) + require.Equal(t, 0, response.StatusCode) + runID := response.Body.(jobs.RunNowResponse).RunId + + // The first poll reports RUNNING, the second the terminal state. + require.Equal(t, jobs.RunLifeCycleStateRunning, getRun(t, workspace, runID).State.LifeCycleState) + assert.Equal(t, jobs.RunResultStateSuccess, getRun(t, workspace, runID).State.ResultState) +} + func TestJobsSubmit_RejectsInvalidGitProvider(t *testing.T) { workspace := NewFakeWorkspace("http://test", "dbapi123") From a3d077df2bcaa867636022185fed218fe95d16d0 Mon Sep 17 00:00:00 2001 From: Rada Kamysheva Date: Tue, 28 Jul 2026 08:19:17 +0000 Subject: [PATCH 2/2] testserver: shorten the comments added by this PR --- .../job_runs/failed_run/databricks.yml | 4 ++-- libs/testserver/jobs.go | 21 +++++++------------ libs/testserver/jobs_test.go | 6 ++---- 3 files changed, 12 insertions(+), 19 deletions(-) diff --git a/acceptance/bundle/resources/job_runs/failed_run/databricks.yml b/acceptance/bundle/resources/job_runs/failed_run/databricks.yml index 659976f5765..8a538daf853 100644 --- a/acceptance/bundle/resources/job_runs/failed_run/databricks.yml +++ b/acceptance/bundle/resources/job_runs/failed_run/databricks.yml @@ -6,8 +6,8 @@ resources: my_job: name: my-job tasks: - # The test server runs this task locally; the script exits non-zero, so - # the task and with it the whole run fail. + # The test server runs this locally; the script exits non-zero, which + # fails the task and with it the run. - task_key: main spark_python_task: python_file: ./fail.py diff --git a/libs/testserver/jobs.go b/libs/testserver/jobs.go index 6730468364a..59e7ee591c3 100644 --- a/libs/testserver/jobs.go +++ b/libs/testserver/jobs.go @@ -20,11 +20,9 @@ import ( const missingJobGitProviderMessage = "git_source.git_provider must be one of: github,gitlab,bitbucketcloud,gitlabenterpriseedition,bitbucketserver,azuredevopsservices,githubenterprise,awscodecommit" -// errNoCodeInWorkspace marks a task whose code the fake workspace does not have, -// so there is nothing to execute locally: an immutable deployment, for example, -// uploads the bundle as a snapshot zip that the fake workspace never unpacks. -// Such a task is left successful, because the gap is in the fake workspace -// rather than in the job under test. +// errNoCodeInWorkspace marks a task there is nothing to execute for, e.g. +// because an immutable deployment uploaded the code as a snapshot zip this +// server never unpacks. The gap is here, not in the job, so the task succeeds. var errNoCodeInWorkspace = errors.New("task code is not in the workspace") // venvPython returns the path to the Python executable in a venv. @@ -410,7 +408,7 @@ func (s *FakeWorkspace) JobsRunNow(req Request) Response { switch { case errors.Is(err, errNoCodeInWorkspace): - // Nothing was executed; see errNoCodeInWorkspace. + // Nothing ran, so the task keeps its SUCCESS state. case err != nil: taskRun.State.ResultState = jobs.RunResultStateFailed s.JobRunOutputs[taskRunId] = jobs.RunOutput{ @@ -876,13 +874,11 @@ func sparkVersionToPython(task jobs.Task) string { return "3.10" } -// terminateRun completes the run and rolls its task outcomes up into the -// run-level state, the way the Jobs API does: a task that JobsRunNow executed -// and that failed fails the whole run. +// terminateRun completes the run, rolling task outcomes up into the run-level +// state the way the Jobs API does: one failed task fails the whole run. func terminateRun(run *jobs.Run) { for i := range run.Tasks { - // Tasks the fake workspace does not execute (jobs/runs/submit) are still - // running at this point; complete them before rolling the run up. + // Tasks that were never executed (jobs/runs/submit) are still running. if run.Tasks[i].State.LifeCycleState != jobs.RunLifeCycleStateTerminated { run.Tasks[i].State.LifeCycleState = jobs.RunLifeCycleStateTerminated run.Tasks[i].State.ResultState = jobs.RunResultStateSuccess @@ -919,8 +915,7 @@ func (s *FakeWorkspace) JobsGetRun(req Request) Response { return Response{StatusCode: 404} } - // Simulate cloud behavior: first poll returns RUNNING, next returns the - // terminal state the tasks add up to. + // Simulate cloud behavior: first poll returns RUNNING, next the terminal state. if run.State.LifeCycleState == jobs.RunLifeCycleStateRunning { // Transition stored state to TERMINATED for the next poll. terminateRun(&run) diff --git a/libs/testserver/jobs_test.go b/libs/testserver/jobs_test.go index 76eb0910ff5..83d1c5595f6 100644 --- a/libs/testserver/jobs_test.go +++ b/libs/testserver/jobs_test.go @@ -144,8 +144,6 @@ func terminatedTask(taskKey string, result jobs.RunResultState) jobs.RunTask { } } -// A run whose task failed must not report SUCCESS: the failure is what a test -// waiting on the run is there to observe. func TestTerminateRun_FailedTaskFailsTheRun(t *testing.T) { run := jobs.Run{Tasks: []jobs.RunTask{ terminatedTask("first", jobs.RunResultStateSuccess), @@ -172,8 +170,8 @@ func TestTerminateRun_CompletesTasksThatAreStillRunning(t *testing.T) { assert.Equal(t, jobs.RunResultStateSuccess, run.Tasks[0].State.ResultState) } -// The fake workspace has nothing to execute for a task whose code it does not -// have, so the task is left successful (see errNoCodeInWorkspace). +// See errNoCodeInWorkspace: a missing notebook is this server's gap, not a +// failure of the job. func TestJobsGetRun_TaskWithoutCodeDoesNotFailTheRun(t *testing.T) { workspace := NewFakeWorkspace("http://test", "dbapi123") jobID := createJob(t, workspace, jobs.Task{