Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion ms_agent/cron/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
36 changes: 36 additions & 0 deletions tests/cron/test_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
37 changes: 37 additions & 0 deletions tests/cron/test_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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)
Expand Down