项目文件夹

文件
2026-07-13 13:12:33 +08:00

400 行
11 KiB
Python

"""Main+systemEvent cron must pin heartbeat reply target via the resolver.
Main-session reminders must not drift to whichever channel touched the session
last. These tests assert the resolver's mapping and that make_system_event_handler
forwards its result through to heartbeat run_once.
"""
from __future__ import annotations
import pytest
import opensquilla.scheduler.handlers as handlers_mod
from opensquilla.scheduler.handlers import (
_resolve_system_event_heartbeat_delivery_override,
make_system_event_handler,
)
from opensquilla.scheduler.types import (
CronJob,
CronWakeMode,
DeliveryConfig,
DeliveryMode,
ReplyTargetSnapshot,
SessionTarget,
)
# --- pure unit tests on the resolver --------------------------------------
def _job(delivery: DeliveryConfig) -> CronJob:
return CronJob(
id="job-1",
cron_expr="* * * * *",
handler_key="system_event",
payload={"kind": "system_event", "text": "x", "agent_id": "main"},
session_target=SessionTarget.MAIN,
delivery=delivery,
)
def test_resolver_uses_explicit_channel_fields_when_channel_mode() -> None:
delivery = DeliveryConfig(
mode=DeliveryMode.CHANNEL,
channel_name="slack",
channel_id="C123",
account_id="T1",
thread_id="th-1",
)
override = _resolve_system_event_heartbeat_delivery_override(_job(delivery))
assert override == {
"channel_name": "slack",
"channel_id": "C123",
"account_id": "T1",
"thread_id": "th-1",
}
def test_resolver_falls_back_to_snapshot_when_mode_none() -> None:
snap = ReplyTargetSnapshot(
channel_name="feishu",
channel_type="feishu",
to="chat-9",
account_id="acc-x",
thread_id="msg-42",
)
delivery = DeliveryConfig(mode=DeliveryMode.NONE, originating_reply_target=snap)
override = _resolve_system_event_heartbeat_delivery_override(_job(delivery))
assert override == {
"channel_name": "feishu",
"channel_id": "chat-9",
"account_id": "acc-x",
"thread_id": "msg-42",
}
def test_resolver_returns_none_when_no_snapshot_and_mode_none() -> None:
override = _resolve_system_event_heartbeat_delivery_override(
_job(DeliveryConfig(mode=DeliveryMode.NONE))
)
assert override is None
def test_resolver_prefers_snapshot_when_origin_mode() -> None:
snap = ReplyTargetSnapshot(
channel_name="feishu",
channel_type="feishu",
to="chat-42",
account_id="",
thread_id="",
)
delivery = DeliveryConfig(
mode=DeliveryMode.ORIGIN,
channel_name="discord",
channel_id="dis-1",
originating_reply_target=snap,
)
override = _resolve_system_event_heartbeat_delivery_override(_job(delivery))
assert override == {
"channel_name": "feishu",
"channel_id": "chat-42",
"account_id": "",
"thread_id": "",
}
def test_resolver_uses_origin_fields_when_origin_without_snapshot() -> None:
delivery = DeliveryConfig(
mode=DeliveryMode.ORIGIN,
channel_name="discord",
channel_id="dis-1",
account_id="",
thread_id="",
)
override = _resolve_system_event_heartbeat_delivery_override(_job(delivery))
assert override == {
"channel_name": "discord",
"channel_id": "dis-1",
"account_id": "",
"thread_id": "",
}
# --- integration: handler forwards override to run_once -------------------
class _StubChain:
async def notify_start(self, job, text):
pass
class _StubSessionManager:
async def get_or_create(self, **kw):
return None
async def append_message(self, *a, **k):
return None
class _CapturingHeartbeat:
def __init__(self) -> None:
self.calls: list[dict] = []
async def run_once(self, **kwargs):
self.calls.append(kwargs)
class _R:
status = "delivered"
reason = ""
delivery_status = "delivered"
return _R()
class _FailingHeartbeat:
def __init__(self) -> None:
self.calls: list[dict] = []
async def run_once(self, **kwargs):
self.calls.append(kwargs)
class _R:
status = "skipped"
reason = "delivery_failed"
delivery_status = "delivery_failed"
return _R()
class _BusyLoop:
def __init__(self) -> None:
self.calls: list[dict] = []
self.requests: list[dict] = []
async def run_once_now(self, **kwargs):
self.calls.append(kwargs)
class _R:
status = "skipped"
reason = "requests-in-flight"
delivery_status = "requests-in-flight"
return _R()
def request_now(self, **kwargs) -> None:
self.requests.append(kwargs)
def _main_job(snapshot: ReplyTargetSnapshot | None) -> CronJob:
delivery = DeliveryConfig(
mode=DeliveryMode.NONE,
originating_reply_target=snapshot,
)
return CronJob(
id="job-1",
cron_expr="* * * * *",
handler_key="system_event",
payload={"kind": "system_event", "text": "stand up", "agent_id": "main"},
session_target=SessionTarget.MAIN,
session_key="main:main",
wake_mode=CronWakeMode.NOW,
delivery=delivery,
timeout_seconds=10.0,
)
async def test_handler_forwards_resolver_output_as_delivery_override(monkeypatch):
"""When the resolver returns an override, the handler must pass it through."""
monkeypatch.setattr(
handlers_mod,
"_build_cron_tool_context",
lambda *a, **kw: object(),
)
snapshot = ReplyTargetSnapshot(
channel_name="feishu",
channel_type="feishu",
to="chat-001",
account_id="acct-xyz",
thread_id="msg-42",
)
job = _main_job(snapshot)
hb = _CapturingHeartbeat()
handler = make_system_event_handler(
delivery_chain=_StubChain(),
session_manager_ref=lambda: _StubSessionManager(),
heartbeat_service_ref=lambda: hb,
heartbeat_loop_ref=lambda: None,
)
await handler(job)
assert hb.calls, "heartbeat run_once must be invoked for wake_mode=NOW"
override = hb.calls[-1].get("delivery_override")
assert override is not None
assert override["channel_name"] == "feishu"
assert override["channel_id"] == "chat-001"
assert override["account_id"] == "acct-xyz"
assert override["thread_id"] == "msg-42"
async def test_handler_fails_pinned_delivery_by_default(monkeypatch):
"""Pinned main-session delivery is required unless best_effort is enabled."""
monkeypatch.setattr(
handlers_mod,
"_build_cron_tool_context",
lambda *a, **kw: object(),
)
snapshot = ReplyTargetSnapshot(
channel_name="feishu",
channel_type="feishu",
to="chat-001",
)
job = _main_job(snapshot)
hb = _FailingHeartbeat()
handler = make_system_event_handler(
delivery_chain=_StubChain(),
session_manager_ref=lambda: _StubSessionManager(),
heartbeat_service_ref=lambda: hb,
heartbeat_loop_ref=lambda: None,
)
with pytest.raises(RuntimeError, match="delivery_failed"):
await handler(job)
async def test_handler_best_effort_pinned_delivery_failure_succeeds(monkeypatch):
monkeypatch.setattr(
handlers_mod,
"_build_cron_tool_context",
lambda *a, **kw: object(),
)
snapshot = ReplyTargetSnapshot(
channel_name="feishu",
channel_type="feishu",
to="chat-001",
)
job = _main_job(snapshot)
job.delivery.best_effort = True
hb = _FailingHeartbeat()
handler = make_system_event_handler(
delivery_chain=_StubChain(),
session_manager_ref=lambda: _StubSessionManager(),
heartbeat_service_ref=lambda: hb,
heartbeat_loop_ref=lambda: None,
)
result = await handler(job)
assert result.delivery_status == "delivery_failed"
async def test_handler_preserves_busy_heartbeat_fallback(monkeypatch):
"""Busy heartbeat fallback queues delivery later and is not a delivery failure."""
monkeypatch.setattr(
handlers_mod,
"_build_cron_tool_context",
lambda *a, **kw: object(),
)
snapshot = ReplyTargetSnapshot(
channel_name="feishu",
channel_type="feishu",
to="chat-001",
)
job = _main_job(snapshot)
loop = _BusyLoop()
handler = make_system_event_handler(
delivery_chain=_StubChain(),
session_manager_ref=lambda: _StubSessionManager(),
heartbeat_service_ref=lambda: _CapturingHeartbeat(),
heartbeat_loop_ref=lambda: loop,
wake_now_busy_max_wait_seconds=0.0,
wake_now_busy_retry_delay_seconds=0.0,
)
result = await handler(job)
assert result.delivery_status == "not-requested"
assert loop.requests
async def test_handler_omits_override_when_resolver_returns_none(monkeypatch):
"""No snapshot + mode NONE → override is not forwarded (legacy behaviour)."""
monkeypatch.setattr(
handlers_mod,
"_build_cron_tool_context",
lambda *a, **kw: object(),
)
job = _main_job(None)
hb = _CapturingHeartbeat()
handler = make_system_event_handler(
delivery_chain=_StubChain(),
session_manager_ref=lambda: _StubSessionManager(),
heartbeat_service_ref=lambda: hb,
heartbeat_loop_ref=lambda: None,
)
await handler(job)
assert hb.calls
assert "delivery_override" not in hb.calls[-1]
class _CapturingLoop:
def __init__(self) -> None:
self.calls: list[dict] = []
async def run_once_now(self, **kwargs):
self.calls.append(kwargs)
class _R:
status = "delivered"
reason = ""
delivery_status = "delivered"
return _R()
async def test_handler_forwards_override_through_loop_run_once_now(monkeypatch):
"""When a heartbeat_loop is available, the override flows through its run_once_now too."""
monkeypatch.setattr(
handlers_mod,
"_build_cron_tool_context",
lambda *a, **kw: object(),
)
snapshot = ReplyTargetSnapshot(
channel_name="feishu",
channel_type="feishu",
to="chat-77",
account_id="",
thread_id="",
)
job = _main_job(snapshot)
loop = _CapturingLoop()
# heartbeat_service is consulted only when loop has no run_once_now;
# provide one anyway so a failure in the loop path doesn't silently fall
# back to a different call site.
hb = _CapturingHeartbeat()
handler = make_system_event_handler(
delivery_chain=_StubChain(),
session_manager_ref=lambda: _StubSessionManager(),
heartbeat_service_ref=lambda: hb,
heartbeat_loop_ref=lambda: loop,
)
await handler(job)
assert loop.calls, "heartbeat_loop.run_once_now must be invoked"
assert hb.calls == [], "service.run_once must not be called when loop path succeeds"
override = loop.calls[-1].get("delivery_override")
assert override == {
"channel_name": "feishu",
"channel_id": "chat-77",
"account_id": "",
"thread_id": "",
}