mirror of
https://github.com/abhi1693/openclaw-mission-control.git
synced 2026-08-14 01:07:51 +00:00
fix(tasks): address PR review for assignee wake behavior
This commit is contained in:
+42
-31
@@ -705,16 +705,18 @@ async def _notify_agent_on_task_assign(
|
||||
board: Board,
|
||||
task: Task,
|
||||
agent: Agent,
|
||||
wake_assignee: bool = True,
|
||||
) -> None:
|
||||
if not agent.openclaw_session_id:
|
||||
return
|
||||
await _wake_agent_online_for_task(
|
||||
session=session,
|
||||
board=board,
|
||||
task=task,
|
||||
agent=agent,
|
||||
reason="assignment",
|
||||
)
|
||||
if wake_assignee:
|
||||
await _wake_agent_online_for_task(
|
||||
session=session,
|
||||
board=board,
|
||||
task=task,
|
||||
agent=agent,
|
||||
reason="assignment",
|
||||
)
|
||||
dispatch = GatewayDispatchService(session)
|
||||
config = await dispatch.optional_gateway_config_for_board(board)
|
||||
if config is None:
|
||||
@@ -2566,21 +2568,24 @@ async def _notify_task_update_assignment_changes(
|
||||
*,
|
||||
update: _TaskUpdateInput,
|
||||
) -> None:
|
||||
board = (
|
||||
await Board.objects.by_id(update.task.board_id).first(session)
|
||||
if update.task.board_id
|
||||
else None
|
||||
)
|
||||
board: Board | None = None
|
||||
|
||||
async def _board() -> Board | None:
|
||||
nonlocal board
|
||||
if board is None and update.task.board_id:
|
||||
board = await Board.objects.by_id(update.task.board_id).first(session)
|
||||
return board
|
||||
|
||||
if (
|
||||
update.task.status == "inbox"
|
||||
and update.task.assigned_agent_id is None
|
||||
and (update.previous_status != "inbox" or update.previous_assigned is not None)
|
||||
):
|
||||
if board:
|
||||
current_board = await _board()
|
||||
if current_board:
|
||||
await _notify_lead_on_task_unassigned(
|
||||
session=session,
|
||||
board=board,
|
||||
board=current_board,
|
||||
task=update.task,
|
||||
)
|
||||
|
||||
@@ -2593,20 +2598,23 @@ async def _notify_task_update_assignment_changes(
|
||||
if assigned_agent is None:
|
||||
return
|
||||
|
||||
if (
|
||||
board
|
||||
and update.task.status == "in_progress"
|
||||
and update.previous_status != "in_progress"
|
||||
):
|
||||
await _wake_agent_online_for_task(
|
||||
session=session,
|
||||
board=board,
|
||||
task=update.task,
|
||||
agent=assigned_agent,
|
||||
reason="status_in_progress",
|
||||
)
|
||||
assignment_changed = update.task.assigned_agent_id != update.previous_assigned
|
||||
entered_in_progress = (
|
||||
update.task.status == "in_progress" and update.previous_status != "in_progress"
|
||||
)
|
||||
|
||||
if update.task.assigned_agent_id == update.previous_assigned:
|
||||
if entered_in_progress and not assignment_changed:
|
||||
current_board = await _board()
|
||||
if current_board:
|
||||
await _wake_agent_online_for_task(
|
||||
session=session,
|
||||
board=current_board,
|
||||
task=update.task,
|
||||
agent=assigned_agent,
|
||||
reason="status_in_progress",
|
||||
)
|
||||
|
||||
if not assignment_changed:
|
||||
return
|
||||
|
||||
if (
|
||||
@@ -2616,10 +2624,11 @@ async def _notify_task_update_assignment_changes(
|
||||
and update.actor.agent
|
||||
and update.actor.agent.is_board_lead
|
||||
):
|
||||
if board:
|
||||
current_board = await _board()
|
||||
if current_board:
|
||||
await _notify_agent_on_task_rework(
|
||||
session=session,
|
||||
board=board,
|
||||
board=current_board,
|
||||
task=update.task,
|
||||
agent=assigned_agent,
|
||||
lead=update.actor.agent,
|
||||
@@ -2633,12 +2642,14 @@ async def _notify_task_update_assignment_changes(
|
||||
):
|
||||
return
|
||||
|
||||
if board:
|
||||
current_board = await _board()
|
||||
if current_board:
|
||||
await _notify_agent_on_task_assign(
|
||||
session=session,
|
||||
board=board,
|
||||
board=current_board,
|
||||
task=update.task,
|
||||
agent=assigned_agent,
|
||||
wake_assignee=True,
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -913,3 +913,205 @@ async def test_non_lead_agent_moves_to_review_without_comment_or_recent_comment_
|
||||
assert exc.value.detail == "Comment is required."
|
||||
finally:
|
||||
await engine.dispose()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_lead_assignment_and_in_progress_wakes_assignee_once(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
async def _fake_send_agent_task_message(**_: Any) -> str | None:
|
||||
return None
|
||||
|
||||
monkeypatch.setattr(tasks_api, "_send_agent_task_message", _fake_send_agent_task_message)
|
||||
|
||||
engine = await _make_engine()
|
||||
try:
|
||||
async with await _make_session(engine) as session:
|
||||
org_id = uuid4()
|
||||
board_id = uuid4()
|
||||
gateway_id = uuid4()
|
||||
lead_id = uuid4()
|
||||
worker_id = uuid4()
|
||||
task_id = uuid4()
|
||||
|
||||
session.add(Organization(id=org_id, name="org"))
|
||||
session.add(
|
||||
Gateway(
|
||||
id=gateway_id,
|
||||
organization_id=org_id,
|
||||
name="gateway",
|
||||
url="https://gateway.local",
|
||||
workspace_root="/tmp/workspace",
|
||||
),
|
||||
)
|
||||
session.add(
|
||||
Board(
|
||||
id=board_id,
|
||||
organization_id=org_id,
|
||||
name="board",
|
||||
slug="board",
|
||||
gateway_id=gateway_id,
|
||||
),
|
||||
)
|
||||
session.add(
|
||||
Agent(
|
||||
id=lead_id,
|
||||
name="lead",
|
||||
board_id=board_id,
|
||||
gateway_id=gateway_id,
|
||||
status="online",
|
||||
is_board_lead=True,
|
||||
openclaw_session_id="session-lead",
|
||||
),
|
||||
)
|
||||
session.add(
|
||||
Agent(
|
||||
id=worker_id,
|
||||
name="worker",
|
||||
board_id=board_id,
|
||||
gateway_id=gateway_id,
|
||||
status="offline",
|
||||
openclaw_session_id="session-worker",
|
||||
),
|
||||
)
|
||||
session.add(
|
||||
Task(
|
||||
id=task_id,
|
||||
board_id=board_id,
|
||||
title="assignment wake",
|
||||
description="",
|
||||
status="inbox",
|
||||
assigned_agent_id=None,
|
||||
),
|
||||
)
|
||||
await session.commit()
|
||||
|
||||
task = (await session.exec(select(Task).where(col(Task.id) == task_id))).first()
|
||||
assert task is not None
|
||||
lead = (await session.exec(select(Agent).where(col(Agent.id) == lead_id))).first()
|
||||
assert lead is not None
|
||||
|
||||
updated = await tasks_api.update_task(
|
||||
payload=TaskUpdate(assigned_agent_id=worker_id, status="in_progress"),
|
||||
task=task,
|
||||
session=session,
|
||||
actor=ActorContext(actor_type="agent", agent=lead),
|
||||
)
|
||||
|
||||
assert updated.status == "in_progress"
|
||||
assert updated.assigned_agent_id == worker_id
|
||||
|
||||
reloaded_worker = (
|
||||
await session.exec(select(Agent).where(col(Agent.id) == worker_id))
|
||||
).first()
|
||||
assert reloaded_worker is not None
|
||||
assert reloaded_worker.status == "online"
|
||||
assert reloaded_worker.last_seen_at is not None
|
||||
|
||||
wake_events = (
|
||||
await session.exec(
|
||||
select(ActivityEvent)
|
||||
.where(col(ActivityEvent.task_id) == task_id)
|
||||
.where(col(ActivityEvent.event_type) == "task.assignee_woken"),
|
||||
)
|
||||
).all()
|
||||
assert len(wake_events) == 1
|
||||
assert wake_events[0].message is not None
|
||||
assert "(assignment)" in wake_events[0].message
|
||||
finally:
|
||||
await engine.dispose()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_entering_in_progress_with_existing_assignee_wakes_assignee(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
async def _fake_send_agent_task_message(**_: Any) -> str | None:
|
||||
return None
|
||||
|
||||
monkeypatch.setattr(tasks_api, "_send_agent_task_message", _fake_send_agent_task_message)
|
||||
|
||||
engine = await _make_engine()
|
||||
try:
|
||||
async with await _make_session(engine) as session:
|
||||
org_id = uuid4()
|
||||
board_id = uuid4()
|
||||
gateway_id = uuid4()
|
||||
worker_id = uuid4()
|
||||
task_id = uuid4()
|
||||
|
||||
session.add(Organization(id=org_id, name="org"))
|
||||
session.add(
|
||||
Gateway(
|
||||
id=gateway_id,
|
||||
organization_id=org_id,
|
||||
name="gateway",
|
||||
url="https://gateway.local",
|
||||
workspace_root="/tmp/workspace",
|
||||
),
|
||||
)
|
||||
session.add(
|
||||
Board(
|
||||
id=board_id,
|
||||
organization_id=org_id,
|
||||
name="board",
|
||||
slug="board",
|
||||
gateway_id=gateway_id,
|
||||
),
|
||||
)
|
||||
session.add(
|
||||
Agent(
|
||||
id=worker_id,
|
||||
name="worker",
|
||||
board_id=board_id,
|
||||
gateway_id=gateway_id,
|
||||
status="offline",
|
||||
openclaw_session_id="session-worker",
|
||||
),
|
||||
)
|
||||
session.add(
|
||||
Task(
|
||||
id=task_id,
|
||||
board_id=board_id,
|
||||
title="status wake",
|
||||
description="",
|
||||
status="inbox",
|
||||
assigned_agent_id=worker_id,
|
||||
),
|
||||
)
|
||||
await session.commit()
|
||||
|
||||
task = (await session.exec(select(Task).where(col(Task.id) == task_id))).first()
|
||||
assert task is not None
|
||||
worker = (await session.exec(select(Agent).where(col(Agent.id) == worker_id))).first()
|
||||
assert worker is not None
|
||||
|
||||
updated = await tasks_api.update_task(
|
||||
payload=TaskUpdate(status="in_progress"),
|
||||
task=task,
|
||||
session=session,
|
||||
actor=ActorContext(actor_type="agent", agent=worker),
|
||||
)
|
||||
|
||||
assert updated.status == "in_progress"
|
||||
assert updated.assigned_agent_id == worker_id
|
||||
|
||||
reloaded_worker = (
|
||||
await session.exec(select(Agent).where(col(Agent.id) == worker_id))
|
||||
).first()
|
||||
assert reloaded_worker is not None
|
||||
assert reloaded_worker.status == "online"
|
||||
assert reloaded_worker.last_seen_at is not None
|
||||
|
||||
wake_events = (
|
||||
await session.exec(
|
||||
select(ActivityEvent)
|
||||
.where(col(ActivityEvent.task_id) == task_id)
|
||||
.where(col(ActivityEvent.event_type) == "task.assignee_woken"),
|
||||
)
|
||||
).all()
|
||||
assert len(wake_events) == 1
|
||||
assert wake_events[0].message is not None
|
||||
assert "(status_in_progress)" in wake_events[0].message
|
||||
finally:
|
||||
await engine.dispose()
|
||||
|
||||
Reference in New Issue
Block a user