diff --git a/README.md b/README.md index f57d8b6..3e612eb 100644 --- a/README.md +++ b/README.md @@ -106,6 +106,7 @@ task_search_log_files.delay( from_date="2021-05-21", until_date="2021-05-21", trigger_validation=True, + parse_queue_name="parse_xlarge", ) ``` @@ -150,6 +151,20 @@ Use `compose=production.yml` or another Compose file when needed: make ps compose=production.yml ``` +Parsing behavior can be adjusted through comma-separated environment variables: + +| Variable | Default | +|---|---| +| `DEFAULT_PARSE_QUEUE` | `parse_small` | +| `PARSING_METADATA_CACHE_COLLECTIONS` | All active log collections | +| `PARSING_METADATA_CACHE_RELEASE_COLLECTIONS` | `scl` | +| `YEAR_PARTITIONED_COLLECTIONS` | `chl,col,mex,scl` | + +The worker image starts one Celery process per container. Its entrypoint accepts +`CELERY_WORKER_QUEUES`, `CELERY_WORKER_CONCURRENCY`, +`CELERY_WORKER_PREFETCH_MULTIPLIER`, `CELERY_WORKER_NAME`, and +`CELERY_WORKER_LOG_LEVEL`. An empty queue setting consumes Celery's default queue. + Run one test path: ```bash @@ -172,14 +187,17 @@ Metadata synchronization keeps sources and documents updated from ArticleMeta, O Configure the default schedule manually in Wagtail/Admin through `django-celery-beat` `PeriodicTask` records. Exact cron times may vary by installation, but the default -operational setup should include: +operational setup should include the entries below. Every task that triggers +parsing can receive an explicit `parse_queue_name`. If omitted, the downstream +worker uses `DEFAULT_PARSE_QUEUE` (`parse_small` by default). The periodic task +itself runs on `load`, while `parse_queue_name` selects the parsing worker. | Task | Suggested schedule | Notes | |---|---|---| -| `[Metadata] Daily Sync Routine (Auto)` | Daily, early morning | Refreshes sources and documents before log processing. Use the `load` queue. | -| `[Log Pipeline] Daily Routine (Auto)` | Daily, after metadata sync | Runs Search -> Validate -> Parse -> Export for new logs. Use the `load` queue. | -| `[Metrics] Resume Log Exports` | Every 15-30 minutes | Retries errored or stale daily metric export jobs. | -| `[Metrics] Resume Stale Parsing Logs` | Every 30-60 minutes | Marks stale `PAR` logs for retry. | +| Individual `[Metadata] Sync ...` tasks | Daily, early morning | Stagger the required ArticleMeta, OPAC, Books, Preprints, and Dataverse collectors on the `load` queue. | +| `[Log Pipeline] 1. Search Logs (Manual)` | Daily, after metadata sync | Create one entry per collection group. Use `load` as the task queue and set `parse_queue_name` for parsing. | +| `[Metrics] Resume Log Exports` | Every 15-30 minutes | Retries errored or stale daily metric export jobs. Accepts `queue_name`; otherwise uses `DEFAULT_PARSE_QUEUE`. | +| `[Metrics] Resume Stale Parsing Logs` | Every 30-60 minutes | Marks stale `PAR` logs for retry. Accepts `queue_name`; otherwise uses `DEFAULT_PARSE_QUEUE`. | | `[Metrics] Cleanup Daily Payloads` | Daily or weekly | Removes old exported daily payload files. | | `[Reports] Populate All Reports` | Daily, after log processing | Refreshes weekly, monthly, and yearly log report tables. | diff --git a/VERSION b/VERSION index e75da3e..00355e2 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -2.3.6 +2.3.7 diff --git a/collection/tasks.py b/collection/tasks.py index 303ecf8..9e1046b 100644 --- a/collection/tasks.py +++ b/collection/tasks.py @@ -1,11 +1,7 @@ -from django.contrib.auth import get_user_model - from collection.models import Collection from config import celery_app from core.utils.request_utils import _get_user -User = get_user_model() - @celery_app.task(bind=True, name="[Collection] Load Collection Data") def task_load_collections(self, user_id=None, username=None): diff --git a/compose/local/django/Dockerfile b/compose/local/django/Dockerfile index aac7972..fb9cae7 100644 --- a/compose/local/django/Dockerfile +++ b/compose/local/django/Dockerfile @@ -1,10 +1,10 @@ -ARG PYTHON_VERSION=3.11-bullseye +ARG PYTHON_VERSION=3.11-bookworm # define an alias for the specfic python version used in this file. -FROM python:${PYTHON_VERSION} as python +FROM python:${PYTHON_VERSION} AS python # Python build stage -FROM python as python-build-stage +FROM python AS python-build-stage ARG BUILD_ENVIRONMENT=local @@ -27,7 +27,7 @@ RUN pip wheel --wheel-dir /usr/src/app/wheels -r ${BUILD_ENVIRONMENT}.txt # Python 'run' stage -FROM python as python-run-stage +FROM python AS python-run-stage ARG BUILD_ENVIRONMENT=local ARG APP_HOME=/app diff --git a/compose/local/django/celery/worker/start b/compose/local/django/celery/worker/start index f0c7efc..993e6c6 100644 --- a/compose/local/django/celery/worker/start +++ b/compose/local/django/celery/worker/start @@ -1,34 +1,20 @@ #!/bin/bash set -o errexit +set -o pipefail set -o nounset -# Worker padrão -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -n worker.default@%h & +worker_args=( + -A config.celery_app + worker + -l "${CELERY_WORKER_LOG_LEVEL:-INFO}" + --concurrency="${CELERY_WORKER_CONCURRENCY:-1}" + --prefetch-multiplier="${CELERY_WORKER_PREFETCH_MULTIPLIER:-1}" + -n "${CELERY_WORKER_NAME:-worker.default@%h}" +) -# Worker para load e validação de dados -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=4 -Q load -n worker.load@%h & +if [ -n "${CELERY_WORKER_QUEUES:-}" ]; then + worker_args+=(-Q "${CELERY_WORKER_QUEUES}") +fi -# Worker para scl -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_xlarge -n worker.parse_xlarge@%h & - -# Worker para chl col mex -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_large -n worker.parse_large@%h & - -# Worker para cri esp psi prt ven -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_medium -n worker.parse_medium@%h & - -# Worker para arg bol cub data ecu per preprints pry rve spa sss sza ury wid -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small -n worker.parse_small@%h & - -# Workers seriais adicionais para backfill paralelo de colecoes pequenas -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_1 -n worker.parse_small_1@%h & -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_2 -n worker.parse_small_2@%h & -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_3 -n worker.parse_small_3@%h & -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_4 -n worker.parse_small_4@%h & -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_5 -n worker.parse_small_5@%h & -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_6 -n worker.parse_small_6@%h & -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_7 -n worker.parse_small_7@%h & -watchgod celery.__main__.main --args -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_8 -n worker.parse_small_8@%h & - -wait +exec watchgod celery.__main__.main --args "${worker_args[@]}" diff --git a/compose/local/django/start b/compose/local/django/start index f076ee5..ba96db4 100644 --- a/compose/local/django/start +++ b/compose/local/django/start @@ -6,4 +6,4 @@ set -o nounset python manage.py migrate -python manage.py runserver_plus 0.0.0.0:8000 +exec python manage.py runserver_plus 0.0.0.0:8000 diff --git a/compose/production/django/Dockerfile b/compose/production/django/Dockerfile index 38d73da..87a2378 100644 --- a/compose/production/django/Dockerfile +++ b/compose/production/django/Dockerfile @@ -1,10 +1,10 @@ -ARG PYTHON_VERSION=3.11-bullseye +ARG PYTHON_VERSION=3.11-bookworm # define an alias for the specfic python version used in this file. -FROM python:${PYTHON_VERSION} as python +FROM python:${PYTHON_VERSION} AS python # Python build stage -FROM python as python-build-stage +FROM python AS python-build-stage ARG BUILD_ENVIRONMENT=production @@ -28,7 +28,7 @@ RUN pip wheel --wheel-dir /usr/src/app/wheels \ # Python 'run' stage -FROM python as python-run-stage +FROM python AS python-run-stage ARG BUILD_ENVIRONMENT=production ARG APP_HOME=/app diff --git a/compose/production/django/celery/worker/start b/compose/production/django/celery/worker/start index 6269dd5..f9d1ad5 100644 --- a/compose/production/django/celery/worker/start +++ b/compose/production/django/celery/worker/start @@ -4,32 +4,17 @@ set -o errexit set -o pipefail set -o nounset -# Worker padrão -celery -A config.celery_app worker -l INFO --concurrency=1 -n worker.default@%h & - -# Worker para load e validação de dados -celery -A config.celery_app worker -l INFO --concurrency=4 -Q load -n worker.load@%h & - -# Worker para parse_scl (coleções extra-grandes) -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_xlarge -n worker.parse_xlarge@%h & - -# Worker para chl col mex (coleções grandes) -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_large -n worker.parse_large@%h & - -# Worker para cri esp psi prt ven (coleções médias) -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_medium -n worker.parse_medium@%h & - -# Worker para arg bol cub data ecu per preprints pry rve spa sss sza ury wid (coleções pequenas) -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small -n worker.parse_small@%h & - -# Workers seriais adicionais para backfill paralelo de colecoes pequenas -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_1 -n worker.parse_small_1@%h & -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_2 -n worker.parse_small_2@%h & -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_3 -n worker.parse_small_3@%h & -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_4 -n worker.parse_small_4@%h & -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_5 -n worker.parse_small_5@%h & -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_6 -n worker.parse_small_6@%h & -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_7 -n worker.parse_small_7@%h & -celery -A config.celery_app worker -l INFO --concurrency=1 -Q parse_small_8 -n worker.parse_small_8@%h & - -wait +worker_args=( + -A config.celery_app + worker + -l "${CELERY_WORKER_LOG_LEVEL:-INFO}" + --concurrency="${CELERY_WORKER_CONCURRENCY:-1}" + --prefetch-multiplier="${CELERY_WORKER_PREFETCH_MULTIPLIER:-1}" + -n "${CELERY_WORKER_NAME:-worker.default@%h}" +) + +if [ -n "${CELERY_WORKER_QUEUES:-}" ]; then + worker_args+=(-Q "${CELERY_WORKER_QUEUES}") +fi + +exec celery "${worker_args[@]}" diff --git a/compose/production/django/start b/compose/production/django/start index d9da8ea..35c1533 100644 --- a/compose/production/django/start +++ b/compose/production/django/start @@ -26,4 +26,5 @@ if compress_enabled; then # NOTE this command will fail if django-compressor is disabled python /app/manage.py compress fi + /usr/local/bin/gunicorn config.wsgi --bind 0.0.0.0:5000 --chdir=/app --timeout 1000 --workers 3 --worker-connections=1000 --worker-class=gevent diff --git a/config/collections.py b/config/collections.py index e64a53d..369b0f7 100644 --- a/config/collections.py +++ b/config/collections.py @@ -1,52 +1,7 @@ -COLLECTION_ACRON3_SIZE_MAP = { - "scl": "xlarge", - "chl": "large", - "col": "large", - "mex": "large", - "cri": "medium", - "esp": "medium", - "psi": "medium", - "prt": "medium", - "ven": "medium", - "arg": "small", - "bol": "small", - "books": "small", - "cub": "small", - "data": "small", - "dom": "small", - "ecu": "small", - "per": "small", - "preprints": "small", - "pry": "small", - "rve": "small", - "rvt": "small", - "spa": "small", - "sss": "small", - "sza": "small", - "ury": "small", - "wid": "small", -} - -COLLECTION_SIZE_SAMPLE_MAP = { - "small": 1.0, - "medium": 0.5, - "large": 0.1, - "xlarge": 0.1, -} - COLLECTION_OPAC_URL_MAP = { "dom": "https://scielo.do/api/v1/counter_dict", "scl": "https://www.scielo.br/api/v1/counter_dict", } - - -def get_collection_size(collection_acronym): - return COLLECTION_ACRON3_SIZE_MAP.get(collection_acronym, "small") - - -def get_collection_parse_queue(collection_acronym): - return f"parse_{get_collection_size(collection_acronym)}" - LOG_MANAGER_SEED_DATA = [ { "acronym": "arg", @@ -74,6 +29,7 @@ def get_collection_parse_queue(collection_acronym): "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "classic", + "sample_size": 0.1, }, { "acronym": "col", @@ -83,6 +39,7 @@ def get_collection_parse_queue(collection_acronym): "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "classic", + "sample_size": 0.1, }, { "acronym": "cri", @@ -92,6 +49,7 @@ def get_collection_parse_queue(collection_acronym): "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "classic", + "sample_size": 0.5, }, { "acronym": "cub", @@ -102,6 +60,25 @@ def get_collection_parse_queue(collection_acronym): "e-mail": "tecnologia@scielo.org", "translator_class": "classic", }, + { + "acronym": "data", + "directory_name": "Dataverse legado", + "path": "/app/logs/bkp-dataverse", + "quantity": 1, + "start_date": "2020-01-01", + "e-mail": "tecnologia@scielo.org", + "translator_class": "dataverse", + "directory_active": False, + }, + { + "acronym": "data", + "directory_name": "BunnyNet Data", + "path": "/app/logs/bkp-bunnynet/data", + "quantity": 1, + "start_date": "2020-01-01", + "e-mail": "tecnologia@scielo.org", + "translator_class": "dataverse", + }, { "acronym": "ecu", "directory_name": "Site clássico", @@ -119,6 +96,7 @@ def get_collection_parse_queue(collection_acronym): "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "classic", + "sample_size": 0.5, }, { "acronym": "mex", @@ -128,6 +106,7 @@ def get_collection_parse_queue(collection_acronym): "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "classic", + "sample_size": 0.1, }, { "acronym": "per", @@ -138,6 +117,15 @@ def get_collection_parse_queue(collection_acronym): "e-mail": "tecnologia@scielo.org", "translator_class": "classic", }, + { + "acronym": "preprints", + "directory_name": "BunnyNet Preprints", + "path": "/app/logs/bkp-bunnynet/preprints", + "quantity": 1, + "start_date": "2020-01-01", + "e-mail": "tecnologia@scielo.org", + "translator_class": "preprints", + }, { "acronym": "prt", "directory_name": "Site clássico", @@ -146,6 +134,7 @@ def get_collection_parse_queue(collection_acronym): "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "classic", + "sample_size": 0.5, }, { "acronym": "pry", @@ -164,6 +153,7 @@ def get_collection_parse_queue(collection_acronym): "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "classic", + "sample_size": 0.5, }, { "acronym": "rve", @@ -182,6 +172,7 @@ def get_collection_parse_queue(collection_acronym): "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "opac", + "sample_size": 0.1, }, { "acronym": "scl", @@ -191,6 +182,7 @@ def get_collection_parse_queue(collection_acronym): "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "opac", + "sample_size": 0.1, }, { "acronym": "sza", @@ -212,21 +204,64 @@ def get_collection_parse_queue(collection_acronym): }, { "acronym": "ven", - "directory_name": "Site clássico", + "directory_name": "Venezuela legado", + "path": "/app/logs/bkp-venezuela", + "quantity": 1, + "start_date": "2020-01-01", + "e-mail": "tecnologia@scielo.org", + "translator_class": "classic", + "directory_active": False, + "sample_size": 0.5, + }, + { + "acronym": "ven", + "directory_name": "Ratchet Venezuela", "path": "/app/logs/bkp-ratchet/scielo.ve", "quantity": 1, "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "classic", + "directory_active": False, + "sample_size": 0.5, + }, + { + "acronym": "ven", + "directory_name": "BunnyNet Venezuela", + "path": "/app/logs/bkp-bunnynet/venezuela", + "quantity": 1, + "start_date": "2020-01-01", + "e-mail": "tecnologia@scielo.org", + "translator_class": "classic", + "sample_size": 0.5, }, { "acronym": "wid", - "directory_name": "SciELO Caribbean", + "directory_name": "Ratchet West Indies", + "path": "/app/logs/bkp-ratchet/scielo.wi", + "quantity": 1, + "start_date": "2020-01-01", + "e-mail": "tecnologia@scielo.org", + "translator_class": "classic", + "directory_active": False, + }, + { + "acronym": "wid", + "directory_name": "BunnyNet Caribbean legado", "path": "/app/logs/bkp-bunnynet/caribbean", "quantity": 1, "start_date": "2020-01-01", "e-mail": "tecnologia@scielo.org", "translator_class": "classic", + "directory_active": False, + }, + { + "acronym": "wid", + "directory_name": "BunnyNet West Indies", + "path": "/app/logs/bkp-bunnynet/westindies", + "quantity": 1, + "start_date": "2020-01-01", + "e-mail": "tecnologia@scielo.org", + "translator_class": "classic", }, { "acronym": "books", diff --git a/config/settings/base.py b/config/settings/base.py index fb0fea5..0ca49ed 100644 --- a/config/settings/base.py +++ b/config/settings/base.py @@ -6,7 +6,34 @@ import environ -from config.collections import COLLECTION_ACRON3_SIZE_MAP # noqa: F401 +DEFAULT_PARSING_METADATA_CACHE_COLLECTIONS = ( + "arg", + "bol", + "books", + "chl", + "col", + "cri", + "cub", + "data", + "dom", + "ecu", + "esp", + "mex", + "per", + "preprints", + "prt", + "pry", + "psi", + "rve", + "scl", + "sss", + "sza", + "ury", + "ven", + "wid", +) +DEFAULT_PARSING_METADATA_CACHE_RELEASE_COLLECTIONS = ("scl",) +DEFAULT_YEAR_PARTITIONED_COLLECTIONS = ("chl", "col", "mex", "scl") ROOT_DIR = Path(__file__).resolve(strict=True).parent.parent.parent # core/ @@ -279,7 +306,7 @@ # https://docs.djangoproject.com/en/dev/ref/settings/#admins ADMINS = [ ("""Rafael JP Damaceno""", "contato@pitangainnovare.com.br"), - ("""Jamil Atta Junior""", "atta.jamil@innolabs.com.br") + ("""Jamil Atta Junior""", "atta.jamil@innolabs.com.br"), ] # https://docs.djangoproject.com/en/dev/ref/settings/#managers MANAGERS = ADMINS @@ -324,7 +351,6 @@ "document.tasks.articlemeta", "document.tasks.dataverse", "document.tasks.opac", - "document.tasks.pipeline", "document.tasks.preprints", "document.tasks.scielo_books", "metrics.tasks.cleanup", @@ -492,10 +518,21 @@ SCIELO_BOOKS_DB_NAME = env("SCIELO_BOOKS_DB_NAME", default="scielobooks_1a") SCIELO_BOOKS_LIMIT = env.int("SCIELO_BOOKS_LIMIT", default=1000) -# Collection size categories +# Log parsing # ------------------------------------------------------------------------------ -SUPPORTED_LOGFILE_EXTENSIONS = env.list("SUPPORTED_LOGFILE_EXTENSIONS", default=[".log", ".gz", ".zip"]) +DEFAULT_PARSE_QUEUE = env("DEFAULT_PARSE_QUEUE", default="parse_small") +SUPPORTED_LOGFILE_EXTENSIONS = env.list( + "SUPPORTED_LOGFILE_EXTENSIONS", default=[".log", ".gz", ".zip"] +) PARSING_METADATA_CACHE_COLLECTIONS = env.list( "PARSING_METADATA_CACHE_COLLECTIONS", - default=list(COLLECTION_ACRON3_SIZE_MAP), + default=list(DEFAULT_PARSING_METADATA_CACHE_COLLECTIONS), +) +PARSING_METADATA_CACHE_RELEASE_COLLECTIONS = env.list( + "PARSING_METADATA_CACHE_RELEASE_COLLECTIONS", + default=list(DEFAULT_PARSING_METADATA_CACHE_RELEASE_COLLECTIONS), +) +YEAR_PARTITIONED_COLLECTIONS = env.list( + "YEAR_PARTITIONED_COLLECTIONS", + default=list(DEFAULT_YEAR_PARTITIONED_COLLECTIONS), ) diff --git a/core/users/tasks.py b/core/users/tasks.py deleted file mode 100644 index 39511e8..0000000 --- a/core/users/tasks.py +++ /dev/null @@ -1,11 +0,0 @@ -from django.contrib.auth import get_user_model - -from config import celery_app - -User = get_user_model() - - -@celery_app.task(bind=True, name="Get users count") -def get_users_count(self): - """A pointless Celery task to demonstrate usage.""" - return User.objects.count() diff --git a/core/users/tests/test_tasks.py b/core/users/tests/test_tasks.py deleted file mode 100644 index 90a5a06..0000000 --- a/core/users/tests/test_tasks.py +++ /dev/null @@ -1,16 +0,0 @@ -import pytest -from celery.result import EagerResult - -from core.users.tasks import get_users_count -from core.users.tests.factories import UserFactory - -pytestmark = pytest.mark.django_db - - -def test_user_count(settings): - """A basic test to execute the get_users_count Celery task.""" - UserFactory.create_batch(3) - settings.CELERY_TASK_ALWAYS_EAGER = True - task_result = get_users_count.delay() - assert isinstance(task_result, EagerResult) - assert task_result.result == 3 diff --git a/document/tasks/pipeline.py b/document/tasks/pipeline.py deleted file mode 100644 index 1073aa8..0000000 --- a/document/tasks/pipeline.py +++ /dev/null @@ -1,28 +0,0 @@ -import logging - -from celery import group -from django.utils.translation import gettext as _ - -from config import celery_app - -from document.tasks.articlemeta import task_load_documents_from_article_meta -from document.tasks.dataverse import task_load_dataset_metadata_into_documents -from document.tasks.opac import task_load_documents_from_opac -from document.tasks.preprints import task_load_preprints_into_documents -from document.tasks.scielo_books import task_sync_documents_from_scielo_books - - -@celery_app.task( - bind=True, name=_("[Metadata] Daily Sync Routine (Auto)"), queue="load" -) -def task_daily_metadata_sync_pipeline(self): - logging.info("Starting Daily Metadata Sync Pipeline") - group( - [ - task_load_documents_from_article_meta.s(), - task_load_documents_from_opac.s(), - task_load_preprints_into_documents.s(), - task_load_dataset_metadata_into_documents.s(), - task_sync_documents_from_scielo_books.s(), - ] - ).apply_async() diff --git a/local.yml b/local.yml index 9b3a047..44fe005 100644 --- a/local.yml +++ b/local.yml @@ -1,28 +1,40 @@ +x-django: &django + image: scielo_usage_local_django + depends_on: + - redis + - postgres + - mailhog + volumes: + - .:/app:z + - /mnt/pidata2/pi/scl/logs:/app/logs + # Uncomment to use local SciELO lib repos for development: + # - ../scielo_log_validator:/app/scielo_log_validator:z + # - ../scielo_usage_counter:/app/scielo_usage_counter:z + env_file: + - ./.envs/.local/.django + - ./.envs/.local/.postgres + environment: &django_environment + USE_LOCAL_SCIELO_LIBS: 0 + ports: + - "8009:8000" + command: /start + +x-celeryworker: &celeryworker + <<: *django + image: scielo_usage_local_celeryworker + ports: [] + command: /start-celeryworker + environment: + <<: *django_environment + CELERY_WORKER_NAME: worker.default@%h + services: - django: &django + django: + <<: *django build: context: . dockerfile: ./compose/local/django/Dockerfile - image: scielo_usage_local_django container_name: scielo_usage_local_django - depends_on: - - redis - - postgres - - mailhog - volumes: - - .:/app:z - - /mnt/pidata2/pi/scl/logs:/app/logs - # Uncomment to use local SciELO lib repos for development: - # - ../scielo_log_validator:/app/scielo_log_validator:z - # - ../scielo_usage_counter:/app/scielo_usage_counter:z - env_file: - - ./.envs/.local/.django - - ./.envs/.local/.postgres - environment: - - USE_LOCAL_SCIELO_LIBS=0 - ports: - - "8009:8000" - command: /start postgres: build: @@ -51,18 +63,74 @@ services: - "6399:6379" celeryworker: - <<: *django - image: scielo_usage_local_celeryworker + <<: *celeryworker + build: + context: . + dockerfile: ./compose/local/django/Dockerfile container_name: scielo_usage_local_celeryworker - depends_on: - - redis - - postgres - - mailhog - ports: [] - command: /start-celeryworker + + celeryworker_load: + <<: *celeryworker + container_name: scielo_usage_local_celeryworker_load + environment: + <<: *django_environment + CELERY_WORKER_CONCURRENCY: 4 + CELERY_WORKER_NAME: worker.load@%h + CELERY_WORKER_QUEUES: load + + celeryworker_tiny: + <<: *celeryworker + container_name: scielo_usage_local_celeryworker_tiny + environment: + <<: *django_environment + CELERY_WORKER_NAME: worker.parse_tiny@%h + CELERY_WORKER_QUEUES: parse_tiny + + celeryworker_small_1: + <<: *celeryworker + container_name: scielo_usage_local_celeryworker_small_1 + environment: + <<: *django_environment + CELERY_WORKER_NAME: worker.parse_small_1@%h + CELERY_WORKER_QUEUES: parse_small + + celeryworker_small_2: + <<: *celeryworker + container_name: scielo_usage_local_celeryworker_small_2 + environment: + <<: *django_environment + CELERY_WORKER_NAME: worker.parse_small_2@%h + CELERY_WORKER_QUEUES: parse_small + + celeryworker_medium: + <<: *celeryworker + container_name: scielo_usage_local_celeryworker_medium + environment: + <<: *django_environment + CELERY_WORKER_NAME: worker.parse_medium@%h + CELERY_WORKER_QUEUES: parse_medium + + celeryworker_large: + <<: *celeryworker + container_name: scielo_usage_local_celeryworker_large + environment: + <<: *django_environment + CELERY_WORKER_NAME: worker.parse_large@%h + CELERY_WORKER_QUEUES: parse_large + + celeryworker_xlarge: + <<: *celeryworker + container_name: scielo_usage_local_celeryworker_xlarge + environment: + <<: *django_environment + CELERY_WORKER_NAME: worker.parse_xlarge@%h + CELERY_WORKER_QUEUES: parse_xlarge celerybeat: <<: *django + build: + context: . + dockerfile: ./compose/local/django/Dockerfile image: scielo_usage_local_celerybeat container_name: scielo_usage_local_celerybeat depends_on: diff --git a/log_manager/tasks.py b/log_manager/tasks.py index a04f22a..b20eb1c 100644 --- a/log_manager/tasks.py +++ b/log_manager/tasks.py @@ -1,9 +1,7 @@ -import logging - from celery import chord +from django.conf import settings from config import celery_app -from config.collections import get_collection_parse_queue from core.utils.request_utils import _get_user from log_manager.services import catalog, validation from metrics.tasks.log_parsing import task_enqueue_log_parsing_jobs @@ -19,13 +17,17 @@ def task_search_log_files( user_id=None, username=None, trigger_validation=False, + parse_queue_name=None, ): """ Search for log files in configured collection directories. When trigger_validation=True, this starts the full Search -> Validate -> Parse - chain. Parse callbacks are routed by collection size. + chain using parse_queue_name or the configured default parse queue. """ + if trigger_validation: + parse_queue_name = parse_queue_name or settings.DEFAULT_PARSE_QUEUE + _get_user(self.request, username=username, user_id=user_id) catalog.catalog_log_files_from_configured_directories( @@ -45,6 +47,7 @@ def task_search_log_files( "user_id": user_id, "username": username, "trigger_parse": True, + "parse_queue_name": parse_queue_name, } ) @@ -67,13 +70,17 @@ def task_validate_log_files( trigger_parse=False, revalidate=False, status_list=None, + parse_queue_name=None, ): """ Validate cataloged log files. When trigger_parse=True, one parse orchestration task is enqueued per - collection and routed to the proper parse_ queue. + collection using parse_queue_name or the configured default parse queue. """ + if trigger_parse: + parse_queue_name = parse_queue_name or settings.DEFAULT_PARSE_QUEUE + log_hashes_by_collection = validation.get_validation_candidate_hashes_by_collection( collections=collections, from_date=from_date, @@ -100,6 +107,7 @@ def task_validate_log_files( days_to_go_back=days_to_go_back, user_id=user_id, username=username, + parse_queue_name=parse_queue_name, ) return @@ -120,15 +128,6 @@ def task_validate_log_file(self, log_file_hash, user_id=None, username=None): validation.validate_log_file_and_update_status(log_file_hash) -@celery_app.task(bind=True, name="[Log Pipeline] Daily Routine (Auto)", queue="load") -def task_daily_log_ingestion_pipeline(self): - """ - Start the daily Search -> Validate -> Parse chain with default parameters. - """ - logging.info("Starting Daily Log Ingestion Pipeline") - task_search_log_files.apply_async(kwargs={"trigger_validation": True}) - - def _build_validation_tasks(log_hashes_by_collection, user_id, username): return { collection_code: [ @@ -140,7 +139,13 @@ def _build_validation_tasks(log_hashes_by_collection, user_id, username): def _enqueue_parse_after_validation( - tasks_by_collection, from_date, until_date, days_to_go_back, user_id, username + tasks_by_collection, + from_date, + until_date, + days_to_go_back, + user_id, + username, + parse_queue_name, ): for collection_code, validation_tasks in tasks_by_collection.items(): if validation_tasks: @@ -152,6 +157,7 @@ def _enqueue_parse_after_validation( days_to_go_back, user_id, username, + parse_queue_name, ) ) else: @@ -163,12 +169,19 @@ def _enqueue_parse_after_validation( days_to_go_back, user_id, username, + parse_queue_name, ) ) def _build_parse_signature( - collection_code, from_date, until_date, days_to_go_back, user_id, username + collection_code, + from_date, + until_date, + days_to_go_back, + user_id, + username, + parse_queue_name, ): apply_kwargs = _build_parse_apply_kwargs( collection_code, @@ -177,28 +190,33 @@ def _build_parse_signature( days_to_go_back, user_id, username, + parse_queue_name, ) parse_callback = task_enqueue_log_parsing_jobs.si(**apply_kwargs["kwargs"]) - if apply_kwargs.get("queue"): - parse_callback.set(queue=apply_kwargs["queue"]) + parse_callback.set(queue=apply_kwargs["queue"]) return parse_callback def _build_parse_apply_kwargs( - collection_code, from_date, until_date, days_to_go_back, user_id, username + collection_code, + from_date, + until_date, + days_to_go_back, + user_id, + username, + parse_queue_name, ): collections = [collection_code] - parse_queue = get_collection_parse_queue(collection_code) apply_kwargs = { "kwargs": { "collections": collections, "from_date": from_date, "until_date": until_date, "days_to_go_back": days_to_go_back, - "queue_name": parse_queue, + "queue_name": parse_queue_name, "user_id": user_id, "username": username, }, - "queue": parse_queue, + "queue": parse_queue_name, } return apply_kwargs diff --git a/log_manager/tests/test_tasks.py b/log_manager/tests/test_tasks.py index 79d1db7..a326ff8 100644 --- a/log_manager/tests/test_tasks.py +++ b/log_manager/tests/test_tasks.py @@ -1,6 +1,6 @@ from unittest.mock import patch -from django.test import TestCase +from django.test import TestCase, override_settings from log_manager import tasks @@ -17,7 +17,7 @@ def test_returns_none_for_empty_date_range(self): self.assertIsNone(result) mocked_signature.assert_not_called() - def test_routes_parse_callback_to_collection_queue(self): + def test_routes_parse_callback_to_configured_queue(self): with patch( "log_manager.tasks.task_enqueue_log_parsing_jobs.apply_async" ) as mocked_apply_async: @@ -26,16 +26,17 @@ def test_routes_parse_callback_to_collection_queue(self): from_date="2024-02-01", until_date="2024-02-02", trigger_parse=True, + parse_queue_name="parse_tiny", ) mocked_apply_async.assert_called_once() - self.assertEqual(mocked_apply_async.call_args.kwargs["queue"], "parse_small") + self.assertEqual(mocked_apply_async.call_args.kwargs["queue"], "parse_tiny") self.assertEqual( mocked_apply_async.call_args.kwargs["kwargs"]["queue_name"], - "parse_small", + "parse_tiny", ) - def test_routes_each_collection_to_its_queue(self): + def test_routes_all_configured_collections_to_same_queue(self): with patch( "log_manager.tasks.task_enqueue_log_parsing_jobs.apply_async" ) as mocked_apply_async: @@ -44,10 +45,85 @@ def test_routes_each_collection_to_its_queue(self): from_date="2024-02-01", until_date="2024-02-02", trigger_parse=True, + parse_queue_name="parse_tiny", ) calls = { call.kwargs["kwargs"]["collections"][0]: call.kwargs["queue"] for call in mocked_apply_async.call_args_list } - self.assertEqual(calls, {"books": "parse_small", "scl": "parse_xlarge"}) + self.assertEqual(calls, {"books": "parse_tiny", "scl": "parse_tiny"}) + + @patch( + "log_manager.tasks.validation.get_validation_candidate_hashes_by_collection", + return_value={"books": ["a" * 32]}, + ) + @patch("log_manager.tasks.chord") + def test_routes_parse_callback_after_validation_to_configured_queue( + self, mocked_chord, mocked_candidates + ): + tasks.task_validate_log_files.run( + collections=["books"], + trigger_parse=True, + parse_queue_name="parse_tiny", + ) + + callback = mocked_chord.return_value.call_args.args[0] + self.assertEqual(callback.options["queue"], "parse_tiny") + self.assertEqual(callback.kwargs["queue_name"], "parse_tiny") + + @override_settings(DEFAULT_PARSE_QUEUE="parse_default") + def test_trigger_parse_uses_default_queue(self): + with patch( + "log_manager.tasks.task_enqueue_log_parsing_jobs.apply_async" + ) as mocked_apply_async: + tasks.task_validate_log_files.run( + collections=["books"], + from_date="2024-02-01", + until_date="2024-02-02", + trigger_parse=True, + ) + + self.assertEqual(mocked_apply_async.call_args.kwargs["queue"], "parse_default") + + +class SearchLogFilesTaskTests(TestCase): + @patch("log_manager.tasks.catalog.catalog_log_files_from_configured_directories") + @patch("log_manager.tasks.task_validate_log_files.apply_async") + def test_propagates_configured_queue_to_validation( + self, mocked_validation_apply_async, mocked_catalog + ): + tasks.task_search_log_files.run( + collections=["data"], + days_to_go_back=7, + trigger_validation=True, + parse_queue_name="parse_tiny", + ) + + mocked_catalog.assert_called_once() + self.assertEqual( + mocked_validation_apply_async.call_args.kwargs["kwargs"][ + "parse_queue_name" + ], + "parse_tiny", + ) + + @override_settings(DEFAULT_PARSE_QUEUE="parse_default") + @patch("log_manager.tasks.catalog.catalog_log_files_from_configured_directories") + @patch("log_manager.tasks.task_validate_log_files.apply_async") + def test_trigger_validation_uses_default_queue( + self, mocked_validation_apply_async, mocked_catalog + ): + tasks.task_search_log_files.run( + collections=["data"], + days_to_go_back=7, + trigger_validation=True, + ) + + mocked_catalog.assert_called_once() + self.assertEqual( + mocked_validation_apply_async.call_args.kwargs["kwargs"][ + "parse_queue_name" + ], + "parse_default", + ) diff --git a/log_manager_config/models.py b/log_manager_config/models.py index f8fc106..80ed804 100644 --- a/log_manager_config/models.py +++ b/log_manager_config/models.py @@ -146,7 +146,7 @@ def load(cls, data, user): config=config, directory_name=item.get("directory_name"), path=item.get("path"), - active=item.get("active", True), + active=item.get("directory_active", item.get("active", True)), translator_class=item.get("translator_class", "classic"), ) diff --git a/log_manager_config/tasks.py b/log_manager_config/tasks.py index 12ceabc..17ffb7c 100644 --- a/log_manager_config/tasks.py +++ b/log_manager_config/tasks.py @@ -1,17 +1,14 @@ import logging -from config import celery_app -from config.collections import ( - COLLECTION_OPAC_URL_MAP, - COLLECTION_SIZE_SAMPLE_MAP, - LOG_MANAGER_SEED_DATA, - get_collection_size, -) from collection.models import Collection +from config import celery_app +from config.collections import COLLECTION_OPAC_URL_MAP, LOG_MANAGER_SEED_DATA from core.utils.request_utils import _get_user - from log_manager_config import models +DEFAULT_VALIDATION_SAMPLE_SIZE = 1.0 +DEFAULT_VALIDATION_BUFFER_SIZE = 2048 + @celery_app.task(bind=True, name="[Log Pipeline] Load Log Manager Settings (Seed)") def task_load_log_manager_collection_settings( @@ -19,9 +16,10 @@ def task_load_log_manager_collection_settings( ): user = _get_user(self.request, username=username, user_id=user_id) - if not data: - data = LOG_MANAGER_SEED_DATA + using_default_data = not data + data = [dict(item) for item in (data or LOG_MANAGER_SEED_DATA)] + if using_default_data: for acronym, opac_url in COLLECTION_OPAC_URL_MAP.items(): try: collection = Collection.objects.get(acron3=acronym) @@ -33,10 +31,9 @@ def task_load_log_manager_collection_settings( collection.updated_by = user collection.save(update_fields=["opac_url", "updated_by", "updated"]) - for i in data: - size = get_collection_size(i["acronym"]) - i["sample_size"] = COLLECTION_SIZE_SAMPLE_MAP.get(size, 1.0) - i["buffer_size"] = 2048 + for item in data: + item.setdefault("sample_size", DEFAULT_VALIDATION_SAMPLE_SIZE) + item.setdefault("buffer_size", DEFAULT_VALIDATION_BUFFER_SIZE) models.LogManagerCollectionConfig.load(data, user) models.CollectionLogDirectory.load(data, user) diff --git a/log_manager_config/tests/test_models.py b/log_manager_config/tests/test_models.py index 6c1dad3..fa8fcb7 100644 --- a/log_manager_config/tests/test_models.py +++ b/log_manager_config/tests/test_models.py @@ -82,3 +82,20 @@ def test_translator_class_defaults_to_classic(self): ) self.assertEqual(directory.translator_class, "classic") + + def test_load_uses_directory_active_independently(self): + CollectionLogDirectory.load( + [ + { + "acronym": "scl", + "directory_name": "legacy logs", + "path": "/data/logs/legacy", + "directory_active": False, + } + ], + self.user, + ) + + directory = CollectionLogDirectory.objects.get(path="/data/logs/legacy") + + self.assertFalse(directory.active) diff --git a/log_manager_config/tests/test_tasks.py b/log_manager_config/tests/test_tasks.py index 229459b..0253fd1 100644 --- a/log_manager_config/tests/test_tasks.py +++ b/log_manager_config/tests/test_tasks.py @@ -6,54 +6,100 @@ from config.collections import COLLECTION_OPAC_URL_MAP, LOG_MANAGER_SEED_DATA from log_manager_config import tasks - EXPECTED_LOG_DIRECTORIES = { - ("arg", "/app/logs/bkp-ratchet/scielo.ar", "classic"), - ("bol", "/app/logs/bkp-ratchet/scielo.bo", "classic"), - ("books", "/app/logs/bkp-bunnynet/books", "books"), - ("chl", "/app/logs/bkp-ratchet/scielo.cl", "classic"), - ("col", "/app/logs/bkp-ratchet/scielo.co", "classic"), - ("cri", "/app/logs/bkp-ratchet/scielo.cr", "classic"), - ("cub", "/app/logs/bkp-ratchet/scielo.cu", "classic"), - ("ecu", "/app/logs/bkp-ratchet/scielo.ec", "classic"), - ("esp", "/app/logs/bkp-ratchet/scielo.es", "classic"), - ("mex", "/app/logs/bkp-ratchet/scielo.mx", "classic"), - ("per", "/app/logs/bkp-ratchet/scielo.pe", "classic"), - ("prt", "/app/logs/bkp-ratchet/scielo.pt", "classic"), - ("pry", "/app/logs/bkp-ratchet/scielo.py", "classic"), - ("psi", "/app/logs/bkp-ratchet/scielo.pepsic", "classic"), - ("rve", "/app/logs/bkp-ratchet/scielo.revenf", "classic"), - ("scl", "/app/logs/bkp-bunnynet/scielo-br", "opac"), - ("scl", "/app/logs/bkp-bunnynet/scielo-br-2", "opac"), - ("sza", "/app/logs/bkp-ratchet/scielo.za", "classic"), - ("ury", "/app/logs/bkp-ratchet/scielo.uy", "classic"), - ("ven", "/app/logs/bkp-ratchet/scielo.ve", "classic"), - ("wid", "/app/logs/bkp-bunnynet/caribbean", "classic"), + ("arg", "/app/logs/bkp-ratchet/scielo.ar", "classic", True), + ("bol", "/app/logs/bkp-ratchet/scielo.bo", "classic", True), + ("books", "/app/logs/bkp-bunnynet/books", "books", True), + ("chl", "/app/logs/bkp-ratchet/scielo.cl", "classic", True), + ("col", "/app/logs/bkp-ratchet/scielo.co", "classic", True), + ("cri", "/app/logs/bkp-ratchet/scielo.cr", "classic", True), + ("cub", "/app/logs/bkp-ratchet/scielo.cu", "classic", True), + ("data", "/app/logs/bkp-dataverse", "dataverse", False), + ("data", "/app/logs/bkp-bunnynet/data", "dataverse", True), + ("ecu", "/app/logs/bkp-ratchet/scielo.ec", "classic", True), + ("esp", "/app/logs/bkp-ratchet/scielo.es", "classic", True), + ("mex", "/app/logs/bkp-ratchet/scielo.mx", "classic", True), + ("per", "/app/logs/bkp-ratchet/scielo.pe", "classic", True), + ("preprints", "/app/logs/bkp-bunnynet/preprints", "preprints", True), + ("prt", "/app/logs/bkp-ratchet/scielo.pt", "classic", True), + ("pry", "/app/logs/bkp-ratchet/scielo.py", "classic", True), + ("psi", "/app/logs/bkp-ratchet/scielo.pepsic", "classic", True), + ("rve", "/app/logs/bkp-ratchet/scielo.revenf", "classic", True), + ("scl", "/app/logs/bkp-bunnynet/scielo-br", "opac", True), + ("scl", "/app/logs/bkp-bunnynet/scielo-br-2", "opac", True), + ("sza", "/app/logs/bkp-ratchet/scielo.za", "classic", True), + ("ury", "/app/logs/bkp-ratchet/scielo.uy", "classic", True), + ("ven", "/app/logs/bkp-venezuela", "classic", False), + ("ven", "/app/logs/bkp-ratchet/scielo.ve", "classic", False), + ("ven", "/app/logs/bkp-bunnynet/venezuela", "classic", True), + ("wid", "/app/logs/bkp-ratchet/scielo.wi", "classic", False), + ("wid", "/app/logs/bkp-bunnynet/caribbean", "classic", False), + ("wid", "/app/logs/bkp-bunnynet/westindies", "classic", True), } -def test_default_seed_matches_active_log_directories(): +def test_default_seed_matches_log_directories(): configured_directories = { - (item["acronym"], item["path"], item["translator_class"]) + ( + item["acronym"], + item["path"], + item["translator_class"], + item.get("directory_active", True), + ) for item in LOG_MANAGER_SEED_DATA } assert configured_directories == EXPECTED_LOG_DIRECTORIES assert { - item["quantity"] - for item in LOG_MANAGER_SEED_DATA - if item["acronym"] == "scl" + item["quantity"] for item in LOG_MANAGER_SEED_DATA if item["acronym"] == "scl" } == {2} assert not { - "data", "dom", - "preprints", "rvt", "spa", "sss", } & {item["acronym"] for item in LOG_MANAGER_SEED_DATA} +def test_default_seed_preserves_validation_sample_policy(): + sample_sizes = { + item["acronym"]: item.get( + "sample_size", + tasks.DEFAULT_VALIDATION_SAMPLE_SIZE, + ) + for item in LOG_MANAGER_SEED_DATA + } + + assert { + acronym for acronym, sample_size in sample_sizes.items() if sample_size == 0.1 + } == {"chl", "col", "mex", "scl"} + assert { + acronym for acronym, sample_size in sample_sizes.items() if sample_size == 0.5 + } == {"cri", "esp", "prt", "psi", "ven"} + + +def test_custom_seed_gets_defaults_without_mutating_input(): + data = [{"acronym": "custom"}] + + with ( + patch.object(tasks, "_get_user", return_value=None), + patch.object(tasks.models.LogManagerCollectionConfig, "load") as load_config, + patch.object(tasks.models.CollectionLogDirectory, "load"), + patch.object(tasks.models.CollectionEmail, "load"), + ): + tasks.task_load_log_manager_collection_settings.run(data=data) + + loaded_data = load_config.call_args.args[0] + assert loaded_data == [ + { + "acronym": "custom", + "sample_size": tasks.DEFAULT_VALIDATION_SAMPLE_SIZE, + "buffer_size": tasks.DEFAULT_VALIDATION_BUFFER_SIZE, + } + ] + assert data == [{"acronym": "custom"}] + + @pytest.mark.django_db def test_default_seed_configures_collection_opac_urls(): scl = Collection.objects.create(acron3="scl") diff --git a/metrics/opensearch/names.py b/metrics/opensearch/names.py index b567d11..7ad0e87 100644 --- a/metrics/opensearch/names.py +++ b/metrics/opensearch/names.py @@ -1,4 +1,4 @@ -from config.collections import get_collection_size +from django.conf import settings def _validate_index_inputs(index_prefix, collection, date): @@ -15,17 +15,22 @@ def extract_access_year(date): return date.split("-")[0] +def _uses_year_partition(collection): + configured_collections = { + value.strip().lower() for value in settings.YEAR_PARTITIONED_COLLECTIONS + } + return collection.lower() in configured_collections + + def generate_month_index_name(index_prefix, collection, date): _validate_index_inputs(index_prefix, collection, date) - size = get_collection_size(collection) - if size in ("xlarge", "large"): + if _uses_year_partition(collection): return f"{index_prefix}_monthly_{collection}_{extract_access_year(date)}" return f"{index_prefix}_monthly_{collection}" def generate_year_index_name(index_prefix, collection, date): _validate_index_inputs(index_prefix, collection, date) - size = get_collection_size(collection) - if size in ("xlarge", "large"): + if _uses_year_partition(collection): return f"{index_prefix}_yearly_{collection}_{extract_access_year(date)}" return f"{index_prefix}_yearly_{collection}" diff --git a/metrics/services/log_parsing_jobs.py b/metrics/services/log_parsing_jobs.py index de5b20f..ede822a 100644 --- a/metrics/services/log_parsing_jobs.py +++ b/metrics/services/log_parsing_jobs.py @@ -1,5 +1,6 @@ +from django.conf import settings + from collection.models import Collection -from config.collections import get_collection_parse_queue from core.utils.date_utils import get_date_obj, get_date_range_str from log_manager import choices from log_manager.models import LogFile @@ -28,6 +29,8 @@ def enqueue_log_parsing_jobs( skip_log_hashes=None, robots_source=None, ): + queue_name = queue_name or settings.DEFAULT_PARSE_QUEUE + from_date, until_date = get_date_range_str(from_date, until_date, days_to_go_back) from_date_obj = get_date_obj(from_date) until_date_obj = get_date_obj(until_date) @@ -119,6 +122,8 @@ def wait_log_parsing_wave( robots_source=None, wave_log_hashes=None, ): + queue_name = queue_name or settings.DEFAULT_PARSE_QUEUE + wave_job_ids = wave_job_ids or wave_log_hashes or [] if DailyMetricJob.objects.filter( pk__in=wave_job_ids, @@ -146,9 +151,8 @@ def wait_log_parsing_wave( apply_kwargs = { "kwargs": kwargs, "countdown": poll_interval_seconds, + "queue": queue_name, } - if queue_name: - apply_kwargs["queue"] = queue_name wait_log_parsing_wave_task.apply_async(**apply_kwargs) return {"wave_completed": False, "reexecution_enqueued": False} @@ -169,10 +173,7 @@ def wait_log_parsing_wave( skip_log_hashes=skip_log_hashes, robots_source=robots_source, ) - apply_kwargs = {"kwargs": kwargs} - if queue_name: - apply_kwargs["queue"] = queue_name - log_parsing_task.apply_async(**apply_kwargs) + log_parsing_task.apply_async(kwargs=kwargs, queue=queue_name) return {"wave_completed": True, "reexecution_enqueued": True} @@ -251,7 +252,7 @@ def _enqueue_collection_daily_jobs( daily_metric_export_task.apply_async( args=(job.pk, track_errors, user_id, username, robots_source), - queue=queue_name or get_collection_parse_queue(collection.acron3), + queue=queue_name, ) result["enqueued_wave_job_ids"].append(job.pk) result["enqueued_jobs"] += 1 @@ -308,10 +309,7 @@ def _schedule_log_parsing_reexecution( robots_source=robots_source, ) - apply_kwargs = {"kwargs": kwargs} - if queue_name: - apply_kwargs["queue"] = queue_name - wait_log_parsing_wave_task.apply_async(**apply_kwargs) + wait_log_parsing_wave_task.apply_async(kwargs=kwargs, queue=queue_name) return True diff --git a/metrics/services/parsing/job_payloads.py b/metrics/services/parsing/job_payloads.py index 94991c5..f4213c7 100644 --- a/metrics/services/parsing/job_payloads.py +++ b/metrics/services/parsing/job_payloads.py @@ -4,7 +4,6 @@ from django.conf import settings -from config.collections import get_collection_size from log_manager.models import LogFile from metrics.counter.access.daily_accumulator import DailyAccessAccumulator from metrics.counter.indexing import converter as index_docs @@ -49,7 +48,7 @@ def build_daily_metric_job_payload(job, robots_list, mmdb, track_errors=False): memory.format_snapshot(), ) - if get_collection_size(job.collection.acron3) == "xlarge": + if metadata_cache.should_release_after_job(job.collection): metadata_cache.clear() gc.collect() logging.info( diff --git a/metrics/services/parsing/metadata_cache.py b/metrics/services/parsing/metadata_cache.py index ac1eb80..d071227 100644 --- a/metrics/services/parsing/metadata_cache.py +++ b/metrics/services/parsing/metadata_cache.py @@ -7,7 +7,6 @@ from django.db import connection, transaction from django.db.models import Count, Max -from config.collections import get_collection_size from document.models import Document from source.models import Source @@ -16,18 +15,23 @@ def is_enabled(collection): - enabled_collections = { - value.lower() - for value in getattr(settings, "PARSING_METADATA_CACHE_COLLECTIONS", []) - } - return collection.acron3.lower() in enabled_collections + return _is_collection_configured( + collection, + "PARSING_METADATA_CACHE_COLLECTIONS", + ) + + +def should_release_after_job(collection): + return _is_collection_configured( + collection, + "PARSING_METADATA_CACHE_RELEASE_COLLECTIONS", + ) def get_url_translation_manager(collection, translator_class, build_manager): global _CACHE_ENTRY acronym = collection.acron3 - size = get_collection_size(acronym) translator_name = translator_class.__name__ with _CACHE_LOCK: @@ -38,9 +42,8 @@ def get_url_translation_manager(collection, translator_class, build_manager): signature = _read_signature(collection) if signature == entry["signature"]: logging.info( - "Parsing metadata cache hit for %s (size=%s, signature=%s).", + "Parsing metadata cache hit for %s (signature=%s).", acronym, - size, signature, ) return _fresh_manager(entry) @@ -57,10 +60,9 @@ def get_url_translation_manager(collection, translator_class, build_manager): logging.info( "Parsing metadata cache %s for %s " - "(size=%s, reason=%s, signature=%s, build_seconds=%.3f).", + "(reason=%s, signature=%s, build_seconds=%.3f).", "miss" if entry is None else "rebuild", acronym, - size, reason, new_entry["signature"], elapsed, @@ -75,6 +77,13 @@ def clear(): _CACHE_ENTRY = None +def _is_collection_configured(collection, setting_name): + configured_collections = { + value.strip().lower() for value in getattr(settings, setting_name, []) + } + return collection.acron3.lower() in configured_collections + + def _get_rebuild_reason(entry, collection, translator_name): if entry is None: return "empty" diff --git a/metrics/services/resume.py b/metrics/services/resume.py index 48253a4..6ccceda 100644 --- a/metrics/services/resume.py +++ b/metrics/services/resume.py @@ -1,8 +1,8 @@ import logging +from django.conf import settings from django.utils import timezone -from config.collections import get_collection_parse_queue from core.utils.date_utils import get_date_obj, get_date_range_str from log_manager import choices from log_manager.models import LogFile @@ -29,6 +29,8 @@ def resume_daily_metric_jobs( username=None, robots_source=None, ): + queue_name = queue_name or settings.DEFAULT_PARSE_QUEUE + from_date, until_date = get_date_range_str(from_date, until_date, days_to_go_back) from_date_obj = get_date_obj(from_date) until_date_obj = get_date_obj(until_date) @@ -79,6 +81,8 @@ def resume_stale_parsing_logs( username=None, robots_source=None, ): + queue_name = queue_name or settings.DEFAULT_PARSE_QUEUE + from_date, until_date = get_date_range_str(from_date, until_date, days_to_go_back) from_date_obj = get_date_obj(from_date) until_date_obj = get_date_obj(until_date) @@ -129,7 +133,7 @@ def _enqueue_resumable_daily_metric_jobs( daily_metric_export_task.apply_async( args=(job.pk, False, user_id, username, robots_source), - queue=queue_name or get_collection_parse_queue(job.collection.acron3), + queue=queue_name, ) resumed_jobs += 1 return resumed_jobs @@ -226,27 +230,23 @@ def _enqueue_log_parsing_retry( username, robots_source, ): - apply_kwargs = { - "kwargs": { - "collections": collections, - "include_logs_with_error": True, - "batch_size": batch_size, - "max_log_files": max_log_files, - "auto_reexecute": False, - "replace": False, - "track_errors": track_errors, - "from_date": from_date, - "until_date": until_date, - "days_to_go_back": None, - "queue_name": queue_name, - "user_id": user_id, - "username": username, - "robots_source": robots_source, - } + kwargs = { + "collections": collections, + "include_logs_with_error": True, + "batch_size": batch_size, + "max_log_files": max_log_files, + "auto_reexecute": False, + "replace": False, + "track_errors": track_errors, + "from_date": from_date, + "until_date": until_date, + "days_to_go_back": None, + "queue_name": queue_name, + "user_id": user_id, + "username": username, + "robots_source": robots_source, } - if queue_name: - apply_kwargs["queue"] = queue_name - log_parsing_task.apply_async(**apply_kwargs) + log_parsing_task.apply_async(kwargs=kwargs, queue=queue_name) def _extract_date_from_validation_dict(validation): diff --git a/metrics/tests/opensearch/test_names.py b/metrics/tests/opensearch/test_names.py index f33dab1..58ad18a 100644 --- a/metrics/tests/opensearch/test_names.py +++ b/metrics/tests/opensearch/test_names.py @@ -1,9 +1,18 @@ import unittest +from django.conf import settings +from django.test import override_settings + from metrics.opensearch.names import generate_month_index_name, generate_year_index_name class TestIndexNames(unittest.TestCase): + def test_year_partition_policy_is_explicit(self): + self.assertEqual( + set(settings.YEAR_PARTITIONED_COLLECTIONS), + {"chl", "col", "mex", "scl"}, + ) + def test_generate_index_names_for_year_and_month(self): self.assertEqual( generate_year_index_name("usage", "scl", "2024-01-15"), @@ -21,3 +30,18 @@ def test_generate_index_names_for_year_and_month(self): generate_month_index_name("usage", "books", "2024-01-15"), "usage_monthly_books", ) + + @override_settings(YEAR_PARTITIONED_COLLECTIONS=[" books ", "SCL"]) + def test_generate_index_names_uses_configured_collections(self): + self.assertEqual( + generate_year_index_name("usage", "books", "2024-01-15"), + "usage_yearly_books_2024", + ) + self.assertEqual( + generate_month_index_name("usage", "scl", "2024-01-15"), + "usage_monthly_scl_2024", + ) + self.assertEqual( + generate_year_index_name("usage", "chl", "2024-01-15"), + "usage_yearly_chl", + ) diff --git a/metrics/tests/parsing/test_metadata_cache.py b/metrics/tests/parsing/test_metadata_cache.py index c076ad2..a967b8f 100644 --- a/metrics/tests/parsing/test_metadata_cache.py +++ b/metrics/tests/parsing/test_metadata_cache.py @@ -4,12 +4,9 @@ import pytest from collection.models import Collection -from config.collections import COLLECTION_ACRON3_SIZE_MAP, LOG_MANAGER_SEED_DATA +from config.collections import LOG_MANAGER_SEED_DATA from document.models import Document -from log_manager_config.models import ( - CollectionLogDirectory, - LogManagerCollectionConfig, -) +from log_manager_config.models import CollectionLogDirectory, LogManagerCollectionConfig from metrics.services.parsing import metadata, metadata_cache from source.models import Source @@ -56,10 +53,18 @@ def parsing_collection(db, settings): def test_all_known_collections_are_enabled_by_default(settings): active_log_collections = {item["acronym"] for item in LOG_MANAGER_SEED_DATA} - assert active_log_collections <= set(COLLECTION_ACRON3_SIZE_MAP) - assert set(settings.PARSING_METADATA_CACHE_COLLECTIONS) == set( - COLLECTION_ACRON3_SIZE_MAP - ) + assert active_log_collections <= set(settings.PARSING_METADATA_CACHE_COLLECTIONS) + + +def test_release_after_job_uses_dedicated_collection_setting(settings): + collection = SimpleNamespace(acron3="scl") + settings.PARSING_METADATA_CACHE_RELEASE_COLLECTIONS = ["scl"] + + assert metadata_cache.should_release_after_job(collection) + + settings.PARSING_METADATA_CACHE_RELEASE_COLLECTIONS = [] + + assert not metadata_cache.should_release_after_job(collection) @pytest.mark.django_db diff --git a/metrics/tests/services/test_tasks.py b/metrics/tests/services/test_tasks.py index abad24e..f85caf7 100644 --- a/metrics/tests/services/test_tasks.py +++ b/metrics/tests/services/test_tasks.py @@ -1,7 +1,7 @@ from datetime import date, timedelta from unittest.mock import patch -from django.test import TestCase +from django.test import TestCase, override_settings from django.utils import timezone from collection.models import Collection @@ -47,6 +47,7 @@ def test_task_enqueue_log_parsing_jobs_enqueues_one_daily_job_per_collection_dat include_logs_with_error=False, from_date="2012-03-01", until_date="2012-03-31", + queue_name="parse_tiny", ) self.assertEqual(result["enqueued_jobs"], 2) @@ -59,7 +60,7 @@ def test_task_enqueue_log_parsing_jobs_enqueues_one_daily_job_per_collection_dat self.assertEqual(jobs[0].input_log_hashes, sorted([first.hash, second.hash])) self.assertEqual(jobs[1].input_log_hashes, [third.hash]) - def test_task_enqueue_log_parsing_jobs_allows_queue_override_and_robots_source( + def test_task_enqueue_log_parsing_jobs_uses_queue_and_robots_source( self, ): self._log_file("1" * 32, "2012-03-10") @@ -96,6 +97,7 @@ def test_task_enqueue_log_parsing_jobs_excludes_error_logs_when_not_requested(se include_logs_with_error=False, from_date="2012-03-01", until_date="2012-03-31", + queue_name="parse_tiny", ) mocked_apply_async.assert_called_once() @@ -122,6 +124,7 @@ def test_task_enqueue_log_parsing_jobs_skip_log_hashes_prevents_reprocessing_sam from_date="2012-03-01", until_date="2012-03-31", skip_log_hashes=[skipped.hash], + queue_name="parse_tiny", ) mocked_apply_async.assert_called_once() @@ -143,6 +146,7 @@ def test_task_enqueue_log_parsing_jobs_max_log_files_counts_files_not_jobs(self) max_log_files=1, from_date="2012-03-01", until_date="2012-03-31", + queue_name="parse_tiny", ) mocked_apply_async.assert_called_once() @@ -172,6 +176,7 @@ def test_wait_log_parsing_wave_rechecks_until_daily_jobs_complete(self): include_logs_with_error=False, max_log_files=2, auto_reexecute=True, + queue_name="parse_tiny", ) self.assertEqual( @@ -203,6 +208,21 @@ def test_wait_log_parsing_wave_preserves_queue_name(self): mocked_wait_apply_async.call_args.kwargs["queue"], "parse_small" ) + @override_settings(DEFAULT_PARSE_QUEUE="parse_default") + def test_task_enqueue_log_parsing_jobs_uses_default_queue(self): + self._log_file("4" * 32, "2012-03-10") + + with patch( + "metrics.tasks.log_parsing.task_build_and_export_daily_metric_job.apply_async" + ) as mocked_apply_async: + task_enqueue_log_parsing_jobs.run( + collections=["books"], + from_date="2012-03-01", + until_date="2012-03-31", + ) + + self.assertEqual(mocked_apply_async.call_args.kwargs["queue"], "parse_default") + class ResumeDailyMetricJobTests(TestCase): def setUp(self): @@ -241,6 +261,16 @@ def test_resume_log_exports_requeues_error_daily_jobs(self): ) self.assertEqual(result["resumed_logs"], 1) + @override_settings(DEFAULT_PARSE_QUEUE="parse_default") + @patch( + "metrics.services.resume._enqueue_resumable_daily_metric_jobs", + return_value=0, + ) + def test_resume_log_exports_uses_default_queue(self, mocked_enqueue): + task_resume_log_exports.run(collections=["books"]) + + self.assertEqual(mocked_enqueue.call_args.kwargs["queue_name"], "parse_default") + def test_resume_log_exports_clears_payload_when_current_logs_change(self): log_file = LogFile.objects.create( hash="2" * 32, @@ -267,6 +297,7 @@ def test_resume_log_exports_clears_payload_when_current_logs_change(self): collections=["books"], from_date="2012-03-01", until_date="2012-03-31", + queue_name="parse_tiny", ) job.refresh_from_db() @@ -301,6 +332,7 @@ def test_resume_log_exports_preserves_payload_when_current_logs_match(self): collections=["books"], from_date="2012-03-01", until_date="2012-03-31", + queue_name="parse_tiny", ) job.refresh_from_db() @@ -325,6 +357,7 @@ def test_resume_log_exports_requeues_stored_payload_without_current_logs(self): collections=["books"], from_date="2012-03-01", until_date="2012-03-31", + queue_name="parse_tiny", ) mocked_apply_async.assert_called_once() @@ -345,6 +378,7 @@ def test_resume_log_exports_skips_jobs_without_logs_or_payload(self): collections=["books"], from_date="2012-03-01", until_date="2012-03-31", + queue_name="parse_tiny", ) mocked_apply_async.assert_not_called() @@ -375,6 +409,7 @@ def test_resume_log_exports_releases_stale_exporting_jobs(self): from_date="2012-03-01", until_date="2012-03-31", stale_after_minutes=60, + queue_name="parse_tiny", ) job.refresh_from_db()