From e4a8194e6de25e8a48e2c61ce7f89a8159fd99d9 Mon Sep 17 00:00:00 2001 From: Dennis Fink Date: Sun, 9 Aug 2026 14:15:29 +0200 Subject: Implement received Webmention handling Add the Flask application setup, database models and migrations, authentication, and configuration for development and testing. Implement asynchronous Webmention verification with Huey, including HTML and plain-text source validation, retries, status tracking, and size limits. Add status, login, and paginated received-Webmention views together with comprehensive tests for forms, views, and receiver tasks. --- webmentions_ssg/tasks/__init__.py | 3 + webmentions_ssg/tasks/consumer.py | 8 ++ webmentions_ssg/tasks/extension.py | 159 ++++++++++++++++++++++++++ webmentions_ssg/tasks/receiver.py | 225 +++++++++++++++++++++++++++++++++++++ 4 files changed, 395 insertions(+) create mode 100644 webmentions_ssg/tasks/__init__.py create mode 100644 webmentions_ssg/tasks/consumer.py create mode 100644 webmentions_ssg/tasks/extension.py create mode 100644 webmentions_ssg/tasks/receiver.py (limited to 'webmentions_ssg/tasks') diff --git a/webmentions_ssg/tasks/__init__.py b/webmentions_ssg/tasks/__init__.py new file mode 100644 index 0000000..5af3987 --- /dev/null +++ b/webmentions_ssg/tasks/__init__.py @@ -0,0 +1,3 @@ +from .extension import Huey + +__all__ = ["Huey"] diff --git a/webmentions_ssg/tasks/consumer.py b/webmentions_ssg/tasks/consumer.py new file mode 100644 index 0000000..5b85d0e --- /dev/null +++ b/webmentions_ssg/tasks/consumer.py @@ -0,0 +1,8 @@ +from .. import HUEY, create_app + +app = create_app() + +# Import the tasks so they are registered with the initialized Huey instance. +from . import receiver # noqa: E402, F401 + +huey = HUEY.huey diff --git a/webmentions_ssg/tasks/extension.py b/webmentions_ssg/tasks/extension.py new file mode 100644 index 0000000..ebdaa76 --- /dev/null +++ b/webmentions_ssg/tasks/extension.py @@ -0,0 +1,159 @@ +from functools import wraps +from typing import Any, Callable +from urllib.parse import urlsplit, urlunsplit + +from flask import Flask + + +class Huey: + def __init__(self, app: Flask | None = None): + self.app: Flask | None = None + self._huey = None + + if app is not None: + self.init_app(app) + + def init_app(self, app: Flask): + config: dict[str, Any] = { + "name": app.import_name, + "results": True, + "store_none": False, + "utc": True, + "immediate": app.config.get("TESTING", False), + **app.config.get("HUEY", {}), + } + + url = app.config.get("HUEY_URL", config.pop("url", "memory://")) + huey_class, storage_kwargs = self.backend_from_url(url) + + self.app = app + self._huey = huey_class( + **config, + **storage_kwargs, + ) + + app.extensions["huey"] = self + + @property + def huey(self): + if self._huey is None: + raise RuntimeError( + "Huey has not been initialized. " + "Call huey.init_app(app) before importing tasks." + ) + return self._huey + + def task(self, *task_args: Any, **task_kwargs: Any): + def decorator(func: Callable): + @wraps(func) + def wrapper(*args: Any, **kwargs: Any): + if self.app is None: + raise RuntimeError("Flask app is not available.") + + with self.app.app_context(): + return func(*args, **kwargs) + + return self.huey.task(*task_args, **task_kwargs)(wrapper) + + return decorator + + def periodic_task(self, *task_args: Any, **task_kwargs: Any): + def decorator(func: Callable): + @wraps(func) + def wrapper(*args: Any, **kwargs: Any): + if self.app is None: + raise RuntimeError("Flask app is not available.") + + with self.app.app_context(): + return func(*args, **kwargs) + + return self.huey.periodic_task(*task_args, **task_kwargs)(wrapper) + + return decorator + + def __getattr__(self, name: str): + """ + Forward unknown attributes to the real Huey instance. + + This lets you still use things like: + huey.enqueue(...) + huey.scheduled() + huey.pending() + """ + return getattr(self.huey, name) + + @staticmethod + def backend_from_url(url: str) -> tuple[Any, dict[str, str]]: + parsed = urlsplit(url) + scheme = parsed.scheme.lower() + + if scheme.startswith("redis") or scheme.startswith("rediss"): + fixed_url = urlunsplit(parsed._replace(scheme=scheme.split("+", 1)[0])) + + if scheme.endswith("priority+expire"): + from huey import PriorityRedisExpireHuey + + return PriorityRedisExpireHuey, {"url": fixed_url} + elif scheme.endswith("priority"): + from huey import PriorityRedisHuey + + return PriorityRedisHuey, {"url": fixed_url} + elif scheme.endswith("expire"): + from huey import RedisExpireHuey + + return RedisExpireHuey, {"url": fixed_url} + else: + from huey import RedisHuey + + return RedisHuey, {"url": url} + + elif scheme == "sqlite": + from huey import SqliteHuey + + prefix = "sqlite:///" + + if not url.startswith(prefix): + raise RuntimeError( + "SQLite Huey URLs must look like sqlite:///var/huey.db" + ) + + filename = url.removeprefix(prefix) + + if not filename: + raise RuntimeError("SQLite Huey URL must include a database path.") + + return SqliteHuey, {"filename": filename} + + elif scheme == "file": + from huey import FileHuey + + prefix = "file:///" + + if not url.startswith(prefix): + raise RuntimeError( + "File Huey URLs must look like file:///var/huey-queue" + ) + + path = url.removeprefix(prefix) + + if not path: + raise RuntimeError("File Huey URL must include a directory path.") + + return FileHuey, {"path": path} + + elif scheme in {"postgres", "postgresql"}: + from huey import PostgresHuey + + return PostgresHuey, {"dsn": url} + + elif scheme == "memory": + from huey import MemoryHuey + + return MemoryHuey, {} + + elif scheme == "blackhole": + from huey import BlackHoleHuey + + return BlackHoleHuey, {} + + raise RuntimeError(f"Unsupported HUEY_URL scheme: {scheme!r}") diff --git a/webmentions_ssg/tasks/receiver.py b/webmentions_ssg/tasks/receiver.py new file mode 100644 index 0000000..ea48299 --- /dev/null +++ b/webmentions_ssg/tasks/receiver.py @@ -0,0 +1,225 @@ +import uuid +from urllib.parse import urljoin + +import httpx +import rfc3987 +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 ReceivedWebmention + + +class VerificationError(Exception): + """The source permanently failed ReceivedWebmention verification.""" + + +class SourceGoneError(VerificationError): + """The source explicitly reports that it has been removed.""" + + +class TemporaryFetchError(Exception): + """Fetching the source may succeed when retried later.""" + + +IRI_PATTERN = rfc3987.get_compiled_pattern("IRI") + + +HTML_URL_ATTRIBUTES = { + "href": {"a", "area", "link"}, + "src": { + "audio", + "embed", + "iframe", + "img", + 'input[type="image" i]', + "script", + "audio source", + "video source", + "track", + "video", + }, + "cite": { + "blockquote", + "del", + "ins", + "q", + }, +} + + +def html_mentions_target( + body: bytes, + source_url: str, + target_url: str, +) -> bool: + """Check valid HTML URL attributes for the exact target URL.""" + + document = BeautifulSoup(body, "html.parser") + base_url = source_url + + if (base_element := document.select_one("base[href]")) is not None and isinstance( + base_href := base_element.get("href"), + str, + ): + base_url = urljoin( + source_url, + base_href.strip(), + ) + + for attribute, selectors in HTML_URL_ATTRIBUTES.items(): + selector = ", ".join( + [ + "{selector}[{attribute}]".format( + selector=selector_string, attribute=attribute + ) + for selector_string in selectors + ] + ) + + for element in document.select(selector): + if not isinstance( + reference := element.get(attribute), + str, + ): + continue + + if ( + urljoin( + base_url, + reference.strip(), + ) + == target_url + ): + return True + + return False + + +def text_mentions_target(body: str, target_url: str) -> bool: + """Check whether plain text contains the exact target IRI.""" + return any(match.group() == target_url for match in IRI_PATTERN.finditer(body)) + + +def fetch_source(source_url: str) -> tuple[httpx.Response, bytes]: + """Fetch a source with limits on redirects, time, and response size.""" + + with httpx.Client( + headers={ + "Accept": "text/html, application/xhtml+xml;q=0.9, text/plain;q=0.8", + "User-Agent": f"{APP_NAME}/{VERSION} ReceivedWebmention", + }, + 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, + ) as client: + with client.stream("GET", source_url) as response: + match response.status_code: + case 200: + pass + case 410: + raise SourceGoneError("Source returned HTTP 410") + case status: + if status in {408, 425, 429} or 500 <= status <= 599: + raise TemporaryFetchError(f"Source returned HTTP {status}") + else: + raise VerificationError(f"Source returned HTTP {status}") + + max_source_bytes = current_app.config.get( + "WEBMENTIONS_SSG_MAX_SOURCE_BYTES", + 1_000_000, + ) + + if (content_length := response.headers.get("Content-Length")) is not None: + try: + if int(content_length) > max_source_bytes: + raise VerificationError("Source 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_source_bytes: + raise VerificationError("Source document is too large") + + return response, bytes(body) + + +def source_mentions_target(source_url: str, target_url: str) -> bool: + """Fetch the source and verify it according to its media type.""" + + response, body = fetch_source(source_url) + + media_type = ( + response.headers.get("Content-Type", "").partition(";")[0].strip().lower() + ) + + match media_type: + case "text/html" | "application/xhtml+xml": + return html_mentions_target(body, str(response.url), target_url) + case "text/plain": + try: + decoded_body = body.decode( + response.encoding or "utf-8", + errors="replace", + ) + except LookupError: + decoded_body = body.decode( + "utf-8", + errors="replace", + ) + return text_mentions_target(decoded_body, target_url) + case _: + raise VerificationError( + f"Unsupported source content type: {media_type or 'missing'}" + ) + + +@huey.task(retries=2, retry_delay=50) +def verify_webmention(webmention_uuid: uuid.UUID) -> None: + """Verify a ReceivedWebmention and store the result.""" + + webmention = db.session.get(ReceivedWebmention, webmention_uuid) + + if webmention is None: + current_app.logger.warning( + "Cannot verify unknown ReceivedWebmention %s", + webmention_uuid, + ) + return + + webmention.status = "verifying" + webmention.failure_reason = None + db.session.commit() + + try: + mentions_target = source_mentions_target(webmention.source, webmention.target) + except SourceGoneError as exc: + webmention.status = "deleted" + webmention.failure_reason = str(exc) + except VerificationError as exc: + webmention.status = "failed" + webmention.failure_reason = str(exc) + except (TemporaryFetchError, httpx.RequestError) as exc: + webmention.status = "failed" + webmention.failure_reason = str(exc) or "Source could not be fetched" + db.session.commit() + + # Huey retries the task because the exception escapes. + raise + else: + if mentions_target: + webmention.status = "verified" + webmention.failure_reason = None + else: + webmention.status = "deleted" + webmention.failure_reason = "Source does not mention target" + + db.session.commit() -- cgit v1.3.1