CI and images / lint (push) Successful in 2s
CI and images / extension-version (push) Successful in 3s
CI and images / frontend-build (push) Successful in 21s
CI and images / backend-lint-and-test (push) Successful in 32s
CI and images / integration (push) Failing after 2m19s
CI and images / sign-extension (push) Skipped
CI and images / build-web (push) Skipped
CI and images / smoke-web (push) Skipped
CI and images / promote (push) Skipped
CI and images / build-agent (push) Skipped
- Recover and Recapture show on every native source. The menu gated them on
a copied platform list ('patreon', 'subscribestar') that went stale when
Discord moved over. Sources now carry `native_ingester` from the backend's
own predicate.
- A running backfill is due on every scheduler tick. Nothing queued a
backfill's next chunk: each one waited for the source's regular interval,
so an armed backfill sat idle until the next check (8h at the default) and
a five-chunk walk took most of two days. The in-flight guard and the
platform lock keep one chunk at a time. A failing source falls back to its
backoff, and a stalled or out-of-budget walk stops being due. It also runs
when the artist has auto-check off, since the operator started it by hand.
- gallery-dl no longer carries Discord: its naming constants, platform
defaults, sidecar-mirroring postprocessor and token injection are gone.
The naming test moves to the native downloader and still renders against
the real gallery-dl sidecar fixture. That is the guard that the files
gallery-dl wrote are found on disk rather than fetched again.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01LVjrnpQjRgHdvq95rASoiR
293 lines
12 KiB
Python
293 lines
12 KiB
Python
"""FC-3d: pure logic for deciding which sources are due for a check.
|
|
|
|
The Celery tick task wraps `select_due_sources` and fires
|
|
`download_source.delay()` per result; this module knows nothing about
|
|
Celery, gallery-dl, or any side effect.
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
from datetime import UTC, datetime, timedelta
|
|
|
|
from sqlalchemy import func, select
|
|
from sqlalchemy.dialects.postgresql import insert as pg_insert
|
|
from sqlalchemy.ext.asyncio import AsyncSession
|
|
from sqlalchemy.orm import selectinload
|
|
|
|
from ..models import AppSetting, Artist, ImportSettings, Source
|
|
from .db_helpers import failing_sources_clause, no_access_sources_clause
|
|
|
|
MIN_INTERVAL_SECONDS = 60
|
|
MAX_INTERVAL_SECONDS = 86400
|
|
MAX_BACKOFF_EXPONENT = 6
|
|
|
|
# AppSetting key stamped every time the Beat tick fires (see scan.py). The
|
|
# tick runs every 60s; the UI flags the scheduler as stalled if the last
|
|
# stamp is older than a few minutes.
|
|
SCHEDULER_LAST_TICK_KEY = "scheduler_last_tick_at"
|
|
|
|
# AppSetting key prefix for per-platform rate-limit cooldowns. When a
|
|
# download surfaces ErrorType.RATE_LIMITED, every other source on the same
|
|
# platform is deferred for PLATFORM_RATE_LIMIT_COOLDOWN_SECONDS so the next
|
|
# scan tick doesn't fire a burst of due same-platform sources back into the
|
|
# same limit. Per-source consecutive_failures backoff still applies on top
|
|
# of this — but this is PREVENTIVE (kills the same-tick burst from N due
|
|
# sources hammering the platform at once), while consecutive_failures is
|
|
# REACTIVE (slows the offender down over many cycles). Operator-confirmed
|
|
# 2026-05-30.
|
|
PLATFORM_COOLDOWN_KEY_PREFIX = "platform_cooldown:"
|
|
PLATFORM_RATE_LIMIT_COOLDOWN_SECONDS = 900 # 15 min
|
|
|
|
|
|
def compute_effective_interval(
|
|
source: Source, artist: Artist, settings: ImportSettings,
|
|
) -> int:
|
|
"""Return the seconds between scheduled checks for one source.
|
|
|
|
Precedence: source.check_interval_override > artist.check_interval_seconds
|
|
> settings.download_schedule_default_seconds. Multiplied by
|
|
2 ** min(consecutive_failures, MAX_BACKOFF_EXPONENT) and clamped to
|
|
[MIN_INTERVAL_SECONDS, MAX_INTERVAL_SECONDS].
|
|
"""
|
|
base = (
|
|
source.check_interval_override
|
|
or artist.check_interval_seconds
|
|
or settings.download_schedule_default_seconds
|
|
)
|
|
exponent = min(max(0, source.consecutive_failures or 0), MAX_BACKOFF_EXPONENT)
|
|
factor = 2 ** exponent
|
|
raw = base * factor
|
|
return max(MIN_INTERVAL_SECONDS, min(MAX_INTERVAL_SECONDS, raw))
|
|
|
|
|
|
async def set_platform_cooldown(
|
|
session: AsyncSession, platform: str,
|
|
seconds: int = PLATFORM_RATE_LIMIT_COOLDOWN_SECONDS,
|
|
) -> None:
|
|
"""Stamp a cooldown expiry on the given platform so select_due_sources
|
|
skips every source on that platform until it expires.
|
|
|
|
Called when a download surfaces ErrorType.RATE_LIMITED so the other
|
|
sources on the same platform don't all retry into the same rate limit.
|
|
Caller is responsible for committing the session.
|
|
|
|
Uses INSERT...ON CONFLICT DO UPDATE so two concurrent workers hitting
|
|
the same platform's rate limit don't race: a SELECT-then-INSERT pattern
|
|
would let the loser's whole transaction (including the source-health
|
|
update + event finalize) roll back on a unique-violation, stranding
|
|
that event. Atomic upsert avoids that.
|
|
"""
|
|
now = datetime.now(UTC)
|
|
expires_at = (now + timedelta(seconds=seconds)).isoformat()
|
|
key = f"{PLATFORM_COOLDOWN_KEY_PREFIX}{platform}"
|
|
stmt = pg_insert(AppSetting.__table__).values(
|
|
key=key, value=expires_at, updated_at=now,
|
|
).on_conflict_do_update(
|
|
index_elements=["key"],
|
|
set_={"value": expires_at, "updated_at": now},
|
|
)
|
|
await session.execute(stmt)
|
|
|
|
|
|
async def active_platform_cooldowns(session: AsyncSession) -> dict[str, datetime]:
|
|
"""Return {platform: expires_at} for platforms whose cooldown is still
|
|
in the future. Expired rows are ignored (a future maintenance sweep can
|
|
delete them; they don't affect routing decisions on their own).
|
|
|
|
Exposed beyond scheduler_service so the manual check endpoint
|
|
(`/api/sources/<id>/check`) can defer bulk retries that would bowl
|
|
into the same rate limit the cooldown is preventing.
|
|
"""
|
|
rows = (await session.execute(
|
|
select(AppSetting.key, AppSetting.value)
|
|
.where(AppSetting.key.startswith(PLATFORM_COOLDOWN_KEY_PREFIX))
|
|
)).all()
|
|
if not rows:
|
|
return {}
|
|
now = datetime.now(UTC)
|
|
active: dict[str, datetime] = {}
|
|
for key, value in rows:
|
|
try:
|
|
expires_at = datetime.fromisoformat(value)
|
|
except (ValueError, TypeError):
|
|
continue
|
|
if expires_at > now:
|
|
active[key[len(PLATFORM_COOLDOWN_KEY_PREFIX):]] = expires_at
|
|
return active
|
|
|
|
|
|
def backfill_ready(source: Source) -> bool:
|
|
"""A deep walk the operator started, with budget left and no failure
|
|
backing it off — due NOW rather than at its next scheduled check.
|
|
|
|
A backfill runs one time-boxed chunk per download (plan #693), and nothing
|
|
queued the next chunk: each waited for the source's regular interval. At
|
|
the 8-hour default a freshly armed backfill sat untouched until the next
|
|
check (the operator armed one on 2026-09-25 and saw nothing happen) and a
|
|
five-chunk walk took most of two days. The tick's in-flight guard keeps
|
|
one chunk at a time per source and the platform lock one walk per
|
|
platform, so "due every tick" means "next chunk as soon as the last one
|
|
ends".
|
|
|
|
The failure gate is what keeps a broken source from retrying every
|
|
minute: any failed chunk raises `consecutive_failures`, which drops the
|
|
source back onto its backed-off interval. A chunk that fails to progress
|
|
twice marks the walk stalled (download_service), which ends it here too.
|
|
"""
|
|
co = source.config_overrides or {}
|
|
return (
|
|
co.get("_backfill_state") == "running"
|
|
and (source.backfill_runs_remaining or 0) > 0
|
|
and not (source.consecutive_failures or 0)
|
|
)
|
|
|
|
|
|
async def select_due_sources(session: AsyncSession) -> list[Source]:
|
|
"""Sources where (enabled, artist.auto_check) and now >= last_checked_at + effective_interval.
|
|
|
|
Never-checked sources (last_checked_at IS NULL) are always due. Sources
|
|
whose platform is currently in a rate-limit cooldown are excluded — the
|
|
cooldown is the preventive half of the burst-prevention pair (per-source
|
|
consecutive_failures backoff handles the offending source itself).
|
|
|
|
A running backfill (`backfill_ready`) is due on every tick, and whether
|
|
or not its artist is on auto-check — the operator started it by hand.
|
|
|
|
Ordering: last_checked_at ASC NULLS FIRST, then id. Never-checked
|
|
sources go first, then the longest-since-checked, so the most overdue
|
|
sources hit Celery's FIFO download queue first. Anti-starvation: if
|
|
queue throughput ever falls below the tick rate, a freshly-rerun source
|
|
can't keep cutting in line ahead of one that hasn't been checked at all.
|
|
Operator-confirmed 2026-05-30.
|
|
"""
|
|
rows = (await session.execute(
|
|
select(Source)
|
|
.options(selectinload(Source.artist))
|
|
.join(Artist, Source.artist_id == Artist.id)
|
|
.where(Source.enabled.is_(True))
|
|
.order_by(Source.last_checked_at.asc().nulls_first(), Source.id)
|
|
)).scalars().all()
|
|
|
|
cooldowns = await active_platform_cooldowns(session)
|
|
settings = await ImportSettings.load(session)
|
|
|
|
now = datetime.now(UTC)
|
|
due: list[Source] = []
|
|
for s in rows:
|
|
if s.platform in cooldowns:
|
|
continue
|
|
if backfill_ready(s):
|
|
due.append(s)
|
|
continue
|
|
if not s.artist.auto_check:
|
|
continue
|
|
interval = compute_effective_interval(s, s.artist, settings)
|
|
if s.last_checked_at is None:
|
|
due.append(s)
|
|
continue
|
|
elapsed = (now - s.last_checked_at).total_seconds()
|
|
if elapsed >= interval:
|
|
due.append(s)
|
|
return due
|
|
|
|
|
|
def compute_next_check_at(
|
|
source: Source, artist: Artist, settings: ImportSettings,
|
|
) -> datetime | None:
|
|
"""Return the projected datetime of the next check, or None if never checked."""
|
|
if backfill_ready(source):
|
|
return datetime.now(UTC)
|
|
if source.last_checked_at is None:
|
|
return None
|
|
interval = compute_effective_interval(source, artist, settings)
|
|
return source.last_checked_at + timedelta(seconds=interval)
|
|
|
|
|
|
async def record_tick(session: AsyncSession) -> None:
|
|
"""Stamp the current time on the SCHEDULER_LAST_TICK_KEY AppSetting.
|
|
|
|
Called once per Beat tick so the UI can prove the scheduler is alive.
|
|
Commits its own write so the stamp survives even if the rest of the
|
|
tick errors out.
|
|
"""
|
|
now_iso = datetime.now(UTC).isoformat()
|
|
row = (await session.execute(
|
|
select(AppSetting).where(AppSetting.key == SCHEDULER_LAST_TICK_KEY)
|
|
)).scalar_one_or_none()
|
|
if row is None:
|
|
session.add(AppSetting(key=SCHEDULER_LAST_TICK_KEY, value=now_iso))
|
|
else:
|
|
row.value = now_iso
|
|
await session.commit()
|
|
|
|
|
|
async def scheduler_status(session: AsyncSession) -> dict:
|
|
"""Summarise scheduler health for the dashboard.
|
|
|
|
Returns last_tick_at (when Beat last fired), next_due_at (earliest
|
|
upcoming scheduled check across enabled auto-check sources), due_now
|
|
(how many are due right now), and auto_sources (total under schedule).
|
|
"""
|
|
last_tick_at = (await session.execute(
|
|
select(AppSetting.value).where(AppSetting.key == SCHEDULER_LAST_TICK_KEY)
|
|
)).scalar_one_or_none()
|
|
|
|
rows = (await session.execute(
|
|
select(Source)
|
|
.options(selectinload(Source.artist))
|
|
.join(Artist, Source.artist_id == Artist.id)
|
|
.where(Source.enabled.is_(True))
|
|
.where(Artist.auto_check.is_(True))
|
|
)).scalars().all()
|
|
settings = await ImportSettings.load(session)
|
|
|
|
now = datetime.now(UTC)
|
|
due_now = 0
|
|
next_due_at: datetime | None = None
|
|
for s in rows:
|
|
if s.last_checked_at is None:
|
|
due_now += 1
|
|
continue
|
|
nca = compute_next_check_at(s, s.artist, settings)
|
|
if nca is None or nca <= now:
|
|
due_now += 1
|
|
elif next_due_at is None or nca < next_due_at:
|
|
next_due_at = nca
|
|
|
|
cooldowns = await active_platform_cooldowns(session)
|
|
|
|
# Ingestion health for the front-door ribbon (#387 B3). Counted over ENABLED
|
|
# sources rather than the auto_check subset walked above: a source that is
|
|
# erroring or paywalled is worth surfacing whether or not a schedule happens
|
|
# to poll it. Two scalar COUNTs, not a second pass over `rows`.
|
|
#
|
|
# Both predicates are the shared ones, so the ribbon and the surfaces it
|
|
# links to cannot disagree about what they are counting.
|
|
failing_sources = (await session.execute(
|
|
select(func.count()).select_from(Source)
|
|
.where(failing_sources_clause())
|
|
)).scalar_one()
|
|
no_access_sources = (await session.execute(
|
|
select(func.count()).select_from(Source)
|
|
.where(Source.enabled.is_(True), no_access_sources_clause())
|
|
)).scalar_one()
|
|
# #387 B4: lets the front door tell "nothing configured yet" (a fresh
|
|
# install — show the on-ramp) apart from "configured, still fetching" (a
|
|
# first run in progress — show what's running). Telling someone to add a
|
|
# source when they already have three and are mid-backfill is worse than
|
|
# saying nothing. Deliberately NOT auto_sources, which counts only what is
|
|
# on a schedule: a source with auto_check off still means "configured".
|
|
total_sources = (await session.execute(
|
|
select(func.count()).select_from(Source).where(Source.enabled.is_(True))
|
|
)).scalar_one()
|
|
|
|
return {
|
|
"last_tick_at": last_tick_at,
|
|
"next_due_at": next_due_at.isoformat() if next_due_at else None,
|
|
"due_now": due_now,
|
|
"auto_sources": len(rows),
|
|
"failing_sources": failing_sources,
|
|
"no_access_sources": no_access_sources,
|
|
"total_sources": total_sources,
|
|
"platform_cooldowns": {p: dt.isoformat() for p, dt in cooldowns.items()},
|
|
}
|