Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
44 changes: 44 additions & 0 deletions migrations/versions/026_repair_subject_release_tracking.py
Original file line number Diff line number Diff line change
@@ -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
107 changes: 107 additions & 0 deletions tests/test_subject_release_migration.py
Original file line number Diff line number Diff line change
@@ -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()
Loading