"""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)."))