diff --git a/proc/tasks.py b/proc/tasks.py index 3dec57af9..2abc0c7dc 100644 --- a/proc/tasks.py +++ b/proc/tasks.py @@ -847,8 +847,10 @@ def task_migrate_and_publish_articles( force_update=force_migrate_document_records or force_migrate_document_files, to_migrate_articles=True, ) - issue_proc_ids = selected_issue_procs.values_list("journal_proc_id", "id").distinct() - + issue_proc_ids = selected_issue_procs.values_list( + "journal_proc_id", "id" + ).distinct().iterator(chunk_size=1000) + for journal_proc_id, issue_proc_id in issue_proc_ids: items_to_process.setdefault(journal_proc_id, []).append(issue_proc_id) else: @@ -856,7 +858,7 @@ def task_migrate_and_publish_articles( journal_proc_ids = JournalProc.select_items( collection_acron_list=collection_acron_list, journal_acron_list=journal_acron_list, - ).values_list("id", flat=True) + ).values_list("id", flat=True).iterator(chunk_size=1000) items_to_process = {journal_proc_id: None for journal_proc_id in journal_proc_ids} total_journals_to_process = len(items_to_process) @@ -869,7 +871,11 @@ def task_migrate_and_publish_articles( task_exec.total_to_process = total_journals_to_process total_processed = 0 - for journal_proc_id, issue_proc_id_list in items_to_process.items(): + # Drena items_to_process ao dispatchar para liberar progressivamente + # as listas de issue_proc_ids (potencialmente grandes em coleções com + # muitos journals × issues), em vez de mantê-lo inteiro em RAM. + while items_to_process: + journal_proc_id, issue_proc_id_list = items_to_process.popitem() kwargs = {} kwargs.update(kwargs_) kwargs["journal_proc_id"] = journal_proc_id @@ -951,8 +957,6 @@ def task_migrate_and_publish_articles_by_journal( force_update=force_import_acron_id_file, ) - qa_api_data = get_api_data(journal_proc.collection, "issue", "QA") - public_api_data = get_api_data(journal_proc.collection, "issue", "PUBLIC") total_processed = 0 total_to_process = 0 @@ -980,14 +984,21 @@ def task_migrate_and_publish_articles_by_journal( exclude_issue_proc_id_list=list(issue_proc_id_list), status_list=status, force_update=force_update, - ).values_list("issue_proc_id", "id").distinct() + ).values_list("issue_proc_id", "id").distinct().iterator(chunk_size=1000) for issue_proc_id, article_proc_id in selected_article_proc_items: issue_proc_and_related_article_proc_id_list.setdefault(issue_proc_id, []).append(article_proc_id) total_to_process = len(issue_proc_and_related_article_proc_id_list) - for issue_proc_id, article_proc_id_list in issue_proc_and_related_article_proc_id_list.items(): + # Drena o dict ao dispatchar para liberar progressivamente as listas + # de article_proc_ids (que podem somar muitos inteiros para journals + # grandes), em vez de mantê-lo inteiro em RAM até o fim do loop. + while issue_proc_and_related_article_proc_id_list: + issue_proc_id, article_proc_id_list = issue_proc_and_related_article_proc_id_list.popitem() total_processed += 1 + # qa_api_data/public_api_data não são propagados: task_sync_issue + # (despachada por _by_issue) cacheia get_api_data internamente, + # evitando login HTTP redundante e mensagens grandes no broker. task_migrate_and_publish_articles_by_issue.delay( user_id=user_id, username=username, @@ -997,9 +1008,7 @@ def task_migrate_and_publish_articles_by_journal( force_update=force_update, force_migrate_document_records=force_migrate_document_records, force_migrate_document_files=force_migrate_document_files, - qa_api_data=qa_api_data, - public_api_data=public_api_data, - ) + ) task_exec.total_processed = total_processed task_exec.total_to_process = total_to_process task_exec.finish() @@ -1023,9 +1032,20 @@ def task_migrate_and_publish_articles_by_issue( force_update=False, force_migrate_document_records=False, force_migrate_document_files=False, - qa_api_data=None, - public_api_data=None, + # qa_api_data e public_api_data foram removidos: nunca eram lidos no + # corpo da função e infláveis (token + credenciais) no payload Celery. + # Aceitos como **kwargs para compatibilidade com mensagens já enfileiradas. + **legacy_kwargs, ): + # Sinaliza kwargs inesperados (típos, etc.) sem quebrar; ignora os legacy + # conhecidos (qa_api_data/public_api_data). + _LEGACY_IGNORED = {"qa_api_data", "public_api_data"} + unknown_kwargs = [k for k in legacy_kwargs if k not in _LEGACY_IGNORED] + if unknown_kwargs: + logging.warning( + "task_migrate_and_publish_articles_by_issue: ignoring unknown kwargs %s", + unknown_kwargs, + ) task_params = { "user_id": user_id, "username": username, @@ -1062,7 +1082,8 @@ def task_migrate_and_publish_articles_by_issue( # (issue_proc.docs_status e issue_proc.files_status estão como DONE) total_articles_to_process = len(article_proc_id_list) article_procs = ArticleProc.objects.select_related( - "issue_proc", + "issue_proc", "issue_proc__journal_proc", + "collection", "sps_pkg", ).filter( id__in=article_proc_id_list ) @@ -1082,6 +1103,9 @@ def task_migrate_and_publish_articles_by_issue( issue_proc_id_list=[issue_proc_id], status_list=status, force_update=force_update, + ).select_related( + "issue_proc", "issue_proc__journal_proc", + "collection", "sps_pkg", ) total_articles_to_process = article_procs.count() task_exec.total_to_process = total_articles_to_process @@ -1089,7 +1113,9 @@ def task_migrate_and_publish_articles_by_issue( task_exec.add_event("Migrate articles") total_processed = 0 exceptions = {} - for article_proc in article_procs: + # iterator(chunk_size) evita carregar todos os ArticleProc em memória, + # reduzindo RAM no worker quando o issue tem muitos artigos. + for article_proc in article_procs.iterator(chunk_size=100): try: article = article_proc.migrate_article(user, force_update) total_processed += 1 @@ -1099,17 +1125,11 @@ def task_migrate_and_publish_articles_by_issue( task_exec.add_exception(exceptions[article_proc.pid]) task_exec.total_processed = total_processed - - article_ids_to_publish = ArticleProc.objects.select_related( - "issue_proc", "sps_pkg", - ).filter( - Q(qa_ws_status__in=status) | Q(public_ws_status__in=status), - issue_proc=issue_proc, - sps_pkg__pid_v3__isnull=False, - ).values_list("id", flat=True) - total_articles_to_publish = article_ids_to_publish.count() - task_exec.add_number("total_articles_to_publish", total_articles_to_publish) + # task_sync_issue é sempre despachada para QA e PUBLIC: além de publicar + # artigos pendentes, ela garante que não haja duplicidade no site — + # condição que não pode ser determinada antecipadamente sem consultar o + # próprio site. for website_label in (QA, PUBLIC): task_exec.add_event(f"Schedule Publish articles / sync issue tasks for {website_label}") task_sync_issue.apply_async( @@ -1172,15 +1192,16 @@ def task_sync_issue( elif website_kind == PUBLIC: query_by_status = Q(public_ws_status__in=status) - article_ids_to_publish = ArticleProc.objects.select_related( - "issue_proc", "sps_pkg", - ).filter( + article_ids_qs = ArticleProc.objects.filter( query_by_status, issue_proc=issue_proc, sps_pkg__pid_v3__isnull=False, ).values_list("id", flat=True) - task_exec.total_to_process = article_ids_to_publish.count() + # .count() + .iterator() em vez de list(qs): evita carregar todos os + # IDs de artigos do issue em RAM (issues grandes podem ter milhares). + # Custo: 1 SELECT COUNT(*) extra; ganho: footprint constante no loop. + task_exec.total_to_process = article_ids_qs.count() total_processed = 0 api_data = get_api_data(issue_proc.collection, "article", website_kind) @@ -1189,7 +1210,7 @@ def task_sync_issue( task_exec.finish() return - for article_proc_id in article_ids_to_publish: + for article_proc_id in article_ids_qs.iterator(chunk_size=500): try: # executa de forma síncrona para evitar muitos processos em paralelo, o que pode causar lentidão e instabilidade no ambiente de origem (ex: site clássico) task_publish_article( diff --git a/publication/api/publication.py b/publication/api/publication.py index 9e3196528..67833f488 100644 --- a/publication/api/publication.py +++ b/publication/api/publication.py @@ -1,6 +1,8 @@ +import copy import json import logging import sys +import time import traceback import urllib @@ -20,13 +22,59 @@ def get_api(collection, content_type, website_kind): return {"error": f"Website {collection} {website_kind} is not enabled ({api_data})"} +# Cache em processo para api_data (inclui token). Objetivo: solicitar token +# apenas 1x por (collection, content_type, website_kind) e revalidar somente +# quando expirar. Estratégia: +# - TTL longo (alinhado ao lifetime típico de JWT) como teto de segurança, +# evitando reter tokens indefinidamente em workers de longa duração; +# - invalidação explícita via invalidate_api_data_cache(...) — chamadores +# que detectarem falha de autenticação no PublicationAPI.post_data podem +# purgar a entrada e forçar novo login na próxima chamada. +# Observação: PublicationAPI.post_data já refaz get_token() na própria +# instância ao receber falha, então um token expirado entre a leitura do +# cache e o uso resulta em 1 retry interno (não em erro propagado). +_API_DATA_CACHE = {} +_API_DATA_CACHE_TTL = 3600 # segundos (1h); ajustar via clear/invalidate se necessário + + +def _api_data_cache_key(collection, content_type, website_kind): + return (getattr(collection, "pk", None), content_type, website_kind) + + +def clear_api_data_cache(): + """Limpa todo o cache em processo de api_data (uso em testes/admin).""" + _API_DATA_CACHE.clear() + + +def invalidate_api_data_cache(collection, content_type, website_kind=None): + """Invalida uma entrada específica do cache. + + Deve ser chamado por consumidores ao detectar falha de autenticação + (token expirado/revogado) no resultado de PublicationAPI.post_data, + para forçar novo login na próxima chamada de get_api_data. + """ + _API_DATA_CACHE.pop( + _api_data_cache_key(collection, content_type, website_kind), None + ) + + def get_api_data(collection, content_type, website_kind=None): + key = _api_data_cache_key(collection, content_type, website_kind) + cached = _API_DATA_CACHE.get(key) + now = time.time() + if cached and now - cached[0] < _API_DATA_CACHE_TTL: + # deepcopy: defesa contra mutação de estruturas aninhadas pelos + # chamadores (ex.: api_data["verify"] = verify em task_publish_articles). + return copy.deepcopy(cached[1]) try: - return get_api(collection, content_type, website_kind) + data = get_api(collection, content_type, website_kind) except WebSiteConfiguration.DoesNotExist: return {"error": f"Website does not exist: {collection} {website_kind}"} except Exception as e: return {"error": f"Unable to get API data for {content_type} {collection} {website_kind}: {type(e)} {e}"} + if isinstance(data, dict) and not data.get("error") and key[0] is not None: + _API_DATA_CACHE[key] = (now, copy.deepcopy(data)) + return data class PublicationAPI: diff --git a/publication/api/test_publication.py b/publication/api/test_publication.py new file mode 100644 index 000000000..aa52d4fd1 --- /dev/null +++ b/publication/api/test_publication.py @@ -0,0 +1,108 @@ +"""Tests for publication.api.publication.get_api_data caching.""" +import unittest +from unittest.mock import patch + +from publication.api import publication as publication_module +from publication.api.publication import ( + clear_api_data_cache, + get_api_data, + invalidate_api_data_cache, +) + + +class _FakeCollection: + def __init__(self, pk): + self.pk = pk + + def __str__(self): + return f"Collection({self.pk})" + + +class GetApiDataCacheTest(unittest.TestCase): + def setUp(self): + clear_api_data_cache() + + def tearDown(self): + clear_api_data_cache() + + def test_caches_successful_response_per_key(self): + collection = _FakeCollection(pk=1) + with patch.object( + publication_module, + "get_api", + return_value={"token": "abc", "post_data_url": "http://x"}, + ) as mocked: + first = get_api_data(collection, "issue", "QA") + second = get_api_data(collection, "issue", "QA") + third = get_api_data(collection, "issue", "PUBLIC") + + # Mesma collection/content_type/website_kind: chamado 1x. + # Chave diferente para PUBLIC: 1x adicional. + self.assertEqual(mocked.call_count, 2) + self.assertEqual(first["token"], "abc") + self.assertEqual(second["token"], "abc") + self.assertEqual(third["token"], "abc") + + def test_returns_copy_so_caller_mutation_does_not_poison_cache(self): + collection = _FakeCollection(pk=2) + with patch.object( + publication_module, + "get_api", + return_value={"token": "t", "post_data_url": "u", "nested": {"x": 1}}, + ): + first = get_api_data(collection, "article", "PUBLIC") + first["verify"] = True # mutação como em task_publish_articles + first["nested"]["x"] = 999 # mutação aninhada + second = get_api_data(collection, "article", "PUBLIC") + + self.assertNotIn("verify", second) + self.assertEqual(second["nested"]["x"], 1) + + def test_does_not_cache_error_responses(self): + collection = _FakeCollection(pk=3) + # Primeira chamada retorna erro, segunda retorna sucesso. + responses = iter([ + {"error": "boom"}, + {"token": "ok", "post_data_url": "u"}, + ]) + with patch.object( + publication_module, + "get_api", + side_effect=lambda *a, **kw: next(responses), + ) as mocked: + err = get_api_data(collection, "issue", "QA") + ok = get_api_data(collection, "issue", "QA") + + self.assertEqual(mocked.call_count, 2) + self.assertIn("error", err) + self.assertEqual(ok["token"], "ok") + + def test_clear_cache_helper(self): + collection = _FakeCollection(pk=4) + with patch.object( + publication_module, + "get_api", + return_value={"token": "z"}, + ) as mocked: + get_api_data(collection, "issue", "QA") + clear_api_data_cache() + get_api_data(collection, "issue", "QA") + + self.assertEqual(mocked.call_count, 2) + + def test_invalidate_cache_forces_relogin_for_specific_key(self): + collection_a = _FakeCollection(pk=5) + collection_b = _FakeCollection(pk=6) + with patch.object( + publication_module, + "get_api", + return_value={"token": "t"}, + ) as mocked: + get_api_data(collection_a, "issue", "QA") + get_api_data(collection_b, "issue", "QA") + # Invalida apenas a entrada de A; B segue cacheada. + invalidate_api_data_cache(collection_a, "issue", "QA") + get_api_data(collection_a, "issue", "QA") # re-login + get_api_data(collection_b, "issue", "QA") # cache hit + + self.assertEqual(mocked.call_count, 3)