from __future__ import annotations import os import uuid import pytest from sqlalchemy.engine import make_url pytestmark = pytest.mark.integration class DefinitionRepository: def __init__(self, definition): self.definition = definition def get(self, uid): return self.definition if str(uid) == self.definition.uid else None class CredentialRepository: def __init__(self, credential): self.credential = credential def get_active(self, _session, _uid, _version): return self.credential @pytest.mark.parametrize( ("database_type", "source_environment"), [ ("postgresql", "TEST_SOURCE_POSTGRES_URL"), ("mysql", "TEST_SOURCE_MYSQL_URL"), ], ) def test_real_database_catalog_collection_persists_attempt_and_evidence( monkeypatch, database_type, source_environment, ): platform_url = os.environ.get("TEST_DATABASE_URL") source_url = os.environ.get(source_environment) if not platform_url or not source_url: pytest.skip( f"TEST_DATABASE_URL and {source_environment} are required" ) monkeypatch.setenv("DATABASE_URL", platform_url) from app import create_app, db from app.config.config import datasource_pool_settings from app.core.data_research.catalog.execution import ( CatalogIngestionExecutor, ) from app.core.data_research.catalog.models import CatalogScope from app.core.data_research.catalog.mysql import MySqlCatalogCollector from app.core.data_research.catalog.postgresql import ( PostgreSqlCatalogCollector, ) from app.core.data_research.catalog.service import CatalogCollectionService from app.core.data_research.ingestion import IngestionService from app.core.data_research.repository import ( SqlAlchemyCatalogSnapshotRepository, SqlAlchemyIngestionJobRepository, ) from app.core.data_source.adapters import adapter_for from app.core.data_source.manager import DataSourceConnectionManager from app.core.data_source.models import ( DataSourceCredential, DataSourceDefinition, ) from app.core.data_source.pool_registry import PoolRegistry from app.models.data_research import ( CatalogSnapshot, EvidenceFragment, IngestionJob, IngestionSource, ) parsed = make_url(source_url) source_uid = str(uuid.uuid4()) definition = DataSourceDefinition( uid=source_uid, name_en=f"acceptance-{database_type}", database_type=database_type, host=parsed.host, port=parsed.port, database=parsed.database, schema="public", credential_ref=source_uid, credential_version=1, ) credential = DataSourceCredential( username=parsed.username, password=parsed.password, ) settings = datasource_pool_settings() manager = DataSourceConnectionManager( definitions=DefinitionRepository(definition), credentials=CredentialRepository(credential), platform_session=lambda: object(), adapter_resolver=adapter_for, registry=PoolRegistry({**settings, "drain_timeout": 30}), settings_resolver=datasource_pool_settings, ) app = create_app() app.config.update(TESTING=True) job_uid = None try: with app.app_context(): db.session.add( IngestionSource( uid=source_uid, source_type="database", name=f"验收 {database_type}", config={ "database_type": database_type, "database": parsed.database, "schema": "public", }, permission_scope={}, status="active", created_by="integration-test", ) ) db.session.commit() ingestion = IngestionService( SqlAlchemyIngestionJobRepository(db.session), commit=db.session.commit, rollback=db.session.rollback, ) job, _created = ingestion.create_job( { "source_uid": source_uid, "job_type": "catalog_collect", "parser_version": "catalog-v1", "parameters": { "include_schemas": ( ["public"] if database_type == "postgresql" else [] ), "exclude_schemas": [], "include_tables": ["acceptance_customers"], "exclude_tables": [], }, "force_rerun": True, }, actor_uid="integration-test", ) job_uid = job.uid collection = CatalogCollectionService( manager, definition_resolver=lambda _uid: definition, collector_resolver=lambda _type: ( PostgreSqlCatalogCollector() if database_type == "postgresql" else MySqlCatalogCollector() ), ) executor = CatalogIngestionExecutor( ingestion, collection, SqlAlchemyCatalogSnapshotRepository(db.session), commit=db.session.commit, rollback=db.session.rollback, ) completed = executor.execute(job.uid) assert completed.status == "awaiting_review" assert completed.attempt_count == 1 assert completed.statistics["asset_count"] == 1 assert completed.statistics["field_count"] == 2 snapshot = db.session.query(CatalogSnapshot).filter_by( job_uid=job.uid ).one() assert snapshot.attempt == 1 assert snapshot.snapshot["assets"][0]["name"] == ( "acceptance_customers" ) evidence = db.session.query(EvidenceFragment).filter_by( job_uid=job.uid ).all() assert { item.locator["column"] for item in evidence } == {"id", "customer_name"} assert all( "password" not in str(item.locator).lower() for item in evidence ) finally: manager.close() with app.app_context(): if job_uid: db.session.query(IngestionJob).filter_by( uid=job_uid ).delete(synchronize_session=False) db.session.query(IngestionSource).filter_by( uid=source_uid ).delete(synchronize_session=False) db.session.commit()