Quellcode durchsuchen

feat(client): retry writes on DEADLINE_EXCEEDED/RESOURCE_EXHAUSTED/UNAVAILABLE (v1.7.0 T11)

Every write RPC (insert, updateIfVersion, upsert, remove, patch) now
retries transient transport failures up to writeRetries times (default
3) with exponential backoff + jitter (default 100ms -> 200ms -> 400ms,
capped at 5000ms, +/-25% jitter).

Retries ONLY on:
  DEADLINE_EXCEEDED   - server queued behind eviction/WAL load
  RESOURCE_EXHAUSTED  - admission control (T7) or gRPC concurrency cap (v1.6.2)
  UNAVAILABLE         - transient connection issue

All other status codes (INVALID_ARGUMENT, FAILED_PRECONDITION, etc.)
are final - no retry.

Version-conflict retry in update() is separate (it's an application-
level concern, not transport). update()'s internal updateIfVersion()
calls benefit from the transport retry; get() inside update() is not
retried (reads are never auto-retried).

This makes the Zoe-incident pattern non-fatal on the client: eviction
spike -> write temporarily returns RESOURCE_EXHAUSTED -> client sleeps
100-400ms -> retry -> success. shadowman-cpp's fire-and-forget
messages->save() now survives eviction bursts automatically.
fszontagh vor 3 Monaten
Ursprung
Commit
9a4651f2f0
2 geänderte Dateien mit 140 neuen und 25 gelöschten Zeilen
  1. 11 0
      client/include/smartbotic/database/client.hpp
  2. 129 25
      client/src/client.cpp

+ 11 - 0
client/include/smartbotic/database/client.hpp

@@ -22,6 +22,17 @@ public:
         std::string address = "localhost:9004";
         uint32_t timeoutMs = 5000;
         uint32_t maxRetries = 3;
+
+        // Write-retry configuration (v1.7.0 T11)
+        // Applied to insert/updateIfVersion/upsert/remove/patch on transient
+        // transport-level failures: DEADLINE_EXCEEDED, RESOURCE_EXHAUSTED,
+        // UNAVAILABLE. Other status codes are final and never retried.
+        // update()'s version-conflict retry (maxRetries above) is separate —
+        // its internal updateIfVersion() calls benefit from transport retry.
+        uint32_t writeRetries = 3;                  // max retry attempts for writes
+        uint32_t writeRetryBackoffMs = 100;         // initial backoff; exponential: 100, 200, 400...
+        uint32_t writeRetryMaxBackoffMs = 5000;     // cap on any single backoff
+        double writeRetryJitter = 0.25;             // ±25% jitter on each backoff
     };
 
     explicit Client(Config config);

+ 129 - 25
client/src/client.cpp

@@ -6,11 +6,40 @@
 #include <spdlog/spdlog.h>
 
 #include <atomic>
+#include <chrono>
+#include <cstdlib>
 #include <mutex>
 #include <thread>
 
 namespace smartbotic::database {
 
+namespace {
+// v1.7.0 T11 — transient transport errors that are safe to retry for idempotent
+// writes (insert with explicit ID, updateIfVersion, upsert, remove, patch).
+// All other codes (INVALID_ARGUMENT, FAILED_PRECONDITION, NOT_FOUND, ...) are
+// final — no retry.
+bool retryableStatus(const grpc::Status& s) {
+    switch (s.error_code()) {
+        case grpc::StatusCode::DEADLINE_EXCEEDED:  // server queued behind eviction/WAL
+        case grpc::StatusCode::RESOURCE_EXHAUSTED: // admission control or gRPC concurrency cap
+        case grpc::StatusCode::UNAVAILABLE:        // transient connection issue
+            return true;
+        default:
+            return false;
+    }
+}
+
+// Exponential backoff with jitter. attempt=0 => baseMs; attempt=1 => baseMs*2; ...
+uint32_t computeBackoffMs(uint32_t baseMs, uint32_t attempt, uint32_t capMs, double jitterPct) {
+    uint64_t exp = static_cast<uint64_t>(baseMs) << attempt;
+    if (exp > capMs) exp = capMs;
+    // ±jitterPct jitter; std::rand() is fine here — not a security sensitive RNG.
+    double jitter = 1.0 + ((double(std::rand()) / RAND_MAX) * 2.0 - 1.0) * jitterPct;
+    if (jitter < 0.1) jitter = 0.1;
+    return static_cast<uint32_t>(exp * jitter);
+}
+} // anonymous namespace
+
 // ===== PIMPL Implementation =====
 
 class Client::Impl {
@@ -85,12 +114,27 @@ public:
         }
 
         smartbotic::databasepb::InsertResponse response;
-        grpc::ClientContext context;
-        setDeadline(context);
-
-        auto status = stub_->Insert(&context, request, &response);
+        grpc::Status status;
+        for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
+            grpc::ClientContext context;
+            setDeadline(context);
+            response.Clear();
+            status = stub_->Insert(&context, request, &response);
+            if (status.ok() || !retryableStatus(status)) {
+                break;
+            }
+            if (attempt < config_.writeRetries) {
+                uint32_t backoffMs = computeBackoffMs(
+                    config_.writeRetryBackoffMs, attempt,
+                    config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
+                spdlog::warn("Client::insert {}; retrying in {}ms (attempt {}/{})",
+                             status.error_message(), backoffMs,
+                             attempt + 1, config_.writeRetries);
+                std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
+            }
+        }
         if (!status.ok()) {
-            spdlog::error("Client::insert failed: {}", status.error_message());
+            spdlog::error("Client::insert failed after retries: {}", status.error_message());
             throw std::runtime_error(status.error_message());
         }
 
@@ -180,12 +224,27 @@ public:
         }
 
         smartbotic::databasepb::UpdateResponse response;
-        grpc::ClientContext context;
-        setDeadline(context);
-
-        auto status = stub_->Update(&context, request, &response);
+        grpc::Status status;
+        for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
+            grpc::ClientContext context;
+            setDeadline(context);
+            response.Clear();
+            status = stub_->Update(&context, request, &response);
+            if (status.ok() || !retryableStatus(status)) {
+                break;
+            }
+            if (attempt < config_.writeRetries) {
+                uint32_t backoffMs = computeBackoffMs(
+                    config_.writeRetryBackoffMs, attempt,
+                    config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
+                spdlog::warn("Client::updateIfVersion {}; retrying in {}ms (attempt {}/{})",
+                             status.error_message(), backoffMs,
+                             attempt + 1, config_.writeRetries);
+                std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
+            }
+        }
         if (!status.ok()) {
-            spdlog::error("Client::updateIfVersion failed: {}", status.error_message());
+            spdlog::error("Client::updateIfVersion failed after retries: {}", status.error_message());
             return false;
         }
 
@@ -203,12 +262,27 @@ public:
         }
 
         smartbotic::databasepb::PatchDocumentResponse response;
-        grpc::ClientContext context;
-        setDeadline(context);
-
-        auto status = stub_->PatchDocument(&context, request, &response);
+        grpc::Status status;
+        for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
+            grpc::ClientContext context;
+            setDeadline(context);
+            response.Clear();
+            status = stub_->PatchDocument(&context, request, &response);
+            if (status.ok() || !retryableStatus(status)) {
+                break;
+            }
+            if (attempt < config_.writeRetries) {
+                uint32_t backoffMs = computeBackoffMs(
+                    config_.writeRetryBackoffMs, attempt,
+                    config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
+                spdlog::warn("Client::patch {}; retrying in {}ms (attempt {}/{})",
+                             status.error_message(), backoffMs,
+                             attempt + 1, config_.writeRetries);
+                std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
+            }
+        }
         if (!status.ok()) {
-            spdlog::error("Client::patch failed: {}", status.error_message());
+            spdlog::error("Client::patch failed after retries: {}", status.error_message());
             return 0;
         }
 
@@ -237,12 +311,27 @@ public:
         }
 
         smartbotic::databasepb::UpsertResponse response;
-        grpc::ClientContext context;
-        setDeadline(context);
-
-        auto status = stub_->Upsert(&context, request, &response);
+        grpc::Status status;
+        for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
+            grpc::ClientContext context;
+            setDeadline(context);
+            response.Clear();
+            status = stub_->Upsert(&context, request, &response);
+            if (status.ok() || !retryableStatus(status)) {
+                break;
+            }
+            if (attempt < config_.writeRetries) {
+                uint32_t backoffMs = computeBackoffMs(
+                    config_.writeRetryBackoffMs, attempt,
+                    config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
+                spdlog::warn("Client::upsert {}; retrying in {}ms (attempt {}/{})",
+                             status.error_message(), backoffMs,
+                             attempt + 1, config_.writeRetries);
+                std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
+            }
+        }
         if (!status.ok()) {
-            spdlog::error("Client::upsert failed: {}", status.error_message());
+            spdlog::error("Client::upsert failed after retries: {}", status.error_message());
             throw std::runtime_error(status.error_message());
         }
 
@@ -255,12 +344,27 @@ public:
         request.set_id(id);
 
         smartbotic::databasepb::DeleteResponse response;
-        grpc::ClientContext context;
-        setDeadline(context);
-
-        auto status = stub_->Delete(&context, request, &response);
+        grpc::Status status;
+        for (uint32_t attempt = 0; attempt <= config_.writeRetries; ++attempt) {
+            grpc::ClientContext context;
+            setDeadline(context);
+            response.Clear();
+            status = stub_->Delete(&context, request, &response);
+            if (status.ok() || !retryableStatus(status)) {
+                break;
+            }
+            if (attempt < config_.writeRetries) {
+                uint32_t backoffMs = computeBackoffMs(
+                    config_.writeRetryBackoffMs, attempt,
+                    config_.writeRetryMaxBackoffMs, config_.writeRetryJitter);
+                spdlog::warn("Client::remove {}; retrying in {}ms (attempt {}/{})",
+                             status.error_message(), backoffMs,
+                             attempt + 1, config_.writeRetries);
+                std::this_thread::sleep_for(std::chrono::milliseconds(backoffMs));
+            }
+        }
         if (!status.ok()) {
-            spdlog::error("Client::remove failed: {}", status.error_message());
+            spdlog::error("Client::remove failed after retries: {}", status.error_message());
             return false;
         }