From c2738aee8ffac50e1d7e0482d29814beb6ed2e3d Mon Sep 17 00:00:00 2001 From: Lewis <233926000+lewismosciski@users.noreply.github.com> Date: Tue, 22 Sep 2026 11:39:37 +0800 Subject: [PATCH] fix(cron): preserve pause when a running job finishes --- ms_agent/cron/manager.py | 5 ++++- tests/cron/test_manager.py | 36 ++++++++++++++++++++++++++++++++++++ tests/cron/test_service.py | 37 +++++++++++++++++++++++++++++++++++++ 3 files changed, 77 insertions(+), 1 deletion(-) diff --git a/ms_agent/cron/manager.py b/ms_agent/cron/manager.py index c44886528..58c1f18f9 100644 --- a/ms_agent/cron/manager.py +++ b/ms_agent/cron/manager.py @@ -156,7 +156,10 @@ def record_result(self, job: CronJobSpec, result: ExecutionResult) -> None: else: next_run = advance_next_run(job.schedule, state.next_run_at or '') state.next_run_at = next_run - state.status = 'scheduled' if next_run else 'completed' + if not next_run: + state.status = 'completed' + elif state.status != 'paused': + state.status = 'scheduled' self._repo.save_job_and_state(job, state) diff --git a/tests/cron/test_manager.py b/tests/cron/test_manager.py index ddcaa774f..c37f42b4b 100644 --- a/tests/cron/test_manager.py +++ b/tests/cron/test_manager.py @@ -87,6 +87,42 @@ def test_mark_running(self, manager): class TestJobManagerRecordResult: + @pytest.mark.parametrize('success', [True, False]) + def test_paused_running_job_stays_paused_after_result(self, manager, success): + job = manager.create_job(schedule_str='every 60s', prompt='pause during run') + job.repeat = RepeatSpec(times=3) + manager.repo.save_job(job) + manager.mark_running(job.id) + assert manager.pause_job(job.id) + old_next = manager.get_job(job.id)[1].next_run_at + + manager.record_result(job, ExecutionResult(success=success, duration_ms=10)) + + stored_job, state = manager.get_job(job.id) + assert state.status == 'paused' + assert state.next_run_at != old_next + assert state.run_count == 1 + assert state.last_status == ('ok' if success else 'error') + assert stored_job.repeat.completed == 1 + assert manager.resume_job(job.id) + assert manager.get_job(job.id)[1].status == 'scheduled' + + @pytest.mark.parametrize('schedule', ['every 60s', '2099-01-01T00:00:00']) + def test_paused_final_run_still_completes(self, manager, schedule): + job = manager.create_job(schedule_str=schedule, prompt='final run') + if job.schedule.kind == 'interval': + job.repeat = RepeatSpec(times=1) + manager.repo.save_job(job) + manager.mark_running(job.id) + assert manager.pause_job(job.id) + + manager.record_result(job, ExecutionResult(success=True)) + + state = manager.get_job(job.id)[1] + assert state.status == 'completed' + assert state.next_run_at is None + assert state.run_count == 1 + def test_record_success(self, manager): job = manager.create_job(schedule_str='every 60s', prompt='ok') result = ExecutionResult(success=True, output='done', duration_ms=500) diff --git a/tests/cron/test_service.py b/tests/cron/test_service.py index 69eae4894..1b189db42 100644 --- a/tests/cron/test_service.py +++ b/tests/cron/test_service.py @@ -6,6 +6,7 @@ import pytest +from ms_agent.cron.manager import JobManager from ms_agent.cron.service import CronService, PidManager from ms_agent.cron.types import CronJobSpec, ExecutionResult @@ -113,6 +114,42 @@ def test_get_output_empty(self, workspace): class TestCronServiceCallbacks: + @pytest.mark.asyncio + async def test_pause_from_another_manager_survives_completion(self, workspace): + service = CronService(workspace=workspace) + job = service.create_job(schedule_str='every 60s', prompt='pause while running') + service.trigger_job(job.id) + started = asyncio.Event() + release = asyncio.Event() + completed = asyncio.Event() + + async def execute(job, config): + started.set() + await release.wait() + return ExecutionResult(success=True, output='ok') + + async def on_complete(job, result): + completed.set() + + service.on_job_complete.append(on_complete) + with patch.object(service._executor, 'execute', side_effect=execute): + try: + assert await service.manual_tick() == 1 + await asyncio.wait_for(started.wait(), timeout=5) + # The CLI uses a separate manager to update the shared jobs file. + assert JobManager(workspace).pause_job(job.id) + release.set() + await asyncio.wait_for(completed.wait(), timeout=5) + + state = service.get_job(job.id)[1] + assert state.status == 'paused' + assert state.run_count == 1 + assert state.last_status == 'ok' + assert service.resume_job(job.id) + finally: + release.set() + await service.stop(force=True) + @pytest.mark.asyncio async def test_on_job_start_callback(self, workspace): service = CronService(workspace=workspace)