diff options
| author | stuppie | 2026-09-04 10:59:38 -0600 |
|---|---|---|
| committer | stuppie | 2026-09-04 10:59:38 -0600 |
| commit | d65c35ba4090a14e32ec88fb503f80a2ba1338e3 (patch) | |
| tree | 0b96c7933db7f5fe8be6b5360b875677cb886191 | |
| parent | 1d2b9e48c2e0b6da9385364fe57d51cd30bd64a5 (diff) | |
| download | generalresearch-d65c35ba4090a14e32ec88fb503f80a2ba1338e3.tar.gz generalresearch-d65c35ba4090a14e32ec88fb503f80a2ba1338e3.zip | |
from_mysql -> from_db. FilePath -> Path. Fix check_model_after on MergeCollectionItem
| -rw-r--r-- | generalresearch/incite/base.py | 17 | ||||
| -rw-r--r-- | generalresearch/incite/collections/base.py | 68 | ||||
| -rw-r--r-- | generalresearch/incite/mergers/base.py | 20 |
3 files changed, 63 insertions, 42 deletions
diff --git a/generalresearch/incite/base.py b/generalresearch/incite/base.py index a06aac9..78247eb 100644 --- a/generalresearch/incite/base.py +++ b/generalresearch/incite/base.py @@ -67,7 +67,6 @@ Items = Sequence[Item] DT_STR = "%Y-%m-%d %H:%M:%S" _dir_adapter = TypeAdapter(DirectoryPath) -_filepath_adapter = TypeAdapter(FilePath) class NFSMount(BaseModel): @@ -704,20 +703,20 @@ class CollectionItemBase(BaseModel): return f"{self.filename}.empty" @property - def path(self) -> FilePath: - return _filepath_adapter.validate_python( + def path(self) -> Path: + return Path( os.path.join(self._collection.archive_path, self.filename) ) @property - def partial_path(self) -> FilePath: - return FilePath( + def partial_path(self) -> Path: + return Path( os.path.join(self._collection.archive_path, self.partial_filename) ) @property - def empty_path(self) -> FilePath: - return FilePath( + def empty_path(self) -> Path: + return Path( os.path.join(self._collection.archive_path, self.empty_filename) ) @@ -782,8 +781,8 @@ class CollectionItemBase(BaseModel): # up as always returning the same tmp filename return f"{self.filename}.{uuid4().hex}" - def tmp_path(self) -> FilePath: - return FilePath( + def tmp_path(self) -> Path: + return Path( os.path.join(self._collection.archive_path, self.tmp_filename()) ) diff --git a/generalresearch/incite/collections/base.py b/generalresearch/incite/collections/base.py index 47bb70a..04dd97e 100644 --- a/generalresearch/incite/collections/base.py +++ b/generalresearch/incite/collections/base.py @@ -1,5 +1,6 @@ from __future__ import annotations +import logging import os import subprocess import time @@ -94,9 +95,18 @@ DFCollectionTypeSchemas = { DFCollectionType.SPECTRUM_SURVEY_TIMESERIES: SpectrumSurveyTimeseriesSchema, } +# 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 +MYSQL_ALLOWED_COLL_TYPES = { + DFCollectionType.INNOVATE_SURVEY_HISTORY, + DFCollectionType.MORNING_SURVEY_TIMESERIES, + DFCollectionType.SAGO_SURVEY_HISTORY, + DFCollectionType.SPECTRUM_SURVEY_TIMESERIES, +} -class DFCollectionItem(CollectionItemBase): +class DFCollectionItem(CollectionItemBase): # --- Properties --- @property def filename(self) -> str: @@ -126,7 +136,7 @@ class DFCollectionItem(CollectionItemBase): connected = True try: self._collection.pg_config.execute_sql_query("""SELECT 1;""") - except AssertionError: + except Exception: connected = False return connected @@ -173,7 +183,7 @@ class DFCollectionItem(CollectionItemBase): def to_dict(self) -> dict[str, Any]: return self._to_dict() - def from_mysql(self, since: datetime | None = None) -> pd.DataFrame | None: + 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" assert self._collection.pg_config is not None @@ -185,14 +195,14 @@ class DFCollectionItem(CollectionItemBase): return self.from_postgres_standard(since=since) def from_mysql_standard(self, since: datetime | None = None) -> pd.DataFrame | None: - - assert ( - self._collection.data_type != DFCollectionType.LEDGER - ), "Can't call from_mysql_standard for Ledger DFCollectionItem" + data_type = self._collection.data_type + assert data_type in MYSQL_ALLOWED_COLL_TYPES, ( + f"Unsupported {data_type=} for mysql" + ) start, finish = self.start, self.finish LOG.debug( - f"{self._collection.data_type.value}.from_mysql(" + f"{data_type.value}.from_mysql(" f"start={start.strftime(DT_STR)}, " f"finish={finish.strftime(DT_STR)})" ) @@ -215,7 +225,7 @@ class DFCollectionItem(CollectionItemBase): """, params=[start, finish], ) - except (Exception,) as e: + except Exception as e: capture_exception(error=e) LOG.error(f"_from_mysql Exception: {e}") return None @@ -229,7 +239,7 @@ class DFCollectionItem(CollectionItemBase): df = self.validate_df(df=df) if df is None: - LOG.warning(f"_from_mysql query results failed validation") + LOG.warning("_from_mysql query results failed validation") # Schema validation can fail... return None @@ -238,13 +248,14 @@ class DFCollectionItem(CollectionItemBase): def from_postgres_standard( self, since: datetime | None = None ) -> pd.DataFrame | None: - assert ( - self._collection.data_type != DFCollectionType.LEDGER - ), "Can't call from_postgres_standard for Ledger DFCollectionItem" + data_type = self._collection.data_type + assert data_type != DFCollectionType.LEDGER, ( + "Can't call from_postgres_standard for Ledger DFCollectionItem" + ) start, finish = self.start, self.finish LOG.debug( - f"{self._collection.data_type.value}.from_postgres(" + f"{data_type.value}.from_postgres(" f"start={start.strftime(DT_STR)}, " f"finish={finish.strftime(DT_STR)})" ) @@ -261,18 +272,18 @@ class DFCollectionItem(CollectionItemBase): res = pg_config.execute_sql_query( query=f""" SELECT {cols_str} - FROM {coll.data_type.value} + FROM {data_type.value} WHERE {order_key} >= %s AND {order_key} < %s; """, params=[start, finish], ) - except (Exception,) as e: + except Exception as e: capture_exception(error=e) LOG.error(f"_from_postgres Exception: {e}") return None if not res: - LOG.warning(f"_from_postgres query returned nothing") + LOG.warning("_from_postgres query returned nothing") # Return an empty df.DataFrame with the correct columns return empty_dataframe_from_schema(coll._schema) @@ -280,20 +291,21 @@ class DFCollectionItem(CollectionItemBase): df = self.validate_df(df=df) if df is None: - LOG.warning(f"_from_postgres query results failed validation") + LOG.warning("_from_postgres query results failed validation") # Schema validation can fail... return None return df def from_postgres_ledger(self) -> pd.DataFrame | None: - assert ( - self._collection.data_type == DFCollectionType.LEDGER - ), "Can only call from_postgres_ledger on Ledger DFCollectionItem" + data_type = self._collection.data_type + assert data_type == DFCollectionType.LEDGER, ( + "Can only call from_postgres_ledger on Ledger DFCollectionItem" + ) start, finish = self.start, self.finish LOG.info( - f"{self._collection.data_type.value}.from_postgres_ledger(" + f"{data_type.value}.from_postgres_ledger(" f"start={start.strftime(DT_STR)}, " f"finish={finish.strftime(DT_STR)})" ) @@ -305,9 +317,7 @@ class DFCollectionItem(CollectionItemBase): offset = 0 res = [] while True: - logging.info( - f"{self._collection.data_type.value}.from_postgres_ledger({limit=}, {offset=})" - ) + logging.info(f"{data_type.value}.from_postgres_ledger({limit=}, {offset=})") chunk = pg_config.execute_sql_query( query=f""" SELECT lt.id AS tx_id, lt.created, lt.ext_description, lt.tag, @@ -352,7 +362,7 @@ class DFCollectionItem(CollectionItemBase): c: Cursor = conn.cursor() for chunk in chunked(tx_ids, n=5_000): c.execute( - query=f""" + query=""" SELECT ltm.transaction_id AS tx_id, ltm.id AS tx_metadata_id, ltm.key, ltm.value @@ -528,9 +538,9 @@ class DFCollectionItem(CollectionItemBase): # Make sure these are in the same dir. b/c the symlink has to be # relative, not an absolute path - assert ( - partial_path.parent == next_numbered_path.parent - ), "Can't have numbered_path in a different directory" + assert partial_path.parent == next_numbered_path.parent, ( + "Can't have numbered_path in a different directory" + ) target = ( next_numbered_path.name ) # this is the symlink's target. it is a relative path (only the name) diff --git a/generalresearch/incite/mergers/base.py b/generalresearch/incite/mergers/base.py index 69c6c46..a0874a7 100644 --- a/generalresearch/incite/mergers/base.py +++ b/generalresearch/incite/mergers/base.py @@ -65,7 +65,6 @@ MergeTypeSchemas = { class MergeCollectionItem(CollectionItemBase): - # --- Properties --- @property @@ -141,9 +140,9 @@ class MergeCollectionItem(CollectionItemBase): compression="brotli", ) client.compute(f, sync=True, priority=2, resources=client_resources) - assert not os.path.exists( - self.path.as_posix() - ), f"already exits!: {self.path.as_posix()}" + assert not os.path.exists(self.path.as_posix()), ( + f"already exits!: {self.path.as_posix()}" + ) if platform == "darwin": subprocess.call(["mv", tmp_path.as_posix(), self.path.as_posix()]) @@ -246,6 +245,19 @@ class MergeCollection(CollectionBase): collection_item_class: type[MergeCollectionItem] = MergeCollectionItem @model_validator(mode="after") + def check_model_after(self) -> Self: + if self.offset is None or self.start is None: + return self + + offset_total_sec = pd.Timedelta(self.offset).total_seconds() + start_total_sec = (datetime.now(tz=UTC) - self.start).total_seconds() + + if offset_total_sec > start_total_sec: + raise ValueError("Offset must be equal to, or smaller the start timestamp") + + return self + + @model_validator(mode="after") def check_start_and_offset_nullable(self) -> Self: if self.offset is None and self.start is None: raise AssertionError("cannot set both start and offset to None") |
