From 01e4c138f3b911133b821d00823550e6a03902e9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pablo=20M=C3=A9ndez=20Hern=C3=A1ndez?= Date: Thu, 27 Aug 2026 10:07:45 +0200 Subject: [PATCH 1/2] Add UpstreamPulp.remote_policy for remotes created during replication replicate() never set Remote.policy, so new remotes defaulted to immediate and downloaded all artifacts. Let UpstreamPulp carry the intended download policy so Capsules can replicate with on_demand. Always include policy in the remote settings dict so that clearing remote_policy (setting it to null) correctly reverts existing remotes back to Remote.IMMEDIATE instead of leaving the old value stale. Assisted-By: Cursor Co-authored-by: Cursor Co-Authored-By: Claude Opus 4.6 --- CHANGES/+remote-policy.feature | 1 + docs/user/guides/replication.md | 4 +- .../0160_upstreampulp_remote_policy.py | 34 ++++++++ pulpcore/app/models/replica.py | 3 + pulpcore/app/serializers/replica.py | 13 ++- pulpcore/app/tasks/replica.py | 33 ++++--- .../tests/functional/api/test_replication.py | 87 +++++++++++++++++++ pulpcore/tests/unit/test_replica.py | 36 +++++++- 8 files changed, 195 insertions(+), 16 deletions(-) create mode 100644 CHANGES/+remote-policy.feature create mode 100644 pulpcore/app/migrations/0160_upstreampulp_remote_policy.py diff --git a/CHANGES/+remote-policy.feature b/CHANGES/+remote-policy.feature new file mode 100644 index 00000000000..712da8f4e85 --- /dev/null +++ b/CHANGES/+remote-policy.feature @@ -0,0 +1 @@ +Added `UpstreamPulp.remote_policy` so remotes created during replication can use `on_demand` or `streamed` instead of defaulting to `immediate`. diff --git a/docs/user/guides/replication.md b/docs/user/guides/replication.md index a8b627e3e42..7bf07c3a4c6 100644 --- a/docs/user/guides/replication.md +++ b/docs/user/guides/replication.md @@ -48,6 +48,7 @@ pulp upstream-pulp create \ | `tls_validation` | Whether to verify the upstream server's TLS certificate. Defaults to `True`. | | `q_select` | A filter expression to select which upstream distributions to replicate. See [Filtering Distributions](#filtering-distributions-with-q_select). | | `policy` | Controls how replication manages local objects. One of `all`, `labeled`, or `nodelete`. See [Replication Policies](#replication-policies). Defaults to `all`. | +| `remote_policy` | Download policy for remotes created during replication. One of `immediate`, `on_demand`, or `streamed`. Distinct from `policy`. When unset, remotes use Pulp's default (`immediate`). | ## Running Replication @@ -151,7 +152,8 @@ pulp upstream-pulp replicate --upstream-pulp "my-upstream" ## Replication Policies The `policy` field controls how replication handles local objects, particularly when upstream -distributions are removed or no longer match a `q_select` filter. +distributions are removed or no longer match a `q_select` filter. It is not the same as a remote's +download policy (`immediate`, `on_demand`, or `streamed`); set that with `remote_policy`. ### `all` (default) diff --git a/pulpcore/app/migrations/0160_upstreampulp_remote_policy.py b/pulpcore/app/migrations/0160_upstreampulp_remote_policy.py new file mode 100644 index 00000000000..e319f3d2151 --- /dev/null +++ b/pulpcore/app/migrations/0160_upstreampulp_remote_policy.py @@ -0,0 +1,34 @@ +from django.db import migrations, models + + +class Migration(migrations.Migration): + + dependencies = [ + ("core", "0159_alter_contentartifact_relative_path_and_more"), + ] + + operations = [ + migrations.AddField( + model_name="upstreampulp", + name="remote_policy", + field=models.TextField( + choices=[ + ("immediate", "When syncing, download all metadata and content now."), + ( + "on_demand", + "When syncing, download metadata, but do not download content now. " + "Instead, download content as clients request it, and save it in Pulp " + "to be served for future client requests.", + ), + ( + "streamed", + "When syncing, download metadata, but do not download content now. " + "Instead,download content as clients request it, but never save it in " + "Pulp. This causes future requests for that same content to have to be " + "downloaded again.", + ), + ], + null=True, + ), + ), + ] diff --git a/pulpcore/app/models/replica.py b/pulpcore/app/models/replica.py index 29b32c4d5c0..4bbd7fa73a4 100644 --- a/pulpcore/app/models/replica.py +++ b/pulpcore/app/models/replica.py @@ -11,6 +11,8 @@ from pulpcore.app.util import get_domain_pk from pulpcore.plugin.models import AutoAddObjPermsMixin, BaseModel, EncryptedTextField +from .repository import Remote + class UpstreamPulp(BaseModel, AutoAddObjPermsMixin): ALL = "all" @@ -59,6 +61,7 @@ class UpstreamPulp(BaseModel, AutoAddObjPermsMixin): sock_read_timeout = models.FloatField( null=True, validators=[MinValueValidator(0.0, "Timeout must be >= 0")] ) + remote_policy = models.TextField(choices=Remote.POLICY_CHOICES, null=True) q_select = models.TextField(null=True) policy = models.TextField(choices=POLICY_CHOICES, default=ALL) diff --git a/pulpcore/app/serializers/replica.py b/pulpcore/app/serializers/replica.py index 7d13425865f..e130ac8afee 100644 --- a/pulpcore/app/serializers/replica.py +++ b/pulpcore/app/serializers/replica.py @@ -3,7 +3,7 @@ from rest_framework import serializers from rest_framework.validators import UniqueValidator -from pulpcore.app.models import UpstreamPulp +from pulpcore.app.models import Remote, UpstreamPulp from pulpcore.app.serializers import ( HiddenFieldsMixin, IdentityField, @@ -122,6 +122,16 @@ class UpstreamPulpSerializer(ModelSerializer, HiddenFieldsMixin): ), min_value=0.0, ) + remote_policy = serializers.ChoiceField( + choices=Remote.POLICY_CHOICES, + help_text=_( + "Download policy for remotes created during replication. One of 'immediate', " + "'on_demand', or 'streamed'. Distinct from 'policy', which controls how replicate " + "manages local objects. Defaults to the Remote default ('immediate') when unset." + ), + required=False, + allow_null=True, + ) pulp_last_updated = serializers.DateTimeField( help_text="Timestamp of the most recent update of the remote.", read_only=True @@ -178,6 +188,7 @@ class Meta: "connect_timeout", "sock_connect_timeout", "sock_read_timeout", + "remote_policy", "pulp_last_updated", "hidden_fields", "q_select", diff --git a/pulpcore/app/tasks/replica.py b/pulpcore/app/tasks/replica.py index 19de2bd60ed..e2bb50295d8 100644 --- a/pulpcore/app/tasks/replica.py +++ b/pulpcore/app/tasks/replica.py @@ -11,7 +11,7 @@ from pulp_glue.common.exceptions import PulpException as GluePulpException from pulpcore.app.apps import PulpAppConfig, pulp_plugin_configs -from pulpcore.app.models import Distribution, Repository, Task, TaskGroup, UpstreamPulp +from pulpcore.app.models import Distribution, Remote, Repository, Task, TaskGroup, UpstreamPulp from pulpcore.app.replica import ReplicaContext, distros_lock_uri from pulpcore.constants import TASK_STATES from pulpcore.exceptions import ExternalServiceError @@ -52,6 +52,24 @@ def _ssl_temp_files(server): pass +def _build_remote_settings(server): + """Build fields copied onto remotes created during replication.""" + remote_settings = { + "ca_cert": server.ca_cert, + "tls_validation": server.tls_validation, + "client_cert": server.client_cert, + "client_key": server.client_key, + "download_concurrency": server.download_concurrency, + "max_retries": server.max_retries, + "total_timeout": server.total_timeout, + "connect_timeout": server.connect_timeout, + "sock_connect_timeout": server.sock_connect_timeout, + "sock_read_timeout": server.sock_read_timeout, + } + remote_settings["policy"] = server.remote_policy or Remote.IMMEDIATE + return remote_settings + + def replicate_distributions(server_pk, q_select=None, **kwargs): server = UpstreamPulp.objects.get(pk=server_pk) with _ssl_temp_files(server) as ssl_files: @@ -75,18 +93,7 @@ def replicate_distributions(server_pk, q_select=None, **kwargs): } ) - remote_settings = { - "ca_cert": server.ca_cert, - "tls_validation": server.tls_validation, - "client_cert": server.client_cert, - "client_key": server.client_key, - "download_concurrency": server.download_concurrency, - "max_retries": server.max_retries, - "total_timeout": server.total_timeout, - "connect_timeout": server.connect_timeout, - "sock_connect_timeout": server.sock_connect_timeout, - "sock_read_timeout": server.sock_read_timeout, - } + remote_settings = _build_remote_settings(server) try: task_group = TaskGroup.current() supported_replicators = [] diff --git a/pulpcore/tests/functional/api/test_replication.py b/pulpcore/tests/functional/api/test_replication.py index f208a4be082..867b36ad068 100644 --- a/pulpcore/tests/functional/api/test_replication.py +++ b/pulpcore/tests/functional/api/test_replication.py @@ -252,6 +252,7 @@ def test_replication_remote_settings_propagation( assert remote.sock_read_timeout == 45.0 assert remote.download_concurrency == 5 assert remote.max_retries == 7 + assert remote.policy == "immediate" # Update all settings and re-replicate to verify propagation on update pulpcore_bindings.UpstreamPulpsApi.partial_update( @@ -281,6 +282,92 @@ def test_replication_remote_settings_propagation( assert remote.max_retries == 2 +@pytest.mark.parallel +def test_replication_remote_policy( + domain_factory, + bindings_cfg, + pulpcore_bindings, + file_bindings, + monitor_task, + monitor_task_group, + pulp_settings, + gen_object_with_cleanup, + file_distribution_factory, + file_publication_factory, + file_repository_factory, + tmp_path, + add_domain_objects_to_cleanup, +): + """Remotes created by replicate() inherit UpstreamPulp.remote_policy when set.""" + source_domain = domain_factory() + add_domain_objects_to_cleanup(source_domain) + + repository = file_repository_factory(pulp_domain=source_domain.name) + file_path = tmp_path / "file.txt" + file_path.write_text("DEADBEEF") + monitor_task( + file_bindings.ContentFilesApi.create( + file=str(file_path), + relative_path="file.txt", + repository=repository.pulp_href, + pulp_domain=source_domain.name, + ).task + ) + publication = file_publication_factory( + pulp_domain=source_domain.name, repository=repository.pulp_href + ) + file_distribution_factory(pulp_domain=source_domain.name, publication=publication.pulp_href) + + replica_domain = domain_factory() + add_domain_objects_to_cleanup(replica_domain) + + upstream_pulp_body = { + "name": str(uuid.uuid4()), + "base_url": bindings_cfg.host, + "api_root": pulp_settings.API_ROOT, + "domain": source_domain.name, + "username": bindings_cfg.username, + "password": bindings_cfg.password, + "remote_policy": "on_demand", + } + upstream_pulp = gen_object_with_cleanup( + pulpcore_bindings.UpstreamPulpsApi, upstream_pulp_body, pulp_domain=replica_domain.name + ) + + response = pulpcore_bindings.UpstreamPulpsApi.replicate( + upstream_pulp.pulp_href, pulpcore_bindings.module.UpstreamPulpReplicate() + ) + monitor_task_group(response.task_group) + + result = file_bindings.RemotesFileApi.list(pulp_domain=replica_domain.name) + assert result.count == 1 + remote = result.results[0] + assert remote.policy == "on_demand" + + pulpcore_bindings.UpstreamPulpsApi.partial_update( + upstream_pulp.pulp_href, {"remote_policy": "streamed"} + ) + response = pulpcore_bindings.UpstreamPulpsApi.replicate( + upstream_pulp.pulp_href, pulpcore_bindings.module.UpstreamPulpReplicate() + ) + monitor_task_group(response.task_group) + + remote = file_bindings.RemotesFileApi.list(pulp_domain=replica_domain.name).results[0] + assert remote.policy == "streamed" + + # Clearing remote_policy should revert remotes back to immediate + pulpcore_bindings.UpstreamPulpsApi.partial_update( + upstream_pulp.pulp_href, {"remote_policy": None} + ) + response = pulpcore_bindings.UpstreamPulpsApi.replicate( + upstream_pulp.pulp_href, pulpcore_bindings.module.UpstreamPulpReplicate() + ) + monitor_task_group(response.task_group) + + remote = file_bindings.RemotesFileApi.list(pulp_domain=replica_domain.name).results[0] + assert remote.policy == "immediate" + + @pytest.mark.parallel def test_replication_with_repo_based_distribution( domain_factory, diff --git a/pulpcore/tests/unit/test_replica.py b/pulpcore/tests/unit/test_replica.py index 8a7c550281e..a07574e7396 100644 --- a/pulpcore/tests/unit/test_replica.py +++ b/pulpcore/tests/unit/test_replica.py @@ -3,8 +3,9 @@ import pytest +from pulpcore.app.models import Remote from pulpcore.app.tasks import replica -from pulpcore.app.tasks.replica import _ssl_temp_files +from pulpcore.app.tasks.replica import _build_remote_settings, _ssl_temp_files def test_ssl_temp_files_keep_all_certs_until_context_exits(tmp_path, monkeypatch): @@ -76,6 +77,7 @@ def test_replicate_distributions_sets_verify_ssl( connect_timeout=5, sock_connect_timeout=5, sock_read_timeout=5, + remote_policy=None, q_select=None, pulp_domain_id="domain-id", pk="server-pk", @@ -118,3 +120,35 @@ def fake_from_config(config): assert isinstance(captured["config"]["verify_ssl"], str) else: assert captured["config"]["verify_ssl"] is False + + +def _fake_server(**overrides): + base = { + "ca_cert": "api-ca", + "tls_validation": True, + "client_cert": "api-cert", + "client_key": "api-key", + "download_concurrency": 10, + "max_retries": 3, + "total_timeout": 30, + "connect_timeout": 5, + "sock_connect_timeout": 5, + "sock_read_timeout": 5, + "remote_policy": None, + } + base.update(overrides) + return SimpleNamespace(**base) + + +def test_build_remote_settings_defaults_to_immediate_when_unset(): + settings = _build_remote_settings(_fake_server()) + + assert settings["policy"] == Remote.IMMEDIATE + assert settings["ca_cert"] == "api-ca" + assert settings["download_concurrency"] == 10 + + +def test_build_remote_settings_includes_policy_when_set(): + settings = _build_remote_settings(_fake_server(remote_policy="on_demand")) + + assert settings["policy"] == "on_demand" From 429c5a349bfb0b6842f013709f667439ecc52c59 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pablo=20M=C3=A9ndez=20Hern=C3=A1ndez?= Date: Thu, 17 Sep 2026 21:07:33 +0200 Subject: [PATCH 2/2] Fix null-clearing test: use PatchedUpstreamPulp model class The generated Python client uses exclude_none=True in model_dump(), which silently drops None values from raw dicts. By constructing the PatchedUpstreamPulp model explicitly, Pydantic tracks remote_policy in model_fields_set and serializes it as null in the PATCH body. Co-Authored-By: Claude Opus 4.6 --- pulpcore/tests/functional/api/test_replication.py | 7 +++++-- 1 file changed, 5 insertions(+), 2 deletions(-) diff --git a/pulpcore/tests/functional/api/test_replication.py b/pulpcore/tests/functional/api/test_replication.py index 867b36ad068..914fa1974fa 100644 --- a/pulpcore/tests/functional/api/test_replication.py +++ b/pulpcore/tests/functional/api/test_replication.py @@ -355,10 +355,13 @@ def test_replication_remote_policy( remote = file_bindings.RemotesFileApi.list(pulp_domain=replica_domain.name).results[0] assert remote.policy == "streamed" - # Clearing remote_policy should revert remotes back to immediate + # Model class needed: raw dict {"remote_policy": None} is dropped by the client. pulpcore_bindings.UpstreamPulpsApi.partial_update( - upstream_pulp.pulp_href, {"remote_policy": None} + upstream_pulp.pulp_href, + pulpcore_bindings.module.PatchedUpstreamPulp(remote_policy=None), ) + upstream_pulp = pulpcore_bindings.UpstreamPulpsApi.read(upstream_pulp.pulp_href) + assert upstream_pulp.remote_policy is None response = pulpcore_bindings.UpstreamPulpsApi.replicate( upstream_pulp.pulp_href, pulpcore_bindings.module.UpstreamPulpReplicate() )