Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
23 changes: 12 additions & 11 deletions docs/capture-storage-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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`
Expand Down
44 changes: 33 additions & 11 deletions docs/integration-api-v1.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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`.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
51 changes: 32 additions & 19 deletions src/dmi/adapters/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."
)


Expand Down Expand Up @@ -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()",
Expand Down
2 changes: 1 addition & 1 deletion src/dmi/adapters/huggingface/adapter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
10 changes: 6 additions & 4 deletions src/dmi/adapters/huggingface/generation.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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
Expand Down
Loading
Loading