Просмотр исходного кода

fix: harden Docker integration acceptance

马小龙 1 месяц назад
Родитель
Сommit
77b72a0166

+ 3 - 0
app/core/data_research/ontology/repository.py

@@ -159,6 +159,9 @@ class SqlAlchemyOntologyRepository:
             created_by=created_by,
             created_by=created_by,
         )
         )
         self.session.add(model)
         self.session.add(model)
+        # The links are written with explicit foreign-key values rather than an
+        # ORM relationship, so persist the parent before flushing link rows.
+        self.session.flush()
         for link in domain_links:
         for link in domain_links:
             self.session.add(
             self.session.add(
                 OntologyDomainLinkModel(
                 OntologyDomainLinkModel(

+ 3 - 0
deployment/app/core/data_research/ontology/repository.py

@@ -159,6 +159,9 @@ class SqlAlchemyOntologyRepository:
             created_by=created_by,
             created_by=created_by,
         )
         )
         self.session.add(model)
         self.session.add(model)
+        # The links are written with explicit foreign-key values rather than an
+        # ORM relationship, so persist the parent before flushing link rows.
+        self.session.flush()
         for link in domain_links:
         for link in domain_links:
             self.session.add(
             self.session.add(
                 OntologyDomainLinkModel(
                 OntologyDomainLinkModel(

+ 14 - 3
deployment/migrations/env.py

@@ -22,6 +22,7 @@ if not database_url:
 # ConfigParser treats percent signs as interpolation markers.
 # ConfigParser treats percent signs as interpolation markers.
 config.set_main_option("sqlalchemy.url", database_url.replace("%", "%%"))
 config.set_main_option("sqlalchemy.url", database_url.replace("%", "%%"))
 target_metadata = None
 target_metadata = None
+MIGRATION_ADVISORY_LOCK_ID = 2026072201
 
 
 
 
 def run_migrations_offline() -> None:
 def run_migrations_offline() -> None:
@@ -42,9 +43,19 @@ def run_migrations_online() -> None:
         poolclass=pool.NullPool,
         poolclass=pool.NullPool,
     )
     )
     with connectable.connect() as connection:
     with connectable.connect() as connection:
-        context.configure(connection=connection, target_metadata=target_metadata)
-        with context.begin_transaction():
-            context.run_migrations()
+        connection.exec_driver_sql(
+            "SELECT pg_advisory_lock(%s)", (MIGRATION_ADVISORY_LOCK_ID,)
+        )
+        connection.commit()
+        try:
+            context.configure(connection=connection, target_metadata=target_metadata)
+            with context.begin_transaction():
+                context.run_migrations()
+        finally:
+            connection.exec_driver_sql(
+                "SELECT pg_advisory_unlock(%s)", (MIGRATION_ADVISORY_LOCK_ID,)
+            )
+            connection.commit()
 
 
 
 
 if context.is_offline_mode():
 if context.is_offline_mode():

+ 14 - 10
docs/validation/data-research-v60-v65-acceptance.md

@@ -4,7 +4,7 @@
 
 
 ## 1. 验收结论
 ## 1. 验收结论
 
 
-V60-V65 的非外部依赖开发项全部完成:多源采集、证据治理、数据元素生命周期、本体 MVP、动态本体、交换/知识/MCP 服务、运营对账、Vue 2 完整工作流及发布副本均已交付。最终状态为“代码验收通过;真实外部服务集成待可用 Docker/发布环境补跑”
+V60-V65 全部完成:多源采集、证据治理、数据元素生命周期、本体 MVP、动态本体、交换/知识/MCP 服务、运营对账、Vue 2 完整工作流及发布副本均已交付。本地 Docker 隔离栈已启动,真实 PostgreSQL/MySQL、Neo4j、MinIO、Kestra、n8n、Runner、后端和前端集成验收通过
 
 
 ## 2. 阶段交付
 ## 2. 阶段交付
 
 
@@ -23,13 +23,16 @@ V60-V65 的非外部依赖开发项全部完成:多源采集、证据治理、
 |---|---|
 |---|---|
 | `PYTHONPATH=. .venv/bin/pytest -q tests/data_research tests/acceptance/test_data_research_v60_v65.py` | `79 passed` |
 | `PYTHONPATH=. .venv/bin/pytest -q tests/data_research tests/acceptance/test_data_research_v60_v65.py` | `79 passed` |
 | `PYTHONPATH=. .venv/bin/pytest -q` | `372 passed, 23 skipped, 59 subtests passed` |
 | `PYTHONPATH=. .venv/bin/pytest -q` | `372 passed, 23 skipped, 59 subtests passed` |
+| Docker 真实服务集成矩阵(含数据源故障恢复、Runner、Kestra、MCP、RBAC、迁移) | `62 passed, 1 skipped` |
+| 数据研发容器 HTTP 端到端(CSV/MinIO/任务/双业务域本体/Outbox) | `1 passed` |
+| 启用全部本地集成开关后的全仓回归 | `396 passed, 1 skipped, 59 subtests passed` |
 | `PATH=<bundled-node-24>:$PATH npm --prefix frontend run build` | 构建成功;0 error,21 条既有 `no-console` 警告,另有依赖年龄/CSS 顺序/包体积警告 |
 | `PATH=<bundled-node-24>:$PATH npm --prefix frontend run build` | 构建成功;0 error,21 条既有 `no-console` 警告,另有依赖年龄/CSS 顺序/包体积警告 |
 | `python scripts/generate_openapi.py --output docs/architecture/OPENAPI.yaml` | 138 个操作 |
 | `python scripts/generate_openapi.py --output docs/architecture/OPENAPI.yaml` | 138 个操作 |
 | `PYTHONPATH=. .venv/bin/pytest -q tests/test_architecture_artifacts.py` | 通过 |
 | `PYTHONPATH=. .venv/bin/pytest -q tests/test_architecture_artifacts.py` | 通过 |
 | `diff -qr app deployment/app` | 无差异 |
 | `diff -qr app deployment/app` | 无差异 |
 | `git diff --check` | 通过 |
 | `git diff --check` | 通过 |
 
 
-验收期间全量测试曾捕获任务令牌尾字符的非规范 Base64URL 别名问题;验证器现强制 JWT 三段采用规范编码,专项测试与全量回归均已通过。
+验收期间全量测试曾捕获任务令牌尾字符的非规范 Base64URL 别名问题;验证器现强制 JWT 三段采用规范编码,专项测试与全量回归均已通过。Docker 集成进一步捕获并修复了后端/Runner 并发 Alembic 迁移竞争,以及本体主记录与业务域关联记录的写入顺序问题;两项均已增加回归测试。
 
 
 ## 4. 安全与治理断言
 ## 4. 安全与治理断言
 
 
@@ -41,14 +44,15 @@ V60-V65 的非外部依赖开发项全部完成:多源采集、证据治理、
 - AI 建议必须携带证据、置信度、模型/提示版本;未决建议不进入草稿或发布。
 - AI 建议必须携带证据、置信度、模型/提示版本;未决建议不进入草稿或发布。
 - 语义查询使用固定深度与固定 Cypher,只返回已发布、业务域范围内的白名单字段。
 - 语义查询使用固定深度与固定 Cypher,只返回已发布、业务域范围内的白名单字段。
 
 
-## 5. 环境受限
+## 5. Docker 集成范围与剩余人工
 
 
-`docker compose -f deploy/docker/docker-compose.yml ps` 返回 Docker daemon 未运行。因此以下真实服务测试在本机保持跳过,不用单元假件结果冒充生产集成结果
+已验证
 
 
-- Alembic 在临时 PostgreSQL 数据库的重复升级/降级与实际表检查。
-- PostgreSQL/MySQL 目录采集及查询超时集成。
-- PostgreSQL Outbox 与 Neo4j 投影重试/幂等的真实跨存储验证。
-- MinIO 原件上传、哈希复用、缺失工件对账。
-- MCP stdio/canary 依赖的外部运行环境及真实服务权限验证。
+- Alembic 在临时 PostgreSQL 数据库中的升级、重复升级、降级,以及两个进程并发升级互斥。
+- PostgreSQL/MySQL 外部数据源并发连接池、凭据轮换、MySQL 停机隔离和恢复。
+- Runner 真实读取两类受治理数据源、任务令牌防重放、Kestra 到 Runner 调用链。
+- Kestra 流程验证/部署、官方 MCP stdio、DataOps MCP 持久化与 V53 L3 链路。
+- 数据研发 CSV 上传到 MinIO、哈希复用、采集任务幂等/取消、双业务域本体创建/校验/发布/导出及 Outbox 事件。
+- RBAC、V54 离线模型降级、跨存储 Outbox 失败重试和幂等消费。
 
 
-这些检查不阻断代码验收,但在生产发布前必须在 Docker 或等价预发布环境补跑,并将结果附加到本报告
+唯一自动跳过项为 V55 的 n8n 真实切换/回滚演练。它要求先在本地 n8n 完成 owner 初始化并生成 `N8N_API_KEY`,属于必须人工提供凭据的破坏性切换演练;当前未伪造该凭据。Docker 测试环境保持运行,完成 n8n 初始化后可单独补跑该用例

+ 14 - 3
migrations/env.py

@@ -22,6 +22,7 @@ if not database_url:
 # ConfigParser treats percent signs as interpolation markers.
 # ConfigParser treats percent signs as interpolation markers.
 config.set_main_option("sqlalchemy.url", database_url.replace("%", "%%"))
 config.set_main_option("sqlalchemy.url", database_url.replace("%", "%%"))
 target_metadata = None
 target_metadata = None
+MIGRATION_ADVISORY_LOCK_ID = 2026072201
 
 
 
 
 def run_migrations_offline() -> None:
 def run_migrations_offline() -> None:
@@ -42,9 +43,19 @@ def run_migrations_online() -> None:
         poolclass=pool.NullPool,
         poolclass=pool.NullPool,
     )
     )
     with connectable.connect() as connection:
     with connectable.connect() as connection:
-        context.configure(connection=connection, target_metadata=target_metadata)
-        with context.begin_transaction():
-            context.run_migrations()
+        connection.exec_driver_sql(
+            "SELECT pg_advisory_lock(%s)", (MIGRATION_ADVISORY_LOCK_ID,)
+        )
+        connection.commit()
+        try:
+            context.configure(connection=connection, target_metadata=target_metadata)
+            with context.begin_transaction():
+                context.run_migrations()
+        finally:
+            connection.exec_driver_sql(
+                "SELECT pg_advisory_unlock(%s)", (MIGRATION_ADVISORY_LOCK_ID,)
+            )
+            connection.commit()
 
 
 
 
 if context.is_offline_mode():
 if context.is_offline_mode():

+ 11 - 2
tests/integration/test_cross_store_consistency.py

@@ -84,6 +84,14 @@ def test_downstream_failure_recovers_and_restart_does_not_duplicate(database_url
                 event_type="governance.updated",
                 event_type="governance.updated",
                 payload={"marker": marker},
                 payload={"marker": marker},
             )
             )
+            session.execute(
+                text(
+                    "UPDATE public.outbox_events "
+                    "SET available_at = TIMESTAMPTZ '1970-01-01' "
+                    "WHERE event_id = CAST(:event_id AS uuid)"
+                ),
+                {"event_id": event_id},
+            )
             session.commit()
             session.commit()
 
 
         with Session(engine) as session:
         with Session(engine) as session:
@@ -108,7 +116,7 @@ def test_downstream_failure_recovers_and_restart_does_not_duplicate(database_url
             session.execute(
             session.execute(
                 text(
                 text(
                     "UPDATE public.outbox_events "
                     "UPDATE public.outbox_events "
-                    "SET available_at = CURRENT_TIMESTAMP "
+                    "SET available_at = TIMESTAMPTZ '1970-01-01' "
                     "WHERE event_id = CAST(:event_id AS uuid)"
                     "WHERE event_id = CAST(:event_id AS uuid)"
                 ),
                 ),
                 {"event_id": event_id},
                 {"event_id": event_id},
@@ -131,7 +139,8 @@ def test_downstream_failure_recovers_and_restart_does_not_duplicate(database_url
             session.execute(
             session.execute(
                 text(
                 text(
                     "UPDATE public.outbox_events "
                     "UPDATE public.outbox_events "
-                    "SET status = 'pending', available_at = CURRENT_TIMESTAMP "
+                    "SET status = 'pending', "
+                    "available_at = TIMESTAMPTZ '1970-01-01' "
                     "WHERE event_id = CAST(:event_id AS uuid)"
                     "WHERE event_id = CAST(:event_id AS uuid)"
                 ),
                 ),
                 {"event_id": event_id},
                 {"event_id": event_id},

+ 282 - 0
tests/integration/test_data_research_docker_api.py

@@ -0,0 +1,282 @@
+"""Opt-in end-to-end acceptance for the containerized data-research API."""
+
+from __future__ import annotations
+
+import os
+import uuid
+
+import pytest
+import requests
+from minio import Minio
+from sqlalchemy import create_engine, text
+
+
+pytestmark = pytest.mark.skipif(
+    os.getenv("RUN_DATA_RESEARCH_DOCKER_API") != "1",
+    reason="set RUN_DATA_RESEARCH_DOCKER_API=1 for Docker API acceptance",
+)
+
+
+def _data(response, expected_status):
+    assert response.status_code == expected_status, response.text
+    payload = response.json()
+    assert payload["code"] == 200, payload
+    return payload["data"]
+
+
+def test_file_ingestion_and_dynamic_ontology_publish_through_live_api():
+    base_url = os.getenv("TEST_BACKEND_URL", "http://127.0.0.1:15500/api")
+    database_url = os.environ["TEST_DATABASE_URL"]
+    engine = create_engine(database_url, pool_pre_ping=True)
+    minio = Minio(
+        os.getenv("TEST_MINIO_HOST", "127.0.0.1:19000"),
+        access_key="dataops-test",
+        secret_key="dataops-test-password",
+        secure=False,
+    )
+    source_uid = str(uuid.uuid4())
+    owner_domain_uid = str(uuid.uuid4())
+    contributor_domain_uid = str(uuid.uuid4())
+    ontology_uid = None
+    artifact_uid = None
+    artifact_object_key = None
+
+    login = requests.post(
+        f"{base_url}/system/auth/login",
+        json={"username": "admin", "password": os.environ["TEST_ADMIN_PASSWORD"]},
+        timeout=20,
+    )
+    token = _data(login, 200)["token"]
+    headers = {"Authorization": f"Bearer {token}"}
+
+    try:
+        with engine.begin() as connection:
+            connection.execute(
+                text(
+                    "INSERT INTO public.ingestion_sources "
+                    "(uid, source_type, name, created_by) "
+                    "VALUES (CAST(:uid AS uuid), 'file', :name, 'docker-api-test')"
+                ),
+                {"uid": source_uid, "name": f"docker-api-{source_uid}"},
+            )
+
+        csv_content = b"code,name\nCUSTOMER_ID,Customer identifier\n"
+        uploaded = requests.post(
+            f"{base_url}/development/v1/sources/files",
+            headers=headers,
+            data={"source_uid": source_uid, "parser_version": "csv-v1"},
+            files={"file": ("elements.csv", csv_content, "text/csv")},
+            timeout=20,
+        )
+        artifact = _data(uploaded, 201)
+        artifact_uid = artifact["uid"]
+
+        duplicate = requests.post(
+            f"{base_url}/development/v1/sources/files",
+            headers=headers,
+            data={"source_uid": source_uid, "parser_version": "csv-v1"},
+            files={"file": ("elements.csv", csv_content, "text/csv")},
+            timeout=20,
+        )
+        assert _data(duplicate, 200)["uid"] == artifact_uid
+
+        job_payload = {
+            "source_uid": source_uid,
+            "artifact_uid": artifact_uid,
+            "job_type": "file_extract",
+            "parser_version": "csv-v1",
+            "parameters": {"delimiter": ","},
+        }
+        created_job = _data(
+            requests.post(
+                f"{base_url}/development/v1/ingestion-jobs",
+                headers=headers,
+                json=job_payload,
+                timeout=20,
+            ),
+            201,
+        )
+        duplicate_job = _data(
+            requests.post(
+                f"{base_url}/development/v1/ingestion-jobs",
+                headers=headers,
+                json=job_payload,
+                timeout=20,
+            ),
+            200,
+        )
+        assert duplicate_job["uid"] == created_job["uid"]
+        cancelled = _data(
+            requests.post(
+                f"{base_url}/development/v1/ingestion-jobs/"
+                f"{created_job['uid']}/cancel",
+                headers=headers,
+                timeout=20,
+            ),
+            200,
+        )
+        assert cancelled["status"] == "cancelled"
+
+        ontology = _data(
+            requests.post(
+                f"{base_url}/development/v1/ontologies",
+                headers=headers,
+                json={
+                    "code": f"DOCKER_{uuid.uuid4().hex[:12].upper()}",
+                    "name": "Docker API customer ontology",
+                    "domain_links": [
+                        {"domain_uid": owner_domain_uid, "role": "owner"},
+                        {
+                            "domain_uid": contributor_domain_uid,
+                            "role": "contributor",
+                        },
+                    ],
+                },
+                timeout=20,
+            ),
+            201,
+        )
+        ontology_uid = ontology["uid"]
+        graph = {
+            "classes": [{"uid": "customer", "name": "Customer"}],
+            "properties": [
+                {
+                    "uid": "customer-id",
+                    "name": "customerId",
+                    "class_uid": "customer",
+                    "required": True,
+                    "cardinality": "1",
+                    "data_element_uid": "customer-id-element",
+                }
+            ],
+            "relations": [],
+            "constraints": [],
+            "domain_links": [
+                {"domain_uid": owner_domain_uid, "role": "owner"},
+                {
+                    "domain_uid": contributor_domain_uid,
+                    "role": "contributor",
+                },
+            ],
+            "element_mappings": [
+                {
+                    "property_uid": "customer-id",
+                    "data_element_uid": "customer-id-element",
+                }
+            ],
+        }
+        draft = _data(
+            requests.patch(
+                f"{base_url}/development/v1/ontologies/{ontology_uid}/graph",
+                headers={**headers, "If-Match": '"0"'},
+                json=graph,
+                timeout=20,
+            ),
+            200,
+        )
+        assert draft["version"] == 1
+        assert _data(
+            requests.post(
+                f"{base_url}/development/v1/ontologies/{ontology_uid}/validate",
+                headers=headers,
+                timeout=20,
+            ),
+            200,
+        ) == []
+        published = _data(
+            requests.post(
+                f"{base_url}/development/v1/ontologies/{ontology_uid}/publish",
+                headers={**headers, "Idempotency-Key": f"docker-{ontology_uid}"},
+                timeout=20,
+            ),
+            200,
+        )
+        assert published["status"] == "published"
+        exported = requests.get(
+            f"{base_url}/development/v1/ontologies/{ontology_uid}/export?format=json",
+            headers=headers,
+            timeout=20,
+        )
+        assert exported.status_code == 200
+        assert exported.json()["graph_document"]["classes"][0]["uid"] == "customer"
+
+        with engine.connect() as connection:
+            artifact_object_key = connection.execute(
+                text(
+                    "SELECT storage_ref FROM public.source_artifacts "
+                    "WHERE uid = CAST(:uid AS uuid)"
+                ),
+                {"uid": artifact_uid},
+            ).scalar_one().removeprefix("minio://dataops-bucket/")
+            outbox_payload = connection.execute(
+                text(
+                    "SELECT payload FROM public.outbox_events "
+                    "WHERE aggregate_id = :uid "
+                    "AND event_type = 'ontology.version_published'"
+                ),
+                {"uid": ontology_uid},
+            ).scalar_one()
+        assert outbox_payload["ontology_uid"] == ontology_uid
+    finally:
+        with engine.begin() as connection:
+            if ontology_uid:
+                connection.execute(
+                    text(
+                        "DELETE FROM public.outbox_consumptions WHERE event_id IN "
+                        "(SELECT event_id FROM public.outbox_events "
+                        "WHERE aggregate_id = :uid)"
+                    ),
+                    {"uid": ontology_uid},
+                )
+                connection.execute(
+                    text("DELETE FROM public.outbox_events WHERE aggregate_id = :uid"),
+                    {"uid": ontology_uid},
+                )
+                connection.execute(
+                    text(
+                        "UPDATE public.ontologies SET active_version_uid = NULL "
+                        "WHERE uid = CAST(:uid AS uuid)"
+                    ),
+                    {"uid": ontology_uid},
+                )
+                for table in (
+                    "ontology_publish_runs",
+                    "ontology_change_sets",
+                    "ontology_domain_links",
+                    "ontology_versions",
+                    "ontologies",
+                ):
+                    connection.execute(
+                        text(
+                            f"DELETE FROM public.{table} "
+                            "WHERE ontology_uid = CAST(:uid AS uuid)"
+                            if table != "ontologies"
+                            else "DELETE FROM public.ontologies "
+                            "WHERE uid = CAST(:uid AS uuid)"
+                        ),
+                        {"uid": ontology_uid},
+                    )
+            connection.execute(
+                text(
+                    "DELETE FROM public.ingestion_jobs "
+                    "WHERE source_uid = CAST(:uid AS uuid)"
+                ),
+                {"uid": source_uid},
+            )
+            connection.execute(
+                text(
+                    "DELETE FROM public.source_artifacts "
+                    "WHERE source_uid = CAST(:uid AS uuid)"
+                ),
+                {"uid": source_uid},
+            )
+            connection.execute(
+                text(
+                    "DELETE FROM public.ingestion_sources "
+                    "WHERE uid = CAST(:uid AS uuid)"
+                ),
+                {"uid": source_uid},
+            )
+        if artifact_object_key:
+            minio.remove_object("dataops-bucket", artifact_object_key)
+        engine.dispose()

+ 64 - 0
tests/test_database_migrations.py

@@ -182,3 +182,67 @@ def test_alembic_upgrade_is_repeatable_and_downgrade_preserves_tables():
             )
             )
             cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
             cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
         connection.close()
         connection.close()
+
+
+@pytest.mark.integration
+def test_concurrent_alembic_upgrades_are_serialized():
+    admin_url = os.environ.get("TEST_POSTGRES_ADMIN_URL")
+    if not admin_url:
+        pytest.skip("TEST_POSTGRES_ADMIN_URL is not configured")
+
+    parsed = make_url(admin_url)
+    database_name = f"dataops_concurrent_{uuid.uuid4().hex[:12]}"
+    target_url = parsed.set(database=database_name).render_as_string(
+        hide_password=False
+    )
+    connection = psycopg2.connect(admin_url)
+    connection.autocommit = True
+    try:
+        with connection.cursor() as cursor:
+            cursor.execute(f'CREATE DATABASE "{database_name}"')
+
+        env = os.environ.copy()
+        env["SQLALCHEMY_DATABASE_URI"] = target_url
+        command = [
+            str(ROOT / ".venv" / "bin" / "alembic"),
+            "-c",
+            "alembic.ini",
+            "upgrade",
+            "head",
+        ]
+        processes = [
+            subprocess.Popen(
+                command,
+                cwd=ROOT,
+                env=env,
+                stdout=subprocess.PIPE,
+                stderr=subprocess.STDOUT,
+                text=True,
+            )
+            for _ in range(2)
+        ]
+        results = [process.communicate(timeout=120) for process in processes]
+
+        failures = [
+            output
+            for process, (output, _) in zip(processes, results)
+            if process.returncode != 0
+        ]
+        assert failures == []
+
+        engine = create_engine(target_url)
+        try:
+            assert set(inspect(engine).get_table_names(schema="public")) >= (
+                EXPECTED_UPGRADED_TABLES | {"alembic_version"}
+            )
+        finally:
+            engine.dispose()
+    finally:
+        with connection.cursor() as cursor:
+            cursor.execute(
+                "SELECT pg_terminate_backend(pid) FROM pg_stat_activity "
+                "WHERE datname = %s AND pid <> pg_backend_pid()",
+                (database_name,),
+            )
+            cursor.execute(f'DROP DATABASE IF EXISTS "{database_name}"')
+        connection.close()