"""Unit tests for navi.core.tasks — TaskManager / TaskJob / BoundedEventQueue."""
import asyncio
import time
import pytest
from navi.core import tasks as tasks_mod
from navi.core.events import TaskUpdate
from navi.core.tasks import (
BoundedEventQueue,
TaskJob,
TaskManager,
args_summary,
get_task_manager,
)
from navi.tools._internal.base import ToolContext, ToolResult, current_stop_event
def patch_settings(monkeypatch, **overrides):
"""Replace the frozen settings singleton with an overridden copy."""
import navi.config as config_mod
from navi.config import Settings
new_settings = Settings(**overrides)
monkeypatch.setattr(config_mod, "settings", new_settings)
return new_settings
@pytest.fixture(autouse=True)
def fresh_manager(monkeypatch):
"""Isolate the module singleton; clean up stray tasks afterwards."""
manager = TaskManager()
monkeypatch.setattr(tasks_mod, "_manager", manager)
yield manager
for job in list(manager._jobs.values()):
if job.task is not None and not job.task.done():
job.task.cancel()
if manager._sweeper_task is not None and not manager._sweeper_task.done():
manager._sweeper_task.cancel()
def make_ctx(**overrides):
defaults = dict(
session_id="s1",
event_sink=None,
stop_event=None,
model="m",
user_id="u1",
user_role="admin",
user_info=None,
cwd="/tmp",
)
defaults.update(overrides)
return ToolContext(**defaults)
class TestLifecycle:
async def test_completed_job(self, fresh_manager):
async def factory(bg_ctx):
return ToolResult(success=True, output="done 42")
job = fresh_manager.submit("s1", "code_exec", {"code": "x"}, factory, make_ctx())
assert job.status == "running"
await job.done.wait()
assert job.status == "completed"
assert job.result.output == "done 42"
assert job.finished_at is not None
async def test_failed_job(self, fresh_manager):
async def factory(bg_ctx):
raise ValueError("boom")
job = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
await job.done.wait()
assert job.status == "failed"
assert "boom" in job.error
assert job.result.success is False
async def test_cancel_running(self, fresh_manager):
started = asyncio.Event()
async def factory(bg_ctx):
started.set()
await asyncio.sleep(30)
job = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
await started.wait()
assert fresh_manager.cancel(job) is True
await job.done.wait()
assert job.status == "cancelled"
assert fresh_manager.cancel(job) is False # already finished
async def test_task_stop_event_is_isolated(self, fresh_manager):
"""The job sees its own stop_event, not any run-level one."""
seen = {}
async def factory(bg_ctx):
seen["ctx_stop"] = bg_ctx.stop_event
seen["var_stop"] = current_stop_event.get()
return ToolResult(success=True, output="ok")
job = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
await job.done.wait()
assert seen["ctx_stop"] is job.stop_event
assert seen["var_stop"] is job.stop_event
async def test_get_is_session_scoped(self, fresh_manager):
async def factory(bg_ctx):
return ToolResult(success=True, output="ok")
job = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
await job.done.wait()
assert fresh_manager.get(job.task_id, "s1") is job
assert fresh_manager.get(job.task_id, "s2") is None
assert [j.task_id for j in fresh_manager.list("s1")] == [job.task_id]
assert fresh_manager.list("s2") == []
class TestCaps:
async def test_per_session_cap(self, fresh_manager, monkeypatch):
patch_settings(monkeypatch, tasks_max_per_session=1)
async def factory(bg_ctx):
await asyncio.sleep(30)
first = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
assert not isinstance(first, str)
second = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
assert isinstance(second, str) and "limit" in second
fresh_manager.cancel(first)
async def test_rate_limit(self, fresh_manager, monkeypatch):
patch_settings(monkeypatch, tasks_rate_limit=1)
async def factory(bg_ctx):
return ToolResult(success=True, output="ok")
first = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
assert not isinstance(first, str)
second = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
assert isinstance(second, str) and "rate limit" in second
async def test_spawn_cap_only_limits_spawn_agent(self, fresh_manager, monkeypatch):
patch_settings(monkeypatch, tasks_max_spawn=1)
async def factory(bg_ctx):
await asyncio.sleep(30)
first = fresh_manager.submit("s1", "spawn_agent", {}, factory, make_ctx())
assert not isinstance(first, str)
# another spawn_agent is rejected…
second = fresh_manager.submit("s1", "spawn_agent", {}, factory, make_ctx())
assert isinstance(second, str) and "subagent" in second
# …but a regular tool is not
third = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
assert not isinstance(third, str)
fresh_manager.cancel(first)
fresh_manager.cancel(third)
class TestUpdateCallback:
async def test_published_on_submit_and_finish(self, fresh_manager):
updates = []
fresh_manager.set_update_callback(updates.append)
async def factory(bg_ctx):
return ToolResult(success=True, output="ok")
job = fresh_manager.submit("s1", "terminal", {"command": "ls"}, factory, make_ctx(),
parent_tool_call_id="tc1")
await job.done.wait()
await asyncio.sleep(0) # let the callback fire settle
statuses = [u.status for u in updates]
assert statuses[0] == "running"
assert statuses[-1] == "completed"
final = updates[-1]
assert isinstance(final, TaskUpdate)
assert final.task_id == job.task_id
assert final.tool == "terminal"
assert final.parent_tool_call_id == "tc1"
assert final.to_wire()["type"] == "task_update"
async def test_broken_callback_does_not_kill_job(self, fresh_manager):
def bad_callback(update):
raise RuntimeError("nope")
fresh_manager.set_update_callback(bad_callback)
async def factory(bg_ctx):
return ToolResult(success=True, output="ok")
job = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
await job.done.wait()
assert job.status == "completed"
class TestBoundedEventQueue:
async def test_drop_oldest_when_full(self):
q = BoundedEventQueue(maxsize=2)
await q.put("a")
q.put_nowait("b")
await q.put("c")
assert list(q._queue) == ["b", "c"]
async def test_put_never_blocks(self):
q = BoundedEventQueue(maxsize=1)
await asyncio.wait_for(q.put("a"), timeout=0.1)
await asyncio.wait_for(q.put("b"), timeout=0.1)
assert q.qsize() == 1
class TestReap:
async def test_reap_drops_expired_finished_jobs(self, fresh_manager, monkeypatch):
patch_settings(monkeypatch, tasks_ttl_sec=0)
async def factory(bg_ctx):
return ToolResult(success=True, output="ok")
job = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
await job.done.wait()
assert fresh_manager.reap() == 1
assert fresh_manager.get(job.task_id, "s1") is None
async def test_reap_keeps_running_jobs(self, fresh_manager, monkeypatch):
patch_settings(monkeypatch, tasks_ttl_sec=0)
async def factory(bg_ctx):
await asyncio.sleep(30)
job = fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
assert fresh_manager.reap() == 0
assert fresh_manager.get(job.task_id, "s1") is job
fresh_manager.cancel(job)
class TestHelpers:
def test_get_task_manager_singleton(self, monkeypatch):
monkeypatch.setattr(tasks_mod, "_manager", None)
assert get_task_manager() is get_task_manager()
def test_args_summary(self):
assert args_summary({"a": 1}) == '{"a": 1}'
assert args_summary({"a": "x" * 500}, limit=10) == '{"a": "xx' + "x"
async def test_task_job_defaults(self):
job = TaskJob(task_id="bt-x", session_id="s1", tool="terminal", args={})
assert job.status == "running"
assert job.subagent_tokens is None
assert job.preview() == "" # running → empty preview
async def test_time_in_submit_window(self, fresh_manager):
async def factory(bg_ctx):
return ToolResult(success=True, output="ok")
fresh_manager.submit("s1", "code_exec", {}, factory, make_ctx())
assert len(fresh_manager._submit_times["s1"]) == 1
assert time.time() - fresh_manager._submit_times["s1"][0] < 5