feat(feishu): 添加飞书入站事件inbox和混合数据库协调功能 - 实现飞书入站事件持久化inbox机制,支持状态管理、租约锁定和重试退避 - 添加混合数据库基线协调工具,确保平台PostgreSQL结构安全对齐 - 增加运行组件心跳检测和readiness就绪检查机制 - 实现app_ticket事件的安全轮换和验证处理 - 添加生产环境运行编排和fail-closed安全机制 - 支持webhook快速确认和长连接独立进程处理 - 完善个人数据擦除时的待处理事件清理功能 ```
339 lines
12 KiB
Python
339 lines
12 KiB
Python
import json
|
|
import sys
|
|
from collections.abc import Iterator
|
|
from types import SimpleNamespace
|
|
from typing import Any
|
|
|
|
import pytest
|
|
from fastapi import HTTPException
|
|
from sqlalchemy import create_engine, func, select
|
|
from sqlalchemy.orm import Session, sessionmaker
|
|
from sqlalchemy.pool import StaticPool
|
|
|
|
from app.application.feishu.events import FeishuEventService, _identifier_digest
|
|
from app.core.config import get_settings
|
|
from app.core.database import Base
|
|
from app.modules.audit.models import AuditLog
|
|
from app.modules.feishu import long_connection
|
|
from app.modules.feishu.app_tickets import FeishuAppTicketService
|
|
from app.modules.feishu.constants import FeishuEventSource
|
|
from app.modules.feishu.models import FeishuAppTicket, FeishuEventReceipt
|
|
from app.modules.feishu.services import FeishuInboundService
|
|
|
|
|
|
@pytest.fixture
|
|
def session_factory(
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> Iterator[sessionmaker[Session]]:
|
|
monkeypatch.setenv("FEISHU_APP_ID", "cli-ticket-app")
|
|
monkeypatch.setenv("FEISHU_APP_SECRET", "ticket-app-secret")
|
|
monkeypatch.setenv("FEISHU_EVENT_TRANSPORT", "long_connection")
|
|
monkeypatch.setenv("FEISHU_VERIFICATION_TOKEN", "ticket-token")
|
|
monkeypatch.setenv("FEISHU_USER_FEATURES_ENABLED", "false")
|
|
get_settings.cache_clear()
|
|
engine = create_engine(
|
|
"sqlite://",
|
|
connect_args={"check_same_thread": False},
|
|
poolclass=StaticPool,
|
|
)
|
|
Base.metadata.create_all(
|
|
engine,
|
|
tables=[
|
|
AuditLog.__table__,
|
|
FeishuEventReceipt.__table__,
|
|
FeishuAppTicket.__table__,
|
|
],
|
|
)
|
|
factory = sessionmaker(bind=engine, expire_on_commit=False)
|
|
try:
|
|
yield factory
|
|
finally:
|
|
Base.metadata.drop_all(
|
|
engine,
|
|
tables=[
|
|
FeishuAppTicket.__table__,
|
|
FeishuEventReceipt.__table__,
|
|
AuditLog.__table__,
|
|
],
|
|
)
|
|
engine.dispose()
|
|
get_settings.cache_clear()
|
|
|
|
|
|
def _v2_ticket_event(
|
|
event_id: str,
|
|
ticket: str,
|
|
*,
|
|
app_id: str = "cli-ticket-app",
|
|
) -> dict[str, Any]:
|
|
return {
|
|
"schema": "2.0",
|
|
"header": {
|
|
"event_id": event_id,
|
|
"event_type": "app_ticket",
|
|
"token": "ticket-token",
|
|
"app_id": app_id,
|
|
},
|
|
"event": {"app_ticket": ticket},
|
|
}
|
|
|
|
|
|
def _v1_ticket_event(
|
|
event_id: str,
|
|
ticket: str,
|
|
*,
|
|
app_id: str = "cli-ticket-app",
|
|
) -> dict[str, Any]:
|
|
return {
|
|
"ts": "1785081600.000",
|
|
"uuid": event_id,
|
|
"token": "ticket-token",
|
|
"type": "app_ticket",
|
|
"event": {
|
|
"app_id": app_id,
|
|
"app_ticket": ticket,
|
|
},
|
|
}
|
|
|
|
|
|
def test_verified_ticket_is_deduplicated_rotated_and_never_leaked(
|
|
session_factory: sessionmaker[Session],
|
|
) -> None:
|
|
first_ticket = "ticket-secret-first"
|
|
rotated_ticket = "ticket-secret-rotated"
|
|
with session_factory() as db:
|
|
service = FeishuEventService(db)
|
|
first = service._handle_verified_event(
|
|
_v2_ticket_event("ticket-event-1", first_ticket),
|
|
source=FeishuEventSource.WEBHOOK,
|
|
)
|
|
stored = db.scalar(select(FeishuAppTicket))
|
|
assert stored is not None
|
|
stored_id = stored.id
|
|
first_received_at = stored.received_at
|
|
assert FeishuAppTicketService(db).get_ticket("cli-ticket-app") == first_ticket
|
|
assert first_ticket not in json.dumps(first, ensure_ascii=False)
|
|
|
|
duplicate_payload = _v2_ticket_event(
|
|
"ticket-event-1",
|
|
"ticket-secret-duplicate-must-not-win",
|
|
)
|
|
duplicate = service._handle_verified_event(
|
|
duplicate_payload,
|
|
source=FeishuEventSource.LONG_CONNECTION,
|
|
)
|
|
assert duplicate["duplicate"] is True
|
|
assert FeishuAppTicketService(db).get_ticket("cli-ticket-app") == first_ticket
|
|
|
|
rotated = service._handle_verified_event(
|
|
_v2_ticket_event("ticket-event-2", rotated_ticket),
|
|
source=FeishuEventSource.LONG_CONNECTION,
|
|
)
|
|
db.expire_all()
|
|
current = db.scalar(select(FeishuAppTicket))
|
|
assert current is not None
|
|
assert current.id == stored_id
|
|
assert current.app_ticket == rotated_ticket
|
|
assert current.received_at >= first_received_at
|
|
assert db.scalar(select(func.count()).select_from(FeishuAppTicket)) == 1
|
|
assert db.scalar(select(func.count()).select_from(FeishuEventReceipt)) == 2
|
|
receipts = list(db.execute(select(FeishuEventReceipt)).scalars())
|
|
assert all(item.status == "succeeded" for item in receipts)
|
|
assert all(item.payload is None for item in receipts)
|
|
assert rotated_ticket not in json.dumps(rotated, ensure_ascii=False)
|
|
|
|
audits = list(db.execute(select(AuditLog)).scalars())
|
|
assert len(audits) == 2
|
|
serialized_audits = json.dumps(
|
|
[
|
|
{
|
|
"request": item.request_payload,
|
|
"response": item.response_payload,
|
|
"target": item.target_id,
|
|
}
|
|
for item in audits
|
|
],
|
|
ensure_ascii=False,
|
|
)
|
|
assert first_ticket not in serialized_audits
|
|
assert rotated_ticket not in serialized_audits
|
|
assert "ticket-secret-duplicate-must-not-win" not in serialized_audits
|
|
|
|
|
|
def test_only_verified_matching_app_ticket_events_can_write(
|
|
session_factory: sessionmaker[Session],
|
|
) -> None:
|
|
with session_factory() as db:
|
|
service = FeishuEventService(db)
|
|
unverified = _v2_ticket_event("unverified-ticket", "unverified-secret")
|
|
unverified["header"]["token"] = "invalid-token"
|
|
with pytest.raises(HTTPException) as unverified_error:
|
|
service.handle_event(unverified, source=FeishuEventSource.WEBHOOK)
|
|
assert unverified_error.value.status_code == 401
|
|
|
|
with pytest.raises(HTTPException) as mismatch_error:
|
|
service._handle_verified_event(
|
|
_v2_ticket_event(
|
|
"wrong-app-ticket",
|
|
"wrong-app-secret",
|
|
app_id="cli-other-app",
|
|
),
|
|
source=FeishuEventSource.WEBHOOK,
|
|
)
|
|
assert mismatch_error.value.status_code == 401
|
|
|
|
missing_ticket = _v2_ticket_event("missing-ticket", "")
|
|
with pytest.raises(HTTPException) as missing_error:
|
|
service._handle_verified_event(
|
|
missing_ticket,
|
|
source=FeishuEventSource.WEBHOOK,
|
|
)
|
|
assert missing_error.value.status_code == 400
|
|
|
|
not_ticket_event = _v2_ticket_event("ordinary-event", "must-not-store")
|
|
not_ticket_event["header"]["event_type"] = "im.message.receive_v1"
|
|
result = service._handle_verified_event(
|
|
not_ticket_event,
|
|
source=FeishuEventSource.WEBHOOK,
|
|
auto_reply=False,
|
|
)
|
|
assert result["handled"] is False
|
|
assert db.scalar(select(FeishuAppTicket)) is None
|
|
|
|
|
|
def test_app_ticket_write_and_receipt_success_are_atomic(
|
|
session_factory: sessionmaker[Session],
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
payload = _v2_ticket_event("atomic-app-ticket", "atomic-ticket-secret")
|
|
original_mark_success = FeishuInboundService._mark_success
|
|
fail_before_commit = True
|
|
|
|
def mark_success(
|
|
atomic_db: Session,
|
|
event_key: str,
|
|
lock_owner: str,
|
|
current: Any,
|
|
) -> None:
|
|
nonlocal fail_before_commit
|
|
if fail_before_commit:
|
|
fail_before_commit = False
|
|
raise RuntimeError("simulated crash before atomic app-ticket commit")
|
|
original_mark_success(atomic_db, event_key, lock_owner, current)
|
|
|
|
monkeypatch.setattr(
|
|
FeishuInboundService,
|
|
"_mark_success",
|
|
staticmethod(mark_success),
|
|
)
|
|
with session_factory() as db:
|
|
service = FeishuEventService(db)
|
|
with pytest.raises(HTTPException) as first_error:
|
|
service._handle_verified_event(
|
|
payload,
|
|
source=FeishuEventSource.WEBHOOK,
|
|
)
|
|
|
|
assert first_error.value.status_code == 503
|
|
assert db.scalar(select(FeishuAppTicket)) is None
|
|
receipt = db.scalar(select(FeishuEventReceipt))
|
|
assert receipt is not None
|
|
assert receipt.status == "retry"
|
|
|
|
result = service._handle_verified_event(
|
|
payload,
|
|
source=FeishuEventSource.WEBHOOK,
|
|
)
|
|
|
|
assert result == {"ok": True, "handled": True}
|
|
assert FeishuAppTicketService(db).get_ticket("cli-ticket-app") == (
|
|
"atomic-ticket-secret"
|
|
)
|
|
db.refresh(receipt)
|
|
assert receipt.status == "succeeded"
|
|
assert receipt.attempt_count == 2
|
|
|
|
|
|
def test_v1_app_ticket_payload_uses_uuid_receipt(
|
|
session_factory: sessionmaker[Session],
|
|
) -> None:
|
|
ticket = "v1-ticket-secret"
|
|
with session_factory() as db:
|
|
result = FeishuEventService(db).handle_event(
|
|
_v1_ticket_event("v1-ticket-uuid", ticket),
|
|
source=FeishuEventSource.WEBHOOK,
|
|
)
|
|
|
|
receipt = db.scalar(select(FeishuEventReceipt))
|
|
assert result == {"ok": True, "handled": True}
|
|
assert receipt is not None
|
|
assert receipt.event_id == _identifier_digest(
|
|
"v1-ticket-uuid",
|
|
"event-id",
|
|
)
|
|
assert receipt.event_key == _identifier_digest(
|
|
"cli-ticket-app:app_ticket:v1-ticket-uuid",
|
|
"event-key",
|
|
)
|
|
assert "v1-ticket-uuid" not in receipt.event_id
|
|
assert "cli-ticket-app" not in receipt.event_key
|
|
assert receipt.status == "succeeded"
|
|
assert receipt.payload is None
|
|
assert FeishuAppTicketService(db).get_ticket("cli-ticket-app") == ticket
|
|
assert ticket not in json.dumps(result, ensure_ascii=False)
|
|
|
|
|
|
def test_long_connection_registers_custom_app_ticket_handler(
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
registrations: dict[str, Any] = {}
|
|
|
|
class FakeBuilder:
|
|
def register_p2_im_message_receive_v1(self, handler: Any) -> "FakeBuilder":
|
|
registrations["message"] = handler
|
|
return self
|
|
|
|
def register_p1_customized_event(
|
|
self,
|
|
event_type: str,
|
|
handler: Any,
|
|
) -> "FakeBuilder":
|
|
registrations[event_type] = handler
|
|
return self
|
|
|
|
def build(self) -> "FakeBuilder":
|
|
return self
|
|
|
|
class FakeDispatcherHandler:
|
|
@staticmethod
|
|
def builder(encrypt_key: str, verification_token: str) -> FakeBuilder:
|
|
registrations["builder_args"] = (encrypt_key, verification_token)
|
|
return FakeBuilder()
|
|
|
|
class FakeClient:
|
|
def __init__(self, **kwargs: Any):
|
|
registrations["client_kwargs"] = kwargs
|
|
|
|
def start(self) -> None:
|
|
registrations["started"] = True
|
|
|
|
fake_lark = SimpleNamespace(
|
|
EventDispatcherHandler=FakeDispatcherHandler,
|
|
LogLevel=SimpleNamespace(WARNING="warning"),
|
|
ws=SimpleNamespace(Client=FakeClient),
|
|
)
|
|
monkeypatch.setitem(sys.modules, "lark_oapi", fake_lark)
|
|
monkeypatch.setenv("FEISHU_APP_ID", "cli-ticket-app")
|
|
monkeypatch.setenv("FEISHU_APP_SECRET", "ticket-app-secret")
|
|
monkeypatch.setenv("FEISHU_EVENT_TRANSPORT", "long_connection")
|
|
monkeypatch.setenv("FEISHU_VERIFICATION_TOKEN", "ticket-token")
|
|
get_settings.cache_clear()
|
|
try:
|
|
long_connection.run_long_connection()
|
|
finally:
|
|
get_settings.cache_clear()
|
|
|
|
assert registrations["app_ticket"] is long_connection._handle_app_ticket_event
|
|
assert registrations["message"] is long_connection._handle_message_event
|
|
assert registrations["started"] is True
|