diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 4bba35c..08b8064 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -147,6 +147,23 @@ layouts). `vp sync` reconciles the index idempotently (content-addressed, so re-running skips what's present) and records per-asset provenance; `vp sync --dry-run` previews found/matched/unmatched/classes per source. +**Multi-provider targets.** With a `target:` and `copy` mode, objects land in a +content-addressed sink. Same-provider transfers are server-side (S3 CopyObject / +GCS rewrite — bytes never transit the client). Cross-provider transfers +(local→S3, S3→GCS, …) *relay* the bytes sync already read to compute the sha256: +one read (needed anyway) + one upload (size-verified against the landed +object), never a second download. Either way an unchanged re-sync stays +metadata-only. + +**Request economy & resilience.** Target membership is resolved with one prefix +listing per run (rclone-style fast-list), never a per-object existence check. +Every remote call is wrapped in exponential-backoff retries +(`sources/retry.py`); exhaustion surfaces as a per-object `IngestFailure`, so a +flaky bucket degrades a sync instead of aborting it. Ingest concurrency is +sized for latency-bound object stores (16+ workers for remote sources) and +tunable with `vp sync --jobs`. Labels fetched during class inference are cached +and replayed at ingest, so no remote label is ever read twice. + ### Geometry model & task coverage (`core/models.py`, `formats/`) The tagged geometry above lets one dataset model cover the common CV tasks. Classification uses the ImageFolder convention; detection uses YOLO/COCO bbox; @@ -223,11 +240,17 @@ The roadmap is sequenced so each phase unblocks the next. - [x] `vp pack --profile training` (WebDataset shards) - [ ] DuckDB index — deferred by decision (see Storage & index) -### Multi-source ingestion ✅ (remote backends pending) +### Multi-source ingestion ✅ - [x] declarative `sources:` + `vp sync` (+ `--dry-run`), local backends, joins, provenance, class reconciliation, idempotent re-sync -- [ ] remote backends via fsspec behind extras (s3/gcs/azure/git, pinned by ref) - and COCO-format sources, plugging into the same resolver layer +- [x] remote YOLO/COCO/ImageFolder sources via fsspec behind extras + (s3/gcs/azure), metadata-only re-sync, content-addressed cloud target +- [x] cross-provider `copy` targets (local↔S3, S3↔GCS, …): server-side copy when + providers match, single-pass verified byte relay when they don't +- [x] hardening: retries with backoff on every remote call, fast-list target + membership (no per-object HEAD), `--jobs` concurrency control, + single-read remote labels +- [ ] git sources pinned by ref ### Task coverage ✅ (beyond detection) - [x] tagged geometry model (bbox | polygon | keypoints | none), backward compatible @@ -259,8 +282,10 @@ The roadmap is sequenced so each phase unblocks the next. - [ ] Hugging Face Datasets export ### Phase C — Reporting & polish +- [x] machine-readable output: every pipeline command takes `--json` and prints + a schema-versioned envelope (`cli/output.py`) — the stable contract for + driving VisionPack from services/UIs/CI; errors are structured too - [ ] HTML validation / stats / drift reports -- [ ] JSON report output for stats and diff - [ ] richer terminal output with `rich` - [ ] move CLI plumbing from `argparse` to `typer` once commands stabilize diff --git a/README.md b/README.md index d169059..29ac90f 100644 --- a/README.md +++ b/README.md @@ -9,7 +9,7 @@ leak-free, ready-to-train dataset. ![Python](https://img.shields.io/badge/python-3.11%2B-blue) ![License](https://img.shields.io/badge/license-Apache--2.0-blue) ![Status](https://img.shields.io/badge/status-active%20development-orange) -![Tests](https://img.shields.io/badge/tests-118%20passing-brightgreen) +![Tests](https://img.shields.io/badge/tests-143%20passing-brightgreen) [Documentation](https://caiowing.github.io/VisionPack/) · [Install](https://caiowing.github.io/VisionPack/installation/) · @@ -176,8 +176,12 @@ See the [Cloud Sync guide](https://caiowing.github.io/VisionPack/cloud-sync/). extra dependencies, scale-proof via LSH bucketing; surfaced in `vp validate`. - **Multi-source sync** — declarative `sources:` + `vp sync`, with per-asset provenance; idempotent re-sync that only pulls what's new. -- **Cloud-native** — sync from and to S3/GCS/Azure without downloading the whole - dataset; server-side `copy` into a content-addressed target, streaming export. +- **Cloud-native, multi-provider** — sync YOLO, COCO, and ImageFolder sources + from S3/GCS/Azure without downloading the whole dataset; server-side `copy` + into a content-addressed target when source and target share a provider, + single-pass verified relay across providers (S3→GCS, local→S3, …); one + fast-list instead of per-object lookups, retries with backoff on every remote + call, tunable concurrency (`--jobs`); streaming export. - **Content-addressed snapshots & diff** — reproducible versions; compare any two. - **Strong validation** — unreadable images, missing/orphan labels, unknown classes, invalid/out-of-bounds boxes, exact + near duplicates, split leakage. @@ -195,6 +199,10 @@ See the [Cloud Sync guide](https://caiowing.github.io/VisionPack/cloud-sync/). (WebDataset shards); exports hardlink from the CAS or stream from the cloud. - **Interoperable I/O** — YOLO (incl. YOLO-seg), COCO, ImageFolder in and out; semantic masks out. +- **Machine-readable everything** — every pipeline command takes `--json` and + prints a stable, schema-versioned envelope on stdout, so services, UIs, and CI + can drive VisionPack without scraping text. See the + [JSON Output guide](https://caiowing.github.io/VisionPack/json-output/). Full command reference and per-command options live in the [CLI guide](https://caiowing.github.io/VisionPack/usage/). @@ -246,7 +254,7 @@ multi-source ingestion (local and cloud) → validation → deterministic splits snapshots → ready-to-train export/packing → evaluation (`vp eval`) and model-in-the-loop labeling (`vp autolabel` / `vp queue`) — works end-to-end across classification, detection, instance/semantic segmentation, and keypoints, -with 118 passing tests. APIs may still shift; feedback and contributions are +with 143 passing tests. APIs may still shift; feedback and contributions are welcome. ```bash diff --git a/docs/cloud-sync.md b/docs/cloud-sync.md index e047d9e..7d4f106 100644 --- a/docs/cloud-sync.md +++ b/docs/cloud-sync.md @@ -33,10 +33,22 @@ pip install "visionpack[azure]" # Azure Blob cause a mismatch. - **Re-sync is metadata-only.** Re-running lists object metadata, sees the etags match, and does nothing — no downloads, no copies. +- **One listing beats many lookups.** Source metadata comes from paginated + LISTs (never a per-object HEAD), and the target CAS is checked with **one + prefix listing per run** instead of a per-object existence check — at 100k + objects that's ~100 LISTs, not 100k HEADs. +- **Transient errors are retried.** Every remote call gets exponential-backoff + retries (throttling, dropped connections, 5xx). If an object still fails, it + is recorded as a per-object failure and the rest of the sync proceeds — the + command exits non-zero so CI can gate on it. {: .note } -v1 is **same-provider** (S3↔S3 or GCS↔GCS). Cross-cloud transfer (S3↔GCS) is on -the roadmap. +Transfers are **server-side within one provider** (S3↔S3, GCS↔GCS) — the bytes +never touch your machine. **Cross-provider** targets (S3→GCS, local→S3, S3→local +…) also work: sync *relays* the bytes it already read to compute the `sha256`, +so a cross-provider copy still costs exactly **one read + one upload**, never a +second download. Relayed uploads are verified against the landed object's +metadata before the index points at them. ## Declare remote sources @@ -75,8 +87,17 @@ Then reconcile. Re-running is idempotent — unchanged objects are skipped entir ```bash vp sync vp sync --source camera-A # just one source +vp sync --jobs 32 # concurrent transfers per source ``` +Remote sources default to **16+ concurrent transfers** (object-store throughput +is latency-bound, so the CPU-derived default would undersize it); tune with +`--jobs`. + +All three source formats work remotely: **YOLO** (images + label dir), +**COCO** (`labels:` points at the instances JSON), and **ImageFolder** +(`root:` is the directory of class subfolders). + ## A content-addressed target Set a `target:` and `copy` mode lands objects in a self-sufficient, @@ -111,7 +132,7 @@ Pick how each source materializes its bytes with `copy:`. | Mode | What it does | Use when | |------|--------------|----------| -| `copy` | Server-side copy into the `target:` content-addressed store. Target is self-sufficient; global dedup. | The common cloud case. | +| `copy` | Copy into the `target:` content-addressed store — server-side when source and target share a provider, single-pass relay when they don't. Target is self-sufficient; global dedup. | The common cloud case. | | `reference` | No copy — the index points straight at the source object. | You control the source bucket and want zero extra storage. | | `ingest` | Download into the **local** CAS (`.vp/objects/`). | Offline / edge work on a remote dataset. | diff --git a/docs/index.md b/docs/index.md index 9e0ee12..43886a5 100644 --- a/docs/index.md +++ b/docs/index.md @@ -60,7 +60,8 @@ packs, and exports. |------|--------|----------| | **Classification** | ImageFolder (folder-per-class) | whole-image label | | **Detection** | YOLO, COCO | bounding box | -| **Instance segmentation** | COCO | polygon | +| **Instance segmentation** | YOLO-seg, COCO | polygon | +| **Semantic segmentation (export)** | — | class-index mask PNGs (`--format masks`) | | **Keypoints / pose** | COCO | keypoints | ## What's in the box @@ -68,10 +69,16 @@ packs, and exports. - **Deterministic, lockable splits** — `stratified` / `random` / `hash`, captured in snapshots. - **Near-duplicate & cross-split leakage detection** — perceptual-hash tier, scale-proof via LSH bucketing. - **Multi-source sync** — declarative `sources:` + `vp sync`, with per-asset provenance. -- **Cloud-native** — sync from and to S3/GCS/Azure without downloading the whole dataset; see {% link cloud-sync.md %}. +- **Cloud-native, multi-provider** — sync YOLO/COCO/ImageFolder sources from and to + S3/GCS/Azure without downloading the whole dataset, including cross-provider + targets; see [Cloud Sync]({% link cloud-sync.md %}). - **Content-addressed snapshots & diff** — reproducible versions; compare any two. - **Strong validation** — unreadable images, missing/orphan labels, unknown classes, bad boxes, exact + near duplicates, split leakage. +- **Benchmarking & model-in-the-loop** — `vp eval` (mAP / accuracy on a locked split), + `vp autolabel` (confident predictions become labels), `vp queue` (what to label next). - **Packing & export** — `archive` (`.tar.zst`) and `training` (WebDataset shards); byte-free exports via hardlinks / streaming manifests. +- **Machine-readable output** — every pipeline command takes `--json` and prints a + stable, schema-versioned envelope; see [JSON Output]({% link json-output.md %}). ## Next steps @@ -79,6 +86,7 @@ packs, and exports. - [Quickstart]({% link quickstart.md %}) — a dataset in 60 seconds - [CLI Guide]({% link usage.md %}) — every command and its options - [Cloud Sync]({% link cloud-sync.md %}) — S3 / GCS / Azure datasets +- [JSON Output]({% link json-output.md %}) — drive VisionPack from other programs {: .note } VisionPack is in **active development** (early but usable). The end-to-end diff --git a/docs/installation.md b/docs/installation.md index 480ed49..47b21df 100644 --- a/docs/installation.md +++ b/docs/installation.md @@ -43,8 +43,9 @@ pip install "visionpack[gcs]" # Google Cloud Storage (gcsfs) pip install "visionpack[azure]" # Azure Blob (adlfs) ``` -You can combine them: `pip install "visionpack[s3,gcs]"`. See {% link cloud-sync.md %} -for declaring remote sources and a cloud target. +You can combine them: `pip install "visionpack[s3,gcs]"`. See +[Cloud Sync]({% link cloud-sync.md %}) for declaring remote sources and a cloud +target. ## Develop from source diff --git a/docs/json-output.md b/docs/json-output.md new file mode 100644 index 0000000..9f0d190 --- /dev/null +++ b/docs/json-output.md @@ -0,0 +1,105 @@ +--- +title: JSON Output +nav_order: 6 +--- + +# JSON Output — drive VisionPack from other programs +{: .no_toc } + +Every pipeline command accepts `--json` and prints **exactly one +machine-readable JSON document to stdout** — no progress bars, no prose. This +is the stable contract for wrapping VisionPack in another program: a backend +service, a UI, CI, a notebook. + +1. TOC +{:toc} + +## The envelope + +Success: + +```json +{ + "schema": 1, + "command": "sync", + "data": { "...": "command-specific payload" } +} +``` + +Failure (the process also exits non-zero): + +```json +{ + "schema": 1, + "command": "diff", + "error": { "type": "VisionPackError", "message": "No snapshot named 'v99'." } +} +``` + +Rules a consumer can rely on: + +- `schema` bumps **only on a breaking change** to the envelope or an existing + `data` shape. New fields may appear without a bump — parse leniently. +- Success has `data`; failure has `error` (never both). Check `error` + + the exit code, not stderr. +- Exit codes keep their CLI meaning: `0` success, `1` domain failure (e.g. + validation errors, ingest failures), `2` command error. + +## Commands and payloads + +| Command | `data` highlights | +|---|---| +| `vp sync --json` | `summaries[]` (per source: `assets_added`, `assets_existing`, `annotations`, `objects`, `failures[]`), `total_assets_added`, `total_failures` | +| `vp sync --dry-run --json` | `plans[]` (per source: `images_found`, `labels_found`, `matched`, `class_names[]`) | +| `vp import ... --json` | `assets`, `annotations`, `objects`, `classes_added`, `recorded_source`, `failures[]` | +| `vp validate --json` | `ok`, `errors`, `warnings`, `issues[]` (severity, code, message, asset_id, path) | +| `vp stats --json` | `stats` (counts, `class_distribution`, `resolutions`), `splits` (per-split breakdowns) | +| `vp split create/lock/list/show --json` | `id`, `strategy`, `locked`, `sets` (name → count); `show` adds `asset_ids` | +| `vp snapshot create/list/show --json` | snapshot records (`version`, `message`, `created_at`, `stats`) | +| `vp diff v1 v2 --json` | `assets_added/removed`, `annotations_added/removed/modified`, `classes_added/removed`, `splits_changed` | +| `vp export --json` | `format`, `output`, per-format counts (`images`, `objects`, `sets`, `streamed`) | +| `vp pack --json` | `profile`, `format`, shard/archive counts and paths | +| `vp fsck --json` | `ok`, `mode`, `checked_assets`, `checked_objects`, `issues[]` | +| `vp eval ... --json` | full result: `task`, `scope`, `metrics` (mAP@50, mAP@50-95, accuracy…), `per_class` | +| `vp autolabel ... --json` | `labeled`, `objects`, `skipped_existing`, `skipped_low_confidence`, `unmatched`, `unknown_classes[]` | +| `vp queue --json` | `total`, `items[]` (`asset_id`, `score`, `reasons[]`) | + +## Example: a pipeline from a script + +```bash +set -e +vp sync --json > sync.json +vp validate --json > validate.json || echo "validation found errors" +vp split create --json > split.json +vp snapshot create -m "auto $(date -I)" --json > snapshot.json + +jq '.data.total_assets_added' sync.json +jq '.data.errors' validate.json +jq '.data.version' snapshot.json +``` + +Or from Python, without parsing text: + +```python +import json, subprocess + +def vp(*argv: str) -> dict: + proc = subprocess.run(["vp", *argv, "--json"], capture_output=True, text=True) + envelope = json.loads(proc.stdout) + if "error" in envelope: + raise RuntimeError(envelope["error"]["message"]) + return envelope["data"] + +added = vp("sync")["total_assets_added"] +report = vp("validate") +``` + +{: .note } +Prefer the JSON contract over importing `visionpack` internals when driving the +tool from another service: the CLI + envelope is the supported integration +surface, and the `schema` field is the compatibility signal. + +## See also + +- [CLI Guide]({% link usage.md %}) — the same commands, human-readable. +- [Cloud Sync]({% link cloud-sync.md %}) — remote sources and targets. diff --git a/docs/usage.md b/docs/usage.md index 804a718..3e81f25 100644 --- a/docs/usage.md +++ b/docs/usage.md @@ -35,6 +35,11 @@ uv run python -m visionpack --help # or run the module directly The examples below use the bare `vp` command; prefix with `uv run` when working from a source checkout. +{: .tip } +Every pipeline command also takes `--json` and prints one machine-readable, +schema-versioned document to stdout — the supported way to drive VisionPack +from another program. See [JSON Output]({% link json-output.md %}). + ## Initialize A Dataset Create a VisionPack project in the current directory: @@ -119,6 +124,7 @@ sources: vp sync --dry-run # preview found / matched / unmatched / classes per source vp sync # ingest; idempotent, records per-asset provenance vp sync --source camera-A # sync just one source +vp sync --jobs 32 # concurrent transfers per source (remote defaults to 16+) ``` Sources can also live in object stores. Remote URIs go anywhere a local path diff --git a/tests/test_cloud_hardening.py b/tests/test_cloud_hardening.py new file mode 100644 index 0000000..4e0ff83 --- /dev/null +++ b/tests/test_cloud_hardening.py @@ -0,0 +1,354 @@ +from __future__ import annotations + +import io +import json +import tempfile +import unittest +from pathlib import Path + +import fsspec +from PIL import Image + +from visionpack.core.errors import VisionPackError +from visionpack.core.project import Project +from visionpack.sources import retry, sync_sources +from visionpack.sources.importer import SourceSyncer +from visionpack.sources.resolver import FsspecResolver +from visionpack.sources.retry import with_retries +from visionpack.sources.schema import Source + + +def _png_bytes(seed: int, size: tuple[int, int] = (32, 24)) -> bytes: + buffer = io.BytesIO() + Image.new("RGB", size, (seed * 7 % 256, seed * 13 % 256, seed * 29 % 256)).save(buffer, format="PNG") + return buffer.getvalue() + + +def _clear_memory_fs(*prefixes: str) -> fsspec.AbstractFileSystem: + fs = fsspec.filesystem("memory") + for prefix in prefixes: + for path in list(fs.find(prefix)): + fs.rm(path) + return fs + + +class RetryTest(unittest.TestCase): + def setUp(self) -> None: + self._delay = retry.BASE_DELAY_SECONDS + retry.BASE_DELAY_SECONDS = 0.0 + + def tearDown(self) -> None: + retry.BASE_DELAY_SECONDS = self._delay + + def test_transient_failure_is_retried_until_success(self) -> None: + calls = {"n": 0} + + def flaky() -> str: + calls["n"] += 1 + if calls["n"] < 3: + raise ConnectionResetError("throttled") + return "ok" + + self.assertEqual(with_retries("read(x)", flaky), "ok") + self.assertEqual(calls["n"], 3) + + def test_exhaustion_becomes_a_visionpack_error(self) -> None: + calls = {"n": 0} + + def always_down() -> None: + calls["n"] += 1 + raise TimeoutError("gateway timeout") + + with self.assertRaises(VisionPackError) as ctx: + with_retries("read(s3://b/k)", always_down) + self.assertIn("read(s3://b/k)", str(ctx.exception)) + self.assertEqual(calls["n"], retry.MAX_ATTEMPTS) + + def test_permanent_errors_are_not_retried(self) -> None: + calls = {"n": 0} + + def missing() -> None: + calls["n"] += 1 + raise FileNotFoundError("no such key") + + with self.assertRaises(FileNotFoundError): + with_retries("read(x)", missing) + self.assertEqual(calls["n"], 1) + + def test_resolver_read_rides_out_transient_provider_errors(self) -> None: + fs = _clear_memory_fs("/retrysrc") + fs.pipe("/retrysrc/a.bin", b"payload") + original = type(fs).cat_file + calls = {"n": 0} + + def flaky(self, path, *args, **kwargs): # noqa: ANN001 + calls["n"] += 1 + if calls["n"] < 3: + raise ConnectionResetError("blip") + return original(self, path, *args, **kwargs) + + type(fs).cat_file = flaky # type: ignore[method-assign] + try: + data = FsspecResolver("memory").read_bytes("memory://retrysrc/a.bin") + finally: + type(fs).cat_file = original # type: ignore[method-assign] + self.assertEqual(data, b"payload") + self.assertEqual(calls["n"], 3) + + +class TargetFastListTest(unittest.TestCase): + """The target CAS membership comes from one prefix listing, not per-object + existence checks — and relayed uploads are size-verified.""" + + def setUp(self) -> None: + self.fs = _clear_memory_fs("/flsrc", "/fldst") + self.fs.pipe("/flsrc/imgs/a.png", _png_bytes(1)) + self.fs.pipe("/flsrc/imgs/b.png", _png_bytes(2)) + self.fs.pipe("/flsrc/lbls/a.txt", b"0 0.5 0.5 0.4 0.4\n") + self.fs.pipe("/flsrc/lbls/b.txt", b"0 0.5 0.5 0.4 0.4\n") + self.fs.pipe("/flsrc/lbls/classes.txt", b"cat\n") + self._delay = retry.BASE_DELAY_SECONDS + retry.BASE_DELAY_SECONDS = 0.0 + + def tearDown(self) -> None: + retry.BASE_DELAY_SECONDS = self._delay + + def _project(self, tmp: str) -> Project: + root = Path(tmp) + Project.init(root, name="fastlist") + project = Project.open(root) + project.manifest.sources = [ + {"name": "s1", "images": "memory://flsrc/imgs", "labels": "memory://flsrc/lbls", "copy": "copy"} + ] + project.manifest.target = "memory://fldst" + project.save_manifest() + return Project.open(root) + + def test_no_per_object_existence_checks_against_the_target_cas(self) -> None: + checked: list[str] = [] + original = FsspecResolver.exists + + def counting(self: FsspecResolver, uri: str) -> bool: + checked.append(uri) + return original(self, uri) + + FsspecResolver.exists = counting # type: ignore[method-assign] + try: + with tempfile.TemporaryDirectory() as tmp: + summary = sync_sources(self._project(tmp))[0] + finally: + FsspecResolver.exists = original # type: ignore[method-assign] + + self.assertEqual(summary.assets_added, 2) + object_heads = [uri for uri in checked if "/objects/sha256/" in uri] + self.assertEqual(object_heads, [], "membership must come from the prefix listing, not per-object checks") + + def test_truncated_relay_upload_is_caught_and_recorded_as_failure(self) -> None: + # Local source -> memory target forces the relay path; corrupt it. + original = FsspecResolver.write_bytes + + def truncating(self: FsspecResolver, uri: str, data: bytes) -> None: + return original(self, uri, data[:-1]) + + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + Project.init(root, name="verify") + imgs = root / "imgs" + imgs.mkdir() + (imgs / "a.png").write_bytes(_png_bytes(3)) + project = Project.open(root) + project.manifest.sources = [{"name": "s1", "images": imgs.as_posix(), "copy": "copy"}] + project.manifest.target = "memory://fldst" + project.save_manifest() + + FsspecResolver.write_bytes = truncating # type: ignore[method-assign] + try: + summary = sync_sources(Project.open(root))[0] + finally: + FsspecResolver.write_bytes = original # type: ignore[method-assign] + + self.assertEqual(summary.assets_added, 0) + self.assertEqual(len(summary.failures), 1) + self.assertIn("landed", summary.failures[0].error) + self.assertEqual(Project.open(root).index.count_assets(), 0) + + +class LabelSingleReadTest(unittest.TestCase): + def test_each_remote_label_is_fetched_exactly_once_when_inferring_classes(self) -> None: + _clear_memory_fs("/lblsrc") + fs = fsspec.filesystem("memory") + fs.pipe("/lblsrc/imgs/a.png", _png_bytes(1)) + fs.pipe("/lblsrc/imgs/b.png", _png_bytes(2)) + # No classes.txt anywhere: class names must be inferred from the labels. + fs.pipe("/lblsrc/lbls/a.txt", b"1 0.5 0.5 0.4 0.4\n") + fs.pipe("/lblsrc/lbls/b.txt", b"0 0.5 0.5 0.4 0.4\n") + + reads: list[str] = [] + original = FsspecResolver.read_bytes + + def counting(self: FsspecResolver, uri: str) -> bytes: + reads.append(uri) + return original(self, uri) + + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + Project.init(root, name="labels") + project = Project.open(root) + project.manifest.sources = [ + {"name": "s1", "images": "memory://lblsrc/imgs", "labels": "memory://lblsrc/lbls"} + ] + project.save_manifest() + + FsspecResolver.read_bytes = counting # type: ignore[method-assign] + try: + summary = sync_sources(Project.open(root))[0] + finally: + FsspecResolver.read_bytes = original # type: ignore[method-assign] + + self.assertEqual(summary.assets_added, 2) + self.assertEqual(summary.objects, 2) + label_reads = [uri for uri in reads if uri.endswith(".txt")] + self.assertEqual(len(label_reads), 2, f"labels must be read once each, got: {label_reads}") + # Inferred names cover the highest class index seen (0 and 1). + names = {item.name for item in Project.open(root).manifest.classes} + self.assertEqual(names, {"class_0", "class_1"}) + + +class RemoteImageFolderTest(unittest.TestCase): + def setUp(self) -> None: + self.fs = _clear_memory_fs("/ifsrc") + for index, class_name in ((1, "ok"), (2, "ok"), (3, "defect")): + self.fs.pipe(f"/ifsrc/{class_name}/img{index}.png", _png_bytes(index)) + + def _project(self, tmp: str) -> Project: + root = Path(tmp) + Project.init(root, name="cls", task="classification") + project = Project.open(root) + project.manifest.sources = [{"name": "folder-a", "format": "imagefolder", "root": "memory://ifsrc"}] + project.save_manifest() + return Project.open(root) + + def test_remote_imagefolder_syncs_with_whole_image_labels(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + summary = sync_sources(self._project(tmp))[0] + self.assertEqual(summary.assets_added, 3) + self.assertEqual(summary.annotations, 3) + self.assertEqual(summary.classes_added, 2) + + project = Project.open(Path(tmp)) + names = {item.name for item in project.manifest.classes} + self.assertEqual(names, {"ok", "defect"}) + for asset in project.index.assets(): + self.assertEqual(asset.source, "folder-a") + annotation = project.index.annotation_for_asset(asset.id) + assert annotation is not None + self.assertEqual(len(annotation.objects), 1) + self.assertIsNone(annotation.objects[0].geometry) + + def test_remote_imagefolder_resync_is_idempotent(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + sync_sources(self._project(tmp)) + summary = sync_sources(Project.open(Path(tmp)))[0] + self.assertEqual(summary.assets_added, 0) + self.assertEqual(summary.assets_existing, 3) + + +class RemoteCocoTest(unittest.TestCase): + def setUp(self) -> None: + self.fs = _clear_memory_fs("/cocosrc") + self.fs.pipe("/cocosrc/imgs/a.png", _png_bytes(1)) + self.fs.pipe("/cocosrc/imgs/sub/b.png", _png_bytes(2)) + document = { + "images": [ + {"id": 1, "file_name": "a.png", "width": 32, "height": 24}, + {"id": 2, "file_name": "sub/b.png", "width": 32, "height": 24}, + {"id": 3, "file_name": "missing.png", "width": 32, "height": 24}, + ], + "annotations": [ + {"id": 10, "image_id": 1, "category_id": 7, "bbox": [2, 3, 10, 8], "iscrowd": 0}, + {"id": 11, "image_id": 1, "category_id": 8, "bbox": [1, 1, 5, 5], "iscrowd": 0}, + {"id": 12, "image_id": 2, "category_id": 8, "bbox": [4, 4, 6, 6], "iscrowd": 0}, + ], + "categories": [{"id": 7, "name": "scratch"}, {"id": 8, "name": "dent"}], + } + self.fs.pipe("/cocosrc/instances.json", json.dumps(document).encode("utf-8")) + + def _project(self, tmp: str) -> Project: + root = Path(tmp) + Project.init(root, name="cocor") + project = Project.open(root) + project.manifest.sources = [ + { + "name": "coco-a", + "format": "coco", + "images": "memory://cocosrc/imgs", + "labels": "memory://cocosrc/instances.json", + } + ] + project.save_manifest() + return Project.open(root) + + def test_remote_coco_syncs_images_and_annotations(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + summary = sync_sources(self._project(tmp))[0] + self.assertEqual(summary.assets_added, 2) + self.assertEqual(summary.annotations, 2) + self.assertEqual(summary.objects, 3) + self.assertEqual(summary.classes_added, 2) + # The image declared in the JSON but absent from the listing is a + # clean per-image failure, not an aborted sync. + self.assertEqual(len(summary.failures), 1) + self.assertIn("missing.png", summary.failures[0].path) + + project = Project.open(Path(tmp)) + names = {item.name for item in project.manifest.classes} + self.assertEqual(names, {"scratch", "dent"}) + class_by_name = {item.name: item.id for item in project.manifest.classes} + annotated = [ + project.index.annotation_for_asset(asset.id) + for asset in project.index.assets() + if asset.source == "coco-a" + ] + all_objects = [obj for ann in annotated if ann for obj in ann.objects] + self.assertEqual(len(all_objects), 3) + self.assertEqual( + sorted(obj.class_id for obj in all_objects), + sorted([class_by_name["scratch"], class_by_name["dent"], class_by_name["dent"]]), + ) + bboxes = {(obj.bbox.x, obj.bbox.y, obj.bbox.width, obj.bbox.height) for obj in all_objects} + self.assertIn((2.0, 3.0, 10.0, 8.0), bboxes) + + def test_remote_coco_resync_is_idempotent(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + sync_sources(self._project(tmp)) + summary = sync_sources(Project.open(Path(tmp)))[0] + self.assertEqual(summary.assets_added, 0) + self.assertEqual(summary.assets_existing, 2) + + +class PoolSizeTest(unittest.TestCase): + def _syncer(self, tmp: str, images_uri: str, max_workers: int | None) -> SourceSyncer: + root = Path(tmp) + Project.init(root, name="jobs") + project = Project.open(root) + source = Source.from_dict({"name": "s1", "images": images_uri}) + return SourceSyncer(project, source, max_workers=max_workers) + + def test_jobs_flag_wins(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + self.assertEqual(self._syncer(tmp, "memory://x/imgs", 3)._pool_size(), 3) + + def test_remote_default_has_a_floor_of_16(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + size = self._syncer(tmp, "memory://x/imgs", None)._pool_size() + assert size is not None + self.assertGreaterEqual(size, 16) + self.assertLessEqual(size, 32) + + def test_local_default_is_the_executor_default(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + self.assertIsNone(self._syncer(tmp, "./imgs", None)._pool_size()) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_cloud_target.py b/tests/test_cloud_target.py index 6e1028d..c66c38d 100644 --- a/tests/test_cloud_target.py +++ b/tests/test_cloud_target.py @@ -107,8 +107,10 @@ def test_remote_asset_cannot_be_materialized_locally(self) -> None: with self.assertRaises(VisionPackError): asset.resolved_path(Path(tmp)) - def test_cross_provider_copy_is_rejected(self) -> None: - # A local source feeding a remote target can't be a server-side copy. + def test_cross_provider_copy_relays_already_read_bytes(self) -> None: + # A local source feeding a remote target can't be a server-side copy, so + # the sync relays the bytes it already read for hashing: one upload per + # object, no server_copy call, target still content-addressed. with tempfile.TemporaryDirectory() as tmp: root = Path(tmp) Project.init(root, name="x") @@ -119,8 +121,66 @@ def test_cross_provider_copy_is_rejected(self) -> None: project.manifest.sources = [{"name": "s1", "images": imgs.as_posix(), "copy": "copy"}] project.manifest.target = "memory://dst" project.save_manifest() - with self.assertRaises(VisionPackError): - sync_sources(Project.open(root)) + + copies: list[str] = [] + original = FsspecResolver.server_copy + + def counting(self: FsspecResolver, src: str, dst: str) -> None: + copies.append(dst) + return original(self, src, dst) + + FsspecResolver.server_copy = counting # type: ignore[method-assign] + try: + summary = sync_sources(Project.open(root))[0] + finally: + FsspecResolver.server_copy = original # type: ignore[method-assign] + + self.assertEqual(summary.assets_added, 1) + self.assertEqual(copies, [], "cross-provider transfer must not attempt a server-side copy") + self.assertEqual(len(_target_objects(self.fs)), 1) + for asset in Project.open(root).index.assets(): + self.assertTrue(asset.path.startswith("memory://dst/objects/sha256/")) + + def test_cross_provider_resync_is_metadata_only(self) -> None: + with tempfile.TemporaryDirectory() as tmp: + root = Path(tmp) + Project.init(root, name="x") + imgs = root / "imgs" + imgs.mkdir() + (imgs / "a.png").write_bytes(_png_bytes(1)) + project = Project.open(root) + project.manifest.sources = [{"name": "s1", "images": imgs.as_posix(), "copy": "copy"}] + project.manifest.target = "memory://dst" + project.save_manifest() + sync_sources(Project.open(root)) + + writes: list[str] = [] + original = FsspecResolver.write_bytes + + def counting(self: FsspecResolver, uri: str, data: bytes) -> None: + writes.append(uri) + return original(self, uri, data) + + FsspecResolver.write_bytes = counting # type: ignore[method-assign] + try: + summary = sync_sources(Project.open(root))[0] + finally: + FsspecResolver.write_bytes = original # type: ignore[method-assign] + + self.assertEqual(summary.assets_added, 0) + self.assertEqual(summary.assets_existing, 1) + self.assertEqual(writes, [], "unchanged cross-provider re-sync re-uploaded objects") + + def test_remote_source_relays_into_local_target(self) -> None: + # memory:// source, local-directory target: the opposite relay direction. + with tempfile.TemporaryDirectory() as tmp, tempfile.TemporaryDirectory() as dst: + project = self._project(tmp, "copy", target=Path(dst).as_posix()) + summary = sync_sources(project)[0] + self.assertEqual(summary.assets_added, 2) + landed = sorted(p for p in Path(dst).rglob("*") if p.is_file()) + self.assertEqual(len(landed), 2) + for path in landed: + self.assertIn("objects/sha256/", path.as_posix()) class ManifestTargetTest(unittest.TestCase): diff --git a/tests/test_json_contract.py b/tests/test_json_contract.py new file mode 100644 index 0000000..0b3a1b6 --- /dev/null +++ b/tests/test_json_contract.py @@ -0,0 +1,161 @@ +from __future__ import annotations + +import io +import json +import os +import tempfile +import unittest +from contextlib import redirect_stdout +from pathlib import Path + +from PIL import Image + +from visionpack.cli.main import main +from visionpack.cli.output import SCHEMA_VERSION +from visionpack.core.project import Project + + +def _png(path: Path, seed: int) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + Image.new("RGB", (32, 24), (seed * 7 % 256, seed * 13 % 256, seed * 29 % 256)).save(path, format="PNG") + + +class JsonContractTest(unittest.TestCase): + """`--json` prints exactly one schema-versioned envelope on stdout. + + This is the machine contract external programs (a frontend backend, CI) + build against, so the assertions pin the envelope keys and each command's + load-bearing data fields. + """ + + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + self.root = Path(self._tmp.name) + Project.init(self.root, name="jsonctl", task="detection") + for index in range(1, 4): + _png(self.root / "raw" / "images" / f"img{index}.png", index) + label = self.root / "raw" / "labels" / f"img{index}.txt" + label.parent.mkdir(parents=True, exist_ok=True) + label.write_text("0 0.5 0.5 0.4 0.4\n", encoding="utf-8") + (self.root / "raw" / "labels" / "classes.txt").write_text("cat\n", encoding="utf-8") + project = Project.open(self.root) + project.manifest.sources = [ + {"name": "cam-a", "images": "./raw/images", "labels": "./raw/labels", "match": "stem", "copy": "ingest"} + ] + project.save_manifest() + self._cwd = os.getcwd() + os.chdir(self.root) + + def tearDown(self) -> None: + os.chdir(self._cwd) + self._tmp.cleanup() + + def _run(self, *argv: str) -> tuple[int, dict]: + buffer = io.StringIO() + with redirect_stdout(buffer): + code = main(list(argv)) + output = buffer.getvalue() + try: + envelope = json.loads(output) + except json.JSONDecodeError as exc: # pragma: no cover - assertion aid + raise AssertionError(f"stdout is not a single JSON document:\n{output}") from exc + self.assertEqual(envelope["schema"], SCHEMA_VERSION) + return code, envelope + + def _sync(self) -> None: + code, _ = self._run("sync", "--json") + self.assertEqual(code, 0) + + def test_sync_dry_run(self) -> None: + code, envelope = self._run("sync", "--dry-run", "--json") + self.assertEqual(code, 0) + self.assertEqual(envelope["command"], "sync") + self.assertTrue(envelope["data"]["dry_run"]) + plan = envelope["data"]["plans"][0] + self.assertEqual(plan["images_found"], 3) + self.assertEqual(plan["matched"], 3) + self.assertEqual(plan["class_names"], ["cat"]) + + def test_sync(self) -> None: + code, envelope = self._run("sync", "--json") + self.assertEqual(code, 0) + data = envelope["data"] + self.assertEqual(data["total_assets_added"], 3) + self.assertEqual(data["total_failures"], 0) + self.assertEqual(data["summaries"][0]["name"], "cam-a") + + def test_stats(self) -> None: + self._sync() + code, envelope = self._run("stats", "--json") + self.assertEqual(code, 0) + self.assertEqual(envelope["data"]["stats"]["assets"], 3) + self.assertIn("splits", envelope["data"]) + + def test_validate(self) -> None: + self._sync() + code, envelope = self._run("validate", "--json") + data = envelope["data"] + self.assertEqual(code, 0 if data["ok"] else 1) + self.assertIn("issues", data) + self.assertEqual(data["errors"], sum(1 for i in data["issues"] if i["severity"] == "error")) + + def test_split_create_list_show_lock(self) -> None: + self._sync() + code, envelope = self._run("split", "create", "--train", "0.5", "--val", "0.25", "--test", "0.25", "--json") + self.assertEqual(code, 0) + self.assertEqual(envelope["command"], "split.create") + self.assertEqual(sum(envelope["data"]["sets"].values()), 3) + + code, envelope = self._run("split", "list", "--json") + self.assertEqual(envelope["data"]["splits"][0]["id"], "default") + + code, envelope = self._run("split", "show", "--json") + self.assertTrue(envelope["data"]["found"]) + self.assertIn("asset_ids", envelope["data"]) + + code, envelope = self._run("split", "lock", "--json") + self.assertTrue(envelope["data"]["locked"]) + + def test_snapshot_create_list_diff(self) -> None: + self._sync() + code, envelope = self._run("snapshot", "create", "-m", "baseline", "--json") + self.assertEqual(code, 0) + self.assertEqual(envelope["command"], "snapshot.create") + version = envelope["data"]["version"] + + code, envelope = self._run("snapshot", "list", "--json") + self.assertEqual(len(envelope["data"]["snapshots"]), 1) + + code, envelope = self._run("snapshot", "create", "-m", "second", "--json") + second = envelope["data"]["version"] + code, envelope = self._run("diff", version, second, "--json") + self.assertEqual(code, 0) + self.assertEqual(envelope["command"], "diff") + self.assertEqual(envelope["data"]["left"], version) + self.assertEqual(envelope["data"]["assets_added"], []) + + def test_export(self) -> None: + self._sync() + code, envelope = self._run("export", "--format", "yolo", "--output", "exports/yolo", "--json") + self.assertEqual(code, 0) + self.assertEqual(envelope["data"]["images"], 3) + self.assertTrue((self.root / "exports" / "yolo").exists()) + + def test_fsck(self) -> None: + self._sync() + code, envelope = self._run("fsck", "--json") + self.assertEqual(code, 0) + self.assertTrue(envelope["data"]["ok"]) + self.assertEqual(envelope["data"]["checked_assets"], 3) + + def test_error_envelope(self) -> None: + code, envelope = self._run("diff", "v98", "v99", "--json") + self.assertEqual(code, 2) + self.assertEqual(envelope["command"], "diff") + self.assertNotIn("data", envelope) + self.assertIn("message", envelope["error"]) + self.assertTrue(envelope["error"]["type"]) + + +if __name__ == "__main__": + unittest.main() diff --git a/visionpack/cli/commands/autolabel.py b/visionpack/cli/commands/autolabel.py index a61e8e2..c56b1f3 100644 --- a/visionpack/cli/commands/autolabel.py +++ b/visionpack/cli/commands/autolabel.py @@ -4,6 +4,7 @@ from pathlib import Path from visionpack.autolabel import apply_predictions +from visionpack.cli.output import emit_json from visionpack.core.project import Project from visionpack.predictions import FORMATS, load_predictions @@ -18,6 +19,7 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> action="store_true", help="Also overwrite assets that already have annotations (default: only unlabeled assets are touched)", ) + parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") parser.set_defaults(func=run) @@ -25,6 +27,12 @@ def run(args: argparse.Namespace) -> int: project = Project.open(".") predictions = load_predictions(project, Path(args.predictions), fmt=args.format) summary = apply_predictions(project, predictions, min_confidence=args.min_confidence, replace=args.replace) + if args.json: + emit_json( + "autolabel", + {**summary, "min_confidence": args.min_confidence, "unknown_classes": sorted(set(predictions.unknown_classes))}, + ) + return 0 print(f"Autolabeled {summary['labeled']} asset(s) with {summary['objects']} object(s) (source recorded as type=model).") if summary["skipped_existing"]: print(f"Skipped {summary['skipped_existing']} already-labeled asset(s); pass --replace to overwrite them.") diff --git a/visionpack/cli/commands/diff.py b/visionpack/cli/commands/diff.py index 87d74d2..056449e 100644 --- a/visionpack/cli/commands/diff.py +++ b/visionpack/cli/commands/diff.py @@ -1,8 +1,8 @@ from __future__ import annotations import argparse -import json +from visionpack.cli.output import emit_json from visionpack.core.project import Project from visionpack.diff import diff_snapshots @@ -11,7 +11,7 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> parser = subparsers.add_parser("diff", help="Diff two snapshots") parser.add_argument("left") parser.add_argument("right") - parser.add_argument("--json", action="store_true", help="Print JSON") + parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") parser.set_defaults(func=run) @@ -19,7 +19,7 @@ def run(args: argparse.Namespace) -> int: project = Project.open(".") result = diff_snapshots(project, args.left, args.right) if args.json: - print(json.dumps(result, indent=2, sort_keys=True)) + emit_json("diff", {"left": args.left, "right": args.right, **result}) return 0 print(f"Diff {args.left} -> {args.right}") for key in ( diff --git a/visionpack/cli/commands/eval.py b/visionpack/cli/commands/eval.py index 01cb2ed..a06bbd2 100644 --- a/visionpack/cli/commands/eval.py +++ b/visionpack/cli/commands/eval.py @@ -1,9 +1,9 @@ from __future__ import annotations import argparse -import json from pathlib import Path +from visionpack.cli.output import emit_json from visionpack.core.project import Project from visionpack.eval import evaluate from visionpack.predictions import FORMATS, load_predictions @@ -32,7 +32,7 @@ def run(args: argparse.Namespace) -> int: conf_threshold=args.conf, ) if args.json: - print(json.dumps(result, indent=2, sort_keys=True)) + emit_json("eval", result) return 0 scope = result["scope"] diff --git a/visionpack/cli/commands/export.py b/visionpack/cli/commands/export.py index 3482c8c..f663147 100644 --- a/visionpack/cli/commands/export.py +++ b/visionpack/cli/commands/export.py @@ -3,6 +3,7 @@ import argparse from pathlib import Path +from visionpack.cli.output import emit_json from visionpack.core.project import Project from visionpack.formats.classification import export_imagefolder from visionpack.formats.coco import export_coco @@ -33,6 +34,7 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> "--snapshot", help="Export the dataset as it was at this snapshot version (e.g. v2) instead of the current state", ) + parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") parser.set_defaults(func=run) @@ -41,6 +43,8 @@ def run(args: argparse.Namespace) -> int: if args.snapshot: project = open_snapshot(project, args.snapshot) output = Path(args.output) + if args.json: + return _run_json(args, project, output) with cli_progress(f"Exporting {args.format}") as callback: if args.format == "coco": summary = export_coco(project, output, split_id=args.split, progress=callback) @@ -72,3 +76,26 @@ def run(args: argparse.Namespace) -> int: if summary.get("skipped"): print(f"Skipped {summary['skipped']} assets not assigned to any set in split {args.split!r}") return 0 + + +def _run_json(args: argparse.Namespace, project: Project, output: Path) -> int: + # No progress bar: stdout carries exactly one JSON document. + if args.format == "coco": + summary = export_coco(project, output, split_id=args.split) + elif args.format == "imagefolder": + summary = export_imagefolder(project, output, split_id=args.split) + elif args.format == "masks": + summary = export_masks(project, output, split_id=args.split) + else: + summary = export_yolo(project, output, split_id=args.split, seg=args.seg) + emit_json( + "export", + { + "format": args.format, + "output": str(output.resolve()), + "split": args.split, + "snapshot": args.snapshot, + **summary, + }, + ) + return 0 diff --git a/visionpack/cli/commands/fsck.py b/visionpack/cli/commands/fsck.py index 5afa1c2..c1a4743 100644 --- a/visionpack/cli/commands/fsck.py +++ b/visionpack/cli/commands/fsck.py @@ -1,7 +1,9 @@ from __future__ import annotations import argparse +from dataclasses import asdict +from visionpack.cli.output import emit_json from visionpack.core.project import Project from visionpack.fsck import run_fsck @@ -18,6 +20,7 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> action="store_true", help="Skip the scan for unreferenced objects in the store", ) + parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") parser.set_defaults(func=run) @@ -25,6 +28,20 @@ def run(args: argparse.Namespace) -> int: project = Project.open(".") report = run_fsck(project, deep=args.deep, check_orphans=not args.no_orphans) mode = "deep" if args.deep else "quick" + if args.json: + emit_json( + "fsck", + { + "ok": report.ok, + "mode": mode, + "checked_assets": report.checked_assets, + "checked_objects": report.checked_objects, + "errors": len(report.errors), + "warnings": len(report.warnings), + "issues": [asdict(issue) for issue in report.issues], + }, + ) + return 0 if report.ok else 1 print( f"fsck ({mode}): checked {report.checked_assets} assets, {report.checked_objects} objects " f"-> {len(report.errors)} errors, {len(report.warnings)} warnings" diff --git a/visionpack/cli/commands/import_.py b/visionpack/cli/commands/import_.py index a7416a7..384e00d 100644 --- a/visionpack/cli/commands/import_.py +++ b/visionpack/cli/commands/import_.py @@ -2,8 +2,10 @@ import argparse import os +from dataclasses import asdict from pathlib import Path +from visionpack.cli.output import emit_json from visionpack.core.errors import VisionPackError from visionpack.core.lock import project_lock from visionpack.core.project import Project @@ -32,6 +34,7 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> help="Do not add this import as a source in visionpack.yaml (use for one-off/throwaway imports)", ) parser.add_argument("--class-map", help="Reserved for explicit class mapping files") + parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") parser.set_defaults(func=run) @@ -58,6 +61,24 @@ def _run_locked(project: Project, args: argparse.Namespace) -> int: importer = YoloImporter(project, Path(args.source), copy_mode=args.copy) label = "YOLO" + if args.json: + summary = importer.run() + recorded = None if args.no_record else _record_source(project, args) + emit_json( + "import", + { + "format": args.format, + "assets": summary.assets, + "annotations": summary.annotations, + "objects": summary.objects, + "classes_added": summary.classes_added, + "orphan_labels": summary.orphan_labels, + "recorded_source": recorded, + "failures": [asdict(failure) for failure in summary.failures], + }, + ) + return 1 if summary.failures else 0 + with cli_progress(f"Importing {label}") as callback: summary = importer.run(progress=callback) diff --git a/visionpack/cli/commands/pack.py b/visionpack/cli/commands/pack.py index 4d4ff86..c81072d 100644 --- a/visionpack/cli/commands/pack.py +++ b/visionpack/cli/commands/pack.py @@ -1,8 +1,10 @@ from __future__ import annotations import argparse +from dataclasses import asdict from pathlib import Path +from visionpack.cli.output import emit_json from visionpack.core.errors import VisionPackError from visionpack.core.project import Project from visionpack.packing import pack_archive, pack_training @@ -19,6 +21,7 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> default=None, help="For training packs: emit per-set shards from this split (default 'default' when given without a value)", ) + parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") parser.set_defaults(func=run) @@ -33,6 +36,9 @@ def run(args: argparse.Namespace) -> int: if fmt == "webdataset": summary = pack_training(project, output=output, profile_name=args.profile, split_id=args.split) + if args.json: + emit_json("pack", {"profile": args.profile, "format": fmt, **asdict(summary)}) + return 0 sets = ", ".join(f"{name}={count}" for name, count in summary.sets.items()) print( f"Packed WebDataset to {summary.path} " @@ -44,6 +50,9 @@ def run(args: argparse.Namespace) -> int: if fmt in {"tar", "tar.zst"}: summary = pack_archive(project, output=output, profile_name=args.profile) + if args.json: + emit_json("pack", {"profile": args.profile, **asdict(summary)}) + return 0 print( f"Packed archive: {summary.path} " f"({summary.format}, {summary.files} files, {summary.assets} assets, {summary.size_bytes} bytes)" diff --git a/visionpack/cli/commands/queue.py b/visionpack/cli/commands/queue.py index c687670..77f10c5 100644 --- a/visionpack/cli/commands/queue.py +++ b/visionpack/cli/commands/queue.py @@ -1,9 +1,9 @@ from __future__ import annotations import argparse -import json from pathlib import Path +from visionpack.cli.output import emit_json from visionpack.core.project import Project from visionpack.curation import rank_for_annotation from visionpack.predictions import FORMATS, load_predictions @@ -34,7 +34,8 @@ def run(args: argparse.Namespace) -> int: ranked = rank_for_annotation(project, predictions, include_labeled=args.include_labeled) if args.json: - print(json.dumps(ranked[: args.limit] if args.limit else ranked, indent=2)) + shown = ranked[: args.limit] if args.limit else ranked + emit_json("queue", {"total": len(ranked), "items": shown}) return 0 if not ranked: diff --git a/visionpack/cli/commands/snapshot.py b/visionpack/cli/commands/snapshot.py index a2560cd..543076d 100644 --- a/visionpack/cli/commands/snapshot.py +++ b/visionpack/cli/commands/snapshot.py @@ -3,6 +3,7 @@ import argparse import json +from visionpack.cli.output import emit_json from visionpack.core.lock import project_lock from visionpack.core.project import Project from visionpack.snapshot import create_snapshot, list_snapshots, load_snapshot @@ -14,13 +15,16 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> create = nested.add_parser("create", help="Create a snapshot") create.add_argument("-m", "--message", required=True, help="Snapshot message") + create.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") create.set_defaults(func=run_create) list_parser = nested.add_parser("list", help="List snapshots") + list_parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") list_parser.set_defaults(func=run_list) show = nested.add_parser("show", help="Show snapshot details") show.add_argument("version", help="Snapshot version, e.g. v1") + show.add_argument("--json", action="store_true", help="Print the machine-readable JSON envelope") show.set_defaults(func=run_show) @@ -28,6 +32,9 @@ def run_create(args: argparse.Namespace) -> int: project = Project.open(".") with project_lock(project.root): snapshot = create_snapshot(project, args.message) + if args.json: + emit_json("snapshot.create", snapshot) + return 0 print(f"Created snapshot {snapshot['version']}: {snapshot['message']}") return 0 @@ -35,6 +42,9 @@ def run_create(args: argparse.Namespace) -> int: def run_list(args: argparse.Namespace) -> int: project = Project.open(".") snapshots = list_snapshots(project) + if args.json: + emit_json("snapshot.list", {"snapshots": snapshots}) + return 0 for item in snapshots: stats = item.get("stats", {}) counts = f"{stats.get('assets', '?')} imgs, {stats.get('objects', '?')} objs" @@ -46,5 +56,9 @@ def run_list(args: argparse.Namespace) -> int: def run_show(args: argparse.Namespace) -> int: project = Project.open(".") - print(json.dumps(load_snapshot(project, args.version), indent=2, sort_keys=True)) + snapshot = load_snapshot(project, args.version) + if args.json: + emit_json("snapshot.show", snapshot) + return 0 + print(json.dumps(snapshot, indent=2, sort_keys=True)) return 0 diff --git a/visionpack/cli/commands/split.py b/visionpack/cli/commands/split.py index d49d424..98fcab8 100644 --- a/visionpack/cli/commands/split.py +++ b/visionpack/cli/commands/split.py @@ -2,11 +2,21 @@ import argparse +from visionpack.cli.output import emit_json from visionpack.core.lock import project_lock from visionpack.core.project import Project from visionpack.split import create_split, lock_split +def _split_payload(split) -> dict: # noqa: ANN001 + return { + "id": split.id, + "strategy": split.strategy, + "locked": split.locked, + "sets": {name: len(ids) for name, ids in split.sets.items()}, + } + + def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> None: parser = subparsers.add_parser("split", help="Create and manage deterministic, versioned splits") nested = parser.add_subparsers(dest="split_command", required=True) @@ -25,17 +35,21 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> create.add_argument("--seed", type=int, default=0, help="Seed mixed into the content hash for reproducibility") create.add_argument("--id", dest="split_id", default="default", help="Split id (default: 'default')") create.add_argument("--force", action="store_true", help="Overwrite even if the split is locked") + create.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") create.set_defaults(func=run_create) lock = nested.add_parser("lock", help="Lock a split so it cannot be changed") lock.add_argument("--id", dest="split_id", default="default", help="Split id (default: 'default')") + lock.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") lock.set_defaults(func=run_lock) list_parser = nested.add_parser("list", help="List splits") + list_parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") list_parser.set_defaults(func=run_list) show = nested.add_parser("show", help="Show set sizes for a split") show.add_argument("--id", dest="split_id", default="default", help="Split id (default: 'default')") + show.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") show.set_defaults(func=run_show) @@ -53,6 +67,9 @@ def run_create(args: argparse.Namespace) -> int: by=args.by, force=args.force, ) + if args.json: + emit_json("split.create", {**_split_payload(split), "seed": args.seed}) + return 0 sizes = ", ".join(f"{name}={len(ids)}" for name, ids in split.sets.items()) print(f"Created split {split.id!r} (strategy={split.strategy}, seed={args.seed}): {sizes}") return 0 @@ -62,6 +79,9 @@ def run_lock(args: argparse.Namespace) -> int: project = Project.open(".") with project_lock(project.root): split = lock_split(project, args.split_id) + if args.json: + emit_json("split.lock", _split_payload(split)) + return 0 print(f"Locked split {split.id!r}. It will be captured as-is in snapshots.") return 0 @@ -69,6 +89,9 @@ def run_lock(args: argparse.Namespace) -> int: def run_list(args: argparse.Namespace) -> int: project = Project.open(".") splits = project.index.splits() + if args.json: + emit_json("split.list", {"splits": [_split_payload(split) for split in splits]}) + return 0 for split in splits: sizes = ", ".join(f"{name}={len(ids)}" for name, ids in split.sets.items()) lock = " [locked]" if split.locked else "" @@ -82,8 +105,14 @@ def run_show(args: argparse.Namespace) -> int: project = Project.open(".") split = next((item for item in project.index.splits() if item.id == args.split_id), None) if split is None: + if args.json: + emit_json("split.show", {"id": args.split_id, "found": False}) + return 1 print(f"No split named {args.split_id!r}. Create one with: vp split create") return 1 + if args.json: + emit_json("split.show", {**_split_payload(split), "found": True, "asset_ids": split.sets}) + return 0 print(f"Split {split.id!r} strategy={split.strategy} locked={split.locked}") for name, ids in split.sets.items(): print(f" {name}: {len(ids)} images") diff --git a/visionpack/cli/commands/stats.py b/visionpack/cli/commands/stats.py index 3f12aee..03c3bb6 100644 --- a/visionpack/cli/commands/stats.py +++ b/visionpack/cli/commands/stats.py @@ -1,8 +1,8 @@ from __future__ import annotations import argparse -import json +from visionpack.cli.output import emit_json from visionpack.core.project import Project from visionpack.stats import collect_stats, split_breakdown @@ -10,7 +10,7 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> None: parser = subparsers.add_parser("stats", help="Show dataset statistics") parser.add_argument("--by", choices=["class", "split"], help="Focus output on one dimension") - parser.add_argument("--json", action="store_true", help="Print JSON") + parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") parser.add_argument("--html", help="Reserved for HTML reports") parser.set_defaults(func=run) @@ -19,7 +19,12 @@ def run(args: argparse.Namespace) -> int: project = Project.open(".") stats = collect_stats(project) if args.json: - print(json.dumps(stats, indent=2, sort_keys=True)) + breakdowns = {} + for split in project.index.splits(): + breakdown = split_breakdown(project, split.id) + if breakdown is not None: + breakdowns[split.id] = breakdown + emit_json("stats", {"stats": stats, "splits": breakdowns}) return 0 if args.by == "class": for class_id, count in stats["class_distribution"].items(): diff --git a/visionpack/cli/commands/sync.py b/visionpack/cli/commands/sync.py index 1fc353d..24daea9 100644 --- a/visionpack/cli/commands/sync.py +++ b/visionpack/cli/commands/sync.py @@ -1,7 +1,9 @@ from __future__ import annotations import argparse +from dataclasses import asdict +from visionpack.cli.output import emit_json from visionpack.core.lock import project_lock from visionpack.core.project import Project from visionpack.progress import cli_progress @@ -19,6 +21,12 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> action="store_true", help="Report what each source would ingest (found/matched/unmatched) without writing", ) + parser.add_argument( + "--jobs", + type=int, + help="Concurrent transfers per source (default: 16+ for remote sources, CPU-derived for local)", + ) + parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") parser.set_defaults(func=run) @@ -27,6 +35,9 @@ def run(args: argparse.Namespace) -> int: if args.dry_run: plans = plan_sources(project, args.source) + if args.json: + emit_json("sync", {"dry_run": True, "plans": [asdict(plan) for plan in plans]}) + return 0 for plan in plans: print(f"[{plan.name}] {plan.format}") print(f" images: {plan.images_uri}") @@ -43,8 +54,22 @@ def run(args: argparse.Namespace) -> int: print(f" classes: {', '.join(plan.class_names) if plan.class_names else '(none discovered)'}") return 0 + # JSON mode keeps stdout to a single document: no progress bar. + progress_factory = None if args.json else cli_progress with project_lock(project.root): - summaries = sync_sources(project, args.source, progress_factory=cli_progress) + summaries = sync_sources(project, args.source, progress_factory=progress_factory, max_workers=args.jobs) + if args.json: + total_failures = sum(len(summary.failures) for summary in summaries) + emit_json( + "sync", + { + "dry_run": False, + "summaries": [asdict(summary) for summary in summaries], + "total_assets_added": sum(summary.assets_added for summary in summaries), + "total_failures": total_failures, + }, + ) + return 1 if total_failures else 0 total_added = 0 total_failures = 0 for summary in summaries: diff --git a/visionpack/cli/commands/validate.py b/visionpack/cli/commands/validate.py index b77d2ac..5164e46 100644 --- a/visionpack/cli/commands/validate.py +++ b/visionpack/cli/commands/validate.py @@ -5,6 +5,7 @@ from dataclasses import asdict from pathlib import Path +from visionpack.cli.output import emit_json from visionpack.core.project import Project from visionpack.validation import validate_project @@ -14,12 +15,25 @@ def register(subparsers: argparse._SubParsersAction[argparse.ArgumentParser]) -> parser.add_argument("--strict", action="store_true", help="Treat missing annotations as errors") parser.add_argument("--fix", action="store_true", help="Reserved for future automatic fixes") parser.add_argument("--report", help="Write JSON validation report") + parser.add_argument("--json", action="store_true", help="Print a machine-readable JSON result") parser.set_defaults(func=run) def run(args: argparse.Namespace) -> int: project = Project.open(".") report = validate_project(project, strict=args.strict) + if args.json: + emit_json( + "validate", + { + "ok": report.ok, + "strict": args.strict, + "errors": len(report.errors), + "warnings": len(report.warnings), + "issues": [asdict(issue) for issue in report.issues], + }, + ) + return 0 if report.ok else 1 print(f"Validation: {len(report.errors)} errors, {len(report.warnings)} warnings") for issue in report.issues[:50]: print(f"[{issue.severity}] {issue.code}: {issue.message}") diff --git a/visionpack/cli/main.py b/visionpack/cli/main.py index 2e42da6..f3dca97 100644 --- a/visionpack/cli/main.py +++ b/visionpack/cli/main.py @@ -64,11 +64,15 @@ def main(argv: list[str] | None = None) -> int: args = parser.parse_args(argv) try: return int(args.func(args) or 0) - except VisionPackError as exc: - print(f"error: {exc}", file=sys.stderr) - return 2 - except FileNotFoundError as exc: - print(f"error: {exc}", file=sys.stderr) + except (VisionPackError, FileNotFoundError) as exc: + # In JSON mode the failure is part of the contract: a structured error + # envelope on stdout (plus the non-zero exit), not prose on stderr. + if getattr(args, "json", False): + from visionpack.cli.output import emit_json_error + + emit_json_error(args.command, exc) + else: + print(f"error: {exc}", file=sys.stderr) return 2 diff --git a/visionpack/cli/output.py b/visionpack/cli/output.py new file mode 100644 index 0000000..929a580 --- /dev/null +++ b/visionpack/cli/output.py @@ -0,0 +1,49 @@ +from __future__ import annotations + +import dataclasses +import json +from pathlib import Path +from typing import Any + +# The machine-readable contract every `--json` flag speaks. Bumped only on a +# breaking change to the envelope or an existing command's `data` shape — +# additive fields do not bump it. Consumers should check `schema` and the +# presence of `error` (plus the process exit code) rather than parse stderr. +SCHEMA_VERSION = 1 + + +def emit_json(command: str, data: Any) -> None: + """Print the success envelope for ``command`` to stdout. + + In JSON mode this must be the only thing on stdout: human-facing prints and + progress bars are suppressed by the callers, so the output is always one + parseable document. + """ + _print({"schema": SCHEMA_VERSION, "command": command, "data": data}) + + +def emit_json_error(command: str, error: Exception) -> None: + """Print the error envelope for a command that raised. + + The process still exits non-zero; the envelope exists so a driving program + gets a structured reason on stdout instead of scraping stderr. + """ + _print( + { + "schema": SCHEMA_VERSION, + "command": command, + "error": {"type": type(error).__name__, "message": str(error)}, + } + ) + + +def _print(envelope: dict[str, Any]) -> None: + print(json.dumps(envelope, indent=2, sort_keys=True, default=_default)) + + +def _default(value: Any) -> Any: + if dataclasses.is_dataclass(value) and not isinstance(value, type): + return dataclasses.asdict(value) + if isinstance(value, Path): + return str(value) + return str(value) diff --git a/visionpack/formats/coco.py b/visionpack/formats/coco.py index fa895c9..6baffb6 100644 --- a/visionpack/formats/coco.py +++ b/visionpack/formats/coco.py @@ -135,7 +135,7 @@ def _process_image( records = annotations_by_image.get(int(image_record["id"]), []) objects: list[ObjectAnnotation] = [] for record in records: - geometry = _geometry_from_record(record, self.project.manifest.task, file_name) + geometry = geometry_from_record(record, self.project.manifest.task, file_name) objects.append( ObjectAnnotation( class_id=category_to_class_id.get(int(record["category_id"]), str(record["category_id"])), @@ -260,7 +260,7 @@ def export_coco( return result -def _geometry_from_record(record: dict[str, Any], task: str, file_name: str) -> Geometry: +def geometry_from_record(record: dict[str, Any], task: str, file_name: str) -> Geometry: """Pick the geometry to keep for a COCO annotation, guided by the task. COCO records often carry several geometries at once (a segmentation plus its diff --git a/visionpack/sources/importer.py b/visionpack/sources/importer.py index dd60bf7..1a008c8 100644 --- a/visionpack/sources/importer.py +++ b/visionpack/sources/importer.py @@ -1,5 +1,6 @@ from __future__ import annotations +import os from collections.abc import Callable from concurrent.futures import ThreadPoolExecutor from contextlib import AbstractContextManager @@ -8,7 +9,7 @@ from typing import Any from visionpack.core.errors import VisionPackError -from visionpack.core.models import Annotation, Asset, utc_now +from visionpack.core.models import Annotation, Asset, ObjectAnnotation, utc_now from visionpack.core.project import Project from visionpack.formats.base import IngestFailure from visionpack.formats.yolo import parse_yolo_label_text @@ -66,14 +67,18 @@ class SourceSyncer: so images already in the store are recognized and not re-added. """ - def __init__(self, project: Project, source: Source) -> None: + def __init__(self, project: Project, source: Source, max_workers: int | None = None) -> None: self.project = project # Manifest paths are project-relative; resolve them against the project # root (not the cwd) so `vp sync` works from any directory. self.source = _rebase_source(source, project.root) + self.max_workers = max_workers # The cloud sink for `copy` mode, built only when one is actually needed # (copy + a declared target); otherwise objects stay in the local CAS. self._target = self._build_target() + # Label texts read once during class inference, replayed at ingest so a + # remote label file never costs two GETs. + self._label_texts: dict[str, str] = {} # -- resolvers ------------------------------------------------------------ @@ -92,21 +97,35 @@ def _build_target(self) -> CloudTarget | None: loc = _rebase_location(Location.parse(self.project.manifest.target), self.project.root) assert loc is not None # parse of a non-empty value is never None target_uri = loc.resolved_uri() - # v1 is same-provider: a server-side copy can't bridge two providers, so - # reject e.g. a local source feeding an s3 target up front. - src_scheme = self._source_scheme() - if src_scheme != scheme_of(target_uri): - raise VisionPackError( - f"Source {self.source.name!r} copies into target {target_uri!r}, but the source " - f"({src_scheme or 'local'}) and target ({scheme_of(target_uri) or 'local'}) are different " - "providers; v1 supports same-provider copy only." - ) - return CloudTarget(base_uri=target_uri, resolver=get_resolver(target_uri, _resolver_options(loc))) + # Same provider -> server-side copy (bytes never transit the client). + # Different providers (local -> s3, s3 -> gcs, ...) -> relay: upload the + # bytes the sync already read for hashing; still a single read. + server_side = self._source_scheme() == scheme_of(target_uri) + return CloudTarget( + base_uri=target_uri, + resolver=get_resolver(target_uri, _resolver_options(loc)), + server_side=server_side, + ) def _source_scheme(self) -> str: loc = self.source.images or self.source.root return scheme_of(loc.resolved_uri()) if loc is not None else "" + def _pool_size(self) -> int | None: + """Worker count for the ingest pool (``--jobs`` wins). + + Object-store throughput is latency-bound, not CPU-bound, so the + CPU-derived executor default undersizes remote syncs on small machines; + remote sources get a floor of 16 concurrent transfers. Local sources + keep the executor default (disk parallelism saturates early). + """ + if self.max_workers is not None: + return self.max_workers + if self._source_scheme() not in ("", "file"): + cpus = os.cpu_count() or 1 + return max(16, min(32, cpus + 4)) + return None + # -- public --------------------------------------------------------------- def plan(self) -> SourcePlan: @@ -164,29 +183,45 @@ def _run_yolo(self, progress: ProgressCallback | None = None) -> SourceSyncSumma images_without_label=result.images_without_label, labels_without_image=len(result.labels_without_image), ) - existing_ids = {asset.id for asset in self.project.index.assets()} # Load every cached probe for this run's images in one query, so the # per-image lookup in the threads below never opens a connection. self.project.index.prime_blob_cache([image.uri for image, _ in result.pairs]) - # Reading bytes, hashing, probing, perceptual-hashing and storing are - # per-image and I/O-bound, so fan them out across threads (matching - # YoloImporter). Index mutation stays on this thread; pool.map preserves - # input order, so the result is deterministic regardless of scheduling. def process( pair: tuple[FileRef, FileRef | None], ) -> tuple[Asset, Annotation | None, int, _CacheWrite | None] | IngestFailure: image_ref, label_ref = pair try: - label_text = label_res.read_bytes(label_ref.uri).decode("utf-8") if label_ref else None + label_text = self._label_text(label_ref, label_res) if label_ref else None origin = label_ref.uri if label_ref else None return self._ingest(image_ref, image_res, label_text, origin, index_to_class_id) except (VisionPackError, OSError) as exc: return IngestFailure(path=image_ref.uri, error=str(exc)) - total = len(result.pairs) - with ThreadPoolExecutor() as pool: - for done, outcome in enumerate(pool.map(process, result.pairs), 1): + self._drain_pool(process, result.pairs, summary, progress) + self._record(summary, images_loc, labels_loc) + return summary + + def _label_text(self, label_ref: FileRef, label_res: Resolver | None) -> str: + # Class inference may already have fetched this label; never GET twice. + cached = self._label_texts.pop(label_ref.uri, None) + if cached is not None: + return cached + assert label_res is not None # a label ref implies a label resolver + return label_res.read_bytes(label_ref.uri).decode("utf-8") + + def _drain_pool(self, process, items, summary: SourceSyncSummary, progress: ProgressCallback | None) -> None: + """Fan ``process`` out over ``items`` and fold outcomes into ``summary``. + + Reading bytes, hashing, probing, perceptual-hashing and storing are + per-item and I/O-bound, so they run across threads. Index mutation + stays on this thread; ``pool.map`` preserves input order, so the result + is deterministic regardless of scheduling. + """ + existing_ids = {asset.id for asset in self.project.index.assets()} + total = len(items) + with ThreadPoolExecutor(max_workers=self._pool_size()) as pool: + for done, outcome in enumerate(pool.map(process, items), 1): if isinstance(outcome, IngestFailure): summary.failures.append(outcome) else: @@ -207,9 +242,6 @@ def process( if progress is not None: progress(done, total) - self._record(summary, images_loc, labels_loc) - return summary - def _ingest( self, image_ref: FileRef, @@ -218,39 +250,49 @@ def _ingest( label_origin: str | None, index_to_class_id: dict[int, str], ) -> tuple[Asset, Annotation | None, int, _CacheWrite | None]: + asset, cache_write = self._ingest_asset(image_ref, image_res) + annotation: Annotation | None = None + object_count = 0 + if label_text is not None: + objects = parse_yolo_label_text( + label_text, label_origin or image_ref.uri, asset.width, asset.height, index_to_class_id + ) + object_count = len(objects) + annotation = Annotation( + id=f"ann_{asset.id}", + asset_id=asset.id, + task=self.project.manifest.task, + format="internal", + objects=objects, + source={"type": "sync", "format": "yolo", "source": self.source.name, "path": label_origin, "imported_at": utc_now()}, + ) + return asset, annotation, object_count, cache_write + + def _ingest_asset(self, image_ref: FileRef, image_res: Resolver) -> tuple[Asset, _CacheWrite | None]: + """Read/probe/store one image and build its ``Asset`` — format-agnostic. + + The shared core of every source format: YOLO wraps it with label + parsing, ImageFolder with a whole-image label, COCO with its record + conversion. All of them inherit the probe cache and the copy-mode / + cloud-target routing for free. + """ probe, cache_write = self._resolve_blob(image_ref, image_res) digest = probe["sha256"] - width, height = probe["width"], probe["height"] - asset_id = f"asset_{digest[:16]}" asset = Asset( - id=asset_id, + id=f"asset_{digest[:16]}", sha256=digest, media_type="image", path=probe["stored_path"], original_path=image_ref.uri, - width=width, - height=height, + width=probe["width"], + height=probe["height"], channels=probe["channels"], format=probe["format"], size_bytes=probe["size_bytes"], phash=probe["phash"], source=self.source.name, ) - - annotation: Annotation | None = None - object_count = 0 - if label_text is not None: - objects = parse_yolo_label_text(label_text, label_origin or image_ref.uri, width, height, index_to_class_id) - object_count = len(objects) - annotation = Annotation( - id=f"ann_{asset_id}", - asset_id=asset_id, - task=self.project.manifest.task, - format="internal", - objects=objects, - source={"type": "sync", "format": "yolo", "source": self.source.name, "path": label_origin, "imported_at": utc_now()}, - ) - return asset, annotation, object_count, cache_write + return asset, cache_write def _resolve_blob(self, image_ref: FileRef, image_res: Resolver) -> tuple[dict, _CacheWrite | None]: """Probe an image, skipping the body read when it is provably unchanged. @@ -288,14 +330,15 @@ def _resolve_blob(self, image_ref: FileRef, image_res: Resolver) -> tuple[dict, def _store_blob(self, image_ref: FileRef, image_res: Resolver, data: bytes, digest: str) -> str: """Materialize the just-read bytes and return the asset ``path``. - ``reference`` keeps no copy; cloud ``copy`` lands the object in the target - CAS server-side (the bytes never transit the client — we only read them - once, above, for the hash); everything else writes the local CAS. + ``reference`` keeps no copy; ``copy`` with a target lands the object in + the target CAS — server-side when source and target share a provider, + otherwise by relaying the bytes we already read for the hash (still one + read total); everything else writes the local CAS. """ if self.source.copy == "reference": return self._reference_path(image_ref, image_res) if self._target is not None: # copy + a declared target - return self._target.ensure_object(image_ref.uri, digest) + return self._target.ensure_object(image_ref.uri, digest, data=data) local = image_res.local_path(image_ref.uri) or Path(image_ref.uri) return self.project.object_store.store(local, digest, self.source.copy, data=data) @@ -341,9 +384,33 @@ def _yolo_class_names( names = found break if not names: - names = _infer_class_names(labels, label_res) + names = self._infer_class_names(labels, label_res) return [self._remap(index, name) for index, name in enumerate(names)] + def _infer_class_names(self, labels: list[FileRef], label_res: Resolver | None) -> list[str]: + """Infer ``class_N`` names from the highest class index used in labels. + + Label bodies are fetched in parallel and cached, so ingest replays them + instead of issuing a second GET per label (labels are tiny; the cache is + drained as ingest consumes it). + """ + if label_res is None or not labels: + return [] + with ThreadPoolExecutor(max_workers=self._pool_size()) as pool: + texts = pool.map(lambda ref: label_res.read_bytes(ref.uri).decode("utf-8"), labels) + self._label_texts = {ref.uri: text for ref, text in zip(labels, texts, strict=True)} + max_class = -1 + for text in self._label_texts.values(): + for line in text.splitlines(): + stripped = line.strip().lstrip("") + if not stripped: + continue + try: + max_class = max(max_class, int(float(stripped.split()[0]))) + except (ValueError, IndexError): + continue + return [f"class_{index}" for index in range(max_class + 1)] + def _explicit_class_names(self) -> list[str]: if self.source.classes is None: return [] @@ -392,10 +459,7 @@ def _run_coco(self, progress: ProgressCallback | None = None) -> SourceSyncSumma images_path = image_res.local_path(images_loc.resolved_uri()) labels_path = label_res.local_path(labels_loc.resolved_uri()) if images_path is None or labels_path is None: - raise VisionPackError( - f"COCO source {self.source.name!r} currently supports local images/labels only " - "(remote COCO arrives with the fsspec backends in Phase 2)." - ) + return self._run_coco_remote(images_loc, labels_loc, image_res, label_res, progress) before = {asset.id for asset in self.project.index.assets()} result = CocoImporter(self.project, labels_path, images_path, copy_mode=self.source.copy).run(progress) added = self._tag_provenance(before) @@ -409,6 +473,98 @@ def _run_coco(self, progress: ProgressCallback | None = None) -> SourceSyncSumma failures=result.failures, ) + def _run_coco_remote( + self, + images_loc: Location, + labels_loc: Location, + image_res: Resolver, + label_res: Resolver, + progress: ProgressCallback | None, + ) -> SourceSyncSummary: + """COCO over any resolver: the JSON is one read, images stream through + the shared blob pipeline (probe cache, copy modes, cloud target).""" + import json + from collections import defaultdict + + from visionpack.core.manifest import class_id_from_name + from visionpack.formats.coco import geometry_from_record + + labels_uri = labels_loc.resolved_uri() + try: + document = json.loads(label_res.read_bytes(labels_uri).decode("utf-8")) + except json.JSONDecodeError as exc: + raise VisionPackError(f"COCO annotation file is not valid JSON: {labels_uri} ({exc})") from exc + if not isinstance(document, dict): + raise VisionPackError(f"COCO annotation file must contain a JSON object at the top level: {labels_uri}") + + categories = {int(cat["id"]): str(cat.get("name", cat["id"])) for cat in document.get("categories", [])} + # Merge classes by (remapped) name, exactly like the local importer. + remapped = {cat_id: self._remap(index, name) for index, (cat_id, name) in enumerate(categories.items())} + classes_added = self.project.manifest.merge_classes(list(remapped.values())) + name_to_id = {item.name: item.id for item in self.project.manifest.classes} + category_to_class_id = { + cat_id: name_to_id.get(name, class_id_from_name(name)) for cat_id, name in remapped.items() + } + + annotations_by_image: dict[int, list[dict]] = defaultdict(list) + for record in document.get("annotations", []): + annotations_by_image[int(record["image_id"])].append(record) + + refs = image_res.list_files(images_loc.resolved_uri(), IMAGE_EXTENSIONS) + # file_name is relative to the images root; fall back to the bare + # basename only when it is unambiguous across the listing. + by_rel = {f"{ref.relkey}{ref.suffix}": ref for ref in refs} + by_name: dict[str, FileRef | None] = {} + for ref in refs: + name = f"{ref.stem}{ref.suffix}" + by_name[name] = None if name in by_name else ref + + def resolve_ref(file_name: str) -> FileRef | None: + rel = file_name.replace("\\", "/").lstrip("./") + return by_rel.get(rel) or by_name.get(rel.rsplit("/", 1)[-1]) + + images = document.get("images", []) + matched = [resolve_ref(str(record.get("file_name"))) for record in images] + self.project.index.prime_blob_cache([ref.uri for ref in matched if ref is not None]) + summary = SourceSyncSummary(name=self.source.name, classes_added=classes_added) + task = self.project.manifest.task + + def process(record: dict) -> tuple[Asset, Annotation | None, int, _CacheWrite | None] | IngestFailure: + file_name = str(record.get("file_name")) + try: + ref = resolve_ref(file_name) + if ref is None: + raise VisionPackError( + f"COCO image not found under {images_loc.resolved_uri()}: file_name={file_name!r} " + f"(image id={record.get('id')})" + ) + asset, cache_write = self._ingest_asset(ref, image_res) + objects = [ + ObjectAnnotation( + class_id=category_to_class_id.get(int(item["category_id"]), str(item["category_id"])), + geometry=geometry_from_record(item, task, file_name), + attributes={"iscrowd": int(item["iscrowd"])} if item.get("iscrowd") else {}, + ) + for item in annotations_by_image.get(int(record["id"]), []) + ] + annotation: Annotation | None = None + if objects: + annotation = Annotation( + id=f"ann_{asset.id}", + asset_id=asset.id, + task=task, + format="internal", + objects=objects, + source={"type": "sync", "format": "coco", "source": self.source.name, "path": labels_uri, "imported_at": utc_now()}, + ) + return asset, annotation, len(objects), cache_write + except (VisionPackError, OSError) as exc: + return IngestFailure(path=file_name, error=str(exc)) + + self._drain_pool(process, images, summary, progress) + self._record(summary, images_loc, labels_loc) + return summary + # -- ImageFolder (classification) ----------------------------------------- def _imagefolder_root(self) -> Location: @@ -443,10 +599,7 @@ def _run_imagefolder(self, progress: ProgressCallback | None = None) -> SourceSy resolver = self._resolver(root) root_path = resolver.local_path(root.resolved_uri()) if root_path is None: - raise VisionPackError( - f"ImageFolder source {self.source.name!r} currently supports a local root only " - "(remote backends arrive with fsspec in Phase 2)." - ) + return self._run_imagefolder_remote(root, resolver, progress) before = {asset.id for asset in self.project.index.assets()} result = ImageFolderImporter(self.project, root_path, copy_mode=self.source.copy).run(progress) added = self._tag_provenance(before) @@ -460,6 +613,47 @@ def _run_imagefolder(self, progress: ProgressCallback | None = None) -> SourceSy failures=result.failures, ) + def _run_imagefolder_remote( + self, root: Location, resolver: Resolver, progress: ProgressCallback | None + ) -> SourceSyncSummary: + """ImageFolder over any resolver: the class is the first path segment + under the root, each image gets a whole-image label (no geometry).""" + refs = resolver.list_files(root.resolved_uri(), IMAGE_EXTENSIONS) + labeled = [(ref, ref.relkey.split("/")[0]) for ref in refs if "/" in ref.relkey] + if not labeled: + raise VisionPackError( + f"No class subdirectories found under {root.resolved_uri()}. " + "Expected layout: //." + ) + raw_names = sorted({name for _, name in labeled}) + remapped = {raw: self._remap(index, raw) for index, raw in enumerate(raw_names)} + classes_added = self.project.manifest.merge_classes(list(remapped.values())) + name_to_id = {item.name: item.id for item in self.project.manifest.classes} + + self.project.index.prime_blob_cache([ref.uri for ref, _ in labeled]) + summary = SourceSyncSummary(name=self.source.name, classes_added=classes_added) + task = self.project.manifest.task + + def process(item: tuple[FileRef, str]) -> tuple[Asset, Annotation | None, int, _CacheWrite | None] | IngestFailure: + ref, raw_name = item + try: + asset, cache_write = self._ingest_asset(ref, resolver) + annotation = Annotation( + id=f"ann_{asset.id}", + asset_id=asset.id, + task=task, + format="internal", + objects=[ObjectAnnotation(class_id=name_to_id[remapped[raw_name]], geometry=None)], + source={"type": "sync", "format": "imagefolder", "source": self.source.name, "path": ref.uri, "imported_at": utc_now()}, + ) + return asset, annotation, 1, cache_write + except (VisionPackError, OSError) as exc: + return IngestFailure(path=ref.uri, error=str(exc)) + + self._drain_pool(process, labeled, summary, progress) + self._record(summary, root, None) + return summary + # -- shared --------------------------------------------------------------- def _tag_provenance(self, before: set[str]) -> int: @@ -575,32 +769,18 @@ def _rebase_source(source: Source, root: Path) -> Source: ) -def _infer_class_names(labels: list[FileRef], label_res: Resolver | None) -> list[str]: - if label_res is None: - return [] - max_class = -1 - for ref in labels: - for line in label_res.read_bytes(ref.uri).decode("utf-8").splitlines(): - stripped = line.strip().lstrip("") - if not stripped: - continue - try: - max_class = max(max_class, int(float(stripped.split()[0]))) - except (ValueError, IndexError): - continue - return [f"class_{index}" for index in range(max_class + 1)] - - def sync_sources( project: Project, source_name: str | None = None, progress_factory: Callable[[str], AbstractContextManager[ProgressCallback | None]] | None = None, + max_workers: int | None = None, ) -> list[SourceSyncSummary]: """Sync the declared sources. ``progress_factory`` (e.g. ``cli_progress``) - yields a fresh progress callback per source so each gets its own bar.""" + yields a fresh progress callback per source so each gets its own bar. + ``max_workers`` overrides the per-source ingest concurrency (``--jobs``).""" summaries: list[SourceSyncSummary] = [] for source in _select_sources(project, source_name): - syncer = SourceSyncer(project, source) + syncer = SourceSyncer(project, source, max_workers=max_workers) if progress_factory is None: summaries.append(syncer.run()) else: diff --git a/visionpack/sources/resolver.py b/visionpack/sources/resolver.py index 1c9e8dd..a51d998 100644 --- a/visionpack/sources/resolver.py +++ b/visionpack/sources/resolver.py @@ -9,6 +9,7 @@ from urllib.request import url2pathname from visionpack.core.errors import VisionPackError +from visionpack.sources.retry import with_retries @dataclass(slots=True) @@ -71,14 +72,25 @@ def stat(self, uri: str) -> ObjectStat: @abstractmethod def read_bytes(self, uri: str) -> bytes: ... + @abstractmethod + def write_bytes(self, uri: str, data: bytes) -> None: + """Write ``data`` to ``uri``, creating parents as the backend requires. + + The upload half of a cross-provider relay: when a server-side copy is + impossible (source and target live on different providers), sync uploads + the bytes it already read for hashing — one read, one write, no second + round-trip. + """ + @abstractmethod def server_copy(self, src_uri: str, dst_uri: str) -> None: """Copy ``src_uri`` to ``dst_uri`` without routing the bytes through us. - Same-provider only (v1): the copy happens inside the backend (S3 + Same-provider only: the copy happens inside the backend (S3 ``CopyObject`` / GCS rewrite / a local filesystem copy), so the client never downloads-then-uploads. Used by ``sync`` to land objects in a - content-addressed target bucket (see docs/SPEC-cloud-sync.md). + content-addressed target bucket; cross-provider transfers go through + :meth:`write_bytes` instead (the relay path). """ @abstractmethod @@ -133,6 +145,11 @@ def stat(self, uri: str) -> ObjectStat: def read_bytes(self, uri: str) -> bytes: return self._path(uri).read_bytes() + def write_bytes(self, uri: str, data: bytes) -> None: + destination = self._path(uri) + destination.parent.mkdir(parents=True, exist_ok=True) + destination.write_bytes(data) + def server_copy(self, src_uri: str, dst_uri: str) -> None: dst = self._path(dst_uri) dst.parent.mkdir(parents=True, exist_ok=True) @@ -155,6 +172,11 @@ class FsspecResolver(Resolver): The provider library (``s3fs``/``gcsfs``/``adlfs``) is an optional extra; it is imported lazily by :func:`get_resolver`, never by the core. Filesystem instances are cached by fsspec itself, so resolving per call is cheap. + + Every network call goes through :func:`with_retries`: transient provider + errors (throttling, dropped connections, 5xx) are retried with exponential + backoff, and exhaustion surfaces as a :class:`VisionPackError` so ingest + records a per-object failure instead of aborting the sync. """ def __init__(self, scheme: str, storage_options: dict[str, Any] | None = None) -> None: @@ -185,17 +207,18 @@ def _stat_from_info(info: dict[str, Any]) -> ObjectStat: def exists(self, uri: str) -> bool: fs, path = self._fs_and_path(uri) - return bool(fs.exists(path)) + return bool(with_retries(f"exists({uri})", lambda: fs.exists(path))) def list_files(self, uri: str, suffixes: set[str] | None = None) -> list[FileRef]: fs, root = self._fs_and_path(uri) - if not fs.exists(root): + if not with_retries(f"exists({uri})", lambda: fs.exists(root)): raise VisionPackError(f"Source location does not exist: {uri}") base = root.rstrip("/") refs: list[FileRef] = [] # detail=True returns size + etag in the same paginated LIST, so a sync # needs no per-object HEAD (the metadata-only listing the spec promises). - for path, info in sorted(fs.find(root, detail=True).items()): + listing = with_retries(f"list({uri})", lambda: fs.find(root, detail=True)) + for path, info in sorted(listing.items()): name = path.rsplit("/", 1)[-1] suffix = "." + name.rsplit(".", 1)[-1].lower() if "." in name else "" if suffixes is not None and suffix not in suffixes: @@ -211,24 +234,28 @@ def list_files(self, uri: str, suffixes: set[str] | None = None) -> list[FileRef def stat(self, uri: str) -> ObjectStat: fs, path = self._fs_and_path(uri) - return self._stat_from_info(fs.info(path)) + return self._stat_from_info(with_retries(f"stat({uri})", lambda: fs.info(path))) def read_bytes(self, uri: str) -> bytes: fs, path = self._fs_and_path(uri) - return fs.cat_file(path) + return with_retries(f"read({uri})", lambda: fs.cat_file(path)) + + def write_bytes(self, uri: str, data: bytes) -> None: + fs, path = self._fs_and_path(uri) + with_retries(f"write({uri})", lambda: fs.pipe_file(path, data)) def server_copy(self, src_uri: str, dst_uri: str) -> None: src_fs, src_path = self._fs_and_path(src_uri) dst_fs, dst_path = self._fs_and_path(dst_uri) - # v1 is same-provider: a cross-provider copy can't be server-side (it - # would have to transit the client), so refuse it loudly rather than - # silently downloading-then-uploading. + # A cross-provider copy can't be server-side (it would have to transit + # the client), so refuse it loudly rather than silently + # downloading-then-uploading — callers relay via write_bytes instead. if type(src_fs) is not type(dst_fs): raise VisionPackError( f"Server-side copy needs source and target on the same provider " - f"(got {src_uri!r} -> {dst_uri!r}); cross-cloud transfer is not supported in v1." + f"(got {src_uri!r} -> {dst_uri!r}); use the relay path (write_bytes) for cross-provider transfer." ) - dst_fs.copy(src_path, dst_path) + with_retries(f"copy({src_uri} -> {dst_uri})", lambda: dst_fs.copy(src_path, dst_path)) def local_path(self, uri: str) -> Path | None: # Remote objects have no local path; ingest works from in-memory bytes. diff --git a/visionpack/sources/retry.py b/visionpack/sources/retry.py new file mode 100644 index 0000000..64cc6ad --- /dev/null +++ b/visionpack/sources/retry.py @@ -0,0 +1,54 @@ +from __future__ import annotations + +import time +from collections.abc import Callable +from typing import TypeVar + +from visionpack.core.errors import VisionPackError + +T = TypeVar("T") + +# Tuned like rclone's low-level retries: enough attempts to ride out a rate +# limit or a dropped connection, short enough that a hard outage fails a sync +# in seconds, not minutes. Module-level so tests (and power users) can adjust. +MAX_ATTEMPTS = 4 +BASE_DELAY_SECONDS = 0.5 + +# Errors that retrying can never fix. FileNotFoundError doubles as fsspec's +# "no such key"; VisionPackError is our own diagnosis, already actionable. +_PERMANENT = ( + FileNotFoundError, + IsADirectoryError, + NotADirectoryError, + PermissionError, + KeyboardInterrupt, + VisionPackError, +) + + +def with_retries(operation: str, fn: Callable[[], T]) -> T: + """Run ``fn``, retrying transient failures with exponential backoff. + + Object-store calls fail transiently all the time (throttling, dropped + connections, 5xx) and every provider library spells those errors + differently, so the classification is by exclusion: anything not provably + permanent is worth retrying — every operation behind this wrapper is + idempotent (reads, listings, full-object writes, server-side copies), so a + retry can duplicate work but never corrupt state. + + When attempts run out the last error is wrapped in a + :class:`VisionPackError` naming the operation, so the ingest loop records a + clean per-object failure instead of crashing the whole sync on an exception + type it doesn't know. + """ + last: Exception | None = None + for attempt in range(MAX_ATTEMPTS): + try: + return fn() + except _PERMANENT: + raise + except Exception as exc: # provider-specific transient errors + last = exc + if attempt + 1 < MAX_ATTEMPTS: + time.sleep(BASE_DELAY_SECONDS * (2**attempt)) + raise VisionPackError(f"{operation} failed after {MAX_ATTEMPTS} attempts: {last}") from last diff --git a/visionpack/sources/target.py b/visionpack/sources/target.py index 4ebb9bd..4e6b4be 100644 --- a/visionpack/sources/target.py +++ b/visionpack/sources/target.py @@ -1,34 +1,96 @@ from __future__ import annotations -from dataclasses import dataclass +import threading +from dataclasses import dataclass, field +from visionpack.core.errors import VisionPackError from visionpack.sources.resolver import Resolver @dataclass(slots=True) class CloudTarget: - """A content-addressed sink objects are copied into, server-side. + """A content-addressed sink objects are copied into. Mirrors the local CAS layout (``objects/sha256///``) in a target bucket so the target is self-sufficient and dedups globally by content: the same image arriving from any source lands on the same key, copied at most once. See docs/SPEC-cloud-sync.md. + + ``server_side`` says whether source and target live on the same provider. + When they do, objects move with a server-side copy (S3 ``CopyObject`` / GCS + rewrite) and never transit the client. When they don't (cross-provider: + S3 -> GCS, local -> S3, ...), the sync *relays* the bytes it already read + for hashing — a single read (needed anyway for the sha256) plus one upload, + never a second download. + + Membership is resolved with one **prefix listing** of the target CAS on + first use (rclone's fast-list approach), not a per-object existence check: + at 100k objects that is ~100 paginated LISTs instead of 100k HEADs. The set + is safe to consult from ingest worker threads. """ base_uri: str resolver: Resolver + server_side: bool = True + _present: set[str] | None = field(default=None, init=False) + _mutex: threading.Lock = field(default_factory=threading.Lock, init=False) def object_uri(self, sha256: str) -> str: return f"{self.base_uri.rstrip('/')}/objects/sha256/{sha256[:2]}/{sha256[2:4]}/{sha256}" - def ensure_object(self, src_uri: str, sha256: str) -> str: - """Server-copy ``src_uri`` into its content-addressed slot, returning it. + def ensure_object(self, src_uri: str, sha256: str, data: bytes | None = None) -> str: + """Land ``src_uri`` in its content-addressed slot, returning that slot. Idempotent: if the slot already holds the object (this run, a prior run, - or another machine), the copy is skipped — content addressing makes the - ``exists`` check a safe dedup, since identical content means identical key. + or another machine), nothing is transferred — content addressing makes + the membership check a safe dedup, since identical content means + identical key. Same-provider transfers are server-side copies; + cross-provider transfers upload ``data`` (the bytes the caller already + read to hash) and verify the landed size. """ destination = self.object_uri(sha256) - if not self.resolver.exists(destination): + if sha256 in self._known(): + return destination + if self.server_side: self.resolver.server_copy(src_uri, destination) + elif data is not None: + self.resolver.write_bytes(destination, data) + self._verify_upload(destination, len(data)) + else: + raise VisionPackError( + f"Cross-provider transfer of {src_uri!r} needs the object bytes to relay, " + "but none were provided." + ) + with self._mutex: + assert self._present is not None # _known() above populated it + self._present.add(sha256) return destination + + def _known(self) -> set[str]: + """The sha256s already present in the target, listed once per sync run. + + Two workers racing on the *same new* hash may both transfer it — the + write is idempotent (same bytes, same key), so that costs one duplicate + upload at worst and never corrupts the CAS. + """ + with self._mutex: + if self._present is None: + prefix = f"{self.base_uri.rstrip('/')}/objects/sha256" + try: + refs = self.resolver.list_files(prefix) + except VisionPackError: + refs = [] # target CAS not created yet: everything is new + # Object keys are bare sha256s (no extension), so stem == hash. + self._present = {ref.stem for ref in refs} + return self._present + + def _verify_upload(self, destination: str, expected_size: int) -> None: + # A relayed upload transits the client, so confirm the provider landed + # every byte before the index starts pointing at the slot. (Server-side + # copies don't need this: the provider guarantees copy integrity.) + landed = self.resolver.stat(destination) + if landed.size != expected_size: + raise VisionPackError( + f"Relay upload to {destination!r} landed {landed.size} bytes, expected {expected_size}; " + "the object was not indexed — re-run the sync." + )