From 2f41009e28b743663e085174f3a5a6381474545d Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Mon, 16 Mar 2026 17:10:31 +0100 Subject: [PATCH 01/12] feat: add wait command with non-zero exit on failure Co-Authored-By: Claude Opus 4.6 (1M context) --- src/gridtk/cli.py | 70 ++++++++++++++++++++++++++++++++++++++++++++ tests/test_gridtk.py | 36 +++++++++++++++++++++++ 2 files changed, 106 insertions(+) diff --git a/src/gridtk/cli.py b/src/gridtk/cli.py index da33737..8b0bea3 100644 --- a/src/gridtk/cli.py +++ b/src/gridtk/cli.py @@ -636,6 +636,76 @@ def delete( session.commit() +@cli.command() +@job_filters +@click.option( + "--interval", + default=10, + type=click.INT, + help="Polling interval in seconds.", +) +@click.pass_context +def wait(ctx, job_ids, states, names, dependents, interval): + """Wait for jobs to finish. Exits with code 1 if any job failed.""" + import time + + from .manager import JobManager + + job_manager: JobManager = ctx.meta["job_manager"] + # Terminal states - jobs in these states won't change + terminal_states = { + "BOOT_FAIL", + "CANCELLED", + "COMPLETED", + "DEADLINE", + "FAILED", + "NODE_FAIL", + "OUT_OF_MEMORY", + "PREEMPTED", + "REVOKED", + "SPECIAL_EXIT", + "TIMEOUT", + } + # Failed states - if any job ends in these, exit code 1 + failed_states = { + "BOOT_FAIL", + "CANCELLED", + "DEADLINE", + "FAILED", + "NODE_FAIL", + "OUT_OF_MEMORY", + "PREEMPTED", + "REVOKED", + "SPECIAL_EXIT", + "TIMEOUT", + } + + while True: + with job_manager as session: + jobs = job_manager.list_jobs( + job_ids=job_ids, states=states, names=names, dependents=dependents + ) + if not jobs: + click.echo("No jobs found.") + return + + all_terminal = all(job.state in terminal_states for job in jobs) + if all_terminal: + any_failed = any(job.state in failed_states for job in jobs) + for job in jobs: + click.echo(f"Job {job.id}: {job.state} ({job.exit_code})") + if any_failed: + ctx.exit(1) + return + + # Show progress + pending = sum(1 for j in jobs if j.state not in terminal_states) + click.echo(f"Waiting for {pending} job(s)... (checking every {interval}s)") + session.commit() + + time.sleep(interval) + + @cli.command() @job_filters @click.option( diff --git a/tests/test_gridtk.py b/tests/test_gridtk.py index 76bf09e..e9921d2 100644 --- a/tests/test_gridtk.py +++ b/tests/test_gridtk.py @@ -732,6 +732,42 @@ def test_submit_json(mock_check_output, runner): assert data["name"] == "gridtk" +@patch("subprocess.check_output") +def test_wait_command(mock_check_output, runner): + # Test wait with COMPLETED job (exit code 0) + with runner.isolated_filesystem(): + submit_job_id = 9876543 + _submit_job( + runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id + ) + mock_check_output.side_effect = _make_side_effect( + [ + json.dumps( + _jobs_sacct_dict([submit_job_id], "COMPLETED", "None", "node001") + ), + ] + ) + result = runner.invoke(cli, ["wait"]) + assert_click_runner_result(result) + assert "Job 1: COMPLETED" in result.output + + # Test wait with FAILED job (exit code 1) + with runner.isolated_filesystem(): + submit_job_id = 9876543 + mock_check_output.side_effect = None + _submit_job( + runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id + ) + mock_check_output.side_effect = _make_side_effect( + [ + _failed_job_sacct_json(submit_job_id), + ] + ) + result = runner.invoke(cli, ["wait"]) + assert_click_runner_result(result, exit_code=1) + assert "Job 1: FAILED" in result.output + + if __name__ == "__main__": import sys From 48784790f6245179409c6311ad28da7a64369707 Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Mon, 16 Mar 2026 17:22:16 +0100 Subject: [PATCH 02/12] fix: ensure wait output is printed before exit Replace ctx.exit(1) with raise SystemExit(1) to ensure click.echo output is flushed before the process exits. ctx.exit(1) can suppress output in some Python/Click versions (e.g. Python 3.9). Co-Authored-By: Claude Opus 4.6 (1M context) --- src/gridtk/cli.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/gridtk/cli.py b/src/gridtk/cli.py index 8b0bea3..cad2dbe 100644 --- a/src/gridtk/cli.py +++ b/src/gridtk/cli.py @@ -695,7 +695,7 @@ def wait(ctx, job_ids, states, names, dependents, interval): for job in jobs: click.echo(f"Job {job.id}: {job.state} ({job.exit_code})") if any_failed: - ctx.exit(1) + raise SystemExit(1) return # Show progress From 40556e3fa7439f2cadccdb7db6888a45831ae35c Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Mon, 16 Mar 2026 17:33:27 +0100 Subject: [PATCH 03/12] fix: commit session before SystemExit in wait command Ensures the DB session is properly committed before raising SystemExit to avoid leaving the database in a bad state. Co-Authored-By: Claude Opus 4.6 (1M context) --- src/gridtk/cli.py | 1 + 1 file changed, 1 insertion(+) diff --git a/src/gridtk/cli.py b/src/gridtk/cli.py index cad2dbe..cdc953e 100644 --- a/src/gridtk/cli.py +++ b/src/gridtk/cli.py @@ -694,6 +694,7 @@ def wait(ctx, job_ids, states, names, dependents, interval): any_failed = any(job.state in failed_states for job in jobs) for job in jobs: click.echo(f"Job {job.id}: {job.state} ({job.exit_code})") + session.commit() if any_failed: raise SystemExit(1) return From 5e045996c9960d612527d4748666131c785311f1 Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Mon, 16 Mar 2026 17:48:09 +0100 Subject: [PATCH 04/12] docs: document wait command in README Co-Authored-By: Claude Opus 4.6 (1M context) --- README.md | 21 +++++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/README.md b/README.md index fc778a0..f5ed8c3 100644 --- a/README.md +++ b/README.md @@ -260,6 +260,27 @@ Here are some useful commands: sacctmgr -n -p list assoc where user=$USER | awk '-F|' '{print " "$2}' ``` +### Waiting for Jobs + +Use `gridtk wait` to block until all jobs finish: +```bash +$ gridtk wait +Waiting for 3 job(s)... (checking every 10s) +Job 1: COMPLETED (0) +Job 2: COMPLETED (0) +Job 3: FAILED (1) +``` + +`gridtk wait` exits with code 1 if any job failed, making it easy to chain: +```bash +$ gridtk submit job.sh && gridtk wait && echo "All done!" +``` + +You can filter which jobs to wait for and change the polling interval: +```bash +$ gridtk wait -j 1,2 --interval 30 +``` + ### Tab Completion GridTK supports tab completion for the `gridtk` command. To enable it, add the following From 44f8e7156b6d3a5d0187f35cc84c48c4fec99a0c Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Tue, 17 Mar 2026 09:44:41 +0100 Subject: [PATCH 05/12] fix: show state breakdown in wait progress (pending vs running) Co-Authored-By: Claude Opus 4.6 (1M context) --- src/gridtk/cli.py | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) diff --git a/src/gridtk/cli.py b/src/gridtk/cli.py index cdc953e..f8c2074 100644 --- a/src/gridtk/cli.py +++ b/src/gridtk/cli.py @@ -699,9 +699,18 @@ def wait(ctx, job_ids, states, names, dependents, interval): raise SystemExit(1) return - # Show progress - pending = sum(1 for j in jobs if j.state not in terminal_states) - click.echo(f"Waiting for {pending} job(s)... (checking every {interval}s)") + # Show progress with state breakdown + from collections import Counter + + active = [j for j in jobs if j.state not in terminal_states] + counts = Counter(j.state for j in active) + breakdown = ", ".join( + f"{n} {s.lower()}" for s, n in sorted(counts.items()) + ) + click.echo( + f"Waiting for {len(active)} job(s): {breakdown}" + f" (checking every {interval}s)" + ) session.commit() time.sleep(interval) From 25dca2b08477750a8ae66dba7d9fa95e67ab85b0 Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Tue, 17 Mar 2026 13:10:38 +0100 Subject: [PATCH 06/12] fix: robust wait test and formatting Use return_value instead of side_effect iterator to avoid mock exhaustion from GC-triggered __del__ calls on some Python versions. Co-Authored-By: Claude Opus 4.6 (1M context) --- src/gridtk/cli.py | 4 +--- tests/test_gridtk.py | 15 ++++----------- 2 files changed, 5 insertions(+), 14 deletions(-) diff --git a/src/gridtk/cli.py b/src/gridtk/cli.py index f8c2074..f8aab02 100644 --- a/src/gridtk/cli.py +++ b/src/gridtk/cli.py @@ -704,9 +704,7 @@ def wait(ctx, job_ids, states, names, dependents, interval): active = [j for j in jobs if j.state not in terminal_states] counts = Counter(j.state for j in active) - breakdown = ", ".join( - f"{n} {s.lower()}" for s, n in sorted(counts.items()) - ) + breakdown = ", ".join(f"{n} {s.lower()}" for s, n in sorted(counts.items())) click.echo( f"Waiting for {len(active)} job(s): {breakdown}" f" (checking every {interval}s)" diff --git a/tests/test_gridtk.py b/tests/test_gridtk.py index e9921d2..a9d9411 100644 --- a/tests/test_gridtk.py +++ b/tests/test_gridtk.py @@ -740,12 +740,9 @@ def test_wait_command(mock_check_output, runner): _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) - mock_check_output.side_effect = _make_side_effect( - [ - json.dumps( - _jobs_sacct_dict([submit_job_id], "COMPLETED", "None", "node001") - ), - ] + # Use return_value so all calls (squeue, sacct, __del__) get valid data + mock_check_output.return_value = json.dumps( + _jobs_sacct_dict([submit_job_id], "COMPLETED", "None", "node001") ) result = runner.invoke(cli, ["wait"]) assert_click_runner_result(result) @@ -758,11 +755,7 @@ def test_wait_command(mock_check_output, runner): _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) - mock_check_output.side_effect = _make_side_effect( - [ - _failed_job_sacct_json(submit_job_id), - ] - ) + mock_check_output.return_value = _failed_job_sacct_json(submit_job_id) result = runner.invoke(cli, ["wait"]) assert_click_runner_result(result, exit_code=1) assert "Job 1: FAILED" in result.output From 5c41ce4d918e477a74923b4273a22ad66e9a09bc Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Tue, 17 Mar 2026 13:13:27 +0100 Subject: [PATCH 07/12] fix: force GC in wait test to prevent __del__ race condition Co-Authored-By: Claude Opus 4.6 (1M context) --- tests/test_gridtk.py | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/tests/test_gridtk.py b/tests/test_gridtk.py index a9d9411..c332bed 100644 --- a/tests/test_gridtk.py +++ b/tests/test_gridtk.py @@ -734,13 +734,16 @@ def test_submit_json(mock_check_output, runner): @patch("subprocess.check_output") def test_wait_command(mock_check_output, runner): + import gc + # Test wait with COMPLETED job (exit code 0) with runner.isolated_filesystem(): submit_job_id = 9876543 _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) - # Use return_value so all calls (squeue, sacct, __del__) get valid data + # Force GC to run __del__ on submit's JobManager before changing mock + gc.collect() mock_check_output.return_value = json.dumps( _jobs_sacct_dict([submit_job_id], "COMPLETED", "None", "node001") ) @@ -755,6 +758,7 @@ def test_wait_command(mock_check_output, runner): _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) + gc.collect() mock_check_output.return_value = _failed_job_sacct_json(submit_job_id) result = runner.invoke(cli, ["wait"]) assert_click_runner_result(result, exit_code=1) From 1446a95d5110f9863aa18a08ac9b2f15b58f0363 Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Tue, 17 Mar 2026 13:14:38 +0100 Subject: [PATCH 08/12] debug: add list assertion before wait --- tests/test_gridtk.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/tests/test_gridtk.py b/tests/test_gridtk.py index c332bed..40091fa 100644 --- a/tests/test_gridtk.py +++ b/tests/test_gridtk.py @@ -747,6 +747,9 @@ def test_wait_command(mock_check_output, runner): mock_check_output.return_value = json.dumps( _jobs_sacct_dict([submit_job_id], "COMPLETED", "None", "node001") ) + # Debug: verify job exists in DB before wait + result = runner.invoke(cli, ["list"]) + assert "gridtk" in result.output, f"list before wait: {result.output!r}" result = runner.invoke(cli, ["wait"]) assert_click_runner_result(result) assert "Job 1: COMPLETED" in result.output From 6952b826236cb8130acc37524a349f76b55aa195 Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Tue, 17 Mar 2026 13:16:13 +0100 Subject: [PATCH 09/12] debug: assert DB exists after submit --- tests/test_gridtk.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/tests/test_gridtk.py b/tests/test_gridtk.py index 40091fa..5dcd3e6 100644 --- a/tests/test_gridtk.py +++ b/tests/test_gridtk.py @@ -742,14 +742,12 @@ def test_wait_command(mock_check_output, runner): _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) - # Force GC to run __del__ on submit's JobManager before changing mock + # Force GC and ensure DB survives __del__ cleanup gc.collect() + assert Path("jobs.sql3").exists(), "DB was deleted by __del__" mock_check_output.return_value = json.dumps( _jobs_sacct_dict([submit_job_id], "COMPLETED", "None", "node001") ) - # Debug: verify job exists in DB before wait - result = runner.invoke(cli, ["list"]) - assert "gridtk" in result.output, f"list before wait: {result.output!r}" result = runner.invoke(cli, ["wait"]) assert_click_runner_result(result) assert "Job 1: COMPLETED" in result.output @@ -762,6 +760,7 @@ def test_wait_command(mock_check_output, runner): runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) gc.collect() + assert Path("jobs.sql3").exists(), "DB was deleted by __del__" mock_check_output.return_value = _failed_job_sacct_json(submit_job_id) result = runner.invoke(cli, ["wait"]) assert_click_runner_result(result, exit_code=1) From d8adb6bb407b034c8798eb17bad1c86c0ddf6fe3 Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Tue, 17 Mar 2026 13:18:41 +0100 Subject: [PATCH 10/12] fix: add gc.collect() to all tests that submit then query Prevents __del__ race condition on CI where GC runs between submit and subsequent commands, deleting the DB. Co-Authored-By: Claude Opus 4.6 (1M context) --- tests/test_gridtk.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/tests/test_gridtk.py b/tests/test_gridtk.py index 5dcd3e6..a827b0e 100644 --- a/tests/test_gridtk.py +++ b/tests/test_gridtk.py @@ -696,12 +696,14 @@ def test_submit_with_dependencies(mock_check_output, runner): @patch("subprocess.check_output") def test_list_json(mock_check_output, runner): + import gc + with runner.isolated_filesystem(): submit_job_id = 9876543 _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) - + gc.collect() mock_check_output.return_value = _pending_job_sacct_json(submit_job_id) result = runner.invoke(cli, ["list", "--json"]) assert_click_runner_result(result) From f4a5d8c25516f1323ea763d16663b8716b8c7340 Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Tue, 17 Mar 2026 13:20:55 +0100 Subject: [PATCH 11/12] debug: trace __del__ and DB existence on CI --- tests/test_gridtk.py | 22 +++++++++++++++------- 1 file changed, 15 insertions(+), 7 deletions(-) diff --git a/tests/test_gridtk.py b/tests/test_gridtk.py index a827b0e..4bfdf87 100644 --- a/tests/test_gridtk.py +++ b/tests/test_gridtk.py @@ -8,6 +8,8 @@ import subprocess import traceback +import gridtk.manager + from pathlib import Path from unittest.mock import Mock, patch @@ -697,13 +699,25 @@ def test_submit_with_dependencies(mock_check_output, runner): @patch("subprocess.check_output") def test_list_json(mock_check_output, runner): import gc + import os + + original_del = gridtk.manager.JobManager.__del__ + + def debug_del(self): + print(f"DEBUG __del__: DB={self.database}, exists={Path(self.database).exists()}") + original_del(self) + print(f"DEBUG __del__ after: exists={Path(self.database).exists()}") with runner.isolated_filesystem(): submit_job_id = 9876543 _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) + print(f"DEBUG after submit: DB exists={os.path.exists('jobs.sql3')}") + gridtk.manager.JobManager.__del__ = debug_del gc.collect() + print(f"DEBUG after gc: DB exists={os.path.exists('jobs.sql3')}") + gridtk.manager.JobManager.__del__ = original_del mock_check_output.return_value = _pending_job_sacct_json(submit_job_id) result = runner.invoke(cli, ["list", "--json"]) assert_click_runner_result(result) @@ -734,19 +748,15 @@ def test_submit_json(mock_check_output, runner): assert data["name"] == "gridtk" +@patch("gridtk.manager.JobManager.__del__", lambda self: None) @patch("subprocess.check_output") def test_wait_command(mock_check_output, runner): - import gc - # Test wait with COMPLETED job (exit code 0) with runner.isolated_filesystem(): submit_job_id = 9876543 _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) - # Force GC and ensure DB survives __del__ cleanup - gc.collect() - assert Path("jobs.sql3").exists(), "DB was deleted by __del__" mock_check_output.return_value = json.dumps( _jobs_sacct_dict([submit_job_id], "COMPLETED", "None", "node001") ) @@ -761,8 +771,6 @@ def test_wait_command(mock_check_output, runner): _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) - gc.collect() - assert Path("jobs.sql3").exists(), "DB was deleted by __del__" mock_check_output.return_value = _failed_job_sacct_json(submit_job_id) result = runner.invoke(cli, ["wait"]) assert_click_runner_result(result, exit_code=1) From a475127707e6162c959acb525a39227c27eed645 Mon Sep 17 00:00:00 2001 From: Amir Mohammadi Date: Tue, 17 Mar 2026 13:24:27 +0100 Subject: [PATCH 12/12] fix: move DB cleanup from __del__ to explicit cleanup_empty_database() MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit __del__ was deleting the DB on GC due to SQLite session isolation — a new session in __del__ couldn't see rows committed by another engine instance. Now __del__ only disposes the engine, and DB cleanup is called explicitly after delete commands. Co-Authored-By: Claude Opus 4.6 (1M context) --- src/gridtk/cli.py | 5 +++-- src/gridtk/manager.py | 6 ++++-- tests/test_gridtk.py | 18 ------------------ 3 files changed, 7 insertions(+), 22 deletions(-) diff --git a/src/gridtk/cli.py b/src/gridtk/cli.py index f8aab02..dbefd0c 100644 --- a/src/gridtk/cli.py +++ b/src/gridtk/cli.py @@ -186,9 +186,10 @@ def cli(ctx, database, logs_dir): @cli.result_callback() def process_result(result, **kwargs): - """Delete the job manager from the context.""" + """Clean up empty databases and dispose the job manager.""" ctx = click.get_current_context() - del ctx.meta["job_manager"] + job_manager = ctx.meta.pop("job_manager") + job_manager.cleanup_empty_database() @cli.command( diff --git a/src/gridtk/manager.py b/src/gridtk/manager.py index dee41af..85e193f 100644 --- a/src/gridtk/manager.py +++ b/src/gridtk/manager.py @@ -302,8 +302,8 @@ def resubmit_jobs(self, **kwargs): self.session.add(job) return jobs - def __del__(self): - # if there are no jobs in the database, delete the database file and the logs directory (if empty) + def cleanup_empty_database(self): + """Delete the database file and logs directory if no jobs remain.""" with self: if ( Path(self.database).exists() @@ -312,4 +312,6 @@ def __del__(self): Path(self.database).unlink() if self.logs_dir.exists() and len(os.listdir(self.logs_dir)) == 0: shutil.rmtree(self.logs_dir) + + def __del__(self): self.engine.dispose() diff --git a/tests/test_gridtk.py b/tests/test_gridtk.py index 4bfdf87..2f1e7ec 100644 --- a/tests/test_gridtk.py +++ b/tests/test_gridtk.py @@ -8,8 +8,6 @@ import subprocess import traceback -import gridtk.manager - from pathlib import Path from unittest.mock import Mock, patch @@ -698,26 +696,11 @@ def test_submit_with_dependencies(mock_check_output, runner): @patch("subprocess.check_output") def test_list_json(mock_check_output, runner): - import gc - import os - - original_del = gridtk.manager.JobManager.__del__ - - def debug_del(self): - print(f"DEBUG __del__: DB={self.database}, exists={Path(self.database).exists()}") - original_del(self) - print(f"DEBUG __del__ after: exists={Path(self.database).exists()}") - with runner.isolated_filesystem(): submit_job_id = 9876543 _submit_job( runner=runner, mock_check_output=mock_check_output, job_id=submit_job_id ) - print(f"DEBUG after submit: DB exists={os.path.exists('jobs.sql3')}") - gridtk.manager.JobManager.__del__ = debug_del - gc.collect() - print(f"DEBUG after gc: DB exists={os.path.exists('jobs.sql3')}") - gridtk.manager.JobManager.__del__ = original_del mock_check_output.return_value = _pending_job_sacct_json(submit_job_id) result = runner.invoke(cli, ["list", "--json"]) assert_click_runner_result(result) @@ -748,7 +731,6 @@ def test_submit_json(mock_check_output, runner): assert data["name"] == "gridtk" -@patch("gridtk.manager.JobManager.__del__", lambda self: None) @patch("subprocess.check_output") def test_wait_command(mock_check_output, runner): # Test wait with COMPLETED job (exit code 0)