Files
FabledCurator/backend/app/services/membership_reconcile.py
T
bvandeusenandClaude Opus 5 7ff8915147
CI / lint (push) Successful in 2s
CI / extension-version (push) Successful in 2s
Build images / sign-extension (push) Successful in 3s
Build images / build-agent (push) Successful in 6s
CI / frontend-build (push) Successful in 27s
CI / backend-lint-and-test (push) Successful in 34s
Build images / build-web (push) Successful in 1m24s
Build images / smoke-web (push) Skipped
Build images / build-ml (push) Successful in 2m52s
Build images / promote (push) Skipped
CI / integration (push) Successful in 3m1s
feat: a source stops pulling once its membership ends, and resumes on resubscribe (3995)
Operator, 2026-09-13: "if I kill a subscription on patreon I would like the pulling to stop on curator as well", with auto-resume chosen. This reverses the 2026-09-11 "report only" decision for lapsed sources.

membership_reconcile.apply_membership_lapses runs in sync_memberships right after each platform's successful sync, so it only ever acts on the roster just written.

It stops a source (enabled=false, with the same failure-state reset as a manual disable, #1285) only when all of these hold:
- the roster is fresh
- the source's matched membership says has_paid_access is False (lapsed, or a free follow)
- the paid-through date has passed, where the platform gives one (Patreon's member.access_expires_at; SubscribeStar gives none, so it stops at once)
- the source is enabled
- the operator hasn't chosen to keep it

It never acts on absence. A source with no matched membership keeps pulling, because a rename or a never-walked source produces the same absence. An unrecognised status is never a lapse either.

It resumes only sources carrying its own `_membership_stopped` marker, once the membership is paid again.

The operator outranks the sweep both ways (SourceService.update):
- turning a stopped source back on marks it `_membership_kept`, so the next sweep leaves it alone until it's paid again
- turning a source off by hand drops the marker, so the sweep never switches it back on

Both are `_`-prefixed app-managed config keys, which operator edits already preserve. No migration.

The roster/fetch line holds. This is a source-level action by the sweep. No download path reads the roster, and the scheduler still selects on `enabled` alone. test_no_fetch_path_can_read_the_roster is unchanged.

UI: SourceRow shows a neutral "Membership ended" chip, with the status and the resume/keep explanation, ahead of the other chips. The sweep's task summary reports stopped/resumed counts.

Tests (tests/test_membership_lapses.py):
- a lapse stops the source with a clean slate and keeps the id cache
- paid-through is honoured
- absence, an unknown status and a stale roster never stop anything
- a resume touches only what the sweep stopped
- a manual on sticks, a manual off drops the marker, and a kept source is released once paid

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01SHQB1YukL3VyvMK8rcbmV9
2026-09-13 22:08:06 -04:00

360 lines
15 KiB
Python

"""Reconciling the learned roster against the sources FC actually tracks.
Milestone 387, step C4. The step the operator asked for; C0-C3 are what make it
trustworthy enough to act on.
## The buckets
1. `subscribed_not_tracked` — you pay for this and FC does not follow it. The
adoption win, and the only bucket carrying an action.
2. `tracked_not_subscribed` — FC follows this and the roster does not show you
paying for it. No longer shown on the card: the operator reversed the
2026-09-11 "report only" call on 2026-09-13. The lapsed half of it now ACTS,
in `apply_membership_lapses` below (#3995). The absent half still only
reports, because absence proves nothing.
3. `matched` — the healthy set. Counted, not listed loudly.
4. `unidentified` — sources this join cannot speak to at all. Reported as
exactly that, because the alternative is filing them under a verdict.
## Why absence is the dangerous direction
Bucket 1 is safe to be wrong about: the cost of offering a source the operator
does not want is one ignored row. Bucket 2 is not. It is computed from an
ABSENCE — no membership matched — and three different things produce that
absence: the subscription genuinely lapsed, the sweep failed, or the creator
renamed and this source has never been walked so no exact id was ever cached.
Two guards follow from that, and they are the substance of this module:
* the whole bucket is gated on `roster_is_fresh`, so a failed or never-run sweep
yields an empty list rather than a confident accusation (C3 built the state
this reads);
* every row carries the BASIS for its claim, so "your membership says former
patron" and "we know this creator's id and it is not in your roster" and "we
only have a URL handle to go on" are three different sentences rather than one
overconfident one.
`has_paid_access` returning None is honoured throughout: unknown is never
rendered as lapsed. That is the whole reason it returns a tri-state.
"""
from __future__ import annotations
from datetime import UTC, datetime
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from ..models import Artist, MembershipSync, PlatformMembership, Source
from .membership_roster import (
get_sync_state,
identity_keys_for_source,
pair_sources_with_memberships,
roster_is_fresh,
url_tail,
)
from .native_ingest_common import has_paid_access
# Why a source appears in `tracked_not_subscribed`. Ordered strongest first —
# the UI renders a different sentence per basis, because collapsing them into
# one would make the weakest claim sound like the strongest.
BASIS_LAPSED = "lapsed" # a matched membership says access ended
BASIS_ABSENT_EXACT = "absent_exact" # exact id known, not in a fresh roster
BASIS_ABSENT_HANDLE = "absent_handle" # only a URL handle to go on
def _membership_row(m: PlatformMembership) -> dict:
return {
"id": m.id,
"platform": m.platform,
"external_campaign_id": m.external_campaign_id,
"display_name": m.display_name or m.vanity_or_none(),
"url": m.url,
"vanity": m.vanity_or_none(),
"status": m.status,
"tier_names": m.tier_names,
"amount_cents": m.amount_cents,
"currency": m.currency,
"paid_access": has_paid_access(
m.platform, m.status,
is_free_member=bool((m.details or {}).get("is_free_member")),
),
}
def _source_row(source: Source, artist: Artist) -> dict:
return {
"id": source.id,
"platform": source.platform,
"url": source.url,
"enabled": source.enabled,
"artist": {"id": artist.id, "name": artist.name, "slug": artist.slug},
}
async def reconcile(
session: AsyncSession, *, platform: str, now: datetime | None = None,
) -> dict:
"""Sort one platform's memberships and sources into the four buckets.
Always returns the COMPLETE shape, including when the roster is not fresh —
a caller reading `len(result["tracked_not_subscribed"])` must not have to
check which keys exist first. `fresh` is what says whether the emptiness
means anything.
"""
state = await get_sync_state(session, platform)
fresh = roster_is_fresh(state, now=now)
memberships = (await session.execute(
select(PlatformMembership).where(PlatformMembership.platform == platform)
)).scalars().all()
rows = (await session.execute(
select(Source, Artist)
.join(Artist, Artist.id == Source.artist_id)
.where(Source.platform == platform)
)).all()
# The join itself lives in `membership_roster` beside `match_kind`, so C5's
# gated-reason annotation pairs sources with memberships by exactly the same
# rule this card sorts them by. Two copies would let the Subscriptions row
# and this card disagree about which creator a source IS.
pairs = pair_sources_with_memberships([s for s, _a in rows], memberships)
matched_membership_ids = {m.id for m, _kind in pairs.values()}
subscribed_not_tracked = []
for m in memberships:
if m.id in matched_membership_ids:
continue
paid = has_paid_access(
m.platform, m.status,
is_free_member=bool((m.details or {}).get("is_free_member")),
)
# A membership FC knows has ENDED is not an adoption opportunity —
# adding it would start a walk that can only fetch what is already
# public. Unknown (None) is still offered: the operator can judge it,
# and refusing to show it would hide a real subscription behind a word
# this code has not been taught.
if paid is False:
continue
subscribed_not_tracked.append(_membership_row(m))
tracked_not_subscribed = []
matched = []
unidentified = []
for source, artist in rows:
pair = pairs.get(source.id)
if pair is not None:
m, kind = pair
paid = has_paid_access(
m.platform, m.status,
is_free_member=bool((m.details or {}).get("is_free_member")),
)
if paid is False:
if not source.enabled:
# Already off. Reporting a source the operator has already
# stopped following is noise, not a finding.
continue
row = _source_row(source, artist)
row["basis"] = BASIS_LAPSED
row["matched_by"] = kind
row["membership"] = _membership_row(m)
tracked_not_subscribed.append(row)
else:
row = _source_row(source, artist)
row["matched_by"] = kind
row["membership"] = _membership_row(m)
matched.append(row)
continue
# No membership matched. Whether that MEANS anything depends entirely on
# how well this source can be identified at all.
has_exact = bool(identity_keys_for_source(source))
if not has_exact and url_tail(source.url) is None:
# Nothing to match on — a sidecar anchor or a URL with no handle.
# Reported as unidentified rather than silently dropped, so the
# counts add up to the source list the operator can see.
unidentified.append(_source_row(source, artist))
continue
if not source.enabled:
# Already off. Telling the operator to stop following something they
# have stopped following is noise, not a finding.
continue
row = _source_row(source, artist)
row["basis"] = BASIS_ABSENT_EXACT if has_exact else BASIS_ABSENT_HANDLE
row["matched_by"] = None
row["membership"] = None
tracked_not_subscribed.append(row)
# THE GATE. Everything above computed the bucket; this decides whether it may
# be shown. A stale or never-run roster makes every absence meaningless, and
# an absence rendered as a verdict is how this feature would tell the
# operator to cancel something they are still paying for.
if not fresh:
tracked_not_subscribed = []
return {
"platform": platform,
"fresh": fresh,
# How many sources exist on this platform at all. The UI needs it to
# decide whether an untrustworthy roster is worth mentioning: with no
# sources here there is nothing to reconcile, and a stale-roster warning
# would be noise on an install that simply has not started yet (that
# empty-install case is C6's, not this card's).
"tracked_total": len(rows),
"last_success_at": (
state.last_success_at.isoformat()
if state is not None and state.last_success_at else None
),
"subscribed_not_tracked": subscribed_not_tracked,
"tracked_not_subscribed": tracked_not_subscribed,
"matched": matched,
"unidentified": unidentified,
}
async def reconcile_all(session: AsyncSession, now: datetime | None = None) -> dict:
"""Every platform the roster knows about, in one payload for the UI.
The platform list is the UNION of platforms with memberships and platforms
with sync state, not just the former. A sweep that has never succeeded has
recorded zero memberships, and deriving the list from memberships alone
would drop exactly that platform from the payload — making a broken
credential indistinguishable from a platform FC was never asked about. That
distinction is the whole reason C3 records sync state.
"""
with_memberships = (await session.execute(
select(PlatformMembership.platform).distinct()
)).scalars().all()
with_state = (await session.execute(
select(MembershipSync.platform)
)).scalars().all()
platforms = set(with_memberships) | set(with_state)
return {
"platforms": [
await reconcile(session, platform=p, now=now) for p in sorted(platforms)
]
}
# ---------------------------------------------------------------------------
# Stop pulling what the account no longer pays for (#3995)
# ---------------------------------------------------------------------------
#
# Operator decision, 2026-09-13, reversing the 2026-09-11 "report only" call
# for this direction: "if I kill a subscription on patreon I would like the
# pulling to stop on curator as well", with automatic resume on resubscribing.
#
# This is a SOURCE-level action taken by the daily sweep, visible on the source
# row and reversible there. It is not a fetch-path decision. The line C5 draws,
# that the roster never decides a POST is inaccessible, still holds: nothing
# here reads per-post access, and no download path reads the roster
# (`test_no_fetch_path_can_read_the_roster`). The scheduler keeps selecting on
# `enabled` alone.
#
# Acts ONLY on positive evidence. A source whose matched membership says access
# has ended is stopped. A source with NO matched membership is left alone,
# because absence has innocent causes: a creator rename, a source never walked
# so no id is cached, a membership the platform stopped listing. Stopping on
# absence would switch off things the operator still pays for.
#
# Two app-managed config_overrides keys carry the state. The `_` prefix is
# already the "FC writes this, an operator edit preserves it" family.
# _membership_stopped set when the sweep stops a source; the sweep resumes
# ONLY sources carrying it, so a source the operator
# switched off by hand is never switched back on
# _membership_kept set by SourceService.update when the operator turns a
# stopped source back ON: a deliberate choice to keep
# pulling a lapsed creator, which the next sweep must
# not undo. Cleared when the membership is paid again.
STOPPED_KEY = "_membership_stopped"
KEPT_KEY = "_membership_kept"
def _access_expires_at(m: PlatformMembership) -> datetime | None:
"""When paid access actually ends, if the platform says.
Patreon keeps a cancelled membership's access until the end of the billing
period and reports that date (`member.access_expires_at`, note #3992).
SubscribeStar's page gives no such date, so a cancelled SubscribeStar
membership stops at once. Returns None when there is no usable date.
"""
details = m.details or {}
raw = details.get("access_expires_at") or (details.get("member") or {}).get("access_expires_at")
if not isinstance(raw, str) or not raw:
return None
try:
parsed = datetime.fromisoformat(raw.replace("Z", "+00:00"))
except ValueError:
return None
return parsed if parsed.tzinfo else parsed.replace(tzinfo=UTC)
async def apply_membership_lapses(
session: AsyncSession, *, platform: str, now: datetime | None = None,
) -> dict:
"""Stop sources whose paid access has ended; resume the ones this stopped.
Refuses to act on a roster that isn't fresh, for the same reason C4 refuses
to draw conclusions from one.
"""
now = now or datetime.now(UTC)
state = await get_sync_state(session, platform)
if not roster_is_fresh(state, now=now):
return {"platform": platform, "skipped": "roster not fresh", "stopped": 0, "resumed": 0}
memberships = (await session.execute(
select(PlatformMembership).where(PlatformMembership.platform == platform)
)).scalars().all()
sources = (await session.execute(
select(Source).where(Source.platform == platform)
)).scalars().all()
pairs = pair_sources_with_memberships(list(sources), list(memberships))
stopped: list[int] = []
resumed: list[int] = []
for source in sources:
pair = pairs.get(source.id)
if pair is None:
continue # absence is never acted on, see above
m, _kind = pair
paid = has_paid_access(
m.platform, m.status,
is_free_member=bool((m.details or {}).get("is_free_member")),
)
co = dict(source.config_overrides or {})
if paid is True:
changed = co.pop(KEPT_KEY, None) is not None
if STOPPED_KEY in co:
co.pop(STOPPED_KEY)
source.enabled = True
resumed.append(source.id)
changed = True
if changed:
source.config_overrides = co
continue
# Unknown status: never a reason to stop something (has_paid_access's
# tri-state exists for exactly this).
if paid is None:
continue
if not source.enabled or co.get(KEPT_KEY):
continue
expires = _access_expires_at(m)
if expires is not None and expires > now:
continue # still inside the paid-through period
co[STOPPED_KEY] = {"at": now.isoformat(), "status": m.status}
source.config_overrides = co
source.enabled = False
# The same clean slate a manual disable gives (SourceService.update,
# #1285), so a stopped source doesn't linger as failing or gated.
source.last_error = None
source.error_type = None
source.consecutive_failures = 0
stopped.append(source.id)
await session.commit()
return {"platform": platform, "stopped": len(stopped), "resumed": len(resumed)}