Skip to content
Open
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
1 change: 1 addition & 0 deletions CHANGES/7929.feature
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
Added `ETag`/`If-None-Match` and `Last-Modified`/`If-Modified-Since` (`304 Not Modified`) support on content-app responses so edge caches can revalidate without re-downloading.
3 changes: 3 additions & 0 deletions docs/admin/guides/configure-pulp/configure-storages.md
Original file line number Diff line number Diff line change
Expand Up @@ -141,6 +141,9 @@ When creating a `Domain` you can use the following payload:
Create your remotes, repositories, and distributions under this domain and the requests will be redirected to the
CloudFront custom domain specified.

For the content response validators and cache policy that apply when a CDN is in front of Pulp, see
the [CDN caching guide](site:pulpcore/docs/admin/learn/cdn-caching/).


Comprehensive options for Amazon S3 can be found in
[`django-storages` docs](https://django-storages.readthedocs.io/en/latest/backends/amazon-S3.html#configuration-settings).
Expand Down
39 changes: 39 additions & 0 deletions docs/admin/learn/cdn-caching.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
# CDN Caching

Pulp's content app implements some behavior to help with external caching processes.

## Cache Validation

For most requests, the Content App returns [`Last-Modified`][last-modified] and [`ETag`][etag] headers for the requested resource; the ETag is usually a digest.
[Shared caches][cache-types], such as CDNs, and private caches, such as clients, can use these validators to decide whether to serve a stored response or request the resource from Pulp again.
To avoid resources being served to unauthorized parties, Pulp uses [`Cache-Control`][cache-control] directives to enforce that:
(a) cache entities always ask for validation before serving cached content;
(b) only Private Caches can store certain sensitive responses (e.g., signed URLs).

A typical flow looks like:

1. A client requests a resource through a shared cache, such as a CDN.
If the cache has no copy, it forwards the request to Pulp.
2. Pulp returns the content with `Cache-Control: public, max-age=0, must-revalidate`, `Last-Modified`, and `ETag`.
The shared cache stores the response and its validators.
3. On a later request, the shared cache must revalidate its stored response with Pulp.
It sends [`If-None-Match`][if-none-match] with the stored ETag, [`If-Modified-Since`][if-modified-since] with the stored date, or both.
When both are present, Pulp checks `If-None-Match`.
4. If the content is unchanged, Pulp returns `304 Not Modified` without a body.
The shared cache serves its stored content to the client.
If the content changed, Pulp returns the full response with updated validators, and the cache replaces its stored copy.

### Known limitations

Pulp does not track the content that a distribution serves over time.
It infers `Last-Modified` based on the assumption that a distribution serves repository versions monotonically.
If that's not true, Pulp might tell the cache entity that a resource has not changed when it has.

The `ETag`/`If-None-Match` mechanism is unaffected by this limitation.

[cache-types]: https://developer.mozilla.org/en-US/docs/Web/HTTP/Guides/Caching#types_of_caches
[last-modified]: https://developer.mozilla.org/en-US/docs/Web/HTTP/Reference/Headers/Last-Modified
[etag]: https://developer.mozilla.org/en-US/docs/Web/HTTP/Reference/Headers/ETag
[cache-control]: https://developer.mozilla.org/en-US/docs/Web/HTTP/Reference/Headers/Cache-Control
[if-none-match]: https://developer.mozilla.org/en-US/docs/Web/HTTP/Reference/Headers/If-None-Match
[if-modified-since]: https://developer.mozilla.org/en-US/docs/Web/HTTP/Reference/Headers/If-Modified-Since
6 changes: 3 additions & 3 deletions pulpcore/app/models/publication.py
Original file line number Diff line number Diff line change
Expand Up @@ -868,7 +868,7 @@ def content_headers_for(self, path):

def get_fallback_ca(self, path):
"""
Return a ContentArtifact for path from the grace-period publication history, or None.
Return the ContentArtifact and RepositoryVersion for path from publication history, or None.

Iterates DistributedPublication records for this distribution from newest to oldest,
trying each publication until the path is found. Handles both pass-through and
Expand All @@ -893,7 +893,7 @@ def get_fallback_ca(self, path):
.first()
)
if ca is not None:
return ca
return ca, pub.repository_version
else:
pa = (
pub.published_artifact.select_related(
Expand All @@ -904,7 +904,7 @@ def get_fallback_ca(self, path):
.first()
)
if pa is not None:
return pa.content_artifact
return pa.content_artifact, pub.repository_version
return None

@hook(BEFORE_CREATE)
Expand Down
20 changes: 20 additions & 0 deletions pulpcore/app/util.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
from django.conf import settings
from django.db import connection
from django.db.models import Model, UUIDField
from django.utils.http import parse_http_date
from rest_framework.reverse import reverse as drf_reverse
from rest_framework.serializers import ValidationError

Expand Down Expand Up @@ -714,6 +715,25 @@ def normalize_http_status(status):
return ""


def check_request_was_modified(request, last_modified, etag=None):
if_none_match = request.headers.get("If-None-Match")
if_modified_since = request.headers.get("If-Modified-Since")

if not if_none_match and not (last_modified and if_modified_since):
return True

if if_none_match:
client_etags = [etag == client_etag.strip() for client_etag in if_none_match.split(",")]
return not any(client_etags)

try:
last_modified_ts = parse_http_date(last_modified)
if_modified_ts = parse_http_date(if_modified_since)
return last_modified_ts > if_modified_ts
except (TypeError, ValueError):
return True


class HashingFileWriter(RawIOBase):
"""
A file-like object that handles writing data to disk with simultaneous
Expand Down
27 changes: 18 additions & 9 deletions pulpcore/cache/cache.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,8 @@
import time
from functools import wraps

from aiohttp.web import FileResponse, HTTPSuccessful, Request, Response, StreamResponse
from aiohttp.web_exceptions import HTTPException, HTTPFound, HTTPNotFound
from aiohttp.web import FileResponse, HTTPSuccessful, Request, Response
from aiohttp.web_exceptions import HTTPException, HTTPFound, HTTPNotFound, HTTPNotModified
from django.conf import settings
from django.http import FileResponse as ApiFileResponse
from django.http import HttpResponse, HttpResponseRedirect
Expand All @@ -18,6 +18,7 @@
get_async_redis_connection,
get_redis_connection,
)
from pulpcore.app.util import check_request_was_modified
from pulpcore.metrics import artifacts_size_counter
from pulpcore.responses import ArtifactResponse

Expand Down Expand Up @@ -352,7 +353,7 @@ async def cached_function(*args, **kwargs):
await self.auth(request, self, bk)
key = self.make_key(request)
# Check cache
response = await self.make_response(key, bk)
response = await self.make_response(key, bk, request)
if response is None:
# Cache miss, create new entry
response = await self.make_entry(
Expand All @@ -373,7 +374,7 @@ def get_request_from_args(self, args):
if isinstance(arg, Request):
return arg

async def make_response(self, key, base_key):
async def make_response(self, key, base_key, request=None):
"""Tries to find the cached entry and turn it into a proper response"""
entry = await self.get(key, base_key)
if not entry:
Expand All @@ -396,21 +397,29 @@ async def make_response(self, key, base_key):
# Bad entry, delete from cache
await self.delete(key, base_key)
return None
response = self.RESPONSE_TYPES[response_type](**entry)

headers = entry.get("headers", {})
if request and not check_request_was_modified(
request, last_modified=headers.get("Last-Modified"), etag=headers.get("ETag")
):
response = HTTPNotModified(
headers={key: headers[key] for key in ("Cache-Control", "ETag") if key in headers}
)
else:
response = self.RESPONSE_TYPES[response_type](**entry)
response.headers.update({"X-PULP-CACHE": "HIT"})
return response

async def make_entry(self, key, base_key, handler, args, kwargs, expires=DEFAULT_EXPIRES_TTL):
"""Gets the response for the request and try to turn it into a cacheable entry"""
try:
response = await handler(*args, **kwargs)
except (HTTPSuccessful, HTTPFound, HTTPNotFound) as e:
except (HTTPSuccessful, HTTPFound, HTTPNotFound, HTTPNotModified) as e:
response = e

original_response = response
if isinstance(response, StreamResponse):
if hasattr(response, "future_response"):
response = response.future_response
if hasattr(response, "future_response"):
response = response.future_response

entry = {"headers": dict(response.headers), "status": response.status}

Expand Down
94 changes: 90 additions & 4 deletions pulpcore/content/handler.py

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I've added these [SRV-XYZ] "label comments" because there are 15 different return cases in the match-and-stream (without counting the serve vs stream variations), and if you wanna think about it, or take notes about this flow, it's a bit hard. Of course AI helps a lot, but still, naming things help us humans understand and manipulate them.

I'm fine with removing if this feels too personal.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What does SRV even stand for? Serve? The problem with acronyms is that no one ever knows what they mean. I would just choose a simple 1-2 word title to label the comments if that will be helpful.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yes, serve hehe I kinda find this helpful to myself, but I can see it might just cause confusion to anyone else. I'll drop it.

Original file line number Diff line number Diff line change
Expand Up @@ -18,10 +18,12 @@
HTTPFound,
HTTPMovedPermanently,
HTTPNotFound,
HTTPNotModified,
HTTPRequestRangeNotSatisfiable,
)
from asgiref.sync import sync_to_async
from django.utils import timezone
from django.utils.http import http_date
from multidict import CIMultiDict
from yarl import URL

Expand Down Expand Up @@ -57,6 +59,7 @@
)
from pulpcore.app.util import ( # noqa: E402
cache_key,
check_request_was_modified,
get_domain,
)
from pulpcore.cache import AsyncContentCache # noqa: E402
Expand All @@ -69,6 +72,12 @@
log = logging.getLogger(__name__)
_current_distribution = ContextVar("current_distribution", default=None)

# The "shared cache" (cdn/edge) should always ask pulp if their cache is still valid
EDGE_CACHE_CONTROL = "public, max-age=0, must-revalidate"

# Don't store sensitive resource on "shared cache", only on "private cache" (client)
PRIVATE_EDGE_CACHE_CONTROL = "private, max-age=0, must-revalidate"


class PathNotResolved(HTTPNotFound):
"""
Expand Down Expand Up @@ -536,6 +545,8 @@ def response_headers(path, distribution=None):
if content_type:
headers["Content-Type"] = content_type

headers["Cache-Control"] = EDGE_CACHE_CONTROL

# Let plugin-Distribution set headers for this path if it wants.
if distribution:
headers.update(distribution.content_headers_for(path))
Expand Down Expand Up @@ -594,6 +605,48 @@ def render_html(directory_list, path="", dates=None, sizes=None):
sizes=sizes,
)

@staticmethod
async def _set_last_modified_header(
headers,
last_modified: datetime | None = None,
ca=None,
rv=None,
):
"""Add the last-modified header to the response headers if not already present.

Pulp doesnt track "content served by a distribution" over time, so we need to use some
heuristics here. Lets call hypothesis 1 (H1) the common case where the distribution D is
serving content from RVs monotonically (either auto-publish, or something else) and that
the handler found a ContentArtifact (CA) matching the path for a request.

If H1 is true, we can use the repository version history as the distributed content history.
In other words, we can use the following strategy for establishing last-modified for the
path/resource in the request with a matching CA/RV:

Given a CA is provided with its RepositoryVersion (RV_n), find the earliest version RV_k
such that CA is in RV_k and use that timestamp as last-modified.

If H1 is not true, (e.g the distribution rolled back the RV it serves, changed repository, etc)
then this can produce incorrect last_modified results (e.g, say resource dist/path/to/resource
was not modified, when in fact it was).
"""

def _find_repo_add_time():
cpk = ca.content_id
rc = rv._content_relationships().filter(content_id=cpk).first()
return rc.pulp_created if rc else rv.pulp_created

if "Last-Modified" not in headers:
if last_modified is None and ca and rv:
last_modified = await sync_to_async(_find_repo_add_time)()

if last_modified is not None:
headers["Last-Modified"] = http_date(last_modified.timestamp())

@staticmethod
def _set_etag_headers(headers, sha):
headers["ETag"] = f'"{sha}"'

async def list_directory(self, repo_version, publication, path):
"""
Generate a set with directory listing of the path.
Expand Down Expand Up @@ -734,6 +787,13 @@ async def _match_and_stream(self, path, request):
content_handler_result = await sync_to_async(distro.content_handler)(original_rel_path)
if content_handler_result is not None:
if isinstance(content_handler_result, ContentArtifact):
# infer the RV which CA returned by content handler probably belongs to
__, rv_candidate, __ = await sync_to_async(
distro.get_repository_publication_and_version
)()
await self._set_last_modified_header(
headers, ca=content_handler_result, rv=rv_candidate
)
if content_handler_result.artifact:
return await self._serve_content_artifact(
content_handler_result, headers, request
Expand Down Expand Up @@ -769,6 +829,10 @@ async def _match_and_stream(self, path, request):
raise HTTPMovedPermanently(f"{request.path}/")
original_rel_path = index_path
headers = self.response_headers(original_rel_path, distro)
# last-modified heuristic assumes distribution doesn't rollback to older publications
await self._set_last_modified_header(
headers, last_modified=publication.pulp_created
)
except ObjectDoesNotExist:
dir_list, dates, sizes = await self.list_directory(None, publication, rel_path)
dir_list.update(
Expand Down Expand Up @@ -798,6 +862,8 @@ async def _match_and_stream(self, path, request):
except ObjectDoesNotExist:
pass
else:
publication_rv = publication.repository_version
await self._set_last_modified_header(headers, ca=ca, rv=publication_rv)
if ca.artifact:
return await self._serve_content_artifact(ca, headers, request)
else:
Expand Down Expand Up @@ -827,6 +893,8 @@ async def _match_and_stream(self, path, request):
except ObjectDoesNotExist:
pass
else:
publication_rv = publication.repository_version
await self._set_last_modified_header(headers, ca=ca, rv=publication_rv)
if ca.artifact:
return await self._serve_content_artifact(ca, headers, request)
else:
Expand All @@ -836,8 +904,10 @@ async def _match_and_stream(self, path, request):

# Grace-period fallback: serve from a recently-superseded publication
if distro.SERVE_FROM_PUBLICATION:
ca = await sync_to_async(distro.get_fallback_ca)(original_rel_path)
if ca is not None:
fallback = await sync_to_async(distro.get_fallback_ca)(original_rel_path)
if fallback is not None:
ca, fallback_rv = fallback
await self._set_last_modified_header(headers, ca=ca, rv=fallback_rv)
if ca.artifact:
return await self._serve_content_artifact(ca, headers, request)
else:
Expand Down Expand Up @@ -885,6 +955,7 @@ async def _match_and_stream(self, path, request):
except ObjectDoesNotExist:
pass
else:
await self._set_last_modified_header(headers, ca=ca, rv=repo_version)
if ca.artifact:
return await self._serve_content_artifact(ca, headers, request)
else:
Expand Down Expand Up @@ -1177,8 +1248,10 @@ async def _serve_content_artifact(self, content_artifact, headers, request):
Returns:
The [aiohttp.web.FileResponse][] for the file.
"""
artifact_file = content_artifact.artifact.file
artifact = content_artifact.artifact
artifact_file = artifact.file
content_length = artifact_file.size
self._set_etag_headers(headers, artifact.sha256)

try:
range_start, range_stop = request.http_range.start, request.http_range.stop
Expand All @@ -1192,10 +1265,23 @@ async def _serve_content_artifact(self, content_artifact, headers, request):
size = artifact_file.size or "*"
raise HTTPRequestRangeNotSatisfiable(headers={"Content-Range": f"bytes */{size}"})

response = self._build_response_from_content_artifact(content_artifact, headers, request)

if not check_request_was_modified(
request, last_modified=headers.get("Last-Modified"), etag=headers.get("ETag")
):
nmod_response = HTTPNotModified(
headers={key: headers[key] for key in ("Cache-Control", "ETag") if key in headers}
)
if settings.CACHE_ENABLED:
nmod_response.future_response = response
raise nmod_response

artifacts_size_counter.add(content_length)

response = self._build_response_from_content_artifact(content_artifact, headers, request)
if isinstance(response, HTTPFound):
# the response is redirect with possibly a signed URL in it
response.headers["Cache-Control"] = PRIVATE_EDGE_CACHE_CONTROL
raise response
else:
return response
Expand Down
Loading
Loading