Compare commits

..
6 Commits
Author SHA1 Message Date
bvandeusen a8f6a464aa Merge pull request 'fix(download): salvage soft-time-limit kills + fix timeout ladder' (#59) from dev into main
CI / lint (push) Successful in 2s
Build images / sign-extension (push) Successful in 3s
CI / backend-lint-and-test (push) Successful in 21s
CI / frontend-build (push) Successful in 24s
Build images / build-web (push) Successful in 2m37s
Build images / build-ml (push) Successful in 3m21s
CI / intimp (push) Successful in 3m40s
CI / intapi (push) Successful in 7m47s
CI / intcore (push) Successful in 8m15s
2026-06-03 19:35:45 -04:00
bvandeusenandClaude Opus 4.8 9cb24c9e1b style(test): fix ruff I001 import order in download task test
CI / lint (push) Successful in 3s
CI / frontend-build (push) Successful in 20s
CI / backend-lint-and-test (push) Successful in 25s
CI / intimp (push) Successful in 3m42s
CI / intapi (push) Successful in 7m38s
CI / intcore (push) Successful in 8m16s
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-03 17:04:29 -04:00
bvandeusenandClaude Opus 4.8 6590dcdb39 fix(download): salvage soft-time-limit kills + fix timeout ladder
CI / frontend-build (push) Successful in 21s
CI / lint (push) Failing after 3s
CI / backend-lint-and-test (push) Successful in 21s
CI / intimp (push) Successful in 3m32s
CI / intapi (push) Successful in 7m22s
CI / intcore (push) Successful in 8m4s
Backfill downloads stranded with empty logs + a generic "stranded by
recovery sweep" error. Root cause: the backfill gallery-dl subprocess
timeout (1170s) exceeded download_source's Celery soft_time_limit (900s),
so SoftTimeLimitExceeded preempted subprocess.TimeoutExpired. The
TimeoutExpired path (which captures partial stdout/stderr and finalizes
the event) never ran, the event was left 'running', and phase 3 never
decremented backfill_runs_remaining — so the source re-ran and
re-stranded every tick (Anduo #39912).

Two layers:
1. Raise download_source limits (soft 900→1350, hard 1200→1500) so both
   subprocess budgets (870 tick / 1170 backfill) sit below the soft
   limit with phase-3 persist headroom. Promote to module constants and
   guard the invariant with a test.
2. Catch SoftTimeLimitExceeded in download_source and finalize the
   in-flight event with a real reason, mirror phase-3 source-health, and
   decrement backfill so a chronically-slow source self-heals to tick
   mode. The existing celery_signals handler only covered TaskRun, not
   DownloadEvent — that was the gap.

Updates stale 900/1200 references in gallery_dl.py + maintenance.py.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-03 16:56:13 -04:00
bvandeusen ab9922ad2e Merge pull request 'feat(artist): "new since last visit" badge + banner' (#58) from dev into main
Build images / sign-extension (push) Successful in 3s
CI / lint (push) Successful in 4s
CI / backend-lint-and-test (push) Successful in 23s
Build images / build-ml (push) Successful in 2m53s
CI / frontend-build (push) Successful in 21s
Build images / build-web (push) Successful in 2m9s
CI / intimp (push) Successful in 3m36s
CI / intapi (push) Successful in 7m32s
CI / intcore (push) Successful in 8m9s
2026-06-03 16:20:54 -04:00
bvandeusen 3162cff96b fix(artist): ruff UP017 + test_directory_card_shape pin
CI / lint (push) Successful in 2s
CI / frontend-build (push) Successful in 27s
CI / backend-lint-and-test (push) Successful in 27s
CI / intimp (push) Successful in 3m29s
CI / intapi (push) Successful in 7m20s
CI / intcore (push) Successful in 8m4s
Two CI bounces on b65e956:
1. ruff UP017 — Python 3.14's preferred form is `datetime.UTC`, not
   `timezone.utc`. Switch the test's two TZ literals.
2. test_directory_card_shape pinned the card key set to the pre-feature
   shape; `unseen_count` was added to the API payload but the pin
   wasn't updated. Same shape as the recurring 'plan-grep-pinned-tests'
   trap — should have grepped tests/ for card.keys() before pushing.
2026-06-03 15:45:59 -04:00
bvandeusen b65e956ad2 feat(artist): "new since last visit" badge + banner
CI / lint (push) Failing after 2s
CI / backend-lint-and-test (push) Successful in 14s
CI / frontend-build (push) Successful in 22s
CI / intimp (push) Successful in 3m39s
CI / intapi (push) Failing after 7m41s
CI / intcore (push) Successful in 8m42s
Per-artist "+N" accent pill on the artists directory and a "N new since
last visit" banner inside ArtistView. Counts new IMAGES (not posts) so
multi-image posts increment correctly.

- alembic 0034: artist_visit (artist_id PK, last_viewed_at NOT NULL).
  Seeds every existing artist with last_viewed_at=NOW() so the badge
  starts at 0 across the board — no noisy "5000 unseen images" on
  first deploy.
- ArtistService.find_or_create autoseeds a visit row alongside new
  artists, so freshly imported content doesn't read as unseen.
- ArtistService.overview reads pre-visit last_viewed_at, counts images
  created since, then atomically UPSERTs last_viewed_at=NOW() via
  postgres ON CONFLICT DO UPDATE (no SELECT-then-INSERT race per
  reference_scalar_one_or_none_duplicates). Returns the pre-update
  count as `unseen_count_at_visit` so the banner has data.
- ArtistDirectoryService.list_artists adds an `unseen_count` aggregate
  to each card via LEFT JOIN artist_visit + conditional COUNT. NULL
  last_viewed_at (artist created before this code shipped) defensively
  counts as "never visited" → all images unseen.
- Frontend: ArtistCard renders an accent pill in the preview-strip
  corner when unseen_count > 0 (capped at 99+); ArtistView shows a
  closable v-alert banner on initial load when
  unseen_count_at_visit > 0, re-arms on slug change.

Single-row-per-artist (no user_id) — rule #47 multi-user ACL is
aspirational; widens to (user_id, artist_id) PK when User lands, per
rule #22.

Scribe plan #597.
2026-06-03 15:27:11 -04:00
13 changed files with 683 additions and 37 deletions
+53
View File
@@ -0,0 +1,53 @@
"""artist_visit: per-artist last-viewed timestamp for the "+N new" badge
Revision ID: 0034
Revises: 0033
Create Date: 2026-06-03
Powers the artists-directory "+N new since last visit" badge + ArtistView
banner. Single row per artist (no user_id yet — rule #47 multi-user ACL
is aspirational; widens to (user_id, artist_id) PK when User lands).
Seed every existing artist with `last_viewed_at = NOW()` so the badge
starts at 0 across the board — no noisy "you have 5000 unseen images"
on first deploy. New artists auto-get a row via
`ArtistService.find_or_create`.
"""
from typing import Sequence, Union
import sqlalchemy as sa
from alembic import op
revision: str = "0034"
down_revision: Union[str, None] = "0033"
branch_labels: Union[str, Sequence[str], None] = None
depends_on: Union[str, Sequence[str], None] = None
def upgrade() -> None:
op.create_table(
"artist_visit",
sa.Column(
"artist_id",
sa.Integer,
sa.ForeignKey("artist.id", ondelete="CASCADE"),
primary_key=True,
),
sa.Column(
"last_viewed_at",
sa.DateTime(timezone=True),
nullable=False,
server_default=sa.text("NOW()"),
),
)
# Seed: every existing artist starts "fully caught up". Without this,
# every operator with N artists would see N badges (worth of every
# image ever imported) on first deploy.
op.execute(
"INSERT INTO artist_visit (artist_id, last_viewed_at) "
"SELECT id, NOW() FROM artist"
)
def downgrade() -> None:
op.drop_table("artist_visit")
+2
View File
@@ -2,6 +2,7 @@
from .app_setting import AppSetting from .app_setting import AppSetting
from .artist import Artist from .artist import Artist
from .artist_visit import ArtistVisit
from .backup_run import BackupRun from .backup_run import BackupRun
from .base import Base from .base import Base
from .credential import Credential from .credential import Credential
@@ -28,6 +29,7 @@ __all__ = [
"Base", "Base",
"AppSetting", "AppSetting",
"Artist", "Artist",
"ArtistVisit",
"BackupRun", "BackupRun",
"Source", "Source",
"Credential", "Credential",
+36
View File
@@ -0,0 +1,36 @@
"""ArtistVisit — per-artist 'last viewed' timestamp.
Powers the "+N new since last visit" badge on the artists directory and
the matching banner on `ArtistView`. One row per artist, single global
operator. When the multi-user model lands, the PK widens to
`(user_id, artist_id)` — currently aspirational only (no User model,
no services/access.py); operator approved skipping `user_id` for now
under rule #22 (breaking changes welcome).
Seed at migration time: every existing artist gets `last_viewed_at = NOW()`
so the badge starts at 0 across the board (no noisy "5000 unseen" on
first deploy). New artists also auto-get a row via
`ArtistService.find_or_create`.
"""
from datetime import datetime
from sqlalchemy import DateTime, ForeignKey, Integer, func
from sqlalchemy.orm import Mapped, mapped_column
from .base import Base
class ArtistVisit(Base):
__tablename__ = "artist_visit"
artist_id: Mapped[int] = mapped_column(
Integer,
ForeignKey("artist.id", ondelete="CASCADE"),
primary_key=True,
)
last_viewed_at: Mapped[datetime] = mapped_column(
DateTime(timezone=True),
nullable=False,
server_default=func.now(),
)
@@ -13,10 +13,10 @@ from __future__ import annotations
import base64 import base64
from dataclasses import dataclass from dataclasses import dataclass
from sqlalchemy import and_, exists, func, or_, select from sqlalchemy import and_, case, exists, func, or_, select
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from ..models import Artist, ImageRecord, Source from ..models import Artist, ArtistVisit, ImageRecord, Source
from .gallery_service import thumbnail_url from .gallery_service import thumbnail_url
_SEP = "|" _SEP = "|"
@@ -58,9 +58,27 @@ class ArtistDirectoryService:
raise ValueError("limit must be between 1 and 200") raise ValueError("limit must be between 1 and 200")
count_col = func.count(ImageRecord.id).label("image_count") count_col = func.count(ImageRecord.id).label("image_count")
# Unseen = images imported since the artist's last_viewed_at.
# NULL last_viewed_at (artist created before alembic 0034 seed
# or before find_or_create autoseed) defensively counts as
# "never visited" → all images unseen. Single grouped query, no
# N+1.
unseen_col = func.count(
case(
(
or_(
ArtistVisit.last_viewed_at.is_(None),
ImageRecord.created_at > ArtistVisit.last_viewed_at,
),
ImageRecord.id,
),
else_=None,
)
).label("unseen_count")
stmt = ( stmt = (
select(Artist, count_col) select(Artist, count_col, unseen_col)
.outerjoin(ImageRecord, ImageRecord.artist_id == Artist.id) .outerjoin(ImageRecord, ImageRecord.artist_id == Artist.id)
.outerjoin(ArtistVisit, ArtistVisit.artist_id == Artist.id)
.group_by(Artist.id) .group_by(Artist.id)
) )
if q: if q:
@@ -94,7 +112,7 @@ class ArtistDirectoryService:
next_cursor = _encode(last_artist.name, last_artist.id) next_cursor = _encode(last_artist.name, last_artist.id)
rows = rows[:limit] rows = rows[:limit]
artist_ids = [a.id for a, _ in rows] artist_ids = [a.id for a, _, _ in rows]
previews = await self._previews(artist_ids) previews = await self._previews(artist_ids)
cards = [ cards = [
@@ -104,9 +122,10 @@ class ArtistDirectoryService:
"slug": artist.slug, "slug": artist.slug,
"is_subscription": bool(artist.is_subscription), "is_subscription": bool(artist.is_subscription),
"image_count": int(image_count), "image_count": int(image_count),
"unseen_count": int(unseen_count),
"preview_thumbnails": previews.get(artist.id, []), "preview_thumbnails": previews.get(artist.id, []),
} }
for artist, image_count in rows for artist, image_count, unseen_count in rows
] ]
return DirectoryPage(cards=cards, next_cursor=next_cursor) return DirectoryPage(cards=cards, next_cursor=next_cursor)
+49
View File
@@ -9,11 +9,13 @@ Dates come from Post.post_date via ImageProvenance.post_id.
from dataclasses import dataclass from dataclasses import dataclass
from sqlalchemy import and_, case, func, or_, select from sqlalchemy import and_, case, func, or_, select
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.exc import IntegrityError from sqlalchemy.exc import IntegrityError
from sqlalchemy.ext.asyncio import AsyncSession from sqlalchemy.ext.asyncio import AsyncSession
from ..models import ( from ..models import (
Artist, Artist,
ArtistVisit,
ImageProvenance, ImageProvenance,
ImageRecord, ImageRecord,
Post, Post,
@@ -122,6 +124,12 @@ class ArtistService:
) )
).scalar_one() ).scalar_one()
# Mark this artist as "visited now"; the returned count is what
# the operator should see in the banner ("N new since last
# visit"). Done LAST so the read aggregates above all see the
# pre-visit state (cosmetic — none depend on visit data).
unseen_at_visit = await self._mark_visited_returning_unseen(aid)
return { return {
"id": artist.id, "id": artist.id,
"name": artist.name, "name": artist.name,
@@ -129,6 +137,7 @@ class ArtistService:
"is_subscription": bool(artist.is_subscription), "is_subscription": bool(artist.is_subscription),
"image_count": int(image_count), "image_count": int(image_count),
"post_count": int(post_count), "post_count": int(post_count),
"unseen_count_at_visit": unseen_at_visit,
"date_range": { "date_range": {
"min": dmin.isoformat() if dmin else None, "min": dmin.isoformat() if dmin else None,
"max": dmax.isoformat() if dmax else None, "max": dmax.isoformat() if dmax else None,
@@ -157,6 +166,39 @@ class ArtistService:
], ],
} }
async def _mark_visited_returning_unseen(self, artist_id: int) -> int:
"""Read pre-visit `last_viewed_at`, count images added since,
then upsert `last_viewed_at = NOW()`. Returns the count BEFORE
the upsert so the banner has data to render.
Postgres UPSERT (`ON CONFLICT DO UPDATE`) keeps the write
atomic — no SELECT-then-INSERT race per
`reference_scalar_one_or_none_duplicates`.
"""
prev = (
await self.session.execute(
select(ArtistVisit.last_viewed_at).where(
ArtistVisit.artist_id == artist_id
)
)
).scalar_one_or_none()
count_stmt = select(func.count(ImageRecord.id)).where(
ImageRecord.artist_id == artist_id
)
if prev is not None:
count_stmt = count_stmt.where(ImageRecord.created_at > prev)
unseen = (await self.session.execute(count_stmt)).scalar_one()
upsert = pg_insert(ArtistVisit.__table__).values(artist_id=artist_id)
upsert = upsert.on_conflict_do_update(
index_elements=["artist_id"],
set_={"last_viewed_at": func.now()},
)
await self.session.execute(upsert)
await self.session.commit()
return int(unseen)
async def images( async def images(
self, slug: str, cursor: str | None, limit: int = 60 self, slug: str, cursor: str | None, limit: int = 60
) -> ArtistImagesPage | None: ) -> ArtistImagesPage | None:
@@ -230,6 +272,13 @@ class ArtistService:
artist = Artist(name=cleaned, slug=slug) artist = Artist(name=cleaned, slug=slug)
self.session.add(artist) self.session.add(artist)
await self.session.flush() await self.session.flush()
# New artist starts "caught up" — seed ArtistVisit so the
# directory's `+N new` badge stays at 0 until real new
# content arrives. Without this, the unseen-count query
# treats NULL last_viewed_at as "never visited" and would
# count every image imported in the same session.
self.session.add(ArtistVisit(artist_id=artist.id))
await self.session.flush()
await sp.commit() await sp.commit()
except IntegrityError: except IntegrityError:
await sp.rollback() await sp.rollback()
+23 -19
View File
@@ -59,29 +59,33 @@ TICK_SKIP_VALUE = "exit:20"
# Source.backfill_runs_remaining > 0 selects this mode; the longer # Source.backfill_runs_remaining > 0 selects this mode; the longer
# timeout below absorbs creators with thousands of posts. # timeout below absorbs creators with thousands of posts.
# #
# 30 seconds shy of Celery's hard `time_limit=1200` on download_source # Sits below download_source's Celery soft_time_limit
# (tasks/download.py:33). subprocess.run MUST raise TimeoutExpired # (DOWNLOAD_SOFT_TIME_LIMIT=1350, tasks/download.py) with ~180s of
# before Celery SIGKILLs the worker — same rationale as the tick # headroom for phase-3 persist. subprocess.run MUST raise TimeoutExpired
# default at line 74. The audit (2026-06-02) caught this at 1800, # before Celery raises SoftTimeLimitExceeded — that exception path
# guaranteeing SIGKILL on any backfill that ran to its subprocess # captures partial stdout/stderr and finalizes the event; the soft-limit
# budget: stdout/stderr lost, backfill_runs_remaining never # path (until the 2026-06-03 fix) did not. Audit history: 1800 guaranteed
# decrements, recovery sweep stamps generic "stranded" 30 min later. # SIGKILL against the old hard limit (Knuxy #38275); 1170 was then sized
# Recreates the exact Knuxy #38275 failure mode the tick 870s default # "30s shy of the hard limit (1200)" but still EXCEEDED the soft limit
# was added to prevent. backfill_runs_remaining=3 still gives ~58 # (900), so SoftTimeLimitExceeded preempted TimeoutExpired and every
# minutes of cumulative walk across three runs for prolific creators. # backfill stranded empty (Anduo #39912). Raising the Celery soft/hard
# limits to 1350/1500 (tasks/download.py) is what made 1170 safe.
# backfill_runs_remaining=3 still gives ~58 minutes of cumulative walk
# across three runs for prolific creators.
BACKFILL_SKIP_VALUE = True BACKFILL_SKIP_VALUE = True
BACKFILL_TIMEOUT_SECONDS = 1170 BACKFILL_TIMEOUT_SECONDS = 1170
# 30 seconds shy of download_source's Celery soft_time_limit (900s, see # Sits well below download_source's Celery soft_time_limit
# tasks/download.py:32). subprocess.run MUST raise TimeoutExpired before # (DOWNLOAD_SOFT_TIME_LIMIT=1350, tasks/download.py). subprocess.run MUST
# Celery raises SoftTimeLimitExceeded — otherwise Celery wins the race, # raise TimeoutExpired before Celery raises SoftTimeLimitExceeded —
# SIGKILLs the worker, in-memory stdout/stderr is lost, and the # otherwise Celery wins the race, SIGKILLs the worker, in-memory
# DownloadEvent ends up empty-logged with "stranded by recovery sweep" # stdout/stderr is lost, and the DownloadEvent ends up empty-logged with
# 18 minutes later (operator-flagged 2026-05-31, Knuxy event #38275). # "stranded by recovery sweep" (operator-flagged 2026-05-31, Knuxy event
# The 30s buffer absorbs scheduler jitter / GC pauses without making # #38275; recurred in backfill mode as Anduo #39912). Per-source bumps
# legitimately-long-running syncs timeout-friendlier. Per-source bumps # still live in source.config_overrides for legitimately long syncs —
# still live in source.config_overrides for legitimately long syncs. # keep any override below the soft limit, or the soft-limit salvage path
# in tasks/download.py (_finalize_soft_limited) is the only safety net.
_DEFAULT_GDL_TIMEOUT_SECONDS = 870 _DEFAULT_GDL_TIMEOUT_SECONDS = 870
+94 -4
View File
@@ -1,12 +1,17 @@
"""download_source Celery task — runs DownloadService for one source.""" """download_source Celery task — runs DownloadService for one source."""
import asyncio import asyncio
import logging
from datetime import UTC, datetime
from pathlib import Path from pathlib import Path
from celery.exceptions import SoftTimeLimitExceeded
from sqlalchemy import select
from sqlalchemy.exc import DBAPIError, OperationalError from sqlalchemy.exc import DBAPIError, OperationalError
from sqlalchemy.orm import Session as SyncSession
from ..celery_app import celery from ..celery_app import celery
from ..models import ImportSettings from ..models import DownloadEvent, ImportSettings, Source
from ..services.credential_crypto import CredentialCrypto from ..services.credential_crypto import CredentialCrypto
from ..services.credential_service import CredentialService from ..services.credential_service import CredentialService
from ..services.download_service import DownloadService from ..services.download_service import DownloadService
@@ -16,9 +21,79 @@ from ..services.thumbnailer import Thumbnailer
from ._async_session import async_session_factory from ._async_session import async_session_factory
from .import_file import _sync_session_factory from .import_file import _sync_session_factory
log = logging.getLogger(__name__)
IMAGES_ROOT = Path("/images") IMAGES_ROOT = Path("/images")
_KEY_PATH = IMAGES_ROOT / "secrets" / "credential_key.b64" _KEY_PATH = IMAGES_ROOT / "secrets" / "credential_key.b64"
# Celery time budget for one download_source run. The ceiling that
# governs *clean* teardown is the SOFT limit: it raises a catchable
# SoftTimeLimitExceeded in-process, whereas the HARD limit SIGKILLs the
# worker (no chance to finalize). Both gallery-dl subprocess budgets
# (gallery_dl.py: _DEFAULT_GDL_TIMEOUT_SECONDS=870 tick,
# BACKFILL_TIMEOUT_SECONDS=1170 backfill) MUST sit below the soft limit
# so subprocess.run raises its own TimeoutExpired first — that path
# captures partial stdout/stderr and finalizes the DownloadEvent. soft is
# max-subprocess (1170) + ~180s phase-3 persist headroom; hard is soft +
# 150s SIGKILL backstop. Audit 2026-06-03 (Anduo #39912): the old
# soft=900 sat BELOW the 1170 backfill budget, so SoftTimeLimitExceeded
# preempted TimeoutExpired and the event stranded empty. The recovery
# sweep's DOWNLOAD_STALL_THRESHOLD_MINUTES (30 min) still trails the new
# 25-min hard kill by 5 min, so it stays a true backstop. Invariant
# guarded by test_timeout_ladder_keeps_subprocess_budgets_under_soft_limit.
DOWNLOAD_SOFT_TIME_LIMIT = 1350
DOWNLOAD_HARD_TIME_LIMIT = 1500
def _finalize_soft_limited(session: SyncSession, source_id: int) -> None:
"""Defense in depth for the soft-time-limit kill path.
A SoftTimeLimitExceeded unwinds download_source before phase 3 can
finalize the DownloadEvent, leaving it 'running' until the recovery
sweep stamps a context-free "stranded" error 30 min later — AND
leaving backfill_runs_remaining undecremented so the source re-runs
and re-strands every tick (Anduo #39912, 2026-06-03). Flip the
in-flight event to error with a real reason, mirror phase 3's
source-health write, and decrement any backfill budget so a
chronically-slow source self-heals back to tick mode.
The caller owns the commit. All mutations are gated on actually
finding a running event, so a benign late soft-limit (phase 3 already
committed) is a no-op.
"""
now = datetime.now(UTC)
ev = session.execute(
select(DownloadEvent)
.where(DownloadEvent.source_id == source_id)
.where(DownloadEvent.status == "running")
.order_by(DownloadEvent.id.desc())
.limit(1)
).scalar_one_or_none()
if ev is None:
return
ev.status = "error"
ev.finished_at = now
ev.error = (
f"killed by Celery soft time limit ({DOWNLOAD_SOFT_TIME_LIMIT}s) "
"before the gallery-dl subprocess returned — the run exceeded its "
"budget and its stdout/stderr were lost with the worker thread. "
"If this recurs, the source is too large for one run; the backfill "
"budget was decremented so the next tick walks less."
)
ev.metadata_ = {
**(ev.metadata_ or {}),
"error_type": "timeout",
"soft_time_limited": True,
}
src = session.get(Source, source_id)
if src is not None:
src.consecutive_failures = (src.consecutive_failures or 0) + 1
src.last_error = "soft time limit exceeded"
src.error_type = "timeout"
src.last_checked_at = now
if (src.backfill_runs_remaining or 0) > 0:
src.backfill_runs_remaining = max(0, src.backfill_runs_remaining - 1)
@celery.task( @celery.task(
name="backend.app.tasks.download.download_source", name="backend.app.tasks.download.download_source",
@@ -29,8 +104,8 @@ _KEY_PATH = IMAGES_ROOT / "secrets" / "credential_key.b64"
retry_backoff_max=120, retry_backoff_max=120,
retry_jitter=True, retry_jitter=True,
max_retries=3, max_retries=3,
soft_time_limit=900, soft_time_limit=DOWNLOAD_SOFT_TIME_LIMIT,
time_limit=1200, time_limit=DOWNLOAD_HARD_TIME_LIMIT,
) )
def download_source(self, source_id: int) -> int: def download_source(self, source_id: int) -> int:
"""Returns the DownloadEvent.id.""" """Returns the DownloadEvent.id."""
@@ -73,4 +148,19 @@ def download_source(self, source_id: int) -> int:
finally: finally:
await async_engine.dispose() await async_engine.dispose()
return asyncio.run(_run()) try:
return asyncio.run(_run())
except SoftTimeLimitExceeded:
# phase 3 never ran — salvage the in-flight event so the operator
# sees a real reason instead of the recovery sweep's generic
# "stranded" 30 min later (Anduo #39912). Best-effort: a failure
# here must not mask the timeout. Re-raise so Celery + the
# task_run signal handler still record the kill.
try:
SyncFactory = _sync_session_factory()
with SyncFactory() as session:
_finalize_soft_limited(session, source_id)
session.commit()
except Exception: # noqa: BLE001 — cleanup must not swallow the kill
log.exception("soft-limit finalize failed for source %s", source_id)
raise
+7 -4
View File
@@ -46,9 +46,12 @@ MAX_RECOVERY_ATTEMPTS = 3
ORPHAN_PENDING_THRESHOLD_MINUTES = 30 ORPHAN_PENDING_THRESHOLD_MINUTES = 30
# DownloadEvent (pending|running) recovery threshold. download_source has # DownloadEvent (pending|running) recovery threshold. download_source has
# time_limit=1200s (20 min); 30 min is 10 min past that, so a legitimately- # time_limit=1500s (25 min, DOWNLOAD_HARD_TIME_LIMIT); 30 min is 5 min past
# running task is never killed by the sweep. Operator-confirmed 2026-05-29 # that, so a legitimately-running task is hard-killed before the sweep ever
# after 43 sources stranded at "last check never" by the in-flight guard. # touches it — the sweep only catches events whose worker died without
# finalizing. Operator-confirmed 2026-05-29 after 43 sources stranded at
# "last check never" by the in-flight guard; budget bumped 2026-06-03 with
# the soft/hard limit raise (Anduo #39912).
DOWNLOAD_STALL_THRESHOLD_MINUTES = 30 DOWNLOAD_STALL_THRESHOLD_MINUTES = 30
OLD_TASK_DAYS = 7 OLD_TASK_DAYS = 7
@@ -535,7 +538,7 @@ def recover_stalled_download_events() -> int:
tasks.scan._tick_due_sources_async) inserts DownloadEvent(status='pending') tasks.scan._tick_due_sources_async) inserts DownloadEvent(status='pending')
and fires download_source.delay(). If that task dies before finalizing the and fires download_source.delay(). If that task dies before finalizing the
event — worker OOM/SIGKILL, lost task, or a gallery-dl that didn't unwind event — worker OOM/SIGKILL, lost task, or a gallery-dl that didn't unwind
on the 1200s hard time_limit — the event stays in-flight forever. The next on the 1500s hard time_limit — the event stays in-flight forever. The next
tick then skips that source because of the in-flight guard (scan.py:168) tick then skips that source because of the in-flight guard (scan.py:168)
and Source.last_checked_at never updates; the operator sees "last check and Source.last_checked_at never updates; the operator sees "last check
never" in the Subscriptions health column, permanently. never" in the Subscriptions health column, permanently.
@@ -8,6 +8,15 @@
<div v-if="card.preview_thumbnails.length === 0" class="fc-artistcard__noimg"> <div v-if="card.preview_thumbnails.length === 0" class="fc-artistcard__noimg">
No preview No preview
</div> </div>
<!-- Accent pill in the corner when this artist has content imported
since the operator last opened their detail view. Caps at 99+
to keep the layout compact; the actual count appears in the
banner inside ArtistView. -->
<span
v-if="(card.unseen_count || 0) > 0"
class="fc-artistcard__unseen"
:aria-label="`${card.unseen_count} new since last visit`"
>+{{ card.unseen_count > 99 ? '99+' : card.unseen_count }}</span>
</div> </div>
<v-card-text class="fc-artistcard__body"> <v-card-text class="fc-artistcard__body">
<div class="fc-artistcard__name">{{ card.name }}</div> <div class="fc-artistcard__name">{{ card.name }}</div>
@@ -37,6 +46,7 @@ function onCardClick() {
<style scoped> <style scoped>
.fc-artistcard { cursor: pointer; } .fc-artistcard { cursor: pointer; }
.fc-artistcard__previews { .fc-artistcard__previews {
position: relative;
display: grid; grid-template-columns: repeat(3, 1fr); display: grid; grid-template-columns: repeat(3, 1fr);
gap: 2px; aspect-ratio: 3 / 1; gap: 2px; aspect-ratio: 3 / 1;
/* Explicit floor + ceiling so tall source images can't escape the /* Explicit floor + ceiling so tall source images can't escape the
@@ -45,6 +55,19 @@ function onCardClick() {
overflow: hidden; overflow: hidden;
background: rgb(var(--v-theme-surface-light)); background: rgb(var(--v-theme-surface-light));
} }
.fc-artistcard__unseen {
position: absolute;
top: 6px; right: 6px;
display: inline-flex; align-items: center;
padding: 2px 8px;
font-size: 11px; font-weight: 700; letter-spacing: 0.02em;
font-variant-numeric: tabular-nums;
color: rgb(var(--v-theme-on-accent, 0, 0, 0));
background: rgb(var(--v-theme-accent));
border-radius: 999px;
box-shadow: 0 1px 3px rgba(0, 0, 0, 0.3);
pointer-events: none;
}
.fc-artistcard__previews img { .fc-artistcard__previews img {
display: block; display: block;
width: 100%; height: 100%; width: 100%; height: 100%;
+25 -1
View File
@@ -20,6 +20,22 @@
:last-added="store.lastAdded" :last-added="store.lastAdded"
/> />
<v-container fluid class="pt-2 pb-4"> <v-container fluid class="pt-2 pb-4">
<!-- "N new since last visit" banner. Visible only on the initial
load that triggered the visit-mark; dismissable via close
button or by switching tabs. Re-entry only re-shows if more
content has arrived (overview returns 0 immediately after a
previous visit). -->
<v-alert
v-if="unseenBanner"
type="info" variant="tonal" density="compact"
class="mb-3" closable
@click:close="unseenBanner = false"
>
<span class="fc-artist__unseen-msg">
<strong>{{ store.overview.unseen_count_at_visit }}</strong>
new since last visit
</span>
</v-alert>
<v-window v-model="tab"> <v-window v-model="tab">
<v-window-item value="posts"> <v-window-item value="posts">
<ArtistPostsTab <ArtistPostsTab
@@ -39,7 +55,7 @@
</template> </template>
<script setup> <script setup>
import { computed, watch } from 'vue' import { computed, ref, watch } from 'vue'
import { useRoute } from 'vue-router' import { useRoute } from 'vue-router'
import { useArtistStore } from '../stores/artist.js' import { useArtistStore } from '../stores/artist.js'
@@ -60,6 +76,10 @@ const { tab, resolve } = useTabQuery(
() => ((store.postCount ?? 0) > 0 ? 'posts' : 'gallery'), () => ((store.postCount ?? 0) > 0 ? 'posts' : 'gallery'),
) )
// One-shot banner — reset on each new artist-slug load so it re-appears
// when navigating between artists that each have unseen content.
const unseenBanner = ref(false)
watch(slug, async (s) => { watch(slug, async (s) => {
if (!s) return if (!s) return
await store.load(s) await store.load(s)
@@ -67,6 +87,7 @@ watch(slug, async (s) => {
? `${store.overview.name} — FabledCurator` ? `${store.overview.name} — FabledCurator`
: 'FabledCurator' : 'FabledCurator'
tab.value = resolve() tab.value = resolve()
unseenBanner.value = (store.overview?.unseen_count_at_visit || 0) > 0
}, { immediate: true }) }, { immediate: true })
</script> </script>
@@ -74,4 +95,7 @@ watch(slug, async (s) => {
.fc-artist__loading { .fc-artist__loading {
display: flex; justify-content: center; padding: 64px 0; display: flex; justify-content: center; padding: 64px 0;
} }
.fc-artist__unseen-msg {
font-size: 14px;
}
</style> </style>
+2 -1
View File
@@ -43,7 +43,8 @@ async def test_directory_card_shape(client, seeded):
body = await resp.get_json() body = await resp.get_json()
card = next(c for c in body["cards"] if c["name"] == "alice-api") card = next(c for c in body["cards"] if c["name"] == "alice-api")
assert set(card.keys()) == { assert set(card.keys()) == {
"id", "name", "slug", "is_subscription", "image_count", "preview_thumbnails", "id", "name", "slug", "is_subscription", "image_count",
"unseen_count", "preview_thumbnails",
} }
assert card["is_subscription"] is True assert card["is_subscription"] is True
assert card["image_count"] == 1 assert card["image_count"] == 1
+175
View File
@@ -0,0 +1,175 @@
"""ArtistVisit unseen-count badge + ArtistService.overview banner data.
Covers:
- Directory cards include `unseen_count`
- LEFT JOIN keeps artists without a visit row (treats NULL as "never
visited" → all images unseen)
- overview() returns `unseen_count_at_visit` and stamps the visit
- find_or_create autoseeds a visit row so freshly imported content
doesn't show up as unseen
- Repeat overview() returns 0 (since the previous visit just stamped
last_viewed_at = NOW())
"""
from datetime import UTC, datetime, timedelta
import pytest
from sqlalchemy import select
from backend.app.models import Artist, ArtistVisit, ImageRecord
from backend.app.services.artist_directory_service import ArtistDirectoryService
from backend.app.services.artist_service import ArtistService
pytestmark = pytest.mark.integration
_LONG_AGO = datetime(2000, 1, 1, tzinfo=UTC)
_RECENTLY = datetime(2099, 1, 1, tzinfo=UTC)
async def _seed_artist(db, name: str) -> Artist:
a = Artist(name=name, slug=name.lower().replace(" ", "-"))
db.add(a)
await db.flush()
return a
async def _seed_visit(db, artist_id: int, when: datetime) -> None:
db.add(ArtistVisit(artist_id=artist_id, last_viewed_at=when))
await db.flush()
async def _seed_image(db, artist_id: int, suffix: str, *, created_at: datetime) -> None:
db.add(ImageRecord(
path=f"/images/visit-{suffix}.jpg",
sha256=f"visit{suffix}".ljust(64, "0")[:64],
size_bytes=10, mime="image/jpeg", width=10, height=10,
origin="downloaded", artist_id=artist_id,
created_at=created_at,
))
await db.flush()
# --- Directory unseen_count -----------------------------------------------
@pytest.mark.asyncio
async def test_directory_unseen_count_zero_when_no_images(db):
await _seed_artist(db, "zoey-empty-visit")
await db.commit()
page = await ArtistDirectoryService(db).list_artists(
q="zoey-empty-visit", platform=None, cursor=None, limit=60,
)
target = next(c for c in page.cards if c["name"] == "zoey-empty-visit")
assert target["unseen_count"] == 0
@pytest.mark.asyncio
async def test_directory_unseen_counts_only_images_after_visit(db):
a = await _seed_artist(db, "alice-visit-mix")
await _seed_visit(db, a.id, _LONG_AGO + timedelta(days=365))
# Two images BEFORE the visit (seen), three AFTER (unseen).
for i in range(2):
await _seed_image(db, a.id, f"old-{i}", created_at=_LONG_AGO)
for i in range(3):
await _seed_image(db, a.id, f"new-{i}", created_at=_RECENTLY)
await db.commit()
page = await ArtistDirectoryService(db).list_artists(
q="alice-visit-mix", platform=None, cursor=None, limit=60,
)
target = next(c for c in page.cards if c["name"] == "alice-visit-mix")
assert target["image_count"] == 5
assert target["unseen_count"] == 3
@pytest.mark.asyncio
async def test_directory_treats_missing_visit_as_never_visited(db):
"""No ArtistVisit row → defensive count of all images as unseen.
Shouldn't happen in practice (migration 0034 seeds existing
artists, find_or_create autoseeds new ones), but the directory
query must not regress to "0 unseen" if a row is missing.
"""
a = await _seed_artist(db, "bob-no-visit")
for i in range(4):
await _seed_image(db, a.id, f"orphan-{i}", created_at=_RECENTLY)
await db.commit()
page = await ArtistDirectoryService(db).list_artists(
q="bob-no-visit", platform=None, cursor=None, limit=60,
)
target = next(c for c in page.cards if c["name"] == "bob-no-visit")
assert target["unseen_count"] == 4
# --- overview() marks visit + returns count -------------------------------
@pytest.mark.asyncio
async def test_overview_returns_unseen_count_at_visit_and_stamps_now(db):
a = await _seed_artist(db, "carol-stamp")
await _seed_visit(db, a.id, _LONG_AGO)
await _seed_image(db, a.id, "stamp-1", created_at=_RECENTLY)
await _seed_image(db, a.id, "stamp-2", created_at=_RECENTLY)
await db.commit()
data = await ArtistService(db).overview("carol-stamp")
assert data is not None
assert data["unseen_count_at_visit"] == 2
# last_viewed_at advanced to NOW() — directly check the row.
visit_at = (await db.execute(
select(ArtistVisit.last_viewed_at).where(ArtistVisit.artist_id == a.id)
)).scalar_one()
assert visit_at > _LONG_AGO
@pytest.mark.asyncio
async def test_overview_repeat_call_returns_zero(db):
a = await _seed_artist(db, "dana-repeat")
await _seed_visit(db, a.id, _LONG_AGO)
await _seed_image(db, a.id, "repeat-1", created_at=_LONG_AGO + timedelta(days=1))
await db.commit()
first = await ArtistService(db).overview("dana-repeat")
assert first is not None
assert first["unseen_count_at_visit"] == 1
# Second call: no new images, visit just stamped → count is 0.
second = await ArtistService(db).overview("dana-repeat")
assert second is not None
assert second["unseen_count_at_visit"] == 0
# --- find_or_create autoseeds the visit row ------------------------------
@pytest.mark.asyncio
async def test_find_or_create_autoseeds_visit_row(db):
artist, created = await ArtistService(db).find_or_create("Eve-Autoseed")
assert created is True
row = (await db.execute(
select(ArtistVisit.last_viewed_at).where(ArtistVisit.artist_id == artist.id)
)).scalar_one_or_none()
assert row is not None
@pytest.mark.asyncio
async def test_find_or_create_existing_does_not_reset_visit(db):
"""Calling find_or_create on an existing artist returns it as-is —
must NOT touch the visit row's timestamp."""
a = await _seed_artist(db, "Frank-Existing")
await _seed_visit(db, a.id, _LONG_AGO)
await db.commit()
artist, created = await ArtistService(db).find_or_create("Frank-Existing")
assert created is False
assert artist.id == a.id
visit_at = (await db.execute(
select(ArtistVisit.last_viewed_at).where(ArtistVisit.artist_id == a.id)
)).scalar_one()
assert visit_at == _LONG_AGO
+170 -3
View File
@@ -1,9 +1,15 @@
"""Smoke test that the FC-3c Celery task is registered and routed correctly. """Tests for the FC-3c download_source Celery task wrapper.
Mirrors test_tasks_register.py — Celery's `include=[...]` is lazy, so Covers registration/routing (smoke) plus the soft-time-limit salvage
the task module must be imported explicitly to trigger registration. path (audit 2026-06-03, Anduo #39912): a SoftTimeLimitExceeded must not
leave the DownloadEvent stranded empty for the recovery sweep.
""" """
from datetime import UTC, datetime
import pytest
from sqlalchemy import select
# Side-effect import: the @celery.task decorator on download_source fires # Side-effect import: the @celery.task decorator on download_source fires
# at module import time and registers the task with the global instance. # at module import time and registers the task with the global instance.
import backend.app.tasks.download # noqa: F401 import backend.app.tasks.download # noqa: F401
@@ -18,3 +24,164 @@ def test_download_source_routes_to_download_queue():
routes = celery.conf.task_routes routes = celery.conf.task_routes
assert "backend.app.tasks.download.*" in routes assert "backend.app.tasks.download.*" in routes
assert routes["backend.app.tasks.download.*"]["queue"] == "download" assert routes["backend.app.tasks.download.*"]["queue"] == "download"
def test_timeout_ladder_keeps_subprocess_budgets_under_soft_limit():
"""Regression guard for Anduo #39912: every gallery-dl subprocess
budget MUST sit below download_source's Celery soft limit so
subprocess.run raises its own TimeoutExpired (which captures partial
logs + finalizes the event) BEFORE Celery's SoftTimeLimitExceeded
preempts it. soft must in turn sit below the hard SIGKILL cap."""
from backend.app.services.gallery_dl import (
_DEFAULT_GDL_TIMEOUT_SECONDS,
BACKFILL_TIMEOUT_SECONDS,
)
from backend.app.tasks.download import (
DOWNLOAD_HARD_TIME_LIMIT,
DOWNLOAD_SOFT_TIME_LIMIT,
)
assert _DEFAULT_GDL_TIMEOUT_SECONDS < DOWNLOAD_SOFT_TIME_LIMIT
assert BACKFILL_TIMEOUT_SECONDS < DOWNLOAD_SOFT_TIME_LIMIT
assert DOWNLOAD_SOFT_TIME_LIMIT < DOWNLOAD_HARD_TIME_LIMIT
def test_decorated_limits_match_module_constants():
"""The @celery.task decorator must use the audited constants, not
drifted literals."""
from backend.app.tasks.download import (
DOWNLOAD_HARD_TIME_LIMIT,
DOWNLOAD_SOFT_TIME_LIMIT,
download_source,
)
assert download_source.soft_time_limit == DOWNLOAD_SOFT_TIME_LIMIT
assert download_source.time_limit == DOWNLOAD_HARD_TIME_LIMIT
def _seed_running_event(db_sync, *, slug, backfill, failures=0):
from backend.app.models import Artist, DownloadEvent, Source
artist = Artist(name=slug, slug=slug)
db_sync.add(artist)
db_sync.flush()
source = Source(
artist_id=artist.id, platform="patreon",
url=f"https://patreon.com/{slug}", enabled=True,
config_overrides={}, backfill_runs_remaining=backfill,
consecutive_failures=failures,
)
db_sync.add(source)
db_sync.flush()
ev = DownloadEvent(
source_id=source.id, status="running",
started_at=datetime.now(UTC),
)
db_sync.add(ev)
db_sync.flush()
return source, ev.id
@pytest.mark.integration
@pytest.mark.asyncio
async def test_finalize_soft_limited_flips_event_and_decrements_backfill(db_sync):
from backend.app.models import DownloadEvent, Source
from backend.app.tasks.download import _finalize_soft_limited
source, event_id = _seed_running_event(
db_sync, slug="anduo", backfill=2, failures=0,
)
_finalize_soft_limited(db_sync, source.id)
status, finished_at, error, meta = db_sync.execute(
select(
DownloadEvent.status, DownloadEvent.finished_at,
DownloadEvent.error, DownloadEvent.metadata_,
).where(DownloadEvent.id == event_id)
).one()
assert status == "error"
assert finished_at is not None
assert "soft time limit" in (error or "").lower()
assert meta.get("error_type") == "timeout"
assert meta.get("soft_time_limited") is True
backfill, failures, error_type = db_sync.execute(
select(
Source.backfill_runs_remaining, Source.consecutive_failures,
Source.error_type,
).where(Source.id == source.id)
).one()
assert backfill == 1 # 2 -> 1, source self-heals toward tick mode
assert failures == 1 # 0 -> 1, mirrors phase-3 source-health write
assert error_type == "timeout"
@pytest.mark.integration
@pytest.mark.asyncio
async def test_finalize_soft_limited_is_noop_without_running_event(db_sync):
"""A benign late soft-limit (phase 3 already committed → no running
event) must not touch source health or backfill budget."""
from backend.app.models import DownloadEvent, Source
from backend.app.tasks.download import _finalize_soft_limited
source, event_id = _seed_running_event(
db_sync, slug="noevent", backfill=3, failures=0,
)
# Simulate phase 3 having already finalized the event.
ev = db_sync.get(DownloadEvent, event_id)
ev.status = "ok"
db_sync.flush()
_finalize_soft_limited(db_sync, source.id)
backfill, failures, error_type = db_sync.execute(
select(
Source.backfill_runs_remaining, Source.consecutive_failures,
Source.error_type,
).where(Source.id == source.id)
).one()
assert backfill == 3 # untouched
assert failures == 0 # untouched
assert error_type is None
@pytest.mark.integration
@pytest.mark.asyncio
async def test_download_source_catches_soft_limit_and_salvages_event(
db_sync, monkeypatch,
):
"""End-to-end wiring: when the inner run raises SoftTimeLimitExceeded,
download_source's handler must flip the in-flight event to error
instead of letting it strand. Uses eager mode + a stubbed asyncio.run
so no real gallery-dl subprocess is spawned."""
from celery.exceptions import SoftTimeLimitExceeded
import backend.app.tasks.download as dl
from backend.app.models import DownloadEvent
monkeypatch.setattr(celery.conf, "task_always_eager", True)
monkeypatch.setattr(celery.conf, "task_eager_propagates", False)
source, event_id = _seed_running_event(
db_sync, slug="anduowire", backfill=1, failures=0,
)
# The task opens a fresh session via _sync_session_factory(); commit
# so that session can see the seeded running event.
db_sync.commit()
def _raise(coro=None, *a, **k):
# Close the un-awaited coroutine so pytest output stays pristine.
if coro is not None and hasattr(coro, "close"):
coro.close()
raise SoftTimeLimitExceeded("simulated soft limit")
monkeypatch.setattr(dl.asyncio, "run", _raise)
with pytest.raises(SoftTimeLimitExceeded):
dl.download_source.delay(source.id).get(propagate=True)
status = db_sync.execute(
select(DownloadEvent.status).where(DownloadEvent.id == event_id)
).scalar_one()
assert status == "error"