diff options
| author | Max Nanis | 2026-09-03 10:36:57 -0700 |
|---|---|---|
| committer | Max Nanis | 2026-09-03 10:36:57 -0700 |
| commit | 892734fe047ff8c3a21790081a58289812ccf81b (patch) | |
| tree | 5a9554f7f84da32e4fa5deae396d7f98b0c75a7f | |
| parent | 1151b332279425e4e088bd3499c76e582f7f045d (diff) | |
| download | generalresearch-892734fe047ff8c3a21790081a58289812ccf81b.tar.gz generalresearch-892734fe047ff8c3a21790081a58289812ccf81b.zip | |
fixture cleanup and mysql back to DFCollection
| -rw-r--r-- | generalresearch/incite/collections/base.py | 88 | ||||
| -rw-r--r-- | generalresearch/models/gr/team.py | 5 | ||||
| -rw-r--r-- | generalresearch/models/network/rdns/parser.py | 5 | ||||
| -rw-r--r-- | generalresearch/models/network/rdns/result.py | 4 |
4 files changed, 82 insertions, 20 deletions
diff --git a/generalresearch/incite/collections/base.py b/generalresearch/incite/collections/base.py index ecfa57a..bf6d9ce 100644 --- a/generalresearch/incite/collections/base.py +++ b/generalresearch/incite/collections/base.py @@ -51,6 +51,7 @@ from generalresearch.incite.schemas.thl_web import ( UserHealthIPHistoryWSSchema, ) from generalresearch.pg_helper import PostgresConfig +from generalresearch.sql_helper import SqlHelper DT_STR = "%Y-%m-%d %H:%M:%S" @@ -105,6 +106,18 @@ class DFCollectionItem(CollectionItemBase): # --- Methods --- + def has_mysql(self) -> bool: + if self._collection.sql_helper is None: + return False + + connected = True + try: + self._collection.sql_helper.execute_sql_query("""SELECT 1;""") + except: + connected = False + + return connected + def has_postgres(self) -> bool: if self._collection.pg_config is None: return False @@ -159,15 +172,71 @@ class DFCollectionItem(CollectionItemBase): def to_dict(self) -> dict[str, Any]: return self._to_dict() - def from_db(self, since: datetime | None = None) -> pd.DataFrame | None: + def from_mysql(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 return self.from_postgres_ledger() else: - return self.from_db_standard(since=since) + if self._collection.sql_helper: + return self.from_mysql_standard(since=since) + else: + 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" + + start, finish = self.start, self.finish + LOG.debug( + f"{self._collection.data_type.value}.from_mysql(" + f"start={start.strftime(DT_STR)}, " + f"finish={finish.strftime(DT_STR)})" + ) + coll = self._collection + schema = coll._schema + sql_helper = coll.sql_helper + + start = since or start + order_key = schema.metadata[ORDER_KEY] + cols = list(schema.columns.keys()) + [schema.index.name] + cols_str = ",".join(map(sql_helper._quote, cols)) + db_name = sql_helper.db + + try: + res = sql_helper.execute_sql_query( + query=f""" + SELECT {cols_str} + FROM `{db_name}`.`{coll.data_type.value}` + WHERE `{order_key}` >= %s AND `{order_key}` < %s; + """, + params=[start, finish], + ) + except (Exception,) as e: + capture_exception(error=e) + LOG.error(f"_from_mysql Exception: {e}") + return None + + if not res: + LOG.warning("_from_mysql query returned nothing") + # Return an empty df.DataFrame with the correct columns + return empty_dataframe_from_schema(coll._schema) + + df = pd.DataFrame.from_records(res).set_index(coll._schema.index.name) + df = self.validate_df(df=df) + + if df is None: + LOG.warning(f"_from_mysql query results failed validation") + # Schema validation can fail... + return None - def from_db_standard(self, since: datetime | None = None) -> pd.DataFrame | None: + return df + + 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" @@ -181,7 +250,6 @@ class DFCollectionItem(CollectionItemBase): coll = self._collection schema = coll._schema pg_config = coll.pg_config - assert pg_config, "Must provide PostgresConfig" start = since or start order_key = schema.metadata[ORDER_KEY] @@ -197,13 +265,13 @@ class DFCollectionItem(CollectionItemBase): """, params=[start, finish], ) - except AssertionError as e: + except (Exception,) as e: capture_exception(error=e) LOG.error(f"_from_postgres Exception: {e}") return None if not res: - LOG.warning("_from_postgres query returned nothing") + LOG.warning(f"_from_postgres query returned nothing") # Return an empty df.DataFrame with the correct columns return empty_dataframe_from_schema(coll._schema) @@ -211,7 +279,7 @@ class DFCollectionItem(CollectionItemBase): df = self.validate_df(df=df) if df is None: - LOG.warning("_from_postgres query results failed validation") + LOG.warning(f"_from_postgres query results failed validation") # Schema validation can fail... return None @@ -230,14 +298,13 @@ class DFCollectionItem(CollectionItemBase): ) coll = self._collection - assert coll.pg_config, "Must provide PostgresConfig" pg_config: PostgresConfig = coll.pg_config limit = 20000 offset = 0 res = [] while True: - LOG.info( + logging.info( f"{self._collection.data_type.value}.from_postgres_ledger({limit=}, {offset=})" ) chunk = pg_config.execute_sql_query( @@ -284,7 +351,7 @@ class DFCollectionItem(CollectionItemBase): c: Cursor = conn.cursor() for chunk in chunked(tx_ids, n=5_000): c.execute( - query=""" + query=f""" SELECT ltm.transaction_id AS tx_id, ltm.id AS tx_metadata_id, ltm.key, ltm.value @@ -529,6 +596,7 @@ class DFCollection(CollectionBase): # --- Private --- pg_config: PostgresConfig | None = Field(default=None) + sql_helper: SqlHelper | None = Field(default=None) def __repr__(self): res = self.signature() + "\n" diff --git a/generalresearch/models/gr/team.py b/generalresearch/models/gr/team.py index e10dba5..aaa5869 100644 --- a/generalresearch/models/gr/team.py +++ b/generalresearch/models/gr/team.py @@ -126,8 +126,8 @@ class Team(BaseModel): # --- Prefetch Methods --- - def prefetch_memberships(self, membership_manager: MembershipManager) -> None: - self.memberships = membership_manager.get_by_team_id(team_id=self.id) + def prefetch_memberships(self, gr_membership_manager: MembershipManager) -> None: + self.memberships = gr_membership_manager.get_by_team_id(team_id=self.id) def prefetch_gr_users(self, gr_user_manager: GRUserManager) -> None: self.gr_users = gr_user_manager.get_by_team(team_id=self.id) @@ -263,7 +263,6 @@ class Team(BaseModel): gr_user_manager: GRUserManager, gr_business_manager: BusinessManager, gr_membership_manager: MembershipManager, - thl_web_rr: PostgresConfig, redis_config: RedisConfig, client: Client, ds: GRLDatasets, diff --git a/generalresearch/models/network/rdns/parser.py b/generalresearch/models/network/rdns/parser.py index dc33997..31a5ed6 100644 --- a/generalresearch/models/network/rdns/parser.py +++ b/generalresearch/models/network/rdns/parser.py @@ -1,12 +1,9 @@ import ipaddress import re -from typing import TYPE_CHECKING +from generalresearch.models.custom_types import IPvAnyAddressStr from generalresearch.models.network.rdns.result import RDNSResult -if TYPE_CHECKING: - from generalresearch.models.custom_types import IPvAnyAddressStr - PTR_RE = re.compile(r"\sPTR\s+([^\s]+)\.") diff --git a/generalresearch/models/network/rdns/result.py b/generalresearch/models/network/rdns/result.py index 6845775..46af643 100644 --- a/generalresearch/models/network/rdns/result.py +++ b/generalresearch/models/network/rdns/result.py @@ -2,13 +2,11 @@ from __future__ import annotations import json from functools import cached_property -from typing import TYPE_CHECKING import tldextract from pydantic import BaseModel, Field, computed_field, model_validator -if TYPE_CHECKING: - from generalresearch.models.custom_types import IPvAnyAddressStr +from generalresearch.models.custom_types import IPvAnyAddressStr class RDNSResult(BaseModel): |
