14 Commits
Author SHA1 Message Date
JakeBreath 2eb7af0d41 Fix the CD wine build and add per-part toggles
CI / Backend tests (push) Successful in 2m13s
CI / Frontend build & lint (push) Successful in 20s
The Windows NSIS step failed under wine for two reasons: no X display
(nodrv_CreateWindow) and missing 32-bit libraries (failed to load
syswow64\ntdll.dll). The job now installs xvfb + wine32:i386 and runs the
desktop build under xvfb-run, with Gecko/Mono lookups disabled.

Also add `images` and `desktop` dispatch inputs so either half of the CD
can be skipped (e.g. desktop-only or images-only releases).
2026-09-23 00:06:55 -05:00
JakeBreath 7cecfeabc6 Use a minimal registry PAT and the job token for releases
CI / Backend tests (push) Successful in 2m23s
CI / Frontend build & lint (push) Successful in 22s
Gitea's container registry rejects the automatic job token
(go-gitea/gitea#23642 is still open), so the image push keeps a PAT with
only the write:package scope; a preflight step fails clearly when the
REGISTRY_USER/REGISTRY_TOKEN secrets are missing. Release creation needs
no PAT: the desktop job asks for contents: write on the job token.
2026-09-22 23:35:55 -05:00
JakeBreath ed6178d12e Make Actions token permissions explicit
The automatic job token creates the desktop release, so the CD desktop job
asks for contents: write; the images job keeps contents: read and asks for
packages: write so the job token can stand in for the scoped registry PAT.
CI stays read-only. The registry token itself remains a write:package-only
PAT (verified login + pull).
2026-09-22 23:30:37 -05:00
JakeBreath 8f9656ac0e Add the manual CD release workflow
One dispatch builds and pushes both images and builds the desktop packages
into a Gitea release (desktop-v<version>, installers + latest*.yml attached,
idempotent on re-run). The live update feed stays a deploy-host operation:
CI has no SSH key for jakerasp, so push_desktop.sh --no-build remains the
way to publish it.

Repo secrets REGISTRY_USER/REGISTRY_TOKEN are set, so the image push uses
the Gitea registry credentials directly.
2026-09-22 23:23:05 -05:00
JakeBreath 72fc42217f Fix CI for the user-scoped runners
- ci.yml: connect to the test MariaDB as root so Django creates the test
  database itself (no client install/grant step), and drop actions/cache
  (cache: pip/npm): Gitea's cache service hangs the job on restore/save.
- publish.yml: prefer the REGISTRY_USER/REGISTRY_TOKEN secrets (as on other
  repos) and fall back to the automatic Actions token.
- AGENTS.md: note the CI layout, the runner labels and the cache caveat.
2026-09-22 23:12:56 -05:00
JakeBreath d9c1e9e521 Add CI and manual image publishing workflows
CI / Backend tests (push) Successful in 12m54s
CI / Frontend build & lint (push) Failing after 4m59s
- .gitea/workflows/ci.yml: on every push/PR, run Django checks + the full
  backend suite against MariaDB/Redis services and the frontend
  lint/type-check/build. Runs on the nitro-ci runner (ubuntu-latest).
- .gitea/workflows/publish.yml: manual dispatch; multi-arch build+push of
  both images as :latest and :<short-sha> with GIT_HASH baked in.
- push_*.sh: non-interactive registry login for CI (REGISTRY_USER/
  REGISTRY_TOKEN) and a PLATFORMS override.
2026-09-22 22:32:37 -05:00
JakeBreath e2697c0a78 Stop throttling signed media and ease the browser's e621 queue
Signed media URLs are fetched by <img>/<video> tags without an
Authorization header, so they were charged to the anonymous 120/min
bucket: past that, galleries and the fish-greeting download got 429 JSON
instead of image bytes. The raw/thumbnail/staged-file/similarity-file
actions are now exempt, and THROTTLE_ENABLED=false removes the general
anon+user limits for private/tailnet deployments (login/register/proxy
guards stay).

The SPA's e621 client also stops self-throttling so hard: 1s gap between
browsing calls (2.5s for the stricter IQDB endpoint) and a 15s cooldown
instead of 60s when e621 answers 429.
2026-09-22 22:32:37 -05:00
JakeBreath 474403ffe2 Upload updates 2026-09-21 09:01:01 -05:00
JakeBreath 98025e9e6d Prune old installers from the remote feed on push
rsync without --delete left every previous version on the server (the
screenshot showed 0.1.0 and 0.1.1 side by side). The push now removes
non-current installers over ssh first, so the remote feed mirrors the local
one whether the transfer uses rsync or tar.
2026-09-20 21:39:43 -05:00
JakeBreath 3183a3bece Prune old desktop builds from release/ on every build
build_desktop.sh now reads the version first and removes anything in
desktop/release/ that is not that version (plus the regenerated unpacked
trees), so a version bump never leaves old installers lying around — the
same rule push_desktop.sh applies to the feed.
2026-09-20 21:16:56 -05:00
JakeBreath 041a9d4471 Keep desktop builds and the feed to one version
Both scripts now read the version from desktop/package.json: the build
report and checksums only cover the current version's artifacts, the feed
copy ignores older files, and publishing prunes previous installers from
the feed (latest*.yml only ever points at the current one).
2026-09-20 21:11:19 -05:00
JakeBreath d84fdd0e98 Bump the desktop app to 0.1.1
The icon set, e621 referrer fix and setup-screen corrections shipped after
0.1.0, and the updater compares versions, so installed 0.1.0 builds would
never have seen them.
2026-09-20 21:09:42 -05:00
JakeBreath c992a63b8f Fix e621 images and the setup screen in the desktop shell
e621's CDN answers cross-site image loads that carry no Referer with a 403
(Chromium sends none from a custom-scheme page, then blocks the response as
ORB), so images never appeared in the desktop app. The main process now
attaches an e621 referrer to requests for its hosts.

The shell also answers /api, /admin, /static and /health with a 404 JSON
instead of the SPA fallback — that fallback made the setup screen's empty-URL
connection test report "Connected" against the shell itself. The setup screen
is now desktop-aware (no same-origin option, no "Use this server", clearer
copy), and the smoke test runs against a throwaway profile and covers both
regressions.
2026-09-20 21:08:56 -05:00
JakeBreath bd2417aff8 Keep a broken Redis from 500ing the whole API
Redis backs the DRF throttles, and the stock RedisCache raises inside the
throttle check when Redis is unreachable or refusing writes (a failed RDB
snapshot disables writes by default) — turning a cache problem into a
blanket 500, which is exactly what took prod down. ResilientRedisCache
treats backend failures as cache misses, logs the first one per worker, and
lets rate limits degrade until Redis is back.
2026-09-20 21:08:38 -05:00
37 changed files with 3238 additions and 848 deletions
+182
View File
@@ -0,0 +1,182 @@
# J621 CD — manual release workflow (Actions tab -> "Run workflow").
#
# One dispatch does everything; each half can be skipped with the `images`
# and `desktop` inputs:
# * builds and pushes the backend + frontend images (multi-arch, :latest
# and :<short-sha>, GIT_HASH baked in for the version pill),
# * builds the desktop packages and attaches them (plus the update
# metadata) to the Gitea release tagged `desktop-v<package.json version>`.
#
# The live update feed (deploy/data/desktop, served by the frontend nginx at
# /desktop/) is not touched here: it is runtime state on the deploy host and
# is still published with `deploy/push_desktop.sh --no-build` from a machine
# that can reach it.
#
# Registry login uses a repo PAT with the minimal write:package scope (the
# Gitea registry rejects the automatic job token, go-gitea/gitea#23642);
# release creation uses the automatic job token. Jobs run on the user-scoped
# nitro-ci runner (ubuntu-latest).
name: CD
on:
workflow_dispatch:
inputs:
images:
description: Build and push the Docker images
required: false
default: "true"
desktop:
description: Build the desktop release
required: false
default: "true"
platforms:
description: Image platforms (comma separated)
required: false
default: linux/amd64,linux/arm64
windows:
description: Also cross-build the Windows installer (needs wine, slow)
required: false
default: "false"
concurrency:
group: cd
cancel-in-progress: false
jobs:
images:
name: Build & push images
if: ${{ inputs.images != 'false' }}
runs-on: ubuntu-latest
# The Gitea container registry does not accept the automatic job token
# (go-gitea/gitea#23642 is still open), so the push uses a repo PAT with
# the minimal write:package scope. Releases use the job token instead.
permissions:
contents: read
env:
REGISTRY_USER: ${{ secrets.REGISTRY_USER }}
REGISTRY_TOKEN: ${{ secrets.REGISTRY_TOKEN }}
PLATFORMS: ${{ inputs.platforms || 'linux/amd64,linux/arm64' }}
steps:
- uses: actions/checkout@v4
with:
fetch-depth: 0
- name: Check the registry credentials
run: |
if [ -z "$REGISTRY_USER" ] || [ -z "$REGISTRY_TOKEN" ]; then
echo "Set the REGISTRY_USER and REGISTRY_TOKEN repo secrets" >&2
echo "(a PAT with the write:package scope)." >&2
exit 1
fi
- name: Register binfmt (multi-arch builds)
run: docker run --privileged --rm tonistiigi/binfmt --install all
- name: Build & push both images
run: |
set -euo pipefail
SHA="$(git rev-parse --short HEAD)"
echo "Publishing $SHA for $PLATFORMS"
PLATFORMS="$PLATFORMS" ./deploy/push_frontend.sh "$SHA"
PLATFORMS="$PLATFORMS" ./deploy/push_backend.sh "$SHA"
desktop:
name: Desktop release
if: ${{ inputs.desktop != 'false' }}
runs-on: ubuntu-latest
# Creating the release and uploading its assets uses the automatic job
# token, so it needs write access to the repository's releases.
permissions:
contents: write
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: "22"
- name: Install packaging tools
run: |
sudo apt-get update
sudo apt-get install -y --no-install-recommends fakeroot libarchive-tools
if [ "${{ inputs.windows }}" = "true" ]; then
# electron-builder runs the 32-bit NSIS installer under wine to
# build the uninstaller: that needs a virtual display (Xvfb) and
# 32-bit wine libraries.
sudo dpkg --add-architecture i386
sudo apt-get update
sudo apt-get install -y --no-install-recommends xvfb wine wine32:i386
fi
- name: Install frontend + desktop dependencies
run: |
npm --prefix frontend ci --no-audit --no-fund
npm --prefix desktop ci --no-audit --no-fund
- name: Build desktop packages
env:
# Keep wine from trying to fetch Gecko/Mono on first run.
WINEDLLOVERRIDES: mscoree,mshtml=
run: |
if [ "${{ inputs.windows }}" = "true" ]; then
xvfb-run -a ./deploy/build_desktop.sh --all
else
./deploy/build_desktop.sh --linux
fi
- name: Add the Gitea release
env:
GITEA_TOKEN: ${{ secrets.GITEA_TOKEN || secrets.GITHUB_TOKEN }}
run: |
set -euo pipefail
VERSION="$(node -p "require('./desktop/package.json').version")"
TAG="desktop-v$VERSION"
API="${{ github.server_url }}/api/v1/repos/${{ github.repository }}"
AUTH="Authorization: token $GITEA_TOKEN"
NOTES="$(printf 'J621 desktop %s\n\n' "$VERSION"
cd desktop/release
sha256sum ./*.deb ./*.pkg.tar.zst ./*.exe 2>/dev/null || true)"
RELEASE_ID="$(curl -sf -H "$AUTH" "$API/releases/tags/$TAG" \
| python3 -c 'import json,sys; print(json.load(sys.stdin).get("id",""))' \
2>/dev/null || true)"
if [ -z "$RELEASE_ID" ]; then
echo "Creating release $TAG"
PAYLOAD="$(python3 - "$TAG" "${{ github.sha }}" "$NOTES" <<'PY'
import json, sys
print(json.dumps({
"tag_name": sys.argv[1],
"name": sys.argv[1],
"body": sys.argv[3],
"target_commitish": sys.argv[2],
}))
PY
)"
RELEASE_ID="$(curl -sf -X POST -H "$AUTH" \
-H "Content-Type: application/json" -d "$PAYLOAD" "$API/releases" \
| python3 -c 'import json,sys; print(json.load(sys.stdin)["id"])')"
else
echo "Release $TAG already exists (id $RELEASE_ID); attaching missing files."
fi
EXISTING="$(curl -sf -H "$AUTH" "$API/releases/$RELEASE_ID/assets" \
| python3 -c 'import json,sys; print("\n".join(a["name"] for a in json.load(sys.stdin)))' \
|| true)"
for FILE in desktop/release/*"$VERSION"*.deb \
desktop/release/*"$VERSION"*.pkg.tar.zst \
desktop/release/latest-linux.yml \
desktop/release/latest.yml \
desktop/release/*"$VERSION"*.exe \
desktop/release/*"$VERSION"*.exe.blockmap; do
[ -e "$FILE" ] || continue
NAME="$(basename "$FILE")"
case "$EXISTING" in
*"$NAME"*) echo " already attached: $NAME"; continue ;;
esac
echo " attaching $NAME"
curl -sf -X POST -H "$AUTH" -H "Content-Type: application/octet-stream" \
--data-binary @"$FILE" "$API/releases/$RELEASE_ID/assets?name=$NAME" >/dev/null
done
echo "Release: ${{ github.server_url }}/${{ github.repository }}/releases/tag/$TAG"
+101
View File
@@ -0,0 +1,101 @@
# J621 CI — runs on every push (and pull request): Django checks + the full
# backend test suite against MariaDB/Redis, and the frontend type-check,
# lint and production build.
#
# Runner: the "nitro-ci" act_runner with the custom `ubuntu-latest` label.
name: CI
# Tests and builds only need to read the repository; the automatic job token
# stays read-only.
permissions:
contents: read
on:
push:
branches: ["**"]
tags-ignore: ["**"]
pull_request:
workflow_dispatch:
concurrency:
group: ci-${{ github.ref }}
cancel-in-progress: true
jobs:
backend:
name: Backend tests
runs-on: ubuntu-latest
services:
mariadb:
image: mariadb:11.4
env:
MARIADB_ROOT_PASSWORD: root
MARIADB_DATABASE: j621
MARIADB_USER: j621
MARIADB_PASSWORD: j621
options: >-
--health-cmd="healthcheck.sh --connect --innodb_initialized"
--health-interval=5s
--health-timeout=5s
--health-retries=12
redis:
image: redis:7-alpine
options: >-
--health-cmd="redis-cli ping"
--health-interval=5s
--health-timeout=5s
--health-retries=12
env:
# Connect as root so Django can create the test database itself;
# everything else mirrors the development defaults.
DB_HOST: mariadb
DB_PORT: "3306"
DB_NAME: j621
DB_USER: root
DB_PASSWORD: root
REDIS_URL: redis://redis:6379/1
steps:
- uses: actions/checkout@v4
- uses: actions/setup-python@v5
with:
python-version: "3.14"
- name: Install backend dependencies
run: pip install -r backend/requirements.txt
- name: Django system checks
working-directory: backend
run: python manage.py check
- name: Backend tests
working-directory: backend
run: >-
python manage.py test
apps.core.tests
apps.library.tests
apps.follows.tests
apps.accounts.tests
frontend:
name: Frontend build & lint
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: "22"
- name: Install frontend dependencies
working-directory: frontend
run: npm ci
- name: Lint
working-directory: frontend
run: npm run lint
- name: Type-check & build
working-directory: frontend
run: npm run build
+16 -1
View File
@@ -49,6 +49,16 @@ Project constraints (do not regress):
- Periodic commands (follow syncs, similarity cleanup, guest blacklist
refresh) run in the composes' `scheduler` service — the backend image with
the j621-scheduler entrypoint, intervals via J621_*_EVERY. No host cron.
- CI/CD lives in .gitea/workflows: ci.yml runs on every push/PR (Django checks
+ the full backend suite against MariaDB/Redis service containers, frontend
lint/type-check/build); cd.yml is manual and builds/pushes both images
multi-arch plus the desktop packages (attached to the Gitea release
`desktop-v<version>`). Jobs run on the user-scoped runners: `ubuntu-latest`
on nitro-ci, `desktop` on msi-mortar-ci. Do not add actions/cache
(`cache: pip`/`npm`) to these workflows: Gitea's cache service hangs the job
on restore/save. The live desktop update feed (deploy/data/desktop) is still
published with `deploy/push_desktop.sh --no-build` from a machine with SSH
to the deploy host — CI has no key for that.
- Security/permission tests live in backend/apps/core/tests and need a
one-time grant: GRANT ALL ON `test_j621`.* TO 'j621'@'%';
@@ -85,7 +95,12 @@ Security hardening (do not weaken):
SECRET_KEY (apps/accounts/crypto.py); rotating SECRET_KEY invalidates them
(and all signed media URLs), so users must re-enter the key.
- API throttles live in REST_FRAMEWORK (env-overridable): anon 120/min,
user 600/min, login 5/min, register 20/hour, e621_proxy 60/hour.
user 600/min, login 5/min, register 20/hour, e621_proxy 60/hour. Signed
media URLs (raw/thumbnail/staged-file/similarity-file actions) are exempt
on purpose: <img>/<video> tags fetch them without an Authorization header,
so a gallery would otherwise drain the anonymous bucket and get 429 JSON
instead of images. THROTTLE_ENABLED=false removes the anon+user limits for
private/tailnet deployments (the login/register/proxy guards stay).
- Only admins (superusers) may grant/revoke the staff role or delete
staff/admin accounts; staff manage regular/uploader accounts only.
- Storage, duplicates, delete, temp-clear, uploads and downloads require
+12 -9
View File
@@ -95,18 +95,21 @@ Files now stage first and are resolved before entering the library.
- [x] `cleanup_temp_uploads` command for old staged files
- [x] **Auto-upload / auto-match**
- [x] MD5 computed on staging; exact duplicates resolve immediately
- [x] MD5 batch-checked against e621; matches auto-complete with post
metadata stored and rating seeded
- [x] **IQDB similarity on upload** (SPA-driven)
- [x] Automatic + manual IQDB checks with candidate posts
- [x] "Visual Similarity Detected" state with candidate picker
- [x] Perceptual-hash comparison against the library (staged uploads are
flagged with their library matches as soon as they land)
- [x] MD5 batch-checked against e621 in chunks of 75; matches auto-complete
with post metadata stored and rating seeded
- [x] **Background pipeline** (server-side)
- [x] Daemon-thread worker runs MD5 → visual similarity → IQDB for every
staged upload, so the work continues after the page or tab is closed
- [x] Durable progress (`UploadRun` + per-file phase flags) polled by the
shell indicator; rate-limited runs retry with backoff
- [x] Perceptual-hash comparison against the library, loaded once per batch
- [x] IQDB candidates stored with one batched enrichment request
- [x] **Upload UI**
- [x] Three-column board: Pending & Unmatched / Visual Similarity Detected /
Auto-uploaded & Indexed
Auto-uploaded & Indexed, private per user (staff included)
- [x] Metadata modal (link to e621 post, IQDB candidates, custom metadata)
- [x] Per-file progress plus batch processing indicator
- [x] Per-file progress plus background pipeline status
- [x] Bulk actions: bulk rate, discard all (pending/visual), dismiss all
## 4. Staff tools
+108
View File
@@ -0,0 +1,108 @@
"""Redis cache that degrades instead of taking the whole API down.
Redis backs the DRF throttles and a few caches (storage stats, guest
blacklist, tag clouds). With Django's stock ``RedisCache``, a Redis that is
unreachable — or merely refusing writes because its RDB snapshot failed, the
default ``stop-writes-on-bgsave-error`` behaviour — raises inside the
throttle check on every request, so a cache outage becomes a blanket 500.
This backend treats cache failures as misses: rate limits and cached values
simply stop working until Redis is back, and the first failure per worker is
logged once so the cause is still visible.
"""
import logging
from django.core.cache.backends.redis import RedisCache
logger = logging.getLogger(__name__)
_warned = False
def _degrade(operation: str, error: Exception, default):
global _warned
if not _warned:
_warned = True
logger.warning(
"Cache unavailable (%s failed: %s) — continuing without it.",
operation,
error,
)
return default
class ResilientRedisCache(RedisCache):
"""``RedisCache`` where a broken Redis behaves like an empty cache."""
def add(self, *args, **kwargs):
try:
return super().add(*args, **kwargs)
except Exception as error: # noqa: BLE001 - any backend failure degrades
return _degrade("add", error, False)
def get(self, key, default=None, version=None):
try:
return super().get(key, default, version)
except Exception as error: # noqa: BLE001
return _degrade("get", error, default)
def set(self, *args, **kwargs):
try:
return super().set(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("set", error, None)
def touch(self, *args, **kwargs):
try:
return super().touch(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("touch", error, False)
def delete(self, *args, **kwargs):
try:
return super().delete(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("delete", error, False)
def get_many(self, *args, **kwargs):
try:
return super().get_many(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("get_many", error, {})
def has_key(self, *args, **kwargs):
try:
return super().has_key(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("has_key", error, False)
def incr(self, *args, **kwargs):
try:
return super().incr(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("incr", error, None)
def set_many(self, *args, **kwargs):
try:
return super().set_many(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("set_many", error, [])
def delete_many(self, *args, **kwargs):
try:
return super().delete_many(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("delete_many", error, None)
def clear(self, *args, **kwargs):
try:
return super().clear(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("clear", error, None)
def close(self, *args, **kwargs):
try:
return super().close(*args, **kwargs)
except Exception as error: # noqa: BLE001
return _degrade("close", error, None)
+36
View File
@@ -0,0 +1,36 @@
"""A broken Redis must degrade the cache, not 500 the API.
A Redis that cannot persist (the default ``stop-writes-on-bgsave-error``)
or is simply unreachable used to raise inside the DRF throttle check on
every request; ``ResilientRedisCache`` treats that as a cache miss.
"""
from unittest import mock
from django.core.cache import cache
from django.core.cache.backends.redis import RedisCacheClient
from django.test import SimpleTestCase
def failing(method: str):
return mock.patch.object(
RedisCacheClient,
method,
side_effect=RuntimeError("redis is down"),
)
class ResilientCacheTests(SimpleTestCase):
def test_get_returns_the_default_when_redis_fails(self):
with failing("get"):
self.assertIsNone(cache.get("j621-cache-test"))
self.assertEqual(cache.get("j621-cache-test", "fallback"), "fallback")
def test_writes_report_failure_without_raising(self):
with failing("set"):
self.assertIsNone(cache.set("j621-cache-test", "value"))
def test_bulk_and_delete_operations_degrade(self):
with failing("get_many"), failing("delete"):
self.assertEqual(cache.get_many(["a", "b"]), {})
self.assertFalse(cache.delete("a"))
+37 -1
View File
@@ -345,7 +345,13 @@ class PrivacyTests(SecurityTestCase):
)
self.assertEqual(self.guest.get(f"/api/uploads/{temp.id}/file/").status_code, 401)
self.assertNotIn(
str(temp.id), self.client_for("sec-uploader").get("/api/uploads/").content.decode()
str(temp.id),
self.client_for("sec-uploader").get("/api/uploads/").content.decode(),
)
# The board is per-user: staff only see their own staged uploads.
self.assertNotIn(
str(temp.id),
self.client_for("sec-staff").get("/api/uploads/").content.decode(),
)
def test_similarity_checks_are_private(self):
@@ -423,6 +429,36 @@ class ThrottleTests(SecurityTestCase):
{self.guest.get("/api/status/").status_code for _ in range(12)}, {200}
)
def test_signed_media_urls_are_not_throttled(self):
"""<img>/<video> tags fetch these without an Authorization header.
Regression: they were charged to the anonymous bucket, so galleries
and the fish-greeting download started returning 429 JSON instead of
the image bytes.
"""
item = self.make_item("throttle-media", owner=self.users["sec-uploader"])
codes = {
self.guest.get(f"/api/files/J-{item.id}/raw/").status_code
for _ in range(150)
}
self.assertEqual(codes, {200})
def test_staged_upload_files_are_not_throttled(self):
temp = TempUpload.objects.create(
user=self.users["sec-uploader"],
file=SimpleUploadedFile("throttle-temp.bin", b"staged"),
original_filename="throttle-temp.bin",
md5=hashlib.md5(b"throttle-temp").hexdigest(),
size=6,
)
signature = signing.dumps(
{"temp": str(temp.id), "user": self.users["sec-uploader"].id},
salt=services.UPLOAD_FILE_SALT,
)
url = f"/api/uploads/{temp.id}/file/?sig={signature}"
codes = {self.guest.get(url).status_code for _ in range(150)}
self.assertEqual(codes, {200})
class RemoteUrlTests(SecurityTestCase):
def test_allowlist(self):
+257 -34
View File
@@ -1,22 +1,39 @@
"""Minimal e621 API client for server-side matching and metadata refresh.
The SPA talks to e621 directly for browsing; this client exists for work the
browser cannot do reliably: long batch scans, and requests tied to a library
item rather than an open page. It uses the requesting user's stored
credentials and a global throttle (e621 asks for at most two requests per
second).
browser cannot do reliably: long batch scans, staged-upload processing and
requests tied to a library item rather than an open page. It uses the
requesting user's stored credentials and a global throttle (e621 asks for at
most two requests per second, one per second sustained).
e621's load balancer also sheds load with 429s (sometimes with an HTML
"shedding" page instead of JSON) and the IQDB endpoint has its own, much
stricter throttle. Every call therefore retries with exponential backoff and
honours ``Retry-After``; only 401/403 are treated as fatal.
"""
import logging
import random
import threading
import time
from pathlib import Path
import requests
from django.conf import settings
logger = logging.getLogger(__name__)
# Seconds between requests, per process. e621 allows 2/s hard and 1/s
# sustained; each gunicorn worker throttles on its own, so leave enough
# headroom that combined traffic does not trip the limit.
REQUEST_INTERVAL = 1.0
# IQDB is throttled far more aggressively than the rest of the API.
IQDB_INTERVAL = 2.0
MAX_ATTEMPTS = 4
BACKOFF_BASE = 2.0
BACKOFF_CAP = 60.0
RETRYABLE_STATUSES = {429, 500, 502, 503, 504}
class E621Error(Exception):
@@ -27,6 +44,14 @@ class E621NotFound(E621Error):
"""The requested post does not exist (HTTP 404)."""
class E621AuthError(E621Error):
"""e621 rejected the stored credentials (401/403)."""
class E621RateLimited(E621Error):
"""e621 shed load or throttled the request after every retry."""
_throttle_lock = threading.Lock()
_last_request_at = 0.0
@@ -35,56 +60,170 @@ def credentials_configured(user):
return bool(user is not None and getattr(user, "e621_configured", False))
def _wait_for_slot():
def _wait_for_slot(interval=REQUEST_INTERVAL):
global _last_request_at
with _throttle_lock:
delay = _last_request_at + REQUEST_INTERVAL - time.monotonic()
delay = _last_request_at + interval - time.monotonic()
if delay > 0:
time.sleep(delay)
_last_request_at = time.monotonic()
def _retry_delay(attempt, response=None):
"""Backoff for a retryable failure, honouring ``Retry-After``."""
if response is not None:
retry_after = response.headers.get("Retry-After")
if retry_after:
try:
return max(float(retry_after), 1.0)
except (TypeError, ValueError):
pass
delay = min(BACKOFF_BASE * (2**attempt), BACKOFF_CAP)
return delay + random.uniform(0, delay * 0.25)
def _request(
user,
method,
path,
*,
params=None,
data=None,
files=None,
timeout=30,
require_auth=True,
interval=REQUEST_INTERVAL,
attempts=MAX_ATTEMPTS,
):
"""One e621 call with retries; returns the parsed JSON payload.
``files`` may be a callable returning the multipart mapping, which is
called once per attempt: streamed uploads consume their file handle, so a
retry needs a freshly opened file.
Raises E621NotFound for 404s, E621AuthError for 401/403 and
E621RateLimited when e621 keeps shedding/throttling after every attempt.
"""
configured = credentials_configured(user)
if require_auth and not configured:
raise E621Error("Configure your e621 credentials in Account first.")
base = (getattr(user, "e621_base_url", "") or "https://e621.net").rstrip("/")
auth = (
(user.e621_username, user.e621_api_key_plain) if configured else None
)
url = f"{base}{path}"
last_error = None
for attempt in range(attempts):
request_files = files() if callable(files) else files
_wait_for_slot(interval)
response = None
try:
response = requests.request(
method,
url,
params=params,
data=data,
files=request_files,
auth=auth,
headers={"User-Agent": settings.USER_AGENT},
timeout=timeout,
)
except requests.RequestException as exc:
last_error = E621Error(f"Could not reach e621: {exc}")
else:
if response.status_code == 404:
raise E621NotFound(f"e621 returned 404 for {path}")
if response.status_code in {401, 403}:
raise E621AuthError(
f"e621 rejected the request ({response.status_code}). "
"Check the stored e621 credentials."
)
if response.status_code == 429:
# Throttles and load-shedding can arrive as JSON ({"message":
# "Throttled: ..."}) or as an HTML page.
message = ""
try:
payload = response.json()
except ValueError:
payload = None
if isinstance(payload, dict):
message = str(
payload.get("message") or payload.get("error") or ""
)
last_error = E621RateLimited(
message or f"e621 throttled the request for {path}"
)
elif response.status_code < 400:
try:
return response.json()
except ValueError:
# An HTML page with a 2xx status.
last_error = E621RateLimited(
f"e621 returned an unexpected {response.status_code} response."
)
elif response.status_code in RETRYABLE_STATUSES:
last_error = E621RateLimited(
f"e621 replied {response.status_code} for {path}"
)
else:
raise E621Error(f"e621 replied {response.status_code} for {path}")
finally:
_close_upload_files(request_files)
if attempt + 1 < attempts:
delay = _retry_delay(attempt, response)
logger.info(
"e621 %s %s failed (%s); retrying in %.1fs",
method,
path,
last_error,
delay,
)
time.sleep(delay)
if last_error is None:
last_error = E621Error("e621 request failed.")
raise last_error
def _close_upload_files(files):
"""Close the handles behind a multipart mapping (see _request)."""
if not isinstance(files, dict):
return
for value in files.values():
handle = value[1] if isinstance(value, tuple) and len(value) > 1 else value
close = getattr(handle, "close", None)
if close is not None:
try:
close()
except Exception: # noqa: BLE001 - closing must never mask errors
pass
def get(user, path, params=None, timeout=30, require_auth=True):
"""GET an e621 API path using the user's credentials.
Reads that e621 serves anonymously (searches, pools, tags) can pass
require_auth=False; matching endpoints keep requiring credentials.
Raises E621NotFound for 404s and E621Error for everything else that isn't
a 2xx, so callers never see requests exceptions.
"""
configured = credentials_configured(user)
if require_auth and not configured:
raise E621Error("Configure your e621 credentials in Account first.")
base = (getattr(user, "e621_base_url", "") or "https://e621.net").rstrip("/")
_wait_for_slot()
try:
response = requests.get(
f"{base}{path}",
return _request(
user,
"GET",
path,
params=params,
auth=(
(user.e621_username, user.e621_api_key_plain)
if configured
else None
),
headers={"User-Agent": settings.USER_AGENT},
timeout=timeout,
require_auth=require_auth,
)
except requests.RequestException as exc:
raise E621Error(f"Could not reach e621: {exc}") from exc
if response.status_code == 404:
raise E621NotFound(f"e621 returned 404 for {path}")
if response.status_code >= 400:
raise E621Error(f"e621 replied {response.status_code} for {path}")
try:
return response.json()
except ValueError as exc:
raise E621Error("e621 returned an unexpected response.") from exc
def find_post_by_md5(user, md5):
"""The e621 post with this exact MD5, or None."""
payload = get(user, "/posts.json", params={"tags": f"md5:{md5}", "limit": 1})
payload = get(
user,
"/posts.json",
params={"tags": f"md5:{md5}", "limit": 1},
require_auth=False,
)
posts = payload.get("posts") if isinstance(payload, dict) else None
if not posts:
return None
@@ -98,3 +237,87 @@ def fetch_post(user, post_id):
if not isinstance(post, dict):
raise E621Error("e621 returned an unexpected post payload.")
return post
def check_md5_batch(user, md5s):
"""Look many MD5s up in one posts.json query.
Returns ``{md5: post}`` for the ones e621 knows; missing MD5s are simply
absent. Works anonymously, like the original app's batch cache command.
"""
wanted = {str(value).strip().lower() for value in md5s if value}
if not wanted:
return {}
values = sorted(wanted)
payload = get(
user,
"/posts.json",
params={
"tags": f"md5:{','.join(values)}",
"limit": min(len(values), 320),
},
require_auth=False,
)
posts = payload.get("posts") if isinstance(payload, dict) else None
found = {}
for post in posts or []:
if not isinstance(post, dict):
continue
file_data = post.get("file") or {}
md5 = str(file_data.get("md5") or "").strip().lower()
if md5 in wanted:
found[md5] = post
return found
def fetch_posts_by_ids(user, ids):
"""Fetch many posts in one query (up to 320 ids). Missing ids are absent."""
values = sorted({int(value) for value in ids})
if not values:
return []
payload = get(
user,
"/posts.json",
params={
"tags": f"id:{','.join(str(value) for value in values)}",
"limit": min(len(values), 320),
},
require_auth=False,
)
posts = payload.get("posts") if isinstance(payload, dict) else None
return [post for post in posts or [] if isinstance(post, dict)]
def iqdb_search(user, path, timeout=60):
"""Reverse-image search one file through e621's IQDB endpoint.
Returns the legacy match list. Uses the extra-strict IQDB interval and
retries through e621's throttle; raises E621RateLimited when it persists.
"""
path = Path(path)
def open_file():
# A fresh handle per attempt: the stream is consumed by the request.
return {"search[file]": (path.name, open(path, "rb"))}
payload = _request(
user,
"POST",
"/iqdb_queries.json",
files=open_file,
timeout=timeout,
require_auth=False,
interval=IQDB_INTERVAL,
)
if isinstance(payload, list):
return payload
if isinstance(payload, dict):
matches = payload.get("matches")
if isinstance(matches, list):
return matches
# e621 answers its throttle with {"success": false, "message": ...}.
message = payload.get("message") or payload.get("error")
if message:
raise E621RateLimited(str(message))
raise E621Error("e621 returned an unexpected IQDB payload.")
return []
@@ -18,6 +18,9 @@ class Command(BaseCommand):
)
def handle(self, *args, **options):
from apps.library.upload_pipeline import reap_stale_claims
reap_stale_claims()
cutoff = timezone.now() - timedelta(hours=options["hours"])
queryset = TempUpload.objects.filter(created_at__lt=cutoff)
if not options["include_completed"]:
@@ -0,0 +1,68 @@
"""Process staged uploads: e621 MD5, visual similarity, IQDB.
Runs synchronously, unlike the daemon thread the API starts on demand. Useful
for tests, for a manual drain after an outage and for the scheduler if a
deployment wants a periodic safety net.
python manage.py process_uploads # every user with queued work, once
python manage.py process_uploads --user 3 # one user
python manage.py process_uploads --loop 60 # keep draining every 60s
"""
import time
from django.core.management.base import BaseCommand
from apps.library.models import TempUpload
from apps.library.upload_pipeline import (
MAX_ATTEMPTS,
OUTSTANDING_Q,
WORK_STATUSES,
reap_stale_claims,
run_pipeline,
)
class Command(BaseCommand):
help = "Run the staged-upload pipeline (MD5 -> visual similarity -> IQDB)."
def add_arguments(self, parser):
parser.add_argument(
"--user",
type=int,
default=None,
help="Only process this user id.",
)
parser.add_argument(
"--loop",
type=int,
default=0,
metavar="SECONDS",
help="Keep draining every SECONDS seconds instead of exiting.",
)
def handle(self, *args, **options):
interval = options["loop"] or 0
while True:
self.drain(user_id=options["user"])
if interval <= 0:
return
time.sleep(interval)
def drain(self, user_id=None):
reap_stale_claims()
queryset = (
TempUpload.objects.filter(status__in=WORK_STATUSES)
.filter(OUTSTANDING_Q)
.filter(attempts__lt=MAX_ATTEMPTS)
)
if user_id is not None:
queryset = queryset.filter(user_id=user_id)
user_ids = list(queryset.values_list("user_id", flat=True).distinct())
if not user_ids:
self.stdout.write("No staged uploads need processing.")
return
for value in user_ids:
self.stdout.write(f"Processing staged uploads for user {value}...")
run_pipeline(value)
self.stdout.write(self.style.SUCCESS(f"Processed {len(user_ids)} queue(s)."))
@@ -0,0 +1,81 @@
# Generated by Django 6.1.1 on 2026-09-21 13:27
import django.db.models.deletion
from django.conf import settings
from django.db import migrations, models
def backfill_pipeline_checks(apps, schema_editor):
"""Mark pre-pipeline rows as already MD5/visual-checked when they were.
Rows with an e621 post id (or already completed) clearly went through the
MD5 phase; rows with visual matches went through the visual phase.
Everything else stays unset so the new pipeline picks it up once after
deploy — a re-check of stale pending uploads is the desired behavior.
"""
TempUpload = apps.get_model("library", "TempUpload")
TempUpload.objects.filter(
models.Q(e621_post_id__isnull=False) | models.Q(status="completed")
).update(e621_checked_at=models.F("updated_at"))
TempUpload.objects.filter(visual_matches__isnull=False).update(
visual_checked_at=models.F("updated_at")
)
class Migration(migrations.Migration):
dependencies = [
('library', '0010_similaritycheck'),
migrations.swappable_dependency(settings.AUTH_USER_MODEL),
]
operations = [
migrations.AddField(
model_name='tempupload',
name='attempts',
field=models.PositiveSmallIntegerField(default=0),
),
migrations.AddField(
model_name='tempupload',
name='claimed_at',
field=models.DateTimeField(blank=True, db_index=True, null=True),
),
migrations.AddField(
model_name='tempupload',
name='e621_checked_at',
field=models.DateTimeField(blank=True, null=True),
),
migrations.AddField(
model_name='tempupload',
name='pipeline_error',
field=models.TextField(blank=True, default=''),
),
migrations.AddField(
model_name='tempupload',
name='visual_checked_at',
field=models.DateTimeField(blank=True, null=True),
),
migrations.CreateModel(
name='UploadRun',
fields=[
('id', models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name='ID')),
('status', models.CharField(choices=[('idle', 'Idle'), ('running', 'Running'), ('paused', 'Paused'), ('error', 'Error')], default='idle', max_length=20)),
('phase', models.CharField(blank=True, choices=[('', 'None'), ('md5', 'e621 MD5'), ('visual', 'Visual similarity'), ('iqdb', 'IQDB')], default='', max_length=20)),
('total', models.IntegerField(default=0)),
('processed', models.IntegerField(default=0)),
('matched', models.IntegerField(default=0)),
('failed', models.IntegerField(default=0)),
('error', models.TextField(blank=True, default='')),
('started_at', models.DateTimeField(blank=True, null=True)),
('updated_at', models.DateTimeField(auto_now=True)),
('user', models.OneToOneField(on_delete=django.db.models.deletion.CASCADE, related_name='upload_run', to=settings.AUTH_USER_MODEL)),
],
options={
'ordering': ['-updated_at'],
},
),
migrations.RunPython(
code=backfill_pipeline_checks, reverse_code=migrations.RunPython.noop
),
]
+65
View File
@@ -164,6 +164,16 @@ class TempUpload(models.Model):
on_delete=models.SET_NULL,
related_name="temp_uploads",
)
# Background pipeline bookkeeping (see apps/library/upload_pipeline.py).
# e621_checked_at/visual_checked_at are set once the corresponding phase
# ran, so "never checked" and "checked, nothing found" stay distinct.
e621_checked_at = models.DateTimeField(null=True, blank=True)
visual_checked_at = models.DateTimeField(null=True, blank=True)
# Worker claim for cross-process mutual exclusion; stale claims are
# reaped and the row queued again.
claimed_at = models.DateTimeField(null=True, blank=True, db_index=True)
pipeline_error = models.TextField(blank=True, default="")
attempts = models.PositiveSmallIntegerField(default=0)
created_at = models.DateTimeField(auto_now_add=True)
updated_at = models.DateTimeField(auto_now=True)
@@ -174,6 +184,61 @@ class TempUpload(models.Model):
return f"{self.original_filename} ({self.status})"
class UploadRun(models.Model):
"""Per-user state of the background upload pipeline.
One row per user acts as the cheap status source the SPA polls and as a
place for batch-level failures (broken e621 credentials, outages) that
would otherwise be repeated on every row.
"""
STATUS_IDLE = "idle"
STATUS_RUNNING = "running"
STATUS_PAUSED = "paused"
STATUS_ERROR = "error"
STATUS_CHOICES = [
(STATUS_IDLE, "Idle"),
(STATUS_RUNNING, "Running"),
(STATUS_PAUSED, "Paused"),
(STATUS_ERROR, "Error"),
]
PHASE_MD5 = "md5"
PHASE_VISUAL = "visual"
PHASE_IQDB = "iqdb"
PHASE_CHOICES = [
("", "None"),
(PHASE_MD5, "e621 MD5"),
(PHASE_VISUAL, "Visual similarity"),
(PHASE_IQDB, "IQDB"),
]
user = models.OneToOneField(
settings.AUTH_USER_MODEL,
on_delete=models.CASCADE,
related_name="upload_run",
)
status = models.CharField(
max_length=20, choices=STATUS_CHOICES, default=STATUS_IDLE
)
phase = models.CharField(
max_length=20, choices=PHASE_CHOICES, blank=True, default=""
)
total = models.IntegerField(default=0)
processed = models.IntegerField(default=0)
matched = models.IntegerField(default=0)
failed = models.IntegerField(default=0)
error = models.TextField(blank=True, default="")
started_at = models.DateTimeField(null=True, blank=True)
updated_at = models.DateTimeField(auto_now=True)
class Meta:
ordering = ["-updated_at"]
def __str__(self):
return f"Upload run for {self.user_id} ({self.status})"
class DownloadTask(models.Model):
"""A background 'Download to Library' job with progress tracking."""
+44
View File
@@ -145,6 +145,11 @@ class TempUploadSerializer(serializers.ModelSerializer):
library_j_id = serializers.SerializerMethodField()
file_url = serializers.SerializerMethodField()
preview_url = serializers.SerializerMethodField()
md5_checked = serializers.SerializerMethodField()
visual_checked = serializers.SerializerMethodField()
iqdb_checked = serializers.SerializerMethodField()
processing = serializers.SerializerMethodField()
similar_count = serializers.SerializerMethodField()
class Meta:
model = TempUpload
@@ -165,6 +170,13 @@ class TempUploadSerializer(serializers.ModelSerializer):
"library_j_id",
"file_url",
"preview_url",
"pipeline_error",
"attempts",
"md5_checked",
"visual_checked",
"iqdb_checked",
"processing",
"similar_count",
"created_at",
"updated_at",
]
@@ -215,6 +227,38 @@ class TempUploadSerializer(serializers.ModelSerializer):
item, user, action, request=self.context.get("request")
)
def get_md5_checked(self, obj):
return obj.e621_checked_at is not None
def get_visual_checked(self, obj):
return obj.visual_checked_at is not None
def get_iqdb_checked(self, obj):
return obj.iqdb_data is not None
def get_processing(self, obj):
return obj.claimed_at is not None
def get_similar_count(self, obj):
return len(obj.iqdb_data or []) + len(obj.visual_matches or [])
class TempUploadListSerializer(TempUploadSerializer):
"""Compact staged-upload row for the board and the status polling.
Drops the heavy post/IQDB payloads (the metadata modal fetches the full
row) while keeping the pipeline flags the board renders per tile.
"""
class Meta(TempUploadSerializer.Meta):
fields = [
field
for field in TempUploadSerializer.Meta.fields
if field
not in {"e621_data", "iqdb_data", "visual_matches", "custom_tags", "custom_notes"}
]
read_only_fields = fields
class DownloadTaskSerializer(serializers.ModelSerializer):
task_id = serializers.UUIDField(source="id", read_only=True)
+6 -1
View File
@@ -136,7 +136,12 @@ class SimilarityCheckViewSet(
self.get_serializer(check).data, status=status.HTTP_201_CREATED
)
@action(detail=True, methods=["get", "head"], permission_classes=[AllowAny])
@action(
detail=True,
methods=["get", "head"],
permission_classes=[AllowAny],
throttle_classes=[],
)
def file(self, request, pk=None):
"""Serve the temp file; accepts a signed URL like staged uploads."""
check = None
+99
View File
@@ -0,0 +1,99 @@
"""The server-side e621 client: batching, retries and IQDB stream handling."""
from pathlib import Path
from tempfile import TemporaryDirectory
from unittest import mock
from django.test import SimpleTestCase
from apps.library import e621
class FakeResponse:
def __init__(self, status_code, payload=None, headers=None):
self.status_code = status_code
self._payload = payload
self.headers = headers or {}
def json(self):
if self._payload is None:
raise ValueError("not json")
return self._payload
class E621ClientTests(SimpleTestCase):
def test_iqdb_search_reopens_the_file_on_retry(self):
bodies = []
def fake_request(method, url, **kwargs):
handle = kwargs["files"]["search[file]"][1]
bodies.append(handle.read())
if len(bodies) == 1:
return FakeResponse(
429, {"success": False, "message": "Throttled"}
)
return FakeResponse(200, [{"post_id": 1, "score": 90.0}])
with TemporaryDirectory() as tmp:
path = Path(tmp) / "x.png"
path.write_bytes(b"image-bytes")
with mock.patch.object(
e621.requests, "request", side_effect=fake_request
), mock.patch.object(e621.time, "sleep"), mock.patch.object(
e621, "_wait_for_slot"
):
result = e621.iqdb_search(None, path)
# Both attempts must send the full body, not the consumed handle.
self.assertEqual(bodies, [b"image-bytes", b"image-bytes"])
self.assertEqual(result, [{"post_id": 1, "score": 90.0}])
def test_check_md5_batch_keys_by_md5(self):
payload = {
"posts": [
{"id": 5, "file": {"md5": "a" * 32}},
{"id": 6, "file": {"md5": "b" * 32}},
]
}
with mock.patch.object(e621, "get", return_value=payload) as getter:
found = e621.check_md5_batch(None, ["A" * 32, "b" * 32])
self.assertEqual(set(found), {"a" * 32, "b" * 32})
params = getter.call_args.kwargs["params"]
self.assertTrue(params["tags"].startswith("md5:"))
self.assertEqual(params["limit"], 2)
def test_auth_errors_are_not_retried(self):
calls = []
def fake_request(*args, **kwargs):
calls.append(1)
return FakeResponse(403, {"error": "nope"})
with mock.patch.object(
e621.requests, "request", side_effect=fake_request
), mock.patch.object(e621, "_wait_for_slot"):
with self.assertRaises(e621.E621AuthError):
e621._request(None, "GET", "/posts.json", require_auth=False)
self.assertEqual(len(calls), 1)
def test_load_shedding_html_raises_rate_limited_after_retries(self):
with mock.patch.object(
e621.requests,
"request",
return_value=FakeResponse(200, None),
), mock.patch.object(e621.time, "sleep"), mock.patch.object(
e621, "_wait_for_slot"
):
with self.assertRaises(e621.E621RateLimited):
e621._request(
None, "GET", "/posts.json", require_auth=False, attempts=2
)
def test_fetch_posts_by_ids_queries_with_id_tag(self):
with mock.patch.object(
e621, "get", return_value={"posts": [{"id": 9}]}
) as getter:
posts = e621.fetch_posts_by_ids(None, [9])
self.assertEqual(posts, [{"id": 9}])
params = getter.call_args.kwargs["params"]
self.assertEqual(params["tags"], "id:9")
+296 -1
View File
@@ -6,17 +6,21 @@ import io
import json
import shutil
import tempfile
from datetime import timedelta
from pathlib import Path
from unittest import mock
from django.contrib.auth import get_user_model
from django.core.files.uploadedfile import SimpleUploadedFile
from django.test import Client, TestCase, override_settings
from django.utils import timezone
from PIL import Image
from rest_framework.authtoken.models import Token
from apps.library.models import MediaItem, TempUpload
from apps.library import upload_pipeline
from apps.library.models import MediaItem, TempUpload, UploadRun
User = get_user_model()
@@ -95,6 +99,7 @@ class TempUploadListTests(TestCase):
self.assertEqual(Client().get("/api/uploads/").status_code, 401)
@override_settings(UPLOAD_PIPELINE_AUTOSTART=False)
class StagedUploadWorkflowTests(TestCase):
@classmethod
def setUpClass(cls):
@@ -400,3 +405,293 @@ class IqdbRecordingTests(TestCase):
self.client, f"/api/uploads/{temp.id}/iqdb/", {"results": "nope"}
)
self.assertEqual(response.status_code, 400)
@override_settings(UPLOAD_PIPELINE_AUTOSTART=False)
class UploadPipelineTests(TestCase):
"""The server-side MD5 -> visual -> IQDB queue and its board API."""
@classmethod
def setUpClass(cls):
super().setUpClass()
cls._tmp = tempfile.mkdtemp(prefix="j621-pipeline-")
cls._watched = Path(cls._tmp) / "library"
cls._watched.mkdir(parents=True, exist_ok=True)
cls._settings = override_settings(
MEDIA_ROOT=cls._tmp, WATCHED_FOLDER=str(cls._watched)
)
cls._settings.enable()
@classmethod
def tearDownClass(cls):
cls._settings.disable()
shutil.rmtree(cls._tmp, ignore_errors=True)
super().tearDownClass()
def setUp(self):
self.uploader = User.objects.create_user(
username="pipe-uploader", password="pipe-pass-123456"
)
self.uploader.role = "uploader"
self.uploader.save(update_fields=["role"])
self.other = User.objects.create_user(
username="pipe-other", password="pipe-pass-123456"
)
self.other.role = "uploader"
self.other.save(update_fields=["role"])
def api_client(self, user):
client = Client()
client.defaults["HTTP_AUTHORIZATION"] = (
f"Token {Token.objects.create(user=user).key}"
)
return client
def stage(self, user=None, label="file", payload=None, filename=None):
user = user or self.uploader
payload = payload or TINY_PNG
name = filename or f"{label}.png"
return TempUpload.objects.create(
user=user,
file=SimpleUploadedFile(name, payload, content_type="image/png"),
original_filename=name,
md5=hashlib.md5(payload + label.encode()).hexdigest(),
size=len(payload),
)
def test_md5_match_auto_imports_the_file(self):
temp = self.stage(label="match")
post = {
"id": 123456,
"rating": "s",
"file": {
"md5": temp.md5,
"url": "https://static1.e621.net/data/m.png",
},
"tags": {"general": ["canine"]},
}
with mock.patch.object(
upload_pipeline.e621,
"check_md5_batch",
return_value={temp.md5: post},
):
upload_pipeline.run_pipeline(self.uploader.id)
temp.refresh_from_db()
self.assertEqual(temp.status, TempUpload.STATUS_COMPLETED)
self.assertEqual(temp.resolution, TempUpload.RESOLUTION_AUTO_MD5)
self.assertEqual(temp.e621_post_id, 123456)
self.assertIsNotNone(temp.library_item_id)
self.assertIsNotNone(temp.e621_checked_at)
self.assertIsNone(temp.claimed_at)
run = UploadRun.objects.get(user=self.uploader)
self.assertEqual(run.status, UploadRun.STATUS_IDLE)
self.assertEqual(run.matched, 1)
self.assertEqual(run.processed, 1)
def test_unmatched_file_runs_every_phase(self):
temp = self.stage(label="nomatch")
raw_iqdb = [
{
"post_id": 777,
"score": 91.0,
"post": {
"id": 777,
"rating": "q",
"md5": "a" * 32,
"score": 5,
"fav_count": 2,
"image_width": 800,
"image_height": 600,
},
}
]
modern = [
{
"id": 777,
"rating": "q",
"fav_count": 4,
"score": {"total": 9},
"preview": {"url": "https://static1.e621.net/data/preview/x.jpg"},
"file": {"md5": "a" * 32, "width": 801, "height": 601},
"tags": {"general": ["canine", "solo"]},
}
]
with mock.patch.object(
upload_pipeline.e621, "check_md5_batch", return_value={}
), mock.patch.object(
upload_pipeline.e621, "iqdb_search", return_value=raw_iqdb
), mock.patch.object(
upload_pipeline.e621, "fetch_posts_by_ids", return_value=modern
):
upload_pipeline.run_pipeline(self.uploader.id)
temp.refresh_from_db()
self.assertIsNotNone(temp.e621_checked_at)
self.assertIsNotNone(temp.visual_checked_at)
self.assertEqual(temp.status, TempUpload.STATUS_VISUAL_MATCH)
self.assertEqual(len(temp.iqdb_data), 1)
candidate = temp.iqdb_data[0]
self.assertEqual(candidate["post_id"], 777)
self.assertEqual(candidate["score_total"], 9)
self.assertEqual(candidate["fav_count"], 4)
self.assertEqual(candidate["width"], 801)
self.assertEqual(
candidate["preview_url"], "https://static1.e621.net/data/preview/x.jpg"
)
self.assertEqual(candidate["tags_preview"], ["canine", "solo"])
run = UploadRun.objects.get(user=self.uploader)
self.assertEqual(run.status, UploadRun.STATUS_IDLE)
self.assertEqual(run.processed, 1)
def test_visual_match_flags_similar_library_items(self):
from apps.library.uploads import complete_temp_upload
seed = self.stage(label="seed")
complete_temp_upload(seed)
self.assertTrue(MediaItem.objects.exists())
temp = self.stage(label="similar")
with mock.patch.object(
upload_pipeline.e621, "check_md5_batch", return_value={}
), mock.patch.object(
upload_pipeline.e621, "iqdb_search", return_value=[]
):
upload_pipeline.run_pipeline(self.uploader.id)
temp.refresh_from_db()
self.assertGreaterEqual(len(temp.visual_matches or []), 1)
self.assertEqual(temp.status, TempUpload.STATUS_VISUAL_MATCH)
def test_iqdb_rate_limit_pauses_the_run(self):
temp = self.stage(label="throttled")
with mock.patch.object(
upload_pipeline.e621, "check_md5_batch", return_value={}
), mock.patch.object(
upload_pipeline.e621,
"iqdb_search",
side_effect=upload_pipeline.e621.E621RateLimited("Throttled"),
):
upload_pipeline.run_pipeline(self.uploader.id)
run = UploadRun.objects.get(user=self.uploader)
self.assertEqual(run.status, UploadRun.STATUS_PAUSED)
self.assertIn("e621", run.error)
temp.refresh_from_db()
self.assertEqual(temp.attempts, 1)
self.assertNotEqual(temp.pipeline_error, "")
self.assertIsNone(temp.iqdb_data)
self.assertIsNone(temp.claimed_at)
client = self.api_client(self.uploader)
response = jpost(client, f"/api/uploads/{temp.id}/retry/", {})
self.assertEqual(response.status_code, 200)
temp.refresh_from_db()
self.assertEqual(temp.attempts, 0)
self.assertEqual(temp.pipeline_error, "")
def test_failed_rows_stop_after_the_attempt_cap(self):
temp = self.stage(label="broken")
for _ in range(upload_pipeline.MAX_ATTEMPTS):
with mock.patch.object(
upload_pipeline.e621, "check_md5_batch", return_value={}
), mock.patch.object(
upload_pipeline.e621,
"iqdb_search",
side_effect=upload_pipeline.e621.E621Error("boom"),
):
upload_pipeline.run_pipeline(self.uploader.id)
temp.refresh_from_db()
self.assertEqual(temp.attempts, upload_pipeline.MAX_ATTEMPTS)
self.assertEqual(upload_pipeline.count_outstanding(self.uploader), 0)
run = UploadRun.objects.get(user=self.uploader)
self.assertGreaterEqual(run.failed, 1)
def test_status_process_and_compact_board_payload(self):
temp = self.stage(label="board")
client = self.api_client(self.uploader)
status = client.get("/api/uploads/status/").json()
self.assertEqual(status["status"], UploadRun.STATUS_IDLE)
self.assertTrue(status["active"])
self.assertEqual(status["outstanding"], 1)
self.assertEqual(status["waiting"]["md5"], 1)
self.assertEqual(client.post("/api/uploads/process/").status_code, 200)
rows = client.get("/api/uploads/").json()
self.assertEqual(len(rows), 1)
row = rows[0]
for key in (
"md5_checked",
"visual_checked",
"iqdb_checked",
"processing",
"similar_count",
"pipeline_error",
):
self.assertIn(key, row)
self.assertNotIn("e621_data", row)
self.assertNotIn("iqdb_data", row)
detail = client.get(f"/api/uploads/{temp.id}/").json()
self.assertIn("e621_data", detail)
self.assertIn("iqdb_data", detail)
def test_discard_bulk_removes_only_own_rows(self):
client = self.api_client(self.uploader)
first = self.stage(label="discard-a")
second = self.stage(label="discard-b")
theirs = self.stage(self.other, label="discard-theirs")
paths = [Path(first.file.path), Path(second.file.path)]
response = jpost(
client,
"/api/uploads/discard-bulk/",
{"temp_ids": [str(first.id), str(second.id), str(theirs.id)]},
)
self.assertEqual(response.status_code, 200)
body = response.json()
self.assertEqual(len(body["discarded"]), 2)
self.assertEqual(len(body["errors"]), 1)
self.assertEqual(body["errors"][0]["error"], "not found")
for path in paths:
self.assertFalse(path.exists())
self.assertFalse(
TempUpload.objects.filter(pk__in=[first.id, second.id]).exists()
)
self.assertTrue(TempUpload.objects.filter(pk=theirs.id).exists())
def test_retry_with_phase_rechecks_iqdb(self):
temp = self.stage(label="recheck")
TempUpload.objects.filter(pk=temp.pk).update(
iqdb_data=[{"post_id": 1}],
status=TempUpload.STATUS_VISUAL_MATCH,
)
client = self.api_client(self.uploader)
response = jpost(
client, f"/api/uploads/{temp.id}/retry/", {"phase": "iqdb"}
)
self.assertEqual(response.status_code, 200)
temp.refresh_from_db()
self.assertIsNone(temp.iqdb_data)
self.assertEqual(temp.status, TempUpload.STATUS_VISUAL_MATCH)
def test_stale_claims_are_released(self):
temp = self.stage(label="stale")
TempUpload.objects.filter(pk=temp.pk).update(
claimed_at=timezone.now() - upload_pipeline.STALE_CLAIM_AFTER
- timedelta(minutes=1)
)
UploadRun.objects.create(user=self.uploader, status=UploadRun.STATUS_RUNNING)
UploadRun.objects.filter(user=self.uploader).update(
updated_at=timezone.now() - upload_pipeline.STALE_RUN_AFTER
- timedelta(minutes=1)
)
released, paused = upload_pipeline.reap_stale_claims()
self.assertEqual(released, 1)
self.assertEqual(paused, 1)
temp.refresh_from_db()
self.assertIsNone(temp.claimed_at)
run = UploadRun.objects.get(user=self.uploader)
self.assertEqual(run.status, UploadRun.STATUS_PAUSED)
+589
View File
@@ -0,0 +1,589 @@
"""Background processing for staged uploads.
Each staged upload runs through three phases, in batches of 75 (the same
lookup size the original app used for its e621 MD5 cache command):
1. e621 MD5 lookup — one ``posts.json`` query per round; byte-identical
matches are imported straight into the library (``auto_md5``).
2. Local visual similarity — perceptual hashes are compared against the
whole library once per round.
3. e621 IQDB — reverse-image search for whatever is still unresolved.
The pipeline runs in a daemon thread started on demand (like the download and
match scans), so the browser can navigate away and the work keeps going.
Progress and pause/error state live in the ``UploadRun`` row; per-file state
lives on ``TempUpload``. Rows are claimed with ``SELECT ... FOR UPDATE SKIP
LOCKED`` so several gunicorn workers cannot process the same file, and stale
claims left by a recycled worker are reaped and picked up again.
"""
import logging
import threading
import time
from datetime import timedelta
from django.conf import settings
from django.contrib.auth import get_user_model
from django.db import connection, transaction
from django.db.models import F, Q
from django.utils import timezone
from . import e621, services
from .models import TempUpload, UploadRun
logger = logging.getLogger(__name__)
# Same batch size as the original app's e621 cache command.
MD5_BATCH_SIZE = 75
# Rows one worker round claims; also the MD5 query size.
CLAIM_SIZE = MD5_BATCH_SIZE
# A claimed row is assumed dead after this long and is queued again.
STALE_CLAIM_AFTER = timedelta(minutes=15)
# A run whose heartbeat stopped this long ago can be taken over.
STALE_RUN_AFTER = timedelta(minutes=15)
# Per-row failures before the pipeline stops retrying automatically.
MAX_ATTEMPTS = 3
# Round-level e621 retries before the run is paused.
ROUND_ATTEMPTS = 3
ROUND_RETRY_SECONDS = 20
WORK_STATUSES = (TempUpload.STATUS_PENDING, TempUpload.STATUS_VISUAL_MATCH)
VIDEO_RE = r"\.(mp4|webm)$"
# A staged upload still needs work when any phase has not run yet. IQDB is
# skipped for videos, which never get iqdb_data, so they must not stay
# "outstanding" forever.
OUTSTANDING_Q = (
Q(e621_checked_at__isnull=True)
| Q(visual_checked_at__isnull=True)
| (Q(iqdb_data__isnull=True) & ~Q(original_filename__iregex=VIDEO_RE))
)
_running_lock = threading.Lock()
_running_users: set[int] = set()
class PipelinePaused(Exception):
"""A round-level failure that should pause the run instead of failing rows."""
def is_video(filename):
return bool(filename) and filename.lower().endswith((".mp4", ".webm"))
def is_finished(temp):
"""True when every phase this file needs has run."""
if temp.status == TempUpload.STATUS_COMPLETED:
return True
if temp.e621_checked_at is None or temp.visual_checked_at is None:
return False
return is_video(temp.original_filename) or temp.iqdb_data is not None
def outstanding_queryset(user):
return (
TempUpload.objects.filter(user=user, status__in=WORK_STATUSES)
.filter(OUTSTANDING_Q)
.filter(attempts__lt=MAX_ATTEMPTS)
)
def count_outstanding(user):
return outstanding_queryset(user).count()
def count_failed(user):
return (
TempUpload.objects.filter(user=user, status__in=WORK_STATUSES)
.filter(attempts__gte=MAX_ATTEMPTS)
.count()
)
def waiting_counts(user):
"""How many files are left per phase (phases overlap by design)."""
base = TempUpload.objects.filter(
user=user, status__in=WORK_STATUSES, attempts__lt=MAX_ATTEMPTS
)
return {
"md5": base.filter(e621_checked_at__isnull=True).count(),
"visual": base.filter(visual_checked_at__isnull=True).count(),
"iqdb": (
base.filter(iqdb_data__isnull=True)
.exclude(original_filename__iregex=VIDEO_RE)
.count()
),
}
def status_payload(user):
"""Cheap state for the shell/upload page to poll."""
run = UploadRun.objects.filter(user=user).first()
outstanding = count_outstanding(user)
failed = count_failed(user)
status = run.status if run is not None else UploadRun.STATUS_IDLE
total = run.total if run is not None else 0
processed = run.processed if run is not None else 0
# A paused run with nothing left to do is not "active" (the user may have
# resolved or discarded the failed rows); failed rows stay visible until
# they are retried or dismissed.
active = (
outstanding > 0 or failed > 0 or status == UploadRun.STATUS_RUNNING
)
return {
"status": status,
"active": active,
"phase": run.phase if run is not None else "",
"total": max(total, processed + failed),
"processed": processed,
"matched": run.matched if run is not None else 0,
"failed": max(failed, run.failed if run is not None else 0),
"error": run.error if run is not None else "",
"outstanding": outstanding,
"waiting": waiting_counts(user),
"updated_at": run.updated_at.isoformat() if run is not None else None,
}
def reap_stale_claims():
"""Queue rows left claimed by a recycled worker and pause dead runs."""
cutoff = timezone.now() - STALE_CLAIM_AFTER
released = TempUpload.objects.filter(claimed_at__lt=cutoff).update(
claimed_at=None
)
paused = UploadRun.objects.filter(
status=UploadRun.STATUS_RUNNING, updated_at__lt=cutoff
).update(
status=UploadRun.STATUS_PAUSED,
phase="",
error="The worker stopped before finishing. Retry to resume.",
updated_at=timezone.now(),
)
if released or paused:
logger.info(
"Reaped %s stale upload claims and %s dead upload runs", released, paused
)
return released, paused
def start_pipeline(user):
"""Start the pipeline for one user in a daemon thread (idempotent)."""
if not getattr(settings, "UPLOAD_PIPELINE_AUTOSTART", True):
return False
if user is None or not getattr(user, "can_upload", False):
return False
user_id = int(user.pk)
with _running_lock:
if user_id in _running_users:
return False
_running_users.add(user_id)
reap_stale_claims()
thread = threading.Thread(target=_thread_entry, args=(user_id,), daemon=True)
thread.start()
return True
def _thread_entry(user_id):
try:
run_pipeline(user_id)
except Exception: # noqa: BLE001 - a thread must never crash the worker
logger.exception("Upload pipeline for user %s crashed", user_id)
finally:
with _running_lock:
_running_users.discard(user_id)
connection.close()
def run_pipeline(user_id):
user = get_user_model().objects.filter(pk=user_id).first()
if user is None or not user.can_upload:
return
now = timezone.now()
# Claim the run row so two gunicorn workers cannot own the same queue.
with transaction.atomic():
run, _ = UploadRun.objects.select_for_update().get_or_create(user=user)
if (
run.status == UploadRun.STATUS_RUNNING
and run.updated_at is not None
and run.updated_at > now - STALE_RUN_AFTER
):
# Another worker owns this run.
return
# Reset the counters when a new queue starts cleanly; otherwise keep
# accumulating so failed rows from an earlier pass stay visible.
live_failed = count_failed(user)
outstanding = count_outstanding(user)
if run.status == UploadRun.STATUS_IDLE and live_failed == 0:
run.total = outstanding
run.processed = 0
run.matched = 0
run.failed = 0
else:
run.total = max(
run.total or 0, run.processed + run.failed + outstanding
)
run.failed = max(run.failed, live_failed)
run.status = UploadRun.STATUS_RUNNING
run.phase = ""
run.error = ""
run.started_at = now
run.save()
try:
while True:
rows = claim_round(user)
if not rows:
break
process_round(run, user, rows)
except PipelinePaused as exc:
run.status = UploadRun.STATUS_PAUSED
run.phase = ""
run.error = str(exc)
run.save()
except Exception as exc: # noqa: BLE001 - surface crashes as a run error
logger.exception("Upload pipeline for user %s failed", user_id)
run.status = UploadRun.STATUS_ERROR
run.phase = ""
run.error = f"The upload pipeline stopped: {exc}"
run.save()
else:
run.status = UploadRun.STATUS_IDLE
run.phase = ""
run.save()
def claim_round(user, size=CLAIM_SIZE):
"""Claim up to ``size`` outstanding rows for this worker."""
now = timezone.now()
with transaction.atomic():
rows = list(
TempUpload.objects.select_for_update(skip_locked=True)
.filter(user=user, status__in=WORK_STATUSES)
.filter(OUTSTANDING_Q)
.filter(claimed_at__isnull=True, attempts__lt=MAX_ATTEMPTS)
.order_by("created_at")[:size]
)
if rows:
TempUpload.objects.filter(pk__in=[row.pk for row in rows]).update(
claimed_at=now
)
for row in rows:
row.claimed_at = now
return rows
def process_round(run, user, rows):
"""Run every phase for one claimed round, then release/account the rows."""
ids = [row.pk for row in rows]
matched = 0
try:
matched += md5_phase(run, user, rows)
rows = refresh(ids)
visual_phase(run, user, rows)
rows = refresh(ids)
iqdb_phase(run, user, rows)
finally:
finalize_round(run, ids, matched)
def refresh(ids):
return list(TempUpload.objects.filter(pk__in=ids))
def _save_run(run, **fields):
for key, value in fields.items():
setattr(run, key, value)
run.save(
update_fields=[*fields.keys(), "updated_at"]
)
def md5_phase(run, user, rows):
"""One e621 MD5 batch query; matches are imported into the library."""
targets = [row for row in rows if row.e621_checked_at is None]
if not targets:
return 0
_save_run(
run,
phase=UploadRun.PHASE_MD5,
total=run.processed + run.failed + count_outstanding(user),
)
posts = _e621_round(
lambda: e621.check_md5_batch(user, [row.md5 for row in targets])
)
by_md5 = {}
for post in posts.values():
file_data = post.get("file") or {}
md5 = str(file_data.get("md5") or "").strip().lower()
if md5:
by_md5[md5] = post
matched = 0
now = timezone.now()
for row in targets:
post = by_md5.get(str(row.md5).strip().lower())
if post is None:
TempUpload.objects.filter(pk=row.pk).update(
e621_checked_at=now, pipeline_error="", updated_at=now
)
continue
try:
trimmed = services.trim_e621_post(post)
if trimmed is None or not trimmed.get("id"):
raise e621.E621Error("e621 returned an unexpected post payload.")
TempUpload.objects.filter(pk=row.pk).update(
e621_post_id=int(trimmed["id"]),
e621_data=trimmed,
resolution=TempUpload.RESOLUTION_AUTO_MD5,
e621_checked_at=now,
pipeline_error="",
updated_at=now,
)
row.refresh_from_db()
from .uploads import complete_temp_upload
complete_temp_upload(row)
matched += 1
except Exception as exc: # noqa: BLE001 - keep going for other files
logger.exception("Could not auto-import staged upload %s", row.pk)
record_failure(row, f"Could not finish the upload: {exc}")
return matched
def visual_phase(run, user, rows):
"""Compare each row's perceptual hashes against the library once."""
from .uploads import build_hash_index, match_hashes
targets = [
row
for row in rows
if row.visual_checked_at is None
and row.status in WORK_STATUSES
and row.file
]
if not targets:
return
_save_run(run, phase=UploadRun.PHASE_VISUAL)
index = build_hash_index()
now = timezone.now()
for row in targets:
try:
hashes = services.compute_visual_hashes(row.file.path)
if not hashes:
TempUpload.objects.filter(pk=row.pk).update(
visual_checked_at=now, pipeline_error="", updated_at=now
)
continue
matches = match_hashes(hashes, index, user=user)
update = {
"visual_matches": matches,
"visual_checked_at": now,
"pipeline_error": "",
"updated_at": now,
}
if matches and row.status == TempUpload.STATUS_PENDING:
update["status"] = TempUpload.STATUS_VISUAL_MATCH
TempUpload.objects.filter(pk=row.pk).update(**update)
except Exception as exc: # noqa: BLE001 - keep going for other files
logger.exception("Visual similarity failed for %s", row.pk)
record_failure(row, f"Visual similarity failed: {exc}")
def iqdb_phase(run, user, rows):
"""Reverse-image search every unresolved image, one e621 query each."""
targets = [
row
for row in rows
if row.iqdb_data is None
and row.status in WORK_STATUSES
and row.file
and not is_video(row.original_filename)
]
if not targets:
return
_save_run(run, phase=UploadRun.PHASE_IQDB)
heartbeat_at = time.monotonic()
for row in targets:
# Keep the run row fresh: a 75-file IQDB round takes minutes and must
# not look like a dead worker to another request.
if time.monotonic() - heartbeat_at > 30:
UploadRun.objects.filter(pk=run.pk).update(updated_at=timezone.now())
heartbeat_at = time.monotonic()
try:
raw = e621.iqdb_search(user, row.file.path)
results = normalize_iqdb_results(user, raw)
except e621.E621AuthError as exc:
record_failure(row, str(exc))
raise PipelinePaused(
"e621 rejected the credentials — fix them in Account and retry."
) from exc
except e621.E621RateLimited as exc:
record_failure(row, str(exc))
raise PipelinePaused(
"e621 is throttling IQDB right now; the queue will resume."
) from exc
except Exception as exc: # noqa: BLE001 - keep going for other files
logger.exception("IQDB search failed for %s", row.pk)
record_failure(row, f"IQDB search failed: {exc}")
continue
now = timezone.now()
update = {
"iqdb_data": results,
"pipeline_error": "",
"updated_at": now,
}
if results and row.status == TempUpload.STATUS_PENDING:
update["status"] = TempUpload.STATUS_VISUAL_MATCH
TempUpload.objects.filter(pk=row.pk).update(**update)
def record_failure(row, message):
"""Count one failed attempt against a row and queue it for a retry."""
TempUpload.objects.filter(pk=row.pk).update(
attempts=F("attempts") + 1,
pipeline_error=str(message)[:2000],
claimed_at=None,
updated_at=timezone.now(),
)
def finalize_round(run, ids, matched):
"""Account finished/failed rows and release the rest of the claims."""
rows = refresh(ids)
finished = {row.pk for row in rows if is_finished(row)}
failed = {row.pk for row in rows if row.attempts >= MAX_ATTEMPTS}
# Release every claim: finished rows must not keep looking "processing"
# to the board, and unfinished rows are re-queued for the next run.
release = [row.pk for row in rows if row.claimed_at is not None]
if release:
TempUpload.objects.filter(pk__in=release).update(claimed_at=None)
_save_run(
run,
processed=run.processed + len(finished),
failed=run.failed + len(failed - finished),
matched=run.matched + matched,
)
def _e621_round(task):
"""Run a round-level e621 call, retrying through rate limits."""
last_error = None
for attempt in range(ROUND_ATTEMPTS):
try:
return task()
except e621.E621AuthError:
raise
except (e621.E621RateLimited, e621.E621Error) as exc:
last_error = exc
if attempt + 1 >= ROUND_ATTEMPTS:
break
delay = ROUND_RETRY_SECONDS * (attempt + 1)
logger.info("e621 round failed (%s); retrying in %ss", exc, delay)
time.sleep(delay)
raise PipelinePaused(
f"e621 is not answering right now ({last_error}); the queue will resume."
) from last_error
def flatten_tag_preview(tags, limit=8):
"""First few tag names from a modern post payload, like the SPA shows."""
if not isinstance(tags, dict):
return []
out = []
for values in tags.values():
if not isinstance(values, list):
continue
for tag in values:
if isinstance(tag, str) and tag not in out:
out.append(tag)
if len(out) >= limit:
return out
return out
def _legacy_iqdb_post(entry):
"""Unwrap the post payload embedded in a legacy IQDB match."""
post = entry.get("post")
if not isinstance(post, dict):
return {}
inner = post.get("posts")
return inner if isinstance(inner, dict) else post
def normalize_iqdb_results(user, raw_results):
"""Shape legacy IQDB matches like the SPA's E621IqdbCandidate entries.
The IQDB payload carries little post data, so candidates are enriched
with one batched ``id:`` lookup before they are stored.
"""
candidates = []
for entry in (raw_results or [])[:10]:
if not isinstance(entry, dict):
continue
post = _legacy_iqdb_post(entry)
post_id = entry.get("post_id")
if not isinstance(post_id, int):
post_id = post.get("id")
score = entry.get("score")
candidates.append(
{
"post_id": post_id if isinstance(post_id, int) else None,
"score": float(score) if isinstance(score, (int, float)) else None,
"preview_url": None,
"rating": (
post.get("rating") if isinstance(post.get("rating"), str) else None
),
"md5": post.get("md5") if isinstance(post.get("md5"), str) else None,
"score_total": (
post.get("score") if isinstance(post.get("score"), int) else None
),
"fav_count": (
post.get("fav_count")
if isinstance(post.get("fav_count"), int)
else None
),
"width": (
post.get("image_width")
if isinstance(post.get("image_width"), int)
else None
),
"height": (
post.get("image_height")
if isinstance(post.get("image_height"), int)
else None
),
"tags_preview": [],
}
)
ids = [entry["post_id"] for entry in candidates if entry["post_id"]]
if not ids:
return services.sanitize_iqdb_results(candidates)
try:
posts = e621.fetch_posts_by_ids(user, ids)
except e621.E621Error as exc:
# Candidates without enrichment still show up; keep them.
logger.info("Could not enrich IQDB candidates: %s", exc)
posts = []
by_id = {post.get("id"): post for post in posts}
for entry in candidates:
post = by_id.get(entry["post_id"])
if not isinstance(post, dict):
continue
file_data = post.get("file") or {}
preview = post.get("preview") or {}
score = post.get("score") or {}
entry["preview_url"] = preview.get("url") or entry["preview_url"]
entry["rating"] = post.get("rating") or entry["rating"]
entry["md5"] = file_data.get("md5") or entry["md5"]
if isinstance(score, dict):
entry["score_total"] = score.get("total")
entry["fav_count"] = post.get("fav_count")
entry["width"] = file_data.get("width")
entry["height"] = file_data.get("height")
entry["tags_preview"] = flatten_tag_preview(post.get("tags"))
return services.sanitize_iqdb_results(candidates)
+188 -23
View File
@@ -29,27 +29,32 @@ from rest_framework.response import Response
from . import services
from .models import MediaItem, TempUpload
from .permissions import CanUpload
from .serializers import TempUploadSerializer
from .tools import HASH_FIELDS, hashes_similarity
from .serializers import TempUploadListSerializer, TempUploadSerializer
from .tools import HASH_FIELDS, hashed_items, hashes_similarity
logger = logging.getLogger(__name__)
def find_library_matches(path, limit=10, user=None, request=None):
"""Library items visually similar to a staged file."""
hashes = services.compute_visual_hashes(path)
if not hashes:
return []
def build_hash_index():
"""Library hash mappings for similarity scans, loaded once per batch.
Only items that actually carry perceptual hashes are included; the old
per-file scan walked every row (including videos and unchecked items).
"""
algorithms = list(HASH_FIELDS)
return [
(item, {field: getattr(item, field, "") for field in algorithms})
for item in hashed_items(algorithms)
]
def match_hashes(hashes, index, limit=10, user=None, request=None):
"""Library items whose perceptual hashes are close to ``hashes``."""
algorithms = list(HASH_FIELDS)
threshold = settings.VISUAL_MATCH_THRESHOLD
matches = []
for item in MediaItem.objects.prefetch_related("locations"):
similarity = hashes_similarity(
hashes,
{field: getattr(item, field, "") for field in algorithms},
algorithms,
threshold,
)
for item, item_hashes in index:
similarity = hashes_similarity(hashes, item_hashes, algorithms, threshold)
if similarity is None:
continue
location = item.locations.first()
@@ -67,6 +72,20 @@ def find_library_matches(path, limit=10, user=None, request=None):
return matches[:limit]
def find_library_matches(path, limit=10, user=None, request=None, index=None):
"""Library items visually similar to a staged file.
Pass a prebuilt ``index`` (see build_hash_index) to reuse it across a
whole batch instead of rescanning the library per file.
"""
hashes = services.compute_visual_hashes(path)
if not hashes:
return []
if index is None:
index = build_hash_index()
return match_hashes(hashes, index, limit=limit, user=user, request=request)
def complete_temp_upload(temp, download_url=None):
"""Index the upload into the library.
@@ -156,12 +175,20 @@ class TempUploadViewSet(
pagination_class = None
http_method_names = ["get", "post", "delete", "head", "options"]
def get_serializer_class(self):
# The board polls the list, so its payload stays small; the metadata
# modal fetches the full row from the detail endpoint.
if self.action == "list":
return TempUploadListSerializer
return TempUploadSerializer
def get_queryset(self):
queryset = TempUpload.objects.select_related("library_item")
user = self.request.user
if not user.is_app_staff:
queryset = queryset.filter(user=user)
return queryset
# Staged uploads are private: everyone, staff included, only sees
# their own board. (The file action still lets staff read bytes by id
# for support purposes.)
return TempUpload.objects.select_related("library_item").filter(
user=self.request.user
)
def create(self, request):
upload = request.FILES.get("file")
@@ -190,10 +217,14 @@ class TempUploadViewSet(
temp.resolution = TempUpload.RESOLUTION_DUPLICATE
temp.library_item = existing
temp.file.delete(save=False)
# Visual similarity is deliberately a separate phase (the
# /visual-match action) so a large batch uploads at full speed and
# the board runs MD5 -> visual -> IQDB over the whole batch.
# Visual similarity and IQDB run in the background pipeline so a large
# batch uploads at full speed and the work survives the browser.
temp.save()
# Kick the server-side pipeline; staging no longer waits on e621 and
# the work continues even if the browser navigates away.
from .upload_pipeline import start_pipeline
start_pipeline(request.user)
return Response(
self.get_serializer(temp).data, status=status.HTTP_201_CREATED
)
@@ -221,7 +252,141 @@ class TempUploadViewSet(
temp.save(update_fields=["visual_matches", "status", "updated_at"])
return Response(self.get_serializer(temp).data)
@action(detail=True, methods=["get", "head"], permission_classes=[AllowAny])
@action(detail=False, methods=["get"])
def status(self, request):
"""Cheap pipeline state for the shell indicator and the upload page."""
from .upload_pipeline import status_payload
return Response(status_payload(request.user))
@action(detail=False, methods=["post"])
def process(self, request):
"""Start (or resume) the pipeline for the caller's staged uploads.
Idempotent: the client calls this after staging files, on page load
and when a paused run should be retried.
"""
from .upload_pipeline import start_pipeline, status_payload
start_pipeline(request.user)
return Response(status_payload(request.user))
@action(detail=True, methods=["post"])
def retry(self, request, pk=None):
"""Queue one staged upload for another pipeline pass."""
from .upload_pipeline import start_pipeline, status_payload
temp = self.get_object()
if temp.status == TempUpload.STATUS_COMPLETED:
return Response(
{"detail": "This upload is already in the library."},
status=status.HTTP_400_BAD_REQUEST,
)
if not temp.file:
return Response(
{"detail": "The staged file is missing."},
status=status.HTTP_400_BAD_REQUEST,
)
update = {
"claimed_at": None,
"attempts": 0,
"pipeline_error": "",
"updated_at": timezone.now(),
}
# An explicit phase re-runs that one check even if it already ran.
phase = str(request.data.get("phase") or "").strip()
if phase == "md5":
update["e621_checked_at"] = None
elif phase == "visual":
update["visual_matches"] = None
update["visual_checked_at"] = None
elif phase == "iqdb":
update["iqdb_data"] = None
if temp.status == TempUpload.STATUS_ERROR:
# A failed import has to go through the MD5 phase again so the
# completion is retried; other errors only re-run missing phases.
update["status"] = TempUpload.STATUS_PENDING
update["e621_checked_at"] = None
elif temp.status not in (
TempUpload.STATUS_PENDING,
TempUpload.STATUS_VISUAL_MATCH,
):
update["status"] = TempUpload.STATUS_PENDING
TempUpload.objects.filter(pk=temp.pk).update(**update)
start_pipeline(request.user)
return Response(status_payload(request.user))
@action(detail=False, methods=["post"], url_path="retry-all")
def retry_all(self, request):
"""Queue every retryable staged upload for another pipeline pass."""
from .upload_pipeline import start_pipeline, status_payload
now = timezone.now()
retryable = self.get_queryset().filter(
status__in=[
TempUpload.STATUS_PENDING,
TempUpload.STATUS_VISUAL_MATCH,
TempUpload.STATUS_ERROR,
]
)
retryable.exclude(file="").update(
claimed_at=None,
attempts=0,
pipeline_error="",
updated_at=now,
)
retryable.exclude(file="").filter(status=TempUpload.STATUS_ERROR).update(
status=TempUpload.STATUS_PENDING,
e621_checked_at=None,
)
start_pipeline(request.user)
return Response(status_payload(request.user))
@action(detail=False, methods=["post"], url_path="discard-bulk")
def discard_bulk(self, request):
"""Discard many staged uploads in one request (the board's "all")."""
ids = request.data.get("temp_ids")
if not isinstance(ids, list) or not ids:
return Response(
{"detail": "temp_ids must be a non-empty list."},
status=status.HTTP_400_BAD_REQUEST,
)
if len(ids) > 1000:
return Response(
{"detail": "Too many ids in one request (max 1000)."},
status=status.HTTP_400_BAD_REQUEST,
)
values = list(dict.fromkeys(str(value) for value in ids))
try:
queryset = self.get_queryset().filter(pk__in=values)
except (ValidationError, ValueError):
return Response(
{"detail": "One or more temp_ids are not valid upload ids."},
status=status.HTTP_400_BAD_REQUEST,
)
by_id = {str(temp.pk): temp for temp in queryset}
discarded: list[str] = []
errors: list[dict[str, str]] = []
for value in values:
temp = by_id.get(value)
if temp is None:
errors.append({"temp_id": value, "error": "not found"})
continue
try:
self.perform_destroy(temp)
discarded.append(value)
except Exception as exc: # noqa: BLE001 - report per-file failures
logger.exception("Could not discard staged upload %s", value)
errors.append({"temp_id": value, "error": str(exc)})
return Response({"discarded": discarded, "errors": errors})
@action(
detail=True,
methods=["get", "head"],
permission_classes=[AllowAny],
throttle_classes=[],
)
def file(self, request, pk=None):
"""Serve the staged file; accepts a signed URL for media tags."""
user = request.user if request.user.is_authenticated else None
+2 -2
View File
@@ -148,7 +148,7 @@ class MediaItemViewSet(
return item
return self.get_object()
@action(detail=True, methods=["get"])
@action(detail=True, methods=["get"], throttle_classes=[])
def raw(self, request, pk=None):
item = self._media_object(request, "raw")
location = item.locations.first()
@@ -161,7 +161,7 @@ class MediaItemViewSet(
request, location.path, download=request.query_params.get("download") == "1"
)
@action(detail=True, methods=["get"])
@action(detail=True, methods=["get"], throttle_classes=[])
def thumbnail(self, request, pk=None):
item = self._media_object(request, "thumbnail")
location = item.locations.first()
+27 -4
View File
@@ -234,6 +234,12 @@ GUEST_BLACKLIST_TTL = int(os.getenv("GUEST_BLACKLIST_TTL", "3600"))
# Similarity threshold for flagging staged uploads that match library items.
VISUAL_MATCH_THRESHOLD = float(os.getenv("VISUAL_MATCH_THRESHOLD", "0.9"))
# Start the staged-upload pipeline when a file is staged (daemon thread in the
# worker). Tests turn this off and drive the pipeline synchronously.
UPLOAD_PIPELINE_AUTOSTART = os.getenv(
"UPLOAD_PIPELINE_AUTOSTART", "true"
).strip().lower() not in {"0", "false", "no", "off"}
# Ephemeral similarity-check uploads are deleted after this many minutes
# (and always on startup).
SIMILARITY_TTL_MINUTES = int(os.getenv("SIMILARITY_TTL_MINUTES", "30"))
@@ -242,13 +248,23 @@ SIMILARITY_TTL_MINUTES = int(os.getenv("SIMILARITY_TTL_MINUTES", "30"))
# workers and management commands (e.g. the mirrored guest blacklist).
CACHES = {
"default": {
"BACKEND": "django.core.cache.backends.redis.RedisCache",
"BACKEND": "apps.core.cache.ResilientRedisCache",
"LOCATION": os.getenv("REDIS_URL", "redis://127.0.0.1:6380/1"),
}
}
# Django REST Framework
# Private / tailnet-only deployments can drop the general anon+user limits
# entirely (THROTTLE_ENABLED=false). The scoped guards below (login, register,
# e621 proxy) and the media endpoints' own protections stay active either way.
THROTTLE_ENABLED = os.getenv("THROTTLE_ENABLED", "true").strip().lower() not in {
"0",
"false",
"no",
"off",
}
REST_FRAMEWORK = {
"DEFAULT_AUTHENTICATION_CLASSES": [
"rest_framework.authentication.TokenAuthentication",
@@ -263,11 +279,18 @@ REST_FRAMEWORK = {
],
"DEFAULT_PAGINATION_CLASS": "config.pagination.StandardPagination",
"PAGE_SIZE": 48,
# Per-IP/per-user rate limits (counted in the shared Redis cache).
"DEFAULT_THROTTLE_CLASSES": [
# Per-IP/per-user rate limits (counted in the shared Redis cache). Signed
# media URLs are deliberately excluded at the view level: <img>/<video>
# tags fetch them without an Authorization header, so a library page would
# otherwise burn the anonymous bucket and start returning JSON 429s.
"DEFAULT_THROTTLE_CLASSES": (
[
"rest_framework.throttling.AnonRateThrottle",
"rest_framework.throttling.UserRateThrottle",
],
]
if THROTTLE_ENABLED
else []
),
"DEFAULT_THROTTLE_RATES": {
# Generous enough for the shell polling (status every 5s, stats every 2s).
"anon": os.getenv("THROTTLE_ANON", "120/min"),
+4
View File
@@ -64,6 +64,10 @@ DB_ROOT_PASSWORD=j621root
# THROTTLE_LOGIN=5/min
# THROTTLE_REGISTER=20/hour
# THROTTLE_E621_PROXY=60/hour
# Tailnet-only / private deployments can drop the general limits entirely.
# Signed media URLs (<img>/<video>) and the login/register/proxy guards are
# exempt from this switch either way.
# THROTTLE_ENABLED=false
# e621 media hosts the backend may fetch from (downloads, proxies)
# E621_MEDIA_HOSTS=static1.e621.net,static2.e621.net,static3.e621.net
+10
View File
@@ -164,6 +164,16 @@ does not. The desktop app's "Check for updates…" menu item reads
`latest-linux.yml` / `latest.yml` from there (see `desktop/README.md`).
Backend-only composes have no frontend, so no feed.
The manual **CD** workflow (Actions tab) builds the desktop packages on the
runner and attaches them plus the update metadata to the Gitea release
`desktop-v<version>`; it does not touch the live feed, because that is
runtime state on the deploy host and CI has no SSH key for it. After a CD run,
publish the feed from a machine that can reach the deploy checkout:
```bash
./push_desktop.sh --no-build --local # or without --local to also copy it
```
## Scheduled jobs
Compose files with a backend also run a **`scheduler`** service — the same
+18 -5
View File
@@ -22,6 +22,8 @@ case "$TARGET" in
;;
esac
VERSION="$(node -p "require('./desktop/package.json').version")"
if [ ! -d desktop/node_modules ]; then
echo "==> Installing desktop dependencies ..."
npm --prefix desktop ci --no-audit --no-fund
@@ -32,18 +34,29 @@ if [ "$TARGET" != "linux" ] && ! command -v wine >/dev/null 2>&1; then
exit 1
fi
# release/ should only ever hold the current version: old installers are
# rebuilt from git when needed, and the unpacked trees are regenerated.
shopt -s nullglob
stale=(desktop/release/*.deb desktop/release/*.pkg.tar.zst desktop/release/"J621 Setup "*.exe
desktop/release/*.blockmap desktop/release/latest*.yml)
for file in "${stale[@]}"; do
case "$(basename "$file")" in
*"$VERSION"*) ;;
*) echo " removing $(basename "$file")"; rm -f "$file" ;;
esac
done
rm -rf desktop/release/linux-unpacked desktop/release/win-unpacked
case "$TARGET" in
linux) npm --prefix desktop run dist:linux ;;
win) npm --prefix desktop run dist:win ;;
all) npm --prefix desktop run dist:all ;;
esac
VERSION="$(node -p "require('./desktop/package.json').version")"
case "$TARGET" in
linux) ARTIFACTS=(-name "*.deb" -o -name "*.pkg.tar.zst") ;;
win) ARTIFACTS=(-name "*.exe") ;;
all) ARTIFACTS=(-name "*.deb" -o -name "*.pkg.tar.zst" -o -name "*.exe") ;;
linux) ARTIFACTS=(-name "*${VERSION}*.deb" -o -name "*${VERSION}*.pkg.tar.zst") ;;
win) ARTIFACTS=(-name "*${VERSION}*.exe") ;;
all) ARTIFACTS=(-name "*${VERSION}*.deb" -o -name "*${VERSION}*.pkg.tar.zst" -o -name "*${VERSION}*.exe") ;;
esac
echo
+8 -2
View File
@@ -10,10 +10,16 @@ REGISTRY="gitea.rainbow-herring.ts.net/jakebreath/j621-backend"
REGISTRY_HOST="$(printf '%s' "$REGISTRY" | cut -d/ -f1)"
SHA="${1:-$(git rev-parse --short HEAD)}"
BUILDER=multiarch
PLATFORMS="linux/amd64,linux/arm64"
PLATFORMS="${PLATFORMS:-linux/amd64,linux/arm64}"
echo "==> Logging in to $REGISTRY_HOST ..."
docker login "$REGISTRY_HOST"
if [ -n "${REGISTRY_USER:-}" ] && [ -n "${REGISTRY_TOKEN:-}" ]; then
# Non-interactive login for CI (workflow passes GITHUB_TOKEN).
printf '%s' "$REGISTRY_TOKEN" | docker login "$REGISTRY_HOST" \
-u "$REGISTRY_USER" --password-stdin
else
docker login "$REGISTRY_HOST"
fi
if ! docker buildx inspect "$BUILDER" >/dev/null 2>&1; then
echo "==> Creating buildx builder '$BUILDER' ..."
+25 -10
View File
@@ -54,20 +54,29 @@ FEED=deploy/data/desktop
mkdir -p "$FEED"
shopt -s nullglob
deb=(desktop/release/*.deb)
zst=(desktop/release/*.pkg.tar.zst)
VERSION="$(node -p "require('./desktop/package.json').version")"
deb=(desktop/release/*"$VERSION"*.deb)
zst=(desktop/release/*"$VERSION"*.pkg.tar.zst)
linux_meta=(desktop/release/latest-linux.yml)
win_exe=("desktop/release/J621 Setup "*.exe)
win_exe=("desktop/release/J621 Setup ${VERSION}"*.exe)
win_meta=(desktop/release/latest.yml)
blockmaps=(desktop/release/*.blockmap)
blockmaps=(desktop/release/*"$VERSION"*.exe.blockmap)
if [ "${#deb[@]}" -eq 0 ] && [ "${#zst[@]}" -eq 0 ]; then
echo "No artifacts in desktop/release/ — run without --no-build first." >&2
echo "No $VERSION artifacts in desktop/release/ — run without --no-build first." >&2
exit 1
fi
# Replace the previous release's metadata before copying the new artifacts.
# Replace the previous release's metadata and drop older installers from the
# feed: latest*.yml only ever points at the current version.
rm -f "$FEED"/latest-linux.yml "$FEED"/latest.yml
old=("$FEED"/*.deb "$FEED"/*.pkg.tar.zst "$FEED"/"J621 Setup "*.exe "$FEED"/*.exe.blockmap)
for file in "${old[@]}"; do
case "$(basename "$file")" in
*"$VERSION"*) ;;
*) echo " pruning $(basename "$file")"; rm -f "$file" ;;
esac
done
cp -f "${deb[@]}" "${zst[@]}" "${linux_meta[@]}" "$FEED"/ 2>/dev/null || true
if [ "${#win_exe[@]}" -gt 0 ]; then
cp -f "${win_exe[@]}" "${win_meta[@]}" "${blockmaps[@]}" "$FEED"/ 2>/dev/null || true
@@ -83,6 +92,16 @@ if [ -n "$REMOTE" ]; then
echo
echo "==> Copying the feed to $REMOTE_TARGET:$REMOTE_FEED ..."
if ! ssh "$REMOTE_TARGET" "mkdir -p '$REMOTE_FEED'"; then
echo "ssh to $REMOTE_TARGET failed; nothing was copied." >&2
exit 1
fi
# The remote feed keeps only the current version, like the local one.
if ! ssh "$REMOTE_TARGET" \
"find '$REMOTE_FEED' -maxdepth 1 -type f \\( -name '*.deb' -o -name '*.pkg.tar.zst' -o -name '*.exe' -o -name '*.exe.blockmap' \\) ! -name '*$VERSION*' -exec rm -f {} +"; then
echo "==> Could not prune old installers on $REMOTE_TARGET (continuing)." >&2
fi
transferred=0
if command -v rsync >/dev/null 2>&1; then
if rsync -a --info=progress2 "$FEED"/ "$REMOTE_TARGET:$REMOTE_FEED/"; then
@@ -92,10 +111,6 @@ if [ -n "$REMOTE" ]; then
fi
fi
if [ "$transferred" -eq 0 ]; then
if ! ssh "$REMOTE_TARGET" "mkdir -p '$REMOTE_FEED'"; then
echo "ssh to $REMOTE_TARGET failed; nothing was copied." >&2
exit 1
fi
if ! tar -C "$FEED" -cf - . |
ssh "$REMOTE_TARGET" "tar -C '$REMOTE_FEED' -xf -"; then
echo "Copy to $REMOTE_TARGET:$REMOTE_FEED failed." >&2
+8 -2
View File
@@ -10,10 +10,16 @@ REGISTRY="gitea.rainbow-herring.ts.net/jakebreath/j621-frontend"
REGISTRY_HOST="$(printf '%s' "$REGISTRY" | cut -d/ -f1)"
SHA="${1:-$(git rev-parse --short HEAD)}"
BUILDER=multiarch
PLATFORMS="linux/amd64,linux/arm64"
PLATFORMS="${PLATFORMS:-linux/amd64,linux/arm64}"
echo "==> Logging in to $REGISTRY_HOST ..."
docker login "$REGISTRY_HOST"
if [ -n "${REGISTRY_USER:-}" ] && [ -n "${REGISTRY_TOKEN:-}" ]; then
# Non-interactive login for CI (workflow passes GITHUB_TOKEN).
printf '%s' "$REGISTRY_TOKEN" | docker login "$REGISTRY_HOST" \
-u "$REGISTRY_USER" --password-stdin
else
docker login "$REGISTRY_HOST"
fi
if ! docker buildx inspect "$BUILDER" >/dev/null 2>&1; then
echo "==> Creating buildx builder '$BUILDER' ..."
+5
View File
@@ -20,6 +20,11 @@ src/preload.ts window.j621Desktop bridge
electron-builder.yml deb + pacman + nsis packaging
```
The main process attaches a `Referer` to requests for e621 hosts: Chromium
sends no referrer from a custom-scheme page, and e621's CDN answers
cross-site image loads without one with a 403 (which Chromium then blocks as
ORB). API calls work either way.
## Development
```bash
+2 -2
View File
@@ -1,12 +1,12 @@
{
"name": "j621-desktop",
"version": "0.1.0",
"version": "0.1.1",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "j621-desktop",
"version": "0.1.0",
"version": "0.1.1",
"license": "LicenseRef-Jake-Labs-Non-Commercial",
"dependencies": {
"electron-updater": "^6.8.9"
+1 -1
View File
@@ -1,7 +1,7 @@
{
"name": "j621-desktop",
"productName": "J621",
"version": "0.1.0",
"version": "0.1.1",
"private": true,
"description": "Desktop shell for the J621 self-hosted media archive",
"author": {
+67 -2
View File
@@ -12,9 +12,9 @@
* the same way it does in a browser. The only desktop-specific bits are the
* origin, external-link handling and (later) updates.
*/
import { app, BrowserWindow, dialog, ipcMain, Menu, protocol, shell } from "electron";
import { app, BrowserWindow, dialog, ipcMain, Menu, protocol, session, shell } from "electron";
import { autoUpdater } from "electron-updater";
import { readFileSync } from "node:fs";
import { readFileSync, rmSync } from "node:fs";
import { readFile, stat, writeFile } from "node:fs/promises";
import path from "node:path";
@@ -24,6 +24,15 @@ const APP_ORIGIN = `${SCHEME}://${HOST}`;
const DEV_SERVER = (process.env.J621_DEV_SERVER ?? "").replace(/\/+$/, "");
const SMOKE = process.argv.includes("--j621-smoke");
if (SMOKE) {
// A throwaway profile keeps the checks deterministic (always the first-run
// setup screen) and never touches real state or leaves a stale
// single-instance lock behind.
const smokeDir = path.join(app.getPath("temp"), "j621-smoke");
rmSync(smokeDir, { recursive: true, force: true });
app.setPath("userData", smokeDir);
}
protocol.registerSchemesAsPrivileged([
{
scheme: SCHEME,
@@ -83,6 +92,29 @@ async function fileResponse(file: string): Promise<Response> {
});
}
/** Paths the nginx proxy sends to Django; there is no backend behind app://. */
const BACKEND_PREFIXES = ["/api/", "/admin/", "/static/", "/health"];
/**
* e621's CDN refuses cross-site image loads that carry no Referer
* (`Sec-Fetch-Site: cross-site` with an empty referrer returns 403), and
* Chromium never sends one for pages on a custom scheme like app://j621.
* Attach a normal e621 referrer to its hosts so images load; API calls are
* unaffected.
*/
function installRefererFix(targetSession: Electron.Session): void {
targetSession.webRequest.onBeforeSendHeaders(
{ urls: ["*://e621.net/*", "*://*.e621.net/*"] },
(details, callback) => {
const headers = details.requestHeaders;
if (!headers.Referer && !headers.referer) {
headers.Referer = "https://e621.net/";
}
callback({ requestHeaders: headers });
},
);
}
async function handleAppRequest(request: Request): Promise<Response> {
const url = new URL(request.url);
if (url.host !== HOST) return new Response("Not found", { status: 404 });
@@ -91,6 +123,21 @@ async function handleAppRequest(request: Request): Promise<Response> {
let pathname = decodeURIComponent(url.pathname);
if (!pathname || pathname === "/") pathname = "/index.html";
// Never answer backend paths with the SPA: that made the setup screen's
// connection test "succeed" against the shell's own origin.
if (
BACKEND_PREFIXES.some(
(prefix) => pathname === prefix.replace(/\/$/, "") || pathname.startsWith(prefix),
)
) {
return new Response(
JSON.stringify({
detail: "No backend is attached to the desktop app; set one in /setup.",
}),
{ status: 404, headers: { "content-type": "application/json" } },
);
}
const target = path.resolve(root, "." + pathname);
if (target !== root && !target.startsWith(root + path.sep)) {
return new Response("Forbidden", { status: 403 });
@@ -426,6 +473,7 @@ function runSmokeTest(win: BrowserWindow): void {
const asset = await fetch("/favicon.svg");
const route = await fetch("/gallery/some/deep/route");
const routeBody = await route.text();
const health = await fetch("/health");
localStorage.setItem("j621.smoke", "ok");
history.pushState({}, "", "/gallery");
return {
@@ -433,6 +481,18 @@ function runSmokeTest(win: BrowserWindow): void {
title: document.title,
assetOk: asset.ok && (await asset.text()).includes("<svg"),
routeOk: route.ok && routeBody.includes('id="root"'),
// /health must not fall back to the SPA (false "connected" test).
healthOk: health.status === 404,
// Testing an empty backend URL on the desktop must say so instead
// of "Connected" against the shell's own origin.
emptyTestOk: await (async () => {
const button = [...document.querySelectorAll("button")].find((b) =>
b.textContent.includes("Test connection"),
);
if (!button) return false;
button.click();
return waitFor(() => document.body.innerText.includes("no backend of its own"));
})(),
mounted: await waitFor(() => (document.querySelector("#root")?.childElementCount ?? 0) > 0),
setupOk: document.body.innerText.includes("Where is your backend?"),
storageOk: localStorage.getItem("j621.smoke") === "ok",
@@ -452,6 +512,10 @@ function runSmokeTest(win: BrowserWindow): void {
routeOk: true,
mounted: true,
setupOk: !DEV_SERVER,
// Both are app://-only checks: the Vite dev server serves the SPA
// for /health and never shows the setup screen.
healthOk: !DEV_SERVER,
emptyTestOk: !DEV_SERVER,
storageOk: true,
historyOk: true,
opfsOk: true,
@@ -517,6 +581,7 @@ if (!gotLock) {
app.whenReady().then(() => {
protocol.handle(SCHEME, handleAppRequest);
installRefererFix(session.defaultSession);
ipcMain.handle("j621:version", () => app.getVersion());
ipcMain.handle("j621:open-external", (_event, url: unknown) => {
+2
View File
@@ -15,6 +15,7 @@ import { ConfirmDialog } from "@/components/ConfirmDialog";
import { StatusFooter } from "@/components/StatusFooter";
import { StatusPill } from "@/components/StatusPill";
import { Toasts } from "@/components/Toasts";
import { UploadIndicator } from "@/components/UploadIndicator";
import { api } from "@/lib/api";
import { hasBackend } from "@/lib/backend";
import { cn } from "@/lib/cn";
@@ -153,6 +154,7 @@ export function AppShell() {
<div className="ml-auto flex items-center gap-3">
{backend ? (
<>
<UploadIndicator />
<StatusPill status={status} />
{user ? (
<>
@@ -0,0 +1,56 @@
import { Link } from "react-router-dom";
import { cn } from "@/lib/cn";
import { useUploadStatus } from "@/lib/uploadStatus";
const PHASE_LABELS: Record<string, string> = {
md5: "e621 MD5",
visual: "Visual similarity",
iqdb: "IQDB",
};
/**
* Shell-wide progress for the server-side upload pipeline.
*
* Staged uploads keep processing after the Upload page is closed, so this
* pill is the "leave and keep an eye on it" view: it lives in the header on
* every page and links back to the board.
*/
export function UploadIndicator() {
const { data: status } = useUploadStatus();
if (!status || !status.active) return null;
const total = Math.max(status.total, status.processed + status.failed);
const done = status.processed + status.failed;
const percent = total > 0 ? Math.min(100, Math.round((done / total) * 100)) : 0;
const broken = status.status === "paused" || status.status === "error";
const label =
status.status === "paused"
? "uploads paused"
: status.status === "error"
? "uploads failed"
: (PHASE_LABELS[status.phase] ?? "processing uploads");
const count = total > 0 ? `${done}/${total}` : `${status.outstanding} left`;
return (
<Link
to="/upload"
title={status.error || "Staged uploads are being processed in the background"}
className={cn(
"flex items-center gap-2 rounded-md border px-2 py-1 font-mono text-[10px] transition",
broken
? "border-ctp-red/40 text-ctp-red hover:bg-ctp-red/10"
: "border-ctp-surface1 text-ctp-subtext0 hover:bg-ctp-surface0 hover:text-ctp-text",
)}
>
<span className="hidden sm:inline">{label}</span>
<span>{count}</span>
<span className="h-1 w-12 overflow-hidden rounded-full bg-ctp-surface0">
<span
className={cn("block h-full", broken ? "bg-ctp-red" : "bg-ctp-teal")}
style={{ width: `${percent}%` }}
/>
</span>
</Link>
);
}
+29 -8
View File
@@ -17,6 +17,16 @@ import { setToken } from "@/lib/api";
type TestResult = { ok: boolean; text: string } | null;
async function testBackend(base: string): Promise<TestResult> {
// The desktop shell has no server of its own: an empty URL is not a
// same-origin backend, and /health on app://j621 is answered by the shell.
if (!base && window.j621Desktop) {
return {
ok: false,
text:
"The desktop app has no backend of its own. Enter your J621 server's " +
"URL, or continue without a backend.",
};
}
const controller = new AbortController();
const timeout = window.setTimeout(() => controller.abort(), 5_000);
try {
@@ -110,10 +120,15 @@ export function SetupPage() {
</div>
<p className="mt-3 text-sm leading-relaxed text-ctp-subtext0">
Enter the origin this app should call for its API. Leave it blank
when the app and the API are served from the same domain. Without a
backend the app runs in local mode: e621 browsing only, with
credentials kept in this browser.
{window.j621Desktop
? "Enter the origin of your J621 server (for example " +
"https://j621.example.com). Without a backend the app runs in " +
"local mode: e621 browsing only, with credentials kept on this " +
"device."
: "Enter the origin this app should call for its API. Leave it " +
"blank when the app and the API are served from the same " +
"domain. Without a backend the app runs in local mode: e621 " +
"browsing only, with credentials kept in this browser."}
</p>
<label className="mt-5 flex flex-col gap-1.5">
@@ -122,7 +137,11 @@ export function SetupPage() {
</span>
<input
className={cn(inputClass, "font-mono")}
placeholder="https://j621.example.com — blank for this server"
placeholder={
window.j621Desktop
? "https://j621.example.com"
: "https://j621.example.com — blank for this server"
}
value={value}
onChange={(event) => {
setValue(event.target.value);
@@ -161,7 +180,7 @@ export function SetupPage() {
{testing ? <Spinner className="h-3.5 w-3.5" /> : null}
{testing ? "Testing…" : "Test connection"}
</Button>
{value.trim() ? (
{value.trim() && !window.j621Desktop ? (
<Button variant="ghost" onClick={() => save("")}>
Use this server
</Button>
@@ -172,8 +191,10 @@ export function SetupPage() {
</div>
<p className="mt-4 text-xs leading-relaxed text-ctp-overlay0">
Stored in this browser only. Changing the backend signs you out of
the previous one.
{window.j621Desktop
? "Stored in this app only."
: "Stored in this browser only."}{" "}
Changing the backend signs you out of the previous one.
</p>
</div>
</div>
File diff suppressed because it is too large Load Diff
+12 -9
View File
@@ -89,18 +89,19 @@ export function effectiveCredentials(
const GIT_HASH = typeof __GIT_HASH__ === "string" ? __GIT_HASH__ : "dev";
const CLIENT_VERSION = `J621/${GIT_HASH} (JakeBreath)`;
// e621 allows 2 requests/second hard, 1/second sustained — and the IQDB
// endpoint is stricter, so stay comfortably under it. Serialize every request
// through a queue with a minimum gap.
// e621 allows 2 requests/second hard, 1/second sustained. Serialize every
// request through a queue with a minimum gap; IQDB is throttled much harder
// by e621, so it gets a wider gap of its own.
let lastRequestAt = 0;
let queue: Promise<unknown> = Promise.resolve();
/** A hung request would block the whole serialized queue forever. */
const REQUEST_TIMEOUT_MS = 20_000;
const REQUEST_GAP_MS = 1500;
const REQUEST_GAP_MS = 1000;
const IQDB_GAP_MS = 2500;
/** A 429 (or a CORS-blocked failure) pauses every e621 call for a while. */
const RATE_LIMIT_COOLDOWN_MS = 60_000;
/** A 429 (or a CORS-blocked failure) pauses every e621 call briefly. */
const RATE_LIMIT_COOLDOWN_MS = 15_000;
const COOLDOWN_KEY = "j621.e621.cooldown";
let cooldownUntil = 0;
@@ -143,10 +144,10 @@ function schedule<T>(task: () => Promise<T>): Promise<T> {
return run;
}
async function throttle(): Promise<void> {
async function throttle(minGap = REQUEST_GAP_MS): Promise<void> {
const wait = Math.max(
0,
lastRequestAt + REQUEST_GAP_MS - Date.now(),
lastRequestAt + minGap - Date.now(),
e621CooldownRemainingMs(),
);
if (wait > 0) {
@@ -168,7 +169,9 @@ export function e621Request<T>(
options: E621RequestOptions = {},
): Promise<T> {
return schedule(async () => {
await throttle();
await throttle(
path.includes("iqdb_queries") ? IQDB_GAP_MS : REQUEST_GAP_MS,
);
const base = credentials.base_url.replace(/\/+$/, "");
const url = new URL(`${base}/${path.replace(/^\/+/, "")}`);
+33 -5
View File
@@ -355,12 +355,17 @@ export interface TempUpload {
status: "pending" | "visual_match" | "completed" | "error";
resolution: "" | "auto_md5" | "duplicate" | "linked" | "custom";
e621_post_id: number | null;
e621_data: E621StoredPost | null;
/** Full post payload: only present on the detail endpoint. */
e621_data?: E621StoredPost | null;
custom_rating: Rating;
custom_tags: string[];
custom_notes: string;
iqdb_data: E621IqdbCandidate[] | null;
visual_matches:
/** Only present on the detail endpoint. */
custom_tags?: string[];
/** Only present on the detail endpoint. */
custom_notes?: string;
/** Only present on the detail endpoint. */
iqdb_data?: E621IqdbCandidate[] | null;
/** Only present on the detail endpoint. */
visual_matches?:
| {
j_id: string;
filename: string;
@@ -371,10 +376,33 @@ export interface TempUpload {
library_j_id: string | null;
file_url: string | null;
preview_url: string | null;
/** Background pipeline state (see the UploadRun status endpoint). */
pipeline_error: string;
attempts: number;
md5_checked?: boolean;
visual_checked?: boolean;
iqdb_checked?: boolean;
processing?: boolean;
similar_count?: number;
created_at: string;
updated_at: string;
}
/** Progress of the server-side staged-upload pipeline. */
export interface UploadStatus {
status: "idle" | "running" | "paused" | "error";
active: boolean;
phase: "" | "md5" | "visual" | "iqdb";
total: number;
processed: number;
matched: number;
failed: number;
error: string;
outstanding: number;
waiting: { md5: number; visual: number; iqdb: number };
updated_at: string | null;
}
export interface StatJob {
kind: "download" | "match";
id: string;
+30
View File
@@ -0,0 +1,30 @@
import { useQuery } from "@tanstack/react-query";
import { api } from "@/lib/api";
import { hasBackend } from "@/lib/backend";
import type { UploadStatus } from "@/lib/types";
import { useAuth } from "@/store/auth";
export const uploadStatusQueryKey = ["upload-status"] as const;
/**
* Poll the server-side upload pipeline while it has work and stop once it is
* idle. Shared by the Upload page and the shell indicator so both show the
* same state without extra requests.
*/
export function useUploadStatus(enabled = true) {
const user = useAuth((state) => state.user);
return useQuery({
queryKey: uploadStatusQueryKey,
queryFn: () => api<UploadStatus>("/api/uploads/status/"),
enabled: enabled && hasBackend() && Boolean(user?.can_upload),
// TanStack takes the smallest interval across the observers, so the
// Upload page and the shell pill can safely share this query.
refetchInterval: (query) => (query.state.data?.active ? 2_000 : false),
});
}
/** Ask the server to (re)start processing the caller's staged uploads. */
export function postProcessUploads(): Promise<UploadStatus> {
return api<UploadStatus>("/api/uploads/process/", { method: "POST" });
}