155 lines
4.1 KiB
Python
155 lines
4.1 KiB
Python
"""
|
|
End-to-end tests for the Scouter main workflow (Mongo load path).
|
|
"""
|
|
|
|
import json
|
|
|
|
import pytest
|
|
from redis import Redis
|
|
|
|
from e2e.conftest import E2E_DATABASE
|
|
from e2e.helpers import (
|
|
count_laborious_rows,
|
|
load_scenario_input,
|
|
make_workflow_id,
|
|
seed_last_data_timestamp,
|
|
seed_raw_collection,
|
|
start_and_await_workflow,
|
|
)
|
|
from scouter.workflow.scouter import Scouter
|
|
from temporalio.testing import WorkflowEnvironment
|
|
from temporalio.worker import Worker
|
|
|
|
|
|
@pytest.mark.e2e
|
|
@pytest.mark.asyncio
|
|
async def test_scenario_1_1_1_happy_path(
|
|
temporal_env: WorkflowEnvironment,
|
|
temporal_worker: Worker,
|
|
mongo_uri: str,
|
|
redis_client: Redis,
|
|
postgres_engine,
|
|
):
|
|
scenario = load_scenario_input('scouter_happy_path')
|
|
workflow_input = scenario['workflow_input']
|
|
seed_raw_collection(
|
|
mongo_uri,
|
|
E2E_DATABASE,
|
|
workflow_input['schedule_name'],
|
|
scenario['raw_documents'],
|
|
)
|
|
|
|
await start_and_await_workflow(
|
|
temporal_env.client,
|
|
Scouter.run,
|
|
workflow_input,
|
|
make_workflow_id('scouter-happy'),
|
|
)
|
|
|
|
model_id = workflow_input['model_id']
|
|
assert count_laborious_rows(postgres_engine, model_id) >= 1
|
|
|
|
key = f"last_data_timestamp:scouter:{workflow_input['schedule_name']}"
|
|
assert redis_client.get(key) is not None
|
|
|
|
|
|
@pytest.mark.e2e
|
|
@pytest.mark.asyncio
|
|
async def test_scenario_1_2_1_incremental_load(
|
|
temporal_env: WorkflowEnvironment,
|
|
temporal_worker: Worker,
|
|
mongo_uri: str,
|
|
redis_client: Redis,
|
|
postgres_engine,
|
|
):
|
|
scenario = load_scenario_input('scouter_incremental_load')
|
|
workflow_input = scenario['workflow_input']
|
|
redis_seed = scenario['redis_seed']
|
|
seed_last_data_timestamp(
|
|
redis_client,
|
|
'scouter',
|
|
workflow_input['schedule_name'],
|
|
redis_seed['last_data_timestamp'],
|
|
)
|
|
seed_raw_collection(
|
|
mongo_uri,
|
|
E2E_DATABASE,
|
|
workflow_input['schedule_name'],
|
|
scenario['raw_documents'],
|
|
)
|
|
|
|
await start_and_await_workflow(
|
|
temporal_env.client,
|
|
Scouter.run,
|
|
workflow_input,
|
|
make_workflow_id('scouter-incremental'),
|
|
)
|
|
|
|
model_id = workflow_input['model_id']
|
|
assert count_laborious_rows(postgres_engine, model_id) == scenario['expected_newer_count']
|
|
|
|
stored = json.loads(
|
|
redis_client.get(f"last_data_timestamp:scouter:{workflow_input['schedule_name']}")
|
|
)
|
|
assert stored == scenario['expected_last_timestamp']
|
|
|
|
|
|
@pytest.mark.e2e
|
|
@pytest.mark.asyncio
|
|
async def test_scenario_1_3_1_empty_mongo_early_exit(
|
|
temporal_env: WorkflowEnvironment,
|
|
temporal_worker: Worker,
|
|
mongo_uri: str,
|
|
redis_client: Redis,
|
|
postgres_engine,
|
|
):
|
|
scenario = load_scenario_input('scouter_empty_mongo')
|
|
workflow_input = scenario['workflow_input']
|
|
seed_raw_collection(
|
|
mongo_uri,
|
|
E2E_DATABASE,
|
|
workflow_input['schedule_name'],
|
|
scenario['raw_documents'],
|
|
)
|
|
|
|
await start_and_await_workflow(
|
|
temporal_env.client,
|
|
Scouter.run,
|
|
workflow_input,
|
|
make_workflow_id('scouter-empty'),
|
|
)
|
|
|
|
assert count_laborious_rows(postgres_engine, workflow_input['model_id']) == 0
|
|
key = f"last_data_timestamp:scouter:{workflow_input['schedule_name']}"
|
|
assert redis_client.get(key) is None
|
|
|
|
|
|
@pytest.mark.e2e
|
|
@pytest.mark.asyncio
|
|
async def test_scenario_1_3_2_no_redis_timestamp_first_run(
|
|
temporal_env: WorkflowEnvironment,
|
|
temporal_worker: Worker,
|
|
mongo_uri: str,
|
|
redis_client: Redis,
|
|
postgres_engine,
|
|
):
|
|
scenario = load_scenario_input('scouter_no_redis_timestamp_first_run')
|
|
workflow_input = scenario['workflow_input']
|
|
seed_raw_collection(
|
|
mongo_uri,
|
|
E2E_DATABASE,
|
|
workflow_input['schedule_name'],
|
|
scenario['raw_documents'],
|
|
)
|
|
|
|
await start_and_await_workflow(
|
|
temporal_env.client,
|
|
Scouter.run,
|
|
workflow_input,
|
|
make_workflow_id('scouter-first-run'),
|
|
)
|
|
|
|
assert count_laborious_rows(postgres_engine, workflow_input['model_id']) >= 1
|
|
key = f"last_data_timestamp:scouter:{workflow_input['schedule_name']}"
|
|
assert redis_client.get(key) is not None
|