Files
awoooi/apps/api/tests/test_knowledge_service_readback_degraded.py

718 lines
22 KiB
Python

import asyncio
import sys
from types import SimpleNamespace
import pytest
from fastapi import FastAPI, HTTPException
from fastapi.testclient import TestClient
from src.api.v1 import knowledge as knowledge_api
from src.core.context import clear_project_context, set_project_context
from src.models.knowledge import (
CategoryCount,
EntrySource,
EntryStatus,
EntryType,
KnowledgeAssetTaxonomyCount,
KnowledgeEntry,
KnowledgeListResponse,
)
from src.services import knowledge_service as knowledge_service_module
from src.services.knowledge_service import KnowledgeService
class _BrokenDbContext:
async def __aenter__(self):
raise RuntimeError("db pool exhausted")
async def __aexit__(self, exc_type, exc, tb):
return False
class _OkDbContext:
async def __aenter__(self):
return object()
async def __aexit__(self, exc_type, exc, tb):
return False
@pytest.mark.asyncio
async def test_knowledge_list_entries_fails_soft_when_readback_breaks(monkeypatch) -> None:
monkeypatch.setattr(
knowledge_service_module,
"get_db_context",
lambda: _BrokenDbContext(),
)
service = KnowledgeService.__new__(KnowledgeService)
response = await service.list_entries(limit=50)
assert response.total == 13
assert len(response.items) == 13
assert response.items[0].id == "source-backed-project-awoooi"
assert response.items[0].source == "ai_extracted"
assert response.items[0].status == "approved"
assert [row.category for row in response.categories] == [
"project",
"product",
"website",
"service",
"package",
"tool",
"log",
"alert",
"playbook",
"rag",
"mcp",
"schedule",
"general",
]
assert all(row.count == 1 for row in response.categories)
assert [row.key for row in response.asset_taxonomy] == [
"project",
"product",
"website",
"service",
"package",
"tool",
"log",
"alert",
"playbook",
"rag",
"mcp",
"schedule",
]
assert all(row.count >= 1 for row in response.asset_taxonomy)
assert response.readback_status == "source_backed_degraded"
assert response.primary_readback_ready is False
assert response.degraded_reason_code == "primary_km_db_timeout_or_pool_exhausted"
assert response.operator_stage == "knowledge_readback_source_backed_ai_controlled_repair"
assert response.next_step == "repair_primary_km_db_readback_then_promote_source_backed_receipts_to_persistent_km"
assert response.writes_on_read is False
assert response.manual_review_required is False
@pytest.mark.asyncio
async def test_knowledge_list_entries_source_backed_filter_and_search(monkeypatch) -> None:
monkeypatch.setattr(
knowledge_service_module,
"get_db_context",
lambda: _BrokenDbContext(),
)
service = KnowledgeService.__new__(KnowledgeService)
response = await service.list_entries(category="alert", q="Telegram", limit=50)
assert response.total == 1
assert response.items[0].id == "source-backed-alert-telegram-monitoring-coverage"
assert response.categories[[row.category for row in response.categories].index("alert")].count == 1
assert response.asset_taxonomy[[row.key for row in response.asset_taxonomy].index("alert")].count == 1
@pytest.mark.asyncio
async def test_knowledge_list_entries_retries_transient_pool_timeout(monkeypatch) -> None:
service = KnowledgeService.__new__(KnowledgeService)
calls = 0
async def fake_primary(**_kwargs):
nonlocal calls
calls += 1
if calls == 1:
raise RuntimeError("db pool exhausted")
return KnowledgeListResponse(items=[], total=1, readback_status="ready")
monkeypatch.setattr(service, "_list_entries_from_primary", fake_primary)
response = await service.list_entries(limit=50)
assert calls == 2
assert response.readback_status == "ready"
assert response.primary_readback_ready is True
assert response.operator_stage == "knowledge_readback_primary_retry_recovered"
@pytest.mark.asyncio
async def test_knowledge_list_entries_bounds_primary_timeout_to_source_backed(monkeypatch) -> None:
service = KnowledgeService.__new__(KnowledgeService)
calls = 0
async def never_finishes_primary(**_kwargs):
nonlocal calls
calls += 1
await asyncio.Event().wait()
monkeypatch.setattr(service, "_list_entries_from_primary", never_finishes_primary)
monkeypatch.setattr(knowledge_service_module, "_PRIMARY_KM_LIST_TIMEOUT_SECONDS", 0.01)
monkeypatch.setattr(knowledge_service_module, "_PRIMARY_KM_LIST_RETRY_TIMEOUT_SECONDS", 0.01)
monkeypatch.setattr(knowledge_service_module, "_PRIMARY_KM_RETRY_DELAY_SECONDS", 0)
response = await service.list_entries(limit=50)
assert calls == 2
assert response.readback_status == "source_backed_degraded"
assert response.primary_readback_ready is False
assert response.degraded_reason_code == "primary_km_db_timeout_or_pool_exhausted"
assert response.total == 13
@pytest.mark.asyncio
async def test_knowledge_list_entries_uses_direct_readback_after_session_pool_timeout(monkeypatch) -> None:
service = KnowledgeService.__new__(KnowledgeService)
primary_calls = 0
direct_calls = 0
primary_entry = KnowledgeEntry(
id="km-primary-live-1",
title="Primary KM live row",
content="Persistent KM row recovered through direct readback",
entry_type=EntryType.RUNBOOK,
category="alert_handling",
tags=["telegram", "ai_agent", "log"],
source=EntrySource.AI_EXTRACTED,
status=EntryStatus.APPROVED,
)
async def pool_exhausted_primary(**_kwargs):
nonlocal primary_calls
primary_calls += 1
raise RuntimeError("db pool exhausted")
async def direct_readback(**kwargs):
nonlocal direct_calls
direct_calls += 1
assert kwargs["limit"] == 50
assert kwargs["project_id"] == "awoooi"
return KnowledgeListResponse(
items=[primary_entry],
total=727,
categories=[
CategoryCount(category="alert_handling", count=113),
CategoryCount(category="AI自動化/Ansible受控修復", count=87),
],
asset_taxonomy=[
KnowledgeAssetTaxonomyCount(key="log", count=727),
KnowledgeAssetTaxonomyCount(key="alert", count=443),
KnowledgeAssetTaxonomyCount(key="mcp", count=36),
],
readback_status="ready_direct_connection_after_session_timeout",
primary_readback_ready=True,
operator_stage="knowledge_readback_direct_connection_recovered",
next_step="repair_session_pool_readback_without_hiding_primary_km",
writes_on_read=False,
manual_review_required=False,
)
monkeypatch.setattr(service, "_list_entries_from_primary", pool_exhausted_primary)
monkeypatch.setattr(service, "_list_entries_from_direct", direct_readback)
monkeypatch.setattr(knowledge_service_module, "_PRIMARY_KM_RETRY_DELAY_SECONDS", 0)
response = await service.list_entries(project_id="awoooi", limit=50)
assert primary_calls == 2
assert direct_calls == 1
assert response.total == 727
assert response.items == [primary_entry]
assert response.readback_status == "ready_direct_connection_after_session_timeout"
assert response.primary_readback_ready is True
assert response.operator_stage == "knowledge_readback_direct_connection_recovered"
assert response.next_step == "repair_session_pool_readback_without_hiding_primary_km"
assert response.categories[0].category == "alert_handling"
assert response.asset_taxonomy[0].key == "log"
assert response.writes_on_read is False
assert response.manual_review_required is False
@pytest.mark.asyncio
async def test_direct_list_keeps_primary_rows_when_asset_taxonomy_times_out(
monkeypatch,
) -> None:
service = KnowledgeService.__new__(KnowledgeService)
class DirectConnection:
async def execute(self, *_args):
return "OK"
async def fetchval(self, *_args):
return 727
async def fetch(self, *_args):
return [
{
"id": "km-primary-direct-1",
"title": "Primary direct KM row",
"content": "Persistent row survives taxonomy pressure",
"entry_type": "runbook",
"category": "alert_handling",
"tags": ["alert", "telegram"],
"source": "ai_extracted",
"status": "approved",
"view_count": 0,
}
]
async def close(self):
return None
connection = DirectConnection()
async def connect(_url):
return connection
async def categories(_conn, *, project_id):
assert project_id == "awoooi"
return [CategoryCount(category="alert_handling", count=113)]
async def taxonomy_timeout(_conn, *, project_id):
assert project_id == "awoooi"
raise TimeoutError("taxonomy statement timeout")
monkeypatch.setitem(
sys.modules,
"asyncpg",
SimpleNamespace(connect=connect),
)
monkeypatch.setattr(
knowledge_service_module.settings,
"DATABASE_URL",
"postgresql+asyncpg://test:test@127.0.0.1:5432/test",
)
monkeypatch.setattr(
service,
"_read_categories_direct_with_conn",
categories,
)
monkeypatch.setattr(
service,
"_read_asset_taxonomy_direct_with_conn",
taxonomy_timeout,
)
response = await service._list_entries_from_direct(
project_id="awoooi",
category=None,
entry_type=None,
status=None,
tags=None,
q=None,
limit=1,
offset=0,
)
assert response is not None
assert response.total == 727
assert [item.id for item in response.items] == ["km-primary-direct-1"]
assert response.primary_readback_ready is True
assert response.readback_status == (
"ready_direct_connection_partial_taxonomy_degraded"
)
assert response.degraded_reason_code == (
"primary_km_taxonomy_readback_degraded"
)
assert response.categories[0].category == "alert_handling"
assert any(row.key == "alert" for row in response.asset_taxonomy)
@pytest.mark.asyncio
async def test_knowledge_list_entries_keeps_primary_items_when_taxonomy_degrades(monkeypatch) -> None:
primary_entry = KnowledgeEntry(
id="km-primary-1",
title="Primary KM row",
content="Persistent AI automation memory row",
entry_type=EntryType.RUNBOOK,
category="alert_handling",
tags=["telegram", "ai_agent"],
source=EntrySource.AI_EXTRACTED,
status=EntryStatus.APPROVED,
)
class _TaxonomyBrokenRepo:
def __init__(self, _db):
pass
async def list_entries(self, **_kwargs):
return [primary_entry], 638
async def get_categories(self):
raise RuntimeError("db pool exhausted")
async def get_asset_taxonomy_counts(self):
raise RuntimeError("db pool exhausted")
monkeypatch.setattr(
knowledge_service_module,
"get_db_context",
lambda: _OkDbContext(),
)
monkeypatch.setattr(
knowledge_service_module,
"KnowledgeDBRepository",
_TaxonomyBrokenRepo,
)
service = KnowledgeService.__new__(KnowledgeService)
response = await service.list_entries(limit=50)
assert response.total == 638
assert response.items == [primary_entry]
assert response.readback_status == "ready_partial_taxonomy_degraded"
assert response.primary_readback_ready is True
assert response.degraded_reason_code == "primary_km_taxonomy_readback_degraded"
assert response.operator_stage == "knowledge_primary_entries_ready_taxonomy_degraded"
assert response.categories[[row.category for row in response.categories].index("alert_handling")].count == 1
assert response.asset_taxonomy[[row.key for row in response.asset_taxonomy].index("alert")].count == 1
@pytest.mark.asyncio
async def test_knowledge_search_fails_soft_to_source_backed_entries(monkeypatch) -> None:
monkeypatch.setattr(
knowledge_service_module,
"get_db_context",
lambda: _BrokenDbContext(),
)
service = KnowledgeService.__new__(KnowledgeService)
entries = await service.search("Telegram", limit=10)
entry_ids = {entry.id for entry in entries}
assert all(entry.id.startswith("source-backed-") for entry in entries)
assert {
"source-backed-service-telegram-alert-receipts",
"source-backed-alert-telegram-monitoring-coverage",
"source-backed-schedule-report-monitoring",
}.issubset(entry_ids)
@pytest.mark.asyncio
async def test_knowledge_categories_fails_soft_when_readback_breaks(monkeypatch) -> None:
monkeypatch.setattr(
knowledge_service_module,
"get_db_context",
lambda: _BrokenDbContext(),
)
service = KnowledgeService.__new__(KnowledgeService)
categories = await service.get_categories()
assert [row.category for row in categories] == [
"project",
"product",
"website",
"service",
"package",
"tool",
"log",
"alert",
"playbook",
"rag",
"mcp",
"schedule",
"general",
]
assert all(row.count == 1 for row in categories)
@pytest.mark.asyncio
async def test_knowledge_categories_use_direct_readback_when_pool_is_busy(monkeypatch) -> None:
service = KnowledgeService.__new__(KnowledgeService)
async def pool_exhausted_categories():
raise RuntimeError("db pool exhausted")
async def direct_categories():
return [
CategoryCount(category="alert_handling", count=113),
CategoryCount(category="AI自動化/Ansible受控修復", count=87),
]
monkeypatch.setattr(service, "_read_primary_categories", pool_exhausted_categories)
monkeypatch.setattr(service, "_read_categories_direct", direct_categories)
categories = await service.get_categories()
assert [(row.category, row.count) for row in categories] == [
("alert_handling", 113),
("AI自動化/Ansible受控修復", 87),
]
@pytest.mark.asyncio
async def test_knowledge_asset_taxonomy_fails_soft_when_readback_breaks(monkeypatch) -> None:
monkeypatch.setattr(
knowledge_service_module,
"get_db_context",
lambda: _BrokenDbContext(),
)
service = KnowledgeService.__new__(KnowledgeService)
taxonomy = await service.get_asset_taxonomy()
assert [row.key for row in taxonomy] == [
"project",
"product",
"website",
"service",
"package",
"tool",
"log",
"alert",
"playbook",
"rag",
"mcp",
"schedule",
]
assert all(row.count >= 1 for row in taxonomy)
@pytest.mark.asyncio
async def test_knowledge_asset_taxonomy_uses_direct_readback_when_pool_is_busy(monkeypatch) -> None:
service = KnowledgeService.__new__(KnowledgeService)
async def pool_exhausted_taxonomy():
raise RuntimeError("db pool exhausted")
async def direct_taxonomy():
return [
KnowledgeAssetTaxonomyCount(key="log", count=727),
KnowledgeAssetTaxonomyCount(key="alert", count=443),
KnowledgeAssetTaxonomyCount(key="mcp", count=36),
]
monkeypatch.setattr(service, "_read_primary_asset_taxonomy", pool_exhausted_taxonomy)
monkeypatch.setattr(service, "_read_asset_taxonomy_direct", direct_taxonomy)
taxonomy = await service.get_asset_taxonomy()
assert [(row.key, row.count) for row in taxonomy] == [
("log", 727),
("alert", 443),
("mcp", 36),
]
@pytest.mark.asyncio
async def test_knowledge_side_readbacks_bound_timeouts_to_source_backed(monkeypatch) -> None:
service = KnowledgeService.__new__(KnowledgeService)
async def never_finishes(*_args):
await asyncio.Event().wait()
monkeypatch.setattr(knowledge_service_module, "_PRIMARY_KM_SIDE_READ_TIMEOUT_SECONDS", 0.01)
monkeypatch.setattr(service, "_read_primary_categories", never_finishes)
monkeypatch.setattr(service, "_read_primary_asset_taxonomy", never_finishes)
monkeypatch.setattr(service, "_search_primary_entries", never_finishes)
categories = await service.get_categories()
taxonomy = await service.get_asset_taxonomy()
entries = await service.search("Telegram", limit=10)
assert [row.category for row in categories] == [
"project",
"product",
"website",
"service",
"package",
"tool",
"log",
"alert",
"playbook",
"rag",
"mcp",
"schedule",
"general",
]
assert [row.key for row in taxonomy] == [
"project",
"product",
"website",
"service",
"package",
"tool",
"log",
"alert",
"playbook",
"rag",
"mcp",
"schedule",
]
assert {entry.id for entry in entries} >= {
"source-backed-service-telegram-alert-receipts",
"source-backed-alert-telegram-monitoring-coverage",
}
@pytest.mark.asyncio
async def test_knowledge_list_binds_explicit_canonical_project_to_all_primary_reads(
monkeypatch,
) -> None:
project_contexts: list[str | None] = []
primary_entry = KnowledgeEntry(
id="km-primary-bound-1",
title="Canonical AWOOOI KM row",
content="Persistent project-bound knowledge",
entry_type=EntryType.RUNBOOK,
category="project",
tags=["awoooi"],
source=EntrySource.AI_EXTRACTED,
status=EntryStatus.APPROVED,
)
class _BoundRepo:
def __init__(self, _db):
pass
async def list_entries(self, **_kwargs):
return [primary_entry], 1
async def get_categories(self):
return [("project", 1)]
async def get_asset_taxonomy_counts(self):
return [("project", 1)]
def bound_db_context(project_id=None):
project_contexts.append(project_id)
return _OkDbContext()
monkeypatch.setattr(
knowledge_service_module,
"get_db_context",
bound_db_context,
)
monkeypatch.setattr(
knowledge_service_module,
"KnowledgeDBRepository",
_BoundRepo,
)
service = KnowledgeService.__new__(KnowledgeService)
response = await service.list_entries(
project_id="awoooi",
limit=50,
)
assert response.items == [primary_entry]
assert response.total == 1
assert response.readback_status == "ready"
assert response.primary_readback_ready is True
assert project_contexts == ["awoooi", "awoooi", "awoooi"]
@pytest.mark.asyncio
async def test_knowledge_read_router_uses_only_canonical_awoooi_project(
monkeypatch,
) -> None:
observed_projects: list[tuple[str, str | None]] = []
class _BoundService:
async def list_entries(self, **kwargs):
observed_projects.append(("list", kwargs.get("project_id")))
return KnowledgeListResponse(items=[], total=0)
async def search(self, _query, _limit, *, project_id=None):
observed_projects.append(("search", project_id))
return []
async def get_categories(self, *, project_id=None):
observed_projects.append(("categories", project_id))
return []
async def get_asset_taxonomy(self, *, project_id=None):
observed_projects.append(("asset_taxonomy", project_id))
return []
monkeypatch.setattr(
knowledge_api,
"get_knowledge_service",
lambda: _BoundService(),
)
await knowledge_api.list_entries(
category=None,
entry_type=None,
status=None,
q=None,
limit=20,
offset=0,
)
await knowledge_api.search_entries(q="Telegram", limit=20)
await knowledge_api.get_categories()
await knowledge_api.get_asset_taxonomy()
assert observed_projects == [
("list", "awoooi"),
("search", "awoooi"),
("categories", "awoooi"),
("asset_taxonomy", "awoooi"),
]
@pytest.mark.asyncio
async def test_explicit_blank_project_fails_closed_without_inheriting_context() -> None:
service = KnowledgeService.__new__(KnowledgeService)
spoofed_tokens = set_project_context(
"spoofed-project",
source="test.knowledge-spoof",
)
try:
with pytest.raises(HTTPException) as exc_info:
async with knowledge_service_module._knowledge_db_context(" "):
pass
direct_response = await service._list_entries_from_direct(
project_id=" ",
category=None,
entry_type=None,
status=None,
tags=None,
q=None,
limit=20,
offset=0,
)
finally:
clear_project_context(spoofed_tokens)
assert exc_info.value.status_code == 401
assert direct_response is None
absent_tokens = set_project_context(None, source="test.knowledge-absent")
try:
absent_response = await service._list_entries_from_direct(
category=None,
entry_type=None,
status=None,
tags=None,
q=None,
limit=20,
offset=0,
)
finally:
clear_project_context(absent_tokens)
assert absent_response is None
def test_knowledge_http_query_cannot_override_canonical_project(monkeypatch) -> None:
observed_projects: list[str | None] = []
class _BoundService:
async def list_entries(self, **kwargs):
observed_projects.append(kwargs.get("project_id"))
return KnowledgeListResponse(items=[], total=0)
monkeypatch.setattr(
knowledge_api,
"get_knowledge_service",
lambda: _BoundService(),
)
app = FastAPI()
app.include_router(knowledge_api.router, prefix="/api/v1")
with TestClient(app) as client:
response = client.get(
"/api/v1/knowledge",
params={"project_id": "evil-project"},
)
assert response.status_code == 200
assert observed_projects == ["awoooi"]