summaryrefslogtreecommitdiff
path: root/alembic/versions/20260411_0004_recent_backfill_scope.py
blob: 8d2998784a5438d220fe8c1feaaaf11fc8df51aa (plain) (blame)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
"""recent backfill scope and actor fields"""

from __future__ import annotations

from alembic import op
import sqlalchemy as sa
from sqlalchemy import inspect


revision = "20260411_0004"
down_revision = "20260411_0003"
branch_labels = None
depends_on = None


def _table_names() -> set[str]:
    return set(inspect(op.get_bind()).get_table_names())


def _column_names(table_name: str) -> set[str]:
    return {column["name"] for column in inspect(op.get_bind()).get_columns(table_name)}


def _service_backfill_needs_upgrade() -> bool:
    if "service_backfill_states" not in _table_names():
        return False
    columns = _column_names("service_backfill_states")
    if "scope" not in columns:
        return True
    unique_constraints = {
        constraint["name"]
        for constraint in inspect(op.get_bind()).get_unique_constraints("service_backfill_states")
    }
    return "uq_service_backfill_state_actor_service_scope" not in unique_constraints


def _upgrade_service_backfill_states() -> None:
    op.execute(
        sa.text(
            """
            CREATE TABLE service_backfill_states__alembic_new (
                id INTEGER NOT NULL PRIMARY KEY,
                actor VARCHAR(255) NOT NULL,
                service VARCHAR(32) NOT NULL,
                scope VARCHAR(16) NOT NULL,
                cursor_json JSON,
                status VARCHAR(32) NOT NULL,
                started_at DATETIME,
                completed_at DATETIME,
                last_error TEXT,
                updated_at DATETIME NOT NULL,
                CONSTRAINT uq_service_backfill_state_actor_service_scope UNIQUE (actor, service, scope)
            )
            """
        )
    )
    op.execute(
        sa.text(
            """
            INSERT INTO service_backfill_states__alembic_new
                (id, actor, service, scope, cursor_json, status, started_at, completed_at, last_error, updated_at)
            SELECT id, actor, service, 'full', cursor_json, status, started_at, completed_at, last_error, updated_at
            FROM service_backfill_states
            """
        )
    )
    op.execute(sa.text("DROP TABLE service_backfill_states"))
    op.execute(sa.text("ALTER TABLE service_backfill_states__alembic_new RENAME TO service_backfill_states"))


def upgrade() -> None:
    if "tracked_actors" in _table_names():
        columns = _column_names("tracked_actors")
        if "recent_backfill_status" not in columns:
            op.add_column("tracked_actors", sa.Column("recent_backfill_status", sa.String(length=32), nullable=True))
            op.execute(sa.text("UPDATE tracked_actors SET recent_backfill_status = 'pending' WHERE recent_backfill_status IS NULL"))
        if "recent_backfill_started_at" not in columns:
            op.add_column("tracked_actors", sa.Column("recent_backfill_started_at", sa.DateTime(timezone=True), nullable=True))
        if "recent_backfill_completed_at" not in columns:
            op.add_column("tracked_actors", sa.Column("recent_backfill_completed_at", sa.DateTime(timezone=True), nullable=True))
        if "last_recent_backfill_error" not in columns:
            op.add_column("tracked_actors", sa.Column("last_recent_backfill_error", sa.Text(), nullable=True))

    if _service_backfill_needs_upgrade():
        _upgrade_service_backfill_states()


def downgrade() -> None:
    pass