FluxIP/tests/test_fleet_repository.py

858 lines
33 KiB
Python

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