diff options
| author | Dennis Fink | 2026-08-15 20:54:53 +0200 |
|---|---|---|
| committer | Dennis Fink | 2026-08-15 20:54:53 +0200 |
| commit | e0f65d7ff582b32f1c7b5508bba8cda807aaae5f (patch) | |
| tree | 6693b9fc259bf642e0392b4d8920fbca7ffd6dda /webmentions_ssg/tasks/scanner.py | |
| parent | b3e12b975064f0873bdb4d2834cae75925387d6b (diff) | |
| download | webmentions-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/scanner.py | 397 |
1 files changed, 397 insertions, 0 deletions
diff --git a/webmentions_ssg/tasks/scanner.py b/webmentions_ssg/tasks/scanner.py new file mode 100644 index 0000000..f80e5f4 --- /dev/null +++ b/webmentions_ssg/tasks/scanner.py @@ -0,0 +1,397 @@ +from collections.abc import Iterator +from dataclasses import dataclass +from datetime import datetime, timezone +from fnmatch import fnmatchcase +from itertools import chain +from pathlib import Path +from typing import Any +from urllib.parse import quote, urldefrag, urljoin, urlsplit + +import mf2py +import sqlalchemy as sa +import xxhash +from bs4 import BeautifulSoup, Tag +from flask import current_app +from sqlalchemy.orm import selectinload + +from .. import DATABASE as db +from .. import HUEY as huey +from ..models import SentWebmention, Source +from ..url_security import is_http_url +from .sender import send_webmention + +REACTION_PROPERTIES = ("in-reply-to", "like-of", "repost-of", "bookmark-of") + + +class SourceScanError(Exception): + """A source document cannot safely be processed.""" + + +@dataclass(frozen=True) +class ScannedSource: + path: str + url: str + content_hash: str + targets: frozenset[str] + + +def source_url_for_path(relative_path: Path, base_url: str) -> str: + """Derive the public source URL from its relative filesystem path.""" + directory = relative_path.parent.as_posix() + + return urljoin(base_url, f"{quote(directory, safe='/')}/") + + +def canonical_url( + document: BeautifulSoup, *, relative_path: Path, base_url: str | None +) -> str | None: + """Return the canonical URL declared by the document.""" + if (link := document.select_one('link[rel~="canonical"][href]')) is not None: + href = link.get("href") + + if not isinstance(href, str) or not (href := href.strip()): + raise SourceScanError("Document contains an empty canonical URL") + + if is_http_url(href): + return href + + if base_url is None: + raise SourceScanError( + "Document contains a relative canonical URL, but " + "WEBMENTIONS_SSG_SOURCE_BASE_URL is not configured" + ) + + resolved = urljoin(source_url_for_path(relative_path, base_url), href) + + if not is_http_url(resolved): + raise SourceScanError( + f"Canonical URL is not a valid HTTP or HTTPS URL: {href!r}" + ) + + return resolved + + +def property_urls(entry: dict[str, Any], property_name: str) -> Iterator[str]: + """Yield URL values from a microformats property.""" + properties = entry.get("properties") + + if not isinstance(properties, dict): + return + + values = properties.get(property_name) + + if not isinstance(values, list): + return + + yield from (value for value in values if isinstance(value, str)) + + +def parse_entry(element: Tag, base_url: str) -> dict[str, Any]: + """Parse the source h-entry with mf2py.""" + parsed = mf2py.parse(doc=str(element), url=base_url) + + if not isinstance(parsed, dict): + raise SourceScanError("Microformats parser did not return a document object") + + items = parsed.get("items") + + if not isinstance(items, list) or len(items) != 1: + raise SourceScanError("Expected exactly one parsed h-entry") + + entry = items[0] + + if not isinstance(entry, dict): + raise SourceScanError("Parsed h-entry is not an object") + + types = entry.get("type") + + if not isinstance(types, list) or "h-entry" not in types: + raise SourceScanError("Parsed microformats item is not an h-entry") + + return entry + + +def primary_entry( + document: BeautifulSoup, source_url: str +) -> tuple[Tag, dict[str, Any]]: + """Return and validate the source h-entry.""" + entries = document.find_all(class_="h-entry") + + if len(entries) != 1: + raise SourceScanError(f"Expected exactly one h-entry, found {len(entries)}") + + entry_element = entries[0] + mf2_entry = parse_entry(entry_element, source_url) + + urls = tuple(property_urls(mf2_entry, "url")) + + if source_url not in urls: + raise SourceScanError( + f"The h-entry u-url does not match the source URL {source_url!r}: {urls!r}" + ) + + return entry_element, mf2_entry + + +def content_element(entry: Tag) -> Tag: + """Return the source h-entry's e-content element.""" + contents = entry.find_all(class_="e-content") + + if len(contents) != 1: + raise SourceScanError(f"Expected exactly one e-content, found {len(contents)}") + + return contents[0] + + +def normalize_target(value: str, *, base_url: str, source_url: str) -> str | None: + """Resolve and validate a possible Webmention target URL.""" + value = value.strip() + + if not value: + return None + + target = urljoin(base_url, value) + + if not is_http_url(target): + return None + + if urldefrag(target)[0] == urldefrag(source_url)[0]: + return None + + return target + + +def iter_targets( + content: Tag, + mf2_entry: dict[str, Any], + *, + base_url: str, + source_url: str, + ignored_hostnames: tuple[str, ...] = (), +) -> Iterator[str]: + """Yield outgoing Webmention targets from an h-entry.""" + hrefs = ( + href + for element in content.find_all("a", href=True) + if isinstance(href := element.get("href"), str) + ) + + reactions = chain.from_iterable( + property_urls(mf2_entry, property_name) for property_name in REACTION_PROPERTIES + ) + + for value in chain(hrefs, reactions): + target = normalize_target(value, base_url=base_url, source_url=source_url) + if target is not None: + hostname = urlsplit(target).hostname or "" + if not any(fnmatchcase(hostname, pattern) for pattern in ignored_hostnames): + yield target + + +def scan_source_file( + path: Path, + *, + root: Path, + base_url: str | None, + ignored_hostnames: tuple[str, ...] = (), +) -> ScannedSource: + """Parse one generated source document.""" + document = BeautifulSoup(path.read_bytes(), "html5lib") + + relative_path = path.relative_to(root) + + source_url = canonical_url(document, relative_path=relative_path, base_url=base_url) + + if source_url is None: + if base_url is None: + raise SourceScanError( + "Document has no canonical URL and " + "WEBMENTIONS_SSG_SOURCE_BASE_URL is not configured" + ) + + source_url = source_url_for_path(relative_path, base_url) + + entry_element, mf2_entry = primary_entry(document, source_url) + + content = content_element(entry_element) + + return ScannedSource( + path=relative_path.as_posix(), + url=source_url, + content_hash=xxhash.xxh3_128_hexdigest(str(entry_element).encode("utf-8")), + targets=frozenset( + iter_targets( + content, + mf2_entry, + base_url=source_url, + source_url=source_url, + ignored_hostnames=ignored_hostnames, + ) + ), + ) + + +def create_source(scanned: ScannedSource, scan_time: datetime) -> Source: + """Create a source from a newly discovered document.""" + source = Source( + path=scanned.path, + url=scanned.url, + content_hash=scanned.content_hash, + revision=1, + last_seen_at=scan_time, + revised_at=scan_time, + deleted_at=None, + ) + + for target in scanned.targets: + source.sent_webmentions.append( + SentWebmention(target=target, active=True, desired_revision=1) + ) + + return source + + +def update_source( + source: Source, scanned: ScannedSource, scan_time: datetime +) -> Source: + """Update a source from a newly scanned revision.""" + if source.url != scanned.url: + raise SourceScanError( + f"Source path {source.path!r} changed public URL " + f"from {source.url!r} to {scanned.url!r}" + ) + + source.last_seen_at = scan_time + + if source.deleted_at is None and source.content_hash == scanned.content_hash: + return source + + source.revision += 1 + source.content_hash = scanned.content_hash + source.revised_at = scan_time + source.deleted_at = None + + existing = {webmention.target: webmention for webmention in source.sent_webmentions} + + for webmention in existing.values(): + webmention.active = webmention.target in scanned.targets + webmention.desired_revision = source.revision + + if not webmention.active and webmention.sent_revision is None: + webmention.processed_revision = source.revision + + for target in scanned.targets: + if target not in existing: + source.sent_webmentions.append( + SentWebmention( + target=target, active=True, desired_revision=source.revision + ) + ) + + return source + + +@huey.lock_task("scan-webmention-sources") +def scan_sources() -> None: + """Scan generated source documents and queue pending Webmentions.""" + directory = current_app.config.get("WEBMENTIONS_SSG_SOURCE_DIRECTORY") + + if directory is None: + return + + root = Path(directory).expanduser().resolve(strict=True) + + if not root.is_dir(): + raise RuntimeError(f"Source directory is not a directory: {root}") + + base_url = current_app.config.get("WEBMENTIONS_SSG_SOURCE_BASE_URL") + + if base_url is not None: + if not is_http_url(base_url): + raise RuntimeError( + "WEBMENTIONS_SSG_SOURCE_BASE_URL must be an absolute HTTP or HTTPS URL" + ) + base_url = f"{base_url.rstrip('/')}/" + + ignored_hostnames = tuple( + pattern.lower().rstrip(".") + for pattern in current_app.config.get("WEBMENTIONS_SSG_IGNORED_HOSTNAMES", ()) + ) + + sources_by_path = { + source.path: source + for source in db.session.scalars( + sa.select(Source).options(selectinload(Source.sent_webmentions)) + ) + } + + seen_paths: set[str] = set() + scanned_count = 0 + scan_time = datetime.now(timezone.utc) + + for index_path in root.rglob("index.html"): + relative_path = index_path.relative_to(root) + + if relative_path.parent == Path("."): + continue + + seen_paths.add(relative_path.as_posix()) + + try: + scanned = scan_source_file( + index_path, + root=root, + base_url=base_url, + ignored_hostnames=ignored_hostnames, + ) + + if (source := sources_by_path.get(scanned.path)) is None: + source = create_source(scanned, scan_time) + db.session.add(source) + else: + source = update_source(source, scanned, scan_time) + + sources_by_path[source.path] = source + scanned_count += 1 + + except (OSError, SourceScanError) as exc: + current_app.logger.warning("Could not scan source %s: %s", index_path, exc) + + for path, source in sources_by_path.items(): + if source.deleted_at is not None or path in seen_paths: + continue + + source.revision += 1 + source.revised_at = scan_time + source.deleted_at = scan_time + + for webmention in source.sent_webmentions: + webmention.active = False + webmention.desired_revision = source.revision + + if webmention.sent_revision is None: + webmention.processed_revision = source.revision + + # Persist the desired state before queueing any work. + db.session.commit() + + pending_count = 0 + + for identifier in db.session.scalars( + sa.select(SentWebmention.uuid) + .where( + sa.or_( + SentWebmention.processed_revision.is_(None), + SentWebmention.processed_revision < SentWebmention.desired_revision, + ) + ) + .order_by(SentWebmention.uuid) + ): + send_webmention(identifier) + pending_count += 1 + + current_app.logger.info( + "Webmention source scan complete: %d sources scanned, %d pending sends", + scanned_count, + pending_count, + ) |
