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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion VERSION
Original file line number Diff line number Diff line change
@@ -1 +1 @@
2.2.1
2.3.0
6 changes: 3 additions & 3 deletions core/utils/date_utils.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

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

Expand Down
25 changes: 13 additions & 12 deletions metrics/counter/access/accumulation.py
Original file line number Diff line number Diff line change
@@ -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):
Expand All @@ -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,
Expand All @@ -30,16 +31,14 @@ 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"]

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 ""),
Expand All @@ -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)
Expand Down Expand Up @@ -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,
Expand Down
33 changes: 33 additions & 0 deletions metrics/counter/access/daily_accumulator.py
Original file line number Diff line number Diff line change
@@ -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)
27 changes: 13 additions & 14 deletions metrics/counter/indexing/converter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
7 changes: 7 additions & 0 deletions metrics/services/daily_metric_exports.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import logging
from time import monotonic

from metrics.models import DailyMetricJob
from metrics.opensearch.client import OpenSearchUsageClient
Expand Down Expand Up @@ -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.")
Expand All @@ -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,
)
27 changes: 20 additions & 7 deletions metrics/services/daily_payloads.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
22 changes: 21 additions & 1 deletion metrics/services/parsing/job_payloads.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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,
Expand All @@ -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

Expand Down Expand Up @@ -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,
Expand All @@ -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()
Expand Down
18 changes: 17 additions & 1 deletion metrics/tests/counter/access/test_accumulation.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
68 changes: 68 additions & 0 deletions metrics/tests/counter/access/test_daily_accumulator.py
Original file line number Diff line number Diff line change
@@ -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"]
Loading
Loading