Skip to content

Native capture storage service: spool to object store to catalog, in-process - #143

Merged
zaoxing merged 14 commits into
mainfrom
feat/native-capture-service
Sep 24, 2026
Merged

zaoxing merged 14 commits into
mainfrom
feat/native-capture-service

Conversation

@zaoxing

@zaoxing zaoxing commented Sep 23, 2026 •

Copy link
Copy Markdown
Collaborator

Based on main. #135 landed as b929b65; this PR's five commits were re-parented onto it with no content change.

What this adds

Every stage of the native capture storage path after the spool already existed in C++: SpoolUploader, S3Client, CatalogWriter, NativeIndexer, CatalogSchema and the reader. Until now, only the conformance_* driver binaries ran them. This PR composes them for a real process:

  • CaptureStorageService (native/csrc/catalog/storage_service.{h,cpp}) runs one background thread that:
    • uploads what the sink staged;
    • indexes it into the ClickHouse catalog;
    • keeps the publisher lease alive;
    • reconciles the bucket against the catalog at start, which covers a crash between upload and index.
  • _dmi_native_store (bindings_store.cpp) exposes the service and the native reader to Python. It links nothing from torch and registers no ring types.
  • dmi.storage.native_capture adds NativeCaptureStorageConfig, NativeCaptureStorage and NativeCaptureReader (search / select / read, returning descriptors plus verified payloads, .tensor()).
  • Engine wiring. With MonitoringConfig(storage_backend="capture", capture_sink_config=..., capture_storage_config=...):
    • create_record_runtime starts the service before the sink opens the spool, because its start sweeps stale .open files.
    • flush_and_wait returns only once every record captured before it is queryable in the catalog.
    • close() seals the sink, drains the service and releases the lease.

Three behaviours the pieces alone did not give (each pinned by a test)

  • An index failure after a verified upload keeps the pack owed. The uploader has already removed it from the spool, so without this an empty spool read as "drained" with packs missing from the catalog. The live test cuts ClickHouse with a TCP switch mid-run.
  • The reconciler validates pack ids as canonical UUIDs before querying. One foreign object under v1/ otherwise failed the whole committed-ids query, and every reconcile pass with it.
  • Listing the spool re-hashes every pending pack, so a cycle lists once, and the loop backs off while cycles fail. Against a dead endpoint over 3 s: 146 cycles without the backoff, at most 15 with it.

Also fixed here: close() sealing the sink first (67319ac)

Stopping the ring releases the sink without flushing it, so the last open pack missed the catalog on close(). A GPU test now calls close() with no flush_and_wait. Before the fix it found 0 of 3 captures; after it, all 3 read back byte-identical.

What this does NOT do yet (tracked, not hidden)

  • No adapter drives this path. HF, vLLM and Megatron never create a capture record runtime, so the path is reachable only by a hand-built record runtime. The docs now say so. Under the capture config, HF silently stores nothing today; a follow-up makes that fail loudly, and the HF/vLLM bridges come after it.
  • One capture process per catalog. The service holds the single publisher lease, so multi-rank needs the planned upload-only role with one publisher.
  • Hardening from the audit, as follow-ups:
    • catalog ClickHouse client auth and TLS (request timeouts landed in this PR);
    • recovery from transient ClickHouse errors (today one outage of about 10 s or more latches "lease lost");
    • the sink overload policy, which today can raise inside the forward;
    • a spool owner lock;
    • the S3 CA and multipart part-size options.

Review checklist

  • libcurl now loads inside the Python process. This reopens the premise of Dynamic-dim shape inference, upload conflict reporting, and libcurl lifetime #135's linkage discussion (Copilot 3991691927, answered in 4049573867: "the catalog client is not part of the normal extension build"). That is now false for _dmi_native_store, which links clickhouse_client.cpp, s3_client.cpp and curl_init.cpp. Please check:
    • EnsureCurlGlobalInit is once-only and never cleaned up. First init is safe against other threads because of std::call_once, and no curl_global_cleanup runs anywhere; tests/test_native_curl_global_lifetime.py guards that.
    • System libcurl 7.68 / OpenSSL 1.1.1f coexists with torch's and Python ssl's OpenSSL in one process.
    • vLLM's spawned multiprocessing workers (not exercised yet).
  • A hazard for [Core] Add recurring D2H window scheduling #131 (draft, another author). [Core] Add recurring D2H window scheduling #131's transport is not None guards merge textually into this PR's new _attach_record_runtime, where transport is no longer bound. Whichever lands second must move those guards into create_record_runtime, where transport still is.
  • Deployment rule for reconcile. Object keys name no catalog, so the reconcile prefix (by default the whole bucket) must belong to one catalog. The periodic pass is off by default, and reconcile_on_start=False supports a shared bucket.

Evidence

  • CPU tier: 2213 passed. CPU-only native goals: clean rebuild, 0 warnings.
  • Wiring suite: 16 tests, with native modules faked. Four engine mutations each fail it: no service flush, no stop on a failed attach, no stop in close, sweep always.
  • Live suite (tests/test_native_capture_storage_live.py, fake S3 plus real ClickHouse): 6 passed.
    • It covers the round trip, parity with the Python reference reader, index-failure retry, reconcile after a crash plus a foreign object, one publisher per catalog, and backoff.
    • The CI live job now builds _dmi_native_store for it.
  • GPU (RTX 4090; CI has no GPU runner): tests/test_native_capture_storage_gpu_e2e.py (2 tests) plus tests/test_native_sink_ring_e2e.py, 8 passed. That run used a local C++20 build of the backend, since Build the torch-including native targets as C++20 #138 isn't on main yet.
  • Real model: Qwen2.5-0.5B through config alone, against a local Garage bucket plus ClickHouse.
    • 12,288 records in 3 packs; flush_and_wait covering sink, upload and index took 0.55 s.
    • 12/12 reference tensors byte-exact through NativeCaptureReader; generation unchanged by capture.
  • Known unrelated local failures: make check's package step fails because this host's uv-managed Python can't create a venv. One ClickHouse live test needs a CREATE USER grant this local server lacks.

Since the independent review (2026-09-24)

An independent review reproduced three majors against a live catalog; all are fixed here, each with a live test that failed first:

  • One pack that cannot be indexed no longer stops all indexing (0ff295a). A pack too big for the indexer's batch budget used to park every queued pack behind it. It is now set aside at once (a pack the indexer refuses on its own, after 5 tries), left in the object store, counted in snapshot()["rejected_packs"], and the next flush raises once, naming it. Later packs index normally.
  • flush honours its deadline when ClickHouse stops answering (677ff88, 0ff295a). The catalog client now has connect and request timeouts (10 s / 60 s by default, configurable), and flush takes the cycle lock with a deadline. Against a stalled ClickHouse, flush(1.0) returns within 3 s; before, it was still blocked at 40 s.
  • The lease holds through an object-store outage (0ff295a). A dedicated thread renews it every TTL/3, independent of the cycle backoff. With S3 dead, a rival retrying for 12 s is always refused; before, it took the catalog at about 106 s.

And the minors (fbf463e..d14ca30):

  • enable_ring_transport() over a record ring now seals the sink and retires the service, as close() does.
  • Nothing new is uploaded while an uploaded pack is still owed to the catalog, so a catalog outage keeps packs in the durable spool.
  • A failed reconcile HEAD is an error (reconcile_head_errors), not a foreign object.
  • The Python config refuses values the native side cannot honour: sub-millisecond polls, non-finite intervals, scheme-less endpoints, https with the insecure flag, negative read limits.
  • A garbage-collected service is destroyed without holding the GIL.
  • The batch-split path is pinned by a live test.

Current evidence: live suite (tests/test_native_capture_storage_live.py) 13 passed; wiring suite 37 passed; CPU tier 2234 passed; GPU capture storage and sink ring suites 8 passed (RTX 4090). CI green at d14ca30.

Copilot AI lite review requested due to automatic review settings September 23, 2026 18:02

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot was unable to review this pull request because the user who requested the review has reached their quota limit.

…error kind

Two prerequisites for running the uploader and indexer in a production
process rather than only in the single-threaded conformance drivers.

S3Client::last_attempts_ was a plain int written by every request.
SpoolUploader shares one client across its worker threads, so concurrent
uploads raced on it -- undefined behaviour even though only the store driver
ever reads the value. It is now std::atomic<int>, relaxed: the counter has
no ordering role, it only must not tear.

NativeIndexer refused an over-budget batch as CatalogError::kValue, the same
kind as every other value error. A caller that can recover by splitting the
batch -- the storage service's reconciler -- would have had to match message
text. The refusal now has its own kind, kBatchTooLarge. The conformance
driver still reports it as ValueError, which is what the Python oracle
raises, so the parity suites are unaffected.

Native store and catalog suites (s3 client, uploader, catalog lease, reader
parity): 138 passed; drivers build with 0 warnings.
Every piece after the spool already existed in C++ -- SpoolUploader,
S3Client, CatalogWriter, NativeIndexer, CatalogSchema, the reader -- but the
conformance drivers were the only thing that ran them. CaptureStorageService
composes them for a production process: one background thread uploads what
the sink staged, indexes it, keeps the publisher lease alive, and reconciles
the bucket against the catalog at start. _dmi_native_store exposes it and the
reader to Python; dmi.storage.native_capture wraps both.

Three behaviours the pieces alone do not give:

- The uploader removes a pack from the spool once its upload verifies,
  before anything indexes it. A pack whose index then fails is kept on a
  retry list, and flush() reports drained only when that list is empty --
  otherwise an empty spool read as success with packs missing from the
  catalog. The live test cuts ClickHouse with a TCP switch to pin this; with
  the retry list's term removed from "drained" it fails.
- The reconciler validates each key's pack id as a canonical UUID before
  asking the catalog about it. One foreign object under v1/ otherwise fails
  the whole committed-ids query (pack_id is a UUID column), and with it
  every reconcile pass. Found by the live test.
- Listing the spool re-hashes every pending pack, so a cycle lists it once,
  not twice, and the background loop doubles its wait while cycles fail.
  Against a dead endpoint: 146 cycles in 3 s without the backoff, <= 15 with.

Object keys name no catalog, so the reconciler's prefix must belong to one
catalog; the periodic pass is off by default and the start pass can be
disabled for a shared bucket. The reader gains search_item_columns() and a
public resolve(), so a caller can name an item's fields and pair each
hydrated payload with its descriptor at the selection's watermark.

The CI live job builds the module beside the four drivers.

Live suite (fake S3 + ClickHouse): 6 passed. ClickHouse live glob: 207
passed, 1 failed for a missing CREATE USER grant on this server (a Python
schema test this change does not touch). cpu-goals rebuild: 0 warnings.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
…config

MonitoringConfig gains capture_storage_config. Set beside
capture_sink_config, it makes the capture backend complete in-process: the
engine starts CaptureStorageService in create_record_runtime, flush_and_wait
returns only once every record captured before it is queryable in the
ClickHouse catalog, and close drains and releases the publisher lease. No
record_sink and no Python upload or index: config selects the whole path.

Ordering is the contract. The service starts BEFORE the default sink opens
the spool, because its start sweeps a crashed sink's stale .open files,
which is safe only while nothing writes there; an explicit record_sink may
already hold the spool open, so it is left unswept. A failed attach stops
the service it started, so the lease is not left held. flush_and_wait gives
the service what remains of the caller's deadline after the sink's flush.
close is best effort and logs; flush_and_wait is the boundary that reports.

Config without a sink config is refused, and the spool is not configured
twice: the service drains capture_sink_config.spool_root.

Wiring suite (cpu, native modules faked): 15 passed, and each of four
mutations to the engine (no service flush, no stop on a failed attach, no
stop in close, sweep always) fails it. GPU: ring -> sink -> service ->
catalog -> NativeCaptureReader, byte-identical for f32, bf16, i64 and a
CPU-direct f16 record, 7 passed with the sink ring suite. Qwen2.5-0.5B on
Garage + ClickHouse through config alone: 12,288 records in 3 packs,
flush_and_wait 0.55 s, 12/12 reference tensors byte-exact on read-back.
CPU tier: 2213 passed.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
close() stopped the ring and then drained the storage service. Stopping the
ring releases the sink without flushing it, so the records of the sink's
open pack were still in memory while the service drained a spool that did
not hold them yet: it reported drained, stopped, and the tail missed the
catalog. It reached the spool later -- on the linger, the sink's destructor
or interpreter exit -- and was lost outright on a hard exit before that.

In record mode with the storage service running, close() now runs a bounded
flush_records_and_wait before stopping the ring. close_flush_timeout_s is one
budget for the whole drain: the service gets what the sink flush left. A
failed sink flush is logged and the ring and service still stop.

GPU (RTX 4090): a new test captures 3 records against max_pack_records=2 and
a 60 s linger, then calls close() with no flush_and_wait. Before the fix,
none were queryable ("selection requires at least one capture"); after it,
all 3 read back byte-identical. The wiring suite pins the order: sink flush,
ring stop, service flush, service stop. Wiring and engine API suites: 52
passed; GPU capture storage and sink ring suites: 8 passed.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
capture-storage-design.md said the capture host runs no ClickHouse client.
Since capture_storage_config it can: the in-process storage service indexes
from the capture process, off the hook path. The design now says where the
indexer runs is a deployment choice, and that the single-publisher catalog
limits the in-process shape to one capture process per catalog until the
planned upload-only role exists.

integration-api-v1.md now says plainly that no adapter creates a capture
record runtime yet, so the path is reached only by hand-built record
runtimes, as the GPU end-to-end test does.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
The catalog client set no curl timeouts, so a server that accepted the
connection and never answered held the calling thread indefinitely -- a
lease renewal, a publish, and through them flush_and_wait and close. It
now takes ClickHouseTimeouts (connect 10 s, whole request 60 s by default)
and sets CURLOPT_NOSIGNAL so the timeouts never use SIGALRM in a
multi-threaded process. The defaults bound the conformance drivers too; the
storage service and its reader pass configured values.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
… S3 outage

Three defects an independent review reproduced against a live catalog:

- One pack too big for the indexer's batch budget stopped all indexing. On
  a single-pack kBatchTooLarge the give-up path parked every queued pack
  with it, and the poison pack went first on every cycle: a later small pack
  was never indexed (0 indexed, 49 splits). Such a pack is now set aside at
  once; a pack the indexer refuses on its own is set aside after
  max_index_attempts (5). A batch that threw -- an outage -- still counts
  against no pack. Set-aside packs stay in the object store, leave the
  flush boundary, show in snapshot()["rejected_packs"], and the next flush
  raises once, naming them, instead of timing out forever.
- flush(timeout) took the cycle lock with no deadline, and the background
  loop held it for a whole cycle, so a catalog that stopped answering held
  flush past its deadline (still blocked after 20 s). The cycle lock is now
  a timed mutex taken until the deadline; with the client timeouts from the
  previous commit, flush overruns by at most its own last cycle.
- The lease was renewed only inside a cycle, and failed cycles back off to
  max_backoff_ns (30 s) -- the lease TTL. With S3 down and ClickHouse
  healthy, a rival took the lease at about 106 s. A lease thread now renews
  every TTL/3 on its own schedule; a lease mutex serialises it against
  publishes and the reconciler's catalog query.

The Python config gains clickhouse_connect_timeout_s and
clickhouse_request_timeout_s.

Tests, each red first against the previous build: a 40-record pack under a
4000-byte budget is set aside and a later pack still indexes (was DID NOT
RAISE); flush(1.0) against a stalled ClickHouse returns within 3 s (was
still blocked at 40 s); a rival retrying start() for 12 s against TTL 3 s
with S3 dead is always refused (was "a second publisher took the lease").
Live ClickHouse suite: 211 passed, 1 failed (the known CREATE USER grant
this local server lacks). CPU tier: 2213 passed. GPU capture storage and
sink ring suites: 8 passed.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
enable_ring_transport() over a record ring stopped the ring and cleared
_record_mode, but left self._capture_storage running and holding the
catalog lease, and never sealed the sink: its open pack's records were
released unflushed. The next create_record_runtime then built a second
service, whose start() the process's own lease refuses.

The teardown now does what close() does: flush the sink within
close_flush_timeout_s, stop the ring, then drain and stop the service.
close() and enable_ring_transport() share the two steps as
_seal_capture_sink and _retire_capture_storage.

Test: test_replacing_a_record_ring_drains_and_stops_the_service was red
first (events were [ring stop, ring create]: no sink flush, no service
flush or stop), green after. Wiring suite: 17 passed.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
NativeCaptureStorageConfig accepted values that fail later, or not at all:

- poll_interval_s below a nanosecond became a zero native wait, a busy
  spin; it must now be at least 0.001.
- inf and NaN passed every check (NaN fails no comparison) for all five
  interval and timeout fields; each must now be finite.
- A scheme-less s3_endpoint was accepted, then read by S3Client as plain
  HTTP; an https:// endpoint with s3_allow_insecure_http=True was
  accepted, then refused by S3Client (s3_client.cpp, host_ cleared) at
  every request. Both are now refused at construction: the endpoint must
  start with http:// or https://, and the insecure flag only pairs with
  http://.

NativeCaptureReader.read now refuses a negative byte_limit and a
non-positive request_limit before resolve's catalog query. The binding
checks them too and passes byte_limit through as int64 instead of
casting it to uint64. The reported "negative limit disables the limit"
did not reproduce: the cast round-trips back to int64 and the hydration
core refused it, as a RuntimeError, after the resolve round trip.

Tests: 20 new wiring cases, 19 red first (16 accepted the bad value; the
three read cases reached resolve), green after; a zero byte_limit still
reads. The live round trip now calls the binding's hydrate with -1 and
request_limit 0 and expects ValueError (a probe of the old build raised
RuntimeError from the core).
Wiring suite: 37 passed.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
The uploader removes a pack from the spool once its upload is verified,
so a pack the catalog could not index lives only in the service's
in-memory retry list. With the catalog down, every cycle kept uploading,
moving each new pack off the durable spool onto that list. If the
process then exited, only a reconcile could recover them, and
reconcile_on_start=False -- the documented shared-bucket setting -- never
runs one.

run_cycle now retries the owed index first and uploads only once
nothing is owed, so during an outage new packs stay in the spool and the
list holds at most one cycle's uploads. Drained is unchanged: owed packs
still fail the cycle. A cycle that already failed now skips the spool
re-listing that could only confirm it is not drained, so a backlog is not
re-hashed on every cycle of an outage. snapshot()["pending_index"] is
updated on every cycle, including one that threw.

Test: test_no_new_upload_while_an_uploaded_pack_is_owed cuts ClickHouse,
stages one pack, waits for it to be owed, stages a second and flushes.
Red first ("2 uploaded but unindexed": the second pack left the spool
too), green after: it stays as the one *.dmi-pack.ready file, and after
restoring ClickHouse both packs index and read back. Live suite: 10
passed.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
The reconciler validates each uncommitted pack's HEAD metadata and skips
what is not a DMI pack. A HEAD that failed -- a 5xx, a 403, a timeout --
returns found=false, so it took the same branch: counted in
reconcile_skipped_objects, with no error recorded. A pack the pass could
not read looked like one it had rightly ignored.

A HEAD error is now recorded (last_error names the key and the cause)
and counted in a new snapshot counter, reconcile_head_errors, exposed in
the binding; it is not counted as skipped. The object is left for the
next pass, as before: periodic when reconcile_interval_s is set,
otherwise the next start.

Test: test_a_failed_head_is_an_error_not_a_foreign_object puts an object
under the fake S3's fault/forbidden/ prefix (HEAD answers 403, the list
is not faulted) and reconciles it at start. Red first
(reconcile_skipped_objects was 1), green after: skipped 0, head errors 1,
last_error names the key. Live suite: 11 passed.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
~CaptureStorageService calls stop(), which joins the cycle and lease
threads. start, flush and stop release the GIL, but destruction did not:
a service garbage-collected without stop() while a cycle waited on the
catalog joined it with the GIL held, and every other Python thread froze
until the cycle's request timed out and the lease release returned.

The binding's holder now deletes the service through a deleter that
releases the GIL around the delete (when it holds it). The service's own
threads never take the GIL, so the join cannot deadlock either way; it
just no longer stops the interpreter.

Test: test_dropping_a_running_service_does_not_hold_the_gil stalls
ClickHouse behind the TCP switch, lets a cycle block on it, then drops the
un-stopped service (del plus gc.collect()) while a thread ticks every
10 ms. Red first: 1 tick in the 3.56 s the destruction took. Green
after. Live suite: 12 passed.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
index_bounded halves a batch the indexer refuses as too large, until the
halves fit, and sets aside only a single pack that cannot fit alone. The
set-aside path had a live test; the split path, which every normal
backlog larger than the budget takes, had none.

test_a_batch_over_the_budget_splits_until_every_pack_indexes stages ten
two-record packs (~720 estimated bytes each) under a 2000-byte budget,
so the ten-pack batch must split while every pack fits alone. All ten
index and read back, with batch_splits > 0 and rejected_packs == 0. It
passes against the current service; there was no behaviour to change.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
integration-api-v1.md said the service runs from create_record_runtime
until close; enable_ring_transport over a record ring now also seals the
sink, drains the service and stops it. capture-storage-design.md now says
that while an uploaded pack is owed to the catalog nothing new is
uploaded, so a catalog outage leaves later packs in the spool.

Claude-Session: https://claude.ai/code/session_01PcY9QS1FAehHTzkjdpvN6Y
@zaoxing
zaoxing force-pushed the feat/native-capture-service branch from d14ca30 to b3dfb0d Compare September 24, 2026 14:46
@zaoxing
zaoxing merged commit 20ada8c into main Sep 24, 2026
2 checks passed
@zaoxing
zaoxing deleted the feat/native-capture-service branch September 24, 2026 15:37
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants