aboutsummaryrefslogtreecommitdiff
path: root/webmentions_ssg/tasks/sender.py
diff options
context:
space:
mode:
Diffstat (limited to '')
-rw-r--r--webmentions_ssg/tasks/sender.py315
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()