| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202 |
- 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()
|