aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--generalresearch/managers/thl/product.py4
-rw-r--r--generalresearch/models/thl/product.py84
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)