diff options
| -rw-r--r-- | generalresearch/managers/thl/product.py | 4 | ||||
| -rw-r--r-- | generalresearch/models/thl/product.py | 84 |
2 files changed, 83 insertions, 5 deletions
diff --git a/generalresearch/managers/thl/product.py b/generalresearch/managers/thl/product.py index abf6b4d..b730f92 100644 --- a/generalresearch/managers/thl/product.py +++ b/generalresearch/managers/thl/product.py @@ -202,7 +202,7 @@ class ProductManager(PostgresManager): order_field: str = "created", descending: bool = False, conn: Connection | None = None, - ) -> tuple[list[Product], int]: + ) -> tuple[list[Product], NonNegativeInt]: products = self.filter_by( product_uuids=product_uuids, business_uuids=business_uuids, @@ -723,7 +723,7 @@ class ProductManager(PostgresManager): """ updates_by_field = {} for product in products: - data = product.model_dump(mode="json", include=self.CACHED_FIELDS) + data = product.model_dump(mode="json", include=set(self.CACHED_FIELDS)) for k, v in data.items(): if k in self.CACHED_FIELDS_JSON and v is not None: v = json.dumps(v) diff --git a/generalresearch/models/thl/product.py b/generalresearch/models/thl/product.py index 816339f..9d826fa 100644 --- a/generalresearch/models/thl/product.py +++ b/generalresearch/models/thl/product.py @@ -72,6 +72,9 @@ if TYPE_CHECKING: from dask.distributed import Client from generalresearch.incite.base import GRLDatasets + from generalresearch.incite.mergers.foundations.enriched_session import ( + EnrichedSessionMerge, + ) from generalresearch.incite.mergers.pop_ledger import PopLedgerMerge from generalresearch.managers.thl.ledger_manager.thl_ledger import ( ThlLedgerManager, @@ -1091,6 +1094,83 @@ class Product(BaseModel, validate_assignment=True): # --- Prebuild --- @staticmethod + def get_enriched_session_metrics_df( + product_ids: Collection[UUIDStr], + client: Client, + enriched_session: EnrichedSessionMerge, + ) -> pd.DataFrame: + """Load all EnrichedSession rows needed to cache a batch of Products.""" + now = pd.Timestamp.now(tz="UTC") + cutoff = now - pd.Timedelta(days=7) + + ddf = enriched_session.ddf( + include_partial=True, + force_rr_latest=False, + columns=[ + "product_id", + "user_id", + "started", + "status", + ], + filters=[ + ("started", ">=", cutoff.to_pydatetime()), + ("started", "<", now.to_pydatetime()), + ("product_id", "in", list(product_ids)), + ], + ) + if ddf is None: + return pd.DataFrame( + columns=[ + "product_id", + "users_active_7d", + "task_completes_7d", + ] + ) + ddf = ddf.assign(is_complete=ddf["status"].eq("c")) + + users = ( + ddf[["product_id", "user_id"]] + .drop_duplicates() + .groupby("product_id") + .size() + .rename("users_active_7d") + ) + completes = ( + ddf.loc[ddf["is_complete"], ["product_id"]] + .groupby("product_id") + .size() + .rename("task_completes_7d") + ) + users_series, completes_series = client.compute( + [users, completes], + sync=True, + ) + metrics_df = ( + users_series.to_frame() + .join( + completes_series.to_frame(), + how="outer", + ) + .fillna(0) + .astype( + { + "users_active_7d": int, + "task_completes_7d": int, + } + ) + ) + return metrics_df + + def prebuild_metrics(self, metrics_df: pd.DataFrame) -> None: + if self.id in metrics_df.index: + row = metrics_df.loc[self.id] + self.users_active_7d = int(row["users_active_7d"]) + self.task_completes_7d = int(row["task_completes_7d"]) + else: + self.users_active_7d = 0 + self.task_completes_7d = 0 + + @staticmethod def get_pop_ledger_df( product_ids: Collection[UUIDStr], client: Client, @@ -1225,9 +1305,7 @@ class Product(BaseModel, validate_assignment=True): cutoff = pd.Timestamp.now(tz="UTC") - timedelta(days=7) balance_7d_df = balance_df.loc[balance_df.index >= cutoff] self.balance_net_7d = ( - 0 - if balance_7d_df.empty - else ProductBalances.from_pandas(balance_7d_df).net + 0 if balance_7d_df.empty else ProductBalances.from_pandas(balance_7d_df).net ) balance = ProductBalances.from_pandas(balance_df) |
