Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 2 additions & 4 deletions .vscode/settings.json
Original file line number Diff line number Diff line change
@@ -1,8 +1,6 @@
{
"python.testing.pytestArgs": [
".",
"--ignore=test_data"
],
"python.testing.pytestArgs": [],
"python.testing.cwd": "${workspaceFolder}",
"python.testing.unittestEnabled": false,
"python.testing.pytestEnabled": true,
"python.analysis.extraPaths": [
Expand Down
Empty file.
Empty file.
3 changes: 1 addition & 2 deletions libs/datapipe-cvat/tests/test_simple_project_example.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,8 +14,7 @@
from datapipe_cvat.cvat_step import create_cvat_client
from PIL import Image

from test_cvat_integration import _require_cvat

from .test_cvat_integration import _require_cvat

pytestmark = pytest.mark.cvat

Expand Down
14 changes: 7 additions & 7 deletions libs/datapipe-label-studio/tests/test_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,17 +10,17 @@
from datapipe.step.datatable_transform import DatatableTransformStep
from datapipe.store.database import TableStoreDB
from datapipe.types import data_to_index
from label_studio_sdk import LabelStudio
from pytest_cases import parametrize, parametrize_with_cases
from sqlalchemy.sql.schema import Column
from sqlalchemy.sql.sqltypes import JSON, String

from datapipe_label_studio.sdk_utils import get_project_by_title
from datapipe_label_studio.upload_predictions_pipeline import (
LabelStudioUploadPredictions,
)
from datapipe_label_studio.upload_tasks_pipeline import LabelStudioUploadTasks
from tests.ls_test_helpers import (
from label_studio_sdk import LabelStudio
from pytest_cases import parametrize, parametrize_with_cases
from sqlalchemy.sql.schema import Column
from sqlalchemy.sql.sqltypes import JSON, String

from .ls_test_helpers import (
DELETE_UNANNOTATED_TASKS_ONLY_ON_UPDATE,
INCLUDE_PARAMS,
INCLUDE_PREDICTIONS,
Expand All @@ -30,7 +30,7 @@
convert_to_ls_input_data,
wrapped_partial,
)
from tests.util import get_project_id, get_project_tasks, wait_until_label_studio_is_up
from .util import get_project_id, get_project_tasks, wait_until_label_studio_is_up


def gen_data_df():
Expand Down
14 changes: 7 additions & 7 deletions libs/datapipe-label-studio/tests/test_pipeline_many_projects.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,11 +9,6 @@
from datapipe.step.batch_transform import BatchTransform
from datapipe.step.datatable_transform import DatatableTransformStep
from datapipe.store.database import TableStoreDB
from label_studio_sdk import LabelStudio
from pytest_cases import parametrize, parametrize_with_cases
from sqlalchemy.sql.schema import Column
from sqlalchemy.sql.sqltypes import JSON, String

from datapipe_label_studio.create_projects_step import CreateLabelStudioProjects
from datapipe_label_studio.sdk_utils import get_project_by_title
from datapipe_label_studio.upload_predictions_pipeline import (
Expand All @@ -22,7 +17,12 @@
from datapipe_label_studio.upload_tasks_pipeline import (
LabelStudioUploadTasksToProjects,
)
from tests.ls_test_helpers import (
from label_studio_sdk import LabelStudio
from pytest_cases import parametrize, parametrize_with_cases
from sqlalchemy.sql.schema import Column
from sqlalchemy.sql.sqltypes import JSON, String

from .ls_test_helpers import (
DELETE_UNANNOTATED_TASKS_ONLY_ON_UPDATE,
INCLUDE_PARAMS,
INCLUDE_PREDICTIONS,
Expand All @@ -32,7 +32,7 @@
convert_to_ls_input_data,
wrapped_partial,
)
from tests.util import get_project_id, get_project_tasks, wait_until_label_studio_is_up
from .util import get_project_id, get_project_tasks, wait_until_label_studio_is_up


def gen_ls_project_setting():
Expand Down
6 changes: 3 additions & 3 deletions libs/datapipe-ml/tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

os.environ["SQLALCHEMY_WARN_20"] = "1"

from tests.helpers.test_env import load_test_env
from .helpers.test_env import load_test_env

load_test_env()

Expand All @@ -13,8 +13,8 @@
from sqlalchemy import Column, create_engine, text
from sqlalchemy.sql.sqltypes import JSON, String

from tests.fixtures.smoke_data import SmokeDataset, make_smoke_dataset
from tests.helpers.dbconn import get_sqlite_dbconnstr
from .fixtures.smoke_data import SmokeDataset, make_smoke_dataset
from .helpers.dbconn import get_sqlite_dbconnstr


@pytest.fixture(
Expand Down
2 changes: 1 addition & 1 deletion libs/datapipe-ml/tests/helpers/checkpoint_kill9.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
import time
from pathlib import Path

from tests.helpers.checkpoint_fixtures import write_corrupt_zip_checkpoint, write_valid_zip_checkpoint
from .checkpoint_fixtures import write_corrupt_zip_checkpoint, write_valid_zip_checkpoint


def _child_atomic_save_before_replace(checkpoint_path: str, ready_queue: mp.Queue) -> None:
Expand Down
2 changes: 1 addition & 1 deletion libs/datapipe-ml/tests/helpers/cloud_smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@

import pytest

from tests.helpers.training_smoke import (
from .training_smoke import (
Workdir,
classification_freeze_step,
classification_train_step,
Expand Down
11 changes: 5 additions & 6 deletions libs/datapipe-ml/tests/helpers/failure_injection.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,22 +7,23 @@
from typing import TYPE_CHECKING

import pytest

from datapipe_ml.frameworks.yolo.checkpoint_sync import max_completed_epoch_from_run_dir
from datapipe_ml.training.sync import manifest_path_for_run, read_checkpoint_manifest

from tests.helpers.failure_injection_bootstrap import (
from .failure_injection_bootstrap import (
FAIL_AFTER_EPOCH_ENV,
FAIL_MODE_ENV,
checkpoint_for_epoch_exists,
configured_fail_after_epoch,
configured_fail_mode,
install_training_failure_hooks as _install_training_failure_hooks,
install_training_failure_hooks_direct,
maybe_fail_after_epoch,
run_dir_for_pipe_death_poll,
run_training_with_failure_hooks,
)
from .failure_injection_bootstrap import (
install_training_failure_hooks as _install_training_failure_hooks,
)

if TYPE_CHECKING:
from _pytest.monkeypatch import MonkeyPatch
Expand Down Expand Up @@ -68,9 +69,7 @@ def _spawn_with_pipe_death(target, *args): # noqa: ANN001
except queue.Empty:
if not p.is_alive():
p.join()
raise RuntimeError(
f"Training subprocess exited before returning a result. exitcode={p.exitcode}"
)
raise RuntimeError(f"Training subprocess exited before returning a result. exitcode={p.exitcode}")
run_dir = run_dir_for_pipe_death_poll()
if fail_after is not None and run_dir and _pipe_death_checkpoint_ready(run_dir, fail_after):
if p.is_alive():
Expand Down
11 changes: 6 additions & 5 deletions libs/datapipe-ml/tests/helpers/recovery_pipe_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -28,15 +28,16 @@ def main(argv: list[str] | None = None) -> int:
if args.phase == "train":
_install_short_lease_status_manager()

from tests.helpers.failure_injection import install_pipe_death_hooks_direct
from tests.helpers.failure_injection_bootstrap import WORK_DIR_ENV
from tests.helpers.training_recovery import (
from datapipe.compute import Pipeline, build_compute, run_steps

from .failure_injection import install_pipe_death_hooks_direct
from .failure_injection_bootstrap import WORK_DIR_ENV
from .training_recovery import (
configure_recovery_steps,
make_recovery_runtime,
recovery_case_by_id,
)
from datapipe.compute import Pipeline, build_compute, run_steps
from tests.helpers.training_smoke import run_pipeline
from .training_smoke import run_pipeline

os.environ[WORK_DIR_ENV] = str(workdir)

Expand Down
14 changes: 7 additions & 7 deletions libs/datapipe-ml/tests/helpers/training_recovery.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,16 +13,16 @@
import pytest
from datapipe.compute import Pipeline, PipelineStep, build_compute
from datapipe.types import IndexDF

from datapipe_ml.frameworks.yolo.checkpoint_sync import max_completed_epoch_from_run_dir
from datapipe_ml.training.runs import active_lease
from datapipe_ml.training.specs import TrainingResumeConfig, TrainingSyncConfig
from datapipe_ml.frameworks.yolo.checkpoint_sync import max_completed_epoch_from_run_dir
from datapipe_ml.training.sync import (
manifest_path_for_run,
read_checkpoint_manifest,
verify_manifest_checkpoint,
)
from tests.helpers.training_smoke import (

from .training_smoke import (
SmokeRuntime,
Workdir,
detection_freeze_step,
Expand Down Expand Up @@ -230,7 +230,7 @@ def real_recovery_cases() -> list:

def recovery_case_by_id(case_id: str) -> RealRecoveryCase:
if case_id in TENSORFLOW_RECOVERY_CASE_IDS:
from tests.helpers.training_recovery_tensorflow import recovery_tensorflow_case_by_id
from .training_recovery_tensorflow import recovery_tensorflow_case_by_id

return recovery_tensorflow_case_by_id(case_id)
for param in real_recovery_torch_cases():
Expand All @@ -242,7 +242,7 @@ def recovery_case_by_id(case_id: str) -> RealRecoveryCase:

def recovery_extra_step_configure(case: RealRecoveryCase) -> dict[type, StepConfigureFn] | None:
if case.id in TENSORFLOW_RECOVERY_CASE_IDS:
from tests.helpers.training_recovery_tensorflow import tensorflow_step_configure
from .training_recovery_tensorflow import tensorflow_step_configure

return tensorflow_step_configure()
return None
Expand Down Expand Up @@ -330,7 +330,7 @@ def make_recovery_runtime(
tmp_path: Path, case: RealRecoveryCase, *, working_dir: Workdir | None = None
) -> tuple[SmokeRuntime, list]:
if case.id in TENSORFLOW_RECOVERY_CASE_IDS:
from tests.helpers.training_recovery_tensorflow import make_recovery_runtime as make_tf_recovery_runtime
from .training_recovery_tensorflow import make_recovery_runtime as make_tf_recovery_runtime

return make_tf_recovery_runtime(tmp_path, case, working_dir=working_dir)
workdir = working_dir if working_dir is not None else tmp_path
Expand Down Expand Up @@ -499,7 +499,7 @@ def invoke_real_train_callable_for_backfill(
step: RecoveryTrainStep,
) -> None:
if case.id in TENSORFLOW_RECOVERY_CASE_IDS:
from tests.helpers.training_recovery_tensorflow import (
from .training_recovery_tensorflow import (
invoke_real_train_callable_for_backfill as invoke_tf_backfill,
)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
import pytest
from datapipe.types import IndexDF

from tests.helpers.training_recovery import (
from .training_recovery import (
TENSORFLOW_RECOVERY_CASE_IDS,
RealRecoveryCase,
RecoveryTrainStep,
Expand All @@ -19,12 +19,13 @@
direct_train_kwargs,
make_runtime,
)
from tests.helpers.training_smoke import Workdir, classification_freeze_step, classification_train_step
from .training_smoke import Workdir, classification_freeze_step, classification_train_step

if TYPE_CHECKING:
from datapipe_ml.frameworks.tensorflow.classification_runner import TF_ClassificationTrainingConfig
from datapipe_ml.tasks.classification.train.tensorflow import Train_Tensorflow_ClassificationModel


def _configure_tf_configs(configs: list[TF_ClassificationTrainingConfig], epochs: int) -> None:
for config in configs:
config.epochs = epochs
Expand Down
35 changes: 25 additions & 10 deletions libs/datapipe-ml/tests/helpers/training_smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@
from pathlib import Path, PurePosixPath
from typing import Iterable

from tests.helpers.cloud_storage import (
from .cloud_storage import (
assert_model_path_under_working_dir,
assert_url_exists,
is_cloud_url,
Expand All @@ -18,8 +18,8 @@

Workdir = str | Path

import pandas as pd
import fsspec
import pandas as pd
from datapipe.compute import (
Catalog,
Pipeline,
Expand All @@ -33,6 +33,7 @@
from sklearn.model_selection import train_test_split
from sqlalchemy import Column
from sqlalchemy.sql.sqltypes import JSON, String

from .dbconn import get_sqlite_dbconnstr

TESTS_DIR = Path(__file__).parents[1]
Expand Down Expand Up @@ -232,7 +233,7 @@ def make_runtime(


def make_cloud_runtime(tmp_path: Path, suffix: str, **kwargs) -> tuple[SmokeRuntime, str]:
from tests.helpers.cloud_storage import cloud_working_dir
from .cloud_storage import cloud_working_dir

workdir = str(cloud_working_dir(suffix))
runtime = make_runtime(tmp_path, working_dir=workdir, **kwargs)
Expand Down Expand Up @@ -455,7 +456,9 @@ def detection_freeze_step(workdir: Workdir):
)


def detection_train_step(workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None):
def detection_train_step(
workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None
):
from datapipe_ml.tasks.detection.train.yolov8 import (
Train_YoloV8_DetectionModel,
YoloV8_TrainingConfig,
Expand Down Expand Up @@ -548,7 +551,9 @@ def detection_train_step_with_local_checkpoint(
)


def detection_train_step_with_augmentations(workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None):
def detection_train_step_with_augmentations(
workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None
):
from datapipe_ml.tasks.detection.train.yolov8 import (
Train_YoloV8_DetectionModel,
YoloV8_TrainingConfig,
Expand Down Expand Up @@ -600,7 +605,9 @@ def detection_train_step_with_augmentations(workdir: Workdir, *, local_scratch:
)


def detection_yolov5_train_step(workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None):
def detection_yolov5_train_step(
workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None
):
from datapipe_ml.tasks.detection.train.yolov5 import (
Train_YoloV5_DetectionModel,
YoloV5_TrainingConfig,
Expand Down Expand Up @@ -642,7 +649,9 @@ def detection_yolov5_train_step(workdir: Workdir, *, local_scratch: Path | None
)


def detection_yolov5_train_step_with_augmentations(workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None):
def detection_yolov5_train_step_with_augmentations(
workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None
):
from datapipe_ml.tasks.detection.train.yolov5 import (
Train_YoloV5_DetectionModel,
YoloV5_TrainingConfig,
Expand Down Expand Up @@ -738,7 +747,9 @@ def segmentation_freeze_step(workdir: Workdir):
)


def segmentation_train_step(workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None):
def segmentation_train_step(
workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None
):
from datapipe_ml.tasks.segmentation.train.yolov8 import (
Train_YoloV8_SegmentationModel,
YoloV8_TrainingConfig,
Expand Down Expand Up @@ -847,7 +858,9 @@ def keypoints_freeze_step(workdir: Workdir):
)


def keypoints_train_step(workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None):
def keypoints_train_step(
workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None
):
from datapipe_ml.tasks.keypoints.train.yolov8 import (
Train_YoloV8_KeypointsModel,
YoloV8_TrainingConfig,
Expand Down Expand Up @@ -974,7 +987,9 @@ def classification_freeze_step(workdir: Workdir):
)


def classification_train_step(workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None):
def classification_train_step(
workdir: Workdir, *, local_scratch: Path | None = None, filedir_fsspec_kwargs: dict | None = None
):
from datapipe_ml.tasks.classification.train.tensorflow import (
TF_ClassificationTrainingConfig,
Train_Tensorflow_ClassificationModel,
Expand Down
2 changes: 1 addition & 1 deletion libs/datapipe-ml/tests/test_app_smoke.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
from sqlalchemy import Column
from sqlalchemy.sql.sqltypes import JSON, String

from tests.fixtures.smoke_data import (
from .fixtures.smoke_data import (
SmokeDataset,
define_ground_truth_for_classification,
generate_smoke_data,
Expand Down
Loading
Loading