diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 1d626f63d..89c5e1463 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -47,14 +47,17 @@ jobs: - name: Checkout Code Repository uses: actions/checkout@v3 + - name: Prepare writable directories + run: sudo install -d -o 1000 -g 1000 -m 0755 logs + - name: Build the Stack - run: docker-compose -f local.yml build + run: docker compose -f local.yml build - name: Run DB Migrations - run: docker-compose -f local.yml run --rm django python manage.py migrate + run: docker compose -f local.yml run --rm django python manage.py migrate - name: Run Django Tests - run: docker-compose -f local.yml run django pytest + run: docker compose -f local.yml run django pytest - name: Tear down the Stack - run: docker-compose -f local.yml down + run: docker compose -f local.yml down diff --git a/article/controller.py b/article/controller.py index 3bf5ac216..00b541c15 100644 --- a/article/controller.py +++ b/article/controller.py @@ -4,6 +4,7 @@ import sys import traceback +from django.db.models import F, Q from packtools.sps.formats.am import am from article.models import Article, ArticleExporter, ArticleFunding, ArticleSource @@ -12,11 +13,8 @@ from core.mongodb import write_item from core.utils.harvesters import AMHarvester, OPACHarvester from institution.models import Sponsor -from journal.models import Journal, SciELOJournal -from pid_provider.choices import ( - PPXML_STATUS_TODO, - PPXML_STATUS_INVALID, -) +from journal.models import SciELOJournal +from pid_provider import choices as pid_provider_choices from pid_provider.models import PidProviderXML from tracker.models import UnexpectedEvent @@ -24,70 +22,6 @@ class ArticleIsNotAvailableError(Exception): ... -# def get_pp_xml_ids( -# collection_acron_list=None, -# journal_acron_list=None, -# from_pub_year=None, -# until_pub_year=None, -# from_updated_date=None, -# until_updated_date=None, -# proc_status_list=None, -# ): -# return select_pp_xml( -# collection_acron_list, -# journal_acron_list, -# from_pub_year, -# until_pub_year, -# from_updated_date, -# until_updated_date, -# proc_status_list=proc_status_list, -# ).values_list("id", flat=True) - - -# def select_pp_xml( -# collection_acron_list=None, -# journal_acron_list=None, -# from_pub_year=None, -# until_pub_year=None, -# from_updated_date=None, -# until_updated_date=None, -# proc_status_list=None, -# params=None, -# ): -# params = params or {} - -# q = Q() -# if journal_acron_list or collection_acron_list: -# issns = Journal.get_issn_list(collection_acron_list, journal_acron_list) -# issn_print_list = issns["issn_print_list"] -# issn_electronic_list = issns["issn_electronic_list"] - -# if issn_print_list or issn_electronic_list: -# q = Q(issn_print__in=issn_print_list) | Q( -# issn_electronic__in=issn_electronic_list -# ) -# elif issn_print_list: -# q = Q(issn_print__in=issn_print_list) -# elif issn_electronic_list: -# q = Q(issn_electronic__in=issn_electronic_list) - -# if from_updated_date: -# params["updated__gte"] = from_updated_date -# if until_updated_date: -# params["updated__lte"] = until_updated_date - -# if from_pub_year: -# params["pub_year__gte"] = from_pub_year -# if until_pub_year: -# params["pub_year__lte"] = until_pub_year - -# if proc_status_list: -# params["proc_status__in"] = proc_status_list - -# logging.info(params) -# return PidProviderXML.objects.filter(q, **params) - - def load_financial_data(row, user): article_findings = [] for institution in row.get("funding_source").split(","): @@ -133,33 +67,25 @@ def export_article_to_articlemeta( ) -> bool: try: + events = [] if not article.classic_available(collection_acron_list): raise ArticleIsNotAvailableError( f"Article {article} {collection_acron_list} (classic) is not available. Unable to export to ArticleMeta." ) + new_available = article.new_available(collection_acron_list).exists() - logging.info(f"Article new {article} {collection_acron_list} {new_available}") + events.append(f"Article new {article} {collection_acron_list} {new_available}") - logging.info( + events.append( f"export_article_to_articlemeta: {article}, collections: {collection_acron_list}, force_update: {force_update}" ) legacy_keys_items = list(article.get_legacy_keys( collection_acron_list, is_active=True )) - logging.info(f"Legacy keys to process: {legacy_keys_items}") + events.append(f"Legacy keys to process: {legacy_keys_items}") if not legacy_keys_items: - UnexpectedEvent.create( - exception=ValueError("No legacy keys found for article"), - detail={ - "operation": "export_article_to_articlemeta", - "article": str(article), - "collection_acron_list": collection_acron_list, - "force_update": force_update, - }, - ) - return + raise ValueError("No legacy keys found for article") - events = [] external_data = { "created_at": article.created.strftime("%Y-%m-%d"), "document_type": article.article_type, @@ -175,14 +101,15 @@ def export_article_to_articlemeta( article_data = {} for legacy_keys in legacy_keys_items: + col = legacy_keys.get("collection") + pid = legacy_keys.get("pid") + item = f"{pid}-{col}" try: exporter = None response = None events = [] data = {} - col = legacy_keys.get("collection") - pid = legacy_keys.get("pid") - + if not article_data: events.append("building articlemeta format for article") article_data = am.build(article.xmltree, external_data) @@ -224,9 +151,7 @@ def export_article_to_articlemeta( json.dumps(data) data["code"] except Exception as e: - logging.exception(e) response = str(data) - logging.info(data) raise e # Export the article to ArticleMeta @@ -256,10 +181,11 @@ def export_article_to_articlemeta( ) else: UnexpectedEvent.create( + action="export_article_to_articlemeta", + item=item, exception=e, exc_traceback=exc_traceback, detail={ - "operation": "export_article_to_articlemeta", "article": str(article), "legacy_keys": str(legacy_keys), "events": events, @@ -270,11 +196,11 @@ def export_article_to_articlemeta( except Exception as e: exc_type, exc_value, exc_traceback = sys.exc_info() UnexpectedEvent.create( + action="export_article_to_articlemeta", + item=str(article), exception=e, exc_traceback=exc_traceback, detail={ - "operation": "export_article_to_articlemeta", - "article": str(article), "collection_acron_list": collection_acron_list, "force_update": force_update, "traceback": traceback.format_exc(), @@ -310,7 +236,13 @@ def bulk_export_articles_to_articlemeta( version: Version identifier for export Returns: - bool: True if the export was successful, False otherwise + bool: True when the batch processing finishes. Errors exporting + individual articles are logged and do not interrupt the batch. + + Raises: + ValueError: If no articles match the provided filters. + Exception: If an unexpected error occurs outside the processing of an + individual article. The error is logged and re-raised. """ try: params = {} @@ -329,21 +261,16 @@ def bulk_export_articles_to_articlemeta( params=params ) if not queryset.exists(): - UnexpectedEvent.create( - exception=ValueError("No articles found for the given filters"), - detail={ - "operation": "bulk_export_articles_to_articlemeta", - "collection_acron_list": collection_acron_list, - "journal_acron_list": journal_acron_list, - "from_pub_year": from_pub_year, - "until_pub_year": until_pub_year, - "from_date": str(from_date) if from_date else None, - "until_date": str(until_date) if until_date else None, - "days_to_go_back": days_to_go_back, - "force_update": force_update, - }, + args = dict( + collection_acron_list=collection_acron_list, + journal_acron_list=journal_acron_list, + from_pub_year=from_pub_year, + until_pub_year=until_pub_year, + from_updated_date=from_date, + until_updated_date=until_date, + params=params, ) - return False + raise ValueError(f"No articles found. Arguments: {args}") for article in queryset.select_related("journal", "journal__official", "pp_xml").iterator(): try: @@ -362,14 +289,13 @@ def bulk_export_articles_to_articlemeta( # Registra erro do article mas continua processando outros exc_type, exc_value, exc_traceback = sys.exc_info() UnexpectedEvent.create( + action="bulk_export_articles_to_articlemeta", + item=str(article), exception=e, exc_traceback=exc_traceback, detail={ - "function": "bulk_export_articles_to_articlemeta", - "article_id": article.id, - "article_pid": getattr(article, "pid", None), - "journal_acron": getattr(article, "journal_acron", None), - "pub_year": getattr(article, "pub_year", None), + "collection_acron_list": collection_acron_list, + "version": version, "force_update": force_update, }, ) @@ -380,10 +306,11 @@ def bulk_export_articles_to_articlemeta( except Exception as e: exc_type, exc_value, exc_traceback = sys.exc_info() UnexpectedEvent.create( + action="bulk_export_articles_to_articlemeta", + item="", exception=e, exc_traceback=exc_traceback, detail={ - "function": "bulk_export_articles_to_articlemeta", "collection_acron_list": collection_acron_list, "journal_acron_list": journal_acron_list, "from_pub_year": from_pub_year, @@ -399,34 +326,37 @@ def bulk_export_articles_to_articlemeta( class ArticleIteratorBuilder: """ - Monta e encadeia iteradores de seleção de artigos para despacho ao pipeline. - - Cada método ``_iter_from_*`` é um gerador que yields kwargs prontos para - ``task_process_article_pipeline``. Os iteradores ativos são determinados - pelos argumentos exclusivos presentes na instância — múltiplos podem estar - ativos simultaneamente. - - Argumentos exclusivos e seus iteradores: - - ========================= ================================================ - Argumento exclusivo Iterador ativado - ========================= ================================================ - proc_status_list _iter_from_pid_provider - data_status_list _iter_from_article - limit / timeout / opac_url _iter_from_harvest - article_source_status_list _iter_from_article_source - (nenhum) _iter_from_pid_provider (padrão) - ========================= ================================================ - - Usage:: - - it = ArticleIteratorBuilder( + Constrói iteradores de seleção de artigos para despacho ao pipeline + (``task_process_article_pipeline``). + + O ``__init__`` guarda apenas os filtros COMUNS a mais de uma fonte + (usuário, coleção/periódico, intervalo de datas/anos, força de + atualização, parâmetros de harvest). Os filtros que são exclusivos de + uma única fonte (``proc_status_list``, ``data_status_list``, + ``article_source_status_list``) NÃO ficam no ``__init__`` — são passados + diretamente ao método correspondente, o que deixa explícito qual fonte + está sendo usada em cada chamada. + + Cada método ``from_*``: + - usa os atributos comuns já armazenados na instância; + - retorna um gerador independente que yields kwargs prontos para + ``task_process_article_pipeline``; + - representa exatamente UM ponto de entrada do pipeline. + + Mutuamente excludentes por construção: a classe não tem mais um + ``__iter__`` que decide sozinha quais iteradores ativar e os roda em + sequência. Quem chama (normalmente ``task_dispatch_articles``) escolhe + explicitamente UM método por execução. Isso elimina a possibilidade de + dois iteradores selecionarem o mesmo artigo/``pp_xml_id`` na mesma + chamada. + + Uso:: + + builder = ArticleIteratorBuilder( user=user, collection_acron_list=["scl"], - proc_status_list=["todo"], - data_status_list=["invalid"], ) - for kwargs in it: + for kwargs in builder.from_pid_provider(proc_status_list=["todo"]): task_process_article_pipeline.delay(**kwargs) """ @@ -439,13 +369,10 @@ def __init__( until_pub_year=None, from_date=None, until_date=None, - proc_status_list=None, - data_status_list=None, - article_source_status_list=None, + force_update=None, limit=None, timeout=None, opac_url=None, - force_update=None, stop=None, ): self.user = user @@ -455,106 +382,135 @@ def __init__( self.until_pub_year = until_pub_year self.from_date = from_date self.until_date = until_date - self.proc_status_list = proc_status_list - self.data_status_list = data_status_list - self.article_source_status_list = article_source_status_list + self.force_update = force_update self.limit = limit self.timeout = timeout self.opac_url = opac_url - self.force_update = force_update self.stop = stop - self._iter_from_harvest_count = 0 - self._iter_from_article_source_count = 0 - self._iter_from_pid_provider_count = 0 - self._iter_from_article_count = 0 - - def __iter__(self): - yield from self._iter_from_harvest() - yield from self._iter_from_article_source() - yield from self._iter_from_pid_provider() - yield from self._iter_from_article() - - logging.info(f"Iterators summary: harvest={self._iter_from_harvest_count}, " - f"article_source={self._iter_from_article_source_count}, " - f"pid_provider={self._iter_from_pid_provider_count}, " - f"article={self._iter_from_article_count}") - # ------------------------------------------------------------------ - # Iteradores de seleção + # from_pid_provider # ------------------------------------------------------------------ + def from_pid_provider(self, proc_status_list=None): + filters = { + "proc_status__in": proc_status_list or [pid_provider_choices.PPXML_STATUS_TODO], + } + if self.from_date: + filters["updated__gte"] = self.from_date + if self.until_date: + filters["updated__lte"] = self.until_date + if self.from_pub_year: + filters["pub_year__gte"] = self.from_pub_year + if self.until_pub_year: + filters["pub_year__lte"] = self.until_pub_year + + params = {} + if self.collection_acron_list: + params["collection__acron3__in"] = self.collection_acron_list + if self.journal_acron_list: + params["journal_acron__in"] = self.journal_acron_list - def _iter_from_pid_provider(self): - """Itera PidProviderXML filtrados por periódico, data e status.""" - journal_issn_groups = ( - Journal.get_journal_issns(self.collection_acron_list, self.journal_acron_list) - or [None] + q = Q() + if params: + journal_issns = SciELOJournal.objects.select_related( + "journal__official" + ).filter( + **params + ).values_list( + "journal__official__issn_print", "journal__official__issn_electronic" + ).distinct() + + issn_list = set() + for issn_print, issn_electronic in journal_issns: + if issn_print: + issn_list.add(issn_print) + if issn_electronic: + issn_list.add(issn_electronic) + + q = Q(issn_print__in=issn_list) | Q(issn_electronic__in=issn_list) + qs = ( + PidProviderXML.objects.filter(q, **filters) + .order_by("-updated") + .values(pp_xml_id=F("id")) + .distinct() ) - for journal_issns in journal_issn_groups: - issn_list = [i for i in journal_issns if i] if journal_issns else None - if journal_issns and not issn_list: - continue - qs = PidProviderXML.get_queryset( - issn_list=issn_list, - from_pub_year=self.from_pub_year, - until_pub_year=self.until_pub_year, - from_updated_date=self.from_date, - until_updated_date=self.until_date, - proc_status_list=self.proc_status_list or [PPXML_STATUS_TODO, PPXML_STATUS_INVALID], - ) - self._iter_from_pid_provider_count += qs.count() - for item in qs.iterator(): - yield {"pp_xml_id": item.id} - logging.info(f"_iter_from_pid_provider: yielded {self._iter_from_pid_provider_count} items") - def _iter_from_article(self): + yield from qs.iterator() + + # ------------------------------------------------------------------ + # from_article + # ------------------------------------------------------------------ + def from_article(self, data_status_list=None): """ - Itera Articles filtrados por data_status. - Yields None para artigos sem pp_xml recuperável (sinaliza skip). + Itera Article pendentes ou com erro. """ - filters = { - "data_status__in": self.data_status_list or [ - choices.DATA_STATUS_PENDING, - choices.DATA_STATUS_UNDEF, - choices.DATA_STATUS_INVALID, - ] - } - journal_id_list = Journal.get_ids( - collection_acron_list=self.collection_acron_list, - journal_acron_list=self.journal_acron_list, - ) - if journal_id_list: - filters["journal__in"] = journal_id_list + filters = {} + if data_status_list: + filters["data_status__in"] = data_status_list + if self.collection_acron_list: + filters["journal__scielojournal__collection__acron3__in"] = ( + self.collection_acron_list + ) + if self.journal_acron_list: + filters["journal__scielojournal__journal_acron__in"] = ( + self.journal_acron_list + ) if self.from_pub_year: - filters["pub_year__gte"] = self.from_pub_year + filters["pub_date_year__gte"] = self.from_pub_year if self.until_pub_year: - filters["pub_year__lte"] = self.until_pub_year + filters["pub_date_year__lte"] = self.until_pub_year if self.from_date: filters["updated__gte"] = self.from_date if self.until_date: filters["updated__lte"] = self.until_date - articles = Article.objects.filter(**filters) - self._iter_from_article_count += articles.count() - for article in articles.iterator(): - if not article.pp_xml: - try: - article.pp_xml = PidProviderXML.get_by_pid_v3(pid_v3=article.pid_v3) - article.save(update_fields=["pp_xml"]) - except Exception as e: - logging.error(f"pp_xml not found for article {article.id}: {e}") - yield None - continue - yield {"pp_xml_id": article.pp_xml.id} - logging.info(f"_iter_from_article: yielded {self._iter_from_article_count} articles") + filters["valid"] = False + base_qs = Article.objects.filter(**filters).distinct() - def _iter_from_harvest(self): - """Itera documentos coletados via OPAC ou ArticleMeta.""" + # Artigos que já têm pp_xml: values() já entrega o dict pronto. + yield from base_qs.filter(pp_xml__isnull=False).values("pp_xml_id").iterator() + # ------------------------------------------------------------------ + # from_article_source + # ------------------------------------------------------------------ + def from_article_source(self, article_source_status_list=None): + """ + Itera ArticleSources pendentes ou com erro. + + Otimização: ``.values(article_source_id=F("id"))`` sobre a + queryset retornada por ``get_queryset_to_complete_data`` já + entrega o dict pronto — sem instanciar cada ArticleSource + completo (não precisamos de mais nenhum campo do objeto para + montar o kwarg de despacho). + """ + params = {} + if article_source_status_list: + params["status__in"] = article_source_status_list + if self.from_date: + params["updated__gte"] = self.from_date + if self.until_date: + params["updated__lte"] = self.until_date + + if self.force_update: + qs = Q() + else: + qs = ( + Q(pid_provider_xml__proc_status__in=pid_provider_choices.PPXML_STATUS_TO_CREATE_OR_UPDATE_ARTICLE_SOURCE) | + Q(pid_provider_xml__isnull=True) | Q(file__isnull=True) + ) + yield from ArticleSource.objects.filter( + qs, + **params, + ).values(article_source_id=F("id")).iterator() + + # ------------------------------------------------------------------ + # from_harvest + # ------------------------------------------------------------------ + def from_harvest(self): + """Itera documentos coletados via OPAC ou ArticleMeta.""" if Collection.objects.count() == 0: Collection.load(self.user) - count = 0 params = {} if self.collection_acron_list: params["collection__acron3__in"] = self.collection_acron_list @@ -572,35 +528,14 @@ def _iter_from_harvest(self): for collection_acron, journal_acron, issn_scielo in collection_and_journal_items: harvester = self._build_harvester(collection_acron, journal_acron, issn_scielo) for document in harvester.harvest_documents(): - count += 1 yield { "xml_url": document["url"], "collection_acron": collection_acron, "pid": document["pid_v2"], "source_date": document.get("processing_date") or document.get("origin_date"), - "is_public": document.get("is_public") + "is_public": document.get("is_public"), + "document": document, } - - self._iter_from_harvest_count = count - logging.info(f"Harvest iterator yielded {count} documents") - - def _iter_from_article_source(self): - """Itera ArticleSources pendentes ou com erro.""" - count = 0 - for article_source in ArticleSource.get_queryset_to_complete_data( - self.from_date, - self.until_date, - self.force_update, - self.article_source_status_list, - ): - count += 1 - yield {"article_source_id": article_source.id} - self._iter_from_article_source_count += count - logging.info(f"ArticleSource iterator yielded {count} items") - - # ------------------------------------------------------------------ - # Helpers privados - # ------------------------------------------------------------------ def _build_harvester(self, collection_acron, journal_acron=None, journal_id=None): """Instancia o harvester adequado para a coleção.""" @@ -619,4 +554,3 @@ def _build_harvester(self, collection_acron, journal_acron=None, journal_id=None if journal_id: kwargs["journal"] = journal_id return AMHarvester("article", collection_acron, **kwargs) - diff --git a/article/migrations/0050_remove_articlesource_am_article_and_more.py b/article/migrations/0050_remove_articlesource_am_article_and_more.py new file mode 100644 index 000000000..a235949a6 --- /dev/null +++ b/article/migrations/0050_remove_articlesource_am_article_and_more.py @@ -0,0 +1,103 @@ +# Generated by Django 5.2.7 on 2026-07-27 20:02 + +import django.db.models.deletion +from django.db import migrations, models + + +BATCH_SIZE = 1000 + + +def backfill_article_source(apps, schema_editor): + ArticleSource = apps.get_model("article", "ArticleSource") + database = schema_editor.connection.alias + + queryset = ( + ArticleSource.objects.using(database) + .select_related("am_article", "pid_provider_xml") + .prefetch_related("pid_provider_xml__collections") + .order_by("pk") + ) + + batch = [] + + for article_source in queryset.iterator(chunk_size=BATCH_SIZE): + am_article = article_source.am_article + pid_provider_xml = article_source.pid_provider_xml + changed = False + + if am_article: + if am_article.collection_id: + article_source.collection_id = am_article.collection_id + changed = True + if am_article.pid: + article_source.pid = am_article.pid + changed = True + + if pid_provider_xml: + if not article_source.pid and pid_provider_xml.v2: + article_source.pid = pid_provider_xml.v2 + changed = True + + if not article_source.collection_id: + collection_ids = [ + collection.pk + for collection in pid_provider_xml.collections.all() + ] + if len(collection_ids) == 1: + article_source.collection_id = collection_ids[0] + changed = True + + if changed: + batch.append(article_source) + + if len(batch) >= BATCH_SIZE: + ArticleSource.objects.using(database).bulk_update( + batch, + ["collection", "pid"], + batch_size=BATCH_SIZE, + ) + batch.clear() + + if batch: + ArticleSource.objects.using(database).bulk_update( + batch, + ["collection", "pid"], + batch_size=BATCH_SIZE, + ) + + +class Migration(migrations.Migration): + + dependencies = [ + ("article", "0049_alter_articlesource_status"), + ("collection", "0007_collection_platform_status"), + ("pid_provider", "0012_pidproviderxml_collections"), + ] + + operations = [ + migrations.AddField( + model_name="articlesource", + name="collection", + field=models.ForeignKey( + blank=True, + help_text="Related Collection instance", + null=True, + on_delete=django.db.models.deletion.SET_NULL, + to="collection.collection", + verbose_name="Collection", + ), + ), + migrations.AddField( + model_name="articlesource", + name="pid", + field=models.CharField(blank=True, max_length=24, null=True), + ), + migrations.RunPython( + backfill_article_source, + migrations.RunPython.noop, + ), + migrations.RemoveField( + model_name="articlesource", + name="am_article", + ), + ] diff --git a/article/models.py b/article/models.py index bef06d7c7..8e9b8b543 100755 --- a/article/models.py +++ b/article/models.py @@ -43,7 +43,7 @@ from institution.models import Publisher, Sponsor from issue.models import Issue, TableOfContents from journal.models import Journal, SciELOJournal -from pid_provider.choices import PPXML_STATUS_DONE +from pid_provider import choices as pid_provider_choices from pid_provider.models import PidProviderXML from pid_provider.provider import PidProvider from location.models import Location @@ -447,8 +447,7 @@ def source(self): try: return descriptive_format(**leg_dict) except Exception as ex: - logging.exception("Erro on article %s, error: %s" % (self.pid_v2, ex)) - return "" + return str(leg_dict) @property def pub_date(self): @@ -532,7 +531,6 @@ def create( handle_multiple=True, ): try: - logging.info(f"create: {pid_v3} {sps_pkg_name}") obj = cls() obj.pid_v3 = pid_v3 obj.sps_pkg_name = sps_pkg_name @@ -554,7 +552,6 @@ def create_or_update( sps_pkg_name=None, handle_multiple=False, ): - logging.info(f"Article.get_or_create: {user} {pid_v3} {sps_pkg_name}") try: return cls.get( pid_v3=pid_v3, @@ -761,7 +758,6 @@ def get_availability( params["collection"] = collection if collection_acron_list: params["collection__acron3__in"] = collection_acron_list - logging.info(f"get_availability {params}") return self.article_availability.filter(available=True, **params) def check_availability(self, user, force_update=False): @@ -770,7 +766,7 @@ def check_availability(self, user, force_update=False): return False if not force_update and self.is_available(): - return True + return self.mark_as_available() urls = [] for item in self.article_availability.all(): @@ -1729,14 +1725,16 @@ class StatusChoices(models.TextChoices): help_text=_("Related PID Provider XML instance"), ) detail = models.JSONField(null=True, blank=True, default=None) - am_article = models.ForeignKey( - AMArticle, + pid = models.CharField(max_length=24, null=True, blank=True) + collection = models.ForeignKey( + Collection, null=True, blank=True, on_delete=models.SET_NULL, - verbose_name=_("Legacy Article"), - help_text=_("Related Legacy Article instance"), + verbose_name=_("Collection"), + help_text=_("Related Collection instance"), ) + base_form_class = CoreAdminModelForm panels = [ @@ -1744,7 +1742,8 @@ class StatusChoices(models.TextChoices): FieldPanel("file", read_only=True), FieldPanel("source_date", read_only=True), FieldPanel("status"), - FieldPanel("am_article", read_only=True), + FieldPanel("collection", read_only=True), + FieldPanel("pid", read_only=True), FieldPanel("pid_provider_xml", read_only=True), FieldPanel("detail", read_only=True), ] @@ -1798,7 +1797,7 @@ def get(cls, url): raise ValueError("ArticleSource.get requires url") @classmethod - def create(cls, user, url=None, source_date=None, am_article=None, force_update=None, auto_solve_pid_conflict=False, is_public=None): + def create(cls, user, url=None, source_date=None, collection=None, pid=None, is_public=None, detail=None): if not url: raise ValueError("ArticleSource.create requires url") @@ -1807,55 +1806,81 @@ def create(cls, user, url=None, source_date=None, am_article=None, force_update= obj.creator = user obj.url = url obj.source_date = source_date - obj.am_article = am_article + obj.collection = collection + obj.pid = pid + obj.detail = detail if is_public is False: obj.status = cls.StatusChoices.NOT_PUBLIC else: obj.status = cls.StatusChoices.PENDING - obj.add_pid_provider(user, force_update, auto_solve_pid_conflict=auto_solve_pid_conflict) + obj.save() return obj except IntegrityError: return cls.get(url=url) @classmethod def create_or_update( - cls, user, url=None, source_date=None, am_article=None, force_update=None, auto_solve_pid_conflict=False, is_public=None + cls, user, url=None, source_date=None, collection=None, pid=None, force_update=None, auto_solve_pid_conflict=False, + is_public=None, detail=None, ): try: - logging.info( - f"ArticleSource.create_or_update {url} {source_date} {am_article} {force_update}" - ) - changed = False obj = cls.get(url=url) - if is_public is False: - obj.status = cls.StatusChoices.NOT_PUBLIC - changed = True - elif is_public is True and obj.status == cls.StatusChoices.NOT_PUBLIC: - obj.status = cls.StatusChoices.PENDING - changed = True - if ( - force_update - or (source_date and source_date != obj.source_date) - or not obj.is_completed - ): - obj.updated_by = user - obj.source_date = source_date - obj.am_article = am_article - obj.add_pid_provider(user, force_update, auto_solve_pid_conflict=auto_solve_pid_conflict) - changed = True - if changed: - obj.save() - return obj + changed = obj.update( + source_date, + collection, + pid, + is_public, + ) except cls.DoesNotExist: - return cls.create( + obj = cls.create( user, url=url, source_date=source_date, - am_article=am_article, - force_update=force_update, - auto_solve_pid_conflict=auto_solve_pid_conflict, + collection=collection, + pid=pid, is_public=is_public, + detail=detail, ) + changed = True + add_pid_provider_changed = obj.add_pid_provider(user, force_update, auto_solve_pid_conflict) + if changed or add_pid_provider_changed: + obj.updated_by = user + obj.save() + return obj + + def update( + self, + source_date, + collection, + pid, + is_public, + ): + changed = False + if self.source_date != source_date: + self.source_date = source_date + changed = True + if self.collection != collection: + self.collection = collection + changed = True + if self.pid != pid: + self.pid = pid + changed = True + if is_public is False and self.status != ArticleSource.StatusChoices.NOT_PUBLIC: + self.status = ArticleSource.StatusChoices.NOT_PUBLIC + changed = True + elif ( + is_public is True + and self.status == ArticleSource.StatusChoices.NOT_PUBLIC + ): + self.status = ArticleSource.StatusChoices.PENDING + changed = True + return changed + + def get_pid_provider_xml_id(self): + try: + return self.pid_provider_xml.id + except AttributeError: + pass @cached_property def xml_with_pre(self): @@ -1866,7 +1891,7 @@ def xml_with_pre(self): pass if self.file and self.file.path and os.path.isfile(self.file.path): try: - return XMLWithPre.from_file(self.file.path) + return list(XMLWithPre.create(path=self.file.path))[0] except Exception as e: pass if self.url: @@ -1882,11 +1907,10 @@ def sps_pkg_name(self): except Exception: pass - def request_xml(self, detail): + def request_xml(self): if not self.url: raise ValueError("URL is required") - logging.info(f"ArticleSource.request_xml for {self.url}") try: xml_with_pre = list(XMLWithPre.create(uri=self.url))[0] self.save_file( @@ -1905,7 +1929,7 @@ def save_file(self, filename, content): try: self.file.delete(save=False) except Exception as e: - logging.exception(e) + pass self.file.save(filename, ContentFile(content)) # Métodos para controle de status @@ -1917,22 +1941,18 @@ def mark_as_processing(self): def mark_as_completed(self): """Marca como concluído""" self.status = self.StatusChoices.COMPLETED - self.save() def mark_as_error(self): """Marca como erro""" self.status = self.StatusChoices.ERROR - self.save() def mark_as_url_error(self): """Marca como erro de URL""" self.status = self.StatusChoices.URL_ERROR - self.save() def mark_as_xml_error(self): """Marca como erro de XML""" self.status = self.StatusChoices.XML_ERROR - self.save() def mark_for_reprocess(self): """Marca para reprocessamento""" @@ -1972,56 +1992,6 @@ def get_needs_processing(cls): status__in=[cls.StatusChoices.PENDING, cls.StatusChoices.REPROCESS] ) - @classmethod - def get_queryset_to_complete_data( - cls, - from_date=None, - until_date=None, - force_update=None, - status_list=None, - params=None, - ): - params = params or {} - if status_list: - params["status__in"] = status_list - if from_date: - params["updated__gte"] = from_date - if until_date: - params["updated__lte"] = until_date - - if force_update: - return cls.objects.filter(**params) - - return cls.objects.filter( - Q(pid_provider_xml__isnull=True) | Q(file__isnull=True), - **params, - ) - - @property - def is_completed(self): - if not self.pid_provider_xml: - logging.info(f"Not completed: ArticleSource {self.url} has no pid_provider_xml") - return False - try: - if not self.pid_provider_xml.xml_with_pre: - logging.info(f"Not completed: ArticleSource {self.url} has pid_provider_xml but no xml_with_pre") - return False - except Exception: - pass - if not self.am_article: - logging.info(f"Not completed: ArticleSource {self.url} has no am_article") - return False - if not self.file: - logging.info(f"Not completed: ArticleSource {self.url} has no file") - return False - if not self.file.path or not os.path.isfile(self.file.path): - logging.info(f"Not completed: ArticleSource {self.url} has file path invalid or file does not exist") - return False - if self.status != ArticleSource.StatusChoices.COMPLETED: - self.status = ArticleSource.StatusChoices.COMPLETED - logging.info(f"Completed: ArticleSource {self.url} is completed") - return True - def add_pid_provider(self, user, force_update=False, auto_solve_pid_conflict=False): """ Executa o pipeline de obtenção de XML e registro de PID para este @@ -2038,117 +2008,71 @@ def add_pid_provider(self, user, force_update=False, auto_solve_pid_conflict=Fal pid_provider_xml), somente a etapa faltante é executada. """ try: - detail = [] - - if self.status == ArticleSource.StatusChoices.NOT_PUBLIC: - if not force_update: - return - - self.status = ArticleSource.StatusChoices.PENDING - - # --- Etapa 1: request_xml --- - has_valid_file = ( - self.file - and self.file.name - and os.path.isfile(self.file.path) - ) - - if force_update or not has_valid_file: - logging.info(f"Requesting XML for {self.url}") - self.request_xml(detail) - logging.info(f"XML requested successfully for {self.url}") - else: - logging.info( - f"Skipping request_xml: file already exists for {self.url}" - ) - detail.append("request_xml skipped (file already exists)") - - # --- Etapa 2: request_pid --- - has_pid_provider = self.pid_provider_xml is not None - - if force_update or not has_pid_provider: - logging.info(f"Requesting PID for {self.url}") - self.request_pid( - user, detail, force_update, auto_solve_pid_conflict - ) - logging.info( - f"PID requested successfully for {self.pid_provider_xml}" - ) - else: - logging.info( - f"Skipping request_pid: pid_provider_xml already set " - f"for {self.url}" - ) - detail.append("request_pid skipped (pid_provider_xml already set)") - - self.detail = detail - self.mark_as_completed() - logging.info(f"ArticleSource {self.status}") - - except XMLException as e: - exc_type, exc_value, exc_traceback = sys.exc_info() - detail.append(str({"error_type": str(type(e)), "error_message": str(e)})) - self.detail = detail - self.mark_as_xml_error() - logging.info(f"ArticleSource {self.url} marked as XML error") - except RequestXMLException as e: - exc_type, exc_value, exc_traceback = sys.exc_info() - detail.append(str({"error_type": str(type(e)), "error_message": str(e)})) - self.detail = detail - self.mark_as_url_error() - logging.info(f"ArticleSource {self.url} marked as URL error") + changed = False + try: + if self.status == ArticleSource.StatusChoices.NOT_PUBLIC: + if not force_update: + return changed + + if force_update or not self.file or not os.path.isfile(self.file.path): + # faz download do xml + self.request_xml() + changed = True + + if force_update or not self.pid_provider_xml: + # atribui pid_provider_xml + self.request_pid( + user, force_update, auto_solve_pid_conflict, + ) + changed = True + + if changed: + self.mark_as_completed() + except XMLException as e: + self.mark_as_xml_error() + raise + except RequestXMLException as e: + self.mark_as_url_error() + raise + except Exception as e: + self.mark_as_error() + raise except Exception as e: - logging.exception(e) exc_type, exc_value, exc_traceback = sys.exc_info() - detail.append(str({"error_type": str(type(e)), "error_message": str(e)})) - self.detail = detail - self.mark_as_error() - - def request_pid(self, user, detail, force_update, auto_solve_pid_conflict): - try: - detail.append("create pid_provider_xml") - - # Instancia o provedor de PIDs - pp = PidProvider() - - # Solicita PID para o arquivo XML/ZIP - logging.info(f"Requesting PID for {self.file.path}") - responses = pp.provide_pid_for_xml_zip( - self.file.path, - user, - filename=self.sps_pkg_name, - origin_date=self.source_date, - force_update=force_update, - is_published=True, - auto_solve_pid_conflict=auto_solve_pid_conflict, - ) + self.detail = self.detail or {} + self.detail.update({ + "error_type": str(type(e)), + "error_msg": str(e), + "traceback": traceback.format_exc() + }) + changed = True + return changed + + def request_pid(self, user, force_update, auto_solve_pid_conflict): + # Instancia o provedor de PIDs + pp = PidProvider() + + # Solicita PID para o arquivo XML/ZIP + responses = pp.provide_pid_for_xml_zip( + self.file.path, + user, + filename=self.sps_pkg_name, + origin_date=self.source_date, + force_update=force_update, + is_published=True, + auto_solve_pid_conflict=auto_solve_pid_conflict, + ) - # Obtém a primeira resposta (assumindo apenas uma) - response = list(responses)[0] - v3 = response.get("v3") - if v3: - # Associa o PidProviderXML ao ArticleSource - self.pid_provider_xml = PidProviderXML.get_by_pid_v3(v3) - if not self.pid_provider_xml: - raise UnableToRegisterPIDError("Failed to obtain or create PID v3") - detail.append("set pid_provider_xml") - else: - # Registra erro se não conseguiu obter v3 - detail.append(str(response)) - except Exception as e: - logging.exception(e) - exc_type, exc_value, exc_traceback = sys.exc_info() - unexpected_event = UnexpectedEvent.create( - exception=e, - exc_traceback=exc_traceback, - detail=dict( - function="article.models.ArticleSource.request_pid", - article_source_id=self.id, - url=self.url, - ), - ) - detail.append(str(unexpected_event.data)) - raise UnableToRegisterPIDError(str(e)) + # Obtém a primeira resposta (assumindo apenas uma) + response = list(responses)[0] + v3 = response.get("v3") + if v3: + # Associa o PidProviderXML ao ArticleSource + self.pid_provider_xml = PidProviderXML.get_by_pid_v3(v3) + if not self.pid_provider_xml: + raise UnableToRegisterPIDError(f"Unable to get by pid v3: {v3}") + else: + raise UnableToRegisterPIDError(response) class ArticleAvailability(CommonControlField): @@ -2194,8 +2118,8 @@ class Meta: # url já tem unique=True (cria índice automaticamente) @classmethod - def get(cls, article, url): - return cls.objects.get(article=article, url=url) + def get(cls, url): + return cls.objects.get(url=url) @classmethod def create( @@ -2221,7 +2145,7 @@ def create( obj.save() return obj except IntegrityError: - return cls.get(article, url) + return cls.get(url) @classmethod def create_or_update( @@ -2237,7 +2161,7 @@ def create_or_update( try: if lang: lang = Language.objects.filter(code2=lang).first() - obj = cls.get(article=article, url=url) + obj = cls.get(url=url) obj.fmt = fmt obj.lang = lang obj.collection = collection @@ -2937,6 +2861,17 @@ def __str__(self): parts.append(str(self.affiliation)) return " - ".join(parts) + @property + def data(self): + return dict( + article=self.article, + declared_name=self.declared_name, + orcid=self.orcid, + given_names=self.given_names, + last_name=self.last_name, + suffix=self.suffix + ) + def get_formatted_fullname(self, use_comma_separator=True, suffix_position="end"): """ Get formatted full name from name components. @@ -3078,7 +3013,7 @@ def create(cls, user, article, declared_name=None, given_names=None, if user: obj.creator = user - + try: obj.save() return obj @@ -3260,8 +3195,6 @@ def add_normalized_affiliation(self, user, organization=None, location=None, user=user, article=self.article ) - # Save to persist the relationship before using it - self.save() # Add normalized affiliation to the ArticleAffiliation self.affiliation.set_normalized( @@ -3272,7 +3205,6 @@ def add_normalized_affiliation(self, user, organization=None, location=None, level_2=level_2, level_3=level_3 ) - self.updated_by = user self.save() return self @@ -3341,7 +3273,6 @@ def create(cls, user, article, name, detail=None): obj.save() return obj except Exception as e: - logging.exception(f"Error creating ArticleEvent: {e}") raise EventSaveError(f"Unable to create article event: {e}") diff --git a/article/tasks.py b/article/tasks.py index bea451cd4..bf3dd356c 100644 --- a/article/tasks.py +++ b/article/tasks.py @@ -5,6 +5,7 @@ from django.utils.translation import gettext_lazy as _ from article import controller +from article.controller import ArticleIteratorBuilder from article.models import Article, ArticleFormat, ArticleSource, AMArticle from article.sources.preprint import harvest_preprints from article.sources.xmlsps import load_article @@ -128,7 +129,6 @@ def task_convert_xml_to_other_formats_for_articles( user = _get_user(self.request, username, user_id) for item in Article.objects.filter(sps_pkg_name__isnull=False).iterator(): - logging.info(item.pid_v3) try: convert_xml_to_other_formats.apply_async( kwargs={ @@ -204,7 +204,6 @@ def convert_xml_to_other_formats( done = True except ArticleFormat.DoesNotExist: done = False - logging.info(f"Done {done}") if not done or force_update: ArticleFormat.generate_formats(user, article=article) @@ -250,7 +249,7 @@ def transfer_license_statements_fk_to_article_license( if not instance.license and first.data: data = first.data instance.license = License.create_or_update(user, license_type=data.get("license_type"), version=data.get("license_version")) - + if not instance.license: continue instance.updated_by = user @@ -260,7 +259,6 @@ def transfer_license_statements_fk_to_article_license( Article.objects.bulk_update( articles_to_update, ["license", "updated_by"] ) - logging.info("The license of model Articles have been updated") def get_researcher_identifier_unnormalized(): @@ -303,7 +301,6 @@ def normalize_stored_email( - Identifica e-mails com formato inválido usando regex - Aplica normalização através de extracts_normalized_email - Executa bulk_update para otimizar performance em lotes - - Registra logs de processamento Examples: # Executar normalização de e-mails @@ -372,7 +369,6 @@ def task_export_articles_to_articlemeta( Side Effects: - Exporta múltiplos artigos para ArticleMeta - Atualiza status de exportação dos artigos - - Registra logs de processamento - Registra UnexpectedEvent em caso de erro Examples: @@ -407,17 +403,18 @@ def task_export_articles_to_articlemeta( days_to_go_back=days_to_go_back, force_update=force_update, ) - + return result - + except Exception as e: exc_type, exc_value, exc_traceback = sys.exc_info() - + UnexpectedEvent.create( + action="task_export_articles_to_articlemeta", + item="", exception=e, exc_traceback=exc_traceback, detail={ - "task": "task_export_articles_to_articlemeta", "collection_acron_list": collection_acron_list, "journal_acron_list": journal_acron_list, "year_of_publication": year_of_publication, @@ -432,7 +429,7 @@ def task_export_articles_to_articlemeta( "task_id": self.request.id if hasattr(self.request, 'id') else None, }, ) - + # Re-raise para que o Celery possa tratar a exceção adequadamente raise @@ -466,7 +463,6 @@ def task_export_article_to_articlemeta( Side Effects: - Exporta artigo específico para ArticleMeta - Atualiza status de exportação do artigo - - Registra logs de processamento - Registra UnexpectedEvent em caso de erro Raises: @@ -484,8 +480,9 @@ def task_export_article_to_articlemeta( - Utiliza controller.export_article_to_articlemeta internamente - Requer que o artigo exista na base local antes da exportação """ + item = pid_v3 or "" + try: - item = pid_v3 if not pid_v3: raise ValueError("task_export_article_to_articlemeta requires pid_v3") @@ -506,6 +503,7 @@ def task_export_article_to_articlemeta( ) except Article.DoesNotExist as exception: return False + except Exception as exception: exc_type, exc_value, exc_traceback = sys.exc_info() UnexpectedEvent.create( @@ -705,6 +703,82 @@ def task_check_article_availability( ) +@celery_app.task(bind=True) +def task_harvest_articles( + self, + username=None, + user_id=None, + collection_acron_list=None, + journal_acron_list=None, + from_date=None, + until_date=None, + force_update=None, + export_to_articlemeta=False, + auto_solve_pid_conflict=None, + limit=None, + timeout=None, + opac_url=None, + stop=None, +): + item = "" + params = { + "collection_acron_list": collection_acron_list, + "journal_acron_list": journal_acron_list, + "from_date": from_date, + "until_date": until_date, + "force_update": force_update, + "export_to_articlemeta": export_to_articlemeta, + "auto_solve_pid_conflict": auto_solve_pid_conflict, + "limit": limit, + "timeout": timeout, + "opac_url": opac_url, + "stop": stop, + } + + try: + items = (collection_acron_list or []) + (journal_acron_list or []) + item = "-".join(items) + user = _get_user(self.request, username=username, user_id=user_id) + + common_kwargs = { + "user_id": user.id, + "username": user.username, + "force_update": force_update, + "export_to_articlemeta": export_to_articlemeta, + "auto_solve_pid_conflict": auto_solve_pid_conflict, + } + + builder = ArticleIteratorBuilder( + user=user, + collection_acron_list=collection_acron_list, + journal_acron_list=journal_acron_list, + from_date=from_date, + until_date=until_date, + force_update=force_update, + limit=limit, + timeout=timeout, + opac_url=opac_url, + stop=stop, + ) + item_iterator = builder.from_harvest() + + for item_kwargs in item_iterator: + if item_kwargs is None: + continue + task_process_article_pipeline.delay(**item_kwargs, **common_kwargs) + + except Exception as e: + exc_type, exc_value, exc_traceback = sys.exc_info() + UnexpectedEvent.create( + action="task_harvest_articles", + item=item, + exception=e, + exc_traceback=exc_traceback, + detail=params, + ) + raise + + @celery_app.task(bind=True) def task_dispatch_articles( self, @@ -724,62 +798,28 @@ def task_dispatch_articles( proc_status_list=None, # --- ativa article --- data_status_list=None, - # --- ativa harvest (qualquer um) --- - limit=None, - timeout=None, - opac_url=None, # --- ativa article_source --- article_source_status_list=None, - verify=None, - stop=None, ): - """ - Tarefa orquestradora que dispara processamento em lote de artigos. - - Utiliza ArticleIteratorBuilder para selecionar artigos baseado em - múltiplos critérios e dispara task_process_article_pipeline para - cada item encontrado, permitindo processamento paralelo. + item = "" + params = { + "collection_acron_list": collection_acron_list, + "journal_acron_list": journal_acron_list, + "from_pub_year": from_pub_year, + "until_pub_year": until_pub_year, + "from_date": from_date, + "until_date": until_date, + "force_update": force_update, + "export_to_articlemeta": export_to_articlemeta, + "auto_solve_pid_conflict": auto_solve_pid_conflict, + "proc_status_list": proc_status_list, + "data_status_list": data_status_list, + "article_source_status_list": article_source_status_list, + } - Args: - self: Instância da tarefa Celery - username (str, optional): Nome do usuário executando a tarefa - user_id (int, optional): ID do usuário executando a tarefa - collection_acron_list (list, optional): Filtro por acrônimos de coleções - journal_acron_list (list, optional): Filtro por acrônimos de periódicos - from_pub_year (int, optional): Ano inicial de publicação - until_pub_year (int, optional): Ano final de publicação - from_date (str, optional): Data inicial (formato ISO) - until_date (str, optional): Data final (formato ISO) - force_update (bool, optional): Força reprocessamento - export_to_articlemeta (bool): Exporta para ArticleMeta após processamento - auto_solve_pid_conflict (bool, optional): Resolve conflitos de PID automaticamente - proc_status_list (list, optional): Status do pid_provider para filtro - data_status_list (list, optional): Status do article para filtro - limit (int, optional): Limite máximo de artigos a processar - timeout (int, optional): Timeout para operações HTTP - opac_url (str, optional): URL base do OPAC para harvest - article_source_status_list (list, optional): Status do article_source para filtro - - Returns: - dict: Resumo com contadores de dispatched/skipped - - Examples: - # Processamento padrão por coleção - task_dispatch_articles.delay(collection_acron_list=["scl"]) - - # Múltiplas fontes simultaneamente - task_dispatch_articles.delay( - proc_status_list=["todo"], - data_status_list=["invalid"], - article_source_status_list=["error"], - limit=500 - ) - - Notes: - - Ver ArticleIteratorBuilder para detalhes sobre iteradores ativados - - Cada artigo encontrado gera uma subtarefa independente - """ try: + items = (collection_acron_list or []) + (journal_acron_list or []) + item = "-".join(items) user = _get_user(self.request, username=username, user_id=user_id) common_kwargs = { @@ -790,9 +830,7 @@ def task_dispatch_articles( "auto_solve_pid_conflict": auto_solve_pid_conflict, } - dispatched = skipped = 0 - - for item_kwargs in controller.ArticleIteratorBuilder( + builder = ArticleIteratorBuilder( user=user, collection_acron_list=collection_acron_list, journal_acron_list=journal_acron_list, @@ -800,50 +838,37 @@ def task_dispatch_articles( until_pub_year=until_pub_year, from_date=from_date, until_date=until_date, - proc_status_list=proc_status_list, - data_status_list=data_status_list, - article_source_status_list=article_source_status_list, - limit=limit, - timeout=timeout, - opac_url=opac_url, force_update=force_update, - stop=stop - ): - if item_kwargs is None: - skipped += 1 - continue - logging.info(f"Dispatching article with kwargs: {item_kwargs}") - task_process_article_pipeline.delay(**item_kwargs, **common_kwargs) - dispatched += 1 + ) - return { - "status": "success", - "dispatched": dispatched, - "skipped": skipped, - } + item_iterators = ( + builder.from_article_source( + article_source_status_list=article_source_status_list + ), + builder.from_pid_provider(proc_status_list=proc_status_list), + builder.from_article(data_status_list=data_status_list) + ) + + for item_iterator in item_iterators: + if not item_iterator: + continue + for item_kwargs in item_iterator: + if item_kwargs is None: + continue + task_process_article_pipeline.delay(**item_kwargs, **common_kwargs) except Exception as e: exc_type, exc_value, exc_traceback = sys.exc_info() UnexpectedEvent.create( + action="task_dispatch_articles", + item=item, exception=e, exc_traceback=exc_traceback, - detail={ - "task": "task_dispatch_articles", - "collection_acron_list": collection_acron_list, - "journal_acron_list": journal_acron_list, - "from_pub_year": from_pub_year, - "until_pub_year": until_pub_year, - "from_date": from_date, - "until_date": until_date, - "proc_status_list": proc_status_list, - "data_status_list": data_status_list, - "article_source_status_list": article_source_status_list, - "force_update": force_update, - "export_to_articlemeta": export_to_articlemeta, - }, + detail=params, ) raise + @celery_app.task(bind=True) def task_process_article_pipeline( self, @@ -865,104 +890,35 @@ def task_process_article_pipeline( user_id=None, username=None, is_public=None, + document=None, ): - """ - Pipeline principal de processamento de artigos com múltiplos pontos de entrada. - - Implementa um pipeline flexível que pode iniciar em diferentes estágios: - - Fluxo A: XML URL → ArticleSource → PidProviderXML → Article - - Fluxo B: ArticleSource existente → PidProviderXML → Article - - Fluxo C: PidProviderXML → Article (entrada direta) - - Args: - self: Instância da tarefa Celery - xml_url (str, optional): URL do XML para fluxo A (requer collection_acron e pid) - collection_acron (str, optional): Acrônimo da coleção (obrigatório com xml_url) - pid (str, optional): PID do artigo (obrigatório com xml_url) - source_date (datetime, optional): Data da fonte para fluxo A - article_source_id (int, optional): ID do ArticleSource para fluxo B - pp_xml_id (int, optional): ID do PidProviderXML para fluxo C - export_to_articlemeta (bool): Se True, exporta para ArticleMeta após processamento - collection_acron_list (list, optional): Lista de coleções para exportação - force_update (bool, optional): Força reprocessamento mesmo se existir - auto_solve_pid_conflict (bool, optional): Resolve conflitos de PID automaticamente - version (str, optional): Versão específica a processar - user_id (int, optional): ID do usuário executando a tarefa - username (str, optional): Nome do usuário executando a tarefa - - Returns: - None - - Side Effects: - - Cria/atualiza ArticleSource (fluxo A) - - Cria/atualiza PidProviderXML - - Cria/atualiza Article - - Verifica disponibilidade do artigo - - Exporta para ArticleMeta se solicitado - - Registra UnexpectedEvent em caso de erro - - Raises: - ValueError: Se nenhum ponto de entrada válido for fornecido - Se xml_url fornecido sem collection_acron ou pid - - Examples: - # Fluxo completo a partir de URL - task_process_article_pipeline.delay( - xml_url="http://example.com/article.xml", - collection_acron="scl", - pid="S1234-56782024000100001", - export_to_articlemeta=True - ) - - # A partir de ArticleSource existente - task_process_article_pipeline.delay( - article_source_id=123, - force_update=True - ) - - # Entrada direta via PidProviderXML - task_process_article_pipeline.delay( - pp_xml_id=456, - export_to_articlemeta=True - ) - """ try: - unexpected_event_item = None + unexpected_event_item = xml_url user = _get_user(self.request, username=username, user_id=user_id) - if xml_url: - unexpected_event_item = xml_url + + article_source = None + if article_source_id: + article_source = ArticleSource.objects.get(id=article_source_id) + elif xml_url: if not collection_acron: raise ValueError("collection_acron is required when xml_url is provided") if not pid: raise ValueError("pid is required when xml_url is provided") - am_article = AMArticle.create_or_update( - pid, Collection.get(collection_acron), None, user - ) - if not am_article: - raise ValueError( - f"Failed to create or update AMArticle with pid: {pid} and collection: {collection_acron}" - ) - article_source = ArticleSource.create_or_update( user=user, url=xml_url, source_date=source_date, + collection=Collection.get(collection_acron), + pid=pid, force_update=force_update, - am_article=am_article, auto_solve_pid_conflict=auto_solve_pid_conflict, is_public=is_public, + detail=document, ) - pp_xml_id = article_source.pid_provider_xml.id - - if article_source_id: - article_source = ArticleSource.objects.get(id=article_source_id) + + if article_source: unexpected_event_item = str(article_source) - article_source.add_pid_provider( - user=user, - force_update=force_update, - auto_solve_pid_conflict=auto_solve_pid_conflict, - ) - pp_xml_id = article_source.pid_provider_xml.id + pp_xml_id = article_source.get_pid_provider_xml_id() if not pp_xml_id: raise ValueError( @@ -981,10 +937,9 @@ def task_process_article_pipeline( pp_xml.collections.set(article.collections) article.check_availability(user, force_update=export_to_articlemeta or force_update) - + if export_to_articlemeta: if not article.is_classic_public or not article.valid: - logging.warning(f"Article {article.pid_v3} is not valid or not public. Skipping export to ArticleMeta.") return task_export_article_to_articlemeta.delay( pid_v3=article.pid_v3, @@ -1006,11 +961,8 @@ def task_process_article_pipeline( "pp_xml_id": pp_xml_id, "pid": pid, "collection_acron": collection_acron, - "source_date": source_date, - "collection_acron_list": collection_acron_list, - "auto_solve_pid_conflict": auto_solve_pid_conflict, - "version": version, "export_to_articlemeta": export_to_articlemeta, "force_update": force_update, }, ) + raise diff --git a/article/tests/test_article_iterator_builder.py b/article/tests/test_article_iterator_builder.py new file mode 100644 index 000000000..2651ff128 --- /dev/null +++ b/article/tests/test_article_iterator_builder.py @@ -0,0 +1,690 @@ +import logging +import unittest +from unittest.mock import MagicMock, patch + +from django.db.models import Q + +# Ajuste este caminho se a classe estiver em outro módulo. +MODULE_PATH = "article.controller" +from article.controller import ArticleIteratorBuilder # noqa: E402 + + +def setUpModule(): + """ + Desliga o logging (abaixo de CRITICAL) para todo o módulo de testes. + + Sem isso, logging.info/logging.error chamados dentro de + ArticleIteratorBuilder alcançam os handlers reais configurados no + settings.LOGGING do Django (ex.: handler que despacha para o + OpenSearch de forma assíncrona) — em ambiente de teste esse host + normalmente não existe/não está acessível, gerando ruído de + "NameResolutionError" / "--- Logging error ---" no output dos testes. + """ + logging.disable(logging.CRITICAL) + + +def tearDownModule(): + logging.disable(logging.NOTSET) + + +# ======================================================================= +# Helpers +# ======================================================================= +def make_builder(**kwargs): + defaults = dict( + user=MagicMock(name="user"), + collection_acron_list=None, + journal_acron_list=None, + from_pub_year=None, + until_pub_year=None, + from_date=None, + until_date=None, + force_update=None, + limit=None, + timeout=None, + opac_url=None, + ) + defaults.update(kwargs) + return ArticleIteratorBuilder(**defaults) + + +class FakeValuesQuerySet: + """ + Simula o resultado final de .values(...)/.values(...=F(...)) — os + itens já são os dicts finais, e .iterator() apenas os devolve. + """ + + def __init__(self, items): + self.items = list(items) + + def values(self, *args, **kwargs): + return self + + def distinct(self): + return self + + def order_by(self, *args, **kwargs): + return self + + def iterator(self): + return iter(self.items) + + +class FakeArticleBaseQuerySet: + """ + Simula o `base_qs` de from_article: .filter(pp_xml__isnull=False) + registra os kwargs recebidos (para asserção) e devolve a própria + instância, permitindo encadear .values(...).iterator(). + """ + + def __init__(self, items=None): + self._items = list(items or []) + self.received_filter_kwargs = None + + def filter(self, **kwargs): + self.received_filter_kwargs = kwargs + return self + + def values(self, *args, **kwargs): + return self + + def iterator(self): + return iter(self._items) + + +# ======================================================================= +# __init__ +# ======================================================================= +class TestInit(unittest.TestCase): + def test_stores_common_attributes(self): + user = MagicMock(name="user") + builder = ArticleIteratorBuilder( + user=user, + collection_acron_list=["scl"], + journal_acron_list=["abc"], + from_pub_year=2020, + until_pub_year=2022, + from_date="2020-01-01", + until_date="2022-12-31", + force_update=True, + limit=10, + timeout=30, + opac_url="www.custom.br", + ) + self.assertIs(builder.user, user) + self.assertEqual(builder.collection_acron_list, ["scl"]) + self.assertEqual(builder.journal_acron_list, ["abc"]) + self.assertEqual(builder.from_pub_year, 2020) + self.assertEqual(builder.until_pub_year, 2022) + self.assertEqual(builder.from_date, "2020-01-01") + self.assertEqual(builder.until_date, "2022-12-31") + self.assertTrue(builder.force_update) + self.assertEqual(builder.limit, 10) + self.assertEqual(builder.timeout, 30) + self.assertEqual(builder.opac_url, "www.custom.br") + + def test_defaults_are_none(self): + builder = ArticleIteratorBuilder(user=MagicMock()) + self.assertIsNone(builder.collection_acron_list) + self.assertIsNone(builder.journal_acron_list) + self.assertIsNone(builder.from_pub_year) + self.assertIsNone(builder.until_pub_year) + self.assertIsNone(builder.from_date) + self.assertIsNone(builder.until_date) + self.assertIsNone(builder.force_update) + self.assertIsNone(builder.limit) + self.assertIsNone(builder.timeout) + self.assertIsNone(builder.opac_url) + + def test_source_specific_filters_are_not_instance_attributes(self): + """ + proc_status_list / data_status_list / article_source_status_list + são exclusivos de cada fonte e não devem existir como atributo de + instância (ficam só como parâmetro do método correspondente). + """ + builder = make_builder() + self.assertFalse(hasattr(builder, "proc_status_list")) + self.assertFalse(hasattr(builder, "data_status_list")) + self.assertFalse(hasattr(builder, "article_source_status_list")) + + +# ======================================================================= +# from_pid_provider +# ======================================================================= +class TestFromPidProvider(unittest.TestCase): + def test_yields_dicts_directly_from_queryset(self): + builder = make_builder() + items = [{"pp_xml_id": 1}, {"pp_xml_id": 2}] + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX: + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet(items) + + result = list(builder.from_pid_provider()) + + self.assertEqual(result, items) + + def test_default_proc_status_used_when_not_provided(self): + builder = make_builder() + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX, \ + patch(f"{MODULE_PATH}.pid_provider_choices") as MockChoices: + MockChoices.PPXML_STATUS_TODO = "todo" + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet([]) + + list(builder.from_pid_provider()) + + _, kwargs = MockPPX.objects.filter.call_args + self.assertEqual(kwargs["proc_status__in"], ["todo"]) + + def test_custom_proc_status_list_used_when_provided(self): + builder = make_builder() + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX: + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet([]) + + list(builder.from_pid_provider(proc_status_list=["custom"])) + + _, kwargs = MockPPX.objects.filter.call_args + self.assertEqual(kwargs["proc_status__in"], ["custom"]) + + def test_date_and_year_filters_applied(self): + builder = make_builder( + from_pub_year=2020, + until_pub_year=2022, + from_date="2020-01-01", + until_date="2022-12-31", + ) + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX: + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet([]) + + list(builder.from_pid_provider()) + + _, kwargs = MockPPX.objects.filter.call_args + self.assertEqual(kwargs["pub_year__gte"], 2020) + self.assertEqual(kwargs["pub_year__lte"], 2022) + self.assertEqual(kwargs["updated__gte"], "2020-01-01") + self.assertEqual(kwargs["updated__lte"], "2022-12-31") + + def test_no_collection_or_journal_filter_skips_scielojournal_query(self): + """ + Sem collection_acron_list nem journal_acron_list, `params` fica + vazio e SciELOJournal não deve ser consultado; `q` permanece um + Q() vazio. + """ + builder = make_builder() + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet([]) + + list(builder.from_pid_provider()) + + MockSciELOJournal.objects.select_related.assert_not_called() + (q_arg,), _ = MockPPX.objects.filter.call_args + self.assertFalse(q_arg) # Q() vazio é falsy + + def test_collection_acron_list_alone_queries_scielojournal(self): + builder = make_builder(collection_acron_list=["scl"]) + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + sj_mock = MockSciELOJournal.objects.select_related.return_value + sj_mock.filter.return_value.values_list.return_value.distinct.return_value = [ + ("1234-5678", "8765-4321"), + ] + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet([{"pp_xml_id": 1}]) + + result = list(builder.from_pid_provider()) + + sj_mock.filter.assert_called_once_with(collection__acron3__in=["scl"]) + self.assertEqual(result, [{"pp_xml_id": 1}]) + + def test_journal_acron_list_alone_queries_scielojournal(self): + builder = make_builder(journal_acron_list=["abc"]) + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + sj_mock = MockSciELOJournal.objects.select_related.return_value + sj_mock.filter.return_value.values_list.return_value.distinct.return_value = [ + ("1234-5678", None), + ] + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet([]) + + list(builder.from_pid_provider()) + + sj_mock.filter.assert_called_once_with(journal_acron__in=["abc"]) + + def test_collection_and_journal_combined_in_single_filter_call(self): + builder = make_builder( + collection_acron_list=["scl"], journal_acron_list=["abc"] + ) + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + sj_mock = MockSciELOJournal.objects.select_related.return_value + sj_mock.filter.return_value.values_list.return_value.distinct.return_value = [] + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet([]) + + list(builder.from_pid_provider()) + + sj_mock.filter.assert_called_once_with( + collection__acron3__in=["scl"], journal_acron__in=["abc"] + ) + + def test_issn_list_filters_out_falsy_values(self): + builder = make_builder(journal_acron_list=["abc"]) + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + sj_mock = MockSciELOJournal.objects.select_related.return_value + # Um periódico só com issn_print, outro só com issn_electronic, + # outro com ambos nulos (deve ser ignorado). + sj_mock.filter.return_value.values_list.return_value.distinct.return_value = [ + ("1111-1111", None), + (None, "2222-2222"), + (None, None), + ] + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet([]) + + list(builder.from_pid_provider()) + + (q_arg,), _ = MockPPX.objects.filter.call_args + q_str = str(q_arg) + self.assertIn("1111-1111", q_str) + self.assertIn("2222-2222", q_str) + + def test_no_issn_found_still_queries_pidproviderxml(self): + """ + Diferente de uma implementação anterior: aqui NÃO há early-return + quando SciELOJournal não encontra nenhum ISSN — o código monta um + Q(issn_print__in=set())|Q(issn_electronic__in=set()) (que não + casa com nada) e segue para PidProviderXML.objects.filter mesmo + assim. + """ + builder = make_builder(journal_acron_list=["nao-existe"]) + + with patch(f"{MODULE_PATH}.PidProviderXML") as MockPPX, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + sj_mock = MockSciELOJournal.objects.select_related.return_value + sj_mock.filter.return_value.values_list.return_value.distinct.return_value = [] + ( + MockPPX.objects.filter.return_value + .order_by.return_value + .values.return_value + .distinct.return_value + ) = FakeValuesQuerySet([]) + + list(builder.from_pid_provider()) + + MockPPX.objects.filter.assert_called_once() + + +# ======================================================================= +# from_article +# ======================================================================= +class TestFromArticle(unittest.TestCase): + def test_no_data_status_list_means_no_filter_key(self): + builder = make_builder() + base_qs = FakeArticleBaseQuerySet() + + with patch(f"{MODULE_PATH}.Article") as MockArticle: + MockArticle.objects.filter.return_value.distinct.return_value = base_qs + + list(builder.from_article()) + + _, kwargs = MockArticle.objects.filter.call_args + self.assertNotIn("data_status__in", kwargs) + + def test_data_status_list_included_when_provided(self): + builder = make_builder() + base_qs = FakeArticleBaseQuerySet() + + with patch(f"{MODULE_PATH}.Article") as MockArticle: + MockArticle.objects.filter.return_value.distinct.return_value = base_qs + + list(builder.from_article(data_status_list=["pending", "invalid"])) + + _, kwargs = MockArticle.objects.filter.call_args + self.assertEqual(kwargs["data_status__in"], ["pending", "invalid"]) + + def test_valid_false_always_applied(self): + builder = make_builder() + base_qs = FakeArticleBaseQuerySet() + + with patch(f"{MODULE_PATH}.Article") as MockArticle: + MockArticle.objects.filter.return_value.distinct.return_value = base_qs + + list(builder.from_article()) + + _, kwargs = MockArticle.objects.filter.call_args + self.assertIs(kwargs["valid"], False) + + def test_collection_and_journal_filters_traverse_scielojournal_join(self): + builder = make_builder( + collection_acron_list=["scl"], journal_acron_list=["abc"] + ) + base_qs = FakeArticleBaseQuerySet() + + with patch(f"{MODULE_PATH}.Article") as MockArticle: + MockArticle.objects.filter.return_value.distinct.return_value = base_qs + + list(builder.from_article()) + + _, kwargs = MockArticle.objects.filter.call_args + self.assertEqual( + kwargs["journal__scielojournal__collection__acron3__in"], ["scl"] + ) + self.assertEqual( + kwargs["journal__scielojournal__journal_acron__in"], ["abc"] + ) + + def test_pub_year_and_date_filters_applied(self): + builder = make_builder( + from_pub_year=2019, + until_pub_year=2021, + from_date="2019-01-01", + until_date="2021-12-31", + ) + base_qs = FakeArticleBaseQuerySet() + + with patch(f"{MODULE_PATH}.Article") as MockArticle: + MockArticle.objects.filter.return_value.distinct.return_value = base_qs + + list(builder.from_article()) + + _, kwargs = MockArticle.objects.filter.call_args + self.assertEqual(kwargs["pub_date_year__gte"], 2019) + self.assertEqual(kwargs["pub_date_year__lte"], 2021) + self.assertEqual(kwargs["updated__gte"], "2019-01-01") + self.assertEqual(kwargs["updated__lte"], "2021-12-31") + + def test_only_articles_with_pp_xml_are_yielded(self): + builder = make_builder() + items = [{"pp_xml_id": 10}, {"pp_xml_id": 20}] + base_qs = FakeArticleBaseQuerySet(items) + + with patch(f"{MODULE_PATH}.Article") as MockArticle: + MockArticle.objects.filter.return_value.distinct.return_value = base_qs + + result = list(builder.from_article()) + + self.assertEqual(base_qs.received_filter_kwargs, {"pp_xml__isnull": False}) + self.assertEqual(result, items) + + +# ======================================================================= +# from_article_source +# ======================================================================= +class TestFromArticleSource(unittest.TestCase): + def test_yields_dicts_directly_from_queryset(self): + builder = make_builder(from_date="d1", until_date="d2", force_update=False) + items = [{"article_source_id": 10}, {"article_source_id": 20}] + + with patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource, \ + patch(f"{MODULE_PATH}.pid_provider_choices") as MockChoices: + MockChoices.PPXML_STATUS_TO_CREATE_OR_UPDATE_ARTICLE_SOURCE = "to_create" + ( + MockArticleSource.objects.filter.return_value.values.return_value + ) = FakeValuesQuerySet(items) + + result = list( + builder.from_article_source(article_source_status_list=["pending"]) + ) + + self.assertEqual(result, items) + + def test_force_update_uses_empty_q(self): + builder = make_builder(force_update=True) + + with patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource: + ( + MockArticleSource.objects.filter.return_value.values.return_value + ) = FakeValuesQuerySet([]) + + list(builder.from_article_source()) + + (q_arg,), _ = MockArticleSource.objects.filter.call_args + self.assertFalse(q_arg) + + +# ======================================================================= +# from_harvest +# ======================================================================= +class TestFromHarvest(unittest.TestCase): + def test_loads_collection_when_empty(self): + builder = make_builder() + + with patch(f"{MODULE_PATH}.Collection") as MockCollection, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + MockCollection.objects.count.return_value = 0 + ( + MockSciELOJournal.objects.select_related.return_value + .filter.return_value + .values_list.return_value + .distinct.return_value + ) = [] + + list(builder.from_harvest()) + + MockCollection.load.assert_called_once_with(builder.user) + + def test_does_not_load_collection_when_not_empty(self): + builder = make_builder() + + with patch(f"{MODULE_PATH}.Collection") as MockCollection, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + MockCollection.objects.count.return_value = 5 + ( + MockSciELOJournal.objects.select_related.return_value + .filter.return_value + .values_list.return_value + .distinct.return_value + ) = [] + + list(builder.from_harvest()) + + MockCollection.load.assert_not_called() + + def test_collection_and_journal_filters_passed_to_scielojournal_query(self): + builder = make_builder(collection_acron_list=["scl"], journal_acron_list=["abc"]) + + with patch(f"{MODULE_PATH}.Collection") as MockCollection, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + MockCollection.objects.count.return_value = 5 + sj_mock = MockSciELOJournal.objects.select_related.return_value + sj_mock.filter.return_value.values_list.return_value.distinct.return_value = [] + + list(builder.from_harvest()) + + sj_mock.filter.assert_called_once_with( + collection__acron3__in=["scl"], journal_acron__in=["abc"] + ) + + def test_builds_harvester_per_item_with_positional_args(self): + builder = make_builder() + + with patch(f"{MODULE_PATH}.Collection") as MockCollection, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + MockCollection.objects.count.return_value = 5 + ( + MockSciELOJournal.objects.select_related.return_value + .filter.return_value + .values_list.return_value + .distinct.return_value + ) = [("scl", "abc", "1234-5678")] + + fake_harvester = MagicMock() + fake_harvester.harvest_documents.return_value = [] + with patch.object(builder, "_build_harvester", return_value=fake_harvester) as mock_build: + list(builder.from_harvest()) + + mock_build.assert_called_once_with("scl", "abc", "1234-5678") + + def test_yields_expected_dict_from_documents(self): + builder = make_builder() + doc = { + "url": "http://x/y.xml", + "pid_v2": "S123", + "processing_date": "2024-01-01", + "is_public": True, + } + + with patch(f"{MODULE_PATH}.Collection") as MockCollection, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + MockCollection.objects.count.return_value = 5 + ( + MockSciELOJournal.objects.select_related.return_value + .filter.return_value + .values_list.return_value + .distinct.return_value + ) = [("scl", "abc", "1234-5678")] + + fake_harvester = MagicMock() + fake_harvester.harvest_documents.return_value = [doc] + with patch.object(builder, "_build_harvester", return_value=fake_harvester): + result = list(builder.from_harvest()) + + self.assertEqual(result, [{ + "xml_url": "http://x/y.xml", + "collection_acron": "scl", + "pid": "S123", + "source_date": "2024-01-01", + "is_public": True, + "document": { + "url": "http://x/y.xml", + "pid_v2": "S123", + "processing_date": "2024-01-01", + "is_public": True, + }, + }]) + + def test_uses_origin_date_when_processing_date_absent(self): + builder = make_builder() + doc = { + "url": "http://x/y.xml", + "pid_v2": "S123", + "origin_date": "2023-05-05", + "is_public": False, + } + + with patch(f"{MODULE_PATH}.Collection") as MockCollection, \ + patch(f"{MODULE_PATH}.SciELOJournal") as MockSciELOJournal: + MockCollection.objects.count.return_value = 5 + ( + MockSciELOJournal.objects.select_related.return_value + .filter.return_value + .values_list.return_value + .distinct.return_value + ) = [("scl", "abc", "1234-5678")] + + fake_harvester = MagicMock() + fake_harvester.harvest_documents.return_value = [doc] + with patch.object(builder, "_build_harvester", return_value=fake_harvester): + result = list(builder.from_harvest()) + + self.assertEqual(result[0]["source_date"], "2023-05-05") + + +# ======================================================================= +# _build_harvester +# ======================================================================= +class TestBuildHarvester(unittest.TestCase): + def test_scl_collection_uses_opac_harvester_with_default_url(self): + builder = make_builder( + opac_url=None, from_date="a", until_date="b", limit=10, timeout=30 + ) + + with patch(f"{MODULE_PATH}.OPACHarvester") as MockOPACHarvester: + builder._build_harvester("scl") + + MockOPACHarvester.assert_called_once_with( + "https://www.scielo.br", "scl", from_date="a", until_date="b", + limit=10, timeout=30, + ) + + def test_scl_collection_uses_provided_opac_url(self): + builder = make_builder(opac_url="www.custom.br") + + with patch(f"{MODULE_PATH}.OPACHarvester") as MockOPACHarvester: + builder._build_harvester("scl") + + args, _ = MockOPACHarvester.call_args + self.assertEqual(args[0], "www.custom.br") + + def test_non_scl_collection_uses_am_harvester(self): + builder = make_builder(from_date="a", until_date="b", limit=5, timeout=15) + + with patch(f"{MODULE_PATH}.AMHarvester") as MockAMHarvester: + builder._build_harvester("mex") + + MockAMHarvester.assert_called_once_with( + "article", "mex", from_date="a", until_date="b", limit=5, timeout=15 + ) + + def test_non_scl_collection_with_journal_id_passes_journal_kwarg(self): + builder = make_builder() + + with patch(f"{MODULE_PATH}.AMHarvester") as MockAMHarvester: + builder._build_harvester("mex", journal_id="1234-5678") + + _, kwargs = MockAMHarvester.call_args + self.assertEqual(kwargs["journal"], "1234-5678") + + def test_non_scl_collection_without_journal_id_has_no_journal_kwarg(self): + builder = make_builder() + + with patch(f"{MODULE_PATH}.AMHarvester") as MockAMHarvester: + builder._build_harvester("mex") + + _, kwargs = MockAMHarvester.call_args + self.assertNotIn("journal", kwargs) + + +if __name__ == "__main__": + unittest.main() \ No newline at end of file diff --git a/article/tests/test_article_pipeline.py b/article/tests/test_article_pipeline.py new file mode 100644 index 000000000..3ef7d426f --- /dev/null +++ b/article/tests/test_article_pipeline.py @@ -0,0 +1,511 @@ +import logging +import unittest +from unittest.mock import MagicMock, patch, call + +from article.models import Article + +# Ajuste este caminho se as tasks estiverem em outro módulo. +MODULE_PATH = "article.tasks" +from article.tasks import ( # noqa: E402 + task_export_article_to_articlemeta, + task_harvest_articles, + task_dispatch_articles, + task_process_article_pipeline, +) + + +def setUpModule(): + """ + Desliga o logging (abaixo de CRITICAL) para todo o módulo de testes, + evitando ruído de handlers reais (ex.: OpenSearch) configurados no + settings.LOGGING do Django durante os testes. + """ + logging.disable(logging.CRITICAL) + + +def tearDownModule(): + logging.disable(logging.NOTSET) + + +def make_user(user_id=1, username="roberta"): + return MagicMock(id=user_id, username=username) + + +# ======================================================================= +# task_harvest_articles +# ======================================================================= +class TestTaskHarvestArticles(unittest.TestCase): + def _patches(self): + return ( + patch(f"{MODULE_PATH}._get_user"), + patch(f"{MODULE_PATH}.ArticleIteratorBuilder"), + patch(f"{MODULE_PATH}.task_process_article_pipeline"), + patch(f"{MODULE_PATH}.UnexpectedEvent"), + ) + + def test_dispatches_pipeline_task_for_each_yielded_item(self): + user = make_user() + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleIteratorBuilder") as MockBuilder, \ + patch(f"{MODULE_PATH}.task_process_article_pipeline") as mock_pipeline: + mock_get_user.return_value = user + MockBuilder.return_value.from_harvest.return_value = iter([ + {"xml_url": "http://a", "collection_acron": "scl", "pid": "S1"}, + {"xml_url": "http://b", "collection_acron": "scl", "pid": "S2"}, + ]) + + task_harvest_articles(collection_acron_list=["scl"]) + + self.assertEqual(mock_pipeline.delay.call_count, 2) + first_call_kwargs = mock_pipeline.delay.call_args_list[0].kwargs + self.assertEqual(first_call_kwargs["xml_url"], "http://a") + self.assertEqual(first_call_kwargs["user_id"], user.id) + self.assertEqual(first_call_kwargs["username"], user.username) + + def test_none_items_are_skipped(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleIteratorBuilder") as MockBuilder, \ + patch(f"{MODULE_PATH}.task_process_article_pipeline") as mock_pipeline: + mock_get_user.return_value = make_user() + MockBuilder.return_value.from_harvest.return_value = iter( + [None, {"xml_url": "http://a"}, None] + ) + + task_harvest_articles() + + self.assertEqual(mock_pipeline.delay.call_count, 1) + + def test_builder_is_constructed_with_expected_kwargs_including_stop(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleIteratorBuilder") as MockBuilder, \ + patch(f"{MODULE_PATH}.task_process_article_pipeline"): + user = make_user() + mock_get_user.return_value = user + MockBuilder.return_value.from_harvest.return_value = iter([]) + + task_harvest_articles( + collection_acron_list=["scl"], + journal_acron_list=["abc"], + from_date="2024-01-01", + until_date="2024-12-31", + force_update=True, + limit=10, + timeout=30, + opac_url="www.custom.br", + stop=5, + ) + + _, kwargs = MockBuilder.call_args + self.assertIs(kwargs["user"], user) + self.assertEqual(kwargs["collection_acron_list"], ["scl"]) + self.assertEqual(kwargs["journal_acron_list"], ["abc"]) + self.assertEqual(kwargs["from_date"], "2024-01-01") + self.assertEqual(kwargs["until_date"], "2024-12-31") + self.assertTrue(kwargs["force_update"]) + self.assertEqual(kwargs["limit"], 10) + self.assertEqual(kwargs["timeout"], 30) + self.assertEqual(kwargs["opac_url"], "www.custom.br") + self.assertEqual(kwargs["stop"], 5) + + def test_get_user_failure_creates_event_and_reraises_original_exception(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.UnexpectedEvent") as MockEvent: + mock_get_user.side_effect = RuntimeError("boom") + + with self.assertRaisesRegex(RuntimeError, "boom"): + task_harvest_articles( + collection_acron_list=["scl", "mex"], + journal_acron_list=["abc"], + ) + + MockEvent.create.assert_called_once() + + def test_exception_after_params_assigned_creates_unexpected_event_and_reraises(self): + """ + Quando a falha ocorre DEPOIS de `params` já estar definido (ex.: + dentro da construção do builder ou da iteração de from_harvest), + o comportamento é o esperado: loga e relança. + """ + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleIteratorBuilder") as MockBuilder, \ + patch(f"{MODULE_PATH}.UnexpectedEvent") as MockEvent: + mock_get_user.return_value = make_user() + MockBuilder.return_value.from_harvest.side_effect = RuntimeError("boom") + + with self.assertRaises(RuntimeError): + task_harvest_articles( + collection_acron_list=["scl", "mex"], + journal_acron_list=["abc"], + ) + + MockEvent.create.assert_called_once() + _, kwargs = MockEvent.create.call_args + self.assertEqual(kwargs["action"], "task_harvest_articles") + self.assertEqual(kwargs["item"], "scl-mex-abc") + self.assertIsInstance(kwargs["exception"], RuntimeError) + self.assertEqual(kwargs["detail"]["collection_acron_list"], ["scl", "mex"]) + + +# ======================================================================= +# task_dispatch_articles +# ======================================================================= +class TestTaskDispatchArticles(unittest.TestCase): + def test_dispatches_across_all_three_sources(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleIteratorBuilder") as MockBuilder, \ + patch(f"{MODULE_PATH}.task_process_article_pipeline") as mock_pipeline: + user = make_user() + mock_get_user.return_value = user + builder = MockBuilder.return_value + builder.from_article_source.return_value = iter([{"article_source_id": 1}]) + builder.from_pid_provider.return_value = iter([{"pp_xml_id": 2}]) + builder.from_article.return_value = iter([{"pp_xml_id": 3}]) + + task_dispatch_articles() + + self.assertEqual(mock_pipeline.delay.call_count, 3) + dispatched_kwargs = [c.kwargs for c in mock_pipeline.delay.call_args_list] + self.assertEqual(dispatched_kwargs[0]["article_source_id"], 1) + self.assertEqual(dispatched_kwargs[1]["pp_xml_id"], 2) + self.assertEqual(dispatched_kwargs[2]["pp_xml_id"], 3) + for kwargs in dispatched_kwargs: + self.assertEqual(kwargs["user_id"], user.id) + self.assertEqual(kwargs["username"], user.username) + + def test_none_items_are_skipped(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleIteratorBuilder") as MockBuilder, \ + patch(f"{MODULE_PATH}.task_process_article_pipeline") as mock_pipeline: + mock_get_user.return_value = make_user() + builder = MockBuilder.return_value + builder.from_article_source.return_value = iter([None]) + builder.from_pid_provider.return_value = iter([None, {"pp_xml_id": 2}]) + builder.from_article.return_value = iter([None]) + + task_dispatch_articles() + + self.assertEqual(mock_pipeline.delay.call_count, 1) + + def test_sources_are_passed_their_specific_status_filters(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleIteratorBuilder") as MockBuilder, \ + patch(f"{MODULE_PATH}.task_process_article_pipeline"): + mock_get_user.return_value = make_user() + builder = MockBuilder.return_value + builder.from_article_source.return_value = iter([]) + builder.from_pid_provider.return_value = iter([]) + builder.from_article.return_value = iter([]) + + task_dispatch_articles( + proc_status_list=["todo"], + data_status_list=["pending"], + article_source_status_list=["error"], + ) + + builder.from_article_source.assert_called_once_with( + article_source_status_list=["error"] + ) + builder.from_pid_provider.assert_called_once_with(proc_status_list=["todo"]) + builder.from_article.assert_called_once_with(data_status_list=["pending"]) + + def test_builder_constructed_without_harvest_only_kwargs(self): + """ + task_dispatch_articles não usa from_harvest, então o builder é + construído sem limit/timeout/opac_url/stop. + """ + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleIteratorBuilder") as MockBuilder, \ + patch(f"{MODULE_PATH}.task_process_article_pipeline"): + mock_get_user.return_value = make_user() + builder = MockBuilder.return_value + builder.from_article_source.return_value = iter([]) + builder.from_pid_provider.return_value = iter([]) + builder.from_article.return_value = iter([]) + + task_dispatch_articles(collection_acron_list=["scl"]) + + _, kwargs = MockBuilder.call_args + for absent_key in ("limit", "timeout", "opac_url", "stop"): + self.assertNotIn(absent_key, kwargs) + + def test_get_user_failure_creates_event_and_reraises_original_exception(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.UnexpectedEvent") as MockEvent: + mock_get_user.side_effect = RuntimeError("boom") + + with self.assertRaisesRegex(RuntimeError, "boom"): + task_dispatch_articles(collection_acron_list=["scl"]) + + MockEvent.create.assert_called_once() + + def test_exception_after_params_assigned_creates_unexpected_event_and_reraises(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleIteratorBuilder") as MockBuilder, \ + patch(f"{MODULE_PATH}.UnexpectedEvent") as MockEvent: + mock_get_user.return_value = make_user() + MockBuilder.return_value.from_article_source.side_effect = RuntimeError("boom") + + with self.assertRaises(RuntimeError): + task_dispatch_articles(collection_acron_list=["scl"]) + + MockEvent.create.assert_called_once() + _, kwargs = MockEvent.create.call_args + self.assertEqual(kwargs["action"], "task_dispatch_articles") + self.assertEqual(kwargs["item"], "scl") + + +# ======================================================================= +# task_process_article_pipeline +# ======================================================================= +class TestTaskProcessArticlePipeline(unittest.TestCase): + def _mock_common(self, MockArticleSource, MockPidProviderXML, mock_load_article): + article_source = MagicMock() + article_source.get_pid_provider_xml_id.return_value = 55 + MockArticleSource.objects.get.return_value = article_source + + pp_xml = MagicMock() + MockPidProviderXML.objects.select_related.return_value.get.return_value = pp_xml + + article = MagicMock( + is_classic_public=True, valid=True, pid_v3="pid-v3", collections=["c1"] + ) + mock_load_article.return_value = article + return article_source, pp_xml, article + + def test_flow_with_article_source_id(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource, \ + patch(f"{MODULE_PATH}.PidProviderXML") as MockPidProviderXML, \ + patch(f"{MODULE_PATH}.load_article") as mock_load_article, \ + patch(f"{MODULE_PATH}.task_export_article_to_articlemeta") as mock_export: + mock_get_user.return_value = make_user() + article_source, pp_xml, article = self._mock_common( + MockArticleSource, MockPidProviderXML, mock_load_article + ) + + task_process_article_pipeline(article_source_id=99) + + MockArticleSource.objects.get.assert_called_once_with(id=99) + mock_load_article.assert_called_once() + pp_xml.collections.set.assert_called_once_with(article.collections) + article.check_availability.assert_called_once() + mock_export.delay.assert_not_called() + + def test_flow_with_xml_url_creates_article_source(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource, \ + patch(f"{MODULE_PATH}.Collection") as MockCollection, \ + patch(f"{MODULE_PATH}.PidProviderXML") as MockPidProviderXML, \ + patch(f"{MODULE_PATH}.load_article") as mock_load_article: + mock_get_user.return_value = make_user() + MockCollection.get.return_value = "collection-obj" + article_source = MagicMock() + article_source.get_pid_provider_xml_id.return_value = 55 + MockArticleSource.create_or_update.return_value = article_source + MockPidProviderXML.objects.select_related.return_value.get.return_value = MagicMock() + mock_load_article.return_value = MagicMock( + is_classic_public=True, valid=True, collections=[] + ) + + task_process_article_pipeline( + xml_url="http://x/y.xml", + collection_acron="scl", + pid="S123", + source_date="2024-01-01", + ) + + MockArticleSource.create_or_update.assert_called_once() + _, kwargs = MockArticleSource.create_or_update.call_args + self.assertEqual(kwargs["url"], "http://x/y.xml") + self.assertEqual(kwargs["pid"], "S123") + self.assertEqual(kwargs["collection"], "collection-obj") + MockCollection.get.assert_called_once_with("scl") + + def test_flow_with_pp_xml_id_direct_skips_article_source(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource, \ + patch(f"{MODULE_PATH}.PidProviderXML") as MockPidProviderXML, \ + patch(f"{MODULE_PATH}.load_article") as mock_load_article: + mock_get_user.return_value = make_user() + MockPidProviderXML.objects.select_related.return_value.get.return_value = MagicMock() + mock_load_article.return_value = MagicMock( + is_classic_public=True, valid=True, collections=[] + ) + + task_process_article_pipeline(pp_xml_id=123) + + MockArticleSource.objects.get.assert_not_called() + MockArticleSource.create_or_update.assert_not_called() + MockPidProviderXML.objects.select_related.return_value.get.assert_called_once_with( + id=123 + ) + + def test_export_dispatched_when_classic_public_and_valid(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource, \ + patch(f"{MODULE_PATH}.PidProviderXML") as MockPidProviderXML, \ + patch(f"{MODULE_PATH}.load_article") as mock_load_article, \ + patch(f"{MODULE_PATH}.task_export_article_to_articlemeta") as mock_export: + user = make_user() + mock_get_user.return_value = user + self._mock_common(MockArticleSource, MockPidProviderXML, mock_load_article) + + task_process_article_pipeline( + article_source_id=1, + export_to_articlemeta=True, + collection_acron_list=["scl"], + ) + + mock_export.delay.assert_called_once_with( + pid_v3="pid-v3", + collection_acron_list=["scl"], + force_update=None, + user_id=user.id, + username=user.username, + ) + + def test_export_skipped_when_not_classic_public(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource, \ + patch(f"{MODULE_PATH}.PidProviderXML") as MockPidProviderXML, \ + patch(f"{MODULE_PATH}.load_article") as mock_load_article, \ + patch(f"{MODULE_PATH}.task_export_article_to_articlemeta") as mock_export: + mock_get_user.return_value = make_user() + article_source, pp_xml, article = self._mock_common( + MockArticleSource, MockPidProviderXML, mock_load_article + ) + article.is_classic_public = False + + result = task_process_article_pipeline( + article_source_id=1, export_to_articlemeta=True + ) + + mock_export.delay.assert_not_called() + self.assertIsNone(result) + + def test_export_skipped_when_not_valid(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource, \ + patch(f"{MODULE_PATH}.PidProviderXML") as MockPidProviderXML, \ + patch(f"{MODULE_PATH}.load_article") as mock_load_article, \ + patch(f"{MODULE_PATH}.task_export_article_to_articlemeta") as mock_export: + mock_get_user.return_value = make_user() + article_source, pp_xml, article = self._mock_common( + MockArticleSource, MockPidProviderXML, mock_load_article + ) + article.valid = False + + task_process_article_pipeline( + article_source_id=1, export_to_articlemeta=True + ) + + mock_export.delay.assert_not_called() + + def test_check_availability_force_update_true_when_export_flag_set(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource, \ + patch(f"{MODULE_PATH}.PidProviderXML") as MockPidProviderXML, \ + patch(f"{MODULE_PATH}.load_article") as mock_load_article, \ + patch(f"{MODULE_PATH}.task_export_article_to_articlemeta"): + mock_get_user.return_value = make_user() + article_source, pp_xml, article = self._mock_common( + MockArticleSource, MockPidProviderXML, mock_load_article + ) + + task_process_article_pipeline( + article_source_id=1, export_to_articlemeta=True, force_update=False + ) + + _, kwargs = article.check_availability.call_args + self.assertTrue(kwargs["force_update"]) + + def test_check_availability_force_update_false_when_no_flags(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.ArticleSource") as MockArticleSource, \ + patch(f"{MODULE_PATH}.PidProviderXML") as MockPidProviderXML, \ + patch(f"{MODULE_PATH}.load_article") as mock_load_article, \ + patch(f"{MODULE_PATH}.task_export_article_to_articlemeta"): + mock_get_user.return_value = make_user() + article_source, pp_xml, article = self._mock_common( + MockArticleSource, MockPidProviderXML, mock_load_article + ) + + task_process_article_pipeline( + article_source_id=1, export_to_articlemeta=False, force_update=False + ) + + _, kwargs = article.check_availability.call_args + self.assertFalse(kwargs["force_update"]) + + def test_generic_exception_is_logged_and_reraised(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.UnexpectedEvent") as MockEvent: + mock_get_user.side_effect = RuntimeError("boom") + + with self.assertRaisesRegex(RuntimeError, "boom"): + task_process_article_pipeline(xml_url="http://x") + + MockEvent.create.assert_called_once() + + def test_missing_collection_acron_validation_error_is_logged_and_reraised(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.UnexpectedEvent") as MockEvent: + mock_get_user.return_value = make_user() + + with self.assertRaisesRegex(ValueError, "collection_acron is required"): + task_process_article_pipeline(xml_url="http://x", pid="S1") + + MockEvent.create.assert_called_once() + _, kwargs = MockEvent.create.call_args + self.assertIn("collection_acron is required", str(kwargs["exception"])) + + def test_no_entry_point_validation_error_is_logged_and_reraised(self): + with patch(f"{MODULE_PATH}._get_user") as mock_get_user, \ + patch(f"{MODULE_PATH}.UnexpectedEvent") as MockEvent: + mock_get_user.return_value = make_user() + + with self.assertRaisesRegex(ValueError, "No valid entry point"): + task_process_article_pipeline() + + MockEvent.create.assert_called_once() + _, kwargs = MockEvent.create.call_args + self.assertIn("No valid entry point", str(kwargs["exception"])) + + +class TestTaskExportArticleToArticleMeta(unittest.TestCase): + def test_export_failure_creates_event_with_action_and_article_item(self): + article = MagicMock(spec=Article) + article.__str__.return_value = "S123456789" + + with patch(f"{MODULE_PATH}.Article.objects.get", return_value=article), \ + patch(f"{MODULE_PATH}._get_user", return_value=make_user()), \ + patch( + f"{MODULE_PATH}.controller.export_article_to_articlemeta", + side_effect=RuntimeError("boom"), + ), \ + patch(f"{MODULE_PATH}.UnexpectedEvent") as MockEvent: + task_export_article_to_articlemeta( + pid_v3="S123456789", + collection_acron_list=["scl"], + force_update=True, + ) + + MockEvent.create.assert_called_once() + _, kwargs = MockEvent.create.call_args + self.assertEqual( + kwargs["action"], + "article.tasks.task_export_article_to_articlemeta", + ) + self.assertEqual(kwargs["item"], "S123456789") + self.assertEqual( + kwargs["detail"], + { + "collection_acron_list": ["scl"], + "pid_v3": "S123456789", + "force_update": True, + }, + ) + + +if __name__ == "__main__": + unittest.main() \ No newline at end of file diff --git a/article/tests/test_models.py b/article/tests/test_models.py index 5c602ad32..d4a2fea16 100644 --- a/article/tests/test_models.py +++ b/article/tests/test_models.py @@ -2,13 +2,21 @@ from unittest.mock import patch from django.contrib.auth import get_user_model -from django.test import TestCase +from django.test import SimpleTestCase, TestCase from django.utils.timezone import make_aware from freezegun import freeze_time from article import choices -from article.models import Article, ArticleAffiliation, ContribCollab, ContribPerson +from article.models import ( + AMArticle, + Article, + ArticleAffiliation, + ArticleSource, + ContribCollab, + ContribPerson, +) from article.tests.test_mixins import ArticleTestMixin +from collection.models import Collection from organization.models import NormAffiliation from organization.tests.test_mixins import OrganizationTestMixin @@ -739,4 +747,113 @@ def test_contrib_person_all_name_fields(self): self.assertEqual(person.given_names, "John Robert") self.assertEqual(person.last_name, "Smith") self.assertEqual(person.suffix, "Jr.") - self.assertEqual(person.declared_name, "Dr. John R. Smith Jr.") \ No newline at end of file + self.assertEqual(person.declared_name, "Dr. John R. Smith Jr.") + + +class ArticleSourceUpdateTest(SimpleTestCase): + def test_public_article_source_returns_from_not_public_to_pending(self): + article_source = ArticleSource( + source_date="2026-08-03", + pid="S123456789", + status=ArticleSource.StatusChoices.NOT_PUBLIC, + ) + + changed = article_source.update( + source_date="2026-08-03", + collection=None, + pid="S123456789", + is_public=True, + ) + + self.assertTrue(changed) + self.assertEqual( + article_source.status, + ArticleSource.StatusChoices.PENDING, + ) + + +class AMArticleGetTest(TestCase): + def setUp(self): + self.user = User.objects.create_user(username="legacy_user") + self.collection = Collection.objects.create( + acron3="scl", + creator=self.user, + ) + self.article1 = Article.objects.create( + pid_v3="article-v3-1", + creator=self.user, + ) + self.article2 = Article.objects.create( + pid_v3="article-v3-2", + creator=self.user, + ) + self.am_article1 = AMArticle.objects.create( + pid="S0000000000000000000001", + collection=self.collection, + new_record=self.article1, + status="done", + creator=self.user, + ) + self.am_article2 = AMArticle.objects.create( + pid="S0000000000000000000002", + collection=self.collection, + new_record=self.article2, + status="done", + creator=self.user, + ) + + def test_get_preserves_generated_legacy_records_without_url_or_data(self): + found = AMArticle.get( + pid=self.am_article1.pid, + collection=self.collection, + ) + + self.assertEqual(found, self.am_article1) + self.assertEqual(AMArticle.objects.count(), 2) + self.assertTrue( + AMArticle.objects.filter(pk=self.am_article2.pk).exists() + ) + + +class ArticleCheckAvailabilityTest(SimpleTestCase): + def test_existing_availability_restores_public_status(self): + article = Article( + valid=True, + data_status=choices.DATA_STATUS_COMPLETED, + is_classic_public=True, + is_new_public=True, + is_public=True, + ) + + with patch.object( + article, + "is_pp_xml_valid", + return_value=True, + ), patch.object( + article, + "is_available", + return_value=True, + ), patch.object( + article, + "classic_available", + ) as mock_classic_available, patch.object( + article, + "new_available", + ) as mock_new_available, patch.object( + article, + "save", + ) as mock_save: + mock_classic_available.return_value.exists.return_value = True + mock_new_available.return_value.exists.return_value = True + + result = article.check_availability( + user=None, + force_update=False, + ) + + self.assertTrue(result) + self.assertEqual( + article.data_status, + choices.DATA_STATUS_PUBLIC, + ) + mock_save.assert_called_once() diff --git a/article/wagtail_hooks.py b/article/wagtail_hooks.py index 08464e780..af830e593 100644 --- a/article/wagtail_hooks.py +++ b/article/wagtail_hooks.py @@ -135,17 +135,20 @@ class ArticleSourceSnippetViewSet(SnippetViewSet): menu_order = 200 list_display = [ - "am_article", + "url", "pid_provider_xml", + "pid", + "collection", "status", "source_date", "updated", ] list_filter = [ "status", - "am_article__collection", + "collection", + ] - search_fields = ["url", "pid_provider_xml__v3", "am_article__collection__acron3"] + search_fields = ["url", "pid_provider_xml__v3", "pid"] ordering = ["-updated"] list_per_page = 25 diff --git a/bigbang/tasks_scheduler.py b/bigbang/tasks_scheduler.py index 48aa115d1..97b2f70d5 100644 --- a/bigbang/tasks_scheduler.py +++ b/bigbang/tasks_scheduler.py @@ -42,8 +42,16 @@ def delete_outdated_tasks(task_list=None): "article.tasks.task_create_pid_provider_xml", "article.tasks.task_fix_journal_articles_status", "article.tasks.task_select_articles_to_export_to_articlemeta", + # task_dispatch_articles foi substituída por 2 tasks especializadas: + # task_harvest_articles e task_dispatch_articles + # (esta última consolida article_source/pid_provider/article em um + # único fluxo sequencial). + "article.tasks.task_dispatch_articles", + "article.tasks.task_dispatch_articles_from_pid_provider", + "article.tasks.task_dispatch_articles_from_article", + "article.tasks.task_dispatch_articles_from_article_source", "issue.tasks.load_issue_from_article_meta", - + # Tarefas de Article sem namespace (legacy) "article_complete_data", "convert_xml_to_other_formats", @@ -70,6 +78,7 @@ def delete_outdated_tasks(task_list=None): "task_create_pid_provider_xml", "task_fix_journal_articles_status", "task_select_articles_to_export_to_articlemeta", + "task_dispatch_articles", ] delete_tasks(task_list) @@ -87,15 +96,16 @@ def schedule_tasks(username): delete_outdated_tasks() # Tarefas de Article mantidas + schedule_task_harvest_articles(username, enabled) schedule_task_dispatch_articles(username, enabled) schedule_task_export_articles_to_articlemeta(username, enabled) schedule_task_fix_article_status(username, enabled) - + # Tarefas de issue schedule_export_issue_to_articlemeta(username, enabled) schedule_export_issues_to_articlemeta(username, enabled) schedule_load_issue_from_articlemeta(username, enabled) - + # Tarefas de journal schedule_export_journal_to_articlemeta(username, enabled) schedule_export_journals_to_articlemeta(username, enabled) @@ -103,7 +113,7 @@ def schedule_tasks(username): schedule_fetch_and_process_journal_logos_in_collection(username, enabled) schedule_load_journal_from_article_meta(username, enabled) schedule_collect_journals_from_am(username, enabled) - + # Tarefas de pid_provider schedule_fix_pid_provider_xmls_status(username, enabled) @@ -115,11 +125,44 @@ def schedule_tasks(username): # TAREFAS DE ARTICLE MANTIDAS # ============================================================================== +def schedule_task_harvest_articles(username, enabled=False): + """ + Agenda o pipeline completo a partir da coleta (harvest) até a + publicação no ArticleMeta. + """ + schedule_task( + task="article.tasks.task_harvest_articles", + name="article.tasks.task_harvest_articles", + kwargs=dict( + username=username, + user_id=None, + collection_acron_list=None, + journal_acron_list=None, + from_date=None, + until_date=None, + force_update=False, + export_to_articlemeta=False, + auto_solve_pid_conflict=False, + limit=None, + timeout=None, + opac_url=None, + stop=None, + ), + description=_("Dispatch articles: harvest -> ... -> ArticleMeta"), + priority=TASK_PRIORITY, + enabled=enabled, + run_once=False, + day_of_week="*", + hour="2", + minute="1", + ) + + def schedule_task_dispatch_articles(username, enabled=False): """ - Agenda a tarefa orquestradora de despacho de artigos para o pipeline. - Substitui as antigas tarefas de seleção (complete_data, load_from_api, - load_from_article_source, load_articles). + Agenda o fluxo sequencial pelas 3 etapas internas (article_source -> + pid_provider -> article), pegando pendentes em cada uma, até a + publicação no ArticleMeta. """ schedule_task( task="article.tasks.task_dispatch_articles", @@ -136,20 +179,20 @@ def schedule_task_dispatch_articles(username, enabled=False): force_update=False, export_to_articlemeta=False, auto_solve_pid_conflict=False, + article_source_status_list=None, proc_status_list=None, data_status_list=None, - limit=None, - timeout=None, - opac_url=None, - article_source_status_list=None, ), - description=_("Dispatch articles to processing pipeline"), + description=_( + "Dispatch pending articles: article_source -> pid_provider -> " + "article -> ArticleMeta" + ), priority=TASK_PRIORITY, enabled=enabled, run_once=False, day_of_week="*", hour="2", - minute="1", + minute="16", ) @@ -183,13 +226,10 @@ def schedule_task_export_articles_to_articlemeta(username, enabled=False): ) - - - def schedule_task_fix_article_status(username, enabled=False): """ Agenda a tarefa de corrigir status dos registros de artigos. - + Permite marcar artigos como inválidos, públicos ou duplicados, além de deduplicar registros conforme necessário. """ @@ -241,7 +281,7 @@ def schedule_bigbang_start(username, enabled=False): def schedule_bigbang_delete_outdated_tasks(username, enabled=False): """ Agenda a tarefa de limpeza de tarefas obsoletas do Article - + Remove tarefas antigas e não utilizadas do módulo Article, mantendo o scheduler limpo e organizado. """ @@ -270,7 +310,7 @@ def schedule_bigbang_delete_outdated_tasks(username, enabled=False): def schedule_load_journal_from_article_meta(username, enabled=False): """ Agenda a tarefa de carga de dados de journals obtidos do AM e Core. - + Configura verify=False para verificação SSL nas requisições HTTP. """ schedule_task( @@ -294,7 +334,7 @@ def schedule_load_journal_from_article_meta(username, enabled=False): def schedule_collect_journals_from_am(username, enabled=False): """ Agenda a tarefa de coleta de journals da fonte AM. - + Configura verify=False para verificação SSL nas requisições HTTP. """ schedule_task( @@ -503,7 +543,7 @@ def schedule_export_issue_to_articlemeta(username, enabled=False): def schedule_fix_pid_provider_xmls_status(username, enabled=False): """ Agenda a tarefa de corrigir status dos XMLs do PID Provider. - + Permite marcar XMLs como inválidos, públicos ou duplicados, além de deduplicar registros conforme necessário. """ diff --git a/core/models.py b/core/models.py index 2f7dc6ac0..52e93246b 100755 --- a/core/models.py +++ b/core/models.py @@ -1284,10 +1284,7 @@ def __str__(self): def get(cls, pid, collection): if not pid and not collection: raise ValueError("Param pid and collection_acron3 is required") - try: - cls.objects.filter(url__isnull=True, data__isnull=True).delete() - except Exception: - pass + try: return cls.objects.get(pid=pid, collection=collection) except cls.MultipleObjectsReturned: diff --git a/journal/models.py b/journal/models.py index b5666d10a..d622b6e8a 100755 --- a/journal/models.py +++ b/journal/models.py @@ -950,21 +950,7 @@ def select_items( params["scielojournal__collection__acron3__in"] = collection_acron_list if journal_acron_list: params["scielojournal__journal_acron__in"] = journal_acron_list - queryset = cls.objects.filter(**params).distinct() - if not queryset.exists(): - UnexpectedEvent.create( - exception=ValueError("No journals found for the given filters"), - detail={ - "operation": "Journal.select_items", - "collection_acron_list": collection_acron_list, - "journal_acron_list": journal_acron_list, - "from_date": from_date, - "until_date": until_date, - "days_to_go_back": days_to_go_back, - "params": params, - }, - ) - return queryset + return cls.objects.filter(**params).distinct() @classmethod def get_journal_issns( @@ -1027,8 +1013,9 @@ def collection_acrons(self): @classmethod def get_ids(cls, collection_acron_list=None, journal_acron_list=None): - qs = cls.select_items(collection_acron_list, journal_acron_list) - return qs.values_list("id", flat=True).distinct() + return cls.select_items( + collection_acron_list, journal_acron_list + ).values_list("id", flat=True).distinct() def select_collections(self, collection_acron_list=None, is_active=None): params = {} diff --git a/pid_provider/choices.py b/pid_provider/choices.py index 9415dab9b..446ec24db 100644 --- a/pid_provider/choices.py +++ b/pid_provider/choices.py @@ -2,16 +2,24 @@ ENDPOINTS = (("fix-pid-v2", "fix-pid-v2"),) +# article is not public PPXML_STATUS_WAIT = "WAIT" +# PPXML_STATUS_IGNORED = "IGNORE" +# ready to create article PPXML_STATUS_TODO = "TODO" +# article was created PPXML_STATUS_DONE = "DONE" PPXML_STATUS_UNDEF = "UNDEF" +# xml is broken PPXML_STATUS_INVALID = "NVALID" +# journal / issue in xml is not registered PPXML_STATUS_UNMATCHED_JOURNAL_OR_ISSUE = "UNMATCH" - +# is duplicated PPXML_STATUS_DUPLICATED = "DUP" +# removed the duplication PPXML_STATUS_DEDUPLICATED = "DEDUP" + PPXML_STATUS = ( (PPXML_STATUS_TODO, _("To do")), (PPXML_STATUS_DONE, _("Done")), @@ -23,3 +31,18 @@ (PPXML_STATUS_DEDUPLICATED, _("deduplicated")), (PPXML_STATUS_UNMATCHED_JOURNAL_OR_ISSUE, _("unmatched journal or issue")), ) + +PPXML_STATUS_TO_CREATE_OR_UPDATE_ARTICLE_SOURCE = [ + PPXML_STATUS_INVALID, + PPXML_STATUS_WAIT, +] +PPXML_STATUS_TO_IGNORE = [ + PPXML_STATUS_INVALID, + PPXML_STATUS_DONE, + PPXML_STATUS_WAIT, + PPXML_STATUS_IGNORED, + PPXML_STATUS_UNDEF, + PPXML_STATUS_DUPLICATED, + PPXML_STATUS_DEDUPLICATED, + PPXML_STATUS_UNMATCHED_JOURNAL_OR_ISSUE, +] \ No newline at end of file