aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorstuppie2026-10-08 16:10:12 -0600
committerstuppie2026-10-08 16:10:12 -0600
commitd95960396373d7c200e96f2e0743ccd3693e79d1 (patch)
treecd3efcdc7e81fcd4d639663bbc8b0190e2701a2e
parent8d6ad7203f6f4595d00edabfe229e5ffcc076916 (diff)
downloadgeneralresearch-d95960396373d7c200e96f2e0743ccd3693e79d1.tar.gz
generalresearch-d95960396373d7c200e96f2e0743ccd3693e79d1.zip
incite add validate_time_coverageHEADv3.7.5master
-rw-r--r--generalresearch/incite/collections/base.py40
-rw-r--r--generalresearch/incite/mergers/base.py46
-rw-r--r--pyproject.toml2
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"