From d88ba6a16736edd285f6a93a4fbde7778d3ce277 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Wed, 23 Sep 2026 16:53:11 -0400 Subject: [PATCH 1/5] Name the storage choices in-memory, persistent and none The user's storage choice now has exactly three values: - "in-memory": records delivered in memory. Until its consumer interface exists this is the ClickHouseRecordSink path, which needs a host engine. - "persistent": the native capture storage path, object store + catalog + ClickHouse. - "none": capture off entirely. The engine allocates no ring, so none of the ring's pinned staging or payload memory. A record runtime, a host engine, an explicit ring_config and adapter attachment are refused, each with a message that names the two choices that capture. dmi.config.USER_STORAGE_CHOICES lists them. "auto" stays the unset default, the inference every caller made before the field existed; the configurator will always emit one of the three. The earlier names still work, with a DeprecationWarning that names the replacement: "native" means in-memory and "capture" means persistent. MonitoringConfig keeps whatever the caller wrote, so an integration that reads the field back still sees its own value. The engine acts on canonical_storage_backend. "none" used to mean capture and transport with no persistence: the ring still ran, and every record then failed at flush_and_wait with "record sink is not configured". Nothing depended on that, and it is now a clear refusal up front. The one test that expected an adapter to attach under "none" now expects the refusal. Docs: the v1 contract lists the three choices, the aliases and what "none" refuses; the capture-storage design doc uses the new names. Tests: a new tests/test_storage_backend_choices.py (13 cpu tests) covers the three choices, auto as the unset default, the aliases and their warning, the unknown-value message, and none allocating no ring and refusing a ring_config, a record runtime, a host engine and attachment. It failed at import before the change. The existing suites moved to the new names. CPU tier: 2238 passed, 1 skipped (no CUDA device), 0 deprecation warnings. GPU (RTX 4090) capture storage, sink ring and record-ring refusal suites: 29 passed with DeprecationWarning as an error, so no test still uses an old name. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- docs/capture-storage-design.md | 15 ++- docs/integration-api-v1.md | 38 ++++-- src/dmi/adapters/base.py | 51 +++++--- src/dmi/adapters/huggingface/adapter.py | 2 +- src/dmi/adapters/huggingface/generation.py | 10 +- src/dmi/config.py | 113 ++++++++++------- src/dmi/engine.py | 46 +++++-- src/dmi/storage/capture/native_sink.py | 2 +- src/dmi/storage/native_capture.py | 2 +- tests/test_engine_runtime_api.py | 18 +-- tests/test_hf_capture_refusal.py | 18 +-- tests/test_integration_api_v1.py | 2 +- tests/test_native_capture_storage_gpu_e2e.py | 6 +- tests/test_native_capture_storage_wiring.py | 10 +- tests/test_native_sink_ring_e2e.py | 4 +- tests/test_storage_backend_choices.py | 127 +++++++++++++++++++ 16 files changed, 337 insertions(+), 127 deletions(-) create mode 100644 tests/test_storage_backend_choices.py diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index 57fd96652..7e6bf0586 100644 --- a/docs/capture-storage-design.md +++ b/docs/capture-storage-design.md @@ -106,7 +106,7 @@ exclusive per record runtime: | ClickHouse role | the store itself | a rebuildable index over the packs | | ClickHouse footprint | one configured table (`offload` by default) | `{prefix}_*` (`dmi_*` by default) | | Selected by | `create_record_runtime(fmt)` with no `record_sink`, plus a `host_engine`/`db_config` on the engine | `create_record_runtime(fmt, record_sink=reference.native_sink)` | -| Declared by | `MonitoringConfig(storage_backend="native")` | `MonitoringConfig(storage_backend="capture")` | +| Declared by | `MonitoringConfig(storage_backend="in-memory")` (formerly `"native"`) | `MonitoringConfig(storage_backend="persistent")` (formerly `"capture"`) | | Status | production | explicitly reference-only; production sinks remain native-only | Both are `ring::RecordSink` implementations and the record engine takes exactly @@ -128,7 +128,7 @@ ClickHouse insert pipeline started, connected and never fed. ```python engine = MonitoringEngine( - config=MonitoringConfig(storage_backend="capture"), + config=MonitoringConfig(storage_backend="persistent"), model_id="...", ring_config=ring_config, ) # a host_engine here is now refused @@ -138,11 +138,12 @@ runtime = engine.create_record_runtime( ) ``` -`"native"` is the mirror image: it requires a host engine and refuses an -explicit sink. `"none"` is capture and transport with no persistence at all. -The default is `"auto"`, which infers the backend from what was passed -- what -every caller did before the field existed, so nothing that predates it -changes. +`"in-memory"` is the mirror image: it requires a host engine and refuses an +explicit sink. `"none"` turns capture off entirely: no ring is allocated, and +nothing that would capture can be attached. The default is `"auto"`, which +infers the backend from what was passed -- what every caller did before the +field existed, so nothing that predates it changes. The earlier names +`"native"` and `"capture"` still work, with a `DeprecationWarning`. Their ClickHouse footprints are disjoint, so the two can share one server: the catalog's schema guard only ever names `{prefix}_*` objects, and `drop_schema` diff --git a/docs/integration-api-v1.md b/docs/integration-api-v1.md index db1ebb673..c6a9845f9 100644 --- a/docs/integration-api-v1.md +++ b/docs/integration-api-v1.md @@ -135,15 +135,32 @@ schedule.should_capture_step( The predicates apply warmup, then offset, then stride. Step selection also honors `capture_prefill`/`capture_decode`; an unknown phase raises `ValueError`. `MonitoringConfig` carries this schedule plus three storage fields: -`storage_backend` (`"auto" | "native" | "capture" | "none"`), -`capture_sink_config` (a `NativeSinkConfig`, or `None`) and +`storage_backend`, `capture_sink_config` (a `NativeSinkConfig`, or `None`) and `capture_storage_config` (a `NativeCaptureStorageConfig` from -`dmi.storage.native_capture`, or `None`). All are acted on by -`MonitoringEngine`, and some combinations are refused at construction -- -`storage_backend="native"` without a host engine, or `"capture"`/`"none"` with -one, or `capture_storage_config` without `capture_sink_config` -- so a caller -setting them should expect `ValueError` rather than a silent choice. -`capture_sink_config` is read only when `storage_backend` is `"capture"`; see +`dmi.storage.native_capture`, or `None`). + +`storage_backend` is the user's storage choice, one of +`dmi.config.USER_STORAGE_CHOICES`: + +- `"in-memory"`: records are delivered in memory. Until its consumer interface + exists, this runs the C++ `ClickHouseRecordSink` and needs a host engine. +- `"persistent"`: the native capture storage path (object store + catalog + + ClickHouse). +- `"none"`: capture off entirely. The engine allocates no ring, and a record + runtime, a host engine, a `ring_config` and adapter attachment are refused. + +Left unset, it is `"auto"`, which infers the path from what was passed, as every +caller did before the field existed. The earlier names `"native"` and +`"capture"` still work, with a `DeprecationWarning`, and mean `"in-memory"` and +`"persistent"`. `MonitoringConfig.canonical_storage_backend` gives the current +name. + +All three fields are acted on by `MonitoringEngine`, and some combinations are +refused at construction -- `"in-memory"` without a host engine, +`"persistent"`/`"none"` with one, or `capture_storage_config` without +`capture_sink_config` -- so a caller setting them should expect `ValueError` +rather than a silent choice. +`capture_sink_config` is read only when `storage_backend` is `"persistent"`; see `docs/capture-storage-design.md` for the writer it selects. With `capture_storage_config` as well, the engine runs the native storage service in-process: from `create_record_runtime` until `close`, or until @@ -156,7 +173,7 @@ drives this path yet: the HF, vLLM and Megatron integrations never create a capture record runtime, so it is reached only by a caller that builds the record runtime and its hook points itself, as `tests/test_native_capture_storage_gpu_e2e.py` does, and the HF adaptor -refuses the capture backend with `ConfigurationError` rather than generating +refuses the persistent backend with `ConfigurationError` rather than generating with nothing stored. The catalog takes one publisher per `(database, table_prefix)`, so a second engine on the same catalog is refused at `create_record_runtime`. @@ -484,7 +501,8 @@ or non-owning PP/TP hooks, installs ring fields on remaining HookPoints, and publishes the selected inventory. It mutates HookPoints and is not transactional; call it before graph capture. An unknown selection raises `ValueError`, and an executable inventory containing -`module=None` raises `RuntimeError`. Under `storage_backend="capture"` it raises +`module=None` raises `RuntimeError`. Under `storage_backend="persistent"`, and +under `"none"`, which turns capture off, it raises `dmi.configuration.ConfigurationError`, as do the HF entry points built on it, `generate_with_monitoring()` and `generate_greedy_with_monitoring()`: no adaptor drives the capture storage path yet, and attaching would install hooks whose diff --git a/src/dmi/adapters/base.py b/src/dmi/adapters/base.py index 75a02e25d..0b1a79011 100644 --- a/src/dmi/adapters/base.py +++ b/src/dmi/adapters/base.py @@ -54,28 +54,40 @@ def _refuse_unwired_capture_storage( engine: object, entry_point: str, adapter_name: str ) -> None: - """Refuse ``storage_backend="capture"`` at an adapter entry point. - - Adapters drive legacy HookPoints, and nothing connects those to the - capture path yet. Under the capture config they run on the engine's - legacy ring, which has no host (the backend refuses one), so its P2P - thread drops every capture: generation "succeeds" and the catalog stays - empty. After ``create_record_runtime`` the same hooks meet a record ring, - which refuses their metadata. Neither stores what the config asked for, - so the refusal comes here, before anything is armed. + """Refuse a storage choice an adapter cannot serve, before anything arms. + + ``storage_backend="none"`` turns capture off, so the engine has no ring + and there is nothing to attach hooks to. + + ``storage_backend="persistent"``: adapters drive legacy HookPoints, and + nothing connects those to the native capture storage path yet. Under + that config they run on the engine's legacy ring, which has no host (the + backend refuses one), so its P2P thread drops every capture: generation + "succeeds" and the catalog stays empty. After ``create_record_runtime`` + the same hooks meet a record ring, which refuses their metadata. Neither + stores what the config asked for. """ - if getattr(engine, "_storage_backend", None) != "capture": + backend = getattr(engine, "_storage_backend", None) + if backend not in ("persistent", "none"): return from ..configuration.errors import ConfigurationError + if backend == "none": + raise ConfigurationError( + f"{entry_point}: config.storage_backend='none' means capture is " + f"off, so the engine has no ring and {adapter_name} has nothing " + "to attach its hooks to. Pick 'in-memory' or 'persistent' to " + "capture, or run the model without a monitoring adapter." + ) raise ConfigurationError( - f"{entry_point}: config.storage_backend='capture' selects the native " - f"capture storage path, but capture storage is not wired to " - f"{adapter_name} yet. Its hooks write to the legacy ring, which " - "under this config has no host and drops every capture, so " - "generation would succeed and store nothing. Capture records through " - "engine.create_record_runtime(...), or use storage_backend='native' " - "with a host engine for monitored generation." + f"{entry_point}: config.storage_backend='persistent' selects the " + f"native capture storage path, but capture storage is not wired to " + f"{adapter_name} yet. Its hooks write to the legacy ring, which under this config has " + "no host and drops every capture, so generation would succeed and " + "store nothing. Capture records through " + "engine.create_record_runtime(...), or use " + "storage_backend='in-memory' with a host engine for monitored " + "generation." ) @@ -207,8 +219,9 @@ def attach_model( establishes the enabled/disabled state across every spec, and each later filter only ever disables further. - Raises ``ConfigurationError`` under ``storage_backend="capture"``, - which no adapter's hooks are wired to yet. + Raises ``ConfigurationError`` under ``storage_backend="persistent"``, + which no adapter's hooks are wired to yet, and under ``"none"``, which + turns capture off. """ _refuse_unwired_capture_storage( self.engine, f"{type(self).__name__}.attach_model()", diff --git a/src/dmi/adapters/huggingface/adapter.py b/src/dmi/adapters/huggingface/adapter.py index 90975e05b..06a7aa1e9 100644 --- a/src/dmi/adapters/huggingface/adapter.py +++ b/src/dmi/adapters/huggingface/adapter.py @@ -372,7 +372,7 @@ def _handle_driver_failure(self, exc: Exception) -> None: """ engine = self.engine if (getattr(engine, "_record_mode", False) - or getattr(engine, "_storage_backend", None) == "capture"): + or getattr(engine, "_storage_backend", None) == "persistent"): raise exc if self._warned_driver_failure: return diff --git a/src/dmi/adapters/huggingface/generation.py b/src/dmi/adapters/huggingface/generation.py index 1aed21d9a..d5633d166 100644 --- a/src/dmi/adapters/huggingface/generation.py +++ b/src/dmi/adapters/huggingface/generation.py @@ -110,8 +110,9 @@ def _generate_with_monitoring_impl( external compilation and injects an equivalent ``CompileConfig`` so HF compiles only the decode path (prefill stays uncompiled). - Raises ``ConfigurationError`` under ``storage_backend="capture"``, which - this adapter is not wired to yet. + Raises ``ConfigurationError`` under ``storage_backend="persistent"``, + which this adapter is not wired to yet, and under ``"none"``, which turns + capture off. """ import types @@ -536,8 +537,9 @@ def generate_greedy_with_monitoring( (``cuda_graphs=False``). Unpadded batches are unaffected. monitoring: if True, install ring transport hooks via HuggingFaceAdapter and call before_forward_manual before each forward pass. Refused - with ``ConfigurationError`` under ``storage_backend="capture"``, - which this adapter is not wired to yet. + with ``ConfigurationError`` under ``storage_backend="persistent"``, + which this adapter is not wired to yet, and under ``"none"``, + which turns capture off. hook_selection: hook selection preset (e.g. "hidden-states", "full"). Only used when monitoring=True. no_strip_left_pad: forwarded to ``HuggingFaceAdapter`` when monitoring=True. If diff --git a/src/dmi/config.py b/src/dmi/config.py index c7b13f1c3..d37dd5beb 100644 --- a/src/dmi/config.py +++ b/src/dmi/config.py @@ -3,10 +3,18 @@ from __future__ import annotations from dataclasses import dataclass, field +import warnings from typing import Literal, Optional, get_args -StorageBackend = Literal["auto", "native", "capture", "none"] +StorageBackend = Literal[ + "in-memory", "persistent", "none", "auto", "native", "capture"] + +# What a user chooses; the configurator emits exactly one of these. "auto" is +# the unset default, and "native"/"capture" are the deprecated earlier names. +USER_STORAGE_CHOICES = ("in-memory", "persistent", "none") + +_DEPRECATED_STORAGE_BACKENDS = {"native": "in-memory", "capture": "persistent"} @dataclass @@ -71,29 +79,36 @@ class MonitoringConfig: schedule: CaptureSchedule = field(default_factory=CaptureSchedule) - # Which of the two storage paths this engine is for. They are mutually - # exclusive per record runtime and share nothing but the record envelope: + # Where captured records go. The user's choice is one of three: + # + # "in-memory" -- records are delivered in memory. Until its consumer + # interface exists this is the C++ + # ``ClickHouseRecordSink``: one ClickHouse row per record + # with the tensor bytes inline. Requires ``host_engine`` + # or ``db_config`` on the engine, and refuses an explicit + # ``record_sink``. + # "persistent" -- the native capture storage path: immutable packs in + # object storage with the ClickHouse catalog as an index + # over them. Defaults to the native pack writer (built + # from ``capture_sink_config``); an explicit + # ``record_sink`` overrides it -- pass the reference sink + # to roll back. Refuses a host engine, which would + # otherwise sit started, connected and unused. + # "none" -- capture off entirely: the engine allocates no ring, and + # refuses a record runtime, a host engine and adapter + # attachment. # - # "native" -- the C++ ``ClickHouseRecordSink``: one ClickHouse row per - # record with the tensor bytes inline, into its own table. - # Requires ``host_engine`` or ``db_config`` on the engine, - # and refuses an explicit ``record_sink``. - # "capture" -- the capture path: immutable packs in object storage - # with the catalog as an index over them. Defaults to the - # native pack writer (built from ``capture_sink_config``); - # an explicit ``record_sink`` overrides it — pass the - # reference sink to roll back. Refuses a host engine, - # which would otherwise sit started, connected and unused. - # "none" -- capture and transport with no persistence at all. - # "auto" -- infer from what was passed, which is what every caller - # did before this field existed and remains the default. + # "auto", the default, is not a choice but its absence: infer from what + # was passed, which is what every caller did before this field existed. + # "native" and "capture" are the earlier names of "in-memory" and + # "persistent"; they still work and warn. # # The point of declaring it is that a mismatch becomes an error instead of # a silent choice: passing a ``record_sink`` while a host engine is # configured writes packs and leaves a ClickHouse insert pipeline running # that nothing feeds, and "auto" cannot tell that apart from intent. - # The capture backend's default writer bounds, when ``storage_backend`` - # is "capture" and the caller passes no ``record_sink``. Typed, not + # The persistent path's default writer bounds, when ``storage_backend`` + # is "persistent" and the caller passes no ``record_sink``. Typed, not # Any: the engine validates it at the boundary, and the type is the # contract's documentation. The native pack # sink is the default writer (D4's flip, post-Checkpoint-B); passing a @@ -106,51 +121,63 @@ class MonitoringConfig: # the sink stages to the object store and indexes it into the ClickHouse # catalog, and ``flush_and_wait`` returns only once they are queryable. # Unset, packs stay in the spool for something else to drain. Needs - # ``storage_backend="capture"`` and ``capture_sink_config``, whose spool - # it drains. + # ``storage_backend="persistent"`` and ``capture_sink_config``, whose + # spool it drains. capture_storage_config: Optional["NativeCaptureStorageConfig"] = None storage_backend: StorageBackend = "auto" + @property + def canonical_storage_backend(self) -> str: + """``storage_backend`` with a deprecated name replaced by its new one.""" + return _DEPRECATED_STORAGE_BACKENDS.get( + self.storage_backend, self.storage_backend) + def __post_init__(self) -> None: - backends = get_args(StorageBackend) - if self.storage_backend not in backends: + if self.storage_backend not in get_args(StorageBackend): raise ValueError( "storage_backend must be one of " - + ", ".join(repr(name) for name in backends) - + f"; got {self.storage_backend!r}" + + ", ".join(repr(name) for name in USER_STORAGE_CHOICES) + + " (or left unset); got " + + repr(self.storage_backend) + ) + replacement = _DEPRECATED_STORAGE_BACKENDS.get(self.storage_backend) + if replacement is not None: + warnings.warn( + f"storage_backend={self.storage_backend!r} is deprecated; use " + f"{replacement!r}, which means the same", + DeprecationWarning, + stacklevel=3, ) + backend = self.canonical_storage_backend # ``capture_sink_config`` is read in exactly one place -- the - # engine's default-writer branch, which tests ``storage_backend == - # "capture"`` literally. Nothing RESOLVES "auto" into "capture": - # "auto" is the pre-field behaviour, and it reaches that branch as - # "auto" and falls straight through. So under any other backend a - # configured sink is not merely unused, it is unreachable, and the - # engine's type check at the boundary passes a config that then - # writes no packs at all -- the silent no-op this field's own design - # note ("a mismatch becomes an error instead of a silent choice") - # exists to prevent. + # engine's default-writer branch, which tests for "persistent" + # literally. Nothing RESOLVES "auto" into "persistent": "auto" is the + # pre-field behaviour, and it reaches that branch as "auto" and falls + # straight through. So under any other backend a configured sink is + # not merely unused, it is unreachable, and the engine's type check + # at the boundary passes a config that then writes no packs at all -- + # the silent no-op this field's own design note ("a mismatch becomes + # an error instead of a silent choice") exists to prevent. # # Deliberate behaviour change: a caller who passes the sink config - # with a non-capture backend gets a startup error where they used to - # get silence. That is the trade the note asks for -- the alternative - # is a run that captures nothing and says so nowhere. - if self.capture_sink_config is not None and ( - self.storage_backend != "capture" - ): + # with a non-persistent backend gets a startup error where they used + # to get silence. That is the trade the note asks for -- the + # alternative is a run that captures nothing and says so nowhere. + if self.capture_sink_config is not None and backend != "persistent": raise ValueError( - "capture_sink_config configures the capture backend's " + "capture_sink_config configures the persistent path's " "default pack writer and is read only when " - "storage_backend='capture'; got storage_backend=" + "storage_backend='persistent'; got storage_backend=" f"{self.storage_backend!r}, under which it would be silently " "ignored and nothing would be written. Set " - "storage_backend='capture', or drop capture_sink_config" + "storage_backend='persistent', or drop capture_sink_config" ) if self.capture_storage_config is not None and ( self.capture_sink_config is None ): raise ValueError( "capture_storage_config drains the spool capture_sink_config " - "stages into; set storage_backend='capture' and " + "stages into; set storage_backend='persistent' and " "capture_sink_config=NativeSinkConfig(...) with it" ) diff --git a/src/dmi/engine.py b/src/dmi/engine.py index e4388b5bd..bbd740b68 100644 --- a/src/dmi/engine.py +++ b/src/dmi/engine.py @@ -119,7 +119,12 @@ def __init__( if host_engine is not None and db_config is not None: raise ValueError("Provide either host_engine or db_config, not both") - self._storage_backend = getattr(config, "storage_backend", "auto") + # The canonical name: a deprecated one ("native", "capture") is + # replaced here, so everything below reads only the current three + # choices and "auto". + self._storage_backend = getattr( + config, "canonical_storage_backend", + getattr(config, "storage_backend", "auto")) # A None config is the ctor's documented no-configuration mode; # say so here rather than behind a getattr default. self._capture_sink_config = ( @@ -143,13 +148,20 @@ def __init__( # The running storage service, while a record runtime is attached. self._capture_storage: Optional[Any] = None host_configured = host_engine is not None or db_config is not None - if self._storage_backend == "native" and not host_configured: + if self._storage_backend == "in-memory" and not host_configured: raise ValueError( - "config.storage_backend='native' selects the C++ " - "ClickHouseRecordSink, which needs a host engine: pass " - "host_engine= or db_config=" + "config.storage_backend='in-memory' delivers records through " + "the C++ ClickHouseRecordSink until its consumer interface " + "exists, which needs a host engine: pass host_engine= or " + "db_config=" ) - if self._storage_backend in ("capture", "none") and host_configured: + if self._storage_backend == "none" and ring_config is not None: + raise ValueError( + "config.storage_backend='none' turns capture off, so the " + "engine allocates no ring, but a ring_config was given. Pick " + "'in-memory' or 'persistent' to capture, or drop ring_config" + ) + if self._storage_backend in ("persistent", "none") and host_configured: raise ValueError( f"config.storage_backend={self._storage_backend!r} does not " "use the C++ ClickHouse host, but host_engine/db_config was " @@ -185,6 +197,10 @@ def __init__( except Exception as exc: raise RuntimeError("Failed to start host_engine") from exc + # Capture off means no ring at all: none of its pinned staging and + # payload memory is allocated. + if self._storage_backend == "none": + enable_ring_transport = False if enable_ring_transport or ring_config is not None: if ring_config is None: ring_config = self._make_default_ring_config( @@ -262,15 +278,15 @@ def _reject_a_sink_the_config_did_not_ask_for( backend = getattr(self, "_storage_backend", "auto") if backend == "auto": return - if backend == "capture" and record_sink is None: + if backend == "persistent" and record_sink is None: raise ValueError( - "config.storage_backend='capture' selects the object-store " + "config.storage_backend='persistent' selects the object-store " "path, whose default writer is the native pack sink — pass " "capture_sink_config=NativeSinkConfig(...) in the config, or " "a record_sink explicitly (the reference sink is the " "documented rollback)" ) - if backend in ("native", "none") and record_sink is not None: + if backend in ("in-memory", "none") and record_sink is not None: raise ValueError( f"config.storage_backend={backend!r} does not use an explicit " "record_sink; passing one would send records to a backend the " @@ -290,6 +306,12 @@ def create_record_runtime( runtime; the two paths are never active at the same time. """ + if getattr(self, "_storage_backend", "auto") == "none": + raise RuntimeError( + "config.storage_backend='none' turns capture off: this engine " + "has no ring and creates no record runtime. Pick 'in-memory' " + "or 'persistent' to capture" + ) transport = self._ring_transport ring_config = self._ring_config if transport is None or ring_config is None: @@ -322,7 +344,7 @@ def create_record_runtime( def _start_capture_storage(self, *, sweep_spool: bool) -> Optional[Any]: config = self._capture_storage_config - if config is None or self._storage_backend != "capture": + if config is None or self._storage_backend != "persistent": return None from .storage.native_capture import NativeCaptureStorage @@ -346,13 +368,13 @@ def _attach_record_runtime( from .records import RecordRuntime ring_config = self._ring_config - # D4's flip: the capture backend's default writer is the native + # D4's flip: the persistent backend's default writer is the native # pack sink, built from the config's bounds. An explicit # record_sink overrides it — the reference sink is the documented # rollback — and the ClickHouse host path is untouched. if ( record_sink is None - and self._storage_backend == "capture" + and self._storage_backend == "persistent" and self._capture_sink_config is not None ): from .storage.capture.native_sink import create_native_pack_sink diff --git a/src/dmi/storage/capture/native_sink.py b/src/dmi/storage/capture/native_sink.py index 6f6310f31..ade74833a 100644 --- a/src/dmi/storage/capture/native_sink.py +++ b/src/dmi/storage/capture/native_sink.py @@ -6,7 +6,7 @@ the native writer instead: envelopes travel the ring's record worker straight into pack assembly with no Python on the capture path. -Selection defaults to this writer under ``storage_backend="capture"`` +Selection defaults to this writer under ``storage_backend="persistent"`` (built from the config's ``capture_sink_config``), since Checkpoint B and C1/C2 passed. An explicit ``record_sink`` at the call site overrides it — pass the reference sink to roll back. Packs staged by either writer are diff --git a/src/dmi/storage/native_capture.py b/src/dmi/storage/native_capture.py index fea9f733d..db2c90b6e 100644 --- a/src/dmi/storage/native_capture.py +++ b/src/dmi/storage/native_capture.py @@ -1,6 +1,6 @@ """The native capture storage path after the spool, and its query side. -Under ``storage_backend="capture"`` the native pack sink stages immutable +Under ``storage_backend="persistent"`` the native pack sink stages immutable packs in a local spool (``capture_sink_config``). Setting ``capture_storage_config`` as well makes the engine run a :class:`NativeCaptureStorage` in-process: one C++ thread that uploads each diff --git a/tests/test_engine_runtime_api.py b/tests/test_engine_runtime_api.py index 09e6fada4..963c82fe5 100644 --- a/tests/test_engine_runtime_api.py +++ b/tests/test_engine_runtime_api.py @@ -642,15 +642,15 @@ def test_storage_backend_rejects_an_unknown_name(): MonitoringConfig(storage_backend="object-store") -@pytest.mark.parametrize("backend", ["auto", "native", "none"]) +@pytest.mark.parametrize("backend", ["auto", "in-memory", "none"]) def test_a_sink_config_under_a_backend_that_cannot_use_it_is_refused( backend: str, ): """The combination that used to write nothing and say nothing. ``capture_sink_config`` is consumed in exactly one branch, guarded by - ``storage_backend == "capture"``. Under any other backend -- the DEFAULT - "auto" included, which is never resolved into "capture" anywhere -- the + ``storage_backend == "persistent"``. Under any other backend -- the DEFAULT + "auto" included, which is never resolved into "persistent" anywhere -- the config passed the engine's type check and then fell through ``create_record_runtime`` with ``record_sink=None``: no packs written, no diagnostic. A startup error instead, naming both fields. @@ -676,7 +676,7 @@ def test_the_capture_backend_still_takes_its_sink_config(): sink_config = NativeSinkConfig(spool_root="/tmp/unused") config = MonitoringConfig( - storage_backend="capture", capture_sink_config=sink_config + storage_backend="persistent", capture_sink_config=sink_config ) assert config.capture_sink_config is sink_config # And a bare config is untouched: the refusal is about a sink config that @@ -689,13 +689,13 @@ def test_native_backend_requires_a_host_engine(): silently transport-only engine.""" with pytest.raises(ValueError, match="needs a host engine"): MonitoringEngine( - config=_config("native"), + config=_config("in-memory"), model_id="m", enable_ring_transport=False, ) -@pytest.mark.parametrize("backend", ["capture", "none"]) +@pytest.mark.parametrize("backend", ["persistent", "none"]) def test_a_non_native_backend_refuses_a_host_engine(backend: str): """Configured together, the host engine starts, connects and is never fed. @@ -717,7 +717,7 @@ def test_capture_backend_requires_the_sink_that_selects_it(monkeypatch): the one that actually routes the records.""" engine, _old, _ring = _engine_with_fake_ring() engine._ring_config = object() - engine._storage_backend = "capture" + engine._storage_backend = "persistent" with pytest.raises(ValueError, match="capture_sink_config"): engine.create_record_runtime(_explicit_sink_format()) @@ -726,7 +726,7 @@ def test_capture_backend_requires_the_sink_that_selects_it(monkeypatch): def test_native_backend_refuses_an_explicit_sink(monkeypatch): engine, _old, _ring = _engine_with_fake_ring() engine._ring_config = object() - engine._storage_backend = "native" + engine._storage_backend = "in-memory" fake_native_module = ModuleType("dmi.transport.native") @@ -771,7 +771,7 @@ def test_capture_backend_defaults_to_the_native_pack_sink(monkeypatch, tmp_path) engine, _old_transport, old_ring = _engine_with_fake_ring() ring_config = object() engine._ring_config = ring_config - engine._storage_backend = "capture" + engine._storage_backend = "persistent" sink_config = NativeSinkConfig(spool_root=str(tmp_path / "spool")) diff --git a/tests/test_hf_capture_refusal.py b/tests/test_hf_capture_refusal.py index f5a3ebe6a..41190b55f 100644 --- a/tests/test_hf_capture_refusal.py +++ b/tests/test_hf_capture_refusal.py @@ -1,4 +1,4 @@ -"""HF under ``storage_backend="capture"`` must fail loudly, never store nothing. +"""HF under ``storage_backend="persistent"`` must fail loudly, never store nothing. Two reproductions from the capture-path audit, both of which "succeeded": @@ -194,12 +194,12 @@ def _assert_untouched(model, engine): # --------------------------------------------------------------------------- -# Link 6: the capture config is refused at every HF entry point +# Link 6: the persistent config is refused at every HF entry point # --------------------------------------------------------------------------- def test_hf_attach_model_refuses_capture_storage(): - engine = _SpyEngine("capture") + engine = _SpyEngine("persistent") model = _TinyHookedLM(engine) with pytest.raises(ConfigurationError, @@ -216,7 +216,7 @@ def test_base_attach_model_refuses_capture_storage_for_any_adapter( via_attach_config): """The legacy HookPoint path is what every adapter's attach installs, so the refusal is in the base class, not only in the HF override.""" - engine = _SpyEngine("capture") + engine = _SpyEngine("persistent") model = _TinyHookedLM(engine) adapter = _StubAdapter(engine, "tiny") @@ -232,7 +232,7 @@ def test_base_attach_model_refuses_capture_storage_for_any_adapter( def test_generate_with_monitoring_refuses_capture_storage(): - engine = _SpyEngine("capture") + engine = _SpyEngine("persistent") model = _TinyHookedLM(engine) input_ids, attention_mask = _inputs() @@ -248,7 +248,7 @@ def test_generate_with_monitoring_refuses_capture_storage(): def test_generate_greedy_with_monitoring_refuses_capture_storage(): - engine = _SpyEngine("capture") + engine = _SpyEngine("persistent") model = _TinyHookedLM(engine) input_ids, attention_mask = _inputs() @@ -264,9 +264,9 @@ def test_generate_greedy_with_monitoring_refuses_capture_storage(): _assert_untouched(model, engine) -@pytest.mark.parametrize("backend", ["auto", "native", "none"]) +@pytest.mark.parametrize("backend", ["auto", "in-memory"]) def test_other_storage_backends_still_attach(backend): - """Only 'capture' is refused; the legacy backends are unchanged.""" + """Only 'persistent' and 'none' are refused; the rest attach as before.""" engine = _SpyEngine(backend) model = _TinyHookedLM(engine) adapter = HuggingFaceAdapter(engine, "tiny") @@ -313,7 +313,7 @@ def test_capture_mode_failure_propagates_out_of_the_prepare_wrapper(): model = _TinyHookedLM(engine) adapter = HuggingFaceAdapter(engine, "tiny") adapter.attach_model(model) - engine._storage_backend = "capture" + engine._storage_backend = "persistent" input_ids, attention_mask = _inputs() try: diff --git a/tests/test_integration_api_v1.py b/tests/test_integration_api_v1.py index 590241cb5..d1c1fd096 100644 --- a/tests/test_integration_api_v1.py +++ b/tests/test_integration_api_v1.py @@ -272,7 +272,7 @@ def test_v1_documents_the_capture_and_record_mode_refusals() -> None: document = (root / "docs" / "integration-api-v1.md").read_text() attach = _document_section(document, "#### `attach_model(") - assert 'storage_backend="capture"' in attach + assert 'storage_backend="persistent"' in attach assert "`dmi.configuration.ConfigurationError`" in attach for entry in ("generate_with_monitoring", "generate_greedy_with_monitoring"): assert f"`{entry}()`" in attach diff --git a/tests/test_native_capture_storage_gpu_e2e.py b/tests/test_native_capture_storage_gpu_e2e.py index 1e76424c0..da94c41b0 100644 --- a/tests/test_native_capture_storage_gpu_e2e.py +++ b/tests/test_native_capture_storage_gpu_e2e.py @@ -1,7 +1,7 @@ """GPU -> Ring -> NativePackSink -> storage service -> catalog -> reader. The whole native capture storage path on a real engine, selected by config -alone: ``storage_backend="capture"`` with ``capture_sink_config`` and +alone: ``storage_backend="persistent"`` with ``capture_sink_config`` and ``capture_storage_config``. ``flush_and_wait`` returning is the promise under test -- every record captured before it is queryable in the catalog and reads back byte-identical to the CUDA tensor it came from -- with no Python between @@ -119,7 +119,7 @@ def test_flush_and_wait_means_queryable_and_byte_identical(fake_s3, tmp_path): poll_interval_s=0.05, ) config = MonitoringConfig( - storage_backend="capture", + storage_backend="persistent", capture_sink_config=NativeSinkConfig( spool_root=str(tmp_path / "spool"), max_pack_records=2, max_linger_ns=60_000_000_000), @@ -201,7 +201,7 @@ def test_close_alone_delivers_the_tail_to_the_catalog(fake_s3, tmp_path): poll_interval_s=0.05, ) config = MonitoringConfig( - storage_backend="capture", + storage_backend="persistent", capture_sink_config=NativeSinkConfig( spool_root=str(tmp_path / "spool"), max_pack_records=2, max_linger_ns=60_000_000_000), diff --git a/tests/test_native_capture_storage_wiring.py b/tests/test_native_capture_storage_wiring.py index 473dab70e..ce59315b4 100644 --- a/tests/test_native_capture_storage_wiring.py +++ b/tests/test_native_capture_storage_wiring.py @@ -38,7 +38,7 @@ def test_storage_config_needs_the_sink_config_whose_spool_it_drains(): from dmi.config import MonitoringConfig with pytest.raises(ValueError, match="capture_sink_config"): - MonitoringConfig(storage_backend="capture", + MonitoringConfig(storage_backend="persistent", capture_storage_config=_storage_config()) @@ -48,7 +48,7 @@ def test_storage_config_is_accepted_beside_a_sink_config(tmp_path): storage = _storage_config() config = MonitoringConfig( - storage_backend="capture", + storage_backend="persistent", capture_sink_config=NativeSinkConfig(spool_root=str(tmp_path)), capture_storage_config=storage, ) @@ -158,7 +158,7 @@ def test_read_accepts_a_zero_byte_limit(monkeypatch): def test_engine_refuses_a_storage_config_of_the_wrong_type(): - config = SimpleNamespace(storage_backend="capture", capture_sink_config=None, + config = SimpleNamespace(storage_backend="persistent", capture_sink_config=None, capture_storage_config={"s3_bucket": "b"}) with pytest.raises(TypeError, match="NativeCaptureStorageConfig"): MonitoringEngine(config=config, enable_ring_transport=False) @@ -197,7 +197,7 @@ def rethrow_if_failed(self): def _capture_engine(monkeypatch, tmp_path, *, fail_ring=False): - """An engine under storage_backend="capture" with both native modules + """An engine under storage_backend="persistent" with both native modules faked. Returns (engine, events, services).""" from dmi.storage.capture.native_sink import NativeSinkConfig @@ -207,7 +207,7 @@ def _capture_engine(monkeypatch, tmp_path, *, fail_ring=False): force_eager=False) engine._ring_engine = SimpleNamespace(stop=lambda: None) engine._ring_config = object() - engine._storage_backend = "capture" + engine._storage_backend = "persistent" engine._capture_sink_config = NativeSinkConfig( spool_root=str(tmp_path / "spool"), spool_max_bytes=1 << 30) engine._capture_storage_config = _storage_config() diff --git a/tests/test_native_sink_ring_e2e.py b/tests/test_native_sink_ring_e2e.py index a44f14f46..458892429 100644 --- a/tests/test_native_sink_ring_e2e.py +++ b/tests/test_native_sink_ring_e2e.py @@ -278,7 +278,7 @@ def metadata_and_bytes(step): def test_the_capture_backend_selects_the_native_sink_by_config(tmp_path: Path): - """The automatic entry point: storage_backend="capture" + capture_sink_config.""" + """The automatic entry point: storage_backend="persistent" + capture_sink_config.""" from dmi.api.v1 import ( HookPointV1, HookSpecV1, MonitoringEngine, TransportSpec, ) @@ -288,7 +288,7 @@ def test_the_capture_backend_selects_the_native_sink_by_config(tmp_path: Path): spool_root = tmp_path / "spool" config = MonitoringConfig( - storage_backend="capture", + storage_backend="persistent", capture_sink_config=NativeSinkConfig( spool_root=str(spool_root), max_pack_records=8, max_linger_ns=1_000_000_000), diff --git a/tests/test_storage_backend_choices.py b/tests/test_storage_backend_choices.py new file mode 100644 index 000000000..d7276461c --- /dev/null +++ b/tests/test_storage_backend_choices.py @@ -0,0 +1,127 @@ +"""The storage choice a user makes: ``in-memory``, ``persistent`` or ``none``. + +``persistent`` is the native capture storage path (object store + catalog + +ClickHouse). ``in-memory`` delivers records in memory; until its consumer +interface exists it runs today's ``ClickHouseRecordSink`` path, which needs +a host engine. ``none`` turns capture off entirely: no ring is allocated, and +nothing that would capture can be attached. + +The earlier names ``native`` and ``capture`` still work, with a +DeprecationWarning, and mean ``in-memory`` and ``persistent``. ``auto`` is the +unset default -- the inference every caller made before the field existed -- +and is not one of the user's choices. +""" + +from __future__ import annotations + +import warnings + +import pytest + +from dmi.config import USER_STORAGE_CHOICES, MonitoringConfig, StorageBackend +from dmi.engine import MonitoringEngine + +pytestmark = pytest.mark.cpu + + +def test_the_user_choices_are_in_memory_persistent_and_none(): + assert USER_STORAGE_CHOICES == ("in-memory", "persistent", "none") + for choice in USER_STORAGE_CHOICES: + with warnings.catch_warnings(): + warnings.simplefilter("error") # none of them is deprecated + assert MonitoringConfig(storage_backend=choice).storage_backend == choice + + +def test_auto_is_the_unset_default_not_a_user_choice(): + assert MonitoringConfig().storage_backend == "auto" + assert "auto" not in USER_STORAGE_CHOICES + + +@pytest.mark.parametrize( + "legacy, canonical", [("native", "in-memory"), ("capture", "persistent")]) +def test_the_old_names_still_work_and_say_what_replaced_them(legacy, canonical): + with pytest.warns(DeprecationWarning, match=f"'{canonical}'"): + config = MonitoringConfig(storage_backend=legacy) + # The config keeps what the caller wrote; the engine acts on the new name. + assert config.storage_backend == legacy + assert config.canonical_storage_backend == canonical + + +def test_an_unknown_choice_names_the_three_real_ones(): + with pytest.raises(ValueError) as refused: + MonitoringConfig(storage_backend="object-store") + message = str(refused.value) + for choice in USER_STORAGE_CHOICES: + assert repr(choice) in message + + +def test_every_accepted_value_is_in_the_type(): + for value in ("in-memory", "persistent", "none", "auto", "native", "capture"): + assert value in StorageBackend.__args__ + + +def test_persistent_takes_the_capture_sink_config(tmp_path): + from dmi.storage.capture.native_sink import NativeSinkConfig + + sink = NativeSinkConfig(spool_root=str(tmp_path)) + config = MonitoringConfig(storage_backend="persistent", capture_sink_config=sink) + assert config.capture_sink_config is sink + with pytest.raises(ValueError, match="storage_backend='persistent'"): + MonitoringConfig(storage_backend="in-memory", capture_sink_config=sink) + + +# --- none: capture off entirely --------------------------------------------- + + +def test_none_allocates_no_ring(monkeypatch): + def _must_not_run(*args, **kwargs): + raise AssertionError("storage_backend='none' allocated a ring") + + monkeypatch.setattr(MonitoringEngine, "enable_ring_transport", _must_not_run) + monkeypatch.setattr(MonitoringEngine, "_make_default_ring_config", + staticmethod(_must_not_run)) + + engine = MonitoringEngine(config=MonitoringConfig(storage_backend="none")) + + assert engine._ring_transport is None + assert engine.capture_enabled is False + + +def test_none_refuses_an_explicit_ring_config(): + with pytest.raises(ValueError, match="capture off"): + MonitoringEngine(config=MonitoringConfig(storage_backend="none"), + ring_config=object()) + + +def test_none_refuses_a_record_runtime(): + engine = MonitoringEngine(config=MonitoringConfig(storage_backend="none")) + + with pytest.raises(RuntimeError, match="storage_backend='none' turns capture off"): + engine.create_record_runtime(object()) + + +def test_none_refuses_a_host_engine(): + with pytest.raises(ValueError, match="'none'"): + MonitoringEngine(config=MonitoringConfig(storage_backend="none"), + model_id="m", host_engine=object()) + + +def test_an_adapter_refuses_to_attach_when_capture_is_off(): + from dmi.adapters.base import _refuse_unwired_capture_storage + from dmi.configuration.errors import ConfigurationError + + engine = MonitoringEngine(config=MonitoringConfig(storage_backend="none")) + + with pytest.raises(ConfigurationError, match="capture is off"): + _refuse_unwired_capture_storage(engine, "attach_model()", "SomeAdapter") + + +def test_an_adapter_refuses_persistent_by_its_new_name(): + from dmi.adapters.base import _refuse_unwired_capture_storage + from dmi.configuration.errors import ConfigurationError + + engine = MonitoringEngine(enable_ring_transport=False) + engine._storage_backend = "persistent" + + with pytest.raises(ConfigurationError, match="storage_backend='persistent'"): + _refuse_unwired_capture_storage(engine, "attach_model()", "SomeAdapter") From 2f51af5ede06a40ba3ce7a6159d11a577bf8e7c9 Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Wed, 23 Sep 2026 19:30:24 -0400 Subject: [PATCH 2/5] Refuse a ring under none after construction too; pin the old names in the engine Two findings from an independent review: - storage_backend="none" skipped only the constructor's ring. The public enable_ring_transport(), which the adapter docs point callers at, still allocated one and made it globally active. It now refuses under "none" with the same message as create_record_runtime. - Nothing tested that a deprecated name reaches the engine: reverting the engine's mapping to the raw field left all 85 related tests green. The mapping is now a module-level canonical_storage_backend() the engine applies to whatever config it is given, so a duck-typed config carrying "native" or "capture" is checked too (before, it slipped past every check). New tests drive the engine with each old name: "native" without a host is refused, "capture" with a host is refused, "capture" with a sink config resolves to "persistent", and a duck-typed "native" is refused. With the old raw read put back, 3 of them fail. CPU tier: 2242 passed, 1 skipped (no CUDA device). GPU (RTX 4090) capture storage, sink ring, record-ring refusal and storage-choice suites: 33 passed with DeprecationWarning as an error. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- src/dmi/config.py | 26 +++++++++------ src/dmi/engine.py | 17 +++++++--- tests/test_storage_backend_choices.py | 46 +++++++++++++++++++++++++++ 3 files changed, 76 insertions(+), 13 deletions(-) diff --git a/src/dmi/config.py b/src/dmi/config.py index d37dd5beb..471071273 100644 --- a/src/dmi/config.py +++ b/src/dmi/config.py @@ -17,6 +17,21 @@ _DEPRECATED_STORAGE_BACKENDS = {"native": "in-memory", "capture": "persistent"} +def canonical_storage_backend(name: str) -> str: + """``name`` with a deprecated storage backend replaced by its new one. + + Unknown names are refused here as well as in ``MonitoringConfig``, so an + engine handed any config-like object acts only on the current names. + """ + if name not in get_args(StorageBackend): + raise ValueError( + "storage_backend must be one of " + + ", ".join(repr(choice) for choice in USER_STORAGE_CHOICES) + + f" (or left unset); got {name!r}" + ) + return _DEPRECATED_STORAGE_BACKENDS.get(name, name) + + @dataclass class CaptureSchedule: """Schedule for step-level and request-level capture.""" @@ -130,17 +145,10 @@ class MonitoringConfig: @property def canonical_storage_backend(self) -> str: """``storage_backend`` with a deprecated name replaced by its new one.""" - return _DEPRECATED_STORAGE_BACKENDS.get( - self.storage_backend, self.storage_backend) + return canonical_storage_backend(self.storage_backend) def __post_init__(self) -> None: - if self.storage_backend not in get_args(StorageBackend): - raise ValueError( - "storage_backend must be one of " - + ", ".join(repr(name) for name in USER_STORAGE_CHOICES) - + " (or left unset); got " - + repr(self.storage_backend) - ) + canonical_storage_backend(self.storage_backend) # refuses unknowns replacement = _DEPRECATED_STORAGE_BACKENDS.get(self.storage_backend) if replacement is not None: warnings.warn( diff --git a/src/dmi/engine.py b/src/dmi/engine.py index bbd740b68..702c40e42 100644 --- a/src/dmi/engine.py +++ b/src/dmi/engine.py @@ -120,10 +120,12 @@ def __init__( raise ValueError("Provide either host_engine or db_config, not both") # The canonical name: a deprecated one ("native", "capture") is - # replaced here, so everything below reads only the current three - # choices and "auto". - self._storage_backend = getattr( - config, "canonical_storage_backend", + # replaced here, for a MonitoringConfig and any config-like object + # alike, so everything below reads only the current three choices + # and "auto". + from .config import canonical_storage_backend + + self._storage_backend = canonical_storage_backend( getattr(config, "storage_backend", "auto")) # A None config is the ctor's documented no-configuration mode; # say so here rather than behind a getattr default. @@ -489,6 +491,8 @@ def enable_ring_transport( ) -> Any: """Switch to ring-based D2H transport. + Refused under ``storage_backend="none"``, which turns capture off. + Creates a RingEngine with the C++ host engine as the submit target so tensor reconstruction, slicing, and DB submission all happen in C++ without the GIL. @@ -505,6 +509,11 @@ def enable_ring_transport( ``self._ring_transport``). Returned so adapters can hold a direct reference instead of reaching through the engine. """ + if getattr(self, "_storage_backend", "auto") == "none": + raise RuntimeError( + "config.storage_backend='none' turns capture off: this engine " + "allocates no ring. Pick 'in-memory' or 'persistent' to capture" + ) _rt = _ring_module() _native_engine = _native_module() diff --git a/tests/test_storage_backend_choices.py b/tests/test_storage_backend_choices.py index d7276461c..5b037645f 100644 --- a/tests/test_storage_backend_choices.py +++ b/tests/test_storage_backend_choices.py @@ -125,3 +125,49 @@ def test_an_adapter_refuses_persistent_by_its_new_name(): with pytest.raises(ConfigurationError, match="storage_backend='persistent'"): _refuse_unwired_capture_storage(engine, "attach_model()", "SomeAdapter") + + +def test_none_refuses_a_ring_enabled_after_construction(): + """The constructor skipping the ring is not enough: the public + enable_ring_transport() would still allocate one and make it active.""" + engine = MonitoringEngine(config=MonitoringConfig(storage_backend="none")) + + with pytest.raises(RuntimeError, match="storage_backend='none' turns capture off"): + engine.enable_ring_transport(object()) + assert engine._ring_transport is None + + +# --- the old names reach the engine, not only the config --------------------- + + +def test_the_engine_acts_on_native_as_in_memory(): + with pytest.warns(DeprecationWarning): + config = MonitoringConfig(storage_backend="native") + with pytest.raises(ValueError, match="needs a host engine"): + MonitoringEngine(config=config, enable_ring_transport=False) + + +def test_the_engine_acts_on_capture_as_persistent(tmp_path): + from dmi.storage.capture.native_sink import NativeSinkConfig + + with pytest.warns(DeprecationWarning): + refused = MonitoringConfig(storage_backend="capture") + with pytest.raises(ValueError, match="does not use the C\\+\\+ ClickHouse host"): + MonitoringEngine(config=refused, model_id="m", host_engine=object(), + enable_ring_transport=False) + + with pytest.warns(DeprecationWarning): + config = MonitoringConfig( + storage_backend="capture", + capture_sink_config=NativeSinkConfig(spool_root=str(tmp_path))) + engine = MonitoringEngine(config=config, enable_ring_transport=False) + assert engine._storage_backend == "persistent" + + +def test_a_duck_typed_config_with_an_old_name_is_still_checked(): + from types import SimpleNamespace + + config = SimpleNamespace(storage_backend="native", capture_sink_config=None, + capture_storage_config=None) + with pytest.raises(ValueError, match="needs a host engine"): + MonitoringEngine(config=config, enable_ring_transport=False) From ea50e4a687cffe63a7708eb4b57d3324a8c0169e Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 09:26:02 -0400 Subject: [PATCH 3/5] Type the canonical storage name, and name the caller in its warnings and errors Three smaller findings from the independent review: - canonical_storage_backend returned a plain str, so every literal comparison in the engine and adapters went unchecked by a type checker. It now returns CanonicalStorageBackend, a Literal of the four names the engine acts on. - The deprecation warning used stacklevel=3, which points at the caller only for direct construction. Through dataclasses.replace() -- which the configurator's _install_schedule uses -- it was attributed to dataclasses.py. The level is now computed by walking past this module, the generated __init__ and dataclasses.py, so both routes name the caller's line. - A refusal quoted the canonical name even when the caller wrote the old one, so "native" produced an error about 'in-memory'. Refusals now read "config.storage_backend='native' (now 'in-memory') ...". Tests: three new ones in test_storage_backend_choices.py, each red first. With the storage, engine and wiring suites, 72 passed. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- src/dmi/config.py | 28 ++++++++++++++++++++---- src/dmi/engine.py | 27 +++++++++++++++++------ tests/test_storage_backend_choices.py | 31 +++++++++++++++++++++++++++ 3 files changed, 76 insertions(+), 10 deletions(-) diff --git a/src/dmi/config.py b/src/dmi/config.py index 471071273..9e0c6547f 100644 --- a/src/dmi/config.py +++ b/src/dmi/config.py @@ -3,6 +3,7 @@ from __future__ import annotations from dataclasses import dataclass, field +import sys import warnings from typing import Literal, Optional, get_args @@ -10,6 +11,10 @@ StorageBackend = Literal[ "in-memory", "persistent", "none", "auto", "native", "capture"] +# The names the engine acts on: a deprecated name is replaced before anything +# compares against these. +CanonicalStorageBackend = Literal["in-memory", "persistent", "none", "auto"] + # What a user chooses; the configurator emits exactly one of these. "auto" is # the unset default, and "native"/"capture" are the deprecated earlier names. USER_STORAGE_CHOICES = ("in-memory", "persistent", "none") @@ -17,7 +22,7 @@ _DEPRECATED_STORAGE_BACKENDS = {"native": "in-memory", "capture": "persistent"} -def canonical_storage_backend(name: str) -> str: +def canonical_storage_backend(name: str) -> CanonicalStorageBackend: """``name`` with a deprecated storage backend replaced by its new one. Unknown names are refused here as well as in ``MonitoringConfig``, so an @@ -29,7 +34,22 @@ def canonical_storage_backend(name: str) -> str: + ", ".join(repr(choice) for choice in USER_STORAGE_CHOICES) + f" (or left unset); got {name!r}" ) - return _DEPRECATED_STORAGE_BACKENDS.get(name, name) + return _DEPRECATED_STORAGE_BACKENDS.get(name, name) # type: ignore[return-value] + + +def _caller_stacklevel() -> int: + """The ``stacklevel`` that points a warning raised in ``__post_init__`` + at user code: past the dataclass-generated ``__init__`` and, for a copy, + past ``dataclasses.replace``.""" + frame = sys._getframe(1) # __post_init__ + level = 1 + while frame is not None and ( + frame.f_code.co_filename in (__file__, "") + or frame.f_code.co_filename.endswith("dataclasses.py") + ): + frame = frame.f_back + level += 1 + return level @dataclass @@ -143,7 +163,7 @@ class MonitoringConfig: storage_backend: StorageBackend = "auto" @property - def canonical_storage_backend(self) -> str: + def canonical_storage_backend(self) -> CanonicalStorageBackend: """``storage_backend`` with a deprecated name replaced by its new one.""" return canonical_storage_backend(self.storage_backend) @@ -155,7 +175,7 @@ def __post_init__(self) -> None: f"storage_backend={self.storage_backend!r} is deprecated; use " f"{replacement!r}, which means the same", DeprecationWarning, - stacklevel=3, + stacklevel=_caller_stacklevel(), ) backend = self.canonical_storage_backend # ``capture_sink_config`` is read in exactly one place -- the diff --git a/src/dmi/engine.py b/src/dmi/engine.py index 702c40e42..efd1fbc3b 100644 --- a/src/dmi/engine.py +++ b/src/dmi/engine.py @@ -125,8 +125,11 @@ def __init__( # and "auto". from .config import canonical_storage_backend - self._storage_backend = canonical_storage_backend( - getattr(config, "storage_backend", "auto")) + # What the caller wrote, kept only to quote it back in a refusal. + self._storage_backend_requested = getattr( + config, "storage_backend", "auto") + self._storage_backend: str = canonical_storage_backend( + self._storage_backend_requested) # A None config is the ctor's documented no-configuration mode; # say so here rather than behind a getattr default. self._capture_sink_config = ( @@ -152,7 +155,8 @@ def __init__( host_configured = host_engine is not None or db_config is not None if self._storage_backend == "in-memory" and not host_configured: raise ValueError( - "config.storage_backend='in-memory' delivers records through " + f"config.storage_backend={self._backend_label()} delivers " + "records through " "the C++ ClickHouseRecordSink until its consumer interface " "exists, which needs a host engine: pass host_engine= or " "db_config=" @@ -165,7 +169,7 @@ def __init__( ) if self._storage_backend in ("persistent", "none") and host_configured: raise ValueError( - f"config.storage_backend={self._storage_backend!r} does not " + f"config.storage_backend={self._backend_label()} does not " "use the C++ ClickHouse host, but host_engine/db_config was " "given. Configured together, the host engine starts, connects " "and is then never fed, because a record runtime is handed " @@ -262,6 +266,15 @@ def set_capture_enabled(self, enabled: bool) -> None: # lifecycle toggle; the next committed step recomputes it. transport.force_eager = False + def _backend_label(self) -> str: + """The storage backend as the caller wrote it, with its current name + when that was a deprecated one: ``'native' (now 'in-memory')``.""" + requested = getattr(self, "_storage_backend_requested", + self._storage_backend) + if requested == self._storage_backend: + return repr(self._storage_backend) + return f"{requested!r} (now {self._storage_backend!r})" + def _reject_a_sink_the_config_did_not_ask_for( self, record_sink: Optional[Any] ) -> None: @@ -282,7 +295,8 @@ def _reject_a_sink_the_config_did_not_ask_for( return if backend == "persistent" and record_sink is None: raise ValueError( - "config.storage_backend='persistent' selects the object-store " + f"config.storage_backend={self._backend_label()} selects the " + "object-store " "path, whose default writer is the native pack sink — pass " "capture_sink_config=NativeSinkConfig(...) in the config, or " "a record_sink explicitly (the reference sink is the " @@ -290,7 +304,8 @@ def _reject_a_sink_the_config_did_not_ask_for( ) if backend in ("in-memory", "none") and record_sink is not None: raise ValueError( - f"config.storage_backend={backend!r} does not use an explicit " + f"config.storage_backend={self._backend_label()} does not use " + "an explicit " "record_sink; passing one would send records to a backend the " "configuration did not select" ) diff --git a/tests/test_storage_backend_choices.py b/tests/test_storage_backend_choices.py index 5b037645f..39eb3e1a5 100644 --- a/tests/test_storage_backend_choices.py +++ b/tests/test_storage_backend_choices.py @@ -171,3 +171,34 @@ def test_a_duck_typed_config_with_an_old_name_is_still_checked(): capture_storage_config=None) with pytest.raises(ValueError, match="needs a host engine"): MonitoringEngine(config=config, enable_ring_transport=False) + + +# --- the details the old names carry ------------------------------------------ + + +def test_the_canonical_name_has_its_own_type(): + from typing import get_args + + from dmi.config import CanonicalStorageBackend + + assert set(get_args(CanonicalStorageBackend)) == { + "in-memory", "persistent", "none", "auto"} + + +def test_the_warning_names_the_callers_line_even_through_replace(): + import dataclasses + + with pytest.warns(DeprecationWarning) as direct: + config = MonitoringConfig(storage_backend="native") + assert direct[0].filename == __file__ + + with pytest.warns(DeprecationWarning) as replaced: + dataclasses.replace(config) + assert replaced[0].filename == __file__ + + +def test_a_refusal_quotes_the_name_the_caller_wrote(): + with pytest.warns(DeprecationWarning): + config = MonitoringConfig(storage_backend="native") + with pytest.raises(ValueError, match=r"'native' \(now 'in-memory'\)"): + MonitoringEngine(config=config, enable_ring_transport=False) From 81dac0b3687b14c4bd6f353df8c918f4dbb2aebf Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 09:26:02 -0400 Subject: [PATCH 4/5] Bring the storage docs in line with the persistent path as it ships The design doc's comparison table still headed its columns "native path" and "capture path", and described the persistent side as the Python reference sink, "explicitly reference-only". Since the native pack writer became the default and the in-process storage service landed, that read as "persistent = Python reference-only". The table now names both paths by their storage choice and describes the native writer and storage service as production, the Python sink as the reference and rollback. The v1 contract's list of refusals left out "none" with a ring_config (ValueError at construction) and did not say that "none" refuses create_record_runtime and enable_ring_transport with RuntimeError and an adaptor's attach_model with ConfigurationError. It now lists each with its exception type. Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- docs/capture-storage-design.md | 8 ++++---- docs/integration-api-v1.md | 14 +++++++++----- 2 files changed, 13 insertions(+), 9 deletions(-) diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index 7e6bf0586..cfd0f4cfc 100644 --- a/docs/capture-storage-design.md +++ b/docs/capture-storage-design.md @@ -99,15 +99,15 @@ upload-only role with a single publisher, which is planned but not built. This document describes ONE of two storage paths, and they are mutually exclusive per record runtime: -| | native path | capture path (this document) | +| | in-memory path | persistent path (this document) | |---|---|---| -| Sink | `ClickHouseRecordSink` (C++) | `ReferencePythonCaptureSink` (C++ bridge) → `CapturePackReferenceSink` (Python) | +| Sink | `ClickHouseRecordSink` (C++), until the in-memory consumer interface exists | `NativePackSink` (C++), the default writer; the Python `CapturePackReferenceSink` remains as the reference and rollback, passed explicitly as `record_sink` | | Durable form | one ClickHouse row per record, tensor bytes inline | immutable packs in object storage | | ClickHouse role | the store itself | a rebuildable index over the packs | | ClickHouse footprint | one configured table (`offload` by default) | `{prefix}_*` (`dmi_*` by default) | -| Selected by | `create_record_runtime(fmt)` with no `record_sink`, plus a `host_engine`/`db_config` on the engine | `create_record_runtime(fmt, record_sink=reference.native_sink)` | +| Selected by | `create_record_runtime(fmt)` with no `record_sink`, plus a `host_engine`/`db_config` on the engine | `create_record_runtime(fmt)` with `capture_sink_config` in the config (the native writer), plus `capture_storage_config` to upload and index in-process | | Declared by | `MonitoringConfig(storage_backend="in-memory")` (formerly `"native"`) | `MonitoringConfig(storage_backend="persistent")` (formerly `"capture"`) | -| Status | production | explicitly reference-only; production sinks remain native-only | +| Status | production | production native writer and storage service; the Python sink is reference-only | Both are `ring::RecordSink` implementations and the record engine takes exactly one of them: `RingEngine.create_record` is handed either the host engine or a diff --git a/docs/integration-api-v1.md b/docs/integration-api-v1.md index c6a9845f9..5384982ca 100644 --- a/docs/integration-api-v1.md +++ b/docs/integration-api-v1.md @@ -155,11 +155,15 @@ caller did before the field existed. The earlier names `"native"` and `"persistent"`. `MonitoringConfig.canonical_storage_backend` gives the current name. -All three fields are acted on by `MonitoringEngine`, and some combinations are -refused at construction -- `"in-memory"` without a host engine, -`"persistent"`/`"none"` with one, or `capture_storage_config` without -`capture_sink_config` -- so a caller setting them should expect `ValueError` -rather than a silent choice. +All three fields are acted on by `MonitoringEngine`, and a mismatch is an error +rather than a silent choice: + +- At construction, `ValueError`: `"in-memory"` without a host engine, + `"persistent"` or `"none"` with one, `"none"` with a `ring_config`, or + `capture_storage_config` without `capture_sink_config`. +- Under `"none"`, `RuntimeError` from `create_record_runtime()` and + `enable_ring_transport()`, and `ConfigurationError` from an adaptor's + `attach_model()`. `capture_sink_config` is read only when `storage_backend` is `"persistent"`; see `docs/capture-storage-design.md` for the writer it selects. With `capture_storage_config` as well, the engine runs the native storage service From 67a07b86327f4e6ac2cf2df9d9255125c4060dfc Mon Sep 17 00:00:00 2001 From: Alan Liu Date: Thu, 24 Sep 2026 09:45:04 -0400 Subject: [PATCH 5/5] Refuse none at every HF entry point, and finish the rename over #146's additions Under storage_backend="none" the adapter refusal was tested only through the private helper. A parametrized test now drives all five entry points -- the HF and base attach_model, attach_config, generate_with_monitoring and generate_greedy_with_monitoring -- and checks each raises "capture is off" before the model runs. It passed first (the behaviour was right, the coverage missing); with "none" taken out of the helper's refusal, all five fail. #146's later commits added a test and a v1 sentence that still spelled the persistent backend "capture"; both now say "persistent". Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y --- docs/integration-api-v1.md | 2 +- tests/test_hf_capture_refusal.py | 35 ++++++++++++++++++++++++++++++-- 2 files changed, 34 insertions(+), 3 deletions(-) diff --git a/docs/integration-api-v1.md b/docs/integration-api-v1.md index 5384982ca..e5e1ec0a9 100644 --- a/docs/integration-api-v1.md +++ b/docs/integration-api-v1.md @@ -589,7 +589,7 @@ This holds while capture is disabled too, so disabling capture never turns a record-mode step into `SKIPPED`: `set_capture_enabled(False)` leaves the installed `HookPoint`s armed, and on a record ring they would fail inside the model forward with a native error that names neither the adaptor nor the cause. -Under `storage_backend="capture"` it raises `ConfigurationError`, also before +Under `storage_backend="persistent"` it raises `ConfigurationError`, also before reserving anything, for the reason `attach_model()` does. The check is repeated here because a step need not come through the base `attach_model()`: an adaptor may override it without calling `super()`, or arm its hooks with diff --git a/tests/test_hf_capture_refusal.py b/tests/test_hf_capture_refusal.py index 41190b55f..6f08a05e3 100644 --- a/tests/test_hf_capture_refusal.py +++ b/tests/test_hf_capture_refusal.py @@ -264,6 +264,37 @@ def test_generate_greedy_with_monitoring_refuses_capture_storage(): _assert_untouched(model, engine) +@pytest.mark.parametrize("entry", [ + "hf_attach_model", "base_attach_model", "attach_config", + "generate_with_monitoring", "generate_greedy_with_monitoring"]) +def test_every_hf_entry_point_refuses_when_capture_is_off(entry): + """storage_backend="none" turns capture off, so there is no ring for the + hooks, and each entry point says so before touching the model.""" + engine = _SpyEngine("none") + model = _TinyHookedLM(engine) + input_ids, attention_mask = _inputs() + + with pytest.raises(ConfigurationError, match="capture is off"): + if entry == "hf_attach_model": + HuggingFaceAdapter(engine, "tiny").attach_model(model) + elif entry == "base_attach_model": + _StubAdapter(engine, "tiny").attach_model(model) + elif entry == "attach_config": + attach_config(_StubAdapter(engine, "tiny"), model, DMIConfig( + observations=ObservationConfig(hooks=["resid_pre"]))) + elif entry == "generate_with_monitoring": + generate_with_monitoring(model, input_ids, + attention_mask=attention_mask) + else: + generate_greedy_with_monitoring( + model, input_ids, attention_mask, + max_new_tokens=2, monitoring=True) + + assert model.generate_calls == 0 + assert model.forward_calls == 0 + _assert_untouched(model, engine) + + @pytest.mark.parametrize("backend", ["auto", "in-memory"]) def test_other_storage_backends_still_attach(backend): """Only 'persistent' and 'none' are refused; the rest attach as before.""" @@ -416,8 +447,8 @@ def test_commit_step_refuses_capture_storage_when_attach_skipped_super( via_attach_config): """The reviewer's reproduction: with the base attach bypassed, the step reserved and published into a ring with no host, and commit_step - returned RESERVED under storage_backend="capture".""" - engine = _SpyEngine("capture") + returned RESERVED under storage_backend="persistent".""" + engine = _SpyEngine("persistent") model = _TinyHookedLM(engine) adapter = _NoSuperAttachAdapter(engine, "tiny") if via_attach_config: