diff options
| author | stuppie | 2026-10-08 16:10:12 -0600 |
|---|---|---|
| committer | stuppie | 2026-10-08 16:10:12 -0600 |
| commit | d95960396373d7c200e96f2e0743ccd3693e79d1 (patch) | |
| tree | cd3efcdc7e81fcd4d639663bbc8b0190e2701a2e | |
| parent | 8d6ad7203f6f4595d00edabfe229e5ffcc076916 (diff) | |
| download | generalresearch-d95960396373d7c200e96f2e0743ccd3693e79d1.tar.gz generalresearch-d95960396373d7c200e96f2e0743ccd3693e79d1.zip | |
| -rw-r--r-- | generalresearch/incite/collections/base.py | 40 | ||||
| -rw-r--r-- | generalresearch/incite/mergers/base.py | 46 | ||||
| -rw-r--r-- | pyproject.toml | 2 |
3 files changed, 85 insertions, 3 deletions
diff --git a/generalresearch/incite/collections/base.py b/generalresearch/incite/collections/base.py index f8159e8..b83fd8a 100644 --- a/generalresearch/incite/collections/base.py +++ b/generalresearch/incite/collections/base.py @@ -4,8 +4,9 @@ import os import subprocess import time import warnings -from datetime import datetime +from datetime import datetime, timedelta from enum import StrEnum +from pathlib import Path from sys import platform from typing import Any @@ -93,6 +94,11 @@ DFCollectionTypeSchemas = { DFCollectionType.SPECTRUM_SURVEY_TIMESERIES: SpectrumSurveyTimeseriesSchema, } +DFCollectionTypeTimestampColumns: dict[DFCollectionType, str] = { + data_type: schema.metadata[ORDER_KEY] + for data_type, schema in DFCollectionTypeSchemas.items() +} + # This is not a technical limitation, it is just a double check # since this should never happen since thl is now and forever # forward on postgres @@ -181,6 +187,38 @@ class DFCollectionItem(CollectionItemBase): def to_dict(self) -> dict[str, Any]: return self._to_dict() + def validate_time_coverage(self, path: Path) -> None: + timestamp_column = DFCollectionTypeTimestampColumns[self._collection.data_type] + df = pd.read_parquet(path, columns=[timestamp_column]) + timestamps = pd.to_datetime(df[timestamp_column]) + coverage = timestamps.max() - timestamps.min() + minimum_coverage = self.interval.length - timedelta(days=2) + + if pd.isna(coverage) or coverage < minimum_coverage: + raise ValueError( + f"{self.name} archive does not cover its interval: " + f"{coverage=} < {minimum_coverage=}" + ) + + def valid_archive( + self, + generic_path: FilePath | None = None, + sample: int | None = None, + ) -> bool: + if not super().valid_archive(generic_path=generic_path, sample=sample): + return False + + path = Path(generic_path or self.path) + if path == self.partial_path or ".partial." in path.name: + return True + + try: + self.validate_time_coverage(path) + except ValueError as e: + LOG.warning(f"Invalid archive time coverage {path=} {e=}") + return False + return True + def from_db(self, since: datetime | None = None) -> pd.DataFrame | None: if self._collection.data_type == DFCollectionType.LEDGER: assert since is None, "Shouldn't pass since for Ledger item" diff --git a/generalresearch/incite/mergers/base.py b/generalresearch/incite/mergers/base.py index a0874a7..6bb41e9 100644 --- a/generalresearch/incite/mergers/base.py +++ b/generalresearch/incite/mergers/base.py @@ -1,7 +1,8 @@ import os.path import subprocess -from datetime import UTC, datetime +from datetime import UTC, datetime, timedelta from enum import StrEnum +from pathlib import Path from sys import platform from typing import Self @@ -63,6 +64,16 @@ MergeTypeSchemas = { MergeType.ENRICHED_TASK_ADJUST: EnrichedTaskAdjustSchema, } +MergeTypeTimestampColumns: dict[MergeType, str | None] = { + MergeType.YM_SURVEY_WALL: "started", + MergeType.YM_WALL_SUMMARY: "date", + MergeType.POP_LEDGER: "time_idx", + MergeType.USER_ID_PRODUCT: None, + MergeType.ENRICHED_WALL: "started", + MergeType.ENRICHED_SESSION: "started", + MergeType.ENRICHED_TASK_ADJUST: "alerted", +} + class MergeCollectionItem(CollectionItemBase): # --- Properties --- @@ -100,6 +111,34 @@ class MergeCollectionItem(CollectionItemBase): res["group_by"] = self._collection.group_by return res + def validate_time_coverage(self, path: Path) -> None: + timestamp_column = MergeTypeTimestampColumns[self._collection.merge_type] + if timestamp_column is None: + return + + df = pd.read_parquet(path, columns=[timestamp_column]) + timestamps = pd.to_datetime(df[timestamp_column]) + coverage = timestamps.max() - timestamps.min() + minimum_coverage = self.interval.length - timedelta(days=2) + + if pd.isna(coverage) or coverage < minimum_coverage: + raise ValueError( + f"{self.name} archive does not cover its interval: " + f"{coverage=} < {minimum_coverage=}" + ) + + def valid_archive(self, generic_path=None, sample: int | None = None) -> bool: + if not super().valid_archive(generic_path=generic_path, sample=sample): + return False + + path = generic_path or self.path + try: + self.validate_time_coverage(path) + except ValueError as e: + LOG.warning(f"Invalid archive time coverage {path=} {e=}") + return False + return True + def to_archive( self, client: Client, @@ -140,6 +179,11 @@ class MergeCollectionItem(CollectionItemBase): compression="brotli", ) client.compute(f, sync=True, priority=2, resources=client_resources) + try: + self.validate_time_coverage(tmp_path) + except Exception: + self.delete_archive(tmp_path) + raise assert not os.path.exists(self.path.as_posix()), ( f"already exits!: {self.path.as_posix()}" ) diff --git a/pyproject.toml b/pyproject.toml index 08ed683..d1b13e7 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "generalresearch" -version = "3.7.4" +version = "3.7.5" description = "Python Utilities for General Research" readme = "README.md" requires-python = ">=3.14" |
