diff --git a/native/csrc/catalog/bindings_store.cpp b/native/csrc/catalog/bindings_store.cpp index 11d8ebdec..90be3dc99 100644 --- a/native/csrc/catalog/bindings_store.cpp +++ b/native/csrc/catalog/bindings_store.cpp @@ -79,6 +79,8 @@ dc::StorageServiceConfig service_config(const py::dict& d) { c.writer.publish_timeout_ns = get(d, "publish_timeout_ns", c.writer.publish_timeout_ns); c.writer.clock_skew_ns = get(d, "clock_skew_ns", c.writer.clock_skew_ns); + c.start_lease_wait_ns = + get(d, "start_lease_wait_ns", c.start_lease_wait_ns); c.indexer.max_packs = get(d, "indexer_max_packs", c.indexer.max_packs); c.indexer.max_estimated_bytes = get(d, "indexer_max_estimated_bytes", c.indexer.max_estimated_bytes); @@ -112,6 +114,12 @@ py::dict snapshot_dict(const dc::StorageServiceSnapshot& s) { out["swept_on_start"] = s.swept_on_start; out["pending_index"] = s.pending_index; out["rejected_packs"] = s.rejected_packs; + out["failed"] = s.failed; + out["lease_state"] = s.lease_state; + // Seconds on the monotonic clock, comparable with time.monotonic() (both + // CLOCK_MONOTONIC on Linux); 0.0 when not quarantined. + out["quarantined_until"] = static_cast(s.quarantined_until_ns) / 1e9; + out["lease_reacquisitions"] = s.lease_reacquisitions; out["last_error"] = s.last_error; return out; } diff --git a/native/csrc/catalog/storage_service.cpp b/native/csrc/catalog/storage_service.cpp index efc1f2a99..63be51074 100644 --- a/native/csrc/catalog/storage_service.cpp +++ b/native/csrc/catalog/storage_service.cpp @@ -3,8 +3,10 @@ #include #include #include +#include #include #include +#include #include #include "catalog/schema.h" @@ -52,6 +54,18 @@ bool is_pack_id(const std::string& value) { return dmi_pack::ParseUuid(value, &bytes, &canonical) && canonical == value; } +bool is_lease_refusal(const CatalogError& exc) { + return exc.kind() == CatalogError::Kind::kHeld || + exc.kind() == CatalogError::Kind::kLease; +} + +// The lease thread's tick, which is also the retry interval for a claim +// another holder refused: a sixth of the TTL, so a renewal due at ttl/3 is +// never more than a tick late. +uint64_t lease_tick_ns(uint64_t ttl_ns) { + return std::max(ttl_ns / 6, 10'000'000ull); +} + } // namespace CaptureStorageService::CaptureStorageService(StorageServiceConfig config) @@ -91,10 +105,28 @@ void CaptureStorageService::start() { std::lock_guard cycle(cycle_mutex_); if (started_) throw std::logic_error("storage service: already started"); + CatalogSchema(clickhouse_, config_.writer.database, config_.writer.table_prefix) + .ensure(&writer_.leases(), config_.schema_retry_sleep_ns); + { + std::lock_guard lease(lease_mutex_); + acquire_lease_at_start(); // throws kHeld if another publisher keeps it + } + + // After the lease, never before, so a second process pointed at this spool + // usually learns that the catalog is held before it can delete a live + // sink's .open files. Only usually: a holder that is quarantined has let + // its row lapse, and a second process can take the lease in that gap. The + // spool itself is not locked; one process per spool is the caller's job + // until the spool gets an owner lock. if (config_.sweep_spool_on_start) { std::vector recovered; std::string error; if (spool_.Recover(&recovered, &error) != dmi_store::SpoolStatus::kOk) { + try { + std::lock_guard lease(lease_mutex_); + if (writer_.held_lease() != nullptr) writer_.release_lease(); + } catch (...) { + } throw std::runtime_error("storage service: spool recovery failed: " + error); } @@ -102,14 +134,6 @@ void CaptureStorageService::start() { state_.swept_on_start = recovered.size(); } - CatalogSchema(clickhouse_, config_.writer.database, config_.writer.table_prefix) - .ensure(&writer_.leases(), config_.schema_retry_sleep_ns); - { - std::lock_guard lease(lease_mutex_); - writer_.acquire_lease(config_.holder); // throws kHeld if another publisher holds it - last_renew_ns_ = steady_ns(); - } - // A failed pass is not fatal -- the bucket is still there next time -- but a // lost lease is, and start() must not return holding a lease stop() will // never release. @@ -117,8 +141,7 @@ void CaptureStorageService::start() { try { reconcile(); } catch (const CatalogError& exc) { - if (exc.kind() == CatalogError::Kind::kHeld || - exc.kind() == CatalogError::Kind::kLease) { + if (is_lease_refusal(exc)) { try { std::lock_guard lease(lease_mutex_); if (writer_.held_lease() != nullptr) writer_.release_lease(); @@ -136,6 +159,11 @@ void CaptureStorageService::start() { { std::lock_guard lock(wake_mutex_); stop_requested_ = false; + kick_ = false; + } + { + std::lock_guard lease(lease_mutex_); + publish_lease_state(); } { std::lock_guard lock(state_mutex_); @@ -157,12 +185,22 @@ void CaptureStorageService::stop() { std::lock_guard cycle(cycle_mutex_); if (started_) { started_ = false; + // A quarantined writer holds no lease, so it writes no tombstone: the + // outcome-unknown statement may still be running, and its row must stay + // live until the TTL keeps a successor out of that window. + bool released = false; try { std::lock_guard lease(lease_mutex_); - if (writer_.held_lease() != nullptr) writer_.release_lease(); + if (writer_.held_lease() != nullptr) { + writer_.release_lease(); + released = true; + } } catch (const std::exception& exc) { record_error(std::string("lease release failed: ") + exc.what()); } + std::lock_guard lock(state_mutex_); + if (!state_.failed) state_.lease_state = released ? "released" : "none"; + state_.quarantined_until_ns = 0; } std::lock_guard lock(state_mutex_); state_.running = false; @@ -216,12 +254,13 @@ void CaptureStorageService::loop() { { std::unique_lock lock(wake_mutex_); wake_.wait_for(lock, std::chrono::nanoseconds(wait_ns), - [this] { return stop_requested_; }); + [this] { return stop_requested_ || kick_; }); if (stop_requested_) return; + kick_ = false; } { std::lock_guard lock(state_mutex_); - if (failure_) return; // the lease is gone; nothing more can publish + if (failure_) return; // another publisher holds the catalog } std::lock_guard cycle(cycle_mutex_); run_cycle(); @@ -238,9 +277,21 @@ void CaptureStorageService::loop() { CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { CycleOutcome outcome; + // The catalog phase needs the lease. Without one -- quarantined after an + // unknown outcome, or refused by another holder -- the cycle uploads + // nothing either: whatever it uploaded it could only owe, in memory. + bool catalog = false; + { + std::lock_guard lease(lease_mutex_); + catalog = ensure_publisher_lease(); + } // Indexes refs, keeping whatever does not index owed: it is already gone // from the spool, so pending_index_ is the only record of it in-process. - const auto index_or_owe = [this](std::vector refs) { + const auto index_or_owe = [this, catalog](std::vector refs) { + if (!catalog) { + pending_index_.insert(pending_index_.end(), refs.begin(), refs.end()); + return; + } std::vector unindexed; try { if (!refs.empty()) index_bounded(std::move(refs), &unindexed); @@ -257,15 +308,21 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { // of it is still owed, the catalog is down or refusing: upload // nothing new, so new packs stay in the durable spool rather than // joining a list that only this process remembers. - if (!pending_index_.empty()) { + if (catalog && !pending_index_.empty()) { std::vector owed; owed.swap(pending_index_); index_or_owe(std::move(owed)); } - // 2. Upload everything the sink has staged. + // 2. Upload everything the sink has staged -- but only with the lease + // and nothing owed. Without the lease an uploaded pack could only be + // owed, and pending_index_ dies with the process: with + // reconcile_on_start off, a crash would leave it in the bucket and + // never in the catalog. Left in the spool it survives the crash. dmi_store::UploadBatchResult batch; - if (pending_index_.empty()) batch = uploader_->UploadPending(-1); + if (catalog && pending_index_.empty()) { + batch = uploader_->UploadPending(-1); + } std::vector to_index; uint64_t uploaded_packs = 0; uint64_t uploaded_bytes = 0; @@ -297,7 +354,7 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { index_or_owe(std::move(to_index)); // 4. Reconcile on its interval. The lease thread keeps the lease alive. - if (config_.reconcile_interval_ns > 0 && + if (catalog && config_.reconcile_interval_ns > 0 && steady_ns() - last_reconcile_ns_ >= config_.reconcile_interval_ns) { reconcile(); last_reconcile_ns_ = steady_ns(); @@ -311,7 +368,10 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { // A cycle that already failed is not drained whatever the spool holds, // so it skips the listing: an empty spool lists for free, but a backlog // would be re-hashed on every cycle of an outage. - outcome.failed = upload_failures != 0 || !pending_index_.empty(); + // Without the lease nothing can be confirmed in the catalog, so the + // cycle is not drained, and it counts towards the backoff. + outcome.failed = + !catalog || upload_failures != 0 || !pending_index_.empty(); bool nothing_pending = !batch.refs.empty(); if (batch.refs.empty() && !outcome.failed) { std::vector pending; @@ -328,13 +388,12 @@ CaptureStorageService::CycleOutcome CaptureStorageService::run_cycle() { std::lock_guard lock(state_mutex_); ++state_.cycles; } catch (const CatalogError& exc) { - if (exc.kind() == CatalogError::Kind::kHeld || - exc.kind() == CatalogError::Kind::kLease) { - latch_failure(std::current_exception(), - std::string("publisher lease lost: ") + exc.what()); - } else { - record_error(exc.what()); - } + // A lease refusal here is not fatal either: the writer has dropped or + // been fenced out of its lease, and the next ensure_publisher_lease() + // decides between a fresh lease and a latch. + record_error(is_lease_refusal(exc) + ? std::string("publisher lease lost: ") + exc.what() + : std::string(exc.what())); } catch (const std::exception& exc) { record_error(exc.what()); } @@ -398,8 +457,7 @@ void CaptureStorageService::index_bounded(std::vector refs, reject(batch.front(), exc.what()); continue; } - const bool lease_lost = exc.kind() == CatalogError::Kind::kHeld || - exc.kind() == CatalogError::Kind::kLease; + const bool lease_lost = is_lease_refusal(exc); give_up(batch, std::string("index failed: ") + exc.what()); if (lease_lost) throw; return; @@ -537,10 +595,10 @@ void CaptureStorageService::keep_lease() { // The lease renews only inside a publish, and the cycle loop backs off up // to max_backoff_ns while the object store is down -- past the lease TTL. // Renewing from the cycle let the lease lapse during an outage and a rival - // take the catalog. This thread renews on its own schedule instead. - const uint64_t ttl = config_.writer.lease_ttl_ns; - const auto tick = std::chrono::nanoseconds( - std::max(ttl / 6, 10'000'000ull)); + // take the catalog. This thread renews on its own schedule instead, and + // takes a fresh lease once a lost one can be replaced. + const auto tick = + std::chrono::nanoseconds(lease_tick_ns(config_.writer.lease_ttl_ns)); while (true) { { std::unique_lock lock(wake_mutex_); @@ -551,20 +609,27 @@ void CaptureStorageService::keep_lease() { std::lock_guard lock(state_mutex_); if (failure_) return; } + std::lock_guard lease(lease_mutex_); + if (writer_.held_lease() == nullptr) { + ensure_publisher_lease(); + continue; + } try { - std::lock_guard lease(lease_mutex_); renew_lease_if_due(); } catch (const CatalogError& exc) { - if (exc.kind() == CatalogError::Kind::kHeld || - exc.kind() == CatalogError::Kind::kLease) { - latch_failure(std::current_exception(), - std::string("publisher lease lost: ") + exc.what()); - return; + if (is_lease_refusal(exc)) { + // The coordinator dropped the lease: a live foreign head refused + // the renewal claim, or the claim was contested. + lease_held_elsewhere(exc); + } else { + record_error(std::string("lease renewal failed: ") + exc.what()); } - record_error(std::string("lease renewal failed: ") + exc.what()); } catch (const std::exception& exc) { + // An unknown outcome: the writer quarantined itself and dropped the + // lease. ensure_publisher_lease() replaces it after the window. record_error(std::string("lease renewal failed: ") + exc.what()); } + publish_lease_state(); } } @@ -575,10 +640,126 @@ void CaptureStorageService::renew_lease_if_due() { if (ttl == 0 || steady_ns() - last_renew_ns_ < ttl / 3) return; writer_.renew_lease(); last_renew_ns_ = steady_ns(); + held_elsewhere_since_ns_ = 0; std::lock_guard lock(state_mutex_); ++state_.lease_renewals; } +void CaptureStorageService::acquire_lease_at_start() { + // A crashed predecessor's lease stays live for up to its TTL. Waiting it + // out here turns a restart inside that window into a short delay instead + // of a failed start; a live publisher keeps renewing, so the wait ends in + // the same refusal as before, naming the holder. + const uint64_t deadline = steady_ns() + config_.start_lease_wait_ns; + const uint64_t poll = std::clamp( + config_.writer.lease_ttl_ns / 10, 50'000'000ull, 500'000'000ull); + while (true) { + try { + writer_.acquire_lease(config_.holder); + last_renew_ns_ = steady_ns(); + return; + } catch (const CatalogError& exc) { + if (!is_lease_refusal(exc) || config_.start_lease_wait_ns == 0) throw; + const uint64_t now = steady_ns(); + if (now >= deadline) { + throw CatalogError( + exc.kind(), + "storage service: waited " + + std::to_string(config_.start_lease_wait_ns / 1'000'000) + + " ms for the publisher lease and it is still held: " + + exc.what()); + } + std::this_thread::sleep_for( + std::chrono::nanoseconds(std::min(poll, deadline - now))); + } + } +} + +bool CaptureStorageService::ensure_publisher_lease() { + if (writer_.held_lease() != nullptr) return true; + { + std::lock_guard lock(state_mutex_); + if (failure_) return false; + } + if (writer_.quarantined()) { + // The unknown-outcome statement may still be running under the lease + // it dropped; no claim, not even a fresh lease_id, until the window + // (one TTL) has passed and that lease's row has expired with it. + publish_lease_state(); + return false; + } + const uint64_t now = steady_ns(); + if (now < next_claim_ns_) return false; + try { + // A fresh lease_id: the coordinator mints one because the quarantine + // or refusal already dropped the old lease. + writer_.acquire_lease(config_.holder); + } catch (const CatalogError& exc) { + if (is_lease_refusal(exc)) { + lease_held_elsewhere(exc); + } else { + record_error(std::string("publisher lease acquisition failed: ") + + exc.what()); + } + publish_lease_state(); + return false; + } catch (const std::exception& exc) { + // Unknown outcome again: the writer quarantined itself for a TTL. + record_error(std::string("publisher lease acquisition failed: ") + + exc.what()); + publish_lease_state(); + return false; + } + last_renew_ns_ = steady_ns(); + held_elsewhere_since_ns_ = 0; + next_claim_ns_ = 0; + { + std::lock_guard lock(state_mutex_); + ++state_.lease_reacquisitions; + } + publish_lease_state(); + { + std::lock_guard lock(wake_mutex_); + kick_ = true; + } + wake_.notify_all(); + return true; +} + +void CaptureStorageService::lease_held_elsewhere(const CatalogError& refusal) { + const uint64_t now = steady_ns(); + const uint64_t ttl = config_.writer.lease_ttl_ns; + if (held_elsewhere_since_ns_ == 0) held_elsewhere_since_ns_ = now; + next_claim_ns_ = now + lease_tick_ns(ttl); + // Our own dropped row is dead within one TTL of the loss, and a handover + // (a rival that stops, releasing with a tombstone) ends sooner still. A + // refusal that has lasted 2 x TTL is a publisher that means to stay. + if (now - held_elsewhere_since_ns_ >= 2 * ttl) { + latch_failure(std::make_exception_ptr(refusal), + std::string("publisher lease held by another publisher " + "for over 2 x TTL: ") + refusal.what()); + } else { + record_error(std::string("publisher lease held elsewhere: ") + + refusal.what()); + } +} + +void CaptureStorageService::publish_lease_state() { + uint64_t until = 0; + const char* lease_state = "reacquiring"; + if (writer_.held_lease() != nullptr) { + lease_state = "held"; + } else if (writer_.quarantined(&until)) { + lease_state = "quarantined"; + } else { + until = 0; + } + std::lock_guard lock(state_mutex_); + if (state_.failed) return; + state_.lease_state = lease_state; + state_.quarantined_until_ns = until; +} + void CaptureStorageService::record_error(const std::string& message) { std::lock_guard lock(state_mutex_); state_.last_error = message.substr(0, 512); @@ -586,9 +767,23 @@ void CaptureStorageService::record_error(const std::string& message) { void CaptureStorageService::latch_failure(std::exception_ptr failure, const std::string& message) { - std::lock_guard lock(state_mutex_); - if (!failure_) failure_ = std::move(failure); - state_.last_error = message.substr(0, 512); + std::string line = message.substr(0, 512); + std::replace(line.begin(), line.end(), '\n', ' '); + { + std::lock_guard lock(state_mutex_); + if (failure_) return; + failure_ = std::move(failure); + state_.last_error = line; + state_.failed = true; + state_.running = false; + state_.lease_state = "failed"; + state_.quarantined_until_ns = 0; + } + // Once, and on one line: the loop and the lease thread both stop here, + // and nothing else will say why indexing did. + std::fprintf(stderr, "dmi capture storage: indexing stopped: %s\n", + line.c_str()); + std::fflush(stderr); } } // namespace dmi_catalog diff --git a/native/csrc/catalog/storage_service.h b/native/csrc/catalog/storage_service.h index 10372b022..965387c91 100644 --- a/native/csrc/catalog/storage_service.h +++ b/native/csrc/catalog/storage_service.h @@ -21,9 +21,23 @@ // periodically when reconcile_interval_ns is non-zero. // // Deployment shape: the service holds the catalog's single publisher lease, so -// run ONE service per (database, table_prefix). A second one fails to acquire -// the lease at start(). That is the in-process mode; a standalone daemon can -// reuse this class unchanged. +// run ONE service per (database, table_prefix). A second one waits up to +// start_lease_wait_ns for the lease at start(), then fails naming the holder. +// That is the in-process mode; a standalone daemon can reuse this class +// unchanged. +// +// The lease through ClickHouse errors. A catalog statement whose outcome is +// unknown (a transport error or timeout on a renewal or a publish) +// quarantines the writer: it drops its lease without a tombstone and refuses +// to publish for one TTL (catalog_writer.cpp). That is recoverable, not +// fatal: while the writer is quarantined or holds no lease the cycle skips +// the catalog phase, keeps pending_index_ and uploads nothing, so new packs +// stay in the durable spool rather than in a list only this process +// remembers. Once the window passes the service acquires a FRESH lease_id, +// as the Python oracle's writer documents +// (clickhouse_catalog.py, publish_snapshot). Only a foreign lease that stays +// live for 2 x TTL is fatal: the service stops (snapshot().failed), writes +// one line to stderr, and flush() rethrows the refusal naming the holder. #pragma once #include @@ -82,10 +96,16 @@ struct StorageServiceConfig { int max_index_attempts = 5; uint64_t schema_retry_sleep_ns = 500'000'000ull; + // How long start() waits for another holder's lease to expire before it + // fails with the lease held. A crashed predecessor's lease stays live for + // up to its TTL; 0 fails at once. + uint64_t start_lease_wait_ns = 0; + // Sweep a crashed sink's stale .open files before anything writes to the // spool. Recover() deletes every .open file this object does not own, so it // is only safe while no writer is live: start() must run before the sink - // opens the spool. + // opens the spool. It runs after the lease is taken, so a start refused + // the catalog never touches the spool. bool sweep_spool_on_start = true; bool reconcile_on_start = true; }; @@ -108,6 +128,15 @@ struct StorageServiceSnapshot { uint64_t swept_on_start = 0; // ready packs Recover() found at start uint64_t pending_index = 0; // uploaded packs awaiting a retried index uint64_t rejected_packs = 0; // set aside: cannot be indexed (see flush) + // A foreign lease outlived 2 x TTL: the service stopped for good. + bool failed = false; + // "none" before start, "held", "quarantined" (an unknown outcome set the + // lease aside), "reacquiring" (no lease, trying for a fresh one), + // "failed", or "released" (stop() wrote the tombstone). + std::string lease_state = "none"; + // steady_clock ns at which the quarantine ends; 0 when not quarantined. + uint64_t quarantined_until_ns = 0; + uint64_t lease_reacquisitions = 0; // fresh leases taken after a loss std::string last_error; }; @@ -118,9 +147,10 @@ class CaptureStorageService { CaptureStorageService(const CaptureStorageService&) = delete; CaptureStorageService& operator=(const CaptureStorageService&) = delete; - // Sweep the spool, ensure the catalog schema, take the publisher lease, + // Ensure the catalog schema, take the publisher lease (waiting up to + // start_lease_wait_ns for another holder's to expire), sweep the spool, // reconcile once, then start the background cycle. Throws if the lease is - // held by another publisher. + // still held by another publisher when the wait ends. void start(); // Run cycles until one finds the spool empty with every uploaded pack @@ -137,8 +167,8 @@ class CaptureStorageService { StorageServiceSnapshot snapshot() const; - // Rethrows a fatal failure latched by the background cycle: losing the - // publisher lease to another holder. + // Rethrows a fatal failure latched by the background cycle: another + // publisher holding the lease for longer than 2 x TTL. void rethrow_if_failed() const; private: @@ -156,6 +186,15 @@ class CaptureStorageService { void reconcile(); void keep_lease(); // the lease thread's body void renew_lease_if_due(); // requires lease_mutex_ + // Takes the lease at start(), waiting for an expiring predecessor. + void acquire_lease_at_start(); // requires lease_mutex_ + // Whether the writer holds a lease, taking a fresh one when it has none + // and is no longer quarantined. Never throws. Requires lease_mutex_. + bool ensure_publisher_lease(); + // Another holder refused a claim or renewal; latches once that has lasted + // 2 x TTL. Call from the catch block. Requires lease_mutex_. + void lease_held_elsewhere(const CatalogError& refusal); + void publish_lease_state(); // requires lease_mutex_ // Sets a pack aside for good; flush() reports it. Requires cycle_mutex_. void reject(const PackRefData& ref, const std::string& reason); void record_error(const std::string& message); @@ -177,6 +216,12 @@ class CaptureStorageService { // while cycles publish. Taken inside cycle_mutex_, never the other way. std::mutex lease_mutex_; uint64_t last_renew_ns_ = 0; // guarded by lease_mutex_ + // When a claim or renewal was first refused by another holder since the + // lease was last held; 0 while none has been. Guarded by lease_mutex_. + uint64_t held_elsewhere_since_ns_ = 0; + // Earliest next claim after a refusal, so flush()'s fast cycles do not + // hammer the lease table. Guarded by lease_mutex_. + uint64_t next_claim_ns_ = 0; uint64_t last_reconcile_ns_ = 0; int failure_streak_ = 0; // consecutive failed cycles, for the backoff // Uploaded, so gone from the spool, but not yet in the catalog. @@ -191,6 +236,9 @@ class CaptureStorageService { std::mutex wake_mutex_; std::condition_variable wake_; bool stop_requested_ = false; + // Set when a lease is re-acquired, so a loop in a long backoff indexes + // what is owed now rather than after its wait. Guarded by wake_mutex_. + bool kick_ = false; bool started_ = false; mutable std::mutex state_mutex_; diff --git a/src/dmi/storage/native_capture.py b/src/dmi/storage/native_capture.py index 3adde5b5c..bfd1e2e8c 100644 --- a/src/dmi/storage/native_capture.py +++ b/src/dmi/storage/native_capture.py @@ -19,7 +19,8 @@ Deployment shape: the catalog has ONE publisher lease per (``database``, ``table_prefix``), so run one capture process per catalog. A second engine -on the same catalog fails at ``create_record_runtime`` with the lease held. +on the same catalog waits ``start_lease_wait_s`` for the lease, then fails +at ``create_record_runtime`` with the lease held, naming the holder. """ from __future__ import annotations @@ -55,6 +56,22 @@ def _load_native_store_extension() -> Any: ) from exc +def _finite(name: str, value: Any) -> None: + if type(value) not in (int, float): + raise TypeError(f"{name} must be float") + if not math.isfinite(value): + raise ValueError(f"{name} must be finite") + + +def _ns(seconds: float) -> int: + return int(round(seconds * 1_000_000_000)) + + +# The native writer's MINIMUM_FENCE_MARGIN_NS: what must remain of the lease +# once the publish statement cap and the skew bound are spent. +_FENCE_MARGIN_NS = 100_000_000 + + def _positive(name: str, value: Any, kind: type) -> None: if type(value) is not kind and not (kind is float and type(value) is int): raise TypeError(f"{name} must be {kind.__name__}") @@ -114,6 +131,25 @@ class NativeCaptureStorageConfig: # open pack, then getting every staged pack into the catalog. close_flush_timeout_s: float = 60.0 + # The publisher lease. A crashed process holds the catalog for up to + # lease_ttl_s; a ClickHouse error of unknown outcome sets the lease aside + # for as long, after which the service takes a fresh one. The TTL must + # exceed publish_timeout_s + clock_skew_s by at least 0.1 s: that margin + # is what keeps a publish statement from outliving its lease. + lease_ttl_s: float = 15.0 + # The server-side cap on each fenced publish statement, whole seconds. + publish_timeout_s: float = 5.0 + # The bound on host clock disagreement across a replicated catalog. + clock_skew_s: float = 0.0 + # How long start() waits for a predecessor's lease to expire before + # failing with it held. None waits lease_ttl_s + publish_timeout_s + + # clock_skew_s, enough to outlast a crashed predecessor that ran with the + # same knobs; 0 fails at once. A predecessor with a longer TTL (the + # native default is 30 s, which processes predating these knobs used) + # can outlast it: set this explicitly for the first restart after such + # a process. + start_lease_wait_s: Optional[float] = None + def __post_init__(self) -> None: for name in ("s3_endpoint", "s3_bucket", "s3_access_key", "s3_secret_key", "store_id", "database", "table_prefix", @@ -163,6 +199,45 @@ def __post_init__(self) -> None: raise ValueError("reconcile_interval_s must be finite") if self.reconcile_interval_s < 0: raise ValueError("reconcile_interval_s must be non-negative") + self._validate_lease() + + def _validate_lease(self) -> None: + for name in ("lease_ttl_s", "publish_timeout_s", "clock_skew_s"): + _finite(name, getattr(self, name)) + if self.start_lease_wait_s is not None: + _finite("start_lease_wait_s", self.start_lease_wait_s) + if self.start_lease_wait_s < 0: + raise ValueError("start_lease_wait_s must be non-negative") + _positive("lease_ttl_s", self.lease_ttl_s, float) + _positive("publish_timeout_s", self.publish_timeout_s, float) + if self.clock_skew_s < 0: + raise ValueError("clock_skew_s must be non-negative") + if float(self.publish_timeout_s) != int(self.publish_timeout_s): + raise ValueError( + "publish_timeout_s must be a whole number of seconds: the " + "catalog sends it as max_execution_time in seconds, where a " + "fraction truncates and 0 means no limit") + # In nanoseconds, as the native writer checks it, so a pairing + # accepted here is never refused at start(). + margin = (_ns(self.lease_ttl_s) - _ns(self.publish_timeout_s) + - _ns(self.clock_skew_s)) + if margin < _FENCE_MARGIN_NS: + raise ValueError( + "lease_ttl_s must exceed publish_timeout_s + clock_skew_s by " + "at least 0.1 s, or a publish statement can still be running " + "when its lease becomes takeable") + + def _lease_native(self) -> dict[str, int]: + wait = self.start_lease_wait_s + if wait is None: + wait = (self.lease_ttl_s + self.publish_timeout_s + + self.clock_skew_s) + return { + "lease_ttl_ns": _ns(self.lease_ttl_s), + "publish_timeout_ns": _ns(self.publish_timeout_s), + "clock_skew_ns": _ns(self.clock_skew_s), + "start_lease_wait_ns": _ns(wait), + } def _native_dict(self) -> dict[str, Any]: return { @@ -210,12 +285,17 @@ def __init__( reconcile_prefix=config.reconcile_prefix, reconcile_interval_ns=int(config.reconcile_interval_s * 1e9), sweep_spool_on_start=sweep_spool, + **config._lease_native(), ) self._config = config self._service = module.StorageService(native) def start(self) -> None: - """Sweep the spool, ensure the catalog schema, take the lease.""" + """Ensure the catalog schema, take the lease, sweep the spool. + + Waits up to ``start_lease_wait_s`` for another holder's lease to + expire, then raises naming the holder. + """ self._service.start() def flush(self, timeout_s: float) -> None: diff --git a/tests/test_native_capture_storage_live.py b/tests/test_native_capture_storage_live.py index 43d2aabca..6e6f051bc 100644 --- a/tests/test_native_capture_storage_live.py +++ b/tests/test_native_capture_storage_live.py @@ -514,7 +514,10 @@ def test_a_failed_head_is_an_error_not_a_foreign_object(fake_s3, tmp_path): def test_one_publisher_per_catalog(fake_s3, tmp_path): with _catalog() as (_client, catalog): - config = _storage_config(fake_s3, catalog.table_prefix) + # No start wait: the refusal is the point here, and waiting out the + # holder is covered by the restart tests. + config = _storage_config(fake_s3, catalog.table_prefix, + start_lease_wait_s=0.0) first = _service(config, tmp_path / "first") second = _service(config, tmp_path / "second") first.start() @@ -770,6 +773,315 @@ def _native(spool, holder): assert snapshot["lease_renewals"] >= 3, snapshot +# --- the publisher lease through ClickHouse errors and restarts --------------- + + +def test_a_catalog_cut_spanning_a_renewal_and_a_publish_recovers( + fake_s3, tmp_path): + """A ClickHouse error of unknown outcome -- a renewal or a publish that + cannot reach the server -- quarantines the catalog writer and drops its + lease without a tombstone. When the quarantine ended the service found + no lease, latched "publisher lease lost", and indexing stopped for good, + although ClickHouse was back and nobody else held the catalog. It now + waits the quarantine out and takes a fresh lease.""" + spool_root = tmp_path / "spool" + switch = _Switch(CLICKHOUSE_HOST, CLICKHOUSE_HTTP_PORT) + with _catalog() as (_client, catalog): + config = _storage_config( + fake_s3, catalog.table_prefix, clickhouse_port=switch.port, + reconcile_on_start=False, lease_ttl_s=3.0, publish_timeout_s=1) + service = _service(config, spool_root) + service.start() + try: + tensors = _stage(spool_root, range(2)) + service.flush(30.0) + snapshot = service.snapshot() + assert snapshot["lease_state"] == "held", snapshot + assert snapshot["failed"] is False, snapshot + assert snapshot["lease_reacquisitions"] == 0, snapshot + assert snapshot["quarantined_until"] == 0.0, snapshot + + switch.cut() + cut_at = time.monotonic() + # Staged during the cut: uploaded, then owed, since the index + # (and its publish) cannot reach the catalog. + tensors.update(_stage(spool_root, range(2, 4))) + # The renewal due every ttl/3 fails and quarantines the lease. + _wait_for(lambda: service.snapshot()["lease_state"] == "quarantined", + timeout_s=5.0) + snapshot = service.snapshot() + assert snapshot["quarantined_until"] > time.monotonic(), snapshot + assert snapshot["running"] is True, snapshot + assert snapshot["index_failures"] >= 1, snapshot # it tried + assert snapshot["pending_index"] == 1, snapshot + with pytest.raises(TimeoutError, match="1 uploaded but unindexed"): + service.flush(0.5) + # Longer than ttl/3, so the cut spans a renewal and a publish. + time.sleep(max(0.0, cut_at + 2.5 - time.monotonic())) + switch.restore() + + # The background loop recovers on its own, without a flush. + _wait_for(lambda: service.snapshot()["indexed_packs"] == 2, + timeout_s=15.0) + snapshot = service.snapshot() + assert snapshot["failed"] is False, snapshot + assert snapshot["running"] is True, snapshot + assert snapshot["lease_state"] == "held", snapshot + assert snapshot["lease_reacquisitions"] >= 1, snapshot + assert snapshot["quarantined_until"] == 0.0, snapshot + + tensors.update(_stage(spool_root, range(4, 6))) + service.flush(30.0) + service.rethrow_if_failed() + snapshot = service.snapshot() + finally: + service.stop() + switch.close() + + assert snapshot["indexed_packs"] == 3, snapshot + assert snapshot["pending_index"] == 0, snapshot + direct = _storage_config(fake_s3, catalog.table_prefix) + captures = _read_all(direct) + assert sorted(captures) == sorted(tensors) + for capture_id, tensor in tensors.items(): + assert captures[capture_id].payload == tensor.numpy().tobytes() + + +def test_a_quarantined_service_leaves_new_packs_in_the_spool( + fake_s3, tmp_path): + """While quarantined the service cannot index, so it must not upload + either: an uploaded but unindexed pack is remembered only in memory, and + a crash then orphans it in the bucket when reconcile_on_start is off. + Before, a quarantined cycle with nothing owed uploaded the whole spool + into that in-memory list; now new packs stay in the durable spool until + the lease is back.""" + spool_root = tmp_path / "spool" + switch = _Switch(CLICKHOUSE_HOST, CLICKHOUSE_HTTP_PORT) + with _catalog() as (_client, catalog): + config = _storage_config( + fake_s3, catalog.table_prefix, clickhouse_port=switch.port, + reconcile_on_start=False, lease_ttl_s=3.0, publish_timeout_s=1) + service = _service(config, spool_root) + service.start() + try: + switch.cut() + # Nothing is owed: the renewal alone fails and quarantines. + _wait_for(lambda: service.snapshot()["lease_state"] == "quarantined", + timeout_s=5.0) + assert service.snapshot()["pending_index"] == 0 + + tensors = _stage(spool_root, range(2)) + assert len(_ready(spool_root)) == 1 + # flush() drives cycles; none of them may upload. + with pytest.raises(TimeoutError): + service.flush(0.5) + snapshot = service.snapshot() + assert snapshot["lease_state"] == "quarantined", snapshot + assert snapshot["uploaded_packs"] == 0, snapshot + assert snapshot["pending_index"] == 0, snapshot + assert len(_ready(spool_root)) == 1 + + switch.restore() + service.flush(30.0) + service.rethrow_if_failed() + snapshot = service.snapshot() + finally: + service.stop() + switch.close() + + assert snapshot["uploaded_packs"] == 1, snapshot + assert snapshot["indexed_packs"] == 1, snapshot + assert _ready(spool_root) == [] + captures = _read_all(_storage_config(fake_s3, catalog.table_prefix)) + assert sorted(captures) == sorted(tensors) + + +def _latch_lines(err: str) -> list[str]: + return [line for line in err.splitlines() if "indexing stopped" in line] + + +def test_a_rival_that_takes_over_during_a_cut_still_latches( + fake_s3, tmp_path, capfd): + """Recovery must not paper over a real takeover: when another publisher + holds the catalog for longer than two lease TTLs, the service stops, + says so once on stderr, and flush raises naming the holder.""" + switch = _Switch(CLICKHOUSE_HOST, CLICKHOUSE_HTTP_PORT) + with _catalog() as (_client, catalog): + knobs = dict(reconcile_on_start=False, lease_ttl_s=3.0, + publish_timeout_s=1) + first = _service(_storage_config( + fake_s3, catalog.table_prefix, clickhouse_port=switch.port, + holder="first-publisher", **knobs), tmp_path / "first") + rival = _service(_storage_config( + fake_s3, catalog.table_prefix, holder="rival-publisher", + start_lease_wait_s=10.0, **knobs), tmp_path / "rival") + first.start() + try: + switch.cut() + rival.start() # waits out the lease the cut first cannot renew + rival_started = time.monotonic() + switch.restore() + + _wait_for(lambda: first.snapshot()["failed"], timeout_s=20.0) + latched_at = time.monotonic() + snapshot = first.snapshot() + assert snapshot["running"] is False, snapshot + assert snapshot["lease_state"] == "failed", snapshot + assert "rival-publisher" in snapshot["last_error"], snapshot + # Only a foreign lease that outlives two TTLs latches. + assert latched_at - rival_started >= 2 * 3.0 - 0.5 + with pytest.raises(RuntimeError, match="rival-publisher"): + first.flush(1.0) + with pytest.raises(RuntimeError, match="rival-publisher"): + first.rethrow_if_failed() + + rival.flush(10.0) + rival_snapshot = rival.snapshot() + assert rival_snapshot["lease_state"] == "held", rival_snapshot + assert rival_snapshot["failed"] is False, rival_snapshot + finally: + first.stop() + rival.stop() + switch.close() + + lines = _latch_lines(capfd.readouterr().err) + assert len(lines) == 1, lines + assert "rival-publisher" in lines[0] + + +def test_a_rival_that_stops_within_two_ttls_does_not_latch( + fake_s3, tmp_path, capfd): + """A foreign lease that is gone again before two TTLs pass is a + handover, not a takeover: the service takes the lease back.""" + switch = _Switch(CLICKHOUSE_HOST, CLICKHOUSE_HTTP_PORT) + spool_root = tmp_path / "first" + with _catalog() as (_client, catalog): + knobs = dict(reconcile_on_start=False, lease_ttl_s=3.0, + publish_timeout_s=1) + config = _storage_config( + fake_s3, catalog.table_prefix, clickhouse_port=switch.port, + holder="first-publisher", **knobs) + first = _service(config, spool_root) + rival = _service(_storage_config( + fake_s3, catalog.table_prefix, holder="rival-publisher", + start_lease_wait_s=10.0, **knobs), tmp_path / "rival") + first.start() + try: + switch.cut() + rival.start() + switch.restore() + time.sleep(2.0) + rival.stop() # releases with a tombstone + + tensors = _stage(spool_root, range(2)) + first.flush(20.0) + snapshot = first.snapshot() + finally: + first.stop() + rival.stop() + switch.close() + + assert snapshot["failed"] is False, snapshot + assert snapshot["lease_state"] == "held", snapshot + assert snapshot["lease_reacquisitions"] >= 1, snapshot + assert sorted(_read_all(_storage_config( + fake_s3, catalog.table_prefix))) == sorted(tensors) + assert _latch_lines(capfd.readouterr().err) == [] + + +_HOLD_LEASE = """ +import json, sys, time +from dmi.storage.native_capture import ( + NativeCaptureStorage, NativeCaptureStorageConfig) + +service = NativeCaptureStorage( + NativeCaptureStorageConfig(**json.loads(sys.argv[1])), + spool_root=sys.argv[2], spool_max_bytes=1 << 40, sweep_spool=True) +service.start() +print("started", flush=True) +time.sleep(600) +""" + + +def test_a_restart_within_the_ttl_of_a_killed_predecessor_succeeds( + fake_s3, tmp_path): + """A SIGKILLed publisher leaves its lease live for up to a TTL, and a + restart in that window failed at once with the lease held. start() now + waits, within start_lease_wait_s, for the lease to expire.""" + import os + import signal + import sys + + with _catalog() as (_client, catalog): + knobs = dict( + s3_endpoint=fake_s3, s3_bucket=BUCKET, s3_region=REGION, + s3_access_key=ACCESS, s3_secret_key=SECRET, + s3_allow_insecure_http=True, clickhouse_host=CLICKHOUSE_HOST, + clickhouse_port=CLICKHOUSE_HTTP_PORT, database=DATABASE, + table_prefix=catalog.table_prefix, reconcile_on_start=False, + lease_ttl_s=5.0, publish_timeout_s=1) + env = dict(os.environ) + env["PYTHONPATH"] = os.pathsep.join( + [str(REPO / "src")] + ([env["PYTHONPATH"]] + if env.get("PYTHONPATH") else [])) + predecessor = subprocess.Popen( + [sys.executable, "-c", _HOLD_LEASE, + json.dumps(dict(knobs, holder="crashed-publisher")), + str(tmp_path / "predecessor")], + stdout=subprocess.PIPE, text=True, env=env) + try: + assert predecessor.stdout.readline().strip() == "started" + finally: + predecessor.send_signal(signal.SIGKILL) + predecessor.wait(timeout=30) + assert predecessor.returncode == -signal.SIGKILL + + from dmi.storage.native_capture import NativeCaptureStorageConfig + + successor = _service(NativeCaptureStorageConfig( + **knobs, holder="restarted-publisher", start_lease_wait_s=8.0), + tmp_path / "successor") + started = time.monotonic() + successor.start() + elapsed = time.monotonic() - started + try: + snapshot = successor.snapshot() + finally: + successor.stop() + + assert snapshot["running"] is True, snapshot + assert snapshot["lease_state"] == "held", snapshot + # It really waited on the dead holder's lease, and no longer than it. + assert 1.0 < elapsed < 8.0, elapsed + + +def test_a_start_refused_the_lease_leaves_the_spool_unswept(fake_s3, tmp_path): + """start() swept the spool before taking the lease, so a second process + pointed at a live process's spool deleted its in-progress .open files + and only then found the catalog held. The lease now comes first.""" + with _catalog() as (_client, catalog): + first = _service(_storage_config( + fake_s3, catalog.table_prefix, reconcile_on_start=False), + tmp_path / "first") + spool_root = tmp_path / "second" + spool_root.mkdir() + in_progress = spool_root / f".{uuid.uuid4()}.1234.open" + in_progress.write_bytes(b"a pack being written") + second = _service(_storage_config( + fake_s3, catalog.table_prefix, reconcile_on_start=False, + start_lease_wait_s=0.0), spool_root) + first.start() + try: + with pytest.raises(RuntimeError, match="held"): + second.start() + assert in_progress.exists() + finally: + first.stop() + second.start() # the lease is free now, and the sweep runs + second.stop() + assert not in_progress.exists() + + # --- https with a private CA ---------------------------------------------------- diff --git a/tests/test_native_capture_storage_wiring.py b/tests/test_native_capture_storage_wiring.py index 061afa3f7..842040a07 100644 --- a/tests/test_native_capture_storage_wiring.py +++ b/tests/test_native_capture_storage_wiring.py @@ -500,3 +500,128 @@ def test_replacing_a_record_ring_drains_and_stops_the_service( engine.create_record_runtime(_record_format()) assert len(services) == 2 assert engine._capture_storage is not None + + +# --- the publisher lease knobs ------------------------------------------------- + + +def test_lease_knobs_default_to_a_fifteen_second_lease(): + config = _storage_config() + assert config.lease_ttl_s == 15.0 + assert config.publish_timeout_s == 5.0 + assert config.clock_skew_s == 0.0 + # None waits out a predecessor for the TTL plus the publish timeout. + assert config.start_lease_wait_s is None + + +@pytest.mark.parametrize("ttl, publish, skew", [ + (5.0, 5.0, 0.0), # no margin at all + (5.09, 5.0, 0.0), # under the 0.1 s the renewed lease needs + (6.0, 5.0, 0.95), # the skew bound spends the rest + (3.0, 4.0, 0.0), # the statement cap outlives the lease +]) +def test_the_lease_must_outlast_the_publish_fence(ttl, publish, skew): + # The native writer's rule, checked here so a bad pairing fails at + # construction rather than at start(): the TTL must exceed the publish + # timeout plus the skew bound by at least 0.1 s. + with pytest.raises(ValueError, match="lease_ttl_s"): + _storage_config(lease_ttl_s=ttl, publish_timeout_s=publish, + clock_skew_s=skew) + + +@pytest.mark.parametrize("ttl, publish, skew", [ + (5.1, 5.0, 0.0), (3.0, 1.0, 0.0), (7.0, 5.0, 1.5), (2, 1, 0)]) +def test_a_lease_that_outlasts_the_fence_is_accepted(ttl, publish, skew): + config = _storage_config(lease_ttl_s=ttl, publish_timeout_s=publish, + clock_skew_s=skew) + assert config.lease_ttl_s == ttl + + +@pytest.mark.parametrize("publish", [1.5, 0.5, 4.999]) +def test_the_publish_timeout_is_whole_seconds(publish): + # The native writer sends it as max_execution_time in whole seconds; a + # fraction would truncate, and 0 is no limit at all. + with pytest.raises(ValueError, match="publish_timeout_s must be a whole"): + _storage_config(lease_ttl_s=30.0, publish_timeout_s=publish) + + +@pytest.mark.parametrize("name, value", [ + ("lease_ttl_s", 0.0), ("lease_ttl_s", -3.0), ("publish_timeout_s", 0), + ("publish_timeout_s", -1.0), ("clock_skew_s", -0.5), + ("start_lease_wait_s", -1.0)]) +def test_lease_knobs_refuse_out_of_range_values(name, value): + with pytest.raises(ValueError, match=name): + _storage_config(**{name: value}) + + +@pytest.mark.parametrize("name", [ + "lease_ttl_s", "publish_timeout_s", "clock_skew_s", "start_lease_wait_s"]) +@pytest.mark.parametrize("value", [float("inf"), float("nan")]) +def test_lease_knobs_must_be_finite(name, value): + with pytest.raises(ValueError, match=f"{name} must be finite"): + _storage_config(**{name: value}) + + +def test_a_zero_start_wait_is_accepted(): + assert _storage_config(start_lease_wait_s=0.0).start_lease_wait_s == 0.0 + + +def test_the_service_gets_the_lease_knobs_in_nanoseconds(monkeypatch, tmp_path): + engine, _events, services = _capture_engine(monkeypatch, tmp_path) + + engine.create_record_runtime(_record_format()) + + config = services[0].config + assert config["lease_ttl_ns"] == 15_000_000_000 + assert config["publish_timeout_ns"] == 5_000_000_000 + assert config["clock_skew_ns"] == 0 + assert config["start_lease_wait_ns"] == 20_000_000_000 + + +def test_the_default_start_wait_outlasts_a_skewed_predecessor(monkeypatch, + tmp_path): + # A replica whose clock lags by clock_skew_s still sees the crashed + # predecessor's row as live for that long after lease_ttl_s. + engine, _events, services = _capture_engine(monkeypatch, tmp_path) + engine._capture_storage_config = _storage_config(clock_skew_s=2.0) + + engine.create_record_runtime(_record_format()) + + assert services[0].config["start_lease_wait_ns"] == 22_000_000_000 + + +def test_explicit_lease_knobs_reach_the_service(monkeypatch, tmp_path): + engine, _events, services = _capture_engine(monkeypatch, tmp_path) + engine._capture_storage_config = _storage_config( + lease_ttl_s=3, publish_timeout_s=1.0, clock_skew_s=0.25, + start_lease_wait_s=0.5) + + engine.create_record_runtime(_record_format()) + + config = services[0].config + assert config["lease_ttl_ns"] == 3_000_000_000 + assert config["publish_timeout_ns"] == 1_000_000_000 + assert config["clock_skew_ns"] == 250_000_000 + assert config["start_lease_wait_ns"] == 500_000_000 + for name in ("lease_ttl_ns", "publish_timeout_ns", "clock_skew_ns", + "start_lease_wait_ns"): + assert type(config[name]) is int, name + + +def test_the_reader_is_not_handed_the_lease_knobs(monkeypatch): + # A reader takes no lease; the knobs are the service's alone. + from dmi.storage import native_capture + + seen = {} + + class _Reader: + def __init__(self, config): + seen.update(config) + + monkeypatch.setattr( + native_capture, "_load_native_store_extension", + lambda: SimpleNamespace(CaptureReader=_Reader, SEARCH_ITEM_COLUMNS=())) + native_capture.NativeCaptureReader(_storage_config(lease_ttl_s=3.0, + publish_timeout_s=1)) + assert not {"lease_ttl_ns", "publish_timeout_ns", "clock_skew_ns", + "start_lease_wait_ns"} & set(seen)