diff options
Diffstat (limited to 'webmentions_ssg/tasks/sender.py')
| -rw-r--r-- | webmentions_ssg/tasks/sender.py | 315 |
1 files changed, 315 insertions, 0 deletions
diff --git a/webmentions_ssg/tasks/sender.py b/webmentions_ssg/tasks/sender.py new file mode 100644 index 0000000..fa9b624 --- /dev/null +++ b/webmentions_ssg/tasks/sender.py @@ -0,0 +1,315 @@ +import re +import uuid +from datetime import datetime, timezone +from urllib.parse import urljoin + +import httpx +from bs4 import BeautifulSoup +from flask import current_app + +from .. import APP_NAME, VERSION +from .. import DATABASE as db +from .. import HUEY as huey +from ..models import SentWebmention, SentWebmentionStatus +from ..url_security import AddressResolutionError, is_http_url, is_public_url + +LINK_SPLIT = re.compile(r",\s*(?=<)") + +DISCOVERY_HEADERS = {"Accept": "text/html, application/xhtml+xml;q=0.9"} + + +class SenderError(Exception): + """A Webmention cannot be sent.""" + + +class TemporarySenderError(SenderError): + """A Webmention send failed for a potentially temporary reason.""" + + +class PermanentSenderError(SenderError): + """A Webmention send failed permanently for this source revision.""" + + +def temporary_http_status(status: int) -> bool: + """Return whether an HTTP status should be retried.""" + return status in {408, 425, 429} or 500 <= status <= 599 + + +def ensure_public_request(request: httpx.Request) -> None: + """Prevent requests to non-public network addresses.""" + try: + if not is_public_url(str(request.url)): + raise PermanentSenderError("Request resolves to a non-public address") + except AddressResolutionError as exc: + raise TemporarySenderError("Request hostname could not be resolved") from exc + + +def resolve_endpoint(response: httpx.Response, href: str) -> str: + """Resolve and validate a discovered Webmention endpoint.""" + endpoint = urljoin(str(response.url), href.strip()) + + if not is_http_url(endpoint): + raise PermanentSenderError(f"Invalid Webmention endpoint: {endpoint!r}") + + return endpoint + + +def parse_link_value(value: str) -> tuple[str, set[str]] | None: + """Parse a Link header value into its target and relations.""" + value = value.strip() + + if not value.startswith("<"): + return None + + href, separator, parameters = value[1:].partition(">") + + if not separator: + return None + + relations: set[str] = set() + + for parameter in parameters.split(";"): + name, separator, value = parameter.partition("=") + + if separator and name.strip().lower() == "rel": + relations.update( + relation.lower() for relation in value.strip().strip("\"'").split() + ) + + return href, relations + + +def endpoint_from_headers(response: httpx.Response) -> str | None: + """Return the first Webmention endpoint advertised by HTTP Link.""" + for header in response.headers.get_list("Link"): + for value in LINK_SPLIT.split(header): + if (link := parse_link_value(value)) is None: + continue + + href, relations = link + + if "webmention" in relations: + return resolve_endpoint(response, href) + + return None + + +def endpoint_from_html(response: httpx.Response, body: bytes) -> str | None: + """Return the first HTML Webmention endpoint in document order.""" + document = BeautifulSoup(body, "html5lib") + + for element in document.find_all(["link", "a"], href=True): + relations = element.get("rel") + + if isinstance(relations, str): + relations = relations.split() + + if not relations: + continue + + if not any( + isinstance(relation, str) and relation.lower() == "webmention" + for relation in relations + ): + continue + + href = element.get("href") + + if not isinstance(href, str): + continue + + return resolve_endpoint(response, href) + + return None + + +def read_target_body(response: httpx.Response) -> bytes: + """Read a target document up to the configured size limit.""" + max_bytes = current_app.config.get("WEBMENTIONS_SSG_MAX_TARGET_BYTES", 1_000_000) + + if (content_length := response.headers.get("Content-Length")) is not None: + try: + if int(content_length) > max_bytes: + raise PermanentSenderError("Target document is too large") + except ValueError: + pass + + body = bytearray() + + for chunk in response.iter_bytes(chunk_size=64 * 1024): + body.extend(chunk) + + if len(body) > max_bytes: + raise PermanentSenderError("Target document is too large") + + return bytes(body) + + +def discover_webmention_endpoint(client: httpx.Client, target: str) -> str | None: + """Discover the Webmention endpoint advertised by a target.""" + head_response = client.head(target, headers=DISCOVERY_HEADERS) + + if endpoint := endpoint_from_headers(head_response): + return endpoint + + with client.stream("GET", target, headers=DISCOVERY_HEADERS) as response: + if endpoint := endpoint_from_headers(response): + return endpoint + + if temporary_http_status(response.status_code): + raise TemporarySenderError(f"Target returned HTTP {response.status_code}") + + if not response.is_success: + raise PermanentSenderError(f"Target returned HTTP {response.status_code}") + + media_type = ( + response.headers.get("Content-Type", "").partition(";")[0].strip().lower() + ) + + if media_type not in {"text/html", "application/xhtml+xml"}: + return None + + return endpoint_from_html(response, read_target_body(response)) + + +def post_webmention( + client: httpx.Client, *, endpoint: str, source: str, target: str +) -> tuple[int, str | None]: + """POST a Webmention and return its status code and status URL.""" + with client.stream( + "POST", endpoint, data={"source": source, "target": target} + ) as response: + status_url = None + + if response.status_code == 201 and ( + location := response.headers.get("Location") + ): + status_url = urljoin(str(response.url), location) + + return (response.status_code, status_url) + + +def attempt_is_current(webmention: SentWebmention, revision: int) -> bool: + """Return whether an attempt still represents the desired revision.""" + db.session.refresh(webmention) + + return webmention.desired_revision == revision and ( + webmention.processed_revision is None + or webmention.processed_revision < revision + ) + + +@huey.task(retries=2, retry_delay=50) +def send_webmention(webmention_uuid: uuid.UUID) -> None: + """Discover a receiver endpoint and send one Webmention.""" + webmention = db.session.get(SentWebmention, webmention_uuid) + + if webmention is None: + current_app.logger.warning( + "Cannot send unknown SentWebmention %s", webmention_uuid + ) + return + + if ( + webmention.processed_revision is not None + and webmention.processed_revision >= webmention.desired_revision + ): + return + + revision = webmention.desired_revision + source = webmention.source.url + target = webmention.target + + webmention.last_attempted_at = datetime.now(timezone.utc) + webmention.endpoint = None + webmention.response_status = None + webmention.status_url = None + + db.session.commit() + + endpoint: str | None = None + response_status: int | None = None + + try: + with httpx.Client( + headers={"User-Agent": (f"{APP_NAME}/{VERSION} SentWebmention")}, + timeout=httpx.Timeout( + current_app.config.get("WEBMENTIONS_SSG_REQUEST_TIMEOUT", 5.0) + ), + follow_redirects=True, + max_redirects=current_app.config.get("WEBMENTIONS_SSG_MAX_REDIRECTS", 20), + trust_env=False, + event_hooks={"request": [ensure_public_request]}, + ) as client: + endpoint = discover_webmention_endpoint(client, target) + + if endpoint is None: + if not attempt_is_current(webmention, revision): + return + + webmention.processed_revision = revision + webmention.status = SentWebmentionStatus.UNSUPPORTED + webmention.failure_reason = "No Webmention endpoint discovered" + webmention.endpoint = None + webmention.response_status = None + + db.session.commit() + return + + (response_status, status_url) = post_webmention( + client, endpoint=endpoint, source=source, target=target + ) + + if not (200 <= response_status <= 299): + message = f"Webmention endpoint returned HTTP {response_status}" + + if temporary_http_status(response_status): + raise TemporarySenderError(message) + + raise PermanentSenderError(message) + + except (TemporarySenderError, httpx.RequestError) as exc: + if not attempt_is_current(webmention, revision): + return + + webmention.status = SentWebmentionStatus.FAILED + webmention.failure_reason = str(exc) or "Webmention request failed" + webmention.endpoint = endpoint + webmention.response_status = response_status + + db.session.commit() + + # Leave processed_revision unchanged. Huey may retry the + # attempt, and the periodic scanner will also rediscover + # outstanding work. + raise + + except PermanentSenderError as exc: + if not attempt_is_current(webmention, revision): + return + + webmention.processed_revision = revision + webmention.status = SentWebmentionStatus.FAILED + webmention.failure_reason = str(exc) + webmention.endpoint = endpoint + webmention.response_status = response_status + + db.session.commit() + return + + if not attempt_is_current(webmention, revision): + return + + webmention.processed_revision = revision + webmention.sent_revision = revision + + webmention.status = SentWebmentionStatus.SENT + webmention.failure_reason = None + + webmention.endpoint = endpoint + webmention.response_status = response_status + webmention.status_url = status_url + + webmention.last_sent_at = datetime.now(timezone.utc) + + db.session.commit() |
