Files
LegacyHUB/app/storage/minio_client.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

120 lines
3.7 KiB
Python

"""Thin wrapper around the MinIO Python SDK with bucket bootstrap and retries."""
from __future__ import annotations
import io
from functools import lru_cache
from pathlib import Path
from typing import Any
from minio import Minio
from minio.error import S3Error
from tenacity import retry, retry_if_exception_type, stop_after_attempt, wait_exponential
from app.config import settings
from app.logging_config import get_logger
logger = get_logger(__name__)
class MinioStorage:
def __init__(self, client: Minio | None = None) -> None:
self.client = client or Minio(
endpoint=settings.minio_endpoint,
access_key=settings.minio_access_key,
secret_key=settings.minio_secret_key,
secure=settings.minio_secure,
region=settings.minio_region,
)
self.originals_bucket = settings.minio_bucket_originals
self.derived_bucket = settings.minio_bucket_derived
self.quarantine_bucket = settings.minio_bucket_quarantine
self.tmp_bucket = settings.minio_bucket_tmp
self.exports_bucket = settings.minio_bucket_exports
def ensure_buckets(self) -> None:
for bucket in (
self.originals_bucket,
self.derived_bucket,
self.quarantine_bucket,
self.tmp_bucket,
self.exports_bucket,
):
if not self.client.bucket_exists(bucket):
logger.info("minio.create_bucket", bucket=bucket)
self.client.make_bucket(bucket)
@retry(
stop=stop_after_attempt(3),
wait=wait_exponential(multiplier=1, min=1, max=10),
retry=retry_if_exception_type(S3Error),
reraise=True,
)
def put_file(
self,
bucket: str,
key: str,
path: Path,
content_type: str = "application/octet-stream",
metadata: dict[str, str] | None = None,
) -> None:
size = path.stat().st_size
with path.open("rb") as f:
self.client.put_object(
bucket_name=bucket,
object_name=key,
data=f,
length=size,
content_type=content_type,
metadata=metadata or {},
)
@retry(
stop=stop_after_attempt(3),
wait=wait_exponential(multiplier=1, min=1, max=10),
retry=retry_if_exception_type(S3Error),
reraise=True,
)
def put_bytes(
self,
bucket: str,
key: str,
data: bytes,
content_type: str = "application/octet-stream",
metadata: dict[str, str] | None = None,
) -> None:
self.client.put_object(
bucket_name=bucket,
object_name=key,
data=io.BytesIO(data),
length=len(data),
content_type=content_type,
metadata=metadata or {},
)
def get_to_path(self, bucket: str, key: str, dest: Path, version_id: str | None = None) -> Path:
dest.parent.mkdir(parents=True, exist_ok=True)
self.client.fget_object(bucket, key, str(dest), version_id=version_id)
return dest
def exists(self, bucket: str, key: str) -> bool:
try:
self.client.stat_object(bucket, key)
return True
except S3Error as exc:
if exc.code in {"NoSuchKey", "NoSuchObject"}:
return False
raise
def health(self) -> dict[str, Any]:
try:
buckets = [b.name for b in self.client.list_buckets()]
return {"status": "ok", "buckets": buckets}
except Exception as exc:
return {"status": "error", "error": str(exc)}
@lru_cache(maxsize=1)
def get_storage() -> MinioStorage:
return MinioStorage()