Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
81 changes: 51 additions & 30 deletions proc/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -847,16 +847,18 @@ 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:
task_exec.add_event(_("Select journals by collection"))
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)
Expand All @@ -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
Expand Down Expand Up @@ -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")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@copilot o objetivo de criar qa_api_data e public_api_data é para que o token seja solicitado apenas 1 vezes no início, depois somente quando expirar. Qual uma boa alternativa para lidar com este propósito?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Bom ponto — o TTL=600s da primeira versão forçava re-login a cada 10 min mesmo com token válido, contrariando o objetivo. Ajustes no commit 1db793c em publication/api/publication.py:

  • TTL elevado para 3600s (1h), alinhado ao lifetime típico de JWT — apenas como teto de segurança para não reter tokens indefinidamente em workers de longa duração.
  • Novo helper invalidate_api_data_cache(collection, content_type, website_kind): o consumidor que detectar falha de autenticação no resultado de PublicationAPI.post_data pode purgar a entrada e forçar novo login na próxima chamada.
  • PublicationAPI.post_data já refaz get_token() automaticamente na própria instância em caso de falha (linha 113 do publication.py original) — então um token expirado entre a leitura do cache e o uso resulta no máximo em 1 retry interno, não em erro propagado.

Resultado: 1 login por (collection, content_type, website_kind) por worker até o token efetivamente expirar (detectado via invalidate_api_data_cache) ou até atingir 1h. Se preferir TTL maior/menor ou parametrizá-lo via Django setting, posso ajustar.

public_api_data = get_api_data(journal_proc.collection, "issue", "PUBLIC")
total_processed = 0
total_to_process = 0

Expand Down Expand Up @@ -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,
Expand All @@ -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()
Expand All @@ -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,
Expand Down Expand Up @@ -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
)
Expand All @@ -1082,14 +1103,19 @@ 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

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
Expand All @@ -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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@copilot remove identificação de que há artigos para publicar e considere que o sync_issue tem que ser executado independentemente se há ou não artigos para publicar pois ele precisa ser exeuctado para garantir que não exista duplicidade no site. E isso não é possível identifcar previamente

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Removida a verificação de total_articles_to_publish e a query associada no commit e75d39b. task_sync_issue volta a ser despachada incondicionalmente para QA e PUBLIC, preservando a garantia de deduplicação no site.

# 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(
Expand Down Expand Up @@ -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)
Expand All @@ -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(
Expand Down
50 changes: 49 additions & 1 deletion publication/api/publication.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
import copy
import json
import logging
import sys
import time
import traceback
import urllib

Expand All @@ -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:
Expand Down
108 changes: 108 additions & 0 deletions publication/api/test_publication.py
Original file line number Diff line number Diff line change
@@ -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)