"""Unit tests for the tasks tool (list/check/wait/cancel of background jobs)."""
import asyncio
import pytest
from navi.core import tasks as tasks_mod
from navi.core.tasks import TaskManager
from navi.tools._internal.base import ToolContext, ToolResult
from navi.tools.tasks import TasksTool
@pytest.fixture(autouse=True)
def fresh_manager(monkeypatch):
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(session_id="s1"):
return ToolContext(
session_id=session_id, event_sink=None, stop_event=None, model=None,
user_id=None, user_role="admin", user_info=None, cwd=None,
)
async def submit_ok(manager, session_id="s1", output="done", delay=0.0):
async def factory(bg_ctx):
if delay:
await asyncio.sleep(delay)
return ToolResult(success=True, output=output)
return manager.submit(session_id, "code_exec", {"code": "x"}, factory,
make_ctx(session_id))
class TestList:
async def test_empty(self):
result = await TasksTool().execute({"action": "list"}, ctx=make_ctx())
assert result.success
assert "No background tasks" in result.output
async def test_lists_session_jobs_only(self, fresh_manager):
await submit_ok(fresh_manager, "s1")
await submit_ok(fresh_manager, "s2")
result = await TasksTool().execute({"action": "list"}, ctx=make_ctx("s1"))
assert result.success
assert "code_exec" in result.output
async def test_list_shows_finished_preview(self, fresh_manager):
job = await submit_ok(fresh_manager, "s1", output="42 files")
await job.done.wait()
result = await TasksTool().execute({"action": "list"}, ctx=make_ctx())
assert "completed" in result.output
assert "42 files" in result.output
class TestCheck:
async def test_check_finished_shows_result(self, fresh_manager):
job = await submit_ok(fresh_manager, output="the answer")
await job.done.wait()
result = await TasksTool().execute(
{"action": "check", "task_id": job.task_id}, ctx=make_ctx())
assert result.success
assert "completed" in result.output
assert "the answer" in result.output
async def test_check_running_shows_progress_placeholder(self, fresh_manager):
job = await submit_ok(fresh_manager, delay=5)
result = await TasksTool().execute(
{"action": "check", "task_id": job.task_id}, ctx=make_ctx())
assert "running" in result.output
fresh_manager.cancel(job)
async def test_check_unknown_task(self):
result = await TasksTool().execute(
{"action": "check", "task_id": "bt-nope"}, ctx=make_ctx())
assert not result.success
assert result.error == "task_not_found"
async def test_task_id_required(self):
result = await TasksTool().execute({"action": "check"}, ctx=make_ctx())
assert not result.success
async def test_check_is_session_scoped(self, fresh_manager):
job = await submit_ok(fresh_manager, session_id="other")
result = await TasksTool().execute(
{"action": "check", "task_id": job.task_id}, ctx=make_ctx("s1"))
assert result.error == "task_not_found"
class TestWait:
async def test_wait_returns_result(self, fresh_manager):
job = await submit_ok(fresh_manager, output="late result", delay=0.05)
result = await TasksTool().execute(
{"action": "wait", "task_id": job.task_id}, ctx=make_ctx())
assert result.success
assert "late result" in result.output
async def test_wait_timeout(self, fresh_manager):
job = await submit_ok(fresh_manager, delay=30)
result = await TasksTool().execute(
{"action": "wait", "task_id": job.task_id, "timeout": 0.05},
ctx=make_ctx())
assert not result.success
assert result.error == "wait_timeout"
assert "still running" in result.output
fresh_manager.cancel(job)
async def test_wait_timeout_capped_at_120(self, fresh_manager):
job = await submit_ok(fresh_manager, delay=0.01)
# timeout above the cap must not blow up; job finishes fast anyway
result = await TasksTool().execute(
{"action": "wait", "task_id": job.task_id, "timeout": 9999},
ctx=make_ctx())
assert result.success
class TestCancel:
async def test_cancel_running(self, fresh_manager):
job = await submit_ok(fresh_manager, delay=30)
result = await TasksTool().execute(
{"action": "cancel", "task_id": job.task_id}, ctx=make_ctx())
assert result.success
await job.done.wait()
assert job.status == "cancelled"
async def test_cancel_finished(self, fresh_manager):
job = await submit_ok(fresh_manager)
await job.done.wait()
result = await TasksTool().execute(
{"action": "cancel", "task_id": job.task_id}, ctx=make_ctx())
assert not result.success
assert result.error == "not_running"
async def test_cancel_confirms_status(self, fresh_manager):
"""A cancel on a live job waits for the terminal state (B9: honest answer)."""
job = await submit_ok(fresh_manager, delay=0.05)
result = await TasksTool().execute(
{"action": "cancel", "task_id": job.task_id}, ctx=make_ctx())
assert result.success
assert "cancelled" in result.output
await job.done.wait()
assert job.status == "cancelled"
async def test_cancel_raced_completion_shows_result(self, fresh_manager, monkeypatch):
"""Completed-right-under-the-cancel reports the real outcome (B9)."""
import time
job = await submit_ok(fresh_manager, delay=30, output="raced done")
# simulate the race: the job completed normally before the
# cancellation could land — terminal state published, so the tool
# must show the real result, not "cancelled".
def raced_cancel(target):
job.result = ToolResult(success=True, output="raced done")
job.status = "completed"
job.finished_at = time.time()
job.done.set()
return True
monkeypatch.setattr(fresh_manager, "cancel", raced_cancel)
result = await TasksTool().execute(
{"action": "cancel", "task_id": job.task_id}, ctx=make_ctx())
assert not result.success
assert result.error == "not_running"
assert "raced done" in result.output
fresh_manager.cancel(job) # cleanup the sleeper task
async def test_cancel_of_done_task_shows_result(self, fresh_manager):
"""Task already done but status not yet flipped: no false 'cancelled'."""
from navi.core.tasks import TaskJob
async def done_coro():
return ToolResult(success=True, output="real result")
job = TaskJob(task_id="bt-deadbeef", session_id="s1", tool="code_exec", args={})
job.status = "completed"
job.result = ToolResult(success=True, output="real result")
job.task = asyncio.create_task(done_coro())
await asyncio.sleep(0) # let it finish
job.done.set()
fresh_manager._jobs[job.task_id] = job
result = await TasksTool().execute(
{"action": "cancel", "task_id": job.task_id}, ctx=make_ctx())
assert not result.success
assert result.error == "not_running"
assert "real result" in result.output
class TestFallbacks:
async def test_session_falls_back_to_contextvar(self, fresh_manager, monkeypatch):
from navi.tools._internal.base import current_session_id
job = await submit_ok(fresh_manager, "ctx-session")
token = current_session_id.set("ctx-session")
try:
await asyncio.sleep(0) # let the background job run
result = await TasksTool().execute(
{"action": "check", "task_id": job.task_id}, ctx=None)
finally:
current_session_id.reset(token)
assert result.success
assert "completed" in result.output