Files
LegacyHUB/app/api/routes_ingestion.py
Vadim Malanov d27dd0ffbb feat: align LegacyHUB with TeamHUB platform contract (D2/D4, assets, security)
Close 12 audit-driven platform-compliance gaps on a single branch.

- D4 dispatch: app/integrations/dispatch_client.py participant `legacyhub`,
  emits LegacyhubDocumentIndexed + AssetDerivativeReady after the indexing
  commit (idempotent uuid5), http_inbox route (reindex/tombstone) with
  audit-based dedupe; docs/dispatch-contract.md. Celery+Redis stays intra-module.
- D2 SSO: app/integrations/identity.py validates X-TeamHub-* + role/scope
  mapper; security.py adds trusted-header enforcement (AUTH_REQUIRE_IDENTITY)
  and a scope check on /search; docker-compose.teamhub.yml (external teamhub_net
  + internal db net, api not host-published); RUNBOOK network/firewall section.
- Asset standard: SearchHit/Citation carry asset_id/owner_module; buckets
  renamed teamhub-legacyhub-* (+quarantine/tmp/exports); purge-by-asset_id with
  legal-hold guard (app/indexing/projection.py); OCR-markdown derivative event.
- audit_log model + Alembic 0003 + record_audit on writes (same transaction).
- Secret masking: app/common/json_logger.py recursive mask wired into structlog
  (+ensure_ascii=False); event payloads redacted before persistence.
- Service X-API-Key mandatory on ingest endpoints (defence-in-depth).
- Port: host API 8000->8050 (collision with SalesHUB/MailHUB resolved),
  container still listens on 8000.
- Config: no plaintext secret defaults; fail-loud in non-dev (no value leak).
- Docs drift: README PG 5440, layered-auth note, 5173 removed from CORS;
  ingest/folder gated by ENABLE_FOLDER_INGEST (410 by default).
- ADRs: layers mapping, shared-core extraction, UI locale (RU-first).

Tests: 78 passing (ruff, compileall, pytest, tsc, vite build, compose config).

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-15 11:44:15 +03:00

123 lines
3.9 KiB
Python

"""Ingestion endpoints."""
from __future__ import annotations
import uuid
from pathlib import Path
from fastapi import APIRouter, Depends, Header, HTTPException, Response, status
from app.api.schemas import (
AssetManifestEnvelope,
IngestFolderRequest,
IngestFolderResponse,
KnowledgeIngestResponse,
)
from app.api.security import require_service_api_key
from app.config import settings
from app.db.audit import record_audit
from app.db.session import session_scope
from app.ingestion.knowledge_ingest import KnowledgeIngestError, accept_knowledge_ingest
from app.logging_config import get_logger
logger = get_logger(__name__)
router = APIRouter(tags=["ingestion"])
@router.post(
"/knowledge-ingest",
response_model=KnowledgeIngestResponse,
dependencies=[Depends(require_service_api_key)],
)
def knowledge_ingest(
req: AssetManifestEnvelope,
response: Response,
idempotency_key: str | None = Header(default=None, alias="Idempotency-Key"),
) -> KnowledgeIngestResponse:
"""Accept a TeamHUB AssetManifest and queue LegacyHUB projection ingest."""
try:
result = accept_knowledge_ingest(req, idempotency_key=idempotency_key)
except KnowledgeIngestError as exc:
raise HTTPException(
status_code=exc.status_code,
detail={"reason_code": exc.reason_code, "message": exc.message},
) from exc
response.status_code = (
status.HTTP_202_ACCEPTED if result.status == "accepted" else status.HTTP_200_OK
)
return result
@router.post(
"/ingest/folder",
response_model=IngestFolderResponse,
deprecated=True,
dependencies=[Depends(require_service_api_key)],
)
def ingest_folder(req: IngestFolderRequest, response: Response) -> IngestFolderResponse:
"""Discover all PDFs under ``path`` and queue them for processing.
The request returns immediately after the discovery pass. Per-document
OCR / extraction / indexing happens asynchronously in Celery workers.
Deprecated for TeamHUB integrations: use ``POST /api/v1/knowledge-ingest``
with an AssetManifest and object-storage reference instead.
"""
response.headers["Deprecation"] = "true"
response.headers["Link"] = '</api/v1/knowledge-ingest>; rel="successor-version"'
if not settings.enable_folder_ingest:
raise HTTPException(
status_code=status.HTTP_410_GONE,
detail="folder ingest is disabled; use POST /api/v1/knowledge-ingest",
)
folder = Path(req.path)
if not folder.exists() or not folder.is_dir():
raise HTTPException(status_code=400, detail=f"Folder not found: {req.path}")
# Lazy import - keeps module load light.
from app.ingestion.scanner import discover_documents # noqa: PLC0415
from app.workers.tasks import process_document # noqa: PLC0415
run_id = uuid.uuid4()
discovered, queued, dups, invalid = 0, 0, 0, 0
for record in discover_documents(folder, recursive=req.recursive, force=req.force):
discovered += 1
if record.duplicate and not req.force:
dups += 1
continue
if not record.document_id:
invalid += 1
continue
process_document.delay(str(record.document_id), str(run_id))
queued += 1
logger.info(
"ingest.folder.queued",
path=str(folder),
discovered=discovered,
queued=queued,
skipped_duplicates=dups,
invalid=invalid,
run_id=str(run_id),
)
with session_scope() as db:
record_audit(
db,
action="ingest.folder.queued",
actor="service:ingest-folder",
entity_type="ingestion_run",
entity_id=str(run_id),
details={"discovered": discovered, "queued": queued, "skipped_duplicates": dups},
)
return IngestFolderResponse(
run_id=run_id,
discovered=discovered,
queued=queued,
skipped_duplicates=dups,
invalid_files=invalid,
)