aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorstuppie2026-09-04 10:59:38 -0600
committerstuppie2026-09-04 10:59:38 -0600
commitd65c35ba4090a14e32ec88fb503f80a2ba1338e3 (patch)
tree0b96c7933db7f5fe8be6b5360b875677cb886191
parent1d2b9e48c2e0b6da9385364fe57d51cd30bd64a5 (diff)
downloadgeneralresearch-d65c35ba4090a14e32ec88fb503f80a2ba1338e3.tar.gz
generalresearch-d65c35ba4090a14e32ec88fb503f80a2ba1338e3.zip
from_mysql -> from_db. FilePath -> Path. Fix check_model_after on MergeCollectionItem
-rw-r--r--generalresearch/incite/base.py17
-rw-r--r--generalresearch/incite/collections/base.py68
-rw-r--r--generalresearch/incite/mergers/base.py20
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")