diff --git a/src/pinky_daemon/scheduler.py b/src/pinky_daemon/scheduler.py index a053541d..8aec3ca0 100644 --- a/src/pinky_daemon/scheduler.py +++ b/src/pinky_daemon/scheduler.py @@ -905,16 +905,26 @@ async def _replay_pending_locked(self, agent_name: str) -> None: elif not current_schedule.enabled and not current_schedule.one_shot: zombie_reason = "schedule disabled" if zombie_reason: + if pending.parked_at != 0: + continue quarantined = self._registry.park_pending_schedule_wake( pending.id, reason=f"terminal replay policy: {zombie_reason}", ) - _log( - f"scheduler: PERSISTED_WAKE_ZOMBIE_QUARANTINED pending " - f"#{pending.id}, schedule #{pending.schedule_id} for " - f"agent '{pending.agent_name}': {zombie_reason}; " - f"quarantined={quarantined}" - ) + if quarantined: + _log( + f"scheduler: PERSISTED_WAKE_ZOMBIE_QUARANTINED pending " + f"#{pending.id}, schedule #{pending.schedule_id} for " + f"agent '{pending.agent_name}': {zombie_reason}; " + "quarantined=True" + ) + else: + _log( + f"scheduler: PERSISTED_WAKE_ZOMBIE_PARK_NOOP pending " + f"#{pending.id}, schedule #{pending.schedule_id} for " + f"agent '{pending.agent_name}': {zombie_reason}; " + "park returned no state change" + ) continue if pending.parked_at == 0: pending_wakes.append(pending) diff --git a/tests/test_scheduler.py b/tests/test_scheduler.py index 6a6e133b..e70a3e65 100644 --- a/tests/test_scheduler.py +++ b/tests/test_scheduler.py @@ -1885,6 +1885,98 @@ async def unconfirmed(agent_name, session_id, prompt): assert "PERSISTED_WAKE_ZOMBIE_QUARANTINED" in error_log assert f"schedule #{zombie.id}" in error_log + @pytest.mark.asyncio + async def test_replay_quarantines_live_zombie_once_then_skips_terminal_row( + self, registry, monkeypatch, capsys + ): + registry.register("oleg") + schedule = registry.add_schedule( + "oleg", "* * * * *", name="zombie", prompt="never deliver" + ) + pending, _ = registry.persist_schedule_wake( + schedule.id, + agent_name="oleg", + schedule_name="zombie", + prompt="never deliver", + fired_at=100.0, + ) + registry.remove_schedule(schedule.id) + park_calls: list[int] = [] + original_park = registry.park_pending_schedule_wake + + def tracked_park(pending_id, **kwargs): + park_calls.append(pending_id) + return original_park(pending_id, **kwargs) + + monkeypatch.setattr( + registry, "park_pending_schedule_wake", tracked_park + ) + attempts: list[str] = [] + + async def confirmed(agent_name, session_id, prompt): + del agent_name, session_id + attempts.append(prompt) + return True + + first_boot = AgentScheduler(registry, wake_callback=confirmed) + await first_boot._replay_pending_locked("oleg") + first_terminal_state = registry.get_schedule_wake_by_fire( + schedule.id, 100.0 + ).to_dict() + + for _ in range(5): + later_boot = AgentScheduler(registry, wake_callback=confirmed) + await later_boot._replay_pending_locked("oleg") + + assert attempts == [] + assert park_calls == [pending.id] + assert registry.get_schedule_wake_by_fire( + schedule.id, 100.0 + ).to_dict() == first_terminal_state + assert first_terminal_state["state"] == "quarantined" + assert first_terminal_state["last_error"].endswith( + "schedule deleted" + ) + assert ( + capsys.readouterr().err.count( + "PERSISTED_WAKE_ZOMBIE_QUARANTINED" + ) + == 1 + ) + + @pytest.mark.asyncio + async def test_replay_logs_live_zombie_park_noop( + self, registry, monkeypatch, capsys + ): + registry.register("oleg") + schedule = registry.add_schedule( + "oleg", "* * * * *", name="zombie", prompt="never deliver" + ) + pending, _ = registry.persist_schedule_wake( + schedule.id, + agent_name="oleg", + schedule_name="zombie", + prompt="never deliver", + fired_at=100.0, + ) + registry.remove_schedule(schedule.id) + monkeypatch.setattr( + registry, + "park_pending_schedule_wake", + lambda pending_id, **kwargs: False, + ) + + scheduler = AgentScheduler(registry) + await scheduler._replay_pending_locked("oleg") + + stored = registry.get_schedule_wake_by_fire(schedule.id, 100.0) + assert stored is not None + assert stored.to_dict() == pending.to_dict() + error_log = capsys.readouterr().err + assert "PERSISTED_WAKE_ZOMBIE_PARK_NOOP" in error_log + assert f"pending #{pending.id}" in error_log + assert "PERSISTED_WAKE_ZOMBIE_QUARANTINED" not in error_log + @pytest.mark.asyncio async def test_terminal_zombie_head_retires_and_fifo_advances( self, registry, capsys