aboutsummaryrefslogtreecommitdiff
path: root/webmentions_ssg/tasks
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
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/consumer.py19
-rw-r--r--webmentions_ssg/tasks/extension.py5
-rw-r--r--webmentions_ssg/tasks/receiver.py65
-rw-r--r--webmentions_ssg/tasks/scanner.py397
-rw-r--r--webmentions_ssg/tasks/sender.py315
5 files changed, 744 insertions, 57 deletions
diff --git a/webmentions_ssg/tasks/consumer.py b/webmentions_ssg/tasks/consumer.py
index 5b85d0e..e53fce5 100644
--- a/webmentions_ssg/tasks/consumer.py
+++ b/webmentions_ssg/tasks/consumer.py
@@ -1,8 +1,19 @@
+from huey import Huey, crontab
+
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
+def create_consumer() -> Huey:
+ app = create_app()
+
+ from . import receiver, scanner, sender # noqa: E402, F401
+
+ schedule = crontab(
+ *app.config["WEBMENTIONS_SSG_SCANNER_SCHEDULE"].split(), strict=True
+ )
+ HUEY.periodic_task(schedule)(scanner.scan_sources)
+
+ return HUEY.huey
+
-huey = HUEY.huey
+huey = create_consumer()
diff --git a/webmentions_ssg/tasks/extension.py b/webmentions_ssg/tasks/extension.py
index ebdaa76..15bc669 100644
--- a/webmentions_ssg/tasks/extension.py
+++ b/webmentions_ssg/tasks/extension.py
@@ -27,10 +27,7 @@ class Huey:
huey_class, storage_kwargs = self.backend_from_url(url)
self.app = app
- self._huey = huey_class(
- **config,
- **storage_kwargs,
- )
+ self._huey = huey_class(**config, **storage_kwargs)
app.extensions["huey"] = self
diff --git a/webmentions_ssg/tasks/receiver.py b/webmentions_ssg/tasks/receiver.py
index 9475866..39def73 100644
--- a/webmentions_ssg/tasks/receiver.py
+++ b/webmentions_ssg/tasks/receiver.py
@@ -10,11 +10,7 @@ from .. import APP_NAME, VERSION
from .. import DATABASE as db
from .. import HUEY as huey
from ..models import ReceivedWebmention
-from ..url_security import (
- AddressResolutionError,
- NonPublicAddressError,
- ensure_public_url,
-)
+from ..url_security import AddressResolutionError, is_public_url
class VerificationError(Exception):
@@ -46,33 +42,20 @@ HTML_URL_ATTRIBUTES = {
"track",
"video",
},
- "cite": {
- "blockquote",
- "del",
- "ins",
- "q",
- },
+ "cite": {"blockquote", "del", "ins", "q"},
}
-def html_mentions_target(
- body: bytes,
- source_url: str,
- target_url: str,
-) -> bool:
+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_href := base_element.get("href"), str
):
- base_url = urljoin(
- source_url,
- base_href.strip(),
- )
+ base_url = urljoin(source_url, base_href.strip())
for attribute, selectors in HTML_URL_ATTRIBUTES.items():
selector = ", ".join(
@@ -85,19 +68,10 @@ def html_mentions_target(
)
for element in document.select(selector):
- if not isinstance(
- reference := element.get(attribute),
- str,
- ):
+ if not isinstance(reference := element.get(attribute), str):
continue
- if (
- urljoin(
- base_url,
- reference.strip(),
- )
- == target_url
- ):
+ if urljoin(base_url, reference.strip()) == target_url:
return True
return False
@@ -111,9 +85,10 @@ def text_mentions_target(body: str, target_url: str) -> bool:
def ensure_public_request(request: httpx.Request) -> None:
"""Prevent requests to non-public network addresses."""
try:
- ensure_public_url(str(request.url))
- except NonPublicAddressError as exc:
- raise VerificationError("Source resolves to a non-public address") from exc
+ if not is_public_url(str(request.url)):
+ raise VerificationError("Source resolves to a non-public address")
+ except ValueError as exc:
+ raise VerificationError("Source URL has no hostname") from exc
except AddressResolutionError as exc:
raise TemporaryFetchError("Source hostname could not be resolved") from exc
@@ -132,9 +107,7 @@ def fetch_source(source_url: str) -> tuple[httpx.Response, bytes]:
follow_redirects=True,
max_redirects=current_app.config.get("WEBMENTIONS_SSG_MAX_REDIRECTS", 20),
trust_env=False,
- event_hooks={
- "request": [ensure_public_request],
- },
+ event_hooks={"request": [ensure_public_request]},
) as client:
with client.stream("GET", source_url) as response:
match response.status_code:
@@ -149,8 +122,7 @@ def fetch_source(source_url: str) -> tuple[httpx.Response, bytes]:
raise VerificationError(f"Source returned HTTP {status}")
max_source_bytes = current_app.config.get(
- "WEBMENTIONS_SSG_MAX_SOURCE_BYTES",
- 1_000_000,
+ "WEBMENTIONS_SSG_MAX_SOURCE_BYTES", 1_000_000
)
if (content_length := response.headers.get("Content-Length")) is not None:
@@ -185,14 +157,10 @@ def source_mentions_target(source_url: str, target_url: str) -> bool:
case "text/plain":
try:
decoded_body = body.decode(
- response.encoding or "utf-8",
- errors="replace",
+ response.encoding or "utf-8", errors="replace"
)
except LookupError:
- decoded_body = body.decode(
- "utf-8",
- errors="replace",
- )
+ decoded_body = body.decode("utf-8", errors="replace")
return text_mentions_target(decoded_body, target_url)
case _:
raise VerificationError(
@@ -208,8 +176,7 @@ def verify_webmention(webmention_uuid: uuid.UUID) -> None:
if webmention is None:
current_app.logger.warning(
- "Cannot verify unknown ReceivedWebmention %s",
- webmention_uuid,
+ "Cannot verify unknown ReceivedWebmention %s", webmention_uuid
)
return
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,
+ )
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()