diff --git a/.github/workflows/pytest.yml b/.github/workflows/pytest.yml index 3e456ca2..041c1958 100644 --- a/.github/workflows/pytest.yml +++ b/.github/workflows/pytest.yml @@ -3,96 +3,107 @@ name: Run Tests on: [push] jobs: - pytest_3_10: + unit: runs-on: ubuntu-latest - timeout-minutes: 30 + timeout-minutes: 10 strategy: fail-fast: false + matrix: + python-version: ["3.10", "3.11", "3.12", "3.13"] steps: - uses: actions/checkout@v4 - - name: set up Python - uses: actions/setup-python@v5 - with: - python-version: "3.10" - - name: set up uv - uses: astral-sh/setup-uv@v3 + - uses: astral-sh/setup-uv@v3 + - name: install dependencies + # --locked fails on a lockfile that no longer matches pyproject.toml + run: uv sync --extra dev --locked --python ${{ matrix.python-version }} + - name: check interpreter + run: uv run --no-sync python -c "import sys; assert sys.version_info[:2] == tuple(map(int, '${{ matrix.python-version }}'.split('.'))), sys.version" + - name: test_unit + # no Learning Loop and no secrets needed, so this suite gates the slow ones + run: uv run --no-sync python -m pytest learning_loop_node/tests/unit -v + + pytest: + needs: + - unit + runs-on: ubuntu-latest + timeout-minutes: 45 + strategy: + fail-fast: false + # these suites share one Learning Loop instance, and test_general creates and deletes + # the project zauberzeug/pytest_nodelib_general by a fixed name -- two entries at once + # would delete each other's project + max-parallel: 1 + matrix: + # the floor requires-python promises, and the version every node ships with; the unit + # job covers 3.11 and 3.13 too + python-version: ["3.10", "3.12"] + steps: + - uses: actions/checkout@v4 + - uses: astral-sh/setup-uv@v3 - name: install test dependencies run: | sudo apt-get update sudo apt-get install libcurl4-openssl-dev libssl-dev jpeginfo - name: install dependencies - run: | - uv sync --extra dev --frozen + run: uv sync --extra dev --locked --python ${{ matrix.python-version }} - name: test_general env: LOOP_HOST: "preview.learning-loop.ai" LOOP_USERNAME: "admin" LOOP_PASSWORD: ${{ secrets.LEARNING_LOOP_ADMIN_PASSWORD }} - run: | - uv run python -m pytest "learning_loop_node/tests/general" -v + run: uv run --no-sync python -m pytest learning_loop_node/tests/general -v - name: test_detector env: LOOP_HOST: "preview.learning-loop.ai" LOOP_USERNAME: "admin" LOOP_PASSWORD: ${{ secrets.LEARNING_LOOP_ADMIN_PASSWORD }} - run: | - uv run python -m pytest learning_loop_node/tests/detector -v + run: uv run --no-sync python -m pytest learning_loop_node/tests/detector -v - name: test_mock_detector env: LOOP_HOST: "preview.learning-loop.ai" LOOP_USERNAME: "admin" LOOP_PASSWORD: ${{ secrets.LEARNING_LOOP_ADMIN_PASSWORD }} - run: | - uv run python -m pytest mock_detector -v - - pytest_3_13: - needs: - - pytest_3_10 - runs-on: ubuntu-latest - timeout-minutes: 30 - strategy: - fail-fast: false - steps: - - uses: actions/checkout@v4 - - name: set up Python - uses: actions/setup-python@v5 - with: - python-version: "3.13" - - name: set up uv - uses: astral-sh/setup-uv@v3 - - name: install test dependencies - run: | - sudo apt-get update - sudo apt-get install libcurl4-openssl-dev libssl-dev jpeginfo - - name: install dependencies - run: | - uv sync --extra dev --frozen + run: uv run --no-sync python -m pytest mock_detector -v - name: test_annotator env: LOOP_HOST: "preview.learning-loop.ai" LOOP_USERNAME: "admin" LOOP_PASSWORD: ${{ secrets.LEARNING_LOOP_ADMIN_PASSWORD }} - run: | - uv run python -m pytest learning_loop_node/tests/annotator -v + run: uv run --no-sync python -m pytest learning_loop_node/tests/annotator -v - name: test_trainer env: LOOP_HOST: "preview.learning-loop.ai" LOOP_USERNAME: "admin" LOOP_PASSWORD: ${{ secrets.LEARNING_LOOP_ADMIN_PASSWORD }} - run: | - uv run python -m pytest learning_loop_node/tests/trainer -v + run: uv run --no-sync python -m pytest learning_loop_node/tests/trainer -v - name: test_mock_trainer env: LOOP_HOST: "preview.learning-loop.ai" LOOP_USERNAME: "admin" LOOP_PASSWORD: ${{ secrets.LEARNING_LOOP_ADMIN_PASSWORD }} + run: uv run --no-sync python -m pytest mock_trainer -v + + all-green: + # the one status check the branch ruleset requires: a matrix job's check run is named + # " ()", so requiring the entries themselves ties the ruleset to the matrix + # and every change to it leaves the ruleset waiting for a name nothing reports any more + needs: + - unit + - pytest + if: always() # a required check that is skipped reports nothing and waits forever + runs-on: ubuntu-latest + steps: + - name: check the jobs it gates + # skipped and cancelled have to fail too -- only success is green run: | - uv run python -m pytest mock_trainer -v + echo "unit: ${{ needs.unit.result }}, pytest: ${{ needs.pytest.result }}" + [ "${{ needs.unit.result }}" = success ] || exit 1 + [ "${{ needs.pytest.result }}" = success ] || exit 1 slack: needs: - - pytest_3_10 - - pytest_3_13 + - unit + - pytest if: always() # also execute when pytest fails runs-on: ubuntu-latest steps: diff --git a/AGENTS.md b/AGENTS.md index 2893b826..716853ff 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -53,15 +53,63 @@ A `DataExchanger` sits on top of the communicator to move images and model zips - **Annotator** — thin: it forwards the loop frontend's `handle_user_input` events into `AnnotatorLogic` and keeps a per-frontend history. +`trainer/subprocess.py`, `trainer/batch_size.py` and `trainer/metrics.py` hold the parts of a +trainer that are not framework-specific. `iterator_cpu_bound` runs a training generator in a +spawned process and yields its progress through a `maxsize=1` queue, so the event loop stays +responsive, CUDA state stays out of the node process, and the training can never run more than +one item ahead of the bookkeeping. `find_batch_size` probes for the largest power-of-two batch +that fits, around a `fits` predicate the trainer supplies — none of it imports torch, so a node +brings its own way of running a step. `macro_f1` scores the confusion matrix +`_get_new_best_training_state` returns. + +`trainer/cuda.py` is the one exception to that framework independence, and holds everything +about a batch-size probe that torch has to answer. `usable_memory_bytes` and `limit_cuda_memory` +turn a `--vram-limit-gb` setting into the budget a probe measures against and the cap that holds +the process to it, and capping an allocator has no NVML equivalent. `probe_batch_size` is the +whole probe for a node whose measurement is a single call — it resolves the limit, falls back +without a card, holds the safety margin and runs the search. A node that must build a throwaway +model first reserves the margin before building it, and so composes the same pieces itself: +`reserve_margin`, `measured_fits` and `find_batch_size`. `measured_fits` is where an +out-of-memory failure is told from a bug — both arrive as the same exception types, and a probe +that confuses them reports the smallest batch size as the card's fault. + +It imports torch, the package does **not** declare it, and only a trainer imports the module — so +the library keeps working where nothing trains. Its unit test installs a stand-in under the name +`torch`, which covers the arithmetic, the guards and the search; whether the cap holds, and what +a real step costs, can only be seen on a card. + +`helpers/entrypoint.py` holds what every node's `main.py` repeats: `node_parser` builds a +configargparse parser with `--host`/`--port`, `run_node` starts uvicorn. A setting is a flag +*and* an environment variable from one declaration — `--conf-threshold` reads +`CONF_THRESHOLD`. The exception is `--host`/`--port`, which read `NODE_HOST`/`NODE_PORT`: the +bare `HOST` already means the loop's address, and a node adopting it would hand it to uvicorn +and fail to bind. A node that used to require a prefix passes `legacy_env_prefix`, and the +prefixed names keep working with a warning. + +`detector/postprocess.py` and `detector/geometry.py` hold the parts of a detector that do *not* +depend on the model: confidence filtering, per-class NMS, box/point clipping, and turning +predictions into the loop's dataclasses. A node should import them rather than write its own — +every node repository had grown its own drifting copy, which is why they live here. Note the two +containers: `to_image_metadata` builds what a **detector** node reports, `to_detections` what a +**trainer**'s auto-detection pass reports. Both go through one routine, so the two paths cannot +drift apart again. + All node state lives under `GLOBALS.data_folder` (`DATA_FOLDER`, default `/data`): `uuids.json` (the node uuid is derived from its name and reused across restarts), `models/` plus the `current_model` symlink, `outbox/`, and the per-project training folders. ## Running and testing -The suites talk to a real Learning Loop instance and read their credentials from a local `.env` -(`LOOP_HOST`, `LOOP_USERNAME`, `LOOP_PASSWORD`). Without a reachable loop they cannot pass — do not -treat their failure as a regression you introduced. +The `unit` suite is self-contained — no Learning Loop, no credentials, no network — so it is the +one to run while iterating, and the one CI gates the others on: + +```bash +python -m pytest learning_loop_node/tests/unit -v +``` + +Every other suite talks to a real Learning Loop instance and reads its credentials from a local +`.env` (`LOOP_HOST`, `LOOP_USERNAME`, `LOOP_PASSWORD`). Without a reachable loop those cannot pass — +do not treat their failure as a regression you introduced. ```bash ./run_tests.sh # all suites @@ -92,7 +140,8 @@ uvx ruff check . A clean tree already reports several hundred ruff findings, so a clean run is not a reachable goal. Compare the count on the files you touched, before and after. -`.github/workflows/pytest.yml` runs the suites, `publish.yml` releases to PyPI on a tagged release. +`.github/workflows/pytest.yml` runs the suites (the `unit` job first, without secrets), +`publish.yml` releases to PyPI on a tagged release. ## Working in this repository diff --git a/README.md b/README.md index 3b3aea5f..dc4ba24b 100644 --- a/README.md +++ b/README.md @@ -36,6 +36,10 @@ You can configure connection to our Learning Loop by specifying the following en Note that organization and project IDs are always lower case and may differ from the names in the Learning Loop which can have uppercase letters. +Where a name has an alias, either spelling works. If both are set to **different** values the +prefixed name wins and a warning names the value used — the variable is never treated as unset, +which would otherwise let `LOOP_HOST` fall back to its default of `learning-loop.ai`. + #### Testing We use github actions for CI. Tests can also be executed locally by running diff --git a/learning_loop_node/detector/categories.py b/learning_loop_node/detector/categories.py new file mode 100644 index 00000000..67e7d27c --- /dev/null +++ b/learning_loop_node/detector/categories.py @@ -0,0 +1,27 @@ +"""Resolving the class indices or names a model emits against ``ModelInformation.categories``.""" + +from ..data_classes import Category, ModelInformation + + +def category_by_index(model_information: ModelInformation, index: int) -> Category: + """Resolve the category a model's class index refers to. + + :raises ValueError: If the index is outside the model's category list. + """ + categories = model_information.categories + if not 0 <= index < len(categories): + raise ValueError( + f'category index {index} is out of range for a model with {len(categories)} categories') + return categories[index] + + +def category_by_name(model_information: ModelInformation, name: str) -> Category: + """Resolve a category by name, for models whose outputs are named rather than indexed. + + :raises ValueError: If no category of that name exists. + """ + for category in model_information.categories: + if category.name == name: + return category + known = ', '.join(category.name for category in model_information.categories) + raise ValueError(f'unknown category name {name!r}; the model knows: {known}') diff --git a/learning_loop_node/detector/detector_node.py b/learning_loop_node/detector/detector_node.py index 90b06fe1..8d9a35e6 100644 --- a/learning_loop_node/detector/detector_node.py +++ b/learning_loop_node/detector/detector_node.py @@ -18,7 +18,6 @@ from ..data_classes import ( AboutResponse, - Category, Context, DetectorStatus, ImageMetadata, @@ -33,6 +32,7 @@ from ..helpers import background_tasks, environment_reader, run from ..helpers.misc import numpy_image_from_dict from ..node import Node +from .categories import category_by_name from .detector_logic import DetectorLogic, DetectorLogicFactory from .exceptions import NodeNeedsRestartError from .inbox_filter.relevance_filter import RelevanceFilter @@ -709,26 +709,28 @@ async def upload_images( await self.outbox.save(image, image_metadata, upload_priority) def add_category_id_to_detections(self, model_info: ModelInformation, image_metadata: ImageMetadata): - def find_category_id_by_name(categories: List[Category], category_name: str): - category_id = [category.id for category in categories if category.name == category_name] - return category_id[0] if category_id else '' - - for box_detection in image_metadata.box_detections: - category_name = box_detection.category_name - category_id = find_category_id_by_name(model_info.categories, category_name) - box_detection.category_id = category_id - for point_detection in image_metadata.point_detections: - category_name = point_detection.category_name - category_id = find_category_id_by_name(model_info.categories, category_name) - point_detection.category_id = category_id - for segmentation_detection in image_metadata.segmentation_detections: - category_name = segmentation_detection.category_name - category_id = find_category_id_by_name(model_info.categories, category_name) - segmentation_detection.category_id = category_id - for classification_detection in image_metadata.classification_detections: - category_name = classification_detection.category_name - category_id = find_category_id_by_name(model_info.categories, category_name) - classification_detection.category_id = category_id + """Resolve each detection's category id from its name, in metadata a client uploaded. + + A name the model does not know costs that one detection its id, not the whole upload. + """ + unknown_names: set[str] = set() + + def category_id_by_name(category_name: str) -> str: + try: + return category_by_name(model_info, category_name).id + except ValueError: + unknown_names.add(category_name) + return '' + + for detection in (*image_metadata.box_detections, + *image_metadata.point_detections, + *image_metadata.segmentation_detections, + *image_metadata.classification_detections): + detection.category_id = category_id_by_name(detection.category_name) + + if unknown_names: + self.log.warning('Model %s knows no category named %s', model_info.version, + ', '.join(sorted(unknown_names))) return image_metadata def register_sio_events(self, sio_client: AsyncClient): diff --git a/learning_loop_node/detector/geometry.py b/learning_loop_node/detector/geometry.py new file mode 100644 index 00000000..22f9c554 --- /dev/null +++ b/learning_loop_node/detector/geometry.py @@ -0,0 +1,39 @@ +"""Box and point clipping shared by every detector node. + +The loop stores a box as its top-left corner plus a size, which is the form :func:`clip_box` +takes and produces. +""" + + +def clip_box( + *, + x1: float, + y1: float, + width: float, + height: float, + img_width: int, + img_height: int, +) -> tuple[int, int, int, int]: + """Clip a top-left-anchored box to the image bounds. + + :return: The clipped ``(x1, y1, width, height)``; the size is never negative. + """ + x2 = x1 + width + y2 = y1 + height + + clipped_x1 = round(max(0.0, x1)) + clipped_y1 = round(max(0.0, y1)) + clipped_x2 = round(min(float(img_width), x2)) + clipped_y2 = round(min(float(img_height), y2)) + + clipped_width = max(clipped_x2 - clipped_x1, 0) + clipped_height = max(clipped_y2 - clipped_y1, 0) + + return clipped_x1, clipped_y1, clipped_width, clipped_height + + +def clip_point(x: float, y: float, img_width: int, img_height: int) -> tuple[float, float]: + """Clamp a point into the image bounds.""" + x = min(max(0, x), img_width) + y = min(max(0, y), img_height) + return x, y diff --git a/learning_loop_node/detector/postprocess.py b/learning_loop_node/detector/postprocess.py new file mode 100644 index 00000000..7157dc77 --- /dev/null +++ b/learning_loop_node/detector/postprocess.py @@ -0,0 +1,219 @@ +"""Model-agnostic detection postprocessing.""" + +import logging +from collections.abc import Sequence +from dataclasses import dataclass + +import numpy as np + +from ..data_classes import ( + BoxDetection, + Detections, + ImageMetadata, + ModelInformation, + PointDetection, +) +from ..enums import CategoryType +from .categories import category_by_index +from .geometry import clip_box, clip_point + +logger = logging.getLogger(__name__) + +MIN_BOX_SIZE: int = 2 + + +@dataclass(kw_only=True, slots=True, frozen=True) +class Prediction: + """One surviving prediction in the model's coordinates: unrounded top-left corner and size + in pixels, and the category as an index into :attr:`ModelInformation.categories`.""" + + x: float + y: float + width: float + height: float + category_index: int + confidence: float + + +def post_process( + boxes: np.ndarray, + scores: np.ndarray, + classes: np.ndarray, + *, + conf_threshold: float, + iou_threshold: float, + origin_h: int, + origin_w: int, +) -> list[Prediction]: + """Filter by confidence, run NMS, return what survives.""" + mask = scores > conf_threshold + boxes = boxes[mask].copy() + scores = scores[mask] + classes = classes[mask] + + if len(scores) == 0: + return [] + + boxes, scores, classes = non_max_suppression( + boxes, scores, classes, + iou_threshold=iou_threshold, origin_h=origin_h, origin_w=origin_w) + + return predictions_from_xyxy(labels=classes, boxes=boxes, scores=scores) + + +def non_max_suppression( + boxes: np.ndarray, + scores: np.ndarray, + classes: np.ndarray, + *, + iou_threshold: float, + origin_h: int, + origin_w: int, +) -> tuple[np.ndarray, np.ndarray, np.ndarray]: + """Clip to image bounds, sort by score descending, apply per-class NMS. + + :return: ``(boxes, scores, classes)`` — filtered arrays in the same order. + """ + boxes = boxes.copy() + boxes[:, 0] = np.clip(boxes[:, 0], 0, origin_w - 1) + boxes[:, 2] = np.clip(boxes[:, 2], 0, origin_w - 1) + boxes[:, 1] = np.clip(boxes[:, 1], 0, origin_h - 1) + boxes[:, 3] = np.clip(boxes[:, 3], 0, origin_h - 1) + + order = np.argsort(-scores) + boxes = boxes[order] + scores = scores[order] + classes = classes[order] + + keep_indices: list[int] = [] + for cls in np.unique(classes): + cls_mask = classes == cls + cls_indices = np.where(cls_mask)[0] + cls_boxes = boxes[cls_mask] + + while len(cls_indices) > 0: + keep_indices.append(cls_indices[0]) + if len(cls_indices) == 1: + break + ious = bbox_iou(cls_boxes[0:1], cls_boxes[1:]) + keep_mask = ious.flatten() <= iou_threshold + cls_indices = cls_indices[1:][keep_mask] + cls_boxes = cls_boxes[1:][keep_mask] + + keep_indices = sorted(keep_indices) + return boxes[keep_indices], scores[keep_indices], classes[keep_indices] + + +def bbox_iou( + box1: np.ndarray, + box2: np.ndarray, +) -> np.ndarray: + """Compute IoU between box1 (1x4) and box2 (Nx4), both in x1y1x2y2 format.""" + b1_x1, b1_y1, b1_x2, b1_y2 = box1[:, 0], box1[:, 1], box1[:, 2], box1[:, 3] + b2_x1, b2_y1, b2_x2, b2_y2 = box2[:, 0], box2[:, 1], box2[:, 2], box2[:, 3] + + inter_x1 = np.maximum(b1_x1, b2_x1) + inter_y1 = np.maximum(b1_y1, b2_y1) + inter_x2 = np.minimum(b1_x2, b2_x2) + inter_y2 = np.minimum(b1_y2, b2_y2) + + inter_area = np.clip(inter_x2 - inter_x1 + 1, 0, None) * np.clip(inter_y2 - inter_y1 + 1, 0, None) + b1_area = (b1_x2 - b1_x1 + 1) * (b1_y2 - b1_y1 + 1) + b2_area = (b2_x2 - b2_x1 + 1) * (b2_y2 - b2_y1 + 1) + + return inter_area / (b1_area + b2_area - inter_area + 1e-16) + + +def predictions_from_xyxy( + *, + labels: Sequence[float], + boxes: Sequence[Sequence[float]], + scores: Sequence[float], +) -> list[Prediction]: + """Convert xyxy model output into predictions, for models that suppress their own overlaps.""" + return [Prediction(x=float(x1), y=float(y1), width=float(x2 - x1), height=float(y2 - y1), + category_index=int(label), confidence=float(score)) + for label, (x1, y1, x2, y2), score in zip(labels, boxes, scores, strict=True)] + + +def to_image_metadata( + predictions: list[Prediction], + model_information: ModelInformation, + im_height: int, + im_width: int, +) -> ImageMetadata: + """Build the container a *detector* node reports.""" + image_metadata = ImageMetadata() + _append_predictions(image_metadata, predictions, model_information, im_height, im_width) + return image_metadata + + +def to_detections( + predictions: list[Prediction], + model_information: ModelInformation, + im_height: int, + im_width: int, + *, + image_id: str | None = None, +) -> Detections: + """Build the container a *trainer*'s auto-detection pass reports.""" + result = Detections(image_id=image_id) + _append_predictions(result, predictions, model_information, im_height, im_width) + return result + + +def _append_predictions( + target: ImageMetadata | Detections, + predictions: list[Prediction], + model_information: ModelInformation, + im_height: int, + im_width: int, +) -> None: + """Resolve each prediction's category and append it to ``target``, clipped to the image.""" + skipped_predictions = [] + + for prediction in predictions: + category = category_by_index(model_information, prediction.category_index) + if prediction.width <= MIN_BOX_SIZE or prediction.height <= MIN_BOX_SIZE: + skipped_predictions.append((category.name, prediction)) + continue + if category.type == CategoryType.Box: + clipped_x1, clipped_y1, clipped_w, clipped_h = clip_box( + x1=prediction.x, + y1=prediction.y, + width=prediction.width, + height=prediction.height, + img_width=im_width, + img_height=im_height, + ) + target.box_detections.append( + BoxDetection( + category_name=category.name, + x=clipped_x1, + y=clipped_y1, + width=clipped_w, + height=clipped_h, + category_id=category.id, + model_name=model_information.version, + confidence=prediction.confidence, + ) + ) + elif category.type == CategoryType.Point: + cx, cy = prediction.x + prediction.width / 2, prediction.y + prediction.height / 2 + cx, cy = clip_point(cx, cy, im_width, im_height) + target.point_detections.append( + PointDetection( + category_name=category.name, + x=cx, + y=cy, + category_id=category.id, + model_name=model_information.version, + confidence=prediction.confidence, + ) + ) + else: + logger.warning('Unsupported category type %s for category %s', category.type, category.name) + + if skipped_predictions: + log_msg = '\n'.join([str(p) for p in skipped_predictions]) + logger.warning('Removed %d small detections from result: \n%s', len(skipped_predictions), log_msg) diff --git a/learning_loop_node/helpers/entrypoint.py b/learning_loop_node/helpers/entrypoint.py new file mode 100644 index 00000000..4c42602e --- /dev/null +++ b/learning_loop_node/helpers/entrypoint.py @@ -0,0 +1,81 @@ +"""The boilerplate every node's ``main.py`` repeats: reading the settings, serving the node. + +Every setting is a flag *and* an environment variable named after it: ``--conf-threshold`` +reads ``CONF_THRESHOLD``. ``--host`` and ``--port`` are the exception, reading ``NODE_HOST`` +and ``NODE_PORT``: the bare ``HOST`` is the address of the loop. +""" + +import logging +import os +from argparse import Action, Namespace + +import configargparse +import uvicorn + +logger = logging.getLogger(__name__) + + +def node_parser(*, description: str, legacy_env_prefix: str = '') -> configargparse.ArgumentParser: + """Build the parser for a node, pre-loaded with the settings every node has. + + :param legacy_env_prefix: A prefix an earlier version of this node required, e.g. + ``'MY_DETECTOR_'``. Both spellings keep working, with a warning: the prefix on the + current name (``MY_DETECTOR_NODE_HOST``) and the prefix on the flag it was originally + applied to (``MY_DETECTOR_HOST``). Leave empty for a node that never used one. + """ + parser = _NodeArgumentParser(description=description, legacy_env_prefix=legacy_env_prefix) + parser.add_argument('--host', default='0.0.0.0', env_var='NODE_HOST', + help='Host interface to bind to') + parser.add_argument('--port', type=int, default=80, env_var='NODE_PORT', help='Port to bind to') + return parser + + +def run_node(app: str, args: Namespace) -> None: + """Serve the node. + + :param app: Import string of the node object, conventionally ``'main:node'``. + """ + reload = os.getenv('UVICORN_RELOAD', 'FALSE').lower() in ('true', '1') + logger.info('Uvicorn reload is set to: %s', reload) + uvicorn.run(app, host=args.host, port=args.port, lifespan='on', reload=reload) + + +class _NodeArgumentParser(configargparse.ArgumentParser): + """Parser whose every setting is also an environment variable named after its flag.""" + + def __init__(self, *, description: str, legacy_env_prefix: str) -> None: + super().__init__(description=description) + self.legacy_env_prefix = legacy_env_prefix + + def add_argument(self, *args, **kwargs) -> Action: # type: ignore[override] + action = super().add_argument(*args, **kwargs) + if getattr(action, 'env_var', None) is None and action.dest != 'help': + action.env_var = action.dest.upper() + return action + + def parse_args(self, *args, **kwargs) -> Namespace: # type: ignore[override] + self._adopt_legacy_env_vars() + return super().parse_args(*args, **kwargs) + + def _adopt_legacy_env_vars(self) -> None: + if not self.legacy_env_prefix: + return + for action in self._actions: + name = getattr(action, 'env_var', None) + if not name or name in os.environ: + continue + for legacy in self._legacy_names(action, name): + if legacy not in os.environ: + continue + os.environ[name] = os.environ[legacy] + logger.warning('%s is deprecated and will stop being read; set %s instead', legacy, name) + break + + def _legacy_names(self, action: Action, name: str) -> list[str]: + """The prefixed names this setting may still be configured under, preferred first. + + A prefix used to be applied to the *flag*, so the old name of ``--host`` is + ``HOST``, not ``NODE_HOST``. + """ + candidates = [self.legacy_env_prefix + name, self.legacy_env_prefix + action.dest.upper()] + return list(dict.fromkeys(candidates)) diff --git a/learning_loop_node/helpers/environment_reader.py b/learning_loop_node/helpers/environment_reader.py index 4833e1b2..002652f0 100644 --- a/learning_loop_node/helpers/environment_reader.py +++ b/learning_loop_node/helpers/environment_reader.py @@ -1,26 +1,32 @@ import logging import os -from typing import List, Optional + +logger = logging.getLogger(__name__) # TODO ignore_errors should default to False, but maybe some tests rely on this behavior -def read_from_env(possible_names: List[str], ignore_errors: bool = True) -> Optional[str]: +def read_from_env(possible_names: list[str], ignore_errors: bool = True) -> str | None: + """Read the first of ``possible_names`` that is set. + + :param possible_names: In order of preference; on a disagreement the first one set wins. + :raises ValueError: If nothing is set or the values disagree, unless ``ignore_errors``. + """ values = [os.environ.get(name, None) for name in possible_names] values = list(filter(None, values)) # Possible error: no values are set if not values: if ignore_errors: - logging.warning('no environment variable set for %s', possible_names) + logger.warning('no environment variable set for %s', possible_names) return None raise ValueError(f'no environment variable set for {possible_names}') # Possible error: multiple values are not None and not equal if len(values) > 1 and len(set(values)) > 1: - if ignore_errors: - logging.warning('different environment variables set for %s: %s', possible_names, values) - return None - raise ValueError(f'different environment variables set for {possible_names}: {values}') + if not ignore_errors: + raise ValueError(f'different environment variables set for {possible_names}: {values}') + logger.warning('different environment variables set for %s: %s - using %s', + possible_names, values, values[0]) return values[0] diff --git a/learning_loop_node/tests/unit/__init__.py b/learning_loop_node/tests/unit/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/learning_loop_node/tests/unit/pytest.ini b/learning_loop_node/tests/unit/pytest.ini new file mode 100644 index 00000000..e3ce9d31 --- /dev/null +++ b/learning_loop_node/tests/unit/pytest.ini @@ -0,0 +1,10 @@ +[pytest] +python_files = test_*.py +asyncio_mode = auto + +cache_dir = /tmp/pytest_cache + +# for debugging tests: +; log_cli_level = INFO +; log_cli_format = %(asctime)s [%(levelname)8s] %(message)s (%(filename)s:%(lineno)s) +; log_cli_date_format=%Y-%m-%d %H:%M:%S diff --git a/learning_loop_node/tests/unit/test_batch_size.py b/learning_loop_node/tests/unit/test_batch_size.py new file mode 100644 index 00000000..9ad3fd3e --- /dev/null +++ b/learning_loop_node/tests/unit/test_batch_size.py @@ -0,0 +1,115 @@ +from collections.abc import Callable + +import pytest + +from ...trainer.batch_size import ( + MIN_TRAIN_STEPS_PER_EPOCH, + batch_count, + dataset_limit, + find_batch_size, + is_out_of_memory, + no_gpu_batch_size, + smaller_pot, +) + + +def test_the_search_doubles_up_to_the_limit(): + fits, calls = _recording(1024) + assert find_batch_size(fits, limit=16) == 16 + assert calls == [1, 2, 4, 8, 16] + + +def test_the_search_backs_off_to_the_last_size_that_fit(): + fits, calls = _recording(20) + assert find_batch_size(fits, limit=64) == 16 + assert calls == [1, 2, 4, 8, 16, 32] # probes 32, fails, keeps 16 + + +def test_a_limit_that_is_not_a_power_of_two_is_rounded_down(): + assert find_batch_size(_fits_up_to(1024), limit=48) == 32 + assert find_batch_size(_fits_up_to(1024), limit=1) == 1 + + +def test_a_machine_that_cannot_take_one_sample_is_an_error(): + fits, calls = _recording(0) + with pytest.raises(RuntimeError, match='batch size 1 does not fit'): + find_batch_size(fits, limit=64) + assert calls == [1], 'must give up instead of probing larger sizes' + + +@pytest.mark.parametrize('capacity', range(1, 130)) +def test_the_result_is_always_the_largest_power_of_two_that_fits(capacity: int): + assert find_batch_size(_fits_up_to(capacity), limit=512) == smaller_pot(capacity) + + +def test_equal_hardware_yields_an_equal_recipe(): + assert find_batch_size(_fits_up_to(37), limit=512) == find_batch_size(_fits_up_to(39), limit=512) + + +def test_smaller_pot(): + assert [smaller_pot(n) for n in (1, 2, 3, 4, 7, 8, 15, 1293)] == [1, 2, 2, 4, 4, 8, 8, 1024] + with pytest.raises(ValueError, match='n must be >= 1'): + smaller_pot(0) + + +def test_batch_count_covers_the_whole_set_without_overshooting_by_a_batch(): + for sample_count in (1, 7, 8, 900, 5000): + for batch_size in (1, 2, 8, 64, 512): + covered = batch_count(sample_count, batch_size) * batch_size + assert covered >= sample_count + assert covered - sample_count < batch_size + + +def test_the_dataset_limit_keeps_enough_steps_per_epoch(): + for sample_count in (8, 20, 47, 100, 1000, 118_000): + limit = smaller_pot(dataset_limit(sample_count)) + assert sample_count // limit >= MIN_TRAIN_STEPS_PER_EPOCH + + +def test_the_dataset_limit_stays_usable_for_a_tiny_set(): + assert dataset_limit(7) == 1 + with pytest.raises(ValueError, match='sample_count must be >= 1'): + dataset_limit(0) + + +def test_memory_still_decides_below_the_dataset_limit(): + assert find_batch_size(_fits_up_to(2), limit=dataset_limit(20)) == 2 + + +def test_without_a_gpu_the_fallback_respects_the_limit(): + assert no_gpu_batch_size(1024, 'probe') == 8 + assert no_gpu_batch_size(2, 'probe') == 2 + + +@pytest.mark.parametrize('message', [ + 'CUDA out of memory. Tried to allocate 20.00 MiB', + 'cuDNN error: CUDNN_STATUS_ALLOC_FAILED', + 'CUDA error: unknown error', +]) +def test_allocation_failures_are_recognised_however_they_surface(message: str): + assert is_out_of_memory(RuntimeError(message)) + + +def test_a_real_bug_is_not_mistaken_for_a_full_card(): + assert not is_out_of_memory(RuntimeError('shape mismatch in forward pass')) + assert not is_out_of_memory(ValueError('bad config')) + + +def test_the_host_running_out_of_memory_counts_too(): + assert is_out_of_memory(MemoryError()) + + +def _recording(capacity: int) -> tuple[Callable[[int], bool], list[int]]: + """A `fits` predicate for a machine of `capacity`, plus the sizes it gets asked about.""" + calls: list[int] = [] + + def fits(batch_size: int) -> bool: + calls.append(batch_size) + return batch_size <= capacity + + return fits, calls + + +def _fits_up_to(capacity: int) -> Callable[[int], bool]: + fits, _ = _recording(capacity) + return fits diff --git a/learning_loop_node/tests/unit/test_categories.py b/learning_loop_node/tests/unit/test_categories.py new file mode 100644 index 00000000..65492898 --- /dev/null +++ b/learning_loop_node/tests/unit/test_categories.py @@ -0,0 +1,32 @@ +import pytest + +from ...data_classes import Category, ModelInformation +from ...detector.categories import category_by_index, category_by_name +from ...enums import CategoryType + +BOX = Category(id='uuid-box', name='car', type=CategoryType.Box) +POINT = Category(id='uuid-point', name='weed', type=CategoryType.Point) + + +def model_information(*categories: Category) -> ModelInformation: + return ModelInformation(id='model-uuid', host='localhost', organization='zauberzeug', + project='pytest', version='1.2', categories=list(categories or (BOX, POINT))) + + +def test_category_is_resolved_by_index(): + assert category_by_index(model_information(), 1) is POINT + + +@pytest.mark.parametrize('index', [-1, 2, 99]) +def test_an_index_outside_the_model_categories_is_an_error(index: int): + with pytest.raises(ValueError, match='out of range'): + category_by_index(model_information(), index) + + +def test_category_is_resolved_by_name(): + assert category_by_name(model_information(), 'weed') is POINT + + +def test_an_unknown_category_name_lists_the_known_ones(): + with pytest.raises(ValueError, match='car, weed'): + category_by_name(model_information(), 'tractor') diff --git a/learning_loop_node/tests/unit/test_cuda.py b/learning_loop_node/tests/unit/test_cuda.py new file mode 100644 index 00000000..c700f3f2 --- /dev/null +++ b/learning_loop_node/tests/unit/test_cuda.py @@ -0,0 +1,270 @@ +"""Tests for the one library module that imports torch. + +There is no torch in the dev extra and no GPU in CI, so a stand-in is installed under the name +``torch`` before the module is imported. Whether the cap actually holds, and what a real step +costs, can only be observed on a card. +""" +from __future__ import annotations + +import importlib +import logging +import sys +import types +from collections.abc import Callable +from typing import Any + +import pytest + +from ...trainer.batch_size import MAX_BATCH_SIZE, NO_GPU_BATCH_SIZE +from ...trainer.exceptions import InsufficientMemoryError + +MODULE = 'learning_loop_node.trainer.cuda' +GIB = 1024**3 + + +# --- the budget and the cap --- + +def test_no_limit_means_the_whole_card(load): + cuda, _ = load(total_gb=8.0) + assert cuda.usable_memory_bytes(0) == 8 * GIB + assert cuda.usable_memory_bytes(-1) == 8 * GIB + + +def test_a_limit_below_the_card_is_the_budget(load): + cuda, _ = load(total_gb=8.0) + assert cuda.usable_memory_bytes(6) == 6 * GIB + + +def test_a_limit_above_the_card_is_clamped_to_it(load): + cuda, _ = load(total_gb=8.0) + assert cuda.usable_memory_bytes(16) == 8 * GIB + + +def test_the_cap_is_the_limits_share_of_the_card(load): + cuda, fake = load(total_gb=8.0) + cuda.limit_cuda_memory(2) + assert fake.capped == [0.25] + + +def test_no_limit_caps_nothing(load): + cuda, fake = load(total_gb=8.0) + cuda.limit_cuda_memory(0) + cuda.limit_cuda_memory(-1) + assert fake.capped == [] + + +def test_nothing_is_capped_without_cuda(load): + cuda, fake = load(cuda_available=False) + cuda.limit_cuda_memory(2) + assert fake.capped == [] + + +def test_a_limit_the_card_cannot_reach_warns_instead_of_capping(load, caplog): + cuda, fake = load(total_gb=8.0) + with caplog.at_level(logging.WARNING): + cuda.limit_cuda_memory(8) + assert fake.capped == [] + assert 'exceeds the card capacity' in caplog.text + + +def test_the_budget_and_the_cap_follow_the_current_device(load): + cuda, fake = load(total_gb=8.0) + cuda.usable_memory_bytes(2) + cuda.limit_cuda_memory(2) + assert fake.asked_devices and all(device is None for device in fake.asked_devices) + + +def test_freeing_empties_the_cache(load): + cuda, fake = load() + cuda.free_cuda_memory() + assert fake.cache_clears == 1 + + +# --- the safety margin --- + +def test_the_margin_is_a_share_of_the_budget(load): + cuda, fake = load(total_gb=8.0) + cuda.reserve_margin(4, probe='probe') + assert fake.allocated == [int(4 * GIB * cuda.SAFETY_MARGIN)] + + +def test_the_margin_is_a_share_of_the_whole_card_when_nothing_is_budgeted(load): + cuda, fake = load(total_gb=8.0) + cuda.reserve_margin(0, probe='probe') + assert fake.allocated == [int(8 * GIB * cuda.SAFETY_MARGIN)] + + +# --- probing --- + +def test_the_probe_keeps_the_largest_size_that_fits(load): + cuda, fake = load() + ran: list[int] = [] + assert cuda.probe_batch_size(_fits_up_to(16, fake, ran), limit=64) == 16 + assert ran == [1, 2, 4, 8, 16, 32], 'only powers of two, and one trial past the answer' + + +def test_the_probe_reserves_the_margin_before_it_measures(load): + cuda, fake = load(total_gb=8.0) + + def run_batch(_: int) -> None: + assert fake.allocated == [int(8 * GIB * cuda.SAFETY_MARGIN)] + + cuda.probe_batch_size(run_batch, limit=2) + + +def test_the_limit_is_rounded_down_to_a_power_of_two(load): + cuda, fake = load() + ran: list[int] = [] + assert cuda.probe_batch_size(_fits_up_to(1024, fake, ran), limit=12) == 8 + + +def test_an_unset_limit_stops_at_the_maximum(load): + cuda, fake = load() + ran: list[int] = [] + assert cuda.probe_batch_size(_fits_up_to(2048, fake, ran)) == MAX_BATCH_SIZE + + +def test_a_probe_without_a_gpu_does_not_run_the_step(load): + cuda, fake = load(cuda_available=False) + ran: list[int] = [] + assert cuda.probe_batch_size(_fits_up_to(1024, fake, ran), limit=32) == NO_GPU_BATCH_SIZE + assert cuda.probe_batch_size(_fits_up_to(1024, fake, ran), limit=2) == 2 + assert not ran, 'nothing may be run without a GPU' + assert not fake.allocated, 'and no margin claimed on a card that is not there' + + +def test_a_batch_size_of_one_that_does_not_fit_is_an_error(load): + cuda, fake = load() + ran: list[int] = [] + with pytest.raises(InsufficientMemoryError): + cuda.probe_batch_size(_fits_up_to(0, fake, ran), limit=32) + + +def test_exhausted_host_memory_also_means_it_does_not_fit(load): + cuda, _ = load() + + def run_batch(batch_size: int) -> None: + if batch_size > 2: + raise MemoryError + + assert cuda.probe_batch_size(run_batch, limit=32) == 2 + + +def test_a_plain_runtime_error_saying_it_is_out_of_memory_does_not_fit(load): + # cuDNN and cuBLAS workspaces arrive as bare RuntimeErrors + cuda, _ = load() + + def run_batch(batch_size: int) -> None: + if batch_size > 2: + raise RuntimeError('cuDNN error: CUDNN_STATUS_ALLOC_FAILED') + + assert cuda.probe_batch_size(run_batch, limit=32) == 2 + + +def test_torchs_own_error_needs_no_recognisable_message(load): + cuda, fake = load() + + def run_batch(batch_size: int) -> None: + if batch_size > 2: + raise fake.OutOfMemoryError('see the memory summary') + + assert cuda.probe_batch_size(run_batch, limit=32) == 2 + + +def test_a_failure_that_is_not_about_memory_is_a_bug_and_propagates(load): + cuda, _ = load() + + def run_batch(batch_size: int) -> None: + if batch_size > 2: + raise RuntimeError('mat1 and mat2 shapes cannot be multiplied') + + with pytest.raises(RuntimeError, match='shapes cannot be multiplied'): + cuda.probe_batch_size(run_batch, limit=32) + + +def test_what_the_step_reports_is_logged_beside_the_peak(load, caplog): + cuda, _ = load(peak_gb=3.0) + with caplog.at_level(logging.INFO): + cuda.probe_batch_size(lambda _: 'validation peaked at 1.00 GB', probe='miniature epoch', limit=1) + assert 'miniature epoch: 1 fits (peak 3.00 GB, margin included); validation peaked at 1.00 GB' in caplog.text + + +def test_only_a_trial_that_ran_out_of_memory_gets_cleaned_up_after(load): + cuda, fake = load() + cleaned: list[int] = [] + ran: list[int] = [] + fits = cuda.measured_fits(_fits_up_to(2, fake, ran), probe='probe', + on_out_of_memory=lambda: cleaned.append(len(ran))) + + assert [fits(size) for size in (1, 2, 4)] == [True, True, False] + assert cleaned == [3], 'once, after the third trial' + + +# --- the card that is not there --- + +@pytest.fixture(name='load') +def load_fixture(monkeypatch: pytest.MonkeyPatch) -> Callable[..., tuple[Any, _FakeTorch]]: + """Import the module against a fake torch; returns it and the stand-in it ran against.""" + def load(*, cuda_available: bool = True, total_gb: float = 8.0, peak_gb: float = 2.0): + fake = _FakeTorch(cuda_available=cuda_available, total_gb=total_gb, peak_gb=peak_gb) + monkeypatch.setitem(sys.modules, 'torch', fake.module) + monkeypatch.delitem(sys.modules, MODULE, raising=False) + return importlib.import_module(MODULE), fake + + yield load + sys.modules.pop(MODULE, None) + + +def _fits_up_to(largest: int, fake: _FakeTorch, ran: list[int]) -> Callable[[int], None]: + """A step that runs out of memory above ``largest``, recording every size it was asked for.""" + def run_batch(batch_size: int) -> None: + ran.append(batch_size) + if batch_size > largest: + raise fake.OutOfMemoryError('tried to allocate 20.00 GiB') + return run_batch + + +class _FakeTorch: + """What the module uses of torch, plus a record of what it asked for.""" + + def __init__(self, *, cuda_available: bool, total_gb: float, peak_gb: float) -> None: + self.capped: list[float] = [] + self.allocated: list[int] = [] + self.asked_devices: list[int | None] = [] + self.cache_clears = 0 + + class OutOfMemoryError(RuntimeError): + """Torch's own; note the module may not rely on its message.""" + + self.OutOfMemoryError = OutOfMemoryError # named as torch spells it + self.module = types.ModuleType('torch') + self.module.uint8 = 'uint8' # type: ignore[attr-defined] + self.module.empty = self._empty # type: ignore[attr-defined] + self.module.cuda = types.SimpleNamespace( # type: ignore[attr-defined] + OutOfMemoryError=OutOfMemoryError, + is_available=lambda: cuda_available, + empty_cache=self._empty_cache, + get_device_properties=self._device_properties(total_gb), + set_per_process_memory_fraction=self._cap, + reset_peak_memory_stats=lambda: None, + max_memory_allocated=lambda: int(peak_gb * GIB), + synchronize=lambda: None, + ) + + def _empty(self, count: int, dtype: str, device: str) -> object: + assert (dtype, device) == ('uint8', 'cuda') + self.allocated.append(count) + return object() + + def _empty_cache(self) -> None: + self.cache_clears += 1 + + def _device_properties(self, total_gb: float) -> Callable[[int | None], Any]: + def get_device_properties(device: int | None = None): + self.asked_devices.append(device) + return types.SimpleNamespace(total_memory=int(total_gb * GIB)) + return get_device_properties + + def _cap(self, fraction: float, device: int | None = None) -> None: + self.asked_devices.append(device) + self.capped.append(fraction) diff --git a/learning_loop_node/tests/unit/test_entrypoint.py b/learning_loop_node/tests/unit/test_entrypoint.py new file mode 100644 index 00000000..9eed399a --- /dev/null +++ b/learning_loop_node/tests/unit/test_entrypoint.py @@ -0,0 +1,105 @@ +import pytest + +from ...helpers.entrypoint import node_parser + +MANAGED = ('WEIGHT_TYPE', 'MY_DETECTOR_WEIGHT_TYPE', 'HOST', 'NODE_HOST', 'NODE_PORT', 'PORT', + 'MY_DETECTOR_HOST', 'MY_DETECTOR_PORT', 'MY_DETECTOR_NODE_HOST') + + +def test_every_node_gets_a_host_and_a_port(): + args = _parser().parse_args([]) + assert (args.host, args.port) == ('0.0.0.0', 80) + + +def test_a_flag_beats_everything(): + args = _parser().parse_args(['--weight-type', 'FP32']) + assert args.weight_type == 'FP32' + + +def test_a_setting_is_read_from_the_variable_named_after_its_flag(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('WEIGHT_TYPE', 'FP32') + assert _parser().parse_args([]).weight_type == 'FP32' + + +def test_the_loop_own_host_is_never_mistaken_for_the_bind_address(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('HOST', 'preview.learning-loop.ai') + assert _parser().parse_args([]).host == '0.0.0.0' + + +def test_the_bind_address_has_a_name_of_its_own(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('NODE_HOST', '127.0.0.1') + monkeypatch.setenv('NODE_PORT', '8080') + args = _parser().parse_args([]) + assert (args.host, args.port) == ('127.0.0.1', 8080) + + +def test_a_node_that_used_a_prefix_still_reads_it(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('MY_DETECTOR_WEIGHT_TYPE', 'FP32') + parser = _parser(legacy_env_prefix='MY_DETECTOR_') + assert parser.parse_args([]).weight_type == 'FP32' + + +def test_the_prefixed_name_warns_which_one_to_use_instead(monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture): + monkeypatch.setenv('MY_DETECTOR_WEIGHT_TYPE', 'FP32') + _parser(legacy_env_prefix='MY_DETECTOR_').parse_args([]) + assert 'MY_DETECTOR_WEIGHT_TYPE' in caplog.text + assert 'WEIGHT_TYPE' in caplog.text + + +def test_the_current_name_wins_over_the_prefixed_one(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('MY_DETECTOR_WEIGHT_TYPE', 'FP32') + monkeypatch.setenv('WEIGHT_TYPE', 'FP16') + assert _parser(legacy_env_prefix='MY_DETECTOR_').parse_args([]).weight_type == 'FP16' + + +def test_a_node_without_a_legacy_prefix_ignores_prefixed_names(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('MY_DETECTOR_WEIGHT_TYPE', 'FP32') + assert _parser().parse_args([]).weight_type == 'FP16' + + +def test_the_bind_address_is_still_read_under_the_name_the_prefix_gave_the_flag( + monkeypatch: pytest.MonkeyPatch): + """A prefix used to be applied to the flag, so `--host` was `HOST`.""" + monkeypatch.setenv('MY_DETECTOR_HOST', '127.0.0.1') + monkeypatch.setenv('MY_DETECTOR_PORT', '8099') + args = _parser(legacy_env_prefix='MY_DETECTOR_').parse_args([]) + assert (args.host, args.port) == ('127.0.0.1', 8099) + + +def test_the_renamed_setting_warns_which_name_to_use_instead(monkeypatch: pytest.MonkeyPatch, + caplog: pytest.LogCaptureFixture): + monkeypatch.setenv('MY_DETECTOR_HOST', '127.0.0.1') + _parser(legacy_env_prefix='MY_DETECTOR_').parse_args([]) + assert 'MY_DETECTOR_HOST' in caplog.text + assert 'NODE_HOST' in caplog.text + + +def test_the_prefixed_current_name_is_honoured_too(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('MY_DETECTOR_NODE_HOST', '10.0.0.1') + assert _parser(legacy_env_prefix='MY_DETECTOR_').parse_args([]).host == '10.0.0.1' + + +def test_the_current_name_wins_over_both_prefixed_spellings(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('NODE_HOST', '10.0.0.1') + monkeypatch.setenv('MY_DETECTOR_NODE_HOST', '127.0.0.2') + monkeypatch.setenv('MY_DETECTOR_HOST', '127.0.0.3') + assert _parser(legacy_env_prefix='MY_DETECTOR_').parse_args([]).host == '10.0.0.1' + + +def test_the_loop_own_host_is_not_adopted_by_a_prefixed_node(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('HOST', 'preview.learning-loop.ai') + assert _parser(legacy_env_prefix='MY_DETECTOR_').parse_args([]).host == '0.0.0.0' + + +@pytest.fixture(autouse=True) +def clean_env(monkeypatch: pytest.MonkeyPatch): + """Every test starts without the variables it is about to set.""" + for name in MANAGED: + monkeypatch.delenv(name, raising=False) + + +def _parser(**kwargs): + parser = node_parser(description='a node', **kwargs) + parser.add_argument('--weight-type', default='FP16') + return parser diff --git a/learning_loop_node/tests/unit/test_environment_reader.py b/learning_loop_node/tests/unit/test_environment_reader.py new file mode 100644 index 00000000..df957bc0 --- /dev/null +++ b/learning_loop_node/tests/unit/test_environment_reader.py @@ -0,0 +1,51 @@ +"""Tests for `helpers/environment_reader`, which resolves every loop setting a node reads.""" + +import pytest + +from ...helpers import environment_reader + +NAMES = ('LOOP_HOST', 'HOST', 'LOOP_ORGANIZATION', 'ORGANIZATION') + + +@pytest.fixture(autouse=True) +def clean_env(monkeypatch: pytest.MonkeyPatch): + for name in NAMES: + monkeypatch.delenv(name, raising=False) + + +def test_the_prefixed_name_is_preferred(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('LOOP_HOST', 'preview.learning-loop.ai') + monkeypatch.setenv('HOST', 'preview.learning-loop.ai') + assert environment_reader.host() == 'preview.learning-loop.ai' + + +def test_either_name_alone_is_read(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('HOST', 'preview.learning-loop.ai') + assert environment_reader.host() == 'preview.learning-loop.ai' + monkeypatch.delenv('HOST') + monkeypatch.setenv('LOOP_HOST', 'other.learning-loop.ai') + assert environment_reader.host() == 'other.learning-loop.ai' + + +def test_a_disagreement_resolves_to_the_preferred_name(monkeypatch: pytest.MonkeyPatch): + """Returning nothing here would let host() fall back to its default.""" + monkeypatch.setenv('LOOP_HOST', 'preview.learning-loop.ai') + monkeypatch.setenv('HOST', 'learning-loop.ai') + assert environment_reader.host(default='learning-loop.ai') == 'preview.learning-loop.ai' + + +def test_a_disagreement_falls_back_to_the_second_name_when_the_first_is_unset( + monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('ORGANIZATION', 'zauberzeug') + assert environment_reader.organization() == 'zauberzeug' + + +def test_a_disagreement_still_raises_when_errors_are_not_ignored(monkeypatch: pytest.MonkeyPatch): + monkeypatch.setenv('LOOP_HOST', 'a') + monkeypatch.setenv('HOST', 'b') + with pytest.raises(ValueError, match='different environment variables'): + environment_reader.read_from_env(['LOOP_HOST', 'HOST'], ignore_errors=False) + + +def test_nothing_set_yields_the_default(): + assert environment_reader.host(default='fallback') == 'fallback' diff --git a/learning_loop_node/tests/unit/test_geometry.py b/learning_loop_node/tests/unit/test_geometry.py new file mode 100644 index 00000000..76fa2c0b --- /dev/null +++ b/learning_loop_node/tests/unit/test_geometry.py @@ -0,0 +1,24 @@ +from ...detector.geometry import clip_box, clip_point + + +def test_box_inside_the_image_is_unchanged(): + assert clip_box(x1=10, y1=20, width=30, height=40, img_width=100, img_height=100) == (10, 20, 30, 40) + + +def test_box_is_clipped_to_the_image_bounds(): + assert clip_box(x1=-20, y1=-20, width=60, height=60, img_width=100, img_height=100) == (0, 0, 40, 40) + assert clip_box(x1=80, y1=80, width=40, height=40, img_width=100, img_height=100) == (80, 80, 20, 20) + + +def test_box_fully_outside_the_image_collapses_to_zero_size(): + # the corner is clamped at the lower bound only; the zero size is what marks the box empty + assert clip_box(x1=200, y1=200, width=10, height=10, img_width=100, img_height=100) == (200, 200, 0, 0) + + +def test_box_corners_are_rounded(): + assert clip_box(x1=10.4, y1=10.6, width=20.0, height=20.0, img_width=100, img_height=100) == (10, 11, 20, 20) + + +def test_point_is_clamped_into_the_image(): + assert clip_point(50, 50, 100, 100) == (50, 50) + assert clip_point(-10, 150, 100, 100) == (0, 100) diff --git a/learning_loop_node/tests/unit/test_metrics.py b/learning_loop_node/tests/unit/test_metrics.py new file mode 100644 index 00000000..b9a0fe6b --- /dev/null +++ b/learning_loop_node/tests/unit/test_metrics.py @@ -0,0 +1,38 @@ +import pytest + +from ...trainer.metrics import category_f1, confusion_matrix_from_counts, macro_f1 + + +def test_the_score_averages_the_categories_instead_of_pooling_them(): + """Counts from a real epoch: face 248/19/59 and hand 39/54/73 - 74% pooled, 62% averaged.""" + assert macro_f1({'face': {'tp': 248, 'fp': 19, 'fn': 59}, + 'hand': {'tp': 39, 'fp': 54, 'fn': 73}}) == pytest.approx(0.622, abs=0.001) + assert macro_f1({'both': {'tp': 287, 'fp': 73, 'fn': 132}}) == pytest.approx(0.737, abs=0.001) + + +def test_a_rare_category_carries_the_same_weight(): + assert macro_f1({'frequent': {'tp': 1000, 'fp': 0, 'fn': 0}, + 'rare': {'tp': 0, 'fp': 0, 'fn': 3}}) == pytest.approx(0.5) + + +def test_a_category_without_any_counts_scores_zero_instead_of_raising(): + assert macro_f1({'unseen': {'tp': 0, 'fp': 0, 'fn': 0}}) == 0.0 + assert category_f1({'tp': 0, 'fp': 0, 'fn': 0}) == 0.0 + + +def test_an_empty_confusion_matrix_scores_zero(): + assert macro_f1({}) == 0.0 + + +def test_a_perfect_category_scores_one(): + assert category_f1({'tp': 10, 'fp': 0, 'fn': 0}) == pytest.approx(1.0) + + +def test_precision_and_recall_are_balanced(): + assert category_f1({'tp': 5, 'fp': 5, 'fn': 0}) == category_f1({'tp': 5, 'fp': 0, 'fn': 5}) + + +def test_counts_can_be_assembled_from_separate_mappings(): + matrix = confusion_matrix_from_counts({'a': 3}, {'a': 1, 'b': 2}, {'b': 4}) + assert matrix == {'a': {'tp': 3, 'fp': 1, 'fn': 0}, + 'b': {'tp': 0, 'fp': 2, 'fn': 4}} diff --git a/learning_loop_node/tests/unit/test_postprocess.py b/learning_loop_node/tests/unit/test_postprocess.py new file mode 100644 index 00000000..93cb6b48 --- /dev/null +++ b/learning_loop_node/tests/unit/test_postprocess.py @@ -0,0 +1,148 @@ +import numpy as np +import pytest + +from ...data_classes import Category, ModelInformation +from ...detector.postprocess import ( + Prediction, + bbox_iou, + non_max_suppression, + post_process, + predictions_from_xyxy, + to_detections, + to_image_metadata, +) +from ...enums import CategoryType + +BOX = Category(id='uuid-box', name='car', type=CategoryType.Box) +POINT = Category(id='uuid-point', name='weed', type=CategoryType.Point) + + +# --- iou and suppression --- + +def test_identical_boxes_have_an_iou_of_one(): + box = np.array([[0, 0, 10, 10]], dtype=np.float32) + assert bbox_iou(box, box)[0] == pytest.approx(1.0) + + +def test_disjoint_boxes_have_an_iou_of_zero(): + a = np.array([[0, 0, 10, 10]], dtype=np.float32) + b = np.array([[100, 100, 110, 110]], dtype=np.float32) + assert bbox_iou(a, b)[0] == pytest.approx(0.0) + + +def test_overlapping_boxes_of_one_class_are_suppressed_keeping_the_best(): + boxes = np.array([[10, 10, 60, 60], [12, 12, 62, 62]], dtype=np.float32) + scores = np.array([0.9, 0.8], dtype=np.float32) + kept_boxes, kept_scores, _ = non_max_suppression( + boxes, scores, np.array([0, 0]), iou_threshold=0.45, origin_h=200, origin_w=200) + assert len(kept_boxes) == 1 + assert kept_scores[0] == pytest.approx(0.9) + + +def test_overlapping_boxes_of_different_classes_both_survive(): + boxes = np.array([[10, 10, 60, 60], [12, 12, 62, 62]], dtype=np.float32) + kept_boxes, _, _ = non_max_suppression( + boxes, np.array([0.9, 0.8], dtype=np.float32), np.array([0, 1]), + iou_threshold=0.45, origin_h=200, origin_w=200) + assert len(kept_boxes) == 2 + + +def test_suppression_clips_boxes_into_the_image(): + boxes = np.array([[-10, -10, 300, 300]], dtype=np.float32) + kept_boxes, _, _ = non_max_suppression( + boxes, np.array([0.9], dtype=np.float32), np.array([0]), + iou_threshold=0.45, origin_h=100, origin_w=100) + assert list(kept_boxes[0]) == [0, 0, 99, 99] + + +# --- post_process --- + +def test_post_process_drops_predictions_below_the_confidence_threshold(): + boxes = np.array([[10, 10, 60, 60], [100, 100, 150, 150]], dtype=np.float32) + result = post_process(boxes, np.array([0.9, 0.1], dtype=np.float32), np.array([0, 0]), + conf_threshold=0.5, iou_threshold=0.45, origin_h=200, origin_w=200) + assert len(result) == 1 + prediction = result[0] + assert (prediction.x, prediction.y, prediction.width, prediction.height) == (10, 10, 50, 50) + assert prediction.category_index == 0 + # the model's own float32 score, not rounded + assert prediction.confidence == pytest.approx(0.9) + + +def test_post_process_on_an_empty_prediction_returns_nothing(): + empty_boxes = np.zeros((0, 4), dtype=np.float32) + assert post_process(empty_boxes, np.zeros(0, dtype=np.float32), np.zeros(0, dtype=int), + conf_threshold=0.5, iou_threshold=0.45, origin_h=10, origin_w=10) == [] + + +def test_already_suppressed_output_keeps_the_model_s_own_precision(): + assert predictions_from_xyxy(labels=[1.0], boxes=[[10.4, 10.6, 60.4, 60.6]], scores=[0.55]) == \ + [Prediction(x=10.4, y=10.6, width=50.0, height=50.0, category_index=1, confidence=0.55)] + + +def test_converting_already_suppressed_output_requires_matching_lengths(): + with pytest.raises(ValueError): + predictions_from_xyxy(labels=[1.0, 2.0], boxes=[[0.0, 0.0, 1.0, 1.0]], scores=[0.5]) + + +# --- building the containers --- + +def test_a_box_category_becomes_a_box_detection(): + metadata = to_image_metadata([Prediction(x=10, y=20, width=30, height=40, category_index=0, confidence=0.9)], model_information(), 200, 200) + assert len(metadata.point_detections) == 0 + detection = metadata.box_detections[0] + assert (detection.x, detection.y, detection.width, detection.height) == (10, 20, 30, 40) + assert (detection.category_name, detection.category_id) == ('car', 'uuid-box') + assert detection.model_name == '1.2' + assert detection.confidence == pytest.approx(0.9) + + +def test_a_point_category_becomes_the_centre_of_the_box(): + metadata = to_image_metadata([Prediction(x=100, y=100, width=40, height=40, category_index=1, confidence=0.7)], model_information(), 200, 200) + assert len(metadata.box_detections) == 0 + detection = metadata.point_detections[0] + assert (detection.x, detection.y) == (120, 120) + assert detection.category_id == 'uuid-point' + + +def test_detections_are_clipped_to_the_image(): + metadata = to_image_metadata([Prediction(x=-20, y=-20, width=60, height=60, category_index=0, confidence=0.5)], model_information(), 200, 200) + detection = metadata.box_detections[0] + assert (detection.x, detection.y, detection.width, detection.height) == (0, 0, 40, 40) + + +@pytest.mark.parametrize('width,height', [(2, 30), (30, 2), (1, 1)]) +def test_boxes_too_small_to_be_useful_are_dropped(width: int, height: int): + metadata = to_image_metadata([Prediction(x=5, y=5, width=width, height=height, category_index=0, confidence=0.5)], model_information(), 200, 200) + assert len(metadata) == 0 + + +def test_a_category_type_the_node_cannot_report_is_skipped(): + classification = Category(id='uuid-cls', name='ripe', type=CategoryType.Classification) + metadata = to_image_metadata([Prediction(x=10, y=10, width=30, height=30, category_index=0, confidence=0.5)], + model_information(classification), 200, 200) + assert len(metadata) == 0 + + +def test_the_trainer_container_carries_the_image_id(): + result = to_detections([Prediction(x=10, y=20, width=30, height=40, category_index=0, confidence=0.9)], model_information(), 200, 200, + image_id='image-uuid') + assert result.image_id == 'image-uuid' + assert len(result.box_detections) == 1 + + +def test_trainer_and_detector_paths_agree_on_the_same_detections(): + predictions = [Prediction(x=-5, y=-5, width=60, height=60, category_index=0, confidence=0.9), Prediction(x=100, y=100, width=40, height=40, category_index=1, confidence=0.7), + Prediction(x=5, y=5, width=1, height=1, category_index=0, confidence=0.5)] + metadata = to_image_metadata(predictions, model_information(), 200, 200) + result = to_detections(predictions, model_information(), 200, 200, image_id='image-uuid') + + assert [(d.x, d.y, d.width, d.height, d.category_id) for d in result.box_detections] == \ + [(d.x, d.y, d.width, d.height, d.category_id) for d in metadata.box_detections] + assert [(d.x, d.y, d.category_id) for d in result.point_detections] == \ + [(d.x, d.y, d.category_id) for d in metadata.point_detections] + + +def model_information(*categories: Category) -> ModelInformation: + return ModelInformation(id='model-uuid', host='localhost', organization='zauberzeug', + project='pytest', version='1.2', categories=list(categories or (BOX, POINT))) diff --git a/learning_loop_node/tests/unit/test_subprocess.py b/learning_loop_node/tests/unit/test_subprocess.py new file mode 100644 index 00000000..b197cc9b --- /dev/null +++ b/learning_loop_node/tests/unit/test_subprocess.py @@ -0,0 +1,53 @@ +import pytest + +from ...trainer.subprocess import iterator_cpu_bound + + +def counting(limit: int): + yield from range(limit) + + +def raising_after(limit: int): + yield from range(limit) + raise RuntimeError('the training crashed') + + +def empty(): + return iter(()) + + +async def _collect(iterator) -> list: + return [item async for item in iterator] + + +async def test_everything_the_generator_yields_arrives_in_order(): + async with iterator_cpu_bound(counting, 5) as iterator: + assert await _collect(iterator) == [0, 1, 2, 3, 4] + + +async def test_a_generator_that_yields_nothing_simply_finishes(): + async with iterator_cpu_bound(empty) as iterator: + assert await _collect(iterator) == [] + + +async def test_a_failure_in_the_process_is_raised_in_the_caller(): + with pytest.raises(RuntimeError, match='the training crashed'): + async with iterator_cpu_bound(raising_after, 2) as iterator: + await _collect(iterator) + + +async def test_what_the_generator_produced_before_failing_still_arrives(): + received = [] + with pytest.raises(RuntimeError): + async with iterator_cpu_bound(raising_after, 3) as iterator: + async for item in iterator: + received.append(item) + assert received == [0, 1, 2] + + +async def test_leaving_early_does_not_leave_the_process_running(): + async with iterator_cpu_bound(counting, 1000) as iterator: + async for item in iterator: + if item == 2: + break + # reaching here without hanging is the test diff --git a/learning_loop_node/trainer/batch_size.py b/learning_loop_node/trainer/batch_size.py new file mode 100644 index 00000000..8fb0418c --- /dev/null +++ b/learning_loop_node/trainer/batch_size.py @@ -0,0 +1,81 @@ +"""Choosing a batch size by probing, rather than configuring one. + +The trainer supplies a ``fits`` predicate that runs a representative step at doubling sizes. +Only powers of two are visited, so equal hardware yields an equal recipe. Nothing here imports a +deep-learning framework. + +Adapted from PyTorch Lightning's ``BatchSizeFinder`` (power-scaling mode). +Copyright The Lightning AI team. Licensed under the Apache License, Version 2.0. +https://github.com/Lightning-AI/pytorch-lightning +""" + +import logging +from collections.abc import Callable + +from .exceptions import InsufficientMemoryError + +logger = logging.getLogger(__name__) + +MAX_BATCH_SIZE = 1024 +"""Where a search stops when its caller sets no bound of its own.""" + +NO_GPU_BATCH_SIZE = 8 +"""Batch size used when there is no GPU to probe.""" + +MIN_TRAIN_STEPS_PER_EPOCH = 8 +"""Fewest optimizer steps an epoch must have; :func:`dataset_limit` is derived from it.""" + + +def find_batch_size(fits: Callable[[int], bool], *, limit: int) -> int: + """Return the largest power-of-two batch size that fits, never exceeding ``limit``. + + :param fits: Runs a representative probe; ``False`` on out-of-memory. + :raises InsufficientMemoryError: If not even a batch size of 1 fits. + """ + limit = smaller_pot(limit) + if not fits(1): + raise InsufficientMemoryError('batch size 1 does not fit in memory') + + size = 1 + while size < limit and fits(size * 2): + size *= 2 + + return size + + +def dataset_limit(sample_count: int) -> int: + """The batch size ceiling that still leaves ``MIN_TRAIN_STEPS_PER_EPOCH`` steps per epoch.""" + if sample_count < 1: + raise ValueError(f'sample_count must be >= 1, got {sample_count}') + return max(1, sample_count // MIN_TRAIN_STEPS_PER_EPOCH) + + +def no_gpu_batch_size(limit: int, probe: str) -> int: + """The batch size to fall back on when there is no GPU to probe.""" + batch_size = min(smaller_pot(limit), NO_GPU_BATCH_SIZE) + logger.warning('%s: CUDA is unavailable; using batch size %d without probing', probe, batch_size) + return batch_size + + +def smaller_pot(n: int) -> int: + """The largest power of two that is <= ``n``.""" + if n < 1: + raise ValueError(f'n must be >= 1, got {n}') + return 1 << (n.bit_length() - 1) + + +def batch_count(sample_count: int, batch_size: int) -> int: + """How many batches a loader yields, the last one possibly short.""" + return -(-sample_count // batch_size) + + +def is_out_of_memory(exception: BaseException) -> bool: + """Whether the exception signals exhausted memory, on the GPU or the host. + + cuDNN and cuBLAS workspace failures raise a plain ``RuntimeError``, so the message has to be + matched too. + """ + if isinstance(exception, MemoryError): + return True + message = str(exception).lower() + return any(text in message for text in ('out of memory', 'alloc_failed', 'cuda error: unknown error')) diff --git a/learning_loop_node/trainer/cuda.py b/learning_loop_node/trainer/cuda.py new file mode 100644 index 00000000..1b126bb0 --- /dev/null +++ b/learning_loop_node/trainer/cuda.py @@ -0,0 +1,137 @@ +"""GPU memory budgeting and batch-size probing for trainers that train in-process. + +``--vram-limit-gb`` -- gigabytes of the card this training may use -- becomes both the budget a +probe measures against and the cap that holds the process to it, on whichever GPU the calling +process is using. The cap is a share of the card's *total* memory, not of what is free. + +This is the one module in the library that imports torch, which the package does not declare, so +only a trainer may import it. The search itself is in :mod:`~learning_loop_node.trainer.batch_size`. +""" +from __future__ import annotations + +import gc +import logging +from collections.abc import Callable + +import torch + +from .batch_size import MAX_BATCH_SIZE, find_batch_size, is_out_of_memory, no_gpu_batch_size, smaller_pot + +logger = logging.getLogger(__name__) + +SAFETY_MARGIN = 0.05 +"""Share of the budget held back while probing, against allocator fragmentation later on.""" + + +def probe_batch_size(run_batch: Callable[[int], str | None], *, probe: str = 'batch-size probe', + limit: int = 0, vram_limit_gb: float = 0) -> int: + """Find the largest power-of-two batch size ``run_batch`` fits into. + + For a probe whose measurement is one call. A probe that has to build a throwaway model first + composes :func:`reserve_margin`, :func:`measured_fits` and ``find_batch_size`` itself. + + :param run_batch: Runs the batch; may return a detail to append to the log line. + :param probe: Names this probe in the log, so a node running several stays readable. + :param limit: Caps the search, rounded down to a power of two; 0 means + :data:`~learning_loop_node.trainer.batch_size.MAX_BATCH_SIZE`. + :param vram_limit_gb: The budget the safety margin is a share of; 0 means the whole card. + :raises InsufficientMemoryError: If not even a batch size of 1 fits. + """ + limit = smaller_pot(limit or MAX_BATCH_SIZE) + + if not torch.cuda.is_available(): + return no_gpu_batch_size(limit, probe) + + margin = reserve_margin(vram_limit_gb, probe=probe) + try: + chosen = find_batch_size(measured_fits(run_batch, probe=probe), limit=limit) + finally: + del margin + free_cuda_memory() + logger.info('%s: selected batch size %d (upper bound %d)', probe, chosen, limit) + return chosen + + +def measured_fits(run_batch: Callable[[int], str | None], *, probe: str, + on_out_of_memory: Callable[[], None] | None = None) -> Callable[[int], bool]: + """Wrap ``run_batch`` into the ``fits`` predicate ``find_batch_size`` searches with. + + An out-of-memory failure is the answer "does not fit"; anything else is re-raised. Both + arrive as the same exception types, so they are told apart by + :func:`~learning_loop_node.trainer.batch_size.is_out_of_memory`, not by ``except``. + + :param run_batch: Runs the batch; may return a detail to append to the log line. + :param on_out_of_memory: Runs after a trial ran out of memory, to drop what it left behind + (an optimizer's gradients, say). + """ + def fits(batch_size: int) -> bool: + free_cuda_memory() + try: + torch.cuda.reset_peak_memory_stats() + detail = run_batch(batch_size) + torch.cuda.synchronize() + logger.info('%s: %d fits (peak %.2f GB, margin included)%s', probe, batch_size, + torch.cuda.max_memory_allocated() / 1024**3, f'; {detail}' if detail else '') + return True + except (torch.cuda.OutOfMemoryError, RuntimeError, MemoryError) as exc: + if not isinstance(exc, torch.cuda.OutOfMemoryError) and not is_out_of_memory(exc): + raise + logger.info('%s: %d does not fit (%s)', probe, batch_size, type(exc).__name__) + if on_out_of_memory is not None: + on_out_of_memory() + return False + + return fits + + +def reserve_margin(vram_limit_gb: float, *, probe: str) -> torch.Tensor: + """Claim :data:`SAFETY_MARGIN` of the budget, so a trial competes against a smaller card. + + Keep the returned tensor alive for as long as the probe runs: releasing it hands the margin + back, and the chosen size is no longer the size that was measured. + + :param vram_limit_gb: The budget the margin is a share of; 0 means the whole card. + """ + margin_bytes = int(usable_memory_bytes(vram_limit_gb) * SAFETY_MARGIN) + logger.info('%s: keeping %.0f MB free as a safety margin', probe, margin_bytes / 1024**2) + return torch.empty(margin_bytes, dtype=torch.uint8, device='cuda') + + +def usable_memory_bytes(vram_limit_gb: float) -> int: + """How much GPU memory this process may allocate, honouring :func:`limit_cuda_memory`. + + :param vram_limit_gb: 0 or less means the whole card. + """ + total_bytes = torch.cuda.get_device_properties(None).total_memory + if vram_limit_gb <= 0: + return total_bytes + return min(total_bytes, int(vram_limit_gb * 1024**3)) + + +def limit_cuda_memory(vram_limit_gb: float) -> None: + """Cap how much of the GPU this process may allocate, to ``vram_limit_gb`` gigabytes. + + Call this once per process that touches the GPU, a spawned training process included: the + cap does not survive the spawn. + + :param vram_limit_gb: 0 or less means no cap. + """ + if vram_limit_gb <= 0 or not torch.cuda.is_available(): + return + + total_bytes = torch.cuda.get_device_properties(None).total_memory + fraction = vram_limit_gb * 1024**3 / total_bytes + total_gb = total_bytes / 1024**3 + + if fraction >= 1.0: + logger.warning('VRAM limit of %.1f GB exceeds the card capacity of %.1f GB; not limiting', + vram_limit_gb, total_gb) + return + + torch.cuda.set_per_process_memory_fraction(fraction, None) + logger.info('Limiting VRAM usage to %.1f GB of %.1f GB (%.0f%%)', vram_limit_gb, total_gb, fraction * 100) + + +def free_cuda_memory() -> None: + gc.collect() + torch.cuda.empty_cache() diff --git a/learning_loop_node/trainer/exceptions.py b/learning_loop_node/trainer/exceptions.py index c3cb54f2..b77548cd 100644 --- a/learning_loop_node/trainer/exceptions.py +++ b/learning_loop_node/trainer/exceptions.py @@ -1,12 +1,16 @@ class CriticalError(Exception): - ''' + """ CriticalError is raised when the training cannot be continued. In this case the trainer jumps to the TrainerState.ReadyForCleanup and tries to upload the latest model. - ''' + """ + + +class InsufficientMemoryError(RuntimeError): + """Raised when not even the smallest unit of work fits in memory.""" class NodeNeedsRestartError(Exception): - ''' + """ NodeNeedsRestartError is raised when the node needs to be restarted. This is e.g. the case when the GPU is not available anymore. - ''' + """ diff --git a/learning_loop_node/trainer/metrics.py b/learning_loop_node/trainer/metrics.py new file mode 100644 index 00000000..0ce6c43d --- /dev/null +++ b/learning_loop_node/trainer/metrics.py @@ -0,0 +1,28 @@ +"""Scoring a training from the confusion matrix the loop stores: one +``{'tp': .., 'fp': .., 'fn': ..}`` per category.""" + +import statistics +from collections.abc import Mapping + + +def macro_f1(confusion_matrix: Mapping[str, Mapping[str, int]]) -> float: + """The unweighted mean of the per-category F1 scores, as the loop's UI shows it by default.""" + scores = [category_f1(counts) for counts in confusion_matrix.values()] + return statistics.mean(scores) if scores else 0.0 + + +def category_f1(counts: Mapping[str, int]) -> float: + """The F1 score of one category, 0.0 where it is undefined rather than a division error.""" + tp, fp, fn = counts['tp'], counts['fp'], counts['fn'] + precision = tp / (tp + fp) if tp + fp else 0.0 + recall = tp / (tp + fn) if tp + fn else 0.0 + return 2 * precision * recall / (precision + recall) if precision + recall else 0.0 + + +def confusion_matrix_from_counts(true_positives: Mapping[str, int], false_positives: Mapping[str, int], + false_negatives: Mapping[str, int]) -> dict[str, dict[str, int]]: + """Assemble the loop's confusion matrix from three per-category count mappings.""" + return {category: {'tp': true_positives.get(category, 0), + 'fp': false_positives.get(category, 0), + 'fn': false_negatives.get(category, 0)} + for category in {*true_positives, *false_positives, *false_negatives}} diff --git a/learning_loop_node/trainer/subprocess.py b/learning_loop_node/trainer/subprocess.py new file mode 100644 index 00000000..118fbf5d --- /dev/null +++ b/learning_loop_node/trainer/subprocess.py @@ -0,0 +1,91 @@ +"""Run a blocking, CPU-bound generator in its own process without blocking the event loop. + +The queue has ``maxsize=1``, so the producer never runs more than one item ahead of the consumer. +The context is spawn on every platform, so ``it`` and its arguments must be picklable and a +process that already initialised CUDA is never forked. An exception raised inside the process is +re-raised in the caller, and the process is killed if the caller leaves the context early. +""" +from __future__ import annotations + +import asyncio +import logging +import multiprocessing +import queue +from collections.abc import AsyncGenerator, Callable, Iterator +from contextlib import asynccontextmanager +from multiprocessing.queues import Queue as MPQueue +from typing import Any, ParamSpec, TypeVar + +logger = logging.getLogger(__name__) + +T = TypeVar('T') +P = ParamSpec('P') + + +@asynccontextmanager +async def iterator_cpu_bound( + it: Callable[P, Iterator[T]], + *args: P.args, + **kwargs: P.kwargs, +) -> AsyncGenerator[AsyncGenerator[T, None], None]: + iterator = _iterator_cpu_bound_inner(it, *args, **kwargs) + try: + yield iterator + finally: + await asyncio.shield(iterator.aclose()) + + +async def _iterator_cpu_bound_inner( + it: Callable[P, Iterator[T]], + *args: P.args, + **kwargs: P.kwargs, +) -> AsyncGenerator[T, None]: + ctx = multiprocessing.get_context('spawn') + state_queue: MPQueue[T | Exception | IteratorDone] = ctx.Queue(maxsize=1) + process = ctx.Process( + target=_iterator_wrapper, + args=(it, state_queue, args, kwargs), + name='iterator_cpu_bound', + ) + + process.start() + + try: + while True: + try: + item = await asyncio.to_thread(state_queue.get, True, 0.5) + except queue.Empty: + if not process.is_alive(): + break + continue + match item: + case IteratorDone(): + break + case Exception() as e: + raise e + case _ as other: + yield other + finally: + if process.is_alive(): + process.kill() + process.join() + + +def _iterator_wrapper( + it: Callable[..., Iterator[T]], + state_queue: MPQueue[T | Exception | IteratorDone], + args: tuple[Any, ...], + kwargs: dict[str, Any], +) -> None: + try: + for data in it(*args, **kwargs): + state_queue.put(data) + except Exception as e: + logger.exception('iterator_cpu_bound child process failed') + state_queue.put(e) + + state_queue.put(IteratorDone()) + + +class IteratorDone: + pass diff --git a/learning_loop_node/trainer/trainer_logic.py b/learning_loop_node/trainer/trainer_logic.py index c5a9bade..b786113f 100644 --- a/learning_loop_node/trainer/trainer_logic.py +++ b/learning_loop_node/trainer/trainer_logic.py @@ -161,27 +161,38 @@ async def stop(self) -> None: @abstractmethod async def _start_training_from_base_model(self) -> None: - '''Should be used to start a training on executer, e.g. self.executor.start(cmd).''' + """Should be used to start a training on executer, e.g. self.executor.start(cmd).""" @abstractmethod async def _start_training_from_scratch(self) -> None: - '''Should be used to start a training from scratch on executer, e.g. self.executor.start(cmd). + """Should be used to start a training from scratch on executer, e.g. self.executor.start(cmd). NOTE base_model_id is now accessible via self.training.base_model_id - the id of a pretrained model provided by self.provided_pretrained_models.''' + the id of a pretrained model provided by self.provided_pretrained_models.""" @abstractmethod def _can_resume(self) -> bool: - '''Override this method to return True if the trainer can resume training.''' + """Override this method to return True if the trainer can resume training.""" @abstractmethod async def _resume(self) -> None: - '''Is called when self.can_resume() returns True. - One may resume the training on a previously trained model stored by self.on_model_published(basic_model).''' + """Is called when self.can_resume() returns True. + One may resume the training on a previously trained model stored by self.on_model_published(basic_model).""" - @abstractmethod def _get_executor_error_from_log(self) -> Optional[str]: - '''Should be used to provide error informations to the Learning Loop by extracting data from self.executor.get_log().''' + """Reports what went wrong to the Learning Loop by reading self.executor's log. + + The default recognises the CUDA failures every trainer hits. Override to add messages a + particular training framework produces, and call super() to keep these. + """ + if self._executor is None: + return None + for line in self._executor.get_log_by_lines(tail=50): + if 'CUDA out of memory' in line: + return 'graphics card is out of memory' + if 'CUDA error: invalid device ordinal' in line: + return 'graphics card not found' + return None @abstractmethod async def _detect(self, model_information: ModelInformation, images: List[str], model_folder: str) -> List[Detections]: - '''Called to run detections on a list of images.''' + """Called to run detections on a list of images.""" diff --git a/pyproject.toml b/pyproject.toml index 47409dff..d05073f5 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -13,6 +13,7 @@ dependencies = [ "python-socketio>=5.16.2,<6.0.0", "aiofiles>=0.7.0", "python-multipart>=0.0.31", + "configargparse>=1.7.1", "psutil>=5.9.0,<8.0.0", "numpy>=2.0,<3.0", "Pillow>=12.3.0,<13.0.0", diff --git a/run_tests.sh b/run_tests.sh index e83d3852..e8ae7b66 100755 --- a/run_tests.sh +++ b/run_tests.sh @@ -6,6 +6,7 @@ set -o allexport; source .env; set +o allexport # Check if argument is provided if [ $# -eq 1 ]; then # Run tests with filter + python -m pytest learning_loop_node/tests/unit -v -s -k "$1" python -m pytest learning_loop_node/tests/annotator -v -s -k "$1" python -m pytest learning_loop_node/tests/detector -v -s -k "$1" python -m pytest learning_loop_node/tests/trainer -v -s -k "$1" @@ -17,6 +18,8 @@ fi # Run the tests +# unit runs first: it is the only suite that needs no Learning Loop +python -m pytest learning_loop_node/tests/unit -v python -m pytest learning_loop_node/tests/annotator -v python -m pytest learning_loop_node/tests/detector -v python -m pytest learning_loop_node/tests/trainer -v diff --git a/uv.lock b/uv.lock index da0c0bf7..1b373e19 100644 --- a/uv.lock +++ b/uv.lock @@ -449,6 +449,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/d1/d6/3965ed04c63042e047cb6a3e6ed1a63a35087b6a609aa3a15ed8ac56c221/colorama-0.4.6-py2.py3-none-any.whl", hash = "sha256:4f1d9991f5acc0ca119f9d443620b77f9d6b33703e51011c16baf57afb285fc6", size = 25335, upload-time = "2022-10-25T02:36:20.889Z" }, ] +[[package]] +name = "configargparse" +version = "1.7.5" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/3f/0b/30328302903c55218ffc5199646d0e9d28348ff26c02ba77b2ffc58d294a/configargparse-1.7.5.tar.gz", hash = "sha256:e3f9a7bb6be34d66b2e3c4a2f58e3045f8dfae47b0dc039f87bcfaa0f193fb0f", size = 53548, upload-time = "2026-03-11T02:19:38.144Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/fe/19/3ba5e1b0bcc7b91aeab6c258afd70e4907d220fed3972febe38feb40db30/configargparse-1.7.5-py3-none-any.whl", hash = "sha256:1e63fdffedf94da9cd435fc13a1cd24777e76879dd2343912c1f871d4ac8c592", size = 27692, upload-time = "2026-03-11T02:19:36.442Z" }, +] + [[package]] name = "dacite" version = "1.9.2" @@ -781,6 +790,7 @@ source = { virtual = "." } dependencies = [ { name = "aiofiles" }, { name = "aiohttp" }, + { name = "configargparse" }, { name = "dacite" }, { name = "fastapi" }, { name = "httpx" }, @@ -813,6 +823,7 @@ requires-dist = [ { name = "aiofiles", specifier = ">=0.7.0" }, { name = "aiohttp", specifier = ">=3.14.3,<4.0.0" }, { name = "autopep8", marker = "extra == 'dev'", specifier = ">=2.0.2,<3.0.0" }, + { name = "configargparse", specifier = ">=1.7.1" }, { name = "dacite", specifier = ">=1.8.1,<2.0.0" }, { name = "debugpy", marker = "extra == 'dev'", specifier = ">=1.6.7.post1,<2.0.0" }, { name = "fastapi", specifier = ">=0.135,<1.0" },