Skip to content
Open
Show file tree
Hide file tree
Changes from 2 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
99 changes: 62 additions & 37 deletions proc/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -951,8 +951,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 +978,17 @@ 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():
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 +998,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 +1022,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 +1072,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 +1093,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,28 +1115,37 @@ 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()

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.

# Materializa uma única vez para evitar count() + iteração separados.
article_ids_to_publish = list(
ArticleProc.objects.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 = len(article_ids_to_publish)
task_exec.add_number("total_articles_to_publish", total_articles_to_publish)

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(
kwargs=dict(
user_id=user_id,
username=username,
issue_proc_id=issue_proc.id,
website_kind=website_label,
status=status,
force_update=force_update,
# Só agenda task_sync_issue se há artigos a publicar; caso contrário
# cada despacho dispararia get_api_data (login HTTP) e queries
# redundantes só para descobrir que não há trabalho.
if total_articles_to_publish:
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(
kwargs=dict(
user_id=user_id,
username=username,
issue_proc_id=issue_proc.id,
website_kind=website_label,
status=status,
force_update=force_update,
)
)
else:
task_exec.add_event(
f"Skip task_sync_issue for issue_proc {issue_proc.id}: no articles to publish"
)

task_exec.finish()
Expand Down Expand Up @@ -1172,15 +1197,15 @@ 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(
query_by_status,
issue_proc=issue_proc,
sps_pkg__pid_v3__isnull=False,
).values_list("id", flat=True)
article_ids_to_publish = list(
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()
task_exec.total_to_process = len(article_ids_to_publish)
total_processed = 0

api_data = get_api_data(issue_proc.collection, "article", website_kind)
Expand Down
30 changes: 29 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,39 @@ 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). Evita um HTTP login a cada
# chamada de get_api_data dentro do mesmo worker. O TTL é curto o suficiente
# para tolerar expiração razoável de token sem precisar de invalidação manual.
_API_DATA_CACHE = {}
_API_DATA_CACHE_TTL = 600 # segundos


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 o cache em processo de api_data (uso em testes/admin)."""
_API_DATA_CACHE.clear()


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
90 changes: 90 additions & 0 deletions publication/api/test_publication.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,90 @@
"""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,
)


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)