diff --git a/pesacheck_meedan_bridge/.env.example b/pesacheck_meedan_bridge/.env.example index c7208058..e5580515 100644 --- a/pesacheck_meedan_bridge/.env.example +++ b/pesacheck_meedan_bridge/.env.example @@ -1,6 +1,27 @@ +# Which CMS to read fact-checks from: ghost or superdesk. +PESACHECK_PROVIDER=ghost +# Page size, and how many articles a run with no stored article fetches. +PESACHECK_POSTS_LIMIT=15 +# Ceiling on one catch-up run, so a checkpoint far in the past (a long outage, +# or the first run after switching provider) can't import years at once. +PESACHECK_MAX_ARTICLES=100 +# Public site, used to build the article URL posted to Check. +PESACHECK_SITE_URL=https://pesacheck.org +# How a Superdesk article's URL is built. The default is the Ghost-era shape, +# which every fact-check already in Check links to. Switch to +# "{site}/fact-checks/{desk}/{slug}" once the new site serves that canonically. +PESACHECK_ARTICLE_URL_TEMPLATE={site}/{slug}/ + +# Ghost (provider: ghost) PESACHECK_URL=https://pesacheck.org PESACHECK_GHOST_CONTENT_API_KEY= -PESACHECK_GHOST_POSTS_LIMIT=15 + +# Superdesk/Publisher (provider: superdesk) +PESACHECK_SUPERDESK_GRAPHQL_URL=https://graphql-staging.pesacheck.org/v1/graphql +PESACHECK_SUPERDESK_TENANT_CODE= +# Optional: shared secret a Cloudflare WAF rule matches to skip bot protection. +PESACHECK_SUPERDESK_PRESHARED_AUTH= + PESACHECK_CHECK_URL=https://check-api.checkmedia.org/api/graphql PESACHECK_CHECK_TOKEN= PESACHECK_CHECK_WORKSPACE_SLUG=pesacheck-tipline-sandbox diff --git a/pesacheck_meedan_bridge/README.md b/pesacheck_meedan_bridge/README.md index 7e0ca90c..da2de07a 100644 --- a/pesacheck_meedan_bridge/README.md +++ b/pesacheck_meedan_bridge/README.md @@ -41,10 +41,48 @@ To run `pex` binary, execute: docker compose exec api-pesacheck_meedan_bridge ./pex ``` +## Where articles come from + +`PESACHECK_PROVIDER` selects the CMS the bridge reads: + +| Provider | Reads | Needs | +|---|---|---| +| `ghost` | Ghost Content API at `{PESACHECK_URL}/ghost/api/content/posts/` | `PESACHECK_GHOST_CONTENT_API_KEY` | +| `superdesk` | Publisher's GraphQL API (`swp_article`), filtered to the tenant and to articles carrying a `Debunk` verdict | `PESACHECK_SUPERDESK_GRAPHQL_URL`, `PESACHECK_SUPERDESK_TENANT_CODE` | + +Only the active provider's settings are required, so switching is a config +change plus a restart. Each article records the provider it came from in the +`source` column, and the fetch position is tracked per provider. + +A provider that has no articles of its own yet — the first run after a switch — +starts from the newest article stored under *any* source, i.e. wherever the +previous provider stopped. Both CMSes hold the same fact-checks under different +ids and URLs, so starting from scratch would re-post them: Check would not +reject them as duplicates, because its signature covers the fact-check URL, +which differs between the two sites. It would also skip anything published +beyond one page since the last run. + +Superdesk articles are posted to Check with a URL built from +`PESACHECK_SITE_URL` and `PESACHECK_ARTICLE_URL_TEMPLATE` (`{site}`, `{desk}`, +`{slug}`); point the site at a preview deployment to test against one. +Publisher has no canonical-URL field — `swp_redirect_route` is empty — so the +URL is assembled here, and the default keeps the Ghost-era `{site}/{slug}/` +shape that every fact-check already in Check links to. The shape matters +beyond the link: Check's duplicate signature covers the fact-check URL, so +changing it makes already-imported articles look new. + +Their language comes from Superdesk directly, and the Check tags are the +language, country, content type and harm type. + ## How articles are picked up -Each run fetches everything published since the newest PesaCheck (Ghost) article -already stored in SQLite, paginating in pages of `PESACHECK_GHOST_POSTS_LIMIT`. +Each run fetches everything published since the newest article already stored +for the active provider, oldest first, in pages of `PESACHECK_POSTS_LIMIT` and +at most `PESACHECK_MAX_ARTICLES` per run. The cap bounds a checkpoint far in +the past — a long outage, or the first run after a provider switch — and +costs nothing: the checkpoint advances as articles are stored, so the next run +picks up where this one stopped. + When nothing is stored yet, only the newest page is fetched, so the first run against a fresh database does not backfill PesaCheck's whole archive. @@ -57,6 +95,12 @@ Articles are tracked by `status`: | `Completed` | Accepted by Check, with the Check ids stored. | | `Duplicate` | Check already has this fact-check (it rejects a repeat of the same content). Terminal: never retried. The Check ids are not recorded, so find the article in Check by title if you need them. | +### Adding another provider + +Add a module exposing `name`, `fetch(since, limit)` returning raw records, and +`parse(record)` returning a `provider_base.Article`; register it in +`providers.py`. `main.py` knows nothing about any particular CMS. + ### Known limitation: backdated articles The cursor is the newest `published_at` we have stored, so an article published diff --git a/pesacheck_meedan_bridge/py/VERSION b/pesacheck_meedan_bridge/py/VERSION index 7e72641b..7db26729 100644 --- a/pesacheck_meedan_bridge/py/VERSION +++ b/pesacheck_meedan_bridge/py/VERSION @@ -1 +1 @@ -0.1.22 +0.1.26 diff --git a/pesacheck_meedan_bridge/py/database.py b/pesacheck_meedan_bridge/py/database.py index 80336be6..ac223752 100644 --- a/pesacheck_meedan_bridge/py/database.py +++ b/pesacheck_meedan_bridge/py/database.py @@ -19,12 +19,22 @@ class PesacheckFeed: check_project_media_id: str = "" check_full_url: str = "" claim_description_id: str = "" + # Which CMS the article came from: "medium", "ghost" or "superdesk". + source: str = "" + language: str = "" class PesacheckDatabase: + # Appended to pesacheck_feeds after the original columns, in this order. + ADDED_COLUMNS = ( + ("source", "TEXT NOT NULL DEFAULT ''"), + ("language", "TEXT NOT NULL DEFAULT ''"), + ) + def __init__(self): self.db_file = settings.PESACHECK_DATABASE_NAME self.create_table() + self.migrate() def create_connection(self): return sqlite3.connect(self.db_file) @@ -46,18 +56,49 @@ def create_table(self): categories TEXT DEFAULT '[]', check_project_media_id TEXT, check_full_url TEXT, - claim_description_id TEXT)""" + claim_description_id TEXT, + source TEXT NOT NULL DEFAULT '', + language TEXT NOT NULL DEFAULT '')""" ) conn.commit() finally: conn.close() + def migrate(self): + """Add columns a database created by an older version is missing.""" + conn = self.create_connection() + try: + cur = conn.cursor() + existing = { + row[1] for row in cur.execute("PRAGMA table_info(pesacheck_feeds)") + } + added = [] + for column, definition in self.ADDED_COLUMNS: + if column not in existing: + cur.execute( + f"ALTER TABLE pesacheck_feeds ADD COLUMN {column} {definition}" + ) + added.append(column) + if "source" in added: + # Rows predating the column: Medium stored the post URL as the + # guid, Ghost stored the post id. + cur.execute( + """UPDATE pesacheck_feeds + SET source = CASE WHEN guid LIKE 'http%' THEN 'medium' + ELSE 'ghost' END + WHERE source = ''""" + ) + conn.commit() + finally: + conn.close() + def insert_pesacheck_feed(self, feed): conn = self.create_connection() sql = """INSERT INTO pesacheck_feeds (title, pubDate, author, guid, link, thumbnail, description, status, categories, - check_project_media_id, check_full_url, claim_description_id) - VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""" + check_project_media_id, check_full_url, claim_description_id, + source, language) + VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""" try: cur = conn.cursor() cur.execute( @@ -75,6 +116,8 @@ def insert_pesacheck_feed(self, feed): feed.check_project_media_id, feed.check_full_url, feed.claim_description_id, + feed.source, + feed.language, ), ) conn.commit() @@ -87,7 +130,8 @@ def update_pesacheck_feed(self, guid, new_feed, expected_status=None): SET title = ?, pubDate = ?, author = ?, link = ?, thumbnail = ?, description = ?, status = ?, categories = ?, check_project_media_id = ?, check_full_url = ?, - claim_description_id = ? WHERE guid = ?""" + claim_description_id = ?, source = ?, language = ? + WHERE guid = ?""" params_tail = [] if expected_status is not None: sql += " AND status = ?" @@ -108,6 +152,8 @@ def update_pesacheck_feed(self, guid, new_feed, expected_status=None): new_feed.check_project_media_id, new_feed.check_full_url, new_feed.claim_description_id, + new_feed.source, + new_feed.language, guid, *params_tail, ), @@ -174,15 +220,22 @@ def feed_exists(self, guid): finally: conn.close() - def get_ghost_pub_dates(self): - # Legacy Medium rows use the post URL as guid; Ghost rows use the post id. + def get_pub_dates(self, source=None): + """Publication dates already stored, for one provider or for all. + + All sources is what the first run of a newly configured provider needs: + the articles exist under the old provider's source, so its position is + the only thing that says where the new one should start. + """ conn = self.create_connection() + sql = "SELECT pubDate FROM pesacheck_feeds WHERE pubDate != ''" + params = () + if source is not None: + sql += " AND source = ?" + params = (source,) try: cur = conn.cursor() - cur.execute( - "SELECT pubDate FROM pesacheck_feeds " - "WHERE guid NOT LIKE 'http%' AND pubDate != ''" - ) + cur.execute(sql, params) return [row[0] for row in cur.fetchall()] finally: conn.close() diff --git a/pesacheck_meedan_bridge/py/main.py b/pesacheck_meedan_bridge/py/main.py index 3ae43cdf..bea129ce 100755 --- a/pesacheck_meedan_bridge/py/main.py +++ b/pesacheck_meedan_bridge/py/main.py @@ -8,19 +8,15 @@ import settings from check_api import DuplicateFactCheckError, post_to_check from database import PesacheckDatabase, PesacheckFeed - - -def is_ghost_feed(feed): - # Legacy Medium rows use the post URL as guid; Ghost rows use the post id. - return not feed.guid.startswith("http") +from provider_ghost import language_from_categories +from providers import get_provider def extract_summary(feed): - # Ghost rows store the plain-text excerpt, which is the fact-check summary. - if is_ghost_feed(feed): + # Providers store a plain-text summary. Legacy Medium rows stored full + # HTML, where the summary preceded the first
. + if feed.source != "medium": return feed.description.strip() or None - # Legacy Medium rows stored full HTML, where the summary preceded the first - #
. tree = lxml.html.fromstring(feed.description) figures = tree.xpath("//figure") if len(figures) == 0 or figures[0].getprevious() is None: @@ -29,90 +25,21 @@ def extract_summary(feed): return summary_text.strip() if summary_text else None -# Identify the bridge instead of defaulting to "python-requests/x.y.z", which -# Cloudflare challenges in front of pesacheck.org. Kept separate from -# py/VERSION, which isn't packaged into the pex. -USER_AGENT = "PesaCheckMeedanBridge/1.0 (+https://pesacheck.org)" - -language_codes = { - "english": "en", - "french": "fr", - "oromo": "om", - "afaan": "om", - "afaan oromoo": "om", - "swahili": "sw", - "kiswahili": "sw", - "amharic": "am", - "somali": "so", - "somaaliga": "so", - "tigrinya": "ti", - "arabic": "ar", -} - - -def parse_ghost_post(post): - # Internal Ghost tags (e.g. #hash-tags) are for site organisation only. - tags = [ - tag["name"] - for tag in post.get("tags") or [] - if tag.get("visibility") == "public" - ] - authors = [author["name"] for author in post.get("authors") or []] - return { - "title": post["title"], - "pubDate": post.get("published_at") or "", - "author": ", ".join(authors), - "guid": post["id"], - "link": post["url"], - "thumbnail": post.get("feature_image") or "", - "description": ( - post.get("custom_excerpt") or post.get("excerpt") or "" - ).strip(), - "categories": tags, - } - - -def get_checkpoint(db): +def get_checkpoint(db, source=None): pub_dates = [] - for pub_date in db.get_ghost_pub_dates(): + for pub_date in db.get_pub_dates(source): try: - pub_dates.append(datetime.fromisoformat(pub_date)) + parsed = datetime.fromisoformat(pub_date) except ValueError: continue + if parsed.tzinfo is None: + # Legacy Medium rows stored "2024-11-18 23:19:22" (UTC, unmarked); + # without this they can't be compared with the others. + parsed = parsed.replace(tzinfo=UTC) + pub_dates.append(parsed) return max(pub_dates) if pub_dates else None -def fetch_from_pesacheck(since=None): - # Without a checkpoint (first run against Ghost), only fetch the newest - # page instead of backfilling PesaCheck's entire archive. - url = f"{settings.PESACHECK_URL.rstrip('/')}/ghost/api/content/posts/" - params = { - "key": settings.PESACHECK_GHOST_CONTENT_API_KEY, - "limit": settings.PESACHECK_GHOST_POSTS_LIMIT, - "order": "published_at desc", - "include": "tags,authors", - "fields": "id,title,url,excerpt,custom_excerpt,feature_image,published_at", - } - if since: - since_utc = since.astimezone(UTC).strftime("%Y-%m-%d %H:%M:%S") - params["filter"] = f"published_at:>='{since_utc}'" - headers = {"Accept-Version": "v5.0", "User-Agent": USER_AGENT} - posts = [] - page = 1 - while page: - params["page"] = page - response = requests.get(url, params=params, headers=headers, timeout=60) - if response.status_code != 200: - raise Exception( - f"An Error Occurred fetching data from pesacheck: {response.text}" - ) - data = response.json() - posts.extend(data.get("posts") or []) - pagination = (data.get("meta") or {}).get("pagination") or {} - page = pagination.get("next") if since else None - return posts - - def store_in_database(feed, db): db.insert_pesacheck_feed(feed) return feed @@ -120,12 +47,8 @@ def store_in_database(feed, db): def build_check_input(feed): categories = json.loads(feed.categories) - codes = [ - language_codes[language.lower()] - for language in categories - if language.lower() in language_codes - ] - language = "en" if not codes else codes[0] + # Rows stored before the language column fall back to the tag names. + language = feed.language or language_from_categories(categories) or "en" claim_description = feed.title summary = extract_summary(feed) or "Not Found" return { @@ -221,32 +144,45 @@ def main(db): ) for pending in db.get_pesacheck_feeds_by_status("Pending"): post_and_record(pending, db, success_posts, duplicates) - from_pesacheck = fetch_from_pesacheck(since=get_checkpoint(db)) - # Oldest first, so the checkpoint never moves past an unstored article. - for post in reversed(from_pesacheck): + provider = get_provider(settings.PESACHECK_PROVIDER) + # On a provider's first run, carry on from wherever the previous one + # stopped: the same fact-checks exist in both CMSes, and re-posting + # them would create duplicate published reports in Check rather than + # being rejected (Check's signature covers the URL, which differs). + since = get_checkpoint(db, provider.name) or get_checkpoint(db) + from_pesacheck = provider.fetch( + since=since, + limit=settings.PESACHECK_POSTS_LIMIT, + max_articles=settings.PESACHECK_MAX_ARTICLES, + ) + # Providers return oldest first, so the checkpoint never moves past an + # unstored article. + for post in from_pesacheck: try: - item = parse_ghost_post(post) + article = provider.parse(post) except Exception as exception: # A malformed post is skipped: holding the checkpoint behind it # would block every later article indefinitely. sentry_sdk.capture_exception(exception) continue try: - if db.feed_exists(item["guid"]): + if db.feed_exists(article.guid): continue feed = PesacheckFeed( - title=item["title"], - pubDate=item["pubDate"], - author=item["author"], - guid=item["guid"], - link=item["link"], - categories=json.dumps(item["categories"]), - thumbnail=item["thumbnail"], - description=item["description"], + title=article.title, + pubDate=article.published_at, + author=article.author, + guid=article.guid, + link=article.url, + categories=json.dumps(article.categories), + thumbnail=article.thumbnail, + description=article.summary, status="Pending", check_project_media_id="", check_full_url="", claim_description_id="", + source=provider.name, + language=article.language, ) store_in_database(feed, db=db) except Exception as exception: diff --git a/pesacheck_meedan_bridge/py/provider_base.py b/pesacheck_meedan_bridge/py/provider_base.py new file mode 100644 index 00000000..425fa411 --- /dev/null +++ b/pesacheck_meedan_bridge/py/provider_base.py @@ -0,0 +1,32 @@ +from dataclasses import dataclass, field + +import lxml.html # nosec B410 + +# Identify the bridge instead of defaulting to "python-requests/x.y.z", which +# Cloudflare challenges in front of pesacheck.org. Kept separate from +# py/VERSION, which isn't packaged into the pex. +USER_AGENT = "PesaCheckMeedanBridge/1.0 (+https://pesacheck.org)" + + +@dataclass +class Article: + """A fact-check as the bridge sees it, whichever CMS it came from.""" + + guid: str + title: str + url: str + published_at: str + # Plain text: each provider strips its own markup. + summary: str = "" + categories: list = field(default_factory=list) + # ISO code ("en", "so", ...); empty when the provider can't tell. + language: str = "" + author: str = "" + thumbnail: str = "" + + +def html_to_text(value): + """Collapse an HTML fragment to plain text. Safe on plain text and None.""" + if not value or not value.strip(): + return "" + return lxml.html.fromstring(value).text_content().strip() diff --git a/pesacheck_meedan_bridge/py/provider_ghost.py b/pesacheck_meedan_bridge/py/provider_ghost.py new file mode 100644 index 00000000..747de0db --- /dev/null +++ b/pesacheck_meedan_bridge/py/provider_ghost.py @@ -0,0 +1,99 @@ +from datetime import UTC + +import requests +import settings +from provider_base import USER_AGENT, Article + +# Tag names PesaCheck publishes under, mapped to the fact-check language. +LANGUAGE_CODES = { + "english": "en", + "french": "fr", + "oromo": "om", + "afaan": "om", + "afaan oromoo": "om", + "swahili": "sw", + "kiswahili": "sw", + "amharic": "am", + "somali": "so", + "somaaliga": "so", + "tigrinya": "ti", + "arabic": "ar", +} + + +def language_from_categories(categories): + for category in categories: + code = LANGUAGE_CODES.get(category.lower()) + if code: + return code + return "" + + +class GhostProvider: + """PesaCheck on Ghost, read through the (public, read-only) Content API.""" + + name = "ghost" + + def __init__(self): + if not settings.PESACHECK_GHOST_CONTENT_API_KEY: + raise ValueError( + "PESACHECK_GHOST_CONTENT_API_KEY is required when " + "PESACHECK_PROVIDER is 'ghost'" + ) + + def fetch(self, since=None, limit=15, max_articles=None): + # Oldest first when catching up from a checkpoint, newest page when + # there is no checkpoint. Returns oldest first either way. + url = f"{settings.PESACHECK_URL.rstrip('/')}/ghost/api/content/posts/" + params = { + "key": settings.PESACHECK_GHOST_CONTENT_API_KEY, + "limit": limit, + "order": "published_at asc" if since else "published_at desc", + "include": "tags,authors", + "fields": "id,title,url,excerpt,custom_excerpt,feature_image,published_at", + } + if since: + since_utc = since.astimezone(UTC).strftime("%Y-%m-%d %H:%M:%S") + params["filter"] = f"published_at:>='{since_utc}'" + headers = {"Accept-Version": "v5.0", "User-Agent": USER_AGENT} + posts = [] + page = 1 + while page: + params["page"] = page + response = requests.get(url, params=params, headers=headers, timeout=60) + if response.status_code != 200: + raise Exception( + f"An Error Occurred fetching data from pesacheck: {response.text}" + ) + data = response.json() + posts.extend(data.get("posts") or []) + if max_articles is not None and len(posts) >= max_articles: + # Break rather than return: the single exit below is what puts + # a newest-first page back into oldest-first order. + posts = posts[:max_articles] + break + pagination = (data.get("meta") or {}).get("pagination") or {} + page = pagination.get("next") if since else None + # The no-checkpoint page came newest first. + return posts if since else list(reversed(posts)) + + def parse(self, post): + # Internal Ghost tags (e.g. #hash-tags) are for site organisation only. + categories = [ + tag["name"] + for tag in post.get("tags") or [] + if tag.get("visibility") == "public" + ] + authors = [author["name"] for author in post.get("authors") or []] + return Article( + guid=post["id"], + title=post["title"], + url=post["url"], + published_at=post.get("published_at") or "", + # Ghost excerpts are already plain text. + summary=(post.get("custom_excerpt") or post.get("excerpt") or "").strip(), + categories=categories, + language=language_from_categories(categories), + author=", ".join(authors), + thumbnail=post.get("feature_image") or "", + ) diff --git a/pesacheck_meedan_bridge/py/provider_superdesk.py b/pesacheck_meedan_bridge/py/provider_superdesk.py new file mode 100644 index 00000000..17739192 --- /dev/null +++ b/pesacheck_meedan_bridge/py/provider_superdesk.py @@ -0,0 +1,189 @@ +import json +from datetime import UTC + +import requests +import settings +from provider_base import USER_AGENT, Article, html_to_text + +# An article is a fact-check iff it carries a verdict from the "Debunk" +# vocabulary. Without this clause the query also returns homepage blocks, team +# profiles and Media Centre entries, which publish to the same routes. +DEBUNK_SCHEME = "Debunk" + +# Subject vocabularies that become Check tags: language, country, the kind of +# fact-check, and the harm it addresses. Order is the tag order. +CATEGORY_SCHEMES = ["Debunklang", "countries", "content_type", "Harm_type"] + +FACT_CHECKS_QUERY = """ + query FactChecks( + $where: swp_article_bool_exp! + $limit: Int! + $offset: Int! + $order: order_by! + ) { + items: swp_article( + where: $where + order_by: { published_at: $order } + limit: $limit + offset: $offset + ) { + id + title + slug + lead + published_at + metadata + swp_route { + slug + } + } + } +""" + + +def parse_metadata(raw): + # Hasura exposes the jsonb `metadata` column as a JSON-encoded string. + if not raw: + return {} + try: + parsed = json.loads(raw) + except ValueError: + return {} + return parsed if isinstance(parsed, dict) else {} + + +def subject_names(metadata, scheme): + return [ + subject["name"] + for subject in metadata.get("subject") or [] + if subject.get("scheme") == scheme and subject.get("name") + ] + + +class SuperdeskProvider: + """PesaCheck on Superdesk, read through Publisher's GraphQL API.""" + + name = "superdesk" + + def __init__(self): + if not settings.PESACHECK_SUPERDESK_GRAPHQL_URL: + raise ValueError( + "PESACHECK_SUPERDESK_GRAPHQL_URL is required when " + "PESACHECK_PROVIDER is 'superdesk'" + ) + if not settings.PESACHECK_SUPERDESK_TENANT_CODE: + raise ValueError( + "PESACHECK_SUPERDESK_TENANT_CODE is required when " + "PESACHECK_PROVIDER is 'superdesk'" + ) + + def build_where(self, since=None): + published_at = {"_is_null": False} + if since: + # published_at is stored without an offset; the API treats it as UTC. + published_at["_gte"] = since.astimezone(UTC).strftime("%Y-%m-%dT%H:%M:%S") + return { + "tenant_code": {"_eq": settings.PESACHECK_SUPERDESK_TENANT_CODE}, + "published_at": published_at, + "_and": [ + { + "swp_article_metadata": { + "swp_article_metadata_subjects": { + "scheme": {"_eq": DEBUNK_SCHEME} + } + } + } + ], + } + + def fetch(self, since=None, limit=15, max_articles=None): + # Oldest first when catching up from a checkpoint, newest page when + # there is no checkpoint. Returns oldest first either way. + headers = {"Content-Type": "application/json", "User-Agent": USER_AGENT} + if settings.PESACHECK_SUPERDESK_PRESHARED_AUTH: + # A Cloudflare WAF rule skips bot protection when this matches. + headers["x-preshared-auth"] = settings.PESACHECK_SUPERDESK_PRESHARED_AUTH + where = self.build_where(since) + order = "asc" if since else "desc" + articles = [] + offset = 0 + while True: + body = { + "query": FACT_CHECKS_QUERY, + "variables": { + "where": where, + "limit": limit, + "offset": offset, + "order": order, + }, + } + response = requests.post( + settings.PESACHECK_SUPERDESK_GRAPHQL_URL, + json=body, + headers=headers, + timeout=60, + ) + if response.status_code != 200: + raise Exception( + f"An Error Occurred fetching data from pesacheck: {response.text}" + ) + data = response.json() + if data.get("errors"): + raise Exception( + f"An Error Occurred fetching data from pesacheck: {response.text}" + ) + page = ((data.get("data") or {}).get("items")) or [] + articles.extend(page) + if max_articles is not None and len(articles) >= max_articles: + # Break rather than return: the single exit below is what puts + # a newest-first page back into oldest-first order. + articles = articles[:max_articles] + break + if not since or len(page) < limit: + break + offset += limit + # The no-checkpoint page came newest first. + return articles if since else list(reversed(articles)) + + def parse(self, article): + metadata = parse_metadata(article.get("metadata")) + categories = [] + for scheme in CATEGORY_SCHEMES: + categories.extend(subject_names(metadata, scheme)) + published_at = article.get("published_at") or "" + if published_at and not published_at.endswith("Z") and "+" not in published_at: + # Naive in the API, UTC in fact; store it unambiguously. + published_at = f"{published_at}+00:00" + return Article( + # The Superdesk guid survives a re-import; the numeric id may not. + guid=str(metadata.get("guid") or article["id"]), + title=article["title"], + url=self.article_url(article), + published_at=published_at, + summary=html_to_text(article.get("lead")), + categories=categories, + language=metadata.get("language") or "", + author=metadata.get("byline") or "", + # Check never receives a thumbnail, so the renditions aren't worth + # resolving here. + thumbnail="", + ) + + def article_url(self, article): + """The public URL posted to Check, per PESACHECK_ARTICLE_URL_TEMPLATE. + + Publisher has no canonical-URL field (swp_redirect_route is empty), so + the URL is built here. It matters beyond the link itself: Check's + duplicate signature covers the fact-check URL, so changing the shape + makes already-imported articles look new. + """ + desk = (article.get("swp_route") or {}).get("slug") or "" + template = settings.PESACHECK_ARTICLE_URL_TEMPLATE + if "{desk}" in template and not desk: + # An article with no route can't have a desk in its path. + template = "{site}/{slug}/" + return template.format( + site=settings.PESACHECK_SITE_URL.rstrip("/"), + desk=desk, + slug=article["slug"], + ) diff --git a/pesacheck_meedan_bridge/py/providers.py b/pesacheck_meedan_bridge/py/providers.py new file mode 100644 index 00000000..374846a6 --- /dev/null +++ b/pesacheck_meedan_bridge/py/providers.py @@ -0,0 +1,19 @@ +from provider_ghost import GhostProvider +from provider_superdesk import SuperdeskProvider + +PROVIDERS = { + GhostProvider.name: GhostProvider, + SuperdeskProvider.name: SuperdeskProvider, +} + + +def get_provider(name): + """The configured CMS the bridge reads fact-checks from.""" + try: + provider = PROVIDERS[name] + except KeyError: + known = ", ".join(sorted(PROVIDERS)) + raise ValueError( + f"Unknown PESACHECK_PROVIDER {name!r}. Known providers: {known}" + ) from None + return provider() diff --git a/pesacheck_meedan_bridge/py/settings.py b/pesacheck_meedan_bridge/py/settings.py index ea5b8ec3..347c6c98 100644 --- a/pesacheck_meedan_bridge/py/settings.py +++ b/pesacheck_meedan_bridge/py/settings.py @@ -19,9 +19,41 @@ profiles_sample_rate=1.0, ) +# Which CMS to read fact-checks from: "ghost" or "superdesk". Each provider +# validates its own settings, so only the active one's are required. +PESACHECK_PROVIDER = env("PESACHECK_PROVIDER", "ghost") + +# Page size, and the number of articles a run with no checkpoint fetches. +# Falls back to the old Ghost-only name so existing deployments keep theirs. +PESACHECK_POSTS_LIMIT = env.int( + "PESACHECK_POSTS_LIMIT", env.int("PESACHECK_GHOST_POSTS_LIMIT", 15) +) + +# Ceiling on one run when catching up from a checkpoint. Without it, a +# checkpoint far in the past (a long outage, or the first run after switching +# provider) would import years of articles, posting all of them to Check. +# Catch-up then continues on the next run, since the checkpoint advances. +PESACHECK_MAX_ARTICLES = env.int("PESACHECK_MAX_ARTICLES", 100) + +# Public site, used to build the article URL posted to Check. +PESACHECK_SITE_URL = env("PESACHECK_SITE_URL", "https://pesacheck.org") + +# How a Superdesk article's public URL is built, as a template over +# {site}, {desk} (the Publisher route slug) and {slug}. The default is the +# Ghost-era shape, which is what every fact-check already in Check links to; +# switch it to "{site}/fact-checks/{desk}/{slug}" once the new site serves +# that as the canonical URL. +PESACHECK_ARTICLE_URL_TEMPLATE = env("PESACHECK_ARTICLE_URL_TEMPLATE", "{site}/{slug}/") + +# Ghost PESACHECK_URL = env("PESACHECK_URL", "https://pesacheck.org") -PESACHECK_GHOST_CONTENT_API_KEY = env("PESACHECK_GHOST_CONTENT_API_KEY") -PESACHECK_GHOST_POSTS_LIMIT = env.int("PESACHECK_GHOST_POSTS_LIMIT", 15) +PESACHECK_GHOST_CONTENT_API_KEY = env("PESACHECK_GHOST_CONTENT_API_KEY", None) + +# Superdesk (Publisher's GraphQL API) +PESACHECK_SUPERDESK_GRAPHQL_URL = env("PESACHECK_SUPERDESK_GRAPHQL_URL", None) +PESACHECK_SUPERDESK_TENANT_CODE = env("PESACHECK_SUPERDESK_TENANT_CODE", None) +# Shared secret a Cloudflare WAF rule matches to skip bot protection. +PESACHECK_SUPERDESK_PRESHARED_AUTH = env("PESACHECK_SUPERDESK_PRESHARED_AUTH", None) PESACHECK_CHECK_URL = env("PESACHECK_CHECK_URL") PESACHECK_CHECK_TOKEN = env("PESACHECK_CHECK_TOKEN") diff --git a/pesacheck_meedan_bridge/py/test_bridge.py b/pesacheck_meedan_bridge/py/test_bridge.py index d6470057..12ecd82c 100644 --- a/pesacheck_meedan_bridge/py/test_bridge.py +++ b/pesacheck_meedan_bridge/py/test_bridge.py @@ -1,7 +1,7 @@ """Tests for the bridge. Run with `pants test pesacheck_meedan_bridge/py::`. -Nothing here touches the network or a real database: Ghost and Check are -mocked, and every test gets its own SQLite file. +Nothing here touches the network or a real database: the CMS APIs and Check +are mocked, and every test gets its own SQLite file. """ import json @@ -20,9 +20,13 @@ os.environ.update( { "PESACHECK_SENTRY_DSN": "", + "PESACHECK_PROVIDER": "ghost", "PESACHECK_URL": "https://pesacheck.org", + "PESACHECK_SITE_URL": "https://pesacheck.org", "PESACHECK_GHOST_CONTENT_API_KEY": "key", - "PESACHECK_GHOST_POSTS_LIMIT": "2", + "PESACHECK_POSTS_LIMIT": "2", + "PESACHECK_SUPERDESK_GRAPHQL_URL": "https://graphql.invalid/v1/graphql", + "PESACHECK_SUPERDESK_TENANT_CODE": "123abc", "PESACHECK_CHECK_URL": "https://check.invalid/graphql", "PESACHECK_CHECK_TOKEN": "t", # nosec B105 - placeholder, nothing is called "PESACHECK_CHECK_WORKSPACE_SLUG": "ws", @@ -33,6 +37,10 @@ import check_api # noqa: E402 import database # noqa: E402 import main # noqa: E402 +import provider_base # noqa: E402 +import provider_ghost # noqa: E402 +import provider_superdesk # noqa: E402 +import providers # noqa: E402 import settings # noqa: E402 @@ -62,6 +70,40 @@ def ghost_response(posts, next_page=None): return resp +def superdesk_article(i, day=1, **extra): + metadata = { + "guid": f"uuid-{i}", + "language": "so", + "byline": "A", + "subject": [ + {"code": "false", "name": "False", "scheme": "Debunk"}, + {"code": "KEN", "name": "Kenya", "scheme": "countrymention1"}, + {"code": "KEN", "name": "Kenya", "scheme": "countries"}, + {"code": "quickread", "name": "Quick Read", "scheme": "content_type"}, + {"code": "debunkso", "name": "Somali", "scheme": "Debunklang"}, + {"code": "sports", "name": "Sports", "scheme": "Harm_type"}, + ], + } + metadata.update(extra.pop("metadata", None) or {}) + article = { + "id": 6000 + i, + "title": f"Post {i}", + "slug": f"post-{i}", + "lead": f"

Summary {i}

", + "published_at": f"2026-09-{day:02d}T10:00:00", + "metadata": json.dumps(metadata), + "swp_route": {"slug": "somali"}, + } + article.update(extra) + return article + + +def superdesk_response(articles): + resp = mock.Mock(status_code=200) + resp.json.return_value = {"data": {"items": articles}} + return resp + + def check_response(i): return { "data": { @@ -77,10 +119,14 @@ def check_response(i): class Base(unittest.TestCase): + provider = "ghost" + def setUp(self): fd, self.db_file = tempfile.mkstemp(suffix=".db", dir=TMP) os.close(fd) settings.PESACHECK_DATABASE_NAME = self.db_file + settings.PESACHECK_PROVIDER = self.provider + self.addCleanup(setattr, settings, "PESACHECK_PROVIDER", "ghost") self.db = database.PesacheckDatabase() self.posted = [] @@ -97,19 +143,21 @@ def fake_post(data): self.sentry_msg = mock.patch.object(main.sentry_sdk, "capture_message").start() self.addCleanup(mock.patch.stopall) - def run_with(self, posts): - with mock.patch.object( - main.requests, "get", return_value=ghost_response(posts) - ): - main.main(self.db) - def rows(self): conn = sqlite3.connect(self.db_file) rows = conn.execute("SELECT guid, status FROM pesacheck_feeds").fetchall() conn.close() return dict(rows) - def add_row(self, guid, pub_date, status="Completed", description="x"): + def sources(self): + conn = sqlite3.connect(self.db_file) + rows = conn.execute("SELECT guid, source FROM pesacheck_feeds").fetchall() + conn.close() + return dict(rows) + + def add_row( + self, guid, pub_date, status="Completed", description="x", source="ghost" + ): self.db.insert_pesacheck_feed( database.PesacheckFeed( title="t", @@ -121,32 +169,63 @@ def add_row(self, guid, pub_date, status="Completed", description="x"): description=description, status=status, categories="[]", + source=source, ) ) + def run_ghost(self, posts, next_page=None): + with mock.patch.object( + provider_ghost.requests, + "get", + return_value=ghost_response(posts, next_page), + ): + main.main(self.db) + + +class TestProviderRegistry(unittest.TestCase): + def test_known_providers(self): + self.assertIsInstance( + providers.get_provider("ghost"), provider_ghost.GhostProvider + ) + self.assertIsInstance( + providers.get_provider("superdesk"), provider_superdesk.SuperdeskProvider + ) + + def test_unknown_provider_names_the_known_ones(self): + with self.assertRaises(ValueError) as ctx: + providers.get_provider("wordpress") + self.assertIn("ghost, superdesk", str(ctx.exception)) + + def test_provider_validates_its_own_settings(self): + with mock.patch.object(settings, "PESACHECK_GHOST_CONTENT_API_KEY", None): + with self.assertRaises(ValueError): + providers.get_provider("ghost") + # A missing Ghost key does not stop Superdesk running. + providers.get_provider("superdesk") + with mock.patch.object(settings, "PESACHECK_SUPERDESK_TENANT_CODE", None): + with self.assertRaises(ValueError): + providers.get_provider("superdesk") -class TestPagination(Base): + +class TestGhostProvider(Base): def test_first_run_fetches_single_page_without_filter(self): with mock.patch.object( - main.requests, + provider_ghost.requests, "get", return_value=ghost_response([ghost_post(2), ghost_post(1)], next_page=2), ) as get: - posts = main.fetch_from_pesacheck(since=None) + posts = provider_ghost.GhostProvider().fetch(since=None, limit=2) self.assertEqual(len(posts), 2) self.assertEqual(get.call_count, 1) self.assertNotIn("filter", get.call_args.kwargs["params"]) - self.assertEqual( - get.call_args.args[0], "https://pesacheck.org/ghost/api/content/posts/" - ) def test_sends_identifying_user_agent(self): with mock.patch.object( - main.requests, "get", return_value=ghost_response([]) + provider_ghost.requests, "get", return_value=ghost_response([]) ) as get: - main.fetch_from_pesacheck() + provider_ghost.GhostProvider().fetch() headers = get.call_args.kwargs["headers"] - self.assertEqual(headers["User-Agent"], main.USER_AGENT) + self.assertEqual(headers["User-Agent"], provider_base.USER_AGENT) self.assertNotIn("python-requests", headers["User-Agent"]) def test_checkpoint_follows_next_until_exhausted(self): @@ -155,46 +234,48 @@ def test_checkpoint_follows_next_until_exhausted(self): ghost_response([ghost_post(3), ghost_post(2)], next_page=3), ghost_response([ghost_post(1)], next_page=None), ] - seen_pages = [] + seen = [] def fake_get(url, params, headers, timeout): - seen_pages.append((params["page"], params["filter"])) + seen.append((params["page"], params["filter"])) return pages[params["page"] - 1] - since = datetime( - 2026, - 9, - 1, - 13, - 0, - tzinfo=datetime.fromisoformat("2026-01-01T00:00+03:00").tzinfo, - ) - with mock.patch.object(main.requests, "get", side_effect=fake_get): - posts = main.fetch_from_pesacheck(since=since) + tz = datetime.fromisoformat("2026-01-01T00:00+03:00").tzinfo + since = datetime(2026, 9, 1, 13, 0, tzinfo=tz) + with mock.patch.object(provider_ghost.requests, "get", side_effect=fake_get): + posts = provider_ghost.GhostProvider().fetch(since=since, limit=2) self.assertEqual( [p["title"] for p in posts], [f"Post {i}" for i in (5, 4, 3, 2, 1)] ) - # Converted to UTC. - self.assertEqual(seen_pages[0], (1, "published_at:>='2026-09-01 10:00:00'")) - self.assertEqual([p for p, _ in seen_pages], [1, 2, 3]) - - def test_checkpoint_ignores_legacy_rows(self): - self.add_row("https://medium.com/p/abc", "2030-01-01 00:00:00") - self.assertIsNone(main.get_checkpoint(self.db)) + self.assertEqual(seen[0], (1, "published_at:>='2026-09-01 10:00:00'")) + + def test_parse_maps_tags_and_language(self): + article = provider_ghost.GhostProvider().parse(ghost_post(1)) + self.assertEqual(article.guid, ghost_post(1)["id"]) + self.assertEqual(article.categories, ["Somali"]) + self.assertEqual(article.language, "so") + self.assertEqual(article.summary, "Summary 1") + self.assertEqual(article.url, "https://pesacheck.org/post-1/") + + def test_checkpoint_is_scoped_to_the_provider(self): + self.add_row("https://medium.com/p/abc", "2030-01-01 00:00:00", source="medium") + self.assertIsNone(main.get_checkpoint(self.db, "ghost")) self.add_row("a" * 24, "2026-09-01T10:00:00.000+00:00") self.add_row("b" * 24, "2026-09-03T10:00:00.000+00:00") self.assertEqual( - main.get_checkpoint(self.db), datetime(2026, 9, 3, 10, tzinfo=UTC) + main.get_checkpoint(self.db, "ghost"), datetime(2026, 9, 3, 10, tzinfo=UTC) ) + self.assertIsNone(main.get_checkpoint(self.db, "superdesk")) def test_burst_larger_than_limit_is_fully_imported_oldest_first(self): self.add_row("0" * 24, "2026-09-01T10:00:00.000+00:00") + # Catching up, so Ghost is asked for oldest first. pages = [ - ghost_response([ghost_post(5, 5), ghost_post(4, 4)], next_page=2), - ghost_response([ghost_post(3, 3), ghost_post(2, 2)], next_page=None), + ghost_response([ghost_post(2, 2), ghost_post(3, 3)], next_page=2), + ghost_response([ghost_post(4, 4), ghost_post(5, 5)], next_page=None), ] with mock.patch.object( - main.requests, + provider_ghost.requests, "get", side_effect=lambda url, params, **kw: pages[params["page"] - 1], ): @@ -202,16 +283,16 @@ def test_burst_larger_than_limit_is_fully_imported_oldest_first(self): self.assertEqual( [d["title"] for d in self.posted], ["Post 2", "Post 3", "Post 4", "Post 5"] ) - self.assertTrue(all(v == "Completed" for v in self.rows().values())) self.assertEqual(self.posted[0]["language"], "so") self.assertEqual(self.posted[0]["set_tags"], ["Somali"]) + self.assertEqual(set(self.sources().values()), {"ghost"}) def test_fetch_error_on_later_page_stores_nothing(self): self.add_row("0" * 24, "2026-09-01T10:00:00.000+00:00") err = mock.Mock(status_code=500, text="boom") - pages = [ghost_response([ghost_post(5, 5)], next_page=2), err] + pages = [ghost_response([ghost_post(2, 2)], next_page=2), err] with mock.patch.object( - main.requests, + provider_ghost.requests, "get", side_effect=lambda url, params, **kw: pages[params["page"] - 1], ): @@ -219,7 +300,424 @@ def test_fetch_error_on_later_page_stores_nothing(self): main.main(self.db) self.assertEqual(self.posted, []) self.assertEqual(len(self.rows()), 1) - self.sentry_exc.assert_called_once() + + +class TestSuperdeskProvider(Base): + provider = "superdesk" + + def run_superdesk(self, articles): + with mock.patch.object( + provider_superdesk.requests, + "post", + return_value=superdesk_response(articles), + ): + main.main(self.db) + + def test_parse_maps_the_staging_shape(self): + article = provider_superdesk.SuperdeskProvider().parse(superdesk_article(1)) + self.assertEqual(article.guid, "uuid-1") + self.assertEqual(article.title, "Post 1") + self.assertEqual(article.url, "https://pesacheck.org/post-1/") + self.assertEqual(article.summary, "Summary 1") # HTML stripped + self.assertEqual(article.language, "so") + self.assertEqual( + article.categories, ["Somali", "Kenya", "Quick Read", "Sports"] + ) + self.assertEqual(article.published_at, "2026-09-01T10:00:00+00:00") + self.assertEqual(article.author, "A") + + def test_url_base_is_configurable(self): + with mock.patch.object( + settings, "PESACHECK_SITE_URL", "https://pesacheck-ui.vercel.app/" + ): + article = provider_superdesk.SuperdeskProvider().parse(superdesk_article(1)) + self.assertEqual(article.url, "https://pesacheck-ui.vercel.app/post-1/") + + def test_url_shape_is_configurable(self): + # Default is the Ghost-era shape, which every fact-check already in + # Check links to; the new site's own shape is one setting away. + with mock.patch.object( + settings, + "PESACHECK_ARTICLE_URL_TEMPLATE", + "{site}/fact-checks/{desk}/{slug}", + ): + article = provider_superdesk.SuperdeskProvider().parse(superdesk_article(1)) + self.assertEqual(article.url, "https://pesacheck.org/fact-checks/somali/post-1") + + def test_an_article_without_a_route_still_gets_a_url(self): + article = superdesk_article(1) + article["swp_route"] = None + with mock.patch.object( + settings, + "PESACHECK_ARTICLE_URL_TEMPLATE", + "{site}/fact-checks/{desk}/{slug}", + ): + parsed = provider_superdesk.SuperdeskProvider().parse(article) + self.assertEqual(parsed.url, "https://pesacheck.org/post-1/") + + def test_falls_back_to_numeric_id_without_a_guid(self): + article = superdesk_article(1) + article["metadata"] = json.dumps({"subject": []}) + self.assertEqual( + provider_superdesk.SuperdeskProvider().parse(article).guid, "6001" + ) + + def test_unparseable_metadata_does_not_crash(self): + article = superdesk_article(1) + article["metadata"] = "not json" + parsed = provider_superdesk.SuperdeskProvider().parse(article) + self.assertEqual(parsed.categories, []) + self.assertEqual(parsed.language, "") + + def test_query_filters_to_fact_checks_for_the_tenant(self): + with mock.patch.object( + provider_superdesk.requests, "post", return_value=superdesk_response([]) + ) as post: + provider_superdesk.SuperdeskProvider().fetch( + since=datetime(2026, 9, 1, 10, tzinfo=UTC), limit=5 + ) + where = post.call_args.kwargs["json"]["variables"]["where"] + self.assertEqual(where["tenant_code"], {"_eq": "123abc"}) + self.assertEqual(where["published_at"]["_gte"], "2026-09-01T10:00:00") + self.assertEqual( + where["_and"][0]["swp_article_metadata"]["swp_article_metadata_subjects"], + {"scheme": {"_eq": "Debunk"}}, + ) + + def test_first_run_fetches_a_single_page(self): + with mock.patch.object( + provider_superdesk.requests, + "post", + return_value=superdesk_response([superdesk_article(i) for i in range(2)]), + ) as post: + articles = provider_superdesk.SuperdeskProvider().fetch(since=None, limit=2) + self.assertEqual(post.call_count, 1) + self.assertEqual(len(articles), 2) + + def test_pages_until_a_short_page(self): + pages = [ + superdesk_response([superdesk_article(1), superdesk_article(2)]), + superdesk_response([superdesk_article(3)]), + ] + offsets = [] + + def fake_post(url, json, headers, timeout): + offsets.append(json["variables"]["offset"]) + return pages[len(offsets) - 1] + + with mock.patch.object( + provider_superdesk.requests, "post", side_effect=fake_post + ): + articles = provider_superdesk.SuperdeskProvider().fetch( + since=datetime(2026, 9, 1, tzinfo=UTC), limit=2 + ) + self.assertEqual(offsets, [0, 2]) + self.assertEqual(len(articles), 3) + + def test_graphql_errors_are_raised(self): + resp = mock.Mock(status_code=200, text='{"errors":[{"message":"boom"}]}') + resp.json.return_value = {"errors": [{"message": "boom"}]} + with mock.patch.object(provider_superdesk.requests, "post", return_value=resp): + with self.assertRaises(Exception): + provider_superdesk.SuperdeskProvider().fetch() + + def test_preshared_auth_header_sent_when_configured(self): + with mock.patch.object( + settings, "PESACHECK_SUPERDESK_PRESHARED_AUTH", "s3cret" + ): + with mock.patch.object( + provider_superdesk.requests, "post", return_value=superdesk_response([]) + ) as post: + provider_superdesk.SuperdeskProvider().fetch() + self.assertEqual(post.call_args.kwargs["headers"]["x-preshared-auth"], "s3cret") + + def test_end_to_end_posts_and_records_the_source(self): + self.run_superdesk([superdesk_article(2, day=3), superdesk_article(1, day=2)]) + self.assertEqual([d["title"] for d in self.posted], ["Post 1", "Post 2"]) + self.assertEqual(self.posted[0]["language"], "so") + self.assertEqual( + self.posted[0]["set_tags"], ["Somali", "Kenya", "Quick Read", "Sports"] + ) + self.assertEqual(self.posted[0]["url"], "https://pesacheck.org/post-1/") + self.assertEqual(set(self.sources().values()), {"superdesk"}) + self.assertEqual( + main.get_checkpoint(self.db, "superdesk"), + datetime(2026, 9, 3, 10, tzinfo=UTC), + ) + + def test_rerun_posts_nothing(self): + self.run_superdesk([superdesk_article(1)]) + self.run_superdesk([superdesk_article(1)]) + self.assertEqual(len(self.posted), 1) + + def test_switching_providers_keeps_separate_checkpoints(self): + self.run_superdesk([superdesk_article(1, day=5)]) + settings.PESACHECK_PROVIDER = "ghost" + self.run_ghost([ghost_post(9, day=2)]) + self.assertEqual([d["title"] for d in self.posted], ["Post 1", "Post 9"]) + self.assertEqual(self.sources(), {"uuid-1": "superdesk", f"{9:024x}": "ghost"}) + self.assertEqual( + main.get_checkpoint(self.db, "superdesk"), + datetime(2026, 9, 5, 10, tzinfo=UTC), + ) + self.assertEqual( + main.get_checkpoint(self.db, "ghost"), datetime(2026, 9, 2, 10, tzinfo=UTC) + ) + + +class TestCutover(Base): + """Switching a live deployment from one CMS to the other.""" + + provider = "superdesk" + + def fetch_since(self, articles=None): + with mock.patch.object( + provider_superdesk.requests, + "post", + return_value=superdesk_response(articles or []), + ) as post: + main.main(self.db) + where = post.call_args.kwargs["json"]["variables"]["where"] + return where["published_at"].get("_gte") + + def test_first_superdesk_run_resumes_where_ghost_stopped(self): + # The same fact-checks exist in both CMSes under different guids and + # URLs, so starting from scratch would re-post them to Check. + self.add_row("a" * 24, "2026-09-20T10:00:00.000+00:00", source="ghost") + self.assertEqual(self.fetch_since(), "2026-09-20T10:00:00") + + def test_an_article_older_than_the_ghost_checkpoint_is_not_reposted(self): + self.add_row("a" * 24, "2026-09-20T10:00:00.000+00:00", source="ghost") + # Superdesk still returns it (the API filter is what excludes it), so + # assert on the window we ask for rather than on the mock's reply. + self.assertEqual(self.fetch_since(), "2026-09-20T10:00:00") + + def test_a_providers_own_checkpoint_wins_once_it_has_one(self): + self.add_row("a" * 24, "2026-09-25T10:00:00.000+00:00", source="ghost") + self.add_row("uuid-9", "2026-09-21T10:00:00.000+00:00", source="superdesk") + # Older, but it's this provider's own position: anything newer from + # Ghost is already stored, and re-fetching from the Ghost date would + # just re-examine rows we have. + self.assertEqual(self.fetch_since(), "2026-09-21T10:00:00") + + def test_a_fresh_database_still_fetches_a_single_page(self): + self.assertIsNone(self.fetch_since()) + + def test_legacy_naive_dates_dont_break_the_global_checkpoint(self): + # Medium rows stored "2024-11-18 23:19:22" with no offset; mixing them + # with offset-aware dates used to raise TypeError in max(). + self.add_row("https://medium.com/p/1", "2024-11-18 23:19:22", source="medium") + self.add_row("a" * 24, "2026-09-20T10:00:00.000+00:00", source="ghost") + self.assertEqual( + main.get_checkpoint(self.db), datetime(2026, 9, 20, 10, tzinfo=UTC) + ) + self.assertEqual( + main.get_checkpoint(self.db, "medium"), + datetime(2024, 11, 18, 23, 19, 22, tzinfo=UTC), + ) + + +class TestCatchUpIsBounded(Base): + """A checkpoint far in the past must not import years in one run.""" + + def test_ghost_asks_oldest_first_when_catching_up(self): + self.add_row("0" * 24, "2026-09-01T10:00:00.000+00:00") + with mock.patch.object( + provider_ghost.requests, "get", return_value=ghost_response([]) + ) as get: + main.main(self.db) + self.assertEqual(get.call_args.kwargs["params"]["order"], "published_at asc") + + def test_ghost_asks_newest_first_without_a_checkpoint(self): + with mock.patch.object( + provider_ghost.requests, "get", return_value=ghost_response([]) + ) as get: + main.main(self.db) + self.assertEqual(get.call_args.kwargs["params"]["order"], "published_at desc") + + def test_newest_page_is_returned_oldest_first(self): + # No checkpoint: the API answers newest first, the caller wants the + # reverse so a partial run still advances the checkpoint safely. + with mock.patch.object( + provider_ghost.requests, + "get", + return_value=ghost_response([ghost_post(2, day=3), ghost_post(1, day=2)]), + ): + posts = provider_ghost.GhostProvider().fetch(since=None, limit=2) + self.assertEqual([p["title"] for p in posts], ["Post 1", "Post 2"]) + + def test_the_cap_keeps_the_oldest_first_contract(self): + # Reachable by raising the page size to at least the cap: the cap + # fires on the first (newest-first) page, which still has to come + # back oldest first or main() would advance the checkpoint past + # articles it hasn't stored. + newest_first = [ + ghost_post(3, day=3), + ghost_post(2, day=2), + ghost_post(1, day=1), + ] + with mock.patch.object( + provider_ghost.requests, "get", return_value=ghost_response(newest_first) + ): + posts = provider_ghost.GhostProvider().fetch( + since=None, limit=3, max_articles=2 + ) + self.assertEqual([p["title"] for p in posts], ["Post 2", "Post 3"]) + + def test_the_cap_keeps_the_oldest_first_contract_for_superdesk(self): + newest_first = [ + superdesk_article(3, day=3), + superdesk_article(2, day=2), + superdesk_article(1, day=1), + ] + with mock.patch.object( + provider_superdesk.requests, + "post", + return_value=superdesk_response(newest_first), + ): + articles = provider_superdesk.SuperdeskProvider().fetch( + since=None, limit=3, max_articles=2 + ) + self.assertEqual([a["title"] for a in articles], ["Post 2", "Post 3"]) + + def test_ghost_stops_at_the_cap(self): + pages = [ + ghost_response([ghost_post(1), ghost_post(2)], next_page=2), + ghost_response([ghost_post(3), ghost_post(4)], next_page=3), + ] + calls = [] + + def fake_get(url, params, headers, timeout): + calls.append(params["page"]) + return pages[params["page"] - 1] + + with mock.patch.object(provider_ghost.requests, "get", side_effect=fake_get): + posts = provider_ghost.GhostProvider().fetch( + since=datetime(2024, 1, 1, tzinfo=UTC), limit=2, max_articles=3 + ) + self.assertEqual(len(posts), 3) + self.assertEqual(calls, [1, 2]) # stopped instead of walking the archive + + def test_superdesk_stops_at_the_cap(self): + pages = [ + superdesk_response([superdesk_article(1), superdesk_article(2)]), + superdesk_response([superdesk_article(3), superdesk_article(4)]), + ] + calls = [] + + def fake_post(url, json, headers, timeout): + calls.append(json["variables"]["offset"]) + return pages[len(calls) - 1] + + with mock.patch.object( + provider_superdesk.requests, "post", side_effect=fake_post + ): + articles = provider_superdesk.SuperdeskProvider().fetch( + since=datetime(2024, 1, 1, tzinfo=UTC), limit=2, max_articles=3 + ) + self.assertEqual(len(articles), 3) + self.assertEqual(calls, [0, 2]) + + def test_superdesk_orders_by_the_direction_it_asked_for(self): + with mock.patch.object( + provider_superdesk.requests, "post", return_value=superdesk_response([]) + ) as post: + provider_superdesk.SuperdeskProvider().fetch( + since=datetime(2024, 1, 1, tzinfo=UTC), limit=2 + ) + self.assertEqual(post.call_args.kwargs["json"]["variables"]["order"], "asc") + with mock.patch.object( + provider_superdesk.requests, "post", return_value=superdesk_response([]) + ) as post: + provider_superdesk.SuperdeskProvider().fetch(since=None, limit=2) + self.assertEqual(post.call_args.kwargs["json"]["variables"]["order"], "desc") + + def test_a_capped_run_advances_the_checkpoint_so_the_next_one_continues(self): + self.add_row("0" * 24, "2026-09-01T10:00:00.000+00:00") + with mock.patch.object(settings, "PESACHECK_MAX_ARTICLES", 1): + with mock.patch.object( + provider_ghost.requests, + "get", + return_value=ghost_response( + [ghost_post(1, day=2), ghost_post(2, day=3)], next_page=None + ), + ): + main.main(self.db) + self.assertEqual([d["title"] for d in self.posted], ["Post 1"]) + self.assertEqual( + main.get_checkpoint(self.db, "ghost"), datetime(2026, 9, 2, 10, tzinfo=UTC) + ) + + +class TestMigration(unittest.TestCase): + """A database created before the source/language columns existed.""" + + OLD_SCHEMA = """CREATE TABLE pesacheck_feeds + (title TEXT NOT NULL, pubDate TEXT NOT NULL, author TEXT NOT NULL, + guid TEXT PRIMARY KEY, link TEXT NOT NULL, thumbnail TEXT NOT NULL, + description TEXT NOT NULL, status TEXT DEFAULT 'Pending', + categories TEXT DEFAULT '[]', check_project_media_id TEXT, + check_full_url TEXT, claim_description_id TEXT)""" + + def setUp(self): + fd, self.db_file = tempfile.mkstemp(suffix=".db", dir=TMP) + os.close(fd) + conn = sqlite3.connect(self.db_file) + conn.execute(self.OLD_SCHEMA) + conn.execute( + "INSERT INTO pesacheck_feeds VALUES " + "('t','2024-11-18 23:19:22','a','https://medium.com/p/1','l','','d'," + "'Completed','[]','','','')" + ) + conn.execute( + "INSERT INTO pesacheck_feeds VALUES " + "('t','2026-09-18T07:23:07.000+00:00','a','6aace4daa4a78b00073f9a6e'," + "'l','','d','Completed','[]','','','')" + ) + conn.commit() + conn.close() + settings.PESACHECK_DATABASE_NAME = self.db_file + + def test_adds_columns_and_backfills_source(self): + db = database.PesacheckDatabase() + conn = sqlite3.connect(self.db_file) + columns = {row[1] for row in conn.execute("PRAGMA table_info(pesacheck_feeds)")} + self.assertIn("source", columns) + self.assertIn("language", columns) + sources = dict(conn.execute("SELECT guid, source FROM pesacheck_feeds")) + conn.close() + self.assertEqual(sources["https://medium.com/p/1"], "medium") + self.assertEqual(sources["6aace4daa4a78b00073f9a6e"], "ghost") + # The Ghost row's date becomes the Ghost checkpoint; Medium's does not. + self.assertEqual( + main.get_checkpoint(db, "ghost"), + datetime.fromisoformat("2026-09-18T07:23:07.000+00:00"), + ) + + def test_is_idempotent_and_rows_still_load(self): + database.PesacheckDatabase() + db = database.PesacheckDatabase() + feeds = db.get_pesacheck_feeds_by_status("Completed") + self.assertEqual(len(feeds), 2) + self.assertEqual({f.source for f in feeds}, {"medium", "ghost"}) + + def test_legacy_medium_summary_still_uses_the_figure_rule(self): + database.PesacheckDatabase() + feed = database.PesacheckFeed( + title="t", + pubDate="", + author="", + guid="https://medium.com/p/2", + link="", + thumbnail="", + description="

The summary.

", + status="Pending", + categories="[]", + source="medium", + ) + self.assertEqual(main.extract_summary(feed), "The summary.") + feed.description = "

Body only

" + self.assertIsNone(main.extract_summary(feed)) class TestDatabaseWrites(Base): @@ -227,9 +725,9 @@ def test_failed_insert_does_not_post(self): with mock.patch.object( self.db, "insert_pesacheck_feed", - side_effect=sqlite3.OperationalError("database is locked"), + side_effect=sqlite3.OperationalError("locked"), ): - self.run_with([ghost_post(1)]) + self.run_ghost([ghost_post(1)]) self.assertEqual(self.posted, []) self.sentry_exc.assert_called_once() @@ -239,7 +737,7 @@ def test_failed_claim_does_not_post(self): "claim_pending_feed", side_effect=sqlite3.OperationalError("locked"), ): - self.run_with([ghost_post(1)]) + self.run_ghost([ghost_post(1)]) self.assertEqual(self.posted, []) self.assertEqual(self.rows(), {ghost_post(1)["id"]: "Pending"}) self.sentry_exc.assert_called_once() @@ -248,33 +746,32 @@ def test_failed_update_after_post_is_not_reposted(self): with mock.patch.object( self.db, "update_pesacheck_feed", - side_effect=sqlite3.OperationalError("disk full"), + side_effect=sqlite3.OperationalError("full"), ): - self.run_with([ghost_post(1)]) + self.run_ghost([ghost_post(1)]) self.assertEqual(len(self.posted), 1) self.assertEqual(self.rows(), {ghost_post(1)["id"]: "Posting"}) - # Next run: not re-posted, and a warning is raised for reconciliation. self.sentry_msg.reset_mock() - self.run_with([ghost_post(1)]) + self.run_ghost([ghost_post(1)]) self.assertEqual(len(self.posted), 1) levels = [c.kwargs.get("level") for c in self.sentry_msg.call_args_list] self.assertIn("warning", levels) def test_check_error_reverts_to_pending_and_retries(self): self.post_mock.side_effect = Exception('{"errors": ["bad"]}') - self.run_with([ghost_post(1)]) + self.run_ghost([ghost_post(1)]) self.assertEqual(self.rows(), {ghost_post(1)["id"]: "Pending"}) self.post_mock.side_effect = lambda data: ( self.posted.append(data), check_response(1), )[1] - self.run_with([]) + self.run_ghost([]) self.assertEqual(self.rows(), {ghost_post(1)["id"]: "Completed"}) self.assertEqual(len(self.posted), 1) def test_check_timeout_leaves_posting(self): self.post_mock.side_effect = requests.Timeout("timed out") - self.run_with([ghost_post(1)]) + self.run_ghost([ghost_post(1)]) self.assertEqual(self.rows(), {ghost_post(1)["id"]: "Posting"}) def test_update_of_missing_row_raises(self): @@ -305,75 +802,12 @@ def test_connection_failure_raises_original_error(self): database.PesacheckDatabase() -class TestMalformedPosts(Base): - def test_malformed_post_skipped_others_processed(self): - bad = ghost_post(2) - del bad["title"] - with mock.patch.object( - main.requests, - "get", - return_value=ghost_response([ghost_post(3), bad, ghost_post(1)]), - ): - main.main(self.db) - self.assertEqual([d["title"] for d in self.posted], ["Post 1", "Post 3"]) - self.assertEqual(self.sentry_exc.call_count, 1) - self.assertIsInstance(self.sentry_exc.call_args.args[0], KeyError) - - -class TestSummary(Base): - def feed(self, guid, description): - return database.PesacheckFeed( - title="t", - pubDate="", - author="", - guid=guid, - link="", - thumbnail="", - description=description, - status="Pending", - categories="[]", - ) - - def test_ghost_excerpt_is_plain_text(self): - self.assertEqual( - main.extract_summary(self.feed("a" * 24, " A claim & more ")), - "A claim & more", - ) - self.assertIsNone(main.extract_summary(self.feed("a" * 24, " "))) - - def test_legacy_with_figure(self): - html = "

The summary.

Body

" - self.assertEqual( - main.extract_summary(self.feed("https://medium.com/p/1", html)), - "The summary.", - ) - - def test_legacy_without_figure_returns_none(self): - html = "

Whole long article body

More body

" - self.assertIsNone( - main.extract_summary(self.feed("https://medium.com/p/1", html)) - ) - - def test_legacy_figure_first_returns_none(self): - html = "

Body

" - self.assertIsNone( - main.extract_summary(self.feed("https://medium.com/p/1", html)) - ) - - def test_legacy_without_figure_posts_not_found(self): - self.add_row( - "https://medium.com/p/1", "", status="Pending", description="

Body

" - ) - with mock.patch.object(main.requests, "get", return_value=ghost_response([])): - main.main(self.db) - self.assertEqual(self.posted[0]["summary"], "Not Found") - - -class TestReviewRound2(Base): +class TestBatchFailures(Base): def test_failed_insert_aborts_batch_and_holds_checkpoint(self): self.add_row("0" * 24, "2026-09-01T10:00:00.000+00:00") - before = main.get_checkpoint(self.db) + before = main.get_checkpoint(self.db, "ghost") older, newer = ghost_post(1, day=2), ghost_post(2, day=3) + # Oldest first, as a catch-up fetch returns them. real_insert = self.db.insert_pesacheck_feed def flaky(feed): @@ -382,34 +816,30 @@ def flaky(feed): return real_insert(feed) with mock.patch.object(self.db, "insert_pesacheck_feed", side_effect=flaky): - with mock.patch.object( - main.requests, "get", return_value=ghost_response([newer, older]) - ): - main.main(self.db) - # Neither article was stored, so the checkpoint still covers both. + self.run_ghost([older, newer]) self.assertNotIn(older["id"], self.rows()) self.assertNotIn(newer["id"], self.rows()) - self.assertEqual(main.get_checkpoint(self.db), before) + self.assertEqual(main.get_checkpoint(self.db, "ghost"), before) self.assertEqual(self.posted, []) - # Next run (database healthy) imports both, oldest first. - with mock.patch.object( - main.requests, "get", return_value=ghost_response([newer, older]) - ): - main.main(self.db) + self.run_ghost([older, newer]) self.assertEqual([d["title"] for d in self.posted], ["Post 1", "Post 2"]) + def test_malformed_post_does_not_block_later_articles(self): + bad = ghost_post(1, day=2) + del bad["title"] + self.run_ghost([bad, ghost_post(2, day=3)]) + self.assertEqual([d["title"] for d in self.posted], ["Post 2"]) + self.assertIsInstance(self.sentry_exc.call_args.args[0], KeyError) + def test_duplicate_post_across_pages_is_handled_once(self): - # A post published mid-pagination can be returned on two pages; the - # feed_exists() check already covers that, since each article is stored - # before the next is looked at. dup = ghost_post(1, day=2) pages = [ - ghost_response([ghost_post(2, day=3), dup], next_page=2), - ghost_response([dup, ghost_post(3, day=1)], next_page=None), + ghost_response([ghost_post(3, day=1), dup], next_page=2), + ghost_response([dup, ghost_post(2, day=3)], next_page=None), ] self.add_row("0" * 24, "2026-09-01T10:00:00.000+00:00") with mock.patch.object( - main.requests, + provider_ghost.requests, "get", side_effect=lambda url, params, **kw: pages[params["page"] - 1], ): @@ -419,98 +849,88 @@ def test_duplicate_post_across_pages_is_handled_once(self): ) self.sentry_exc.assert_not_called() - def test_malformed_post_does_not_block_later_articles(self): - bad = ghost_post(1, day=2) - del bad["title"] - good = ghost_post(2, day=3) - with mock.patch.object( - main.requests, "get", return_value=ghost_response([good, bad]) - ): - main.main(self.db) - self.assertEqual([d["title"] for d in self.posted], ["Post 2"]) - self.assertIsInstance(self.sentry_exc.call_args.args[0], KeyError) - def test_bad_categories_leave_row_pending_without_posting(self): self.add_row("a" * 24, "", status="Pending") conn = sqlite3.connect(self.db_file) conn.execute("UPDATE pesacheck_feeds SET categories = 'not-json'") conn.commit() conn.close() - with mock.patch.object(main.requests, "get", return_value=ghost_response([])): - main.main(self.db) + self.run_ghost([]) self.assertEqual(self.posted, []) self.assertEqual(self.rows(), {"a" * 24: "Pending"}) - def test_check_rejection_reverts_to_pending(self): - self.post_mock.side_effect = Exception("Mutation rejected") - with mock.patch.object( - main.requests, "get", return_value=ghost_response([ghost_post(1)]) - ): - main.main(self.db) - self.assertEqual(self.rows(), {ghost_post(1)["id"]: "Pending"}) - def test_whole_run_failure_raises(self): with mock.patch.object( - main.requests, "get", side_effect=requests.ConnectionError("down") + provider_ghost.requests, "get", side_effect=requests.ConnectionError("down") ): with self.assertRaises(requests.ConnectionError): main.main(self.db) - # The end-of-run message is still sent. self.assertTrue(self.sentry_msg.called) class TestDuplicates(Base): - duplicate_payload = { - "errors": [ - { - "message": ( - "PG::UniqueViolation: ERROR: duplicate key value violates " - 'unique constraint "index_fact_checks_on_signature"' - ) - } - ] - } - def test_duplicate_is_marked_terminally_and_not_retried(self): self.post_mock.side_effect = check_api.DuplicateFactCheckError("dup") - with mock.patch.object( - main.requests, "get", return_value=ghost_response([ghost_post(1)]) - ): - main.main(self.db) + self.run_ghost([ghost_post(1)]) self.assertEqual(self.rows(), {ghost_post(1)["id"]: "Duplicate"}) # The in-memory feed agrees with the row, so later reads can't drift. self.assertEqual( self.db.get_pesacheck_feeds_by_status("Duplicate")[0].status, "Duplicate" ) - # Next run: not retried, and no exception reported for it. self.sentry_exc.reset_mock() self.post_mock.reset_mock() - with mock.patch.object(main.requests, "get", return_value=ghost_response([])): - main.main(self.db) + self.run_ghost([]) self.post_mock.assert_not_called() self.sentry_exc.assert_not_called() def test_duplicates_are_summarised_once_not_reported_as_errors(self): self.post_mock.side_effect = check_api.DuplicateFactCheckError("dup") - with mock.patch.object( - main.requests, - "get", - return_value=ghost_response([ghost_post(2, day=3), ghost_post(1, day=2)]), - ): - main.main(self.db) + self.run_ghost([ghost_post(2, day=3), ghost_post(1, day=2)]) self.sentry_exc.assert_not_called() message = self.sentry_msg.call_args.args[0] - self.assertIn("Posted 0 PesaCheck article(s)", message) self.assertIn("Skipped 2 PesaCheck article(s) Check already has", message) def test_pending_duplicate_from_an_earlier_run_is_resolved(self): self.add_row("a" * 24, "", status="Pending") self.post_mock.side_effect = check_api.DuplicateFactCheckError("dup") - with mock.patch.object(main.requests, "get", return_value=ghost_response([])): - main.main(self.db) + self.run_ghost([]) self.assertEqual(self.rows(), {"a" * 24: "Duplicate"}) +class TestSummary(Base): + def feed(self, source, description): + return database.PesacheckFeed( + title="t", + pubDate="", + author="", + guid="g", + link="", + thumbnail="", + description=description, + status="Pending", + categories="[]", + source=source, + ) + + def test_provider_summary_is_used_as_is(self): + self.assertEqual( + main.extract_summary(self.feed("ghost", " A claim & more ")), + "A claim & more", + ) + self.assertIsNone(main.extract_summary(self.feed("superdesk", " "))) + + def test_legacy_without_figure_posts_not_found(self): + self.add_row( + "https://medium.com/p/1", + "", + status="Pending", + description="

Body

", + source="medium", + ) + self.run_ghost([]) + self.assertEqual(self.posted[0]["summary"], "Not Found") + + class TestConcurrentRuns(Base): """Two runs overlapping on the same row (cron firing while one is late).""" @@ -552,7 +972,7 @@ def test_a_failure_while_marking_duplicate_keeps_the_duplicate_error(self): "update_pesacheck_feed_status", side_effect=sqlite3.OperationalError("locked"), ): - self.run_with([ghost_post(1)]) + self.run_ghost([ghost_post(1)]) # The run still knows it was a duplicate, and says so once. message = self.sentry_msg.call_args.args[0] self.assertIn("Skipped 1 PesaCheck article(s)", message) @@ -595,6 +1015,17 @@ def test_odd_error_shapes_do_not_crash(self): class TestCheckApiValidation(unittest.TestCase): + duplicate_payload = { + "errors": [ + { + "message": ( + "PG::UniqueViolation: ERROR: duplicate key value violates " + 'unique constraint "index_fact_checks_on_signature"' + ) + } + ] + } + def call_with(self, payload, status=200): resp = mock.Mock(status_code=status, text=json.dumps(payload)) resp.json.return_value = payload @@ -621,7 +1052,7 @@ def test_rejects_null_project_media(self): def test_raises_duplicate_error_for_signature_violation(self): with self.assertRaises(check_api.DuplicateFactCheckError): - self.call_with(TestDuplicates.duplicate_payload) + self.call_with(self.duplicate_payload) def test_other_graphql_errors_are_not_duplicates(self): with self.assertRaises(Exception) as ctx: