Skip to content

Commit fc37121

Browse files
pablomhcursoragent
andcommitted
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. Assisted-By: Cursor Grok 4.6 Co-authored-by: Cursor <cursoragent@cursor.com>
1 parent 4b9ff21 commit fc37121

8 files changed

Lines changed: 183 additions & 14 deletions

File tree

CHANGES/+remote-policy.feature

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1 @@
1+
Added `UpstreamPulp.remote_policy` so remotes created during replication can use `on_demand` or `streamed` instead of defaulting to `immediate`.

docs/user/guides/replication.md

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -48,6 +48,7 @@ pulp upstream-pulp create \
4848
| `tls_validation` | Whether to verify the upstream server's TLS certificate. Defaults to `True`. |
4949
| `q_select` | A filter expression to select which upstream distributions to replicate. See [Filtering Distributions](#filtering-distributions-with-q_select). |
5050
| `policy` | Controls how replication manages local objects. One of `all`, `labeled`, or `nodelete`. See [Replication Policies](#replication-policies). Defaults to `all`. |
51+
| `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`). |
5152

5253
## Running Replication
5354

@@ -151,7 +152,8 @@ pulp upstream-pulp replicate --upstream-pulp "my-upstream"
151152
## Replication Policies
152153

153154
The `policy` field controls how replication handles local objects, particularly when upstream
154-
distributions are removed or no longer match a `q_select` filter.
155+
distributions are removed or no longer match a `q_select` filter. It is not the same as a remote's
156+
download policy (`immediate`, `on_demand`, or `streamed`); set that with `remote_policy`.
155157

156158
### `all` (default)
157159

Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
from django.db import migrations, models
2+
3+
4+
class Migration(migrations.Migration):
5+
6+
dependencies = [
7+
("core", "0156_alter_contentartifact_relative_path_and_more"),
8+
]
9+
10+
operations = [
11+
migrations.AddField(
12+
model_name="upstreampulp",
13+
name="remote_policy",
14+
field=models.TextField(
15+
choices=[
16+
("immediate", "When syncing, download all metadata and content now."),
17+
(
18+
"on_demand",
19+
"When syncing, download metadata, but do not download content now. "
20+
"Instead, download content as clients request it, and save it in Pulp "
21+
"to be served for future client requests.",
22+
),
23+
(
24+
"streamed",
25+
"When syncing, download metadata, but do not download content now. "
26+
"Instead,download content as clients request it, but never save it in "
27+
"Pulp. This causes future requests for that same content to have to be "
28+
"downloaded again.",
29+
),
30+
],
31+
null=True,
32+
),
33+
),
34+
]

pulpcore/app/models/replica.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,8 @@
1111
from pulpcore.app.util import get_domain_pk
1212
from pulpcore.plugin.models import AutoAddObjPermsMixin, BaseModel, EncryptedTextField
1313

14+
from .repository import Remote
15+
1416

1517
class UpstreamPulp(BaseModel, AutoAddObjPermsMixin):
1618
ALL = "all"
@@ -59,6 +61,7 @@ class UpstreamPulp(BaseModel, AutoAddObjPermsMixin):
5961
sock_read_timeout = models.FloatField(
6062
null=True, validators=[MinValueValidator(0.0, "Timeout must be >= 0")]
6163
)
64+
remote_policy = models.TextField(choices=Remote.POLICY_CHOICES, null=True)
6265

6366
q_select = models.TextField(null=True)
6467
policy = models.TextField(choices=POLICY_CHOICES, default=ALL)

pulpcore/app/serializers/replica.py

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,7 +3,7 @@
33
from rest_framework import serializers
44
from rest_framework.validators import UniqueValidator
55

6-
from pulpcore.app.models import UpstreamPulp
6+
from pulpcore.app.models import Remote, UpstreamPulp
77
from pulpcore.app.serializers import (
88
HiddenFieldsMixin,
99
IdentityField,
@@ -122,6 +122,16 @@ class UpstreamPulpSerializer(ModelSerializer, HiddenFieldsMixin):
122122
),
123123
min_value=0.0,
124124
)
125+
remote_policy = serializers.ChoiceField(
126+
choices=Remote.POLICY_CHOICES,
127+
help_text=_(
128+
"Download policy for remotes created during replication. One of 'immediate', "
129+
"'on_demand', or 'streamed'. Distinct from 'policy', which controls how replicate "
130+
"manages local objects. Defaults to the Remote default ('immediate') when unset."
131+
),
132+
required=False,
133+
allow_null=True,
134+
)
125135

126136
pulp_last_updated = serializers.DateTimeField(
127137
help_text="Timestamp of the most recent update of the remote.", read_only=True
@@ -178,6 +188,7 @@ class Meta:
178188
"connect_timeout",
179189
"sock_connect_timeout",
180190
"sock_read_timeout",
191+
"remote_policy",
181192
"pulp_last_updated",
182193
"hidden_fields",
183194
"q_select",

pulpcore/app/tasks/replica.py

Lines changed: 21 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,26 @@ def user_agent():
2828
return f"pulpcore/{pulp_version} ({python}, {system}) (pulp-glue {pulp_glue_version})"
2929

3030

31+
def _build_remote_settings(server):
32+
"""Build fields copied onto remotes created during replication."""
33+
remote_settings = {
34+
"ca_cert": server.ca_cert,
35+
"tls_validation": server.tls_validation,
36+
"client_cert": server.client_cert,
37+
"client_key": server.client_key,
38+
"download_concurrency": server.download_concurrency,
39+
"max_retries": server.max_retries,
40+
"total_timeout": server.total_timeout,
41+
"connect_timeout": server.connect_timeout,
42+
"sock_connect_timeout": server.sock_connect_timeout,
43+
"sock_read_timeout": server.sock_read_timeout,
44+
}
45+
# Omit policy when unset so new remotes keep Remote.policy's default (immediate).
46+
if server.remote_policy is not None:
47+
remote_settings["policy"] = server.remote_policy
48+
return remote_settings
49+
50+
3151
def replicate_distributions(server_pk, q_select=None, **kwargs):
3252
server = UpstreamPulp.objects.get(pk=server_pk)
3353

@@ -58,18 +78,7 @@ def replicate_distributions(server_pk, q_select=None, **kwargs):
5878
}
5979
)
6080

61-
remote_settings = {
62-
"ca_cert": server.ca_cert,
63-
"tls_validation": server.tls_validation,
64-
"client_cert": server.client_cert,
65-
"client_key": server.client_key,
66-
"download_concurrency": server.download_concurrency,
67-
"max_retries": server.max_retries,
68-
"total_timeout": server.total_timeout,
69-
"connect_timeout": server.connect_timeout,
70-
"sock_connect_timeout": server.sock_connect_timeout,
71-
"sock_read_timeout": server.sock_read_timeout,
72-
}
81+
remote_settings = _build_remote_settings(server)
7382

7483
try:
7584
task_group = TaskGroup.current()

pulpcore/tests/functional/api/test_replication.py

Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -281,6 +281,80 @@ def test_replication_remote_settings_propagation(
281281
assert remote.max_retries == 2
282282

283283

284+
@pytest.mark.parallel
285+
def test_replication_remote_policy(
286+
domain_factory,
287+
bindings_cfg,
288+
pulpcore_bindings,
289+
file_bindings,
290+
monitor_task,
291+
monitor_task_group,
292+
pulp_settings,
293+
gen_object_with_cleanup,
294+
file_distribution_factory,
295+
file_publication_factory,
296+
file_repository_factory,
297+
tmp_path,
298+
add_domain_objects_to_cleanup,
299+
):
300+
"""Remotes created by replicate() inherit UpstreamPulp.remote_policy when set."""
301+
source_domain = domain_factory()
302+
add_domain_objects_to_cleanup(source_domain)
303+
304+
repository = file_repository_factory(pulp_domain=source_domain.name)
305+
file_path = tmp_path / "file.txt"
306+
file_path.write_text("DEADBEEF")
307+
monitor_task(
308+
file_bindings.ContentFilesApi.create(
309+
file=str(file_path),
310+
relative_path="file.txt",
311+
repository=repository.pulp_href,
312+
pulp_domain=source_domain.name,
313+
).task
314+
)
315+
publication = file_publication_factory(
316+
pulp_domain=source_domain.name, repository=repository.pulp_href
317+
)
318+
file_distribution_factory(pulp_domain=source_domain.name, publication=publication.pulp_href)
319+
320+
replica_domain = domain_factory()
321+
add_domain_objects_to_cleanup(replica_domain)
322+
323+
upstream_pulp_body = {
324+
"name": str(uuid.uuid4()),
325+
"base_url": bindings_cfg.host,
326+
"api_root": pulp_settings.API_ROOT,
327+
"domain": source_domain.name,
328+
"username": bindings_cfg.username,
329+
"password": bindings_cfg.password,
330+
"remote_policy": "on_demand",
331+
}
332+
upstream_pulp = gen_object_with_cleanup(
333+
pulpcore_bindings.UpstreamPulpsApi, upstream_pulp_body, pulp_domain=replica_domain.name
334+
)
335+
336+
response = pulpcore_bindings.UpstreamPulpsApi.replicate(
337+
upstream_pulp.pulp_href, pulpcore_bindings.module.UpstreamPulpReplicate()
338+
)
339+
monitor_task_group(response.task_group)
340+
341+
result = file_bindings.RemotesFileApi.list(pulp_domain=replica_domain.name)
342+
assert result.count == 1
343+
remote = result.results[0]
344+
assert remote.policy == "on_demand"
345+
346+
pulpcore_bindings.UpstreamPulpsApi.partial_update(
347+
upstream_pulp.pulp_href, {"remote_policy": "streamed"}
348+
)
349+
response = pulpcore_bindings.UpstreamPulpsApi.replicate(
350+
upstream_pulp.pulp_href, pulpcore_bindings.module.UpstreamPulpReplicate()
351+
)
352+
monitor_task_group(response.task_group)
353+
354+
remote = file_bindings.RemotesFileApi.list(pulp_domain=replica_domain.name).results[0]
355+
assert remote.policy == "streamed"
356+
357+
284358
@pytest.mark.parallel
285359
def test_replication_with_repo_based_distribution(
286360
domain_factory,
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
from types import SimpleNamespace
2+
3+
from pulpcore.app.tasks.replica import _build_remote_settings
4+
5+
6+
def _server(**overrides):
7+
base = {
8+
"ca_cert": "api-ca",
9+
"tls_validation": True,
10+
"client_cert": "api-cert",
11+
"client_key": "api-key",
12+
"download_concurrency": 10,
13+
"max_retries": 3,
14+
"total_timeout": 30,
15+
"connect_timeout": 5,
16+
"sock_connect_timeout": 5,
17+
"sock_read_timeout": 5,
18+
"remote_policy": None,
19+
}
20+
base.update(overrides)
21+
return SimpleNamespace(**base)
22+
23+
24+
def test_build_remote_settings_omits_policy_when_unset():
25+
settings = _build_remote_settings(_server())
26+
27+
assert "policy" not in settings
28+
assert settings["ca_cert"] == "api-ca"
29+
assert settings["download_concurrency"] == 10
30+
31+
32+
def test_build_remote_settings_includes_policy_when_set():
33+
settings = _build_remote_settings(_server(remote_policy="on_demand"))
34+
35+
assert settings["policy"] == "on_demand"

0 commit comments

Comments
 (0)