aboutsummaryrefslogtreecommitdiff
path: root/webmentions_ssg/tasks
diff options
context:
space:
mode:
Diffstat (limited to 'webmentions_ssg/tasks')
-rw-r--r--webmentions_ssg/tasks/__init__.py3
-rw-r--r--webmentions_ssg/tasks/consumer.py8
-rw-r--r--webmentions_ssg/tasks/extension.py159
-rw-r--r--webmentions_ssg/tasks/receiver.py225
4 files changed, 395 insertions, 0 deletions
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()