Skip to content
Merged
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
22 changes: 16 additions & 6 deletions src/pinky_daemon/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
92 changes: 92 additions & 0 deletions tests/test_scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading