Develop #11

Merged
sakibcc merged 4 commits from develop into main 2026-08-11 15:08:13 +08:00
Showing only changes of commit 8f5f504368 - Show all commits
@@ -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",
]