diff --git a/VERSION b/VERSION index 8e8299d..437459c 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -2.4.2 +2.5.0 diff --git a/metrics/legacy_matomo/__init__.py b/metrics/legacy_matomo/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/metrics/legacy_matomo/importer.py b/metrics/legacy_matomo/importer.py new file mode 100644 index 0000000..a40e2d1 --- /dev/null +++ b/metrics/legacy_matomo/importer.py @@ -0,0 +1,282 @@ +import gzip +import json +from functools import partial +from itertools import islice + +from django.conf import settings +from opensearchpy import helpers + +from log_manager_config.choices import OpenSearchPartitionStrategy +from metrics.legacy_matomo.manifest import build_migration_id +from metrics.legacy_matomo.opensearch_actions import ( + build_analytics_increment_action, + build_counter_increment_action, +) +from metrics.legacy_matomo.routing import build_usage_index_target +from metrics.opensearch.names import generate_yearly_write_alias + +DOCUMENT_INDEX_LOOKUP_BATCH_SIZE = 1000 +DEFAULT_PROGRESS_INTERVAL = 100000 + + +def iter_document_items(path): + with gzip.open(path, "rt", encoding="utf-8") as source: + for line in source: + row = json.loads(line) + yield row["_id"], row["_source"] + + +def _collection_config(collection): + config = getattr(collection, "log_manager_config", None) + if config is None: + raise ValueError( + "Collection %s has no log manager configuration." % collection.acron3 + ) + return config + + +def _dataset_target(collection, config, manifest, dataset): + return build_usage_index_target( + index_prefix=settings.OPENSEARCH_INDEX_NAME, + collection=collection.acron3, + dataset=dataset, + partition_strategy=config.opensearch_partition_strategy, + access_date=manifest["month"] + "-01", + ) + + +def _inspect_alias(search_client, alias_name): + client = search_client.client + if not client.indices.exists_alias(name=alias_name): + return { + "alias": alias_name, + "exists": False, + "indexes": [], + "write_indexes": [], + } + + indexes = client.indices.get_alias(name=alias_name) + write_indexes = [ + index_name + for index_name, index_data in indexes.items() + if index_data.get("aliases", {}).get(alias_name, {}).get("is_write_index") + ] + return { + "alias": alias_name, + "exists": True, + "indexes": sorted(indexes), + "write_indexes": sorted(write_indexes), + } + + +def _write_index(search_client, alias_name): + alias = _inspect_alias(search_client, alias_name) + if len(alias["write_indexes"]) == 1: + return alias["write_indexes"][0] + if len(alias["indexes"]) == 1: + return alias["indexes"][0] + raise RuntimeError("Alias %s has no unique write index." % alias_name) + + +def _existing_document_indexes(search_client, alias_name, doc_ids): + response = search_client.client.search( + index=alias_name, + body={ + "size": len(doc_ids) * 2, + "_source": False, + "query": {"ids": {"values": doc_ids}}, + }, + ) + indexes = {} + for hit in response.get("hits", {}).get("hits", []): + doc_id = hit["_id"] + if doc_id in indexes and indexes[doc_id] != hit["_index"]: + raise RuntimeError( + "Document %s exists in multiple indexes behind alias %s." + % (doc_id, alias_name) + ) + indexes[doc_id] = hit["_index"] + return indexes + + +def _update_document_items( + search_client, + alias_name, + document_items, + action_builder, + dataset=None, + progress_callback=None, + progress_interval=DEFAULT_PROGRESS_INTERVAL, + stop_controller=None, +): + if progress_interval <= 0: + raise ValueError("Progress interval must be greater than zero.") + + write_index = _write_index(search_client, alias_name) + succeeded = 0 + reported = 0 + next_progress = progress_interval + batch_size = min( + search_client.bulk_chunk_size, + DOCUMENT_INDEX_LOOKUP_BATCH_SIZE, + ) + + while True: + if stop_controller: + stop_controller.raise_if_requested() + batch = list(islice(document_items, batch_size)) + if not batch: + if progress_callback and succeeded != reported: + progress_callback(dataset, "import", succeeded) + return succeeded + + existing_indexes = _existing_document_indexes( + search_client, + alias_name, + [doc_id for doc_id, _document in batch], + ) + batch_succeeded, _failed = helpers.bulk( + search_client.client, + ( + action_builder( + index_name=existing_indexes.get(doc_id, write_index), + doc_id=doc_id, + document=document, + ) + for doc_id, document in batch + ), + chunk_size=search_client.bulk_chunk_size, + ) + succeeded += batch_succeeded + if succeeded >= next_progress: + if progress_callback: + progress_callback(dataset, "import", succeeded) + reported = succeeded + while next_progress <= succeeded: + next_progress += progress_interval + if stop_controller: + stop_controller.raise_if_requested() + + +def _validate_alias_topology(search_client, target, access_date): + read_alias = _inspect_alias(search_client, target["read_alias"]) + write_alias = _inspect_alias(search_client, target["write_alias"]) + + if target["partition_strategy"] == OpenSearchPartitionStrategy.YEARLY: + if read_alias["write_indexes"]: + raise ValueError( + "Yearly target has a writable continuous alias: %s." + % target["read_alias"] + ) + else: + yearly_alias = generate_yearly_write_alias(target["read_alias"], access_date) + yearly = _inspect_alias(search_client, yearly_alias) + if yearly["exists"]: + raise ValueError( + "Rollover target conflicts with yearly alias: %s." % yearly_alias + ) + + if write_alias["exists"]: + has_unique_writer = len(write_alias["write_indexes"]) == 1 + has_single_implicit_writer = ( + not write_alias["write_indexes"] and len(write_alias["indexes"]) == 1 + ) + if not has_unique_writer and not has_single_implicit_writer: + raise ValueError( + "Target alias has no unique write index: %s." % target["write_alias"] + ) + + return { + "dataset": target["dataset"], + "partition_strategy": target["partition_strategy"], + "read_alias": read_alias, + "write_alias": write_alias, + } + + +def build_import_plan(search_client, collection, manifest): + config = _collection_config(collection) + access_date = manifest["month"] + "-01" + datasets = [] + for dataset in ("counter", "analytics"): + target = _dataset_target(collection, config, manifest, dataset) + datasets.append(_validate_alias_topology(search_client, target, access_date)) + + return { + "collection": collection.acron3, + "month": manifest["month"], + "primary_shards": config.opensearch_primary_shards, + "datasets": datasets, + } + + +def import_dataset( + search_client, + collection, + config, + manifest, + dataset, + progress_callback=None, + progress_interval=DEFAULT_PROGRESS_INTERVAL, + stop_controller=None, +): + target = _dataset_target(collection, config, manifest, dataset) + write_alias = search_client.prepare_usage_index( + alias_name=target["read_alias"], + mappings=target["mappings"], + partition_strategy=target["partition_strategy"], + access_date=manifest["month"] + "-01", + primary_shards=config.opensearch_primary_shards, + ) + action_builder = partial( + ( + build_analytics_increment_action + if dataset == "analytics" + else build_counter_increment_action + ), + migration_id=build_migration_id(manifest, dataset), + source_days=manifest["source_days"], + ) + imported = _update_document_items( + search_client=search_client, + alias_name=write_alias, + document_items=iter_document_items(manifest[dataset]["path"]), + action_builder=action_builder, + dataset=dataset, + progress_callback=progress_callback, + progress_interval=progress_interval, + stop_controller=stop_controller, + ) + search_client.rollover_usage_index( + alias_name=write_alias, + mappings=target["mappings"], + primary_shards=config.opensearch_primary_shards, + read_alias=( + target["read_alias"] if write_alias != target["read_alias"] else None + ), + ) + return imported + + +def import_manifest( + search_client, + collection, + manifest, + progress_callback=None, + progress_interval=DEFAULT_PROGRESS_INTERVAL, + stop_controller=None, +): + config = _collection_config(collection) + return { + dataset: import_dataset( + search_client, + collection, + config, + manifest, + dataset, + progress_callback=progress_callback, + progress_interval=progress_interval, + stop_controller=stop_controller, + ) + for dataset in ("counter", "analytics") + } diff --git a/metrics/legacy_matomo/manifest.py b/metrics/legacy_matomo/manifest.py new file mode 100644 index 0000000..ee5f3b1 --- /dev/null +++ b/metrics/legacy_matomo/manifest.py @@ -0,0 +1,67 @@ +import json +from datetime import date + +ARTICLE_METRIC_PROFILE = "article-usage-v1" +DEFAULT_HISTORICAL_CUTOFF = date(2025, 12, 31) +HISTORICAL_CUTOFFS = { + "scl": date(2026, 6, 30), +} +MANIFEST_SCHEMA_VERSION = 2 +MIGRATION_NAME = "legacy-matomo" +SOURCE_COLLECTION_ALIASES = { + ("nbr", "scl"), +} +UNSUPPORTED_COLLECTIONS = {"books", "data"} + + +def load_manifest(path): + with open(path, encoding="utf-8") as source: + return json.load(source) + + +def build_migration_id(manifest, dataset): + file_hash = manifest[dataset]["uncompressed_sha256"] + return "%s:v%d:%s:%s:%s:%s" % ( + MIGRATION_NAME, + MANIFEST_SCHEMA_VERSION, + manifest["target"]["collection"], + manifest["month"], + dataset, + file_hash, + ) + + +def validate_migration_scope(manifest): + if manifest.get("metric_profile") != ARTICLE_METRIC_PROFILE: + raise ValueError("Unsupported migration metric profile.") + + source_collection = manifest["source"]["collection"] + target_collection = manifest["target"]["collection"] + if target_collection in UNSUPPORTED_COLLECTIONS: + raise ValueError( + "Collection %s is not supported by the article migration." + % target_collection + ) + if ( + source_collection != target_collection + and (source_collection, target_collection) not in SOURCE_COLLECTION_ALIASES + ): + raise ValueError( + "Unsupported source-to-target collection mapping: %s -> %s." + % (source_collection, target_collection) + ) + + historical_cutoff = HISTORICAL_CUTOFFS.get( + target_collection, + DEFAULT_HISTORICAL_CUTOFF, + ) + invalid = [ + value + for value in manifest["source_days"] + if date.fromisoformat(value) > historical_cutoff + ] + if invalid: + raise ValueError( + "Historical import for %s is limited to %s; found %s." + % (target_collection, historical_cutoff.isoformat(), invalid[0]) + ) diff --git a/metrics/legacy_matomo/opensearch_actions.py b/metrics/legacy_matomo/opensearch_actions.py new file mode 100644 index 0000000..7adaa95 --- /dev/null +++ b/metrics/legacy_matomo/opensearch_actions.py @@ -0,0 +1,162 @@ +from datetime import date + +from metrics.opensearch.painless import ( + ANNUAL_MASK_BITS, + ANNUAL_MASK_BUCKETS, + METRIC_FIELDS, +) + +IDEMPOTENT_COUNTER_INCREMENT_SCRIPT = """ +if (ctx._source.applied_migrations == null) { + ctx._source.applied_migrations = []; +} +if (ctx._source.applied_migrations.contains(params.migration_id)) { + ctx.op = 'none'; + return; +} +if (ctx._source.applied_days == null) { + ctx._source.applied_days = []; +} +for (day in params.source_days) { + if (ctx._source.applied_days.contains(day)) { + throw new IllegalStateException('Historical migration overlaps applied day ' + day); + } +} +for (entry in params.document.entrySet()) { + if (!params.metric_fields.contains(entry.getKey()) + && !'applied_days'.equals(entry.getKey()) + && !'applied_migrations'.equals(entry.getKey()) + && !'daily_metrics'.equals(entry.getKey())) { + if (!ctx._source.containsKey(entry.getKey()) || ctx._source[entry.getKey()] != entry.getValue()) { + ctx._source[entry.getKey()] = entry.getValue(); + } + } +} +for (field in params.metric_fields) { + def currentValue = ctx._source.containsKey(field) ? ctx._source[field] : 0; + def increment = params.document.containsKey(field) ? params.document[field] : 0; + ctx._source[field] = currentValue + increment; +} +if (ctx._source.daily_metrics == null) { + ctx._source.daily_metrics = new HashMap(); +} +for (dayEntry in params.document.daily_metrics.entrySet()) { + ctx._source.daily_metrics[dayEntry.getKey()] = dayEntry.getValue(); +} +ctx._source.applied_days.addAll(params.source_days); +ctx._source.applied_migrations.add(params.migration_id); +""" + +IDEMPOTENT_ANALYTICS_INCREMENT_SCRIPT = """ +if (ctx._source.applied_migrations == null) { + ctx._source.applied_migrations = []; +} +if (ctx._source.applied_migrations.contains(params.migration_id)) { + ctx.op = 'none'; + return; +} +if (ctx._source.applied_day_masks == null) { + ctx._source.applied_day_masks = params.empty_day_masks; +} +for (int index = 0; index < params.day_masks.size(); index++) { + if ((ctx._source.applied_day_masks[index] & params.day_masks[index]) != 0) { + throw new IllegalStateException('Historical migration overlaps an applied analytics day'); + } +} +for (entry in params.document.entrySet()) { + if (!params.metric_fields.contains(entry.getKey()) + && !'applied_day_masks'.equals(entry.getKey()) + && !'applied_migrations'.equals(entry.getKey())) { + if (!ctx._source.containsKey(entry.getKey()) || ctx._source[entry.getKey()] != entry.getValue()) { + ctx._source[entry.getKey()] = entry.getValue(); + } + } +} +for (field in params.metric_fields) { + def currentValue = ctx._source.containsKey(field) ? ctx._source[field] : 0; + def increment = params.document.containsKey(field) ? params.document[field] : 0; + ctx._source[field] = currentValue + increment; +} +for (int index = 0; index < params.day_masks.size(); index++) { + ctx._source.applied_day_masks[index] |= params.day_masks[index]; +} +ctx._source.applied_migrations.add(params.migration_id); +""" + + +def build_counter_increment_action( + index_name, + doc_id, + document, + migration_id, + source_days, +): + month = document["month"] + document_source_days = [ + "%s-%s" % (month, day) for day in sorted(document["daily_metrics"]) + ] + if set(document_source_days) - set(source_days): + raise ValueError( + "Counter document contains days outside the migration source period." + ) + + return { + "_op_type": "update", + "_index": index_name, + "_id": doc_id, + "retry_on_conflict": 5, + "scripted_upsert": True, + "script": { + "lang": "painless", + "source": IDEMPOTENT_COUNTER_INCREMENT_SCRIPT, + "params": { + "document": document, + "migration_id": migration_id, + "source_days": document_source_days, + "metric_fields": list(METRIC_FIELDS), + }, + }, + "upsert": { + "applied_days": [], + "applied_migrations": [], + }, + } + + +def build_analytics_increment_action( + index_name, + doc_id, + document, + migration_id, + source_days, +): + day_masks = [0] * ANNUAL_MASK_BUCKETS + for access_day in source_days: + access_date = date.fromisoformat(access_day) + day_offset = access_date.timetuple().tm_yday - 1 + day_masks[day_offset // ANNUAL_MASK_BITS] |= 1 << ( + day_offset % ANNUAL_MASK_BITS + ) + + return { + "_op_type": "update", + "_index": index_name, + "_id": doc_id, + "retry_on_conflict": 5, + "scripted_upsert": True, + "script": { + "lang": "painless", + "source": IDEMPOTENT_ANALYTICS_INCREMENT_SCRIPT, + "params": { + "document": document, + "migration_id": migration_id, + "day_masks": day_masks, + "empty_day_masks": [0] * ANNUAL_MASK_BUCKETS, + "metric_fields": list(METRIC_FIELDS), + }, + }, + "upsert": { + "applied_day_masks": [0] * ANNUAL_MASK_BUCKETS, + "applied_migrations": [], + }, + } diff --git a/metrics/legacy_matomo/operations.py b/metrics/legacy_matomo/operations.py new file mode 100644 index 0000000..20d6ed3 --- /dev/null +++ b/metrics/legacy_matomo/operations.py @@ -0,0 +1,59 @@ +import json +import os +import signal +from pathlib import Path + + +class MigrationInterrupted(RuntimeError): + pass + + +class StopController: + def __init__(self): + self.requested = False + self.previous_handlers = {} + + def install(self): + for signal_number in (signal.SIGINT, signal.SIGTERM): + self.previous_handlers[signal_number] = signal.getsignal(signal_number) + signal.signal(signal_number, self._request_stop) + + def restore(self): + for signal_number, handler in self.previous_handlers.items(): + signal.signal(signal_number, handler) + self.previous_handlers = {} + + def _request_stop(self, _signal_number, _frame): + self.requested = True + + def raise_if_requested(self): + if self.requested: + raise MigrationInterrupted( + "Migration interruption requested; the current batch was completed." + ) + + +def write_json_atomic(path, value): + destination = Path(path) + destination.parent.mkdir(parents=True, exist_ok=True) + temporary = destination.with_name(".%s.%d.tmp" % (destination.name, os.getpid())) + try: + with open(temporary, "w", encoding="utf-8") as output: + json.dump(value, output, ensure_ascii=False, indent=2, sort_keys=True) + output.write("\n") + output.flush() + os.fsync(output.fileno()) + os.replace(temporary, destination) + finally: + if temporary.exists(): + temporary.unlink() + + +def validate_report_destination(path, protected_paths): + if not path: + return + + destination = Path(path).resolve() + protected = {Path(value).resolve() for value in protected_paths} + if destination in protected: + raise ValueError("Report path would overwrite a migration input file.") diff --git a/metrics/legacy_matomo/routing.py b/metrics/legacy_matomo/routing.py new file mode 100644 index 0000000..9aec996 --- /dev/null +++ b/metrics/legacy_matomo/routing.py @@ -0,0 +1,53 @@ +from copy import deepcopy + +from log_manager_config.choices import OpenSearchPartitionStrategy +from metrics.opensearch.mappings import get_index_mappings +from metrics.opensearch.names import ( + generate_analytics_index_name, + generate_month_index_name, + generate_yearly_write_alias, +) + +MIGRATION_TRACKING_MAPPING = { + "type": "keyword", + "index": False, + "doc_values": False, +} + + +def _migration_mappings(dataset): + mappings = deepcopy(get_index_mappings(dataset)) + mappings["properties"]["applied_migrations"] = MIGRATION_TRACKING_MAPPING + return mappings + + +def build_usage_index_target( + index_prefix, + collection, + dataset, + partition_strategy, + access_date, +): + if dataset == "counter": + read_alias = generate_month_index_name(index_prefix, collection) + elif dataset == "analytics": + read_alias = generate_analytics_index_name(index_prefix, collection) + else: + raise ValueError("Unsupported usage dataset: %s." % dataset) + + if partition_strategy == OpenSearchPartitionStrategy.YEARLY: + write_alias = generate_yearly_write_alias(read_alias, access_date) + elif partition_strategy == OpenSearchPartitionStrategy.ROLLOVER: + write_alias = read_alias + else: + raise ValueError( + "Unsupported OpenSearch partition strategy: %s." % partition_strategy + ) + + return { + "dataset": dataset, + "mappings": _migration_mappings(dataset), + "partition_strategy": partition_strategy, + "read_alias": read_alias, + "write_alias": write_alias, + } diff --git a/metrics/legacy_matomo/validation.py b/metrics/legacy_matomo/validation.py new file mode 100644 index 0000000..515ea2a --- /dev/null +++ b/metrics/legacy_matomo/validation.py @@ -0,0 +1,298 @@ +import gzip +import hashlib +import json +import os +import sqlite3 +import tempfile +from datetime import date +from pathlib import Path + +from metrics.legacy_matomo.manifest import ( + MANIFEST_SCHEMA_VERSION, + MIGRATION_NAME, + validate_migration_scope, +) +from metrics.opensearch.keys import metric_key +from metrics.opensearch.painless import METRIC_FIELDS + +ARTICLE_METRIC_DIMENSIONS = { + "metric_scope": "item", + "data_type": "Article", + "parent_data_type": "Journal", + "article_version": None, + "access_type": "Open", + "access_method": "Regular", +} +SQLITE_INSERT_BATCH_SIZE = 10000 +DEFAULT_PROGRESS_INTERVAL = 100000 + + +class UniqueIdStore: + def __init__(self, temporary_directory=None): + descriptor, self.path = tempfile.mkstemp( + prefix="legacy-matomo-ids-", + suffix=".sqlite3", + dir=temporary_directory, + ) + os.close(descriptor) + self.connection = sqlite3.connect(self.path) + self.connection.execute("PRAGMA journal_mode=OFF") + self.connection.execute("PRAGMA synchronous=OFF") + self.connection.execute( + "CREATE TABLE ids (value TEXT PRIMARY KEY) WITHOUT ROWID" + ) + self.pending = [] + + def add(self, value): + self.pending.append((value,)) + if len(self.pending) >= SQLITE_INSERT_BATCH_SIZE: + self.flush() + + def flush(self): + if not self.pending: + return + try: + self.connection.executemany("INSERT INTO ids VALUES (?)", self.pending) + except sqlite3.IntegrityError as exc: + raise ValueError("Migration payload contains a duplicate _id.") from exc + self.pending = [] + + def close(self): + try: + self.flush() + finally: + self.connection.close() + os.unlink(self.path) + + +def _metric_values(document): + values = [] + for field in METRIC_FIELDS: + value = document.get(field) + if not isinstance(value, int) or isinstance(value, bool) or value < 0: + raise ValueError("Invalid metric value for %s." % field) + values.append(value) + return values + + +def _add_metrics(target, values): + for index, value in enumerate(values): + target[index] += value + + +def _validate_dimensions(document): + for field, expected in ARTICLE_METRIC_DIMENSIONS.items(): + if document.get(field) != expected: + raise ValueError("Invalid counter dimension %s." % field) + + +def _expected_metric_id(collection, document, period, analytics): + return metric_key( + collection=collection, + source_identifier=document["source_key"], + document_identifier=document["document_key"], + period=period, + projection="country_language" if analytics else None, + country_code=document.get("country_code"), + content_language=document.get("content_language"), + **ARTICLE_METRIC_DIMENSIONS, + ) + + +def validate_document_file( + file_summary, + collection, + source_days, + analytics=False, + temporary_directory=None, + progress_callback=None, + progress_interval=DEFAULT_PROGRESS_INTERVAL, + stop_controller=None, +): + path = Path(file_summary["path"]) + expected_dataset = "analytics" if analytics else "counter" + if file_summary.get("dataset") != expected_dataset: + raise ValueError("Invalid dataset summary for %s." % path) + + expected_period = ( + file_summary["source_month"][:4] if analytics else file_summary["month"] + ) + source_day_suffixes = {value[-2:] for value in source_days} + totals = [0] * len(METRIC_FIELDS) + payload_hash = hashlib.sha256() + documents = 0 + unique_ids = UniqueIdStore(temporary_directory) + + try: + with gzip.open(path, "rb") as source: + for line_number, line in enumerate(source, 1): + payload_hash.update(line) + row = json.loads(line) + document = row.get("_source") + if not row.get("_id") or not isinstance(document, dict): + raise ValueError( + "Invalid JSONL row at %s:%d." % (path, line_number) + ) + unique_ids.add(row["_id"]) + _validate_dimensions(document) + if row["_id"] != _expected_metric_id( + collection, document, expected_period, analytics + ): + raise ValueError( + "Invalid metric key at %s:%d." % (path, line_number) + ) + + if analytics: + _validate_analytics_document( + document, expected_period, path, line_number + ) + else: + _validate_counter_document( + document, + collection, + expected_period, + source_day_suffixes, + path, + line_number, + ) + + _add_metrics(totals, _metric_values(document)) + documents += 1 + if documents % progress_interval == 0: + if progress_callback: + progress_callback(expected_dataset, "validation", documents) + if stop_controller: + stop_controller.raise_if_requested() + if progress_callback and documents % progress_interval: + progress_callback(expected_dataset, "validation", documents) + if stop_controller: + stop_controller.raise_if_requested() + finally: + unique_ids.close() + + actual_totals = dict(zip(METRIC_FIELDS, totals)) + if documents != file_summary["documents"]: + raise ValueError("Document count mismatch for %s." % path) + if actual_totals != file_summary["totals"]: + raise ValueError("Metric totals mismatch for %s." % path) + if payload_hash.hexdigest() != file_summary["uncompressed_sha256"]: + raise ValueError("SHA-256 mismatch for %s." % path) + + return { + "path": str(path), + "documents": documents, + "totals": actual_totals, + "uncompressed_sha256": payload_hash.hexdigest(), + } + + +def _validate_analytics_document(document, expected_year, path, line_number): + if document.get("year") != expected_year: + raise ValueError("Invalid analytics year at %s:%d." % (path, line_number)) + if not document.get("country_code") or not document.get("content_language"): + raise ValueError( + "Incomplete analytics projection at %s:%d." % (path, line_number) + ) + + +def _validate_counter_document( + document, + collection, + expected_month, + source_day_suffixes, + path, + line_number, +): + if document.get("collection") != collection: + raise ValueError("Invalid collection at %s:%d." % (path, line_number)) + if document.get("month") != expected_month: + raise ValueError("Invalid counter month at %s:%d." % (path, line_number)) + daily_metrics = document.get("daily_metrics") + if not isinstance(daily_metrics, dict) or not daily_metrics: + raise ValueError( + "Missing counter daily metrics at %s:%d." % (path, line_number) + ) + if not set(daily_metrics).issubset(source_day_suffixes): + raise ValueError("Unexpected daily metric at %s:%d." % (path, line_number)) + + daily_totals = [0] * len(METRIC_FIELDS) + for values in daily_metrics.values(): + _add_metrics(daily_totals, _metric_values(values)) + if daily_totals != _metric_values(document): + raise ValueError( + "Counter daily totals mismatch at %s:%d." % (path, line_number) + ) + + +def validate_manifest( + manifest, + temporary_directory=None, + progress_callback=None, + progress_interval=DEFAULT_PROGRESS_INTERVAL, + stop_controller=None, +): + if progress_interval <= 0: + raise ValueError("Progress interval must be greater than zero.") + if manifest.get("schema_version") != MANIFEST_SCHEMA_VERSION: + raise ValueError("Unsupported migration manifest schema version.") + if manifest.get("migration") != MIGRATION_NAME: + raise ValueError("Unsupported migration manifest.") + validate_migration_scope(manifest) + + collection = manifest["target"]["collection"] + month = manifest["month"] + date.fromisoformat(month + "-01") + source_days = manifest["source_days"] + if not source_days: + raise ValueError("Migration manifest has no source days.") + if source_days != sorted(set(source_days)): + raise ValueError("Migration source days must be unique and sorted.") + for value in source_days: + date.fromisoformat(value) + if value[:7] != month: + raise ValueError("Source day outside migration month: %s." % value) + + empty_days = manifest.get("empty_days", []) + if set(source_days) & set(empty_days): + raise ValueError("Source days and empty days overlap.") + for value in empty_days: + date.fromisoformat(value) + if value[:7] != month: + raise ValueError("Empty day outside migration month: %s." % value) + + if manifest["counter"].get("month") != month: + raise ValueError("Counter summary month differs from manifest.") + if manifest["analytics"].get("source_month") != month: + raise ValueError("Analytics summary month differs from manifest.") + + counter = validate_document_file( + manifest["counter"], + collection, + source_days, + analytics=False, + temporary_directory=temporary_directory, + progress_callback=progress_callback, + progress_interval=progress_interval, + stop_controller=stop_controller, + ) + analytics = validate_document_file( + manifest["analytics"], + collection, + source_days, + analytics=True, + temporary_directory=temporary_directory, + progress_callback=progress_callback, + progress_interval=progress_interval, + stop_controller=stop_controller, + ) + if counter["totals"] != analytics["totals"]: + raise ValueError("Counter and analytics totals do not reconcile.") + + return { + "status": "valid", + "collection": collection, + "month": month, + "source_days": len(source_days), + "counter": counter, + "analytics": analytics, + } diff --git a/metrics/management/__init__.py b/metrics/management/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/metrics/management/commands/__init__.py b/metrics/management/commands/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/metrics/management/commands/import_legacy_matomo.py b/metrics/management/commands/import_legacy_matomo.py new file mode 100644 index 0000000..ed68fd9 --- /dev/null +++ b/metrics/management/commands/import_legacy_matomo.py @@ -0,0 +1,172 @@ +import json +from datetime import datetime, timezone +from time import monotonic + +from django.core.exceptions import ObjectDoesNotExist +from django.core.management.base import BaseCommand, CommandError + +from collection.models import Collection +from metrics.legacy_matomo.importer import build_import_plan, import_manifest +from metrics.legacy_matomo.manifest import load_manifest +from metrics.legacy_matomo.operations import ( + MigrationInterrupted, + StopController, + validate_report_destination, + write_json_atomic, +) +from metrics.legacy_matomo.validation import validate_manifest +from metrics.opensearch.client import OpenSearchUsageClient + + +class Command(BaseCommand): + help = "Validate or import a legacy Matomo migration manifest." + + def add_arguments(self, parser): + parser.add_argument("--manifest", required=True) + mode = parser.add_mutually_exclusive_group() + mode.add_argument("--preflight", action="store_true") + mode.add_argument("--execute", action="store_true") + parser.add_argument("--report") + parser.add_argument("--temporary-directory") + parser.add_argument("--progress-every", type=int, default=100000) + + def _progress(self, dataset, stage, documents): + self.stdout.write("[%s] %s: %d documents" % (stage, dataset, documents)) + self.stdout.flush() + + def _finish_report(self, report, path, status, started, error=None): + report["status"] = status + report["finished_at"] = datetime.now(timezone.utc).isoformat() + report["duration_seconds"] = round(monotonic() - started, 3) + if error: + report["error"] = str(error) + if path: + write_json_atomic(path, report) + + def handle(self, *args, **options): + started = monotonic() + report = { + "manifest": options["manifest"], + "started_at": datetime.now(timezone.utc).isoformat(), + } + report_path = None + stop_controller = StopController() + stop_controller.install() + + try: + manifest = load_manifest(options["manifest"]) + validate_report_destination( + options["report"], + [ + options["manifest"], + manifest["counter"]["path"], + manifest["analytics"]["path"], + ], + ) + report_path = options["report"] + report["collection"] = manifest.get("target", {}).get("collection") + report["month"] = manifest.get("month") + validation = validate_manifest( + manifest, + temporary_directory=options["temporary_directory"], + progress_callback=self._progress, + progress_interval=options["progress_every"], + stop_controller=stop_controller, + ) + report["validation"] = validation + + if not options["preflight"] and not options["execute"]: + self.stdout.write(json.dumps(validation, indent=2, sort_keys=True)) + self.stdout.write( + self.style.SUCCESS( + "Dry-run validation completed; nothing imported." + ) + ) + self._finish_report( + report, + report_path, + "validated", + started, + ) + return + + collection = Collection.objects.get(acron3=manifest["target"]["collection"]) + search_client = OpenSearchUsageClient() + if not search_client.ping(): + raise ValueError("OpenSearch preflight ping failed.") + plan = build_import_plan(search_client, collection, manifest) + report["plan"] = plan + if not options["execute"]: + self.stdout.write( + json.dumps( + {"validation": validation, "plan": plan}, + indent=2, + sort_keys=True, + ) + ) + self.stdout.write( + self.style.SUCCESS( + "Read-only preflight completed; nothing imported." + ) + ) + self._finish_report( + report, + report_path, + "preflight-completed", + started, + ) + return + imported = import_manifest( + search_client, + collection, + manifest, + progress_callback=self._progress, + progress_interval=options["progress_every"], + stop_controller=stop_controller, + ) + report["imported"] = imported + self._finish_report( + report, + report_path, + "completed", + started, + ) + except MigrationInterrupted as exc: + self._finish_report( + report, + report_path, + "stopped", + started, + error=exc, + ) + raise CommandError(str(exc)) from exc + except ( + OSError, + Collection.DoesNotExist, + ObjectDoesNotExist, + KeyError, + TypeError, + ValueError, + ) as exc: + self._finish_report( + report, + report_path, + "failed", + started, + error=exc, + ) + raise CommandError(str(exc)) from exc + except Exception as exc: + self._finish_report( + report, + report_path, + "failed", + started, + error=exc, + ) + raise + finally: + stop_controller.restore() + + self.stdout.write(json.dumps(imported, indent=2, sort_keys=True)) + self.stdout.write(self.style.SUCCESS("Historical Matomo import completed.")) diff --git a/metrics/tests/opensearch/test_routing.py b/metrics/tests/opensearch/test_routing.py new file mode 100644 index 0000000..4754150 --- /dev/null +++ b/metrics/tests/opensearch/test_routing.py @@ -0,0 +1,39 @@ +from django.test import SimpleTestCase + +from metrics.legacy_matomo.routing import build_usage_index_target +from metrics.opensearch.mappings import MONTH_INDEX_MAPPINGS + + +class UsageIndexRoutingTests(SimpleTestCase): + def test_yearly_strategy_partitions_both_datasets_by_access_year(self): + counter = build_usage_index_target( + "usage", "scl", "counter", "yearly", "2025-08-01" + ) + analytics = build_usage_index_target( + "usage", "scl", "analytics", "yearly", "2025-08-01" + ) + + self.assertEqual(counter["read_alias"], "usage_monthly_scl") + self.assertEqual(counter["write_alias"], "usage_monthly_scl_2025") + self.assertIn("applied_migrations", counter["mappings"]["properties"]) + self.assertNotIn("applied_migrations", MONTH_INDEX_MAPPINGS["properties"]) + self.assertEqual(analytics["read_alias"], "usage_yearly_analytics_scl") + self.assertEqual( + analytics["write_alias"], + "usage_yearly_analytics_scl_2025", + ) + + def test_rollover_strategy_uses_continuous_aliases(self): + counter = build_usage_index_target( + "usage", "books", "counter", "rollover", "2025-08-01" + ) + analytics = build_usage_index_target( + "usage", "books", "analytics", "rollover", "2025-08-01" + ) + + self.assertEqual(counter["write_alias"], counter["read_alias"]) + self.assertEqual(analytics["write_alias"], analytics["read_alias"]) + + def test_rejects_unknown_dataset(self): + with self.assertRaisesMessage(ValueError, "Unsupported usage dataset"): + build_usage_index_target("usage", "scl", "unknown", "yearly", "2025-08-01") diff --git a/metrics/tests/services/legacy_matomo.py b/metrics/tests/services/legacy_matomo.py new file mode 100644 index 0000000..f651df7 --- /dev/null +++ b/metrics/tests/services/legacy_matomo.py @@ -0,0 +1,116 @@ +import gzip +import hashlib +import json +import tempfile +from pathlib import Path + +from django.test import SimpleTestCase + +from metrics.legacy_matomo.manifest import ARTICLE_METRIC_PROFILE +from metrics.opensearch.keys import document_key, metric_key, source_key + +ARTICLE_DIMENSIONS = { + "metric_scope": "item", + "data_type": "Article", + "parent_data_type": "Journal", + "article_version": None, + "access_type": "Open", + "access_method": "Regular", +} +METRICS = { + "total_requests": 3, + "total_investigations": 4, + "unique_requests": 1, + "unique_investigations": 2, +} + + +def write_jsonl(path, metric_id, document): + encoded = ( + json.dumps( + {"_id": metric_id, "_source": document}, + separators=(",", ":"), + sort_keys=True, + ) + + "\n" + ).encode() + with gzip.GzipFile(filename=path, mode="wb", mtime=0) as output: + output.write(encoded) + return { + "path": str(path), + "documents": 1, + "totals": METRICS, + "uncompressed_sha256": hashlib.sha256(encoded).hexdigest(), + } + + +class LegacyManifestTestCase(SimpleTestCase): + def setUp(self): + self.temporary_directory = tempfile.TemporaryDirectory() + root = Path(self.temporary_directory.name) + source_identifier = source_key("scl", "journal", "0103-2100") + document_identifier = document_key("scl", "article", "pid") + counter_document = { + "collection": "scl", + "source_key": source_identifier, + "document_key": document_identifier, + "month": "2025-08", + **{ + key: value + for key, value in ARTICLE_DIMENSIONS.items() + if value is not None + }, + **METRICS, + "daily_metrics": {"01": METRICS}, + } + counter_id = metric_key( + collection="scl", + source_identifier=source_identifier, + document_identifier=document_identifier, + period="2025-08", + **ARTICLE_DIMENSIONS, + ) + analytics_document = { + "year": "2025", + "source_key": source_identifier, + "document_key": document_identifier, + **{ + key: value + for key, value in ARTICLE_DIMENSIONS.items() + if value is not None + }, + "country_code": "BR", + "content_language": "pt", + **METRICS, + } + analytics_id = metric_key( + collection="scl", + source_identifier=source_identifier, + document_identifier=document_identifier, + period="2025", + projection="country_language", + country_code="BR", + content_language="pt", + **ARTICLE_DIMENSIONS, + ) + counter = write_jsonl(root / "counter.jsonl.gz", counter_id, counter_document) + counter.update({"month": "2025-08", "dataset": "counter"}) + analytics = write_jsonl( + root / "analytics.jsonl.gz", analytics_id, analytics_document + ) + analytics.update({"source_month": "2025-08", "dataset": "analytics"}) + self.manifest = { + "schema_version": 2, + "migration": "legacy-matomo", + "metric_profile": ARTICLE_METRIC_PROFILE, + "source": {"database": "matomo", "collection": "nbr"}, + "month": "2025-08", + "source_days": ["2025-08-01"], + "empty_days": [], + "target": {"collection": "scl"}, + "counter": counter, + "analytics": analytics, + } + + def tearDown(self): + self.temporary_directory.cleanup() diff --git a/metrics/tests/services/test_import_legacy_matomo_command.py b/metrics/tests/services/test_import_legacy_matomo_command.py new file mode 100644 index 0000000..8c04d39 --- /dev/null +++ b/metrics/tests/services/test_import_legacy_matomo_command.py @@ -0,0 +1,67 @@ +import json +from io import StringIO +from pathlib import Path + +from django.core.management import CommandError, call_command + +from metrics.legacy_matomo.operations import ( + validate_report_destination, + write_json_atomic, +) +from metrics.tests.services.legacy_matomo import LegacyManifestTestCase + + +class ImportLegacyMatomoCommandTests(LegacyManifestTestCase): + def test_writes_report_atomically(self): + report_path = Path(self.temporary_directory.name) / "reports" / "result.json" + + write_json_atomic(report_path, {"status": "completed"}) + + with open(report_path, encoding="utf-8") as source: + self.assertEqual(json.load(source), {"status": "completed"}) + self.assertEqual(list(report_path.parent.glob("*.tmp")), []) + + def test_rejects_report_that_would_overwrite_an_input(self): + input_path = Path(self.manifest["counter"]["path"]) + + with self.assertRaisesMessage(ValueError, "overwrite"): + validate_report_destination(input_path, [input_path]) + + def test_management_command_writes_validation_report(self): + root = Path(self.temporary_directory.name) + manifest_path = root / "manifest.json" + report_path = root / "report.json" + work_path = root / "work" + work_path.mkdir() + with open(manifest_path, "w", encoding="utf-8") as output: + json.dump(self.manifest, output) + + call_command( + "import_legacy_matomo", + manifest=str(manifest_path), + report=str(report_path), + temporary_directory=str(work_path), + progress_every=1, + stdout=StringIO(), + ) + + with open(report_path, encoding="utf-8") as source: + report = json.load(source) + self.assertEqual(report["status"], "validated") + self.assertEqual(report["collection"], "scl") + + def test_management_command_does_not_overwrite_manifest_with_report(self): + manifest_path = Path(self.temporary_directory.name) / "manifest.json" + with open(manifest_path, "w", encoding="utf-8") as output: + json.dump(self.manifest, output) + original = manifest_path.read_bytes() + + with self.assertRaisesMessage(CommandError, "overwrite"): + call_command( + "import_legacy_matomo", + manifest=str(manifest_path), + report=str(manifest_path), + stdout=StringIO(), + ) + + self.assertEqual(manifest_path.read_bytes(), original) diff --git a/metrics/tests/services/test_legacy_matomo_importer.py b/metrics/tests/services/test_legacy_matomo_importer.py new file mode 100644 index 0000000..4ff250c --- /dev/null +++ b/metrics/tests/services/test_legacy_matomo_importer.py @@ -0,0 +1,157 @@ +from types import SimpleNamespace +from unittest.mock import Mock, patch + +from django.test import override_settings + +from metrics.legacy_matomo import importer +from metrics.legacy_matomo.opensearch_actions import ( + build_analytics_increment_action, + build_counter_increment_action, +) +from metrics.tests.services.legacy_matomo import METRICS, LegacyManifestTestCase + + +class LegacyMatomoImporterTests(LegacyManifestTestCase): + def test_counter_action_tracks_only_document_days(self): + document = { + "month": "2025-08", + "daily_metrics": {"01": METRICS}, + **METRICS, + } + + action = build_counter_increment_action( + "usage_monthly_scl_2025", + "key", + document, + "migration-with-hash", + ["2025-08-01", "2025-08-02"], + ) + + params = action["script"]["params"] + self.assertEqual(params["migration_id"], "migration-with-hash") + self.assertEqual(params["source_days"], ["2025-08-01"]) + + def test_analytics_action_builds_all_day_masks(self): + action = build_analytics_increment_action( + "usage_yearly_analytics_scl_2025", + "key", + METRICS, + "migration-with-hash", + ["2025-08-01", "2025-08-31"], + ) + + day_masks = action["script"]["params"]["day_masks"] + august_first_offset = 212 + august_last_offset = 242 + self.assertTrue( + day_masks[august_first_offset // 63] & (1 << (august_first_offset % 63)) + ) + self.assertTrue( + day_masks[august_last_offset // 63] & (1 << (august_last_offset % 63)) + ) + + @override_settings(OPENSEARCH_INDEX_NAME="usage") + def test_preflight_builds_yearly_targets_without_writes(self): + collection = SimpleNamespace( + acron3="scl", + log_manager_config=SimpleNamespace( + opensearch_partition_strategy="yearly", + opensearch_primary_shards=3, + ), + ) + search_client = Mock() + search_client.client.indices.exists_alias.return_value = False + + plan = importer.build_import_plan(search_client, collection, self.manifest) + + self.assertEqual(plan["primary_shards"], 3) + self.assertEqual( + [item["write_alias"]["alias"] for item in plan["datasets"]], + ["usage_monthly_scl_2025", "usage_yearly_analytics_scl_2025"], + ) + self.assertFalse(search_client.prepare_usage_index.called) + + @override_settings(OPENSEARCH_INDEX_NAME="usage") + def test_preflight_rejects_rollover_with_existing_yearly_alias(self): + collection = SimpleNamespace( + acron3="scl", + log_manager_config=SimpleNamespace( + opensearch_partition_strategy="rollover", + opensearch_primary_shards=1, + ), + ) + search_client = Mock() + search_client.client.indices.exists_alias.side_effect = ( + lambda name: name.endswith("_2025") + ) + search_client.client.indices.get_alias.side_effect = lambda name: { + name + + "-000001": { + "aliases": {name: {"is_write_index": True}}, + } + } + + with self.assertRaisesMessage(ValueError, "conflicts with yearly alias"): + importer.build_import_plan(search_client, collection, self.manifest) + + @patch("metrics.legacy_matomo.importer.helpers.bulk") + def test_historical_updates_follow_existing_backing_index(self, bulk): + search_client = Mock() + search_client.bulk_chunk_size = 500 + search_client.client.indices.exists_alias.return_value = True + search_client.client.indices.get_alias.return_value = { + "usage_monthly_scl_2025-000001": { + "aliases": { + "usage_monthly_scl_2025": {"is_write_index": False}, + } + }, + "usage_monthly_scl_2025-000002": { + "aliases": { + "usage_monthly_scl_2025": {"is_write_index": True}, + } + }, + } + search_client.client.search.return_value = { + "hits": { + "hits": [ + { + "_id": "existing", + "_index": "usage_monthly_scl_2025-000001", + } + ] + } + } + actions = [] + progress = [] + + def consume(_client, action_items, chunk_size): + actions.extend(action_items) + return len(actions), 0 + + bulk.side_effect = consume + + imported = importer._update_document_items( + search_client, + "usage_monthly_scl_2025", + iter([("existing", {}), ("new", {})]), + lambda index_name, doc_id, document: { + "_index": index_name, + "_id": doc_id, + "_source": document, + }, + dataset="counter", + progress_callback=lambda dataset, stage, documents: progress.append( + (dataset, stage, documents) + ), + progress_interval=1, + ) + + self.assertEqual(imported, 2) + self.assertEqual( + [action["_index"] for action in actions], + [ + "usage_monthly_scl_2025-000001", + "usage_monthly_scl_2025-000002", + ], + ) + self.assertEqual(progress, [("counter", "import", 2)]) diff --git a/metrics/tests/services/test_legacy_matomo_validation.py b/metrics/tests/services/test_legacy_matomo_validation.py new file mode 100644 index 0000000..1fdd81a --- /dev/null +++ b/metrics/tests/services/test_legacy_matomo_validation.py @@ -0,0 +1,121 @@ +import gzip +import hashlib +from pathlib import Path +from unittest.mock import Mock + +from metrics.legacy_matomo import validation +from metrics.legacy_matomo.manifest import build_migration_id, validate_migration_scope +from metrics.tests.services.legacy_matomo import METRICS, LegacyManifestTestCase + + +class LegacyMatomoValidationTests(LegacyManifestTestCase): + def test_validates_reconciled_manifest(self): + result = validation.validate_manifest(self.manifest) + + self.assertEqual(result["status"], "valid") + self.assertEqual(result["counter"]["totals"], METRICS) + self.assertEqual(result["analytics"]["totals"], METRICS) + + def test_scl_scope_accepts_nbr_alias_through_june_2026(self): + self.manifest["source_days"] = ["2026-06-30"] + + validate_migration_scope(self.manifest) + + def test_scl_scope_rejects_days_after_historical_cutoff(self): + self.manifest["source_days"] = ["2026-07-01"] + + with self.assertRaisesMessage(ValueError, "limited to 2026-06-30"): + validate_migration_scope(self.manifest) + + def test_other_article_collection_is_limited_to_2025(self): + self.manifest["source"]["collection"] = "spa" + self.manifest["target"]["collection"] = "spa" + self.manifest["source_days"] = ["2026-01-01"] + + with self.assertRaisesMessage(ValueError, "limited to 2025-12-31"): + validate_migration_scope(self.manifest) + + def test_rejects_non_article_collection(self): + self.manifest["source"]["collection"] = "books" + self.manifest["target"]["collection"] = "books" + + with self.assertRaisesMessage(ValueError, "not supported"): + validate_migration_scope(self.manifest) + + def test_rejects_undeclared_collection_alias(self): + self.manifest["source"]["collection"] = "old-spa" + self.manifest["target"]["collection"] = "spa" + + with self.assertRaisesMessage(ValueError, "old-spa -> spa"): + validate_migration_scope(self.manifest) + + def test_rejects_unknown_metric_profile(self): + self.manifest["metric_profile"] = "books" + + with self.assertRaisesMessage(ValueError, "metric profile"): + validate_migration_scope(self.manifest) + + def test_migration_identity_includes_payload_hash(self): + migration_id = build_migration_id(self.manifest, "counter") + + self.assertIn(self.manifest["counter"]["uncompressed_sha256"], migration_id) + + def test_rejects_manifest_without_source_days(self): + self.manifest["source_days"] = [] + + with self.assertRaisesMessage(ValueError, "no source days"): + validation.validate_manifest(self.manifest) + + def test_rejects_duplicate_payload_ids(self): + path = Path(self.manifest["counter"]["path"]) + with gzip.open(path, "rb") as source: + line = source.read() + with gzip.GzipFile(filename=path, mode="wb", mtime=0) as output: + output.write(line) + output.write(line) + self.manifest["counter"]["documents"] = 2 + self.manifest["counter"]["totals"] = { + key: value * 2 for key, value in METRICS.items() + } + self.manifest["counter"]["uncompressed_sha256"] = hashlib.sha256( + line + line + ).hexdigest() + + with self.assertRaisesMessage(ValueError, "duplicate _id"): + validation.validate_manifest(self.manifest) + + def test_validation_reports_progress_and_uses_requested_temporary_directory(self): + temporary_directory = Path(self.temporary_directory.name) / "work" + temporary_directory.mkdir() + progress = [] + + validation.validate_manifest( + self.manifest, + temporary_directory=temporary_directory, + progress_callback=lambda dataset, stage, documents: progress.append( + (dataset, stage, documents) + ), + progress_interval=1, + ) + + self.assertEqual( + progress, + [ + ("counter", "validation", 1), + ("analytics", "validation", 1), + ], + ) + self.assertEqual(list(temporary_directory.iterdir()), []) + + def test_validation_honors_stop_request(self): + stop_controller = Mock() + stop_controller.raise_if_requested.side_effect = RuntimeError( + "interruption requested" + ) + + with self.assertRaisesMessage(RuntimeError, "interruption requested"): + validation.validate_manifest( + self.manifest, + progress_interval=1, + stop_controller=stop_controller, + )