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.py273
1 files changed, 131 insertions, 142 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 8038d3b..edf90f7 100644
--- a/tests/incite/collections/test_df_collection_item_thl_web.py
+++ b/tests/incite/collections/test_df_collection_item_thl_web.py
@@ -5,12 +5,12 @@ 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
from uuid import uuid4
import dask.dataframe as dd
import pandas as pd
import pytest
+from dask.distributed import Client as DaskClient
from distributed import Client, Scheduler, Worker
# noinspection PyUnresolvedReferences
@@ -21,20 +21,19 @@ from faker import Faker
from pandera.pandas import DataFrameSchema
from pydantic import FilePath
-from generalresearch.incite.base import CollectionItemBase
+from generalresearch.incite.base import CollectionItemBase, GRLDatasets
from generalresearch.incite.collections import (
+ DFCollection,
DFCollectionItem,
DFCollectionType,
)
from generalresearch.incite.schemas import ARCHIVE_AFTER
+from generalresearch.managers.thl.ledger_manager.thl_ledger import ThlLedgerManager
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
-
fake = Faker()
df_collections = [
@@ -72,7 +71,12 @@ class TestDFCollectionItemBase:
)
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)
@@ -89,37 +93,59 @@ class TestDFCollectionItemProperties:
)
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)
@@ -138,11 +164,8 @@ class TestDFCollectionItemMethod:
def test_has_mysql(
self,
- df_collection,
+ df_collection: DFCollection,
thl_web_rr: PostgresConfig,
- offset: str,
- duration: timedelta,
- df_collection_data_type,
delete_df_collection: Callable[..., None],
):
delete_df_collection(coll=df_collection)
@@ -168,12 +191,6 @@ class TestDFCollectionItemMethod:
@pytest.mark.skip
def test_update_partial_archive(
self,
- df_collection,
- offset: str,
- duration: timedelta,
- thl_web_rw: PostgresConfig,
- df_collection_data_type,
- delete_df_collection: Callable[..., None],
):
# for i in collection.items:
# assert i.update_partial_archive()
@@ -183,28 +200,12 @@ class TestDFCollectionItemMethod:
@pytest.mark.skip
def test_create_partial_archive(
self,
- df_collection,
- offset: str,
- duration: str,
- create_main_accounts: Callable[..., None],
- thl_web_rw: PostgresConfig,
- thl_lm,
- df_collection_data_type,
- user_factory: Callable[..., User],
- product: product: Product,
- client_no_amm,
- incite_item_factory,
- delete_df_collection: Callable[..., None],
- mnt_filepath: GRLDatasets,
):
assert 1 + 1 == 2
def test_dict(
self,
- df_collection_data_type,
- offset: str,
- duration: timedelta,
- df_collection,
+ df_collection: DFCollection,
delete_df_collection: Callable[..., None],
):
delete_df_collection(coll=df_collection)
@@ -225,15 +226,15 @@ 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: Callable[..., None],
thl_web_rw: PostgresConfig,
user_factory: Callable[..., User],
- product: product: Product,
- incite_item_factory,
+ product: Product,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
):
@@ -253,12 +254,14 @@ class TestDFCollectionItemMethod:
if df_collection.data_type == DFCollectionType.LEDGER:
assert df is None
else:
+ assert isinstance(df, pd.DataFrame)
assert df.empty
assert set(df.columns) == set(df_collection._schema.columns.keys())
incite_item_factory(user=u1, item=item)
df = item.from_mysql()
+ 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:
@@ -270,13 +273,13 @@ class TestDFCollectionItemMethod:
def test_from_mysql_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: Product,
- incite_item_factory,
+ product: Product,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
):
@@ -293,7 +296,7 @@ class TestDFCollectionItemMethod:
# 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()
+ _ = item.from_mysql_standard()
assert (
"Can't call from_mysql_standard for Ledger DFCollectionItem"
in str(cm.value)
@@ -304,32 +307,34 @@ class TestDFCollectionItemMethod:
# Unlike .from_mysql_ledger(), .from_mysql_standard() will return
# back and empty df with the correct columns in place
df = item.from_mysql_standard()
+ assert isinstance(df, pd.DataFrame)
assert df.empty
assert set(df.columns) == set(df_collection._schema.columns.keys())
incite_item_factory(user=u1, item=item)
df = item.from_mysql_standard()
+ assert isinstance(df, pd.DataFrame)
assert not df.empty
assert set(df.columns) == set(df_collection._schema.columns.keys())
assert df.shape[0] > 0
def test_from_mysql_ledger(
self,
- df_collection,
+ df_collection: DFCollection,
user: User,
create_main_accounts: Callable[..., None],
offset: str,
duration: timedelta,
thl_web_rw: PostgresConfig,
- thl_lm,
- df_collection_data_type,
+ thl_ledger_manager: ThlLedgerManager,
+ df_collection_data_type: DFCollectionType,
user_factory: Callable[..., User],
- product: product: Product,
- client_no_amm,
- incite_item_factory,
+ product: Product,
+ client_no_amm: DaskClient,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
- mnt_filepath,
+ mnt_filepath: GRLDatasets,
):
if df_collection.data_type != DFCollectionType.LEDGER:
@@ -370,17 +375,17 @@ class TestDFCollectionItemMethod:
def test_to_archive(
self,
- df_collection,
+ df_collection: DFCollection,
user: User,
offset: str,
duration: timedelta,
- df_collection_data_type,
+ df_collection_data_type: DFCollectionType,
user_factory: Callable[..., User],
- product: product: Product,
- client_no_amm,
- incite_item_factory,
+ product: Product,
+ client_no_amm: DaskClient,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
- mnt_filepath,
+ mnt_filepath: GRLDatasets,
):
if df_collection.data_type in unsupported_mock_types:
@@ -407,17 +412,17 @@ class TestDFCollectionItemMethod:
def test__to_archive(
self,
- df_collection_data_type,
- df_collection,
+ df_collection_data_type: DFCollectionType,
+ df_collection: DFCollection,
user_factory: Callable[..., User],
- product: product: Product,
+ product: Product,
offset: str,
duration: timedelta,
- client_no_amm,
+ client_no_amm: DaskClient,
user: User,
- incite_item_factory,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
- mnt_filepath,
+ 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
@@ -480,19 +485,19 @@ class TestDFCollectionItemMethod:
@pytest.mark.skip
def test_to_archive_numbered_partial(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
):
pass
@pytest.mark.skip
def test_initial_load(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
):
pass
@pytest.mark.skip
def test_clear_corrupt_archive(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
):
pass
@@ -505,34 +510,40 @@ class TestDFCollectionItemMethodBase:
@pytest.mark.skip
def test_path_exists(
- self, df_collection_data_type, offset: str, duration: timedelta
+ self,
):
pass
@pytest.mark.skip
def test_next_numbered_path(
- self, df_collection_data_type, offset: str, duration: timedelta
+ self,
):
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,
):
pass
@pytest.mark.skip
- def test_tmp_path(self, df_collection_data_type, offset: str, duration: timedelta):
+ def test_tmp_path(
+ self,
+ ):
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
@@ -549,7 +560,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()
@@ -557,7 +569,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
@@ -594,7 +607,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
@@ -617,7 +631,8 @@ 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
aa = schema.metadata[ARCHIVE_AFTER]
@@ -635,12 +650,13 @@ class TestDFCollectionItemMethodBase:
@pytest.mark.skip
def test_set_empty(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
):
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
@@ -664,18 +680,19 @@ class TestDFCollectionItemMethodBase:
@pytest.mark.skip
def test_validate_df(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
):
pass
@pytest.mark.skip
def test_from_archive(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
):
pass
def test__to_dict(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
+ df_collection: DFCollection,
):
for item in df_collection.items:
@@ -694,19 +711,19 @@ class TestDFCollectionItemMethodBase:
@pytest.mark.skip
def test_delete_partial(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
):
pass
@pytest.mark.skip
def test_cleanup_partials(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
):
pass
@pytest.mark.skip
def test_delete_dangling_partials(
- self, df_collection_data_type, df_collection, offset: str, duration: timedelta
+ self,
):
pass
@@ -726,7 +743,9 @@ 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, s, w, 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)}"
@@ -750,17 +769,12 @@ 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: Product,
- incite_item_factory,
+ product: Product,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
- mnt_filepath: GRLDatasets,
):
if df_collection.data_type in unsupported_mock_types:
@@ -799,17 +813,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: Product,
- df_collection_data_type,
- incite_item_factory,
+ product: Product,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
- mnt_filepath: GRLDatasets,
):
"""A functional test to write some Parquet files for the
DFCollection and then confirm that the files get written
@@ -846,16 +854,12 @@ 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: Product,
- offset: str,
- duration: timedelta,
- df_collection_data_type,
- incite_item_factory,
+ product: Product,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
- mnt_filepath: GRLDatasets,
):
delete_df_collection(coll=df_collection)
@@ -885,7 +889,8 @@ class TestDFCollectionItemFunctionalTest:
@pytest.mark.skip
def test_get_items(
- self, df_collection, product: product: Product, offset: str, duration: timedelta
+ self,
+ df_collection: DFCollection,
):
with pytest.warns(expected_warning=ResourceWarning) as cm:
df_collection.get_items_last365()
@@ -898,16 +903,11 @@ class TestDFCollectionItemFunctionalTest:
def test_saving_protections(
self,
- client_no_amm,
- df_collection_data_type,
- df_collection,
- incite_item_factory,
+ df_collection: DFCollection,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
user_factory: Callable[..., User],
- product: product: Product,
- offset: str,
- duration: timedelta,
- mnt_filepath: GRLDatasets,
+ product: Product,
):
"""Don't allow creating an archive for data that will likely be
overwritten or updated
@@ -939,15 +939,8 @@ class TestDFCollectionItemFunctionalTest:
def test_empty_item(
self,
- client_no_amm,
- df_collection_data_type,
- df_collection,
- incite_item_factory,
+ df_collection: DFCollection,
delete_df_collection: Callable[..., None],
- user: User,
- offset: str,
- duration: timedelta,
- mnt_filepath: GRLDatasets,
):
delete_df_collection(coll=df_collection)
@@ -967,16 +960,12 @@ class TestDFCollectionItemFunctionalTest:
def test_file_touching(
self,
- client_no_amm,
- df_collection_data_type,
- df_collection,
- incite_item_factory,
+ client_no_amm: DaskClient,
+ df_collection: DFCollection,
+ incite_item_factory: Callable[..., None],
delete_df_collection: Callable[..., None],
user_factory: Callable[..., User],
- product: product: Product,
- offset: str,
- duration: timedelta,
- mnt_filepath,
+ product: Product,
):
delete_df_collection(coll=df_collection)