Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
33 commits
Select commit Hold shift + click to select a range
648f6ce
Stage the sink's open pack when the ring releases it
zaoxing Sep 29, 2026
fa250a1
Let a cancel cut uploads short: abort the transfer, wake the backoff
zaoxing Sep 29, 2026
5f882cd
Return flush(timeout) on time, and stop() promptly while an upload st…
zaoxing Sep 29, 2026
af4f4f2
Say what close() drains, and what its budget does and does not bound
zaoxing Sep 29, 2026
7a60f35
Bound a sink flush by its own timeout while another is in flight
zaoxing Sep 29, 2026
b13c7ec
Cut index reads at stop() and past a flush; end a pass at an unanswer…
zaoxing Sep 29, 2026
0121538
Start no index batch past a flush's deadline but the first
zaoxing Sep 29, 2026
efe9279
Run no loop cycle between a flush and the stop() after it
zaoxing Sep 29, 2026
90d00b0
Say what close() and flush() now overrun their budget by
zaoxing Sep 29, 2026
15db912
Do not hash the spool in a flush that is out of time
zaoxing Sep 29, 2026
b61db70
Book an upload as cancelled only when the cancel ended it
zaoxing Sep 29, 2026
7189a15
Keep the lease test's last cycle busy with work stop() cannot cut
zaoxing Sep 29, 2026
7a6617d
Refuse, not hold, the connection a closed switch still accepts
zaoxing Sep 29, 2026
4d90b5a
Pin that a cancel wakes both upload backoffs, not just ends them
zaoxing Sep 29, 2026
f9807ce
Test that the release backstop is bounded and fails soft
zaoxing Sep 29, 2026
53c9c2b
Pin that a flush never runs the reconcile, only the loop does
zaoxing Sep 29, 2026
80df337
Test the cancel paths a stop and a restart go through
zaoxing Sep 29, 2026
58453c8
Say that close()'s stop can wait for a renewal, and drop "within"
zaoxing Sep 29, 2026
5fecd88
Stop a spool listing between packs at stop() or a flush's deadline
zaoxing Sep 29, 2026
f563341
Index a full batch first, not the remainder, past a flush's deadline
zaoxing Sep 29, 2026
7a4b312
Drop the flush-out-of-time spool check the cancelled listing replaced
zaoxing Sep 29, 2026
d677f07
Upload and index a chunk at a time, so a cut cycle owes one chunk at …
zaoxing Sep 29, 2026
9acff4d
Start no index batch once the reads are cut, and no reconcile query
zaoxing Sep 29, 2026
e49739f
Book a failed index read as cancelled only when the cancel cut it
zaoxing Sep 29, 2026
41398b8
Pin that a cut PUT or part on the last attempt books a cancel
zaoxing Sep 29, 2026
2172aea
Loosen the listing-cut test's time bound to what it has to separate
zaoxing Sep 29, 2026
301cdb2
Assert that a reconcile listing stop() cuts ends quietly
zaoxing Sep 29, 2026
8bc3fb7
Say, and pin, that a cut multipart upload's abort outlasts a request …
zaoxing Sep 29, 2026
bb774c4
Say close() drains with a budget of close_flush_timeout_s, not within it
zaoxing Sep 29, 2026
7525aac
Say in cancel.h that the service owns two Cancellations
zaoxing Sep 29, 2026
1303d85
Name the S3 retries the client actually makes, and its backoff cap
zaoxing Sep 29, 2026
119cea8
Cap the S3 backoff's shift, not only its result
zaoxing Sep 29, 2026
8044c5e
Reflow three lines the chunking and read-cancel commits left long
zaoxing Sep 29, 2026
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
16 changes: 16 additions & 0 deletions docs/integration-api-v1.md
Original file line number Diff line number Diff line change
Expand Up @@ -330,6 +330,22 @@ Shutdown exceptions are suppressed, so an integration requiring an
authoritative final read must ensure every worker reaches this close path and
should separately check native host failures.

Under `storage_backend="persistent"`, the native pack sink stages the pack it
still has open when the stopping ring releases it, so the last records reach
the spool without a flush. With `capture_storage_config` set, `close()` first
drains capture, best effort, with a budget of `close_flush_timeout_s`: it
flushes the sink, stops the ring, waits for the storage service to get the
staged packs into the catalog, and stops the service. What misses the budget is
logged and left for the next start (the spool, or the reconcile);
`flush_and_wait` is the call that raises when captures are not queryable in
time. Past the budget the drain starts no upload and at most one index batch,
so against a catalog or object store that stops answering, `close()` outlasts
the budget by up to about two `clickhouse_request_timeout_s` (the request in
flight, and the lease release), plus up to about 6 s when the budget cuts a
multipart upload, whose abort nothing cuts; `close_flush_timeout_s` documents
the full bound. The same bounds hold for `flush_and_wait(timeout_s)`, less the
lease release.

Closing does not disable or uninstall HookPoints: they retain hook IDs and the
old payload tensor. Treat the attached model as terminal too. A later CUDA
forward—especially after another engine becomes active—can combine stale hook
Expand Down
6 changes: 3 additions & 3 deletions native/Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -367,15 +367,15 @@ build/conformance_sign: csrc/store/s3_sign.cpp csrc/store/conformance_sign.cpp c
-Icsrc/store -Icsrc/common -lcrypto

# A2 store conformance driver + fault-matrix target. Needs libcurl (above).
build/conformance_store: csrc/store/s3_sign.cpp csrc/store/s3_client.cpp csrc/store/spool.cpp csrc/store/uploader.cpp csrc/store/conformance_store.cpp csrc/store/s3_sign.h csrc/store/s3_client.h csrc/store/spool.h csrc/store/uploader.h csrc/common/json.cpp csrc/common/json.h csrc/common/curl_init.cpp csrc/common/curl_init.h | check-libcurl
build/conformance_store: csrc/store/s3_sign.cpp csrc/store/s3_client.cpp csrc/store/spool.cpp csrc/store/uploader.cpp csrc/store/conformance_store.cpp csrc/store/s3_sign.h csrc/store/s3_client.h csrc/store/cancel.h csrc/store/spool.h csrc/store/uploader.h csrc/common/json.cpp csrc/common/json.h csrc/common/curl_init.cpp csrc/common/curl_init.h | check-libcurl
mkdir -p $(BUILD_DIR)
$(CXX) -std=c++17 -O2 -Wall -Wextra -o $@ \
csrc/store/s3_sign.cpp csrc/store/s3_client.cpp csrc/store/spool.cpp csrc/store/uploader.cpp csrc/store/conformance_store.cpp csrc/common/json.cpp csrc/common/curl_init.cpp \
-Icsrc/store -Icsrc/common $(CURL_CPPFLAGS) $(CURL_LDFLAGS) -lcrypto -lcurl -lpthread

# B1 catalog conformance driver: lease coordinator + version allocator over
# ClickHouse's HTTP interface. No torch/pybind; needs libcurl like the store.
build/conformance_catalog: csrc/catalog/clickhouse_client.cpp csrc/catalog/lease_coordinator.cpp csrc/catalog/version_allocator.cpp csrc/catalog/catalog_writer.cpp csrc/catalog/pack_index.cpp csrc/catalog/indexer.cpp csrc/catalog/schema.cpp csrc/catalog/reader.cpp csrc/catalog/hydration.cpp csrc/catalog/conformance_catalog.cpp csrc/common/json.cpp csrc/pack/pack_builder.cpp csrc/store/s3_sign.cpp csrc/store/s3_client.cpp csrc/catalog/clickhouse_client.h csrc/catalog/lease_coordinator.h csrc/catalog/version_allocator.h csrc/catalog/catalog_writer.h csrc/catalog/sql_escape.h csrc/catalog/pack_index.h csrc/catalog/indexer.h csrc/catalog/schema.h csrc/catalog/reader.h csrc/catalog/hydration.h csrc/common/json.h csrc/store/s3_client.h csrc/common/curl_init.cpp csrc/common/curl_init.h | check-libcurl
build/conformance_catalog: csrc/catalog/clickhouse_client.cpp csrc/catalog/lease_coordinator.cpp csrc/catalog/version_allocator.cpp csrc/catalog/catalog_writer.cpp csrc/catalog/pack_index.cpp csrc/catalog/indexer.cpp csrc/catalog/schema.cpp csrc/catalog/reader.cpp csrc/catalog/hydration.cpp csrc/catalog/conformance_catalog.cpp csrc/common/json.cpp csrc/pack/pack_builder.cpp csrc/store/s3_sign.cpp csrc/store/s3_client.cpp csrc/catalog/clickhouse_client.h csrc/catalog/lease_coordinator.h csrc/catalog/version_allocator.h csrc/catalog/catalog_writer.h csrc/catalog/sql_escape.h csrc/catalog/pack_index.h csrc/catalog/indexer.h csrc/catalog/schema.h csrc/catalog/reader.h csrc/catalog/hydration.h csrc/common/json.h csrc/store/s3_client.h csrc/store/cancel.h csrc/common/curl_init.cpp csrc/common/curl_init.h | check-libcurl
mkdir -p $(BUILD_DIR)
$(CXX) -std=c++17 -O2 -Wall -Wextra -o $@ \
csrc/catalog/clickhouse_client.cpp csrc/catalog/lease_coordinator.cpp \
Expand All @@ -388,7 +388,7 @@ build/conformance_catalog: csrc/catalog/clickhouse_client.cpp csrc/catalog/lease
$(CURL_CPPFLAGS) $(CURL_LDFLAGS) -lcrypto -lcurl -lpthread

# A3a spool conformance driver: no network, no torch; runs in any CI job.
build/conformance_spool: csrc/store/spool.cpp csrc/store/conformance_spool.cpp csrc/store/spool.h csrc/common/json.cpp csrc/common/json.h
build/conformance_spool: csrc/store/spool.cpp csrc/store/conformance_spool.cpp csrc/store/spool.h csrc/store/cancel.h csrc/common/json.cpp csrc/common/json.h
mkdir -p $(BUILD_DIR)
$(CXX) -std=c++17 -O2 -Wall -Wextra -o $@ \
csrc/store/spool.cpp csrc/store/conformance_spool.cpp csrc/common/json.cpp \
Expand Down
1 change: 1 addition & 0 deletions native/csrc/catalog/bindings_store.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -116,6 +116,7 @@ py::dict snapshot_dict(const dc::StorageServiceSnapshot& s) {
out["uploaded_packs"] = s.uploaded_packs;
out["uploaded_bytes"] = s.uploaded_bytes;
out["upload_failures"] = s.upload_failures;
out["cancelled_uploads"] = s.cancelled_uploads;
out["indexed_packs"] = s.indexed_packs;
out["indexed_rows"] = s.indexed_rows;
out["index_failures"] = s.index_failures;
Expand Down
61 changes: 61 additions & 0 deletions native/csrc/catalog/conformance_catalog.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@
#include <vector>

#include "../common/json.h"
#include "../store/cancel.h"
#include "../store/s3_client.h"
#include "catalog_writer.h"
#include "clickhouse_client.h"
Expand Down Expand Up @@ -466,6 +467,66 @@ std::string respond(const std::string& line, Session* session) {
out += "]";
return prefix + "true" + out + "}";
}
if (op == "read_pack_rows") {
// Session-less: one pack's descriptor rows, read through the object
// store as the indexer reads them, the client holding a Cancellation
// -- whose deadline passes cancel_after_ms from now, as a flush arms
// one, or which is cancelled once the first exchange has returned
// (cancel_after_exchange), as a cancel that comes in while a failure
// is reported. Pins how a failed read is told apart -- the store's
// answer, the store not answering, and whether a cancel is what cut
// it -- without a catalog.
dmi_store::S3Config s3_config;
s3_config.endpoint = jc::FindString(line, "endpoint");
s3_config.bucket = jc::FindString(line, "bucket");
s3_config.region = jc::FindString(line, "region");
s3_config.access_key = jc::FindString(line, "access");
s3_config.secret_key = jc::FindString(line, "secret");
s3_config.allow_insecure_http = jc::FindBool(line, "insecure");
if (jc::HasKey(line, "read_timeout")) {
s3_config.read_timeout_s =
static_cast<int>(field_int(line, "read_timeout"));
}
if (jc::HasKey(line, "max_attempts")) {
s3_config.max_attempts =
static_cast<int>(field_int(line, "max_attempts"));
}
dmi_store::Cancellation cancel; // outlives the client
dmi_store::S3Client s3(s3_config);
if (jc::HasKey(line, "cancel_after_ms")) {
cancel.set_deadline(
dmi_store::Cancellation::NowNs() +
static_cast<uint64_t>(field_int(line, "cancel_after_ms")) *
1'000'000ull);
s3.set_cancellation(&cancel);
}
if (jc::FindBool(line, "cancel_after_exchange")) {
s3.set_cancellation(&cancel);
s3.SetAfterExchangeHookForTesting([&cancel] { cancel.Cancel(); });
}
const std::string element = jc::FindObject(line, "ref");
dmi_catalog::PackRefData ref;
ref.pack_id = jc::FindString(element, "pack_id");
ref.store_id = jc::FindString(element, "store_id");
ref.object_key = jc::FindString(element, "object_key");
ref.object_bytes =
static_cast<uint64_t>(field_int(element, "object_bytes"));
ref.checksum = jc::FindString(element, "checksum");
ref.record_count =
static_cast<uint64_t>(field_int(element, "record_count"));
try {
out = ",\"rows\":" +
std::to_string(dmi_catalog::read_pack_descriptor_rows(&s3, ref)
.size());
} catch (const dmi_catalog::StoreUnavailableError& e) {
std::string message;
escape_into(e.what(), &message);
return prefix + "false,\"error\":\"StoreUnavailable\",\"cancelled\":" +
(e.cancelled() ? "true" : "false") + ",\"message\":" +
message + "}";
}
return prefix + "true" + out + "}";
}
if (session->writer == nullptr) {
return prefix + "false,\"what\":\"call open first\"}";
}
Expand Down
4 changes: 4 additions & 0 deletions native/csrc/catalog/indexer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -239,6 +239,10 @@ void NativeIndexer::read(IndexPlan* planned) {
try {
rows = read_pack_descriptor_rows(s3_, ref);
} catch (const CatalogError& e) {
if (config_.end_read_when_store_unavailable &&
dynamic_cast<const StoreUnavailableError*>(&e) != nullptr) {
throw;
}
std::string message = e.what();
if (message.size() > 512) message.resize(512);
result.failures.push_back(
Expand Down
11 changes: 10 additions & 1 deletion native/csrc/catalog/indexer.h
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,14 @@ struct IndexerConfig {
int max_rows_per_insert = 10'000;
uint64_t max_estimated_bytes = 128ull * 1024 * 1024;
int max_publish_attempts = 8;
// read() ends at a pack the object store did not answer for
// (StoreUnavailableError propagates) instead of failing that one pack
// and reading the next. The oracle, CatalogIndexer, fails the pack and
// moves on, and so does the default; the storage service sets it, since
// against a store that stopped answering every further pack costs the
// client's full timeouts and fails the same way, and an outage is not
// the pack's fault to count against it.
bool end_read_when_store_unavailable = false;
// Test seam, unset in production, called with each allocated version
// just before the publish that carries it. The publish wedges in
// CatalogWriter exist for the same reason: a version race needs the
Expand Down Expand Up @@ -82,7 +90,8 @@ class NativeIndexer {
// Deduplicates refs and reads the replay guard (the catalog).
IndexPlan plan(const std::vector<PackRefData>& refs);
// Reads each pending pack's descriptor rows (the object store only).
// Throws kBatchTooLarge past max_estimated_bytes.
// Throws kBatchTooLarge past max_estimated_bytes, and
// StoreUnavailableError under end_read_when_store_unavailable.
void read(IndexPlan* plan);
// Allocates a version, writes the descriptors, publishes and commits the
// inventory (the catalog). Requires read().
Expand Down
25 changes: 19 additions & 6 deletions native/csrc/catalog/pack_index.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,17 @@ constexpr uint64_t kMaxLayerNumber = 0x7FFFFFFFull; // 2^31 - 1
throw CatalogError(CatalogError::Kind::kValue, what);
}

// A range read that failed: the store's answer about the pack, or no
// answer at all (StoreUnavailableError) -- the cancel's only when the read
// says the Cancellation cut it. One set once the store had failed the read
// (a flush's deadline passing just after a 503) cut nothing: that failure
// is the store's, and is recorded as such.
[[noreturn]] void range_read_failed(const std::string& what, bool unavailable,
bool cancelled) {
if (unavailable) throw StoreUnavailableError(what, cancelled);
throw CatalogError(CatalogError::Kind::kValue, what);
}

// `CaptureMetadata.from_mapping` wraps every ValueError `__post_init__`
// raises as `invalid capture metadata: {exc}`, so the refusals that come
// from the metadata model's own bounds carry that prefix and the ones that
Expand Down Expand Up @@ -433,10 +444,12 @@ std::vector<std::string> read_pack_descriptor_rows(
std::vector<uint8_t> trailer;
const uint64_t trailer_offset = ref.object_bytes - kTrailerSize;
if (charge) charge(kTrailerSize);
bool unavailable = false;
bool cancelled = false;
if (!s3->GetRange(ref.object_key, trailer_offset, kTrailerSize, &trailer,
&error)) {
throw CatalogError(CatalogError::Kind::kValue,
"pack trailer read failed: " + error);
&error, &unavailable, &cancelled)) {
range_read_failed("pack trailer read failed: " + error, unavailable,
cancelled);
}
if (trailer.size() != kTrailerSize) format_error("pack trailer is truncated");
const auto read_u64 = [&](size_t at) {
Expand Down Expand Up @@ -482,9 +495,9 @@ std::vector<std::string> read_pack_descriptor_rows(
// about to be read -- not the object's size, which only bounds it.
if (charge) charge(footer_length);
if (!s3->GetRange(ref.object_key, footer_offset, footer_length, &footer,
&error)) {
throw CatalogError(CatalogError::Kind::kValue,
"pack footer read failed: " + error);
&error, &unavailable, &cancelled)) {
range_read_failed("pack footer read failed: " + error, unavailable,
cancelled);
}
if (dmi_pack::Crc32(footer.data(), footer.size()) != footer_crc) {
throw CatalogError(CatalogError::Kind::kValue, "footer checksum mismatch");
Expand Down
22 changes: 21 additions & 1 deletion native/csrc/catalog/pack_index.h
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,28 @@
#include <vector>

#include "../store/s3_client.h"
#include "lease_coordinator.h"

namespace dmi_catalog {

// read_pack_descriptor_rows could not get the object store's answer about
// the pack: a transport error or timeout, a retryable status on every
// attempt, or the S3 client's Cancellation cutting the read (cancelled()
// says which; a cancel that came in once the store had failed it does not
// count, as in S3Client). It says
// nothing about the pack itself, unlike every other refusal the read makes.
// A CatalogError of kind kValue like those, so a caller that treats every
// unreadable pack alike still does.
class StoreUnavailableError : public CatalogError {
public:
StoreUnavailableError(const std::string& what, bool cancelled)
: CatalogError(Kind::kValue, what), cancelled_(cancelled) {}
bool cancelled() const { return cancelled_; }

private:
bool cancelled_;
};

struct PackRefData {
std::string pack_id;
std::string store_id;
Expand All @@ -29,7 +48,8 @@ struct PackRefData {
// Reads the pack the ref names and renders one descriptor VALUES row per
// record — the 33 capture_raw columns in schema order, without
// index_version (the batch's own version). Throws CatalogError (kValue)
// on any format violation, mirroring PackFormatError / PackIntegrityError.
// on any format violation, mirroring PackFormatError / PackIntegrityError,
// and StoreUnavailableError when the store did not answer a range read.
// The bucket is the S3Client's own config; a bucket parameter here would
// only invite a caller to believe passing a different one redirects the
// read.
Expand Down
Loading
Loading