aboutsummaryrefslogtreecommitdiff
path: root/webmentions_ssg/tasks/sender.py
diff options
context:
space:
mode:
authorDennis Fink2026-08-15 20:54:53 +0200
committerDennis Fink2026-08-15 20:54:53 +0200
commite0f65d7ff582b32f1c7b5508bba8cda807aaae5f (patch)
tree6693b9fc259bf642e0392b4d8920fbca7ffd6dda /webmentions_ssg/tasks/sender.py
parentb3e12b975064f0873bdb4d2834cae75925387d6b (diff)
downloadwebmentions-ssg-e0f65d7ff582b32f1c7b5508bba8cda807aaae5f.tar.gz
webmentions-ssg-e0f65d7ff582b32f1c7b5508bba8cda807aaae5f.zip
feat(sender): add outgoing Webmention pipeline
Scan generated h-entry documents for Webmention targets and track source revisions and delivery state in the database. Discover target endpoints, send Webmentions asynchronously, and reconcile updated, removed, restored, and deleted sources across scans. Add authenticated views for inspecting sent Webmentions and make the scanner schedule configurable.
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()