From ceba57eeabb855b0c4e6f4aa628d047fae55bc7d Mon Sep 17 00:00:00 2001 From: Gerrod Ubben Date: Thu, 27 Aug 2026 10:44:02 -0400 Subject: [PATCH] Make content_ids migration incremental and restartable Perform the cache populating entirely in SQL and add per repository commits. Assisted By: Cursor Grok 4.6 Co-authored-by: Cursor --- CHANGES/+content-ids-migration.bugfix | 1 + ...152_alter_repositoryversion_content_ids.py | 68 +++++++++++----- .../unit/models/test_content_ids_migration.py | 77 +++++++++++++++++++ 3 files changed, 127 insertions(+), 19 deletions(-) create mode 100644 CHANGES/+content-ids-migration.bugfix create mode 100644 pulpcore/tests/unit/models/test_content_ids_migration.py diff --git a/CHANGES/+content-ids-migration.bugfix b/CHANGES/+content-ids-migration.bugfix new file mode 100644 index 00000000000..9ab90a2b53c --- /dev/null +++ b/CHANGES/+content-ids-migration.bugfix @@ -0,0 +1 @@ +Filled missing repository version `content_ids` in SQL from content membership, committing per repository so large databases are less likely to hit statement timeouts. diff --git a/pulpcore/app/migrations/0152_alter_repositoryversion_content_ids.py b/pulpcore/app/migrations/0152_alter_repositoryversion_content_ids.py index fc4ba269b16..8bba79ab876 100644 --- a/pulpcore/app/migrations/0152_alter_repositoryversion_content_ids.py +++ b/pulpcore/app/migrations/0152_alter_repositoryversion_content_ids.py @@ -1,35 +1,65 @@ # Generated by Django 5.2.13 on 2026-06-08 21:54 import django.contrib.postgres.fields -from django.db import migrations, models +from django.db import migrations, models, transaction + +# Same membership as RepositoryVersion._content_relationships(). +_POPULATE_SQL = """ +UPDATE __RV__ AS rv +SET content_ids = COALESCE(( + SELECT ARRAY_AGG(rc.content_id) + FROM __RC__ AS rc + INNER JOIN __RV__ AS va ON va.pulp_id = rc.version_added_id + LEFT JOIN __RV__ AS vr ON vr.pulp_id = rc.version_removed_id + WHERE rc.repository_id = rv.repository_id + AND va.number <= rv.number + AND (vr.pulp_id IS NULL OR vr.number > rv.number) +), '{}'::uuid[]) +WHERE rv.repository_id = %s + AND rv.content_ids IS NULL +""" def populate_content_ids(apps, schema_editor): - RepositoryVersion = apps.get_model('core', 'RepositoryVersion') - RepositoryContent = apps.get_model('core', 'RepositoryContent') - repo_versions = [] - for rv in RepositoryVersion.objects.filter(content_ids=None).iterator(chunk_size=2000): - rv.content_ids = list(RepositoryContent.objects.filter( - repository_id=rv.repository_id, version_added__number__lte=rv.number - ).exclude(version_removed__number__lte=rv.number).values_list("content_id", flat=True)) - repo_versions.append(rv) - if len(repo_versions) >= 2000: - RepositoryVersion.objects.bulk_update(repo_versions, ['content_ids']) - repo_versions = [] - if repo_versions: - RepositoryVersion.objects.bulk_update(repo_versions, ['content_ids']) + """Fill missing content_ids from content membership, committing per repository.""" + RepositoryVersion = apps.get_model("core", "RepositoryVersion") + RepositoryContent = apps.get_model("core", "RepositoryContent") + connection = schema_editor.connection + qn = schema_editor.quote_name + rv = qn(RepositoryVersion._meta.db_table) + rc = qn(RepositoryContent._meta.db_table) + sql = _POPULATE_SQL.replace("__RV__", rv).replace("__RC__", rc) + + with connection.cursor() as cursor: + cursor.execute(f"SELECT DISTINCT repository_id FROM {rv} WHERE content_ids IS NULL") + repo_ids = [row[0] for row in cursor.fetchall()] + + for repo_id in repo_ids: + with transaction.atomic(using=connection.alias): + with connection.cursor() as cursor: + cursor.execute(sql, [repo_id]) + class Migration(migrations.Migration): + # Per-repository commits inside RunPython; AlterField runs afterwards. + atomic = False dependencies = [ - ('core', '0151_upstreampulp_connect_timeout_and_more'), + ("core", "0151_upstreampulp_connect_timeout_and_more"), ] operations = [ - migrations.RunPython(populate_content_ids, reverse_code=migrations.RunPython.noop, elidable=True), + migrations.RunPython( + populate_content_ids, + reverse_code=migrations.RunPython.noop, + elidable=True, + atomic=False, + ), migrations.AlterField( - model_name='repositoryversion', - name='content_ids', - field=django.contrib.postgres.fields.ArrayField(base_field=models.UUIDField(), default=list, size=None), + model_name="repositoryversion", + name="content_ids", + field=django.contrib.postgres.fields.ArrayField( + base_field=models.UUIDField(), default=list, size=None + ), ), ] diff --git a/pulpcore/tests/unit/models/test_content_ids_migration.py b/pulpcore/tests/unit/models/test_content_ids_migration.py new file mode 100644 index 00000000000..7bbeb1fc874 --- /dev/null +++ b/pulpcore/tests/unit/models/test_content_ids_migration.py @@ -0,0 +1,77 @@ +import importlib +from types import SimpleNamespace +from uuid import uuid4 + +import pytest +from django.apps import apps +from django.db import connection + +from pulpcore.app.models import RepositoryVersion +from pulpcore.plugin.models import Content, Repository + +populate_content_ids = importlib.import_module( + "pulpcore.app.migrations.0152_alter_repositoryversion_content_ids" +).populate_content_ids + + +def _run_populate(): + schema_editor = SimpleNamespace( + connection=connection, + quote_name=connection.ops.quote_name, + ) + populate_content_ids(apps, schema_editor) + + +def _membership_ids(version): + return set(version._content_relationships().values_list("content_id", flat=True)) + + +@pytest.fixture +def repository(db): + repository = Repository.objects.create(name=uuid4()) + repository.CONTENT_TYPES = [Content] + return repository + + +def test_populate_content_ids_from_content_relationships(repository): + contents = [Content(pulp_type="core.content") for _ in range(4)] + Content.objects.bulk_create(contents) + pks = [c.pk for c in contents] + + version0 = repository.latest_version() + with repository.new_version() as version1: + version1.add_content(Content.objects.filter(pk__in=pks[:3])) + with repository.new_version() as version2: + version2.remove_content(Content.objects.filter(pk__in=pks[:1])) + with repository.new_version() as version3: + version3.add_content(Content.objects.filter(pk__in=pks[3:])) + + # A second repository whose cache is already populated must be left alone. + other = Repository.objects.create(name=uuid4()) + other.CONTENT_TYPES = [Content] + with other.new_version() as other_v1: + other_v1.add_content(Content.objects.filter(pk__in=pks[:2])) + other_v1.refresh_from_db() + other_ids_before = list(other_v1.content_ids) + + rv_table = connection.ops.quote_name(RepositoryVersion._meta.db_table) + with connection.cursor() as cursor: + # Flush deferred triggers from the version inserts so ALTER TABLE can run + # inside the test transaction. + cursor.execute("SET CONSTRAINTS ALL IMMEDIATE") + cursor.execute(f"ALTER TABLE {rv_table} ALTER COLUMN content_ids DROP NOT NULL") + # Leave version3 populated (the post-3.83 case) and null the older versions. + cursor.execute( + f"UPDATE {rv_table} SET content_ids = NULL WHERE repository_id = %s AND number < %s", + [repository.pk, version3.number], + ) + + _run_populate() + + for version in (version0, version1, version2, version3): + version.refresh_from_db() + assert version.content_ids is not None + assert set(version.content_ids) == _membership_ids(version) + + other_v1.refresh_from_db() + assert list(other_v1.content_ids) == other_ids_before