diff --git a/migrations/versions/026_repair_subject_release_tracking.py b/migrations/versions/026_repair_subject_release_tracking.py new file mode 100644 index 00000000..4903d572 --- /dev/null +++ b/migrations/versions/026_repair_subject_release_tracking.py @@ -0,0 +1,44 @@ +"""Repair databases that applied 022 before subject release tracking was added. + +Revision ID: 026 +Revises: 025 +Create Date: 2026-09-13 +""" + +from collections.abc import Sequence + +import sqlalchemy as sa +from alembic import op + + +revision: str = "026" +down_revision: str = "025" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + columns = sa.inspect(op.get_bind()).get_columns("automation_runs") + if any(column["name"] == "subject_released_at" for column in columns): + return # The newer version of 022 already created the column and index. + + op.add_column( + "automation_runs", + sa.Column("subject_released_at", sa.DateTime(timezone=True), nullable=True), + ) + op.drop_index("ix_automation_runs_subject", table_name="automation_runs") + where = sa.text("subject_key IS NOT NULL AND subject_released_at IS NULL") + op.create_index( + "ix_automation_runs_subject", + "automation_runs", + ["automation_id", "subject_key", "created_at"], + unique=False, + postgresql_where=where, + sqlite_where=where, + ) + + +def downgrade() -> None: + # Current 022 already owns this schema. Preserve it (and release history) + # when returning to 022; that revision's downgrade removes both columns. + pass diff --git a/tests/test_subject_release_migration.py b/tests/test_subject_release_migration.py new file mode 100644 index 00000000..4c05f96c --- /dev/null +++ b/tests/test_subject_release_migration.py @@ -0,0 +1,107 @@ +"""Upgrade both schemas that were deployed with the same revision 022.""" + +from datetime import datetime +from pathlib import Path +from uuid import uuid4 + +import pytest +import sqlalchemy as sa +from alembic import command +from alembic.config import Config +from alembic.migration import MigrationContext +from alembic.operations import Operations +from sqlalchemy.exc import OperationalError +from sqlalchemy.orm import Session + +from openhands.automation.models import Automation, AutomationRun, AutomationRunStatus + + +@pytest.mark.parametrize("schema", ["legacy_022", "current_022", "fresh"]) +def test_subject_release_upgrade_preserves_runs(schema, tmp_path, monkeypatch): + root = Path(__file__).parent.parent + url = f"sqlite:///{tmp_path / 'upgrade.db'}" + monkeypatch.setenv("AUTOMATION_DB_URL", url) + config = Config(str(root / "alembic.ini")) + config.set_main_option("script_location", str(root / "migrations")) + engine = sa.create_engine(url) + try: + target = "021" if schema == "legacy_022" else "022" + command.upgrade(config, "head" if schema == "fresh" else target) + if schema == "legacy_022": + # Exact 022 DDL from c756d24c91bc12da8bb329be815824bbc77dafe5: + # the original index had a PostgreSQL predicate, but no SQLite one. + with engine.begin() as connection: + operations = Operations(MigrationContext.configure(connection)) + operations.add_column( + "automation_runs", sa.Column("subject_key", sa.String(500)) + ) + operations.create_index( + "ix_automation_runs_subject", + "automation_runs", + ["automation_id", "subject_key", "created_at"], + postgresql_where=sa.text("subject_key IS NOT NULL"), + ) + command.stamp(config, "022") + # This is the deployed failure state: later migrations applied, + # but Alembic could not see that revision 022's schema was stale. + command.upgrade(config, "025") + with Session(engine) as session: + with pytest.raises(OperationalError, match="subject_released_at"): + session.scalars(sa.select(AutomationRun)).all() + + automation_id, run_id = uuid4(), uuid4() + released_at = None if schema == "legacy_022" else datetime(2026, 9, 13) + with engine.begin() as connection: + connection.execute( + sa.insert(Automation).values( + id=automation_id, + user_id=uuid4(), + org_id=uuid4(), + name="existing automation", + trigger={}, + tarball_path="fixture.tar.gz", + entrypoint="python main.py", + ) + ) + values = {"subject_released_at": released_at} if released_at else {} + connection.execute( + sa.insert(AutomationRun).values( + id=run_id, + automation_id=automation_id, + status=AutomationRunStatus.COMPLETED, + subject_key="repo/issue/1", + **values, + ) + ) + + command.upgrade(config, "head") + command.upgrade(config, "head") # Normal subsequent service startup. + with Session(engine) as session: + run = session.scalars(sa.select(AutomationRun)).one() + assert run.id == run_id + assert run.subject_key == "repo/issue/1" + assert run.status == AutomationRunStatus.COMPLETED + assert run.subject_released_at == released_at + index = next( + item + for item in sa.inspect(engine).get_indexes("automation_runs") + if item["name"] == "ix_automation_runs_subject" + ) + assert str(index.get("dialect_options", {}).get("sqlite_where")) == ( + "subject_key IS NOT NULL AND subject_released_at IS NULL" + ) + command.downgrade(config, "022") + # Current ORM models include fields from later migrations, so inspect + # revision 022 with SQL instead of trying to load it through the model. + with engine.connect() as connection: + downgraded = connection.execute( + sa.text("SELECT subject_released_at FROM automation_runs") + ).scalar_one() + assert (downgraded is None) == (released_at is None) + command.upgrade(config, "head") + with Session(engine) as session: + run = session.get(AutomationRun, run_id) + assert run is not None + assert run.subject_released_at == released_at + finally: + engine.dispose()