aboutsummaryrefslogtreecommitdiff
path: root/webmentions_ssg/tasks
diff options
context:
space:
mode:
Diffstat (limited to 'webmentions_ssg/tasks')
-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()