diff --git a/zhixing-server/src/zhixing_server/modules/market_data/presentation/integrity.py b/zhixing-server/src/zhixing_server/modules/market_data/presentation/integrity.py new file mode 100644 index 0000000..ad48ae7 --- /dev/null +++ b/zhixing-server/src/zhixing_server/modules/market_data/presentation/integrity.py @@ -0,0 +1,217 @@ +"""HTTP presentation for user-triggered market-data integrity checks.""" + +from datetime import date, datetime +from typing import Annotated, Literal + +from fastapi import APIRouter, BackgroundTasks, Depends, HTTPException, Query, status +from pydantic import BaseModel, Field + +from zhixing_server.bootstrap.config import Settings, get_settings +from zhixing_server.modules.market_data.application.integrity import RunMarketIntegrityCheck +from zhixing_server.modules.market_data.domain.integrity import ( + IntegrityCheckInProgress, + IntegrityCheckNoData, + IntegrityCheckPage, + IntegrityCheckQuery, + IntegrityCheckStoreError, + IntegrityIssue, +) +from zhixing_server.modules.market_data.domain.ports import MarketDataRepositoryError +from zhixing_server.modules.market_data.infrastructure.csv_snapshot import CsvSnapshotStore +from zhixing_server.modules.market_data.infrastructure.integrity import ( + PostgresIntegrityCheckStore, + PostgresIntegritySnapshotReader, +) +from zhixing_server.modules.market_data.presentation.home import get_market_data_repository + +integrity_router = APIRouter() + +IntegrityStatusValue = Literal["no_data", "running", "passed", "issues_found", "failed"] + + +class IntegrityCheckAcceptedResponse(BaseModel): + """Small 202 response returned after the running slot is claimed.""" + + check_id: str + status: Literal["running"] + window_start: date + window_end: date + + +class IntegrityIssueResponse(BaseModel): + """One safe issue in a paginated integrity report.""" + + issue_key: str + item_kind: str + item_key: str + issue_type: str + message: str + + @classmethod + def from_domain(cls, issue: IntegrityIssue) -> "IntegrityIssueResponse": + """Translate a domain issue without exposing storage timestamps.""" + + return cls( + issue_key=issue.issue_key, + item_kind=issue.item_kind, + item_key=issue.item_key, + issue_type=issue.issue_type, + message=issue.message, + ) + + +class IntegrityCheckResponse(BaseModel): + """Progress, terminal state, and one bounded issue page.""" + + check_id: str | None + status: IntegrityStatusValue + window_start: date | None = None + window_end: date | None = None + target_count: int = Field(default=0, ge=0) + checked_count: int = Field(default=0, ge=0) + issue_count: int = Field(default=0, ge=0) + error_type: str | None = None + error_message: str | None = None + created_at: datetime | None = None + finished_at: datetime | None = None + page: int = Field(default=1, ge=1) + page_size: int = Field(default=10, ge=1, le=100) + issues_total: int = Field(default=0, ge=0) + issues: list[IntegrityIssueResponse] = Field(default_factory=lambda: []) + + +def get_market_integrity_service( + settings: Annotated[Settings, Depends(get_settings)], +) -> RunMarketIntegrityCheck: + """Compose the check from local PostgreSQL and formal CSV adapters only.""" + + repository = get_market_data_repository(settings) + return RunMarketIntegrityCheck( + PostgresIntegritySnapshotReader(repository), + CsvSnapshotStore(settings.market_data_csv_root), + PostgresIntegrityCheckStore(repository), + lock_key=settings.market_data_advisory_lock_key, + ) + + +@integrity_router.post( + "", + response_model=IntegrityCheckAcceptedResponse, + status_code=status.HTTP_202_ACCEPTED, +) +def trigger_integrity_check( + background_tasks: BackgroundTasks, + service: Annotated[RunMarketIntegrityCheck, Depends(get_market_integrity_service)], +) -> IntegrityCheckAcceptedResponse: + """Claim one report and schedule its read-only comparison.""" + + try: + run = service.prepare() + except IntegrityCheckInProgress as exc: + raise _http_error(409, "integrity_check_in_progress", str(exc)) from exc + except IntegrityCheckNoData as exc: + raise _http_error(422, "integrity_check_no_data", str(exc)) from exc + except (IntegrityCheckStoreError, MarketDataRepositoryError) as exc: + raise _http_error(503, "integrity_storage_unavailable", str(exc)) from exc + if run.window is None: + raise _http_error(503, "integrity_storage_unavailable", "integrity check window is missing") + background_tasks.add_task(service.execute, run.id) + return IntegrityCheckAcceptedResponse( + check_id=run.id, + status="running", + window_start=run.window.start, + window_end=run.window.end, + ) + + +@integrity_router.get("/latest", response_model=IntegrityCheckResponse) +def get_latest_integrity_check( + service: Annotated[RunMarketIntegrityCheck, Depends(get_market_integrity_service)], + page: Annotated[int, Query(ge=1)] = 1, + page_size: Annotated[int, Query(ge=1, le=100)] = 10, +) -> IntegrityCheckResponse: + """Return the latest report or an explicit ``no_data`` result.""" + + query = IntegrityCheckQuery(page=page, page_size=page_size) + try: + result = service.get_latest(query=query) + except (IntegrityCheckStoreError, MarketDataRepositoryError) as exc: + raise _http_error(503, "integrity_storage_unavailable", str(exc)) from exc + if result is None: + return _no_data_response(query) + return _page_response(result) + + +@integrity_router.get("/{check_id}", response_model=IntegrityCheckResponse) +def get_integrity_check( + check_id: str, + service: Annotated[RunMarketIntegrityCheck, Depends(get_market_integrity_service)], + page: Annotated[int, Query(ge=1)] = 1, + page_size: Annotated[int, Query(ge=1, le=100)] = 10, +) -> IntegrityCheckResponse: + """Return one report for polling and issue pagination.""" + + query = IntegrityCheckQuery(page=page, page_size=page_size) + try: + result = service.get(check_id, query=query) + except (IntegrityCheckStoreError, MarketDataRepositoryError) as exc: + raise _http_error(503, "integrity_storage_unavailable", str(exc)) from exc + if result is None: + raise _http_error( + 404, + "integrity_check_not_found", + f"integrity check not found: {check_id}", + ) + return _page_response(result) + + +def _page_response(result: IntegrityCheckPage) -> IntegrityCheckResponse: + """Translate the application page into the stable HTTP shape.""" + + run = result.run + return IntegrityCheckResponse( + check_id=run.id, + status=run.status, + window_start=run.window_start, + window_end=run.window_end, + target_count=run.target_count, + checked_count=run.checked_count, + issue_count=run.issue_count, + error_type=run.error_type, + error_message=run.error_message, + created_at=run.created_at, + finished_at=run.finished_at, + page=result.page, + page_size=result.page_size, + issues_total=result.issues_total, + issues=[IntegrityIssueResponse.from_domain(issue) for issue in result.issues], + ) + + +def _no_data_response(query: IntegrityCheckQuery) -> IntegrityCheckResponse: + """Return an explicit empty latest response without fabricating a run.""" + + return IntegrityCheckResponse( + check_id=None, + status="no_data", + page=query.page, + page_size=query.page_size, + ) + + +def _http_error(code: int, error_type: str, message: str) -> HTTPException: + """Use the repository's established explicit error envelope.""" + + return HTTPException( + status_code=code, + detail={"code": error_type, "message": " ".join(message.split())[:500]}, + ) + + +__all__ = [ + "IntegrityCheckAcceptedResponse", + "IntegrityCheckResponse", + "IntegrityIssueResponse", + "get_market_integrity_service", + "integrity_router", +]