From e1e84fbad3b0f5808d56daae236945a022763a86 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Wed, 23 Sep 2026 14:27:33 -0400 Subject: [PATCH 01/13] Drive the real native sink through to the catalog in the live job Every live suite staged its packs through the conformance_sink driver, which feeds the sink core JSON and base64, so the torch-facing NativePackSink that the ring hands envelopes to never reached a catalog in CI. tests/test_native_capture_chain_live.py drives it with torch CPU tensors (through its attach/submit_envelope test entry points) into a spool, runs CaptureStorageService in-process against ClickHouse and the signature-verifying fake S3, and reads the payloads back through NativeCaptureReader, requiring the bytes to match. No CUDA, no ring, no conformance driver. It also sends one 68 MiB pack up as a multipart upload, twice. The fake S3 accepts parts of any size, takes the part list on trust and never checks the completion ETags, so against it the test asserts the part sizes the client sent. Against MinIO, S3's 5 MiB part floor is enforced by the server; the MinIO case skips only when DMI_MINIO_ENDPOINT is unset. The clickhouse-live job now builds build/_dmi_native_sink and starts MinIO (quay.io image pinned by digest, via docker run because a service container takes no command) and sets that variable, so its skip gate turns a lost variable into a red job. Evidence, local ClickHouse 127.0.0.1 and a MinIO binary extracted from the pinned image on 127.0.0.1:19100: 3 passed; with the endpoint unset, 2 passed, 1 skipped. With the client's multipart chunk shrunk to 4 MiB (not committed), the fake-S3 case failed on the part count (18 != 5) and the MinIO case on CompleteMultipartUpload EntityTooSmall. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 61 +++- tests/test_native_capture_chain_live.py | 393 ++++++++++++++++++++++++ 2 files changed, 449 insertions(+), 5 deletions(-) create mode 100644 tests/test_native_capture_chain_live.py diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index 2b62f483b..a740934e0 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -143,10 +143,11 @@ jobs: # It runs HERE, in the cpu job, and not in the live job, because it # has to collect the WHOLE tree and only this job's environment can # import the whole tree: `make -C native cpu-goals` above builds - # CPU_ONLY_GOALS, which includes build/_dmi_native_sink, while the - # live job builds only the four conformance drivers and so hits a - # module-level skip in the two suites that import that .so. Doing the - # tree walk in the live job is precisely the regression this replaces. + # CPU_ONLY_GOALS, while the live job builds only what its suites + # drive; when this step was written that was the four conformance + # drivers, and the tree walk there hit a module-level skip in the two + # suites that import build/_dmi_native_sink. Doing the tree walk in + # the live job is precisely the regression this replaces. # This job already collects the same tree (`make check` -> `pytest -m # cpu`, no path), so the walk costs seconds and adds no new imports. # @@ -323,6 +324,15 @@ jobs: # test_native_capture_storage_live.py drives in-process. It takes # only pybind11's headers from the torch this job installs and links # nothing from it. + # + # And _dmi_native_sink, the torch-facing NativePackSink the ring + # hands envelopes to. Every other suite here stages packs through + # the conformance_sink DRIVER, JSON and base64 in, so the real + # adapter never reached a catalog in CI; + # test_native_capture_chain_live.py drives it with torch CPU tensors + # through the service into ClickHouse and reads the bytes back. It + # links torch_cpu from the wheel this job installs (C++20, which the + # Makefile already passes for this one target). run: | sudo apt-get update -q sudo apt-get install -y -q libcurl4-openssl-dev @@ -332,13 +342,51 @@ jobs: build/conformance_sink \ build/conformance_spool \ build/_dmi_native_store \ + build/_dmi_native_sink \ CURL_INCDIR=/usr/include/x86_64-linux-gnu \ CURL_LIBDIR=/usr/lib/x86_64-linux-gnu + - name: Start MinIO for the multipart capture case + # A real S3 implementation for the one case the in-repo fake cannot + # judge: a pack of 64 MiB or more goes up as a multipart upload, and + # the fake assembles parts of any size, where S3 refuses every part + # but the last under 5 MiB (EntityTooSmall at CompleteMultipartUpload). + # Verified locally with the minio binary from this image: with the + # client's part shrunk to 4 MiB, the fake-S3 case fails only on its + # own part-size assertion, and the MinIO case on exactly that error. + # + # A `docker run` step, not a `services:` entry: a service container + # takes no command, and this image does nothing without `server + # `. quay.io because the Docker Hub repository is gone; pinned + # by digest, the last release published as an image. Plain HTTP, on + # 9100 because ClickHouse holds 9000. B7a moves it to HTTPS with a CA + # generated at job start. + run: | + docker run -d --name minio -p 9100:9000 \ + -e MINIO_ROOT_USER=dmi-ci-access \ + -e MINIO_ROOT_PASSWORD=dmi-ci-secret-key \ + quay.io/minio/minio:RELEASE.2025-09-07T16-13-09Z@sha256:14cea493d9a34af32f524e538b8346cf79f3321eff8e708c1e2960462bd8936e \ + server /data + for attempt in $(seq 1 30); do + if curl -sf http://127.0.0.1:9100/minio/health/live; then + echo "MinIO is up" + exit 0 + fi + sleep 2 + done + docker logs minio + echo "MinIO did not become healthy" >&2 + exit 1 + - name: Run the ClickHouse live suites env: DMI_CLICKHOUSE_HOST: 127.0.0.1 DMI_CLICKHOUSE_PORT: "9000" + # Unset, the MinIO multipart case skips, and the gate below turns + # that skip into a failure: the variable is load-bearing. + DMI_MINIO_ENDPOINT: http://127.0.0.1:9100 + DMI_MINIO_ACCESS_KEY: dmi-ci-access + DMI_MINIO_SECRET_KEY: dmi-ci-secret-key # `python -m pytest`, matching the Makefile. `not garage` drops the one # suite that also needs a Garage S3 endpoint; everything else under # these markers needs nothing but the server above. @@ -354,7 +402,10 @@ jobs: # collection skip lands in the JUnit as ``, and the gate # below then fails the job -- correctly, by its own doctrine -- over # two files this suite does not even want. Measured in CI: `201 - # collected, 2 skipped`. + # collected, 2 skipped`. (This job has since started building the + # sink, for the capture-chain suite, so those two import here now. + # The glob stays: the argument is about any module that imports an + # artifact this job does not build, not about those two.) # # So the RUN stays scoped to the files it needs to import. The hole # a glob leaves -- a `manual`/`clickhouse` test named outside the diff --git a/tests/test_native_capture_chain_live.py b/tests/test_native_capture_chain_live.py new file mode 100644 index 000000000..b8d1a8049 --- /dev/null +++ b/tests/test_native_capture_chain_live.py @@ -0,0 +1,393 @@ +"""The whole native capture chain on CPU: sink -> spool -> service -> reader. + +The other live suites each start somewhere in the middle. The storage-service +suite stages its packs through the conformance_sink DRIVER, which feeds the +sink core JSON and base64 rather than tensors, and the adapter suite stops at +the spool. Nothing drove the real torch-facing NativePackSink into a real +catalog, so a disagreement between what the adapter stages and what the +service, the indexer and the reader expect had nowhere to show up. + +Here the REAL sink (_dmi_native_sink, entered through its test-only +`attach` / `submit_envelope`, the same `submit` the ring calls) stages torch +CPU tensors into a spool, the in-process CaptureStorageService +(_dmi_native_store) uploads and indexes them into a live ClickHouse, and +NativeCaptureReader reads them back; the payload bytes must equal the tensors +that went in. No CUDA, no ring engine, no conformance driver, no Python on +the data path. + +The object store is the in-repo signature-verifying fake S3, and, when +DMI_MINIO_ENDPOINT names one, a real MinIO as well. The multipart case needs +both: a pack of 64 MiB or more crosses the client's multipart threshold, and +the fake S3 enforces NOTHING about part sizes -- it accepts a part of any +length, takes the part list on trust, and never checks the ETags in the +completion body -- so against it the test checks the parts the client sent +itself. MinIO holds the client to S3's rules: every part but the last is at +least 5 MiB, or CompleteMultipartUpload fails with EntityTooSmall. + +Build: make -C native build/_dmi_native_sink build/_dmi_native_store \\ + PYTHON=/bin/python +""" + +from __future__ import annotations + +import json +import math +import sys +import uuid +from contextlib import contextmanager +from os import environ +from pathlib import Path + +import pytest + +# Module-level so the fake-S3 fixture registers in this module. +from tests.test_native_s3_client import ( # noqa: E402 + ACCESS, BUCKET, REGION, SECRET, STATE, fake_s3, +) + +REPO = Path(__file__).resolve().parents[1] +BUILD = REPO / "native" / "build" +SINK_BUILT = bool(sorted(BUILD.glob("_dmi_native_sink*.so"))) +STORE_BUILT = bool(sorted(BUILD.glob("_dmi_native_store*.so"))) + +# A per-test skipif rather than a module-level skip, and the sink is imported +# inside the tests rather than at module scope: this module is collected by +# the cpu job's live-glob walk too, which must be able to import it. In the +# live job both extensions are built, so this never fires there -- and if it +# did, that job's skip gate would fail it. +pytestmark = [ + pytest.mark.manual, + pytest.mark.clickhouse, + pytest.mark.skipif( + not (SINK_BUILT and STORE_BUILT), + reason="the native sink and store modules are not built; run " + "`make -C native build/_dmi_native_sink build/_dmi_native_store " + "PYTHON=/bin/python`", + ), +] + +CLICKHOUSE_HOST = environ.get("DMI_CLICKHOUSE_HOST", "127.0.0.1") +CLICKHOUSE_HTTP_PORT = int(environ.get("DMI_CLICKHOUSE_HTTP_PORT", "8123")) +DATABASE = environ.get("DMI_CLICKHOUSE_DATABASE", "default") + +MINIO_ENDPOINT = environ.get("DMI_MINIO_ENDPOINT", "") +MINIO_ACCESS = environ.get("DMI_MINIO_ACCESS_KEY", "minioadmin") +MINIO_SECRET = environ.get("DMI_MINIO_SECRET_KEY", "minioadmin") + +LAYOUT = "capture_pack_reference_v1" +# at::ScalarType numeric values (c10/core/ScalarType.h -- stable ABI). +ATEN = {"float16": 5, "bfloat16": 15, "float32": 6, "int64": 4} + +MiB = 1 << 20 +# S3Config's defaults (native/csrc/store/s3_client.h), which the service +# does not override: a pack at or over the threshold goes up in parts. +MULTIPART_THRESHOLD = 64 * MiB +MULTIPART_CHUNK = 16 * MiB +# S3's floor for every part but the last. +S3_MIN_PART = 5 * MiB + + +def _native_sink(): + sys.path.insert(0, str(BUILD)) + try: + import _dmi_native_sink + finally: + sys.path.remove(str(BUILD)) + return _dmi_native_sink + + +def _metadata(index: int, dtype: str, shape: tuple[int, ...]) -> dict: + from dmi.storage.capture import CaptureMetadata + + return CaptureMetadata( + capture_id=f"chain-{index:04d}", tenant_id="t", experiment_id="e", + run_id="r", session_id="s", request_id=f"q{index}", + sequence_id=f"n{index}", model_id="m", model_revision="mr", + adapter_revision=None, capture_policy_version="v", + hook_name="resid_post", layer_number=index % 2, producer_rank=0, + step_number=index, token_start=index, token_end=index + 1, + batch_position=0, dtype=dtype, shape=shape, + captured_at_ns=1_700_000_000_000_000_000 + index, + ).to_mapping() + + +def _dtype_name(tensor) -> str: + return str(tensor.dtype).removeprefix("torch.") + + +def _raw(tensor): + """The tensor's bytes as a flat uint8 tensor, in memory order.""" + import torch + + return tensor.contiguous().view(-1).view(torch.uint8) + + +class _Envelope: + """One ring envelope: a flat payload and the rows that slice it.""" + + def __init__(self): + self.rows: list[dict] = [] + self.parts = [] + self.offset = 0 + self.expected: dict[str, object] = {} + + def add(self, index: int, tensor): + import torch + + # The sink refuses a slice that is not dtype-aligned, as the ring's + # own layout never produces one: pad up to the element width. + pad = -self.offset % tensor.element_size() + if pad: + self.parts.append(torch.zeros(pad, dtype=torch.uint8)) + self.offset += pad + dtype = _dtype_name(tensor) + length = tensor.numel() * tensor.element_size() + self.rows.append({ + "metadata_json": json.dumps( + _metadata(index, dtype, tuple(tensor.shape))), + "offset": self.offset, "length": length, + "dtype": ATEN[dtype], "shape": list(tensor.shape), + }) + self.parts.append(_raw(tensor)) + self.offset += length + self.expected[f"chain-{index:04d}"] = tensor + + def payload(self): + import torch + + return torch.cat(self.parts) + + +def _open_sink(spool_root: Path, **overrides): + config = dict(spool_root=str(spool_root), layout=LAYOUT, + # Only flush seals a pack: the multipart case must land as + # ONE pack, whatever the submission pace on a busy runner. + max_linger_ns=600 * 10**9) + config.update(overrides) + sink = _native_sink().NativePackSink(**config) + return sink, sink.attach() + + +@contextmanager +def _catalog(): + """A fresh table prefix, dropped afterwards through the reference writer.""" + clickhouse_driver = pytest.importorskip("clickhouse_driver") + from dmi.storage.capture.clickhouse_catalog import ( + ClickHouseCatalogConfig, ClickHouseCatalogWriter, + ) + + client = clickhouse_driver.Client( + host=CLICKHOUSE_HOST, + port=int(environ.get("DMI_CLICKHOUSE_PORT", "9000"))) + config = ClickHouseCatalogConfig( + database=DATABASE, table_prefix=f"dmi_chain_{uuid.uuid4().hex}") + try: + yield config.table_prefix + finally: + ClickHouseCatalogWriter(client, config).drop_schema() + + +def _storage_config(endpoint, bucket, access, secret, prefix): + from dmi.storage.native_capture import NativeCaptureStorageConfig + + return NativeCaptureStorageConfig( + s3_endpoint=endpoint, s3_bucket=bucket, s3_region=REGION, + s3_access_key=access, s3_secret_key=secret, + s3_allow_insecure_http=True, clickhouse_host=CLICKHOUSE_HOST, + clickhouse_port=CLICKHOUSE_HTTP_PORT, database=DATABASE, + table_prefix=prefix, poll_interval_s=0.05) + + +def _run_chain(config, spool_root: Path, envelopes, *, sink_overrides=None): + """Service first, as the engine orders it (its start sweeps the spool, + which is safe only while no sink writes there), then the sink; flush + both, and return the snapshots and what the reader reads back.""" + from dmi.storage.native_capture import ( + NativeCaptureReader, NativeCaptureStorage, + ) + + service = NativeCaptureStorage(config, spool_root=str(spool_root), + spool_max_bytes=1 << 40, sweep_spool=True) + service.start() + try: + sink, _lease = _open_sink(spool_root, **(sink_overrides or {})) + for envelope in envelopes: + sink.submit_envelope(LAYOUT, envelope.rows, envelope.payload()) + assert sink.flush_and_wait(120.0) + sink.rethrow_if_failed() + sink_snapshot = sink.snapshot() + service.flush(120.0) + service.rethrow_if_failed() + service_snapshot = service.snapshot() + finally: + service.stop() + + reader = NativeCaptureReader(config) + selection = reader.select(tenant_id="t") + captures = {capture.descriptor["capture_id"]: capture + for capture in reader.read(selection, byte_limit=1 << 30)} + return sink_snapshot, service_snapshot, captures + + +def _assert_bytes_equal(captures, envelopes): + expected = {} + for envelope in envelopes: + expected.update(envelope.expected) + assert sorted(captures) == sorted(expected) + for capture_id, tensor in expected.items(): + capture = captures[capture_id] + assert capture.payload == _raw(tensor).numpy().tobytes(), capture_id + assert capture.descriptor["shape"] == tuple(tensor.shape) + assert capture.descriptor["dtype"] == _dtype_name(tensor) + + +def _large_envelopes(count: int, elements: int): + """`count` float32 captures of `elements` each, one envelope apiece.""" + import torch + + envelopes = [] + for index in range(count): + envelope = _Envelope() + generator = torch.Generator().manual_seed(index) + envelope.add(index, torch.randn(elements, generator=generator)) + envelopes.append(envelope) + return envelopes + + +def test_the_real_sink_reaches_the_catalog_and_reads_back_exactly( + fake_s3, tmp_path): + import torch + + # Several rows slicing one payload, as the ring presents them, in the + # dtypes a model captures; three envelopes, so more than one pack. + envelopes = [] + for batch in range(3): + envelope = _Envelope() + base = batch * 4 + envelope.add(base, torch.arange(6, dtype=torch.float16).reshape(2, 3) + + base) + envelope.add(base + 1, torch.linspace(-1, 1, 8, + dtype=torch.bfloat16)) + envelope.add(base + 2, torch.randn(4, 5, generator=torch.Generator() + .manual_seed(batch))) + envelope.add(base + 3, torch.arange(3, dtype=torch.int64) * base) + envelopes.append(envelope) + + spool_root = tmp_path / "spool" + with _catalog() as prefix: + config = _storage_config(fake_s3, BUCKET, ACCESS, SECRET, prefix) + sink_snapshot, snapshot, captures = _run_chain( + config, spool_root, envelopes, + sink_overrides={"max_pack_records": 5}) + + assert sink_snapshot["persisted_records"] == 12, sink_snapshot + assert sink_snapshot["dropped_records"] == 0, sink_snapshot + assert sink_snapshot["failures"] == 0, sink_snapshot + assert snapshot["uploaded_packs"] == 3, snapshot # 5 + 5 + 2 records + assert snapshot["indexed_packs"] == 3, snapshot + assert snapshot["indexed_rows"] == 12, snapshot + assert snapshot["pending_index"] == 0, snapshot + assert list(spool_root.rglob("*.dmi-pack.ready")) == [] + _assert_bytes_equal(captures, envelopes) + for capture_id, capture in captures.items(): + expected = next(e.expected[capture_id] for e in envelopes + if capture_id in e.expected) + assert torch.equal(capture.tensor(), expected), capture_id + + +def _one_large_pack(captures, expected_records): + packs = {c.descriptor["object_key"]: c.descriptor for c in + captures.values()} + assert len(packs) == 1, sorted(packs) + ((key, descriptor),) = packs.items() + assert descriptor["pack_record_count"] == expected_records + assert descriptor["object_bytes"] >= MULTIPART_THRESHOLD, descriptor + return key, descriptor["object_bytes"] + + +# 17 captures of 4 MiB: one pack just over the 64 MiB threshold, so five +# parts -- four full 16 MiB chunks and a short last one. +LARGE_COUNT = 17 +LARGE_ELEMENTS = MiB # float32: 4 MiB each +LARGE_SINK = {"max_pack_bytes": 128 * MiB, "max_queue_bytes": 256 * MiB, + "max_pack_records": 10_000} + + +def test_a_multipart_pack_through_the_fake_s3(fake_s3, tmp_path): + envelopes = _large_envelopes(LARGE_COUNT, LARGE_ELEMENTS) + with _catalog() as prefix: + config = _storage_config(fake_s3, BUCKET, ACCESS, SECRET, prefix) + sink_snapshot, snapshot, captures = _run_chain( + config, tmp_path / "spool", envelopes, sink_overrides=LARGE_SINK) + + assert sink_snapshot["persisted_records"] == LARGE_COUNT, sink_snapshot + assert sink_snapshot["dropped_records"] == 0, sink_snapshot + assert snapshot["uploaded_packs"] == 1, snapshot + assert snapshot["indexed_rows"] == LARGE_COUNT, snapshot + key, object_bytes = _one_large_pack(captures, LARGE_COUNT) + + # The fake assembles whatever parts it is sent, so the S3 rules are + # checked here, on what the client actually sent. + with STATE.lock: + calls = list(STATE.calls) + parts = [call["body_len"] for call in calls + if call["method"] == "PUT" and "partNumber=" in call["path"]] + assert len(parts) == math.ceil(object_bytes / MULTIPART_CHUNK), parts + assert sum(parts) == object_bytes + assert all(size == MULTIPART_CHUNK for size in parts[:-1]), parts + assert all(size >= S3_MIN_PART for size in parts[:-1]), parts + assert STATE.objects[key]["etag"].endswith('-multipart"') + _assert_bytes_equal(captures, envelopes) + + +@contextmanager +def _minio_bucket(): + """A fresh bucket on the MinIO endpoint, emptied and removed afterwards.""" + import botocore.session + + client = botocore.session.get_session().create_client( + "s3", endpoint_url=MINIO_ENDPOINT, region_name=REGION, + aws_access_key_id=MINIO_ACCESS, aws_secret_access_key=MINIO_SECRET) + bucket = f"dmi-chain-{uuid.uuid4().hex[:16]}" + client.create_bucket(Bucket=bucket) + try: + yield client, bucket + finally: + for page in client.get_paginator("list_objects_v2").paginate( + Bucket=bucket): + for item in page.get("Contents", []): + client.delete_object(Bucket=bucket, Key=item["Key"]) + for upload in client.list_multipart_uploads(Bucket=bucket).get( + "Uploads", []): + client.abort_multipart_upload( + Bucket=bucket, Key=upload["Key"], UploadId=upload["UploadId"]) + client.delete_bucket(Bucket=bucket) + + +# A skip allowed ONLY for the endpoint being unset: CI's live job starts a +# MinIO and sets it, and that job's gate fails any skip, so a lost variable +# there is a red job rather than a silent gap. +@pytest.mark.skipif( + not MINIO_ENDPOINT, + reason="DMI_MINIO_ENDPOINT is unset; the live CI job starts MinIO and " + "sets it") +def test_a_multipart_pack_through_minio(tmp_path): + envelopes = _large_envelopes(LARGE_COUNT, LARGE_ELEMENTS) + with _minio_bucket() as (client, bucket), _catalog() as prefix: + config = _storage_config(MINIO_ENDPOINT, bucket, MINIO_ACCESS, + MINIO_SECRET, prefix) + sink_snapshot, snapshot, captures = _run_chain( + config, tmp_path / "spool", envelopes, sink_overrides=LARGE_SINK) + + assert sink_snapshot["persisted_records"] == LARGE_COUNT, sink_snapshot + assert snapshot["uploaded_packs"] == 1, snapshot + assert snapshot["indexed_rows"] == LARGE_COUNT, snapshot + key, object_bytes = _one_large_pack(captures, LARGE_COUNT) + + # A multipart object's ETag is "-": the + # pack really went up in parts, and MinIO accepted every one of them. + head = client.head_object(Bucket=bucket, Key=key) + assert head["ContentLength"] == object_bytes + parts = math.ceil(object_bytes / MULTIPART_CHUNK) + assert head["ETag"].strip('"').endswith(f"-{parts}"), head["ETag"] + _assert_bytes_equal(captures, envelopes) From d6458e0a994d5095fe12dcfbce579336196d06ab Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Wed, 23 Sep 2026 14:36:05 -0400 Subject: [PATCH 02/13] Compile the full native backend and the ring tests in CI Nothing in CI compiled _native_backend: the cpu job builds only the CPU-only goals and dry-runs `host`, so the .cu sources, bindings.cpp and the torch/CUDA link could break unnoticed, as the C++20 floor of torch 2.14 did (#138). The new native-backend-compile job builds clickhouse-cpp (docs/install.md step 5), runs `make -C native all`, imports the result, compiles all six tests/native/ring binaries and runs the two that need no device. It installs the CUDA 13.0 pieces from NVIDIA's apt repository on the runner instead of using an nvidia/cuda devel container. A container job pulls its image before any step can free disk, and that image is 3.68 GiB compressed before the ~4.4 GB that torch 2.14+cu130 takes installed. The package list (nvcc, cudart-dev, cublas/cusparse/cusolver dev, about 2.7 GB by Installed-Size) is what a real build reads: the backend's -MMD files and `nvcc -M` of the ring tests, mapped with dpkg -S. #138 is not in this base, so a step makes #138's three -std=c++17 -> c++20 edits with sed; after #138 merges it is a no-op and should be deleted. Local evidence: the sed step run against this Makefile changes exactly lines 125, 130 and 175 (the #138 diff) and against #138's head changes nothing; with those edits, CUDA_HOME=/usr/local/cuda-13.0, SM_ARCH=sm_89 and the libcuda stub on LIBRARY_PATH, `make all` built and passed check-link, and the backend imported with no GPU visible. The stock flags stop at bindings.cpp on "#error C++20 or later compatible compiler is required". Ring tests: all six built with CXX_STD=c++20, and test_record_consumer (14 passed) and test_clickhouse_record_sink (21 passed) ran with CUDA_VISIBLE_DEVICES empty. Not verified until CI runs: the apt install on ubuntu-24.04, disk headroom, gcc 13 (local is 11.4), and the job as a whole. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 148 ++++++++++++++++++++++++++++ 1 file changed, 148 insertions(+) diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index a740934e0..95563f21d 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -452,3 +452,151 @@ jobs: print(f" skipped: {case.get('name')}: {mark.get('message')}") sys.exit(f"{skipped} live test(s) skipped; a skip is not a pass") EOF + + native-backend-compile: + # The full ring backend, _native_backend, which nothing in CI compiled: + # the cpu job builds CPU_ONLY_GOALS, and the one build test there + # (tests/test_cpu_native_build.py) dry-runs `host`. So the .cu sources, + # bindings.cpp and the torch/CUDA link could break on any PR and stay + # broken until someone built on a GPU host -- which is how the C++20 + # floor torch 2.14 imposes went unnoticed (#138). This job compiles and + # links `make -C native all`, imports the result, and builds the ring + # test binaries, running the two that need no device. + # + # No GPU is needed to COMPILE, only the toolkit, and this runner has + # none: everything that needs a device stays with the GPU job (D0). + # + # On the runner, not in an nvidia/cuda:*-devel container. A container + # job pulls its image before the first step runs, so nothing can free + # disk first, and the image is 3.68 GiB COMPRESSED + # (nvidia/cuda:13.0.1-devel-ubuntu24.04, amd64, per its registry + # manifest) -- several times that unpacked -- before the ~4.4 GB that + # torch 2.14+cu130 and its nvidia-* wheels take installed (measured in a + # cu130 venv). That does not reliably fit the ~14 GB a hosted runner + # guarantees. Installing only the packages the build reads is ~2.7 GB: + # the header set comes from the -MMD dependency files of a real build, + # mapped to packages with dpkg -S on a host whose toolkit came from the + # same apt repository, and the sizes are those packages' Installed-Size. + runs-on: ubuntu-24.04 + timeout-minutes: 60 + env: + # Explicit, so the Makefile's resolver takes this toolkit rather than + # searching; it still checks the major against torch.version.cuda. + CUDA_HOME: /usr/local/cuda-13.0 + # Explicit because `native`, the default, asks the GPU this runner does + # not have. Ada, the architecture the GPU hosts build for. + SM_ARCH: sm_89 + + steps: + - name: Free runner disk + # The two largest preinstalled trees this job never uses (the + # Android SDK and .NET); `df` either side so a disk failure is + # readable from the log. + run: | + df -h / + sudo rm -rf /usr/local/lib/android /usr/share/dotnet + df -h / + + - uses: actions/checkout@v4 + + - name: Fetch the clickhouse-cpp submodule + # Only this one. `make all` links clickhouse-cpp and needs none of the + # framework submodules, and `recursive` would also clone + # DMI-Megatron-Integration's nested Megatron-LM onto a disk this job + # is rationing. clickhouse-cpp has no submodules of its own. + run: git submodule update --init third_party/clickhouse-cpp + + - uses: actions/setup-python@v5 + with: + python-version: "3.11" + # No pip cache: the torch CUDA wheel set is gigabytes, and caching + # it would spend the repository's cache quota to save a download. + + - name: Install the CUDA 13.0 compiler and the headers torch includes + # NVIDIA's apt repository, and only what the build reads: nvcc + # (pulling cccl, crt and nvvm), cudart-dev (pulling driver-dev, whose + # stubs/libcuda.so satisfies the backend's -lcuda with no driver + # installed), and cublas, cusparse and cusolver for their headers, + # which ATen's CUDA context includes. 13.0 because torch 2.14's CUDA + # build here is cu130 and the resolver requires the same major. + run: | + curl -fsSLO https://developer.download.nvidia.com/compute/cuda/repos/ubuntu2404/x86_64/cuda-keyring_1.1-1_all.deb + sudo dpkg -i cuda-keyring_1.1-1_all.deb + sudo apt-get update -q + sudo apt-get install -y -q --no-install-recommends \ + cuda-nvcc-13-0 cuda-cudart-dev-13-0 \ + libcublas-dev-13-0 libcusparse-dev-13-0 libcusolver-dev-13-0 + "$CUDA_HOME/bin/nvcc" --version + df -h / + + - name: Install the CUDA build of torch + # 2.14, the release the GPU hosts run (the plan's pin), from the cu130 + # index. A CUDA build is required, not the CPU wheel the other jobs + # use: the Makefile reads torch.version.cuda to pick the toolkit and + # links libtorch_cuda and libc10_cuda. + run: | + python -m pip install --upgrade pip + python -m pip install --no-cache-dir "torch==2.14.*" \ + --index-url https://download.pytorch.org/whl/cu130 + python -c "import torch; print(torch.__version__, torch.version.cuda)" + df -h / + + - name: Build the ClickHouse C++ client + # docs/install.md step 5, verbatim. + run: | + cmake -S third_party/clickhouse-cpp -B third_party/clickhouse-cpp/build \ + -DCMAKE_BUILD_TYPE=Release \ + -DCMAKE_POSITION_INDEPENDENT_CODE=ON + cmake --build third_party/clickhouse-cpp/build -j"$(nproc)" + + - name: "Compile the torch-including targets as C++20 until #138 lands" + # torch 2.14's headers open with `#error C++20 or later compatible + # compiler is required` (reproduced: the stock flags stop at + # bindings.cpp with exactly that). #138 moves CXXFLAGS, HOST_CXXFLAGS + # and NVCCFLAGS to -std=c++20; until it merges, this makes the same + # three edits to the checkout -- they are the lines #138 changes and + # nothing else. Once it has merged, the sed matches nothing, the diff + # below is empty, and this step should be deleted. The Makefile's + # variables are `:=` assignments, so a command-line override would + # replace the whole flag set rather than the standard; hence sed. + run: | + sed -i \ + -e '/^CXXFLAGS :=/s/-std=c++17/-std=c++20/' \ + -e '/^HOST_CXXFLAGS :=/s/-std=c++17/-std=c++20/' \ + -e '/^NVCCFLAGS :=/,/-std=/s/-std=c++17/-std=c++20/' \ + native/Makefile + git diff native/Makefile + test "$(grep -cE '^(HOST_)?CXXFLAGS := -std=c\+\+20|^ -std=c\+\+20 \\$' native/Makefile)" = 3 + + - name: Compile the full native backend + # LIBRARY_PATH puts the toolkit's libcuda stub where the linker finds + # -lcuda; the result does not end up NEEDing libcuda (checked with + # readelf on a local build: the link is --as-needed and nothing calls + # the driver API), so check-link's ldd has no driver to find either. + env: + LIBRARY_PATH: /usr/local/cuda-13.0/targets/x86_64-linux/lib/stubs + run: make -C native -j"$(nproc)" all PYTHON=python + + - name: Import the backend + # A shared library links with undefined symbols; only loading it + # proves the link. No device is needed to load it. + run: | + PYTHONPATH=src python -c " + from dmi.transport import native + module = native._load_extension() + print(module.__file__) + assert hasattr(module, 'RecordSink'), 'no RecordSink in the backend' + " + + - name: Build the ring tests and run the host-only ones + # All six binaries are compiled, so a .cu test that no longer builds + # fails here. Only test_record_consumer and + # test_clickhouse_record_sink RUN: they use ATen CPU tensors and no + # CUDA call, and pass with no device visible (14 and 21 checks, + # measured with CUDA_VISIBLE_DEVICES empty). The other four launch + # kernels and belong to the GPU job. CXX_STD because this Makefile's + # default is c++17, which torch 2.14's headers refuse. + run: | + make -C tests/native/ring -j"$(nproc)" all CXX_STD=c++20 PYTHON=python + tests/native/ring/build/test_record_consumer + tests/native/ring/build/test_clickhouse_record_sink From 177a6e322f0537cee04c18cb23bb32b1fbab5f96 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Wed, 23 Sep 2026 14:37:09 -0400 Subject: [PATCH 03/13] Make src/dmi pass ruff's undefined-name check ruff F821 over src/dmi reported 14 names, in two groups, and neither is a NameError today: - config.py annotates capture_sink_config and capture_storage_config with the strings "NativeSinkConfig" and "NativeCaptureStorageConfig" and imports neither. With `from __future__ import annotations` nothing evaluates them at runtime, but typing.get_type_hints(MonitoringConfig) raises NameError on the first one (reproduced; nothing in src/ calls it today). They are now imported under TYPE_CHECKING, so a static tool resolves them and the module still loads no storage package at runtime (checked: importing dmi.config loads no dmi.storage module). get_type_hints still raises, as before; B2's move of NativeSinkConfig into live code is where a runtime import can go. - hooks/specs.py compares against HOOK_TYPE_Q and eleven siblings that the module binds through a globals() loop over _HOOK_DEFS, which a linter cannot see. Each of those ten lines now carries its own `# noqa: F821` with a comment saying why: per line rather than per file, so the rest of the module stays checked. A typo inside one of those ten lines would now be suppressed too; that is the cost. No behaviour change. ruff 0.16.5 `check --select F821 src/dmi`: 14 errors before, "All checks passed!" after; tests/test_tp_shapes.py, test_hook_spec_flags.py and test_config.py: 63 passed. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- src/dmi/config.py | 9 ++++++++- src/dmi/hooks/specs.py | 24 ++++++++++++++---------- 2 files changed, 22 insertions(+), 11 deletions(-) diff --git a/src/dmi/config.py b/src/dmi/config.py index 9e0c6547f..e68c7ec3e 100644 --- a/src/dmi/config.py +++ b/src/dmi/config.py @@ -5,7 +5,14 @@ from dataclasses import dataclass, field import sys import warnings -from typing import Literal, Optional, get_args +from typing import TYPE_CHECKING, Literal, Optional, get_args + +if TYPE_CHECKING: + # For the string annotations on MonitoringConfig only. Importing them at + # runtime would make this dependency-free module load the storage + # packages; the engine imports them where it validates the values. + from .storage.capture.native_sink import NativeSinkConfig + from .storage.native_capture import NativeCaptureStorageConfig StorageBackend = Literal[ diff --git a/src/dmi/hooks/specs.py b/src/dmi/hooks/specs.py index 4667b71af..f9efe3670 100644 --- a/src/dmi/hooks/specs.py +++ b/src/dmi/hooks/specs.py @@ -237,33 +237,37 @@ def compute_hook_shape( if hook_type in _HIDDEN_DIM_TYPES: return b + [q_len, cfg.hidden_dim] - if hook_type == HOOK_TYPE_Q: + # The HOOK_TYPE_* names are bound by the globals() loop over + # _HOOK_DEFS above, which static analysis cannot follow: hence the + # F821 suppressions, one per line rather than one for the file, so the + # rest of the module stays checked. + if hook_type == HOOK_TYPE_Q: # noqa: F821 return b + [q_len, cfg.num_heads // tp, cfg.head_dim] - if hook_type in (HOOK_TYPE_K, HOOK_TYPE_V): + if hook_type in (HOOK_TYPE_K, HOOK_TYPE_V): # noqa: F821 kv_heads = max(1, cfg.num_kv_heads // tp) # GQA: may replicate return b + [q_len, kv_heads, cfg.head_dim] - if hook_type == HOOK_TYPE_Z: + if hook_type == HOOK_TYPE_Z: # noqa: F821 # Packed/flattened convention flattens heads into a single # trailing dim -> [q_len, num_heads * head_dim]. # Batched convention keeps four dims -> [batch, q_len, num_heads, head_dim]. if batch == 0: return [q_len, (cfg.num_heads // tp) * cfg.head_dim] return b + [q_len, cfg.num_heads // tp, cfg.head_dim] - if hook_type in (HOOK_TYPE_ATTN_SCORES, HOOK_TYPE_PATTERN): + if hook_type in (HOOK_TYPE_ATTN_SCORES, HOOK_TYPE_PATTERN): # noqa: F821 return b + [cfg.num_heads // tp, q_len, kv_dim] - if hook_type == HOOK_TYPE_MLP_POST: + if hook_type == HOOK_TYPE_MLP_POST: # noqa: F821 if cfg.intermediate_dim == 0: return [] # intermediate_dim unknown -- skip this hook return b + [q_len, cfg.intermediate_dim // tp] - if hook_type == HOOK_TYPE_ROUTER_LOGITS: + if hook_type == HOOK_TYPE_ROUTER_LOGITS: # noqa: F821 return (b + [q_len, cfg.num_experts]) if cfg.num_experts > 0 else [] - if hook_type == HOOK_TYPE_TOPK_IDS: + if hook_type == HOOK_TYPE_TOPK_IDS: # noqa: F821 return (b + [q_len, cfg.top_k]) if cfg.top_k > 0 else [] - if hook_type == HOOK_TYPE_TOPK_WEIGHTS: + if hook_type == HOOK_TYPE_TOPK_WEIGHTS: # noqa: F821 return (b + [q_len, cfg.top_k]) if cfg.top_k > 0 else [] - if hook_type == HOOK_TYPE_TOKEN_IDS: + if hook_type == HOOK_TYPE_TOKEN_IDS: # noqa: F821 return b + [q_len] - if hook_type == HOOK_TYPE_FINAL_LOGITS: + if hook_type == HOOK_TYPE_FINAL_LOGITS: # noqa: F821 # compute_logits returns fewer rows than q_len when the framework # only materializes the last-token logits per request. # From 75414d6dd26f4be869bccfdcc9517ed4d6474612 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Wed, 23 Sep 2026 14:37:28 -0400 Subject: [PATCH 04/13] Lint src/dmi for undefined names in the cpu job An undefined name is a NameError that waits for its line to run, and the cpu job had nothing that would see one first: `make check` compiles with `python -m compileall`, which accepts any name. The new step runs ruff's F821 rule, and only that rule, over src/dmi, pinned to ruff 0.16.5 with GitHub annotations. The 14 existing hits are handled in the previous commit, in the source, not excluded here. Local evidence with ruff 0.16.5: the step's command passes on this tree, and with an undefined name appended to src/dmi/config.py (not committed) it exits 1 naming the file and line. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index 95563f21d..cba87dad0 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -46,6 +46,23 @@ jobs: python -m pip install "boto3>=1.40,<2" python -m pip install --no-deps -e . + - name: Fail on an undefined name in src/dmi + # An undefined name is a NameError that waits for its line to run, + # and nothing here would see it first: `make check`'s compile stage + # is `python -m compileall`, which accepts any name, and a branch + # the tests do not reach never raises. ruff's F821 finds it + # statically. That one rule only, not a style gate. + # + # Enabling it found 14, none of them a NameError today, and each is + # handled in the source where a reviewer sees it rather than + # excluded here: config.py's string annotations now import their + # types under TYPE_CHECKING, and hooks/specs.py's HOOK_TYPE_* names, + # bound by a globals() loop no linter can follow, carry a per-line + # noqa saying so. Pinned, so a ruff release cannot move the gate. + run: | + python -m pip install "ruff==0.16.5" + ruff check --select F821 --no-cache --output-format=github src/dmi + - name: Compile every CPU-only native target # The C++ was compiled NOWHERE in CI until this step. `make check`'s # compile stage is `python -m compileall` -- Python only -- and the From 9da4cfeaf7b3ce8325c4934051b33a41c3ff4fb5 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Wed, 23 Sep 2026 19:05:29 -0400 Subject: [PATCH 05/13] Check every catalog metadata field in the capture chain test The chain test compared only capture_id, payload bytes, shape and dtype, so a metadata field the torch sink's parser mangled went unnoticed. A reviewer changed native/csrc/sink/record_row.cpp to store step_number + 7 and the chain test still passed, as did every CPU suite. This chain is the only CI test that reaches record_row.cpp, so it has to catch that. Each capture's read-back descriptor is now compared field by field against the CaptureMetadata the test submitted for it: all 21 fields, with shape as a tuple and a NULL adapter_revision as the empty string the native reader returns for it. The integer fields now also take distinct values (token_start and batch_position vary apart from step_number), so a parser that wires one field into another's slot fails too. producer_rank stays fixed, at 1, because the sink seals a pack when the rank changes. The fake-S3 multipart check counted every logged PUT carrying a partNumber, so one transient retry would have failed the part count. It now keys parts by part number, keeping the last attempt. Evidence, from a scratch copy of HEAD built with make build/_dmi_native_sink build/_dmi_native_store: - step_number + 7 mutant: 2 failed, 1 skipped (MinIO unset); both fail on {'step_number': 7} != {'step_number': 0} for chain-0000. - unmutated rebuild: 2 passed, 1 skipped. Same in this worktree. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- tests/test_native_capture_chain_live.py | 64 ++++++++++++++++++++----- 1 file changed, 51 insertions(+), 13 deletions(-) diff --git a/tests/test_native_capture_chain_live.py b/tests/test_native_capture_chain_live.py index b8d1a8049..95c578346 100644 --- a/tests/test_native_capture_chain_live.py +++ b/tests/test_native_capture_chain_live.py @@ -37,6 +37,7 @@ from contextlib import contextmanager from os import environ from pathlib import Path +from urllib.parse import parse_qs, urlsplit import pytest @@ -97,6 +98,11 @@ def _native_sink(): def _metadata(index: int, dtype: str, shape: tuple[int, ...]) -> dict: + """The capture's metadata. The integer fields differ from one another + across the records, so a parser that wires one field into another's + slot reads back wrong rather than coincidentally right. producer_rank + alone stays fixed: the sink seals a pack when the rank changes, and the + pack counts below assume one rank.""" from dmi.storage.capture import CaptureMetadata return CaptureMetadata( @@ -104,9 +110,10 @@ def _metadata(index: int, dtype: str, shape: tuple[int, ...]) -> dict: run_id="r", session_id="s", request_id=f"q{index}", sequence_id=f"n{index}", model_id="m", model_revision="mr", adapter_revision=None, capture_policy_version="v", - hook_name="resid_post", layer_number=index % 2, producer_rank=0, - step_number=index, token_start=index, token_end=index + 1, - batch_position=0, dtype=dtype, shape=shape, + hook_name="resid_post", layer_number=index % 2, + producer_rank=1, step_number=index, token_start=100 + index, + token_end=102 + index, batch_position=index % 5, dtype=dtype, + shape=shape, captured_at_ns=1_700_000_000_000_000_000 + index, ).to_mapping() @@ -130,6 +137,7 @@ def __init__(self): self.parts = [] self.offset = 0 self.expected: dict[str, object] = {} + self.metadata: dict[str, dict] = {} def add(self, index: int, tensor): import torch @@ -142,15 +150,16 @@ def add(self, index: int, tensor): self.offset += pad dtype = _dtype_name(tensor) length = tensor.numel() * tensor.element_size() + metadata = _metadata(index, dtype, tuple(tensor.shape)) self.rows.append({ - "metadata_json": json.dumps( - _metadata(index, dtype, tuple(tensor.shape))), + "metadata_json": json.dumps(metadata), "offset": self.offset, "length": length, "dtype": ATEN[dtype], "shape": list(tensor.shape), }) self.parts.append(_raw(tensor)) self.offset += length - self.expected[f"chain-{index:04d}"] = tensor + self.expected[metadata["capture_id"]] = tensor + self.metadata[metadata["capture_id"]] = metadata def payload(self): import torch @@ -229,16 +238,38 @@ def _run_chain(config, spool_root: Path, envelopes, *, sink_overrides=None): return sink_snapshot, service_snapshot, captures -def _assert_bytes_equal(captures, envelopes): - expected = {} +def _catalog_form(name: str, value): + """A submitted metadata value as NativeCaptureReader returns it: shape + as a tuple, and the Nullable adapter_revision's NULL as the empty + string (the native reader's sentinel for it; see parse_tsv_tuple).""" + if name == "shape": + return tuple(value) + if name == "adapter_revision" and value is None: + return "" + return value + + +def _assert_read_back_exactly(captures, envelopes): + """Every capture reads back with its bytes AND every metadata field it + was submitted with. This chain is the only CI test that drives the torch + sink's metadata parser (native/csrc/sink/record_row.cpp), so a field it + mangles has to show up here.""" + expected, metadata = {}, {} for envelope in envelopes: expected.update(envelope.expected) + metadata.update(envelope.metadata) assert sorted(captures) == sorted(expected) for capture_id, tensor in expected.items(): capture = captures[capture_id] assert capture.payload == _raw(tensor).numpy().tobytes(), capture_id assert capture.descriptor["shape"] == tuple(tensor.shape) assert capture.descriptor["dtype"] == _dtype_name(tensor) + # Every submitted field, not a chosen few: a field the reader stops + # returning is a KeyError here rather than a silent pass. + submitted = metadata[capture_id] + read_back = {name: capture.descriptor[name] for name in submitted} + assert read_back == {name: _catalog_form(name, value) + for name, value in submitted.items()}, capture_id def _large_envelopes(count: int, elements: int): @@ -288,7 +319,7 @@ def test_the_real_sink_reaches_the_catalog_and_reads_back_exactly( assert snapshot["indexed_rows"] == 12, snapshot assert snapshot["pending_index"] == 0, snapshot assert list(spool_root.rglob("*.dmi-pack.ready")) == [] - _assert_bytes_equal(captures, envelopes) + _assert_read_back_exactly(captures, envelopes) for capture_id, capture in captures.items(): expected = next(e.expected[capture_id] for e in envelopes if capture_id in e.expected) @@ -330,14 +361,21 @@ def test_a_multipart_pack_through_the_fake_s3(fake_s3, tmp_path): # checked here, on what the client actually sent. with STATE.lock: calls = list(STATE.calls) - parts = [call["body_len"] for call in calls - if call["method"] == "PUT" and "partNumber=" in call["path"]] + # By part number, the last attempt at each winning: the log holds every + # HTTP attempt, and a retried part would otherwise count twice. + by_number = {} + for call in calls: + query = parse_qs(urlsplit(call["path"]).query) + if call["method"] == "PUT" and "partNumber" in query: + by_number[int(query["partNumber"][0])] = call["body_len"] + assert sorted(by_number) == list(range(1, len(by_number) + 1)), by_number + parts = [by_number[number] for number in sorted(by_number)] assert len(parts) == math.ceil(object_bytes / MULTIPART_CHUNK), parts assert sum(parts) == object_bytes assert all(size == MULTIPART_CHUNK for size in parts[:-1]), parts assert all(size >= S3_MIN_PART for size in parts[:-1]), parts assert STATE.objects[key]["etag"].endswith('-multipart"') - _assert_bytes_equal(captures, envelopes) + _assert_read_back_exactly(captures, envelopes) @contextmanager @@ -390,4 +428,4 @@ def test_a_multipart_pack_through_minio(tmp_path): assert head["ContentLength"] == object_bytes parts = math.ceil(object_bytes / MULTIPART_CHUNK) assert head["ETag"].strip('"').endswith(f"-{parts}"), head["ETag"] - _assert_bytes_equal(captures, envelopes) + _assert_read_back_exactly(captures, envelopes) From be05d376e560065e27877a4a503c9ba0b34848a3 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 09:25:55 -0400 Subject: [PATCH 06/13] Drop the unused E402 noqa from the capture chain test The fake-S3 import sits among the module's top-level imports, before any statement, so E402 never fires on it and the directive suppresses nothing. ruff 0.16.5 with --select E402,RUF100 reports it as an unused noqa (unused: E402); with it removed the same check passes. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- tests/test_native_capture_chain_live.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/test_native_capture_chain_live.py b/tests/test_native_capture_chain_live.py index 95c578346..725cb52da 100644 --- a/tests/test_native_capture_chain_live.py +++ b/tests/test_native_capture_chain_live.py @@ -42,7 +42,7 @@ import pytest # Module-level so the fake-S3 fixture registers in this module. -from tests.test_native_s3_client import ( # noqa: E402 +from tests.test_native_s3_client import ( ACCESS, BUCKET, REGION, SECRET, STATE, fake_s3, ) From 7d6eca3717118df194bb20502ca35a597a7d43ca Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 09:26:17 -0400 Subject: [PATCH 07/13] Warn when the C++20 sed step has nothing left to change The step exists only until #138 moves the Makefile to -std=c++20, and it had no way to say when that happened: after the merge the sed matches nothing, and the closing grep still counts 3, because the three lines already read c++20. The step would pass silently forever. An unchanged native/Makefile after the sed is exactly that state (before #138 the sed always rewrites all three lines), so the step now emits a ::warning:: asking for its own deletion. Rehearsed on a scratch copy of the Makefile: at this branch's Makefile the diff is 3 lines and no warning; with those edits committed, as #138 would leave them, the warning prints and the count is still 3. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index cba87dad0..6f312d0b6 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -576,12 +576,21 @@ jobs: # below is empty, and this step should be deleted. The Makefile's # variables are `:=` assignments, so a command-line override would # replace the whole flag set rather than the standard; hence sed. + # + # The count check alone cannot say that: with #138 merged the three + # lines already read c++20, so it still counts 3 and the step passes + # having done nothing. An unchanged Makefile is that state and only + # that state -- before #138 the sed always rewrites all three -- so + # it raises a warning on the run summary rather than failing the job. run: | sed -i \ -e '/^CXXFLAGS :=/s/-std=c++17/-std=c++20/' \ -e '/^HOST_CXXFLAGS :=/s/-std=c++17/-std=c++20/' \ -e '/^NVCCFLAGS :=/,/-std=/s/-std=c++17/-std=c++20/' \ native/Makefile + if git diff --quiet native/Makefile; then + echo "::warning title=Delete the C++20 sed step::native/Makefile already compiles as C++20, so #138 has landed and this step changed nothing. Delete it from python-checks.yml." + fi git diff native/Makefile test "$(grep -cE '^(HOST_)?CXXFLAGS := -std=c\+\+20|^ -std=c\+\+20 \\$' native/Makefile)" = 3 From 66cc82f0ee1b0c8a105d8f2618432bcd09f549e9 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 09:32:05 -0400 Subject: [PATCH 08/13] Check the capture sink binds the backend's RecordSink in CI _dmi_native_sink picks its RecordSink base once, when it loads: the _native_backend registration when that backend is loadable, module-local stand-ins otherwise. Only the first is attachable, since create_record_runtime checks isinstance against the backend's class. The cpu and live jobs build the sink with no backend, and native-backend-compile built the backend with no sink, so the backend branch of test_the_sink_derives_from_the_engines_record_sink ran in no job. native-backend-compile now builds build/_dmi_native_sink after the backend, with the same torch and PYTHON, and checks the binding in both import orders: through dmi's loaders (backend first, as the engine does) and as a bare import (the sink finds the backend itself). Each asserts RING_TYPES_ARE_STANDINS is false and NativePackSink subclasses the backend's RecordSink. Then it runs the adapter test and refuses a skip. The assertions are needed because the test alone is not enough. Rehearsed on a scratch copy of this branch with a copied torch 2.14+cu130 backend: with the backend present, both checks print an MRO through _native_backend.RecordSink and the test passes (1 passed). With the backend removed, both assertions fail on the stand-ins, and the test still passes, because its stand-in branch is a pass by design. The rehearsal also found that the bare import needs torch imported first, since the .so has no rpath to libtorch_cpu, so the check does that, as the suite does. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 52 ++++++++++++++++++++++++++++- 1 file changed, 51 insertions(+), 1 deletion(-) diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index 6f312d0b6..f905e045e 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -477,7 +477,8 @@ jobs: # bindings.cpp and the torch/CUDA link could break on any PR and stay # broken until someone built on a GPU host -- which is how the C++20 # floor torch 2.14 imposes went unnoticed (#138). This job compiles and - # links `make -C native all`, imports the result, and builds the ring + # links `make -C native all`, imports the result, checks that the capture + # sink built beside it derives from its RecordSink, and builds the ring # test binaries, running the two that need no device. # # No GPU is needed to COMPILE, only the toolkit, and this runner has @@ -614,6 +615,55 @@ jobs: assert hasattr(module, 'RecordSink'), 'no RecordSink in the backend' " + - name: Build the capture sink beside the backend + # _dmi_native_sink decides its RecordSink base once, when it loads: + # it derives NativePackSink from _native_backend's registration when + # that backend is loadable, and registers module-local stand-ins of + # its own when it is not. Only the first is attachable -- + # create_record_runtime checks isinstance against the BACKEND's + # RecordSink -- yet the cpu and live jobs build the sink with no + # backend, so the stand-in branch was the only one CI ever took. + # This is the one job that has a backend to bind to. Same torch and + # PYTHON as the backend build above. + run: make -C native build/_dmi_native_sink PYTHON=python + + - name: Check the sink binds the backend's RecordSink + # Two import orders, one process each, because the binding is fixed + # for the process by whichever order ran. Through dmi's loaders, the + # order the engine uses: backend first, then the sink. And the bare + # `import _dmi_native_sink` the adapter suite does, where the module + # has to find the backend itself through dmi.transport.native. Both + # assert the stand-ins were NOT taken, which is what the pytest run + # after them cannot: its test passes on either branch by design, so + # on its own it would stay green with the binding broken. It still + # runs, because here it takes the test's backend branch, which no + # other job reaches; the grep refuses a skip as a pass. + run: | + PYTHONPATH=src python -c " + from dmi.storage.capture.native_sink import _load_native_sink_extension + from dmi.transport import native + backend = native._load_extension() + sink = _load_native_sink_extension() + assert not sink.RING_TYPES_ARE_STANDINS, 'the sink bound stand-ins' + assert issubclass(sink.NativePackSink, backend.RecordSink), sink.NativePackSink.__mro__ + assert not hasattr(sink, 'RecordSink'), 'the sink registered a second RecordSink' + print('loaders:', sink.NativePackSink.__mro__) + " + PYTHONPATH=src:native/build python -c " + import torch # first, as the suite does: the .so has no rpath to libtorch_cpu + import _dmi_native_sink as sink + from dmi.transport import native + assert not sink.RING_TYPES_ARE_STANDINS, 'the sink bound stand-ins' + assert issubclass(sink.NativePackSink, native._load_extension().RecordSink), sink.NativePackSink.__mro__ + print('bare import:', sink.NativePackSink.__mro__) + " + python -m pip install pytest + python -m pytest -q -p no:cacheprovider \ + "tests/test_native_adapter_torch.py::test_the_sink_derives_from_the_engines_record_sink" \ + > sink-test.out || { cat sink-test.out; exit 1; } + cat sink-test.out + grep -Eq '^1 passed in ' sink-test.out + - name: Build the ring tests and run the host-only ones # All six binaries are compiled, so a .cu test that no longer builds # fails here. Only test_record_consumer and From 69f107c79d51b5d91817e481230de49e8ca755a0 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 09:33:16 -0400 Subject: [PATCH 09/13] Bound the MinIO health probe and drop the last-release claim The health loop ran curl -sf with no timeout, so a socket that accepts and never answers would have held the first probe until the 30-minute job timeout, and the loop's 30 attempts would never have run. Each probe is now capped with --max-time 2. Checked locally against a listener that accepts and sleeps: curl gives up after 2.01 s with exit 28, which the loop treats as not up yet and retries. The comment also called the pinned release "the last release published as an image". Review found later tags on quay.io, for example RELEASE.2025-09-07T16-13-09Z.hotfix.7aa24e772. (I could not list the tags from here: quay's API wants a token.) The pin is still right, so the comment now says only what it does: tag and digest together keep an upstream retag from changing what the job tests against. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index f905e045e..72db1fc1f 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -375,9 +375,13 @@ jobs: # A `docker run` step, not a `services:` entry: a service container # takes no command, and this image does nothing without `server # `. quay.io because the Docker Hub repository is gone; pinned - # by digest, the last release published as an image. Plain HTTP, on - # 9100 because ClickHouse holds 9000. B7a moves it to HTTPS with a CA - # generated at job start. + # by tag and digest, so a retag upstream cannot change what this + # job tests against. Plain HTTP, on 9100 because ClickHouse holds + # 9000. B7a moves it to HTTPS with a CA generated at job start. + # + # Each health probe is capped at 2 s: a socket that accepts and never + # answers would otherwise hold one curl until the job timeout, and + # the loop's 30 attempts would never get past the first. run: | docker run -d --name minio -p 9100:9000 \ -e MINIO_ROOT_USER=dmi-ci-access \ @@ -385,7 +389,7 @@ jobs: quay.io/minio/minio:RELEASE.2025-09-07T16-13-09Z@sha256:14cea493d9a34af32f524e538b8346cf79f3321eff8e708c1e2960462bd8936e \ server /data for attempt in $(seq 1 30); do - if curl -sf http://127.0.0.1:9100/minio/health/live; then + if curl -sf --max-time 2 http://127.0.0.1:9100/minio/health/live; then echo "MinIO is up" exit 0 fi From 988054405edc1846b342a0d563c9105ed84b6c61 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 09:33:41 -0400 Subject: [PATCH 10/13] Describe the runner disk step as a precaution, with measured figures The job comment justified freeing disk with "the ~14 GB a hosted runner guarantees", but the runner this job actually got had far more. The df output from run 35933626430 (job 107425626899) shows 87 GB free before the step, 103 GB after it, 99 GB after the CUDA apt install and 94 GB after torch. Freeing the Android SDK and .NET took about two minutes of that run. The step stays, since the free space belongs to the runner image and can change, but both comments now call it a precaution and quote only the figures that log shows. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 17 ++++++++++++----- 1 file changed, 12 insertions(+), 5 deletions(-) diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index 72db1fc1f..8aa4115a2 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -494,8 +494,10 @@ jobs: # (nvidia/cuda:13.0.1-devel-ubuntu24.04, amd64, per its registry # manifest) -- several times that unpacked -- before the ~4.4 GB that # torch 2.14+cu130 and its nvidia-* wheels take installed (measured in a - # cu130 venv). That does not reliably fit the ~14 GB a hosted runner - # guarantees. Installing only the packages the build reads is ~2.7 GB: + # cu130 venv). The free disk on a hosted runner is the image's to + # decide, not this repository's, so the job keeps its footprint small + # rather than depend on it. Installing only the packages the build + # reads is ~2.7 GB: # the header set comes from the -MMD dependency files of a real build, # mapped to packages with dpkg -S on a host whose toolkit came from the # same apt repository, and the sizes are those packages' Installed-Size. @@ -511,9 +513,14 @@ jobs: steps: - name: Free runner disk - # The two largest preinstalled trees this job never uses (the - # Android SDK and .NET); `df` either side so a disk failure is - # readable from the log. + # A precaution, not a measured need. On run 35933626430 of this job + # the runner had 87 GB free before this step, 103 GB after it, and + # still 94 GB once the CUDA toolkit and torch were installed, so the + # job fits without it today. It stays because that figure belongs to + # the runner image and can shrink between runs, and it removes only + # the two largest preinstalled trees this job never uses (the Android + # SDK and .NET). It cost about two minutes on that run. `df` either + # side, so a disk failure can be read from the log. run: | df -h / sudo rm -rf /usr/local/lib/android /usr/share/dotnet From ede71cb76d1941b0739ca96780fac63487ea407d Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 10:17:29 -0400 Subject: [PATCH 11/13] Pull MinIO from Bitnami's archive, and let a warning pass the sink check Two CI failures on the previous push, neither in the code under test: - The MinIO step's image could no longer be pulled. The MinIO project has withdrawn its public images: quay.io/minio/minio and docker.io/minio/minio both now answer anonymous pulls with 401, including `latest`, although this job pulled the pinned quay image fine a day earlier. The step now runs bitnamilegacy/minio:2025.7.23-debian-12-r5, pinned by digest -- a frozen archive, which suits a fixture whose S3 multipart rules do not change. The Bitnami image starts its own server from MINIO_ROOT_USER/_PASSWORD, so the `server /data` argument is gone, and the health loop allows about two minutes for its entrypoint's setup. The live skip gate then failed only because the suite never ran. - The sink-binding check passed ("1 passed, 1 warning in 1.50s") but the step greps for '^1 passed in ', and torch's NumPy warning changed the summary line. The step now accepts a summary with warnings and still fails on any skip. Validated locally: the workflow parses, and the new pass condition accepts "1 passed" and "1 passed, 1 warning" and refuses "1 skipped". The image's digest, amd64 manifest and entrypoint were read from Docker Hub's registry; running it needs docker, which only CI has. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 33 ++++++++++++++++++----------- 1 file changed, 21 insertions(+), 12 deletions(-) diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index 8aa4115a2..ab0a3e421 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -372,23 +372,29 @@ jobs: # client's part shrunk to 4 MiB, the fake-S3 case fails only on its # own part-size assertion, and the MinIO case on exactly that error. # - # A `docker run` step, not a `services:` entry: a service container - # takes no command, and this image does nothing without `server - # `. quay.io because the Docker Hub repository is gone; pinned - # by tag and digest, so a retag upstream cannot change what this - # job tests against. Plain HTTP, on 9100 because ClickHouse holds - # 9000. B7a moves it to HTTPS with a CA generated at job start. + # Bitnami's archived MinIO build. The MinIO project no longer + # publishes public images: quay.io/minio/minio and docker.io/ + # minio/minio both answer anonymous pulls with 401 (seen 2026-09-24, + # after this job had pulled the quay image fine). bitnamilegacy is a + # frozen archive, which suits a test fixture -- the S3 multipart + # rules it checks do not change -- and it is pinned by tag and + # digest, so nothing upstream can change what this job tests + # against. If it too disappears, Garage is the fallback. The image + # starts its own server from MINIO_ROOT_USER/_PASSWORD, so this is a + # `docker run` step only for the port map and the logs on failure. + # Plain HTTP, on 9100 because ClickHouse holds 9000. B7a moves it to + # HTTPS with a CA generated at job start. # # Each health probe is capped at 2 s: a socket that accepts and never - # answers would otherwise hold one curl until the job timeout, and - # the loop's 30 attempts would never get past the first. + # answers would otherwise hold one curl until the job timeout. The + # Bitnami entrypoint runs its setup before the server listens, so the + # loop allows about two minutes. run: | docker run -d --name minio -p 9100:9000 \ -e MINIO_ROOT_USER=dmi-ci-access \ -e MINIO_ROOT_PASSWORD=dmi-ci-secret-key \ - quay.io/minio/minio:RELEASE.2025-09-07T16-13-09Z@sha256:14cea493d9a34af32f524e538b8346cf79f3321eff8e708c1e2960462bd8936e \ - server /data - for attempt in $(seq 1 30); do + docker.io/bitnamilegacy/minio:2025.7.23-debian-12-r5@sha256:6dabb4a2088c9a79908de3bc05f4586c23ad2182c8908e7e3acbf61c1467fb20 + for attempt in $(seq 1 60); do if curl -sf --max-time 2 http://127.0.0.1:9100/minio/health/live; then echo "MinIO is up" exit 0 @@ -673,7 +679,10 @@ jobs: "tests/test_native_adapter_torch.py::test_the_sink_derives_from_the_engines_record_sink" \ > sink-test.out || { cat sink-test.out; exit 1; } cat sink-test.out - grep -Eq '^1 passed in ' sink-test.out + # "1 passed in ..." or "1 passed, N warnings in ..." -- a warning + # (torch's NumPy one appears here) is not a failure; a skip is. + grep -Eq '^1 passed(,| in )' sink-test.out + ! grep -Eq '[0-9]+ skipped' sink-test.out - name: Build the ring tests and run the host-only ones # All six binaries are compiled, so a .cu test that no longer builds From 8d7c3d1c36e0b54ecf77d2e3226538bd1490ac6f Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 10:33:20 -0400 Subject: [PATCH 12/13] Drop MinIO from the live job MinIO was there for one check: that a multipart pack follows real S3's rules (every part but the last at least 5 MiB, and a "-N" ETag), which the in-repo fake S3 does not enforce. Nothing in DMI runs against MinIO -- the object store it uses is Garage -- and the MinIO project has withdrawn its public images, so the job had come to depend on a frozen third-party archive for a store DMI does not use. The CI step, its DMI_MINIO_* settings and the MinIO test case are removed. The fake-S3 multipart case stays: it asserts the parts the client sends (every part but the last is exactly the client's chunk and at least 5 MiB) and reads the pack back byte-exact. Checks against a real store belong to the Garage suites. Chain test against local ClickHouse: 2 passed, 0 skipped; no tables left. The workflow parses, and the only step removed from any job is the MinIO one. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 47 ---------------- tests/test_native_capture_chain_live.py | 72 +++---------------------- 2 files changed, 8 insertions(+), 111 deletions(-) diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index ab0a3e421..3f94aa539 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -363,57 +363,10 @@ jobs: CURL_INCDIR=/usr/include/x86_64-linux-gnu \ CURL_LIBDIR=/usr/lib/x86_64-linux-gnu - - name: Start MinIO for the multipart capture case - # A real S3 implementation for the one case the in-repo fake cannot - # judge: a pack of 64 MiB or more goes up as a multipart upload, and - # the fake assembles parts of any size, where S3 refuses every part - # but the last under 5 MiB (EntityTooSmall at CompleteMultipartUpload). - # Verified locally with the minio binary from this image: with the - # client's part shrunk to 4 MiB, the fake-S3 case fails only on its - # own part-size assertion, and the MinIO case on exactly that error. - # - # Bitnami's archived MinIO build. The MinIO project no longer - # publishes public images: quay.io/minio/minio and docker.io/ - # minio/minio both answer anonymous pulls with 401 (seen 2026-09-24, - # after this job had pulled the quay image fine). bitnamilegacy is a - # frozen archive, which suits a test fixture -- the S3 multipart - # rules it checks do not change -- and it is pinned by tag and - # digest, so nothing upstream can change what this job tests - # against. If it too disappears, Garage is the fallback. The image - # starts its own server from MINIO_ROOT_USER/_PASSWORD, so this is a - # `docker run` step only for the port map and the logs on failure. - # Plain HTTP, on 9100 because ClickHouse holds 9000. B7a moves it to - # HTTPS with a CA generated at job start. - # - # Each health probe is capped at 2 s: a socket that accepts and never - # answers would otherwise hold one curl until the job timeout. The - # Bitnami entrypoint runs its setup before the server listens, so the - # loop allows about two minutes. - run: | - docker run -d --name minio -p 9100:9000 \ - -e MINIO_ROOT_USER=dmi-ci-access \ - -e MINIO_ROOT_PASSWORD=dmi-ci-secret-key \ - docker.io/bitnamilegacy/minio:2025.7.23-debian-12-r5@sha256:6dabb4a2088c9a79908de3bc05f4586c23ad2182c8908e7e3acbf61c1467fb20 - for attempt in $(seq 1 60); do - if curl -sf --max-time 2 http://127.0.0.1:9100/minio/health/live; then - echo "MinIO is up" - exit 0 - fi - sleep 2 - done - docker logs minio - echo "MinIO did not become healthy" >&2 - exit 1 - - name: Run the ClickHouse live suites env: DMI_CLICKHOUSE_HOST: 127.0.0.1 DMI_CLICKHOUSE_PORT: "9000" - # Unset, the MinIO multipart case skips, and the gate below turns - # that skip into a failure: the variable is load-bearing. - DMI_MINIO_ENDPOINT: http://127.0.0.1:9100 - DMI_MINIO_ACCESS_KEY: dmi-ci-access - DMI_MINIO_SECRET_KEY: dmi-ci-secret-key # `python -m pytest`, matching the Makefile. `not garage` drops the one # suite that also needs a Garage S3 endpoint; everything else under # these markers needs nothing but the server above. diff --git a/tests/test_native_capture_chain_live.py b/tests/test_native_capture_chain_live.py index 725cb52da..134746806 100644 --- a/tests/test_native_capture_chain_live.py +++ b/tests/test_native_capture_chain_live.py @@ -15,14 +15,14 @@ that went in. No CUDA, no ring engine, no conformance driver, no Python on the data path. -The object store is the in-repo signature-verifying fake S3, and, when -DMI_MINIO_ENDPOINT names one, a real MinIO as well. The multipart case needs -both: a pack of 64 MiB or more crosses the client's multipart threshold, and -the fake S3 enforces NOTHING about part sizes -- it accepts a part of any -length, takes the part list on trust, and never checks the ETags in the -completion body -- so against it the test checks the parts the client sent -itself. MinIO holds the client to S3's rules: every part but the last is at -least 5 MiB, or CompleteMultipartUpload fails with EntityTooSmall. +The object store is the in-repo signature-verifying fake S3. For the +multipart case -- a pack of 64 MiB or more crosses the client's multipart +threshold -- note that the fake enforces NOTHING about part sizes: it accepts +a part of any length, takes the part list on trust, and never checks the +ETags in the completion body. So the test checks the parts the client sent +itself against S3's rule: every part but the last is at least 5 MiB, or a +real store refuses CompleteMultipartUpload with EntityTooSmall. The real +object store DMI runs against, Garage, is exercised by the Garage suites. Build: make -C native build/_dmi_native_sink build/_dmi_native_store \\ PYTHON=/bin/python @@ -71,9 +71,6 @@ CLICKHOUSE_HTTP_PORT = int(environ.get("DMI_CLICKHOUSE_HTTP_PORT", "8123")) DATABASE = environ.get("DMI_CLICKHOUSE_DATABASE", "default") -MINIO_ENDPOINT = environ.get("DMI_MINIO_ENDPOINT", "") -MINIO_ACCESS = environ.get("DMI_MINIO_ACCESS_KEY", "minioadmin") -MINIO_SECRET = environ.get("DMI_MINIO_SECRET_KEY", "minioadmin") LAYOUT = "capture_pack_reference_v1" # at::ScalarType numeric values (c10/core/ScalarType.h -- stable ABI). @@ -376,56 +373,3 @@ def test_a_multipart_pack_through_the_fake_s3(fake_s3, tmp_path): assert all(size >= S3_MIN_PART for size in parts[:-1]), parts assert STATE.objects[key]["etag"].endswith('-multipart"') _assert_read_back_exactly(captures, envelopes) - - -@contextmanager -def _minio_bucket(): - """A fresh bucket on the MinIO endpoint, emptied and removed afterwards.""" - import botocore.session - - client = botocore.session.get_session().create_client( - "s3", endpoint_url=MINIO_ENDPOINT, region_name=REGION, - aws_access_key_id=MINIO_ACCESS, aws_secret_access_key=MINIO_SECRET) - bucket = f"dmi-chain-{uuid.uuid4().hex[:16]}" - client.create_bucket(Bucket=bucket) - try: - yield client, bucket - finally: - for page in client.get_paginator("list_objects_v2").paginate( - Bucket=bucket): - for item in page.get("Contents", []): - client.delete_object(Bucket=bucket, Key=item["Key"]) - for upload in client.list_multipart_uploads(Bucket=bucket).get( - "Uploads", []): - client.abort_multipart_upload( - Bucket=bucket, Key=upload["Key"], UploadId=upload["UploadId"]) - client.delete_bucket(Bucket=bucket) - - -# A skip allowed ONLY for the endpoint being unset: CI's live job starts a -# MinIO and sets it, and that job's gate fails any skip, so a lost variable -# there is a red job rather than a silent gap. -@pytest.mark.skipif( - not MINIO_ENDPOINT, - reason="DMI_MINIO_ENDPOINT is unset; the live CI job starts MinIO and " - "sets it") -def test_a_multipart_pack_through_minio(tmp_path): - envelopes = _large_envelopes(LARGE_COUNT, LARGE_ELEMENTS) - with _minio_bucket() as (client, bucket), _catalog() as prefix: - config = _storage_config(MINIO_ENDPOINT, bucket, MINIO_ACCESS, - MINIO_SECRET, prefix) - sink_snapshot, snapshot, captures = _run_chain( - config, tmp_path / "spool", envelopes, sink_overrides=LARGE_SINK) - - assert sink_snapshot["persisted_records"] == LARGE_COUNT, sink_snapshot - assert snapshot["uploaded_packs"] == 1, snapshot - assert snapshot["indexed_rows"] == LARGE_COUNT, snapshot - key, object_bytes = _one_large_pack(captures, LARGE_COUNT) - - # A multipart object's ETag is "-": the - # pack really went up in parts, and MinIO accepted every one of them. - head = client.head_object(Bucket=bucket, Key=key) - assert head["ContentLength"] == object_bytes - parts = math.ceil(object_bytes / MULTIPART_CHUNK) - assert head["ETag"].strip('"').endswith(f"-{parts}"), head["ETag"] - _assert_read_back_exactly(captures, envelopes) From 6b08e743155415b274397429f6e6d55552d56737 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 11:10:25 -0400 Subject: [PATCH 13/13] Build the documented native build in the compile job, now that #138 and #144 landed - #138 is on main, so the temporary step that applied its C++20 flags with sed changed nothing and would only raise its "delete me" warning. Deleted, as its own comment asked. - #144 puts the capture extensions in `make -C native all`, and the store needs libcurl's headers; without them the build now stops before anything else, by design. The job installs libcurl4-openssl-dev and pkg-config and lets the Makefile find libcurl through pkg-config -- its default lookup, with no CURL_* override -- so it builds exactly what docs/install.md tells a user to run. It lists the three extensions it expects in src/dmi afterwards. Validated locally: the workflow parses and the job's steps are as listed; ruff 0.16.5 F821 passes on the merged src/dmi; with #144's Makefile the CPU targets build and install as links, and the chain test passes (2 passed). The CUDA build itself runs only in CI. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- .github/workflows/python-checks.yml | 43 ++++++++--------------------- 1 file changed, 12 insertions(+), 31 deletions(-) diff --git a/.github/workflows/python-checks.yml b/.github/workflows/python-checks.yml index 3f94aa539..d17f12600 100644 --- a/.github/workflows/python-checks.yml +++ b/.github/workflows/python-checks.yml @@ -537,42 +537,23 @@ jobs: -DCMAKE_POSITION_INDEPENDENT_CODE=ON cmake --build third_party/clickhouse-cpp/build -j"$(nproc)" - - name: "Compile the torch-including targets as C++20 until #138 lands" - # torch 2.14's headers open with `#error C++20 or later compatible - # compiler is required` (reproduced: the stock flags stop at - # bindings.cpp with exactly that). #138 moves CXXFLAGS, HOST_CXXFLAGS - # and NVCCFLAGS to -std=c++20; until it merges, this makes the same - # three edits to the checkout -- they are the lines #138 changes and - # nothing else. Once it has merged, the sed matches nothing, the diff - # below is empty, and this step should be deleted. The Makefile's - # variables are `:=` assignments, so a command-line override would - # replace the whole flag set rather than the standard; hence sed. + - name: Compile the full native build + # `all` is the documented build: the backend plus the capture + # extensions (#144), and the store needs libcurl's headers. With the + # dev package installed, the Makefile finds it through pkg-config -- + # its default lookup, so no CURL_* override here: this job builds + # exactly what docs/install.md tells a user to run. # - # The count check alone cannot say that: with #138 merged the three - # lines already read c++20, so it still counts 3 and the step passes - # having done nothing. An unchanged Makefile is that state and only - # that state -- before #138 the sed always rewrites all three -- so - # it raises a warning on the run summary rather than failing the job. - run: | - sed -i \ - -e '/^CXXFLAGS :=/s/-std=c++17/-std=c++20/' \ - -e '/^HOST_CXXFLAGS :=/s/-std=c++17/-std=c++20/' \ - -e '/^NVCCFLAGS :=/,/-std=/s/-std=c++17/-std=c++20/' \ - native/Makefile - if git diff --quiet native/Makefile; then - echo "::warning title=Delete the C++20 sed step::native/Makefile already compiles as C++20, so #138 has landed and this step changed nothing. Delete it from python-checks.yml." - fi - git diff native/Makefile - test "$(grep -cE '^(HOST_)?CXXFLAGS := -std=c\+\+20|^ -std=c\+\+20 \\$' native/Makefile)" = 3 - - - name: Compile the full native backend # LIBRARY_PATH puts the toolkit's libcuda stub where the linker finds # -lcuda; the result does not end up NEEDing libcuda (checked with # readelf on a local build: the link is --as-needed and nothing calls # the driver API), so check-link's ldd has no driver to find either. env: LIBRARY_PATH: /usr/local/cuda-13.0/targets/x86_64-linux/lib/stubs - run: make -C native -j"$(nproc)" all PYTHON=python + run: | + sudo apt-get install -y -q libcurl4-openssl-dev pkg-config + make -C native -j"$(nproc)" all PYTHON=python + ls -l src/dmi/_native_backend*.so src/dmi/_dmi_native_sink*.so src/dmi/_dmi_native_store*.so - name: Import the backend # A shared library links with undefined symbols; only loading it @@ -593,8 +574,8 @@ jobs: # create_record_runtime checks isinstance against the BACKEND's # RecordSink -- yet the cpu and live jobs build the sink with no # backend, so the stand-in branch was the only one CI ever took. - # This is the one job that has a backend to bind to. Same torch and - # PYTHON as the backend build above. + # This is the one job that has a backend to bind to. `all` above has + # already built it; the step stays so the dependency is explicit. run: make -C native build/_dmi_native_sink PYTHON=python - name: Check the sink binds the backend's RecordSink