Stop storing completed uploads; announce them through a live feed

Every auto-matched, duplicate or manually resolved upload left a completed
TempUpload row on the board until it was dismissed by hand, so the rows
accumulated without bound and the bulk dismiss (capped at 1000 ids) failed
once there were more. The original app never stored these: they are
notifications, not records.

- complete_temp_upload now appends {filename, J-ID, resolution, post} to a
  bounded recent_completions feed on UploadRun and deletes the staged row
- staging duplicates never create a board record either; the create response
  carries the J-ID and preview so the SPA can show the card immediately
- status_payload returns the feed (newest first, signed thumbnails) for the
  live board; finalize_round counts deleted matches in processed
- resolve/link-bulk return synthetic completion payloads
- migration 0012 adds the field and purges the existing completed backlog
  (and any stray staged files) on deploy
This commit is contained in:
2026-09-23 18:24:38 -05:00
parent 37085c5dac
commit 9ababb8b48
5 changed files with 237 additions and 45 deletions
@@ -0,0 +1,34 @@
# Generated by Django 6.1.1 on 2026-09-23
from django.db import migrations, models
def purge_completed_uploads(apps, schema_editor):
"""Completed staged uploads are notifications, not records.
The board used to keep every auto-uploaded/indexed row until it was
dismissed by hand, so they accumulated without bound (and bulk dismissal
is capped at 1000 ids). Completions now live in a small per-run feed, so
the old rows are removed here.
"""
TempUpload = apps.get_model("library", "TempUpload")
for temp in TempUpload.objects.filter(status="completed").iterator():
if temp.file:
temp.file.delete(save=False)
temp.delete()
class Migration(migrations.Migration):
dependencies = [
("library", "0011_upload_pipeline"),
]
operations = [
migrations.AddField(
model_name="uploadrun",
name="recent_completions",
field=models.JSONField(blank=True, default=list),
),
migrations.RunPython(purge_completed_uploads, migrations.RunPython.noop),
]
+5
View File
@@ -229,6 +229,11 @@ class UploadRun(models.Model):
matched = models.IntegerField(default=0)
failed = models.IntegerField(default=0)
error = models.TextField(blank=True, default="")
# Rolling feed of recent completions for the upload board. Completed
# staged uploads are deleted as soon as they are indexed; this only tells
# the live page "filename -> J-x" while it watches. Bounded, never
# dismissed, and ignored by fresh page loads.
recent_completions = models.JSONField(default=list, blank=True)
started_at = models.DateTimeField(null=True, blank=True)
updated_at = models.DateTimeField(auto_now=True)
+81 -21
View File
@@ -170,12 +170,18 @@ class StagedUploadWorkflowTests(TestCase):
self.assertEqual(len(body["resolved"]), 2)
self.assertEqual(body["errors"], [])
for temp in (first, second):
temp.refresh_from_db()
self.assertEqual(temp.status, TempUpload.STATUS_COMPLETED)
self.assertIsNotNone(temp.library_item_id)
self.assertEqual(temp.library_item.rating, "q")
self.assertEqual(temp.library_item.uploaded_by_id, self.uploader.id)
# Completed uploads are notifications now: the staged rows are gone
# and the live feed carries the filename -> J-ID mapping.
self.assertFalse(
TempUpload.objects.filter(pk__in=[first.id, second.id]).exists()
)
feed = UploadRun.objects.get(user=self.uploader).recent_completions
self.assertEqual(
{entry["filename"] for entry in feed}, {"one.png", "two.png"}
)
for item in MediaItem.objects.all():
self.assertEqual(item.rating, "q")
self.assertEqual(item.uploaded_by_id, self.uploader.id)
untouched.refresh_from_db()
self.assertEqual(untouched.status, TempUpload.STATUS_PENDING)
@@ -231,6 +237,25 @@ class StagedUploadWorkflowTests(TestCase):
self.assertEqual(response.status_code, 200)
return seed
def test_duplicate_upload_returns_a_completion_without_a_record(self):
client = self.api_client(self.uploader)
first = self.upload_via_api(client, "same.png")
self.assertEqual(first.status_code, 201)
# Index the first upload so the second one is byte-identical.
seed = TempUpload.objects.get(pk=first.json()["temp_id"])
self.resolve_bulk(client, [seed.id], "s")
response = self.upload_via_api(client, "same.png")
self.assertEqual(response.status_code, 201)
body = response.json()
self.assertEqual(body["status"], TempUpload.STATUS_COMPLETED)
self.assertEqual(body["resolution"], TempUpload.RESOLUTION_DUPLICATE)
self.assertTrue(body["library_j_id"].startswith("J-"))
# Duplicates never become board records; the feed announces them.
self.assertFalse(TempUpload.objects.filter(pk=body["temp_id"]).exists())
feed = UploadRun.objects.get(user=self.uploader).recent_completions
self.assertEqual(feed[-1]["filename"], "same.png")
def test_upload_defers_visual_similarity_to_its_phase(self):
client = self.api_client(self.uploader)
self.seed_library_item(client, "seed-defer")
@@ -312,14 +337,23 @@ class StagedUploadWorkflowTests(TestCase):
body = response.json()
self.assertEqual(body["errors"], [])
self.assertEqual(len(body["updated"]), 2)
for temp, post_id in ((first, 900001), (second, 900002)):
temp.refresh_from_db()
self.assertEqual(temp.status, TempUpload.STATUS_COMPLETED)
self.assertEqual(temp.library_item_id is not None, True)
self.assertEqual(temp.e621_post_id, post_id)
self.assertEqual(temp.resolution, TempUpload.RESOLUTION_AUTO_MD5)
self.assertEqual(temp.library_item.e621_post_id, post_id)
self.assertEqual(MediaItem.objects.count(), 2)
for entry, post_id in zip(body["updated"], (900001, 900002)):
self.assertEqual(entry["resolution"], TempUpload.RESOLUTION_AUTO_MD5)
self.assertEqual(entry["e621_post_id"], post_id)
self.assertTrue(entry["library_j_id"].startswith("J-"))
# Indexed uploads no longer leave a board record; the completion feed
# carries them for the live page instead.
self.assertFalse(TempUpload.objects.exists())
self.assertEqual(
{item.e621_post_id for item in MediaItem.objects.all()},
{900001, 900002},
)
feed = UploadRun.objects.get(user=self.uploader).recent_completions
self.assertEqual(len(feed), 2)
self.assertEqual(
{entry["filename"] for entry in feed},
{"bulk-link-1.png", "bulk-link-2.png"},
)
class IqdbRecordingTests(TestCase):
@@ -461,6 +495,7 @@ class UploadPipelineTests(TestCase):
def test_md5_match_auto_imports_the_file(self):
temp = self.stage(label="match")
temp_id = temp.id
post = {
"id": 123456,
"rating": "s",
@@ -477,17 +512,42 @@ class UploadPipelineTests(TestCase):
):
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)
# The indexed upload leaves no board record; the feed reports it.
self.assertFalse(TempUpload.objects.filter(pk=temp_id).exists())
item = MediaItem.objects.get(e621_post_id=123456)
self.assertEqual(item.uploaded_by_id, self.uploader.id)
run = UploadRun.objects.get(user=self.uploader)
self.assertEqual(run.status, UploadRun.STATUS_IDLE)
self.assertEqual(run.matched, 1)
self.assertEqual(run.processed, 1)
self.assertEqual(len(run.recent_completions), 1)
entry = run.recent_completions[0]
self.assertEqual(entry["id"], str(temp_id))
self.assertEqual(entry["item_id"], item.id)
self.assertEqual(entry["filename"], "match.png")
self.assertEqual(entry["resolution"], TempUpload.RESOLUTION_AUTO_MD5)
def test_status_reports_the_completion_feed(self):
temp = self.stage(label="feed")
client = self.api_client(self.uploader)
post = {
"id": 654321,
"rating": "s",
"file": {"md5": temp.md5, "url": "https://static1.e621.net/data/f.png"},
}
with mock.patch.object(
upload_pipeline.e621,
"check_md5_batch",
return_value={temp.md5: post},
):
upload_pipeline.run_pipeline(self.uploader.id)
body = client.get("/api/uploads/status/").json()
completions = body["recent_completions"]
self.assertEqual(len(completions), 1)
self.assertEqual(completions[0]["filename"], "feed.png")
self.assertEqual(completions[0]["j_id"], f"J-{MediaItem.objects.get().id}")
self.assertIn("/thumbnail/", completions[0]["thumbnail_url"])
def test_unmatched_file_runs_every_phase(self):
temp = self.stage(label="nomatch")
+73 -14
View File
@@ -29,7 +29,7 @@ from django.db.models import F, Q
from django.utils import timezone
from . import e621, services
from .models import TempUpload, UploadRun
from .models import MediaItem, TempUpload, UploadRun
logger = logging.getLogger(__name__)
@@ -46,6 +46,9 @@ MAX_ATTEMPTS = 3
# Round-level e621 retries before the run is paused.
ROUND_ATTEMPTS = 3
ROUND_RETRY_SECONDS = 20
# Completions kept in the live feed. The board only shows what happened while
# the page was open, so a bounded rolling window is plenty.
COMPLETION_FEED_LIMIT = 200
WORK_STATUSES = (TempUpload.STATUS_PENDING, TempUpload.STATUS_VISUAL_MATCH)
VIDEO_RE = r"\.(mp4|webm)$"
@@ -116,7 +119,7 @@ def waiting_counts(user):
}
def status_payload(user):
def status_payload(user, request=None):
"""Cheap state for the shell/upload page to poll."""
run = UploadRun.objects.filter(user=user).first()
outstanding = count_outstanding(user)
@@ -141,10 +144,67 @@ def status_payload(user):
"error": run.error if run is not None else "",
"outstanding": outstanding,
"waiting": waiting_counts(user),
"recent_completions": completion_payload(run, user, request=request),
"updated_at": run.updated_at.isoformat() if run is not None else None,
}
def completion_payload(run, user, request=None):
"""The live completion feed: filename -> J-ID for freshly indexed uploads.
Completed ``TempUpload`` rows are deleted, so this is the only place the
board learns about them. It is a notification feed, not durable state:
bounded, never dismissed, and ignored by fresh page loads.
"""
entries = list(run.recent_completions or []) if run is not None else []
if not entries:
return []
ids = [entry.get("item_id") for entry in entries if entry.get("item_id")]
items = MediaItem.objects.in_bulk(ids)
out = []
for entry in reversed(entries): # newest first
item = items.get(entry.get("item_id"))
if item is None:
continue
out.append(
{
"id": entry.get("id"),
"filename": entry.get("filename"),
"j_id": f"J-{item.id}",
"resolution": entry.get("resolution", ""),
"post_id": entry.get("post_id"),
"thumbnail_url": services.signed_media_url(
item, user, "thumbnail", request=request
),
"at": entry.get("at"),
}
)
return out
def record_completion(temp, item, resolution=""):
"""Append one completion to the owner's feed; never fails an import."""
entry = {
"id": str(temp.pk),
"filename": temp.original_filename,
"item_id": item.pk,
"resolution": resolution or temp.resolution or "",
"post_id": temp.e621_post_id,
"at": timezone.now().isoformat(),
}
try:
with transaction.atomic():
run, _ = UploadRun.objects.select_for_update().get_or_create(
user_id=temp.user_id
)
feed = list(run.recent_completions or [])
feed.append(entry)
run.recent_completions = feed[-COMPLETION_FEED_LIMIT:]
run.save(update_fields=["recent_completions", "updated_at"])
except Exception: # noqa: BLE001 - a notification must not break an import
logger.exception("Could not record the completion of %s", temp.pk)
def reap_stale_claims():
"""Queue rows left claimed by a recycled worker and pause dead runs."""
cutoff = timezone.now() - STALE_CLAIM_AFTER
@@ -237,20 +297,17 @@ def run_pipeline(user_id):
break
process_round(run, user, rows)
except PipelinePaused as exc:
run.status = UploadRun.STATUS_PAUSED
run.phase = ""
run.error = str(exc)
run.save()
_save_run(run, status=UploadRun.STATUS_PAUSED, phase="", error=str(exc))
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()
_save_run(
run,
status=UploadRun.STATUS_ERROR,
phase="",
error=f"The upload pipeline stopped: {exc}",
)
else:
run.status = UploadRun.STATUS_IDLE
run.phase = ""
run.save()
_save_run(run, status=UploadRun.STATUS_IDLE, phase="")
def claim_round(user, size=CLAIM_SIZE):
@@ -460,7 +517,9 @@ def finalize_round(run, ids, matched):
TempUpload.objects.filter(pk__in=release).update(claimed_at=None)
_save_run(
run,
processed=run.processed + len(finished),
# Completed rows are deleted as they are imported, so they cannot be
# seen in the refreshed rows; count the matches explicitly.
processed=run.processed + len(finished) + matched,
failed=run.failed + len(failed - finished),
matched=run.matched + matched,
)
+44 -10
View File
@@ -158,6 +158,13 @@ def complete_temp_upload(temp, download_url=None):
item.save(update_fields=update_fields + ["updated_at"])
temp.save()
from .upload_pipeline import record_completion
record_completion(temp, item)
# The board learns about completions from the live feed, so the record is
# deleted as soon as the file is indexed: nothing left to dismiss. Delete
# through the queryset so callers keep ``temp.pk`` for their response.
TempUpload.objects.filter(pk=temp.pk).delete()
return item
@@ -190,6 +197,28 @@ class TempUploadViewSet(
user=self.request.user
)
def _completed_payload(self, temp, item):
"""Synthetic row for an upload that is indexed immediately.
Duplicates and resolved uploads never leave a board record; the SPA
turns this response (or the live completion feed) into a "J-x
uploaded" card that lives only in the page session.
"""
return {
"temp_id": str(temp.pk),
"original_filename": temp.original_filename,
"md5": temp.md5,
"size": temp.size,
"status": TempUpload.STATUS_COMPLETED,
"resolution": temp.resolution,
"e621_post_id": temp.e621_post_id,
"library_j_id": f"J-{item.id}",
"file_url": None,
"preview_url": services.signed_media_url(
item, self.request.user, "thumbnail", request=self.request
),
}
def create(self, request):
upload = request.FILES.get("file")
if upload is None:
@@ -217,6 +246,13 @@ class TempUploadViewSet(
temp.resolution = TempUpload.RESOLUTION_DUPLICATE
temp.library_item = existing
temp.file.delete(save=False)
temp.save()
from .upload_pipeline import record_completion
record_completion(temp, existing)
payload = self._completed_payload(temp, existing)
TempUpload.objects.filter(pk=temp.pk).delete()
return Response(payload, status=status.HTTP_201_CREATED)
# 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()
@@ -257,7 +293,7 @@ class TempUploadViewSet(
"""Cheap pipeline state for the shell indicator and the upload page."""
from .upload_pipeline import status_payload
return Response(status_payload(request.user))
return Response(status_payload(request.user, request=request))
@action(detail=False, methods=["post"])
def process(self, request):
@@ -269,7 +305,7 @@ class TempUploadViewSet(
from .upload_pipeline import start_pipeline, status_payload
start_pipeline(request.user)
return Response(status_payload(request.user))
return Response(status_payload(request.user, request=request))
@action(detail=True, methods=["post"])
def retry(self, request, pk=None):
@@ -314,7 +350,7 @@ class TempUploadViewSet(
update["status"] = TempUpload.STATUS_PENDING
TempUpload.objects.filter(pk=temp.pk).update(**update)
start_pipeline(request.user)
return Response(status_payload(request.user))
return Response(status_payload(request.user, request=request))
@action(detail=False, methods=["post"], url_path="retry-all")
def retry_all(self, request):
@@ -340,7 +376,7 @@ class TempUploadViewSet(
e621_checked_at=None,
)
start_pipeline(request.user)
return Response(status_payload(request.user))
return Response(status_payload(request.user, request=request))
@action(detail=False, methods=["post"], url_path="discard-bulk")
def discard_bulk(self, request):
@@ -494,7 +530,7 @@ class TempUploadViewSet(
temp.save()
try:
complete_temp_upload(temp, download_url=download_url)
item = complete_temp_upload(temp, download_url=download_url)
except Exception as exc: # noqa: BLE001 - report completion failures
logger.exception("Could not complete staged upload %s", temp.id)
temp.status = TempUpload.STATUS_ERROR
@@ -503,8 +539,7 @@ class TempUploadViewSet(
{"detail": f"Could not finish the upload: {exc}"},
status=status.HTTP_400_BAD_REQUEST,
)
temp.refresh_from_db()
return Response(self.get_serializer(temp).data)
return Response(self._completed_payload(temp, item))
@action(detail=False, methods=["post"], url_path="link-bulk")
def link_bulk(self, request):
@@ -576,15 +611,14 @@ class TempUploadViewSet(
candidate_url = ""
temp.save()
try:
complete_temp_upload(temp, download_url=candidate_url or None)
item = complete_temp_upload(temp, download_url=candidate_url or None)
except Exception as exc: # noqa: BLE001 - report per-file failures
logger.exception("Could not complete staged upload %s", temp.id)
temp.status = TempUpload.STATUS_ERROR
temp.save(update_fields=["status", "updated_at"])
errors.append({"temp_id": temp_id, "error": str(exc)})
continue
temp.refresh_from_db()
updated.append(self.get_serializer(temp).data)
updated.append(self._completed_payload(temp, item))
return Response({"updated": updated, "errors": errors})