feat(feishu): 添加飞书入站事件inbox和混合数据库协调功能 - 实现飞书入站事件持久化inbox机制,支持状态管理、租约锁定和重试退避 - 添加混合数据库基线协调工具,确保平台PostgreSQL结构安全对齐 - 增加运行组件心跳检测和readiness就绪检查机制 - 实现app_ticket事件的安全轮换和验证处理 - 添加生产环境运行编排和fail-closed安全机制 - 支持webhook快速确认和长连接独立进程处理 - 完善个人数据擦除时的待处理事件清理功能 ```
362 lines
10 KiB
Python
362 lines
10 KiB
Python
from collections.abc import Iterator
|
|
from datetime import datetime, timedelta
|
|
from types import SimpleNamespace
|
|
from uuid import uuid4
|
|
|
|
import pytest
|
|
from fastapi import FastAPI
|
|
from fastapi.testclient import TestClient
|
|
from sqlalchemy import create_engine
|
|
from sqlalchemy.orm import Session, sessionmaker
|
|
from sqlalchemy.pool import StaticPool
|
|
|
|
import app.core.background.task_queue.subscriptions as subscription_queue
|
|
from app.application.scheduling import create_scheduler
|
|
from app.core.config import get_settings
|
|
from app.core.database import Base, get_db
|
|
from app.modules.feishu_users.models import FeishuUser
|
|
from app.modules.observability.constants import (
|
|
ObservabilityKey,
|
|
ObservabilityStatus,
|
|
)
|
|
from app.modules.observability.routes import router as observability_router
|
|
from app.modules.subscriptions.constants import (
|
|
PushDeliveryStatus,
|
|
PushSubscriptionStatus,
|
|
SubscriptionScheduleType,
|
|
SubscriptionTargetType,
|
|
)
|
|
from app.modules.subscriptions.models import PushDelivery, PushSubscription
|
|
from app.modules.subscriptions.routes import router as subscriptions_router
|
|
|
|
|
|
@pytest.fixture(autouse=True)
|
|
def _reset_settings() -> Iterator[None]:
|
|
get_settings.cache_clear()
|
|
yield
|
|
get_settings.cache_clear()
|
|
|
|
|
|
@pytest.fixture
|
|
def session_factory(
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> Iterator[sessionmaker[Session]]:
|
|
monkeypatch.setenv("API_KEY", "subscription-service-key")
|
|
monkeypatch.setenv("API_KEYS", "[]")
|
|
get_settings.cache_clear()
|
|
engine = create_engine(
|
|
"sqlite://",
|
|
connect_args={"check_same_thread": False},
|
|
poolclass=StaticPool,
|
|
)
|
|
Base.metadata.create_all(engine)
|
|
factory = sessionmaker(bind=engine, expire_on_commit=False)
|
|
try:
|
|
yield factory
|
|
finally:
|
|
engine.dispose()
|
|
|
|
|
|
def _seed_subscription(
|
|
db: Session,
|
|
*,
|
|
suffix: str,
|
|
status_value: str = PushSubscriptionStatus.ACTIVE,
|
|
) -> PushSubscription:
|
|
owner = FeishuUser(
|
|
code=f"FSU-{suffix}",
|
|
tenant_key=f"tenant-{suffix}",
|
|
open_id=f"open-{suffix}",
|
|
)
|
|
db.add(owner)
|
|
db.flush()
|
|
subscription = PushSubscription(
|
|
code=f"SUB-{suffix}",
|
|
owner_id=owner.id,
|
|
target_type=SubscriptionTargetType.USER,
|
|
target_id=owner.open_id,
|
|
prompt="生成一条个人提醒",
|
|
schedule_type=SubscriptionScheduleType.DAILY,
|
|
schedule_config={"hour": 9, "minute": 0},
|
|
timezone="Asia/Shanghai",
|
|
next_run_at=(
|
|
datetime(2030, 1, 1, 1, 0)
|
|
if status_value == PushSubscriptionStatus.ACTIVE
|
|
else None
|
|
),
|
|
status=status_value,
|
|
consented_at=datetime(2026, 7, 26, 1, 0),
|
|
)
|
|
db.add(subscription)
|
|
db.commit()
|
|
db.refresh(subscription)
|
|
return subscription
|
|
|
|
|
|
def _seed_delivery(
|
|
db: Session,
|
|
subscription: PushSubscription,
|
|
*,
|
|
suffix: str,
|
|
status_value: str,
|
|
) -> PushDelivery:
|
|
delivery = PushDelivery(
|
|
code=f"DEL-{suffix}",
|
|
subscription_id=subscription.id,
|
|
scheduled_for=datetime(2026, 7, 27, 1, 0),
|
|
idempotency_key=uuid4().hex + uuid4().hex,
|
|
message_uuid=str(uuid4()),
|
|
status=status_value,
|
|
attempt_count=2,
|
|
next_attempt_at=datetime(2030, 1, 1, 1, 0),
|
|
last_error="temporary provider failure",
|
|
)
|
|
db.add(delivery)
|
|
db.commit()
|
|
db.refresh(delivery)
|
|
return delivery
|
|
|
|
|
|
def _test_app(
|
|
session_factory: sessionmaker[Session],
|
|
*,
|
|
include_subscriptions: bool = False,
|
|
include_observability: bool = False,
|
|
) -> FastAPI:
|
|
app = FastAPI()
|
|
if include_subscriptions:
|
|
app.include_router(
|
|
subscriptions_router,
|
|
prefix="/api/v1/subscriptions",
|
|
)
|
|
if include_observability:
|
|
app.include_router(observability_router, prefix="/api/v1")
|
|
|
|
def override_db() -> Iterator[Session]:
|
|
with session_factory() as db:
|
|
yield db
|
|
|
|
app.dependency_overrides[get_db] = override_db
|
|
return app
|
|
|
|
|
|
def test_subscription_internal_apis_require_key_and_serialize_records(
|
|
session_factory: sessionmaker[Session],
|
|
) -> None:
|
|
with session_factory() as db:
|
|
subscription = _seed_subscription(db, suffix="api")
|
|
delivery = _seed_delivery(
|
|
db,
|
|
subscription,
|
|
suffix="api",
|
|
status_value=PushDeliveryStatus.RETRY,
|
|
)
|
|
|
|
client = TestClient(
|
|
_test_app(session_factory, include_subscriptions=True)
|
|
)
|
|
assert client.get("/api/v1/subscriptions").status_code == 401
|
|
assert client.get("/api/v1/subscriptions/deliveries").status_code == 401
|
|
assert (
|
|
client.get(
|
|
"/api/v1/subscriptions",
|
|
headers={"X-API-Key": "invalid-key"},
|
|
).status_code
|
|
== 401
|
|
)
|
|
|
|
headers = {"X-API-Key": "subscription-service-key"}
|
|
subscriptions_response = client.get(
|
|
"/api/v1/subscriptions",
|
|
headers=headers,
|
|
)
|
|
deliveries_response = client.get(
|
|
"/api/v1/subscriptions/deliveries",
|
|
headers=headers,
|
|
)
|
|
|
|
assert subscriptions_response.status_code == 200
|
|
subscriptions_body = subscriptions_response.json()
|
|
assert subscriptions_body["total"] == 1
|
|
subscription_item = subscriptions_body["items"][0]
|
|
assert set(subscription_item) == {
|
|
"code",
|
|
"owner_id",
|
|
"target_type",
|
|
"target_id",
|
|
"prompt",
|
|
"schedule_type",
|
|
"schedule_config",
|
|
"timezone",
|
|
"next_run_at",
|
|
"status",
|
|
"consented_at",
|
|
"last_run_at",
|
|
"created_at",
|
|
"updated_at",
|
|
}
|
|
assert subscription_item["code"] == subscription.code
|
|
assert subscription_item["schedule_config"] == {"hour": 9, "minute": 0}
|
|
assert subscription_item["next_run_at"] is not None
|
|
|
|
assert deliveries_response.status_code == 200
|
|
deliveries_body = deliveries_response.json()
|
|
assert deliveries_body["total"] == 1
|
|
delivery_item = deliveries_body["items"][0]
|
|
assert set(delivery_item) == {
|
|
"code",
|
|
"subscription_id",
|
|
"scheduled_for",
|
|
"idempotency_key",
|
|
"message_uuid",
|
|
"status",
|
|
"attempt_count",
|
|
"next_attempt_at",
|
|
"provider_message_id",
|
|
"last_error",
|
|
"sent_at",
|
|
"created_at",
|
|
"updated_at",
|
|
}
|
|
assert delivery_item["code"] == delivery.code
|
|
assert delivery_item["message_uuid"] == delivery.message_uuid
|
|
assert delivery_item["attempt_count"] == 2
|
|
assert "rendered_content" not in delivery_item
|
|
|
|
|
|
@pytest.mark.parametrize(
|
|
("features_enabled", "job_expected"),
|
|
[(False, False), (True, True)],
|
|
)
|
|
def test_subscription_scheduler_job_follows_feature_flag(
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
features_enabled: bool,
|
|
job_expected: bool,
|
|
) -> None:
|
|
monkeypatch.setenv(
|
|
"FEISHU_USER_FEATURES_ENABLED",
|
|
str(features_enabled).lower(),
|
|
)
|
|
get_settings.cache_clear()
|
|
|
|
scheduler = create_scheduler()
|
|
job = scheduler.get_job("subscription_delivery_cycle")
|
|
|
|
assert scheduler.running is False
|
|
assert (job is not None) is job_expected
|
|
if job is not None:
|
|
assert job.trigger.interval == timedelta(minutes=1)
|
|
assert scheduler._job_defaults["coalesce"] is True
|
|
assert scheduler._job_defaults["max_instances"] == 1
|
|
|
|
|
|
def test_subscription_queue_runs_inline_without_celery(
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
monkeypatch.setenv("TASK_QUEUE_ENABLED", "false")
|
|
get_settings.cache_clear()
|
|
calls: list[str] = []
|
|
|
|
def fake_cycle(*, actor: str) -> dict[str, list[str]]:
|
|
calls.append(actor)
|
|
return {"created": [], "processed": []}
|
|
|
|
monkeypatch.setattr(
|
|
subscription_queue,
|
|
"run_subscription_cycle",
|
|
fake_cycle,
|
|
)
|
|
|
|
result = subscription_queue.enqueue_subscription_cycle(actor="runtime-test")
|
|
|
|
assert result == {
|
|
"queued": False,
|
|
"mode": "inline",
|
|
"task_name": "subscriptions.run_cycle",
|
|
"result": {"created": [], "processed": []},
|
|
}
|
|
assert calls == ["runtime-test"]
|
|
|
|
|
|
def test_subscription_queue_uses_celery_when_enabled(
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
from app.tasks import celery_app
|
|
|
|
monkeypatch.setenv("TASK_QUEUE_ENABLED", "true")
|
|
get_settings.cache_clear()
|
|
calls: dict[str, object] = {}
|
|
|
|
class FakeSignature:
|
|
def apply_async(self) -> SimpleNamespace:
|
|
calls["applied"] = True
|
|
return SimpleNamespace(id="queued-task-id")
|
|
|
|
def fake_signature(
|
|
task_name: str,
|
|
*,
|
|
kwargs: dict[str, str],
|
|
) -> FakeSignature:
|
|
calls["task_name"] = task_name
|
|
calls["kwargs"] = kwargs
|
|
return FakeSignature()
|
|
|
|
monkeypatch.setattr(celery_app, "signature", fake_signature)
|
|
monkeypatch.setattr(
|
|
subscription_queue,
|
|
"run_subscription_cycle",
|
|
lambda **_: pytest.fail("Celery dispatch must not run inline"),
|
|
)
|
|
|
|
result = subscription_queue.enqueue_subscription_cycle(actor="runtime-test")
|
|
|
|
assert result == {
|
|
"queued": True,
|
|
"mode": "celery",
|
|
"task_name": "subscriptions.run_cycle",
|
|
"task_id": "queued-task-id",
|
|
}
|
|
assert calls == {
|
|
"task_name": "subscriptions.run_cycle",
|
|
"kwargs": {"actor": "runtime-test"},
|
|
"applied": True,
|
|
}
|
|
|
|
|
|
def test_ready_route_returns_503_for_processable_delivery_without_credentials(
|
|
session_factory: sessionmaker[Session],
|
|
monkeypatch: pytest.MonkeyPatch,
|
|
) -> None:
|
|
monkeypatch.setenv("FEISHU_APP_TYPE", "self")
|
|
monkeypatch.setenv("FEISHU_APP_ID", "")
|
|
monkeypatch.setenv("FEISHU_APP_SECRET", "")
|
|
monkeypatch.setenv("FEISHU_USER_FEATURES_ENABLED", "true")
|
|
get_settings.cache_clear()
|
|
with session_factory() as db:
|
|
subscription = _seed_subscription(
|
|
db,
|
|
suffix="ready",
|
|
status_value=PushSubscriptionStatus.COMPLETED,
|
|
)
|
|
_seed_delivery(
|
|
db,
|
|
subscription,
|
|
suffix="ready",
|
|
status_value=PushDeliveryStatus.PENDING,
|
|
)
|
|
|
|
client = TestClient(
|
|
_test_app(session_factory, include_observability=True)
|
|
)
|
|
response = client.get("/api/v1/health/ready")
|
|
|
|
assert response.status_code == 503
|
|
body = response.json()
|
|
assert body[ObservabilityKey.STATUS] == ObservabilityStatus.DEGRADED
|
|
check = body[ObservabilityKey.CHECKS]["feishu_subscriptions"]
|
|
assert check[ObservabilityKey.STATUS] == ObservabilityStatus.DEGRADED
|
|
assert check["active"] == 0
|
|
assert check["processable_deliveries"] == 1
|
|
assert check["reasons"] == ["credentials_missing"]
|
|
assert "生成一条个人提醒" not in response.text
|