"""Fetch and store followed tag/pool feeds from e621. Used by the two periodic management commands (`sync_followed_tags`, `sync_followed_pools`) and by the follow API for the first fill of a new subscription. One e621 search per unique tag/pool is shared across every user following it. """ import logging from django.contrib.auth import get_user_model from django.db.models import Q from django.utils import timezone from apps.library import e621 from apps.library.services import trim_e621_post from .models import FollowedPool, FollowedPost, FollowedTag logger = logging.getLogger(__name__) FEED_LIMIT = 24 PRUNE_KEEP = 200 def preferred_fetch_user(username=None): """A configured account for sync calls; None means run anonymously. e621 serves post searches and pool reads anonymously, so a missing account is not fatal — credentials mainly raise rate limits. """ User = get_user_model() if username: return User.objects.filter(username=username).first() return ( User.objects.filter( Q(is_superuser=True) | Q(is_staff=True) | Q(role=User.ROLE_STAFF) ) .exclude(e621_username="") .exclude(e621_api_key="") .order_by("id") .first() ) def search_posts(user, term, limit=FEED_LIMIT): """Newest posts for a search term (``tag`` or ``pool:``).""" payload = e621.get( user, "/posts.json", params={"tags": term, "limit": limit}, require_auth=False, ) posts = payload.get("posts") if isinstance(payload, dict) else None return posts or [] def cover_payload(post): """Lightweight card cover: enough to render a thumbnail and rating.""" if not isinstance(post, dict): return None preview = post.get("preview") or {} score = post.get("score") or {} return { "post_id": post.get("id"), "rating": post.get("rating"), "preview_url": preview.get("url"), "score": score.get("total"), "fav_count": post.get("fav_count"), } def feed_payload(post): """Compact stored post for the feed (tags included for the tag cloud).""" trimmed = trim_e621_post(post) if trimmed is None: return None return { "id": trimmed.get("id"), "created_at": trimmed.get("created_at"), "rating": trimmed.get("rating"), "tags": trimmed.get("tags") or {}, "score": trimmed.get("score") or {}, "fav_count": trimmed.get("fav_count"), "file": trimmed.get("file") or {}, "preview": trimmed.get("preview") or {}, "uploader_name": trimmed.get("uploader_name"), } def _store_posts(follow, posts, tag=None, pool=None): """Insert unseen rows for new posts. Returns how many were added.""" existing = set(follow.posts.values_list("post_id", flat=True)) rows = [] for post in posts: post_id = post.get("id") if isinstance(post, dict) else None if not isinstance(post_id, int) or post_id in existing: continue data = feed_payload(post) if data is None: continue rows.append( FollowedPost( user=follow.user, post_id=post_id, data=data, tag=tag, pool=pool, ) ) if rows: FollowedPost.objects.bulk_create(rows, ignore_conflicts=True) _prune(follow) return len(rows) def _prune(follow): """Keep only the newest PRUNE_KEEP posts per follow.""" keep = list( follow.posts.order_by("-post_id").values_list("id", flat=True)[ :PRUNE_KEEP ] ) follow.posts.exclude(id__in=keep).delete() def sync_tag_follow(follow: FollowedTag, fetch_user=None, limit=FEED_LIMIT, posts=None): """Fetch new posts for one followed tag; updates cover and sync time.""" if posts is None: posts = search_posts(fetch_user, follow.tag, limit) added = _store_posts(follow, posts, tag=follow) updates = ["last_synced_at"] cover = cover_payload(posts[0]) if posts else None if cover is not None: follow.cover_data = cover updates.append("cover_data") follow.last_synced_at = timezone.now() follow.save(update_fields=updates) return added def sync_pool_follow( follow: FollowedPool, fetch_user=None, limit=FEED_LIMIT, posts=None, refresh_meta=True, ): """Fetch new posts for one followed pool; updates cover and metadata.""" if posts is None: posts = search_posts(fetch_user, f"pool:{follow.pool_id}", limit) added = _store_posts(follow, posts, pool=follow) updates = ["last_synced_at"] cover = cover_payload(posts[0]) if posts else None if cover is not None: follow.cover_data = cover updates.append("cover_data") if refresh_meta or not follow.name: try: pool = e621.get( fetch_user, f"/pools/{follow.pool_id}.json", require_auth=False, ) except e621.E621Error as exc: logger.warning("Could not refresh pool %s: %s", follow.pool_id, exc) else: if isinstance(pool, dict) and pool.get("id"): follow.name = str(pool.get("name") or "")[:200] follow.post_count = int(pool.get("post_count") or 0) updates += ["name", "post_count"] follow.last_synced_at = timezone.now() follow.save(update_fields=updates) return added