from __future__ import annotations import os import re from collections.abc import Iterator from importlib.util import module_from_spec, spec_from_file_location from pathlib import Path from types import SimpleNamespace from typing import Any import pytest from alembic.config import Config from commercial_migration_assertions import ( _assert_commercial_head_schema, _assert_commercial_runtime_invariants, ) from financial_connector_migration_assertions import ( _assert_financial_connector_head_schema, _assert_financial_connector_runtime_invariants, ) from release_telemetry_migration_assertions import ( _assert_release_telemetry_head_schema, _assert_release_telemetry_runtime_invariants, ) from savings_migration_assertions import ( _assert_savings_head_schema, _assert_savings_runtime_invariants, ) from sqlalchemy import create_engine, inspect, text from sqlalchemy.engine import Engine, make_url from sqlalchemy.exc import IntegrityError from sqlalchemy.pool import NullPool from alembic import command from app.core.config import get_settings from app.db.migration_preflight import ( LEGACY_ADOPTABLE_HISTORICAL_CASE_TABLES, MigrationPreflightError, validate_migration_state, ) from app.db.schema_ownership import MIGRATION_OWNED_TABLES, create_legacy_schema from app.models.ai_memory import MemoryEntry from app.models.approval_task import ApprovalTask, ApprovalTaskEvent from app.models.risk_disposition import RiskDisposition, RiskDispositionEvent from app.models.risk_observation import RiskObservation MIGRATION_TEST_DATABASE_URL = os.getenv("MIGRATION_TEST_DATABASE_URL", "").strip() LEGACY_PROBE_TABLE = "legacy_migration_probe_records" HEAD_REVISION = "20260717_0028" SERVER_DIR = Path(__file__).resolve().parents[1] ALEMBIC_INI_PATH = SERVER_DIR / "alembic.ini" class _UnsupportedDialectOperationGuard: def __init__(self, dialect_name: str = "sqlite") -> None: self.bind = SimpleNamespace(dialect=SimpleNamespace(name=dialect_name)) self.mutation_calls: list[str] = [] def get_bind(self) -> SimpleNamespace: return self.bind def __getattr__(self, name: str) -> Any: self.mutation_calls.append(name) raise AssertionError(f"unsupported dialect attempted migration operation: {name}") def _load_migration_module(filename: str) -> Any: path = SERVER_DIR / "alembic" / "versions" / filename spec = spec_from_file_location(f"migration_test_{path.stem}", path) assert spec is not None and spec.loader is not None module = module_from_spec(spec) spec.loader.exec_module(module) return module def _normalize_probe_component(value: str) -> str: return re.sub(r"[^a-z0-9]+", "-", value.lower()).strip("-") def _is_disposable_probe_host(value: str) -> bool: markers = ("migration-probe", "disposable-probe") return value in markers or any( value.startswith(f"{marker}-") or value.startswith(f"x-financial-{marker}-") for marker in markers ) def _is_disposable_probe_database(value: str) -> bool: markers = ("migration-probe", "disposable-probe") return value in markers or any(value.startswith(f"{marker}-") for marker in markers) def _require_disposable_probe_url(raw_url: str) -> str: try: parsed = make_url(raw_url) except Exception as exc: # pragma: no cover - SQLAlchemy 提供具体解析异常 raise RuntimeError("MIGRATION_TEST_DATABASE_URL 不是有效的数据库 URL") from exc if parsed.get_backend_name() != "postgresql": raise RuntimeError("迁移测试只允许连接 PostgreSQL 一次性数据库") host = _normalize_probe_component(parsed.host or "") database = _normalize_probe_component(parsed.database or "") if not _is_disposable_probe_host(host): raise RuntimeError("迁移测试数据库主机名必须使用 migration-probe 或 disposable-probe 前缀") if not _is_disposable_probe_database(database): raise RuntimeError("迁移测试数据库名必须使用 migration-probe 或 disposable-probe 前缀") return raw_url @pytest.fixture(scope="module") def migration_database_url() -> Iterator[str]: if not MIGRATION_TEST_DATABASE_URL: pytest.skip("仅在显式配置 MIGRATION_TEST_DATABASE_URL 时运行一次性 PostgreSQL 迁移测试") database_url = _require_disposable_probe_url(MIGRATION_TEST_DATABASE_URL) previous_database_url = os.environ.get("DATABASE_URL") os.environ["DATABASE_URL"] = database_url get_settings.cache_clear() try: resolved_url = get_settings().resolved_database_url if make_url(resolved_url) != make_url(database_url): raise RuntimeError("运行时数据库 URL 未解析到 MIGRATION_TEST_DATABASE_URL") yield database_url finally: if previous_database_url is None: os.environ.pop("DATABASE_URL", None) else: os.environ["DATABASE_URL"] = previous_database_url get_settings.cache_clear() def _alembic_config(database_url: str) -> Config: config = Config(str(ALEMBIC_INI_PATH)) config.set_main_option("sqlalchemy.url", database_url.replace("%", "%%")) return config def _upgrade_head(database_url: str) -> None: _upgrade_revision(database_url, "head") def _upgrade_revision(database_url: str, revision: str) -> None: get_settings.cache_clear() command.upgrade(_alembic_config(database_url), revision) def _downgrade_base(database_url: str) -> None: _downgrade_revision(database_url, "base") def _downgrade_revision(database_url: str, revision: str) -> None: get_settings.cache_clear() command.downgrade(_alembic_config(database_url), revision) def _table_names(engine: Engine) -> set[str]: return set(inspect(engine).get_table_names(schema="public")) def _assert_unique_constraint( engine: Engine, table_name: str, constraint_name: str, expected_columns: tuple[str, ...], ) -> None: constraints = { str(item["name"]): tuple(item["column_names"]) for item in inspect(engine).get_unique_constraints(table_name, schema="public") } assert constraints.get(constraint_name) == expected_columns def _assert_indexes( engine: Engine, table_name: str, expected_indexes: dict[str, tuple[str, ...]], ) -> None: indexes = { str(item["name"]): tuple(item["column_names"]) for item in inspect(engine).get_indexes(table_name, schema="public") } for index_name, expected_columns in expected_indexes.items(): assert indexes.get(index_name) == expected_columns def _assert_postgresql_index_predicate( engine: Engine, table_name: str, index_name: str, *expected_fragments: str, ) -> None: with engine.connect() as connection: index_definition = str( connection.scalar( text( "SELECT indexdef FROM pg_indexes " "WHERE schemaname = 'public' " "AND tablename = :table_name AND indexname = :index_name" ), {"table_name": table_name, "index_name": index_name}, ) or "" ).lower() for fragment in expected_fragments: assert fragment.lower() in index_definition def _assert_check_constraint( engine: Engine, table_name: str, constraint_name: str, ) -> None: constraints = { str(item["name"]) for item in inspect(engine).get_check_constraints(table_name, schema="public") } assert constraint_name in constraints def _assert_cascade_foreign_key(engine: Engine, table_name: str) -> None: foreign_keys = inspect(engine).get_foreign_keys(table_name, schema="public") matching = [ item for item in foreign_keys if item["constrained_columns"] == ["expense_case_id"] and item["referred_table"] == "expense_cases" and item["referred_columns"] == ["id"] ] assert len(matching) == 1 assert str(matching[0].get("options", {}).get("ondelete", "")).upper() == "CASCADE" def _assert_composite_foreign_key( engine: Engine, table_name: str, constrained_columns: tuple[str, ...], referred_table: str, referred_columns: tuple[str, ...] = ("tenant_id", "id"), ) -> None: foreign_keys = inspect(engine).get_foreign_keys(table_name, schema="public") matching = [ item for item in foreign_keys if tuple(item["constrained_columns"]) == constrained_columns and item["referred_table"] == referred_table and tuple(item["referred_columns"]) == referred_columns ] assert len(matching) == 1 assert str(matching[0].get("options", {}).get("ondelete", "")).upper() == "RESTRICT" def _assert_no_foreign_key( engine: Engine, table_name: str, constrained_columns: tuple[str, ...], referred_table: str, ) -> None: foreign_keys = inspect(engine).get_foreign_keys(table_name, schema="public") assert not any( tuple(item["constrained_columns"]) == constrained_columns and item["referred_table"] == referred_table for item in foreign_keys ) def _assert_head_schema(engine: Engine) -> None: names = _table_names(engine) assert MIGRATION_OWNED_TABLES.issubset(names) assert "alembic_version" in names with engine.connect() as connection: assert connection.scalar(text("SELECT version_num FROM alembic_version")) == HEAD_REVISION _assert_unique_constraint( engine, "expense_cases", "uq_expense_cases_tenant_case_no", ("tenant_id", "case_no"), ) _assert_unique_constraint( engine, "expense_case_links", "uq_expense_case_links_resource", ("resource_type", "resource_id"), ) _assert_unique_constraint( engine, "business_events", "uq_business_events_tenant_case_id", ("tenant_id", "expense_case_id", "id"), ) _assert_unique_constraint( engine, "business_events", "uq_business_event_idempotency", ("tenant_id", "aggregate_type", "aggregate_id", "event_type", "idempotency_key"), ) _assert_unique_constraint( engine, "auth_sessions", "uq_auth_sessions_token_hash", ("token_hash",), ) _assert_unique_constraint( engine, "ai_decisions", "uq_ai_decisions_tenant_case_id", ("tenant_id", "expense_case_id", "id"), ) _assert_unique_constraint( engine, "ai_decisions", "uq_ai_decisions_tenant_idempotency", ("tenant_id", "idempotency_key"), ) _assert_unique_constraint( engine, "ai_decisions", "uq_ai_decisions_tenant_preview_decision", ("tenant_id", "preview_decision_id"), ) _assert_unique_constraint( engine, "ai_application_preview_decisions", "uq_ai_application_preview_decisions_issue_request", ("tenant_id", "actor_id", "auth_session_id", "issue_request_id"), ) _assert_unique_constraint( engine, "ai_decision_feedback", "uq_ai_decision_feedback_tenant_id", ("tenant_id", "id"), ) _assert_unique_constraint( engine, "ai_decision_feedback", "uq_ai_decision_feedback_tenant_idempotency", ("tenant_id", "idempotency_key"), ) _assert_unique_constraint( engine, "workflow_outcomes", "uq_workflow_outcomes_tenant_id", ("tenant_id", "id"), ) _assert_unique_constraint( engine, "workflow_outcomes", "uq_workflow_outcomes_tenant_idempotency", ("tenant_id", "idempotency_key"), ) _assert_unique_constraint( engine, "memory_entries", "uq_memory_entries_generation", ("tenant_id", "scope_type", "scope_id", "scene", "field_key", "generation"), ) _assert_unique_constraint( engine, "memory_evidence_links", "uq_memory_evidence_links_entry_case", ("tenant_id", "memory_entry_id", "expense_case_id"), ) _assert_unique_constraint( engine, "attachment_association_jobs", "uq_attachment_association_jobs_owner_dedupe", ("tenant_id", "owner_username", "dedupe_key", "generation"), ) _assert_unique_constraint( engine, "risk_observations", "uq_risk_observations_tenant_key", ("tenant_id", "observation_key"), ) _assert_unique_constraint( engine, "risk_observations", "uq_risk_observations_tenant_id", ("tenant_id", "id"), ) _assert_unique_constraint( engine, "few_shot_samples", "uq_few_shot_samples_tenant_key", ("tenant_id", "sample_key"), ) _assert_unique_constraint( engine, "approval_action_ledgers", "uq_approval_action_ledger_request", ("tenant_id", "actor_id", "request_id"), ) _assert_indexes( engine, "approval_action_ledgers", { "ix_approval_action_ledger_claim_action": ( "tenant_id", "claim_id", "action", ) }, ) _assert_check_constraint( engine, "approval_action_ledgers", "ck_approval_action_ledger_action", ) _assert_check_constraint( engine, "approval_action_ledgers", "ck_approval_action_ledger_completion", ) _assert_unique_constraint( engine, "risk_dispositions", "uq_risk_dispositions_tenant_observation", ("tenant_id", "observation_id"), ) _assert_unique_constraint( engine, "risk_dispositions", "uq_risk_dispositions_tenant_id", ("tenant_id", "id"), ) _assert_composite_foreign_key( engine, "risk_dispositions", ("tenant_id", "observation_id"), "risk_observations", ) _assert_indexes( engine, "risk_dispositions", { "ix_risk_dispositions_tenant_lifecycle_due": ( "tenant_id", "lifecycle_status", "due_at", ), "ix_risk_dispositions_assignee": ("tenant_id", "assignee"), "ix_risk_dispositions_tenant_waiver_expiry": ( "tenant_id", "lifecycle_status", "waiver_expires_at", ), }, ) _assert_check_constraint( engine, "risk_dispositions", "ck_risk_dispositions_adjudication", ) _assert_check_constraint( engine, "risk_dispositions", "ck_risk_dispositions_lifecycle", ) _assert_check_constraint( engine, "risk_dispositions", "ck_risk_dispositions_version", ) _assert_check_constraint( engine, "risk_dispositions", "ck_risk_dispositions_waiver_request", ) _assert_check_constraint( engine, "risk_dispositions", "ck_risk_dispositions_waiver_decision", ) _assert_check_constraint( engine, "risk_dispositions", "ck_risk_dispositions_waiver_lifecycle", ) _assert_unique_constraint( engine, "risk_disposition_events", "uq_risk_disposition_events_tenant_request", ("tenant_id", "request_id"), ) _assert_unique_constraint( engine, "risk_disposition_events", "uq_risk_disposition_events_version", ("disposition_id", "version"), ) _assert_indexes( engine, "risk_disposition_events", { "ix_risk_disposition_events_disposition_id": ("disposition_id",), "ix_risk_disposition_events_tenant_observation_time": ( "tenant_id", "observation_id", "created_at", ), }, ) _assert_check_constraint( engine, "risk_disposition_events", "ck_risk_disposition_events_action", ) _assert_check_constraint( engine, "risk_disposition_events", "ck_risk_disposition_events_version", ) _assert_composite_foreign_key( engine, "risk_disposition_events", ("tenant_id", "disposition_id"), "risk_dispositions", ) _assert_composite_foreign_key( engine, "risk_disposition_events", ("tenant_id", "observation_id"), "risk_observations", ) risk_disposition_event_columns = { str(item["name"]) for item in inspect(engine).get_columns( "risk_disposition_events", schema="public", ) } assert "response_json" in risk_disposition_event_columns risk_disposition_columns = { str(item["name"]) for item in inspect(engine).get_columns("risk_dispositions", schema="public") } assert { "waiver_requester_id", "waiver_requester_name", "waiver_requested_at", "waiver_reason", "waiver_scope", "waiver_expires_at", "waiver_conditions_json", "waiver_decision", "waiver_decider_id", "waiver_decider_name", "waiver_decided_at", "waiver_decision_reason", }.issubset(risk_disposition_columns) with engine.connect() as connection: append_only_trigger_count = int( connection.scalar( text( "SELECT COUNT(*) FROM pg_trigger trigger " "JOIN pg_class relation ON relation.oid = trigger.tgrelid " "WHERE relation.relname = 'risk_disposition_events' " "AND trigger.tgname = 'trg_risk_disposition_events_append_only' " "AND NOT trigger.tgisinternal" ) ) or 0 ) assert append_only_trigger_count == 1 _assert_unique_constraint( engine, "approval_tasks", "uq_approval_tasks_tenant_id", ("tenant_id", "id"), ) _assert_unique_constraint( engine, "approval_tasks", "uq_approval_tasks_tenant_node_entry", ("tenant_id", "node_entry_key"), ) _assert_unique_constraint( engine, "approval_tasks", "uq_approval_tasks_node_participant", ("tenant_id", "node_instance_id", "assignee_kind", "assignee_key"), ) _assert_unique_constraint( engine, "approval_task_events", "uq_approval_task_events_actor_request", ("tenant_id", "actor_id", "request_id"), ) _assert_unique_constraint( engine, "approval_task_events", "uq_approval_task_events_task_version", ("tenant_id", "task_id", "result_task_version"), ) _assert_composite_foreign_key( engine, "approval_tasks", ("tenant_id", "parent_task_id"), "approval_tasks", ) _assert_composite_foreign_key( engine, "approval_tasks", ("tenant_id", "expense_case_id"), "expense_cases", ) _assert_composite_foreign_key( engine, "approval_task_events", ("tenant_id", "task_id"), "approval_tasks", ) _assert_no_foreign_key( engine, "approval_tasks", ("claim_id",), "expense_claims", ) _assert_indexes( engine, "approval_tasks", { "uq_approval_tasks_open_root_per_claim": ("tenant_id", "claim_id"), "ix_approval_tasks_personal_inbox": ( "tenant_id", "assignee_kind", "assignee_key", "status", "due_at", ), "ix_approval_tasks_tenant_queue": ( "tenant_id", "status", "priority_score", "due_at", ), "ix_approval_tasks_tenant_claim": ( "tenant_id", "claim_id", "node_sequence", ), "ix_approval_tasks_tenant_node": ( "tenant_id", "node_instance_id", "sequence_order", ), }, ) _assert_postgresql_index_predicate( engine, "approval_tasks", "uq_approval_tasks_open_root_per_claim", "task_kind", "root", "status", "waiting", "pending", ) _assert_indexes( engine, "approval_task_events", { "ix_approval_task_events_tenant_task_time": ( "tenant_id", "task_id", "occurred_at", ), "ix_approval_task_events_tenant_node_time": ( "tenant_id", "node_instance_id", "occurred_at", ), }, ) for constraint_name in ( "ck_approval_tasks_task_kind", "ck_approval_tasks_status", "ck_approval_tasks_version", "ck_approval_tasks_priority_score", "ck_approval_tasks_evidence_completeness", "ck_approval_task_events_type", "ck_approval_task_events_actor_type", "ck_approval_task_events_version", ): table_name = ( "approval_task_events" if constraint_name.startswith("ck_approval_task_events") else "approval_tasks" ) _assert_check_constraint(engine, table_name, constraint_name) approval_event_columns = { str(item["name"]) for item in inspect(engine).get_columns("approval_task_events", schema="public") } assert "response_json" in approval_event_columns with engine.connect() as connection: approval_append_only_trigger_count = int( connection.scalar( text( "SELECT COUNT(*) FROM pg_trigger trigger " "JOIN pg_class relation ON relation.oid = trigger.tgrelid " "WHERE relation.relname = 'approval_task_events' " "AND trigger.tgname = 'trg_approval_task_events_append_only' " "AND NOT trigger.tgisinternal" ) ) or 0 ) assert approval_append_only_trigger_count == 1 _assert_check_constraint( engine, "memory_entries", "ck_memory_entries_expired_fields", ) _assert_check_constraint( engine, "memory_entries", "ck_memory_entries_scope_origin", ) _assert_check_constraint( engine, "memory_entries", "ck_memory_entries_enterprise_scope", ) _assert_check_constraint( engine, "memory_entries", "ck_memory_entries_management_audit", ) _assert_check_constraint( engine, "memory_entries", "ck_memory_entries_management_idempotency_pair", ) _assert_check_constraint( engine, "memory_entries", "ck_memory_entries_revoke_idempotency_pair", ) _assert_check_constraint( engine, "attachment_association_jobs", "ck_attachment_association_jobs_running_lease", ) _assert_check_constraint( engine, "attachment_association_jobs", "ck_attachment_association_jobs_generation", ) _assert_savings_head_schema(engine) _assert_commercial_head_schema(engine) _assert_financial_connector_head_schema(engine) _assert_release_telemetry_head_schema(engine) _assert_indexes( engine, "expense_cases", { "ix_expense_cases_tenant_stage": ("tenant_id", "current_stage"), "ix_expense_cases_tenant_status": ("tenant_id", "status"), }, ) _assert_postgresql_index_predicate( engine, "memory_entries", "uq_memory_entries_active_scope", "status", "'active'", "scope_type", "department", "enterprise", ) _assert_indexes( engine, "expense_case_links", {"ix_expense_case_links_tenant_case": ("tenant_id", "expense_case_id")}, ) _assert_indexes( engine, "business_events", { "ix_business_events_aggregate": ("aggregate_type", "aggregate_id"), "ix_business_events_outbox": ("delivery_status", "occurred_at"), "ix_business_events_tenant_case_time": ( "tenant_id", "expense_case_id", "occurred_at", ), }, ) _assert_indexes( engine, "auth_sessions", { "ix_auth_sessions_principal_active": ( "principal_type", "revoked_at", "expires_at", ), "ix_auth_sessions_tenant_username": ("tenant_id", "username"), }, ) _assert_indexes( engine, "attachment_association_jobs", { "ix_attachment_association_jobs_owner_time": ( "tenant_id", "owner_username", "created_at", ), "ix_attachment_association_jobs_status_lease": ( "status", "lease_expires_at", ), }, ) _assert_indexes( engine, "ai_application_preview_decisions", { "ix_ai_application_preview_decisions_actor_status_expiry": ( "tenant_id", "actor_id", "status", "expires_at", ), "ix_ai_application_preview_decisions_conversation": ( "tenant_id", "conversation_id", "created_at", ), }, ) _assert_indexes( engine, "ai_decisions", { "ix_ai_decisions_tenant_case_time": ( "tenant_id", "expense_case_id", "created_at", ), "ix_ai_decisions_tenant_subject": ( "tenant_id", "subject_type", "subject_id", ), }, ) _assert_indexes( engine, "ai_decision_feedback", { "ix_ai_decision_feedback_tenant_decision_time": ( "tenant_id", "decision_id", "created_at", ), }, ) _assert_indexes( engine, "workflow_outcomes", { "ix_workflow_outcomes_tenant_case_time": ( "tenant_id", "expense_case_id", "effective_at", ), }, ) _assert_indexes( engine, "memory_entries", { "ix_memory_entries_scope_lookup": ( "tenant_id", "scope_type", "scope_id", "scene", "field_key", "status", ), "uq_memory_entries_management_request": ( "tenant_id", "management_request_id", ), "uq_memory_entries_revoke_request": ( "tenant_id", "revoke_request_id", ), "uq_memory_entries_active_scope": ( "tenant_id", "scope_type", "scope_id", "scene", "field_key", ), }, ) _assert_indexes( engine, "memory_evidence_links", { "ix_memory_evidence_links_entry_time": ( "tenant_id", "memory_entry_id", "created_at", ), }, ) _assert_indexes( engine, "risk_observations", { "ix_risk_observations_tenant_status": ( "tenant_id", "status", "created_at", ), }, ) _assert_indexes( engine, "few_shot_samples", { "ix_few_shot_samples_tenant_rule_lookup": ( "tenant_id", "scene", "policy_ref", "rule_version", "status", ), }, ) _assert_cascade_foreign_key(engine, "expense_case_links") _assert_cascade_foreign_key(engine, "business_events") _assert_composite_foreign_key( engine, "ai_decisions", ("tenant_id", "preview_decision_id"), "ai_application_preview_decisions", ) _assert_composite_foreign_key( engine, "ai_decisions", ("tenant_id", "expense_case_id"), "expense_cases", ) _assert_no_foreign_key( engine, "risk_observations", ("claim_id",), "expense_claims", ) _assert_composite_foreign_key( engine, "ai_decisions", ("tenant_id", "expense_case_id", "business_event_id"), "business_events", ("tenant_id", "expense_case_id", "id"), ) _assert_composite_foreign_key( engine, "ai_decision_feedback", ("tenant_id", "decision_id"), "ai_decisions", ) _assert_composite_foreign_key( engine, "workflow_outcomes", ("tenant_id", "expense_case_id"), "expense_cases", ) _assert_composite_foreign_key( engine, "workflow_outcomes", ("tenant_id", "expense_case_id", "decision_id"), "ai_decisions", ("tenant_id", "expense_case_id", "id"), ) _assert_composite_foreign_key( engine, "workflow_outcomes", ("tenant_id", "expense_case_id", "business_event_id"), "business_events", ("tenant_id", "expense_case_id", "id"), ) _assert_composite_foreign_key( engine, "memory_entries", ("tenant_id", "superseded_by_id"), "memory_entries", ) _assert_composite_foreign_key( engine, "memory_evidence_links", ("tenant_id", "memory_entry_id"), "memory_entries", ) _assert_composite_foreign_key( engine, "memory_evidence_links", ("tenant_id", "expense_case_id"), "expense_cases", ) _assert_composite_foreign_key( engine, "memory_evidence_links", ("tenant_id", "decision_id"), "ai_decisions", ) _assert_composite_foreign_key( engine, "memory_evidence_links", ("tenant_id", "feedback_id"), "ai_decision_feedback", ) _assert_composite_foreign_key( engine, "memory_evidence_links", ("tenant_id", "outcome_id"), "workflow_outcomes", ) def _assert_runtime_cascade(engine: Engine) -> None: with engine.begin() as connection: connection.execute( text( """ INSERT INTO expense_cases ( id, tenant_id, case_no, scene_code, title, current_stage, status ) VALUES ( 'migration-probe-case', 'migration-probe', 'CASE-MIGRATION-PROBE', 'reimbursement', '迁移级联验证', 'claiming', 'active' ) """ ) ) connection.execute( text( """ INSERT INTO expense_case_links ( id, tenant_id, expense_case_id, resource_type, resource_id, relation_type ) VALUES ( 'migration-probe-link', 'migration-probe', 'migration-probe-case', 'expense_claim', 'migration-probe-claim', 'claim' ) """ ) ) connection.execute( text( """ INSERT INTO business_events ( id, tenant_id, expense_case_id, aggregate_type, aggregate_id, event_type, event_version, idempotency_key, correlation_id, actor_id, actor_type, payload_json, delivery_status, delivery_attempts ) VALUES ( 'migration-probe-event', 'migration-probe', 'migration-probe-case', 'expense_claim', 'migration-probe-claim', 'claim_draft_created', 1, 'migration-probe-idempotency', 'migration-probe-correlation', 'migration-probe-user', 'user', '{}', 'pending', 0 ) """ ) ) connection.execute(text("DELETE FROM expense_cases WHERE id = 'migration-probe-case'")) assert ( connection.scalar( text("SELECT COUNT(*) FROM expense_case_links WHERE id = 'migration-probe-link'") ) == 0 ) assert ( connection.scalar( text("SELECT COUNT(*) FROM business_events WHERE id = 'migration-probe-event'") ) == 0 ) def _assert_learning_ledger_tenant_boundary(engine: Engine) -> None: with engine.begin() as connection: connection.execute( text( """ INSERT INTO expense_cases ( id, tenant_id, case_no, scene_code, title, current_stage, status ) VALUES ( 'learning-probe-case', 'learning-probe', 'CASE-LEARNING-PROBE', 'travel', '学习闭环迁移验证', 'claiming', 'active' ), ( 'learning-probe-case-b', 'learning-probe', 'CASE-LEARNING-PROBE-B', 'travel', '学习闭环同租户第二费用单', 'claiming', 'active' ) """ ) ) connection.execute( text( """ INSERT INTO business_events ( id, tenant_id, expense_case_id, aggregate_type, aggregate_id, event_type, event_version, idempotency_key, correlation_id, actor_id, actor_type, payload_json, delivery_status, delivery_attempts ) VALUES ( 'learning-probe-event', 'learning-probe', 'learning-probe-case', 'expense_claim', 'learning-probe-claim', 'claim_draft_created', 1, 'learning-probe-event-key', 'learning-probe-correlation', 'learning-probe-user', 'user', '{}', 'pending', 0 ) """ ) ) connection.execute( text( """ INSERT INTO ai_decisions ( id, tenant_id, expense_case_id, business_event_id, expense_claim_id, agent_run_id, correlation_id, subject_type, subject_id, decision_type, decision_source, status, automation_mode, confidence, suggestion_json, evidence_json, version_json, schema_version, idempotency_key, content_fingerprint ) VALUES ( 'learning-probe-decision', 'learning-probe', 'learning-probe-case', 'learning-probe-event', 'learning-probe-claim', NULL, 'learning-probe-correlation', 'expense_claim', 'learning-probe-claim', 'expense_application_prefill', 'hybrid', 'executed', 'prefill', 0.9, '{}', '{}', '{}', 1, 'learning-probe-decision-key', 'sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa' ) """ ) ) connection.execute( text( """ INSERT INTO ai_decision_feedback ( id, tenant_id, decision_id, expense_claim_id, correlation_id, feedback_type, action_type, actor_id, actor_type, evidence_source, verification_status, final_value_json, changed_fields_json, idempotency_key, content_fingerprint ) VALUES ( 'learning-probe-feedback', 'learning-probe', 'learning-probe-decision', 'learning-probe-claim', 'learning-probe-correlation', 'accepted', 'save_draft', 'learning-probe-user', 'user', 'client_action_confirmation', 'client_observed', '{}', '[]', 'learning-probe-feedback-key', 'sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb' ) """ ) ) connection.execute( text( """ INSERT INTO workflow_outcomes ( id, tenant_id, expense_case_id, decision_id, business_event_id, expense_claim_id, correlation_id, outcome_type, outcome_status, actor_id, actor_type, result_json, idempotency_key, content_fingerprint ) VALUES ( 'learning-probe-outcome', 'learning-probe', 'learning-probe-case', 'learning-probe-decision', 'learning-probe-event', 'learning-probe-claim', 'learning-probe-correlation', 'draft_saved', 'recorded', 'learning-probe-user', 'user', '{}', 'learning-probe-outcome-key', 'sha256:cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc' ) """ ) ) with pytest.raises(IntegrityError): with engine.begin() as connection: connection.execute( text( """ INSERT INTO ai_decisions ( id, tenant_id, expense_case_id, expense_claim_id, correlation_id, subject_type, subject_id, decision_type, decision_source, status, automation_mode, suggestion_json, evidence_json, version_json, schema_version, idempotency_key, content_fingerprint ) VALUES ( 'cross-tenant-decision', 'other-tenant', 'learning-probe-case', 'cross-tenant-claim', 'cross-tenant-correlation', 'expense_claim', 'cross-tenant-claim', 'expense_application_prefill', 'heuristic', 'executed', 'prefill', '{}', '{}', '{}', 1, 'cross-tenant-key', 'sha256:dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd' ) """ ) ) with pytest.raises(IntegrityError): with engine.begin() as connection: connection.execute( text( """ INSERT INTO ai_decision_feedback ( id, tenant_id, decision_id, expense_claim_id, correlation_id, feedback_type, action_type, actor_id, actor_type, evidence_source, verification_status, training_eligible, final_value_json, changed_fields_json, idempotency_key, content_fingerprint ) VALUES ( 'unverified-training-feedback', 'learning-probe', 'learning-probe-decision', 'learning-probe-claim', 'unverified-training-correlation', 'accepted', 'save_draft', 'learning-probe-user', 'user', 'client_action_confirmation', 'client_observed', TRUE, '{}', '[]', 'unverified-training-feedback-key', 'sha256:9999999999999999999999999999999999999999999999999999999999999999' ) """ ) ) with pytest.raises(IntegrityError): with engine.begin() as connection: connection.execute( text( """ INSERT INTO ai_decisions ( id, tenant_id, expense_case_id, business_event_id, expense_claim_id, correlation_id, subject_type, subject_id, decision_type, decision_source, status, automation_mode, suggestion_json, evidence_json, version_json, schema_version, idempotency_key, content_fingerprint ) VALUES ( 'cross-case-event-decision', 'learning-probe', 'learning-probe-case-b', 'learning-probe-event', 'cross-case-event-claim', 'cross-case-event-correlation', 'expense_claim', 'cross-case-event-claim', 'expense_application_prefill', 'heuristic', 'executed', 'prefill', '{}', '{}', '{}', 1, 'cross-case-event-key', 'sha256:eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee' ) """ ) ) with pytest.raises(IntegrityError): with engine.begin() as connection: connection.execute( text( """ INSERT INTO workflow_outcomes ( id, tenant_id, expense_case_id, decision_id, expense_claim_id, correlation_id, outcome_type, outcome_status, actor_id, actor_type, result_json, idempotency_key, content_fingerprint ) VALUES ( 'cross-case-decision-outcome', 'learning-probe', 'learning-probe-case-b', 'learning-probe-decision', 'cross-case-decision-claim', 'cross-case-decision-correlation', 'draft_saved', 'recorded', 'learning-probe-user', 'user', '{}', 'cross-case-decision-key', 'sha256:ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff' ) """ ) ) def _create_legacy_sentinel(engine: Engine) -> None: with engine.begin() as connection: connection.execute( text( f""" CREATE TABLE {LEGACY_PROBE_TABLE} ( id VARCHAR(64) PRIMARY KEY, payload VARCHAR(255) NOT NULL ) """ ) ) connection.execute( text( f""" INSERT INTO {LEGACY_PROBE_TABLE} (id, payload) VALUES ('legacy-sentinel', 'must-survive-migration-cycle') """ ) ) def _create_hierarchical_memory_downgrade_probe(engine: Engine) -> None: """验证 0007 在真实组织记忆存在时仍可安全降级。""" with engine.begin() as connection: connection.execute( text( """ INSERT INTO memory_entries ( id, tenant_id, scope_type, scope_id, origin_type, managed_by, managed_at, management_reason, policy_version, scene, field_key, generation, value_json, value_fingerprint, status, evidence_count, approved_evidence_count, confidence, last_evidence_at, candidate_expires_at, activated_at, active_expires_at ) VALUES ( 'hierarchical-memory-downgrade-probe', 'migration-probe', 'enterprise', 'migration-probe', 'admin_managed', 'migration-admin', CURRENT_TIMESTAMP, '验证组织记忆降级', 'expense_application_transport_org_memory.v1', 'travel_application', 'transport_mode', 1, CAST('{"value":"火车"}' AS JSON), 'sha256:hierarchical-memory-downgrade-probe', 'active', 0, 0, 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + INTERVAL '180 days', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + INTERVAL '180 days' ) """ ) ) assert ( connection.scalar( text( "SELECT COUNT(*) FROM memory_entries " "WHERE id = 'hierarchical-memory-downgrade-probe'" ) ) == 1 ) def _create_duplicate_active_organization_memory_probe(engine: Engine) -> None: with engine.begin() as connection: connection.execute( text( """ INSERT INTO memory_entries ( id, tenant_id, scope_type, scope_id, origin_type, managed_by, managed_at, management_reason, policy_version, scene, field_key, generation, value_json, value_fingerprint, status, evidence_count, approved_evidence_count, confidence, last_evidence_at, candidate_expires_at, activated_at, active_expires_at ) VALUES ( 'duplicate-active-organization-a', 'migration-probe', 'department', 'department-probe', 'admin_managed', 'migration-admin', CURRENT_TIMESTAMP, '验证组织生效记忆重复预检', 'expense_application_transport_org_memory.v1', 'travel_application', 'transport_mode', 1, CAST('{"value":"火车"}' AS JSON), 'sha256:duplicate-active-organization-a', 'active', 0, 0, 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + INTERVAL '180 days', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + INTERVAL '180 days' ), ( 'duplicate-active-organization-b', 'migration-probe', 'department', 'department-probe', 'admin_managed', 'migration-admin', CURRENT_TIMESTAMP, '验证组织生效记忆重复预检', 'expense_application_transport_org_memory.v1', 'travel_application', 'transport_mode', 2, CAST('{"value":"飞机"}' AS JSON), 'sha256:duplicate-active-organization-b', 'active', 0, 0, 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + INTERVAL '180 days', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP + INTERVAL '180 days' ) """ ) ) def _delete_duplicate_active_organization_memory_probe(engine: Engine) -> None: with engine.begin() as connection: connection.execute( text( "DELETE FROM memory_entries " "WHERE id IN (" "'duplicate-active-organization-a', " "'duplicate-active-organization-b'" ")" ) ) def _create_enriched_few_shot_downgrade_probe(engine: Engine) -> None: with engine.begin() as connection: connection.execute( text( """ INSERT INTO few_shot_samples ( id, tenant_id, sample_key, scene, policy_ref, rule_version, payload_json ) VALUES ( 'enriched-few-shot-downgrade-probe', 'default', 'enriched:few-shot:downgrade-probe', 'expense_reimbursement', 'TRAVEL-001', 'v2', CAST('{}' AS JSON) ) """ ) ) def _delete_enriched_few_shot_downgrade_probe(engine: Engine) -> None: with engine.begin() as connection: connection.execute( text("DELETE FROM few_shot_samples WHERE id = 'enriched-few-shot-downgrade-probe'") ) def _create_historical_case_downgrade_probe(engine: Engine) -> None: with engine.begin() as connection: connection.execute( text( """ INSERT INTO risk_observations ( id, tenant_id, observation_key, subject_type, subject_key, risk_type, risk_signal, risk_level, contribution_scores_json, baseline_json, evidence_json, graph_node_keys_json, graph_edge_keys_json, policy_refs_json, similar_case_claim_ids_json, ontology_json, decision_trace_json ) VALUES ( 'historical-downgrade-observation', 'default', 'historical:downgrade:observation', 'expense_claim', 'claim:historical-downgrade', 'duplicate_invoice', 'duplicate_invoice', 'high', CAST('{}' AS JSON), CAST('{}' AS JSON), CAST('[]' AS JSON), CAST('[]' AS JSON), CAST('[]' AS JSON), CAST('[]' AS JSON), CAST('[]' AS JSON), CAST('{}' AS JSON), CAST('{}' AS JSON) ) """ ) ) connection.execute( text( """ INSERT INTO risk_observation_feedback ( id, observation_id, feedback_type, payload_json ) VALUES ( 'historical-downgrade-feedback', 'historical-downgrade-observation', 'confirm', CAST('{}' AS JSON) ) """ ) ) connection.execute( text( """ INSERT INTO few_shot_samples ( id, tenant_id, sample_key, source_observation_id, scene, policy_ref, rule_version, payload_json ) VALUES ( 'historical-downgrade-sample', 'default', 'historical:downgrade:sample', 'historical-downgrade-observation', 'expense_reimbursement', '', '', CAST('{}' AS JSON) ) """ ) ) def _assert_historical_case_downgrade_probe(engine: Engine) -> None: inspector = inspect(engine) risk_columns = {str(item["name"]) for item in inspector.get_columns("risk_observations")} sample_columns = {str(item["name"]) for item in inspector.get_columns("few_shot_samples")} assert "tenant_id" not in risk_columns assert {"tenant_id", "policy_ref", "rule_version"}.isdisjoint(sample_columns) assert not any( item["constrained_columns"] == ["claim_id"] and item["referred_table"] == "expense_claims" and item["referred_columns"] == ["id"] for item in inspector.get_foreign_keys("risk_observations") ) with engine.connect() as connection: assert ( connection.scalar( text( "SELECT COUNT(*) FROM risk_observations " "WHERE id = 'historical-downgrade-observation'" ) ) == 1 ) assert ( connection.scalar( text( "SELECT COUNT(*) FROM risk_observation_feedback " "WHERE id = 'historical-downgrade-feedback'" ) ) == 1 ) assert ( connection.scalar( text( "SELECT COUNT(*) FROM few_shot_samples WHERE id = 'historical-downgrade-sample'" ) ) == 1 ) def _create_legacy_connector_payload_probe(engine: Engine) -> None: with engine.begin() as connection: connection.execute( text( """ INSERT INTO financial_connector_configs ( id, tenant_id, provider, environment, key_version, secret_ref, allowed_event_types_json, clock_skew_seconds, status, created_by ) VALUES ( 'legacy-connector-config', 'tenant-legacy', 'legacy-bank', 'mock', 'v1', 'server/legacy', '["payment_settled"]', 300, 'active', 'migration-probe' ) """ ) ) connection.execute( text( """ INSERT INTO financial_connector_events ( id, tenant_id, config_id, provider, environment, direction, external_event_id, event_type, occurred_at, key_version, verification_level, request_fingerprint, content_hash, processing_status, correlation_id, normalized_payload_json, response_json ) VALUES ( 'legacy-connector-event', 'tenant-legacy', 'legacy-connector-config', 'legacy-bank', 'mock', 'inbound', 'legacy-external', 'payment_settled', now(), 'v1', 'simulated', 'sha256:aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa', 'sha256:bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb', 'processed', 'legacy-correlation', '{"claim_reference": "BX-SENSITIVE-001", "amount": "66.00"}', '{"accepted": true, "reconciliation_case_id": "legacy-case"}' ) """ ) ) def _assert_and_delete_legacy_connector_payload_probe(engine: Engine) -> None: with engine.begin() as connection: row = connection.execute( text( """ SELECT config.version, event.normalized_payload_json::jsonb ? 'claim_reference', event.normalized_payload_json ->> 'amount', event.response_json ->> 'projection_scope' FROM financial_connector_configs AS config JOIN financial_connector_events AS event ON event.config_id = config.id WHERE config.id = 'legacy-connector-config' """ ) ).one() assert row == (1, False, "66.00", "legacy_nonproduction_effect_unknown") connection.execute( text( "ALTER TABLE financial_connector_events " "DISABLE TRIGGER trg_financial_connector_events_append_only" ) ) connection.execute( text("DELETE FROM financial_connector_events WHERE id = 'legacy-connector-event'") ) connection.execute( text( "ALTER TABLE financial_connector_events " "ENABLE TRIGGER trg_financial_connector_events_append_only" ) ) connection.execute( text( "DELETE FROM financial_connector_configs " "WHERE id = 'legacy-connector-config'" ) ) def _assert_legacy_sentinel(engine: Engine) -> None: assert LEGACY_PROBE_TABLE in _table_names(engine) with engine.connect() as connection: payload = connection.scalar( text(f"SELECT payload FROM {LEGACY_PROBE_TABLE} WHERE id = 'legacy-sentinel'") ) assert payload == "must-survive-migration-cycle" def _assert_base_schema(engine: Engine) -> None: names = _table_names(engine) unsafe_owned_tables = ( MIGRATION_OWNED_TABLES - LEGACY_ADOPTABLE_HISTORICAL_CASE_TABLES ).intersection(names) assert not unsafe_owned_tables assert LEGACY_ADOPTABLE_HISTORICAL_CASE_TABLES.issubset(names) assert "alembic_version" in names with engine.connect() as connection: assert connection.scalar(text("SELECT COUNT(*) FROM alembic_version")) == 0 @pytest.mark.parametrize( ("filename", "direction"), [ ("20260716_0008_tenant_safe_historical_cases.py", "upgrade"), ("20260716_0008_tenant_safe_historical_cases.py", "downgrade"), ("20260716_0009_organization_memory_idempotency.py", "upgrade"), ("20260716_0009_organization_memory_idempotency.py", "downgrade"), ("20260716_0010_approval_action_protocol.py", "upgrade"), ("20260716_0010_approval_action_protocol.py", "downgrade"), ("20260716_0011_risk_disposition.py", "upgrade"), ("20260716_0011_risk_disposition.py", "downgrade"), ("20260716_0012_risk_disposition_response_snapshot.py", "upgrade"), ("20260716_0012_risk_disposition_response_snapshot.py", "downgrade"), ("20260716_0013_approval_tasks.py", "upgrade"), ("20260716_0013_approval_tasks.py", "downgrade"), ("20260716_0014_risk_waiver_decision.py", "upgrade"), ("20260716_0014_risk_waiver_decision.py", "downgrade"), ("20260716_0015_savings_value_ledger.py", "upgrade"), ("20260716_0015_savings_value_ledger.py", "downgrade"), ("20260716_0016_commercial_metering.py", "upgrade"), ("20260716_0016_commercial_metering.py", "downgrade"), ("20260716_0017_financial_connector_reconciliation.py", "upgrade"), ("20260716_0017_financial_connector_reconciliation.py", "downgrade"), ("20260716_0018_agent_asset_release_telemetry.py", "upgrade"), ("20260716_0018_agent_asset_release_telemetry.py", "downgrade"), ("20260716_0019_commercial_runtime_reservations.py", "upgrade"), ("20260716_0019_commercial_runtime_reservations.py", "downgrade"), ("20260716_0020_financial_connector_config_lifecycle.py", "upgrade"), ("20260716_0020_financial_connector_config_lifecycle.py", "downgrade"), ("20260716_0021_commercial_billing_periods.py", "upgrade"), ("20260716_0021_commercial_billing_periods.py", "downgrade"), ("20260716_0022_financial_connector_operational_events.py", "upgrade"), ("20260716_0022_financial_connector_operational_events.py", "downgrade"), ("20260716_0023_agent_asset_release_blind_audit.py", "upgrade"), ("20260716_0023_agent_asset_release_blind_audit.py", "downgrade"), ("20260717_0024_commercial_resource_quantity_bases.py", "upgrade"), ("20260717_0024_commercial_resource_quantity_bases.py", "downgrade"), ("20260717_0025_tenant_identity_foundation.py", "upgrade"), ("20260717_0025_tenant_identity_foundation.py", "downgrade"), ("20260717_0026_agent_asset_tenant_security.py", "upgrade"), ("20260717_0026_agent_asset_tenant_security.py", "downgrade"), ("20260717_0027_knowledge_tenant_security.py", "upgrade"), ("20260717_0027_knowledge_tenant_security.py", "downgrade"), ("20260717_0028_hermes_ontology_tenant_security.py", "upgrade"), ("20260717_0028_hermes_ontology_tenant_security.py", "downgrade"), ], ) def test_postgresql_only_migrations_reject_other_dialects_before_mutation( filename: str, direction: str, ) -> None: migration = _load_migration_module(filename) operation_guard = _UnsupportedDialectOperationGuard() migration.op = operation_guard with pytest.raises(RuntimeError, match="only supports PostgreSQL"): getattr(migration, direction)() assert operation_guard.mutation_calls == [] def test_risk_disposition_snapshot_migration_refuses_lossy_downgrade() -> None: migration = _load_migration_module("20260716_0012_risk_disposition_response_snapshot.py") operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="contains 1 immutable snapshot"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_approval_task_migration_refuses_non_empty_audit_chain_downgrade() -> None: migration = _load_migration_module("20260716_0013_approval_tasks.py") operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="audit chain is not empty"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_risk_waiver_migration_refuses_lossy_audit_downgrade() -> None: migration = _load_migration_module("20260716_0014_risk_waiver_decision.py") operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="immutable waiver audit data exists"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_savings_ledger_migration_refuses_non_empty_fact_downgrade() -> None: migration = _load_migration_module("20260716_0015_savings_value_ledger.py") operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="immutable value facts exist"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_commercial_migration_refuses_non_empty_contract_or_fact_downgrade() -> None: migration = _load_migration_module("20260716_0016_commercial_metering.py") operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="contracts or immutable facts exist"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_financial_connector_migration_refuses_non_empty_fact_downgrade() -> None: migration = _load_migration_module( "20260716_0017_financial_connector_reconciliation.py" ) operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="configurations or immutable facts exist"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_release_telemetry_migration_refuses_non_empty_fact_downgrade() -> None: migration = _load_migration_module( "20260716_0018_agent_asset_release_telemetry.py" ) operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="immutable observations or labels exist"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_runtime_reservation_migration_refuses_non_empty_hold_downgrade() -> None: migration = _load_migration_module( "20260716_0019_commercial_runtime_reservations.py" ) operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="operational quota holds exist"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_connector_lifecycle_migration_refuses_lossy_downgrade() -> None: migration = _load_migration_module( "20260716_0020_financial_connector_config_lifecycle.py" ) operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="versioned configuration state"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_commercial_billing_migration_refuses_non_empty_history_downgrade() -> None: migration = _load_migration_module( "20260716_0021_commercial_billing_periods.py" ) operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="immutable periods or audit facts exist"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_connector_operational_migration_refuses_non_empty_history_downgrade() -> None: migration = _load_migration_module( "20260716_0022_financial_connector_operational_events.py" ) operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="immutable operational facts exist"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_release_blind_audit_migration_refuses_non_empty_evidence_downgrade() -> None: migration = _load_migration_module( "20260716_0023_agent_asset_release_blind_audit.py" ) operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="immutable audit evidence exists"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_commercial_resource_basis_migration_refuses_lossy_downgrade() -> None: migration = _load_migration_module( "20260717_0024_commercial_resource_quantity_bases.py" ) operation_guard = _UnsupportedDialectOperationGuard(dialect_name="postgresql") operation_guard.bind.scalar = lambda _statement: 1 migration.op = operation_guard with pytest.raises(RuntimeError, match="resource reservations exist"): migration.downgrade() assert operation_guard.mutation_calls == [] def test_head_model_declares_soft_claim_reference_and_organization_only_active_index() -> None: claim_column = RiskObservation.__table__.c.claim_id assert not claim_column.foreign_keys claim_relationship = inspect(RiskObservation).relationships["claim"] assert claim_relationship.viewonly is True active_scope_index = next( index for index in MemoryEntry.__table__.indexes if index.name == "uq_memory_entries_active_scope" ) predicate = str(active_scope_index.dialect_options["postgresql"]["where"]) assert "status = 'active'" in predicate assert "scope_type IN ('department', 'enterprise')" in predicate def test_approval_task_model_declares_tenant_safety_and_open_root_invariant() -> None: claim_column = ApprovalTask.__table__.c.claim_id assert not claim_column.foreign_keys assert inspect(ApprovalTask).relationships["claim"].viewonly is True constraint_names = { constraint.name for constraint in ApprovalTask.__table__.constraints } assert { "uq_approval_tasks_tenant_id", "uq_approval_tasks_tenant_node_entry", "uq_approval_tasks_node_participant", "fk_approval_tasks_tenant_parent", "ck_approval_tasks_status", "ck_approval_tasks_version", "ck_approval_tasks_priority_score", "ck_approval_tasks_evidence_completeness", }.issubset(constraint_names) open_root_index = next( index for index in ApprovalTask.__table__.indexes if index.name == "uq_approval_tasks_open_root_per_claim" ) assert open_root_index.unique is True predicate = str(open_root_index.dialect_options["postgresql"]["where"]) assert "task_kind = 'root'" in predicate assert "status IN ('waiting', 'pending')" in predicate event_constraint_names = { constraint.name for constraint in ApprovalTaskEvent.__table__.constraints } assert { "uq_approval_task_events_actor_request", "uq_approval_task_events_task_version", "fk_approval_task_events_tenant_task", "ck_approval_task_events_version", }.issubset(event_constraint_names) assert ApprovalTaskEvent.__table__.c.response_json.nullable is False def test_risk_waiver_model_declares_decision_metadata_constraints() -> None: constraint_names = { constraint.name for constraint in RiskDisposition.__table__.constraints } assert { "ck_risk_dispositions_waiver_request", "ck_risk_dispositions_waiver_decision", "ck_risk_dispositions_waiver_lifecycle", }.issubset(constraint_names) lifecycle_constraint = next( constraint for constraint in RiskDisposition.__table__.constraints if constraint.name == "ck_risk_dispositions_lifecycle" ) assert "waived" in str(lifecycle_constraint.sqltext) assert "waiver_rejected" in str(lifecycle_constraint.sqltext) waiver_index = next( index for index in RiskDisposition.__table__.indexes if index.name == "ix_risk_dispositions_tenant_waiver_expiry" ) assert tuple(column.name for column in waiver_index.columns) == ( "tenant_id", "lifecycle_status", "waiver_expires_at", ) event_action_constraint = next( constraint for constraint in RiskDispositionEvent.__table__.constraints if constraint.name == "ck_risk_disposition_events_action" ) assert "approve_waiver" in str(event_action_constraint.sqltext) assert "reject_waiver" in str(event_action_constraint.sqltext) @pytest.mark.parametrize( "database_url", [ "sqlite+pysqlite:///:memory:", "postgresql+psycopg://postgres:postgres@x-financial-local-postgres:5432/x_financial", "postgresql+psycopg://probe:probe@migration-probe-123:5432/x_financial", ], ) def test_disposable_database_guard_rejects_unsafe_urls(database_url: str) -> None: with pytest.raises(RuntimeError): _require_disposable_probe_url(database_url) def test_alembic_migration_cycle_on_disposable_postgres( migration_database_url: str, ) -> None: engine = create_engine(migration_database_url, poolclass=NullPool) try: assert _table_names(engine) == set(), "迁移测试必须从全新空库开始" assert validate_migration_state(engine).revision is None _upgrade_revision(migration_database_url, "20260716_0008") _create_duplicate_active_organization_memory_probe(engine) with pytest.raises(RuntimeError, match="duplicate active memories"): _upgrade_head(migration_database_url) assert validate_migration_state(engine).revision == "20260716_0008" memory_columns = { str(item["name"]) for item in inspect(engine).get_columns("memory_entries", schema="public") } assert { "management_request_id", "management_payload_fingerprint", "revoke_request_id", "revoke_payload_fingerprint", }.isdisjoint(memory_columns) _delete_duplicate_active_organization_memory_probe(engine) _upgrade_revision(migration_database_url, "20260716_0019") _create_legacy_connector_payload_probe(engine) _upgrade_head(migration_database_url) _assert_and_delete_legacy_connector_payload_probe(engine) _assert_head_schema(engine) assert validate_migration_state(engine).revision == HEAD_REVISION _assert_savings_runtime_invariants(engine) _assert_commercial_runtime_invariants(engine) _assert_financial_connector_runtime_invariants(engine) _assert_release_telemetry_runtime_invariants(engine) with engine.begin() as connection: connection.execute( text( """ INSERT INTO profile_baseline_snapshots ( id, tenant_id, baseline_key, baseline_type, dimension_type, dimension_id, metric_key, unit, baseline_value, sample_count, method, query_fingerprint, data_quality_status, data_quality_score, algorithm_version, frozen_at, frozen_by ) VALUES ( 'downgrade-refusal-baseline', 'tenant-a', 'downgrade-refusal-baseline', 'manual', 'tenant', 'tenant-a', 'expense_amount', 'currency', 1, 0, 'migration_probe', 'downgrade-refusal-fingerprint', 'complete', 1, 'migration-probe-v1', now(), 'migration-probe' ) """ ) ) with pytest.raises(RuntimeError, match="immutable value facts exist"): _downgrade_revision(migration_database_url, "20260716_0014") assert validate_migration_state(engine).revision == HEAD_REVISION with engine.begin() as connection: assert connection.scalar( text( "SELECT COUNT(*) FROM profile_baseline_snapshots " "WHERE id = 'downgrade-refusal-baseline'" ) ) == 1 connection.execute( text( "DELETE FROM profile_baseline_snapshots " "WHERE id = 'downgrade-refusal-baseline'" ) ) create_legacy_schema(engine) assert "expense_claims" in _table_names(engine) _upgrade_head(migration_database_url) _assert_head_schema(engine) _assert_learning_ledger_tenant_boundary(engine) _assert_runtime_cascade(engine) _create_legacy_sentinel(engine) _downgrade_revision(migration_database_url, "20260716_0008") _create_enriched_few_shot_downgrade_probe(engine) with pytest.raises(RuntimeError, match="contain policy_ref or rule_version data"): _downgrade_revision(migration_database_url, "20260716_0007") assert validate_migration_state(engine).revision == "20260716_0008" with engine.connect() as connection: assert connection.execute( text( "SELECT policy_ref, rule_version FROM few_shot_samples " "WHERE id = 'enriched-few-shot-downgrade-probe'" ) ).one() == ("TRAVEL-001", "v2") _delete_enriched_few_shot_downgrade_probe(engine) _upgrade_head(migration_database_url) _create_hierarchical_memory_downgrade_probe(engine) _create_historical_case_downgrade_probe(engine) _downgrade_base(migration_database_url) _assert_base_schema(engine) _assert_legacy_sentinel(engine) _assert_historical_case_downgrade_probe(engine) assert validate_migration_state(engine).revision is None with engine.begin() as connection: connection.execute(text("CREATE TABLE expense_cases (id VARCHAR(36) PRIMARY KEY)")) with pytest.raises(MigrationPreflightError, match="migration-owned tables exist"): validate_migration_state(engine) with engine.begin() as connection: connection.execute(text("DROP TABLE expense_cases")) assert validate_migration_state(engine).revision is None _upgrade_head(migration_database_url) _assert_head_schema(engine) _assert_legacy_sentinel(engine) finally: engine.dispose()