Files
ModelForge/backend/tests/test_migration_engine_m13.py

582 lines
20 KiB
Python

from __future__ import annotations
import hashlib
import uuid
import pytest
from sqlalchemy import create_engine
from sqlalchemy.orm import Session
from modelforge_api.domain.migration_contracts import (
AdapterContract,
BatchReport,
CutoverPrepare,
CutoverReport,
MigrationClass,
MigrationPlanCreate,
PreflightReport,
ReconciliationReport,
RollbackReport,
ShadowReport,
StateAction,
ValidationReport,
assert_migration_transition,
)
from modelforge_api.persistence.models import (
ArtifactSet,
Base,
CapabilityDeployment,
EmbeddingSpace,
LifecycleApprovalRequest,
MigrationEvent,
MigrationPlan,
Project,
ProjectBinding,
)
from modelforge_api.services.migration_adapters import REQUIRED_REINDEX_OPERATIONS
from modelforge_api.services.migration_engine import (
MigrationEngineError,
MigrationEngineService,
)
from tests.test_lifecycle_m12 import approved_lab_plan
def digest(value: str) -> str:
return hashlib.sha256(value.encode()).hexdigest()
@pytest.fixture
def session() -> Session:
engine = create_engine("sqlite+pysqlite:///:memory:")
Base.metadata.create_all(engine)
with Session(engine) as value:
yield value
def planned_migration(
session: Session,
*,
source: str = "rag_source_v1",
target: str = "rag_shadow_v2",
migration_class: MigrationClass = MigrationClass.REQUIRES_REINDEX,
schema_steps: list[dict[str, object]] | None = None,
schema_operations: frozenset[str] = frozenset(),
) -> tuple[MigrationEngineService, object]:
_lifecycle, _subject, lifecycle_plan = approved_lab_plan(session, target=f"m13-{uuid.uuid4()}")
deployment = session.get(CapabilityDeployment, lifecycle_plan.candidate_deployment_id)
approval = session.get(LifecycleApprovalRequest, lifecycle_plan.approval_request_id)
assert deployment is not None and approval is not None
artifact_set = session.get(ArtifactSet, deployment.artifact_set_id)
assert artifact_set is not None
project = Project(
key=f"m13-project-{uuid.uuid4()}",
name="M13 isolated rehearsal",
description="M13 test project",
active=True,
)
session.add(project)
session.flush()
binding = ProjectBinding(
project_id=project.id,
capability_contract_id=deployment.capability_contract_id,
channel="experiment",
priority="background",
optional=False,
fallback_policy={"mode": "fail_closed"},
migration_support="reindex",
slo={},
benchmark_requirements=["retrieval"],
)
assert deployment.embedding_space_id is not None
target_space = session.get(EmbeddingSpace, deployment.embedding_space_id)
assert target_space is not None
session.add(binding)
session.commit()
service = MigrationEngineService(session)
service.ensure_defaults()
policy = next(
item for item in service.validation_policies() if item.key == "isolated-lab-rehearsal"
)
adapter = AdapterContract(
key="examplerag.qdrant-reindex",
version="1",
operations=REQUIRED_REINDEX_OPERATIONS,
fingerprint=digest("examplerag.qdrant-reindex@1"),
schema_operations=schema_operations,
)
request = MigrationPlanCreate(
project_id=project.id,
project_binding_id=binding.id,
capability_contract_id=deployment.capability_contract_id,
migration_class=migration_class,
environment="LAB",
adapter=adapter,
source_identity={
"fingerprint": digest(source),
"collection": source,
"embedding_space_id": "source-space-v1",
},
target_identity={
"fingerprint": digest(target),
"collection": target,
"embedding_space_id": str(target_space.id),
"capability_deployment_id": str(deployment.id),
"artifact_set_id": str(deployment.artifact_set_id),
"runtime_profile_id": str(deployment.runtime_profile_id),
"model_revision_id": str(artifact_set.revision_id),
"dimension": 1024,
"distance_metric": "COSINE",
"document_semantics": {"prefix": ""},
"query_semantics": {"prefix": "query"},
},
source_data_target=source,
target_shadow_target=target,
source_space_ref="embedding-space-v1",
target_space_id=target_space.id,
corpus_revision="corpus-exact-1",
migration_policy_revision="m13-policy-1",
validation_policy_revision_id=policy.id,
lifecycle_approval_id=approval.id,
rollback_target_ref=source,
total_expected_items=4,
batch_size=2,
max_in_flight_batches=1,
concurrency=1,
target_storage={"backend": "qdrant", "collection": target},
shadow_policy={"minimum_queries": 2},
cutover_policy={"maximum_error_rate": 0.0},
rollback_retention_days=30,
environment_fingerprint=digest("isolated-test-environment"),
idempotency_key=f"m13-plan-{uuid.uuid4()}",
created_by="test-operator",
irreversible=migration_class is MigrationClass.SCHEMA_BREAKING,
schema_steps=schema_steps or [],
)
return service, service.create_plan(request)
def preflight(service: MigrationEngineService, plan: object) -> object:
return service.preflight(
plan.id,
PreflightReport(
expected_version=plan.version,
adapter_fingerprint=plan.adapter.fingerprint,
source_fingerprint=plan.source_identity["fingerprint"],
source_exists=True,
source_healthy=True,
source_count=plan.total_expected_items,
target_conflict_free=True,
target_space_valid=True,
capability_healthy=True,
project_credential_valid=True,
storage_sufficient=True,
scheduler_capacity=True,
adapter_available=True,
rollback_source_retained=True,
evaluation_suite_available=True,
lifecycle_approval_current=True,
evidence={"probe": "exact"},
),
)
def batch(
service: MigrationEngineService, plan: object, number: int, *, retryable: bool = False
) -> object:
return service.record_batch(
plan.id,
BatchReport(
expected_version=plan.version,
generation=plan.generation,
batch_number=number,
cursor_start=str(number * 2),
cursor_end=str(number * 2 + 2),
item_count=2,
completed_items=0 if retryable else 2,
failed_items=2 if retryable else 0,
retryable_items=2 if retryable else 0,
permanent_failed_items=0,
item_fingerprint=digest(f"items-{number}"),
result_fingerprint=digest(f"result-{number}-{'retry' if retryable else 'ok'}"),
output_shape_valid=True,
finite=True,
target_space_matches=True,
destination_committed=not retryable,
content_hashes_match=True,
duration_ms=10.0,
error_code="WRITE_FAILED" if retryable else None,
bounded_errors=[{"code": "temporary"}] if retryable else [],
),
)
def complete_backfill(service: MigrationEngineService, plan: object) -> object:
plan = preflight(service, plan)
plan = service.start_backfill(
plan.id,
StateAction(expected_version=plan.version, actor="operator", reason="begin"),
)
batch(service, plan, 0)
plan = service.plan(plan.id)
batch(service, plan, 1)
return service.plan(plan.id)
def validate_and_shadow(service: MigrationEngineService, plan: object) -> object:
validation = service.validate(
plan.id,
ValidationReport(
expected_version=plan.version,
generation=plan.generation,
expected_count=4,
actual_count=4,
missing_count=0,
duplicate_count=0,
malformed_count=0,
non_finite_count=0,
wrong_dimension_count=0,
content_hash_mismatch_count=0,
wrong_space_count=0,
index_schema_matches=True,
distance_metric_matches=True,
payload_integrity=True,
target_fingerprint=digest("target-validated"),
evaluation_run_ids=[uuid.uuid4()],
comparable=True,
critical_regressions=0,
latency_regression_ratio=0.9,
project_fit_eligible=False,
external_validation_satisfied=False,
security_approved=True,
evidence={"suite": "isolated"},
),
)
assert validation.technical_cutover_eligible is True
assert validation.project_promotion_eligible is False
plan = service.plan(plan.id)
plan = service.start_shadow(
plan.id,
StateAction(expected_version=plan.version, actor="operator", reason="shadow"),
)
return service.complete_shadow(
plan.id,
ShadowReport(
expected_version=plan.version,
generation=plan.generation,
request_count=4,
source_error_count=0,
target_error_count=0,
source_latency_p95_ms=20,
target_latency_p95_ms=19,
critical_regressions=0,
metrics={"recall_delta": 0.0},
evidence_refs=["evaluation:isolated"],
result_fingerprint=digest(f"shadow-{plan.id}"),
),
)
def ready_for_cutover(service: MigrationEngineService, plan: object) -> object:
return validate_and_shadow(service, complete_backfill(service, plan))
def test_transition_graph_rejects_shortcuts() -> None:
assert_migration_transition("READY", "BACKFILLING")
with pytest.raises(ValueError, match="invalid migration transition"):
assert_migration_transition("PLANNED", "CUTOVER_COMMITTED")
@pytest.mark.parametrize("migration_class", [MigrationClass.TRANSPARENT, MigrationClass.BEHAVIORAL])
def test_embedding_space_change_requires_reindex(
session: Session, migration_class: MigrationClass
) -> None:
with pytest.raises(MigrationEngineError) as raised:
planned_migration(session, migration_class=migration_class)
assert raised.value.code == "REINDEX_REQUIRED"
def test_plan_is_idempotent_and_immutable(session: Session) -> None:
service, plan = planned_migration(session)
stored = session.get(MigrationPlan, plan.id)
assert stored is not None
with pytest.raises(ValueError, match="immutable"):
stored.target_shadow_target = "other-target"
session.commit()
session.rollback()
assert service.plan(plan.id).target_shadow_target == "rag_shadow_v2"
def test_preflight_fails_closed_on_changed_source(session: Session) -> None:
service, plan = planned_migration(session)
report = PreflightReport(
expected_version=plan.version,
adapter_fingerprint=plan.adapter.fingerprint,
source_fingerprint=digest("unexpected-source"),
source_exists=True,
source_healthy=True,
source_count=4,
target_conflict_free=True,
target_space_valid=True,
capability_healthy=True,
project_credential_valid=True,
storage_sufficient=True,
scheduler_capacity=True,
adapter_available=True,
rollback_source_retained=True,
evaluation_suite_available=True,
lifecycle_approval_current=True,
)
assert service.preflight(plan.id, report).failure_code == "SOURCE_CHANGED"
def test_backfill_pause_resume_retry_and_idempotency(session: Session) -> None:
service, plan = planned_migration(session)
plan = preflight(service, plan)
plan = service.start_backfill(
plan.id, StateAction(expected_version=plan.version, actor="worker", reason="start")
)
first = batch(service, plan, 0, retryable=True)
assert first.status == "RETRYABLE_FAILED"
plan = service.plan(plan.id)
plan = service.pause_backfill(
plan.id, StateAction(expected_version=plan.version, actor="operator", reason="interrupt")
)
plan = service.start_backfill(
plan.id, StateAction(expected_version=plan.version, actor="operator", reason="resume")
)
completed = batch(service, plan, 0)
replay = service.record_batch(
plan.id,
BatchReport(
expected_version=plan.version,
generation=plan.generation,
batch_number=0,
cursor_start="0",
cursor_end="2",
item_count=2,
completed_items=2,
failed_items=0,
retryable_items=0,
permanent_failed_items=0,
item_fingerprint=digest("items-0"),
result_fingerprint=digest("result-0-ok"),
output_shape_valid=True,
finite=True,
target_space_matches=True,
destination_committed=True,
content_hashes_match=True,
duration_ms=10,
),
)
assert replay.id == completed.id
current = service.plan(plan.id)
assert current.completed_items == 2
assert current.retryable_items == 0
def test_validation_keeps_lab_technical_and_project_eligibility_separate(
session: Session,
) -> None:
service, plan = planned_migration(session)
validation = service.validate(
plan.id,
ValidationReport(
expected_version=complete_backfill(service, plan).version,
generation=1,
expected_count=4,
actual_count=4,
missing_count=0,
duplicate_count=0,
malformed_count=0,
non_finite_count=0,
wrong_dimension_count=0,
content_hash_mismatch_count=0,
wrong_space_count=0,
index_schema_matches=True,
distance_metric_matches=True,
payload_integrity=True,
target_fingerprint=digest("validated"),
evaluation_run_ids=[uuid.uuid4()],
comparable=True,
critical_regressions=1,
latency_regression_ratio=None,
project_fit_eligible=False,
external_validation_satisfied=False,
security_approved=True,
),
)
assert validation.technical_cutover_eligible is True
assert validation.project_promotion_eligible is False
assert "PROJECT_FIT_NOT_ELIGIBLE" in validation.blockers
def test_atomic_cutover_and_exact_rollback(session: Session) -> None:
service, plan = planned_migration(session)
plan = ready_for_cutover(service, plan)
operation = service.prepare_cutover(
plan.id,
CutoverPrepare(
expected_version=plan.version,
idempotency_key=f"cutover-{uuid.uuid4()}",
actor="operator",
expected_external_source=plan.source_data_target,
observed_external_source=plan.source_data_target,
external_state_fingerprint=digest("before"),
configuration_version="cfg-1",
),
)
plan = service.plan(plan.id)
operation = service.report_cutover(
plan.id,
CutoverReport(
expected_version=plan.version,
operation_id=operation.id,
generation=plan.generation,
external_source_before=plan.source_data_target,
external_target_after=plan.target_shadow_target,
external_state_fingerprint=digest("after"),
switch_duration_ms=3.5,
target_reachable=True,
expected_identity=True,
capability_healthy=True,
project_read_path_healthy=True,
error_rate=0,
smoke_query_count=2,
smoke_error_count=0,
evidence={"alias": "verified"},
),
)
assert operation.stage == "COMMITTED"
plan = service.plan(plan.id)
operation = service.rollback(
plan.id,
RollbackReport(
expected_version=plan.version,
operation_id=operation.id,
generation=plan.generation,
restored_external_target=plan.source_data_target,
external_state_fingerprint=digest("restored"),
elapsed_ms=2.0,
source_reachable=True,
exact_identity_restored=True,
capability_healthy=True,
project_read_path_healthy=True,
evidence={"alias": "restored"},
),
)
assert operation.stage == "ROLLED_BACK"
assert service.plan(plan.id).state == "ROLLED_BACK"
def test_stale_source_is_rejected_before_cutover(session: Session) -> None:
service, plan = planned_migration(session)
plan = ready_for_cutover(service, plan)
with pytest.raises(MigrationEngineError, match="external source changed") as raised:
service.prepare_cutover(
plan.id,
CutoverPrepare(
expected_version=plan.version,
idempotency_key=f"cutover-{uuid.uuid4()}",
actor="operator",
expected_external_source=plan.source_data_target,
observed_external_source="concurrent-writer-target",
external_state_fingerprint=digest("stale"),
configuration_version="cfg-2",
),
)
assert raised.value.code == "STALE_SOURCE"
def test_crash_reconciliation_uses_observed_external_truth(session: Session) -> None:
service, plan = planned_migration(session)
plan = ready_for_cutover(service, plan)
operation = service.prepare_cutover(
plan.id,
CutoverPrepare(
expected_version=plan.version,
idempotency_key=f"cutover-{uuid.uuid4()}",
actor="operator",
expected_external_source=plan.source_data_target,
observed_external_source=plan.source_data_target,
external_state_fingerprint=digest("before-crash"),
configuration_version="cfg-crash",
),
)
assert service.pending_reconciliation_count() == 1
operation = service.reconcile(
ReconciliationReport(
operation_id=operation.id,
generation=plan.generation,
observed_external_target=plan.target_shadow_target,
external_state_fingerprint=digest("observed-after-restart"),
target_healthy=True,
source_healthy=True,
switch_duration_ms=17.5,
smoke_query_count=3,
smoke_error_count=0,
evidence={"adapter": "isolated-test"},
)
)
assert operation.stage == "COMMITTED"
assert operation.switch_duration_ms == 17.5
assert operation.health_evidence["smoke_query_count"] == 3
assert service.pending_reconciliation_count() == 0
def test_schema_breaking_plan_cannot_inject_arbitrary_execution(session: Session) -> None:
with pytest.raises(MigrationEngineError) as raised:
planned_migration(
session,
migration_class=MigrationClass.SCHEMA_BREAKING,
schema_steps=[{"operation": "shell"}],
)
assert raised.value.code == "ARBITRARY_EXECUTION_DENIED"
def test_schema_breaking_plan_stops_at_typed_manual_boundary(session: Session) -> None:
step_key = "catalog.expand-v2"
service, plan = planned_migration(
session,
migration_class=MigrationClass.SCHEMA_BREAKING,
schema_operations=frozenset({step_key}),
schema_steps=[
{
"operation": "expand",
"adapter_step": step_key,
"preconditions": ["schema-v1-present"],
"required_application_versions": {"catalog": ">=2.0"},
"compatibility_window": "v1-v2-dual-read",
"rollback_feasible": False,
"irreversible": True,
}
],
)
assert plan.state == "PLANNED"
assert plan.irreversible is True
with pytest.raises(MigrationEngineError) as raised:
service.prepare_cutover(
plan.id,
CutoverPrepare(
expected_version=plan.version,
idempotency_key=f"schema-cutover-{uuid.uuid4()}",
actor="operator",
expected_external_source=plan.source_data_target,
observed_external_source=plan.source_data_target,
external_state_fingerprint=digest("schema-boundary"),
configuration_version="schema-v1",
),
)
assert raised.value.code == "SCHEMA_BREAKING_AUTO_EXECUTION_DENIED"
def test_migration_events_are_append_only(session: Session) -> None:
service, plan = planned_migration(session)
event = service.events(plan.id, 1)[0]
stored = session.get(MigrationEvent, event.id)
assert stored is not None
with pytest.raises(ValueError, match="append-only"):
stored.reason = "rewritten"
session.commit()