diff --git a/VERSION b/VERSION index c043eea..276cbf9 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -2.2.1 +2.3.0 diff --git a/core/utils/date_utils.py b/core/utils/date_utils.py index 4f3df0e..49be589 100644 --- a/core/utils/date_utils.py +++ b/core/utils/date_utils.py @@ -105,7 +105,7 @@ def truncate_datetime_to_hour(dt): Returns: datetime: The truncated datetime object. """ - dt = _coerce_datetime(dt) + dt = coerce_datetime(dt) if dt is None: return None @@ -122,14 +122,14 @@ def extract_minute_second_key(dt): Returns: str: A string in the format "MM:SS" representing the minute and second. """ - dt = _coerce_datetime(dt) + dt = coerce_datetime(dt) if dt is None: return None return f"{dt.minute:02}:{dt.second:02}" -def _coerce_datetime(dt): +def coerce_datetime(dt): if isinstance(dt, datetime): return dt diff --git a/metrics/counter/access/accumulation.py b/metrics/counter/access/accumulation.py index bed2407..6b009d5 100644 --- a/metrics/counter/access/accumulation.py +++ b/metrics/counter/access/accumulation.py @@ -1,7 +1,7 @@ import re from urllib.parse import unquote, urlparse -from core.utils.date_utils import extract_minute_second_key, truncate_datetime_to_hour +from core.utils.date_utils import coerce_datetime def accumulate(results, counter_access, line): @@ -12,14 +12,15 @@ def accumulate(results, counter_access, line): client_name = line.get("client_name") client_version = line.get("client_version") - local_datetime = line.get("local_datetime") + local_datetime = coerce_datetime(line.get("local_datetime")) ip_address = line.get("ip_address") - access_datetime = truncate_datetime_to_hour(local_datetime) - ms_key = extract_minute_second_key(local_datetime) - if access_datetime is None or ms_key is None: + if local_datetime is None: raise ValueError("Invalid local_datetime in parsed log line.") + access_datetime = local_datetime.replace(minute=0, second=0, microsecond=0) + second_of_hour = local_datetime.minute * 60 + local_datetime.second + user_session_id = _generate_user_session_id( client_name, client_version, @@ -30,7 +31,7 @@ def accumulate(results, counter_access, line): counter_access=counter_access, line=line, access_datetime=access_datetime, - minute_second_key=ms_key, + second_of_hour=second_of_hour, user_session_id=user_session_id, ) item_access_id = raw_record["id"] @@ -38,8 +39,6 @@ def accumulate(results, counter_access, line): if item_access_id not in results: results[item_access_id] = raw_record["data"] - _increment_timestamp_count(results[item_access_id]["click_timestamps"], ms_key) - access_url_key = access_url or "|".join( [ str(counter_access.get("pid_generic") or ""), @@ -51,11 +50,15 @@ def accumulate(results, counter_access, line): "click_timestamps_by_url", {} ) url_timestamps = timestamps_by_url.setdefault(access_url_key, {}) - _increment_timestamp_count(url_timestamps, ms_key) + _increment_timestamp_count(url_timestamps, second_of_hour) def _build_record( - counter_access, line, access_datetime, minute_second_key, user_session_id + counter_access, + line, + access_datetime, + second_of_hour, + user_session_id, ): collection = counter_access.get("collection") source_key = _source_key(counter_access, collection) @@ -91,9 +94,7 @@ def _build_record( "document": _document_metadata(counter_access), "title_pid_generic": counter_access.get("title_pid_generic") or pid_generic, "user_session_id": user_session_id, - "click_timestamps": {minute_second_key: 0}, "click_timestamps_by_url": {}, - "access_url": counter_access.get("access_url"), "media_format": media_format, "content_language": content_language, "content_type": content_type, diff --git a/metrics/counter/access/daily_accumulator.py b/metrics/counter/access/daily_accumulator.py new file mode 100644 index 0000000..7032ccb --- /dev/null +++ b/metrics/counter/access/daily_accumulator.py @@ -0,0 +1,33 @@ +class DailyAccessAccumulator(dict): + def __init__(self): + super().__init__() + self._documents = {} + self._sources = {} + self._sessions = {} + + def __setitem__(self, key, value): + source_key = value.get("source_key") + source = value.get("source") + if source_key and source: + value["source"] = self._sources.setdefault(source_key, source) + + document_key = ( + value.get("document_type"), + value.get("pid_v2"), + value.get("pid_v3"), + value.get("pid_generic"), + value.get("title_pid_generic"), + ) + document = value.get("document") + document_identifiers = document_key[1:] + if document is not None and any(document_identifiers): + value["document"] = self._documents.setdefault(document_key, document) + + user_session_id = value.get("user_session_id") + if user_session_id: + value["user_session_id"] = self._sessions.setdefault( + user_session_id, + user_session_id, + ) + + super().__setitem__(key, value) diff --git a/metrics/counter/indexing/converter.py b/metrics/counter/indexing/converter.py index 4b4ab1f..79dba6c 100644 --- a/metrics/counter/indexing/converter.py +++ b/metrics/counter/indexing/converter.py @@ -18,27 +18,26 @@ def convert(data): if not isinstance(data, dict): return {"month": {}, "year": {}} - month_data = {} - month_unique_state = _initialize_unique_state() - year_data = {} - year_unique_state = _initialize_unique_state() + month_data = _convert_granularity(data, "month") + year_data = _convert_granularity(data, "year") + + return {"month": month_data, "year": year_data} + + +def _convert_granularity(data, granularity): + converted_data = {} + unique_state = _initialize_unique_state() for value in data.values(): pipeline = _get_pipeline(value) pipeline.accumulate( - data=month_data, - unique_state=month_unique_state, + data=converted_data, + unique_state=unique_state, value=value, - granularity="month", - ) - pipeline.accumulate( - data=year_data, - unique_state=year_unique_state, - value=value, - granularity="year", + granularity=granularity, ) - return {"month": month_data, "year": year_data} + return converted_data def _get_pipeline(value): diff --git a/metrics/services/daily_metric_exports.py b/metrics/services/daily_metric_exports.py index 8933b3d..7a14a21 100644 --- a/metrics/services/daily_metric_exports.py +++ b/metrics/services/daily_metric_exports.py @@ -1,4 +1,5 @@ import logging +from time import monotonic from metrics.models import DailyMetricJob from metrics.opensearch.client import OpenSearchUsageClient @@ -60,6 +61,7 @@ def _load_or_build_payload(job, track_errors, robots_source): def _export_payload(job, payload): + opensearch_started = monotonic() search_client = OpenSearchUsageClient() if not search_client.ping(): raise RuntimeError("OpenSearch client is not available.") @@ -69,3 +71,8 @@ def _export_payload(job, payload): job=job, payload=payload, ) + logging.info( + "Daily metric job %s OpenSearch export completed in %.3f seconds.", + job.pk, + monotonic() - opensearch_started, + ) diff --git a/metrics/services/daily_payloads.py b/metrics/services/daily_payloads.py index f908f1f..0db8f22 100644 --- a/metrics/services/daily_payloads.py +++ b/metrics/services/daily_payloads.py @@ -32,16 +32,29 @@ def write_payload(storage_path, payload): resolved_path = resolve_storage_path(storage_path) resolved_path.parent.mkdir(parents=True, exist_ok=True) - payload_json = json.dumps( - payload, ensure_ascii=True, sort_keys=True, separators=(",", ":") + encoder = json.JSONEncoder( + ensure_ascii=True, + sort_keys=True, + separators=(",", ":"), ) - payload_hash = hashlib.sha256(payload_json.encode("utf-8")).hexdigest() - + payload_hash = hashlib.sha256() tmp_path = resolved_path.with_suffix(f"{resolved_path.suffix}.tmp") - tmp_path.write_text(payload_json, encoding="utf-8") - tmp_path.replace(resolved_path) - return payload_hash + try: + with tmp_path.open("wb") as output: + for chunk in encoder.iterencode(payload): + encoded_chunk = chunk.encode("utf-8") + payload_hash.update(encoded_chunk) + output.write(encoded_chunk) + tmp_path.replace(resolved_path) + except Exception: + try: + tmp_path.unlink() + except FileNotFoundError: + pass + raise + + return payload_hash.hexdigest() def read_payload(storage_path): diff --git a/metrics/services/parsing/job_payloads.py b/metrics/services/parsing/job_payloads.py index fa30b3b..e3dc0ba 100644 --- a/metrics/services/parsing/job_payloads.py +++ b/metrics/services/parsing/job_payloads.py @@ -1,8 +1,10 @@ +import logging from time import monotonic from django.conf import settings from log_manager.models import LogFile +from metrics.counter.access.daily_accumulator import DailyAccessAccumulator from metrics.counter.indexing import converter as index_docs from metrics.services import daily_payloads from metrics.services.parsing.environment import setup_parsing_environment @@ -19,12 +21,13 @@ def build_daily_metric_job_payload(job, robots_list, mmdb, track_errors=False): input_log_hashes = sorted(job.input_log_hashes or []) log_files = _get_job_log_files(job, input_log_hashes) - results = {} + results = DailyAccessAccumulator() summary = _initial_summary(log_files, input_log_hashes) mark_logs_as_parsing(log_files) clear_discarded_lines(log_files) + parsing_started = monotonic() for log_file in log_files: log_summary = _parse_log_file_into_results( log_file=log_file, @@ -34,8 +37,19 @@ def build_daily_metric_job_payload(job, robots_list, mmdb, track_errors=False): track_errors=track_errors, ) _merge_log_summary(summary, log_summary) + logging.info( + "Daily metric job %s parsing completed in %.3f seconds.", + job.pk, + monotonic() - parsing_started, + ) + conversion_started = monotonic() documents = index_docs.convert(results) + logging.info( + "Daily metric job %s conversion completed in %.3f seconds.", + job.pk, + monotonic() - conversion_started, + ) payload = _write_job_payload(job, documents, summary) return payload @@ -125,6 +139,7 @@ def _merge_log_summary(summary, log_summary): def _write_job_payload(job, documents, summary): + serialization_started = monotonic() storage_path = daily_payloads.build_daily_storage_path( job.collection, job.access_date, @@ -137,6 +152,11 @@ def _write_job_payload(job, documents, summary): "summary": summary, } payload_hash = daily_payloads.write_payload(storage_path, payload) + logging.info( + "Daily metric job %s serialization completed in %.3f seconds.", + job.pk, + monotonic() - serialization_started, + ) job.input_log_hashes = summary["input_log_hashes"] job.storage_path = storage_path.as_posix() diff --git a/metrics/tests/counter/access/test_accumulation.py b/metrics/tests/counter/access/test_accumulation.py index ccf9044..a9eda44 100644 --- a/metrics/tests/counter/access/test_accumulation.py +++ b/metrics/tests/counter/access/test_accumulation.py @@ -161,7 +161,23 @@ def test_same_url_within_window_produces_single_url_bucket(self): raw = next(iter(results.values())) self.assertEqual( raw["click_timestamps_by_url"], - {"/id/c2248/03": {"00:05": 1, "00:20": 1}}, + {"/id/c2248/03": {5: 1, 20: 1}}, + ) + + def test_parses_datetime_string_and_stores_integer_seconds(self): + results = {} + + accumulation.accumulate( + results, + self._book_counter_access(), + self._line(local_datetime="2024-01-15 10:01:05"), + ) + + raw = next(iter(results.values())) + self.assertNotIn("click_timestamps", raw) + self.assertEqual( + raw["click_timestamps_by_url"], + {"BOOK:Q7GTD|html|full_text": {65: 1}}, ) def test_generates_session_id_from_client_ip_datetime(self): diff --git a/metrics/tests/counter/access/test_daily_accumulator.py b/metrics/tests/counter/access/test_daily_accumulator.py new file mode 100644 index 0000000..c02c84d --- /dev/null +++ b/metrics/tests/counter/access/test_daily_accumulator.py @@ -0,0 +1,68 @@ +from metrics.counter.access.daily_accumulator import DailyAccessAccumulator + + +def _record(session_id, **overrides): + record = { + "collection": "scl", + "source_key": "1234-5678", + "document_type": "article", + "pid_v2": "S123456782026000100001", + "pid_v3": "abc123", + "pid_generic": None, + "title_pid_generic": None, + "source": {"source_id": "1234-5678", "main_title": "Journal"}, + "document": {"title": "Article"}, + "user_session_id": session_id, + } + record.update(overrides) + return record + + +def test_interns_repeated_metadata_and_sessions(): + accumulator = DailyAccessAccumulator() + first_session = "|".join(["Firefox", "1", "127.0.0.1", "2026-08-25", "10"]) + second_session = "|".join(["Firefox", "1", "127.0.0.1", "2026-08-25", "10"]) + assert first_session is not second_session + + accumulator["first"] = _record(first_session) + accumulator["second"] = _record(second_session) + + assert accumulator["first"]["source"] is accumulator["second"]["source"] + assert accumulator["first"]["document"] is accumulator["second"]["document"] + assert ( + accumulator["first"]["user_session_id"] + is accumulator["second"]["user_session_id"] + ) + + +def test_interns_empty_document_metadata_for_the_same_document(): + accumulator = DailyAccessAccumulator() + + accumulator["first"] = _record("first", document={}) + accumulator["second"] = _record("second", document={}) + + assert accumulator["first"]["document"] is accumulator["second"]["document"] + + +def test_does_not_share_metadata_between_distinct_documents(): + accumulator = DailyAccessAccumulator() + + accumulator["first"] = _record("first", pid_v3="first", document={}) + accumulator["second"] = _record("second", pid_v3="second", document={}) + + assert accumulator["first"]["document"] is not accumulator["second"]["document"] + + +def test_does_not_share_documents_without_identifiers(): + accumulator = DailyAccessAccumulator() + identifiers = { + "pid_v2": None, + "pid_v3": None, + "pid_generic": None, + "title_pid_generic": None, + } + + accumulator["first"] = _record("first", document={}, **identifiers) + accumulator["second"] = _record("second", document={}, **identifiers) + + assert accumulator["first"]["document"] is not accumulator["second"]["document"] diff --git a/metrics/tests/services/test_daily_payloads.py b/metrics/tests/services/test_daily_payloads.py new file mode 100644 index 0000000..98dfff6 --- /dev/null +++ b/metrics/tests/services/test_daily_payloads.py @@ -0,0 +1,55 @@ +import hashlib +import json +import tempfile +from pathlib import Path + +from django.test import SimpleTestCase, override_settings + +from metrics.services import daily_payloads + + +class DailyPayloadTests(SimpleTestCase): + def setUp(self): + self.temporary_directory = tempfile.TemporaryDirectory() + self.settings_override = override_settings( + MEDIA_ROOT=self.temporary_directory.name + ) + self.settings_override.enable() + + def tearDown(self): + self.settings_override.disable() + self.temporary_directory.cleanup() + + def test_write_payload_preserves_canonical_bytes_and_hash(self): + payload = { + "collection": "scl", + "documents": {"รก": {"total_requests": 2}}, + "summary": {"valid_lines": 1}, + } + expected = json.dumps( + payload, + ensure_ascii=True, + sort_keys=True, + separators=(",", ":"), + ).encode("utf-8") + + payload_hash = daily_payloads.write_payload( + Path("scl/2026/08/2026-08-25.json"), + payload, + ) + + resolved_path = daily_payloads.resolve_storage_path( + Path("scl/2026/08/2026-08-25.json") + ) + self.assertEqual(resolved_path.read_bytes(), expected) + self.assertEqual(payload_hash, hashlib.sha256(expected).hexdigest()) + + def test_write_payload_removes_temporary_file_after_serialization_error(self): + storage_path = Path("scl/2026/08/2026-08-25.json") + + with self.assertRaises(TypeError): + daily_payloads.write_payload(storage_path, {"invalid": object()}) + + resolved_path = daily_payloads.resolve_storage_path(storage_path) + self.assertFalse(resolved_path.exists()) + self.assertFalse(resolved_path.with_suffix(".json.tmp").exists()) diff --git a/requirements/base.txt b/requirements/base.txt index cbc05bb..1228eda 100644 --- a/requirements/base.txt +++ b/requirements/base.txt @@ -69,7 +69,7 @@ git+https://github.com/scieloorg/scielo_log_validator@2.0.0#egg=scielo_log_valid git+https://github.com/scieloorg/scielo_scholarly_data@v0.1.4#egg=scielo_scholarly_data # SciELO Usage COUNTER -git+https://github.com/scieloorg/scielo_usage_counter@2.1.0#egg=scielo_usage_counter +git+https://github.com/scieloorg/scielo_usage_counter@2.2.0#egg=scielo_usage_counter # Device Detector device-detector==0.10 # https://github.com/thinkwelltwd/device_detector