aboutsummaryrefslogtreecommitdiff
path: root/tests/incite/collections/test_df_collection_item_thl_web.py
diff options
context:
space:
mode:
Diffstat (limited to 'tests/incite/collections/test_df_collection_item_thl_web.py')
-rw-r--r--tests/incite/collections/test_df_collection_item_thl_web.py454
1 files changed, 218 insertions, 236 deletions
diff --git a/tests/incite/collections/test_df_collection_item_thl_web.py b/tests/incite/collections/test_df_collection_item_thl_web.py
index 8b8bcbe..eeabb41 100644
--- a/tests/incite/collections/test_df_collection_item_thl_web.py
+++ b/tests/incite/collections/test_df_collection_item_thl_web.py
@@ -1,17 +1,19 @@
from __future__ import annotations
-from collections.abc import Generator
-from datetime import datetime, timedelta, timezone
+from collections.abc import Callable, Generator
+from datetime import UTC, datetime, timedelta
from itertools import product as iter_product
from os.path import join as pjoin
from pathlib import Path, PurePath
-from typing import TYPE_CHECKING, Callable
+from typing import TYPE_CHECKING
from uuid import uuid4
import dask.dataframe as dd
import pandas as pd
import pytest
-from distributed import Client, Scheduler, Worker
+from dask.distributed import Client as DaskClient
+from dask.distributed import Scheduler as DaskScheduler
+from dask.distributed import Worker as DaskWorker
# noinspection PyUnresolvedReferences
from distributed.utils_test import (
@@ -22,18 +24,21 @@ from pandera.pandas import DataFrameSchema
from pydantic import FilePath
from generalresearch.incite.base import CollectionItemBase
-from generalresearch.incite.collections import (
- DFCollectionItem,
+from generalresearch.incite.collections.base import (
DFCollectionType,
)
from generalresearch.incite.schemas import ARCHIVE_AFTER
-from generalresearch.models.thl.product import Product
-from generalresearch.models.thl.user import User
from generalresearch.pg_helper import PostgresConfig
from generalresearch.sql_helper import PostgresDsn
if TYPE_CHECKING:
from generalresearch.incite.base import GRLDatasets
+ from generalresearch.incite.collections.base import (
+ DFCollection,
+ DFCollectionItem,
+ )
+ from generalresearch.models.thl.product import Product
+ from generalresearch.models.thl.user import User
fake = Faker()
@@ -52,12 +57,11 @@ unsupported_mock_types = {
}
-def combo_object() -> Generator[str, None, None]:
- for x in iter_product(
+def combo_object() -> Generator[tuple[DFCollectionType, str]]:
+ yield from iter_product(
df_collections,
- ["15min", "45min", "1H"],
- ):
- yield from x
+ ["15min", "45min", "1h"],
+ )
class TestDFCollectionItemBase:
@@ -71,8 +75,12 @@ class TestDFCollectionItemBase:
argnames="df_collection_data_type, offset", argvalues=combo_object()
)
class TestDFCollectionItemProperties:
-
- def test_filename(self, df_collection_data_type, df_collection, offset: str):
+ def test_filename(
+ self,
+ df_collection_data_type: DFCollectionType,
+ df_collection: DFCollection,
+ offset: str,
+ ):
for i in df_collection.items:
assert isinstance(i.filename, str)
@@ -88,38 +96,59 @@ class TestDFCollectionItemProperties:
argnames="df_collection_data_type, offset", argvalues=combo_object()
)
class TestDFCollectionItemPropertiesBase:
-
- def test_name(self, df_collection_data_type, offset: str, df_collection):
+ def test_name(
+ self,
+ df_collection: DFCollection,
+ ):
for i in df_collection.items:
assert isinstance(i.name, str)
- def test_finish(self, df_collection_data_type, offset: str, df_collection):
+ def test_finish(
+ self,
+ df_collection: DFCollection,
+ ):
for i in df_collection.items:
assert isinstance(i.finish, datetime)
- def test_interval(self, df_collection_data_type, offset: str, df_collection):
+ def test_interval(
+ self,
+ df_collection: DFCollection,
+ ):
for i in df_collection.items:
assert isinstance(i.interval, pd.Interval)
def test_partial_filename(
- self, df_collection_data_type, offset: str, df_collection
+ self,
+ df_collection: DFCollection,
):
for i in df_collection.items:
assert isinstance(i.partial_filename, str)
- def test_empty_filename(self, df_collection_data_type, offset: str, df_collection):
+ def test_empty_filename(
+ self,
+ df_collection: DFCollection,
+ ):
for i in df_collection.items:
assert isinstance(i.empty_filename, str)
- def test_path(self, df_collection_data_type, offset: str, df_collection):
+ def test_path(
+ self,
+ df_collection: DFCollection,
+ ):
for i in df_collection.items:
assert isinstance(i.path, FilePath)
- def test_partial_path(self, df_collection_data_type, offset: str, df_collection):
+ def test_partial_path(
+ self,
+ df_collection: DFCollection,
+ ):
for i in df_collection.items:
assert isinstance(i.partial_path, FilePath)
- def test_empty_path(self, df_collection_data_type, offset: str, df_collection):
+ def test_empty_path(
+ self,
+ df_collection: DFCollection,
+ ):
for i in df_collection.items:
assert isinstance(i.empty_path, FilePath)
@@ -135,26 +164,25 @@ class TestDFCollectionItemPropertiesBase:
),
)
class TestDFCollectionItemMethod:
-
- def test_has_mysql(
+ def test_has_postgres(
self,
- df_collection,
- thl_web_rr: PostgresConfig,
+ df_collection_data_type: DFCollectionType,
offset: str,
duration: timedelta,
- df_collection_data_type,
- delete_df_collection,
+ delete_df_collection: Callable[..., None],
+ df_collection: DFCollection,
+ thl_web_rr: PostgresConfig,
):
delete_df_collection(coll=df_collection)
df_collection.pg_config = None
for i in df_collection.items:
- assert not i.has_mysql()
+ assert not i.has_postgres()
# Confirm that the regular connection should work as expected
df_collection.pg_config = thl_web_rr
for i in df_collection.items:
- assert i.has_mysql()
+ assert i.has_postgres()
# Make a fake connection and confirm it does NOT work
df_collection.pg_config = PostgresConfig(
@@ -163,17 +191,14 @@ class TestDFCollectionItemMethod:
statement_timeout=1,
)
for i in df_collection.items:
- assert not i.has_mysql()
+ assert not i.has_postgres()
@pytest.mark.skip
def test_update_partial_archive(
self,
- df_collection,
+ df_collection_data_type: DFCollectionType,
offset: str,
duration: timedelta,
- thl_web_rw: PostgresConfig,
- df_collection_data_type,
- delete_df_collection,
):
# for i in collection.items:
# assert i.update_partial_archive()
@@ -183,29 +208,16 @@ class TestDFCollectionItemMethod:
@pytest.mark.skip
def test_create_partial_archive(
self,
- df_collection,
+ df_collection_data_type: DFCollectionType,
offset: str,
- duration: str,
- create_main_accounts,
- thl_web_rw: PostgresConfig,
- thl_lm,
- df_collection_data_type,
- user_factory: Callable[..., User],
- product: Product,
- client_no_amm,
- incite_item_factory,
- delete_df_collection,
- mnt_filepath: GRLDatasets,
+ duration: timedelta,
):
- assert 1 + 1 == 2
+ pass
def test_dict(
self,
- df_collection_data_type,
- offset: str,
- duration: timedelta,
- df_collection,
- delete_df_collection,
+ df_collection: DFCollection,
+ delete_df_collection: Callable[..., None],
):
delete_df_collection(coll=df_collection)
@@ -225,18 +237,17 @@ class TestDFCollectionItemMethod:
def test_from_mysql(
self,
- df_collection_data_type,
- df_collection,
+ df_collection_data_type: DFCollectionType,
+ df_collection: DFCollection,
offset: str,
duration: timedelta,
- create_main_accounts,
+ create_main_accounts: Callable[..., None],
thl_web_rw: PostgresConfig,
user_factory: Callable[..., User],
product: Product,
- incite_item_factory,
- delete_df_collection,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
):
- from generalresearch.models.thl.user import User
if df_collection.data_type in unsupported_mock_types:
return
@@ -249,38 +260,32 @@ class TestDFCollectionItemMethod:
for item in df_collection.items:
# Unlike .from_mysql_ledger(), .from_mysql_standard() will return
# back and empty df with the correct columns in place
- delete_df_collection(coll=df_collection)
- df = item.from_mysql()
if df_collection.data_type == DFCollectionType.LEDGER:
- assert df is None
- else:
- assert df.empty
- assert set(df.columns) == set(df_collection._schema.columns.keys())
+ continue
+ delete_df_collection(coll=df_collection)
+ df = item.from_postgres_standard()
+ assert isinstance(df, pd.DataFrame)
+ assert df.empty
+ assert set(df.columns) == set(df_collection.type_schema.columns.keys())
incite_item_factory(user=u1, item=item)
- df = item.from_mysql()
+ df = item.from_postgres_standard()
+ assert isinstance(df, pd.DataFrame)
assert not df.empty
- assert set(df.columns) == set(df_collection._schema.columns.keys())
- if df_collection.data_type == DFCollectionType.LEDGER:
- # The number of rows in this dataframe will change depending
- # on the mocking of data. It's because if the account has
- # user wallet on, then there will be more transactions for
- # example.
- assert df.shape[0] > 0
+ assert set(df.columns) == set(df_collection.type_schema.columns.keys())
- def test_from_mysql_standard(
+ def test_from_postgres_standard(
self,
- df_collection_data_type,
- df_collection,
+ df_collection_data_type: DFCollectionType,
+ df_collection: DFCollection,
offset: str,
duration: timedelta,
user_factory: Callable[..., User],
product: Product,
- incite_item_factory,
- delete_df_collection,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
):
- from generalresearch.models.thl.user import User
if df_collection.data_type in unsupported_mock_types:
return
@@ -292,48 +297,31 @@ class TestDFCollectionItemMethod:
item: DFCollectionItem
if df_collection.data_type == DFCollectionType.LEDGER:
- # We're using parametrize, so this If statement is just to
- # confirm other Item Types will always raise an assertion
- with pytest.raises(expected_exception=AssertionError) as cm:
- res = item.from_mysql_standard()
- assert (
- "Can't call from_mysql_standard for Ledger DFCollectionItem"
- in str(cm.value)
- )
-
continue
# Unlike .from_mysql_ledger(), .from_mysql_standard() will return
# back and empty df with the correct columns in place
- df = item.from_mysql_standard()
+ df = item.from_postgres_standard()
+ assert isinstance(df, pd.DataFrame)
assert df.empty
- assert set(df.columns) == set(df_collection._schema.columns.keys())
+ assert set(df.columns) == set(df_collection.type_schema.columns.keys())
incite_item_factory(user=u1, item=item)
- df = item.from_mysql_standard()
+ df = item.from_postgres_standard()
+ assert isinstance(df, pd.DataFrame)
assert not df.empty
- assert set(df.columns) == set(df_collection._schema.columns.keys())
+ assert set(df.columns) == set(df_collection.type_schema.columns.keys())
assert df.shape[0] > 0
- def test_from_mysql_ledger(
+ def test_from_postgres_ledger(
self,
- df_collection,
- user: User,
- create_main_accounts,
- offset: str,
- duration: timedelta,
- thl_web_rw: PostgresConfig,
- thl_lm,
- df_collection_data_type,
+ df_collection: DFCollection,
user_factory: Callable[..., User],
product: Product,
- client_no_amm,
- incite_item_factory,
- delete_df_collection,
- mnt_filepath,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
):
- from generalresearch.models.thl.user import User
if df_collection.data_type != DFCollectionType.LEDGER:
return
@@ -348,14 +336,14 @@ class TestDFCollectionItemMethod:
# Okay, now continue with the actual Ledger Item tests... we need
# to ensure that this item.start - item.finish range hasn't had
# any prior transactions created within that range.
- assert item.from_mysql_ledger() is None
+ assert item.from_postgres_ledger() is None
# Create main accounts doesn't matter because it doesn't
# add any transactions to the db
- assert item.from_mysql_ledger() is None
+ assert item.from_postgres_ledger() is None
incite_item_factory(user=u1, item=item)
- df = item.from_mysql_ledger()
+ df = item.from_postgres_ledger()
assert isinstance(df, pd.DataFrame)
# Not only is this a np.int64 to int comparison, but I also know it
@@ -373,19 +361,12 @@ class TestDFCollectionItemMethod:
def test_to_archive(
self,
- df_collection,
- user: User,
- offset: str,
- duration: timedelta,
- df_collection_data_type,
+ df_collection: DFCollection,
user_factory: Callable[..., User],
product: Product,
- client_no_amm,
- incite_item_factory,
- delete_df_collection,
- mnt_filepath,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
):
- from generalresearch.models.thl.user import User
if df_collection.data_type in unsupported_mock_types:
return
@@ -400,7 +381,7 @@ class TestDFCollectionItemMethod:
# Load up the data that we'll be using for various to_archive
# methods.
- df = item.from_mysql()
+ df = item.from_postgres_standard()
ddf = dd.from_pandas(df, npartitions=1)
# (1) Write the basic archive, the issue is that because it's
@@ -411,17 +392,12 @@ class TestDFCollectionItemMethod:
def test__to_archive(
self,
- df_collection_data_type,
- df_collection,
+ df_collection: DFCollection,
user_factory: Callable[..., User],
product: Product,
- offset: str,
- duration: timedelta,
- client_no_amm,
- user: User,
- incite_item_factory,
- delete_df_collection,
- mnt_filepath,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
+ mnt_filepath: GRLDatasets,
):
"""We already have a test for the "non-private" version of this,
which primarily just uses the respective Client to determine if
@@ -443,7 +419,7 @@ class TestDFCollectionItemMethod:
# Load up the data that we'll be using for various to_archive
# methods. Will always be empty pd.DataFrames for now...
- df = item.from_mysql()
+ df = item.from_db()
ddf = dd.from_pandas(df, npartitions=1)
# (1) Confirm a missing ddf (shouldn't bc of type hint) should
@@ -484,19 +460,28 @@ class TestDFCollectionItemMethod:
@pytest.mark.skip
def test_to_archive_numbered_partial(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
@pytest.mark.skip
def test_initial_load(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
@pytest.mark.skip
def test_clear_corrupt_archive(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
@@ -506,37 +491,36 @@ class TestDFCollectionItemMethod:
argvalues=list(iter_product(df_collections, ["12h", "10D"], [timedelta(days=15)])),
)
class TestDFCollectionItemMethodBase:
-
- @pytest.mark.skip
- def test_path_exists(
- self, df_collection_data_type, offset: str, duration: timedelta
- ):
- pass
-
- @pytest.mark.skip
- def test_next_numbered_path(
- self, df_collection_data_type, offset: str, duration: timedelta
- ):
- pass
-
@pytest.mark.skip
def test_search_highest_numbered_path(
- self, df_collection_data_type, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
@pytest.mark.skip
def test_tmp_filename(
- self, df_collection_data_type, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
@pytest.mark.skip
- def test_tmp_path(self, df_collection_data_type, offset: str, duration: timedelta):
+ def test_tmp_path(
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
+ ):
pass
def test_is_empty(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection: DFCollection,
):
"""
test_has_empty was merged into this because item.has_empty is
@@ -553,7 +537,8 @@ class TestDFCollectionItemMethodBase:
assert item.has_empty()
def test_has_partial_archive(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection: DFCollection,
):
for item in df_collection.items:
assert not item.has_partial_archive()
@@ -561,7 +546,8 @@ class TestDFCollectionItemMethodBase:
assert item.has_partial_archive()
def test_has_archive(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection: DFCollection,
):
for item in df_collection.items:
# (1) Originally, nothing exists... so let's just make a file and
@@ -598,7 +584,8 @@ class TestDFCollectionItemMethodBase:
assert item.has_archive(include_empty=True)
def test_delete_archive(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection: DFCollection,
):
for item in df_collection.items:
item: DFCollectionItem
@@ -621,9 +608,11 @@ class TestDFCollectionItemMethodBase:
assert not item.partial_path.exists()
def test_should_archive(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection: DFCollection,
):
- schema: DataFrameSchema = df_collection._schema
+ schema: DataFrameSchema = df_collection.type_schema
+ assert schema.metadata
aa = schema.metadata[ARCHIVE_AFTER]
# It shouldn't be None, it can be timedelta(seconds=0)
@@ -632,19 +621,23 @@ class TestDFCollectionItemMethodBase:
for item in df_collection.items:
item: DFCollectionItem
- if datetime.now(tz=timezone.utc) > item.finish + aa:
+ if datetime.now(tz=UTC) > item.finish + aa:
assert item.should_archive()
else:
assert not item.should_archive()
@pytest.mark.skip
def test_set_empty(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
def test_valid_archive(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection: DFCollection,
):
# Originally, nothing has been saved or anything.. so confirm it
# always comes back as None
@@ -668,18 +661,28 @@ class TestDFCollectionItemMethodBase:
@pytest.mark.skip
def test_validate_df(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
@pytest.mark.skip
def test_from_archive(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
def test__to_dict(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
+ df_collection: DFCollection,
):
for item in df_collection.items:
@@ -698,30 +701,39 @@ class TestDFCollectionItemMethodBase:
@pytest.mark.skip
def test_delete_partial(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
@pytest.mark.skip
def test_cleanup_partials(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
@pytest.mark.skip
def test_delete_dangling_partials(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+ duration: timedelta,
):
pass
@gen_cluster(client=True, nthreads=[("127.0.0.1", 1)])
-async def test_client(client, s, worker):
+async def test_client(client: DaskClient, s: DaskScheduler, worker: DaskWorker):
"""c,s,a are all required - the secondary Worker (b) is not required"""
- assert isinstance(client, Client)
- assert isinstance(s, Scheduler)
- assert isinstance(worker, Worker)
+ assert isinstance(client, DaskClient)
+ assert isinstance(s, DaskScheduler)
+ assert isinstance(worker, DaskWorker)
@pytest.mark.parametrize(
@@ -730,12 +742,18 @@ async def test_client(client, s, worker):
)
@gen_cluster(client=True, nthreads=[("127.0.0.1", 1)])
@pytest.mark.anyio
-async def test_client_parametrize(c, s, w, df_collection_data_type, offset: str):
+async def test_client_parametrize(
+ c: DaskClient,
+ s: DaskScheduler,
+ w: DaskWorker,
+ df_collection_data_type: DFCollectionType,
+ offset: str,
+):
"""c,s,a are all required - the secondary Worker (b) is not required"""
- assert isinstance(c, Client), f"c is not Client, it's {type(c)}"
- assert isinstance(s, Scheduler), f"s is not Scheduler, it's {type(s)}"
- assert isinstance(w, Worker), f"w is not Worker, it's {type(w)}"
+ assert isinstance(c, DaskClient), f"c is not Client, it's {type(c)}"
+ assert isinstance(s, DaskScheduler), f"s is not Scheduler, it's {type(s)}"
+ assert isinstance(w, DaskWorker), f"w is not Worker, it's {type(w)}"
assert df_collection_data_type is not None
assert isinstance(offset, str)
@@ -751,22 +769,15 @@ async def test_client_parametrize(c, s, w, df_collection_data_type, offset: str)
argvalues=list(iter_product(df_collections, ["12h", "10D"], [timedelta(days=15)])),
)
class TestDFCollectionItemFunctionalTest:
-
def test_to_archive_and_ddf(
self,
- df_collection_data_type,
- offset: str,
- duration: timedelta,
- client_no_amm,
- df_collection,
- user: User,
+ client_no_amm: DaskClient,
+ df_collection: DFCollection,
user_factory: Callable[..., User],
product: Product,
- incite_item_factory,
- delete_df_collection,
- mnt_filepath: GRLDatasets,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
):
- from generalresearch.models.thl.user import User
if df_collection.data_type in unsupported_mock_types:
return
@@ -804,17 +815,11 @@ class TestDFCollectionItemFunctionalTest:
def test_filesize_estimate(
self,
- df_collection,
- user: User,
- offset: str,
- duration: timedelta,
- client_no_amm,
+ df_collection: DFCollection,
user_factory: Callable[..., User],
product: Product,
- df_collection_data_type,
- incite_item_factory,
- delete_df_collection,
- mnt_filepath: GRLDatasets,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
):
"""A functional test to write some Parquet files for the
DFCollection and then confirm that the files get written
@@ -828,8 +833,6 @@ class TestDFCollectionItemFunctionalTest:
import pyarrow.parquet as pq
- from generalresearch.models.thl.user import User
-
if df_collection.data_type in unsupported_mock_types:
return
delete_df_collection(coll=df_collection)
@@ -853,18 +856,13 @@ class TestDFCollectionItemFunctionalTest:
def test_to_archive_client(
self,
- client_no_amm,
- df_collection,
+ client_no_amm: DaskClient,
+ df_collection: DFCollection,
user_factory: Callable[..., User],
product: Product,
- offset: str,
- duration: timedelta,
- df_collection_data_type,
- incite_item_factory,
- delete_df_collection,
- mnt_filepath: GRLDatasets,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
):
- from generalresearch.models.thl.user import User
delete_df_collection(coll=df_collection)
df_collection._client = client_no_amm
@@ -880,7 +878,7 @@ class TestDFCollectionItemFunctionalTest:
# Load up the data that we'll be using for various to_archive
# methods. Will always be empty pd.DataFrames for now...
- df = item.from_mysql()
+ df = item.from_db()
ddf = dd.from_pandas(df, npartitions=1)
assert isinstance(ddf, dd.DataFrame)
@@ -893,7 +891,8 @@ class TestDFCollectionItemFunctionalTest:
@pytest.mark.skip
def test_get_items(
- self, df_collection, product: Product, offset: str, duration: timedelta
+ self,
+ df_collection: DFCollection,
):
with pytest.warns(expected_warning=ResourceWarning) as cm:
df_collection.get_items_last365()
@@ -906,27 +905,22 @@ class TestDFCollectionItemFunctionalTest:
def test_saving_protections(
self,
- client_no_amm,
- df_collection_data_type,
- df_collection,
- incite_item_factory,
- delete_df_collection,
+ df_collection: DFCollection,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
user_factory: Callable[..., User],
product: Product,
- offset: str,
- duration: timedelta,
- mnt_filepath: GRLDatasets,
):
"""Don't allow creating an archive for data that will likely be
overwritten or updated
"""
- from generalresearch.models.thl.user import User
if df_collection.data_type in unsupported_mock_types:
return
u1: User = user_factory(product=product)
- schema: DataFrameSchema = df_collection._schema
+ schema: DataFrameSchema = df_collection.type_schema
+ assert schema.metadata
aa = schema.metadata[ARCHIVE_AFTER]
assert isinstance(aa, timedelta)
@@ -948,21 +942,14 @@ class TestDFCollectionItemFunctionalTest:
def test_empty_item(
self,
- client_no_amm,
- df_collection_data_type,
- df_collection,
- incite_item_factory,
- delete_df_collection,
- user: User,
- offset: str,
- duration: timedelta,
- mnt_filepath: GRLDatasets,
+ df_collection: DFCollection,
+ delete_df_collection: Callable[..., None],
):
delete_df_collection(coll=df_collection)
for item in df_collection.items:
assert not item.has_empty()
- df: pd.DataFrame = item.from_mysql()
+ df: pd.DataFrame = item.from_db()
# We do this check b/c the Ledger returns back None and
# I don't want it to fail when we go to make a ddf
@@ -976,18 +963,13 @@ class TestDFCollectionItemFunctionalTest:
def test_file_touching(
self,
- client_no_amm,
- df_collection_data_type,
- df_collection,
- incite_item_factory,
- delete_df_collection,
+ client_no_amm: DaskClient,
+ df_collection: DFCollection,
+ incite_item_factory: Callable[..., None],
+ delete_df_collection: Callable[..., None],
user_factory: Callable[..., User],
product: Product,
- offset: str,
- duration: timedelta,
- mnt_filepath,
):
- from generalresearch.models.thl.user import User
delete_df_collection(coll=df_collection)
df_collection._client = client_no_amm