diff --git a/docs/capture-storage-design.md b/docs/capture-storage-design.md index 57fd96652..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)` | -| Declared by | `MonitoringConfig(storage_backend="native")` | `MonitoringConfig(storage_backend="capture")` | -| Status | production | explicitly reference-only; production sinks remain native-only | +| 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 | 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 @@ -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..e5e1ec0a9 100644 --- a/docs/integration-api-v1.md +++ b/docs/integration-api-v1.md @@ -135,15 +135,36 @@ 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 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 in-process: from `create_record_runtime` until `close`, or until @@ -156,7 +177,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 +505,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 @@ -567,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/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..9e0c6547f 100644 --- a/src/dmi/config.py +++ b/src/dmi/config.py @@ -3,10 +3,53 @@ from __future__ import annotations from dataclasses import dataclass, field +import sys +import warnings from typing import Literal, Optional, get_args -StorageBackend = Literal["auto", "native", "capture", "none"] +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") + +_DEPRECATED_STORAGE_BACKENDS = {"native": "in-memory", "capture": "persistent"} + + +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 + 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) # 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 @@ -71,29 +114,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 +156,56 @@ 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) -> CanonicalStorageBackend: + """``storage_backend`` with a deprecated name replaced by its new one.""" + return canonical_storage_backend(self.storage_backend) + def __post_init__(self) -> None: - backends = get_args(StorageBackend) - if self.storage_backend not in backends: - raise ValueError( - "storage_backend must be one of " - + ", ".join(repr(name) for name in backends) - + f"; got {self.storage_backend!r}" + canonical_storage_backend(self.storage_backend) # refuses unknowns + 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=_caller_stacklevel(), ) + 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..efd1fbc3b 100644 --- a/src/dmi/engine.py +++ b/src/dmi/engine.py @@ -119,7 +119,17 @@ 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, 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 + + # 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 = ( @@ -143,15 +153,23 @@ 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=" + 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=" ) - if self._storage_backend in ("capture", "none") and host_configured: + if self._storage_backend == "none" and ring_config is not None: raise ValueError( - f"config.storage_backend={self._storage_backend!r} does not " + "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._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 " @@ -185,6 +203,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( @@ -244,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: @@ -262,17 +293,19 @@ 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 " + 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 " "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 " + 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" ) @@ -290,6 +323,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 +361,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 +385,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 @@ -467,6 +506,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. @@ -483,6 +524,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/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..6f08a05e3 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,40 @@ def test_generate_greedy_with_monitoring_refuses_capture_storage(): _assert_untouched(model, engine) -@pytest.mark.parametrize("backend", ["auto", "native", "none"]) +@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 '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 +344,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: @@ -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: 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..39eb3e1a5 --- /dev/null +++ b/tests/test_storage_backend_choices.py @@ -0,0 +1,204 @@ +"""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") + + +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) + + +# --- 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)