from __future__ import annotations import shutil from datetime import timedelta from pathlib import Path import pytest from pydantic import ValidationError from app.core.time import from_iso, to_iso, utc_now from app.database.database import Database from app.fleet.repository import FleetRepository from app.fleet.schemas import ManagedInstanceCreate, ManagedInstanceUpdate MIGRATIONS_PATH = Path(__file__).resolve().parents[1] / "migrations" @pytest.fixture def database(tmp_path: Path) -> Database: result = Database(tmp_path / "fleet.db", MIGRATIONS_PATH) result.migrate() return result @pytest.fixture def repository(database: Database) -> FleetRepository: return FleetRepository(database) def instance_values( instance_id: str, *, display_name: str | None = None, region: str = "us-east-1", lightsail_name: str | None = None, record_name: str | None = None, enabled: bool = True, ) -> dict[str, object]: suffix = instance_id.replace("_", "-") return { "id": instance_id, "display_name": display_name or f"Proxy {suffix}", "aws_region": region, "lightsail_instance_name": lightsail_name or f"node-{suffix}", "cloudflare_zone_name": "example.com", "cloudflare_zone_id": f"zone-{suffix}", "cloudflare_record_name": record_name or f"{suffix}.example.com", "socks_port": 1080, "proxy_health_check": True, "health_timeout_seconds": 120, "release_grace_seconds": 75, "enabled": enabled, } def test_migration_creates_fleet_contract_and_validates_schema(database: Database) -> None: with database.connect() as connection: tables = { row["name"] for row in connection.execute( "SELECT name FROM sqlite_master WHERE type = 'table'" ).fetchall() } item_columns = { row["name"]: row for row in connection.execute("PRAGMA table_info(fleet_run_items)") } instance_columns = { row["name"]: row for row in connection.execute("PRAGMA table_info(managed_instances)") } assert { "managed_instances", "instance_groups", "instance_group_members", "fleet_runs", "fleet_run_items", "fleet_events", "fleet_operation_locks", "managed_instance_secrets", } <= tables assert item_columns["stage_started_at"]["notnull"] == 1 assert item_columns["attempt_count"]["dflt_value"] == "0" assert instance_columns["release_grace_seconds"]["dflt_value"] == "75" assert instance_columns["proxy_protocol"]["dflt_value"] == "'socks5'" assert instance_columns["proxy_auth_configured"]["dflt_value"] == "0" assert instance_columns["cloudflare_proxied"]["dflt_value"] == "0" assert item_columns["proxy_protocol"]["dflt_value"] == "'socks5'" assert item_columns["cloudflare_proxied"]["dflt_value"] == "0" normalized = ManagedInstanceCreate.model_validate( { **instance_values("schema"), "cloudflare_zone_name": "Example.COM.", "cloudflare_record_name": "S5.Example.COM.", } ) assert normalized.cloudflare_zone_name == "example.com" assert normalized.cloudflare_record_name == "s5.example.com" assert ( ManagedInstanceCreate.model_validate( {**instance_values("group-trim"), "group_id": " primary "} ).group_id == "primary" ) assert ( ManagedInstanceCreate.model_validate( {**instance_values("group-blank"), "group_id": " "} ).group_id is None ) normalized_update = ManagedInstanceUpdate.model_validate( {"config_version": 1, "group_id": " secondary "} ) assert normalized_update.group_id == "secondary" detached_update = ManagedInstanceUpdate.model_validate({"config_version": 1, "group_id": " "}) assert detached_update.group_id is None assert detached_update.model_dump(exclude_unset=True)["group_id"] is None with pytest.raises(ValidationError): ManagedInstanceCreate.model_validate( {**instance_values("bad-grace"), "release_grace_seconds": 59} ) with pytest.raises(ValidationError): ManagedInstanceCreate.model_validate( { **instance_values("bad-zone"), "cloudflare_record_name": "proxy.example.net", } ) def test_instance_crud_partial_uniqueness_status_and_soft_archive( repository: FleetRepository, ) -> None: created = repository.create_instance(instance_values("one")) assert created.id == "one" assert created.config_version == 1 assert created.release_grace_seconds == 75 with pytest.raises(RuntimeError, match="INSTANCE_DISPLAY_NAME_CONFLICT"): repository.create_instance( instance_values( "duplicate-name", display_name="proxy one", lightsail_name="other-node", record_name="other.example.com", ) ) with pytest.raises(RuntimeError, match="INSTANCE_AWS_TARGET_CONFLICT"): repository.create_instance( instance_values( "duplicate-aws", display_name="Other display", lightsail_name="node-one", record_name="third.example.com", ) ) with pytest.raises(RuntimeError, match="INSTANCE_DNS_RECORD_CONFLICT"): repository.create_instance( instance_values( "duplicate-dns", display_name="DNS duplicate", lightsail_name="dns-node", record_name="ONE.EXAMPLE.COM", ) ) checked_at = to_iso() status = repository.update_instance_status( "one", {"last_known_ip": "198.51.100.5", "last_checked_at": checked_at} ) assert status.last_known_ip == "198.51.100.5" assert status.last_checked_at == checked_at assert status.config_version == 1 updated = repository.update_instance("one", {"socks_port": 2080}) assert updated.socks_port == 2080 assert updated.config_version == 2 with pytest.raises(RuntimeError, match="INSTANCE_CONFIG_VERSION_CONFLICT"): repository.update_instance("one", {"config_version": 1, "socks_port": 3080}) assert repository.get_instance("one").socks_port == 2080 # type: ignore[union-attr] archived = repository.archive_instance("one") assert archived.archived_at is not None assert archived.enabled is False assert repository.get_instance("one") is None assert repository.get_instance("one", include_archived=True) == archived replacement = repository.create_instance( instance_values( "replacement", display_name="Proxy one", lightsail_name="node-one", record_name="ONE.EXAMPLE.COM", ) ) assert replacement.id == "replacement" assert [item.id for item in repository.list_instances()] == ["replacement"] def test_group_members_are_atomic_unique_and_schedule_changes_are_explicit( repository: FleetRepository, ) -> None: first = repository.create_instance(instance_values("first")) second = repository.create_instance(instance_values("second", enabled=False)) group = repository.create_group( { "id": "primary", "name": "Primary", "enabled": True, "interval_minutes": 30, "member_ids": [first.id, second.id], } ) assert group.member_ids == [first.id, second.id] assert group.next_run_at is not None original_next = group.next_run_at renamed = repository.update_group("primary", {"name": "Primary proxies"}) assert renamed.next_run_at == original_next assert renamed.config_version == 2 interval_changed = repository.update_group("primary", {"interval_minutes": 45}) assert interval_changed.next_run_at is not None assert interval_changed.next_run_at != original_next seconds_until_due = (from_iso(interval_changed.next_run_at) - utc_now()).total_seconds() assert 44 * 60 < seconds_until_due <= 45 * 60 with pytest.raises(RuntimeError, match="GROUP_CONFIG_VERSION_CONFLICT"): repository.update_group( "primary", { "config_version": 1, "name": "Stale update", "member_ids": [first.id], }, ) unchanged = repository.get_group("primary") assert unchanged is not None assert unchanged.name == "Primary proxies" assert unchanged.member_ids == [first.id, second.id] reordered = repository.update_group( "primary", { "name": "Primary proxies", "enabled": True, "interval_minutes": 45, "member_ids": [second.id, first.id], }, ) assert reordered.member_ids == [second.id, first.id] assert reordered.next_run_at == interval_changed.next_run_at other_group = repository.create_group( {"id": "secondary", "name": "Secondary", "enabled": False} ) with pytest.raises(RuntimeError, match="INSTANCE_ALREADY_GROUPED"): repository.update_group( other_group.id, { "name": "Should roll back", "enabled": False, "interval_minutes": 60, "member_ids": [first.id], }, ) assert repository.get_group(other_group.id).member_ids == [] # type: ignore[union-attr] assert repository.get_group(other_group.id).name == "Secondary" # type: ignore[union-attr] with pytest.raises(RuntimeError, match="GROUP_REQUIRES_ENABLED_MEMBER"): repository.update_group( "primary", { "name": "Must also roll back", "enabled": True, "interval_minutes": 45, "member_ids": [second.id], }, ) assert repository.get_group("primary").member_ids == [ # type: ignore[union-attr] second.id, first.id, ] assert repository.get_group("primary").name == "Primary proxies" # type: ignore[union-attr] with pytest.raises(RuntimeError, match="GROUP_REQUIRES_ENABLED_MEMBER"): repository.update_instance(first.id, {"enabled": False}) disabled = repository.update_group("primary", {"enabled": False}) assert disabled.next_run_at is None renamed_disabled = repository.update_group("primary", {"name": "Dormant"}) assert renamed_disabled.next_run_at is None enabled_again = repository.update_group("primary", {"enabled": True}) assert enabled_again.next_run_at is not None def test_group_entry_points_bump_changed_instance_versions_and_preserve_cross_cas( repository: FleetRepository, ) -> None: first = repository.create_instance(instance_values("cross-cas-first")) second = repository.create_instance(instance_values("cross-cas-second")) third = repository.create_instance(instance_values("cross-cas-third")) fourth = repository.create_instance(instance_values("cross-cas-fourth")) target = repository.create_group( {"id": "cross-cas-target", "name": "Cross CAS target", "enabled": False} ) group = repository.create_group( { "id": "cross-cas-group", "name": "Cross CAS group", "enabled": False, "member_ids": [first.id, second.id], } ) first_after_create = repository.get_instance(first.id) second_after_create = repository.get_instance(second.id) assert first_after_create is not None assert second_after_create is not None assert first_after_create.config_version == first.config_version + 1 assert second_after_create.config_version == second.config_version + 1 assert repository.get_instance(third.id) == third assert repository.get_instance(fourth.id) == fourth reordered = repository.update_group( group.id, { "config_version": group.config_version, "member_ids": [second.id, first.id], }, ) assert reordered.member_ids == [second.id, first.id] assert repository.get_instance(first.id) == first_after_create assert repository.get_instance(second.id) == second_after_create changed = repository.update_group( group.id, { "config_version": reordered.config_version, "member_ids": [second.id, third.id], }, ) first_after_remove = repository.get_instance(first.id) second_retained = repository.get_instance(second.id) third_after_add = repository.get_instance(third.id) assert changed.member_ids == [second.id, third.id] assert first_after_remove is not None assert first_after_remove.config_version == first_after_create.config_version + 1 assert second_retained == second_after_create assert third_after_add is not None assert third_after_add.config_version == third.config_version + 1 with pytest.raises(RuntimeError, match="INSTANCE_CONFIG_VERSION_CONFLICT"): repository.update_instance( first.id, { "config_version": first_after_create.config_version, "group_id": target.id, }, ) assert repository.get_instance(first.id) == first_after_remove assert repository.get_group(group.id) == changed assert repository.get_group(target.id) == target saved = repository.save_group_members(group.id, [third.id, fourth.id]) second_after_remove = repository.get_instance(second.id) third_retained = repository.get_instance(third.id) fourth_after_add = repository.get_instance(fourth.id) assert saved.member_ids == [third.id, fourth.id] assert second_after_remove is not None assert second_after_remove.config_version == second_after_create.config_version + 1 assert third_retained == third_after_add assert fourth_after_add is not None assert fourth_after_add.config_version == fourth.config_version + 1 saved_reordered = repository.save_group_members(group.id, [fourth.id, third.id]) assert saved_reordered.member_ids == [fourth.id, third.id] assert repository.get_instance(third.id) == third_retained assert repository.get_instance(fourth.id) == fourth_after_add archived = repository.archive_group(group.id) third_after_archive = repository.get_instance(third.id) fourth_after_archive = repository.get_instance(fourth.id) assert archived.member_ids == [] assert third_after_archive is not None assert third_after_archive.config_version == third_retained.config_version + 1 assert fourth_after_archive is not None assert fourth_after_archive.config_version == fourth_after_add.config_version + 1 def test_instance_writes_assign_move_and_remove_group_with_version_updates( repository: FleetRepository, ) -> None: first_group = repository.create_group( {"id": "assignment-first", "name": "Assignment first", "enabled": False} ) second_group = repository.create_group( {"id": "assignment-second", "name": "Assignment second", "enabled": False} ) instance = repository.create_instance( {**instance_values("assigned-on-create"), "group_id": first_group.id} ) after_create = repository.get_group(first_group.id) assert after_create is not None assert after_create.member_ids == [instance.id] assert after_create.config_version == first_group.config_version + 1 detached = repository.update_instance( instance.id, {"config_version": instance.config_version, "group_id": None}, ) first_after_detach = repository.get_group(first_group.id) assert detached.config_version == instance.config_version + 1 assert first_after_detach is not None assert first_after_detach.member_ids == [] assert first_after_detach.config_version == after_create.config_version + 1 joined = repository.update_instance( instance.id, {"config_version": detached.config_version, "group_id": first_group.id}, ) first_after_join = repository.get_group(first_group.id) assert joined.config_version == detached.config_version + 1 assert first_after_join is not None assert first_after_join.member_ids == [instance.id] assert first_after_join.config_version == first_after_detach.config_version + 1 moved = repository.update_instance( instance.id, {"config_version": joined.config_version, "group_id": second_group.id}, ) first_after_move = repository.get_group(first_group.id) second_after_move = repository.get_group(second_group.id) assert moved.config_version == joined.config_version + 1 assert first_after_move is not None assert first_after_move.member_ids == [] assert first_after_move.config_version == first_after_join.config_version + 1 assert second_after_move is not None assert second_after_move.member_ids == [instance.id] assert second_after_move.config_version == second_group.config_version + 1 removed = repository.update_instance( instance.id, {"config_version": moved.config_version, "group_id": None}, ) second_after_remove = repository.get_group(second_group.id) assert removed.config_version == moved.config_version + 1 assert removed.socks_port == instance.socks_port assert second_after_remove is not None assert second_after_remove.member_ids == [] assert second_after_remove.config_version == second_after_move.config_version + 1 def test_instance_group_assignment_rejects_invalid_archived_and_stale_targets_atomically( repository: FleetRepository, ) -> None: with pytest.raises(RuntimeError, match="GROUP_NOT_FOUND"): repository.create_instance( {**instance_values("missing-group-create"), "group_id": "missing-group"} ) assert repository.get_instance("missing-group-create") is None archived_group = repository.create_group( {"id": "archived-target", "name": "Archived target", "enabled": False} ) repository.archive_group(archived_group.id) with pytest.raises(RuntimeError, match="GROUP_NOT_FOUND"): repository.create_instance( {**instance_values("archived-group-create"), "group_id": archived_group.id} ) assert repository.get_instance("archived-group-create") is None source = repository.create_group( {"id": "atomic-source", "name": "Atomic source", "enabled": False} ) target = repository.create_group( {"id": "atomic-target", "name": "Atomic target", "enabled": False} ) instance = repository.create_instance( {**instance_values("atomic-member"), "group_id": source.id} ) instance_before = repository.get_instance(instance.id) source_before = repository.get_group(source.id) target_before = repository.get_group(target.id) with pytest.raises(RuntimeError, match="INSTANCE_CONFIG_VERSION_CONFLICT"): repository.update_instance( instance.id, {"config_version": instance.config_version - 1, "group_id": target.id}, ) with pytest.raises(RuntimeError, match="GROUP_NOT_FOUND"): repository.update_instance( instance.id, {"config_version": instance.config_version, "group_id": "missing-group"}, ) with pytest.raises(RuntimeError, match="GROUP_NOT_FOUND"): repository.update_instance( instance.id, {"config_version": instance.config_version, "group_id": archived_group.id}, ) assert repository.get_instance(instance.id) == instance_before assert repository.get_group(source.id) == source_before assert repository.get_group(target.id) == target_before def test_moving_last_enabled_member_out_of_enabled_group_rolls_back( repository: FleetRepository, ) -> None: instance = repository.create_instance(instance_values("last-enabled-member")) source = repository.create_group( { "id": "enabled-source", "name": "Enabled source", "enabled": True, "member_ids": [instance.id], } ) target = repository.create_group( {"id": "disabled-target", "name": "Disabled target", "enabled": False} ) instance_before = repository.get_instance(instance.id) source_before = repository.get_group(source.id) target_before = repository.get_group(target.id) for group_id in (target.id, None): with pytest.raises(RuntimeError, match="GROUP_REQUIRES_ENABLED_MEMBER"): repository.update_instance( instance.id, { "config_version": instance_before.config_version, "group_id": group_id, }, ) assert repository.get_instance(instance.id) == instance_before assert repository.get_group(source.id) == source_before assert repository.get_group(target.id) == target_before def test_archiving_removes_members_without_deleting_instances_or_cloud_resources( repository: FleetRepository, ) -> None: instance = repository.create_instance(instance_values("archive-member")) group = repository.create_group( { "id": "archive-group", "name": "Archive group", "enabled": True, "member_ids": [instance.id], } ) archived_instance = repository.archive_instance(instance.id) current_group = repository.get_group(group.id) assert archived_instance.archived_at is not None assert current_group is not None assert current_group.member_ids == [] assert current_group.enabled is False assert current_group.next_run_at is None new_instance = repository.create_instance(instance_values("new-member")) repository.save_group_members(group.id, [new_instance.id]) archived_group = repository.archive_group(group.id) assert archived_group.archived_at is not None assert archived_group.member_ids == [] assert repository.get_group(group.id) is None assert repository.get_instance(new_instance.id) is not None reused = repository.create_group( {"id": "replacement-group", "name": "Archive group", "enabled": False} ) assert reused.id == "replacement-group" def test_archiving_instance_updates_remaining_group_members_and_version( repository: FleetRepository, ) -> None: first = repository.create_instance(instance_values("archive-first")) second = repository.create_instance(instance_values("archive-second")) group = repository.create_group( { "id": "archive-with-remaining-member", "name": "Archive with remaining member", "enabled": True, "member_ids": [first.id, second.id], } ) repository.archive_instance(first.id) updated_group = repository.get_group(group.id) assert updated_group is not None assert updated_group.member_ids == [second.id] assert updated_group.config_version == group.config_version + 1 assert updated_group.enabled is True assert updated_group.next_run_at == group.next_run_at def test_dns_sync_lock_blocks_external_writes_allows_owner_and_expires( database: Database, repository: FleetRepository, ) -> None: instance = repository.create_instance(instance_values("dns-locked")) group = repository.create_group( { "id": "dns-locked-group", "name": "DNS locked group", "enabled": True, "member_ids": [instance.id], } ) instance_after_group_create = repository.get_instance(instance.id) assert instance_after_group_create is not None owner_id = "dns-sync-owner" now = utc_now() with database.connect() as connection, connection: connection.execute( """ INSERT INTO fleet_operation_locks( slot, kind, owner_id, lease_until, created_at ) VALUES (1, 'dns_sync', ?, ?, ?) """, ( owner_id, to_iso(now + timedelta(minutes=5)), to_iso(now), ), ) blocked_writes = ( lambda: repository.create_instance(instance_values("blocked-by-dns")), lambda: repository.update_instance(instance.id, {"socks_port": 2080}), lambda: repository.update_instance_status(instance.id, {"last_checked_at": to_iso()}), lambda: repository.archive_instance(instance.id), lambda: repository.create_group({"name": "Blocked group", "enabled": False}), lambda: repository.update_group(group.id, {"name": "Blocked rename"}), lambda: repository.save_group_members(group.id, [instance.id]), lambda: repository.archive_group(group.id), ) for write in blocked_writes: with pytest.raises(RuntimeError, match="FLEET_RUN_ACTIVE"): write() status = repository.update_instance_status( instance.id, {"last_known_ip": "198.51.100.44", "last_checked_at": to_iso()}, operation_owner_id=owner_id, ) updated = repository.update_instance( instance.id, {"cloudflare_zone_id": "resolved-zone-id"}, operation_owner_id=owner_id, ) assert status.last_known_ip == "198.51.100.44" assert updated.cloudflare_zone_id == "resolved-zone-id" assert updated.config_version == instance_after_group_create.config_version + 1 with database.connect() as connection, connection: connection.execute( "UPDATE fleet_operation_locks SET lease_until = ? WHERE owner_id = ?", (to_iso(utc_now() - timedelta(seconds=1)), owner_id), ) created_after_expiry = repository.create_instance(instance_values("after-expiry")) assert created_after_expiry.id == "after-expiry" with database.connect() as connection: lock = connection.execute( "SELECT 1 FROM fleet_operation_locks WHERE owner_id = ?", (owner_id,) ).fetchone() assert lock is None def test_all_repository_writes_are_blocked_by_active_fleet_run( database: Database, repository: FleetRepository, ) -> None: instance = repository.create_instance(instance_values("locked")) now = to_iso() with database.connect() as connection, connection: connection.execute( """ INSERT INTO fleet_runs( id, target_type, target_id, target_name, trigger, status, active_slot, total_items, succeeded_items, started_at, updated_at ) VALUES ('active', 'instance', ?, ?, 'manual', 'running', 1, 1, 0, ?, ?) """, (instance.id, instance.display_name, now, now), ) assert repository.get_instance(instance.id) is not None with pytest.raises(RuntimeError, match="FLEET_RUN_ACTIVE"): repository.create_instance(instance_values("blocked-create")) with pytest.raises(RuntimeError, match="FLEET_RUN_ACTIVE"): repository.update_instance(instance.id, {"socks_port": 2080}) with pytest.raises(RuntimeError, match="FLEET_RUN_ACTIVE"): repository.update_instance_status(instance.id, {"last_checked_at": now}) with pytest.raises(RuntimeError, match="FLEET_RUN_ACTIVE"): repository.create_group({"name": "Blocked", "enabled": False}) def test_run_item_records_include_recovery_timing_fields( database: Database, repository: FleetRepository, ) -> None: instance = repository.create_instance(instance_values("run-item")) now = to_iso() with database.connect() as connection, connection: connection.execute( """ INSERT INTO fleet_runs( id, target_type, target_id, target_name, trigger, status, total_items, succeeded_items, started_at, updated_at ) VALUES ('run', 'instance', ?, ?, 'manual', 'queued', 1, 0, ?, ?) """, (instance.id, instance.display_name, now, now), ) connection.execute( """ INSERT INTO fleet_run_items( id, run_id, instance_id, position, status, stage, stage_started_at, attempt_count, config_version, instance_display_name, aws_region, lightsail_instance_name, cloudflare_zone_name, cloudflare_zone_id, cloudflare_record_name, socks_port, proxy_health_check, health_timeout_seconds, release_grace_seconds, started_at, updated_at ) VALUES ( 'item', 'run', ?, 0, 'queued', 'preflight', ?, 2, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ? ) """, ( instance.id, now, instance.config_version, instance.display_name, instance.aws_region, instance.lightsail_instance_name, instance.cloudflare_zone_name, instance.cloudflare_zone_id, instance.cloudflare_record_name, instance.socks_port, int(instance.proxy_health_check), instance.health_timeout_seconds, instance.release_grace_seconds, now, now, ), ) item = repository.get_fleet_run_item("item") assert item is not None assert item.stage_started_at == now assert item.attempt_count == 2 assert item.proxy_health_check is True assert repository.list_fleet_run_items("run") == [item] def test_legacy_configuration_is_migrated_and_old_active_run_is_failed( tmp_path: Path, ) -> None: staged_migrations = tmp_path / "migrations" staged_migrations.mkdir() shutil.copy2(MIGRATIONS_PATH / "001_initial.sql", staged_migrations) shutil.copy2(MIGRATIONS_PATH / "002_operation_locks.sql", staged_migrations) database = Database(tmp_path / "legacy.db", staged_migrations) database.migrate() now = to_iso() with database.connect() as connection, connection: connection.execute( """ UPDATE app_settings SET config_version = 9, aws_region = 'eu-west-1', lightsail_instance_name = 'legacy-node', cloudflare_zone_name = 'example.com', cloudflare_zone_id = 'legacy-zone', cloudflare_record_name = 'legacy.example.com', socks_port = 2080, proxy_health_check = 0, health_timeout_seconds = 45 WHERE id = 1 """ ) connection.execute( """ UPDATE schedule_state SET enabled = 1, interval_minutes = 90, next_run_at = ?, last_run_at = ? WHERE id = 1 """, (to_iso(utc_now() + timedelta(minutes=90)), now), ) connection.execute( """ INSERT INTO rotation_runs( id, trigger, status, stage, stage_started_at, active_slot, config_version, lease_owner, lease_until, started_at, updated_at ) VALUES ( 'legacy-active', 'manual', 'running', 'waiting_new_ip', ?, 1, 9, 'old-worker', ?, ?, ? ) """, (now, to_iso(utc_now() + timedelta(minutes=5)), now, now), ) connection.execute( """ INSERT INTO operation_locks(slot, kind, owner_id, lease_until, created_at) VALUES (1, 'rotation', 'legacy-active', NULL, ?) """, (now,), ) shutil.copy2(MIGRATIONS_PATH / "003_multi_instance_static_ip.sql", staged_migrations) shutil.copy2(MIGRATIONS_PATH / "004_rotation_rollback.sql", staged_migrations) shutil.copy2(MIGRATIONS_PATH / "005_credential_accounts.sql", staged_migrations) shutil.copy2( MIGRATIONS_PATH / "006_proxy_health_and_cloudflare_mode.sql", staged_migrations, ) database.migrate() repository = FleetRepository(database) instance = repository.get_instance("legacy-instance") group = repository.get_group("legacy-group") assert instance is not None assert instance.enabled is False assert instance.config_version == 9 assert instance.lightsail_instance_name == "legacy-node" assert instance.release_grace_seconds == 75 assert group is not None assert group.enabled is False assert group.interval_minutes == 90 assert group.next_run_at is None assert group.last_run_at == now assert group.member_ids == [instance.id] with database.connect() as connection: legacy_run = connection.execute( """ SELECT status, stage, active_slot, lease_owner, lease_until, error_code, finished_at FROM rotation_runs WHERE id = 'legacy-active' """ ).fetchone() event = connection.execute( "SELECT level, details_json FROM rotation_events WHERE run_id = 'legacy-active'" ).fetchone() lock_count = connection.execute("SELECT COUNT(*) FROM operation_locks").fetchone()[0] assert dict(legacy_run) == { "status": "failed", "stage": "failed", "active_slot": None, "lease_owner": None, "lease_until": None, "error_code": "WORKFLOW_UPGRADED", "finished_at": legacy_run["finished_at"], } assert legacy_run["finished_at"] is not None assert dict(event) == { "level": "warning", "details_json": '{"code":"WORKFLOW_UPGRADED"}', } assert lock_count == 0