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
6 changes: 6 additions & 0 deletions docs/integration-api-v1.md
Original file line number Diff line number Diff line change
Expand Up @@ -210,6 +210,12 @@ not use a `readonly=1` profile: the reader sends its query limits
`publish_timeout_s` (10 s at the default 5 s publish timeout). A refused connection is retried for any
statement; a reset, or a 5xx that is not a permanent ClickHouse error (such
as a row limit or a denied grant), for reads only; and a timeout never.
While the storage service holds the publisher lease, every request it makes,
retries included, must also be answered by the lease deadline: `lease_ttl_s`
less `clock_skew_s` and 0.1 s after the claim that stamped the lease row was
sent. Each request of a claim made without a lease has
min(`clickhouse_request_timeout_s`, `lease_ttl_s` / 3). `clock_skew_s` must
therefore be at most `lease_ttl_s` / 2 - 0.3 s.
The schedule's default factory creates a distinct `CaptureSchedule` for each
config instance.
`MonitoringEngine` stores the config, while concrete adaptors decide whether
Expand Down
2 changes: 2 additions & 0 deletions native/csrc/catalog/bindings_store.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,8 @@ py::dict snapshot_dict(const dc::StorageServiceSnapshot& s) {
// CLOCK_MONOTONIC on Linux); 0.0 when not quarantined.
out["quarantined_until"] = static_cast<double>(s.quarantined_until_ns) / 1e9;
out["lease_reacquisitions"] = s.lease_reacquisitions;
out["lease_timeouts"] = s.lease_timeouts;
out["lease_timeout_error"] = s.lease_timeout_error;
out["last_error"] = s.last_error;
return out;
}
Expand Down
22 changes: 21 additions & 1 deletion native/csrc/catalog/catalog_writer.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -299,12 +299,20 @@ PublisherLease CatalogWriter::acquire_lease(const std::string& holder) {
require_owned_by_this_process();
const std::lock_guard<std::recursive_mutex> serial(serial_);
require_not_quarantined();
const bool renewing = leases_->lease() != nullptr;
try {
return leases_->acquire(holder);
} catch (const CatalogError&) {
throw;
} catch (const std::exception&) {
quarantine();
// An unknown outcome only if a claim row may have been written: a
// renewal's, or a claim whose INSERT may have reached the server. One
// that failed before that -- its head read timed out, or could not
// connect -- wrote nothing, and a quarantine would only keep a writer
// with nothing to wait out from trying again. Deliberately unlike the
// Python oracle (clickhouse_catalog.py, acquire_publisher_lease), which
// quarantines on any error here.
if (renewing || leases_->claim_insert_sent()) quarantine();
throw;
}
}
Expand All @@ -315,6 +323,18 @@ PublisherLease CatalogWriter::renew_lease() {
return renew_for_publish();
}

void CatalogWriter::abandon_lease() {
require_owned_by_this_process();
const std::lock_guard<std::recursive_mutex> serial(serial_);
if (leases_->lease() != nullptr) quarantine();
}

bool CatalogWriter::refused_by_own_claims() const {
require_owned_by_this_process();
const std::lock_guard<std::recursive_mutex> serial(serial_);
return leases_->refused_by_own_claims();
}

void CatalogWriter::release_lease() {
require_owned_by_this_process();
const std::lock_guard<std::recursive_mutex> serial(serial_);
Expand Down
22 changes: 22 additions & 0 deletions native/csrc/catalog/catalog_writer.h
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,28 @@ class CatalogWriter {
PublisherLease renew_lease();
void release_lease();
const PublisherLease* held_lease() const { return leases_->lease(); }
// The held lease's timing, on the steady clock, 0 while none is held: when
// the claim INSERT that stamped its row was sent -- the lease thread's
// last renewal, or a publish's -- and the deadline every request made
// under it has to be answered by (lease_coordinator.h). Lock-free like
// held_lease(): the storage service reads the deadline before each request
// it sends while holding its lease lock.
uint64_t lease_sent_ns() const {
const PublisherLease* held = leases_->lease();
return held != nullptr ? held->sent_ns : 0;
}
uint64_t lease_deadline_ns() const {
const PublisherLease* held = leases_->lease();
return held != nullptr ? held->deadline_ns : 0;
}
// Gives up a lease whose deadline has passed unrenewed. Its row may still
// be live, and a statement of the pass that ran out of time may still be
// running, so it goes as an unknown outcome does: no tombstone, and
// quarantined for a TTL.
void abandon_lease();
// Whether the last lease refusal came only from this writer's own earlier
// claim rows (LeaseCoordinator::refused_by_own_claims).
bool refused_by_own_claims() const;
uint64_t allocate_version();
uint64_t max_version(const std::string& table,
const std::string& column) const;
Expand Down
222 changes: 203 additions & 19 deletions native/csrc/catalog/clickhouse_client.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -98,6 +98,19 @@ bool transient_clickhouse_error(int code) {
}
}

// The ClickHouse errors that say the server gave up on a statement for a
// time limit the request set: TIMEOUT_EXCEEDED (159, max_execution_time;
// answered with a 408), DEADLOCK_AVOIDED (473, lock_acquire_timeout) and
// UNKNOWN_STATUS_OF_INSERT (319, insert_quorum_timeout -- which the server
// also answers when it lost its Keeper session mid-insert, an unknown
// outcome either way). The lease INSERTs set all three from the time the
// client gives them (lease_coordinator.cpp), so the server often gives up
// first: that is the request running out of time as surely as the client's
// own timeout, and execute() says so (ClickHouseError::timed_out).
bool server_time_limit(int code) {
return code == 159 || code == 473 || code == 319;
}

// libcurl takes whole milliseconds, and 0 means "its default" -- no bound at
// all for the whole request. A positive timeout therefore rounds UP, so one
// below a millisecond still bounds the request (at 1 ms); validate() has
Expand All @@ -121,6 +134,19 @@ std::string url_encode(const std::string& value) {
return out;
}

// The innermost RequestDeadline on this thread; each links to the one it
// nests in.
thread_local RequestDeadline* innermost_deadline = nullptr;
// Set while a before_request hook runs on this thread.
thread_local bool in_before_request = false;

// "60 s", "0.25 s": a timeout for a message.
std::string seconds_text(double seconds) {
char out[32];
std::snprintf(out, sizeof(out), "%g s", seconds);
return out;
}

bool has_header_breaking_byte(const std::string& value) {
return value.find_first_of(std::string("\r\n\0", 3)) != std::string::npos;
}
Expand Down Expand Up @@ -219,6 +245,63 @@ std::vector<Row> parse_tsv(const std::string& body) {

} // namespace

uint64_t steady_now_ns() {
return static_cast<uint64_t>(
std::chrono::duration_cast<std::chrono::nanoseconds>(
std::chrono::steady_clock::now().time_since_epoch())
.count());
}

RequestDeadline::RequestDeadline(uint64_t deadline_ns, std::string bound)
: fixed_ns_(deadline_ns), bound_(std::move(bound)),
outer_(innermost_deadline) {
innermost_deadline = this;
}

RequestDeadline::RequestDeadline(std::function<uint64_t()> deadline_ns,
std::string bound,
std::function<void()> before_request)
: moving_ns_(std::move(deadline_ns)), bound_(std::move(bound)),
before_request_(std::move(before_request)),
outer_(innermost_deadline) {
innermost_deadline = this;
}

RequestDeadline::~RequestDeadline() { innermost_deadline = outer_; }

uint64_t RequestDeadline::current(std::string* bound) {
uint64_t tightest = 0;
for (const RequestDeadline* scope = innermost_deadline; scope != nullptr;
scope = scope->outer_) {
const uint64_t at = scope->moving_ns_ ? scope->moving_ns_()
: scope->fixed_ns_;
if (at != 0 && (tightest == 0 || at < tightest)) {
tightest = at;
if (bound != nullptr) *bound = scope->bound_;
}
}
return tightest;
}

void RequestDeadline::note_outcome(const std::string& timeout) {
for (RequestDeadline* scope = innermost_deadline; scope != nullptr;
scope = scope->outer_) {
scope->last_timeout_ = timeout;
}
}

void RequestDeadline::before_request() {
const RequestDeadline* scope = innermost_deadline;
if (scope == nullptr || !scope->before_request_ || in_before_request) {
return;
}
struct Running {
Running() { in_before_request = true; }
~Running() { in_before_request = false; }
} running;
scope->before_request_();
}

void validate(const ClickHouseConnection& c) {
if (c.scheme != "http" && c.scheme != "https") {
throw ClickHouseError("clickhouse scheme must be \"http\" or \"https\"");
Expand Down Expand Up @@ -399,21 +482,46 @@ ClickHouseClient::~ClickHouseClient() = default;

std::vector<Row> ClickHouseClient::execute(
const std::string& query, const Params& params,
const std::map<std::string, std::string>& settings, int* attempts) const {
const std::map<std::string, std::string>& settings, int* attempts,
const AttemptSettings& per_attempt) const {
// First, so that what it does -- a lease renewal -- moves the deadline
// before this request reads it.
RequestDeadline::before_request();
// The outcome every scope in force records (RequestDeadline::last_timeout):
// a timeout's message, noted where one is thrown, or empty however else
// this call ends.
struct Outcome {
bool noted = false;
~Outcome() {
if (!noted) RequestDeadline::note_outcome("");
}
} outcome;
const auto timed_out = [&outcome](const std::string& error, bool sent) {
RequestDeadline::note_outcome(error);
outcome.noted = true;
return ClickHouseError(error, true, sent);
};
const std::string statement = substitute(query, params);
const bool read = is_read_statement(statement);

// Settings ride as URL parameters; the statement is the POST body
// (GET-with-query is evaluated as readonly — writes are refused). A
// caller's own wait_end_of_query wins over the default added here.
std::map<std::string, std::string> url_settings = settings;
if (read) url_settings.emplace("wait_end_of_query", "1");
std::string url = connection_.scheme + "://" + connection_.host + ":" +
std::to_string(connection_.port) + "/?";
for (const auto& [key, value] : url_settings) {
url += url_encode(key) + "=" + url_encode(value) + "&";
}
url.pop_back();
// caller's own wait_end_of_query wins over the default added here. Built
// per attempt: per_attempt's settings depend on the time each one has.
const auto url_for = [&](long request_ms) {
std::map<std::string, std::string> url_settings = settings;
if (per_attempt) {
per_attempt(static_cast<uint64_t>(request_ms), &url_settings);
}
if (read) url_settings.emplace("wait_end_of_query", "1");
std::string url = connection_.scheme + "://" + connection_.host + ":" +
std::to_string(connection_.port) + "/?";
for (const auto& [key, value] : url_settings) {
url += url_encode(key) + "=" + url_encode(value) + "&";
}
url.pop_back();
return url;
};

// Credentials as headers: ClickHouse reads X-ClickHouse-User/-Key, and
// unlike URL parameters or userinfo they do not end up in access logs or
Expand All @@ -431,7 +539,8 @@ std::vector<Row> ClickHouseClient::execute(
}
}

const auto perform = [&]() {
const auto perform = [&](const std::string& url, long connect_ms,
long request_ms) {
Attempt attempt;
CURL* curl = curl_easy_init();
if (curl == nullptr) throw ClickHouseError("libcurl init failed");
Expand Down Expand Up @@ -470,23 +579,78 @@ std::vector<Row> ClickHouseClient::execute(
}
// NOSIGNAL: timeouts must not use SIGALRM in a multi-threaded process.
curl_easy_setopt(curl, CURLOPT_NOSIGNAL, 1L);
curl_easy_setopt(curl, CURLOPT_CONNECTTIMEOUT_MS,
timeout_ms(connection_.timeouts.connect_s));
curl_easy_setopt(curl, CURLOPT_TIMEOUT_MS,
timeout_ms(connection_.timeouts.request_s));
curl_easy_setopt(curl, CURLOPT_CONNECTTIMEOUT_MS, connect_ms);
curl_easy_setopt(curl, CURLOPT_TIMEOUT_MS, request_ms);
attempt.code = curl_easy_perform(curl);
curl_easy_getinfo(curl, CURLINFO_RESPONSE_CODE, &attempt.status);
curl_easy_cleanup(curl);
attempt.detail = error_buffer;
return attempt;
};

// Whether any attempt may have reached the server (ClickHouseError::sent).
bool sent = false;
std::string last_error; // the previous attempt's, for a retry not sent
for (int number = 1;; ++number) {
const Attempt attempt = perform();
// The client's timeouts, cut to the RequestDeadline in force: read
// afresh for every attempt, since its owner may have moved it.
long connect_ms = timeout_ms(connection_.timeouts.connect_s);
long request_ms = timeout_ms(connection_.timeouts.request_s);
std::string bound;
const uint64_t deadline_ns = RequestDeadline::current(&bound);
bool by_deadline = false;
if (deadline_ns != 0) {
const uint64_t now_ns = steady_now_ns();
// Rounded down, so the attempt never outlasts the deadline.
const uint64_t left_ms =
deadline_ns > now_ns ? (deadline_ns - now_ns) / 1'000'000 : 0;
if (left_ms == 0) {
std::string error =
"Timeout: not sent, because " + bound + " had already passed";
if (number > 1) {
error += "; the attempt before failed with: " + last_error;
}
throw timed_out(error, sent);
}
if (left_ms < static_cast<uint64_t>(request_ms)) {
request_ms = static_cast<long>(left_ms);
by_deadline = true;
}
connect_ms = std::min(connect_ms, request_ms);
}

const Attempt attempt =
perform(url_for(request_ms), connect_ms, request_ms);
if (attempt.code == CURLE_OK && attempt.status == 200) {
if (attempts != nullptr) *attempts = number;
return parse_tsv(attempt.body);
}
if (attempt.code == CURLE_OK || !never_connected(attempt.code)) {
sent = true;
}
// Named in a timeout's message, so whoever reads it knows which knob to
// turn: the deadline that cut the attempt short, or the client's own
// timeouts.
const auto bounded_by = [&] {
return by_deadline
? " -- bounded by " + bound
: " -- bounded by the client's timeouts "
"(clickhouse_request_timeout_s = " +
seconds_text(connection_.timeouts.request_s) +
", clickhouse_connect_timeout_s = " +
seconds_text(connection_.timeouts.connect_s) + ")";
};
if (attempt.code == CURLE_OPERATION_TIMEDOUT) {
// Never retried (above).
std::string error = std::string("curl: ") +
curl_easy_strerror(attempt.code);
if (!attempt.detail.empty()) error += ": " + attempt.detail;
error += bounded_by();
if (number > 1) {
error += " (attempt " + std::to_string(number) + ")";
}
throw timed_out(error, sent);
}
std::string error;
bool retry = false;
if (attempt.code != CURLE_OK) {
Expand All @@ -503,20 +667,40 @@ std::vector<Row> ClickHouseClient::execute(
// say, and is retried as before. Reads only, either way.
int code = attempt.exception_code;
if (code < 0) code = body_exception_code(attempt.body);
if (server_time_limit(code)) {
// A timeout too, and like one never retried.
error += " -- the server's own time limit (max_execution_time, "
"lock_acquire_timeout or insert_quorum_timeout) ran out" +
bounded_by();
if (number > 1) {
error += " (attempt " + std::to_string(number) + ")";
}
throw timed_out(error, sent);
}
retry = read && attempt.status >= 500 && attempt.status < 600 &&
(code < 0 || transient_clickhouse_error(code));
}
if (!retry || number >= connection_.max_attempts) {
if (number > 1) {
error += " (after " + std::to_string(number) + " attempts)";
}
throw ClickHouseError(error);
throw ClickHouseError(error, false, sent);
}
// 100 ms, doubling, capped at 1 s: enough for a restarting server or a
// flapping connection, short beside the request timeout it adds to.
const int shift = std::min(number - 1, 4);
std::this_thread::sleep_for(std::chrono::milliseconds(
std::min(100 << shift, 1000)));
const uint64_t backoff_ms =
static_cast<uint64_t>(std::min(100 << shift, 1000));
if (deadline_ns != 0 &&
steady_now_ns() + backoff_ms * 1'000'000 >= deadline_ns) {
throw ClickHouseError(
error + " (after " + std::to_string(number) +
" attempt(s); not retried, because " + bound +
" leaves no time for another)",
false, sent);
}
last_error = error;
std::this_thread::sleep_for(std::chrono::milliseconds(backoff_ms));
}
}

Expand Down
Loading
Loading